Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_steel/src/srv/ws/term.rs

17.7 KiB, 147 runs

created by r1870400018:12100, which is this file's identity for as long as the history lasts, whatever it is later renamed to

download · who wrote it · its history

1//! Terminal bridge — manages tmux sessions and bridges PTY I/O to
2//! WebSocket binary frames.
3//!
4//! Two layers:
5//!
6//! 1. **Management commands** (sync, via the text WS handler):
7//! `term_new`, `term_list`, `term_close`, `term_set_name`.
8//! These are thin wrappers around the `tmux` CLI and return
9//! syntax-protocol responses.
10//!
11//! 2. **Terminal I/O bridge** (async, via a separate WS path
12//! `/term/<session>`): creates a PTY pair, spawns `tmux attach`
13//! with the slave end as stdin/stdout/stderr, and bridges bytes
14//! between the PTY master and the WebSocket. The client sends
15//! keystrokes as binary WS frames; the server pushes terminal
16//! output as binary WS frames.
17//!
18//! tmux manages the PTY for the child process (e.g. Goose). We
19//! create our own PTY for the `tmux attach` process so it sees a
20//! real terminal. The tmux session persists after the WebSocket
21//! disconnects; reconnecting reattaches.
22//!
23//! # Platforms
24//!
25//! The management layer is portable: it runs a command and reads its output. The
26//! **bridge is Unix-only**, and deliberately so. It rests on a pseudo-terminal
27//! whose slave end becomes a child's three standard streams, and Windows has no
28//! counterpart to that pair of ideas -- its console pseudo-terminal is a
29//! different mechanism with a different lifetime, and a shim over it would be a
30//! second implementation rather than a translation. On a platform without a
31//! pseudo-terminal the bridge therefore refuses by name, and the refusal says
32//! which half is missing, so a caller learns it at the handshake rather than
33//! from a connection that goes quiet.
34//!
35//! This is the only place in the crate that reaches for `nix`, which is why it
36//! is the whole of the Windows task for the server.
37//!
38//! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\
39//! Anthropic Claude
40
41use oxedyne_fe2o3_core::prelude::*;
42use oxedyne_fe2o3_jdat::{
43 prelude::*,
44 id::NumIdDat,
45};
46use oxedyne_fe2o3_net::ws::core::{WebSocket, WebSocketMessage};
47use oxedyne_fe2o3_iop_crypto::enc::Encrypter;
48use oxedyne_fe2o3_iop_db::api::Database;
49use oxedyne_fe2o3_iop_hash::api::Hasher;
50
51use std::process::Command;
52
53#[cfg(unix)]
54use std::{
55 fs::File,
56 os::fd::AsFd,
57 process::Stdio,
58};
59
60#[cfg(unix)]
61use tokio::{
62 process::Command as TokioCommand,
63};
64
65
66// ┌───────────────────────────────────────────────────────────────┐
67// │ TerminalManager │
68// └───────────────────────────────────────────────────────────────┘
69
70/// Manages terminal sessions via tmux.
71///
72/// All management methods are synchronous because they are called from the sync
73/// `handle_text` WS handler. The tmux CLI calls are quick, so blocking briefly
74/// is acceptable. The terminal I/O bridge (`handle_terminal_websocket`) is
75/// async and runs on a separate WS connection.
76#[derive(Clone, Debug)]
77pub struct TerminalManager {
78 session_prefix: String, // e.g. "goose-"
79 launch_command: String, // e.g. "goose session"
80}
81
82impl TerminalManager {
83
84 pub fn new(session_prefix: &str, launch_command: &str) -> Self {
85 Self {
86 session_prefix: session_prefix.to_string(),
87 launch_command: launch_command.to_string(),
88 }
89 }
90
91 /// Runs the launch command in a new tmux session, named like "goose-3".
92 pub fn new_session(&self) -> Outcome<String> {
93 let max = ok!(self.list_session_nums()).iter().max().copied().unwrap_or(0);
94 let name = fmt!("{}{}", self.session_prefix, max + 1);
95 let status = res!(Command::new("tmux")
96 .arg("new-session")
97 .arg("-d")
98 .arg("-s")
99 .arg(&name)
100 .arg(&self.launch_command)
101 .status());
102 if !status.success() {
103 return Err(err!(
104 "tmux new-session failed for '{}'.", name;
105 IO, System));
106 }
107 Ok(name)
108 }
109
110 /// The JDAT map is `{ "sessions": [{ "name": "goose-1" }, ...] }`.
111 pub fn list_sessions_dat(&self) -> Outcome<Dat> {
112 let names = res!(self.list_session_names());
113 let mut arr = Vec::new();
114 for n in names {
115 let mut m = DaticleMap::new();
116 m.insert(dat!("name"), dat!(n));
117 arr.push(Dat::Map(m));
118 }
119 let mut out = DaticleMap::new();
120 out.insert(dat!("sessions"), Dat::List(arr));
121 Ok(Dat::Map(out))
122 }
123
124 pub fn close_session(&self, name: &str) -> Outcome<()> {
125 let status = res!(Command::new("tmux")
126 .arg("kill-session")
127 .arg("-t")
128 .arg(name)
129 .status());
130 if !status.success() {
131 return Err(err!(
132 "tmux kill-session failed for '{}'.", name;
133 IO, System));
134 }
135 Ok(())
136 }
137
138 pub fn set_session_name(&self, old: &str, new: &str) -> Outcome<()> {
139 let status = res!(Command::new("tmux")
140 .arg("rename-session")
141 .arg("-t")
142 .arg(old)
143 .arg(new)
144 .status());
145 if !status.success() {
146 return Err(err!(
147 "tmux rename-session failed '{}' -> '{}'.", old, new;
148 IO, System));
149 }
150 Ok(())
151 }
152
153 // ── Internal helpers ──────────────────────────────────────
154
155 fn list_session_names(&self) -> Outcome<Vec<String>> {
156 let output = match Command::new("tmux")
157 .arg("list-sessions")
158 .arg("-F")
159 .arg("#{session_name}")
160 .output()
161 {
162 Ok(o) => o,
163 Err(_) => return Ok(Vec::new()),
164 };
165 if !output.status.success() {
166 return Ok(Vec::new());
167 }
168 let names: Vec<String> = String::from_utf8_lossy(&output.stdout)
169 .lines()
170 .filter(|l| l.starts_with(&self.session_prefix))
171 .map(|l| l.to_string())
172 .collect();
173 Ok(names)
174 }
175
176 fn list_session_nums(&self) -> Outcome<Vec<u32>> {
177 let names = res!(self.list_session_names());
178 let nums: Vec<u32> = names
179 .iter()
180 .filter_map(|n| {
181 n.strip_prefix(&self.session_prefix)
182 .and_then(|s| s.parse().ok())
183 })
184 .collect();
185 Ok(nums)
186 }
187}
188
189
190// ┌───────────────────────────────────────────────────────────────┐
191// │ Terminal I/O bridge — PTY ↔ WebSocket │
192// └───────────────────────────────────────────────────────────────┘
193
194/// Bridges a WebSocket connection to a tmux session via a PTY.
195///
196/// Creates a PTY pair, spawns `tmux attach -t <session>` with the slave end as
197/// stdin/stdout/stderr, so tmux sees a real terminal, and bridges bytes between
198/// the PTY master and the WebSocket: client keystrokes arrive as binary frames
199/// and are written to the master, and what the master reads is pushed back as
200/// binary frames.
201///
202/// The tmux session persists after the WebSocket disconnects; reconnecting with
203/// the same session name reattaches.
204#[cfg(unix)]
205pub async fn handle_terminal_websocket<
206 const UIDL: usize,
207 UID: NumIdDat<UIDL> + 'static,
208 ENC: Encrypter + 'static,
209 KH: Hasher + 'static,
210 DB: Database<UIDL, UID, ENC, KH> + 'static,
211 S: tokio::io::AsyncRead + tokio::io::AsyncWrite + Unpin,
212>(
213 mut stream: S,
214 session_name: String,
215 request: oxedyne_fe2o3_net::http::msg::HttpMessage,
216 id: &String,
217)
218 -> Outcome<()>
219{
220 // ── WebSocket handshake ───────────────────────────────────
221 let mut ws: WebSocket<
222 '_,
223 UIDL, UID, ENC, KH, DB,
224 S,
225 oxedyne_fe2o3_net::ws::handler::WebSocketEchoHandler,
226 > = WebSocket::new_server(
227 &mut stream,
228 oxedyne_fe2o3_net::ws::handler::WebSocketEchoHandler,
229 crate::srv::constant::WEBSOCKET_CHUNK_SIZE,
230 crate::srv::constant::WEBSOCKET_CHUNKING_THRESHOLD,
231 );
232 match ws.connect_as_server(request).await {
233 Ok(()) => (),
234 Err(e) => return Err(err!(e,
235 "{}: Terminal WS handshake failed.", id;
236 IO, Network, Wire)),
237 }
238
239 // ── Create PTY ────────────────────────────────────────────
240 let winsize = nix::pty::Winsize {
241 ws_row: 24,
242 ws_col: 80,
243 ws_xpixel: 0,
244 ws_ypixel: 0,
245 };
246 let pty = match nix::pty::openpty(&Some(winsize), &None) {
247 Ok(p) => p,
248 Err(e) => {
249 let err_msg = WebSocketMessage::Text(
250 fmt!("error \"PTY creation failed: {}\"", e));
251 let _ = ws.send(&err_msg).await;
252 return Err(err!(e,
253 "{}: Failed to create PTY for '{}'.", id, session_name;
254 IO, System));
255 }
256 };
257
258 // pty.master and pty.slave are OwnedFd (safe FD wrappers).
259
260 // Dup the slave for each stdio stream using nix's safe API.
261 let slave_stdin = match nix::unistd::dup(pty.slave.as_fd()) {
262 Ok(f) => f,
263 Err(e) => return Err(err!(e, "{}: dup slave stdin failed.", id; IO, System)),
264 };
265 let slave_stdout = match nix::unistd::dup(pty.slave.as_fd()) {
266 Ok(f) => f,
267 Err(e) => return Err(err!(e, "{}: dup slave stdout failed.", id; IO, System)),
268 };
269 let slave_stderr = match nix::unistd::dup(pty.slave.as_fd()) {
270 Ok(f) => f,
271 Err(e) => return Err(err!(e, "{}: dup slave stderr failed.", id; IO, System)),
272 };
273
274 // ── Spawn tmux attach with the PTY slave ──────────────────
275 let mut cmd = TokioCommand::new("tmux");
276 cmd.arg("attach-session")
277 .arg("-t")
278 .arg(&session_name)
279 .stdin(Stdio::from(slave_stdin))
280 .stdout(Stdio::from(slave_stdout))
281 .stderr(Stdio::from(slave_stderr))
282 .env("TERM", "xterm-256color");
283
284 let mut child = match cmd.spawn() {
285 Ok(c) => c,
286 Err(e) => {
287 let err_msg = WebSocketMessage::Text(
288 fmt!("error \"tmux attach failed: {}\"", e));
289 let _ = ws.send(&err_msg).await;
290 return Err(err!(e,
291 "{}: Failed to spawn tmux attach for '{}'.", id, session_name;
292 IO, System));
293 }
294 };
295
296 // ── Set master to non-blocking ────────────────────────────
297 //
298 // We must do this before wrapping the master in AsyncFd.
299 // nix::fcntl::fcntl takes AsFd, so we pass the OwnedFd
300 // before it is consumed by the File conversion below.
301 if let Err(e) = nix::fcntl::fcntl(
302 &pty.master,
303 nix::fcntl::FcntlArg::F_SETFL(nix::fcntl::OFlag::O_NONBLOCK),
304 ) {
305 warn!("{}: Failed to set master non-blocking: {}", id, e);
306 }
307
308 // ── Wrap the master for async I/O ─────────────────────────
309 //
310 // Convert the master OwnedFd to a File for AsyncFd. The
311 // OwnedFd is consumed and the File takes ownership of the FD.
312 let master_file: File = pty.master.into();
313 let master_async = match tokio::io::unix::AsyncFd::new(master_file) {
314 Ok(a) => a,
315 Err(e) => {
316 let _ = child.kill().await;
317 return Err(err!(e,
318 "{}: Failed to create AsyncFd for PTY master.", id;
319 IO, System));
320 }
321 };
322
323 // ── Bidirectional pipe loop ───────────────────────────────
324 let mut buf = vec![0u8; 16384];
325
326 loop {
327 tokio::select! {
328 // ── Read from WebSocket → write to PTY master ──────
329 ws_read = ws.read() => {
330 match ws_read {
331 Ok(Some(WebSocketMessage::Binary(byts))) => {
332 use std::io::Write;
333 if let Err(e) = (&*master_async.get_ref()).write(&byts) {
334 if e.kind() != std::io::ErrorKind::WouldBlock {
335 warn!("{}: PTY write error: {}", id, e);
336 break;
337 }
338 }
339 }
340 Ok(Some(WebSocketMessage::Text(txt))) => {
341 use std::io::Write;
342 if let Err(e) = (&*master_async.get_ref()).write(txt.as_bytes()) {
343 if e.kind() != std::io::ErrorKind::WouldBlock {
344 warn!("{}: PTY write error: {}", id, e);
345 break;
346 }
347 }
348 }
349 Ok(Some(WebSocketMessage::Close(_, _))) => {
350 debug!("{}: Terminal WS client closed.", id);
351 break;
352 }
353 Ok(None) => {
354 debug!("{}: Terminal WS client disconnected.", id);
355 break;
356 }
357 Ok(_) => (),
358 Err(e) => {
359 warn!("{}: Terminal WS read error: {}", id, e);
360 break;
361 }
362 }
363 }
364 // ── Read from PTY master → send as binary WS ───────
365 readable = master_async.readable() => {
366 let mut guard = match readable {
367 Ok(g) => g,
368 Err(e) => {
369 warn!("{}: PTY readable error: {}", id, e);
370 break;
371 }
372 };
373 use std::io::Read;
374 match guard.try_io(|afd| afd.get_ref().read(&mut buf)) {
375 Ok(Ok(0)) => {
376 debug!("{}: PTY EOF for '{}'.", id, session_name);
377 break;
378 }
379 Ok(Ok(n)) => {
380 let msg = WebSocketMessage::Binary(buf[..n].to_vec());
381 if let Err(e) = ws.send(&msg).await {
382 warn!("{}: WS send error: {}", id, e);
383 break;
384 }
385 }
386 Ok(Err(ref e)) if e.kind() == std::io::ErrorKind::WouldBlock => {
387 guard.clear_ready();
388 }
389 Ok(Err(e)) => {
390 warn!("{}: PTY read error: {}", id, e);
391 break;
392 }
393 Err(_) => {
394 guard.clear_ready();
395 }
396 }
397 }
398 }
399 }
400
401 // ── Cleanup ───────────────────────────────────────────────
402 // Kill the tmux attach process. The tmux session itself
403 // persists for reconnection.
404 let _ = child.kill().await;
405 let _ = child.wait().await;
406 let _ = ws.send(&WebSocketMessage::Close(
407 None, Some("session ended".to_string()))).await;
408 debug!("{}: Terminal bridge closed for '{}'.", id, session_name);
409 Ok(())
410}
411
412/// The same, on a platform with no pseudo-terminal: the handshake is completed
413/// and the bridge then refuses by name.
414///
415/// The handshake happens first on purpose. A client that has asked to attach to
416/// a terminal is owed the reason it cannot, and a WebSocket that is closed
417/// before it opens carries no reason at all -- the browser reports a failed
418/// connection and nothing about why. So the socket is established, the refusal
419/// is sent as text on it in the same shape as every other error this bridge
420/// reports, and only then does the call return an error for the log.
421#[cfg(not(unix))]
422pub async fn handle_terminal_websocket<
423 const UIDL: usize,
424 UID: NumIdDat<UIDL> + 'static,
425 ENC: Encrypter + 'static,
426 KH: Hasher + 'static,
427 DB: Database<UIDL, UID, ENC, KH> + 'static,
428 S: tokio::io::AsyncRead + tokio::io::AsyncWrite + Unpin,
429>(
430 mut stream: S,
431 session_name: String,
432 request: oxedyne_fe2o3_net::http::msg::HttpMessage,
433 id: &String,
434)
435 -> Outcome<()>
436{
437 let mut ws: WebSocket<
438 '_,
439 UIDL, UID, ENC, KH, DB,
440 S,
441 oxedyne_fe2o3_net::ws::handler::WebSocketEchoHandler,
442 > = WebSocket::new_server(
443 &mut stream,
444 oxedyne_fe2o3_net::ws::handler::WebSocketEchoHandler,
445 crate::srv::constant::WEBSOCKET_CHUNK_SIZE,
446 crate::srv::constant::WEBSOCKET_CHUNKING_THRESHOLD,
447 );
448 match ws.connect_as_server(request).await {
449 Ok(()) => (),
450 Err(e) => return Err(err!(e,
451 "{}: Terminal WS handshake failed.", id;
452 IO, Network, Wire)),
453 }
454 let msg = WebSocketMessage::Text(fmt!(
455 "error \"A terminal session needs a pseudo-terminal, which this platform \
456 does not provide.\""));
457 let _ = ws.send(&msg).await;
458 let _ = ws.send(&WebSocketMessage::Close(
459 None, Some("no pseudo-terminal".to_string()))).await;
460 Err(err!(
461 "{}: A bridge to terminal session '{}' was asked for on a platform with no \
462 pseudo-terminal. The management commands work here; the bridge does not.",
463 id, session_name;
464 Unimplemented, System))
465}