Skip to content
Merged
Show file tree
Hide file tree
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
6 changes: 6 additions & 0 deletions aimdb-bench/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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/
Expand Down
264 changes: 263 additions & 1 deletion aimdb-bench/benches/b0_alloc_connector.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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};

Expand All @@ -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,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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::<Reading>("out.record", move |reg| {
outbound_link(reg, serializer, topic)
});
})
.await;
let mut outbound = OutboundRoutes::new(&db, SCHEME).unwrap();
let producer = db.producer::<Reading>("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::<Reading>(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::<Reading>(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<Reading>],
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::<Reading>("out.record", |reg| {
outbound_link(reg, Serializer::Scratch, Topic::Static)
});
})
.await;
let mut outbound = OutboundRoutes::new(&db, SCHEME).unwrap();
let producer = db.producer::<Reading>("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() {
Expand Down Expand Up @@ -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,
),
]
});

Expand Down
Loading
Loading