---
title: Inline step message ownership
description: Why inline step_started now records its owning queue message ID, why ownership is bounded by a lease, and why the alternatives, from queue serialization to heartbeats, were rejected. Fixes duplicate inline step execution (issue
---

# Inline step message ownership



# Inline step message ownership

> Fixes [#2780](https://github.com/vercel/workflow/issues/2780) (duplicate
> inline step execution when a hook or wait wakes a run mid-step). This page
> is the design defense for the change: what it does, why each choice was
> made, and why the alternatives were rejected. Kill switch:
> `WORKFLOW_INLINE_OWNERSHIP=0`; lease tuning:
> `WORKFLOW_INLINE_OWNERSHIP_LEASE_SECONDS` (see
> [Runtime tuning](/docs/configuration/runtime-tuning)).

## What this change does

The lazy `step_started` that creates an inline step now also records the queue message ID
of the invocation running its body (`eventData.ownerMessageId`). While that ownership is
active (until the step's first `step_retrying` or terminal event), a replay triggered by
anything other than the owning message does **not** requeue the step. It instead ensures a
*delayed backstop wake* exists, timed to the remainder of an ownership lease. Only a
handler processing the owning message (the original invocation, or the queue's redelivery
of it after a crash) may re-execute the step before that lease expires.

## The bug being fixed

An inline step has no queue message: that is the point of inline execution (it saves the
dispatch round-trip). But the runtime's crash-recovery rule queued every created,
non-terminal step **unconditionally** on each replay, relying on the queue's
`idempotencyKey = correlationId` to dedupe repeats. For inline steps there is no prior
message to dedupe against, so any wake that replayed the run mid-step (`hook_received`,
an elapsed wait continuation, a cancellation) enqueued a *first* message for that
correlation ID. Its consumer sent a bare `step_started` on the `running` step (which is
allowed: retries legitimately re-start non-terminal steps) and ran the body a second time,
concurrently with the original. One `step_completed` won the terminal write; every side
effect had already happened twice.

Note what is *not* broken: the unconditional requeue is correct for eager steps and is
itself the crash-recovery mechanism for a handler that wrote `step_created` and died
before enqueueing. Any fix must suppress the requeue *only* while a live invocation is
demonstrably running the body, which turns the bug into a liveness problem.

## The core problem is liveness, and why we did not build a liveness mechanism

The missing primitive is the ability to distinguish "this attempt is in flight in a live
invocation" from "this attempt died with its process." Stateless compute offers no death
signal: a crashed function emits nothing, and the event log looks identical either way
(`step_started`, no terminal event).

The classic answers are heartbeats or a lock store: the executing invocation periodically
renews a claim, and recovery waits for the claim to go stale. We rejected building one:

* **It adds World API surface.** Every backend (Vercel, local filesystem, community
  Postgres/Turso/Redis worlds) would have to implement a renewable-claim store with
  expiry semantics. The event log is the World contract's one source of truth; a second,
  mutable liveness store beside it is a large contract change for a race fix.
* **It adds steady-state write load.** Heartbeats cost a write per interval per in-flight
  step, paid by every healthy run to detect the rare crashed one.
* **It does not remove the hard part.** A heartbeat still needs an expiry to survive a
  crashed heartbeater: that expiry *is* a lease. Any liveness design degenerates to
  "claim + bounded staleness"; the machinery around it is overhead.

Instead, the design reuses a liveness signal the system already has: **the queue's
delivery state**. An un-acked queue message is precisely "work that a live invocation may
be processing, which will be redelivered if the processor died." Stamping the owning
message ID on the step makes the queue's own at-least-once machinery serve as the claim,
the crash detector, and the recovery driver, with zero new World surface and zero
steady-state writes beyond one field on an event we already write.

### Why the identity is the queue `messageId`

`createQueueHandler` already delivers `meta.messageId`, and one enqueued message keeps its
ID across redelivery attempts. That stability is exactly the property recovery needs: "a
delivery whose ID matches the stamp" means "the queue redelivered the work that crashed":
permission to re-execute. The requirement is now documented in the
[Queue contract](/docs/api-reference/workflow-runtime/world/queue); a World whose queue
mints fresh IDs per delivery degrades gracefully (the owner check never matches, so
crashed steps recover via the delayed backstop instead of immediately): it never wedges
and never duplicates.

### Why ownership lives in the event log

Ownership state is derived per-replay from the step's events, not held in memory or in a
side store. Every replayer (the owner's redelivery, a hook wake, the backstop) computes
the same answer from the same log, which is the workflow runtime's existing consistency
model. The rules are chosen so the log alone is sufficient:

* **Latest `step_started` wins.** A stamped start (inline execution, or owner recovery)
  sets the owner; an unstamped bare start (a retry attempt driven by a queued step
  message, or an older runtime) clears it. This is why owner recovery must *re-stamp* its
  bare start: an unstamped recovery start would read as "unowned" to a later wake, which
  would immediately requeue the step the owner is re-running, reintroducing the bug on
  the recovery path.
* **`step_retrying` lapses ownership permanently** for the correlation ID. From the first
  retry on, the step is queue-owned: the retry handoff enqueues a real step message, and
  the ordinary `idempotencyKey = correlationId` dedupe works again. Extending ownership
  across retries was rejected deliberately: retry backoffs are delay-dominated (the
  owning invocation would hold compute open doing nothing), each attempt would pay a
  replay to re-derive state, and transferring ownership between attempts adds a state
  machine where the queue-owned path already recovers correctly.
* **Eager steps are untouched.** Their execution is owned by their own step message and
  its idempotency dedupe; stamping the orchestrator's ID on them would claim work the
  orchestrator is not performing.

## Why ownership requires a lease

Ownership cannot be unconditional. Note first what the lease is *not* for: an owner that
crash-loops through the SDK's delivery budget does not wedge the run even without one.
The flow handler fails the run when it receives an over-budget delivery
(`metadata.attempt > maxQueueDeliveries`). The lease exists because that check, and owner
recovery itself, both depend on assumptions the World contract does not actually promise:

* **Owner-message loss the SDK never observes.** The exhaustion check only fires if the
  queue *delivers* the over-budget attempt. At-least-once doesn't guarantee that: a
  queue-side redrive policy can dead-letter the message below the SDK's budget (a
  community SQS world with a small `maxReceiveCount`), retention can expire it, an
  operator can purge it. The run is then still `running`, the step stamped and
  non-terminal, and no message exists. With unbounded ownership every future wake defers
  to a ghost, a permanent wedge; with the lease, the already-armed backstop (or the
  first wake after expiry) recovers the step.
* **Worlds with unstable message IDs: there the lease is the *entire* recovery
  mechanism.** If a queue mints a fresh ID per delivery, the owner check never matches,
  including on the crashed owner's own redelivery. Unbounded ownership would defer
  forever to a stamp no delivery can ever match, while each deferring replay acks its own
  message: the run drains to zero messages while still `running`. The lease is what
  makes the graceful-degradation claim in the Queue contract true.
* **Insurance on the ack invariant.** Correctness leans on "no path acks the owning
  message while an owned step is non-terminal" (see the decision-table invariants). That
  is an audited property of today's code, not of the contract; a future refactor or a
  queue implementation bug could violate it. The lease caps the cost of any such bug at
  one bounded stall instead of a permanent wedge.
* **Failure granularity for poison steps.** Even on the well-behaved exhaustion path,
  lease expiry lets a poison step execute and fail on the background path, a
  *step-level*, `catch`-able failure the workflow can handle. Unbounded ownership funnels
  the same poison into run-level "exceeded max deliveries", which kills the whole run
  uncatchably, and only after the full backed-off delivery budget.

In short: unbounded ownership is a bet that message-ID stability, queue delivery
guarantees, and the ack invariant all hold, forever, on every World. The lease caps the
cost of losing any of those bets at a single bounded stall.

So ownership is honored only for a bounded window: `leaseRemaining = min(lease, max(0,
lastStartedAt + lease − now))`, anchored at the latest `step_started`'s server-assigned
timestamp. Within the window, non-owner replays defer (and arm the backstop); after it,
dispatch falls back to the pre-existing immediate enqueue. The lease is the upper bound on
"how long a dead owner can delay recovery," and equally the lower bound on "how long a
live owner is protected from duplicates."

The upper clamp exists for clock skew: `lastStartedAt` is server-stamped while `now` is
the local clock, so a client running behind the server would otherwise compute a remainder
*longer* than the lease, and above 900s, a `delaySeconds` that SQS-backed queues reject
outright.

### Why a fixed constant, and why 860 seconds

The correct lease is "longer than any invocation can possibly live". Beyond that point
the owner is provably dead on platforms that kill invocations. Ideally we would derive it
from the workflow route's resolved `maxDuration`. **No such signal exists**: builders emit
`maxDuration: 'max'`, which the platform resolves per-plan at deploy time; there is no
environment variable, request-context deadline, or build-time value to read. Deriving the
lease was rejected because there is nothing to derive it from.

860s is justified by a platform rule rather than a measurement: durations above 800s
require explicit per-function numeric configuration, so `'max'` resolves to ≤ 800s for
any builder-emitted workflow route, and 860 dominates it with headroom. The constant is
configurable through an environment variable (`WORKFLOW_INLINE_OWNERSHIP_LEASE_SECONDS`, clamped from 1 to 900, where 900 is the
queue's maximum per-message delay, so a single delayed backstop message always suffices
and no delay chaining is needed). The code comment carries the 30-minute-`maxDuration`
beta caveat so the constant is revisited when the platform ceiling moves.

On worlds with **no** invocation kill bound (world-local's single process, self-hosted
deployments), no constant is a death proof, which is why the in-process single-flight
below is a required layer, not an optimization.

## Why the non-owner action is a delayed backstop wake: not a skip, and not a step message

The naive non-owner behavior is to *skip* the requeue. Rejected: a pure skip makes
lease expiry useless, because nothing is scheduled to *observe* the expiry. If the owner
dies and no external wake happens to arrive later, the run wedges. The escape hatch has to
be folded into the suppression itself.

So the non-owner enqueues a **plain run continuation** (no `stepId`) with `delaySeconds =
leaseRemaining`. When it fires, it replays the run and re-enters the same dispatch
decision table, which handles every state the step can be in by then: terminal → nothing
pending; queue-owned after `step_retrying` → normal keyed dispatch; owner dead with lease
expired → immediate dispatch, preserving *step-level* failure semantics for poison steps
(the step fails and the workflow's `catch` sees it, rather than the run dying on a
delivery-budget backstop); lease refreshed by owner recovery → re-arm for the new
remainder.

Two shapes of this backstop were tried and rejected by hard evidence, and both lessons are
now encoded in `backstopIdempotencyKey`:

1. **The backstop must not be the step's own message.** The first implementation enqueued
   the step message itself (keyed `correlationId`) with the lease delay. But the owner's
   retry handoff enqueues the step under that *same* key with a \~1s backoff, and the
   pending backstop absorbed it, turning a 1-second retry into a full-lease stall. Caught
   by the abort-mid-flight e2e wedging on every world-local lane.
2. **The backstop key must change when ownership is re-stamped.** Queues dedupe an
   idempotency key for the original message's lifetime, *including while a delivery is
   in flight*. With a fixed `${correlationId}:backstop` key, a backstop firing during a
   lease that owner recovery had refreshed could never publish its own replacement (the
   re-arm deduped against the in-flight backstop itself and was dropped); if the
   recovered owner then died with its redelivery budget exhausted, no escape hatch
   remained. Caught in review. The key is therefore scoped to the **ownership epoch**:
   `${correlationId}:backstop:${lastStartedAt}`. Wakes within one epoch still dedupe to a
   single pending backstop; each owner-recovery re-stamp opens a new epoch with a fresh
   key; pending backstops stay bounded by the owning message's redelivery budget.

## The dispatch decision table

For each pending step in the replay's suspension set that is not designated for lazy
inline execution:

| Step state                                                                                                          | Action                                                                                                                                                   |
| ------------------------------------------------------------------------------------------------------------------- | -------------------------------------------------------------------------------------------------------------------------------------------------------- |
| Uncreated (`!hasCreatedEvent`)                                                                                      | Unchanged: lazy-inline candidate or eager create+queue                                                                                                   |
| Created, ownership active, `owner === myMessageId`                                                                  | **Execute in this invocation** (owned recovery, alongside the lazy-inline batch; input hydrates from the step entity; the bare `step_started` re-stamps) |
| Created, ownership active, `owner !== myMessageId`                                                                  | **Ensure backstop wake**: plain run continuation, `delaySeconds = leaseRemaining`, epoch-scoped idempotency key                                          |
| Created, ownership lapsed (`step_retrying` seen), never owned (eager / old events), lease expired, or kill-switched | Unchanged: immediate enqueue, `idempotencyKey = correlationId`                                                                                           |
| Terminal                                                                                                            | Unchanged: not in the pending set                                                                                                                        |

Two invariants keep the table sound:

* **No path acks the owning message while an owned step is non-terminal.** Ack means
  handler return or `reinvoke()` (which acks and continues under a *new* message ID); an
  acked owner can never be redelivered, so recovery would fall to the backstop for the
  full lease. All inline bodies are awaited before any ack path, a dev assertion
  (`error`-level log) guards the ordering against refactors, and turbo's `reinvoke()`
  paths are safe by construction: turbo requires delivery attempt 1, while owned-recovery
  steps can only exist on redeliveries (attempt ≥ 2), since a previous delivery of the same
  message must have stamped them. That mutual exclusion is documented at turbo's
  engagement gate.
* **Owned recovery must not be silently orphaned by early returns.** The background-step
  fast path used to return without replaying when other steps were still pending; it now
  falls through to the main loop when one of those pending steps is owned by the arriving
  message, so a redelivered owner actually performs its recovery instead of leaving it to
  the backstop.

## Why in-process single-flight is a required layer

The lease bounds *cross-instance* duplication only on platforms that kill invocations. On
world-local (one process, no kill bound) a delayed backstop can fire while the owning
execution is still mid-body *in the same process*, and on Fluid Compute, an owner
redelivery and a backstop can land on the same instance. A module-level map keyed
`runId:correlationId` absorbs both: the loser awaits the winner's settlement and then acks
**without executing**.

The loser must not ack-and-skip early (before the winner settles): a crash after an early
ack would consume the loser's message while the winner's outcome is unknown, potentially
orphaning the step with no message left to drive it. Awaiting settlement first keeps the
at-least-once envelope intact: if the loser's own invocation hits its deadline while
waiting, its message redelivers and re-checks, degrading gracefully to polling.

Cross-instance duplicates on multi-instance *self-hosted* worlds during steps longer than
the lease remain the documented residual (mitigate by raising the lease environment variable). This equals
the trade-off every lease-based system makes; eliminating it entirely requires the
heartbeat machinery rejected above.

## Alternatives considered and rejected

* **Flow-route queue concurrency = 1 (serialize all run messages)**: Removes the
  parallelism the wake mechanism exists for: `Promise.race(step, sleep)` works because
  the wait continuation fires in a *separate* invocation while the inline step blocks
  its handler. With one slot, the sleep could never win. It is also a queue-backend
  feature the World contract does not guarantee, and it serializes unrelated work
  (hooks, cancellations) behind long step bodies. Noted as a long-term option only if
  ownership proves unmaintainable.
* **Inline-eligibility latch (never inline while hooks/waits are open).** The cheapest
  hotfix (the condition is already computed for turbo's latch), but it permanently
  forfeits inline execution for exactly the workflows that use hooks, taxing every run to
  prevent a race that needs an in-flight step to matter. And it is incomplete:
  cancellation can wake *any* run mid-step, hooks or not.
* **Inline retries / transferring ownership across attempts.** Rejected above: backoffs
  are delay-dominated, attempts would pay replay costs to stay inline, and the
  queue-owned retry path already recovers correctly. Ownership deliberately ends at
  `step_retrying`.
* **Heartbeat / lock-store liveness.** Rejected above: new World surface for every
  backend, steady-state write amplification, and it still needs a lease to survive a
  crashed heartbeater. All cost, same bound.
* **Deriving the lease from the route's `maxDuration`.** Nothing to derive from: builders
  emit `'max'`, resolved per-plan by the platform at deploy; no runtime or build-time API
  exposes the resolved value.
* **Fixing it in the backend.** The server sees the same event log and has the same
  liveness blind spot; it would need its own claim mechanism, and world-local plus every
  community World would remain broken. The dispatch semantics live in the runtime, so the
  fix does too.

## Wire format, compatibility, and rollout

* `ownerMessageId` is an optional field on the `step_started` eventData schema
  (`@workflow/world`). In `@workflow/world-vercel` it rides the v4 frame meta; the
  compile-time wire guard (`assertEventDataWireContractExhaustive`) fails the build if a
  schema field is not explicitly routed, which is what forces the split/merge sites to be
  handled. The backend persists it verbatim on the event row (the lazy-input strip on
  `step_started` leaves it untouched, and the synthetic `step_created` never carries it)
  and re-emits it on event lists.
* **Deploy order: backend before SDK.** An older backend silently drops the meta field →
  replays see unowned steps → exactly today's behavior. Safe, but a pointless window to
  ship into.
* **Version skew is a non-issue**: runs are pinned to their deployment, and old events
  lack the field → unowned → current behavior. Nothing to migrate.
* **Kill switch**: `WORKFLOW_INLINE_OWNERSHIP=0` reverts dispatch to the unconditional
  immediate requeue. Stamping continues (it is inert data), so the switch is purely a
  dispatch-behavior toggle.
* **Degraded modes are all "today's behavior", never worse**: unstable message IDs (owner
  check never matches → backstop-lease recovery), missing/unusable event timestamps
  (lease remaining 0 → immediate enqueue), old backend (field dropped → unowned).

## Accepted residual risks

* **Redelivery-while-alive** (visibility lapse or heartbeat partition inside the queue):
  the owner check passes on a redelivery racing the live owner → duplicate. This equals
  the queue's at-least-once envelope (the floor for any client-side design) and is
  strictly rarer than the every-wake duplication being fixed. The in-process single-flight
  absorbs the same-instance case.
* **Multi-instance self-hosted worlds with steps longer than the lease** (see
  single-flight section): raise the lease environment variable.
* **A backstop per ownership epoch**: an owner crash-looping through its redelivery
  budget arms up to one delayed wake per re-stamp. Bounded by the queue's delivery
  budget; each fires as a cheap replay no-op if the step completed.

## Observability

* Span attributes on `workflow.execute`:
  `workflow.inline_ownership.owned_recovery_steps` (crash recovery re-executed owned
  steps) and `workflow.inline_ownership.backstop_wakes_armed` (a replay suppressed an
  immediate requeue).
* Always-printed `warn` logs when owned recovery runs (a prior delivery died mid-body)
  and when the single-flight absorbs a would-be duplicate (a burst of these means leases
  are expiring under live executions: raise the lease environment variable). Invariant violations log at
  `error`. Backstop arming logs at `debug` (`DEBUG=workflow:runtime:*`), since it can
  legitimately fire on every wake replay during a long inline step.

## How the design is protected by tests

* The **#2780 repro** (`workbench/vitest/test/inline-step-ownership.test.ts`): hook
  resume mid-inline-step → side-effect marker fires exactly once. Verified bidirectional:
  with `WORKFLOW_INLINE_OWNERSHIP=0` the marker fires twice.
* **Unit**: the ownership state machine (stamp → wake sees owner → retrying clears → bare
  start clears → re-stamp restores), lease math including the clock-skew clamp, and the
  backstop key's epoch behavior, including a regression test walking the full
  owner-recovery re-arm sequence against a dedupe model matching world-local's in-flight
  key retention; single-flight winner/loser semantics.
* **Wire**: schema round-trip tests on both sides, plus backend integration tests
  asserting the field survives materialization and the lazy-input strip, and never leaks
  into synthetic `step_created`.
* **End-to-end (E2E)**: The full suite including the abort/cancellation lanes that caught backstop
  shape #1: world-vercel prod lanes are mandatory for sign-off, since world-local's
  synchronous single-process behavior masks distributed races.

## Open questions

1. Confirm platform behavior: `'max'` resolves ≤ 800s even for accounts in the 30-minute
   beta, and users cannot raise a generated workflow route's `maxDuration` via
   `vercel.json` `functions` config without the builder seeing it (if they can, the
   builder should detect the override and fail or bake it into the lease).
2. Does world-postgres's handler meta deliver a stable ID across retries? If not:
   document degraded (backstop-lease) mode for community worlds, don't block on it.


---

For a semantic overview of all documentation, see [/sitemap.md](/sitemap.md)

For an index of all available documentation, see [/llms.txt](/llms.txt)

For agent-facing discovery, including API and MCP surfaces, see [/agents.md](/agents.md)