feat(mqtt): pull embedded outbound publishes from OutboundRoutes (design 054 §4.7) - #287
Merged
Merged
Conversation
…ign 054 §4.7) Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Stage 10 of the 054 implementation plan. The embedded MQTT session now pulls outbound messages from
OutboundRouteswhen it can send. The action channel,MqttSinkand the per-routepump_sinktasks are gone, so the embedded connector contributes exactly one future: its session task. This stage also lands four follow-ups from the #283 review: counting skipped publishes (task 1),with_write_buffer(task 2), size checks at build (task 6), and the broker's CONNACK limit (task 7).Change
Core (
outbound/routes.rs):RouteStats::rejectedandOutboundRoutes::reject(id): a connector reports a message it took but could not send.sentnow reads "staged and handed to the connector, including any it then rejected".New
publish_opts.rs:PublishOpts::parse(&RouteInfo)readsqos(0/1/2, default 1) andretain(true/false, default false).embeddedfor now. Stage 11 shares it with the native backend.connector.rs:with_write_buffer(bytes)onMqttConnector<Embedded<_>>andMqttConnector<EmbeddedTls<_>>, default 4,096. The size travels throughbuild_*,setup_*andrun_sessions/run_tlsto the ring.embedded/mod.rs— newprepare_outbound, run inbuild():OutboundRoutes::new(db, "mqtt")and parses each route's options. Routes asking forqos=2get a warning once each (log_warn!anddefmt); this client sends them at QoS 1, as before.AimdbMqttAction(public, so breaking),ActionChannel,CHANNEL_SIZE,MqttSink,collect_pumps, thepump_sinkcall,map_qos,opt_u8/opt_bool, andwarn_unsupported_qos(with its use ofcollect_outbound_routes).ActionChannelalias is removed too.session_loop.rs:roomwaker). It then callspoll_stage.Ready(None)is latched inoutbound_done, and a losing arm takes nothing.publish_staged: after theselect, it takes the staged message, builds the PUBLISH with the route'sqos/retain, then writes it into the ring and updates the state.outbound.reject(id)when the frame doesn't fit the ring, or exceeds the broker's Maximum Packet Size from its CONNACK. The CONNACK property is read indrain_packetsand kept per session.connect_packetis shared by the session and the build check.subscribe_lenandpublish_frame_lensize packets at build.write_ring::fits_ringis the same rule without a ring.performis gone. At most once still holds: a message is taken from its record buffer before it is written.Diff: +844 / −378 overall, including core, the new parser and tests. In
aimdb-mqtt-connector/srcexcluding the parser it's +470 / −363.Tests
New:
after_a_stall_a_single_latest_record_sends_only_its_newest_value(session_loop): during a held PUBACK, values 1–9 are produced 20 ms apart, and the broker then receives["0", "9"].src/stashed), the same test sees["0", "1", "2", …].an_invalid_qos_or_retain_fails_the_build:qos=3,qos=abcandretain=yeseach fail the build, naming the route.a_route_too_large_for_the_write_ring_fails_the_build,a_connect_too_large_…anda_subscribe_too_large_…: each fails at the default ring and builds withwith_write_buffer(8192).publish_stagedon a real staged message:sent 1, rejected 1). This proof flipped.publish_frame_len_matches_what_the_client_state_encodes.PublishOptsunit tests.a_rejected_message_is_counted_beside_sent(Tokio adapter) andrejectwith an unknown id.Updated:
a_qos2_route_is_visible_to_the_build_time_scanreadsOutboundRoutes::routes().proof_an_oversize_publish_is_dropped_silentlyis renamedan_oversize_publish_is_skipped_and_the_session_stays_up.an_idle_session_wakes_at_the_ping_cadence(inbound-only connector) still passes, soReady(None)doesn't spin.embedded-tlslibsession_looptokio_brokerembassy_brokertls_sessiontls_brokerbackend_paritywrite_ring_proofs--features stdaimdb-coreaimdb-tokio-adapterNot done here
embassy_brokertoo. That test runs its own executor and fake broker over embassy-net, so a held PUBACK would need a second scripted broker. The question concerns the buffers, not MQTT, so I propose a follow-up test ofOutboundRoutesdirectly over the Embassy adapter's buffers, likeoutage_semantics_per_buffer_typefor Tokio.RouteStatsfrom outside the connector.OutboundRouteslives inside the session task, so only the connector can callstats(id). The new count is asserted at unit level. Users need a read path, for exampleAimDb::outbound_stats(scheme)backed by shared counters. That's a core API decision, so I've left it open.Verification
make clippy(whole workspace, including embedded targets), every core and MQTT doc leg (-D warnings) andcargo fmt --all --check: clean.embassy-mqtt-connector-demoandweather-station-gammabuild forthumbv8m.main-none-eabihf.feat/054-connector-boundary.🤖 Generated with Claude Code