import type { ExpressionBuilder, Selectable } from 'kysely' import * as pg from 'pg' import type * as Db from '../Db.js' import type * as db_Schema from '../Schema.js' /** Columns of the `routes_idempotency_requests` table. */ export type Table = db_Schema.RoutesIdempotencyRequest /** A stored idempotency claim. */ export type Record = Selectable /** Route creation operations sharing the idempotency table. */ export const operations = { depositAddress: 'deposit_address', transfer: 'transfer', } as const /** Route creation operation owning an idempotency claim. */ export type Operation = (typeof operations)[keyof typeof operations] /** * Claims `(apiKeyId, keyHash)` for one creation attempt. A live claim rejects * mismatches, blocks pending work, resumes checkpoints, or replays completion. * Deposit checkpoints are shared across rotated keys for the same scoped request. */ export async function claim(db: Db.Db, input: claim.Input): Promise { const now = new Date() const requestHashes = [input.requestHash, ...(input.compatibleRequestHashes ?? [])] const record: Record = { apiKeyId: input.apiKeyId, createdAt: now.toISOString(), expiresAt: new Date(now.getTime() + input.ttlMs).toISOString(), irrevocable: input.irrevocable, keyHash: input.keyHash, matchHash: null, operation: input.operation, orgId: input.orgId, requestHash: input.requestHash, response: null, status: 'pending', transferId: null, } const inserted = await db.kysely .insertInto('routes_idempotency_requests') .values(record) .onConflict((oc) => oc.columns(['apiKeyId', 'keyHash']).doNothing()) .returningAll() .executeTakeFirst() if (inserted) { if (input.operation !== operations.depositAddress) return { createdAt: inserted.createdAt, type: 'claimed' } // Expired unsubsidized checkpoints stop sharing; irrevocable subsidized results remain recoverable after key rotation. const recoverable = await db.kysely .selectFrom('routes_idempotency_requests') .selectAll() .where((eb) => eb.or([eb('apiKeyId', '!=', inserted.apiKeyId), eb('keyHash', '!=', inserted.keyHash)]), ) .where((eb) => eb.or([ eb.and([eb('operation', '=', operations.depositAddress), eb('orgId', '=', input.orgId)]), eb.and([ eb('apiKeyId', '=', input.apiKeyId), eb('operation', 'is', null), eb('orgId', 'is', null), ]), ]), ) .where('requestHash', 'in', requestHashes) .where('status', '=', 'provisioned') .where((eb) => eb.or([eb('expiresAt', '>', now.toISOString()), eb('irrevocable', '=', true)])) .orderBy('createdAt', 'asc') .forUpdate() .executeTakeFirst() if (!recoverable) return { createdAt: inserted.createdAt, type: 'claimed' } await db.kysely .updateTable('routes_idempotency_requests') .set({ irrevocable: recoverable.irrevocable || input.recoveredIrrevocable === true, operation: input.operation, orgId: input.orgId, requestHash: input.requestHash, }) .where('apiKeyId', '=', recoverable.apiKeyId) .where('createdAt', '=', recoverable.createdAt) .where('keyHash', '=', recoverable.keyHash) .where('status', '=', 'provisioned') .execute() return { checkpointFresh: recoverable.expiresAt > now.toISOString(), createdAt: inserted.createdAt, recoveredClaim: { apiKeyId: recoverable.apiKeyId, createdAt: recoverable.createdAt, keyHash: recoverable.keyHash, }, response: recoverable.response ?? '', type: 'resume', } } // Reclaim an expired row in place; the guard keeps two callers from both // reclaiming (only one UPDATE matches the still-expired predicate). const reclaimed = await db.kysely .updateTable('routes_idempotency_requests') .set(record) .where('apiKeyId', '=', input.apiKeyId) .where('keyHash', '=', input.keyHash) .where('expiresAt', '<=', now.toISOString()) .where('irrevocable', '=', false) .where('status', '!=', 'provisioned') .returningAll() .executeTakeFirst() if (reclaimed) return { createdAt: reclaimed.createdAt, type: 'claimed' } const existing = await db.kysely .selectFrom('routes_idempotency_requests') .selectAll() .where('apiKeyId', '=', input.apiKeyId) .where('keyHash', '=', input.keyHash) .executeTakeFirst() // The row vanished between the insert and the read (released claim); retry. if (!existing) return claim(db, input) if (!requestHashes.includes(existing.requestHash)) return { type: 'mismatch' } if (existing.orgId !== null && existing.orgId !== input.orgId) return { type: 'mismatch' } if (existing.operation !== null && existing.operation !== input.operation) return { type: 'mismatch' } // Claims created before organization reservations existed adopt ownership when retried. if (existing.orgId === null || existing.operation === null) { const adopted = await db.kysely .updateTable('routes_idempotency_requests') .set({ irrevocable: existing.irrevocable || (existing.status === 'provisioned' && input.recoveredIrrevocable === true), operation: input.operation, orgId: input.orgId, requestHash: input.requestHash, }) .where('apiKeyId', '=', input.apiKeyId) .where('keyHash', '=', input.keyHash) .where((eb) => eb.or([eb('operation', 'is', null), eb('orgId', 'is', null)])) .returning('apiKeyId') .executeTakeFirst() if (!adopted) return claim(db, input) } if (existing.status === 'pending') return { type: 'pending' } if (existing.status === 'provisioned') return { checkpointFresh: existing.expiresAt > now.toISOString(), createdAt: existing.createdAt, response: existing.response ?? '', type: 'resume', } return { response: existing.response ?? '', type: 'replay' } } export declare namespace claim { /** Claim identity and retention. */ type Input = { /** API key id (`key_…`) the claim is scoped to. */ apiKeyId: string /** Older fingerprints accepted only when reading an existing claim. */ compatibleRequestHashes?: readonly string[] | undefined /** Whether provider success creates liability that must not expire. */ irrevocable: boolean /** SHA-256 of the caller Idempotency-Key. */ keyHash: string /** Route creation operation owning the claim. */ operation: Operation /** Organization owning the route creation. */ orgId: string /** Whether an adopted provider checkpoint carries irrevocable liability. */ recoveredIrrevocable?: boolean | undefined /** SHA-256 fingerprint of the canonical validated request. */ requestHash: string /** Lease duration for the pending execution, in milliseconds. */ ttlMs: number } /** Outcome of a claim attempt. */ type Result = | { createdAt: string; type: 'claimed' } | { type: 'mismatch' } | { type: 'pending' } | { response: string; type: 'replay' } | { /** Whether the checkpoint still carries current quote terms. */ checkpointFresh: boolean createdAt: string recoveredClaim?: release.Input | undefined response: string type: 'resume' } } /** Returns whether the database rejected a claim for a deleted owner. */ export function isOwnerFenceError(cause: unknown): boolean { return ( cause instanceof pg.DatabaseError && cause.code === '23503' && [ 'ownerless route idempotency creation is disabled', 'route idempotency owner was deleted', ].includes(cause.message) ) } /** Marks an owned live claim as potentially creating irrevocable provider liability. */ export async function markProviderAttempt( db: Db.Db, input: markProviderAttempt.Input, ): Promise { const now = new Date() return Boolean( await db.kysely .updateTable('routes_idempotency_requests') .set({ expiresAt: new Date(now.getTime() + input.ttlMs).toISOString(), irrevocable: true, }) .where('apiKeyId', '=', input.apiKeyId) .where('createdAt', '=', input.createdAt) .where('expiresAt', '>', now.toISOString()) .where('keyHash', '=', input.keyHash) .where('status', '=', 'pending') .returning('apiKeyId') .executeTakeFirst(), ) } export declare namespace markProviderAttempt { /** Owned pending claim and renewed liability lease. */ type Input = release.Input & { /** Renewed pending lease duration, in milliseconds. */ ttlMs: number } } /** Counts pending or provisioned deposit-address claims for an organization creator. */ export async function countInFlightDepositAddresses( db: Db.Db, owner: countInFlightDepositAddresses.Owner, ): Promise { const now = new Date().toISOString() let query = db.kysely .selectFrom('routes_idempotency_requests') .select(({ fn }) => fn.countAll().as('count')) .where((eb) => inFlight(eb, now)) .where('operation', '=', operations.depositAddress) const userId = owner.userId query = userId ? query.where('orgId', 'in', (eb) => eb.selectFrom('organizations').select('id').where('userId', '=', userId), ) : query.where('orgId', '=', owner.orgId) const result = await query.executeTakeFirstOrThrow() return Number(result.count) } export declare namespace countInFlightDepositAddresses { /** Canonical creator scope, with organization fallback for unowned rows. */ type Owner = { /** Organization id (`org_…`). */ orgId: string /** Canonical creator id (`usr_…`). */ userId: string | null } } /** Reserves one canonical reusable match across idempotency keys. */ export async function reserveDepositAddressMatch( db: Db.Db, input: reserveDepositAddressMatch.Input, ): Promise { const now = new Date() const pending = await db.kysely .updateTable('routes_idempotency_requests') .set({ expiresAt: new Date(now.getTime() + input.ttlMs).toISOString(), matchHash: input.matchHash, }) .where('apiKeyId', '=', input.apiKeyId) .where('createdAt', '=', input.createdAt) .where('keyHash', '=', input.keyHash) .where('expiresAt', '>', now.toISOString()) .where('status', '=', 'pending') .returningAll() .executeTakeFirst() const claimed = pending ?? (await db.kysely .updateTable('routes_idempotency_requests') .set({ matchHash: input.matchHash }) .where('apiKeyId', '=', input.apiKeyId) .where('createdAt', '=', input.createdAt) .where('keyHash', '=', input.keyHash) .where('status', '=', 'provisioned') .returningAll() .executeTakeFirst()) if (!claimed) return { type: 'lost' } const provisioned = await db.kysely .selectFrom('routes_idempotency_requests') .selectAll() .where((eb) => eb.or([eb('apiKeyId', '!=', input.apiKeyId), eb('keyHash', '!=', input.keyHash)]), ) .where('matchHash', '=', input.matchHash) .where('operation', '=', operations.depositAddress) .where('orgId', '=', input.orgId) .where('status', '=', 'provisioned') .where((eb) => eb.or([eb('expiresAt', '>', now.toISOString()), eb('irrevocable', '=', true)])) .orderBy('createdAt', 'asc') .forUpdate() .executeTakeFirst() if (provisioned && claimed.status === 'pending') { return { checkpointFresh: provisioned.expiresAt > now.toISOString(), claim: { apiKeyId: provisioned.apiKeyId, createdAt: provisioned.createdAt, keyHash: provisioned.keyHash, }, response: provisioned.response ?? '', type: 'resume', } } const inProgress = await db.kysely .selectFrom('routes_idempotency_requests') .select('apiKeyId') .where((eb) => eb.or([eb('apiKeyId', '!=', input.apiKeyId), eb('keyHash', '!=', input.keyHash)]), ) .where('matchHash', '=', input.matchHash) .where('operation', '=', operations.depositAddress) .where('orgId', '=', input.orgId) .where((eb) => inFlight(eb, now.toISOString())) .forUpdate() .executeTakeFirst() if (claimed.status === 'provisioned') return { type: 'provisioned' } if (!provisioned && !inProgress) return { type: 'reserved' } if (claimed.status === 'pending') await release(db, input) return { type: 'pending' } } export declare namespace reserveDepositAddressMatch { /** Claim and canonical match identity. */ type Input = release.Input & { /** SHA-256 fingerprint of the canonical reusable match. */ matchHash: string /** Organization owning the reusable match. */ orgId: string /** Renewed pending lease duration, in milliseconds. */ ttlMs: number } /** Match reservation, recovery, or ownership outcome. */ type Result = | { type: 'lost' } | { type: 'pending' } | { type: 'provisioned' } | { type: 'reserved' } | { checkpointFresh: boolean; claim: release.Input; response: string; type: 'resume' } } /** Returns whether provider provisioning may still create an address for an organization. */ export async function hasInFlightDepositAddress(db: Db.Db, orgId: string): Promise { const now = new Date().toISOString() return Boolean( await db.kysely .selectFrom('routes_idempotency_requests') .select('apiKeyId') .where((eb) => inFlight(eb, now)) .where('operation', '=', operations.depositAddress) .where('orgId', '=', orgId) .limit(1) .executeTakeFirst(), ) } /** Returns whether a rolling deployment left route work without attributable ownership. */ export async function hasUnattributedInFlightClaim(db: Db.Db): Promise { const now = new Date().toISOString() return Boolean( await db.kysely .selectFrom('routes_idempotency_requests') .select('apiKeyId') .where('status', 'in', ['pending', 'provisioned']) .where((eb) => eb.or([eb('operation', 'is', null), eb('orgId', 'is', null)])) .where((eb) => eb.or([eb('expiresAt', '>', now), eb('irrevocable', '=', true)])) .limit(1) .executeTakeFirst(), ) } /** Stores a provider result before committing its API resource and replay response. */ export function checkpoint(db: Db.Db, input: checkpoint.Input): Promise { const now = new Date() return db.kysely .updateTable('routes_idempotency_requests') .set({ expiresAt: new Date(now.getTime() + input.replayTtlMs).toISOString(), response: input.response, status: 'provisioned', }) .where('apiKeyId', '=', input.apiKeyId) .where('createdAt', '=', input.createdAt) .where('keyHash', '=', input.keyHash) .where('status', '=', 'pending') .returningAll() .executeTakeFirst() } export declare namespace checkpoint { /** Provider result and owned claim required for durable recovery. */ type Input = { /** API key id (`key_…`) the claim is scoped to. */ apiKeyId: string /** Creation time returned by the successful claim. */ createdAt: string /** SHA-256 of the caller Idempotency-Key. */ keyHash: string /** Recovery checkpoint retention in milliseconds. */ replayTtlMs: number /** Serialized provider result required to resume resource creation. */ response: string } } /** Completes an owned pending or provisioned claim with its replay response. */ export function complete(db: Db.Db, input: complete.Input): Promise { const now = new Date() return db.kysely .updateTable('routes_idempotency_requests') .set({ expiresAt: new Date(now.getTime() + input.replayTtlMs).toISOString(), irrevocable: false, response: input.response, status: 'completed', transferId: input.transferId ?? null, }) .where('apiKeyId', '=', input.apiKeyId) .where('createdAt', '=', input.createdAt) .where('keyHash', '=', input.keyHash) .where('status', 'in', ['pending', 'provisioned']) .returningAll() .executeTakeFirst() } export declare namespace complete { /** Completion fields. */ type Input = { /** API key id (`key_…`) the claim is scoped to. */ apiKeyId: string /** Creation time returned by the successful claim. */ createdAt: string /** SHA-256 of the caller Idempotency-Key. */ keyHash: string /** Successful-response replay retention from completion, in milliseconds. */ replayTtlMs: number /** Serialized success response replayed for the retention window. */ response: string /** Transfer id (`rtr_…`) when the request created a route transfer. */ transferId?: string | undefined } } /** Releases a pending claim after a retryable failure. */ export async function release(db: Db.Db, input: release.Input): Promise { await db.kysely .deleteFrom('routes_idempotency_requests') .where('apiKeyId', '=', input.apiKeyId) .where('createdAt', '=', input.createdAt) .where('keyHash', '=', input.keyHash) .where('status', '=', 'pending') .execute() } export declare namespace release { /** Pending claim identity. */ type Input = { /** API key id (`key_…`) the claim is scoped to. */ apiKeyId: string /** Creation time returned by the successful claim. */ createdAt: string /** SHA-256 of the caller Idempotency-Key. */ keyHash: string } } function inFlight(eb: ExpressionBuilder, now: string) { return eb.and([ eb('status', 'in', ['pending', 'provisioned']), eb.or([eb('expiresAt', '>', now), eb('irrevocable', '=', true)]), ]) }