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)]),
])
}