use std::{
str::FromStr,
sync::Arc,
time::{Duration, Instant},
};
use anyhow::Result;
use clap::Parser;
use dittolive_ditto::{fs::TempRoot, prelude::*, preview::peer_pubkey::PeerPubkey};
#[derive(Parser, Debug)]
struct Args {
#[arg(short, long, default_value = "100")]
number: usize,
#[arg(short, long, default_value = "64")]
size: usize,
#[arg(short, long, default_value = "1.0")]
warmup: f64,
}
#[tokio::main]
async fn main() -> Result<()> {
let args = Args::parse();
let root = TempRoot::new();
let ditto = Ditto::open_sync(
DittoConfig::new(
"bus_example_app",
DittoConfigConnect::SmallPeersOnly { private_key: None },
)
.with_persistence_directory(root.root_path()),
)?;
ditto.set_license_from_env("DITTO_LICENSE")?;
ditto.update_transport_config(|tc| {
tc.enable_all_peer_to_peer();
});
ditto.sync().start()?;
let ditto = Arc::new(ditto);
let remote_peer = loop {
let peers = ditto.presence().graph().remote_peers;
if let Some(peer) = peers.iter().next() {
break PeerPubkey::from_str(&peer.peer_key).unwrap();
}
tokio::time::sleep(Duration::from_secs(1)).await;
};
let mut stream = ditto
.datastreams()
.connect(remote_peer.clone(), "example")
.on_receive_factory(tokio::sync::mpsc::unbounded_channel)
.finish_async()
.await
.unwrap();
let mut samples = Vec::with_capacity(args.number);
let static_payload = Box::leak(vec![0u8; args.size].into_boxed_slice());
let start = Instant::now();
while start.elapsed() < Duration::from_secs_f64(args.warmup) {
let _ = stream.message(&*static_payload).send();
stream.recv().await;
}
for _ in 0..args.number {
let now = Instant::now();
let _ = stream.message(&*static_payload).send();
stream.recv().await;
samples.push(now.elapsed());
}
for (i, s) in samples.iter().enumerate() {
println!("{} bytes: seq={} time={:#?}", args.size, i, s);
}
Ok(())
}