refactor(mesh): simplify code and fix flaky integration tests - #628
Conversation
|
Note Reviews pausedIt looks like this branch is under active development. To avoid overwhelming you with review comments due to an influx of new commits, CodeRabbit has automatically paused this review. You can configure this behavior by changing the Use the following commands to manage reviews:
Use the checkboxes below for quick actions:
📝 WalkthroughWalkthroughRemoved actor/owner arguments from store Changes
Sequence Diagram(s)sequenceDiagram
participant Test as "Test"
participant Server as "MeshServer"
participant Gossip as "GossipService"
participant Listener as "TcpListener"
participant Stores as "State Stores"
Test->>Server: build server & pass TcpListener
Server->>Server: start_inner(Some(listener))
Server->>Gossip: spawn serve_ping_with_listener(listener, shutdown)
Gossip->>Listener: accept connections
Listener-->>Gossip: connection stream
Gossip->>Server: deliver ping/update
Server->>Stores: apply snapshot/updates (insert without actor)
Stores-->>Server: ack
Estimated code review effort🎯 4 (Complex) | ⏱️ ~45 minutes Possibly related PRs
Suggested reviewers
Poem
🚥 Pre-merge checks | ✅ 3✅ Passed checks (3 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches🧪 Generate unit tests (beta)
Tip Try Coding Plans. Let us write the prompt for your AI agent so you can ship faster (with fewer bugs). Comment |
Summary of ChangesHello, I'm Gemini Code Assist1! I'm currently reviewing this pull request and will post my feedback shortly. In the meantime, here's a summary to help you and other reviewers quickly get up to speed! This pull request refactors the mesh service by removing an unused Highlights
Changelog
Activity
Using Gemini Code AssistThe full guide for Gemini Code Assist can be found on our documentation page, here are some quick tips. Invoking Gemini You can request assistance from Gemini at any point by creating a comment using either
Customization To customize Gemini Code Assist for GitHub experience, repository maintainers can create a configuration file and/or provide a custom code review style guide (such as PEP-8 for Python) by creating and adding files to a Limitations & Feedback Gemini Code Assist may make mistakes. Please leave feedback on any instances where its feedback is incorrect or counter productive. You can react with 👍 and 👎 on @gemini-code-assist comments. If you're interested in giving your feedback about your experience with Gemini Code Assist for Github and other Google products, sign up here. You can also get AI-powered code generation, chat, as well as code reviews directly in the IDE at no cost with the Gemini Code Assist IDE Extension. Footnotes
|
There was a problem hiding this comment.
Code Review
This is a strong pull request that significantly improves the codebase by simplifying logic and enhancing test reliability. The removal of the unused _actor parameter across numerous files is a great cleanup. The fixes for flaky tests, particularly replacing sleep with a polling mechanism and addressing the TOCTOU port race, are excellent and demonstrate a deep understanding of testing concurrent systems. The other refactorings, like introducing helper functions and optimizing data structures, further improve code quality. I have one suggestion to reduce code duplication in one of the refactored tests.
There was a problem hiding this comment.
Actionable comments posted: 4
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@mesh/src/service.rs`:
- Around line 659-677: Extract the helper functions bind_node and wait_for from
mesh/src/service.rs into a shared test utilities module (e.g., tests/util or
mesh/tests/common) and replace the duplicated implementations in
mesh/src/tests/comprehensive.rs with imports from that module; specifically move
the async fn bind_node() -> (TcpListener, SocketAddr) and async fn
wait_for<F>(condition: F, timeout: Duration, msg: &str) where F: Fn() -> bool
into the common test utils, update both files to use the shared functions via
use/import, and ensure any needed visibility (pub(crate) or pub) and required
imports (TcpListener, SocketAddr, Duration, tokio) are adjusted so tests
compile.
- Around line 385-390: The start_with_listener path can advertise self.self_addr
while actually listening on a different address; modify
start_inner/start_with_listener to validate that when a listener is passed (the
TcpListener from start_with_listener), listener.local_addr()? (use
tokio::net::TcpListener::local_addr()) matches self.self_addr (SocketAddr)
before proceeding; if the addresses differ, return an Err (or propagate a clear
error) rather than starting the server so callers cannot advertise a mismatched
address; update start_with_listener to call start_inner with the listener and
ensure the validation occurs at the top of start_inner.
In `@mesh/src/stores.rs`:
- Around line 223-231: The matching in StoreType::to_proto should not use
hardcoded integers; instead map to the generated proto enum values to keep a
single source of truth. Change the implementation of StoreType::to_proto to
convert each variant to the corresponding generated proto enum (e.g.,
proto::StoreType::Membership, proto::StoreType::App, etc.) and return its as i32
(or implement/from the generated enum) so the numeric discriminants always come
from the proto definition rather than literals.
In `@mesh/src/tests/comprehensive.rs`:
- Around line 377-388: The hardcoded tokio::time::sleep should be replaced with
a condition-based wait that ensures the incremental sync cycle and its mark_sent
bookkeeping have completed before the second update; specifically, replace the
sleep with a wait_for that checks handler_b.stores.app.get("shared_key") exists
and matches the version/value currently held by handler_a (e.g., compare the
stored value/version from handler_b to handler_a) so you only proceed when
handler_b has fully observed the first write rather than waiting a fixed 2s.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: b52c7375-23e3-49c6-971a-c47f821c4274
📒 Files selected for processing (13)
mesh/src/controller.rsmesh/src/crdt_kv/crdt.rsmesh/src/crdt_kv/operation.rsmesh/src/crdt_kv/replica.rsmesh/src/crdt_kv/tests.rsmesh/src/incremental.rsmesh/src/node_state_machine.rsmesh/src/ping_server.rsmesh/src/rate_limit_window.rsmesh/src/service.rsmesh/src/stores.rsmesh/src/sync.rsmesh/src/tests/comprehensive.rs
💤 Files with no reviewable changes (1)
- mesh/src/node_state_machine.rs
What changed: - stores.rs: Remove unused `_actor: String` parameter from `define_state_store!` macro's `insert` method - controller.rs, ping_server.rs, sync.rs, service.rs, node_state_machine.rs, rate_limit_window.rs, incremental.rs, tests/comprehensive.rs: Update all ~36 `insert` call sites to drop the removed actor argument - crdt_kv/operation.rs: Change `snapshot_and_truncate` return type from BTreeMap to HashMap (matches `latest_operations_by_key`) - crdt_kv/crdt.rs: Replace 4 copy-pasted 5-arm StoreType match blocks with `StoreType::to_proto()` helper - crdt_kv/replica.rs: Relax Lamport clock atomics from SeqCst to AcqRel/Acquire (sufficient for single-counter monotonicity) - crdt_kv/tests.rs: Minor test cleanup - stores.rs: Add `StoreType::to_proto()` and `StoreType::from_proto()` conversion methods; add `policy_key()` helper - ping_server.rs: Add `serve_ping_with_listener` accepting a pre-bound TcpListener via `TcpIncoming` to eliminate TOCTOU port race in tests - service.rs: Add `MeshServer::start_with_listener` and 4-arg `mesh_run!` macro variant threading the listener through; rewrite test infrastructure with `bind_node()` and `wait_for()` helpers; remove dead `print_state` function - tests/comprehensive.rs: Replace `get_node_addr()`/`find_free_port()` with `bind_node()` returning `(TcpListener, SocketAddr)`; add `wait_for()` polling helper (100ms interval); convert all integration tests from sleep-then-assert to poll-with-timeout; un-ignore `test_three_node_cluster_formation` and `test_multi_node_data_propagation`; keep failure-detection tests ignored with clearer rationale; remove dead `print_cluster_state` Why: - The `_actor` parameter was vestigial — never read, just stored as `_actor` with a backwards-compat comment. Removing it simplifies every store write call. - BTreeMap→HashMap: `latest_operations_by_key` results are iterated without needing key ordering, so HashMap's O(1) amortized insert is more appropriate than BTreeMap's O(log n). - Integration tests were flaky due to a TOCTOU port race (`TcpListener` dropped before tonic could bind) and fixed `sleep(N)` + assert patterns that were timing-sensitive. How: - Port race eliminated by passing the pre-bound TcpListener all the way through to tonic's `serve_with_incoming_shutdown` via `TcpIncoming::from(TcpListener)`. - Sleep-based assertions replaced with `wait_for(|| condition, timeout, msg)` that polls every 100ms, making tests both faster (no unnecessary waiting) and more reliable (generous timeouts). - Tests relying on SWIM failure detection for hard-shutdown nodes remain ignored because they need many gossip rounds and are inherently timing-sensitive under parallel CI load. Refs: mesh crate simplification and test reliability Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
66dee58 to
9b29a61
Compare
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 9b29a619c6
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
| pub fn insert( | ||
| &self, | ||
| key: String, | ||
| value: $value_type, | ||
| _actor: String, | ||
| ) -> Result<Option<$value_type>, serde_json::Error> { |
There was a problem hiding this comment.
Restore insert-call compatibility across workspace crates
Changing StateStore::insert to two parameters here leaves existing workspace callers uncompilable; I checked model_gateway/tests/mesh_integration_test.rs and it still calls stores.membership.insert(..., actor) and stores.app.insert(..., actor) (for example around lines 108, 129, and 376), which now produces arity/type errors when model_gateway tests are built. Please update those call sites in the same change or provide a temporary compatibility path.
Useful? React with 👍 / 👎.
| pub fn apply_remote_rate_limit_counter_value_with_actor_and_timestamp( | ||
| &self, | ||
| key: String, |
There was a problem hiding this comment.
Keep removed rate-limit actor helper or migrate callers
This refactor removes the apply_remote_rate_limit_counter_value_with_actor convenience API, but downstream callers were not migrated; I checked model_gateway/tests/mesh_integration_test.rs and it still invokes that method multiple times (around lines 160-190), so building those tests now fails with a missing-method error. Either restore a forwarding wrapper or update all callers in the same commit.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
♻️ Duplicate comments (3)
mesh/src/stores.rs (1)
222-231:⚠️ Potential issue | 🟠 MajorAvoid hardcoded proto discriminants in
StoreType::to_proto.These literals can drift from the wire enum and introduce subtle cross-version incompatibilities.
♻️ Suggested fix
impl StoreType { /// Convert to proto StoreType (i32) pub fn to_proto(self) -> i32 { match self { - StoreType::Membership => 0, - StoreType::App => 1, - StoreType::Worker => 2, - StoreType::Policy => 3, - StoreType::RateLimit => 4, + StoreType::Membership => super::service::gossip::StoreType::Membership as i32, + StoreType::App => super::service::gossip::StoreType::App as i32, + StoreType::Worker => super::service::gossip::StoreType::Worker as i32, + StoreType::Policy => super::service::gossip::StoreType::Policy as i32, + StoreType::RateLimit => super::service::gossip::StoreType::RateLimit as i32, } } }🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@mesh/src/stores.rs` around lines 222 - 231, The function StoreType::to_proto uses hardcoded numeric discriminants; replace those literals with the generated proto enum variants to avoid drift — update StoreType::to_proto to map each Rust variant to the corresponding generated proto enum (e.g. proto::StoreType::Membership, ::App, ::Worker, ::Policy, ::RateLimit) and return their as i32 (or use proto::StoreType::from(...) / Into<i32> if available) instead of raw numbers so the wire values are driven by the proto definition rather than hardcoded constants.mesh/src/tests/comprehensive.rs (1)
373-384:⚠️ Potential issue | 🟠 MajorRemove the fixed 2s sleep from the sync path.
This hardcoded pause is still timing-sensitive under CI load and can reintroduce flakes.
♻️ Suggested refactor
- wait_for( - || handler_b.stores.app.get("shared_key").is_some(), - Duration::from_secs(15), - "shared_key synced to node B", - ) - .await; - - // Allow the sync cycle to settle before writing a second update. - // The incremental collector runs on a 1s interval and the mark_sent - // bookkeeping must complete before the next version can be detected. - tokio::time::sleep(Duration::from_secs(2)).await; + wait_for( + || { + handler_b + .stores + .app + .get("shared_key") + .is_some_and(|v| v.value == b"shared_value") + }, + Duration::from_secs(15), + "shared_key synced to node B with expected value", + ) + .await;🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@mesh/src/tests/comprehensive.rs` around lines 373 - 384, Remove the hardcoded tokio::time::sleep(Duration::from_secs(2)).await and instead poll for the sync bookkeeping to finish using the existing wait_for helper: after the initial wait_for confirming "shared_key synced to node B", call wait_for with a closure that asserts the incremental-mark_sent bookkeeping/collector has completed (for example, compare versions or timestamps between handler_a.stores.app.get("shared_key") and handler_b.stores.app.get("shared_key") or wait until handler_b sees the same version/value that handler_a has) so the next update can be detected reliably; replace the sleep line with this wait_for-based check referencing the wait_for function and handler_a/handler_b stores.mesh/src/service.rs (1)
385-390:⚠️ Potential issue | 🟠 MajorValidate listener address against advertised
self_addrbefore startup.
start_with_listenercan still serve on one socket while advertising another address, which can break peer connectivity. Please reject startup when they differ.Suggested fix
pub async fn start_with_listener(self, listener: tokio::net::TcpListener) -> Result<()> { self.start_inner(Some(listener)).await } async fn start_inner(self, listener: Option<tokio::net::TcpListener>) -> Result<()> { + if let Some(ref tcp_listener) = listener { + let bound_addr = tcp_listener + .local_addr() + .map_err(|e| anyhow::anyhow!("Failed to read listener local addr: {e}"))?; + if bound_addr != self.self_addr { + return Err(anyhow::anyhow!( + "Listener/self_addr mismatch: listener={}, self_addr={}", + bound_addr, + self.self_addr + )); + } + } log::info!("Mesh server listening on {}", self.self_addr); let self_name = self.self_name.clone(); let self_address = self.self_addr;Also applies to: 423-431
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@mesh/src/service.rs` around lines 385 - 390, Ensure the provided TcpListener's bound socket address matches the service's advertised self_addr before allowing startup: in start_with_listener (and similarly in start_inner call sites), retrieve the listener local_addr() and compare it to self_addr (or normalize host/port as your config requires); if they differ, return an Err early with a clear message instead of proceeding to start_inner/start, so the server cannot serve on a socket that mismatches its advertised address.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Duplicate comments:
In `@mesh/src/service.rs`:
- Around line 385-390: Ensure the provided TcpListener's bound socket address
matches the service's advertised self_addr before allowing startup: in
start_with_listener (and similarly in start_inner call sites), retrieve the
listener local_addr() and compare it to self_addr (or normalize host/port as
your config requires); if they differ, return an Err early with a clear message
instead of proceeding to start_inner/start, so the server cannot serve on a
socket that mismatches its advertised address.
In `@mesh/src/stores.rs`:
- Around line 222-231: The function StoreType::to_proto uses hardcoded numeric
discriminants; replace those literals with the generated proto enum variants to
avoid drift — update StoreType::to_proto to map each Rust variant to the
corresponding generated proto enum (e.g. proto::StoreType::Membership, ::App,
::Worker, ::Policy, ::RateLimit) and return their as i32 (or use
proto::StoreType::from(...) / Into<i32> if available) instead of raw numbers so
the wire values are driven by the proto definition rather than hardcoded
constants.
In `@mesh/src/tests/comprehensive.rs`:
- Around line 373-384: Remove the hardcoded
tokio::time::sleep(Duration::from_secs(2)).await and instead poll for the sync
bookkeeping to finish using the existing wait_for helper: after the initial
wait_for confirming "shared_key synced to node B", call wait_for with a closure
that asserts the incremental-mark_sent bookkeeping/collector has completed (for
example, compare versions or timestamps between
handler_a.stores.app.get("shared_key") and
handler_b.stores.app.get("shared_key") or wait until handler_b sees the same
version/value that handler_a has) so the next update can be detected reliably;
replace the sleep line with this wait_for-based check referencing the wait_for
function and handler_a/handler_b stores.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: 27801e11-307e-493c-b411-a4804a076b5d
📒 Files selected for processing (13)
mesh/src/controller.rsmesh/src/crdt_kv/crdt.rsmesh/src/crdt_kv/operation.rsmesh/src/crdt_kv/replica.rsmesh/src/crdt_kv/tests.rsmesh/src/incremental.rsmesh/src/node_state_machine.rsmesh/src/ping_server.rsmesh/src/rate_limit_window.rsmesh/src/service.rsmesh/src/stores.rsmesh/src/sync.rsmesh/src/tests/comprehensive.rs
💤 Files with no reviewable changes (1)
- mesh/src/node_state_machine.rs
- Use proto enum values in StoreType::to_proto() instead of hardcoded integers, referencing gossip::StoreType variants directly - Add listener/self_addr mismatch validation in start_with_listener - Move bind_node/wait_for helpers to shared test_utils module and make it pub(crate) so both service.rs and comprehensive.rs tests can import from a single source - Deduplicate wait_for calls in test_state_synchronization by extracting check_statuses closure and iterating over handlers - Restore apply_remote_rate_limit_counter_value and apply_remote_rate_limit_counter_value_with_actor convenience methods in sync.rs that were accidentally removed - Remove actor argument from model_gateway insert call sites to match updated store API - Tighten first wait_for in data sync test to check exact value match instead of just existence Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
There was a problem hiding this comment.
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (1)
mesh/src/stores.rs (1)
222-250:⚠️ Potential issue | 🟠 MajorDo not coerce unknown proto store types to
Membership.
from_protocurrently defaults unknown values toMembership, which can route updates into the wrong store instead of rejecting/ignoring invalid input.🛡️ Proposed fix
- pub fn from_proto(proto_value: i32) -> Self { + pub fn from_proto(proto_value: i32) -> Option<Self> { match proto_value { - 0 => StoreType::Membership, - 1 => StoreType::App, - 2 => StoreType::Worker, - 3 => StoreType::Policy, - 4 => StoreType::RateLimit, + 0 => Some(StoreType::Membership), + 1 => Some(StoreType::App), + 2 => Some(StoreType::Worker), + 3 => Some(StoreType::Policy), + 4 => Some(StoreType::RateLimit), unknown => { tracing::warn!( proto_value = unknown, - "Unknown StoreType proto value, defaulting to Membership" + "Unknown StoreType proto value" ); - StoreType::Membership + None } } }🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@mesh/src/stores.rs` around lines 222 - 250, from_proto currently coerces unknown proto integers into StoreType::Membership which can misroute updates; change the signature of from_proto(proto_value: i32) to return Option<StoreType> (or Result<StoreType, SomeError>) and map known values to Some(StoreType::...), and on unknown values log a warning and return None (or Err) instead of defaulting to Membership; then update all call sites that use from_proto to handle the None/Err path by rejecting or ignoring the invalid input rather than inserting into the Membership store.
♻️ Duplicate comments (1)
mesh/src/tests/comprehensive.rs (1)
360-364:⚠️ Potential issue | 🟠 MajorRemove the remaining fixed 2s sleep from sync verification flow.
This hardcoded delay is still timing-sensitive and can flake under parallel CI load; keep this path condition-driven like the rest of the file.
♻️ Proposed refactor
- // Allow the sync cycle to settle before writing a second update. - // The incremental collector runs on a 1s interval and the mark_sent - // bookkeeping must complete before the next version can be detected. - tokio::time::sleep(Duration::from_secs(2)).await; + // Ensure first write is fully converged before issuing second update. + wait_for( + || { + let a = handler_a.stores.app.get("shared_key"); + let b = handler_b.stores.app.get("shared_key"); + matches!( + (a, b), + (Some(a_state), Some(b_state)) + if a_state.version == b_state.version && a_state.value == b_state.value + ) + }, + Duration::from_secs(15), + "first shared_key write fully converged before second update", + ) + .await;🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@mesh/src/tests/comprehensive.rs` around lines 360 - 364, Remove the fixed tokio::time::sleep(Duration::from_secs(2)).await and replace it with a condition-driven wait that polls for the actual sync condition (e.g., the next version being detected or the mark_sent bookkeeping completing) with a bounded timeout; implement this by looping with short sleeps (e.g., tokio::time::sleep(Duration::from_millis(100)).await) and breaking when the expected condition is true, or fail the test if tokio::time::timeout elapses — target the sync verifier used in this test (the incremental collector / mark_sent check) instead of a fixed 2s delay.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Outside diff comments:
In `@mesh/src/stores.rs`:
- Around line 222-250: from_proto currently coerces unknown proto integers into
StoreType::Membership which can misroute updates; change the signature of
from_proto(proto_value: i32) to return Option<StoreType> (or Result<StoreType,
SomeError>) and map known values to Some(StoreType::...), and on unknown values
log a warning and return None (or Err) instead of defaulting to Membership; then
update all call sites that use from_proto to handle the None/Err path by
rejecting or ignoring the invalid input rather than inserting into the
Membership store.
---
Duplicate comments:
In `@mesh/src/tests/comprehensive.rs`:
- Around line 360-364: Remove the fixed
tokio::time::sleep(Duration::from_secs(2)).await and replace it with a
condition-driven wait that polls for the actual sync condition (e.g., the next
version being detected or the mark_sent bookkeeping completing) with a bounded
timeout; implement this by looping with short sleeps (e.g.,
tokio::time::sleep(Duration::from_millis(100)).await) and breaking when the
expected condition is true, or fail the test if tokio::time::timeout elapses —
target the sync verifier used in this test (the incremental collector /
mark_sent check) instead of a fixed 2s delay.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: c3e9a2a7-e91e-4c69-9a1f-6061d06a191f
📒 Files selected for processing (7)
mesh/src/service.rsmesh/src/stores.rsmesh/src/sync.rsmesh/src/tests/comprehensive.rsmesh/src/tests/mod.rsmesh/src/tests/test_utils.rsmodel_gateway/tests/mesh_integration_test.rs
💤 Files with no reviewable changes (1)
- model_gateway/tests/mesh_integration_test.rs
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@mesh/src/tests/test_utils.rs`:
- Around line 33-40: The current wait_for helper (using start, timeout,
condition, msg) can return true after the timeout because it only asserts
timeout while condition() is false; change it to compute a strict deadline
(e.g., let deadline = start + timeout) and check the current time against that
deadline before invoking or accepting condition() so any condition that becomes
true after the deadline fails; update the loop to panic/error with the same
message when now >= deadline (or use start.elapsed() >= timeout) and optionally
sleep for the remaining time or a fixed short interval between checks to avoid
busy-waiting.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: 1f87f97f-a794-401d-969a-9a3e1a9790f9
📒 Files selected for processing (3)
mesh/src/controller.rsmesh/src/sync.rsmesh/src/tests/test_utils.rs
| let start = std::time::Instant::now(); | ||
| while !condition() { | ||
| assert!( | ||
| start.elapsed() <= timeout, | ||
| "Timeout after {timeout:?} waiting for: {msg}" | ||
| ); | ||
| tokio::time::sleep(Duration::from_millis(100)).await; | ||
| } |
There was a problem hiding this comment.
wait_for can pass even after the timeout boundary.
The loop checks timeout only while condition() is false, so a condition that flips true slightly after the deadline can still pass. That weakens the helper’s timeout contract.
⏱️ Suggested fix (strict timeout semantics)
pub async fn wait_for<F>(condition: F, timeout: Duration, msg: &str)
where
F: Fn() -> bool,
{
- let start = std::time::Instant::now();
- while !condition() {
- assert!(
- start.elapsed() <= timeout,
- "Timeout after {timeout:?} waiting for: {msg}"
- );
- tokio::time::sleep(Duration::from_millis(100)).await;
- }
+ let start = std::time::Instant::now();
+ let deadline = start + timeout;
+ loop {
+ if condition() {
+ return;
+ }
+ assert!(
+ std::time::Instant::now() <= deadline,
+ "Timeout after {timeout:?} waiting for: {msg}"
+ );
+ let remaining = deadline.saturating_duration_since(std::time::Instant::now());
+ tokio::time::sleep(remaining.min(Duration::from_millis(100))).await;
+ }
}🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@mesh/src/tests/test_utils.rs` around lines 33 - 40, The current wait_for
helper (using start, timeout, condition, msg) can return true after the timeout
because it only asserts timeout while condition() is false; change it to compute
a strict deadline (e.g., let deadline = start + timeout) and check the current
time against that deadline before invoking or accepting condition() so any
condition that becomes true after the deadline fails; update the loop to
panic/error with the same message when now >= deadline (or use start.elapsed()
>= timeout) and optionally sleep for the remaining time or a fixed short
interval between checks to avoid busy-waiting.
- Remove unused LocalStoreType import in controller.rs incremental sender block - Remove needless borrows on model_id in sync.rs policy_key calls - Replace manual panic-in-if with assert! in test_utils wait_for and inline format args Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
e16bdbd to
aeb7a84
Compare
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@mesh/src/tests/comprehensive.rs`:
- Around line 358-361: Replace the hardcoded
tokio::time::sleep(Duration::from_secs(2)).await with a condition-based wait
that polls for the incremental sync to complete (e.g., loop with
tokio::time::sleep or tokio::select and timeout) rather than a fixed delay;
check the actual sync/mark_sent state exposed by the incremental collector or
add a small helper like wait_for_sync_cycle()/wait_until_mark_sent() that
repeatedly queries the collector/state (with a configurable timeout and short
backoff) and returns once bookkeeping is complete, failing the test on timeout.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: 0d2b5e88-20e7-45aa-9957-b6e6503ac61e
📒 Files selected for processing (5)
mesh/src/controller.rsmesh/src/service.rsmesh/src/sync.rsmesh/src/tests/comprehensive.rsmesh/src/tests/test_utils.rs
| // Allow the sync cycle to settle before writing a second update. | ||
| // The incremental collector runs on a 1s interval and the mark_sent | ||
| // bookkeeping must complete before the next version can be detected. | ||
| tokio::time::sleep(Duration::from_secs(2)).await; |
There was a problem hiding this comment.
🧹 Nitpick | 🔵 Trivial
Hardcoded 2s sleep remains between write operations.
The comment explains the rationale (incremental collector runs on 1s interval, mark_sent bookkeeping must complete), but a condition-based wait would be more robust. Consider polling for a sync-cycle completion signal if available, though this may require additional instrumentation.
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@mesh/src/tests/comprehensive.rs` around lines 358 - 361, Replace the
hardcoded tokio::time::sleep(Duration::from_secs(2)).await with a
condition-based wait that polls for the incremental sync to complete (e.g., loop
with tokio::time::sleep or tokio::select and timeout) rather than a fixed delay;
check the actual sync/mark_sent state exposed by the incremental collector or
add a small helper like wait_for_sync_cycle()/wait_until_mark_sent() that
repeatedly queries the collector/state (with a configurable timeout and short
backoff) and returns once bookkeeping is complete, failing the test on timeout.
Summary
_actorparameter from storeinsertcalls (~36 call sites across 10 files)What changed
Code simplification
_actor: Stringparameter fromdefine_state_store!macro'sinsertmethod; addStoreType::to_proto()/from_proto()helpers andpolicy_key()helperinsertcall sites to drop the removed actor argumentsnapshot_and_truncatereturn type fromBTreeMaptoHashMap— key ordering isn't needed andHashMapgives O(1) amortized insertStoreTypematch blocks withStoreType::to_proto()SeqCsttoAcqRel/Acquire(sufficient for single-counter monotonicity)print_stateandprint_cluster_statefunctionsTest reliability fixes
serve_ping_with_listener()accepting a pre-boundTcpListenerviaTcpIncomingto eliminate TOCTOU port raceMeshServer::start_with_listener()and 4-argmesh_run!macro variant threading the listener throughget_node_addr()/find_free_port()withbind_node()returning(TcpListener, SocketAddr); addwait_for()polling helper (100ms interval with configurable timeout)sleep(N)+ assert towait_for(|| condition, timeout, msg)polling patternTest coverage changes
test_three_node_cluster_formation,test_multi_node_data_propagation— these only test cluster formation and data sync (no failure detection), now reliable with bind_node + pollingtest_state_synchronization,test_five_node_cluster_with_failure— these rely on SWIM failure detection for hard-shutdown nodes which needs many gossip rounds and is inherently timing-sensitive under parallel CI load. Both pass reliably when run in isolation.Why
The
_actorparameter was never read — just stored as_actorwith a backwards-compatibility comment. Removing it simplifies every store write.Integration tests were flaky due to two root causes:
TcpListenerwas dropped before tonic could bind, creating a window where another process could claim the portsleep(2-5s)then asserted, which was both slow and unreliableHow
TcpListeneralive and passing it through to tonic'sserve_with_incoming_shutdownviaTcpIncoming::from(TcpListener)wait_for(|| condition, timeout, msg)that polls every 100ms — tests are both faster (no unnecessary waiting) and more reliable (generous timeouts up to 60s)Test plan
cargo check -p smg-meshpasses (only pre-existing unused import warning)cargo test -p smg-mesh— 147 passed, 0 failed, 2 ignored (5/5 consecutive green runs)--ignoredSummary by CodeRabbit
New Features
Improvements
Tests