Skip to content

perf(engine): raise GC thresholds in the EngineCore and ModelRunner too - #1980

Merged
valarLip merged 1 commit into
mainfrom
jun/gc_worker
Aug 23, 2026
Merged

valarLip merged 1 commit into
mainfrom
jun/gc_worker

Conversation

@junhaha666

@junhaha666 junhaha666 commented Aug 21, 2026 •

Copy link
Copy Markdown
Collaborator

perf(engine): raise GC thresholds in the EngineCore and ModelRunner too

What

ATOM_GC_THRESHOLD only ever reached the API server. _tune_gc() was called
from its FastAPI lifespan, but the EngineCore and ModelRunner processes are
spawned -- a fresh interpreter each, created before the lifespan ever runs -- so
they stayed on CPython's default 700,10,10. GC thresholds are per-interpreter,
so this moves the helper to atom.utils.tune_gc and has every process that runs
a hot loop call it.

Five files, +55/-32. Inert unless ATOM_GC_THRESHOLD is set, exactly as before.

Why

At c=2048 the GPUs go idle for a few hundred milliseconds every few hundred
decode steps, on all four DP ranks at once, with no event of any kind in the
trace for the whole window -- not a HIP call, not a gloo barrier. The period
tracks step count rather than wall clock, which points at allocation. It is a
gen-2 collection in the EngineCore: that process stops, the scheduler stops with
it, and the ModelRunner workers have nothing to run.

Raising the thresholds is close to free here. gen-2 is stop-the-world and
rescans every tracked container, so its cost tracks live objects, and in these
processes almost everything alive is the model, the compiled graph and the
tokenizer -- none of it collectable. Reference counting does the real work; the
collector only breaks cycles.

Measurements

DeepSeek-V4-Flash, tp=4, --enable-dp-attention (dp=4, 512 seqs/rank), ISL/OSL
1024/1024, c=2048, 4096 requests. Baseline is this PR's parent with
ATOM_GC_THRESHOLD already set
, so the API server alone is tuned -- which is
all the current code can do.

api_server only all three
output throughput 29009 30018 tok/s
duration 144.6s 139.7s
GPU busy 91.9% 94.2%
idle per rank 8.2s 5.2s
gaps in 100-700ms, 4 ranks 56 / 21.3s 0

What is left afterwards is one wave-boundary bubble per rank, where the engine
has drained and is waiting on the closed-loop client. It opens on
dummy_decode[bs=1] and closes on prefill[bs=1], not on a resumed decode, so
it is a supply gap and not a pause -- and not something this PR can address.

How the attribution was reached

This needed an ablation rather than a trace, and the first attempt got it wrong,
so it is worth spelling out.

The torch profiler only covers the ModelRunner workers. A pause in the EngineCore
stalls the scheduler, so the workers sit idle waiting for the next step -- which
in a worker-only trace is indistinguishable from a pause in the worker
itself, since neither one emits an event. Reading the trace alone, the obvious
conclusion (the worker's own GC) is the wrong one.

Ablated with only the API server tuned as the baseline, gaps of 400-700ms summed
over four ranks:

EngineCore worker freezes
default default many
default tuned 9.25s — workers alone change nothing
tuned default 2.14s — this is the one that matters
tuned tuned 0s

The worker call covers one freeze per rank at the wave boundary, where the
prefill burst churns enough objects to trigger gen-2 there.

Not in this PR, but worth flagging

ATOM_GC_THRESHOLD has no default, and leaving it unset is expensive on its own:
18871 tok/s at 73.5% GPU busy, against 29009 at 91.9% with it set. Any
deployment not setting it explicitly is in the first state. Giving it a non-empty
default would change GC behaviour for every deployment, so it belongs in its own
change -- but it is a bigger lever than this PR.

Risk

  • No behaviour change when ATOM_GC_THRESHOLD is unset (tune_gc() returns on
    its first line).
  • Same value applies to all three process kinds. Their object profiles differ, so
    a future split into separate variables is plausible; there was no evidence for
    it here.
  • Long-running behaviour is unverified. The observation window was ~5 minutes of
    load. Delaying collection does not accumulate in these runs -- splitting each
    run into thirds shows no growth in gap size toward the end, and with the
    thresholds raised gen-2 stops firing entirely rather than firing rarely and
    expensively -- but a 30-minute soak watching the RSS slope has not been done.

@github-actions

Copy link
Copy Markdown
Contributor

🏷️ CI Guide

Runs automatically on every eligible PR before approval:

  • ✅ Pre Checkin: Black, Ruff, catalog schema validation, non-GPU unit tests

Heavy model tests:

  • ✅ Run after the PR is approved and Pre Checkin passes
  • ✅ Run immediately when an approval review is submitted
  • ✅ Can be requested before approval with labels
Label Tests
ci:full Run all heavy PR model tests: native ATOM, vLLM, and SGLang
ci:atom Run native ATOM model accuracy tests
ci:vllm Run ATOM vLLM OOT model accuracy tests
ci:sglang Run ATOM SGLang model accuracy tests

Heavy jobs are skipped when the PR is not approved and no matching ci:* label is present.
Add labels via the sidebar or gh pr edit 1980 --add-label <label>

At c=2048 the GPUs go idle for a few hundred milliseconds every few hundred
decode steps, on all four DP ranks at once, with no event of any kind in
the trace for the whole window -- not a HIP call, not a gloo barrier. The
period tracks step count rather than wall clock, which points at
allocation, and the cause is a gen-2 collection in the EngineCore.

_tune_gc lives in the API server's FastAPI lifespan. Subprocesses are
spawned, so they get a fresh interpreter, and they are created before the
lifespan ever runs -- ATOM_GC_THRESHOLD never reaches them and they stay on
CPython's 700,10,10. Thresholds are per-interpreter, so the function moves
to atom.utils and every process that runs a hot loop calls it.

Attribution needed an ablation rather than a trace. The torch profiler only
covers the ModelRunner workers, and a pause in the EngineCore stalls the
scheduler, so the workers sit idle waiting for the next step -- in a
worker-only trace that is indistinguishable from a pause in the worker
itself, since neither emits an event. Ablated with only the API server
tuned as the baseline, gaps of 400-700ms summed over four ranks: workers
only 9.25s (no effect at all), EngineCore only 2.14s, both 0s. The worker
call covers one freeze per rank at the wave boundary, where the prefill
burst churns enough objects to trigger gen-2 there.

Measured on DeepSeek-V4-Flash tp4 dp4, c=2048, 4096 requests, against this
commit's parent with ATOM_GC_THRESHOLD already set, so the API server alone
is tuned -- which is all the current code can do:

                              api_server only   all three
  output throughput             29009            30018 tok/s
  duration                      144.6s           139.7s
  GPU busy                       91.9%            94.2%
  idle per rank                   8.2s             5.2s
  gaps in 100-700ms, 4 ranks   56 / 21.3s        0

What is left in the trace afterwards is one wave-boundary bubble per rank,
where the engine has drained and is waiting on the closed-loop client: it
opens on dummy_decode[bs=1] and closes on prefill[bs=1], not on a resumed
decode, so it is a supply gap and not a pause.

Inert unless ATOM_GC_THRESHOLD is set, as before. Worth noting separately
that leaving it unset is expensive on its own: 18871 tok/s at 73.5% GPU
busy, against 29009 at 91.9% with it set.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
@valarLip
valarLip merged commit 84157ef into main Aug 23, 2026
34 of 51 checks passed
@valarLip
valarLip deleted the jun/gc_worker branch August 23, 2026 04:48
valarLip added a commit that referenced this pull request Aug 23, 2026
…this

Enumerating which processes get the GC policy has now gone stale twice. #1980
reached the API server but not the disaggregated EngineCores, whose
`run_engine` does not call the base one. The first version of this branch
reached those but not atomesh, which builds its own engine in
`launch_atom_standalone` and never runs the FastAPI lifespan the API server
applies the policy from -- so that frontend had neither `tune_gc` (since
#1980) nor the freeze, while the EngineCore and workers it spawns had both.

It is the API server's counterpart, with the same profile: the tokenizer, the
per-request accumulators, ~632k tracked objects and a 265 ms gen-2 pause.

Rather than fix the third instance and wait for a fourth, the rule is now
tested: `test_every_serving_frontend_applies_the_gc_policy` walks every module
that builds an engine and serves, and fails naming the ones that do not apply
the policy. Written the obvious way it was vacuous -- an unused import
satisfies a search for the name -- so it requires the call.

The prose enumerations in `tune_gc` and the env-var docs are replaced by the
rule they kept failing to track.
valarLip added a commit that referenced this pull request Aug 23, 2026
)

* perf(engine): stop the collector walking a heap it never reclaims

At steady state the garbage collector in these processes finds nothing.
Measured on DeepSeek-V4-Flash-DSpark tp1 with a `gc.callbacks` probe: over
1754 gen-0, 158 gen-1 and 50 gen-2 passes after startup, it reclaimed
**zero** objects. Everything the hot loop allocates is acyclic and reference
counting takes it. What each pass does cost, being stop-the-world and
proportional to the live heap:

    EngineCore   268 ms      of 819,923 tracked
    TP0 worker   979 ms      of 959,694
    api_server   265 ms      of 632,424

Almost all of that heap is startup state -- model, compiled graph, tokenizer,
KV block pool -- so `gc.freeze()` after warmup removes the work rather than
merely spacing it out, which is all raising thresholds (#1980) can do. With it
on, gen-2 stops firing entirely in all three processes for the rest of the
run, while gen-0/1 keep going at under 2 ms: this narrows the collector, it
does not disable it, so a cycle written by later code is still caught.

`unfreeze_gc_heap` on shutdown is not optional. A frozen object is invisible
to the collector, so an engine destroyed inside a live interpreter -- a test,
an RL loop that rebuilds it -- would leave its weights and KV cache
unreachable *and* uncollectable, which presents as a GPU memory leak with
nothing about GC in the symptom. vLLM hit this and handles it the same way.

Placement is one call per process at the boundary where startup ends:

  - api_server, in the lifespan just before it yields
  - EngineCore, at the end of `__init__` -- KV cache allocated, graphs
    captured, BlockPool built (94,763 `Block` objects on a V4-Flash tp1)
  - each ModelRunner worker, over an RPC the EngineCore sends, because only
    the caller knows warmup is over: weights, compile and capture all arrive
    as RPCs it issues, and `--enforce-eager` has no capture step to hook

Disaggregated decode is the exception and is left partly covered. Its block
pool is built in `_init_disagg`, after `super().__init__()` has already sent
READY, so a freeze there could make queued requests permanent. It keeps the
base freeze and forgoes the pool.

`tune_gc` moves into the new module and stays as the fallback for
`ATOM_GC_FREEZE=0`. `ATOM_GC_DEBUG` is added alongside: it logs every
collection's generation, duration and reclaimed count, which is the only way
to see these pauses -- a stall in the EngineCore idles the workers, and an
idle worker emits no trace event at all. It is expensive (~90s of extra
startup) and off by default.

Two gaps found next door and closed, since the fix needed a shared per-process
setup anyway: `PrefillEngineCore.run_engine` and `DecodeEngineCore.run_engine`
override the base without calling it, so a disaggregated deployment got
neither `enable_orphan_reaping` nor `tune_gc`, and neither ever set a process
title -- both showed as bare `python` in `ps`. `_setup_engine_process` now
owns identity, reaping and GC policy for all three, which also collapses the
base's three `set_process_title` branches into one and makes the process title
and the GC log lines agree by construction.

Throughput is unchanged at this scale (32098 vs 32122 tok/s, and 32098 with
freezing off), because gen-2 fires roughly once per seven minutes in steady
state here. This is a tail-latency change; it pays at the concurrency #1980
measured at, which is not reproduced here.

* fix(engine): one name per worker process, not three

The worker side had the drift the EngineCore side just lost. Three places
produced its name and only one knew about data parallelism:

  - `set_process_title` in `AsyncIOProc.__init__`, dp-aware (`DP0TP0`)
  - the GC debug callback, `f"TP{rank}"`
  - `ModelRunner.freeze_gc_heap`, `f"TP{rank}"`

So under dp>1 every rank's worker logged as `TP0`, which is exactly the case
worth telling apart -- and `ps` disagreed with both log lines.

`engine_process_name` and `worker_process_name` now live next to
`set_process_title`, since naming a process is what that function is for, and
all five call sites take their string from them. `AsyncIOProc` names itself
once, before anything logs, which also retires a `try/except Exception` that
existed only to repeat the fallback the helper now expresses directly.

`EngineCore._process_name_for` moves out to join them rather than staying a
method, so a reader asking "what are these processes called" finds one answer.

* fix(entrypoints): the atomesh frontend was the third process to miss this

Enumerating which processes get the GC policy has now gone stale twice. #1980
reached the API server but not the disaggregated EngineCores, whose
`run_engine` does not call the base one. The first version of this branch
reached those but not atomesh, which builds its own engine in
`launch_atom_standalone` and never runs the FastAPI lifespan the API server
applies the policy from -- so that frontend had neither `tune_gc` (since
#1980) nor the freeze, while the EngineCore and workers it spawns had both.

It is the API server's counterpart, with the same profile: the tokenizer, the
per-request accumulators, ~632k tracked objects and a 265 ms gen-2 pause.

Rather than fix the third instance and wait for a fourth, the rule is now
tested: `test_every_serving_frontend_applies_the_gc_policy` walks every module
that builds an engine and serves, and fails naming the ones that do not apply
the policy. Written the obvious way it was vacuous -- an unused import
satisfies a search for the name -- so it requires the call.

The prose enumerations in `tune_gc` and the env-var docs are replaced by the
rule they kept failing to track.

* fix(engine): restore the decode freeze, dropped on a misreading

An earlier commit on this branch removed `_freeze_after_startup` from
`DecodeEngineCore` on the grounds that READY had already gone out, so freezing
could catch requests the input thread had queued and make them permanent.

That is not what the code does. `DecodeEngineCore.__init__` sets
`_ready_deferred = True` before calling `super().__init__()`, precisely so the
base's ready signal is suppressed; the real one is sent nine lines *after* the
block pool is built, once the kvcache IPC import and graph capture are done.
Nothing can have been admitted at that point, so there is no window and the
freeze is safe -- and it is what covers decode's 94,763 `Block` objects.

The comment left in its place asserted the false premise in so many words,
which would have taught the next reader the same wrong thing.
valarLip added a commit that referenced this pull request Aug 23, 2026
… one (#1990)

* fix(engine): a fourth consumer of token ids, and a guard for the next one

`Sequence.token_ids` became an `array("i")` in #1989. That commit named
three places that compare token ids and fixed each. Sweeping the rest of
the chain turned up a fourth consumer of a different kind -- one that
serializes rather than compares -- and it fails silently too.

`publish_loaded_prefix` hands `_hash_block_tokens`' output straight into a
`BlockStored` event, which is msgpack-encoded. msgspec has no encoding for
`array.array`, and `KVEventPublisher.publish` counts encode failures rather
than raising, so with KV events enabled the whole event stream goes dark
after one log line. The other two BlockStored sites accumulate into a list
via `extend` and were never affected.

Rather than convert at each site, the two boundaries where the type is
load-bearing now assert it:

  - `Block.update` takes an array and refuses a list. A list there costs
    twice and neither cost announces itself: it never compares equal to the
    arrays the other publish paths store, so every hit on that block reads
    as a collision; and it hands the collector one traversal slot per token.
  - `_make_block_stored` takes a list and refuses an array, for the encoder.

The rest of the chain follows the same axis: everything that grows once per
token is an array now. `Sequence.output_tokens` (the completion half of
`token_ids`, kept in step with it by every writer), `Sequence.logprobs`, the
API server's three per-request accumulators, the streaming detokenizer's
buffer, and atomesh's two. `tokenizer.decode` takes an array as is, so only
the JSON edges convert back -- and only `LLMEngine`'s, because the HTTP
response builders read `text` and the counters and never `token_ids`. A test
pins that, so a builder that starts forwarding the key fails there rather
than in production.

Four annotations said `list[int]` where an array now flows.

Measured on a live V4-Flash-DSpark tp1 (94,681 blocks) with a `gc.callbacks`
probe: the EngineCore's stop-the-world gen-2 pause is 288.6ms with list
storage and 242.8ms with arrays, and with arrays it stops growing as the
prefix cache fills. Throughput is unchanged at this scale -- gen-2 fires
about once per seven minutes in steady state here, so the pause matters at
the concurrency #1980 measured at, not at this one.

* fix(engine): two more token-id annotations that said list

`compute_hash` was annotated `list[int] | array.array` on the reasoning that
one caller still passed a list. It does not. All six production call sites
pass an `array("i")`:

  - `block_manager.py` x5, from `_hash_block_tokens` / `seq.block(i)`, both
    slices of `Sequence.token_ids`
  - `policy.py`, from `_chained_prefix_hashes(seq.token_ids, ...)`

The union came from reading `_chained_prefix_hashes(token_ids: list[int], ...)`
and taking the signature for the fact -- but that annotation is stale in the
same way, since `connector.py` hands it `seq.token_ids`. Both are narrowed to
what actually arrives.

The int64 pin stays. Its reason changes rather than disappearing: nothing
passes a list today, so it is no longer reconciling two live callers, but it is
what stops a caller who does pass one from silently computing a different
digest for the same tokens.

Only tests pass lists now, and `test_compute_hash_does_not_depend_on_its_
argument_type` keeps asserting the two agree -- the annotation says what
callers should pass, that test says the function is robust if they do not.

* perf(engine): stop the collector walking a heap it never reclaims (#1991)

* perf(engine): stop the collector walking a heap it never reclaims

At steady state the garbage collector in these processes finds nothing.
Measured on DeepSeek-V4-Flash-DSpark tp1 with a `gc.callbacks` probe: over
1754 gen-0, 158 gen-1 and 50 gen-2 passes after startup, it reclaimed
**zero** objects. Everything the hot loop allocates is acyclic and reference
counting takes it. What each pass does cost, being stop-the-world and
proportional to the live heap:

    EngineCore   268 ms      of 819,923 tracked
    TP0 worker   979 ms      of 959,694
    api_server   265 ms      of 632,424

Almost all of that heap is startup state -- model, compiled graph, tokenizer,
KV block pool -- so `gc.freeze()` after warmup removes the work rather than
merely spacing it out, which is all raising thresholds (#1980) can do. With it
on, gen-2 stops firing entirely in all three processes for the rest of the
run, while gen-0/1 keep going at under 2 ms: this narrows the collector, it
does not disable it, so a cycle written by later code is still caught.

`unfreeze_gc_heap` on shutdown is not optional. A frozen object is invisible
to the collector, so an engine destroyed inside a live interpreter -- a test,
an RL loop that rebuilds it -- would leave its weights and KV cache
unreachable *and* uncollectable, which presents as a GPU memory leak with
nothing about GC in the symptom. vLLM hit this and handles it the same way.

Placement is one call per process at the boundary where startup ends:

  - api_server, in the lifespan just before it yields
  - EngineCore, at the end of `__init__` -- KV cache allocated, graphs
    captured, BlockPool built (94,763 `Block` objects on a V4-Flash tp1)
  - each ModelRunner worker, over an RPC the EngineCore sends, because only
    the caller knows warmup is over: weights, compile and capture all arrive
    as RPCs it issues, and `--enforce-eager` has no capture step to hook

Disaggregated decode is the exception and is left partly covered. Its block
pool is built in `_init_disagg`, after `super().__init__()` has already sent
READY, so a freeze there could make queued requests permanent. It keeps the
base freeze and forgoes the pool.

`tune_gc` moves into the new module and stays as the fallback for
`ATOM_GC_FREEZE=0`. `ATOM_GC_DEBUG` is added alongside: it logs every
collection's generation, duration and reclaimed count, which is the only way
to see these pauses -- a stall in the EngineCore idles the workers, and an
idle worker emits no trace event at all. It is expensive (~90s of extra
startup) and off by default.

Two gaps found next door and closed, since the fix needed a shared per-process
setup anyway: `PrefillEngineCore.run_engine` and `DecodeEngineCore.run_engine`
override the base without calling it, so a disaggregated deployment got
neither `enable_orphan_reaping` nor `tune_gc`, and neither ever set a process
title -- both showed as bare `python` in `ps`. `_setup_engine_process` now
owns identity, reaping and GC policy for all three, which also collapses the
base's three `set_process_title` branches into one and makes the process title
and the GC log lines agree by construction.

Throughput is unchanged at this scale (32098 vs 32122 tok/s, and 32098 with
freezing off), because gen-2 fires roughly once per seven minutes in steady
state here. This is a tail-latency change; it pays at the concurrency #1980
measured at, which is not reproduced here.

* fix(engine): one name per worker process, not three

The worker side had the drift the EngineCore side just lost. Three places
produced its name and only one knew about data parallelism:

  - `set_process_title` in `AsyncIOProc.__init__`, dp-aware (`DP0TP0`)
  - the GC debug callback, `f"TP{rank}"`
  - `ModelRunner.freeze_gc_heap`, `f"TP{rank}"`

So under dp>1 every rank's worker logged as `TP0`, which is exactly the case
worth telling apart -- and `ps` disagreed with both log lines.

`engine_process_name` and `worker_process_name` now live next to
`set_process_title`, since naming a process is what that function is for, and
all five call sites take their string from them. `AsyncIOProc` names itself
once, before anything logs, which also retires a `try/except Exception` that
existed only to repeat the fallback the helper now expresses directly.

`EngineCore._process_name_for` moves out to join them rather than staying a
method, so a reader asking "what are these processes called" finds one answer.

* fix(entrypoints): the atomesh frontend was the third process to miss this

Enumerating which processes get the GC policy has now gone stale twice. #1980
reached the API server but not the disaggregated EngineCores, whose
`run_engine` does not call the base one. The first version of this branch
reached those but not atomesh, which builds its own engine in
`launch_atom_standalone` and never runs the FastAPI lifespan the API server
applies the policy from -- so that frontend had neither `tune_gc` (since
#1980) nor the freeze, while the EngineCore and workers it spawns had both.

It is the API server's counterpart, with the same profile: the tokenizer, the
per-request accumulators, ~632k tracked objects and a 265 ms gen-2 pause.

Rather than fix the third instance and wait for a fourth, the rule is now
tested: `test_every_serving_frontend_applies_the_gc_policy` walks every module
that builds an engine and serves, and fails naming the ones that do not apply
the policy. Written the obvious way it was vacuous -- an unused import
satisfies a search for the name -- so it requires the call.

The prose enumerations in `tune_gc` and the env-var docs are replaced by the
rule they kept failing to track.

* fix(engine): restore the decode freeze, dropped on a misreading

An earlier commit on this branch removed `_freeze_after_startup` from
`DecodeEngineCore` on the grounds that READY had already gone out, so freezing
could catch requests the input thread had queued and make them permanent.

That is not what the code does. `DecodeEngineCore.__init__` sets
`_ready_deferred = True` before calling `super().__init__()`, precisely so the
base's ready signal is suppressed; the real one is sent nine lines *after* the
block pool is built, once the kvcache IPC import and graph capture are done.
Nothing can have been admitted at that point, so there is no window and the
freeze is safe -- and it is what covers decode's 94,763 `Block` objects.

The comment left in its place asserted the false premise in so many words,
which would have taught the next reader the same wrong thing.

* style(entrypoints): make atomesh/server.py clean, not just no-worse

CI runs ruff through reviewdog with `-filter-mode=diff_context`, so it reports
anything landing on a changed line or in the context around it -- pre-existing
or not. Adding an import inside this file's already-unsorted block pulled a
long-standing `I001` into review and failed the check.

Comparing whole-file findings against `main` is the wrong local check for
that: identical counts still fail if one of them is near your edit. Cheaper to
leave the file with none.

Sorting the imports is the fix CI asked for. The other two go with it:
`initialize_standalone_service` declared `global engine, tokenizer` but only
reads them, which `global` is not for; and the blind except around the version
banner gets the reason it always had.

* fix(tests): the worker-RPC contract test needed no aiter, and no import

CI's non-GPU runner has no aiter, and `atom.model_engine.model_runner` imports
it. Importing `ModelRunner` inside the test to read `freeze_gc_heap` off it
therefore failed there with `ImportError: cannot import name
destroy_dist_env`.

The import-without-aiter probe missed it because that probe models collection:
it imports each test module and stops. An import inside a test body only runs
when the test does.

The contract is about what the source says -- an RPC target must return
something or `call_func(..., wait_out=True)` hangs -- so it is now read with
`ast` straight from the file. No import, no aiter, no GPU, and the whole module
drops from 5.7s to 1.6s. Re-verified by running every test file this branch
touches with `aiter` and `triton` blocked at the import hook: 142 passed.
valarLip pushed a commit that referenced this pull request Aug 26, 2026
…2020)

* perf(frontend): merge_chunk rebuilt the whole backlog on every merge

`StreamOutputCollector` folds an arriving chunk into the one still waiting
whenever the consumer lags, and `merge_chunk` rebuilt `token_ids` and
re-concatenated `text` from scratch each time. That is quadratic in the
backlog depth, and it runs on the event loop: the further the frontend falls
behind the engine the deeper the backlog gets, the more every merge costs,
and the further behind it falls. Past the point where the engine outruns
delivery the frontend does not degrade, it collapses.

The first chunk's list is the engine's own `output_tokens` and still must not
be appended to, so the copy happens once, into a list this module owns and
may extend in place from then on. `text` is popped before the concatenation
so the string has a single reference and CPython resizes it in place instead
of allocating a new one per merge.

DSR1 fp4, tp8 with DP attention, c=8192, ISL 1024 / OSL 4096, one wave,
otherwise identical builds:

                        before      after
  output throughput     42,970  ->  57,735 tok/s   (+34.4%)
  total throughput      53,713  ->  72,169 tok/s
  mean TPOT              179.9  ->   130.9 ms
  median ITL            17,475  ->     912 ms

The engine's last decode step lands at 568.6s in both; the run used to end at
793s and now ends at 581s, so what the quadratic term was buying was 225
seconds of the frontend draining a backlog after the GPUs had gone idle.

Same numbers seen from the other side: reverting the engine-side GC tuning
from #1980 also gains 10.2% on this config, because a slower engine keeps the
backlog shallow enough that the quadratic term stays cheap. That is worth
noting only as evidence for the mechanism -- the fix belongs here, not there.

gsm8k 3-shot 0.9507 flexible / 0.9492 strict, 1319/1319. Merge output is
byte-identical to the unmerged path at backlog depths from 1 to 500, and to
the previous implementation over 4000 randomized merge sequences, with the
engine's list never mutated.

* perf(frontend): copy in put_nowait, not around a premise that does not hold

Review found the premise false. The docstring this change inherited says the
first chunk's list is the engine's own `output_tokens` and must not be appended
to; `scheduler.py` is the only `RequestOutput` site in the repo and already
hands out a per-step copy five lines above it, and the object is unpickled
fresh after crossing ZMQ besides. The ownership tag, the exact-class test and
the conditional copy were all guarding a list nothing else reads.

The copy is kept as defence in depth -- the collector should not extend a list
handed to it across two files and a process boundary on the strength of a
comment -- but it moves to `put_nowait`, which is where the producer's chunk is
first taken and already the only caller. That is one copy per stall rather than
a branch and a type test per merge, and it is faster than the previous
implementation at every backlog depth rather than from depth 20 on.

Three behaviours the old rebuild gave away for free and the tag version had
dropped:

  - a delta with no text no longer loses the text already accumulated. The
    string is popped to keep its refcount at one, so the write-back is in a
    `finally` and a `TypeError` from the concatenation leaves `into` intact.
  - `token_ids` is de-aliased from the producer's list on the way in now,
    unconditionally, rather than on the first merge that carries ids.
  - a chunk that never went through `put_nowait` now raises `KeyError` in
    `merge_chunk` instead of silently growing a key. The precondition is
    pinned by a test rather than left to be rediscovered.

The docstring said "rebuilt rather than extended" over code that extends; it is
rewritten. The pop carries a comment, since restoring `into["text"] = a + b`
reads as a cleanup and brings the quadratic term back with no test or lint
signal.

Tests: the harness the first commit described now exists. Two of the new cases
fail with the source change reverted, which the suite could not do before.

* style(frontend): trim the comments to the why
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.

2 participants