| 1 | //! `/v1/jobs` — the client-facing shell job surface. |
| 2 | //! |
| 3 | //! One authority: every job lives on the thread's shared `ShellManager`, the |
| 4 | //! same manager the thread's engine uses for model-launched shell work. Jobs |
| 5 | //! created here carry an `api:{thread_id}` owner scope so an engine's |
| 6 | //! per-session completion drain never claims client-launched work as model |
| 7 | //! evidence — and the model's jobs never appear "client-owned" here. |
| 8 | |
| 9 | use std::collections::HashMap; |
| 10 | |
| 11 | use axum::Json; |
| 12 | use axum::extract::{Path, Query, State}; |
| 13 | use axum::http::StatusCode; |
| 14 | use base64::Engine as _; |
| 15 | use serde::{Deserialize, Serialize}; |
| 16 | |
| 17 | use crate::tools::shell::{ |
| 18 | PtyDimensions, ShellJobSnapshot, ShellManager, ShellOutputChunk, ShellOutputStream, |
| 19 | ShellResult, ShellStatus, |
| 20 | }; |
| 21 | |
| 22 | use super::{ApiError, RuntimeApiState, map_thread_err}; |
| 23 | |
| 24 | /// `owner_session_id` prefix for jobs launched through this API. A real |
| 25 | /// engine session id can never carry it, so `*_for_session` drains stay |
| 26 | /// model-owned and API listings can tell the two apart. |
| 27 | const API_JOB_SCOPE_PREFIX: &str = "api:"; |
| 28 | |
| 29 | const COMMAND_MAX_BYTES: usize = 32 * 1024; |
| 30 | const ENV_MAX_ENTRIES: usize = 64; |
| 31 | const ENV_KEY_MAX_BYTES: usize = 128; |
| 32 | const ENV_VALUE_MAX_BYTES: usize = 8 * 1024; |
| 33 | const OUTPUT_CHUNK_DEFAULT: usize = 64 * 1024; |
| 34 | const OUTPUT_CHUNK_MAX: usize = 512 * 1024; |
| 35 | const OUTPUT_WAIT_MAX_MS: u64 = 30_000; |
| 36 | const STDIN_MAX_BYTES: usize = 64 * 1024; |
| 37 | const JOB_ID_MAX_BYTES: usize = 128; |
| 38 | |
| 39 | fn api_job_scope(thread_id: &str) -> String { |
| 40 | format!("{API_JOB_SCOPE_PREFIX}{thread_id}") |
| 41 | } |
| 42 | |
| 43 | fn job_owner(snapshot: &ShellJobSnapshot) -> &'static str { |
| 44 | if snapshot.owner_agent_id.is_some() { |
| 45 | "subagent" |
| 46 | } else if snapshot.owner_session_id.starts_with(API_JOB_SCOPE_PREFIX) { |
| 47 | "client" |
| 48 | } else { |
| 49 | "agent" |
| 50 | } |
| 51 | } |
| 52 | |
| 53 | fn map_job_err(error: anyhow::Error) -> ApiError { |
| 54 | let message = error.to_string(); |
| 55 | if message.ends_with("not found") { |
| 56 | ApiError::not_found(message) |
| 57 | } else { |
| 58 | ApiError::internal(message) |
| 59 | } |
| 60 | } |
| 61 | |
| 62 | #[derive(Debug, Serialize)] |
| 63 | pub(super) struct JobView { |
| 64 | #[serde(flatten)] |
| 65 | snapshot: ShellJobSnapshot, |
| 66 | thread_id: String, |
| 67 | /// `client` = launched through this API, `agent` = launched by the model's |
| 68 | /// shell tool, `subagent` = owned by a delegated agent. |
| 69 | owner: &'static str, |
| 70 | /// Null for historical records whose original transport is not known. |
| 71 | tty: Option<bool>, |
| 72 | terminal_size: Option<PtyDimensions>, |
| 73 | } |
| 74 | |
| 75 | impl JobView { |
| 76 | fn new(snapshot: ShellJobSnapshot, thread_id: String, manager: &ShellManager) -> Self { |
| 77 | let owner = job_owner(&snapshot); |
| 78 | let terminal_size = manager.job_terminal_size(&snapshot.job_id); |
| 79 | let tty = (!snapshot.stale).then_some(terminal_size.is_some()); |
| 80 | Self { |
| 81 | snapshot, |
| 82 | thread_id, |
| 83 | owner, |
| 84 | tty, |
| 85 | terminal_size, |
| 86 | } |
| 87 | } |
| 88 | } |
| 89 | |
| 90 | #[derive(Debug, Serialize)] |
| 91 | pub(super) struct JobListResponse { |
| 92 | jobs: Vec<JobView>, |
| 93 | } |
| 94 | |
| 95 | /// `GET /v1/jobs` — every live and known-stale job across all threads. |
| 96 | pub(super) async fn list_jobs( |
| 97 | State(state): State<RuntimeApiState>, |
| 98 | ) -> Result<Json<JobListResponse>, ApiError> { |
| 99 | let managers = state.runtime_threads.shell_managers_snapshot().await; |
| 100 | let jobs = tokio::task::spawn_blocking(move || { |
| 101 | let mut jobs = Vec::new(); |
| 102 | for (thread_id, manager) in managers { |
| 103 | let mut guard = manager.lock().unwrap_or_else(|e| e.into_inner()); |
| 104 | jobs.extend( |
| 105 | guard |
| 106 | .list_jobs() |
| 107 | .into_iter() |
| 108 | .map(|snapshot| JobView::new(snapshot, thread_id.clone(), &guard)), |
| 109 | ); |
| 110 | } |
| 111 | jobs |
| 112 | }) |
| 113 | .await |
| 114 | .map_err(|_| ApiError::internal("job listing failed"))?; |
| 115 | Ok(Json(JobListResponse { jobs })) |
| 116 | } |
| 117 | |
| 118 | async fn thread_manager( |
| 119 | state: &RuntimeApiState, |
| 120 | thread_id: &str, |
| 121 | create: bool, |
| 122 | ) -> Result<crate::tools::shell::SharedShellManager, ApiError> { |
| 123 | state |
| 124 | .runtime_threads |
| 125 | .thread_shell_manager(thread_id, create) |
| 126 | .await |
| 127 | .map_err(map_thread_err)? |
| 128 | .ok_or_else(|| ApiError::not_found(format!("thread {thread_id} has no jobs"))) |
| 129 | } |
| 130 | |
| 131 | /// `GET /v1/threads/{id}/jobs` — all jobs owned by one thread's manager: |
| 132 | /// model-launched, subagent-launched, and client-launched together. |
| 133 | pub(super) async fn list_thread_jobs( |
| 134 | State(state): State<RuntimeApiState>, |
| 135 | Path(thread_id): Path<String>, |
| 136 | ) -> Result<Json<JobListResponse>, ApiError> { |
| 137 | let Some(manager) = state |
| 138 | .runtime_threads |
| 139 | .thread_shell_manager(&thread_id, false) |
| 140 | .await |
| 141 | .map_err(map_thread_err)? |
| 142 | else { |
| 143 | return Ok(Json(JobListResponse { jobs: Vec::new() })); |
| 144 | }; |
| 145 | let jobs = tokio::task::spawn_blocking(move || { |
| 146 | let mut guard = manager.lock().unwrap_or_else(|e| e.into_inner()); |
| 147 | guard |
| 148 | .list_jobs() |
| 149 | .into_iter() |
| 150 | .map(|snapshot| JobView::new(snapshot, thread_id.clone(), &guard)) |
| 151 | .collect() |
| 152 | }) |
| 153 | .await |
| 154 | .map_err(|_| ApiError::internal("job listing failed"))?; |
| 155 | Ok(Json(JobListResponse { jobs })) |
| 156 | } |
| 157 | |
| 158 | #[derive(Deserialize)] |
| 159 | #[serde(deny_unknown_fields)] |
| 160 | pub(super) struct CreateJobRequest { |
| 161 | command: String, |
| 162 | /// Working directory. Omitted = the thread's workspace. |
| 163 | cwd: Option<String>, |
| 164 | /// Bounds the foreground-wait contract inside the manager; background |
| 165 | /// jobs are never killed at timeout. |
| 166 | timeout_ms: Option<u64>, |
| 167 | /// Run under a PTY: stderr merges into stdout and the command sees a |
| 168 | /// terminal. Required for interactive programs. |
| 169 | #[serde(default)] |
| 170 | tty: bool, |
| 171 | #[serde(default)] |
| 172 | env: HashMap<String, String>, |
| 173 | } |
| 174 | |
| 175 | #[derive(Debug, Serialize)] |
| 176 | pub(super) struct CreateJobResponse { |
| 177 | job: JobView, |
| 178 | } |
| 179 | |
| 180 | /// Resolve a job `cwd` against the thread workspace, never the server's cwd. |
| 181 | /// Outside trust mode the directory must stay inside the workspace once |
| 182 | /// symlinks resolve. This is stricter than the shell tool's `resolve_path`: |
| 183 | /// the route has no `workspace_follow_symlinks` or `/trust add` roots and no |
| 184 | /// `~` expansion, so a symlink leading out of the workspace is refused. |
| 185 | /// |
| 186 | /// The manager receives the canonical directory that was checked, not the |
| 187 | /// caller's spelling: a symlink retargeted between this check and the launch |
| 188 | /// cannot move the job elsewhere. The manager takes a `&str`, so a directory |
| 189 | /// whose name is not UTF-8 is refused rather than lossily renamed into a |
| 190 | /// different one. On Windows the `\\?\` verbatim prefix `canonicalize` adds |
| 191 | /// is dropped ([`plain_windows_path`]) because `cmd.exe` refuses it as a |
| 192 | /// current directory. Runs on `tokio::fs`, so no path resolution blocks a |
| 193 | /// runtime worker. |
| 194 | async fn resolve_job_cwd( |
| 195 | workspace: &std::path::Path, |
| 196 | raw: &str, |
| 197 | trust_mode: bool, |
| 198 | ) -> Result<String, ApiError> { |
| 199 | let requested = std::path::Path::new(raw); |
| 200 | let candidate = if requested.is_absolute() { |
| 201 | requested.to_path_buf() |
| 202 | } else { |
| 203 | workspace.join(requested) |
| 204 | }; |
| 205 | let not_a_directory = || ApiError::bad_request("cwd must be an existing directory"); |
| 206 | let resolved = tokio::fs::canonicalize(&candidate) |
| 207 | .await |
| 208 | .map_err(|_| not_a_directory())?; |
| 209 | if !tokio::fs::metadata(&resolved) |
| 210 | .await |
| 211 | .is_ok_and(|meta| meta.is_dir()) |
| 212 | { |
| 213 | return Err(not_a_directory()); |
| 214 | } |
| 215 | if !trust_mode { |
| 216 | let root = tokio::fs::canonicalize(workspace) |
| 217 | .await |
| 218 | .map_err(|_| ApiError::internal("thread workspace is unavailable"))?; |
| 219 | if !resolved.starts_with(&root) { |
| 220 | return Err(ApiError::forbidden( |
| 221 | "cwd must stay inside the thread workspace", |
| 222 | )); |
| 223 | } |
| 224 | } |
| 225 | let resolved = resolved |
| 226 | .into_os_string() |
| 227 | .into_string() |
| 228 | .map_err(|_| ApiError::bad_request("cwd must resolve to a UTF-8 path"))?; |
| 229 | Ok(if cfg!(windows) { |
| 230 | plain_windows_path(&resolved) |
| 231 | } else { |
| 232 | resolved |
| 233 | }) |
| 234 | } |
| 235 | |
| 236 | /// `\\?\C:\dir` -> `C:\dir` and `\\?\UNC\host\share` -> `\\host\share`; any |
| 237 | /// other spelling is returned unchanged. A string rule so every host tests it. |
| 238 | fn plain_windows_path(path: &str) -> String { |
| 239 | if let Some(rest) = path.strip_prefix(r"\\?\UNC\") { |
| 240 | return format!(r"\\{rest}"); |
| 241 | } |
| 242 | match path.strip_prefix(r"\\?\") { |
| 243 | Some(rest) |
| 244 | if rest.as_bytes().first().is_some_and(u8::is_ascii_alphabetic) |
| 245 | && rest.as_bytes().get(1) == Some(&b':') => |
| 246 | { |
| 247 | rest.to_string() |
| 248 | } |
| 249 | _ => path.to_string(), |
| 250 | } |
| 251 | } |
| 252 | |
| 253 | /// `POST /v1/threads/{id}/jobs` — launch a client-owned background job under |
| 254 | /// the thread's own sandbox policy. The client asking is the approval; the |
| 255 | /// thread's posture still bounds what the job may touch. |
| 256 | pub(super) async fn create_thread_job( |
| 257 | State(state): State<RuntimeApiState>, |
| 258 | Path(thread_id): Path<String>, |
| 259 | Json(request): Json<CreateJobRequest>, |
| 260 | ) -> Result<(StatusCode, Json<CreateJobResponse>), ApiError> { |
| 261 | let command = request.command.trim(); |
| 262 | if command.is_empty() { |
| 263 | return Err(ApiError::bad_request("command is required")); |
| 264 | } |
| 265 | if command.len() > COMMAND_MAX_BYTES { |
| 266 | return Err(ApiError::bad_request(format!( |
| 267 | "command must be at most {COMMAND_MAX_BYTES} bytes" |
| 268 | ))); |
| 269 | } |
| 270 | if request.env.len() > ENV_MAX_ENTRIES { |
| 271 | return Err(ApiError::bad_request(format!( |
| 272 | "env may carry at most {ENV_MAX_ENTRIES} entries" |
| 273 | ))); |
| 274 | } |
| 275 | for (key, value) in &request.env { |
| 276 | if key.len() > ENV_KEY_MAX_BYTES || key.contains(['=', '\0']) { |
| 277 | return Err(ApiError::bad_request("invalid env key")); |
| 278 | } |
| 279 | if value.len() > ENV_VALUE_MAX_BYTES || value.contains('\0') { |
| 280 | return Err(ApiError::bad_request("invalid env value")); |
| 281 | } |
| 282 | } |
| 283 | |
| 284 | let thread = state |
| 285 | .runtime_threads |
| 286 | .get_thread(&thread_id) |
| 287 | .await |
| 288 | .map_err(map_thread_err)?; |
| 289 | if !thread.allow_shell { |
| 290 | return Err(ApiError::forbidden( |
| 291 | "this thread does not allow shell commands", |
| 292 | )); |
| 293 | } |
| 294 | state |
| 295 | .runtime_threads |
| 296 | .validate_shell_access_policy( |
| 297 | &thread.workspace, |
| 298 | state.config_path.as_deref(), |
| 299 | state.config_profile.as_deref(), |
| 300 | ) |
| 301 | .await |
| 302 | .map_err(|error| ApiError::forbidden(error.to_string()))?; |
| 303 | // The manager resolves a relative `working_dir` against the server's own |
| 304 | // cwd, so hand it the directory resolved against the thread workspace. |
| 305 | let request_cwd = match request.cwd.as_deref() { |
| 306 | Some(cwd) => Some(resolve_job_cwd(&thread.workspace, cwd, thread.trust_mode).await?), |
| 307 | None => None, |
| 308 | }; |
| 309 | let policy = state |
| 310 | .runtime_threads |
| 311 | .thread_job_sandbox_policy(&thread) |
| 312 | .await; |
| 313 | let manager = thread_manager(&state, &thread_id, true).await?; |
| 314 | let scope = api_job_scope(&thread_id); |
| 315 | let request_timeout = request.timeout_ms; |
| 316 | let request_tty = request.tty; |
| 317 | let request_env = request.env; |
| 318 | let command = command.to_string(); |
| 319 | let job = tokio::task::spawn_blocking(move || -> Result<JobView, ApiError> { |
| 320 | let mut guard = manager.lock().unwrap_or_else(|e| e.into_inner()); |
| 321 | let result = guard |
| 322 | .execute_with_options_env_for_session( |
| 323 | &command, |
| 324 | request_cwd.as_deref(), |
| 325 | request_timeout.unwrap_or(120_000), |
| 326 | true, |
| 327 | None, |
| 328 | request_tty, |
| 329 | Some(policy), |
| 330 | request_env, |
| 331 | &scope, |
| 332 | ) |
| 333 | .map_err(|error| ApiError::internal(format!("job launch failed: {error}")))?; |
| 334 | let task_id = result.task_id.clone().unwrap_or_default(); |
| 335 | let snapshot = guard |
| 336 | .inspect_job(&task_id) |
| 337 | .map_err(|error| { |
| 338 | ApiError::internal(format!("job launched but is not tracked: {error}")) |
| 339 | })? |
| 340 | .snapshot; |
| 341 | Ok(JobView::new(snapshot, thread_id, &guard)) |
| 342 | }) |
| 343 | .await |
| 344 | .map_err(|_| ApiError::internal("job launch failed"))??; |
| 345 | Ok((StatusCode::CREATED, Json(CreateJobResponse { job }))) |
| 346 | } |
| 347 | |
| 348 | #[derive(Debug, Serialize)] |
| 349 | pub(super) struct JobDetailResponse { |
| 350 | job: JobView, |
| 351 | stdout_tail: String, |
| 352 | stderr_tail: String, |
| 353 | } |
| 354 | |
| 355 | /// `GET /v1/threads/{id}/jobs/{job_id}` — snapshot plus the retained output |
| 356 | /// tails. For the full stream, follow `output` with a cursor instead. |
| 357 | pub(super) async fn get_thread_job( |
| 358 | State(state): State<RuntimeApiState>, |
| 359 | Path((thread_id, job_id)): Path<(String, String)>, |
| 360 | ) -> Result<Json<JobDetailResponse>, ApiError> { |
| 361 | if job_id.len() > JOB_ID_MAX_BYTES { |
| 362 | return Err(ApiError::not_found("job not found")); |
| 363 | } |
| 364 | let manager = thread_manager(&state, &thread_id, false).await?; |
| 365 | tokio::task::spawn_blocking(move || { |
| 366 | let mut guard = manager.lock().unwrap_or_else(|e| e.into_inner()); |
| 367 | let detail = guard.inspect_job(&job_id).map_err(map_job_err)?; |
| 368 | Ok(Json(JobDetailResponse { |
| 369 | job: JobView::new(detail.snapshot, thread_id, &guard), |
| 370 | stdout_tail: detail.stdout, |
| 371 | stderr_tail: detail.stderr, |
| 372 | })) |
| 373 | }) |
| 374 | .await |
| 375 | .map_err(|_| ApiError::internal("job inspect failed"))? |
| 376 | } |
| 377 | |
| 378 | #[derive(Deserialize)] |
| 379 | #[serde(deny_unknown_fields)] |
| 380 | pub(super) struct JobOutputQuery { |
| 381 | /// `stdout` (default) or `stderr`; PTY jobs merge stderr into stdout. |
| 382 | #[serde(default)] |
| 383 | stream: Option<String>, |
| 384 | /// Absolute byte offset into the stream's lifetime output. |
| 385 | #[serde(default)] |
| 386 | cursor: Option<usize>, |
| 387 | /// Per-request byte ceiling, default 64 KiB, max 512 KiB. |
| 388 | #[serde(default)] |
| 389 | max_bytes: Option<usize>, |
| 390 | /// Long-poll bound for new bytes on a running job, max 30s. |
| 391 | #[serde(default)] |
| 392 | wait_ms: Option<u64>, |
| 393 | /// `base64` (default, exact bytes) or `text` (lossy UTF-8). |
| 394 | #[serde(default)] |
| 395 | format: Option<String>, |
| 396 | } |
| 397 | |
| 398 | #[derive(Debug, Serialize)] |
| 399 | pub(super) struct JobOutputResponse { |
| 400 | job_id: String, |
| 401 | stream: &'static str, |
| 402 | /// Absolute offset of `data[0]`; exceeds `cursor` when the bounded buffer |
| 403 | /// already discarded that prefix (`dropped` reports the cutoff). |
| 404 | offset: usize, |
| 405 | /// Next cursor: pass it back to continue the stream. |
| 406 | next_cursor: usize, |
| 407 | /// Total bytes the stream has produced, including discarded bytes. |
| 408 | total: usize, |
| 409 | /// Leading bytes permanently discarded by the in-flight bound. |
| 410 | dropped: usize, |
| 411 | encoding: &'static str, |
| 412 | data: String, |
| 413 | status: ShellStatus, |
| 414 | exit_code: Option<i64>, |
| 415 | /// Terminal status and no bytes remain past `next_cursor`. |
| 416 | done: bool, |
| 417 | } |
| 418 | |
| 419 | /// `GET /v1/threads/{id}/jobs/{job_id}/output` — the resumable byte stream. |
| 420 | /// Reads are non-consuming: several clients may hold independent cursors, and |
| 421 | /// polling here never steals output from the engine's own delta consumer. |
| 422 | pub(super) async fn get_thread_job_output( |
| 423 | State(state): State<RuntimeApiState>, |
| 424 | Path((thread_id, job_id)): Path<(String, String)>, |
| 425 | Query(query): Query<JobOutputQuery>, |
| 426 | ) -> Result<Json<JobOutputResponse>, ApiError> { |
| 427 | if job_id.len() > JOB_ID_MAX_BYTES { |
| 428 | return Err(ApiError::not_found("job not found")); |
| 429 | } |
| 430 | let (stream, stream_name) = match query.stream.as_deref().unwrap_or("stdout") { |
| 431 | "stdout" => (ShellOutputStream::Stdout, "stdout"), |
| 432 | "stderr" => (ShellOutputStream::Stderr, "stderr"), |
| 433 | _ => return Err(ApiError::bad_request("stream must be stdout or stderr")), |
| 434 | }; |
| 435 | let cursor = query.cursor.unwrap_or(0); |
| 436 | let max_bytes = query.max_bytes.unwrap_or(OUTPUT_CHUNK_DEFAULT); |
| 437 | if !(1..=OUTPUT_CHUNK_MAX).contains(&max_bytes) { |
| 438 | return Err(ApiError::bad_request(format!( |
| 439 | "max_bytes must be between 1 and {OUTPUT_CHUNK_MAX}" |
| 440 | ))); |
| 441 | } |
| 442 | let wait_ms = query.wait_ms.unwrap_or(0).min(OUTPUT_WAIT_MAX_MS); |
| 443 | let format = query.format.as_deref().unwrap_or("base64"); |
| 444 | if !matches!(format, "base64" | "text") { |
| 445 | return Err(ApiError::bad_request("format must be base64 or text")); |
| 446 | } |
| 447 | let manager = thread_manager(&state, &thread_id, false).await?; |
| 448 | let chunk = tokio::task::spawn_blocking({ |
| 449 | let job_id = job_id.clone(); |
| 450 | move || -> Result<ShellOutputChunk, ApiError> { |
| 451 | // Wait between non-consuming snapshots, never while owning the |
| 452 | // thread's shared ShellManager. Input, resize, and kill stay live. |
| 453 | let deadline = std::time::Instant::now() + std::time::Duration::from_millis(wait_ms); |
| 454 | loop { |
| 455 | let chunk = { |
| 456 | let mut guard = manager.lock().unwrap_or_else(|e| e.into_inner()); |
| 457 | guard |
| 458 | .read_output_chunk(&job_id, stream, cursor, max_bytes, 0) |
| 459 | .map_err(map_job_err)? |
| 460 | }; |
| 461 | let remaining = deadline.saturating_duration_since(std::time::Instant::now()); |
| 462 | if chunk.total > cursor |
| 463 | || chunk.status != ShellStatus::Running |
| 464 | || remaining.is_zero() |
| 465 | { |
| 466 | break Ok(chunk); |
| 467 | } |
| 468 | std::thread::sleep(remaining.min(std::time::Duration::from_millis(50))); |
| 469 | } |
| 470 | } |
| 471 | }) |
| 472 | .await |
| 473 | .map_err(|_| ApiError::internal("job output read failed"))??; |
| 474 | Ok(Json(encode_chunk(&job_id, stream_name, chunk, format))) |
| 475 | } |
| 476 | |
| 477 | fn encode_chunk( |
| 478 | job_id: &str, |
| 479 | stream_name: &'static str, |
| 480 | chunk: ShellOutputChunk, |
| 481 | format: &str, |
| 482 | ) -> JobOutputResponse { |
| 483 | let (encoding, data) = match format { |
| 484 | "text" => ("utf-8", String::from_utf8_lossy(&chunk.bytes).into_owned()), |
| 485 | _ => ( |
| 486 | "base64", |
| 487 | base64::engine::general_purpose::STANDARD.encode(&chunk.bytes), |
| 488 | ), |
| 489 | }; |
| 490 | let done = chunk.status != ShellStatus::Running && chunk.next_offset >= chunk.total; |
| 491 | JobOutputResponse { |
| 492 | job_id: job_id.to_string(), |
| 493 | stream: stream_name, |
| 494 | offset: chunk.offset, |
| 495 | next_cursor: chunk.next_offset, |
| 496 | total: chunk.total, |
| 497 | dropped: chunk.dropped, |
| 498 | encoding, |
| 499 | data, |
| 500 | status: chunk.status, |
| 501 | exit_code: chunk.exit_code, |
| 502 | done, |
| 503 | } |
| 504 | } |
| 505 | |
| 506 | #[derive(Deserialize)] |
| 507 | #[serde(deny_unknown_fields)] |
| 508 | pub(super) struct JobStdinRequest { |
| 509 | /// UTF-8 text (default) or base64 for arbitrary bytes. |
| 510 | data: String, |
| 511 | #[serde(default)] |
| 512 | encoding: Option<String>, |
| 513 | /// Close stdin after writing (EOF). |
| 514 | #[serde(default)] |
| 515 | close: bool, |
| 516 | } |
| 517 | |
| 518 | /// `POST /v1/threads/{id}/jobs/{job_id}/stdin` — write to a running job's |
| 519 | /// stdin. Works for PTY and piped jobs alike. |
| 520 | pub(super) async fn write_thread_job_stdin( |
| 521 | State(state): State<RuntimeApiState>, |
| 522 | Path((thread_id, job_id)): Path<(String, String)>, |
| 523 | Json(request): Json<JobStdinRequest>, |
| 524 | ) -> Result<StatusCode, ApiError> { |
| 525 | if job_id.len() > JOB_ID_MAX_BYTES { |
| 526 | return Err(ApiError::not_found("job not found")); |
| 527 | } |
| 528 | let input = match request.encoding.as_deref().unwrap_or("utf-8") { |
| 529 | "utf-8" => { |
| 530 | if request.data.len() > STDIN_MAX_BYTES { |
| 531 | return Err(ApiError::bad_request(format!( |
| 532 | "data must be at most {STDIN_MAX_BYTES} bytes" |
| 533 | ))); |
| 534 | } |
| 535 | request.data.into_bytes() |
| 536 | } |
| 537 | "base64" => { |
| 538 | if request.data.len() > STDIN_MAX_BYTES * 2 { |
| 539 | return Err(ApiError::bad_request("data exceeds the stdin limit")); |
| 540 | } |
| 541 | let bytes = base64::engine::general_purpose::STANDARD |
| 542 | .decode(&request.data) |
| 543 | .map_err(|_| ApiError::bad_request("data is not valid base64"))?; |
| 544 | if bytes.len() > STDIN_MAX_BYTES { |
| 545 | return Err(ApiError::bad_request(format!( |
| 546 | "data must be at most {STDIN_MAX_BYTES} decoded bytes" |
| 547 | ))); |
| 548 | } |
| 549 | bytes |
| 550 | } |
| 551 | _ => return Err(ApiError::bad_request("encoding must be utf-8 or base64")), |
| 552 | }; |
| 553 | let close = request.close; |
| 554 | let manager = thread_manager(&state, &thread_id, false).await?; |
| 555 | tokio::task::spawn_blocking(move || { |
| 556 | let mut guard = manager.lock().unwrap_or_else(|e| e.into_inner()); |
| 557 | guard |
| 558 | .write_stdin_bytes(&job_id, &input, close) |
| 559 | .map_err(map_job_err) |
| 560 | }) |
| 561 | .await |
| 562 | .map_err(|_| ApiError::internal("job stdin write failed"))??; |
| 563 | Ok(StatusCode::NO_CONTENT) |
| 564 | } |
| 565 | |
| 566 | #[derive(Debug, Serialize)] |
| 567 | pub(super) struct KillJobResponse { |
| 568 | job: JobView, |
| 569 | result: ShellResult, |
| 570 | } |
| 571 | |
| 572 | /// `POST /v1/threads/{id}/jobs/{job_id}/kill` — bounded SIGTERM → SIGKILL |
| 573 | /// escalation on the whole process group; the final snapshot rides along. |
| 574 | pub(super) async fn kill_thread_job( |
| 575 | State(state): State<RuntimeApiState>, |
| 576 | Path((thread_id, job_id)): Path<(String, String)>, |
| 577 | ) -> Result<Json<KillJobResponse>, ApiError> { |
| 578 | if job_id.len() > JOB_ID_MAX_BYTES { |
| 579 | return Err(ApiError::not_found("job not found")); |
| 580 | } |
| 581 | let manager = thread_manager(&state, &thread_id, false).await?; |
| 582 | tokio::task::spawn_blocking(move || { |
| 583 | let mut guard = manager.lock().unwrap_or_else(|e| e.into_inner()); |
| 584 | let result = guard.kill(&job_id).map_err(map_job_err)?; |
| 585 | let snapshot = guard.inspect_job(&job_id).map_err(map_job_err)?.snapshot; |
| 586 | Ok(Json(KillJobResponse { |
| 587 | job: JobView::new(snapshot, thread_id, &guard), |
| 588 | result, |
| 589 | })) |
| 590 | }) |
| 591 | .await |
| 592 | .map_err(|_| ApiError::internal("job kill failed"))? |
| 593 | } |
| 594 | |
| 595 | /// `POST /v1/threads/{id}/jobs/{job_id}/resize` — resize this existing PTY. |
| 596 | pub(super) async fn resize_thread_job( |
| 597 | State(state): State<RuntimeApiState>, |
| 598 | Path((thread_id, job_id)): Path<(String, String)>, |
| 599 | Json(size): Json<PtyDimensions>, |
| 600 | ) -> Result<Json<JobDetailResponse>, ApiError> { |
| 601 | let size = size |
| 602 | .validate() |
| 603 | .map_err(|error| ApiError::bad_request(error.to_string()))?; |
| 604 | if job_id.len() > JOB_ID_MAX_BYTES { |
| 605 | return Err(ApiError::not_found("job not found")); |
| 606 | } |
| 607 | let manager = thread_manager(&state, &thread_id, false).await?; |
| 608 | tokio::task::spawn_blocking(move || { |
| 609 | let mut guard = manager.lock().unwrap_or_else(|e| e.into_inner()); |
| 610 | let detail = guard.inspect_job(&job_id).map_err(map_job_err)?; |
| 611 | if detail.snapshot.stale || detail.snapshot.status != ShellStatus::Running { |
| 612 | return Err(ApiError { |
| 613 | status: StatusCode::CONFLICT, |
| 614 | message: "This job is no longer running".into(), |
| 615 | code: None, |
| 616 | }); |
| 617 | } |
| 618 | if guard.job_terminal_size(&job_id).is_none() { |
| 619 | return Err(ApiError::bad_request("This job is not a PTY")); |
| 620 | } |
| 621 | guard.resize_pty(&job_id, size).map_err(map_job_err)?; |
| 622 | let detail = guard.inspect_job(&job_id).map_err(map_job_err)?; |
| 623 | Ok(Json(JobDetailResponse { |
| 624 | job: JobView::new(detail.snapshot, thread_id, &guard), |
| 625 | stdout_tail: detail.stdout, |
| 626 | stderr_tail: detail.stderr, |
| 627 | })) |
| 628 | }) |
| 629 | .await |
| 630 | .map_err(|_| ApiError::internal("PTY resize failed"))? |
| 631 | } |
| 632 | |
| 633 | #[cfg(test)] |
| 634 | mod tests { |
| 635 | use super::*; |
| 636 | |
| 637 | fn canonical(path: &std::path::Path) -> String { |
| 638 | let canonical = path.canonicalize().unwrap().to_string_lossy().into_owned(); |
| 639 | if cfg!(windows) { |
| 640 | plain_windows_path(&canonical) |
| 641 | } else { |
| 642 | canonical |
| 643 | } |
| 644 | } |
| 645 | |
| 646 | #[tokio::test] |
| 647 | async fn job_cwd_resolves_against_the_thread_workspace() { |
| 648 | let tmp = tempfile::tempdir().unwrap(); |
| 649 | let workspace = tmp.path().join("workspace"); |
| 650 | std::fs::create_dir_all(workspace.join("packages/app")).unwrap(); |
| 651 | let outside = tmp.path().join("outside"); |
| 652 | std::fs::create_dir_all(&outside).unwrap(); |
| 653 | |
| 654 | // Relative and absolute requests both come back as the canonical |
| 655 | // directory that was checked, never the caller's spelling. |
| 656 | let resolved = resolve_job_cwd(&workspace, "packages/app", false) |
| 657 | .await |
| 658 | .unwrap(); |
| 659 | assert!(std::path::Path::new(&resolved).is_absolute(), "{resolved}"); |
| 660 | assert_eq!(resolved, canonical(&workspace.join("packages/app"))); |
| 661 | let absolute = workspace.join("packages"); |
| 662 | assert_eq!( |
| 663 | resolve_job_cwd(&workspace, &absolute.to_string_lossy(), false) |
| 664 | .await |
| 665 | .unwrap(), |
| 666 | canonical(&absolute) |
| 667 | ); |
| 668 | |
| 669 | let missing = resolve_job_cwd(&workspace, "packages/none", false) |
| 670 | .await |
| 671 | .unwrap_err(); |
| 672 | assert_eq!(missing.status, StatusCode::BAD_REQUEST); |
| 673 | let outside_raw = outside.to_string_lossy().into_owned(); |
| 674 | for raw in [outside_raw.as_str(), "../outside"] { |
| 675 | let error = resolve_job_cwd(&workspace, raw, false).await.unwrap_err(); |
| 676 | assert_eq!(error.status, StatusCode::FORBIDDEN, "{raw}"); |
| 677 | } |
| 678 | #[cfg(unix)] |
| 679 | { |
| 680 | std::os::unix::fs::symlink(&outside, workspace.join("escape")).unwrap(); |
| 681 | let error = resolve_job_cwd(&workspace, "escape", false) |
| 682 | .await |
| 683 | .unwrap_err(); |
| 684 | assert_eq!(error.status, StatusCode::FORBIDDEN); |
| 685 | } |
| 686 | assert_eq!( |
| 687 | resolve_job_cwd(&workspace, "../outside", true) |
| 688 | .await |
| 689 | .unwrap(), |
| 690 | canonical(&outside) |
| 691 | ); |
| 692 | } |
| 693 | |
| 694 | /// The cwd handed to the launcher is the directory whose containment was |
| 695 | /// checked: retargeting the requested symlink afterwards does not move |
| 696 | /// the child out of the workspace. |
| 697 | #[cfg(unix)] |
| 698 | #[tokio::test] |
| 699 | async fn job_cwd_launches_in_the_checked_directory_after_a_retarget() { |
| 700 | let tmp = tempfile::tempdir().unwrap(); |
| 701 | let workspace = tmp.path().join("workspace"); |
| 702 | let inside = workspace.join("inside"); |
| 703 | std::fs::create_dir_all(&inside).unwrap(); |
| 704 | let outside = tmp.path().join("outside"); |
| 705 | std::fs::create_dir_all(&outside).unwrap(); |
| 706 | let link = workspace.join("link"); |
| 707 | std::os::unix::fs::symlink(&inside, &link).unwrap(); |
| 708 | |
| 709 | let resolved = resolve_job_cwd(&workspace, "link", false).await.unwrap(); |
| 710 | assert_eq!(resolved, canonical(&inside)); |
| 711 | |
| 712 | std::fs::remove_file(&link).unwrap(); |
| 713 | std::os::unix::fs::symlink(&outside, &link).unwrap(); |
| 714 | let output = std::process::Command::new("pwd") |
| 715 | .arg("-P") |
| 716 | .current_dir(&resolved) |
| 717 | .output() |
| 718 | .unwrap(); |
| 719 | assert!(output.status.success(), "{output:?}"); |
| 720 | assert_eq!( |
| 721 | String::from_utf8_lossy(&output.stdout).trim(), |
| 722 | canonical(&inside) |
| 723 | ); |
| 724 | } |
| 725 | |
| 726 | /// A directory whose name is not UTF-8 is refused before launch rather |
| 727 | /// than lossily renamed into a `U+FFFD` sibling. Linux only: macOS and |
| 728 | /// Windows filesystems do not store non-UTF-8 names. |
| 729 | #[cfg(target_os = "linux")] |
| 730 | #[tokio::test] |
| 731 | async fn job_cwd_refuses_a_directory_name_that_is_not_utf8() { |
| 732 | use std::os::unix::ffi::OsStrExt as _; |
| 733 | let tmp = tempfile::tempdir().unwrap(); |
| 734 | let workspace = tmp.path().join("workspace"); |
| 735 | let raw_name = std::ffi::OsStr::from_bytes(b"dir-\xff"); |
| 736 | std::fs::create_dir_all(workspace.join(raw_name)).unwrap(); |
| 737 | std::fs::create_dir_all(workspace.join("dir-\u{fffd}")).unwrap(); |
| 738 | std::os::unix::fs::symlink(workspace.join(raw_name), workspace.join("alias")).unwrap(); |
| 739 | |
| 740 | let error = resolve_job_cwd(&workspace, "alias", false) |
| 741 | .await |
| 742 | .unwrap_err(); |
| 743 | assert_eq!(error.status, StatusCode::BAD_REQUEST); |
| 744 | } |
| 745 | |
| 746 | #[test] |
| 747 | fn plain_windows_path_drops_only_the_verbatim_prefix() { |
| 748 | assert_eq!(plain_windows_path(r"\\?\C:\ws\app"), r"C:\ws\app"); |
| 749 | assert_eq!( |
| 750 | plain_windows_path(r"\\?\UNC\host\share\ws"), |
| 751 | r"\\host\share\ws" |
| 752 | ); |
| 753 | for unchanged in [r"C:\ws", r"\\host\share", r"\\?\Volume{0}\ws", "/tmp/ws"] { |
| 754 | assert_eq!(plain_windows_path(unchanged), unchanged); |
| 755 | } |
| 756 | } |
| 757 | } |
| 758 |