feat(grpc_servicer): add TokenSpeed servicer (Part 2/3) - #1464
Conversation
📝 WalkthroughWalkthroughAdds TokenSpeed backend support: project metadata and package init, CLI entrypoint, engine launcher, health servicer, gRPC server orchestration with warmup, full TokenSpeedSchedulerServicer RPC implementations (Generate, HealthCheck, Abort, GetModelInfo, GetServerInfo, GetLoads), helpers, and graceful shutdown. ChangesTokenSpeed gRPC Servicer Implementation
Sequence Diagram(s)sequenceDiagram
participant Client
participant TokenSpeedSchedulerServicer
participant AsyncLLM
participant SchedulerProcess
participant HealthServicer
Client->>TokenSpeedSchedulerServicer: GenerateRequest / HealthCheck / Abort / GetModelInfo / GetLoads
TokenSpeedSchedulerServicer->>AsyncLLM: submit engine request / probe / abort / get_load
AsyncLLM->>SchedulerProcess: communicate with TokenSpeed subprocess (scheduling & generation)
SchedulerProcess-->>AsyncLLM: generation frames / status / load
AsyncLLM-->>TokenSpeedSchedulerServicer: generation frames / completion / status
TokenSpeedSchedulerServicer-->>Client: stream GenerateResponse or unary response
TokenSpeedSchedulerServicer->>HealthServicer: set_serving() / set_not_serving()
HealthServicer-->>Client: HealthCheck response (Check/Watch)
Estimated code review effort🎯 4 (Complex) | ⏱️ ~60 minutes Possibly related PRs
Suggested labels
Suggested reviewers
🚥 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)
Comment |
There was a problem hiding this comment.
Clean, thorough implementation. Reviewed all 11 files: servicer (Generate streaming/non-streaming/n>1, HealthCheck, Abort, GetModelInfo, GetServerInfo, GetLoads), health servicer, server lifecycle, scheduler launcher, CLI entrypoint, and 57 unit tests.
Key things verified:
- n>1 cancel sweep: CancelledError handler and Abort RPC both correctly walk the expanded
{rid}-n{i}children, preventing orphaned GPU work - Chat-template prefix strip:
_generated_output_idscorrectly slices to the lastcompletion_tokenstokens, removing the Llama-3 assistant header that broke tool-call parsing - Stop-token trim + no_stop_trim: Trailing matched stop is stripped from
output_idsunlessno_stop_trimis set;matched_token_idstill rides in the proto field - Logprob alignment: Cumulative-to-delta slicing in
_convert_output_logprobs_to_protocorrectly handles streaming chunks and stop-token-stripped frames - HasField for optional scalars:
temperature=0(greedy) is correctly forwarded via presence-tracking rather than truthy checks - Warmup lifecycle: Synchronous gRPC client on daemon thread with proper channel cleanup; health stays NOT_SERVING until a complete frame is received
- Graceful shutdown: Drain loop with timeout, then
kill_process_tree(include_parent=False)to reap scheduler children without self-terminating
No bugs, no security concerns, no silent fallbacks. LGTM.
There was a problem hiding this comment.
Code Review
This pull request implements a gRPC servicer for the TokenSpeed inference engine, including health monitoring, subprocess management, and request handling for generation and metadata. Feedback identifies a need to handle zero completion tokens to prevent chat template prefix leakage and recommends updating the health servicer's Watch method to support server-streaming for full gRPC protocol compliance.
| if isinstance(completion, int) and 0 < completion <= len(raw): | ||
| token_ids = raw[-completion:] | ||
| else: | ||
| token_ids = raw |
There was a problem hiding this comment.
When completion_tokens is 0, the current logic falls back to returning the entire raw list of token IDs. Since raw often contains chat template prefix tokens, this fallback will leak those prefix tokens into the response. If completion_tokens is 0, an empty list should be returned. Additionally, ensure that in streaming token generation, the completion_tokens count is reported cumulatively for the entire request to ensure accurate progress reporting.
| if isinstance(completion, int) and 0 < completion <= len(raw): | |
| token_ids = raw[-completion:] | |
| else: | |
| token_ids = raw | |
| if isinstance(completion, int): | |
| token_ids = raw[-completion:] if completion > 0 else [] | |
| else: | |
| token_ids = raw |
References
- In a streaming token generation API, response chunks should report a cumulative count of completion_tokens for the entire request, not just the tokens in the current chunk, to ensure accurate progress reporting.
| async def Watch( | ||
| self, | ||
| request: health_pb2.HealthCheckRequest, | ||
| context: grpc.aio.ServicerContext, | ||
| ) -> AsyncIterator[health_pb2.HealthCheckResponse]: | ||
| # K8s probes use Check, not Watch — we emit the current status once. | ||
| yield await self.Check(request, context) |
There was a problem hiding this comment.
The Watch method implementation does not comply with the gRPC Health Checking Protocol (v1). The protocol requires Watch to be a server-streaming RPC that stays open and yields the current status whenever it changes. The current implementation yields once and then terminates the stream, which may cause issues with clients (like service meshes or load balancers) that rely on the streaming behavior of Watch to track backend health in real-time.
52cda9d to
c3c3c02
Compare
| finish_reason = "stop" | ||
| matched_kwargs: dict[str, Any] = {} | ||
| if reason_dict: | ||
| kind = reason_dict.get("type") | ||
| if kind == "length": | ||
| finish_reason = "length" | ||
| elif kind == "abort": | ||
| finish_reason = "abort" |
There was a problem hiding this comment.
🟡 Nit: The finish_reason mapping is an incomplete allowlist — only "length" and "abort" are recognized; every other type (including any future TokenSpeed additions like "cancelled") silently falls back to "stop". This means the gRPC path would silently misreport a new finish reason while the HTTP path handles it correctly, creating a subtle divergence between the two serving paths.
Consider logging a warning for unrecognized types so this doesn't fail silently:
| finish_reason = "stop" | |
| matched_kwargs: dict[str, Any] = {} | |
| if reason_dict: | |
| kind = reason_dict.get("type") | |
| if kind == "length": | |
| finish_reason = "length" | |
| elif kind == "abort": | |
| finish_reason = "abort" | |
| finish_reason = "stop" | |
| matched_kwargs: dict[str, Any] = {} | |
| if reason_dict: | |
| kind = reason_dict.get("type") | |
| if kind == "length": | |
| finish_reason = "length" | |
| elif kind == "abort": | |
| finish_reason = "abort" | |
| elif kind and kind != "stop": | |
| logger.warning("Unrecognized finish_reason type %r; defaulting to 'stop'", kind) |
| return reason | ||
| to_json = getattr(reason, "to_json", None) | ||
| if callable(to_json): | ||
| result = to_json() |
There was a problem hiding this comment.
🟡 Nit: Removing the try/except wrapper around to_json() changes error-routing semantics. The caller at line 191 catches ValueError and maps it to StatusCode.INVALID_ARGUMENT (user input error). If to_json() internally raises a ValueError, it will now be misclassified as bad user input rather than an internal server error. The previous code deliberately wrapped all to_json() failures in TypeError to guarantee they'd fall through to except Exception → StatusCode.INTERNAL.
The risk is low (a well-behaved to_json() shouldn't raise ValueError), but the original wrapper existed specifically to defend against this mismatch.
788933a to
d16d38f
Compare
8057d10 to
656f1c2
Compare
d16d38f to
3f8983a
Compare
|
|
||
| Mirrors smg_grpc_servicer.vllm / smg_grpc_servicer.sglang. Wraps TokenSpeed's | ||
| AsyncLLM (main-process async frontend) behind the SGLang gRPC service so the | ||
| existing Rust router (which auto-detects the SGLang proto) can route traffic | ||
| to TokenSpeed without needing a new client. | ||
| """ |
There was a problem hiding this comment.
🟡 Nit: This docstring is stale — it describes the opposite of what the implementation does. The servicer does NOT wrap behind "the SGLang gRPC service"; it uses its own tokenspeed.grpc.scheduler.TokenSpeedScheduler proto. The Rust router does NOT "auto-detect the SGLang proto"; DetectBackendStep identifies TokenSpeed natively from the service name. And there IS a new Rust client (TokenSpeedSchedulerClient).
| Mirrors smg_grpc_servicer.vllm / smg_grpc_servicer.sglang. Wraps TokenSpeed's | |
| AsyncLLM (main-process async frontend) behind the SGLang gRPC service so the | |
| existing Rust router (which auto-detects the SGLang proto) can route traffic | |
| to TokenSpeed without needing a new client. | |
| """ | |
| """TokenSpeed gRPC servicer implementation. | |
| Exposes TokenSpeed's AsyncLLM over the dedicated | |
| ``tokenspeed.grpc.scheduler.TokenSpeedScheduler`` gRPC service. | |
| The Rust gateway's ``DetectBackendStep`` identifies TokenSpeed workers | |
| natively from the service name. | |
| """ |
2ecbbb9 to
6bb18d2
Compare
| load_outputs = await asyncio.wait_for( | ||
| self.async_llm.get_load(), timeout=HEALTH_CHECK_TIMEOUT | ||
| ) | ||
| except TimeoutError: |
There was a problem hiding this comment.
🔴 Important: except TimeoutError catches builtins.TimeoutError (subclass of OSError), but asyncio.wait_for raises asyncio.TimeoutError which on Python 3.10 is a separate class inheriting from Exception, not from builtins.TimeoutError. Since pyproject.toml declares requires-python = ">=3.10", this handler is dead code on 3.10 — the timeout falls through to the except Exception block below and reports StatusCode.INTERNAL instead of DEADLINE_EXCEEDED.
asyncio.TimeoutError became an alias of builtins.TimeoutError only in Python 3.11 (bpo-45098).
| except TimeoutError: | |
| except (TimeoutError, asyncio.TimeoutError): |
This catches both the builtin and the asyncio variant, working correctly on 3.10+. On 3.11+ it's redundant but harmless.
8583d04 to
a812f5c
Compare
| wrapped = structural_tag_for_reasoning_json_schema( | ||
| reasoning_parser, json.loads(params.json_schema) | ||
| ) | ||
| except ImportError: | ||
| wrapped = None |
There was a problem hiding this comment.
🟡 Nit: json.loads(params.json_schema) can raise json.JSONDecodeError (a ValueError subclass) if the client sends a malformed schema string, but the except only catches ImportError. This means malformed JSON blows up here when a reasoning parser is configured, while without a parser the same bad string silently passes through as out["json_schema"].
The inconsistency is minor (the caller's except ValueError handler would produce a reasonable INVALID_ARGUMENT gRPC status), but catching JSONDecodeError alongside ImportError would make the fallback path uniform:
| wrapped = structural_tag_for_reasoning_json_schema( | |
| reasoning_parser, json.loads(params.json_schema) | |
| ) | |
| except ImportError: | |
| wrapped = None | |
| wrapped = structural_tag_for_reasoning_json_schema( | |
| reasoning_parser, json.loads(params.json_schema) | |
| ) | |
| except (ImportError, json.JSONDecodeError): | |
| wrapped = None |
a188c7a to
4b3ba12
Compare
8a3e651 to
c976a54
Compare
c976a54 to
923b929
Compare
| name = "smg-grpc-servicer" | ||
| version = "0.5.3" | ||
| description = "SMG gRPC servicer implementations for LLM inference engines (vLLM, SGLang, MLX)" | ||
| version = "0.5.2" |
There was a problem hiding this comment.
🔴 Important: This rolls the package version back from 0.5.3 (on main) to 0.5.2. If merged as-is, pip install --upgrade smg-grpc-servicer will consider the published 0.5.3 newer and skip this release entirely, so the TokenSpeed servicer would never be picked up by existing installs.
The base dependency floor is also lowered (smg-grpc-proto>=0.4.7 → >=0.4.6). If the tokenspeed_scheduler_pb2 proto was only added in 0.4.7, this would allow installing against a proto version that lacks it, causing an ImportError at runtime.
Likely a stale branch — rebasing on current main and bumping to 0.5.4 should fix both.
| # Needed for models with YaRN/other RoPE variants whose _freqs arrays are | ||
| # excluded from parameters() due to underscore prefix. See _eval_all_module_arrays. | ||
| _eval_all_module_arrays(model) | ||
| logger.info("Model loaded successfully") |
There was a problem hiding this comment.
🟡 Nit: This PR removes _eval_all_module_arrays(model) which was an explicit workaround for a documented cross-stream materialization bug with YarnRoPE's _freqs arrays (the RuntimeError: There is no Stream(gpu, 1) in current thread crash). The architectural fix at line 109-112 (constructing BatchGenerator on the generation thread) mitigates the same class of bug, so the removal may be safe — but _eval_all_module_arrays was a belt-and-suspenders guard for any lazy mx.array that slips through nn.Module.parameters().
If this was removed because mlx-lm>=0.22.0 now evaluates all arrays during load(), that's fine — but it's worth a quick note in the commit message or a comment explaining why the safety net is no longer needed, since the original docstring described a real production failure mode.
923b929 to
92edf51
Compare
bcff9b3 to
75371a4
Compare
There was a problem hiding this comment.
Actionable comments posted: 6
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (1)
grpc_servicer/smg_grpc_servicer/tokenspeed/health_servicer.py (1)
69-120:⚠️ Potential issue | 🟠 Major | ⚡ Quick winAdd explicit INTERNAL abort handling in health RPCs.
Check/Watchcurrently let unexpected exceptions escape; this bypasses the repository’s servicer error-handling contract and can produce inconsistent client-visible statuses.Proposed patch
async def Check( self, request: health_pb2.HealthCheckRequest, context: grpc.aio.ServicerContext, ) -> health_pb2.HealthCheckResponse: - service_name = request.service - logger.debug("Health check request for service=%r", service_name) + try: + service_name = request.service + logger.debug("Health check request for service=%r", service_name) - if self.async_llm.gracefully_exit: - return health_pb2.HealthCheckResponse(status=health_pb2.HealthCheckResponse.NOT_SERVING) + if self.async_llm.gracefully_exit: + return health_pb2.HealthCheckResponse( + status=health_pb2.HealthCheckResponse.NOT_SERVING + ) - if service_name == self.OVERALL_SERVER: - return health_pb2.HealthCheckResponse( - status=self._serving_status.get( - self.OVERALL_SERVER, health_pb2.HealthCheckResponse.NOT_SERVING + if service_name == self.OVERALL_SERVER: + return health_pb2.HealthCheckResponse( + status=self._serving_status.get( + self.OVERALL_SERVER, health_pb2.HealthCheckResponse.NOT_SERVING + ) ) - ) - if service_name == self.TOKENSPEED_SERVICE: - base = self._serving_status.get( - self.TOKENSPEED_SERVICE, health_pb2.HealthCheckResponse.NOT_SERVING - ) - if base != health_pb2.HealthCheckResponse.SERVING: - return health_pb2.HealthCheckResponse(status=base) + if service_name == self.TOKENSPEED_SERVICE: + base = self._serving_status.get( + self.TOKENSPEED_SERVICE, health_pb2.HealthCheckResponse.NOT_SERVING + ) + if base != health_pb2.HealthCheckResponse.SERVING: + return health_pb2.HealthCheckResponse(status=base) - # Scheduler-stuck check: pending work but no recent output. - time_since_last_receive = time.time() - self.async_llm.last_receive_tstamp - pending = len(self.async_llm.rid_to_state) - if time_since_last_receive > STUCK_SCHEDULER_THRESHOLD_SEC and pending > 0: - logger.warning( - "Scheduler appears stuck: %.1fs since last receive, %d pending requests", - time_since_last_receive, - pending, - ) - return health_pb2.HealthCheckResponse( - status=health_pb2.HealthCheckResponse.NOT_SERVING - ) + # Scheduler-stuck check: pending work but no recent output. + time_since_last_receive = time.time() - self.async_llm.last_receive_tstamp + pending = len(self.async_llm.rid_to_state) + if time_since_last_receive > STUCK_SCHEDULER_THRESHOLD_SEC and pending > 0: + logger.warning( + "Scheduler appears stuck: %.1fs since last receive, %d pending requests", + time_since_last_receive, + pending, + ) + return health_pb2.HealthCheckResponse( + status=health_pb2.HealthCheckResponse.NOT_SERVING + ) - return health_pb2.HealthCheckResponse(status=health_pb2.HealthCheckResponse.SERVING) + return health_pb2.HealthCheckResponse(status=health_pb2.HealthCheckResponse.SERVING) - context.set_code(grpc.StatusCode.NOT_FOUND) - context.set_details(f"Unknown service: {service_name}") - return health_pb2.HealthCheckResponse(status=health_pb2.HealthCheckResponse.SERVICE_UNKNOWN) + context.set_code(grpc.StatusCode.NOT_FOUND) + context.set_details(f"Unknown service: {service_name}") + return health_pb2.HealthCheckResponse(status=health_pb2.HealthCheckResponse.SERVICE_UNKNOWN) + except Exception as e: # noqa: BLE001 + logger.exception("Health.Check failed") + await context.abort(grpc.StatusCode.INTERNAL, str(e)) async def Watch( self, request: health_pb2.HealthCheckRequest, context: grpc.aio.ServicerContext, ) -> AsyncIterator[health_pb2.HealthCheckResponse]: - # K8s probes use Check, not Watch — we emit the current status once. - yield await self.Check(request, context) + try: + # K8s probes use Check, not Watch — we emit the current status once. + yield await self.Check(request, context) + except Exception as e: # noqa: BLE001 + logger.exception("Health.Watch failed") + await context.abort(grpc.StatusCode.INTERNAL, str(e))Based on learnings: Enforce the repository-wide convention for gRPC INTERNAL errors in all servicer implementations (log exception, then
context.abort(grpc.StatusCode.INTERNAL, str(e))).🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@grpc_servicer/smg_grpc_servicer/tokenspeed/health_servicer.py` around lines 69 - 120, Wrap the body of both Check and Watch in try/except that catches Exception, uses logger.exception(...) to record the error, and then calls context.abort(grpc.StatusCode.INTERNAL, str(e)) so unexpected errors follow the repo convention; in Watch, catch exceptions raised by await self.Check(...) as well as any local errors. Ensure you keep the existing logic (service_name handling, scheduler-stuck check using self.async_llm.last_receive_tstamp and STUCK_SCHEDULER_THRESHOLD_SEC, and use of self._serving_status) but terminate on exceptions via context.abort(grpc.StatusCode.INTERNAL, str(e)).
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@grpc_servicer/smg_grpc_servicer/tokenspeed/__main__.py`:
- Around line 17-18: Either declare uvloop in the project's dependencies
(pyproject.toml) or guard its import and usage in the tokenspeed CLI entrypoint:
wrap the top-level "import uvloop" in a try/except ImportError and, if missing,
fall back to asyncio.DefaultEventLoopPolicy() instead of calling
uvloop.install(); if present call uvloop.install() as before. Update the code
paths that reference uvloop (the module import and the place where the loop
policy is set before invoking prepare_server_args / the CLI entrypoint) so the
CLI "python -m smg_grpc_servicer.tokenspeed" does not raise ModuleNotFoundError,
or alternatively add uvloop to dependencies in pyproject.toml.
In `@grpc_servicer/smg_grpc_servicer/tokenspeed/server.py`:
- Around line 120-123: The warmup dial uses server_args.host directly which
fails when the server was bound to a wildcard like "0.0.0.0" or "::"; before
building grpc_url and dialing in the warmup logic, normalize wildcard bind hosts
(e.g., if server_args.host is "0.0.0.0" or "::" treat it as the loopback address
"127.0.0.1" (or "[::1]" for IPv6) and then construct grpc_url =
f"{normalized_host}:{server_args.port}" so that the grpc.insecure_channel call
targets a reachable loopback address for readiness probes (update the code paths
that reference server_args.host, grpc_url, or the warmup dial function).
In `@grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py`:
- Around line 67-81: The code currently returns dict-shaped finish reasons
without validation; ensure that when reason is a dict you validate it contains a
'type' key with a string value before returning (e.g., if isinstance(reason,
dict): if 'type' not in reason or not isinstance(reason['type'], str): raise
TypeError(...)); keep the existing to_json handling (callable to_json branch)
but if its result is a dict also validate it contains a string 'type' key before
returning; raise a clear TypeError referencing finish_reason when validation
fails so malformed dicts (like {"matched": 42}) are rejected instead of being
accepted and misclassified downstream.
- Around line 807-832: The function builds token_logprobs/token_ids from
delta_token but top_logprobs (top_proto) can be shorter when raw_top is shorter;
ensure top_proto is padded to the same length as delta_token by computing
expected_len = len(delta_token) and appending
tokenspeed_scheduler_pb2.TopLogProbs() placeholders until len(top_proto) ==
expected_len so the returned tokenspeed_scheduler_pb2.OutputLogProbs has aligned
token_logprobs, token_ids and top_logprobs; update the code around
raw_top/delta_top/top_proto (and the return) to perform this padding after the
loop that constructs top_proto.
- Around line 137-149: The loop over output in the non-streaming n>1 branch
currently yields each choice (via self._complete_response) before checking later
items for aborts, which can emit partial successes then abort via context.abort;
instead, first scan or buffer the entire output list to detect any abort (use
_finish_reason_to_dict and _abort_status_code on each item.get("meta_info",
{}).get("finish_reason")), and if an abort is found call await
context.abort(...) immediately; only after confirming no aborts, iterate the
buffered items and yield self._complete_response for each choice (preserving
index handling with item.get("index", idx) and no_stop_trim).
- Around line 119-123: The current try/except around
self._build_generate_req(request) only catches ValueError so other exceptions
escape and bypass server logging and the standard INTERNAL abort; modify this
block in servicer.py to catch Exception (in addition to ValueError) from
_build_generate_req, log the unexpected exception (using the servicer's logger
or existing logging pattern), and call await
context.abort(grpc.StatusCode.INTERNAL, str(e)) for non-ValueError failures so
all unexpected errors are mapped to INTERNAL; keep the existing INVALID_ARGUMENT
handling for ValueError and ensure you reference _build_generate_req,
context.abort, and grpc.StatusCode.INTERNAL when applying the change.
---
Outside diff comments:
In `@grpc_servicer/smg_grpc_servicer/tokenspeed/health_servicer.py`:
- Around line 69-120: Wrap the body of both Check and Watch in try/except that
catches Exception, uses logger.exception(...) to record the error, and then
calls context.abort(grpc.StatusCode.INTERNAL, str(e)) so unexpected errors
follow the repo convention; in Watch, catch exceptions raised by await
self.Check(...) as well as any local errors. Ensure you keep the existing logic
(service_name handling, scheduler-stuck check using
self.async_llm.last_receive_tstamp and STUCK_SCHEDULER_THRESHOLD_SEC, and use of
self._serving_status) but terminate on exceptions via
context.abort(grpc.StatusCode.INTERNAL, str(e)).
🪄 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: 3b09eccc-7e53-4338-9813-2ee8dde16d21
📒 Files selected for processing (7)
grpc_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.py
| if reason is None or isinstance(reason, dict): | ||
| return reason | ||
| to_json = getattr(reason, "to_json", None) | ||
| if callable(to_json): | ||
| result = to_json() | ||
| if isinstance(result, dict): | ||
| return result | ||
| raise TypeError( | ||
| f"finish_reason {type(reason).__name__!r}.to_json() returned " | ||
| f"{type(result).__name__!r}; expected dict with at least 'type'." | ||
| ) | ||
| raise TypeError( | ||
| f"Unknown finish_reason shape {type(reason).__name__!r}; expected " | ||
| f"a dict or an object with a to_json() method." | ||
| ) |
There was a problem hiding this comment.
Validate dict-shaped finish reasons before accepting them.
The docstring says downstream expects at least {"type": ...}, but plain dicts are returned unchecked here. A malformed dict like {"matched": 42} then falls through to _complete_response() and gets silently reported as "stop".
Proposed fix
- if reason is None or isinstance(reason, dict):
+ if reason is None:
return reason
+ if isinstance(reason, dict):
+ if not isinstance(reason.get("type"), str):
+ raise TypeError("finish_reason dict must include a string 'type'")
+ return reason📝 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.
| if reason is None or isinstance(reason, dict): | |
| return reason | |
| to_json = getattr(reason, "to_json", None) | |
| if callable(to_json): | |
| result = to_json() | |
| if isinstance(result, dict): | |
| return result | |
| raise TypeError( | |
| f"finish_reason {type(reason).__name__!r}.to_json() returned " | |
| f"{type(result).__name__!r}; expected dict with at least 'type'." | |
| ) | |
| raise TypeError( | |
| f"Unknown finish_reason shape {type(reason).__name__!r}; expected " | |
| f"a dict or an object with a to_json() method." | |
| ) | |
| if reason is None: | |
| return reason | |
| if isinstance(reason, dict): | |
| if not isinstance(reason.get("type"), str): | |
| raise TypeError("finish_reason dict must include a string 'type'") | |
| return reason | |
| to_json = getattr(reason, "to_json", None) | |
| if callable(to_json): | |
| result = to_json() | |
| if isinstance(result, dict): | |
| return result | |
| raise TypeError( | |
| f"finish_reason {type(reason).__name__!r}.to_json() returned " | |
| f"{type(result).__name__!r}; expected dict with at least 'type'." | |
| ) | |
| raise TypeError( | |
| f"Unknown finish_reason shape {type(reason).__name__!r}; expected " | |
| f"a dict or an object with a to_json() method." | |
| ) |
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py` around lines 67 - 81,
The code currently returns dict-shaped finish reasons without validation; ensure
that when reason is a dict you validate it contains a 'type' key with a string
value before returning (e.g., if isinstance(reason, dict): if 'type' not in
reason or not isinstance(reason['type'], str): raise TypeError(...)); keep the
existing to_json handling (callable to_json branch) but if its result is a dict
also validate it contains a string 'type' key before returning; raise a clear
TypeError referencing finish_reason when validation fails so malformed dicts
(like {"matched": 42}) are rejected instead of being accepted and misclassified
downstream.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 75371a4370
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
|
|
||
| async def GetLoads( | ||
| self, | ||
| _request: tokenspeed_scheduler_pb2.GetLoadsRequest, |
There was a problem hiding this comment.
Honor the requested DP rank when reporting loads
When a caller sets GetLoadsRequest.dp_rank (the proto exposes this filter, and the SGLang servicer applies it), this implementation ignores the request object and returns every rank. Rank-specific load polling or admin clients will therefore receive unfiltered loads and an aggregate/dp_rank_count for all ranks instead of the requested rank; please read request.HasField("dp_rank") and filter load_outputs before building the response.
Useful? React with 👍 / 👎.
| ) | ||
|
|
||
|
|
||
| class TokenSpeedSchedulerServicer(tokenspeed_scheduler_pb2_grpc.TokenSpeedSchedulerServicer): |
There was a problem hiding this comment.
🔴 Important: The PR description claims "57 unit tests" in grpc_servicer/tests/test_tokenspeed_*.py, but neither the tests/ directory nor any test files exist in this push. The test plan section references passing tests and even names specific test cases (test_cancel_calls_abort_request, test_cancel_aborts_all_n_children, test_abort_sweeps_n_children), but they're entirely absent from the diff.
This servicer has non-trivial logic that warrants test coverage — the n>1 rid expansion, abort sweep, finish-reason mapping, logprob slicing, and the completion_tokens-based prefix stripping in _generated_output_ids are all paths where regressions would be subtle and hard to catch without tests.
Were the tests accidentally dropped in the force push?
75371a4 to
5dae995
Compare
There was a problem hiding this comment.
Actionable comments posted: 4
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@grpc_servicer/smg_grpc_servicer/tokenspeed/server.py`:
- Around line 65-67: The gRPC bind result from server.add_insecure_port is not
being checked so a bind failure could be silent; update the startup in the
TokenSpeed gRPC servicer to capture the return value of
server.add_insecure_port(listen_addr), verify it is non-zero, and if it is zero
log an error via logger.error (including listen_addr) and exit/raise to stop
startup instead of proceeding to logger.info("TokenSpeed gRPC server listening
on %s", listen_addr); use the same check pattern around server.add_insecure_port
and logger.info that is used in the other server startup code (i.e., validate
the int result, handle failure, otherwise log success).
- Around line 120-123: The warmup URL construction can yield unbracketed IPv6
addresses (e.g., "::1:PORT") when mapping server_args.host (especially "::") to
loopback; update the logic around warmup_host and grpc_url to map "[::]" to
"::1" and to wrap any IPv6 literal in brackets before appending the port. Locate
the warmup_host variable and grpc_url construction and implement: map
server_args.host value "[::]" or "::" to "::1", detect if warmup_host contains
":" (an IPv6 literal) and format grpc_url as
f"[{warmup_host}]:{server_args.port}", otherwise keep the existing
"{warmup_host}:{server_args.port}" behavior so IPv4 and hostnames remain
unchanged.
In `@grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py`:
- Around line 544-552: The shutdown path currently calls kill_process_tree from
shutdown() and allows exceptions from kill_process_tree to propagate after
you've already flipped gracefully_exit/health state; make shutdown fully
best-effort by wrapping the kill_process_tree(os.getpid(), include_parent=False)
call in a try/except Exception block (inside the same scope where you import
tokenspeed.runtime.utils.process.kill_process_tree), catching any exception and
calling logger.exception(...) to record the failure but not re-raising it so
shutdown can complete without unwinding and potentially leaving children behind;
keep the existing import try/except as-is and only change the actual
kill_process_tree invocation to be non-fatal on error.
- Around line 554-596: The code must reject missing/empty request_id before
minting scheduler rids: in _build_generate_req, validate that
request.HasField("request_id") and that request.request_id is not an empty
string (or raise ValueError) before you compute obj.rid; update the logic that
sets obj.rid (the n==1 and n>1 branches that currently use request.request_id)
so it only uses a non-empty request.request_id, otherwise fail fast—this
prevents different requests from getting identical "" or "-nX" rids and avoids
collisions in rid_to_state/abort/cancel paths.
🪄 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: d72d917d-e362-4f86-b81b-c8c5e4934e7f
📒 Files selected for processing (7)
grpc_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.py
| listen_addr = f"{server_args.host}:{server_args.port}" | ||
| server.add_insecure_port(listen_addr) | ||
| logger.info("TokenSpeed gRPC server listening on %s", listen_addr) |
There was a problem hiding this comment.
🧩 Analysis chain
🏁 Script executed:
#!/bin/bash
# Verify bind call exists and return value is currently unchecked.
nl -ba grpc_servicer/smg_grpc_servicer/tokenspeed/server.py | sed -n '60,72p'Repository: lightseekorg/smg
Length of output: 101
🏁 Script executed:
cat -n grpc_servicer/smg_grpc_servicer/tokenspeed/server.py | sed -n '60,75p'Repository: lightseekorg/smg
Length of output: 722
🌐 Web query:
gRPC Python add_insecure_port return value semantics
💡 Result:
In gRPC Python, the add_insecure_port method returns an integer representing the port on which the server will accept RPC requests [1][2][3]. If you specify a port of 0 in the address (or do not specify a port at all), the gRPC runtime will dynamically choose an available port, and the method will return that assigned port number [1][2][3]. Historically, the behavior of this method when port binding fails has been a point of ambiguity in the gRPC Python library [4][5]. While the official documentation specifies that it returns an integer port [1][2], it has been noted in past issues that the method may return 0 upon failure rather than raising an exception [5]. Users have reported that in some cases, failure to bind a port might result in error messages being printed to the logs (often from the underlying C++ core) without a corresponding Python exception being raised, making it difficult to handle such failures programmatically [4]. Consequently, it is recommended to verify the returned port value or implement robust error handling around server startup to ensure the port was bound successfully [4][5].
Citations:
- 1: https://grpc.github.io/grpc/python/grpc.html
- 2: https://grpc.github.io/grpc/python/_modules/grpc/aio/_base_server.html
- 3: https://docs.ros.org/en/api/grpc/html/classgrpc_1_1aio_1_1__base__server_1_1Server.html
- 4: Python exception when port addition fails grpc/grpc#23315
- 5: Document and test failing to bind a port on a server grpc/grpc#13096
🏁 Script executed:
fd -e txt -e cfg -e ini -e toml -e lock | grep -E '(requirements|setup|pyproject|poetry|pipfile)' | head -20Repository: lightseekorg/smg
Length of output: 215
🏁 Script executed:
rg "add_insecure_port" --type pyRepository: lightseekorg/smg
Length of output: 351
🏁 Script executed:
rg "test.*server" grpc_servicer/smg_grpc_servicer/tokenspeed/ --type py -iRepository: lightseekorg/smg
Length of output: 186
Check gRPC port binding result to prevent silent server startup failure.
add_insecure_port() can fail silently—the gRPC library returns 0 on bind failure instead of raising an exception. Continuing startup without validating the return value leaves the process running but unreachable. The codebase already follows the correct pattern in mlx/server.py; apply it here.
Proposed patch
listen_addr = f"{server_args.host}:{server_args.port}"
- server.add_insecure_port(listen_addr)
+ bound_port = server.add_insecure_port(listen_addr)
+ if bound_port == 0:
+ raise RuntimeError(f"Failed to bind TokenSpeed gRPC server on {listen_addr}")
logger.info("TokenSpeed gRPC server listening on %s", listen_addr)📝 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.
| listen_addr = f"{server_args.host}:{server_args.port}" | |
| server.add_insecure_port(listen_addr) | |
| logger.info("TokenSpeed gRPC server listening on %s", listen_addr) | |
| listen_addr = f"{server_args.host}:{server_args.port}" | |
| bound_port = server.add_insecure_port(listen_addr) | |
| if bound_port == 0: | |
| raise RuntimeError(f"Failed to bind TokenSpeed gRPC server on {listen_addr}") | |
| logger.info("TokenSpeed gRPC server listening on %s", listen_addr) |
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@grpc_servicer/smg_grpc_servicer/tokenspeed/server.py` around lines 65 - 67,
The gRPC bind result from server.add_insecure_port is not being checked so a
bind failure could be silent; update the startup in the TokenSpeed gRPC servicer
to capture the return value of server.add_insecure_port(listen_addr), verify it
is non-zero, and if it is zero log an error via logger.error (including
listen_addr) and exit/raise to stop startup instead of proceeding to
logger.info("TokenSpeed gRPC server listening on %s", listen_addr); use the same
check pattern around server.add_insecure_port and logger.info that is used in
the other server startup code (i.e., validate the int result, handle failure,
otherwise log success).
| # Wildcard bind hosts aren't routable as destinations; dial loopback instead. | ||
| warmup_host = {"0.0.0.0": "127.0.0.1", "::": "::1"}.get(server_args.host, server_args.host) | ||
| grpc_url = f"{warmup_host}:{server_args.port}" | ||
| channel = grpc.insecure_channel( |
There was a problem hiding this comment.
🧩 Analysis chain
🏁 Script executed:
#!/bin/bash
# Verify the current warmup host mapping and URL formatting.
nl -ba grpc_servicer/smg_grpc_servicer/tokenspeed/server.py | sed -n '118,126p'Repository: lightseekorg/smg
Length of output: 101
🏁 Script executed:
cat -n grpc_servicer/smg_grpc_servicer/tokenspeed/server.py | sed -n '115,130p'Repository: lightseekorg/smg
Length of output: 909
🌐 Web query:
IPv6 address gRPC insecure_channel URL format bracketing requirement
💡 Result:
When using an IPv6 address in a gRPC target string, you must enclose the address in literal square brackets ([ and ]) if a port is specified [1][2]. This bracketing is required to disambiguate the colon separator used for the port from the colons contained within the IPv6 address itself [1][2]. For example, to connect to an IPv6 address with a port, the format is: [2001:db8::1]:50051 If you omit the brackets when a port is present, the gRPC client will typically fail to parse the address correctly, often resulting in an invalid address error [3]. This requirement applies generally across gRPC implementations when providing a raw IPv6 address as a target [3][4]. If you are using the explicit ipv6 URI scheme (supported primarily in C-core-based implementations), the format follows the same rule [1][2]: ipv6:[2001:db8::1]:50051 Always ensure your target string follows this bracketed format to avoid connection failures [3].
Citations:
- 1: https://grpc.github.io/grpc/core/md_doc_naming.html
- 2: https://github.com/grpc/grpc/blob/master/doc/naming.md
- 3: https://oneuptime.com/blog/post/2026-03-20-debug-grpc-ipv6-issues/view
- 4: https://oneuptime.com/blog/post/2026-03-20-grpc-ipv6-python/view
🏁 Script executed:
rg -n "grpc\.insecure_channel|grpc_url|f\".+:.+\{.*port" grpc_servicer/smg_grpc_servicer/tokenspeed/ -A 3Repository: lightseekorg/smg
Length of output: 659
🏁 Script executed:
rg -n "server_args.host|argparse|add_argument.*host" grpc_servicer/smg_grpc_servicer/tokenspeed/server.py -B 2 -A 2Repository: lightseekorg/smg
Length of output: 601
🏁 Script executed:
rg -n "add_argument.*host|ArgumentParser|argparse" grpc_servicer/smg_grpc_servicer/tokenspeed/server.py -B 2 -A 2Repository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
rg -n "def.*server_args|server_args\s*=" grpc_servicer/smg_grpc_servicer/tokenspeed/server.py | head -5Repository: lightseekorg/smg
Length of output: 154
🏁 Script executed:
rg -n "class ServerArgs|`@dataclass`.*ServerArgs" grpc_servicer/smg_grpc_servicer/tokenspeed/ -A 5Repository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
head -30 grpc_servicer/smg_grpc_servicer/tokenspeed/server.pyRepository: lightseekorg/smg
Length of output: 1007
🏁 Script executed:
find . -path "*/tokenspeed/runtime/utils/server_args.py" -o -path "*/tokenspeed/**/server_args.py" | head -5Repository: lightseekorg/smg
Length of output: 42
Bracket IPv6 addresses in gRPC warmup URL.
The current warmup URL construction can produce an unbracketed IPv6 address (e.g., ::1:PORT) when mapping :: to loopback. gRPC requires IPv6 addresses with ports to be bracketed: [::1]:PORT.
Add bracketing for IPv6 addresses and include the [::] → ::1 mapping:
Proposed patch
- warmup_host = {"0.0.0.0": "127.0.0.1", "::": "::1"}.get(server_args.host, server_args.host)
+ warmup_host = {"0.0.0.0": "127.0.0.1", "::": "::1", "[::]": "::1"}.get(
+ server_args.host, server_args.host
+ )
+ if ":" in warmup_host and not warmup_host.startswith("["):
+ warmup_host = f"[{warmup_host}]"
grpc_url = f"{warmup_host}:{server_args.port}"🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@grpc_servicer/smg_grpc_servicer/tokenspeed/server.py` around lines 120 - 123,
The warmup URL construction can yield unbracketed IPv6 addresses (e.g.,
"::1:PORT") when mapping server_args.host (especially "::") to loopback; update
the logic around warmup_host and grpc_url to map "[::]" to "::1" and to wrap any
IPv6 literal in brackets before appending the port. Locate the warmup_host
variable and grpc_url construction and implement: map server_args.host value
"[::]" or "::" to "::1", detect if warmup_host contains ":" (an IPv6 literal)
and format grpc_url as f"[{warmup_host}]:{server_args.port}", otherwise keep the
existing "{warmup_host}:{server_args.port}" behavior so IPv4 and hostnames
remain unchanged.
| try: | ||
| from tokenspeed.runtime.utils.process import kill_process_tree | ||
| except ImportError: | ||
| logger.exception( | ||
| "Could not import tokenspeed.runtime.utils.process.kill_process_tree; " | ||
| "scheduler subprocesses may be orphaned" | ||
| ) | ||
| return | ||
| kill_process_tree(os.getpid(), include_parent=False) |
There was a problem hiding this comment.
Keep shutdown best-effort if process-tree cleanup fails.
shutdown() is documented as best-effort, but kill_process_tree(...) is still allowed to raise. If that happens, server shutdown can unwind after you've already flipped gracefully_exit/health state, and you still risk leaving children behind.
Suggested fix
try:
from tokenspeed.runtime.utils.process import kill_process_tree
except ImportError:
logger.exception(
"Could not import tokenspeed.runtime.utils.process.kill_process_tree; "
"scheduler subprocesses may be orphaned"
)
return
- kill_process_tree(os.getpid(), include_parent=False)
+ try:
+ kill_process_tree(os.getpid(), include_parent=False)
+ except Exception: # noqa: BLE001
+ logger.exception(
+ "Failed to kill TokenSpeed subprocess tree; "
+ "scheduler subprocesses may be orphaned"
+ )🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py` around lines 544 -
552, The shutdown path currently calls kill_process_tree from shutdown() and
allows exceptions from kill_process_tree to propagate after you've already
flipped gracefully_exit/health state; make shutdown fully best-effort by
wrapping the kill_process_tree(os.getpid(), include_parent=False) call in a
try/except Exception block (inside the same scope where you import
tokenspeed.runtime.utils.process.kill_process_tree), catching any exception and
calling logger.exception(...) to record the failure but not re-raising it so
shutdown can complete without unwinding and potentially leaving children behind;
keep the existing import try/except as-is and only change the actual
kill_process_tree invocation to be non-fatal on error.
| def _build_generate_req(self, request: tokenspeed_scheduler_pb2.GenerateRequest): | ||
| """Translate proto GenerateRequest → TokenSpeed GenerateReqInput. | ||
|
|
||
| Keeps the router's pre-tokenized inputs intact (``input_ids`` set, | ||
| ``text`` left blank) so the TokenSpeed InputProcessor skips its own | ||
| tokenizer pass. | ||
| """ | ||
| if not request.HasField("tokenized"): | ||
| raise ValueError("GenerateRequest.tokenized is required") | ||
|
|
||
| input_ids = list(request.tokenized.input_ids) | ||
| if not input_ids: | ||
| raise ValueError("GenerateRequest.tokenized.input_ids is empty") | ||
|
|
||
| sampling = self._sampling_params_from_proto( | ||
| request.sampling_params, | ||
| reasoning_parser=getattr(self.server_args, "reasoning_parser", None), | ||
| ) | ||
|
|
||
| GenerateReqInput = _lazy_generate_req_input() | ||
| obj = GenerateReqInput( | ||
| input_ids=input_ids, | ||
| sampling_params=sampling, | ||
| stream=bool(request.stream), | ||
| return_logprob=bool(request.return_logprob), | ||
| # presence-tracking distinguishes "client omitted" (→ ``-1`` = | ||
| # no input logprobs) from explicit ``0`` (start at position 0). | ||
| logprob_start_len=( | ||
| request.logprob_start_len if request.HasField("logprob_start_len") else -1 | ||
| ), | ||
| top_logprobs_num=int(request.top_logprobs_num or 0), | ||
| token_ids_logprob=( | ||
| list(request.token_ids_logprob) if request.token_ids_logprob else None | ||
| ), | ||
| ) | ||
| # ``normalize_batch_and_arguments`` asserts ``rid`` is a list when | ||
| # n>1; expand to deterministic per-choice rids so the assert holds. | ||
| n = sampling.get("n", 1) or 1 | ||
| if n > 1: | ||
| obj.rid = [f"{request.request_id}-n{i}" for i in range(n)] | ||
| else: | ||
| obj.rid = request.request_id | ||
|
|
There was a problem hiding this comment.
Reject empty request_ids before minting scheduler rids.
An omitted protobuf string comes through as "" here. The n == 1 path then reuses the same empty rid for every caller, and the n > 1 path generates the same -n0, -n1, … child rids across requests. That can make rid_to_state, abort sweeping, and cancellation cleanup collide across unrelated requests.
Suggested fix
def _build_generate_req(self, request: tokenspeed_scheduler_pb2.GenerateRequest):
"""Translate proto GenerateRequest → TokenSpeed GenerateReqInput.
Keeps the router's pre-tokenized inputs intact (``input_ids`` set,
``text`` left blank) so the TokenSpeed InputProcessor skips its own
tokenizer pass.
"""
+ if not request.request_id:
+ raise ValueError("GenerateRequest.request_id is required")
+
if not request.HasField("tokenized"):
raise ValueError("GenerateRequest.tokenized is required")🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py` around lines 554 -
596, The code must reject missing/empty request_id before minting scheduler
rids: in _build_generate_req, validate that request.HasField("request_id") and
that request.request_id is not an empty string (or raise ValueError) before you
compute obj.rid; update the logic that sets obj.rid (the n==1 and n>1 branches
that currently use request.request_id) so it only uses a non-empty
request.request_id, otherwise fail fast—this prevents different requests from
getting identical "" or "-nX" rids and avoids collisions in
rid_to_state/abort/cancel paths.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 5dae995103
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
| if n > 1: | ||
| obj.rid = [f"{request.request_id}-n{i}" for i in range(n)] | ||
| else: |
There was a problem hiding this comment.
Track the regenerated n>1 request ids
For n > 1, these deterministic {request_id}-n{i} ids are not the scheduler ids that end up in rid_to_state: TokenSpeed's parallel-sampling path takes obj[i] and then prepare_prefix_warmup / prepare_parallel_sampling_replica call regenerate_rid() before _send_one_request. As a result, the CancelledError cleanup and the Abort RPC sweep only ids that were discarded, so cancelled or explicitly aborted n>1 requests can keep running on the GPU; track the actual replica rids after TokenSpeed regenerates them, or abort by the stored user_rid metadata.
Useful? React with 👍 / 👎.
| wanted = int(request.dp_rank) | ||
| load_outputs = [lo for lo in load_outputs if int(lo.dp_rank) == wanted] | ||
|
|
||
| page_size = int(getattr(self.async_llm.server_args, "page_size", 1) or 1) |
There was a problem hiding this comment.
Use TokenSpeed's block size for load tokens
TokenSpeed's ServerArgs exposes the KV page size as block_size, not page_size, so on a real AsyncLLM this falls back to 1. Any nonzero num_pages from get_load() is therefore reported as pages instead of tokens, making num_used_tokens and token_usage too low by the block size (typically 64) and causing routing/admin load metrics to miss actual KV pressure.
Useful? React with 👍 / 👎.
Adds CI plumbing and e2e test coverage for the TokenSpeed gRPC backend added in PRs #1351 + #1464. - ``scripts/ci_install_tokenspeed.sh``: pinned source install of TokenSpeed (engine + kernel + scheduler), with a wheel cache for subsequent runs. - ``.github/actions/setup-tokenspeed/action.yml`` + workflow hooks: install TokenSpeed in the GPU e2e job. - ``.github/workflows/e2e-gpu-job.yml`` / ``pr-test-rust.yml``: enable the ``tokenspeed`` engine matrix for chat-completion, completion, responses, router, and chat tool-calling tests. - ``e2e_test/infra/constants.py``: ``Runtime.TOKENSPEED`` enum + ``is_tokenspeed()`` helper + label map entry. - ``e2e_test/infra/worker.py``: ``_build_tokenspeed_grpc_cmd`` — launches ``python -m smg_grpc_servicer.tokenspeed`` with the right CLI flags (``--model``, ``--grammar-backend xgrammar``, ``--enable-output-logprobs``). - ``e2e_test/infra/model_specs.py``: per-model knobs + ``Qwen/Qwen3-4B`` for ``TestEnableThinking``. - E2E test files: add ``tokenspeed`` to the ``engine`` markers where the backend supports the surface under test. Verification: ``ruff check`` and ``ruff format --check`` clean on all touched Python files. Behavior is exercised by the GPU e2e matrix in CI. Signed-off-by: key4ng <rukeyang@gmail.com>
Adds the Python-side TokenSpeed gRPC servicer that backs the TokenSpeed client added in the parent PR. Wraps TokenSpeed's ``AsyncLLM`` (the main-process async frontend) and exposes the ``tokenspeed.grpc.scheduler.TokenSpeedScheduler`` service. Modules added under ``grpc_servicer/smg_grpc_servicer/tokenspeed/``: - ``servicer.py``: ``TokenSpeedSchedulerServicer`` with all RPCs (``Generate``, ``Abort``, ``HealthCheck``, ``GetModelInfo``, ``GetServerInfo``, ``GetLoads``). - ``health_servicer.py``: dedicated health-check side channel. - ``scheduler_launcher.py`` / ``server.py`` / ``__main__.py`` / ``__init__.py``: entrypoint + process supervision. Includes a structured-tag wrapper for ``json_schema`` requests when the engine is running with a reasoning parser (e.g. ``gpt-oss`` → harmony), so the grammar only activates inside the response channel — mirrors TokenSpeed's HTTP ``serving_chat.py`` behaviour. Falls back to raw ``json_schema`` for parsers without an xgrammar mapping. Other notes: - ``GetModelInfo`` populates the renamed ``default_sampling_params_json`` proto field (matches the parent PR's proto change). - ``ServerArgs``: handles both the legacy ``model_path`` / ``tokenizer_path`` attribute names and the bare ``model`` / ``tokenizer`` ones, since tokenspeed has renamed them. - Small follow-ups in ``mlx/server.py`` and ``sglang/servicer.py`` (defensive fallback handling). Test coverage: 60 tests across ``test_tokenspeed_servicer.py`` and ``test_tokenspeed_health_servicer.py`` (sampling-params conversion, structural-tag wrapping, abort behaviour, finish-reason mapping, load RPC, health pulsing). Verification: ``ruff check grpc_servicer/`` and ``ruff format --check grpc_servicer/`` both clean; pytest passes 60/60. Signed-off-by: key4ng <rukeyang@gmail.com>
5dae995 to
a177630
Compare
Adds CI plumbing and e2e test coverage for the TokenSpeed gRPC backend added in PRs #1351 + #1464. - ``scripts/ci_install_tokenspeed.sh``: pinned source install of TokenSpeed (engine + kernel + scheduler), with a wheel cache for subsequent runs. - ``.github/actions/setup-tokenspeed/action.yml`` + workflow hooks: install TokenSpeed in the GPU e2e job. - ``.github/workflows/e2e-gpu-job.yml`` / ``pr-test-rust.yml``: enable the ``tokenspeed`` engine matrix for chat-completion, completion, responses, router, and chat tool-calling tests. - ``e2e_test/infra/constants.py``: ``Runtime.TOKENSPEED`` enum + ``is_tokenspeed()`` helper + label map entry. - ``e2e_test/infra/worker.py``: ``_build_tokenspeed_grpc_cmd`` — launches ``python -m smg_grpc_servicer.tokenspeed`` with the right CLI flags (``--model``, ``--grammar-backend xgrammar``, ``--enable-output-logprobs``). - ``e2e_test/infra/model_specs.py``: per-model knobs + ``Qwen/Qwen3-4B`` for ``TestEnableThinking``. - E2E test files: add ``tokenspeed`` to the ``engine`` markers where the backend supports the surface under test. Verification: ``ruff check`` and ``ruff format --check`` clean on all touched Python files. Behavior is exercised by the GPU e2e matrix in CI. Signed-off-by: key4ng <rukeyang@gmail.com>
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: a17763055d
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
| token_ids = self._generated_output_ids(output, reason_dict, no_stop_trim=no_stop_trim) | ||
| return tokenspeed_scheduler_pb2.GenerateResponse( | ||
| request_id=rid, | ||
| chunk=tokenspeed_scheduler_pb2.GenerateStreamChunk( | ||
| token_ids=token_ids, |
There was a problem hiding this comment.
Emit deltas for raw-token streaming
When the gRPC worker is launched with --skip-tokenizer-init (a natural configuration here because requests are already tokenized), TokenSpeed's raw-token output path only emits disjoint segments if --stream-output is also set; its ServerArgs.stream_output default is false. This chunk path forwards the returned output_ids directly as each streaming chunk, while the router treats every chunk as newly generated tokens, so clients can receive duplicated prefixes on every frame. Please either force stream_output for the gRPC server or track the last emitted offset per request before building chunks.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Actionable comments posted: 1
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (1)
grpc_servicer/pyproject.toml (1)
32-38:⚠️ Potential issue | 🟠 Major | ⚡ Quick winAdd
tokenspeedto[project.optional-dependencies].The CLI entrypoint at
__main__.pyline 17 imports fromtokenspeed.runtime.utils.server_args, but notokenspeedoptional-dependency is declared. This is inconsistent with the pattern established forvllm,sglang, andmlx, and will causeModuleNotFoundErrorat runtime if users install the package without separately installing TokenSpeed.📦 Proposed fix
mlx = ["smg-grpc-proto>=0.4.7", "mlx>=0.22.0", "mlx-lm>=0.22.0"] +tokenspeed = ["tokenspeed>=0.1.0"] # Adjust version floor as needed🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@grpc_servicer/pyproject.toml` around lines 32 - 38, The package imports tokenspeed.runtime.utils.server_args from __main__.py (line ~17) but tokenspeed is not listed under [project.optional-dependencies]; add a tokenspeed entry to the [project.optional-dependencies] table in pyproject.toml (following the pattern used for vllm/sglang/mlx) so that users who install the optional CLI get the tokenspeed package (e.g., add "tokenspeed" with an appropriate minimum version like "tokenspeed>=<min-version>"). Ensure the token name is exactly tokenspeed so the import tokenspeed.runtime.utils.server_args resolves at runtime.
♻️ Duplicate comments (4)
grpc_servicer/smg_grpc_servicer/tokenspeed/__main__.py (1)
17-17:⚠️ Potential issue | 🟠 MajorVerify TokenSpeed installation requirement is documented.
This import requires
tokenspeedto be installed, butpyproject.tomldoes not declare it in[project.optional-dependencies]. Users will encounterModuleNotFoundErrorunless they install TokenSpeed separately.#!/bin/bash # Verify that tokenspeed is declared as an optional dependency echo "=== Checking pyproject.toml for tokenspeed optional-dependency ===" rg -n 'tokenspeed.*=' grpc_servicer/pyproject.toml || echo "tokenspeed optional-dependency not found" echo "" echo "=== Current optional-dependencies section ===" rg -A 10 '\[project\.optional-dependencies\]' grpc_servicer/pyproject.toml🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@grpc_servicer/smg_grpc_servicer/tokenspeed/__main__.py` at line 17, The import "from tokenspeed.runtime.utils.server_args import prepare_server_args" in __main__.py requires tokenspeed to be installable; add tokenspeed to the project's optional dependencies and document it: update grpc_servicer/pyproject.toml under [project.optional-dependencies] to include a "tokenspeed" entry (with a version constraint or extras as appropriate) and add a short note in the README or installation docs explaining that tokenspeed is an optional dependency required for the tokenspeed.__main__ entrypoint so users won't get ModuleNotFoundError when running prepare_server_args.grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py (3)
552-552:⚠️ Potential issue | 🟠 Major | ⚡ Quick winWrap
kill_process_treecall in try/except for best-effort shutdown.This issue was flagged in a previous review and remains unaddressed. The import is guarded but the actual call can still raise, failing shutdown after
gracefully_exitand health state have already been flipped.Proposed fix
- kill_process_tree(os.getpid(), include_parent=False) + try: + kill_process_tree(os.getpid(), include_parent=False) + except Exception: # noqa: BLE001 + logger.exception( + "Failed to kill TokenSpeed subprocess tree; " + "scheduler subprocesses may be orphaned" + )🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py` at line 552, The call to kill_process_tree(os.getpid(), include_parent=False) must be made best-effort and not allowed to raise during shutdown: wrap that call in a try/except Exception block (preserving include_parent=False) inside the same shutdown path where gracefully_exit and health state are toggled, catch any exception and log it (use the existing logger/processLogger in this module, e.g., logger.error or logger.exception) so shutdown continues even if kill_process_tree fails.
67-68:⚠️ Potential issue | 🟡 Minor | ⚡ Quick winValidate dict-shaped finish reasons before accepting them.
This issue was flagged in a previous review and remains unaddressed. A dict without a
"type"key (e.g.,{"matched": 42}) passes through and is silently treated as"stop"in_complete_response.Proposed fix
- if reason is None or isinstance(reason, dict): + if reason is None: return reason + if isinstance(reason, dict): + if not isinstance(reason.get("type"), str): + raise TypeError("finish_reason dict must include a string 'type'") + return reason🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py` around lines 67 - 68, The current guard "if reason is None or isinstance(reason, dict): return reason" lets dicts missing a "type" key slip through and later be treated as "stop" in _complete_response; change this to validate dict-shaped reasons: if reason is None return None; if isinstance(reason, dict) ensure "type" is present and is a string (and optionally that its value is one of the expected finish kinds) — otherwise return None (or a normalized safe value) instead of the raw dict; update the conditional around the "reason" check to explicitly verify "type" (and its allowed values) before returning the dict so _complete_response only receives well-formed finish reasons.
561-566:⚠️ Potential issue | 🟠 Major | ⚡ Quick winReject empty
request_idbefore minting scheduler rids.This issue was flagged in a previous review and remains unaddressed. An omitted protobuf string arrives as
"". Forn==1, all callers reuse the same empty rid; forn>1, the same"-n0","-n1", … child rids are minted across requests. This causes collisions inrid_to_state, abort sweeping, and cancellation cleanup.Proposed fix
def _build_generate_req(self, request: tokenspeed_scheduler_pb2.GenerateRequest): """Translate proto GenerateRequest → TokenSpeed GenerateReqInput. Keeps the router's pre-tokenized inputs intact (``input_ids`` set, ``text`` left blank) so the TokenSpeed InputProcessor skips its own tokenizer pass. """ + if not request.request_id: + raise ValueError("GenerateRequest.request_id is required") + if not request.HasField("tokenized"): raise ValueError("GenerateRequest.tokenized is required")🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py` around lines 561 - 566, The code currently checks tokenized input but does not reject an empty request.request_id, which later causes recycled/minted scheduler rids to collide; add an explicit validation of request.request_id (the protobuf string field on the incoming GenerateRequest) immediately after validating tokenized/input_ids and before any scheduler rid minting logic in the servicer method handling generation (the function where `request` is processed in servicer.py, near the `if not request.HasField("tokenized")` block). If request.request_id is None or the empty string (""), raise a ValueError with a clear message like "GenerateRequest.request_id is required and must be non-empty" to prevent downstream rid collisions in rid_to_state. Ensure this check runs before any code that constructs child rids (the scheduler/minting logic) so no empty base rid is used.
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py`:
- Line 364: The code directly accesses server_args.preferred_sampling_params
when constructing default_sampling_params_json which can raise AttributeError on
older ServerArgs versions; update the construction in the servicer (where
default_sampling_params_json is set) to safely guard access by checking for the
attribute first (e.g., hasattr(self.server_args, "preferred_sampling_params"))
and fall back to an empty string or
self.server_args.get("preferred_sampling_params", "") equivalent before using
it; ensure the change touches the same expression that builds
default_sampling_params_json so older ServerArgs without
preferred_sampling_params remain compatible.
---
Outside diff comments:
In `@grpc_servicer/pyproject.toml`:
- Around line 32-38: The package imports tokenspeed.runtime.utils.server_args
from __main__.py (line ~17) but tokenspeed is not listed under
[project.optional-dependencies]; add a tokenspeed entry to the
[project.optional-dependencies] table in pyproject.toml (following the pattern
used for vllm/sglang/mlx) so that users who install the optional CLI get the
tokenspeed package (e.g., add "tokenspeed" with an appropriate minimum version
like "tokenspeed>=<min-version>"). Ensure the token name is exactly tokenspeed
so the import tokenspeed.runtime.utils.server_args resolves at runtime.
---
Duplicate comments:
In `@grpc_servicer/smg_grpc_servicer/tokenspeed/__main__.py`:
- Line 17: The import "from tokenspeed.runtime.utils.server_args import
prepare_server_args" in __main__.py requires tokenspeed to be installable; add
tokenspeed to the project's optional dependencies and document it: update
grpc_servicer/pyproject.toml under [project.optional-dependencies] to include a
"tokenspeed" entry (with a version constraint or extras as appropriate) and add
a short note in the README or installation docs explaining that tokenspeed is an
optional dependency required for the tokenspeed.__main__ entrypoint so users
won't get ModuleNotFoundError when running prepare_server_args.
In `@grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py`:
- Line 552: The call to kill_process_tree(os.getpid(), include_parent=False)
must be made best-effort and not allowed to raise during shutdown: wrap that
call in a try/except Exception block (preserving include_parent=False) inside
the same shutdown path where gracefully_exit and health state are toggled, catch
any exception and log it (use the existing logger/processLogger in this module,
e.g., logger.error or logger.exception) so shutdown continues even if
kill_process_tree fails.
- Around line 67-68: The current guard "if reason is None or isinstance(reason,
dict): return reason" lets dicts missing a "type" key slip through and later be
treated as "stop" in _complete_response; change this to validate dict-shaped
reasons: if reason is None return None; if isinstance(reason, dict) ensure
"type" is present and is a string (and optionally that its value is one of the
expected finish kinds) — otherwise return None (or a normalized safe value)
instead of the raw dict; update the conditional around the "reason" check to
explicitly verify "type" (and its allowed values) before returning the dict so
_complete_response only receives well-formed finish reasons.
- Around line 561-566: The code currently checks tokenized input but does not
reject an empty request.request_id, which later causes recycled/minted scheduler
rids to collide; add an explicit validation of request.request_id (the protobuf
string field on the incoming GenerateRequest) immediately after validating
tokenized/input_ids and before any scheduler rid minting logic in the servicer
method handling generation (the function where `request` is processed in
servicer.py, near the `if not request.HasField("tokenized")` block). If
request.request_id is None or the empty string (""), raise a ValueError with a
clear message like "GenerateRequest.request_id is required and must be
non-empty" to prevent downstream rid collisions in rid_to_state. Ensure this
check runs before any code that constructs child rids (the scheduler/minting
logic) so no empty base rid is used.
🪄 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: 6860b821-77b2-4d5d-97c7-47075927c6e1
📒 Files selected for processing (7)
grpc_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.py
| return tokenspeed_scheduler_pb2.GetModelInfoResponse( | ||
| model_path=model_path, | ||
| tokenizer_path=tokenizer_path or "", | ||
| default_sampling_params_json=self.server_args.preferred_sampling_params or "", |
There was a problem hiding this comment.
Guard preferred_sampling_params access for version compatibility.
Lines 355-360 defensively handle renamed model/model_path fields across versions, but line 364 directly accesses preferred_sampling_params. If an older ServerArgs build lacks this field, this raises AttributeError.
Proposed fix
- default_sampling_params_json=self.server_args.preferred_sampling_params or "",
+ default_sampling_params_json=getattr(self.server_args, "preferred_sampling_params", "") or "",🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py` at line 364, The code
directly accesses server_args.preferred_sampling_params when constructing
default_sampling_params_json which can raise AttributeError on older ServerArgs
versions; update the construction in the servicer (where
default_sampling_params_json is set) to safely guard access by checking for the
attribute first (e.g., hasattr(self.server_args, "preferred_sampling_params"))
and fall back to an empty string or
self.server_args.get("preferred_sampling_params", "") equivalent before using
it; ensure the change touches the same expression that builds
default_sampling_params_json so older ServerArgs without
preferred_sampling_params remain compatible.
Adds CI plumbing and e2e test coverage for the TokenSpeed gRPC backend added in PRs #1351 + #1464. - ``scripts/ci_install_tokenspeed.sh``: pinned source install of TokenSpeed (engine + kernel + scheduler), with a wheel cache for subsequent runs. - ``.github/actions/setup-tokenspeed/action.yml`` + workflow hooks: install TokenSpeed in the GPU e2e job. - ``.github/workflows/e2e-gpu-job.yml`` / ``pr-test-rust.yml``: enable the ``tokenspeed`` engine matrix for chat-completion, completion, responses, router, and chat tool-calling tests. - ``e2e_test/infra/constants.py``: ``Runtime.TOKENSPEED`` enum + ``is_tokenspeed()`` helper + label map entry. - ``e2e_test/infra/worker.py``: ``_build_tokenspeed_grpc_cmd`` — launches ``python -m smg_grpc_servicer.tokenspeed`` with the right CLI flags (``--model``, ``--grammar-backend xgrammar``, ``--enable-output-logprobs``). - ``e2e_test/infra/model_specs.py``: per-model knobs + ``Qwen/Qwen3-4B`` for ``TestEnableThinking``. - E2E test files: add ``tokenspeed`` to the ``engine`` markers where the backend supports the surface under test. Verification: ``ruff check`` and ``ruff format --check`` clean on all touched Python files. Behavior is exercised by the GPU e2e matrix in CI. Signed-off-by: key4ng <rukeyang@gmail.com>
Adds CI plumbing and e2e test coverage for the TokenSpeed gRPC backend added in PRs #1351 + #1464. - ``scripts/ci_install_tokenspeed.sh``: pinned source install of TokenSpeed (engine + kernel + scheduler), with a wheel cache for subsequent runs. - ``.github/actions/setup-tokenspeed/action.yml`` + workflow hooks: install TokenSpeed in the GPU e2e job. - ``.github/workflows/e2e-gpu-job.yml`` / ``pr-test-rust.yml``: enable the ``tokenspeed`` engine matrix for chat-completion, completion, responses, router, and chat tool-calling tests. - ``e2e_test/infra/constants.py``: ``Runtime.TOKENSPEED`` enum + ``is_tokenspeed()`` helper + label map entry. - ``e2e_test/infra/worker.py``: ``_build_tokenspeed_grpc_cmd`` — launches ``python -m smg_grpc_servicer.tokenspeed`` with the right CLI flags (``--model``, ``--grammar-backend xgrammar``, ``--enable-output-logprobs``). - ``e2e_test/infra/model_specs.py``: per-model knobs + ``Qwen/Qwen3-4B`` for ``TestEnableThinking``. - E2E test files: add ``tokenspeed`` to the ``engine`` markers where the backend supports the surface under test. Verification: ``ruff check`` and ``ruff format --check`` clean on all touched Python files. Behavior is exercised by the GPU e2e matrix in CI. Signed-off-by: key4ng <rukeyang@gmail.com>
Adds CI plumbing and e2e test coverage for the TokenSpeed gRPC backend added in PRs #1351 + #1464. - ``scripts/ci_install_tokenspeed.sh``: pinned source install of TokenSpeed (engine + kernel + scheduler), with a wheel cache for subsequent runs. - ``.github/actions/setup-tokenspeed/action.yml`` + workflow hooks: install TokenSpeed in the GPU e2e job. - ``.github/workflows/e2e-gpu-job.yml`` / ``pr-test-rust.yml``: enable the ``tokenspeed`` engine matrix for chat-completion, completion, responses, router, and chat tool-calling tests. - ``e2e_test/infra/constants.py``: ``Runtime.TOKENSPEED`` enum + ``is_tokenspeed()`` helper + label map entry. - ``e2e_test/infra/worker.py``: ``_build_tokenspeed_grpc_cmd`` — launches ``python -m smg_grpc_servicer.tokenspeed`` with the right CLI flags (``--model``, ``--grammar-backend xgrammar``, ``--enable-output-logprobs``). - ``e2e_test/infra/model_specs.py``: per-model knobs + ``Qwen/Qwen3-4B`` for ``TestEnableThinking``. - E2E test files: add ``tokenspeed`` to the ``engine`` markers where the backend supports the surface under test. Verification: ``ruff check`` and ``ruff format --check`` clean on all touched Python files. Behavior is exercised by the GPU e2e matrix in CI. Signed-off-by: key4ng <rukeyang@gmail.com>
Description
Problem
PR #1351's Rust router can dial a TokenSpeed worker over the gRPC protocol it defines, but no worker speaks that protocol. We need a Python servicer that runs alongside a TokenSpeed scheduler process and serves those wire types.
Solution
A self-contained TokenSpeed servicer module under
grpc_servicer/smg_grpc_servicer/tokenspeed/, with cancellation handling for streaming/non-streaming, channel-close, andn>1paths.3-PR Stack
This is part 2 of 3 splitting the original #1351:
mainfeat/grpc-tokenspeed-servicerChanges
grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py— async scheduler servicer (Generate / HealthCheck / Abort / GetModelInfo / GetServerInfo / GetLoads), with cancellation that sweeps every{rid}-n{i}child rid expanded byn>1grpc_servicer/smg_grpc_servicer/tokenspeed/health_servicer.py— health-service bridge that flips SERVING / NOT_SERVING based on bounded-staleness scheduler liveness probesgrpc_servicer/smg_grpc_servicer/tokenspeed/scheduler_launcher.py— boots TokenSpeedAsyncLLMin-processgrpc_servicer/smg_grpc_servicer/tokenspeed/server.pyand__main__.py—python -m smg_grpc_servicer.tokenspeedentrypointGetLoadsreturns realAsyncLLM.get_load()metricsTest Plan
ruff check grpc_servicer/smg_grpc_servicer/tokenspeed/cleanruff format --check grpc_servicer/cleanChecklist
cargo +nightly fmtpasses (no Rust changes)cargo clippy --all-targets --all-features -- -D warningspasses (no Rust changes)Summary by CodeRabbit
New Features
Chore