use std::io::{BufRead, BufReader, Read, Write};
use std::net::{TcpListener, TcpStream};
use std::time::{Duration, Instant};
use turbo_debug_console::proto::StreamKind;
use turbo_debug_console::registry::{Server, ServerEvent};
fn hello(port: u16, line: &str) -> String {
let mut s = TcpStream::connect(("127.0.0.1", port)).unwrap();
s.write_all(format!("{line}\n").as_bytes()).unwrap();
let mut r = BufReader::new(s);
let mut reply = String::new();
r.read_line(&mut reply).unwrap();
reply.trim_end().to_string()
}
fn wait_for<T>(server: &Server, mut f: impl FnMut(&ServerEvent) -> Option<T>) -> T {
let deadline = Instant::now() + Duration::from_secs(1);
while Instant::now() < deadline {
while let Ok(ev) = server.events().try_recv() {
if let Some(v) = f(&ev) {
return v;
}
}
std::thread::sleep(Duration::from_millis(10));
}
panic!("timed out waiting for event");
}
#[test]
fn hello_allocates_a_data_port_and_opens_a_session() {
let server = Server::bind(0).unwrap();
let reply = hello(server.control_port(), "HELLO 1 tokens alpha");
let port: u16 = reply.strip_prefix("PORT ").unwrap().parse().unwrap();
assert_ne!(port, server.control_port());
let (name, ev_port) = wait_for(&server, |ev| match ev {
ServerEvent::Opened { name, port, .. } => Some((name.clone(), *port)),
_ => None,
});
assert_eq!(name, "alpha");
assert_eq!(ev_port, port);
let mut data = TcpStream::connect(("127.0.0.1", port)).unwrap();
data.write_all(b"tokens").unwrap();
let got = wait_for(&server, |ev| match ev {
ServerEvent::Bytes { data, .. } => Some(data.clone()),
_ => None,
});
assert_eq!(got, b"tokens");
}
#[test]
fn the_same_name_returns_the_same_port_and_reconnects() {
let server = Server::bind(0).unwrap();
let first = hello(server.control_port(), "HELLO 1 tokens beta");
let port: u16 = first.strip_prefix("PORT ").unwrap().parse().unwrap();
let mut data = TcpStream::connect(("127.0.0.1", port)).unwrap();
data.write_all(b"a").unwrap();
drop(data);
wait_for(&server, |ev| {
matches!(ev, ServerEvent::Disconnected { .. }).then_some(())
});
let second = hello(server.control_port(), "HELLO 1 tokens beta");
assert_eq!(first, second, "a known name must keep its port");
let mut data = TcpStream::connect(("127.0.0.1", port)).unwrap();
data.write_all(b"b").unwrap();
let reattached = wait_for(&server, |ev| match ev {
ServerEvent::Attached { reattached, .. } => Some(*reattached),
_ => None,
});
assert!(reattached, "a repeat attach must report reattached: true");
}
#[test]
fn the_ordinary_first_attach_reports_attached_not_reattached() {
let server = Server::bind(0).unwrap();
let reply = hello(server.control_port(), "HELLO 1 tokens first-timer");
let port: u16 = reply.strip_prefix("PORT ").unwrap().parse().unwrap();
let mut data = TcpStream::connect(("127.0.0.1", port)).unwrap();
data.write_all(b"hi").unwrap();
let reattached = wait_for(&server, |ev| match ev {
ServerEvent::Attached { reattached, .. } => Some(*reattached),
_ => None,
});
assert!(
!reattached,
"a brand-new session's first attach must not report reattached: true"
);
}
#[test]
fn a_dropped_data_socket_leaves_the_session_listening() {
let server = Server::bind(0).unwrap();
let reply = hello(server.control_port(), "HELLO 1 tokens gamma");
let port: u16 = reply.strip_prefix("PORT ").unwrap().parse().unwrap();
let data = TcpStream::connect(("127.0.0.1", port)).unwrap();
drop(data);
wait_for(&server, |ev| {
matches!(ev, ServerEvent::Disconnected { .. }).then_some(())
});
let mut again = TcpStream::connect(("127.0.0.1", port)).unwrap();
again.write_all(b"back").unwrap();
let got = wait_for(&server, |ev| match ev {
ServerEvent::Bytes { data, .. } => Some(data.clone()),
_ => None,
});
assert_eq!(got, b"back");
}
#[test]
fn a_second_live_writer_is_refused() {
let server = Server::bind(0).unwrap();
let reply = hello(server.control_port(), "HELLO 1 tokens delta");
let port: u16 = reply.strip_prefix("PORT ").unwrap().parse().unwrap();
let mut first = TcpStream::connect(("127.0.0.1", port)).unwrap();
first.write_all(b"x").unwrap();
wait_for(&server, |ev| {
matches!(ev, ServerEvent::Bytes { .. }).then_some(())
});
let mut second = TcpStream::connect(("127.0.0.1", port)).unwrap();
second
.set_read_timeout(Some(Duration::from_secs(1)))
.unwrap();
let mut buf = Vec::new();
second.read_to_end(&mut buf).unwrap();
assert!(
String::from_utf8_lossy(&buf).starts_with("ERR"),
"a duplicate writer must be closed with a banner, got {buf:?}"
);
}
#[test]
fn bad_names_are_refused_and_the_server_stays_up() {
let server = Server::bind(0).unwrap();
assert_eq!(
hello(server.control_port(), "HELLO 1 tokens "),
"ERR bad name"
);
assert_eq!(
hello(
server.control_port(),
&format!("HELLO 1 tokens {}", "x".repeat(65))
),
"ERR bad name"
);
assert!(hello(server.control_port(), "HELLO 1 tokens ok").starts_with("PORT "));
}
#[test]
fn a_good_versioned_handshake_is_accepted() {
let server = Server::bind(0).unwrap();
assert!(hello(server.control_port(), "HELLO 1 tokens versioned").starts_with("PORT "));
}
#[test]
fn an_unsupported_version_is_refused_with_its_number() {
let server = Server::bind(0).unwrap();
assert_eq!(
hello(server.control_port(), "HELLO 2 versioned"),
"ERR unsupported protocol version 2"
);
}
#[test]
fn a_hello_with_no_version_is_a_hard_error_not_an_assumed_v1() {
let server = Server::bind(0).unwrap();
assert_eq!(
hello(server.control_port(), "HELLO no-version"),
"ERR missing protocol version"
);
}
#[test]
fn a_non_numeric_version_is_refused() {
let server = Server::bind(0).unwrap();
assert_eq!(
hello(server.control_port(), "HELLO v1 bad-version"),
"ERR bad protocol version"
);
}
#[test]
fn the_anonymous_fallback_still_works_when_the_first_line_is_not_hello_at_all() {
let server = Server::bind(0).unwrap();
let mut s = TcpStream::connect(("127.0.0.1", server.control_port())).unwrap();
s.write_all(b"just a raw capture\n").unwrap();
let name = wait_for(&server, |ev| match ev {
ServerEvent::Opened { name, .. } => Some(name.clone()),
_ => None,
});
assert!(name.starts_with("anon-"), "got {name}");
let got = wait_for(&server, |ev| match ev {
ServerEvent::Bytes { data, .. } => Some(data.clone()),
_ => None,
});
assert_eq!(got, b"just a raw capture\n");
}
#[test]
fn a_non_hello_first_line_becomes_an_anonymous_session() {
let server = Server::bind(0).unwrap();
let mut s = TcpStream::connect(("127.0.0.1", server.control_port())).unwrap();
s.write_all(b"just some tokens\n").unwrap();
let name = wait_for(&server, |ev| match ev {
ServerEvent::Opened { name, .. } => Some(name.clone()),
_ => None,
});
assert!(name.starts_with("anon-"), "got {name}");
let got = wait_for(&server, |ev| match ev {
ServerEvent::Bytes { data, .. } => Some(data.clone()),
_ => None,
});
assert_eq!(got, b"just some tokens\n");
}
#[test]
fn a_reconnect_clears_idle_since_so_the_session_never_reaps_while_reused() {
let mut server = Server::bind(0).unwrap();
let reply = hello(server.control_port(), "HELLO 1 tokens epsilon");
let port: u16 = reply.strip_prefix("PORT ").unwrap().parse().unwrap();
let data = TcpStream::connect(("127.0.0.1", port)).unwrap();
drop(data);
wait_for(&server, |ev| {
matches!(ev, ServerEvent::Disconnected { .. }).then_some(())
});
let second = hello(server.control_port(), "HELLO 1 tokens epsilon");
assert_eq!(reply, second);
server.reap(Duration::from_secs(0));
let deadline = Instant::now() + Duration::from_millis(300);
while Instant::now() < deadline {
if let Ok(ev) = server.events().try_recv() {
assert!(
!matches!(ev, ServerEvent::Closed { .. }),
"reconnected session must not be reaped, got {ev:?}"
);
}
std::thread::sleep(Duration::from_millis(10));
}
let mut again = TcpStream::connect(("127.0.0.1", port)).unwrap();
again.write_all(b"still-here").unwrap();
let got = wait_for(&server, |ev| match ev {
ServerEvent::Bytes { data, .. } => Some(data.clone()),
_ => None,
});
assert_eq!(got, b"still-here");
}
#[test]
fn reaping_sends_closed_and_releases_the_session_listener_port() {
let mut server = Server::bind(0).unwrap();
let reply = hello(server.control_port(), "HELLO 1 tokens zeta");
let port: u16 = reply.strip_prefix("PORT ").unwrap().parse().unwrap();
server.reap(Duration::from_secs(0));
let closed_id = wait_for(&server, |ev| match ev {
ServerEvent::Closed { id } => Some(*id),
_ => None,
});
assert!(closed_id > 0);
let deadline = Instant::now() + Duration::from_secs(2);
loop {
match TcpListener::bind(("127.0.0.1", port)) {
Ok(_) => break,
Err(e) if Instant::now() < deadline => {
std::thread::sleep(Duration::from_millis(10));
let _ = e;
}
Err(e) => panic!("port {port} was not released after reap: {e}"),
}
}
}
#[test]
fn close_session_releases_the_port_even_though_the_session_is_still_live() {
let mut server = Server::bind(0).unwrap();
let reply = hello(server.control_port(), "HELLO 1 tokens closed-window");
let port: u16 = reply.strip_prefix("PORT ").unwrap().parse().unwrap();
let id = wait_for(&server, |ev| match ev {
ServerEvent::Opened { id, .. } => Some(*id),
_ => None,
});
let mut data = TcpStream::connect(("127.0.0.1", port)).unwrap();
data.write_all(b"still streaming").unwrap();
wait_for(&server, |ev| {
matches!(ev, ServerEvent::Bytes { .. }).then_some(())
});
assert_eq!(server.live_count(), 1);
server.close_session(id);
assert_eq!(
server.live_count(),
0,
"the session must be gone, not merely idle"
);
let deadline = Instant::now() + Duration::from_secs(2);
loop {
match TcpListener::bind(("127.0.0.1", port)) {
Ok(_) => break,
Err(e) if Instant::now() < deadline => {
std::thread::sleep(Duration::from_millis(10));
let _ = e;
}
Err(e) => panic!("port {port} was not released after close_session: {e}"),
}
}
let reopened = hello(server.control_port(), "HELLO 1 tokens closed-window");
let reopened_port: u16 = reopened.strip_prefix("PORT ").unwrap().parse().unwrap();
assert_ne!(
reopened_port, port,
"closing a session must not leave it reusable under its old identity"
);
}
#[test]
fn close_then_immediate_redial_ends_up_attached_with_one_writer_and_no_spurious_disconnect() {
let server = Server::bind(0).unwrap();
let reply = hello(server.control_port(), "HELLO 1 tokens race");
let port: u16 = reply.strip_prefix("PORT ").unwrap().parse().unwrap();
let id = wait_for(&server, |ev| match ev {
ServerEvent::Opened { id, .. } => Some(*id),
_ => None,
});
for i in 0..20 {
let mut first = TcpStream::connect(("127.0.0.1", port)).unwrap();
first.write_all(b"x").unwrap();
wait_for(&server, |ev| {
matches!(ev, ServerEvent::Bytes { .. }).then_some(())
});
drop(first);
let mut second = TcpStream::connect(("127.0.0.1", port)).unwrap();
let marker = format!("iter-{i}");
second.write_all(marker.as_bytes()).unwrap();
second
.set_read_timeout(Some(Duration::from_millis(200)))
.unwrap();
let mut banner = [0u8; 64];
match second.read(&mut banner) {
Ok(0) | Err(_) => {}
Ok(n) => {
let text = String::from_utf8_lossy(&banner[..n]);
assert!(
!text.starts_with("ERR"),
"redial rejected on iteration {i}: {text}"
);
}
}
let mut seen = Vec::new();
let mut saw_marker = false;
let find_marker_deadline = Instant::now() + Duration::from_secs(1);
while Instant::now() < find_marker_deadline && !saw_marker {
if let Ok(ev) = server.events().try_recv() {
saw_marker =
matches!(&ev, ServerEvent::Bytes { data, .. } if data == marker.as_bytes());
seen.push(ev);
}
}
assert!(saw_marker, "iteration {i}: never saw the marker bytes");
let settle_deadline = Instant::now() + Duration::from_millis(50);
while Instant::now() < settle_deadline {
if let Ok(ev) = server.events().try_recv() {
seen.push(ev);
}
}
for ev in &seen {
if let ServerEvent::Bytes { data, .. } = ev
&& data != marker.as_bytes()
&& !data.is_empty()
{
panic!(
"iteration {i}: saw interleaved bytes {data:?} alongside marker \
{marker:?} (two writers attached at once)"
);
}
}
let attached_at = seen.iter().position(
|ev| matches!(ev, ServerEvent::Attached { id: eid, reattached: true } if *eid == id),
);
if let Some(attached_at) = attached_at {
assert!(
!seen[attached_at + 1..]
.iter()
.any(|ev| matches!(ev, ServerEvent::Disconnected { .. })),
"iteration {i}: spurious Disconnected after the redial's Attached \
(old guard tore down the new attachment): {seen:?}"
);
}
assert_eq!(
server.live_count(),
1,
"iteration {i}: exactly one writer must be attached"
);
drop(second);
wait_for(&server, |ev| {
matches!(ev, ServerEvent::Disconnected { .. }).then_some(())
});
}
}
#[test]
fn an_anonymous_session_becomes_reapable_after_its_connection_ends() {
let mut server = Server::bind(0).unwrap();
let mut s = TcpStream::connect(("127.0.0.1", server.control_port())).unwrap();
s.write_all(b"raw bytes\n").unwrap();
let id = wait_for(&server, |ev| match ev {
ServerEvent::Opened { id, .. } => Some(*id),
_ => None,
});
drop(s);
wait_for(&server, |ev| {
matches!(ev, ServerEvent::Disconnected { id: eid } if *eid == id).then_some(())
});
server.reap(Duration::from_secs(0));
let closed_id = wait_for(&server, |ev| match ev {
ServerEvent::Closed { id } => Some(*id),
_ => None,
});
assert_eq!(closed_id, id);
}
#[test]
fn a_trace_hello_opens_a_session_with_the_trace_kind() {
let server = Server::bind(0).unwrap();
let reply = hello(server.control_port(), "HELLO 1 trace myapp");
assert!(reply.starts_with("PORT "));
let kind = wait_for(&server, |ev| match ev {
ServerEvent::Opened { kind, .. } => Some(*kind),
_ => None,
});
assert_eq!(kind, StreamKind::Trace);
}
#[test]
fn a_tokens_hello_opens_a_session_with_the_tokens_kind() {
let server = Server::bind(0).unwrap();
hello(server.control_port(), "HELLO 1 tokens build");
let kind = wait_for(&server, |ev| match ev {
ServerEvent::Opened { kind, .. } => Some(*kind),
_ => None,
});
assert_eq!(kind, StreamKind::Tokens);
}
#[test]
fn an_unknown_stream_kind_is_refused_with_its_name() {
let server = Server::bind(0).unwrap();
assert_eq!(
hello(server.control_port(), "HELLO 1 bogus myapp"),
"ERR unknown stream kind bogus"
);
}
#[test]
fn the_old_two_field_hello_form_is_a_hard_error() {
let server = Server::bind(0).unwrap();
assert_eq!(
hello(server.control_port(), "HELLO 1 build-agent"),
"ERR missing stream kind"
);
}
#[test]
fn the_anonymous_fallback_defaults_to_the_tokens_kind() {
let server = Server::bind(0).unwrap();
let mut s = TcpStream::connect(("127.0.0.1", server.control_port())).unwrap();
s.write_all(b"just a raw capture\n").unwrap();
let kind = wait_for(&server, |ev| match ev {
ServerEvent::Opened { kind, .. } => Some(*kind),
_ => None,
});
assert_eq!(kind, StreamKind::Tokens);
}