Repository navigation
[KV Offload] Expose OffloadingManager.lock()/unlock() and service the tiering control plane off the scheduler thread - #58168
Conversation
on_schedule_end is documented as running at the end of each scheduler step, but build_connector_meta called it ~40 lines early, before the preempted-job flush, before _build_store_jobs (which calls prepare_store, and so evicts), and before on_request_finished. Move the call to just before the return so the contract holds. This matters for callers that want to bound a step's exclusive access to the manager: releasing at on_schedule_end today leaves prepare_store and on_request_finished outside the window. prepare_store's poll becomes unconditional. Its fresh poll previously came from the per-step gate being reset mid-build_connector_meta; with on_schedule_end at the end, the gate would still be set by the step's first lookup() and the poll would be skipped, leaving eviction decisions on stale completion data. _flush_pending_cascades now runs after this step's prepare_store, so keys it parked are re-parked and retried next step. That is safe -- _maybe_finalize_request holds off while pending_cascade_keys is non-empty -- and is noted in the code. Test asserts on_schedule_end is the last manager call of every step; without the move it observes ['on_schedule_end', 'prepare_store']. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: Liran Schour <lirans@il.ibm.com>
…rol thread Expose exclusive access to an OffloadingManager as a public API, and use it for the first case that needs it. The API. lock()/unlock(), plus a locked() helper, guard a manager's whole state subtree: the manager, the tiers it composes, and tier front-end state. Manager methods never acquire it themselves, so a caller holds it across as many operations as its invariants need. That matters because several manager operations only make sense in sequence -- prepare_load() documents that its keys are offloaded, a fact lookup() established one hook earlier -- and a caller has to be able to span both. OffloadingConnectorScheduler is the first caller. It holds one region per engine step, opened by the first manager-touching hook and closed at the end of build_connector_meta, so the manager is free while the engine runs the model. Hooks outside schedule() take a short region instead. has_pending_push_work must use the short one: the engine asks it before schedule() and returns early when the answer is False, so a step-scoped acquire there would never reach the release and would hold the manager for good -- the state an idle engine settles into. It does need the lock, since has_pending_work() iterates _req_state while another thread may insert or delete. take_events now drains the manager inside the region instead of lazily. manager.take_events() is a generator whose backing list is cleared once exhausted, so iterating it lazily held the manager for as long as the consumer took and discarded events appended in the meantime, unpublished. The consumer. The p2p tier answers a remote peer, which the engine does not drive, so its control plane only advanced from on_schedule_end and a peer's lookup or fetch waited for a step boundary -- seconds on a saturated rank, against a transfer of well under one. Tiers now opt in with needs_control_plane_thread, and the tiering manager runs one thread that polls them and lets them serve, once per control_plane_interval_s (default 1 ms, 0 disables). _process_finished_jobs takes a tier filter so a round does not drag tiers that expect the scheduler thread along with it. The thread yields unconditionally after unlock(), which is what bounds the scheduler's wait to one round: Python locks are not fair, so a release-then-reacquire loop could starve it. A failed round is logged rather than fatal, and on_schedule_end still serves every tier, so the fallback is per-step servicing. shutdown() stops and joins the thread before any tier teardown, since closing a transport under a live round can end the process rather than raise. Docstrings that claimed single-threaded tier access are corrected, including the eviction-race argument in _pin_and_register_hits, which still holds -- now because serve_external_requests runs inside one hold of the lock. One assumption is recorded rather than resolved, in nixl.py: strict serialization is sufficient only if libnixl and UCX keep no state tied to the creating thread. The agent runs with NIXL_THREAD_SYNC_NONE and the config offers no way to ask for STRICT without also starting NIXL's listen thread. This needs a real two-node soak before it is trusted. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: Liran Schour <lirans@il.ibm.com>
|
Documentation preview: https://vllm--58168.org.readthedocs.build/en/58168/ |
The note was written against the locally installed NIXL 0.6.0 and generalized wrongly. On 1.4.1, which is what deployments run, nixl_agent_config does accept sync_mode, and NIXL_THREAD_SYNC_RW does exist. More importantly the open question is now answered rather than flagged. makeMtType never produces mt_mode_t::CONTEXT, the only value mapping to UCS_THREAD_MODE_SINGLE and so the only UCX mode with creator-thread affinity; the weakest mode NIXL requests is UCS_THREAD_MODE_SERIALIZED, documented as "multiple threads can access, but only one at a time", and this transport's default prog thread gets it MULTI with a hard library check. Serialized use from a second thread is supported. What does need stating is that the manager lock is the only protection, since syncMode NONE makes NIXL's own agent lock a no-op, and that NIXL binds a shared UCX worker per calling thread -- harmless only while there is one shared worker, as there is here. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: Liran Schour <lirans@il.ibm.com>
orozery
left a comment
There was a problem hiding this comment.
cc @mkhazraee wdyt?
This should help kvcr tier as well.
The purpose is to assure serve_external_requests is called more frequently (e.g. 1ms).
Today it is only called from the main thread, which may get to every O(s) if the model runner is busy.
| return self._manager_lock.acquire() | ||
| return self._manager_lock.acquire(timeout=timeout) | ||
|
|
||
| def unlock(self) -> None: |
There was a problem hiding this comment.
Let's replace lock/unlock/locked with a single lock context manager property.
|
|
||
| class OffloadingManager(ABC): | ||
| def __init__(self) -> None: | ||
| self._manager_lock = threading.Lock() |
There was a problem hiding this comment.
I think we should remove this.
The default should be no lock.
The CPU tier will implement its own threading lock.
| # sets how long a peer's request can sit unserved; 1 ms is three orders of | ||
| # magnitude below the model step it hides, and matches the pause the p2p tier's | ||
| # own drain loops and EngineCoreProc's GIL-yield sleep already use. | ||
| _CONTROL_PLANE_INTERVAL_S = 0.001 |
There was a problem hiding this comment.
spec.py imports this, and it's a user-facing default.
So it should be public.
We should check if this is the best place here.
There was a problem hiding this comment.
Made it public as DEFAULT_TIER_POLL_INTERVAL_S.
spec.py imports it from there as the fallback for the user-facing tier_poll_interval_s key.
| # scheduler thread; this pause is what bounds its wait to one round. | ||
| self._control_plane_stop.wait(self._control_plane_interval_s) | ||
|
|
||
| def serve_control_plane(self) -> None: |
There was a problem hiding this comment.
How about we rename serve_control_plane / _control_plane_interval_s to poll_tiers() / tier_poll_interval_s?
| and existing sessions are polled even when no requests are scheduled. | ||
| """ | ||
|
|
||
| needs_control_plane_thread: ClassVar[bool] = True |
There was a problem hiding this comment.
how about serves_external_requests?
Address review feedback on the lock API: - Replace lock()/unlock()/locked() with a single `lock` property returning a context manager. Callers write `with manager.lock:`. - Drop the lock from the base class. The default is a no-op nullcontext() for managers only ever entered from one thread; CPUOffloadingManager and TieringOffloadingManager each own a threading.Lock and expose it. The tiering manager's control-plane loop and shutdown() still need a bounded wait, so they call acquire(timeout=)/release() on the lock they own directly; the timeout is no longer part of the public API. The scheduler's step-scoped region spans schedule() to the end of build_connector_meta(), which no single with-block can cover, so it is now carried by an ExitStack that holds the entered lock between _acquire_step() and _release_step(). Signed-off-by: Liran Schour <lirans@il.ibm.com>
Address review feedback: - serve_control_plane() -> poll_tiers() - control_plane_interval_s -> tier_poll_interval_s, for the constructor argument, the attribute, and the user-facing kv_connector_extra_config key (docs updated). - _CONTROL_PLANE_INTERVAL_S -> DEFAULT_TIER_POLL_INTERVAL_S. spec.py imports it and it is the documented default of a user-facing key, so it is public. It stays in tiering/manager.py next to the thread it paces: spec.py already imports manager.py, so moving it to spec.py would leave the manager's constructor without a default to reference. The remaining control-plane constants stay private; nothing outside the module reads them. Signed-off-by: Liran Schour <lirans@il.ibm.com>
…uests Address review feedback: the opt-in flag now states what the tier does -- answers a counterpart the engine does not drive -- rather than what the manager should do for it, and matches the serve_external_requests() method it gates. The flag's docstring now says that opting in moves both get_finished_jobs() and serve_external_requests() onto the polling thread, and that overriding serve_external_requests() alone does not opt in. _control_plane_tiers is renamed _external_serving_tiers to match. KVCR implements serve_external_requests() but stays opted out, now explicitly: opting in would also poll the KVCR library from the polling thread, which it is not yet verified to tolerate. Signed-off-by: Liran Schour <lirans@il.ibm.com>
A tier served by the polling thread can look keys up through its parent -- the p2p server role does so for every peer lookup -- and TieringOffloadingManager.lookup() runs the once-per-step _maybe_process_finished_jobs(). Between steps the gate is clear, so the first such lookup in a round: - polled every secondary tier from the polling thread, including tiers that did not set serves_external_requests; - re-entered the serving tier's get_finished_jobs() from inside its own serve_external_requests(), which for p2p means _poll_once() adding or reaping sessions while serve iterates over them; - left the gate set, so the next step's lookups and on_schedule_end()'s catch-all poll skipped their poll. poll_tiers() now holds the gate shut for the duration of the round and restores it afterwards. The round has just polled its own tiers, so the skipped poll loses nothing, and the scheduler's per-step gate is left as it found it. Add a regression test with a tier that does a parent lookup while serving, and narrow the opt-in test's docstring: ParentManager fan-out can still reach tiers that did not opt in, as tiering/base.py documents. Signed-off-by: Liran Schour <lirans@il.ibm.com>
…rule The NixlTransport threading note says the manager lock must cover every agent entry point. Peer registration moved onto a worker thread, as NixlConnector's handshake executor already does and as vllm-project#57139 proposes for this transport, deliberately runs outside that lock. Say so, so the note stays accurate when that lands. Signed-off-by: Liran Schour <lirans@il.ibm.com>
|
This pull request has merge conflicts that must be resolved before it can be |
…er-lock-api # Conflicts: # vllm/distributed/kv_transfer/kv_connector/v1/offloading/scheduler.py # vllm/v1/kv_offload/tiering/p2p/manager.py Signed-off-by: Liran Schour <lirans@il.ibm.com>
|
On my reproducer this PR fixes #55179: with a busy source rank the P2P lookup delay goes from 2.8-2.9 s to 0.6-0.7 s, the same value as with an idle source, and nothing else in the run got worse. I ran this PR against the #55179 reproducer on a current nightly (
The busy rows land in the idle band with the PR, and it matches what I measured for the poller approach in #57165, so I think #57165 can be closed in favor of this one. Decode-side TPOT on the probed decode rank ( The first pull after boot is a different problem (#57133: the destination rank's own first dummy batch JIT) and this PR does not change it: 25.5 s without, 25.2 s with, both on fresh nodes. One row I cannot interpret yet: a gateway-routed follow-up with all eight source ranks loaded waited 9.7 s on its lookup with the PR (TTFT 11.8 s, 102K tokens pulled). n=1, no counterpart without the PR on this nightly, no JIT on the destination at that time, no lock-timeout warning from the control-plane thread. I will look at it separately; the direct-injection rows above do not show it. One note on the setup, in case it matters for reproducing: prefill ran with |
Purpose
Two things, in priority order.
1. Expose exclusive access to an
OffloadingManageras a public API.lock()/unlock(), plus alocked()helper, guard a manager's whole state subtree: the manager, the tiers it composes, and any front-end state those tiers own. Manager methods never acquire it themselves, so a caller holds it across as many operations as its own invariants need. That property is the point — several manager operations only make sense in sequence.prepare_load()documents that its keys are offloaded, a factlookup()established one hook earlier, so a caller has to be able to span both. Internal per-method locking cannot offer that.2. Use it for the first case that needs it (#55179). The
p2psecondary tier answers a remote peer, which the engine does not drive, so its control plane only advanced fromon_schedule_end. A peer's lookup or fetch therefore waited for a model-step boundary — seconds on a saturated rank, gating a transfer of well under one. Tiers now opt in withneeds_control_plane_thread, andTieringOffloadingManagerruns one thread that polls them and lets them serve between steps.Not a duplicate of the open PRs on #55179
This overlaps open work and I want that on the record rather than buried.
TieringOffloadingManagerlock()/unlock()onOffloadingManager, held by callers_enter_step/_exit_step/_observation_lock+*_unlockedvariantsEngineCorepoller threadKVConnectorBase_V1.poll_pending_workEngineCore+ connector base API + executor changesThis supersedes my own #58113, which I will close. The difference is not cosmetic:
*_unlockedvariant of every method_SecondaryTierFacingParentcan re-enter (lookup,on_new_request,create_store_job,on_request_finished,reset_cache) because the parent handle's caller already holds it. That duplication is most of its+393lines inmanager.py, and it is precisely what prevents an outside caller from saying "hold this across several operations." Moving the lock to the caller removes every one of those variants.on_schedule_end, which is not the end of the step.build_connector_metacalledon_schedule_end~40 lines early, before_build_store_jobs→prepare_store(which evicts) and beforeon_request_finished. So itsprepare_storere-acquired via_enter_step()and nothing released it: on every step that builds store jobs — the common case under load — the lock was held acrossfuture.result(), starving the control thread in exactly the window the change exists to exploit. That plausibly explains [KVConnector][Tiering] Drive the tiering control plane from a dedicated thread #58113's own "parity-within-noise" result, which was measured through the bug, and that conclusion should be re-tested.Relative to #57165 / #55962: same goal, but no new
KVConnectorBase_V1hook and noEngineCoreor executor changes. In #55179 @nilig asked for a read on "a thread touching the connector from insideEngineCore, even phase-exclusive, versus a control thread owned by the tier" — this is a working implementation of the second option, offered as input to that question. If maintainers prefer either existing PR, close this one.Changes
Commit 1,
on_schedule_endruns last — independent of all locking, reviewable on its own. Its docstring already claimed "called once at the end of each scheduler step"; it was called ~40 lines early. Moving it makes the contract true and gives callers a sound release point. Two ordering consequences, both handled:prepare_store's poll becomes unconditional (the per-step gate was previously reset mid-build_connector_meta, which is where its fresh completion data came from — without this, eviction decisions would run on data up to a step stale), and_flush_pending_cascadesnow re-parks keys this step parked, which is safe because_maybe_finalize_requestholds off whilepending_cascade_keysis non-empty.Commit 2, the API and its consumer.
OffloadingConnectorSchedulerholds one region per engine step, opened by the first manager-touching hook and closed at the end ofbuild_connector_meta, so the manager is free while the engine runs the model. Hooks outsideschedule()take a short region.Two placements are load-bearing and easy to get wrong:
has_pending_push_workmust use the short region.EngineCore.step()returns early atcore.py:597whenhas_requests()is false, so a step-scoped acquire there would never reach the release and would hold the manager for good — the state an idle engine settles into. It does still need the lock:has_pending_work()iterates_req_statewhile the control thread may insert or delete viaparent.on_new_request/on_request_finished.take_eventsdrains inside the region.manager.take_events()is a generator whose backing list is cleared once exhausted, so iterating it lazily both pinned the manager for the consumer's duration and discarded events a concurrentcomplete_writeappended, unpublished. It stays a generator so the caller's truthiness check atsched/scheduler.pyis unchanged.The thread yields unconditionally after
unlock(), which is what bounds the scheduler's wait to one round — Python locks are not fair, so a release-then-reacquire loop could starve it. A failed round is logged, not fatal, andon_schedule_endstill serves every tier, so the fallback is per-step servicing._process_finished_jobstakes a tier filter so a round does not drag tiers that expect the scheduler thread.shutdown()stops and joins the thread before any tier teardown, since closing a transport under a live round can end the process rather than raise.Config:
control_plane_interval_sinkv_connector_extra_config, default 1 ms,0disables the thread and restores per-step servicing.Docstrings that claimed single-threaded tier access are corrected, including the eviction-race argument in
_pin_and_register_hits— which still holds, now becauseserve_external_requestsruns inside one hold of the lock.async_lookup.py's "owned exclusively by the scheduler thread" claim was outright false onceparent.lookupcan fan out tofs.lookupfrom the control thread.Test results
Note for reviewers: a whole-directory run of
tests/v1/kv_offload/produces a large number of CUDA IMA errors that are pre-existing and unrelated (tests/v1/kv_offload/cpu/poisons the rest of the session). Run the subdirectories separately.12 new tests. The two that matter most were both verified to fail without the change, not just to pass with it:
test_on_schedule_end_is_the_last_manager_call_of_a_step— without the reorder it observes['on_schedule_end', 'prepare_store']. This is the check [KVConnector][Tiering] Drive the tiering control plane from a dedicated thread #58113's suite lacks, because its_step()helper driveson_schedule_endwith no followingprepare_store.test_control_thread_and_scheduler_never_overlap— withlock()/unlock()stubbed to no-ops it detects 86 overlaps. It also asserts both threads made progress, so neither a thread that never runs nor one that starves the scheduler passes it.I wrote and then deleted a third test ("no eviction between a served lookup and its pin") after finding it still passed with the lock neutered: LRU keeps evicting the fresh keys rather than the constantly-looked-up hit keys, so it could never fail. The overlap test is the real guard for that property, and a test that cannot fail is worse than no test.
Model evaluation — the 4-case matrix from #55179
Run on two H200 pods, Llama-3.2-1B,
--enforce-eager,--block-size 16, 96 GiB CPUprimary tier,
secondary_tiers: [{"type": "p2p"}],max-model-len 32768.Topology. Source = the pod holding the KV in its CPU offload tier. Consumer = the
other pod, which probes the source over P2P (
remote_kv_source) and fetches. Everymeasured fetch uses a freshly generated ~27k-token prompt that is populated on the source
first, so nothing hits a GPU prefix cache and the fetch is real work: each one moved
~1720 chunks (
kv_offload_tiering_chunk_hitsdelta), i.e. the full prefix.Busy = source saturated with 8 concurrent ~27k-token prefills. That took its
request latency from 0.018 s idle to 1.40 s, so the control plane really is gated
behind ~1.4 s of model work — the same regime as the issue's report.
Cold = no P2P session yet. A session lives as long as the engine, so only the first
fetch after a restart is genuinely cold; both pods were restarted between every phase and
the cold fetch measured first.
Signal = the consumer's
vllm:kv_offload_lookup_async_delay_secondsdelta for themeasured request (one observation per fetch). Warm figures are the mean of 3 reps.
7b077f36a8)The per-step gating penalty, which is what this PR removes:
Saturating the source costs baseline 4.8x on the consumer's lookup delay. With the control
thread it costs 1.17x — the stall is essentially gone, and a busy source now behaves like
an idle one.
Supporting numbers on busy-warm: per-chunk tiering delay
(
kv_offload_tiering_lookup_async_delay_seconds) 0.100 s → 0.0174 s (5.8x), and end-to-endTTFT for the fetch 0.310 s → 0.144 s (2.2x). Idle is parity within noise in every metric,
as expected — there is no model step to hide behind.
Per-rep spread, busy-warm: baseline
0.1929 / 0.2532 / 0.2127, this PR0.0388 / 0.0427 / 0.0786. The distributions do not overlap — the worst rep here is stillbetter than the best baseline rep.
No source-side slowdown. The concern that holding the manager lock slows the scheduler
does not show up: on the busy source, request latency was 1.4174 s (baseline) vs 1.4039 s
(this PR), and the populate-request mean 1.7056 s vs 1.7100 s. The background load
generator completed 158 vs 153 requests over comparable windows, a ~3% difference that is
within single-run noise and not corroborated by either latency measure. A dedicated
vllm bench servethroughput comparison is still worth running before merge; I did notrun one.
On #58113
Arm B is #58113 at
b8ea135c71, measured on the same pods in the same session. It landsbetween baseline and this PR: busy-warm 0.129 s vs baseline 0.220 s vs 0.0534 s here.
That is what the release-point bug predicts. #58113's thread is not dead — on steps that
build no store jobs the lock is released normally, so it recovers some of the win (1.7x
over baseline). But on any step that does build store jobs,
prepare_storere-acquiresafter
on_schedule_endreleased and nothing releases it again, so the lock is held acrossfuture.result()and the thread is shut out for exactly that step. Under load most stepsbuild store jobs, which is why it recovers only about a third of what is available (4.1x
over baseline here, 2.4x over #58113).
I am reporting this as consistent with the bug rather than as direct proof of the
mechanism; the mechanism itself is demonstrated by the unit test, which observes the call
order
['on_schedule_end', 'prepare_store']without the reorder in commit 1.Caveat on absolute magnitudes
These numbers are much smaller than the issue's 13.6 s, and deliberately so: this is
Llama-3.2-1B on a single H200 per side, not GLM-5.2-FP8 on DP8. What reproduces here is the
mechanism and its scaling — the consumer's lookup delay is proportional to the source's
step time under baseline, and roughly independent of it with the control thread. The issue's
13.6 s came from a far slower source; the ratio, not the absolute, is what transfers.
Raw per-rep JSON for all six runs is available on request.
NIXL thread-safety: checked, and it holds
Checked against the NIXL 1.4.1 sources (
ai-dynamo/nixlat tagv1.4.1), which is thedeployed version. This corrects an earlier version of this description, which said
NIXL_THREAD_SYNC_RWdoes not exist and that the config offers no way to request a syncmode. Both statements were true of 0.6.0, which is what I had installed; neither is true of
1.4.1.
Serialized use from a second thread is supported.
makeMtType(
src/plugins/ucx/ucx_utils.cpp:383-392) returns onlyWORKERorSINGLE, neverCONTEXT-- andmt_mode_t::CONTEXTis the only value that maps toUCS_THREAD_MODE_SINGLE(ucx_utils.cpp:476-492), the one UCX mode with creator-threadaffinity. Grepping confirms
CONTEXTnever appears as a produced value, only in switcharms. So the weakest mode NIXL ever requests is
UCS_THREAD_MODE_SERIALIZED, which UCXdocuments as "Multiple threads can access, but only one at a time" -- precisely this
design. And because
enable_prog_threaddefaults on for this transport, it actually getsUCS_THREAD_MODE_MULTI, with a hard library-support check that throws if unavailable.The manager lock is load-bearing, not redundant. With
syncMode = NONE-- which is whatthis transport gets, since it sets neither
enable_listennorsync_mode-- NIXL's internalagent lock compiles to no-ops (
src/core/sync.h:36-39), leaving ~23 agent entry pointsunguarded over shared, non-atomic containers. The external mutex is therefore the only
protection, and it must cover every entry point plus the lifetime of handles passed between
calls. It does.
Two caveats worth recording:
tlsSharedWorkerMap,ucx_backend.cpp:821-825). That is harmless only while there is exactly one sharedworker, which holds for this configuration (
num_workersunset ⇒numSharedWorkers_ == 1,so every thread maps to worker 0) but would not if a
num_workersUCX backend param wereever set. The chosen worker also travels on the request handle rather than being
re-derived per call, so a transfer prepped on one thread and polled on another still uses
the same worker.
sync_mode(includingNIXL_THREAD_SYNC_RW, value 2),which would make NIXL's own locking real as belt-and-braces, and is what every upstream
multi-threaded test uses. I have not added it, because the parameter does not exist on
0.6.0 and passing it unconditionally would break older environments; it needs the same
kind of version gate
capture_telemetryalready needs. Worth a follow-up.Honest limit:
~/nixlis a 1.4.1 source checkout, not a build, and only 0.6.0 is installedlocally, so every 1.4.1 claim above is static source reading rather than runtime
introspection. A runtime confirmation (log the thread mode
ucp_worker_queryactuallygrants, or run the two-thread pattern under TSan) would close the remaining gap. The
end-to-end matrix above ran on the deployed 1.4.1 across ~14k P2P chunk fetches with no
transfer failures, which is evidence but not a substitute.
pyzmq is fine under strict serialization: libzmq permits migrating a socket between threads
given a full memory barrier, which lock acquire/release provides, and nothing is bound
thread-locally at socket creation. The one hazard is
zmq_ctx_destroyracing a socket call,which the shutdown ordering handles.
AI assistance
AI assistance was used for this change (Claude Code): the code-path analysis, the implementation, and the tests. I am reviewing every changed line. The two items that were blocking this leaving draft are now done: the 4-case matrix has been run on H200 hardware (results above) and the NIXL question has been checked against 1.4.1. What remains before undraft is a dedicated throughput comparison and a human pass over the diff.
🤖 Generated with Claude Code