Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
159 changes: 121 additions & 38 deletions docs/design/054-zero-alloc-connector-boundary.md
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,9 @@
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
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.

Expand Down Expand Up @@ -182,7 +184,7 @@ impl InboundDispatch {

```rust
/// Every outbound link of one scheme, read by the connector's transport task.
pub struct OutboundRoutes { /* routes, round-robin cursor, scratch, staged message */ }
pub struct OutboundRoutes { /* routes, ready set, 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]`).
Expand Down Expand Up @@ -222,8 +224,8 @@ impl OutboundRoutes {
pub fn routes(&self) -> &[RouteInfo];
pub fn stats(&self, id: RouteId) -> Option<RouteStats>;

/// Take the next ready value from any route, round-robin, and serialize it
/// into the scratch. Lends nothing. `Ready(Some(id))`: a message from
/// Take the next ready value from a route that woke, in FIFO order, 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<Option<RouteId>>;
/// Lend the staged message and clear it. `None` if nothing is staged.
Expand All @@ -247,11 +249,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, 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.
- 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.
- 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.
Expand All @@ -271,9 +272,51 @@ places; the message is lent once the poll has returned.
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.
- Cost of a wake-up: one `poll_recv` per route that woke, not per open
route. Measured (§5.1) with one busy route among N and the producer in
another task: 590–630 ns per message from 1 to 256 routes, against about
470 ns for today's task per route and 450 ns (1 route) to 32,230 ns (256
routes) for polling every route.

**Ready set.** Polling every open route on each wake-up costs about 120 ns
per idle route per message once the transport parks between messages (it
scans once to find the value and once to park again); at 8 routes that is
already 2.7 times slower than today's task per route (§5.1). So each route's
reader is polled with a waker of its own, and only routes that woke are
polled:

```rust
struct ReadySet {
queue: spin::Mutex<VecDeque<RouteId>>, // capacity: every route
queued: Vec<AtomicBool>, // a route is queued at most once
task: AtomicWaker, // the transport task
}

struct RouteWake { id: RouteId, set: Arc<ReadySet> }

impl Wake for RouteWake {
fn wake_by_ref(self: &Arc<Self>) {
self.set.push(self.id); // no-op if already queued
self.set.task.wake();
}
fn wake(self: Arc<Self>) { self.wake_by_ref() }
}
```

- The route wakers are built once, in `OutboundRoutes::new`; polling a
route uses `Context::from_waker(&route_waker)`. Every route starts queued,
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).

**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
Expand Down Expand Up @@ -394,16 +437,23 @@ Document this next to the buffer types; no API change.
`[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.

This holds where the connector owns its send path (embedded MQTT, the
WebSocket server). The native MQTT backend hands each message to
`rumqttc`, whose request channel (inbound topics + 10 deep) still queues:
in the outage test (§5.1) it sent all ten values both before and after
this change. Past that depth `publish` waits, and the rest stay in the
record buffers.
- **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. 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.** 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.
- **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
Expand Down Expand Up @@ -556,7 +606,14 @@ what the prototype measured (§5.1):
| `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.
path in the measured window. The prototype measured the round-robin and
parked rows both with polling every route and with the ready set (§4.2):
0 either way.

Beside the gated rows, an informational bench times one busy route among
1, 8, 64 and 256 (§5.1), so a change that brings back a per-route scan
shows up even though it allocates nothing. It is not a gate: timings need a
quiet runner.

The bench asserts exact values, so every change updates `EXPECTED` and the
baseline together. It compares the rounded per-message figure: a one-off
Expand Down Expand Up @@ -588,27 +645,50 @@ 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:
`recv`. 300 round trips after 100 of warm-up, same harness on both trees;
latency as the median and range of 5 runs. Live heap is what the database
thread holds before the measured window; the peak is what the window adds.

| Backend | `main`: allocs / bytes / µs | Prototype: allocs / bytes / µs |
| Backend | `main` | Prototype |
|---|---|---|
| Embedded: allocations | 9 (302 B) | 0 |
| Embedded: latency | 56.1 µs (49.3–58.3) | 44.4 µs (44.0–46.1) |
| Embedded: live heap / window peak | 41,346 B / +73 B | 45,646 B / +0 B |
| Native: allocations | 16 (2,406 B) | 11 (557 B) |

The write ring (2 KB) is most of the 4.3 KB the embedded backend now holds
up front; in exchange nothing is allocated while traffic flows. Nine
allocations at 14–23 ns each (host glibc) account for about 0.2 µs of the
12 µs latency gain; the rest is the removed task and channel hand-offs.
Timings come from one 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).

Outage: ten values into a single-latest record while the broker accepted
the connection but did not answer, then the broker answered. Received:

| Backend | `main` | Prototype |
|---|---|---|
| Embedded | 9 / 302 / 58.6 | 0 / 0 / 48.3–49.3 |
| Native | 16 / 2,406 / 44,013 | 11 / 557 / 44,120 |
| Embedded | `0, 1, … 9` | `9` |
| Native | `0, 1, … 9` | `0, 1, … 9` |

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: one busy route among N idle ones, the producer and the
transport in separate Tokio tasks, so every message pays one wake-up. ns
per message, median of 5 runs:

Wake-up cost with one hot route among N (host, release build, 20,000 pulls,
ns per produce and pull):
| Routes | Task per route (`main`) | One task, poll every route | One task, ready set |
|---|---|---|---|
| 1 | 478 | 450 | 622 |
| 8 | 467 | 1,273 | 615 |
| 64 | 454 | 7,916 | 590 |
| 256 | 489 | 32,230 | 631 |

| Buffer | 1 | 8 | 64 | 256 |
|---|---|---|---|---|
| SPMC ring | 128 | 681 | 4,954 | 21,016 |
| Single-latest | 316 | 673 | 3,353 | 12,592 |
In a single task with a value always ready (no parking), polling every
route cost about 80 ns per idle route, and replacing today's two-allocation
path with `next` saved 34 ns per message (164 → 130 ns).

Behaviour covered by tests on real Tokio buffers: round-robin order, a hot
Behaviour covered by tests on real Tokio buffers, run against both polling
every route and the ready set: 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()`,
Expand Down Expand Up @@ -713,6 +793,11 @@ The Zenoh connector (053) is not implemented yet; it is written against
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).
10. **Poll every route on each wake-up (previous revision).** No per-route
wakers and the cheapest option with one route (450 ns against 622 ns),
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).

## 8. Open questions

Expand All @@ -721,11 +806,9 @@ 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: 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.
- **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.
- **Outage semantics per buffer type.** Confirmed for the Tokio buffers
(§4.6). Still to confirm for the Embassy buffers, with the loopback
harness.
Expand Down
Loading