Skip to content

fix: MQTT compliance — message delivery, will messages, QoS 2, sessions - #37

Merged
jamesainslie merged 4 commits into
mainfrom
fix-mqtt-compliance
Mar 30, 2026
Merged

jamesainslie merged 4 commits into
mainfrom
fix-mqtt-compliance

Conversation

@jamesainslie

Copy link
Copy Markdown
Contributor

Summary

  • Issue rmq-mqtt: Message delivery not implemented — subscribers receive nothing #16: Complete MQTT broker overhaul making subscribers functional:
    • AMQP-to-MQTT message delivery via background delivery tasks polling queues
    • Will messages stored from CONNECT, published on unexpected disconnect
    • Retained messages sent on subscribe
    • QoS 2 complete flow (PubRel → PubComp)
    • Session persistence for clean_session=false
    • Queue naming collision fix (mqtt:{client_id}:{topic})
    • Error handling (no more silent let _ =)

Test plan

  • 5 new tests (delivery round-trip, will messages, QoS 2, session restore, retained)
  • 27 total MQTT tests pass
  • All workspace tests pass
  • Pre-push CI green

Closes #16

Change queue naming from `mqtt.{client_id}.{topic}` to
`mqtt:{client_id}:{topic}`. Colons cannot appear in MQTT client IDs
or topic names, eliminating ambiguous splits when client IDs or topics
contain dots.
All `let _ = self.vhost.*()` calls now log errors via tracing::warn
instead of silently swallowing them. Covers publish, declare_queue,
bind_queue, and delete_queue operations.
- Add will: Option<WillMessage> field for last-will-and-testament
- Add clean_disconnect flag to distinguish clean vs unexpected disconnect
- Add delivery_tasks: HashMap<String, JoinHandle<()>> for per-subscription
  background delivery tasks
- Cancel delivery tasks on unsubscribe and implement Drop to clean up
- Make topic_matches_filter public for use in broker
Refactor broker to use shared writer (Arc<tokio::sync::Mutex<W>>)
enabling background delivery tasks to write PUBLISH packets to clients.

Changes:
- AMQP-to-MQTT message delivery: spawn per-subscription task that polls
  the AMQP queue via shift(), converts to MQTT PUBLISH, writes to client
- Will messages: store from CONNECT, publish on unexpected disconnect
  (not clean DISCONNECT)
- Retained messages on subscribe: send matching retained messages as
  PUBLISH with retain=true when a new subscription is created
- QoS 2 completion: PubRel handler responds with PubComp, completing
  the PUBLISH->PUBREC->PUBREL->PUBCOMP flow
- Session persistence: clean_session=false stores session in broker map
  on disconnect, restores subscriptions on reconnect; clean_session=true
  removes any existing session
- Fix parking_lot::Mutex deadlock: extract lock().remove() result before
  match to prevent MutexGuard from living across match arms

Tests added:
- test_message_delivery_round_trip: publish->subscribe delivery via AMQP
- test_will_message_delivery: will published on unexpected disconnect
- test_qos2_flow: PUBLISH->PUBREC->PUBREL->PUBCOMP round-trip
- test_session_restore: clean_session=false restores, true clears
- test_retained_messages_on_subscribe: retained messages sent on subscribe
@jamesainslie
jamesainslie merged commit abb446b into main Mar 30, 2026
7 checks passed
@github-actions

Copy link
Copy Markdown

Benchmark Results

Metric main PR Delta Status
Publish 4p (msg/s) 106325 105440 -0.8% 🟢
E2E 4p/4c (msg/s) 41589 43196 3.9% 🟢

🟢 PASSED: No significant regression.

Median of 3 runs. Threshold: -10%.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

rmq-mqtt: Message delivery not implemented — subscribers receive nothing

1 participant