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
15 changes: 15 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
4 changes: 2 additions & 2 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

24 changes: 20 additions & 4 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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) |
Expand Down
2 changes: 1 addition & 1 deletion crates/mqdb-cli/Cargo.toml
Original file line number Diff line number Diff line change
@@ -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"
Expand Down
23 changes: 23 additions & 0 deletions crates/mqdb-cli/src/commands/cluster.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<String> = 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<u64, usize> = HashMap::new();
let mut with_replicas = 0;
Expand Down
2 changes: 1 addition & 1 deletion crates/mqdb-cluster/Cargo.toml
Original file line number Diff line number Diff line change
@@ -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"
Expand Down
4 changes: 4 additions & 0 deletions crates/mqdb-cluster/src/cluster/db_handler/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -92,6 +92,10 @@ impl ClusterTransport for MockTransport {
Ok(())
}

async fn direct_peers(&self) -> Option<Vec<NodeId>> {
None
}

fn recv(&self) -> Option<InboundMessage> {
self.inbox.lock().unwrap().pop_front()
}
Expand Down
4 changes: 4 additions & 0 deletions crates/mqdb-cluster/src/cluster/mqtt_transport.rs
Original file line number Diff line number Diff line change
Expand Up @@ -422,6 +422,10 @@ impl ClusterTransport for MqttTransport {
Err(TransportError::PartitionNotFound(partition))
}

async fn direct_peers(&self) -> Option<Vec<NodeId>> {
None
}

fn recv(&self) -> Option<InboundMessage> {
if let Ok(mut requeue) = self.requeue_buffer.try_lock()
&& let Some(msg) = requeue.pop_front()
Expand Down
31 changes: 31 additions & 0 deletions crates/mqdb-cluster/src/cluster/node_controller/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -649,6 +649,37 @@ impl<T: ClusterTransport> NodeController<T> {
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<NodeId> {
let mut members: std::collections::BTreeSet<NodeId> =
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<NodeId> {
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<NodeId> {
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
Expand Down
54 changes: 54 additions & 0 deletions crates/mqdb-cluster/src/cluster/node_controller/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,7 @@ struct MockTransport {
node_id: NodeId,
inbox: Arc<Mutex<VecDeque<InboundMessage>>>,
outbox: Arc<Mutex<Vec<(NodeId, ClusterMessage)>>>,
linked: Arc<Mutex<Option<Vec<NodeId>>>>,
}

impl MockTransport {
Expand All @@ -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)),
}
}

Expand Down Expand Up @@ -114,6 +116,10 @@ impl ClusterTransport for MockTransport {
Ok(())
}

async fn direct_peers(&self) -> Option<Vec<NodeId>> {
self.linked.lock().unwrap().clone()
}

fn recv(&self) -> Option<InboundMessage> {
self.inbox.lock().unwrap().pop_front()
}
Expand Down Expand Up @@ -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());
}
4 changes: 4 additions & 0 deletions crates/mqdb-cluster/src/cluster/quic_transport.rs
Original file line number Diff line number Diff line change
Expand Up @@ -363,6 +363,10 @@ impl ClusterTransport for QuicDirectTransport {
Err(TransportError::PartitionNotFound(partition))
}

async fn direct_peers(&self) -> Option<Vec<NodeId>> {
Some(self.peers.read().await.keys().copied().collect())
}

fn recv(&self) -> Option<InboundMessage> {
if let Ok(mut requeue) = self.requeue_buffer.try_lock()
&& let Some(msg) = requeue.pop_front()
Expand Down
4 changes: 4 additions & 0 deletions crates/mqdb-cluster/src/cluster/raft/coordinator/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,10 @@ impl ClusterTransport for MockTransport {
Ok(())
}

async fn direct_peers(&self) -> Option<Vec<NodeId>> {
None
}

fn recv(&self) -> Option<InboundMessage> {
None
}
Expand Down
4 changes: 4 additions & 0 deletions crates/mqdb-cluster/src/cluster/transport.rs
Original file line number Diff line number Diff line change
Expand Up @@ -487,6 +487,10 @@ pub trait ClusterTransport: Send + Sync + Debug + Clone {
message: ClusterMessage,
) -> impl std::future::Future<Output = Result<(), TransportError>> + 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<Output = Option<Vec<NodeId>>> + Send;

fn recv(&self) -> Option<InboundMessage>;

fn pending_count(&self) -> usize;
Expand Down
10 changes: 10 additions & 0 deletions crates/mqdb-cluster/src/cluster_agent/admin.rs
Original file line number Diff line number Diff line change
Expand Up @@ -241,6 +241,14 @@ impl ClusteredAgent {
let raft_status = self.rx_raft_status.borrow().clone();

let alive_nodes: Vec<u16> = ctrl.alive_nodes().iter().map(|n| n.get()).collect();
let members = ctrl.known_members();
let unlinked_nodes: Vec<u16> = ctrl
.unlinked_from(&members)
.await
.iter()
.map(|n| n.get())
.collect();
let known_members: Vec<u16> = members.iter().map(|n| n.get()).collect();

let mut partitions = Vec::new();
for partition in PartitionId::all() {
Expand Down Expand Up @@ -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
}))
Expand Down
24 changes: 23 additions & 1 deletion crates/mqdb-cluster/src/cluster_agent/event_loop.rs
Original file line number Diff line number Diff line change
@@ -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,
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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;
}
Expand Down Expand Up @@ -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<u16> = 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;
Expand Down
1 change: 1 addition & 0 deletions crates/mqdb-cluster/src/cluster_agent/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
Loading
Loading