From d6297cd4efd3ed33115db69c85778ea9b1dbe800 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Fabr=C3=ADcio=20Bracht?= Date: Fri, 2 Oct 2026 16:06:44 -0700 Subject: [PATCH 1/2] default dev clusters to a full peer mesh and wait for a settled partition map --- crates/mqdb-cli/src/cli_types/dev.rs | 2 +- crates/mqdb-cli/src/commands/dev/cluster.rs | 2 +- crates/mqdb-cli/src/commands/dev/tests.rs | 138 ++++++++++++++++++-- docs/distributed-design.md | 8 +- 4 files changed, 131 insertions(+), 19 deletions(-) diff --git a/crates/mqdb-cli/src/cli_types/dev.rs b/crates/mqdb-cli/src/cli_types/dev.rs index 7b807ee2..8eba1477 100644 --- a/crates/mqdb-cli/src/cli_types/dev.rs +++ b/crates/mqdb-cli/src/cli_types/dev.rs @@ -83,7 +83,7 @@ pub(crate) enum DevAction { #[arg( long, value_name = "TYPE", - help = "Topology: partial (default), upper, or full" + help = "Topology: full (default), partial, or upper" )] topology: Option, #[arg( diff --git a/crates/mqdb-cli/src/commands/dev/cluster.rs b/crates/mqdb-cli/src/commands/dev/cluster.rs index 19313886..d772a4e5 100644 --- a/crates/mqdb-cli/src/commands/dev/cluster.rs +++ b/crates/mqdb-cli/src/commands/dev/cluster.rs @@ -53,7 +53,7 @@ pub(crate) fn cmd_dev_start_cluster( }; let passwd_path = passwd.unwrap_or_else(|| generated_passwd.as_deref().expect("just created")); - let topology_name = topology.unwrap_or("partial"); + let topology_name = topology.unwrap_or("full"); println!("Using {topology_name} mesh topology (bind: {bind_host})"); for node_id in 1..=nodes { diff --git a/crates/mqdb-cli/src/commands/dev/tests.rs b/crates/mqdb-cli/src/commands/dev/tests.rs index 743bf804..c6d05451 100644 --- a/crates/mqdb-cli/src/commands/dev/tests.rs +++ b/crates/mqdb-cli/src/commands/dev/tests.rs @@ -42,9 +42,13 @@ fn wait_for_cluster_ready(nodes: u8, timeout_secs: u64) -> bool { } if all_ready { - std::thread::sleep(Duration::from_secs(2)); - println!("Cluster ready!"); - return true; + let settled = wait_for_partitions_settled(nodes, timeout_secs); + if settled { + println!("Cluster ready!"); + } else { + println!("Warning: running tests while partitions may still be moving"); + } + return settled; } std::thread::sleep(Duration::from_millis(500)); @@ -537,6 +541,13 @@ fn run_pubsub_test(pub_port: u16, sub_port: u16, topic: &str, msg: &str) -> bool } } +fn full_mesh_peers(node_id: u8, nodes: u8) -> Vec { + (1..=nodes) + .filter(|&n| n != node_id) + .map(|n| format!("{}@127.0.0.1:{}", n, 1882 + u16::from(n))) + .collect() +} + fn wait_for_auth_cluster(nodes: u8, timeout_secs: u64) -> bool { let start = Instant::now(); let timeout = Duration::from_secs(timeout_secs); @@ -572,8 +583,7 @@ fn wait_for_auth_cluster(nodes: u8, timeout_secs: u64) -> bool { } if all_ready { - std::thread::sleep(Duration::from_secs(2)); - return true; + return wait_for_partitions_settled(nodes, timeout_secs); } std::thread::sleep(Duration::from_millis(500)); @@ -582,6 +592,112 @@ fn wait_for_auth_cluster(nodes: u8, timeout_secs: u64) -> bool { false } +type PartitionAssignment = (u64, Option, Vec); + +fn partition_map_on(port: u16) -> Option> { + let output = Command::new("timeout") + .args([ + "3", + "mosquitto_rr", + "-h", + "127.0.0.1", + "-p", + &port.to_string(), + "-u", + "admin", + "-P", + "admin", + "-t", + "$SYS/mqdb/cluster/status", + "-e", + &format!("resp/dev-settle-{port}"), + "-m", + "{}", + "-W", + "2", + ]) + .output() + .ok()?; + let status: serde_json::Value = serde_json::from_slice(&output.stdout).ok()?; + let data = &status["data"]; + let fully_linked = data["unlinked_nodes"].as_array().is_some_and(Vec::is_empty); + if !fully_linked { + return None; + } + let mut assignments: Vec = data["partitions"] + .as_array()? + .iter() + .map(|p| { + let replicas = p["replicas"] + .as_array() + .map(|r| r.iter().filter_map(serde_json::Value::as_u64).collect()) + .unwrap_or_default(); + ( + p["id"].as_u64().unwrap_or(0), + p["primary"].as_u64(), + replicas, + ) + }) + .collect(); + assignments.sort_unstable(); + Some(assignments) +} + +fn partition_map_complete(map: &[PartitionAssignment], nodes: u8) -> bool { + map.len() == usize::from(mqdb_core::NUM_PARTITIONS) + && map + .iter() + .all(|(_, primary, replicas)| primary.is_some() && (nodes < 2 || !replicas.is_empty())) + && (1..=u64::from(nodes)) + .all(|node| map.iter().any(|(_, primary, _)| *primary == Some(node))) +} + +fn wait_for_partitions_settled(nodes: u8, timeout_secs: u64) -> bool { + const STABLE_POLLS: u32 = 3; + let start = Instant::now(); + let timeout = Duration::from_secs(timeout_secs.max(30)); + let mut previous: Option> = None; + let mut stable_polls = 0; + + while start.elapsed() < timeout { + let maps: Vec>> = (0..nodes) + .map(|i| partition_map_on(1883 + u16::from(i))) + .collect(); + let agreed = match maps.first() { + Some(Some(first)) if partition_map_complete(first, nodes) => maps + .iter() + .all(|m| m.as_ref() == Some(first)) + .then(|| first.clone()), + _ => None, + }; + match agreed { + Some(map) if previous.as_ref() == Some(&map) => stable_polls += 1, + Some(map) => { + previous = Some(map); + stable_polls = 1; + } + None => { + previous = None; + stable_polls = 0; + } + } + if stable_polls >= STABLE_POLLS { + println!( + "Partition map settled on all {nodes} nodes after {:.1}s", + start.elapsed().as_secs_f64() + ); + return true; + } + std::thread::sleep(Duration::from_secs(1)); + } + + println!( + "Warning: partition map did not settle within {}s", + timeout.as_secs() + ); + false +} + #[allow(clippy::too_many_lines)] fn run_test_ownership(nodes: u8, _ports: &[u16], license: Option<&Path>) { println!("=== Ownership Enforcement Test ({nodes} nodes) ===\n"); @@ -646,9 +762,7 @@ fn run_test_ownership(nodes: u8, _ports: &[u16], license: Option<&Path>) { cmd.args(["--license", lic_str]); } - let peers: Vec = (1..node_id) - .map(|n| format!("{}@127.0.0.1:{}", n, 1882 + u16::from(n))) - .collect(); + let peers = full_mesh_peers(node_id, nodes); if !peers.is_empty() { cmd.args(["--peers", &peers.join(",")]); } @@ -1032,9 +1146,7 @@ fn run_test_sharing(nodes: u8, _ports: &[u16], license: Option<&Path>) { cmd.args(["--license", lic_str]); } - let peers: Vec = (1..node_id) - .map(|n| format!("{}@127.0.0.1:{}", n, 1882 + u16::from(n))) - .collect(); + let peers = full_mesh_peers(node_id, nodes); if !peers.is_empty() { cmd.args(["--peers", &peers.join(",")]); } @@ -1495,9 +1607,7 @@ fn start_presence_cluster( cmd.args(["--license", lic_str]); } - let peers: Vec = (1..node_id) - .map(|n| format!("{}@127.0.0.1:{}", n, 1882 + u16::from(n))) - .collect(); + let peers = full_mesh_peers(node_id, nodes); if !peers.is_empty() { cmd.args(["--peers", &peers.join(",")]); } diff --git a/docs/distributed-design.md b/docs/distributed-design.md index 0eb476a9..7119c0a0 100644 --- a/docs/distributed-design.md +++ b/docs/distributed-design.md @@ -496,9 +496,11 @@ The fix detects bridge clients by their ID pattern `node-X-to-node-Y` (set in `c | Topology | Description | Bridges (3 nodes) | |----------|-------------|-------------------| -| `partial` (default) | To lower-numbered nodes | N1:0, N2:1, N3:2 | +| `full` (default) | All-to-all (duplicates) | N1:2, N2:2, N3:2 | +| `partial` | To lower-numbered nodes | N1:0, N2:1, N3:2 | | `upper` | To higher-numbered nodes | N1:2, N2:1, N3:0 | -| `full` | All-to-all (duplicates) | N1:2, N2:2, N3:2 | + +`full` is the default because it is the only topology in which every node is configured with every other node, which the cluster requires. `partial` and `upper` start nodes with an incomplete peer list; they remain available to reproduce membership problems such as #155 (a node that starts with no peers elects itself and two Raft leaders can then be elected in one term). **Bridge Direction Semantics**: @@ -593,7 +595,7 @@ Full mesh topology experiences Raft leader flapping (every 4-5s) due to unreliab - Heartbeats are NOT RECEIVED by followers consistently - Election timeout (3-5s) expires, triggering new elections -Use **partial** or **upper** topology for production until this is resolved. +This applied to the deprecated MQTT bridge transport and was resolved by making QUIC the default (11.16). Do **not** use `partial` or `upper` for production: every node must be configured with every other node (#155). **Key Findings**: - **0 peers = best pubsub**: Nodes with 0 bridge connections achieve 146-154k msg/s vs 5-15k for nodes with peers From 5b495db1a13a2e936b522fa6a23dba3c9a26a2a5 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Fabr=C3=ADcio=20Bracht?= Date: Fri, 2 Oct 2026 19:25:45 -0700 Subject: [PATCH 2/2] keep partial as the default dev topology without quic --- crates/mqdb-cli/src/cli_types/dev.rs | 2 +- crates/mqdb-cli/src/commands/dev/cluster.rs | 2 +- docs/distributed-design.md | 6 +++--- 3 files changed, 5 insertions(+), 5 deletions(-) diff --git a/crates/mqdb-cli/src/cli_types/dev.rs b/crates/mqdb-cli/src/cli_types/dev.rs index 8eba1477..47cac7d3 100644 --- a/crates/mqdb-cli/src/cli_types/dev.rs +++ b/crates/mqdb-cli/src/cli_types/dev.rs @@ -83,7 +83,7 @@ pub(crate) enum DevAction { #[arg( long, value_name = "TYPE", - help = "Topology: full (default), partial, or upper" + help = "Topology: full (default with QUIC), partial (default with --no-quic), or upper" )] topology: Option, #[arg( diff --git a/crates/mqdb-cli/src/commands/dev/cluster.rs b/crates/mqdb-cli/src/commands/dev/cluster.rs index d772a4e5..d8e3c1d0 100644 --- a/crates/mqdb-cli/src/commands/dev/cluster.rs +++ b/crates/mqdb-cli/src/commands/dev/cluster.rs @@ -53,7 +53,7 @@ pub(crate) fn cmd_dev_start_cluster( }; let passwd_path = passwd.unwrap_or_else(|| generated_passwd.as_deref().expect("just created")); - let topology_name = topology.unwrap_or("full"); + let topology_name = topology.unwrap_or(if no_quic { "partial" } else { "full" }); println!("Using {topology_name} mesh topology (bind: {bind_host})"); for node_id in 1..=nodes { diff --git a/docs/distributed-design.md b/docs/distributed-design.md index 7119c0a0..0ae2ff6a 100644 --- a/docs/distributed-design.md +++ b/docs/distributed-design.md @@ -496,11 +496,11 @@ The fix detects bridge clients by their ID pattern `node-X-to-node-Y` (set in `c | Topology | Description | Bridges (3 nodes) | |----------|-------------|-------------------| -| `full` (default) | All-to-all (duplicates) | N1:2, N2:2, N3:2 | -| `partial` | To lower-numbered nodes | N1:0, N2:1, N3:2 | +| `full` (default with QUIC) | All-to-all (duplicates) | N1:2, N2:2, N3:2 | +| `partial` (default with `--no-quic`) | To lower-numbered nodes | N1:0, N2:1, N3:2 | | `upper` | To higher-numbered nodes | N1:2, N2:1, N3:0 | -`full` is the default because it is the only topology in which every node is configured with every other node, which the cluster requires. `partial` and `upper` start nodes with an incomplete peer list; they remain available to reproduce membership problems such as #155 (a node that starts with no peers elects itself and two Raft leaders can then be elected in one term). +`full` is the default with QUIC because it is the only topology in which every node is configured with every other node, which the cluster requires. With the deprecated MQTT bridge transport (`--no-quic`) the default stays `partial`, since a full mesh of bridges either amplifies (`Both`) or flaps (`Out`, issue 11.16). `partial` and `upper` start nodes with an incomplete peer list; they remain available to reproduce membership problems such as #155 (a node that starts with no peers elects itself and two Raft leaders can then be elected in one term). **Bridge Direction Semantics**: