import { sql, type JSONColumnType, 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' import * as RoutesSubsidies from './routesSubsidies.js' /** Columns of the `routes_deposits` table. */ export type Table = Omit< db_Schema.RoutesDeposit, 'providerRequestIds' | 'providerTransactionHashes' > & { /** Private provider request identifiers stored as JSON. */ providerRequestIds: JSONColumnType /** Private provider transaction references stored as JSON. */ providerTransactionHashes: JSONColumnType } /** A stored route deposit. */ export type Record = Selectable /** Ownership filter inherited from the deposit address. */ export type Owner = { /** Key environment. */ environment: Record['environment'] /** Owning organization id (`org_…`). */ orgId: string /** Attributed project id for project-scoped API keys. */ projectId?: string | undefined } /** Optional filters shared by route deposit list and count queries. */ export type Filters = { /** Reusable source-chain address that received the funds. */ depositAddress?: string | undefined /** Canonical destination token keys accepted by the filter. */ destinationTokenKeys?: readonly string[] | undefined /** Route provider id. */ providerId?: string | undefined /** Tempo account that receives completed deposits. */ recipient?: string | undefined /** Canonical source chain ids accepted by the filter. */ sourceChainIds?: readonly string[] | undefined /** Canonical source token keys accepted by the filter. */ sourceTokenKeys?: readonly string[] | undefined /** Deposit lifecycle status. */ status?: Record['status'] | undefined } /** Aggregates active deposit states for bounded operational metrics. */ export function summarizeStates(db: Db.Db, options: summarizeStates.Options) { const reason = sql`routes_deposits.status_reason->>'code'` return db.kysely .selectFrom('routes_deposits') .innerJoin( 'routes_deposit_addresses', 'routes_deposit_addresses.id', 'routes_deposits.depositAddressId', ) .select(({ fn }) => [ 'routes_deposits.environment', 'routes_deposit_addresses.providerId', 'routes_deposits.status', fn.countAll().as('count'), sql`min(coalesce(routes_deposits.status_updated_at, routes_deposits.updated_at))`.as( 'oldestStatusUpdatedAt', ), reason.as('reason'), ]) .where('routes_deposits.status', 'in', options.statuses) .groupBy([ 'routes_deposits.environment', 'routes_deposit_addresses.providerId', 'routes_deposits.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 one provider-observed source transfer as a route deposit. */ export function insert(db: Db.Db, record: Record): Promise { return db.kysely .insertInto('routes_deposits') .values({ ...record, providerRequestIds: JSON.stringify(record.providerRequestIds), providerTransactionHashes: JSON.stringify(record.providerTransactionHashes), }) .returningAll() .executeTakeFirstOrThrow() } /** Inserts a chain-observed source transfer or returns its existing deposit. */ export async function insertOrGetSource( db: Db.Db, record: Record & { sourceTransferIndex: number }, ): Promise { const inserted = await db.kysely .insertInto('routes_deposits') .values({ ...record, providerRequestIds: JSON.stringify(record.providerRequestIds), providerTransactionHashes: JSON.stringify(record.providerTransactionHashes), }) .onConflict((oc) => oc.columns(['depositAddressId', 'sourceTransactionHash', 'sourceTransferIndex']).doNothing(), ) .returningAll() .executeTakeFirst() if (inserted) return { record: inserted, type: 'created' } const existing = await db.kysely .selectFrom('routes_deposits') .selectAll() .where('depositAddressId', '=', record.depositAddressId) .where('sourceTransactionHash', '=', record.sourceTransactionHash) .where('sourceTransferIndex', '=', record.sourceTransferIndex) .executeTakeFirstOrThrow() return { record: existing, type: 'existing' } } export declare namespace insertOrGetSource { /** Result of idempotently persisting one verified source transfer. */ type Result = { record: Record; type: 'created' | 'existing' } } /** Serializes source attribution for one address and transaction. */ export function withSourceTransaction( db: Db.Db, options: withSourceTransaction.Options, ): Promise { return db.transaction(async (tx) => { const key = `${options.depositAddressId}:${options.sourceTransactionHash.toLowerCase()}` await sql`SELECT pg_advisory_xact_lock(hashtextextended(${key}, 0))`.execute(tx.kysely) return options.fn(tx) }) } export declare namespace withSourceTransaction { /** Source identity and transaction-scoped work protected by its advisory lock. */ type Options = { /** Route deposit address that received the transfer. */ depositAddressId: string /** Database work serialized for the source transaction. */ fn: (db: Db.Db) => Promise /** Source transaction whose deposit attribution must not race. */ sourceTransactionHash: string } } /** Serializes subsidy preparation and persistence with organization deletion. */ export function withSubsidySettlement( db: Db.Db, options: withSubsidySettlement.Options, ): Promise { return db.transaction(async (tx) => { await RoutesSubsidies.lockOrganization(tx, options.orgId) return options.fn(tx) }) } export declare namespace withSubsidySettlement { /** Organization and transaction-scoped work protected by its subsidy lock. */ type Options = { /** Database work serialized with organization deletion. */ fn: (db: Db.Db) => Promise /** Organization whose subsidy liability may change. */ orgId: string } } /** Gets deposits created from one provider-observed source transaction. */ export function listByProviderObservation( db: Db.Db, options: listByProviderObservation.Options, ): Promise { return db.kysely .selectFrom('routes_deposits') .selectAll() .where('depositAddressId', '=', options.depositAddressId) .where('providerRequestId', '=', options.providerRequestId) .where('sourceTransactionHash', '=', options.sourceTransactionHash) .orderBy('providerTransferIndex', 'asc') .execute() } export declare namespace listByProviderObservation { /** Provider observation fields that identify one source transaction. */ type Options = { /** Route deposit address that received the transfer. */ depositAddressId: string /** Provider request that reported the transfer. */ providerRequestId: string /** Provider-observed source transaction reference. */ sourceTransactionHash: string } } /** Attaches first-observation timestamps to every deposit from one provider request. */ export function recordRequestObservation( db: Db.Db, options: recordRequestObservation.Options, ): Promise { return db.kysely .updateTable('routes_deposits') .set({ ...(options.pollObservedAt ? { pollObservedAt: sql`LEAST(COALESCE(poll_observed_at, ${options.pollObservedAt}), ${options.pollObservedAt})`, } : {}), ...(options.webhookReceivedAt ? { webhookReceivedAt: sql`LEAST(COALESCE(webhook_received_at, ${options.webhookReceivedAt}), ${options.webhookReceivedAt})`, } : {}), }) .where('depositAddressId', '=', options.depositAddressId) .where('providerRequestId', '=', options.providerRequestId) .returningAll() .execute() } export declare namespace recordRequestObservation { /** First observations associated with one provider request. */ type Options = { /** Route deposit address id (`rda_…`). */ depositAddressId: string /** When polling first observed the provider request. */ pollObservedAt?: string | undefined /** Provider request identifier. */ providerRequestId: string /** When Tempo first received an authenticated provider webhook. */ webhookReceivedAt?: string | undefined } } /** Gets deposits associated with one address and source transaction. */ export function listBySourceTransaction( db: Db.Db, options: listBySourceTransaction.Options, ): Promise { return db.kysely .selectFrom('routes_deposits') .selectAll() .where('depositAddressId', '=', options.depositAddressId) .where('sourceTransactionHash', '=', options.sourceTransactionHash) .orderBy('sourceTransferIndex', 'asc') .execute() } export declare namespace listBySourceTransaction { /** Source transaction fields scoped to one deposit address. */ type Options = { /** Route deposit address that received the transfer. */ depositAddressId: string /** Provider-observed source transaction reference. */ sourceTransactionHash: string } } /** Lists a bounded page of source transfers awaiting provider attribution. */ export function listUnattributed(db: Db.Db, options: listUnattributed.Options): Promise { let query = db.kysely .selectFrom('routes_deposits') .selectAll() .where('depositAddressId', '=', options.depositAddressId) .where('providerRequestId', 'is', null) .where('createdAt', '<=', options.before) if (options.cursor) query = query.where( sql`(created_at, id) > (${options.cursor.createdAt}, ${options.cursor.id})`, ) return query.orderBy('createdAt', 'asc').orderBy('id', 'asc').limit(options.limit).execute() } export declare namespace listUnattributed { /** Fixed sweep boundary and page position for one address. */ type Options = { /** Inclusive creation-time boundary retained until the sweep finishes. */ before: string /** Last deposit attempted in this sweep. */ cursor?: listByAddress.Cursor | undefined /** Route deposit address whose transfers are reconciled. */ depositAddressId: string /** Maximum deposits returned. */ limit: number } } /** Returns whether an address has a source transfer awaiting provider attribution. */ export async function hasUnattributed(db: Db.Db, depositAddressId: string): Promise { return Boolean( await db.kysely .selectFrom('routes_deposits') .select('id') .where('depositAddressId', '=', depositAddressId) .where('providerRequestId', 'is', null) .limit(1) .executeTakeFirst(), ) } /** Gets a route deposit without ownership scoping. */ export function get(db: Db.Db, id: string): Promise { return db.kysely.selectFrom('routes_deposits').selectAll().where('id', '=', id).executeTakeFirst() } /** Gets the earliest source transfer detected for one reusable address. */ export function getFirstByAddress(db: Db.Db, addressId: string): Promise { return db.kysely .selectFrom('routes_deposits') .selectAll() .where('depositAddressId', '=', addressId) .orderBy('createdAt', 'asc') .orderBy('id', 'asc') .limit(1) .executeTakeFirst() } /** Gets an owner-visible route deposit. */ export function getOwned(db: Db.Db, owner: Owner, id: string): Promise { let query = db.kysely .selectFrom('routes_deposits') .selectAll() .where('environment', '=', owner.environment) .where('id', '=', id) .where('orgId', '=', owner.orgId) if (owner.projectId !== undefined) query = query.where('projectId', '=', owner.projectId) return query.executeTakeFirst() } /** Lists owner-visible route deposits newest-first. */ export function listByOwner(db: Db.Db, options: listByOwner.Options): Promise { const { cursor, limit } = options let builder = query(db, options).selectAll('routes_deposits') if (cursor) builder = builder.where((eb) => eb.or([ eb('routes_deposits.createdAt', '<', cursor.createdAt), eb.and([ eb('routes_deposits.createdAt', '=', cursor.createdAt), eb('routes_deposits.id', '<', cursor.id), ]), ]), ) return builder .orderBy('routes_deposits.createdAt', 'desc') .orderBy('routes_deposits.id', 'desc') .limit(limit) .execute() } export declare namespace listByOwner { /** Ownership and paging options for route deposits. */ type Options = Filters & { /** Exclusive lower bound from the previous page. */ cursor?: listByAddress.Cursor | undefined /** Maximum rows to return. */ limit: number /** Ownership filter derived from the authenticated API key. */ owner: Owner } } /** Lists filtered route deposits, newest update first. */ export function list(db: Db.Db, options: list.Options): Promise { let query = db.kysely .selectFrom('routes_deposits') .innerJoin( 'routes_deposit_addresses', 'routes_deposit_addresses.id', 'routes_deposits.depositAddressId', ) .selectAll('routes_deposits') if (options.depositAddressId !== undefined) query = query.where('routes_deposits.depositAddressId', '=', options.depositAddressId) if (options.environments !== undefined) query = query.where('routes_deposits.environment', 'in', options.environments) else if (options.environment !== undefined) query = query.where('routes_deposits.environment', '=', options.environment) if (options.hasSubsidy === true) query = query.where((eb) => eb.or([ sql`(routes_deposits.provider_subsidy->'amount'->>'baseUnits')::numeric > 0`, eb.and([ eb('routes_deposits.settlementTransactionHash', 'is not', null), sql`(routes_deposits.subsidy_amount->>'baseUnits')::numeric > 0`, ]), ]), ) if (options.orgId !== undefined) query = query.where('routes_deposits.orgId', '=', options.orgId) if (options.projectId !== undefined) query = query.where('routes_deposits.projectId', '=', options.projectId) if (options.providerIds !== undefined) query = query.where('routes_deposit_addresses.providerId', 'in', options.providerIds) else if (options.providerId !== undefined) query = query.where('routes_deposit_addresses.providerId', '=', options.providerId) if (options.sourceChainIds !== undefined) query = query.where('routes_deposits.sourceChainId', 'in', options.sourceChainIds) else if (options.sourceChainId !== undefined) query = query.where('routes_deposits.sourceChainId', '=', options.sourceChainId) if (options.statuses !== undefined) query = query.where('routes_deposits.status', 'in', options.statuses) else if (options.status !== undefined) query = query.where('routes_deposits.status', '=', options.status) if (options.query !== undefined) { const exact = options.query const address = /^0x[0-9a-fA-F]{40}$/.test(options.query) ? options.query.toLowerCase() : options.query const transactionRef = /^(?:0x)?[0-9a-fA-F]{64}$/.test(options.query) ? options.query.toLowerCase() : options.query query = query.where((eb) => eb.or([ eb('routes_deposit_addresses.address', '=', address), eb('routes_deposits.id', '=', exact), eb('routes_deposits.settlementTransactionHash', '=', transactionRef), eb('routes_deposits.sourceTransactionHash', '=', transactionRef), sql`routes_deposits.provider_transaction_hashes @> ${JSON.stringify([transactionRef])}::jsonb`, sql`routes_deposits.snapshot->'destinationTransactionHashes' @> ${JSON.stringify([transactionRef])}::jsonb`, sql`routes_deposits.snapshot->'refundTransactionHashes' @> ${JSON.stringify([transactionRef])}::jsonb`, ]), ) } const cursor = options.cursor if (cursor !== undefined) query = query.where((eb) => eb.or([ eb('routes_deposits.updatedAt', '<', cursor.updatedAt), eb.and([ eb('routes_deposits.updatedAt', '=', cursor.updatedAt), eb('routes_deposits.id', '<', cursor.id), ]), ]), ) return query .orderBy('routes_deposits.updatedAt', 'desc') .orderBy('routes_deposits.id', 'desc') .limit(options.limit) .execute() } export declare namespace list { /** Deposit filters and page bound. */ type Options = { /** Exclusive lower bound from the previous page. */ cursor?: Cursor | undefined /** Restricts deposits to one deposit address. */ depositAddressId?: string | undefined /** Restricts deposits to one route environment. */ environment?: Record['environment'] | undefined /** Restricts deposits to any listed route environment. */ environments?: readonly Record['environment'][] | undefined /** Restricts deposits to records with a positive provider or Tempo subsidy. */ hasSubsidy?: true | undefined /** Maximum rows returned. */ limit: number /** Restricts deposits to one organization. */ orgId?: string | undefined /** Restricts deposits to one project. */ projectId?: string | undefined /** Restricts deposits to one provider. */ providerId?: string | undefined /** Restricts deposits to any listed provider. */ providerIds?: readonly string[] | undefined /** Exact deposit id, deposit address, or transaction reference. */ query?: string | undefined /** Restricts deposits to one source chain. */ sourceChainId?: string | undefined /** Restricts deposits to any listed source chain. */ sourceChainIds?: readonly string[] | undefined /** Restricts deposits to one lifecycle status. */ status?: Record['status'] | undefined /** Restricts deposits to any listed lifecycle status. */ statuses?: readonly Record['status'][] | undefined } /** Cursor fields for the last deposit on the previous page. */ type Cursor = { /** Deposit id, used as a deterministic tie-breaker. */ id: string /** Latest material update time. */ updatedAt: string } } /** Summarizes deposit work across the selected filters. */ export async function summarize( db: Db.Db, options: summarize.Options = {}, ): Promise { let query = db.kysely .selectFrom('routes_deposits') .innerJoin( 'routes_deposit_addresses', 'routes_deposit_addresses.id', 'routes_deposits.depositAddressId', ) if (options.environments !== undefined) query = query.where('routes_deposits.environment', 'in', options.environments) else if (options.environment !== undefined) query = query.where('routes_deposits.environment', '=', options.environment) if (options.orgId !== undefined) query = query.where('routes_deposits.orgId', '=', options.orgId) if (options.projectId !== undefined) query = query.where('routes_deposits.projectId', '=', options.projectId) if (options.providerIds !== undefined) query = query.where('routes_deposit_addresses.providerId', 'in', options.providerIds) else if (options.providerId !== undefined) query = query.where('routes_deposit_addresses.providerId', '=', options.providerId) if (options.sourceChainIds !== undefined) query = query.where('routes_deposits.sourceChainId', 'in', options.sourceChainIds) else if (options.sourceChainId !== undefined) query = query.where('routes_deposits.sourceChainId', '=', options.sourceChainId) const active = ['action-required', 'bridging', 'detected', 'refunding', 'settling'] as const const row = await query .select((eb) => [ eb.fn .countAll() .filterWhere('routes_deposits.status', '=', 'action-required') .as('actionRequired'), eb.fn .countAll() .filterWhere('routes_deposits.status', 'in', [ 'bridging', 'detected', 'refunding', 'settling', ]) .as('inProgress'), eb.fn .min(eb.fn.coalesce('routes_deposits.statusUpdatedAt', 'routes_deposits.updatedAt')) .filterWhere('routes_deposits.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 deposits to one route environment. */ environment?: Record['environment'] | undefined /** Restricts deposits to any listed route environment. */ environments?: readonly Record['environment'][] | undefined /** Restricts deposits to one organization. */ orgId?: string | undefined /** Restricts deposits to one project. */ projectId?: string | undefined /** Restricts deposits to one provider. */ providerId?: string | undefined /** Restricts deposits to any listed provider. */ providerIds?: readonly string[] | undefined /** Restricts deposits to one source chain. */ sourceChainId?: string | undefined /** Restricts deposits to any listed source chain. */ sourceChainIds?: readonly string[] | undefined } /** Current deposit work totals. */ type Result = { /** Deposits that require operator action. */ actionRequired: number /** Deposits currently owned by platform processing or recovery. */ inProgress: number /** When the oldest active deposit entered its current status. */ oldestActiveAt: string | null } } /** Counts owner-visible route deposits through the list's shared filters. */ export async function count(db: Db.Db, options: count.Options): Promise { const rows = query(db, options) .select('routes_deposits.id') .limit(options.cap + 1) .as('rows') const result = await db.kysely .selectFrom(rows) .select(sql`count(*)::int`.as('count')) .executeTakeFirstOrThrow() return result.count } export declare namespace count { /** Ownership, filter, and cap options for counting route deposits. */ type Options = Filters & { /** Maximum exact count before returning the cap plus one. */ cap: number /** Ownership filter derived from the authenticated API key. */ owner: Owner } } function query(db: Db.Db, options: query.Options) { const owner = options.owner let builder = db.kysely .selectFrom('routes_deposits') .innerJoin( 'routes_deposit_addresses', 'routes_deposit_addresses.id', 'routes_deposits.depositAddressId', ) .where('routes_deposit_addresses.environment', '=', owner.environment) .where('routes_deposit_addresses.orgId', '=', owner.orgId) .where('routes_deposits.environment', '=', owner.environment) .where('routes_deposits.orgId', '=', owner.orgId) if (owner.projectId !== undefined) builder = builder .where('routes_deposit_addresses.projectId', '=', owner.projectId) .where('routes_deposits.projectId', '=', owner.projectId) if (options.depositAddress !== undefined) { const address = /^0x[0-9a-fA-F]{40}$/.test(options.depositAddress) ? options.depositAddress.toLowerCase() : options.depositAddress builder = builder.where('routes_deposit_addresses.address', '=', address) } if (options.destinationTokenKeys !== undefined) builder = options.destinationTokenKeys.length ? builder.where( 'routes_deposit_addresses.destinationTokenKey', 'in', options.destinationTokenKeys, ) : builder.where(sql`false`) if (options.providerId !== undefined) builder = builder.where('routes_deposit_addresses.providerId', '=', options.providerId) if (options.recipient !== undefined) builder = builder.where( sql`lower(routes_deposit_addresses.recipient) = ${options.recipient.toLowerCase()}`, ) if (options.sourceChainIds !== undefined) builder = options.sourceChainIds.length ? builder.where('routes_deposit_addresses.sourceChainId', 'in', options.sourceChainIds) : builder.where(sql`false`) if (options.sourceTokenKeys !== undefined) builder = options.sourceTokenKeys.length ? builder.where('routes_deposit_addresses.sourceTokenKey', 'in', options.sourceTokenKeys) : builder.where(sql`false`) if (options.status !== undefined) builder = builder.where('routes_deposits.status', '=', options.status) return builder } declare namespace query { type Options = Filters & { owner: Owner } } /** Lists one owner-visible address's deposits newest-first. */ export function listByAddress( db: Db.Db, owner: Owner, addressId: string, options: listByAddress.Options, ): Promise { const { cursor, limit } = options let query = db.kysely .selectFrom('routes_deposits') .selectAll() .where('depositAddressId', '=', addressId) .where('environment', '=', owner.environment) .where('orgId', '=', owner.orgId) if (owner.projectId !== undefined) query = query.where('projectId', '=', owner.projectId) if (cursor) query = query.where((eb) => eb.or([ eb('createdAt', '<', cursor.createdAt), eb.and([eb('createdAt', '=', cursor.createdAt), eb('id', '<', cursor.id)]), ]), ) return query.orderBy('createdAt', 'desc').orderBy('id', 'desc').limit(limit).execute() } export declare namespace listByAddress { /** Cursor fields for the last deposit returned by the previous page. */ type Cursor = { /** Deposit creation time. */ createdAt: string /** Deposit id, used as a deterministic tie-breaker. */ id: string } /** Paging options for one address's deposits. */ type Options = { /** Exclusive lower bound from the previous page. */ cursor?: Cursor | undefined /** Maximum rows to return. */ limit: number } } /** Lists deposits for one reusable source-chain address newest-first. */ export function listByDepositAddress( db: Db.Db, options: listByDepositAddress.Options, ): Promise { const { cursor, depositAddress, limit, owner } = options const address = /^0x[0-9a-fA-F]{40}$/.test(depositAddress) ? depositAddress.toLowerCase() : depositAddress let query = db.kysely .selectFrom('routes_deposits') .innerJoin( 'routes_deposit_addresses', 'routes_deposit_addresses.id', 'routes_deposits.depositAddressId', ) .selectAll('routes_deposits') .where('routes_deposit_addresses.address', '=', address) .where('routes_deposit_addresses.environment', '=', owner.environment) .where('routes_deposit_addresses.orgId', '=', owner.orgId) .where('routes_deposits.environment', '=', owner.environment) .where('routes_deposits.orgId', '=', owner.orgId) if (owner.projectId !== undefined) query = query .where('routes_deposit_addresses.projectId', '=', owner.projectId) .where('routes_deposits.projectId', '=', owner.projectId) if (cursor) query = query.where((eb) => eb.or([ eb('routes_deposits.createdAt', '<', cursor.createdAt), eb.and([ eb('routes_deposits.createdAt', '=', cursor.createdAt), eb('routes_deposits.id', '<', cursor.id), ]), ]), ) return query .orderBy('routes_deposits.createdAt', 'desc') .orderBy('routes_deposits.id', 'desc') .limit(limit) .execute() } export declare namespace listByDepositAddress { /** Deposit address filter, ownership, and paging options. */ type Options = { /** Exclusive lower bound from the previous page. */ cursor?: listByAddress.Cursor | undefined /** Reusable source-chain address that received the funds. */ depositAddress: string /** Maximum rows to return. */ limit: number /** Ownership filter derived from the authenticated API key. */ owner: Owner } } /** Lists one recipient's owner-visible deposits newest-first. */ export function listByRecipient(db: Db.Db, options: listByRecipient.Options): Promise { const { cursor, limit, owner, recipient } = options let query = db.kysely .selectFrom('routes_deposits') .selectAll() .where('environment', '=', owner.environment) .where('orgId', '=', owner.orgId) .where(sql`lower(snapshot->>'recipient') = ${recipient.toLowerCase()}`) if (owner.projectId !== undefined) query = query.where('projectId', '=', owner.projectId) if (cursor) query = query.where((eb) => eb.or([ eb('createdAt', '<', cursor.createdAt), eb.and([eb('createdAt', '=', cursor.createdAt), eb('id', '<', cursor.id)]), ]), ) return query.orderBy('createdAt', 'desc').orderBy('id', 'desc').limit(limit).execute() } export declare namespace listByRecipient { /** Recipient filter, ownership, and paging options. */ type Options = { /** Exclusive lower bound from the previous page. */ cursor?: listByAddress.Cursor | undefined /** Maximum rows to return. */ limit: number /** Ownership filter derived from the authenticated API key. */ owner: Owner /** Tempo account that receives completed deposits. */ recipient: string } } /** Returns provider-paid and settled Tempo subsidies for one deposit. */ export function subsidies(record: Record) { return [ ...(record.providerSubsidy ? [record.providerSubsidy] : []), ...(record.settlementTransactionHash && record.subsidyAmount ? [ { amount: record.subsidyAmount, provider: 'tempo' as const, token: record.snapshot.destinationToken, }, ] : []), ] } /** Lists verified production subsidies awaiting a Stripe acknowledgement. */ export function listUnreportedSubsidies( db: Db.Db, options: listUnreportedSubsidies.Options = {}, ): Promise { let query = db.kysely .selectFrom('routes_deposits') .selectAll() .where('environment', '=', 'production') .where((eb) => eb.or([ eb('status', '=', 'completed'), eb.and([ eb('status', 'in', ['action-required', 'settling']), eb('tempoGasPaid', 'is not', null), ]), ]), ) .where((eb) => eb.or([ sql`(provider_subsidy->'amount'->>'baseUnits')::numeric > 0`, eb.and([ eb('settlementTransactionHash', 'is not', null), eb('subsidyAmount', 'is not', null), 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) if (options.orgId !== undefined) query = query.where('orgId', '=', options.orgId) return query.execute() } export declare namespace listUnreportedSubsidies { /** Options for the Routes subsidy reporting scan. */ type Options = { /** Maximum rows to return. */ limit?: number | undefined /** Organization whose pending usage is listed. */ orgId?: string | undefined } } /** Stores the effective timestamp after Stripe acknowledges a Routes subsidy meter event. */ export async function markSubsidyReported(db: Db.Db, id: string, meteredAt: string): Promise { await db.kysely .updateTable('routes_deposits') .set({ subsidyMeterReportedAt: meteredAt }) .where('id', '=', id) .where('subsidyMeterReportedAt', 'is', null) .execute() } /** Returns whether one organization has a Routes subsidy awaiting Stripe settlement. */ export async function hasUnreportedSubsidies( db: Db.Db, options: hasUnreportedSubsidies.Options, ): Promise { return Boolean( await db.kysely .selectFrom('routes_deposits') .select('id') .where('environment', '=', 'production') .where('orgId', '=', options.orgId) .where('status', 'in', ['action-required', 'settling', 'completed']) .where((eb) => eb.or([ eb.and([ sql`(provider_subsidy->'amount'->>'baseUnits')::numeric > 0`, eb.or([ eb('status', '=', 'completed'), eb.and([ eb('status', 'in', ['action-required', 'settling']), eb('tempoGasPaid', 'is not', null), ]), ]), ]), eb.and([ eb('settlementTransactionHash', 'is not', null), eb('subsidyAmount', 'is not', null), sql`(subsidy_amount->>'baseUnits')::numeric > 0`, ]), ]), ) .where('subsidyMeterReportedAt', 'is', null) .limit(1) .executeTakeFirst(), ) } export declare namespace hasUnreportedSubsidies { /** Organization filter for pending Routes subsidy usage. */ type Options = { /** Organization whose pending usage is checked. */ orgId: string } } /** Reads billable Routes subsidy usage for one organization and time window. */ export async function subsidyUsage( db: Db.Db, options: subsidyUsage.Options, ): Promise { // Sum the same per-row ceilings used by Stripe reporting so admin usage matches the invoice. const providerAmount = sql`coalesce(ceil((provider_subsidy->'amount'->>'baseUnits')::numeric * power(10::numeric, ${Fees.tokenDecimals} - (provider_subsidy->'amount'->>'decimals')::integer)), 0)` const tempoAmount = sql`case when settlement_transaction_hash is not null then coalesce(ceil((subsidy_amount->>'baseUnits')::numeric * power(10::numeric, ${Fees.tokenDecimals} - (subsidy_amount->>'decimals')::integer)), 0) else 0 end` const amount = sql`${providerAmount} + ${tempoAmount}` const meteredAt = sql`coalesce(subsidy_meter_reported_at, status_updated_at, updated_at)` const row = await db.kysely .selectFrom('routes_deposits') .select([ sql`coalesce(sum(${amount}), 0)`.as('amount'), sql`count(*)`.as('count'), ]) .where('orgId', '=', options.orgId) .where('environment', '=', options.environment) .where((eb) => eb.or([ eb('status', '=', 'completed'), eb.and([ eb('status', 'in', ['action-required', 'settling']), eb('tempoGasPaid', 'is not', null), ]), ]), ) .where((eb) => eb.or([ sql`(provider_subsidy->'amount'->>'baseUnits')::numeric > 0`, eb.and([ eb('settlementTransactionHash', 'is not', null), eb('subsidyAmount', 'is not', null), 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 Routes 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 } /** Billable Routes subsidy usage in fee-payer meter units. */ type Result = { /** USD subsidy total in six-decimal meter units. */ amount: bigint /** Metered provider fees and Tempo top-ups. */ count: number } } /** Applies a version-guarded deposit update. */ export function update(db: Db.Db, options: update.Options): Promise { return db.kysely .updateTable('routes_deposits') .set({ providerDeliveredAt: options.providerDeliveredAt, providerOutputAmount: options.providerOutputAmount, providerRequestId: options.providerRequestId, providerRequestIds: JSON.stringify(options.providerRequestIds), providerState: options.providerState, providerSubsidy: options.providerSubsidy, providerTransactionHashes: JSON.stringify(options.providerTransactionHashes), providerTransferIndex: options.providerTransferIndex, retryState: options.retryState, settlementTransaction: options.settlementTransaction, settlementTransactionHash: options.settlementTransactionHash, snapshot: options.snapshot, sourceTransferIndex: options.sourceTransferIndex, status: options.status, statusReason: options.statusReason, statusUpdatedAt: options.statusUpdatedAt, subsidyAmount: options.subsidyAmount, tempoGasPaid: options.tempoGasPaid, updatedAt: options.updatedAt, version: options.version, }) .where('id', '=', options.id) .where('version', '=', options.expectedVersion) .returningAll() .executeTakeFirst() } export declare namespace update { /** Guarded deposit update fields. */ type Options = { /** Stored version required for the update. */ expectedVersion: number /** Route deposit id (`rdp_…`). */ id: string /** When provider delivery was first verified. */ providerDeliveredAt: Record['providerDeliveredAt'] /** Provider output amount confirmed on Tempo. */ providerOutputAmount: Record['providerOutputAmount'] /** Provider request that first reported the deposit. */ providerRequestId: Record['providerRequestId'] /** Private provider request identifiers. */ providerRequestIds: Record['providerRequestIds'] /** Bounded private provider state. */ providerState: Record['providerState'] /** Provider-paid subsidy amount and token. */ providerSubsidy: Record['providerSubsidy'] /** Provider-leg transaction references. */ providerTransactionHashes: Record['providerTransactionHashes'] /** Deposit position within the provider request. */ providerTransferIndex: Record['providerTransferIndex'] /** Bounded private retry state. */ retryState: Record['retryState'] /** Persisted settlement transaction bytes for exactly-once rebroadcast. */ settlementTransaction: Record['settlementTransaction'] /** Persisted settlement transaction hash. */ settlementTransactionHash: Record['settlementTransactionHash'] /** Replacement public snapshot. */ snapshot: Record['snapshot'] /** Verified transfer position in the source transaction, or null. */ sourceTransferIndex: Record['sourceTransferIndex'] /** Replacement lifecycle status. */ status: Record['status'] /** Replacement customer-safe reason. */ statusReason: Record['statusReason'] /** When the deposit entered its current status. */ statusUpdatedAt: Record['statusUpdatedAt'] /** Destination token amount supplied by Tempo. */ subsidyAmount: Record['subsidyAmount'] /** Tempo gas paid for settlement, in base units. */ tempoGasPaid: Record['tempoGasPaid'] /** New material update time. */ updatedAt: string /** New material version. */ version: number } }