返回 CodeWhale
terminal.rs
根目录 / crates / tui / src / runtime_api / terminal.rs
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
443 lines RUST