import { sql, type Selectable } from 'kysely' import * as Fees from '../../internal/Fees.js' import type * as Db from '../Db.js' import type * as db_Schema from '../Schema.js' /** Columns of the `routes_transfers` table. */ export type Table = db_Schema.RoutesTransfer /** A stored route transfer. */ export type Record = Selectable /** Time retained for late source registration after a subsidized quote expires. */ export const subsidyCommitmentRecoveryTtlMs = 86_400_000 /** Returns whether an awaiting-source subsidy commitment is past its registration window. */ export function isSubsidyCommitmentExpired(record: Record): boolean { const amount = record.subsidyCommitmentAmount if (!amount || BigInt(amount.baseUnits) <= 0n) return false const recoverableAfter = new Date(Date.now() - subsidyCommitmentRecoveryTtlMs).toISOString() return record.quoteExpiresAt <= recoverableAfter } /** * Ownership filter derived from the API-key principal, never from the * request. A `projectId` narrows reads to that project; organization-attributed * keys omit it and see every project in their organization and environment. */ export type Owner = { /** Key environment. */ environment: Record['environment'] /** Owning organization id (`org_…`). */ orgId: string /** Attributed project id, for project-attributed keys. */ projectId?: string | undefined } /** Aggregates active transfer states for bounded operational metrics. */ export function summarizeStates(db: Db.Db, options: summarizeStates.Options) { const reason = sql`status_reason->>'code'` return db.kysely .selectFrom('routes_transfers') .select(({ fn }) => [ 'environment', 'method', 'providerId', 'status', fn.countAll().as('count'), sql`min(coalesce(status_updated_at, updated_at))`.as('oldestStatusUpdatedAt'), reason.as('reason'), ]) .where('status', 'in', options.statuses) .groupBy(['environment', 'method', 'providerId', 'status']) .groupBy(reason) .execute() } export declare namespace summarizeStates { /** Active-state aggregation options. */ type Options = { /** Lifecycle statuses still owned by platform processing or recovery. */ statuses: readonly Record['status'][] } } /** Inserts a transfer row. Callers commit the matching event in the same transaction. */ export function insert(db: Db.Db, record: Record): Promise { return db.kysely .insertInto('routes_transfers') .values(record) .returningAll() .executeTakeFirstOrThrow() } /** Gets a transfer by id without ownership scoping (domain and reconciler use). */ export function get(db: Db.Db, id: string): Promise { return db.kysely .selectFrom('routes_transfers') .selectAll() .where('id', '=', id) .executeTakeFirst() } /** Lists processing transfers for one provider, oldest first. */ export function listProcessing(db: Db.Db, options: listProcessing.Options): Promise { let query = db.kysely .selectFrom('routes_transfers') .selectAll() .where('providerId', '=', options.providerId) .where('status', '=', 'processing') if (options.cursor !== undefined) query = query.where('id', '>', options.cursor) return query.orderBy('id', 'asc').limit(options.limit).execute() } export declare namespace listProcessing { /** Provider filter for one recovery scan. */ type Options = { /** Exclusive lower bound from the prior recovery page. */ cursor?: string | undefined /** Maximum transfers returned in one recovery page. */ limit: number /** Route provider id. */ providerId: string } } /** Gets and locks a transfer for one atomic domain mutation. */ export function getForUpdate(db: Db.Db, id: string): Promise { return db.kysely .selectFrom('routes_transfers') .selectAll() .where('id', '=', id) .forUpdate() .executeTakeFirst() } /** Persists the first verified provider-delivery boundary without a public version change. */ export function recordProviderDelivery( db: Db.Db, options: recordProviderDelivery.Options, ): Promise { return db.kysely .updateTable('routes_transfers') .set({ providerDeliveredAt: options.deliveredAt }) .where('id', '=', options.id) .where('providerDeliveredAt', 'is', null) .where('status', '=', 'processing') .returningAll() .executeTakeFirst() } export declare namespace recordProviderDelivery { /** Provider delivery boundary write. */ type Options = { /** First verified provider-delivery time. */ deliveredAt: string /** Route transfer id (`rtr_…`). */ id: string } } /** Gets an owner-visible transfer by id; `undefined` for absent or foreign rows. */ export function getOwned(db: Db.Db, owner: Owner, id: string): Promise { let query = db.kysely .selectFrom('routes_transfers') .selectAll() .where('id', '=', id) .where('orgId', '=', owner.orgId) .where('environment', '=', owner.environment) if (owner.projectId !== undefined) query = query.where('projectId', '=', owner.projectId) return query.executeTakeFirst() } /** * Lists an owner's transfers newest-first. Ids embed a timestamp, so lexical * key order is chronological and the cursor is the last returned id. */ export function listByOwner( db: Db.Db, owner: Owner, options: listByOwner.Options, ): Promise { let query = db.kysely .selectFrom('routes_transfers') .selectAll() .where('orgId', '=', owner.orgId) .where('environment', '=', owner.environment) if (owner.projectId !== undefined) query = query.where('projectId', '=', owner.projectId) if (options.status !== undefined) query = query.where('status', '=', options.status) if (options.cursor !== undefined) query = query.where('id', '<', options.cursor) return query.orderBy('id', 'desc').limit(options.limit).execute() } export declare namespace listByOwner { /** Paging and filter options. */ type Options = { /** Exclusive lower bound: the last id of the previous page. */ cursor?: string | undefined /** Maximum rows to return (fetch one extra to detect a further page). */ limit: number /** Restricts to one lifecycle status. */ status?: Record['status'] | undefined } } /** * Counts an owner's transfers through a capped subquery sharing the list * filter, so the count cannot scan unboundedly. A result above `cap` means * the true count was truncated. */ export async function count(db: Db.Db, owner: Owner, options: count.Options): Promise { const row = await db.kysely .selectFrom((eb) => { let query = eb .selectFrom('routes_transfers') .select('id') .where('orgId', '=', owner.orgId) .where('environment', '=', owner.environment) if (owner.projectId !== undefined) query = query.where('projectId', '=', owner.projectId) if (options.status !== undefined) query = query.where('status', '=', options.status) return query.limit(options.cap + 1).as('page') }) .select(({ fn }) => fn.countAll().as('count')) .executeTakeFirstOrThrow() return Number(row.count) } export declare namespace count { /** Count options sharing the list filter. */ type Options = { /** Row ceiling; a result above it signals a truncated count. */ cap: number /** Restricts to one lifecycle status. */ status?: Record['status'] | undefined } } /** Lists filtered route transfers, newest update first. */ export function list(db: Db.Db, options: list.Options): Promise { let query = db.kysely.selectFrom('routes_transfers').selectAll() if (options.environments !== undefined) query = query.where('environment', 'in', options.environments) else if (options.environment !== undefined) query = query.where('environment', '=', options.environment) if (options.hasSubsidy === true) query = query .where('subsidyAmount', 'is not', null) .where(sql`(subsidy_amount->>'baseUnits')::numeric > 0`) if (options.methods !== undefined) query = query.where('method', 'in', options.methods) else if (options.method !== undefined) query = query.where('method', '=', options.method) if (options.orgId !== undefined) query = query.where('orgId', '=', options.orgId) if (options.projectId !== undefined) query = query.where('projectId', '=', options.projectId) if (options.providerIds !== undefined) query = query.where('providerId', 'in', options.providerIds) else if (options.providerId !== undefined) query = query.where('providerId', '=', options.providerId) if (options.sourceChainIds !== undefined) query = query.where(sql`snapshot->'sourceChain'->>'id'`, 'in', options.sourceChainIds) else if (options.sourceChainId !== undefined) query = query.where(sql`snapshot->'sourceChain'->>'id' = ${options.sourceChainId}`) if (options.statuses !== undefined) query = query.where('status', 'in', options.statuses) else if (options.status !== undefined) query = query.where('status', '=', options.status) if (options.query !== undefined) { const exact = options.query const value = /^(?:0x)?[0-9a-fA-F]{64}$/.test(options.query) ? options.query.toLowerCase() : options.query query = query.where((eb) => eb.or([ eb('id', '=', exact), eb.exists( eb .selectFrom('routes_transfer_transactions') .select('routes_transfer_transactions.transferId') .whereRef('routes_transfer_transactions.transferId', '=', 'routes_transfers.id') .where('routes_transfer_transactions.transactionRef', '=', value), ), ]), ) } const cursor = options.cursor if (cursor !== undefined) query = query.where((eb) => eb.or([ eb('updatedAt', '<', cursor.updatedAt), eb.and([eb('updatedAt', '=', cursor.updatedAt), eb('id', '<', cursor.id)]), ]), ) return query.orderBy('updatedAt', 'desc').orderBy('id', 'desc').limit(options.limit).execute() } export declare namespace list { /** Transfer filters and page bound. */ type Options = { /** Exclusive lower bound from the previous page. */ cursor?: Cursor | undefined /** Restricts transfers to one route environment. */ environment?: Record['environment'] | undefined /** Restricts transfers to any listed route environment. */ environments?: readonly Record['environment'][] | undefined /** Restricts transfers to records with a positive Tempo subsidy. */ hasSubsidy?: true | undefined /** Maximum rows returned. */ limit: number /** Restricts transfers to one route method. */ method?: Record['method'] | undefined /** Restricts transfers to any listed route method. */ methods?: readonly Record['method'][] | undefined /** Restricts transfers to one organization. */ orgId?: string | undefined /** Restricts transfers to one project. */ projectId?: string | undefined /** Restricts transfers to one provider. */ providerId?: string | undefined /** Restricts transfers to any listed provider. */ providerIds?: readonly string[] | undefined /** Exact transfer id or transaction reference. */ query?: string | undefined /** Restricts transfers to one source chain. */ sourceChainId?: string | undefined /** Restricts transfers to any listed source chain. */ sourceChainIds?: readonly string[] | undefined /** Restricts transfers to one lifecycle status. */ status?: Record['status'] | undefined /** Restricts transfers to any listed lifecycle status. */ statuses?: readonly Record['status'][] | undefined } /** Cursor fields for the last transfer on the previous page. */ type Cursor = { /** Transfer id, used as a deterministic tie-breaker. */ id: string /** Latest material update time. */ updatedAt: string } } /** Summarizes transfer work across the selected filters. */ export async function summarize( db: Db.Db, options: summarize.Options = {}, ): Promise { let query = db.kysely.selectFrom('routes_transfers') if (options.environments !== undefined) query = query.where('environment', 'in', options.environments) else if (options.environment !== undefined) query = query.where('environment', '=', options.environment) if (options.methods !== undefined) query = query.where('method', 'in', options.methods) else if (options.method !== undefined) query = query.where('method', '=', options.method) if (options.orgId !== undefined) query = query.where('orgId', '=', options.orgId) if (options.projectId !== undefined) query = query.where('projectId', '=', options.projectId) if (options.providerIds !== undefined) query = query.where('providerId', 'in', options.providerIds) else if (options.providerId !== undefined) query = query.where('providerId', '=', options.providerId) if (options.sourceChainIds !== undefined) query = query.where(sql`snapshot->'sourceChain'->>'id'`, 'in', options.sourceChainIds) else if (options.sourceChainId !== undefined) query = query.where(sql`snapshot->'sourceChain'->>'id' = ${options.sourceChainId}`) const active = ['action-required', 'processing', 'refunding'] as const const row = await query .select((eb) => [ eb.fn.countAll().filterWhere('status', '=', 'action-required').as('actionRequired'), eb.fn .countAll() .filterWhere('status', 'in', ['processing', 'refunding']) .as('inProgress'), eb.fn .min(eb.fn.coalesce('statusUpdatedAt', 'updatedAt')) .filterWhere('status', 'in', active) .as('oldestActiveAt'), ]) .executeTakeFirstOrThrow() return { actionRequired: Number(row.actionRequired), inProgress: Number(row.inProgress), oldestActiveAt: row.oldestActiveAt, } } export declare namespace summarize { /** Attribution and route filters. */ type Options = { /** Restricts transfers to one route environment. */ environment?: Record['environment'] | undefined /** Restricts transfers to any listed route environment. */ environments?: readonly Record['environment'][] | undefined /** Restricts transfers to one route method. */ method?: Record['method'] | undefined /** Restricts transfers to any listed route method. */ methods?: readonly Record['method'][] | undefined /** Restricts transfers to one organization. */ orgId?: string | undefined /** Restricts transfers to one project. */ projectId?: string | undefined /** Restricts transfers to one provider. */ providerId?: string | undefined /** Restricts transfers to any listed provider. */ providerIds?: readonly string[] | undefined /** Restricts transfers to one source chain. */ sourceChainId?: string | undefined /** Restricts transfers to any listed source chain. */ sourceChainIds?: readonly string[] | undefined } /** Current transfer work totals. */ type Result = { /** Transfers that require operator action. */ actionRequired: number /** Transfers currently owned by platform processing or recovery. */ inProgress: number /** When the oldest active transfer entered its current status. */ oldestActiveAt: string | null } } /** Returns the Tempo-funded subsidy recorded for one outbound transfer. */ export function subsidies(record: Record) { return record.subsidyAmount ? [ { amount: record.subsidyAmount, provider: 'tempo' as const, token: record.snapshot.destinationToken, }, ] : [] } /** Lists completed production transfer subsidies awaiting Stripe acknowledgement. */ export function listUnreportedSubsidies( db: Db.Db, options: listUnreportedSubsidies.Options = {}, ): Promise { let query = db.kysely .selectFrom('routes_transfers') .selectAll() .where('environment', '=', 'production') .where('status', '=', 'completed') .where('subsidyTransactionHash', 'is not', null) .where('subsidyAmount', 'is not', null) .where(sql`(subsidy_amount->>'baseUnits')::numeric > 0`) .where('subsidyMeterReportedAt', 'is', null) .orderBy('statusUpdatedAt', 'asc') .orderBy('id', 'asc') if (options.limit !== undefined) query = query.limit(options.limit) return query.execute() } export declare namespace listUnreportedSubsidies { /** Options for the outbound Routes subsidy reporting scan. */ type Options = { /** Maximum rows to return. */ limit?: number | undefined } } /** Stores the effective timestamp after Stripe acknowledges an outbound subsidy. */ export async function markSubsidyReported(db: Db.Db, id: string, meteredAt: string): Promise { await db.kysely .updateTable('routes_transfers') .set({ subsidyMeterReportedAt: meteredAt }) .where('id', '=', id) .where('subsidyMeterReportedAt', 'is', null) .execute() } /** Returns whether an organization has an outbound subsidy awaiting resolution or reporting. */ export async function hasUnreportedSubsidies(db: Db.Db, orgId: string): Promise { return Boolean( await db.kysely .selectFrom('routes_transfers') .select('id') .where('orgId', '=', orgId) .where((eb) => eb.or([ eb.and([ eb('environment', '=', 'production'), eb('subsidyTransactionHash', 'is not', null), eb('subsidyAmount', 'is not', null), sql`(subsidy_amount->>'baseUnits')::numeric > 0`, eb('subsidyMeterReportedAt', 'is', null), ]), eb.and([ eb('status', 'in', ['awaiting-source', 'processing', 'action-required', 'refunding']), eb('subsidyCommitmentAmount', 'is not', null), sql`(subsidy_commitment_amount->>'baseUnits')::numeric > 0`, ]), ]), ) .limit(1) .executeTakeFirst(), ) } /** Reads outbound subsidy usage for one organization and time window. */ export async function subsidyUsage( db: Db.Db, options: subsidyUsage.Options, ): Promise { const amount = sql`ceil((subsidy_amount->>'baseUnits')::numeric * power(10::numeric, ${Fees.tokenDecimals} - (subsidy_amount->>'decimals')::integer))` const meteredAt = sql`coalesce(subsidy_meter_reported_at, status_updated_at, updated_at)` const row = await db.kysely .selectFrom('routes_transfers') .select([ sql`coalesce(sum(${amount}), 0)`.as('amount'), sql`count(*)`.as('count'), ]) .where('orgId', '=', options.orgId) .where('environment', '=', options.environment) .where('status', '=', 'completed') .where('subsidyTransactionHash', 'is not', null) .where('subsidyAmount', 'is not', null) .where(sql`(subsidy_amount->>'baseUnits')::numeric > 0`) .where(meteredAt, '>=', options.from) .where(meteredAt, '<=', options.to) .executeTakeFirstOrThrow() return { amount: BigInt(row.amount), count: Number(row.count) } } export declare namespace subsidyUsage { /** Organization and meter-event window for outbound subsidy usage. */ type Options = { /** Key environment to include. */ environment: Record['environment'] /** Inclusive meter-event lower bound. */ from: string /** Organization whose subsidy usage is read. */ orgId: string /** Inclusive meter-event upper bound. */ to: string } /** Outbound subsidy usage in six-decimal meter units. */ type Result = { /** USD subsidy total in six-decimal meter units. */ amount: bigint /** Metered outbound subsidies. */ count: number } } /** Returns token and native inventory promised by active transfer commitments. */ export async function subsidyReservations( db: Db.Db, options: subsidyReservations.Options, ): Promise { const tokenAmount = sql`coalesce(sum(CASE WHEN subsidy_token_address = ${options.tokenAddress.toLowerCase()} THEN (subsidy_commitment_amount->>'baseUnits')::numeric ELSE 0 END), 0)` let query = db.kysely .selectFrom('routes_transfers') .select([ sql`coalesce(sum(subsidy_native_amount::numeric), 0)`.as('nativeAmount'), tokenAmount.as('tokenAmount'), ]) .where('status', 'in', ['awaiting-source', 'processing', 'action-required', 'refunding']) .where('subsidyAccount', '=', options.account.toLowerCase()) .where('subsidyChainId', '=', options.chainId) .where('subsidyCommitmentAmount', 'is not', null) .where(sql`(subsidy_commitment_amount->>'baseUnits')::numeric > 0`) if (options.excludeId !== undefined) query = query.where('id', '!=', options.excludeId) const row = await query.executeTakeFirstOrThrow() return { nativeAmount: BigInt(row.nativeAmount), tokenAmount: BigInt(row.tokenAmount) } } /** Lists abandoned subsidy commitments eligible for bounded expiry. */ export function listExpiredSubsidyCommitments( db: Db.Db, options: listExpiredSubsidyCommitments.Options, ): Promise { return db.kysely .selectFrom('routes_transfers') .selectAll() .where('status', '=', 'awaiting-source') .where('quoteExpiresAt', '<=', options.expiredBefore) .where('subsidyCommitmentAmount', 'is not', null) .where(sql`(subsidy_commitment_amount->>'baseUnits')::numeric > 0`) .orderBy('quoteExpiresAt', 'asc') .orderBy('id', 'asc') .limit(options.limit) .execute() } export declare namespace listExpiredSubsidyCommitments { /** Expired commitment scan options. */ type Options = { /** Latest quote expiry eligible for cleanup. */ expiredBefore: string /** Maximum rows returned. */ limit: number } } export declare namespace subsidyReservations { /** Signer, chain, and token identifying one inventory pool. */ type Options = { /** Destination subsidy signer address. */ account: string /** Destination EVM chain CAIP-2 id. */ chainId: string /** Transfer excluded while converting its commitment into a signed transaction. */ excludeId?: string | undefined /** Destination ERC-20 token address. */ tokenAddress: string } /** Inventory promised by signed transactions that have not completed. */ type Result = { /** Native gas amount reserved on the signer chain. */ nativeAmount: bigint /** Destination token amount reserved for the requested token. */ tokenAmount: bigint } } /** * Applies a version-guarded update; `undefined` when the row moved past * `expectedVersion` (or vanished) so the caller retries against fresh state. */ export function update(db: Db.Db, options: update.Options): Promise { return db.kysely .updateTable('routes_transfers') .set({ providerDeliveredAt: options.providerDeliveredAt, providerState: options.providerState, snapshot: options.snapshot, status: options.status, statusReason: options.statusReason, statusUpdatedAt: options.statusUpdatedAt, subsidyAccount: options.subsidyAccount, subsidyAmount: options.subsidyAmount, subsidyChainId: options.subsidyChainId, subsidyCommitmentAmount: options.subsidyCommitmentAmount, subsidyNativeAmount: options.subsidyNativeAmount, subsidyTokenAddress: options.subsidyTokenAddress, subsidyTransactionHash: options.subsidyTransactionHash, updatedAt: options.updatedAt, version: options.version, }) .where('id', '=', options.id) .where('version', '=', options.expectedVersion) .returningAll() .executeTakeFirst() } export declare namespace update { /** Guarded update fields. */ type Options = { /** The version the update was computed against. */ expectedVersion: number /** Transfer id (`rtr_…`). */ id: string /** Replacement private provider state; omitted keeps the stored value. */ providerState?: Record['providerState'] | undefined /** When provider destination delivery was first verified. */ providerDeliveredAt: Record['providerDeliveredAt'] /** Replacement public snapshot. */ snapshot: Record['snapshot'] /** New lifecycle status. */ status: Record['status'] /** Customer-safe reason, or null to clear. */ statusReason: Record['statusReason'] /** When the transfer entered its current status. */ statusUpdatedAt: Record['statusUpdatedAt'] /** Destination subsidy signer address. */ subsidyAccount: Record['subsidyAccount'] /** Destination token amount supplied by Tempo. */ subsidyAmount: Record['subsidyAmount'] /** Destination subsidy chain id. */ subsidyChainId: Record['subsidyChainId'] /** Maximum destination subsidy authorized at action creation. */ subsidyCommitmentAmount: Record['subsidyCommitmentAmount'] /** Native gas amount reserved by the signed transaction. */ subsidyNativeAmount: Record['subsidyNativeAmount'] /** Destination subsidy token address. */ subsidyTokenAddress: Record['subsidyTokenAddress'] /** Destination subsidy transaction hash. */ subsidyTransactionHash: Record['subsidyTransactionHash'] /** Mutation timestamp (ISO 8601). */ updatedAt: string /** The incremented version. */ version: number } }