Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
72 changes: 72 additions & 0 deletions apps/node/src/adapters/node-orchestrator.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,72 @@
import type { JobOrchestrator } from '@codraoss/core/ports';
import type { ReviewJobMessage } from '@codraoss/schema';
import { runReview } from '@codraoss/core';
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.
//
// 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.
//
// 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<void> {
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 === '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;
}

await sleep(delaySeconds);
continue;
}

// 'ack': finished, or not ours to run.
logger.info('Review job finished', { jobId: id, phase });
return;
}
}
}
63 changes: 63 additions & 0 deletions apps/node/src/adapters/review-runtime.ts
Original file line number Diff line number Diff line change
@@ -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 },
};
}
68 changes: 68 additions & 0 deletions apps/node/src/adapters/telemetry.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,68 @@
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;
// 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());
}

// 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),
});
}
},
};
}
17 changes: 10 additions & 7 deletions apps/node/src/api-deps.ts
Original file line number Diff line number Diff line change
@@ -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({
Expand All @@ -12,16 +14,17 @@ export function createNodeApiDeps(env: NodeAppBindings) {
enqueueReviewJob: async (input) => {
await env.REVIEW_QUEUE.send(input);
},
terminateJobWorkflow: async (_job) => {
logger.warn('[STUB] terminateJobWorkflow called');
},
// 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 () => {},

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 No-op terminateJobWorkflow fails to prevent race conditions

The terminateJobWorkflow function is implemented as a no-op, which the accompanying comment explicitly states "cannot honour that [stop the old run racing the new one]". This means that when a stop or rerun command is issued, the currently running review job will not be immediately terminated. This introduces a race condition where the old job might continue processing for an unspecified duration, potentially conflicting with a newly started job, consuming resources unnecessarily, and leading to inconsistent results or wasted computation.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Correct, and deliberate — it is called out in the PR description as a known gap rather than an oversight. Worth splitting into the two cases though:

Between phases is already handled. The engine resolves the job at the start of every phase and returns ack once the row is terminal or superseded, so a stopped job stops at the next phase boundary without any help from terminateJobWorkflow.

Mid-phase is the real gap. A phase already in flight keeps running to its boundary, so a stop can still cost model tokens and post comments after the user asked for it. Nothing in the app can close that: the engine would have to poll a cancel signal at its own checkpoints, which is a change in @codraoss/core affecting the Cloudflare path too.

On Cloudflare this is covered because REVIEW_WORKFLOW.get(id).terminate() kills the instance outright. There is no equivalent on Node — the review is just a function running inside a BullMQ job.

I left it out rather than half-implement a cancel that only fires where the engine already stops on its own. Happy to open a separate issue for a cooperative cancel signal if you want it tracked, or to take it in this PR if you would rather it not ship with the gap.

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,

Expand Down
122 changes: 102 additions & 20 deletions apps/node/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,32 +4,57 @@ 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 { InMemoryOrchestrator, InMemorySessionStore } from '@codraoss/core/ports';
import { closeDb } from '@codraoss/db/client';
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
redisClient.on('error', (err) => {
logger.error('[Redis Error]', err);
});

const reviewQueue = new Queue('codra-reviews', { 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 stubs = {
const env: NodeAppBindings = createNodeEnv({
SESSION_STORE: new InMemorySessionStore(),
APP_KV: new RedisKVAdapter(redisClient),
REVIEW_QUEUE: new RedisQueueAdapter(reviewQueue),
REVIEW_ORCHESTRATOR: new InMemoryOrchestrator(),
};
// 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) },
});

const env = createNodeEnv(stubs);
const orchestrator = new NodeOrchestrator(env);

import fs from 'node:fs';
import { serveStatic } from '@hono/node-server/serve-static';
Expand Down Expand Up @@ -63,18 +88,75 @@ 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),
};
// 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) => {
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.
// 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) {
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();
// 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);
process.exit(1);
}
}

process.on('SIGTERM', () => void shutdown('SIGTERM'));
process.on('SIGINT', () => void shutdown('SIGINT'));



Expand Down
3 changes: 3 additions & 0 deletions apps/node/src/queue.ts
Original file line number Diff line number Diff line change
@@ -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';
Loading
Loading