#![cfg(not(target_arch = "wasm32"))]
use beam::adapters::{MemoryStorage, OutgoingWebsocketManager, WsServer, WsServerConfig};
use beam::{Config, Node, Value};
use std::time::{Duration, Instant};
use tokio::time::sleep;
async fn start_relay(port: u16) -> Node {
let ws_config = WsServerConfig {
port,
cert_path: None,
key_path: None,
};
let node = Node::new_with_config(
Config::default(),
vec![Box::new(MemoryStorage::new())],
vec![Box::new(WsServer::new_with_config(
Config::default(),
ws_config,
))],
);
loop {
if tokio::net::TcpStream::connect(format!("127.0.0.1:{}", port))
.await
.is_ok()
{
break;
}
sleep(Duration::from_millis(50)).await;
}
sleep(Duration::from_millis(200)).await;
node
}
async fn connect_client(port: u16) -> Node {
let client = OutgoingWebsocketManager::new(
Config::default(),
vec![format!("ws://127.0.0.1:{}/ws", port)],
);
let node = Node::new_with_config(
Config::default(),
vec![Box::new(MemoryStorage::new())],
vec![Box::new(client)],
);
sleep(Duration::from_millis(300)).await;
node
}
async fn run_bench(senders: usize, messages: usize, port: u16) {
let mut relay = start_relay(port).await;
let mut subscriber_node = connect_client(port).await;
let mut sender_nodes = Vec::new();
for _ in 0..senders {
let node = connect_client(port).await;
sender_nodes.push(node);
}
sleep(Duration::from_millis(500)).await;
let before = relay.metrics().snapshot();
let start = Instant::now();
for (sender_idx, sender_node) in sender_nodes.iter_mut().enumerate() {
for i in 0..messages {
let key = format!("bench/{}_{}", sender_idx, i);
let _ = sender_node
.get(&key)
.put(Value::Text(format!("msg_{}", i)))
.await;
}
}
let send_elapsed = start.elapsed();
let stabilize_start = Instant::now();
loop {
let snap = relay.metrics().snapshot();
sleep(Duration::from_millis(500)).await;
let snap2 = relay.metrics().snapshot();
if snap.messages_relayed == snap2.messages_relayed {
break;
}
if stabilize_start.elapsed() > Duration::from_secs(30) {
break;
}
}
let total_elapsed = start.elapsed();
let after = relay.metrics().snapshot();
let relayed = after.messages_relayed - before.messages_relayed;
let ws_sent = after.ws_messages_sent - before.ws_messages_sent;
let ws_recv = after.ws_messages_received - before.ws_messages_received;
let parsed = after.messages_parsed - before.messages_parsed;
let dropped_dup = after.messages_dropped_dup - before.messages_dropped_dup;
let fanout = after.subscriber_fanout_total - before.subscriber_fanout_total;
let serialized = after.serialization_calls - before.serialization_calls;
let throughput = if total_elapsed.as_secs_f64() > 0.0 {
relayed as f64 / total_elapsed.as_secs_f64()
} else {
0.0
};
let send_rate = if send_elapsed.as_secs_f64() > 0.0 {
(senders * messages) as f64 / send_elapsed.as_secs_f64()
} else {
0.0
};
println!("\n============================================================");
println!(" RELAY THROUGHPUT BENCHMARK");
println!("============================================================");
println!(" Senders: {}", senders);
println!(" Messages/sender: {}", messages);
println!(" Total sent: {}", senders * messages);
println!(
" Send phase: {:.3}s ({:.0} puts/sec)",
send_elapsed.as_secs_f64(),
send_rate
);
println!(" Total elapsed: {:.3}s", total_elapsed.as_secs_f64());
println!();
println!(" --- Relay Hot-Path Counters ---");
println!(" ws_messages_received: {}", ws_recv);
println!(" messages_parsed: {}", parsed);
println!(" messages_dropped_dup: {}", dropped_dup);
println!(" messages_relayed: {}", relayed);
println!(" subscriber_fanout: {}", fanout);
println!(" serialization_calls: {}", serialized);
println!(" ws_messages_sent: {}", ws_sent);
println!(
" dropped_sends: {}",
after.dropped_sends - before.dropped_sends
);
println!();
println!(" Throughput: {:.0} msgs/sec (relayed)", throughput);
println!(
" Fanout ratio: {:.1}",
if relayed > 0 {
fanout as f64 / relayed as f64
} else {
0.0
}
);
println!(
" Dedup rate: {:.1}%",
if parsed > 0 {
100.0 * dropped_dup as f64 / parsed as f64
} else {
0.0
}
);
println!("============================================================\n");
for mut node in sender_nodes {
node.stop();
}
subscriber_node.stop();
relay.stop();
}
#[tokio::test]
#[ignore = "benchmark — run with --release --ignored --nocapture"]
async fn relay_throughput_1_sender_10k() {
run_bench(1, 10_000, 9970).await;
}
#[tokio::test]
#[ignore = "benchmark — run with --release --ignored --nocapture"]
async fn relay_throughput_1_sender_50k() {
run_bench(1, 50_000, 9972).await;
}
#[tokio::test]
#[ignore = "benchmark — run with --release --ignored --nocapture"]
async fn relay_throughput_10_senders_5k_each() {
run_bench(10, 5_000, 9974).await;
}