From 2a927acb8654832c7c600d51c4715d1852dad13d Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Fabr=C3=ADcio=20Bracht?= Date: Mon, 14 Sep 2026 20:57:08 -0300 Subject: [PATCH 1/2] add opt-in presence feed for client connect and disconnect in agent mode --- CHANGELOG.md | 13 ++ Cargo.lock | 4 +- README.md | 1 + crates/mqdb-agent/Cargo.toml | 2 +- crates/mqdb-agent/src/agent/broker.rs | 16 ++ crates/mqdb-agent/src/agent/handlers.rs | 4 + crates/mqdb-agent/src/agent/mod.rs | 32 ++++ crates/mqdb-agent/src/agent/tasks.rs | 55 ++++++ crates/mqdb-agent/src/lib.rs | 1 + crates/mqdb-agent/src/presence.rs | 208 +++++++++++++++++++++++ crates/mqdb-agent/src/topic_rules.rs | 24 +++ crates/mqdb-agent/tests/presence_test.rs | 171 +++++++++++++++++++ crates/mqdb-cli/Cargo.toml | 2 +- crates/mqdb-cli/src/cli_types/agent.rs | 6 + crates/mqdb-cli/src/commands/agent.rs | 5 + crates/mqdb-cli/src/main.rs | 1 + docs/design/hold-reclaim.md | 50 ++++-- 17 files changed, 579 insertions(+), 16 deletions(-) create mode 100644 crates/mqdb-agent/src/presence.rs create mode 100644 crates/mqdb-agent/tests/presence_test.rs diff --git a/CHANGELOG.md b/CHANGELOG.md index d06efccd..cca5af44 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,19 @@ 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-12 — mqdb-agent 0.8.28, mqdb-cli 0.8.37 + +### Added + +- **Presence feed (agent mode), opt-in via `--presence` / `MQDB_PRESENCE`.** When enabled, the broker publishes a retained QoS 0 JSON message to `$DB/_presence/{client_id}` on every client connect and disconnect: `{client_id, user_id, event, unexpected, ts}`. `event` is `connect` or `disconnect`; `unexpected` distinguishes an abnormal drop from a clean DISCONNECT. This is the third primitive of on-disconnect hold reclaim (`docs/design/hold-reclaim.md`): a janitor subscribes to `$DB/_presence/#`, applies a grace period, and releases an abandoned hold with a version-guarded (`_expected_version`) delete. Unlike Last Will, presence fires on **both** clean and unexpected disconnects. Internal `mqdb-` clients never produce presence events, and the feed stays silent unless explicitly enabled. +- **`$DB/_presence/#` is a `ReadOnly` protected topic**, so a non-admin janitor may subscribe while no client — not even an admin — can publish to it and forge a presence message; only the internal service publisher can. Without this rule the `_`-prefixed topic would also have denied the janitor's subscribe outright. + +### Notes + +- Presence is a **hint**, not a ledger: messages may be dropped under load (the queue is bounded) and can be reordered across a flap, so `ts` is authoritative. Correctness still rests on the TTL backstop plus a returning holder re-asserting its hold; see the invariant model in `specs/AbandonedHoldReclaim.tla`. +- Because messages are retained, the broker keeps one retained message per distinct `client_id`, overwritten in place on each transition. That is bounded by the client population, but a deployment that uses a fresh random `client_id` for every session will accumulate retained entries — prefer stable client ids when enabling presence. +- Cluster-mode presence is a follow-up. The cluster's wildcard matcher deliberately drops every `$`-prefixed topic (`cluster/topic_trie.rs`), so a `$DB/_presence/#` subscriber cannot be resolved as a cross-node target; cluster support will broadcast presence to all nodes for local delivery rather than relying on that routing. + ## 2026-09-07 — mqdb-agent 0.8.27, mqdb-cluster 0.4.13, mqdb-vault 0.1.5, mqdb-cli 0.8.36 ### Changed diff --git a/Cargo.lock b/Cargo.lock index 2d9149f3..7de93c44 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1420,7 +1420,7 @@ dependencies = [ [[package]] name = "mqdb-agent" -version = "0.8.27" +version = "0.8.28" dependencies = [ "arc-swap", "argon2", @@ -1455,7 +1455,7 @@ dependencies = [ [[package]] name = "mqdb-cli" -version = "0.8.36" +version = "0.8.37" dependencies = [ "base64", "bebytes", diff --git a/README.md b/README.md index f21de7c4..e63d17be 100644 --- a/README.md +++ b/README.md @@ -844,6 +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_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 0e1cf0d7..c2c159eb 100644 --- a/crates/mqdb-agent/Cargo.toml +++ b/crates/mqdb-agent/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "mqdb-agent" -version = "0.8.27" +version = "0.8.28" edition.workspace = true license = "Apache-2.0" authors.workspace = true diff --git a/crates/mqdb-agent/src/agent/broker.rs b/crates/mqdb-agent/src/agent/broker.rs index 8bbb730b..388a848b 100644 --- a/crates/mqdb-agent/src/agent/broker.rs +++ b/crates/mqdb-agent/src/agent/broker.rs @@ -3,6 +3,9 @@ use super::MqdbAgent; use crate::broker_defaults::{BROKER_MAX_CLIENTS, BROKER_MAX_PACKET_SIZE, SESSION_EXPIRY_SECS}; +use crate::presence::{ + PRESENCE_CHANNEL_CAPACITY, PRESENCE_TOPIC_PREFIX, PresenceEvent, PresenceEventHandler, +}; use crate::topic_protection::TopicProtectionAuthProvider; use mqtt5::broker::auth::{CompositeAuthProvider, ComprehensiveAuthProvider}; use mqtt5::broker::config::{ @@ -121,6 +124,19 @@ impl MqdbAgent { )) } + pub(super) fn apply_presence_handler( + &self, + config: &mut BrokerConfig, + ) -> Option> { + if !self.presence { + return None; + } + let (sender, receiver) = flume::bounded(PRESENCE_CHANNEL_CAPACITY); + config.event_handler = Some(Arc::new(PresenceEventHandler::new(sender))); + info!("presence feed enabled on {PRESENCE_TOPIC_PREFIX}"); + Some(receiver) + } + pub(super) fn apply_transport_config(&self, config: &mut BrokerConfig) { if let (Some(cert_file), Some(key_file)) = (&self.quic_cert_file, &self.quic_key_file) { let quic_config = QuicConfig::new(cert_file.clone(), key_file.clone()) diff --git a/crates/mqdb-agent/src/agent/handlers.rs b/crates/mqdb-agent/src/agent/handlers.rs index 71b488f5..cbfb0949 100644 --- a/crates/mqdb-agent/src/agent/handlers.rs +++ b/crates/mqdb-agent/src/agent/handlers.rs @@ -194,6 +194,10 @@ pub(super) async fn handle_message(ctx: &MessageContext<'_>, message: Message) { return; } + if topic.starts_with(crate::presence::PRESENCE_TOPIC_PREFIX) { + return; + } + if let Some(admin_op) = parse_admin_topic(topic) { let admin_ctx = AdminContext { db, diff --git a/crates/mqdb-agent/src/agent/mod.rs b/crates/mqdb-agent/src/agent/mod.rs index f34c829e..3cf31ba4 100644 --- a/crates/mqdb-agent/src/agent/mod.rs +++ b/crates/mqdb-agent/src/agent/mod.rs @@ -39,6 +39,7 @@ pub struct MqdbAgent { pub(super) ownership_config: Arc, pub(super) scope_config: Arc, pub(super) scoped_events: bool, + pub(super) presence: bool, pub(super) vault_backend: Arc, #[cfg(feature = "http-api")] pub(super) auth_rate_limiter: Arc, @@ -76,6 +77,7 @@ impl MqdbAgent { ownership_config: Arc::new(mqdb_core::types::OwnershipConfig::default()), scope_config: Arc::new(mqdb_core::types::ScopeConfig::default()), scoped_events: false, + presence: false, vault_backend: Arc::new(NoopVaultBackend), #[cfg(feature = "http-api")] auth_rate_limiter: Arc::new(RateLimiter::new(10)), @@ -207,6 +209,12 @@ impl MqdbAgent { self } + #[must_use] + pub fn with_presence(mut self, enabled: bool) -> Self { + self.presence = enabled; + self + } + #[must_use] pub fn with_license_expiry(mut self, expires_at: u64) -> Self { self.license_expires_at = Some(expires_at); @@ -259,6 +267,7 @@ impl MqdbAgent { self.build_broker_config().await?; self.apply_transport_config(&mut config); + let presence_rx = self.apply_presence_handler(&mut config); let broker = mqtt5::broker::MqttBroker::with_config(config).await?; let (mut broker, auth_providers) = Self::apply_auth_providers( @@ -291,6 +300,14 @@ impl MqdbAgent { service_username.clone(), service_password.clone(), ); + let presence_task = presence_rx.map(|rx| { + self.spawn_presence_task( + bind_addr, + service_username.clone(), + service_password.clone(), + rx, + ) + }); let http_task: Option> = { #[cfg(feature = "http-api")] { @@ -315,6 +332,9 @@ impl MqdbAgent { let _ = self.shutdown_tx.send(()); let _ = handler_task.await; let _ = event_task.await; + if let Some(presence) = presence_task { + let _ = presence.await; + } if let Some(http) = http_task { let _ = http.await; } @@ -344,6 +364,7 @@ impl MqdbAgent { self.build_broker_config().await?; self.apply_transport_config(&mut config); + let presence_rx = self.apply_presence_handler(&mut config); let broker = mqtt5::broker::MqttBroker::with_config(config).await?; let (mut broker, auth_providers) = Self::apply_auth_providers( @@ -379,6 +400,14 @@ impl MqdbAgent { service_username.clone(), service_password.clone(), ); + let presence_task = presence_rx.map(|rx| { + self.spawn_presence_task( + bind_addr, + service_username.clone(), + service_password.clone(), + rx, + ) + }); let http_task: Option> = { #[cfg(feature = "http-api")] { @@ -414,6 +443,9 @@ impl MqdbAgent { let _ = shutdown_tx.send(()); let _ = handler_task.await; let _ = event_task.await; + if let Some(presence) = presence_task { + let _ = presence.await; + } if let Some(http) = http_task { let _ = http.await; } diff --git a/crates/mqdb-agent/src/agent/tasks.rs b/crates/mqdb-agent/src/agent/tasks.rs index 0dc24c11..eba60302 100644 --- a/crates/mqdb-agent/src/agent/tasks.rs +++ b/crates/mqdb-agent/src/agent/tasks.rs @@ -249,6 +249,61 @@ impl MqdbAgent { }) } + pub(super) fn spawn_presence_task( + &self, + presence_addr: SocketAddr, + presence_service_username: Option, + presence_service_password: Option, + presence_rx: flume::Receiver, + ) -> tokio::task::JoinHandle<()> { + let mut presence_shutdown_rx = self.shutdown_tx.subscribe(); + + tokio::spawn(async move { + tokio::time::sleep(Duration::from_millis(200)).await; + + let client = MqttClient::new("mqdb-presence-publisher"); + let addr = format!("{}:{}", presence_addr.ip(), presence_addr.port()); + + if let Err(e) = connect_mqtt_client( + &client, + "mqdb-presence-publisher", + &addr, + presence_service_username, + presence_service_password, + ) + .await + { + error!("Failed to connect presence publisher: {e}"); + return; + } + + loop { + tokio::select! { + presence = presence_rx.recv_async() => { + let Ok(presence) = presence else { + debug!("Presence channel closed"); + break; + }; + let options = mqtt5::types::PublishOptions { + retain: true, + ..Default::default() + }; + if let Err(e) = client + .publish_with_options(&presence.topic(), presence.payload(), options) + .await + { + warn!("Failed to publish presence: {e}"); + } + } + _ = presence_shutdown_rx.recv() => { + debug!("Presence publisher shutting down"); + break; + } + } + } + }) + } + #[cfg(feature = "http-api")] pub(super) fn spawn_http_task( &self, diff --git a/crates/mqdb-agent/src/lib.rs b/crates/mqdb-agent/src/lib.rs index 04a0eb45..cd049916 100644 --- a/crates/mqdb-agent/src/lib.rs +++ b/crates/mqdb-agent/src/lib.rs @@ -11,6 +11,7 @@ pub mod db_helpers; pub mod dedup; pub mod dispatcher; pub mod outbox_processor; +pub mod presence; pub mod rate_limiter; pub mod runtime; pub mod session; diff --git a/crates/mqdb-agent/src/presence.rs b/crates/mqdb-agent/src/presence.rs new file mode 100644 index 00000000..6aab75f8 --- /dev/null +++ b/crates/mqdb-agent/src/presence.rs @@ -0,0 +1,208 @@ +// Copyright 2025-2026 LabOverWire. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +use std::future::Future; +use std::pin::Pin; + +use mqtt5::broker::events::{BrokerEventHandler, ClientConnectEvent, ClientDisconnectEvent}; +use serde::Serialize; +use tracing::debug; + +pub const PRESENCE_TOPIC_PREFIX: &str = "$DB/_presence/"; + +pub const PRESENCE_CHANNEL_CAPACITY: usize = 1024; + +const INTERNAL_CLIENT_PREFIX: &str = "mqdb-"; + +#[must_use] +pub fn is_internal_client(client_id: &str) -> bool { + client_id.starts_with(INTERNAL_CLIENT_PREFIX) +} + +#[must_use] +pub fn presence_topic(client_id: &str) -> String { + format!("{PRESENCE_TOPIC_PREFIX}{client_id}") +} + +fn current_time_ms() -> u64 { + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .map_or(0, |d| d.as_secs() * 1000 + u64::from(d.subsec_millis())) +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)] +#[serde(rename_all = "lowercase")] +pub enum PresenceState { + Connect, + Disconnect, +} + +#[derive(Debug, Clone, Serialize)] +pub struct PresenceEvent { + pub client_id: String, + pub user_id: Option, + pub event: PresenceState, + pub unexpected: bool, + pub ts: u64, +} + +impl PresenceEvent { + #[must_use] + pub fn connect(client_id: &str, user_id: Option<&str>) -> Self { + Self { + client_id: client_id.to_string(), + user_id: user_id.map(str::to_string), + event: PresenceState::Connect, + unexpected: false, + ts: current_time_ms(), + } + } + + #[must_use] + pub fn disconnect(client_id: &str, user_id: Option<&str>, unexpected: bool) -> Self { + Self { + client_id: client_id.to_string(), + user_id: user_id.map(str::to_string), + event: PresenceState::Disconnect, + unexpected, + ts: current_time_ms(), + } + } + + #[must_use] + pub fn topic(&self) -> String { + presence_topic(&self.client_id) + } + + #[must_use] + pub fn payload(&self) -> Vec { + serde_json::to_vec(self).unwrap_or_default() + } +} + +pub struct PresenceEventHandler { + sender: flume::Sender, +} + +impl PresenceEventHandler { + #[must_use] + pub fn new(sender: flume::Sender) -> Self { + Self { sender } + } + + fn emit(&self, presence: PresenceEvent) { + if self.sender.try_send(presence).is_err() { + debug!("presence queue full or closed, dropping presence event"); + } + } +} + +impl BrokerEventHandler for PresenceEventHandler { + fn on_client_connect<'a>( + &'a self, + event: ClientConnectEvent, + ) -> Pin + Send + 'a>> { + Box::pin(async move { + if is_internal_client(&event.client_id) { + return; + } + self.emit(PresenceEvent::connect( + &event.client_id, + event.user_id.as_deref(), + )); + }) + } + + fn on_client_disconnect<'a>( + &'a self, + event: ClientDisconnectEvent, + ) -> Pin + Send + 'a>> { + Box::pin(async move { + if is_internal_client(&event.client_id) { + return; + } + self.emit(PresenceEvent::disconnect( + &event.client_id, + event.user_id.as_deref(), + event.unexpected, + )); + }) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn connect_payload_shape() { + let presence = PresenceEvent::connect("seat-holder-7", Some("alice")); + let value: serde_json::Value = + serde_json::from_slice(&presence.payload()).expect("payload is valid JSON"); + assert_eq!(value["client_id"], "seat-holder-7"); + assert_eq!(value["user_id"], "alice"); + assert_eq!(value["event"], "connect"); + assert_eq!(value["unexpected"], false); + assert!(value["ts"].as_u64().is_some_and(|ts| ts > 0)); + } + + #[test] + fn disconnect_payload_preserves_unexpected_and_anonymous_user() { + let presence = PresenceEvent::disconnect("seat-holder-7", None, true); + let value: serde_json::Value = + serde_json::from_slice(&presence.payload()).expect("payload is valid JSON"); + assert_eq!(value["event"], "disconnect"); + assert_eq!(value["unexpected"], true); + assert!(value["user_id"].is_null()); + } + + #[test] + fn topic_is_namespaced_per_client() { + let presence = PresenceEvent::connect("seat-holder-7", None); + assert_eq!(presence.topic(), "$DB/_presence/seat-holder-7"); + } + + #[test] + fn internal_clients_are_recognized() { + assert!(is_internal_client("mqdb-admin-1")); + assert!(is_internal_client("mqdb-forward-2")); + assert!(is_internal_client("mqdb-presence-publisher")); + assert!(!is_internal_client("seat-holder-7")); + assert!(!is_internal_client("janitor")); + } + + #[tokio::test] + async fn handler_emits_connect_and_skips_internal_clients() { + let (tx, rx) = flume::bounded(PRESENCE_CHANNEL_CAPACITY); + let handler = PresenceEventHandler::new(tx); + + handler + .on_client_disconnect(ClientDisconnectEvent { + client_id: "mqdb-admin-1".into(), + user_id: None, + reason: mqtt5::types::ReasonCode::Success, + unexpected: false, + }) + .await; + assert!( + rx.is_empty(), + "internal clients must not produce presence events" + ); + + handler + .on_client_disconnect(ClientDisconnectEvent { + client_id: "seat-holder-7".into(), + user_id: Some("alice".into()), + reason: mqtt5::types::ReasonCode::Success, + unexpected: false, + }) + .await; + let presence = rx.try_recv().expect("a presence event was queued"); + assert_eq!(presence.client_id, "seat-holder-7"); + assert_eq!(presence.event, PresenceState::Disconnect); + assert!( + !presence.unexpected, + "a clean disconnect must still emit presence" + ); + } +} diff --git a/crates/mqdb-agent/src/topic_rules.rs b/crates/mqdb-agent/src/topic_rules.rs index 0fe48f2f..5f5e1c62 100644 --- a/crates/mqdb-agent/src/topic_rules.rs +++ b/crates/mqdb-agent/src/topic_rules.rs @@ -50,6 +50,10 @@ pub const PROTECTED_TOPICS: &[TopicRule] = &[ pattern: "$DB/u/#", tier: ProtectionTier::ReadOnly, }, + TopicRule { + pattern: "$DB/_presence/#", + tier: ProtectionTier::ReadOnly, + }, TopicRule { pattern: "$DB/_sub/#", tier: ProtectionTier::WriteOnly, @@ -392,6 +396,26 @@ mod tests { ); } + #[test] + fn check_access_presence_read_only() { + assert_eq!( + check_topic_access("$DB/_presence/seat-holder-7", false, false), + Ok(()), + "a non-admin janitor must be able to subscribe to presence" + ); + assert_eq!(check_topic_access("$DB/_presence/#", false, false), Ok(())); + assert_eq!( + check_topic_access("$DB/_presence/seat-holder-7", true, false), + Err(BlockReason::ReadOnlyTopic), + "a client must not be able to forge a presence message" + ); + assert_eq!( + check_topic_access("$DB/_presence/seat-holder-7", true, true), + Err(BlockReason::ReadOnlyTopic), + "not even an admin may publish presence; only the internal service bypass may" + ); + } + #[test] fn check_access_user_namespace_read_only() { assert_eq!( diff --git a/crates/mqdb-agent/tests/presence_test.rs b/crates/mqdb-agent/tests/presence_test.rs new file mode 100644 index 00000000..66f15fe3 --- /dev/null +++ b/crates/mqdb-agent/tests/presence_test.rs @@ -0,0 +1,171 @@ +// Copyright 2025-2026 LabOverWire. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +mod common; + +use common::next_test_port; +use mqdb_agent::{Database, MqdbAgent}; +use mqtt5::client::MqttClient; +use serde_json::Value; +use std::net::SocketAddr; +use std::time::Duration; +use tempfile::TempDir; + +async fn start_agent_with_presence( + port: u16, + presence: bool, +) -> (TempDir, tokio::task::JoinHandle<()>) { + let tmp = TempDir::new().unwrap(); + let db = Database::open(tmp.path()).await.unwrap(); + let addr: SocketAddr = format!("127.0.0.1:{port}").parse().unwrap(); + let agent = MqdbAgent::new(db) + .with_bind_address(addr) + .with_anonymous(true) + .with_presence(presence); + let (handle, mut ready_rx, _shutdown) = agent.start().await.unwrap(); + let _ = ready_rx.changed().await; + (tmp, handle) +} + +async fn subscribe_to_presence(port: u16) -> (MqttClient, flume::Receiver<(String, Value)>) { + let janitor = MqttClient::new("janitor"); + janitor + .connect(&format!("mqtt://127.0.0.1:{port}")) + .await + .expect("a non-admin janitor must be able to connect"); + + let (tx, rx) = flume::bounded(64); + janitor + .subscribe("$DB/_presence/#", move |msg| { + if msg.topic.ends_with("/janitor") { + return; + } + if let Ok(value) = serde_json::from_slice::(&msg.payload) { + let _ = tx.try_send((msg.topic.to_string(), value)); + } + }) + .await + .expect("a non-admin janitor must be allowed to subscribe to presence"); + + tokio::time::sleep(Duration::from_millis(150)).await; + (janitor, rx) +} + +async fn next_presence(rx: &flume::Receiver<(String, Value)>) -> (String, Value) { + tokio::time::timeout(Duration::from_secs(3), rx.recv_async()) + .await + .expect("a presence message should arrive") + .expect("presence channel stays open") +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn presence_reports_connect_and_clean_disconnect() { + let port = next_test_port(); + let (_tmp, agent_handle) = start_agent_with_presence(port, true).await; + let (janitor, rx) = subscribe_to_presence(port).await; + + let holder = MqttClient::new("seat-holder-7"); + holder + .connect(&format!("mqtt://127.0.0.1:{port}")) + .await + .unwrap(); + + let (topic, connect_event) = next_presence(&rx).await; + assert_eq!(topic, "$DB/_presence/seat-holder-7"); + assert_eq!(connect_event["client_id"], "seat-holder-7"); + assert_eq!(connect_event["event"], "connect"); + assert!(connect_event["ts"].as_u64().is_some_and(|ts| ts > 0)); + + holder.disconnect().await.unwrap(); + + let (topic, disconnect_event) = next_presence(&rx).await; + assert_eq!(topic, "$DB/_presence/seat-holder-7"); + assert_eq!( + disconnect_event["event"], "disconnect", + "a clean disconnect must still emit presence, unlike LWT which is unexpected-only" + ); + assert_eq!(disconnect_event["unexpected"], false); + + janitor.disconnect().await.unwrap(); + agent_handle.abort(); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn presence_is_retained_for_late_subscribers() { + let port = next_test_port(); + let (_tmp, agent_handle) = start_agent_with_presence(port, true).await; + + let holder = MqttClient::new("seat-holder-late"); + holder + .connect(&format!("mqtt://127.0.0.1:{port}")) + .await + .unwrap(); + tokio::time::sleep(Duration::from_millis(400)).await; + + let (janitor, rx) = subscribe_to_presence(port).await; + let (topic, retained) = next_presence(&rx).await; + assert_eq!(topic, "$DB/_presence/seat-holder-late"); + assert_eq!( + retained["event"], "connect", + "a janitor subscribing after the fact must still learn the current state" + ); + + janitor.disconnect().await.unwrap(); + holder.disconnect().await.unwrap(); + agent_handle.abort(); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn presence_clients_cannot_forge_presence() { + let port = next_test_port(); + let (_tmp, agent_handle) = start_agent_with_presence(port, true).await; + let (janitor, rx) = subscribe_to_presence(port).await; + + let attacker = MqttClient::new("attacker"); + attacker + .connect(&format!("mqtt://127.0.0.1:{port}")) + .await + .unwrap(); + + let forged = br#"{"client_id":"seat-holder-7","event":"disconnect","unexpected":true,"ts":1}"#; + let _ = attacker + .publish("$DB/_presence/seat-holder-7", forged.to_vec()) + .await; + + tokio::time::sleep(Duration::from_millis(400)).await; + + let forged_delivered = rx + .try_iter() + .any(|(_, value)| value["ts"].as_u64() == Some(1)); + assert!( + !forged_delivered, + "the presence topic is ReadOnly, so a client must never be able to forge a presence message" + ); + + attacker.disconnect().await.unwrap(); + janitor.disconnect().await.unwrap(); + agent_handle.abort(); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn presence_is_off_by_default() { + let port = next_test_port(); + let (_tmp, agent_handle) = start_agent_with_presence(port, false).await; + let (janitor, rx) = subscribe_to_presence(port).await; + + let holder = MqttClient::new("seat-holder-quiet"); + holder + .connect(&format!("mqtt://127.0.0.1:{port}")) + .await + .unwrap(); + tokio::time::sleep(Duration::from_millis(400)).await; + + assert!( + rx.is_empty(), + "presence must stay silent unless it is explicitly enabled" + ); + + holder.disconnect().await.unwrap(); + janitor.disconnect().await.unwrap(); + agent_handle.abort(); +} diff --git a/crates/mqdb-cli/Cargo.toml b/crates/mqdb-cli/Cargo.toml index 28618466..b9a8b0a9 100644 --- a/crates/mqdb-cli/Cargo.toml +++ b/crates/mqdb-cli/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "mqdb-cli" -version = "0.8.36" +version = "0.8.37" 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 f2d75aaa..24ad3129 100644 --- a/crates/mqdb-cli/src/cli_types/agent.rs +++ b/crates/mqdb-cli/src/cli_types/agent.rs @@ -113,6 +113,12 @@ pub(crate) struct AgentStartFields { help = "Route change events to per-recipient topics ($DB/u/{user}/events/#) for ownership-enabled entities; subscribers must subscribe to their own namespace" )] pub(crate) scoped_events: bool, + #[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_OTLP_ENDPOINT", diff --git a/crates/mqdb-cli/src/commands/agent.rs b/crates/mqdb-cli/src/commands/agent.rs index bf1b67f7..ca9ae563 100644 --- a/crates/mqdb-cli/src/commands/agent.rs +++ b/crates/mqdb-cli/src/commands/agent.rs @@ -31,6 +31,7 @@ pub(crate) struct AgentStartArgs { pub(crate) ownership: Option, pub(crate) ownership_derive: Option, pub(crate) scoped_events: bool, + pub(crate) presence: bool, pub(crate) event_scope: Option, pub(crate) passphrase_file: Option, pub(crate) passphrase_data: Option, @@ -160,6 +161,10 @@ pub(crate) async fn cmd_agent_start( agent = agent.with_scoped_events(true); } + if args.presence { + agent = agent.with_presence(true); + } + #[cfg(feature = "opentelemetry")] if let Some(ref endpoint) = args.otlp_endpoint { let telemetry_config = mqtt5::telemetry::TelemetryConfig::new(&args.otel_service_name) diff --git a/crates/mqdb-cli/src/main.rs b/crates/mqdb-cli/src/main.rs index 7aafdbb7..075e0553 100644 --- a/crates/mqdb-cli/src/main.rs +++ b/crates/mqdb-cli/src/main.rs @@ -245,6 +245,7 @@ async fn dispatch_agent(action: AgentAction) -> Result<(), Box`, which the `mqdb-forward-` bypass does not match). So `_presence` + must be added to that allow-list. +- `TopicTrie::match_topic` returns nothing for any `$`-prefixed topic, so a + `$DB/_presence/#` subscriber can never be resolved as a cross-node target. Cross-node + delivery therefore **broadcasts** presence to every node, each publishing locally (the + broker's own matcher is spec-correct and does match `$DB/_presence/#`), rather than + routing via `route_and_forward_publish`. A happy side effect: every node ends up holding + the retained set, so a janitor sees the full picture from whichever node it attaches to. +- The janitor must **not** be named `mqdb-*`: `on_client_subscribe` early-returns on that + prefix, so its subscription would never be registered cluster-wide. ## Phasing From d6a40d006df6a95d02ccbcf4c9c48f94ddaf6a93 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Fabr=C3=ADcio=20Bracht?= Date: Fri, 18 Sep 2026 15:23:52 -0300 Subject: [PATCH 2/2] suppress presence disconnect on session takeover and fix publisher address and payload handling --- CHANGELOG.md | 7 +- crates/mqdb-agent/src/agent/tasks.rs | 11 ++- crates/mqdb-agent/src/presence.rs | 138 +++++++++++++++++++++++++-- docs/design/hold-reclaim.md | 12 ++- 4 files changed, 155 insertions(+), 13 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index cca5af44..29ee408b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -13,8 +13,11 @@ Each entry lists the date and the crate versions that were released. ### Notes -- Presence is a **hint**, not a ledger: messages may be dropped under load (the queue is bounded) and can be reordered across a flap, so `ts` is authoritative. Correctness still rests on the TTL backstop plus a returning holder re-asserting its hold; see the invariant model in `specs/AbandonedHoldReclaim.tla`. -- Because messages are retained, the broker keeps one retained message per distinct `client_id`, overwritten in place on each transition. That is bounded by the client population, but a deployment that uses a fresh random `client_id` for every session will accumulate retained entries — prefer stable client ids when enabling presence. +- Presence is a **hint**, not a ledger: messages may be dropped under load (the queue is bounded) and can be reordered across a flap, so `ts` is authoritative. The two directions are not equally benign — a dropped `disconnect` merely delays reclaim to the TTL backstop, whereas a dropped `connect` can leave a live client retained as disconnected. Correctness therefore still rests on the TTL backstop plus a returning holder re-asserting its hold; see the invariant model in `specs/AbandonedHoldReclaim.tla`. +- A session takeover (a client reconnecting with the same `client_id`) fires the displaced connection's disconnect *after* the new connection's connect, so presence tracks live connections per client and only reports `disconnect` when the last one closes. Without that, a reconnecting holder would be retained as disconnected and could be reclaimed while still connected. +- Because messages are retained, the broker keeps one retained message per distinct `client_id`, overwritten in place on each transition. That is bounded by the client population, but a deployment that uses a fresh random `client_id` for every session will accumulate retained entries — prefer stable client ids when enabling presence. Note that the `ReadOnly` rule also prevents operators from clearing an entry by publishing an empty payload; only the internal service publisher can overwrite one. +- Any authenticated subscriber may read the whole feed: it is not scoped per user, so every subscriber sees every `client_id`, `user_id` and connection timing. Enable `--presence` only where that is acceptable, and restrict `$DB/_presence/#` by ACL if subscribers are not all trusted. +- Client ids that cannot form a valid topic (those containing `+` or `#`) are skipped, and internal `mqdb-`-prefixed clients are excluded by client-id prefix, so a client naming itself `mqdb-…` simply omits itself from the feed and falls back to the TTL backstop. - Cluster-mode presence is a follow-up. The cluster's wildcard matcher deliberately drops every `$`-prefixed topic (`cluster/topic_trie.rs`), so a `$DB/_presence/#` subscriber cannot be resolved as a cross-node target; cluster support will broadcast presence to all nodes for local delivery rather than relying on that routing. ## 2026-09-07 — mqdb-agent 0.8.27, mqdb-cluster 0.4.13, mqdb-vault 0.1.5, mqdb-cli 0.8.36 diff --git a/crates/mqdb-agent/src/agent/tasks.rs b/crates/mqdb-agent/src/agent/tasks.rs index eba60302..76a9a862 100644 --- a/crates/mqdb-agent/src/agent/tasks.rs +++ b/crates/mqdb-agent/src/agent/tasks.rs @@ -262,7 +262,7 @@ impl MqdbAgent { tokio::time::sleep(Duration::from_millis(200)).await; let client = MqttClient::new("mqdb-presence-publisher"); - let addr = format!("{}:{}", presence_addr.ip(), presence_addr.port()); + let addr = resolve_connect_address(presence_addr); if let Err(e) = connect_mqtt_client( &client, @@ -284,12 +284,19 @@ impl MqdbAgent { debug!("Presence channel closed"); break; }; + let payload = match presence.payload() { + Ok(payload) => payload, + Err(e) => { + error!("Failed to serialize presence: {e}"); + continue; + } + }; let options = mqtt5::types::PublishOptions { retain: true, ..Default::default() }; if let Err(e) = client - .publish_with_options(&presence.topic(), presence.payload(), options) + .publish_with_options(&presence.topic(), payload, options) .await { warn!("Failed to publish presence: {e}"); diff --git a/crates/mqdb-agent/src/presence.rs b/crates/mqdb-agent/src/presence.rs index 6aab75f8..82c348b1 100644 --- a/crates/mqdb-agent/src/presence.rs +++ b/crates/mqdb-agent/src/presence.rs @@ -1,8 +1,10 @@ // Copyright 2025-2026 LabOverWire. All rights reserved. // SPDX-License-Identifier: Apache-2.0 +use std::collections::HashMap; use std::future::Future; use std::pin::Pin; +use std::sync::Mutex; use mqtt5::broker::events::{BrokerEventHandler, ClientConnectEvent, ClientDisconnectEvent}; use serde::Serialize; @@ -24,6 +26,11 @@ pub fn presence_topic(client_id: &str) -> String { format!("{PRESENCE_TOPIC_PREFIX}{client_id}") } +#[must_use] +pub fn is_publishable_client_id(client_id: &str) -> bool { + mqtt5::is_valid_topic_name(&presence_topic(client_id)) +} + fn current_time_ms() -> u64 { std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) @@ -74,20 +81,46 @@ impl PresenceEvent { presence_topic(&self.client_id) } - #[must_use] - pub fn payload(&self) -> Vec { - serde_json::to_vec(self).unwrap_or_default() + /// # Errors + /// Returns an error if the presence event cannot be serialized to JSON. + pub fn payload(&self) -> Result, serde_json::Error> { + serde_json::to_vec(self) } } pub struct PresenceEventHandler { sender: flume::Sender, + live_connections: Mutex>, } impl PresenceEventHandler { #[must_use] pub fn new(sender: flume::Sender) -> Self { - Self { sender } + Self { + sender, + live_connections: Mutex::new(HashMap::new()), + } + } + + 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; + } + } + + fn is_last_connection(&self, client_id: &str) -> bool { + let Ok(mut live) = self.live_connections.lock() else { + return true; + }; + let Some(remaining) = live.get_mut(client_id) else { + return true; + }; + *remaining = remaining.saturating_sub(1); + if *remaining == 0 { + live.remove(client_id); + return true; + } + false } fn emit(&self, presence: PresenceEvent) { @@ -103,9 +136,10 @@ impl BrokerEventHandler for PresenceEventHandler { event: ClientConnectEvent, ) -> Pin + Send + 'a>> { Box::pin(async move { - if is_internal_client(&event.client_id) { + if is_internal_client(&event.client_id) || !is_publishable_client_id(&event.client_id) { return; } + self.register_connection(&event.client_id); self.emit(PresenceEvent::connect( &event.client_id, event.user_id.as_deref(), @@ -118,7 +152,14 @@ impl BrokerEventHandler for PresenceEventHandler { event: ClientDisconnectEvent, ) -> Pin + Send + 'a>> { Box::pin(async move { - if is_internal_client(&event.client_id) { + if is_internal_client(&event.client_id) || !is_publishable_client_id(&event.client_id) { + return; + } + if !self.is_last_connection(&event.client_id) { + debug!( + client_id = %event.client_id, + "session taken over by a newer connection, suppressing presence disconnect" + ); return; } self.emit(PresenceEvent::disconnect( @@ -138,7 +179,8 @@ mod tests { fn connect_payload_shape() { let presence = PresenceEvent::connect("seat-holder-7", Some("alice")); let value: serde_json::Value = - serde_json::from_slice(&presence.payload()).expect("payload is valid JSON"); + serde_json::from_slice(&presence.payload().expect("serializes")) + .expect("payload is valid JSON"); assert_eq!(value["client_id"], "seat-holder-7"); assert_eq!(value["user_id"], "alice"); assert_eq!(value["event"], "connect"); @@ -150,7 +192,8 @@ mod tests { fn disconnect_payload_preserves_unexpected_and_anonymous_user() { let presence = PresenceEvent::disconnect("seat-holder-7", None, true); let value: serde_json::Value = - serde_json::from_slice(&presence.payload()).expect("payload is valid JSON"); + serde_json::from_slice(&presence.payload().expect("serializes")) + .expect("payload is valid JSON"); assert_eq!(value["event"], "disconnect"); assert_eq!(value["unexpected"], true); assert!(value["user_id"].is_null()); @@ -171,6 +214,85 @@ mod tests { assert!(!is_internal_client("janitor")); } + fn connect_event(client_id: &str) -> ClientConnectEvent { + ClientConnectEvent { + client_id: client_id.into(), + user_id: None, + clean_start: true, + session_expiry_interval: 0, + will_topic: None, + will_payload: None, + will_qos: None, + will_retain: None, + } + } + + fn disconnect_event(client_id: &str, unexpected: bool) -> ClientDisconnectEvent { + ClientDisconnectEvent { + client_id: client_id.into(), + user_id: None, + reason: mqtt5::types::ReasonCode::Success, + unexpected, + } + } + + #[test] + fn wildcard_client_ids_are_not_publishable() { + assert!(is_publishable_client_id("seat-holder-7")); + assert!(!is_publishable_client_id("#")); + assert!(!is_publishable_client_id("+")); + assert!(!is_publishable_client_id("seat+holder")); + } + + #[tokio::test] + async fn session_takeover_does_not_report_a_live_client_as_disconnected() { + let (tx, rx) = flume::bounded(PRESENCE_CHANNEL_CAPACITY); + let handler = PresenceEventHandler::new(tx); + + handler + .on_client_connect(connect_event("seat-holder-7")) + .await; + let _ = rx.try_recv().expect("first connect is reported"); + + handler + .on_client_connect(connect_event("seat-holder-7")) + .await; + let _ = rx.try_recv().expect("the takeover connect is reported"); + + handler + .on_client_disconnect(disconnect_event("seat-holder-7", true)) + .await; + assert!( + rx.is_empty(), + "the displaced connection must not retain a live client as disconnected" + ); + + handler + .on_client_disconnect(disconnect_event("seat-holder-7", true)) + .await; + let presence = rx + .try_recv() + .expect("the last connection reports disconnect"); + assert_eq!(presence.event, PresenceState::Disconnect); + assert!(presence.unexpected); + } + + #[tokio::test] + async fn wildcard_client_ids_emit_nothing() { + let (tx, rx) = flume::bounded(PRESENCE_CHANNEL_CAPACITY); + let handler = PresenceEventHandler::new(tx); + + handler.on_client_connect(connect_event("#")).await; + handler + .on_client_disconnect(disconnect_event("#", true)) + .await; + + assert!( + rx.is_empty(), + "a wildcard client id would form an invalid publish topic" + ); + } + #[tokio::test] async fn handler_emits_connect_and_skips_internal_clients() { let (tx, rx) = flume::bounded(PRESENCE_CHANNEL_CAPACITY); diff --git a/docs/design/hold-reclaim.md b/docs/design/hold-reclaim.md index 442f10b6..f9f0b79c 100644 --- a/docs/design/hold-reclaim.md +++ b/docs/design/hold-reclaim.md @@ -130,7 +130,17 @@ invisible in the feed. Retained means one retained message per distinct `client_id`, overwritten in place. Bounded by the client population; a deployment minting a fresh random `client_id` per -session will accumulate entries. +session will accumulate entries, and the `ReadOnly` rule blocks operators from clearing +one with an empty payload — only the internal publisher can overwrite it. + +Two failure modes are asymmetric and worth stating plainly. A dropped `disconnect` only +defers reclaim to the TTL backstop, but a dropped `connect` leaves a live client retained +as disconnected. Likewise a **session takeover** fires the displaced connection's +disconnect *after* the new connection's connect (mqtt5 `register_client` precedes +`fire_connect_event`, and `fire_disconnect_event` runs unconditionally regardless of +`session_taken_over`), so the handler counts live connections per `client_id` and reports +`disconnect` only when the last one closes. Presence is also unscoped: every subscriber +sees every client id, user id and connection timing. Agent (done): `presence.rs` holds the shared payload plus a `PresenceEventHandler` implementing `BrokerEventHandler` (the agent had none — `config.event_handler` was always