| 1 | //! `/v1/terminal/{name}` — the Engine's terminal byte stream. |
| 2 | //! |
| 3 | //! The owner is [`crate::tools::terminal_session`]: the same PTY-backed shell |
| 4 | //! the agent's terminal tools drive, so a client attaching here sees the |
| 5 | //! session the model is already working in rather than a second shell beside |
| 6 | //! it. These routes never create a session — a name with no live session is a |
| 7 | //! 404, because conjuring a shell from an HTTP request would give the app a |
| 8 | //! terminal the Engine does not know about. |
| 9 | //! |
| 10 | //! Authentication is the `/v1` route layer's bearer token (see |
| 11 | //! [`super::auth`]); nothing here re-implements or bypasses it. |
| 12 | //! |
| 13 | //! Wire shape deliberately follows `/v1/threads/{id}/jobs/{job_id}/output`: |
| 14 | //! `cursor` / `max_bytes` / `format` in, `offset` / `next_cursor` / `total` / |
| 15 | //! `dropped` out. Two byte streams in one product should not speak two |
| 16 | //! dialects. |
| 17 | //! |
| 18 | //! Known limitations, recorded because a reader will otherwise assume them: |
| 19 | //! |
| 20 | //! - **No long poll.** `wait_ms` is not accepted; a client polls the cursor. |
| 21 | //! The jobs route can block because a job owns a notification; a terminal |
| 22 | //! session's ring has no wake-up channel yet, and inventing one here would |
| 23 | //! be a second mechanism rather than a reuse. |
| 24 | //! - **No scrollback recovery.** `dropped` reports what the 512 KiB ring |
| 25 | //! discarded; those bytes are gone with the process, not on disk. |
| 26 | //! - **Live sessions only.** Persistence is identity and lifecycle, never |
| 27 | //! output, so a restarted Engine reports no session rather than pretending to |
| 28 | //! reattach (#34 acceptance: "Restart truthfully reports lost live PTYs"). |
| 29 | //! - **Unix only.** The owner is `#[cfg(all(unix, not(target_env = "ohos")))]` end to end; on Windows these |
| 30 | //! routes do not exist yet. ConPTY qualification is its own slice. |
| 31 | |
| 32 | use axum::Json; |
| 33 | use axum::extract::{Path, Query, State}; |
| 34 | #[cfg(all(unix, not(target_env = "ohos")))] |
| 35 | use base64::Engine as _; |
| 36 | use serde::{Deserialize, Serialize}; |
| 37 | |
| 38 | // The owner does not exist on ohos (`tools/mod.rs`), so neither does any |
| 39 | // handler that drives it; the stubs below answer there instead. |
| 40 | #[cfg(all(unix, not(target_env = "ohos")))] |
| 41 | use crate::tools::terminal_session; |
| 42 | |
| 43 | use super::{ApiError, RuntimeApiState}; |
| 44 | |
| 45 | /// Default per-response ceiling; the owner clamps to its own `READ_LIMIT`. |
| 46 | #[cfg(all(unix, not(target_env = "ohos")))] |
| 47 | const TERMINAL_CHUNK_DEFAULT: usize = 64 * 1024; |
| 48 | /// Session names come from the agent's tools; this only bounds the echo. |
| 49 | #[cfg(all(unix, not(target_env = "ohos")))] |
| 50 | const TERMINAL_NAME_MAX_BYTES: usize = 128; |
| 51 | /// One input frame. Interactive typing is bytes, not uploads. |
| 52 | #[cfg(all(unix, not(target_env = "ohos")))] |
| 53 | const TERMINAL_INPUT_MAX_BYTES: usize = 64 * 1024; |
| 54 | #[cfg(all(unix, not(target_env = "ohos")))] |
| 55 | const TERMINAL_DIMENSION_MAX: u16 = 1000; |
| 56 | |
| 57 | // The request shape is the contract on every platform; only the Unix |
| 58 | // handlers read it, so the 501 builds expect the fields to stay unread. |
| 59 | #[cfg_attr(any(not(unix), target_env = "ohos"), expect(dead_code))] |
| 60 | #[derive(Deserialize)] |
| 61 | #[serde(deny_unknown_fields)] |
| 62 | pub(super) struct TerminalOutputQuery { |
| 63 | /// Absolute byte offset into the session's lifetime output. |
| 64 | #[serde(default)] |
| 65 | cursor: Option<u64>, |
| 66 | /// Per-response byte ceiling, default 64 KiB, clamped by the owner. |
| 67 | #[serde(default)] |
| 68 | max_bytes: Option<usize>, |
| 69 | /// `base64` (default, exact bytes) or `text` (lossy UTF-8). |
| 70 | #[serde(default)] |
| 71 | format: Option<String>, |
| 72 | } |
| 73 | |
| 74 | #[derive(Debug, Serialize)] |
| 75 | pub(super) struct TerminalOutputResponse { |
| 76 | name: String, |
| 77 | /// Absolute offset of `data[0]`; exceeds `cursor` when the ring already |
| 78 | /// discarded that prefix (`dropped` reports the cutoff). |
| 79 | offset: u64, |
| 80 | /// Pass back as `cursor` to continue. |
| 81 | next_cursor: u64, |
| 82 | /// Everything the session has produced, including discarded bytes. |
| 83 | total: u64, |
| 84 | /// Leading bytes the bounded ring permanently discarded. |
| 85 | dropped: u64, |
| 86 | encoding: &'static str, |
| 87 | data: String, |
| 88 | /// False once the shell has exited and no bytes remain past `next_cursor`. |
| 89 | running: bool, |
| 90 | exit_code: Option<i64>, |
| 91 | } |
| 92 | |
| 93 | #[cfg_attr(any(not(unix), target_env = "ohos"), expect(dead_code))] |
| 94 | #[derive(Deserialize)] |
| 95 | #[serde(deny_unknown_fields)] |
| 96 | pub(super) struct TerminalInputRequest { |
| 97 | data: String, |
| 98 | /// `base64` (default, exact bytes) or `text`. |
| 99 | #[serde(default)] |
| 100 | encoding: Option<String>, |
| 101 | } |
| 102 | |
| 103 | #[derive(Debug, Serialize)] |
| 104 | pub(super) struct TerminalWriteResponse { |
| 105 | name: String, |
| 106 | written: usize, |
| 107 | } |
| 108 | |
| 109 | #[cfg_attr(any(not(unix), target_env = "ohos"), expect(dead_code))] |
| 110 | #[derive(Deserialize)] |
| 111 | #[serde(deny_unknown_fields)] |
| 112 | pub(super) struct TerminalResizeRequest { |
| 113 | rows: u16, |
| 114 | cols: u16, |
| 115 | } |
| 116 | |
| 117 | #[derive(Debug, Serialize)] |
| 118 | pub(super) struct TerminalResizeResponse { |
| 119 | name: String, |
| 120 | rows: u16, |
| 121 | cols: u16, |
| 122 | } |
| 123 | |
| 124 | #[derive(Debug, Serialize)] |
| 125 | pub(super) struct TerminalKillResponse { |
| 126 | name: String, |
| 127 | killed: bool, |
| 128 | } |
| 129 | |
| 130 | /// `base64` keeps bytes exact; `text` is the lossy convenience form. |
| 131 | #[cfg(all(unix, not(target_env = "ohos")))] |
| 132 | fn chunk_encoding(format: &str) -> Result<&'static str, ApiError> { |
| 133 | match format { |
| 134 | "base64" => Ok("base64"), |
| 135 | "text" => Ok("text"), |
| 136 | _ => Err(ApiError::bad_request("format must be base64 or text")), |
| 137 | } |
| 138 | } |
| 139 | |
| 140 | #[cfg(all(unix, not(target_env = "ohos")))] |
| 141 | fn encode_bytes(bytes: &[u8], encoding: &str) -> String { |
| 142 | if encoding == "base64" { |
| 143 | base64::engine::general_purpose::STANDARD.encode(bytes) |
| 144 | } else { |
| 145 | String::from_utf8_lossy(bytes).into_owned() |
| 146 | } |
| 147 | } |
| 148 | |
| 149 | #[cfg(all(unix, not(target_env = "ohos")))] |
| 150 | fn decode_bytes(data: &str, encoding: &str) -> Result<Vec<u8>, ApiError> { |
| 151 | let bytes = match encoding { |
| 152 | "base64" => base64::engine::general_purpose::STANDARD |
| 153 | .decode(data) |
| 154 | .map_err(|_| ApiError::bad_request("data is not valid base64"))?, |
| 155 | "text" => data.as_bytes().to_vec(), |
| 156 | _ => return Err(ApiError::bad_request("encoding must be base64 or text")), |
| 157 | }; |
| 158 | if bytes.len() > TERMINAL_INPUT_MAX_BYTES { |
| 159 | return Err(ApiError::bad_request(format!( |
| 160 | "input exceeds {TERMINAL_INPUT_MAX_BYTES} bytes" |
| 161 | ))); |
| 162 | } |
| 163 | Ok(bytes) |
| 164 | } |
| 165 | |
| 166 | #[cfg(all(unix, not(target_env = "ohos")))] |
| 167 | fn bounded_max_bytes(requested: Option<usize>) -> Result<usize, ApiError> { |
| 168 | let max_bytes = requested.unwrap_or(TERMINAL_CHUNK_DEFAULT); |
| 169 | if !(1..=terminal_session::READ_LIMIT).contains(&max_bytes) { |
| 170 | return Err(ApiError::bad_request(format!( |
| 171 | "max_bytes must be between 1 and {}", |
| 172 | terminal_session::READ_LIMIT |
| 173 | ))); |
| 174 | } |
| 175 | Ok(max_bytes) |
| 176 | } |
| 177 | |
| 178 | #[cfg(all(unix, not(target_env = "ohos")))] |
| 179 | fn bounded_dimension(value: u16, field: &str) -> Result<u16, ApiError> { |
| 180 | if !(1..=TERMINAL_DIMENSION_MAX).contains(&value) { |
| 181 | return Err(ApiError::bad_request(format!( |
| 182 | "{field} must be between 1 and {TERMINAL_DIMENSION_MAX}" |
| 183 | ))); |
| 184 | } |
| 185 | Ok(value) |
| 186 | } |
| 187 | |
| 188 | /// Resolve a live session or 404. Never creates one — see the module docs. |
| 189 | #[cfg(all(unix, not(target_env = "ohos")))] |
| 190 | fn open_session( |
| 191 | state: &RuntimeApiState, |
| 192 | name: &str, |
| 193 | ) -> Result<terminal_session::SharedSession, ApiError> { |
| 194 | if name.is_empty() || name.len() > TERMINAL_NAME_MAX_BYTES { |
| 195 | return Err(ApiError::not_found("terminal session not found")); |
| 196 | } |
| 197 | terminal_session::lookup(name, &state.workspace) |
| 198 | .ok_or_else(|| ApiError::not_found(format!("no live terminal session named '{name}'"))) |
| 199 | } |
| 200 | |
| 201 | /// How many terminal route calls may occupy blocking threads at once. The |
| 202 | /// rest wait in `with_session` asynchronously, so a client that disconnects |
| 203 | /// while queued simply disappears instead of holding a pool thread, and one |
| 204 | /// stuck write (a child that stopped reading, lock held) can stall terminal |
| 205 | /// routes but never the runtime's other blocking work. |
| 206 | #[cfg(all(unix, not(target_env = "ohos")))] |
| 207 | static ROUTE_GATE: tokio::sync::Semaphore = tokio::sync::Semaphore::const_new(8); |
| 208 | |
| 209 | /// Run one operation against the locked session on the blocking pool. |
| 210 | /// |
| 211 | /// The session mutex and the PTY behind it are synchronous: the agent's own |
| 212 | /// tool holds the lock across a whole command, and a write to a child that |
| 213 | /// stopped reading blocks until the kernel buffer drains. Neither may park a |
| 214 | /// runtime worker (#6149), so a route never touches the session inline, and |
| 215 | /// `ROUTE_GATE` bounds how many such touches can be in flight. |
| 216 | #[cfg(all(unix, not(target_env = "ohos")))] |
| 217 | async fn with_session<T>( |
| 218 | session: terminal_session::SharedSession, |
| 219 | operation: impl FnOnce(&mut terminal_session::TerminalSession) -> Result<T, ApiError> |
| 220 | + Send |
| 221 | + 'static, |
| 222 | ) -> Result<T, ApiError> |
| 223 | where |
| 224 | T: Send + 'static, |
| 225 | { |
| 226 | let permit = ROUTE_GATE |
| 227 | .acquire() |
| 228 | .await |
| 229 | .map_err(|_| ApiError::internal("terminal route gate closed"))?; |
| 230 | tokio::task::spawn_blocking(move || { |
| 231 | let _permit = permit; |
| 232 | let mut guard = session |
| 233 | .lock() |
| 234 | .map_err(|_| ApiError::internal("terminal session lock poisoned"))?; |
| 235 | operation(&mut guard) |
| 236 | }) |
| 237 | .await |
| 238 | .map_err(|error| ApiError::internal(error.to_string()))? |
| 239 | } |
| 240 | |
| 241 | /// `GET /v1/terminal/{name}/output` — the resumable byte stream. |
| 242 | /// |
| 243 | /// Reads are non-consuming: several clients may hold independent cursors, and |
| 244 | /// polling here never steals output from the agent's own consuming read. |
| 245 | #[cfg(all(unix, not(target_env = "ohos")))] |
| 246 | pub(super) async fn terminal_output( |
| 247 | State(state): State<RuntimeApiState>, |
| 248 | Path(name): Path<String>, |
| 249 | Query(query): Query<TerminalOutputQuery>, |
| 250 | ) -> Result<Json<TerminalOutputResponse>, ApiError> { |
| 251 | let session = open_session(&state, &name)?; |
| 252 | let encoding = chunk_encoding(query.format.as_deref().unwrap_or("base64"))?; |
| 253 | let max_bytes = bounded_max_bytes(query.max_bytes)?; |
| 254 | let cursor = query.cursor.unwrap_or(0); |
| 255 | let (chunk, exit) = with_session(session, move |session| { |
| 256 | let chunk = terminal_session::read_session_since(session, cursor, max_bytes) |
| 257 | .map_err(ApiError::internal)?; |
| 258 | let exit = terminal_session::session_exit_status(session).map_err(ApiError::internal)?; |
| 259 | Ok((chunk, exit)) |
| 260 | }) |
| 261 | .await?; |
| 262 | let running = exit.is_none(); |
| 263 | Ok(Json(TerminalOutputResponse { |
| 264 | name, |
| 265 | offset: chunk.offset, |
| 266 | next_cursor: chunk.next_cursor, |
| 267 | total: chunk.total, |
| 268 | dropped: chunk.dropped, |
| 269 | encoding, |
| 270 | data: encode_bytes(&chunk.bytes, encoding), |
| 271 | // A gap means bytes were lost; `running` alone must not imply there is |
| 272 | // nothing behind us, so drain state is reported independently. |
| 273 | running, |
| 274 | exit_code: exit.map(|status| i64::from(status.exit_code())), |
| 275 | })) |
| 276 | } |
| 277 | |
| 278 | /// `POST /v1/terminal/{name}/input` — bytes into the live shell. |
| 279 | /// |
| 280 | /// Input attribution is the caller's: this route is the client's writer, and |
| 281 | /// the agent's writer is `terminal_send`. Nothing here re-labels one as the |
| 282 | /// other. |
| 283 | #[cfg(all(unix, not(target_env = "ohos")))] |
| 284 | pub(super) async fn terminal_input( |
| 285 | State(state): State<RuntimeApiState>, |
| 286 | Path(name): Path<String>, |
| 287 | Json(request): Json<TerminalInputRequest>, |
| 288 | ) -> Result<Json<TerminalWriteResponse>, ApiError> { |
| 289 | let session = open_session(&state, &name)?; |
| 290 | let bytes = decode_bytes( |
| 291 | &request.data, |
| 292 | request.encoding.as_deref().unwrap_or("base64"), |
| 293 | )?; |
| 294 | let written = bytes.len(); |
| 295 | with_session(session, move |session| { |
| 296 | terminal_session::write_bytes(session, &bytes).map_err(ApiError::internal) |
| 297 | }) |
| 298 | .await?; |
| 299 | Ok(Json(TerminalWriteResponse { name, written })) |
| 300 | } |
| 301 | |
| 302 | /// `POST /v1/terminal/{name}/resize` — the window the child should draw for. |
| 303 | #[cfg(all(unix, not(target_env = "ohos")))] |
| 304 | pub(super) async fn terminal_resize( |
| 305 | State(state): State<RuntimeApiState>, |
| 306 | Path(name): Path<String>, |
| 307 | Json(request): Json<TerminalResizeRequest>, |
| 308 | ) -> Result<Json<TerminalResizeResponse>, ApiError> { |
| 309 | let session = open_session(&state, &name)?; |
| 310 | let rows = bounded_dimension(request.rows, "rows")?; |
| 311 | let cols = bounded_dimension(request.cols, "cols")?; |
| 312 | with_session(session, move |session| { |
| 313 | terminal_session::resize_session(session, rows, cols).map_err(ApiError::internal) |
| 314 | }) |
| 315 | .await?; |
| 316 | Ok(Json(TerminalResizeResponse { name, rows, cols })) |
| 317 | } |
| 318 | |
| 319 | /// `POST /v1/terminal/{name}/kill` — end the shell. |
| 320 | /// |
| 321 | /// The exit itself is observed through `output` (`running` / `exit_code`), |
| 322 | /// so a client that kills and then polls learns the truth instead of an |
| 323 | /// optimistic acknowledgement. |
| 324 | #[cfg(all(unix, not(target_env = "ohos")))] |
| 325 | pub(super) async fn terminal_kill( |
| 326 | State(state): State<RuntimeApiState>, |
| 327 | Path(name): Path<String>, |
| 328 | ) -> Result<Json<TerminalKillResponse>, ApiError> { |
| 329 | let session = open_session(&state, &name)?; |
| 330 | with_session(session, |session| { |
| 331 | terminal_session::kill_session(session).map_err(ApiError::internal) |
| 332 | }) |
| 333 | .await?; |
| 334 | Ok(Json(TerminalKillResponse { name, killed: true })) |
| 335 | } |
| 336 | |
| 337 | /// Windows build: the owner is `#[cfg(all(unix, not(target_env = "ohos")))]` end to end, so the contract |
| 338 | /// exists but cannot be served. These answer 501 rather than 404 so a client |
| 339 | /// can tell "this Engine build cannot do terminals" apart from "that session |
| 340 | /// is gone" — and so the ConPTY slice has one place to replace. |
| 341 | #[cfg(any(not(unix), target_env = "ohos"))] |
| 342 | mod platform { |
| 343 | use super::*; |
| 344 | |
| 345 | fn unsupported() -> ApiError { |
| 346 | ApiError::not_implemented( |
| 347 | "terminal sessions are Unix-only in this build; native Windows PTY support is not implemented yet", |
| 348 | ) |
| 349 | } |
| 350 | |
| 351 | pub(crate) async fn terminal_output( |
| 352 | State(_): State<RuntimeApiState>, |
| 353 | Path(_): Path<String>, |
| 354 | Query(_): Query<TerminalOutputQuery>, |
| 355 | ) -> Result<Json<TerminalOutputResponse>, ApiError> { |
| 356 | Err(unsupported()) |
| 357 | } |
| 358 | |
| 359 | pub(crate) async fn terminal_input( |
| 360 | State(_): State<RuntimeApiState>, |
| 361 | Path(_): Path<String>, |
| 362 | Json(_): Json<TerminalInputRequest>, |
| 363 | ) -> Result<Json<TerminalWriteResponse>, ApiError> { |
| 364 | Err(unsupported()) |
| 365 | } |
| 366 | |
| 367 | pub(crate) async fn terminal_resize( |
| 368 | State(_): State<RuntimeApiState>, |
| 369 | Path(_): Path<String>, |
| 370 | Json(_): Json<TerminalResizeRequest>, |
| 371 | ) -> Result<Json<TerminalResizeResponse>, ApiError> { |
| 372 | Err(unsupported()) |
| 373 | } |
| 374 | |
| 375 | pub(crate) async fn terminal_kill( |
| 376 | State(_): State<RuntimeApiState>, |
| 377 | Path(_): Path<String>, |
| 378 | ) -> Result<Json<TerminalKillResponse>, ApiError> { |
| 379 | Err(unsupported()) |
| 380 | } |
| 381 | } |
| 382 | |
| 383 | #[cfg(any(not(unix), target_env = "ohos"))] |
| 384 | pub(super) use platform::{terminal_input, terminal_kill, terminal_output, terminal_resize}; |
| 385 | |
| 386 | #[cfg(all(test, unix, not(target_env = "ohos")))] |
| 387 | mod tests { |
| 388 | use super::*; |
| 389 | |
| 390 | #[test] |
| 391 | fn encodings_round_trip_exact_bytes_and_stay_lossy_only_on_request() { |
| 392 | // Non-UTF-8 bytes survive base64 and are the reason it is the default. |
| 393 | let raw = [0xf0, 0x9f, 0x90, 0x8b, 0x00, 0xff]; |
| 394 | let encoded = encode_bytes(&raw, "base64"); |
| 395 | assert_eq!(decode_bytes(&encoded, "base64").unwrap(), raw); |
| 396 | // The lossy form is explicit and cannot be mistaken for fidelity. |
| 397 | let text = encode_bytes(&raw, "text"); |
| 398 | assert!(text.contains('\u{fffd}')); |
| 399 | assert_eq!(decode_bytes(&text, "text").unwrap(), text.as_bytes()); |
| 400 | } |
| 401 | |
| 402 | #[test] |
| 403 | fn encoding_names_are_closed_sets() { |
| 404 | for good in ["base64", "text"] { |
| 405 | assert_eq!(chunk_encoding(good).unwrap(), good); |
| 406 | } |
| 407 | for bad in ["utf8", "raw", "Base64", ""] { |
| 408 | assert!(chunk_encoding(bad).is_err(), "{bad} must not be accepted"); |
| 409 | assert!(decode_bytes("", bad).is_err(), "{bad} must not decode"); |
| 410 | } |
| 411 | // A base64 decoder that ignores padding would accept junk bytes. |
| 412 | assert!(decode_bytes("not base64!!", "base64").is_err()); |
| 413 | } |
| 414 | |
| 415 | #[test] |
| 416 | fn chunk_and_dimension_bounds_reject_the_edges() { |
| 417 | assert_eq!(bounded_max_bytes(None).unwrap(), TERMINAL_CHUNK_DEFAULT); |
| 418 | assert_eq!( |
| 419 | bounded_max_bytes(Some(terminal_session::READ_LIMIT)).unwrap(), |
| 420 | terminal_session::READ_LIMIT |
| 421 | ); |
| 422 | assert!(bounded_max_bytes(Some(0)).is_err()); |
| 423 | assert!(bounded_max_bytes(Some(terminal_session::READ_LIMIT + 1)).is_err()); |
| 424 | assert_eq!(bounded_dimension(24, "rows").unwrap(), 24); |
| 425 | assert!(bounded_dimension(0, "rows").is_err()); |
| 426 | assert!(bounded_dimension(TERMINAL_DIMENSION_MAX + 1, "cols").is_err()); |
| 427 | } |
| 428 | |
| 429 | #[test] |
| 430 | fn input_is_bounded_before_it_reaches_the_pty() { |
| 431 | let too_much = |
| 432 | base64::engine::general_purpose::STANDARD |
| 433 | .encode(vec![b'a'; TERMINAL_INPUT_MAX_BYTES + 1]); |
| 434 | assert!(decode_bytes(&too_much, "base64").is_err()); |
| 435 | let at_limit = |
| 436 | base64::engine::general_purpose::STANDARD.encode(vec![b'a'; TERMINAL_INPUT_MAX_BYTES]); |
| 437 | assert_eq!( |
| 438 | decode_bytes(&at_limit, "base64").unwrap().len(), |
| 439 | TERMINAL_INPUT_MAX_BYTES |
| 440 | ); |
| 441 | } |
| 442 | } |
| 443 |