Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion crates/mqdb-cli/src/cli_types/dev.rs
Original file line number Diff line number Diff line change
Expand Up @@ -83,7 +83,7 @@ pub(crate) enum DevAction {
#[arg(
long,
value_name = "TYPE",
help = "Topology: partial (default), upper, or full"
help = "Topology: full (default with QUIC), partial (default with --no-quic), or upper"
)]
topology: Option<String>,
#[arg(
Expand Down
2 changes: 1 addition & 1 deletion crates/mqdb-cli/src/commands/dev/cluster.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(if no_quic { "partial" } else { "full" });
println!("Using {topology_name} mesh topology (bind: {bind_host})");

for node_id in 1..=nodes {
Expand Down
138 changes: 124 additions & 14 deletions crates/mqdb-cli/src/commands/dev/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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));
Expand Down Expand Up @@ -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<String> {
(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);
Expand Down Expand Up @@ -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));
Expand All @@ -582,6 +592,112 @@ fn wait_for_auth_cluster(nodes: u8, timeout_secs: u64) -> bool {
false
}

type PartitionAssignment = (u64, Option<u64>, Vec<u64>);

fn partition_map_on(port: u16) -> Option<Vec<PartitionAssignment>> {
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<PartitionAssignment> = 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<Vec<PartitionAssignment>> = None;
let mut stable_polls = 0;

while start.elapsed() < timeout {
let maps: Vec<Option<Vec<PartitionAssignment>>> = (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");
Expand Down Expand Up @@ -646,9 +762,7 @@ fn run_test_ownership(nodes: u8, _ports: &[u16], license: Option<&Path>) {
cmd.args(["--license", lic_str]);
}

let peers: Vec<String> = (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(",")]);
}
Expand Down Expand Up @@ -1032,9 +1146,7 @@ fn run_test_sharing(nodes: u8, _ports: &[u16], license: Option<&Path>) {
cmd.args(["--license", lic_str]);
}

let peers: Vec<String> = (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(",")]);
}
Expand Down Expand Up @@ -1495,9 +1607,7 @@ fn start_presence_cluster(
cmd.args(["--license", lic_str]);
}

let peers: Vec<String> = (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(",")]);
}
Expand Down
8 changes: 5 additions & 3 deletions docs/distributed-design.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 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` | All-to-all (duplicates) | N1:2, N2:2, N3:2 |

`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**:

Expand Down Expand Up @@ -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
Expand Down
Loading