use std::{
net::TcpStream,
path::{Path, PathBuf},
sync::{
Arc, Mutex,
atomic::{AtomicBool, Ordering},
mpsc,
},
thread,
time::{Duration, Instant},
};
use anyhow::{Context, bail};
use super::types::AxvisorHttpProbeConfig;
pub(crate) type HostHttpProbeFn = Box<dyn FnOnce() -> anyhow::Result<()> + Send + 'static>;
const CONNECT_RETRY_INTERVAL: Duration = Duration::from_millis(100);
const QMP_CONNECT_RETRY_INTERVAL: Duration = Duration::from_millis(100);
const QMP_CONNECT_RETRIES: usize = 10;
const QMP_EXIT_WAIT: Duration = Duration::from_secs(4);
const QMP_READ_POLL_INTERVAL: Duration = Duration::from_millis(100);
const QMP_QUIT_RETRIES: usize = 4;
pub(crate) struct HostHttpProbeGuard {
stop: Arc<AtomicBool>,
result: Arc<Mutex<Option<anyhow::Result<()>>>>,
thread: Option<thread::JoinHandle<()>>,
}
impl HostHttpProbeGuard {
pub(crate) fn start(
config: &AxvisorHttpProbeConfig,
host_port: u16,
case_name: &str,
qmp_socket: Option<PathBuf>,
stop: Arc<AtomicBool>,
probe: HostHttpProbeFn,
) -> anyhow::Result<Self> {
let addr = format!("127.0.0.1:{host_port}");
let connect_timeout = Duration::from_secs(config.connect_timeout_secs);
let thread_stop = stop.clone();
let result = Arc::new(Mutex::new(None));
let thread_result = result.clone();
let case_name = case_name.to_string();
let (ready_tx, ready_rx) = mpsc::channel();
let thread_addr = addr.clone();
let thread_case_name = case_name.clone();
let thread = thread::spawn(move || {
let _ = ready_tx.send(());
let verdict = (|| -> anyhow::Result<()> {
wait_for_port_ready(&thread_addr, connect_timeout, &thread_stop).with_context(
|| {
format!(
"guest HTTP server never became reachable within {connect_timeout:?}"
)
},
)?;
probe()
})();
*thread_result.lock().unwrap() = Some(verdict);
if let Some(socket) = qmp_socket
&& let Err(err) = request_qmp_quit(&socket)
{
eprintln!(
" host http probe: {thread_case_name}: failed to quit QEMU via QMP: {err:#}"
);
}
});
if ready_rx.recv_timeout(Duration::from_secs(1)).is_err() {
stop.store(true, Ordering::Release);
bail!("host http probe for `{case_name}` did not become ready");
}
println!(" host http probe: {addr} -> guest:{}", config.guest_port);
Ok(Self {
stop,
result,
thread: Some(thread),
})
}
pub(crate) fn take_result(&self) -> Option<anyhow::Result<()>> {
self.result.lock().unwrap().take()
}
}
impl Drop for HostHttpProbeGuard {
fn drop(&mut self) {
self.stop.store(true, Ordering::Release);
if let Some(thread) = self.thread.take() {
let _ = thread.join();
}
}
}
fn wait_for_port_ready(
addr: &str,
connect_timeout: Duration,
stop: &AtomicBool,
) -> anyhow::Result<()> {
let started = Instant::now();
loop {
if stop.load(Ordering::Acquire) {
bail!("host http probe stopped");
}
if started.elapsed() >= connect_timeout {
bail!("timed out after {connect_timeout:?}");
}
if TcpStream::connect(addr).is_ok() {
return Ok(());
}
thread::sleep(CONNECT_RETRY_INTERVAL);
}
}
#[cfg(unix)]
fn request_qmp_quit(socket: &Path) -> anyhow::Result<()> {
use std::{
io::{ErrorKind, Read, Write},
os::unix::net::UnixStream,
};
fn connect_with_retries(socket: &Path) -> anyhow::Result<UnixStream> {
let mut last_err = None;
for _ in 0..QMP_CONNECT_RETRIES {
match UnixStream::connect(socket) {
Ok(stream) => return Ok(stream),
Err(err) => {
last_err = Some(err);
thread::sleep(QMP_CONNECT_RETRY_INTERVAL);
}
}
}
bail!(
"failed to connect QMP socket {}: {}",
socket.display(),
last_err.as_ref().expect("at least one connect attempted")
)
}
fn qmp_handshake_quit(stream: &mut UnixStream) -> std::io::Result<()> {
stream
.set_read_timeout(Some(Duration::from_millis(200)))
.ok();
stream
.set_write_timeout(Some(Duration::from_millis(200)))
.ok();
let mut buf = [0_u8; 512];
let _ = stream.read(&mut buf); stream.write_all(b"{\"execute\":\"qmp_capabilities\"}\r\n")?;
buf.fill(0);
let _ = stream.read(&mut buf); stream.write_all(b"{\"execute\":\"quit\"}\r\n")?;
stream.flush()
}
fn socket_connectable(socket: &Path) -> bool {
UnixStream::connect(socket).is_ok()
}
for _ in 0..QMP_QUIT_RETRIES {
let mut stream = connect_with_retries(socket)?;
if let Err(err) = qmp_handshake_quit(&mut stream) {
if socket_connectable(socket) {
bail!("failed to send QMP quit: {err}");
}
return Ok(());
}
let wait_started = Instant::now();
loop {
if wait_started.elapsed() >= QMP_EXIT_WAIT {
break;
}
let mut buf = [0_u8; 512];
match stream.read(&mut buf) {
Ok(0) => return Ok(()), Err(err) if matches!(err.kind(), ErrorKind::WouldBlock | ErrorKind::TimedOut) => {
if !socket_connectable(socket) {
return Ok(()); }
}
Err(_) => return Ok(()), Ok(_) => {} }
thread::sleep(QMP_READ_POLL_INTERVAL);
}
}
Ok(())
}
#[cfg(not(unix))]
fn request_qmp_quit(_socket: &Path) -> anyhow::Result<()> {
bail!("QMP unix sockets are not supported on this host")
}