Repository navigation
Conversation
Introduces TokenSpeed as a first-class worker alongside SGLang, vLLM, TensorRT-LLM, and MLX — with a dedicated scheduler proto/service so the two backends can diverge freely over time. Python servicer (grpc_servicer/smg_grpc_servicer/tokenspeed/): - Scheduler launcher wraps tokenspeed.AsyncLLM + scheduler subprocess - Async gRPC servicer implements all 9 RPCs (Generate, Embed, Abort, HealthCheck, GetModelInfo, GetServerInfo, GetLoads, GetTokenizer, SubscribeKvEvents) using shared SGLang messages via proto import - Health servicer advertises tokenspeed.grpc.scheduler.TokenSpeedScheduler - Standalone server mirrors the SGLang servicer's lifecycle + warmup - pyproject.toml adds a "tokenspeed" optional-dep group - 47 unit tests under grpc_servicer/tests/ cover launcher, servicer, and health probe paths Rust router integration: - New crates/grpc_client/proto/tokenspeed_scheduler.proto (+ mirror under python/smg_grpc_proto/proto/) for the dedicated service, with messages imported from sglang_scheduler.proto today - build.rs compiles the proto into a separate OUT_DIR with extern_path back to the SGLang types to avoid duplication - TokenSpeedSchedulerClient in crates/grpc_client/src/tokenspeed_scheduler.rs with the same surface as SglangSchedulerClient (sampling_params builders are now pub(crate) and reused) - AbortOnDropStream is refactored onto a client-agnostic AbortDispatcher closure so both clients can share it - protocols::worker::RuntimeType gains a TokenSpeed variant (Display / FromStr) - GrpcClient enum, reachability probe, detect_backend ordering, harmony request building, multimodal dispatch, and embedding request building all learn the new arm E2E infra: - Runtime.TOKENSPEED added to constants.py - worker.py gets _build_tokenspeed_grpc_cmd and launches via python -m smg_grpc_servicer.tokenspeed (with xgrammar grammar backend) - model_specs.py registers Qwen/Qwen3-4B and Qwen/Qwen3-30B-A3B - test_function_calling.py runs against tokenspeed; the suite temporarily pins Qwen3-30B-A3B + qwen parser because TokenSpeed does not support LlamaForCausalLM CI: - New .github/actions/setup-tokenspeed/action.yml + scripts/ci_install_tokenspeed.sh - e2e-gpu-job.yml and pr-test-rust.yml wire in the new engine matrix entry E2E results on Qwen3-30B-A3B: 12 passed / 4 failed / 2 skipped, matching vLLM on the same model/suite (10/6/2 on vLLM). The 4 failures reproduce identically on both backends; root cause is the SMG gateway's tool_choice=required/specific constraint translation layer and is unrelated to this change.
📝 WalkthroughWalkthroughThis PR introduces TokenSpeed as a new inference engine integrated via a gRPC servicer. It adds a complete TokenSpeed gRPC server implementation wrapping TokenSpeed's AsyncLLM backend using the SGLang gRPC protocol for compatibility with the existing router. Changes span proto definitions, Rust gRPC client, Python servicer with health checks, model gateway dispatch logic, E2E test infrastructure, and CI/workflow automation. Changes
Sequence Diagram(s)sequenceDiagram
actor User
participant Main as __main__.py
participant Server as server.py
participant AsyncLLM as AsyncLLM<br/>(TokenSpeed)
participant Servicer as TokenSpeedScheduler<br/>Servicer
participant Health as Health<br/>Servicer
User->>Main: python -m smg_grpc_servicer.tokenspeed
Main->>Main: parse args, set uvloop
Main->>Server: await serve_grpc(server_args)
Server->>AsyncLLM: launch_engine() + AsyncLLM init
AsyncLLM-->>Server: initialized
Server->>Server: create gRPC server
Server->>Servicer: register TokenSpeedSchedulerServicer
Server->>Health: register TokenSpeedHealthServicer
Health->>Health: set_not_serving()
Server->>Server: start gRPC server
Server->>Server: spawn warmup thread
rect rgba(100, 150, 255, 0.5)
Note over Server: Warmup Phase
Server->>Servicer: GetModelInfo()
Servicer->>AsyncLLM: fetch model info
AsyncLLM-->>Servicer: response
Servicer-->>Server: success (with timeout)
alt is_generation
Server->>Servicer: Generate(one-token probe)
else
Server->>Servicer: Embed(probe)
end
Servicer->>AsyncLLM: process request
AsyncLLM-->>Servicer: response
Servicer-->>Server: complete
end
Server->>Health: set_serving()
Server->>User: ✅ Server ready, blocking on signal
sequenceDiagram
participant Client as gRPC Client<br/>(Router)
participant Servicer as TokenSpeedScheduler<br/>Servicer
participant AsyncLLM as AsyncLLM
participant Engine as TokenSpeed<br/>Engine
Client->>Servicer: Generate(GenerateRequest)
rect rgba(100, 150, 255, 0.5)
Note over Servicer: Request Conversion
Servicer->>Servicer: parse proto, build TokenSpeed input
Servicer->>Servicer: expand rid for n>1 (per-choice)
end
Servicer->>AsyncLLM: add_request(rid, req)
rect rgba(150, 200, 150, 0.5)
Note over Servicer,Engine: Streaming Loop
loop for each response batch
Engine-->>AsyncLLM: next token(s) + metadata
AsyncLLM-->>Servicer: yield output dict
Servicer->>Servicer: extract choice index
Servicer->>Servicer: trim output_ids (stop token removal)
Servicer->>Servicer: normalize finish_reason
Servicer-->>Client: stream GenerateResponse
end
end
rect rgba(200, 150, 100, 0.5)
Note over Servicer: Cleanup
Servicer->>Servicer: remove from rid_to_state
Servicer->>Servicer: if cancellation: spawn abort RPC
end
Estimated code review effort🎯 4 (Complex) | ⏱️ ~60 minutes The changes span proto definitions, Rust client generation, Python servicer with dense logic (827 lines), health checks with scheduler liveness detection, gateway integration across multiple files, and comprehensive test coverage. The heterogeneous nature (Rust, Python, proto, CI, tests) combined with non-trivial request conversion, output-id trimming, finish-reason mapping, and per-choice index routing requires careful review across multiple domains. Possibly related PRs
Suggested labels
Suggested reviewers
Poem
🚥 Pre-merge checks | ✅ 4 | ❌ 1❌ Failed checks (1 warning)
✅ Passed checks (4 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches📝 Generate docstrings
🧪 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 |
|
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:
|
There was a problem hiding this comment.
Code Review
This pull request implements comprehensive support for the TokenSpeed backend, including a new gRPC servicer, proto definitions, and integration into the Rust client and model gateway. The changes also extend the E2E testing framework and CI scripts to accommodate TokenSpeed. Review feedback identifies critical issues in the sampling parameter translation logic, where zero-valued temperatures are incorrectly omitted, preventing greedy decoding, and default top_k values from the Rust client may not be correctly mapped for the TokenSpeed engine.
| if params.top_p: | ||
| out["top_p"] = params.top_p |
There was a problem hiding this comment.
The check if params.temperature: will evaluate to False when temperature is 0.0 (greedy decoding), causing it to be omitted from the out dictionary. This results in the engine using its default temperature (likely 1.0) instead of the requested greedy decoding. Since 0.0 is a valid and common value for temperature, this heuristic prevents greedy decoding from working. If the proto field is not optional, you might need to always forward it or use a different sentinel value. If you can change the proto, marking these fields as optional is the recommended fix.
| if params.top_k: | ||
| out["top_k"] = params.top_k | ||
| if params.min_p: | ||
| out["min_p"] = params.min_p |
There was a problem hiding this comment.
The SGLang gRPC client (Rust) defaults top_k to -1 to indicate it is disabled. The current logic if params.top_k: will evaluate to True for -1, forwarding it to TokenSpeed. If TokenSpeed expects 0 to disable top_k (as suggested by the comment) and might not handle -1 correctly, you should explicitly map -1 to 0 or ensure only positive values are forwarded if that's what the engine expects.
| out["frequency_penalty"] = params.frequency_penalty | ||
| if params.presence_penalty: |
There was a problem hiding this comment.
🔴 Important: temperature=0.0 is falsy in Python, so this check silently drops explicit greedy-decoding requests. The router sets temperature: request.temperature.unwrap_or(1.0) (see sglang_scheduler.rs:467), so when a user sends temperature=0 the proto carries 0.0 — but if params.temperature: evaluates to False and the field is never forwarded. TokenSpeed's default is likely 1.0, meaning the user gets random sampling instead of deterministic output.
The same issue applies to top_p in theory (top_p=0.0 is nonsensical so it's not a practical concern), but temperature=0.0 is one of the most common sampling overrides.
Proto3 non-optional scalars can't distinguish "not set" from "set to zero". Since the router always populates temperature (defaulting to 1.0 when the user omits it), you can safely forward it unconditionally:
| out["frequency_penalty"] = params.frequency_penalty | |
| if params.presence_penalty: | |
| out["temperature"] = params.temperature | |
| if params.top_p: |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 960d4eb8eb
ℹ️ About Codex in GitHub
Codex has been enabled to automatically 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 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
| finally: | ||
| channel.close() | ||
|
|
||
| health_servicer.set_serving() |
There was a problem hiding this comment.
Keep health NOT_SERVING when warmup request fails
The warmup routine marks the service as SERVING even after the Generate/Embed probe throws, because the exception path only logs and then still calls set_serving(). In that failure mode (e.g., model init succeeds enough for GetModelInfo but inference path is broken), health checks report ready and the router can start sending production traffic to a worker that already failed its startup probe.
Useful? React with 👍 / 👎.
| if params.temperature: | ||
| out["temperature"] = params.temperature | ||
| if params.top_p: | ||
| out["top_p"] = params.top_p |
There was a problem hiding this comment.
Preserve zero-valued sampling params when translating proto
This translation drops temperature=0.0 and top_p=0.0 because it only forwards truthy float values. 0.0 is a meaningful user setting (especially for deterministic decoding), so TokenSpeed requests can silently fall back to backend defaults instead of honoring caller intent. That creates behavior drift versus other backends for identical API requests.
Useful? React with 👍 / 👎.
|
|
||
| # Step 3: kernel (CUDA compile — the expensive one). Try the cached wheel first. | ||
| CACHED_KERNEL_WHEEL=$(find "$WHEEL_CACHE" -name "tokenspeed_kernel-*.whl" 2>/dev/null | head -1 || true) | ||
| if [ -n "$CACHED_KERNEL_WHEEL" ] && [ -f "$CACHED_KERNEL_WHEEL" ]; then | ||
| echo "Installing cached tokenspeed-kernel wheel: $CACHED_KERNEL_WHEEL" | ||
| uv pip install "$CACHED_KERNEL_WHEEL" --no-build-isolation |
There was a problem hiding this comment.
🟡 Nit: This block is meant to "cache the built wheel" but it only prints the dist-info path and never actually copies any .whl file into $WHEEL_CACHE. Every CI run will rebuild the kernel from source (~30 min). The find check at line 70 will never find a cached wheel.
If uv doesn't leave a .whl around, you could use pip wheel --no-deps tokenspeed-kernel/python/ -w "$WHEEL_CACHE" as a post-install step to produce a reusable wheel, or drop the caching logic until it's wired up.
There was a problem hiding this comment.
Actionable comments posted: 12
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (3)
e2e_test/infra/constants.py (1)
67-76: 🧹 Nitpick | 🔵 TrivialStale docstring on
get_runtime.The
ENV_RUNTIMEcomment on line 56 was correctly updated to enumerateSGLANG/VLLM/TRTLLM/TOKENSPEED, butget_runtime's docstring at line 68 still says "sglang or vllm". Worth updating for consistency.📝 Proposed fix
def get_runtime() -> str: - """Get the current test runtime (sglang or vllm). + """Get the current test runtime (sglang, vllm, trtllm, or tokenspeed).🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@e2e_test/infra/constants.py` around lines 67 - 76, Update the stale docstring on get_runtime to reflect all supported runtimes: SGLANG, VLLM, TRTLLM, and TOKENSPEED; keep references to ENV_RUNTIME, DEFAULT_RUNTIME, and _RUNTIME_CACHE intact and ensure the Returns section documents that the function reads E2E_RUNTIME and defaults to DEFAULT_RUNTIME when unset.crates/grpc_client/src/lib.rs (1)
1-4: 🧹 Nitpick | 🔵 TrivialCrate-level doc is stale: mention TokenSpeed.
The re-exports now publicly expose
TokenSpeedSchedulerClient/tokenspeed_proto, but the module header still only lists SGLang, vLLM, TensorRT-LLM, and MLX.📝 Proposed doc tweak
-//! gRPC clients for SGLang, vLLM, TensorRT-LLM, and MLX backends +//! gRPC clients for SGLang, vLLM, TensorRT-LLM, MLX, and TokenSpeed backends //! //! This crate provides gRPC client implementations for communicating with -//! SGLang scheduler, vLLM engine, TensorRT-LLM engine, and MLX engine backends. +//! SGLang scheduler, vLLM engine, TensorRT-LLM engine, MLX engine, and +//! TokenSpeed scheduler backends.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@crates/grpc_client/src/lib.rs` around lines 1 - 4, The crate-level documentation is stale and doesn't mention the newly exposed TokenSpeed client; update the top-level doc comment to include TokenSpeed (e.g., reference TokenSpeedSchedulerClient and the tokenspeed_proto re-export) alongside SGLang, vLLM, TensorRT-LLM, and MLX so the crate docs accurately reflect public re-exports like TokenSpeedSchedulerClient and tokenspeed_proto.model_gateway/src/workflow/steps/local/detect_backend.rs (1)
29-32: 🧹 Nitpick | 🔵 TrivialFunction docstring is stale: probe order no longer matches.
Line 32 still says "Otherwise tries sglang → vllm → trtllm → mlx", but the loop on line 55 now tries
tokenspeedfirst. Update the docstring for consistency with the module-level doc on lines 6-9.📝 Proposed fix
-/// If `runtime_hint` is provided (from explicit config), tries that first. -/// Otherwise tries sglang → vllm → trtllm → mlx. +/// If `runtime_hint` is provided (from explicit config), tries that first. +/// Otherwise tries tokenspeed → sglang → vllm → trtllm → mlx.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@model_gateway/src/workflow/steps/local/detect_backend.rs` around lines 29 - 32, The function docstring for the backend detection is stale and lists the wrong probe order; update the docstring above the detect-backend function in detect_backend.rs to reflect the actual probe sequence used by the loop (tokenspeed before sglang, then vllm → trtllm → mlx), or align it with the module-level docstring; specifically mention "tries tokenspeed → sglang → vllm → trtllm → mlx" (and that an explicit runtime_hint is tried first) so the text matches the code where the loop tests tokenspeed first.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In @.github/workflows/pr-test-rust.yml:
- Around line 450-453: The workflow currently sets a 60-minute timeout for the
tokenspeed matrix entry (engine: tokenspeed, timeout: 60) which is fine; as a
follow-up optimization, add caching around the TokenSpeed kernel/scheduler build
by caching the build output keyed on the installer script hash
(scripts/ci_install_tokenspeed.sh) so repeated runs skip rebuilding—implement a
cache action step that saves/restores the built artifacts using a cache key that
includes the script file hash and relevant platform identifiers.
In `@crates/grpc_client/src/tokenspeed_scheduler.rs`:
- Around line 212-379: The code uses the magic literal -1 for logprob_start_len
in build_generate_request_from_chat and build_generate_request_from_completion
while build_plain_generate_request uses unwrap_or(-1); define a shared named
constant (e.g. DEFAULT_LOGPROB_START_LEN: i32 = -1) and replace the literal
occurrences in build_generate_request_from_chat,
build_generate_request_from_completion and the unwrap_or in
build_plain_generate_request with that constant to make the intent explicit and
keep parity with SglangSchedulerClient helpers.
In `@e2e_test/infra/worker.py`:
- Around line 271-277: The docstring for _build_tokenspeed_grpc_cmd is
incorrect: TokenSpeed is auto-detected as "tokenspeed", not as SGLang; update
the docstring to state that this launches the SMG-hosted TokenSpeed gRPC server
(smg_grpc_servicer.tokenspeed) which advertises the tokenspeed.grpc.scheduler
service and is auto-detected by the Rust router as "tokenspeed" (remove the
SGLang claim), so that the function _build_tokenspeed_grpc_cmd and its
description accurately reflect the detection behavior.
In `@grpc_servicer/smg_grpc_servicer/tokenspeed/__main__.py`:
- Around line 33-38: The comment about scheduler processes reading env vars is
misplaced; it should either be moved to document prepare_server_args(argv)
(since prepare_server_args is the authoritative env/resource setup path) or
removed if it isn't describing that behavior. Update the file so the comment
sits immediately above the prepare_server_args(argv) call (or delete it), and
ensure lines containing asyncio.set_event_loop_policy(uvloop.EventLoopPolicy())
and asyncio.run(serve_grpc(server_args)) no longer carry that env-var
explanation.
In `@grpc_servicer/smg_grpc_servicer/tokenspeed/scheduler_launcher.py`:
- Around line 16-43: The code currently imports and unpacks the private symbol
_launch_subprocesses inside launch_engine, which can silently break if its
return tuple shape changes; to fix, either pin the tokenspeed dependency in
grpc_servicer/pyproject.toml to a specific compatible version, or add a
defensive arity check immediately after calling _launch_subprocesses in
launch_engine (e.g., verify the returned value is a tuple/list and has length 3)
and raise a clear RuntimeError if it does not, referencing _launch_subprocesses
and launch_engine in the error message so failures fail fast and are easy to
debug.
In `@grpc_servicer/smg_grpc_servicer/tokenspeed/server.py`:
- Around line 207-213: The warmup exception is being swallowed so
health_servicer.set_serving() is called regardless; change the warmup error
handling in the try/except around the warmup routine so that on any Exception
you log the error (logger.warning("TokenSpeed warmup failed: %s", e)) and then
return (or set health_servicer.set_not_serving()) instead of falling through to
health_servicer.set_serving(); keep channel.close() in the finally block so the
channel is always closed. Ensure the final placement of
health_servicer.set_serving() occurs only after a successful warmup (mirror the
early return behavior used in the "not connected" branch).
In `@grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py`:
- Around line 609-613: Fix the docstring typo that incorrectly says "bool3" —
update the comment that documents the proto field `log_metrics` in servicer
functions within tokenspeed (the comment block in
grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py referring to
`log_metrics`) to read "bool" instead of "bool3" so the docstring accurately
describes the plain boolean scalar for `log_metrics`.
- Around line 142-143: Remove the unreachable bare "return" statements that
follow "await context.abort(...)" calls in the tokenspeed servicer methods
(e.g., Generate, Embed, GetTokenizer and other handlers referenced). Locate each
usage of "await context.abort(grpc.StatusCode.INVALID_ARGUMENT, ...)" in
servicer methods and delete the subsequent "return" statement(s) so the code
follows the repo convention that context.abort raises grpc.aio.AbortError and
the explicit returns are dead code.
- Line 305: The health-check request id generated as rid =
f"HEALTH_CHECK_{time.time()}" is not guaranteed unique across concurrent probes
and can cause self.async_llm.abort_request(rid) to abort another probe; change
the rid generation in the health-check handler to use a higher-entropy
identifier (e.g., time.time_ns() or uuid.uuid4()) and import uuid if using UUIDs
so each health-check has a unique rid before calling
self.async_llm.start/abort_request.
- Around line 338-361: The probe task created by
asyncio.create_task(_drive_probe()) may be cancelled but never awaited, leaking
the coroutine; modify the finally block to, after calling task.cancel() when not
task.done(), await the task inside a contextlib.suppress(asyncio.CancelledError,
Exception) to ensure the _drive_probe() coroutine is finished/cleaned up before
proceeding, and add import contextlib at the top; keep the existing
self.async_llm.abort_request(rid) call for best-effort scheduler cleanup.
In `@grpc_servicer/tests/test_tokenspeed_servicer.py`:
- Around line 163-180: FakeAsyncLLM.generate_request currently assumes obj.rid
is hashable and does self.rid_to_state[rid] = _FakeState(), which fails when
obj.rid is a list for sampling_params.n > 1; change generate_request to detect
if rid is a list and, if so, iterate over each child id and set
self.rid_to_state[child_rid] = _FakeState() (and later mark each
child_rid.finished = True) instead of using the list as a dict key; keep
existing behavior for scalar rid, and ensure last_receive_tstamp and
generate_fn/output yielding logic remain the same.
In `@scripts/ci_install_tokenspeed.sh`:
- Around line 86-96: The CI is not actually populating the wheel cache: the uv
pip install tokenspeed-kernel/python/ call installs in-place and the python3 -c
block that inspects tokenspeed_kernel only prints diagnostic info and never
copies a built .whl into $WHEEL_CACHE, so every run rebuilds the kernel; fix by
after the install locating the built wheel inside uv's cache (or the pip/uv
wheel cache directory) and copying that .whl into $WHEEL_CACHE (or alternatively
remove the broken caching logic), updating the script sections that reference
WHEEL_CACHE and the python3 -c diagnostics/tokenspeed_kernel lookup to perform a
concrete copy of the discovered wheel into $WHEEL_CACHE so subsequent CI runs
reuse the cached wheel.
---
Outside diff comments:
In `@crates/grpc_client/src/lib.rs`:
- Around line 1-4: The crate-level documentation is stale and doesn't mention
the newly exposed TokenSpeed client; update the top-level doc comment to include
TokenSpeed (e.g., reference TokenSpeedSchedulerClient and the tokenspeed_proto
re-export) alongside SGLang, vLLM, TensorRT-LLM, and MLX so the crate docs
accurately reflect public re-exports like TokenSpeedSchedulerClient and
tokenspeed_proto.
In `@e2e_test/infra/constants.py`:
- Around line 67-76: Update the stale docstring on get_runtime to reflect all
supported runtimes: SGLANG, VLLM, TRTLLM, and TOKENSPEED; keep references to
ENV_RUNTIME, DEFAULT_RUNTIME, and _RUNTIME_CACHE intact and ensure the Returns
section documents that the function reads E2E_RUNTIME and defaults to
DEFAULT_RUNTIME when unset.
In `@model_gateway/src/workflow/steps/local/detect_backend.rs`:
- Around line 29-32: The function docstring for the backend detection is stale
and lists the wrong probe order; update the docstring above the detect-backend
function in detect_backend.rs to reflect the actual probe sequence used by the
loop (tokenspeed before sglang, then vllm → trtllm → mlx), or align it with the
module-level docstring; specifically mention "tries tokenspeed → sglang → vllm →
trtllm → mlx" (and that an explicit runtime_hint is tried first) so the text
matches the code where the loop tests tokenspeed first.
🪄 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: 3214c17e-34f1-4133-aee6-fcad33da32bf
📒 Files selected for processing (34)
.github/actions/setup-tokenspeed/action.yml.github/workflows/e2e-gpu-job.yml.github/workflows/pr-test-rust.ymlcrates/grpc_client/build.rscrates/grpc_client/proto/tokenspeed_scheduler.protocrates/grpc_client/python/smg_grpc_proto/__init__.pycrates/grpc_client/src/lib.rscrates/grpc_client/src/sglang_scheduler.rscrates/grpc_client/src/tokenspeed_scheduler.rscrates/protocols/src/worker.rse2e_test/chat_completions/test_function_calling.pye2e_test/infra/constants.pye2e_test/infra/model_specs.pye2e_test/infra/worker.pygrpc_servicer/pyproject.tomlgrpc_servicer/smg_grpc_servicer/tokenspeed/__init__.pygrpc_servicer/smg_grpc_servicer/tokenspeed/__main__.pygrpc_servicer/smg_grpc_servicer/tokenspeed/health_servicer.pygrpc_servicer/smg_grpc_servicer/tokenspeed/scheduler_launcher.pygrpc_servicer/smg_grpc_servicer/tokenspeed/server.pygrpc_servicer/smg_grpc_servicer/tokenspeed/servicer.pygrpc_servicer/tests/__init__.pygrpc_servicer/tests/conftest.pygrpc_servicer/tests/test_tokenspeed_health_servicer.pygrpc_servicer/tests/test_tokenspeed_servicer.pymodel_gateway/src/routers/grpc/client.rsmodel_gateway/src/routers/grpc/common/stages/request_execution.rsmodel_gateway/src/routers/grpc/harmony/stages/request_building.rsmodel_gateway/src/routers/grpc/multimodal.rsmodel_gateway/src/routers/grpc/regular/stages/embedding/request_building.rsmodel_gateway/src/workflow/steps/classify.rsmodel_gateway/src/workflow/steps/local/detect_backend.rsmodel_gateway/src/workflow/steps/util.rsscripts/ci_install_tokenspeed.sh
| # TokenSpeed builds kernel (CUDA) + scheduler (C++/CMake) from | ||
| # source, so first run takes ~30 min; cached runs are faster. | ||
| - engine: tokenspeed | ||
| timeout: 60 |
There was a problem hiding this comment.
🧹 Nitpick | 🔵 Trivial
LGTM — 60-min timeout reasonable for from-source builds; consider caching in a follow-up.
The matrix addition is well-scoped and the timeout budget matches the documented build cost. If the installer script ends up dominating CI time, a follow-up that caches the TokenSpeed kernel/scheduler build outputs (keyed on scripts/ci_install_tokenspeed.sh hash) would materially speed up most runs. Not a blocker for this PR.
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In @.github/workflows/pr-test-rust.yml around lines 450 - 453, The workflow
currently sets a 60-minute timeout for the tokenspeed matrix entry (engine:
tokenspeed, timeout: 60) which is fine; as a follow-up optimization, add caching
around the TokenSpeed kernel/scheduler build by caching the build output keyed
on the installer script hash (scripts/ci_install_tokenspeed.sh) so repeated runs
skip rebuilding—implement a cache action step that saves/restores the built
artifacts using a cache key that includes the script file hash and relevant
platform identifiers.
| #[expect( | ||
| clippy::unused_self, | ||
| reason = "receiver kept for API parity with SglangSchedulerClient" | ||
| )] | ||
| pub fn build_embed_request( | ||
| &self, | ||
| request_id: String, | ||
| original_text: Option<String>, | ||
| token_ids: Vec<u32>, | ||
| ) -> proto::EmbedRequest { | ||
| proto::EmbedRequest { | ||
| request_id, | ||
| tokenized: Some(proto::TokenizedInput { | ||
| original_text: original_text.unwrap_or_default(), | ||
| input_ids: token_ids, | ||
| }), | ||
| ..Default::default() | ||
| } | ||
| } | ||
|
|
||
| #[expect( | ||
| clippy::unused_self, | ||
| reason = "receiver kept for API parity with SglangSchedulerClient" | ||
| )] | ||
| pub fn build_generate_request_from_chat( | ||
| &self, | ||
| request_id: String, | ||
| body: &ChatCompletionRequest, | ||
| processed_text: String, | ||
| token_ids: Vec<u32>, | ||
| multimodal_inputs: Option<proto::MultimodalInputs>, | ||
| tool_call_constraint: Option<(String, String)>, | ||
| ) -> Result<proto::GenerateRequest, String> { | ||
| let sampling_params = SglangSchedulerClient::build_grpc_sampling_params_from_chat( | ||
| body, | ||
| tool_call_constraint, | ||
| )?; | ||
| Ok(proto::GenerateRequest { | ||
| request_id, | ||
| tokenized: Some(proto::TokenizedInput { | ||
| original_text: processed_text, | ||
| input_ids: token_ids, | ||
| }), | ||
| mm_inputs: multimodal_inputs, | ||
| sampling_params: Some(sampling_params), | ||
| return_logprob: body.logprobs, | ||
| logprob_start_len: -1, | ||
| top_logprobs_num: body.top_logprobs.unwrap_or(0) as i32, | ||
| return_hidden_states: body.return_hidden_states, | ||
| stream: body.stream, | ||
| ..Default::default() | ||
| }) | ||
| } | ||
|
|
||
| #[expect( | ||
| clippy::unused_self, | ||
| reason = "receiver kept for API parity with SglangSchedulerClient" | ||
| )] | ||
| pub fn build_plain_generate_request( | ||
| &self, | ||
| request_id: String, | ||
| body: &GenerateRequest, | ||
| original_text: Option<String>, | ||
| token_ids: Vec<u32>, | ||
| ) -> Result<proto::GenerateRequest, String> { | ||
| let sampling_params = | ||
| SglangSchedulerClient::build_sampling_params_from_plain(body.sampling_params.as_ref())?; | ||
| Ok(proto::GenerateRequest { | ||
| request_id, | ||
| tokenized: Some(proto::TokenizedInput { | ||
| original_text: original_text.unwrap_or_default(), | ||
| input_ids: token_ids, | ||
| }), | ||
| sampling_params: Some(sampling_params), | ||
| return_logprob: body.return_logprob.unwrap_or(false), | ||
| logprob_start_len: body.logprob_start_len.unwrap_or(-1), | ||
| top_logprobs_num: body.top_logprobs_num.unwrap_or(0), | ||
| token_ids_logprob: body.token_ids_logprob.clone().unwrap_or_default(), | ||
| return_hidden_states: body.return_hidden_states, | ||
| stream: body.stream, | ||
| log_metrics: body.log_metrics, | ||
| ..Default::default() | ||
| }) | ||
| } | ||
|
|
||
| #[expect( | ||
| clippy::unused_self, | ||
| reason = "receiver kept for API parity with SglangSchedulerClient" | ||
| )] | ||
| pub fn build_generate_request_from_responses( | ||
| &self, | ||
| request_id: String, | ||
| body: &ResponsesRequest, | ||
| processed_text: String, | ||
| token_ids: Vec<u32>, | ||
| constraint: Option<(String, String)>, | ||
| ) -> Result<proto::GenerateRequest, String> { | ||
| let sampling_params = | ||
| SglangSchedulerClient::build_grpc_sampling_params_from_responses(body, constraint)?; | ||
| Ok(proto::GenerateRequest { | ||
| request_id, | ||
| tokenized: Some(proto::TokenizedInput { | ||
| original_text: processed_text, | ||
| input_ids: token_ids, | ||
| }), | ||
| sampling_params: Some(sampling_params), | ||
| stream: body.stream.unwrap_or(false), | ||
| ..Default::default() | ||
| }) | ||
| } | ||
|
|
||
| #[expect( | ||
| clippy::unused_self, | ||
| reason = "receiver kept for API parity with SglangSchedulerClient" | ||
| )] | ||
| pub fn build_generate_request_from_messages( | ||
| &self, | ||
| request_id: String, | ||
| body: &CreateMessageRequest, | ||
| processed_text: String, | ||
| token_ids: Vec<u32>, | ||
| multimodal_inputs: Option<proto::MultimodalInputs>, | ||
| tool_call_constraint: Option<(String, String)>, | ||
| ) -> Result<proto::GenerateRequest, String> { | ||
| let sampling_params = SglangSchedulerClient::build_grpc_sampling_params_from_messages( | ||
| body, | ||
| tool_call_constraint, | ||
| )?; | ||
| Ok(proto::GenerateRequest { | ||
| request_id, | ||
| tokenized: Some(proto::TokenizedInput { | ||
| original_text: processed_text, | ||
| input_ids: token_ids, | ||
| }), | ||
| mm_inputs: multimodal_inputs, | ||
| sampling_params: Some(sampling_params), | ||
| stream: body.stream.unwrap_or(false), | ||
| ..Default::default() | ||
| }) | ||
| } | ||
|
|
||
| #[expect( | ||
| clippy::unused_self, | ||
| reason = "receiver kept for API parity with SglangSchedulerClient" | ||
| )] | ||
| pub fn build_generate_request_from_completion( | ||
| &self, | ||
| request_id: String, | ||
| body: &CompletionRequest, | ||
| original_text: String, | ||
| token_ids: Vec<u32>, | ||
| ) -> Result<proto::GenerateRequest, String> { | ||
| let sampling_params = | ||
| SglangSchedulerClient::build_grpc_sampling_params_from_completion(body)?; | ||
| Ok(proto::GenerateRequest { | ||
| request_id, | ||
| tokenized: Some(proto::TokenizedInput { | ||
| original_text, | ||
| input_ids: token_ids, | ||
| }), | ||
| sampling_params: Some(sampling_params), | ||
| return_logprob: body.logprobs.is_some(), | ||
| logprob_start_len: -1, | ||
| top_logprobs_num: body.logprobs.unwrap_or(0) as i32, | ||
| stream: body.stream, | ||
| ..Default::default() | ||
| }) | ||
| } |
There was a problem hiding this comment.
🧹 Nitpick | 🔵 Trivial
Builder delegation to SGLang sampling-param helpers looks correct; consider one readability nit.
Delegation to SglangSchedulerClient::build_grpc_sampling_params_* keeps the two wire-identical backends from drifting in sampling-param handling, and the module-level doc comment already sets expectations for the eventual divergence copy-and-specialize policy.
Optional nit: logprob_start_len: -1 appears as a magic literal in build_generate_request_from_chat (line 258) and build_generate_request_from_completion (line 374). The plain builder uses body.logprob_start_len.unwrap_or(-1) (line 287). A named const DEFAULT_LOGPROB_START_LEN: i32 = -1; would make the "start from last prompt token" intent explicit — but only worth doing if SGLang's builder is updated in lockstep to preserve the intentional parity called out in the module docs.
🤖 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 212 - 379, The
code uses the magic literal -1 for logprob_start_len in
build_generate_request_from_chat and build_generate_request_from_completion
while build_plain_generate_request uses unwrap_or(-1); define a shared named
constant (e.g. DEFAULT_LOGPROB_START_LEN: i32 = -1) and replace the literal
occurrences in build_generate_request_from_chat,
build_generate_request_from_completion and the unwrap_or in
build_plain_generate_request with that constant to make the intent explicit and
keep parity with SglangSchedulerClient helpers.
| def _build_tokenspeed_grpc_cmd(self, model_path: str, tp_size: int, spec: dict) -> list[str]: | ||
| """Build TokenSpeed gRPC server command. | ||
|
|
||
| Launches the SMG-hosted TokenSpeed gRPC server | ||
| (``smg_grpc_servicer.tokenspeed``) which wraps TokenSpeed's AsyncLLM | ||
| behind the SGLang proto. Auto-detected as SGLang by the Rust router. | ||
| """ |
There was a problem hiding this comment.
Docstring is stale/incorrect: TokenSpeed is auto-detected as tokenspeed, not SGLang.
Per model_gateway/src/workflow/steps/local/detect_backend.rs (lines 6-9, 50-55), TokenSpeed advertises its own tokenspeed.grpc.scheduler service and the detection probe returns "tokenspeed" (probed before sglang). The "Auto-detected as SGLang by the Rust router" sentence appears to reflect an earlier iteration of the design (SGLang-proto reuse) and will mislead readers debugging router routing decisions.
📝 Proposed docstring fix
- Launches the SMG-hosted TokenSpeed gRPC server
- (``smg_grpc_servicer.tokenspeed``) which wraps TokenSpeed's AsyncLLM
- behind the SGLang proto. Auto-detected as SGLang by the Rust router.
+ Launches the SMG-hosted TokenSpeed gRPC server
+ (``smg_grpc_servicer.tokenspeed``) which wraps TokenSpeed's AsyncLLM.
+ Exposes the ``tokenspeed.grpc.scheduler`` service (reusing SGLang's
+ message types) and is auto-detected as ``tokenspeed`` by the Rust
+ router's gRPC probe order.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@e2e_test/infra/worker.py` around lines 271 - 277, The docstring for
_build_tokenspeed_grpc_cmd is incorrect: TokenSpeed is auto-detected as
"tokenspeed", not as SGLang; update the docstring to state that this launches
the SMG-hosted TokenSpeed gRPC server (smg_grpc_servicer.tokenspeed) which
advertises the tokenspeed.grpc.scheduler service and is auto-detected by the
Rust router as "tokenspeed" (remove the SGLang claim), so that the function
_build_tokenspeed_grpc_cmd and its description accurately reflect the detection
behavior.
| server_args = prepare_server_args(argv) | ||
| # The scheduler processes will read these env vars; make sure we ran | ||
| # through TokenSpeed's shared env/resource setup path instead of | ||
| # duplicating it here. | ||
| asyncio.set_event_loop_policy(uvloop.EventLoopPolicy()) | ||
| asyncio.run(serve_grpc(server_args)) |
There was a problem hiding this comment.
🧹 Nitpick | 🔵 Trivial
Comment is misplaced / doesn't describe the lines that follow.
The "scheduler processes will read these env vars..." comment sits immediately above asyncio.set_event_loop_policy(uvloop.EventLoopPolicy()) / asyncio.run(...), which have nothing to do with env-var propagation or TokenSpeed's shared env setup. If the intent is to document prepare_server_args(argv) (line 33) as the authoritative env/resource setup path, move the comment above that call; otherwise drop it.
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@grpc_servicer/smg_grpc_servicer/tokenspeed/__main__.py` around lines 33 - 38,
The comment about scheduler processes reading env vars is misplaced; it should
either be moved to document prepare_server_args(argv) (since prepare_server_args
is the authoritative env/resource setup path) or removed if it isn't describing
that behavior. Update the file so the comment sits immediately above the
prepare_server_args(argv) call (or delete it), and ensure lines containing
asyncio.set_event_loop_policy(uvloop.EventLoopPolicy()) and
asyncio.run(serve_grpc(server_args)) no longer carry that env-var explanation.
| from tokenspeed.runtime.engine.async_llm import AsyncLLM | ||
| from tokenspeed.runtime.entrypoints.engine import _launch_subprocesses | ||
| from tokenspeed.runtime.server_args import PortArgs, ServerArgs | ||
|
|
||
| logger = logging.getLogger(__name__) | ||
|
|
||
|
|
||
| def launch_engine( | ||
| server_args: ServerArgs, | ||
| port_args: PortArgs | None = None, | ||
| ) -> tuple[AsyncLLM, dict[str, Any]]: | ||
| """Launch TokenSpeed scheduler subprocess(es) and the main-process AsyncLLM. | ||
|
|
||
| Returns: | ||
| A tuple ``(async_llm, scheduler_info)``. ``async_llm`` is the live | ||
| :class:`AsyncLLM` that the gRPC servicer will drive. ``scheduler_info`` | ||
| is the dict rank-0 sent back once its scheduler was ready (contains | ||
| e.g. ``max_total_num_tokens``, ``max_req_input_len``, ...). | ||
|
|
||
| Raises: | ||
| RuntimeError: If rank-0 scheduler fails to initialize. The original | ||
| ``_launch_subprocesses`` surfaces this by re-raising the EOF/assertion | ||
| error — we propagate it unchanged. | ||
| """ | ||
| async_llm, _template_manager, scheduler_info = _launch_subprocesses( | ||
| server_args=server_args, | ||
| port_args=port_args, | ||
| ) |
There was a problem hiding this comment.
🧩 Analysis chain
🏁 Script executed:
#!/bin/bash
# Confirm _launch_subprocesses signature/return shape in the pinned tokenspeed version.
rg -nP --type=py -C3 '\bdef\s+_launch_subprocesses\s*\('
# Also see what version of tokenspeed is required:
rg -nP 'tokenspeed' grpc_servicer/pyproject.toml || trueRepository: lightseekorg/smg
Length of output: 129
🌐 Web query:
Does the tokenspeed PyPI package document tokenspeed.runtime.entrypoints.engine._launch_subprocesses as a stable public API?
💡 Result:
No, the tokenspeed PyPI package does not document tokenspeed.runtime.entrypoints.engine._launch_subprocesses as a stable public API.
Citations:
🏁 Script executed:
cat grpc_servicer/pyproject.tomlRepository: lightseekorg/smg
Length of output: 1754
🏁 Script executed:
cat grpc_servicer/smg_grpc_servicer/tokenspeed/scheduler_launcher.pyRepository: lightseekorg/smg
Length of output: 2539
Brittle dependency on TokenSpeed private API _launch_subprocesses.
The import at line 17 and tuple unpacking at line 40 depend on an underscore-prefixed symbol that is not documented as a stable public API. Upstream can change the return signature (e.g., from 3-tuple to 2-tuple) without a deprecation cycle, causing a silent unpacking failure: ValueError: not enough values to unpack.
The None guard (line 48–52) handles rank-0 vs non-rank-0 logic but does not validate the tuple shape. No version pinning exists in pyproject.toml (tokenspeed is installed out-of-tree as editable).
Mitigations:
- Explicitly pin tokenspeed version in
grpc_servicer/pyproject.toml, or - Add a lightweight import-time arity check (e.g.,
len(result) == 3) so failure is "fails fast on import" rather than "fails deep inside unpacking."
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@grpc_servicer/smg_grpc_servicer/tokenspeed/scheduler_launcher.py` around
lines 16 - 43, The code currently imports and unpacks the private symbol
_launch_subprocesses inside launch_engine, which can silently break if its
return tuple shape changes; to fix, either pin the tokenspeed dependency in
grpc_servicer/pyproject.toml to a specific compatible version, or add a
defensive arity check immediately after calling _launch_subprocesses in
launch_engine (e.g., verify the returned value is a tuple/list and has length 3)
and raise a clear RuntimeError if it does not, referencing _launch_subprocesses
and launch_engine in the error message so failures fail fast and are easy to
debug.
| request with ``log_metrics=False`` so health checks don't skew | ||
| Prometheus counters. | ||
| """ | ||
| rid = f"HEALTH_CHECK_{time.time()}" |
There was a problem hiding this comment.
Health-check rid can collide across concurrent probes.
f"HEALTH_CHECK_{time.time()}" relies on sub-second wall-clock resolution to remain unique. Two overlapping health probes (e.g., from different clients) can produce the same rid, so the self.async_llm.abort_request(rid) in the finally block can end up aborting the other probe's request. Consider time.time_ns() or a UUID:
🛠️ Proposed fix
- rid = f"HEALTH_CHECK_{time.time()}"
+ rid = f"HEALTH_CHECK_{uuid.uuid4().hex}"(plus import uuid at the top).
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py` at line 305, The
health-check request id generated as rid = f"HEALTH_CHECK_{time.time()}" is not
guaranteed unique across concurrent probes and can cause
self.async_llm.abort_request(rid) to abort another probe; change the rid
generation in the health-check handler to use a higher-entropy identifier (e.g.,
time.time_ns() or uuid.uuid4()) and import uuid if using UUIDs so each
health-check has a unique rid before calling self.async_llm.start/abort_request.
| task = asyncio.create_task(_drive_probe()) | ||
| try: | ||
| while time.time() - tic < HEALTH_CHECK_TIMEOUT: | ||
| await asyncio.sleep(0.5) | ||
| # Any scheduler push after we started counts as healthy. | ||
| if self.async_llm.last_receive_tstamp > tic: | ||
| return sglang_scheduler_pb2.HealthCheckResponse( | ||
| healthy=True, | ||
| message="Health check passed", | ||
| ) | ||
| if task.done(): | ||
| return sglang_scheduler_pb2.HealthCheckResponse( | ||
| healthy=bool(task.result()), | ||
| message=( | ||
| "Health check passed" | ||
| if task.result() | ||
| else "Scheduler returned no output" | ||
| ), | ||
| ) | ||
| finally: | ||
| if not task.done(): | ||
| task.cancel() | ||
| # Best-effort cleanup: the probe rid shouldn't linger. | ||
| self.async_llm.abort_request(rid) |
There was a problem hiding this comment.
Cancelled probe task is never awaited — minor resource-hygiene concern.
task.cancel() only requests cancellation; without a subsequent await task (wrapped in suppress(asyncio.CancelledError, Exception)), the probe coroutine may still be mid-generate_request when we return. The subsequent self.async_llm.abort_request(rid) partially mitigates this by asking the scheduler to drop the rid, but the awaitable itself leaks to the event loop until GC. A short await-with-suppress cleans up reliably.
🛠️ Proposed fix
finally:
if not task.done():
task.cancel()
+ with contextlib.suppress(asyncio.CancelledError, Exception):
+ await task
# Best-effort cleanup: the probe rid shouldn't linger.
self.async_llm.abort_request(rid)(plus import contextlib at the top).
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py` around lines 338 -
361, The probe task created by asyncio.create_task(_drive_probe()) may be
cancelled but never awaited, leaking the coroutine; modify the finally block to,
after calling task.cancel() when not task.done(), await the task inside a
contextlib.suppress(asyncio.CancelledError, Exception) to ensure the
_drive_probe() coroutine is finished/cleaned up before proceeding, and add
import contextlib at the top; keep the existing
self.async_llm.abort_request(rid) call for best-effort scheduler cleanup.
| # ``log_metrics`` on the proto is a plain bool3 scalar — there's | ||
| # no unset/zero-default distinction. Leaving tokenspeed's default | ||
| # (True) in place matches SGLang's behaviour where the router | ||
| # never opts out of metrics at this layer. | ||
| ) |
There was a problem hiding this comment.
Typo in docstring: bool3 → bool.
✏️ Fix typo
- # ``log_metrics`` on the proto is a plain bool3 scalar — there's
+ # ``log_metrics`` on the proto is a plain bool scalar — there's📝 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.
| # ``log_metrics`` on the proto is a plain bool3 scalar — there's | |
| # no unset/zero-default distinction. Leaving tokenspeed's default | |
| # (True) in place matches SGLang's behaviour where the router | |
| # never opts out of metrics at this layer. | |
| ) | |
| # ``log_metrics`` on the proto is a plain bool scalar — there's | |
| # no unset/zero-default distinction. Leaving tokenspeed's default | |
| # (True) in place matches SGLang's behaviour where the router | |
| # never opts out of metrics at this layer. | |
| ) |
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py` around lines 609 -
613, Fix the docstring typo that incorrectly says "bool3" — update the comment
that documents the proto field `log_metrics` in servicer functions within
tokenspeed (the comment block in
grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py referring to
`log_metrics`) to read "bool" instead of "bool3" so the docstring accurately
describes the plain boolean scalar for `log_metrics`.
| def abort_request(self, rid: str) -> None: | ||
| self.aborted_rids.append(rid) | ||
| self.rid_to_state.pop(rid, None) | ||
|
|
||
| async def generate_request(self, obj): | ||
| # Record the request so tests can assert on what was forwarded. | ||
| rid = getattr(obj, "rid", None) or "no-rid" | ||
| self.rid_to_state[rid] = _FakeState() | ||
| if self.generate_fn is not None: | ||
| async for out in self.generate_fn(obj): | ||
| self.last_receive_tstamp = 9999.0 # anything > tic | ||
| yield out | ||
| return | ||
| for out in self.outputs: | ||
| self.last_receive_tstamp = 9999.0 | ||
| yield out | ||
| self.rid_to_state[rid].finished = True | ||
|
|
There was a problem hiding this comment.
🧩 Analysis chain
🏁 Script executed:
#!/bin/bash
# Confirm how the servicer expands rid for n>1 and whether FakeAsyncLLM's current
# scalar-rid assumption breaks when Generate is invoked with n>1.
rg -nP -C3 '\bobj\.rid\s*=\s*\[' grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py
rg -nP -C5 'async def generate_request' grpc_servicer/tests/test_tokenspeed_servicer.pyRepository: lightseekorg/smg
Length of output: 821
🏁 Script executed:
rg -n 'def test_cancel_aborts_all_n_children' grpc_servicer/tests/test_tokenspeed_servicer.pyRepository: lightseekorg/smg
Length of output: 112
🏁 Script executed:
# Find the test and get its full implementation
rg -n 'def test_cancel_aborts_all_n_children' -A 50 grpc_servicer/tests/test_tokenspeed_servicer.pyRepository: lightseekorg/smg
Length of output: 2061
🏁 Script executed:
# Check if _build_generate_req is used in the test and what sampling params it uses
rg -n '_build_generate_req|sampling_params' grpc_servicer/tests/test_tokenspeed_servicer.py | head -40Repository: lightseekorg/smg
Length of output: 987
🏁 Script executed:
# Check if test has any pytest markers indicating it's skipped/xfailed
rg -B5 'async def test_cancel_aborts_all_n_children' grpc_servicer/tests/test_tokenspeed_servicer.py | head -20Repository: lightseekorg/smg
Length of output: 209
🏁 Script executed:
# Check the current full implementation of FakeAsyncLLM to see if generate_request was fixed
rg -n 'class FakeAsyncLLM' -A 100 grpc_servicer/tests/test_tokenspeed_servicer.py | grep -A 20 'async def generate_request'Repository: lightseekorg/smg
Length of output: 807
🏁 Script executed:
# Check servicer.Generate to understand when generate_request is called
rg -n 'def Generate|async def Generate' -A 30 grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py | head -80Repository: lightseekorg/smg
Length of output: 1591
🏁 Script executed:
# Check for any TODO/FIXME/xfail markers related to this test or FakeAsyncLLM
rg -n 'TODO|FIXME|xfail' grpc_servicer/tests/test_tokenspeed_servicer.py | grep -E '(163|170|634|FakeAsyncLLM|rid_to_state)'Repository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
# Check entire FakeAsyncLLM class definition to see all its methods and state management
rg -n 'class FakeAsyncLLM' -A 150 grpc_servicer/tests/test_tokenspeed_servicer.py | head -180Repository: lightseekorg/smg
Length of output: 5821
🏁 Script executed:
# Check if there's a recent fix or commit that addresses this bug
git log --oneline --all -20 -- grpc_servicer/tests/test_tokenspeed_servicer.py | head -20Repository: lightseekorg/smg
Length of output: 118
🏁 Script executed:
# Check the abort_request flow in servicer to see if it handles list rids
rg -n 'def.*abort|abort_request|_cancel_request' -A 10 grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.pyRepository: lightseekorg/smg
Length of output: 3452
🏁 Script executed:
# Check if there's any special handling for list rid in servicer.Generate
rg -n 'async for output in self.async_llm.generate_request' -A 30 grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py | head -50Repository: lightseekorg/smg
Length of output: 2760
🏁 Script executed:
# Let's check if there's a comment or note in the test file about this being known/expected
rg -n 'test_cancel_aborts_all_n_children' -B 10 grpc_servicer/tests/test_tokenspeed_servicer.py | grep -E '(TODO|FIXME|BUG|xfail|skip)'Repository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
# Double-check the actual servicer code that handles list rids during abort cleanup
rg -n 'if isinstance(expanded_rid' -B 5 -A 10 grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.pyRepository: lightseekorg/smg
Length of output: 150
🏁 Script executed:
# Check the cleanup section in Generate where expanded_rid is handled
rg -n 'expanded_rid' -B 3 -A 15 grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.pyRepository: lightseekorg/smg
Length of output: 3012
🏁 Script executed:
# Verify the exact lines in the finally block that handle list rid cleanup
sed -n '206,220p' grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.pyRepository: lightseekorg/smg
Length of output: 807
🏁 Script executed:
# Check if the test file has any recent changes or if the test is actually running
git log --oneline --all -10 -- grpc_servicer/tests/test_tokenspeed_servicer.pyRepository: lightseekorg/smg
Length of output: 118
🏁 Script executed:
# Check if there's any xfail or skip decorator on test_cancel_aborts_all_n_children
sed -n '630,665p' grpc_servicer/tests/test_tokenspeed_servicer.pyRepository: lightseekorg/smg
Length of output: 1362
🏁 Script executed:
# Final check: see if the test file imports or uses xfail anywhere near this test
rg -n 'xfail|skip|mark\.' grpc_servicer/tests/test_tokenspeed_servicer.py | head -30Repository: lightseekorg/smg
Length of output: 871
🏁 Script executed:
# Check the commit that added this file to see if it was supposed to be WIP
git show 960d4eb --stat -- grpc_servicer/tests/test_tokenspeed_servicer.py 2>/dev/null | head -20Repository: lightseekorg/smg
Length of output: 1117
Fix FakeAsyncLLM.generate_request to handle list rid from n > 1 requests.
When sampling_params.n > 1, the servicer expands obj.rid to a list of per-choice IDs (e.g., ["rid-1-n0", "rid-1-n1", "rid-1-n2"]). In FakeAsyncLLM.generate_request (lines 169–170), the code assumes rid is a scalar and uses it as a dict key: self.rid_to_state[rid] = _FakeState(). For list rid, this raises TypeError: unhashable type: 'list' before the generator executes, causing test_cancel_aborts_all_n_children to hang at await started.wait() instead of exercising the cancel path.
🛠️ Suggested fix
async def generate_request(self, obj):
# Record the request so tests can assert on what was forwarded.
- rid = getattr(obj, "rid", None) or "no-rid"
- self.rid_to_state[rid] = _FakeState()
+ rid_attr = getattr(obj, "rid", None) or "no-rid"
+ rids = rid_attr if isinstance(rid_attr, list) else [rid_attr]
+ for r in rids:
+ self.rid_to_state[r] = _FakeState()
if self.generate_fn is not None:
async for out in self.generate_fn(obj):
self.last_receive_tstamp = 9999.0 # anything > tic
yield out
return
for out in self.outputs:
self.last_receive_tstamp = 9999.0
yield out
- self.rid_to_state[rid].finished = True
+ for r in rids:
+ self.rid_to_state[r].finished = True🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@grpc_servicer/tests/test_tokenspeed_servicer.py` around lines 163 - 180,
FakeAsyncLLM.generate_request currently assumes obj.rid is hashable and does
self.rid_to_state[rid] = _FakeState(), which fails when obj.rid is a list for
sampling_params.n > 1; change generate_request to detect if rid is a list and,
if so, iterate over each child id and set self.rid_to_state[child_rid] =
_FakeState() (and later mark each child_rid.finished = True) instead of using
the list as a dict key; keep existing behavior for scalar rid, and ensure
last_receive_tstamp and generate_fn/output yielding logic remain the same.
| echo "Building tokenspeed-kernel from source (this takes ~30 min the first time)..." | ||
| MAX_JOBS="${MAX_JOBS:-16}" FLASHINFER_CUDA_ARCH_LIST="9.0a 10.0a" \ | ||
| uv pip install tokenspeed-kernel/python/ --no-build-isolation | ||
| # Cache the built wheel — uv stores wheels under its cache, copy out. | ||
| mkdir -p "$WHEEL_CACHE" | ||
| python3 -c "import tokenspeed_kernel, os, shutil, glob; \ | ||
| d = os.path.dirname(tokenspeed_kernel.__file__); \ | ||
| site = os.path.dirname(d); \ | ||
| whls = glob.glob(os.path.join(site, 'tokenspeed_kernel-*.dist-info')); \ | ||
| print('kernel install dir:', whls)" || true | ||
| fi |
There was a problem hiding this comment.
🧩 Analysis chain
🌐 Web query:
Does uv buildsupport building wheels for a package with--no-build-isolation and CUDA compilation prerequisites?
💡 Result:
No, uv build does not support the --no-build-isolation flag. The official documentation for uv build (results 19,21,22,23) lists options like --sdist, --wheel, --build-constraints, --python, but no --no-build-isolation. This flag is supported in installation commands like uv pip install and uv add (results 2,3,4,14), where it allows using the current environment for packages needing pre-installed dependencies like Torch for CUDA packages (results 1,34). For uv build, which acts as a PEP 517 build frontend invoking the backend (results 20,21), build isolation is handled differently; uv prefers augmenting isolated environments via [tool.uv.extra-build-dependencies] for CUDA prereqs like matching Torch versions (result 1). The uv build backend (uv_build) supports only pure Python wheels, requiring alternative backends (e.g., setuptools) for C/CUDA extensions (result 20). Users must ensure CUDA toolkit/prereqs are system-installed before running uv build on the host machine.
Citations:
- 1: https://docs.astral.sh/uv/guides/package/
- 2: https://docs.astral.sh/uv/concepts/build-backend/
- 3: https://docs.astral.sh/uv/concepts/projects/build/
- 4: https://docs.astral.sh/uv/reference/cli/
- 5: https://www.mintlify.com/astral-sh/uv/cli/build
- 6: Can
uvmakeno-build-isolationinstallation easier? astral-sh/uv#15248 - 7: UV does not appear to respect --no-build-isolation flag or pyproject.toml build dependencies astral-sh/uv#14043
- 8: Add environment variable for
--no-build-isolationastral-sh/uv#3309 - 9:
--no-build-isolationhas no effect inuv addastral-sh/uv#6367 - 10: Add support for
--no-build-isolationastral-sh/uv#2258 - 11: Support "--no-build-isolation" when installing packages astral-sh/uv#1715
🏁 Script executed:
# Read the script file to see context around the problematic lines
cat -n scripts/ci_install_tokenspeed.sh | sed -n '75,105p'Repository: lightseekorg/smg
Length of output: 1799
Wheel cache is not actually populated — 30-minute kernel build will run on every CI invocation.
The comment on line 89 claims the built wheel is cached, but the python3 -c block only prints diagnostic info; it never copies a .whl into $WHEEL_CACHE. The uv pip install tokenspeed-kernel/python/ installs the package in-place and doesn't leave a standalone wheel in a discoverable location, so the find "$WHEEL_CACHE" probe on line 81 will keep returning empty on every subsequent run.
Note: The suggested fix in the original comment is incorrect. uv build does not support the --no-build-isolation flag, so the proposed uv build ... --no-build-isolation command will fail. To fix this, either drop the caching logic entirely, or after the in-place install, manually locate and copy the built wheel from the uv cache directory to $WHEEL_CACHE for reuse. Given the documented ~30 min compile cost, restoring functional caching is important.
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@scripts/ci_install_tokenspeed.sh` around lines 86 - 96, The CI is not
actually populating the wheel cache: the uv pip install
tokenspeed-kernel/python/ call installs in-place and the python3 -c block that
inspects tokenspeed_kernel only prints diagnostic info and never copies a built
.whl into $WHEEL_CACHE, so every run rebuilds the kernel; fix by after the
install locating the built wheel inside uv's cache (or the pip/uv wheel cache
directory) and copying that .whl into $WHEEL_CACHE (or alternatively remove the
broken caching logic), updating the script sections that reference WHEEL_CACHE
and the python3 -c diagnostics/tokenspeed_kernel lookup to perform a concrete
copy of the discovered wheel into $WHEEL_CACHE so subsequent CI runs reuse the
cached wheel.
Summary
Introduces TokenSpeed as a first-class worker alongside SGLang, vLLM, TRT-LLM and MLX — with a dedicated
tokenspeed.grpc.scheduler.TokenSpeedSchedulerproto/service so the two backends can diverge freely over time (per @SimoLin's review note on the first iteration).Python servicer —
grpc_servicer/smg_grpc_servicer/tokenspeed/tokenspeed.AsyncLLM+ scheduler subprocessGenerate,Embed,Abort,HealthCheck,GetModelInfo,GetServerInfo,GetLoads,GetTokenizer,SubscribeKvEvents) using shared SGLang messages via proto import (import sglang_scheduler.proto)tokenspeed.grpc.scheduler.TokenSpeedSchedulerserver.pymirrors the SGLang servicer's lifecycle (AsyncLLM launch → gRPC bind → background warmup → SIGTERM/SIGINT)pyproject.tomladds atokenspeedoptional-dep groupgrpc_servicer/tests/cover launcher, servicer and health paths — all passingRust router integration
crates/grpc_client/proto/tokenspeed_scheduler.proto(+ mirror underpython/smg_grpc_proto/proto/) for the dedicated servicebuild.rscompiles it into a separateOUT_DIRwithextern_path(".sglang.grpc.scheduler", "crate::sglang_scheduler::proto")so messages are reused, not duplicatedTokenSpeedSchedulerClientincrates/grpc_client/src/tokenspeed_scheduler.rs— same surface asSglangSchedulerClient;build_grpc_sampling_params_*are promoted topub(crate)and delegatedAbortOnDropStreamrefactored onto a client-agnosticAbortDispatcher = Arc<dyn Fn(String) + Send + Sync>closure so both clients share itprotocols::worker::RuntimeTypegains aTokenSpeedvariant (Display/FromStr)GrpcClientenum, reachability probe, detect_backend ordering, harmony request building, multimodal dispatch, and embedding request building all grow the new arm (exhaustive match-checked)E2E infra
Runtime.TOKENSPEEDadded toe2e_test/infra/constants.pye2e_test/infra/worker.pygrows_build_tokenspeed_grpc_cmd(launches viapython -m smg_grpc_servicer.tokenspeed, forcing--grammar-backend xgrammar)e2e_test/infra/model_specs.pyregistersQwen/Qwen3-4BandQwen/Qwen3-30B-A3Btest_function_calling.pyruns againsttokenspeed; suite temporarily pinsQwen3-30B-A3B+qwenparser because TokenSpeed does not supportLlamaForCausalLMyetCI
.github/actions/setup-tokenspeed/action.yml+scripts/ci_install_tokenspeed.she2e-gpu-job.yml/pr-test-rust.ymlgain the new engine matrix entryE2E evidence
Ran the full
TestOpenAIServerFunctionCallingsuite onQwen/Qwen3-30B-A3Bwith a real TokenSpeed worker on a GCP B200 node, driven through the Rust router:tokenspeedvllm(same model, same suite)Rust router correctly labels the worker as
runtime_type: \"tokenspeed\"(native probe via the dedicated service — no marker hack).The 4 failures reproduce identically on both backends; root cause is the SMG gateway's
tool_choice=required/specificconstraint translation layer, not this change. Non-regression.Test plan
cargo checkacross the workspace — clean (all backends exhaustive)pytest grpc_servicer/tests/— 47 passedSupersedes #1351 (earlier iteration that reused SGLang's proto). Will close #1351 once this lands / gets the first review pass.
Summary by CodeRabbit
Release Notes
New Features
Chores