Repository navigation
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:
📝 Walkthrough📝 Walkthrough🚥 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)
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
| syntax = "proto3"; | ||
|
|
||
| package tokenspeed.grpc.scheduler; | ||
|
|
||
| import "google/protobuf/timestamp.proto"; | ||
| import "google/protobuf/struct.proto"; |
There was a problem hiding this comment.
🔴 Important: This proto should import "common.proto" and reference the shared types for GetTokenizer and SubscribeKvEvents RPCs, like every other backend proto does.
Currently this file redefines GetTokenizerRequest, GetTokenizerChunk, SubscribeKvEventsRequest, KvEventBatch, KvCacheEvent, KvBlocksStored, KvBlock, KvBlocksRemoved, and KvCacheCleared locally (lines 389–422). These are byte-identical copies of the types in common.proto. Without the import, the tonic-generated client expects tokenspeed.grpc.scheduler.* types, which means:
- The shared
crate::impl_get_tokenizer!()andcrate::impl_subscribe_kv_events!()macros can't be used (they passcommon_proto::*types). subscribe_kv_eventsreturnsStreaming<proto::KvEventBatch>instead ofStreaming<common_proto::KvEventBatch>, breaking type compatibility with any code that handles KV events generically across backends.
See sglang_scheduler.proto lines 7, 34, 37–38 for the pattern:
import "common.proto";
// ...
rpc GetTokenizer(smg.grpc.common.GetTokenizerRequest) returns (stream smg.grpc.common.GetTokenizerChunk);
rpc SubscribeKvEvents(smg.grpc.common.SubscribeKvEventsRequest) returns (stream smg.grpc.common.KvEventBatch);Then the local type definitions for KV/tokenizer messages (lines 389–422) can be removed.
| } | ||
|
|
||
| /// Get load metrics. | ||
| pub async fn get_loads( | ||
| &self, | ||
| req: proto::GetLoadsRequest, | ||
| ) -> Result<proto::GetLoadsResponse, tonic::Status> { | ||
| let request = Request::new(req); | ||
| let mut client = self.client.clone(); | ||
| let response = client.get_loads(request).await?; | ||
| Ok(response.into_inner()) | ||
| } | ||
|
|
||
| /// Get tokenizer artifacts. | ||
| pub async fn get_tokenizer( | ||
| &self, | ||
| ) -> Result< | ||
| crate::tokenizer_bundle::StreamBundle, | ||
| Box<dyn std::error::Error + Send + Sync>, | ||
| > { | ||
| let request = Request::new(proto::GetTokenizerRequest {}); | ||
| let mut client = self.client.clone(); | ||
| crate::tokenizer_bundle::collect_bundle_from_rpc( | ||
| client.get_tokenizer(request), | ||
| |chunk| (chunk.data, chunk.sha256), | ||
| Duration::from_secs(120), | ||
| ) | ||
| .await | ||
| } | ||
|
|
||
| /// Subscribe to KV cache events. | ||
| pub async fn subscribe_kv_events( | ||
| &self, | ||
| start_sequence_number: u64, | ||
| ) -> Result<Streaming<proto::KvEventBatch>, tonic::Status> { | ||
| let request = Request::new(proto::SubscribeKvEventsRequest { | ||
| start_sequence_number, | ||
| }); | ||
| let mut client = self.client.clone(); | ||
| let response = client.subscribe_kv_events(request).await?; | ||
| Ok(response.into_inner()) | ||
| } | ||
| } |
There was a problem hiding this comment.
🔴 Important: get_tokenizer and subscribe_kv_events are manually reimplemented here instead of using the shared macros that all other backends use. Once the proto is fixed to import common.proto (see proto comment), replace these two methods with:
crate::impl_get_tokenizer!();
crate::impl_subscribe_kv_events!();The current subscribe_kv_events returns Streaming<proto::KvEventBatch> (a TokenSpeed-local type) instead of Streaming<common_proto::KvEventBatch>. This will cause a compile error when this client is wired into the routing layer, which handles KV events generically across backends using the common type.
| let mut client = self.client.clone(); | ||
| let response = client.health_check(request).await?; | ||
| Ok(response.into_inner()) | ||
| } | ||
|
|
||
| /// Abort a request. | ||
| pub async fn abort_request( | ||
| &self, | ||
| request_id: String, | ||
| reason: String, | ||
| ) -> Result<proto::AbortResponse, tonic::Status> { | ||
| let request = Request::new(proto::AbortRequest { | ||
| request_id, | ||
| reason, |
There was a problem hiding this comment.
🟡 Nit: abort_request returns Result<proto::AbortResponse, tonic::Status>, while all other backends (SglangSchedulerClient, VllmEngineClient) return Result<(), tonic::Status>. This works for the AbortOnDropStream caller (which only checks is_err()), but will be an API inconsistency when integrating with generic backend dispatch code.
Consider matching the existing pattern:
| let mut client = self.client.clone(); | |
| let response = client.health_check(request).await?; | |
| Ok(response.into_inner()) | |
| } | |
| /// Abort a request. | |
| pub async fn abort_request( | |
| &self, | |
| request_id: String, | |
| reason: String, | |
| ) -> Result<proto::AbortResponse, tonic::Status> { | |
| let request = Request::new(proto::AbortRequest { | |
| request_id, | |
| reason, | |
| pub async fn abort_request( | |
| &self, | |
| request_id: String, | |
| reason: String, | |
| ) -> Result<(), tonic::Status> { | |
| let request = Request::new(proto::AbortRequest { | |
| request_id, | |
| reason, | |
| }); | |
| let mut client = self.client.clone(); | |
| let _response = client.abort(request).await?; | |
| Ok(()) | |
| } |
There was a problem hiding this comment.
Code Review
This pull request introduces the TokenSpeedScheduler gRPC client and its corresponding protobuf definitions. The implementation includes a new AbortOnDropStream wrapper to ensure that inference requests are aborted if the client stream is dropped before completion. A review comment suggests optimizing poll_next to mark the stream as completed upon natural termination or error, preventing redundant abort calls during drop.
| type Item = Result<proto::GenerateResponse, tonic::Status>; | ||
|
|
||
| fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> { | ||
| Pin::new(&mut self.inner).poll_next(cx) |
There was a problem hiding this comment.
The current implementation of poll_next does not mark the stream as completed when it finishes naturally or encounters an error. This results in a redundant Abort RPC call being triggered in the Drop implementation for every request, even those that completed successfully. Marking the stream as completed upon termination avoids unnecessary network calls and potential server-side errors for aborting already-finished requests.
let res = Pin::new(&mut self.inner).poll_next(cx);
if matches!(res, Poll::Ready(None) | Poll::Ready(Some(Err(_)))) {
self.aborted.store(true, Ordering::Release);
}
resThere was a problem hiding this comment.
Actionable comments posted: 4
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@crates/grpc_client/proto/tokenspeed_scheduler.proto`:
- Around line 224-233: EmbedRequest has missing field numbers 3 and 6 which
should be explicitly reserved to avoid accidental reuse; update the EmbedRequest
message definition to add a reserved declaration for those field numbers (e.g.,
reserve 3 and 6) so it mirrors the safety used in EmbedResponse and prevents
future collisions when adding fields.
- Around line 93-112: GenerateRequest intentionally skips field number 15 (used
by SGLang's scheduler for lora_id); add a reserved 15; declaration inside the
GenerateRequest message to document and protect this intentional gap. Update the
GenerateRequest message definition to include reserved 15; alongside the
existing fields (e.g., near input_embeds = 14 and data_parallel_rank = 16) so
the protobuf prevents accidental reuse of field number 15.
In `@crates/grpc_client/src/tokenspeed_scheduler.rs`:
- Around line 198-251: The other RPC methods (health_check, abort_request,
get_model_info, get_server_info, get_loads) are missing trace context injection
like generate() and embed(); update each of these methods to create the Request
as before but then inject the current trace/span context into the outbound gRPC
metadata (the same mechanism used in generate() and embed(), e.g., call the
trace injection helper that writes into request.metadata_mut()) before calling
client.<method>(), ensuring you reference the existing functions health_check,
abort_request, get_model_info, get_server_info, and get_loads so each forwards
trace headers in the same way as generate/embed.
- Around line 126-130: The code currently converts endpoints with
endpoint.strip_prefix("grpc://") to "http://{addr}" via the http_endpoint
variable; update this logic to also detect "grpcs://" and convert it to
"https://{addr}" (ensure you check for "grpcs://" before "grpc://" or handle
both branches explicitly). Apply the same change to the other gRPC client
implementations that use the same pattern so TLS endpoints are consistently
mapped (look for the http_endpoint construction using
endpoint.strip_prefix("grpc://") in tokenspeed_scheduler.rs, vllm_engine.rs,
sglang_scheduler.rs, mlx_engine.rs, and trtllm_service.rs).
🪄 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: 136eca63-c481-4f71-8491-4947df276565
📒 Files selected for processing (4)
crates/grpc_client/build.rscrates/grpc_client/proto/tokenspeed_scheduler.protocrates/grpc_client/src/lib.rscrates/grpc_client/src/tokenspeed_scheduler.rs
| let http_endpoint = if let Some(addr) = endpoint.strip_prefix("grpc://") { | ||
| format!("http://{addr}") | ||
| } else { | ||
| endpoint.to_string() | ||
| }; |
There was a problem hiding this comment.
🧩 Analysis chain
🏁 Script executed:
#!/bin/bash
# Check how other backend clients handle endpoint schemes
rg -n 'strip_prefix.*grpc' crates/grpc_client/src/ --type rust -A 3Repository: lightseekorg/smg
Length of output: 1825
🏁 Script executed:
#!/bin/bash
# Search for grpcs usage anywhere in codebase
echo "=== Searching for grpcs references ==="
rg -n 'grpcs' . --type rust -i
echo -e "\n=== Searching for https endpoint references ==="
rg -n 'https://' crates/grpc_client/src/ --type rust
echo -e "\n=== Searching for TLS-related comments or TODOs ==="
rg -n 'TLS|tls|secure|https' crates/grpc_client/src/ --type rust | head -20
echo -e "\n=== Checking test files for endpoint scheme handling ==="
find crates/grpc_client -name '*.rs' -type f | xargs grep -l 'strip_prefix\|endpoint' | grep -i testRepository: lightseekorg/smg
Length of output: 675
🏁 Script executed:
#!/bin/bash
# Check what model_gateway does with grpcs
echo "=== model_gateway handling of grpcs ==="
sed -n '780,795p' ./model_gateway/src/main.rs
echo -e "\n=== Python bindings handling of grpcs ==="
sed -n '467,485p' ./bindings/python/src/lib.rs
# Also check if there's any conversion logic for grpcs elsewhere
echo -e "\n=== Searching for any https conversion or tls handling ==="
rg -n 'https|grpcs' ./model_gateway/src/main.rs -A 5 -B 5 | head -40Repository: lightseekorg/smg
Length of output: 1608
🏁 Script executed:
#!/bin/bash
# Check how http_endpoint is used in tokenspeed_scheduler
echo "=== How http_endpoint is used in tokenspeed_scheduler ==="
sed -n '126,145p' crates/grpc_client/src/tokenspeed_scheduler.rs
# Check what HTTP client library is being used
echo -e "\n=== HTTP client imports in tokenspeed_scheduler ==="
head -30 crates/grpc_client/src/tokenspeed_scheduler.rs
# Search for actual HTTP calls using the endpoint
echo -e "\n=== Where http_endpoint is used in the file ==="
rg -n 'http_endpoint' crates/grpc_client/src/tokenspeed_scheduler.rs -A 2Repository: lightseekorg/smg
Length of output: 2321
🏁 Script executed:
#!/bin/bash
# Check Cargo.toml for tonic version
echo "=== tonic dependency version ==="
grep -n 'tonic' Crates/grpc_client/Cargo.toml 2>/dev/null || grep -n 'tonic' crates/grpc_client/Cargo.toml
# Check if there are any tests for endpoint handling
echo -e "\n=== Test files in grpc_client ==="
find crates/grpc_client -name '*test*' -o -name '*tests*' | head -10
# Search for any documentation or comments about TLS/grpcs
echo -e "\n=== TLS/grpcs comments in grpc_client ==="
rg -n 'tls|TLS|grpc|grpcs' crates/grpc_client/src/ --type rust | grep -i 'comment\|todo\|note\|doc\|//\|/\*' | head -20
# Check if grpcs:// is used as input anywhere
echo -e "\n=== References to grpcs in tests or configs ==="
rg -n 'grpcs://' . --type rust --type tomlRepository: lightseekorg/smg
Length of output: 1915
Handle grpcs:// → https:// scheme conversion for TLS endpoints.
All gRPC client implementations intentionally convert the custom grpc:// scheme to http:// for tonic (see vllm_engine.rs:145 comment). However, none handle grpcs:// → https:// for TLS connections, despite other parts of the codebase recognizing grpcs:// (e.g., model_gateway/src/main.rs:786, bindings/python/src/lib.rs:473).
For consistency, add TLS support:
Suggested change
let http_endpoint = if let Some(addr) = endpoint.strip_prefix("grpc://") {
format!("http://{addr}")
+ } else if let Some(addr) = endpoint.strip_prefix("grpcs://") {
+ format!("https://{addr}")
} else {
endpoint.to_string()
};This gap exists in all five gRPC client implementations (tokenspeed_scheduler.rs, vllm_engine.rs, sglang_scheduler.rs, mlx_engine.rs, trtllm_service.rs).
📝 Committable suggestion
‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.
| let http_endpoint = if let Some(addr) = endpoint.strip_prefix("grpc://") { | |
| format!("http://{addr}") | |
| } else { | |
| endpoint.to_string() | |
| }; | |
| let http_endpoint = if let Some(addr) = endpoint.strip_prefix("grpc://") { | |
| format!("http://{addr}") | |
| } else if let Some(addr) = endpoint.strip_prefix("grpcs://") { | |
| format!("https://{addr}") | |
| } else { | |
| endpoint.to_string() | |
| }; |
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@crates/grpc_client/src/tokenspeed_scheduler.rs` around lines 126 - 130, The
code currently converts endpoints with endpoint.strip_prefix("grpc://") to
"http://{addr}" via the http_endpoint variable; update this logic to also detect
"grpcs://" and convert it to "https://{addr}" (ensure you check for "grpcs://"
before "grpc://" or handle both branches explicitly). Apply the same change to
the other gRPC client implementations that use the same pattern so TLS endpoints
are consistently mapped (look for the http_endpoint construction using
endpoint.strip_prefix("grpc://") in tokenspeed_scheduler.rs, vllm_engine.rs,
sglang_scheduler.rs, mlx_engine.rs, and trtllm_service.rs).
|
🔴 Important — Security: Committed key material in This PR adds a 32-byte binary file at Details:
Action needed:
|
There was a problem hiding this comment.
Actionable comments posted: 2
♻️ Duplicate comments (1)
crates/grpc_client/src/tokenspeed_scheduler.rs (1)
139-143: 🧹 Nitpick | 🔵 TrivialHandle
grpcs://→https://scheme conversion for TLS endpoints.The endpoint normalization only handles
grpc://→http://. For TLS support, also handlegrpcs://:let http_endpoint = if let Some(addr) = endpoint.strip_prefix("grpc://") { format!("http://{addr}") + } else if let Some(addr) = endpoint.strip_prefix("grpcs://") { + format!("https://{addr}") } else { endpoint.to_string() };This gap exists in all gRPC client implementations and should be addressed uniformly.
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@crates/grpc_client/src/tokenspeed_scheduler.rs` around lines 139 - 143, The endpoint normalization in tokenspeed_scheduler.rs currently only maps "grpc://" to "http://" for the local variable http_endpoint; update that logic to also map "grpcs://" to "https://" (i.e., detect strip_prefix("grpcs://") and format!("https://{addr}") ) so TLS gRPC endpoints are correctly converted; apply the same change across all gRPC client implementations that perform this normalization to ensure consistent handling of secure endpoints.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@crates/grpc_client/proto/tokenspeed_scheduler.proto`:
- Around line 317-333: The SchedulerLoad protobuf message is missing a reserved
declaration for the skipped field number 14; add a reserved 14; statement inside
the SchedulerLoad message (between the speculative = 13 and disaggregation = 15
field entries) so the gap is explicitly documented and prevents future reuse of
field number 14.
In `@docs/PROJECT_ANALYSIS_CN.md`:
- Around line 37-75: The markdown has multiple fenced code blocks missing
language specifiers; add appropriate specifiers (e.g., use ```text for
conceptual flow blocks and leave pure ASCII-art diagrams without a specifier)
and insert a blank line immediately before the fenced block labeled "MCP 循环流程"
so the renderer treats it as a separate code block; specifically, update the
fenced blocks that represent conceptual flows (the ones noted around the MCP
循环流程 and the two conceptual-flow blocks referenced in the review) to use a
language tag like text, ensure a blank line precedes the MCP 循环流程 fence, and do
not modify the pure ASCII art diagram fences referenced in the review.
---
Duplicate comments:
In `@crates/grpc_client/src/tokenspeed_scheduler.rs`:
- Around line 139-143: The endpoint normalization in tokenspeed_scheduler.rs
currently only maps "grpc://" to "http://" for the local variable http_endpoint;
update that logic to also map "grpcs://" to "https://" (i.e., detect
strip_prefix("grpcs://") and format!("https://{addr}") ) so TLS gRPC endpoints
are correctly converted; apply the same change across all gRPC client
implementations that perform this normalization to ensure consistent handling of
secure endpoints.
🪄 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: f6e68149-472f-47e2-8fbb-1b1d4de47f4a
📒 Files selected for processing (4)
crates/grpc_client/proto/tokenspeed_scheduler.protocrates/grpc_client/src/tokenspeed_scheduler.rsdocs/PROJECT_ANALYSIS_CN.mdlogs/security/.security-key
| ``` | ||
| ┌─────────────────────────────────────────────────────────────────┐ | ||
| │ Client Applications │ | ||
| │ (OpenAI SDK / REST / gRPC / WebSocket) │ | ||
| └────────────────────────────┬────────────────────────────────────┘ | ||
| │ | ||
| ▼ | ||
| ┌─────────────────────────────────────────────────────────────────┐ | ||
| │ SMG Gateway Layer │ | ||
| │ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐ │ | ||
| │ │Rate Limit│ │OIDC/JWT │ │ WASM │ │ OTel │ │ | ||
| │ │ │ │Auth │ │ Plugins │ │ Tracing │ │ | ||
| │ └──────────┘ └──────────┘ └──────────┘ └──────────┘ │ | ||
| ├─────────────────────────────────────────────────────────────────┤ | ||
| │ Router Layer │ | ||
| │ ┌────────────────┐ ┌───────────────┐ ┌───────────────────┐ │ | ||
| │ │ gRPC Path │ │ HTTP Path │ │ 3rd Party Path │ │ | ||
| │ │ (Full Server) │ │ (Smart Proxy)│ │ (Unified Router) │ │ | ||
| │ └────────┬───────┘ └───────┬───────┘ └─────────┬─────────┘ │ | ||
| │ │ │ │ │ | ||
| │ ┌────────┴─────────────────┴────────────────────┴──────────┐ │ | ||
| │ │ Load Balancing Policy Engine │ │ | ||
| │ │ cache_aware | bucket | power_of_two | consistent_hashing │ │ | ||
| │ │ prefix_hash | manual | round_robin | random │ │ | ||
| │ └───────────────────────────────────────────────────────────┘ │ | ||
| ├─────────────────────────────────────────────────────────────────┤ | ||
| │ Resilience Layer │ | ||
| │ ┌──────────────┐ ┌───────┐ ┌────────────┐ ┌────────────────┐ │ | ||
| │ │Circuit Breaker│ │ Retry │ │Health Check│ │ Token Bucket │ │ | ||
| │ └──────────────┘ └───────┘ └────────────┘ └────────────────┘ │ | ||
| └────────────────────────────┬────────────────────────────────────┘ | ||
| │ | ||
| ┌──────────────┼──────────────┐ | ||
| ▼ ▼ ▼ | ||
| ┌──────────┐ ┌──────────┐ ┌──────────────┐ | ||
| │ vLLM │ │ SGLang │ │ TensorRT-LLM │ | ||
| │(gRPC/HTTP)│ │(gRPC/HTTP)│ │ (gRPC/HTTP) │ | ||
| └──────────┘ └──────────┘ └──────────────┘ | ||
| ``` |
There was a problem hiding this comment.
🧹 Nitpick | 🔵 Trivial
Add language specifiers to fenced code blocks where applicable.
Static analysis indicates multiple fenced code blocks are missing language specifiers. While ASCII diagrams can remain unspecified, consider adding appropriate specifiers for improved readability:
- Line 287: Add blank line before the fenced code block (MCP 循环流程)
- Lines 554, 574: These represent conceptual flows and could use
textor remain as-is
For pure ASCII art diagrams (lines 37, 87, 102, 176, 240, 261, 419, 437, 451), leaving them without a language specifier is acceptable since they're visual representations rather than code.
🧰 Tools
🪛 markdownlint-cli2 (0.22.0)
[warning] 37-37: Fenced code blocks should have a language specified
(MD040, fenced-code-language)
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@docs/PROJECT_ANALYSIS_CN.md` around lines 37 - 75, The markdown has multiple
fenced code blocks missing language specifiers; add appropriate specifiers
(e.g., use ```text for conceptual flow blocks and leave pure ASCII-art diagrams
without a specifier) and insert a blank line immediately before the fenced block
labeled "MCP 循环流程" so the renderer treats it as a separate code block;
specifically, update the fenced blocks that represent conceptual flows (the ones
noted around the MCP 循环流程 and the two conceptual-flow blocks referenced in the
review) to use a language tag like text, ensure a blank line precedes the MCP
循环流程 fence, and do not modify the pure ASCII art diagram fences referenced in
the review.
84c4138 to
0682c08
Compare
Add gRPC client support for the TokenSpeed inference engine, alongside existing SGLang, vLLM, TRT-LLM, and MLX clients. New files: - proto/tokenspeed_scheduler.proto: TokenSpeedScheduler service (9 RPCs) with TokenSpeed-specific fields (spec_verify_ct, accept_draft_tokens). Uses common.proto for shared GetTokenizer/KvEvents types. - src/tokenspeed_scheduler.rs: TokenSpeedSchedulerClient with AbortOnDropStream (auto-abort on drop with mark_completed), trace injection, and all RPC wrappers. Uses impl_get_tokenizer!() and impl_subscribe_kv_events!() shared macros. Modified files: - build.rs: compile tokenspeed_scheduler.proto - src/lib.rs: export TokenSpeedSchedulerClient + tokenspeed_proto Tests (13 new, 59 total): - Proto type compilation and construction - GenerateRequest/Response with sampling params + oneof variants - TokenSpeed-specific fields (spec_verify_ct, accept_draft_tokens) - matched_stop oneof (string + token_id) - SamplingParams proto3 defaults + constraint oneof - EmbedRequest, AbortRequest, GetModelInfoResponse - LogProbs structure, DisaggregatedParams - Client error path (invalid endpoint) Related: lightseekorg/tokenspeed#120, lightseekorg/tokenspeed#185 Signed-off-by: yetone <yetoneful@gmail.com>
34ec451 to
27ddce7
Compare
Addresses review feedback on #1167: - Inject trace metadata on the remaining unary RPCs (health_check, abort_request, get_model_info, get_server_info, get_loads). Factor the injector call into a small inject_trace helper so every RPC shares the same behavior instead of conditionally forwarding in two places. - Accept grpcs:// as a TLS-enabled alias for the endpoint scheme. The previous mapping only covered grpc:// → http://, so TLS endpoints passed via grpcs:// silently fell through to the raw string and failed.
|
Hi @yetone, the DCO sign-off check has failed. All commits must include a To fix existing commits: # Sign off the last N commits (replace N with the number of unsigned commits)
git rebase HEAD~N --signoff
git push --force-with-leaseTo sign off future commits automatically:
|
| // Accept `grpc://` (→ http) and `grpcs://` (→ https) prefixes so TLS | ||
| // endpoints work without the caller having to rewrite the scheme. | ||
| let http_endpoint = if let Some(addr) = endpoint.strip_prefix("grpcs://") { | ||
| format!("https://{addr}") |
There was a problem hiding this comment.
🟡 Nit: The grpcs:// → https:// conversion requires tonic's TLS support, but grpc_client's Cargo.toml only enables ["gzip", "transport"] for tonic (no tls-ring). This works today because model_gateway also depends on smg-mesh, which enables tls-ring — Cargo feature unification means tonic gets TLS in the final binary.
However, this is an implicit dependency. If grpc_client is ever used in a binary without mesh, grpcs:// endpoints will fail at runtime with no compile-time warning. Consider adding tls-ring explicitly:
# crates/grpc_client/Cargo.toml
tonic = { workspace = true, features = ["tls-ring"] }This matches what crates/mesh/Cargo.toml already does.
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 `@crates/grpc_client/src/tokenspeed_scheduler.rs`:
- Around line 500-504: The test test_client_connect_invalid_endpoint currently
only asserts a generic error; change it to assert the error kind/text is the
expected invalid-endpoint failure by calling result.unwrap_err() and checking
its message or variant (e.g., assert!(err.to_string().contains("invalid") ||
matches!(err, tonic::transport::Error::InvalidUri(_))), so the test specifically
verifies TokenSpeedSchedulerClient::connect("invalid://endpoint") fails due to
an invalid URI/scheme rather than any other transport/runtime error.
🪄 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: db285db5-010b-474e-adc1-e713c42facd0
📒 Files selected for processing (4)
crates/grpc_client/build.rscrates/grpc_client/proto/tokenspeed_scheduler.protocrates/grpc_client/src/lib.rscrates/grpc_client/src/tokenspeed_scheduler.rs
| #[tokio::test] | ||
| async fn test_client_connect_invalid_endpoint() { | ||
| let result = TokenSpeedSchedulerClient::connect("invalid://endpoint").await; | ||
| assert!(result.is_err()); | ||
| } |
There was a problem hiding this comment.
🧹 Nitpick | 🔵 Trivial
Minor: consider an explicit error kind on the invalid-endpoint test.
assert!(result.is_err()) passes as long as connect("invalid://endpoint") fails for any reason, including an unrelated change in Channel::from_shared/runtime behavior later. Without a connect_timeout configured (intentionally deferred per the shared backend pattern), this test can also be slow-to-fail in CI if the URI ever parses as a hostname that attempts a real DNS lookup.
Not blocking — just something to tighten in a follow-up once a uniform connect_timeout lands across the backend clients.
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@crates/grpc_client/src/tokenspeed_scheduler.rs` around lines 500 - 504, The
test test_client_connect_invalid_endpoint currently only asserts a generic
error; change it to assert the error kind/text is the expected invalid-endpoint
failure by calling result.unwrap_err() and checking its message or variant
(e.g., assert!(err.to_string().contains("invalid") || matches!(err,
tonic::transport::Error::InvalidUri(_))), so the test specifically verifies
TokenSpeedSchedulerClient::connect("invalid://endpoint") fails due to an invalid
URI/scheme rather than any other transport/runtime error.
|
Hi @yetone, the DCO sign-off check has failed. All commits must include a To fix existing commits: # Sign off the last N commits (replace N with the number of unsigned commits)
git rebase HEAD~N --signoff
git push --force-with-leaseTo sign off future commits automatically:
|
| .ok() | ||
| .and_then(|v| { | ||
| v.split(',') | ||
| .map(|s| s.trim().parse::<u32>().ok()) | ||
| .collect::<Option<Vec<_>>>() | ||
| }) | ||
| .unwrap_or_else(|| vec![1, 2, 3]) | ||
| } | ||
|
|
||
| fn sampling_params(max_new_tokens: u32) -> proto::SamplingParams { |
There was a problem hiding this comment.
🟡 Nit: If TOKENSPEED_TEST_INPUT_IDS is set but malformed (e.g. "abc" or "1,,2"), this silently falls back to vec![1, 2, 3]. In CI, that means a misconfigured env var looks like a passing test with default IDs — exactly the "silent fallback to default when validation should fail loudly" pattern.
Consider panicking when the var is set but unparseable, matching how endpoint() treats its env var:
| .ok() | |
| .and_then(|v| { | |
| v.split(',') | |
| .map(|s| s.trim().parse::<u32>().ok()) | |
| .collect::<Option<Vec<_>>>() | |
| }) | |
| .unwrap_or_else(|| vec![1, 2, 3]) | |
| } | |
| fn sampling_params(max_new_tokens: u32) -> proto::SamplingParams { | |
| fn input_ids() -> Vec<u32> { | |
| match std::env::var("TOKENSPEED_TEST_INPUT_IDS") { | |
| Ok(v) => v | |
| .split(',') | |
| .map(|s| { | |
| s.trim() | |
| .parse::<u32>() | |
| .unwrap_or_else(|e| panic!("bad TOKENSPEED_TEST_INPUT_IDS token {s:?}: {e}")) | |
| }) | |
| .collect(), | |
| Err(_) => vec![1, 2, 3], | |
| } | |
| } |
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 `@crates/grpc_client/tests/tokenspeed_scheduler_integration.rs`:
- Around line 154-162: The test get_loads_returns_empty_stub contains a fragile
assertion assert_eq!(resp.loads.len(), 0) tied to the current server stub;
change it to only assert the RPC succeeded and validate any returned entries are
well-formed so the test survives when the server returns real data: call
client.get_loads(proto::GetLoadsRequest::default()), ensure the response (resp)
is Ok, then iterate resp.loads and perform light structural checks (e.g.,
required fields are present/non-empty or within expected ranges) instead of
asserting the length equals zero.
🪄 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: b3bfeb31-264a-4174-9903-32976a5f114a
📒 Files selected for processing (3)
crates/grpc_client/src/lib.rscrates/grpc_client/src/tokenspeed_scheduler.rscrates/grpc_client/tests/tokenspeed_scheduler_integration.rs
There was a problem hiding this comment.
Actionable comments posted: 2
♻️ Duplicate comments (1)
crates/grpc_client/tests/tokenspeed_scheduler_integration.rs (1)
162-170:⚠️ Potential issue | 🟡 MinorBrittle assertion tied to current server stub.
assert_eq!(resp.loads.len(), 0)will start failing as soon as the companion TokenSpeed server returns realGetLoadsentries — the inline comment already flags this as stub behavior. Prefer asserting the RPC succeeds (and optionally that each returned entry is well-formed) so this test survives the stub being replaced.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@crates/grpc_client/tests/tokenspeed_scheduler_integration.rs` around lines 162 - 170, The test get_loads_returns_empty_stub contains a brittle assertion assert_eq!(resp.loads.len(), 0) tied to the server stub; change it to only assert the RPC succeeded (i.e., keep the .await.expect call) and replace the hardcoded zero-length check on resp.loads with a more resilient verification: either no assertion on length or assert that every returned load in resp.loads is well-formed (e.g., required fields are present/non-empty or within valid ranges) by iterating resp.loads and validating fields on each entry; update the assertions to reference resp and its loads so the test survives when the server returns real entries.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@crates/grpc_client/tests/tokenspeed_scheduler_integration.rs`:
- Around line 153-158: The uptime assertion in test
get_server_info_reports_uptime_and_type is too weak; replace the
assert!(resp.uptime_seconds >= 0.0) with a stricter check such as
assert!(resp.uptime_seconds > 0.0) (or remove it entirely) to ensure the server
has been up long enough to respond; update the assertion referencing
resp.uptime_seconds accordingly in that test function.
- Around line 255-291: The test currently calls client.abort_request(...) and
then drop(stream), which triggers AbortOnDropStream::drop to spawn a second
Abort RPC because mark_completed() isn't set; to avoid the fire-and-forget spawn
race, drain the stream to completion (e.g. loop awaiting stream.next() until it
returns None) or call the existing helper that sets mark_completed() on the
AbortOnDropStream before dropping; locate the test function
abort_request_out_of_band_returns_ok, the Generate stream obtained from
client.generate(...) and ensure you consume the stream to None (or invoke
mark_completed) after abort_request and before drop(stream) so only the explicit
abort RPC is sent.
---
Duplicate comments:
In `@crates/grpc_client/tests/tokenspeed_scheduler_integration.rs`:
- Around line 162-170: The test get_loads_returns_empty_stub contains a brittle
assertion assert_eq!(resp.loads.len(), 0) tied to the server stub; change it to
only assert the RPC succeeded (i.e., keep the .await.expect call) and replace
the hardcoded zero-length check on resp.loads with a more resilient
verification: either no assertion on length or assert that every returned load
in resp.loads is well-formed (e.g., required fields are present/non-empty or
within valid ranges) by iterating resp.loads and validating fields on each
entry; update the assertions to reference resp and its loads so the test
survives when the server returns real entries.
🪄 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: f26a1319-355e-4a18-9851-712809a4991e
📒 Files selected for processing (1)
crates/grpc_client/tests/tokenspeed_scheduler_integration.rs
Addresses review feedback on #1167: - Inject trace metadata on the remaining unary RPCs (health_check, abort_request, get_model_info, get_server_info, get_loads). Factor the injector call into a small inject_trace helper so every RPC shares the same behavior instead of conditionally forwarding in two places. - Accept grpcs:// as a TLS-enabled alias for the endpoint scheme. The previous mapping only covered grpc:// → http://, so TLS endpoints passed via grpcs:// silently fell through to the raw string and failed. Signed-off-by: yetone <yetoneful@gmail.com>
Signed-off-by: yetone <yetoneful@gmail.com>
…erClient The unit tests in src/tokenspeed_scheduler.rs only construct proto types; they never exercised the actual RPC path (connection, streaming, abort, trace injection) against a real server. This PR's unit tests therefore gave zero coverage for grpc:// / grpcs:// prefix parsing, the per-RPC inject_trace helper, AbortOnDropStream behavior, or any wrapper method. Add a `tests/tokenspeed_scheduler_integration.rs` suite gated on the `TOKENSPEED_GRPC_ENDPOINT` env var (skipped via `#[ignore]` by default) that connects to a live TokenSpeed server and exercises each path: * health_check_reports_healthy * get_model_info_reports_populated_fields * get_server_info_reports_uptime_and_type * get_loads_returns_empty_stub * generate_non_streaming_returns_completion_with_tokens * generate_streaming_terminates_with_complete * abort_request_out_of_band_returns_ok * abort_on_drop_fires_and_client_stays_usable * trace_injector_fires_on_every_unary_rpc * trace_injector_fires_on_generate_and_embed The trace injector tests use a custom `TraceInjector` that counts every inject() call, giving a hard assertion that every RPC wrapper actually invokes the injector — direct evidence that the new inject_trace helper is wired on health_check / abort_request / get_model_info / get_server_info / get_loads as well as the original generate / embed paths. Verified end-to-end against `openai/gpt-oss-20b` running in this PR's companion server (lightseekorg/tokenspeed#185) on an NV B200: **10 / 10 tests passed**. Signed-off-by: yetone <yetoneful@gmail.com>
… not auto-detected by clippy) Signed-off-by: yetone <yetoneful@gmail.com>
Signed-off-by: yetone <yetoneful@gmail.com>
c2364f8 to
786313d
Compare
Four fixes, one per reviewer comment: * Cargo.toml: pin `tonic` with `tls-ring` feature. `grpcs://` → `https://` endpoint conversion relies on tonic TLS support, which today is only brought in transitively through `smg-mesh`'s feature set. Any binary that links `grpc_client` without `mesh` would currently fail at runtime. Make the dependency explicit. (claude[bot]) * input_ids(): panic on malformed `TOKENSPEED_TEST_INPUT_IDS` instead of silently falling back to the default. A misconfigured CI env var passing tests with wrong inputs is strictly worse than a loud failure. (claude[bot]) * get_server_info_reports_uptime_and_type: tighten `>= 0.0` → `> 0.0`; by the time this test runs the server has served at least one prior RPC, so any non-positive uptime signals a real bug. (coderabbitai) * get_loads_returns_empty_stub → get_loads_returns_shape: don't pin `loads.len() == 0`; that would break the day the server-side stub is replaced with real per-DP-rank metrics. Validate shape instead. (coderabbitai) * abort_request_out_of_band_returns_ok: call `stream.mark_completed()` before dropping after an explicit Abort so `AbortOnDropStream::Drop` doesn't spawn a second fire-and-forget Abort RPC that could outlive the test's runtime. (coderabbitai) Also add `clippy::panic` to the file-level expect list (same rationale as `clippy::expect_used`: `allow-panic-in-tests` doesn't detect `#[tokio::test]`). Signed-off-by: yetone <yetoneful@gmail.com>
Three additions that fill the holes the previous self-review called out: 1. **`grpcs://` scheme conversion is now unit-tested.** Extract the inline `grpc://` / `grpcs://` prefix stripper into a private `normalize_endpoint` fn and add unit tests covering every rewrite (`grpc://` → `http://`, `grpcs://` → `https://`, passthrough for `http://`/`https://`/garbage, path+query preservation). Previously a regression in the TLS mapping would have been invisible without a real TLS server. 2. **Shared `common_proto` types are compile-time asserted.** Two `async fn _...` stubs in the unit-test module return `Streaming<common_proto::KvEventBatch>` / `tokenizer_bundle::StreamBundle` from the client's wrappers. If the proto ever regresses to local `GetTokenizerChunk` / `KvEventBatch` types, these stop compiling — catching what the claude[bot] review was worried about. 3. **Trace headers are verified on the wire, not just on the injector.** New `tests/tokenspeed_scheduler_mock.rs` starts an in-process mock `TokenSpeedSchedulerServer` implementation and runs the Rust client against it. The mock records every RPC's inbound `MetadataMap`; the tests assert that a sentinel `x-tokenspeed-test-trace` header survives the wire on every unary RPC (`HealthCheck`, `GetModelInfo`, `GetServerInfo`, `GetLoads`, `Abort`) **and** on streaming RPCs (`Generate`, `Embed`). This closes the gap that a pure "injector-was-called" counter could not reach without modifying the real server. The mock also exercises: - `grpc://` end-to-end handshake (complementing #1's string-level test) - `subscribe_kv_events` round-trip typed as `common_proto::KvEventBatch` - `get_tokenizer` through the shared `StreamBundle` helper + sha256 All mock tests run in vanilla `cargo test` (no env var gating, no remote server), giving CI permanent coverage for paths that previously only the live integration suite could touch. Verified locally: 15 unit tests (was 13; +2 for normalize_endpoint) and 5 mock-server tests passing. Signed-off-by: yetone <yetoneful@gmail.com>
`constant_time_eq` 0.4.3 (published 2026-04-19) bumps its minimum
rustc to 1.95.0; the workspace toolchain is still 1.90. Because this
workspace does not commit `Cargo.lock`, every CI run resolves fresh
and started pulling 0.4.3 today, breaking both `unit-tests` and
`build-wheel` with:
error: rustc 1.90.0 is not supported by the following package:
constant_time_eq@0.4.3 requires rustc 1.95.0
Pin to `=0.4.2` here. The crate is a transitive dep of both `rustls`
and `zip`, both of which `grpc_client` now pulls in (rustls via the
newly-enabled `tls-ring` feature, zip via the tokenizer-bundle
helper), so pinning at this crate's root flows out to every
dependent in the workspace via cargo's unified resolver. Remove
once the workspace MSRV bumps to 1.95+.
Signed-off-by: yetone <yetoneful@gmail.com>
Wraps TokenSpeed's AsyncLLM behind the SGLang proto so the existing Rust router auto-detects and routes to TokenSpeed workers with zero gateway-side changes. Why the SGLang proto? TokenSpeed and SGLang share the tokenized-request shape, sampling- param naming, and output dict format, so reusing the proto avoids inventing a new one — the previous tokenspeed_scheduler.proto attempt in lightseekorg/tokenspeed#185 and #1167 was closed in favor of this approach. Module layout (grpc_servicer/smg_grpc_servicer/tokenspeed/): - servicer.py TokenSpeedSchedulerServicer — Generate, Embed, HealthCheck, Abort, GetModelInfo, GetServerInfo, GetTokenizer, GetLoads - health_servicer.py grpc.health.v1.Health bridge advertising the SGLang service name for router auto-detection - scheduler_launcher.py thin wrapper around TokenSpeed's _launch_subprocesses - server.py python -m entrypoint that boots the scheduler, AsyncLLM, grpc.aio server, and warmup probe - __main__.py CLI shim e2e_test infra: - Runtime.TOKENSPEED enum (infra/constants.py) - _build_tokenspeed_grpc_cmd so e2e fixtures can spawn tokenspeed workers via python -m smg_grpc_servicer.tokenspeed - Qwen/Qwen3-4B and Qwen/Qwen3-30B-A3B model specs for coverage - TestOpenAIServerFunctionCalling now runs against tokenspeed too; model temporarily switched to Qwen/Qwen3-30B-A3B + qwen parser because TokenSpeed's registry does not cover plain LlamaForCausalLM (only LlamaForCausalLMMoE / LlamaForCausalLMEagle3) Test evidence (single B200, Qwen/Qwen3-30B-A3B): - Unit: 47 / 47 PASSED (grpc_servicer/tests/) - E2E: tokenspeed 12 passed / 4 failed / 2 skipped vllm (ref) 10 passed / 6 failed / 2 skipped The 4 e2e failures are identical across backends (test_function_call_ required / specific × openai/smg clients). Root cause is upstream of the engine — in smg gateway's tool_choice=required/specific constraint translation path — not this integration. Closes #120 Signed-off-by: yetone <yetoneful@gmail.com>
Wraps TokenSpeed's AsyncLLM behind the SGLang proto so the existing Rust router auto-detects and routes to TokenSpeed workers with zero gateway-side changes. Why the SGLang proto? TokenSpeed and SGLang share the tokenized-request shape, sampling- param naming, and output dict format, so reusing the proto avoids inventing a new one — the previous tokenspeed_scheduler.proto attempt in lightseekorg/tokenspeed#185 and #1167 was closed in favor of this approach. Module layout (grpc_servicer/smg_grpc_servicer/tokenspeed/): - servicer.py TokenSpeedSchedulerServicer — Generate, Embed, HealthCheck, Abort, GetModelInfo, GetServerInfo, GetTokenizer, GetLoads - health_servicer.py grpc.health.v1.Health bridge advertising the SGLang service name for router auto-detection - scheduler_launcher.py thin wrapper around TokenSpeed's _launch_subprocesses - server.py python -m entrypoint that boots the scheduler, AsyncLLM, grpc.aio server, and warmup probe - __main__.py CLI shim e2e_test infra: - Runtime.TOKENSPEED enum (infra/constants.py) - _build_tokenspeed_grpc_cmd so e2e fixtures can spawn tokenspeed workers via python -m smg_grpc_servicer.tokenspeed - Qwen/Qwen3-4B and Qwen/Qwen3-30B-A3B model specs for coverage - TestOpenAIServerFunctionCalling now runs against tokenspeed too; model temporarily switched to Qwen/Qwen3-30B-A3B + qwen parser because TokenSpeed's registry does not cover plain LlamaForCausalLM (only LlamaForCausalLMMoE / LlamaForCausalLMEagle3) Test evidence (single B200, Qwen/Qwen3-30B-A3B): - Unit: 47 / 47 PASSED (grpc_servicer/tests/) - E2E: tokenspeed 12 passed / 4 failed / 2 skipped vllm (ref) 10 passed / 6 failed / 2 skipped The 4 e2e failures are identical across backends (test_function_call_ required / specific × openai/smg clients). Root cause is upstream of the engine — in smg gateway's tool_choice=required/specific constraint translation path — not this integration. Closes #120 Signed-off-by: yetone <yetoneful@gmail.com>
Wraps TokenSpeed's AsyncLLM behind the SGLang proto so the existing Rust router auto-detects and routes to TokenSpeed workers with zero gateway-side changes. Why the SGLang proto? TokenSpeed and SGLang share the tokenized-request shape, sampling- param naming, and output dict format, so reusing the proto avoids inventing a new one — the previous tokenspeed_scheduler.proto attempt in lightseekorg/tokenspeed#185 and #1167 was closed in favor of this approach. Module layout (grpc_servicer/smg_grpc_servicer/tokenspeed/): - servicer.py TokenSpeedSchedulerServicer — Generate, Embed, HealthCheck, Abort, GetModelInfo, GetServerInfo, GetTokenizer, GetLoads - health_servicer.py grpc.health.v1.Health bridge advertising the SGLang service name for router auto-detection - scheduler_launcher.py thin wrapper around TokenSpeed's _launch_subprocesses - server.py python -m entrypoint that boots the scheduler, AsyncLLM, grpc.aio server, and warmup probe - __main__.py CLI shim e2e_test infra: - Runtime.TOKENSPEED enum (infra/constants.py) - _build_tokenspeed_grpc_cmd so e2e fixtures can spawn tokenspeed workers via python -m smg_grpc_servicer.tokenspeed - Qwen/Qwen3-4B and Qwen/Qwen3-30B-A3B model specs for coverage - TestOpenAIServerFunctionCalling keeps its original Llama-3.2-1B + llama parser coverage for sglang/vllm/trtllm; a new TestOpenAIServerFunctionCallingTokenSpeed subclass runs the same test bodies against Qwen/Qwen3-4B + qwen parser for tokenspeed, since TokenSpeed's model registry does not cover plain LlamaForCausalLM (only LlamaForCausalLMMoE / LlamaForCausalLMEagle3). Test evidence (single B200, Qwen/Qwen3-30B-A3B): - Unit: 47 / 47 PASSED (grpc_servicer/tests/) - E2E: tokenspeed 12 passed / 4 failed / 2 skipped vllm (ref) 10 passed / 6 failed / 2 skipped The 4 e2e failures are identical across backends (test_function_call_ required / specific × openai/smg clients). Root cause is upstream of the engine — in smg gateway's tool_choice=required/specific constraint translation path — not this integration. Closes #120 Signed-off-by: yetone <yetoneful@gmail.com>
Wraps TokenSpeed's AsyncLLM behind the SGLang proto so the existing Rust router auto-detects and routes to TokenSpeed workers with zero gateway-side changes. Why the SGLang proto? TokenSpeed and SGLang share the tokenized-request shape, sampling- param naming, and output dict format, so reusing the proto avoids inventing a new one — the previous tokenspeed_scheduler.proto attempt in lightseekorg/tokenspeed#185 and #1167 was closed in favor of this approach. Module layout (grpc_servicer/smg_grpc_servicer/tokenspeed/): - servicer.py TokenSpeedSchedulerServicer — Generate, Embed, HealthCheck, Abort, GetModelInfo, GetServerInfo, GetTokenizer, GetLoads - health_servicer.py grpc.health.v1.Health bridge advertising the SGLang service name for router auto-detection - scheduler_launcher.py thin wrapper around TokenSpeed's _launch_subprocesses - server.py python -m entrypoint that boots the scheduler, AsyncLLM, grpc.aio server, and warmup probe - __main__.py CLI shim e2e_test infra: - Runtime.TOKENSPEED enum (infra/constants.py) - _build_tokenspeed_grpc_cmd so e2e fixtures can spawn tokenspeed workers via python -m smg_grpc_servicer.tokenspeed - Qwen/Qwen3-4B and Qwen/Qwen3-30B-A3B model specs for coverage - TestOpenAIServerFunctionCalling now runs the shared Llama-3.2-1B- Instruct fixture against tokenspeed alongside sglang/vllm/trtllm. TokenSpeed support for dense ``LlamaForCausalLM`` is wired up via lightseekorg/tokenspeed#357 (a new ``tokenspeed.runtime.models.llama`` module registered in the model registry); ci_install_tokenspeed.sh pins to that branch until the upstream PR merges. Test evidence (single B200, shared Llama-3.2-1B-Instruct fixture): - Unit: 47 / 47 PASSED (grpc_servicer/tests/) - Registry: ``LlamaForCausalLM`` resolves to the new dense class in the ``tokenspeed:asyncllm`` image; weights load cleanly (``Load weight end. type=LlamaForCausalLM, dtype=torch.bfloat16``) - End-to-end generation smoke against the tokenspeed image surfaced an unrelated detokenizer/sampler issue in the ``tokenspeed:asyncllm`` build itself (reproduces on stock ``Qwen3-0.6B`` with both the ``greedy`` and ``flashinfer`` sampling backends). The CI e2e matrix rebuilds TokenSpeed from source, so the relevant end-to-end pass lives in this PR's GitHub Actions run rather than the local image. Closes #120 Signed-off-by: yetone <yetoneful@gmail.com>
Wraps TokenSpeed's AsyncLLM behind the SGLang proto so the existing Rust router auto-detects and routes to TokenSpeed workers with zero gateway-side changes. Why the SGLang proto? TokenSpeed and SGLang share the tokenized-request shape, sampling- param naming, and output dict format, so reusing the proto avoids inventing a new one — the previous tokenspeed_scheduler.proto attempt in lightseekorg/tokenspeed#185 and #1167 was closed in favor of this approach. Module layout (grpc_servicer/smg_grpc_servicer/tokenspeed/): - servicer.py TokenSpeedSchedulerServicer — Generate, Embed, HealthCheck, Abort, GetModelInfo, GetServerInfo, GetTokenizer, GetLoads - health_servicer.py grpc.health.v1.Health bridge advertising the SGLang service name for router auto-detection - scheduler_launcher.py thin wrapper around TokenSpeed's _launch_subprocesses - server.py python -m entrypoint that boots the scheduler, AsyncLLM, grpc.aio server, and warmup probe - __main__.py CLI shim e2e_test infra: - Runtime.TOKENSPEED enum (infra/constants.py) - _build_tokenspeed_grpc_cmd so e2e fixtures can spawn tokenspeed workers via python -m smg_grpc_servicer.tokenspeed - Qwen/Qwen3-4B and Qwen/Qwen3-30B-A3B model specs for coverage - TestOpenAIServerFunctionCalling now runs the shared Llama-3.2-1B- Instruct fixture against tokenspeed alongside sglang/vllm/trtllm. TokenSpeed support for dense ``LlamaForCausalLM`` is wired up via lightseekorg/tokenspeed#357 (a new ``tokenspeed.runtime.models.llama`` module registered in the model registry); ci_install_tokenspeed.sh pins to that branch until the upstream PR merges. Test evidence (single B200, shared Llama-3.2-1B-Instruct fixture): - Unit: 47 / 47 PASSED (grpc_servicer/tests/) - Registry: ``LlamaForCausalLM`` resolves to the new dense class in the ``tokenspeed:asyncllm`` image; weights load cleanly (``Load weight end. type=LlamaForCausalLM, dtype=torch.bfloat16``) - End-to-end generation smoke against the tokenspeed image surfaced an unrelated detokenizer/sampler issue in the ``tokenspeed:asyncllm`` build itself (reproduces on stock ``Qwen3-0.6B`` with both the ``greedy`` and ``flashinfer`` sampling backends). The CI e2e matrix rebuilds TokenSpeed from source, so the relevant end-to-end pass lives in this PR's GitHub Actions run rather than the local image. Closes #120 Signed-off-by: yetone <yetoneful@gmail.com>
Wraps TokenSpeed's AsyncLLM behind the SGLang proto so the existing Rust router auto-detects and routes to TokenSpeed workers with zero gateway-side changes. Why the SGLang proto? TokenSpeed and SGLang share the tokenized-request shape, sampling- param naming, and output dict format, so reusing the proto avoids inventing a new one — the previous tokenspeed_scheduler.proto attempt in lightseekorg/tokenspeed#185 and #1167 was closed in favor of this approach. Module layout (grpc_servicer/smg_grpc_servicer/tokenspeed/): - servicer.py TokenSpeedSchedulerServicer — Generate, Embed, HealthCheck, Abort, GetModelInfo, GetServerInfo, GetTokenizer, GetLoads - health_servicer.py grpc.health.v1.Health bridge advertising the SGLang service name for router auto-detection - scheduler_launcher.py thin wrapper around TokenSpeed's _launch_subprocesses - server.py python -m entrypoint that boots the scheduler, AsyncLLM, grpc.aio server, and warmup probe - __main__.py CLI shim e2e_test infra: - Runtime.TOKENSPEED enum (infra/constants.py) - _build_tokenspeed_grpc_cmd so e2e fixtures can spawn tokenspeed workers via python -m smg_grpc_servicer.tokenspeed - Qwen/Qwen3-4B and Qwen/Qwen3-30B-A3B model specs for coverage - TestOpenAIServerFunctionCalling now runs the shared Llama-3.2-1B- Instruct fixture against tokenspeed alongside sglang/vllm/trtllm. TokenSpeed support for dense ``LlamaForCausalLM`` is wired up via lightseekorg/tokenspeed#357 (a new ``tokenspeed.runtime.models.llama`` module registered in the model registry); ci_install_tokenspeed.sh pins to that branch until the upstream PR merges. Test evidence (single B200, shared Llama-3.2-1B-Instruct fixture): - Unit: 47 / 47 PASSED (grpc_servicer/tests/) - Registry: ``LlamaForCausalLM`` resolves to the new dense class in the ``tokenspeed:asyncllm`` image; weights load cleanly (``Load weight end. type=LlamaForCausalLM, dtype=torch.bfloat16``) - End-to-end generation smoke against the tokenspeed image surfaced an unrelated detokenizer/sampler issue in the ``tokenspeed:asyncllm`` build itself (reproduces on stock ``Qwen3-0.6B`` with both the ``greedy`` and ``flashinfer`` sampling backends). The CI e2e matrix rebuilds TokenSpeed from source, so the relevant end-to-end pass lives in this PR's GitHub Actions run rather than the local image. Closes #120 Signed-off-by: yetone <yetoneful@gmail.com>
Wraps TokenSpeed's AsyncLLM behind the SGLang proto so the existing Rust router auto-detects and routes to TokenSpeed workers with zero gateway-side changes. Why the SGLang proto? TokenSpeed and SGLang share the tokenized-request shape, sampling- param naming, and output dict format, so reusing the proto avoids inventing a new one — the previous tokenspeed_scheduler.proto attempt in lightseekorg/tokenspeed#185 and #1167 was closed in favor of this approach. Module layout (grpc_servicer/smg_grpc_servicer/tokenspeed/): - servicer.py TokenSpeedSchedulerServicer — Generate, Embed, HealthCheck, Abort, GetModelInfo, GetServerInfo, GetTokenizer, GetLoads - health_servicer.py grpc.health.v1.Health bridge advertising the SGLang service name for router auto-detection - scheduler_launcher.py thin wrapper around TokenSpeed's _launch_subprocesses - server.py python -m entrypoint that boots the scheduler, AsyncLLM, grpc.aio server, and warmup probe - __main__.py CLI shim e2e_test infra: - Runtime.TOKENSPEED enum (infra/constants.py) - _build_tokenspeed_grpc_cmd so e2e fixtures can spawn tokenspeed workers via python -m smg_grpc_servicer.tokenspeed - Qwen/Qwen3-4B and Qwen/Qwen3-30B-A3B model specs for coverage - TestOpenAIServerFunctionCalling now runs the shared Llama-3.2-1B- Instruct fixture against tokenspeed alongside sglang/vllm/trtllm. TokenSpeed support for dense ``LlamaForCausalLM`` is wired up via lightseekorg/tokenspeed#357 (a new ``tokenspeed.runtime.models.llama`` module registered in the model registry); ci_install_tokenspeed.sh pins to that branch until the upstream PR merges. Test evidence (single B200, shared Llama-3.2-1B-Instruct fixture): - Unit: 47 / 47 PASSED (grpc_servicer/tests/) - Registry: ``LlamaForCausalLM`` resolves to the new dense class in the ``tokenspeed:asyncllm`` image; weights load cleanly (``Load weight end. type=LlamaForCausalLM, dtype=torch.bfloat16``) - End-to-end generation smoke against the tokenspeed image surfaced an unrelated detokenizer/sampler issue in the ``tokenspeed:asyncllm`` build itself (reproduces on stock ``Qwen3-0.6B`` with both the ``greedy`` and ``flashinfer`` sampling backends). The CI e2e matrix rebuilds TokenSpeed from source, so the relevant end-to-end pass lives in this PR's GitHub Actions run rather than the local image. Closes #120 Signed-off-by: yetone <yetoneful@gmail.com>
Wraps TokenSpeed's AsyncLLM behind the SGLang proto so the existing Rust router auto-detects and routes to TokenSpeed workers with zero gateway-side changes. Why the SGLang proto? TokenSpeed and SGLang share the tokenized-request shape, sampling- param naming, and output dict format, so reusing the proto avoids inventing a new one — the previous tokenspeed_scheduler.proto attempt in lightseekorg/tokenspeed#185 and #1167 was closed in favor of this approach. Module layout (grpc_servicer/smg_grpc_servicer/tokenspeed/): - servicer.py TokenSpeedSchedulerServicer — Generate, Embed, HealthCheck, Abort, GetModelInfo, GetServerInfo, GetTokenizer, GetLoads - health_servicer.py grpc.health.v1.Health bridge advertising the SGLang service name for router auto-detection - scheduler_launcher.py thin wrapper around TokenSpeed's _launch_subprocesses - server.py python -m entrypoint that boots the scheduler, AsyncLLM, grpc.aio server, and warmup probe - __main__.py CLI shim e2e_test infra: - Runtime.TOKENSPEED enum (infra/constants.py) - _build_tokenspeed_grpc_cmd so e2e fixtures can spawn tokenspeed workers via python -m smg_grpc_servicer.tokenspeed - Qwen/Qwen3-4B and Qwen/Qwen3-30B-A3B model specs for coverage - TestOpenAIServerFunctionCalling now runs the shared Llama-3.2-1B- Instruct fixture against tokenspeed alongside sglang/vllm/trtllm. TokenSpeed support for dense ``LlamaForCausalLM`` is wired up via lightseekorg/tokenspeed#357 (a new ``tokenspeed.runtime.models.llama`` module registered in the model registry); ci_install_tokenspeed.sh pins to that branch until the upstream PR merges. Test evidence (single B200, shared Llama-3.2-1B-Instruct fixture): - Unit: 47 / 47 PASSED (grpc_servicer/tests/) - Registry: ``LlamaForCausalLM`` resolves to the new dense class in the ``tokenspeed:asyncllm`` image; weights load cleanly (``Load weight end. type=LlamaForCausalLM, dtype=torch.bfloat16``) - End-to-end generation smoke against the tokenspeed image surfaced an unrelated detokenizer/sampler issue in the ``tokenspeed:asyncllm`` build itself (reproduces on stock ``Qwen3-0.6B`` with both the ``greedy`` and ``flashinfer`` sampling backends). The CI e2e matrix rebuilds TokenSpeed from source, so the relevant end-to-end pass lives in this PR's GitHub Actions run rather than the local image. Closes #120 Signed-off-by: yetone <yetoneful@gmail.com>
Summary
Adds a gRPC client for the TokenSpeed inference engine, alongside the existing SGLang / vLLM / TRT-LLM / MLX clients. This is the SMG Router side of the TokenSpeed gRPC integration in lightseekorg/tokenspeed#185.
Layout
proto/tokenspeed_scheduler.protoTokenSpeedSchedulerservice: 9 RPCs + TokenSpeed-specific fields (spec_verify_ct,accept_draft_tokens); imports sharedcommon.protoforGetTokenizer/KvEventstypessrc/tokenspeed_scheduler.rsTokenSpeedSchedulerClientwithAbortOnDropStream, trace injection on every RPC,grpcs://TLS support, sharedimpl_get_tokenizer!()/impl_subscribe_kv_events!()macrostests/tokenspeed_scheduler_integration.rsTOKENSPEED_GRPC_ENDPOINT)build.rstokenspeed_scheduler.protosrc/lib.rsRPCs
GenerateAbortOnDropStream→ fire-and-forget Abort on drop,mark_completedon natural end)EmbedAbortResult<(), tonic::Status>matching other backends)GetModelInfoHealthCheckGetServerInfoGetLoadsGetTokenizerSubscribeKvEventsAll RPC methods inject OpenTelemetry trace context via
BoxedTraceInjector.Test evidence
cargo test -p smg-grpc-client tokenspeed_scheduler)tests/tokenspeed_scheduler_integration.rs, real server on B200)The integration suite spins up
tokenspeed serve --grpc-port(the companion PR) inside the same container and exercises every RPC wrapper against a liveopenai/gpt-oss-20bmodel. It directly validates code paths that unit tests cannot reach:connect("grpc://…")handshake with keepalive / window settingshealth_check/get_model_info/get_server_info/get_loadsagainst a real server (each populated correctly)generatenon-streaming round-trip → at least one output token returnedgeneratestreaming → terminates with aCompleteframeabort_requestout-of-band reaches the server and returnsOk(())AbortOnDropStream::Dropfires the fire-and-forget Abort RPC; client remains usable for subsequent RPCs (no state leak)TraceInjectorasserts thatinjectis called exactly once per RPC (5 unary + 2 stream) — direct evidence the newinject_tracehelper is wired on every endpointRun locally:
Companion PR
AsyncLLM(EngineClient) directly. Latest B200 result withopenai/gpt-oss-20b: 38 / 38 unit + 14 / 14 E2E passing.Build verification
Related: lightseekorg/tokenspeed#120
Summary by CodeRabbit
New Features
Tests