Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/postgres-post-ceiling-job-attempts.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
'@workflow/world-postgres': patch
---

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.
5 changes: 5 additions & 0 deletions packages/core/src/runtime/constants.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,11 @@ import { runtimeLogger } from '../logger.js';
// backend outage; conversely, spanning the full 24h window would require a
// substantially higher cap here, not a higher per-hop ceiling, since VQS
// clamps every hop at 900s.)
//
// world-postgres sizes its Graphile job attempt cap from this value
// (`CORE_MAX_DELIVERIES_EXCEEDED_ATTEMPT` = this + 1, plus headroom for
// post-ceiling redeliveries of the terminal write). Update it there too if
// this changes.
export const MAX_QUEUE_DELIVERIES = 48;

/**
Expand Down
47 changes: 45 additions & 2 deletions packages/world-postgres/src/queue.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -486,7 +486,7 @@ describe('postgres queue http execution', () => {
}),
expect.objectContaining({
jobKey: 'step_01ABC',
maxAttempts: 49,
maxAttempts: 73,
runAt: new Date('2024-01-01T00:00:05.000Z'),
})
);
Expand Down Expand Up @@ -576,6 +576,36 @@ describe('postgres queue http execution', () => {
}
});

it('keeps post-ceiling headroom when transferring a job on the default cap', async () => {
const queue = buildQueue(
{ connectionString: 'postgres://test', enableInvoke: true },
pool
);
await queue.start();
const fetchMock = vi
.spyOn(nodeHttp, 'nodeHttpFetch')
.mockResolvedValue(Response.json({ ok: true }));
try {
const payload = buildMessageData('__wkf_workflow_example', {
runId: 'run_a',
});
// Delivery 49 is where core records MAX_DELIVERIES_EXCEEDED. A job with
// no stored cap takes the default, so the transferred job must still
// have attempts left for core's post-ceiling redeliveries (73 - 49 + 1).
await getTaskHandler('workflow_flows')(payload, {
job: { attempts: 49 },
});
expect(fetchMock).not.toHaveBeenCalled();
expect(workerUtilsMock.addJob).toHaveBeenCalledWith(
'workflow_flows_executor',
expect.objectContaining({ attempt: 49, attemptOffset: 48 }),
expect.objectContaining({ maxAttempts: 25 })
);
} finally {
fetchMock.mockRestore();
}
});

it('does not acknowledge a legacy transfer when enqueue fails', async () => {
const queue = buildQueue(
{ connectionString: 'postgres://test', enableInvoke: true },
Expand Down Expand Up @@ -843,10 +873,23 @@ describe('postgres queue http execution', () => {
}),
expect.objectContaining({
jobKey: 'step_01ABC',
maxAttempts: 49,
maxAttempts: 73,
})
);
});

it('leaves job attempts for redeliveries past core max deliveries', async () => {
// Core records MAX_DELIVERIES_EXCEEDED on delivery 49 and throws when that
// terminal write fails transiently. The job must still have attempts left
// for the redelivery, or the run is stranded `running`.
const queue = buildQueue({ connectionString: 'postgres://test' }, pool);
await queue.start();

await queue.queue('__wkf_workflow_example', { runId: 'run_01ABC' });

const [, , options] = vi.mocked(workerUtilsMock.addJob).mock.calls[0];
expect(options?.maxAttempts).toBeGreaterThan(49);
});
});

function buildQueue(
Expand Down
16 changes: 14 additions & 2 deletions packages/world-postgres/src/queue.ts
Original file line number Diff line number Diff line change
Expand Up @@ -137,8 +137,20 @@ export function getDeliveryTimeouts() {
};
}
const COMPLETED_IDEMPOTENCY_CACHE_LIMIT = 10_000;
// Core records MAX_DELIVERIES_EXCEEDED on delivery 49.
const MAX_GRAPHILE_JOB_ATTEMPTS = 49;
// Core records MAX_DELIVERIES_EXCEEDED on delivery 49 (MAX_QUEUE_DELIVERIES +
// 1). Past that delivery, core only retries its terminal write, and it throws
// on a transient failure (429 / 5xx / transport) so the queue redelivers
// rather than acking and leaving the run `running`. Those redeliveries need
// attempts left on the job: with a cap of exactly 49, the first post-ceiling
// throw would retire the job and strand the run anyway. Graphile's retry
// backoff is exp(min(attempts, 10)) seconds, so each extra attempt waits ~6h
// and 24 of them keep retrying the terminal write for ~6 days.
// Mirrors `MAX_QUEUE_DELIVERIES + 1` in @workflow/core (runtime/constants.ts),
// which this package does not depend on. Keep the two in sync.
const CORE_MAX_DELIVERIES_EXCEEDED_ATTEMPT = 49;
const POST_CEILING_RETRY_ATTEMPTS = 24;
const MAX_GRAPHILE_JOB_ATTEMPTS =
CORE_MAX_DELIVERIES_EXCEEDED_ATTEMPT + POST_CEILING_RETRY_ATTEMPTS;
const EXECUTOR_JOB_HEADER = 'x-workflow-postgres-executor-job';
const EXECUTOR_WORKER_HEADER = 'x-workflow-postgres-executor-worker';
const EXECUTOR_ATTEMPT_HEADER = 'x-workflow-postgres-executor-attempt';
Expand Down
Loading