Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
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
7 changes: 7 additions & 0 deletions .changeset/sparkly-deer-smoke.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
---
"@workflow/world": minor
"@workflow/core": minor
"@workflow/world-postgres": minor
---

Add optional executor invocation delivery for hooks, with a Postgres input/result mailbox and run-scoped Graphile executor queues.
23 changes: 23 additions & 0 deletions docs/content/worlds/v5/postgres.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -80,6 +80,29 @@ The migration is idempotent and can safely be run as a post-deployment lifecycle

## Starting the World

### Experimental invocation delivery

Set `WORKFLOW_POSTGRES_INVOKE=1` (default: off), or pass `enableInvoke: true`
to the Postgres World's factory, to enable synchronous hook-input delivery.
Apply the database migrations and upgrade the participating workers first.

The World advertises `capabilities.invoke`. Hook callers send inputs to the
run's executor, which inspects them, writes accepted hook events, and then
responds. Graphile's `workflow_flows_executor` task uses a per-run named queue
(`workflow_flows:<runId>:executor`); only one executor job for a run is active at
a time, while other runs and the separate step jobs remain concurrent. A custom
job prefix replaces `workflow_` in these names. Each invoke schedules a wake,
even if a live executor already handled its input.

Inputs/results are stored in a Postgres table and polled every 50ms. The default
invoke response timeout is 30 seconds, and the encoded input/result limit is
1 MiB each. A timeout does not undo processing. Event persistence precedes
response persistence, but the two are not atomic: a crash in between may cause
duplicate hook events. Result cleanup, stale-executor fencing and routing across
incompatible code versions are not implemented by this experimental mode.

### Server startup

To subscribe to the graphile-worker queue, your workflow app needs to start the world on server start. Here are examples for a few frameworks:

<Callout type="warn">
Expand Down
9 changes: 9 additions & 0 deletions packages/core/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -6,3 +6,12 @@ Queued step messages carry immutable run identity in `runContext`, so step
execution skips the initial `runs.get` and fetches the run row only when
continuing into workflow replay. Messages from older producers without
`runContext` retain the initial fetch.

When a World advertises `capabilities.invoke`, `resumeHook()` sends the serialized
input through `world.invoke()` and awaits the executor's decision. The executor
services the optional queue-handler invocation feed alongside execution, validates
the hook, writes its event, and then responds. The event and response writes are
sequential, not atomic. Other Worlds retain the existing producer-write/wake path.
Input admission remains active during inline steps; workflow code advances at its
existing replay boundaries. A retained Node VM can continue over newly committed
inputs within the executor's bounded idle window.
44 changes: 41 additions & 3 deletions packages/core/src/runtime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -91,6 +91,7 @@ import {
stepDispatchIdempotencyKey,
withHealthCheck,
} from './runtime/helpers.js';
import { withInvocationFeed } from './runtime/invocations.js';
import {
handleReplayBudgetExhausted,
ReplayBudget,
Expand Down Expand Up @@ -579,6 +580,7 @@ function getRetentionDecision({
suspension,
serializationBlockerCount,
hookContinuation = false,
invocationContinuation = false,
}: {
suspension: WorkflowSuspension;
serializationBlockerCount: number;
Expand All @@ -589,6 +591,7 @@ function getRetentionDecision({
* write of its own. See the policy above.
*/
hookContinuation?: boolean;
invocationContinuation?: boolean;
}): RetentionDecision {
if (!isVmRetentionEnabled()) {
return { retain: false, reason: 'disabled' };
Expand All @@ -602,7 +605,11 @@ function getRetentionDecision({
if (hookContinuation) {
return { retain: true };
}
if (suspension.stepCount === 0 && suspension.attributeCount === 0) {
if (
!invocationContinuation &&
suspension.stepCount === 0 &&
suspension.attributeCount === 0
) {
return { retain: false, reason: 'no_replay_driver' };
}
return { retain: true };
Expand Down Expand Up @@ -713,7 +720,7 @@ export function workflowEntrypoint(
const handler = (worldHandlers: WorldHandlers) =>
worldHandlers.createQueueHandler(
workflowPrefix,
async (message_, metadata) => {
withInvocationFeed(async (message_, metadata, invocations) => {
// T2 of the hook-resume TTR window (see runtime/resume-latency.ts):
// the instant this consumer began, before message parsing. Only used
// when the message turns out to carry resume timing; taking it
Expand Down Expand Up @@ -1960,6 +1967,24 @@ export function workflowEntrypoint(
return { timeoutSeconds: stepResult.timeoutSeconds };
}

// Invoke-capable worlds serialize orchestrator deliveries,
// not step bodies. Return the replay to that executor lane.
if (
world.capabilities?.invoke &&
stepResult.type !== 'gone'
) {
await queueMessage(
world,
getWorkflowQueueName(workflowName, namespace),
{
runId,
traceCarrier: await nextTraceCarrier(),
requestedAt: new Date(),
}
);
return;
}

// If step had pending ops (stream writes), break and let
// waitUntil flush them, so can't continue inline.
if (
Expand Down Expand Up @@ -2892,6 +2917,7 @@ export function workflowEntrypoint(

// Main replay loop
while (true) {
const invocationRevision = invocations?.revision ?? 0;
loopIteration++;

// Replay-budget check: bail out (retry or fail) if
Expand Down Expand Up @@ -3649,6 +3675,7 @@ export function workflowEntrypoint(
hookContinuation:
suspensionResult.hasAwaitedHookCreation ||
suspensionResult.hasHookConflict,
invocationContinuation: invocations !== undefined,
})
: undefined;
if (retentionDecision?.retain === false) {
Expand Down Expand Up @@ -4329,6 +4356,17 @@ export function workflowEntrypoint(
// committed hook_created.
return await reinvoke(0);
}
if (
invocations &&
Date.now() - invocationStartTime <
noInlineReplayAfterMs &&
(await invocations.waitForActivity(
invocationRevision
))
) {
eventLog = nextEventLogLoad(eventLog);
continue;
}
return;
}

Expand Down Expand Up @@ -5214,7 +5252,7 @@ export function workflowEntrypoint(
}
);
});
}
})
);

let cachedHandler: ((req: Request) => Promise<Response>) | undefined;
Expand Down
136 changes: 136 additions & 0 deletions packages/core/src/runtime/invocations.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,136 @@
import { HookNotFoundError } from '@workflow/errors';
import type { Invocation, World } from '@workflow/world';
import { describe, expect, it, vi } from 'vitest';
import { handleInvocation, InvocationPump } from './invocations.js';

function fixture() {
const create = vi.fn().mockResolvedValue({});
const getHook = vi
.fn()
.mockResolvedValue({ hookId: 'h', token: 't', runId: 'r', specVersion: 7 });
const world = {
hooks: { get: getHook },
runs: { get: vi.fn().mockResolvedValue({ status: 'running' }) },
events: { create },
} as unknown as World;
const input: Invocation = {
id: 'request',
payload: {
type: 'hook_resume',
version: 1,
hookId: 'h',
token: 't',
payload: new Uint8Array([1, 2]),
},
respond: vi.fn().mockResolvedValue(undefined),
};
return { world, input, create, getHook };
}

describe('runner invocation decisions', () => {
it('does not accept until persistence finishes', async () => {
const { world, input, create } = fixture();
let commit!: () => void;
create.mockImplementation(
() =>
new Promise<void>((resolve) => {
commit = resolve;
})
);
const decision = handleInvocation(world, 'r', input);
await vi.waitFor(() => expect(create).toHaveBeenCalledOnce());
expect(input.respond).not.toHaveBeenCalled();
commit();
await expect(decision).resolves.toBe(true);
expect(input.respond).toHaveBeenCalledWith({ status: 'accepted' });
});

it('rejects mismatched hook/run identity before writing', async () => {
const { world, input, create, getHook } = fixture();
getHook.mockResolvedValue({ runId: 'another-run', token: 't' });
await expect(handleInvocation(world, 'r', input)).resolves.toBe(false);
expect(create).not.toHaveBeenCalled();
expect(input.respond).toHaveBeenCalledWith({
status: 'rejected',
code: 'HOOK_NOT_FOUND',
});
});

it('distinguishes a storage lifecycle rejection from infrastructure failure', async () => {
const { world, input, create } = fixture();
create.mockRejectedValueOnce(new HookNotFoundError('h'));
await handleInvocation(world, 'r', input);
expect(input.respond).toHaveBeenCalledWith({
status: 'rejected',
code: 'HOOK_NOT_FOUND',
});
vi.mocked(input.respond).mockClear();
const error = new Error('database unavailable');
create.mockRejectedValueOnce(error);
await expect(handleInvocation(world, 'r', input)).rejects.toBe(error);
expect(input.respond).not.toHaveBeenCalled();
});

it('does not turn a failed response write into a rejection after the event committed', async () => {
const { world, input, create } = fixture();
const error = new Error('response storage failed');
vi.mocked(input.respond).mockRejectedValue(error);
await expect(handleInvocation(world, 'r', input)).rejects.toBe(error);
expect(create).toHaveBeenCalledOnce();
expect(input.respond).toHaveBeenCalledExactlyOnceWith({
status: 'accepted',
});
});

it('rejects unknown input protocols without journaling', async () => {
const { world, input, create } = fixture();
input.payload = { type: 'unknown' };
await handleInvocation(world, 'r', input);
expect(create).not.toHaveBeenCalled();
expect(input.respond).toHaveBeenCalledWith({
status: 'rejected',
code: 'INVALID_INPUT',
});
});

it('services two inputs sequentially and waits for in-flight persistence on close', async () => {
const { world, input, create } = fixture();
const second = { ...input, id: 'second', respond: vi.fn() };
const inputs = [input, second];
let stopped = false;
const source: AsyncIterableIterator<Invocation> = {
[Symbol.asyncIterator]() {
return this;
},
next: async () =>
inputs.length && !stopped
? { done: false, value: inputs.shift()! }
: { done: true, value: undefined },
return: async () => {
stopped = true;
return { done: true, value: undefined };
},
};
let finishSecond!: () => void;
create.mockResolvedValueOnce({}).mockImplementationOnce(
() =>
new Promise<void>((resolve) => {
finishSecond = resolve;
})
);
const pump = new InvocationPump(source, (item) =>
handleInvocation(world, 'r', item)
);
await vi.waitFor(() => expect(create).toHaveBeenCalledTimes(2));
expect(pump.revision).toBe(1);
let closed = false;
const closing = pump.close().then(() => {
closed = true;
});
await Promise.resolve();
expect(closed).toBe(false);
finishSecond();
await closing;
expect(second.respond).toHaveBeenCalledWith({ status: 'accepted' });
});
});
Loading
Loading