diff --git a/docs/design/054-zero-alloc-connector-boundary.md b/docs/design/054-zero-alloc-connector-boundary.md index 16dcdb84..89fdda0d 100644 --- a/docs/design/054-zero-alloc-connector-boundary.md +++ b/docs/design/054-zero-alloc-connector-boundary.md @@ -1,21 +1,26 @@ # 054 — Zero-allocation connector boundary -**Status:** 📝 Proposed — pull model: each connector's own task drives both -directions; core runs no per-route tasks. Revised 2026-10-03 after a -throwaway prototype (§5.1), which met every allocation target and changed +**Status:** 📝 Proposed, ready for implementation — pull model: each +connector's own task drives both directions; core runs no per-route tasks. +Revised 2026-10-03 after a throwaway prototype (§5.1), which met every allocation target and changed the outbound pull into two steps (§4.2), the topic-writer builder (§4.3) and the embedded write ring's sizing and wake-ups (§4.7). A second round made the ready set part of the design (§4.2): polling every route was already -2.7 times slower than today's task per route at 8 routes. Earlier revisions -(per-route push pumps with `RoutePublisher`, per-route record rings) are in -this branch's history and summarised in §7. +2.7 times slower than today's task per route at 8 routes. A third round +(2026-10-03, checked against the tree at `137ad2d`) made the ready set +lock-free (§4.2), covered the session client connectors and KNX (§4.8, +§4.9), fixed the embedded write-ring default (§4.7) and completed the +migration list (§6). 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, 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` and its closure form `with_topic_fn` (§4.3). +SPI and to two user-facing APIs: `with_topic_provider` is replaced by +`with_topic_writer` and its closure form `with_topic_fn` (§4.3), and +`KnxConnector::new` loses its `&'static Channels` argument (§4.9). `link_from`, `link_to`, `with_deserializer`, `with_match_deserializer`, -`with_serializer`, `with_serializer_into` and `Reader::recv` do not change; +`with_topic_resolver`, `with_serializer`, `with_serializer_into` and `Reader::recv` do not change; `Reader` gains `poll_recv`. **Compatibility:** none kept. No adapters for the old traits; all in-tree @@ -71,11 +76,15 @@ native one): **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`. +server, AimX session engine). 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 KNX +command channel, one socket). The session client connectors +(`SessionClientConnector`, used by the UDS, TCP and serial clients, and the +WebSocket client) use `pump_client` instead: one task per outbound route +over `SerializedReader::recv` (the owned path, never `recv_into`) and one +per inbound subscription. They and the WebSocket server 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 @@ -98,7 +107,10 @@ Nothing gates the numbers yet. client library's API requires, documented per backend. - G3. `b0_alloc_connector` runs in CI and gates the result. - G4. The connector SPI is two objects, one per direction, and core runs no - tasks of its own for a connector. No compatibility layers. + generic per-route or per-connector pump tasks (`pump_source`, + `pump_sink`). Connectors that core itself implements (the session + client and server connectors) own their tasks like any other connector. + No compatibility layers. **Non-goals** @@ -111,6 +123,10 @@ Nothing gates the numbers yet. - 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. +- Allocations inside the AimX session engine. `ClientHandle::write` takes + an owned topic and an `Arc<[u8]>` payload for its command channel; the + session client connectors keep that API (§4.8) and are treated like + `rumqttc` under G2. - Remote-access JSON paths (tracked by `b0_alloc_remote_json`). ## 3. Approach @@ -249,10 +265,10 @@ places; the message is lent once the poll has returned. - If a message is already staged, its route is returned at once. A staged message that was never taken is neither lost nor overwritten. -- Otherwise it takes routes from its ready set, oldest first, and polls - only those routes' typed readers (`poll_recv`, 037). The first value that - is ready is serialized (topic writer, then serializer) into the scratch - and staged. +- Otherwise it takes routes from its ready set in round-robin order and + polls only those routes' typed readers (`poll_recv`, 037). The first + value that is ready is serialized (topic writer, then serializer) into + the scratch and staged. - A value is taken from its buffer only in a call that returns `Ready`. `poll_stage` is therefore safe as a `select` arm: a losing arm takes nothing. @@ -286,37 +302,61 @@ reader is polled with a waker of its own, and only routes that woke are polled: ```rust +use portable_atomic::AtomicU32; // already a core dependency +use futures_util::task::AtomicWaker; // already a core dependency, no_std + struct ReadySet { - queue: spin::Mutex>, // capacity: every route - queued: Vec, // a route is queued at most once - task: AtomicWaker, // the transport task + ready: Box<[AtomicU32]>, // one bit per route, ceil(routes / 32) words + task: AtomicWaker, // the transport task } struct RouteWake { id: RouteId, set: Arc } impl Wake for RouteWake { fn wake_by_ref(self: &Arc) { - self.set.push(self.id); // no-op if already queued + self.set.ready[self.id / 32].fetch_or(1 << (self.id % 32), Release); self.set.task.wake(); } fn wake(self: Arc) { self.wake_by_ref() } } + +// OutboundRoutes also keeps `cursor: RouteId` (plain field, `&mut self`) +// and `open: Box<[u32]>` (bitmap of routes that are not closed). ``` +- **Lock-free, so a waker may fire from any context.** A route's waker is + called by whoever produces into its record. On Embassy that can be a + task on an `InterruptExecutor` (049) that preempts the transport task. A + lock held by the transport task while the waker runs would deadlock a + single core, which rules out the spin-locked FIFO the prototype used + (§7, alternative 11). `fetch_or` and `AtomicWaker` take no lock. - The route wakers are built once, in `OutboundRoutes::new`; polling a - route uses `Context::from_waker(&route_waker)`. Every route starts queued, + route uses `Context::from_waker(&route_waker)`. Every bit starts set, since no reader has registered a waker yet. -- `poll_stage` registers the task's waker before it looks at the queue, so - a route that wakes after the queue reads empty still wakes the task. -- A route that staged a value goes to the back of the queue: it may hold - more, and its waker will not fire again for values already in its buffer. - A route that returns `Pending` leaves the queue until its waker fires; a - closed route leaves it for good, and `Ready(None)` follows the last one. -- Oldest-first order is the fairness rule (§4.6). -- Nothing allocates per message: the queue has room for every route and a - route is queued at most once, the lock is a spin lock (`no_std`), and - cloning a route waker is a reference-count increment. Measured: 0 - allocations, parked or not (§5). +- `poll_stage` registers the task's waker before it reads the bitmap, so a + route that wakes after the bitmap reads empty still wakes the task. +- **Clear, then poll.** `poll_stage` clears a route's bit + (`fetch_and(!bit, Acquire)`) before polling its reader. A wake that lands + during the poll sets the bit again, which costs at most one spurious + re-poll and never loses a wake-up. +- **Round-robin.** The scan starts at the route after `cursor`, wraps once, + and skips words that are zero. A route that staged a value has its bit + set again (it may hold more, and its waker will not fire again for values + already in its buffer), and `cursor` moves to it, so every other ready + route is served before it is served again. A route that returns `Pending` + stays clear until its waker fires; a closed route is removed from `open` + for good, and `Ready(None)` follows the last one. +- Round-robin over ready routes is the fairness rule (§4.6). +- Cost per wake-up: one `poll_recv` per route that woke, plus a scan of + `ceil(routes / 32)` words (8 at 256 routes). +- Nothing allocates per message: the bitmap is allocated once, and cloning + a route waker is a reference-count increment. + +The prototype measured a spin-locked FIFO (one `VecDeque` of route ids and +a queued flag per route), not this bitmap. Both poll only routes that woke, +so §5.1's flat 590–630 ns is expected to hold; the +`outbound_next_round_robin` and `outbound_next_parked` gates (§5) and the +informational 1/8/64/256-route bench confirm it during step 2 of §6. **Where it is built.** In the connector's `build()`, not when its transport task starts. Per-route configuration (§4.4) and frame sizes (§4.7) are then @@ -327,10 +367,21 @@ task first polls instead of missing it. **Naming.** `aimdb_core::Outbound` already names the session envelope, so the pulled message is `OutboundMessage`. -`SerializedReader`, `SerializedPayload`, `RecvSerializedIntoFuture` and -`pump_sink` are removed. The typed reader, topic writer and serializer stay -inside core as the per-route state of `OutboundRoutes`, behind a -crate-private poll-shaped trait, so nothing is boxed or stored per message. +**Thread-safety.** `OutboundRoutes: Send` (moved into a connector's +`Send` task; `&mut self` everywhere, so not `Sync`). `InboundDispatch: +Send + Sync + Clone`. Both hold only `Send` readers and `Send + Sync` +(de)serializers, as today's pumps do. + +Removed from the public API: `pump_sink`, `SerializedReader`, +`SerializedSource`, `SerializedValue`, `SerializedValueInto`, +`SerializedPayload`, `RecvSerializedFuture` and `RecvSerializedIntoFuture`. +`AimDb::collect_outbound_routes` and `OutboundRoute` become crate-private; +their callers outside core (the embedded MQTT `warn_unsupported_qos`, and +the tests listed in §6) read `OutboundRoutes::routes()` instead. The typed +reader, topic writer and serializer stay inside core as the per-route state +of `OutboundRoutes`, behind a crate-private poll-shaped trait built by the +link's existing source factory (`SourceFactoryFn`, also crate-private), so +nothing is boxed or stored per message. ### 4.3 Outbound: topics are written, not returned @@ -397,8 +448,16 @@ while let Some(msg) = outbound.next().await { - 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 build; an invalid link (for - MQTT, a `qos` above 2) fails the build instead of every publish. +- Configuration is parsed once per route at build; an invalid link fails + the build instead of every publish. For MQTT that is a `qos` that is not + 0, 1 or 2, or a `retain` that is not `true` or `false`. Today a value + that does not parse silently falls back to the default (QoS 1, no + retain); after this change it is a build error, on both backends. The + embedded backend keeps its once-per-route warning for `qos=2`, which it + sends at QoS 1, now emitted from the same parse. +- `transport.rs` keeps `ConnectorConfig` (now `RouteInfo::config`, built by + `ConnectorConfig::from_query` as today) and `PublishError` (KNX's + `GroupWrite::try_new` returns it); only the `Connector` trait goes. - A connector that must send while also reading (one socket, one task) puts `poll_stage` in its `select`, gated by its own readiness, takes the message with `take_staged` after the `select` returns, and disarms the arm @@ -408,9 +467,10 @@ while let Some(msg) = outbound.next().await { buffer has closed (the database is dropped). The connector flushes its transport as it sees fit. -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. +Removed from the SPI: `Connector`, `Source`, `pump_source`, `pump_sink`, +the `Serialized*` types (§4.2) and `TopicProvider`. `pump_client` stays, +rewritten (§4.8). The `Vec` connectors return today shrinks to +their own transport tasks. ### 4.5 Values with heap data (guidance) @@ -449,11 +509,20 @@ Document this next to the buffer types; no API change. 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.** Oldest-first over routes that woke (§4.2); a busy route - goes to the back after each message, so it cannot starve the others. - Measured: three routes with three values each are pulled 0, 1, 2, 0, 1, - 2, …; with 200 values queued on one route and one on another, the second - is served first or second. Priorities are not in scope. +- **Fairness.** Round-robin over routes that woke (§4.2); after each + message from a busy route every other ready route is served once before + it, so it cannot starve the others. Measured on the prototype's FIFO, and + required of the bitmap by the same tests: three routes with three values + each are pulled 0, 1, 2, 0, 1, 2, …; with 200 values queued on one route + and one on another, the second is served first or second. Priorities are + not in scope. +- **Embedded MQTT: outbound PUBLISH size is bounded.** Today `encode` + allocates each packet at its exact length, so any size goes out. With the + write ring (§4.7) a PUBLISH frame is at most `capacity / 2 − + CONTROL_RESERVE`: 1,984 bytes at the 4,096-byte default. A route whose + declared capacities exceed that fails the build; an owned-serializer + payload over it is skipped, logged at warn level and counted by the + connector. The ring size is a connector setting (§4.7). - **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 @@ -558,12 +627,24 @@ let action_arm = poll_fn(|cx| { `payload_capacity`, from `RouteInfo`) plus `CONTROL_RESERVE` must fit in `capacity / 2`. A route that does not is rejected when the connector builds; an owned payload over the same bound is skipped and counted. -- **Sizes.** The prototype used a 2,048-byte ring and a 64-byte reserve, - which caps a PUBLISH at 960 bytes (§8). +- **Sizes.** Default ring 4,096 bytes, the same as the existing + `BUFFER_SIZE` the inbound side uses, and a 64-byte reserve: a PUBLISH of + up to 1,984 bytes. The prototype used 2,048 bytes (960-byte cap); + allocation counts do not depend on the size. The embedded builder gets + `with_write_buffer(bytes)` to change it, checked at build against + `2 × (max_frame + CONTROL_RESERVE)` of every route. The ring is + allocated once per connector and reused across reconnects. - `bbqueue = { version = "0.7", default-features = false, features = - ["alloc"] }`; `AtomicCoord` where the target has atomic pointers. On - `thumbv6m` the critical-section coordination (`CsCoord`) is needed - instead; not yet built (§8). + ["alloc"] }`, with `AtomicCoord`. Every embedded target this repository + builds (`thumbv7em`, `thumbv8m.main`) has atomic CAS; `thumbv6m` is not + a supported target (`aimdb-core` needs `alloc::sync::Arc`, which needs + CAS). +- *Files.* Besides `session_loop.rs`, the embedded backend's `mod.rs` + (`MqttSink`, `MqttSource`, `AimdbMqttEvent`, `AimdbMqttAction`, the + channel types and `CHANNEL_SIZE`), `manager.rs` (the channel aliases), + `session.rs` and `tls.rs` (both pass `events` and `actions` into + `run_session`; they pass `&InboundDispatch` and `&mut OutboundRoutes` + instead). 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 @@ -587,6 +668,87 @@ through `TokioTcpDialer` (052). (2,406 B) to 11 (557 B); what remains is the owned topic and payload and `rumqttc`'s own. +### 4.8 Session client connectors (`pump_client`) + +`SessionClientConnector` (the UDS, TCP and serial clients) and the +WebSocket client both call `pump_client`. It stays public, with a new +signature: + +```rust +pub fn pump_client( + db: &AimDb, + scheme: &str, + inbound: InboundDispatch, + handle: &ClientHandle, +) -> DbResult>>; // fallible: OutboundRoutes::new +``` + +- *Outbound.* One task, replacing one per route: + + ```rust + while let Some(msg) = outbound.next().await { + let payload: Payload = match msg.payload { + OutboundPayload::Borrowed(b) => Payload::from(b), + OutboundPayload::Owned(v) => Payload::from(v), + }; + if handle.write(msg.topic, payload).is_err() { + break; // engine stopped: every handle is gone + } + } + ``` + + `ClientHandle::write` is synchronous: it enqueues onto the engine's + bounded command queue (`max_offline_queue`) with `force_send`, which + displaces the oldest queued command when the queue is full. That is + today's behaviour and does not change; like `rumqttc`'s request channel + (§4.6), the engine's queue, not the record buffer, decides what survives + an outage. The owned topic `String` and the `Arc<[u8]>` payload stay + (non-goal, §2). The path moves from `recv` (owned serializer every time) + to the scratch path, so `with_serializer_into` links no longer allocate a + `Vec` first. +- *Inbound.* One task per subscription stays: each `ClientHandle::subscribe` + stream is its own channel, and merging them is out of scope. Each task + calls `inbound.dispatch(id, &update.data)` where it calls `router.route` + today. These are the connector's own tasks (G4). +- *Callers.* `SessionClientConnector::build` and the WebSocket client's + `build` change one line each (`InboundDispatch::new` and `?` on + `pump_client`). The UDS, TCP and serial crates do not change. +- The session *server* side (`SessionServerConnector`, `serve`, + `AimxDispatch`) does not use links and does not change. + +### 4.9 KNX connector + +Today the connection task talks to core through two `'static` channels +the application supplies (`Channels`: telegrams in, `GroupWrite` +commands out), drained and filled by `pump_source` and one `pump_sink` per +route. With both pumps gone, nothing is left on the other end of either +channel, so both go: + +- `Channels`, `TelegramChannel`, `CommandChannel` and `DEFAULT_QUEUE` are + removed, as is the `N` parameter of `KnxConnector`. + `KnxConnector::new(binder, delay, gateway_url)` drops its `channels` + argument; the two examples drop their `static KNX_CHANNELS`. This is the + second user-facing break (Scope). +- *Inbound.* `TelegramSink::try_send(String, Payload) -> bool` becomes + `deliver(&self, topic: &str, payload: &[u8])`, implemented over + `InboundDispatch`. `NeutralIo::forward` formats the group address into a + stack buffer (`core::fmt::Write` over `[u8; 16]`; the longest address, + `31/7/255`, is 8 bytes) instead of `addr.to_string()`. The `Vec` the + tunnel engine hands to `forward` is the engine's own and stays (§8). +- *Outbound.* `CommandSource` is removed. `drive_connection` takes `&mut + OutboundRoutes`, and `cmd_arm` becomes a `poll_fn` over `poll_stage`, + armed only while connected, with `Ready(None)` latched as in §4.7. + After the `select`, `take_staged` → `GroupWrite::try_new(msg.topic, + payload)` → `engine.handle_command`. An invalid group address or an + oversize payload is logged, counted by the connector and skipped, as + `KnxSink::publish` rejects it today. The arm-order swap that keeps + inbound traffic from starving commands stays. +- *Behaviour.* Commands produced while the tunnel is connecting or backing + off wait in their record buffers (§4.6) instead of the 32-deep command + channel. +- *Tests.* The task-count tests expect one future (the connection task) + with or without outbound routes. + ## 5. Measurement and gate `b0_alloc_connector` is rewritten against the new interfaces. Targets, and @@ -620,19 +782,32 @@ baseline together. It compares the rounded per-message figure: a one-off allocation can land in a 2,000-message window (one was seen in the prototype) without being a per-message cost. -A connector-level row per MQTT backend records G2: the embedded rows target -0 allocations in both directions and record the one outbound copy; the -native outbound row records the owned topic, the payload copy and whatever -`rumqttc` adds. The prototype measured these over loopback TCP with an echo -broker on its own thread, counting only allocations made on the database -thread (a counting allocator with a thread-local filter). Allocation counts -do not depend on socket options; latency does (§5.1). +**Connector-level rows (G2).** A new integration test, +`aimdb-mqtt-connector/tests/alloc_round_trip.rs`, runs §5.1's round trip +(produce → PUBLISH QoS 1 → echo → dispatch → `recv`) over loopback TCP +against the echo broker in `tests/common`. It counts only allocations made +on the database thread (a counting global allocator with a thread-local +filter) over 300 round trips after 100 of warm-up. It runs under the +existing `_test-backend-parity` feature, which builds both backends, from a +new `make test` line beside `backend_parity`, so CI's `make test` gates it. -**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 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. +| Row | Asserted | +|---|---| +| embedded, allocations per round trip | 0 | +| native, allocations per round trip | ≤ 11 (the prototype's figure; `rumqttc`'s share can only be bounded) | + +The test also prints bytes per round trip and the embedded backend's live +heap, and records the one outbound copy (scratch into frame) in a comment +beside the embedded row. Allocation counts do not depend on socket +options; latency does (§5.1), so latency is printed, not asserted. + +**Gate.** A Makefile target `bench-gate` runs `cargo bench -p aimdb-bench +--bench b0_alloc_connector` (host only, a few seconds) and fails on any +assertion. A `bench-gate` job in `.github/workflows/ci.yml` runs it and is +added to `comprehensive-check`'s `needs`. 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. The informational timing bench is not part of the job. ### 5.1 Prototype @@ -701,36 +876,42 @@ the new paths. | Who | Change | |---|---| -| Application code | `.with_topic_provider(p)` → `.with_topic_writer(capacity, w)` for a writer type, or `.with_topic_fn(capacity, closure)`. Choose `capacity` as the longest topic the writer produces. Nothing else. | -| Connector authors | Build `InboundDispatch` and `OutboundRoutes` in `build()` and move them into the transport task; call `dispatch` on receive; pull with `next`, or `poll_stage` + `take_staged` in a `select`, when the transport can send; stop on `None`. `Source`, `pump_source`, `pump_sink`, `Connector::publish` are gone. | +| Application code | `.with_topic_provider(p)` → `.with_topic_writer(capacity, w)` for a writer type, or `.with_topic_fn(capacity, closure)`. Choose `capacity` as the longest topic the writer produces. KNX: `KnxConnector::new(binder, delay, url)` without the `Channels` argument, and drop the `static` that held it. An MQTT `qos`/`retain` value that does not parse now fails the build (§4.4). | +| Connector authors | Build `InboundDispatch` and `OutboundRoutes` in `build()` and move them into the transport task; call `dispatch` on receive; pull with `next`, or `poll_stage` + `take_staged` in a `select`, when the transport can send; stop on `None`. `Source`, `pump_source`, `pump_sink`, `Connector`, the `Serialized*` types and `collect_outbound_routes` are gone; session clients call the new `pump_client` (§4.8). | 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 moves from `router.route` | `OutboundRoutes`, `TopicWriter`, `with_topic_writer`/`with_topic_fn`, `Reader::poll_recv`; 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.rs` and `decouple_record_keys_topics.rs` → `TopicWriter`; KNX task-count tests drop the per-route publisher; bench rewritten (§5) | +| `aimdb-core` | `InboundDispatch` (`router.rs`, `builder.rs`); remove `Source`, `pump_source` (`session/pump.rs`, `session/mod.rs`), public `inbound_router`; `Router` crate-private | `OutboundRoutes`, `RouteInfo`, `OutboundMessage`, `OutboundPayload`, `RouteStats`, ready set (§4.2); `TopicWriter`, `TopicBuf`, `TopicOverflow`, `with_topic_writer`/`with_topic_fn` (`typed_api.rs`, `connector.rs`); `Reader::poll_recv` (`buffer/reader.rs`); remove `pump_sink`, `Connector` (`transport.rs`), the `Serialized*` types, `TopicProvider`; `collect_outbound_routes`/`OutboundRoute` crate-private | +| `aimdb-core` session client | `pump_client` and `inbound_pump` take `InboundDispatch` (`session/client.rs`) | `pump_client` outbound becomes one task over `OutboundRoutes` (§4.8); `SessionClientConnector::build` (`session/connector.rs`) | +| `aimdb-mqtt-connector` embedded | `dispatch` in `drain_packets`; remove `MqttSource`, `AimdbMqttEvent`, `EventChannel`, `Received::Event` | action arm pulls from `OutboundRoutes`; `bbqueue` write ring; `with_write_buffer`; remove `MqttSink`, `AimdbMqttAction::Publish`, `ActionChannel`, `CHANNEL_SIZE`; `warn_unsupported_qos` reads `RouteInfo`. Files: `session_loop.rs`, `mod.rs`, `manager.rs`, `session.rs`, `tls.rs` (§4.7) | +| `aimdb-mqtt-connector` native | `dispatch` in the event-loop task; remove `MqttEventLoopSource` | publish loop task over `AsyncClient`; remove `MqttSink`; `qos`/`retain` parsed at build (`native.rs`) | +| `aimdb-knx-connector` | `TelegramSink::deliver` over `InboundDispatch`; remove `KnxSource`, `TelegramChannel` | `cmd_arm` pulls via `poll_stage`; remove `KnxSink`, `CommandSource`, `CommandChannel`, `Channels`, `DEFAULT_QUEUE`, `N` (`connector.rs`, `client.rs`, `lib.rs`; §4.9) | +| `aimdb-websocket-connector` | server: `WsDispatch` holds `InboundDispatch` (`dispatch.rs`, `builder.rs`); client: through `pump_client` | server: one broadcast loop pulls (`record_index` from `RouteInfo.config`), replacing `WsBusSink` (`server/connector.rs`); client: through `pump_client` (`client/builder.rs`) | +| `aimdb-uds-connector`, `aimdb-tcp-connector`, `aimdb-serial-connector` | — (via `SessionClientConnector`) | — | +| `aimdb-embassy-adapter` | remove `EmbassySourceRaw`, `EmbassySource` (`connectors.rs`) | remove `EmbassySinkRaw`, `EmbassySink` (`connectors.rs`); update the module docs in `connectors.rs` and `send_wrapper.rs` that name the pumps | +| examples | — | `tokio-knx-connector-demo`, `embassy-knx-connector-demo`: drop `Channels` | +| tests | `aimdb-data-contracts/src/link_codec.rs` (inbound half), `aimdb-mqtt-connector/tests/link_ext_tests.rs`, core's `typed_api.rs` router tests (`pump_source_routes_through_the_given_router` is removed) → `InboundDispatch` | `link_codec.rs` outbound half and `link_ext_tests.rs` → `OutboundRoutes::routes()` and `next`; `aimdb-mqtt-connector/tests/session_loop.rs` (`warn_unsupported_qos` test) → `RouteInfo`; `topic_provider_tests.rs` (mqtt, knx), websocket `e2e.rs` and `decouple_record_keys_topics.rs` → `TopicWriter`; core's `fused_reader_*` tests in `typed_api.rs` → the same cases against `OutboundRoutes` (buffer errors, serializer skips, scratch and owned fallback, invalid length, written topic); KNX task-count tests expect one future; `aimdb-client/tests/pump_client.rs` → new `pump_client` signature; new `alloc_round_trip.rs` (§5) | +| bench, docs | — | `b0_alloc_connector` rewritten and baseline replaced (§5); `aimdb-bench/README.md`; `CHANGELOG.md` entries in core, mqtt, knx, websocket, embassy-adapter | **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: `InboundDispatch`, `OutboundRoutes`, - `TopicWriter`; bench rewritten. +1. On main, independently: the `bench-gate` job for today's bench (§5), + and the native MQTT topic move (`publish.topic` instead of `.clone()`). +2. On the feature branch: `Reader::poll_recv`, `InboundDispatch`, + `OutboundRoutes` with the ready set, `TopicWriter`; bench rewritten. + Confirm the round-robin, parked and 1/8/64/256-route numbers (§4.2) + before any connector moves. 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`); + `queue_lossy`, `write_out`); 2. the action arm pulls from `OutboundRoutes`; `ActionChannel` and `MqttSink` removed. -4. MQTT native, KNX, WebSocket, AimX session client, Embassy adapter, - tests. +4. MQTT native, KNX, WebSocket, `pump_client`, Embassy adapter, examples, + tests; `alloc_round_trip.rs`. 5. Merge once the gate and every connector's tests pass. The prototype supports landing it as one change. While `with_topic_provider` @@ -798,28 +979,40 @@ The Zenoh connector (053) is not implemented yet; it is written against but about 120 ns per idle route per message: 1,273 ns at 8 routes and 32,230 ns at 256, where today's task per route stays near 470 ns. The ready set stays flat at 590–630 ns (§5.1). +11. **Ready set as a spin-locked FIFO (the prototype).** A `spin::Mutex` + around a `VecDeque` of route ids, with a queued flag per route. Strict + oldest-first order, and what §5.1 measured. But route wakers run in the + producer's context, and on Embassy that can be an `InterruptExecutor` + (049) preempting the transport task while it holds the lock, which + deadlocks a single core. The bitmap (§4.2) takes no lock and gives + round-robin order instead, which meets the same fairness tests. +12. **Keep KNX's `Channels` with a core-side pump into them.** Keeps + `KnxConnector::new` unchanged, but brings back the per-route task and + the queue between the record buffers and the transport that this + design removes everywhere else (§4.9). ## 8. Open questions +None of these blocks implementation. Each is a measurement or a follow-up; +§6 can proceed with the decisions above. + - **(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. -- **Ready set on an MCU.** The ready set (§4.2) keeps the host cost flat, - about 140 ns above a task per route. Its spin lock and waker hops have not - been measured on the STM32H5 rig, nor on Embassy's buffers. + publish loop. Measure on the STM32H5 rig (037's B3) and with + `alloc_round_trip.rs`'s printed latency; decide whether to document a + budget. Measurement, after merge. +- **Ready set on an MCU.** The bitmap (§4.2) is confirmed on the host by + §6 step 2. Its atomics and waker hops have not been measured on the + STM32H5 rig, nor on Embassy's buffers. Measurement, after merge. - **Outage semantics per buffer type.** Confirmed for the Tokio buffers - (§4.6). Still to confirm for the Embassy buffers, with the loopback - harness. -- **Write ring sizing (embedded).** 2,048 bytes and a 64-byte reserve - worked in the prototype and cap a PUBLISH at 960 bytes (§4.7). Decide the - default, and whether a connector can set its own size. -- **`thumbv6m`.** Needs bbqueue's critical-section coordination (`CsCoord`) - for the write ring; not built yet. `thumbv7em` builds. -- **Other connectors' own copies.** KNX, the WebSocket client and the AimX - session were not prototyped or measured at connector level; §5's - per-connector rows will show them. + (§4.6). Confirm for the Embassy buffers with the loopback harness in §6 + step 3; a difference is documented on the link builder, not designed + around. +- **Other connectors' own copies.** KNX's tunnel engine hands `forward` an + owned `Vec` per telegram, and the AimX session engine takes owned + topics and payloads (§4.8). Both are outside G1; a follow-up can make + either borrow. - **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. @@ -828,11 +1021,16 @@ The Zenoh connector (053) is not implemented yet; it is written against directions, every exchange then waits on the Nagle and delayed-ACK interaction, about 40 ms per round trip (§5.1). Track separately. +Resolved in the 2026-10-03 revision: write-ring sizing (4,096-byte default, +`with_write_buffer`; §4.7) and `thumbv6m` (not a supported target; §4.7). + ## 9. References - [037 — Zero-allocation consume path](./037-zero-alloc-consume-path.md) - [045 — Per-link codec selection](./045-per-link-codec-selection.md) (`with_serializer_into` scratch path) +- [049 — Real-time cadence](./049-real-time-cadence.md) + (`InterruptExecutor` producers, §4.2) - [052 — Runtime-neutral connectors](./052-runtime-neutral-connectors.md) (`pump_source`, `pump_sink`, `StreamDialer`) - [053 — Zenoh connector](./053-zenoh-connector.md)