#![allow(dead_code)]
use kovan_channel::flavors::unbounded::{Receiver, Sender, channel};
use std::io::{BufRead, BufReader, Read, Write};
use std::net::{SocketAddr, TcpListener, TcpStream};
use std::process::{Child, ChildStderr, ChildStdout, Command, Stdio};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::thread;
use std::time::{Duration, Instant};
pub fn pick_port() -> u16 {
let l = TcpListener::bind("127.0.0.1:0").expect("bind ephemeral");
let p = l.local_addr().expect("local_addr").port();
drop(l);
p
}
pub struct ChildGuard(Option<Child>);
impl ChildGuard {
pub fn new(c: Child) -> Self {
Self(Some(c))
}
}
impl Drop for ChildGuard {
fn drop(&mut self) {
if let Some(mut c) = self.0.take() {
let _ = c.kill();
let _ = c.wait();
}
}
}
impl std::ops::Deref for ChildGuard {
type Target = Child;
fn deref(&self) -> &Child {
self.0.as_ref().expect("child taken")
}
}
impl std::ops::DerefMut for ChildGuard {
fn deref_mut(&mut self) -> &mut Child {
self.0.as_mut().expect("child taken")
}
}
pub struct StdoutWatcher {
pub marker_seen: Arc<AtomicBool>,
lines_rx: Receiver<String>,
}
impl StdoutWatcher {
pub fn spawn(stdout: ChildStdout, marker: String) -> Self {
let marker_seen = Arc::new(AtomicBool::new(false));
let (tx, rx): (Sender<String>, Receiver<String>) = channel();
let seen = marker_seen.clone();
thread::spawn(move || {
let reader = BufReader::new(stdout);
for line in reader.lines() {
let Ok(l) = line else {
break;
};
if l.contains(&marker) {
seen.store(true, Ordering::Release);
}
tx.send(l);
}
});
Self {
marker_seen,
lines_rx: rx,
}
}
pub fn drain_lines(&self) -> Vec<String> {
let mut out = Vec::new();
while let Some(l) = self.lines_rx.try_recv() {
out.push(l);
}
out
}
}
pub fn wait_for_listening(
watcher: &StdoutWatcher,
port: u16,
timeout: Duration,
stderr: Option<&mut ChildStderr>,
) -> Result<(), String> {
let addr = SocketAddr::from(([127, 0, 0, 1], port));
let deadline = Instant::now() + timeout;
while Instant::now() < deadline {
if watcher.marker_seen.load(Ordering::Acquire) {
return Ok(());
}
if TcpStream::connect_timeout(&addr, Duration::from_millis(100)).is_ok() {
return Ok(());
}
thread::sleep(Duration::from_millis(50));
}
let lines = watcher.drain_lines();
let mut err_buf = String::new();
if let Some(s) = stderr {
let _ = s.read_to_string(&mut err_buf);
}
Err(format!(
"daemon never became ready on :{port} within {timeout:?}\n\
--- buffered stdout ({n} lines) ---\n{stdout}\n\
--- captured stderr ---\n{stderr}",
n = lines.len(),
stdout = lines.join("\n"),
stderr = err_buf,
))
}
pub fn listening_marker_js(port: u16) -> String {
format!("console.log('LISTENING:{port}');")
}
pub fn spawn_and_wait_listening(
mut cmd: Command,
port: u16,
timeout: Duration,
) -> (ChildGuard, StdoutWatcher) {
cmd.stdout(Stdio::piped()).stderr(Stdio::piped());
let mut child = cmd.spawn().expect("spawn burn");
let stdout = child.stdout.take().expect("stdout piped");
let watcher = StdoutWatcher::spawn(stdout, format!("LISTENING:{port}"));
let mut guard = ChildGuard::new(child);
if let Err(e) = wait_for_listening(&watcher, port, timeout, guard.stderr.as_mut()) {
panic!("{e}");
}
(guard, watcher)
}
pub fn http_get(port: u16, path: &str) -> String {
let mut stream = TcpStream::connect(("127.0.0.1", port)).expect("connect");
stream.set_read_timeout(Some(Duration::from_secs(5))).ok();
let req = format!("GET {path} HTTP/1.1\r\nHost: 127.0.0.1:{port}\r\nConnection: close\r\n\r\n");
stream.write_all(req.as_bytes()).expect("write");
let mut resp = String::new();
stream.read_to_string(&mut resp).expect("read");
resp
}
pub fn http_post(port: u16, path: &str, body: &str, content_type: &str) -> String {
let mut stream = TcpStream::connect(("127.0.0.1", port)).expect("connect");
stream.set_read_timeout(Some(Duration::from_secs(5))).ok();
let req = format!(
"POST {path} HTTP/1.1\r\nHost: 127.0.0.1:{port}\r\nContent-Type: {content_type}\r\n\
Content-Length: {len}\r\nConnection: close\r\n\r\n{body}",
len = body.len(),
);
stream.write_all(req.as_bytes()).expect("write");
let mut resp = String::new();
stream.read_to_string(&mut resp).expect("read");
resp
}
pub fn wait_for_listener(port: u16, timeout: Duration) -> bool {
let addr = SocketAddr::from(([127, 0, 0, 1], port));
let start = Instant::now();
while start.elapsed() < timeout {
if TcpStream::connect_timeout(&addr, Duration::from_millis(100)).is_ok() {
return true;
}
thread::sleep(Duration::from_millis(50));
}
false
}
pub fn extract_body(resp: &str) -> &str {
resp.split_once("\r\n\r\n").map(|(_, b)| b).unwrap_or("")
}
pub fn run_burn_capped(
args: &[&str],
envs: &[(&str, &str)],
timeout: Duration,
) -> std::process::Output {
let mut cmd = Command::new(env!("CARGO_BIN_EXE_burn"));
cmd.args(args);
for (k, v) in envs {
cmd.env(k, v);
}
let mut child = cmd
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.expect("spawn burn");
let mut out_pipe = child.stdout.take().expect("stdout pipe");
let mut err_pipe = child.stderr.take().expect("stderr pipe");
let out_h = thread::spawn(move || {
let mut b = Vec::new();
let _ = out_pipe.read_to_end(&mut b);
b
});
let err_h = thread::spawn(move || {
let mut b = Vec::new();
let _ = err_pipe.read_to_end(&mut b);
b
});
let deadline = Instant::now() + timeout;
loop {
if let Some(status) = child.try_wait().expect("try_wait burn") {
let stdout = out_h.join().unwrap_or_default();
let stderr = err_h.join().unwrap_or_default();
return std::process::Output {
status,
stdout,
stderr,
};
}
if Instant::now() >= deadline {
let _ = child.kill();
let _ = child.wait();
let stdout = out_h.join().unwrap_or_default();
let stderr = err_h.join().unwrap_or_default();
panic!(
"burn subprocess wedged: no exit within {timeout:?}, killed.\n\
The runtime event loop hung, so the in-program watchdog never fired.\n\
args: {args:?}\n--- stdout so far ---\n{}\n--- stderr so far ---\n{}",
String::from_utf8_lossy(&stdout),
String::from_utf8_lossy(&stderr),
);
}
thread::sleep(Duration::from_millis(25));
}
}