diff --git a/aimdb-bench/Cargo.toml b/aimdb-bench/Cargo.toml index 06136f8d..1025a4f9 100644 --- a/aimdb-bench/Cargo.toml +++ b/aimdb-bench/Cargo.toml @@ -67,6 +67,12 @@ harness = false name = "b0_alloc_connector" harness = false +# Outbound wake-up time per message, 1 to 256 routes. +# Informational: not in `bench-gate`. +[[bench]] +name = "b1_outbound_wakeup" +harness = false + [features] default = ["std"] # Gates `profiles`/`reports`/`harness` and their criterion/serde_json/ diff --git a/aimdb-bench/benches/b0_alloc_connector.rs b/aimdb-bench/benches/b0_alloc_connector.rs index 45778fcd..b79df06a 100644 --- a/aimdb-bench/benches/b0_alloc_connector.rs +++ b/aimdb-bench/benches/b0_alloc_connector.rs @@ -13,6 +13,11 @@ //! `SerializedReader::recv_into` followed by `Connector::publish` on a no-op //! connector — for the scratch and owned serializers, with a static and a //! dynamic (`TopicProvider`) topic. +//! - **`InboundDispatch` and `OutboundRoutes`:** the same inbound cases through +//! `InboundDispatch::dispatch`, and `OutboundRoutes::next` with a static +//! topic, a written topic and the owned serializer, plus eight routes that +//! are all ready (round-robin) and one pull that parks before every value +//! (the waker path). //! //! Buffers, ingest and routing allocate nothing (design 037, and the `route` //! row here); every non-zero row is a cost of the connector interface. The @@ -37,7 +42,10 @@ use aimdb_core::connector::{ }; use aimdb_core::session::{pump_source, Payload, Source}; use aimdb_core::transport::{Connector, ConnectorConfig, PublishError}; -use aimdb_core::{AimDb, AimDbBuilder, BoxFut, DbResult, ExactGrammar, RuntimeContext, StringKey}; +use aimdb_core::{ + AimDb, AimDbBuilder, BoxFut, DbResult, ExactGrammar, InboundDispatch, OutboundRoutes, + RuntimeContext, StringKey, +}; use aimdb_mqtt_connector::MqttGrammar; use aimdb_tokio_adapter::{TokioAdapter, TokioRecordRegistrarExt}; @@ -63,8 +71,20 @@ const EXPECTED: &[(&str, u64)] = &[ ("outbound_scratch_static_topic", 2), ("outbound_scratch_dynamic_topic", 3), ("outbound_owned_static_topic", 3), + ("inbound_dispatch", 0), + ("inbound_dispatch_pattern", 0), + ("inbound_dispatch_keyed_known", 0), + ("inbound_dispatch_keyed_new", 1), + ("outbound_next_static_topic", 0), + ("outbound_next_written_topic", 0), + ("outbound_next_owned", 1), + ("outbound_next_round_robin", 0), + ("outbound_next_parked", 0), ]; +/// Routes in the `outbound_next_round_robin` row. +const ROUND_ROBIN_ROUTES: usize = 8; + #[derive(Clone, Copy, Debug)] struct Reading { id: u32, @@ -242,6 +262,40 @@ async fn measure_route() -> (u64, u64) { snapshot() } +/// `measure_route` through `InboundDispatch`. +async fn measure_dispatch() -> (u64, u64) { + let db = inbound_db().await; + let inbound = InboundDispatch::new(&db, SCHEME, &ExactGrammar).unwrap(); + let payload = [1u8; 8]; + for _ in 0..WARMUP_ITERS { + inbound.dispatch("in/target", &payload); + } + reset(); + for _ in 0..MEASURE_ITERS { + inbound.dispatch("in/target", black_box(&payload)); + } + snapshot() +} + +/// `measure_pattern_route` through `InboundDispatch`. +async fn measure_pattern_dispatch( + keyed: bool, + warmup: &[String], + measured: &[String], +) -> (u64, u64) { + let db = pattern_db(keyed).await; + let inbound = InboundDispatch::new(&db, SCHEME, &MqttGrammar).unwrap(); + let payload = [1u8; 8]; + for topic in warmup { + inbound.dispatch(topic, &payload); + } + reset(); + for topic in measured { + inbound.dispatch(black_box(topic), black_box(&payload)); + } + snapshot() +} + /// Yields `remaining` copies of one message, then ends. struct MinimalSource { topic: String, @@ -372,6 +426,158 @@ async fn measure_outbound(serializer: Serializer, dynamic_topic: bool) -> (u64, snapshot() } +// --- OutboundRoutes --------------------------------------------------------- + +#[derive(Clone, Copy)] +enum Topic { + Static, + Written, +} + +/// One outbound link on `record` with the given serializer and topic. +fn outbound_link( + reg: &mut aimdb_core::RecordRegistrar<'_, Reading>, + serializer: Serializer, + topic: Topic, +) { + reg.buffer(BufferCfg::SpmcRing { capacity: 64 }); + let mut link = reg + .link_to("bench://out/default") + .with_serializer(|_ctx, r: &Reading| Ok(encode(r).to_vec())); + if let Serializer::Scratch = serializer { + link = link.with_serializer_into(SCRATCH_CAPACITY, |_ctx, r: &Reading, buf| { + let bytes = encode(r); + let dst = buf + .get_mut(..bytes.len()) + .ok_or(SerializeError::BufferTooSmall)?; + dst.copy_from_slice(&bytes); + Ok(bytes.len()) + }); + } + if let Topic::Written = topic { + link = link.with_topic_fn(16, |r, out| { + use std::fmt::Write as _; + write!(out, "out/{}", r.id)?; + Ok(true) + }); + } + link.finish(); +} + +/// Pulls one message and hands its topic and bytes to `black_box`. +async fn pull_one(outbound: &mut OutboundRoutes) { + let msg = outbound.next().await.expect("route open"); + black_box((msg.topic, msg.payload.as_slice())); +} + +async fn measure_next(serializer: Serializer, topic: Topic) -> (u64, u64) { + let db = build_db(|b| { + b.configure::("out.record", move |reg| { + outbound_link(reg, serializer, topic) + }); + }) + .await; + let mut outbound = OutboundRoutes::new(&db, SCHEME).unwrap(); + let producer = db.producer::("out.record").expect("producer"); + + for i in 0..WARMUP_ITERS { + producer.produce(reading(i)); + pull_one(&mut outbound).await; + } + reset(); + for i in 0..MEASURE_ITERS { + producer.produce(reading(i)); + pull_one(&mut outbound).await; + } + snapshot() +} + +/// One value on each of `ROUND_ROBIN_ROUTES` routes, then as many pulls. +async fn measure_next_round_robin() -> (u64, u64) { + let db = build_db(|b| { + for i in 0..ROUND_ROBIN_ROUTES { + b.configure::(StringKey::intern(format!("out.rr{i}")), |reg| { + outbound_link(reg, Serializer::Scratch, Topic::Static) + }); + } + }) + .await; + let mut outbound = OutboundRoutes::new(&db, SCHEME).unwrap(); + let producers: Vec<_> = (0..ROUND_ROBIN_ROUTES) + .map(|i| { + db.producer::(format!("out.rr{i}")) + .expect("producer") + }) + .collect(); + for i in 0..WARMUP_ITERS / ROUND_ROBIN_ROUTES { + round_robin(&mut outbound, &producers, i).await; + } + reset(); + for i in 0..MEASURE_ITERS / ROUND_ROBIN_ROUTES { + round_robin(&mut outbound, &producers, i).await; + } + snapshot() +} + +async fn round_robin( + outbound: &mut OutboundRoutes, + producers: &[aimdb_core::Producer], + i: usize, +) { + for p in producers { + p.produce(reading(i)); + } + for _ in 0..producers.len() { + pull_one(outbound).await; + } +} + +/// The transport pulls in its own task and parks before every value; the +/// producer writes one value, then yields until it was pulled. +async fn measure_next_parked() -> (u64, u64) { + use std::sync::atomic::{AtomicUsize, Ordering}; + + let db = build_db(|b| { + b.configure::("out.record", |reg| { + outbound_link(reg, Serializer::Scratch, Topic::Static) + }); + }) + .await; + let mut outbound = OutboundRoutes::new(&db, SCHEME).unwrap(); + let producer = db.producer::("out.record").expect("producer"); + let pulled = Arc::new(AtomicUsize::new(0)); + let total = WARMUP_ITERS + MEASURE_ITERS; + let transport = { + let pulled = pulled.clone(); + tokio::spawn(async move { + for _ in 0..total { + pull_one(&mut outbound).await; + pulled.fetch_add(1, Ordering::Release); + } + }) + }; + let send = |i: usize| { + producer.produce(reading(i)); + let pulled = pulled.clone(); + async move { + while pulled.load(Ordering::Acquire) <= i { + tokio::task::yield_now().await; + } + } + }; + + for i in 0..WARMUP_ITERS { + send(i).await; + } + reset(); + for i in WARMUP_ITERS..total { + send(i).await; + } + let counted = snapshot(); + transport.await.expect("transport"); + counted +} + // --- Driver ----------------------------------------------------------------- fn main() { @@ -432,6 +638,62 @@ fn main() { "SpmcRing", measure_outbound(Serializer::Owned, false).await, ), + ("inbound_dispatch", "SpmcRing", measure_dispatch().await), + ( + "inbound_dispatch_pattern", + "SpmcRing", + measure_pattern_dispatch( + false, + &pattern_topics(WARMUP_ITERS, 0, false), + &pattern_topics(MEASURE_ITERS, 0, false), + ) + .await, + ), + ( + "inbound_dispatch_keyed_known", + "SpmcRing", + measure_pattern_dispatch( + true, + &pattern_topics(WARMUP_ITERS, 0, false), + &pattern_topics(MEASURE_ITERS, 0, false), + ) + .await, + ), + ( + "inbound_dispatch_keyed_new", + "SpmcRing", + measure_pattern_dispatch( + true, + &pattern_topics(WARMUP_ITERS, 0, true), + &pattern_topics(MEASURE_ITERS, WARMUP_ITERS, true), + ) + .await, + ), + ( + "outbound_next_static_topic", + "SpmcRing", + measure_next(Serializer::Scratch, Topic::Static).await, + ), + ( + "outbound_next_written_topic", + "SpmcRing", + measure_next(Serializer::Scratch, Topic::Written).await, + ), + ( + "outbound_next_owned", + "SpmcRing", + measure_next(Serializer::Owned, Topic::Static).await, + ), + ( + "outbound_next_round_robin", + "SpmcRing", + measure_next_round_robin().await, + ), + ( + "outbound_next_parked", + "SpmcRing", + measure_next_parked().await, + ), ] }); diff --git a/aimdb-bench/benches/b1_outbound_wakeup.rs b/aimdb-bench/benches/b1_outbound_wakeup.rs new file mode 100644 index 00000000..bc769a36 --- /dev/null +++ b/aimdb-bench/benches/b1_outbound_wakeup.rs @@ -0,0 +1,240 @@ +//! B1-Outbound — time per message when the transport parks before every +//! value. Informational: not part of `bench-gate`, since timings need a quiet +//! host. +//! +//! One busy route among N (1, 8, 64, 256), the producer and the transport in +//! separate tasks on a current-thread Tokio runtime. The producer writes one +//! value and waits until it was pulled, so every message pays one park and +//! one wake-up of the transport. Two columns: +//! +//! - **task per route:** one task per route awaiting `Reader::recv`, the +//! shape of the per-route pumps. +//! - **OutboundRoutes:** one task pulling with `OutboundRoutes::next`, which +//! polls only the routes that woke. +//! +//! A per-route scan would show as `OutboundRoutes` growing with N. Compare +//! columns within one run, not across hosts. +//! +//! Run `cargo bench -p aimdb-bench --bench b1_outbound_wakeup`. + +use std::future::poll_fn; +use std::hint::black_box; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::sync::Arc; +use std::task::Poll; +use std::time::Instant; + +use aimdb_bench::alloc::{reset, snapshot}; +use aimdb_core::buffer::BufferCfg; +use aimdb_core::connector::{ConnectorBuilder, SerializeError}; +use aimdb_core::{AimDb, AimDbBuilder, DbResult, OutboundRoutes, Producer, StringKey}; +use aimdb_tokio_adapter::{TokioAdapter, TokioRecordRegistrarExt}; +use futures::task::AtomicWaker; + +#[global_allocator] +static GLOBAL: aimdb_bench::alloc::CountingAllocator = + aimdb_bench::alloc::CountingAllocator(std::alloc::System); + +const ROUTE_COUNTS: &[usize] = &[1, 8, 64, 256]; +const WARMUP: u64 = 2_000; +const MEASURE: u64 = 20_000; +const RUNS: usize = 5; +const RING: usize = 64; +const SCHEME: &str = "bench"; + +/// Registers the scheme so `link_to` succeeds; drives nothing. +struct NoopConnectorBuilder; + +impl ConnectorBuilder for NoopConnectorBuilder { + #[allow(clippy::type_complexity)] + fn build<'a>( + &'a self, + _db: &'a AimDb, + ) -> std::pin::Pin< + Box< + dyn std::future::Future< + Output = DbResult< + Vec + Send>>>, + >, + > + Send + + 'a, + >, + > { + Box::pin(async { Ok(Vec::new()) }) + } + + fn scheme(&self) -> &str { + SCHEME + } +} + +/// `n` SPMC records `r0..r{n}`, each linked to `bench://r{i}` with a scratch +/// serializer. +async fn build_db(n: usize) -> AimDb { + let runtime = Arc::new(TokioAdapter::new().expect("tokio adapter")); + let mut builder = AimDbBuilder::new() + .runtime(runtime) + .with_connector(NoopConnectorBuilder); + for i in 0..n { + builder.configure::(StringKey::intern(format!("r{i}")), move |reg| { + reg.buffer(BufferCfg::SpmcRing { capacity: RING }) + .link_to(&format!("{SCHEME}://r{i}")) + .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(); + }); + } + let (db, runner) = builder.build().await.expect("build"); + // Nothing in the runner is needed; keep it alive for the process. + std::mem::forget(runner); + db +} + +/// Values pulled so far, and the producer's waker. +struct Ack { + count: AtomicU64, + waker: AtomicWaker, +} + +impl Ack { + fn new() -> Arc { + Arc::new(Self { + count: AtomicU64::new(0), + waker: AtomicWaker::new(), + }) + } + + fn bump(&self) { + self.count.fetch_add(1, Ordering::Release); + self.waker.wake(); + } + + async fn wait_for(&self, target: u64) { + poll_fn(|cx| { + self.waker.register(cx.waker()); + if self.count.load(Ordering::Acquire) >= target { + Poll::Ready(()) + } else { + Poll::Pending + } + }) + .await + } +} + +/// Writes one value, waits until it was pulled. Returns ns and allocations +/// per message over the measured window. +async fn drive_producer(producer: Producer, ack: Arc) -> (f64, f64) { + for i in 0..WARMUP { + producer.produce(i); + ack.wait_for(i + 1).await; + } + reset(); + let start = Instant::now(); + for i in WARMUP..WARMUP + MEASURE { + producer.produce(i); + ack.wait_for(i + 1).await; + } + let ns = start.elapsed().as_nanos() as f64 / MEASURE as f64; + let (allocs, _) = snapshot(); + (ns, allocs as f64 / MEASURE as f64) +} + +async fn task_per_route(n: usize) -> (f64, f64) { + let db = build_db(n).await; + let busy = n - 1; + let producer = db.producer::(format!("r{busy}")).expect("producer"); + let ack = Ack::new(); + let mut tasks = Vec::new(); + for i in 0..n { + let mut reader = db.subscribe::(format!("r{i}")).expect("subscribe"); + let ack = ack.clone(); + tasks.push(tokio::spawn(async move { + loop { + let v = reader.recv().await.expect("recv"); + black_box((i, v)); + ack.bump(); + if v + 1 == WARMUP + MEASURE { + break; + } + } + })); + } + let result = drive_producer(producer, ack).await; + tasks.pop().expect("busy task").await.expect("busy task"); + for task in tasks { + task.abort(); + } + result +} + +async fn outbound_routes(n: usize) -> (f64, f64) { + let db = build_db(n).await; + let busy = n - 1; + let producer = db.producer::(format!("r{busy}")).expect("producer"); + let mut outbound = OutboundRoutes::new(&db, SCHEME).expect("routes"); + let ack = Ack::new(); + let transport = { + let ack = ack.clone(); + tokio::spawn(async move { + loop { + let msg = outbound.next().await.expect("route open"); + let v = u64::from_le_bytes(msg.payload.as_slice().try_into().expect("8 bytes")); + black_box((msg.route.id, msg.topic, v)); + ack.bump(); + if v + 1 == WARMUP + MEASURE { + break; + } + } + }) + }; + let result = drive_producer(producer, ack).await; + transport.await.expect("transport"); + result +} + +fn median(mut v: Vec) -> f64 { + v.sort_by(|a, b| a.total_cmp(b)); + v[v.len() / 2] +} + +/// Median of `RUNS` runs, each on a fresh database. +fn row(runtime: &tokio::runtime::Runtime, name: &str, n: usize, run: F) -> f64 +where + F: Fn(usize) -> Fut, + Fut: std::future::Future, +{ + let mut ns = Vec::with_capacity(RUNS); + let mut allocs = 0.0; + for _ in 0..RUNS { + let (t, a) = runtime.block_on(run(n)); + ns.push(t); + allocs = a; + } + let lo = ns.iter().copied().fold(f64::MAX, f64::min); + let hi = ns.iter().copied().fold(0.0, f64::max); + let m = median(ns); + println!( + "{name:<16} {n:>4} routes {m:>7.0} ns/msg ({lo:.0}–{hi:.0}) {allocs:.3} allocs/msg" + ); + m +} + +fn main() { + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .expect("runtime"); + + println!("=== B1 outbound wake-up: one busy route among N, current-thread runtime ==="); + for &n in ROUTE_COUNTS { + row(&runtime, "task per route", n, task_per_route); + row(&runtime, "OutboundRoutes", n, outbound_routes); + println!(); + } +} diff --git a/aimdb-bench/data/baselines/b0_alloc_connector.json b/aimdb-bench/data/baselines/b0_alloc_connector.json index fbd67fcf..ae351427 100644 --- a/aimdb-bench/data/baselines/b0_alloc_connector.json +++ b/aimdb-bench/data/baselines/b0_alloc_connector.json @@ -70,5 +70,86 @@ "batch_size": 2000, "allocs_per_msg": 3.0, "bytes_per_msg": 81.0 + }, + { + "profile": "inbound_dispatch", + "buffer_type": "SpmcRing", + "total_allocs": 0, + "total_bytes": 0, + "batch_size": 2000, + "allocs_per_msg": 0.0, + "bytes_per_msg": 0.0 + }, + { + "profile": "inbound_dispatch_pattern", + "buffer_type": "SpmcRing", + "total_allocs": 0, + "total_bytes": 0, + "batch_size": 2000, + "allocs_per_msg": 0.0, + "bytes_per_msg": 0.0 + }, + { + "profile": "inbound_dispatch_keyed_known", + "buffer_type": "SpmcRing", + "total_allocs": 0, + "total_bytes": 0, + "batch_size": 2000, + "allocs_per_msg": 0.0, + "bytes_per_msg": 0.0 + }, + { + "profile": "inbound_dispatch_keyed_new", + "buffer_type": "SpmcRing", + "total_allocs": 2010, + "total_bytes": 373456, + "batch_size": 2000, + "allocs_per_msg": 1.005, + "bytes_per_msg": 186.728 + }, + { + "profile": "outbound_next_static_topic", + "buffer_type": "SpmcRing", + "total_allocs": 0, + "total_bytes": 0, + "batch_size": 2000, + "allocs_per_msg": 0.0, + "bytes_per_msg": 0.0 + }, + { + "profile": "outbound_next_written_topic", + "buffer_type": "SpmcRing", + "total_allocs": 0, + "total_bytes": 0, + "batch_size": 2000, + "allocs_per_msg": 0.0, + "bytes_per_msg": 0.0 + }, + { + "profile": "outbound_next_owned", + "buffer_type": "SpmcRing", + "total_allocs": 2000, + "total_bytes": 16000, + "batch_size": 2000, + "allocs_per_msg": 1.0, + "bytes_per_msg": 8.0 + }, + { + "profile": "outbound_next_round_robin", + "buffer_type": "SpmcRing", + "total_allocs": 0, + "total_bytes": 0, + "batch_size": 2000, + "allocs_per_msg": 0.0, + "bytes_per_msg": 0.0 + }, + { + "profile": "outbound_next_parked", + "buffer_type": "SpmcRing", + "total_allocs": 0, + "total_bytes": 0, + "batch_size": 2000, + "allocs_per_msg": 0.0, + "bytes_per_msg": 0.0 } ] \ No newline at end of file