From 282c08aadbdb7d4cd2314825e7e225ba682ad596 Mon Sep 17 00:00:00 2001 From: Andrei Date: Tue, 22 Sep 2026 19:00:04 +0100 Subject: [PATCH 1/5] feat(node): build the review runtime for the Node target The engine takes a ReviewRuntime, and apps/worker's composition root builds one out of Cloudflare bindings. Node had none, so createReviewRuntime threw and getOrFetchRawDiffForCompletedJob returned an empty string, which left the dashboard's diff view blank. Adds the Node composition root. The stores take a DbEnv rather than bindings, so most of it is a re-wire; the two pieces that were worker-local are the telemetry sink, rebuilt here on @codraoss/db's InstanceIdStore, and the repo config loader, which comes from @codraoss/api/platform. No Workers AI binding, so Cloudflare-hosted models are unavailable on this target and every other provider is reached over HTTP as before. --- apps/node/src/adapters/review-runtime.ts | 63 +++++++++++++++++++++++ apps/node/src/adapters/telemetry.ts | 65 ++++++++++++++++++++++++ apps/node/src/api-deps.ts | 14 ++--- 3 files changed, 135 insertions(+), 7 deletions(-) create mode 100644 apps/node/src/adapters/review-runtime.ts create mode 100644 apps/node/src/adapters/telemetry.ts diff --git a/apps/node/src/adapters/review-runtime.ts b/apps/node/src/adapters/review-runtime.ts new file mode 100644 index 00000000..ac901957 --- /dev/null +++ b/apps/node/src/adapters/review-runtime.ts @@ -0,0 +1,63 @@ +import type { ReviewRuntime } from '@codraoss/core/ports'; +import type { DbEnv } from '@codraoss/db/env'; +import { TokenTracker } from '@codraoss/core/token-tracker'; +import { FormatterService } from '@codraoss/core/formatter'; +import { GitHubService } from '@codraoss/provider-github'; +import { isRetryableModelError, ModelRunner, nextChainIndexOf } from '@codraoss/models'; +import { getResolvedModelConfig } from '@codraoss/db/model-configs'; +import { + makeFileReviewStore, + makeJobStore, + makeLearningStore, + makeModelConfigReader, + makeReviewSettingsReader, + makeWebhookDeliveryReader, +} from '@codraoss/db/repositories'; +import { loadRepoConfig } from '@codraoss/api/platform'; +import type { NodeAppBindings } from '../env'; +import { makeTelemetrySink } from './telemetry'; + +// The Node composition root: the one place Postgres, Redis and the GitHub/model services are wired +// to the engine's ports, mirroring apps/worker/src/adapters/index.ts. @codraoss/core sees this +// object and nothing else. +// +// Thinner than the Worker's because the stores take a DbEnv rather than Cloudflare bindings, and +// because there is no Workers AI binding: Cloudflare-hosted models are unavailable on this target, +// every other provider is reached over HTTP exactly as it is on Workers. +export function createReviewRuntime(env: NodeAppBindings): ReviewRuntime { + const dbEnv: DbEnv = { HYPERDRIVE: env.HYPERDRIVE, APP_KV: env.APP_KV, workerMode: false }; + + return { + kv: { + get: (key) => env.APP_KV.get(key), + put: (key, value, options) => env.APP_KV.put(key, value, options), + }, + clock: { now: () => Date.now() }, + ids: { randomUUID: () => crypto.randomUUID() }, + + botUsername: env.BOT_USERNAME, + + jobs: makeJobStore(dbEnv), + fileReviews: makeFileReviewStore(dbEnv), + settings: makeReviewSettingsReader(dbEnv), + webhooks: makeWebhookDeliveryReader(dbEnv), + learning: makeLearningStore(dbEnv), + modelConfigs: makeModelConfigReader(dbEnv), + repoConfig: { loadRepoConfig: (input) => loadRepoConfig(env.APP_KV, dbEnv, input) }, + telemetry: makeTelemetrySink(env, dbEnv), + + createTokenTracker: () => new TokenTracker(), + createGitHub: (installationId, tracker) => new GitHubService(env, installationId, tracker), + createModel: (jobId, tracker) => new ModelRunner({ + kv: env.APP_KV, + secretStore: { getSecret: async (key) => (env[key as keyof NodeAppBindings] as string | undefined) ?? process.env[key] ?? null }, + getConfig: (modelId) => getResolvedModelConfig(dbEnv, modelId), + tracker, + jobId, + }), + createFormatter: () => new FormatterService(env.APP_URL), + + githubClients: { forInstallation: (installationId) => new GitHubService(env, installationId) }, + modelErrors: { isRetryableModelError, nextChainIndexOf }, + }; +} diff --git a/apps/node/src/adapters/telemetry.ts b/apps/node/src/adapters/telemetry.ts new file mode 100644 index 00000000..7e2d6d13 --- /dev/null +++ b/apps/node/src/adapters/telemetry.ts @@ -0,0 +1,65 @@ +import type { DbEnv } from '@codraoss/db/env'; +import type { ReviewTelemetryEvent, TelemetrySink } from '@codraoss/core/ports'; +import { makeInstanceIdStore } from '@codraoss/db/repositories'; +import { logger } from '@codraoss/api/logger'; +import type { NodeAppBindings } from '../env'; +import pkg from '../../../../package.json' with { type: 'json' }; + +const DEFAULT_TELEMETRY_URL = 'https://codra.run/api/telemetry'; +const DEFAULT_TELEMETRY_SECRET = 'codra-telemetry-v1-secret-8f9a2b5c'; +const TELEMETRY_TIMEOUT_MS = 5000; + +function isDisabled(env: NodeAppBindings): boolean { + const flag = String(process.env.TELEMETRY_DISABLED ?? '').toLowerCase(); + if (flag === 'true' || flag === '1') return true; + return ['test', 'local'].includes(String(env.ENVIRONMENT ?? '').toLowerCase()); +} + +// Mirrors apps/worker/src/core/telemetry.ts. Kept per app rather than shared because the version +// string comes from the repo-root package.json, which no package-relative path can reach. +export function makeTelemetrySink(env: NodeAppBindings, dbEnv: DbEnv): TelemetrySink { + const instanceIds = makeInstanceIdStore(dbEnv); + + return { + // Swallows every error: telemetry must never fail a review. + async send(event: ReviewTelemetryEvent) { + try { + if (isDisabled(env)) return; + + // Drops the stub models used in tests so they cannot skew the aggregate. + const modelsUsed = event.modelsUsed + .map((model) => model.replace(/^(google|cloudflare|openai|anthropic|openrouter|nvidia):/i, '').trim()) + .filter((model) => model && !model.toLowerCase().includes('test')); + + if (event.modelsUsed.length > 0 && modelsUsed.length === 0) return; + + const controller = new AbortController(); + const timeout = setTimeout(() => controller.abort(), TELEMETRY_TIMEOUT_MS); + + try { + await fetch(process.env.TELEMETRY_API_URL ?? DEFAULT_TELEMETRY_URL, { + method: 'POST', + headers: { + 'Content-Type': 'application/json', + Authorization: `Bearer ${process.env.TELEMETRY_SECRET ?? DEFAULT_TELEMETRY_SECRET}`, + }, + body: JSON.stringify({ + ...event, + modelsUsed, + instanceId: await instanceIds.getOrCreateInstanceId(), + prsReviewed: 1, + codraVersion: pkg.version, + }), + signal: controller.signal, + }); + } finally { + clearTimeout(timeout); + } + } catch (error) { + logger.debug('Failed to send anonymous telemetry event', { + error: error instanceof Error ? error.message : String(error), + }); + } + }, + }; +} diff --git a/apps/node/src/api-deps.ts b/apps/node/src/api-deps.ts index 00d52bcd..0cf0f9e6 100644 --- a/apps/node/src/api-deps.ts +++ b/apps/node/src/api-deps.ts @@ -1,6 +1,8 @@ import { createSharedApiDeps } from '@codraoss/api'; +import { getOrFetchRawDiffForCompletedJob } from '@codraoss/core'; import type { NodeAppBindings } from './env'; import { logger } from '@codraoss/api/logger'; +import { createReviewRuntime } from './adapters/review-runtime'; export function createNodeApiDeps(env: NodeAppBindings) { return createSharedApiDeps({ @@ -12,16 +14,14 @@ export function createNodeApiDeps(env: NodeAppBindings) { enqueueReviewJob: async (input) => { await env.REVIEW_QUEUE.send(input); }, - terminateJobWorkflow: async (_job) => { - logger.warn('[STUB] terminateJobWorkflow called'); - }, + // No durable instance to kill: a Node review runs inside the BullMQ job, and cancelling it is + // the job store's business, which the caller has already handled by the time this runs. + terminateJobWorkflow: async () => {}, scheduleBestEffortJobMaintenance: () => { // In node, this is a long running process, we can just spawn a promise. }, - createReviewRuntime: () => { - throw new Error('Review runtime not implemented for Node yet.'); - }, - getOrFetchRawDiffForCompletedJob: async () => '', + createReviewRuntime: () => createReviewRuntime(env), + getOrFetchRawDiffForCompletedJob, logger, getSecret: async (key) => process.env[key] ?? null, From d7c3cc011346a5f144ddc1a7d03555189c175351 Mon Sep 17 00:00:00 2001 From: Andrei Date: Tue, 22 Sep 2026 19:00:04 +0100 Subject: [PATCH 2/5] feat(node): add the Node job orchestrator Closes #106. Runs the engine's phase machine in a plain loop. The Cloudflare orchestrator hibernates between phases because a Workers invocation has a fixed subrequest budget and next_phase exists to hand the next phase a fresh one; Node has no such budget, so the same contract collapses into awaiting the delay and calling again in-process, and freshInstance needs no handling. The loop runs inside runWithDb so every phase shares one AsyncLocalStorage-scoped connection. Errors are left to propagate: BullMQ owns retries and dead-lettering. --- apps/node/src/adapters/node-orchestrator.ts | 52 +++++++++++++++++++++ 1 file changed, 52 insertions(+) create mode 100644 apps/node/src/adapters/node-orchestrator.ts diff --git a/apps/node/src/adapters/node-orchestrator.ts b/apps/node/src/adapters/node-orchestrator.ts new file mode 100644 index 00000000..88495a38 --- /dev/null +++ b/apps/node/src/adapters/node-orchestrator.ts @@ -0,0 +1,52 @@ +import type { JobOrchestrator } from '@codraoss/core/ports'; +import type { ReviewJobMessage } from '@codraoss/schema'; +import { runReview } from '@codraoss/core'; +import { runWithDb } from '@codraoss/db/client'; +import { logger } from '@codraoss/api/logger'; +import type { NodeAppBindings } from '../env'; +import { createReviewRuntime } from './review-runtime'; + +const DEFAULT_RETRY_DELAY_SECONDS = 60; + +const sleep = (seconds: number) => new Promise((resolve) => setTimeout(resolve, seconds * 1000)); + +// Drives the engine's phase machine in a plain loop. +// +// The Cloudflare orchestrator has to hibernate between phases: a Workers invocation gets a fixed +// subrequest budget, so `next_phase` exists to hand the next phase a fresh one. Node has no such +// budget, so the same contract collapses into awaiting the delay and calling again in-process, and +// `freshInstance` needs no special handling -- the next iteration is already a clean slate for +// everything the flag was protecting. +export class NodeOrchestrator implements JobOrchestrator { + constructor(private readonly env: NodeAppBindings) {} + + async startReviewJob(id: string, params: ReviewJobMessage): Promise { + // The whole run shares one AsyncLocalStorage-scoped connection, so every query inside every + // phase resolves against the same client. + await runWithDb(this.env, async () => { + const runtime = createReviewRuntime(this.env); + let phase: ReviewJobMessage['phase'] = params.phase ?? 'prepare'; + + for (;;) { + const result = await runReview(runtime, { ...params, phase }); + + if (result.action === 'next_phase') { + phase = result.phase; + if (result.delaySeconds > 0) await sleep(result.delaySeconds); + continue; + } + + if (result.action === 'retry') { + // No work happened: admission was throttled or the lease is held elsewhere. Same phase, + // after the delay the engine asked for. + await sleep(result.delaySeconds ?? DEFAULT_RETRY_DELAY_SECONDS); + continue; + } + + // 'ack': finished, or not ours to run. + logger.info('Review job finished', { jobId: id, phase }); + return; + } + }); + } +} From 14aee3b277b707559b5ece4bdadfb8fb5d9db923 Mon Sep 17 00:00:00 2001 From: Andrei Date: Tue, 22 Sep 2026 19:00:11 +0100 Subject: [PATCH 3/5] feat(node): run the BullMQ review worker in the app process Closes #107. The worker consumes codra-reviews and hands each message to the orchestrator. Payloads that fail schema validation throw UnrecoverableError so a permanently broken message is failed once instead of spending the backoff schedule. START_API and START_WORKER both default on, so one container is the whole deployment; setting either to 'false' splits the API and the review workers across containers that scale independently, and setting both is a startup error. SIGTERM and SIGINT close the listener, drop idle keep-alive sockets, and wait for worker.close() so the review in flight finishes before the process goes away. Shutdown deliberately avoids process.exit() on the happy path, which discards buffered stdout and loses the shutdown log lines; a 30s unref'd timer is the backstop. The worker owns its own Redis connection and closes it by hand, since BullMQ only closes connections it opened itself. --- apps/node/src/index.ts | 95 ++++++++++++++++++++++++++++++++--------- apps/node/src/queue.ts | 3 ++ apps/node/src/worker.ts | 61 ++++++++++++++++++++++++++ 3 files changed, 140 insertions(+), 19 deletions(-) create mode 100644 apps/node/src/queue.ts create mode 100644 apps/node/src/worker.ts diff --git a/apps/node/src/index.ts b/apps/node/src/index.ts index c1dc488f..5208ce93 100644 --- a/apps/node/src/index.ts +++ b/apps/node/src/index.ts @@ -5,14 +5,22 @@ dotenv.config({ path: path.resolve(process.cwd(), '.dev.vars') }); // Root level import { serve } from '@hono/node-server'; import { createApiRouter } from '@codraoss/api'; import { runWithDb } from '@codraoss/db/client'; -import { InMemoryOrchestrator, InMemorySessionStore } from '@codraoss/core/ports'; +import { InMemorySessionStore } from '@codraoss/core/ports'; import { createNodeApiDeps } from './api-deps'; -import { createNodeEnv } from './env'; +import { createNodeEnv, type NodeAppBindings } from './env'; import { logger } from '@codraoss/api/logger'; import Redis from 'ioredis'; import { RedisKVAdapter } from './adapters/redis-kv'; import { Queue } from 'bullmq'; import { RedisQueueAdapter } from './adapters/redis-queue'; +import { NodeOrchestrator } from './adapters/node-orchestrator'; +import { REVIEW_QUEUE_NAME } from './queue'; +import { startWorker } from './worker'; + +// Both default on, so a single container is the whole deployment. Setting one to 'false' splits the +// API and the review workers across containers that scale independently. +const runApi = process.env.START_API !== 'false'; +const runWorker = process.env.START_WORKER !== 'false'; const redisUrl = process.env.REDIS_URL || 'redis://localhost:6379'; const redisClient = new Redis(redisUrl, { maxRetriesPerRequest: null }); // maxRetriesPerRequest: null is required for bullmq @@ -20,16 +28,17 @@ redisClient.on('error', (err) => { logger.error('[Redis Error]', err); }); -const reviewQueue = new Queue('codra-reviews', { connection: redisClient }); +const reviewQueue = new Queue(REVIEW_QUEUE_NAME, { connection: redisClient }); -const stubs = { +const env: NodeAppBindings = createNodeEnv({ SESSION_STORE: new InMemorySessionStore(), APP_KV: new RedisKVAdapter(redisClient), REVIEW_QUEUE: new RedisQueueAdapter(reviewQueue), - REVIEW_ORCHESTRATOR: new InMemoryOrchestrator(), -}; + // Set below: the orchestrator needs the env it is being attached to. + REVIEW_ORCHESTRATOR: { startReviewJob: async () => {} }, +}); -const env = createNodeEnv(stubs); +env.REVIEW_ORCHESTRATOR = new NodeOrchestrator(env); import fs from 'node:fs'; import { serveStatic } from '@hono/node-server/serve-static'; @@ -63,18 +72,66 @@ app.use('/*.svg', serveStatic({ root: process.cwd().endsWith('node') ? '../../di app.use('/*.ico', serveStatic({ root: process.cwd().endsWith('node') ? '../../dist/client' : 'dist/client' })); const port = parseInt(process.env.PORT || '3000', 10); -serve({ - fetch: async (request) => { - const apiEnv = { - ...envWithAssets, - deps: createNodeApiDeps(env), - }; - try { return await runWithDb(env, () => app.fetch(request, apiEnv as any)); } catch (e) { console.error('SERVE ERROR:', e); throw e; } - }, - port, -}, (info) => { - logger.info(`Codra Node server running on http://localhost:${info.port}`); -}); +const server = runApi + ? serve({ + fetch: async (request) => { + const apiEnv = { + ...envWithAssets, + deps: createNodeApiDeps(env), + }; + try { return await runWithDb(env, () => app.fetch(request, apiEnv as any)); } catch (e) { console.error('SERVE ERROR:', e); throw e; } + }, + port, + }, (info) => { + logger.info(`Codra Node server running on http://localhost:${info.port}`); + }) + : null; + +const reviewWorker = runWorker ? startWorker(env, redisUrl) : null; +if (reviewWorker) logger.info(`Codra review worker listening on queue ${REVIEW_QUEUE_NAME}`); +if (!server && !reviewWorker) { + logger.error('START_API and START_WORKER are both false: nothing to run'); + process.exit(1); +} + +// worker.close() waits for the review in flight to finish, so a redeploy does not abandon a job +// mid-phase with its lease still held. +const SHUTDOWN_TIMEOUT_MS = 30_000; + +let shuttingDown = false; +async function shutdown(signal: string) { + if (shuttingDown) return; + shuttingDown = true; + logger.info(`Received ${signal}, shutting down`); + + // No process.exit() on the happy path: it discards buffered stdout, so the shutdown log lines are + // lost exactly when an operator needs them. Closing every handle lets the loop drain and Node + // exit 0 on its own; the timer is the backstop for a close that never resolves, and unref() keeps + // it from holding the process open by itself. + const forceExit = setTimeout(() => { + logger.error('Shutdown timed out, exiting'); + process.exit(1); + }, SHUTDOWN_TIMEOUT_MS); + forceExit.unref(); + + try { + // Stop accepting, drop sockets sitting idle on keep-alive, then let the worker finish the review + // it is holding before the remaining sockets are cut. + server?.close(); + if (server && 'closeIdleConnections' in server) server.closeIdleConnections(); + await reviewWorker?.close(); + if (server && 'closeAllConnections' in server) server.closeAllConnections(); + await reviewQueue.close(); + redisClient.disconnect(); + clearTimeout(forceExit); + } catch (error) { + logger.error('Error during shutdown', error); + process.exit(1); + } +} + +process.on('SIGTERM', () => void shutdown('SIGTERM')); +process.on('SIGINT', () => void shutdown('SIGINT')); diff --git a/apps/node/src/queue.ts b/apps/node/src/queue.ts new file mode 100644 index 00000000..1ec1d185 --- /dev/null +++ b/apps/node/src/queue.ts @@ -0,0 +1,3 @@ +// Shared between the producer in index.ts and the consumer in worker.ts: they must agree on the +// name or enqueued reviews are never picked up. +export const REVIEW_QUEUE_NAME = 'codra-reviews'; diff --git a/apps/node/src/worker.ts b/apps/node/src/worker.ts new file mode 100644 index 00000000..38aa0038 --- /dev/null +++ b/apps/node/src/worker.ts @@ -0,0 +1,61 @@ +import { UnrecoverableError, Worker, type Job } from 'bullmq'; +import Redis from 'ioredis'; +import { reviewJobMessageSchema } from '@codraoss/schema'; +import { logger } from '@codraoss/api/logger'; +import { NodeOrchestrator } from './adapters/node-orchestrator'; +import type { NodeAppBindings } from './env'; +import { REVIEW_QUEUE_NAME } from './queue'; + +export interface ReviewWorker { + worker: Worker; + close(): Promise; +} + +// A review holds its connection for the length of the job, so the worker gets its own rather than +// sharing the producer's: BullMQ issues blocking commands here that would stall the API's enqueues. +export function startWorker(env: NodeAppBindings, redisUrl: string): ReviewWorker { + const connection = new Redis(redisUrl, { maxRetriesPerRequest: null }); + connection.on('error', (err) => logger.error('Worker Redis connection error', err)); + + const worker = new Worker( + REVIEW_QUEUE_NAME, + async (job: Job) => { + const parsed = reviewJobMessageSchema.safeParse(job.data); + if (!parsed.success) { + // An unparseable payload never becomes valid on a retry, so the job is failed outright + // rather than spending the backoff schedule on it. + logger.error('Discarding review job with an invalid payload', { + jobId: job.id, + issues: parsed.error.issues, + }); + throw new UnrecoverableError('Invalid review job payload'); + } + + await new NodeOrchestrator(env).startReviewJob(job.id ?? parsed.data.deliveryId, parsed.data); + }, + { connection }, + ); + + worker.on('completed', (job) => { + logger.info('Review job completed', { jobId: job.id }); + }); + + worker.on('failed', (job, err) => { + logger.error('Review job failed', { + jobId: job?.id, + attemptsMade: job?.attemptsMade, + error: err instanceof Error ? err.message : String(err), + }); + }); + + return { + worker, + // BullMQ closes the connections it opened itself, never one handed to it, so this disconnects + // the socket above by hand. Without it the process keeps an open handle after SIGTERM and never + // exits on its own. + async close() { + await worker.close(); + connection.disconnect(); + }, + }; +} From be5e3d4fa1f22d97c2c6d8eabe9383c282de1784 Mon Sep 17 00:00:00 2001 From: Andrei Date: Tue, 22 Sep 2026 19:41:20 +0100 Subject: [PATCH 4/5] fix(node): address self-review findings on the worker and orchestrator Connection handling, found by measuring: 20 review jobs left 20 idle Postgres connections and the process would not exit on SIGTERM. runWithDb builds a pool per call and never ends it, which is correct on Workers, where an invocation may not reuse another's I/O, and wrong for a long-lived server. The orchestrator and the request path now use the process-wide pool instead; the same 20 jobs hold one connection. Ending that pool needed an API, so @codraoss/db gains closeDb(), which closes the pooled clients only -- tracking runWithDb's would retain a handle per request on Workers. Shutdown calls it; without it the open socket kept the event loop alive after everything else had closed. - Bound the retry loop. 'retry' means no work happened, and the loop had no ceiling, so a stale 'running' row that nothing reaps on this target (job maintenance is a no-op here) could spin a worker slot forever. After five rounds the message goes back to the queue. - Carry result.jobId into later phases, so they take getJobForProcessing rather than re-resolving the webhook delivery, which for comment events meant another GitHub call at every phase boundary. - Worker concurrency defaults to 4 and reads WORKER_CONCURRENCY. BullMQ's default of 1 let one review hold the container through its in-process sleeps. - Shutdown budget is 5 minutes, not 30 seconds, and reads SHUTDOWN_TIMEOUT_MS. A review runs for minutes, so the old backstop killed one on every redeploy. - Restore the NODE_ENV/VITEST guard on telemetry, which the Worker has and this sink had dropped: a test touching it would have posted a real event. - Correct the terminateJobWorkflow comment, which had the ordering backwards. Stop and rerun call it before cancelling; with a no-op, stop is eventually-effective, not immediate. --- apps/node/src/adapters/node-orchestrator.ts | 68 +++++++++++++-------- apps/node/src/adapters/telemetry.ts | 3 + apps/node/src/api-deps.ts | 7 ++- apps/node/src/index.ts | 15 ++++- apps/node/src/worker.ts | 8 ++- packages/db/src/client.ts | 21 ++++++- 6 files changed, 89 insertions(+), 33 deletions(-) diff --git a/apps/node/src/adapters/node-orchestrator.ts b/apps/node/src/adapters/node-orchestrator.ts index 88495a38..399a48cc 100644 --- a/apps/node/src/adapters/node-orchestrator.ts +++ b/apps/node/src/adapters/node-orchestrator.ts @@ -1,13 +1,18 @@ import type { JobOrchestrator } from '@codraoss/core/ports'; import type { ReviewJobMessage } from '@codraoss/schema'; import { runReview } from '@codraoss/core'; -import { runWithDb } from '@codraoss/db/client'; import { logger } from '@codraoss/api/logger'; import type { NodeAppBindings } from '../env'; import { createReviewRuntime } from './review-runtime'; const DEFAULT_RETRY_DELAY_SECONDS = 60; +// A 'retry' means no work happened: admission was throttled, or another runner holds the lease. +// Waiting it out in-process is cheaper than a Redis round trip, but only for a few rounds -- a +// worker slot spinning here is a slot not running reviews, and at concurrency 1 that is the whole +// queue. Past this the message goes back to the queue so something else can make progress. +const MAX_CONSECUTIVE_RETRIES = 5; + const sleep = (seconds: number) => new Promise((resolve) => setTimeout(resolve, seconds * 1000)); // Drives the engine's phase machine in a plain loop. @@ -17,36 +22,51 @@ const sleep = (seconds: number) => new Promise((resolve) => setTimeout(resolve, // budget, so the same contract collapses into awaiting the delay and calling again in-process, and // `freshInstance` needs no special handling -- the next iteration is already a clean slate for // everything the flag was protecting. +// +// Deliberately not wrapped in runWithDb. That helper builds a fresh postgres pool per call and never +// ends it, which is right on Workers, where an invocation may not reuse another's I/O, and wrong +// here: it left one connection per review job open for the pool's max_lifetime, and the open handles +// kept the process alive through SIGTERM. Outside a runWithDb scope, getDb serves the process-wide +// pool instead, which is what a long-lived server wants. export class NodeOrchestrator implements JobOrchestrator { constructor(private readonly env: NodeAppBindings) {} async startReviewJob(id: string, params: ReviewJobMessage): Promise { - // The whole run shares one AsyncLocalStorage-scoped connection, so every query inside every - // phase resolves against the same client. - await runWithDb(this.env, async () => { - const runtime = createReviewRuntime(this.env); - let phase: ReviewJobMessage['phase'] = params.phase ?? 'prepare'; - - for (;;) { - const result = await runReview(runtime, { ...params, phase }); - - if (result.action === 'next_phase') { - phase = result.phase; - if (result.delaySeconds > 0) await sleep(result.delaySeconds); - continue; - } + const runtime = createReviewRuntime(this.env); + let message = params; + let phase: ReviewJobMessage['phase'] = params.phase ?? 'prepare'; + let consecutiveRetries = 0; + + for (;;) { + const result = await runReview(runtime, { ...message, phase }); - if (result.action === 'retry') { - // No work happened: admission was throttled or the lease is held elsewhere. Same phase, - // after the delay the engine asked for. - await sleep(result.delaySeconds ?? DEFAULT_RETRY_DELAY_SECONDS); - continue; + if (result.action === 'next_phase') { + consecutiveRetries = 0; + phase = result.phase; + // Carrying the id forward puts later phases on the cheap getJobForProcessing path, instead + // of re-resolving the webhook delivery -- which for a comment event means another call to + // GitHub at every phase boundary. + if (result.jobId) message = { ...message, jobId: result.jobId }; + if (result.delaySeconds > 0) await sleep(result.delaySeconds); + continue; + } + + if (result.action === 'retry') { + const delaySeconds = result.delaySeconds ?? DEFAULT_RETRY_DELAY_SECONDS; + + if (++consecutiveRetries > MAX_CONSECUTIVE_RETRIES) { + logger.warn('Review job still not admitted, returning it to the queue', { jobId: id, phase }); + await this.env.REVIEW_QUEUE.send({ ...message, phase }, { delaySeconds }); + return; } - // 'ack': finished, or not ours to run. - logger.info('Review job finished', { jobId: id, phase }); - return; + await sleep(delaySeconds); + continue; } - }); + + // 'ack': finished, or not ours to run. + logger.info('Review job finished', { jobId: id, phase }); + return; + } } } diff --git a/apps/node/src/adapters/telemetry.ts b/apps/node/src/adapters/telemetry.ts index 7e2d6d13..5163dc5a 100644 --- a/apps/node/src/adapters/telemetry.ts +++ b/apps/node/src/adapters/telemetry.ts @@ -12,6 +12,9 @@ const TELEMETRY_TIMEOUT_MS = 5000; function isDisabled(env: NodeAppBindings): boolean { const flag = String(process.env.TELEMETRY_DISABLED ?? '').toLowerCase(); if (flag === 'true' || flag === '1') return true; + // ENVIRONMENT defaults to 'development', so without the NODE_ENV/VITEST check a test that + // exercises this sink would post a real event. + if (process.env.NODE_ENV === 'test' || process.env.VITEST) return true; return ['test', 'local'].includes(String(env.ENVIRONMENT ?? '').toLowerCase()); } diff --git a/apps/node/src/api-deps.ts b/apps/node/src/api-deps.ts index 0cf0f9e6..355d7a7b 100644 --- a/apps/node/src/api-deps.ts +++ b/apps/node/src/api-deps.ts @@ -14,8 +14,11 @@ export function createNodeApiDeps(env: NodeAppBindings) { enqueueReviewJob: async (input) => { await env.REVIEW_QUEUE.send(input); }, - // No durable instance to kill: a Node review runs inside the BullMQ job, and cancelling it is - // the job store's business, which the caller has already handled by the time this runs. + // There is no durable instance to terminate: a Node review runs inside its BullMQ job. The stop + // and rerun routes call this BEFORE cancelling, to stop the old run racing the new one, and a + // no-op cannot honour that -- the loop keeps working the current phase until it next reads the + // job's status and sees it is terminal. Closing that gap needs a cooperative cancel signal the + // loop polls; until then, stop is eventually-effective rather than immediate. terminateJobWorkflow: async () => {}, scheduleBestEffortJobMaintenance: () => { // In node, this is a long running process, we can just spawn a promise. diff --git a/apps/node/src/index.ts b/apps/node/src/index.ts index 5208ce93..9d592341 100644 --- a/apps/node/src/index.ts +++ b/apps/node/src/index.ts @@ -4,7 +4,7 @@ dotenv.config({ path: path.resolve(process.cwd(), '../../.dev.vars') }); // Fall dotenv.config({ path: path.resolve(process.cwd(), '.dev.vars') }); // Root level import { serve } from '@hono/node-server'; import { createApiRouter } from '@codraoss/api'; -import { runWithDb } from '@codraoss/db/client'; +import { closeDb } from '@codraoss/db/client'; import { InMemorySessionStore } from '@codraoss/core/ports'; import { createNodeApiDeps } from './api-deps'; import { createNodeEnv, type NodeAppBindings } from './env'; @@ -79,7 +79,9 @@ const server = runApi ...envWithAssets, deps: createNodeApiDeps(env), }; - try { return await runWithDb(env, () => app.fetch(request, apiEnv as any)); } catch (e) { console.error('SERVE ERROR:', e); throw e; } + // No runWithDb: it opens a postgres pool per call and never ends it, so wrapping the + // request path leaked a connection per request. getDb falls back to the process-wide pool. + try { return await app.fetch(request, apiEnv as any); } catch (e) { console.error('SERVE ERROR:', e); throw e; } }, port, }, (info) => { @@ -96,7 +98,11 @@ if (!server && !reviewWorker) { // worker.close() waits for the review in flight to finish, so a redeploy does not abandon a job // mid-phase with its lease still held. -const SHUTDOWN_TIMEOUT_MS = 30_000; +// Long enough to outlast a review phase: worker.close() waits for the job in flight, and model +// calls plus the engine's inter-phase sleeps run to minutes. Too short a budget and every redeploy +// kills a running review, leaving its row 'running' with the lease still held. +// Docker caps this independently -- raise stop_grace_period past it, or the container is killed first. +const SHUTDOWN_TIMEOUT_MS = Number(process.env.SHUTDOWN_TIMEOUT_MS ?? 300_000); let shuttingDown = false; async function shutdown(signal: string) { @@ -123,6 +129,9 @@ async function shutdown(signal: string) { if (server && 'closeAllConnections' in server) server.closeAllConnections(); await reviewQueue.close(); redisClient.disconnect(); + // The pooled Postgres socket is an active handle: without ending it the process stays up after + // everything else has closed. + await closeDb(); clearTimeout(forceExit); } catch (error) { logger.error('Error during shutdown', error); diff --git a/apps/node/src/worker.ts b/apps/node/src/worker.ts index 38aa0038..428c41bd 100644 --- a/apps/node/src/worker.ts +++ b/apps/node/src/worker.ts @@ -6,6 +6,12 @@ import { NodeOrchestrator } from './adapters/node-orchestrator'; import type { NodeAppBindings } from './env'; import { REVIEW_QUEUE_NAME } from './queue'; +// BullMQ defaults to 1, which lets a single review monopolise the container: the engine sleeps +// in-process between phases and while polling async model batches, and none of that is work the slot +// could not spend on another job. 4 matches the engine's highest admission level; the engine still +// decides how many actually run at once. +const WORKER_CONCURRENCY = Number(process.env.WORKER_CONCURRENCY ?? 4); + export interface ReviewWorker { worker: Worker; close(): Promise; @@ -33,7 +39,7 @@ export function startWorker(env: NodeAppBindings, redisUrl: string): ReviewWorke await new NodeOrchestrator(env).startReviewJob(job.id ?? parsed.data.deliveryId, parsed.data); }, - { connection }, + { connection, concurrency: WORKER_CONCURRENCY }, ); worker.on('completed', (job) => { diff --git a/packages/db/src/client.ts b/packages/db/src/client.ts index 020f0235..b778c0db 100644 --- a/packages/db/src/client.ts +++ b/packages/db/src/client.ts @@ -9,7 +9,7 @@ type DbClient = { const dbStorage = new AsyncLocalStorage(); -function createDbClient(env: DbEnv): DbClient { +function createDbClient(env: DbEnv, onOpen?: (close: () => Promise) => void): DbClient { const sql = postgres(env.HYPERDRIVE.connectionString, { max: 5, fetch_types: false, @@ -17,6 +17,8 @@ function createDbClient(env: DbEnv): DbClient { onnotice: () => {}, }); + onOpen?.(() => sql.end({ timeout: 5 })); + return { async query(sqlText: string, params: unknown[] = []) { return (await sql.unsafe(sqlText, params.map(normalizeParam) as any[], { prepare: false })) as T[]; @@ -65,6 +67,10 @@ export function runWithDb(env: DbEnv, fn: () => T): T { // Module-scoped Map pools connections outside runWithDb, but must self-heal when request context changes. const fallbackClients = new Map(); +// Parallel to fallbackClients: the postgres handle behind each, so a long-lived host can end them on +// shutdown. Without this the open socket keeps Node's event loop alive after SIGTERM. +const fallbackCloses = new Map Promise>(); + // Catch dead request context I/O errors and terminated connections. function isStaleConnectionError(error: unknown): boolean { const message = error instanceof Error ? error.message : String(error); @@ -79,12 +85,21 @@ export function getDb(env: DbEnv) { const connectionString = env.HYPERDRIVE.connectionString; let client = fallbackClients.get(connectionString); if (!client) { - client = createDbClient(env); + client = createDbClient(env, (close) => fallbackCloses.set(connectionString, close)); fallbackClients.set(connectionString, client); } return client; } +// Only the pooled clients are closable. A runWithDb client belongs to one invocation and is dropped +// with it, and holding those open here would retain a handle per request on Workers. +export async function closeDb(): Promise { + const closing = [...fallbackCloses.values()].map((close) => close()); + fallbackCloses.clear(); + fallbackClients.clear(); + await Promise.allSettled(closing); +} + // Retries op once on a fresh client if the module-cached socket has gone stale. async function withStaleConnectionRecovery(env: DbEnv, op: (db: DbClient) => Promise): Promise { const inScope = dbStorage.getStore() !== undefined; @@ -95,7 +110,7 @@ async function withStaleConnectionRecovery(env: DbEnv, op: (db: DbClient) => const connectionString = env.HYPERDRIVE.connectionString; fallbackClients.delete(connectionString); - const fresh = createDbClient(env); + const fresh = createDbClient(env, (close) => fallbackCloses.set(connectionString, close)); fallbackClients.set(connectionString, fresh); return op(fresh); } From 0c2a4a59b705c6d13ea12e03d54d3fdcb8867aeb Mon Sep 17 00:00:00 2001 From: Andrei Date: Tue, 22 Sep 2026 21:30:32 +0100 Subject: [PATCH 5/5] fix(node): configure job retries and consume the orchestrator port The queue had no attempts or backoff, and BullMQ defaults to attempts: 0 with a retry check of `attemptsMade + 1 < opts.attempts`, so every job got exactly one run. That mattered because runReview's try block does not cover resolving the job, the admission query or claiming the lease: a transient Postgres error there escapes, and with no retry the row is stranded where nothing on this target revisits it. The queue now mirrors the Cloudflare driver's five attempts with 60s exponential backoff. UnrecoverableError still bypasses it, so an unparseable payload fails once. Completed and failed jobs are trimmed so a long-lived install does not grow Redis without bound. - The worker now drives reviews through env.REVIEW_ORCHESTRATOR instead of constructing its own. The field was wired and read by nothing, and the stub it was initialised with would have silently swallowed any job that arrived before the real orchestrator was assigned. - Release the replaced connection's closer in the stale-connection branch rather than overwriting the entry, so the map cannot strand a handle if closeDb() is ever used from a worker-mode host. --- apps/node/src/index.ts | 24 ++++++++++++++++++++---- apps/node/src/worker.ts | 5 +++-- packages/db/src/client.ts | 3 +++ 3 files changed, 26 insertions(+), 6 deletions(-) diff --git a/apps/node/src/index.ts b/apps/node/src/index.ts index 9d592341..c632b5c3 100644 --- a/apps/node/src/index.ts +++ b/apps/node/src/index.ts @@ -28,17 +28,33 @@ redisClient.on('error', (err) => { logger.error('[Redis Error]', err); }); -const reviewQueue = new Queue(REVIEW_QUEUE_NAME, { connection: redisClient }); +// Without these a job gets exactly one attempt: BullMQ defaults to attempts: 0, and its retry check +// is `attemptsMade + 1 < opts.attempts`. The engine's own try block does not cover resolving the job, +// the admission query or claiming the lease, so a transient Postgres error there would escape and +// strand the row with nothing to revisit it. Mirrors the Cloudflare driver, which runs each phase +// under `retries: { limit: 5, delay: '60 seconds', backoff: 'exponential' }`. UnrecoverableError +// still bypasses this, so an unparseable payload fails once. +// Completed and failed jobs are trimmed so a long-lived install does not grow Redis without bound. +const reviewQueue = new Queue(REVIEW_QUEUE_NAME, { + connection: redisClient, + defaultJobOptions: { + attempts: 5, + backoff: { type: 'exponential', delay: 60_000 }, + removeOnComplete: { count: 1000 }, + removeOnFail: { count: 5000 }, + }, +}); const env: NodeAppBindings = createNodeEnv({ SESSION_STORE: new InMemorySessionStore(), APP_KV: new RedisKVAdapter(redisClient), REVIEW_QUEUE: new RedisQueueAdapter(reviewQueue), - // Set below: the orchestrator needs the env it is being attached to. - REVIEW_ORCHESTRATOR: { startReviewJob: async () => {} }, + // Delegates rather than holding the instance, because the orchestrator needs the env being built + // here. The closure only runs once a job is in hand, by which point `orchestrator` is assigned. + REVIEW_ORCHESTRATOR: { startReviewJob: (id, params) => orchestrator.startReviewJob(id, params) }, }); -env.REVIEW_ORCHESTRATOR = new NodeOrchestrator(env); +const orchestrator = new NodeOrchestrator(env); import fs from 'node:fs'; import { serveStatic } from '@hono/node-server/serve-static'; diff --git a/apps/node/src/worker.ts b/apps/node/src/worker.ts index 428c41bd..6f47dc8f 100644 --- a/apps/node/src/worker.ts +++ b/apps/node/src/worker.ts @@ -2,7 +2,6 @@ import { UnrecoverableError, Worker, type Job } from 'bullmq'; import Redis from 'ioredis'; import { reviewJobMessageSchema } from '@codraoss/schema'; import { logger } from '@codraoss/api/logger'; -import { NodeOrchestrator } from './adapters/node-orchestrator'; import type { NodeAppBindings } from './env'; import { REVIEW_QUEUE_NAME } from './queue'; @@ -37,7 +36,9 @@ export function startWorker(env: NodeAppBindings, redisUrl: string): ReviewWorke throw new UnrecoverableError('Invalid review job payload'); } - await new NodeOrchestrator(env).startReviewJob(job.id ?? parsed.data.deliveryId, parsed.data); + // Goes through the port on env rather than constructing an orchestrator here, so the wiring + // in index.ts is the single place that decides what drives a review. + await env.REVIEW_ORCHESTRATOR.startReviewJob(job.id ?? parsed.data.deliveryId, parsed.data); }, { connection, concurrency: WORKER_CONCURRENCY }, ); diff --git a/packages/db/src/client.ts b/packages/db/src/client.ts index b778c0db..1ea7bdcc 100644 --- a/packages/db/src/client.ts +++ b/packages/db/src/client.ts @@ -110,6 +110,9 @@ async function withStaleConnectionRecovery(env: DbEnv, op: (db: DbClient) => const connectionString = env.HYPERDRIVE.connectionString; fallbackClients.delete(connectionString); + // Release the socket being replaced before its closer is overwritten; it is already stale, so a + // failure to end it is nothing to act on. + void fallbackCloses.get(connectionString)?.().catch(() => {}); const fresh = createDbClient(env, (close) => fallbackCloses.set(connectionString, close)); fallbackClients.set(connectionString, fresh); return op(fresh);