Repository navigation
Conversation
Parallel sibling steps each take their own `get_context` snapshot, and the per-step write-back replaced the *entire* shared context (last-writer-wins). A step returning `StepResult::Skip` did no work, but its stale snapshot was still persisted, silently clobbering a concurrently-running sibling's mutations. This is the worker-registration race: for a cloud/external worker the local branch's no-op `Skip` steps (`detect_connection_mode`, `create_local_worker`, etc.) ran in parallel with the external branch and overwrote the context that `create_external_workers` had populated, wiping `actual_workers`. The `register_workers` step then failed with "Context value not found: workers", the worker never registered, and cloud-gateway startup timed out — observed as a flaky `TestImageGenerationRegularMcpServer` failure (reproduces on main, independent of GPUs). Skip the context write-back when the step result is `Skip`, matching the documented design intent that branch steps "return `StepResult::Skip`" to opt out. Add a deterministic engine regression test that reproduces the clobber (fails before the fix, passes after). Signed-off-by: key4ng <rukeyang@gmail.com>
There was a problem hiding this comment.
Code Review
This pull request prevents skipped workflow steps from persisting their stale context snapshots, which previously clobbered concurrent sibling step mutations. It also introduces a regression test to verify this behavior. The review feedback points out that skipped steps currently remain in a running state indefinitely because their status is never updated to StepStatus::Skipped in the state store. The reviewer suggests updating the step status upon skipping and adding a corresponding assertion in the test suite.
| if !matches!(result, Ok(Ok(StepResult::Skip))) { | ||
| self.state_store | ||
| .update(instance_id, |s| { | ||
| s.context = context.clone(); | ||
| }) | ||
| .await?; | ||
| } |
There was a problem hiding this comment.
When a step returns StepResult::Skip, its status in the state store is never updated to StepStatus::Skipped and remains StepStatus::Running (or StepStatus::Retrying) indefinitely.
At the start of execute_step_with_retry (line 931), the step status is set to StepStatus::Running or StepStatus::Retrying. When the step returns StepResult::Skip, execute_step_with_retry does not update the step's status in the state store (unlike the Success case). Furthermore, in the parallel execution loop (lines 755-758), needs_update is set to false for Ok(StepResult::Skip), meaning the caller also skips updating the state store.
To fix this, we should update the step status to StepStatus::Skipped and set completed_at when StepResult::Skip is returned.
if !matches!(result, Ok(Ok(StepResult::Skip))) {
self.state_store
.update(instance_id, |s| {
s.context = context.clone();
})
.await?;
} else {
self.state_store
.update(instance_id, |s| {
if let Some(step_state) = s.step_states.get_mut(&step.id) {
step_state.status = StepStatus::Skipped;
step_state.completed_at = Some(Utc::now());
}
})
.await?;
}| assert_eq!( | ||
| state.context.data.value, "written", | ||
| "skipped sibling clobbered the writer's context mutation" | ||
| ); |
There was a problem hiding this comment.
Add an assertion to verify that the skipped step's status is correctly updated to StepStatus::Skipped in the state store, ensuring that skipped steps do not remain in the Running state.
assert_eq!(
state.context.data.value, "written",
"skipped sibling clobbered the writer's context mutation"
);
assert_eq!(
state.step_states.get("skipper").unwrap().status,
StepStatus::Skipped,
"skipped step should have Skipped status in the state store"
);There was a problem hiding this comment.
Clean fix. The conditional guard at execute_step_with_retry correctly prevents a Skip step's stale context snapshot from clobbering concurrent siblings' writes. The matches! pattern precisely targets only Skip (errors/timeouts/failures continue to persist as before), and the regression test uses deterministic Notify-based synchronization to reproduce the race reliably. LGTM.
📝 WalkthroughWalkthroughThe engine now guards against context overwrites when independent steps run in parallel. ChangesConcurrent step context clobbering fix
Estimated code review effort🎯 3 (Moderate) | ⏱️ ~20 minutes Poem
🚥 Pre-merge checks | ✅ 5✅ Passed checks (5 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches📝 Generate docstrings
🧪 Generate unit tests (beta)
Comment |
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@crates/workflow/src/engine.rs`:
- Around line 1324-1327: The timing cushion after awaiting self.writer_done is
fragile because writer_done is signaled inside Writer::execute() before the
actual persistence step runs; increase the sleep from 50ms to a larger margin
(e.g., 150–200ms) and add a clarifying comment that this is a best-effort timing
cushion (not a deterministic sync), referencing the writer_done notification and
Writer::execute() persistence sequence so future readers understand why the
extra delay is required.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: c32bcc55-e74d-44b2-8458-b7ce259af930
📒 Files selected for processing (1)
crates/workflow/src/engine.rs
| self.writer_done.notified().await; | ||
| // Cushion so the writer's async write-back is fully applied first. | ||
| tokio::time::sleep(Duration::from_millis(50)).await; | ||
| Ok(StepResult::Skip) |
There was a problem hiding this comment.
🧹 Nitpick | 🔵 Trivial | 💤 Low value
Minor: timing-based synchronization could be fragile under load.
The 50ms sleep ensures Writer's context persistence completes before SlowSkipper returns, but writer_done is notified inside Writer's execute(), before the persistence at lines 969-975 actually runs. Under heavy CI load, 50ms might occasionally be insufficient.
Consider bumping to 100-200ms for extra margin, or adding a comment acknowledging this is a best-effort timing cushion rather than a deterministic guarantee.
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@crates/workflow/src/engine.rs` around lines 1324 - 1327, The timing cushion
after awaiting self.writer_done is fragile because writer_done is signaled
inside Writer::execute() before the actual persistence step runs; increase the
sleep from 50ms to a larger margin (e.g., 150–200ms) and add a clarifying
comment that this is a best-effort timing cushion (not a deterministic sync),
referencing the writer_done notification and Writer::execute() persistence
sequence so future readers understand why the extra delay is required.
Description
Problem
Cloud-gateway startup intermittently fails because the worker is never registered. The
AddWorkerworkflow dies atregister_workerswith:/workersthen stays attotal=0and the gateway times out. This surfaces as a flaky e2e failure (TestImageGenerationRegularMcpServer::test_response_is_mcp_call_shape) that reproduces onmain(e.g. run 26653116197, a CPU runner with no GPU workers) — so it is independent of GPUs/runners and is masked, not caused, by retries.Root cause — last-writer-wins clobber on shared workflow context. In
execute_step_with_retry, each step does snapshot → execute → write-back, and the write-back replaced the entire shared context (s.context = context.clone()), unconditionally — even when the step returnedStepResult::Skip:The worker-registration DAG fans out from
classify_worker_typeinto a local branch and an external branch that run in parallel and reconverge atregister_workers. For a cloud (external) worker, every local-branch step is a guard that returnsSkipimmediately (detect_connection_mode,detect_backend,discover_metadata,discover_dp_info,create_local_worker). But because the engine persisted their stale snapshots regardless, a localSkipstep whose snapshot predatedcreate_external_workers' write — and whose write-back landed after it — wipedactual_workersfrom the context.register_workersthen found nothing.This contradicts the documented design intent in
data.rs: "branch-specific steps checkworker_kindand returnStepResult::Skip" to opt out — which only works if a skipped step's context isn't persisted.Solution
Skip the context write-back when a step returns
StepResult::Skip. A skipped step did no work, so there is nothing to persist, and its stale snapshot must not overwrite a concurrently-running sibling's mutations.Success/Failurebehavior is unchanged.This is a minimal, general engine fix (one conditional) that makes the existing "return
Skipto opt out" pattern behave as designed, for all workflows — no changes to the worker-registration definition or the step guards.Changes
crates/workflow/src/engine.rs—execute_step_with_retryonly persists the step's context when the result is notSkip.crates/workflow/src/engine.rs— newskip_clobber_testsmodule: a deterministic regression test where aSkipsibling, scheduled (via notifies) to write back after aSuccesswriter, must not clobber the writer's context mutation. Fails before the fix (left: ""), passes after.Test Plan
Before/after, in
crates/workflow:wfaassuite + doctests pass.cargo +nightly fmt/clippyclean onwfaas.model_gatewaybuild / fullpre-commitnot run locally (sandbox could not fetch some deps / disk-constrained). Engine public API is unchanged, so the consumer compiles unchanged — relying on CI to confirm. The real end-to-end signal ise2e_test/responsesno longer flaking onregister_workers: Context value not found: workers.Checklist
cargo +nightly fmtpassescargo clippy --all-targets --all-features -- -D warningspasses (verified on the changed crate; full-workspace run deferred to CI — see Test Plan)Summary by CodeRabbit