import type * as RoutesObservability from '../../internal/routes/Observability.js' import type * as Analytics from '../Analytics.js' import type * as RequestEvents from './requestEvents.js' /** Columns of `routes_lifecycle_intervals`: one completed routes stage interval. */ export type Table = { /** UTC time when the stage ended. */ ended_at: string /** Deployment environment that recorded the interval. */ environment: string /** API-key environment that owns the route journey. */ routes_environment: 'production' | 'sandbox' /** Route method used by the journey. */ method: 'deposit_address' | 'transaction' /** Route provider that owns the journey. */ provider: string /** UTC time when the interval was recorded. */ recorded_at: string /** Stable route resource id. */ resource_id: string /** Service that recorded the interval. */ service: string /** Completed lifecycle stage. */ stage: RoutesObservability.LifecycleInterval['stage'] /** UTC time when the stage began. */ started_at: string } /** Discriminated queue event carrying one route lifecycle interval. */ export type Event = { /** Stored interval row. */ data: Table /** Queue event discriminator. */ type: 'routes:lifecycle:interval' } /** Builds a queue sink for completed route lifecycle intervals. */ export function createQueueSink( queue: createQueueSink.Queue, options: createQueueSink.Options, ): createQueueSink.Sink { return async (interval) => { await queue.send({ data: { ended_at: timestamp(interval.endedAt), environment: options.environment, routes_environment: interval.routesEnvironment, method: interval.method, provider: interval.provider, recorded_at: timestamp(new Date().toISOString()), resource_id: interval.resourceId, service: options.service, stage: interval.stage, started_at: timestamp(interval.startedAt), }, type: 'routes:lifecycle:interval', }) } } export declare namespace createQueueSink { /** Context added to every stored interval. */ type Options = { /** Deployment environment that records intervals. */ environment: string /** Service that records intervals. */ service: string } /** Minimal queue producer used by the interval sink. */ type Queue = { /** Enqueues one lifecycle interval. */ send(message: Event): Promise } /** Records one completed route lifecycle interval. */ type Sink = (interval: RoutesObservability.LifecycleInterval) => Promise } /** Returns whether a queue body is a route lifecycle interval event. */ export function isEvent(value: unknown): value is Event { return ( typeof value === 'object' && value !== null && 'type' in value && value.type === 'routes:lifecycle:interval' ) } /** Inserts queued interval rows and applies per-message queue disposition. */ export async function insertMessages( analytics: Analytics.Analytics, messages: readonly insertMessages.Message[], options: insertMessages.Options = {}, ): Promise { if (messages.length === 0) return try { await analytics.insert( 'routes_lifecycle_intervals', messages.map((message) => message.body.data), ) for (const message of messages) message.ack() report({ messages: messages.length, outcome: 'acked' }) } catch (cause) { console.error('ClickHouse route lifecycle insert failed', cause) for (const message of messages) message.retry() report({ cause, messages: messages.length, outcome: 'retried' }) } function report(result: RequestEvents.insertMessages.Result) { try { options.onResult?.(result) } catch (cause) { console.error('Analytics queue result callback failed', cause) } } } export declare namespace insertMessages { /** Minimal queue message carrying one lifecycle event. */ type Message = { /** Acknowledges the message. */ ack(): void /** Queued lifecycle event. */ body: Event /** Marks the message for redelivery. */ retry(): void } /** Queue insert observability options. */ type Options = { /** Receives the final insert disposition. */ onResult?: ((result: RequestEvents.insertMessages.Result) => void) | undefined } } function timestamp(value: string) { return new Date(value).toISOString().replace('T', ' ').replace('Z', '') }