use std::net::SocketAddr;
use std::time::Duration;
use weida::{CursorLevel, Error, Identity, Runtime, RuntimeConfig, TransferMeta, Trust};
pub struct Bound {
pub url: String,
pub runtime: Runtime,
pub listener: weida::Listener,
_binding: weida::Binding,
}
pub async fn bind(path: &str) -> Result<Bound, Error> {
let identity = Identity::generate()?;
let fingerprint = identity.fingerprint()?;
let runtime = Runtime::new(RuntimeConfig::default())?;
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,
})
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct Answered {
pub reply: Vec<u8>,
pub peer_was_named: bool,
}
pub async fn hello() -> Result<Answered, Error> {
let bound = bind("/hello").await?;
let replier = bound.listener.replier("/hello")?;
let serving = tokio::spawn(async move {
let mut request = replier.accept().await?;
let body = request.body().read_capped(1024).await?;
let mut reply = request.reply(TransferMeta::default()).await?;
reply
.write_all(format!("hello, {}", String::from_utf8_lossy(&body)).as_bytes())
.await?;
reply.finish()?;
Ok::<(), Error>(())
});
let client = Runtime::new(RuntimeConfig::default())?;
let requester = client.requester(Trust::by_address());
requester.connect(&bound.url).await?;
let reply = requester.request(b"world").await?.collect(1024).await?;
serving.await.map_err(|e| Error::Runtime(e.to_string()))??;
client.shutdown().await;
Ok(Answered {
reply,
peer_was_named: bound.url.contains("sha256:"),
})
}
pub const STREAMED: usize = 64 * 1024 * 1024;
pub const CHUNK: usize = 64 * 1024;
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct Streamed {
pub bytes: u64,
pub checksum: u64,
pub collect_refused: bool,
pub sender_was_told: Option<String>,
}
pub async fn gigabyte() -> Result<Streamed, Error> {
use tokio::io::AsyncReadExt;
let bound = bind("/bulk").await?;
let puller = bound.listener.puller("/bulk")?;
let reading = tokio::spawn(async move {
let mut transfer = puller.recv().await?;
let mut buffer = vec![0u8; CHUNK];
let (mut bytes, mut checksum) = (0u64, 0xcbf2_9ce4_8422_2325u64);
loop {
let read = transfer.read(&mut buffer).await?;
if read == 0 {
break;
}
bytes += read as u64;
for byte in &buffer[..read] {
checksum ^= u64::from(*byte);
checksum = checksum.wrapping_mul(0x100_0000_01b3);
}
}
let second = puller.recv().await?;
let refused = matches!(second.collect(1024 * 1024).await, Err(Error::LimitExceeded));
Ok::<(u64, u64, bool), Error>((bytes, checksum, refused))
});
let client = Runtime::new(RuntimeConfig::default())?;
let pusher = client.pusher(Trust::by_address());
pusher.connect(&bound.url).await?;
let block = vec![0x5au8; CHUNK];
let mut transfer = pusher.open(TransferMeta::default()).await?;
for _ in 0..STREAMED / CHUNK {
transfer.write_all(&block).await?;
}
transfer.finish()?;
let mut capped = pusher.open(TransferMeta::default()).await?;
let mut sender_was_told = None;
for _ in 0..STREAMED / CHUNK {
if let Err(e) = capped.write_all(&block).await {
sender_was_told = Some(e.to_string());
break;
}
}
if sender_was_told.is_none() {
if let Err(e) = capped.finish()?.delivered().await {
sender_was_told = Some(e.to_string());
}
}
let (bytes, checksum, collect_refused) =
reading.await.map_err(|e| Error::Runtime(e.to_string()))??;
client.shutdown().await;
Ok(Streamed {
bytes,
checksum,
collect_refused,
sender_was_told,
})
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum Outcome {
Delivered,
Refused(String),
Indeterminate,
}
pub fn classify(result: Result<(), Error>) -> Outcome {
match result {
Ok(()) => Outcome::Delivered,
Err(Error::Indeterminate) => Outcome::Indeterminate,
Err(e) if e.is_definite_failure() => Outcome::Refused(e.to_string()),
Err(e) => Outcome::Refused(e.to_string()),
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct Outcomes {
pub push_served: Outcome,
pub push_unserved: Outcome,
pub request_unserved: Outcome,
}
pub async fn outcome() -> Result<Outcomes, Error> {
let bound = bind("/served").await?;
let puller = bound.listener.puller("/served")?;
let drain = tokio::spawn(async move {
while let Ok(transfer) = puller.recv().await {
let _ = transfer.collect(1024).await;
}
});
let client = Runtime::new(RuntimeConfig::default())?;
let served = client.pusher(Trust::by_address());
served.connect(&bound.url).await?;
let mut transfer = served.open(TransferMeta::default()).await?;
transfer.write_all(b"for a path somebody serves").await?;
let push_served = classify(transfer.finish()?.delivered().await);
let missing_url = bound.url.replace("/served", "/nobody-serves-this");
let missing = client.pusher(Trust::by_address());
missing.connect(&missing_url).await?;
let mut transfer = missing.open(TransferMeta::default()).await?;
let write = transfer.write_all(b"for a path nobody serves").await;
let push_unserved = match write {
Err(e) => classify(Err(e)),
Ok(()) => classify(transfer.finish()?.delivered().await),
};
let asking = client.requester(Trust::by_address());
asking.connect(&missing_url).await?;
let request_unserved = match asking.request(b"is anybody there").await {
Ok(_) => Outcome::Delivered,
Err(e) => classify(Err(e)),
};
drain.abort();
client.shutdown().await;
Ok(Outcomes {
push_served,
push_unserved,
request_unserved,
})
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub struct HowFar {
pub delivered: bool,
pub reported: Option<u64>,
pub sent: u64,
}
pub fn stage() -> CursorLevel {
CursorLevel::application(CursorLevel::APPLICATION_FLOOR).expect("16 is at the floor")
}
pub async fn how_far() -> Result<HowFar, Error> {
let bound = bind("/work").await?;
let puller = bound.listener.puller("/work")?;
let working = tokio::spawn(async move {
let transfer = puller.recv().await?;
let mut reporter = transfer.reporter().expect("the sender ordered a report");
let body = transfer.collect(64 * 1024).await?;
reporter.report(stage(), body.len() as u64).await?;
reporter.finish().await?;
Ok::<usize, Error>(body.len())
});
let client = Runtime::new(RuntimeConfig::default())?;
let pusher = client.pusher(Trust::by_address());
pusher.connect(&bound.url).await?;
let payload = b"eight billion of these".to_vec();
let meta = TransferMeta::default().with_report([stage()]);
let mut transfer = pusher.open(meta).await?;
let mut cursors = transfer.cursors().expect("a report was ordered");
transfer.write_all(&payload).await?;
let delivered = transfer.finish()?.delivered().await.is_ok();
let reported = cursors.changed().await.and_then(|set| set.offset(stage()));
working.await.map_err(|e| Error::Runtime(e.to_string()))??;
client.shutdown().await;
Ok(HowFar {
delivered,
reported,
sent: payload.len() as u64,
})
}
#[derive(Clone, Copy, Debug)]
pub struct ReceiptCost {
pub without: Duration,
pub with: Duration,
pub rounds: usize,
}
pub async fn receipt_cost(rounds: usize) -> Result<ReceiptCost, Error> {
let bound = bind("/cost").await?;
let puller = bound.listener.puller("/cost")?;
let draining = tokio::spawn(async move {
while let Ok(transfer) = puller.recv().await {
let _ = transfer.collect(64 * 1024).await;
}
});
let client = Runtime::new(RuntimeConfig::default())?;
let pusher = client.pusher(Trust::by_address());
pusher.connect(&bound.url).await?;
let payload = vec![0x5au8; 1024];
let started = std::time::Instant::now();
for _ in 0..rounds {
pusher.send(&payload).await?;
}
let without = started.elapsed() / rounds as u32;
let started = std::time::Instant::now();
for _ in 0..rounds {
let mut transfer = pusher.open(TransferMeta::default()).await?;
transfer.write_all(&payload).await?;
transfer.finish()?.delivered().await?;
}
let with = started.elapsed() / rounds as u32;
draining.abort();
client.shutdown().await;
Ok(ReceiptCost {
without,
with,
rounds,
})
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
println!("§1.1 hello");
let answered = hello().await?;
println!(
" reply {:?}, the address named the peer: {}",
String::from_utf8_lossy(&answered.reply),
answered.peer_was_named
);
println!("§1.2 the same four calls carry {} MiB", STREAMED >> 20);
let streamed = gigabyte().await?;
println!(
" {} bytes folded through a {} KiB buffer, checksum {:#x}",
streamed.bytes,
CHUNK >> 10,
streamed.checksum
);
println!(
" the same payload under a 1 MiB cap: reader refused it ({}), sender was told {:?}",
streamed.collect_refused, streamed.sender_was_told
);
println!("§1.3 the outcome is a value");
let outcomes = outcome().await?;
println!(" push, served path: {:?}", outcomes.push_served);
println!(
" push, unserved path: {:?} <- a race by design, see 0005",
outcomes.push_unserved
);
println!(
" request, same path: {:?} <- definite, every time",
outcomes.request_unserved
);
println!("§1.4 how far did it get");
let far = how_far().await?;
println!(
" transport receipt: {}, application reported {:?} of {} bytes",
far.delivered, far.reported, far.sent
);
println!("§1.5 what the receipt costs");
let cost = receipt_cost(8).await?;
println!(
" 1 KiB, {} rounds each: {:?} without the receipt, {:?} with it",
cost.rounds, cost.without, cost.with
);
Ok(())
}