apps/keeper/src/jobs/claims.ts

Collects creator fees into the vault. Shown whole, as it was in the repository when this site was built. Line numbers link: add #L12 to the address.

apps/keeper/src/jobs/claims.ts159 lines
1// claims: collect each coin's pump.fun creator fees into the vault and record them in the ledger.2// Each coin's fees wait in its own creator vault (the pump PDA of its sharing config), so one transaction per coin3// (PumpSwap transfer_creator_fees_to_pump when graduated + distribute_creator_fees) pays the vault exactly that4// coin's share; the vault's balance change in that transaction is the claim. recordFeeClaim writes fee_claims and5// moves the coin's reward index in one database transaction (SPEC §3).6// Anyone may call distribute_creator_fees (a creator claiming their cut on pump.fun pays the vault too), so every7// run also reads the vault's new transactions and records any distribute that paid it; (tx, coin) is unique, so the8// keeper's own claims are never counted twice.9// Intents (kind 'claim'): built → sent (signature stored before the transaction leaves) → confirmed | failed | expired.10import { and, eq, inArray, sql } from 'drizzle-orm';11import { intents } from '@solary/db';12import { event, getSetting, lamportsToCents, latestSolUsd, setSetting, sol, solUsdAt, usd, type JobCtx, type Stats } from '../ctx';13import { createIntent, holdReason, intentState, openIntents, setIntent } from '../intents';14import { recordFeeClaim, type RecordedFeeClaim } from '../ledger/claims';15import { words } from '../ledger/words';16import { vaultPendingLamports, type DistributeOutcome } from '../solana/adapter';17import { coinRef, handles, keeperCoins, settingBig, settingNum, type CoinRow } from './common';18 19const KIND = 'claim';20 21async function announce(ctx: JobCtx, r: RecordedFeeClaim, collectionName: string | null, sig: string): Promise<void> {22  await event(ctx.db, {23    coin: r.coin,24    collection: r.collection,25    action: 'fee_claim',26    message: words.claimed(r.symbol, r.lamports + r.creatorLamports, r.holders, collectionName),27    amounts: { lamports: r.lamports.toString(), creator_lamports: r.creatorLamports.toString(), holders_lamports: r.holders.toString(), buyback_lamports: r.buyback.toString(), platform_lamports: r.platform.toString(), usd_cents: r.usdCents },28    sig,29  });30}31 32/** Record every coin a confirmed distribute transaction paid. Returns the claims written (0 when already recorded). */33async function record(ctx: JobCtx, o: DistributeOutcome, names: Map<string, CoinRow>, intentId: number | null): Promise<number> {34  const price = await solUsdAt(ctx.db, o.at);35  let n = 0;36  for (const c of o.coins) {37    if (!names.has(c.mint)) continue;38    const r = await recordFeeClaim(ctx.db, { coin: c.mint, tx: o.signature, lamports: c.lamports, creatorLamports: c.creatorLamports, solUsd: price, at: o.at, intentId });39    if (r) {40      n++;41      await announce(ctx, r, names.get(c.mint)!.collectionName, o.signature);42    }43  }44  if (intentId && n === 0) await setIntent(ctx, intentId, 'confirmed');45  return n;46}47 48async function reconcile(ctx: JobCtx, names: Map<string, CoinRow>): Promise<{ confirmed: number; failed: number; expired: number }> {49  let confirmed = 0;50  let failed = 0;51  let expired = 0;52  for (const i of await openIntents(ctx, KIND)) {53    const st = await intentState(ctx, i);54    if (st === 'confirmed') {55      const o = await ctx.sol.distributeOutcome(i.sig);56      if (o) {57        await record(ctx, o, names, i.id);58        confirmed++;59      } else if (ctx.sol.mode === 'fake') await setIntent(ctx, i.id, 'expired'); // the fake chain forgets on restart60    } else if (st === 'failed' || st === 'expired') {61      await setIntent(ctx, i.id, st);62      if (st === 'failed') failed++;63      else expired++;64    }65  }66  return { confirmed, failed, expired };67}68 69/** Distributes someone else sent (or ours that were never recorded), from the vault's transaction history. */70async function ingestHistory(ctx: JobCtx, names: Map<string, CoinRow>): Promise<number> {71  if (ctx.sol.mode === 'fake') return 0;72  const cursorKey = `keeper.${ctx.sol.mode}.vault_cursor`;73  const cursor = await getSetting<string | null>(ctx.db, cursorKey, null);74  const known = new Set(75    (76      (await ctx.db.execute(sql`SELECT sig FROM intents WHERE sig IS NOT NULL AND kind IN ('payout', 'sweep', 'buyback', 'burn') AND created_at > now() - interval '30 days'`)) as unknown as Array<{ sig: string }>77    ).map((r) => r.sig),78  );79  const { outcomes, newest } = await ctx.sol.vaultDistributes(cursor, (s) => known.has(s));80  let n = 0;81  const unmatched = await getSetting<string[]>(ctx.db, 'claims.unmatched', []);82  for (const o of outcomes) {83    n += await record(ctx, o, names, null);84    for (const c of o.coins) if (!names.has(c.mint)) unmatched.push(`${o.signature}:${c.mint}`);85  }86  // retry distributes of coins the index job had not added yet87  const still: string[] = [];88  for (const u of [...new Set(unmatched)].slice(-200)) {89    const [sig, mint] = u.split(':') as [string, string];90    if (!names.has(mint)) {91      still.push(u);92      continue;93    }94    const o = await ctx.sol.distributeOutcome(sig);95    if (o) n += await record(ctx, { ...o, coins: o.coins.filter((c) => c.mint === mint) }, names, null);96  }97  await setSetting(ctx.db, 'claims.unmatched', still);98  if (newest && newest !== cursor) await setSetting(ctx.db, cursorKey, newest);99  return n;100}101 102export async function claimsJob(ctx: JobCtx): Promise<Stats> {103  const { db, sol: adapter } = ctx;104  const all = await keeperCoins(ctx);105  const names = new Map(all.map((c) => [c.mint, c]));106  const rec = await reconcile(ctx, names);107  const ingested = await ingestHistory(ctx, names);108 109  // claim coins that pay somebody: active, paired with a collection (or the official coin)110  const list = all.filter((c) => c.status === 'active' && (c.collection || c.official) && handles(ctx, c.mint));111  if (!list.length) return { coins: 0, ingested, reconciled: rec.confirmed, reconcile_failed: rec.failed, expired: rec.expired };112  const price = await latestSolUsd(db);113  const minLamports = await settingBig(ctx, 'claims.min_lamports', ctx.env.fast ? 5_000_000n : 20_000_000n);114  const perRun = await settingNum(ctx, 'claims.per_run', 20);115  const priority = await settingBig(ctx, 'tx.priority_micro_lamports', 50_000n);116 117  const states = await adapter.coinStates(list.map(coinRef));118  const open = new Set(119    (await db.select({ payload: intents.payload }).from(intents).where(and(eq(intents.kind, KIND), inArray(intents.status, ['built', 'sent'])))).map((r) => String(r.payload.mint)),120  );121  const due = list122    .map((c) => ({ c, s: states.get(c.mint) }))123    .filter((x): x is { c: CoinRow; s: NonNullable<typeof x.s> } => !!x.s && x.s.exists && x.s.usesSharingConfig && x.s.vaultBps > 0)124    .map((x) => ({ ...x, ours: vaultPendingLamports(x.s) }))125    .filter((x) => x.ours >= minLamports && !open.has(x.c.mint))126    .sort((a, b) => (a.ours > b.ours ? -1 : 1))127    .slice(0, perRun);128 129  const hold = holdReason(ctx);130  let claimed = 0;131  let would = 0;132  let failed = 0;133  for (const { c, s, ours } of due) {134    const worth = price ? usd(lamportsToCents(ours, price.usd)) : 'unpriced';135    const prepared = await adapter.prepareClaim(s, priority);136    if (!prepared.ok) {137      failed++;138      ctx.log(`$${c.symbol}: claim simulation failed: ${prepared.error}`);139      continue;140    }141    if (hold) {142      // dry runs and pauses record nothing: the log line is the whole result143      would++;144      ctx.log(`${hold}: would claim ${sol(ours)} (${worth}) for $${c.symbol}${s.graduated ? ' (PumpSwap + curve)' : ''}; simulated OK, ${prepared.unitsConsumed} CU`);145      continue;146    }147    const intentId = await createIntent(ctx, KIND, { mint: c.mint, symbol: c.symbol, expected_lamports: ours, graduated: s.graduated });148    const res = await adapter.send(prepared, (sig, lvbh) => setIntent(ctx, intentId, 'sent', { sig, payload: { lastValidBlockHeight: lvbh } }));149    if (res.status === 'confirmed') {150      const o = await adapter.distributeOutcome(res.signature);151      if (o && (await record(ctx, o, names, intentId)) > 0) claimed++;152    } else if (res.status === 'failed' || res.status === 'expired') {153      failed++;154      await setIntent(ctx, intentId, res.status, { payload: { error: res.error ?? res.status } });155      ctx.log(`$${c.symbol}: claim ${res.status}: ${res.error}`);156    } else ctx.log(`$${c.symbol}: claim ${res.signature.slice(0, 12)}… not confirmed yet; the next run resolves it`);157  }158  return { coins: list.length, due: due.length, claimed, would_claim: would, failed, ingested, reconciled: rec.confirmed, reconcile_failed: rec.failed, expired: rec.expired };159}