/** * Async-promote queue worker supervisor. * * PR #3 of the async-promote-queue series. PR #1 shipped the queue * library; PR #2 exposed it to the agent + HTTP routes. This module * drains the queue: an in-process queue signal wakes N worker slots, * with periodic polling retained for durable cross-process writes and * retry deadlines. A claimed slot calls `agent.assertion.promote(...)`, * then records success or classified failure back into the queue. * * Open questions from the plan §10 are resolved here: * * 1. **Worker model**: in-process enqueue/recover signals wake all available * slots immediately. One periodic poll remains as a durable fallback for * cross-process writes and retry eligibility — the daemon still owns the * recovery clock. * * 2. **Backoff curve**: defined in `async-promote-queue-utils.ts`, * inherited by the queue itself; the worker just calls `fail()` with * a classification and the queue handles backoff bookkeeping. * * 3. **Error classification**: see `async-promote-error-classification.ts`. * Seeded from the rc.10 Graphify import patterns (see * `INTEGRATION_NOTES_GRAPHIFY.md` and `dkg-graphify-rc10-test/FINDINGS_v2.md`). * * 4. **Telemetry**: the supervisor emits `memoryGraphChanged` on every * `succeeded` transition that promoted >0 triples, mirroring the * sync `/promote` route. State transitions to `queued`/`running`/ * `failed_retrying` are NOT emitted — they're queue internals. * * Shutdown semantics: RFC §6.2 says do NOT mark `running → queued` * on shutdown. We stop polling, let in-flight jobs complete (or * timeout), and rely on `recoverOnStartup()` at the next boot to * decide what to do with any leases the old worker held. */ import type { DKGAgent } from '@origintrail-official/dkg-agent'; import { type AsyncPromoteQueue, type PromoteJob, type PromoteRequest } from '@origintrail-official/dkg-publisher'; import { type ClassifiedPromoteError } from './async-promote-error-classification.js'; export { classifyPromoteError } from './async-promote-error-classification.js'; export type { ClassifiedPromoteError } from './async-promote-error-classification.js'; /** * Convenience type for the daemon's existing `emitMemoryGraphChanged` * callback. Kept local so this module doesn't have to import the * `MemoryGraphChangedEvent` shape from `routes/context.ts` and pull in * its full route surface. */ export interface PromoteMemoryGraphChangedEvent { contextGraphId: string; layers: ('wm' | 'swm')[]; subGraphName?: string; operation: string; source: string; counts?: { triples?: number; }; } /** * Logging is strictly best-effort: sinks may be synchronous or asynchronous, * but worker progress never waits for them and sink failures are discarded. */ export type PromoteWorkerLogger = (message: string) => void | Promise; /** * Internal logger shape: normalization makes every worker call fire-and-forget. * Worker code holds only this type and calls it directly. A public * `PromoteWorkerLogger` is assignable to it, so the rule that keeps a raw sink * out is structural: normalize at the entry point, pass nothing else down. */ export type PromoteWorkerSyncLogger = (message: string) => void; /** * Normalize a public logger once at the worker boundary. The worker's internal * paths can then emit diagnostics without each call having to know whether the * configured sink is synchronous, asynchronous, or hostile. */ export declare function normalizePromoteWorkerLogger(configured: PromoteWorkerLogger | undefined): PromoteWorkerSyncLogger; export interface PromoteWorkerConfig { /** The host DKG agent — provides the queue + the sync `promote` call. */ agent: DKGAgent; /** Number of concurrent worker loops (default 4 — RFC §4.5). */ workerConcurrency?: number; /** Durable fallback poll interval (default 100ms, including with queue wake support). */ pollIntervalMs?: number; /** * Heartbeat interval. Must be SHORTER than the queue's `leaseMs` * (default 15min); default 60s = 15× safety margin. */ heartbeatIntervalMs?: number; /** * Max time `stop()` will wait for in-flight jobs to complete before * returning. After the timeout, the in-flight `agent.assertion.promote` * continues in the background but the supervisor stops tracking it — * the next boot's `recoverOnStartup()` reconciles. */ shutdownTimeoutMs?: number; /** Deterministic time source for tests. */ now?: () => number; /** Deterministic randomness source for claim-failure jitter tests. */ random?: () => number; /** Defaults to `console.warn`. The daemon passes its own logger. */ log?: PromoteWorkerLogger; /** Retry interval for queue-only outcome bookkeeping (default 5s). */ bookkeepingRetryIntervalMs?: number; /** Maximum queue-only bookkeeping recovery window (default 10min). */ bookkeepingRetryBudgetMs?: number; /** * Interval of the post-commit recovery sweep (default 30s). The sweep also * runs once right after `recoverOnStartup()`; `0` disables only the * periodic repetition. */ postCommitRecoveryIntervalMs?: number; /** Deterministic sleep hook for tests. */ sleep?: (ms: number) => Promise; /** Defaults to a no-op. The daemon passes its `memoryGraphChanged` emitter. */ emitMemoryGraphChanged?: (event: PromoteMemoryGraphChangedEvent) => void; /** * Worker-id prefix used when minting each loop's `workerId`. Defaults * to `daemon-`. Tests inject a stable prefix. */ workerIdPrefix?: string; } export interface PromoteWorkerSupervisor { /** Run `recoverOnStartup()` once, then spawn the polling loops. */ start(): Promise; /** Stop polling and wait (up to `shutdownTimeoutMs`) for in-flight jobs. */ stop(): Promise; /** * Test-only: drive one full poll across every worker slot * synchronously. Returns the number of jobs picked up by this round. */ tickOnce(): Promise; /** * Observability — counts of completed runs since start. Counters reset * on every `start()` call (a new supervisor lifecycle). */ getCounters(): PromoteWorkerCounters; } export interface PromoteWorkerCounters { succeeded: number; failedTerminal: number; failedRetrying: number; /** * Codex #665: jobs whose promote ran successfully but whose post-promote * bookkeeping (commit-marker write / queue.succeed) failed mid-flight. * These remain in `running` state until next startup recovery; operators * MUST inspect SWM/VM before any explicit `/recover`. */ partialPromoteAmbiguity: number; /** Number of `runJob` invocations that started (regardless of outcome). */ attempted: number; /** Set when shuttingDown was hit mid-job; ops can correlate with abandoned counts at next startup. */ interruptedAtShutdown: number; /** Terminal post-commit failures the recovery sweep requeued for the idempotent replay. */ postCommitRequeued: number; /** Post-commit failures whose replay budget is spent; they wait for an operator `recover`. */ postCommitExhausted: number; } /** * Per-job execution — extracted so tests can exercise the exact * try/catch/heartbeat shape without spinning the supervisor. * * Resolves once the job's lifecycle transition (succeed or fail) has * been written back to the queue. Heartbeats run in the background * for the duration; both `succeed` and `fail` clear the lease so any * in-flight heartbeat after that point throws a `PromoteJobLeaseError`, * which the heartbeat catcher logs and ignores. */ export declare function runPromoteJob(args: { job: PromoteJob; queue: AsyncPromoteQueue; workerId: string; runPromote: (request: PromoteRequest, markPromoteStarted: () => Promise) => Promise<{ promotedCount: number; }>; now: () => number; heartbeatIntervalMs: number; bookkeepingRetryIntervalMs?: number; bookkeepingRetryBudgetMs?: number; sleep?: (ms: number) => Promise; shutdownSignal?: AbortSignal; log: PromoteWorkerLogger; emitMemoryGraphChanged?: (event: PromoteMemoryGraphChangedEvent) => void; }): Promise<{ outcome: 'succeeded' | 'failed_retrying' | 'failed_terminal' | 'partial_promote_ambiguity'; error?: ClassifiedPromoteError; }>; export declare function createPromoteWorkerSupervisor(config: PromoteWorkerConfig): PromoteWorkerSupervisor; //# sourceMappingURL=async-promote-worker.d.ts.map