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
22 changes: 22 additions & 0 deletions .github/workflows/rust.yml
Original file line number Diff line number Diff line change
Expand Up @@ -167,6 +167,28 @@ jobs:
- name: Check WASM compilation
run: cargo check --target wasm32-unknown-unknown --all-features -p mqtt5-wasm

- name: Install Node.js
uses: actions/setup-node@v4
with:
node-version: 22

- name: Install wasm-bindgen-test-runner
run: |
version=$(awk '/^name = "wasm-bindgen"$/ { getline; gsub(/"/, "", $3); print $3; exit }' Cargo.lock)
test -n "$version"
installed=$(wasm-bindgen-test-runner --version 2>/dev/null || true)
if [ "$installed" != "wasm-bindgen-test-runner $version" ]; then
cargo install wasm-bindgen-cli --version "$version" --locked
fi

- name: Run WASM tests
env:
CARGO_TARGET_WASM32_UNKNOWN_UNKNOWN_RUNNER: wasm-bindgen-test-runner
WASM_BINDGEN_TEST_TIMEOUT: "120"
run: |
cargo test -p mqtt5-wasm --target wasm32-unknown-unknown
cargo test -p mqtt5-wasm --target wasm32-unknown-unknown --features broker

- name: Build WASM package
run: |
cd crates/mqtt5-wasm
Expand Down
120 changes: 120 additions & 0 deletions CHANGELOG.md

Large diffs are not rendered by default.

30 changes: 26 additions & 4 deletions Makefile.toml
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,7 @@ echo ""
echo "🌐 WASM BUILD (mqtt5-wasm crate)"
echo " cargo make wasm-check Check WASM compilation"
echo " cargo make wasm-clippy Run clippy for WASM target"
echo " cargo make wasm-test Run WASM crate tests (native)"
echo " cargo make wasm-test Run WASM crate tests under Node"
echo " cargo make wasm-build Build WASM package with wasm-pack"
echo " cargo make wasm-verify Run all WASM checks"
echo " cargo make wasm-examples Build WASM and copy to examples"
Expand Down Expand Up @@ -252,10 +252,32 @@ command = "cargo"
args = ["clippy", "--target", "wasm32-unknown-unknown", "-p", "mqtt5-wasm", "--features", "broker", "--", "-D", "warnings", "-W", "clippy::pedantic"]
description = "Run clippy for WASM target (strict, pedantic)"

[tasks.wasm-test-runner]
script = '''
#!/usr/bin/env bash
set -euo pipefail
version=$(awk '/^name = "wasm-bindgen"$/ { getline; gsub(/"/, "", $3); print $3; exit }' Cargo.lock)
if [ -z "$version" ]; then
echo "wasm-bindgen not found in Cargo.lock" >&2
exit 1
fi
installed=$(wasm-bindgen-test-runner --version 2>/dev/null || true)
if [ "$installed" != "wasm-bindgen-test-runner $version" ]; then
cargo install wasm-bindgen-cli --version "$version" --locked
fi
'''
description = "Install the wasm-bindgen-test-runner matching the Cargo.lock wasm-bindgen version"

[tasks.wasm-test]
command = "cargo"
args = ["test", "-p", "mqtt5-wasm"]
description = "Run WASM crate tests (native)"
dependencies = ["wasm-test-runner"]
env = { CARGO_TARGET_WASM32_UNKNOWN_UNKNOWN_RUNNER = "wasm-bindgen-test-runner", WASM_BINDGEN_TEST_TIMEOUT = "120" }
script = '''
#!/usr/bin/env bash
set -euo pipefail
cargo test -p mqtt5-wasm --target wasm32-unknown-unknown
cargo test -p mqtt5-wasm --target wasm32-unknown-unknown --features broker
'''
description = "Run WASM crate tests under Node with wasm-bindgen-test-runner"

[tasks.wasm-verify]
dependencies = ["wasm-check", "wasm-clippy", "wasm-test"]
Expand Down
40 changes: 37 additions & 3 deletions WASM_USAGE.md
Original file line number Diff line number Diff line change
Expand Up @@ -262,15 +262,23 @@ client.destroy();
```javascript
await client.publish(topic, payloadBytes);

await client.publishWithOptions(topic, payloadBytes, publishOptions);
const qosUsed = await client.publishWithOptions(topic, payloadBytes, publishOptions);
// resolves with the QoS the message was sent at (0, 1 or 2)

const packetId = await client.publishQos1(topic, payloadBytes, callback);
// callback(reasonCode) called when PUBACK received
// callback(reasonCode, qosUsed) called when PUBACK received

const packetId = await client.publishQos2(topic, payloadBytes, callback);
// callback(reasonCode) called when PUBCOMP received
// callback(reasonCode, qosUsed) called when PUBCOMP received
```

A publish never goes out above the server's Maximum QoS (from CONNACK). A higher requested QoS is downgraded to the server maximum and the QoS actually used is reported: `publishWithOptions` resolves with it, and the `publishQos1`/`publishQos2` callback receives it as the second argument. When the server maximum is 0, the message is sent at QoS 0, `publishQos1`/`publishQos2` resolve with packet identifier `0` (none) and the callback is called immediately with `(0, 0)`. Retain Available and Maximum Packet Size are not adjusted: a publish that violates them is rejected before anything is sent.

`disconnect()` ends this instance's use of the connection and never leaves an acknowledgement pending:

- If the session ends with the connection (see [Session Lifetime](#session-lifetime)), the session state is discarded: pending `publishWithOptions` promises reject with `Publish not acknowledged: indeterminate: session ended with the connection; may have been delivered` and `publishQos1`/`publishQos2` callbacks receive that `indeterminate: ...` string.
- If the session outlives the connection, the session state is kept. Pending `publishWithOptions` promises reject with `Publish not acknowledged: disconnected; message remains in session and will be resent on resume by this client instance`, and `publishQos1`/`publishQos2` callbacks receive that `disconnected; ...` string. A later `connect*()` with `cleanStart = false` on the same `MqttClient` instance resumes the session (when the server reports Session Present = 1) and resends those messages; their acknowledgements are then handled silently.

#### Subscribing

```javascript
Expand Down Expand Up @@ -572,6 +580,8 @@ const opts = new ConnectOptions();

opts.keepAlive = 60; // default: 60 seconds
opts.cleanStart = true; // default: true
opts.resumeExistingSession = false; // default: false; with cleanStart = false, accept
// Session Present = 1 without local session state
opts.username = 'alice'; // default: null
opts.set_password(encoder.encode('pw')); // accepts Uint8Array
opts.protocolVersion = 5; // 4 (v3.1.1) or 5 (v5.0), default: 5
Expand Down Expand Up @@ -605,6 +615,30 @@ opts.clearCodecRegistry();
await client.connectWithOptions('ws://broker:8000/mqtt', opts);
```

#### Session Lifetime

The client keeps its session state (unacknowledged QoS 1/2 publishes and QoS 2 releases) for as long as the session outlives the network connection:

- MQTT v5.0: while the Session Expiry Interval is greater than 0. The value the server returns in CONNACK replaces the requested `sessionExpiryInterval`.
- MQTT v3.1.1 (`protocolVersion = 4`): while `cleanStart = false` (CleanSession = 0). With `cleanStart = true` the session ends with the connection.

When the session ends with the connection, pending publishes are settled on connection loss or `disconnect()` with `indeterminate: session ended with the connection; may have been delivered`. A `connect*()` with `cleanStart = true` discards any session state the instance still holds before sending CONNECT (MQTT-3.1.2-4): nothing is resent, pending publishes settle with `indeterminate: session discarded by clean start; may have been delivered`, and quarantined packet identifiers are released. When it outlives the connection, a connection loss leaves them pending and they are resent when the session resumes (automatic reconnect, or a later `connect*()` with `cleanStart = false` on the same `MqttClient` instance). Automatic reconnects send `cleanStart = false` only while the client holds session state (or `resumeExistingSession` is set), so a v3.1.1 `cleanStart = true` client never turns into a persistent session.

On the next CONNACK after a `cleanStart = false` CONNECT on the same `MqttClient` instance, the unacknowledged messages are handled according to Session Present and the limits of the new connection (Retain Available, Maximum Packet Size, Maximum QoS):

| CONNACK | Unacknowledged message | Outcome |
|---|---|---|
| Session Present = 1 | QoS 1/2 PUBLISH within the new limits | Resent with the same packet identifier, in the original order (DUP = 1 when it was already sent on this session) |
| Session Present = 1 | QoS 1/2 PUBLISH outside the new limits | Not resent and removed from the session; settled with `indeterminate: may have been delivered; not resent because the new connection's limits do not allow it (<reason>)`. A QoS 2 packet identifier abandoned this way is not reused for new messages until a connection reports Session Present = 0 |
| Session Present = 1 | QoS 2 awaiting PUBCOMP | PUBREL resent |
| Session Present = 0 (server lost the session) | QoS 1 PUBLISH within the new limits | Sent again as a new message (DUP = 0), in the original order; its promise or callback settles normally when acknowledged |
| Session Present = 0 | QoS 1 PUBLISH outside the new limits | Not sent; settled with `indeterminate: session lost; may have been delivered; not resent because the new connection's limits do not allow it (<reason>)` |
| Session Present = 0 | QoS 2 PUBLISH or PUBREL | Not sent; settled with `indeterminate: session lost; may have been delivered` |

`publishWithOptions` promises reject with `Publish not acknowledged: ` followed by that text, and `publishQos1`/`publishQos2` callbacks receive the text as their only argument. Messages whose promise or callback was already settled by `disconnect()` are handled the same way without a further notification.

A CONNACK with Session Present = 1 is rejected (the client sends DISCONNECT 0x82 on v5.0 and closes the connection) unless the client holds session state or `resumeExistingSession` is set.

### PublishOptions API

The JavaScript class is exported as `PublishOptions` (Rust type: `WasmPublishOptions`).
Expand Down
24 changes: 24 additions & 0 deletions crates/mqtt5-conformance/CONFORMANCE_DIARY.md
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,30 @@

## Diary Entries

### 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.

**The offline-queue question was settled by three independent TLA+ models, not by argument.** All three found the same defects in the current behaviour: silent loss at flush, silent loss on Session Present=0, and replay resending packets the new CONNACK forbids. All three arrived at the same design: per-publish outcomes (Delivered / Rejected / Indeterminate) that are never keyed by packet id, checks at enqueue, reject-and-continue at flush, re-checks on replay, and quarantine of abandoned QoS 2 ids. Holding a message fails liveness, skipping ahead breaks ordering, and transforming it silently loses the caller's intent. One model showed that reporting failures by packet id misattributes them after the id is reused (ABA). The user chose downgrade-and-report for Maximum QoS, and requeueing QoS 1 when the server loses the session. A Clean Start=1 connect discards unacked outbound state, as `[MQTT-3.1.2-4]` requires. The consolidated spec is in `specs/tla/offline-queue/`. New tests are in `crates/mqtt5/tests/conf_client_offline_queue.rs`: 11 fake-broker tests, all failing on the prior tree.

**Second quorum round on the outcome code**: three re-reviewers checked the implementation against the models. The protocol behaviour held. Two defects were in how outcomes are settled: a connection loss never closed the send quota, so a parked flush task kept handles alive forever, and a reader aborted between releasing session state and settling the outcome left a delivered publish unsettled. A live publish returned `Err` on connection loss although its message stayed in the session and was resent. It now returns a handle like a queued one. A QoS 2 message that received PUBREC Success before the session was lost is reported delivered. Each fix has a test that fails on the previous code.

**Tooling lesson**: tla-mcp 0.9.4 passed a `~>` negative control vacuously and ignored missing fairness. Liveness results are trusted only as `[]<>` properties with negative controls that fail, backed by an ENABLED-based progress invariant.

### Client-side audit: the suite only ever tested brokers, and our own clients failed ~40 MUSTs (2026-09-23)

**Trigger**: checking a third-party client's claim of full MQTT v5 conformance. This suite's SUT is always a broker, so it could not answer the question. The 149 statements in `conformance.toml` with `applies_to = "Client"` or `"Both"` were audited instead with raw-byte fake-broker tests that drive the real client and record what it puts on the wire. After the third-party client had been tested, the same audit ran against our own `MqttClient` and the `mqtt5-wasm` client.

**Result**: our native client failed about 40 client MUST statements. The third-party client failed 7. Resend on session resume did not exist. Packet identifiers were reused while in flight. There was no topic or filter validation. Topic Alias Maximum and Retain Available were never enforced. The offline queue bypassed flow control. Protocol errors left the socket half-open. WebSocket reads assumed one packet per frame. The wasm client had most of the same defects and also never sent PUBACK. All of them are fixed in mqtt5 0.41.0 and mqtt5-wasm 2.0.0.

**Where the tests live**: `crates/mqtt5/tests/conf_client_{a,b,c,d}.rs` (native) and `crates/mqtt5-wasm/tests/conformance_client.rs` (wasm, MessagePort fake broker under Node). They are not yet registered in this crate's manifest or runner. A client-side SUT mode for this suite is the natural next step.

**Manifest drift bit the audit**: statement IDs in `conformance.toml` were used to label findings, and several were wrong. For example, the manifest files "no session state + Session Present=1 → close" under 3.2.2-5, but it is 3.2.2-4. All client-test names use IDs from `mqtt-v5.0-statement-texts.txt`. `known-text-drift.txt` is real debt with consequences outside this crate.

**Decision recorded**: `[MQTT-3.2.2-4]` is enforced strictly by default. A fresh client can still resume a broker-held session through an explicit `ConnectOptions::resume_existing_session` opt-in, which the deferred-ack crash-recovery pattern needs. This crate's in-process test client sets the opt-in, because it checks `session_present` as an observer of the broker.

**Lesson**: a conformance suite that only tests one side of the protocol says nothing about the other. Run any check we would point at someone else's implementation against our own first.

### External-broker ack timeouts were a lost-wakeup race in the test client, not broker timing (2026-09-07)

**Trigger**: issue #146. `deferred_qos2_zero_quota_still_serves_control_plane [MQTT-4.9.0-3]` failed
Expand Down
26 changes: 15 additions & 11 deletions crates/mqtt5-conformance/src/test_client/inprocess.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@ use super::{MessageQueue, ReceivedMessage, Subscription, TestClientError};
use crate::sut::SutHandle;
use mqtt5::MqttClient;
use mqtt5_protocol::types::{ConnectOptions, PublishOptions, SubscribeOptions};
use std::sync::{Arc, Mutex};
use std::sync::{Arc, Mutex, PoisonError};

/// In-process backing for [`crate::test_client::TestClient`].
pub struct InProcessTestClient {
Expand Down Expand Up @@ -35,6 +35,7 @@ impl InProcessTestClient {

let wrapper = mqtt5::ConnectOptions {
protocol_options: options.clone(),
resume_existing_session: true,
..mqtt5::ConnectOptions::default()
};
let client = MqttClient::with_options(wrapper.clone());
Expand All @@ -61,28 +62,31 @@ impl InProcessTestClient {
/// Publishes with the given [`PublishOptions`].
///
/// # Errors
/// Returns an error if the broker rejects the publish or the client
/// is disconnected.
/// Returns an error if the broker rejects the publish, the client is
/// disconnected, or the publish was not acknowledged before the connection
/// ended or the acknowledgement wait elapsed.
pub async fn publish_with_options(
&self,
topic: &str,
payload: &[u8],
options: PublishOptions,
) -> Result<(), TestClientError> {
self.client
match self
.client
.publish_with_options(topic, payload.to_vec(), options)
.await?;
Ok(())
.await?
{
mqtt5::PublishResult::Sent(_) => Ok(()),
mqtt5::PublishResult::Queued(_) => {
Err(TestClientError::Timeout("publish acknowledgement"))
}
}
}

/// Subscribes to `filter` and returns a [`Subscription`] handle.
///
/// # Errors
/// Returns an error if the broker rejects the subscription.
///
/// # Panics
/// Panics from the delivery callback if the internal mutex has been
/// poisoned.
pub async fn subscribe(
&self,
filter: &str,
Expand All @@ -95,7 +99,7 @@ impl InProcessTestClient {
.subscribe_with_options(filter, options, move |msg| {
messages_cb
.lock()
.unwrap()
.unwrap_or_else(PoisonError::into_inner)
.push(ReceivedMessage::from_message(msg));
})
.await?;
Expand Down
2 changes: 1 addition & 1 deletion crates/mqtt5-protocol/Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[package]
name = "mqtt5-protocol"
version = "0.15.1"
version = "0.15.2"
edition.workspace = true
rust-version.workspace = true
authors.workspace = true
Expand Down
Loading
Loading