feat(grpc): add subscribe_kv_events to all backend clients - #558
Conversation
|
No actionable comments were generated in the recent review. 🎉 ℹ️ Recent review infoConfiguration used: Organization UI Review profile: ASSERTIVE Plan: Pro 📒 Files selected for processing (5)
📝 WalkthroughWalkthroughAdds a shared macro Changes
Sequence Diagram(s)sequenceDiagram
participant App
participant GrpcClient as "GrpcClient (gateway)"
participant Backend as "Engine Client\n(Sglang/Vllm/Trtllm)"
participant Engine as "gRPC Engine Server"
App->>GrpcClient: subscribe_kv_events(start_seq)
GrpcClient->>Backend: subscribe_kv_events(start_seq)
Backend->>Engine: gRPC SubscribeKvEventsRequest(start_sequence_number)
Engine-->>Backend: stream KvEventBatch...
Backend-->>GrpcClient: forward stream
GrpcClient-->>App: tonic::Streaming<KvEventBatch>
Estimated code review effort🎯 3 (Moderate) | ⏱️ ~20 minutes Poem
🚥 Pre-merge checks | ✅ 3✅ Passed checks (3 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches
🧪 Generate unit tests (beta)
Comment |
Summary of ChangesHello, I'm Gemini Code Assist1! I'm currently reviewing this pull request and will post my feedback shortly. In the meantime, here's a summary to help you and other reviewers quickly get up to speed! This pull request significantly enhances the Highlights
Changelog
Activity
Using Gemini Code AssistThe full guide for Gemini Code Assist can be found on our documentation page, here are some quick tips. Invoking Gemini You can request assistance from Gemini at any point by creating a comment using either
Customization To customize Gemini Code Assist for GitHub experience, repository maintainers can create a configuration file and/or provide a custom code review style guide (such as PEP-8 for Python) by creating and adding files to a Limitations & Feedback Gemini Code Assist may make mistakes. Please leave feedback on any instances where its feedback is incorrect or counter productive. You can react with 👍 and 👎 on @gemini-code-assist comments. If you're interested in giving your feedback about your experience with Gemini Code Assist for Github and other Google products, sign up here. You can also get AI-powered code generation, chat, as well as code reviews directly in the IDE at no cost with the Gemini Code Assist IDE Extension. Footnotes
|
There was a problem hiding this comment.
Code Review
This pull request effectively adds the subscribe_kv_events functionality to all backend gRPC clients by using a shared macro. This is a great approach to reduce code duplication and maintain consistency. The implementation is clean and follows existing patterns in the codebase. The suggestion to improve naming consistency in the unified GrpcClient wrapper aligns with prioritizing external API specifications.
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@grpc_client/src/lib.rs`:
- Around line 67-74: The subscribe_kv_events wrapper currently builds a
tonic::Request and calls client.subscribe_kv_events without injecting trace
context, breaking distributed traces; update the code that constructs
$crate::common_proto::SubscribeKvEventsRequest into a tonic::Request and inject
the current trace/OTel context into the request's metadata (using the same
metadata propagation helper or approach used by the other RPC wrappers like
generate/embed) before calling client.subscribe_kv_events(request). Ensure you
modify the code around Request::new, the SubscribeKvEventsRequest construction,
and the client.subscribe_kv_events call so the request carries the trace headers
exactly as the other traced RPC entrypoints do.
ℹ️ Review info
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
📒 Files selected for processing (5)
grpc_client/src/lib.rsgrpc_client/src/sglang_scheduler.rsgrpc_client/src/trtllm_service.rsgrpc_client/src/vllm_engine.rsmodel_gateway/src/routers/grpc/client.rs
Add impl_subscribe_kv_events!() macro in grpc_client/src/lib.rs following the existing impl_get_tokenizer!() pattern. The macro provides a shared subscribe_kv_events() method that returns tonic::Streaming<KvEventBatch> for long-lived event streams. The macro is invoked in all three backend clients: - grpc_client/src/sglang_scheduler.rs - grpc_client/src/vllm_engine.rs - grpc_client/src/trtllm_service.rs Add unified subscribe_kv_events() to GrpcClient enum in model_gateway/src/routers/grpc/client.rs for backend-agnostic event subscription. Refs: DESIGN.md Section 7 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
1545371 to
1fb7cad
Compare
…uting Add a positional indexer that uses DashMap<(position, ContentHash), SeqEntry> for O(1) random access to any depth position, replacing tree pointer-chasing with direct positional lookup. Jump search skips positions in strides of jump_size (default 64), yielding amortized O(D/J + W) matching complexity. What changed: - kv_index/Cargo.toml: add rustc-hash (FxHash) and xxhash-rust (XXH3) deps - kv_index/src/event_tree.rs: new 1250-line module with PositionalIndexer - kv_index/src/lib.rs: export new types (PositionalIndexer, ContentHash, SequenceHash, StoredBlock, OverlapScores, WorkerId, compute_content_hash) Key design decisions: - Dual-hash scheme: ContentHash (XXH3-64, position-independent, from token IDs) for indexing; SequenceHash (position-aware, from backend proto block_hash) for disambiguation - SeqEntry enum with Single/Multi optimization avoids HashMap allocation in the common single-sequence-hash case - Per-worker LevelIndex reverse lookup enables O(1) block removal - FxHashMap/FxHashSet throughout for 3-5x faster hashing on non-adversarial data - DashMap entries cleaned up when last worker is removed (no memory leak) - TOCTOU-safe worker lookup uses graceful fallback instead of unwrap - tracing::warn on unresolvable parent hash, tracing::debug on untracked worker removal - All methods take &self with internal DashMap + parking_lot::RwLock for thread safety Public API: - apply_stored(worker, blocks, parent_seq_hash): process store events - apply_removed(worker, seq_hashes): process remove events - apply_cleared(worker): process cache-clear events - remove_worker(worker): full worker removal - find_matches(content_hashes) -> OverlapScores: jump-search matching - compute_content_hash(token_ids) -> ContentHash: XXH3 hashing - current_size() -> usize: total blocks across all workers Tests: 43 tests covering store/match/remove/clear operations, jump search with various configurations (jump_size=1, 3, 4, 64), concurrent read+write, DashMap cleanup verification, hash computation, and edge cases. Refs: #557, #558 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
| tonic::Streaming<$crate::common_proto::KvEventBatch>, | ||
| Box<dyn std::error::Error + Send + Sync>, | ||
| > { | ||
| let request = tonic::Request::new($crate::common_proto::SubscribeKvEventsRequest { |
There was a problem hiding this comment.
I wish rust have a way to check this kind of imports. It technically follows the current setting but still a bit ugly.
Summary
impl_subscribe_kv_events!()shared macro for KV event subscription across all three backend clientssubscribe_kv_events()toGrpcClientenum for backend-agnostic usageWhat changed
grpc_client/src/lib.rsimpl_subscribe_kv_events!()macro (followsimpl_get_tokenizer!()pattern)grpc_client/src/sglang_scheduler.rssubscribe_kv_events()methodgrpc_client/src/vllm_engine.rssubscribe_kv_events()methodgrpc_client/src/trtllm_service.rssubscribe_kv_events()methodmodel_gateway/src/routers/grpc/client.rssubscribe_kv_events()onGrpcClientenumWhy
The gateway needs to subscribe to real-time KV cache events from backends for event-driven cache-aware routing. All three backends expose the same
SubscribeKvEventsRPC (added in #557), so a shared macro eliminates duplication — same pattern as the existingimpl_get_tokenizer!().How
lib.rsgenerates asubscribe_kv_events(start_sequence_number) -> Result<Streaming<KvEventBatch>>methodimplblockGrpcClient::subscribe_kv_events()dispatches to the correct backend variant$crate::common_proto::*types for the request/response (shared proto types)Test plan
cargo build -p smg-grpc-clientcompilescargo build -p smgcompilesRefs:
.claude/kv-event/DESIGN.mdSection 7Summary by CodeRabbit