use core::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::{
string::String,
sync::{Arc, OnceLock},
time::Duration,
};
use anyhow::{Context, Result};
use axvisor::console_mux::HostOutputQueue;
use axvm::VMId;
use {
ax_std::os::arceos::api::task::AxCpuMask,
ax_std::os::arceos::api::task::ax_set_current_affinity,
ax_std::os::arceos::modules::ax_runtime::task::sync::irq::IrqWaitCell,
ax_std::os::arceos::modules::ax_runtime::task::sync::irq::IrqWorkerWaiter,
ax_std::os::arceos::modules::ax_runtime::task::thread::current::current_thread_handle,
ax_std::os::arceos::sync::NoPreemptMutex,
};
mod layout;
mod delivery;
use delivery::{DeliveryFrame, DeliveryQueue};
use layout::{ConsoleLane, Endpoint, plan_endpoints};
const CONSOLE_LANE_COUNT: usize = ConsoleLane::COUNT;
const OUTPUT_QUEUE_CAPACITY: usize = 64 * 1024;
const OUTPUT_BATCH_CAPACITY: usize = 512;
const OUTPUT_FRAME_TARGET_CAPACITY: usize = 4096;
const OUTPUT_DELIVERY_QUEUE_CAPACITY: usize = 64 * 1024;
const OUTPUT_DELIVERY_BATCH_CAPACITY: usize = 4096;
const OUTPUT_COALESCE_WINDOW: Duration = Duration::from_millis(10);
const MANAGEMENT_LINE_CAPACITY: usize = 256;
const MANAGEMENT_CPU_ID: usize = 0;
static OUTPUT_HUB: NetworkOutputHub = NetworkOutputHub::new();
static ENDPOINTS: OnceLock<Vec<Endpoint>> = OnceLock::new();
struct NetworkOutputHub {
lanes: [NetworkOutputLane; CONSOLE_LANE_COUNT],
}
struct NetworkOutputLane {
queue: NoPreemptMutex<HostOutputQueue<OUTPUT_QUEUE_CAPACITY>>,
connected: AtomicBool,
session: AtomicUsize,
ready: IrqWaitCell,
}
struct BrowserOutputDelivery {
queue: NoPreemptMutex<DeliveryQueue<OUTPUT_DELIVERY_QUEUE_CAPACITY>>,
closed: AtomicBool,
ready: IrqWaitCell,
}
impl NetworkOutputHub {
const fn new() -> Self {
Self {
lanes: [
NetworkOutputLane::new(),
NetworkOutputLane::new(),
NetworkOutputLane::new(),
NetworkOutputLane::new(),
],
}
}
fn submit(&self, lane: ConsoleLane, bytes: &[u8]) {
self.lanes[lane.index()].submit(bytes);
}
fn is_connected(&self, lane: ConsoleLane) -> bool {
self.lanes[lane.index()].connected.load(Ordering::Acquire)
}
fn begin_session(&self, lane: ConsoleLane) -> Option<usize> {
self.lanes[lane.index()].begin_session()
}
fn end_session(&self, lane: ConsoleLane, session: usize) {
self.lanes[lane.index()].end_session(session);
}
fn receive(
&self,
lane: ConsoleLane,
session: usize,
waiter: &IrqWorkerWaiter,
) -> Result<Option<NetworkOutputBatch>> {
self.lanes[lane.index()].receive(session, waiter)
}
fn take_batch(&self, lane: ConsoleLane, session: usize) -> Option<NetworkOutputBatch> {
self.lanes[lane.index()].take_batch(session)
}
}
impl NetworkOutputLane {
const fn new() -> Self {
Self {
queue: NoPreemptMutex::new(HostOutputQueue::new()),
connected: AtomicBool::new(false),
session: AtomicUsize::new(0),
ready: IrqWaitCell::new(),
}
}
fn submit(&self, bytes: &[u8]) {
if bytes.is_empty() || !self.connected.load(Ordering::Acquire) {
return;
}
let submitted = {
let mut queue = self.queue.lock();
if !self.connected.load(Ordering::Acquire) {
false
} else {
queue.enqueue(bytes);
true
}
};
if submitted {
let _result = self.ready.notify();
}
}
fn begin_session(&self) -> Option<usize> {
let mut queue = self.queue.lock();
self.connected
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
.ok()?;
*queue = HostOutputQueue::new();
Some(self.session.fetch_add(1, Ordering::AcqRel).wrapping_add(1))
}
fn end_session(&self, session: usize) {
let mut queue = self.queue.lock();
if self.session.load(Ordering::Acquire) == session {
self.connected.store(false, Ordering::Release);
*queue = HostOutputQueue::new();
drop(queue);
let _result = self.ready.notify();
}
}
fn take_batch(&self, session: usize) -> Option<NetworkOutputBatch> {
let mut queue = self.queue.lock();
if !self.connected.load(Ordering::Acquire)
|| self.session.load(Ordering::Acquire) != session
{
return None;
}
let mut batch = NetworkOutputBatch::new();
batch.dropped_bytes = queue.take_dropped_bytes();
batch.len = queue.dequeue(&mut batch.bytes);
(!batch.is_empty()).then_some(batch)
}
fn receive(
&self,
session: usize,
waiter: &IrqWorkerWaiter,
) -> Result<Option<NetworkOutputBatch>> {
loop {
if let Some(batch) = self.take_batch(session) {
return Ok(Some(batch));
}
if !self.connected.load(Ordering::Acquire)
|| self.session.load(Ordering::Acquire) != session
{
return Ok(None);
}
waiter
.wait(&self.ready)
.context("failed to wait for browser console output")?;
}
}
}
impl BrowserOutputDelivery {
const fn new() -> Self {
Self {
queue: NoPreemptMutex::new(DeliveryQueue::new()),
closed: AtomicBool::new(false),
ready: IrqWaitCell::new(),
}
}
fn enqueue(&self, bytes: &[u8]) {
if self.closed.load(Ordering::Acquire) {
return;
}
let submitted = {
let mut queue = self.queue.lock();
if self.closed.load(Ordering::Acquire) {
false
} else {
queue.enqueue(bytes);
true
}
};
if submitted {
let _result = self.ready.notify();
}
}
fn close(&self) {
self.closed.store(true, Ordering::Release);
let _result = self.ready.notify();
}
fn receive(&self, waiter: &IrqWorkerWaiter) -> Result<Option<Vec<u8>>> {
loop {
let mut bytes = [0; OUTPUT_DELIVERY_BATCH_CAPACITY];
let (len, dropped_bytes) = self.queue.lock().dequeue(&mut bytes);
if len != 0 || dropped_bytes != 0 {
let mut frame = DeliveryFrame::with_capacity(len + 96);
frame.append(&bytes[..len], dropped_bytes);
return Ok(Some(frame.into_bytes()));
}
if self.closed.load(Ordering::Acquire) {
return Ok(None);
}
waiter
.wait(&self.ready)
.context("failed to wait for browser console delivery")?;
}
}
}
struct NetworkOutputBatch {
bytes: [u8; OUTPUT_BATCH_CAPACITY],
len: usize,
dropped_bytes: usize,
}
impl NetworkOutputBatch {
const fn new() -> Self {
Self {
bytes: [0; OUTPUT_BATCH_CAPACITY],
len: 0,
dropped_bytes: 0,
}
}
const fn is_empty(&self) -> bool {
self.len == 0 && self.dropped_bytes == 0
}
}
struct ActiveSession {
lane: ConsoleLane,
session: usize,
}
impl ActiveSession {
fn install(lane: ConsoleLane) -> Result<Self> {
let session = OUTPUT_HUB.begin_session(lane).with_context(|| {
format!("{} console already has an active session", lane_name(lane))
})?;
Ok(Self { lane, session })
}
}
impl Drop for ActiveSession {
fn drop(&mut self) {
OUTPUT_HUB.end_session(self.lane, self.session);
}
}
pub(crate) fn start() -> Result<()> {
ENDPOINTS
.set(build_startup_endpoints())
.map_err(|_| anyhow::anyhow!("browser console endpoints were already initialized"))
}
fn build_startup_endpoints() -> Vec<Endpoint> {
let guests = crate::manager::AxvmManager::vm_list()
.into_iter()
.map(|vm| (vm.id(), vm.name()))
.collect();
plan_endpoints(guests)
}
fn endpoints() -> &'static [Endpoint] {
ENDPOINTS.get().map(Vec::as_slice).unwrap_or(&[])
}
fn endpoint_for_lane(lane: ConsoleLane) -> Option<&'static Endpoint> {
endpoints().iter().find(|endpoint| endpoint.lane == lane)
}
fn endpoint_for_route(route: &str) -> Option<&'static Endpoint> {
endpoints().iter().find(|endpoint| endpoint.route == route)
}
fn lane_name(lane: ConsoleLane) -> String {
endpoint_for_lane(lane)
.map(|endpoint| endpoint.display_name.clone())
.unwrap_or_else(|| format!("console lane {}", lane.index()))
}
pub(crate) fn pin_current_task() {
ax_set_current_affinity(AxCpuMask::one_shot(MANAGEMENT_CPU_ID))
.expect("web console management CPU affinity must be valid");
}
pub(crate) fn has_console_route(route: &str) -> bool {
endpoint_for_route(route).is_some()
}
pub(crate) fn console_descriptions() -> Vec<ConsoleDescription> {
endpoints()
.iter()
.map(|endpoint| ConsoleDescription {
route: endpoint.route.clone(),
display_name: endpoint.display_name.clone(),
})
.collect()
}
pub(crate) struct ConsoleDescription {
pub(crate) route: String,
pub(crate) display_name: String,
}
pub(crate) fn submit_management_output(bytes: &[u8]) {
OUTPUT_HUB.submit(ConsoleLane::MANAGEMENT, bytes);
}
pub(crate) fn guest_output_connected(vm_id: VMId) -> bool {
endpoints()
.iter()
.find(|endpoint| endpoint.vm_id == Some(vm_id))
.is_some_and(|endpoint| OUTPUT_HUB.is_connected(endpoint.lane))
}
pub(crate) fn submit_guest_output(vm_id: VMId, bytes: &[u8]) {
let Some(lane) = endpoints()
.iter()
.find(|endpoint| endpoint.vm_id == Some(vm_id))
.map(|endpoint| endpoint.lane)
else {
return;
};
OUTPUT_HUB.submit(lane, bytes);
}
pub(crate) fn open_browser_console(
route: &str,
) -> Result<(BrowserConsoleInput, BrowserConsoleOutput)> {
let endpoint = endpoint_for_route(route)
.cloned()
.with_context(|| format!("unknown console endpoint `{route}`"))?;
let active_session = ActiveSession::install(endpoint.lane)?;
let lane = active_session.lane;
let session = active_session.session;
let delivery = start_output_dispatcher(lane, session)?;
Ok((
BrowserConsoleInput {
endpoint,
editor: ManagementLineEditor::new(),
_active_session: active_session,
},
BrowserConsoleOutput {
delivery,
waiter: None,
},
))
}
fn start_output_dispatcher(
lane: ConsoleLane,
session: usize,
) -> Result<Arc<BrowserOutputDelivery>> {
let delivery = Arc::new(BrowserOutputDelivery::new());
let dispatcher_delivery = Arc::clone(&delivery);
let task_name = format!("{}-browser-console-dispatcher", lane_name(lane));
let worker_name = task_name.clone();
std::thread::Builder::new()
.name(task_name.clone())
.spawn(move || {
if let Err(error) = run_output_dispatcher(lane, session, &dispatcher_delivery) {
warn!("{worker_name} stopped: {error:#}");
}
dispatcher_delivery.close();
})
.with_context(|| format!("failed to start {task_name}"))?;
Ok(delivery)
}
fn run_output_dispatcher(
lane: ConsoleLane,
session: usize,
delivery: &BrowserOutputDelivery,
) -> Result<()> {
let current = current_thread_handle()
.context("failed to bind browser console dispatcher to its worker")?;
let waiter = IrqWorkerWaiter::new(current.wake_handle());
while let Some(frame) = receive_output_frame(lane, session, &waiter)? {
delivery.enqueue(&frame);
}
Ok(())
}
fn receive_output_frame(
lane: ConsoleLane,
session: usize,
waiter: &IrqWorkerWaiter,
) -> Result<Option<Vec<u8>>> {
let Some(first) = OUTPUT_HUB.receive(lane, session, waiter)? else {
return Ok(None);
};
let mut frame = DeliveryFrame::with_capacity(OUTPUT_FRAME_TARGET_CAPACITY);
frame.append(&first.bytes[..first.len], first.dropped_bytes);
std::thread::sleep(OUTPUT_COALESCE_WINDOW);
while frame.len() < OUTPUT_FRAME_TARGET_CAPACITY {
let Some(batch) = OUTPUT_HUB.take_batch(lane, session) else {
break;
};
frame.append(&batch.bytes[..batch.len], batch.dropped_bytes);
}
Ok(Some(frame.into_bytes()))
}
pub(crate) struct BrowserConsoleInput {
endpoint: Endpoint,
editor: ManagementLineEditor,
_active_session: ActiveSession,
}
impl BrowserConsoleInput {
pub(crate) fn greeting(&self) -> String {
if let Some(vm_id) = self.endpoint.vm_id {
format!("[Axvisor] browser console attached to VM {vm_id}\r\n")
} else {
format!(
"Welcome to AxVisor Browser Shell!\r\nType 'help' for commands.\r\n{}",
crate::shell::network_prompt()
)
}
}
pub(crate) fn route(&mut self, bytes: &[u8]) -> bool {
if let Some(vm_id) = self.endpoint.vm_id {
crate::guest_console::route_network_input(vm_id, bytes);
true
} else {
self.editor.process(bytes)
}
}
}
pub(crate) struct BrowserConsoleOutput {
delivery: Arc<BrowserOutputDelivery>,
waiter: Option<IrqWorkerWaiter>,
}
impl BrowserConsoleOutput {
pub(crate) fn receive(&mut self) -> Result<Option<Vec<u8>>> {
if self.waiter.is_none() {
let current = current_thread_handle()
.context("failed to bind browser console output to its worker")?;
self.waiter = Some(IrqWorkerWaiter::new(current.wake_handle()));
}
let waiter = self
.waiter
.as_ref()
.expect("browser console waiter was initialized above");
self.delivery.receive(waiter)
}
}
struct ManagementLineEditor {
line: [u8; MANAGEMENT_LINE_CAPACITY],
len: usize,
previous_was_cr: bool,
}
impl ManagementLineEditor {
const fn new() -> Self {
Self {
line: [0; MANAGEMENT_LINE_CAPACITY],
len: 0,
previous_was_cr: false,
}
}
fn process(&mut self, bytes: &[u8]) -> bool {
for &byte in bytes {
if byte == b'\n' && self.previous_was_cr {
self.previous_was_cr = false;
continue;
}
self.previous_was_cr = byte == b'\r';
match byte {
b'\r' | b'\n' => {
submit_management_output(b"\r\n");
let command = String::from_utf8_lossy(&self.line[..self.len]);
self.len = 0;
if !crate::shell::run_network_command(&command) {
submit_management_output(b"Goodbye!\r\n");
return false;
}
submit_management_output(crate::shell::network_prompt().as_bytes());
}
b'\x08' | b'\x7f' if self.len != 0 => {
self.len -= 1;
submit_management_output(b"\x08 \x08");
}
0x20..=0x7e if self.len < self.line.len() => {
self.line[self.len] = byte;
self.len += 1;
submit_management_output(&[byte]);
}
0x20..=0x7e => submit_management_output(b"\x07"),
_ => {}
}
}
true
}
}