From 5cf41bc1e7c9ce8dfc42b380c2c125f55ee8ce6a Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Alexander=20Schn=C3=B6rch?= Date: Mon, 5 Oct 2026 17:27:59 +0000 Subject: [PATCH 1/4] ci: run CI and the docs check for the 054 feature branch Co-Authored-By: Claude Opus 5.5 --- .github/workflows/ci.yml | 4 ++-- .github/workflows/docs.yml | 4 ++-- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index a4349083..7f528614 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -2,9 +2,9 @@ name: CI on: push: - branches: [ main, develop ] + branches: [ main, develop, feat/054-connector-boundary ] pull_request: - branches: [ main, develop ] + branches: [ main, develop, feat/054-connector-boundary ] env: CARGO_TERM_COLOR: always diff --git a/.github/workflows/docs.yml b/.github/workflows/docs.yml index 63086e73..7a55d83d 100644 --- a/.github/workflows/docs.yml +++ b/.github/workflows/docs.yml @@ -2,9 +2,9 @@ name: Documentation on: push: - branches: [ main ] + branches: [ main, feat/054-connector-boundary ] pull_request: - branches: [ main ] + branches: [ main, feat/054-connector-boundary ] permissions: contents: read From 2d0ae150d00c47d585545e24378f257b6ed1e0dd Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Alexander=20Schn=C3=B6rch?= Date: Mon, 5 Oct 2026 17:32:39 +0000 Subject: [PATCH 2/4] =?UTF-8?q?test(embassy):=20outage=20semantics=20of=20?= =?UTF-8?q?OutboundRoutes=20on=20the=20Embassy=20buffers=20(design=20054?= =?UTF-8?q?=20=C2=A78)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Opus 5.5 --- .../tests/outbound_routes.rs | 104 ++++++++++++++++++ 1 file changed, 104 insertions(+) create mode 100644 aimdb-embassy-adapter/tests/outbound_routes.rs diff --git a/aimdb-embassy-adapter/tests/outbound_routes.rs b/aimdb-embassy-adapter/tests/outbound_routes.rs new file mode 100644 index 00000000..0b82ed87 --- /dev/null +++ b/aimdb-embassy-adapter/tests/outbound_routes.rs @@ -0,0 +1,104 @@ +//! What a connector pulling from `OutboundRoutes` sees after an outage, on the +//! Embassy buffers: the same cases as the Tokio adapter's +//! `outage_semantics_per_buffer_type`, driven on the host with a no-op waker. +#![cfg(all(feature = "embassy-sync", feature = "embassy-time"))] + +use std::future::Future; +use std::pin::Pin; +use std::sync::Arc; +use std::task::{Context, Poll, Waker}; + +use aimdb_core::buffer::DynBuffer; +use aimdb_core::connector::ConnectorBuilder; +use aimdb_core::executor::test_support::NoopRuntimeOps; +use aimdb_core::{AimDb, AimDbBuilder, DbResult, OutboundRoutes}; +use aimdb_embassy_adapter::EmbassyBuffer; +use futures::executor::block_on; + +// No-op defmt logger + host time driver, so the binary links. +aimdb_embassy_adapter::host_test_stubs!(); + +type Futures = Vec + Send + 'static>>>; + +/// Lets `link_to("test://…")` register; drives nothing. +struct TestConnector; + +impl ConnectorBuilder for TestConnector { + fn build<'a>( + &'a self, + _db: &'a AimDb, + ) -> Pin> + Send + 'a>> { + Box::pin(async { Ok(Vec::new()) }) + } + fn scheme(&self) -> &str { + "test" + } +} + +const KEYS: [&str; 3] = ["r0", "r1", "r2"]; + +/// One record per buffer, keyed `r0`, `r1`, …, each linked to `test://r{i}`. +fn db(buffers: Vec>>) -> AimDb { + let mut builder = AimDbBuilder::new() + .runtime(Arc::new(NoopRuntimeOps)) + .with_connector(TestConnector); + for (i, buffer) in buffers.into_iter().enumerate() { + builder.configure::(KEYS[i], move |reg| { + reg.buffer_raw(buffer) + .link_to(&format!("test://r{i}")) + .with_serializer(|_ctx, v: &u32| Ok(v.to_le_bytes().to_vec())) + .finish(); + }); + } + block_on(builder.build()).expect("build").0 +} + +/// Every message available now, as (route, value). +fn drain(o: &mut OutboundRoutes) -> Vec<(usize, u32)> { + let mut cx = Context::from_waker(Waker::noop()); + let mut out = Vec::new(); + while let Poll::Ready(Some(m)) = o.poll_next(&mut cx) { + out.push(( + m.route.id, + u32::from_le_bytes(m.payload.as_slice().try_into().unwrap()), + )); + } + out +} + +fn values(got: &[(usize, u32)], route: usize) -> Vec { + got.iter() + .filter(|(id, _)| *id == route) + .map(|(_, v)| *v) + .collect() +} + +#[test] +fn outage_semantics_per_buffer_type() { + let db = db(vec![ + Box::new(EmbassyBuffer::::new_watch()), + Box::new(EmbassyBuffer::::new_mailbox()), + Box::new(EmbassyBuffer::::new_spmc()), + ]); + let mut o = OutboundRoutes::new(&db, "test").unwrap(); + for v in 0..5 { + for key in KEYS { + db.produce::(key, v).unwrap(); + } + } + let got = drain(&mut o); + assert_eq!(values(&got, 0), [4], "single-latest"); + assert_eq!(values(&got, 1), [4], "mailbox"); + assert_eq!(values(&got, 2), [0, 1, 2, 3, 4], "spmc ring"); +} + +#[test] +fn an_spmc_ring_that_overflows_reports_lag_then_recovers() { + let db = db(vec![Box::new(EmbassyBuffer::::new_spmc())]); + let mut o = OutboundRoutes::new(&db, "test").unwrap(); + for v in 0..10 { + db.produce::("r0", v).unwrap(); + } + assert_eq!(values(&drain(&mut o), 0), [6, 7, 8, 9]); + assert_eq!(o.stats(0).unwrap().lagged, 6); +} From 8beaeccb5ba504dca143f775fcb8af372daa442c Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Alexander=20Schn=C3=B6rch?= Date: Mon, 5 Oct 2026 17:38:23 +0000 Subject: [PATCH 3/4] =?UTF-8?q?test(mqtt):=20gate=20allocations=20per=20ro?= =?UTF-8?q?und=20trip=20on=20both=20backends=20(design=20054=20=C2=A75)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Embedded: 0 per round trip. Native: at most 11, rumqttc's own. Co-Authored-By: Claude Opus 5.5 --- Makefile | 4 + .../tests/alloc_round_trip.rs | 264 ++++++++++++++++++ aimdb-mqtt-connector/tests/common/mod.rs | 78 ++++++ 3 files changed, 346 insertions(+) create mode 100644 aimdb-mqtt-connector/tests/alloc_round_trip.rs diff --git a/Makefile b/Makefile index 15881391..dc085e25 100644 --- a/Makefile +++ b/Makefile @@ -249,6 +249,8 @@ test: cargo test --package aimdb-mqtt-connector --no-default-features --features "_test-tokio-broker" --test tokio_broker @printf "$(YELLOW) → Testing MQTT connector (both backends, one broker, one process)$(NC)\n" cargo test --package aimdb-mqtt-connector --no-default-features --features "_test-backend-parity" --test backend_parity + @printf "$(YELLOW) → Testing MQTT connector (allocations per round trip, both backends)$(NC)\n" + cargo test --package aimdb-mqtt-connector --no-default-features --features "_test-backend-parity" --test alloc_round_trip @printf "$(YELLOW) → Testing MQTT connector (mqtts:// against a pinned self-signed root)$(NC)\n" cargo test --package aimdb-mqtt-connector --no-default-features --features "_test-tls-broker" --test tls_broker @printf "$(YELLOW) → Testing MQTT connector (event-driven session: wake cadence, partial packets, QoS 1)$(NC)\n" @@ -400,6 +402,8 @@ clippy: cargo clippy --package aimdb-mqtt-connector --no-default-features --features "_test-tokio-broker" --test tokio_broker -- -D warnings @printf "$(YELLOW) → Clippy on MQTT connector (backend parity)$(NC)\n" cargo clippy --package aimdb-mqtt-connector --no-default-features --features "_test-backend-parity" --test backend_parity -- -D warnings + @printf "$(YELLOW) → Clippy on MQTT connector (allocations per round trip)$(NC)\n" + cargo clippy --package aimdb-mqtt-connector --no-default-features --features "_test-backend-parity" --test alloc_round_trip -- -D warnings @printf "$(YELLOW) → Clippy on MQTT connector (mqtts:// host smoke)$(NC)\n" cargo clippy --package aimdb-mqtt-connector --no-default-features --features "_test-tls-broker" --test tls_broker -- -D warnings @printf "$(YELLOW) → Clippy on MQTT connector (event-driven session criteria)$(NC)\n" diff --git a/aimdb-mqtt-connector/tests/alloc_round_trip.rs b/aimdb-mqtt-connector/tests/alloc_round_trip.rs new file mode 100644 index 00000000..fa3a6381 --- /dev/null +++ b/aimdb-mqtt-connector/tests/alloc_round_trip.rs @@ -0,0 +1,264 @@ +//! Allocations per MQTT round trip, per backend (`_test-backend-parity`). +//! +//! One round trip: produce → PUBLISH QoS 1 → the broker's PUBACK and echo → +//! inbound dispatch → the client's PUBACK → reader `recv`. The database, its +//! connector and the produce/recv loop run on one thread with a current-thread +//! runtime; a counting allocator counts only on that thread, so the broker on +//! another thread is not measured. +//! +//! The embedded backend's one remaining copy per message is topic and payload +//! from `OutboundRoutes`' scratch into the encoded frame; it allocates nothing. +//! The native backend's count is `rumqttc`'s: `AsyncClient::publish` takes an +//! owned topic and payload, and builds its own request. +#![cfg(feature = "_test-backend-parity")] + +use std::alloc::{GlobalAlloc, Layout, System}; +use std::cell::Cell; +use std::sync::Arc; +use std::time::{Duration, Instant}; + +use tokio::net::TcpListener; + +use aimdb_core::buffer::BufferCfg; +use aimdb_core::connector::{ConnectorBuilder, SerializeError}; +use aimdb_core::AimDbBuilder; +use aimdb_mqtt_connector::MqttConnector; +use aimdb_tokio_adapter::net::TokioNet; +use aimdb_tokio_adapter::{TokioAdapter, TokioRecordRegistrarExt}; + +mod common; +use common::echo_broker; + +// Each test binary defines these exactly once. +#[defmt::global_logger] +struct HostTestLogger; +unsafe impl defmt::Logger for HostTestLogger { + fn acquire() {} + unsafe fn flush() {} + unsafe fn release() {} + unsafe fn write(_bytes: &[u8]) {} +} +#[defmt::panic_handler] +fn defmt_panic() -> ! { + core::panic!("defmt panic in host test") +} +defmt::timestamp!("{=u64:us}", 0); + +struct HostClock; +impl embassy_time_driver::Driver for HostClock { + fn now(&self) -> u64 { + use std::sync::OnceLock; + static START: OnceLock = OnceLock::new(); + let start = START.get_or_init(Instant::now); + (start.elapsed().as_micros() * u128::from(embassy_time_driver::TICK_HZ) / 1_000_000) as u64 + } + fn schedule_wake(&self, _at: u64, waker: &core::task::Waker) { + waker.wake_by_ref(); + } +} +embassy_time_driver::time_driver_impl!(static HOST_CLOCK: HostClock = HostClock); + +// --------------------------------------------------------------------------- +// Counting allocator: per-thread counters, so tests may run in parallel. +// --------------------------------------------------------------------------- + +struct Counting; + +thread_local! { + /// Set on the database thread; nothing else is counted. + static COUNT_HERE: Cell = const { Cell::new(false) }; + /// Set only during the measured round trips. + static WINDOW: Cell = const { Cell::new(false) }; + static ALLOCS: Cell = const { Cell::new(0) }; + static BYTES: Cell = const { Cell::new(0) }; + /// Bytes this thread allocated and has not freed. + static LIVE: Cell = const { Cell::new(0) }; +} + +fn add(cell: &'static std::thread::LocalKey>, n: usize) { + cell.with(|c| c.set(c.get() + n)); +} + +unsafe impl GlobalAlloc for Counting { + unsafe fn alloc(&self, layout: Layout) -> *mut u8 { + if COUNT_HERE.with(Cell::get) { + add(&LIVE, layout.size()); + if WINDOW.with(Cell::get) { + add(&ALLOCS, 1); + add(&BYTES, layout.size()); + } + } + unsafe { System.alloc(layout) } + } + unsafe fn dealloc(&self, ptr: *mut u8, layout: Layout) { + if COUNT_HERE.with(Cell::get) { + LIVE.with(|c| c.set(c.get().saturating_sub(layout.size()))); + } + unsafe { System.dealloc(ptr, layout) } + } +} + +#[global_allocator] +static GLOBAL: Counting = Counting; + +// --------------------------------------------------------------------------- +// The round trip. +// --------------------------------------------------------------------------- + +const WARMUP: u64 = 100; +const MEASURED: u64 = 300; +const TOPIC: &str = "mqtt://rt/ping"; + +struct Report { + allocs: usize, + bytes: usize, + live_before: usize, + median: Duration, + min: Duration, + max: Duration, +} + +/// Round trips through `connector` on a fresh thread; `COUNT_HERE` is on for +/// that thread only. +fn measure(connector: impl ConnectorBuilder + 'static) -> Report { + std::thread::spawn(move || { + COUNT_HERE.with(|c| c.set(true)); + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .expect("runtime"); + runtime.block_on(round_trips(connector)) + }) + .join() + .expect("database thread") +} + +async fn round_trips(connector: impl ConnectorBuilder + 'static) -> Report { + let mut builder = AimDbBuilder::new() + .runtime(Arc::new(TokioAdapter)) + .with_connector(connector); + builder.configure::("ping", |reg| { + reg.buffer(BufferCfg::SingleLatest) + .link_to(TOPIC) + .with_serializer(|_ctx, v: &u64| Ok(v.to_le_bytes().to_vec())) + .with_serializer_into(8, |_ctx, v: &u64, out| { + out.get_mut(..8) + .ok_or(SerializeError::BufferTooSmall)? + .copy_from_slice(&v.to_le_bytes()); + Ok(8) + }) + .finish(); + }); + builder.configure::("pong", |reg| { + reg.buffer(BufferCfg::SingleLatest) + .link_from(TOPIC) + .with_deserializer(|_ctx, data: &[u8]| { + data.try_into() + .map(u64::from_le_bytes) + .map_err(|_| String::from("not 8 bytes")) + }) + .finish(); + }); + let (db, runner) = builder.build().await.expect("build db"); + tokio::spawn(runner.run()); + let producer = db.producer::("ping").expect("producer"); + let mut pong = db.subscribe::("pong").expect("subscribe"); + + // Warm-up also waits out the connect and subscribe. + tokio::time::timeout(Duration::from_secs(60), async { + for n in 0..WARMUP { + round_trip(&producer, &mut pong, n).await; + } + }) + .await + .expect("warm-up round trips"); + + let mut latencies = Vec::with_capacity(MEASURED as usize); + let live_before = LIVE.with(Cell::get); + ALLOCS.with(|c| c.set(0)); + BYTES.with(|c| c.set(0)); + WINDOW.with(|w| w.set(true)); + for n in WARMUP..WARMUP + MEASURED { + let start = Instant::now(); + round_trip(&producer, &mut pong, n).await; + latencies.push(start.elapsed()); + } + WINDOW.with(|w| w.set(false)); + + latencies.sort(); + Report { + allocs: ALLOCS.with(Cell::get), + bytes: BYTES.with(Cell::get), + live_before, + median: latencies[latencies.len() / 2], + min: latencies[0], + max: latencies[latencies.len() - 1], + } +} + +/// Produce `n` and wait until it comes back through the broker. +async fn round_trip( + producer: &aimdb_core::Producer, + pong: &mut aimdb_core::buffer::Reader, + n: u64, +) { + producer.produce(n); + while pong.recv().await.expect("pong open") != n {} +} + +/// An echo broker on its own thread; returns its port. +fn broker() -> u16 { + let (tx, rx) = std::sync::mpsc::channel(); + std::thread::spawn(move || { + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .expect("runtime"); + runtime.block_on(async move { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + tx.send(listener.local_addr().unwrap().port()).unwrap(); + echo_broker(listener).await; + }); + }); + rx.recv().expect("broker port") +} + +fn print(backend: &str, r: &Report) { + println!( + "{backend:<8} {:>6.2} allocs/round trip {:>8.1} bytes/round trip live heap before: {} B latency median {:?} ({:?}–{:?})", + r.allocs as f64 / MEASURED as f64, + r.bytes as f64 / MEASURED as f64, + r.live_before, + r.median, + r.min, + r.max + ); +} + +/// The embedded backend allocates nothing per round trip. +#[test] +fn the_embedded_backend_allocates_nothing_per_round_trip() { + let port = broker(); + let report = + measure(MqttConnector::new(format!("mqtt://127.0.0.1:{port}")).transport(TokioNet::tcp())); + print("embedded", &report); + assert_eq!( + report.allocs, 0, + "{} allocations in {MEASURED} round trips", + report.allocs + ); +} + +/// The native backend's allocations are `rumqttc`'s; the bound is what the +/// prototype measured. +#[test] +fn the_native_backend_stays_within_rumqttcs_allocations() { + let port = broker(); + let report = measure(MqttConnector::new(format!("mqtt://127.0.0.1:{port}"))); + print("native", &report); + let per_round_trip = report.allocs as f64 / MEASURED as f64; + assert!( + per_round_trip <= 11.0, + "{per_round_trip:.2} allocations per round trip" + ); +} diff --git a/aimdb-mqtt-connector/tests/common/mod.rs b/aimdb-mqtt-connector/tests/common/mod.rs index b6a28d5f..eb4c12fc 100644 --- a/aimdb-mqtt-connector/tests/common/mod.rs +++ b/aimdb-mqtt-connector/tests/common/mod.rs @@ -716,3 +716,81 @@ impl Delay for CountingDialer { Delay::sleep(&self.inner, d) } } + +// =========================================================================== +// The echo broker: relays a client's publish back to it when it subscribed to +// that topic, as a real broker would. +// =========================================================================== + +/// A QoS 1 PUBLISH in the protocol version the client connected with. +fn publish_qos1_for(topic: &str, payload: &[u8], packet_id: u16, v5: bool) -> Vec { + if v5 { + return publish_qos1(topic, payload, packet_id); + } + let mut rest = Vec::new(); + rest.extend_from_slice(&(topic.len() as u16).to_be_bytes()); + rest.extend_from_slice(topic.as_bytes()); + rest.extend_from_slice(&packet_id.to_be_bytes()); + rest.extend_from_slice(payload); + let mut packet = vec![0x32]; + varint(rest.len(), &mut packet); + packet.extend_from_slice(&rest); + packet +} + +/// Serve one connection: CONNACK, SUBACK (recording the topics), PUBACK every +/// QoS 1 publish, and send a publish on a subscribed topic back at QoS 1. +async fn serve_echo(mut socket: TcpStream) { + let _ = socket.set_nodelay(true); + let mut buf = Vec::new(); + let mut v5 = true; + let mut subscribed: Vec = Vec::new(); + let mut next_id: u16 = 0; + loop { + let Some((first, body)) = read_packet(&mut socket, &mut buf).await else { + return; + }; + let reply: Vec = match first >> 4 { + 1 => { + v5 = is_v5(&body); + if v5 { + vec![0x20, 0x03, 0x00, 0x00, 0x00] + } else { + vec![0x20, 0x02, 0x00, 0x00] + } + } + 8 => suback(&body, v5, &mut subscribed), + 3 => { + let Some((topic, payload, packet_id)) = parse_publish(first, &body, v5) else { + return; + }; + let mut reply = Vec::new(); + if let Some(id) = packet_id { + reply.extend_from_slice(&[0x40, 0x02, id[0], id[1]]); + } + if subscribed.iter().any(|t| t == &topic) { + next_id = next_id.wrapping_add(1).max(1); + reply.extend_from_slice(&publish_qos1_for(&topic, &payload, next_id, v5)); + } + reply + } + 12 => vec![0xD0, 0x00], + 14 => return, + // PUBACKs for the echoes, and anything else: nothing to answer. + _ => continue, + }; + if socket.write_all(&reply).await.is_err() { + return; + } + } +} + +/// Accept forever, serving each connection with [`serve_echo`]. +pub async fn echo_broker(listener: TcpListener) { + loop { + let Ok((socket, _)) = listener.accept().await else { + return; + }; + tokio::spawn(serve_echo(socket)); + } +} From b4d5f42d369314d5d185a9eee1ab7ef9849c5126 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Alexander=20Schn=C3=B6rch?= Date: Mon, 5 Oct 2026 17:55:13 +0000 Subject: [PATCH 4/4] Revert "ci: run CI and the docs check for the 054 feature branch" CI ran green on this PR; the trigger is not kept on the feature branch. This reverts commit 5cf41bc. Co-Authored-By: Claude Opus 5.5 --- .github/workflows/ci.yml | 4 ++-- .github/workflows/docs.yml | 4 ++-- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 7f528614..a4349083 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -2,9 +2,9 @@ name: CI on: push: - branches: [ main, develop, feat/054-connector-boundary ] + branches: [ main, develop ] pull_request: - branches: [ main, develop, feat/054-connector-boundary ] + branches: [ main, develop ] env: CARGO_TERM_COLOR: always diff --git a/.github/workflows/docs.yml b/.github/workflows/docs.yml index 7a55d83d..63086e73 100644 --- a/.github/workflows/docs.yml +++ b/.github/workflows/docs.yml @@ -2,9 +2,9 @@ name: Documentation on: push: - branches: [ main, feat/054-connector-boundary ] + branches: [ main ] pull_request: - branches: [ main, feat/054-connector-boundary ] + branches: [ main ] permissions: contents: read