From 09b9e46da57b5a0ef719d96d4ca7896be0ee4aad Mon Sep 17 00:00:00 2001 From: Claude Date: Sat, 3 Oct 2026 06:19:47 +0000 Subject: [PATCH] docs(054): revise after the pull-model prototype MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A throwaway prototype built §4 against main (core API, both MQTT backends, the WebSocket server's outbound path, the embedded write ring) and measured it. Every allocation target held; the doc changes where the code disagreed. - §4.2: the lending poll_next cannot be wrapped in poll_fn (captured variable cannot escape FnMut), so neither next() nor a select arm compiled. Pull in two steps: poll_stage, then take_staged. Re-poll a route after a skip or lag so its waker stays registered; Ready(None) is final and must disarm a select arm; Tokio's coop budget makes Pending not mean empty. RouteId is a usize, RouteInfo carries the route's topic and payload capacities, RouteStats counts skips, the message type is OutboundMessage (Outbound is taken), and OutboundRoutes is built in build(). - §4.3: with_topic_fn beside with_topic_writer, since an unannotated closure fails inference (E0282) under a generic TopicWriter bound; From for TopicOverflow; overflow detected by TopicBuf even if the writer ignores the error. - §4.7: an empty bipbuffer whose pointers have moved only guarantees about half its capacity, so max_frame + reserve must fit capacity / 2. The has_room gate is a probe grant re-checked on every poll; PUBACKs wait for room (the reserve covers one packet) and are encoded before state.receive (Puback borrows the state); two AtomicWakers replace "polling notifier only"; Ready(None) is latched, or an inbound-only session never yields. - §5/§5.1: prototype measurements. Embedded MQTT round trip 9 -> 0 allocations (302 B -> 0), 58.6 -> ~49 us on loopback; native 16 -> 11; wake-up cost about 80 ns per idle SPMC route on the host; new outbound_next_parked row. - §4.6, §6, §7, §8: outage semantics confirmed on Tokio buffers; the no-compatibility decision supported (an ignored with_topic_provider made ten WebSocket tests time out instead of failing to compile); two more rejected alternatives; open questions updated, TCP_NODELAY noted as out of scope. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01HXDnQEx2THUSTASYkThxCw --- .../054-zero-alloc-connector-boundary.md | 448 +++++++++++++----- 1 file changed, 328 insertions(+), 120 deletions(-) diff --git a/docs/design/054-zero-alloc-connector-boundary.md b/docs/design/054-zero-alloc-connector-boundary.md index 9cc3932a..1fb3137b 100644 --- a/docs/design/054-zero-alloc-connector-boundary.md +++ b/docs/design/054-zero-alloc-connector-boundary.md @@ -1,16 +1,20 @@ # 054 — Zero-allocation connector boundary -**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. +**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 +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). 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` (§4.3). `link_from`, `link_to`, `with_deserializer`, -`with_match_deserializer`, `with_serializer`, `with_serializer_into` and -`Reader::recv` do not change. +by `with_topic_writer` and its closure form `with_topic_fn` (§4.3). +`link_from`, `link_to`, `with_deserializer`, `with_match_deserializer`, +`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 connectors migrate in one change (§6). @@ -54,8 +58,9 @@ connector interface: | A5 | out | `TopicProvider::topic` returns `Option` | `connector.rs` | | A6 | out | owned serializer returns `Vec` (only without `with_serializer_into`) | `typed_api.rs` | -**Connector-side copies and allocations** (read from code, not yet measured -at connector level): +**Connector-side copies and allocations** (read from code; §5.1 measures the +totals per round trip: 9 allocations on the embedded backend, 16 on the +native one): | Backend | Inbound | Outbound | |---|---|---| @@ -177,15 +182,21 @@ impl InboundDispatch { ```rust /// Every outbound link of one scheme, read by the connector's transport task. -pub struct OutboundRoutes { /* routes, round-robin cursor, scratch buffers */ } +pub struct OutboundRoutes { /* routes, round-robin cursor, scratch, staged message */ } + +/// Dense index, `0..routes().len()`. A plain `usize`, so a connector indexes +/// its own per-route tables with it (`opts[msg.route.id]`). +pub type RouteId = usize; pub struct RouteInfo { - pub id: RouteId, // dense index, 0..len + pub id: RouteId, pub default_topic: Arc, // from the link URL pub config: ConnectorConfig, // qos, retain, record_index, … + pub topic_capacity: usize, // §4.3; 0 = default topic only + pub payload_capacity: usize, // `with_serializer_into`; 0 = owned only } -pub struct Outbound<'a> { +pub struct OutboundMessage<'a> { pub route: &'a RouteInfo, /// The written topic (`TopicWriter`) or the route's default. pub topic: &'a str, @@ -200,42 +211,83 @@ pub enum OutboundPayload<'a> { Owned(Vec), } +/// Values taken from a route's buffer, by outcome. +pub struct RouteStats { pub sent: u64, pub lagged: u64, pub topic_overflow: u64, pub serialize_failed: u64 } + 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. + /// Subscribes every outbound link of `scheme`. Allocates the scratch once: + /// the largest topic capacity plus the largest payload capacity. pub fn new(db: &AimDb, scheme: &str) -> DbResult; - /// Routes, for parsing per-route configuration once at start. + /// Routes, for parsing per-route configuration once at build. 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>; + pub fn stats(&self, id: RouteId) -> Option; + + /// Take the next ready value from any route, round-robin, and serialize it + /// into the scratch. Lends nothing. `Ready(Some(id))`: a message from + /// route `id` is staged. `Ready(None)`: every route is closed (final). + pub fn poll_stage(&mut self, cx: &mut Context<'_>) -> Poll>; + /// Lend the staged message and clear it. `None` if nothing is staged. + pub fn take_staged(&mut self) -> Option>; + + /// `poll_stage`, then `take_staged`, for hand-written `poll` code. + pub fn poll_next(&mut self, cx: &mut Context<'_>) -> Poll>>; + /// `poll_fn(|cx| self.poll_stage(cx)).await?; self.take_staged()`. + pub async fn next(&mut self) -> Option>; } ``` -**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. +**Why two steps.** A lending `poll_next` cannot be wrapped in `poll_fn`. An +`FnMut` closure cannot return a borrow of state it captured by `&mut`, since +it may be called again, so neither `next = poll_fn(|cx| self.poll_next(cx))` +nor a `poll_fn` select arm over it compiles ("captured variable cannot escape +`FnMut` closure body"). `poll_stage` returns no borrow and works in both +places; the message is lent once the poll has returned. + +**Semantics of `poll_stage`.** + +- 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, 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 staged. 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 + `poll_stage` 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. +- After a value that is skipped (topic overflow, serializer error) or a + `BufferLagged`, the same route is polled again until it returns `Pending` + or a usable value. Moving on instead would leave that reader without a + registered waker, and its next value would not wake the task. +- Skips and lag are logged and counted per route (`RouteStats`); a closed + buffer closes the route. None of these end the other routes. +- `Ready(None)` is final, and a connector with no outbound links gets it on + the first poll. A `select` arm must be disarmed after it: an arm that + resolves at once on every iteration keeps the loop from ever yielding + (§4.7). +- `Pending` does not mean every route is empty. On Tokio a broadcast reader + returns `Pending`, with a self-wake, once the task's cooperative budget + (128 polls) is spent. An awaiting loop is unaffected; a synchronous "drain + until `Pending`" must not read it as empty. +- One message is staged at a time (`&mut self`). The connector finishes with + it (copies, encodes, or awaits its send) before staging the next. +- Cost of a wake-up: one `poll_recv` per open route. Measured on the host + (§5.1): about 80 ns per idle SPMC-ring route, so about 5 µs per message + at 64 routes and 21 µs at 256. §8 notes a ready set for many routes. + +**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 +validated before anything runs, and the value moves into the transport task. +Cursors start at build, so an SPMC ring keeps what is produced before the +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 + serializer pair stays inside -core as the per-route state of `OutboundRoutes`. +`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. ### 4.3 Outbound: topics are written, not returned @@ -247,30 +299,50 @@ pub trait TopicWriter: Send + Sync { fn write_topic(&self, value: &T, out: &mut TopicBuf<'_>) -> Result; } -// 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) +// Closures implement it through a blanket impl. +impl TopicWriter for F +where + F: Fn(&T, &mut TopicBuf<'_>) -> Result + Send + Sync, +{ /* … */ } + +pub struct TopicOverflow; +// So that `write!(out, …)?` works inside a writer. +impl From for TopicOverflow { /* … */ } + +// TopicBuf implements core::fmt::Write over the scratch topic region, plus +// push_str, len, capacity and as_str. A write past the route's capacity is +// refused whole, so the bytes are valid UTF-8 by construction. ``` -- `.with_topic_writer(capacity, writer)` on the outbound builder replaces - `.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. -- Removes A5. +- Two builder methods replace `.with_topic_provider(provider)`, the one + user-facing break in this design: + - `.with_topic_writer(capacity, writer)` takes any `TopicWriter`, for + reusable writer types. + - `.with_topic_fn(capacity, |v, out| { write!(out, "sensors/{}/{}", v.site, v.id)?; Ok(true) })` + takes a closure. It exists for type inference: under the generic + `W: TopicWriter` bound, that unannotated closure fails with E0282 + ("type annotations needed"). A parameter bound directly on the `Fn` + signature gives the closure its higher-ranked signature. +- Overflow is detected by `TopicBuf`, not by the writer's return value. A + writer that ignores the error (`let _ = write!(…); Ok(true)`) still has + its value skipped. The value is skipped and counted per link + (`RouteStats::topic_overflow`); the topic is never truncated. +- Every migrated provider needs a capacity: the longest topic the writer + can produce. `TopicProvider` never had to say. +- Removes A5. Measured: 0 allocations with a written topic (§5.1). ### 4.4 Connector shape -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: +A connector builds both objects in its `build()` (§4.2) and drives them from +its transport task. For a client with an async send, the outbound half is a +plain loop: ```rust +// In build(): a bad link fails the build. let mut outbound = OutboundRoutes::new(&db, "mqtt")?; let opts: Vec = outbound.routes().iter().map(PublishOpts::parse).collect::>()?; +// In the transport task: 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 { @@ -282,12 +354,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 start; an invalid link is - reported before anything runs. +- Configuration is parsed once per route at build; an invalid link (for + MQTT, a `qos` above 2) fails the build instead of every publish. - 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. + 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 + after `Ready(None)` (§4.7). +- Shutdown: `next` returns `None` once every route is closed: at once when + the scheme has no outbound links, otherwise when every outbound record's + 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` @@ -313,15 +389,21 @@ Document this next to the buffer types; no API change. 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. + last, an SPMC ring its backlog or a lag report. Confirmed on the Tokio + buffers (§5.1): after five values with nobody pulling, single-latest sends + `[4]`, a mailbox `[4]` and an SPMC ring `[0, 1, 2, 3, 4]`; ten values into + a four-slot ring give one lag report, then `[6, 7, 8, 9]`. 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. + 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. - **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 @@ -330,7 +412,8 @@ Document this next to the buffer types; no API change. `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. + they already do for a slow consumer. The loopback test for outbound + progress under an inbound flood passes unchanged on the prototype. - **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 @@ -346,53 +429,98 @@ 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 action arm pulls from `OutboundRoutes` -instead of the `ActionChannel`, under the same gate as today: +*Outbound.* The session loop's action arm stages from `OutboundRoutes`, +under the same gate as today plus room in the write ring. Room is checked on +every poll of the arm, not once when it is built, because PUBACKs and pings +can take ring space while the arm is parked: ```rust let action_ready = connected - && !state.waiting_for_responses() // one in-flight slot - && next_topic >= subscribe_topics.len() - && write_ring.has_room(max_frame + CONTROL_RESERVE); + && !outbound_done // latched on Ready(None) + && !state.waiting_for_responses() // one in-flight slot + && next_topic >= subscribe_topics.len(); let action_arm = poll_fn(|cx| { - if action_ready { outbound.poll_next(cx) } else { Poll::Pending } + if !action_ready { + return Poll::Pending; + } + if !ring.has_room(max_frame + CONTROL_RESERVE) { + ring.room.register(cx.waker()); + if !ring.has_room(max_frame + CONTROL_RESERVE) { + return Poll::Pending; + } + } + match outbound.poll_stage(cx) { + Poll::Ready(None) => { outbound_done = true; Poll::Pending } + other => other, + } }); -// on Some(msg): +// after the select, on a staged route: +// msg = outbound.take_staged() // 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) +// ring.put(&packet, CONTROL_RESERVE).await // the one copy +// state.publish_update(&packet) ``` - `MqttSink`, `AimdbMqttAction::Publish` and the `ActionChannel` are removed. Subscriptions keep their current path (the loop already places them itself). +- **`Ready(None)` is latched.** A connector with only inbound links has no + outbound routes. Passing `Ready(None)` through would resolve the arm on + every iteration, and the session would never yield: in the prototype the + idle-session test then saw no pings at all. - **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. + `bbqueue` 0.7 stream queue (`BBQueue`), + 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. +- **Wake-ups.** bbqueue's polling notifier wakes nothing, and `select3` + polls `write_out` before the session loop, so bytes the loop commits + would sit until something else woke the task. Two `AtomicWaker`s + (embassy-sync) carry the signals: `data`, from a commit to `write_out`, + and `room`, from a release to the loop. +- **`has_room` is a probe.** bbqueue 0.7 has no free-space query, so the + gate takes `grant_exact(n)` and drops it uncommitted. A probe that wraps + commits an early wraparound (the unused tail is skipped until the reader + passes it); no data changes. A count of bytes in flight would not do: it + says nothing about whether a contiguous grant exists. - **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 + CONTROL_RESERVE)` and commits only `frame_len`, so the next control packet + finds room. The reserve covers one packet. A coalesced burst of QoS 1 + PUBLISHes needs a PUBACK each, so the PUBACK path waits for room instead + of assuming it. PINGREQ stays lossy, as today. +- **PUBACK before the state update.** `receive_produce_response(&self)` + returns a `Puback<'_>` that borrows the client state, so the PUBACK is + encoded into the ring (after waiting for room, if it must) before + `state.receive` takes `&mut`. +- **Reconnect.** At the start of each session the loop drains the write + ring through its consumer, so no bytes from an old session reach a new + socket ahead of its CONNECT. 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. +- **Oversize.** A bipbuffer only grants contiguous space. Once its read and + write pointers have both moved to `k`, an empty ring can grant at most + `max(capacity − k, k − 1)`, about half its capacity in the worst case: a + 2,048-byte ring drained at offset 1,024 cannot grant 1,100 bytes. A larger + frame could wait forever. So the largest frame a route can produce + (`max_frame`: PUBLISH overhead + `max(default topic, topic_capacity)` + + `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). - `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. + ["alloc"] }`; `AtomicCoord` where the target has atomic pointers. On + `thumbv6m` the critical-section coordination (`CsCoord`) is needed + instead; not yet built (§8). 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). +its TX buffer, which TLS needs anyway to encrypt. Measured (§5.1): a full +round trip went from 9 allocations (302 B) to 0, and on loopback TCP from +58.6 to about 49 µs. This also gives std users an allocation-free MQTT path +through `TokioTcpDialer` (052). #### Native backend (`rumqttc`) @@ -403,31 +531,45 @@ an allocation-free MQTT path through `TokioTcpDialer` (052). `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. + Per-route `qos` and `retain` are parsed at build. - `rumqttc`'s own allocations (building `Publish`, its request channel) are - outside G1. + outside G1. Measured (§5.1): a round trip went from 16 allocations + (2,406 B) to 11 (557 B); what remains is the owned topic and payload and + `rumqttc`'s own. ## 5. Measurement and gate -`b0_alloc_connector` is rewritten against the new interfaces. Targets: +`b0_alloc_connector` is rewritten against the new interfaces. Targets, and +what the prototype measured (§5.1): -| Row | Before | After | -|---|---|---| -| `inbound_dispatch` (was `inbound_route`) | 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_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 | +| Row | Before | After | Prototype | +|---|---|---|---| +| `inbound_dispatch` (was `inbound_route`) | 0 | 0 | 0 | +| `inbound_dispatch_pattern` | 0 | 0 | 0 | +| `inbound_dispatch_keyed_known` | 0 | 0 | 0 | +| `inbound_dispatch_keyed_new` | 1 | 1 (by design, 055 §5.6) | 1 | +| `inbound_pump_source_minimal` | 2 | removed (`Source` is gone) | — | +| `outbound_next_static_topic` (was `outbound_scratch_static_topic`) | 2 | 0 | 0 | +| `outbound_next_written_topic` (replaces the `TopicProvider` row) | 3 | 0 | 0 | +| `outbound_next_owned` (was `outbound_owned_static_topic`) | 3 | 1 (A6, by choice of serializer) | 1 | +| `outbound_next_round_robin` (new; 8 routes, all ready) | — | 0 | 0 | +| `outbound_next_parked` (new; every pull parks and is woken) | — | 0 | 0 | + +Only `outbound_next_parked` parks before each pull, which puts the waker +path in the measured window. 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 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. +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). **Gate.** A CI job runs `cargo bench -p aimdb-bench --bench b0_alloc_connector` (host only, a few seconds) and fails on any assertion. @@ -435,25 +577,65 @@ 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. +### 5.1 Prototype + +A throwaway prototype (not merged) built §4 against `main` at `dfe6adc`: +`InboundDispatch`, `OutboundRoutes` and `TopicWriter` in core beside the +old SPI, both MQTT backends and the WebSocket server's outbound path on the +new interfaces, and the embedded write ring. KNX, the WebSocket client, the +AimX session client and the Embassy adapter stayed on the old SPI, which +was not deleted, and no CI job was added. + +Connector round trip: produce, outbound PUBLISH at QoS 1, the broker's +PUBACK and echo PUBLISH, inbound dispatch, the client's PUBACK, reader +`recv`. 300 round trips after 100 of warm-up, same harness on both trees: + +| Backend | `main`: allocs / bytes / µs | Prototype: allocs / bytes / µs | +|---|---|---| +| Embedded | 9 / 302 / 58.6 | 0 / 0 / 48.3–49.3 | +| Native | 16 / 2,406 / 44,013 | 11 / 557 / 44,120 | + +Timings are single runs on a shared host. The embedded runs set +`TCP_NODELAY` on the dialer; nothing sets it for `rumqttc`, so the native +round trip sits on the 40 ms Nagle and delayed-ACK floor on both trees +(§8). + +Wake-up cost with one hot route among N (host, release build, 20,000 pulls, +ns per produce and pull): + +| Buffer | 1 | 8 | 64 | 256 | +|---|---|---|---|---| +| SPMC ring | 128 | 681 | 4,954 | 21,016 | +| Single-latest | 316 | 673 | 3,353 | 12,592 | + +Behaviour covered by tests on real Tokio buffers: round-robin order, a hot +route not starving a quiet one, written and default topics, topic overflow +(including a writer that ignores the error), the owned-serializer fallback, +lag, outage semantics per buffer type (§4.6), a dropped pending `next()`, +500 values pulled under a `select` that loses every third poll, and a +skipped value not losing the wake-up. The existing MQTT loopback, TLS, +reconnect and QoS 1 flood tests and the WebSocket end-to-end tests pass on +the new paths. + ## 6. Migration | Who | Change | |---|---| -| Application code | `.with_topic_provider(p)` → `.with_topic_writer(capacity, w)`. Nothing else. | -| 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. | +| 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. | 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`; remove `pump_sink`, `Connector::publish`, `SerializedReader`, `TopicProvider`; AimX session client pulls in its session task | +| `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 → `TopicWriter`; KNX task-count tests drop the per-route publisher; bench rewritten (§5) | +| 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) | **Order:** @@ -471,6 +653,11 @@ work lands as one change (a feature branch merged once): tests. 5. Merge once the gate and every connector's tests pass. +The prototype supports landing it as one change. While `with_topic_provider` +still existed but the pull path ignored it, ten WebSocket tests timed out +instead of failing to compile. Removing the old SPI in the same change makes +every unmigrated caller a compile error. + The Zenoh connector (053) is not implemented yet; it is written against `InboundDispatch` and `OutboundRoutes` from the start. @@ -517,6 +704,15 @@ The Zenoh connector (053) is not implemented yet; it is written against 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. **A lending `poll_next` as the only pull (previous revision).** One call + instead of two, but it cannot be wrapped in `poll_fn`: an `FnMut` closure + cannot return a borrow of state it captured by `&mut`. Neither `next()` + nor a `select` arm could use it, so the prototype split it into + `poll_stage` and `take_staged` (§4.2). +9. **Embedded: gate on a count of bytes in flight.** Cheaper than a probe + grant, but a bipbuffer's free bytes say nothing about whether a + contiguous grant of that size exists. The probe answers the question the + gate asks (§4.7). ## 8. Open questions @@ -525,18 +721,29 @@ The Zenoh connector (053) is not implemented yet; it is written against 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 +- **Many routes.** A wake-up polls every open route: about 80 ns per idle + SPMC-ring route on the host (§5.1), likely an order of magnitude more on + an MCU. Measure on the STM32H5 rig; if a realistic link count costs more + than the message itself, 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. +- **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. - **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. +- **`TCP_NODELAY` (not part of this design).** Neither `TokioNet::tcp()` nor + the native backend's `rumqttc` options set it. With QoS 1 traffic in both + directions, every exchange then waits on the Nagle and delayed-ACK + interaction, about 40 ms per round trip (§5.1). Track separately. ## 9. References @@ -549,6 +756,7 @@ 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) - [`bbqueue`](https://github.com/jamesmunns/bbqueue) 0.7 +- `embassy_sync::waitqueue::AtomicWaker` (write-ring wake-ups, §4.7) - `aimdb-bench/benches/b0_alloc_connector.rs` - `aimdb-mqtt-connector/src/embedded/session_loop.rs` (`client_loop`, `drain_packets`, `perform`, `encode`, `write_out`)