import { C as ContrailConfig, D as Database, L as Logger, I as IngestEvent } from './types-CmjW-xL4.js'; import { AtprotoDid, AtprotoAudience } from '@atcute/lexicons/syntax'; import { S as ServiceOAuthScope } from './service-auth-contract-CER8k0ek.js'; declare const MAX_CHANGE_BATCH_CHANGES = 500; declare const MAX_CHANGE_BATCH_BYTES = 512000; interface ChangeLogCostPlan { enabled: boolean; consumers: number; coveragePairs: number; projectionStateWrites: number; changeHeadWrites: number; changeBatchWrites: number; acknowledgementWrites: number; /** Total expected rows written by one relevant projection transaction. */ relevantProjectionWrites: number; } /** Conservative write-amplification report for one bounded projection batch. */ declare function getChangeLogCostPlan(config: ContrailConfig, mutationUris?: number): ChangeLogCostPlan; interface RecordChange { id: string; kind: "record"; operation: "put" | "delete"; uri: string; did: string; collection: string; rkey: string; cid: string | null; version: { sourceId: string; sourceEpoch: string | null; sourceRevision: string | null; sourceTimeUs: number; sourceCursor: string | null; }; } interface ChangeLogState { generation: string; head: string; retainedFloor: string; createdAt: number; } declare function getChangeLogState(db: Database): Promise; interface CurrentBootstrapClaimOptions { pageSize?: number; leaseMs?: number; /** @internal Deterministic clock seam for conformance tests. */ now?: number; } interface CurrentSnapshotClaim { kind: "snapshot"; consumerId: string; generation: string; bootstrapToken: string; collection: string; fromUri: string | null; throughUri: string; pageId: string; records: CurrentRecord[]; attempt: number; leaseExpiresAt: number; readonly leaseOwner: string; } interface CurrentActivationClaim { kind: "activation"; consumerId: string; generation: string; bootstrapToken: string; target: string; attempt: number; leaseExpiresAt: number; readonly leaseOwner: string; } interface CurrentBootstrapStatus { consumerId: string; generation: string; state: string; anchor: string | null; scanCollection: string | null; scanCursor: string | null; target: string | null; token: string | null; position: string; } /** Claim one stable URI-keyset page. Empty collections advance internally; * exhausting the final collection atomically pins the catch-up target. */ declare function claimCurrentSnapshotPage(db: Database, config: ContrailConfig, consumerId: string, options?: CurrentBootstrapClaimOptions): Promise; declare function acknowledgeCurrentSnapshotPage(db: Database, claim: CurrentSnapshotClaim, options?: { now?: number; }): Promise; declare function failCurrentSnapshotPage(db: Database, claim: CurrentSnapshotClaim, failure: ChangeFailure, options?: { now?: number; }): Promise; declare function renewCurrentBootstrapClaim(db: Database, claim: T, options?: { leaseMs?: number; now?: number; }): Promise; /** Claim the idempotent destination-activation step after fixed-target catch-up. */ declare function claimCurrentActivation(db: Database, consumerId: string, options?: { leaseMs?: number; now?: number; }): Promise; declare function completeCurrentActivation(db: Database, claim: CurrentActivationClaim, options?: { now?: number; }): Promise; declare function failCurrentActivation(db: Database, claim: CurrentActivationClaim, failure: ChangeFailure, options?: { now?: number; }): Promise; declare function getCurrentBootstrapStatus(db: Database, consumerId: string): Promise; declare class ChangeConsumerNotFoundError extends Error { } declare class ChangeClaimTooLargeError extends Error { } declare class ChangeLeaseLostError extends Error { } declare class ChangeHistoryGapError extends Error { } declare class ChangeGenerationMismatchError extends Error { } interface ChangeClaimOptions { maxBatches?: number; maxChanges?: number; maxBytes?: number; leaseMs?: number; /** Maximum irrelevant position ranges core may acknowledge without invoking * a handler in one claim call. */ maxAutoAdvanceRanges?: number; /** @internal Deterministic clock seam for conformance tests. */ now?: number; } interface ChangeClaim { consumerId: string; generation: string; from: string; through: string; changes: RecordChange[]; attempt: number; leaseExpiresAt: number; /** Generation-scoped destination token during current-state catch-up. */ readonly bootstrapToken?: string; /** Fixed catch-up target for a current-state bootstrap claim. */ readonly bootstrapTarget?: string; /** Opaque CAS capability. Do not pass this field to delivery handlers. */ readonly leaseOwner: string; } interface CurrentRecord { uri: string; did: string; collection: string; rkey: string; cid: string | null; value: unknown; timeUs: number; indexedAt: number; } interface DeliveryBatch { consumerId: string; cursor: { generation: string; from: string; through: string; }; changes: RecordChange[]; currentRecords: CurrentRecord[]; absentUris: string[]; /** Generation-scoped destination token during current bootstrap catch-up. */ destinationToken?: string; } interface ChangeFailure { /** Stable sanitized category. Raw destination errors are never persisted. */ code: string; /** Runtime-computed retry eligibility. Null makes the claim immediately due. */ nextAttemptAt: number | null; } interface ChangeConsumerStatus { id: string; generation: string; position: string; bootstrapState: string; initialMode: string; requiredForActivation: boolean; attempts: number; nextAttemptAt: number | null; lastSuccessAt: number | null; lastErrorCode: string | null; lastErrorAt: number | null; leased: boolean; leaseExpiresAt: number | null; backlogBatches: number; backlogChanges: number; backlogBytes: number; oldestPendingAt: number | null; generationMatches: boolean; } interface ChangeLogStatus { enabled: boolean; state: ChangeLogState | null; rows: number; changes: number; bytes: number; oldestRetainedAt: number | null; consumers: ChangeConsumerStatus[]; } /** Verify that schema initialization registered one configured consumer. */ declare function registerChangeConsumer(db: Database, config: ContrailConfig, consumerId: string): Promise; declare function claimChanges(db: Database, consumerId: string, options?: ChangeClaimOptions): Promise; /** Claim changes only through the fixed target of a current-state bootstrap. */ declare function claimCurrentBootstrapChanges(db: Database, consumerId: string, options?: ChangeClaimOptions): Promise; /** Hydrate the newest canonical state for each coalesced claimed URI. */ declare function hydrateChanges(db: Database, config: ContrailConfig, claim: ChangeClaim): Promise; declare function acknowledgeChanges(db: Database, claim: ChangeClaim, options?: { now?: number; }): Promise; declare function renewChangeClaim(db: Database, claim: ChangeClaim, options?: { leaseMs?: number; now?: number; }): Promise; declare function failChanges(db: Database, claim: ChangeClaim, failure: ChangeFailure, options?: { now?: number; }): Promise<{ attempts: number; nextAttemptAt: number | null; }>; declare function retryChangeConsumer(db: Database, consumerId: string, options?: { now?: number; }): Promise; declare function getChangesStatus(db: Database): Promise; interface RequiredChangeConsumerReadiness { ready: boolean; through: string; pending: Array<{ id: string; state: string; position: string; }>; } /** Check activation-gating consumers against one fixed log target. */ declare function getRequiredChangeConsumerReadiness(db: Database, through?: string): Promise; declare function assertRequiredChangeConsumersReady(db: Database, through?: string): Promise; interface SkipChangeConsumerOptions { through: string; reason: string; confirm: boolean; /** @internal Deterministic clock seam for conformance tests. */ now?: number; } /** Explicit audited data-loss operation for one ready consumer. */ declare function skipChangeConsumer(db: Database, consumerId: string, options: SkipChangeConsumerOptions): Promise; interface PruneChangesOptions { maxBatches?: number; /** Delete only batches older than this timestamp. */ olderThan?: number; } interface PruneChangesResult { pruned: number; retainedFloor: string; safeThrough: string; done: boolean; } /** Consumer-aware bounded pruning. The slowest durable position, including a * current bootstrap anchor, is the hard safety boundary. */ declare function pruneChanges(db: Database, options?: PruneChangesOptions): Promise; type DatabaseResolver = (db?: Database) => Database; /** Bound low-level API exposed as `contrail.changes`. */ declare class ChangeConsumers { private readonly config; private readonly database; constructor(config: ContrailConfig, database: DatabaseResolver); register(consumerId: string, db?: Database): Promise; claim(consumerId: string, options?: ChangeClaimOptions, db?: Database): Promise; claimBootstrapChanges(consumerId: string, options?: ChangeClaimOptions, db?: Database): Promise; claimSnapshotPage(consumerId: string, options?: CurrentBootstrapClaimOptions, db?: Database): Promise; ackSnapshotPage(claim: CurrentSnapshotClaim, options?: { now?: number; }, db?: Database): Promise; failSnapshotPage(claim: CurrentSnapshotClaim, failure: ChangeFailure, options?: { now?: number; }, db?: Database): Promise; claimActivation(consumerId: string, options?: { leaseMs?: number; now?: number; }, db?: Database): Promise; completeActivation(claim: CurrentActivationClaim, options?: { now?: number; }, db?: Database): Promise; renewBootstrap(claim: T, options?: { leaseMs?: number; now?: number; }, db?: Database): Promise; failActivation(claim: CurrentActivationClaim, failure: ChangeFailure, options?: { now?: number; }, db?: Database): Promise; bootstrapStatus(consumerId: string, db?: Database): Promise; hydrate(claim: ChangeClaim, db?: Database): Promise; ack(claim: ChangeClaim, options?: { now?: number; }, db?: Database): Promise; renew(claim: ChangeClaim, options?: { leaseMs?: number; now?: number; }, db?: Database): Promise; fail(claim: ChangeClaim, failure: ChangeFailure, options?: { now?: number; }, db?: Database): Promise<{ attempts: number; nextAttemptAt: number | null; }>; retry(consumerId: string, options?: { now?: number; }, db?: Database): Promise; status(db?: Database): Promise; readiness(through?: string, db?: Database): Promise; skip(consumerId: string, options: SkipChangeConsumerOptions, db?: Database): Promise; prune(options?: PruneChangesOptions, db?: Database): Promise; } interface DeliveryContext { env: Env; signal: AbortSignal; attempt: number; } type DeliveryHandler = (batch: DeliveryBatch, context: DeliveryContext) => Promise; type SnapshotDeliveryPage = Omit; type ActivationDelivery = Omit; interface CurrentBootstrapRuntimeHandler { snapshot: (page: SnapshotDeliveryPage, context: DeliveryContext) => Promise; activate: (activation: ActivationDelivery, context: DeliveryContext) => Promise; } type DeliveryHandlers = Record>; type CurrentBootstrapRuntimeHandlers = Record>; 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; } interface DeliverySliceResult { steps: number; delivered: number; snapshotPages: number; activations: number; failures: number; consumerErrors: Record; deadlineReached: boolean; } declare function validateDeliveryHandlers(config: ContrailConfig, deliveries: DeliveryHandlers, bootstraps?: CurrentBootstrapRuntimeHandlers): void; /** Run fair round-robin delivery work without coupling failures to ingestion. */ declare function runChangeDeliverySlice(options: { changes: ChangeConsumers; config: ContrailConfig; db: Database; env: Env; deliveries: DeliveryHandlers; bootstraps?: CurrentBootstrapRuntimeHandlers; runtime?: DeliveryRuntimeOptions; }): Promise; /** Persistent fair supervisor. Source ingestion remains a separate task and is * never cancelled by a destination outage. */ declare function runPersistentChangeDeliveries(options: { changes: ChangeConsumers; config: ContrailConfig; db: Database; env: Env; deliveries: DeliveryHandlers; bootstraps?: CurrentBootstrapRuntimeHandlers; runtime?: DeliveryRuntimeOptions & { idleMs?: number; }; }): Promise; /** Fixed estimate for the URI, source position, revision, CID, and other * metadata retained beside each serialized record body. Deletes consume only * this allowance. The byte budget is intentionally an admission threshold: the * candidate that reaches it is retained, bounding overshoot to one record plus * this allowance. */ declare const SCHEDULED_INGEST_METADATA_BYTES = 512; /** Conservative defaults for one D1/Worker scheduled drain. Persistent * ingestion has its own streaming lifecycle and does not use these limits. */ declare const DEFAULT_SCHEDULED_INGEST_BUDGET: Readonly<{ maxDrainMs: 25000; maxCandidates: 250; maxIdentityUpdates: 250; maxSerializedBytes: number; }>; interface ScheduledIngestBudget { maxDrainMs: number; maxCandidates: number; maxIdentityUpdates: number; maxSerializedBytes: number; } interface ScheduledIngestOptions { /** Maximum wall time spent requesting source items. Default: 25 seconds. */ maxDrainMs?: number; /** Maximum unique commit candidates retained by one drain. Default: 250. */ maxCandidates?: number; /** Maximum distinct handle updates retained by one drain. Default: 250. */ maxIdentityUpdates?: number; /** UTF-8 record bytes plus metadata allowances. Default: 4 MiB. */ maxSerializedBytes?: number; /** @deprecated Compatibility alias for maxDrainMs. */ timeoutMs?: number; } type ScheduledIngestStopReason = "head" | "idle" | "count" | "bytes" | "drain-time" | "cancelled"; interface ScheduledIngestCollectionStats { observedSourceItems: number; commitObservations: number; identityObservations: number; retainedCandidates: number; retainedIdentityUpdates: number; identityUpdatesOmitted: number; exactDuplicatesDropped: number; cursorBoundaryDuplicatesDropped: number; resumeOverlapDropped: number; sourceScopeFiltered: number; sourceInconsistencies: number; serializedCandidateBytes: number; startingCursor: number | null; lastAccountedCursor: number | null; safeEndingCursor: number | null; stopReason: ScheduledIngestStopReason; connections: number; connectionCloses: number; connectionErrors: number; diagnosticSamples: string[]; diagnosticSamplesOmitted: number; } /** Resolve and validate every scheduled collection threshold. */ declare function resolveScheduledIngestBudget(value?: ScheduledIngestOptions | ScheduledIngestBudget | number): ScheduledIngestBudget; /** Distinct actors pruned per ingest tick by the rolling feed sweep. Each * actor costs a handful of index-backed O(cap) deletes, so this bounds the * prune's per-tick CPU regardless of how large feed_items grows. */ declare const FEED_PRUNE_SWEEP_ACTORS = 500; /** How long after a completed full pass the recovery sweep becomes due again, * even when no feed-relevant records are ingested, so over-cap rows that * predate a config change (e.g. a lowered cap) or a bulk import still drain. * Once due, the pass advances one slice per tick, so a full pass *completes* * roughly every (this interval + lap time), where lap time is * ceil(actors / FEED_PRUNE_SWEEP_ACTORS) ticks. Steady-state pruning is driven * by ingest; this is only the safety net. */ declare const FEED_PRUNE_RECOVERY_INTERVAL_MS: number; /** `_contrail_meta` key for the persisted optimize cadence (so recycled cron * isolates don't re-run it every tick — the in-memory-state bug we hit with * the feed prune). Shared by the persistent loop. */ declare const OPTIMIZE_LAST_MS_KEY = "optimize_last_ms"; /** Run the opt-in planner-stat maintenance if enabled and its persisted * interval has elapsed. Bounded + no-op on Postgres (see optimizeDatabase). * Wrapped by callers so a pragma-unsupported environment can't break ingest. */ declare function maybeOptimize(db: Database, config: ContrailConfig, log: Logger): Promise; /** One bounded feed-prune slice: advances the persisted rolling cursor by up to * {@link FEED_PRUNE_SWEEP_ACTORS} actors and reports whether the slice reached * the end of the actor list (i.e. a full pass just completed and the cursor * wrapped). Callers decide WHEN to sweep (ingest-dirty vs recovery); this owns * the slice + cursor mechanics so the cron loop, the persistent loop, and the * notify path all prune identically. No-op (done) when no feed caps apply. */ declare function runFeedPruneSlice(db: Database, config: ContrailConfig): Promise<{ pruned: number; done: boolean; }>; /** Gate and run one feed-prune slice against the *persisted* recovery clock — * shared by the recycling cron isolate and the stateless `notifyOfUpdate` path * (the long-lived persistent loop uses its in-memory clocks instead). Slices * when `feedTouched` (a feed-mutating record was just ingested) or when a full * pass is overdue, and records pass completion so the recovery clock measures * from a real drain (one slice per tick until the cursor wraps) rather than * resetting on a single slice. * * The slice advances the shared rolling cursor, which is NOT necessarily the * actor the mutation touched: a fan-out follower the cursor has already passed * is pruned by the next pass, up to about one recovery interval later, not * instantly. That is * the deliberate trade for a per-tick cost bounded by `FEED_PRUNE_SWEEP_ACTORS` * rather than by fan-out size (a popular author has unboundedly many followers). * feed_items is a soft cache, so a follower sitting a few rows over cap until * the next slice is harmless. No-op when feeds are unconfigured. */ declare function runGatedFeedPrune(db: Database, config: ContrailConfig, feedTouched: boolean): Promise; /** Mutable state that persists across ingest cycles within the same process. */ interface IngestState { cachedKnownDids?: Set; schemaInitialized: boolean; /** Wall-clock of the last feed sweep slice — used only by the long-lived * persistent loop to throttle ingest-driven slices; the recycling cron * isolate persists its clocks in `_contrail_meta` instead. */ lastFeedSweepMs: number; /** Wall-clock at which the persistent loop last *completed a full sweep pass* * over every actor. Drives the recovery interval (a fresh pass becomes due * {@link FEED_PRUNE_RECOVERY_INTERVAL_MS} after the last one completed, then * laps one slice per tick), independent of the ingest-driven throttle above. * The cron isolate persists the equivalent in * `_contrail_meta`. */ lastFullFeedPassMs: number; /** Set by the persistent loop when a flushed batch ingested a feed-mutating * record, so the next sweep window knows there may be prune work. Cleared * when the sweep runs. The cron path makes the same decision per-tick from * its `events` array and doesn't need the flag. */ feedDirty: boolean; } declare function createIngestState(): IngestState; declare function ingestEvents(config: ContrailConfig, cursor: number | null, budgetInput?: ScheduledIngestOptions | ScheduledIngestBudget | number, knownDids?: Set): Promise<{ events: IngestEvent[]; lastCursor: number | null; cursorObservations: Set; identityUpdates: Map; stats: ScheduledIngestCollectionStats; }>; declare function runIngestCycle(db: Database, config: ContrailConfig, budgetInput?: ScheduledIngestOptions | ScheduledIngestBudget | number, state?: IngestState): Promise; interface PublicServiceOptions { /** Canonical public HTTPS origin, for example `https://api.example.com`. */ endpoint: string; /** Permit plain HTTP only on a loopback host for local development. */ allowInsecureHttp?: boolean; } interface PublicServiceCollection { alias: string; nsid: string; methods: string[]; queryable: string[]; searchable: string[]; relations: string[]; references: string[]; } interface PublicServiceProtectedMethod { id: string; type: "query" | "procedure"; } interface PublicServiceAuthContract { type: "atproto-service-auth"; serviceDid: AtprotoDid; audience: AtprotoAudience; scope: ServiceOAuthScope; methods: PublicServiceProtectedMethod[]; } interface PublicServiceManifest { format: "contrail.service"; version: 2; endpoint: string; namespace: string; lexicons: { url: string; digest: string; }; status: { url: string; }; collections: PublicServiceCollection[]; methods: string[]; serviceAuth?: PublicServiceAuthContract | null; } interface LexiconDocument { lexicon?: number; id: string; [key: string]: unknown; } interface PublicServiceDescription { endpoint: string; lexicons: LexiconDocument[]; manifest: PublicServiceManifest; canonicalLexicons: string; } declare function normalizePublicServiceEndpoint(value: string, options?: { allowInsecureHttp?: boolean; }): string; /** Whether this endpoint is the origin that hosts the base DID's document, and * therefore whether Contrail should publish one at `/.well-known/did.json`. */ declare function hostsServiceDidDocument(serviceDid: AtprotoDid, endpoint: string): boolean; declare function validatePublicServiceAuthEndpoint(config: ContrailConfig, options: PublicServiceOptions): void; declare function canonicalJson(value: unknown): string; declare function sha256(value: string): Promise; declare function normalizeLexiconDocuments(values: readonly object[]): LexiconDocument[]; declare function validatePublicServiceLexicons(config: ContrailConfig, values: readonly object[]): LexiconDocument[]; declare function digestLexiconDocuments(values: readonly object[]): Promise<{ lexicons: LexiconDocument[]; canonicalLexicons: string; digest: string; }>; declare function describePublicService(config: ContrailConfig, options: PublicServiceOptions, values: readonly object[]): Promise; declare function validateServiceManifest(manifest: PublicServiceManifest, values: readonly object[]): LexiconDocument[]; declare function isPublicServiceAuthContract(value: unknown): value is PublicServiceAuthContract | null | undefined; declare function isPublicServiceManifest(value: unknown): value is PublicServiceManifest; export { acknowledgeCurrentSnapshotPage as $, type ActivationDelivery as A, MAX_CHANGE_BATCH_CHANGES as B, type CurrentBootstrapRuntimeHandlers as C, type DeliveryHandlers as D, type PruneChangesOptions as E, FEED_PRUNE_RECOVERY_INTERVAL_MS as F, type PruneChangesResult as G, type PublicServiceAuthContract as H, type IngestState as I, type PublicServiceCollection as J, type PublicServiceDescription as K, type LexiconDocument as L, MAX_CHANGE_BATCH_BYTES as M, type PublicServiceManifest as N, OPTIMIZE_LAST_MS_KEY as O, type PublicServiceOptions as P, type PublicServiceProtectedMethod as Q, type RecordChange as R, type ScheduledIngestOptions as S, type RequiredChangeConsumerReadiness as T, SCHEDULED_INGEST_METADATA_BYTES as U, type ScheduledIngestBudget as V, type ScheduledIngestCollectionStats as W, type ScheduledIngestStopReason as X, type SkipChangeConsumerOptions as Y, type SnapshotDeliveryPage as Z, acknowledgeChanges as _, type DeliveryRuntimeOptions as a, assertRequiredChangeConsumersReady as a0, canonicalJson as a1, claimChanges as a2, claimCurrentActivation as a3, claimCurrentBootstrapChanges as a4, claimCurrentSnapshotPage as a5, completeCurrentActivation as a6, createIngestState as a7, describePublicService as a8, digestLexiconDocuments as a9, runPersistentChangeDeliveries as aA, sha256 as aB, skipChangeConsumer as aC, validateDeliveryHandlers as aD, validatePublicServiceAuthEndpoint as aE, validatePublicServiceLexicons as aF, validateServiceManifest as aG, failChanges as aa, failCurrentActivation as ab, failCurrentSnapshotPage as ac, getChangeLogCostPlan as ad, getChangeLogState as ae, getChangesStatus as af, getCurrentBootstrapStatus as ag, getRequiredChangeConsumerReadiness as ah, hostsServiceDidDocument as ai, hydrateChanges as aj, ingestEvents as ak, isPublicServiceAuthContract as al, isPublicServiceManifest as am, maybeOptimize as an, normalizeLexiconDocuments as ao, normalizePublicServiceEndpoint as ap, pruneChanges as aq, registerChangeConsumer as ar, renewChangeClaim as as, renewCurrentBootstrapClaim as at, resolveScheduledIngestBudget as au, retryChangeConsumer as av, runChangeDeliverySlice as aw, runFeedPruneSlice as ax, runGatedFeedPrune as ay, runIngestCycle as az, ChangeConsumers as b, type ChangeClaim as c, type ChangeClaimOptions as d, ChangeClaimTooLargeError as e, ChangeConsumerNotFoundError as f, type ChangeConsumerStatus as g, type ChangeFailure as h, ChangeGenerationMismatchError as i, ChangeHistoryGapError as j, ChangeLeaseLostError as k, type ChangeLogCostPlan as l, type ChangeLogState as m, type ChangeLogStatus as n, type CurrentActivationClaim as o, type CurrentBootstrapClaimOptions as p, type CurrentBootstrapRuntimeHandler as q, type CurrentBootstrapStatus as r, type CurrentRecord as s, type CurrentSnapshotClaim as t, DEFAULT_SCHEDULED_INGEST_BUDGET as u, type DeliveryBatch as v, type DeliveryContext as w, type DeliveryHandler as x, type DeliverySliceResult as y, FEED_PRUNE_SWEEP_ACTORS as z };