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

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

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

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.28"
version = "0.8.29"
edition.workspace = true
license = "Apache-2.0"
authors.workspace = true
Expand Down
52 changes: 36 additions & 16 deletions crates/mqdb-agent/src/presence.rs
Original file line number Diff line number Diff line change
Expand Up @@ -88,40 +88,60 @@ impl PresenceEvent {
}
}

pub struct PresenceEventHandler {
sender: flume::Sender<PresenceEvent>,
live_connections: Mutex<HashMap<String, u32>>,
#[derive(Default)]
pub struct LiveConnections {
counts: Mutex<HashMap<String, u32>>,
}

impl PresenceEventHandler {
impl LiveConnections {
#[must_use]
pub fn new(sender: flume::Sender<PresenceEvent>) -> 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<PresenceEvent>,
live_connections: LiveConnections,
}

impl PresenceEventHandler {
#[must_use]
pub fn new(sender: flume::Sender<PresenceEvent>) -> 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() {
Expand Down
2 changes: 1 addition & 1 deletion crates/mqdb-cli/Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[package]
name = "mqdb-cli"
version = "0.8.37"
version = "0.8.38"
publish = false
edition.workspace = true
license = "AGPL-3.0-only"
Expand Down
6 changes: 6 additions & 0 deletions crates/mqdb-cli/src/cli_types/agent.rs
Original file line number Diff line number Diff line change
Expand Up @@ -257,6 +257,12 @@ pub(crate) struct ClusterStartFields {
help = "Scope events by entity field (e.g. diagrams=diagramId)"
)]
pub(crate) event_scope: Option<String>,
#[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",
Expand Down
2 changes: 2 additions & 0 deletions crates/mqdb-cli/src/cli_types/dev.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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")]
Expand Down
4 changes: 4 additions & 0 deletions crates/mqdb-cli/src/commands/cluster.rs
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@ pub(crate) struct ClusterStartArgs {
pub(crate) oauth: OAuthArgs,
pub(crate) ownership: Option<String>,
pub(crate) event_scope: Option<String>,
pub(crate) presence: bool,
pub(crate) passphrase_file: Option<PathBuf>,
pub(crate) passphrase_data: Option<String>,
pub(crate) license: Option<PathBuf>,
Expand Down Expand Up @@ -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);
}
Expand Down
Loading
Loading