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
}
}