feat: add rust browser terminal prototype
This commit is contained in:
20
server/Cargo.toml
Normal file
20
server/Cargo.toml
Normal file
@@ -0,0 +1,20 @@
|
||||
[package]
|
||||
name = "server"
|
||||
version = "0.1.0"
|
||||
edition = "2024"
|
||||
|
||||
[dependencies]
|
||||
app = { path = "../app", default-features = false, features = ["ssr"] }
|
||||
app_config.workspace = true
|
||||
axum.workspace = true
|
||||
futures-util.workspace = true
|
||||
leptos = { workspace = true, features = ["ssr"] }
|
||||
leptos_axum.workspace = true
|
||||
portable-pty.workspace = true
|
||||
serde_json.workspace = true
|
||||
tokio.workspace = true
|
||||
tower.workspace = true
|
||||
tower-http.workspace = true
|
||||
tracing.workspace = true
|
||||
tracing-subscriber.workspace = true
|
||||
dotenvy = "0.15"
|
||||
93
server/src/main.rs
Normal file
93
server/src/main.rs
Normal file
@@ -0,0 +1,93 @@
|
||||
use app::app::App;
|
||||
use app::shell::shell;
|
||||
use app_config::SiteConfig;
|
||||
use axum::Router;
|
||||
use leptos::logging::log;
|
||||
use leptos::prelude::*;
|
||||
use leptos_axum::{LeptosRoutes, generate_route_list, handle_server_fns};
|
||||
|
||||
mod terminal;
|
||||
|
||||
#[tokio::main(flavor = "multi_thread")]
|
||||
async fn main() -> Result<(), Box<dyn std::error::Error>> {
|
||||
load_env_files();
|
||||
|
||||
tracing_subscriber::fmt()
|
||||
.with_env_filter(
|
||||
tracing_subscriber::EnvFilter::from_default_env().add_directive("info".parse()?),
|
||||
)
|
||||
.init();
|
||||
|
||||
let conf = get_configuration(None).map_err(|err| {
|
||||
tracing::error!("Failed to get Leptos configuration: {err:?}");
|
||||
err
|
||||
})?;
|
||||
let addr = conf.leptos_options.site_addr;
|
||||
let leptos_options = conf.leptos_options;
|
||||
let routes = generate_route_list(App);
|
||||
|
||||
let leptos_router = Router::new()
|
||||
.leptos_routes(&leptos_options, routes, {
|
||||
let leptos_options = leptos_options.clone();
|
||||
move || shell(leptos_options.clone())
|
||||
})
|
||||
.fallback(leptos_axum::file_and_error_handler(shell))
|
||||
.with_state(leptos_options);
|
||||
|
||||
let terminal_routes = Router::new().route(
|
||||
"/terminal/ws",
|
||||
axum::routing::get(terminal::ws::terminal_ws_handler),
|
||||
);
|
||||
|
||||
let app = if SiteConfig::BASE_PATH.is_empty() {
|
||||
terminal_routes.merge(leptos_router)
|
||||
} else {
|
||||
Router::new()
|
||||
.route(
|
||||
"/rustui/api/{*fn_name}",
|
||||
axum::routing::post(handle_server_fns),
|
||||
)
|
||||
.route("/api/{*fn_name}", axum::routing::post(handle_server_fns))
|
||||
.merge(terminal_routes.clone())
|
||||
.nest(SiteConfig::BASE_PATH, terminal_routes.merge(leptos_router))
|
||||
};
|
||||
|
||||
log!(
|
||||
"BASE_PATH demo listening on http://{addr}{}",
|
||||
SiteConfig::BASE_PATH
|
||||
);
|
||||
let listener = tokio::net::TcpListener::bind(&addr).await?;
|
||||
axum::serve(listener, app.into_make_service()).await?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn load_env_files() {
|
||||
for filename in [".env", ".env.local"] {
|
||||
let _ = dotenvy::from_filename(filename);
|
||||
}
|
||||
|
||||
let env = std::env::var("NEXT_PUBLIC_ENV").unwrap_or_else(|_| "development".to_owned());
|
||||
|
||||
for filename in [format!(".env.{env}"), format!(".env.{env}.local")] {
|
||||
let _ = dotenvy::from_filename(&filename);
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn root_path_gets_base_path() {
|
||||
assert_eq!(SiteConfig::with_base_path("/"), "/rustui");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn nested_path_gets_base_path() {
|
||||
assert_eq!(
|
||||
SiteConfig::with_base_path("/deep/nested"),
|
||||
"/rustui/deep/nested"
|
||||
);
|
||||
}
|
||||
}
|
||||
2
server/src/terminal/mod.rs
Normal file
2
server/src/terminal/mod.rs
Normal file
@@ -0,0 +1,2 @@
|
||||
pub mod pty_session;
|
||||
pub mod ws;
|
||||
149
server/src/terminal/pty_session.rs
Normal file
149
server/src/terminal/pty_session.rs
Normal file
@@ -0,0 +1,149 @@
|
||||
use std::io::{Read, Write};
|
||||
use std::sync::{Arc, Mutex};
|
||||
|
||||
use portable_pty::{
|
||||
ChildKiller, CommandBuilder, MasterPty, PtySize, native_pty_system,
|
||||
};
|
||||
use tokio::sync::mpsc;
|
||||
|
||||
#[derive(Clone, Copy, Debug)]
|
||||
pub struct TerminalSize {
|
||||
pub cols: u16,
|
||||
pub rows: u16,
|
||||
pub pixel_width: u16,
|
||||
pub pixel_height: u16,
|
||||
}
|
||||
|
||||
impl Default for TerminalSize {
|
||||
fn default() -> Self {
|
||||
Self {
|
||||
cols: 80,
|
||||
rows: 24,
|
||||
pixel_width: 720,
|
||||
pixel_height: 432,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl TerminalSize {
|
||||
fn into_pty_size(self) -> PtySize {
|
||||
PtySize {
|
||||
cols: self.cols,
|
||||
rows: self.rows,
|
||||
pixel_width: self.pixel_width,
|
||||
pixel_height: self.pixel_height,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub struct PtySession {
|
||||
master: Box<dyn MasterPty + Send>,
|
||||
writer: Arc<Mutex<Box<dyn Write + Send>>>,
|
||||
killer: Arc<Mutex<Box<dyn ChildKiller + Send + Sync>>>,
|
||||
}
|
||||
|
||||
impl PtySession {
|
||||
pub fn spawn(
|
||||
size: TerminalSize,
|
||||
) -> Result<(Self, mpsc::UnboundedReceiver<Vec<u8>>), String> {
|
||||
let pty_system = native_pty_system();
|
||||
let pair = pty_system
|
||||
.openpty(size.into_pty_size())
|
||||
.map_err(|error| format!("failed to open pty: {error}"))?;
|
||||
|
||||
let shell = std::env::var("SHELL")
|
||||
.ok()
|
||||
.filter(|value| !value.trim().is_empty())
|
||||
.unwrap_or_else(|| "bash".to_owned());
|
||||
|
||||
let command = CommandBuilder::new(shell);
|
||||
let child = pair
|
||||
.slave
|
||||
.spawn_command(command)
|
||||
.map_err(|error| format!("failed to spawn shell: {error}"))?;
|
||||
drop(pair.slave);
|
||||
|
||||
let mut reader = pair
|
||||
.master
|
||||
.try_clone_reader()
|
||||
.map_err(|error| format!("failed to clone pty reader: {error}"))?;
|
||||
let writer = pair
|
||||
.master
|
||||
.take_writer()
|
||||
.map_err(|error| format!("failed to take pty writer: {error}"))?;
|
||||
let killer = child.clone_killer();
|
||||
|
||||
let (tx, rx) = mpsc::unbounded_channel();
|
||||
let child_handle = Arc::new(Mutex::new(child));
|
||||
let child_for_thread = child_handle.clone();
|
||||
|
||||
std::thread::spawn(move || {
|
||||
let mut buffer = [0_u8; 4096];
|
||||
|
||||
loop {
|
||||
match reader.read(&mut buffer) {
|
||||
Ok(0) => break,
|
||||
Ok(read) => {
|
||||
if tx.send(buffer[..read].to_vec()).is_err() {
|
||||
break;
|
||||
}
|
||||
}
|
||||
Err(error) => {
|
||||
let _ = tx.send(
|
||||
format!("\r\n[terminal read error] {error}\r\n").into_bytes(),
|
||||
);
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
let exit_code = child_for_thread
|
||||
.lock()
|
||||
.ok()
|
||||
.and_then(|mut child| child.wait().ok())
|
||||
.map(|status| status.exit_code() as i32);
|
||||
|
||||
let payload = serde_json::to_vec(&app::terminal::protocol::ServerTerminalMessage::Exit {
|
||||
code: exit_code,
|
||||
})
|
||||
.unwrap_or_else(|_| {
|
||||
br#"{"type":"exit","code":null}"#.to_vec()
|
||||
});
|
||||
let _ = tx.send(payload);
|
||||
});
|
||||
|
||||
Ok((
|
||||
Self {
|
||||
master: pair.master,
|
||||
writer: Arc::new(Mutex::new(writer)),
|
||||
killer: Arc::new(Mutex::new(killer)),
|
||||
},
|
||||
rx,
|
||||
))
|
||||
}
|
||||
|
||||
pub fn resize(&self, size: TerminalSize) -> Result<(), String> {
|
||||
self.master
|
||||
.resize(size.into_pty_size())
|
||||
.map_err(|error| format!("failed to resize pty: {error}"))
|
||||
}
|
||||
|
||||
pub fn write(&self, data: &str) -> Result<(), String> {
|
||||
let mut writer = self
|
||||
.writer
|
||||
.lock()
|
||||
.map_err(|_| "failed to lock pty writer".to_owned())?;
|
||||
writer
|
||||
.write_all(data.as_bytes())
|
||||
.map_err(|error| format!("failed to write to pty: {error}"))?;
|
||||
writer
|
||||
.flush()
|
||||
.map_err(|error| format!("failed to flush pty writer: {error}"))
|
||||
}
|
||||
|
||||
pub fn kill(&self) {
|
||||
if let Ok(mut killer) = self.killer.lock() {
|
||||
let _ = killer.kill();
|
||||
}
|
||||
}
|
||||
}
|
||||
103
server/src/terminal/ws.rs
Normal file
103
server/src/terminal/ws.rs
Normal file
@@ -0,0 +1,103 @@
|
||||
use app::terminal::protocol::{ClientTerminalMessage, ServerTerminalMessage};
|
||||
use axum::extract::ws::{Message, WebSocket, WebSocketUpgrade};
|
||||
use axum::response::Response;
|
||||
use futures_util::{SinkExt, StreamExt};
|
||||
use tokio::sync::mpsc;
|
||||
|
||||
use super::pty_session::{PtySession, TerminalSize};
|
||||
|
||||
pub async fn terminal_ws_handler(ws: WebSocketUpgrade) -> Response {
|
||||
ws.on_upgrade(handle_terminal_socket)
|
||||
}
|
||||
|
||||
async fn handle_terminal_socket(socket: WebSocket) {
|
||||
let Ok((session, mut output_rx)) = PtySession::spawn(TerminalSize::default()) else {
|
||||
return;
|
||||
};
|
||||
|
||||
let (mut sender, mut receiver) = socket.split();
|
||||
let (server_tx, mut server_rx) = mpsc::unbounded_channel::<ServerTerminalMessage>();
|
||||
|
||||
let output_forwarder = {
|
||||
let server_tx = server_tx.clone();
|
||||
tokio::spawn(async move {
|
||||
while let Some(chunk) = output_rx.recv().await {
|
||||
if let Ok(message) = serde_json::from_slice::<ServerTerminalMessage>(&chunk) {
|
||||
if server_tx.send(message).is_err() {
|
||||
break;
|
||||
}
|
||||
continue;
|
||||
}
|
||||
|
||||
let payload = String::from_utf8_lossy(&chunk).into_owned();
|
||||
if server_tx
|
||||
.send(ServerTerminalMessage::Output { data: payload })
|
||||
.is_err()
|
||||
{
|
||||
break;
|
||||
}
|
||||
}
|
||||
})
|
||||
};
|
||||
|
||||
let send_task = tokio::spawn(async move {
|
||||
while let Some(message) = server_rx.recv().await {
|
||||
let Ok(payload) = serde_json::to_string(&message) else {
|
||||
continue;
|
||||
};
|
||||
|
||||
if sender.send(Message::Text(payload.into())).await.is_err() {
|
||||
break;
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
while let Some(Ok(message)) = receiver.next().await {
|
||||
match message {
|
||||
Message::Text(text) => {
|
||||
let Ok(client_message) =
|
||||
serde_json::from_str::<ClientTerminalMessage>(text.as_str())
|
||||
else {
|
||||
let _ = server_tx.send(ServerTerminalMessage::Error {
|
||||
message: "invalid terminal message".to_owned(),
|
||||
});
|
||||
continue;
|
||||
};
|
||||
|
||||
match client_message {
|
||||
ClientTerminalMessage::Input { data } => {
|
||||
if let Err(error) = session.write(&data) {
|
||||
let _ = server_tx
|
||||
.send(ServerTerminalMessage::Error { message: error });
|
||||
}
|
||||
}
|
||||
ClientTerminalMessage::Resize {
|
||||
cols,
|
||||
rows,
|
||||
pixel_width,
|
||||
pixel_height,
|
||||
} => {
|
||||
if let Err(error) = session.resize(TerminalSize {
|
||||
cols,
|
||||
rows,
|
||||
pixel_width,
|
||||
pixel_height,
|
||||
}) {
|
||||
let _ = server_tx
|
||||
.send(ServerTerminalMessage::Error { message: error });
|
||||
}
|
||||
}
|
||||
ClientTerminalMessage::Ping => {
|
||||
let _ = server_tx.send(ServerTerminalMessage::Pong);
|
||||
}
|
||||
}
|
||||
}
|
||||
Message::Close(_) => break,
|
||||
_ => {}
|
||||
}
|
||||
}
|
||||
|
||||
session.kill();
|
||||
output_forwarder.abort();
|
||||
send_task.abort();
|
||||
}
|
||||
Reference in New Issue
Block a user