From 5c0ba37eda80d39e6ceca59bb1d5f4942f858995 Mon Sep 17 00:00:00 2001 From: "Somhairle H. Marisol" Date: Thu, 17 Sep 2026 14:32:37 +0800 Subject: chore: establish Strategy Lab source baseline (development, not release) --- server/src/worker.rs | 482 +++++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 482 insertions(+) create mode 100644 server/src/worker.rs (limited to 'server/src/worker.rs') diff --git a/server/src/worker.rs b/server/src/worker.rs new file mode 100644 index 0000000..8e28de5 --- /dev/null +++ b/server/src/worker.rs @@ -0,0 +1,482 @@ +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"); + } +} -- cgit v1.2.3