Skip to content
Closed
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/late-hooks-rearm.md
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
---
"@workflow/core": patch
---

Keep late-claimed buffered hook payloads from being preempted by workflow suspension in retained VMs.
46 changes: 45 additions & 1 deletion packages/core/src/async-deserialization-ordering.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -4,7 +4,7 @@ import * as nanoid from 'nanoid';
import { monotonicFactory } from 'ulid';
import { afterEach, describe, expect, it, vi } from 'vitest';
import { EventsConsumer } from './events-consumer.js';
import type { WorkflowOrchestratorContext } from './private.js';
import { isDeliveryIdle, type WorkflowOrchestratorContext } from './private.js';
import { dehydrateStepReturnValue } from './serialization.js';
import { createUseStep } from './step.js';
import { createContext } from './vm/index.js';
Expand DownExpand Up@@ -685,4 +685,48 @@ describe('async deserialization ordering', () => {

expect(ctx.pendingDeliveryBarriers?.size).toBe(0);
});

it('should restore delivery protection when a buffered hook payload is claimed after its barrier retires', async () => {
const payload = await dehydrateStepReturnValue(
'buffered',
'wrun_test',
undefined
);
const ctx = setupWorkflowContext([
{
eventId: 'evnt_0',
runId: 'wrun_test',
eventType: 'hook_received',
correlationId: 'hook_01K11TFZ62YS0YYFDQ3E8B9YCV',
eventData: { payload },
createdAt: new Date(),
},
]);

// Model another payload hydrating in the same replay so the idle safety
// net cannot retire this hook's barrier before the test observes it.
ctx.pendingDeliveries = 1;
const createHook = createCreateHook(ctx);
const hook = createHook<string>();

// The payload arrived before a consumer requested it, so it starts with
// an unarmed barrier and remains buffered after the idle safety net retires
// that barrier.
await vi.waitFor(() => {
expect(ctx.pendingDeliveryBarriers?.size).toBe(1);
});
await ctx.promiseQueue;
expect(ctx.pendingDeliveryBarriers?.size).toBe(1);
ctx.pendingDeliveries = 0;
await vi.waitFor(() => {
expect(ctx.pendingDeliveryBarriers?.size).toBe(0);
});

// Claiming the buffered payload commits it to reaching this consumer. Its
// delivery must become non-idle again until the claim resolves.
const delivery = hook.then((value) => value);
expect(isDeliveryIdle(ctx)).toBe(false);
await expect(delivery).resolves.toBe('buffered');
expect(isDeliveryIdle(ctx)).toBe(true);
});
});
12 changes: 12 additions & 0 deletions packages/core/src/delivery-barrier-dispenser.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -117,6 +117,18 @@ const resumeAtA = new Date('2026-07-27T12:00:05.000Z');
const resumeAtB = new Date('2026-07-27T12:00:06.000Z');

describe('barrier safety-net dispenser', () => {
it('rejects a second barrier owner for the same event index', () => {
const ctx = setupWorkflowContext([]);
const barrier = registerDeliveryBarrier(ctx, 0, 'hook', { armed: false });

expect(() => registerDeliveryBarrier(ctx, 0, 'step')).toThrowError(
'Delivery barrier already registered at event index 0'
);

barrier.markDelivered();
expect(isDeliveryIdle(ctx)).toBe(true);
});

it('suspends only after deliveries parked behind an unclaimed payload have run', async () => {
const ops: Promise<any>[] = [];
const payload = await dehydrateStepReturnValue(
Expand Down
109 changes: 72 additions & 37 deletions packages/core/src/private.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -2,6 +2,7 @@
* Utils used by the bundler when transforming code
*/

import { WorkflowRuntimeError } from '@workflow/errors';
import { withResolvers } from '@workflow/utils';
import type { CryptoKey } from './encryption.js';
import type { EventsConsumer } from './events-consumer.js';
Expand DownExpand Up@@ -215,8 +216,8 @@ export type DeliveryKind = 'hook' | 'wait' | 'step';

interface DeliveryBarrierEntry {
kind: DeliveryKind;
/** Resolves once this delivery has resolved to the workflow. */
delivered: Promise<void>;
/** Resolves once this delivery is handed to the workflow or retired. */
released: Promise<void>;
/**
* Whether this delivery is committed to reaching the workflow without any
* further action by workflow code. True for wait completions and step
Expand All@@ -229,11 +230,15 @@ interface DeliveryBarrierEntry {
* once a consumer takes the payload.
*/
armed: boolean;
/** Whether this entry has been removed and its `released` promise settled. */
retired: boolean;
/**
* Retire this entry: resolve `delivered` and remove it from the registry,
* exactly as `markDelivered` would. Called only by the context's safety-net
* dispenser ({@link ensureBarrierSafetyNet}), and only on the lowest-index
* entry at delivery idle. Idempotent.
* Retire this entry: resolve `released` and remove it from the registry,
* without marking the handle delivered to the workflow. A safety-retired
* buffered payload may therefore install a fresh entry if it is claimed by
* a retained VM later. Called only by the context's safety-net dispenser
* ({@link ensureBarrierSafetyNet}), and only on the lowest-index entry at
* delivery idle. Idempotent.
*/
retire: () => void;
}
Expand DownExpand Up@@ -396,7 +401,7 @@ export async function awaitEarlierDeliveries(
) {
continue;
}
earlier.push(entry.delivered);
earlier.push(entry.released);
}
if (earlier.length > 0) {
await Promise.all(earlier);
Expand All@@ -422,7 +427,7 @@ export async function awaitEarlierDeliveries(
export interface DeliveryBarrier {
/**
* Mark this delivery as delivered to the workflow. Resolves its
* `delivered` promise so any later-in-log delivery gated on it (via
* `released` promise so any later-in-log delivery gated on it (via
* {@link awaitEarlierDeliveries}) may proceed, and removes it from the
* registry. Idempotent.
*/
Expand DownExpand Up@@ -468,41 +473,71 @@ export function registerDeliveryBarrier(
return { markDelivered: () => {}, arm: () => {} };
}

let done = false;
const { promise, resolve } = withResolvers<void>();

const finish = () => {
if (done) {
return;
}
done = true;
if (barriers.get(eventIndex) === entry) {
barriers.delete(eventIndex);
const install = (armed: boolean): DeliveryBarrierEntry => {
if (barriers.has(eventIndex)) {
throw new WorkflowRuntimeError(
`Delivery barrier already registered at event index ${eventIndex}`
);
}
resolve();
const { promise, resolve } = withResolvers<void>();
const entry: DeliveryBarrierEntry = {
kind,
released: promise,
armed,
retired: false,
retire: () => {
if (entry.retired) {
return;
}
entry.retired = true;
if (barriers.get(eventIndex) === entry) {
barriers.delete(eventIndex);
}
resolve();
},
};
barriers.set(eventIndex, entry);

// Safety net: if this delivery is never delivered to the workflow (its
// branch was not taken / the run is suspending, or a buffered hook payload
// is only claimed after a later delivery the workflow is still waiting
// on), it is retired at idle so a later delivery gated on it cannot
// deadlock and the registry cannot leak an entry per abandoned delivery.
// Retirement goes through the context's single ordered dispenser rather
// than a per-barrier idle poll. See {@link ensureBarrierSafetyNet} for why
// the ORDER of these retirements is load-bearing.
ensureBarrierSafetyNet(ctx);
return entry;
};

const entry: DeliveryBarrierEntry = {
kind,
delivered: promise,
armed: options.armed ?? true,
retire: finish,
};
barriers.set(eventIndex, entry);

// Safety net: if this delivery is never delivered to the workflow (its
// branch was not taken / the run is suspending, or a buffered hook payload
// is only claimed after a later delivery the workflow is still waiting on),
// it is retired at idle so a later delivery gated on it cannot deadlock and
// the registry cannot leak an entry per abandoned delivery. Retirement goes
// through the context's single ordered dispenser rather than a per-barrier
// idle poll — see {@link ensureBarrierSafetyNet} for why the ORDER of these
// retirements is load-bearing.
ensureBarrierSafetyNet(ctx);
let entry = install(options.armed ?? true);
let deliveredToWorkflow = false;

return {
markDelivered: finish,
markDelivered: () => {
if (deliveredToWorkflow) {
return;
}
deliveredToWorkflow = true;
entry.retire();
},
arm: () => {
if (deliveredToWorkflow) {
return;
}
// The idle safety net may retire an unclaimed buffered hook payload
// while a retained VM keeps its `claim()` closure alive. If workflow
// code later claims that payload, replace the settled entry so delivery
// remains non-idle until the claim reaches the workflow.
if (entry.retired) {
entry = install(true);
return;
}
if (barriers.get(eventIndex) !== entry) {
throw new WorkflowRuntimeError(
`Delivery barrier lost ownership of event index ${eventIndex}`
);
}
entry.armed = true;
},
};
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Add copy buttons to all
 blocks
(function() {
function addCopyButtons() {
document.querySelectorAll('pre code').forEach(function(codeBlock) {
if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;
codeBlock.parentElement.setAttribute('data-copy-added', 'true');
var btn = document.createElement('button');
btn.textContent = 'Copy';
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;';
btn.onmouseover = function() { this.style.opacity = '1'; };
btn.onmouseout = function() { this.style.opacity = '0.7'; };
btn.onclick = function() {
navigator.clipboard.writeText(codeBlock.textContent).then(function() {
btn.textContent = 'Copied!';
setTimeout(function() { btn.textContent = 'Copy'; }, 1500);
});
};
codeBlock.parentElement.style.position = 'relative';
codeBlock.parentElement.appendChild(btn);
});
}
addCopyButtons();
// Re-run on dynamic content
var observer = new MutationObserver(addCopyButtons);
observer.observe(document.body, { childList: true, subtree: true });
})();
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Backport #3879: fix(core): re-arm late-claimed hook deliveries by github-actions[bot] · Pull Request #3905 · vercel/workflow · GitHub
Skip to content
Closed
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/late-hooks-rearm.md
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
---
"@workflow/core": patch
---

Keep late-claimed buffered hook payloads from being preempted by workflow suspension in retained VMs.
46 changes: 45 additions & 1 deletion packages/core/src/async-deserialization-ordering.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -4,7 +4,7 @@ import * as nanoid from 'nanoid';
import { monotonicFactory } from 'ulid';
import { afterEach, describe, expect, it, vi } from 'vitest';
import { EventsConsumer } from './events-consumer.js';
import type { WorkflowOrchestratorContext } from './private.js';
import { isDeliveryIdle, type WorkflowOrchestratorContext } from './private.js';
import { dehydrateStepReturnValue } from './serialization.js';
import { createUseStep } from './step.js';
import { createContext } from './vm/index.js';
Expand DownExpand Up@@ -685,4 +685,48 @@ describe('async deserialization ordering', () => {

expect(ctx.pendingDeliveryBarriers?.size).toBe(0);
});

it('should restore delivery protection when a buffered hook payload is claimed after its barrier retires', async () => {
const payload = await dehydrateStepReturnValue(
'buffered',
'wrun_test',
undefined
);
const ctx = setupWorkflowContext([
{
eventId: 'evnt_0',
runId: 'wrun_test',
eventType: 'hook_received',
correlationId: 'hook_01K11TFZ62YS0YYFDQ3E8B9YCV',
eventData: { payload },
createdAt: new Date(),
},
]);

// Model another payload hydrating in the same replay so the idle safety
// net cannot retire this hook's barrier before the test observes it.
ctx.pendingDeliveries = 1;
const createHook = createCreateHook(ctx);
const hook = createHook<string>();

// The payload arrived before a consumer requested it, so it starts with
// an unarmed barrier and remains buffered after the idle safety net retires
// that barrier.
await vi.waitFor(() => {
expect(ctx.pendingDeliveryBarriers?.size).toBe(1);
});
await ctx.promiseQueue;
expect(ctx.pendingDeliveryBarriers?.size).toBe(1);
ctx.pendingDeliveries = 0;
await vi.waitFor(() => {
expect(ctx.pendingDeliveryBarriers?.size).toBe(0);
});

// Claiming the buffered payload commits it to reaching this consumer. Its
// delivery must become non-idle again until the claim resolves.
const delivery = hook.then((value) => value);
expect(isDeliveryIdle(ctx)).toBe(false);
await expect(delivery).resolves.toBe('buffered');
expect(isDeliveryIdle(ctx)).toBe(true);
});
});
12 changes: 12 additions & 0 deletions packages/core/src/delivery-barrier-dispenser.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -117,6 +117,18 @@ const resumeAtA = new Date('2026-07-27T12:00:05.000Z');
const resumeAtB = new Date('2026-07-27T12:00:06.000Z');

describe('barrier safety-net dispenser', () => {
it('rejects a second barrier owner for the same event index', () => {
const ctx = setupWorkflowContext([]);
const barrier = registerDeliveryBarrier(ctx, 0, 'hook', { armed: false });

expect(() => registerDeliveryBarrier(ctx, 0, 'step')).toThrowError(
'Delivery barrier already registered at event index 0'
);

barrier.markDelivered();
expect(isDeliveryIdle(ctx)).toBe(true);
});

it('suspends only after deliveries parked behind an unclaimed payload have run', async () => {
const ops: Promise<any>[] = [];
const payload = await dehydrateStepReturnValue(
Expand Down
109 changes: 72 additions & 37 deletions packages/core/src/private.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -2,6 +2,7 @@
* Utils used by the bundler when transforming code
*/

import { WorkflowRuntimeError } from '@workflow/errors';
import { withResolvers } from '@workflow/utils';
import type { CryptoKey } from './encryption.js';
import type { EventsConsumer } from './events-consumer.js';
Expand DownExpand Up@@ -215,8 +216,8 @@ export type DeliveryKind = 'hook' | 'wait' | 'step';

interface DeliveryBarrierEntry {
kind: DeliveryKind;
/** Resolves once this delivery has resolved to the workflow. */
delivered: Promise<void>;
/** Resolves once this delivery is handed to the workflow or retired. */
released: Promise<void>;
/**
* Whether this delivery is committed to reaching the workflow without any
* further action by workflow code. True for wait completions and step
Expand All@@ -229,11 +230,15 @@ interface DeliveryBarrierEntry {
* once a consumer takes the payload.
*/
armed: boolean;
/** Whether this entry has been removed and its `released` promise settled. */
retired: boolean;
/**
* Retire this entry: resolve `delivered` and remove it from the registry,
* exactly as `markDelivered` would. Called only by the context's safety-net
* dispenser ({@link ensureBarrierSafetyNet}), and only on the lowest-index
* entry at delivery idle. Idempotent.
* Retire this entry: resolve `released` and remove it from the registry,
* without marking the handle delivered to the workflow. A safety-retired
* buffered payload may therefore install a fresh entry if it is claimed by
* a retained VM later. Called only by the context's safety-net dispenser
* ({@link ensureBarrierSafetyNet}), and only on the lowest-index entry at
* delivery idle. Idempotent.
*/
retire: () => void;
}
Expand DownExpand Up@@ -396,7 +401,7 @@ export async function awaitEarlierDeliveries(
) {
continue;
}
earlier.push(entry.delivered);
earlier.push(entry.released);
}
if (earlier.length > 0) {
await Promise.all(earlier);
Expand All@@ -422,7 +427,7 @@ export async function awaitEarlierDeliveries(
export interface DeliveryBarrier {
/**
* Mark this delivery as delivered to the workflow. Resolves its
* `delivered` promise so any later-in-log delivery gated on it (via
* `released` promise so any later-in-log delivery gated on it (via
* {@link awaitEarlierDeliveries}) may proceed, and removes it from the
* registry. Idempotent.
*/
Expand DownExpand Up@@ -468,41 +473,71 @@ export function registerDeliveryBarrier(
return { markDelivered: () => {}, arm: () => {} };
}

let done = false;
const { promise, resolve } = withResolvers<void>();

const finish = () => {
if (done) {
return;
}
done = true;
if (barriers.get(eventIndex) === entry) {
barriers.delete(eventIndex);
const install = (armed: boolean): DeliveryBarrierEntry => {
if (barriers.has(eventIndex)) {
throw new WorkflowRuntimeError(
`Delivery barrier already registered at event index ${eventIndex}`
);
}
resolve();
const { promise, resolve } = withResolvers<void>();
const entry: DeliveryBarrierEntry = {
kind,
released: promise,
armed,
retired: false,
retire: () => {
if (entry.retired) {
return;
}
entry.retired = true;
if (barriers.get(eventIndex) === entry) {
barriers.delete(eventIndex);
}
resolve();
},
};
barriers.set(eventIndex, entry);

// Safety net: if this delivery is never delivered to the workflow (its
// branch was not taken / the run is suspending, or a buffered hook payload
// is only claimed after a later delivery the workflow is still waiting
// on), it is retired at idle so a later delivery gated on it cannot
// deadlock and the registry cannot leak an entry per abandoned delivery.
// Retirement goes through the context's single ordered dispenser rather
// than a per-barrier idle poll. See {@link ensureBarrierSafetyNet} for why
// the ORDER of these retirements is load-bearing.
ensureBarrierSafetyNet(ctx);
return entry;
};

const entry: DeliveryBarrierEntry = {
kind,
delivered: promise,
armed: options.armed ?? true,
retire: finish,
};
barriers.set(eventIndex, entry);

// Safety net: if this delivery is never delivered to the workflow (its
// branch was not taken / the run is suspending, or a buffered hook payload
// is only claimed after a later delivery the workflow is still waiting on),
// it is retired at idle so a later delivery gated on it cannot deadlock and
// the registry cannot leak an entry per abandoned delivery. Retirement goes
// through the context's single ordered dispenser rather than a per-barrier
// idle poll — see {@link ensureBarrierSafetyNet} for why the ORDER of these
// retirements is load-bearing.
ensureBarrierSafetyNet(ctx);
let entry = install(options.armed ?? true);
let deliveredToWorkflow = false;

return {
markDelivered: finish,
markDelivered: () => {
if (deliveredToWorkflow) {
return;
}
deliveredToWorkflow = true;
entry.retire();
},
arm: () => {
if (deliveredToWorkflow) {
return;
}
// The idle safety net may retire an unclaimed buffered hook payload
// while a retained VM keeps its `claim()` closure alive. If workflow
// code later claims that payload, replace the settled entry so delivery
// remains non-idle until the claim reaches the workflow.
if (entry.retired) {
entry = install(true);
return;
}
if (barriers.get(eventIndex) !== entry) {
throw new WorkflowRuntimeError(
`Delivery barrier lost ownership of event index ${eventIndex}`
);
}
entry.armed = true;
},
};
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Force GitHub README to respect dark mode (function() { var style = document.createElement('style'); style.textContent = ' .markdown-body { color-scheme: dark light; } .markdown-body pre { background: #161b22 !important; } .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; } .markdown-body table th, .markdown-body table td { border-color: #30363d !important; } .markdown-body img { background: #0d1117; } .markdown-body blockquote { border-left-color: #8b949e; } .markdown-body hr { border-color: #30363d; } '; document.head.appendChild(style); })(); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' Backport #3879: fix(core): re-arm late-claimed hook deliveries by github-actions[bot] · Pull Request #3905 · vercel/workflow · GitHub
Skip to content
Closed
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/late-hooks-rearm.md
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
---
"@workflow/core": patch
---

Keep late-claimed buffered hook payloads from being preempted by workflow suspension in retained VMs.
46 changes: 45 additions & 1 deletion packages/core/src/async-deserialization-ordering.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -4,7 +4,7 @@ import * as nanoid from 'nanoid';
import { monotonicFactory } from 'ulid';
import { afterEach, describe, expect, it, vi } from 'vitest';
import { EventsConsumer } from './events-consumer.js';
import type { WorkflowOrchestratorContext } from './private.js';
import { isDeliveryIdle, type WorkflowOrchestratorContext } from './private.js';
import { dehydrateStepReturnValue } from './serialization.js';
import { createUseStep } from './step.js';
import { createContext } from './vm/index.js';
Expand DownExpand Up@@ -685,4 +685,48 @@ describe('async deserialization ordering', () => {

expect(ctx.pendingDeliveryBarriers?.size).toBe(0);
});

it('should restore delivery protection when a buffered hook payload is claimed after its barrier retires', async () => {
const payload = await dehydrateStepReturnValue(
'buffered',
'wrun_test',
undefined
);
const ctx = setupWorkflowContext([
{
eventId: 'evnt_0',
runId: 'wrun_test',
eventType: 'hook_received',
correlationId: 'hook_01K11TFZ62YS0YYFDQ3E8B9YCV',
eventData: { payload },
createdAt: new Date(),
},
]);

// Model another payload hydrating in the same replay so the idle safety
// net cannot retire this hook's barrier before the test observes it.
ctx.pendingDeliveries = 1;
const createHook = createCreateHook(ctx);
const hook = createHook<string>();

// The payload arrived before a consumer requested it, so it starts with
// an unarmed barrier and remains buffered after the idle safety net retires
// that barrier.
await vi.waitFor(() => {
expect(ctx.pendingDeliveryBarriers?.size).toBe(1);
});
await ctx.promiseQueue;
expect(ctx.pendingDeliveryBarriers?.size).toBe(1);
ctx.pendingDeliveries = 0;
await vi.waitFor(() => {
expect(ctx.pendingDeliveryBarriers?.size).toBe(0);
});

// Claiming the buffered payload commits it to reaching this consumer. Its
// delivery must become non-idle again until the claim resolves.
const delivery = hook.then((value) => value);
expect(isDeliveryIdle(ctx)).toBe(false);
await expect(delivery).resolves.toBe('buffered');
expect(isDeliveryIdle(ctx)).toBe(true);
});
});
12 changes: 12 additions & 0 deletions packages/core/src/delivery-barrier-dispenser.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -117,6 +117,18 @@ const resumeAtA = new Date('2026-07-27T12:00:05.000Z');
const resumeAtB = new Date('2026-07-27T12:00:06.000Z');

describe('barrier safety-net dispenser', () => {
it('rejects a second barrier owner for the same event index', () => {
const ctx = setupWorkflowContext([]);
const barrier = registerDeliveryBarrier(ctx, 0, 'hook', { armed: false });

expect(() => registerDeliveryBarrier(ctx, 0, 'step')).toThrowError(
'Delivery barrier already registered at event index 0'
);

barrier.markDelivered();
expect(isDeliveryIdle(ctx)).toBe(true);
});

it('suspends only after deliveries parked behind an unclaimed payload have run', async () => {
const ops: Promise<any>[] = [];
const payload = await dehydrateStepReturnValue(
Expand Down
109 changes: 72 additions & 37 deletions packages/core/src/private.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -2,6 +2,7 @@
* Utils used by the bundler when transforming code
*/

import { WorkflowRuntimeError } from '@workflow/errors';
import { withResolvers } from '@workflow/utils';
import type { CryptoKey } from './encryption.js';
import type { EventsConsumer } from './events-consumer.js';
Expand DownExpand Up@@ -215,8 +216,8 @@ export type DeliveryKind = 'hook' | 'wait' | 'step';

interface DeliveryBarrierEntry {
kind: DeliveryKind;
/** Resolves once this delivery has resolved to the workflow. */
delivered: Promise<void>;
/** Resolves once this delivery is handed to the workflow or retired. */
released: Promise<void>;
/**
* Whether this delivery is committed to reaching the workflow without any
* further action by workflow code. True for wait completions and step
Expand All@@ -229,11 +230,15 @@ interface DeliveryBarrierEntry {
* once a consumer takes the payload.
*/
armed: boolean;
/** Whether this entry has been removed and its `released` promise settled. */
retired: boolean;
/**
* Retire this entry: resolve `delivered` and remove it from the registry,
* exactly as `markDelivered` would. Called only by the context's safety-net
* dispenser ({@link ensureBarrierSafetyNet}), and only on the lowest-index
* entry at delivery idle. Idempotent.
* Retire this entry: resolve `released` and remove it from the registry,
* without marking the handle delivered to the workflow. A safety-retired
* buffered payload may therefore install a fresh entry if it is claimed by
* a retained VM later. Called only by the context's safety-net dispenser
* ({@link ensureBarrierSafetyNet}), and only on the lowest-index entry at
* delivery idle. Idempotent.
*/
retire: () => void;
}
Expand DownExpand Up@@ -396,7 +401,7 @@ export async function awaitEarlierDeliveries(
) {
continue;
}
earlier.push(entry.delivered);
earlier.push(entry.released);
}
if (earlier.length > 0) {
await Promise.all(earlier);
Expand All@@ -422,7 +427,7 @@ export async function awaitEarlierDeliveries(
export interface DeliveryBarrier {
/**
* Mark this delivery as delivered to the workflow. Resolves its
* `delivered` promise so any later-in-log delivery gated on it (via
* `released` promise so any later-in-log delivery gated on it (via
* {@link awaitEarlierDeliveries}) may proceed, and removes it from the
* registry. Idempotent.
*/
Expand DownExpand Up@@ -468,41 +473,71 @@ export function registerDeliveryBarrier(
return { markDelivered: () => {}, arm: () => {} };
}

let done = false;
const { promise, resolve } = withResolvers<void>();

const finish = () => {
if (done) {
return;
}
done = true;
if (barriers.get(eventIndex) === entry) {
barriers.delete(eventIndex);
const install = (armed: boolean): DeliveryBarrierEntry => {
if (barriers.has(eventIndex)) {
throw new WorkflowRuntimeError(
`Delivery barrier already registered at event index ${eventIndex}`
);
}
resolve();
const { promise, resolve } = withResolvers<void>();
const entry: DeliveryBarrierEntry = {
kind,
released: promise,
armed,
retired: false,
retire: () => {
if (entry.retired) {
return;
}
entry.retired = true;
if (barriers.get(eventIndex) === entry) {
barriers.delete(eventIndex);
}
resolve();
},
};
barriers.set(eventIndex, entry);

// Safety net: if this delivery is never delivered to the workflow (its
// branch was not taken / the run is suspending, or a buffered hook payload
// is only claimed after a later delivery the workflow is still waiting
// on), it is retired at idle so a later delivery gated on it cannot
// deadlock and the registry cannot leak an entry per abandoned delivery.
// Retirement goes through the context's single ordered dispenser rather
// than a per-barrier idle poll. See {@link ensureBarrierSafetyNet} for why
// the ORDER of these retirements is load-bearing.
ensureBarrierSafetyNet(ctx);
return entry;
};

const entry: DeliveryBarrierEntry = {
kind,
delivered: promise,
armed: options.armed ?? true,
retire: finish,
};
barriers.set(eventIndex, entry);

// Safety net: if this delivery is never delivered to the workflow (its
// branch was not taken / the run is suspending, or a buffered hook payload
// is only claimed after a later delivery the workflow is still waiting on),
// it is retired at idle so a later delivery gated on it cannot deadlock and
// the registry cannot leak an entry per abandoned delivery. Retirement goes
// through the context's single ordered dispenser rather than a per-barrier
// idle poll — see {@link ensureBarrierSafetyNet} for why the ORDER of these
// retirements is load-bearing.
ensureBarrierSafetyNet(ctx);
let entry = install(options.armed ?? true);
let deliveredToWorkflow = false;

return {
markDelivered: finish,
markDelivered: () => {
if (deliveredToWorkflow) {
return;
}
deliveredToWorkflow = true;
entry.retire();
},
arm: () => {
if (deliveredToWorkflow) {
return;
}
// The idle safety net may retire an unclaimed buffered hook payload
// while a retained VM keeps its `claim()` closure alive. If workflow
// code later claims that payload, replace the settled entry so delivery
// remains non-idle until the claim reaches the workflow.
if (entry.retired) {
entry = install(true);
return;
}
if (barriers.get(eventIndex) !== entry) {
throw new WorkflowRuntimeError(
`Delivery barrier lost ownership of event index ${eventIndex}`
);
}
entry.armed = true;
},
};
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Highlight search terms from Google/DuckDuckGo/Bing referrer (function() { var ref = document.referrer; var terms = []; if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) { var url = new URL(ref); var q = url.searchParams.get('q') || url.searchParams.get('p'); if (q) { terms = q.split(/\s+/).filter(function(t) { return t.length > 2; }); } } if (terms.length === 0) return; var style = document.createElement('style'); style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }'; document.head.appendChild(style); function highlight(node) { if (node.nodeType === 3) { // text node var text = node.textContent; var found = false; terms.forEach(function(term) { var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\]\\]/g, '\\') + ')', 'gi'); if (regex.test(text)) { found = true; var frag = document.createDocumentFragment(); var parts = text.split(regex); parts.forEach(function(part, i) { if (i % 2 === 0) { frag.appendChild(document.createTextNode(part)); } else { var span = document.createElement('span'); span.className = 'userscript-highlight'; span.textContent = part; frag.appendChild(span); } }); node.parentNode.replaceChild(frag, node); } }); } else if (node.nodeType === 1 && node.childNodes) { // element var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT']; if (!skipTags.includes(node.tagName)) { Array.from(node.childNodes).forEach(highlight); } } } highlight(document.body); // Re-highlight on dynamic content var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1 || node.nodeType === 3) highlight(node); }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' Backport #3879: fix(core): re-arm late-claimed hook deliveries by github-actions[bot] · Pull Request #3905 · vercel/workflow · GitHub
Skip to content
Closed
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/late-hooks-rearm.md
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
---
"@workflow/core": patch
---

Keep late-claimed buffered hook payloads from being preempted by workflow suspension in retained VMs.
46 changes: 45 additions & 1 deletion packages/core/src/async-deserialization-ordering.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -4,7 +4,7 @@ import * as nanoid from 'nanoid';
import { monotonicFactory } from 'ulid';
import { afterEach, describe, expect, it, vi } from 'vitest';
import { EventsConsumer } from './events-consumer.js';
import type { WorkflowOrchestratorContext } from './private.js';
import { isDeliveryIdle, type WorkflowOrchestratorContext } from './private.js';
import { dehydrateStepReturnValue } from './serialization.js';
import { createUseStep } from './step.js';
import { createContext } from './vm/index.js';
Expand DownExpand Up@@ -685,4 +685,48 @@ describe('async deserialization ordering', () => {

expect(ctx.pendingDeliveryBarriers?.size).toBe(0);
});

it('should restore delivery protection when a buffered hook payload is claimed after its barrier retires', async () => {
const payload = await dehydrateStepReturnValue(
'buffered',
'wrun_test',
undefined
);
const ctx = setupWorkflowContext([
{
eventId: 'evnt_0',
runId: 'wrun_test',
eventType: 'hook_received',
correlationId: 'hook_01K11TFZ62YS0YYFDQ3E8B9YCV',
eventData: { payload },
createdAt: new Date(),
},
]);

// Model another payload hydrating in the same replay so the idle safety
// net cannot retire this hook's barrier before the test observes it.
ctx.pendingDeliveries = 1;
const createHook = createCreateHook(ctx);
const hook = createHook<string>();

// The payload arrived before a consumer requested it, so it starts with
// an unarmed barrier and remains buffered after the idle safety net retires
// that barrier.
await vi.waitFor(() => {
expect(ctx.pendingDeliveryBarriers?.size).toBe(1);
});
await ctx.promiseQueue;
expect(ctx.pendingDeliveryBarriers?.size).toBe(1);
ctx.pendingDeliveries = 0;
await vi.waitFor(() => {
expect(ctx.pendingDeliveryBarriers?.size).toBe(0);
});

// Claiming the buffered payload commits it to reaching this consumer. Its
// delivery must become non-idle again until the claim resolves.
const delivery = hook.then((value) => value);
expect(isDeliveryIdle(ctx)).toBe(false);
await expect(delivery).resolves.toBe('buffered');
expect(isDeliveryIdle(ctx)).toBe(true);
});
});
12 changes: 12 additions & 0 deletions packages/core/src/delivery-barrier-dispenser.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -117,6 +117,18 @@ const resumeAtA = new Date('2026-07-27T12:00:05.000Z');
const resumeAtB = new Date('2026-07-27T12:00:06.000Z');

describe('barrier safety-net dispenser', () => {
it('rejects a second barrier owner for the same event index', () => {
const ctx = setupWorkflowContext([]);
const barrier = registerDeliveryBarrier(ctx, 0, 'hook', { armed: false });

expect(() => registerDeliveryBarrier(ctx, 0, 'step')).toThrowError(
'Delivery barrier already registered at event index 0'
);

barrier.markDelivered();
expect(isDeliveryIdle(ctx)).toBe(true);
});

it('suspends only after deliveries parked behind an unclaimed payload have run', async () => {
const ops: Promise<any>[] = [];
const payload = await dehydrateStepReturnValue(
Expand Down
109 changes: 72 additions & 37 deletions packages/core/src/private.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -2,6 +2,7 @@
* Utils used by the bundler when transforming code
*/

import { WorkflowRuntimeError } from '@workflow/errors';
import { withResolvers } from '@workflow/utils';
import type { CryptoKey } from './encryption.js';
import type { EventsConsumer } from './events-consumer.js';
Expand DownExpand Up@@ -215,8 +216,8 @@ export type DeliveryKind = 'hook' | 'wait' | 'step';

interface DeliveryBarrierEntry {
kind: DeliveryKind;
/** Resolves once this delivery has resolved to the workflow. */
delivered: Promise<void>;
/** Resolves once this delivery is handed to the workflow or retired. */
released: Promise<void>;
/**
* Whether this delivery is committed to reaching the workflow without any
* further action by workflow code. True for wait completions and step
Expand All@@ -229,11 +230,15 @@ interface DeliveryBarrierEntry {
* once a consumer takes the payload.
*/
armed: boolean;
/** Whether this entry has been removed and its `released` promise settled. */
retired: boolean;
/**
* Retire this entry: resolve `delivered` and remove it from the registry,
* exactly as `markDelivered` would. Called only by the context's safety-net
* dispenser ({@link ensureBarrierSafetyNet}), and only on the lowest-index
* entry at delivery idle. Idempotent.
* Retire this entry: resolve `released` and remove it from the registry,
* without marking the handle delivered to the workflow. A safety-retired
* buffered payload may therefore install a fresh entry if it is claimed by
* a retained VM later. Called only by the context's safety-net dispenser
* ({@link ensureBarrierSafetyNet}), and only on the lowest-index entry at
* delivery idle. Idempotent.
*/
retire: () => void;
}
Expand DownExpand Up@@ -396,7 +401,7 @@ export async function awaitEarlierDeliveries(
) {
continue;
}
earlier.push(entry.delivered);
earlier.push(entry.released);
}
if (earlier.length > 0) {
await Promise.all(earlier);
Expand All@@ -422,7 +427,7 @@ export async function awaitEarlierDeliveries(
export interface DeliveryBarrier {
/**
* Mark this delivery as delivered to the workflow. Resolves its
* `delivered` promise so any later-in-log delivery gated on it (via
* `released` promise so any later-in-log delivery gated on it (via
* {@link awaitEarlierDeliveries}) may proceed, and removes it from the
* registry. Idempotent.
*/
Expand DownExpand Up@@ -468,41 +473,71 @@ export function registerDeliveryBarrier(
return { markDelivered: () => {}, arm: () => {} };
}

let done = false;
const { promise, resolve } = withResolvers<void>();

const finish = () => {
if (done) {
return;
}
done = true;
if (barriers.get(eventIndex) === entry) {
barriers.delete(eventIndex);
const install = (armed: boolean): DeliveryBarrierEntry => {
if (barriers.has(eventIndex)) {
throw new WorkflowRuntimeError(
`Delivery barrier already registered at event index ${eventIndex}`
);
}
resolve();
const { promise, resolve } = withResolvers<void>();
const entry: DeliveryBarrierEntry = {
kind,
released: promise,
armed,
retired: false,
retire: () => {
if (entry.retired) {
return;
}
entry.retired = true;
if (barriers.get(eventIndex) === entry) {
barriers.delete(eventIndex);
}
resolve();
},
};
barriers.set(eventIndex, entry);

// Safety net: if this delivery is never delivered to the workflow (its
// branch was not taken / the run is suspending, or a buffered hook payload
// is only claimed after a later delivery the workflow is still waiting
// on), it is retired at idle so a later delivery gated on it cannot
// deadlock and the registry cannot leak an entry per abandoned delivery.
// Retirement goes through the context's single ordered dispenser rather
// than a per-barrier idle poll. See {@link ensureBarrierSafetyNet} for why
// the ORDER of these retirements is load-bearing.
ensureBarrierSafetyNet(ctx);
return entry;
};

const entry: DeliveryBarrierEntry = {
kind,
delivered: promise,
armed: options.armed ?? true,
retire: finish,
};
barriers.set(eventIndex, entry);

// Safety net: if this delivery is never delivered to the workflow (its
// branch was not taken / the run is suspending, or a buffered hook payload
// is only claimed after a later delivery the workflow is still waiting on),
// it is retired at idle so a later delivery gated on it cannot deadlock and
// the registry cannot leak an entry per abandoned delivery. Retirement goes
// through the context's single ordered dispenser rather than a per-barrier
// idle poll — see {@link ensureBarrierSafetyNet} for why the ORDER of these
// retirements is load-bearing.
ensureBarrierSafetyNet(ctx);
let entry = install(options.armed ?? true);
let deliveredToWorkflow = false;

return {
markDelivered: finish,
markDelivered: () => {
if (deliveredToWorkflow) {
return;
}
deliveredToWorkflow = true;
entry.retire();
},
arm: () => {
if (deliveredToWorkflow) {
return;
}
// The idle safety net may retire an unclaimed buffered hook payload
// while a retained VM keeps its `claim()` closure alive. If workflow
// code later claims that payload, replace the settled entry so delivery
// remains non-idle until the claim reaches the workflow.
if (entry.retired) {
entry = install(true);
return;
}
if (barriers.get(eventIndex) !== entry) {
throw new WorkflowRuntimeError(
`Delivery barrier lost ownership of event index ${eventIndex}`
);
}
entry.armed = true;
},
};
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Strip utm_, fbclid, gclid, etc. from all links on page (function() { var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content', 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid', 'ref', 'ref_src', 'source', 'medium', 'campaign']; function cleanUrl(url) { try { var u = new URL(url, window.location.origin); var changed = false; trackingParams.forEach(function(p) { if (u.searchParams.has(p)) { u.searchParams.delete(p); changed = true; } }); return changed ? u.toString() : url; } catch (e) { return url; } } function cleanLinks() { document.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } cleanLinks(); var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1) { if (node.tagName === 'A') cleanLinks(); node.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + ' Backport #3879: fix(core): re-arm late-claimed hook deliveries by github-actions[bot] · Pull Request #3905 · vercel/workflow · GitHub
Skip to content
Closed
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/late-hooks-rearm.md
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
---
"@workflow/core": patch
---

Keep late-claimed buffered hook payloads from being preempted by workflow suspension in retained VMs.
46 changes: 45 additions & 1 deletion packages/core/src/async-deserialization-ordering.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -4,7 +4,7 @@ import * as nanoid from 'nanoid';
import { monotonicFactory } from 'ulid';
import { afterEach, describe, expect, it, vi } from 'vitest';
import { EventsConsumer } from './events-consumer.js';
import type { WorkflowOrchestratorContext } from './private.js';
import { isDeliveryIdle, type WorkflowOrchestratorContext } from './private.js';
import { dehydrateStepReturnValue } from './serialization.js';
import { createUseStep } from './step.js';
import { createContext } from './vm/index.js';
Expand DownExpand Up@@ -685,4 +685,48 @@ describe('async deserialization ordering', () => {

expect(ctx.pendingDeliveryBarriers?.size).toBe(0);
});

it('should restore delivery protection when a buffered hook payload is claimed after its barrier retires', async () => {
const payload = await dehydrateStepReturnValue(
'buffered',
'wrun_test',
undefined
);
const ctx = setupWorkflowContext([
{
eventId: 'evnt_0',
runId: 'wrun_test',
eventType: 'hook_received',
correlationId: 'hook_01K11TFZ62YS0YYFDQ3E8B9YCV',
eventData: { payload },
createdAt: new Date(),
},
]);

// Model another payload hydrating in the same replay so the idle safety
// net cannot retire this hook's barrier before the test observes it.
ctx.pendingDeliveries = 1;
const createHook = createCreateHook(ctx);
const hook = createHook<string>();

// The payload arrived before a consumer requested it, so it starts with
// an unarmed barrier and remains buffered after the idle safety net retires
// that barrier.
await vi.waitFor(() => {
expect(ctx.pendingDeliveryBarriers?.size).toBe(1);
});
await ctx.promiseQueue;
expect(ctx.pendingDeliveryBarriers?.size).toBe(1);
ctx.pendingDeliveries = 0;
await vi.waitFor(() => {
expect(ctx.pendingDeliveryBarriers?.size).toBe(0);
});

// Claiming the buffered payload commits it to reaching this consumer. Its
// delivery must become non-idle again until the claim resolves.
const delivery = hook.then((value) => value);
expect(isDeliveryIdle(ctx)).toBe(false);
await expect(delivery).resolves.toBe('buffered');
expect(isDeliveryIdle(ctx)).toBe(true);
});
});
12 changes: 12 additions & 0 deletions packages/core/src/delivery-barrier-dispenser.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -117,6 +117,18 @@ const resumeAtA = new Date('2026-07-27T12:00:05.000Z');
const resumeAtB = new Date('2026-07-27T12:00:06.000Z');

describe('barrier safety-net dispenser', () => {
it('rejects a second barrier owner for the same event index', () => {
const ctx = setupWorkflowContext([]);
const barrier = registerDeliveryBarrier(ctx, 0, 'hook', { armed: false });

expect(() => registerDeliveryBarrier(ctx, 0, 'step')).toThrowError(
'Delivery barrier already registered at event index 0'
);

barrier.markDelivered();
expect(isDeliveryIdle(ctx)).toBe(true);
});

it('suspends only after deliveries parked behind an unclaimed payload have run', async () => {
const ops: Promise<any>[] = [];
const payload = await dehydrateStepReturnValue(
Expand Down
109 changes: 72 additions & 37 deletions packages/core/src/private.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -2,6 +2,7 @@
* Utils used by the bundler when transforming code
*/

import { WorkflowRuntimeError } from '@workflow/errors';
import { withResolvers } from '@workflow/utils';
import type { CryptoKey } from './encryption.js';
import type { EventsConsumer } from './events-consumer.js';
Expand DownExpand Up@@ -215,8 +216,8 @@ export type DeliveryKind = 'hook' | 'wait' | 'step';

interface DeliveryBarrierEntry {
kind: DeliveryKind;
/** Resolves once this delivery has resolved to the workflow. */
delivered: Promise<void>;
/** Resolves once this delivery is handed to the workflow or retired. */
released: Promise<void>;
/**
* Whether this delivery is committed to reaching the workflow without any
* further action by workflow code. True for wait completions and step
Expand All@@ -229,11 +230,15 @@ interface DeliveryBarrierEntry {
* once a consumer takes the payload.
*/
armed: boolean;
/** Whether this entry has been removed and its `released` promise settled. */
retired: boolean;
/**
* Retire this entry: resolve `delivered` and remove it from the registry,
* exactly as `markDelivered` would. Called only by the context's safety-net
* dispenser ({@link ensureBarrierSafetyNet}), and only on the lowest-index
* entry at delivery idle. Idempotent.
* Retire this entry: resolve `released` and remove it from the registry,
* without marking the handle delivered to the workflow. A safety-retired
* buffered payload may therefore install a fresh entry if it is claimed by
* a retained VM later. Called only by the context's safety-net dispenser
* ({@link ensureBarrierSafetyNet}), and only on the lowest-index entry at
* delivery idle. Idempotent.
*/
retire: () => void;
}
Expand DownExpand Up@@ -396,7 +401,7 @@ export async function awaitEarlierDeliveries(
) {
continue;
}
earlier.push(entry.delivered);
earlier.push(entry.released);
}
if (earlier.length > 0) {
await Promise.all(earlier);
Expand All@@ -422,7 +427,7 @@ export async function awaitEarlierDeliveries(
export interface DeliveryBarrier {
/**
* Mark this delivery as delivered to the workflow. Resolves its
* `delivered` promise so any later-in-log delivery gated on it (via
* `released` promise so any later-in-log delivery gated on it (via
* {@link awaitEarlierDeliveries}) may proceed, and removes it from the
* registry. Idempotent.
*/
Expand DownExpand Up@@ -468,41 +473,71 @@ export function registerDeliveryBarrier(
return { markDelivered: () => {}, arm: () => {} };
}

let done = false;
const { promise, resolve } = withResolvers<void>();

const finish = () => {
if (done) {
return;
}
done = true;
if (barriers.get(eventIndex) === entry) {
barriers.delete(eventIndex);
const install = (armed: boolean): DeliveryBarrierEntry => {
if (barriers.has(eventIndex)) {
throw new WorkflowRuntimeError(
`Delivery barrier already registered at event index ${eventIndex}`
);
}
resolve();
const { promise, resolve } = withResolvers<void>();
const entry: DeliveryBarrierEntry = {
kind,
released: promise,
armed,
retired: false,
retire: () => {
if (entry.retired) {
return;
}
entry.retired = true;
if (barriers.get(eventIndex) === entry) {
barriers.delete(eventIndex);
}
resolve();
},
};
barriers.set(eventIndex, entry);

// Safety net: if this delivery is never delivered to the workflow (its
// branch was not taken / the run is suspending, or a buffered hook payload
// is only claimed after a later delivery the workflow is still waiting
// on), it is retired at idle so a later delivery gated on it cannot
// deadlock and the registry cannot leak an entry per abandoned delivery.
// Retirement goes through the context's single ordered dispenser rather
// than a per-barrier idle poll. See {@link ensureBarrierSafetyNet} for why
// the ORDER of these retirements is load-bearing.
ensureBarrierSafetyNet(ctx);
return entry;
};

const entry: DeliveryBarrierEntry = {
kind,
delivered: promise,
armed: options.armed ?? true,
retire: finish,
};
barriers.set(eventIndex, entry);

// Safety net: if this delivery is never delivered to the workflow (its
// branch was not taken / the run is suspending, or a buffered hook payload
// is only claimed after a later delivery the workflow is still waiting on),
// it is retired at idle so a later delivery gated on it cannot deadlock and
// the registry cannot leak an entry per abandoned delivery. Retirement goes
// through the context's single ordered dispenser rather than a per-barrier
// idle poll — see {@link ensureBarrierSafetyNet} for why the ORDER of these
// retirements is load-bearing.
ensureBarrierSafetyNet(ctx);
let entry = install(options.armed ?? true);
let deliveredToWorkflow = false;

return {
markDelivered: finish,
markDelivered: () => {
if (deliveredToWorkflow) {
return;
}
deliveredToWorkflow = true;
entry.retire();
},
arm: () => {
if (deliveredToWorkflow) {
return;
}
// The idle safety net may retire an unclaimed buffered hook payload
// while a retained VM keeps its `claim()` closure alive. If workflow
// code later claims that payload, replace the settled entry so delivery
// remains non-idle until the claim reaches the workflow.
if (entry.retired) {
entry = install(true);
return;
}
if (barriers.get(eventIndex) !== entry) {
throw new WorkflowRuntimeError(
`Delivery barrier lost ownership of event index ${eventIndex}`
);
}
entry.armed = true;
},
};
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Auto-enable theater mode on YouTube (function() { function tryTheater() { var btn = document.querySelector('button[aria-label="Theater mode"], ytd-player #player button[title="Theater mode"]'); if (btn && !btn.classList.contains('activated')) { btn.click(); } } // Try immediately tryTheater(); // Try after navigation (SPA) var lastUrl = location.href; setInterval(function() { if (location.href !== lastUrl) { lastUrl = location.href; setTimeout(tryTheater, 500); } }, 1000); // Also try on player load var observer = new MutationObserver(tryTheater); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' Backport #3879: fix(core): re-arm late-claimed hook deliveries by github-actions[bot] · Pull Request #3905 · vercel/workflow · GitHub
Skip to content
Closed
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/late-hooks-rearm.md
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
---
"@workflow/core": patch
---

Keep late-claimed buffered hook payloads from being preempted by workflow suspension in retained VMs.
46 changes: 45 additions & 1 deletion packages/core/src/async-deserialization-ordering.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -4,7 +4,7 @@ import * as nanoid from 'nanoid';
import { monotonicFactory } from 'ulid';
import { afterEach, describe, expect, it, vi } from 'vitest';
import { EventsConsumer } from './events-consumer.js';
import type { WorkflowOrchestratorContext } from './private.js';
import { isDeliveryIdle, type WorkflowOrchestratorContext } from './private.js';
import { dehydrateStepReturnValue } from './serialization.js';
import { createUseStep } from './step.js';
import { createContext } from './vm/index.js';
Expand DownExpand Up@@ -685,4 +685,48 @@ describe('async deserialization ordering', () => {

expect(ctx.pendingDeliveryBarriers?.size).toBe(0);
});

it('should restore delivery protection when a buffered hook payload is claimed after its barrier retires', async () => {
const payload = await dehydrateStepReturnValue(
'buffered',
'wrun_test',
undefined
);
const ctx = setupWorkflowContext([
{
eventId: 'evnt_0',
runId: 'wrun_test',
eventType: 'hook_received',
correlationId: 'hook_01K11TFZ62YS0YYFDQ3E8B9YCV',
eventData: { payload },
createdAt: new Date(),
},
]);

// Model another payload hydrating in the same replay so the idle safety
// net cannot retire this hook's barrier before the test observes it.
ctx.pendingDeliveries = 1;
const createHook = createCreateHook(ctx);
const hook = createHook<string>();

// The payload arrived before a consumer requested it, so it starts with
// an unarmed barrier and remains buffered after the idle safety net retires
// that barrier.
await vi.waitFor(() => {
expect(ctx.pendingDeliveryBarriers?.size).toBe(1);
});
await ctx.promiseQueue;
expect(ctx.pendingDeliveryBarriers?.size).toBe(1);
ctx.pendingDeliveries = 0;
await vi.waitFor(() => {
expect(ctx.pendingDeliveryBarriers?.size).toBe(0);
});

// Claiming the buffered payload commits it to reaching this consumer. Its
// delivery must become non-idle again until the claim resolves.
const delivery = hook.then((value) => value);
expect(isDeliveryIdle(ctx)).toBe(false);
await expect(delivery).resolves.toBe('buffered');
expect(isDeliveryIdle(ctx)).toBe(true);
});
});
12 changes: 12 additions & 0 deletions packages/core/src/delivery-barrier-dispenser.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -117,6 +117,18 @@ const resumeAtA = new Date('2026-07-27T12:00:05.000Z');
const resumeAtB = new Date('2026-07-27T12:00:06.000Z');

describe('barrier safety-net dispenser', () => {
it('rejects a second barrier owner for the same event index', () => {
const ctx = setupWorkflowContext([]);
const barrier = registerDeliveryBarrier(ctx, 0, 'hook', { armed: false });

expect(() => registerDeliveryBarrier(ctx, 0, 'step')).toThrowError(
'Delivery barrier already registered at event index 0'
);

barrier.markDelivered();
expect(isDeliveryIdle(ctx)).toBe(true);
});

it('suspends only after deliveries parked behind an unclaimed payload have run', async () => {
const ops: Promise<any>[] = [];
const payload = await dehydrateStepReturnValue(
Expand Down
109 changes: 72 additions & 37 deletions packages/core/src/private.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -2,6 +2,7 @@
* Utils used by the bundler when transforming code
*/

import { WorkflowRuntimeError } from '@workflow/errors';
import { withResolvers } from '@workflow/utils';
import type { CryptoKey } from './encryption.js';
import type { EventsConsumer } from './events-consumer.js';
Expand DownExpand Up@@ -215,8 +216,8 @@ export type DeliveryKind = 'hook' | 'wait' | 'step';

interface DeliveryBarrierEntry {
kind: DeliveryKind;
/** Resolves once this delivery has resolved to the workflow. */
delivered: Promise<void>;
/** Resolves once this delivery is handed to the workflow or retired. */
released: Promise<void>;
/**
* Whether this delivery is committed to reaching the workflow without any
* further action by workflow code. True for wait completions and step
Expand All@@ -229,11 +230,15 @@ interface DeliveryBarrierEntry {
* once a consumer takes the payload.
*/
armed: boolean;
/** Whether this entry has been removed and its `released` promise settled. */
retired: boolean;
/**
* Retire this entry: resolve `delivered` and remove it from the registry,
* exactly as `markDelivered` would. Called only by the context's safety-net
* dispenser ({@link ensureBarrierSafetyNet}), and only on the lowest-index
* entry at delivery idle. Idempotent.
* Retire this entry: resolve `released` and remove it from the registry,
* without marking the handle delivered to the workflow. A safety-retired
* buffered payload may therefore install a fresh entry if it is claimed by
* a retained VM later. Called only by the context's safety-net dispenser
* ({@link ensureBarrierSafetyNet}), and only on the lowest-index entry at
* delivery idle. Idempotent.
*/
retire: () => void;
}
Expand DownExpand Up@@ -396,7 +401,7 @@ export async function awaitEarlierDeliveries(
) {
continue;
}
earlier.push(entry.delivered);
earlier.push(entry.released);
}
if (earlier.length > 0) {
await Promise.all(earlier);
Expand All@@ -422,7 +427,7 @@ export async function awaitEarlierDeliveries(
export interface DeliveryBarrier {
/**
* Mark this delivery as delivered to the workflow. Resolves its
* `delivered` promise so any later-in-log delivery gated on it (via
* `released` promise so any later-in-log delivery gated on it (via
* {@link awaitEarlierDeliveries}) may proceed, and removes it from the
* registry. Idempotent.
*/
Expand DownExpand Up@@ -468,41 +473,71 @@ export function registerDeliveryBarrier(
return { markDelivered: () => {}, arm: () => {} };
}

let done = false;
const { promise, resolve } = withResolvers<void>();

const finish = () => {
if (done) {
return;
}
done = true;
if (barriers.get(eventIndex) === entry) {
barriers.delete(eventIndex);
const install = (armed: boolean): DeliveryBarrierEntry => {
if (barriers.has(eventIndex)) {
throw new WorkflowRuntimeError(
`Delivery barrier already registered at event index ${eventIndex}`
);
}
resolve();
const { promise, resolve } = withResolvers<void>();
const entry: DeliveryBarrierEntry = {
kind,
released: promise,
armed,
retired: false,
retire: () => {
if (entry.retired) {
return;
}
entry.retired = true;
if (barriers.get(eventIndex) === entry) {
barriers.delete(eventIndex);
}
resolve();
},
};
barriers.set(eventIndex, entry);

// Safety net: if this delivery is never delivered to the workflow (its
// branch was not taken / the run is suspending, or a buffered hook payload
// is only claimed after a later delivery the workflow is still waiting
// on), it is retired at idle so a later delivery gated on it cannot
// deadlock and the registry cannot leak an entry per abandoned delivery.
// Retirement goes through the context's single ordered dispenser rather
// than a per-barrier idle poll. See {@link ensureBarrierSafetyNet} for why
// the ORDER of these retirements is load-bearing.
ensureBarrierSafetyNet(ctx);
return entry;
};

const entry: DeliveryBarrierEntry = {
kind,
delivered: promise,
armed: options.armed ?? true,
retire: finish,
};
barriers.set(eventIndex, entry);

// Safety net: if this delivery is never delivered to the workflow (its
// branch was not taken / the run is suspending, or a buffered hook payload
// is only claimed after a later delivery the workflow is still waiting on),
// it is retired at idle so a later delivery gated on it cannot deadlock and
// the registry cannot leak an entry per abandoned delivery. Retirement goes
// through the context's single ordered dispenser rather than a per-barrier
// idle poll — see {@link ensureBarrierSafetyNet} for why the ORDER of these
// retirements is load-bearing.
ensureBarrierSafetyNet(ctx);
let entry = install(options.armed ?? true);
let deliveredToWorkflow = false;

return {
markDelivered: finish,
markDelivered: () => {
if (deliveredToWorkflow) {
return;
}
deliveredToWorkflow = true;
entry.retire();
},
arm: () => {
if (deliveredToWorkflow) {
return;
}
// The idle safety net may retire an unclaimed buffered hook payload
// while a retained VM keeps its `claim()` closure alive. If workflow
// code later claims that payload, replace the settled entry so delivery
// remains non-idle until the claim reaches the workflow.
if (entry.retired) {
entry = install(true);
return;
}
if (barriers.get(eventIndex) !== entry) {
throw new WorkflowRuntimeError(
`Delivery barrier lost ownership of event index ${eventIndex}`
);
}
entry.armed = true;
},
};
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Remove or un-stick sticky/fixed headers that block content (function() { function unstick() { document.querySelectorAll('header, nav, [role="banner"], .header, .navbar, .sticky, .fixed-top, [style*="position: fixed"], [style*="position:sticky"]').forEach(function(el) { if (el.style.position === 'fixed' || el.style.position === 'sticky' || getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') { el.style.position = 'static'; el.style.top = 'auto'; el.style.zIndex = 'auto'; } }); } unstick(); var observer = new MutationObserver(unstick); observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] }); })(); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' Backport #3879: fix(core): re-arm late-claimed hook deliveries by github-actions[bot] · Pull Request #3905 · vercel/workflow · GitHub
Skip to content
Closed
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/late-hooks-rearm.md
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
---
"@workflow/core": patch
---

Keep late-claimed buffered hook payloads from being preempted by workflow suspension in retained VMs.
46 changes: 45 additions & 1 deletion packages/core/src/async-deserialization-ordering.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -4,7 +4,7 @@ import * as nanoid from 'nanoid';
import { monotonicFactory } from 'ulid';
import { afterEach, describe, expect, it, vi } from 'vitest';
import { EventsConsumer } from './events-consumer.js';
import type { WorkflowOrchestratorContext } from './private.js';
import { isDeliveryIdle, type WorkflowOrchestratorContext } from './private.js';
import { dehydrateStepReturnValue } from './serialization.js';
import { createUseStep } from './step.js';
import { createContext } from './vm/index.js';
Expand DownExpand Up@@ -685,4 +685,48 @@ describe('async deserialization ordering', () => {

expect(ctx.pendingDeliveryBarriers?.size).toBe(0);
});

it('should restore delivery protection when a buffered hook payload is claimed after its barrier retires', async () => {
const payload = await dehydrateStepReturnValue(
'buffered',
'wrun_test',
undefined
);
const ctx = setupWorkflowContext([
{
eventId: 'evnt_0',
runId: 'wrun_test',
eventType: 'hook_received',
correlationId: 'hook_01K11TFZ62YS0YYFDQ3E8B9YCV',
eventData: { payload },
createdAt: new Date(),
},
]);

// Model another payload hydrating in the same replay so the idle safety
// net cannot retire this hook's barrier before the test observes it.
ctx.pendingDeliveries = 1;
const createHook = createCreateHook(ctx);
const hook = createHook<string>();

// The payload arrived before a consumer requested it, so it starts with
// an unarmed barrier and remains buffered after the idle safety net retires
// that barrier.
await vi.waitFor(() => {
expect(ctx.pendingDeliveryBarriers?.size).toBe(1);
});
await ctx.promiseQueue;
expect(ctx.pendingDeliveryBarriers?.size).toBe(1);
ctx.pendingDeliveries = 0;
await vi.waitFor(() => {
expect(ctx.pendingDeliveryBarriers?.size).toBe(0);
});

// Claiming the buffered payload commits it to reaching this consumer. Its
// delivery must become non-idle again until the claim resolves.
const delivery = hook.then((value) => value);
expect(isDeliveryIdle(ctx)).toBe(false);
await expect(delivery).resolves.toBe('buffered');
expect(isDeliveryIdle(ctx)).toBe(true);
});
});
12 changes: 12 additions & 0 deletions packages/core/src/delivery-barrier-dispenser.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -117,6 +117,18 @@ const resumeAtA = new Date('2026-07-27T12:00:05.000Z');
const resumeAtB = new Date('2026-07-27T12:00:06.000Z');

describe('barrier safety-net dispenser', () => {
it('rejects a second barrier owner for the same event index', () => {
const ctx = setupWorkflowContext([]);
const barrier = registerDeliveryBarrier(ctx, 0, 'hook', { armed: false });

expect(() => registerDeliveryBarrier(ctx, 0, 'step')).toThrowError(
'Delivery barrier already registered at event index 0'
);

barrier.markDelivered();
expect(isDeliveryIdle(ctx)).toBe(true);
});

it('suspends only after deliveries parked behind an unclaimed payload have run', async () => {
const ops: Promise<any>[] = [];
const payload = await dehydrateStepReturnValue(
Expand Down
109 changes: 72 additions & 37 deletions packages/core/src/private.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -2,6 +2,7 @@
* Utils used by the bundler when transforming code
*/

import { WorkflowRuntimeError } from '@workflow/errors';
import { withResolvers } from '@workflow/utils';
import type { CryptoKey } from './encryption.js';
import type { EventsConsumer } from './events-consumer.js';
Expand DownExpand Up@@ -215,8 +216,8 @@ export type DeliveryKind = 'hook' | 'wait' | 'step';

interface DeliveryBarrierEntry {
kind: DeliveryKind;
/** Resolves once this delivery has resolved to the workflow. */
delivered: Promise<void>;
/** Resolves once this delivery is handed to the workflow or retired. */
released: Promise<void>;
/**
* Whether this delivery is committed to reaching the workflow without any
* further action by workflow code. True for wait completions and step
Expand All@@ -229,11 +230,15 @@ interface DeliveryBarrierEntry {
* once a consumer takes the payload.
*/
armed: boolean;
/** Whether this entry has been removed and its `released` promise settled. */
retired: boolean;
/**
* Retire this entry: resolve `delivered` and remove it from the registry,
* exactly as `markDelivered` would. Called only by the context's safety-net
* dispenser ({@link ensureBarrierSafetyNet}), and only on the lowest-index
* entry at delivery idle. Idempotent.
* Retire this entry: resolve `released` and remove it from the registry,
* without marking the handle delivered to the workflow. A safety-retired
* buffered payload may therefore install a fresh entry if it is claimed by
* a retained VM later. Called only by the context's safety-net dispenser
* ({@link ensureBarrierSafetyNet}), and only on the lowest-index entry at
* delivery idle. Idempotent.
*/
retire: () => void;
}
Expand DownExpand Up@@ -396,7 +401,7 @@ export async function awaitEarlierDeliveries(
) {
continue;
}
earlier.push(entry.delivered);
earlier.push(entry.released);
}
if (earlier.length > 0) {
await Promise.all(earlier);
Expand All@@ -422,7 +427,7 @@ export async function awaitEarlierDeliveries(
export interface DeliveryBarrier {
/**
* Mark this delivery as delivered to the workflow. Resolves its
* `delivered` promise so any later-in-log delivery gated on it (via
* `released` promise so any later-in-log delivery gated on it (via
* {@link awaitEarlierDeliveries}) may proceed, and removes it from the
* registry. Idempotent.
*/
Expand DownExpand Up@@ -468,41 +473,71 @@ export function registerDeliveryBarrier(
return { markDelivered: () => {}, arm: () => {} };
}

let done = false;
const { promise, resolve } = withResolvers<void>();

const finish = () => {
if (done) {
return;
}
done = true;
if (barriers.get(eventIndex) === entry) {
barriers.delete(eventIndex);
const install = (armed: boolean): DeliveryBarrierEntry => {
if (barriers.has(eventIndex)) {
throw new WorkflowRuntimeError(
`Delivery barrier already registered at event index ${eventIndex}`
);
}
resolve();
const { promise, resolve } = withResolvers<void>();
const entry: DeliveryBarrierEntry = {
kind,
released: promise,
armed,
retired: false,
retire: () => {
if (entry.retired) {
return;
}
entry.retired = true;
if (barriers.get(eventIndex) === entry) {
barriers.delete(eventIndex);
}
resolve();
},
};
barriers.set(eventIndex, entry);

// Safety net: if this delivery is never delivered to the workflow (its
// branch was not taken / the run is suspending, or a buffered hook payload
// is only claimed after a later delivery the workflow is still waiting
// on), it is retired at idle so a later delivery gated on it cannot
// deadlock and the registry cannot leak an entry per abandoned delivery.
// Retirement goes through the context's single ordered dispenser rather
// than a per-barrier idle poll. See {@link ensureBarrierSafetyNet} for why
// the ORDER of these retirements is load-bearing.
ensureBarrierSafetyNet(ctx);
return entry;
};

const entry: DeliveryBarrierEntry = {
kind,
delivered: promise,
armed: options.armed ?? true,
retire: finish,
};
barriers.set(eventIndex, entry);

// Safety net: if this delivery is never delivered to the workflow (its
// branch was not taken / the run is suspending, or a buffered hook payload
// is only claimed after a later delivery the workflow is still waiting on),
// it is retired at idle so a later delivery gated on it cannot deadlock and
// the registry cannot leak an entry per abandoned delivery. Retirement goes
// through the context's single ordered dispenser rather than a per-barrier
// idle poll — see {@link ensureBarrierSafetyNet} for why the ORDER of these
// retirements is load-bearing.
ensureBarrierSafetyNet(ctx);
let entry = install(options.armed ?? true);
let deliveredToWorkflow = false;

return {
markDelivered: finish,
markDelivered: () => {
if (deliveredToWorkflow) {
return;
}
deliveredToWorkflow = true;
entry.retire();
},
arm: () => {
if (deliveredToWorkflow) {
return;
}
// The idle safety net may retire an unclaimed buffered hook payload
// while a retained VM keeps its `claim()` closure alive. If workflow
// code later claims that payload, replace the settled entry so delivery
// remains non-idle until the claim reaches the workflow.
if (entry.retired) {
entry = install(true);
return;
}
if (barriers.get(eventIndex) !== entry) {
throw new WorkflowRuntimeError(
`Delivery barrier lost ownership of event index ${eventIndex}`
);
}
entry.armed = true;
},
};
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Universal Dark Mode - works on any site (function() { var enabled = true; function applyDarkMode() { if (!enabled) return; // Create style element if it doesn't exist var style = document.getElementById('universal-dark-mode-style'); if (!style) { style = document.createElement('style'); style.id = 'universal-dark-mode-style'; document.head.appendChild(style); } // Dark mode CSS - inverts colors but preserves images/video style.textContent = ' /* Invert everything except media */ html { filter: invert(1) hue-rotate(180deg) !important; background: #1a1a2e !important; } /* Restore images, videos, iframes, canvas */ img, video, iframe, canvas, svg, picture, [style*="background-image"] { filter: invert(1) hue-rotate(180deg) !important; } /* Preserve specific elements that should not be inverted */ .no-dark-mode, .no-dark-mode *, [data-theme="light"], [data-theme="light"], .ace_editor, .ace_editor *, .CodeMirror, .CodeMirror *, .monaco-editor, .monaco-editor *, .markdown-body pre, .markdown-body pre *, .highlight, .highlight *, pre code, pre code * { filter: none !important; } /* Fix common UI elements */ .modal, .popup, .dropdown-menu, .tooltip, .popover { filter: invert(1) hue-rotate(180deg) !important; background: #2d2d44 !important; border-color: #444 !important; } /* Scrollbars */ ::-webkit-scrollbar { background: #1a1a2e !important; } ::-webkit-scrollbar-thumb { background: #444 !important; } ::-webkit-scrollbar-thumb:hover { background: #555 !important; } /* Selection */ ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; } ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; } '; } function removeDarkMode() { var style = document.getElementById('universal-dark-mode-style'); if (style) style.remove(); } // Toggle with Alt+Shift+D document.addEventListener('keydown', function(e) { if (e.altKey && e.shiftKey && e.key === 'D') { e.preventDefault(); enabled = !enabled; if (enabled) { applyDarkMode(); console.log('[Universal Dark Mode] Enabled'); } else { removeDarkMode(); console.log('[Universal Dark Mode] Disabled'); } } }); // Apply on load applyDarkMode(); // Re-apply on dynamic content var observer = new MutationObserver(function(mutations) { if (enabled && !document.getElementById('universal-dark-mode-style')) { applyDarkMode(); } }); observer.observe(document.head, { childList: true }); console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle'); })(); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })(); Backport #3879: fix(core): re-arm late-claimed hook deliveries by github-actions[bot] · Pull Request #3905 · vercel/workflow · GitHub
Skip to content
Closed
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/late-hooks-rearm.md
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
---
"@workflow/core": patch
---

Keep late-claimed buffered hook payloads from being preempted by workflow suspension in retained VMs.
46 changes: 45 additions & 1 deletion packages/core/src/async-deserialization-ordering.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -4,7 +4,7 @@ import * as nanoid from 'nanoid';
import { monotonicFactory } from 'ulid';
import { afterEach, describe, expect, it, vi } from 'vitest';
import { EventsConsumer } from './events-consumer.js';
import type { WorkflowOrchestratorContext } from './private.js';
import { isDeliveryIdle, type WorkflowOrchestratorContext } from './private.js';
import { dehydrateStepReturnValue } from './serialization.js';
import { createUseStep } from './step.js';
import { createContext } from './vm/index.js';
Expand DownExpand Up@@ -685,4 +685,48 @@ describe('async deserialization ordering', () => {

expect(ctx.pendingDeliveryBarriers?.size).toBe(0);
});

it('should restore delivery protection when a buffered hook payload is claimed after its barrier retires', async () => {
const payload = await dehydrateStepReturnValue(
'buffered',
'wrun_test',
undefined
);
const ctx = setupWorkflowContext([
{
eventId: 'evnt_0',
runId: 'wrun_test',
eventType: 'hook_received',
correlationId: 'hook_01K11TFZ62YS0YYFDQ3E8B9YCV',
eventData: { payload },
createdAt: new Date(),
},
]);

// Model another payload hydrating in the same replay so the idle safety
// net cannot retire this hook's barrier before the test observes it.
ctx.pendingDeliveries = 1;
const createHook = createCreateHook(ctx);
const hook = createHook<string>();

// The payload arrived before a consumer requested it, so it starts with
// an unarmed barrier and remains buffered after the idle safety net retires
// that barrier.
await vi.waitFor(() => {
expect(ctx.pendingDeliveryBarriers?.size).toBe(1);
});
await ctx.promiseQueue;
expect(ctx.pendingDeliveryBarriers?.size).toBe(1);
ctx.pendingDeliveries = 0;
await vi.waitFor(() => {
expect(ctx.pendingDeliveryBarriers?.size).toBe(0);
});

// Claiming the buffered payload commits it to reaching this consumer. Its
// delivery must become non-idle again until the claim resolves.
const delivery = hook.then((value) => value);
expect(isDeliveryIdle(ctx)).toBe(false);
await expect(delivery).resolves.toBe('buffered');
expect(isDeliveryIdle(ctx)).toBe(true);
});
});
12 changes: 12 additions & 0 deletions packages/core/src/delivery-barrier-dispenser.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -117,6 +117,18 @@ const resumeAtA = new Date('2026-07-27T12:00:05.000Z');
const resumeAtB = new Date('2026-07-27T12:00:06.000Z');

describe('barrier safety-net dispenser', () => {
it('rejects a second barrier owner for the same event index', () => {
const ctx = setupWorkflowContext([]);
const barrier = registerDeliveryBarrier(ctx, 0, 'hook', { armed: false });

expect(() => registerDeliveryBarrier(ctx, 0, 'step')).toThrowError(
'Delivery barrier already registered at event index 0'
);

barrier.markDelivered();
expect(isDeliveryIdle(ctx)).toBe(true);
});

it('suspends only after deliveries parked behind an unclaimed payload have run', async () => {
const ops: Promise<any>[] = [];
const payload = await dehydrateStepReturnValue(
Expand Down
109 changes: 72 additions & 37 deletions packages/core/src/private.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -2,6 +2,7 @@
* Utils used by the bundler when transforming code
*/

import { WorkflowRuntimeError } from '@workflow/errors';
import { withResolvers } from '@workflow/utils';
import type { CryptoKey } from './encryption.js';
import type { EventsConsumer } from './events-consumer.js';
Expand DownExpand Up@@ -215,8 +216,8 @@ export type DeliveryKind = 'hook' | 'wait' | 'step';

interface DeliveryBarrierEntry {
kind: DeliveryKind;
/** Resolves once this delivery has resolved to the workflow. */
delivered: Promise<void>;
/** Resolves once this delivery is handed to the workflow or retired. */
released: Promise<void>;
/**
* Whether this delivery is committed to reaching the workflow without any
* further action by workflow code. True for wait completions and step
Expand All@@ -229,11 +230,15 @@ interface DeliveryBarrierEntry {
* once a consumer takes the payload.
*/
armed: boolean;
/** Whether this entry has been removed and its `released` promise settled. */
retired: boolean;
/**
* Retire this entry: resolve `delivered` and remove it from the registry,
* exactly as `markDelivered` would. Called only by the context's safety-net
* dispenser ({@link ensureBarrierSafetyNet}), and only on the lowest-index
* entry at delivery idle. Idempotent.
* Retire this entry: resolve `released` and remove it from the registry,
* without marking the handle delivered to the workflow. A safety-retired
* buffered payload may therefore install a fresh entry if it is claimed by
* a retained VM later. Called only by the context's safety-net dispenser
* ({@link ensureBarrierSafetyNet}), and only on the lowest-index entry at
* delivery idle. Idempotent.
*/
retire: () => void;
}
Expand DownExpand Up@@ -396,7 +401,7 @@ export async function awaitEarlierDeliveries(
) {
continue;
}
earlier.push(entry.delivered);
earlier.push(entry.released);
}
if (earlier.length > 0) {
await Promise.all(earlier);
Expand All@@ -422,7 +427,7 @@ export async function awaitEarlierDeliveries(
export interface DeliveryBarrier {
/**
* Mark this delivery as delivered to the workflow. Resolves its
* `delivered` promise so any later-in-log delivery gated on it (via
* `released` promise so any later-in-log delivery gated on it (via
* {@link awaitEarlierDeliveries}) may proceed, and removes it from the
* registry. Idempotent.
*/
Expand DownExpand Up@@ -468,41 +473,71 @@ export function registerDeliveryBarrier(
return { markDelivered: () => {}, arm: () => {} };
}

let done = false;
const { promise, resolve } = withResolvers<void>();

const finish = () => {
if (done) {
return;
}
done = true;
if (barriers.get(eventIndex) === entry) {
barriers.delete(eventIndex);
const install = (armed: boolean): DeliveryBarrierEntry => {
if (barriers.has(eventIndex)) {
throw new WorkflowRuntimeError(
`Delivery barrier already registered at event index ${eventIndex}`
);
}
resolve();
const { promise, resolve } = withResolvers<void>();
const entry: DeliveryBarrierEntry = {
kind,
released: promise,
armed,
retired: false,
retire: () => {
if (entry.retired) {
return;
}
entry.retired = true;
if (barriers.get(eventIndex) === entry) {
barriers.delete(eventIndex);
}
resolve();
},
};
barriers.set(eventIndex, entry);

// Safety net: if this delivery is never delivered to the workflow (its
// branch was not taken / the run is suspending, or a buffered hook payload
// is only claimed after a later delivery the workflow is still waiting
// on), it is retired at idle so a later delivery gated on it cannot
// deadlock and the registry cannot leak an entry per abandoned delivery.
// Retirement goes through the context's single ordered dispenser rather
// than a per-barrier idle poll. See {@link ensureBarrierSafetyNet} for why
// the ORDER of these retirements is load-bearing.
ensureBarrierSafetyNet(ctx);
return entry;
};

const entry: DeliveryBarrierEntry = {
kind,
delivered: promise,
armed: options.armed ?? true,
retire: finish,
};
barriers.set(eventIndex, entry);

// Safety net: if this delivery is never delivered to the workflow (its
// branch was not taken / the run is suspending, or a buffered hook payload
// is only claimed after a later delivery the workflow is still waiting on),
// it is retired at idle so a later delivery gated on it cannot deadlock and
// the registry cannot leak an entry per abandoned delivery. Retirement goes
// through the context's single ordered dispenser rather than a per-barrier
// idle poll — see {@link ensureBarrierSafetyNet} for why the ORDER of these
// retirements is load-bearing.
ensureBarrierSafetyNet(ctx);
let entry = install(options.armed ?? true);
let deliveredToWorkflow = false;

return {
markDelivered: finish,
markDelivered: () => {
if (deliveredToWorkflow) {
return;
}
deliveredToWorkflow = true;
entry.retire();
},
arm: () => {
if (deliveredToWorkflow) {
return;
}
// The idle safety net may retire an unclaimed buffered hook payload
// while a retained VM keeps its `claim()` closure alive. If workflow
// code later claims that payload, replace the settled entry so delivery
// remains non-idle until the claim reaches the workflow.
if (entry.retired) {
entry = install(true);
return;
}
if (barriers.get(eventIndex) !== entry) {
throw new WorkflowRuntimeError(
`Delivery barrier lost ownership of event index ${eventIndex}`
);
}
entry.armed = true;
},
};
Expand Down
Loading