Skip to content

fix(worker): a job that completed while the queue was unreachable is neither reported failed nor run again - #578

Merged
edgehero merged 1 commit into
mainfrom
fix/completed-job-lost-lock
Oct 4, 2026
Merged

edgehero merged 1 commit into
mainfrom
fix/completed-job-lost-lock

Conversation

@edgehero

@edgehero edgehero commented Oct 4, 2026 •

Copy link
Copy Markdown
Owner

A job that finished while Valkey was unreachable could be reported failed, and a scheduled one could run a second time, paid. Found while testing #577; no issue was filed for it.

What happened

When Valkey stays away longer than BullMQ's lock renewal window while a job is finishing, the processor still writes its run record, but BullMQ can no longer mark the job completed: moveToCompleted fails with "Missing lock" and the job stays active without a lock. The stall checker then acts on it:

  • a plain job gets BullMQ's deferred failure "job stalled more than allowable limit", so the worker posted the failure comment and paged the operator for a job that had completed;
  • a scheduled job is put back to wait and runs again, a second paid run, and its record overwrites the first one.

The fix

The run record is already a durable store: one file per job id carrying the attempt that wrote it, so no new store is needed.

  • When a job fails with exactly that stall reason and a record for that same attempt exists with an outcome other than failed, the worker logs job_lost_lock_after_completion as the job's terminal line and posts nothing. A recorded paid policy stop still pages, because the first finish never reached the completed handler that would have paged it.
  • When a stalled job comes back to the processor and its record already settles this attempt, the recorded result is returned without starting a container, reserving anything or writing a new record.
  • On a fleet with PI_WORKER_NAME, the run mirror in Valkey is the second place to look when the local file does not settle the attempt. Any fault there counts as no record, which keeps today's behaviour.
  • A record only counts when its job id and attempt match and it started no earlier than five minutes before the job was created (the two times come from two hosts' clocks), so a reused job id never matches an old record in practice. When a stall failure finds a record it does not accept, job_lost_lock_record_rejected names why (unreadable, id-mismatch, other-attempt, failed or older-than-job) and where it was read, so a comment that is still posted can be explained.

Specs

AMENDED: DES-TERMINAL-COMMENTS-AND-FAILURE-HOOK (the residual that accepted one duplicate is replaced by the record check; what remains is named: a fleet host that exits before Valkey returns or another host failing the job before the writer reconnects, clock skew beyond the tolerance, a run that dies between its comment and its record, and a suppressed job still listed in BullMQ's failed set while its record says completed), CONST-RETRY-INFRA-ONLY (Why), REQ-LOCAL-JOB-VISIBILITY (the new terminal line). UNCHANGED, checked: INT-RUN-HISTORY-FILE-CONTRACT, INT-ON-FAILURE-HOOK-CONTRACT.

Tests

The listener cases (suppressed, attempt mismatch, a failed record, another failure reason, the mirror fallback, the policy page), the processor short circuit, the record reader never throwing, BullMQ's stall reason pinned against the installed source, and a live Valkey test that deletes a running job's lock for a plain and a scheduled job and proves exactly one run, no failure comment and no page. Each rule was checked by reverting it and watching its test fail.

@edgehero
edgehero force-pushed the fix/completed-job-lost-lock branch from 712f0d0 to 1dbc5f9 Compare October 4, 2026 13:43
…neither reported failed nor run again

When Valkey is unreachable for longer than BullMQ's lock renewal window
while a job's processor finishes, the processor still writes the run
record, but moveToCompleted is refused with "Missing lock". The job stays
active without a lock, and the stall check (maxStalledCount 0) takes it
back:

- a plain job gets the deferred failure "job stalled more than allowable
  limit", so the failed listener posted the failure comment and paged the
  operator through PI_ON_FAILURE for a job that completed;
- a scheduled job is moved back to wait and ran again, paid, and its
  second record overwrote the first.

The run record is the store that already knows the job finished, so it is
read back rather than adding an idempotence store:

- run-history.mjs: makeReadRecord reads <sanitizedJobId>.json by the
  writer's own name (null for no file, UNREADABLE_RECORD for a file that
  is not a record), and recordVerdict / makeSettledRecord decide whether a
  record settles this attempt: same id, attempt equal, an outcome other
  than failed, and a startedAt no earlier than the job's creation less
  RECORD_CLOCK_SKEW_MS (5 minutes: the two instants come from two hosts'
  clocks, and a producer ahead of the worker must not turn the check off,
  while a reused id's old record is older by far more). A record found and
  refused is reported with a fixed reason (unreadable, id-mismatch,
  other-attempt, failed, older-than-job) and its source (local, mirror),
  and both callers log it as job_lost_lock_record_rejected so an operator
  can see why the comment was posted or the job ran. The processor's gate
  logs it once per job id, stall count, attempt and source (a bounded set
  of REJECTED_SEEN_MAX, 1000): it runs before every deferral and BullMQ
  never resets a stall count, so a held job meets it on each pickup with
  the same verdict. Never throws.
- run-mirror.mjs: readMirroredRecord, one bounded GET by the sanitized
  id. On a declared fleet the lookup falls back to it when the local file
  does not settle the attempt, since the host that meets the stalled job
  need not be the one that ran it. A mirror fault is no record; a value
  that is there but is not a record is UNREADABLE_RECORD, as locally. The
  mirror write itself survives the outage (the shared client's offline
  queue sends it on reconnect).
- start.mjs: on a terminal failure whose message is exactly BullMQ's stall
  reason (STALLED_FAILED_REASON in index.mjs, pinned by a test against the
  installed bullmq's moveStalledJobsToWait), the failed listener asks the
  record with attempt = attemptsMade (BullMQ has already added one). When
  it settles, the terminal line is job_lost_lock_after_completion, nothing
  is posted, and the hook fires only when the record is a paid policy stop
  the completed listener would have paged for (one shared predicate),
  because that listener never saw the first finish. Every other failure
  runs the unchanged body synchronously.
- index.mjs: the processor's first gate. A job with stalledCounter > 0
  whose record settles attempt attemptsMade + 1 logs
  job_lost_lock_after_completion and returns the recorded result (with the
  record's own budgetReserved, so the completed listener pages exactly as
  it would have) without a container, a reservation or a new record. No
  record keeps today's path. The per-scheduler stall guard still counts
  the stall, since the queue really lost the lock.

Specs: DES-TERMINAL-COMMENTS-AND-FAILURE-HOOK Residuals amended (the
"bounded at one duplicate, not worth the idempotence store" residual is
replaced by the record-read suppression; the remaining residuals are
named: a fleet without a shared logs dir whose mirror holds no copy,
because the writing host exited before Valkey returned or another host
failed the job before the writer reconnected; a clock skew beyond the
tolerance; a run that dies between its comment and its record; and a
suppressed job still sitting in BullMQ's failed set with the stall reason,
so admin views and getJobCounts show it failed while the record says it
completed). CONST-RETRY-INFRA-ONLY Why
amended (a scheduled job that stalls after its completion was recorded is
returned without running). REQ-LOCAL-JOB-VISIBILITY Acceptance amended
(job_lost_lock_after_completion is a terminal line, and
job_lost_lock_record_rejected names a refused record).
INT-RUN-HISTORY-FILE-CONTRACT and INT-ON-FAILURE-HOOK-CONTRACT UNCHANGED,
checked.

Tests: processor gate units (lost-lock.test.mjs, including the record's
own budgetReserved: a free model-not-allowed refusal stays false), the
reader, the verdicts, the tolerance and the lookup's reports
(run-history), the mirror read (run-mirror), start-wiring cases
(suppressed; attempt mismatch, failed record, other reason, reused id and
no record all still comment and page, each refused record logged with its
reason; a skew within the tolerance still suppresses; an unreadable record
is named, from the file or the mirror; five pickups of a held job log one
rejection; a paid policy record pages once; the mirror fallback), and a live-Valkey test that
deletes bull:<queue>:<id>:lock while the processor runs, for a plain and a
scheduled job: one run each, no failure comment, no page.

Signed-off-by: Rob Boerman <robboerman@live.nl>
@edgehero
edgehero force-pushed the fix/completed-job-lost-lock branch from 1dbc5f9 to 452dc70 Compare October 4, 2026 13:52
@edgehero
edgehero merged commit 2ab56a5 into main Oct 4, 2026
7 checks passed
@edgehero
edgehero deleted the fix/completed-job-lost-lock branch October 4, 2026 20:10
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant