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}