Skip to content

Batched completeJobs/failJobs are never retried: the retry loop in batch() exits after the first failure #636

Description

@frvi

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

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions