export type RuntimePostgresLane = 'receipts' | 'sheets'; export type RuntimePostgresLaneTelemetry = { active: number; queued: number; acquired: number; waited: number; timedOut: number; rejected: number; cancelled: number; totalWaitMs: number; maxWaitMs: number; }; export type RuntimePostgresAdmissionSnapshot = Record< RuntimePostgresLane, RuntimePostgresLaneTelemetry >; export class RuntimePostgresAdmissionTimeoutError extends Error { constructor( readonly lane: RuntimePostgresLane, readonly queueDepth: number, readonly timeoutMs: number, readonly active: number, readonly queuedBeforeTimeout: number, readonly waitedMs: number, ) { super( `Runtime Postgres ${lane} admission timed out after ${timeoutMs}ms ` + `(${active} active, ${queuedBeforeTimeout} queued before timeout; ${queueDepth} queued after removal).`, ); this.name = 'RuntimePostgresAdmissionTimeoutError'; } } export class RuntimePostgresAdmissionQueueFullError extends Error { constructor( readonly lane: RuntimePostgresLane, readonly queueDepth: number, ) { super( `Runtime Postgres ${lane} admission queue is full (${queueDepth} queued).`, ); this.name = 'RuntimePostgresAdmissionQueueFullError'; } } export class RuntimePostgresAdmissionCancelledError extends Error { constructor(readonly lane: RuntimePostgresLane) { super(`Runtime Postgres ${lane} admission was cancelled.`); this.name = 'RuntimePostgresAdmissionCancelledError'; } } export function isRuntimePostgresAdmissionCapacityError( error: unknown, ): error is | RuntimePostgresAdmissionTimeoutError | RuntimePostgresAdmissionQueueFullError { return ( error instanceof RuntimePostgresAdmissionTimeoutError || error instanceof RuntimePostgresAdmissionQueueFullError ); } export function isRuntimePostgresAdmissionError( error: unknown, ): error is | RuntimePostgresAdmissionTimeoutError | RuntimePostgresAdmissionQueueFullError | RuntimePostgresAdmissionCancelledError { return ( isRuntimePostgresAdmissionCapacityError(error) || error instanceof RuntimePostgresAdmissionCancelledError ); } type Waiter = { enqueuedAt: number; timeout: ReturnType | null; resolve: (release: () => void) => void; reject: (error: Error) => void; signal: AbortSignal | null; onAbort: (() => void) | null; settled: boolean; }; type LaneState = RuntimePostgresLaneTelemetry & { waiters: Waiter[] }; function createLaneState(): LaneState { return { active: 0, queued: 0, acquired: 0, waited: 0, timedOut: 0, rejected: 0, cancelled: 0, totalWaitMs: 0, maxWaitMs: 0, waiters: [], }; } export class RuntimePostgresAdmission { readonly #maxActivePerLane: number; readonly #maxQueuedPerLane: Record; readonly #acquireTimeoutMs: Record; readonly #lanes: Record = { receipts: createLaneState(), sheets: createLaneState(), }; constructor(input: { maxActivePerLane: number; maxQueuedPerLane: number | Record; acquireTimeoutMs: | number | null | Record; }) { this.#maxActivePerLane = Math.max(1, Math.floor(input.maxActivePerLane)); const normalizeQueueLimit = (value: number): number => Math.max(0, Math.floor(value)); this.#maxQueuedPerLane = typeof input.maxQueuedPerLane === 'object' ? { receipts: normalizeQueueLimit(input.maxQueuedPerLane.receipts), sheets: normalizeQueueLimit(input.maxQueuedPerLane.sheets), } : { receipts: normalizeQueueLimit(input.maxQueuedPerLane), sheets: normalizeQueueLimit(input.maxQueuedPerLane), }; const normalizeTimeout = (value: number | null): number | null => value === null ? null : Math.max(1, Math.floor(value)); const acquireTimeoutMs = input.acquireTimeoutMs; this.#acquireTimeoutMs = acquireTimeoutMs !== null && typeof acquireTimeoutMs === 'object' ? { receipts: normalizeTimeout(acquireTimeoutMs.receipts), sheets: normalizeTimeout(acquireTimeoutMs.sheets), } : { receipts: normalizeTimeout( typeof acquireTimeoutMs === 'number' ? acquireTimeoutMs : null, ), sheets: normalizeTimeout( typeof acquireTimeoutMs === 'number' ? acquireTimeoutMs : null, ), }; } async acquire( lane: RuntimePostgresLane, options: { signal?: AbortSignal } = {}, ): Promise<() => void> { const state = this.#lanes[lane]; if (options.signal?.aborted) { state.cancelled += 1; throw new RuntimePostgresAdmissionCancelledError(lane); } if (state.active < this.#maxActivePerLane && state.waiters.length === 0) { state.active += 1; state.acquired += 1; return this.#releaseFor(lane); } if (state.waiters.length >= this.#maxQueuedPerLane[lane]) { state.rejected += 1; throw new RuntimePostgresAdmissionQueueFullError( lane, state.waiters.length, ); } state.waited += 1; const acquireTimeoutMs = this.#acquireTimeoutMs[lane]; return await new Promise<() => void>((resolve, reject) => { const waiter: Waiter = { enqueuedAt: Date.now(), timeout: null, resolve, reject, signal: options.signal ?? null, onAbort: null, settled: false, }; const removeWaiter = (): boolean => { if (waiter.settled) return false; const index = state.waiters.indexOf(waiter); if (index < 0) return false; waiter.settled = true; state.waiters.splice(index, 1); state.queued = state.waiters.length; if (waiter.timeout) clearTimeout(waiter.timeout); if (waiter.signal && waiter.onAbort) { waiter.signal.removeEventListener('abort', waiter.onAbort); } return true; }; waiter.onAbort = () => { if (!removeWaiter()) return; state.cancelled += 1; reject(new RuntimePostgresAdmissionCancelledError(lane)); }; waiter.signal?.addEventListener('abort', waiter.onAbort, { once: true }); if (acquireTimeoutMs !== null) { waiter.timeout = setTimeout(() => { const queuedBeforeTimeout = state.waiters.length; if (!removeWaiter()) return; state.timedOut += 1; reject( new RuntimePostgresAdmissionTimeoutError( lane, state.waiters.length, acquireTimeoutMs, state.active, queuedBeforeTimeout, Math.max(0, Date.now() - waiter.enqueuedAt), ), ); }, acquireTimeoutMs); } state.waiters.push(waiter); state.queued = state.waiters.length; if (waiter.signal?.aborted) waiter.onAbort(); }); } snapshot(): RuntimePostgresAdmissionSnapshot { const snapshotLane = (state: LaneState): RuntimePostgresLaneTelemetry => ({ active: state.active, queued: state.waiters.length, acquired: state.acquired, waited: state.waited, timedOut: state.timedOut, rejected: state.rejected, cancelled: state.cancelled, totalWaitMs: state.totalWaitMs, maxWaitMs: state.maxWaitMs, }); return { receipts: snapshotLane(this.#lanes.receipts), sheets: snapshotLane(this.#lanes.sheets), }; } #releaseFor(lane: RuntimePostgresLane): () => void { let released = false; return () => { if (released) return; released = true; const state = this.#lanes[lane]; state.active -= 1; const waiter = state.waiters.shift(); state.queued = state.waiters.length; if (!waiter) return; waiter.settled = true; if (waiter.timeout) clearTimeout(waiter.timeout); if (waiter.signal && waiter.onAbort) { waiter.signal.removeEventListener('abort', waiter.onAbort); } const waitMs = Math.max(0, Date.now() - waiter.enqueuedAt); state.totalWaitMs += waitMs; state.maxWaitMs = Math.max(state.maxWaitMs, waitMs); state.active += 1; state.acquired += 1; waiter.resolve(this.#releaseFor(lane)); }; } }