feat(grpc): add vLLM KV cache event support (SubscribeKvEvents bridge) - #1652
Conversation
|
Note Reviews pausedIt looks like this branch is under active development. To avoid overwhelming you with review comments due to an influx of new commits, CodeRabbit has automatically paused this review. You can configure this behavior by changing the Use the following commands to manage reviews:
Use the checkboxes below for quick actions:
📝 WalkthroughWalkthroughBridges vLLM KV-cache ZMQ events into SMG via a new conversion module and an async gRPC server-streaming RPC. Adds endpoint-rank rewriting, deterministic int64 hash reduction, token slicing, config resolution, comprehensive unit/integration tests, pyzmq/msgspec extras, and docs describing vLLM gRPC launch with KV events enabled. ChangesvLLM KV-cache event streaming integration
Sequence Diagram(s)sequenceDiagram
participant vLLM as vLLM (ZMQ PUB)
participant Stream as stream_kv_events
participant Servicer as VllmEngineServicer
participant Client as gRPC Client
vLLM->>Stream: publish multipart (topic, seq, payload)
Stream->>Stream: decode payload -> convert to KvEventBatch
Stream->>Servicer: yield KvEventBatch
Servicer->>Client: stream KvEventBatch over SubscribeKvEvents
Estimated code review effort🎯 4 (Complex) | ⏱️ ~45 minutes Possibly related PRs
Suggested reviewers
Poem
🚥 Pre-merge checks | ✅ 4 | ❌ 1❌ Failed checks (1 warning)
✅ Passed checks (4 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches🧪 Generate unit tests (beta)
Comment |
There was a problem hiding this comment.
Code Review
This pull request adds support for event-driven routing with vLLM by bridging its ZMQ KV cache events to gRPC server-streaming. It introduces a new SubscribeKvEvents endpoint in the vLLM servicer, helper functions for event conversion and ZMQ streaming, corresponding unit and integration tests, documentation, and a cache-aware routing benchmark script. Feedback on the changes highlights two key improvements: replacing asyncio.wait_for with socket.poll on the ZMQ socket to prevent potential message loss or resource leaks under load, and defensively replacing 0.0.0.0 with 127.0.0.1 in endpoint resolution to ensure cross-platform reliability.
Important
The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 122cd08f6b
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
There was a problem hiding this comment.
Clean implementation that faithfully mirrors the SGLang servicer while properly accommodating vLLM's type differences (bytes/int block hashes, config access patterns). The engine-free helper module design enables thorough testing without GPU/vLLM. No issues found.
There was a problem hiding this comment.
Actionable comments posted: 7
🤖 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 `@docs/getting-started/kv-events-cache-aware.md`:
- Around line 98-105: The vLLM example shows the vllm serve command with
--kv-events-config but omits the SMG gRPC mode flags needed for KV-event
subscriptions; update the snippet or surrounding text to include the minimal
gRPC-mode option(s) so the worker actually runs in SMG gRPC mode (e.g., add the
appropriate gRPC-mode flag(s) to the vllm serve invocation or add an inline note
referencing the gRPC Workers guide), and ensure you reference the exact command
and option names shown (vllm serve and --kv-events-config) so readers know to
run the worker in SMG gRPC mode for KV events to work.
In `@grpc_servicer/pyproject.toml`:
- Line 33: Update the vllm dependency entry in pyproject.toml that currently
lists vllm = ["vllm>=0.19.0", "pyzmq>=25.0.0", "msgspec>=0.18.0"] by validating
and, if acceptable, raising the minimum floors for pyzmq and msgspec (e.g., to
pyzmq>=27.1.0 and msgspec>=0.21.1) or to the lowest versions you require; ensure
you verify compatibility with vllm and the rest of the dependency graph, run the
project's test suite and dependency resolution (poetry/pip) to confirm no
conflicts, and update the constraint literals for pyzmq and msgspec in the same
dependency list once confirmed.
In `@grpc_servicer/smg_grpc_servicer/vllm/servicer.py`:
- Around line 949-961: The RPC currently only handles asyncio.CancelledError, so
unexpected exceptions from stream_kv_events will escape and surface as UNKNOWN;
update the try/except to also catch generic exceptions (except
asyncio.CancelledError) around the async for in servicer method that calls
stream_kv_events, log the error with logger (include the exception), and call
context.abort(grpc.StatusCode.INTERNAL, str(e)) to follow the repo convention;
keep the existing finally block that closes sub_socket (linger=0) and the
CancelledError pass behavior.
- Around line 909-956: SubscribeKvEvents currently ignores
request.start_sequence_number causing false resume semantics; update
SubscribeKvEvents to read and validate request.start_sequence_number and enforce
it before yielding events: if start_sequence_number is provided (>=0) either
pass it into stream_kv_events (if stream_kv_events supports a start arg) or
implement a short-circuit filter around the async for loop that discards
received KVEventBatch items until their sequence >=
request.start_sequence_number, and return grpc.INVALID_ARGUMENT for negative
values; reference SubscribeKvEvents, request.start_sequence_number,
stream_kv_events, and the yielded proto_batch to locate where to apply the
logic.
In `@grpc_servicer/tests/test_vllm_kv_events_stream.py`:
- Around line 46-77: The test opens ZMQ sockets (pub/sub) and currently only
closes them on the success path; wrap the async send/consume/await logic inside
a try/finally so pub.close(linger=0) and sub.close(linger=0) always run on
success or failure; update both tests that use ctx.socket, pub, sub (the consume
coroutine / fake_decode usage and the similar block at lines 84-117) to ensure
sockets are closed in the finally block (or use a context manager/helper that
guarantees close) so resources are released even when asyncio.wait_for or
assertions raise.
In `@scripts/bench_kv_events_cache_hit.py`:
- Around line 83-96: The pct helper can raise IndexError when called with an
empty list (used for ttfts and totals in the output JSON); update the pct(xs, p)
function to defensively handle an empty xs (e.g., if not xs: return 0.0) and
keep the existing behavior for len(xs)==1 and len(xs)>1 so
ttft_p50_ms/ttft_p99_ms/total_p50_ms don't crash when all requests produced no
timings.
- Around line 55-61: The SSE loop currently treats each chunk from resp.content
as a full message which can split or combine SSE events; fix by implementing a
buffer-based parser: keep a bytes buffer, append each chunk from resp.content,
then split buffer on b'\n\n' to extract complete events, leaving the remainder
in buffer for the next chunk; for each complete event, parse lines starting with
b"data:" to get the payload, set ttft (using ttft and start) on the first
non-empty data payload, and detect the terminal payload b"[DONE]" to break;
update the loop that currently uses resp.content, ttft, start, and line to use
this buffered-splitting approach so message boundaries are handled correctly.
🪄 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: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: 37efc491-9cfc-4598-b20d-56e7fbc3e902
📒 Files selected for processing (7)
docs/getting-started/kv-events-cache-aware.mdgrpc_servicer/pyproject.tomlgrpc_servicer/smg_grpc_servicer/vllm/kv_events.pygrpc_servicer/smg_grpc_servicer/vllm/servicer.pygrpc_servicer/tests/test_vllm_kv_events.pygrpc_servicer/tests/test_vllm_kv_events_stream.pyscripts/bench_kv_events_cache_hit.py
Signed-off-by: key4ng <rukeyang@gmail.com>
Signed-off-by: key4ng <rukeyang@gmail.com>
Signed-off-by: key4ng <rukeyang@gmail.com>
Signed-off-by: key4ng <rukeyang@gmail.com>
Signed-off-by: key4ng <rukeyang@gmail.com>
Signed-off-by: key4ng <rukeyang@gmail.com>
Signed-off-by: key4ng <rukeyang@gmail.com>
Signed-off-by: key4ng <rukeyang@gmail.com>
122cd08 to
7fc376a
Compare
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 7fc376aa88
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
There was a problem hiding this comment.
Actionable comments posted: 1
♻️ Duplicate comments (4)
grpc_servicer/tests/test_vllm_kv_events_stream.py (2)
46-79:⚠️ Potential issue | 🟡 Minor | ⚡ Quick winGuarantee socket cleanup on failure paths.
pub.close()/sub.close()run only on success. Ifwait_foror assertions fail, sockets remain open and can interfere with subsequent tests.Suggested fix pattern
async def test_stream_yields_batches_in_order(): ctx = zmq.asyncio.Context.instance() pub = ctx.socket(zmq.PUB) - port = pub.bind_to_random_port("tcp://127.0.0.1") sub = ctx.socket(zmq.SUB) - sub.subscribe(b"") - sub.connect(f"tcp://127.0.0.1:{port}") - await asyncio.sleep(0.2) # allow SUB connection to establish before publishing - - # decode ignores payload bytes and returns a prebuilt fake batch keyed by seq. - def fake_decode(payload: bytes): - return KVEventBatch(int.from_bytes(payload, "big")) - - collected = [] - - async def consume(): - async for batch in kv_events.stream_kv_events( - sub, fake_decode, send_initial_metadata=_noop, is_cancelled=lambda: len(collected) >= 2 - ): - collected.append(batch) - - async def _noop(): - return None - - consumer = asyncio.create_task(consume()) - for seq in (1, 2): - await pub.send_multipart([b"kv", _seq_bytes(seq), seq.to_bytes(8, "big")]) - await asyncio.sleep(0.05) - await asyncio.wait_for(consumer, timeout=5) - - pub.close(linger=0) - sub.close(linger=0) - - assert [b.sequence_number for b in collected] == [1, 2] - assert collected[0].events[0].stored.blocks[0].block_hash == 1 + try: + port = pub.bind_to_random_port("tcp://127.0.0.1") + sub.subscribe(b"") + sub.connect(f"tcp://127.0.0.1:{port}") + await asyncio.sleep(0.2) + + def fake_decode(payload: bytes): + return KVEventBatch(int.from_bytes(payload, "big")) + + collected = [] + + async def consume(): + async for batch in kv_events.stream_kv_events( + sub, fake_decode, send_initial_metadata=_noop, is_cancelled=lambda: len(collected) >= 2 + ): + collected.append(batch) + + async def _noop(): + return None + + consumer = asyncio.create_task(consume()) + for seq in (1, 2): + await pub.send_multipart([b"kv", _seq_bytes(seq), seq.to_bytes(8, "big")]) + await asyncio.sleep(0.05) + await asyncio.wait_for(consumer, timeout=5) + + assert [b.sequence_number for b in collected] == [1, 2] + assert collected[0].events[0].stored.blocks[0].block_hash == 1 + finally: + pub.close(linger=0) + sub.close(linger=0)🤖 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 `@grpc_servicer/tests/test_vllm_kv_events_stream.py` around lines 46 - 79, Wrap the test body that creates pub, sub and the consumer task in a try/finally so sockets and the consumer are always cleaned up on failures: after creating pub, sub and starting consumer (asyncio.create_task(consume())), move the publish/await logic into the try block and in finally ensure pub.close(linger=0), sub.close(linger=0) and cancel/await the consumer task (consumer.cancel() and await it handling asyncio.CancelledError) so kv_events.stream_kv_events/consume/fake_decode resources are guaranteed to be released even if wait_for or assertions fail.
84-119:⚠️ Potential issue | 🟡 Minor | ⚡ Quick winGuarantee socket cleanup on failure paths.
Same issue as the previous test:
pub.close()/sub.close()run only on success. Wrap in try/finally to ensure cleanup.Suggested fix pattern
async def test_stream_skips_short_frames_and_decode_errors(): ctx = zmq.asyncio.Context.instance() pub = ctx.socket(zmq.PUB) - port = pub.bind_to_random_port("tcp://127.0.0.1") sub = ctx.socket(zmq.SUB) - sub.subscribe(b"") - sub.connect(f"tcp://127.0.0.1:{port}") - await asyncio.sleep(0.2) - - def fake_decode(payload: bytes): - if payload == b"bad": - raise ValueError("boom") - return KVEventBatch(7) - - collected = [] - - async def _noop(): - return None - - async def consume(): - async for batch in kv_events.stream_kv_events( - sub, fake_decode, send_initial_metadata=_noop, is_cancelled=lambda: len(collected) >= 1 - ): - collected.append(batch) - - consumer = asyncio.create_task(consume()) - await pub.send_multipart([b"kv", _seq_bytes(1)]) # short frame -> skipped - await asyncio.sleep(0.05) - await pub.send_multipart([b"kv", _seq_bytes(2), b"bad"]) # decode error -> skipped - await asyncio.sleep(0.05) - await pub.send_multipart([b"kv", _seq_bytes(3), b"ok"]) # good -> yielded - await asyncio.wait_for(consumer, timeout=5) - - pub.close(linger=0) - sub.close(linger=0) - - assert [b.sequence_number for b in collected] == [3] + try: + port = pub.bind_to_random_port("tcp://127.0.0.1") + sub.subscribe(b"") + sub.connect(f"tcp://127.0.0.1:{port}") + await asyncio.sleep(0.2) + + def fake_decode(payload: bytes): + if payload == b"bad": + raise ValueError("boom") + return KVEventBatch(7) + + collected = [] + + async def _noop(): + return None + + async def consume(): + async for batch in kv_events.stream_kv_events( + sub, fake_decode, send_initial_metadata=_noop, is_cancelled=lambda: len(collected) >= 1 + ): + collected.append(batch) + + consumer = asyncio.create_task(consume()) + await pub.send_multipart([b"kv", _seq_bytes(1)]) # short frame -> skipped + await asyncio.sleep(0.05) + await pub.send_multipart([b"kv", _seq_bytes(2), b"bad"]) # decode error -> skipped + await asyncio.sleep(0.05) + await pub.send_multipart([b"kv", _seq_bytes(3), b"ok"]) # good -> yielded + await asyncio.wait_for(consumer, timeout=5) + + assert [b.sequence_number for b in collected] == [3] + finally: + pub.close(linger=0) + sub.close(linger=0)🤖 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 `@grpc_servicer/tests/test_vllm_kv_events_stream.py` around lines 84 - 119, The test currently only closes ZMQ sockets on the success path; wrap the send/receive and consumer awaiting logic in a try/finally so pub.close(linger=0) and sub.close(linger=0) always run (and cancel the consumer task if still running). Concretely, enclose the code that creates consumer = asyncio.create_task(consume()), the send_multipart calls and await asyncio.wait_for(consumer, ...) in a try block and move the pub.close and sub.close calls into the finally block (optionally checking consumer.cancel() and awaiting it) to guarantee socket cleanup even on failures in kv_events.stream_kv_events, fake_decode, or the test awaits.grpc_servicer/smg_grpc_servicer/vllm/servicer.py (2)
1046-1058:⚠️ Potential issue | 🟠 Major | ⚡ Quick winMap unexpected stream failures to INTERNAL with server-side logging.
Only
asyncio.CancelledErroris handled here. Other runtime errors will escape this RPC and typically surface asUNKNOWNinstead of the repo'sINTERNAL + str(e)convention.Suggested fix
try: async for proto_batch in stream_kv_events( sub_socket, decoder.decode, lambda: context.send_initial_metadata(()), context.cancelled, ): yield proto_batch except asyncio.CancelledError: pass + except Exception as e: + logger.exception("SubscribeKvEvents failed") + await context.abort(grpc.StatusCode.INTERNAL, str(e)) finally: sub_socket.close(linger=0) logger.info("SubscribeKvEvents: stream closed")🤖 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 `@grpc_servicer/smg_grpc_servicer/vllm/servicer.py` around lines 1046 - 1058, The current SubscribeKvEvents stream loop only catches asyncio.CancelledError and lets other exceptions surface as UNKNOWN; update the try/except to catch Exception around the async for (in the function containing the stream_kv_events call) and on any non-cancellation error call logger.exception(...) and convert the error into a gRPC INTERNAL status before closing the stream (use context.abort(grpc.StatusCode.INTERNAL, str(e)) or the repo's INTERNAL + str(e) pattern), ensuring sub_socket.close(linger=0) still runs in finally; reference the async generator call to stream_kv_events, the context.cancelled callback, and sub_socket.close so you update the correct block.Source: Learnings
1006-1023:⚠️ Potential issue | 🟠 Major | ⚡ Quick winHandle
start_sequence_numberexplicitly to avoid API-contract drift.
SubscribeKvEventscurrently ignoresrequest.start_sequence_number, even though the request contract exposes resume behavior. That can mislead clients into assuming resume semantics are honored.Suggested fix
async def SubscribeKvEvents( self, request: common_pb2.SubscribeKvEventsRequest, context: grpc.aio.ServicerContext, ) -> AsyncIterator[common_pb2.KvEventBatch]: + if request.start_sequence_number != 0: + await context.abort( + grpc.StatusCode.UNIMPLEMENTED, + "start_sequence_number resume is not supported yet for vLLM KV events", + ) + if self._kv_events_config is None: await context.abort( grpc.StatusCode.UNIMPLEMENTED, "KV cache events not enabled. Start vLLM with " "--kv-events-config " '\'{"enable_kv_cache_events": true, "publisher": "zmq"}\'', )🤖 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 `@grpc_servicer/smg_grpc_servicer/vllm/servicer.py` around lines 1006 - 1023, SubscribeKvEvents currently ignores request.start_sequence_number causing clients to be unable to resume; update SubscribeKvEvents to honor request.start_sequence_number by validating it (e.g., non-negative) and using it to initialize or drive the ZMQ subscription/streaming logic in this function: if the ZMQ publisher API supports starting from a sequence, pass start_sequence_number into that call, otherwise filter/skipping emitted KvEventBatch items until the publisher sequence >= request.start_sequence_number before yielding to the gRPC stream; ensure invalid values return an appropriate gRPC error via context.abort and document/respect resume semantics in SubscribeKvEvents.
🤖 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 `@grpc_servicer/smg_grpc_servicer/vllm/servicer.py`:
- Line 1023: Remove the redundant `return` statement that immediately follows a
`context.abort()` call in the vllm servicer method (the occurrence where
`context.abort()` is used with gRPC to terminate the request); `context.abort()`
raises `grpc.aio.AbortError` so the `return` is unnecessary and should be
deleted to match the other `context.abort()` usages in this file (e.g., those
without trailing `return`).
---
Duplicate comments:
In `@grpc_servicer/smg_grpc_servicer/vllm/servicer.py`:
- Around line 1046-1058: The current SubscribeKvEvents stream loop only catches
asyncio.CancelledError and lets other exceptions surface as UNKNOWN; update the
try/except to catch Exception around the async for (in the function containing
the stream_kv_events call) and on any non-cancellation error call
logger.exception(...) and convert the error into a gRPC INTERNAL status before
closing the stream (use context.abort(grpc.StatusCode.INTERNAL, str(e)) or the
repo's INTERNAL + str(e) pattern), ensuring sub_socket.close(linger=0) still
runs in finally; reference the async generator call to stream_kv_events, the
context.cancelled callback, and sub_socket.close so you update the correct
block.
- Around line 1006-1023: SubscribeKvEvents currently ignores
request.start_sequence_number causing clients to be unable to resume; update
SubscribeKvEvents to honor request.start_sequence_number by validating it (e.g.,
non-negative) and using it to initialize or drive the ZMQ subscription/streaming
logic in this function: if the ZMQ publisher API supports starting from a
sequence, pass start_sequence_number into that call, otherwise filter/skipping
emitted KvEventBatch items until the publisher sequence >=
request.start_sequence_number before yielding to the gRPC stream; ensure invalid
values return an appropriate gRPC error via context.abort and document/respect
resume semantics in SubscribeKvEvents.
In `@grpc_servicer/tests/test_vllm_kv_events_stream.py`:
- Around line 46-79: Wrap the test body that creates pub, sub and the consumer
task in a try/finally so sockets and the consumer are always cleaned up on
failures: after creating pub, sub and starting consumer
(asyncio.create_task(consume())), move the publish/await logic into the try
block and in finally ensure pub.close(linger=0), sub.close(linger=0) and
cancel/await the consumer task (consumer.cancel() and await it handling
asyncio.CancelledError) so kv_events.stream_kv_events/consume/fake_decode
resources are guaranteed to be released even if wait_for or assertions fail.
- Around line 84-119: The test currently only closes ZMQ sockets on the success
path; wrap the send/receive and consumer awaiting logic in a try/finally so
pub.close(linger=0) and sub.close(linger=0) always run (and cancel the consumer
task if still running). Concretely, enclose the code that creates consumer =
asyncio.create_task(consume()), the send_multipart calls and await
asyncio.wait_for(consumer, ...) in a try block and move the pub.close and
sub.close calls into the finally block (optionally checking consumer.cancel()
and awaiting it) to guarantee socket cleanup even on failures in
kv_events.stream_kv_events, fake_decode, or the test awaits.
🪄 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: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: 68d76101-0785-444a-aaf3-042104057b9c
📒 Files selected for processing (7)
docs/getting-started/kv-events-cache-aware.mdgrpc_servicer/pyproject.tomlgrpc_servicer/smg_grpc_servicer/vllm/kv_events.pygrpc_servicer/smg_grpc_servicer/vllm/servicer.pygrpc_servicer/tests/test_vllm_kv_events.pygrpc_servicer/tests/test_vllm_kv_events_stream.pyscripts/bench_kv_events_cache_hit.py
- stream_kv_events: poll() before recv_multipart() instead of asyncio.wait_for (cancelling a zmq.asyncio recv future can drop a dequeued message) - endpoint_for_rank: rewrite 0.0.0.0 -> 127.0.0.1 (not connectable on macOS/Windows) - SubscribeKvEvents: map unexpected errors to INTERNAL + logger.exception (repo convention); drop redundant return after context.abort() - tests: close ZMQ sockets in try/finally; cover 0.0.0.0 endpoint - bench: guard pct() against empty samples Signed-off-by: key4ng <rukeyang@gmail.com>
Plain 'vllm serve' starts the HTTP server and cannot stream KV events; show 'python -m vllm.entrypoints.grpc_server' (the SMG gRPC worker entrypoint). Signed-off-by: key4ng <rukeyang@gmail.com>
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: f8be95c519
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
There was a problem hiding this comment.
♻️ Duplicate comments (1)
grpc_servicer/smg_grpc_servicer/vllm/servicer.py (1)
1006-1023:⚠️ Potential issue | 🟠 Major | ⚡ Quick winHonor or reject
start_sequence_numberbefore opening the live stream.
SubscribeKvEventsstill ignoresrequest.start_sequence_number, so a client asking to resume atNsilently gets the current live stream instead. That breaks the subscribe contract and can corrupt downstream gap-detection/reconnect state because the yieldedsequence_numbervalues no longer match the caller’s requested starting point.Suggested minimal fix
async def SubscribeKvEvents( self, request: common_pb2.SubscribeKvEventsRequest, context: grpc.aio.ServicerContext, ) -> AsyncIterator[common_pb2.KvEventBatch]: """Bridge vLLM's in-process ZMQ KV cache events to gRPC server-streaming. Mirrors sglang/servicer.py:SubscribeKvEvents. Uses the ZMQ publisher's native sequence numbers directly as gRPC sequence numbers. """ + if request.start_sequence_number < 0: + await context.abort( + grpc.StatusCode.INVALID_ARGUMENT, + "start_sequence_number must be >= 0", + ) + if request.start_sequence_number != 0: + await context.abort( + grpc.StatusCode.UNIMPLEMENTED, + "start_sequence_number resume is not supported yet for vLLM KV events", + ) + if self._kv_events_config is None: await context.abort( grpc.StatusCode.UNIMPLEMENTED, "KV cache events not enabled. Start vLLM with " "--kv-events-config "🤖 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 `@grpc_servicer/smg_grpc_servicer/vllm/servicer.py` around lines 1006 - 1023, SubscribeKvEvents currently ignores request.start_sequence_number and always opens the live ZMQ stream, breaking resume semantics; update SubscribeKvEvents to validate and honor or reject request.start_sequence_number before opening the live stream: check request.start_sequence_number against the publisher's current/earliest available sequence (using the same ZMQ publisher client or sequence tracking used elsewhere in the class), and if the requested start is greater than the latest available sequence or otherwise invalid, call await context.abort with an appropriate grpc.StatusCode (e.g., OUT_OF_RANGE or UNIMPLEMENTED) and message; if the start_sequence_number is valid and corresponds to cached/historical events, send those events first (or use the publisher’s replay API) and only then attach to the live ZMQ publisher for subsequent events so yielded sequence_number values align with request.start_sequence_number.
🤖 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.
Duplicate comments:
In `@grpc_servicer/smg_grpc_servicer/vllm/servicer.py`:
- Around line 1006-1023: SubscribeKvEvents currently ignores
request.start_sequence_number and always opens the live ZMQ stream, breaking
resume semantics; update SubscribeKvEvents to validate and honor or reject
request.start_sequence_number before opening the live stream: check
request.start_sequence_number against the publisher's current/earliest available
sequence (using the same ZMQ publisher client or sequence tracking used
elsewhere in the class), and if the requested start is greater than the latest
available sequence or otherwise invalid, call await context.abort with an
appropriate grpc.StatusCode (e.g., OUT_OF_RANGE or UNIMPLEMENTED) and message;
if the start_sequence_number is valid and corresponds to cached/historical
events, send those events first (or use the publisher’s replay API) and only
then attach to the live ZMQ publisher for subsequent events so yielded
sequence_number values align with request.start_sequence_number.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: 6dbca646-92f2-4245-b100-d8d2b8ee15b6
📒 Files selected for processing (6)
docs/getting-started/kv-events-cache-aware.mdgrpc_servicer/smg_grpc_servicer/vllm/kv_events.pygrpc_servicer/smg_grpc_servicer/vllm/servicer.pygrpc_servicer/tests/test_vllm_kv_events.pygrpc_servicer/tests/test_vllm_kv_events_stream.pyscripts/bench_kv_events_cache_hit.py
E2E + A/B/C benchmark on real GPUs (2× H100, vLLM 0.22.1)Ran the two remaining Test Plan items (manual E2E + A/B/C benchmark) against the latest commit ( Setup: 2× H100 80GB · vLLM 0.22.1 launched in SMG gRPC mode ( ✅ Functional / E2E
Benchmark (n=6/condition, concurrency 48, rotated condition order, burn-in discarded, exact permutation tests)Two pressure regimes (KV cache ≈160k tok/worker; prefix working set sized to force eviction): Regime A —
Regime B —
Findings
Net: the bridge is correct and brings vLLM to KV-event parity with SGLang; cache-aware routing gives a clear p99-tail benefit over round-robin. The event-driven path's advantage over the approximate tree isn't observable in this single-gateway setup — it should surface where the approximation diverges from the worker's real cache state (multi-gateway / mesh, cross-instance routing, restarts), which this bench doesn't exercise. Methodology notes
|
Follow-up: multi-gateway is where event-driven routing wins (~7× throughput)The previous comment found event-driven ≈ approximate tree in a single gateway — expected, since a lone gateway's approximate tree is fed by the very router making the decisions, so its cache assumption is already accurate. The regime where event-driven should matter is multiple gateway replicas sharing a worker pool (the normal production topology), where each gateway's approximate tree is blind to what the other replicas routed. Tested that here. Setup: Result — warm steady-state, balanced condition order, n=4/condition:
Complete separation — every event-driven run beat every approx run (exact permutation p=0.029, the floor for n=4 vs 4). Verified each run: events-on → each gateway logs Why: with events off and no mesh, the two gateways independently (and blindly) place the same prefix on different workers → the prefix gets duplicated across both → effective working set doubles → cache thrashes (44% hit) → the 32B prefill is recomputed on every miss → stuck at ~17 rps. With events on, both gateways see the workers' real KV state and converge each prefix onto one worker → cache fits (84% hit) → ~129 rps. Summary across regimes:
So the feature's value shows up exactly in the multi-replica deployment it's designed for.
Method notes
|
Update: the ~7.3× multi-gateway figure reflects an extreme stress setup, not a normal workloadFollowing review feedback, I re-ran the multi-gateway benchmark from my earlier comment under more realistic conditions. The ~7.3× came from a deliberately stressed configuration — a small/capped KV cache with several gateway replicas crowded onto it, plus a benchmark client that drove the gateways in a phased rather than concurrent pattern. That's not representative of a normal workload. Under more typical conditions the advantage is real but much smaller. Corrected numbers below. Under a more realistic load, the gap is much smaller
The real, defensible result — scaling gateway replicas4 workers, concurrent client, sweeping replica count:
The takeaway: event-driven throughput stays flat as you add gateway replicas (every replica routes from the workers' real KV state), while the approximate tree degrades as replicas multiply (each only knows its own routing). The size of the gap depends on cache headroom and replica count — it's large only when many replicas share a cache too small to absorb the overlap. How to frame the value
Test design & how the numbers were reachedWhat the test isolates. Cache-aware routing has to answer one question per request: which worker already holds this prefix? Two ways to answer it — the approximate tree infers it from each gateway's own routing history ("I sent P to W1, so W1 has it"); event-driven reads the worker's real Why multiple gateways. A single gateway is the sole author of its workers' cache, so its approximation is essentially always right → no difference (confirmed in the single-gateway comment: event ≈ approx). With ≥2 gateways sharing the same workers, gateway B is blind to what gateway A routed — unless it reads the workers' events. So the test runs N independent gateways (no Setup. DeepSeek-R1-Distill-Qwen-32B (dense, full attention), TP=2 workers, CUDA graphs (no How it went.
Method notes
|
…cript - Trim verbose docstrings/comments in kv_events.py and servicer.py and remove cross-references to the SGLang servicer. - Reword the vLLM docs section to stand on its own. - Remove scripts/bench_kv_events_cache_hit.py. Signed-off-by: key4ng <rukeyang@gmail.com>
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 847fe902c0
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
| zmq_ctx = zmq.asyncio.Context.instance() | ||
| sub_socket = zmq_ctx.socket(zmq.SUB) | ||
| sub_socket.subscribe(config.topic.encode("utf-8")) | ||
| sub_socket.connect(pub_endpoint) |
There was a problem hiding this comment.
Bind the SUB when vLLM is configured to connect
When config.endpoint is a concrete TCP address such as tcp://127.0.0.1:5557 or tcp://0.0.0.0:5557, vLLM's ZmqEventPublisher._socket_setup connects its PUB socket instead of binding it; this code also calls connect(), so no side binds and the SubscribeKvEvents stream remains idle for that supported endpoint form. Mirror vLLM's bind/connect heuristic for the subscriber side, or reject/normalize non-binding endpoints before advertising KV event support.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Actionable comments posted: 2
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (4)
docs/getting-started/kv-events-cache-aware.md (1)
237-242: 🧹 Nitpick | 🔵 Trivial | ⚡ Quick winUpdate reference section to include vLLM servicer bridge.
The reference section (line 241) cites only the SGLang servicer. The PR modifies
grpc_servicer/smg_grpc_servicer/vllm/servicer.py, so add a parallel reference line for vLLM.📝 Suggested addition
- Servicer bridge: `grpc_servicer/smg_grpc_servicer/sglang/servicer.py` (`SubscribeKvEvents`) +- Servicer bridge (vLLM): `grpc_servicer/smg_grpc_servicer/vllm/servicer.py` (`SubscribeKvEvents`) - SGLang upstream config: `python/sglang/srt/disaggregation/kv_events.py` (class `KVEventsConfig`)🤖 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 `@docs/getting-started/kv-events-cache-aware.md` around lines 237 - 242, Update the reference list to include the new vLLM servicer bridge entry added in the PR: add a parallel line referencing grpc_servicer/smg_grpc_servicer/vllm/servicer.py and the relevant servicer method (e.g., SubscribeKvEvents) alongside the existing SGLang servicer entry so the docs cite both servicer implementations; ensure the style matches the other lines (backticked path and method/class name).grpc_servicer/smg_grpc_servicer/vllm/kv_events.py (1)
167-170:⚠️ Potential issue | 🟡 Minor | ⚡ Quick winReject malformed sequence frames before converting them.
Line 170 currently parses any
frames[1]length into asequence_number. A truncated or extra-long frame turns into a bogus batch sequence and can trigger false gap handling downstream. Validate the sequence frame is exactly 8 bytes before callingint.from_bytes().Suggested fix
# ZMQ multipart: [topic, 8-byte big-endian seq, msgpack payload]. if len(frames) < 3: continue + if len(frames[1]) != 8: + logger.warning( + "Skipping vLLM KV event frame with invalid sequence size: %d", + len(frames[1]), + ) + continue zmq_seq = int.from_bytes(frames[1], "big")🤖 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 `@grpc_servicer/smg_grpc_servicer/vllm/kv_events.py` around lines 167 - 170, The code reads the ZMQ sequence frame blindly into zmq_seq using int.from_bytes(frames[1], "big"); update the validation to first ensure frames[1] is exactly 8 bytes long and reject the message (continue) if not, so malformed or truncated/oversized sequence frames can't produce bogus sequence numbers; modify the logic around the frames length check in kv_events.py (the block handling the ZMQ multipart frames and the zmq_seq assignment) to validate len(frames[1]) == 8 before calling int.from_bytes() and log or skip the frame when it fails validation.grpc_servicer/smg_grpc_servicer/vllm/servicer.py (2)
1031-1035:⚠️ Potential issue | 🟠 Major | ⚡ Quick winDon’t expose
SubscribeKvEventsfor multi-DP workers with only rank 0 wired.These lines knowingly subscribe to rank 0 only, but
KvEventBatch.dp_rankexists so clients can ingest per-rank state. On a multi-DP worker this publishes a partial cache view while still advertising the RPC as supported, so the gateway will index incomplete KV state. Please gate this RPC off for multi-DP configs until merged/per-rank streaming is implemented.🤖 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 `@grpc_servicer/smg_grpc_servicer/vllm/servicer.py` around lines 1031 - 1035, The SubscribeKvEvents RPC currently forces pub_endpoint = endpoint_for_rank(config.endpoint, 0) which exposes only rank-0 KV events; gate this RPC for multi-data-parallel workers by detecting a multi-DP config (e.g., check a data-parallel size flag on config such as config.dp_size or config.data_parallel_size) and return an unimplemented/disabled response instead of wiring rank-0 only streaming when more than one DP rank is present. Concretely, in the SubscribeKvEvents handler (the code that computes pub_endpoint and subscribes) add a guard: if the config indicates dp_size > 1 then short-circuit with an appropriate gRPC unimplemented/unsupported error (or log and close the stream) so clients cannot subscribe to a partial cache view until per-rank streaming (KvEventBatch.dp_rank / per-rank endpoints) is implemented.
1024-1045:⚠️ Potential issue | 🟠 MajorMove ZMQ socket/decoder setup into the guarded
try/finallyinSubscribeKvEvents.
SubscribeKvEventsconstructszmq_ctx,sub_socket, and themsgspecdecoder before thetry:that logs/aborts withgrpc.StatusCode.INTERNALand before thefinally:that closessub_socket. If any of that setup fails, the RPC can escape withoutcontext.abort(...)and without socket cleanup.Suggested fix
- zmq_ctx = zmq.asyncio.Context.instance() - sub_socket = zmq_ctx.socket(zmq.SUB) - sub_socket.subscribe(config.topic.encode("utf-8")) - sub_socket.connect(pub_endpoint) - logger.info("SubscribeKvEvents: connected to ZMQ endpoint %s", pub_endpoint) - - decoder = msgspec.msgpack.Decoder(KVEventBatch) - + sub_socket = None try: + zmq_ctx = zmq.asyncio.Context.instance() + sub_socket = zmq_ctx.socket(zmq.SUB) + sub_socket.subscribe(config.topic.encode("utf-8")) + sub_socket.connect(pub_endpoint) + logger.info("SubscribeKvEvents: connected to ZMQ endpoint %s", pub_endpoint) + + decoder = msgspec.msgpack.Decoder(KVEventBatch) async for proto_batch in stream_kv_events( sub_socket, decoder.decode, lambda: context.send_initial_metadata(()), context.cancelled, @@ except Exception as e: logger.exception("SubscribeKvEvents failed") await context.abort(grpc.StatusCode.INTERNAL, str(e)) finally: - sub_socket.close(linger=0) + if sub_socket is not None: + sub_socket.close(linger=0) logger.info("SubscribeKvEvents: stream closed")🤖 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 `@grpc_servicer/smg_grpc_servicer/vllm/servicer.py` around lines 1024 - 1045, Move the ZMQ and decoder setup into the guarded try/finally inside SubscribeKvEvents so failures are handled with context.abort and sockets are always cleaned up: create zmq_ctx, sub_socket (and call sub_socket.subscribe/connect) and the msgspec.msgpack.Decoder(KVEventBatch) after entering the try block and ensure the existing finally still closes sub_socket; if any of those setups raise, call context.abort(grpc.StatusCode.INTERNAL, ...) as the other error handling in SubscribeKvEvents does.Source: Learnings
🤖 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 `@docs/getting-started/kv-events-cache-aware.md`:
- Around line 94-107: The docs currently omit vLLM's KV cache block-size flag
and whether SMG will learn it from events; update the vLLM worker example or
explanatory text to either (a) explicitly add the engine flag `--block-size <N>`
to the `python -m vllm.entrypoints.grpc_server` example so the PagedAttention KV
block granularity is fixed, or (b) state that SMG will infer block size from KV
events like `BlockStored` and explain that those events include block-size
metadata—mention `--block-size` and `BlockStored` by name so readers know where
to configure or expect the size.
- Around line 102-106: Update the "Reference" section to include the vLLM bridge
entry alongside the existing SGLang servicer and remove the separate block-size
guidance: mention that the gRPC entrypoint usage is "python -m
vllm.entrypoints.grpc_server" with the --kv-events-config option (showing keys
like enable_kv_cache_events, publisher, endpoint, topic) and state that
KvEventMonitor will overwrite the initial --block-size from the first KV event
so no separate page-size/block-size CLI flag is needed in the guide; add a
link/reference to the vLLM servicer modules (vllm.servicer and vllm.kv_events)
next to the existing SGLang servicer reference (sglang.servicer) so readers can
find the vLLM implementation.
---
Outside diff comments:
In `@docs/getting-started/kv-events-cache-aware.md`:
- Around line 237-242: Update the reference list to include the new vLLM
servicer bridge entry added in the PR: add a parallel line referencing
grpc_servicer/smg_grpc_servicer/vllm/servicer.py and the relevant servicer
method (e.g., SubscribeKvEvents) alongside the existing SGLang servicer entry so
the docs cite both servicer implementations; ensure the style matches the other
lines (backticked path and method/class name).
In `@grpc_servicer/smg_grpc_servicer/vllm/kv_events.py`:
- Around line 167-170: The code reads the ZMQ sequence frame blindly into
zmq_seq using int.from_bytes(frames[1], "big"); update the validation to first
ensure frames[1] is exactly 8 bytes long and reject the message (continue) if
not, so malformed or truncated/oversized sequence frames can't produce bogus
sequence numbers; modify the logic around the frames length check in
kv_events.py (the block handling the ZMQ multipart frames and the zmq_seq
assignment) to validate len(frames[1]) == 8 before calling int.from_bytes() and
log or skip the frame when it fails validation.
In `@grpc_servicer/smg_grpc_servicer/vllm/servicer.py`:
- Around line 1031-1035: The SubscribeKvEvents RPC currently forces pub_endpoint
= endpoint_for_rank(config.endpoint, 0) which exposes only rank-0 KV events;
gate this RPC for multi-data-parallel workers by detecting a multi-DP config
(e.g., check a data-parallel size flag on config such as config.dp_size or
config.data_parallel_size) and return an unimplemented/disabled response instead
of wiring rank-0 only streaming when more than one DP rank is present.
Concretely, in the SubscribeKvEvents handler (the code that computes
pub_endpoint and subscribes) add a guard: if the config indicates dp_size > 1
then short-circuit with an appropriate gRPC unimplemented/unsupported error (or
log and close the stream) so clients cannot subscribe to a partial cache view
until per-rank streaming (KvEventBatch.dp_rank / per-rank endpoints) is
implemented.
- Around line 1024-1045: Move the ZMQ and decoder setup into the guarded
try/finally inside SubscribeKvEvents so failures are handled with context.abort
and sockets are always cleaned up: create zmq_ctx, sub_socket (and call
sub_socket.subscribe/connect) and the msgspec.msgpack.Decoder(KVEventBatch)
after entering the try block and ensure the existing finally still closes
sub_socket; if any of those setups raise, call
context.abort(grpc.StatusCode.INTERNAL, ...) as the other error handling in
SubscribeKvEvents does.
🪄 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: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: b15d9e6f-7f40-4f82-be23-fd86c764ce94
📒 Files selected for processing (3)
docs/getting-started/kv-events-cache-aware.mdgrpc_servicer/smg_grpc_servicer/vllm/kv_events.pygrpc_servicer/smg_grpc_servicer/vllm/servicer.py
Per review: hoist msgspec/zmq/KVEventBatch out of SubscribeKvEvents to the top of servicer.py (the module already imports vllm at import time). Signed-off-by: key4ng <rukeyang@gmail.com>
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 2bbfd54f03
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
| blocks = [] | ||
| for i, block_hash in enumerate(event.block_hashes): | ||
| start = i * block_size | ||
| end = start + block_size | ||
| block = common_pb2.KvBlock( | ||
| block_hash=to_int64(block_hash), | ||
| token_ids=list(event.token_ids[start:end]), | ||
| block_size=block_size, |
There was a problem hiding this comment.
Keep token slices aligned with compacted vLLM hashes
When vLLM skips null/masked blocks (e.g. sliding-window/sparse attention or Mamba align-mode prefix caching), BlockStored.block_hashes is compacted to only the stored blocks while event.token_ids still spans the full [num_cached_blocks, num_full_blocks) token range; this loop then pairs the first emitted hash with token_ids[0:block_size] even if that block was skipped. Fresh evidence beyond the earlier group-key discussion is vLLM's cache_full_blocks behavior: it skips masked/null blocks before appending to new_hashes, but passes the full contiguous token_ids slice in the event. In those cases SMG computes the content hash from the wrong tokens, so cache-aware routing learns blocks under incorrect token identities; consume vLLM's per-block alignment metadata or skip events whose len(block_hashes) * block_size != len(token_ids) rather than slicing by compacted index.
Useful? React with 👍 / 👎.
Matches grpc_servicer/README.md — the --grpc flag loads smg-grpc-servicer, so the worker serves gRPC (and KV events) rather than HTTP. Signed-off-by: key4ng <rukeyang@gmail.com>
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 7521d9cb79
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
Clarify SMG learns block size from BlockStored events; --block-size pins a non-default value. Signed-off-by: key4ng <rukeyang@gmail.com>
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 7854bcb851
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
| # ZMQ multipart: [topic, 8-byte big-endian seq, msgpack payload]. | ||
| if len(frames) < 3: | ||
| continue | ||
| zmq_seq = int.from_bytes(frames[1], "big") |
There was a problem hiding this comment.
Avoid forwarding zero-based KV sequence numbers
vLLM's ZMQ publisher numbers batches from 0, but model_gateway::worker::kv_event_monitor::process_stream uses last_seq == 0 as its uninitialized sentinel and only performs stale/gap checks when last_seq > 0. With this bridge forwarding the native sequence unchanged, after accepting batch 0 the gateway still treats the next batch as the first one; if batch 1 is dropped and batch 2 arrives, it is accepted instead of triggering GapDetected, leaving the cache index missing events. Offset vLLM sequences before building the proto, or change the gateway to track initialization separately.
Useful? React with 👍 / 👎.
Final benchmark — cache-aware routing & vLLM event-driven KV routing (supersedes my earlier benchmark comments)My earlier benchmark comments on this PR used a homegrown closed-loop client and a deliberately shrunk KV cache, which produced inflated / confounded numbers (the "7.3×" etc.). I redid everything with a standard open-loop methodology and a real load balancer; this comment supersedes those. Methodology
Results — achieved throughput (req/s) · e2e p99 · worker prefix-cache hit%, at a saturating offered rate
What it shows
Honest caveats
exact params
|
…bridge) Implements TokenSpeedSchedulerServicer.SubscribeKvEvents so SMG's cache_aware router can route against TokenSpeed workers' actual KV-cache state instead of the approximate token tree. The entire Rust/SMG consumer side (proto RPC, gRPC client dispatch, KvEventMonitor, PositionalIndexer, cache_aware) and the TokenSpeed ZMQ publisher (--kv-events-config) already existed; only the Python servicer bridge was missing — mirroring the vLLM bridge in #1652. - Promote the engine-neutral ZMQ->proto conversion (to_int64, endpoint_for_rank, convert_event, convert_batch, stream_kv_events) from vllm/kv_events.py to a shared smg_grpc_servicer/kv_events.py; vllm/kv_events.py re-exports them and keeps its vLLM-specific resolver. convert_batch now reads the DP rank from data_parallel_rank (vLLM) or attn_dp_rank (TokenSpeed). - Add tokenspeed/kv_events.py resolver: parses server_args.kv_events_config (JSON string) and returns the ZMQ endpoint iff enabled with publisher=zmq. - Implement SubscribeKvEvents in the TokenSpeed servicer (rank-0 subscription, no replay), matching the vLLM bridge. - Add the [tokenspeed] pyproject extra (pyzmq, msgspec). - Tests (engine-free, no TokenSpeed install): resolver + shared conversion unit tests and an in-process ZMQ PUB/SUB integration test using msgspec structs that mirror the TokenSpeed wire layout. - Docs: TokenSpeed worker launch section in kv-events-cache-aware.md. No Rust, proto, or tokenspeed-lib changes (the pinned tokenspeed already ships the KV-event publisher). Signed-off-by: key4ng <rukeyang@gmail.com>
…bridge) Implements TokenSpeedSchedulerServicer.SubscribeKvEvents so SMG's cache_aware router can route against TokenSpeed workers' actual KV-cache state instead of the approximate token tree. The entire Rust/SMG consumer side (proto RPC, gRPC client dispatch, KvEventMonitor, PositionalIndexer, cache_aware) and the TokenSpeed ZMQ publisher (--kv-events-config) already existed; only the Python servicer bridge was missing — mirroring the vLLM bridge in #1652. - Promote the engine-neutral ZMQ->proto conversion (to_int64, endpoint_for_rank, convert_event, convert_batch, stream_kv_events) from vllm/kv_events.py to a shared smg_grpc_servicer/kv_events.py; vllm/kv_events.py re-exports them and keeps its vLLM-specific resolver. convert_batch now reads the DP rank from data_parallel_rank (vLLM) or attn_dp_rank (TokenSpeed). - Add tokenspeed/kv_events.py resolver: parses server_args.kv_events_config (JSON string) and returns the ZMQ endpoint iff enabled with publisher=zmq. - Implement SubscribeKvEvents in the TokenSpeed servicer (rank-0 subscription, no replay), matching the vLLM bridge. - Add the [tokenspeed] pyproject extra (pyzmq, msgspec). - Tests (engine-free, no TokenSpeed install): resolver + shared conversion unit tests and an in-process ZMQ PUB/SUB integration test using msgspec structs that mirror the TokenSpeed wire layout. - Docs: TokenSpeed worker launch section in kv-events-cache-aware.md. No Rust, proto, or tokenspeed-lib changes (the pinned tokenspeed already ships the KV-event publisher). Signed-off-by: key4ng <rukeyang@gmail.com>
Description
Problem
KV-event-driven cache-aware routing — where the gateway routes each request to the worker whose KV cache already holds the longest prefix, using the worker's actual cache state — only worked for SGLang. vLLM gRPC workers silently fell back to the approximate token tree because the vLLM gRPC servicer never implemented
SubscribeKvEvents(the RPC returnedUNIMPLEMENTED, and SMG'sKvEventMonitordisabled the subscription for that worker — seemodel_gateway/src/worker/kv_event_monitor.rs:379).Notably, the entire SMG/Rust side is already engine-agnostic: the proto RPC (
common.protoSubscribeKvEvents/KvEventBatch), all three engine gRPC clients (impl_subscribe_kv_events!), the dispatch enum,KvEventMonitor, thePositionalIndexer, and thecache_awarepolicy. The only missing piece for vLLM was the Python servicer bridge.Solution
Implement
SubscribeKvEventsonVllmEngineServicer, bridging vLLM's in-process ZMQ KV-cache events (msgpack-encodedBlockStored/BlockRemoved/AllBlocksCleared) to the engine-neutralcommon.protoKvEventBatchstream that SMG already consumes.All conversion logic lives in a new vLLM-free helper module (
vllm/kv_events.py), so it is fully unit-testable without a GPU or a vLLM install (events are dispatched by class name; the ZMQ loop takes an injectable decoder). The servicer reads its config fromself.engine.vllm_config.kv_events_config(same pattern already used forkv_transfer_config), so no changes to the out-of-repo launch wiring are needed. No Rust changes.Changes
grpc_servicer/smg_grpc_servicer/vllm/kv_events.py(new):to_int64(accepts vLLM'sbytes | intblock hash),endpoint_for_rank,resolve_kv_events_config,convert_event/convert_batch, and thestream_kv_eventsZMQ→proto loop.grpc_servicer/smg_grpc_servicer/vllm/servicer.py: addSubscribeKvEvents+ resolve config in__init__.grpc_servicer/pyproject.toml: pinpyzmq/msgspecin the[vllm]extra.docs/getting-started/kv-events-cache-aware.md: document the vLLM worker launch (gRPC entrypoint +--kv-events-config).tests/test_vllm_kv_events.py(unit) +tests/test_vllm_kv_events_stream.py(real in-process ZMQ PUB/SUB integration).Test Plan
Verified in CI-equivalent environment (no GPU / vLLM not installed):
pytest grpc_servicer/tests/test_vllm_kv_events.py grpc_servicer/tests/test_vllm_kv_events_stream.py→ 28 passed, 1 skipped (the skip is the vLLM-gated servicer-wiring test).ruff check/ruff format --check→ clean.python -m py_compileon the servicer → clean.Remaining (requires a GPU + real vLLM worker — to run before merge):
python -m vllm.entrypoints.grpc_server) with--kv-events-config '{"enable_kv_cache_events": true, "publisher": "zmq", ...}'+ SMG--policy cache_aware; confirm the gateway logsStarting KV event subscription→Learned block_size, and that a repeated prefix routes to the worker already holding it.Follow-ups (out of scope): TRT support (no servicer + polling-only); multi-DP-rank subscription (currently rank-0 only); ZMQ replay-on-reconnect.
Checklist
ruff check/ruff formatpass (Python-only change; no Rust touched)pytest grpc_servicer/tests/(new tests) pass: 28 passed, 1 skippedSummary by CodeRabbit
New Features
Documentation
Chores