Skip to content

Commit 97dccc9

Browse files
shalabhcvercel[bot]VaguelySerious
authored
[world] Add optional invoke method (hook payloads only for now) (#4168)
* feat(world-postgres): add optional executor invocation for hooks Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * perf(world-postgres): notify invocation input and result waiters Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * fix(world-postgres): verify serialized executor deliveries Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * refactor(world): return invocation results and harden Postgres delivery Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * chore: reconcile main before signed merge Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * fix(world-postgres): replay drained inputs and return invocation errors Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * Fix: Transient/infra errors thrown during hook-resume invocation are captured and stored as a permanent, non-retryable outcome instead of being retried, which can leave a `hook_received` write uncommitted and suspend the workflow forever. This commit fixes the issue reported at packages/world-postgres/src/queue.ts:322 ## The bug Commit `f944ec9` reworked the executor `invoke` path in `packages/world-postgres/src/queue.ts`: ```ts const outcome = await captureInvocationOutcome(() => handler(message, metadata)); await invocations.respondOutcome(input.runId, input.requestId, outcome); return; // delivery acked (HTTP 200), no retry ``` `captureInvocationOutcome` (`packages/errors/src/invocation.ts`) caught **every** error unconditionally and returned `{ ok: false, error }`. `respondOutcome` then persisted it (`result_version = 1`, `responded_at = now()`) and the handler returned normally, so Graphile acked the job — **no redelivery**. For an `invoke` message the wrapped handler routes only to `handleInvocation` (`packages/core/src/runtime/invocations.ts`), which performs pure world I/O: `world.hooks.get`, `world.runs.get`, and the durable `world.events.create` that commits the `hook_received` event resuming the workflow. `f944ec9` also changed `handleInvocation` to **throw** on every failure (it previously returned `{status:'rejected'}` for hook-gone). ### Concrete failure mode 1. `resumeHook` invokes with a stable `requestId`. 2. The executor runs `handleInvocation`; `world.events.create` fails with a **transient** error (DB connection reset, deadlock, `40001` serialization failure — surfaced as a raw `DrizzleQueryError`/pg error, not a `WorkflowError`). 3. `captureInvocationOutcome` catches it → `{ ok:false, error }`. 4. `respondOutcome` stores the error permanently (`responded_at` set) and commits. 5. Handler returns → HTTP 200 → Graphile acks → **no retry**. 6. The `hook_received` event was never committed, so the run stays suspended. 7. `world.invoke()`'s waiter reads the responded row and (via `unwrapInvocationOutcome`) throws the rehydrated error to the caller immediately. **Before** `f944ec9`, a thrown error propagated out of the queue handler into `createQueueHandler`'s retry logic (Graphile redelivery), so a transient blip was retried transparently until the resume committed. The new code converts that into a permanent stored failure. The intended contract (per `packages/world-postgres/test/invoke.test.ts` — *"returns persisted $name outcomes instead of timing out"*, which only exercises **known** `WorkflowWorldError`/`HookNotFoundError`/`WorkflowRunNotFoundError`/`RunExpiredError`/`EntityConflictError`) was to persist **terminal/deterministic** business errors so callers get the right error class instead of a 30s timeout — not to persist transient/unknown failures. ## The fix Distinguish terminal errors (capture + store) from transient/unknown errors (re-throw so the delivery layer retries): * Added `isTerminalInvocationError` in `packages/errors/src/invocation.ts`. An error is terminal only when it is a class the `@workflow/errors` package owns (checked by `name` against the module registry, matching the existing `deserializeWorkflowError` approach) **and** it does not look transient — i.e. it has no `retryAfter` and no `status` of `>= 500` / `408` / `425` / `429`. Unknown/infra failures (raw `Error`, `DrizzleQueryError`, thrown non-Errors) and transient Workflow errors are therefore **not** terminal. * Gave `captureInvocationOutcome` an optional `shouldCapture` predicate (defaulting to capture-all, preserving its generic transport contract and existing unit tests). When it returns `false`, the error is re-thrown. * Both `invoke`-path call sites in `packages/world-postgres/src/queue.ts` now pass `isTerminalInvocationError`, so transient failures propagate to `createQueueHandler`'s Graphile retry (as before `f944ec9`) while terminal errors are still persisted (keeping the `invoke.test.ts` `deliveries === 1` behavior). The known terminal cases thrown by `handleInvocation` remain captured: `INVALID_INPUT` (400), `HookNotFoundError`, `INVOCATION_DATA_EXPIRED` (410), `WorkflowRunNotFoundError`, `RunExpiredError`, and the deterministic `EntityConflictError` from correlated-event dedup. Added unit tests covering both the terminal and transient classifications. Type-checking passes for `@workflow/errors` and `@workflow/world-postgres`, and the `invocation.test.ts` suite passes. Co-authored-by: Vercel <vercel[bot]@users.noreply.github.com> Co-authored-by: VaguelySerious <mittgfu@gmail.com> * chore: reconcile main and invocation retry coverage Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * style(errors): format invocation error fixture Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * docs: consolidate invocation changeset and trim Vercel notes Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * docs(world): define invoke single-runner guarantee Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * docs: apply reviewed invoke documentation wording Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> --------- Co-authored-by: vercel[bot] <35613825+vercel[bot]@users.noreply.github.com> Co-authored-by: Vercel <vercel[bot]@users.noreply.github.com> Co-authored-by: VaguelySerious <mittgfu@gmail.com>
1 parent f1f5b7d commit 97dccc9

45 files changed

Lines changed: 4558 additions & 45 deletions

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

‎.changeset/sparkly-deer-smoke.md‎

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,10 @@
1+
---
2+
"@workflow/world": minor
3+
"@workflow/core": minor
4+
"@workflow/errors": minor
5+
"@workflow/world-postgres": minor
6+
"@workflow/world-local": patch
7+
"@workflow/world-vercel": patch
8+
---
9+
10+
Add an optional `invoke` method and capability to the World interface. It routes a payload to the runner handling the specified `runId` and returns a promise for its response. In world-postgres, this uses a regular queue roundtrip with a run-scoped queue. The runtime uses `invoke` when available, initially to resume hooks.

‎docs/content/worlds/v5/building-a-world.mdx‎

Lines changed: 40 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -239,13 +239,52 @@ interface Queue {
239239
messages: readonly { message: QueuePayload; opts?: QueueOptions }[]
240240
): Promise<QueueBatchResult[]>;
241241

242+
// Optional request/response delivery, enabled with capabilities.invoke.
243+
invoke?(runId: string, payload: unknown, options?: InvokeOptions): Promise<unknown>;
244+
242245
createQueueHandler(
243246
queueNamePrefix: QueuePrefix,
244-
handler: (message: unknown, meta: { attempt: number; queueName: ValidQueueName; messageId: MessageId }) => Promise<void | { timeoutSeconds: number }>
247+
handler: (message: unknown, meta: { attempt: number; queueName: ValidQueueName; messageId: MessageId }) => Promise<unknown>
245248
): (req: Request) => Promise<Response>;
246249
}
247250
```
248251

252+
### Invocation delivery
253+
254+
Implement the optional `invoke` method to send an input to a workflow runner and
255+
return the runner's response. Enable `capabilities.invoke` only when your World
256+
supports this operation.
257+
258+
Your World must guarantee at most one active workflow runner per `runId` across
259+
all worker processes. Different runs may execute concurrently. A replacement
260+
runner can take over after the previous runner stops.
261+
262+
Route inputs to the active runner, or start or resume the runner when none is
263+
active. Call the existing `createQueueHandler` callback with
264+
`{ runId, invoke: true, requestId, input }`. The callback validates the input,
265+
waits for required event writes, and returns a result. Deliver that result to the
266+
original caller. Support input processing while the runner awaits step work,
267+
without starting a second runner for the same run.
268+
269+
`requestId` identifies the logical input across retries and is separate from a
270+
queue delivery ID. `InvokeOptions` accepts an `idempotencyKey` and a `timeoutMs`
271+
response-wait limit.
272+
273+
Treat invocation return values as response data, including any `timeoutSeconds`
274+
property. Ordinary workflow wake results use `timeoutSeconds` to schedule another
275+
execution.
276+
277+
You can transport results with `InvocationOutcome`, using `{ ok: true, value }`
278+
for success or `{ ok: false, error: SerializedWorkflowError }` for an error. The
279+
helpers in `@workflow/errors/invocation` capture handler results and restore
280+
errors with their known Workflow classes and diagnostic fields. Return the value
281+
or throw the restored error from `invoke()`.
282+
283+
Keep failures to store or deliver a response distinct from handler errors. A
284+
transport failure leaves processing unknown, and a handler error may follow
285+
committed event writes. If your World advertises `hookResumeDedup`, enforce that
286+
capability when retrying a hook event whose response was lost.
287+
249288
### Batched publishing
250289

251290
`queueBatch` is optional. Implement it when your transport can accept several messages in one round trip, and the runtime will use it to dispatch a wide `Promise.all` fan-out. A fan-out otherwise costs one round trip per branch, paid before the first branch's step body runs, so it lands directly on time-to-first-step.

‎docs/content/worlds/v5/postgres.mdx‎

Lines changed: 54 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -80,6 +80,60 @@ The migration is idempotent and can safely be run as a post-deployment lifecycle
8080

8181
## Starting the World
8282

83+
### Experimental invocation delivery
84+
85+
Set `WORKFLOW_POSTGRES_INVOKE=1`, or pass `enableInvoke: true` to the Postgres
86+
World's factory, to send hook inputs to the run's executor and wait for its
87+
response. Invocation is off by default. Apply the migrations and upgrade
88+
producers and workers together before enabling it.
89+
90+
The executor validates the hook input and writes its event before responding.
91+
Workflow user code can consume the event later. The World stores inputs and
92+
responses in `workflow.workflow_invocations`. Inserting an input and enqueueing
93+
a workflow execution job share one transaction. Retries also enqueue an execution
94+
job, including when a retained response already exists, so committed events can
95+
be replayed after a runner failure.
96+
97+
With the default job prefix, Graphile uses `workflow_flows_executor` for workflow
98+
execution and `workflow_flows` for step jobs and health checks. Execution jobs
99+
use a named queue per run, `workflow_flows:<runId>:executor`. Graphile permits
100+
one active job on each named queue across workers. Other runs, steps, and health
101+
checks can run concurrently. `queueConcurrency` limits total active jobs per
102+
worker process and defaults to 50. A custom job prefix replaces `workflow_` in
103+
these names.
104+
105+
The HTTP receiver checks each execution request against its active Graphile job
106+
before processing pending inputs. Updated workers and receivers move legacy
107+
orchestration requests to the run's named queue. Before acknowledging an execution
108+
job, the World finishes any input already being processed and replays newly
109+
committed events. These entry checks do not stop an old handler from writing
110+
after its job is reclaimed.
111+
112+
The caller receives a stored return value or a restored Workflow error. Terminal
113+
errors are retained as outcomes; transient or unrecognized failures leave inputs
114+
pending for Graphile to retry. Migration 0022 preserves the meaning of responses
115+
stored by earlier previews. Failure to store or read a response leaves the
116+
processing outcome unknown. Durable resume identities deduplicate the hook event
117+
when response storage must be retried.
118+
119+
A shared `LISTEN/NOTIFY` connection signals changes to pending inputs and
120+
responses. Notifications contain identifiers, and readers fetch data from the
121+
table. Use a session-capable connection for `LISTEN`. Readers fall back to polling
122+
at 1s intervals if notifications fail, and connection retries use a 1s backoff.
123+
Notification failure and recovery are logged once per transition without
124+
payloads or connection details.
125+
126+
The response-wait timeout defaults to 30s and can be set with `timeoutMs`. Encoded
127+
inputs and outcomes are limited to 1 MiB each. A timeout does not cancel
128+
processing. `$retention: 0` purges invocation data and prevents late writes from
129+
restoring it. Callers whose result has expired receive a 410 error with code
130+
`INVOCATION_DATA_EXPIRED`, even if the hook event committed earlier.
131+
132+
Automatic cleanup of other results, stopping writes from reclaimed handlers, and
133+
routing requests across incompatible worker versions remain unimplemented.
134+
135+
### Server startup
136+
83137
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:
84138

85139
<Callout type="warn">

‎packages/core/README.md‎

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -8,3 +8,19 @@ Queued step messages carry immutable run identity in `runContext`, so step
88
execution skips the initial `runs.get` and fetches the run row only when
99
continuing into workflow replay. Messages from older producers without
1010
`runContext` retain the initial fetch.
11+
12+
When a World advertises `capabilities.invoke`, `resumeHook()` sends serialized
13+
hook inputs through `world.invoke()` and waits for the executor's decision. The
14+
World calls the existing handler with `invoke: true`, `requestId`, and the hook
15+
input. Core validates the input, waits for its event write, and returns a
16+
decision. Core shares a run's activity state between workflow execution and these
17+
input calls.
18+
19+
The event write completes before the response is stored. On Worlds with
20+
`hookResumeDedup`, a stable request identity makes retries reuse the same hook
21+
event, including after hook disposal or run completion. Core can process inputs
22+
while inline steps wait. Workflow code observes committed inputs at replay
23+
boundaries, reusing a retained Node virtual machine when available.
24+
25+
Errors returned by the executor propagate through `world.invoke()` to the caller
26+
of `resumeHook()`.

‎packages/core/src/runtime.ts‎

Lines changed: 41 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -91,6 +91,7 @@ import {
9191
stepDispatchIdempotencyKey,
9292
withHealthCheck,
9393
} from './runtime/helpers.js';
94+
import { withRunInputs } from './runtime/invocations.js';
9495
import {
9596
handleReplayBudgetExhausted,
9697
ReplayBudget,
@@ -120,7 +121,7 @@ import {
120121
} from './runtime/suspension-handler.js';
121122
import { useQuickJSVm } from './runtime/vm-mode.js';
122123
import { getWaitContinuationDispatch } from './runtime/wait-continuation.js';
123-
import { getWorld, type WorldHandlers } from './runtime/world.js';
124+
import { getWorld } from './runtime/world.js';
124125
import { dehydrateRunError } from './serialization.js';
125126
import { remapErrorStack } from './source-map.js';
126127
import * as Attribute from './telemetry/semantic-conventions.js';
@@ -584,6 +585,7 @@ function getRetentionDecision({
584585
suspension,
585586
serializationBlockerCount,
586587
hookContinuation = false,
588+
invocationContinuation = false,
587589
}: {
588590
suspension: WorkflowSuspension;
589591
serializationBlockerCount: number;
@@ -594,6 +596,7 @@ function getRetentionDecision({
594596
* write of its own. See the policy above.
595597
*/
596598
hookContinuation?: boolean;
599+
invocationContinuation?: boolean;
597600
}): RetentionDecision {
598601
if (!isVmRetentionEnabled()) {
599602
return { retain: false, reason: 'disabled' };
@@ -607,7 +610,11 @@ function getRetentionDecision({
607610
if (hookContinuation) {
608611
return { retain: true };
609612
}
610-
if (suspension.stepCount === 0 && suspension.attributeCount === 0) {
613+
if (
614+
!invocationContinuation &&
615+
suspension.stepCount === 0 &&
616+
suspension.attributeCount === 0
617+
) {
611618
return { retain: false, reason: 'no_replay_driver' };
612619
}
613620
return { retain: true };
@@ -715,10 +722,10 @@ export function workflowEntrypoint(
715722
const namespace = resolveQueueNamespace(options?.namespace);
716723
const workflowPrefix = getQueueTopicPrefix('workflow', namespace);
717724

718-
const handler = (worldHandlers: WorldHandlers) =>
725+
const handler = (worldHandlers: World) =>
719726
worldHandlers.createQueueHandler(
720727
workflowPrefix,
721-
async (message_, metadata) => {
728+
withRunInputs(worldHandlers)(async (message_, metadata, activity) => {
722729
// T2 of the hook-resume TTR window (see runtime/resume-latency.ts):
723730
// the instant this consumer began, before message parsing. Only used
724731
// when the message turns out to carry resume timing; taking it
@@ -1965,6 +1972,24 @@ export function workflowEntrypoint(
19651972
return { timeoutSeconds: stepResult.timeoutSeconds };
19661973
}
19671974

1975+
// Invoke-capable worlds serialize orchestrator deliveries,
1976+
// not step bodies. Return the replay to that executor lane.
1977+
if (
1978+
world.capabilities?.invoke &&
1979+
stepResult.type !== 'gone'
1980+
) {
1981+
await queueMessage(
1982+
world,
1983+
getWorkflowQueueName(workflowName, namespace),
1984+
{
1985+
runId,
1986+
traceCarrier: await nextTraceCarrier(),
1987+
requestedAt: new Date(),
1988+
}
1989+
);
1990+
return;
1991+
}
1992+
19681993
// If step had pending ops (stream writes), break and let
19691994
// waitUntil flush them, so can't continue inline.
19701995
if (
@@ -2897,6 +2922,7 @@ export function workflowEntrypoint(
28972922

28982923
// Main replay loop
28992924
while (true) {
2925+
const invocationRevision = activity?.revision ?? 0;
29002926
loopIteration++;
29012927

29022928
// Replay-budget check: bail out (retry or fail) if
@@ -3654,6 +3680,7 @@ export function workflowEntrypoint(
36543680
hookContinuation:
36553681
suspensionResult.hasAwaitedHookCreation ||
36563682
suspensionResult.hasHookConflict,
3683+
invocationContinuation: activity !== undefined,
36573684
})
36583685
: undefined;
36593686
if (retentionDecision?.retain === false) {
@@ -4334,6 +4361,15 @@ export function workflowEntrypoint(
43344361
// committed hook_created.
43354362
return await reinvoke(0);
43364363
}
4364+
if (
4365+
activity &&
4366+
Date.now() - invocationStartTime <
4367+
noInlineReplayAfterMs &&
4368+
(await activity.waitForActivity(invocationRevision))
4369+
) {
4370+
eventLog = nextEventLogLoad(eventLog);
4371+
continue;
4372+
}
43374373
return;
43384374
}
43394375

@@ -5229,7 +5265,7 @@ export function workflowEntrypoint(
52295265
}
52305266
);
52315267
});
5232-
}
5268+
})
52335269
);
52345270

52355271
let cachedHandler: ((req: Request) => Promise<Response>) | undefined;

0 commit comments

Comments
 (0)