Repository navigation
Conversation
…enerate responses Signed-off-by: key4ng <rukeyang@gmail.com>
…alue when present Signed-off-by: key4ng <rukeyang@gmail.com>
…e, advertises the RL control endpoint Signed-off-by: key4ng <rukeyang@gmail.com>
…ty keys into worker labels Signed-off-by: key4ng <rukeyang@gmail.com>
…transport Signed-off-by: key4ng <rukeyang@gmail.com>
…less of data transport Signed-off-by: key4ng <rukeyang@gmail.com>
… rl.control_url label Signed-off-by: key4ng <rukeyang@gmail.com>
…ity row, expose control_url in the Python client Signed-off-by: key4ng <rukeyang@gmail.com>
Signed-off-by: key4ng <rukeyang@gmail.com>
…rsion on gRPC generates Signed-off-by: key4ng <rukeyang@gmail.com>
Signed-off-by: key4ng <rukeyang@gmail.com>
…ol endpoint Signed-off-by: key4ng <rukeyang@gmail.com>
…p rl.control_url when RL is on Signed-off-by: key4ng <rukeyang@gmail.com>
…d launch guide Signed-off-by: key4ng <rukeyang@gmail.com>
…trol app Signed-off-by: key4ng <rukeyang@gmail.com>
… refit round-trip Signed-off-by: key4ng <rukeyang@gmail.com>
is_scheduler_paused is pre-existing on the engine and its signature is not pinned by this branch's spec, which writes it without an await. A synchronous implementation raised TypeError inside a broad except that silently reported is_paused=False on every discovery poll. Await the result only when it is actually awaitable. Signed-off-by: key4ng <rukeyang@gmail.com>
An operator PATCH can set rl.control_url to an empty or whitespace-only
string, which previously produced control_url: Some("") in discovery
and a 502 upstream_unreachable instead of the 422 no_control_endpoint
the situation deserves. Filter blank labels the same way capability.rs
and ProtoGenerateComplete::weight_version already treat empty as unset.
Also fix the base_url doc comment in the wire type, which still claimed
to be the address control calls are sent to; that is now control_url's
job, and this field's twin in crates/rl/src/view.rs was already fixed.
Signed-off-by: key4ng <rukeyang@gmail.com>
…sion The four /generate SSE sites (chunk and complete, both regular and logprobs paths) inlined chunk.weight_version().unwrap_or(&ctx.weight_version) instead of calling the shared precedence helper. The two agreed only because ProtoGenerate*::weight_version() already filters empty strings; a future change to the helper would not have reached the SSE path. Signed-off-by: key4ng <rukeyang@gmail.com>
_stamp_rl_control_labels retried every PATCH failure, including a permanent 401/403/404, once a second for the full 300s deadline, logging a warning each time. Give up on the target immediately for a 4xx status and keep the existing retry behavior for everything else. Also gate the stamping thread on connection_mode == "zmq" explicitly. This was previously safe only because TokenspeedWorkerLauncher.build_command raises for any other mode, so a gRPC launch died before the orchestrator reached the stamping code; make the coupling explicit instead of load-bearing-by-accident. The new tests exercise smg.serve, which needs the native smg.smg_rs extension; they were verified by careful reading and could not be run in this environment (ModuleNotFoundError: No module named 'smg.smg_rs'). Signed-off-by: key4ng <rukeyang@gmail.com>
crates/rl/README.md never stated that an rl.control_url label on an HTTP worker is ignored (it always controls through itself); one clause closes that. MockWorker::authorizations() claimed to be index-aligned with bodies(), but record_authorization is only called from the RL control handler while record_request has five call sites, so the alignment does not hold. Narrow the doc instead of changing the recorder's behavior. Signed-off-by: key4ng <rukeyang@gmail.com>
…tests Signed-off-by: key4ng <rukeyang@gmail.com>
…tatic row and docs The companion TokenSpeed branch's scheduler has no receive path for a disk or tensor weight update: update_weights_from_disk and update_weights_from_tensor both answer HTTP 501 and advertise rl.update_from=distributed. Only the trainer-driven NCCL broadcast (update_weights_from_distributed) works. Drop 'disk' from TokenSpeed's static capability row so pre-advertisement builds report only what they can actually do, and update the NOTES.md drift log, crates/rl/README.md, and docs/guides/rl-tokenspeed.md to point callers at the distributed refit path instead of refit_from_disk.py. Signed-off-by: key4ng <rukeyang@gmail.com>
… without taking the engine down Move test_refit_from_disk_is_visible_on_the_next_generate out of the shared mixin into TestRlControlPlaneSglang: TokenSpeed does not support a disk refit. Add TokenSpeed-only coverage that discovery advertises update_from=["distributed"] (engine-advertised, source=label), and that a disk refit fan-out comes back 207 with a single upstream_error/501 failure mentioning 'distributed', while the engine keeps serving afterward. Signed-off-by: key4ng <rukeyang@gmail.com>
TokenSpeed implements no disk or tensor refit, so the only way to refit it through SMG is slime's path: a trainer at rank 0 of a NCCL group broadcasting each weight into the engines. refit_from_trainer.py drives that end to end -- pause, per-worker init_weights_update_group with the right rank_offset, chunked update_weights_from_distributed fan-outs alongside the broadcasts, destroy, resume -- and checks the next /generate reports the new meta_info.weight_version. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: key4ng <rukeyang@gmail.com>
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: key4ng <rukeyang@gmail.com>
A TokenSpeed engine launched without an explicit parallelism flag leaves attn_tp_size unset, so discovery reports tp_size null and the example silently assumed 1. That is right for a TP-1 fleet and wrong for any other, and a wrong rank layout deadlocks the weight-update group instead of failing. Keep the fallback, name it in the warning, and add --tp-size to override it. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: key4ng <rukeyang@gmail.com>
Records acceptance item 1 as blocked on the engine, not on SMG: discovery, the per-worker proxy, the fan-outs, pause/resume and the group setup and teardown all passed, but TokenSpeed answers update_weights_from_distributed with HTTP 200 in 64 ms without entering the NCCL collective, so the trainer deadlocks in comm init and the engine loads its uninitialized receive buffers into the live model. Includes the py-spy stacks from both sides and an isolated three-rank reproduction that rules out the trainer-side code. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: key4ng <rukeyang@gmail.com>
The lead owns everything on GPUs 4-7 -- the slime-rl container, the TokenSpeed engines on ports 312xx and an SMG on 31100. Record that, and add a cleanup helper that stops only the gateway on 30100 and the engines on 30106/30107 and then reports the compute apps on GPUs 0-2 by UUID. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: key4ng <rukeyang@gmail.com>
The engine's world group is created with a bound device_id, so torch's _new_process_group_helper sets split_from and builds the weight-update communicator with ncclCommSplit off the engine's own size-1 comm instead of ncclCommInitRank against the trainer's unique id. A broadcast on a size-1 communicator is a local no-op, which is the 64 ms round trip and the uninitialized buffers. Replaces the three candidates with the mechanism, adds the resume order, and warns that the node's diagnostic edit must be reverted before the engine fix is synced over it. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: key4ng <rukeyang@gmail.com>
C1: the layout test mutated sys.path and sys.modules for the whole pytest process, so every later test importing smg.smg_rs failed. Scope both to a module-scoped fixture that snapshots and restores them, load the script once, and pin the invariant with a test. I2: teardown problems were printed and dropped, so a refit whose destroy_weights_update_group failed still exited 0 with a live group. They now reach the exit code, after the /generate result is printed. I3: --master-address is the trainer host as the engines see it, and the 127.0.0.1 default silently stalls the rendezvous for --timeout against remote engines. Document it in both guides and refuse to start when --smg is remote and --master-address is loopback. I4: the docs promised cleanup that does not survive the NCCL path. Say which failures are cleaned up, say plainly that an engine leaving the collective takes the trainer down with torch's watchdog, and print the continue_generation recovery snippet. I5: a refit that fails after any chunk leaves a part-old, part-new model behind a cache only the last chunk would have flushed. Name those engines on stderr and flush on the teardown path, best effort. M6 applies _call_failed to fan-out results, M7 rejects worker-set drift between the pinned rank layout and each fan-out, M8 fixes a backwards message, M11 retries /generate inside the error handling and survives a null meta_info, M12 adds the --timeout help, M13 drops noqa for a rule that is not enabled. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: key4ng <rukeyang@gmail.com>
TokenSpeed spells its tensor-parallel width attn_tp_size, which TOKENSPEED_GRPC_KEYS already lifts into the worker's labels, but discovery only read tp_size -- so a TokenSpeed engine reported tp_size null and a trainer laying out NCCL ranks had nothing to go on. Read attn_tp_size when tp_size is absent or unparseable; an explicit tp_size still wins. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: key4ng <rukeyang@gmail.com>
…h_cache results The /generate retry added for M11 passed the full --timeout to every one of its five attempts, so a slow or black-holed endpoint could hold the script for 5 * --timeout plus the sleeps -- about 50 minutes on the default, and contradicting what --timeout's help promises. Share one deadline across the attempts and the pauses between them, and say so in the help. The teardown's flush_cache fan-out also inspected only res.failed, the gap M6 closed for _await_fanout. Factor that check into _fanout_problems and use it in both places so the two cannot drift apart again. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: key4ng <rukeyang@gmail.com>
Re-ran acceptance item 1 end to end against the fixed engine. Two trainer-driven NCCL refits landed and the next /generate reported weight_version 1 then 2; worst fan-out overhead was 1 ms over twelve fan-outs against a 5 ms budget; a killed engine gave a 207 naming it upstream_unreachable with the breaker closed and no leaked running requests, and returned 200 after restart; --enable-rl off left /v1/rl 404 with an empty body, /workers intact and zero smg_rl_ series beside 338 other smg_ series. Cleanup freed GPUs 0-2 and left the other lane alone. Keeps the root-cause section and records the Task 19 probes verbatim. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: key4ng <rukeyang@gmail.com>
A model-less SGLang-native /generate (what slime and verl send) is read as model "unknown" so the gateway can route it over the whole fleet. The typed HTTP path then re-serialized the request and forwarded `"model":"unknown"` to the engine. SGLang ignores the field, which hid the problem in the M0 run, but a model-aware upstream such as TokenSpeed's `ts serve` sidecar (an HTTP front for its own gateway) resolves the name and answers 404 model_not_found for every slime rollout request. GenerateRequest.model now skips serialization when it holds the wildcard placeholder, so the forwarded body matches what the client sent. Covered by a protocol round-trip test and an HTTP router test that inspects the body the upstream stub receives. Signed-off-by: key4ng <rukeyang@gmail.com>
Task 19's live probes against a TokenSpeed gRPC fleet found two gaps between SMG's gRPC /generate path and what slime (and SGLang itself) expect: - A model-less /generate 404s (model_not_found for "unknown") because slime never sends `model`. When the fleet serves exactly one model there's no ambiguity, so route_generate_impl now defaults the wildcard placeholder to the registry's single served model before canonicalizing; zero or several served models keep the existing 404. route_chat_impl is untouched: ChatCompletionRequest.model has no default-to-"unknown" fallback, so a live client can't reach it with the wildcard placeholder the way /generate can. - A single-prompt /generate answered with a one-element list instead of the bare object SGLang returns, so slime's `output["meta_info"]` indexing broke. generate_response_body() now mirrors SGLang: a single prompt (`text` or flat `input_ids`) with n<=1 and exactly one response collapses to an object; batches and n>1 stay a list. Signed-off-by: key4ng <rukeyang@gmail.com>
…er review Review of 123074b found three Important issues: - generate_response_body built a serde_json::Value with serde_json::json!, which is a hidden .unwrap() (clippy can't see through the macro) and serializes the response twice (once into a Value tree, once out of it). Replaced with an untagged GenerateBody enum (One(Box<GenerateResponse>) | Many(Vec<GenerateResponse>)) serialized directly by axum::Json, same as before this feature existed. - execute_generate cloned the request Arc to inspect its shape after dispatch, which keeps the parsed request (up to 8k+ input_ids) alive for the whole multi-second generation instead of freeing it at the build boundary as the pipeline's request-free-dispatch invariant intends. generate_response_body now takes the two scalars it actually needs (single_prompt, n), read off the request before it moves into the context. - The model-less default matched on worker_registry.get_models() directly, but an untagged worker (no model card, no model_id label) registers under the wildcard placeholder itself. That silently defeated the default in a mixed tagged/untagged fleet, and turned it into a 500 (tokenizer_not_found) in an untagged-only fleet instead of the intended 404. The wildcard is now filtered out of the candidate set before the single-model match. Also: a one-line comment noting the batch-input_ids arm is unreachable over gRPC today (the preparation stage 400s it first), and a documentation correction -- a list-valued `text` fails JSON parsing (422) on every path, not just gRPC, and gRPC's own 400 on batch input_ids is gRPC-specific (HTTP forwards batch input_ids to the real engine). Signed-off-by: key4ng <rukeyang@gmail.com>
…n the H200 acceptance write-up Signed-off-by: key4ng <rukeyang@gmail.com>
…h refused refit routes Signed-off-by: key4ng <rukeyang@gmail.com>
|
Important Draft PR not reviewedDraft PRs are not automatically reviewed by default.
To automatically review draft PRs, update your CodeRabbit configuration: reviews:
auto_review:
drafts: trueComment |
| /// The version to report on a generate response: what the engine stamped on | ||
| /// this very response beats the dispatch-time label (M2's table, then the | ||
| /// registration label), which beats the historical `"default"`. | ||
| pub(crate) fn effective_weight_version(reported: Option<&str>, dispatch: Option<&str>) -> String { |
There was a problem hiding this comment.
🟡 Nit: effective_weight_version returns an owned String, so the four call sites in regular/streaming.rs (lines 1135, 1168, 1363, 1410) now allocate per streamed chunk where the old code embedded &ctx.weight_version with no allocation — this is the per-token hot path. All inputs are already &strs that outlive the json! construction, so the helper can return a borrowed value and let the one non-streaming caller (processor.rs) add .to_string():
pub(crate) fn effective_weight_version<'a>(
reported: Option<&'a str>,
dispatch: Option<&'a str>,
) -> &'a str {
reported.filter(|v| !v.is_empty()).or(dispatch).unwrap_or("default")
}| .weight_version | ||
| .clone() | ||
| .unwrap_or_else(|| "default".to_string()), | ||
| weight_version: response_formatting::effective_weight_version( |
There was a problem hiding this comment.
🟡 Nit: the /generate path now prefers the engine-stamped version, but the OpenAI-shaped responses built from the same ProtoGenerateComplete still report only the dispatch-time label as system_fingerprint (processor.rs:321, processor.rs:945, and the four maybe_system_fingerprint(dispatch.weight_version…) sites in harmony/). After a trainer-driven NCCL refit on TokenSpeed the dispatch label never updates, so /generate will report v2 while /v1/chat/completions keeps reporting the stale registration-time value indefinitely. If that split is intentional scope (slime only reads /generate), fine — but the same effective_weight_version(complete.weight_version(), …) one-liner would keep the two surfaces consistent, and it may be worth a code comment either way so the asymmetry reads as deliberate.
There was a problem hiding this comment.
Reviewed the full diff: RL control-endpoint model (crates/rl), gateway adapter and client-cache wiring, proto weight_version plumbing across the servicer and the four streaming sites, the single-prompt/model-less /generate shape, smg serve label stamping, and the tests/example.
No blocking issues found. The control-URL wildcard resolution, the no_control_endpoint error surface, the stale-stub guard in the servicer, and the untagged-worker filter in the model-less default are all carefully handled and well tested.
Summary of findings:
- 🔴 Important: 0
- 🟡 Nit: 2 — per-chunk
Stringallocation in the streaming hot path viaeffective_weight_version;system_fingerprinton the OpenAI-shaped gRPC responses still reports only the dispatch label, so it diverges from/generateafter a TokenSpeed refit - 🟣 Pre-existing: 0
Description
Problem
SMG's M1 RL control plane could only drive HTTP engines: control calls went to the worker's data URL, so a TokenSpeed engine served over gRPC or ZMQ answered every
/v1/rlcall with 422unsupported_connection_mode, its capabilities showed all-false, and the gRPC generate path reportedweight_versiononly from dispatch-time labels because the proto carried no engine-reported version. Running a real trainer (slime) against TokenSpeed engines behind SMG also surfaced two data-plane gaps: the HTTP passthrough forwarded SMG's internal"model":"unknown"placeholder to the engine, and the gRPC/generaterequired amodelfield and answered a single prompt as a list, neither of which slime's SGLang-native client can use.Solution
A worker's control endpoint is modeled as a well-known label (
rl.control_url), discovered from TokenSpeed's gRPC server info (or supplied on registration), with wildcard hosts resolved to the worker's host. The RL crate proxies control calls to that URL for gRPC/ZMQ workers and keeps HTTP workers controlling themselves; a non-HTTP worker without an endpoint now fails with 422no_control_endpoint. Capabilities come fromrl.*labels the engine advertises, with a static TokenSpeed row (pause_modes: wait,abort,update_from: distributed) as the fallback. The proto gains an optionalweight_versionon generate chunks and completions, the servicer stamps it, and the gateway prefers the engine-reported version over the dispatch label.smg servestampsrl.control_urlon the TokenSpeed workers it launches. For slime's data plane, the wildcard model is never serialized upstream, and the gRPC/generateanswers a single prompt as one object and defaults an absentmodelto the fleet's single served model. A trainer-side NCCL refit example (examples/rl/refit_from_trainer.py) covers TokenSpeed, whose engines refuse disk and tensor refits with 501.Companion TokenSpeed PR: lightseekorg/tokenspeed
rl/control-endpoint(control-app host/auth/route semantics, distributed-only advertisement, NCCL split guard, SGLang-shaped sidecar/get_server_info).Design:
docs/superpowers/specs/2026-09-21-rl-tokenspeed-control-endpoint-design.md(local); acceptance write-up:docs/superpowers/rl/tokenspeed/acceptance-h200.md.Changes
crates/rl:control_urlon the worker view,resolve_control_url,no_control_endpointerror, capability labels with the static TokenSpeed row,tp_sizefalling back toattn_tp_size;model_gateway/src/rl_adapter.rshands out a control client for gRPC/ZMQ workers.crates/grpc_clientproto 0.4.21:GenerateStreamChunk.weight_version = 7,GenerateComplete.weight_version = 12;grpc_servicer0.12.1 stamps them, merges the engine'srl_advertisement()into server info, redacts secrets;TOKENSPEED_GRPC_KEYSliftsrl.*keys into labels.effective_weight_version(engine > dispatch label > "default") in the processor and the four streaming sites; single-prompt object shape and single-model default for model-less requests.GenerateRequest.modelskips serialization when it holds the wildcard placeholder.bindings/python:Worker.control_url,smg servestampsrl.control_urlfor TokenSpeed launches (--rl-control-port).examples/rl/refit_from_trainer.py+ layout test;mock_workergainsserver_args/weight_version; new integration testrl_tokenspeed_control_endpoint_test.rs; e2etest_rl_control_plane.pyparametrized over SGLang-HTTP and TokenSpeed-gRPC.docs/guides/rl-tokenspeed.md,crates/rl/{README,NOTES,COUPLING}.md,examples/rl/README.md.Test Plan
On an 8x H200 node (details and every command in
docs/superpowers/rl/tokenspeed/acceptance-h200.md):cargo +nightly fmt --all -- --check;cargo clippy --all-targets --all-features -- -D warningson Rust 1.98.0 (CI toolchain, OpenCV container);cargo test -p smg-rl -p openai-protocol -p smgincl.rl_control_plane_test,rl_tokenspeed_control_endpoint_test,grpc_pd_fanout_test,tenant_rate_limiting_grpc_test,zmq_backend_test: 2376 passed, 0 failed.grpc_servicertests 19 passed; bindings tests, e2e helper tests, ruff, mypy (one pre-existing error insimple_eval_common.py).pytest e2e_test/router/test_rl_control_plane.py -k TokenSpeed6/6 (TokenSpeed gRPC engines) and-k 'Sglang or Disabled'6/6 (SGLang HTTP,lmsysorg/sglang:v0.5.20).control_urland label-sourced capabilities; two trainer-driven NCCL refits through/v1/rlwith the next/generatereportingweight_version1 then 2; fan-out overhead at most 1 ms over 12 fan-outs; killing an engine mid-fan-out gives 207 naming it (upstream_unreachable), breakers closed, load 0, restart gives 200;--enable-rloff gives 404 on/v1/rl/*with nosmg_rl_series.ts serveengines: slime's own sgl-router vs SMG--enable-rl: reward over the last 50 steps 0.771 vs 0.778 (σ ≈ 0.08), rollout throughput 6949 vs 7741 tok/GPU/s, step time 6.42 vs 6.79 s (+5.7 %, one run per arm). A labeled two-file slime patch (skip registration, abort via engine addresses) ran slime's rollouts over SMG's gRPC data plane for 10 steps.Before: gRPC TokenSpeed workers answered 422 on every
/v1/rlroute, capabilities all false,weight_versionnever left "default"; slime's model-less/generatethrough SMG got 404model_not_foundon both families.Checklist
cargo +nightly fmtpassescargo clippy --all-targets --all-features -- -D warningspasses🤖 Generated with Claude Code