use std::sync::Arc;
use std::sync::Mutex;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;
use arcbox_virtio::console::{ConsoleIo, SocketConsole};
use crate::device::{DeviceManager, DeviceType};
const INT_VRING: u32 = 1;
const POLL_ACTIVE: Duration = Duration::from_millis(10);
const POLL_IDLE: Duration = Duration::from_millis(250);
pub struct ConsoleRxWorkerContext {
pub device_manager: Arc<DeviceManager>,
pub socket: Arc<Mutex<SocketConsole>>,
pub running: Arc<AtomicBool>,
pub exit_vcpus: Arc<dyn Fn() + Send + Sync>,
}
pub fn console_rx_worker_loop(ctx: ConsoleRxWorkerContext) {
tracing::info!("debug-console RX worker started");
let mut buf = [0u8; 1024];
while ctx.running.load(Ordering::Relaxed) {
let (n, connected) = match ctx.socket.lock() {
Ok(mut s) => {
let n = s.read(&mut buf).unwrap_or(0);
(n, s.is_connected())
}
Err(_) => (0, false),
};
let injected = ctx.device_manager.console_inject_input(&buf[..n]);
if n > 0 {
tracing::debug!(
bytes = n,
connected,
injected,
"debug-console: operator input"
);
}
if injected {
ctx.device_manager
.raise_interrupt_for(DeviceType::VirtioConsole, INT_VRING);
(ctx.exit_vcpus)();
}
std::thread::sleep(if connected { POLL_ACTIVE } else { POLL_IDLE });
}
tracing::info!("debug-console RX worker stopped");
}