diff --git a/CHANGELOG.md b/CHANGELOG.md index 72a75ae0..ee1b92c8 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,21 @@ All notable changes to this project will be documented in this file. Each entry lists the date and the crate versions that were released. +## 2026-09-20 — mqdb-cluster 0.4.15, mqdb-cli 0.8.39 + +### Fixed + +- **A cluster started from a single seed node no longer fails silently.** Inter-node messages travel over direct peer connections and are never relayed, so the cluster has always required a full peer mesh — but the README told users to start node 2 and node 3 each with only `--peers "1@..."`, which leaves node 2 and node 3 unable to reach each other. On such a cluster, cross-node publishes between them are dropped, and a read on one cannot see data whose partition primary is the other, with no error, warning or metric anywhere. The documentation now shows every node peering with all the others and states the requirement next to `--peers`. + +### Added + +- **Incomplete-mesh detection.** Each node compares the cluster members it knows about (from the partition map and voter gossip) against the peers it actually has a connection to, and logs a warning naming every unreachable node once a minute while the mesh is incomplete. `mqdb cluster status` reports the same as an `UNLINKED` line, and the status payload gained `known_members` and `unlinked_nodes`. Verified on a live 3-node hub-and-spoke cluster: the two spokes each name the other, the hub stays silent, and a full mesh produces no warning. + +### Notes + +- This makes the failure visible; it does not remove it. A node still cannot dial a peer it was not configured with, because nothing propagates node addresses — auto-dialing discovered nodes is a follow-up that needs addresses in the gossip. +- Detection covers a node that was **never linked**, which is the misconfiguration case. A link that dies in flight is not yet detected: a peer is never removed from the connection map when its stream fails, so it still counts as linked. Removing it safely needs a per-connection token (a reconnecting peer reuses the same node id and must not have its fresh connection dropped by the old one's cleanup), which belongs with the auto-dial and dial-retry work. + ## 2026-09-19 — mqdb-agent 0.8.29, mqdb-cluster 0.4.14, mqdb-cli 0.8.38 ### Added diff --git a/Cargo.lock b/Cargo.lock index 94b59c8b..ae973298 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1455,7 +1455,7 @@ dependencies = [ [[package]] name = "mqdb-cli" -version = "0.8.38" +version = "0.8.39" dependencies = [ "base64", "bebytes", @@ -1482,7 +1482,7 @@ dependencies = [ [[package]] name = "mqdb-cluster" -version = "0.4.14" +version = "0.4.15" dependencies = [ "arc-swap", "bebytes", diff --git a/README.md b/README.md index 78fd7e04..8e8d8d2b 100644 --- a/README.md +++ b/README.md @@ -702,19 +702,35 @@ MQDB supports distributed clustering with automatic failover and partition rebal --quic-cert test_certs/server.pem --quic-key test_certs/server.key --quic-ca test_certs/ca.pem \ --license /path/to/license.key -# Node 2 (joins via peer) +# Node 2 (peers with node 1) ./target/release/mqdb cluster start --node-id 2 --bind 127.0.0.1:1884 \ --db /tmp/mqdb-node2 --peers "1@127.0.0.1:1883" \ --quic-cert test_certs/server.pem --quic-key test_certs/server.key --quic-ca test_certs/ca.pem \ --license /path/to/license.key -# Node 3 (joins via peer) +# Node 3 (peers with BOTH existing nodes) ./target/release/mqdb cluster start --node-id 3 --bind 127.0.0.1:1885 \ - --db /tmp/mqdb-node3 --peers "1@127.0.0.1:1883" \ + --db /tmp/mqdb-node3 --peers "1@127.0.0.1:1883,2@127.0.0.1:1884" \ --quic-cert test_certs/server.pem --quic-key test_certs/server.key --quic-ca test_certs/ca.pem \ --license /path/to/license.key ``` +> **Every node must end up peered with every other node.** Inter-node messages are delivered over +> direct peer connections and are never relayed, so a node with no direct link to another silently +> misses its subscriptions, presence and cross-node publishes, and cannot read data whose partition +> primary it is. +> +> A node dials only the nodes listed in its `--peers`, and only at startup — it never retries a +> failed dial and never dials a node it learns about later. Because a dial only succeeds against a +> node that is already listening, list the nodes started before it, as above; the links to nodes +> started afterwards are created when *those* nodes dial in. **This also applies to restarts:** a +> node restarted with only its original seed peer will not be re-dialled by the nodes that started +> before it, so restart it with `--peers` listing every other node in the cluster. +> +> Run `mqdb cluster status` after starting or restarting a node — it reports any `UNLINKED` nodes — +> and the broker logs a warning every minute while the mesh is incomplete. Note that this detects a +> node that was never linked; a link that dies in flight is not currently detected. + ### Cluster CLI Commands ```bash @@ -749,7 +765,7 @@ mosquitto_pub -h 127.0.0.1 -p 1883 -t "events/test" -m "hello" -i publisher1 | `--node-name` | Human-readable node name | | `--bind` | MQTT listener address (default: 0.0.0.0:1883) | | `--db` | Database directory path | -| `--peers` | Peer nodes to join (format: id@host:port) | +| `--peers` | All other cluster nodes, comma-separated (format: id@host:port). The cluster requires a full mesh: every node must be peered with every other node. | | `--quic-cert` | TLS certificate for QUIC transport (must have serverAuth + clientAuth EKU) | | `--quic-key` | TLS private key for QUIC transport | | `--quic-ca` | CA certificate for mTLS peer verification (required for mTLS) | diff --git a/crates/mqdb-cli/Cargo.toml b/crates/mqdb-cli/Cargo.toml index 42d8a50f..5675158f 100644 --- a/crates/mqdb-cli/Cargo.toml +++ b/crates/mqdb-cli/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "mqdb-cli" -version = "0.8.38" +version = "0.8.39" publish = false edition.workspace = true license = "AGPL-3.0-only" diff --git a/crates/mqdb-cli/src/commands/cluster.rs b/crates/mqdb-cli/src/commands/cluster.rs index 467565c7..32b331be 100644 --- a/crates/mqdb-cli/src/commands/cluster.rs +++ b/crates/mqdb-cli/src/commands/cluster.rs @@ -336,6 +336,29 @@ pub(crate) async fn cmd_cluster_status( println!("│ Nodes: 1 (this node only) │"); } + if let Some(unlinked) = data + .get("unlinked_nodes") + .and_then(serde_json::Value::as_array) + && !unlinked.is_empty() + { + let nodes: Vec = unlinked + .iter() + .filter_map(serde_json::Value::as_u64) + .map(|n| n.to_string()) + .collect(); + let mut listed = nodes.join(", "); + if listed.chars().count() > 20 { + listed = format!("{} +{} more", nodes[..3].join(", "), nodes.len() - 3); + } + println!("├─────────────────────────────────────────┤"); + println!( + "│ {:<39} │", + format!("UNLINKED: [{listed}] — mesh incomplete") + ); + println!("│ {:<39} │", "start every node with --peers listing"); + println!("│ {:<39} │", "all other nodes; see README clustering"); + } + if let Some(partitions) = data.get("partitions").and_then(serde_json::Value::as_array) { let mut primary_counts: HashMap = HashMap::new(); let mut with_replicas = 0; diff --git a/crates/mqdb-cluster/Cargo.toml b/crates/mqdb-cluster/Cargo.toml index 8d937fc2..fb3a9a33 100644 --- a/crates/mqdb-cluster/Cargo.toml +++ b/crates/mqdb-cluster/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "mqdb-cluster" -version = "0.4.14" +version = "0.4.15" publish = false edition.workspace = true license = "AGPL-3.0-only" diff --git a/crates/mqdb-cluster/src/cluster/db_handler/tests.rs b/crates/mqdb-cluster/src/cluster/db_handler/tests.rs index f49d6646..af8ff805 100644 --- a/crates/mqdb-cluster/src/cluster/db_handler/tests.rs +++ b/crates/mqdb-cluster/src/cluster/db_handler/tests.rs @@ -92,6 +92,10 @@ impl ClusterTransport for MockTransport { Ok(()) } + async fn direct_peers(&self) -> Option> { + None + } + fn recv(&self) -> Option { self.inbox.lock().unwrap().pop_front() } diff --git a/crates/mqdb-cluster/src/cluster/mqtt_transport.rs b/crates/mqdb-cluster/src/cluster/mqtt_transport.rs index c5ecde32..ba44a847 100644 --- a/crates/mqdb-cluster/src/cluster/mqtt_transport.rs +++ b/crates/mqdb-cluster/src/cluster/mqtt_transport.rs @@ -422,6 +422,10 @@ impl ClusterTransport for MqttTransport { Err(TransportError::PartitionNotFound(partition)) } + async fn direct_peers(&self) -> Option> { + None + } + fn recv(&self) -> Option { if let Ok(mut requeue) = self.requeue_buffer.try_lock() && let Some(msg) = requeue.pop_front() diff --git a/crates/mqdb-cluster/src/cluster/node_controller/mod.rs b/crates/mqdb-cluster/src/cluster/node_controller/mod.rs index 52261703..29c03bed 100644 --- a/crates/mqdb-cluster/src/cluster/node_controller/mod.rs +++ b/crates/mqdb-cluster/src/cluster/node_controller/mod.rs @@ -649,6 +649,37 @@ impl NodeController { self.heartbeat.alive_nodes() } + /// Cluster members this node knows about, excluding itself: every node that owns a partition + /// as primary or replica, plus any node named in voter gossip. A node can appear here without + /// being reachable, which is exactly the case `unlinked_nodes` reports. + pub fn known_members(&self) -> Vec { + let mut members: std::collections::BTreeSet = + self.partition_map().all_nodes().into_iter().collect(); + members.extend(self.heartbeat.voters().iter().copied()); + members.remove(&self.node_id); + members.into_iter().collect() + } + + /// Known cluster members this node has no direct link to. Broadcasts are a single hop and + /// targeted sends are not routed, so any node listed here silently misses subscriptions, + /// presence, client locations and forwarded publishes. + pub async fn unlinked_nodes(&self) -> Vec { + self.unlinked_from(&self.known_members()).await + } + + /// `unlinked_nodes` against an already-computed member set, so a caller that needs both does + /// not walk every partition twice. + pub async fn unlinked_from(&self, members: &[NodeId]) -> Vec { + let Some(linked) = self.transport().direct_peers().await else { + return Vec::new(); + }; + members + .iter() + .copied() + .filter(|node| !linked.contains(node)) + .collect() + } + pub fn become_primary(&mut self, partition: PartitionId, epoch: Epoch) { let state = self .replicas diff --git a/crates/mqdb-cluster/src/cluster/node_controller/tests.rs b/crates/mqdb-cluster/src/cluster/node_controller/tests.rs index b5a8cdcd..88bd42cc 100644 --- a/crates/mqdb-cluster/src/cluster/node_controller/tests.rs +++ b/crates/mqdb-cluster/src/cluster/node_controller/tests.rs @@ -57,6 +57,7 @@ struct MockTransport { node_id: NodeId, inbox: Arc>>, outbox: Arc>>, + linked: Arc>>>, } impl MockTransport { @@ -65,6 +66,7 @@ impl MockTransport { node_id, inbox: Arc::new(Mutex::new(VecDeque::new())), outbox: Arc::new(Mutex::new(Vec::new())), + linked: Arc::new(Mutex::new(None)), } } @@ -114,6 +116,10 @@ impl ClusterTransport for MockTransport { Ok(()) } + async fn direct_peers(&self) -> Option> { + self.linked.lock().unwrap().clone() + } + fn recv(&self) -> Option { self.inbox.lock().unwrap().pop_front() } @@ -2557,3 +2563,51 @@ async fn three_way_circular_fk_cascade_terminates() { "should find referencing records in the cycle" ); } + +#[tokio::test] +async fn unlinked_nodes_names_only_members_missing_a_link() { + let node1 = NodeId::validated(1).unwrap(); + let node2 = NodeId::validated(2).unwrap(); + let node3 = NodeId::validated(3).unwrap(); + + let transport = MockTransport::new(node1); + *transport.linked.lock().unwrap() = Some(vec![node2]); + let mut ctrl = create_test_controller(node1, transport); + ctrl.seed_unique_voters([node1, node2, node3].into_iter().collect()); + + assert_eq!( + ctrl.unlinked_nodes().await, + vec![node3], + "node3 is a known member with no direct connection, node2 is linked, self is excluded" + ); +} + +#[tokio::test] +async fn unlinked_nodes_is_empty_for_a_broker_mediated_transport() { + let node1 = NodeId::validated(1).unwrap(); + let node2 = NodeId::validated(2).unwrap(); + let node3 = NodeId::validated(3).unwrap(); + + let mut ctrl = create_test_controller(node1, MockTransport::new(node1)); + ctrl.seed_unique_voters([node1, node2, node3].into_iter().collect()); + + assert_eq!( + ctrl.known_members(), + vec![node2, node3], + "voter gossip tells a node which members exist, excluding itself" + ); + + assert!( + ctrl.unlinked_nodes().await.is_empty(), + "a broker-mediated transport reports no direct peers, so nothing can be called unlinked" + ); +} + +#[tokio::test] +async fn unlinked_nodes_is_empty_without_voter_gossip() { + let node1 = NodeId::validated(1).unwrap(); + let ctrl = create_test_controller(node1, MockTransport::new(node1)); + + assert!(ctrl.known_members().is_empty()); + assert!(ctrl.unlinked_nodes().await.is_empty()); +} diff --git a/crates/mqdb-cluster/src/cluster/quic_transport.rs b/crates/mqdb-cluster/src/cluster/quic_transport.rs index a0942716..713cd1ad 100644 --- a/crates/mqdb-cluster/src/cluster/quic_transport.rs +++ b/crates/mqdb-cluster/src/cluster/quic_transport.rs @@ -363,6 +363,10 @@ impl ClusterTransport for QuicDirectTransport { Err(TransportError::PartitionNotFound(partition)) } + async fn direct_peers(&self) -> Option> { + Some(self.peers.read().await.keys().copied().collect()) + } + fn recv(&self) -> Option { if let Ok(mut requeue) = self.requeue_buffer.try_lock() && let Some(msg) = requeue.pop_front() diff --git a/crates/mqdb-cluster/src/cluster/raft/coordinator/tests.rs b/crates/mqdb-cluster/src/cluster/raft/coordinator/tests.rs index 5435ec05..1e092fe3 100644 --- a/crates/mqdb-cluster/src/cluster/raft/coordinator/tests.rs +++ b/crates/mqdb-cluster/src/cluster/raft/coordinator/tests.rs @@ -54,6 +54,10 @@ impl ClusterTransport for MockTransport { Ok(()) } + async fn direct_peers(&self) -> Option> { + None + } + fn recv(&self) -> Option { None } diff --git a/crates/mqdb-cluster/src/cluster/transport.rs b/crates/mqdb-cluster/src/cluster/transport.rs index 15d56cbc..bf7a4f5b 100644 --- a/crates/mqdb-cluster/src/cluster/transport.rs +++ b/crates/mqdb-cluster/src/cluster/transport.rs @@ -487,6 +487,10 @@ pub trait ClusterTransport: Send + Sync + Debug + Clone { message: ClusterMessage, ) -> impl std::future::Future> + Send; + /// Node ids this transport can deliver to directly, or `None` when the transport is + /// broker-mediated and has no notion of a direct peer. + fn direct_peers(&self) -> impl std::future::Future>> + Send; + fn recv(&self) -> Option; fn pending_count(&self) -> usize; diff --git a/crates/mqdb-cluster/src/cluster_agent/admin.rs b/crates/mqdb-cluster/src/cluster_agent/admin.rs index a8b81516..ee5b1584 100644 --- a/crates/mqdb-cluster/src/cluster_agent/admin.rs +++ b/crates/mqdb-cluster/src/cluster_agent/admin.rs @@ -241,6 +241,14 @@ impl ClusteredAgent { let raft_status = self.rx_raft_status.borrow().clone(); let alive_nodes: Vec = ctrl.alive_nodes().iter().map(|n| n.get()).collect(); + let members = ctrl.known_members(); + let unlinked_nodes: Vec = ctrl + .unlinked_from(&members) + .await + .iter() + .map(|n| n.get()) + .collect(); + let known_members: Vec = members.iter().map(|n| n.get()).collect(); let mut partitions = Vec::new(); for partition in PartitionId::all() { @@ -274,6 +282,8 @@ impl ClusteredAgent { "store_client_locations": stores.client_locations.len(), "store_db_data": stores.db_data.len(), "alive_nodes": alive_nodes, + "known_members": known_members, + "unlinked_nodes": unlinked_nodes, "partition_count": NUM_PARTITIONS, "partitions": partitions })) diff --git a/crates/mqdb-cluster/src/cluster_agent/event_loop.rs b/crates/mqdb-cluster/src/cluster_agent/event_loop.rs index b901b505..373ce2e9 100644 --- a/crates/mqdb-cluster/src/cluster_agent/event_loop.rs +++ b/crates/mqdb-cluster/src/cluster_agent/event_loop.rs @@ -1,7 +1,7 @@ // Copyright 2025-2026 LabOverWire. All rights reserved. // SPDX-License-Identifier: AGPL-3.0-only -use super::{AdminRequest, ClusteredAgent}; +use super::{AdminRequest, ClusteredAgent, MESH_CHECK_INTERVAL_SECS}; use crate::cluster::{ ClusterMessage, ClusterTransport, NodeController, NodeId, PartitionId, ProcessingBatch, TopicSubscriptionBroadcast, WildcardBroadcast, @@ -78,6 +78,7 @@ impl ClusteredAgent { let mut ttl_cleanup_interval = interval(Duration::from_secs(TTL_CLEANUP_INTERVAL_SECS)); let mut wildcard_reconciliation_interval = interval(Duration::from_mins(1)); let mut subscription_reconciliation_interval = interval(Duration::from_mins(5)); + let mut mesh_check_interval = interval(Duration::from_secs(MESH_CHECK_INTERVAL_SECS)); let mut retained_sync_cleanup_interval = interval(Duration::from_secs(RETAINED_SYNC_CLEANUP_INTERVAL_SECS)); let mut cascade_retry_interval = tokio::time::interval_at( @@ -145,6 +146,9 @@ impl ClusteredAgent { _ = subscription_reconciliation_interval.tick() => { self.handle_subscription_reconciliation().await; } + _ = mesh_check_interval.tick() => { + self.warn_on_unlinked_nodes().await; + } _ = retained_sync_cleanup_interval.tick() => { Self::handle_retained_sync_cleanup(&synced_retained_topics).await; } @@ -553,6 +557,24 @@ impl ClusteredAgent { } } + async fn warn_on_unlinked_nodes(&self) { + let unlinked = { + let ctrl = self.controller.read().await; + ctrl.unlinked_nodes().await + }; + if unlinked.is_empty() { + return; + } + let nodes: Vec = unlinked.iter().map(|node| node.get()).collect(); + tracing::warn!( + unlinked_nodes = ?nodes, + "cluster mesh is incomplete: these known nodes have no direct connection to this node. \ + Broadcasts are a single hop and targeted sends are not routed, so subscriptions, \ + presence, client locations and cross-node publishes will not reach them. \ + Start every node with --peers listing all other nodes." + ); + } + async fn handle_wildcard_reconciliation(&self) { let now = current_time_ms(); let ctrl = self.controller.read().await; diff --git a/crates/mqdb-cluster/src/cluster_agent/mod.rs b/crates/mqdb-cluster/src/cluster_agent/mod.rs index b05f4abf..6f7f408e 100644 --- a/crates/mqdb-cluster/src/cluster_agent/mod.rs +++ b/crates/mqdb-cluster/src/cluster_agent/mod.rs @@ -27,6 +27,7 @@ const CLEANUP_INTERVAL_SECS: u64 = 3600; const TTL_CLEANUP_INTERVAL_SECS: u64 = 60; const UNIQUE_RECONCILE_INTERVAL_SECS: u64 = 10; const RETAINED_SYNC_CLEANUP_INTERVAL_SECS: u64 = 30; +const MESH_CHECK_INTERVAL_SECS: u64 = 60; const RETAINED_SYNC_TTL_SECS: u64 = 5; const RAFT_CHANNEL_CAPACITY: usize = 4096; diff --git a/crates/mqdb-cluster/src/cluster_agent/transport.rs b/crates/mqdb-cluster/src/cluster_agent/transport.rs index 61c77403..fff3e9b7 100644 --- a/crates/mqdb-cluster/src/cluster_agent/transport.rs +++ b/crates/mqdb-cluster/src/cluster_agent/transport.rs @@ -97,6 +97,15 @@ impl ClusterTransport for ClusterTransportKind { } } + async fn direct_peers(&self) -> Option> { + match self { + #[cfg(feature = "mqtt-bridge")] + #[allow(deprecated)] + Self::Mqtt(t) => t.direct_peers().await, + Self::Quic(t) => t.direct_peers().await, + } + } + fn recv(&self) -> Option { match self { #[cfg(feature = "mqtt-bridge")] diff --git a/crates/mqdb-cluster/tests/simulation/transport.rs b/crates/mqdb-cluster/tests/simulation/transport.rs index e325553f..82247335 100644 --- a/crates/mqdb-cluster/tests/simulation/transport.rs +++ b/crates/mqdb-cluster/tests/simulation/transport.rs @@ -304,6 +304,10 @@ impl ClusterTransport for SimulatedTransport { } } + async fn direct_peers(&self) -> Option> { + None + } + fn recv(&self) -> Option { let msg = self.network.receive(self.node_id.get())?; let message = Self::deserialize_message(&msg.payload)?; diff --git a/docs/distributed-design.md b/docs/distributed-design.md index 476d9b69..cf7a613d 100644 --- a/docs/distributed-design.md +++ b/docs/distributed-design.md @@ -1508,12 +1508,17 @@ mqdb cluster start \ ### 10.2 Join Existing Cluster +`--peers` must list **every** other node, not just one seed. Messages travel over direct peer +connections and are never relayed, so any node pair without a direct link cannot exchange +subscriptions, presence, client locations or forwarded publishes, and cannot serve each other's +partitions. + ```bash mqdb cluster start \ - --node-id 2 \ - --bind 127.0.0.1:1884 \ - --db /var/lib/mqdb/node2 \ - --peers "1@127.0.0.1:1883" + --node-id 3 \ + --bind 127.0.0.1:1885 \ + --db /var/lib/mqdb/node3 \ + --peers "1@127.0.0.1:1883,2@127.0.0.1:1884" ``` ### 10.3 Bridge Configuration diff --git a/docs/testing/11-cluster-tests.md b/docs/testing/11-cluster-tests.md index 60b26d06..9057a046 100644 --- a/docs/testing/11-cluster-tests.md +++ b/docs/testing/11-cluster-tests.md @@ -250,9 +250,10 @@ mqdb dev kill --node 3 mqdb create partition_test --data '{"during": "partition"}' \ --broker 127.0.0.1:1883 --user admin --pass admin -# 5. Restart Node 3 +# 5. Restart Node 3 — list every other node, otherwise the nodes that started +# before it will not re-dial it and the 2-3 link stays down mqdb cluster start --node-id 3 --bind 127.0.0.1:1885 --db /tmp/mqdb-test-3 \ - --peers 1@127.0.0.1:1883 \ + --peers 1@127.0.0.1:1883,2@127.0.0.1:1884 \ --quic-cert test_certs/server.pem --quic-key test_certs/server.key --quic-ca test_certs/ca.pem \ --passwd /tmp/mqdb-test-passwd --admin-users admin \ --license /path/to/license.key &