fix(world-postgres): serialize keyed deliveries across workers - #3657

Closed
joeyhotz wants to merge 1 commit into
vercel:mainfrom
joeyhotz:fix/postgres-keyed-job-serialization
Closed

fix(world-postgres): serialize keyed deliveries across workers#3657
joeyhotz wants to merge 1 commit into
vercel:mainfrom
joeyhotz:fix/postgres-keyed-job-serialization

Conversation

@joeyhotz

Copy link
Copy Markdown
Contributor

Description

@workflow/world-postgres currently passes QueueOptions.idempotencyKey to Graphile Worker as jobKey. With Graphile's default replace mode, adding the same key while its current job is locked clears the locked job's key, exhausts that row, and inserts a successor. The existing inflightMessages map prevents overlap inside one process, but separate Postgres World processes can claim the current job and its successor at the same time.

We observed this after a self-hosted deployment scaled out: startup recovery replayed active runs and redispatched their keyed steps. In one staging burst, 7 of 34 materialisation steps received a second start without an intervening failure (41 executions total). The same startup-recovery/duplicate-start signature appeared in production. This behavior is present in @workflow/world-postgres 4.3.3 and current main.

This PR assigns every keyed Graphile job to a deterministic named queue. A named queue is Graphile's native cross-worker concurrency primitive: while one worker owns the queue, a replacement remains pending. After the first delivery commits and releases the queue, the successor can replay the terminal state instead of overlapping the original side effect.

The implementation deliberately:

  • Keeps Graphile's default replace mode. This preserves the durable successor created by delayed rescheduling and crash recovery.
  • Applies only to messages that carry idempotencyKey. Unkeyed orchestrator messages and public World interfaces are unchanged.
  • Uses 2,048 stable SHA-256 buckets per configured jobPrefix. Exact per-key queue names would grow Graphile's persistent queue table without bound; hashing the job task name keeps prefixes isolated and the physical name below Graphile's 128-character limit.
  • Applies the same mapping when migrating pg-boss jobs that actually carry MessageData.idempotencyKey. Legacy unkeyed rows used messageId as singleton_key and remain unqueued.
  • Does not run live GC_JOB_QUEUES; Graphile 0.16.6 cleanup raced concurrent producers in a local stress test and left 398 of 7,680 jobs orphaned (5.18%).

jobKeyMode: 'unsafe_dedupe' is not safe here. The handler enqueues a delayed successor before its current locked job returns; unsafe_dedupe would discard that successor, after which completing the current row can strand the workflow. It can also suppress the only replacement for a locked job that later dies on its final attempt.

This serializes competing deliveries; it does not claim exactly-once execution. The runtime's terminal-state replay remains the duplicate-suppression layer after the queue releases.

Operational trade-offs:

CaseBehavior
First mixed-version rollout or rollbackExisting/outgoing jobs without queueName do not own the new queue. Deployments must drain old keyed jobs and producers for a fully protected handoff.
Hard process deathA bucket can remain locked until Graphile's four-hour stale-lock recovery (or an explicit worker unlock). Graceful shutdown releases it normally.
Bucket collisionUnrelated keyed jobs can serialize. At the default concurrency of 50, 2,048 buckets produce about 0.60 expected colliding pairs, or roughly 1.2% expected slot loss.

A shuffled three-repetition Graphile 0.16.6 no-op benchmark at concurrency 50 / pool 8 measured median throughput of 181.5 jobs/s with no named queues, 141.1 jobs/s with all 2,048 buckets populated, and 108.2 jobs/s with 4,096. The VPS timings were noisy, so these are directional rather than a production capacity claim; 2,048 was consistently the better balance. Dropping to 1,024 would double expected collision loss at the default concurrency to roughly 2.4%.

This complements #3119 and #3162 but does not replace them. Those address selecting parked runs and accumulating startup-recovery jobs; this PR closes the separate locked-keyed-delivery concurrency gap. The queueName seam is independent of the HTTP loopback and can be carried through the in-process execution refactor in #3322.

How did you test your changes?

Added unit coverage for:

  • Keyed producer jobs receiving the deterministic named queue.
  • Delayed handler reschedules retaining the same jobKey, named queue, runAt, and 49-attempt budget.
  • Unkeyed messages remaining outside named queues.
  • Custom jobPrefix values receiving isolated queue scopes.
  • pg-boss migration distinguishing real idempotency keys from message-ID singleton keys.

Added a real-PostgreSQL Testcontainers regression with two independent pools and two createQueue() instances. It blocks worker A's HTTP delivery, enqueues the same key through worker B, waits through the polling window, and asserts that both rows share one named queue while only the first is locked. Before the fix, the same test observed two concurrent HTTP deliveries (maxActiveRequests: 2); with the fix it observes one (maxActiveRequests: 1) and then drains both jobs sequentially.

Verification:

  • pnpm exec vitest run packages/world-postgres/src (3 files, 30 tests)
  • Real-PostgreSQL concurrency regression against an isolated local PostgreSQL 17 database (1 test)
  • pnpm --filter @workflow/world-postgres typecheck
  • pnpm --filter @workflow/world-postgres build
  • Biome checks over the changed TypeScript files (only the two pre-existing warnings in queue.ts / queue.test.ts)
  • pnpm changeset status

Additional recovery probe with Graphile Worker 0.16.6: a child worker was killed while holding the named queue, a same-key replacement was added, and after forceUnlockWorkers only the replacement executed.

PR Checklist - Required to merge

  • 📦 pnpm changeset was run to create a changelog for this PR
    • Patch changeset for @workflow/world-postgres.
  • 🔒 DCO sign-off passes (run git commit --signoff on your commits)
  • 📝 Ping @vercel/workflow in a comment once the PR is ready, and the above checklist is complete

Draft for author review. No reviewers have been requested.

Assign idempotent jobs to bounded Graphile named queues so a locked delivery and its durable replacement cannot execute concurrently across worker processes. Keep default replacement semantics for delayed retries and crash recovery.
Signed-off-by: Joey Hotz <joeyhotz1@gmail.com>
@changeset-bot

Copy link
Copy Markdown

🦋 Changeset detected

Latest commit: 163dbf4

The changes in this PR will be included in the next version bump.

This PR includes changesets to release 1 package
NameType
@workflow/world-postgresPatch

Not sure what this means? Click here to learn what changesets are.

Click here if you're a maintainer who wants to add another changeset to this PR

@vercel

vercelBot commented Aug 19, 2026

Copy link
Copy Markdown
Contributor

@joeyhotz is attempting to deploy a commit to the Vercel Labs Team on Vercel.

A member of the Team first needs to authorize it.

Sign up for freeto 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

@joeyhotz
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content

fix(world-postgres): serialize keyed deliveries across workers - #3657

Closed
joeyhotz wants to merge 1 commit into
vercel:mainfrom
joeyhotz:fix/postgres-keyed-job-serialization
Closed

fix(world-postgres): serialize keyed deliveries across workers#3657
joeyhotz wants to merge 1 commit into
vercel:mainfrom
joeyhotz:fix/postgres-keyed-job-serialization

Conversation

@joeyhotz

Copy link
Copy Markdown
Contributor

Description

@workflow/world-postgres currently passes QueueOptions.idempotencyKey to Graphile Worker as jobKey. With Graphile's default replace mode, adding the same key while its current job is locked clears the locked job's key, exhausts that row, and inserts a successor. The existing inflightMessages map prevents overlap inside one process, but separate Postgres World processes can claim the current job and its successor at the same time.

We observed this after a self-hosted deployment scaled out: startup recovery replayed active runs and redispatched their keyed steps. In one staging burst, 7 of 34 materialisation steps received a second start without an intervening failure (41 executions total). The same startup-recovery/duplicate-start signature appeared in production. This behavior is present in @workflow/world-postgres 4.3.3 and current main.

This PR assigns every keyed Graphile job to a deterministic named queue. A named queue is Graphile's native cross-worker concurrency primitive: while one worker owns the queue, a replacement remains pending. After the first delivery commits and releases the queue, the successor can replay the terminal state instead of overlapping the original side effect.

The implementation deliberately:

  • Keeps Graphile's default replace mode. This preserves the durable successor created by delayed rescheduling and crash recovery.
  • Applies only to messages that carry idempotencyKey. Unkeyed orchestrator messages and public World interfaces are unchanged.
  • Uses 2,048 stable SHA-256 buckets per configured jobPrefix. Exact per-key queue names would grow Graphile's persistent queue table without bound; hashing the job task name keeps prefixes isolated and the physical name below Graphile's 128-character limit.
  • Applies the same mapping when migrating pg-boss jobs that actually carry MessageData.idempotencyKey. Legacy unkeyed rows used messageId as singleton_key and remain unqueued.
  • Does not run live GC_JOB_QUEUES; Graphile 0.16.6 cleanup raced concurrent producers in a local stress test and left 398 of 7,680 jobs orphaned (5.18%).

jobKeyMode: 'unsafe_dedupe' is not safe here. The handler enqueues a delayed successor before its current locked job returns; unsafe_dedupe would discard that successor, after which completing the current row can strand the workflow. It can also suppress the only replacement for a locked job that later dies on its final attempt.

This serializes competing deliveries; it does not claim exactly-once execution. The runtime's terminal-state replay remains the duplicate-suppression layer after the queue releases.

Operational trade-offs:

CaseBehavior
First mixed-version rollout or rollbackExisting/outgoing jobs without queueName do not own the new queue. Deployments must drain old keyed jobs and producers for a fully protected handoff.
Hard process deathA bucket can remain locked until Graphile's four-hour stale-lock recovery (or an explicit worker unlock). Graceful shutdown releases it normally.
Bucket collisionUnrelated keyed jobs can serialize. At the default concurrency of 50, 2,048 buckets produce about 0.60 expected colliding pairs, or roughly 1.2% expected slot loss.

A shuffled three-repetition Graphile 0.16.6 no-op benchmark at concurrency 50 / pool 8 measured median throughput of 181.5 jobs/s with no named queues, 141.1 jobs/s with all 2,048 buckets populated, and 108.2 jobs/s with 4,096. The VPS timings were noisy, so these are directional rather than a production capacity claim; 2,048 was consistently the better balance. Dropping to 1,024 would double expected collision loss at the default concurrency to roughly 2.4%.

This complements #3119 and #3162 but does not replace them. Those address selecting parked runs and accumulating startup-recovery jobs; this PR closes the separate locked-keyed-delivery concurrency gap. The queueName seam is independent of the HTTP loopback and can be carried through the in-process execution refactor in #3322.

How did you test your changes?

Added unit coverage for:

  • Keyed producer jobs receiving the deterministic named queue.
  • Delayed handler reschedules retaining the same jobKey, named queue, runAt, and 49-attempt budget.
  • Unkeyed messages remaining outside named queues.
  • Custom jobPrefix values receiving isolated queue scopes.
  • pg-boss migration distinguishing real idempotency keys from message-ID singleton keys.

Added a real-PostgreSQL Testcontainers regression with two independent pools and two createQueue() instances. It blocks worker A's HTTP delivery, enqueues the same key through worker B, waits through the polling window, and asserts that both rows share one named queue while only the first is locked. Before the fix, the same test observed two concurrent HTTP deliveries (maxActiveRequests: 2); with the fix it observes one (maxActiveRequests: 1) and then drains both jobs sequentially.

Verification:

  • pnpm exec vitest run packages/world-postgres/src (3 files, 30 tests)
  • Real-PostgreSQL concurrency regression against an isolated local PostgreSQL 17 database (1 test)
  • pnpm --filter @workflow/world-postgres typecheck
  • pnpm --filter @workflow/world-postgres build
  • Biome checks over the changed TypeScript files (only the two pre-existing warnings in queue.ts / queue.test.ts)
  • pnpm changeset status

Additional recovery probe with Graphile Worker 0.16.6: a child worker was killed while holding the named queue, a same-key replacement was added, and after forceUnlockWorkers only the replacement executed.

PR Checklist - Required to merge

  • 📦 pnpm changeset was run to create a changelog for this PR
    • Patch changeset for @workflow/world-postgres.
  • 🔒 DCO sign-off passes (run git commit --signoff on your commits)
  • 📝 Ping @vercel/workflow in a comment once the PR is ready, and the above checklist is complete

Draft for author review. No reviewers have been requested.

Assign idempotent jobs to bounded Graphile named queues so a locked delivery and its durable replacement cannot execute concurrently across worker processes. Keep default replacement semantics for delayed retries and crash recovery.
Signed-off-by: Joey Hotz <joeyhotz1@gmail.com>
@changeset-bot

Copy link
Copy Markdown

🦋 Changeset detected

Latest commit: 163dbf4

The changes in this PR will be included in the next version bump.

This PR includes changesets to release 1 package
NameType
@workflow/world-postgresPatch

Not sure what this means? Click here to learn what changesets are.

Click here if you're a maintainer who wants to add another changeset to this PR

@vercel

vercelBot commented Aug 19, 2026

Copy link
Copy Markdown
Contributor

@joeyhotz is attempting to deploy a commit to the Vercel Labs Team on Vercel.

A member of the Team first needs to authorize it.

Sign up for freeto 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

@joeyhotz
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

fix(world-postgres): serialize keyed deliveries across workers - #3657

Closed
joeyhotz wants to merge 1 commit into
vercel:mainfrom
joeyhotz:fix/postgres-keyed-job-serialization
Closed

fix(world-postgres): serialize keyed deliveries across workers#3657
joeyhotz wants to merge 1 commit into
vercel:mainfrom
joeyhotz:fix/postgres-keyed-job-serialization

Conversation

@joeyhotz

Copy link
Copy Markdown
Contributor

Description

@workflow/world-postgres currently passes QueueOptions.idempotencyKey to Graphile Worker as jobKey. With Graphile's default replace mode, adding the same key while its current job is locked clears the locked job's key, exhausts that row, and inserts a successor. The existing inflightMessages map prevents overlap inside one process, but separate Postgres World processes can claim the current job and its successor at the same time.

We observed this after a self-hosted deployment scaled out: startup recovery replayed active runs and redispatched their keyed steps. In one staging burst, 7 of 34 materialisation steps received a second start without an intervening failure (41 executions total). The same startup-recovery/duplicate-start signature appeared in production. This behavior is present in @workflow/world-postgres 4.3.3 and current main.

This PR assigns every keyed Graphile job to a deterministic named queue. A named queue is Graphile's native cross-worker concurrency primitive: while one worker owns the queue, a replacement remains pending. After the first delivery commits and releases the queue, the successor can replay the terminal state instead of overlapping the original side effect.

The implementation deliberately:

  • Keeps Graphile's default replace mode. This preserves the durable successor created by delayed rescheduling and crash recovery.
  • Applies only to messages that carry idempotencyKey. Unkeyed orchestrator messages and public World interfaces are unchanged.
  • Uses 2,048 stable SHA-256 buckets per configured jobPrefix. Exact per-key queue names would grow Graphile's persistent queue table without bound; hashing the job task name keeps prefixes isolated and the physical name below Graphile's 128-character limit.
  • Applies the same mapping when migrating pg-boss jobs that actually carry MessageData.idempotencyKey. Legacy unkeyed rows used messageId as singleton_key and remain unqueued.
  • Does not run live GC_JOB_QUEUES; Graphile 0.16.6 cleanup raced concurrent producers in a local stress test and left 398 of 7,680 jobs orphaned (5.18%).

jobKeyMode: 'unsafe_dedupe' is not safe here. The handler enqueues a delayed successor before its current locked job returns; unsafe_dedupe would discard that successor, after which completing the current row can strand the workflow. It can also suppress the only replacement for a locked job that later dies on its final attempt.

This serializes competing deliveries; it does not claim exactly-once execution. The runtime's terminal-state replay remains the duplicate-suppression layer after the queue releases.

Operational trade-offs:

CaseBehavior
First mixed-version rollout or rollbackExisting/outgoing jobs without queueName do not own the new queue. Deployments must drain old keyed jobs and producers for a fully protected handoff.
Hard process deathA bucket can remain locked until Graphile's four-hour stale-lock recovery (or an explicit worker unlock). Graceful shutdown releases it normally.
Bucket collisionUnrelated keyed jobs can serialize. At the default concurrency of 50, 2,048 buckets produce about 0.60 expected colliding pairs, or roughly 1.2% expected slot loss.

A shuffled three-repetition Graphile 0.16.6 no-op benchmark at concurrency 50 / pool 8 measured median throughput of 181.5 jobs/s with no named queues, 141.1 jobs/s with all 2,048 buckets populated, and 108.2 jobs/s with 4,096. The VPS timings were noisy, so these are directional rather than a production capacity claim; 2,048 was consistently the better balance. Dropping to 1,024 would double expected collision loss at the default concurrency to roughly 2.4%.

This complements #3119 and #3162 but does not replace them. Those address selecting parked runs and accumulating startup-recovery jobs; this PR closes the separate locked-keyed-delivery concurrency gap. The queueName seam is independent of the HTTP loopback and can be carried through the in-process execution refactor in #3322.

How did you test your changes?

Added unit coverage for:

  • Keyed producer jobs receiving the deterministic named queue.
  • Delayed handler reschedules retaining the same jobKey, named queue, runAt, and 49-attempt budget.
  • Unkeyed messages remaining outside named queues.
  • Custom jobPrefix values receiving isolated queue scopes.
  • pg-boss migration distinguishing real idempotency keys from message-ID singleton keys.

Added a real-PostgreSQL Testcontainers regression with two independent pools and two createQueue() instances. It blocks worker A's HTTP delivery, enqueues the same key through worker B, waits through the polling window, and asserts that both rows share one named queue while only the first is locked. Before the fix, the same test observed two concurrent HTTP deliveries (maxActiveRequests: 2); with the fix it observes one (maxActiveRequests: 1) and then drains both jobs sequentially.

Verification:

  • pnpm exec vitest run packages/world-postgres/src (3 files, 30 tests)
  • Real-PostgreSQL concurrency regression against an isolated local PostgreSQL 17 database (1 test)
  • pnpm --filter @workflow/world-postgres typecheck
  • pnpm --filter @workflow/world-postgres build
  • Biome checks over the changed TypeScript files (only the two pre-existing warnings in queue.ts / queue.test.ts)
  • pnpm changeset status

Additional recovery probe with Graphile Worker 0.16.6: a child worker was killed while holding the named queue, a same-key replacement was added, and after forceUnlockWorkers only the replacement executed.

PR Checklist - Required to merge

  • 📦 pnpm changeset was run to create a changelog for this PR
    • Patch changeset for @workflow/world-postgres.
  • 🔒 DCO sign-off passes (run git commit --signoff on your commits)
  • 📝 Ping @vercel/workflow in a comment once the PR is ready, and the above checklist is complete

Draft for author review. No reviewers have been requested.

Assign idempotent jobs to bounded Graphile named queues so a locked delivery and its durable replacement cannot execute concurrently across worker processes. Keep default replacement semantics for delayed retries and crash recovery.
Signed-off-by: Joey Hotz <joeyhotz1@gmail.com>
@changeset-bot

Copy link
Copy Markdown

🦋 Changeset detected

Latest commit: 163dbf4

The changes in this PR will be included in the next version bump.

This PR includes changesets to release 1 package
NameType
@workflow/world-postgresPatch

Not sure what this means? Click here to learn what changesets are.

Click here if you're a maintainer who wants to add another changeset to this PR

@vercel

vercelBot commented Aug 19, 2026

Copy link
Copy Markdown
Contributor

@joeyhotz is attempting to deploy a commit to the Vercel Labs Team on Vercel.

A member of the Team first needs to authorize it.

Sign up for freeto 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

@joeyhotz
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

fix(world-postgres): serialize keyed deliveries across workers - #3657

Closed
joeyhotz wants to merge 1 commit into
vercel:mainfrom
joeyhotz:fix/postgres-keyed-job-serialization
Closed

fix(world-postgres): serialize keyed deliveries across workers#3657
joeyhotz wants to merge 1 commit into
vercel:mainfrom
joeyhotz:fix/postgres-keyed-job-serialization

Conversation

@joeyhotz

Copy link
Copy Markdown
Contributor

Description

@workflow/world-postgres currently passes QueueOptions.idempotencyKey to Graphile Worker as jobKey. With Graphile's default replace mode, adding the same key while its current job is locked clears the locked job's key, exhausts that row, and inserts a successor. The existing inflightMessages map prevents overlap inside one process, but separate Postgres World processes can claim the current job and its successor at the same time.

We observed this after a self-hosted deployment scaled out: startup recovery replayed active runs and redispatched their keyed steps. In one staging burst, 7 of 34 materialisation steps received a second start without an intervening failure (41 executions total). The same startup-recovery/duplicate-start signature appeared in production. This behavior is present in @workflow/world-postgres 4.3.3 and current main.

This PR assigns every keyed Graphile job to a deterministic named queue. A named queue is Graphile's native cross-worker concurrency primitive: while one worker owns the queue, a replacement remains pending. After the first delivery commits and releases the queue, the successor can replay the terminal state instead of overlapping the original side effect.

The implementation deliberately:

  • Keeps Graphile's default replace mode. This preserves the durable successor created by delayed rescheduling and crash recovery.
  • Applies only to messages that carry idempotencyKey. Unkeyed orchestrator messages and public World interfaces are unchanged.
  • Uses 2,048 stable SHA-256 buckets per configured jobPrefix. Exact per-key queue names would grow Graphile's persistent queue table without bound; hashing the job task name keeps prefixes isolated and the physical name below Graphile's 128-character limit.
  • Applies the same mapping when migrating pg-boss jobs that actually carry MessageData.idempotencyKey. Legacy unkeyed rows used messageId as singleton_key and remain unqueued.
  • Does not run live GC_JOB_QUEUES; Graphile 0.16.6 cleanup raced concurrent producers in a local stress test and left 398 of 7,680 jobs orphaned (5.18%).

jobKeyMode: 'unsafe_dedupe' is not safe here. The handler enqueues a delayed successor before its current locked job returns; unsafe_dedupe would discard that successor, after which completing the current row can strand the workflow. It can also suppress the only replacement for a locked job that later dies on its final attempt.

This serializes competing deliveries; it does not claim exactly-once execution. The runtime's terminal-state replay remains the duplicate-suppression layer after the queue releases.

Operational trade-offs:

CaseBehavior
First mixed-version rollout or rollbackExisting/outgoing jobs without queueName do not own the new queue. Deployments must drain old keyed jobs and producers for a fully protected handoff.
Hard process deathA bucket can remain locked until Graphile's four-hour stale-lock recovery (or an explicit worker unlock). Graceful shutdown releases it normally.
Bucket collisionUnrelated keyed jobs can serialize. At the default concurrency of 50, 2,048 buckets produce about 0.60 expected colliding pairs, or roughly 1.2% expected slot loss.

A shuffled three-repetition Graphile 0.16.6 no-op benchmark at concurrency 50 / pool 8 measured median throughput of 181.5 jobs/s with no named queues, 141.1 jobs/s with all 2,048 buckets populated, and 108.2 jobs/s with 4,096. The VPS timings were noisy, so these are directional rather than a production capacity claim; 2,048 was consistently the better balance. Dropping to 1,024 would double expected collision loss at the default concurrency to roughly 2.4%.

This complements #3119 and #3162 but does not replace them. Those address selecting parked runs and accumulating startup-recovery jobs; this PR closes the separate locked-keyed-delivery concurrency gap. The queueName seam is independent of the HTTP loopback and can be carried through the in-process execution refactor in #3322.

How did you test your changes?

Added unit coverage for:

  • Keyed producer jobs receiving the deterministic named queue.
  • Delayed handler reschedules retaining the same jobKey, named queue, runAt, and 49-attempt budget.
  • Unkeyed messages remaining outside named queues.
  • Custom jobPrefix values receiving isolated queue scopes.
  • pg-boss migration distinguishing real idempotency keys from message-ID singleton keys.

Added a real-PostgreSQL Testcontainers regression with two independent pools and two createQueue() instances. It blocks worker A's HTTP delivery, enqueues the same key through worker B, waits through the polling window, and asserts that both rows share one named queue while only the first is locked. Before the fix, the same test observed two concurrent HTTP deliveries (maxActiveRequests: 2); with the fix it observes one (maxActiveRequests: 1) and then drains both jobs sequentially.

Verification:

  • pnpm exec vitest run packages/world-postgres/src (3 files, 30 tests)
  • Real-PostgreSQL concurrency regression against an isolated local PostgreSQL 17 database (1 test)
  • pnpm --filter @workflow/world-postgres typecheck
  • pnpm --filter @workflow/world-postgres build
  • Biome checks over the changed TypeScript files (only the two pre-existing warnings in queue.ts / queue.test.ts)
  • pnpm changeset status

Additional recovery probe with Graphile Worker 0.16.6: a child worker was killed while holding the named queue, a same-key replacement was added, and after forceUnlockWorkers only the replacement executed.

PR Checklist - Required to merge

  • 📦 pnpm changeset was run to create a changelog for this PR
    • Patch changeset for @workflow/world-postgres.
  • 🔒 DCO sign-off passes (run git commit --signoff on your commits)
  • 📝 Ping @vercel/workflow in a comment once the PR is ready, and the above checklist is complete

Draft for author review. No reviewers have been requested.

Assign idempotent jobs to bounded Graphile named queues so a locked delivery and its durable replacement cannot execute concurrently across worker processes. Keep default replacement semantics for delayed retries and crash recovery.
Signed-off-by: Joey Hotz <joeyhotz1@gmail.com>
@changeset-bot

Copy link
Copy Markdown

🦋 Changeset detected

Latest commit: 163dbf4

The changes in this PR will be included in the next version bump.

This PR includes changesets to release 1 package
NameType
@workflow/world-postgresPatch

Not sure what this means? Click here to learn what changesets are.

Click here if you're a maintainer who wants to add another changeset to this PR

@vercel

vercelBot commented Aug 19, 2026

Copy link
Copy Markdown
Contributor

@joeyhotz is attempting to deploy a commit to the Vercel Labs Team on Vercel.

A member of the Team first needs to authorize it.

Sign up for freeto 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

@joeyhotz
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content

fix(world-postgres): serialize keyed deliveries across workers - #3657

Closed
joeyhotz wants to merge 1 commit into
vercel:mainfrom
joeyhotz:fix/postgres-keyed-job-serialization
Closed

fix(world-postgres): serialize keyed deliveries across workers#3657
joeyhotz wants to merge 1 commit into
vercel:mainfrom
joeyhotz:fix/postgres-keyed-job-serialization

Conversation

@joeyhotz

Copy link
Copy Markdown
Contributor

Description

@workflow/world-postgres currently passes QueueOptions.idempotencyKey to Graphile Worker as jobKey. With Graphile's default replace mode, adding the same key while its current job is locked clears the locked job's key, exhausts that row, and inserts a successor. The existing inflightMessages map prevents overlap inside one process, but separate Postgres World processes can claim the current job and its successor at the same time.

We observed this after a self-hosted deployment scaled out: startup recovery replayed active runs and redispatched their keyed steps. In one staging burst, 7 of 34 materialisation steps received a second start without an intervening failure (41 executions total). The same startup-recovery/duplicate-start signature appeared in production. This behavior is present in @workflow/world-postgres 4.3.3 and current main.

This PR assigns every keyed Graphile job to a deterministic named queue. A named queue is Graphile's native cross-worker concurrency primitive: while one worker owns the queue, a replacement remains pending. After the first delivery commits and releases the queue, the successor can replay the terminal state instead of overlapping the original side effect.

The implementation deliberately:

  • Keeps Graphile's default replace mode. This preserves the durable successor created by delayed rescheduling and crash recovery.
  • Applies only to messages that carry idempotencyKey. Unkeyed orchestrator messages and public World interfaces are unchanged.
  • Uses 2,048 stable SHA-256 buckets per configured jobPrefix. Exact per-key queue names would grow Graphile's persistent queue table without bound; hashing the job task name keeps prefixes isolated and the physical name below Graphile's 128-character limit.
  • Applies the same mapping when migrating pg-boss jobs that actually carry MessageData.idempotencyKey. Legacy unkeyed rows used messageId as singleton_key and remain unqueued.
  • Does not run live GC_JOB_QUEUES; Graphile 0.16.6 cleanup raced concurrent producers in a local stress test and left 398 of 7,680 jobs orphaned (5.18%).

jobKeyMode: 'unsafe_dedupe' is not safe here. The handler enqueues a delayed successor before its current locked job returns; unsafe_dedupe would discard that successor, after which completing the current row can strand the workflow. It can also suppress the only replacement for a locked job that later dies on its final attempt.

This serializes competing deliveries; it does not claim exactly-once execution. The runtime's terminal-state replay remains the duplicate-suppression layer after the queue releases.

Operational trade-offs:

CaseBehavior
First mixed-version rollout or rollbackExisting/outgoing jobs without queueName do not own the new queue. Deployments must drain old keyed jobs and producers for a fully protected handoff.
Hard process deathA bucket can remain locked until Graphile's four-hour stale-lock recovery (or an explicit worker unlock). Graceful shutdown releases it normally.
Bucket collisionUnrelated keyed jobs can serialize. At the default concurrency of 50, 2,048 buckets produce about 0.60 expected colliding pairs, or roughly 1.2% expected slot loss.

A shuffled three-repetition Graphile 0.16.6 no-op benchmark at concurrency 50 / pool 8 measured median throughput of 181.5 jobs/s with no named queues, 141.1 jobs/s with all 2,048 buckets populated, and 108.2 jobs/s with 4,096. The VPS timings were noisy, so these are directional rather than a production capacity claim; 2,048 was consistently the better balance. Dropping to 1,024 would double expected collision loss at the default concurrency to roughly 2.4%.

This complements #3119 and #3162 but does not replace them. Those address selecting parked runs and accumulating startup-recovery jobs; this PR closes the separate locked-keyed-delivery concurrency gap. The queueName seam is independent of the HTTP loopback and can be carried through the in-process execution refactor in #3322.

How did you test your changes?

Added unit coverage for:

  • Keyed producer jobs receiving the deterministic named queue.
  • Delayed handler reschedules retaining the same jobKey, named queue, runAt, and 49-attempt budget.
  • Unkeyed messages remaining outside named queues.
  • Custom jobPrefix values receiving isolated queue scopes.
  • pg-boss migration distinguishing real idempotency keys from message-ID singleton keys.

Added a real-PostgreSQL Testcontainers regression with two independent pools and two createQueue() instances. It blocks worker A's HTTP delivery, enqueues the same key through worker B, waits through the polling window, and asserts that both rows share one named queue while only the first is locked. Before the fix, the same test observed two concurrent HTTP deliveries (maxActiveRequests: 2); with the fix it observes one (maxActiveRequests: 1) and then drains both jobs sequentially.

Verification:

  • pnpm exec vitest run packages/world-postgres/src (3 files, 30 tests)
  • Real-PostgreSQL concurrency regression against an isolated local PostgreSQL 17 database (1 test)
  • pnpm --filter @workflow/world-postgres typecheck
  • pnpm --filter @workflow/world-postgres build
  • Biome checks over the changed TypeScript files (only the two pre-existing warnings in queue.ts / queue.test.ts)
  • pnpm changeset status

Additional recovery probe with Graphile Worker 0.16.6: a child worker was killed while holding the named queue, a same-key replacement was added, and after forceUnlockWorkers only the replacement executed.

PR Checklist - Required to merge

  • 📦 pnpm changeset was run to create a changelog for this PR
    • Patch changeset for @workflow/world-postgres.
  • 🔒 DCO sign-off passes (run git commit --signoff on your commits)
  • 📝 Ping @vercel/workflow in a comment once the PR is ready, and the above checklist is complete

Draft for author review. No reviewers have been requested.

Assign idempotent jobs to bounded Graphile named queues so a locked delivery and its durable replacement cannot execute concurrently across worker processes. Keep default replacement semantics for delayed retries and crash recovery.
Signed-off-by: Joey Hotz <joeyhotz1@gmail.com>
@changeset-bot

Copy link
Copy Markdown

🦋 Changeset detected

Latest commit: 163dbf4

The changes in this PR will be included in the next version bump.

This PR includes changesets to release 1 package
NameType
@workflow/world-postgresPatch

Not sure what this means? Click here to learn what changesets are.

Click here if you're a maintainer who wants to add another changeset to this PR

@vercel

vercelBot commented Aug 19, 2026

Copy link
Copy Markdown
Contributor

@joeyhotz is attempting to deploy a commit to the Vercel Labs Team on Vercel.

A member of the Team first needs to authorize it.

Sign up for freeto 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

@joeyhotz
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

fix(world-postgres): serialize keyed deliveries across workers - #3657

Closed
joeyhotz wants to merge 1 commit into
vercel:mainfrom
joeyhotz:fix/postgres-keyed-job-serialization
Closed

fix(world-postgres): serialize keyed deliveries across workers#3657
joeyhotz wants to merge 1 commit into
vercel:mainfrom
joeyhotz:fix/postgres-keyed-job-serialization

Conversation

@joeyhotz

Copy link
Copy Markdown
Contributor

Description

@workflow/world-postgres currently passes QueueOptions.idempotencyKey to Graphile Worker as jobKey. With Graphile's default replace mode, adding the same key while its current job is locked clears the locked job's key, exhausts that row, and inserts a successor. The existing inflightMessages map prevents overlap inside one process, but separate Postgres World processes can claim the current job and its successor at the same time.

We observed this after a self-hosted deployment scaled out: startup recovery replayed active runs and redispatched their keyed steps. In one staging burst, 7 of 34 materialisation steps received a second start without an intervening failure (41 executions total). The same startup-recovery/duplicate-start signature appeared in production. This behavior is present in @workflow/world-postgres 4.3.3 and current main.

This PR assigns every keyed Graphile job to a deterministic named queue. A named queue is Graphile's native cross-worker concurrency primitive: while one worker owns the queue, a replacement remains pending. After the first delivery commits and releases the queue, the successor can replay the terminal state instead of overlapping the original side effect.

The implementation deliberately:

  • Keeps Graphile's default replace mode. This preserves the durable successor created by delayed rescheduling and crash recovery.
  • Applies only to messages that carry idempotencyKey. Unkeyed orchestrator messages and public World interfaces are unchanged.
  • Uses 2,048 stable SHA-256 buckets per configured jobPrefix. Exact per-key queue names would grow Graphile's persistent queue table without bound; hashing the job task name keeps prefixes isolated and the physical name below Graphile's 128-character limit.
  • Applies the same mapping when migrating pg-boss jobs that actually carry MessageData.idempotencyKey. Legacy unkeyed rows used messageId as singleton_key and remain unqueued.
  • Does not run live GC_JOB_QUEUES; Graphile 0.16.6 cleanup raced concurrent producers in a local stress test and left 398 of 7,680 jobs orphaned (5.18%).

jobKeyMode: 'unsafe_dedupe' is not safe here. The handler enqueues a delayed successor before its current locked job returns; unsafe_dedupe would discard that successor, after which completing the current row can strand the workflow. It can also suppress the only replacement for a locked job that later dies on its final attempt.

This serializes competing deliveries; it does not claim exactly-once execution. The runtime's terminal-state replay remains the duplicate-suppression layer after the queue releases.

Operational trade-offs:

CaseBehavior
First mixed-version rollout or rollbackExisting/outgoing jobs without queueName do not own the new queue. Deployments must drain old keyed jobs and producers for a fully protected handoff.
Hard process deathA bucket can remain locked until Graphile's four-hour stale-lock recovery (or an explicit worker unlock). Graceful shutdown releases it normally.
Bucket collisionUnrelated keyed jobs can serialize. At the default concurrency of 50, 2,048 buckets produce about 0.60 expected colliding pairs, or roughly 1.2% expected slot loss.

A shuffled three-repetition Graphile 0.16.6 no-op benchmark at concurrency 50 / pool 8 measured median throughput of 181.5 jobs/s with no named queues, 141.1 jobs/s with all 2,048 buckets populated, and 108.2 jobs/s with 4,096. The VPS timings were noisy, so these are directional rather than a production capacity claim; 2,048 was consistently the better balance. Dropping to 1,024 would double expected collision loss at the default concurrency to roughly 2.4%.

This complements #3119 and #3162 but does not replace them. Those address selecting parked runs and accumulating startup-recovery jobs; this PR closes the separate locked-keyed-delivery concurrency gap. The queueName seam is independent of the HTTP loopback and can be carried through the in-process execution refactor in #3322.

How did you test your changes?

Added unit coverage for:

  • Keyed producer jobs receiving the deterministic named queue.
  • Delayed handler reschedules retaining the same jobKey, named queue, runAt, and 49-attempt budget.
  • Unkeyed messages remaining outside named queues.
  • Custom jobPrefix values receiving isolated queue scopes.
  • pg-boss migration distinguishing real idempotency keys from message-ID singleton keys.

Added a real-PostgreSQL Testcontainers regression with two independent pools and two createQueue() instances. It blocks worker A's HTTP delivery, enqueues the same key through worker B, waits through the polling window, and asserts that both rows share one named queue while only the first is locked. Before the fix, the same test observed two concurrent HTTP deliveries (maxActiveRequests: 2); with the fix it observes one (maxActiveRequests: 1) and then drains both jobs sequentially.

Verification:

  • pnpm exec vitest run packages/world-postgres/src (3 files, 30 tests)
  • Real-PostgreSQL concurrency regression against an isolated local PostgreSQL 17 database (1 test)
  • pnpm --filter @workflow/world-postgres typecheck
  • pnpm --filter @workflow/world-postgres build
  • Biome checks over the changed TypeScript files (only the two pre-existing warnings in queue.ts / queue.test.ts)
  • pnpm changeset status

Additional recovery probe with Graphile Worker 0.16.6: a child worker was killed while holding the named queue, a same-key replacement was added, and after forceUnlockWorkers only the replacement executed.

PR Checklist - Required to merge

  • 📦 pnpm changeset was run to create a changelog for this PR
    • Patch changeset for @workflow/world-postgres.
  • 🔒 DCO sign-off passes (run git commit --signoff on your commits)
  • 📝 Ping @vercel/workflow in a comment once the PR is ready, and the above checklist is complete

Draft for author review. No reviewers have been requested.

Assign idempotent jobs to bounded Graphile named queues so a locked delivery and its durable replacement cannot execute concurrently across worker processes. Keep default replacement semantics for delayed retries and crash recovery.
Signed-off-by: Joey Hotz <joeyhotz1@gmail.com>
@changeset-bot

Copy link
Copy Markdown

🦋 Changeset detected

Latest commit: 163dbf4

The changes in this PR will be included in the next version bump.

This PR includes changesets to release 1 package
NameType
@workflow/world-postgresPatch

Not sure what this means? Click here to learn what changesets are.

Click here if you're a maintainer who wants to add another changeset to this PR

@vercel

vercelBot commented Aug 19, 2026

Copy link
Copy Markdown
Contributor

@joeyhotz is attempting to deploy a commit to the Vercel Labs Team on Vercel.

A member of the Team first needs to authorize it.

Sign up for freeto 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

@joeyhotz
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

fix(world-postgres): serialize keyed deliveries across workers - #3657

Closed
joeyhotz wants to merge 1 commit into
vercel:mainfrom
joeyhotz:fix/postgres-keyed-job-serialization
Closed

fix(world-postgres): serialize keyed deliveries across workers#3657
joeyhotz wants to merge 1 commit into
vercel:mainfrom
joeyhotz:fix/postgres-keyed-job-serialization

Conversation

@joeyhotz

Copy link
Copy Markdown
Contributor

Description

@workflow/world-postgres currently passes QueueOptions.idempotencyKey to Graphile Worker as jobKey. With Graphile's default replace mode, adding the same key while its current job is locked clears the locked job's key, exhausts that row, and inserts a successor. The existing inflightMessages map prevents overlap inside one process, but separate Postgres World processes can claim the current job and its successor at the same time.

We observed this after a self-hosted deployment scaled out: startup recovery replayed active runs and redispatched their keyed steps. In one staging burst, 7 of 34 materialisation steps received a second start without an intervening failure (41 executions total). The same startup-recovery/duplicate-start signature appeared in production. This behavior is present in @workflow/world-postgres 4.3.3 and current main.

This PR assigns every keyed Graphile job to a deterministic named queue. A named queue is Graphile's native cross-worker concurrency primitive: while one worker owns the queue, a replacement remains pending. After the first delivery commits and releases the queue, the successor can replay the terminal state instead of overlapping the original side effect.

The implementation deliberately:

  • Keeps Graphile's default replace mode. This preserves the durable successor created by delayed rescheduling and crash recovery.
  • Applies only to messages that carry idempotencyKey. Unkeyed orchestrator messages and public World interfaces are unchanged.
  • Uses 2,048 stable SHA-256 buckets per configured jobPrefix. Exact per-key queue names would grow Graphile's persistent queue table without bound; hashing the job task name keeps prefixes isolated and the physical name below Graphile's 128-character limit.
  • Applies the same mapping when migrating pg-boss jobs that actually carry MessageData.idempotencyKey. Legacy unkeyed rows used messageId as singleton_key and remain unqueued.
  • Does not run live GC_JOB_QUEUES; Graphile 0.16.6 cleanup raced concurrent producers in a local stress test and left 398 of 7,680 jobs orphaned (5.18%).

jobKeyMode: 'unsafe_dedupe' is not safe here. The handler enqueues a delayed successor before its current locked job returns; unsafe_dedupe would discard that successor, after which completing the current row can strand the workflow. It can also suppress the only replacement for a locked job that later dies on its final attempt.

This serializes competing deliveries; it does not claim exactly-once execution. The runtime's terminal-state replay remains the duplicate-suppression layer after the queue releases.

Operational trade-offs:

CaseBehavior
First mixed-version rollout or rollbackExisting/outgoing jobs without queueName do not own the new queue. Deployments must drain old keyed jobs and producers for a fully protected handoff.
Hard process deathA bucket can remain locked until Graphile's four-hour stale-lock recovery (or an explicit worker unlock). Graceful shutdown releases it normally.
Bucket collisionUnrelated keyed jobs can serialize. At the default concurrency of 50, 2,048 buckets produce about 0.60 expected colliding pairs, or roughly 1.2% expected slot loss.

A shuffled three-repetition Graphile 0.16.6 no-op benchmark at concurrency 50 / pool 8 measured median throughput of 181.5 jobs/s with no named queues, 141.1 jobs/s with all 2,048 buckets populated, and 108.2 jobs/s with 4,096. The VPS timings were noisy, so these are directional rather than a production capacity claim; 2,048 was consistently the better balance. Dropping to 1,024 would double expected collision loss at the default concurrency to roughly 2.4%.

This complements #3119 and #3162 but does not replace them. Those address selecting parked runs and accumulating startup-recovery jobs; this PR closes the separate locked-keyed-delivery concurrency gap. The queueName seam is independent of the HTTP loopback and can be carried through the in-process execution refactor in #3322.

How did you test your changes?

Added unit coverage for:

  • Keyed producer jobs receiving the deterministic named queue.
  • Delayed handler reschedules retaining the same jobKey, named queue, runAt, and 49-attempt budget.
  • Unkeyed messages remaining outside named queues.
  • Custom jobPrefix values receiving isolated queue scopes.
  • pg-boss migration distinguishing real idempotency keys from message-ID singleton keys.

Added a real-PostgreSQL Testcontainers regression with two independent pools and two createQueue() instances. It blocks worker A's HTTP delivery, enqueues the same key through worker B, waits through the polling window, and asserts that both rows share one named queue while only the first is locked. Before the fix, the same test observed two concurrent HTTP deliveries (maxActiveRequests: 2); with the fix it observes one (maxActiveRequests: 1) and then drains both jobs sequentially.

Verification:

  • pnpm exec vitest run packages/world-postgres/src (3 files, 30 tests)
  • Real-PostgreSQL concurrency regression against an isolated local PostgreSQL 17 database (1 test)
  • pnpm --filter @workflow/world-postgres typecheck
  • pnpm --filter @workflow/world-postgres build
  • Biome checks over the changed TypeScript files (only the two pre-existing warnings in queue.ts / queue.test.ts)
  • pnpm changeset status

Additional recovery probe with Graphile Worker 0.16.6: a child worker was killed while holding the named queue, a same-key replacement was added, and after forceUnlockWorkers only the replacement executed.

PR Checklist - Required to merge

  • 📦 pnpm changeset was run to create a changelog for this PR
    • Patch changeset for @workflow/world-postgres.
  • 🔒 DCO sign-off passes (run git commit --signoff on your commits)
  • 📝 Ping @vercel/workflow in a comment once the PR is ready, and the above checklist is complete

Draft for author review. No reviewers have been requested.

Assign idempotent jobs to bounded Graphile named queues so a locked delivery and its durable replacement cannot execute concurrently across worker processes. Keep default replacement semantics for delayed retries and crash recovery.
Signed-off-by: Joey Hotz <joeyhotz1@gmail.com>
@changeset-bot

Copy link
Copy Markdown

🦋 Changeset detected

Latest commit: 163dbf4

The changes in this PR will be included in the next version bump.

This PR includes changesets to release 1 package
NameType
@workflow/world-postgresPatch

Not sure what this means? Click here to learn what changesets are.

Click here if you're a maintainer who wants to add another changeset to this PR

@vercel

vercelBot commented Aug 19, 2026

Copy link
Copy Markdown
Contributor

@joeyhotz is attempting to deploy a commit to the Vercel Labs Team on Vercel.

A member of the Team first needs to authorize it.

Sign up for freeto 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

@joeyhotz
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content

fix(world-postgres): serialize keyed deliveries across workers - #3657

Closed
joeyhotz wants to merge 1 commit into
vercel:mainfrom
joeyhotz:fix/postgres-keyed-job-serialization
Closed

fix(world-postgres): serialize keyed deliveries across workers#3657
joeyhotz wants to merge 1 commit into
vercel:mainfrom
joeyhotz:fix/postgres-keyed-job-serialization

Conversation

@joeyhotz

Copy link
Copy Markdown
Contributor

Description

@workflow/world-postgres currently passes QueueOptions.idempotencyKey to Graphile Worker as jobKey. With Graphile's default replace mode, adding the same key while its current job is locked clears the locked job's key, exhausts that row, and inserts a successor. The existing inflightMessages map prevents overlap inside one process, but separate Postgres World processes can claim the current job and its successor at the same time.

We observed this after a self-hosted deployment scaled out: startup recovery replayed active runs and redispatched their keyed steps. In one staging burst, 7 of 34 materialisation steps received a second start without an intervening failure (41 executions total). The same startup-recovery/duplicate-start signature appeared in production. This behavior is present in @workflow/world-postgres 4.3.3 and current main.

This PR assigns every keyed Graphile job to a deterministic named queue. A named queue is Graphile's native cross-worker concurrency primitive: while one worker owns the queue, a replacement remains pending. After the first delivery commits and releases the queue, the successor can replay the terminal state instead of overlapping the original side effect.

The implementation deliberately:

  • Keeps Graphile's default replace mode. This preserves the durable successor created by delayed rescheduling and crash recovery.
  • Applies only to messages that carry idempotencyKey. Unkeyed orchestrator messages and public World interfaces are unchanged.
  • Uses 2,048 stable SHA-256 buckets per configured jobPrefix. Exact per-key queue names would grow Graphile's persistent queue table without bound; hashing the job task name keeps prefixes isolated and the physical name below Graphile's 128-character limit.
  • Applies the same mapping when migrating pg-boss jobs that actually carry MessageData.idempotencyKey. Legacy unkeyed rows used messageId as singleton_key and remain unqueued.
  • Does not run live GC_JOB_QUEUES; Graphile 0.16.6 cleanup raced concurrent producers in a local stress test and left 398 of 7,680 jobs orphaned (5.18%).

jobKeyMode: 'unsafe_dedupe' is not safe here. The handler enqueues a delayed successor before its current locked job returns; unsafe_dedupe would discard that successor, after which completing the current row can strand the workflow. It can also suppress the only replacement for a locked job that later dies on its final attempt.

This serializes competing deliveries; it does not claim exactly-once execution. The runtime's terminal-state replay remains the duplicate-suppression layer after the queue releases.

Operational trade-offs:

CaseBehavior
First mixed-version rollout or rollbackExisting/outgoing jobs without queueName do not own the new queue. Deployments must drain old keyed jobs and producers for a fully protected handoff.
Hard process deathA bucket can remain locked until Graphile's four-hour stale-lock recovery (or an explicit worker unlock). Graceful shutdown releases it normally.
Bucket collisionUnrelated keyed jobs can serialize. At the default concurrency of 50, 2,048 buckets produce about 0.60 expected colliding pairs, or roughly 1.2% expected slot loss.

A shuffled three-repetition Graphile 0.16.6 no-op benchmark at concurrency 50 / pool 8 measured median throughput of 181.5 jobs/s with no named queues, 141.1 jobs/s with all 2,048 buckets populated, and 108.2 jobs/s with 4,096. The VPS timings were noisy, so these are directional rather than a production capacity claim; 2,048 was consistently the better balance. Dropping to 1,024 would double expected collision loss at the default concurrency to roughly 2.4%.

This complements #3119 and #3162 but does not replace them. Those address selecting parked runs and accumulating startup-recovery jobs; this PR closes the separate locked-keyed-delivery concurrency gap. The queueName seam is independent of the HTTP loopback and can be carried through the in-process execution refactor in #3322.

How did you test your changes?

Added unit coverage for:

  • Keyed producer jobs receiving the deterministic named queue.
  • Delayed handler reschedules retaining the same jobKey, named queue, runAt, and 49-attempt budget.
  • Unkeyed messages remaining outside named queues.
  • Custom jobPrefix values receiving isolated queue scopes.
  • pg-boss migration distinguishing real idempotency keys from message-ID singleton keys.

Added a real-PostgreSQL Testcontainers regression with two independent pools and two createQueue() instances. It blocks worker A's HTTP delivery, enqueues the same key through worker B, waits through the polling window, and asserts that both rows share one named queue while only the first is locked. Before the fix, the same test observed two concurrent HTTP deliveries (maxActiveRequests: 2); with the fix it observes one (maxActiveRequests: 1) and then drains both jobs sequentially.

Verification:

  • pnpm exec vitest run packages/world-postgres/src (3 files, 30 tests)
  • Real-PostgreSQL concurrency regression against an isolated local PostgreSQL 17 database (1 test)
  • pnpm --filter @workflow/world-postgres typecheck
  • pnpm --filter @workflow/world-postgres build
  • Biome checks over the changed TypeScript files (only the two pre-existing warnings in queue.ts / queue.test.ts)
  • pnpm changeset status

Additional recovery probe with Graphile Worker 0.16.6: a child worker was killed while holding the named queue, a same-key replacement was added, and after forceUnlockWorkers only the replacement executed.

PR Checklist - Required to merge

  • 📦 pnpm changeset was run to create a changelog for this PR
    • Patch changeset for @workflow/world-postgres.
  • 🔒 DCO sign-off passes (run git commit --signoff on your commits)
  • 📝 Ping @vercel/workflow in a comment once the PR is ready, and the above checklist is complete

Draft for author review. No reviewers have been requested.

Assign idempotent jobs to bounded Graphile named queues so a locked delivery and its durable replacement cannot execute concurrently across worker processes. Keep default replacement semantics for delayed retries and crash recovery.
Signed-off-by: Joey Hotz <joeyhotz1@gmail.com>
@changeset-bot

Copy link
Copy Markdown

🦋 Changeset detected

Latest commit: 163dbf4

The changes in this PR will be included in the next version bump.

This PR includes changesets to release 1 package
NameType
@workflow/world-postgresPatch

Not sure what this means? Click here to learn what changesets are.

Click here if you're a maintainer who wants to add another changeset to this PR

@vercel

vercelBot commented Aug 19, 2026

Copy link
Copy Markdown
Contributor

@joeyhotz is attempting to deploy a commit to the Vercel Labs Team on Vercel.

A member of the Team first needs to authorize it.

Sign up for freeto 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

@joeyhotz