feat: implement Rtty server core (PTY + alacritty engine + WebSocket)

This commit is contained in:
2026-08-01 16:35:46 +08:00
parent 01b3bea1c6
commit f25506ec5b
9 changed files with 714 additions and 25 deletions

View File

@@ -22,10 +22,11 @@ tracing = "0.1"
tracing-subscriber = { version = "0.3", features = ["env-filter"] }
anyhow = "1.0"
dashmap = "6.0" # 高并发线程安全的 SessionMap
futures-util = "0.3"
[profile.release]
opt-level = 3
lto = true
codegen-units = 1
strip = true
strip = true

View File

@@ -0,0 +1,115 @@
//! 服务端配置模块。
//!
//! 支持通过环境变量(`RTTY_HOST`、`RTTY_PORT`、`RTTY_SHELL` 等)覆盖默认值,
//! 未来可扩展为从配置文件加载。
use std::env;
/// Rtty 服务端配置。
#[derive(Debug, Clone)]
pub struct ServerConfig {
/// 监听地址。
pub host: String,
/// 监听端口。
pub port: u16,
/// 默认启动的 Shell 程序。
pub shell: String,
/// 终端默认列数。
pub cols: u16,
/// 终端默认行数。
pub rows: u16,
/// 回滚缓冲区的最大行数。
pub max_scrollback: usize,
/// 单个会话允许的最大并发客户端数0 表示不限制)。
pub max_clients: usize,
}
impl Default for ServerConfig {
fn default() -> Self {
Self {
host: "0.0.0.0".into(),
port: 8080,
shell: default_shell(),
cols: 120,
rows: 32,
max_scrollback: 10_000,
max_clients: 16,
}
}
}
impl ServerConfig {
/// 从环境变量加载配置,未设置的项回落到默认值。
pub fn from_env() -> Self {
let mut cfg = ServerConfig::default();
if let Ok(v) = env::var("RTTY_HOST") {
cfg.host = v;
}
if let Ok(v) = env::var("RTTY_PORT")
&& let Ok(p) = v.parse()
{
cfg.port = p;
}
if let Ok(v) = env::var("RTTY_SHELL") {
cfg.shell = v;
}
if let Ok(v) = env::var("RTTY_COLS")
&& let Ok(c) = v.parse()
{
cfg.cols = c;
}
if let Ok(v) = env::var("RTTY_ROWS")
&& let Ok(r) = v.parse()
{
cfg.rows = r;
}
if let Ok(v) = env::var("RTTY_MAX_SCROLLBACK")
&& let Ok(n) = v.parse()
{
cfg.max_scrollback = n;
}
if let Ok(v) = env::var("RTTY_MAX_CLIENTS")
&& let Ok(n) = v.parse()
{
cfg.max_clients = n;
}
cfg
}
/// 返回 `host:port` 形式的监听地址。
pub fn bind_addr(&self) -> String {
format!("{}:{}", self.host, self.port)
}
}
/// 根据当前平台选择默认 Shell。
fn default_shell() -> String {
#[cfg(windows)]
{
env::var("COMSPEC").unwrap_or_else(|_| "cmd.exe".into())
}
#[cfg(not(windows))]
{
for shell in ["zsh", "bash", "sh"] {
if command_exists(shell) {
return shell.to_string();
}
}
"sh".to_string()
}
}
/// 检查某个命令是否存在于 PATH 中。
#[cfg(not(windows))]
fn command_exists(cmd: &str) -> bool {
use std::process::Command;
Command::new(cmd)
.arg("--version")
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.status()
.map(|s| s.success())
.unwrap_or(false)
}

View File

@@ -2,33 +2,51 @@ mod config;
mod terminal;
mod ws;
use axum::{routing::get, Router};
use std::sync::Arc;
use axum::{routing::get, Router};
use dashmap::DashMap;
use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt};
use config::ServerConfig;
use ws::handler::Session;
/// 全局应用状态。
pub struct AppState {
// 预留全局 Session 管理句柄
/// 会话表session_id -> Session
pub sessions: DashMap<String, Arc<Session>>,
/// 服务端配置。
pub config: Arc<ServerConfig>,
}
#[tokio::main]
async fn main() -> anyhow::Result<()> {
// 初始化日志
// 初始化日志
tracing_subscriber::registry()
.with(tracing_subscriber::EnvFilter::try_from_default_env().unwrap_or_else(|_| "rtty_server=debug".into()))
.with(
tracing_subscriber::EnvFilter::try_from_default_env()
.unwrap_or_else(|_| "rtty_server=debug,tower_http=info".into()),
)
.with(tracing_subscriber::fmt::layer())
.init();
let state = Arc::new(AppState {});
let config = Arc::new(ServerConfig::from_env());
let state = Arc::new(AppState {
sessions: DashMap::new(),
config: config.clone(),
});
let app = Router::new()
.route("/ws", get(ws::handler::ws_route))
.with_state(state);
let addr = "0.0.0.0:8080";
let addr = config.bind_addr();
tracing::info!("🚀 Rtty Server running on ws://{}", addr);
tracing::info!(" shell: {}", config.shell);
tracing::info!(" default size: {}x{}", config.cols, config.rows);
let listener = tokio::net::TcpListener::bind(addr).await?;
let listener = tokio::net::TcpListener::bind(&addr).await?;
axum::serve(listener, app).await?;
Ok(())
}
}

View File

@@ -0,0 +1,144 @@
//! 终端状态引擎。
//!
//! 封装 [`alacritty_terminal::Term`] 作为整个系统的“真相源”,负责:
//! - 将 PTY 读到的原始字节流喂给 ANSI 解析器,维护 2D 屏幕网格与滚动历史;
//! - 把终端回写给 Shell 的数据(如光标位置上报、标题查询响应)转发到 PTY
//! - 为移动端生成语义化 JSON 快照,为断线重连提供状态恢复。
use std::io::Write;
use std::sync::{Arc, Mutex};
use alacritty_terminal::event::{Event, EventListener};
use alacritty_terminal::grid::Dimensions;
use alacritty_terminal::index::{Column, Point};
use alacritty_terminal::term::cell::Flags;
use alacritty_terminal::term::{point_to_viewport, viewport_to_point, Config, Term};
use alacritty_terminal::vte::ansi::Processor;
use crate::ws::protocol::MobileSnapshot;
/// 终端尺寸,实现 alacritty 的 [`Dimensions`]。
#[derive(Debug, Clone, Copy)]
pub struct TermSize {
pub columns: usize,
pub rows: usize,
}
impl Dimensions for TermSize {
fn total_lines(&self) -> usize {
self.rows
}
fn screen_lines(&self) -> usize {
self.rows
}
fn columns(&self) -> usize {
self.columns
}
}
/// 事件监听器:把 alacritty 要求回写给 Shell 的数据写入 PTY。
#[derive(Clone)]
pub struct SessionListener {
writer: Arc<Mutex<Box<dyn Write + Send>>>,
}
impl SessionListener {
pub fn new(writer: Arc<Mutex<Box<dyn Write + Send>>>) -> Self {
Self { writer }
}
}
impl EventListener for SessionListener {
fn send_event(&self, event: Event) {
if let Event::PtyWrite(text) = event {
let mut w = match self.writer.lock() {
Ok(w) => w,
Err(_) => return,
};
let _ = w.write_all(text.as_bytes());
let _ = w.flush();
}
}
}
/// 终端引擎:持有 `Term` 与 ANSI 解析器。
pub struct TerminalEngine {
term: Term<SessionListener>,
parser: Processor,
}
impl TerminalEngine {
/// 创建一个新的终端引擎。
pub fn new(size: TermSize, writer: Arc<Mutex<Box<dyn Write + Send>>>) -> Self {
let listener = SessionListener::new(writer);
let term = Term::new(Config::default(), &size, listener);
Self { term, parser: Processor::new() }
}
/// 将 PTY 读到的字节流喂给解析器。
pub fn feed(&mut self, bytes: &[u8]) {
for &byte in bytes {
self.parser.advance(&mut self.term, byte);
}
}
/// 调整终端网格尺寸。
pub fn resize(&mut self, cols: u16, rows: u16) {
let size = TermSize { columns: cols as usize, rows: rows as usize };
self.term.resize(size);
}
/// 生成移动端语义化快照。
pub fn snapshot(&self) -> MobileSnapshot {
let grid = self.term.grid();
let columns = grid.columns();
let rows = grid.screen_lines();
let display_offset = grid.display_offset();
let cursor_point = grid.cursor.point;
let cursor = point_to_viewport(display_offset, cursor_point);
let mut lines = Vec::with_capacity(rows);
for line_idx in 0..rows {
// 将视口内逻辑行号转换为网格行号。
let grid_point = viewport_to_point(display_offset, Point::new(line_idx, Column(0)));
let row = &grid[grid_point.line];
let mut text = String::new();
for cell in row {
// 跳过全角字符的占位空格,避免语义化文本中出现多余空白。
if cell.flags.contains(Flags::WIDE_CHAR_SPACER) {
continue;
}
// 跳过空格的连续尾部会在上层处理;这里仍收集可见字符。
if cell.c != ' ' || !text.is_empty() {
text.push(cell.c);
}
}
lines.push(text.trim_end().to_string());
}
MobileSnapshot {
cursor_x: cursor.map(|c| c.column.0).unwrap_or(0),
cursor_y: cursor.map(|c| c.line).unwrap_or(0),
cols: columns,
rows,
display_offset,
scrollback_lines: grid.history_size(),
lines,
}
}
/// 当前网格列数。
pub fn columns(&self) -> usize {
self.term.grid().columns()
}
/// 当前网格行数。
pub fn rows(&self) -> usize {
self.term.grid().screen_lines()
}
}

View File

@@ -1,2 +1,2 @@
mod pty;
mod engine;
pub mod engine;
pub mod pty;

View File

@@ -0,0 +1,98 @@
//! PTY 会话封装。
//!
//! 基于 [`portable_pty`] 创建伪终端并拉起子进程Shell提供读写与 resize 能力。
//! 读取端由服务端读取任务独占持有,写入端与 master 句柄可被多个客户端共享。
use std::io::{Read, Write};
use std::sync::{Arc, Mutex};
use anyhow::{Context, Result};
use portable_pty::{Child, MasterPty, PtyPair, PtySize};
/// 一个已启动的 PTY 会话。
pub struct PtySession {
/// Master 端句柄,用于 resize / 获取大小。
master: Box<dyn MasterPty + Send>,
/// 从 Slave 端读取输出的流。独占,由读取任务持有。
reader: Option<Box<dyn Read + Send>>,
/// 写入 Slave 端的流。可被多个客户端共享(加锁)。
writer: Arc<Mutex<Box<dyn Write + Send>>>,
/// 子进程句柄,用于检测退出 / 终止。
child: Box<dyn Child + Send + Sync>,
/// 当前尺寸。
size: PtySize,
}
impl PtySession {
/// 以指定 Shell 与初始尺寸创建一个 PTY 会话。
pub fn new(shell: &str, cols: u16, rows: u16) -> Result<Self> {
let size = PtySize { rows, cols, pixel_width: 0, pixel_height: 0 };
let pty_system = portable_pty::native_pty_system();
let pair = pty_system.openpty(size).context("failed to open pty")?;
let child = spawn_child(&pair, shell).context("failed to spawn shell")?;
let reader = pair.master.try_clone_reader().context("failed to clone pty reader")?;
let writer = pair.master.take_writer().context("failed to take pty writer")?;
Ok(Self {
master: pair.master,
reader: Some(reader),
writer: Arc::new(Mutex::new(writer)),
child,
size,
})
}
/// 取出读取端,供 PTY 读取任务独占使用。取走后不可再次调用。
pub fn take_reader(&mut self) -> Option<Box<dyn Read + Send>> {
self.reader.take()
}
/// 获取共享写入端。
pub fn writer(&self) -> &Arc<Mutex<Box<dyn Write + Send>>> {
&self.writer
}
/// 将数据写入 Slave 端(发送给 Shell
pub fn write(&self, data: &[u8]) -> Result<()> {
let mut w = self
.writer
.lock()
.map_err(|_| anyhow::anyhow!("pty writer poisoned"))?;
w.write_all(data)?;
w.flush()?;
Ok(())
}
/// 调整 PTY 尺寸(通知内核与子进程)。
pub fn resize(&mut self, cols: u16, rows: u16) -> Result<()> {
self.size.rows = rows;
self.size.cols = cols;
self.master.resize(self.size).context("failed to resize pty")
}
/// 当前 PTY 尺寸。
pub fn size(&self) -> (u16, u16) {
(self.size.cols, self.size.rows)
}
/// 子进程是否已退出。
pub fn try_wait(&mut self) -> Result<Option<portable_pty::ExitStatus>> {
self.child.try_wait().context("failed to poll child")
}
/// 终止子进程。
pub fn kill(&mut self) -> Result<()> {
self.child.kill().context("failed to kill child")
}
}
/// 拉起子进程,并把 Stdio 重定向到 PTY 的 Slave 端。
fn spawn_child(pair: &PtyPair, shell: &str) -> Result<Box<dyn Child + Send + Sync>> {
let mut cmd = portable_pty::CommandBuilder::new(shell);
// 让 Shell 以交互方式运行。
cmd.env("TERM", "xterm-256color");
let child = pair.slave.spawn_command(cmd)?;
Ok(child)
}

View File

@@ -0,0 +1,263 @@
//! WebSocket 处理器与会话管理。
//!
//! 每个会话对应一个 PTY + 一个 alacritty 终端引擎。PTY 读取任务将输出广播给所有
//! 订阅者PC 端收原始 ANSI 二进制帧,移动端收 JSON 快照)。客户端消息写入 PTY。
use std::collections::HashMap;
use std::io::Read;
use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use anyhow::Result;
use axum::extract::ws::{Message, WebSocket, WebSocketUpgrade};
use axum::extract::{Query, State};
use axum::response::IntoResponse;
use futures_util::{SinkExt, StreamExt};
use tokio::sync::broadcast;
use crate::config::ServerConfig;
use crate::terminal::engine::{TermSize, TerminalEngine};
use crate::terminal::pty::PtySession;
use crate::ws::protocol::{ClientMessage, MobileSnapshot, ServerMessage};
/// 会话全局 ID 计数器。
static SESSION_SEQ: AtomicU64 = AtomicU64::new(0);
/// 客户端全局 ID 计数器。
static CLIENT_SEQ: AtomicU64 = AtomicU64::new(0);
/// 广播给订阅者的输出事件。
#[derive(Clone)]
pub enum OutputEvent {
/// 原始 ANSI 字节流PC 端渲染)。
Raw(Vec<u8>),
/// 语义化快照(移动端渲染)。
Snapshot(MobileSnapshot),
/// PTY 已关闭。
Exit,
}
/// 一个终端会话。
pub struct Session {
pub id: String,
/// PTY 会话(写入 / resize / kill
pub pty: Arc<Mutex<PtySession>>,
/// alacritty 终端引擎。
pub engine: Arc<Mutex<TerminalEngine>>,
/// 输出广播通道。
pub output: broadcast::Sender<OutputEvent>,
/// 当前持有控制权的客户端 ID。
pub control: Arc<Mutex<Option<String>>>,
/// 当前连接的客户端数。
pub clients: AtomicUsize,
}
/// WebSocket 路由入口。
pub async fn ws_route(
ws: WebSocketUpgrade,
State(state): State<Arc<crate::AppState>>,
Query(params): Query<HashMap<String, String>>,
) -> impl IntoResponse {
ws.on_upgrade(move |socket| handle_socket(socket, state, params))
}
/// 建立连接并处理双向消息。
async fn handle_socket(
socket: WebSocket,
state: Arc<crate::AppState>,
params: HashMap<String, String>,
) {
let client_id = format!("client-{}", CLIENT_SEQ.fetch_add(1, Ordering::Relaxed));
let requested = params.get("session").cloned().unwrap_or_default();
let session = match get_or_create_session(&state, &requested) {
Ok(s) => s,
Err(e) => {
let msg = ServerMessage::Error { message: format!("session error: {e}") };
let _ = send_text(socket, msg.to_json()).await;
return;
}
};
let (mut tx, mut rx) = socket.split();
// 发送就绪消息。
let (cols, rows) = session.pty.lock().unwrap().size();
let ready = ServerMessage::Ready { id: session.id.clone(), cols, rows };
if tx.send(Message::Text(ready.to_json())).await.is_err() {
return;
}
// 发送当前快照(断线重连 / 新加入状态恢复)。
let snap = session.engine.lock().unwrap().snapshot();
let init = ServerMessage::MobileSnapshot { data: snap };
if tx.send(Message::Text(init.to_json())).await.is_err() {
return;
}
// 订阅输出。
let mut out_rx = session.output.subscribe();
session.clients.fetch_add(1, Ordering::SeqCst);
loop {
tokio::select! {
// 服务端输出 -> 客户端
out = out_rx.recv() => {
match out {
Ok(OutputEvent::Raw(bytes)) => {
if tx.send(Message::Binary(bytes)).await.is_err() {
break;
}
}
Ok(OutputEvent::Snapshot(s)) => {
let msg = ServerMessage::MobileSnapshot { data: s };
if tx.send(Message::Text(msg.to_json())).await.is_err() {
break;
}
}
Ok(OutputEvent::Exit) => {
let _ = tx.send(Message::Text(ServerMessage::SessionClosed.to_json())).await;
break;
}
Err(_) => break,
}
}
// 客户端输入 -> 服务端
msg = rx.next() => {
match msg {
Some(Ok(Message::Text(text))) => {
if let Some(reply) = handle_client_message(&session, &client_id, &text)
&& tx.send(Message::Text(reply)).await.is_err()
{
break;
}
}
Some(Ok(Message::Binary(bytes))) => {
// 兼容:客户端也可能用二进制发送输入。
let _ = session.pty.lock().unwrap().write(&bytes);
}
Some(Ok(_)) => {}
Some(Err(_)) | None => break,
}
}
}
}
session.clients.fetch_sub(1, Ordering::SeqCst);
}
/// 处理单条客户端 JSON 消息,返回需要回发给该客户端的 JSON若无则 None
fn handle_client_message(
session: &Session,
client_id: &str,
text: &str,
) -> Option<String> {
let msg: ClientMessage = match serde_json::from_str(text) {
Ok(m) => m,
Err(_) => return None,
};
match msg {
ClientMessage::Input { data } => {
let _ = session.pty.lock().unwrap().write(data.as_bytes());
None
}
ClientMessage::Resize { cols, rows } => {
// 先调整 PTY 内核尺寸,再调整终端网格。
let mut pty = session.pty.lock().unwrap();
let _ = pty.resize(cols, rows);
drop(pty);
session.engine.lock().unwrap().resize(cols, rows);
None
}
ClientMessage::ClaimControl => {
let mut control = session.control.lock().unwrap();
let granted = control.is_none();
if granted {
*control = Some(client_id.to_string());
}
let holder = control.clone();
drop(control);
Some(ServerMessage::ControlResponse { granted, holder }.to_json())
}
ClientMessage::Ping => None,
}
}
/// 获取或创建会话。
fn get_or_create_session(
state: &crate::AppState,
requested: &str,
) -> Result<Arc<Session>> {
if !requested.is_empty()
&& let Some(s) = state.sessions.get(requested)
{
return Ok(s.clone());
}
let id = format!("session-{}", SESSION_SEQ.fetch_add(1, Ordering::Relaxed));
let session = create_session(id.clone(), &state.config)?;
state.sessions.insert(id.clone(), session.clone());
Ok(session)
}
/// 创建会话并启动 PTY 读取任务。
fn create_session(id: String, config: &ServerConfig) -> Result<Arc<Session>> {
let mut pty = PtySession::new(&config.shell, config.cols, config.rows)?;
let writer = pty.writer().clone();
let engine = Arc::new(Mutex::new(TerminalEngine::new(
TermSize { columns: config.cols as usize, rows: config.rows as usize },
writer,
)));
let reader = pty.take_reader().ok_or_else(|| anyhow::anyhow!("pty reader already taken"))?;
let (output, _) = broadcast::channel(256);
let session = Arc::new(Session {
id,
pty: Arc::new(Mutex::new(pty)),
engine,
output,
control: Arc::new(Mutex::new(None)),
clients: AtomicUsize::new(0),
});
start_pty_reader(session.clone(), reader);
Ok(session)
}
/// 启动阻塞的 PTY 读取任务,读取输出并广播。
fn start_pty_reader(session: Arc<Session>, mut reader: Box<dyn Read + Send>) {
tokio::task::spawn_blocking(move || {
let mut buf = vec![0u8; 8192];
loop {
let n = match reader.read(&mut buf) {
Ok(0) => break,
Ok(n) => n,
Err(_) => break,
};
let chunk = &buf[..n];
// 喂给解析器并生成快照。
let snapshot = {
let mut eng = match session.engine.lock() {
Ok(e) => e,
Err(_) => break,
};
eng.feed(chunk);
eng.snapshot()
};
// 广播原始字节 + 快照。
let _ = session.output.send(OutputEvent::Raw(chunk.to_vec()));
let _ = session.output.send(OutputEvent::Snapshot(snapshot));
}
let _ = session.output.send(OutputEvent::Exit);
});
}
/// 发送一条文本消息,忽略错误。
async fn send_text(socket: WebSocket, text: String) -> Result<()> {
let (mut tx, _rx) = socket.split();
tx.send(Message::Text(text)).await?;
Ok(())
}

View File

@@ -1,2 +1,2 @@
pub(crate) mod handler;
mod protocol;
pub mod protocol;

View File

@@ -1,27 +1,77 @@
//! WebSocket 通信协议。
//!
//! 服务端与客户端之间通过 WebSocket 交换消息:
//! - 客户端发送 [`ClientMessage`](输入、调整尺寸、声明控制权);
//! - 服务端发送 [`ServerMessage`](就绪、移动端快照、控制权授予、错误、会话结束)。
//!
//! 另外,为 PC 端提供**二进制**通道PTY 输出的原始 ANSI 字节流直接以 WebSocket
//! 二进制帧下发,供 `flutter_alacritty` 等渲染引擎消费,保证 100% 工业级兼容。
use serde::{Deserialize, Serialize};
/// 客户端发送给 Rtty 服务端的控制指令
/// 客户端发送给 Rtty 服务端的控制指令
#[derive(Debug, Deserialize)]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum ClientMessage {
/// 键盘/文本输入
/// 键盘/文本输入,原样写入 PTY。
Input { data: String },
/// 客户端请求 Resize
/// 客户端请求调整终端尺寸。
Resize { cols: u16, rows: u16 },
/// 客户端申明控制权 (用来解决多端控制冲突)
/// 客户端申明控制权用来解决多端控制冲突)。
ClaimControl,
/// 客户端心跳,保持连接活跃。
Ping,
}
/// 服务端推送给客户端的消息 (分别适配 PC 和 移动端)
/// 移动端语义化屏显数据(已解耦、去除 ANSI 序列)。
#[derive(Debug, Clone, Serialize)]
pub struct MobileSnapshot {
/// 光标列(相对视口)。
pub cursor_x: usize,
/// 光标行(相对视口)。
pub cursor_y: usize,
/// 网格列数。
pub cols: usize,
/// 网格行数。
pub rows: usize,
/// 当前回滚显示偏移。
pub display_offset: usize,
/// 滚动历史中的总行数。
pub scrollback_lines: usize,
/// 视口内的逐行文本(已按词、去尾空白)。
pub lines: Vec<String>,
}
/// 服务端推送给客户端的消息JSON 文本帧)。
#[derive(Debug, Serialize)]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum ServerMessage {
/// 原生 ANSI 字节流 (供 PC 端 xterm/flutter_alacritty 渲染)
RawData { data: String },
/// 移动端专属:解耦后的语义化 JSON 屏显数据
MobileSnapshot {
cursor_x: usize,
cursor_y: usize,
lines: Vec<String>,
/// 会话已就绪:携带会话 ID 与初始尺寸。
Ready {
id: String,
cols: u16,
rows: u16,
},
}
/// 移动端专属:解耦后的语义化屏显快照。
MobileSnapshot {
data: MobileSnapshot,
},
/// 控制权授予结果。
ControlResponse {
granted: bool,
holder: Option<String>,
},
/// 服务端错误。
Error {
message: String,
},
/// 会话已结束PTY 关闭)。
SessionClosed,
}
impl ServerMessage {
/// 序列化为 JSON 文本。
pub fn to_json(&self) -> String {
serde_json::to_string(self).unwrap_or_else(|_| "{}".to_string())
}
}