use std::sync::Arc;
use axum::{
extract::{
Path as UrlPath, State,
ws::{Message, WebSocket, WebSocketUpgrade},
},
http::StatusCode,
response::{IntoResponse, Response},
};
use tokio::sync::broadcast::error::RecvError;
use crate::{
console::{Console, ConsoleEvent},
frame::{ToConsole, ToViewer},
registry::Registry,
};
const BACKLOG_LINES: usize = 2000;
pub async fn attach(
upgrade: WebSocketUpgrade,
UrlPath(name): UrlPath<String>,
State(registry): State<Arc<Registry>>,
) -> Response {
match registry.resolve(Some(&name)) {
Ok(console) => upgrade.on_upgrade(move |socket| pump(socket, console)),
Err(e) => (StatusCode::BAD_REQUEST, e).into_response(),
}
}
async fn pump(mut socket: WebSocket, console: Arc<Console>) {
let mut events = console.subscribe();
let hello = ToViewer::Hello {
console: console.name().to_string(),
baud: console.baud(),
connected: console.connected(),
backlog: console.snapshot(BACKLOG_LINES),
};
if !send(&mut socket, &hello).await {
return;
}
loop {
tokio::select! {
event = events.recv() => match event {
Ok(event) => {
if !send(&mut socket, &frame_of(&console, event)).await {
return;
}
}
Err(RecvError::Lagged(missed)) => {
let note = ToViewer::System {
text: format!("{missed} lines dropped, viewer fell behind"),
};
if !send(&mut socket, ¬e).await {
return;
}
}
Err(RecvError::Closed) => return,
},
incoming = socket.recv() => match incoming {
Some(Ok(Message::Text(text))) => {
if !apply(&console, &text) {
return;
}
}
Some(Ok(Message::Close(_))) | None => return,
Some(Ok(_)) => {}
Some(Err(_)) => return,
},
}
}
}
fn apply(console: &Arc<Console>, text: &str) -> bool {
let outcome = match serde_json::from_str::<ToConsole>(text) {
Ok(ToConsole::Line { text }) => console.queue_line(&text),
Ok(ToConsole::Ctrl { ctrl }) => console.queue_ctrl(ctrl),
Err(e) => {
console.push_system(&format!("viewer sent a frame that made no sense: {e}"));
return true;
}
};
match outcome {
Ok(()) => true,
Err(e) => {
console.push_system(&e);
true
}
}
}
fn frame_of(console: &Arc<Console>, event: ConsoleEvent) -> ToViewer {
match event {
ConsoleEvent::Rx(bytes) => ToViewer::Rx {
data: String::from_utf8_lossy(&bytes).into_owned(),
},
ConsoleEvent::Echo { origin, text } => ToViewer::Echo { origin, text },
ConsoleEvent::System(text) => ToViewer::System { text },
ConsoleEvent::Connected => ToViewer::Connected {
connected: console.connected(),
},
}
}
async fn send(socket: &mut WebSocket, frame: &ToViewer) -> bool {
let Ok(text) = serde_json::to_string(frame) else {
return true;
};
socket.send(Message::Text(text.into())).await.is_ok()
}