Repository navigation
feat(discovery): level-triggered informer reconcile with multi-port pods - #2061
Conversation
📝 WalkthroughSummary by CodeRabbit
WalkthroughThe PR adds configurable worker-port annotations to Python and model-gateway discovery configuration. Kubernetes discovery resolves multiple worker ports and reconciles workers through reflector state, with UID ownership, fallback handling, router event updates, and expanded tests. ChangesWorker Port Discovery
Estimated code review effort: 4 (Complex) | ~60 minutes Possibly related PRs
Suggested labels: Suggested reviewers: Sequence Diagram(s)sequenceDiagram
participant KubernetesReflector
participant ServiceDiscovery
participant WorkerRegistry
participant RouterWatcher
participant MeshNodeStatus
KubernetesReflector->>ServiceDiscovery: pod events and synchronized store
ServiceDiscovery->>WorkerRegistry: reconcile UID-owned workers
RouterWatcher->>ServiceDiscovery: router pod event
ServiceDiscovery->>MeshNodeStatus: apply Alive, Suspected, or Down state
🚥 Pre-merge checks | ✅ 5✅ Passed checks (5 passed)
✨ Finishing Touches📝 Generate docstrings
🧪 Generate unit tests (beta)
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
|
👋 The PR description doesn't fully follow
Please update the PR description so reviewers have the context they need. |
There was a problem hiding this comment.
Actionable comments posted: 4
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (1)
bindings/python/src/lib.rs (1)
904-1034: 🎯 Functional Correctness | 🟠 Major | ⚡ Quick win🔴 Important: Append the new configuration field instead of inserting it mid-list in both Python-facing argument lists. Both files place
worker_ports_annotationdirectly afterbootstrap_port_annotation, which shifts the positional index of every later argument. Both files carry an explicit append-only rule for exactly this reason (bindings/python/src/lib.rslines 486-488 and 979-981;bindings/python/src/smg/router_args.pyline 204). In-repo callers use keywords, so only external positional callers break, and they break silently by binding values to the wrong parameters.
bindings/python/src/lib.rs#L904-L1034: moveworker_ports_annotationafterworker_startup_delayin the#[pyo3(signature = ...)]block, in thefn newparameter list, and in thestruct Routerfield list.bindings/python/src/smg/router_args.py#L95-L95: move theworker_ports_annotationdataclass field belowworker_startup_delay, under the append-only comment at line 204.🤖 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 `@bindings/python/src/lib.rs` around lines 904 - 1034, Append worker_ports_annotation after worker_startup_delay in all three bindings/python/src/lib.rs locations: the PyO3 signature, fn new parameters, and Router field list; preserve the existing append-only ordering. Also move the dataclass field in bindings/python/src/smg/router_args.py:95 below worker_startup_delay under the append-only comment, ensuring existing positional argument indices remain unchanged.
🧹 Nitpick comments (5)
model_gateway/src/service_discovery.rs (3)
368-436: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low value🟡 Nit: Port parsing and alignment rules look correct; one ambiguity is worth a comment.
parse_port_listrejects0, empty parts, and values above 65535, so a trailing comma or junk entry rejects the whole annotation and logs a warning.resolve_bootstrap_portskeepsbootstrap_ports.len() == ports.len()on every path, whichcompute_desired_staterelies on when it indexes by position.When
num_ports == 1, the broadcast arm and the exact-match arm both apply. The result is identical, so this is not a defect. Add a short comment stating that the single-value arm intentionally also covers the one-port case, so a later edit does not reorder the arms and change meaning.Note that a typo in the worker-ports annotation falls back to
config.port. A multi-server pod then registers one worker on a port that may not listen. The warning log makes this diagnosable, and the health checker removes the dead worker, so the behavior is acceptable.🤖 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 `@model_gateway/src/service_discovery.rs` around lines 368 - 436, Add a brief comment in resolve_bootstrap_ports above the single-value match arm explaining that it intentionally also handles the one-port case, preserving the current broadcast behavior and documenting why its order relative to the exact-length arm matters.
826-839: 📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick win🟡 Nit: Record a metric when a removal submission fails.
The addition loop records
REGISTRATION_FAILEDon submit errors (lines 857-860). The removal loop only logs. A failed removal keeps traffic on a worker whose pod is gone, so it deserves the same signal as a failed addition. Add a deregistration failure label, or reuse an existing one.The next reconcile pass retries the removal, so this is an observability gap rather than a correctness defect.
As per coding guidelines: "Run the silent-failure-hunter agent on changed files to detect swallowed errors, inappropriate fallbacks, and missing error propagation."
🤖 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 `@model_gateway/src/service_discovery.rs` around lines 826 - 839, The worker removal loop in the job submission flow records only successful deregistration metrics. Update the Err branch of the match around job_queue.submit in the removal handling to record the appropriate deregistration failure metric, reusing an existing failure label if available, while preserving the current error log and retry behavior.Source: Coding guidelines
1263-1275: 📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick win🟡 Nit: Add a multi-port Decode case to cover the non-Encode/Prefill bootstrap branch.
compute_desired_stateindexesbootstrap_portsby the position of each entry inports, so the two vectors must stay the same length. Thevec![None; ports.len()]branch at line 319 is only exercised with a single port today. A Decode pod that carriessmg.ai/worker-portswith two values would pin that invariant.#[test] fn test_from_pod_decode_multi_port_bootstrap_all_none() { let config = create_pd_config(); let mut pod = create_pd_k8s_pod("decode-0", "10.0.0.2", "decode", Some(9080)); pod.metadata .annotations .as_mut() .unwrap() .insert("smg.ai/worker-ports".to_string(), "8080,8081".to_string()); let info = PodInfo::from_pod(&pod, Some(&config)).unwrap(); assert_eq!(info.ports, vec![8080, 8081]); assert_eq!(info.bootstrap_ports, vec![None, None]); }As per coding guidelines: "Run the pr-test-analyzer agent to verify that tests adequately cover new or changed functionality."
🤖 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 `@model_gateway/src/service_discovery.rs` around lines 1263 - 1275, Add a test near test_pod_info_from_pod_with_pd_config_decode that creates a Decode pod with the smg.ai/worker-ports annotation containing two ports, calls PodInfo::from_pod with the PD config, and asserts both ports are parsed and bootstrap_ports contains a matching [None, None] entry for each port. Run the pr-test-analyzer agent to verify coverage.Source: Coding guidelines
model_gateway/src/main.rs (1)
1350-1350: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low value🟡 Nit: Share one constant for the default annotation key.
The literal
"smg.ai/worker-ports"now appears inmodel_gateway/src/config/types.rs(line 648),model_gateway/src/service_discovery.rs(line 197),model_gateway/src/main.rs(lines 1350 and 1614),bindings/python/src/lib.rs(line 904), andbindings/python/src/smg/router_args.py(lines 95 and 1356). The Rust copies can reference one exported constant so a future rename cannot drift between the paths.Export the constant from
service_discovery.rsnext toPOD_UID_LABELand reference it fromtypes.rsand bothmain.rspaths.🤖 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 `@model_gateway/src/main.rs` at line 1350, Define and export a shared default worker-ports annotation constant in service_discovery.rs alongside POD_UID_LABEL, then replace the duplicated "smg.ai/worker-ports" literals in config types.rs and both relevant main.rs paths with that constant. Keep the existing annotation key value and update Rust references to use the exported symbol.bindings/python/src/smg/router_args.py (1)
1354-1356: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low value🟡 Nit: Preserve annotation values supplied through
from_cli_args. The field loop accepts these values, but lines 1355–1356 overwrite them unconditionally. If the annotation keys are fixed for CLI construction, document that these fields are ignored. Otherwise, add CLI flags and remove the overwrite. Keep theDiscoveryConfigfields because Python bindings and serialized configuration use them.🤖 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 `@bindings/python/src/smg/router_args.py` around lines 1354 - 1356, Update the annotation assignments in the from_cli_args construction flow so values already provided through the field loop are not unconditionally overwritten; add corresponding CLI flags and preserve those parsed values, while retaining the DiscoveryConfig fields for Python bindings and serialization.Source: Coding guidelines
🤖 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 `@bindings/python/src/lib.rs`:
- Line 904: Move worker_ports_annotation from its current mid-list position to
the end of the _Router signature, after worker_startup_delay, in the
#[pyo3(signature = ...)] block, fn new parameter list, and struct field list.
Keep its existing type and default unchanged so all pre-existing positional
arguments retain their indexes.
In `@model_gateway/src/main.rs`:
- Around line 1613-1614: Update to_server_config so both
bootstrap_port_annotation and worker_ports_annotation use the corresponding
configured values from router_config.discovery, retaining their current literals
only as fallbacks. Preserve the existing handling of router_selector and
model_id_source while ensuring custom annotation keys reach
ServiceDiscoveryConfig.
In `@model_gateway/src/service_discovery.rs`:
- Around line 577-615: Replace the final tokio::join! coordinating driver and
reconciler with tokio::select! so the discovery task terminates as soon as
either future completes. Preserve the existing driver stream-end and reconciler
wait_until_ready failure handling, allowing the JoinHandle from
start_service_discovery to observe termination for supervision or restart.
- Around line 742-763: Update compute_actions so removal deduplication preserves
every worker revision for each URL instead of collapsing all ranks to one URL.
Group removals by the URL and revision used by FindWorkersToRemoveStep, or
otherwise ensure all DP-rank revisions are removed atomically while retaining
the existing desired-state filtering.
---
Outside diff comments:
In `@bindings/python/src/lib.rs`:
- Around line 904-1034: Append worker_ports_annotation after
worker_startup_delay in all three bindings/python/src/lib.rs locations: the PyO3
signature, fn new parameters, and Router field list; preserve the existing
append-only ordering. Also move the dataclass field in
bindings/python/src/smg/router_args.py:95 below worker_startup_delay under the
append-only comment, ensuring existing positional argument indices remain
unchanged.
---
Nitpick comments:
In `@bindings/python/src/smg/router_args.py`:
- Around line 1354-1356: Update the annotation assignments in the from_cli_args
construction flow so values already provided through the field loop are not
unconditionally overwritten; add corresponding CLI flags and preserve those
parsed values, while retaining the DiscoveryConfig fields for Python bindings
and serialization.
In `@model_gateway/src/main.rs`:
- Line 1350: Define and export a shared default worker-ports annotation constant
in service_discovery.rs alongside POD_UID_LABEL, then replace the duplicated
"smg.ai/worker-ports" literals in config types.rs and both relevant main.rs
paths with that constant. Keep the existing annotation key value and update Rust
references to use the exported symbol.
In `@model_gateway/src/service_discovery.rs`:
- Around line 368-436: Add a brief comment in resolve_bootstrap_ports above the
single-value match arm explaining that it intentionally also handles the
one-port case, preserving the current broadcast behavior and documenting why its
order relative to the exact-length arm matters.
- Around line 826-839: The worker removal loop in the job submission flow
records only successful deregistration metrics. Update the Err branch of the
match around job_queue.submit in the removal handling to record the appropriate
deregistration failure metric, reusing an existing failure label if available,
while preserving the current error log and retry behavior.
- Around line 1263-1275: Add a test near
test_pod_info_from_pod_with_pd_config_decode that creates a Decode pod with the
smg.ai/worker-ports annotation containing two ports, calls PodInfo::from_pod
with the PD config, and asserts both ports are parsed and bootstrap_ports
contains a matching [None, None] entry for each port. Run the pr-test-analyzer
agent to verify coverage.
🪄 Autofix
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: CHILL
Plan: Pro Plus
Run ID: 140597dc-7138-4a01-b294-dace75765775
📒 Files selected for processing (8)
bindings/python/src/lib.rsbindings/python/src/smg/router.pybindings/python/src/smg/router_args.pybindings/python/tests/test_router_config.pymodel_gateway/src/config/types.rsmodel_gateway/src/main.rsmodel_gateway/src/observability/metrics.rsmodel_gateway/src/service_discovery.rs
💤 Files with no reviewable changes (1)
- model_gateway/src/observability/metrics.rs
| decode_selector = HashMap::new(), | ||
| router_selector = HashMap::new(), | ||
| bootstrap_port_annotation = String::from("sglang.ai/bootstrap-port"), | ||
| worker_ports_annotation = String::from("smg.ai/worker-ports"), |
There was a problem hiding this comment.
🎯 Functional Correctness | 🟠 Major | ⚡ Quick win
🔴 Important: Move worker_ports_annotation to the end of the signature; the mid-list insertion breaks positional callers.
The parameter is inserted after bootstrap_port_annotation, which shifts the index of every later positional argument by one. The file states this rule twice:
- Lines 486-488: "New parameters MUST be appended here (not inserted mid-list) to avoid breaking external Python callers that pass
_Router(...)positionally." - Lines 979-981: "Appended last (not inserted mid-list) so every pre-existing positional argument keeps its index."
Router.from_args calls _Router(**args_dict) with keywords, so in-repo use is unaffected. Any external caller that passes arguments positionally silently binds worker_ports_annotation to what was previously model_id_from, and so on down the list. Append the parameter after worker_startup_delay in the #[pyo3(signature = ...)] block, in the fn new parameter list, and in the struct field list.
🐛 Proposed fix to preserve positional-argument compatibility
bootstrap_port_annotation = String::from("sglang.ai/bootstrap-port"),
- worker_ports_annotation = String::from("smg.ai/worker-ports"),
model_id_from = None, model_aliases = HashMap::new(),
worker_startup_delay = 0,
+ worker_ports_annotation = String::from("smg.ai/worker-ports"),
))]Apply the same move in fn new (remove line 1034, append after worker_startup_delay: u64, at line 1120) and in the struct field list (remove line 411, append after worker_startup_delay: u64, at line 497).
Also applies to: 1034-1034
🤖 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 `@bindings/python/src/lib.rs` at line 904, Move worker_ports_annotation from
its current mid-list position to the end of the _Router signature, after
worker_startup_delay, in the #[pyo3(signature = ...)] block, fn new parameter
list, and struct field list. Keep its existing type and default unchanged so all
pre-existing positional arguments retain their indexes.
| bootstrap_port_annotation: "sglang.ai/bootstrap-port".to_string(), | ||
| worker_ports_annotation: "smg.ai/worker-ports".to_string(), |
There was a problem hiding this comment.
🎯 Functional Correctness | 🟠 Major | ⚡ Quick win
🔴 Important: to_server_config ignores the configured annotation key and hardcodes the default.
DiscoveryConfig::worker_ports_annotation is deserializable from a config file, so an operator can set a custom key there. to_server_config builds the ServiceDiscoveryConfig that the discovery runtime actually consumes, and it writes the literal "smg.ai/worker-ports" instead of reading router_config.discovery. The same function already reads router_selector and model_id_source from router_config.discovery (lines 1574-1601), so the inconsistency is limited to the annotation keys.
Result: a custom worker_ports_annotation in a config file is accepted, validated, and then silently dropped. Pods annotated with the custom key are treated as single-port pods. Read the value from router_config.discovery with the literal only as the fallback.
🟣 Pre-existing: bootstrap_port_annotation on line 1613 has the same defect. Fix both together.
🐛 Proposed fix to honor the configured annotation keys
- let (router_selector, router_mesh_port_annotation) = router_config
+ let (
+ router_selector,
+ router_mesh_port_annotation,
+ bootstrap_port_annotation,
+ worker_ports_annotation,
+ ) = router_config
.discovery
.as_ref()
.map(|d| {
(
d.router_selector.clone(),
d.router_mesh_port_annotation.clone(),
+ d.bootstrap_port_annotation.clone(),
+ d.worker_ports_annotation.clone(),
)
})
- .unwrap_or_else(|| (HashMap::new(), "sglang.ai/mesh-port".to_string()));
+ .unwrap_or_else(|| {
+ (
+ HashMap::new(),
+ "sglang.ai/mesh-port".to_string(),
+ "sglang.ai/bootstrap-port".to_string(),
+ "smg.ai/worker-ports".to_string(),
+ )
+ });- bootstrap_port_annotation: "sglang.ai/bootstrap-port".to_string(),
- worker_ports_annotation: "smg.ai/worker-ports".to_string(),
+ bootstrap_port_annotation,
+ worker_ports_annotation,📝 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.
| bootstrap_port_annotation: "sglang.ai/bootstrap-port".to_string(), | |
| worker_ports_annotation: "smg.ai/worker-ports".to_string(), | |
| let ( | |
| router_selector, | |
| router_mesh_port_annotation, | |
| bootstrap_port_annotation, | |
| worker_ports_annotation, | |
| ) = router_config | |
| .discovery | |
| .as_ref() | |
| .map(|d| { | |
| ( | |
| d.router_selector.clone(), | |
| d.router_mesh_port_annotation.clone(), | |
| d.bootstrap_port_annotation.clone(), | |
| d.worker_ports_annotation.clone(), | |
| ) | |
| }) | |
| .unwrap_or_else(|| { | |
| ( | |
| HashMap::new(), | |
| "sglang.ai/mesh-port".to_string(), | |
| "sglang.ai/bootstrap-port".to_string(), | |
| "smg.ai/worker-ports".to_string(), | |
| ) | |
| }); | |
| bootstrap_port_annotation, | |
| worker_ports_annotation, |
🤖 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 `@model_gateway/src/main.rs` around lines 1613 - 1614, Update to_server_config
so both bootstrap_port_annotation and worker_ports_annotation use the
corresponding configured values from router_config.discovery, retaining their
current literals only as fallbacks. Preserve the existing handling of
router_selector and model_id_source while ensuring custom annotation keys reach
ServiceDiscoveryConfig.
| let driver = { | ||
| let notify = Arc::clone(¬ify); | ||
| async move { | ||
| while let Some(item) = stream.next().await { | ||
| match item { | ||
| Ok(_) => notify.notify_one(), | ||
| Err(e) => { | ||
| error!("K8s worker watcher error (auto-retrying with backoff): {e}"); | ||
| } | ||
| }); | ||
| if let Err(e) = handle.await { | ||
| error!( | ||
| "Periodic reconciliation task panicked: {} -- restarting after {}s", | ||
| e, | ||
| reconcile_interval.as_secs() | ||
| ); | ||
| time::sleep(reconcile_interval).await; | ||
| } else { | ||
| break; | ||
| } | ||
| } | ||
| }); | ||
| error!("K8s worker watcher stream ended; discovery no longer receives updates"); | ||
| } | ||
| }; | ||
|
|
||
| let reconciler = async { | ||
| if store.wait_until_ready().await.is_err() { | ||
| error!("K8s worker watcher dropped before initial sync; reconciliation disabled"); | ||
| return; | ||
| } | ||
| info!( | ||
| "Periodic reconciliation enabled | interval: {}s", | ||
| config.check_interval.as_secs() | ||
| "K8s worker store synced, reconciling on change and every {}s", | ||
| config_arc.check_interval.as_secs() | ||
| ); | ||
| } | ||
|
|
||
| let mut retry_delay = Duration::from_secs(1); | ||
| const MAX_RETRY_DELAY: Duration = Duration::from_secs(300); | ||
|
|
||
| loop { | ||
| let watcher_config = build_watcher_config("worker", &config_arc.list_label_selector()); | ||
| let watcher_stream = watcher(pods.clone(), watcher_config).applied_objects(); | ||
|
|
||
| let config_clone = Arc::clone(&config_arc); | ||
| let tracked_pods_clone = Arc::clone(&tracked_pods); | ||
|
|
||
| let filtered_stream = watcher_stream.filter_map(move |obj_res| { | ||
| let config_inner = Arc::clone(&config_clone); | ||
|
|
||
| async move { | ||
| match obj_res { | ||
| Ok(pod) => { | ||
| if PodInfo::should_include(&pod, &config_inner) { | ||
| Some(Ok(pod)) | ||
| } else { | ||
| None | ||
| } | ||
| } | ||
| Err(e) => Some(Err(e)), | ||
| } | ||
| } | ||
| }); | ||
|
|
||
| let tracked_pods_clone2 = Arc::clone(&tracked_pods_clone); | ||
| let app_context_clone = Arc::clone(&app_context); | ||
| let config_clone2 = Arc::clone(&config_arc); | ||
|
|
||
| let watcher_ok = filtered_stream | ||
| .try_for_each(move |pod| { | ||
| let tracked_pods_inner = Arc::clone(&tracked_pods_clone2); | ||
| let app_context_inner = Arc::clone(&app_context_clone); | ||
| let config_inner = Arc::clone(&config_clone2); | ||
|
|
||
| async move { | ||
| let pod_info = PodInfo::from_pod(&pod, Some(&config_inner)); | ||
|
|
||
| if let Some(pod_info) = pod_info { | ||
| if pod.metadata.deletion_timestamp.is_some() { | ||
| handle_pod_deletion( | ||
| &pod_info, | ||
| tracked_pods_inner, | ||
| app_context_inner, | ||
| port, | ||
| ) | ||
| .await; | ||
| } else { | ||
| handle_pod_event( | ||
| &pod_info, | ||
| tracked_pods_inner, | ||
| app_context_inner, | ||
| port, | ||
| config_inner.disaggregated_mode, | ||
| ) | ||
| .await; | ||
| } | ||
| } | ||
| Ok(()) | ||
| let mut interval = time::interval(config_arc.check_interval); | ||
| interval.set_missed_tick_behavior(time::MissedTickBehavior::Delay); | ||
| loop { | ||
| tokio::select! { | ||
| () = notify.notified() => { | ||
| // Coalesce event bursts (relists, rollouts) into one pass. | ||
| time::sleep(Duration::from_secs(1)).await; | ||
| } | ||
| }) | ||
| .await; | ||
|
|
||
| match watcher_ok { | ||
| Ok(()) => { | ||
| retry_delay = Duration::from_secs(1); | ||
| } | ||
| Err(err) => { | ||
| error!("Error in Kubernetes watcher: {}", err); | ||
| warn!( | ||
| "Retrying in {} seconds with exponential backoff", | ||
| retry_delay.as_secs() | ||
| ); | ||
| time::sleep(retry_delay).await; | ||
|
|
||
| retry_delay = std::cmp::min(retry_delay * 2, MAX_RETRY_DELAY); | ||
| _ = interval.tick() => {} | ||
| } | ||
| reconcile_workers(&store, &config_arc, &app_context).await; | ||
| } | ||
| }; | ||
|
|
||
| warn!( | ||
| "Kubernetes watcher exited, restarting in {} seconds", | ||
| config_arc.check_interval.as_secs() | ||
| ); | ||
| time::sleep(config_arc.check_interval).await; | ||
| } | ||
| tokio::join!(driver, reconciler); |
There was a problem hiding this comment.
🩺 Stability & Availability | 🟠 Major | ⚡ Quick win
🔴 Important: Couple the driver and reconciler lifetimes so a dead watcher does not leave discovery silently frozen.
tokio::join! waits for both futures, and the reconciler loop never returns. Two failure paths leave the task alive but non-functional:
- The watcher stream ends (line 588). The driver returns, the reconciler keeps ticking every
check_interval, and every pass reconciles against a store that no longer receives updates. Deleted pods are never removed and new pods are never registered. The gauge keeps reporting the frozen desired count, so the failure looks healthy on dashboards. wait_until_readyfails (line 593). The reconciler returns immediately, but the driver keeps the task alive forever with no reconciliation.
In both cases the task never completes, so the JoinHandle returned by start_service_discovery cannot signal failure and no supervisor can restart discovery. Replace join! with select! so the task terminates when either half stops. The caller can then observe termination and restart or fail the process.
🛡️ Proposed fix to terminate the task when either half stops
- tokio::join!(driver, reconciler);
+ // Either half stopping means discovery can no longer converge, so end
+ // the task and let the caller observe it through the JoinHandle.
+ tokio::select! {
+ () = driver => {}
+ () = reconciler => {}
+ }
+ error!("K8s worker discovery stopped; workers are no longer reconciled");📝 Committable suggestion
‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.
| let driver = { | |
| let notify = Arc::clone(¬ify); | |
| async move { | |
| while let Some(item) = stream.next().await { | |
| match item { | |
| Ok(_) => notify.notify_one(), | |
| Err(e) => { | |
| error!("K8s worker watcher error (auto-retrying with backoff): {e}"); | |
| } | |
| }); | |
| if let Err(e) = handle.await { | |
| error!( | |
| "Periodic reconciliation task panicked: {} -- restarting after {}s", | |
| e, | |
| reconcile_interval.as_secs() | |
| ); | |
| time::sleep(reconcile_interval).await; | |
| } else { | |
| break; | |
| } | |
| } | |
| }); | |
| error!("K8s worker watcher stream ended; discovery no longer receives updates"); | |
| } | |
| }; | |
| let reconciler = async { | |
| if store.wait_until_ready().await.is_err() { | |
| error!("K8s worker watcher dropped before initial sync; reconciliation disabled"); | |
| return; | |
| } | |
| info!( | |
| "Periodic reconciliation enabled | interval: {}s", | |
| config.check_interval.as_secs() | |
| "K8s worker store synced, reconciling on change and every {}s", | |
| config_arc.check_interval.as_secs() | |
| ); | |
| } | |
| let mut retry_delay = Duration::from_secs(1); | |
| const MAX_RETRY_DELAY: Duration = Duration::from_secs(300); | |
| loop { | |
| let watcher_config = build_watcher_config("worker", &config_arc.list_label_selector()); | |
| let watcher_stream = watcher(pods.clone(), watcher_config).applied_objects(); | |
| let config_clone = Arc::clone(&config_arc); | |
| let tracked_pods_clone = Arc::clone(&tracked_pods); | |
| let filtered_stream = watcher_stream.filter_map(move |obj_res| { | |
| let config_inner = Arc::clone(&config_clone); | |
| async move { | |
| match obj_res { | |
| Ok(pod) => { | |
| if PodInfo::should_include(&pod, &config_inner) { | |
| Some(Ok(pod)) | |
| } else { | |
| None | |
| } | |
| } | |
| Err(e) => Some(Err(e)), | |
| } | |
| } | |
| }); | |
| let tracked_pods_clone2 = Arc::clone(&tracked_pods_clone); | |
| let app_context_clone = Arc::clone(&app_context); | |
| let config_clone2 = Arc::clone(&config_arc); | |
| let watcher_ok = filtered_stream | |
| .try_for_each(move |pod| { | |
| let tracked_pods_inner = Arc::clone(&tracked_pods_clone2); | |
| let app_context_inner = Arc::clone(&app_context_clone); | |
| let config_inner = Arc::clone(&config_clone2); | |
| async move { | |
| let pod_info = PodInfo::from_pod(&pod, Some(&config_inner)); | |
| if let Some(pod_info) = pod_info { | |
| if pod.metadata.deletion_timestamp.is_some() { | |
| handle_pod_deletion( | |
| &pod_info, | |
| tracked_pods_inner, | |
| app_context_inner, | |
| port, | |
| ) | |
| .await; | |
| } else { | |
| handle_pod_event( | |
| &pod_info, | |
| tracked_pods_inner, | |
| app_context_inner, | |
| port, | |
| config_inner.disaggregated_mode, | |
| ) | |
| .await; | |
| } | |
| } | |
| Ok(()) | |
| let mut interval = time::interval(config_arc.check_interval); | |
| interval.set_missed_tick_behavior(time::MissedTickBehavior::Delay); | |
| loop { | |
| tokio::select! { | |
| () = notify.notified() => { | |
| // Coalesce event bursts (relists, rollouts) into one pass. | |
| time::sleep(Duration::from_secs(1)).await; | |
| } | |
| }) | |
| .await; | |
| match watcher_ok { | |
| Ok(()) => { | |
| retry_delay = Duration::from_secs(1); | |
| } | |
| Err(err) => { | |
| error!("Error in Kubernetes watcher: {}", err); | |
| warn!( | |
| "Retrying in {} seconds with exponential backoff", | |
| retry_delay.as_secs() | |
| ); | |
| time::sleep(retry_delay).await; | |
| retry_delay = std::cmp::min(retry_delay * 2, MAX_RETRY_DELAY); | |
| _ = interval.tick() => {} | |
| } | |
| reconcile_workers(&store, &config_arc, &app_context).await; | |
| } | |
| }; | |
| warn!( | |
| "Kubernetes watcher exited, restarting in {} seconds", | |
| config_arc.check_interval.as_secs() | |
| ); | |
| time::sleep(config_arc.check_interval).await; | |
| } | |
| tokio::join!(driver, reconciler); | |
| let driver = { | |
| let notify = Arc::clone(¬ify); | |
| async move { | |
| while let Some(item) = stream.next().await { | |
| match item { | |
| Ok(_) => notify.notify_one(), | |
| Err(e) => { | |
| error!("K8s worker watcher error (auto-retrying with backoff): {e}"); | |
| } | |
| } | |
| } | |
| error!("K8s worker watcher stream ended; discovery no longer receives updates"); | |
| } | |
| }; | |
| let reconciler = async { | |
| if store.wait_until_ready().await.is_err() { | |
| error!("K8s worker watcher dropped before initial sync; reconciliation disabled"); | |
| return; | |
| } | |
| info!( | |
| "K8s worker store synced, reconciling on change and every {}s", | |
| config_arc.check_interval.as_secs() | |
| ); | |
| let mut interval = time::interval(config_arc.check_interval); | |
| interval.set_missed_tick_behavior(time::MissedTickBehavior::Delay); | |
| loop { | |
| tokio::select! { | |
| () = notify.notified() => { | |
| // Coalesce event bursts (relists, rollouts) into one pass. | |
| time::sleep(Duration::from_secs(1)).await; | |
| } | |
| _ = interval.tick() => {} | |
| } | |
| reconcile_workers(&store, &config_arc, &app_context).await; | |
| } | |
| }; | |
| // Either half stopping means discovery can no longer converge, so end | |
| // the task and let the caller observe it through the JoinHandle. | |
| tokio::select! { | |
| () = driver => {} | |
| () = reconciler => {} | |
| } | |
| error!("K8s worker discovery stopped; workers are no longer reconciled"); |
🤖 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 `@model_gateway/src/service_discovery.rs` around lines 577 - 615, Replace the
final tokio::join! coordinating driver and reconciler with tokio::select! so the
discovery task terminates as soon as either future completes. Preserve the
existing driver stream-end and reconciler wait_until_ready failure handling,
allowing the JoinHandle from start_service_discovery to observe termination for
supervision or restart.
| fn compute_actions(desired: &DesiredState, registered: &[OwnedWorker]) -> ReconcileActions { | ||
| let mut actions = ReconcileActions::default(); | ||
|
|
||
| // Layer 4: Record failed registration | ||
| Metrics::record_discovery_registration( | ||
| metrics_labels::DISCOVERY_KUBERNETES, | ||
| metrics_labels::REGISTRATION_FAILED, | ||
| ); | ||
| let mut registered_uid: HashMap<&str, &str> = HashMap::new(); | ||
| for worker in registered { | ||
| registered_uid.insert(worker.url.as_str(), worker.pod_uid.as_str()); | ||
| if !desired.present.contains(&worker.url) { | ||
| actions.remove.push(worker.clone()); | ||
| } | ||
| } | ||
| // DP-rank expansions share one canonical URL; remove it once. | ||
| actions.remove.sort_unstable_by(|a, b| a.url.cmp(&b.url)); | ||
| actions.remove.dedup_by(|a, b| a.url == b.url); | ||
|
|
||
| match tracked_pods.lock() { | ||
| Ok(mut tracker) => { | ||
| tracker.remove(pod_info); | ||
| } | ||
| Err(e) => { | ||
| error!( | ||
| "Lock poisoned during rollback for {}: {} -- tracked state is now inconsistent", | ||
| worker_url, e | ||
| ); | ||
| } | ||
| } | ||
| } | ||
| } | ||
| } else { | ||
| debug!( | ||
| "JobQueue not initialized, skipping async worker addition for: {}", | ||
| worker_url | ||
| ); | ||
| } | ||
| } else { | ||
| // Pod already tracked - this is a duplicate event | ||
| Metrics::record_discovery_registration( | ||
| metrics_labels::DISCOVERY_KUBERNETES, | ||
| metrics_labels::REGISTRATION_DUPLICATE, | ||
| ); | ||
| for worker in &desired.addable { | ||
| match registered_uid.get(worker.url.as_str()) { | ||
| Some(uid) if *uid == worker.pod_uid => {} | ||
| _ => actions.add.push(worker.clone()), | ||
| } | ||
| } | ||
| actions | ||
| } |
There was a problem hiding this comment.
🗄️ Data Integrity & Integration | 🟡 Minor | ⚡ Quick win
🧩 Analysis chain
🏁 Script executed:
#!/bin/bash
# Inspect RemoveWorker handling and the expected_revision guard.
rg -n -C10 'RemoveWorker' --type=rust -g '!**/tests/**'
# Inspect Upsert registration handling.
rg -n -C10 'WorkerRegistrationMode' --type=rust -g '!**/tests/**'
# Inspect registry removal by URL and DP-rank expansion.
rg -nP -C8 'fn (remove_by_url|remove_inner|get_by_url)\s*\(' --type=rustRepository: smg-project/smg
Length of output: 153
🏁 Script executed:
#!/bin/bash
set -u
printf '%s\n' '--- candidate files ---'
git ls-files | rg '(^|/)(service_discovery|registry|job|worker).*\.(rs|py|go)$|Cargo\.toml$'
printf '%s\n' '--- symbols and terms ---'
rg -n -C6 'RemoveWorker|Remove.*Worker|expected_revision|WorkerRegistrationMode|Upsert|upsert|DP.?rank|canonical.*url|register' --glob '*.rs' --glob '*.py' --glob '*.go' . | head -n 1200Repository: smg-project/smg
Length of output: 50372
🏁 Script executed:
#!/bin/bash
set -u
printf '%s\n' '--- service discovery structure ---'
ast-grep outline model_gateway/src/service_discovery.rs
printf '%s\n' '--- service discovery relevant terms ---'
rg -n -C12 'compute_actions|ReconcileActions|OwnedWorker|desired|registered|add|remove|upsert|revision|pod_uid|canonical_host_port' model_gateway/src/service_discovery.rs
printf '%s\n' '--- all exact reconciliation and registry symbols ---'
rg -n -C8 'ReconcileActions|compute_actions|OwnedWorker|expected_revision|RemoveWorker|WorkerRegistrationMode|RegistrationMode|Upsert' model_gateway/src model_gateway/tests crates/protocols/src --glob '*.rs' | head -n 1600Repository: smg-project/smg
Length of output: 50372
🏁 Script executed:
#!/bin/bash
set -u
printf '%s\n' '--- workflow files ---'
git ls-files | rg '(^|/)workflow(/|.*\.(rs|toml)$)|job_queue|worker.*workflow|workflow.*worker'
printf '%s\n' '--- focused job and registration definitions ---'
rg -n -C12 'enum Job|RemoveWorker|enum WorkerRegistrationMode|WorkerRegistrationMode|AddWorker' --glob '*.rs' --glob '!**/tests/**' model_gateway crates | head -n 2000
printf '%s\n' '--- registry mutation APIs ---'
rg -n -C12 'fn (register|register_or_replace|remove|replace|upsert)|pub enum WorkerOrigin|origin_of|get_all_with_ids|revision\(' model_gateway/src/worker model_gateway/src/workflow crates --glob '*.rs' | head -n 2000Repository: smg-project/smg
Length of output: 50372
🏁 Script executed:
#!/bin/bash
set -u
printf '%s\n' '--- job queue structure ---'
ast-grep outline model_gateway/src/workflow/job_queue.rs
printf '%s\n' '--- job queue definitions and RemoveWorker execution ---'
rg -n -C18 'pub enum Job|RemoveWorker|AddWorker|expected_revision|remove_by_url|drain_settle' model_gateway/src/workflow/job_queue.rs
printf '%s\n' '--- registry removal and replacement definitions ---'
rg -n -C18 'fn remove_by_url|fn register_or_replace|fn register\(|fn remove\(|fn resolve_url_to_id|fn get_by_url|url_to_id|worker_id' model_gateway/src/worker/registry.rs | head -n 2200
printf '%s\n' '--- registration workflow expansion ---'
rg -n -C18 'dp_info|workers\s*=|register_worker|register_workers|create_worker_registration_workflow|expand' model_gateway/src/workflow/steps/local model_gateway/src/workflow/steps/shared model_gateway/src/workflow/engines.rs model_gateway/src/workflow/job_queue.rsRepository: smg-project/smg
Length of output: 50372
🏁 Script executed:
#!/bin/bash
set -u
printf '%s\n' '--- worker removal workflow files ---'
git ls-files model_gateway/src/workflow | rg 'remov|worker'
printf '%s\n' '--- removal workflow symbols ---'
rg -n -C20 'create_worker_removal_workflow_data|WorkerRemoval|DrainWorkersStep|expected_revision|remove_by_url|remove\(' model_gateway/src/workflow model_gateway/src/worker/registry.rs --glob '*.rs' | head -n 2600
printf '%s\n' '--- exact registry removal implementation ---'
sed -n '1108,1245p' model_gateway/src/worker/registry.rsRepository: smg-project/smg
Length of output: 50372
🏁 Script executed:
#!/bin/bash
set -u
printf '%s\n' '--- removal-related files and constructors ---'
git ls-files model_gateway/src | rg 'worker.*(remov|workflow)|remov.*worker|workflow.*(mod|engine|data)'
rg -l 'create_worker_removal_workflow_data|FindWorkersToRemoveStep|RemoveFromWorkerRegistryStep|DrainWorkersStep|expected_revision' model_gateway/src --glob '*.rs'
printf '%s\n' '--- removal workflow constructor references ---'
rg -n -C15 'create_worker_removal_workflow_data|FindWorkersToRemoveStep|RemoveFromWorkerRegistryStep|DrainWorkersStep' model_gateway/src --glob '*.rs' | head -n 1800
printf '%s\n' '--- find-worker and removal step bodies ---'
sed -n '1,260p' model_gateway/src/workflow/steps/local/find_workers_to_remove.rs
sed -n '1,260p' model_gateway/src/workflow/steps/local/remove_from_worker_registry.rs
sed -n '1,260p' model_gateway/src/workflow/steps/local/drain_workers.rsRepository: smg-project/smg
Length of output: 50373
🏁 Script executed:
#!/bin/bash
set -u
printf '%s\n' '--- URL matching and revision filtering ---'
sed -n '1,180p' model_gateway/src/workflow/steps/local/find_workers_to_remove.rs
sed -n '1,125p' model_gateway/src/workflow/steps/local/mod.rs
printf '%s\n' '--- registry removal step ---'
sed -n '1,150p' model_gateway/src/workflow/steps/local/remove_from_worker_registry.rs
printf '%s\n' '--- DP-aware configuration source ---'
rg -n -C8 'dp_aware|disaggregated_mode|service_discovery' model_gateway/src/config model_gateway/src/service_discovery.rs model_gateway/src/app_context.rs --glob '*.rs' | head -n 1800Repository: smg-project/smg
Length of output: 50373
🏁 Script executed:
#!/bin/bash
set -u
printf '%s\n' '--- revision lifecycle ---'
rg -n -C12 'revision\(|revision:|fetch_add|AtomicU64|inherit_shared_state_from|set_revision' model_gateway/src/worker model_gateway/src/workflow --glob '*.rs' | head -n 2200
printf '%s\n' '--- DP expansion and worker URL construction ---'
rg -n -C15 'dp_info|@\\{|format!.*@|worker_url.*@|rank.*url|data.?parallel|workers_to_remove' model_gateway/src/workflow/steps/local model_gateway/src/workflow/steps/shared model_gateway/src/worker --glob '*.rs' | head -n 2600
printf '%s\n' '--- existing DP removal tests ---'
rg -n -C20 'expected_revision|dp_aware|find_workers_by_url|remov.*rank|rank.*remov' model_gateway/src/workflow model_gateway/tests --glob '*.rs' | head -n 2200Repository: smg-project/smg
Length of output: 50372
🏁 Script executed:
#!/bin/bash
set -u
printf '%s\n' '--- concrete worker revision implementation ---'
rg -n -C10 'fn revision\(&self\)|fn inherit_shared_state_from\(&self' model_gateway/src/worker/worker.rs model_gateway/src/worker --glob '*.rs'
printf '%s\n' '--- local registration expansion implementation ---'
rg -n -C14 'DpInfo|dp_rank|rank|@|actual_workers|workers\s*=' model_gateway/src/workflow/steps/local/create_worker.rs model_gateway/src/workflow/steps/local/discover_dp.rs model_gateway/src/workflow/steps/mod.rs model_gateway/src/workflow/steps/shared/register.rs --glob '*.rs' | head -n 2600
printf '%s\n' '--- removal tests in focused files ---'
rg -n -C16 'find_workers_by_url|expected_revision|dp_aware|WorkerRemovalRequest' model_gateway/src/workflow/steps/local model_gateway/src/workflow/steps --glob '*.rs' | tail -n 1800Repository: smg-project/smg
Length of output: 50372
🟡 Nit: Preserve each DP-rank revision during removal. FindWorkersToRemoveStep removes only ranks matching the single expected_revision selected after deduplication. Group removals by (url, revision) or remove all ranks atomically. Otherwise, ranks with different revisions remain routable until the next discovery pass.
🤖 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 `@model_gateway/src/service_discovery.rs` around lines 742 - 763, Update
compute_actions so removal deduplication preserves every worker revision for
each URL instead of collapsing all ranks to one URL. Group removals by the URL
and revision used by FindWorkersToRemoveStep, or otherwise ensure all DP-rank
revisions are removed atomically while retaining the existing desired-state
filtering.
Replace the raw-watcher K8s service discovery (tracked-pod shadow set, periodic full-LIST sweep, manual backoff loops, deletion inference from deletion_timestamp on Modified events) with the informer pattern: a reflector Store feeds a level-triggered reconciler that diffs desired workers against the worker registry and submits Add/Remove jobs. Failed or missed work is retried on the next pass by construction. Pods may now run multiple engine servers: the smg.ai/worker-ports annotation lists the data ports (absent = single worker at the configured discovery port), and sglang.ai/bootstrap-port accepts an aligned list for PD pods. Discovery-created workers carry smg.ai/pod-name and smg.ai/pod-uid labels; the reconciler only touches workers bearing the pod-uid label, so manually added workers are never removed. Terminating pods leave the desired set at grace-period start. Removals carry expected_revision so duplicate submissions become graceful skips. Metric label values 'duplicate' and 'pod_deleted' are removed (all deregistrations now emit 'reconciled'); the discovery gauge counts desired workers instead of pods. Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
a1a1b44 to
2acb222
Compare
|
Note GitHub couldn't provide a complete incremental comparison for this pull request, so CodeRabbit is performing a full review instead. This review may take a little longer. |
Multi-agent correctness review — results and fixesRan a 16-agent review (7 lens-specific finders + adversarial verifiers who must trace each claim end-to-end through real code before upholding it; ~1.9M tokens). 14 raw findings → 5 confirmed, 2 refuted as pre-existing-and-improved, 5 low-severity capped unverified. All confirmed findings are fixed in the amended commit; every fix is mutation-tested where applicable. Confirmed and fixed
Also applied from the unverified batch: Refuted (for the record)Two claims that a stale removal can deregister a same-URL replacement mid-drain were traced as pre-existing shared-machinery behavior on main — where the old discovery submitted removals with Known remaining (accepted)Discovery counters record job submissions, not outcomes — now bounded by the one-in-flight gate but still submission-semantics; PD mode with disjoint role selectors watches unfiltered (pre-existing scope, now cheaper with pruning). Verification after fixes: fmt/clippy clean, |
Description
Rewrites K8s service discovery from the raw-watcher design to the informer pattern (roadmap #2031 §5):
watcher().default_backoff().reflect(writer)maintains a reflectorStore, and a level-triggered reconciler diffs the desired worker set (from the store snapshot) against the worker registry, submittingAddWorker/RemoveWorkerjobs for the gap. Every watch event (coalesced 1s) and a periodiccheck_intervaltick trigger one pass; the periodic pass is now an in-memory diff with zero API calls.Deleted (the watcher-gap machinery): the
tracked_podsshadow set and its (name, uid) identity impls, the periodic full-LISTreconcile_podssweep + panic-restart supervisor, both manual watcher restart/backoff loops,.applied_objects()+deletion_timestampdelete inference, same-name/new-UID StatefulSet eviction, and duplicate-event suppression. Net −439 lines while adding the feature below.Multi-port pods (new): a pod running multiple engine servers lists them in the
smg.ai/worker-portsannotation (e.g."8080,8081,8082,8083"); each port becomes an independent worker with its own health/circuit-breaker/policy state. Absent annotation = single worker at--service-discovery-port(fully backward compatible).sglang.ai/bootstrap-portaccepts an aligned list for PD pods (one value broadcasts; N values zip; mismatch warns and is ignored).Ownership correlation (fixes the scenario in #2031 comment by @neverCase): every discovery-created worker carries
smg.ai/pod-name/smg.ai/pod-uidlabels through the label pipeline. The reconciler only ever adds/removes workers bearing the pod-uid label, so manually added workers are untouchable — and a worker registered by a zombie workflow after its pod died is now removed on the next pass instead of living forever. Failed registrations are retried by construction (registry lacks the worker → re-added next pass).Behavioral notes:
expected_revisionas a stale-guard. The reflector prunesmanagedFieldsbefore the Store.Event::Delete(hard-deleted routers marked Down promptly) but deliberately no Store — mesh SWIM/CRDT owns convergence.Breaking change
Test Plan
cargo +nightly fmt --all --checkclean;cargo clippy --all-targets -- -D warningsexit 0cargo test -p smg: 1930 passed, 0 failed (21 binaries);cargo test --workspace --exclude smg: 2362 passed, 0 failed (80 binaries)cargo build -p smg-python/-p smg-golanggreen; Python bindings exposeworker_ports_annotation(both constructors, RouterArgs, docs)Stacked follow-ups: scripted fake-API-server integration tests (
model_gateway/tests/), then kind-based real-cluster e2e.