use std::net::SocketAddr;
use criterion::{Criterion, Throughput, criterion_group, criterion_main};
use std::hint::black_box;
use weida::{Identity, Listener, Publisher, Puller, Runtime, RuntimeConfig, TransferMeta, Trust};
use weida_protocol::{DataHeader, FrameKind, ReportMode, encode_frame};
const PAYLOAD: usize = 1024;
struct Harness {
runtime: Runtime,
trust: Trust,
addr: SocketAddr,
listener: Listener,
_binding: weida::Binding,
}
impl Harness {
fn client(&self) -> Runtime {
Runtime::new(RuntimeConfig::default()).expect("client runtime")
}
fn trust(&self) -> Trust {
self.trust.clone()
}
fn url(&self, path: &str) -> String {
format!("weida://127.0.0.1:{}{}", self.addr.port(), path)
}
}
async fn harness() -> Harness {
let identity = Identity::generate().expect("identity");
let trust = Trust::pin(identity.fingerprint().expect("fingerprint"));
let runtime = Runtime::new(RuntimeConfig::default()).expect("runtime");
let listener = runtime.listener();
let binding = listener
.bind_quic("127.0.0.1:0".parse().expect("loopback"), identity)
.await
.expect("bind");
let addr = binding.local_addr();
Harness {
runtime,
trust,
addr,
listener,
_binding: binding,
}
}
fn tokio_runtime() -> tokio::runtime::Runtime {
tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()
.expect("tokio runtime")
}
fn drain(puller: Puller) {
tokio::spawn(async move {
while let Ok(transfer) = puller.recv().await {
let _ = transfer.collect(64 * 1024).await;
}
});
}
fn bench_push(c: &mut Criterion) {
let rt = tokio_runtime();
let harness = rt.block_on(harness());
let (best_effort, delivered) = rt.block_on(async {
drain(harness.listener.puller("/best").expect("puller best"));
drain(
harness
.listener
.puller("/delivered")
.expect("puller delivered"),
);
let client = harness.client();
let best_effort = client.pusher(harness.trust());
best_effort
.connect(&harness.url("/best"))
.await
.expect("connect best");
let delivered = client.pusher(harness.trust());
delivered
.connect(&harness.url("/delivered"))
.await
.expect("connect delivered");
(best_effort, delivered)
});
let payload = vec![0x61u8; PAYLOAD];
let mut group = c.benchmark_group("push");
group.throughput(Throughput::Bytes(PAYLOAD as u64));
group.bench_function("push_1kib_best_effort", |b| {
b.to_async(&rt).iter(|| async {
best_effort.send(black_box(&payload)).await.expect("send");
})
});
group.bench_function("push_1kib_delivered", |b| {
b.to_async(&rt).iter(|| async {
let mut transfer = delivered.open(TransferMeta::default()).await.expect("open");
transfer
.write_all(black_box(&payload))
.await
.expect("write");
transfer
.finish()
.expect("finish")
.delivered()
.await
.expect("delivered");
})
});
group.finish();
rt.block_on(async { harness.runtime.clone().shutdown().await });
}
fn bench_fanout(c: &mut Criterion) {
const SUBSCRIBERS: usize = 8;
let rt = tokio_runtime();
let harness = rt.block_on(harness());
let (publisher, mut subscribers) = rt.block_on(async {
let publisher: Publisher = harness.listener.publisher("/md").expect("publisher");
let mut subscribers = Vec::with_capacity(SUBSCRIBERS);
for _ in 0..SUBSCRIBERS {
let client = harness.client();
let sub = client.subscriber(harness.trust());
sub.connect(&harness.url("/md")).await.expect("connect");
sub.subscribe("").await.expect("subscribe");
subscribers.push((client, sub));
}
while publisher.filter_count() != SUBSCRIBERS {
tokio::time::sleep(std::time::Duration::from_millis(1)).await;
}
(publisher, subscribers)
});
let payload = vec![0x7au8; PAYLOAD];
let mut group = c.benchmark_group("fanout");
group.throughput(Throughput::Bytes((SUBSCRIBERS * PAYLOAD) as u64));
group.sample_size(10);
group.bench_function("pub_1kib_8_subscribers", |b| {
b.to_async(&rt).iter(|| async {
let sent = publisher
.publish("px.eur", black_box(payload.clone()))
.expect("publish");
assert_eq!(sent, SUBSCRIBERS, "a subscriber dropped the message");
for (_, sub) in &subscribers {
let transfer = sub.recv().await.expect("recv");
let body = transfer.collect(64 * 1024).await.expect("collect");
debug_assert_eq!(body.len(), PAYLOAD);
}
black_box(sent)
})
});
group.finish();
subscribers.clear();
rt.block_on(async { harness.runtime.clone().shutdown().await });
}
fn bench_filters(c: &mut Criterion) {
const FILTERS: usize = 64;
const TOPIC: &str = "px.eur.spot";
let rt = tokio_runtime();
let harness = rt.block_on(harness());
let shapes: [(&str, usize); 5] = [
("literal_x1", 1),
("literal_x64", FILTERS),
("one_segment_x1", 1),
("one_segment_x64", FILTERS),
("rest_x64", FILTERS),
];
let mut group = c.benchmark_group("filters");
for (name, count) in shapes {
let filters: Vec<String> = (0..count)
.map(|i| {
if name.starts_with("literal") {
format!("zz{i}.eur.spot")
} else if name.starts_with("one_segment") {
format!("zz{i}.*.spot")
} else {
format!("zz{i}.eur.#")
}
})
.collect();
let path = format!("/md-{name}");
let (publisher, _client, _sub) = rt.block_on(async {
let publisher: Publisher = harness.listener.publisher(&path).expect("publisher");
let client = harness.client();
let sub = client.subscriber(harness.trust());
sub.connect(&harness.url(&path)).await.expect("connect");
for filter in &filters {
sub.subscribe(filter).await.expect("subscribe");
}
while publisher.filter_count() != count {
tokio::time::sleep(std::time::Duration::from_millis(1)).await;
}
(publisher, client, sub)
});
group.bench_function(name, |b| {
b.iter(|| {
let sent = publisher
.publish(black_box(TOPIC), &b"x"[..])
.expect("publish");
assert_eq!(sent, 0, "the bench topic must match no filter");
sent
})
});
}
let prefix_filters: Vec<String> = (0..FILTERS).map(|i| format!("zz{i}.eur.spot")).collect();
group.bench_function("prefix_reference", |b| {
b.iter(|| {
let mut matched = 0usize;
for filter in &prefix_filters {
if black_box(TOPIC).starts_with(filter.as_str()) {
matched += 1;
}
}
matched
})
});
group.finish();
rt.block_on(async { harness.runtime.clone().shutdown().await });
}
fn bench_header_cost(c: &mut Criterion) {
const SMALL: usize = 64;
const SEQUENCE: u64 = 1 << 20;
const PRODUCER: &str =
"sha256:22ed30a800000000000000000000000000000000000000000000000000009f25";
assert_eq!(PRODUCER.len(), 71, "the producer key is a 71-byte tstr");
let rt = tokio_runtime();
let harness = rt.block_on(harness());
let (lean_pusher, rich_pusher, traced_pusher) = rt.block_on(async {
drain(harness.listener.puller("/keys0").expect("puller keys0"));
drain(harness.listener.puller("/keys2").expect("puller keys2"));
drain(harness.listener.puller("/traced").expect("puller traced"));
let client = harness.client();
let lean = client.pusher(harness.trust());
lean.connect(&harness.url("/keys0"))
.await
.expect("connect keys0");
let rich = client.pusher(harness.trust());
rich.connect(&harness.url("/keys2"))
.await
.expect("connect keys2");
let traced = client.pusher(harness.trust());
traced
.connect(&harness.url("/traced"))
.await
.expect("connect traced");
(lean, rich, traced)
});
let payload = vec![0x61u8; SMALL];
let lean_meta = TransferMeta::default();
let rich_meta = TransferMeta::default()
.with_content_len(SEQUENCE)
.with_content_type(PRODUCER);
let traced_meta = TransferMeta::default().with_trace(weida::new_trace());
let lean_bytes = wire_bytes("/keys0", &lean_meta, SMALL);
let rich_bytes = wire_bytes("/keys2", &rich_meta, SMALL);
let traced_bytes = wire_bytes("/traced", &traced_meta, SMALL);
eprintln!(
"header cost: {lean_bytes} B/message minimal, {rich_bytes} B/message with two extra keys \
(+{} B, +{:.1} % over a {SMALL}-byte payload), {traced_bytes} B/message with a \
caller-supplied traceparent (+{} B)",
rich_bytes - lean_bytes,
(rich_bytes - lean_bytes) as f64 * 100.0 / lean_bytes as f64,
traced_bytes - lean_bytes,
);
let mut group = c.benchmark_group("header");
group.throughput(Throughput::Elements(1));
for (name, meta, pusher) in [
("push_64b_keys0", &lean_meta, &lean_pusher),
("push_64b_keys2", &rich_meta, &rich_pusher),
("push_64b_traced", &traced_meta, &traced_pusher),
] {
group.bench_function(name, |b| {
b.to_async(&rt).iter(|| async {
let mut transfer = pusher.open(meta.clone()).await.expect("open");
transfer
.write_all(black_box(&payload))
.await
.expect("write");
transfer.finish().expect("finish");
})
});
}
group.finish();
rt.block_on(async { harness.runtime.clone().shutdown().await });
}
fn wire_bytes(endpoint: &str, meta: &TransferMeta, payload: usize) -> usize {
let header = DataHeader {
endpoint: Some(endpoint.to_owned()),
content_len: meta.content_len,
content_type: meta.content_type.clone(),
traceparent: meta.trace.map(|t| t.to_traceparent()),
tracestate: None,
topic: None,
sequence: None,
producer: None,
achieved: meta.achieved,
report_id: None,
report: Vec::new(),
report_mode: ReportMode::default(),
};
encode_frame(FrameKind::Data, &header.encode()).len() + payload
}
fn bench_dedup_key(c: &mut Criterion) {
use std::collections::HashMap;
const ENTRIES: usize = 4096;
const SCOPES: usize = 8;
#[derive(Clone, PartialEq, Eq, Hash)]
struct Owned {
producer: Option<[u8; 32]>,
scope: Box<str>,
sequence: u64,
}
let scopes: Vec<String> = (0..SCOPES).map(|i| format!("/md/instrument-{i}")).collect();
let producer = Some([7u8; 32]);
type BySequence = HashMap<(Option<[u8; 32]>, u64), u64>;
let mut owned: HashMap<Owned, u64> = HashMap::new();
let mut borrowed: HashMap<Box<str>, BySequence> = HashMap::new();
for i in 0..ENTRIES {
let scope = &scopes[i % SCOPES];
let sequence = i as u64;
owned.insert(
Owned {
producer,
scope: scope.as_str().into(),
sequence,
},
sequence,
);
borrowed
.entry(scope.as_str().into())
.or_default()
.insert((producer, sequence), sequence);
}
let present = (ENTRIES / 2) as u64;
let absent = ENTRIES as u64 * 3;
let scope = scopes[(present as usize) % SCOPES].as_str();
let mut group = c.benchmark_group("dedup_key");
group.bench_function("owned_hit", |b| {
b.iter(|| {
let key = Owned {
producer,
scope: black_box(scope).into(),
sequence: black_box(present),
};
black_box(owned.contains_key(&key))
})
});
group.bench_function("owned_miss", |b| {
b.iter(|| {
let key = Owned {
producer,
scope: black_box(scope).into(),
sequence: black_box(absent),
};
black_box(owned.contains_key(&key))
})
});
let prebuilt_hit = Owned {
producer,
scope: scope.into(),
sequence: present,
};
let prebuilt_miss = Owned {
producer,
scope: scope.into(),
sequence: absent,
};
group.bench_function("prebuilt_hit", |b| {
b.iter(|| black_box(owned.contains_key(black_box(&prebuilt_hit))))
});
group.bench_function("prebuilt_miss", |b| {
b.iter(|| black_box(owned.contains_key(black_box(&prebuilt_miss))))
});
group.bench_function("borrowed_hit", |b| {
b.iter(|| {
black_box(borrowed.get(black_box(scope)).is_some_and(|by_sequence| {
by_sequence.contains_key(&(producer, black_box(present)))
}))
})
});
group.bench_function("borrowed_miss", |b| {
b.iter(|| {
black_box(borrowed.get(black_box(scope)).is_some_and(|by_sequence| {
by_sequence.contains_key(&(producer, black_box(absent)))
}))
})
});
group.finish();
}
criterion_group!(
benches,
bench_push,
bench_fanout,
bench_filters,
bench_header_cost,
bench_dedup_key
);
criterion_main!(benches);