Skip to content

Process reviews on Node: BullMQ worker, job orchestrator, and review runtime - #123

Closed
onel wants to merge 5 commits into
devarshishimpi:devfrom
onel:feat/node-bullmq-worker
Closed

onel wants to merge 5 commits into
devarshishimpi:devfrom
onel:feat/node-bullmq-worker

Conversation

@onel

@onel onel commented Sep 22, 2026

Copy link
Copy Markdown

Closes #107. Closes #106.

Queued reviews are now processed on the Node target: the BullMQ worker consumes codra-reviews and drives the engine to completion.

Why both issues

#107 alone doesn't build. It calls NodeOrchestrator from #106, and #106 calls a createReviewRuntime that didn't exist for Node — apps/node/src/api-deps.ts threw Review runtime not implemented for Node yet. No sub-issue covered building one, so this PR does all three, in separate commits:

Commit
282c08a Node review runtime and telemetry sink
d7c3cc0 NodeOrchestrator (#106)
14aee3b BullMQ worker and process lifecycle (#107)
be5e3d4 Self-review fixes

api-deps.ts no longer has throwing stubs, so the dashboard's diff view works on Node too.

Notes on the implementation

The runtime mirrors apps/worker/src/adapters/. The stores take a DbEnv rather than Cloudflare bindings, so most of it is a re-wire. Two pieces were worker-local: the telemetry sink, rebuilt on @codraoss/db's existing InstanceIdStore, and the repo config loader, which comes from @codraoss/api/platform. There's no Workers AI binding, so Cloudflare-hosted models are unavailable on this target; every other provider is reached over HTTP as before.

The orchestrator runs the phase machine in a plain loop. The Cloudflare driver 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 that collapses into awaiting the delay and calling again in-process, and freshInstance needs no handling. Errors propagate: BullMQ owns retries and dead-lettering.

Lifecycle. START_API and START_WORKER both default on, so one container is the whole deployment; either set to false splits API and workers across containers, and both false is a startup error. SIGTERM closes the listener, drops idle keep-alive sockets, and waits for worker.close() so the review in flight finishes.

One change outside apps/node

@codraoss/db gains closeDb(). It closes the pooled fallback clients only — tracking runWithDb's would retain a handle per request on Workers — and nothing on the Worker path calls it. It's needed because a long-lived process has to end that socket to exit. Happy to solve it inside the app instead if you'd rather keep the package untouched.

Self-review

I ran a review over the diff and fixed six of seven findings. The one worth calling out was measured, not theorised:

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. That's right on Workers, where an invocation may not reuse another's I/O, and wrong for a server that stays up. It was already leaking one per HTTP request on dev, via the request path in index.ts. Both now use the process-wide pool, and shutdown ends it. Same 20 jobs, one connection.

Also fixed: an unbounded retry loop that could wedge a worker slot forever, since nothing reaps stale running rows on this target; result.jobId dropped between phases, which made every phase boundary re-resolve the webhook delivery and, for comment events, call GitHub again; BullMQ concurrency left at its default of 1; a 30s shutdown backstop that would have killed a running review on every redeploy; a missing NODE_ENV/VITEST guard on telemetry; and a comment about terminateJobWorkflow that had the ordering backwards.

Not fixed, and worth its own issue. terminateJobWorkflow is a no-op on Node — there is no durable instance to terminate. /stop and /rerun call it before cancelling, to stop the old run racing the new one. With a no-op the loop keeps working the current phase until it next reads the job's status, so stop is eventually-effective rather than immediate. Closing that needs a cooperative cancel signal the loop polls.

Job maintenance (scheduleBestEffortJobMaintenance, lease recovery) is still a no-op here, inherited from the scaffold. The retry bound stops it wedging a worker, but expired leases are never reaped on this target.

Verification

447 tests, lint, typecheck and typecheck:all pass. Against real Postgres and Redis:

  • 20 jobs enqueued and dequeued; each ran through the orchestrator into the engine and the database
  • an invalid payload failed once via UnrecoverableError instead of spending the backoff schedule (attemptsMade=1 with attempts: 3)
  • START_WORKER=false serves the API with no worker; START_API=false runs the worker with no listener; both false exits 1 with a clear message
  • SIGTERM after 10 completed jobs: exits 0, shutdown logs intact, one Postgres connection held

Relationship to #121

Stacked on #121 (Docker stack). Both touch apps/node/src/index.ts in different places. #121 should land first.

One follow-up for #121 once this merges: its compose file needs a stop_grace_period longer than a review phase, or Docker will SIGKILL the container at 10s and cut a running review short. The telemetry block I removed from its .env.docker.example can also come back, since this adds a working sink.

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.
Closes devarshishimpi#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.
Closes devarshishimpi#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.
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.

@codra-app-personal codra-app-personal Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

Codra Review

✅ Nothing to flag. Reviewed 8 files (496 changed lines) and found no issues worth raising.

Reviewed commit: be5e3d4fa1

ℹ️ About Codra in GitHub

Your team has set up Codra to review pull requests in this repo. Reviews are triggered when you:

  • Open a pull request for review
  • Mark a draft as ready

Every review posts a summary here. A clean pass also gets a 👍 on the pull request itself.

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.
@onel

onel commented Sep 22, 2026

Copy link
Copy Markdown
Author

Pushed 0c2a4a5 after a second review pass. One of the findings contradicts a claim in the description above, so flagging it rather than quietly fixing it.

Jobs were never retried. The description says errors propagate because "BullMQ owns retries and dead-lettering." That was wrong: the queue set no attempts, BullMQ defaults to attempts: 0, and its retry check is attemptsMade + 1 < opts.attempts, which is 1 < 0. Every job got exactly one run.

That matters more than it first looks. runReview's try block doesn't cover resolveQueuedJob, the admission query, or claimJobLease — so a transient Postgres error in any of those escapes the engine, and with no retry the job row is stranded in queued with nothing on this target to revisit it. The Cloudflare driver doesn't have this exposure: it runs each phase under retries: { limit: 5, delay: '60 seconds', backoff: 'exponential' }.

The queue now sets the same five attempts with 60s exponential backoff. UnrecoverableError still bypasses it, so an unparseable payload fails once. I also added removeOnComplete/removeOnFail bounds, since BullMQ keeps finished jobs forever by default and a long-lived install would grow Redis without limit.

Two smaller ones:

  • The worker constructed its own NodeOrchestrator while env.REVIEW_ORCHESTRATOR was wired and read by nothing. Worse, that field was initialised with a no-op stub before the real orchestrator replaced it, so anything arriving in between would have been silently dropped. The worker now goes through the port, and the stub is gone.
  • In packages/db/src/client.ts, the stale-connection branch overwrote the tracked closer for a connection string, discarding the reference to the socket it replaced. Inert today, since only the Node shutdown path calls closeDb() and Node never reaches that branch, but it would strand a handle if closeDb() were ever used from a worker-mode host.

Verified: 447 tests, lint, typecheck, typecheck:all pass. Confirmed the retry defaults reach added jobs, and re-ran the functional checks against real Postgres and Redis — 10 jobs processed, one connection held, SIGTERM exits 0.

@codra-app-personal codra-app-personal Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

Codra Review

Here are some automated review suggestions for this pull request.

Reviewed commit: 0c2a4a59b7

ℹ️ About Codra in GitHub

Your team has set up Codra to review pull requests in this repo. Reviews are triggered when you:

  • Open a pull request for review
  • Mark a draft as ready

If Codra has suggestions, it will comment; otherwise it will react with 👍.

Comment thread apps/node/src/worker.ts
// 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);

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P1 Unexpected Concurrency if Environment Variable is Empty String

The current parsing of WORKER_CONCURRENCY can lead to unintended behavior. If process.env.WORKER_CONCURRENCY is set to an empty string (""), Number("") evaluates to 0. BullMQ interprets a concurrency value of 0 as 'no limit' (unlimited concurrency), which contradicts the explicit default of 4 and could lead to resource exhaustion if the worker processes too many jobs simultaneously. This differs from the intent of providing a sensible default or adhering to a specific concurrency level.

Suggested change
const WORKER_CONCURRENCY = Number(process.env.WORKER_CONCURRENCY ?? 4);
const rawConcurrency = process.env.WORKER_CONCURRENCY;
const parsedConcurrency = Number(rawConcurrency);
const WORKER_CONCURRENCY = (Number.isFinite(parsedConcurrency) && parsedConcurrency >= 0) ? parsedConcurrency : 4;

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.

Fixed in f082257, thanks — real bug, though the mechanism is not the one described here, and the suggested patch would not have fixed it.

BullMQ has no "unlimited" concurrency mode. Its setter rejects anything below 1:

// node_modules/bullmq/dist/cjs/classes/worker.js:224
set concurrency(concurrency) {
  if (typeof concurrency !== 'number' || concurrency < 1 || !isFinite(concurrency)) {
    throw new Error('concurrency must be a finite number greater than 0');
  }

So 0 does not mean unlimited — it throws from the Worker constructor and the app dies at boot. Confirmed against the installed version: WORKER_CONCURRENCY set to "", "abc" or "0" all exit 1 with that error before the server listens.

The suggestion here keeps 0, since Number.isFinite(0) && 0 >= 0 is true, so it would still hit the same throw.

The fix takes an integer of 1 or more and otherwise falls back to the default. SHUTDOWN_TIMEOUT_MS had the same flaw in the other direction — a blank value gave a 0ms shutdown budget, force-exiting the instant SIGTERM landed and abandoning the review in flight — so both now go through one helper.

Comment thread apps/node/src/api-deps.ts
// 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.

@devarshishimpi

Copy link
Copy Markdown
Owner

Hey @onel Closing this. Just saw your PR, I have one in working as well. #122
Also, have to mention this PR is highly vibe-coded so it would pretty difficult to merge regardless.

onel added a commit to onel/codra that referenced this pull request Sep 22, 2026
Number('') is 0 and Number('abc') is NaN, so WORKER_CONCURRENCY set to either
was read as a real value. BullMQ's concurrency setter rejects anything below 1
and throws from the Worker constructor, so the app died at boot with 'concurrency
must be a finite number greater than 0' rather than using its default.
SHUTDOWN_TIMEOUT_MS had the same flaw in the other direction: a blank value gave
a 0ms budget, force-exiting the instant SIGTERM arrived and abandoning the review
in flight.

Both now go through a helper that takes an integer of 1 or more and otherwise
keeps the default.

Raised by the Codra review on devarshishimpi#123, though for a different reason: it read a
concurrency of 0 as BullMQ's 'unlimited'. BullMQ has no such mode.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants