Conversation
…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`.
There was a problem hiding this comment.
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()andStepSettingsto attach per-step overrides tocontext.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.
`.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.
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
There was a problem hiding this comment.
🟡 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 withperRunKey(context, ...)as done inshared.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:
submitPendingStepreturnsok({ result: "no-pending-step" })when empty anderr(error)when submission fails; it neither returnsundefinednor throws in those cases. Describe theOk/Erroutcomes 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
addStepSettingsdeliberately omits an emptyretryDelay(so noWF_StepConfigmarker is published), but this comparison still treats""as set. Acontext.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
|
On the suppressed comments:
|
There was a problem hiding this comment.
🟡 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
| public get flowControl(): NormalizedFlowControl | undefined { | ||
| return this.effectiveConfig.flowControl; | ||
| } |
| if ( | ||
| stepSettings.retryDelay !== undefined && | ||
| stepSettings.retryDelay !== effectiveConfig.retryDelay | ||
| ) { |
| * - 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. |
|
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 After: the result is serialized when the step finishes, and a detached snapshot is held for publication. The saved result and replay both retain 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; 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. |
Summary
Flow control and retries can now be set on a single step of a run, as the third
argument of
context.run: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.mdis the reference for the vocabulary, the two decisions and theper-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:
A step config request is a hidden helper message carrying nothing but
configuration. QStash gives it its own call type (
stepConfig), keeps it out ofthe 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 featureset, 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 andlooping 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:
is publishing anyway.
Deferred submission
context.run,sleepandsleepUntilproduce their result in-process. Ratherthan 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:
The last row matters: a held result that never reached QStash would make its step
run twice.
onStepsubmits before it rethrows, and there is an end-to-end testfor it.
Reading the delivery's configuration
The effective configuration is also exposed, since the SDK had to parse it anyway:
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 offa delivery, and
describeStepSettingsMismatch, which is the comparison thedecision 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-stepmachinery. The held step is one value in two states (
held/submitting)rather than two fields that could disagree;
submitPendingStepandflushPendingStepreturn neverthrow results carrying a{ result: "submitted-step", abort } | { result: "no-pending-step" }union.src/serve/index.ts—onStepruns the route function, flushes any heldstep, then decides what to throw. Each exit is commented with when it is
reached.
src/qstash/submit-steps.ts—publishStepConfigRequest, published withcontent-based deduplication so a repeat is a no-op rather than a new message.
src/qstash/headers.ts—getHeaderstakes the step's settings and buildstheir headers in a bucket spread last, never prefixed.
src/context/context.ts,src/context/steps.ts,src/types.ts— thethird argument,
StepSettings, and the three getters.src/workflow-parser.ts,src/serve/authorization.ts— the derivedcontexts (failure, callback, authentication) get the delivery's configuration
too, instead of reporting defaults.
src/constants.ts— the new call type, the guard marker, andUpstash-Max-Retries. Note thatUpstash-Retriesmeans the retry limit whenpublishing 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 andtables, sequence diagrams for
context.run/call/waitForEvent/waitForWebhookbefore and after, the loop backstops, and the normalizationrules.
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 alocal 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: 1while the run itself is triggered withparallelism: 5. The step increments a shared counter, holds it, and decrementsit — 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 acontext.run, the one case costing no extra request.call— after acontext.call.invoke— after acontext.invoke.parallel— four steps in one parallel group, in two pairs, each pairsharing a key of
parallelism: 1, so the two keys have to hold their pairsback independently.
And two which are not about concurrency:
normalization— walks a range of settings shapes and asserts each comesback 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'sresult. 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
stepConfigentries in a row, and the entry limit) are described in the glossaryalong with what each does to the run.