Skip to content

Commit efb58a0

Browse files
fix(world-postgres): leave job attempts past core's max-deliveries ceiling (#4428)
* fix(world-postgres): leave job attempts past core's max-deliveries ceiling Core writes run_failed (MAX_DELIVERIES_EXCEEDED) on delivery 49 and, since #4421, throws when that write fails transiently so the queue redelivers. Graphile jobs were capped at exactly 49 attempts, so that throw retired the job and the run stayed running. Add 24 attempts of headroom (~6h apart under Graphile's exp(min(attempts, 10))s backoff). Fixes #4427 Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Co-Authored-By: Pranay Prakash <1797812+pranaygp@users.noreply.github.com> * fix(world-postgres): address review on post-ceiling job attempts - Cross-reference core's MAX_QUEUE_DELIVERIES and world-postgres' CORE_MAX_DELIVERIES_EXCEEDED_ATTEMPT so the two cannot drift silently. - Changeset notes that jobs enqueued before this release keep their stored 49-attempt cap, and executor transfers carry it forward. - Test a transfer of a default-cap job at delivery 49: it keeps 25 attempts (attemptOffset 48) for core's post-ceiling redeliveries. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Co-Authored-By: Pranay Prakash <1797812+pranaygp@users.noreply.github.com> * ci: retrigger (world-postgres pg-teardown and world-local storage flakes) Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Co-Authored-By: Pranay Prakash <1797812+pranaygp@users.noreply.github.com> --------- Co-authored-by: vercel[bot] <35613825+vercel[bot]@users.noreply.github.com> Co-authored-by: Pranay Prakash <1797812+pranaygp@users.noreply.github.com>
1 parent ee1a09b commit efb58a0

4 files changed

Lines changed: 69 additions & 4 deletions

File tree

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
1+
---
2+
'@workflow/world-postgres': patch
3+
---
4+
5+
Leave job attempts past core's max-deliveries ceiling, so a run whose terminal `run_failed` write fails transiently at the ceiling is retried instead of stranded when the job runs out of attempts. Jobs enqueued before this release keep their stored cap of 49 attempts (executor transfers carry it forward); only jobs enqueued after upgrading get the headroom.

‎packages/core/src/runtime/constants.ts‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,11 @@ import { runtimeLogger } from '../logger.js';
2020
// backend outage; conversely, spanning the full 24h window would require a
2121
// substantially higher cap here, not a higher per-hop ceiling, since VQS
2222
// clamps every hop at 900s.)
23+
//
24+
// world-postgres sizes its Graphile job attempt cap from this value
25+
// (`CORE_MAX_DELIVERIES_EXCEEDED_ATTEMPT` = this + 1, plus headroom for
26+
// post-ceiling redeliveries of the terminal write). Update it there too if
27+
// this changes.
2328
export const MAX_QUEUE_DELIVERIES = 48;
2429

2530
/**

‎packages/world-postgres/src/queue.test.ts‎

Lines changed: 45 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -486,7 +486,7 @@ describe('postgres queue http execution', () => {
486486
}),
487487
expect.objectContaining({
488488
jobKey: 'step_01ABC',
489-
maxAttempts: 49,
489+
maxAttempts: 73,
490490
runAt: new Date('2024-01-01T00:00:05.000Z'),
491491
})
492492
);
@@ -576,6 +576,36 @@ describe('postgres queue http execution', () => {
576576
}
577577
});
578578

579+
it('keeps post-ceiling headroom when transferring a job on the default cap', async () => {
580+
const queue = buildQueue(
581+
{ connectionString: 'postgres://test', enableInvoke: true },
582+
pool
583+
);
584+
await queue.start();
585+
const fetchMock = vi
586+
.spyOn(nodeHttp, 'nodeHttpFetch')
587+
.mockResolvedValue(Response.json({ ok: true }));
588+
try {
589+
const payload = buildMessageData('__wkf_workflow_example', {
590+
runId: 'run_a',
591+
});
592+
// Delivery 49 is where core records MAX_DELIVERIES_EXCEEDED. A job with
593+
// no stored cap takes the default, so the transferred job must still
594+
// have attempts left for core's post-ceiling redeliveries (73 - 49 + 1).
595+
await getTaskHandler('workflow_flows')(payload, {
596+
job: { attempts: 49 },
597+
});
598+
expect(fetchMock).not.toHaveBeenCalled();
599+
expect(workerUtilsMock.addJob).toHaveBeenCalledWith(
600+
'workflow_flows_executor',
601+
expect.objectContaining({ attempt: 49, attemptOffset: 48 }),
602+
expect.objectContaining({ maxAttempts: 25 })
603+
);
604+
} finally {
605+
fetchMock.mockRestore();
606+
}
607+
});
608+
579609
it('does not acknowledge a legacy transfer when enqueue fails', async () => {
580610
const queue = buildQueue(
581611
{ connectionString: 'postgres://test', enableInvoke: true },
@@ -843,10 +873,23 @@ describe('postgres queue http execution', () => {
843873
}),
844874
expect.objectContaining({
845875
jobKey: 'step_01ABC',
846-
maxAttempts: 49,
876+
maxAttempts: 73,
847877
})
848878
);
849879
});
880+
881+
it('leaves job attempts for redeliveries past core max deliveries', async () => {
882+
// Core records MAX_DELIVERIES_EXCEEDED on delivery 49 and throws when that
883+
// terminal write fails transiently. The job must still have attempts left
884+
// for the redelivery, or the run is stranded `running`.
885+
const queue = buildQueue({ connectionString: 'postgres://test' }, pool);
886+
await queue.start();
887+
888+
await queue.queue('__wkf_workflow_example', { runId: 'run_01ABC' });
889+
890+
const [, , options] = vi.mocked(workerUtilsMock.addJob).mock.calls[0];
891+
expect(options?.maxAttempts).toBeGreaterThan(49);
892+
});
850893
});
851894

852895
function buildQueue(

‎packages/world-postgres/src/queue.ts‎

Lines changed: 14 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -137,8 +137,20 @@ export function getDeliveryTimeouts() {
137137
};
138138
}
139139
const COMPLETED_IDEMPOTENCY_CACHE_LIMIT = 10_000;
140-
// Core records MAX_DELIVERIES_EXCEEDED on delivery 49.
141-
const MAX_GRAPHILE_JOB_ATTEMPTS = 49;
140+
// Core records MAX_DELIVERIES_EXCEEDED on delivery 49 (MAX_QUEUE_DELIVERIES +
141+
// 1). Past that delivery, core only retries its terminal write, and it throws
142+
// on a transient failure (429 / 5xx / transport) so the queue redelivers
143+
// rather than acking and leaving the run `running`. Those redeliveries need
144+
// attempts left on the job: with a cap of exactly 49, the first post-ceiling
145+
// throw would retire the job and strand the run anyway. Graphile's retry
146+
// backoff is exp(min(attempts, 10)) seconds, so each extra attempt waits ~6h
147+
// and 24 of them keep retrying the terminal write for ~6 days.
148+
// Mirrors `MAX_QUEUE_DELIVERIES + 1` in @workflow/core (runtime/constants.ts),
149+
// which this package does not depend on. Keep the two in sync.
150+
const CORE_MAX_DELIVERIES_EXCEEDED_ATTEMPT = 49;
151+
const POST_CEILING_RETRY_ATTEMPTS = 24;
152+
const MAX_GRAPHILE_JOB_ATTEMPTS =
153+
CORE_MAX_DELIVERIES_EXCEEDED_ATTEMPT + POST_CEILING_RETRY_ATTEMPTS;
142154
const EXECUTOR_JOB_HEADER = 'x-workflow-postgres-executor-job';
143155
const EXECUTOR_WORKER_HEADER = 'x-workflow-postgres-executor-worker';
144156
const EXECUTOR_ATTEMPT_HEADER = 'x-workflow-postgres-executor-attempt';

0 commit comments

Comments
 (0)