use std::sync::Arc; use std::os::unix::fs::PermissionsExt; use tokio::process::Command; use crate::error::{AppError, AppResult}; use crate::state::AppState; const NONROOT_UID: u64 = 65534; const MAX_STDOUT: usize = 64_000; const MAX_STDERR: usize = 16_000; #[derive(Debug)] pub struct ContainerResult { pub status: Option, pub stdout: String, pub stderr: String, } impl ContainerResult { pub fn ok(&self) -> bool { self.status == Some(0) } } /// Docker run arguments (after `docker`). Isolation contract: /// no caps, no privilege escalation, non-root uid 65534, read-only root /// filesystem, small tmpfs, pids/memory/cpu limits, network only when the /// task requires the data provider; backtest always runs with --network=none. pub fn docker_args( network: bool, mounts: &[(String, String, bool)], image: &str, name: &str, args: &[String], ) -> Vec { let mut a: Vec = [ "run", "--rm", // Signals reaching the runner CLI must never be proxied into the // container: on service stop the container is expected to die via the // app's exact-name cleanup (cancel/timeout/startup), not via a CLI // relay; a proxying CLI was observed lingering under systemd // final-sigterm (TimeoutStopSec exhaustion). "--sig-proxy=false", "--name", name, "--cap-drop=ALL", "--security-opt=no-new-privileges", CONCAT_USER, "--read-only", "--tmpfs=/tmp:rw,size=256m,mode=1777", "--pids-limit=128", "--memory=2g", "--cpus=2", ] .iter() .map(|s| s.to_string()) .collect(); // explicit network choice; only data fetch/search uses the bridge a.push(if network { "--network=bridge".into() } else { "--network=none".into() }); for (src, dst, ro) in mounts { a.push("-v".into()); a.push(format!("{}:{}{}", src, dst, if *ro { ":ro" } else { "" })); } a.push(image.to_string()); a.extend_from_slice(args); a } async fn drain_bounded(rd: R, max: usize) -> String where R: tokio::io::AsyncRead + Unpin, { use tokio::io::AsyncReadExt; let mut buf = Vec::with_capacity(1024); let mut chunk = [0u8; 8192]; let mut reader = rd; loop { match reader.read(&mut chunk).await { Ok(0) => break, Ok(n) => { // drain everything, but keep only the tail-relevant bounded prefix if buf.len() < max { let take = n.min(max - buf.len()); buf.extend_from_slice(&chunk[..take]); } if buf.len() >= max { // continue draining the pipe without buffering the rest let mut sink = [0u8; 8192]; loop { match reader.read(&mut sink).await { Ok(0) | Err(_) => break, Ok(_) => {} } } break; } } Err(_) => break, } } truncate(&String::from_utf8_lossy(&buf), max) } async fn execute_docker(full: &[String], name: &str, timeout_secs: u64) -> AppResult { execute_docker_named("docker", full, name, timeout_secs).await } async fn execute_docker_named( docker_bin: &str, full: &[String], name: &str, timeout_secs: u64, ) -> AppResult { let mut child = Command::new(docker_bin) .args(full) // Kill the runner process when the future that owns it is dropped. A // graceful server stop drops the jobs task; without this the docker // CLI child would linger inside the systemd cgroup and block the // unit stop until TimeoutStopSec forced SIGKILL. Bounded lifetime. .kill_on_drop(true) .stdout(std::process::Stdio::piped()) .stderr(std::process::Stdio::piped()) .spawn() .map_err(|e| AppError::internal(format!("failed to spawn worker container: {e}")))?; let stdout = child.stdout.take().expect("stdout piped"); let stderr = child.stderr.take().expect("stderr piped"); let stdout_task = tokio::spawn(drain_bounded(stdout, MAX_STDOUT)); let stderr_task = tokio::spawn(drain_bounded(stderr, MAX_STDERR)); let wait_res = tokio::time::timeout( std::time::Duration::from_secs(timeout_secs), child.wait(), ) .await; let status = match wait_res { Ok(Ok(st)) => st.code(), Ok(Err(e)) => { return Err(AppError::internal(format!("worker process error: {e}")).with_code("runner_failed")) } Err(_) => { // Timeout: kill the specific container by name so user code cannot // linger; then reap the docker client process. kill_container(name).await; let _ = child.wait().await; return Err(AppError::internal(format!( "worker container timed out after {timeout_secs}s and was killed: {name}" )) .with_code("runner_timeout")); } }; Ok(ContainerResult { status, stdout: stdout_task.await.unwrap_or_default(), stderr: stderr_task.await.unwrap_or_default(), }) } /// Kill and remove the named container. Returns true when docker succeeded. pub async fn cancel_container(name: &str) -> bool { kill_container(name).await } async fn kill_container(name: &str) -> bool { Command::new("docker") .args(["kill", name]) .stdout(std::process::Stdio::null()) .stderr(std::process::Stdio::null()) .output() .await .ok(); Command::new("docker") .args(["rm", "-f", name]) .stdout(std::process::Stdio::null()) .stderr(std::process::Stdio::null()) .output() .await .map(|o| o.status.success()) .unwrap_or(false) } pub async fn docker_available() -> bool { tokio::process::Command::new("docker") .args(["version", "--format", "ok"]) .stdout(std::process::Stdio::null()) .stderr(std::process::Stdio::null()) .output() .await .map(|o| o.status.success()) .unwrap_or(false) } /// Mounts need world permissions: the container runs as uid 65534 while host /// ownership is the server user. Best effort only. fn prepare_mounts(mounts: &[(String, String, bool)]) { for (src, _dst, ro) in mounts { let p = std::path::Path::new(src); if !p.is_dir() { continue; } let mode = if *ro { 0o755 } else { 0o777 }; let _ = std::fs::set_permissions(p, std::fs::Permissions::from_mode(mode)); // Files inside ro input dirs must be world readable; output files are // written by the container with its umask. if *ro { if let Ok(rd) = std::fs::read_dir(p) { for e in rd.flatten() { let fmode = if e.path().is_file() { std::fs::Permissions::from_mode(0o644) } else { std::fs::Permissions::from_mode(0o755) }; let _ = std::fs::set_permissions(e.path(), fmode); } } } } } /// Run the worker image with a fixed container name so cancel maps to one /// specific container id (never a global prune). pub async fn run_named( cx: &Arc, network: bool, mounts: &[(String, String, bool)], args: &[String], name: &str, timeout_secs: u64, ) -> AppResult { let cfg = &cx.cfg; prepare_mounts(mounts); let full = docker_args(network, mounts, &cfg.worker_image, name, args); execute_docker(&full, name, timeout_secs).await } /// Instrument catalog search through the worker container (network enabled). /// Returns the JSON array printed by `worker.main search`; failures are real /// errors, never an empty fake success. Results are cached briefly per query so /// repeated keystrokes reuse the actual provider identity (full item payloads). pub async fn search_instruments( cx: &Arc, query: &str, limit: i64, ) -> AppResult> { let query = query.trim().to_string(); if query.is_empty() { return Ok(Vec::new()); } let limit = if (1..=100).contains(&limit) { limit } else { 50 }; let cache_key = format!("q={query}&limit={limit}"); if let Some(items) = catalog_cache_get(&cache_key) { return Ok(items); } let args = vec![ "python".into(), "-m".into(), "worker.main".into(), "search".into(), "--query".into(), query, "--limit".into(), limit.to_string(), ]; // The worker enforces <=4s per HTTP source; the container including startup // is bounded here. This is the outer bound for the whole search round trip. let res = run_named(cx, true, &[], &args, &random_name(), 12).await?; if !res.ok() { return Err(AppError::bad("search_failed", truncate(&res.stderr, 500))); } // The contract is a single JSON object envelope on stdout: // {"status":"ready"|"failed","items":[...],"error":{...},...} // A `failed` status is surfaced as an error, never as an empty success. let envelope: serde_json::Value = serde_json::from_str(res.stdout.trim()) .map_err(|e| AppError::internal(format!("invalid search envelope: {e}")))?; let status = envelope .get("status") .and_then(|v| v.as_str()) .ok_or_else(|| AppError::internal("invalid search envelope: missing status"))?; if status != "ready" { let err = envelope.get("error").cloned().unwrap_or(serde_json::Value::Null); let code = err .get("code") .and_then(|v| v.as_str()) .unwrap_or("search_unavailable"); let message = err .get("message") .and_then(|v| v.as_str()) .unwrap_or("instrument search providers unavailable"); // provider error codes are dynamic; they ride in the message so the // HTTP layer keeps static error codes return Err(AppError::bad( "search_unavailable", format!("[{code}] {message}"), )); } let items: Vec = serde_json::from_value( envelope.get("items").cloned().unwrap_or(serde_json::Value::Null), ) .map_err(|e| AppError::internal(format!("invalid search envelope items: {e}")))?; catalog_cache_put(&cache_key, &items); Ok(items) } /// Small in-process TTL cache for catalog search results (provider identity). const CATALOG_TTL_SECS: u64 = 300; const CATALOG_MAX_ENTRIES: usize = 128; fn catalog_cache() -> &'static tokio::sync::Mutex)>> { static MAP: std::sync::OnceLock)>>> = std::sync::OnceLock::new(); MAP.get_or_init(|| tokio::sync::Mutex::new(std::collections::HashMap::new())) } fn catalog_cache_get(key: &str) -> Option> { // Instant checks must not block behind stdio work; try_lock is fine here. let map = catalog_cache().try_lock().ok()?; let (at, items) = map.get(key)?; if at.elapsed() < std::time::Duration::from_secs(CATALOG_TTL_SECS) { Some(items.clone()) } else { None } } fn catalog_cache_put(key: &str, items: &[serde_json::Value]) { if let Ok(mut map) = catalog_cache().try_lock() { if map.len() >= CATALOG_MAX_ENTRIES { map.clear(); } map.insert(key.to_string(), (std::time::Instant::now(), items.to_vec())); } } fn random_name() -> String { use rand::RngCore; let mut buf = [0u8; 8]; rand::rngs::OsRng.fill_bytes(&mut buf); format!("sl-run-{}", hex::encode(buf)) } pub fn truncate(s: &str, n: usize) -> String { s.chars().take(n).collect() } const CONCAT_USER: &str = "--user=65534:65534"; #[cfg(test)] mod tests { use super::*; #[test] fn run_cli_never_proxies_signals_and_names_the_worker() { // Cause-establishing regression: systemd restart stall happened because // the runner CLI swallowed SIGTERM while relaying to the container. // The CLI (and every other container invokation) must use sig-proxy=false. for network in [true, false] { let full = docker_args( network, &[("/job/output".into(), "/output".into(), false)], "strategy-lab-worker:local", "sl-run-t", &["python".into()], ); let s = full.join(" "); assert!(s.contains("--sig-proxy=false"), "{s}"); assert!(s.contains("--rm --sig-proxy=false --name sl-run-t")); assert!(s.contains(format!("--user={NONROOT_UID}:").as_str())); } } #[tokio::test] async fn runner_child_is_reaped_when_the_job_future_is_dropped() { // Focused cause test: a runner child that ignores SIGTERM must not // outlive its owning future (systemd restart stall cause). We emulate // with a stub runner in a UNIQUE tempdir that ignores SIGTERM and then // `exec`s into sleep so the stub PID IS the sleep process: kill_on_drop // removes it entirely, leaving no grandchild orphan. Wait supervision // uses exact PID checks (`/proc/`), never process-name scans. let td = tempfile::tempdir_in("/tmp/opencode").unwrap(); let stub_path = td.path().join("runner-stub.sh"); let pid_file = td.path().join("stub.pid"); std::fs::write(&stub_path, format!( "#!/bin/sh\ntrap '' TERM INT\necho $$ > {}\nexec sleep 500\n", pid_file.display() )).unwrap(); { use std::os::unix::fs::PermissionsExt; std::fs::set_permissions(&stub_path, std::fs::Permissions::from_mode(0o755)).unwrap(); } let stub = stub_path.clone(); let task = tokio::spawn(async move { let _res = execute_docker_named( stub.to_str().unwrap(), &["long-running-stub".into()], "sl-run-stub", 5, ).await; }); // wait until the stub published its exact PID let mut stub_pid: Option = None; for _ in 0..50 { if let Ok(s) = std::fs::read_to_string(&pid_file) { stub_pid = s.trim().parse().ok(); } if stub_pid.is_some() { break; } tokio::time::sleep(std::time::Duration::from_millis(100)).await; } let stub_pid: i32 = match stub_pid { Some(p) => p, None => panic!("stub never published its PID; test setup broken"), }; let proc_dir = format!("/proc/{stub_pid}"); // sanity: stub alive; and after `exec` it IS the sleep grandchild assert!(std::path::Path::new(&proc_dir).exists(), "stub pid {stub_pid} must be alive before the drop"); let cmd = std::fs::read_to_string(format!("{proc_dir}/cmdline")).unwrap_or_default(); assert!(cmd.contains("sleep"), "exec replace failed; test would leave an orphan: {cmd:?}"); task.abort(); // drops the execute_docker future mid-flight -> kill_on_drop -> SIGKILL // exact-PID supervision: gone == /proc/ has vanished let deadline = std::time::Instant::now() + std::time::Duration::from_secs(3); let mut gone = false; while std::time::Instant::now() < deadline { if !std::path::Path::new(&proc_dir).exists() { gone = true; break; } tokio::time::sleep(std::time::Duration::from_millis(100)).await; } assert!(gone, "runner stub pid {stub_pid} was not reaped when its future was dropped"); // no grandchild either: the exec'd sleep adopted the same PID, then was // SIGKILLed with the rest; nothing named-scan was used. let _ = stub_path; // file removed with the unique tempdir at scope end } fn mounts() -> Vec<(String, String, bool)> { vec![ ("/job/input".into(), "/input".into(), true), ("/job/output".into(), "/output".into(), false), ("/data/objects/x.csv".into(), "/data/x.csv".into(), true), ] } #[test] fn backtest_container_flags_are_bounded_isolated_nonroot() { let args = vec!["python".to_string(), "-m".to_string(), "worker.main".to_string()]; let full = docker_args(false, &mounts(), "strategy-lab-worker:local", "sl-run-x", &args); let s = full.join(" "); assert!(s.contains("strategy-lab-worker:local")); assert!(s.contains("--network=none"), "runner must not use network: {s}"); assert!(s.contains(format!("--user={NONROOT_UID}:").as_str()), "{s}"); assert!(s.contains("--cap-drop=ALL")); assert!(s.contains("--security-opt=no-new-privileges")); assert!(s.contains("--read-only")); assert!(s.contains("--pids-limit=128")); assert!(s.contains("--memory=2g")); assert!(s.contains("--cpus=2")); assert!(s.contains("--tmpfs=/tmp:")); assert!(!s.contains("/var/run/docker.sock"), "no docker socket in worker: {s}"); // mount directions preserved assert!(s.contains("/job/input:/input:ro")); assert!(s.contains("/job/output:/output")); assert!(s.contains("/data/objects/x.csv:/data/x.csv:ro")); } #[test] fn fetch_container_has_network_and_same_isolation() { let full = docker_args(true, &[], "img", "sl-fetch-1", &["python".into()]); let s = full.join(" "); assert!(!s.contains("--network=none")); assert!(s.contains("--cap-drop=ALL") && s.contains(CONCAT_USER)); assert!(s.contains("--rm --sig-proxy=false --name sl-fetch-1")); } #[test] fn truncate_is_char_safe() { let s = "中文内容"; let t = truncate(s, 4); assert!(t.chars().count() <= 4); assert_eq!(truncate("short", 100), "short"); } }