Skip to content

fix(timing): roll back delayed SPAWN children at the sending cutoff - #50

Open
FrankD412 wants to merge 4 commits into
SemiAnalysisAI:masterfrom
FrankD412:fdinatale/port-delayed-spawn-rollback
Open

FrankD412 wants to merge 4 commits into
SemiAnalysisAI:masterfrom
FrankD412:fdinatale/port-delayed-spawn-rollback

Conversation

@FrankD412

@FrankD412 FrankD412 commented Oct 9, 2026 •

Copy link
Copy Markdown

Port of ai-dynamo#1540 (both commits, head 7fd44e9a9). The upstream PR is still open, so if it changes in review this port will be updated to match.

Why

A duration-bounded agentic run (AgentX scenario, semianalysis_cc_traces_weka_062126_256k, --benchmark-duration 3600 --benchmark-grace-period 7200) finishes every request at the duration cutoff, then idles for the whole --benchmark-grace-period before the phase completes:

Phase profiling (profiling) sending complete | sent=694, completed=693, in_flight=1 | ... | timeout_triggered=True
Waiting for event 'profiling phase credits returned' with timeout of 7199.99812217s
... done=694 ok=694 err=0 ... queue=0r/0w     <- still waiting ~2 h later, nothing in flight

Per-request metrics such as TTFT and ITL are unaffected, but the stall distorts every throughput metric. Request, output-token and total-token throughput are divided by the phase's request window, and that window stays open until the grace period expires. A stalled run therefore reports throughput understated by a factor of (duration + grace) / duration, 3× with a 3600 s duration and a 7200 s grace period, even though it completed the same requests as a run that did not stall. benchmark_duration still reports the configured duration, so nothing in the export flags the distortion. The run also holds its allocation idle for the whole grace period.

Cause. A SPAWN child whose first request comes after its branch start is dispatched later, by BranchOrchestrator._start_delayed_first_turn through the shared LoopScheduler. Its join, descendant and tree bookkeeping is registered when the child spawns. At the sending cutoff, PhaseRunner calls cancel_all_pending(), which closes the timer coroutine without running it. The child's rollback in _dispatch_first_turn_after_offset therefore never runs. So has_pending_branch_work() stays True and the runner waits out the full grace period. Cleanup then logs BranchOrchestrator leaked state at cleanup.

What

  • src/aiperf/timing/branch_orchestrator.py: track delayed children in _pending_delayed_children from scheduling until they dispatch. PhaseRunner already calls expire_replay_deadlines() right after cancel_all_pending(). That method now first runs _abandon_delayed_children(), which rolls back every child that never dispatched. It uses the same _rollback_failed_first_turn + _finalize_failed_dispatches path a refused dispatch after the cutoff takes. These children count as children_truncated, and the joins waiting on them complete. Children already in flight are unchanged, and so is the accelerated-warmup -> profiling handoff (preserve_branch_handoff skips expire_replay_deadlines()).
  • tests/unit/timing/test_branch_orchestrator_dispatch_offset.py: 5 tests using the real LoopScheduler, cancelling and then expiring in the same order as PhaseRunner.
  • tests/integration/test_weka_flat_split_e2e.py: test_child_delayed_past_duration_cutoff_does_not_hold_phase_for_grace.
  • docs/benchmark-modes/dag.md: the branch-stats paragraph now covers delayed children cut off at the duration limit.

The upstream diff applies cleanly here. This repo's branch_orchestrator.py has diverged from aiperf main, but every method the fix calls exists and behaves the same way here.

Second commit: children whose next turn is cancelled at the cutoff

The first commit was not enough. On the SemiAnalysis cluster (Qwen3.8-Flash-Next-NVFP4, AgentX scenario, conc 1, duration 3600 s, grace 7200 s), 2 of 3 runs with it still waited out the full grace period. Its rollback did fire (children_truncated=17, the same 17 delayed children), but one child still leaked: leaked state at cleanup: 0 active_joins, 0 future_joins, 1 tracked children, 1 parents with descendants. Without the fix, all 3 runs stalled with 17-19 leaked children.

Cause. Three more call sites put a DAG child's next turn on the shared LoopScheduler, and none of them tell the orchestrator if the cutoff cancels it:

Site What it schedules
agentic_replay._dispatch_next_turn a child continuation with think-time delay_ms
agentic_replay._dispatch_snapshot_for_profiling a live child seeded from the trajectory snapshot at profiling start
request_rate.handle_credit_return a child continuation with delay_ms

cancel_all_pending() closes those timers unrun, so the child never returns another credit or reaches on_child_stopped, and has_pending_branch_work() stays True.

Fix.

  • BranchOrchestrator.park_child_turn() records a child whose next turn waits on a timer.
  • The timer body claims the turn with unpark_child_turn() before dispatching.
  • expire_replay_deadlines() stops every child still parked at the cutoff through on_child_stopped, so it counts in children_truncated and its parent and tree drain.
  • Parking only happens in phases where PhaseRunner calls expire_replay_deadlines(). The accelerated-warmup handoff, where that call is skipped, never parks.

Tests.

  • 3 orchestrator unit tests: stop at the cutoff, release of a gated join, and no stop once the timer has fired.
  • Strategy unit tests for both continuation sites.
  • Two integration tests in test_weka_flat_split_e2e.py:
    • test_child_continuation_past_duration_cutoff_does_not_hold_phase_for_grace: a worker whose turn 1 is recorded about 4000 s after turn 0.
    • test_snapshot_seeded_child_turn_past_duration_cutoff_does_not_hold_phase_for_grace: an AgentX trajectory sampled at t* ≈ 0.75 s, so profiling seeds the worker's turn 1 about 4000 s out.
Check Result
New unit tests on the first commit only 7 fail
Continuation integration test on the first commit only elapsed=25.00s | grace_period_timeout=True, leaked state at cleanup: ... 4 tracked children
Snapshot integration test with only the seed-site change reverted elapsed=25.00s | grace_period_timeout=True, leaked state at cleanup: ... 1 tracked children, 1 parents with descendants (the cluster signature)
Both integration tests with the fix elapsed=5.00s, no warning
Mutation: orchestrator stop loop disabled the 2 orchestrator cutoff tests fail
tests/unit/{timing,credit,common,records} 4338 passed. 1 failure (test_records_manager.py::test_completed_excludes_analyzer_metrics, UnknownScenarioError) also fails on agentx-harness master 89b2186
test_weka_flat_split_e2e.py (-m integration) 16 passed
pre-commit on changed files clean

Not yet verified: a rerun of the cluster study with this commit.

Verification of the first commit (macOS)

Check Result
New unit tests without the branch_orchestrator.py change 4 of 5 fail (as on aiperf main)
New integration test without the change fails, children_truncated=0
tests/unit/timing/ with the change 1288 passed
tests/integration/test_weka_flat_split_e2e.py (-m integration) 14 passed
ruff check / ruff format --check on changed files clean

Not run locally: the full tests/unit/ suite, or a real GPU rerun.


Note

Medium Risk
Changes phase-completion and DAG drain logic in BranchOrchestrator and replay strategies; incorrect rollback could truncate children early or leave joins stuck, but behavior is heavily tested at cutoff boundaries.

Overview
Fixes duration-bounded DAG runs stalling through the full --benchmark-grace-period when cancelled scheduler timers left phantom branch work (has_pending_branch_work() stayed true).

BranchOrchestrator now tracks delayed SPAWN turn-0 children in _pending_delayed_children and rolls them back via _abandon_delayed_children() when expire_replay_deadlines() runs at the sending cutoff (same path as a refused dispatch: children_truncated, joins drain). Delayed-dispatch fallback tasks are keyed by child id so cutoff cancellation does not disturb in-flight dispatches. park_child_turn / unpark_child_turn cover children whose next turn waits on a think-time or snapshot-seed timer; cutoff stops still-parked children through on_child_stopped.

AgenticReplayStrategy and RequestRateStrategy park delayed DAG child continuations before schedule_later and only dispatch if unpark_child_turn succeeds.

Docs clarify that cutoff-abandoned delayed children count as children_truncated. New unit and Weka e2e tests assert phases finish at the duration cutoff without grace timeouts or leaked orchestrator state.

Reviewed by Cursor Bugbot for commit 066fec3. Bugbot is set up for automated code reviews on this repo. Configure here.

Port of ai-dynamo#1540 (head 0ae0014).

A duration-bounded agentic run could finish every request at the
--benchmark-duration cutoff and then idle for the whole
--benchmark-grace-period. Delayed SPAWN children register their join,
descendant and tree bookkeeping at spawn time; PhaseRunner's
cancel_all_pending() at the sending cutoff closes their timer coroutines
without running them, so the rollback never ran and
has_pending_branch_work() stayed True until the grace period expired.

Track delayed children until they dispatch and roll back the ones that
never did in expire_replay_deadlines(), via the same
_rollback_failed_first_turn + _finalize_failed_dispatches path a
post-cutoff dispatch refusal takes. They count as children_truncated and
the joins they gated drain.

Signed-off-by: Francesco Di Natale <3429989+FrankD412@users.noreply.github.com>
@github-actions github-actions Bot added the fix label Oct 9, 2026
@github-actions

github-actions Bot commented Oct 9, 2026 •

Copy link
Copy Markdown

Try out this PR

Quick install:

pip install --upgrade --force-reinstall git+https://github.com/ai-dynamo/aiperf.git@066fec3d4b351e62ee519f035ad08249ccdd6350

Recommended with virtual environment (using uv):

uv venv --python 3.12 && source .venv/bin/activate
uv pip install --upgrade --force-reinstall git+https://github.com/ai-dynamo/aiperf.git@066fec3d4b351e62ee519f035ad08249ccdd6350

Last updated for commit: 066fec3 • Browse code

…g cutoff

The delayed-turn-0 rollback did not cover every child the sending cutoff
can strand. Runs with that fix still waited out the whole grace period
with one leaked child: `leaked state at cleanup: 0 active_joins,
0 future_joins, 1 tracked children, 1 parents with descendants`.

Three more call sites put a DAG child's next turn on the shared
LoopScheduler:

- agentic_replay._dispatch_next_turn: a child continuation with
  think-time delay_ms
- agentic_replay._dispatch_snapshot_for_profiling: a live child seeded
  from the trajectory snapshot at profiling start
- request_rate.handle_credit_return: a child continuation with delay_ms

PhaseRunner's cancel_all_pending() closes those timers unrun, so the
child never returns another credit or reaches on_child_stopped, and
has_pending_branch_work() stays True.

BranchOrchestrator now records these children with park_child_turn().
The timer body claims its turn with unpark_child_turn() before
dispatching. expire_replay_deadlines() stops every child still parked
at the cutoff via on_child_stopped, so it counts in children_truncated
and its parent and tree drain. Parking only happens in phases where
PhaseRunner calls expire_replay_deadlines() (it is skipped only for the
accelerated-warmup handoff, which never parks).

Signed-off-by: Frank Di Natale <3429989+FrankD412@users.noreply.github.com>
Signed-off-by: Francesco Di Natale <3429989+FrankD412@users.noreply.github.com>
Without a shared scheduler, _abandon_delayed_children cancelled every
task in _delayed_dispatch_tasks whenever any delayed child was still
pending. A task that had already popped its child from
_pending_delayed_children and was awaiting dispatch_first_turn was
cancelled too, so neither its own refusal rollback nor the abandon
rollback covered that child, and its bookkeeping held the phase until
the grace period expired.

Key the fallback tasks by child x_correlation_id and cancel only the
tasks of children being rolled back. A task already dispatching
finishes and settles its own result. Production always passes the
shared scheduler, so this only affects isolated callers.

Also drop review-flagged docstrings that restated test and helper
names.

Signed-off-by: Francesco Di Natale <3429989+FrankD412@users.noreply.github.com>
At the start of profiling, _dispatch_snapshot_for_profiling restores
each live child stream from the warmup snapshot and sent its next turn
through issue_credit, ignoring the result. issue_credit's False covers
both "refused" and "issued, was the final credit", so a child refused
after leaving the parked set (for example by --request-count or the
duration cutoff) was never reported to the orchestrator: its parent's
join stayed pending and the phase waited out the grace period.

Route both the delayed and the immediate snapshot child paths through
_issue_child_continuation_or_drain, the dispatcher normal child
continuations already use. It surfaces REJECTED explicitly and calls
on_child_stopped (counted in children_truncated, join released), keeps
DEFERRED admissions, and registers the replay barrier's refusal
cleanup. Root dispatch is unchanged, and request and conversation
counting are unchanged: the counter ignores counts_toward_phase_target
for child turns.

Signed-off-by: Francesco Di Natale <3429989+FrankD412@users.noreply.github.com>
@FrankD412
FrankD412 marked this pull request as ready for review October 10, 2026 03:41

This branch has not been deployed

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

Labels

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant