Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 16 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
4 changes: 2 additions & 2 deletions Cargo.lock

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

1 change: 1 addition & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 | — |

Expand Down
2 changes: 1 addition & 1 deletion crates/mqdb-agent/Cargo.toml
Original file line number Diff line number Diff line change
@@ -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
Expand Down
16 changes: 16 additions & 0 deletions crates/mqdb-agent/src/agent/broker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::{
Expand Down Expand Up @@ -121,6 +124,19 @@ impl MqdbAgent {
))
}

pub(super) fn apply_presence_handler(
&self,
config: &mut BrokerConfig,
) -> Option<flume::Receiver<PresenceEvent>> {
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}<client_id>");
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())
Expand Down
4 changes: 4 additions & 0 deletions crates/mqdb-agent/src/agent/handlers.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
32 changes: 32 additions & 0 deletions crates/mqdb-agent/src/agent/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@ pub struct MqdbAgent {
pub(super) ownership_config: Arc<mqdb_core::types::OwnershipConfig>,
pub(super) scope_config: Arc<mqdb_core::types::ScopeConfig>,
pub(super) scoped_events: bool,
pub(super) presence: bool,
pub(super) vault_backend: Arc<dyn VaultBackend>,
#[cfg(feature = "http-api")]
pub(super) auth_rate_limiter: Arc<RateLimiter>,
Expand Down Expand Up @@ -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)),
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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<tokio::task::JoinHandle<()>> = {
#[cfg(feature = "http-api")]
{
Expand All @@ -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;
}
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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<tokio::task::JoinHandle<()>> = {
#[cfg(feature = "http-api")]
{
Expand Down Expand Up @@ -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;
}
Expand Down
62 changes: 62 additions & 0 deletions crates/mqdb-agent/src/agent/tasks.rs
Original file line number Diff line number Diff line change
Expand Up @@ -249,6 +249,68 @@ impl MqdbAgent {
})
}

pub(super) fn spawn_presence_task(
&self,
presence_addr: SocketAddr,
presence_service_username: Option<String>,
presence_service_password: Option<String>,
presence_rx: flume::Receiver<crate::presence::PresenceEvent>,
) -> 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,
Expand Down
1 change: 1 addition & 0 deletions crates/mqdb-agent/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
Loading
Loading