import { sql, type ColumnType, type JSONColumnType, type Selectable } from 'kysely' import { customAlphabet } from 'nanoid' import type { Address, Hex } from 'viem' import * as Campaigns from '../../internal/rewards/Campaigns.js' import type * as Db from '../Db.js' import type * as db_Schema from '../Schema.js' const id = customAlphabet('0123456789abcdefghijklmnopqrstuvwxyz', 24) /** Columns of the `reward_runs` table. */ export type Table = Omit & { /** Bigint chain id, read as a string from pg and written as a number. */ chainId: ColumnType /** Configuration snapshot consumed by the run. */ config: JSONColumnType /** Resumable calculation and delivery evidence, compacted after completion. */ evidence: JSONColumnType /** Latest cumulative statement. */ statement: JSONColumnType } /** A stored reward execution. */ export type Record = Omit & { /** Completed calculation evidence, or null before calculation finishes. */ evidence: Campaigns.RunEvidence | null } /** Compact execution fields used by the campaign administration surface. */ export type ExecutionSummary = Pick< Record, | 'endsAt' | 'error' | 'evidence' | 'fundedAssets' | 'id' | 'liability' | 'mintedEarnShares' | 'phase' | 'rootVersion' | 'startsAfter' | 'statementHash' | 'updatedAt' > /** Compact delivered-run fields used for public rates and source fallback. */ export type DeliveredSummary = Pick /** Compact delivered-run fields keyed by EarnVault. */ export type DeliveredVaultSummary = DeliveredSummary & Pick /** Creates or returns the unique contiguous run starting after one boundary. */ export async function ensure(db: Db.Db, options: ensure.Options): Promise { return db.transaction(async (tx) => { const campaign = await tx.kysely .selectFrom('reward_campaigns') .select(['config', 'deliveredThrough', 'pendingEffectiveAt']) .where('chainId', '=', String(options.chainId)) .where('vaultAddress', '=', normalize(options.vaultAddress)) .forUpdate() .executeTakeFirstOrThrow() if ( campaign.deliveredThrough !== String(options.startsAfter) || canonical(campaign.config) !== canonical(options.config) || (campaign.pendingEffectiveAt !== null && options.endsAt > Number(campaign.pendingEffectiveAt)) ) throw new Error('reward campaign changed before its run snapshot was created') const now = new Date().toISOString() const inserted = await tx.kysely .insertInto('reward_runs') .values({ chainId: options.chainId, config: JSON.stringify(options.config), createdAt: now, endsAt: String(options.endsAt), error: null, evidence: null, fence: '0', fundedAssets: null, id: `rrn_${id()}`, leaseExpiresAt: null, liability: null, mintedEarnShares: null, phase: 'indexing', root: null, rootVersion: null, startsAfter: String(options.startsAfter), statement: null, statementHash: null, updatedAt: now, vaultAddress: normalize(options.vaultAddress), }) .onConflict((oc) => oc.columns(['chainId', 'vaultAddress', 'startsAfter']).doNothing()) .returningAll() .executeTakeFirst() const row = inserted ?? (await tx.kysely .selectFrom('reward_runs') .selectAll() .where('chainId', '=', String(options.chainId)) .where('vaultAddress', '=', normalize(options.vaultAddress)) .where('startsAfter', '=', String(options.startsAfter)) .executeTakeFirstOrThrow()) const record = toRecord(row) // `startsAfter` is the durable execution identity. Once a run exists its // closing boundary is immutable: a retry may happen after the wall clock // has exposed more intervals, but it must finish this original range // before the campaign can plan later work. if (canonical(record.config) !== canonical(options.config)) throw new Error('existing reward run conflicts with the requested interval snapshot') return record }) } export declare namespace ensure { /** Identity and snapshot for one contiguous run. */ type Options = { /** Chain containing the vault. */ chainId: number /** Configuration snapshot. */ config: Campaigns.Config /** Exclusive final boundary. */ endsAt: number /** Previously delivered boundary. */ startsAfter: number /** EarnVault address. */ vaultAddress: Address } } /** Acquires an expired or unowned run lease and increments its fencing token. */ export async function acquire(db: Db.Db, options: acquire.Options): Promise { const now = new Date() const row = await db.kysely .updateTable('reward_runs') .set({ fence: sql`fence + 1`, leaseExpiresAt: new Date(now.getTime() + options.leaseMilliseconds).toISOString(), updatedAt: now.toISOString(), }) .where('id', '=', options.id) .where((eb) => eb.or([eb('leaseExpiresAt', 'is', null), eb('leaseExpiresAt', '<=', now.toISOString())]), ) .where('phase', 'not in', ['delivered', 'failed']) .returningAll() .executeTakeFirst() return row ? toRecord(row) : undefined } export declare namespace acquire { /** Lease selector and lifetime. */ type Options = { /** Run id. */ id: string /** Lease duration. */ leaseMilliseconds: number } } /** Extends an active run lease without changing its fencing token. */ export async function renew(db: Db.Db, options: renew.Options): Promise { const now = new Date().toISOString() return Boolean( await db.kysely .updateTable('reward_runs') .set({ leaseExpiresAt: options.leaseExpiresAt, updatedAt: new Date().toISOString(), }) .where('id', '=', options.id) .where('fence', '=', options.fence) .where('leaseExpiresAt', 'is not', null) .where('leaseExpiresAt', '>', now) .returning('id') .executeTakeFirst(), ) } export declare namespace renew { /** Fenced lease renewal. */ type Options = { /** Expected lease fencing token. */ fence: string /** Run id. */ id: string /** Replacement lease expiry. */ leaseExpiresAt: string } } /** Persists resumable calculation evidence under the current run fence and phase. */ export async function writeCheckpoint( db: Db.Db, options: writeCheckpoint.Options, ): Promise { const row = await db.kysely .updateTable('reward_runs') .set({ evidence: JSON.stringify(options.evidence), updatedAt: new Date().toISOString(), }) .where('id', '=', options.id) .where('fence', '=', options.fence) .where('phase', '=', 'calculating') .returningAll() .executeTakeFirst() return row ? toRecord(row) : undefined } export declare namespace writeCheckpoint { /** Fenced calculation checkpoint. */ type Options = { /** Resumable calculation evidence. */ evidence: Campaigns.RunCheckpointEvidence /** Expected lease fencing token. */ fence: string /** Run id. */ id: string } } /** Reads resumable calculation evidence without treating completed evidence as a checkpoint. */ export async function getCheckpoint( db: Db.Db, options: getCheckpoint.Options, ): Promise { const row = await db.kysely .selectFrom('reward_runs') .select('evidence') .where('id', '=', options.id) .executeTakeFirst() return Campaigns.isRunCheckpointEvidence(row?.evidence) ? row.evidence : undefined } export declare namespace getCheckpoint { /** Run containing resumable evidence. */ type Options = { /** Run id. */ id: string } } /** Applies one fenced phase transition. */ export async function transition( db: Db.Db, options: transition.Options, ): Promise { const { evidence, statement, ...patch } = options.patch const row = await db.kysely .updateTable('reward_runs') .set({ ...patch, ...(evidence !== undefined ? { evidence: evidence ? JSON.stringify(evidence) : null } : {}), ...(statement !== undefined ? { statement: statement ? JSON.stringify(statement) : null } : {}), phase: options.phase, updatedAt: new Date().toISOString(), }) .where('id', '=', options.id) .where('fence', '=', String(options.fence)) .returningAll() .executeTakeFirst() return row ? toRecord(row) : undefined } /** Records an execution error and releases the fenced lease without regressing its phase. */ export async function releaseWithError( db: Db.Db, options: releaseWithError.Options, ): Promise { const row = await db.kysely .updateTable('reward_runs') .set({ error: options.error, leaseExpiresAt: null, updatedAt: new Date().toISOString(), }) .where('id', '=', options.id) .where('fence', '=', String(options.fence)) .returningAll() .executeTakeFirst() return row ? toRecord(row) : undefined } export declare namespace releaseWithError { /** Fenced lease and diagnostic. */ type Options = { /** Actionable error text. */ error: string /** Expected lease fencing token. */ fence: string /** Run id. */ id: string } } export declare namespace transition { /** Fenced phase transition. */ type Options = { /** Expected lease fencing token. */ fence: string /** Run id. */ id: string /** Next execution phase. */ phase: Record['phase'] /** Durable outputs produced by the completed phase. */ patch: Partial<{ error: string | null evidence: Campaigns.RunEvidence | null fundedAssets: string | null liability: string | null mintedEarnShares: string | null root: Hex | null rootVersion: string | null statement: Campaigns.Statement | null statementHash: Hex | null }> } } /** Reads one run. */ export async function get(db: Db.Db, runId: string): Promise { const row = await db.kysely .selectFrom('reward_runs') .selectAll() .where('id', '=', runId) .executeTakeFirst() return row ? toRecord(row) : undefined } /** Reads the newest run without loading its potentially large Merkle statement. */ export async function latestSummary( db: Db.Db, options: latestSummary.Options, ): Promise { const row = await db.kysely .selectFrom('reward_runs') .select([ 'endsAt', 'error', 'evidence', 'fundedAssets', 'id', 'liability', 'mintedEarnShares', 'phase', 'rootVersion', 'startsAfter', 'statementHash', 'updatedAt', ]) .where('chainId', '=', String(options.chainId)) .where('vaultAddress', '=', normalize(options.vaultAddress)) .orderBy('startsAfter', 'desc') .executeTakeFirst() if (!row) return undefined return { ...row, evidence: Campaigns.isRunCheckpointEvidence(row.evidence) ? null : row.evidence, } } export declare namespace latestSummary { /** Campaign identity. */ type Options = { /** Chain containing the vault. */ chainId: number /** EarnVault address. */ vaultAddress: string } } /** Reads the latest confirmed statement for one campaign. */ export async function latestStatement( db: Db.Db, options: latestStatement.Options, ): Promise { const row = await db.kysely .selectFrom('reward_runs') .selectAll() .where('chainId', '=', String(options.chainId)) .where('vaultAddress', '=', normalize(options.vaultAddress)) .where('statement', 'is not', null) .where((eb) => eb.or([ eb('phase', 'in', ['paying', 'delivered']), eb.and([ eb('phase', '=', 'statement'), sql`EXISTS ( SELECT 1 FROM reward_transaction_attempts attempt WHERE attempt.run_id = reward_runs.id AND attempt.state = 'confirmed' AND attempt.intent ->> 'operation' IN ('publish', 'settle') )`, ]), ]), ) .orderBy('endsAt', 'desc') .executeTakeFirst() if (!row) return undefined const record = toRecord(row) return record.phase === 'statement' && record.rootVersion !== null ? { ...record, rootVersion: (BigInt(record.rootVersion) + 1n).toString() } : record } /** Reads the latest delivered run for rate fallback and monitoring. */ export async function latestDelivered( db: Db.Db, options: latestDelivered.Options, ): Promise { return (await db.kysely .selectFrom('reward_runs') .select(['config', 'endsAt', 'evidence']) .where('chainId', '=', String(options.chainId)) .where('vaultAddress', '=', normalize(options.vaultAddress)) .where('phase', '=', 'delivered') .orderBy('endsAt', 'desc') .executeTakeFirst()) as DeliveredSummary | undefined } /** Reads the latest delivered run for each bounded EarnVault. */ export async function listLatestDelivered( db: Db.Db, options: listLatestDelivered.Options, ): Promise { const vaultAddresses = [...new Set(options.vaultAddresses.map(normalize))] if (vaultAddresses.length === 0) return [] return (await db.kysely .selectFrom('reward_runs') .distinctOn('vaultAddress') .select(['config', 'endsAt', 'evidence', 'vaultAddress']) .where('chainId', '=', String(options.chainId)) .where('vaultAddress', 'in', vaultAddresses) .where('phase', '=', 'delivered') .orderBy('vaultAddress', 'asc') .orderBy('endsAt', 'desc') .execute()) as DeliveredVaultSummary[] } export declare namespace listLatestDelivered { /** Bounded delivered-run selector. */ type Options = { /** Chain containing every vault. */ chainId: number /** EarnVault addresses. */ vaultAddresses: readonly string[] } } /** Clears superseded proof-bearing statements while retaining their accounting evidence. */ export async function pruneStatements(db: Db.Db, options: pruneStatements.Options): Promise { await db.kysely .updateTable('reward_runs') .set({ statement: null }) .where('chainId', '=', String(options.chainId)) .where('vaultAddress', '=', normalize(options.vaultAddress)) .where('id', '!=', options.keepRunId) .where('statement', 'is not', null) .execute() } export declare namespace pruneStatements { /** Campaign and run whose full statement remains live. */ type Options = { /** Chain containing the vault. */ chainId: number /** Run retaining the latest full statement. */ keepRunId: string /** EarnVault address. */ vaultAddress: string } } export declare namespace latestDelivered { /** Campaign identity. */ type Options = { /** Chain containing the vault. */ chainId: number /** EarnVault address. */ vaultAddress: string } } export declare namespace latestStatement { /** Campaign identity. */ type Options = { /** Chain containing the vault. */ chainId: number /** EarnVault address. */ vaultAddress: string } } function normalize(address: string): Address { return address.toLowerCase() as Address } function toRecord(row: Selectable): Record { return { ...row, chainId: Number(row.chainId) } as Record } function canonical(value: unknown): string { if (Array.isArray(value)) return `[${value.map(canonical).join(',')}]` if (value && typeof value === 'object') return `{${Object.entries(value) .sort(([left], [right]) => left.localeCompare(right)) .map(([key, entry]) => `${JSON.stringify(key)}:${canonical(entry)}`) .join(',')}}` return JSON.stringify(value) }