Skip to content

DX-2809: step-level flow control and retries - #206

Open
CahidArda wants to merge 36 commits into
mainfrom
DX-2809-step-level-flow-control
Open

CahidArda wants to merge 36 commits into
mainfrom
DX-2809-step-level-flow-control

Conversation

@CahidArda

@CahidArda CahidArda commented Jul 7, 2026 •

Copy link
Copy Markdown
Collaborator

Summary

Flow control and retries can now be set on a single step of a run, as the third
argument of context.run:

await context.run(
  "charge customer",
  () => chargeCustomer(),
  {
    flowControl: { key: "payment-provider", parallelism: 3 },
    retries: 5,
    retryDelay: "1000",
  }
);

They override, for that step only, whatever the run was triggered with. Anything
left unset falls back to the trigger configuration.

Needs the server side: qstash-server branch DX-2809.

src/GLOSSARY.md is the reference for the vocabulary, the two decisions and the
per-step-kind sequence diagrams. This description is the short version.

The problem

A step's settings have to be on the message whose delivery executes that
step
— QStash applies flow control and retries when it delivers, so by the
time the endpoint is running the step, it is too late to ask for anything.

But the SDK only learns a step exists by replaying the workflow function, which
happens inside a delivery that was published before that step was known. The
message that would have needed the settings is already gone.

An earlier iteration tried to solve this by predicting the next step at publish
time, through three separate discovery mechanisms. It could not be made correct —
branching, loops, and anything non-deterministic between steps all defeat
prediction — and it left a list of documented limitations.

The approach

Don't predict. Publish everything as usual, then correct course when a step
turns out to need something different.

QStash reports on every delivery what it actually applied — flow control key and
value, retry limit, retry delay. The SDK compares a step's settings against that
effective configuration the moment it is about to run the step:

Effective configuration Guard marker Action
matches the step's — execute
differs absent publish a step config request, abort
differs present execute anyway, and warn

A step config request is a hidden helper message carrying nothing but
configuration. QStash gives it its own call type (stepConfig), keeps it out of
the step logs, and delivers it back to the endpoint — and that delivery, now
carrying the right settings, executes the step.

The third row is the bounded failure mode. The guard marker
(Upstash-Workflow-Step-Config) is set by QStash from the message's own feature
set, so it cannot disagree with what QStash actually did. Its presence means "you
already asked" — so a settings mismatch that survives a step config request runs
the step with the wrong settings and an onWarning, rather than asking again and
looping invisibly. It caught two real normalization bugs during development that
no unit test would have.

That is the whole mechanism. Two cases skip the extra request entirely:

  • Parallel steps carry their settings on their own plan steps, which QStash
    is publishing anyway.
  • A step following one whose result the SDK computed itself — see below.

Deferred submission

context.run, sleep and sleepUntil produce their result in-process. Rather
than ending the invocation the moment such a step finishes, the SDK now holds
the result
and lets the route function carry on until it reveals what comes
next. If the next thing is a step with settings, they ride on the held step's
submission — the message QStash was about to receive anyway — and no step config
request is needed.

The held result is submitted on every path out, including the ones that are easy
to get wrong:

What the continuation reaches Attached to the pending submission
a single step that step's settings, if it has any
a parallel group none — plan steps carry their own
route function returns none
route function throws / cancels none; the error resurfaces deterministically on the next replay

The last row matters: a held result that never reached QStash would make its step
run twice. onStep submits before it rethrows, and there is an end-to-end test
for it.

Reading the delivery's configuration

The effective configuration is also exposed, since the SDK had to parse it anyway:

context.flowControl  // { key, parallelism, rate, period } | undefined
context.retries      // number | undefined
context.retryDelay   // string | undefined

These report what QStash applied to the delivery in hand, which is what the
step-settings comparison runs against.

Changes

  • src/qstash/step-config.ts (new) — parses the effective configuration off
    a delivery, and describeStepSettingsMismatch, which is the comparison the
    decision above rests on. Both sides normalize first: QStash joins the flow
    control value without spaces and reports periods in whole seconds, while the
    SDK sends durations and joins with ", ". It also mirrors QStash's defaults
    (unset period is one second, unset parallelism and rate are zero) — a mismatch
    here makes every delivery look wrong, so it has its own unit tests and a
    live CI route.
  • src/context/auto-executor.ts — the decision, and the pending-step
    machinery. The held step is one value in two states (held / submitting)
    rather than two fields that could disagree; submitPendingStep and
    flushPendingStep return neverthrow results carrying a
    { result: "submitted-step", abort } | { result: "no-pending-step" } union.
  • src/serve/index.ts — onStep runs the route function, flushes any held
    step, then decides what to throw. Each exit is commented with when it is
    reached.
  • src/qstash/submit-steps.ts — publishStepConfigRequest, published with
    content-based deduplication so a repeat is a no-op rather than a new message.
  • src/qstash/headers.ts — getHeaders takes the step's settings and builds
    their headers in a bucket spread last, never prefixed.
  • src/context/context.ts, src/context/steps.ts, src/types.ts — the
    third argument, StepSettings, and the three getters.
  • src/workflow-parser.ts, src/serve/authorization.ts — the derived
    contexts (failure, callback, authentication) get the delivery's configuration
    too, instead of reporting defaults.
  • src/constants.ts — the new call type, the guard marker, and
    Upstash-Max-Retries. Note that Upstash-Retries means the retry limit when
    publishing but the retry count on a delivery, which is why the limit needed a
    header of its own.
  • src/GLOSSARY.md (new) — vocabulary, the two decisions as flowcharts and
    tables, sequence diagrams for context.run / call / waitForEvent /
    waitForWebhook before and after, the loop backstops, and the normalization
    rules.

Testing

Unit tests for the executor's decision, the pending-step lifecycle, the
normalization round trip, the header construction, and the derived contexts.

Seven live CI routes (examples/ci/app/test-routes/flow-control/), run against a
local qstash-server carrying the matching server change. Five are concurrency
tests: a coordinator invokes three workers in parallel, each worker reaching a
step with parallelism: 1 while the run itself is triggered with
parallelism: 5. The step increments a shared counter, holds it, and decrements
it — so the value a worker sees is how many workers were inside the step at once,
and it must never exceed one. They differ in what the worker does before that
step, which is what decides how the delivery reaching it was produced:

  • first-step — no earlier step to carry the settings.
  • step — after a context.run, the one case costing no extra request.
  • call — after a context.call.
  • invoke — after a context.invoke.
  • parallel — four steps in one parallel group, in two pairs, each pair
    sharing a key of parallelism: 1, so the two keys have to hold their pairs
    back independently.

And two which are not about concurrency:

  • normalization — walks a range of settings shapes and asserts each comes
    back in a form the SDK recognizes, plus an exact request count. The count is
    what would catch a step being republished in a loop, which the assertions alone
    would not.
  • invoke-failure — a worker throws while the SDK is still holding a step's
    result. The coordinator checks that the invoke reports the failure, and reads
    the worker's own log to confirm the held result reached QStash before the error
    did.

Removing the step-level settings makes all five concurrency routes fail, so they
are measuring the feature rather than passing by construction.

Loop safety was checked by building an SDK that ignores the guard marker and
loops on purpose: the run stops after ten endpoint invocations instead of running
away. The three backstops (content deduplication, the server rejecting two
stepConfig entries in a row, and the entry limit) are described in the glossary
along with what each does to the run.

@linear-code

linear-code Bot commented Jul 7, 2026

Copy link
Copy Markdown

DX-2809

@CahidArda
CahidArda marked this pull request as draft July 14, 2026 08:31
…discovery

A step's settings must be on the request whose delivery executes the
step. Previously three mechanisms tried to get that request right at
publish time: deferred submission (executing a step and continuing the
workflow function to reveal the next step), a discovery replay when a
context.call result arrived, and a marker header on context.invoke
results. Three cases were still unreachable and every workflow paid for
the machinery whether or not it used withSettings.

Instead, publish everything as usual. When the executor is about to run
a step which has settings in a delivery that wasn't gated by them, it
publishes a hidden `discovery` request carrying the settings and aborts;
the step executes when QStash delivers that request. The normal replay
is the discovery, so no separate discovery mode is needed.

Loop prevention reads the published redeliveries off the run's steps
(hidden `discovery` entries carry their target step id) rather than a
request header, so the decision is deterministic across replays and
QStash retries. Parallel steps are unchanged: each carries its settings
on its own plan step, with no extra request.

This removes all three limitations (first step of a run, step after
waitForEvent, step after a parallel group), and workflows which don't
use withSettings behave exactly as before. The cost is one extra message
and one extra endpoint invocation per step which uses withSettings.

Removes the deferred submission bookkeeping, WorkflowDiscoveryAbort,
discovery mode, src/serve/discovery.ts and the invoke result marker;
handleThirdPartyCallResult no longer fetches the run's steps.

The qstash-server side is unchanged: WF_StepConfig and the hidden
`discovery` call type are exactly what this needs.
It lives in utils.ts since the move that broke the circular import; the
re-export from workflow-requests.ts only kept the old import path alive.
Replaces the step-config bookkeeping with a direct comparison. QStash
already reports the configuration it applied to a delivery, so the
executor can compare a step's settings against it and decide what to do,
rather than tracking which steps it has already published a request for.

On reaching a step with step-level settings which has not executed yet:

  effective config | guard marker | action
  -----------------|--------------|--------------------------------------
  matches          | -            | execute
  differs          | absent       | publish a step config request, abort
  differs          | present      | execute anyway, and warn

The third row is a bounded failure mode rather than a normal path. It
means a configuration this SDK published did not come back in a form it
recognizes, which is a normalization bug: executing with the wrong
settings is preferred over publishing a second step config request, since
that would loop forever and invisibly. The guard marker is emitted from
getStepSettingsHeaders as a forwarded header, so every producer of
step-level settings sets it in one place.

Steps whose result is available in-process (run, sleep, sleepUntil) now
hold their result and let the route function continue to the next step, so
that step's settings can ride on the submission instead of needing a step
config request. Consecutive gated steps therefore cost one request for the
first and none after it. The cost is that code after a step also runs in
the invocation which executed it, and that an error or context.cancel()
after a step is deferred by one delivery.

Normalization is the delicate part, since both sides format the same
values differently. Two divergences were found by the guard marker firing
during a live run, not by unit tests: QStash defaults an unset flow
control period to one second and reports it back, and `Upstash-Retries` on
a delivery counts the retries that already happened rather than reporting
the limit (the limit now arrives as `Upstash-Max-Retries`, added on the
server side). Hence the live normalization matrix, which asserts both the
round trip and an exact request count for a range of shapes.

Also drops `timeout` from StepSettings: QStash does not report it on a
delivery, so a mismatch in it could never be detected and the setting
would silently never apply. Renames the `discovery` call type to
`stepConfig` to match the server. Adds src/GLOSSARY.md describing the
request kinds and the rules above.
…kflow

The step-level flow control tests need concurrent runs to show that the
gate serializes them, and got that by calling initiateTest three times for
the same route — the only tests in the suite doing so, and the reason they
could not live in TEST_ROUTES with everything else.

Move the fan-out into the workflow. Each route is now a coordinator which
invokes three workers in parallel; the workers are the concurrency, and
the coordinator collects what they observed and saves the run's result. So
every route is triggered once, TEST_ROUTES holds them like any other
route, and ci.test.ts goes back to its original form.

The shared pieces (the increment step, its settings, the CI header
forwarding, the assertion) move to a shared module, leaving each route to
define only what its worker does before the gated step — which is the part
that differs between them.

Also asserts the guard marker in the normalization route: true inside a
gated step, and false inside a step with no settings, so a marker which
stopped arriving, or one set unconditionally, fails the suite.
A step config request was published without deduplication, so a delivery
which published one and then failed would publish it again on retry. That
second request is redundant — the first already produced the gated
delivery — and it trips the server rule which rejects two step config
requests in a row, turning a retry into a 400 the SDK reports as "tried to
append to a cancelled workflow".

Deduplicate on content instead, which collapses the repeat into the
original. The body already carries the target step, so the requests of two
different steps stay distinct, and the deduplication hash covers the run
id, so runs never collide with each other.

Verified against a local qstash-server by building an SDK which ignores
the guard marker and loops on purpose: the run stops after one repeat
instead of running away.
The three backstops stop a loop in three different ways and only the entry
limit ends the run visibly: deduplication and the two-in-a-row rejection
both leave the run stalled, since no new message is published and every
delivery still answers 200. Nothing fails and no callback fires, so the
only symptom is a run which never completes.

Also warn when a step config request comes back deduplicated. That is
expected when the delivery which published one is being retried, and the
original still produces the gated delivery — but it is also what a client
asking twice for the same step sees, and in that case the warning is the
only trace left behind.
The step settings headers were built by an exported helper and spread
over the result of getHeaders at each of the three call sites, so every
producer of them had to remember to do that. Pass the settings to
getHeaders instead and let it assemble them, like every other part of a
workflow request.

They go in a group of their own, applied last and never prefixed: they
configure the message itself, so they must win over the settings the run
was triggered with, and must not pick up the `Upstash-Callback-` prefix
a call step gives the rest of its workflow headers.

Also makes the effective configuration a private field of the context,
passed to the executor directly rather than reached through the context.
It is the executor that compares against it; the context exposes only
`flowControl`, `retries` and `retryDelay`. The CI route which asserts the
guard marker arrives now reaches in through a named helper rather than a
`@ts-expect-error`, so it keeps working and says why.
triggerRouteFunction is generic control flow — run a step, clean up,
handle cancellation — and had grown a workflow context parameter for the
sole purpose of submitting a step whose result was still held. That is a
concern of running the route function, not of the surrounding control
flow, and it pushed the parameter onto every call site including the
tests.

Do it where the route function is invoked instead. triggerRouteFunction
goes back to what it was, byte for byte, and the tests which end their
route function right after a step now submit the held result themselves,
which is what serve does.
The executor tracked a held step across two fields whose relationship
only existed in a comment: pendingSubmission for one not yet submitted,
pendingFlush for the submission of one already on its way. Nothing stopped
them being set together, and the flush was a Promise<never> the caller
awaited for its rejection, so `if (flush) await flush` read as if it might
carry on when in fact it never did.

One field with two states instead, and the abort is returned rather than
thrown from inside a promise nobody names: submitPendingStep hands back
the WorkflowAbort to throw, or undefined when nothing was held, so the
control flow is visible at the call site. A submission failure still
throws, so that error reaches the caller instead of an abort claiming the
step was submitted.

Also runs executeStep before branching on whether the step's result can be
held, which is what submitSingleStep did internally anyway, so the two
callers compose executeStep and submitStepResult directly and
submitSingleStep goes away.

And renames gateStep to publishedStepConfigRequest, after what its
boolean means.
`if (await this.publishedStepConfigRequest(...))` said nothing at the call
site about what either branch meant, and the name described the true case
only. Return the outcome the way the rest of the codebase reports one —
`Ok<"execute-step" | "step-config-requested">` alongside an error channel,
as `tryAuthentication` does — so both branches are named where they are
read.

The error channel is not decoration: publishing a step config request can
fail, and swallowing that would leave the caller believing the step was
gated when nothing was submitted. Covered by a test which fails if the
error is turned back into an outcome.
Same treatment as applyStepSettings: submitPendingStep returns
`Ok<WorkflowAbort | undefined>` alongside an error channel rather than
returning the abort and throwing on failure, so both ways the call can end
are visible where it is read.

The error channel matters here for the same reason: if a failed submission
were reported as "nothing was held", the caller would carry on and execute
the next step in an invocation whose previous step QStash never received —
the step would then run again on the next delivery. Covered by a test which
fails if the error is turned back into an outcome.

flushPendingStep stays as it is: it exists to reach the executor through
the context and end the invocation, so unwrapping the result and throwing
is its whole job.
flushPendingStep was still unwrapping the result and throwing, so its
callers could not see that submitting a held step is what ends their
invocation — the one place that mattered most, since serve calls it on
both the way out and the way out of a throw.

Returning the result instead collapses those two calls into one. The
route function's outcome is captured first, the held step is submitted
once, and the order the three ways out are checked in states the
precedence which the two-call version only implied: a submission failure,
then the held step's abort, then whatever the route function threw.
recreateUserHeaders was moved to utils.ts, and processRawSteps and
deduplicateSteps to a new raw-steps.ts, to break import cycles that the
first version of this feature created: workflow-requests.ts reached into
the discovery replay, and needed the raw step parsing for the call result
discovery. Neither exists any more, so neither move is needed, and the
diff was carrying three relocated functions for nothing.

Put them back. utils.ts is identical to main again, raw-steps.ts is gone,
and nothing in this PR relocates a function. The functions themselves are
unchanged, apart from the comment in processRawSteps which now names the
stepConfig call type its filter also skips.

Fewer import cycles than before the move, not more: madge counts 33 where
the moved layout had 34.
`Ok<WorkflowAbort | undefined>` left the two cases to be told apart by
whether a value happened to be there, so every caller read as a
truthiness check on something whose absence meant "there was no step to
submit" — a meaning the type never stated.

Say it:

    { result: "submitted-step"; abort: WorkflowAbort } | { result: "no-pending-step" }

Callers switch on the result and reach for the abort only in the branch
that has one, so the branch which does nothing is now as legible as the
branch which ends the invocation.
A step whose result is available in-process is held so the route function
can continue. If the function then throws, the step has already run, so
its result still has to reach QStash — otherwise the next delivery would
run it a second time. Nothing covered that: the failureFunction routes
throw from *inside* a step, where executeStep fails before anything is
held, which is the case where there is correctly nothing to submit.

Drives two deliveries: the first holds the step and throws, and still
submits; the second replays the memoized step and lets the throw surface
with nothing left to submit. The step function runs once across both.

Fails if the submission is skipped when the route function threw.
The comments described the mechanism but left the reader to work out
which of the four exits a given run takes. Put the explanation at each
exit instead, saying what has to have happened to arrive there: the
route function not running to the end, a held step failing to publish, a
held step reaching QStash, a throw with nothing held, and a clean return
with every step already recorded.

Also pins one of the claims: reaching a further step submits the held one
and throws the abort from there, so serve flushes after the route
function has already seen it. That flush reports the same abort rather
than "nothing was held", which is what stops an invocation whose abort
the route function swallowed from carrying on as if the step never ran.
context.flowControl, retries and retryDelay report the configuration
QStash applied to the request in hand, but only the context serve builds
for the run was given it. The context the authorization dry-run runs on,
and the one a failure function receives, both reported nothing — so a
route function reading them before its first step, or a failure function
reading them at all, saw defaults instead of what QStash actually applied.

serve now hoists the configuration and hands it to tryAuthentication,
which passes it to the disabled context; handleFailure reads it from the
failure callback delivery, which is the request in hand there. The
trigger-side context in the client keeps reporting nothing, correctly:
there is no delivery behind it.

Both are covered: the route function reads them during the dry-run, and
the failure function reads them from a callback carrying flow control,
a retry limit and a retry delay.
…tted

The existing test called flushPendingStep directly, so it pinned the
executor's error channel but not what a delivery actually does with it.
Drive it through serve instead: QStash answers 500 to the submission, the
endpoint answers 500 so the delivery is retried, and the retry runs the
step again because the run has no record of it.

Swallowing that failure turns out to be worse than losing the step's
result, which is what makes the test worth having: the route function
returned, so the run would be reported finished and deleted with the step
never recorded. Breaking the error channel fails this test on exactly
that — the second delivery issues the delete instead of the submission.
The delivery answers 500 without ever reaching onCleanup, so the run is
not deleted when a held step cannot be submitted. Worth saying in the
test, since what proves it is indirect: the mock asserts the method and
url of every request it is given, so the delete would fail it.
Adds a Flows section to the glossary: a paragraph on what changed between
the SDK and QStash, and diagrams per step kind showing the same run
before, now without step-level settings, and now with them.

The point the diagrams make is that only one thing moved. A step whose
result the SDK produces itself no longer submits and stops — it holds the
result and lets the route function carry on — so the deliveries are the
same and only the moment of submission changed. That held message is what
the next step's settings travel on. Everything else (call, invoke, the
waits) is untouched, because its result comes from QStash and there is
nothing to hold, which is exactly why a step after one of those needs a
step config request of its own.
applyStepSettings existed to hand runSingle one of two outcomes, and
needed a result type and an error channel to do it. Inlined, the decision
is three lines of ordinary control flow next to the step it is about:
compare, and either ask for a gated delivery or warn and carry on.

The neverthrow result goes with it. It was reconstructing, across a
function boundary, what a throw already does — publishStepConfigRequest
failing now propagates out of runSingle by itself, which is what the
error channel was rebuilding by hand.

No test changed: a failing publish still surfaces, and the guard still
executes the step with a warning rather than asking twice.
Both branches only apply when the step's settings and the delivery's
disagree, so gate on that once and let the inner check say which of the
two it is: ask for a gated delivery, or run the step with a warning
because we already asked.
Adds a flowchart for each decision the step-level settings work turns on,
placed beside the rules it draws rather than among the sequence diagrams:
what the SDK does on reaching a step whose settings differ from the ones
its delivery is under, and what a context.run does once it finishes —
holding its result and letting what the route function reaches next decide
what that result is submitted with.

Renames the vocabulary. "Gated" meant nothing until you had read the
glossary, and the reader has already met the step config request, the
stepConfig call type and Upstash-Workflow-Step-Config by the time they
need it — so a delivery carrying a step's own settings is a
**step-configured delivery**, and the other kind needs no name of its own:
it is an **ordinary delivery**. Applied across the glossary, the diagrams,
the code comments and the CI routes, where wasGated and isGatedDelivery
became withinStepParallelism and carriesStepSettings, which is what they
actually asked.
The step after a parallel group always needs a step config request, and
nothing covered it. Unlike a step after a single step, its settings cannot
ride on anything: the submission which completes the group would have to
know what follows it, which means carrying on past the group — and no
delivery of a parallel step can tell whether it is the last, since its own
result is unrecorded and so is at least one sibling's.

Says that where the question comes up, next to the parallel submission
which could otherwise look like an oversight, and in the glossary.

Also fixes fallout from the previous rename: `ungatedConfig` matched the
`gatedConfig` rule first and became `unstepConfigured`. The two fixtures
are now `stepConfiguredDelivery` and `ordinaryDelivery`.

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

Adds step-level configuration overrides (flow control / retries / retry delay) for individual workflow steps via context.run(...).withSettings(...), including new execution/republish mechanics to ensure settings are applied to the delivery that executes the step.

Changes:

  • Introduces RunStepPromise.withSettings() and StepSettings to attach per-step overrides to context.run.
  • Implements deferred submission + step-config request publishing so per-step settings can be applied on the correct subsequent delivery (including normalization/mismatch detection via “effective config”).
  • Adds unit and CI example routes validating header overlay, mismatch handling, and concurrency enforcement.

Reviewed changes

Copilot reviewed 30 out of 30 changed files in this pull request and generated 4 comments.

Show a summary per file
File Description
src/workflow-requests.ts Adds internal flushPendingStep helper to submit deferred step results.
src/workflow-requests.test.ts Updates tests to flush deferred submissions when route ends after a step.
src/workflow-parser.ts Plumbs “effective config” (delivery-applied config) into failure handling/context.
src/workflow-parser.test.ts Adds tests for skipping stepConfig entries and for failure-function effective config.
src/types.ts Adds StepSettings + RunStepPromise types and documents per-step settings behavior.
src/serve/serve.test.ts Adds tests ensuring held steps are still submitted on throw/cancel and on submit failures.
src/serve/index.ts Ensures any held/deferred step is flushed on all route-function exit paths.
src/serve/authorization.ts Passes effective delivery config into the “disabled” auth-check context run.
src/serve/authorization.test.ts Updates auth tests for deferred submission and effective-config visibility.
src/qstash/submit-steps.ts Splits step execution vs submission; adds step-config request publishing and settings overlay for submissions/plan steps.
src/qstash/step-config.ts Adds parsing/normalization of delivery-applied QStash config + mismatch description helpers.
src/qstash/step-config.test.ts Adds unit tests for normalization, effective-config parsing, and mismatch detection.
src/qstash/headers.ts Adds step-settings header overlay + WF_StepConfig feature flag handling.
src/GLOSSARY.md Adds contributor-focused protocol/architecture glossary for step config + deferred submission.
src/context/steps.ts Adds stepSettings to lazy steps and marks which step types support deferred submission.
src/context/context.ts Changes run to return a chainable RunStepPromise and exposes effective delivery config on context.
src/context/context.test.ts Updates tests to match deferred submission behavior (flush triggers abort).
src/context/auto-executor.ts Implements pending-step holding/flushing, step-config request logic, and settings mismatch handling.
src/context/auto-executor.test.ts Adds comprehensive tests for step-level settings, pending-step behavior, and step-config requests.
src/constants.ts Adds constants for feature sets, step-config call type, effective-config headers, and step-config marker.
examples/ci/app/test-routes/flow-control/step/workflows/[...]/route.ts Adds CI route validating step-level flow control when the configured step follows a run.
examples/ci/app/test-routes/flow-control/shared.ts Shared helpers for CI concurrency assertions and common step-level settings.
examples/ci/app/test-routes/flow-control/README.md Documents CI routes and how step-level settings are applied in different step sequences.
examples/ci/app/test-routes/flow-control/normalization/route.ts Adds CI route validating settings round-tripping/normalization and expected call counts.
examples/ci/app/test-routes/flow-control/invoke/workflows/[...]/route.ts Adds CI route validating step-level settings after an invoke.
examples/ci/app/test-routes/flow-control/first-step/workflows/[...]/route.ts Adds CI route validating step-level settings on a run’s first step (requires step-config request).
examples/ci/app/test-routes/flow-control/call/workflows/[...]/route.ts Adds CI route validating step-level settings after a call.
examples/ci/app/test-routes/flow-control/call/target/route.ts Adds a simple call target endpoint used by the CI call route.
examples/ci/app/ci/upstash/redis.ts Tightens result polling to wait for the expected call count (accounts for extra replay at tail).
examples/ci/app/ci/constants.ts Registers the new flow-control CI routes in the test route list.

💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.

Comment thread src/workflow-parser.ts Outdated
Comment thread examples/ci/app/test-routes/flow-control/invoke/workflows/[...]/route.ts Outdated
Comment thread examples/ci/app/test-routes/flow-control/call/workflows/[...]/route.ts Outdated
Comment thread src/types.ts
`.withSettings()` on the returned promise meant `context.run` could not
just return a promise, and read as though the settings were attached
after the fact rather than being part of the step. They are now the third
argument, which also drops the `RunStepPromise` type.

Two CI routes cover the cases where a step with settings does not cost a
step config request:

- `parallel`: four steps in one parallel group, in two pairs which each
  share a flow control key of `parallelism: 1`. The settings ride on the
  plan steps, and the two keys have to hold their pairs back
  independently.
- `invoke-failure`: a worker throws while the SDK still holds the result
  of a step with settings. The held result has to reach QStash before the
  error does, which the coordinator checks in the worker's own log.

The flow control keys and counters the routes share are now scoped to the
test run. A flow control key is server-side state, so a fixed one made
re-runs queue behind whatever the previous run left in flight, and a
shared counter made one run observe concurrency another run caused.
The failure callback's `effectiveConfig` had lost the subject of its
sentence, and the `call` route's doc comment had a clause spliced into
the middle of another. Both now say what the equivalent comments in
`serve/index.ts` and the sibling routes say.
@CahidArda CahidArda changed the title DX-2809: step-level flow control and retries via context.run().withSettings() DX-2809: step-level flow control and retries Aug 25, 2026
Two steps per group only distinguishes "held back" from "not held back";
three also catches a group which lets a second step through but not a
third. Six plan steps means six result deliveries, so the expected call
count is 13 rather than 11, and the group size is now a constant so the
two places that depend on it cannot drift apart.

Drops the glossary. Its vocabulary lives in the comments at the places it
described, which is where it stays accurate.
qstash-server#1058 renamed the header which reports the configured retry
limit on a delivery: it lands on Upstash-Retries, not Upstash-Max-Retries
(commit b07eaf78, which sets it last and unconditionally, overwriting the
retries-so-far value that header used to carry). Upstash-Retried is still
the header for how many retries have already happened.

Reading a header QStash never sends left context.retries undefined on
every delivery, which failed flow-control/normalization in CI on
`expect(context.retries, 0)`, and — worse — silently disabled step-level
retries: the `effectiveConfig.retries !== undefined` guard in
describeStepSettingsMismatch short-circuits, so a step which sets only
`retries` never reported a mismatch and never got a step config request.

MAX_RETRIES_HEADER is dropped rather than repointed, since it would now
hold the same string as RETRIES_HEADER.

Also caps the retry budget in the two redis.test.ts cases which assert a
call count mismatch. They pass a count that can never match, so the
wait-for-expected-call-count loop polled for its full 40s budget and blew
vitest's 5s default timeout.

Claude-Session: https://claude.ai/code/session_018rmVdGUi5NiBGuxi7nfgnC
@CahidArda
CahidArda marked this pull request as ready for review September 1, 2026 10:37
@CahidArda
CahidArda requested a lite review from Copilot September 11, 2026 09:40

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Changes recommended

Critical issues remain in deferred execution, parallel configuration checks, and retry-limit parsing.

Once you've addressed the issues Copilot identified, you can request another Copilot review.

Review details

Suppressed comments (5)

examples/ci/app/ci/constants.ts:159

  • The comment says these steps are arranged in “two pairs,” but this route actually defines two groups of three (STEPS_PER_GROUP = 3). Correct the description so future changes to the test do not rely on an inaccurate topology.
    // the steps with settings are in a parallel group, in two pairs which
    // each share a flow control key

examples/ci/app/test-routes/flow-control/parallel/route.ts:32

  • These keys and the Redis counters are global to every invocation of this CI route, unlike the other flow-control routes that use perRunKey. Concurrent test runs will share QStash slots and measurements, which can serialize runs or contaminate their observed counts and make the test flaky. Scope both values with perRunKey(context, ...) as done in shared.ts.
const counterOf = (group: string) => `wf-step-flow-control-parallel-${group}-active-counter`;
const flowControlOf = (group: string) => `ci-step-flow-control-parallel-${group}`;

src/context/auto-executor.ts:309

  • This documentation no longer matches the neverthrow implementation: submitPendingStep returns ok({ result: "no-pending-step" }) when empty and err(error) when submission fails; it neither returns undefined nor throws in those cases. Describe the Ok/Err outcomes so callers do not implement the wrong control flow.
   * Returns undefined when no step is being held, which is the caller's
   * signal that there is nothing to end. Throws instead of returning
   * when the submission itself fails, so that error reaches the caller
   * rather than an abort claiming the step was submitted.

src/qstash/step-config.ts:238

  • addStepSettings deliberately omits an empty retryDelay (so no WF_StepConfig marker is published), but this comparison still treats "" as set. A context.run(..., { retryDelay: "" }) therefore publishes an unmarked helper, sees the same mismatch again, and can loop/deduplicate forever. Treat an empty string as unset here too.
  if (
    stepSettings.retryDelay !== undefined &&
    stepSettings.retryDelay !== effectiveConfig.retryDelay
  ) {

src/types.ts:409

  • This public API documentation says every non-parallel step with settings costs an extra request, but the new deferred-submission path avoids that cost when the preceding run/sleep/notify result is produced in-process. Document that case explicitly and reserve the hidden-request description for steps reached after calls, invokes, waits, or other non-deferred steps.
 * - otherwise, the SDK can only learn about the step once the workflow
 *   function has been replayed, which happens in a delivery that was
 *   published before the step was known. So instead of executing the
 *   step in that delivery, a hidden request carrying the
 *   settings is published and the step executes when QStash delivers it.
  • Files reviewed: 32/32 changed files
  • Comments generated: 3
  • Review effort level: Lite

Comment thread src/context/auto-executor.ts
Comment thread src/context/auto-executor.ts
Comment thread src/qstash/step-config.ts
@CahidArda

Copy link
Copy Markdown
Collaborator Author

On the suppressed comments:

  • retryDelay: "" — right, the header builder skips it but the comparison doesn't. I'll treat empty as unset there too.
  • submitPendingStep doc, the "two pairs" comment, the StepSettings doc — stale, I'll update them.
  • Parallel route keys — already scoped per run at the call sites.

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Changes recommended

Empty retry delays can stall runs, and mutable exposed configuration can corrupt matching decisions.

Get a fresh assessment by requesting another Copilot review.

Review details
  • Files reviewed: 32/32 changed files
  • Comments generated: 3
  • Review effort level: Balanced

Comment thread src/context/context.ts
Comment on lines +197 to +199
public get flowControl(): NormalizedFlowControl | undefined {
return this.effectiveConfig.flowControl;
}
Comment thread src/qstash/step-config.ts
Comment on lines +235 to +238
if (
stepSettings.retryDelay !== undefined &&
stepSettings.retryDelay !== effectiveConfig.retryDelay
) {
Comment thread src/types.ts
Comment on lines +405 to +412
* - otherwise, the SDK can only learn about the step once the workflow
* function has been replayed, which happens in a delivery that was
* published before the step was known. So instead of executing the
* step in that delivery, a hidden request carrying the
* settings is published and the step executes when QStash delivers it.
* This costs one extra message and one extra endpoint invocation per
* step which carries settings. QStash hides the request from the
* step logs.
@ytkimirti

Copy link
Copy Markdown
Contributor

Fixed two issues in separate commits:

1. Preserve step results before deferred publication (a764037)

const source = { value: 1 };
const snapshot = await context.run("snapshot", () => source);

source.value = 2;

await context.run("consume", () => snapshot.value);

Before: the executor held a reference to source until publication. Changing it after the step returned caused { value: 2 } to be saved, even though the step originally returned { value: 1 }. Replay then used the changed value.

After: the result is serialized when the step finishes, and a detached snapshot is held for publication. The saved result and replay both retain { value: 1 }. Mutating the object returned to the caller also cannot change the saved result, and the original object's toJSON() is called only once.

2. Normalize retry settings before comparison and publication (adbcf30)

await context.run("work", doWork, { retryDelay: "" });
// Whitespace-only values behave the same way:
await context.run("more-work", doMoreWork, { retryDelay: "   " });

Before: an empty delay counted as an override during comparison but was omitted from the outgoing headers. This could repeatedly request a configuration delivery without its guard marker, leaving the workflow stalled after deduplication.

After: empty or whitespace-only delays mean no override, so the trigger's retry delay is inherited. Nonempty delays are trimmed consistently. Invalid retry counts and non-string delay values are rejected early; retries: 0 remains valid.

Added regression tests for snapshot mutation/replay, custom serialization, and retry settings across helper requests, deferred publication, and parallel steps. The second commit addresses the empty retryDelay review comment.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants