Repository navigation
Process reviews on Node: BullMQ worker, job orchestrator, and review runtime #123
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Closed
Closed
Changes from all commits
Commits
Show all changes
5 commits
Select commit
Hold shift + click to select a range
282c08a
feat(node): build the review runtime for the Node target
onel d7c3cc0
feat(node): add the Node job orchestrator
onel 14aee3b
feat(node): run the BullMQ review worker in the app process
onel be5e3d4
fix(node): address self-review findings on the worker and orchestrator
onel 0c2a4a5
fix(node): configure job retries and consume the orchestrator port
onel File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 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; | ||
| } | ||
| } | ||
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 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 }, | ||
| }; | ||
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 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), | ||
| }); | ||
| } | ||
| }, | ||
| }; | ||
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 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'; |
Oops, something went wrong.
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
terminateJobWorkflowfails to prevent race conditionsThe
terminateJobWorkflowfunction 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.There was a problem hiding this comment.
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
ackonce the row is terminal or superseded, so a stopped job stops at the next phase boundary without any help fromterminateJobWorkflow.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/coreaffecting 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.