apps/keeper/src/jobs/payouts.ts

Automatic payouts and holder claims. 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/payouts.ts328 lines
1// payouts: send NFT holders what their NFTs are owed, in SOL, from the vault (SPEC §3).2//  1. reconcile: payout transactions left `sent` by an earlier run → confirmed (settle) | failed | expired (stay owed)3//  2. resume: transfers inserted but never sent (a crash, a pause, the per-pass limit) go out now4//  3. claims: a holder asked on /earnings#wallet → all of that wallet's owed across coins, ignoring the automatic threshold5//     but not the 0.0001 SOL minimum after gas6//  4. automatic runs per coin: owed grouped by current owner, planTransfers, then autoPayoutState (holders keep ≥ 90%7//     after gas, at least payout.min_lamports, at least payout.min_interval_s since the coin's last payout). Wallets8//     under payout.auto_min_net_lamports wait for a later run (or claim).9// Only wallet owners (owner_is_wallet) are paid; burnt NFTs never; an NFT in an open transfer is left out of new runs.10// Transfers are packed TRANSFERS_PER_TX per transaction; each transaction is one intent (kind 'payout') whose11// payload names its transfers, so a crash at any point is resolved by the next run.12import { and, asc, eq, gte, inArray, lt } from 'drizzle-orm';13import { claimRequests, keeperEvents, payoutTransfers } from '@solary/db';14import { formatSol } from '@solary/core';15import { event, type JobCtx, type Stats } from '../ctx';16import { createIntent, holdReason, intentState, openIntents } from '../intents';17import {18  ConservationError,19  autoDecision,20  confirmTransfers,21  failTransfers,22  insertRun,23  loadOwed,24  loadWalletOwed,25  markSent,26  planRun,27  planWalletClaim,28  type FinishedRun,29} from '../ledger/payouts';30import { words } from '../ledger/words';31import { handles, keeperCoins, settingBig, settingNum, vaultObligations, type CoinRow } from './common';32 33const KIND = 'payout';34 35interface Group {36  transferIds: number[];37  transfers: Array<{ to: string; lamports: bigint }>;38  claimRequestIds: number[];39}40 41interface Pass {42  ctx: JobCtx;43  priority: number;44  coins: Map<string, CoinRow>;45  txBudget: number;46  sent: number;47  confirmed: number;48  failed: number;49  paidLamports: bigint;50}51 52async function announce(p: Pass, finished: FinishedRun[]): Promise<void> {53  for (const f of finished) {54    if (f.kind !== 'auto' || f.status !== 'done') continue;55    const c = p.coins.get(f.coin);56    await event(p.ctx.db, {57      coin: f.coin,58      collection: c?.collection ?? null,59      action: 'payout',60      message: words.paid(f.net, f.recipients, c?.collectionName ?? 'NFT', c?.symbol ?? f.coin.slice(0, 6), f.failed),61      amounts: { run: f.runId, recipients: f.recipients, nfts: f.nfts, gross_lamports: f.gross.toString(), net_lamports: f.net.toString(), failed: f.failed },62    });63  }64}65 66/** One transaction: simulate, record, send, then settle or release. */67async function sendGroup(p: Pass, g: Group): Promise<{ status: 'confirmed' | 'failed' | 'pending'; sig?: string }> {68  const { ctx } = p;69  p.txBudget--;70  const prepared = await ctx.sol.preparePayout(g.transfers, BigInt(p.priority));71  const intentId = await createIntent(ctx, KIND, { transfers: g.transferIds, claims: g.claimRequestIds, recipients: g.transfers.length, lamports: prepared.lamports });72  if (!prepared.ok) {73    p.failed += g.transferIds.length;74    ctx.log(`payout simulation failed: ${prepared.error}`);75    await announce(p, await failTransfers(ctx.db, { transferIds: g.transferIds, status: 'failed', intentId, claimRequestIds: g.claimRequestIds, note: prepared.error }));76    return { status: 'failed' };77  }78  p.sent++;79  const res = await ctx.sol.send(prepared, (sig, lvbh) => markSent(ctx.db, { transferIds: g.transferIds, sig, intentId, lastValidBlockHeight: lvbh }));80  if (res.status === 'confirmed') {81    const c = await confirmTransfers(ctx.db, { transferIds: g.transferIds, sig: res.signature, intentId, claimRequestIds: g.claimRequestIds });82    p.confirmed += c.confirmed;83    p.paidLamports += c.net;84    await announce(p, c.finished);85    return { status: 'confirmed', sig: res.signature };86  }87  if (res.status === 'failed' || res.status === 'expired') {88    p.failed += g.transferIds.length;89    ctx.log(`payout ${res.status}: ${res.error ?? ''}`);90    await announce(p, await failTransfers(ctx.db, { transferIds: g.transferIds, status: res.status, intentId, claimRequestIds: g.claimRequestIds, note: res.error }));91    return { status: 'failed' };92  }93  ctx.log(`payout ${res.signature.slice(0, 12)}… not confirmed yet; the next run resolves it`);94  return { status: 'pending', sig: res.signature };95}96 97async function reconcile(p: Pass): Promise<{ confirmed: number; failed: number }> {98  let confirmed = 0;99  let failed = 0;100  for (const i of await openIntents(p.ctx, KIND)) {101    const ids = (i.payload.transfers as number[] | undefined) ?? [];102    const claims = (i.payload.claims as number[] | undefined) ?? [];103    const st = await intentState(p.ctx, i);104    if (st === 'confirmed') {105      const c = await confirmTransfers(p.ctx.db, { transferIds: ids, sig: i.sig, intentId: i.id, claimRequestIds: claims });106      confirmed += c.confirmed;107      await announce(p, c.finished);108    } else if (st === 'failed' || st === 'expired') {109      failed += ids.length;110      await announce(p, await failTransfers(p.ctx.db, { transferIds: ids, status: st, intentId: i.id, claimRequestIds: claims }));111    }112  }113  return { confirmed, failed };114}115 116/** Transfers inserted but never broadcast: send them (or expire them after a day). */117async function resume(p: Pass): Promise<number> {118  const { db } = p.ctx;119  const built = await db120    .select({ id: payoutTransfers.id, run: payoutTransfers.run, coin: payoutTransfers.coin, owner: payoutTransfers.owner, kind: payoutTransfers.kind, batch: payoutTransfers.batch, net: payoutTransfers.netLamports, createdAt: payoutTransfers.createdAt })121    .from(payoutTransfers)122    .where(eq(payoutTransfers.status, 'built'))123    .orderBy(asc(payoutTransfers.id))124    .limit(2000);125  const mine = built.filter((b) => handles(p.ctx, b.coin));126  const stale = mine.filter((b) => Date.now() - b.createdAt.getTime() > 86_400_000);127  if (stale.length) await announce(p, await failTransfers(db, { transferIds: stale.map((s) => s.id), status: 'expired', note: 'not sent within a day' }));128  const fresh = mine.filter((b) => !stale.includes(b));129  const groups = new Map<string, typeof fresh>();130  for (const b of fresh) {131    const k = b.kind === 'claim' ? `claim:${b.owner}:${b.batch}` : `auto:${b.run}:${b.batch}`;132    groups.set(k, [...(groups.get(k) ?? []), b]);133  }134  let n = 0;135  for (const [k, list] of groups) {136    if (p.txBudget <= 0) break;137    let claimIds: number[] = [];138    if (k.startsWith('claim:')) {139      claimIds = (140        await db.select({ id: claimRequests.id }).from(claimRequests).where(and(eq(claimRequests.owner, list[0]!.owner), eq(claimRequests.status, 'queued')))141      ).map((r) => r.id);142    }143    await sendGroup(p, { transferIds: list.map((b) => b.id), transfers: list.map((b) => ({ to: b.owner, lamports: b.net })), claimRequestIds: claimIds });144    n++;145  }146  return n;147}148 149/** Can the vault cover `lamports` on top of a reserve? Alerts (at most every 6 h) when it cannot. */150async function vaultCovers(p: Pass, lamports: bigint): Promise<boolean> {151  const reserve = await settingBig(p.ctx, 'vault.reserve_lamports', 10_000_000n);152  const balance = await p.ctx.sol.vaultBalance();153  if (balance >= lamports + reserve) return true;154  const [recent] = await p.ctx.db155    .select({ id: keeperEvents.id })156    .from(keeperEvents)157    .where(and(eq(keeperEvents.action, 'vault_short'), gte(keeperEvents.createdAt, new Date(Date.now() - 6 * 3_600_000))))158    .limit(1);159  if (!recent) {160    const owed = await vaultObligations(p.ctx);161    await event(p.ctx.db, { action: 'vault_short', level: 'alert', message: words.vaultShort(balance, owed.total), amounts: { balance: balance.toString(), owed: owed.total.toString() } });162  }163  return false;164}165 166function groupsOf(transferIds: number[], planned: ReadonlyArray<{ owner: string; net: bigint; batch: number }>, claimRequestIds: number[] = []): Group[] {167  const by = new Map<number, Group>();168  planned.forEach((t, i) => {169    const g = by.get(t.batch) ?? { transferIds: [], transfers: [], claimRequestIds };170    g.transferIds.push(transferIds[i]!);171    g.transfers.push({ to: t.owner, lamports: t.net });172    by.set(t.batch, g);173  });174  return [...by.entries()].sort((a, b) => a[0] - b[0]).map(([, g]) => g);175}176 177async function claimPass(p: Pass, hold: string | null): Promise<{ paid: number; nothing: number; waiting: number }> {178  const { ctx } = p;179  const minNet = await settingBig(ctx, 'payout.min_net_lamports', 100_000n);180  const queued = await ctx.db.select().from(claimRequests).where(eq(claimRequests.status, 'queued')).orderBy(asc(claimRequests.createdAt)).limit(20);181  let paid = 0;182  let nothing = 0;183  let waiting = 0;184  for (const q of queued) {185    if (p.txBudget <= 0) break;186    const [inflight] = await ctx.db187      .select({ id: payoutTransfers.id })188      .from(payoutTransfers)189      .where(and(eq(payoutTransfers.owner, q.owner), inArray(payoutTransfers.status, ['built', 'sent'])))190      .limit(1);191    if (inflight) {192      waiting++;193      continue;194    }195    const owed = (await loadWalletOwed(ctx.db, q.owner)).filter((o) => handles(ctx, o.coin));196    const exists = (await ctx.sol.walletsExist([q.owner])).has(q.owner);197    const plan = planWalletClaim(owed, { priorityMicroLamports: p.priority, minNetLamports: minNet, walletExists: exists });198    if (!plan.transfers.length) {199      nothing++;200      await ctx.db201        .update(claimRequests)202        .set({ status: 'nothing', note: owed.length ? 'owed is under the minimum after gas' : 'nothing owed', doneAt: new Date() })203        .where(eq(claimRequests.id, q.id));204      ctx.log(words.claimNothing(q.owner));205      continue;206    }207    if (hold) {208      ctx.log(`${hold}: would pay ${formatSol(plan.net)} to ${q.owner.slice(0, 6)}… on request (${plan.transfers.length} ${plan.transfers.length === 1 ? 'coin' : 'coins'})`);209      continue;210    }211    if (!(await vaultCovers(p, plan.gross))) break;212    const ids: number[] = [];213    const planned: Array<{ owner: string; net: bigint; batch: number }> = [];214    for (const t of plan.transfers) {215      const r = await insertRun(ctx.db, { coin: t.coin, kind: 'claim', index: t.index, transfers: [{ owner: q.owner, gross: t.gross, gas: t.gas, net: t.net, items: t.items, batch: t.batch }] });216      ids.push(r.transferIds[0]!);217      planned.push({ owner: q.owner, net: t.net, batch: t.batch });218    }219    let ok = 0n;220    let sig: string | undefined;221    for (const g of groupsOf(ids, planned, [q.id])) {222      const st = await sendGroup(p, g);223      if (st.status === 'confirmed') {224        ok += g.transfers.reduce((a, t) => a + t.lamports, 0n);225        sig ??= st.sig;226      }227    }228    if (ok > 0n) {229      paid++;230      await event(ctx.db, { action: 'payout_claim', message: words.paidClaim(ok, q.owner, plan.transfers.length), amounts: { owner: q.owner, net_lamports: ok.toString(), coins: plan.transfers.length }, sig });231    }232  }233  return { paid, nothing, waiting };234}235 236async function autoPass(p: Pass, hold: string | null): Promise<{ runs: number; would: number; notReady: number; refused: number }> {237  const { ctx } = p;238  const minLamports = await settingBig(ctx, 'payout.min_lamports', 10_000_000n);239  const minIntervalS = await settingNum(ctx, 'payout.min_interval_s', ctx.env.fast ? 60 : 3600);240  const autoMinNet = await settingBig(ctx, 'payout.auto_min_net_lamports', 1_000_000n);241  let runs = 0;242  let would = 0;243  let notReady = 0;244  let refused = 0;245  // alert coins (fees stopped reaching the vault) still pay out what their holders already earned; paused ones wait246  const list = [...p.coins.values()].filter((c) => c.status !== 'paused' && c.collection && c.collectionStatus === 'listed' && c.holdersBps > 0);247  // longest-waiting first248  list.sort((a, b) => (a.lastPayoutAt?.getTime() ?? 0) - (b.lastPayoutAt?.getTime() ?? 0));249  for (const c of list) {250    if (p.txBudget <= 0) break;251    if (c.lastPayoutAt && Date.now() - c.lastPayoutAt.getTime() < minIntervalS * 1000) continue;252    const { index, lines } = await loadOwed(ctx.db, { mint: c.mint, collection: c.collection!, accIndex: c.accIndex });253    const draft = planRun(lines, index, { priorityMicroLamports: p.priority, minNetLamports: autoMinNet });254    if (!draft.transfers.length) {255      notReady++;256      continue;257    }258    const exists = await ctx.sol.walletsExist(draft.transfers.map((t) => t.owner));259    const plan = planRun(lines, index, { priorityMicroLamports: p.priority, minNetLamports: autoMinNet, walletExists: (o) => exists.has(o) });260    const d = autoDecision(plan, { priorityMicroLamports: p.priority, minLamports, minIntervalS, lastPayoutAt: c.lastPayoutAt, now: Date.now() });261    if (!d.ready) {262      notReady++;263      continue;264    }265    if (hold) {266      would++;267      ctx.log(`${hold}: would pay ${formatSol(plan.net)} to ${plan.transfers.length} ${c.collectionName} holders for $${c.symbol} (${Math.ceil(plan.transfers.length / 18)} transactions)`);268      continue;269    }270    if (!(await vaultCovers(p, plan.gross))) break;271    let ins: { runId: number; transferIds: number[] };272    try {273      ins = await insertRun(ctx.db, { coin: c.mint, kind: 'auto', index, transfers: plan.transfers });274    } catch (e) {275      if (!(e instanceof ConservationError)) throw e;276      refused++;277      await event(ctx.db, { coin: c.mint, collection: c.collection, action: 'payout_refused', level: 'alert', message: `Payout for $${c.symbol} refused: it would pay holders more than the coin's fees gave them` });278      continue;279    }280    runs++;281    for (const g of groupsOf(ins.transferIds, plan.transfers)) {282      if (p.txBudget <= 0) break; // the rest is resumed next run283      await sendGroup(p, g);284    }285  }286  return { runs, would, notReady, refused };287}288 289export async function payoutsJob(ctx: JobCtx): Promise<Stats> {290  const all = await keeperCoins(ctx);291  const p: Pass = {292    ctx,293    priority: await settingNum(ctx, 'payout.priority_micro_lamports', 20_000),294    coins: new Map(all.map((c) => [c.mint, c])),295    txBudget: await settingNum(ctx, 'payout.max_txs_per_pass', ctx.env.fast ? 12 : 60),296    sent: 0,297    confirmed: 0,298    failed: 0,299    paidLamports: 0n,300  };301  const rec = await reconcile(p);302  const hold = holdReason(ctx);303  const resumed = hold ? 0 : await resume(p);304  const claims = await claimPass(p, hold);305  const auto = await autoPass(p, hold);306  // expire stale claim requests nobody can pay (wallet owns nothing any more)307  await ctx.db308    .update(claimRequests)309    .set({ status: 'nothing', note: 'expired after 7 days', doneAt: new Date() })310    .where(and(eq(claimRequests.status, 'queued'), lt(claimRequests.createdAt, new Date(Date.now() - 7 * 86_400_000))));311  return {312    hold,313    reconciled: rec.confirmed,314    reconcile_failed: rec.failed,315    resumed,316    claims_paid: claims.paid,317    claims_nothing: claims.nothing,318    claims_waiting: claims.waiting,319    runs: auto.runs,320    would_run: auto.would,321    not_ready: auto.notReady,322    refused: auto.refused,323    txs: p.sent,324    transfers_confirmed: p.confirmed,325    transfers_failed: p.failed,326    paid_lamports: p.paidLamports.toString(),327  };328}