Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
23 changes: 23 additions & 0 deletions components/src/dynamo/vllm/instrumented_scheduler.py
Original file line number Diff line number Diff line change
Expand Up @@ -106,6 +106,7 @@
from vllm.v1.core.sched.async_scheduler import AsyncScheduler
from vllm.v1.core.sched.output import CachedRequestData, NewRequestData, SchedulerOutput
from vllm.v1.core.single_type_kv_cache_manager import CrossAttentionManager
from vllm.v1.engine.core import EngineCore
from vllm.v1.kv_cache_interface import MambaSpec
from vllm.v1.request import Request, RequestStatus

Expand Down Expand Up @@ -6812,3 +6813,25 @@ def _bench_write_results(self) -> None:
dest,
len(self._bench_results),
)


# TODO(upstream-vllm): remove once vLLM exposes a way to update scheduler
# identity after engine construction. In snapshot mode the engine is built
# before the Dynamo runtime exists, so the FPM worker_id is baked as "".
def _install_fpm_worker_id_utility() -> None:
if hasattr(EngineCore, "set_fpm_worker_id"):
return

def set_fpm_worker_id(self, new_worker_id: str) -> None:
scheduler = self.scheduler
if not isinstance(scheduler, InstrumentedScheduler):
raise RuntimeError(
f"scheduler is {type(scheduler).__name__}, not InstrumentedScheduler"
)
Comment thread
RealNicolasBourbaki marked this conversation as resolved.
scheduler._fpm_worker_id = new_worker_id
scheduler._publisher._worker_id = new_worker_id

EngineCore.set_fpm_worker_id = set_fpm_worker_id


_install_fpm_worker_id_utility()
175 changes: 175 additions & 0 deletions components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,9 @@

import hashlib
import json
import subprocess
import sys
import textwrap
import threading
import time
import uuid
Expand Down Expand Up @@ -6855,3 +6858,175 @@ def test_random_kda_allows_hybrid_warm_chains_without_expert_parallelism(monkeyp
stub._bench_random_kda = True
assert stub._kvwarm_warm_eligible()
assert stub._kvwarm_meta["skip_reason"] is None


# ---------------------------------------------------------------------------
# FPM worker_id propagation into the EngineCore child (snapshot restore)
# ---------------------------------------------------------------------------

FPM_UTILITY_NAME = "set_fpm_worker_id"


def _fpm_utility():
from vllm.v1.engine.core import EngineCore

return getattr(EngineCore, FPM_UTILITY_NAME)


def _fpm_scheduler_stub(worker_id: str = ""):
"""``InstrumentedScheduler`` carrying only the two FPM identity fields."""
scheduler = object.__new__(InstrumentedScheduler)
scheduler._fpm_worker_id = worker_id
scheduler._publisher = SimpleNamespace(_worker_id=worker_id)
return scheduler


def test_fpm_utility_installed_on_engine_core_base_class():
"""Patched on the base class, so every EngineCore variant inherits it."""
from vllm.v1.engine.core import EngineCore, EngineCoreProc

assert FPM_UTILITY_NAME in vars(EngineCore)
assert hasattr(EngineCoreProc, FPM_UTILITY_NAME)


def test_fpm_utility_install_is_idempotent():
"""A second install must not rebind an already-patched class."""
before = _fpm_utility()

instrumented_scheduler_module._install_fpm_worker_id_utility()

assert _fpm_utility() is before


def test_fpm_utility_updates_scheduler_and_publisher():
"""Active samples use the scheduler's id; idle heartbeats use the publisher's."""
scheduler = _fpm_scheduler_stub()
engine_core = SimpleNamespace(scheduler=scheduler)

_fpm_utility()(engine_core, "8465209922961459")

assert scheduler._fpm_worker_id == "8465209922961459"
assert scheduler._publisher._worker_id == "8465209922961459"


def test_fpm_utility_overwrites_a_previously_set_id():
"""A pod may be restored more than once; the id must follow the new runtime."""
scheduler = _fpm_scheduler_stub(worker_id="1111111111111111")
engine_core = SimpleNamespace(scheduler=scheduler)

_fpm_utility()(engine_core, "2222222222222222")

assert scheduler._fpm_worker_id == "2222222222222222"
assert scheduler._publisher._worker_id == "2222222222222222"


@pytest.mark.parametrize(
"scheduler", [None, SimpleNamespace()], ids=["missing", "foreign"]
)
def test_fpm_utility_rejects_non_instrumented_scheduler(scheduler):
"""Raise rather than no-op, so the parent sees the failure."""
engine_core = SimpleNamespace(scheduler=scheduler)

with pytest.raises(RuntimeError, match="not InstrumentedScheduler"):
_fpm_utility()(engine_core, "8465209922961459")


def test_fpm_utility_argument_is_not_msgspec_converted():
"""vLLM converts msgspec.Struct-annotated args; the id must stay a plain str."""
from inspect import isclass, signature

import msgspec

annotation = signature(_fpm_utility()).parameters["new_worker_id"].annotation

assert not (isclass(annotation) and issubclass(annotation, msgspec.Struct))


@pytest.mark.slow
@pytest.mark.timeout(300)
def test_scheduler_cls_resolution_installs_the_patch():
"""Resolving ``--scheduler-cls`` is what installs the patch in the child.

Runs in a fresh interpreter: this module already imported the scheduler.
"""
script = textwrap.dedent(
"""
from vllm.utils.import_utils import resolve_obj_by_qualname
from vllm.v1.engine.core import EngineCore

assert not hasattr(EngineCore, "set_fpm_worker_id"), (
"patch present before scheduler_cls resolution"
)

resolved = resolve_obj_by_qualname(
"dynamo.vllm.instrumented_scheduler.InstrumentedScheduler"
)

assert resolved.__name__ == "InstrumentedScheduler"
assert hasattr(EngineCore, "set_fpm_worker_id"), (
"resolving scheduler_cls did not install the FPM worker_id utility"
)
"""
)

result = subprocess.run(
[sys.executable, "-c", script],
capture_output=True,
text=True,
timeout=280,
)

assert result.returncode == 0, result.stderr


def test_fpm_utility_via_vllm_dispatch_retargets_active_and_heartbeat_ids():
"""Through vLLM's own utility dispatch, both FPM payload kinds carry the new id."""
import queue

import zmq
from vllm.v1.engine import EngineCoreRequestType
from vllm.v1.engine.core import EngineCoreProc, EngineShutdownState

from dynamo.common.forward_pass_metrics import decode

ctx = zmq.Context.instance()
sub = ctx.socket(zmq.SUB)
sub.setsockopt(zmq.SUBSCRIBE, b"")
port = sub.bind_to_random_port("tcp://127.0.0.1")
sub.unbind(sub.getsockopt(zmq.LAST_ENDPOINT))
publisher = instrumented_scheduler_module._FpmPublisherThread(
f"tcp://127.0.0.1:{port}", worker_id="", dp_rank=0
)
sub.connect(f"tcp://127.0.0.1:{port}")
try:
scheduler = object.__new__(InstrumentedScheduler)
scheduler._fpm_worker_id = ""
scheduler._fpm_dp_rank = 0
scheduler._publisher = publisher
engine = object.__new__(EngineCoreProc)
engine.scheduler = scheduler
engine.shutdown_state = EngineShutdownState.RUNNING
engine.output_queue = queue.Queue()

engine._handle_client_request(
EngineCoreRequestType.UTILITY,
(0, 7, FPM_UTILITY_NAME, ("8465209922961459",)),
)

_client, outputs = engine.output_queue.get_nowait()
assert outputs.utility_output.failure_message is None
active = InstrumentedScheduler._extract_metrics(
scheduler,
None,
None,
0.0,
scheduled=instrumented_scheduler_module.ScheduledRequestMetrics(),
)
assert active.worker_id == "8465209922961459"
assert sub.poll(timeout=5000), "no idle heartbeat within 5s"
heartbeat = decode(sub.recv_multipart()[2])
assert heartbeat is not None
assert heartbeat.worker_id == "8465209922961459"
finally:
publisher.shutdown()
sub.close(linger=0)
Loading
Loading