Repository navigation
[Bugfix] Reclaim unconsumed SHM segments and sender state on request abort - #4349
gagandhakrey wants to merge 14 commits into
Conversation
Signed-off-by: Gagan Dhakrey <gagandhakrey@gmail.com>
|
Codex usage limits have been reached for code reviews. Please check with the admins of this repo to increase the limits by adding credits. |
|
@princepride @yuanheng-zhao @yenuo26 |
yuanheng-zhao
left a comment
There was a problem hiding this comment.
Thanks for contributing! Left some comments
| "connector_cleanup": True, | ||
| "external_req_id": external_req_id, | ||
| # save_loop's error log reads task["request_id"]. | ||
| "request_id": external_req_id, |
There was a problem hiding this comment.
Not sure if these three keys are necessary. Shall we shrink the code, for example, appending {"connector_cleanup": external_req_id} and update related code about sending request?
There was a problem hiding this comment.
could merge connector_cleanup/external_req_id into one key , keeping request_id though save_loop's error log relies on it for a useful id on failure, minor style choice, can shrink if you want
|
PTAL @Shirley125 , |
…on-abort # Conflicts: # vllm_omni/distributed/omni_connectors/transfer_adapter/chunk_transfer_adapter.py Signed-off-by: Gagan Dhakrey <gagandhakrey@gmail.com>
…o fix/shm-segment-leak-on-abort Signed-off-by: Gagan Dhakrey <gagandhakrey@gmail.com>
request is already known non-None and Request.external_req_id is always a real attribute, so getattr(..., None) is dead defensiveness. Signed-off-by: Gagan Dhakrey <gagandhakrey@gmail.com>
0dfcea1 to
b9902ed
Compare
|
Production data point supporting this PR: we independently root-caused the sender-state leak on a 2-stage TTS pipeline (details in #6352). Constant ~45 KB per aborted request across 4 configs; sustained-load A/B shows 258 → 60 MiB/h with a minimal abort-time |
|
Thanks for this — the leak is real, and the save-queue FIFO design is the Request changes:
Should-fix:
The FIFO enqueue, the “don't unlink on natural finish” invariant, and |
…on-abort Conflicts resolved: * shm_connector._get_data_with_lock: main independently fixed the same lock-file leak by gating removal on a `deserialized` flag. Kept this branch's `consumed` gate (the segment is already unlinked by shm_read_bytes, so the lock file must go even when deserialization fails and no retry can succeed), with main's type annotations and FileNotFoundError handling. * chunk_transfer_adapter.finish_requests: took main's restructured skip logic (`request is None` plus the resumable-segment-stop carve-out from vllm-project#6360) and its cleanup_receiver()/_held_non_active handling, while keeping the aborted_external_ids collection and the connector-cleanup enqueue. A resumable segment stop is now fully reclaimed here, so it is swept like any other abort. Signed-off-by: Gagan Dhakrey <gagandhakrey@gmail.com>
The branch tip already carried a merge of an older main (3c62c53), which is a subset of the main merged in the previous commit. No content change. Signed-off-by: Gagan Dhakrey <gagandhakrey@gmail.com>
Signed-off-by: Gagan Dhakrey <gagandhakrey@gmail.com>
|
Thanks @amy-why-3459 for detailed review All three blockers are done, plus the should-fixes. Blockers (Resolved)1. Rebased — now
2. Gated on
3. Should-fixes (done)Key matching — now structural (
Task shape — one correction to your note.
On #6352
Ready for re-review. (this comment is AI polished) |
Omni ReviewBot triage noteAutomated triage of commit
These are automated triage suggestions only — the final decision belongs to the maintainers. |
Omni ReviewBot: no human activity for 8 days@gagandhakrey this pull request has had no human commit, comment or review since 2026-08-31. Per repository policy it may be closed if it stays inactive. To keep it moving, any one of these is enough: push an update, reply to the open blocker, or post the current plan and timeline. |
|
We carried this PR's approach into a v0.24 deployment for pre-production load testing Production-scale throughput and volume figures are omitted as a matter of policy. Everything load-bearing below is normalized — per abort, per request, per token — which is also the form that makes the invariance argument work. The load range quoted for the four configurations is the one already published in #6352. Stage-0 sender-abort test (review point 3)That path is the one our load testing reproduces. Conditions:
Retained state per abort is one utterance's codec frame history plus the ref tensor. We Two constraints we ran intoSweep cost scales with the pending-key set. Sender-side Reconstructing the keys avoids it entirely: keys are To be fair to the scanning form, reconstruction is not race-free either: a Reclaiming at sender-finish time is not safe for segments. What we shippedA minimal abort-only The gating matches review point 2: normally-finished requests must keep their state until I can contribute the stage_id=0 sender-abort test case as a PR against this branch if that |
|
@gagandhakrey this PR currently has merge conflicts with Please rebase onto the latest git fetch origin main
git rebase origin/main
# resolve conflicts, then:
git add -A
git rebase --continue
git push --force-with-leaseOnce it is mergeable again we can continue review. Thanks! |
|
The unread-SHM cleanup is still useful, but I found a correctness issue that should be addressed before merging.
An isolated reproduction using this PR’s unmodified cleanup methods and real POSIX shared memory confirms the deletion. Checking that all suffix fields are numeric would not resolve this ambiguity. Could we track exact request-to-key ownership instead, and add this case as a regression test? Also, current Finally, the latest update in #6352 identifies the remaining memory growth as separate issues. The description should distinguish those from the abort-path leak rather than imply this PR resolves all host-memory growth. |
|
Please resolve conflicts |
…on-abort Resolve conflicts with the sender-generation fencing from vllm-project#6670: - shm_connector: keep main's non-blocking flock and BlockingIOError retry path alongside the `consumed` lock-file cleanup. - chunk_transfer_adapter: dispatch the connector_cleanup task before main's generation-fenced send body, and take main's cleanup_sender docstring (it is now safe to call from the scheduler thread). - _run_connector_cleanup no longer calls cleanup_sender: finish_requests already retires the generation inline, and a second call from the save thread would cancel a successor generation that reused the external id. Add a regression test for that case. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
|
resolved @Gaohan123 _run_connector_cleanup no longer calls cleanup_sender(), since finish_requests now does it inline and a second call could cancel a request reusing the same external id. Test added |
Omni ReviewBot: no human activity for 7 days@gagandhakrey this pull request has had no human commit, comment or review since 2026-09-24. Please confirm the current plan and next step. The author or a maintainer decides whether to change the PR state. To keep it moving, any one of these is enough: push an update, reply to the open blocker, or post the current plan and timeline. |
Omni ReviewBot routing recordAssigned Strict on cursor (cursor-grok-4.6-high) under experiment |
Omni ReviewBot attempt recordReview attempt ended as timeout (Strict attempt outlived its budget; falling back to direct/cursor/auto). |
1 similar comment
Omni ReviewBot attempt recordReview attempt ended as timeout (Strict attempt outlived its budget; falling back to direct/cursor/auto). |
vllm-omni-review-bot
left a comment
There was a problem hiding this comment.
Omni ReviewBot review
Changes since the previous review
- 0 new inline finding(s); 1 finding(s) below.
CI at
f2c09b53affe(2026-10-09T13:06:43.964448+00:00): required check(s) blocking:buildkite/vllm-omni(missing).
Note: The assigned review arm
strict/cursor/cursor-grok-4.6-highcould not complete this review, so it was produced by the fallback armdirect/cursor/auto. It is excluded from the routing experiment.
Full review analysis
PR description
On abort or error, finish_requests still reclaims sender state inline, then enqueues a connector_cleanup task so the save thread unlinks that request's unconsumed shared-memory segments only after chunks already queued for it. SharedMemoryConnector.cleanup now matches structural key shapes instead of a bare prefix or suffix, and a consumed read deletes the lock file even when deserialization fails. Natural finish and FINISHED_STOPPED do not unlink segments. The PR text still says the sweep also calls cleanup_sender(); at this head that second call is gone so a reused external id is not retired twice.
Change flow
flowchart LR
Abort["[EXISTING] finish_requests abort or error"]:::existing
Gate["[CHANGED] FINISHED_ABORTED or ERROR gate"]:::changed
Task["[NEW] connector_cleanup save-queue task"]:::new
Match["[CHANGED] cleanup key match and unlink"]:::changed
Shm["[EXISTING] /dev/shm segment and lock file"]:::existing
Abort --> Gate --> Task --> Match --> Shm
classDef existing fill:#e5e7eb,stroke:#6b7280,color:#111827
classDef changed fill:#fef3c7,stroke:#d97706,color:#451a03,stroke-width:2px
classDef new fill:#dcfce7,stroke:#16a34a,color:#052e16,stroke-width:2px
classDef removed fill:#fee2e2,stroke:#dc2626,color:#450a0a,stroke-width:2px
Findings
- [P1] Abort cleanup still unlinks another request's chunk key —
vllm_omni/distributed/omni_connectors/connectors/shm_connector.py:175
Existing thread: #4349 (comment)
Evidence for Abort cleanup still unlinks another request's chunk key
_key_belongs_to treats a four-field remainder as a rank-aware key (len(fields) in (2, 4) and only the last field decimal). Chunk puts are {external_req_id}_{stage_id}_{chunk_id} (chunk_transfer_adapter.py around the connector_put_key format). For live ids abc and abc_1_2, key abc_1_2_0_0 starts with abc_, splits into 1,2,0,0, and cleanup("abc") unlinks that segment and its lock file. external_req_id is the client-facing id, so underscore ids are in contract. The new tests only cover abc vs abc_def_0_0 (three remainder fields) and cleanup("0"); they never put abc_1_2_0_0. That unlink drops a chunk the other request's consumer is still waiting on, which is worse than the leak. Record the exact keys at put and delete those. Requiring every suffix field to be numeric does not separate this key from the rank-aware shape.
🤖 This review was generated by InferMatrix Copilot, an open-source repo-maintenance agent for PR review, CI debugging and issue triage. Try it on your own repo, and ⭐ star it if it helped!
Omni ReviewBot: finding feedback[p1] Abort cleanup still unlinks another request's chunk key — See the review for details. If you are the PR author and disagree, react 👎 here; the maintainer will see your disagreement. |
Purpose
Fixes a
/dev/shmand host-memory leak on request abort in the SHM chunk-transfer path.What breaks
SharedMemoryConnector.cleanup()has no callers on current main. SHM segments are only unlinked insideshm_read_bytes()during successful consumption, so a chunk that is written but never read — client disconnect, downstream early exit, consumer failure — stays allocated until process exit, along with its/dev/shm/shm_{key}_lockfile.lock./dev/shmis physical RAM. Once exhausted, every subsequentput()fails and chunk transfer breaks pipeline-wide, not per-request.Aborted requests additionally retain sender-side state (
put_req_chunk,request_payload,code_prompt_token_ids,requests_num_chunks_sent,ramp_chunk_count,_pending_streaming_prefills), which can pin CPU tensors held by TTS input processors.cleanup_sender()is reachable only from the natural-finish path today.This is issue #6352, measured at ~45 KB/abort and 258 → 60 MiB/h under sustained load.
Repro
Fix
1. Abort-time connector sweep.
finish_requests()enqueues aconnector_cleanuptask onto the save queue;_send_single_request()dispatches it to_run_connector_cleanup(), which callsconnector.cleanup(external_req_id)and thencleanup_sender().Enqueued rather than called inline: both
put()and the sweep run on the save-loop thread, so FIFO ordering guarantees no chunk is written after its own sweep.2. Gated on
finished_status.Only
FINISHED_ABORTED/FINISHED_ERRORsweep. Any other terminal status may still have an unconsumed terminal chunk downstream, and unlinking it would leave the next stage waiting on a finish marker that no longer exists.The
is_finished()guard alone is no longer sufficient — #6360's resumable-FINISHED_STOPPEDcarve-out lets a finished request fall through it.3. Structural cleanup-key matching.
cleanup()previously matched any key withrequest_idas a bare_-delimited prefix or suffix. Both collide:abcsweptabc_def's live segments, andcleanup("0")swept chunk 0 of every request in flight._key_belongs_to()now matches only the real key shapes:{id}{id}_{stage}_{chunk}{id}_{stage}_{chunk}_{from_rank}_{to_rank}omni_{from}_to_{to}_kv_cache_{id}4. Lock-file leak on consume.
_get_data_with_lock()now removes the lock file whenever the segment was consumed, not only when deserialization succeeded.Once
shm_read_bytes()has unlinked the segment, no retry can succeed, so a deserialization failure must not strand the lock file.Scope
_make_key(key, from_stage, to_stage), so a bare-idcleanup()is a no-op there; their equivalent leak is out of scope._pending_keysgrowth on the successful path is not addressed here —put()adds to the producer's set whileget()prunes the consumer's set in a different process. Abort-path growth is bounded by this PR; the natural-finish case needs a separate change.Note for reviewers
Fix 3 required re-pointing an existing test.
TestCleanup::test_cleanup_removes_unconsumed_segmentasserted the generic_-suffix match (keycleanup_req_42, swept bycleanup("req_42")) — exactly the behaviour being removed.It now uses the real KV key shape, with three added regressions covering the collisions. This is a deliberate narrowing of
cleanup()'s contract, not an incidental edit.Test Plan
pytest -sv tests/distributed/omni_connectors/test_shm_connector.py \ tests/distributed/omni_connectors/test_chunk_transfer_adapter.py \ -m 'core_model and cpu'Adapter (
test_chunk_transfer_adapter.py) — 6 tests, 7 cases:FINISHED_STOPPEDdoes not enqueue cleanup (finished_statusgate).stage_id=0abort is swept, parametrized overFINISHED_ABORTED/FINISHED_ERROR.connector.cleanup().Connector (
test_shm_connector.py) — 5 tests:abcdoes not sweepabc_def.cleanup("0")does not sweep other requests' chunk 0.Additional Updates
Also updated the title with the
[Bugfix]prefix and refreshed the description. Four claims had gone stale against current main after the rebase, notably "producer_pending_keysremains bounded", which only holds on the abort path and is now explicitly scoped out.