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}