diff --git a/.changeset/durable-projection-log.md b/.changeset/durable-projection-log.md index 27dd926..619eb38 100644 --- a/.changeset/durable-projection-log.md +++ b/.changeset/durable-projection-log.md @@ -4,4 +4,4 @@ Add the first transactional projection change-log milestone. Optional static consumer definitions now create a fresh-generation log, durable registrations, and collection/phase coverage ledger. Winning logical URI changes append compact references atomically with canonical records, derived projections, tombstones, and source checkpoints; disabled configurations create no log tables or append writes. -Harden all projection writers with transaction-time predecessor guards and bounded conflict retries so overlapping cron, persistent, notify, and backfill work cannot commit stale canonical or derived state. Add independent bounded consumer leases, filtered/coalesced claims, set-oriented current-state hydration, CAS acknowledgement, failure backoff, lease renewal, private status, and manual retry APIs. Add crash-safe current-state snapshot/tail/activation bootstrap, safe additive consumers over existing coverage, required-consumer readiness gates, consumer-aware bounded pruning, and audited explicit skip operations. Private `contrail changes` commands cover status, retry, prune, and skip. Enabling or expanding log coverage on a populated generation fails closed pending explicit quiet-boundary migration tooling. +Harden all projection writers with transaction-time predecessor guards and bounded conflict retries so overlapping cron, persistent, notify, and backfill work cannot commit stale canonical or derived state. Add independent bounded consumer leases, filtered/coalesced claims, set-oriented current-state hydration, CAS acknowledgement, failure backoff, lease renewal, private status, and manual retry APIs. Add crash-safe current-state snapshot/tail/activation bootstrap, safe additive consumers over existing coverage, required-consumer readiness gates, consumer-aware bounded pruning, and audited explicit skip operations. Private `contrail changes` commands cover status, retry, prune, and skip. Add fair bounded Worker delivery after ingestion/retries, best-effort immediate notify wakes, runtime handler validation, deadline cancellation, isolated retry scheduling, and a persistent delivery supervisor. Enabling or expanding log coverage on a populated generation fails closed pending explicit quiet-boundary migration tooling. diff --git a/packages/contrail/README.md b/packages/contrail/README.md index 09f1052..6eac71f 100644 --- a/packages/contrail/README.md +++ b/packages/contrail/README.md @@ -96,6 +96,34 @@ const config = { Static definitions contain no handlers, URLs, clients, credentials, or secrets. Contrail registers them with a random database-generation ID and collection/phase coverage ledger. A winning logical put/delete appends one compact URI/version reference in the same transaction as canonical and derived state plus the source checkpoint. Duplicate, stale, same-CID, absent-delete, rejected, and rolled-back mutations append nothing. Record bodies are hydrated from current state by the later delivery layer rather than copied into the log. +Cloudflare Workers bind handlers and runtime-only secrets separately from static policy: + +```ts +export default createWorker(config, { + deliveries: { + search: async (batch, { env, signal }) => { + await updateSearch(env.SEARCH_KEY, batch, signal); + }, + }, + changeBootstraps: { + search: { + snapshot: async (page, { env, signal }) => { + await writeCandidateIndex(env.SEARCH_KEY, page, signal); + }, + activate: async (activation, { env, signal }) => { + await activateCandidateIdempotently( + env.SEARCH_KEY, + activation.bootstrapToken, + signal, + ); + }, + }, + }, +}); +``` + +Scheduled execution runs ingestion, due historical retries, then bounded fair delivery rounds. One consumer failure is persisted with backoff and does not stop ingestion or another consumer. A successful notify request schedules a best-effort one-round wake through `ExecutionContext`; projection success never depends on it. Persistent deployments run `contrail.runPersistentDeliveries()` alongside `runPersistent()`. + Low-level delivery uses bounded leases and compare-and-swap acknowledgement: ```ts diff --git a/packages/contrail/src/contrail.ts b/packages/contrail/src/contrail.ts index 859d054..8cccb88 100644 --- a/packages/contrail/src/contrail.ts +++ b/packages/contrail/src/contrail.ts @@ -15,6 +15,12 @@ import { } from "./core/validation"; import { getIngestDiagnostics } from "./core/diagnostics"; import { ChangeConsumers } from "./core/changes"; +import { + runPersistentChangeDeliveries, + type CurrentBootstrapRuntimeHandlers, + type DeliveryHandlers, + type DeliveryRuntimeOptions, +} from "./core/delivery"; import { optimizeDatabase } from "./core/db/optimize"; import { assertServingSourceCompatibility, @@ -160,6 +166,25 @@ export class Contrail { await Promise.all(tasks); } + /** Run a persistent fair change-delivery supervisor. Run this alongside + * `runPersistent()`; destination failures never stop source ingestion. */ + async runPersistentDeliveries( + options: { + env: Env; + deliveries: DeliveryHandlers; + bootstraps?: CurrentBootstrapRuntimeHandlers; + runtime?: DeliveryRuntimeOptions & { idleMs?: number }; + }, + db?: Database, + ): Promise { + await runPersistentChangeDeliveries({ + changes: this.changes, + config: this.config, + db: this.getDb(db), + ...options, + }); + } + /** Run *only* the labeler ingestion cycle. Escape hatch for callers who * want to run record and label ingestion in separate processes / workers. * `ingest()` already covers the typical case. */ diff --git a/packages/contrail/src/core/delivery.ts b/packages/contrail/src/core/delivery.ts new file mode 100644 index 0000000..65a2283 --- /dev/null +++ b/packages/contrail/src/core/delivery.ts @@ -0,0 +1,503 @@ +import type { ContrailConfig, Database, Logger } from "./types"; +import type { + ChangeClaimOptions, + ChangeConsumers, + DeliveryBatch, +} from "./changes"; +import type { + CurrentActivationClaim, + CurrentSnapshotClaim, +} from "./change-bootstrap"; + +export interface DeliveryContext { + env: Env; + signal: AbortSignal; + attempt: number; +} + +export type DeliveryHandler = ( + batch: DeliveryBatch, + context: DeliveryContext, +) => Promise; + +export type SnapshotDeliveryPage = Omit; +export type ActivationDelivery = Omit; + +export interface CurrentBootstrapRuntimeHandler { + snapshot: ( + page: SnapshotDeliveryPage, + context: DeliveryContext, + ) => Promise; + activate: ( + activation: ActivationDelivery, + context: DeliveryContext, + ) => Promise; +} + +export type DeliveryHandlers = Record>; +export type CurrentBootstrapRuntimeHandlers = Record< + string, + CurrentBootstrapRuntimeHandler +>; + +export interface DeliveryRuntimeOptions { + maxRounds?: number; + maxDurationMs?: number; + claim?: ChangeClaimOptions; + baseRetryMs?: number; + maxRetryMs?: number; + jitter?: number; + signal?: AbortSignal; + logger?: Logger; + /** @internal Deterministic runtime seams. */ + clock?: () => number; + random?: () => number; +} + +export interface DeliverySliceResult { + steps: number; + delivered: number; + snapshotPages: number; + activations: number; + failures: number; + consumerErrors: Record; + deadlineReached: boolean; +} + +function boundedInteger( + value: number | undefined, + fallback: number, + minimum: number, + maximum: number, + label: string, +): number { + const result = value ?? fallback; + if (!Number.isSafeInteger(result) || result < minimum || result > maximum) { + throw new TypeError(`${label} must be an integer between ${minimum} and ${maximum}`); + } + return result; +} + +export function validateDeliveryHandlers( + config: ContrailConfig, + deliveries: DeliveryHandlers, + bootstraps: CurrentBootstrapRuntimeHandlers = {}, +): void { + const consumers = config.changes?.consumers ?? {}; + for (const id of Object.keys(consumers)) { + if (typeof deliveries[id] !== "function") { + throw new Error(`Missing runtime delivery handler for change consumer ${id}`); + } + if ( + consumers[id]!.initial === "current" && + (!bootstraps[id] || + typeof bootstraps[id].snapshot !== "function" || + typeof bootstraps[id].activate !== "function") + ) { + throw new Error( + `Missing current-state bootstrap handlers for change consumer ${id}`, + ); + } + } + for (const id of Object.keys(deliveries)) { + if (!consumers[id]) { + throw new Error(`Runtime delivery handler ${id} has no static consumer definition`); + } + } + for (const id of Object.keys(bootstraps)) { + if (!consumers[id] || consumers[id]!.initial !== "current") { + throw new Error(`Runtime bootstrap handler ${id} has no current consumer definition`); + } + } +} + +function publicSnapshot(claim: CurrentSnapshotClaim): SnapshotDeliveryPage { + const { leaseOwner: _leaseOwner, ...page } = claim; + return page; +} + +function publicActivation(claim: CurrentActivationClaim): ActivationDelivery { + const { leaseOwner: _leaseOwner, ...activation } = claim; + return activation; +} + +async function withDeadline( + parent: AbortSignal | undefined, + deadline: number, + clock: () => number, + callback: (signal: AbortSignal) => Promise, +): Promise { + const controller = new AbortController(); + const abort = () => controller.abort(parent?.reason ?? new Error("Delivery cancelled")); + if (parent?.aborted) abort(); + else parent?.addEventListener("abort", abort, { once: true }); + const remaining = Math.max(0, deadline - clock()); + const timer = setTimeout( + () => controller.abort(new Error("Delivery deadline reached")), + remaining, + ); + try { + return await callback(controller.signal); + } finally { + clearTimeout(timer); + parent?.removeEventListener("abort", abort); + } +} + +function retryAt( + attempt: number, + now: number, + base: number, + maximum: number, + jitter: number, + random: () => number, +): number { + const exponential = Math.min(maximum, base * 2 ** Math.min(20, attempt - 1)); + const factor = 1 + (random() * 2 - 1) * jitter; + return now + Math.max(1, Math.round(exponential * factor)); +} + +interface RuntimeState { + changes: ChangeConsumers; + config: ContrailConfig; + db: Database; + env: Env; + deliveries: DeliveryHandlers; + bootstraps: CurrentBootstrapRuntimeHandlers; + claim: ChangeClaimOptions; + deadline: number; + baseRetryMs: number; + maxRetryMs: number; + jitter: number; + signal?: AbortSignal; + logger: Logger; + clock: () => number; + random: () => number; +} + +async function persistFailure( + state: RuntimeState, + consumerId: string, + claim: Parameters[0], + attempt: number, +): Promise { + const now = state.clock(); + try { + await state.changes.fail( + claim, + { + code: "handler_error", + nextAttemptAt: retryAt( + attempt, + now, + state.baseRetryMs, + state.maxRetryMs, + state.jitter, + state.random, + ), + }, + { now }, + state.db, + ); + } catch (error) { + state.logger.warn( + `[changes] consumer=${consumerId} could not persist failure: ${error}`, + ); + } +} + +async function runNormal( + state: RuntimeState, + consumerId: string, + bootstrap: boolean, +): Promise<"empty" | "delivered" | "failed"> { + const now = state.clock(); + const claimOptions = { ...state.claim, now }; + const claim = bootstrap + ? await state.changes.claimBootstrapChanges( + consumerId, + claimOptions, + state.db, + ) + : await state.changes.claim(consumerId, claimOptions, state.db); + if (!claim) return "empty"; + try { + const batch = await state.changes.hydrate(claim, state.db); + await withDeadline(state.signal, state.deadline, state.clock, (signal) => + state.deliveries[consumerId]!(batch, { + env: state.env, + signal, + attempt: claim.attempt, + }), + ); + await state.changes.ack( + claim, + { now: state.clock() }, + state.db, + ); + return "delivered"; + } catch (error) { + state.logger.warn(`[changes] consumer=${consumerId} delivery failed: ${error}`); + await persistFailure( + state as RuntimeState, + consumerId, + claim, + claim.attempt, + ); + return "failed"; + } +} + +async function runCurrent( + state: RuntimeState, + consumerId: string, +): Promise<"empty" | "delivered" | "snapshot" | "activation" | "failed" | "progressed"> { + const before = await state.changes.bootstrapStatus(consumerId, state.db); + if (before.state === "ready") { + return runNormal(state, consumerId, false); + } + if (before.state === "pending" || before.state === "scanning") { + const claim = await state.changes.claimSnapshotPage( + consumerId, + { ...state.claim, now: state.clock() }, + state.db, + ); + if (!claim) { + const after = await state.changes.bootstrapStatus(consumerId, state.db); + return after.state !== before.state ? "progressed" : "empty"; + } + try { + await withDeadline(state.signal, state.deadline, state.clock, (signal) => + state.bootstraps[consumerId]!.snapshot(publicSnapshot(claim), { + env: state.env, + signal, + attempt: claim.attempt, + }), + ); + await state.changes.ackSnapshotPage( + claim, + { now: state.clock() }, + state.db, + ); + return "snapshot"; + } catch (error) { + state.logger.warn(`[changes] consumer=${consumerId} snapshot failed: ${error}`); + try { + await state.changes.failSnapshotPage( + claim, + { + code: "snapshot_handler_error", + nextAttemptAt: retryAt( + claim.attempt, + state.clock(), + state.baseRetryMs, + state.maxRetryMs, + state.jitter, + state.random, + ), + }, + { now: state.clock() }, + state.db, + ); + } catch (failureError) { + state.logger.warn(`[changes] consumer=${consumerId} snapshot failure state lost: ${failureError}`); + } + return "failed"; + } + } + if (before.state === "catching-up") { + const result = await runNormal(state, consumerId, true); + if (result !== "empty") return result; + const after = await state.changes.bootstrapStatus(consumerId, state.db); + return after.state !== before.state ? "progressed" : "empty"; + } + if (before.state === "activating") { + const claim = await state.changes.claimActivation( + consumerId, + { now: state.clock() }, + state.db, + ); + if (!claim) return "empty"; + try { + await withDeadline(state.signal, state.deadline, state.clock, (signal) => + state.bootstraps[consumerId]!.activate(publicActivation(claim), { + env: state.env, + signal, + attempt: claim.attempt, + }), + ); + await state.changes.completeActivation( + claim, + { now: state.clock() }, + state.db, + ); + return "activation"; + } catch (error) { + state.logger.warn(`[changes] consumer=${consumerId} activation failed: ${error}`); + try { + await state.changes.failActivation( + claim, + { + code: "activation_handler_error", + nextAttemptAt: retryAt( + claim.attempt, + state.clock(), + state.baseRetryMs, + state.maxRetryMs, + state.jitter, + state.random, + ), + }, + { now: state.clock() }, + state.db, + ); + } catch (failureError) { + state.logger.warn(`[changes] consumer=${consumerId} activation failure state lost: ${failureError}`); + } + return "failed"; + } + } + return "empty"; +} + +/** Run fair round-robin delivery work without coupling failures to ingestion. */ +export async function runChangeDeliverySlice(options: { + changes: ChangeConsumers; + config: ContrailConfig; + db: Database; + env: Env; + deliveries: DeliveryHandlers; + bootstraps?: CurrentBootstrapRuntimeHandlers; + runtime?: DeliveryRuntimeOptions; +}): Promise { + const bootstraps = options.bootstraps ?? {}; + validateDeliveryHandlers(options.config, options.deliveries, bootstraps); + const runtime = options.runtime ?? {}; + const clock = runtime.clock ?? Date.now; + const maxRounds = boundedInteger(runtime.maxRounds, 4, 1, 100, "maxRounds"); + const maxDurationMs = boundedInteger( + runtime.maxDurationMs, + 15_000, + 1, + 10 * 60_000, + "maxDurationMs", + ); + const baseRetryMs = boundedInteger( + runtime.baseRetryMs, + 1_000, + 1, + 60 * 60_000, + "baseRetryMs", + ); + const maxRetryMs = boundedInteger( + runtime.maxRetryMs, + 60 * 60_000, + baseRetryMs, + 48 * 60 * 60_000, + "maxRetryMs", + ); + const jitter = runtime.jitter ?? 0.2; + if (!Number.isFinite(jitter) || jitter < 0 || jitter > 1) { + throw new TypeError("jitter must be between 0 and 1"); + } + const deadline = clock() + maxDurationMs; + const state: RuntimeState = { + changes: options.changes, + config: options.config, + db: options.db, + env: options.env, + deliveries: options.deliveries, + bootstraps, + claim: runtime.claim ?? {}, + deadline, + baseRetryMs, + maxRetryMs, + jitter, + signal: runtime.signal, + logger: runtime.logger ?? options.config.logger ?? console, + clock, + random: runtime.random ?? Math.random, + }; + const result: DeliverySliceResult = { + steps: 0, + delivered: 0, + snapshotPages: 0, + activations: 0, + failures: 0, + consumerErrors: {}, + deadlineReached: false, + }; + const consumers = Object.entries(options.config.changes?.consumers ?? {}).sort( + ([left], [right]) => left.localeCompare(right), + ); + + for (let round = 0; round < maxRounds; round++) { + let progressed = false; + for (const [consumerId, consumer] of consumers) { + if (runtime.signal?.aborted || clock() >= deadline) { + result.deadlineReached = true; + return result; + } + try { + const outcome = consumer.initial === "current" + ? await runCurrent(state, consumerId) + : await runNormal(state, consumerId, false); + if (outcome !== "empty") { + progressed = true; + result.steps++; + } + if (outcome === "delivered") result.delivered++; + else if (outcome === "snapshot") result.snapshotPages++; + else if (outcome === "activation") result.activations++; + else if (outcome === "failed") result.failures++; + } catch (error) { + result.failures++; + result.consumerErrors[consumerId] = + error instanceof Error ? error.message : String(error); + state.logger.error(`[changes] consumer=${consumerId} runtime failed: ${error}`); + } + } + if (!progressed) break; + } + result.deadlineReached = clock() >= deadline; + return result; +} + +function waitFor(ms: number, signal?: AbortSignal): Promise { + if (signal?.aborted) return Promise.resolve(); + return new Promise((resolve) => { + const timer = setTimeout(done, ms); + function done() { + clearTimeout(timer); + signal?.removeEventListener("abort", done); + resolve(); + } + signal?.addEventListener("abort", done, { once: true }); + }); +} + +/** Persistent fair supervisor. Source ingestion remains a separate task and is + * never cancelled by a destination outage. */ +export async function runPersistentChangeDeliveries(options: { + changes: ChangeConsumers; + config: ContrailConfig; + db: Database; + env: Env; + deliveries: DeliveryHandlers; + bootstraps?: CurrentBootstrapRuntimeHandlers; + runtime?: DeliveryRuntimeOptions & { idleMs?: number }; +}): Promise { + const idleMs = boundedInteger( + options.runtime?.idleMs, + 1_000, + 1, + 60_000, + "idleMs", + ); + while (!options.runtime?.signal?.aborted) { + const result = await runChangeDeliverySlice(options); + if (result.steps === 0) { + await waitFor(idleMs, options.runtime?.signal); + } + } +} diff --git a/packages/contrail/src/core/router/diagnostics.ts b/packages/contrail/src/core/router/diagnostics.ts index 5a7927d..a080011 100644 --- a/packages/contrail/src/core/router/diagnostics.ts +++ b/packages/contrail/src/core/router/diagnostics.ts @@ -3,6 +3,7 @@ import type { ContrailConfig, Database } from "../types"; import { getCollectionShortNames, recordsTableName, nsidForShortName } from "../types"; import { getLastCursor, getServingSourcePosition } from "../db"; import { getBackfillStatus } from "../status"; +import { getRequiredChangeConsumerReadiness } from "../changes"; export interface CursorStatus { cursor: number | null; @@ -48,9 +49,12 @@ export async function getStatusOverview(db: Database, config: ContrailConfig) { } } - const [ingestion, backfill] = await Promise.all([ + const [ingestion, backfill, requiredDelivery] = await Promise.all([ getCursorStatus(db), getBackfillStatus(db, config), + config.changes && Object.keys(config.changes.consumers).length > 0 + ? getRequiredChangeConsumerReadiness(db) + : Promise.resolve({ ready: true, through: "0", pending: [] }), ]); return { @@ -59,6 +63,11 @@ export async function getStatusOverview(db: Database, config: ContrailConfig) { collections, ingestion, backfill, + delivery: { + required: requiredDelivery.ready ? "ready" as const : "catching_up" as const, + pending: requiredDelivery.pending.length, + through: requiredDelivery.through, + }, }; } diff --git a/packages/contrail/src/core/router/index.ts b/packages/contrail/src/core/router/index.ts index a4e9720..8f66072 100644 --- a/packages/contrail/src/core/router/index.ts +++ b/packages/contrail/src/core/router/index.ts @@ -53,6 +53,9 @@ export function createApp( const overview = await getStatusOverview(db, config); if (!options.publicService) return c.json(overview); c.header("cache-control", "public, max-age=15, stale-while-revalidate=45"); + const hasRequiredDelivery = Object.values( + config.changes?.consumers ?? {}, + ).some((consumer) => consumer.requiredForActivation === true); return c.json({ status: overview.status, serving: "ready", @@ -63,6 +66,9 @@ export function createApp( seconds_ago: overview.ingestion.seconds_ago, }, backfill: overview.backfill, + ...(hasRequiredDelivery + ? { required_delivery: overview.delivery.required } + : {}), }); }); app.get("/health", (c) => c.json({ status: "ok" })); diff --git a/packages/contrail/src/index.ts b/packages/contrail/src/index.ts index 36b6e5d..dfd17b6 100644 --- a/packages/contrail/src/index.ts +++ b/packages/contrail/src/index.ts @@ -37,6 +37,7 @@ export { export type { ChangeLogState, RecordChange } from "./core/change-log"; export * from "./core/changes"; export * from "./core/change-bootstrap"; +export * from "./core/delivery"; export * from "./core/validation"; export * from "./core/search"; export * from "./core/constellation"; diff --git a/packages/contrail/src/worker/index.ts b/packages/contrail/src/worker/index.ts index a75f0a0..d38ba04 100644 --- a/packages/contrail/src/worker/index.ts +++ b/packages/contrail/src/worker/index.ts @@ -18,6 +18,13 @@ import { Contrail } from "../contrail.js"; import { createHandler } from "../server.js"; import type { ContrailConfig, Database } from "../core/types.js"; import type { BackfillRetryOptions } from "../core/backfill.js"; +import { + runChangeDeliverySlice, + validateDeliveryHandlers, + type CurrentBootstrapRuntimeHandlers, + type DeliveryHandlers, + type DeliveryRuntimeOptions, +} from "../core/delivery.js"; import { normalizePublicServiceEndpoint, validatePublicServiceAuthEndpoint, @@ -25,7 +32,9 @@ import { type PublicServiceOptions, } from "../public-service.js"; -export interface CreateWorkerOptions { +type WorkerEnv = Record; + +export interface CreateWorkerOptions { /** D1 binding name in wrangler env. Default: `"DB"`. */ binding?: string; /** Exact generated/pinned bundle exposed for type generation and used by @@ -36,18 +45,33 @@ export interface CreateWorkerOptions { /** Bounded pending-account retry slice after each scheduled ingest. Enabled * by default; pass `false` to disable or options to tune its budget. */ backfillRetries?: BackfillRetryOptions | false; + /** Runtime delivery handlers, kept separate from static consumer policy. */ + deliveries?: DeliveryHandlers; + /** Snapshot and activation handlers for `initial: "current"` consumers. */ + changeBootstraps?: CurrentBootstrapRuntimeHandlers; + /** Bounded scheduled delivery policy. Set false only when another runtime + * owns all configured consumers. */ + delivery?: DeliveryRuntimeOptions | false; /** Runs once per isolate, after schema init, before handling the first * request. Use for app-specific setup that needs a live DB handle. */ - onInit?: (env: Record, db: Database) => void | Promise; + onInit?: (env: Env, db: Database) => void | Promise; } -type WorkerEnv = Record; - -export function createWorker( +export function createWorker( config: ContrailConfig, - options: CreateWorkerOptions = {} + options: CreateWorkerOptions = {} ) { const binding = options.binding ?? "DB"; + const deliveryEnabled = + options.delivery !== false && + Object.keys(config.changes?.consumers ?? {}).length > 0; + if (deliveryEnabled) { + validateDeliveryHandlers( + config, + options.deliveries ?? {}, + options.changeBootstraps ?? {}, + ); + } if (options.publicService) { normalizePublicServiceEndpoint(options.publicService.endpoint); validatePublicServiceLexicons(config, options.lexicons ?? []); @@ -60,7 +84,7 @@ export function createWorker( }); let ready = false; - const ensureReady = async (env: WorkerEnv, db: Database): Promise => { + const ensureReady = async (env: Env, db: Database): Promise => { if (ready) return; await contrail.init(db); await options.onInit?.(env, db); @@ -68,14 +92,52 @@ export function createWorker( }; return { - async fetch(request: Request, env: WorkerEnv): Promise { + async fetch( + request: Request, + env: Env, + ctx?: ExecutionContext, + ): Promise { const db = env[binding] as Database; await ensureReady(env, db); - return (await handle(request, db)) as Response; + const response = (await handle(request, db)) as Response; + const notifyPath = `/xrpc/${contrail.config.namespace}.notifyOfUpdate`; + if ( + deliveryEnabled && + ctx && + response.ok && + request.method === "POST" && + new URL(request.url).pathname === notifyPath + ) { + ctx.waitUntil( + runChangeDeliverySlice({ + changes: contrail.changes, + config: contrail.config, + db, + env, + deliveries: options.deliveries!, + bootstraps: options.changeBootstraps, + runtime: { + ...(options.delivery || {}), + maxRounds: 1, + maxDurationMs: Math.min( + options.delivery && options.delivery.maxDurationMs + ? options.delivery.maxDurationMs + : 5_000, + 5_000, + ), + }, + }).catch((error) => { + contrail.config.logger?.error( + `[changes] immediate delivery wake failed: ${error}`, + ); + }), + ); + } + return response; }, async scheduled( _event: ScheduledEvent, - env: WorkerEnv, + env: Env, ctx: ExecutionContext ): Promise { const db = env[binding] as Database; @@ -84,9 +146,36 @@ export function createWorker( // failures. A database lease prevents overlap with a manual backfill. ctx.waitUntil( (async () => { - await contrail.ingest({}, db); + try { + await contrail.ingest({}, db); + } catch (error) { + contrail.config.logger?.error(`[ingest] scheduled cycle failed: ${error}`); + } if (options.backfillRetries !== false) { - await contrail.retryBackfill(options.backfillRetries, db); + try { + await contrail.retryBackfill(options.backfillRetries, db); + } catch (error) { + contrail.config.logger?.error( + `[backfill] scheduled retry slice failed: ${error}`, + ); + } + } + if (deliveryEnabled) { + try { + await runChangeDeliverySlice({ + changes: contrail.changes, + config: contrail.config, + db, + env, + deliveries: options.deliveries!, + bootstraps: options.changeBootstraps, + runtime: options.delivery || undefined, + }); + } catch (error) { + contrail.config.logger?.error( + `[changes] scheduled delivery slice failed: ${error}`, + ); + } } })() ); diff --git a/packages/contrail/tests/delivery.test.ts b/packages/contrail/tests/delivery.test.ts new file mode 100644 index 0000000..2835fe5 --- /dev/null +++ b/packages/contrail/tests/delivery.test.ts @@ -0,0 +1,333 @@ +import { describe, expect, it } from "vitest"; +import { + createIngestEvent, + getChangesStatus, + ingestRecords, + initSchema, + resolveConfig, + runChangeDeliverySlice, + runPersistentChangeDeliveries, + validateDeliveryHandlers, + type ChangeConsumerConfig, + type Database, + type DeliveryHandlers, +} from "../src/index"; +import { createSqliteDatabase } from "../src/adapters/sqlite"; +import { Contrail } from "../src/contrail"; + +const EVENT = "com.example.event"; +const logger = { log() {}, warn() {}, error() {} }; + +function config(consumers: Record) { + return resolveConfig({ + namespace: "com.example", + profiles: [], + logger, + collections: { event: { collection: EVENT } }, + changes: { consumers }, + }); +} + +function event(rkey: string, time: number) { + const did = "did:plc:alice"; + return createIngestEvent({ + uri: `at://${did}/${EVENT}/${rkey}`, + did, + collection: EVENT, + rkey, + operation: "update", + cid: `cid-${rkey}-${time}`, + value: { name: rkey }, + timeUs: time, + indexedAt: time + 10_000, + source: { + id: "source", + epoch: "epoch", + time_us: time, + revision: String(time), + cursor: String(time), + }, + }); +} + +async function append(db: Database, resolved: ReturnType, rkey: string, time: number) { + await ingestRecords(db, [event(rkey, time)], resolved); +} + +describe("change delivery runtime", () => { + it("fails startup for missing, extra, or incomplete runtime handlers", () => { + const resolved = config({ + search: { collections: [EVENT], initial: "current" }, + }); + expect(() => validateDeliveryHandlers(resolved, {})).toThrow( + "Missing runtime delivery handler", + ); + expect(() => + validateDeliveryHandlers(resolved, { search: async () => {} }), + ).toThrow("Missing current-state bootstrap handlers"); + expect(() => + validateDeliveryHandlers( + resolved, + { search: async () => {}, extra: async () => {} }, + { + search: { + snapshot: async () => {}, + activate: async () => {}, + }, + }, + ), + ).toThrow("extra has no static consumer definition"); + }); + + it("runs one bounded claim per consumer per fair round", async () => { + const db = createSqliteDatabase(":memory:"); + const resolved = config({ + alpha: { collections: [EVENT], initial: "history" }, + beta: { collections: [EVENT], initial: "history" }, + gamma: { collections: [EVENT], initial: "history" }, + }); + const contrail = new Contrail({ ...resolved, db }); + await contrail.init(); + await append(db, resolved, "one", 1); + await append(db, resolved, "two", 2); + + const order: string[] = []; + const handlers: DeliveryHandlers<{}> = Object.fromEntries( + ["alpha", "beta", "gamma"].map((id) => [ + id, + async () => { + order.push(id); + }, + ]), + ); + const result = await runChangeDeliverySlice({ + changes: contrail.changes, + config: resolved, + db, + env: {}, + deliveries: handlers, + runtime: { + maxRounds: 2, + claim: { maxBatches: 1 }, + jitter: 0, + }, + }); + expect(order).toEqual([ + "alpha", + "beta", + "gamma", + "alpha", + "beta", + "gamma", + ]); + expect(result).toMatchObject({ + delivered: 6, + failures: 0, + steps: 6, + }); + }); + + it("isolates one failing consumer and resumes it after persisted backoff", async () => { + const db = createSqliteDatabase(":memory:"); + const resolved = config({ + broken: { collections: [EVENT], initial: "history" }, + healthy: { collections: [EVENT], initial: "history" }, + }); + const contrail = new Contrail({ ...resolved, db }); + await contrail.init(); + await append(db, resolved, "one", 1); + + let fail = true; + const handlers = { + broken: async () => { + if (fail) throw new Error("destination down"); + }, + healthy: async () => {}, + }; + let current = 1_000; + const first = await runChangeDeliverySlice({ + changes: contrail.changes, + config: resolved, + db, + env: {}, + deliveries: handlers, + runtime: { + maxRounds: 1, + baseRetryMs: 100, + maxRetryMs: 100, + jitter: 0, + clock: () => current, + }, + }); + expect(first).toMatchObject({ delivered: 1, failures: 1 }); + let status = await getChangesStatus(db); + expect(status.consumers.find((item) => item.id === "broken")).toMatchObject({ + position: "0", + attempts: 1, + nextAttemptAt: 1_100, + }); + expect(status.consumers.find((item) => item.id === "healthy")).toMatchObject({ + position: "1", + }); + + fail = false; + current = 1_101; + const resumed = await runChangeDeliverySlice({ + changes: contrail.changes, + config: resolved, + db, + env: {}, + deliveries: handlers, + runtime: { + maxRounds: 1, + jitter: 0, + clock: () => current, + }, + }); + expect(resumed.delivered).toBe(1); + status = await getChangesStatus(db); + expect(status.consumers.find((item) => item.id === "broken")).toMatchObject({ + position: "1", + attempts: 0, + }); + }); + + it("drives current snapshot, catch-up, and idempotent activation", async () => { + const db = createSqliteDatabase(":memory:"); + const resolved = config({ + search: { + collections: [EVENT], + initial: "current", + requiredForActivation: true, + }, + }); + const contrail = new Contrail({ ...resolved, db }); + await contrail.init(); + await append(db, resolved, "one", 1); + expect( + await ( + await contrail.app().fetch(new Request("http://localhost/status")) + ).json(), + ).toMatchObject({ + delivery: { required: "catching_up", pending: 1 }, + }); + + const calls: string[] = []; + const result = await runChangeDeliverySlice({ + changes: contrail.changes, + config: resolved, + db, + env: { secret: "runtime-only" }, + deliveries: { + search: async (batch, context) => { + calls.push(`tail:${batch.cursor.through}:${context.env.secret}`); + }, + }, + bootstraps: { + search: { + snapshot: async (page) => { + expect(page).not.toHaveProperty("leaseOwner"); + calls.push(`snapshot:${page.records.length}`); + }, + activate: async (activation) => { + expect(activation).not.toHaveProperty("leaseOwner"); + calls.push(`activate:${activation.target}`); + }, + }, + }, + runtime: { maxRounds: 6, jitter: 0 }, + }); + expect(calls).toEqual([ + "snapshot:1", + "tail:1:runtime-only", + "activate:1", + ]); + expect(result).toMatchObject({ + snapshotPages: 1, + delivered: 1, + activations: 1, + failures: 0, + }); + expect(await contrail.changes.bootstrapStatus("search")).toMatchObject({ + state: "ready", + position: "1", + }); + expect( + await ( + await contrail.app().fetch(new Request("http://localhost/status")) + ).json(), + ).toMatchObject({ + delivery: { required: "ready", pending: 0 }, + }); + }); + + it("aborts destination work at the runtime deadline", async () => { + const db = createSqliteDatabase(":memory:"); + const resolved = config({ + webhook: { collections: [EVENT], initial: "history" }, + }); + const contrail = new Contrail({ ...resolved, db }); + await contrail.init(); + await append(db, resolved, "one", 1); + let aborted = false; + const result = await runChangeDeliverySlice({ + changes: contrail.changes, + config: resolved, + db, + env: {}, + deliveries: { + webhook: async (_batch, { signal }) => { + await new Promise((resolve) => { + if (signal.aborted) { + aborted = true; + resolve(); + return; + } + signal.addEventListener( + "abort", + () => { + aborted = true; + resolve(); + }, + { once: true }, + ); + }); + }, + }, + runtime: { maxRounds: 1, maxDurationMs: 5 }, + }); + expect(aborted).toBe(true); + expect(result.deadlineReached).toBe(true); + }); + + it("runs a persistent supervisor until cancellation", async () => { + const db = createSqliteDatabase(":memory:"); + const resolved = config({ + webhook: { collections: [EVENT], initial: "history" }, + }); + const contrail = new Contrail({ ...resolved, db }); + await contrail.init(); + await append(db, resolved, "one", 1); + const controller = new AbortController(); + let deliveries = 0; + await runPersistentChangeDeliveries({ + changes: contrail.changes, + config: resolved, + db, + env: {}, + deliveries: { + webhook: async () => { + deliveries++; + controller.abort(); + }, + }, + runtime: { + signal: controller.signal, + idleMs: 1, + maxRounds: 1, + }, + }); + expect(deliveries).toBe(1); + expect((await getChangesStatus(db)).consumers[0].position).toBe("1"); + }); +}); diff --git a/packages/contrail/tests/worker.test.ts b/packages/contrail/tests/worker.test.ts index bd69244..49cedd6 100644 --- a/packages/contrail/tests/worker.test.ts +++ b/packages/contrail/tests/worker.test.ts @@ -2,7 +2,13 @@ import { describe, it, expect, vi } from "vitest"; import { createWorker } from "../src/worker"; import { Contrail } from "../src/contrail"; import { createSqliteDatabase } from "../src/adapters/sqlite"; -import { saveCursor, type ContrailConfig } from "../src/index"; +import { + createIngestEvent, + ingestRecords, + resolveConfig, + saveCursor, + type ContrailConfig, +} from "../src/index"; const MINIMAL_CONFIG: ContrailConfig = { namespace: "com.example", @@ -71,6 +77,24 @@ describe("createWorker", () => { ).not.toThrow(); }); + it("rejects configured consumers without matching runtime handlers", () => { + const config: ContrailConfig = { + ...MINIMAL_CONFIG, + changes: { + consumers: { + webhook: { + collections: ["community.lexicon.calendar.event"], + initial: "history", + }, + }, + }, + }; + expect(() => createWorker(config)).toThrow( + "Missing runtime delivery handler", + ); + expect(() => createWorker(config, { delivery: false })).not.toThrow(); + }); + it("returns an object with fetch + scheduled handlers", () => { const worker = createWorker(MINIMAL_CONFIG); expect(typeof worker.fetch).toBe("function"); @@ -497,6 +521,151 @@ describe("createWorker", () => { ).toThrow("public method requires a matching query Lexicon"); }); + it("scheduled handler isolates ingest failure before retry and delivery", async () => { + const order: string[] = []; + const ingest = vi + .spyOn(Contrail.prototype, "ingest") + .mockImplementation(async () => { + order.push("ingest"); + throw new Error("source unavailable"); + }); + const retry = vi + .spyOn(Contrail.prototype, "retryBackfill") + .mockImplementation(async () => { + order.push("retry"); + return { + attempted: 0, + completed: 0, + failed: 0, + records: 0, + skipped: false, + }; + }); + const config: ContrailConfig = { + namespace: "com.example", + profiles: [], + logger: { log() {}, warn() {}, error() {} }, + collections: { + event: { collection: "community.lexicon.calendar.event" }, + }, + changes: { + consumers: { + webhook: { + collections: ["community.lexicon.calendar.event"], + initial: "history", + }, + }, + }, + }; + + try { + const db = createSqliteDatabase(":memory:"); + const worker = createWorker(config, { + deliveries: { + webhook: async () => { + order.push("delivery"); + }, + }, + }); + const env = { DB: db }; + await worker.fetch(new Request("http://localhost/health"), env); + const resolved = resolveConfig(config); + await ingestRecords( + db, + [ + createIngestEvent({ + did: "did:plc:alice", + collection: "community.lexicon.calendar.event", + rkey: "one", + operation: "create", + cid: "cid-one", + value: { name: "one" }, + timeUs: 1, + }), + ], + resolved, + ); + const waitUntil = vi.fn(); + const ctx = { + waitUntil, + passThroughOnException: vi.fn(), + } as unknown as ExecutionContext; + + await worker.scheduled({} as ScheduledEvent, env, ctx); + await waitUntil.mock.calls[0][0]; + expect(order).toEqual(["ingest", "retry", "delivery"]); + } finally { + ingest.mockRestore(); + retry.mockRestore(); + } + }); + + it("schedules a best-effort delivery wake after successful notify", async () => { + const config: ContrailConfig = { + namespace: "com.example", + profiles: [], + notify: true, + collections: { + event: { collection: "community.lexicon.calendar.event" }, + }, + changes: { + consumers: { + webhook: { + collections: ["community.lexicon.calendar.event"], + initial: "history", + }, + }, + }, + }; + const delivered = vi.fn(async () => {}); + const db = createSqliteDatabase(":memory:"); + const worker = createWorker(config, { + deliveries: { webhook: delivered }, + }); + const env = { DB: db }; + await worker.fetch(new Request("http://localhost/health"), env); + await ingestRecords( + db, + [ + createIngestEvent({ + did: "did:plc:alice", + collection: "community.lexicon.calendar.event", + rkey: "one", + operation: "create", + cid: "cid-one", + value: { name: "one" }, + timeUs: 1, + }), + ], + resolveConfig(config), + ); + + const handler = vi + .spyOn(Contrail.prototype, "handler") + .mockReturnValue(async () => Response.json({ ok: true })); + try { + const waitUntil = vi.fn(); + const ctx = { + waitUntil, + passThroughOnException: vi.fn(), + } as unknown as ExecutionContext; + const response = await worker.fetch( + new Request("http://localhost/xrpc/com.example.notifyOfUpdate", { + method: "POST", + body: JSON.stringify({ uri: "at://did:plc:alice/community.lexicon.calendar.event/one" }), + }), + env, + ctx, + ); + expect(response.status).toBe(200); + expect(waitUntil).toHaveBeenCalledTimes(1); + await waitUntil.mock.calls[0][0]; + expect(delivered).toHaveBeenCalledTimes(1); + } finally { + handler.mockRestore(); + } + }); + it("scheduled handler runs live ingest then a bounded backfill retry slice", async () => { const order: string[] = []; const ingest = vi