diff --git a/CHANGELOG.md b/CHANGELOG.md index d06efcc..29ee408 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,22 @@ 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. 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 ### Changed diff --git a/Cargo.lock b/Cargo.lock index 2d9149f..7de93c4 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 f21de7c..e63d17b 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 0e1cf0d..c2c159e 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 8bbb730..388a848 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 71b488f..cbfb094 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 f34c829..3cf31ba 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 0dc24c1..76a9a86 100644 --- a/crates/mqdb-agent/src/agent/tasks.rs +++ b/crates/mqdb-agent/src/agent/tasks.rs @@ -249,6 +249,68 @@ 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 = resolve_connect_address(presence_addr); + + 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 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(), 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 04a0eb4..cd04991 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 0000000..82c348b --- /dev/null +++ b/crates/mqdb-agent/src/presence.rs @@ -0,0 +1,330 @@ +// 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; +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}") +} + +#[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) + .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) + } + + /// # 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, + 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) { + 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) || !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(), + )); + }) + } + + fn on_client_disconnect<'a>( + &'a self, + event: ClientDisconnectEvent, + ) -> Pin + Send + 'a>> { + Box::pin(async move { + 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( + &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("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"); + 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("serializes")) + .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")); + } + + 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); + 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 0fe48f2..5f5e1c6 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 0000000..66f15fe --- /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 2861846..b9a8b0a 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 f2d75aa..24ad312 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 bf1b67f..ca9ae56 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 7aafdbb..075e055 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