Skip to content

[KVConnector][P2P] Register peers off the scheduler thread - #57139

Open
liranschour wants to merge 9 commits into
vllm-project:mainfrom
liranschour:fix/55179-p2p-async-peer-registration
Open

liranschour wants to merge 9 commits into
vllm-project:mainfrom
liranschour:fix/55179-p2p-async-peer-registration

Conversation

@liranschour

Copy link
Copy Markdown
Contributor

Addresses #55179.

Problem

P2PSession._on_connect called NixlTransport.add_remote_peer inline before replying with ConnectAckMsg. Registration is add_remote_agent plus a prep_xfer_dlist over every block of the peer's region — 19,599 blocks per rank for the 64 GiB CPU tier in the reported deployment — so a cold handshake ran that work inside a scheduling iteration on both peers, and the lookup queued behind the acknowledgement waited for it.

Scope: registration only, not the per-step gating

The issue proposes two independent changes. This PR implements only the first, because the reported numbers point at registration as the dominant term.

Within one scheduler step, TieringManager.on_schedule_end runs get_finished_jobs() → _poll_once() (drain + dispatch) before serve_external_requests(parent), so receive-and-answer collapse into a single tick. The busy source's critical path is therefore about four ticks: ConnectMsg → ack, LookupMsg → LookupRespMsg, FetchMsg → write_blocks, and NIXL completion → TransferDoneMsg. The destination was idle in the measurement, so its ticks are effectively free.

That gives 13.6 s ≈ 2·T_reg + 4·T_step. At the issue's own stated ~1 s/step, per-step gating accounts for ~4 s, leaving ~9.6 s for two registrations. For gating alone to explain 13.6 s you would need ~3.4 s/step, which contradicts the same description. Per-step gating is real as a mechanism but is the minority term here, so it is deliberately left out of this PR and should be justified by its own measurement.

Change

  • DataTransport.add_remote_peer_async — new method whose default implementation runs add_remote_peer inline and returns an already-resolved future. Transports with cheap registration are unaffected and need no override.
  • NixlTransport overrides it with a single-worker ThreadPoolExecutor, mirroring the NIXL connector's _handshake_initiation_executor. add_remote_peer is split into _build_peer_handles (the expensive agent wireup plus per-block descriptors, run unlocked — holding a lock across it would stall the scheduler thread for exactly as long as the work being moved off it) and _publish_peer (the atomic install).
  • P2PSession starts the registration in _on_connect and sends ConnectAckMsg from poll() once the future resolves; a failed registration marks the connection dead so the manager reaps the session.

The acknowledgement is the safety barrier

ConnectAckMsg is what clears the peer's _send_ready gate, so the peer cannot send a FetchMsg until it arrives. Gating the ack on registration completion preserves the invariant that any peer permitted to fetch is a peer already registered here, so write_blocks can never observe an unregistered peer as a result of this change. The wire protocol is unchanged.

Generation guard

Registrations carry a per-peer generation, bumped on every add and every remove. A worker that finishes building handles under a stale generation releases them instead of publishing. This closes the reap race: _reap_dead_sessions calls session.close() and then _data.remove_remote_peer(pid), and without the guard a registration landing after that would leave a registered peer with no session plus leaked NIXL handles. A pending registration also counts as session pending work, so the engine keeps ticking until the acknowledgement goes out.

No NIXL_THREAD_SYNC_STRICT

The issue suggests opening the agent with NIXL_THREAD_SYNC_STRICT or _RW. That is not reachable through nixl_agent_config: the Python wrapper hard-codes agent_config.syncMode, selecting STRICT only when enable_listen is true, which would additionally start a listener thread and bind a port that this tier does not use.

It is also unnecessary. NixlConnector already runs add_remote_agent, get_xfer_descs and prep_xfer_dlist on its single-worker handshake executor concurrently with transfer submission on the main thread, with the default syncMode=NONE. This PR adopts the threading model vLLM already ships rather than introducing a new one, which also removes the need to measure STRICT's cost on the transfer path.

Why this does not duplicate an existing PR

Draft #55962 covers this ground as part of a larger change that also services the control plane during model steps, and does so by adding a poll_pending_work hook to KVConnectorBase_V1 plus changes to EngineCore and both executors.

This PR is materially different in approach and scope:

  • No new connector API. Nothing is added to KVConnectorBase_V1. The audit behind this is that only ServerRole.serve_external_requests and its four helpers need the ParentManager handle — the whole client role, the entire message-dispatch table, and all NIXL submit/poll/cancel work are already parent-free, and on_lookup already defers resolution to the existing serve_external_requests window. No fresh scheduler-thread window is required.
  • Registration only. It does not attempt the per-step gating change, for the reasons in the scope section above.
  • Different files. It touches none of vllm/v1/engine/core.py, vllm/v1/executor/{abstract,uniproc_executor,multiproc_executor}.py, kv_connector/v1/base.py, multi_connector.py, or offloading_connector.py.
  • No STRICT sync mode, per the section above.

Landing this first also makes #55962's remaining half measurable on its own terms: with registration off the critical path, whatever delay is left on a busy source is the per-step gating component.

Tests

Run from venv/ (this checkout's working environment):

python -m pytest tests/v1/kv_offload/tiering/ -q
  → 377 passed, 11 skipped
  (baseline on the untouched tree: 371 passed, 11 skipped; +6 new tests)

python -m pytest tests/v1/kv_offload/tiering/p2p/ -q
  → 220 passed

pre-commit run --files <6 changed files>          → all hooks pass (incl. DCO sign-off)
pre-commit run mypy-3.12 --files <3 source files> --hook-stage manual → passed

New tests, one behavior each:

  • ConnectAckMsg is withheld while registration is in flight, and sent once it resolves.
  • A pending registration counts as pending work — without this nothing drives the acknowledgement out.
  • A failed registration rejects the peer instead of acknowledging it.
  • A duplicate ConnectMsg mid-registration is a protocol error.
  • Registration provably runs off the caller thread.
  • A registration landing after remove_remote_peer is discarded and its handles released.

The session test fake gained a defer_registration mode so the acknowledgement gating is exercised deterministically rather than relying on thread timing.

Model evaluation

No eval was run, because this change cannot affect model output: it reorders when a control-plane acknowledgement is sent and moves NIXL registration to another thread. Which blocks are transferred, and their contents, are untouched, and the wire protocol is unchanged.

Performance validation is still outstanding. The change targets a timing property that unit tests cannot measure, and no multi-pod deploy has been run yet. The measurement that validates it is TTFT on a busy source with a cold peer versus a warm peer: if cold-peer TTFT drops toward warm-peer, registration was the dominant cost. This PR stays in draft until that result is in.

AI assistance

AI assistance (Claude Code) was used for the analysis, implementation and tests in this PR. The submitting human is reviewing every changed line and is running the hardware validation described above before marking it ready for review.

🤖 Generated with Claude Code

P2PSession._on_connect called NixlTransport.add_remote_peer inline before
replying with ConnectAckMsg. Registration is add_remote_agent plus a
prep_xfer_dlist over every block of the peer's region — 19,599 blocks per
rank for a 64 GiB CPU tier — so a cold handshake stalled a scheduling
iteration on both peers and delayed the lookup queued behind the
acknowledgement.

Add DataTransport.add_remote_peer_async, defaulting to the existing inline
call so transports with cheap registration are unaffected. NixlTransport
overrides it with a single-worker executor, mirroring the NIXL connector's
_handshake_initiation_executor. The session starts the registration in
_on_connect and sends ConnectAckMsg from poll() once the future resolves;
a failed registration marks the connection dead instead. Because the peer
treats the acknowledgement as permission to send, a peer that may fetch is
still a peer we have already registered — the wire protocol is unchanged.

Registrations carry a per-peer generation, bumped on every add and remove,
so a worker that finishes building handles for a peer that was reaped or
re-registered meanwhile releases those handles instead of reviving a stale
peer. A pending registration counts as session pending work so the engine
keeps ticking until the acknowledgement goes out.

Signed-off-by: Liran Schour <lirans@il.ibm.com>
@nilig

nilig commented Sep 16, 2026

Copy link
Copy Markdown

I ran this PR on the same cell and case chain I used for #55962, to check whether registration alone covers #55179.

Setup: GLM-5.2-FP8, two prefill pods (DP8/TP1, 8xH200 each) and a decode pair, MultiConnector(NixlConnector, OffloadingConnector) with a 64 GiB CPU tier per rank and the P2P secondary tier, NIXL 1.4.1. Image = my 6f91edf9 hotfix stack plus a clean backport of this PR at 7b077f36a, nothing else. The metric is the destination's kv_offload_lookup_async_delay_seconds: time from the first deferred lookup until the scheduler observes a resolved result. The busy cases load the source rank with three 115K-token prefills before the pull. One boot, n=1 per case.

case baseline (no patch) #55962 as first submitted (async registration active, wait callback inert) this PR #55962 + c5bfb39d3 (wait callback active)
first pull, cold destination 10.5-10.7 s 10.63 s 10.56 s 10.74 s
idle-warm 1 0.63 s 0.63 s 0.58 s 0.57 s
idle-warm 2 0.94 s 0.94 s 0.59 s 0.91 s
busy-warm 1 2.46 s 2.47 s 1.96 s 0.65 s
busy-warm 2 2.88 s 2.05 s 2.17 s 0.60 s
busy-cold 0.99-3.19 s 3.19 s 3.69 s 0.88 s

The busy rows stay in the baseline band with this PR, the same as the registration-only #55962 arm, and drop to 0.6-0.9 s only once the control plane is serviced during the model wait. Idle rows are the same on every tree.

On the model in the description: I measured registration directly with trace stamps on both peers in an earlier round, and add_remote_agent plus prep_xfer_dlist for 19,599 blocks took 14-40 ms per side, not seconds. The 13.6 s in the issue is the cold first pull, and that turned out to be the destination rank's first execute_dummy_batch() compiling kernels (#57133); it is the same ~10.5 s on every tree here, this PR included, and no control-plane change can remove it. What #55179's busy stall is made of, per the same trace, is step waits on the source: the LookupMsg sits in the server queue until the next step (0.55-0.61 s), the FetchMsg is dispatched at the next step (0.45-0.50 s), and each handshake hop costs a step while the source is busy. That is why moving registration off the thread leaves the busy numbers where they were.

So I would keep this change (the reap-race generation guard is worth having on its own), but not as the fix for #55179: that needs the per-step gating addressed.

liranschour added a commit to liranschour/vllm that referenced this pull request Oct 1, 2026
…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>
@liranschour
liranschour marked this pull request as ready for review October 1, 2026 08:38

@claude claude Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Claude Code Review

This pull request is from a fork — automated review is disabled. A repository maintainer can comment @claude review to run a one-time review.

ruff's pydocstyle D413 requires a blank line after the last docstring
section. The new add_remote_peer_async docstring ended its Returns
block directly against the closing quotes, which made the pinned
ruff-check pre-commit hook rewrite the file and fail CI.

Matches the existing style of every other Returns: block in this file.

Signed-off-by: Liran Schour <lirans@il.ibm.com>

@Etelis Etelis left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM

@Etelis Etelis added the ready ONLY add when PR is ready to merge/full CI is needed label Oct 7, 2026
@Etelis

Etelis commented Oct 7, 2026

Copy link
Copy Markdown
Contributor

/ci run

@github-actions

github-actions Bot commented Oct 7, 2026

Copy link
Copy Markdown

✅ Triggered Buildkite CI #93394 for commit f8c18ce0147b.

@github-actions

github-actions Bot commented Oct 7, 2026

Copy link
Copy Markdown

✅ @liranschour, CI is now available for this PR.

  • /ci run starts upstream CI; /amd-ci run starts AMD CI only.
  • Your branch must contain every commit currently on its upstream target branch. Merge or rebase onto the latest target branch, then rerun the command. Append --allow-stale to a run command to test an outdated branch at your own risk.
  • /ci retry retries failed jobs in the CI build for the current PR head. If the current head has no CI build, it starts a new CI build for the current head containing only jobs that failed in the latest earlier CI build for this PR.
  • /amd-ci retry retries failed jobs in AMD CI for the current PR head. Use /amd-ci run when the current head has no AMD CI build.
  • /ci cancel cancels scheduled or running CI builds for this PR branch; /amd-ci cancel does the same for AMD CI only.

@vllm-agent

vllm-agent commented Oct 7, 2026 •

Copy link
Copy Markdown
Contributor

CI selector (shadow): 1 test steps (1 jobs) instead of 62 (78 jobs)

Shadow mode: this changes nothing about what CI runs. It shows what the evidence-based selector would pick for this PR, next to today's rules. How it works.

Feedback welcome: reply here if it would skip a step this change needs, or runs something unrelated.

steps (jobs) Today's rules Selector Would skip Would add
NVIDIA, CPU and others 62 (78) 1 (1) 61 (77) 0 (0)
AMD mirrors 54 (65) 1 (1) 53 (64) 0 (0)
Selector would run (1)
  • v1-kv-offload
Would skip (today's rules run them) (61)
  • ascend-npu-test
  • async-engine-inputs-utils-worker
  • basic-correctness ×2
  • basic-correctness-cpu-offload
  • basic-correctness-cumem
  • basic-correctness-prefetch-offload
  • basic-correctness-sleep-mode
  • basic-models-test-other-cpu
  • basic-models-tests-initialization
  • basic-models-tests-other
  • benchmarks-cli-test
  • cpu-language-generation-and-pooling-model-tests ×3
  • cpu-multimodal-config
  • cpu-params-env-tokenizers-parser
  • cpu-reasoning-renderers
  • cpu-tool-parsers
  • e2e-core-1-gpu
  • e2e-core-large-memory
  • e2e-scheduling-1-gpu
  • e2e-scheduling-accuracy-1-gpu
  • entrypoints-integration-api-server ×4
  • entrypoints-integration-api-server-generate
  • entrypoints-integration-api-server-openai-chat_completion
  • entrypoints-integration-api-server-openai-completion
  • entrypoints-integration-llm
  • entrypoints-integration-multimodal
  • entrypoints-integration-pooling
  • entrypoints-integration-responses-api
  • entrypoints-integration-speech_to_text
  • kernels-fla-ops-test-b200
  • kernels-mhc-test-b200
  • kernels-root-misc-test-b200
  • kv-offload-large
  • kv-offload-medium
  • kv-offload-small
  • language-models-tests-granite-l4-compatibility
  • language-models-tests-hybrid ×2
  • language-models-tests-standard
  • metrics-tracing-2-gpus
  • multi-modal-models-standard-1-qwen2
  • multi-modal-models-standard-2-qwen3-gemma
  • multi-modal-models-standard-3-llava-qwen2-vl
  • multi-modal-models-standard-4-other-whisper
  • multi-modal-processor ×4
  • multi-modal-processor-cpu ×4
  • pytorch-compilation-dynamic-shapes
  • pytorch-compilation-passes-unit-tests
  • pytorch-compilation-unit-tests
  • pytorch-compilation-unit-tests-h100
  • pytorch-fullgraph-cudagraph-l4-compatibility
  • pytorch-fullgraph-test
  • regression
  • rl-entrypoints-tests
  • v1-core
  • v1-executor-worker
  • v1-kv-connectors ×4
  • v1-logits-oracle
  • v1-metrics-lmeval
  • v1-others-cpu
  • v1-sample
  • v1-spec-decode
Would add (today's rules do not run them) (0)

none

AMD mirrors: would skip (53)
  • async-engine-inputs-utils-worker
  • basic-correctness ×2
  • basic-correctness-cpu-offload
  • basic-correctness-cumem
  • basic-correctness-prefetch-offload
  • basic-correctness-sleep-mode
  • basic-models-tests-initialization
  • basic-models-tests-other
  • benchmarks-cli-test
  • e2e-core-1-gpu
  • e2e-core-large-memory
  • e2e-scheduling-1-gpu
  • e2e-scheduling-accuracy-1-gpu
  • entrypoints-integration-api-server ×4
  • entrypoints-integration-api-server-generate
  • entrypoints-integration-api-server-openai-chat_completion
  • entrypoints-integration-api-server-openai-completion
  • entrypoints-integration-llm
  • entrypoints-integration-multimodal
  • entrypoints-integration-pooling
  • entrypoints-integration-responses-api
  • entrypoints-integration-speech_to_text
  • kernels-fla-ops-test-b200
  • kernels-mhc-test-b200
  • kernels-root-misc-test-b200
  • kv-offload-large
  • kv-offload-medium
  • kv-offload-small
  • language-models-tests-granite-l4-compatibility
  • language-models-tests-hybrid ×2
  • language-models-tests-standard
  • metrics-tracing-2-gpus
  • multi-modal-models-standard-1-qwen2
  • multi-modal-models-standard-2-qwen3-gemma
  • multi-modal-models-standard-3-llava-qwen2-vl
  • multi-modal-models-standard-4-other-whisper
  • multi-modal-processor ×4
  • platform-tests
  • pytorch-compilation-dynamic-shapes
  • pytorch-compilation-passes-unit-tests
  • pytorch-compilation-unit-tests
  • pytorch-compilation-unit-tests-h100
  • pytorch-fullgraph-cudagraph-l4-compatibility
  • pytorch-fullgraph-test
  • regression
  • rl-entrypoints-tests
  • v1-core
  • v1-executor-worker
  • v1-kv-connectors ×4
  • v1-logits-oracle
  • v1-metrics-lmeval
  • v1-sample
  • v1-spec-decode
AMD mirrors: would add (0)

none

6 changed files · base 0f112d1800 · head b6b44522b6 · Python record: build 93039 at 1e5d0ea888 · kernel record: table a982c81 (build 93415), map a982c81 · not counted: 10 build steps, 5 A100 steps the generator no longer emits, 7 optional steps the selector would also run

@liranschour

Copy link
Copy Markdown
Contributor Author

/ci run

@github-actions

github-actions Bot commented Oct 8, 2026

Copy link
Copy Markdown

❌ This PR is 14 commits behind upstream main. Your branch must contain every commit currently on upstream main. No new CI build was started. Merge or rebase onto the latest main, then rerun /ci run. To test this branch at your own risk, use /ci run --allow-stale.

@liranschour

Copy link
Copy Markdown
Contributor Author

/ci run

@github-actions

Copy link
Copy Markdown

✅ Triggered Buildkite CI #94025 for commit 1128deede723.

This branch has not been deployed

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

Labels

ready ONLY add when PR is ready to merge/full CI is needed

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants