Skip to content

feat(rl): control endpoints for gRPC and ZMQ workers - #2667

Open
key4ng wants to merge 24 commits into
rl/ts-1-weight-version-provenancefrom
rl/ts-2-control-endpoint
Open

key4ng wants to merge 24 commits into
rl/ts-1-weight-version-provenancefrom
rl/ts-2-control-endpoint

Conversation

@key4ng

@key4ng key4ng commented Sep 23, 2026 •

Copy link
Copy Markdown
Member

Description

Part 2 of 3, stacked on the weight_version provenance PR (base branch rl/ts-1-weight-version-provenance). Companion engine PR: lightseekorg/tokenspeed#1736.

Problem

The M1 RL control plane sent control calls to the worker's data URL, so a TokenSpeed engine served over gRPC or ZMQ answered every /v1/rl route with 422 unsupported_connection_mode, and its capabilities showed all-false because nothing read the engine's advertisement.

Solution

A worker's control endpoint is modeled separately from its data transport as a well-known label, rl.control_url, discovered from the engine's server info (PR 1) or supplied on registration, with wildcard hosts (0.0.0.0, ::) resolved to the worker's host. The RL crate proxies control calls to that URL for gRPC/ZMQ workers through a client borrowed from the gateway's worker client cache; HTTP workers keep controlling themselves and ignore the label. A non-HTTP worker without an endpoint fails with 422 no_control_endpoint. Capabilities come from the rl.* labels (source: "label"), with a static TokenSpeed row (pause_modes: wait,abort, update_from: distributed) as the fallback; TokenSpeed refuses disk and tensor refits with 501, so the row and the docs say distributed only. Discovery reports control_url, and the Python client exposes Worker.control_url. smg serve wires no control endpoint for ZMQ TokenSpeed workers: a headless engine starts no control app (see crates/rl/NOTES.md).

Changes

  • crates/rl: control.rs (resolve_control_url), view.rs (control_url, control_client), error.rs (NoControlEndpoint), proxy.rs, fanout.rs, discovery.rs, capability.rs; crates/protocols/src/rl.rs (RlWorkerEntry.control_url).
  • model_gateway/src/rl_adapter.rs (RegistryRlView::new(registry, client_cache) holding one control client per pool config, CONTROL_URL_LABEL, blank label treated as absent), app_context.rs; crates/protocols/src/worker.rs (HttpPoolConfig: Eq + Hash).
  • model_gateway/tests/common/mock_worker.rs (Config.server_args, weight_version), model_gateway/tests/rl_tokenspeed_control_endpoint_test.rs (a TokenSpeed gRPC mock registered through POST /workers, so its server info becomes the labels; a recording control app; mixed-fleet 200 and 207 cases), and the two mock fields set in grpc_context_length_test.rs from main.
  • bindings/python/src/smg/rl.py (Worker.control_url) and its test.
  • Docs: docs/guides/rl-tokenspeed.md, crates/rl/{README,NOTES,COUPLING}.md.

Review follow-ups

  • smg serve --connection-mode zmq no longer passes --rl-control-* to TokenSpeed or stamps rl.control_url: a headless engine builds no AsyncLLM, and only AsyncLLM starts the control app, so the label pointed at a port nothing served (502 instead of the honest 422). RL control on TokenSpeed is gRPC-only for now; NOTES.md has the row.
  • The guide's recipe works against the merged engine: a routable --rl-control-host (the engine advertises no rl.control_url for a wildcard bind; loopback is right only on the gateway's machine) and registration through POST /workers with the key instead of re-registering a --worker-urls worker (a 409).
  • RegistryRlView holds one strong control client per pool config (the gateway cache's entries are weak and gRPC workers hold none, so every control call rebuilt a client); RlWorkerInfo.control_client is a Result so a failed build is logged once and named in the 422.
  • The attn_tp_size fallback in discovery was dead (normalize_grpc_keys already folds it) and is reverted; resolve_control_url leaves an ipc:// worker's label unchanged; control is a private module.
  • The integration test registers through POST /workers and adds the mixed-fleet 200 case (each worker hit once); the mock's unused --server-arg/--weight-version CLI flags are gone.

Test Plan

  • cargo test -p smg-rl -p smg incl. rl_control_plane_test and rl_tokenspeed_control_endpoint_test; cargo clippy --all-targets --all-features -- -D warnings on Rust 1.98.0; cd bindings/python && pytest tests/test_rl_client.py tests/test_serve.py.
  • Live on an H200 (acceptance item 1, full write-up kept outside the repo): two TokenSpeed gRPC engines discovered with control_url and label-sourced capabilities; per-worker init_weights_update_group, update_weights_from_distributed fan-outs and destroy_weights_update_group all proxied to the engines' control apps; fan-out overhead at most 1 ms over 12 fan-outs; killing an engine mid-fan-out gives 207 naming it (upstream_unreachable) with breakers closed and load 0, restart gives 200; --enable-rl off gives 404 on /v1/rl/* and no smg_rl_ series.

Before: POST /v1/rl/engine/flush_cache on a gRPC TokenSpeed worker → 422 unsupported_connection_mode.

Checklist
  • cargo +nightly fmt passes
  • cargo clippy --all-targets --all-features -- -D warnings passes
  • (Optional) Documentation updated
  • (Optional) Please join us on Slack #sig-smg to discuss, review, and merge PRs

🤖 Generated with Claude Code

@coderabbitai

coderabbitai Bot commented Sep 23, 2026 •

Copy link
Copy Markdown

Important

Review skipped

Auto reviews are disabled on base/target branches other than the default branch.

Please check the settings in the CodeRabbit UI or the .coderabbit.yaml file in this repository. To trigger a single review, invoke the @coderabbitai review command.

⚙️ Run configuration
  • Configuration used: Organization UI
  • Review profile: CHILL
  • Plan: Team
  • Run ID: 264ac04c-b6d1-44e6-8ad7-bec348613858

You can disable this status message by setting the reviews.review_status to false in the CodeRabbit configuration file.

Use the checkbox below for a quick retry:

  • 🔍 Trigger review
  • Autopilot · Keep fixing CodeRabbit findings and required CI, and resolving merge conflicts

Comment @coderabbitai help to get the list of available commands.

@github-actions github-actions Bot added documentation Improvements or additions to documentation python-bindings Python bindings changes dependencies Dependency updates grpc gRPC client and router changes tests Test changes protocols Protocols crate changes model-gateway Model gateway crate changes labels Sep 23, 2026
Comment thread docs/guides/rl-tokenspeed.md Outdated
Comment on lines +21 to +23
Register the engines with their control key so the proxy authenticates:

curl -X POST http://smg:30000/workers -H 'content-type: application/json' \

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Nit: This POST /workers will answer 409 WORKER_ALREADY_EXISTS, not attach the key. The launch example above already registers grpc://rollout-1:30000 via --worker-urls, and a duplicate POST for a registered URL is a hard conflict (WorkerServiceError::Conflict in model_gateway/src/worker/service.rs). A reader following the guide verbatim can't set the worker's api_key this way. Either drop the workers from --worker-urls and register them only through POST /workers, or show the worker update route (PATCH /workers/{id}) here instead.

Comment thread bindings/python/src/smg/serve.py Outdated
)
for _, port in self.workers
]
gateway_url = f"http://127.0.0.1:{router_args.port}"

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Nit: The gateway URL is hardcoded to 127.0.0.1, but RouterArgs.host is configurable — a router launched with --host <interface-ip> binds only that interface, so the stamping thread polls loopback for the whole 300 s deadline and then gives up with just a warning, leaving the ZMQ workers without rl.control_url (every /v1/rl control call then 422s with no_control_endpoint). Consider using router_args.host when it isn't a wildcard (0.0.0.0/::), falling back to 127.0.0.1 otherwise.

Comment thread crates/rl/src/control.rs Outdated
Comment on lines +20 to +27
let worker_host = worker_url
.split_once("://")
.map_or(worker_url, |(_, r)| r)
.split('/')
.next()
.map(|authority| split_host_port(authority).0)
.unwrap_or("");
let worker_host = worker_host.split('@').next().unwrap_or(worker_host);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Nit: A worker URL with no authority yields an empty worker_host, producing an unparseable URL like http://:40100. This is exactly the shape of every smg serve ZMQ TokenSpeed worker (ipc:///tmp/engine-30000): split_once("://") leaves /tmp/engine-30000, and split('/').next() returns "". So an operator who supplies a wildcard rl.control_url label on a ZMQ worker (the documented path for older engines) gets http://:30400, which later fails as upstream_unreachable with a confusing URL. Consider returning advertised unchanged when the worker URL yields no host — same contract as the schemeless case — so the error at least names the engine's own advertisement.

Comment on lines +271 to +272
/// Record the `authorization` header of one RL control request. Called
/// alongside [`record_request`] so the two vectors stay index-aligned.

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Nit: This comment contradicts the one on authorizations() (line 193–196): that one says the vector is not aligned with bodies, while this one says the two "stay index-aligned". The authorizations() doc is the correct one — record_request is also called from the generate/chat handlers (lines 745, 830, 1786, 1807), so bodies can interleave non-RL entries. Also note both records are skipped for a bodyless control request (the calls sit inside if let Some(Json(body))), so "the header of each RL control request received" in the authorizations() doc overstates slightly. Worth rewording one side so the next reader doesn't have to work out which comment lies.

@claude claude Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Reviewed the full diff against the stack base (rl/ts-1-weight-version-provenance). The control-endpoint model is clean — control transport decoupled from data transport, wildcard-host resolution unit-tested, auth isolation verified (caller bearer never forwarded, worker key used), and failure modes are loud (no_control_endpoint 422 with an actionable message). Left 4 🟡 nits (409 in the docs registration flow, hardcoded loopback gateway URL in label stamping, empty-host wildcard resolution for ipc:// workers, contradictory test-helper comments); none blocking. 0 🔴, 4 🟡, 0 🟣.

@key4ng
key4ng force-pushed the rl/ts-2-control-endpoint branch 2 times, most recently from b91d84f to ee3315d Compare September 26, 2026 00:15
@key4ng
key4ng force-pushed the rl/ts-1-weight-version-provenance branch from 5f1f476 to 0647569 Compare September 26, 2026 00:15
Comment thread bindings/python/src/smg/serve.py Outdated
Comment on lines +444 to +447
"--rl-control-host",
"127.0.0.1",
"--rl-control-port",
str(_rl_control_port(port)),

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Nit: --rl-control-host/--rl-control-port are passed to tokenspeed.cli serve unconditionally, but these flags only exist on engines that carry the companion PR (lightseekorg/tokenspeed#1736). An installed TokenSpeed that speaks the headless ZMQ wire but predates that PR will exit at argparse with "unrecognized arguments", so smg serve --backend tokenspeed stops launching engines that worked before this change — while the docs' "Older engines" section suggests pre-advertisement builds remain usable. If lockstep versioning with the engine is intended, consider stating the minimum TokenSpeed version in docs/guides/rl-tokenspeed.md; otherwise gate the flags (e.g. probe tokenspeed.cli serve --help once, or an --rl-control-port 0-style opt-out).

@key4ng
key4ng force-pushed the rl/ts-2-control-endpoint branch from ee3315d to 346c01a Compare September 26, 2026 00:34
@key4ng
key4ng force-pushed the rl/ts-1-weight-version-provenance branch from 51ccf27 to 3966e73 Compare September 29, 2026 23:14
@key4ng
key4ng force-pushed the rl/ts-2-control-endpoint branch from 346c01a to b0b6e2e Compare September 29, 2026 23:14
@key4ng
key4ng force-pushed the rl/ts-1-weight-version-provenance branch from 3966e73 to 63aafe8 Compare October 5, 2026 04:55
@key4ng
key4ng force-pushed the rl/ts-2-control-endpoint branch 2 times, most recently from a7371aa to 7da7910 Compare October 5, 2026 05:39
@key4ng
key4ng marked this pull request as ready for review October 6, 2026 03:17
@key4ng
key4ng force-pushed the rl/ts-2-control-endpoint branch from 7da7910 to 6b6f42d Compare October 6, 2026 17:26

The gateway, with no startup workers:

smg launch --policy cache_aware --enable-rl --disable-health-check --disable-circuit-breaker \

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Nit: Dropping --worker-urls fixes the 409, but it also changes which router the gateway builds. The connection mode comes only from the startup URLs: determine_connection_mode (model_gateway/src/main.rs:1401) returns ConnectionMode::Http for an empty list. Without --enable-igw, Gateway then builds one HTTP regular router (gateway.rs:117-126, factory.rs:105). The engines registered afterwards through POST /workers are grpc://, so /v1/rl/* still works (the RL crate reads the registry directly), but the rollout /generate traffic goes to the HTTP router and never reaches them over gRPC.

Two ways to keep both halves working:

  • Keep --worker-urls grpc://rollout-1:30000 grpc://rollout-2:30000 and set the key with PATCH /workers/{id} {"api_key": "..."}. WorkerUpdateRequest.api_key exists for this.
  • Or add --enable-igw, so the gateway builds a gRPC router for workers registered at runtime.

Comment thread crates/rl/src/proxy.rs
let Some(base) = worker.control_url.as_deref() else {
return Err(no_control_endpoint(
worker,
"no `rl.control_url` label: launch the engine with a routable --rl-control-host (a wildcard bind is not advertised), or set the label at registration",

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Nit: This hint is now TokenSpeed-gRPC-specific, but every non-HTTP worker without the label gets it, and for two of those cases it points the operator at a fix that doesn't work:

  • ZMQ TokenSpeed workers. The NOTES.md row added in this push says headless engines build no AsyncLLM, so --rl-control-host/--rl-control-port do nothing there. Relaunching with a routable host can't help.
  • SGLang/vLLM gRPC workers. Those engines have no --rl-control-host flag at all.

Consider phrasing it engine-neutrally and putting the label path first, e.g. "no rl.control_url label: set it at registration or through the worker update route (a TokenSpeed gRPC engine advertises it when launched with a routable --rl-control-host)". Another option is to branch on worker.runtime/connection_mode.

@key4ng
key4ng force-pushed the rl/ts-2-control-endpoint branch from ff05ac2 to 6c999a8 Compare October 6, 2026 21:48
key4ng added 24 commits October 6, 2026 15:20
…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>
…rsion on gRPC generates

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>
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>
_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>
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.

Signed-off-by: key4ng <rukeyang@gmail.com>
…h refused refit routes

Signed-off-by: key4ng <rukeyang@gmail.com>
… that landed on main

Signed-off-by: key4ng <rukeyang@gmail.com>
Signed-off-by: key4ng <rukeyang@gmail.com>
…re into the 422

The gateway's worker client cache keeps weak handles, and a gRPC or ZMQ
worker holds no HTTP client of its own, so every discovery and every
control call on such a fleet rebuilt a reqwest client (TLS roots parsed
again each time). RegistryRlView now keeps one strong handle per pool
config, shared by the workers that use it and pruned when no registered
non-HTTP worker needs it; HttpPoolConfig derives Eq and Hash for the key.

RlWorkerInfo.control_client is a Result: a failed build is logged once
and its reason rides on the worker, so the 422 names it instead of
pointing at the gateway log. The no-label hint now says what the usual
cause is (a wildcard --rl-control-host, which the engine does not
advertise) rather than suggesting an engine upgrade.

resolve_control_url takes the worker's base URL and leaves the label
unchanged when the worker has no host (an ipc:// ZMQ worker), instead of
producing http://:port; the control module is private, its one function
re-exported.

The attn_tp_size fallback in discovery is gone: normalize_grpc_keys
already folds it into tp_size before labels reach the crate, so the
code and its three tests exercised nothing.

Signed-off-by: key4ng <rukeyang@gmail.com>
…ut over endpoints hits each worker once

The integration test registered its workers with BasicWorkerBuilder and
hand-copied labels, so the mock engine's server_args and the label
lifting in discovery fed no assertion. The fleet now registers through
POST /workers the way an operator does: the control URL and the rl.*
capabilities come out of the mock's own server info, the api_key is the
registration's, and the worker whose engine advertises nothing is shown
to fall back to the static TokenSpeed row. The mock advertises what a
current engine does (update_from distributed,mooncake).

Added the spec's mixed-fleet case: a fan-out whose selector covers only
workers with endpoints answers 200 and each control app records exactly
one request. The mock worker's --server-arg and --weight-version CLI
flags are dropped (tests set the fields on Config directly; nothing
passed the flags), and the recorder comment no longer claims the
authorization and body vectors are index-aligned.

Signed-off-by: key4ng <rukeyang@gmail.com>
smg serve passed --rl-control-host/--rl-control-port to the headless
TokenSpeed engine and stamped rl.control_url on the worker. A headless
engine (launch_scheduler_headless) builds no AsyncLLM, and only AsyncLLM
starts the control app, so the label pointed at a port nothing served:
discovery showed a control URL and every control call failed 502
upstream_unreachable, where the worker without the label gets the honest
422. The launcher no longer passes the flags, the stamping thread and its
tests are gone, and NOTES.md records the drift.

Signed-off-by: key4ng <rukeyang@gmail.com>
The guide started each engine with --rl-control-host 0.0.0.0 and promised
a control URL; the merged engine advertises no rl.control_url for a
wildcard bind, so discovery showed null and every call was 422. It also
re-registered a --worker-urls worker through POST /workers, a 409. The
recipe now uses a routable --rl-control-host, says that wildcard hosts
are not advertised and that loopback is right only on the gateway's own
machine, and registers the engines with their key through POST /workers
only. The README says the same about advertisement. The tp_size sentence
describes what discovery does today (attn_tp_size folded into tp_size;
the newest engines nest parallelism under mapping.*, which it does not
read yet). COUPLING.md describes this PR's files only, and NOTES.md no
longer claims the engine advertises a disk refit source.

Signed-off-by: key4ng <rukeyang@gmail.com>
reqwest keeps a CA bundle as bytes and only parses it when a connection
is made, so a bogus bundle still produced a client and the test's
expect_err panicked. An unparsable client identity is rejected when the
client is built, which is the failure the test means to observe.

Signed-off-by: key4ng <rukeyang@gmail.com>
clippy (field_reassign_with_default) rejects assigning a field on a value
just created with Default::default(); CI runs with -D warnings.

Signed-off-by: key4ng <rukeyang@gmail.com>
@key4ng
key4ng force-pushed the rl/ts-2-control-endpoint branch from 6c999a8 to 5f12276 Compare October 6, 2026 22:20

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

dependencies Dependency updates documentation Improvements or additions to documentation grpc gRPC client and router changes model-gateway Model gateway crate changes protocols Protocols crate changes python-bindings Python bindings changes tests Test changes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant