summaryrefslogtreecommitdiff
path: root/server/src/worker.rs
diff options
context:
space:
mode:
authorSomhairle H. Marisol <[email protected]>2026-09-17 14:32:37 +0800
committerSomhairle H. Marisol <[email protected]>2026-09-17 14:32:37 +0800
commit5c0ba37eda80d39e6ceca59bb1d5f4942f858995 (patch)
tree948723f9cedf7ccb0707fa6ee516bd30fe20fd10 /server/src/worker.rs
downloadstrategy-lab-5c0ba37eda80d39e6ceca59bb1d5f4942f858995.tar.gz
chore: establish Strategy Lab source baseline (development, not release)
Diffstat (limited to 'server/src/worker.rs')
-rw-r--r--server/src/worker.rs482
1 files changed, 482 insertions, 0 deletions
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<i32>,
+ 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<String> {
+ let mut a: Vec<String> = [
+ "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<R>(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<ContainerResult> {
+ 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<ContainerResult> {
+ 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<AppState>,
+ network: bool,
+ mounts: &[(String, String, bool)],
+ args: &[String],
+ name: &str,
+ timeout_secs: u64,
+) -> AppResult<ContainerResult> {
+ 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<AppState>,
+ query: &str,
+ limit: i64,
+) -> AppResult<Vec<serde_json::Value>> {
+ 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::Value> = 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<std::collections::HashMap<String, (std::time::Instant, Vec<serde_json::Value>)>> {
+ static MAP: std::sync::OnceLock<tokio::sync::Mutex<std::collections::HashMap<String, (std::time::Instant, Vec<serde_json::Value>)>>> =
+ std::sync::OnceLock::new();
+ MAP.get_or_init(|| tokio::sync::Mutex::new(std::collections::HashMap::new()))
+}
+
+fn catalog_cache_get(key: &str) -> Option<Vec<serde_json::Value>> {
+ // 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/<pid>`), 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<i32> = 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/<pid> 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");
+ }
+}