fix: address unimplemented/broken issues across server and desktop

This commit is contained in:
2026-08-01 18:58:12 +08:00
parent 93ee801f27
commit 7e4236e2f3
7 changed files with 239 additions and 24 deletions

View File

@@ -320,12 +320,12 @@ class _StatusBar extends StatelessWidget {
), ),
), ),
SizedBox(width: 20), SizedBox(width: 20),
_Hint('鼠标拖选复制', RttyTheme.textFaint), _Hint('拖选复制', RttyTheme.textFaint),
SizedBox(width: 14), SizedBox(width: 14),
_Hint('右键粘贴', RttyTheme.textFaint), _Hint('Ctrl+Shift+C / V 复制粘贴', RttyTheme.textFaint),
Spacer(), Spacer(),
Text( Text(
'FLUTTER × XTERM × ALCATTERM ENGINE', 'FLUTTER × XTERM TERMINAL',
style: TextStyle(color: RttyTheme.textFaint, fontSize: 10), style: TextStyle(color: RttyTheme.textFaint, fontSize: 10),
), ),
], ],

View File

@@ -2,18 +2,22 @@
// //
// 运行前需先启动 Rust 服务端cd .. && cargo run // 运行前需先启动 Rust 服务端cd .. && cargo run
// 运行flutter test test/e2e_ws_test.dart // 运行flutter test test/e2e_ws_test.dart
import 'dart:io';
import 'package:flutter_test/flutter_test.dart'; import 'package:flutter_test/flutter_test.dart';
import 'package:xterm/xterm.dart'; import 'package:xterm/xterm.dart';
import '../lib/src/rtty_client.dart'; import 'package:rtty_desktop/src/rtty_client.dart';
import 'support.dart';
void main() { void main() {
TestWidgetsFlutterBinding.ensureInitialized(); TestWidgetsFlutterBinding.ensureInitialized();
test('RttyClient connects to Rust server and renders command output', test('RttyClient connects to Rust server and renders command output',
() async { () async {
if (!await isServerUp('127.0.0.1', 8080)) {
markTestSkipped('Rtty server not running');
return;
}
final terminal = Terminal(); final terminal = Terminal();
final client = RttyClient(terminal); final client = RttyClient(terminal);

View File

@@ -0,0 +1,77 @@
// 多端解耦 + 控制权强制回归验证(连真实 Rust 服务端)。
import 'dart:convert';
import 'package:flutter_test/flutter_test.dart';
import 'package:web_socket_channel/io.dart';
import 'support.dart';
void main() {
TestWidgetsFlutterBinding.ensureInitialized();
test('多端解耦 + 控制权强制 + 会话清理', () async {
if (!await isServerUp('127.0.0.1', 8080)) {
markTestSkipped('Rtty server not running');
return;
}
final base = 'ws://127.0.0.1:8080/ws';
// 1. 桌面端client=desktop收 raw 二进制,不收 snapshot。
final desktop = IOWebSocketChannel.connect(Uri.parse('$base?client=desktop'));
final dOut = <String>[];
desktop.stream.listen((m) {
if (m is List<int>) {
dOut.add('binary:${utf8.decode(m, allowMalformed: true)}');
} else if (m is String) {
dOut.add('text:${m.length > 100 ? m.substring(0, 100) : m}');
} else {
dOut.add('$m');
}
});
await Future<void>.delayed(const Duration(milliseconds: 700));
final dHasBinary = dOut.any((e) => e.startsWith('binary:'));
final dHasSnapshot = dOut.any((e) => e.contains('mobile_snapshot'));
expect(dHasBinary, isTrue, reason: '桌面端应收到原始 ANSI 二进制 $dOut');
expect(dHasSnapshot, isFalse, reason: '桌面端不应收到语义快照 $dOut');
// 2. 移动端client=mobile收 snapshot不收 binary。
final mobile = IOWebSocketChannel.connect(Uri.parse('$base?client=mobile'));
final mOut = <String>[];
mobile.stream.listen((m) {
if (m is List<int>) {
mOut.add('binary:${utf8.decode(m, allowMalformed: true)}');
} else if (m is String) {
mOut.add('text:${m.length > 100 ? m.substring(0, 100) : m}');
} else {
mOut.add('$m');
}
});
await Future<void>.delayed(const Duration(milliseconds: 700));
final mHasSnapshot = mOut.any((e) => e.contains('mobile_snapshot'));
final mHasBinary = mOut.any((e) => e.startsWith('binary:'));
expect(mHasSnapshot, isTrue, reason: '移动端应收到语义快照 $mOut');
expect(mHasBinary, isFalse, reason: '移动端不应收到二进制 $mOut');
// 3. 桌面端发命令,确认 raw 输出回显。
desktop.sink.add(jsonEncode({'type': 'input', 'data': 'echo D_E2E\r\n'}));
await Future<void>.delayed(const Duration(milliseconds: 900));
expect(dOut.join('\n').contains('D_E2E'), isTrue, reason: '桌面端命令回显');
// 4. 控制权强制:桌面端 claim 后,移动端输入被拒。
desktop.sink.add(jsonEncode({'type': 'claim_control'}));
await Future<void>.delayed(const Duration(milliseconds: 400));
mobile.sink
.add(jsonEncode({'type': 'input', 'data': 'echo BLOCKED_SHOULD_NOT_SHOW\r\n'}));
await Future<void>.delayed(const Duration(milliseconds: 900));
final both = (dOut.join('\n') + mOut.join('\n'));
expect(both.contains('BLOCKED_SHOULD_NOT_SHOW'), isFalse,
reason: '非控制者输入应被阻止');
// 5. 会话清理exit 触发会话移除。
desktop.sink.add(jsonEncode({'type': 'input', 'data': 'exit\r\n'}));
await Future<void>.delayed(const Duration(milliseconds: 1300));
desktop.sink.close();
mobile.sink.close();
});
}

14
desktop/test/support.dart Normal file
View File

@@ -0,0 +1,14 @@
// 集成测试共享工具。
import 'package:web_socket_channel/io.dart';
/// 探测 Rtty 服务端是否可达;返回 false 时集成测试应被跳过。
Future<bool> isServerUp(String host, int port) async {
try {
final ws = IOWebSocketChannel.connect(Uri.parse('ws://$host:$port/ws'));
await ws.ready.timeout(const Duration(seconds: 2));
ws.sink.close();
return true;
} catch (_) {
return false;
}
}

View File

@@ -22,6 +22,8 @@ pub struct ServerConfig {
pub max_scrollback: usize, pub max_scrollback: usize,
/// 单个会话允许的最大并发客户端数0 表示不限制)。 /// 单个会话允许的最大并发客户端数0 表示不限制)。
pub max_clients: usize, pub max_clients: usize,
/// 会话在无客户端连接后保留的秒数;超时则清理(支持断线重连窗口)。
pub idle_timeout_secs: u64,
} }
impl Default for ServerConfig { impl Default for ServerConfig {
@@ -34,6 +36,7 @@ impl Default for ServerConfig {
rows: 32, rows: 32,
max_scrollback: 10_000, max_scrollback: 10_000,
max_clients: 16, max_clients: 16,
idle_timeout_secs: 60,
} }
} }
} }
@@ -74,6 +77,11 @@ impl ServerConfig {
{ {
cfg.max_clients = n; cfg.max_clients = n;
} }
if let Ok(v) = env::var("RTTY_IDLE_TIMEOUT")
&& let Ok(n) = v.parse()
{
cfg.idle_timeout_secs = n;
}
cfg cfg
} }

View File

@@ -71,9 +71,16 @@ pub struct TerminalEngine {
impl TerminalEngine { impl TerminalEngine {
/// 创建一个新的终端引擎。 /// 创建一个新的终端引擎。
pub fn new(size: TermSize, writer: Arc<Mutex<Box<dyn Write + Send>>>) -> Self { ///
/// `scrollback` 为滚动历史的行数上限,对应 alacritty 的 `scrolling_history`。
pub fn new(
size: TermSize,
writer: Arc<Mutex<Box<dyn Write + Send>>>,
scrollback: usize,
) -> Self {
let listener = SessionListener::new(writer); let listener = SessionListener::new(writer);
let term = Term::new(Config::default(), &size, listener); let config = Config { scrolling_history: scrollback, ..Config::default() };
let term = Term::new(config, &size, listener);
Self { term, parser: Processor::new() } Self { term, parser: Processor::new() }
} }

View File

@@ -25,6 +25,25 @@ static SESSION_SEQ: AtomicU64 = AtomicU64::new(0);
/// 客户端全局 ID 计数器。 /// 客户端全局 ID 计数器。
static CLIENT_SEQ: AtomicU64 = AtomicU64::new(0); static CLIENT_SEQ: AtomicU64 = AtomicU64::new(0);
/// 客户端类型:决定它订阅哪一路输出流(多端解耦)。
#[derive(Clone, Copy, PartialEq, Eq)]
enum ClientKind {
/// 桌面端:订阅原始 ANSI 二进制流。
Desktop,
/// 移动端:订阅语义化 JSON 快照。
Mobile,
}
impl ClientKind {
/// 从 WebSocket 查询参数解析客户端类型,默认桌面端。
fn from_params(params: &HashMap<String, String>) -> Self {
match params.get("client").map(|s| s.as_str()) {
Some("mobile") => Self::Mobile,
_ => Self::Desktop,
}
}
}
/// 广播给订阅者的输出事件。 /// 广播给订阅者的输出事件。
#[derive(Clone)] #[derive(Clone)]
pub enum OutputEvent { pub enum OutputEvent {
@@ -67,6 +86,7 @@ async fn handle_socket(
params: HashMap<String, String>, params: HashMap<String, String>,
) { ) {
let client_id = format!("client-{}", CLIENT_SEQ.fetch_add(1, Ordering::Relaxed)); let client_id = format!("client-{}", CLIENT_SEQ.fetch_add(1, Ordering::Relaxed));
let kind = ClientKind::from_params(&params);
let requested = params.get("session").cloned().unwrap_or_default(); let requested = params.get("session").cloned().unwrap_or_default();
let session = match get_or_create_session(&state, &requested) { let session = match get_or_create_session(&state, &requested) {
@@ -87,10 +107,22 @@ async fn handle_socket(
return; return;
} }
// 发送当前快照(断线重连 / 新加入状态恢复)。 // 移动端在连接时获取一次语义化快照作为初始状态(桌面端走原始 ANSI 流)。
let snap = session.engine.lock().unwrap().snapshot(); if kind == ClientKind::Mobile {
let init = ServerMessage::MobileSnapshot { data: snap }; let snap = session.engine.lock().unwrap().snapshot();
if tx.send(Message::Text(init.to_json())).await.is_err() { let init = ServerMessage::MobileSnapshot { data: snap };
if tx.send(Message::Text(init.to_json())).await.is_err() {
return;
}
}
// 订阅前检查并发上限。
let max_clients = state.config.max_clients;
if max_clients > 0 && session.clients.load(Ordering::SeqCst) >= max_clients {
let msg = ServerMessage::Error {
message: format!("session {} is at max client capacity", session.id),
};
let _ = tx.send(Message::Text(msg.to_json())).await;
return; return;
} }
@@ -100,18 +132,22 @@ async fn handle_socket(
loop { loop {
tokio::select! { tokio::select! {
// 服务端输出 -> 客户端 // 服务端输出 -> 客户端(按客户端类型过滤,实现多端解耦)。
out = out_rx.recv() => { out = out_rx.recv() => {
match out { match out {
Ok(OutputEvent::Raw(bytes)) => { Ok(OutputEvent::Raw(bytes)) => {
if tx.send(Message::Binary(bytes)).await.is_err() { if kind == ClientKind::Desktop
&& tx.send(Message::Binary(bytes)).await.is_err()
{
break; break;
} }
} }
Ok(OutputEvent::Snapshot(s)) => { Ok(OutputEvent::Snapshot(s)) => {
let msg = ServerMessage::MobileSnapshot { data: s }; if kind == ClientKind::Mobile {
if tx.send(Message::Text(msg.to_json())).await.is_err() { let msg = ServerMessage::MobileSnapshot { data: s };
break; if tx.send(Message::Text(msg.to_json())).await.is_err() {
break;
}
} }
} }
Ok(OutputEvent::Exit) => { Ok(OutputEvent::Exit) => {
@@ -143,6 +179,23 @@ async fn handle_socket(
} }
session.clients.fetch_sub(1, Ordering::SeqCst); session.clients.fetch_sub(1, Ordering::SeqCst);
// 若已无客户端安排空闲超时清理兜底机制Windows ConPTY 下进程退出
// 检测不可靠,依赖超时确保会话最终被释放)。
if session.clients.load(Ordering::SeqCst) == 0 {
spawn_idle_cleanup(session, state);
}
}
/// 空闲超时清理:若无客户端重连,则终止会话并释放资源。
fn spawn_idle_cleanup(session: Arc<Session>, state: Arc<crate::AppState>) {
let timeout = state.config.idle_timeout_secs;
tokio::spawn(async move {
tokio::time::sleep(std::time::Duration::from_secs(timeout)).await;
if session.clients.load(Ordering::SeqCst) == 0 {
cleanup_session(&session, &state);
}
});
} }
/// 处理单条客户端 JSON 消息,返回需要回发给该客户端的 JSON若无则 None /// 处理单条客户端 JSON 消息,返回需要回发给该客户端的 JSON若无则 None
@@ -158,7 +211,15 @@ fn handle_client_message(
match msg { match msg {
ClientMessage::Input { data } => { ClientMessage::Input { data } => {
let _ = session.pty.lock().unwrap().write(data.as_bytes()); // 控制权强制:若存在控制者且不是当前客户端,则拒绝写入,
// 避免多端同时输入互相干扰。
let blocked = {
let control = session.control.lock().unwrap();
control.as_deref().is_some_and(|h| h != client_id)
};
if !blocked {
let _ = session.pty.lock().unwrap().write(data.as_bytes());
}
None None
} }
ClientMessage::Resize { cols, rows } => { ClientMessage::Resize { cols, rows } => {
@@ -185,7 +246,7 @@ fn handle_client_message(
/// 获取或创建会话。 /// 获取或创建会话。
fn get_or_create_session( fn get_or_create_session(
state: &crate::AppState, state: &Arc<crate::AppState>,
requested: &str, requested: &str,
) -> Result<Arc<Session>> { ) -> Result<Arc<Session>> {
if !requested.is_empty() if !requested.is_empty()
@@ -195,18 +256,23 @@ fn get_or_create_session(
} }
let id = format!("session-{}", SESSION_SEQ.fetch_add(1, Ordering::Relaxed)); let id = format!("session-{}", SESSION_SEQ.fetch_add(1, Ordering::Relaxed));
let session = create_session(id.clone(), &state.config)?; let session = create_session(id.clone(), &state.config, state.clone())?;
state.sessions.insert(id.clone(), session.clone()); state.sessions.insert(id.clone(), session.clone());
Ok(session) Ok(session)
} }
/// 创建会话并启动 PTY 读取任务。 /// 创建会话并启动 PTY 读取任务。
fn create_session(id: String, config: &ServerConfig) -> Result<Arc<Session>> { fn create_session(
id: String,
config: &ServerConfig,
state: Arc<crate::AppState>,
) -> Result<Arc<Session>> {
let mut pty = PtySession::new(&config.shell, config.cols, config.rows)?; let mut pty = PtySession::new(&config.shell, config.cols, config.rows)?;
let writer = pty.writer().clone(); let writer = pty.writer().clone();
let engine = Arc::new(Mutex::new(TerminalEngine::new( let engine = Arc::new(Mutex::new(TerminalEngine::new(
TermSize { columns: config.cols as usize, rows: config.rows as usize }, TermSize { columns: config.cols as usize, rows: config.rows as usize },
writer, writer,
config.max_scrollback,
))); )));
let reader = pty.take_reader().ok_or_else(|| anyhow::anyhow!("pty reader already taken"))?; let reader = pty.take_reader().ok_or_else(|| anyhow::anyhow!("pty reader already taken"))?;
@@ -220,12 +286,22 @@ fn create_session(id: String, config: &ServerConfig) -> Result<Arc<Session>> {
clients: AtomicUsize::new(0), clients: AtomicUsize::new(0),
}); });
start_pty_reader(session.clone(), reader); start_pty_reader(session.clone(), reader, state);
Ok(session) Ok(session)
} }
/// 启动阻塞的 PTY 读取任务,读取输出并广播。 /// 启动阻塞的 PTY 读取任务,读取输出并广播。
fn start_pty_reader(session: Arc<Session>, mut reader: Box<dyn Read + Send>) { ///
/// 同时启动一个进程监控任务,通过轮询 `try_wait` 检测子进程退出——
/// 因为 Windows ConPTY 下 shell 退出可能不触发 PTY EOF仅靠 EOF 无法可靠清理。
fn start_pty_reader(
session: Arc<Session>,
mut reader: Box<dyn Read + Send>,
state: Arc<crate::AppState>,
) {
// 进程退出监控:一旦子进程退出即清理会话。
spawn_exit_monitor(session.clone(), state.clone());
tokio::task::spawn_blocking(move || { tokio::task::spawn_blocking(move || {
let mut buf = vec![0u8; 8192]; let mut buf = vec![0u8; 8192];
loop { loop {
@@ -251,10 +327,39 @@ fn start_pty_reader(session: Arc<Session>, mut reader: Box<dyn Read + Send>) {
let _ = session.output.send(OutputEvent::Snapshot(snapshot)); let _ = session.output.send(OutputEvent::Snapshot(snapshot));
} }
let _ = session.output.send(OutputEvent::Exit); // PTY 读到 EOFUnix 场景):触发会话结束并清理。
cleanup_session(&session, &state);
}); });
} }
/// 轮询子进程退出状态,退出后清理会话。
fn spawn_exit_monitor(session: Arc<Session>, state: Arc<crate::AppState>) {
tokio::task::spawn_blocking(move || {
// 最多监控 5 分钟,避免无意义长驻。
for _ in 0..3000 {
let exited = match session.pty.lock() {
Ok(mut pty) => pty.try_wait().map(|s| s.is_some()).unwrap_or(false),
Err(_) => false,
};
if exited {
cleanup_session(&session, &state);
return;
}
std::thread::sleep(std::time::Duration::from_millis(100));
}
});
}
/// 触发会话结束事件并清理资源(幂等)。
fn cleanup_session(session: &Session, state: &crate::AppState) {
let _ = session.output.send(OutputEvent::Exit);
if let Ok(mut pty) = session.pty.lock() {
let _ = pty.kill();
}
state.sessions.remove(&session.id);
tracing::info!("session {} closed", session.id);
}
/// 发送一条文本消息,忽略错误。 /// 发送一条文本消息,忽略错误。
async fn send_text(socket: WebSocket, text: String) -> Result<()> { async fn send_text(socket: WebSocket, text: String) -> Result<()> {
let (mut tx, _rx) = socket.split(); let (mut tx, _rx) = socket.split();