#![cfg(test)]
use wasm_bindgen::{JsCast, prelude::*};
use wasm_bindgen_test::*;
wasm_bindgen_test_configure!(run_in_node_experimental);
async fn eval_promise(js: &str) -> JsValue {
let promise = js_sys::eval(js)
.expect("eval failed")
.dyn_into::<js_sys::Promise>()
.expect("eval did not return a Promise");
wasm_bindgen_futures::JsFuture::from(promise)
.await
.expect("promise rejected")
}
async fn sleep(ms: u64) {
eval_promise(&format!("new Promise(r => setTimeout(r, {}))", ms)).await;
}
struct Relay {
port: u16,
proc: JsValue,
}
impl Relay {
async fn start(port: u16) -> Self {
let proc = eval_promise(&format!(
r#"
import('node:child_process').then(cp => {{
const child = cp.spawn('target/debug/beam', [
'start', '--port', '{port}',
'--memory-storage', 'true', '--redb-storage', 'false',
'--allow-public-space', 'true'
], {{
cwd: '/home/guan/src/beam',
stdio: ['ignore', 'pipe', 'pipe']
}});
child.stdout.on('data', d => process.stderr.write('[relay] ' + d));
child.stderr.on('data', d => process.stderr.write('[relay] ' + d));
return child;
}})
"#
))
.await;
eval_promise(&format!(
r#"
new Promise((resolve, reject) => {{
import('node:net').then(net => {{
const deadline = Date.now() + 10000;
const tryConnect = () => {{
const sock = net.connect({port}, '127.0.0.1');
sock.on('connect', () => {{ sock.destroy(); resolve(); }});
sock.on('error', () => {{
if (Date.now() > deadline) reject(new Error('relay did not start'));
else setTimeout(tryConnect, 50);
}});
}};
tryConnect();
}});
}})
"#
))
.await;
Self { port, proc }
}
}
impl Drop for Relay {
fn drop(&mut self) {
let f = js_sys::Function::new_with_args("p", "if (p && p.kill) p.kill('SIGTERM');");
let _ = f.call1(&JsValue::UNDEFINED, &self.proc);
}
}
#[wasm_bindgen_test]
fn smoke_test() {
assert_eq!(2 + 2, 4);
}
#[wasm_bindgen_test(async)]
async fn local_put_get_roundtrip() {
use crate::wasm::Beam;
let mut beam = Beam::new();
beam.put("chat.001", "hello world");
sleep(100).await;
let result = wasm_bindgen_futures::JsFuture::from(beam.get("chat.001"))
.await
.expect("get should resolve");
assert_eq!(result.as_string(), Some("hello world".to_string()));
beam.stop();
}
#[wasm_bindgen_test(async)]
async fn relay_connect() {
use crate::wasm::Beam;
let _relay = Relay::start(4960).await;
let mut beam = Beam::new();
beam.connect("ws://127.0.0.1:4960");
sleep(500).await;
beam.stop();
}
#[wasm_bindgen_test(async)]
async fn relay_put_echo() {
use crate::wasm::Beam;
let _relay = Relay::start(4961).await;
let mut beam = Beam::new();
beam.connect("ws://127.0.0.1:4961");
sleep(500).await;
beam.put("chat.relay_test", "relay payload");
sleep(300).await;
let result = wasm_bindgen_futures::JsFuture::from(beam.get("chat.relay_test"))
.await
.expect("get should resolve");
assert_eq!(result.as_string(), Some("relay payload".to_string()));
beam.stop();
}
#[wasm_bindgen_test(async)]
#[ignore = "requires browser WebSocket API — web_sys::WebSocket events don't fire in Node.js test runner"]
async fn two_clients_cross_talk() {
use crate::wasm::Beam;
let _relay = Relay::start(4970).await;
let mut client1 = Beam::new();
client1.connect("ws://127.0.0.1:4970");
let mut client2 = Beam::new();
client2.connect("ws://127.0.0.1:4970");
sleep(1000).await;
js_sys::eval(
r#"
globalThis.__received = [];
globalThis.__on_msg = function(val) {
globalThis.__received.push(val);
};
"#,
)
.unwrap();
let callback = js_sys::eval("globalThis.__on_msg")
.unwrap()
.dyn_into::<js_sys::Function>()
.unwrap();
client2.on("chat", callback);
sleep(200).await;
client1.put("chat.42", "cross-talk!");
sleep(1000).await;
client1.stop();
client2.stop();
let received = js_sys::eval("JSON.stringify(globalThis.__received)")
.unwrap()
.as_string()
.unwrap_or_default();
assert!(
received.contains("cross-talk"),
"client2 should have received 'cross-talk!' but got: {}",
received
);
}
#[wasm_bindgen_test(async)]
#[ignore = "requires browser WebSocket API — web_sys::WebSocket events don't fire in Node.js test runner"]
async fn bidirectional_cross_talk() {
use crate::wasm::Beam;
let _relay = Relay::start(4980).await;
let mut client1 = Beam::new();
client1.connect("ws://127.0.0.1:4980");
let mut client2 = Beam::new();
client2.connect("ws://127.0.0.1:4980");
sleep(1000).await;
js_sys::eval(
r#"
globalThis.__c1_received = [];
globalThis.__c2_received = [];
globalThis.__c1_cb = v => globalThis.__c1_received.push(v);
globalThis.__c2_cb = v => globalThis.__c2_received.push(v);
"#,
)
.unwrap();
let c1_cb = js_sys::eval("globalThis.__c1_cb")
.unwrap()
.dyn_into::<js_sys::Function>()
.unwrap();
let c2_cb = js_sys::eval("globalThis.__c2_cb")
.unwrap()
.dyn_into::<js_sys::Function>()
.unwrap();
client1.on("chat", c1_cb);
client2.on("chat", c2_cb);
sleep(200).await;
client1.put("chat.001", "from_client_1");
sleep(500).await;
client2.put("chat.002", "from_client_2");
sleep(500).await;
client1.stop();
client2.stop();
let c1 = js_sys::eval("JSON.stringify(globalThis.__c1_received)")
.unwrap()
.as_string()
.unwrap_or_default();
let c2 = js_sys::eval("JSON.stringify(globalThis.__c2_received)")
.unwrap()
.as_string()
.unwrap_or_default();
assert!(
c2.contains("from_client_1"),
"client2 should have received 'from_client_1' but got: {}",
c2
);
assert!(
c1.contains("from_client_2"),
"client1 should have received 'from_client_2' but got: {}",
c1
);
}
#[wasm_bindgen_test]
fn wasm_bench_parse_throughput() {
let json =
r##"{"#":"bench/test","put":{"bench/test":{"msg":"hello world","time":1234567890.123}}}"##;
let iterations = 10_000;
let start = web_time::Instant::now();
for _ in 0..iterations {
let _: serde_json::Value = serde_json::from_str(json).unwrap();
}
let elapsed = start.elapsed();
let per_op_ns = elapsed.as_nanos() as f64 / iterations as f64;
console_log!(
"WASM parse small JSON: {:.0} ns/op ({:.0} ops/sec)",
per_op_ns,
1_000_000_000.0 / per_op_ns
);
}
#[wasm_bindgen_test]
fn wasm_bench_serialize_throughput() {
let obj = serde_json::json!({
"#": "bench/test",
"put": {
"bench/test": {
"msg": "hello world",
"time": 1234567890.123
}
}
});
let iterations = 10_000;
let start = web_time::Instant::now();
for _ in 0..iterations {
let _ = obj.to_string();
}
let elapsed = start.elapsed();
let per_op_ns = elapsed.as_nanos() as f64 / iterations as f64;
console_log!(
"WASM serialize small JSON: {:.0} ns/op ({:.0} ops/sec)",
per_op_ns,
1_000_000_000.0 / per_op_ns
);
}
#[wasm_bindgen_test]
fn wasm_bench_get_parse_throughput() {
let json = r##"{"#":"bench/test","get":{"#":"bench/test"}}"##;
let iterations = 10_000;
let start = web_time::Instant::now();
for _ in 0..iterations {
let _: serde_json::Value = serde_json::from_str(json).unwrap();
}
let elapsed = start.elapsed();
let per_op_ns = elapsed.as_nanos() as f64 / iterations as f64;
console_log!(
"WASM parse Get: {:.0} ns/op ({:.0} ops/sec)",
per_op_ns,
1_000_000_000.0 / per_op_ns
);
}
async fn fetch_metrics(ws_port: u16) -> JsValue {
let http_port = ws_port + 1;
eval_promise(&format!(
r#"
fetch("http://127.0.0.1:{http_port}/metrics")
.then(r => r.json())
"#
))
.await
}
fn get_u64(obj: &JsValue, key: &str) -> u64 {
let val = js_sys::Reflect::get(obj, &js_sys::JsString::from(key)).expect("field missing");
val.as_f64().unwrap_or(0.0) as u64
}
async fn run_wasm_relay_throughput(port: u16, count: usize) {
use crate::wasm::Beam;
let _relay = Relay::start(port).await;
let mut beam = Beam::new();
beam.connect(&format!("ws://127.0.0.1:{port}"));
sleep(500).await;
let before = fetch_metrics(port).await;
let start = web_time::Instant::now();
for i in 0..count {
let key = format!("bench/{}", i);
beam.put(&key, &format!("msg_{}", i));
}
let send_elapsed = start.elapsed();
let stabilize_deadline = web_time::Instant::now() + web_time::Duration::from_secs(30);
let mut last_relayed = get_u64(&before, "messages_relayed");
loop {
sleep(500).await;
let snap = fetch_metrics(port).await;
let now_relayed = get_u64(&snap, "messages_relayed");
if now_relayed == last_relayed {
break;
}
last_relayed = now_relayed;
if web_time::Instant::now() > stabilize_deadline {
break;
}
}
let total_elapsed = start.elapsed();
let after = fetch_metrics(port).await;
let relayed = get_u64(&after, "messages_relayed") - get_u64(&before, "messages_relayed");
let ws_sent = get_u64(&after, "ws_messages_sent") - get_u64(&before, "ws_messages_sent");
let ws_recv =
get_u64(&after, "ws_messages_received") - get_u64(&before, "ws_messages_received");
let parsed = get_u64(&after, "messages_parsed") - get_u64(&before, "messages_parsed");
let dropped_dup =
get_u64(&after, "messages_dropped_dup") - get_u64(&before, "messages_dropped_dup");
let fanout =
get_u64(&after, "subscriber_fanout_total") - get_u64(&before, "subscriber_fanout_total");
let serialized =
get_u64(&after, "serialization_calls") - get_u64(&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 {
count as f64 / send_elapsed.as_secs_f64()
} else {
0.0
};
console_log!(
"\n============================================================\n \
WASM RELAY THROUGHPUT BENCHMARK\n \
Messages: {} (fire-and-forget puts)\n \
Send phase: {:.3}s ({:.0} puts/sec)\n \
Total elapsed: {:.3}s\n \
--- Relay Hot-Path Counters ---\n \
ws_messages_received: {}\n \
messages_parsed: {}\n \
messages_dropped_dup: {}\n \
messages_relayed: {}\n \
subscriber_fanout: {}\n \
serialization_calls: {}\n \
ws_messages_sent: {}\n \
Throughput: {:.0} msgs/sec (relayed)\n \
Fanout ratio: {:.1}\n \
Dedup rate: {:.1}%\n\
============================================================\n",
count,
send_elapsed.as_secs_f64(),
send_rate,
total_elapsed.as_secs_f64(),
ws_recv,
parsed,
dropped_dup,
relayed,
fanout,
serialized,
ws_sent,
throughput,
if relayed > 0 {
fanout as f64 / relayed as f64
} else {
0.0
},
if parsed > 0 {
100.0 * dropped_dup as f64 / parsed as f64
} else {
0.0
},
);
assert!(relayed > 0, "relay should have processed messages");
beam.stop();
}
#[wasm_bindgen_test(async)]
#[ignore = "requires browser WebSocket API — web_sys::WebSocket events don't fire in Node.js test runner"]
async fn wasm_relay_throughput_1k() {
run_wasm_relay_throughput(4971, 1_000).await;
}
#[wasm_bindgen_test(async)]
#[ignore = "requires browser WebSocket API — web_sys::WebSocket events don't fire in Node.js test runner"]
async fn wasm_relay_throughput_5k() {
run_wasm_relay_throughput(4972, 5_000).await;
}