Skip to content

feat(grpc): add TokenSpeed KV cache event support (SubscribeKvEvents bridge) - #1771

Merged
slin1237 merged 5 commits into
mainfrom
feat/tokenspeed-kv-events
Jun 18, 2026
Merged

slin1237 merged 5 commits into
mainfrom
feat/tokenspeed-kv-events

Conversation

@key4ng

@key4ng key4ng commented Jun 18, 2026 •

Copy link
Copy Markdown
Member

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 — worked for SGLang and vLLM but not for TokenSpeed. TokenSpeed gRPC workers silently fell back to the approximate token tree because TokenSpeedSchedulerServicer never implemented SubscribeKvEvents (the RPC returned UNIMPLEMENTED, so SMG's KvEventMonitor disabled the subscription for that worker).

The entire Rust/SMG side is already engine-agnostic — the proto RPC (common.proto SubscribeKvEvents/KvEventBatch), the TokenSpeed gRPC client (impl_subscribe_kv_events!), GrpcClient::TokenSpeed dispatch, KvEventMonitor, the PositionalIndexer, and the cache_aware policy. The only missing piece was the Python servicer bridge.

Solution

Implement SubscribeKvEvents on TokenSpeedSchedulerServicer, bridging TokenSpeed's in-process ZMQ KV-cache events (msgpack BlockStored/BlockRemoved/AllBlocksCleared) to the engine-neutral common.proto KvEventBatch stream SMG already consumes. No Rust, proto, or TokenSpeed-engine changes.

The ZMQ→proto conversion is engine-neutral, so it is promoted to a shared module reused by both the vLLM and TokenSpeed bridges; each engine keeps only its own config resolver.

Changes

  • grpc_servicer/smg_grpc_servicer/kv_events.py (new): engine-neutral ZMQ→proto helpers (to_int64, endpoint_for_rank, convert_event, convert_batch, stream_kv_events), promoted from vllm/kv_events.py. convert_batch reads the DP rank from data_parallel_rank (vLLM) or attn_dp_rank (TokenSpeed).
  • grpc_servicer/smg_grpc_servicer/vllm/kv_events.py: re-exports the shared helpers; keeps the vLLM-specific resolver.
  • grpc_servicer/smg_grpc_servicer/tokenspeed/kv_events.py (new): TokenSpeed config resolver — parses server_args.kv_events_config and returns the ZMQ endpoint iff events are enabled with publisher=zmq.
  • grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py: implement SubscribeKvEvents (rank-0 subscription, no replay), mirroring the vLLM bridge.
  • grpc_servicer/pyproject.toml: add the [tokenspeed] extra (pyzmq, msgspec).
  • Tests: engine-free unit tests for the resolver + shared conversion, and an in-process ZMQ PUB/SUB integration test using msgspec structs that mirror the TokenSpeed wire layout.
  • Docs: TokenSpeed worker launch section in docs/getting-started/kv-events-cache-aware.md.

Test Plan

cargo +nightly fmt, ruff check/ruff format, and pytest grpc_servicer/tests/ pass for the new/changed files. Existing test_vllm_kv_events*.py remain green (the shared-module refactor is import-compatible).

Manual E2E (H100)

Built SMG + the servicer from this branch, ran live TokenSpeed gRPC workers (DeepSeek-R1-Distill-Qwen-32B) under SMG --policy cache_aware:

  • gateway logs Starting KV event subscription → KV event stream connected start_seq=0 per worker (RPC streams instead of UNIMPLEMENTED)
  • Learned block_size from KV event ... block_size=64 (block size learned from real BlockStored events through the bridge)
  • repeated-prefix requests route stickily to the worker holding the prefix
  • a worker launched without --kv-events-config → Backend does not implement SubscribeKvEvents, disabling KV event subscription (graceful fallback)

Benchmark (event-driven vs approximate token tree)

vllm bench serve --dataset-name prefix_repetition (600 prefixes × 1000-token shared prefix + 50-token suffix, 128-token output, 3000 requests, open-loop Poisson, saturating rate 32), through 4 SMG gateways (nginx round-robin, no mesh) in front of 3 TokenSpeed workers (DeepSeek-R1-Distill-Qwen-32B, TP=2). TokenSpeed workers launched with --disable-kvstore and --kv-events-config '{"enable_kv_cache_events": true, "publisher": "zmq", ...}'. 8×H100, 3000/3000 successful per run.

metric approximate tree (events off) event-driven (this PR) delta
prefix-cache hit rate 24.4% 48.5% +24 pts (~2×)
Mean TTFT 1699 ms 923 ms −46%
Median TTFT 258 ms 141 ms −45%
P95 TTFT 10256 ms 4724 ms −54%
P99 TTFT 23227 ms 19396 ms −16%
Request throughput 29.3 req/s 30.3 req/s +3.4%

Event-driven cache-aware routing roughly doubles the prefix-cache hit rate and halves TTFT versus the approximate token tree for TokenSpeed workers. The win is concentrated in prefill/TTFT (expected — prefix-cache reuse cuts prefill work); end-to-end latency converges at this scale since the 128-token decode dominates.

Checklist
  • cargo +nightly fmt passes
  • cargo clippy --all-targets --all-features -- -D warnings passes (no Rust changes)
  • (Optional) Documentation updated
  • (Optional) Please join us on Slack #sig-smg to discuss, review, and merge PRs

🤖 Generated with Claude Code

Summary by CodeRabbit

Release Notes

  • New Features

    • Added TokenSpeed KV-cache event streaming via a new gRPC subscription RPC for cache-aware routing.
  • Documentation

    • Added setup guidance for TokenSpeed KV-event publishing and updated KV-events reference links.
  • Enhancements

    • Added a shared KV-events conversion/streaming layer and refactored vLLM integration to reuse it.
    • Introduced optional TokenSpeed dependencies needed for KV-event bridging.
  • Tests

    • Added unit and end-to-end tests covering KV-events config resolution, conversion, and stream decoding.

@github-actions github-actions Bot added documentation Improvements or additions to documentation dependencies Dependency updates grpc gRPC client and router changes tests Test changes labels Jun 18, 2026
@coderabbitai

coderabbitai Bot commented Jun 18, 2026 •

Copy link
Copy Markdown

Review Change Stack

Note

Reviews paused

It 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 reviews.auto_review.auto_pause_after_reviewed_commits setting.

Use the following commands to manage reviews:

  • @coderabbitai resume to resume automatic reviews.
  • @coderabbitai review to trigger a single review.

Use the checkboxes below for quick actions:

  • ▶️ Resume reviews
  • 🔍 Trigger review

No actionable comments were generated in the recent review. 🎉

ℹ️ Recent review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: ASSERTIVE

Plan: Pro

Run ID: c3effa2e-006e-49f8-8d8b-9565fe20e108

📥 Commits

Reviewing files that changed from the base of the PR and between 23844a9 and b51305b.

📒 Files selected for processing (1)
  • docs/getting-started/kv-events-cache-aware.md

📝 Walkthrough

Walkthrough

Extracts engine-neutral KV-event ZMQ→proto conversion helpers (to_int64, endpoint_for_rank, convert_event, convert_batch, stream_kv_events) into a new shared smg_grpc_servicer/kv_events.py. The vLLM module is converted to a re-export shim. A new SubscribeKvEvents RPC is added to TokenSpeedSchedulerServicer, backed by a TokenSpeed-specific JSON config resolver and pyzmq/msgspec optional dependencies, with unit and integration tests and documentation.

Changes

TokenSpeed KV-event bridge

Layer / File(s) Summary
Engine-neutral KV-event conversion and streaming
grpc_servicer/smg_grpc_servicer/kv_events.py
New shared module implementing to_int64, endpoint_for_rank, convert_event, convert_batch, and async stream_kv_events; ZMQ→proto translation loop polls with timeout, validates multipart frame structure, extracts 8-byte sequence numbers, isolates per-frame decode failures.
vLLM kv_events refactored to compatibility shim
grpc_servicer/smg_grpc_servicer/vllm/kv_events.py
Local conversion/streaming implementation removed; five helpers re-exported from the new shared module via __all__; adds resolve_kv_events_config reading engine.vllm_config.kv_events_config.
TokenSpeed config resolver and optional deps
grpc_servicer/smg_grpc_servicer/tokenspeed/kv_events.py, grpc_servicer/pyproject.toml
ResolvedKvEventsConfig dataclass and resolve_kv_events_config parse raw JSON from server_args.kv_events_config, applying endpoint/topic defaults and returning None when disabled or misconfigured; pyproject.toml gains the tokenspeed optional group with pyzmq>=25.0.0 and msgspec>=0.18.0.
TokenSpeed servicer SubscribeKvEvents RPC
grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py
Imports msgspec, zmq, and KV-event bridge helpers; stores _kv_events_config during __init__; adds SubscribeKvEvents server-streaming RPC that gates on config, connects ZMQ SUB for rank 0, decodes msgpack KVEventBatch messages, yields converted proto batches, and closes the socket in finally.
Test infrastructure and unit tests
grpc_servicer/tests/conftest.py, grpc_servicer/tests/test_tokenspeed_kv_events.py
conftest.py adds sys.path control for in-repo imports; unit tests cover config resolution edge cases, convert_event oneof variant mapping, and attn_dp_rank→dp_rank conversion with optional field handling.
Integration test for stream decoding
grpc_servicer/tests/test_tokenspeed_kv_events_stream.py
Exercises stream_kv_events end-to-end over a real ZMQ PUB/SUB socket pair with msgspec-encoded multipart frames, verifying sequence number extraction, batch ordering, and proto field correctness.
Documentation updates
docs/getting-started/kv-events-cache-aware.md
Adds a "Launch a TokenSpeed worker" section with install steps, --kv-events-config field reference table, and block-size learning notes; expands the Reference section with links to the shared module, TokenSpeed bridge, config resolver, and upstream KV-events config.

Sequence Diagram(s)

sequenceDiagram
  participant Client as gRPC Client
  participant Servicer as TokenSpeedSchedulerServicer
  participant Config as resolve_kv_events_config
  participant ZMQ as ZMQ SUB Socket
  participant Shared as stream_kv_events / convert_batch

  Client->>Servicer: SubscribeKvEvents(request)
  Servicer->>Config: resolve_kv_events_config(server_args)
  Config-->>Servicer: ResolvedKvEventsConfig (endpoint, topic) or None
  alt config is None
    Servicer-->>Client: UNIMPLEMENTED abort
  else config present
    Servicer->>ZMQ: connect(endpoint_for_rank(endpoint, rank=0))
    Servicer->>ZMQ: subscribe(topic)
    Servicer->>Shared: stream_kv_events(sub_socket, msgspec.decode, ...)
    Shared-->>Client: send_initial_metadata()
    loop until cancelled
      ZMQ-->>Shared: multipart [seq_bytes, msgpack payload]
      Shared->>Shared: convert_batch(decoded_batch)
      Shared-->>Client: yield KvEventBatch proto
    end
    Servicer->>ZMQ: close()
  end
Loading

Estimated code review effort

🎯 4 (Complex) | ⏱️ ~60 minutes

Possibly related PRs

  • lightseekorg/smg#558: Adds the client-side subscribe_kv_events plumbing that calls the server-side SubscribeKvEvents RPC this PR implements.
  • lightseekorg/smg#677: Extends the Rust KvEventMonitor to "learn block_size" from the BlockStored proto events that this PR's conversion helpers construct from engine KV event batches.
  • lightseekorg/smg#1652: Originally implemented the vLLM-side KV-event helpers that this PR extracts into the shared engine-neutral module.

Suggested reviewers

  • CatherineSue
  • slin1237
  • njhill

Poem

🐇 A ZMQ wire, a proto in hand,
The TokenSpeed signals now flow as planned.
One shared module for vLLM and new,
Blocks stored and removed all converted through.
Hop hop, the cache-aware routing is grand! 🎉

🚥 Pre-merge checks | ✅ 4 | ❌ 1

❌ Failed checks (1 warning)

Check name Status Explanation Resolution
Docstring Coverage ⚠️ Warning Docstring coverage is 25.64% which is insufficient. The required threshold is 80.00%. Write docstrings for the functions missing them to satisfy the coverage threshold.
✅ Passed checks (4 passed)
Check name Status Explanation
Description Check ✅ Passed Check skipped - CodeRabbit’s high-level summary is enabled.
Title check ✅ Passed The pull request title clearly and specifically summarizes the main change: adding TokenSpeed KV cache event support via the SubscribeKvEvents bridge.
Linked Issues check ✅ Passed Check skipped because no linked issues were found for this pull request.
Out of Scope Changes check ✅ Passed Check skipped because no linked issues were found for this pull request.

✏️ Tip: You can configure your own custom pre-merge checks in the settings.

✨ Finishing Touches
🧪 Generate unit tests (beta)
  • Create PR with unit tests
  • Commit unit tests in branch feat/tokenspeed-kv-events

Comment @coderabbitai help to get the list of available commands and usage tips.

@gemini-code-assist gemini-code-assist Bot 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.

Code Review

This pull request introduces support for TokenSpeed KV-event cache-aware routing by implementing the SubscribeKvEvents gRPC stream. It refactors the existing engine-neutral ZMQ-to-proto conversion and streaming logic into a shared kv_events.py module, which is now utilized by both vLLM and TokenSpeed bridges. Additionally, it adds config resolution, optional package dependencies, documentation, and unit/integration tests. Feedback on the implementation highlights the need to lazy-import KVEventBatch inside the streaming function rather than at the top level to prevent import errors when TokenSpeed is not installed, and to wrap the socket connection and subscription setup inside the try block to ensure proper resource cleanup in the finally block.

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.

Comment thread grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py
Comment on lines +577 to +607
config = self._kv_events_config

# DP attention publishes one PUB socket per rank (port + rank) with
# independent sequence counters; subscribing to several on one socket
# interleaves them and breaks gap detection. Subscribe to rank 0 only.
pub_endpoint = endpoint_for_rank(config.endpoint, 0)

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)

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: # noqa: BLE001
logger.exception("SubscribeKvEvents failed")
await context.abort(grpc.StatusCode.INTERNAL, str(e))
finally:
sub_socket.close(linger=0)
logger.info("SubscribeKvEvents: stream closed")

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.

medium

Wrap the socket subscription, connection, and decoder initialization inside the try block. If any of these steps fail (e.g., connection error or missing tokenspeed package), the socket will still be properly closed in the finally block, preventing resource leaks. Additionally, lazy-import KVEventBatch here to keep the module importable without tokenspeed installed.

        config = self._kv_events_config

        # DP attention publishes one PUB socket per rank (port + rank) with
        # independent sequence counters; subscribing to several on one socket
        # interleaves them and breaks gap detection. Subscribe to rank 0 only.
        pub_endpoint = endpoint_for_rank(config.endpoint, 0)

        zmq_ctx = zmq.asyncio.Context.instance()
        sub_socket = zmq_ctx.socket(zmq.SUB)
        try:
            sub_socket.subscribe(config.topic.encode("utf-8"))
            sub_socket.connect(pub_endpoint)
            logger.info("SubscribeKvEvents: connected to ZMQ endpoint %s", pub_endpoint)

            from tokenspeed.runtime.pd.kv_events import KVEventBatch
            decoder = msgspec.msgpack.Decoder(KVEventBatch)

            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:  # noqa: BLE001
            logger.exception("SubscribeKvEvents failed")
            await context.abort(grpc.StatusCode.INTERNAL, str(e))
        finally:
            sub_socket.close(linger=0)
            logger.info("SubscribeKvEvents: stream closed")
References
  1. To support optional dependencies, import modules from the extra inside the functions that use them (lazy import) rather than at the top level to prevent ImportError.
  2. When managing resources that are acquired and need to be released, ensure that all acquired resources are tracked immediately and consistently, especially when exceptions might occur during subsequent operations, to prevent resource leaks.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

Applied the try-wrapping in 23844a9 — subscribe/connect and the Decoder(...) now run inside the try, so the finally always closes the socket (e.g. on a malformed endpoint). Kept KVEventBatch as a top-level import per the project's import convention (the module already requires tokenspeed).

Comment on lines +569 to +575
if self._kv_events_config is None:
await context.abort(
grpc.StatusCode.UNIMPLEMENTED,
"KV cache events not enabled. Start TokenSpeed with "
"--kv-events-config "
'\'{"enable_kv_cache_events": true, "publisher": "zmq"}\'',
)

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

🟡 Nit: Missing defensive return after context.abort(). Every other early-abort in this file (e.g., lines 146–147, 150–151, 489) follows the abort → return pattern. While context.abort() raises AbortError so the return is technically unreachable, it serves as defense-in-depth and keeps the pattern consistent — without it, if abort ever doesn't raise, config would be None and config.endpoint on line 582 would AttributeError.

Suggested change
if self._kv_events_config is None:
await context.abort(
grpc.StatusCode.UNIMPLEMENTED,
"KV cache events not enabled. Start TokenSpeed with "
"--kv-events-config "
'\'{"enable_kv_cache_events": true, "publisher": "zmq"}\'',
)
await context.abort(
grpc.StatusCode.UNIMPLEMENTED,
"KV cache events not enabled. Start TokenSpeed with "
"--kv-events-config "
'\'{"enable_kv_cache_events": true, "publisher": "zmq"}\'',
)
return

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

Done in 23844a9 — added return after the UNIMPLEMENTED abort to match the abort → return pattern used elsewhere in this file.

Comment on lines +600 to +604
except asyncio.CancelledError:
pass
except Exception as e: # noqa: BLE001
logger.exception("SubscribeKvEvents failed")
await context.abort(grpc.StatusCode.INTERNAL, str(e))

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

🟡 Nit: Missing except grpc.aio.AbortError: raise guard — the Generate method (line 222) has one so that an AbortError propagated from inside the try block (e.g., from send_initial_metadata on an already-aborted RPC) doesn't get caught here, logged as "SubscribeKvEvents failed", and re-wrapped as INTERNAL.

Suggested change
except asyncio.CancelledError:
pass
except Exception as e: # noqa: BLE001
logger.exception("SubscribeKvEvents failed")
await context.abort(grpc.StatusCode.INTERNAL, str(e))
except asyncio.CancelledError:
pass
except grpc.aio.AbortError:
raise
except Exception as e: # noqa: BLE001
logger.exception("SubscribeKvEvents failed")
await context.abort(grpc.StatusCode.INTERNAL, str(e))

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

Done in 23844a9 — added except grpc.aio.AbortError: raise so an abort propagated from inside the try isn't caught, logged as "SubscribeKvEvents failed", and re-wrapped as INTERNAL. Matches Generate.

@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.

Clean, well-structured PR. The refactoring to promote engine-neutral ZMQ→proto helpers to a shared module is done carefully — vLLM re-exports preserve backwards compatibility, and the TokenSpeed bridge mirrors the vLLM one closely. Config resolution, conversion, and streaming are all well-tested without requiring an engine install.

Two nits on exception-handling consistency in SubscribeKvEvents (defensive return after abort + AbortError re-raise guard), both matching patterns already established in the same file's Generate method.

0 🔴 Important · 2 🟡 Nit · 0 🟣 Pre-existing

@coderabbitai coderabbitai 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.

Actionable comments posted: 1

Caution

Some comments are outside the diff and can’t be posted inline due to platform limitations.

⚠️ Outside diff range comments (1)
grpc_servicer/pyproject.toml (1)

32-41: ⚠️ Potential issue | 🟠 Major

Update to current versions of pyzmq and msgspec in the tokenspeed optional-dependency group.

The floor constraints pyzmq>=25.0.0 and msgspec>=0.18.0 are significantly outdated. As of June 2026, the latest stable versions are pyzmq==27.1.0 and msgspec==0.21.1 respectively. The declared versions are from 2023 (approximately 2–3 years old) and, while no critical security vulnerabilities are documented for these specific versions, they lack numerous bug fixes, performance improvements, and refinements released since.

Update the floor constraints to more recent versions—ideally pyzmq>=27.0.0 and msgspec>=0.21.0, or at minimum closer to current releases.

🤖 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/pyproject.toml` around lines 32 - 41, The tokenspeed
optional-dependency group in pyproject.toml has outdated floor constraints for
its dependencies. Update the pyzmq constraint from pyzmq>=25.0.0 to
pyzmq>=27.0.0 and the msgspec constraint from msgspec>=0.18.0 to msgspec>=0.21.0
to align with more recent stable versions that include bug fixes and performance
improvements instead of versions from 2–3 years ago.
🤖 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/tests/test_tokenspeed_kv_events_stream.py`:
- Around line 72-106: The test uses fixed timing delays (asyncio.sleep calls at
line 72 and within the batches loop at line 105) which cause nondeterministic
failures due to ZMQ slow-joiner behavior dropping early frames. Replace the
initial 0.2 second sleep before publishing with a handshake mechanism that
confirms the subscriber in the consume() function is actively listening (for
example, by having the consumer signal readiness or by sending and awaiting an
initial acknowledgment message). Similarly, replace the 0.05 second sleep
between batch sends with deterministic synchronization, such as waiting for each
batch to be collected before sending the next one, to ensure the consumer
processes each message before the publisher sends the next batch.

---

Outside diff comments:
In `@grpc_servicer/pyproject.toml`:
- Around line 32-41: The tokenspeed optional-dependency group in pyproject.toml
has outdated floor constraints for its dependencies. Update the pyzmq constraint
from pyzmq>=25.0.0 to pyzmq>=27.0.0 and the msgspec constraint from
msgspec>=0.18.0 to msgspec>=0.21.0 to align with more recent stable versions
that include bug fixes and performance improvements instead of versions from 2–3
years ago.
🪄 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: 007ac31d-5c69-4fdd-b396-79954e8935fb

📥 Commits

Reviewing files that changed from the base of the PR and between 1288d16 and f4c0b1b.

📒 Files selected for processing (9)
  • docs/getting-started/kv-events-cache-aware.md
  • docs/proposals/2026-06-16-tokenspeed-kv-events.md
  • grpc_servicer/pyproject.toml
  • grpc_servicer/smg_grpc_servicer/kv_events.py
  • grpc_servicer/smg_grpc_servicer/tokenspeed/kv_events.py
  • grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py
  • grpc_servicer/smg_grpc_servicer/vllm/kv_events.py
  • grpc_servicer/tests/test_tokenspeed_kv_events.py
  • grpc_servicer/tests/test_tokenspeed_kv_events_stream.py

Comment thread grpc_servicer/tests/test_tokenspeed_kv_events_stream.py
@key4ng
key4ng force-pushed the feat/tokenspeed-kv-events branch from f4c0b1b to 64310da Compare June 18, 2026 01:23

@coderabbitai coderabbitai 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.

Actionable comments posted: 5

🤖 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/kv_events.py`:
- Around line 154-155: In the frame validation block where len(frames) < 3 is
checked, add a debug-level log statement before the continue statement to
capture information about the malformed frame. Include relevant context such as
the frame length and any identifying details that would help diagnose protocol
mismatches during debugging, without impacting production performance.

In `@grpc_servicer/smg_grpc_servicer/tokenspeed/kv_events.py`:
- Around line 68-71: After assigning the default values for endpoint and topic
in the configuration resolution logic, add type validation to ensure both
endpoint and topic are strings. If either endpoint or topic is not a string
type, log a warning message indicating that KV-event streaming is being disabled
due to invalid configuration types, and return None from the function to cleanly
disable this feature instead of allowing invalid types to propagate to
SubscribeKvEvents where they would cause runtime errors later.

In `@grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py`:
- Around line 584-607: The socket setup code including zmq context
instantiation, socket creation, subscription, and connection to the endpoint are
all executed outside the try/finally block in the SubscribeKvEvents method. This
means if any of those operations fail, the exception handling logic and socket
cleanup will not be executed. Move all socket setup code (zmq_ctx instantiation,
sub_socket socket creation, subscribe call, and connect call) into the try block
before the stream_kv_events call, and add a guard in the finally block to check
if sub_socket was successfully created before attempting to close it.
- Around line 24-29: Move the imports of msgspec and zmq-related modules (zmq,
zmq.asyncio) from module scope into the SubscribeKvEvents method where they are
actually used. Wrap these imports in a try-except block within SubscribeKvEvents
to catch ImportError when optional dependencies are not installed, and use gRPC
abort with UNIMPLEMENTED status when the import fails. Additionally, move the
socket setup code (currently at lines 584-587) from before the try block into
the try-except block so socket creation failures are also properly handled and
reported through gRPC error handling instead of causing an unhandled exception
at import time.

In `@grpc_servicer/tests/test_tokenspeed_kv_events.py`:
- Around line 57-59: The test test_none_when_publisher_null is passing the
string "null" instead of an actual None value for the publisher parameter, which
tests the wrong code path. Change the publisher parameter in the _cfg function
call from the string "null" to None (or None value as appropriate for your
language) to properly validate the JSON null contract behavior where publisher
should default to "zmq" when enabled.
🪄 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: e330510c-db99-496c-9f5e-876d7941a406

📥 Commits

Reviewing files that changed from the base of the PR and between f4c0b1b and 64310da.

📒 Files selected for processing (8)
  • docs/getting-started/kv-events-cache-aware.md
  • grpc_servicer/pyproject.toml
  • grpc_servicer/smg_grpc_servicer/kv_events.py
  • grpc_servicer/smg_grpc_servicer/tokenspeed/kv_events.py
  • grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py
  • grpc_servicer/smg_grpc_servicer/vllm/kv_events.py
  • grpc_servicer/tests/test_tokenspeed_kv_events.py
  • grpc_servicer/tests/test_tokenspeed_kv_events_stream.py

Comment thread grpc_servicer/smg_grpc_servicer/kv_events.py
Comment thread grpc_servicer/smg_grpc_servicer/tokenspeed/kv_events.py
Comment thread grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py
Comment thread grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py
Comment thread grpc_servicer/tests/test_tokenspeed_kv_events.py Outdated
key4ng added 4 commits June 18, 2026 10:17
…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>
…repo root

test_vllm_kv_events* load vllm/kv_events.py by file path, which now
re-exports from the shared smg_grpc_servicer.kv_events module. When pytest
runs from the repo root (CI), grpc_servicer/ is not on sys.path so that
package import fails at collection. Add a tests/conftest.py that prepends
this repo's grpc_servicer/ to sys.path (taking precedence over any stale
installed copy).

Signed-off-by: key4ng <rukeyang@gmail.com>
Signed-off-by: key4ng <rukeyang@gmail.com>
- defensive return after the UNIMPLEMENTED abort (matches Generate/others)
- re-raise grpc.aio.AbortError instead of re-wrapping it as INTERNAL
- move socket subscribe/connect/decoder into the try so the finally always
  closes the socket (e.g. on a malformed endpoint)
- resolver: reject non-string endpoint/topic and disable cleanly instead of
  failing later as an opaque INTERNAL
- tests: cover explicit JSON-null publisher (-> zmq) and non-string config

Signed-off-by: key4ng <rukeyang@gmail.com>
@key4ng
key4ng force-pushed the feat/tokenspeed-kv-events branch from 5acf68c to 23844a9 Compare June 18, 2026 17:26
@key4ng

key4ng commented Jun 18, 2026

Copy link
Copy Markdown
Member Author

Rebased onto latest main. Addressed the review feedback in 23844a96:

Applied

  • SubscribeKvEvents: defensive return after the UNIMPLEMENTED abort, and a except grpc.aio.AbortError: raise guard — matches the existing Generate pattern (claude).
  • Moved socket subscribe/connect/decoder into the try so finally always closes the socket, e.g. on a malformed endpoint (gemini, coderabbit).
  • Resolver now rejects non-string endpoint/topic and disables cleanly instead of surfacing later as an opaque INTERNAL (coderabbit).
  • Added tests for the explicit JSON-null publisher (→ zmq) and non-string config (coderabbit).

Skipped (with reason)

  • Lazy-import KVEventBatch / top-level msgspec/zmq (gemini, coderabbit): the servicer already requires tokenspeed to import, and the codebase prefers top-level imports; pyzmq/msgspec are declared in the [tokenspeed] extra.
  • Bump pyzmq/msgspec floors (coderabbit): kept aligned with the existing [vllm] extra — these are minimum constraints, pip resolves to the latest, and 25.x/0.18 already provide the APIs used.
  • Stream-test fixed sleeps (coderabbit): mirrors the established test_vllm_kv_events_stream.py; the 0.2s slow-joiner wait is the existing pattern.
  • Debug-log on short frames (coderabbit): low-value in the shared module; decode errors already warn.

@chatgpt-codex-connector chatgpt-codex-connector 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.

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: 23844a9642

ℹ️ 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".

Comment thread grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py
Signed-off-by: key4ng <rukeyang@gmail.com>
@slin1237
slin1237 merged commit 51b114b into main Jun 18, 2026
76 checks passed
@slin1237
slin1237 deleted the feat/tokenspeed-kv-events branch June 18, 2026 21:22
@key4ng

key4ng commented Jun 18, 2026

Copy link
Copy Markdown
Member Author

Gateway-sweep benchmark — 1 / 2 / 4 gateways

Replicated the PR-1652 methodology for TokenSpeed: round_robin vs cache_aware (approximate token tree, events off) vs cache_aware (event-driven, events on), across 1 / 2 / 4 gateways. 3 workers DeepSeek-R1-Distill-Qwen-32B (TP=2), --disable-kvstore, saturating offered rate 32, 1800 requests/run, 0 failed requests in all runs.

1 gateway

policy throughput (req/s) prefix-cache hit mean TTFT (ms) median TTFT (ms) P99 TTFT (ms) P99 E2EL (ms)
round_robin 11.58 17.4% 15034 136 104857 132618
cache_aware (approx) 28.85 26.5% 1657 184 24086 42808
cache_aware (event) 27.52 24.6% 1796 314 22173 42239

2 gateways

policy throughput (req/s) prefix-cache hit mean TTFT (ms) median TTFT (ms) P99 TTFT (ms) P99 E2EL (ms)
round_robin 10.82 15.1% 16087 168 120670 146749
cache_aware (approx) 28.72 21.6% 1896 268 24398 42985
cache_aware (event) 28.82 23.7% 1334 295 16918 42988

4 gateways

policy throughput (req/s) prefix-cache hit mean TTFT (ms) median TTFT (ms) P99 TTFT (ms) P99 E2EL (ms)
round_robin 11.55 16.0% 15213 185 103526 133447
cache_aware (approx) 25.74 16.5% 2359 431 26040 47420
cache_aware (event) 28.32 26.5% 1537 290 22028 42224

What it shows

  1. cache_aware ≫ round_robin everywhere — ~2.5× throughput and ~5× better P99 TTFT (round_robin saturates at ~11 req/s with a runaway tail; cache_aware sustains ~28).
  2. Event-driven vs approximate tree scales with gateway count:
    • 1 gateway: no advantage (event marginally lower) — the approximate tree is accurate when a single gateway sees all traffic.
    • 2 gateways: event-driven pulls ahead — mean TTFT 1.33s vs 1.90s (−30%), P99 TTFT 16.9s vs 24.4s.
    • 4 gateways: event-driven clearly wins — hit 26.5% vs 16.5%, throughput 28.3 vs 25.7 req/s, mean TTFT 1.54s vs 2.36s (−35%).
  3. Robustness is the point: event-driven barely moves as gateways increase (hit 24.6 → 23.7 → 26.5%, throughput 27.5 → 28.8 → 28.3), while the approximate tree degrades (hit 26.5 → 21.6 → 16.5%, throughput 28.9 → 28.7 → 25.7) because each gateway's local tree fragments across replicas. Event-driven routes every gateway against the worker's actual KV state, so it stays consolidated.

Absolute hit rates are modest (16–27%) because --disable-kvstore limits caching to the device tier (~307k tokens/worker), so the 600k-token working set doesn't fully fit even with perfect routing. The relative pattern (event-driven holds, approximate tree degrades with gateway count) matches the vLLM result.

Test setup & methodology

Hardware: 8×H100 80GB, CUDA 13.

Workers (3, TP=2, GPUs 0–5):

CUDA_VISIBLE_DEVICES=0,1 python -m smg_grpc_servicer.tokenspeed \
  --model deepseek-ai/DeepSeek-R1-Distill-Qwen-32B --tensor-parallel-size 2 \
  --host 127.0.0.1 --port 31001 --disable-kvstore \
  --kv-events-config '{"enable_kv_cache_events": true, "publisher": "zmq", "endpoint": "tcp://*:6001", "topic": "kv-events"}'
# (workers 2/3 on GPUs 2,3 / 4,5, ports 31002/31003, zmq 6002/6003)
# events-off runs (round_robin, approx) launch the same workers without --kv-events-config

Gateways (N ∈ {1,2,4}):

RUST_LOG=warn ./target/release/smg launch \
  --worker-urls grpc://127.0.0.1:31001 grpc://127.0.0.1:31002 grpc://127.0.0.1:31003 \
  --model-path deepseek-ai/DeepSeek-R1-Distill-Qwen-32B --policy {round_robin|cache_aware} \
  --block-size 64 --host 127.0.0.1 --port 30000 --prometheus-port 29001
# N gateways on ports 30000/30010/30020/30030, fronted by nginx round-robin
# (proxy_buffering off, upstream keepalive, worker_rlimit_nofile raised to avoid FD exhaustion)

Client:

vllm bench serve --backend openai --model deepseek-ai/DeepSeek-R1-Distill-Qwen-32B \
  --base-url http://127.0.0.1:<nginx-or-gateway-port> --endpoint /v1/completions \
  --dataset-name prefix_repetition \
  --prefix-repetition-num-prefixes 600 --prefix-repetition-prefix-len 1000 \
  --prefix-repetition-suffix-len 50 --prefix-repetition-output-len 128 \
  --num-prompts 1800 --request-rate 32 --burstiness 1.0 --ignore-eos \
  --percentile-metrics ttft,e2el --metric-percentiles 95,99

Strategy:

  • Cold worker restart per config so every run starts from an empty cache (fair comparison).
  • Per-run event verification: with events on, the gateways log KV event stream connected ×(gateways × workers); with events off, 0 — confirming the policy under test.
  • Saturating rate: offered 32 req/s is past the cluster knee (round_robin tops out ~11 req/s; cache_aware ~28–29), so the better-routing policy converts its cache advantage into lower TTFT / higher throughput.
  • Metrics: throughput and TTFT/E2EL percentiles from the vllm bench serve summary; cluster prefix-cache hit rate computed as Σ#cached-token / Σ(#cached-token + #new-token) over the workers' prefill logs for the run.
  • Open-loop Poisson arrivals (--burstiness 1.0), prefix_repetition dataset = 600 distinct 1000-token prefixes each repeated (1800 prompts = 3×), the standard prefix-caching workload.

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

Labels

dependencies Dependency updates documentation Improvements or additions to documentation grpc gRPC client and router changes tests Test changes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants