use std::net::SocketAddr;
use std::sync::Arc;
use std::time::{Duration, Instant};
use weida::{Error, Identity, Limits, Runtime, RuntimeConfig, Subscriber, Trust};
pub struct Bound {
pub url: String,
pub runtime: Runtime,
pub listener: weida::Listener,
_binding: weida::Binding,
}
pub async fn bind_with(path: &str, config: RuntimeConfig) -> Result<Bound, Error> {
let identity = Identity::generate()?;
let fingerprint = identity.fingerprint()?;
let runtime = Runtime::new(config)?;
let listener = runtime.listener();
let binding = listener
.bind_quic(
"127.0.0.1:0"
.parse::<SocketAddr>()
.expect("a loopback literal"),
identity,
)
.await?;
let url = format!(
"weida://{fingerprint}@127.0.0.1:{}{path}",
binding.local_addr().port()
);
Ok(Bound {
url,
runtime,
listener,
_binding: binding,
})
}
pub async fn bind(path: &str) -> Result<Bound, Error> {
bind_with(path, RuntimeConfig::default()).await
}
pub async fn subscribers_with(
url: &str,
count: usize,
config: RuntimeConfig,
) -> Result<Vec<(Runtime, Arc<Subscriber>)>, Error> {
let mut set = Vec::with_capacity(count);
for _ in 0..count {
let client = Runtime::new(config.clone())?;
let sub = client.subscriber(Trust::by_address());
sub.connect(url).await?;
sub.subscribe("").await?;
set.push((client, Arc::new(sub)));
}
Ok(set)
}
pub async fn subscribers(
url: &str,
count: usize,
) -> Result<Vec<(Runtime, Arc<Subscriber>)>, Error> {
subscribers_with(url, count, RuntimeConfig::default()).await
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct Selection {
pub push_per_peer: Vec<usize>,
pub fanout_per_subscriber: Vec<usize>,
pub bus_per_member: Vec<usize>,
}
pub async fn selection() -> Result<Selection, Error> {
const PUSHES: usize = 6;
let mut servers = Vec::new();
for _ in 0..3 {
servers.push(bind("/work").await?);
}
let mut pullers = Vec::new();
for server in &servers {
pullers.push(server.listener.puller("/work")?);
}
let client = Runtime::new(RuntimeConfig::default())?;
let pusher = client.pusher(Trust::by_address());
for server in &servers {
pusher.connect(&server.url).await?;
}
for message in 0..PUSHES {
pusher.send(format!("job {message}").as_bytes()).await?;
}
let mut push_per_peer = Vec::with_capacity(pullers.len());
for puller in &pullers {
let mut mine = 0;
while let Ok(Ok(transfer)) =
tokio::time::timeout(Duration::from_millis(250), puller.recv()).await
{
transfer.collect(1024).await?;
mine += 1;
}
push_per_peer.push(mine);
}
client.shutdown().await;
for server in servers {
server.runtime.shutdown().await;
}
let bound = bind("/prices").await?;
let publisher = bound.listener.publisher("/prices")?;
let subs = subscribers(&bound.url, 3).await?;
while publisher.filter_count() != subs.len() {
tokio::time::sleep(Duration::from_millis(1)).await;
}
for message in 0..2 {
publisher.publish("px.eur", format!("tick {message}"))?;
}
let mut fanout_per_subscriber = Vec::with_capacity(subs.len());
for (_, sub) in &subs {
let mut mine = 0;
while let Ok(Ok(transfer)) =
tokio::time::timeout(Duration::from_millis(250), sub.recv()).await
{
transfer.collect(1024).await?;
mine += 1;
}
fanout_per_subscriber.push(mine);
}
drop(subs);
bound.runtime.shutdown().await;
let mut buses = Vec::new();
for index in 0..3 {
buses.push(bind(&format!("/bus{index}")).await?);
}
let mut members = Vec::new();
for (index, bus) in buses.iter().enumerate() {
members.push(
bus.listener
.bus(&format!("/bus{index}"), Trust::by_address())?,
);
}
for (index, member) in members.iter().enumerate() {
for (other, bus) in buses.iter().enumerate() {
if other != index {
member.connect(&bus.url).await?;
}
}
}
for (index, member) in members.iter().enumerate() {
member
.send(format!("hello from {index}").as_bytes())
.await?;
}
let mut bus_per_member = Vec::with_capacity(members.len());
for member in &members {
let mut mine = 0;
while let Ok(Ok(transfer)) =
tokio::time::timeout(Duration::from_millis(250), member.recv()).await
{
transfer.collect(1024).await?;
mine += 1;
}
bus_per_member.push(mine);
}
drop(members);
for bus in buses {
bus.runtime.shutdown().await;
}
Ok(Selection {
push_per_peer,
fanout_per_subscriber,
bus_per_member,
})
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct SlowReader {
pub published: usize,
pub healthy_received: usize,
pub reached_both: usize,
pub dropped_while_paced: u64,
pub absorbed: usize,
pub dropped: u64,
pub cause: &'static str,
}
pub async fn slow_reader() -> Result<SlowReader, Error> {
const PACED: usize = 512;
const PAYLOAD: usize = 4096;
let bound = bind("/feed").await?;
let publisher = bound.listener.publisher("/feed")?;
let subs = subscribers(&bound.url, 2).await?;
while publisher.filter_count() != subs.len() {
tokio::time::sleep(Duration::from_millis(1)).await;
}
let reading = Arc::clone(&subs[0].1);
let payload = vec![0x7au8; PAYLOAD];
let mut healthy_received = 0;
let mut reached_both = 0;
for _ in 0..PACED {
if publisher.publish("px.eur", payload.clone())? == 2 {
reached_both += 1;
}
match tokio::time::timeout(Duration::from_secs(5), reading.recv()).await {
Ok(Ok(transfer)) => {
transfer.collect(64 * 1024).await?;
healthy_received += 1;
}
_ => break,
}
}
let dropped_while_paced = publisher.dropped();
let mut absorbed = 0usize;
while publisher.dropped() == 0 && absorbed < 64 * 1024 {
publisher.publish("px.eur", payload.clone())?;
absorbed += 1;
if absorbed.is_multiple_of(256) {
tokio::task::yield_now().await;
}
}
let drops = publisher.dropped_on("px.eur");
let cause = match drops {
Some(ref drops) if drops.subscriber_budget > 0 => "the byte budget",
Some(ref drops) if drops.subscriber_queue > 0 => "the queue",
Some(_) => "no parked connection",
None => "nothing was dropped",
};
let dropped = publisher.dropped();
drop(subs);
bound.runtime.shutdown().await;
Ok(SlowReader {
published: PACED,
healthy_received,
reached_both,
dropped_while_paced,
absorbed,
dropped,
cause,
})
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct Ceiling {
pub budget: usize,
pub accepted: usize,
pub bytes: usize,
pub bound: &'static str,
}
pub async fn ceilings() -> Result<Vec<Ceiling>, Error> {
const PAYLOAD: usize = 1024;
let mut found = Vec::new();
for budget in [64 * 1024usize, 8 * 1024 * 1024] {
let mut config = RuntimeConfig::default();
config.limits = Limits {
subscriber_buffer_bytes: budget,
..config.limits
};
let bound = bind_with("/bulk", config.clone()).await?;
let publisher = bound.listener.publisher("/bulk")?;
let subs = subscribers_with(&bound.url, 1, config).await?;
while publisher.filter_count() != subs.len() {
tokio::time::sleep(Duration::from_millis(1)).await;
}
let payload = vec![0x7au8; PAYLOAD];
let mut accepted = 0usize;
while publisher.dropped() == 0 && accepted < 64 * 1024 {
publisher.publish("px.eur", payload.clone())?;
accepted += 1;
if accepted.is_multiple_of(256) {
tokio::task::yield_now().await;
}
}
let drops = publisher.dropped_on("px.eur");
let bound_name = match drops {
Some(ref drops) if drops.subscriber_budget > 0 => "the byte budget",
Some(ref drops) if drops.subscriber_queue > 0 => "the queue",
Some(_) => "no parked connection",
None => "nothing was dropped",
};
found.push(Ceiling {
budget,
accepted,
bytes: accepted * PAYLOAD,
bound: bound_name,
});
drop(subs);
bound.runtime.shutdown().await;
}
Ok(found)
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct Width {
pub subscribers: usize,
pub per_publish: Duration,
pub copies: usize,
pub dropped: u64,
}
pub async fn width(widths: &[usize], messages: usize) -> Result<Vec<Width>, Error> {
const PAYLOAD: usize = 1024;
let mut found = Vec::with_capacity(widths.len());
for &count in widths {
let bound = bind(&format!("/wide{count}")).await?;
let publisher = bound.listener.publisher(&format!("/wide{count}"))?;
let subs = subscribers(&bound.url, count).await?;
while publisher.filter_count() != count {
tokio::time::sleep(Duration::from_millis(1)).await;
}
let mut drains = Vec::with_capacity(count);
for (_, sub) in &subs {
let sub = Arc::clone(sub);
drains.push(tokio::spawn(async move {
let mut mine = 0;
for _ in 0..messages {
let Ok(Ok(transfer)) =
tokio::time::timeout(Duration::from_secs(10), sub.recv()).await
else {
break;
};
if transfer.collect(64 * 1024).await.is_err() {
break;
}
mine += 1;
}
mine
}));
}
let payload = vec![0x7au8; PAYLOAD];
let started = Instant::now();
for _ in 0..messages {
let reached = publisher.publish("px.eur", payload.clone())?;
debug_assert_eq!(reached, count);
}
let per_publish = started.elapsed() / messages as u32;
let mut copies = 0;
for drain in drains {
copies += drain.await.unwrap_or(0);
}
found.push(Width {
subscribers: count,
per_publish,
copies,
dropped: publisher.dropped(),
});
drop(subs);
bound.runtime.shutdown().await;
}
Ok(found)
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct Surveyed {
pub asked: usize,
pub answered: usize,
pub missing: usize,
}
pub async fn survey_with_a_silent_peer() -> Result<Surveyed, Error> {
const ASKED: usize = 3;
let mut servers = Vec::new();
for index in 0..ASKED {
servers.push(bind(&format!("/poll{index}")).await?);
}
let mut answering = Vec::new();
for (index, server) in servers.iter().enumerate() {
let respondent = server.listener.respondent(&format!("/poll{index}"))?;
let silent = index == ASKED - 1;
answering.push(tokio::spawn(async move {
let mut held = Vec::new();
while let Ok(mut question) = respondent.accept().await {
if silent {
held.push(question);
continue;
}
let Ok(body) = question.body().read_capped(1024).await else {
break;
};
let Ok(mut reply) = question.reply(weida::TransferMeta::default()).await else {
break;
};
if reply.write_all(&body).await.is_err() {
break;
}
if reply.finish().is_err() {
break;
}
}
}));
}
let client = Runtime::new(RuntimeConfig::default())?;
let surveyor = client.surveyor(Trust::by_address());
for server in &servers {
surveyor.connect(&server.url).await?;
}
let mut run = surveyor
.survey(b"who is there", Duration::from_millis(500))
.await?;
let mut answered = 0;
while let Some(answer) = run.next(1024).await {
if answer.is_ok() {
answered += 1;
}
}
for task in answering {
task.abort();
}
client.shutdown().await;
for server in servers {
server.runtime.shutdown().await;
}
Ok(Surveyed {
asked: ASKED,
answered,
missing: ASKED - answered,
})
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
println!("§2.1 selection is the pattern's");
let selected = selection().await?;
println!(
" 6 pushes over 3 pullers: {:?}",
selected.push_per_peer
);
println!(
" 2 publishes to 3 subscribers: {:?}",
selected.fanout_per_subscriber
);
println!(
" 3 BUS members, one send each: {:?} <- never itself",
selected.bus_per_member
);
println!("§2.2 one slow reader");
let slow = slow_reader().await?;
println!(
" paced: {} published, all {} enqueued for both, the reader got {}, \
{} dropped",
slow.published, slow.reached_both, slow.healthy_received, slow.dropped_while_paced
);
println!(
" flat out: {} more absorbed by a subscriber that has never read, then \
{} dropped, cause: {}",
slow.absorbed, slow.dropped, slow.cause
);
println!("§2.3 which ceiling binds");
for ceiling in ceilings().await? {
println!(
" budget {:>4} KiB: {:>5} messages ({:>4} KiB) accepted, then {}",
ceiling.budget >> 10,
ceiling.accepted,
ceiling.bytes >> 10,
ceiling.bound
);
}
println!("§2.4 what width costs the publisher");
let widths = width(&[1, 16, 64], 32).await?;
for measured in &widths {
println!(
" {:>3} subscribers: {:>10?} per publish, {:>5} copies, {} dropped",
measured.subscribers, measured.per_publish, measured.copies, measured.dropped
);
}
if let (Some(one), Some(many)) = (widths.first(), widths.last()) {
let marginal = many.per_publish.saturating_sub(one.per_publish)
/ (many.subscribers.saturating_sub(one.subscribers)).max(1) as u32;
println!(" marginal cost per subscriber: {marginal:?}");
}
if cfg!(debug_assertions) {
println!(
" ^ debug build: run with --release to compare with the measured \
figures (B-247)"
);
}
println!("§2.5 a silent peer");
let surveyed = survey_with_a_silent_peer().await?;
println!(
" asked {}, answered {}, silent {} within the deadline",
surveyed.asked, surveyed.answered, surveyed.missing
);
Ok(())
}