#![cfg(test)]
use beam::adapters::*;
use beam::{Config, Node, Value};
use std::process::{Child, Command, Stdio};
use std::time::Duration;
use tokio::net::TcpStream;
use tokio::time::{sleep, timeout};
struct GunRelay {
child: Child,
ws_port: u16,
api_port: u16,
}
impl GunRelay {
fn spawn(ws_port: u16) -> Self {
let api_port = ws_port + 1;
let child = Command::new("node")
.arg("gun_relay.js")
.arg(ws_port.to_string())
.arg(api_port.to_string())
.current_dir("tests/wire-live")
.stdout(Stdio::null())
.stderr(Stdio::null())
.spawn()
.expect("failed to spawn gun_relay.js — is Node.js installed?");
Self {
child,
ws_port,
api_port,
}
}
async fn wait_for_ready(&self, timeout_ms: u64) {
let start = std::time::Instant::now();
let limit = Duration::from_millis(timeout_ms);
while start.elapsed() < limit {
if TcpStream::connect(format!("127.0.0.1:{}", self.api_port))
.await
.is_ok()
{
return;
}
sleep(Duration::from_millis(100)).await;
}
panic!("Gun relay API did not become ready within {}ms", timeout_ms);
}
async fn api_get(&self, path: &str) -> String {
let url = format!("http://127.0.0.1:{}{}", self.api_port, path);
reqwest_text(&url).await
}
async fn api_post(&self, path: &str, body: &str) -> String {
let url = format!("http://127.0.0.1:{}{}", self.api_port, path);
reqwest_post(&url, body).await
}
}
impl Drop for GunRelay {
fn drop(&mut self) {
let _ = self.child.kill();
let _ = self.child.wait();
}
}
async fn reqwest_text(url: &str) -> String {
let url = url.strip_prefix("http://").unwrap_or(url);
let (host_port, path) = url.split_once('/').unwrap_or((url, "/"));
let (host, port_str) = host_port.rsplit_once(':').unwrap_or((host_port, "80"));
let port: u16 = port_str.parse().unwrap_or(80);
let mut stream = TcpStream::connect(format!("{}:{}", host, port))
.await
.expect("failed to connect to API");
use tokio::io::{AsyncReadExt, AsyncWriteExt};
let request = format!(
"GET {} HTTP/1.1\r\nHost: {}\r\nConnection: close\r\n\r\n",
path, host
);
stream.write_all(request.as_bytes()).await.unwrap();
let mut buf = Vec::new();
stream.read_to_end(&mut buf).await.unwrap();
let response = String::from_utf8_lossy(&buf).to_string();
if let Some(idx) = response.find("\r\n\r\n") {
response[idx + 4..].to_string()
} else {
response
}
}
async fn reqwest_post(url: &str, body: &str) -> String {
let url = url.strip_prefix("http://").unwrap_or(url);
let (host_port, path) = url.split_once('/').unwrap_or((url, "/"));
let (host, port_str) = host_port.rsplit_once(':').unwrap_or((host_port, "80"));
let port: u16 = port_str.parse().unwrap_or(80);
let mut stream = TcpStream::connect(format!("{}:{}", host, port))
.await
.expect("failed to connect to API");
use tokio::io::{AsyncReadExt, AsyncWriteExt};
let request = format!(
"POST {} HTTP/1.1\r\nHost: {}\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}",
path,
host,
body.len(),
body
);
stream.write_all(request.as_bytes()).await.unwrap();
let mut buf = Vec::new();
stream.read_to_end(&mut buf).await.unwrap();
let response = String::from_utf8_lossy(&buf).to_string();
if let Some(idx) = response.find("\r\n\r\n") {
response[idx + 4..].to_string()
} else {
response
}
}
#[tokio::test]
#[ignore = "requires Node.js + Gun.js installed in tests/wire-live/"]
async fn beam_put_gun_receives() {
let relay = GunRelay::spawn(9871);
relay.wait_for_ready(10000).await;
let config = Config::default();
let ws_client = OutgoingWebsocketManager::new(
config.clone(),
vec![format!("ws://127.0.0.1:{}/gun", relay.ws_port)],
);
let mut beam = Node::new_with_config(
config,
vec![Box::new(MemoryStorage::new())],
vec![Box::new(ws_client.clone())],
);
wait_for_connected(&ws_client, 1, 10000).await;
beam.get("beamtest/put1")
.get("name")
.put("Alice".into())
.await
.unwrap();
sleep(Duration::from_secs(2)).await;
let response = relay.api_get("/get?soul=beamtest/put1&key=name").await;
assert!(
response.contains("Alice"),
"Gun.js did not receive BEAM's put. Response: {}",
response
);
beam.stop();
}
#[tokio::test]
#[ignore = "requires Node.js + Gun.js installed in tests/wire-live/"]
async fn gun_put_beam_receives() {
let relay = GunRelay::spawn(9873);
relay.wait_for_ready(10000).await;
let config = Config::default();
let ws_client = OutgoingWebsocketManager::new(
config.clone(),
vec![format!("ws://127.0.0.1:{}/gun", relay.ws_port)],
);
let mut beam = Node::new_with_config(
config,
vec![Box::new(MemoryStorage::new())],
vec![Box::new(ws_client.clone())],
);
wait_for_connected(&ws_client, 1, 10000).await;
let mut sub = beam.get("beamtest/put2").get("name").on();
sleep(Duration::from_secs(1)).await;
relay
.api_post(
"/put",
r#"{"soul":"beamtest/put2","key":"name","value":"Bob"}"#,
)
.await;
let result = timeout(Duration::from_secs(15), sub.recv())
.await
.expect("timeout waiting for Gun.js put to propagate to BEAM")
.expect("subscription channel closed");
match result {
Value::Text(s) => assert_eq!(s, "Bob"),
other => panic!("expected Value::Text, got {:?}", other),
}
beam.stop();
}
#[tokio::test]
#[ignore = "requires Node.js + Gun.js installed in tests/wire-live/"]
async fn bidirectional_convergence() {
let relay = GunRelay::spawn(9875);
relay.wait_for_ready(10000).await;
let config = Config::default();
let ws_client = OutgoingWebsocketManager::new(
config.clone(),
vec![format!("ws://127.0.0.1:{}/gun", relay.ws_port)],
);
let mut beam = Node::new_with_config(
config,
vec![Box::new(MemoryStorage::new())],
vec![Box::new(ws_client.clone())],
);
wait_for_connected(&ws_client, 1, 10000).await;
beam.get("beamtest/conv")
.get("from_beam")
.put("beam_data".into())
.await
.unwrap();
relay
.api_post(
"/put",
r#"{"soul":"beamtest/conv","key":"from_gun","value":"gun_data"}"#,
)
.await;
sleep(Duration::from_secs(3)).await;
let beam_val = relay.api_get("/get?soul=beamtest/conv&key=from_beam").await;
assert!(
beam_val.contains("beam_data"),
"Gun.js missing BEAM's data. Response: {}",
beam_val
);
let mut gun_sub = beam.get("beamtest/conv").get("from_gun").on();
let result = timeout(Duration::from_secs(10), gun_sub.recv())
.await
.expect("timeout waiting for Gun.js data to reach BEAM")
.expect("subscription channel closed");
match result {
Value::Text(s) => assert_eq!(s, "gun_data"),
other => panic!("expected Value::Text, got {:?}", other),
}
beam.stop();
}
#[tokio::test]
#[ignore = "requires Node.js + Gun.js installed in tests/wire-live/"]
async fn reconnection_sync() {
let relay = GunRelay::spawn(9877);
relay.wait_for_ready(10000).await;
let config = Config::default();
let ws_client = OutgoingWebsocketManager::new(
config.clone(),
vec![format!("ws://127.0.0.1:{}/gun", relay.ws_port)],
);
let mut beam = Node::new_with_config(
config,
vec![Box::new(MemoryStorage::new())],
vec![Box::new(ws_client.clone())],
);
wait_for_connected(&ws_client, 1, 10000).await;
beam.get("beamtest/recon")
.get("phase1")
.put("first".into())
.await
.unwrap();
sleep(Duration::from_secs(2)).await;
let v1 = relay.api_get("/get?soul=beamtest/recon&key=phase1").await;
assert!(v1.contains("first"), "initial sync failed: {}", v1);
beam.stop();
sleep(Duration::from_secs(2)).await;
relay
.api_post(
"/put",
r#"{"soul":"beamtest/recon","key":"phase2","value":"second"}"#,
)
.await;
sleep(Duration::from_secs(1)).await;
let ws_client2 = OutgoingWebsocketManager::new(
Config::default(),
vec![format!("ws://127.0.0.1:{}/gun", relay.ws_port)],
);
let mut beam2 = Node::new_with_config(
Config::default(),
vec![Box::new(MemoryStorage::new())],
vec![Box::new(ws_client2.clone())],
);
wait_for_connected(&ws_client2, 1, 10000).await;
let mut sub = beam2.get("beamtest/recon").get("phase2").on();
let result = timeout(Duration::from_secs(15), sub.recv())
.await
.expect("timeout waiting for reconnection sync")
.expect("channel closed");
match result {
Value::Text(s) => assert_eq!(s, "second"),
other => panic!("expected Value::Text, got {:?}", other),
}
beam2.stop();
}
async fn wait_for_connected(client: &OutgoingWebsocketManager, expected: usize, timeout_ms: u64) {
let start = std::time::Instant::now();
let limit = Duration::from_millis(timeout_ms);
while start.elapsed() < limit {
if client.connected_count().await >= expected {
return;
}
sleep(Duration::from_millis(50)).await;
}
panic!(
"OutgoingWebsocketManager not connected within {}ms (expected {} connections)",
timeout_ms, expected
);
}