diff --git a/CHANGELOG.md b/CHANGELOG.md index 29ee408b..72a75ae0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,20 @@ 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-19 — mqdb-agent 0.8.29, mqdb-cluster 0.4.14, mqdb-cli 0.8.38 + +### Added + +- **Presence feed in cluster mode**, behind the same opt-in `--presence` / `MQDB_PRESENCE` flag as agent mode. Each node publishes a retained QoS 0 message to `$DB/_presence/{client_id}` on connect and disconnect for the clients attached to it, and broadcasts that message to its peers, which publish it locally. Every node that receives it ends up holding the retained set, so a janitor sees the current state immediately on subscribe. A node only honours presence it receives if it was itself started with `--presence`. +- **`mqdb dev test --presence`** — a multi-node E2E suite covering same-node presence, cross-node presence in both directions, retained state reaching a late janitor on a different node, and a client's attempt to forge a presence message being rejected. + +### Notes + +- Cross-node delivery deliberately uses a broadcast rather than the normal publish routing. `TopicTrie::match_topic` returns no matches for any `$`-prefixed topic, so a `$DB/_presence/#` subscriber can never be resolved as a remote target; routing would silently deliver nothing across nodes while appearing to work on a single node. `_presence` is also added to the `$DB/` pass-through list in the cluster publish handler so the message is delivered to subscribers instead of being taken for a database operation. +- A session takeover reports its disconnect after the replacing connection's connect, so the cluster tracks live connections per `client_id` and reports `disconnect` only when the last one closes, exactly as agent mode does. A client that reconnects to a *different* node is covered too: a node suppresses its `disconnect` when the cluster already shows that client connected elsewhere. Without both, a reconnecting holder would be retained as disconnected and could be reclaimed while still connected. +- **Presence reaches the nodes a node is peered with.** Like the existing wildcard and topic-subscription broadcasts, a presence broadcast is a single hop over the sender's peer connections and is not relayed, so cluster-wide delivery requires a full peer mesh — every node peered with every other. This is the cluster's existing assumption (cross-node pub/sub behaves the same way), not something specific to presence; every `mqdb dev start-cluster` topology forms a full mesh. A cluster where each node is seeded from only one other node will deliver presence to a subset of nodes. +- Presence is emitted by the node a client is attached to, so if that node dies, no `disconnect` is emitted and its clients stay retained as connected. Reclaiming those holds falls to the TTL backstop. The retained publish is also best-effort: under a connect storm it can be dropped, leaving a stale retained value. + ## 2026-09-12 — mqdb-agent 0.8.28, mqdb-cli 0.8.37 ### Added diff --git a/Cargo.lock b/Cargo.lock index 7de93c44..94b59c8b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1420,7 +1420,7 @@ dependencies = [ [[package]] name = "mqdb-agent" -version = "0.8.28" +version = "0.8.29" dependencies = [ "arc-swap", "argon2", @@ -1455,7 +1455,7 @@ dependencies = [ [[package]] name = "mqdb-cli" -version = "0.8.37" +version = "0.8.38" dependencies = [ "base64", "bebytes", @@ -1482,7 +1482,7 @@ dependencies = [ [[package]] name = "mqdb-cluster" -version = "0.4.13" +version = "0.4.14" dependencies = [ "arc-swap", "bebytes", diff --git a/README.md b/README.md index e63d17be..78fd7e04 100644 --- a/README.md +++ b/README.md @@ -844,7 +844,7 @@ Every CLI flag can be set via environment variable. This is the primary configur | `MQDB_OWNERSHIP` | Ownership config (`entity=field` pairs) | — | | `MQDB_OWNERSHIP_DERIVE` | Child-entity access derivation (`child=field>parent` pairs, agent mode only) | — | | `MQDB_SCOPED_EVENTS` | Route change events to per-recipient topics for ownership-enabled entities (agent mode only) | `false` | -| `MQDB_PRESENCE` | Publish client connect/disconnect presence to `$DB/_presence/{client_id}` (retained, QoS 0; agent mode only) | `false` | +| `MQDB_PRESENCE` | Publish client connect/disconnect presence to `$DB/_presence/{client_id}` (retained, QoS 0; agent and cluster) | `false` | | `MQDB_EVENT_SCOPE` | Scope events by entity field | — | | `MQDB_WS_BIND` | WebSocket bind address | — | diff --git a/crates/mqdb-agent/Cargo.toml b/crates/mqdb-agent/Cargo.toml index c2c159eb..f59d6f14 100644 --- a/crates/mqdb-agent/Cargo.toml +++ b/crates/mqdb-agent/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "mqdb-agent" -version = "0.8.28" +version = "0.8.29" edition.workspace = true license = "Apache-2.0" authors.workspace = true diff --git a/crates/mqdb-agent/src/presence.rs b/crates/mqdb-agent/src/presence.rs index 82c348b1..4461ab9a 100644 --- a/crates/mqdb-agent/src/presence.rs +++ b/crates/mqdb-agent/src/presence.rs @@ -88,40 +88,60 @@ impl PresenceEvent { } } -pub struct PresenceEventHandler { - sender: flume::Sender, - live_connections: Mutex>, +#[derive(Default)] +pub struct LiveConnections { + counts: Mutex>, } -impl PresenceEventHandler { +impl LiveConnections { #[must_use] - pub fn new(sender: flume::Sender) -> Self { - Self { - sender, - live_connections: Mutex::new(HashMap::new()), - } + pub fn new() -> Self { + Self::default() } - fn register_connection(&self, client_id: &str) { - if let Ok(mut live) = self.live_connections.lock() { - *live.entry(client_id.to_string()).or_insert(0) += 1; + pub fn register(&self, client_id: &str) { + if let Ok(mut counts) = self.counts.lock() { + *counts.entry(client_id.to_string()).or_insert(0) += 1; } } - fn is_last_connection(&self, client_id: &str) -> bool { - let Ok(mut live) = self.live_connections.lock() else { + pub fn is_last(&self, client_id: &str) -> bool { + let Ok(mut counts) = self.counts.lock() else { return true; }; - let Some(remaining) = live.get_mut(client_id) else { + let Some(remaining) = counts.get_mut(client_id) else { return true; }; *remaining = remaining.saturating_sub(1); if *remaining == 0 { - live.remove(client_id); + counts.remove(client_id); return true; } false } +} + +pub struct PresenceEventHandler { + sender: flume::Sender, + live_connections: LiveConnections, +} + +impl PresenceEventHandler { + #[must_use] + pub fn new(sender: flume::Sender) -> Self { + Self { + sender, + live_connections: LiveConnections::new(), + } + } + + fn register_connection(&self, client_id: &str) { + self.live_connections.register(client_id); + } + + fn is_last_connection(&self, client_id: &str) -> bool { + self.live_connections.is_last(client_id) + } fn emit(&self, presence: PresenceEvent) { if self.sender.try_send(presence).is_err() { diff --git a/crates/mqdb-cli/Cargo.toml b/crates/mqdb-cli/Cargo.toml index b9a8b0a9..42d8a50f 100644 --- a/crates/mqdb-cli/Cargo.toml +++ b/crates/mqdb-cli/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "mqdb-cli" -version = "0.8.37" +version = "0.8.38" publish = false edition.workspace = true license = "AGPL-3.0-only" diff --git a/crates/mqdb-cli/src/cli_types/agent.rs b/crates/mqdb-cli/src/cli_types/agent.rs index 24ad3129..07dca0d1 100644 --- a/crates/mqdb-cli/src/cli_types/agent.rs +++ b/crates/mqdb-cli/src/cli_types/agent.rs @@ -257,6 +257,12 @@ pub(crate) struct ClusterStartFields { help = "Scope events by entity field (e.g. diagrams=diagramId)" )] pub(crate) event_scope: Option, + #[arg( + long, + env = "MQDB_PRESENCE", + help = "Publish client connect/disconnect presence to $DB/_presence/{client_id} (retained, QoS 0) for subscribers such as an abandoned-hold janitor" + )] + pub(crate) presence: bool, #[arg( long, env = "MQDB_PASSPHRASE_FILE", diff --git a/crates/mqdb-cli/src/cli_types/dev.rs b/crates/mqdb-cli/src/cli_types/dev.rs index f5349bba..7b807ee2 100644 --- a/crates/mqdb-cli/src/cli_types/dev.rs +++ b/crates/mqdb-cli/src/cli_types/dev.rs @@ -48,6 +48,8 @@ pub(crate) enum DevAction { ownership: bool, #[arg(long, help = "Run diagram-sharing tests (needs a license)")] sharing: bool, + #[arg(long, help = "Run presence feed tests (needs a license)")] + presence: bool, #[arg(long, help = "Run constraint stress tests")] stress_constraints: bool, #[arg(long, help = "Run all test suites")] diff --git a/crates/mqdb-cli/src/commands/cluster.rs b/crates/mqdb-cli/src/commands/cluster.rs index 76106305..467565c7 100644 --- a/crates/mqdb-cli/src/commands/cluster.rs +++ b/crates/mqdb-cli/src/commands/cluster.rs @@ -43,6 +43,7 @@ pub(crate) struct ClusterStartArgs { pub(crate) oauth: OAuthArgs, pub(crate) ownership: Option, pub(crate) event_scope: Option, + pub(crate) presence: bool, pub(crate) passphrase_file: Option, pub(crate) passphrase_data: Option, pub(crate) license: Option, @@ -206,6 +207,9 @@ pub(crate) async fn cmd_cluster_start( .map_err(|e| format!("invalid --event-scope: {e}"))?; config = config.with_scope_config(scope_config); } + if args.presence { + config = config.with_presence(true); + } if let Some(ref info) = license_info { config = config.with_license_expiry(info.expires_at); } diff --git a/crates/mqdb-cli/src/commands/dev/tests.rs b/crates/mqdb-cli/src/commands/dev/tests.rs index 78c85f92..743bf804 100644 --- a/crates/mqdb-cli/src/commands/dev/tests.rs +++ b/crates/mqdb-cli/src/commands/dev/tests.rs @@ -64,6 +64,7 @@ pub(crate) fn cmd_dev_test( lwt: bool, ownership: bool, sharing: bool, + presence: bool, stress_constraints: bool, all: bool, nodes: u8, @@ -78,6 +79,7 @@ pub(crate) fn cmd_dev_test( && !lwt && !ownership && !sharing + && !presence && !stress_constraints); wait_for_cluster_ready(nodes, 10); @@ -115,6 +117,10 @@ pub(crate) fn cmd_dev_test( run_test_sharing(nodes, &ports, license); } + if presence { + run_test_presence(nodes, &ports, license); + } + if stress_constraints { run_test_stress_constraints(nodes, &ports); } @@ -1437,3 +1443,224 @@ fn run_test_stress_constraints(nodes: u8, ports: &[u16]) { println!("\nResults: {passed} passed, {failed} failed\n"); } + +fn start_presence_cluster( + nodes: u8, + exe: &Path, + passwd_path: &str, + license: Option<&Path>, +) -> bool { + let db_prefix = "/tmp/mqdb-test-presence"; + + let _ = std::fs::remove_file(passwd_path); + for (user, pass) in [("alice", "alice"), ("admin", "admin")] { + let _ = Command::new(exe) + .args(["passwd", user, "-b", pass, "-f", passwd_path]) + .output(); + } + + let _ = Command::new("pkill").args(["-f", "mqdb cluster"]).status(); + std::thread::sleep(Duration::from_secs(1)); + for i in 1..=nodes { + let _ = std::fs::remove_dir_all(format!("{db_prefix}-{i}")); + } + + let quic_cert = PathBuf::from("test_certs/server.pem"); + let quic_key = PathBuf::from("test_certs/server.key"); + let quic_ca = PathBuf::from("test_certs/ca.pem"); + + for node_id in 1..=nodes { + let port = 1882 + u16::from(node_id); + let db_path = format!("{db_prefix}-{node_id}"); + let _ = std::fs::create_dir_all(&db_path); + + let mut cmd = Command::new(exe); + cmd.args([ + "cluster", + "start", + "--node-id", + &node_id.to_string(), + "--bind", + &format!("127.0.0.1:{port}"), + "--db", + &db_path, + "--admin-users", + "admin", + "--passwd", + passwd_path, + "--presence", + ]); + + if let Some(lic_str) = license.and_then(|p| p.to_str()) { + cmd.args(["--license", lic_str]); + } + + let peers: Vec = (1..node_id) + .map(|n| format!("{}@127.0.0.1:{}", n, 1882 + u16::from(n))) + .collect(); + if !peers.is_empty() { + cmd.args(["--peers", &peers.join(",")]); + } + + if quic_cert.exists() && quic_key.exists() { + cmd.args([ + "--quic-cert", + quic_cert.to_str().unwrap_or(""), + "--quic-key", + quic_key.to_str().unwrap_or(""), + ]); + if quic_ca.exists() { + cmd.args(["--quic-ca", quic_ca.to_str().unwrap_or("")]); + } + #[cfg(feature = "dev-insecure")] + cmd.arg("--quic-insecure"); + } + + cmd.env( + "RUST_LOG", + std::env::var("RUST_LOG").unwrap_or_else(|_| "info".to_string()), + ); + if let Ok(log_file) = std::fs::File::create(format!("{db_path}/mqdb.log")) { + if let Ok(log_clone) = log_file.try_clone() { + cmd.stdout(log_clone); + } + cmd.stderr(log_file); + } + + println!(" Starting node {node_id} on port {port} (auth + presence)..."); + let _ = cmd.spawn(); + std::thread::sleep(Duration::from_millis(500)); + } + + println!(" Waiting for authenticated cluster..."); + wait_for_auth_cluster(nodes, 15) +} + +fn presence_line_has(out: &str, client_id: &str, event: &str) -> bool { + out.lines() + .any(|line| line.contains(client_id) && line.contains(event)) +} + +fn capture_presence(sub_port: u16, holder_id: &str, holder_port: u16) -> String { + let _ = std::fs::remove_file("/tmp/mqdb-presence-out.txt"); + let script = format!( + "mosquitto_sub -i 'janitor-{sub_port}' -h 127.0.0.1 -p {sub_port} -u alice -P alice -t '$DB/_presence/#' -v > /tmp/mqdb-presence-out.txt & \ + sleep 2; \ + mosquitto_pub -i '{holder_id}' -h 127.0.0.1 -p {holder_port} -u alice -P alice -t 'presence/noop' -m 'x'; \ + sleep 2; wait" + ); + let _ = Command::new("timeout") + .args(["8", "sh", "-c", &script]) + .output(); + std::fs::read_to_string("/tmp/mqdb-presence-out.txt").unwrap_or_default() +} + +fn run_test_presence(nodes: u8, _ports: &[u16], license: Option<&Path>) { + println!("=== Cluster Presence Feed E2E ({nodes} nodes) ===\n"); + + let exe = std::env::current_exe().unwrap_or_else(|_| PathBuf::from("mqdb")); + let passwd_path = "/tmp/mqdb-test-presence-passwd"; + let mut passed = 0u32; + let mut failed = 0u32; + + if !start_presence_cluster(nodes, &exe, passwd_path, license) { + println!(" Cluster failed to become ready"); + let _ = Command::new("pkill").args(["-f", "mqdb cluster"]).status(); + println!("\nResults: 0 passed, 1 failed\n"); + return; + } + println!(" Authenticated cluster ready!\n"); + + let ts = std::time::UNIX_EPOCH.elapsed().map_or(0, |d| d.as_millis()); + + let mut check = |name: &str, ok: bool| { + if ok { + println!(" {name}: ✓"); + passed += 1; + } else { + println!(" {name}: ✗"); + failed += 1; + } + }; + + let same_id = format!("seat-holder-same-{ts}"); + let out = capture_presence(1883, &same_id, 1883); + check( + "same-node connect presence", + presence_line_has(&out, &same_id, "\"event\":\"connect\""), + ); + check( + "same-node disconnect presence", + presence_line_has(&out, &same_id, "\"event\":\"disconnect\""), + ); + + if nodes >= 2 { + let cross_id = format!("seat-holder-cross-{ts}"); + let out = capture_presence(1883, &cross_id, 1884); + check( + "cross-node connect presence (janitor n1, holder n2)", + presence_line_has(&out, &cross_id, "\"event\":\"connect\""), + ); + check( + "cross-node disconnect presence (janitor n1, holder n2)", + presence_line_has(&out, &cross_id, "\"event\":\"disconnect\""), + ); + } + + if nodes >= 3 { + let far_id = format!("seat-holder-far-{ts}"); + let out = capture_presence(1885, &far_id, 1884); + check( + "cross-node presence (janitor n3, holder n2)", + out.contains(&far_id), + ); + } + + let retained_id = format!("seat-holder-retained-{ts}"); + let _ = Command::new("timeout") + .args([ + "5", + "sh", + "-c", + &format!( + "mosquitto_pub -i '{retained_id}' -h 127.0.0.1 -p 1884 -u alice -P alice -t 'presence/noop' -m 'x'" + ), + ]) + .output(); + std::thread::sleep(Duration::from_secs(2)); + let late = Command::new("timeout") + .args([ + "5", + "sh", + "-c", + "mosquitto_sub -i 'janitor-late' -h 127.0.0.1 -p 1883 -u alice -P alice -t '$DB/_presence/#' -v -W 3", + ]) + .output() + .map_or_else(|_| String::new(), |o| String::from_utf8_lossy(&o.stdout).to_string()); + check( + "retained presence reaches a late janitor on another node", + late.contains(&retained_id), + ); + + let _ = Command::new("timeout") + .args([ + "5", + "sh", + "-c", + "mosquitto_pub -i 'attacker' -h 127.0.0.1 -p 1883 -u alice -P alice -t '$DB/_presence/victim' -m '{\"forged\":true}'", + ]) + .output(); + let victim = Command::new("timeout") + .args([ + "5", + "sh", + "-c", + "mosquitto_sub -i 'janitor-forge' -h 127.0.0.1 -p 1883 -u alice -P alice -t '$DB/_presence/victim' -v -W 3", + ]) + .output() + .map_or_else(|_| String::new(), |o| String::from_utf8_lossy(&o.stdout).to_string()); + check("a client cannot forge presence", !victim.contains("forged")); + + let _ = Command::new("pkill").args(["-f", "mqdb cluster"]).status(); + println!("\nResults: {passed} passed, {failed} failed\n"); +} diff --git a/crates/mqdb-cli/src/main.rs b/crates/mqdb-cli/src/main.rs index 075e0553..6b14d76e 100644 --- a/crates/mqdb-cli/src/main.rs +++ b/crates/mqdb-cli/src/main.rs @@ -285,6 +285,7 @@ async fn dispatch_cluster(action: ClusterAction) -> Result<(), Box Result<(), Box Result<(), Box)>>>; +type SentMessages = Arc>>; +type PresenceHarness = ( + crate::cluster::event_handler::ClusterEventHandler, + RetainedPublishes, + SentMessages, +); + #[derive(Debug, Clone)] struct MockTransport { node_id: NodeId, inbox: Arc>>, outbox: Arc>>, + retained_publishes: RetainedPublishes, } impl MockTransport { @@ -45,6 +54,7 @@ impl MockTransport { node_id, inbox: Arc::new(Mutex::new(VecDeque::new())), outbox: Arc::new(Mutex::new(Vec::new())), + retained_publishes: Arc::new(Mutex::new(Vec::new())), } } } @@ -100,7 +110,12 @@ impl ClusterTransport for MockTransport { async fn queue_local_publish(&self, _topic: String, _payload: Vec, _qos: u8) {} - async fn queue_local_publish_retained(&self, _topic: String, _payload: Vec, _qos: u8) {} + async fn queue_local_publish_retained(&self, topic: String, payload: Vec, _qos: u8) { + self.retained_publishes + .lock() + .unwrap() + .push((topic, payload)); + } } fn setup_controller_with_partition(partition: PartitionId) -> NodeController { @@ -2379,6 +2394,7 @@ async fn cluster_share_and_read_forward_to_resource_primary() { node_id: node1, inbox: Arc::new(Mutex::new(VecDeque::new())), outbox: Arc::clone(&outbox), + retained_publishes: Arc::new(Mutex::new(Vec::new())), }; let mut ctrl = create_test_controller(node1, transport); ctrl.set_ownership(Arc::clone(&ownership)); @@ -2467,6 +2483,7 @@ async fn cluster_forwarded_read_graded_against_grant_on_primary() { node_id: node1, inbox: Arc::new(Mutex::new(VecDeque::new())), outbox: Arc::clone(&outbox), + retained_publishes: Arc::new(Mutex::new(Vec::new())), }; let mut ctrl = create_test_controller(node1, transport); ctrl.set_ownership(Arc::clone(&ownership)); @@ -2835,3 +2852,147 @@ async fn cluster_identity_inflight_resolution_completes_when_not_unshared() { "ok" ); } + +fn presence_handler_with_transport() -> PresenceHarness { + let node1 = NodeId::validated(1).unwrap(); + let transport = MockTransport::new(node1); + let retained = Arc::clone(&transport.retained_publishes); + let outbox = Arc::clone(&transport.outbox); + let ctrl = create_test_controller(node1, transport); + let handler = crate::cluster::event_handler::ClusterEventHandler::new( + node1, + Arc::new(tokio::sync::RwLock::new(ctrl)), + ) + .with_presence(true); + (handler, retained, outbox) +} + +fn cluster_connect_event(client_id: &str) -> mqtt5::broker::events::ClientConnectEvent { + mqtt5::broker::events::ClientConnectEvent { + client_id: client_id.into(), + user_id: Some("alice".into()), + clean_start: true, + session_expiry_interval: 0, + will_topic: None, + will_payload: None, + will_qos: None, + will_retain: None, + } +} + +#[tokio::test] +async fn cluster_presence_publishes_locally_and_broadcasts() { + use mqtt5::broker::events::BrokerEventHandler; + + let (handler, retained, outbox) = presence_handler_with_transport(); + handler + .on_client_connect(cluster_connect_event("seat-holder-7")) + .await; + + let published = retained.lock().unwrap().clone(); + assert_eq!(published.len(), 1, "presence must be delivered locally"); + assert_eq!(published[0].0, "$DB/_presence/seat-holder-7"); + let value: serde_json::Value = serde_json::from_slice(&published[0].1).unwrap(); + assert_eq!(value["event"], "connect"); + assert_eq!(value["user_id"], "alice"); + + let broadcasts = outbox.lock().unwrap(); + let presence_broadcasts: Vec<_> = broadcasts + .iter() + .filter_map(|(_, msg)| match msg { + ClusterMessage::PresenceBroadcast(b) => Some(b), + _ => None, + }) + .collect(); + assert_eq!( + presence_broadcasts.len(), + 1, + "cross-node delivery relies on a broadcast, not wildcard routing" + ); + assert_eq!( + presence_broadcasts[0].topic_str(), + "$DB/_presence/seat-holder-7" + ); +} + +#[tokio::test] +async fn cluster_presence_skips_internal_clients_and_takeover() { + use mqtt5::broker::events::BrokerEventHandler; + + let (handler, retained, _outbox) = presence_handler_with_transport(); + + handler + .on_client_connect(cluster_connect_event("mqdb-admin-1")) + .await; + assert!( + retained.lock().unwrap().is_empty(), + "internal clients must not appear in the presence feed" + ); + + handler + .on_client_connect(cluster_connect_event("seat-holder-7")) + .await; + handler + .on_client_connect(cluster_connect_event("seat-holder-7")) + .await; + retained.lock().unwrap().clear(); + + handler + .on_client_disconnect(mqtt5::broker::events::ClientDisconnectEvent { + client_id: "seat-holder-7".into(), + user_id: Some("alice".into()), + reason: mqtt5::types::ReasonCode::Success, + unexpected: true, + }) + .await; + assert!( + retained.lock().unwrap().is_empty(), + "a displaced connection must not report a live client as disconnected" + ); +} + +#[tokio::test] +async fn received_presence_broadcast_is_published_locally() { + let node1 = NodeId::validated(1).unwrap(); + let node2 = NodeId::validated(2).unwrap(); + let transport = MockTransport::new(node1); + let retained = Arc::clone(&transport.retained_publishes); + let mut ctrl = create_test_controller(node1, transport); + ctrl.set_presence(true); + + let broadcast = crate::cluster::protocol::PresenceBroadcast::try_new( + "$DB/_presence/remote-holder", + br#"{"client_id":"remote-holder","event":"disconnect"}"#, + ) + .expect("presence fits in a broadcast"); + ctrl.handle_presence_broadcast(node2, &broadcast).await; + + let published = retained.lock().unwrap().clone(); + assert_eq!( + published.len(), + 1, + "a node must locally publish presence it receives" + ); + assert_eq!(published[0].0, "$DB/_presence/remote-holder"); +} + +#[tokio::test] +async fn received_presence_broadcast_is_ignored_when_presence_is_off() { + let node1 = NodeId::validated(1).unwrap(); + let node2 = NodeId::validated(2).unwrap(); + let transport = MockTransport::new(node1); + let retained = Arc::clone(&transport.retained_publishes); + let ctrl = create_test_controller(node1, transport); + + let broadcast = crate::cluster::protocol::PresenceBroadcast::try_new( + "$DB/_presence/remote-holder", + br#"{"client_id":"remote-holder","event":"connect"}"#, + ) + .expect("presence fits in a broadcast"); + ctrl.handle_presence_broadcast(node2, &broadcast).await; + + assert!( + retained.lock().unwrap().is_empty(), + "a node that did not opt into presence must not retain a peer's presence" + ); +} diff --git a/crates/mqdb-cluster/src/cluster/event_handler/broker_events.rs b/crates/mqdb-cluster/src/cluster/event_handler/broker_events.rs index bd277f26..cb9a0729 100644 --- a/crates/mqdb-cluster/src/cluster/event_handler/broker_events.rs +++ b/crates/mqdb-cluster/src/cluster/event_handler/broker_events.rs @@ -37,6 +37,15 @@ impl BrokerEventHandler for ClusterEventHandler BrokerEventHandler for ClusterEventHandler BrokerEventHandler for ClusterEventHandler BrokerEventHandler for ClusterEventHandler { synced_retained_topics: Arc>>, db_handler: DbRequestHandler, vault_key_store: Arc, + presence: bool, + live_connections: mqdb_agent::presence::LiveConnections, } impl ClusterEventHandler { @@ -39,9 +41,17 @@ impl ClusterEventHandler { synced_retained_topics: Arc::new(RwLock::new(HashMap::new())), db_handler: DbRequestHandler::new(node_id), vault_key_store: Arc::new(VaultKeyStore::new()), + presence: false, + live_connections: mqdb_agent::presence::LiveConnections::new(), } } + #[must_use] + pub fn with_presence(mut self, enabled: bool) -> Self { + self.presence = enabled; + self + } + #[must_use] pub fn with_ownership( mut self, diff --git a/crates/mqdb-cluster/src/cluster/event_handler/routing.rs b/crates/mqdb-cluster/src/cluster/event_handler/routing.rs index 7c50d7ca..4c3dbfd7 100644 --- a/crates/mqdb-cluster/src/cluster/event_handler/routing.rs +++ b/crates/mqdb-cluster/src/cluster/event_handler/routing.rs @@ -2,6 +2,7 @@ // SPDX-License-Identifier: AGPL-3.0-only use super::super::node_controller::NodeController; +use super::super::protocol::PresenceBroadcast; use super::super::transport::{ClusterMessage, ClusterTransport}; use super::super::{ ForwardTarget, ForwardedPublish, LwtPublisher, NodeId, PublishRouter, SubscriptionType, @@ -13,6 +14,42 @@ use std::collections::HashMap; use tracing::{debug, trace, warn}; impl ClusterEventHandler { + pub(super) async fn client_is_live_elsewhere(&self, client_id: &str) -> bool { + let ctrl = self.controller.read().await; + Self::resolve_connected_node(&ctrl, client_id).is_some_and(|node| node != self.node_id) + } + + pub(super) async fn emit_presence(&self, presence: &mqdb_agent::presence::PresenceEvent) { + if !mqdb_agent::presence::is_publishable_client_id(&presence.client_id) { + return; + } + let payload = match presence.payload() { + Ok(payload) => payload, + Err(e) => { + warn!(error = %e, "failed to serialize presence event"); + return; + } + }; + let topic = presence.topic(); + + let transport = { + let ctrl = self.controller.read().await; + ctrl.transport().clone() + }; + + transport + .queue_local_publish_retained(topic.clone(), payload.clone(), 0) + .await; + + let Some(broadcast) = PresenceBroadcast::try_new(&topic, &payload) else { + warn!(topic, "presence message too large to broadcast"); + return; + }; + let _ = transport + .broadcast(ClusterMessage::PresenceBroadcast(broadcast)) + .await; + } + pub(super) async fn broadcast_topic_subscription( ctrl: &NodeController, broadcast: TopicSubscriptionBroadcast, diff --git a/crates/mqdb-cluster/src/cluster/node_controller/broadcast.rs b/crates/mqdb-cluster/src/cluster/node_controller/broadcast.rs index 3f166b8c..6b25abe0 100644 --- a/crates/mqdb-cluster/src/cluster/node_controller/broadcast.rs +++ b/crates/mqdb-cluster/src/cluster/node_controller/broadcast.rs @@ -11,6 +11,24 @@ use super::super::{Epoch, NodeId, PartitionId, SubscriptionType, WildcardStoreEr use super::NodeController; impl NodeController { + pub(crate) async fn handle_presence_broadcast( + &self, + from: NodeId, + broadcast: &crate::cluster::protocol::PresenceBroadcast, + ) { + if !self.presence { + return; + } + let topic = broadcast.topic_str(); + if topic.is_empty() { + return; + } + tracing::debug!(topic, from = from.get(), "received presence broadcast"); + self.transport() + .queue_local_publish_retained(topic.to_string(), broadcast.payload().to_vec(), 0) + .await; + } + pub(crate) fn handle_wildcard_broadcast( &mut self, from: NodeId, diff --git a/crates/mqdb-cluster/src/cluster/node_controller/mod.rs b/crates/mqdb-cluster/src/cluster/node_controller/mod.rs index 1ffa6eef..52261703 100644 --- a/crates/mqdb-cluster/src/cluster/node_controller/mod.rs +++ b/crates/mqdb-cluster/src/cluster/node_controller/mod.rs @@ -454,6 +454,7 @@ pub struct NodeController { pub(super) pending_vault_decrypts: HashMap, pub(super) pending_constraints: Arc, pub(super) ownership: Arc, + pub(super) presence: bool, pub(super) vault_key_store: Arc, #[cfg(feature = "http-api")] pub(super) identity_crypto: Option>, @@ -529,6 +530,7 @@ impl NodeController { pending_vault_decrypts: HashMap::new(), pending_constraints: Arc::new(pending::PendingConstraintState::new()), ownership: Arc::new(OwnershipConfig::default()), + presence: false, vault_key_store: Arc::new(VaultKeyStore::new()), #[cfg(feature = "http-api")] identity_crypto: None, @@ -539,6 +541,10 @@ impl NodeController { self.ownership = ownership; } + pub fn set_presence(&mut self, enabled: bool) { + self.presence = enabled; + } + pub fn set_vault_key_store(&mut self, store: Arc) { self.vault_key_store = store; } @@ -1293,6 +1299,9 @@ impl NodeController { ClusterMessage::BatchReadRequest(request) => { self.handle_batch_read_and_respond(from, request).await; } + ClusterMessage::PresenceBroadcast(broadcast) => { + self.handle_presence_broadcast(from, broadcast).await; + } ClusterMessage::WildcardBroadcast(broadcast) => { self.handle_wildcard_broadcast(from, broadcast); } diff --git a/crates/mqdb-cluster/src/cluster/protocol/broadcast.rs b/crates/mqdb-cluster/src/cluster/protocol/broadcast.rs index d53ebbad..9b776de4 100644 --- a/crates/mqdb-cluster/src/cluster/protocol/broadcast.rs +++ b/crates/mqdb-cluster/src/cluster/protocol/broadcast.rs @@ -201,3 +201,41 @@ impl TopicSubscriptionBroadcast { self.qos } } + +#[derive(Debug, Clone, BeBytes)] +pub struct PresenceBroadcast { + version: u8, + topic_len: u16, + #[FromField(topic_len)] + topic: Vec, + payload_len: u16, + #[FromField(payload_len)] + payload: Vec, +} + +impl PresenceBroadcast { + pub const VERSION: u8 = 1; + + #[must_use] + pub fn try_new(topic: &str, payload: &[u8]) -> Option { + let topic_len = u16::try_from(topic.len()).ok()?; + let payload_len = u16::try_from(payload.len()).ok()?; + Some(Self { + version: Self::VERSION, + topic_len, + topic: topic.as_bytes().to_vec(), + payload_len, + payload: payload.to_vec(), + }) + } + + #[must_use] + pub fn topic_str(&self) -> &str { + std::str::from_utf8(&self.topic).unwrap_or("") + } + + #[must_use] + pub fn payload(&self) -> &[u8] { + &self.payload + } +} diff --git a/crates/mqdb-cluster/src/cluster/protocol/mod.rs b/crates/mqdb-cluster/src/cluster/protocol/mod.rs index 8cd7cc7f..1121966d 100644 --- a/crates/mqdb-cluster/src/cluster/protocol/mod.rs +++ b/crates/mqdb-cluster/src/cluster/protocol/mod.rs @@ -14,7 +14,7 @@ mod unique; #[cfg(test)] mod tests; -pub use broadcast::{TopicSubscriptionBroadcast, WildcardBroadcast, WildcardOp}; +pub use broadcast::{PresenceBroadcast, TopicSubscriptionBroadcast, WildcardBroadcast, WildcardOp}; pub use db_messages::{JsonDbRequest, JsonDbResponse}; pub use fk::{FkCheckRequest, FkCheckResponse, FkReverseLookupRequest, FkReverseLookupResponse}; pub use heartbeat::Heartbeat; diff --git a/crates/mqdb-cluster/src/cluster/transport.rs b/crates/mqdb-cluster/src/cluster/transport.rs index 837c35a4..15d56cbc 100644 --- a/crates/mqdb-cluster/src/cluster/transport.rs +++ b/crates/mqdb-cluster/src/cluster/transport.rs @@ -4,10 +4,11 @@ use super::protocol::{ BatchReadRequest, BatchReadResponse, CatchupRequest, CatchupResponse, FkCheckRequest, FkCheckResponse, FkReverseLookupRequest, FkReverseLookupResponse, ForwardedPublish, Heartbeat, - JsonDbRequest, JsonDbResponse, QueryRequest, QueryResponse, ReplicationAck, ReplicationWrite, - TopicSubscriptionBroadcast, UniqueCommitRequest, UniqueCommitResponse, UniqueReassertRequest, - UniqueReleaseRequest, UniqueReleaseResponse, UniqueReplicateAck, UniqueReserveRequest, - UniqueReserveResponse, UniqueSealRequest, UniqueSealResponse, WildcardBroadcast, + JsonDbRequest, JsonDbResponse, PresenceBroadcast, QueryRequest, QueryResponse, ReplicationAck, + ReplicationWrite, TopicSubscriptionBroadcast, UniqueCommitRequest, UniqueCommitResponse, + UniqueReassertRequest, UniqueReleaseRequest, UniqueReleaseResponse, UniqueReplicateAck, + UniqueReserveRequest, UniqueReserveResponse, UniqueSealRequest, UniqueSealResponse, + WildcardBroadcast, }; use super::raft::{ AppendEntriesRequest, AppendEntriesResponse, PartitionUpdate, RequestVoteRequest, @@ -47,6 +48,7 @@ pub enum ClusterMessage { QueryResponse(QueryResponse), BatchReadRequest(BatchReadRequest), BatchReadResponse(BatchReadResponse), + PresenceBroadcast(PresenceBroadcast), WildcardBroadcast(WildcardBroadcast), TopicSubscriptionBroadcast(TopicSubscriptionBroadcast), PartitionUpdate(PartitionUpdate), @@ -99,6 +101,7 @@ impl ClusterMessage { Self::QueryResponse(_) => 51, Self::BatchReadRequest(_) => 52, Self::BatchReadResponse(_) => 53, + Self::PresenceBroadcast(_) => 62, Self::WildcardBroadcast(_) => 60, Self::TopicSubscriptionBroadcast(_) => 61, Self::PartitionUpdate(_) => 70, @@ -145,6 +148,7 @@ impl ClusterMessage { Self::QueryResponse(_) => "QueryResponse", Self::BatchReadRequest(_) => "BatchReadRequest", Self::BatchReadResponse(_) => "BatchReadResponse", + Self::PresenceBroadcast(_) => "PresenceBroadcast", Self::WildcardBroadcast(_) => "WildcardBroadcast", Self::TopicSubscriptionBroadcast(_) => "TopicSubscriptionBroadcast", Self::PartitionUpdate(_) => "PartitionUpdate", @@ -194,6 +198,7 @@ impl ClusterMessage { Self::QueryResponse(resp) => buf.extend_from_slice(&resp.to_bytes()), Self::BatchReadRequest(req) => buf.extend_from_slice(&req.to_bytes()), Self::BatchReadResponse(resp) => buf.extend_from_slice(&resp.to_bytes()), + Self::PresenceBroadcast(b) => buf.extend_from_slice(&b.to_be_bytes()), Self::WildcardBroadcast(b) => buf.extend_from_slice(&b.to_be_bytes()), Self::TopicSubscriptionBroadcast(b) => buf.extend_from_slice(&b.to_be_bytes()), Self::PartitionUpdate(u) => buf.extend_from_slice(&u.to_be_bytes()), @@ -295,6 +300,10 @@ impl ClusterMessage { let (broadcast, _) = TopicSubscriptionBroadcast::try_from_be_bytes(data).ok()?; Some(Self::TopicSubscriptionBroadcast(broadcast)) } + 62 => { + let (broadcast, _) = PresenceBroadcast::try_from_be_bytes(data).ok()?; + Some(Self::PresenceBroadcast(broadcast)) + } 70 => { let (update, _) = PartitionUpdate::try_from_be_bytes(data).ok()?; Some(Self::PartitionUpdate(update)) diff --git a/crates/mqdb-cluster/src/cluster_agent/broker.rs b/crates/mqdb-cluster/src/cluster_agent/broker.rs index cb423952..0ac7a0ac 100644 --- a/crates/mqdb-cluster/src/cluster_agent/broker.rs +++ b/crates/mqdb-cluster/src/cluster_agent/broker.rs @@ -132,7 +132,8 @@ impl ClusteredAgent { ClusterEventHandler::new(self.node_id, self.controller.clone()) .with_ownership(ownership) .with_scope_config(Arc::clone(&self.scope_config)) - .with_vault_key_store(Arc::clone(&self.vault_key_store)), + .with_vault_key_store(Arc::clone(&self.vault_key_store)) + .with_presence(self.presence), ); self.configure_broker_with_auth( event_handler, @@ -544,7 +545,8 @@ impl ClusteredAgent { ClusterEventHandler::new(self.node_id, self.controller.clone()) .with_ownership(Arc::clone(&self.ownership)) .with_scope_config(Arc::clone(&self.scope_config)) - .with_vault_key_store(Arc::clone(&self.vault_key_store)), + .with_vault_key_store(Arc::clone(&self.vault_key_store)) + .with_presence(self.presence), ); let synced_retained_topics = event_handler.synced_retained_topics(); diff --git a/crates/mqdb-cluster/src/cluster_agent/config.rs b/crates/mqdb-cluster/src/cluster_agent/config.rs index 8238b344..92cda1f9 100644 --- a/crates/mqdb-cluster/src/cluster_agent/config.rs +++ b/crates/mqdb-cluster/src/cluster_agent/config.rs @@ -37,6 +37,7 @@ impl ClusterConfig { http_config: None, ownership: mqdb_core::types::OwnershipConfig::default(), scope_config: mqdb_core::types::ScopeConfig::default(), + presence: false, passphrase: None, license_expires_at: None, } @@ -200,6 +201,12 @@ impl ClusterConfig { self } + #[must_use] + pub fn with_presence(mut self, enabled: bool) -> Self { + self.presence = enabled; + self + } + #[must_use] pub fn with_passphrase(mut self, passphrase: String) -> Self { self.passphrase = Some(passphrase); diff --git a/crates/mqdb-cluster/src/cluster_agent/init.rs b/crates/mqdb-cluster/src/cluster_agent/init.rs index 1d6d24cc..d5e494a7 100644 --- a/crates/mqdb-cluster/src/cluster_agent/init.rs +++ b/crates/mqdb-cluster/src/cluster_agent/init.rs @@ -155,6 +155,7 @@ impl ClusteredAgent { tx_raft_events.clone(), ); controller.set_ownership(Arc::clone(&ownership_arc)); + controller.set_presence(config.presence); controller.set_vault_key_store(Arc::clone(&vault_key_store)); #[cfg(feature = "http-api")] controller.set_identity_crypto( @@ -243,6 +244,7 @@ impl ClusteredAgent { http_config: config.http_config, ownership: ownership_arc, scope_config: Arc::new(config.scope_config), + presence: config.presence, auth_providers: None, vault_key_store, license_expires_at: config.license_expires_at, diff --git a/crates/mqdb-cluster/src/cluster_agent/mod.rs b/crates/mqdb-cluster/src/cluster_agent/mod.rs index 46e9ac36..b05f4abf 100644 --- a/crates/mqdb-cluster/src/cluster_agent/mod.rs +++ b/crates/mqdb-cluster/src/cluster_agent/mod.rs @@ -123,6 +123,7 @@ pub struct ClusterConfig { pub http_config: Option, pub ownership: mqdb_core::types::OwnershipConfig, pub scope_config: mqdb_core::types::ScopeConfig, + pub presence: bool, pub passphrase: Option, pub license_expires_at: Option, } @@ -164,6 +165,7 @@ pub struct ClusteredAgent { http_config: Option, ownership: Arc, scope_config: Arc, + presence: bool, auth_providers: Option>, vault_key_store: Arc, license_expires_at: Option, diff --git a/docs/design/hold-reclaim.md b/docs/design/hold-reclaim.md index f9f0b79c..db61dddb 100644 --- a/docs/design/hold-reclaim.md +++ b/docs/design/hold-reclaim.md @@ -155,7 +155,7 @@ forging a presence message, leaving only the internal-service username bypass; a short-circuits `is_internal_entity_topic`, which would otherwise deny a non-admin janitor's *subscribe* outright because `_presence` is `_`-prefixed. -Cluster (follow-up): **do not** copy the LWT template. Two verified constraints shape it: +Cluster (done): **do not** copy the LWT template. Two verified constraints shape it: - `on_client_publish` swallows any `$DB/` topic outside `{_health,_admin,_sub,_resp}` (`PublishAction::Handled`), and a broker-injected `queue_local_publish` *does* re-enter that handler (on QUIC it publishes as @@ -170,6 +170,27 @@ Cluster (follow-up): **do not** copy the LWT template. Two verified constraints - The janitor must **not** be named `mqdb-*`: `on_client_subscribe` early-returns on that prefix, so its subscription would never be registered cluster-wide. +Cluster limitations, all verified against a live 3-node cluster: +- **A presence broadcast is one hop.** `transport.broadcast` walks only the sender's + connected peers and is never relayed, so cluster-wide presence requires a **full peer + mesh**. This is pre-existing and not presence-specific — wildcard and topic-subscription + broadcasts use the same call, and plain cross-node pub/sub fails on the same topology. + Verified: on a hub-and-spoke cluster (nodes 2 and 3 each peered only to node 1) a janitor + on node 3 receives nothing for a client on node 2, and neither does a plain subscriber. + Every `mqdb dev start-cluster` topology happens to form a full mesh, which is why the E2E + does not catch it. +- **A node only honours presence it opted into.** `handle_presence_broadcast` returns early + unless the receiving node was started with `--presence`. +- **Presence retained state is never replicated.** `on_retained_set` skips the presence + prefix alongside `$SYS/` and `_mqdb/`, because the broadcast already puts the message on + every node; replicating it too would make N nodes race writes to the same key and let + arrival order pick the winner. +- **A dead node's clients stay retained as connected**, since presence is only emitted by + the node a client is attached to. Those holds are reclaimed by the TTL backstop, not by + presence. +- The retained local publish is a bounded `try_send`, so under a connect storm a presence + message can be dropped and the retained value left stale. + ## Phasing | PR | Scope | Notes | @@ -177,7 +198,8 @@ Cluster (follow-up): **do not** copy the LWT template. Two verified constraints | 1a ✅ | TTL backstop fix — **agent** | `background.rs`: guard release + `expect_value` on scanned bytes + per-entity batches. Standalone data-loss + seat-lockout bug. | | 1b ✅ | TTL backstop fix — **cluster** | Rewire `handle_ttl_cleanup` through the replicated delete path (primary-gated, guard release, change event, version precondition) — a larger change than 1a, split out to isolate risk. | | 2 ✅ | Client CAS — both paths + new error | `_expected_version` reserved payload key, terminal `PreconditionFailed` (412); agent `update_with_expected`/`delete_with_expected` + both cluster write paths; retry-loop-doesn't-defeat-CAS counter-tests. | -| 3 | Presence feed — both modes + topic rule | Gated on decisions 1–2 below. | +| 3a ✅ | Presence feed — **agent** | `presence.rs` handler + `spawn_presence_task` + `$DB/_presence/# ReadOnly`; opt-in `--presence`. | +| 3b ✅ | Presence feed — **cluster** | `PresenceBroadcast` to all nodes (the `$`-topic wildcard bail makes routing unusable) + `_presence` pass-through; `mqdb dev test --presence`. | | app | Janitor + reassert-on-reconnect contract + short keepalive | Out of mqdb (application). | The TTL fix is split into 1a (agent) and 1b (cluster) because the cluster sweep does a