From f09c1a29b59c09cbb16e982b62e862bd2dfe6b0a Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Fabr=C3=ADcio=20Bracht?= Date: Thu, 1 Oct 2026 20:56:11 -0700 Subject: [PATCH] upgrade mqtt5 to 0.45.1 --- CHANGELOG.md | 28 ++++++++++++++++++++++++++++ Cargo.lock | 16 ++++++++-------- Cargo.toml | 2 +- crates/mqdb-agent/Cargo.toml | 2 +- crates/mqdb-cli/Cargo.toml | 2 +- crates/mqdb-cluster/Cargo.toml | 2 +- crates/mqdb-vault/Cargo.toml | 2 +- docs/design/hold-reclaim.md | 12 ++++++------ 8 files changed, 47 insertions(+), 19 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index ee1b92c..dbbea9a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,34 @@ 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-10-01 — mqdb-agent 0.8.30, mqdb-cluster 0.4.16, mqdb-vault 0.1.6, mqdb-cli 0.8.40 + +### Changed + +- **Upgraded `mqtt5` from 0.39.2 to 0.45.1** (and `mqtt5-protocol` from 0.15.0 to 0.15.2). No MQDB source changes were needed. The broker-side changes MQDB users will notice: + - **Agent mode migrates its MQTT session storage on first start, one way.** `/mqtt_storage` moves to storage format 2 (sessions in one append-only `sessions/sessions.log`), and older builds refuse a migrated directory. **Back up `/mqtt_storage` before upgrading**; rolling back means restoring that backup. Verified upgrading a 0.39.2 data directory: records, the persistent session and a message queued while the client was offline all survive. + - **Session expiry follows MQTT v5 §3.1.2.11.2.** An MQTT v5 CONNECT without a Session Expiry Interval now ends its session at disconnect, so a client that resumes with Clean Start 0 must send a non-zero Session Expiry Interval. + - **MQDB's session expiry setting (1 hour) is now a maximum.** It was previously ignored. A client that asks for longer, or an MQTT v3.1.1 CleanSession=0 client, is granted 1 hour. + +### Fixed + +- **A timed-out send no longer corrupts the stream to a peer** (#143, #145). Inter-node sends wrote each frame with a timeout around `write_all`, which is not cancellation-safe: a timeout mid-frame left a partial length-prefixed frame on the stream, and every later frame to that peer was misread. Each peer connection now has a writer task that owns its send stream and writes whole frames with no timeout. Sends go through a bounded per-peer queue; when it is full the whole frame is dropped instead of truncated, and a slow peer no longer blocks a broadcast to the others. +- **Heartbeats and Raft traffic are no longer delayed or dropped behind bulk data** (#146, #147). Each peer has two queues: a control queue (256 frames: heartbeats, Raft, death and drain notices, partition updates), drained first, and the bulk queue (1024 frames) for everything else. A full bulk queue no longer drops control frames, which under sustained load could cause false node-death detection and election churn. Both queues still share one stream, so a control frame can wait behind the one bulk frame being written. +- **A peer whose connection dies is removed and re-dialled** (#146, #148). A dead connection used to stay in the peer map forever and was never re-dialled, so the peer stayed unreachable, and still counted as linked, until it happened to connect inbound. A failed write now removes the entry, guarded by a per-connection generation so a stale failure cannot remove a newer connection, and every 60 seconds each `--peers` peer the heartbeat reports as not alive is dialled again. The design is model-checked in `specs/ClusterRedial.tla`. +- **Peer re-dial is faster and cannot be stalled by one unreachable peer** (#149, #150). A re-dial now also runs when a node death is detected, so a peer that is back within the dead-detection window reconnects in about 15–17 seconds on a 5-node cluster instead of up to 75. Only peers the heartbeat reports as dead or unknown are re-dialled, so a briefly late but live peer is left alone, and dials run concurrently with a 5-second timeout each. +- **A node that missed a subscription broadcast no longer deletes the replicated subscription** (#141, #151). Subscription reconciliation trusted the local topic index, which is filled only by a one-hop broadcast. A node that held a client's replicated subscription record but missed the broadcast deleted the record, passed the loss on to later replicas, and never routed to that subscriber. The record is now the source of truth for clients whose session partition the node holds as primary or replica: reconciliation adds missing index entries from it and never changes it. Copies held by former owners are ignored. Exact response topics (`resp/…`) are skipped in reconciliation and in startup recovery, since they are never broadcast. The unused wildcard retry timer (`WildcardPendingStore`) is removed. Model: `specs/ClusterSubReconcile.tla`. + +### Changed (dev tooling) + +- **Dev clusters use a full peer mesh** (#156). `mqdb dev start-cluster` now defaults to the `full` topology, and the clusters started by the ownership, sharing and presence suites give every node every other node in `--peers`, which the cluster requires. The previous default, `partial`, started each node with only lower-numbered peers; on that topology a node that starts with no peers elects itself, two Raft leaders can be elected in one term, and partitions can keep moving for over a minute (#155). `partial` remains the default with `--no-quic` and can be selected explicitly. +- **`mqdb dev test` waits for the partition map to settle before writing** (#156). Readiness used to mean only that every node accepted a connection. The suites now wait until three consecutive polls show every node reporting the same partition map, with every partition assigned and every started node holding primaries. On a full mesh this takes 5–8 s on 5 nodes. + +### Notes + +- **`mqdb-cluster` API:** `WildcardPendingStore`, `PendingWildcard` and `WILDCARD_RECONCILIATION_INTERVAL_MS` are no longer exported, and `ReconciliationResult` replaces `subscriptions_added`/`subscriptions_removed` with `index_entries_restored`. `SubscriptionCache::reconcile` takes a `holds_partition` predicate. +- Still open from this release's cluster work: a node that holds no copy of a client's subscription record and misses the broadcast is never repaired (#140); subscription record writes are built from the handling node's local copy (#152); a record created just before its partition moves is missing on the new primary (#153); two Raft leaders can be elected in one term when nodes start without the full peer list (#155). +- The broker's `ClientConnectEvent.clean_start` now reports whether a session was resumed, not the CONNECT flag, so a client's first connect with Clean Start 0 is still reported as clean. Cluster mode uses that flag to decide whether to clear a client's subscriptions on disconnect; moving that decision to the granted session expiry is a follow-up. + ## 2026-09-20 — mqdb-cluster 0.4.15, mqdb-cli 0.8.39 ### Fixed diff --git a/Cargo.lock b/Cargo.lock index ecc86e1..d17876a 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1494,7 +1494,7 @@ dependencies = [ [[package]] name = "mqdb-agent" -version = "0.8.29" +version = "0.8.30" dependencies = [ "arc-swap", "argon2", @@ -1529,7 +1529,7 @@ dependencies = [ [[package]] name = "mqdb-cli" -version = "0.8.39" +version = "0.8.40" dependencies = [ "base64 0.22.1", "bebytes", @@ -1556,7 +1556,7 @@ dependencies = [ [[package]] name = "mqdb-cluster" -version = "0.4.15" +version = "0.4.16" dependencies = [ "arc-swap", "bebytes", @@ -1612,7 +1612,7 @@ dependencies = [ [[package]] name = "mqdb-vault" -version = "0.1.5" +version = "0.1.6" dependencies = [ "base64 0.22.1", "mqdb-agent", @@ -1630,9 +1630,9 @@ dependencies = [ [[package]] name = "mqtt5" -version = "0.39.2" +version = "0.45.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c61d9de73200ed2112ed25ebcdd5e74b22aa5c89065fd3554e111e3353e9d5f0" +checksum = "9455d95f74762a1bb0bfa762cb9a46609779f09182863932f3ab6be55661df06" dependencies = [ "argon2", "base64 0.22.1", @@ -1678,9 +1678,9 @@ dependencies = [ [[package]] name = "mqtt5-protocol" -version = "0.15.0" +version = "0.15.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "92d117c9886f0dfe05735a1c24eb7e6e51cf0a915312350477a2fa6e70392de1" +checksum = "64a8f82fd77dfd1edeee276d5529786b3bd3f3581b0095f2e1814bd7fb21f747" dependencies = [ "bebytes", "bytes", diff --git a/Cargo.toml b/Cargo.toml index af3644d..75d7092 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -33,7 +33,7 @@ zeroize = "1.8.2" clap = { version = "4.5", features = ["derive", "env"] } comfy-table = "7.2" fjall = "3.0" -mqtt5 = "0.39" +mqtt5 = "0.45" tokio = { version = "1.49", features = ["full"] } tracing-subscriber = { version = "0.3", features = ["env-filter"] } lru = "0.16" diff --git a/crates/mqdb-agent/Cargo.toml b/crates/mqdb-agent/Cargo.toml index f59d6f1..81ec767 100644 --- a/crates/mqdb-agent/Cargo.toml +++ b/crates/mqdb-agent/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "mqdb-agent" -version = "0.8.29" +version = "0.8.30" edition.workspace = true license = "Apache-2.0" authors.workspace = true diff --git a/crates/mqdb-cli/Cargo.toml b/crates/mqdb-cli/Cargo.toml index 5675158..a30af07 100644 --- a/crates/mqdb-cli/Cargo.toml +++ b/crates/mqdb-cli/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "mqdb-cli" -version = "0.8.39" +version = "0.8.40" publish = false edition.workspace = true license = "AGPL-3.0-only" diff --git a/crates/mqdb-cluster/Cargo.toml b/crates/mqdb-cluster/Cargo.toml index 7fd1a6e..9bd645e 100644 --- a/crates/mqdb-cluster/Cargo.toml +++ b/crates/mqdb-cluster/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "mqdb-cluster" -version = "0.4.15" +version = "0.4.16" publish = false edition.workspace = true license = "AGPL-3.0-only" diff --git a/crates/mqdb-vault/Cargo.toml b/crates/mqdb-vault/Cargo.toml index 89527ef..760f179 100644 --- a/crates/mqdb-vault/Cargo.toml +++ b/crates/mqdb-vault/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "mqdb-vault" -version = "0.1.5" +version = "0.1.6" publish = false edition.workspace = true license = "AGPL-3.0-only" diff --git a/docs/design/hold-reclaim.md b/docs/design/hold-reclaim.md index db61ddd..5bd1504 100644 --- a/docs/design/hold-reclaim.md +++ b/docs/design/hold-reclaim.md @@ -136,9 +136,9 @@ 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 *after* the new connection's connect (mqtt5 registers the new connection during +its connect handshake, before `fire_connect_event`, and `fire_disconnect_event` runs on every +connection exit, a takeover included), 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. @@ -212,9 +212,9 @@ larger change than the agent's batch tweak. 1. **Per-connection vs per-user binding.** Per-connection false-reclaims a user who still has another live connection. **Rec: per-user** (`_bound_user_id` + user-level presence aggregation). -2. **mqtt-lib `user_id`.** ✅ RESOLVED & LANDED — the workspace is pinned to the published - `mqtt5 = "0.39"` (`0.39.2`), replacing the git dependency (done in its own step, ahead of - PR 3). Both `ClientConnectEvent` and `ClientDisconnectEvent` carry +2. **mqtt-lib `user_id`.** ✅ RESOLVED & LANDED — the workspace depends on the published + `mqtt5` crate (`0.39.2` when this landed, ahead of PR 3; now `0.45`), replacing the git + dependency. Both `ClientConnectEvent` and `ClientDisconnectEvent` carry `pub user_id: Option>` (mirrors `ClientPublishEvent::user_id`; the authenticated username, `None` for anonymous). Read as `event.user_id.as_deref()`. `ClientDisconnectEvent` also has `client_id: Arc` + `unexpected: bool`, so the presence payload