apps/keeper/src/ledger/payouts.ts
Plans, records and settles payout runs. Shown whole, as it was in the repository when this site was built. Line numbers link: add #L12 to the address.
apps/keeper/src/ledger/payouts.ts447 lines
1// Payout runs (SPEC §3). Per coin: snapshot index = acc_index; group what each NFT is owed by its current owner;2// planTransfers (gas shared equally, net = gross − gas; skip under minNet, skip non-existent wallets below rent);3// insert payout_runs + one payout_transfers row per wallet (items = [[nft, lamports]]), packed TRANSFERS_PER_TX per4// transaction (`batch`). On confirm of a transaction: its transfers → confirmed, every item upserts nft_settlements5// (settled_index = run.index, paid_lamports += lamports), coins.holders_paid_lamports += gross. On failure: transfers →6// failed/expired, settlements untouched (still owed). An NFT in an open transfer is excluded from other runs.7// Conservation: Σ confirmed + in-flight gross never exceeds Σ fee_claims.holders_lamports of the coin (checked when a8// run is inserted). The keeper's payouts job and the demo seed both go through these functions.9import { and, eq, inArray, sql } from 'drizzle-orm';10import { claimRequests, coins, intents, payoutRuns, payoutTransfers, type Db } from '@solary/db';11import { AUTO_PAYOUT_KEEP_BPS, RENT_EXEMPT_WALLET, TRANSFERS_PER_TX, autoPayoutState, owed, planTransfers, type AutoPayoutState } from '@solary/core';12 13export type Item = [nft: string, lamports: bigint];14 15export interface NftOwedRow {16 id: string;17 owner: string | null;18 ownerIsWallet: boolean;19 burnt: boolean;20 multiplier: number;21 /** the coin's index at this NFT's last payout (0 = never paid) */22 settledIndex: bigint;23}24 25export interface OwedLine {26 owner: string;27 gross: bigint;28 items: Item[];29}30 31/** Pure: what each wallet is owed for one coin at `index`. Burnt, ownerless, program-owned and excluded NFTs wait. */32export function owedLines(rows: readonly NftOwedRow[], index: bigint, exclude: ReadonlySet<string> = new Set()): OwedLine[] {33 const by = new Map<string, OwedLine>();34 for (const r of rows) {35 if (r.burnt || !r.owner || !r.ownerIsWallet || exclude.has(r.id)) continue;36 const o = owed(r.multiplier, index, r.settledIndex);37 if (o <= 0n) continue;38 const line = by.get(r.owner) ?? { owner: r.owner, gross: 0n, items: [] };39 line.gross += o;40 line.items.push([r.id, o]);41 by.set(r.owner, line);42 }43 return [...by.values()];44}45 46export interface PlannedTransfer {47 owner: string;48 gross: bigint;49 gas: bigint;50 net: bigint;51 items: Item[];52 /** which transaction of the run carries it */53 batch: number;54}55 56export interface RunPlan {57 index: bigint;58 transfers: PlannedTransfer[];59 skipped: OwedLine[];60 gross: bigint;61 gas: bigint;62 net: bigint;63 nfts: number;64}65 66export interface PlanOpts {67 priorityMicroLamports: number;68 minNetLamports: bigint;69 walletExists?: (owner: string) => boolean;70}71 72/** Pure: a run from owed lines. Transfers keep planTransfers' order (largest first) and are packed per transaction. */73export function planRun(lines: readonly OwedLine[], index: bigint, opts: PlanOpts): RunPlan {74 const byOwner = new Map(lines.map((l) => [l.owner, l]));75 const p = planTransfers(76 lines.map((l) => ({ owner: l.owner, gross: l.gross })),77 { priorityMicroLamports: opts.priorityMicroLamports, minNetLamports: opts.minNetLamports, walletExists: opts.walletExists },78 );79 const transfers = p.transfers.map((t, i) => ({ ...t, items: byOwner.get(t.owner)!.items, batch: Math.floor(i / TRANSFERS_PER_TX) }));80 const sum = (f: (t: PlannedTransfer) => bigint) => transfers.reduce((a, t) => a + f(t), 0n);81 const skipped = p.skipped.map((s) => byOwner.get(s.owner)!);82 return { index, transfers, skipped, gross: sum((t) => t.gross), gas: sum((t) => t.gas), net: sum((t) => t.net), nfts: transfers.reduce((a, t) => a + t.items.length, 0) };83}84 85export interface AutoOpts {86 priorityMicroLamports: number;87 /** payout.min_lamports */88 minLamports: bigint;89 /** payout.min_interval_s */90 minIntervalS: number;91 lastPayoutAt: Date | null;92 now: number;93}94 95/** Pure: should a coin's automatic run go out now? (holders keep ≥ 90% after gas, the floor, the interval) */96export function autoDecision(plan: RunPlan, o: AutoOpts): { ready: boolean; reason: string; state: AutoPayoutState } {97 const state = autoPayoutState(plan.gross, plan.transfers.length, { priorityMicroLamports: o.priorityMicroLamports, minLamports: o.minLamports });98 const since = o.lastPayoutAt ? (o.now - o.lastPayoutAt.getTime()) / 1000 : Infinity;99 if (plan.transfers.length === 0) return { ready: false, reason: 'no wallet is owed enough yet', state };100 if (since < o.minIntervalS) return { ready: false, reason: `last payout ${Math.round(since)}s ago (minimum ${o.minIntervalS}s)`, state };101 if (!state.ready) return { ready: false, reason: `pending ${plan.gross} of ${state.threshold} lamports`, state };102 // planTransfers already shares gas; this is the SPEC's 90%-after-gas rule stated once more for the record103 if (plan.net * 10_000n < plan.gross * BigInt(AUTO_PAYOUT_KEEP_BPS)) return { ready: false, reason: 'gas would take more than 10%', state };104 return { ready: true, reason: 'ready', state };105}106 107export interface WalletCoinOwed {108 coin: string;109 index: bigint;110 gross: bigint;111 items: Item[];112}113 114export interface WalletClaimTransfer extends WalletCoinOwed {115 gas: bigint;116 net: bigint;117 batch: number;118}119 120/**121 * Pure: a holder's claim across coins, one transfer per coin packed TRANSFERS_PER_TX per transaction, gas shared122 * over the transfers. Coins under minNet stay owed. A wallet that doesn't exist yet needs the whole claim to reach123 * the rent-exempt minimum.124 */125export function planWalletClaim(126 lines: readonly WalletCoinOwed[],127 opts: { priorityMicroLamports: number; minNetLamports: bigint; walletExists: boolean },128): { transfers: WalletClaimTransfer[]; skipped: WalletCoinOwed[]; net: bigint; gross: bigint } {129 const by = new Map(lines.map((l) => [l.coin, l]));130 const p = planTransfers(131 lines.map((l) => ({ owner: l.coin, gross: l.gross })),132 { priorityMicroLamports: opts.priorityMicroLamports, minNetLamports: opts.minNetLamports },133 );134 let transfers = p.transfers.map((t, i) => ({ ...by.get(t.owner)!, gas: t.gas, net: t.net, batch: Math.floor(i / TRANSFERS_PER_TX) }));135 let skipped = p.skipped.map((s) => by.get(s.owner)!);136 const net = transfers.reduce((a, t) => a + t.net, 0n);137 if (!opts.walletExists && net < BigInt(RENT_EXEMPT_WALLET)) {138 skipped = [...lines];139 transfers = [];140 }141 return { transfers, skipped, net: transfers.reduce((a, t) => a + t.net, 0n), gross: transfers.reduce((a, t) => a + t.gross, 0n) };142}143 144// ---------------------------------------------------------------- database145 146type Tx = Parameters<Parameters<Db['transaction']>[0]>[0];147type Exec = Db | Tx;148 149const rowsOf = <T>(r: unknown): T[] => r as T[];150 151/** NFT ids in open (built or sent) transfers of a coin: they stay out of new runs until resolved. */152async function openItems(db: Exec, coin: string): Promise<Set<string>> {153 const rows = await db154 .select({ items: payoutTransfers.items })155 .from(payoutTransfers)156 .where(and(eq(payoutTransfers.coin, coin), inArray(payoutTransfers.status, ['built', 'sent'])));157 const out = new Set<string>();158 for (const r of rows) for (const [nft] of r.items) out.add(nft);159 return out;160}161 162/** What every wallet is owed for a coin right now (open transfers excluded). */163export async function loadOwed(db: Db, coin: { mint: string; collection: string; accIndex: string | bigint }): Promise<{ index: bigint; lines: OwedLine[] }> {164 const index = BigInt(coin.accIndex);165 const exclude = await openItems(db, coin.mint);166 const rows = rowsOf<{ id: string; owner: string; multiplier: number; settled: string | null }>(167 await db.execute(sql`168 SELECT n.id, n.owner, n.multiplier, s.settled_index AS settled169 FROM nfts n LEFT JOIN nft_settlements s ON s.coin = ${coin.mint} AND s.nft = n.id170 WHERE n.collection = ${coin.collection} AND NOT n.burnt AND n.owner IS NOT NULL AND n.owner_is_wallet171 AND COALESCE(s.settled_index, 0) < ${index.toString()}::numeric`),172 );173 const lines = owedLines(174 rows.map((r) => ({ id: r.id, owner: r.owner, ownerIsWallet: true, burnt: false, multiplier: Number(r.multiplier), settledIndex: BigInt(r.settled ?? '0') })),175 index,176 exclude,177 );178 return { index, lines };179}180 181/** What one wallet is owed across every coin paired with its NFTs (open transfers excluded). */182export async function loadWalletOwed(db: Db, owner: string): Promise<Array<WalletCoinOwed & { symbol: string; collection: string }>> {183 const rows = rowsOf<{ id: string; multiplier: number; mint: string; symbol: string; collection: string; acc_index: string; settled: string | null }>(184 await db.execute(sql`185 SELECT n.id, n.multiplier, c.mint, c.symbol, c.collection, c.acc_index, s.settled_index AS settled186 FROM nfts n187 JOIN coins c ON c.collection = n.collection188 LEFT JOIN nft_settlements s ON s.coin = c.mint AND s.nft = n.id189 WHERE n.owner = ${owner} AND n.owner_is_wallet AND NOT n.burnt AND c.acc_index > COALESCE(s.settled_index, 0)`),190 );191 const byCoin = new Map<string, { symbol: string; collection: string; index: bigint; rows: NftOwedRow[] }>();192 for (const r of rows) {193 const e = byCoin.get(r.mint) ?? { symbol: r.symbol, collection: r.collection, index: BigInt(r.acc_index), rows: [] };194 e.rows.push({ id: r.id, owner, ownerIsWallet: true, burnt: false, multiplier: Number(r.multiplier), settledIndex: BigInt(r.settled ?? '0') });195 byCoin.set(r.mint, e);196 }197 const out: Array<WalletCoinOwed & { symbol: string; collection: string }> = [];198 for (const [mint, e] of byCoin) {199 const exclude = await openItems(db, mint);200 const [line] = owedLines(e.rows, e.index, exclude);201 if (line && line.gross > 0n) out.push({ coin: mint, symbol: e.symbol, collection: e.collection, index: e.index, gross: line.gross, items: line.items });202 }203 return out;204}205 206export class ConservationError extends Error {}207 208const itemsJson = (items: Item[]): Array<[string, string]> => items.map(([n, l]) => [n, l.toString()]);209 210/**211 * Insert a run and its transfers (status built) for one coin. Refuses (ConservationError) when the coin's confirmed212 * plus in-flight plus this run's gross would exceed what its holders ever received.213 */214export async function insertRun(215 db: Db,216 r: { coin: string; kind: 'auto' | 'claim'; index: bigint; transfers: ReadonlyArray<Omit<PlannedTransfer, 'batch'> & { batch: number }>; at?: Date; solUsd?: string | null },217): Promise<{ runId: number; transferIds: number[] }> {218 const at = r.at ?? new Date();219 const gross = r.transfers.reduce((a, t) => a + t.gross, 0n);220 return db.transaction(async (t) => {221 const [c] = await t222 .select({ holders: coins.holdersLamports, paid: coins.holdersPaidLamports })223 .from(coins)224 .where(eq(coins.mint, r.coin))225 .for('update');226 if (!c) throw new Error(`unknown coin ${r.coin}`);227 const [open] = rowsOf<{ g: string | null }>(228 await t.execute(sql`SELECT SUM(gross_lamports)::text AS g FROM payout_transfers WHERE coin = ${r.coin} AND status IN ('built', 'sent')`),229 );230 const inflight = BigInt(open?.g ?? '0');231 if (c.paid + inflight + gross > c.holders) {232 throw new ConservationError(`run for ${r.coin} would pay ${c.paid + inflight + gross} of ${c.holders} lamports ever received for holders`);233 }234 const [run] = await t235 .insert(payoutRuns)236 .values({237 coin: r.coin,238 kind: r.kind,239 status: 'building',240 index: r.index.toString(),241 recipients: r.transfers.length,242 nfts: r.transfers.reduce((a, x) => a + x.items.length, 0),243 grossLamports: gross,244 gasLamports: r.transfers.reduce((a, x) => a + x.gas, 0n),245 netLamports: r.transfers.reduce((a, x) => a + x.net, 0n),246 createdAt: at,247 })248 .returning({ id: payoutRuns.id });249 const ids: number[] = [];250 for (let i = 0; i < r.transfers.length; i += 500) {251 const chunk = r.transfers.slice(i, i + 500);252 const got = await t253 .insert(payoutTransfers)254 .values(255 chunk.map((x) => ({256 run: run!.id,257 coin: r.coin,258 owner: x.owner,259 kind: r.kind,260 nfts: x.items.length,261 grossLamports: x.gross,262 gasLamports: x.gas,263 netLamports: x.net,264 items: itemsJson(x.items),265 batch: x.batch,266 status: 'built' as const,267 solUsd: r.solUsd ?? null,268 createdAt: at,269 })),270 )271 .returning({ id: payoutTransfers.id });272 ids.push(...got.map((g) => g.id));273 }274 return { runId: run!.id, transferIds: ids };275 });276}277 278/** The transaction carrying these transfers is signed and about to leave: record its signature first. */279export async function markSent(db: Db, m: { transferIds: number[]; sig: string; intentId: number; lastValidBlockHeight: number }): Promise<void> {280 await db.transaction(async (t) => {281 await t.update(payoutTransfers).set({ status: 'sent', tx: m.sig }).where(inArray(payoutTransfers.id, m.transferIds));282 const runs = await t.selectDistinct({ run: payoutTransfers.run }).from(payoutTransfers).where(inArray(payoutTransfers.id, m.transferIds));283 await t284 .update(payoutRuns)285 .set({ status: 'sending' })286 .where(and(inArray(payoutRuns.id, runs.map((r) => r.run)), eq(payoutRuns.status, 'building')));287 await t288 .update(intents)289 .set({ status: 'sent', sig: m.sig, updatedAt: new Date(), payload: sql`${intents.payload} || ${JSON.stringify({ lastValidBlockHeight: m.lastValidBlockHeight })}::jsonb` })290 .where(eq(intents.id, m.intentId));291 });292}293 294export interface FinishedRun {295 runId: number;296 coin: string;297 kind: 'auto' | 'claim';298 status: 'done' | 'failed';299 recipients: number;300 nfts: number;301 gross: bigint;302 net: bigint;303 failed: number;304}305 306/** Close runs with nothing left in flight; their totals become what was actually confirmed. */307async function finishRuns(t: Tx, runIds: number[], at: Date): Promise<FinishedRun[]> {308 if (!runIds.length) return [];309 const rows = rowsOf<{ run: string; coin: string; kind: 'auto' | 'claim'; open: string; ok: string; bad: string; nfts: string; gross: string; gas: string; net: string }>(310 await t.execute(sql`311 SELECT run, MIN(coin) AS coin, MIN(kind::text) AS kind,312 COUNT(*) FILTER (WHERE status IN ('built', 'sent')) AS open,313 COUNT(*) FILTER (WHERE status = 'confirmed') AS ok,314 COUNT(*) FILTER (WHERE status IN ('failed', 'expired')) AS bad,315 COALESCE(SUM(nfts) FILTER (WHERE status = 'confirmed'), 0) AS nfts,316 COALESCE(SUM(gross_lamports) FILTER (WHERE status = 'confirmed'), 0) AS gross,317 COALESCE(SUM(gas_lamports) FILTER (WHERE status = 'confirmed'), 0) AS gas,318 COALESCE(SUM(net_lamports) FILTER (WHERE status = 'confirmed'), 0) AS net319 FROM payout_transfers WHERE run IN ${sql`(${sql.join(runIds.map((id) => sql`${id}`), sql`, `)})`}320 GROUP BY run`),321 );322 const out: FinishedRun[] = [];323 for (const r of rows) {324 if (Number(r.open) > 0) continue;325 const status = Number(r.ok) > 0 ? 'done' : 'failed';326 const upd = await t327 .update(payoutRuns)328 .set({329 status,330 finishedAt: at,331 recipients: Number(r.ok),332 nfts: Number(r.nfts),333 grossLamports: BigInt(r.gross),334 gasLamports: BigInt(r.gas),335 netLamports: BigInt(r.net),336 })337 .where(and(eq(payoutRuns.id, Number(r.run)), inArray(payoutRuns.status, ['building', 'sending'])))338 .returning({ id: payoutRuns.id });339 if (upd.length)340 out.push({ runId: Number(r.run), coin: r.coin, kind: r.kind, status, recipients: Number(r.ok), nfts: Number(r.nfts), gross: BigInt(r.gross), net: BigInt(r.net), failed: Number(r.bad) });341 }342 return out;343}344 345export interface ConfirmResult {346 confirmed: number;347 gross: bigint;348 net: bigint;349 finished: FinishedRun[];350}351 352/**353 * A payout transaction confirmed: its transfers → confirmed, settlements up to each run's index, holders_paid and354 * last_payout_at on the coins, claim requests paid, runs closed. Idempotent: transfers already confirmed are skipped.355 */356export async function confirmTransfers(357 db: Db,358 c: { transferIds: number[]; sig: string; at?: Date; intentId?: number | null; claimRequestIds?: number[] },359): Promise<ConfirmResult> {360 const at = c.at ?? new Date();361 return db.transaction(async (t) => {362 const rows = await t363 .update(payoutTransfers)364 .set({ status: 'confirmed', tx: c.sig, confirmedAt: at })365 .where(and(inArray(payoutTransfers.id, c.transferIds), inArray(payoutTransfers.status, ['built', 'sent'])))366 .returning({ run: payoutTransfers.run, coin: payoutTransfers.coin, owner: payoutTransfers.owner, gross: payoutTransfers.grossLamports, net: payoutTransfers.netLamports, items: payoutTransfers.items });367 if (c.intentId) await t.update(intents).set({ status: 'confirmed', sig: c.sig, updatedAt: new Date() }).where(eq(intents.id, c.intentId));368 if (!rows.length) return { confirmed: 0, gross: 0n, net: 0n, finished: [] };369 const runIds = [...new Set(rows.map((r) => r.run))];370 const runIndex = new Map((await t.select({ id: payoutRuns.id, index: payoutRuns.index }).from(payoutRuns).where(inArray(payoutRuns.id, runIds))).map((r) => [r.id, r.index]));371 372 // settlements, aggregated per (coin, nft)373 const settle = new Map<string, { coin: string; nft: string; index: bigint; lamports: bigint }>();374 for (const r of rows) {375 const index = BigInt(runIndex.get(r.run)!);376 for (const [nft, l] of r.items) {377 const k = `${r.coin}:${nft}`;378 const s = settle.get(k) ?? { coin: r.coin, nft, index: 0n, lamports: 0n };379 s.index = index > s.index ? index : s.index;380 s.lamports += BigInt(l);381 settle.set(k, s);382 }383 }384 const list = [...settle.values()];385 for (let i = 0; i < list.length; i += 5000) {386 // one JSON parameter instead of five per row: a wallet with hundreds of NFTs settles them all at once387 const rows = JSON.stringify(list.slice(i, i + 5000).map((s) => ({ c: s.coin, n: s.nft, i: s.index.toString(), l: s.lamports.toString() })));388 await t.execute(sql`389 INSERT INTO nft_settlements (coin, nft, settled_index, paid_lamports, updated_at)390 SELECT x.c, x.n, x.i::numeric, x.l::bigint, ${at.toISOString()}::timestamptz FROM jsonb_to_recordset(${rows}::jsonb) AS x(c text, n text, i text, l text)391 ON CONFLICT (coin, nft) DO UPDATE SET392 settled_index = GREATEST(nft_settlements.settled_index, excluded.settled_index),393 paid_lamports = nft_settlements.paid_lamports + excluded.paid_lamports,394 updated_at = excluded.updated_at`);395 }396 397 // coins398 const perCoin = new Map<string, bigint>();399 for (const r of rows) perCoin.set(r.coin, (perCoin.get(r.coin) ?? 0n) + r.gross);400 for (const [coin, gross] of perCoin) {401 await t402 .update(coins)403 .set({404 holdersPaidLamports: sql`${coins.holdersPaidLamports} + ${gross.toString()}::bigint`,405 lastPayoutAt: sql`GREATEST(${coins.lastPayoutAt}, ${at.toISOString()}::timestamptz)`,406 })407 .where(eq(coins.mint, coin));408 }409 410 // claim requests paid by this transaction411 if (c.claimRequestIds?.length) {412 const reqs = await t.select({ id: claimRequests.id, owner: claimRequests.owner }).from(claimRequests).where(inArray(claimRequests.id, c.claimRequestIds));413 for (const q of reqs) {414 const net = rows.filter((r) => r.owner === q.owner).reduce((a, r) => a + r.net, 0n);415 await t416 .update(claimRequests)417 .set({ status: 'paid', lamports: sql`${claimRequests.lamports} + ${net.toString()}::bigint`, doneAt: at })418 .where(eq(claimRequests.id, q.id));419 }420 }421 422 const finished = await finishRuns(t, runIds, at);423 return { confirmed: rows.length, gross: rows.reduce((a, r) => a + r.gross, 0n), net: rows.reduce((a, r) => a + r.net, 0n), finished };424 });425}426 427/** A payout transaction failed or expired: its transfers stay owed (settlements untouched). */428export async function failTransfers(429 db: Db,430 f: { transferIds: number[]; status: 'failed' | 'expired'; at?: Date; intentId?: number | null; claimRequestIds?: number[]; note?: string },431): Promise<FinishedRun[]> {432 const at = f.at ?? new Date();433 return db.transaction(async (t) => {434 const rows = await t435 .update(payoutTransfers)436 .set({ status: f.status })437 .where(and(inArray(payoutTransfers.id, f.transferIds), inArray(payoutTransfers.status, ['built', 'sent'])))438 .returning({ run: payoutTransfers.run });439 if (f.intentId) await t.update(intents).set({ status: f.status, updatedAt: new Date() }).where(eq(intents.id, f.intentId));440 if (f.claimRequestIds?.length)441 await t442 .update(claimRequests)443 .set({ status: 'failed', note: (f.note ?? `payout ${f.status}`).slice(0, 300), doneAt: at })444 .where(and(inArray(claimRequests.id, f.claimRequestIds), eq(claimRequests.status, 'queued')));445 return finishRuns(t, [...new Set(rows.map((r) => r.run))], at);446 });447}