From 909fbb7a2749e77b4b1fdedbf71a0bc22c959e3f Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 2 Oct 2026 08:20:25 +0000 Subject: [PATCH 1/5] docs(054): drop backward compatibility, two-ring embedded MQTT outbound - One mechanism per direction: InboundDispatch inbound, RoutePublisher outbound. Remove Source, pump_source, public inbound_router, Connector::publish, TopicProvider and the compatibility adapters; with_topic_provider becomes with_topic_writer. - Embedded MQTT outbound: route publishers copy raw records into a record ring; the session loop stays the only encoder and writes into a per-connection bip-buffer write ring with a control-packet reserve. - Pump never ends a route on a publish error; document QoS 1 ack-then-drop and ingest latency on both MQTT backends. - Single-branch migration with per-crate table; CI gate lands first. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01EyVFXVWAec2aouRceoGdpE --- .../054-zero-alloc-connector-boundary.md | 378 +++++++++++------- 1 file changed, 243 insertions(+), 135 deletions(-) diff --git a/docs/design/054-zero-alloc-connector-boundary.md b/docs/design/054-zero-alloc-connector-boundary.md index ab1cd02f..834861fa 100644 --- a/docs/design/054-zero-alloc-connector-boundary.md +++ b/docs/design/054-zero-alloc-connector-boundary.md @@ -1,15 +1,20 @@ # 054 — Zero-allocation connector boundary -**Status:** 📝 Proposed — revised 2026-09-30 against the tree after 055 -landed (router grammar, keys) and after reading both MQTT backends' data -paths. +**Status:** 📝 Proposed — revised 2026-10-02: no backward compatibility +(adapters, `Source` and `TopicProvider` removed; one migration), embedded +MQTT outbound reworked so the session loop stays the only encoder (§4.6). +Previous revision 2026-09-30, after 055 landed. **Scope:** the per-message path between `aimdb-core` and connectors, in both directions: `Source`/`pump_source`, `SerializedReader`, `Connector::publish`, -`TopicProvider`, and the MQTT connector as the first adopter. SPI break in -`aimdb-core` (next major). User-facing APIs (`link_from`, `link_to`, +`TopicProvider`, and every in-tree connector. Breaking change to the +connector SPI and to one user-facing builder method: `with_topic_provider` +is replaced by `with_topic_writer` (§4.3). `link_from`, `link_to`, `with_deserializer`, `with_match_deserializer`, `with_serializer`, -`Reader::recv`) do not change. +`with_serializer_into` and `Reader::recv` do not change. + +**Compatibility:** none kept. There are no adapters for the old traits; all +in-tree connectors migrate in one change (§6). **Builds on:** [037 — Zero-allocation consume path](./037-zero-alloc-consume-path.md), which made buffers and the consume path allocation-free and left the @@ -55,12 +60,17 @@ at connector level): | Backend | Inbound | Outbound | |---|---|---| -| Native (`rumqttc`) | topic `publish.topic.clone()` + payload `Arc::from(publish.payload)` | `destination.to_string()` + `payload.to_vec()`, required by `AsyncClient::publish`'s owned API | +| Native (`rumqttc`) | topic `publish.topic.clone()` + payload `Arc::from(publish.payload.as_ref())` | `destination.to_string()` + `payload.to_vec()`, required by `AsyncClient::publish`'s owned API | | Embedded (`mountain-mqtt` codec, own session loop) | topic `topic_name.to_string()` + payload `Payload::from` | `destination.to_string()` + `payload.to_vec()` into `AimdbMqttAction::Publish`; then `session_loop::encode` allocates a `Vec` per packet for the `Channel, 4>` write queue | -The native topic clone is avoidable today (the `Publish` is owned; move the -`String`). The embedded write queue is a third per-message allocation that -the interface change alone does not remove (§4.6). +Both native inbound copies disappear with §4.1. The embedded write queue is +a third per-message allocation that the interface change alone does not +remove (§4.6). + +Not every connector goes through `Source` today: the AimX session client +(TCP, UDS, serial) and both WebSocket builders already call +`Router::route` with borrows from their own frames. MQTT (both backends), +KNX and the Embassy adapter's glue use `Source` + `pump_source`. **Values with heap data.** Each reader receives its own clone of `T` (`T: Clone` delivery). A value with one `String` field costs one allocation @@ -82,12 +92,18 @@ gates the numbers yet. - G2. The MQTT connector adds no copies or allocations beyond what its client library's API requires, documented per backend. - G3. `b0_alloc_connector` runs in CI and gates the result. +- G4. One mechanism per direction: a single inbound entry point and a single + outbound publisher trait, with no compatibility layers. **Non-goals** - Topic grammars, wildcards and keys; 055 owns them and this design only carries them through the new inbound entry point. - Making `T` delivery cheaper for heap-containing types (§4.5 is guidance). +- Removing owned serializers (A6). `with_serializer` stays: it is used far + more than `with_serializer_into` in-tree (63 call sites against 7), and + requiring a declared capacity everywhere belongs in its own design. Users + who need zero allocations choose `with_serializer_into`. - Allocations inside third-party client libraries. - Remote-access JSON paths (tracked by `b0_alloc_remote_json`). @@ -98,25 +114,32 @@ and `poll` do not. Per-message boxed futures exist only to satisfy trait signatures, so the signatures change to poll form. Owned values that exist only to cross the interface become borrows. +Because no compatibility is kept, each direction gets exactly one shape: +inbound connectors push borrowed messages into `InboundDispatch`; outbound +connectors implement `RoutePublisher`. `Source`, `pump_source`, +`Connector::publish`, `TopicProvider` and the async `SerializedReader` +methods are removed, not adapted. + ## 4. Design ### 4.1 Inbound: connectors push borrowed messages Connectors hold each received message by reference: inside `rumqttc`'s `Event::Incoming(Publish)`, inside the embedded session loop's decoded -`ApplicationMessage`. `Router::route(&str, &[u8], ctx)` already takes -borrows and, after 055, matches through the connector's grammar without -allocating. The interface in between is what forces owned copies (A1, A2). +`ApplicationMessage`, inside a decoded AimX or WebSocket frame. +`Router::route(&str, &[u8], ctx)` already takes borrows and, after 055, +matches through the connector's grammar without allocating. The `Source` +interface in between is what forces owned copies (A1, A2). -Core exposes the router directly to connectors: +Core exposes the router to connectors through one type: ```rust /// Built once per connector from the inbound links of its scheme. pub struct InboundDispatch { /* Router + RuntimeContext */ } impl InboundDispatch { - /// Same validation as `AimDb::inbound_router` (055 §5.3): every link the - /// grammar or its key cannot compile is reported. + /// Same validation as today's `AimDb::inbound_router` (055 §5.3): every + /// link the grammar or its key cannot compile is reported. pub fn new( db: &AimDb, scheme: &str, @@ -129,32 +152,46 @@ impl InboundDispatch { /// Filters to subscribe at the transport (`Router::subscriptions`). /// Called once at connect time; allocates, which is fine there. pub fn subscriptions(&self) -> Vec>; + /// Number of inbound routes, for logging (`Router::route_count`). + pub fn route_count(&self) -> usize; } ``` - Connectors call `dispatch` wherever they hold the borrow. No future, no copy, no channel. `TopicMatch<'a>` already borrows the topic, so match-aware deserializers are unchanged (055 §7). -- `InboundDispatch` replaces `inbound_router` + `pump_source` for migrated - connectors. -- `Source` stays for transports that naturally produce owned frames, but in - poll form (037 pattern), which removes A1: - `fn poll_next(&mut self, cx: &mut Context<'_>) -> Poll>`. - A2 remains for those transports by construction. -- Running ingest on the connector's own task is already today's behaviour - (`pump_source` runs it inline on the reader future); it stays synchronous - and non-blocking. - -**Behaviour change (embedded).** Today the session loop parks on -`events.send(event).await` when `pump_source` falls behind, so the -`EventChannel` (`CHANNEL_SIZE` messages) absorbs bursts and throttles the -socket. With `dispatch`, a burst goes straight to the record buffers, which -drop on overflow as they already do for a slow consumer. Records that need -burst tolerance size their buffer for it. +- `InboundDispatch` is the only inbound entry point. `Source`, + `pump_source` and `AimDb::inbound_router` are removed; `Router` becomes + crate-private. A transport that produces owned frames calls `dispatch` + with borrows of its own frame, so A2 no longer exists as a category. +- For the AimX session client and the WebSocket connectors this is a rename: + they already call `router.route(topic, payload, &ctx)`. +- `InboundDispatch` is cheap to share (`Arc` inside or `Clone`), since the + WebSocket server dispatches from one task per client. + +**Behaviour changes.** + +- *Embedded MQTT: burst absorption.* Today the session loop parks on + `events.send(event).await` when `pump_source` falls behind, so the + `EventChannel` (`CHANNEL_SIZE` messages) absorbs bursts and throttles the + socket. With `dispatch`, a burst goes straight to the record buffers, + which drop on overflow as they already do for a slow consumer. Records + that need burst tolerance size their buffer for it. +- *QoS 1 inbound is acknowledge-then-maybe-drop.* `drain_packets` queues + the PUBACK before the message is delivered, and the broker does not + resend an acknowledged message. Today an acknowledged message has at least + reached the event channel; with `dispatch` it can be dropped at once by a + full record buffer. QoS 1 at the transport therefore means "reached + AimDB", not "reached every record". This is documented on the MQTT + connector; moving the PUBACK after `dispatch` does not change it, because + `dispatch` does not report per-record drops. +- *Ingest runs on the transport task.* For embedded MQTT that is the session + loop; for native MQTT it is the task polling `rumqttc`'s event loop. A + slow user deserializer delays keep-alive and ack handling in both (§8). ### 4.2 Outbound: poll-based reader -`SerializedReader::recv_into` becomes a poll method over the buffer reader's +`SerializedReader` has one method, in poll form, over the buffer reader's existing `poll_recv` (037). Topic and payload are both written into storage the pump owns: @@ -164,9 +201,12 @@ pub struct OutboundScratch { pub topic: String, // new; capacity fixed once per route, never grown } -pub enum Dest { Default, Written, Owned(String) } // Owned: TopicProvider adapter - -pub struct OutboundFrame { pub dest: Dest, pub payload: SerializedPayload } +pub struct OutboundFrame { + /// `true`: `scratch.topic` holds the destination (a `TopicWriter` wrote + /// it). `false`: use the link's default topic. + pub written_topic: bool, + pub payload: SerializedPayload, +} pub trait SerializedReader: Send { fn poll_recv_into( @@ -178,17 +218,18 @@ pub trait SerializedReader: Send { } ``` -- Removes A3. `SerializedPayload::Owned` stays as the fallback for - serializers without an into-slice path (A6 remains by choice of the user). -- `Dest::Written` means `scratch.topic` holds the destination. The scratch - is a `String` so the pump borrows it as `&str` without re-validating. -- The async `recv` / `recv_into` pair is removed rather than adapted: a poll - method cannot wrap an async one without storing a boxed future, and the - only implementors are in this repository (the fused reader in - `typed_api.rs` and the bench). +- Removes A3. `SerializedPayload::Owned` stays for owned serializers (A6, + §2 non-goals). +- The scratch topic is a `String` so the pump borrows it as `&str` without + re-validating. +- The async `recv` / `recv_into` pair and `RecvSerializedIntoFuture` are + removed. The only implementors are in this repository (the fused reader + in `typed_api.rs` and the bench). ### 4.3 Outbound: topics are written, not returned +`TopicProvider` is replaced by: + ```rust pub trait TopicWriter: Send + Sync { /// Write the destination into `out`. `Ok(false)` = use the link's default. @@ -200,12 +241,14 @@ pub trait TopicWriter: Send + Sync { // usage: write!(out, "sensors/{}/{}", v.site, v.id)?; Ok(true) ``` -- `.with_topic_writer(capacity, writer)` on the outbound builder; capacity is - the topic scratch size for that route. +- `.with_topic_writer(capacity, writer)` on the outbound builder replaces + `.with_topic_provider(provider)`; capacity is the topic scratch size for + that route. This is the one user-facing break in this design. - An overflow skips the message and increments a per-link counter; the topic is never truncated. -- `TopicProvider` keeps working through an adapter that yields - `Dest::Owned` (A5 remains for code that has not migrated). +- Closures implement `TopicWriter` through a blanket impl, so a one-line + provider stays a one-line writer. +- Removes A5. ### 4.4 Outbound: per-route publishers in poll form @@ -235,35 +278,45 @@ pub trait RoutePublisher: Send { - This is the `futures::Sink` shape: backpressure through `poll_ready`, a synchronous hand-over, no future per message. `&mut self` per route means - no locks. + no locks inside the trait; a publisher that feeds shared connector state + (§4.6) owns that synchronisation. - `start_send` borrows `dest` and `payload` from the pump's scratch; the publisher decides how to keep them (a pre-allocated ring, a client queue). - Configuration (`qos`, `retain`, …) is parsed once in `route_publisher`, - not per message as the MQTT sinks do today. -- **Error attribution.** A publisher that completes a hand-over + not per message as the MQTT sinks do today. An invalid configuration is + reported there, once, instead of failing every message. +- **Errors never end a route.** A publisher that completes a hand-over asynchronously reports its failure from the next `poll_ready` or - `poll_flush`, i.e. one message late. The publisher logs the failed - destination itself (it still owns the copy); the pump only logs that the - route saw an error. + `poll_flush`, i.e. one message late, and clears it once reported. The + publisher logs the failed destination itself (it still owns the copy); + the pump logs that the route saw an error and continues, as `pump_sink` + does today. - **Shutdown.** When the reader ends (buffer closed), the pump awaits `poll_flush` once before dropping the publisher, so the last message is not silently abandoned. -- A provided `BoxedPublisher` adapter wraps an old-style async `publish` - for third-party connectors during migration, keeping today's cost. +- There is no adapter for the old `publish`. Every in-tree connector + implements `RoutePublisher` directly (§6). With 4.2–4.4 the pump's per-message loop is: ```rust loop { // Readiness first (Sink order): a blocked transport shows up as buffer - // lag, not as one stale message held in the pump. - poll_fn(|cx| publisher.poll_ready(cx)).await?; + // lag, not as one stale message held in the pump. An error here belongs + // to an earlier message. + if let Err(e) = poll_fn(|cx| publisher.poll_ready(cx)).await { + log_warn!("route '{}': earlier publish failed: {:?}", default_topic, e); + continue; + } let frame = match poll_fn(|cx| reader.poll_recv_into(cx, &ctx, &mut scratch)).await { Ok(frame) => frame, Err(DbError::BufferLagged { .. }) => continue, Err(_) => break, }; - publisher.start_send(dest_of(&frame, &scratch, &default_topic), payload_of(&frame, &scratch))?; + let dest = if frame.written_topic { scratch.topic.as_str() } else { &default_topic }; + if let Err(e) = publisher.start_send(dest, payload_of(&frame, &scratch)) { + log_error!("route '{}': publish to '{}' failed: {:?}", default_topic, dest, e); + } } let _ = poll_fn(|cx| publisher.poll_flush(cx)).await; ``` @@ -282,34 +335,75 @@ Document this next to the buffer types; no API change. ### 4.6 MQTT connector -**Embedded backend** (`mountain-mqtt` codec, AimDB's own `session_loop`): - -- *Inbound.* `client_loop` calls `dispatch(message.topic_name, - message.payload)` where it decodes the `ApplicationMessage` today. The - `EventChannel`, `AimdbMqttEvent` and `MqttSource` are removed. -- *Outbound.* Removing the action copies is not enough: `perform` then - encodes each packet into a fresh `Vec` for the write queue. Both are - replaced by one byte ring allocated at build: - - `start_send` encodes the complete PUBLISH packet (fixed header, - topic, packet id, payload) directly into the ring, using - `MqttLenWriter` for the size and `MqttBufWriter` over the reserved - slot. The packet id is reserved from the client state at this point - (§8). - - `write_out` writes contiguous ring slices to the socket and releases - them; the ring replaces `Channel, 4>` for all packets - (publish, subscribe, ping). - - `poll_ready` = ring has space for the route's worst-case frame - (topic capacity + payload scratch + header). - - The session loop no longer sees publishes as actions; it only keeps - client state (QoS acks, keep-alive). -- Result: zero allocations per message on both directions, which also - gives std users a fully allocation-free MQTT path through - `TokioTcpDialer` (052). - -**Native backend** (`rumqttc`): - -- *Inbound.* `dispatch(&publish.topic, &publish.payload)` in the - event-loop task; `MqttEventLoopSource` is removed. +#### Embedded backend (`mountain-mqtt` codec, AimDB's own `session_loop`) + +*Inbound.* `drain_packets` calls `dispatch(publish.topic_name(), +publish.payload())` where it decodes the publish today. `EventChannel`, +`AimdbMqttEvent`, `Received::Event` and `MqttSource` are removed. + +*Outbound.* The session loop's client state (`ClientStateNoQueue`) decides +when a PUBLISH may go out, and it is the only thing that may assign packet +ids: + +```rust +// session_loop.rs today: the action arm is armed only when +let action_ready = connected + && !state.waiting_for_responses() // one in-flight slot + && next_topic >= subscribe_topics.len() + && !outbound.is_full(); +``` + +So route publishers do not encode packets. The hand-over is split across two +pre-allocated rings, and the session loop stays the only encoder: + +```text +RoutePublisher × N ──copy──▶ record ring ──▶ session loop ──encode──▶ write ring ──▶ write_out ──▶ socket + (qos, retain, topic, (per connector, (sole owner of (per connection, + payload) survives ClientStateNoQueue, bip buffer) + reconnects) same gate as today) +``` + +- **Record ring** (replaces `ActionChannel`). One per connector, a byte ring + of variable-length records `{qos, retain, topic_len, payload_len, topic, + payload}`, allocated at build. `start_send` copies one record in; + `poll_ready` = room for this route's largest record (header + topic + capacity + payload scratch capacity, known at `route_publisher` time). + Publishers waiting for room register in a multi-waker; the session loop + waits on a signal raised by `start_send`. The ring lives as long as the + connector, so records queued during an outage are sent after reconnecting, + as actions are today. +- **Session loop.** Where `perform` runs today, the loop peeks the next + record when `action_ready` holds and the write ring has room for its + encoded size (`MqttLenWriter`), calls `state.publish_packet(topic, + payload, qos, retain)`, encodes the packet straight into the write ring + (`MqttBufWriter` over the reserved slot), calls `state.publish_update`, + then releases the record. Packet ids, the one-in-flight rule and the + connect/subscribe ordering are untouched. +- **Write ring** (replaces `Channel, 4>` and `encode`'s `Vec`). One + per connection, cleared on reconnect so no bytes from an old session + reach a new socket. All packets go through it (CONNECT, SUBSCRIBE, + PUBLISH, PUBACK, PINGREQ). `MqttBufWriter` needs a contiguous slice, so + the ring is a bip buffer: a frame that does not fit before the end skips + to the start. `write_out` writes contiguous slices to the socket and + releases them. +- **Control packets.** The session loop encodes a PUBLISH only when the + write ring keeps a reserve for control packets (PUBACK, PINGREQ), so a + stalled socket full of publishes cannot block an acknowledgement. PINGREQ + stays lossy, as today. +- **Oversize frames.** `route_publisher` rejects a route whose largest + record cannot fit in the record ring, or whose encoded frame cannot fit in + the write ring minus the control reserve. + +Cost: two copies of topic and payload per message (into the record, then +into the frame), zero allocations. This is the price of keeping the state +machine single-owner; it also gives std users an allocation-free MQTT path +through `TokioTcpDialer` (052). + +#### Native backend (`rumqttc`) + +- *Inbound.* `dispatch(&publish.topic, &publish.payload)` in the event-loop + task; `MqttEventLoopSource` is removed. This removes both the topic clone + and the `Arc` payload copy. - *Outbound.* `AsyncClient` takes owned `String` and payload, so the two copies stay (G2). Two options for the publisher, to be settled when the native backend migrates: @@ -328,24 +422,21 @@ Document this next to the buffer types; no API change. ## 5. Measurement and gate -`b0_alloc_connector` gains rows for the new interfaces. Targets after this -design: +`b0_alloc_connector` is rewritten against the new interfaces. Targets after +this design: | Row | Before | After | |---|---|---| -| `inbound_route` | 0 | 0 | -| `inbound_route_pattern` | 0 | 0 | -| `inbound_route_keyed_known` | 0 | 0 | -| `inbound_route_keyed_new` | 1 | 1 (by design, 055 §5.6) | -| `inbound_dispatch` (new) | — | 0 | -| `inbound_dispatch_pattern` (new) | — | 0 | -| `inbound_pump_source_minimal` (poll `Source`) | 2 | 1 (the owned topic) | +| `inbound_dispatch` (was `inbound_route`) | 0 | 0 | +| `inbound_dispatch_pattern` (was `inbound_route_pattern`) | 0 | 0 | +| `inbound_dispatch_keyed_known` | 0 | 0 | +| `inbound_dispatch_keyed_new` | 1 | 1 (by design, 055 §5.6) | +| `inbound_pump_source_minimal` | 2 | removed (`Source` is gone) | | `outbound_scratch_static_topic` | 2 | 0 | -| `outbound_scratch_written_topic` (new) | — | 0 | -| `outbound_scratch_dynamic_topic` (`TopicProvider` adapter) | 3 | 1 | -| `outbound_owned_static_topic` | 3 | 1 | +| `outbound_scratch_written_topic` (replaces `outbound_scratch_dynamic_topic`) | 3 | 0 | +| `outbound_owned_static_topic` | 3 | 1 (A6, by choice of serializer) | -The bench asserts exact values, so every step updates `EXPECTED` and the +The bench asserts exact values, so every change updates `EXPECTED` and the baseline together. A connector-level row per MQTT backend (loopback broker for native, the existing loopback harness for embedded) records G2; the embedded rows target 0 in both directions, the native outbound row records @@ -353,32 +444,42 @@ the two copies plus whatever `rumqttc` adds. **Gate.** A CI job runs `cargo bench -p aimdb-bench --bench b0_alloc_connector` (host only, a few seconds) and fails on any assertion. -It is added in step (1) of §6, runs on every PR from then on, and becomes a -required check in step (4). It is the counting-allocator bench, so its -results are deterministic and do not need a quiet runner. +It lands first, against today's interfaces (§6 step 1), and is a required +check from then on. It is the counting-allocator bench, so its results are +deterministic and do not need a quiet runner. ## 6. Migration | Who | Change | |---|---| -| Application code | None | -| Connector authors | `inbound_router` + `pump_source` → `InboundDispatch` (or a poll `Source`); `Connector::publish` → `route_publisher` + `RoutePublisher` (or the `BoxedPublisher` adapter) | +| Application code | `.with_topic_provider(p)` → `.with_topic_writer(capacity, w)`. Nothing else. | +| Connector authors | `inbound_router` + `pump_source` / `router.route` → `InboundDispatch`; `Connector::publish` → `route_publisher` + `RoutePublisher`. No adapters. | | Custom `SerializedReader` implementors | `recv` / `recv_into` → `poll_recv_into` (none known outside this repository) | -| `TopicProvider` users | None; optional move to `TopicWriter` | -In-tree connectors to migrate and measure: MQTT (both backends), KNX, -WebSocket (client and server), the AimX session client (used by TCP, UDS and -serial), and the Embassy adapter's connector glue. The Zenoh connector -(053) is not implemented yet; it is written against `InboundDispatch` and -`RoutePublisher` from the start. +Without adapters, changing the core traits breaks every connector at once, +so the work lands as one change (a feature branch merged once): -**Order:** (1) core interfaces with adapters for the old forms, bench rows -added, CI job running the bench; (2) MQTT embedded (including the write-ring -rework of §4.6), then native; (3) remaining connectors, one PR each; -(4) remove the adapters' use in-tree and make the gate a required check. - -The native topic move (`publish.topic` instead of `.clone()`, §1) needs no -interface change and can land before step (1). +| Crate | Inbound | Outbound | +|---|---|---| +| `aimdb-core` | `InboundDispatch`; remove `Source`, `pump_source`, public `inbound_router`; AimX session client (TCP, UDS, serial) moves from `router.route` | poll `SerializedReader`, `TopicWriter`, `RoutePublisher`, new pump; remove `Connector::publish`, `TopicProvider` | +| `aimdb-mqtt-connector` embedded | `dispatch` in `drain_packets` | record ring + write ring (§4.6) | +| `aimdb-mqtt-connector` native | `dispatch` in event-loop task | `ReusableBoxFuture` publisher | +| `aimdb-knx-connector` | `Source` → `dispatch` | `RoutePublisher` | +| `aimdb-websocket-connector` | rename (`router.route` → `dispatch`), client and server | server `RoutePublisher` | +| `aimdb-embassy-adapter` | `send_wrapper.rs` wrappers for `Source` / `Connector` removed or reduced to `RoutePublisher` | — | +| tests, bench | codec tests that build a router to ingest (`aimdb-data-contracts/src/link_codec.rs`, `aimdb-mqtt-connector/tests/link_ext_tests.rs`) → `InboundDispatch` (`route_count`, `subscriptions`, `dispatch`) | `topic_provider_tests.rs` (mqtt, knx), websocket e2e → `TopicWriter`; bench rewritten (§5) | + +**Order:** + +1. On main, independently: the CI gate for today's bench, and the native + MQTT topic move (`publish.topic` instead of `.clone()`). +2. On the feature branch: core interfaces and pump, bench rewritten. +3. MQTT embedded (§4.6), then native. +4. KNX, WebSocket, Embassy adapter, data contracts, tests. +5. Merge once the gate and every connector's tests pass. + +The Zenoh connector (053) is not implemented yet; it is written against +`InboundDispatch` and `RoutePublisher` from the start. ## 7. Alternatives considered @@ -396,25 +497,32 @@ interface change and can land before step (1). 4. **Keep `Connector::publish` and pool its futures.** Futures of different connectors have different sizes; pooling adds unsafe layout handling for no gain over `poll_ready`/`start_send`. -5. **Keep the embedded action channel, with fixed-size `heapless` topic and +5. **Embedded: keep the action channel, with fixed-size `heapless` topic and payload.** Removes the action allocations but not the encode `Vec`, and - sizes every slot for the largest route. Encoding straight into the write - ring does both. + sizes every slot for the largest route. The record ring (§4.6) holds + variable-length records instead. +6. **Embedded: route publishers encode PUBLISH into the write ring + (previous revision of this design).** One copy instead of two, but + publishers would need the session's client state: whether it is + connected and subscribed, whether the single in-flight slot is free, and + the next packet id. Frames encoded during an outage would reach a new + socket ahead of CONNECT. Sharing the state across tasks needs a lock, and + resetting the ring on reconnect loses queued messages. Keeping the session + loop as the only encoder avoids all of it for one extra memcpy. +7. **Keep compatibility adapters (`BoxedPublisher`, poll `Source`, + `TopicProvider` adapter) for a staged migration.** Allows one PR per + connector, but keeps two mechanisms per direction and their bench rows. + Every connector is in this repository, so one migration is cheaper. ## 8. Open questions -- **Session-task latency (embedded).** Ingest moves into the MQTT session - task. It is synchronous and bounded, but a slow user deserializer now - delays keep-alive handling. Measure on the STM32H5 rig (037's B3). -- **Ring sizing (embedded).** One shared ring for all routes and control - packets, or reserved space for control packets so a stalled publisher - cannot delay a PINGREQ; default size; how a frame larger than the ring is - rejected (at `route_publisher` time, from the route's capacities, is - preferred). -- **Packet ids (embedded).** Encoding in `start_send` means the publisher - reserves a QoS 1/2 packet id outside the session loop. Either the client - state is shared behind the ring's lock, or packet ids are patched into the - frame by the session loop before it releases the slot to `write_out`. +- **Ingest latency on the transport task.** Embedded: a slow user + deserializer delays keep-alive and ack handling in the session loop. + Native: it delays polling `rumqttc`'s event loop. Measure on the STM32H5 + rig (037's B3) and with the native loopback row; decide whether to + document a deserializer budget. +- **Ring sizing (embedded).** Defaults for the record ring and the write + ring, and the size of the write ring's control reserve. - **Other connectors' own copies.** KNX, WebSocket and the AimX session were not measured at connector level; §5's per-connector rows will show them. @@ -429,5 +537,5 @@ interface change and can land before step (1). - [055 — Wildcard inbound links](./055-wildcard-inbound-links.md) (grammar, keys, `TopicMatch`, §7 relation to this design) - `aimdb-bench/benches/b0_alloc_connector.rs` -- `aimdb-mqtt-connector/src/embedded/session_loop.rs` (`perform`, `encode`, - `write_out`) +- `aimdb-mqtt-connector/src/embedded/session_loop.rs` (`client_loop`, + `drain_packets`, `perform`, `encode`, `write_out`) From 076006ba7b5ae1578f272a812ad1b173c70c0e8e Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 2 Oct 2026 08:32:05 +0000 Subject: [PATCH 2/5] docs(054): publishers lend buffers; per-route lock-free record rings - RoutePublisher becomes poll_ready / buffers / commit / poll_flush: the reader serializes topic and payload straight into storage the publisher reserved, removing the scratch-to-publisher copy. route_publisher takes a RouteSpec and returns a Result. - Embedded MQTT: one single-producer, single-consumer record ring per route, so no lock; the session loop encodes round-robin into the write ring, leaving one AimDB copy per message. - Native MQTT moves owned serializer payloads into rumqttc without copying. - Alternatives: hand-encoded zero-copy frames (8) and borrowed-bytes start_send with a shared ring (9) rejected; ring sizing and bbqueue fit added to open questions. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01EyVFXVWAec2aouRceoGdpE --- .../054-zero-alloc-connector-boundary.md | 290 ++++++++++++------ 1 file changed, 189 insertions(+), 101 deletions(-) diff --git a/docs/design/054-zero-alloc-connector-boundary.md b/docs/design/054-zero-alloc-connector-boundary.md index 834861fa..41752478 100644 --- a/docs/design/054-zero-alloc-connector-boundary.md +++ b/docs/design/054-zero-alloc-connector-boundary.md @@ -1,8 +1,10 @@ # 054 — Zero-allocation connector boundary **Status:** 📝 Proposed — revised 2026-10-02: no backward compatibility -(adapters, `Source` and `TopicProvider` removed; one migration), embedded -MQTT outbound reworked so the session loop stays the only encoder (§4.6). +(adapters, `Source` and `TopicProvider` removed; one migration); publishers +lend the buffers the reader serializes into (§4.2, §4.4); embedded MQTT +outbound uses one lock-free record ring per route and keeps the session loop +as the only encoder (§4.6). Previous revision 2026-09-30, after 055 landed. **Scope:** the per-message path between `aimdb-core` and connectors, in both @@ -189,23 +191,35 @@ impl InboundDispatch { loop; for native MQTT it is the task polling `rumqttc`'s event loop. A slow user deserializer delays keep-alive and ack handling in both (§8). -### 4.2 Outbound: poll-based reader +### 4.2 Outbound: the publisher lends the buffers -`SerializedReader` has one method, in poll form, over the buffer reader's -existing `poll_recv` (037). Topic and payload are both written into storage -the pump owns: +The serialized topic and payload are written straight into storage the +publisher owns. For the embedded MQTT backend that storage is a slot in the +route's record ring (§4.6), so the bytes are never copied between +serialization and encoding. For other connectors it is a per-route scratch +the publisher allocates once, which replaces the pump's scratch of today. ```rust -pub struct OutboundScratch { - pub payload: Vec, // allocated once per route (exists today) - pub topic: String, // new; capacity fixed once per route, never grown +/// Where the reader writes one message. Borrowed from the publisher. +pub struct OutboundBuffers<'a> { + pub topic: TopicBuf<'a>, // capacity = the route's topic capacity (§4.3) + pub payload: &'a mut [u8], // capacity = the route's payload capacity } +/// What the reader wrote. pub struct OutboundFrame { - /// `true`: `scratch.topic` holds the destination (a `TopicWriter` wrote - /// it). `false`: use the link's default topic. - pub written_topic: bool, - pub payload: SerializedPayload, + /// `Some(len)`: a `TopicWriter` wrote `len` bytes into `buffers.topic`. + /// `None`: use the route's default topic. + pub topic: Option, + pub payload: CommittedPayload, +} + +pub enum CommittedPayload { + /// `len` bytes written into `buffers.payload` (`with_serializer_into`). + InPlace(usize), + /// An owned serializer's bytes (`with_serializer`, A6). The publisher + /// moves or copies them as its transport needs. + Owned(Vec), } pub trait SerializedReader: Send { @@ -213,18 +227,20 @@ pub trait SerializedReader: Send { &mut self, cx: &mut Context<'_>, ctx: &RuntimeContext, - scratch: &mut OutboundScratch, + buffers: OutboundBuffers<'_>, ) -> Poll>; } ``` -- Removes A3. `SerializedPayload::Owned` stays for owned serializers (A6, - §2 non-goals). -- The scratch topic is a `String` so the pump borrows it as `&str` without - re-validating. -- The async `recv` / `recv_into` pair and `RecvSerializedIntoFuture` are - removed. The only implementors are in this repository (the fused reader - in `typed_api.rs` and the bench). +- Removes A3. The reader runs the topic writer and the serializer only when + a value is ready, so it writes into `buffers` at most once per `Ready`. +- `CommittedPayload::Owned` keeps owned serializers working (A6, §2 + non-goals). It is also the cheaper path for transports whose client takes + an owned payload: the native MQTT publisher moves the `Vec` into + `rumqttc` instead of copying it (§4.6). +- The async `recv` / `recv_into` pair, `RecvSerializedIntoFuture` and + `SerializedPayload` are removed. The only implementors are in this + repository (the fused reader in `typed_api.rs` and the bench). ### 4.3 Outbound: topics are written, not returned @@ -236,14 +252,16 @@ pub trait TopicWriter: Send + Sync { fn write_topic(&self, value: &T, out: &mut TopicBuf<'_>) -> Result; } -// TopicBuf implements core::fmt::Write over the route's topic scratch and -// refuses any write past its fixed capacity. +// TopicBuf implements core::fmt::Write over a byte slice the publisher lends +// (§4.2) and refuses any write past its end. Only `&str` goes in, so the +// bytes are valid UTF-8 by construction. // usage: write!(out, "sensors/{}/{}", v.site, v.id)?; Ok(true) ``` - `.with_topic_writer(capacity, writer)` on the outbound builder replaces - `.with_topic_provider(provider)`; capacity is the topic scratch size for - that route. This is the one user-facing break in this design. + `.with_topic_provider(provider)`; capacity is the topic space the publisher + reserves per message on that route. This is the one user-facing break in + this design. - An overflow skips the message and increments a per-link counter; the topic is never truncated. - Closures implement `TopicWriter` through a blanket impl, so a one-line @@ -256,19 +274,33 @@ pub trait TopicWriter: Send + Sync { (A4). It is replaced by a publisher created once per route: ```rust +/// Everything a publisher needs to know about its route, fixed at build. +pub struct RouteSpec<'a> { + pub default_topic: &'a str, + pub config: &'a ConnectorConfig, // qos, retain, … from the link URL + pub topic_capacity: usize, // 0 without a TopicWriter + /// `Some(n)`: into-slice serializer with an n-byte scratch. + /// `None`: owned serializer; the publisher picks its own limit. + pub payload_capacity: Option, +} + pub trait Connector: Send + Sync { - /// Called once per outbound route when the pump starts. - fn route_publisher(&self, config: &ConnectorConfig) -> Box; + /// Called once per outbound route while the connector builds its pumps, + /// before its transport tasks start. + fn route_publisher(&self, route: &RouteSpec<'_>) + -> Result, PublishError>; } pub trait RoutePublisher: Send { - /// Ready to accept one message (e.g. queue space). Also drives any - /// hand-over still in progress from the previous `start_send`. + /// Ready to accept one message: storage for it is reserved. Also drives + /// any hand-over still in progress from the previous `commit`. fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll>; - /// Hand one message over. Synchronous; the publisher copies what it - /// needs into storage it owns. - fn start_send(&mut self, dest: &str, payload: &[u8]) -> Result<(), PublishError>; - /// Complete every hand-over started so far. Default: `Ready(Ok(()))`. + /// The storage reserved by the last successful `poll_ready`. Calling it + /// again before `commit` returns the same storage. + fn buffers(&mut self) -> OutboundBuffers<'_>; + /// The message written into `buffers` is complete. Synchronous. + fn commit(&mut self, frame: OutboundFrame) -> Result<(), PublishError>; + /// Complete every hand-over committed so far. Default: `Ready(Ok(()))`. fn poll_flush(&mut self, cx: &mut Context<'_>) -> Poll> { let _ = cx; Poll::Ready(Ok(())) @@ -276,21 +308,24 @@ pub trait RoutePublisher: Send { } ``` -- This is the `futures::Sink` shape: backpressure through `poll_ready`, a - synchronous hand-over, no future per message. `&mut self` per route means - no locks inside the trait; a publisher that feeds shared connector state - (§4.6) owns that synchronisation. -- `start_send` borrows `dest` and `payload` from the pump's scratch; the - publisher decides how to keep them (a pre-allocated ring, a client queue). -- Configuration (`qos`, `retain`, …) is parsed once in `route_publisher`, - not per message as the MQTT sinks do today. An invalid configuration is - reported there, once, instead of failing every message. +- This is the `futures::Sink` shape with the buffer lent out instead of + passed in (the grant pattern of `bbqueue`): backpressure through + `poll_ready`, a synchronous hand-over, no future per message, no copy into + the publisher. +- A reservation is held while the reader waits for the next value. That is + why it belongs to one route: a publisher must not reserve in storage other + routes share (§4.6 gives each route its own ring). +- `route_publisher` parses configuration (`qos`, `retain`, …) once, not per + message as the MQTT sinks do today, and rejects a route it cannot serve + (invalid configuration, a message that can never fit) before anything + runs. It returns a `Result` for that reason. - **Errors never end a route.** A publisher that completes a hand-over asynchronously reports its failure from the next `poll_ready` or `poll_flush`, i.e. one message late, and clears it once reported. The - publisher logs the failed destination itself (it still owns the copy); - the pump logs that the route saw an error and continues, as `pump_sink` - does today. + publisher logs the failed destination itself; the pump logs that the route + saw an error and continues, as `pump_sink` does today. A message that does + not fit (topic overflow, an owned payload larger than the publisher's + limit) is skipped and counted, and the reservation is reused. - **Shutdown.** When the reader ends (buffer closed), the pump awaits `poll_flush` once before dropping the publisher, so the last message is not silently abandoned. @@ -308,19 +343,22 @@ loop { log_warn!("route '{}': earlier publish failed: {:?}", default_topic, e); continue; } - let frame = match poll_fn(|cx| reader.poll_recv_into(cx, &ctx, &mut scratch)).await { + // Serialize straight into the publisher's reserved storage. + let frame = match poll_fn(|cx| reader.poll_recv_into(cx, &ctx, publisher.buffers())).await { Ok(frame) => frame, - Err(DbError::BufferLagged { .. }) => continue, + Err(DbError::BufferLagged { .. }) => continue, // reservation kept Err(_) => break, }; - let dest = if frame.written_topic { scratch.topic.as_str() } else { &default_topic }; - if let Err(e) = publisher.start_send(dest, payload_of(&frame, &scratch)) { - log_error!("route '{}': publish to '{}' failed: {:?}", default_topic, dest, e); + if let Err(e) = publisher.commit(frame) { + log_error!("route '{}': publish failed: {:?}", default_topic, e); } } let _ = poll_fn(|cx| publisher.poll_flush(cx)).await; ``` +The pump no longer owns any per-route buffer; `RouteSpec` carries the sizes +the old pump scratch was built from. + ### 4.5 Values with heap data (guidance) Delivery clones `T` once per reader. For records with more than one reader @@ -353,50 +391,73 @@ let action_ready = connected && !outbound.is_full(); ``` -So route publishers do not encode packets. The hand-over is split across two -pre-allocated rings, and the session loop stays the only encoder: +So route publishers do not encode packets. Each route serializes into its +own record ring; the session loop stays the only encoder and writes into +one write ring: ```text -RoutePublisher × N ──copy──▶ record ring ──▶ session loop ──encode──▶ write ring ──▶ write_out ──▶ socket - (qos, retain, topic, (per connector, (sole owner of (per connection, - payload) survives ClientStateNoQueue, bip buffer) - reconnects) same gate as today) +route 1 pump ──serialize in place──▶ record ring 1 ─┐ +route 2 pump ──serialize in place──▶ record ring 2 ─┼─▶ session loop ──encode (1 copy)──▶ write ring ──▶ write_out ──▶ socket +route N pump ──serialize in place──▶ record ring N ─┘ (sole owner of (per connection, + (one producer each, lock-free, ClientStateNoQueue, bip buffer) + survives reconnects) same gate as today) ``` -- **Record ring** (replaces `ActionChannel`). One per connector, a byte ring - of variable-length records `{qos, retain, topic_len, payload_len, topic, - payload}`, allocated at build. `start_send` copies one record in; - `poll_ready` = room for this route's largest record (header + topic - capacity + payload scratch capacity, known at `route_publisher` time). - Publishers waiting for room register in a multi-waker; the session loop - waits on a signal raised by `start_send`. The ring lives as long as the - connector, so records queued during an outage are sent after reconnecting, - as actions are today. -- **Session loop.** Where `perform` runs today, the loop peeks the next - record when `action_ready` holds and the write ring has room for its - encoded size (`MqttLenWriter`), calls `state.publish_packet(topic, - payload, qos, retain)`, encodes the packet straight into the write ring - (`MqttBufWriter` over the reserved slot), calls `state.publish_update`, - then releases the record. Packet ids, the one-in-flight rule and the - connect/subscribe ordering are untouched. +- **Record rings** (replace `ActionChannel`). One per route, created in + `route_publisher` while the connector builds its pumps, so the session + loop starts with a fixed list of ring consumers. Each ring is a + single-producer, single-consumer bip buffer of variable-length records + `{topic_len, payload_len, topic, payload}`; `qos`, `retain` and the + default topic are stored once per ring, not per record. + - `poll_ready` reserves one contiguous slot of `topic_capacity + + payload_capacity` (plus a small record header) and waits if the ring has + no room. + - `buffers` lends the slot's topic and payload regions; the serializer + and `TopicWriter` write directly into them. + - `commit` writes the record header with the real lengths and releases the + unused tail of the slot. `CommittedPayload::Owned` is copied into the + slot (A6 keeps its copy); an owned payload larger than the slot is + skipped and counted. + - Records queued during an outage stay in their rings and are sent after + reconnecting, as actions are today. +- **Synchronisation.** No lock. Each record ring has exactly one producer + (its route's pump) and one consumer (the session loop), so head and tail + are atomics. The pump parks on the ring's single waker when it is full; + `commit` raises one signal the session loop waits on. Neither + `ClientStateNoQueue` nor the write ring is shared: the session loop owns + the state and is the write ring's only writer, `write_out` its only + reader. +- **Session loop.** Where `perform` runs today, when `action_ready` holds + the loop takes the oldest committed record from the rings in round-robin + order, checks the write ring has room for its encoded size + (`MqttLenWriter`), calls `state.publish_packet(topic, payload, qos, + retain)` with slices borrowed from the record ring, encodes the packet + straight into the write ring (`MqttBufWriter` over the reserved slot), + calls `state.publish_update`, then releases the record. Round-robin keeps + one busy route from starving the others; packet ids, the one-in-flight + rule and the connect/subscribe ordering are untouched. - **Write ring** (replaces `Channel, 4>` and `encode`'s `Vec`). One per connection, cleared on reconnect so no bytes from an old session reach a new socket. All packets go through it (CONNECT, SUBSCRIBE, PUBLISH, PUBACK, PINGREQ). `MqttBufWriter` needs a contiguous slice, so - the ring is a bip buffer: a frame that does not fit before the end skips - to the start. `write_out` writes contiguous slices to the socket and - releases them. + it is a bip buffer as well. `write_out` writes contiguous slices to the + socket and releases them. - **Control packets.** The session loop encodes a PUBLISH only when the write ring keeps a reserve for control packets (PUBACK, PINGREQ), so a stalled socket full of publishes cannot block an acknowledgement. PINGREQ stays lossy, as today. -- **Oversize frames.** `route_publisher` rejects a route whose largest - record cannot fit in the record ring, or whose encoded frame cannot fit in - the write ring minus the control reserve. - -Cost: two copies of topic and payload per message (into the record, then -into the frame), zero allocations. This is the price of keeping the state -machine single-owner; it also gives std users an allocation-free MQTT path +- **Sizes and oversize routes.** A route's record ring holds at least two + slots, so the pump can fill the next record while the loop encodes the + previous one. For an owned serializer (`payload_capacity: None`), the slot + uses the connector's maximum payload, a builder option. `route_publisher` + rejects a route whose encoded frame cannot fit in the write ring minus the + control reserve. + +Cost per message: zero allocations and one copy of topic and payload (record +ring into the encoded frame), in addition to the network stack's own copy +into its TX buffer, which TLS needs anyway to encrypt. The remaining copy is +the price of keeping packet encoding inside the session loop (§7, +alternative 8). This also gives std users an allocation-free MQTT path through `TokioTcpDialer` (052). #### Native backend (`rumqttc`) @@ -404,14 +465,16 @@ through `TokioTcpDialer` (052). - *Inbound.* `dispatch(&publish.topic, &publish.payload)` in the event-loop task; `MqttEventLoopSource` is removed. This removes both the topic clone and the `Arc` payload copy. -- *Outbound.* `AsyncClient` takes owned `String` and payload, so the two - copies stay (G2). Two options for the publisher, to be settled when the - native backend migrates: - 1. `start_send` builds the `AsyncClient::publish` future and stores it in +- *Outbound.* `AsyncClient` takes an owned `String` and payload (G2). The + publisher lends a per-route topic and payload scratch from `buffers`; + `commit` builds the owned topic and copies the in-place payload, or moves a + `CommittedPayload::Owned` `Vec` straight in without copying. Two options + for the hand-over, to be settled when the native backend migrates: + 1. `commit` builds the `AsyncClient::publish` future and stores it in a per-route `tokio_util::sync::ReusableBoxFuture` (037's Tokio technique); `poll_ready` / `poll_flush` drive it. Keeps backpressure, removes the box, reports errors one message late (§4.4). - 2. `start_send` calls the synchronous `AsyncClient::try_publish`. No + 2. `commit` calls the synchronous `AsyncClient::try_publish`. No stored future, errors attributed immediately, but no readiness signal: a full request channel drops the message and counts it. @@ -439,8 +502,9 @@ this design: The bench asserts exact values, so every change updates `EXPECTED` and the baseline together. A connector-level row per MQTT backend (loopback broker for native, the existing loopback harness for embedded) records G2; the -embedded rows target 0 in both directions, the native outbound row records -the two copies plus whatever `rumqttc` adds. +embedded rows target 0 allocations in both directions and record the one +outbound copy (§4.6); the native outbound row records the owned topic, the +payload copy and whatever `rumqttc` adds. **Gate.** A CI job runs `cargo bench -p aimdb-bench --bench b0_alloc_connector` (host only, a few seconds) and fails on any assertion. @@ -453,7 +517,7 @@ deterministic and do not need a quiet runner. | Who | Change | |---|---| | Application code | `.with_topic_provider(p)` → `.with_topic_writer(capacity, w)`. Nothing else. | -| Connector authors | `inbound_router` + `pump_source` / `router.route` → `InboundDispatch`; `Connector::publish` → `route_publisher` + `RoutePublisher`. No adapters. | +| Connector authors | `inbound_router` + `pump_source` / `router.route` → `InboundDispatch`; `Connector::publish` → `route_publisher(&RouteSpec)` + `RoutePublisher` (`poll_ready`, `buffers`, `commit`, `poll_flush`); the publisher owns the per-route buffers the pump used to own. No adapters. | | Custom `SerializedReader` implementors | `recv` / `recv_into` → `poll_recv_into` (none known outside this repository) | Without adapters, changing the core traits breaks every connector at once, @@ -462,7 +526,7 @@ so the work lands as one change (a feature branch merged once): | Crate | Inbound | Outbound | |---|---|---| | `aimdb-core` | `InboundDispatch`; remove `Source`, `pump_source`, public `inbound_router`; AimX session client (TCP, UDS, serial) moves from `router.route` | poll `SerializedReader`, `TopicWriter`, `RoutePublisher`, new pump; remove `Connector::publish`, `TopicProvider` | -| `aimdb-mqtt-connector` embedded | `dispatch` in `drain_packets` | record ring + write ring (§4.6) | +| `aimdb-mqtt-connector` embedded | `dispatch` in `drain_packets` | per-route record rings + write ring (§4.6) | | `aimdb-mqtt-connector` native | `dispatch` in event-loop task | `ReusableBoxFuture` publisher | | `aimdb-knx-connector` | `Source` → `dispatch` | `RoutePublisher` | | `aimdb-websocket-connector` | rename (`router.route` → `dispatch`), client and server | server `RoutePublisher` | @@ -496,23 +560,41 @@ The Zenoh connector (053) is not implemented yet; it is written against clients expose messages. Push (`dispatch`) fits them directly. 4. **Keep `Connector::publish` and pool its futures.** Futures of different connectors have different sizes; pooling adds unsafe layout handling for - no gain over `poll_ready`/`start_send`. + no gain over `poll_ready`/`commit`. 5. **Embedded: keep the action channel, with fixed-size `heapless` topic and payload.** Removes the action allocations but not the encode `Vec`, and sizes every slot for the largest route. The record ring (§4.6) holds variable-length records instead. 6. **Embedded: route publishers encode PUBLISH into the write ring - (previous revision of this design).** One copy instead of two, but - publishers would need the session's client state: whether it is - connected and subscribed, whether the single in-flight slot is free, and - the next packet id. Frames encoded during an outage would reach a new - socket ahead of CONNECT. Sharing the state across tasks needs a lock, and - resetting the ring on reconnect loses queued messages. Keeping the session - loop as the only encoder avoids all of it for one extra memcpy. + (first revision of this design).** Publishers would need the session's + client state: whether it is connected and subscribed, whether the single + in-flight slot is free, and the next packet id. Frames encoded during an + outage would reach a new socket ahead of CONNECT. Sharing the state + across tasks needs a lock, and resetting the ring on reconnect loses + queued messages. It has the same single copy as §4.6, so it gains nothing + for its problems. 7. **Keep compatibility adapters (`BoxedPublisher`, poll `Source`, `TopicProvider` adapter) for a staged migration.** Allows one PR per connector, but keeps two mechanisms per direction and their bench rows. Every connector is in this repository, so one migration is cheaper. +8. **Embedded: zero AimDB copies by hand-encoding PUBLISH in the record + ring.** Lay each record out as a frame with gaps (fixed header and topic + length in front, a hole for the packet id), let the session loop patch + only the packet id, and have `write_out` send straight from the record + rings. Removes the last copy, but AimDB would encode PUBLISH outside + `mountain-mqtt`'s codec and must stay byte-identical to it (MQTT 5 + properties included), `write_out` would interleave N record rings and the + control ring through a descriptor queue, and records would stay reserved + until the socket write completes, so rings grow. The copy it saves is one + memcpy of topic and payload, small next to the network stack's own copy + and TLS encryption. Revisit only if the STM32H5 rig shows the copy + matters. +9. **Publisher takes borrowed bytes (`start_send(&str, &[u8])`, second + revision of this design).** Simpler trait, but the pump serializes into + its own scratch and every publisher copies from it, so the embedded + backend pays two copies. With one shared record ring for all routes it + also needs a lock among the route pumps. Lending the buffers (§4.4) and + one ring per route remove both. ## 8. Open questions @@ -521,8 +603,14 @@ The Zenoh connector (053) is not implemented yet; it is written against Native: it delays polling `rumqttc`'s event loop. Measure on the STM32H5 rig (037's B3) and with the native loopback row; decide whether to document a deserializer budget. -- **Ring sizing (embedded).** Defaults for the record ring and the write - ring, and the size of the write ring's control reserve. +- **Ring sizing (embedded).** Slots per record ring (at least two), the + default maximum payload for owned-serializer routes, the write ring size + and its control reserve. Per-route rings cost more RAM than one shared + ring; measure with the demo applications' route counts. +- **Ring implementation.** The record rings need a single-producer, + single-consumer bip buffer with grants and runtime-chosen sizes. Check + whether `bbqueue` fits (`no_std`, waker integration, storage sized at + build rather than by const generic) before writing one. - **Other connectors' own copies.** KNX, WebSocket and the AimX session were not measured at connector level; §5's per-connector rows will show them. From 4a2c47ab50abd53e2e9fc324da909a03532a1235 Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 2 Oct 2026 08:38:42 +0000 Subject: [PATCH 3/5] docs(054): build the MQTT rings on bbqueue 0.7 - Record rings use bbqueue's framed grants (grant, commit with unused tail returned, keep/release on read); the write ring uses stream grants with grant_exact(frame + reserve) to keep a control-packet reserve. - BoxedSlice storage sized per route at build, ArcBBQueue halves, AtomicCoord; a custom Notifier with atomic waker slots fits poll_ready instead of the maitake-sync async notifier. - Spell out the synchronisation: no lock in the per-message path, only waker slots; RMW atomics needed, portable-atomic fallback on thumbv6m. - Open questions: bbqueue pre-1.0 churn, write ring draining on reconnect, inbound chunk channel as a follow-up. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01EyVFXVWAec2aouRceoGdpE --- .../054-zero-alloc-connector-boundary.md | 84 ++++++++++++++++--- 1 file changed, 73 insertions(+), 11 deletions(-) diff --git a/docs/design/054-zero-alloc-connector-boundary.md b/docs/design/054-zero-alloc-connector-boundary.md index 41752478..b750cbbe 100644 --- a/docs/design/054-zero-alloc-connector-boundary.md +++ b/docs/design/054-zero-alloc-connector-boundary.md @@ -420,13 +420,24 @@ route N pump ──serialize in place──▶ record ring N ─┘ (sole owne skipped and counted. - Records queued during an outage stay in their rings and are sent after reconnecting, as actions are today. -- **Synchronisation.** No lock. Each record ring has exactly one producer - (its route's pump) and one consumer (the session loop), so head and tail - are atomics. The pump parks on the ring's single waker when it is full; - `commit` raises one signal the session loop waits on. Neither - `ClientStateNoQueue` nor the write ring is shared: the session loop owns - the state and is the write ring's only writer, `write_out` its only - reader. +- **Synchronisation.** No lock in the per-message path. + - Each record ring has exactly one producer (its route's pump) and one + consumer (the session loop). The ring is split into a producer half, + owned by the `RoutePublisher`, and a consumer half, owned by the + session loop, so the type system enforces it. Head and tail are + atomics. + - The write ring has one writer (the session loop) and one reader + (`write_out`), which run in the same task under one `select3`. + - `ClientStateNoQueue` and the list of record rings belong to the session + loop alone. The list is fixed before the loop starts; a reconnect starts + a new loop only after the old one has ended. + - Wakers are the only shared state left: each ring wakes its own pump when + space is released, and every ring's `commit` wakes the session loop. + Those are atomic waker slots, not locks held while data is copied. + - Atomics: the rings need atomic read-modify-write (`swap`, + `fetch_add`), which `thumbv7em`, `thumbv8m.main` and std targets have. + On targets without it (`thumbv6m`), the rings and wakers fall back to + `portable-atomic` or a critical section. - **Session loop.** Where `perform` runs today, when `action_ready` holds the loop takes the oldest committed record from the rings in round-robin order, checks the write ring has room for its encoded size @@ -446,6 +457,45 @@ route N pump ──serialize in place──▶ record ring N ─┘ (sole owne write ring keeps a reserve for control packets (PUBACK, PINGREQ), so a stalled socket full of publishes cannot block an acknowledgement. PINGREQ stays lossy, as today. +- **Implementation: `bbqueue` 0.7.** Both ring kinds are `bbqueue` + queues (single-producer single-consumer bip buffers with grants, Miri + tested): + - *Record rings* use the framed interface with a `u16` length header. + `poll_ready` calls `FramedProducer::grant(slot_len)` and stores the + returned `FramedGrantW` in the publisher; `buffers` lends regions of + it; `commit` calls `FramedGrantW::commit(used)`, which writes the header + and returns the unused tail to the ring. The session loop's + `FramedConsumer::read` yields one record; if the write ring has no room + yet it calls `keep()` (the record stays queued), otherwise it encodes + and calls `release()`. Records are limited to 64 KiB. + - *Write ring* uses the stream interface. The session loop takes + `grant_exact(frame_len + control_reserve)` and commits only + `frame_len`, which keeps the reserve free without any accounting of its + own; `write_out` takes `read()` grants and releases them after + `write_all`. + - Storage is `BoxedSlice` (feature `alloc`): each ring is allocated once + at build, sized from its `RouteSpec`. Queues are held as + `ArcBBQueue`, so the halves carry no lifetimes; taking a grant clones + the `Arc` (a counter increment, no allocation). + - Coordination is `AtomicCoord` (lock-free) where the target has atomic + pointers; `CsCoord` or the `portable-atomic` features otherwise. + - Notifier: `bbqueue`'s async support is an `AsyncNotifier` with `async + fn`s backed by `maitake-sync`, which does not fit `poll_ready`. AimDB + implements the public `Notifier` trait instead, with two atomic waker + slots per ring: `wake_one_producer` wakes the route's pump, + `wake_one_consumer` wakes the session loop. `poll_ready` registers its + waker, then tries `grant`, and returns `Pending` on + `InsufficientSize`; the session loop does the same for `read`. + - Dependency: `bbqueue = { version = "0.7", default-features = false, + features = ["alloc"] }`, which leaves out `maitake-sync` and + `critical-section`. + - Record layout inside a frame: `[topic_len: u16][topic][payload]`. The + publisher lends the payload region at a fixed offset after the route's + topic capacity, so a route with a `TopicWriter` leaves up to + `topic_capacity` unused bytes per queued record. If that matters, the + reader can serialize the payload first and write the topic after it + (`[payload_len][topic_len][payload][topic]`), at the cost of a less + direct `OutboundBuffers` contract. - **Sizes and oversize routes.** A route's record ring holds at least two slots, so the pump can fill the next record while the loop encodes the previous one. For an owned serializer (`payload_capacity: None`), the slot @@ -607,10 +657,19 @@ The Zenoh connector (053) is not implemented yet; it is written against default maximum payload for owned-serializer routes, the write ring size and its control reserve. Per-route rings cost more RAM than one shared ring; measure with the demo applications' route counts. -- **Ring implementation.** The record rings need a single-producer, - single-consumer bip buffer with grants and runtime-chosen sizes. Check - whether `bbqueue` fits (`no_std`, waker integration, storage sized at - build rather than by const generic) before writing one. +- **`bbqueue` stability.** 0.6 (January 2026) was a rewrite and 0.7 + followed in March; the author plans a few more minor breaking changes + before 1.0. Pin to `0.7` and budget for one upgrade. The pieces used here + (framed and stream grants, `BoxedSlice`, `AtomicCoord`, the `Notifier` + trait) are the core of the crate, not experimental features. +- **Write ring across reconnects.** `bbqueue` has no public `reset`. Either + the session drains the write ring through its consumer when a connection + ends, or each connection allocates a fresh write ring (one allocation per + reconnect, not per message). Draining is preferred on embedded. +- **Inbound chunk channel.** The embedded `read_into` side + (`Channel` into `PacketReader`) could use a stream `bbqueue` + too, reading socket bytes straight into a grant. Out of scope here; + noted for a follow-up. - **Other connectors' own copies.** KNX, WebSocket and the AimX session were not measured at connector level; §5's per-connector rows will show them. @@ -625,5 +684,8 @@ The Zenoh connector (053) is not implemented yet; it is written against - [055 — Wildcard inbound links](./055-wildcard-inbound-links.md) (grammar, keys, `TopicMatch`, §7 relation to this design) - `aimdb-bench/benches/b0_alloc_connector.rs` +- [`bbqueue`](https://github.com/jamesmunns/bbqueue) 0.7, and the + [lock-free ring buffer design](https://ferrous-systems.com/blog/lock-free-ring-buffer/) + it is based on - `aimdb-mqtt-connector/src/embedded/session_loop.rs` (`client_loop`, `drain_packets`, `perform`, `encode`, `write_out`) From dc42e3f3e351009fde289e76173340d92649aafb Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 2 Oct 2026 20:30:14 +0000 Subject: [PATCH 4/5] docs(054): SendBytes for copying transports; split embedded rework in two steps - Core ScratchPublisher + SendBytes: native MQTT, KNX, WebSocket and the AimX session implement send/send_owned; only embedded MQTT lends buffers. - Embedded MQTT migrates in two steps (write ring, then record rings), with step 1 as a fallback stopping point. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01EyVFXVWAec2aouRceoGdpE --- .../054-zero-alloc-connector-boundary.md | 77 +++++++++++++++---- 1 file changed, 64 insertions(+), 13 deletions(-) diff --git a/docs/design/054-zero-alloc-connector-boundary.md b/docs/design/054-zero-alloc-connector-boundary.md index b750cbbe..392618c6 100644 --- a/docs/design/054-zero-alloc-connector-boundary.md +++ b/docs/design/054-zero-alloc-connector-boundary.md @@ -329,8 +329,42 @@ pub trait RoutePublisher: Send { - **Shutdown.** When the reader ends (buffer closed), the pump awaits `poll_flush` once before dropping the publisher, so the last message is not silently abandoned. -- There is no adapter for the old `publish`. Every in-tree connector - implements `RoutePublisher` directly (§6). +- There is no adapter for the old `publish`. + +**Connectors that copy anyway: `SendBytes`.** Most transports copy the +message into their own client or socket buffer, so lending them storage +buys nothing. Core provides a `ScratchPublisher` that owns a per-route topic +and payload scratch (sized from the `RouteSpec`), implements +`RoutePublisher`, and on `commit` calls a much smaller trait: + +```rust +pub trait SendBytes: Send { + /// Transport readiness, e.g. client queue space. Default: always ready. + fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll> { + let _ = cx; + Poll::Ready(Ok(())) + } + /// Send one message, copying what the transport needs. + fn send(&mut self, topic: &str, payload: &[u8]) -> Result<(), PublishError>; + /// Send an owned serializer's payload (A6). Default: borrow and `send`. + /// Override when the transport takes ownership (native MQTT moves it in). + fn send_owned(&mut self, topic: &str, payload: Vec) -> Result<(), PublishError> { + self.send(topic, &payload) + } + fn poll_flush(&mut self, cx: &mut Context<'_>) -> Poll> { + let _ = cx; + Poll::Ready(Ok(())) + } +} + +// connector: Ok(Box::new(ScratchPublisher::new(route, MyTransport { .. }))) +``` + +Native MQTT, KNX, WebSocket and the AimX session client implement +`SendBytes`: today's `publish` body without the boxed future. Only the +embedded MQTT backend implements `RoutePublisher` directly, because only +there does lending remove a copy. This is the normal path for copying +transports, not a compatibility adapter. With 4.2–4.4 the pump's per-message loop is: @@ -516,15 +550,15 @@ through `TokioTcpDialer` (052). task; `MqttEventLoopSource` is removed. This removes both the topic clone and the `Arc` payload copy. - *Outbound.* `AsyncClient` takes an owned `String` and payload (G2). The - publisher lends a per-route topic and payload scratch from `buffers`; - `commit` builds the owned topic and copies the in-place payload, or moves a - `CommittedPayload::Owned` `Vec` straight in without copying. Two options - for the hand-over, to be settled when the native backend migrates: - 1. `commit` builds the `AsyncClient::publish` future and stores it in + backend implements `SendBytes` behind `ScratchPublisher`: `send` builds + the owned topic and copies the payload, `send_owned` moves an owned + serializer's `Vec` straight in without copying. Two options for the + hand-over, to be settled when the native backend migrates: + 1. `send` builds the `AsyncClient::publish` future and stores it in a per-route `tokio_util::sync::ReusableBoxFuture` (037's Tokio technique); `poll_ready` / `poll_flush` drive it. Keeps backpressure, removes the box, reports errors one message late (§4.4). - 2. `commit` calls the synchronous `AsyncClient::try_publish`. No + 2. `send` calls the synchronous `AsyncClient::try_publish`. No stored future, errors attributed immediately, but no readiness signal: a full request channel drops the message and counts it. @@ -575,11 +609,11 @@ so the work lands as one change (a feature branch merged once): | Crate | Inbound | Outbound | |---|---|---| -| `aimdb-core` | `InboundDispatch`; remove `Source`, `pump_source`, public `inbound_router`; AimX session client (TCP, UDS, serial) moves from `router.route` | poll `SerializedReader`, `TopicWriter`, `RoutePublisher`, new pump; remove `Connector::publish`, `TopicProvider` | +| `aimdb-core` | `InboundDispatch`; remove `Source`, `pump_source`, public `inbound_router`; AimX session client (TCP, UDS, serial) moves from `router.route` | poll `SerializedReader`, `TopicWriter`, `RoutePublisher`, `SendBytes` + `ScratchPublisher`, new pump; AimX session client implements `SendBytes`; remove `Connector::publish`, `TopicProvider` | | `aimdb-mqtt-connector` embedded | `dispatch` in `drain_packets` | per-route record rings + write ring (§4.6) | -| `aimdb-mqtt-connector` native | `dispatch` in event-loop task | `ReusableBoxFuture` publisher | -| `aimdb-knx-connector` | `Source` → `dispatch` | `RoutePublisher` | -| `aimdb-websocket-connector` | rename (`router.route` → `dispatch`), client and server | server `RoutePublisher` | +| `aimdb-mqtt-connector` native | `dispatch` in event-loop task | `SendBytes` with a `ReusableBoxFuture` | +| `aimdb-knx-connector` | `Source` → `dispatch` | `SendBytes` | +| `aimdb-websocket-connector` | rename (`router.route` → `dispatch`), client and server | server `SendBytes` | | `aimdb-embassy-adapter` | `send_wrapper.rs` wrappers for `Source` / `Connector` removed or reduced to `RoutePublisher` | — | | tests, bench | codec tests that build a router to ingest (`aimdb-data-contracts/src/link_codec.rs`, `aimdb-mqtt-connector/tests/link_ext_tests.rs`) → `InboundDispatch` (`route_count`, `subscriptions`, `dispatch`) | `topic_provider_tests.rs` (mqtt, knx), websocket e2e → `TopicWriter`; bench rewritten (§5) | @@ -588,7 +622,24 @@ so the work lands as one change (a feature branch merged once): 1. On main, independently: the CI gate for today's bench, and the native MQTT topic move (`publish.topic` instead of `.clone()`). 2. On the feature branch: core interfaces and pump, bench rewritten. -3. MQTT embedded (§4.6), then native. +3. MQTT embedded (§4.6), in two steps, each tested on its own with the + loopback harness: + 1. *Write ring.* Replace `encode`'s `Vec` and `Channel, 4>` with + the stream `bbqueue` write ring, draining on reconnect. The + `ActionChannel` stays; the embedded publisher is temporarily a + `SendBytes` that copies into an action. Touches `perform`, + `drain_packets`, `queue` and `write_out` only, and removes the + per-packet allocation. + 2. *Record rings.* Replace the `ActionChannel` with per-route framed + rings, the custom `Notifier` and round-robin, and switch the embedded + publisher to `RoutePublisher` with lent buffers. + + Step 2 is where the session loop's complexity grows. If it turns out + not to be worth it, stop after step 1: the core interfaces do not + change, and the embedded backend keeps one action allocation per + message until a later change (or uses a shared record ring with one + extra copy, alternative 9). + Then native MQTT. 4. KNX, WebSocket, Embassy adapter, data contracts, tests. 5. Merge once the gate and every connector's tests pass. From 64cbfdebd0ef6c34c79ebd49893bda24de8c7f8e Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 2 Oct 2026 20:36:52 +0000 Subject: [PATCH 5/5] docs(054): rewrite around a pull model Each connector's transport task drives both directions: it pushes received messages into InboundDispatch and pulls serialized messages from OutboundRoutes when its transport can send. Core runs no per-route pumps; the record buffers are the outbound queue. - Removes RoutePublisher, SendBytes, ScratchPublisher, the per-route record rings and the custom notifier from the design. - Embedded MQTT: the action arm pulls from OutboundRoutes under today's gate; a bbqueue write ring replaces the per-packet Vec. One copy, zero allocations, no lock. - Native MQTT: a plain publish loop awaiting AsyncClient; errors are attributed to the message that caused them. - Behaviour changes documented: outbound queueing moves into record buffers, one send at a time per connector, round-robin fairness. - The push design is kept as rejected alternative 4. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01EyVFXVWAec2aouRceoGdpE --- .../054-zero-alloc-connector-boundary.md | 776 +++++++----------- 1 file changed, 294 insertions(+), 482 deletions(-) diff --git a/docs/design/054-zero-alloc-connector-boundary.md b/docs/design/054-zero-alloc-connector-boundary.md index 392618c6..9cc3932a 100644 --- a/docs/design/054-zero-alloc-connector-boundary.md +++ b/docs/design/054-zero-alloc-connector-boundary.md @@ -1,22 +1,19 @@ # 054 — Zero-allocation connector boundary -**Status:** 📝 Proposed — revised 2026-10-02: no backward compatibility -(adapters, `Source` and `TopicProvider` removed; one migration); publishers -lend the buffers the reader serializes into (§4.2, §4.4); embedded MQTT -outbound uses one lock-free record ring per route and keeps the session loop -as the only encoder (§4.6). -Previous revision 2026-09-30, after 055 landed. +**Status:** 📝 Proposed — rewritten 2026-10-02 around a pull model: each +connector's own task drives both directions; core runs no per-route tasks. +Earlier revisions (per-route push pumps with `RoutePublisher`, per-route +record rings) are in this branch's history and summarised in §7. **Scope:** the per-message path between `aimdb-core` and connectors, in both -directions: `Source`/`pump_source`, `SerializedReader`, `Connector::publish`, -`TopicProvider`, and every in-tree connector. Breaking change to the -connector SPI and to one user-facing builder method: `with_topic_provider` -is replaced by `with_topic_writer` (§4.3). `link_from`, `link_to`, -`with_deserializer`, `with_match_deserializer`, `with_serializer`, -`with_serializer_into` and `Reader::recv` do not change. +directions, and every in-tree connector. Breaking change to the connector +SPI and to one user-facing builder method: `with_topic_provider` is replaced +by `with_topic_writer` (§4.3). `link_from`, `link_to`, `with_deserializer`, +`with_match_deserializer`, `with_serializer`, `with_serializer_into` and +`Reader::recv` do not change. -**Compatibility:** none kept. There are no adapters for the old traits; all -in-tree connectors migrate in one change (§6). +**Compatibility:** none kept. No adapters for the old traits; all in-tree +connectors migrate in one change (§6). **Builds on:** [037 — Zero-allocation consume path](./037-zero-alloc-consume-path.md), which made buffers and the consume path allocation-free and left the @@ -45,8 +42,8 @@ current-thread runtime, 2,000 messages after warm-up, no-op connector. Buffers, the consume path (037) and routing/ingest allocate nothing in steady state. The keyed-new row is by design: the first sighting of a key -interns its name (055 §5.6), and a known key costs nothing. Every other -non-zero row is caused by the connector interface: +interns its name (055 §5.6). Every other non-zero row is caused by the +connector interface: | # | Direction | Cause | Where | |---|---|---|---| @@ -65,14 +62,13 @@ at connector level): | Native (`rumqttc`) | topic `publish.topic.clone()` + payload `Arc::from(publish.payload.as_ref())` | `destination.to_string()` + `payload.to_vec()`, required by `AsyncClient::publish`'s owned API | | Embedded (`mountain-mqtt` codec, own session loop) | topic `topic_name.to_string()` + payload `Payload::from` | `destination.to_string()` + `payload.to_vec()` into `AimdbMqttAction::Publish`; then `session_loop::encode` allocates a `Vec` per packet for the `Channel, 4>` write queue | -Both native inbound copies disappear with §4.1. The embedded write queue is -a third per-message allocation that the interface change alone does not -remove (§4.6). - -Not every connector goes through `Source` today: the AimX session client -(TCP, UDS, serial) and both WebSocket builders already call -`Router::route` with borrows from their own frames. MQTT (both backends), -KNX and the Embassy adapter's glue use `Source` + `pump_source`. +**Structure today.** Every connector already owns at least one transport +task (MQTT event loop or session loop, KNX connection task, WebSocket +server, AimX session). Core adds one `pump_source` task per connector and +one `pump_sink` task per outbound route; the outbound pumps all funnel into +the same transport (one MQTT action channel, one socket). The AimX session +client and both WebSocket builders already call `Router::route` with +borrows instead of using `Source`. **Values with heap data.** Each reader receives its own clone of `T` (`T: Clone` delivery). A value with one `String` field costs one allocation @@ -81,8 +77,8 @@ row). This is a property of the buffer contract, not of the connector boundary; §4.5 covers it as guidance. **CI today.** The bench asserts its `EXPECTED` values when run, but CI only -builds it (`make` runs `cargo build --package aimdb-bench --benches`). Nothing -gates the numbers yet. +builds it (`make` runs `cargo build --package aimdb-bench --benches`). +Nothing gates the numbers yet. ## 2. Goals and non-goals @@ -94,33 +90,45 @@ gates the numbers yet. - G2. The MQTT connector adds no copies or allocations beyond what its client library's API requires, documented per backend. - G3. `b0_alloc_connector` runs in CI and gates the result. -- G4. One mechanism per direction: a single inbound entry point and a single - outbound publisher trait, with no compatibility layers. +- G4. The connector SPI is two objects, one per direction, and core runs no + tasks of its own for a connector. No compatibility layers. **Non-goals** - Topic grammars, wildcards and keys; 055 owns them and this design only - carries them through the new inbound entry point. + carries them through the inbound entry point. - Making `T` delivery cheaper for heap-containing types (§4.5 is guidance). - Removing owned serializers (A6). `with_serializer` stays: it is used far more than `with_serializer_into` in-tree (63 call sites against 7), and - requiring a declared capacity everywhere belongs in its own design. Users - who need zero allocations choose `with_serializer_into`. + requiring a declared capacity everywhere belongs in its own design. +- Parallel sends across routes of one connector. Every in-tree transport + serialises sends anyway (one socket, one MQTT in-flight slot); §4.6. - Allocations inside third-party client libraries. - Remote-access JSON paths (tracked by `b0_alloc_remote_json`). ## 3. Approach -The same move as 037: object safety and `async fn` conflict, object safety -and `poll` do not. Per-message boxed futures exist only to satisfy trait -signatures, so the signatures change to poll form. Owned values that exist -only to cross the interface become borrows. +The connector's transport task is the only task that touches its +transport. It already receives inbound messages; under this design it also +decides when to send. Core gives it two objects built from the database and +the connector's scheme: + +- `InboundDispatch`: the connector pushes each received message, as + borrows, and core deserializes and produces it synchronously. +- `OutboundRoutes`: when the transport can take a message, the connector + pulls the next ready one from any of its routes, already serialized into + storage `OutboundRoutes` owns. + +The record buffers (SPMC ring, mailbox, single-latest) are the outbound +queue. Nothing is queued between them and the transport, so no per-route +task, ring, channel or publisher trait is needed. Backpressure is the +connector not pulling; a record whose link falls behind lags exactly as a +slow reader does today. -Because no compatibility is kept, each direction gets exactly one shape: -inbound connectors push borrowed messages into `InboundDispatch`; outbound -connectors implement `RoutePublisher`. `Source`, `pump_source`, -`Connector::publish`, `TopicProvider` and the async `SerializedReader` -methods are removed, not adapted. +Per-message boxed futures disappear because no trait signature returns a +future any more: the connector's loop is concrete code and awaits its own +client directly. Owned values that only existed to cross the interface +become borrows. ## 4. Design @@ -128,12 +136,9 @@ methods are removed, not adapted. Connectors hold each received message by reference: inside `rumqttc`'s `Event::Incoming(Publish)`, inside the embedded session loop's decoded -`ApplicationMessage`, inside a decoded AimX or WebSocket frame. -`Router::route(&str, &[u8], ctx)` already takes borrows and, after 055, -matches through the connector's grammar without allocating. The `Source` -interface in between is what forces owned copies (A1, A2). - -Core exposes the router to connectors through one type: +publish, inside a decoded AimX or WebSocket frame. `Router::route(&str, +&[u8], ctx)` already takes borrows and, after 055, matches through the +connector's grammar without allocating. ```rust /// Built once per connector from the inbound links of its scheme. @@ -159,88 +164,78 @@ impl InboundDispatch { } ``` -- Connectors call `dispatch` wherever they hold the borrow. No future, no - copy, no channel. `TopicMatch<'a>` already borrows the topic, so - match-aware deserializers are unchanged (055 §7). -- `InboundDispatch` is the only inbound entry point. `Source`, - `pump_source` and `AimDb::inbound_router` are removed; `Router` becomes - crate-private. A transport that produces owned frames calls `dispatch` - with borrows of its own frame, so A2 no longer exists as a category. -- For the AimX session client and the WebSocket connectors this is a rename: - they already call `router.route(topic, payload, &ctx)`. -- `InboundDispatch` is cheap to share (`Arc` inside or `Clone`), since the +- `TopicMatch<'a>` already borrows the topic, so match-aware deserializers + are unchanged (055 §7). +- `Source`, `pump_source` and `AimDb::inbound_router` are removed; `Router` + becomes crate-private. A transport that produces owned frames calls + `dispatch` with borrows of its own frame, so A1 and A2 disappear. +- For the AimX session client and the WebSocket connectors this is a + rename. `InboundDispatch` is `Clone` (shared `Arc` inside), since the WebSocket server dispatches from one task per client. -**Behaviour changes.** - -- *Embedded MQTT: burst absorption.* Today the session loop parks on - `events.send(event).await` when `pump_source` falls behind, so the - `EventChannel` (`CHANNEL_SIZE` messages) absorbs bursts and throttles the - socket. With `dispatch`, a burst goes straight to the record buffers, - which drop on overflow as they already do for a slow consumer. Records - that need burst tolerance size their buffer for it. -- *QoS 1 inbound is acknowledge-then-maybe-drop.* `drain_packets` queues - the PUBACK before the message is delivered, and the broker does not - resend an acknowledged message. Today an acknowledged message has at least - reached the event channel; with `dispatch` it can be dropped at once by a - full record buffer. QoS 1 at the transport therefore means "reached - AimDB", not "reached every record". This is documented on the MQTT - connector; moving the PUBACK after `dispatch` does not change it, because - `dispatch` does not report per-record drops. -- *Ingest runs on the transport task.* For embedded MQTT that is the session - loop; for native MQTT it is the task polling `rumqttc`'s event loop. A - slow user deserializer delays keep-alive and ack handling in both (§8). - -### 4.2 Outbound: the publisher lends the buffers - -The serialized topic and payload are written straight into storage the -publisher owns. For the embedded MQTT backend that storage is a slot in the -route's record ring (§4.6), so the bytes are never copied between -serialization and encoding. For other connectors it is a per-route scratch -the publisher allocates once, which replaces the pump's scratch of today. +### 4.2 Outbound: connectors pull serialized messages ```rust -/// Where the reader writes one message. Borrowed from the publisher. -pub struct OutboundBuffers<'a> { - pub topic: TopicBuf<'a>, // capacity = the route's topic capacity (§4.3) - pub payload: &'a mut [u8], // capacity = the route's payload capacity +/// Every outbound link of one scheme, read by the connector's transport task. +pub struct OutboundRoutes { /* routes, round-robin cursor, scratch buffers */ } + +pub struct RouteInfo { + pub id: RouteId, // dense index, 0..len + pub default_topic: Arc, // from the link URL + pub config: ConnectorConfig, // qos, retain, record_index, … } -/// What the reader wrote. -pub struct OutboundFrame { - /// `Some(len)`: a `TopicWriter` wrote `len` bytes into `buffers.topic`. - /// `None`: use the route's default topic. - pub topic: Option, - pub payload: CommittedPayload, +pub struct Outbound<'a> { + pub route: &'a RouteInfo, + /// The written topic (`TopicWriter`) or the route's default. + pub topic: &'a str, + pub payload: OutboundPayload<'a>, } -pub enum CommittedPayload { - /// `len` bytes written into `buffers.payload` (`with_serializer_into`). - InPlace(usize), - /// An owned serializer's bytes (`with_serializer`, A6). The publisher - /// moves or copies them as its transport needs. +pub enum OutboundPayload<'a> { + /// Serialized into `OutboundRoutes`' scratch (`with_serializer_into`). + Borrowed(&'a [u8]), + /// An owned serializer's bytes (`with_serializer`, A6). Transports that + /// take ownership move it in; others borrow it. Owned(Vec), } -pub trait SerializedReader: Send { - fn poll_recv_into( - &mut self, - cx: &mut Context<'_>, - ctx: &RuntimeContext, - buffers: OutboundBuffers<'_>, - ) -> Poll>; +impl OutboundRoutes { + /// Subscribes every outbound link of `scheme`; each route's cursor starts + /// here, so build it when the transport task starts. Allocates the + /// scratch once, sized to the largest route's topic and payload capacity. + pub fn new(db: &AimDb, scheme: &str) -> DbResult; + /// Routes, for parsing per-route configuration once at start. + pub fn routes(&self) -> &[RouteInfo]; + /// The next ready message from any route, round-robin. + /// `Ready(None)`: every route is closed. + pub fn poll_next(&mut self, cx: &mut Context<'_>) -> Poll>>; + /// `poll_fn(|cx| self.poll_next(cx))`. + pub async fn next(&mut self) -> Option>; } ``` -- Removes A3. The reader runs the topic writer and the serializer only when - a value is ready, so it writes into `buffers` at most once per `Ready`. -- `CommittedPayload::Owned` keeps owned serializers working (A6, §2 - non-goals). It is also the cheaper path for transports whose client takes - an owned payload: the native MQTT publisher moves the `Vec` into - `rumqttc` instead of copying it (§4.6). -- The async `recv` / `recv_into` pair, `RecvSerializedIntoFuture` and - `SerializedPayload` are removed. The only implementors are in this - repository (the fused reader in `typed_api.rs` and the bench). +**Semantics of `poll_next`.** + +- Starting after the route that produced the last message, it polls each + open route's typed reader (`poll_recv`, 037). The first value that is + ready is serialized (topic writer, then serializer) into the scratch and + returned. Pending readers register `cx`, so any of them waking wakes the + transport task. +- A value is taken from its buffer only in a call that returns `Ready`. + `poll_next` is therefore safe as a `select` arm: a losing arm takes + nothing. +- `BufferLagged` is logged and that route is polled again; a closed buffer + closes the route; topic overflow and serializer errors skip the value, + log, and count per link. None of these end the other routes. +- One message is borrowed at a time (`&mut self`). The connector finishes + with it (copies, encodes, or awaits its send) before pulling the next. +- Cost of a wake-up: one `poll_recv` per open route. Fine for tens of + routes; §8 notes a ready set for many more. + +`SerializedReader`, `SerializedPayload`, `RecvSerializedIntoFuture` and +`pump_sink` are removed. The typed reader + serializer pair stays inside +core as the per-route state of `OutboundRoutes`. ### 4.3 Outbound: topics are written, not returned @@ -252,146 +247,51 @@ pub trait TopicWriter: Send + Sync { fn write_topic(&self, value: &T, out: &mut TopicBuf<'_>) -> Result; } -// TopicBuf implements core::fmt::Write over a byte slice the publisher lends -// (§4.2) and refuses any write past its end. Only `&str` goes in, so the +// TopicBuf implements core::fmt::Write over the scratch topic region and +// refuses any write past the route's capacity. Only `&str` goes in, so the // bytes are valid UTF-8 by construction. // usage: write!(out, "sensors/{}/{}", v.site, v.id)?; Ok(true) ``` - `.with_topic_writer(capacity, writer)` on the outbound builder replaces - `.with_topic_provider(provider)`; capacity is the topic space the publisher - reserves per message on that route. This is the one user-facing break in + `.with_topic_provider(provider)`. This is the one user-facing break in this design. - An overflow skips the message and increments a per-link counter; the topic is never truncated. -- Closures implement `TopicWriter` through a blanket impl, so a one-line - provider stays a one-line writer. +- Closures implement `TopicWriter` through a blanket impl. - Removes A5. -### 4.4 Outbound: per-route publishers in poll form +### 4.4 Connector shape -`Connector::publish` is shared by all routes and returns a boxed future -(A4). It is replaced by a publisher created once per route: +A connector builds both objects when its transport task starts and drives +them from that task. For a client with an async send, the outbound half is +a plain loop: ```rust -/// Everything a publisher needs to know about its route, fixed at build. -pub struct RouteSpec<'a> { - pub default_topic: &'a str, - pub config: &'a ConnectorConfig, // qos, retain, … from the link URL - pub topic_capacity: usize, // 0 without a TopicWriter - /// `Some(n)`: into-slice serializer with an n-byte scratch. - /// `None`: owned serializer; the publisher picks its own limit. - pub payload_capacity: Option, -} +let mut outbound = OutboundRoutes::new(&db, "mqtt")?; +let opts: Vec = outbound.routes().iter().map(PublishOpts::parse).collect::>()?; -pub trait Connector: Send + Sync { - /// Called once per outbound route while the connector builds its pumps, - /// before its transport tasks start. - fn route_publisher(&self, route: &RouteSpec<'_>) - -> Result, PublishError>; -} - -pub trait RoutePublisher: Send { - /// Ready to accept one message: storage for it is reserved. Also drives - /// any hand-over still in progress from the previous `commit`. - fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll>; - /// The storage reserved by the last successful `poll_ready`. Calling it - /// again before `commit` returns the same storage. - fn buffers(&mut self) -> OutboundBuffers<'_>; - /// The message written into `buffers` is complete. Synchronous. - fn commit(&mut self, frame: OutboundFrame) -> Result<(), PublishError>; - /// Complete every hand-over committed so far. Default: `Ready(Ok(()))`. - fn poll_flush(&mut self, cx: &mut Context<'_>) -> Poll> { - let _ = cx; - Poll::Ready(Ok(())) +while let Some(msg) = outbound.next().await { + let o = &opts[msg.route.id]; + if let Err(e) = client.publish(msg.topic, o.qos, o.retain, msg.payload).await { + log_error!("publish to '{}' failed: {:?}", msg.topic, e); } } ``` -- This is the `futures::Sink` shape with the buffer lent out instead of - passed in (the grant pattern of `bbqueue`): backpressure through - `poll_ready`, a synchronous hand-over, no future per message, no copy into - the publisher. -- A reservation is held while the reader waits for the next value. That is - why it belongs to one route: a publisher must not reserve in storage other - routes share (§4.6 gives each route its own ring). -- `route_publisher` parses configuration (`qos`, `retain`, …) once, not per - message as the MQTT sinks do today, and rejects a route it cannot serve - (invalid configuration, a message that can never fit) before anything - runs. It returns a `Result` for that reason. -- **Errors never end a route.** A publisher that completes a hand-over - asynchronously reports its failure from the next `poll_ready` or - `poll_flush`, i.e. one message late, and clears it once reported. The - publisher logs the failed destination itself; the pump logs that the route - saw an error and continues, as `pump_sink` does today. A message that does - not fit (topic overflow, an owned payload larger than the publisher's - limit) is skipped and counted, and the reservation is reused. -- **Shutdown.** When the reader ends (buffer closed), the pump awaits - `poll_flush` once before dropping the publisher, so the last message is - not silently abandoned. -- There is no adapter for the old `publish`. - -**Connectors that copy anyway: `SendBytes`.** Most transports copy the -message into their own client or socket buffer, so lending them storage -buys nothing. Core provides a `ScratchPublisher` that owns a per-route topic -and payload scratch (sized from the `RouteSpec`), implements -`RoutePublisher`, and on `commit` calls a much smaller trait: +- The `.await` is concrete code in the connector's own task: no boxed + future, no stored future, errors attributed to the message that caused + them. +- Configuration is parsed once per route at start; an invalid link is + reported before anything runs. +- A connector that must send while also reading (one socket, one task) + puts `poll_next` in its `select`, gated by its own readiness (§4.7). +- Shutdown: when `next` returns `None` every outbound record is gone; the + connector flushes its transport as it sees fit. -```rust -pub trait SendBytes: Send { - /// Transport readiness, e.g. client queue space. Default: always ready. - fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll> { - let _ = cx; - Poll::Ready(Ok(())) - } - /// Send one message, copying what the transport needs. - fn send(&mut self, topic: &str, payload: &[u8]) -> Result<(), PublishError>; - /// Send an owned serializer's payload (A6). Default: borrow and `send`. - /// Override when the transport takes ownership (native MQTT moves it in). - fn send_owned(&mut self, topic: &str, payload: Vec) -> Result<(), PublishError> { - self.send(topic, &payload) - } - fn poll_flush(&mut self, cx: &mut Context<'_>) -> Poll> { - let _ = cx; - Poll::Ready(Ok(())) - } -} - -// connector: Ok(Box::new(ScratchPublisher::new(route, MyTransport { .. }))) -``` - -Native MQTT, KNX, WebSocket and the AimX session client implement -`SendBytes`: today's `publish` body without the boxed future. Only the -embedded MQTT backend implements `RoutePublisher` directly, because only -there does lending remove a copy. This is the normal path for copying -transports, not a compatibility adapter. - -With 4.2–4.4 the pump's per-message loop is: - -```rust -loop { - // Readiness first (Sink order): a blocked transport shows up as buffer - // lag, not as one stale message held in the pump. An error here belongs - // to an earlier message. - if let Err(e) = poll_fn(|cx| publisher.poll_ready(cx)).await { - log_warn!("route '{}': earlier publish failed: {:?}", default_topic, e); - continue; - } - // Serialize straight into the publisher's reserved storage. - let frame = match poll_fn(|cx| reader.poll_recv_into(cx, &ctx, publisher.buffers())).await { - Ok(frame) => frame, - Err(DbError::BufferLagged { .. }) => continue, // reservation kept - Err(_) => break, - }; - if let Err(e) = publisher.commit(frame) { - log_error!("route '{}': publish failed: {:?}", default_topic, e); - } -} -let _ = poll_fn(|cx| publisher.poll_flush(cx)).await; -``` - -The pump no longer owns any per-route buffer; `RouteSpec` carries the sizes -the old pump scratch was built from. +Removed from the SPI: `Connector::publish`, `Source`, `pump_source`, +`pump_sink`, `SerializedReader` and `TopicProvider`. The `Vec` +connectors return today shrinks to their own transport tasks. ### 4.5 Values with heap data (guidance) @@ -405,7 +305,40 @@ at message rate: Document this next to the buffer types; no API change. -### 4.6 MQTT connector +### 4.6 Behaviour changes + +- **Outbound queueing moves into the record buffers.** Today a message + leaves its buffer as soon as its route's pump runs and waits in the + connector (MQTT action channel, `CHANNEL_SIZE` deep) until the transport + can send. Now it stays in the record buffer until the connector pulls. + During an outage or a slow link, what is sent afterwards depends on the + record's buffer: single-latest sends only the newest value, a mailbox its + last, an SPMC ring its backlog or a lag report. Document this per buffer + type on the link builder. +- **One send at a time per connector.** Routes of one connector share the + transport task, so a slow send delays all of them. That is already true + underneath for every in-tree transport; the per-route pumps only hid it. + A future transport with genuinely parallel sends spawns its own workers + from its loop. +- **Fairness.** Round-robin over ready routes; a busy route cannot starve + the others. Priorities are not in scope. +- **Ingest and serialization run on the transport task.** Inbound + deserialization already did for `pump_source`'s connectors; outbound + serialization now does too. A slow user (de)serializer delays the + transport's keep-alive and ack handling (§8). +- **Embedded MQTT: burst absorption.** Today the session loop parks on + `events.send(event).await` when `pump_source` falls behind, so the + `EventChannel` absorbs bursts and throttles the socket. With `dispatch`, + a burst goes straight to the record buffers, which drop on overflow as + they already do for a slow consumer. +- **QoS 1 inbound is acknowledge-then-maybe-drop.** `drain_packets` queues + the PUBACK before the message is delivered, and the broker does not + resend an acknowledged message. With `dispatch` an acknowledged message + can be dropped at once by a full record buffer. QoS 1 at the transport + means "reached AimDB", not "reached every record"; documented on the MQTT + connector. + +### 4.7 MQTT connector #### Embedded backend (`mountain-mqtt` codec, AimDB's own `session_loop`) @@ -413,182 +346,88 @@ Document this next to the buffer types; no API change. publish.payload())` where it decodes the publish today. `EventChannel`, `AimdbMqttEvent`, `Received::Event` and `MqttSource` are removed. -*Outbound.* The session loop's client state (`ClientStateNoQueue`) decides -when a PUBLISH may go out, and it is the only thing that may assign packet -ids: +*Outbound.* The session loop's action arm pulls from `OutboundRoutes` +instead of the `ActionChannel`, under the same gate as today: ```rust -// session_loop.rs today: the action arm is armed only when let action_ready = connected - && !state.waiting_for_responses() // one in-flight slot + && !state.waiting_for_responses() // one in-flight slot && next_topic >= subscribe_topics.len() - && !outbound.is_full(); -``` - -So route publishers do not encode packets. Each route serializes into its -own record ring; the session loop stays the only encoder and writes into -one write ring: - -```text -route 1 pump ──serialize in place──▶ record ring 1 ─┐ -route 2 pump ──serialize in place──▶ record ring 2 ─┼─▶ session loop ──encode (1 copy)──▶ write ring ──▶ write_out ──▶ socket -route N pump ──serialize in place──▶ record ring N ─┘ (sole owner of (per connection, - (one producer each, lock-free, ClientStateNoQueue, bip buffer) - survives reconnects) same gate as today) + && write_ring.has_room(max_frame + CONTROL_RESERVE); +let action_arm = poll_fn(|cx| { + if action_ready { outbound.poll_next(cx) } else { Poll::Pending } +}); +// on Some(msg): +// packet = state.publish_packet(msg.topic, payload, qos[msg.route.id], retain[..]) +// encode packet into the write ring (the one copy), state.publish_update(&packet) ``` -- **Record rings** (replace `ActionChannel`). One per route, created in - `route_publisher` while the connector builds its pumps, so the session - loop starts with a fixed list of ring consumers. Each ring is a - single-producer, single-consumer bip buffer of variable-length records - `{topic_len, payload_len, topic, payload}`; `qos`, `retain` and the - default topic are stored once per ring, not per record. - - `poll_ready` reserves one contiguous slot of `topic_capacity + - payload_capacity` (plus a small record header) and waits if the ring has - no room. - - `buffers` lends the slot's topic and payload regions; the serializer - and `TopicWriter` write directly into them. - - `commit` writes the record header with the real lengths and releases the - unused tail of the slot. `CommittedPayload::Owned` is copied into the - slot (A6 keeps its copy); an owned payload larger than the slot is - skipped and counted. - - Records queued during an outage stay in their rings and are sent after - reconnecting, as actions are today. -- **Synchronisation.** No lock in the per-message path. - - Each record ring has exactly one producer (its route's pump) and one - consumer (the session loop). The ring is split into a producer half, - owned by the `RoutePublisher`, and a consumer half, owned by the - session loop, so the type system enforces it. Head and tail are - atomics. - - The write ring has one writer (the session loop) and one reader - (`write_out`), which run in the same task under one `select3`. - - `ClientStateNoQueue` and the list of record rings belong to the session - loop alone. The list is fixed before the loop starts; a reconnect starts - a new loop only after the old one has ended. - - Wakers are the only shared state left: each ring wakes its own pump when - space is released, and every ring's `commit` wakes the session loop. - Those are atomic waker slots, not locks held while data is copied. - - Atomics: the rings need atomic read-modify-write (`swap`, - `fetch_add`), which `thumbv7em`, `thumbv8m.main` and std targets have. - On targets without it (`thumbv6m`), the rings and wakers fall back to - `portable-atomic` or a critical section. -- **Session loop.** Where `perform` runs today, when `action_ready` holds - the loop takes the oldest committed record from the rings in round-robin - order, checks the write ring has room for its encoded size - (`MqttLenWriter`), calls `state.publish_packet(topic, payload, qos, - retain)` with slices borrowed from the record ring, encodes the packet - straight into the write ring (`MqttBufWriter` over the reserved slot), - calls `state.publish_update`, then releases the record. Round-robin keeps - one busy route from starving the others; packet ids, the one-in-flight - rule and the connect/subscribe ordering are untouched. -- **Write ring** (replaces `Channel, 4>` and `encode`'s `Vec`). One - per connection, cleared on reconnect so no bytes from an old session - reach a new socket. All packets go through it (CONNECT, SUBSCRIBE, - PUBLISH, PUBACK, PINGREQ). `MqttBufWriter` needs a contiguous slice, so - it is a bip buffer as well. `write_out` writes contiguous slices to the - socket and releases them. -- **Control packets.** The session loop encodes a PUBLISH only when the - write ring keeps a reserve for control packets (PUBACK, PINGREQ), so a - stalled socket full of publishes cannot block an acknowledgement. PINGREQ - stays lossy, as today. -- **Implementation: `bbqueue` 0.7.** Both ring kinds are `bbqueue` - queues (single-producer single-consumer bip buffers with grants, Miri - tested): - - *Record rings* use the framed interface with a `u16` length header. - `poll_ready` calls `FramedProducer::grant(slot_len)` and stores the - returned `FramedGrantW` in the publisher; `buffers` lends regions of - it; `commit` calls `FramedGrantW::commit(used)`, which writes the header - and returns the unused tail to the ring. The session loop's - `FramedConsumer::read` yields one record; if the write ring has no room - yet it calls `keep()` (the record stays queued), otherwise it encodes - and calls `release()`. Records are limited to 64 KiB. - - *Write ring* uses the stream interface. The session loop takes - `grant_exact(frame_len + control_reserve)` and commits only - `frame_len`, which keeps the reserve free without any accounting of its - own; `write_out` takes `read()` grants and releases them after - `write_all`. - - Storage is `BoxedSlice` (feature `alloc`): each ring is allocated once - at build, sized from its `RouteSpec`. Queues are held as - `ArcBBQueue`, so the halves carry no lifetimes; taking a grant clones - the `Arc` (a counter increment, no allocation). - - Coordination is `AtomicCoord` (lock-free) where the target has atomic - pointers; `CsCoord` or the `portable-atomic` features otherwise. - - Notifier: `bbqueue`'s async support is an `AsyncNotifier` with `async - fn`s backed by `maitake-sync`, which does not fit `poll_ready`. AimDB - implements the public `Notifier` trait instead, with two atomic waker - slots per ring: `wake_one_producer` wakes the route's pump, - `wake_one_consumer` wakes the session loop. `poll_ready` registers its - waker, then tries `grant`, and returns `Pending` on - `InsufficientSize`; the session loop does the same for `read`. - - Dependency: `bbqueue = { version = "0.7", default-features = false, - features = ["alloc"] }`, which leaves out `maitake-sync` and - `critical-section`. - - Record layout inside a frame: `[topic_len: u16][topic][payload]`. The - publisher lends the payload region at a fixed offset after the route's - topic capacity, so a route with a `TopicWriter` leaves up to - `topic_capacity` unused bytes per queued record. If that matters, the - reader can serialize the payload first and write the topic after it - (`[payload_len][topic_len][payload][topic]`), at the cost of a less - direct `OutboundBuffers` contract. -- **Sizes and oversize routes.** A route's record ring holds at least two - slots, so the pump can fill the next record while the loop encodes the - previous one. For an owned serializer (`payload_capacity: None`), the slot - uses the connector's maximum payload, a builder option. `route_publisher` - rejects a route whose encoded frame cannot fit in the write ring minus the - control reserve. - -Cost per message: zero allocations and one copy of topic and payload (record -ring into the encoded frame), in addition to the network stack's own copy -into its TX buffer, which TLS needs anyway to encrypt. The remaining copy is -the price of keeping packet encoding inside the session loop (§7, -alternative 8). This also gives std users an allocation-free MQTT path -through `TokioTcpDialer` (052). +- `MqttSink`, `AimdbMqttAction::Publish` and the `ActionChannel` are + removed. Subscriptions keep their current path (the loop already places + them itself). +- **Write ring** (replaces `encode`'s `Vec` and `Channel, 4>`). A + `bbqueue` 0.7 stream queue, allocated once at build. The session loop + writes every packet (CONNECT, SUBSCRIBE, PUBLISH, PUBACK, PINGREQ) with + `grant_exact`, sizing it with `MqttLenWriter` and encoding with + `MqttBufWriter` over the grant; `write_out` takes `read()` grants and + releases them after `write_all`. Both ends run in the same task under the + existing `select3`, so the ring needs no lock. +- **Control reserve.** For a PUBLISH the loop takes `grant_exact(frame_len + + CONTROL_RESERVE)` and commits only `frame_len`, so PUBACK and PINGREQ + always find room. PINGREQ stays lossy, as today. +- **Reconnect.** When a connection ends, the loop drains the write ring + through its consumer before the next CONNECT, so no bytes from an old + session reach a new socket. Unsent outbound values are still in their + record buffers (§4.6). +- **Oversize.** `max_frame` comes from the largest route's capacities; a + route whose frame cannot fit in the write ring minus the reserve is + rejected when the connector builds. An owned payload larger than the + write ring allows is skipped and counted. +- `bbqueue = { version = "0.7", default-features = false, features = + ["alloc"] }`; `AtomicCoord` where the target has atomic pointers, + `portable-atomic` features on `thumbv6m`. Only the polling notifier is + needed: `write_out` and the session loop share a task, and the loop wakes + `write_out` with a waker it already holds. + +Cost per message: zero allocations and one copy of topic and payload +(scratch into the encoded frame), besides the network stack's own copy into +its TX buffer, which TLS needs anyway to encrypt. This also gives std users +an allocation-free MQTT path through `TokioTcpDialer` (052). #### Native backend (`rumqttc`) - *Inbound.* `dispatch(&publish.topic, &publish.payload)` in the event-loop task; `MqttEventLoopSource` is removed. This removes both the topic clone and the `Arc` payload copy. -- *Outbound.* `AsyncClient` takes an owned `String` and payload (G2). The - backend implements `SendBytes` behind `ScratchPublisher`: `send` builds - the owned topic and copies the payload, `send_owned` moves an owned - serializer's `Vec` straight in without copying. Two options for the - hand-over, to be settled when the native backend migrates: - 1. `send` builds the `AsyncClient::publish` future and stores it in - a per-route `tokio_util::sync::ReusableBoxFuture` (037's Tokio - technique); `poll_ready` / `poll_flush` drive it. Keeps backpressure, - removes the box, reports errors one message late (§4.4). - 2. `send` calls the synchronous `AsyncClient::try_publish`. No - stored future, errors attributed immediately, but no readiness signal: - a full request channel drops the message and counts it. - - Option 1 matches today's behaviour (`publish` waits for channel space) - and is the default. +- *Outbound.* One task runs the loop of §4.4 over `AsyncClient::publish`. + `rumqttc` takes an owned `String` and payload (G2): the topic is copied, + a borrowed payload is copied, an `OutboundPayload::Owned` `Vec` is moved + in. `publish` waits for request-channel space, which is the backpressure. - `rumqttc`'s own allocations (building `Publish`, its request channel) are outside G1. ## 5. Measurement and gate -`b0_alloc_connector` is rewritten against the new interfaces. Targets after -this design: +`b0_alloc_connector` is rewritten against the new interfaces. Targets: | Row | Before | After | |---|---|---| | `inbound_dispatch` (was `inbound_route`) | 0 | 0 | -| `inbound_dispatch_pattern` (was `inbound_route_pattern`) | 0 | 0 | +| `inbound_dispatch_pattern` | 0 | 0 | | `inbound_dispatch_keyed_known` | 0 | 0 | | `inbound_dispatch_keyed_new` | 1 | 1 (by design, 055 §5.6) | | `inbound_pump_source_minimal` | 2 | removed (`Source` is gone) | -| `outbound_scratch_static_topic` | 2 | 0 | -| `outbound_scratch_written_topic` (replaces `outbound_scratch_dynamic_topic`) | 3 | 0 | -| `outbound_owned_static_topic` | 3 | 1 (A6, by choice of serializer) | +| `outbound_next_static_topic` (was `outbound_scratch_static_topic`) | 2 | 0 | +| `outbound_next_written_topic` (replaces the `TopicProvider` row) | 3 | 0 | +| `outbound_next_owned` (was `outbound_owned_static_topic`) | 3 | 1 (A6, by choice of serializer) | +| `outbound_next_round_robin` (new; 8 routes, all ready) | — | 0 | The bench asserts exact values, so every change updates `EXPECTED` and the baseline together. A connector-level row per MQTT backend (loopback broker -for native, the existing loopback harness for embedded) records G2; the +for native, the existing loopback harness for embedded) records G2: the embedded rows target 0 allocations in both directions and record the one -outbound copy (§4.6); the native outbound row records the owned topic, the -payload copy and whatever `rumqttc` adds. +outbound copy; the native outbound row records the owned topic, the payload +copy and whatever `rumqttc` adds. **Gate.** A CI job runs `cargo bench -p aimdb-bench --bench b0_alloc_connector` (host only, a few seconds) and fails on any assertion. @@ -601,128 +440,103 @@ deterministic and do not need a quiet runner. | Who | Change | |---|---| | Application code | `.with_topic_provider(p)` → `.with_topic_writer(capacity, w)`. Nothing else. | -| Connector authors | `inbound_router` + `pump_source` / `router.route` → `InboundDispatch`; `Connector::publish` → `route_publisher(&RouteSpec)` + `RoutePublisher` (`poll_ready`, `buffers`, `commit`, `poll_flush`); the publisher owns the per-route buffers the pump used to own. No adapters. | -| Custom `SerializedReader` implementors | `recv` / `recv_into` → `poll_recv_into` (none known outside this repository) | +| Connector authors | Build `InboundDispatch` and `OutboundRoutes` in the transport task; call `dispatch` on receive, pull with `next`/`poll_next` when the transport can send. `Source`, `pump_source`, `pump_sink`, `Connector::publish` are gone. | -Without adapters, changing the core traits breaks every connector at once, -so the work lands as one change (a feature branch merged once): +Without adapters the core change breaks every connector at once, so the +work lands as one change (a feature branch merged once): | Crate | Inbound | Outbound | |---|---|---| -| `aimdb-core` | `InboundDispatch`; remove `Source`, `pump_source`, public `inbound_router`; AimX session client (TCP, UDS, serial) moves from `router.route` | poll `SerializedReader`, `TopicWriter`, `RoutePublisher`, `SendBytes` + `ScratchPublisher`, new pump; AimX session client implements `SendBytes`; remove `Connector::publish`, `TopicProvider` | -| `aimdb-mqtt-connector` embedded | `dispatch` in `drain_packets` | per-route record rings + write ring (§4.6) | -| `aimdb-mqtt-connector` native | `dispatch` in event-loop task | `SendBytes` with a `ReusableBoxFuture` | -| `aimdb-knx-connector` | `Source` → `dispatch` | `SendBytes` | -| `aimdb-websocket-connector` | rename (`router.route` → `dispatch`), client and server | server `SendBytes` | -| `aimdb-embassy-adapter` | `send_wrapper.rs` wrappers for `Source` / `Connector` removed or reduced to `RoutePublisher` | — | -| tests, bench | codec tests that build a router to ingest (`aimdb-data-contracts/src/link_codec.rs`, `aimdb-mqtt-connector/tests/link_ext_tests.rs`) → `InboundDispatch` (`route_count`, `subscriptions`, `dispatch`) | `topic_provider_tests.rs` (mqtt, knx), websocket e2e → `TopicWriter`; bench rewritten (§5) | +| `aimdb-core` | `InboundDispatch`; remove `Source`, `pump_source`, public `inbound_router`; AimX session client moves from `router.route` | `OutboundRoutes`, `TopicWriter`; remove `pump_sink`, `Connector::publish`, `SerializedReader`, `TopicProvider`; AimX session client pulls in its session task | +| `aimdb-mqtt-connector` embedded | `dispatch` in `drain_packets` | action arm pulls from `OutboundRoutes`; `bbqueue` write ring | +| `aimdb-mqtt-connector` native | `dispatch` in event-loop task | publish loop task over `AsyncClient` | +| `aimdb-knx-connector` | `Source` → `dispatch` | pull in its connection task | +| `aimdb-websocket-connector` | rename (`router.route` → `dispatch`), client and server | server: broadcast loop pulls (`record_index` from `RouteInfo.config`) | +| `aimdb-embassy-adapter` | `send_wrapper.rs` wrappers for `Source` / `Connector` removed | — | +| tests, bench | codec tests that build a router (`aimdb-data-contracts/src/link_codec.rs`, `aimdb-mqtt-connector/tests/link_ext_tests.rs`) → `InboundDispatch` | `topic_provider_tests.rs` (mqtt, knx), websocket e2e → `TopicWriter`; KNX task-count tests drop the per-route publisher; bench rewritten (§5) | **Order:** 1. On main, independently: the CI gate for today's bench, and the native MQTT topic move (`publish.topic` instead of `.clone()`). -2. On the feature branch: core interfaces and pump, bench rewritten. -3. MQTT embedded (§4.6), in two steps, each tested on its own with the - loopback harness: - 1. *Write ring.* Replace `encode`'s `Vec` and `Channel, 4>` with - the stream `bbqueue` write ring, draining on reconnect. The - `ActionChannel` stays; the embedded publisher is temporarily a - `SendBytes` that copies into an action. Touches `perform`, - `drain_packets`, `queue` and `write_out` only, and removes the - per-packet allocation. - 2. *Record rings.* Replace the `ActionChannel` with per-route framed - rings, the custom `Notifier` and round-robin, and switch the embedded - publisher to `RoutePublisher` with lent buffers. - - Step 2 is where the session loop's complexity grows. If it turns out - not to be worth it, stop after step 1: the core interfaces do not - change, and the embedded backend keeps one action allocation per - message until a later change (or uses a shared record ring with one - extra copy, alternative 9). - Then native MQTT. -4. KNX, WebSocket, Embassy adapter, data contracts, tests. +2. On the feature branch: `InboundDispatch`, `OutboundRoutes`, + `TopicWriter`; bench rewritten. +3. MQTT embedded in two steps, each tested with the loopback harness: + 1. the `bbqueue` write ring, replacing `encode`'s `Vec` and the + `Vec` channel (touches `perform`, `drain_packets`, `queue`, + `write_out`); + 2. the action arm pulls from `OutboundRoutes`; `ActionChannel` and + `MqttSink` removed. +4. MQTT native, KNX, WebSocket, AimX session client, Embassy adapter, + tests. 5. Merge once the gate and every connector's tests pass. The Zenoh connector (053) is not implemented yet; it is written against -`InboundDispatch` and `RoutePublisher` from the start. +`InboundDispatch` and `OutboundRoutes` from the start. ## 7. Alternatives considered 1. **`async fn` in traits with generic pumps.** Removes boxes by - monomorphization, but the pumps are type-erased per route (`dyn - SerializedSource`), and `Send` bounds on the returned futures need extra - machinery. Poll form keeps `dyn` and is already the codebase's pattern - (037). + monomorphization, but the pumps are type-erased per route, and `Send` + bounds on the returned futures need extra machinery. The pull model has + no trait that returns a future at all. 2. **`ReusableBoxFuture` everywhere.** One allocation per route instead of per message, but it needs `unsafe` or `tokio-util` on `no_std` (037 §3.2 rejected the hand-rolled version) and keeps an indirection per poll. 3. **Lending async `Source` (`async fn next(&mut self) -> Option>`).** The borrow must outlive `.await`, which conflicts with how both MQTT clients expose messages. Push (`dispatch`) fits them directly. -4. **Keep `Connector::publish` and pool its futures.** Futures of different - connectors have different sizes; pooling adds unsafe layout handling for - no gain over `poll_ready`/`commit`. -5. **Embedded: keep the action channel, with fixed-size `heapless` topic and - payload.** Removes the action allocations but not the encode `Vec`, and - sizes every slot for the largest route. The record ring (§4.6) holds - variable-length records instead. -6. **Embedded: route publishers encode PUBLISH into the write ring - (first revision of this design).** Publishers would need the session's - client state: whether it is connected and subscribed, whether the single - in-flight slot is free, and the next packet id. Frames encoded during an - outage would reach a new socket ahead of CONNECT. Sharing the state - across tasks needs a lock, and resetting the ring on reconnect loses - queued messages. It has the same single copy as §4.6, so it gains nothing - for its problems. -7. **Keep compatibility adapters (`BoxedPublisher`, poll `Source`, - `TopicProvider` adapter) for a staged migration.** Allows one PR per - connector, but keeps two mechanisms per direction and their bench rows. - Every connector is in this repository, so one migration is cheaper. -8. **Embedded: zero AimDB copies by hand-encoding PUBLISH in the record - ring.** Lay each record out as a frame with gaps (fixed header and topic - length in front, a hole for the packet id), let the session loop patch - only the packet id, and have `write_out` send straight from the record - rings. Removes the last copy, but AimDB would encode PUBLISH outside - `mountain-mqtt`'s codec and must stay byte-identical to it (MQTT 5 - properties included), `write_out` would interleave N record rings and the - control ring through a descriptor queue, and records would stay reserved - until the socket write completes, so rings grow. The copy it saves is one - memcpy of topic and payload, small next to the network stack's own copy - and TLS encryption. Revisit only if the STM32H5 rig shows the copy - matters. -9. **Publisher takes borrowed bytes (`start_send(&str, &[u8])`, second - revision of this design).** Simpler trait, but the pump serializes into - its own scratch and every publisher copies from it, so the embedded - backend pays two copies. With one shared record ring for all routes it - also needs a lock among the route pumps. Lending the buffers (§4.4) and - one ring per route remove both. +4. **Per-route push pumps with a poll-form publisher (earlier revisions of + this design).** Core keeps one pump task per outbound route; each + connector implements a `futures::Sink`-shaped `RoutePublisher` + (`poll_ready`, `start_send` or `buffers`/`commit`, `poll_flush`), with a + `SendBytes` helper for copying transports. Reaches the same allocation + numbers, but every connector implements a publisher trait, native MQTT + needs a stored `ReusableBoxFuture` and reports errors one message late, + and the embedded backend needs a queue between the pumps and the session + loop because the pumps run independently of it. That queue is either a + shared ring behind a lock with one extra copy, or one lock-free record + ring per route with lent buffers, a custom `bbqueue` notifier and + round-robin inside the session loop. The pull model removes the pumps, + so none of that is needed: the record buffers are the queue. +5. **Embedded: route publishers encode PUBLISH into the write ring (first + revision).** Publishers would need the session's client state (connected + and subscribed, the single in-flight slot, the next packet id), frames + encoded during an outage would reach a new socket ahead of CONNECT, and + sharing the state across tasks needs a lock. In the pull model the + session loop encodes, so the question does not arise. +6. **Keep compatibility adapters for a staged migration.** Allows one PR per + connector, but keeps two mechanisms per direction. Every connector is in + this repository, so one migration is cheaper. +7. **Embedded: zero AimDB copies by hand-encoding PUBLISH.** Serialize the + topic and payload straight into a write-ring grant at offsets left for + the fixed header, topic length and packet id, then fill those in. Removes + the last copy, but AimDB would encode PUBLISH outside `mountain-mqtt`'s + codec and must stay byte-identical to it (MQTT 5 properties included). + The copy it saves is one memcpy of topic and payload, small next to the + network stack's copy and TLS encryption. Revisit only if the STM32H5 rig + shows the copy matters; the pull model makes it a local change to the + session loop. ## 8. Open questions -- **Ingest latency on the transport task.** Embedded: a slow user - deserializer delays keep-alive and ack handling in the session loop. - Native: it delays polling `rumqttc`'s event loop. Measure on the STM32H5 - rig (037's B3) and with the native loopback row; decide whether to - document a deserializer budget. -- **Ring sizing (embedded).** Slots per record ring (at least two), the - default maximum payload for owned-serializer routes, the write ring size - and its control reserve. Per-route rings cost more RAM than one shared - ring; measure with the demo applications' route counts. -- **`bbqueue` stability.** 0.6 (January 2026) was a rewrite and 0.7 - followed in March; the author plans a few more minor breaking changes - before 1.0. Pin to `0.7` and budget for one upgrade. The pieces used here - (framed and stream grants, `BoxedSlice`, `AtomicCoord`, the `Notifier` - trait) are the core of the crate, not experimental features. -- **Write ring across reconnects.** `bbqueue` has no public `reset`. Either - the session drains the write ring through its consumer when a connection - ends, or each connection allocates a fresh write ring (one allocation per - reconnect, not per message). Draining is preferred on embedded. -- **Inbound chunk channel.** The embedded `read_into` side - (`Channel` into `PacketReader`) could use a stream `bbqueue` - too, reading socket bytes straight into a grant. Out of scope here; - noted for a follow-up. +- **(De)serialization latency on the transport task.** Embedded: a slow + user serializer or deserializer delays keep-alive and ack handling in the + session loop. Native: it delays polling `rumqttc`'s event loop or the + publish loop. Measure on the STM32H5 rig (037's B3) and with the native + loopback row; decide whether to document a budget. +- **Many routes.** A wake-up polls every open route. If a connector has + hundreds of outbound links, give `OutboundRoutes` a ready set fed by + per-route wakers instead of polling all of them. +- **Outage semantics per buffer type.** §4.6 changes what is sent after an + outage. Confirm the documented behaviour per buffer type with the + loopback harness, including SPMC lag reporting. +- **Write ring sizing (embedded).** Default size and `CONTROL_RESERVE`. - **Other connectors' own copies.** KNX, WebSocket and the AimX session were not measured at connector level; §5's per-connector rows will show them. +- **Inbound chunk channel (embedded).** The `read_into` side (`Channel` into `PacketReader`) could read socket bytes straight into a `bbqueue` + grant. Out of scope; noted for a follow-up. ## 9. References @@ -734,9 +548,7 @@ The Zenoh connector (053) is not implemented yet; it is written against - [053 — Zenoh connector](./053-zenoh-connector.md) - [055 — Wildcard inbound links](./055-wildcard-inbound-links.md) (grammar, keys, `TopicMatch`, §7 relation to this design) +- [`bbqueue`](https://github.com/jamesmunns/bbqueue) 0.7 - `aimdb-bench/benches/b0_alloc_connector.rs` -- [`bbqueue`](https://github.com/jamesmunns/bbqueue) 0.7, and the - [lock-free ring buffer design](https://ferrous-systems.com/blog/lock-free-ring-buffer/) - it is based on - `aimdb-mqtt-connector/src/embedded/session_loop.rs` (`client_loop`, `drain_packets`, `perform`, `encode`, `write_out`)