Summary
batch() in src/main.ts wraps batched completions and failures in a retry loop (BATCH_RETRY_OPTIONS, 20 attempts, with shorter delays for retryable codes such as 40P01 deadlock_detected). But the loop's catch ends with break, so it leaves the loop after the first failure and throws. None of the retry code runs, and a single transient error (a deadlock, a serialization failure, a dropped connection) makes the worker pool shut down with Could not completeJobs; queue is in an inconsistent state; aborting.
try {
const result = await rawCallback(specs);
allgood();
return result;
} catch (e) {
holdup();
lastError = coerceError(e);
break; // <- exits the retry loop
}
The comment on the failure path ("Tell other callers to wait until we're successful again") and the allgood() call on success also assume the loop tries again: with the break, the backpressure raised by holdup() is never released.
The loop has never retried in a release. 5a75d76 ("Give 'batch' retry logic with backpressure") introduced it with throw e; in the catch (0.17.x), and eda912b ("Manually fix legit lint errors") changed that to break; (0.18.0); both leave the loop on the first failure. It affects any setup with completeJobBatchDelay or failJobBatchDelay >= 0.
One way to get such a deadlock in practice: a transaction re-queues two running jobs by job_key (replace mode updates their rows) while the worker completes the same two jobs in one batch, locking the rows in the opposite order. A retry once the other transaction has committed would normally go through; instead the pool stops.
Steps to reproduce
Fail the first completion attempt with 40P01; later attempts succeed:
test("a batched completion that hits a deadlock is retried", () =>
withOptions(async (options) => {
const { pgPool } = options;
await reset(pgPool, options);
await pgPool.query(`
create sequence public.completion_attempts;
create or replace function public.deadlock_once() returns trigger
language plpgsql as $$
begin
if nextval('public.completion_attempts') = 1 then
raise exception 'simulated deadlock' using errcode = '40P01';
end if;
return old;
end;
$$;
create trigger deadlock_once
before delete on ${ESCAPED_GRAPHILE_WORKER_SCHEMA}._private_jobs
for each row execute function public.deadlock_once();
`);
const runner = await run({
...options,
taskList: { job1: async () => {} },
preset: { worker: { completeJobBatchDelay: 0 } },
});
let settled = false;
runner.promise.then(() => (settled = true), () => (settled = true));
try {
await runner.addJob("job1", { id: "1" });
let remaining = 1;
for (let i = 0; i < 50 && remaining > 0; i++) {
await sleep(100);
remaining = (
await pgPool.query(`select count(*)::int as n from ${ESCAPED_GRAPHILE_WORKER_SCHEMA}._private_jobs`)
).rows[0].n;
}
expect({ remaining, runnerStopped: settled }).toEqual({ remaining: 0, runnerStopped: false });
} finally {
await runner.stop().catch(() => {});
}
}));
On main (4cda192) there is no retry; the pool stops at once and the job is left locked:
[core] ERROR: Failed to complete jobs '1':
[core] WARNING: Runner stopping (reason: worker pool exited with error: error: simulated deadlock)
- "remaining": 0,
- "runnerStopped": false,
+ "remaining": 1,
+ "runnerStopped": true,
Possible fix
Remove the break, so a failed attempt falls through to the next iteration; the loop already throws lastError once maxAttempts is reached:
} catch (e) {
// Tell other callers to wait until we're successful again (i.e. apply backpressure)
holdup();
lastError = coerceError(e);
- break;
}
With it, the test passes:
completeJobs: attempt 1/20 failed with code "40P01"; retrying after 43ms. Error: error: simulated deadlock
and the rest of the suite is unchanged locally. Note that non-retryable errors are then also retried up to 20 times with the default backoff before the pool gives up, which is what BATCH_RETRY_OPTIONS describes.
Additional context
Summary
batch()insrc/main.tswraps batched completions and failures in a retry loop (BATCH_RETRY_OPTIONS, 20 attempts, with shorter delays for retryable codes such as40P01deadlock_detected). But the loop'scatchends withbreak, so it leaves the loop after the first failure and throws. None of the retry code runs, and a single transient error (a deadlock, a serialization failure, a dropped connection) makes the worker pool shut down withCould not completeJobs; queue is in an inconsistent state; aborting.The comment on the failure path ("Tell other callers to wait until we're successful again") and the
allgood()call on success also assume the loop tries again: with thebreak, the backpressure raised byholdup()is never released.The loop has never retried in a release. 5a75d76 ("Give 'batch' retry logic with backpressure") introduced it with
throw e;in thecatch(0.17.x), and eda912b ("Manually fix legit lint errors") changed that tobreak;(0.18.0); both leave the loop on the first failure. It affects any setup withcompleteJobBatchDelayorfailJobBatchDelay>= 0.One way to get such a deadlock in practice: a transaction re-queues two running jobs by
job_key(replace mode updates their rows) while the worker completes the same two jobs in one batch, locking the rows in the opposite order. A retry once the other transaction has committed would normally go through; instead the pool stops.Steps to reproduce
Fail the first completion attempt with
40P01; later attempts succeed:On
main(4cda192) there is no retry; the pool stops at once and the job is left locked:Possible fix
Remove the
break, so a failed attempt falls through to the next iteration; the loop already throwslastErroroncemaxAttemptsis reached:} catch (e) { // Tell other callers to wait until we're successful again (i.e. apply backpressure) holdup(); lastError = coerceError(e); - break; }With it, the test passes:
and the rest of the suite is unchanged locally. Note that non-retryable errors are then also retried up to 20 times with the default backoff before the pool gives up, which is what
BATCH_RETRY_OPTIONSdescribes.Additional context
breakdefeats that on the batched path.runner.promiseresolves when the worker pool exits with an error #635: once the pool stops,runner.promiseresolves instead of rejecting.