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
2 changes: 1 addition & 1 deletion ARCHITECTURE.md
Original file line number Diff line number Diff line change
Expand Up @@ -127,7 +127,7 @@ The client validates acknowledgment reason codes: PUBACK (QoS 1) returns `MqttEr

Each accepted connection spawns a **ClientHandler** that directly reads and writes packets, manages client session state, and handles the MQTT protocol. Outbound control packets (CONNACK, SUBACK, PUBACK, PUBREC, AUTH, DISCONNECT) are written through a choke point that honours the client's advertised Maximum Packet Size: if a packet would exceed it, the Reason String is omitted and the packet re-encoded, and only if it still does not fit is it discarded, so the broker never exceeds the client's limit (`[MQTT-3.2.2-19]`, `[MQTT-3.1.2-24]`). A CONNECT advertising a Maximum Packet Size of 0 is rejected as a Protocol Error. The **MessageRouter** performs subscription matching using MQTT-compliant topic wildcards (`+`, `#`), protects system topics (`$SYS/#` excluded from `#`), and supports shared subscriptions (`$share/group/topic`).

The **storage backend** persists sessions, retained messages, queued messages, and inflight messages. The file-based backend uses percent-encoded filenames with atomic writes and fsync; the memory backend stores everything in-process.
The **storage backend** persists sessions, retained messages, queued messages, and inflight messages. The memory backend stores everything in-process. The file-based backend stores retained, queued and inflight messages in percent-encoded files with atomic writes, and all sessions in one append-only log, `sessions/sessions.log`, with group commit: concurrent session writes are visible to later readers at once, one flush appends and fsyncs every pending write, and each write is acknowledged only after the flush that covers it and every earlier write. Each record is one line carrying a CRC-32 of its body and an explicit record type (put or remove), so a damaged record is detected rather than misread. A failed flush rejects every pending write and restores the last durable state in memory, and the bytes it may have left are truncated away before the failure is reported; if that truncation fails, the log is rewritten, and no later write is appended until the repair succeeds. The log is compacted into a new file (fsynced, renamed, directory fsynced) when it is larger than both 1 MB and twice its live size. At startup an incomplete last line from an unfinished write is discarded; a damaged complete record is skipped and replay continues, since every record carries the whole session, and the original log is first copied to `sessions.log.corrupt-<unix millis>`. Startup then rewrites the log; if that fails (full disk, read-only directory), the broker serves the replayed sessions and refuses session writes until a later write succeeds in rewriting it. Storage version 1 directories (one file per session) are migrated on open. On Windows, which has no directory fsync, appends are durable through `FlushFileBuffers`, but the rename that installs a compacted log relies on NTFS metadata journaling and may be lost by a power failure right after a compaction.

### Broker Data Flow

Expand Down
113 changes: 113 additions & 0 deletions CHANGELOG.md

Large diffs are not rendered by default.

2 changes: 1 addition & 1 deletion WASM_USAGE.md
Original file line number Diff line number Diff line change
Expand Up @@ -550,7 +550,7 @@ The JavaScript class is exported as `BrokerConfig` (Rust type: `WasmBrokerConfig
const config = new BrokerConfig();

config.maxClients = 1000; // default: 1000
config.sessionExpiryIntervalSecs = 3600; // default: 3600
config.sessionExpiryIntervalSecs = 3600; // maximum granted; default: 4294967295 (no limit)
config.maxPacketSize = 268435456; // default: 268435456 (256MB)
config.topicAliasMaximum = 65535; // default: 65535
config.retainAvailable = true; // default: true
Expand Down
28 changes: 28 additions & 0 deletions crates/mqtt5-conformance/CONFORMANCE_DIARY.md
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,34 @@

## Diary Entries

### Absent Session Expiry now means 0, and DISCONNECT can change the Session Expiry (2026-09-24)

**Trigger**: the quorum review of PR #170 and issue #171, which was folded into it. The broker stored an absent CONNECT Session Expiry Interval as "never expires", although §3.1.2.11.2 says an absent value is 0. A client that left the property out kept its session, subscriptions and queued messages forever. It also broke the Will timing from #154: with the Will bounded by session end, "never" meant the delay was always honoured instead of the Will going out at disconnect. The broker also ignored a Session Expiry Interval sent on DISCONNECT (§3.14.2.2.2), and a resumed session kept the Session Expiry of the connection that created it instead of taking the resuming CONNECT's value.

**Fix**: the Session Expiry is worked out once, from the CONNECT, when a session is created or resumed. MQTT v5: the property's value, or 0 when absent. MQTT v3.1.1 has no property, so CleanSession=1 gives 0 and CleanSession=0 keeps the session with no expiry, as before. A Session Expiry on DISCONNECT replaces the stored value before the session-end and Will logic runs. If CONNECT had 0 and DISCONNECT sends a non-zero value, the server sends DISCONNECT 0x82 and closes. The spec says such a DISCONNECT is not valid, so it is not a normal disconnection: the Will is published, and the session still ends because its expiry stays 0.

**New tests**, all failing on the tree before the fix:
- `absent_session_expiry_discards_session_at_disconnect` (`[MQTT-4.1.0-2]`): connect without the property, subscribe at QoS 1, disconnect, publish while offline, then reconnect with Clean Start 0. It expects Session Present 0 and no delivery of the offline message.
- `disconnect_session_expiry_zero_discards_session` (`[MQTT-4.1.0-2]`): CONNECT 300, DISCONNECT 0, then Session Present 0.
- `disconnect_session_expiry_extends_session` (`[MQTT-3.1.2-23]`): CONNECT 1, DISCONNECT 300, wait 2.5 s, then Session Present 1.
- `disconnect_session_expiry_after_zero_is_protocol_error` (`[MQTT-4.13.1-1]`): CONNECT 0, DISCONNECT 300, then DISCONNECT 0x82 and the connection closes.

**IDs**: neither "absent means 0" nor the DISCONNECT 0 to non-zero rule has its own normative statement in `mqtt-v5.0-statement-texts.txt`; both are prose in §3.1.2.11.2 and §3.14.2.2.2. The tests are filed under the statements they exercise: session discard after the interval (MQTT-4.1.0-2), session storage when the interval is above 0 (MQTT-3.1.2-23), and closing the connection on a Protocol Error (MQTT-4.13.1-1). The manifest text for all three matches the statement file.

**Knock-on**: tests elsewhere in the workspace had resumed sessions without setting a Session Expiry. Several of them only asserted inside `if session_present`, so after the fix they would have passed without checking anything. They now set an explicit expiry and assert Session Present.

### Delayed Will was never cancelled, and the MQTT-3.1.3-9 test could not see it (2026-09-24)

**Trigger**: issue #154. With a Will Delay Interval above zero, the broker spawned a detached task that slept for the delay and then published the Will unconditionally. A client that reconnected inside the delay still had its Will published, which violates `[MQTT-3.1.3-9]` and the "new Network Connection ... before the Will Delay Interval has elapsed" clause of `[MQTT-3.1.2-8]`.

**Why the suite passed anyway**: `will_delay_reconnect_suppresses_will` used a 5 s delay but stopped watching about 2.3 s after the drop, so the stale Will always arrived after the assertion. The test was vacuous. It now uses a 2 s delay and waits 4 s after the reconnect. On the unfixed broker it fails with the Will received. A positive control, `will_delay_elapsed_publishes_will`, runs the same setup without a reconnect and asserts that nothing arrives at 1.2 s and that the Will arrives once the delay has elapsed. It passes on both the old and the fixed broker, which shows the negative test fails for the right reason and not because Wills are never delivered.

**Fix**: the router keeps one pending delayed Will per client id, tagged with the generation of the connection that armed it. A connection arms its Will before it releases its router entry, and only if it still owns that entry. `register_session` removes any pending Will for the client id while it holds the clients write lock, so any new connection (Clean Start 0 or 1, takeover included) cancels it. When the timer fires, the task has to claim the entry by generation before it publishes. Claim and cancel both remove the entry under one mutex, so exactly one of them wins. The Will fires at min(Will Delay Interval, Session Expiry Interval), so a Session Expiry of 0 publishes at once and a shorter expiry publishes when the session ends (`[MQTT-3.1.2-8]`, §3.1.3.2.2). A published Will, and a Will deleted by DISCONNECT 0x00, is also removed from the stored session (`[MQTT-3.1.2-10]`).

**Manifest**: both tests are listed under MQTT-3.1.3-9, whose manifest text matches `mqtt-v5.0-statement-texts.txt`. The manifest entry labelled MQTT-3.1.2-8 carries the Will Retain text ("If the Will Flag is set to 0, then Will Retain MUST be set to 0"), not the Will publication statement, so neither test is cited there. That drift is left as it was.

**Broker-side coverage**: `crates/mqtt5/tests/will_delay.rs` covers resume and clean-start reconnects, no reconnect, Session Expiry 0 and 2 against longer delays, reconnect-then-drop, DISCONNECT 0x00 and 0x04, and takeover with and without a delay. Six of these fail on the unfixed broker. The other four are regression guards for behaviour that was already correct.

### Quorum review of the client fixes, and a TLA+-verified outcome model for the offline queue (2026-09-23)

**Trigger**: a five-reviewer quorum review of PR #164 before merge. Most findings came with a failing test. The worst was a regression in the wasm client: a v3.1.1 persistent session could never reconnect, because the session-lifetime check used Session Expiry (always 0 in 3.1.1) and the new strict `[MQTT-3.2.2-4]` check then rejected the broker's Session Present=1 forever. Other findings: QUIC teardown left the connection open after the client's own DISCONNECT; a publish waiting across a reconnect was sent under the old server's limits; and the offline queue dropped messages silently after `publish()` had returned success. All were fixed in the same PR.
Expand Down
8 changes: 4 additions & 4 deletions crates/mqtt5-conformance/conformance.toml
Original file line number Diff line number Diff line change
Expand Up @@ -408,7 +408,7 @@ level = "Must"
applies_to = "Both"
text = "The Client and Server MUST store the Session State after the Network Connection is closed if the Session Expiry Interval is greater than 0"
status = "Tested"
test_names = ["session_stored_when_expiry_positive"]
test_names = ["session_stored_when_expiry_positive", "disconnect_session_expiry_extends_session"]

[[sections."3.1".statements]]
id = "MQTT-3.1.2-24"
Expand Down Expand Up @@ -519,7 +519,7 @@ level = "MustNot"
applies_to = "Server"
text = "If a new Network Connection to this Session is made before the Will Delay Interval has passed, the Server MUST NOT send the Will Message"
status = "Tested"
test_names = ["will_delay_reconnect_suppresses_will"]
test_names = ["will_delay_reconnect_suppresses_will", "will_delay_elapsed_publishes_will"]

[[sections."3.1".statements]]
id = "MQTT-3.1.3-10"
Expand Down Expand Up @@ -1795,7 +1795,7 @@ level = "Must"
applies_to = "Server"
text = "The Server MUST discard the Session State when the Network Connection is closed and the Session Expiry Interval has passed"
status = "Tested"
test_names = ["session_discarded_after_expiry"]
test_names = ["session_discarded_after_expiry", "absent_session_expiry_discards_session_at_disconnect", "disconnect_session_expiry_zero_discards_session"]

# ===========================================================================
# Section 4.2 -- Network Connections
Expand Down Expand Up @@ -2359,7 +2359,7 @@ level = "Must"
applies_to = "Server"
text = "When a Server detects a Malformed Packet or Protocol Error, and a Reason Code is given in the specification, it MUST close the Network Connection"
status = "Tested"
test_names = ["malformed_packet_closes_connection"]
test_names = ["malformed_packet_closes_connection", "disconnect_session_expiry_after_zero_is_protocol_error"]

[[sections."4.13".statements]]
id = "MQTT-4.13.2-1"
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -379,7 +379,7 @@ async fn will_delay_reconnect_suppresses_will(sut: SutHandle) {
.await
.unwrap();
raw.send_raw(&RawPacketBuilder::connect_with_will_delay(
&client_id, 5, 60,
&client_id, 2, 60,
))
.await
.unwrap();
Expand All @@ -398,7 +398,7 @@ async fn will_delay_reconnect_suppresses_will(sut: SutHandle) {
.await
.expect("reconnect failed");

tokio::time::sleep(Duration::from_secs(2)).await;
tokio::time::sleep(Duration::from_secs(4)).await;

assert_eq!(
subscription.count(),
Expand All @@ -410,6 +410,56 @@ async fn will_delay_reconnect_suppresses_will(sut: SutHandle) {
subscriber.disconnect().await.expect("disconnect failed");
}

/// Positive control for `[MQTT-3.1.3-9]`: without a reconnect, the delayed
/// Will is published once the Will Delay Interval elapses and not before.
#[conformance_test(
ids = ["MQTT-3.1.3-9"],
requires = ["transport.tcp"],
)]
async fn will_delay_elapsed_publishes_will(sut: SutHandle) {
let client_id = unique_client_id("wdpub");
let will_topic = format!("will/{client_id}");

let subscriber = TestClient::connect_with_prefix(&sut, "wdpub-sub")
.await
.unwrap();
let subscription = subscriber
.subscribe(&will_topic, SubscribeOptions::default())
.await
.expect("subscribe failed");
tokio::time::sleep(Duration::from_millis(100)).await;

let mut raw = RawMqttClient::connect_tcp(sut.expect_tcp_addr())
.await
.unwrap();
raw.send_raw(&RawPacketBuilder::connect_with_will_delay(
&client_id, 2, 60,
))
.await
.unwrap();
let connack = raw.expect_connack(TIMEOUT).await;
assert!(connack.is_some(), "Must receive CONNACK");
let (_, reason) = connack.unwrap();
assert_eq!(reason, 0x00, "Connection must succeed");

drop(raw);
tokio::time::sleep(Duration::from_millis(1200)).await;

assert_eq!(
subscription.count(),
0,
"Will message must not be sent before the Will Delay Interval elapses"
);
assert!(
subscription
.wait_for_messages(1, Duration::from_secs(4))
.await,
"Will message must be sent once the Will Delay Interval elapses without a reconnect"
);

subscriber.disconnect().await.expect("disconnect failed");
}

/// `[MQTT-3.1.3-10]` The User Property is part of the Will Properties and
/// the Server MUST maintain the order of User Properties when publishing the
/// Will Message.
Expand Down Expand Up @@ -542,6 +592,79 @@ async fn session_discarded_after_expiry(sut: SutHandle) {
client2.disconnect().await.expect("disconnect failed");
}

/// `[MQTT-4.1.0-2]` An absent Session Expiry Interval means 0 (§3.1.2.11.2),
/// so the Session ends when the Network Connection closes and the Server
/// discards its Session State: subscriptions are gone and messages published
/// while the client was offline are not delivered.
#[conformance_test(
ids = ["MQTT-4.1.0-2"],
requires = ["transport.tcp", "max_qos>=1"],
)]
async fn absent_session_expiry_discards_session_at_disconnect(sut: SutHandle) {
let client_id = unique_client_id("sess-absent");
let topic = format!("sess/{client_id}");

let opts = ConnectOptions::new(&client_id).with_clean_start(true);
let client1 = TestClient::connect_with_options(&sut, opts)
.await
.expect("first connect failed");
client1
.subscribe(
&topic,
SubscribeOptions {
qos: mqtt5_protocol::QoS::AtLeastOnce,
..SubscribeOptions::default()
},
)
.await
.expect("subscribe failed");
client1.disconnect().await.expect("disconnect failed");
tokio::time::sleep(Duration::from_millis(200)).await;

let publisher = TestClient::connect_with_prefix(&sut, "sess-absent-pub")
.await
.unwrap();
publisher
.publish_with_options(
&topic,
b"offline",
mqtt5_protocol::types::PublishOptions {
qos: mqtt5_protocol::QoS::AtLeastOnce,
..Default::default()
},
)
.await
.expect("publish failed");

let mut raw = RawMqttClient::connect_tcp(sut.expect_tcp_addr())
.await
.unwrap();
let reconnect = ConnectOptions::new(&client_id)
.with_clean_start(false)
.with_session_expiry_interval(300);
let mut buf = bytes::BytesMut::new();
mqtt5_protocol::packet::MqttPacket::encode(
&mqtt5_protocol::packet::connect::ConnectPacket::new(reconnect),
&mut buf,
)
.expect("CONNECT encodes");
raw.send_raw(&buf).await.unwrap();
let connack = raw
.expect_connack_packet(TIMEOUT)
.await
.expect("Must receive CONNACK");
assert!(
!connack.session_present,
"[MQTT-4.1.0-2] an absent Session Expiry Interval is 0, so no session survives the disconnect"
);
assert!(
raw.expect_publish(Duration::from_secs(1)).await.is_none(),
"[MQTT-4.1.0-2] a message published while offline must not be delivered to a discarded session"
);

publisher.disconnect().await.expect("disconnect failed");
}

/// `[MQTT-2.1.3-1]` Where a flag bit is marked as Reserved, it is reserved
/// for future use and MUST be set to the value listed.
#[conformance_test(
Expand Down
Loading
Loading