Repository navigation
fix(vllm): propagate FPM worker_id into snapshot-restored EngineCore - #12545
Conversation
|
👋 Hi RealNicolasBourbaki! Thank you for contributing to ai-dynamo/dynamo. Just a reminder: The 🚀 |
This comment has been minimized.
This comment has been minimized.
8ae8825 to
e4cfc65
Compare
e4cfc65 to
75a3499
Compare
WalkthroughAdds a fallback ChangesFPM worker ID synchronization
Estimated code review effort: 3 (Moderate) | ~20 minutes 🚥 Pre-merge checks | ✅ 5✅ Passed checks (5 passed)
✨ Finishing Touches 💡 1🛠️ Fix failing CI checks 💡
Comment |
There was a problem hiding this comment.
Actionable comments posted: 2
🧹 Nitpick comments (1)
components/src/dynamo/vllm/instrumented_scheduler.py (1)
3672-3679: 📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick winReplace the defensive
getattrwith direct attribute access.
EngineCorealways assignsself.schedulerin its constructor, and the utility RPC only runs on a constructed engine. Direct access keeps the failure loud if vLLM ever renames the attribute. TheRuntimeErrorstill covers the foreign-scheduler case.♻️ Proposed refactor
def set_fpm_worker_id(self, new_worker_id: str) -> None: - scheduler = getattr(self, "scheduler", None) + """Update the FPM worker id on the scheduler and its publisher.""" + scheduler = self.scheduler if not isinstance(scheduler, InstrumentedScheduler): raise RuntimeError( f"scheduler is {type(scheduler).__name__}, not InstrumentedScheduler" )Note:
test_fpm_utility_rejects_non_instrumented_schedulerpasses aSimpleNamespace(scheduler=None), so themissingcase still raisesRuntimeErrorafter this change.As per coding guidelines: "Do not use defensive
getattr(obj, "attr", default)when the object's type is known and the attribute is part of its definition; use direct attribute access instead."🤖 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 `@components/src/dynamo/vllm/instrumented_scheduler.py` around lines 3672 - 3679, Update set_fpm_worker_id to access self.scheduler directly instead of using getattr with a default. Retain the existing isinstance validation and RuntimeError for non-InstrumentedScheduler values, including scheduler=None, then continue updating _fpm_worker_id and _publisher._worker_id.Sources: Coding guidelines, Path instructions
🤖 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 `@components/src/dynamo/vllm/worker_factory.py`:
- Around line 1018-1021: Update the FPM relay setup comments in
_create_decode_worker and _create_prefill_worker to reflect that the worker ID
is pushed to the restored child after synchronization, and is empty only if the
utility call fails. Do not retain the outdated checkpoint-mode explanation.
- Around line 71-77: Update _sync_fpm_worker_id to wrap call_utility_async in
asyncio.wait_for using the established bounded timeout, while allowing
asyncio.CancelledError to propagate. In the exception logging path, include
new_worker_id with lazy logger formatting and retain exc_info=True.
---
Nitpick comments:
In `@components/src/dynamo/vllm/instrumented_scheduler.py`:
- Around line 3672-3679: Update set_fpm_worker_id to access self.scheduler
directly instead of using getattr with a default. Retain the existing isinstance
validation and RuntimeError for non-InstrumentedScheduler values, including
scheduler=None, then continue updating _fpm_worker_id and _publisher._worker_id.
🪄 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: Path: .coderabbit.yaml
Review profile: CHILL
Plan: Enterprise
Run ID: 1d6ad441-421f-4e9d-9247-52eccb87a5f7
📒 Files selected for processing (5)
components/src/dynamo/vllm/instrumented_scheduler.pycomponents/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.pycomponents/src/dynamo/vllm/tests/test_vllm_worker_factory.pycomponents/src/dynamo/vllm/worker_factory.pytests/report_pytest_markers.py
Automated evidence record — validation completeValidation status: complete Evidence summary: [5/5 validated] AI review assessment: sound. This assessment is advisory only — it is machine-generated and does not substitute for human review or for a maintainer's approval. Validation result: complete — pass. Evidence audit: complete [5/5 validated] — the evidence table is grounded in recorded runs. Evidence [5/5 validated]Generated from validation/registry.jsonl — do not edit by hand.
|
plan.md# Plan — work item 12591: vLLM snapshot-restored workers publish FPM with an empty `worker_id`
Route: review
Template: existing-request
Review request: https://github.com/ai-dynamo/dynamo/pull/12545
Engine: vllm
## Disposition
No `Disposition:` line is written. The planner brief permits
`Disposition: already-resolved` only when provider metadata proves a **merged**
PR carries the behavior on the target branch, and explicitly forbids it for an
open PR or a partial implementation. PR #12545 is `OPEN`, `mergeable:
CONFLICTING`, and `reviewDecision: REVIEW_REQUIRED`; issue #12591 is still
`OPEN` with `closedAt: null`. Nothing about this fix exists on `main` today —
I verified the defective code is still present at
`components/src/dynamo/vllm/worker_factory.py`, including the very TODO the
issue quotes. The work item is real and unresolved; it is simply already being
solved by somebody else, in the open, with an approach that holds up.
## Note on the template deviation
The lead handed me `debug-investigation`, and my caller instruction was to keep
it "unless discovery proves an existing PR fully implements the change."
Discovery did prove that, in substance: PR #12545 patches all three
snapshot-restore sites that exist in this checkout, updates both of the two
identity fields that FPM payloads are stamped from, leaves the cold-start path
untouched, and ships unit tests. Holding `debug-investigation` would direct the
printer to author a competing fix over the same four files — precisely the
outcome the planner brief tells me to avoid ("Opening a second one over the same
files just leaves a maintainer to work out which of the two to take"). The
`existing-request` overlay is also the only one whose printer contract says to
begin at the supplied request's head, preserve attribution, and reconcile with a
newer base by merging rather than rebasing — which is exactly the mechanic this
work item needs. The two overlays require the same five packets, so nothing
downstream is disturbed by the swap. I record the deviation openly rather than
making it silently.
## User intent
The caller wants snapshot-restored vLLM workers to publish forward-pass metrics
under the same `worker_id` they register in MDC, so the planner stops freezing
on `worker_count_mismatch` and can scale replicas back down when load drops.
Read plainly, that is an ordinary "fix this bug" request. The caller supplied no
review request and recorded none; the prior factory attempt on this work item
died during planning and left nothing behind. Discovery then found that a human
contributor already has an open, sound, near-complete fix for exactly this bug.
The intent is therefore best served by getting **that** fix over the line rather
than by racing it.
## Non-goals
- Not opening a second pull request for issue #12591. One already exists and is
approach-correct.
- Not rewriting #12545's history. A push into its branch must be a
fast-forward, so reconciliation with the newer base is a merge, never a
rebase.
- Not changing the planner side. `_reconcile_fpm_worker_count`
(`components/src/dynamo/planner/core/state_machine.py:465`) is behaving
correctly given the input it is handed; the defect is upstream of it. Relaxing
that check would suppress the symptom rather than fix the cause.
- Not changing the `known_workers` filter in
`lib/bindings/python/rust/llm/fpm.rs`. It is doing its job.
- Not attempting the end-to-end capture → restore → scale-down loop. That is a
multi-component Kubernetes behavior (snapshot-agent, planner, several DGD
replicas) and this sandbox has one GPU, no Docker socket, and no cluster.
- Not touching SGLang or TensorRT-LLM. `$GLAMR_CODING_ENGINE` is `vllm`.
## Discovery
Everything below is something I read in this checkout
(`/home/sandbox/workspace/wi-20260814T172511Z-12591/repo`, `main` at
`c7aa0bf99ce91363700fc85edab418d3166a5943`) or fetched with `gh`. Network
egress and `gh` authentication both worked; no host call was denied, so every
negative recorded here is a checked negative, not an unverified one.
### The reporter's root cause, verified — with corrected line numbers
The issue's line numbers are stale. The mechanism is right; the coordinates are
not.
- `components/src/dynamo/vllm/instrumented_scheduler.py:1470` (issue said 1220)
— `self._fpm_worker_id = os.environ.get(ENV_FPM_WORKER_ID, "")`, inside
`InstrumentedScheduler.__init__`, which begins at line 1447. The env var name
`DYN_FPM_WORKER_ID` is defined at line 134.
- That value is then copied a second time, at line 1487, into
`_FpmPublisherThread(worker_id=...)`. **This is the detail that matters most
and that the issue does not mention: there are two independent identity
fields, not one.** `_FpmPublisherThread.__init__` stores it as `self._worker_id`
(line 1368) and stamps it onto every *idle heartbeat* at line 1423, while
`InstrumentedScheduler._extract_metrics` stamps `self._fpm_worker_id` onto
every *active* sample at line 1719. Any fix that updates only one of the two
leaves an idle replica invisible — and an idle replica is exactly the
scale-down case the issue is about.
- `components/src/dynamo/vllm/snapshot.py` — 53 lines, read in full. Line 41
calls `setup_vllm_engine(config)`, which spawns the EngineCore child and
constructs the scheduler; only afterwards does line 49 call
`snapshot_controller.wait_for_restore()`. Confirms the ordering the issue
claims: the child exists before the Dynamo runtime does.
- `components/src/dynamo/vllm/worker_factory.py` — there are exactly **three**
`if snapshot_engine is not None:` branches, at lines 678, 1130, and 1419
(`_create_realtime_worker`, `_create_decode_worker`, `_create_prefill_worker`).
Each computes `fpm_worker_id = str(generate_endpoint.connection_id())` and then
assigns `os.environ[ENV_FPM_WORKER_ID] = fpm_worker_id` at lines 686, 1138,
and 1431 — into the *parent's* environment, after the child has already
forked. The TODO the issue quotes is still there, immediately above line 1431
(issue said 1284).
- `components/src/dynamo/vllm/main.py:605-606` — the cold-start path sets
`os.environ[ENV_FPM_WORKER_ID]` inside `setup_vllm_engine`, *before* engine
construction. This confirms cold start is genuinely unaffected and that any
fix must be scoped to the snapshot branch only.
### The downstream mechanism, traced end to end
The issue attributes the freeze to the count check alone. The real chain has an
earlier link, which I found by following the consumer:
1. `lib/bindings/python/rust/llm/fpm.rs:900` — `get_recent_stats()` filters
`latest_stats` against `known_workers`:
`.filter(|entry| self.known_workers.contains(&entry.key().0))`. `known_workers`
is maintained from the MDC discovery watch (insert on Added, remove on
Removed — see the Task 2 comment at line ~757). A restored worker registers in
MDC under its real `connection_id`, but publishes FPM under `""`. `""` is never
in `known_workers`, so **every sample from a restored worker is discarded at
read time**, before the planner ever sees it.
2. `components/src/dynamo/planner/environment/metrics_provider/runtime_provider.py:187`
consumes `subscriber.get_recent_stats()` and keys by `(worker_id, dp_rank)`.
3. `components/src/dynamo/planner/core/state_machine.py:465-489` —
`_reconcile_fpm_worker_count` then sees fewer distinct worker ids than
`dgd_count`, logs `Worker count mismatch: DGD=..., FPM=...`, and returns
`False`.
4. `components/src/dynamo/planner/core/load_scaling.py:75` (and 137, 144, 276) —
the caller sets `self._diag_load_reason = "worker_count_mismatch"` and
returns `None`, so no scaling decision is produced at all.
This also explains the asymmetry the reporter observed ("scale-up works,
scale-down never happens") more precisely than the issue does: at one replica the
counts can still line up, so the planner can act; once it has scaled up, the
restored replicas contribute nothing to `fpm_stats` and the count can never match
again, so the deployment is stuck at its high-water mark. Worth stating in the
review, because it tells a maintainer why the bug looks intermittent.
### The existing pull request — the decisive find
`gh pr list --state all --search "12591"` returned exactly one hit, and the same
PR was the top hit for every semantic search I ran (`DYN_FPM_WORKER_ID`,
`_fpm_worker_id`, `InstrumentedScheduler`, `snapshot worker_id`, `FPM worker id`,
`snapshot restore FPM`):
- **https://github.com/ai-dynamo/dynamo/pull/12545** — *"fix(vllm): propagate FPM
worker_id into snapshot-restored EngineCore"*
- Author `RealNicolasBourbaki` (Nianheng Nicole Wu), `is_bot: false`. A human
contributor, **not** this factory: the branch is `fix-fpm-worker-id`, which
does not match the factory's `<type>/<slug>--<hash>` convention, and
`input.md` records no review requests.
- State `OPEN`, base `main`, not a draft, body says `Closes #12591`.
- Head `8f006cf9c3ff92b6ccf550b16f291d4fbcee861f` (a "Merge branch 'main'"
commit), base `39c9e49252f9991dba691b1bd318e8d5362445ad`. Created 2026-08-02,
last updated 2026-08-03 — eleven days stale as of today.
- `mergeable: CONFLICTING`, `reviewDecision: REVIEW_REQUIRED`.
- +287 / −4 across five files: `instrumented_scheduler.py` (+24),
`worker_factory.py` (+26/−4), `tests/test_vllm_instrumented_scheduler.py`
(+169), `tests/test_vllm_worker_factory.py` (+67), `tests/report_pytest_markers.py` (+1).
- Comments and reviews: `copy-pr-bot` gating message, a Datadog comment
reporting **2 failed pipeline jobs including `Pre Merge | pre-commit`**, plus
advisory bot reviews from `coderabbitai` (2 actionable + 1 nitpick) and
`devin-ai-integration` (5 potential issues). **No human review.**
I read the full diff. The approach: a module-level `_install_fpm_worker_id_utility()`
in `instrumented_scheduler.py` attaches a `set_fpm_worker_id` method to
`vllm.v1.engine.core.EngineCore` (guarded by `hasattr`, so it is idempotent and
defers to a native vLLM implementation should one ever appear); the parent then
calls it over vLLM's existing utility RPC via a new
`_sync_fpm_worker_id(engine_client, fpm_worker_id)` helper in `worker_factory.py`,
invoked at all three restore sites. It updates **both** identity fields, which —
per the discovery above — is the non-obvious thing it had to get right. The
monkeypatch reaches the child because the child imports the module while
resolving Dynamo's `--scheduler-cls` dotted path.
**Verdict on the approach: sound.** It uses vLLM's own sanctioned channel for
reaching a live EngineCore rather than the shared memory or process restart the
old TODO assumed were necessary, and it correctly retires that TODO. Being stale
is not grounds to start over; the brief is explicit that only a *wrong* approach
justifies competing, and this one is not wrong.
### One substantive gap I found in #12545
Probing the installed vLLM in this sandbox
(`/opt/dynamo/venv/bin/python`, vLLM 0.22.0):
```
AsyncMPClient.call_utility_async(self, method, *args)
-> self._call_utility_async(method, *args, engine=self.core_engine)
DPAsyncMPClient: 'call_utility_async' in vars() -> False
```
`call_utility_async` targets a **single** `EngineIdentity` (`self.core_engine`),
and `DPAsyncMPClient` / `DPLBAsyncMPClient` inherit it unmodified. In a data-parallel
deployment there is one EngineCore child — and therefore one
`InstrumentedScheduler` and one `_FpmPublisherThread` — per DP rank, each
resolving its own rank via `InstrumentedScheduler._resolve_dp_rank`
(`instrumented_scheduler.py:1504-1517`, which reads `data_parallel_index`
precisely because each child is distinct). So #12545 as written appears to sync
**only one DP rank**; ranks 1..n−1 would keep publishing `""`. That still trips
`_reconcile_fpm_worker_count`, this time at its `dp_sizes` consistency check
(line 478-481) or its coverage check (line 484-488).
Honesty about the limits of this observation: it was made against **vLLM 0.22.0**,
the version installed in this sandbox, whereas `pyproject.toml:63` pins
`vllm[flashinfer,runai,otel]==0.26.0`. The client-class structure is very likely
unchanged, but the claim must be re-confirmed against 0.26.0 before it is
asserted to the PR author. It is offered as a review finding to verify, not as a
settled fact.
### Existing test coverage — RUN / IGNORE verdicts
- `components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py`
(3649 lines) — `grep -c worker_id` returns **0**. No existing coverage of FPM
worker identity whatsoever. **RUN** (regression home; new tests here are not
duplicative).
- `components/src/dynamo/vllm/tests/test_vllm_worker_factory.py` (1092 lines) —
FPM references exist (lines 59, 257, 278, 286, 449-547) but all concern
benchmark payloads and DP-rank result merging, not worker identity. Snapshot
coverage exists at lines 745-804 (`test_passes_snapshot_engine`) and asserts
only that the engine tuple is threaded through. **RUN** (must not regress);
**IGNORE** as duplicate coverage — it does not touch `worker_id`.
- `components/src/dynamo/planner/tests/monitoring/test_decision_state_enums.py`
is the only test file mentioning `worker_count_mismatch`, and it covers the
diagnostic enum, not the reconciliation logic. **IGNORE** for duplication
purposes.
- No test anywhere asserts the `worker_id` field on a published
`ForwardPassMetrics`. This is the coverage hole the fix should close.
Collection check, run read-only from `/tmp` with `-p no:cacheprovider` so the
checkout stayed untouched: both files fail to collect with
`ModuleNotFoundError: No module named 'dynamo'`, and
`/opt/dynamo/venv/bin/python -c "import dynamo"` fails the same way
(`pip show ai-dynamo` → not found). Recipe 00 is therefore genuinely required
before any test run, notwithstanding the note in `compute-env.md`. Per Recipe
00's own text this is a build step, not an infrastructure block.
### Git history over the touched files
`git log --oneline -15` on `instrumented_scheduler.py`, `worker_factory.py`, and
`snapshot.py` shows active, ongoing work and **no prior attempt at this fix and
no revert of one**: `df08bcd56` (#12142 classify/pooling workers), `86408fe23`
(#12829 shut down EngineCore on decode worker exit), `1d0faad16` (#12836),
`e07fa63e8` (#12021 DP self-benchmark capacity), `705444746` (#12358 steady-state
decode), `be65a2136` (#12184 exit vLLM source without engine teardown),
`26e88e91f` (#12202 vLLM bump to 0.26.0). Nine commits touched these files since
2026-08-01 — which is simply why #12545 conflicts: the three restore sites moved
from lines 636/997/1281 at its base to 686/1138/1431 on today's `main`. The
conflict is positional drift, not a design disagreement.
Related open work in the same FPM area, none of it overlapping this fix: #13054,
#12176, #13047, #12627, #12177. Closed-unmerged neighbours #12020 and #12982 are
both about self-benchmark measurement noise, not worker identity.
### Explicitly checked negatives
- `gh pr list --state all --head "fix/vllm-snapshot-worker-id-fpm--5a9b49e37643"`
→ `[]`. `git ls-remote --heads origin <that branch>` → empty.
**No factory-created request exists for this work item**, which rules out
`Route: babysitter`. The prior run recorded in the Slack thread died during
planning and produced nothing.
- `gh issue view 12591` → `state: OPEN`, `closedAt: null`. Nothing merged.
## Chosen approach
Help #12545 land instead of competing with it.
The planner brief is unambiguous about this case: when somebody else already has
an open request doing what you were about to do, the better move is to help it
land, and only a *wrong* approach justifies starting over. #12545's approach is
not wrong — it is the mechanism I would have chosen, and it already handles the
two-identity-fields subtlety that a from-scratch fix would most likely have
missed. What is blocking it is entirely mechanical: an eleven-day-old branch
that now conflicts because nine commits moved the code underneath it, a red
`Pre Merge | pre-commit` job, no human review, and one design question about DP
fan-out that no human has yet asked.
Concretely, under the `existing-request` overlay, working the request **in
place** (no `Review action: reproduce`, so no competing request is opened):
1. **Materialize** #12545's head `8f006cf9c3ff92b6ccf550b16f291d4fbcee861f` onto
the working branch, starting from that head rather than from base — that is
what preserves the option of pushing back into the author's branch as a
fast-forward. Reconcile with today's `main`
(`c7aa0bf99ce91363700fc85edab418d3166a5943`) by **merging**, never rebasing.
The conflicts are positional: re-land the three `_sync_fpm_worker_id` calls
immediately after the `os.environ[ENV_FPM_WORKER_ID] = fpm_worker_id`
assignments now at lines 686, 1138, and 1431, and drop the stale TODO above
1431. Preserve Nianheng Wu's attribution in the commits and in `change.md`.
2. **Green the lint**, since a red `pre-commit` is a hard blocker to merge and
is cheap to fix here.
3. **Add the behavioural regression test** described below. #12545's own tests
are structurally reasonable but lean toward asserting the installation
(`FPM_UTILITY_NAME in vars(EngineCore)`, "install is idempotent", "annotation
is not a msgspec Struct"), which is closer to mirroring the implementation
than to pinning behavior. The gap worth closing is the one no test in the
repository covers: the `worker_id` actually carried on a published
`ForwardPassMetrics`.
4. **Raise the DP fan-out question** as a review finding, clearly labelled as
observed against vLLM 0.22.0 and requiring confirmation against the pinned
0.26.0. Whether it is fixed here or deferred is the author's and maintainer's
call; surfacing it is ours.
The workflow decides push-versus-comment from whether #12545's branch is one we
may push to. I do not decide that, and the publisher — not I — performs it.
### Reproducible symptom
The full symptom (deploy, checkpoint, restore, load, watch replicas fail to
retire) needs Kubernetes and eight GPUs. But the entire causal chain reduces to
two pure-Python observations that run in this sandbox with no GPU, and together
they reproduce the mechanism faithfully:
1. An `InstrumentedScheduler` constructed with `DYN_FPM_WORKER_ID` unset — the
exact state of a snapshot-restored child — yields
`_extract_metrics(...).worker_id == ""`.
2. Feeding the resulting `{("", 0): ...}` into
`_reconcile_fpm_worker_count(fpm_stats, dgd_count=2, label="decode")` returns
`False`, which is precisely what makes the caller set
`_diag_load_reason = "worker_count_mismatch"` and return no decision.
That is the reported failure, start to finish, at unit cost.
### Ranked hypotheses and their discriminating observations
| # | Hypothesis | Smallest distinguishing observation |
|---|---|---|
| H1 (primary) | The child's scheduler holds `_fpm_worker_id = ""` because `__init__` read the env before the parent knew the id, and nothing updates it after restore. | Construct the scheduler with the env unset; assert `_extract_metrics(...).worker_id == ""`. Confirmed structurally by `snapshot.py:41` preceding `:49`. |
| H2 | The id is present in the child but dropped downstream by the `known_workers` filter in `fpm.rs:900` (e.g. an id-format mismatch against MDC). | Compare the `worker_id` on the wire against `generate_endpoint.connection_id()`. If they are equal, H1 is false and the defect is a format mismatch instead. |
| H3 | `_reconcile_fpm_worker_count` is itself over-strict and would freeze even with correct ids. | Pure-function probe: N distinct ids with `dgd_count=N` must return `True`; all-`""` must return `False`. If the first also returns `False`, the planner is at fault, not the worker. |
| H4 (post-fix residual) | Only one DP rank is synced, so DP>1 deployments still freeze. | `call_utility_async` targets `self.core_engine` only, and `DPAsyncMPClient` does not override it (verified on vLLM 0.22.0; re-confirm on 0.26.0). |
H1 is primary and is what the fix addresses. H3 is the cheapest of the four and
should be run first, because a `True` result on distinct ids is what licenses
leaving the planner alone — it converts "we chose not to touch the planner" from
an assumption into a recorded fact.
### Regression test
Behavioural, not tautological — each assertion is on an emitted payload field,
and each fails if the production change is reverted (see
`learnings/no-tautological-tests.md`). Home:
`components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py`, which has
zero `worker_id` coverage today, using the `object.__new__` stub convention the
module already uses.
- **Active-sample identity.** A scheduler in the restored state (`_fpm_worker_id
= ""`) that is then synced with a realistic connection id must have
`_extract_metrics(...).worker_id` equal to that id. Revert the fix → the
payload still carries `""` → fails. Asserts the published field, not the
attribute assignment.
- **Idle-heartbeat identity.** The heartbeat that `_FpmPublisherThread._run`
constructs at line 1423 must carry the same id. This is the field the planner
sees for an *idle* replica, which is exactly the scale-down case, and it is
what distinguishes "updated one field" from "updated both".
- **Planner acceptance (the chain closed).** `_reconcile_fpm_worker_count` must
return `False` for `{("",0)}` with `dgd_count=2` and `True` for two distinct
ids with `dgd_count=2`. This ties the worker-side fix to the user-visible
outcome and simultaneously discharges H3.
### Negative control (must remain unchanged)
- **Cold start.** `setup_vllm_engine` sets `ENV_FPM_WORKER_ID` before engine
construction (`main.py:605-606`), so a cold-started scheduler must still read
the real id in `__init__` and the restore-only RPC must **not** be invoked on
the non-snapshot branch. Assert both.
- **Engine still boots.** `_install_fpm_worker_id_utility()` runs at module
import in *every* vLLM worker, cold start included — not only in snapshot
mode. A vLLM worker must still start and serve a request with it installed.
This is the one part of the negative control that is worth spending GPU on.
## Rejected alternatives
- **Write our own fix (`Route: implementation` + `debug-investigation`).**
Rejected. It would produce a second PR over the same four files for the same
issue, leaving a maintainer to arbitrate — the exact anti-pattern the planner
brief names. #12545's approach is sound, so staleness alone does not license
competing.
- **`Route: babysitter`.** Rejected on fact, not preference: `gh pr list --head
fix/vllm-snapshot-worker-id-fpm--5a9b49e37643` returns `[]` and the branch does
not exist on the remote. #12545 belongs to a human contributor and was not
created by this factory, so it is not a maintenance target.
- **`Disposition: already-resolved`.** Rejected. The brief forbids it for an
open PR, and `main` demonstrably still contains the defect and the TODO.
- **Re-read `os.environ` on every FPM publish** instead of pushing the id over
RPC. Rejected: the env var is set in the *parent*; the forked child's
`os.environ` is a private copy that the parent's assignment never reaches. This
is the same misconception the original TODO encoded. It would also add an env
lookup to a hot path.
- **Restart the EngineCore process after restore** to re-read the environment.
Rejected: it discards the entire point of snapshot restore (fast warm start)
and is what the old TODO wrongly assumed was necessary.
- **Relax `_reconcile_fpm_worker_count`** to tolerate blank ids. Rejected: it
suppresses the symptom while leaving the planner unable to attribute metrics
to replicas, and the debug-investigation overlay directs the reviewer to reject
exactly this shape of fix.
- **Broadcast the sync to every DP rank as our own new design.** Not rejected on
merit — but it is the author's call to make on their own PR, so it is raised
as a review finding rather than unilaterally rewritten.
- **On the parameter-shape question the brief raises:** `_sync_fpm_worker_id` is
a new module-level helper with no existing callers, and #12545 adds no
parameter to any widely-called function. No caller migration is implied, so no
optional-versus-required trade-off arises and `cargo check --workspace` is not
needed as a caller sweep. There are no Rust changes.
## Validation strategy
Python-only change under `components/src/dynamo/vllm/` plus one line in
`tests/report_pytest_markers.py`. Sized accordingly, and weighted toward recipes
that run something rather than read something.
- **`00-dynamo-editable-install`** — mandatory, and here not merely ceremonial: I
verified empirically that `import dynamo` fails in `/opt/dynamo/venv` and both
touched test files fail collection with `ModuleNotFoundError: No module named
'dynamo'`. Every test below imports `dynamo.*`. No GPU.
- **`01-python-lint`** — directly load-bearing, not routine: #12545's
`Pre Merge | pre-commit` job is **red**, and a red pre-commit is a hard merge
blocker. Greening it is a concrete part of helping the request land. No GPU.
- **`03-python-unit-tests-mocker`** — the primary evidence. Runs
`components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py` and
`test_vllm_worker_factory.py` (both marked RUN above), including #12545's
existing tests and the new behavioural regression tests, plus the
`_reconcile_fpm_worker_count` probe from
`components/src/dynamo/planner/core/state_machine.py`. Proves the claim the
change actually makes — that a restored worker's *published* FPM payload
carries the real `worker_id`, on both the active and idle paths — and, via the
planner probe, that this is sufficient for reconciliation to succeed. It also
carries the unit half of the negative control (cold start unchanged, no RPC on
the cold branch). No GPU.
- **`07-agg-smoke`** — the integration-level negative control. The fix installs
an unconditional module-level monkeypatch on `vllm.v1.engine.core.EngineCore`
that executes in *every* vLLM worker, cold start included; a regression there
would break all vLLM serving, and no unit test can catch that. One
A100-80GB is sufficient for a small model. **Hardware:** 1 GPU, available.
**Version caveat:** the sandbox carries vLLM 0.22.0 while `pyproject.toml:63`
pins 0.26.0. If the smoke fails for that skew rather than for the change, the
validator should record `blocked / vllm-version-skew` with the exact evidence
rather than reporting a false `fail`; the unit-level cold-start control in
recipe 03 stands on its own in that event.
- **`05-code-inspection`** — nominated for the one claim that genuinely cannot be
executed here, not as a substitute for running things. Two items: (a) the
end-to-end capture → restore → load-drop → scale-down loop, which needs the
snapshot-agent, the planner, and multiple DGD replicas on Kubernetes and gets a
recorded disposition; (b) the DP fan-out finding, re-checked against the
pinned vLLM 0.26.0 client sources before it is asserted to the author.
What this sandbox can and cannot prove, stated plainly: it **can** prove that
the published FPM payload carries the correct `worker_id` on both the active and
the idle path, that a correctly-populated `fpm_stats` makes
`_reconcile_fpm_worker_count` return `True` where the blank-id form returns
`False`, that cold start is untouched, that a vLLM worker still boots with the
patch installed, and that pre-commit is green. It **cannot** prove that a real
snapshot-restored pod on a real cluster retires a replica under falling load;
that requires the multi-component Kubernetes loop and is recorded as a
disposition, never claimed as evidence.
```validation-recipes
00-dynamo-editable-install
01-python-lint
03-python-unit-tests-mocker
05-code-inspection
07-agg-smoke
```
## Required deliverables
Per the `existing-request` overlay, five packets:
- `plan.md` — this file.
- `change.md` — narrative of the reconciliation with today's `main`, the lint
fix, and the added regression tests, preserving Nianheng Wu's attribution and
naming #12545 as the origin of the fix.
- `change.diff` — nonempty, materialized from #12545's head
`8f006cf9c3ff92b6ccf550b16f291d4fbcee861f` merged with `main`
`c7aa0bf99ce91363700fc85edab418d3166a5943`, never rebased.
- `change-validation.md` — recorded runs of the five recipes above, with
`## Investigation outcome:` kept separate from the verdict, ending in exactly
`## Verdict: pass`, `fail`, or `blocked`.
- `review.md` — advisory assessment ending in exactly `## Assessment: sound` or
`needs_changes`, including the DP fan-out finding as feedback for the PR
author.
Publication is the publisher's alone: it works #12545 in place, pushing into its
branch if permitted and otherwise commenting. No agency role — this planner
included — pushes, opens, or comments on any review request. |
change.md# change.md — wi-20260814T172511Z-12591
Work on existing PR **[#12545 — `fix(vllm): propagate FPM worker_id into
snapshot-restored EngineCore`](https://github.com/ai-dynamo/dynamo/pull/12545)**,
worked **in place** (no reproduction, no second PR).
## Attribution
The fix is **Nianheng Wu's** (`RealNicolasBourbaki` / `NianhengWu
<nianheng.wu@t-systems.com>`), authored in PR #12545. PR #12545 is the origin of
this change. Their three commits and their merge commit are preserved byte for
byte on the working branch — nothing was amended, squashed, reordered, or
rebased. My contribution sits strictly on top as two additional commits.
## Branch and provenance
| | |
|---|---|
| Working branch | `fix/vllm-snapshot-worker-id-fpm--5a9b49e37643` |
| Branched at | `8f006cf9c3ff92b6ccf550b16f291d4fbcee861f` (the PR head) |
| PR head SHA confirmed | **Yes — exact match.** `gh pr view 12545 --json headRefOid` returns `8f006cf9c3ff92b6ccf550b16f291d4fbcee861f` |
| Local base `main` | `c7aa0bf99ce91363700fc85edab418d3166a5943` |
| PR's original base | `39c9e49252f9991dba691b1bd318e8d5362445ad` (merge-base) |
| Reconciled with newer base by | **merge, not rebase** |
Because the branch starts at the PR head and history was only ever added to, a
push into the author's `fix-fpm-worker-id` branch would be a **fast-forward**.
(Nothing was pushed — see *Scope*.)
### Commits on the branch
```
485b1e25f svc-glamr test(vllm): assert restored worker id on the emitted FPM payload <- authored here
feb6761f7 svc-glamr Merge branch 'main' into fix/vllm-snapshot-worker-id-fpm--... <- authored here
8f006cf9c Nianheng Nicole Wu Merge branch 'main' into fix-fpm-worker-id <- inherited
75a3499d9 NianhengWu fix(vllm): appease black and stub vllm.v1.engine.core... <- inherited
76dfb8ccf NianhengWu fix: add test units <- inherited
64f2d020d NianhengWu fix: add syncs between main and child proc for fpm... <- inherited
```
Both commits I authored were made with `git commit --signoff` and carry a DCO
`Signed-off-by` trailer.
## change.diff — what range it covers
`change.diff` is exactly:
```
git diff main..HEAD # main = c7aa0bf99, HEAD = 485b1e25f
```
`main` is an ancestor of `HEAD` (the merge, not a rebase), so the two-dot and
three-dot forms are **byte-identical**; I verified this rather than assuming it.
The diff is the complete net change against the current `main` — the author's
fix *and* my additions together, which is what a reviewer wants to read.
- 706 lines, 5 files, +607 / −4.
- `git apply --check change.diff` **passes against `main`** (verified by
checking out `main`, running the check, and returning to the branch).
## The four requested corrections
### 1. Resolve merge conflicts — done, but smaller than predicted
The plan predicted conflicts at the three restore sites in `worker_factory.py`
and gave line numbers 686 / 1138 / 1431. **I verified the line numbers against
the checkout myself rather than trusting them, and they had moved.** Git
auto-merged all three production hunks; the *only* real conflict was the import
list in `test_vllm_worker_factory.py` (HEAD added `_sync_fpm_worker_id`, `main`
added `_DecodeWorkerLifecycle`). Resolved by keeping both in isort order.
I did not trust the auto-merge. I re-read all three sites and confirmed the
post-merge state in `components/src/dynamo/vllm/worker_factory.py`:
| Site | `if snapshot_engine is not None:` | `os.environ[ENV_FPM_WORKER_ID] = ...` | `await _sync_fpm_worker_id(...)` |
|---|---|---|---|
| `_create_realtime_worker` | 692 | 700 | 703 |
| `_create_decode_worker` | 1148 | 1156 | 1159 |
| `_create_prefill_worker` | 1441 | 1449 | 1452 |
Each `_sync_fpm_worker_id` call sits immediately after its `os.environ`
assignment, inside the restore branch, as required. The shift from the plan's
numbers (+14 / +18 / +21) is explained by the module-level `_sync_fpm_worker_id`
helper the PR adds at line 75. The stale TODO above the third site is gone
(`grep "TODO: The scheduler"` returns nothing).
### 2. Green the lint — **the premise was stale; this was a no-op**
The plan states PR #12545's `Pre Merge | pre-commit` job is red and is a hard
merge blocker. **That is no longer true, and I am reporting it rather than
inventing work.** At head `8f006cf9c`:
```
gh pr view 12545 --json statusCheckRollup -> pre-commit: SUCCESS
```
Every check at the PR head is SUCCESS or SKIPPED. The author's own commit
`75a3499d9 fix(vllm): appease black and stub vllm.v1.engine.core for marker
collection` fixed it. The planner appears to have read an older CI comment.
I confirmed independently that lint is green **after** my merge and **after** my
test additions:
```
pre-commit run --files components/src/dynamo/vllm/worker_factory.py \
components/src/dynamo/vllm/instrumented_scheduler.py \
components/src/dynamo/vllm/tests/test_vllm_worker_factory.py \
components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py \
components/src/dynamo/vllm/main.py
```
All hooks pass, including `Report pytest markers (static + inherited)` — the
hook most likely to break on new tests. (`isort`, `black`, `flake8`, `ruff`,
`codespell`, whitespace/EOF hooks: all Passed. No `mypy` hook is configured in
`.pre-commit-config.yaml`.) I formatted only the files I changed.
### 3. Regression tests — added
The PR's existing FPM tests assert that `set_fpm_worker_id` **writes two
attributes**. That is structural, not behavioural: what the planner actually
consumes is the `worker_id` field of a `ForwardPassMetrics` message on the wire.
The new tests close that gap. Every one starts from the state a
snapshot-restored child really boots into (`""` on both identity fields) and
asserts on an **emitted payload field**.
**`components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py`**
(+364)
- `test_restored_worker_stamps_synced_id_on_active_samples` — plan's
*active-sample identity*. Asserts on `_extract_metrics(...)`'s returned
`ForwardPassMetrics` (the exact struct `_publish_or_record_metrics` publishes),
including that the id rides on a genuine sample (`num_decode_requests == 4`,
`wall_time == 0.0123`) rather than a blank struct.
- `test_restored_worker_stamps_synced_id_on_idle_heartbeats` — plan's
*idle-heartbeat identity*. Drives the **real** `_FpmPublisherThread._run` loop
and decodes the bytes handed to `send_multipart` back through
`dynamo.common.forward_pass_metrics.decode`, so it covers construction *and*
serialization. This is the path a scale-down candidate is on.
- `test_sync_keeps_dp_rank_and_counter_intact` — the sync must change identity
only; `(worker_id, dp_rank)` is the planner's fleet key and `counter_id` is
the per-key sequence.
- `test_planner_accepts_a_restored_fleet_only_after_the_sync` — plan's
*planner acceptance*, closing the chain end to end.
**`components/src/dynamo/vllm/tests/test_vllm_worker_factory.py`** (+192) — the
plan's **negative control**, as a matched pair through one harness
(`_drive_prefill_worker`) so it is a genuine control and not a lone assertion:
- `test_cold_start_does_not_issue_the_restore_rpc` — cold start is unchanged:
no RPC, and the id still routes through `setup_vllm_engine(fpm_worker_id=...)`.
- `test_restored_worker_issues_the_rpc_with_the_new_connection_id` — the restore
branch issues exactly one RPC, against the *restored* client.
#### These are not tautologies — I proved it
Per `no-tautological-tests.md`, I ran **four revert probes** against the
production code, confirmed the expected tests fail, and restored the file each
time (`git diff` clean afterwards, verified):
| Probe (production code temporarily broken) | Result |
|---|---|
| Delete `scheduler._publisher._worker_id = new_worker_id` (the **subtle half-fix** the plan warns about) | **4 failed** — incl. idle-heartbeat + planner tests. Active-sample test still passes, exactly as the two-identity-field analysis predicts. |
| Comment out `_install_fpm_worker_id_utility()` (full revert) | **11 failed** (all FPM tests) |
| Delete the `await _sync_fpm_worker_id(...)` from the prefill restore branch | positive control fails |
| Hoist that call out of the `if` so it fires on cold start too | cold-start control fails |
The half-fix probe is the important one: the planner test failed with the
**literal warning from the original issue** —
`Worker count mismatch: DGD=2, FPM=1 for decode`. The test reproduces the
reported symptom, not a paraphrase of it.
The planner test is a real regression test, not a planner characterization test:
the `worker_id`s it feeds `_reconcile_fpm_worker_count` are **decoded from
payloads emitted by the real publisher loop**, so reverting the vLLM fix
collapses the two fleet keys into one and fails it. It asserts both directions
(unsynced fleet → `False`, synced fleet → `True`).
#### Deviation from the plan: one guarded import
The plan puts all three regression tests in
`test_vllm_instrumented_scheduler.py`. I kept them there, **but guarded the
planner import**:
```python
pytest.importorskip("sklearn", reason="planner extra not installed")
pytest.importorskip("aiconfigurator_core", reason="planner extra not installed")
```
Reason: importing `dynamo.planner.core.state_machine` pulls in `scikit-learn`
(via `perf_model/base.py`) and `aiconfigurator-core` (via
`perf_model/engine_query.py`), neither of which is installed in the
`pytest.mark.vllm` CI job. An unguarded import would turn a currently-green
vLLM unit job **red** — the opposite of helping #12545 land. With the guard the
test runs wherever the planner extra exists and skips cleanly elsewhere. **I ran
it un-skipped locally, with the deps installed, and it passes and is
revert-sensitive** (see the probe table). The validator will see it *skip*; that
is expected, and the skip reason is explicit in `-rs` output.
### 4. Nothing else — respected
- Did **not** touch the planner, `fpm.rs`, SGLang, or TensorRT-LLM. The only
planner contact is a read-only call to a `@staticmethod` from a test.
- Did **not** fix the DP fan-out question. Left as a review finding (below).
- Production files (`worker_factory.py`, `instrumented_scheduler.py`) are
**unmodified by me** — the +24 / +30 in the diff are entirely the author's,
arriving via the branch. My two commits touch only test files.
## Commands run and their results
Everything below was run in `/opt/dynamo/venv`.
| Command | Result |
|---|---|
| `pytest test_vllm_instrumented_scheduler.py test_vllm_worker_factory.py` | **197 passed, 1 skipped** (skip = the guarded planner test, post-venv-restore) |
| same, with planner deps installed | **198 passed** |
| `pytest -k "fpm_utility or restored_worker or sync_keeps or planner_accepts"` | 11 passed |
| 4 × revert probes | failed as designed (table above), production files restored |
| `pre-commit run --files <5 changed files>` | all hooks Passed |
| `gh pr view 12545 --json statusCheckRollup` | all checks SUCCESS/SKIPPED — plan's "lint is red" is stale |
| `git apply --check change.diff` (on `main`) | applies cleanly |
### Things I tried that did not work / needed fixing
1. **`ModuleNotFoundError: No module named 'sklearn'`** importing
`dynamo.planner.core.state_machine`. Fixed with
`uv pip install scikit-learn==1.7.2` (the version pinned in
`container/deps/requirements.planner.txt:23`).
2. **`ModuleNotFoundError: No module named 'aiconfigurator_core'`**, same import
chain. Fixed with `uv pip install aiconfigurator-core==0.11.0.dev20260728`
(pinned at `pyproject.toml:85`) — **but this silently downgraded `numpy`
2.3.5 → 1.26.4 and `scipy` 1.18.0 → 1.17.1 in the shared venv.** See below.
3. **First `env` / `git branch -a` invocation dumped ~138 KB** and was truncated;
re-run as targeted queries.
4. **First draft of the negative control referenced a helper that did not
exist** (`_mock_vllm_config()`); replaced with a local `_fpm_engine_tuple()`.
### Shared venv: perturbed, then restored — please read
`/opt/dynamo/venv` is shared with the validator (recipes 03 / 07). My planner
probe downgraded numpy and scipy in it. **I restored it and verified the
restoration:**
```
uv pip uninstall aiconfigurator-core scikit-learn joblib threadpoolctl
uv pip install numpy==2.3.5 scipy==1.18.0
```
Post-restore verification: `numpy 2.3.5`, `scipy 1.18.0`, `vllm 0.22.0`,
`torch 2.11.0+cu129`, `import dynamo._core` OK, `import dynamo.vllm.main` and
`import dynamo.frontend` OK. The venv is back to its pre-probe state.
**Pre-existing, not caused by me:** `uv pip check` reports `msgspec`,
`prometheus-client`, `pyzmq`, `transformers`, `zstandard` as "not installed".
Four of those five import fine (a uv metadata artifact). `zstandard` is
genuinely absent, but it is used only lazily inside
`components/src/dynamo/common/metadata_upload.py` (behind an explicit
"install the extra" error) and one SGLang test — not on the vLLM serving path.
`uv pip uninstall` does not remove dependencies, so my uninstall could not have
caused it.
## Empirical finding: why the planner is correctly left alone
The plan's non-goal "do not change the planner" was an assumption. I discharged
it as **fact** by calling
`PlannerScalingState._reconcile_fpm_worker_count` directly:
- fleet `{('', 0)}` with `dgd_count=2` → **`False`**, logging
`Worker count mismatch: DGD=2, FPM=1 for decode` — the exact reported symptom.
- fleet with two distinct ids and `dgd_count=2` → **`True`**.
The planner is behaving correctly on the input it is given. The bug is entirely
upstream in what vLLM emits, so the fix belongs where #12545 puts it. This is
now encoded as `test_planner_accepts_a_restored_fleet_only_after_the_sync`.
## For the reviewer — open finding, deliberately not fixed
**DP fan-out (H4).** `_sync_fpm_worker_id` calls
`engine_client.engine_core.call_utility_async(...)`, which targets a **single**
`EngineIdentity` (`self.core_engine`). `DPAsyncMPClient` / `DPLBAsyncMPClient`
do not override `call_utility_async` (observed on the in-sandbox **vLLM
0.22.0**; `pyproject.toml` pins **0.26.0**, so this must be re-confirmed against
the pinned version before acting). If that holds, a snapshot-restored worker
with `dp_size > 1` would sync **rank 0 only**, leaving ranks 1..N−1 publishing
`""` — which `_reconcile_fpm_worker_count` rejects at its
`len(fpm_stats) != dgd_count * dp_size` check, i.e. a partial fix for DP
deployments.
Per instruction this is raised **as a question for the author**, not rewritten
unilaterally. It does not block the non-DP fix, which is complete and now tested.
`test_sync_keeps_dp_rank_and_counter_intact` pins the single-rank contract so a
future fan-out change has a fixed point to build on.
## Scope
Nothing was pushed. No review request was opened, and no comment was posted. No
files were written outside the work-item root and the checkout. No inherited
history was rewritten. The branch name is exactly
`$GLAMR_CODING_BRANCH_NAME`.
MR URL: https://github.com/ai-dynamo/dynamo/pull/12545 |
change-validation.md# Validation — PR #12545 (FPM worker identity on snapshot-restored vLLM workers)
Validated against the printer's working branch (HEAD `485b1e25f`) in the sandbox
described by `compute-env.md`: one A100-SXM4-80GB, CUDA 12.9, Python 3.12.3 at
`/opt/dynamo/venv`, no Docker daemon. Every command in the evidence trail was run
through the customs recorder; nothing was executed outside it. The checkout was
never modified — each recipe asserts `git status --porcelain` empty at its end,
and every run reports `checkout clean: yes`.
## What the change does, and what had to be proven
A vLLM `EngineCore` child restored from a snapshot is forked before the Dynamo
runtime exists, so `InstrumentedScheduler.__init__` reads `DYN_FPM_WORKER_ID`
(`instrumented_scheduler.py:134`) as the empty string at
`instrumented_scheduler.py:1470`. Downstream, the Rust FPM subscriber drops any
sample whose worker is not a known worker (`lib/bindings/python/rust/llm/fpm.rs:908`),
the planner's metrics provider consumes exactly that filtered map
(`planner/environment/metrics_provider/runtime_provider.py:185`),
`_reconcile_fpm_worker_count` (`planner/core/state_machine.py:465`) then sees
disagreeing counts, and `load_scaling.py:74` abandons the scaling decision. The
restored replica is invisible to the planner, and the scale-down never happens.
The fix installs `set_fpm_worker_id` on the `EngineCore` base class at module
import time, unconditionally and idempotently (`instrumented_scheduler.py:4098-4116`),
and the parent invokes it over vLLM's utility RPC from `_sync_fpm_worker_id`
(`worker_factory.py:75`) at all three snapshot-restore sites — each an
`os.environ[ENV_FPM_WORKER_ID] = ...` immediately followed by the sync
(`worker_factory.py:700/703`, `1156/1159`, `1449/1452`).
The load-bearing detail is that there are **two** independent identity fields, not
one: `InstrumentedScheduler._fpm_worker_id`, stamped on active samples
(`:1487`, `:1719`), and `_FpmPublisherThread._worker_id`, stamped on idle
heartbeats (`:1368`, `:1423`). A fix touching only the first leaves an idle
replica publishing an empty id — and idle is precisely the scale-down case the
bug is about. The utility sets both (`:4110-4111`).
## What the recipes established
**Recipe 00** confirmed the editable install: `dynamo` and `dynamo._core` both
import from `/opt/dynamo/venv`. Idempotent here because `cli.ts` pre-built the
bindings, but recorded rather than assumed.
**Recipe 01** ran `pre-commit --hook-stage manual` over the five changed files.
Green, with no tracked-file modification afterwards. Worth stating plainly: the
plan asserted PR #12545's `Pre Merge | pre-commit` was red. That premise is stale
— the author's own commit `75a3499d9` fixed it, upstream now reports the check
green (see Recipe 05 below), and the local run agrees. Nothing to fix here.
**Recipe 03** is the primary evidence: `197 passed, 1 skipped` across
`test_vllm_instrumented_scheduler.py` and `test_vllm_worker_factory.py`, and
`25 passed, 1 skipped` when narrowed to the FPM worker-id tests by name.
The tests are not vacuous, and I established that empirically rather than by
reading them. Both reverts were simulated **at runtime** through `pytest -p`
plugins loaded from `/tmp`, so the checkout stayed untouched:
- *Full revert* (`set_fpm_worker_id` removed from `EngineCore`, installer
neutered — the pre-#12545 state): **11 failed**, 186 passed.
- *Half-fix revert* (utility present but the publisher's identity field no longer
updated): **3 failed**, 194 passed — and exactly the right three.
`test_fpm_utility_updates_scheduler_and_publisher` fails
`assert '' == '8465209922961459'` (`:3729`),
`test_fpm_utility_overwrites_a_previously_set_id` fails
`assert '1111111111111111' == '2222222222222222'` (`:3740`), and
`test_restored_worker_stamps_synced_id_on_idle_heartbeats` fails
`assert '' == '8465209922961459'` (`:3899`) — while every active-sample test
still passes.
That second probe is the one that matters. It reproduces the plausible incomplete
fix and shows the suite catches it on the idle-heartbeat path specifically. The
two-fields claim is proven, not asserted.
**Recipe 05** traced every behavioural claim to a `file:line` citation (the five
claims above), checked for parallel patterns, and read CI. Findings: the three
restore sites each pair the env assignment with the sync; `main.py:606` is a
fourth writer but a **cold-start** site — the assignment precedes
`AsyncLLM.from_vllm_config(...)` at ~`:624`, so the forked child inherits the
correct value and needs no sync, i.e. not a missed restore path;
`EngineCore.set_fpm_worker_id = set_fpm_worker_id` (`:4113`) is the only
monkeypatch onto vLLM internals in this component; and the SGLang and TensorRT-LLM
backends carry no parallel identity-restore pattern, so there is nothing analogous
left unfixed. Upstream CI shows every Required check passing (`pre-commit pass
1m18s`), with `rust-tests`/`operator`/`rust-clippy`/`snapshot`/Fern checks
skipping as expected for a Python-only change. `CHANGES_REQUESTED` count is `0`;
review states are `["COMMENTED"]`.
**Recipe 07** is the integration-level negative control, and it is the reason this
recipe was worth the GPU time. The fix installs an *unconditional* module-level
monkeypatch that executes in every vLLM worker, cold start included — a
regression there would break all vLLM serving and no unit test would catch it. A
worker booted with the patch installed, registered at 61 s, served at 66 s, and
returned HTTP 200 with `choices[0].message.content` (gibberish is expected under
`--load-format dummy`). Only two WARNING lines, both benign and unrelated.
Teardown was clean: no ERROR/Exception/Traceback lines, and the card returned to
`baseline_used_mib=1 after_used_mib=1`.
Two documented deviations from the Recipe 07 text, both recipe-vs-tree drift on
`main` and neither related to the change under validation. Each is proved from
source inside the run before it is applied: `--store-kv` no longer exists and is
now `--discovery-backend` (`common/configuration/groups/runtime_args.py:171`,
`frontend/frontend_args.py:395`); and the `/health` gate must match
`"instance_id"`, not `"endpoint_id"`, because `endpoint_id()` is a *method*
feeding the `endpoints` array and not a serialized field of `Instance`
(`lib/llm/src/http/service/health.rs:77-97`, `lib/runtime/src/component.rs:109-134`).
## What this environment could not prove
Three gaps, stated plainly rather than papered over. None of them is a failure;
all three are limits of the sandbox.
1. **The end-to-end loop is untested here.** Capture → restore → load drop →
planner scale-down needs a multi-replica Kubernetes fleet with the snapshot
agent and planner running. Recipe 05 traced that path statically with
citations, but no execution in this sandbox touches it. The headline claim —
"the planner now scales down a restored fleet" — is argued, not demonstrated.
2. **Everything was proven against vLLM 0.22.0; `pyproject.toml:63` pins
`vllm[flashinfer,runai,otel]==0.26.0`.** This matters more than usual: the fix
monkeypatches a vLLM internal (`EngineCore`) and calls a vLLM internal RPC
(`call_utility_async`). Both are version-sensitive private surfaces. The smoke
passing on 0.22.0 is real evidence that the patch is well-formed, but it is
not evidence about the pinned version. Note that the smoke did *not* fail on
the version skew, so the contingency to record `blocked` for 07 never
triggered.
3. **`test_planner_accepts_a_restored_fleet_only_after_the_sync`
(`test_vllm_instrumented_scheduler.py:3924`) did not run** — skipped with
"planner extra not installed". I confirmed the skip is genuine and not
self-inflicted: `sklearn: ABSENT`, `aiconfigurator_core: ABSENT`, with
`numpy 2.3.5 | scipy 1.18.0 | vllm 0.22.0`. I deliberately did **not** install
those deps, because doing so previously downgraded numpy/scipy in the shared
venv. This is the single test that binds the fix to the planner-visible
outcome, and it is the one test that did not execute. Combined with the PR's
CI showing the test jobs skipping, I have no evidence that this test has
executed anywhere — which the reviewer should weigh, since it is exactly the
assertion the change exists to make true.
## Finding for the reviewer: DP fan-out is likely incomplete
This is a code-level finding, raised because it is actionable, and it also
corrects a claim in `change.md`.
`change.md` states that neither DP client overrides `call_utility_async`. Read
against the sandbox's vLLM 0.22.0
(`vllm/v1/engine/core_client.py`), that is half right:
- `AsyncMPClient.call_utility_async` (`:1070`) targets a **single** engine
(`engine=self.core_engine`).
- `DPLBAsyncMPClient` (`:1349`) **does** override at `:1412` and fans out to
every engine via `asyncio.gather` over `self.core_engines`.
- `DPAsyncMPClient` (`:1169`) does **not** override, so it inherits the
single-engine version.
`DPAsyncMPClient` is the internal-DP-load-balancing case — the mode Dynamo's own
worker announces, warning "vLLM selects internal DP load balancing" at
`components/src/dynamo/vllm/handlers.py:1071`, whose branch comment reads
"internal load balancing, the worker is responsible for all DP ranks"
(`:1069`). If that reading holds, then for a restored worker with
`dp_size > 1` under internal LB only one engine receives the corrected id and the
remaining DP ranks keep publishing the empty string. Because
`_reconcile_fpm_worker_count` keys on `(worker_id, dp_rank)`, a partial identity
restore could still starve the planner — the original bug, narrowed rather than
removed.
Two caveats on this finding, both mine to declare: it was read against 0.22.0 and
not the pinned 0.26.0, and I did not execute a `dp_size > 1` restore to confirm
it. It is a strong code-level signal, not a reproduced defect.
## Honesty notes on the evidence trail
The registry contains failed runs that precede the green ones; customs takes the
most recent run per recipe, but the earlier entries are visible and I would
rather name them than let a reader guess. Three failed Recipe 07 runs: the first
used the recipe's stale `--store-kv`; the second used the recipe's stale
`"endpoint_id"` health gate and produced a false negative on a worker that had in
fact booted and served; the third raced the frontend's model-group commit and got
a legitimate 404, fixed by gating on `/v1/models`. The second run also orphaned
an `EngineCore` child holding 75,880 MiB, which I reaped and then prevented with
a teardown that waits for the child before escalating. Two failed Recipe 03
entries are pure shell-quoting mistakes of my own (a heredoc and a nested-quote
error, exit 127), fixed by moving the logic into a script. All five of these are
my tooling errors or recipe drift, not product defects.
## Interpretation
Every planned recipe executed and cleared, the new tests were run in their real
binary against the changed source, and the test infrastructure was exercised under
both polarities — a positive run with the fix in place and two negative runs that
confirm the suite fails when the fix is reverted, including the sharp half-fix
case that isolates the idle-heartbeat path. Lint is green, upstream CI is green,
and no reviewer has requested changes. The integration smoke confirms the
unconditional monkeypatch does not regress ordinary vLLM serving, which was the
main regression risk the change carried.
What remains open is coverage, not correctness: the Kubernetes end-to-end loop and
the pinned engine version are outside this sandbox, the one planner-binding test
did not run, and the DP fan-out gap deserves an answer before merge. None of
those is something the implementer can resolve by re-running here, and none is a
defect this validation observed. They belong to the reviewer's judgement, and the
DP finding in particular is worth raising on the PR.
## Verdict: pass |
glamr-agent
left a comment
There was a problem hiding this comment.
Automated AI review — advisory. The inline comments below are the findings. The whole review is in the review.md comment on this pull request.
Assessment: soundThe production change is small, correctly targeted, and does the whole job. A The detail that makes it correct rather than nearly correct is that the worker identity I have three findings, all Findings
No finding is What I verified independentlyThe printer's scope claim holds. The base is current and the merge is inert. The branch merges The new tests are behavioural, not tautological. The added helpers drive the real The evidence trail shows real execution. Customs reports Recipe 07 was worth its GPU time. The fix installs an unconditional module-level The The stale One accuracy nit on the evidence trail itself, raised for completeness and not as a What remains open — for the human reviewer, not the authorThese are limits of the validation environment, already declared in
Findings not posted inline
|
tedzhouhk
left a comment
There was a problem hiding this comment.
Agentic review: four inline findings after cross-checking the PR against Dynamo's snapshot lifecycle and the pinned vLLM 0.26.0 implementation.
|
@tedzhouhk Hi I made changes according to your comment, could you please take a look again? Thanks :) |
|
@RealNicolasBourbaki, the current head
@tedzhouhk, once the branch and validation are refreshed, could you take another look? We should advance the existing fix rather than reopen the already-addressed timeout/gating discussion. |
|
@RealNicolasBourbaki if you're on dynamo slack, please ping me Hongkuan Zhou and I can give it another round of review. |
|
@RealNicolasBourbaki if can you can contact with @tedzhouhk on slack, maybe this pr can can forwarded. happy to hear what you like :) |
Thanks! I've got some capacity and will work on this PR further today :) Will ping on slack |
0333b34 to
e7156ea
Compare
e7156ea to
0119d9c
Compare
0119d9c to
e8bf5d1
Compare
e8bf5d1 to
9e62423
Compare
A vLLM EngineCore child restored from a snapshot is forked before the Dynamo runtime exists, so InstrumentedScheduler reads DYN_FPM_WORKER_ID as "" and publishes forward-pass metrics under an empty worker id. The Rust FPM subscriber drops samples that belong to no known worker, the planner then sees fewer FPM workers than registered workers, and _reconcile_fpm_worker_count abandons every scaling decision: the restored replica is invisible to the planner and never scales down. Install a set_fpm_worker_id utility on the EngineCore base class at import time and invoke it over vLLM's utility RPC from the parent at each of the three snapshot-restore sites (realtime, decode, prefill). The utility sets both identity fields -- InstrumentedScheduler._fpm_worker_id, stamped on active samples, and _FpmPublisherThread._worker_id, stamped on idle heartbeats -- since an idle replica is exactly the scale-down case. The sync is gated on the snapshot having captured an InstrumentedScheduler: snapshot mode is independent of FPM and benchmark mode, and a foreign scheduler has no publisher to retarget. Failure is fatal rather than logged, because a worker that registers with an empty id reproduces the same worker_count_mismatch bug silently. A failed sync shuts the restored engine down on every path: decode through its lifecycle, realtime and prefill through _sync_fpm_worker_id_or_shutdown. The RPC is bounded by asyncio.wait_for, as vLLM's call_utility_async awaits its future with no deadline of its own and a wedged rank would otherwise hang the restore before endpoint registration. Signed-off-by: NianhengWu <nianheng.wu@t-systems.com>
9e62423 to
59d2bc0
Compare
tedzhouhk
left a comment
There was a problem hiding this comment.
Reviewed 59d2bc0. The previous restore cleanup and test-harness concerns are addressed across realtime, decode, and prefill. The updated description distinguishes pinned-runtime unit validation and the reported patched/control snapshot acceptance results. Local scheduler-focused tests passed (9); full local collection was blocked by the older environment and missing redis dependency. Approval is based on the code review and documented validation, with full CI still required.
|
/ok to test 59d2bc0 |
@tedzhouhk, there was an error processing your request: See the following link for more information: https://docs.gha-runners.nvidia.com/cpr/e/2/ |
|
/ok to test 59d2bc0 |
✅ Dynamo PR CI passed — run 36918728538 (attempt 1) on
|
Head branch was pushed to by a user without write access
61bf05a to
59d2bc0
Compare
Merges origin/main at 2fde30b. Only #12545 (propagate FPM worker_id into snapshot-restored EngineCore) touched this branch's files, and it conflicted on imports only: - instrumented_scheduler.py: keep main's EngineCore import and this branch's FullAttentionSpec and SlidingWindowSpec imports. - tests/test_vllm_instrumented_scheduler.py: keep this branch's logging import and main's subprocess, sys and textwrap imports. Signed-off-by: Yiming Liu <yimingl@nvidia.com> Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Merges origin/main at ce45db9. Four files conflicted, all resolved by keeping both sides: - instrumented_scheduler.py imports: this branch's importlib.metadata version import and main's itertools chain (vLLM 0.31.0 bump, #15643). - worker_factory.py imports from instrumented_scheduler: main's InstrumentedScheduler (#12545) and this branch's benchmark_content_point_key. - test_vllm_worker_factory.py imports: main's Config (#15441), ENV_FPM_WORKER_ID, InstrumentedScheduler and FPM_SET_WORKER_ID_METHOD_NAME (#12545), next to this branch's FpmBenchmarkWorkerExtension, ENV_FPM_BENCHMARK_OUTPUT_PATH, benchmark_content_point_key and the probe and restore timeouts. - test_vllm_instrumented_scheduler.py: both sides appended tests at the end of the file. This branch's engine provenance and measurement tests come first, then main's FPM worker_id propagation tests (#12545). Main's vLLM 0.31.0 bump deleted test_vllm_kv_cache_metadata_compat.py, which the realseed_prefix_cache fixture's comment named as the importorskip probe that its function-local KVCacheManager import protects. The comment now names test_vllm_dcp_kv_events.py, which still probes vllm.v1.core.kv_cache_manager. No code change for vLLM 0.31.0: the vLLM interfaces this branch uses (collective_rpc by method name, worker extension mixing, cudagraph_metrics, CUDAGraphStat, the stat loggers, the fields the engine probe and the provenance capture read) are unchanged from 0.30.0. Signed-off-by: Yiming Liu <yimingl@nvidia.com> Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Overview:
In snapshot/restore mode the vLLM
EngineCorechild is created before the Dynamo runtime exists, soInstrumentedSchedulerreadsDYN_FPM_WORKER_IDas"". The parent learns its realworker_idonly after restore, but the already-running child never re-reads it. The restored engine therefore publishes every forward-pass metric (FPM) under an empty id, the FPM subscriber drops those samples as coming from no known worker, and the planner cannot see restored replicas, so it never scales them down. More details in #12591.This PR adds the missing channel: a
set_fpm_worker_idutility onEngineCorethat the parent calls over vLLM's existing utility RPC after restore, so the running child picks up the real id. If that sync fails, the restore fails.Details:
instrumented_scheduler.py— the utility, installed in the child_install_fpm_worker_id_utility()runs at module import and attachesset_fpm_worker_idtovllm.v1.engine.core.EngineCore. It reaches the child becauseEngineCore.__init__resolves the--scheduler-clsthat Dynamo already sets (get_scheduler_cls()→resolve_obj_by_qualname()→import_module()).scheduler._fpm_worker_id(stamped on active samples) andscheduler._publisher._worker_id(stamped on idle heartbeats). An idle replica is exactly the scale-down case.InstrumentedScheduler, so the parent sees the failure instead of a silent no-op.worker_factory.py— the caller in the parent_sync_fpm_worker_id()runs on all three snapshot-restore paths (_create_realtime_worker,_create_decode_worker,_create_prefill_worker), afterfactory.bind_endpoint(...)and before model registration.InstrumentedScheduler. That only happens with a user-supplied--scheduler-cls, which has no FPM publisher.asyncio.wait_for(30 s), becausecall_utility_async()has no deadline of its own._DecodeWorkerLifecycle; realtime and prefill through_sync_fpm_worker_id_or_shutdown().TODOin_create_prefill_worker.Cold start is untouched:
main.pysetsDYN_FPM_WORKER_IDbefore engine construction, so no RPC is made.Validation:
Unit tests — PR head, vLLM 0.30.0 (the pinned version)
test_vllm_worker_factory.py+test_vllm_instrumented_scheduler.py: 370 passed. pre-commit is clean._ReachedRegistrationmarker raised by the stubbedregister_vllm_model. The oldcontextlib.suppress(Exception)is gone, and removing the sync from any one path makes these tests fail.test_every_restore_path_syncs_the_endpoint_id_before_registration: realtime, decode and prefill each send the endpoint id, before registering.test_failed_sync_fails_the_restore_before_registration_and_shuts_down_the_engine: 3 paths × {rejected, timeout, unresolvable scheduler class}. The restore fails, the model is never registered or served, andengine_client.shutdownis called exactly once.test_restore_without_instrumented_scheduler_registers_without_syncandtest_cold_start_registers_without_syncing_the_fpm_worker_id: both register without an RPC, on all 3 paths.test_fpm_utility_via_vllm_dispatch_retargets_active_and_heartbeat_idssends a utility request throughEngineCoreProc._handle_client_request, then checks the id on an active sample and on a real idle heartbeat received over ZMQ.getattron the engine; a failed utility comes back as a bareException(failure_message);call_utility_asynchas no deadline; DP internal and hybrid load balancing send the call to every local rank, and external load balancing has exactly one rank per client.End-to-end acceptance on Kubernetes — run by @RealNicolasBourbaki
Setup:
kubernetes-operator-nightly:20260929-51b83df, with checkpoint enabled.vllm-runtime-nightly:20260929-51b83df(vLLM 0.30.0) plus Dynamo's checkpoint-placeholder layer. The patched image adds this PR's diff; the control image doesn't. The nightly predates938d89b(kvwarm warm-up for hybrid models) ininstrumented_scheduler.py, which is unrelated to this fix.mode: agg,optimization_target: throughput,min_endpoint: 1;replicas: 2,experimental.checkpoint,startupPolicy: WaitForCheckpoint. Both workers start from the snapshot.Results:
worker_idinstance_id""worker_idunder load (90 s, 32 concurrent)instance_id(2,232 samples)""(~13k samples)2 -> 1HOLD … decode=2 … load_reason=no_fpm_datafor the whole runPlanner log, patched:
On the control run the subscriber drops every sample, so the planner logs
no_fpm_datarather thanworker_count_mismatch. The scaling outcome is the same: restored replicas are never scaled down.Where should the reviewer start?
components/src/dynamo/vllm/worker_factory.py:_snapshot_uses_instrumented_scheduler,_sync_fpm_worker_idand_sync_fpm_worker_id_or_shutdown, plus the three call sites.components/src/dynamo/vllm/instrumented_scheduler.py(end of file): the utility, about 20 lines.components/src/dynamo/vllm/tests/test_vllm_worker_factory.py: the restore-path harness (_enter_restore_path,_ReachedRegistration).Related Issues
🔗 This PR is linked to an issue: