Skip to content

feat(mqtt): dispatch embedded inbound publishes from the session loop (design 054 §4.7) - #286

Merged
lxsaah merged 1 commit into
feat/054-connector-boundaryfrom
feat/054-s09-inbound-dispatch
Oct 4, 2026
Merged

lxsaah merged 1 commit into
feat/054-connector-boundaryfrom
feat/054-s09-inbound-dispatch

Conversation

@lxsaah

@lxsaah lxsaah commented Oct 4, 2026 •

Copy link
Copy Markdown
Contributor

Stage 9 of the 054 implementation plan. The embedded MQTT session delivers each inbound publish straight into its records from drain_packets, while it still holds the decoded packet. The event channel, MqttSource and the pump_source task are gone. Outbound still goes through the action channel and pump_sink until stage 10.

Change

  • embedded/mod.rs:
    • build_plain and build_tls build InboundDispatch::new(db, "mqtt", &MqttGrammar) instead of inbound_router, and subscribe its subscriptions().
    • The dispatcher moves into the session task through setup_manager/setup_tls_manager and run_sessions/run_tls.
    • Removed: AimdbMqttEvent with its conversion, EventChannel, MqttSource, and the pump_source call.
    • Breaking: aimdb_mqtt_connector::embedded::AimdbMqttEvent was public.
  • manager.rs: the EventChannel alias and the crate-private FromApplicationMessage trait are removed.
  • session_loop.rs:
    • run_session, client_loop and drain_packets take &InboundDispatch, named dispatch because the raw-chunk channel is already called inbound.
    • A new deliver calls dispatch.dispatch(publish.topic_name(), publish.payload()) for Publish and PublishAndPuback. It keeps the empty-topic check, the liveness-only events and the DISCONNECT error. Received is gone.
    • feat(mqtt): encode embedded packets into one bbqueue write ring (design 054 §4.7) #283 review nit: PUBACK room is now waited for only when the head packet is a PUBLISH above QoS 0, read from its first byte (PacketReader::head_needs_ack), not before every packet.
  • connector.rs: the Embedded rustdoc says what QoS 1 inbound means now: the PUBACK goes out before delivery, so it means "reached AimDB". A record whose buffer is full drops the message, and the broker does not resend.
  • setup_tls_manager takes 8 arguments now and gets #[allow(clippy::too_many_arguments)], as the session loop's long signatures already do. Stage 10 replaces the action channel there.

Net: +170 / −171 in aimdb-mqtt-connector, tests included.

Tests

  • New the_embedded_backend_dispatches_inbound_on_its_session_task (tokio_broker): with an inbound link and no outbound ones, build() returns one future. Against the previous code the same test sees 2, which I checked by running it on the old src/.
  • New only_a_publish_above_qos_0_needs_an_ack (packet_reader).
  • The session-size test builds an InboundDispatch over an empty database (NoopRuntimeOps) and still passes the ≤ 8,192-byte bound.

Existing suites, unchanged and passing:

Suite Passed
session_loop (QoS 1 delivery, inbound flood, ping cadence) 6
tokio_broker (reconnect/resubscribe, round trip, Maximum Packet Size, new count test) 5
embassy_broker 1
tls_session 3
tls_broker 3
backend_parity 8
write_ring_proofs 1
--features std 40
embedded-tls lib 43

Verification

  • Every aimdb-mqtt-connector clippy leg (host, test targets, thumbv7em including embassy-tls + defmt) and every doc leg (-D warnings) from the Makefile: clean. cargo fmt --all --check: clean.
  • embassy-mqtt-connector-demo and weather-station-gamma build for thumbv8m.main-none-eabihf.
  • CI does not run on PRs into feat/054-connector-boundary.

🤖 Generated with Claude Code

… (design 054 §4.7)

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
@lxsaah
lxsaah merged commit 588f682 into feat/054-connector-boundary Oct 4, 2026
@lxsaah
lxsaah deleted the feat/054-s09-inbound-dispatch branch October 4, 2026 19:40
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.

1 participant