Repository navigation
test(discovery): end-to-end reconcile tests against a scripted fake API server - #2062
Conversation
📝 WalkthroughSummary by CodeRabbit
WalkthroughThe change adds test-only Kubernetes client injection and end-to-end service-discovery tests. The tests cover pod LIST/WATCH handling, worker registration, reconciliation, watch recovery, deletion draining, registration retry, and router mesh-cluster state updates. ChangesKubernetes service discovery testing
Estimated code review effort: 4 (Complex) | ~45 minutes Sequence Diagram(s)sequenceDiagram
participant FakeKubernetesApi
participant ServiceDiscovery
participant WorkerRegistry
participant RouterMeshState
FakeKubernetesApi->>ServiceDiscovery: Serve pod LIST and WATCH responses
ServiceDiscovery->>WorkerRegistry: Register or reconcile pod-owned workers
FakeKubernetesApi->>ServiceDiscovery: Emit pod changes or watch failures
ServiceDiscovery->>FakeKubernetesApi: Reconnect or re-LIST
ServiceDiscovery->>WorkerRegistry: Drain or remove stale workers
ServiceDiscovery->>RouterMeshState: Set mesh cluster Alive or Down
Possibly related PRs
Suggested reviewers: 🚥 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.
Thorough review complete — no issues found. Clean, well-structured PR:
- Production refactoring: The
run_service_discoveryextraction is a clean mechanical split. Theenabledguard stays in the public entry point, return types are correct, andringinitialization is properly placed. - Feature gating:
test-utilis only activated by the self dev-dependency in[dev-dependencies]— production builds compile no test surface. - Test quality: The fake K8s API server faithfully models LIST + streaming WATCH with 410-Gone expiry. All four scenarios (multi-port lifecycle, graceful termination, zombie cleanup, desync convergence) are well-motivated and cover the key reconciler behaviors. Proper reuse of existing
common::mock_workerandcommon::create_test_contextinfrastructure.
0 🔴 Important · 0 🟡 Nit · 0 🟣 Pre-existing
a1a1b44 to
2acb222
Compare
9747100 to
1a898a3
Compare
1a898a3 to
763101b
Compare
| let _ = self | ||
| .state | ||
| .watch_tx | ||
| .lock() | ||
| .unwrap() | ||
| .send(format!("{status}\n")); | ||
| *self.state.watch_tx.lock().unwrap() = broadcast::channel(64).0; |
There was a problem hiding this comment.
🟡 Nit: Two separate lock acquisitions on watch_tx create a window where a concurrent watch request could subscribe to the old channel (between the send and the replace), then immediately get a Closed error instead of the 410. Unlikely to cause a flake in practice since the kube watcher must process the 410 before reconnecting, but holding a single guard is trivially safer:
| let _ = self | |
| .state | |
| .watch_tx | |
| .lock() | |
| .unwrap() | |
| .send(format!("{status}\n")); | |
| *self.state.watch_tx.lock().unwrap() = broadcast::channel(64).0; | |
| let mut guard = self.state.watch_tx.lock().unwrap(); | |
| let _ = guard.send(format!("{status}\n")); | |
| *guard = broadcast::channel(64).0; |
There was a problem hiding this comment.
🧹 Nitpick comments (2)
model_gateway/src/service_discovery.rs (1)
466-478: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low value🟡 Nit: Document the Tokio runtime requirement.
start_service_discovery_with_clientis synchronous, butrun_service_discoverycallstask::spawn. If a caller invokes this function outside a Tokio runtime,task::spawnpanics. Every current caller is a#[tokio::test]function, so the panic is not reachable today. Add a# Panicsnote so the requirement is visible at the call site.♻️ Proposed doc addition
/// Run discovery against an injected client. Compiled only for this crate's /// own integration tests, which point it at a scripted API server. +/// +/// # Panics +/// +/// Panics if called outside a Tokio runtime. #[cfg(feature = "test-util")] pub fn start_service_discovery_with_client(🤖 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 466 - 478, Add a Rustdoc # Panics section to start_service_discovery_with_client documenting that it panics when called outside an active Tokio runtime because run_service_discovery uses task::spawn. Keep the existing function behavior unchanged.model_gateway/tests/k8s_discovery_test.rs (1)
156-164: 📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick win🟡 Nit: Make dropped watch events fail loudly.
The
Laggedarm discards the number of skipped events and continues. If the broadcast buffer ever overflows, the watcher silently misses pod events. The test then fails with await_fortimeout that gives no indication that the fixture dropped the events.The current tests emit only a few events, so overflow is not reachable today. Surfacing the condition keeps future high-volume scenarios diagnosable.
♻️ Proposed change to surface dropped events
loop { match rx.recv().await { Ok(line) => return Some((Ok::<Bytes, Infallible>(Bytes::from(line)), rx)), - Err(broadcast::error::RecvError::Lagged(_)) => continue, + Err(broadcast::error::RecvError::Lagged(skipped)) => { + panic!("fake k8s watch stream dropped {skipped} events; increase the broadcast capacity"); + } Err(broadcast::error::RecvError::Closed) => return None, } }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/tests/k8s_discovery_test.rs` around lines 156 - 164, Update the futures::stream::unfold receiver handling so broadcast::error::RecvError::Lagged no longer silently continues; propagate or surface the skipped-event error through the stream using the existing error type, while preserving normal Ok line delivery and Closed termination behavior.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.
Nitpick comments:
In `@model_gateway/src/service_discovery.rs`:
- Around line 466-478: Add a Rustdoc # Panics section to
start_service_discovery_with_client documenting that it panics when called
outside an active Tokio runtime because run_service_discovery uses task::spawn.
Keep the existing function behavior unchanged.
In `@model_gateway/tests/k8s_discovery_test.rs`:
- Around line 156-164: Update the futures::stream::unfold receiver handling so
broadcast::error::RecvError::Lagged no longer silently continues; propagate or
surface the skipped-event error through the stream using the existing error
type, while preserving normal Ok line delivery and Closed termination behavior.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: CHILL
Plan: Pro Plus
Run ID: 2c508773-dc31-44ba-97d3-9a87865242d3
📒 Files selected for processing (3)
model_gateway/Cargo.tomlmodel_gateway/src/service_discovery.rsmodel_gateway/tests/k8s_discovery_test.rs
…PI server Add an integration binary that runs the real watcher, reflector store, reconciler, JobQueue, registration workflow, and worker registry against a scripted in-process pods endpoint (LIST + line-delimited WATCH with per-object resourceVersions, 410 expiry, and stream-sever support) and mock engine servers. Steady-state scenarios: multi-port pod registers one labeled worker per engine and removes all on pod deletion; a terminating pod drains at grace-period start while its object still exists; a zombie worker whose pod is gone is cleaned up while a manually added worker is untouched; live annotation edits grow and shrink the worker set; a router pod drives mesh cluster state Alive/Down through the injected ClusterState. Failure and recovery scenarios: a registration that keeps failing while the engine port is closed leaves no phantom worker and succeeds once the engine appears; a pod created after startup registers only when it turns Ready (pure watch-event path); a same-URL pod replacement converges to the new pod uid through the revision-guarded remove + re-register; a silent deletion during a 410 desync converges through the re-LIST; a severed watch stream reconnects and processes subsequent events; a 2s drain window is honored across reconcile passes ticking every 300ms (mutation-checked: disabling the reconciler's in-flight gate fails it). The client-injection seam (start_service_discovery_with_client) is compiled only under the new test-util feature, which a self dev-dependency enables for this crate's own test targets — production builds carry no test surface. Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
763101b to
596fde7
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. |
There was a problem hiding this comment.
🧹 Nitpick comments (3)
model_gateway/tests/k8s_discovery_test.rs (3)
173-181: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low value🟡 Nit: The watch stream drops lagged events silently.
RecvError::Laggedis skipped withcontinue. A lagged receiver means the fixture dropped scripted pod events, and the affected test then fails as an opaque 15-20s timeout instead of reporting the cause. Capacity is 64 and each test sends a few events, so lag is unlikely, but a diagnostic here is cheap.Consider emitting the count, for example
eprintln!("fake k8s watch lagged, dropped {n} events"), before you continue.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/tests/k8s_discovery_test.rs` around lines 173 - 181, Update the watch stream’s RecvError::Lagged branch inside futures::stream::unfold to emit a diagnostic containing the number of dropped events before continuing; keep the existing retry behavior and Closed handling unchanged.Source: Coding guidelines
136-158: 📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick win🟡 Nit:
expire_watchcan drop the 410 event if no watch is subscribed yet.
sever_watch_and_await_reconnectwaits forreceiver_count() > 0before it returns.expire_watchhas no equivalent guard. If the reflector has finished the initial LIST but has not yet issued the WATCH request, the 410 line goes to a sender with zero receivers and is lost. The subsequent sender swap then leaves the watcher on a stream that never re-LISTs, andwatch_interruption_relists_and_convergesfails only after its 20s timeout with an unclear message.The window is small because the test first waits for registration, so this is a flake-hardening suggestion, not a current failure.
♻️ Wait for a subscriber before sending the expiry event
- fn expire_watch(&self) { + async fn expire_watch(&self) { + let deadline = std::time::Instant::now() + Duration::from_secs(10); + while self.state.watch_tx.lock().unwrap().receiver_count() == 0 { + assert!( + std::time::Instant::now() < deadline, + "no watch stream to expire" + ); + tokio::time::sleep(Duration::from_millis(25)).await; + } let status = serde_json::json!({Update the single call site at line 461 to
fake.expire_watch().await;.🤖 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/tests/k8s_discovery_test.rs` around lines 136 - 158, Make expire_watch asynchronous and wait until the watch_tx sender has at least one receiver before sending the 410 event; then replace the call site in watch_interruption_relists_and_converges with an awaited invocation. Preserve the existing event send and sender replacement behavior after the subscriber guard.
588-591: 🩺 Stability & Availability | 🔵 Trivial | 💤 Low value🟡 Nit: The reserve-then-release port pattern can race with concurrent tests.
The test binds an ephemeral port, releases it, and rebinds it more than four seconds later at line 618. Cargo runs the tests in this binary concurrently, and
start_mock_enginealso binds port 0. If the kernel assigns the released port to another test in that window, two failures become possible: the negative assertion at line 605 sees a healthy foreign listener, orengine.start().unwrap()at line 618 panics with an address-in-use error.The kernel normally cycles the ephemeral range before reuse, so this is unlikely. Consider starting
MockWorkerwithport: 0first, recording its port, stopping it, and reusing that port — or accept the risk and document it.🤖 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/tests/k8s_discovery_test.rs` around lines 588 - 591, Eliminate the reserve-then-release port race in the test around MockWorker and start_mock_engine by obtaining an assigned port from a started MockWorker configured with port 0, then stop it and reuse that recorded port for the later assertions and engine startup. Preserve the test’s intended free-port and negative-listener checks without relying on an unprotected four-second gap.
🤖 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.
Nitpick comments:
In `@model_gateway/tests/k8s_discovery_test.rs`:
- Around line 173-181: Update the watch stream’s RecvError::Lagged branch inside
futures::stream::unfold to emit a diagnostic containing the number of dropped
events before continuing; keep the existing retry behavior and Closed handling
unchanged.
- Around line 136-158: Make expire_watch asynchronous and wait until the
watch_tx sender has at least one receiver before sending the 410 event; then
replace the call site in watch_interruption_relists_and_converges with an
awaited invocation. Preserve the existing event send and sender replacement
behavior after the subscriber guard.
- Around line 588-591: Eliminate the reserve-then-release port race in the test
around MockWorker and start_mock_engine by obtaining an assigned port from a
started MockWorker configured with port 0, then stop it and reuse that recorded
port for the later assertions and engine startup. Preserve the test’s intended
free-port and negative-listener checks without relying on an unprotected
four-second gap.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: CHILL
Plan: Pro Plus
Run ID: 1503a637-49da-4712-85f6-cc489dbaceca
📒 Files selected for processing (3)
model_gateway/Cargo.tomlmodel_gateway/src/service_discovery.rsmodel_gateway/tests/k8s_discovery_test.rs
🚧 Files skipped from review as they are similar to previous changes (2)
- model_gateway/Cargo.toml
- model_gateway/src/service_discovery.rs
Stacked on #2061.
Description
Adds
model_gateway/tests/k8s_discovery_test.rs: the real watcher → reflectorStore→ reconciler →JobQueue→ registration workflow → worker registry pipeline runs end to end; only two things are doubles — a scripted in-process pods endpoint (LIST + line-delimited WATCH, with 410-Gone expiry to force re-LISTs) andtests/commonmock engine servers on ephemeral ports.Scenarios:
smg.ai/worker-ports: "p1,p2"registers one worker per engine, each stamped with the pod-uid ownership label; deleting the pod removes both.deletion_timestampon a still-Ready pod removes its worker while the pod object still exists (drain at grace-period start).The client-injection seam (
start_service_discovery_with_client) is gated behind a newtest-utilcargo feature that a self dev-dependency enables for this crate's own test targets only — production builds compile no test surface (cargo build -p smgverified without the feature).Test Plan
cargo test -p smg --test k8s_discovery_test: 13 passed, 0 failed (~5s)cargo test -p smg: 1943 passed, 0 failed (22 binaries)cargo clippy --all-targets -- -D warningsexit 0;cargo +nightly fmt --all --checkcleancargo build -p smg(no test-util) green — seam absent from production builds