Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.
| Every agent can have its own computer. A session that needs one wakes it: a home of its own on g1t cloud, one per agent and never shared, where it runs commands, reads and writes files and keeps what it made, with each session working in its own folder under a shared home; after ten idle minutes it sleeps, its home kept as a snapshot and restored when it wakes, and Reset wipes the home while memory and artifacts stay. Its shell and files are abilities with the usual choices, Alone, Alone when asked, Ask first or Never, offered only inside sessions and never to a chat reply; every command shows on the session with its output, and the agent's new Computer tab shows the state, the disk used of the five gigabytes included, the recent commands, and Wake, Put to sleep and Reset. Machine time counts only while it is awake, on the sandbox lines of the ledger that name the agent and who asked, held to the same spend caps as the session; the disk itself costs nothing in this version. The runner gained a long-lived supervisor that answers the computer's requests inside the container, and the runner service a computer per agent that keeps its snapshot in the agent homes bucket when one is attached, and says so when none is. The REST API and the agent tool can read a computer, wake it, put it to sleep and reset it. The agents, abilities, sessions, runners, billing and deploy guides say how it works and what an operator sets up; pinning a computer to your own runner, its browser and take-over come next. | 1 | //! An agent's own computer (`MODE=computer`): the long-lived process inside |
| 2 | //! the container that is a workspace agent's persistent home | |
| 3 | //! (docs.g1t.sh/guides/agents/, "Its computer"). | |
| 4 | //! | |
| 5 | //! Unlike every other mode, which does one job and exits, this one serves | |
| 6 | //! the runner's `AgentComputer` Durable Object over HTTP on `PORT` for as | |
| 7 | //! long as the computer is awake: | |
| 8 | //! | |
| 9 | //! - `GET /health`: uptime, disk used under the home, the cap. | |
| 10 | //! - `POST /exec`: runs a command under `bash -lc` with the home as `$HOME`, | |
| 11 | //! streaming stdout and stderr lines as NDJSON, then one closing line | |
| 12 | //! with the exit code and how long it took. Output past `MAX_OUTPUT_BYTES` | |
| 13 | //! is dropped and the closing line says so; a command past its timeout is | |
| 14 | //! killed and reported as exit 124. | |
| 15 | //! - `GET /files?path=` and `PUT /files?path=`: small files under the home. | |
| 16 | //! - `POST /restore`: a tar+zstd of the home, unpacked into it (how the | |
| 17 | //! Durable Object puts a snapshot back when it wakes the computer). | |
| 18 | //! - `POST /snapshot`: the home as tar+zstd, streamed; refused with the | |
| 19 | //! largest top-level directories when the home is over `DISK_CAP_BYTES`. | |
| 20 | //! - `POST /stop`: exits cleanly, once the snapshot was taken. | |
| 21 | //! | |
| 22 | //! Every path is checked to lie under the home (`within_home`); nothing | |
| 23 | //! outside it is read, written or run from. `G1T_HOME_RESTORE`, an HTTPS | |
| 24 | //! address with a one-time token in `G1T_HOME_RESTORE_TOKEN`, restores the | |
| 25 | //! home at start for installations that serve snapshots that way; g1t's | |
| 26 | //! cloud streams it in through `/restore` instead. Requests carry the token | |
| 27 | //! in `G1T_COMPUTER_TOKEN` as a bearer, so nothing else on the network | |
| 28 | //! drives the computer. | |
| 29 | //! | |
| 30 | //! Only the pure parts (path checks, framing, the closing line) are used | |
| 31 | //! off Linux, by tests. | |
| 32 | #![cfg_attr(not(unix), allow(dead_code))] | |
| 33 | ||
| 34 | use std::collections::BTreeMap; | |
| 35 | use std::io::{self, BufRead, BufReader, Read, Write}; | |
| 36 | use std::net::{TcpListener, TcpStream}; | |
| 37 | use std::path::{Path, PathBuf}; | |
| 38 | use std::process::{Child, Command, Stdio}; | |
| 39 | use std::sync::atomic::{AtomicUsize, Ordering}; | |
| 40 | use std::sync::mpsc; | |
| 41 | use std::sync::Arc; | |
| 42 | use std::time::{Duration, Instant}; | |
| 43 | ||
| 44 | use serde_json::{Value, json}; | |
| 45 | ||
| 46 | use crate::docker::http::{self, Body, Head}; | |
| 47 | ||
| 48 | /// Where the agent lives: its files, clones, caches and notes. | |
| 49 | pub(crate) const HOME: &str = "/home/agent"; | |
| 50 | /// Where a session's work goes by default, under the home. | |
| 51 | pub(crate) const SESSIONS_DIR: &str = "sessions"; | |
| 52 | /// The port the Durable Object reaches the computer on. | |
| 53 | pub(crate) const PORT: u16 = 8787; | |
| 54 | /// How much the home may hold: 5 GB, a hard cap in v1 (no disk charge). | |
| 55 | pub(crate) const DISK_CAP_BYTES: u64 = 5_000_000_000; | |
| 56 | /// The most output one command streams back; the rest is dropped and said so. | |
| 57 | pub(crate) const MAX_OUTPUT_BYTES: usize = 1024 * 1024; | |
| 58 | /// The most a file read or written through `/files` may be. | |
| 59 | pub(crate) const MAX_FILE_BYTES: u64 = 5 * 1024 * 1024; | |
| 60 | /// A command's timeout when none is given, and the longest one may ask for. | |
| 61 | pub(crate) const DEFAULT_TIMEOUT_SECONDS: u64 = 120; | |
| 62 | pub(crate) const MAX_TIMEOUT_SECONDS: u64 = 1800; | |
| 63 | /// The exit code of a command that was killed at its timeout, as `timeout(1)` reports it. | |
| 64 | pub(crate) const TIMED_OUT_EXIT: i32 = 124; | |
| 65 | ||
| 66 | /// What the process knows while it serves. | |
| 67 | struct Server { | |
| 68 | token: Option<String>, | |
| 69 | started: Instant, | |
| 70 | /// Commands running now: the Durable Object does not put a busy computer to sleep. | |
| 71 | running: AtomicUsize, | |
| 72 | } | |
| 73 | ||
| 74 | /// `MODE=computer`: prepares the home, restores it if asked, then serves until `/stop`. | |
| 75 | pub fn main() -> i32 { | |
| 76 | let server = Arc::new(Server { | |
| 77 | token: std::env::var("G1T_COMPUTER_TOKEN").ok().filter(|t| !t.is_empty()), | |
| 78 | started: Instant::now(), | |
| 79 | running: AtomicUsize::new(0), | |
| 80 | }); | |
| 81 | if let Err(error) = prepare_home(Path::new(HOME)) { | |
| 82 | eprintln!("g1t-runner: the home could not be prepared: {error}"); | |
| 83 | return 2; | |
| 84 | } | |
| 85 | if let Ok(url) = std::env::var("G1T_HOME_RESTORE") { | |
| 86 | if !url.trim().is_empty() { | |
| 87 | let token = std::env::var("G1T_HOME_RESTORE_TOKEN").unwrap_or_default(); | |
| 88 | match restore_from_url(&url, &token) { | |
| 89 | Ok(()) => eprintln!("g1t-runner: the home was restored"), | |
| 90 | Err(error) => eprintln!("g1t-runner: the home could not be restored, so it starts fresh: {error}"), | |
| 91 | } | |
| 92 | } | |
| 93 | } | |
| 94 | let listener = match TcpListener::bind(("0.0.0.0", PORT)) { | |
| 95 | Ok(listener) => listener, | |
| 96 | Err(error) => { | |
| 97 | eprintln!("g1t-runner: could not listen on {PORT}: {error}"); | |
| 98 | return 2; | |
| 99 | } | |
| 100 | }; | |
| 101 | eprintln!("g1t-runner: the computer is awake on {PORT}"); | |
| 102 | for client in listener.incoming() { | |
| 103 | let Ok(client) = client else { continue }; | |
| 104 | let server = server.clone(); | |
| 105 | std::thread::spawn(move || { | |
| 106 | if let Err(error) = serve(client, &server) { | |
| 107 | eprintln!("g1t-runner: a request failed: {error}"); | |
| 108 | } | |
| 109 | }); | |
| 110 | } | |
| 111 | 0 | |
| 112 | } | |
| 113 | ||
| 114 | /// Makes sure the home exists and is ours to write, with its sessions | |
| 115 | /// directory. The image runs as `node` with sudo, and `/home` is root's. | |
| 116 | fn prepare_home(home: &Path) -> io::Result<()> { | |
| 117 | if std::fs::create_dir_all(home).is_err() || !writable(home) { | |
| 118 | let user = Command::new("id").arg("-un").output().ok().map(|o| String::from_utf8_lossy(&o.stdout).trim().to_owned()).unwrap_or_default(); | |
| 119 | let made = Command::new("sudo").args(["-n", "mkdir", "-p"]).arg(home).status().is_ok_and(|s| s.success()); | |
| 120 | let owned = Command::new("sudo").args(["-n", "chown", "-R", &format!("{user}:{user}")]).arg(home).status().is_ok_and(|s| s.success()); | |
| 121 | if !made || !owned || !writable(home) { | |
| 122 | return Err(io::Error::new(io::ErrorKind::PermissionDenied, format!("{} is not writable", home.display()))); | |
| 123 | } | |
| 124 | } | |
| 125 | std::fs::create_dir_all(home.join(SESSIONS_DIR)) | |
| 126 | } | |
| 127 | ||
| 128 | fn writable(dir: &Path) -> bool { | |
| 129 | let probe = dir.join(".g1t-probe"); | |
| 130 | let ok = std::fs::write(&probe, b"").is_ok(); | |
| 131 | let _ = std::fs::remove_file(&probe); | |
| 132 | ok | |
| 133 | } | |
| 134 | ||
| 135 | /// Downloads a tar+zstd snapshot and unpacks it into the home. | |
| 136 | fn restore_from_url(url: &str, token: &str) -> Result<(), String> { | |
| 137 | let response = ureq::get(url) | |
| 138 | .set("Authorization", &format!("Bearer {token}")) | |
| 139 | .timeout(Duration::from_secs(600)) | |
| 140 | .call() | |
| 141 | .map_err(|error| error.to_string())?; | |
| 142 | let mut reader = response.into_reader(); | |
| 143 | let mut tar = untar(Path::new(HOME)).map_err(|error| error.to_string())?; | |
| 144 | { | |
| 145 | let mut stdin = tar.stdin.take().ok_or("tar has no stdin")?; | |
| 146 | io::copy(&mut reader, &mut stdin).map_err(|error| error.to_string())?; | |
| 147 | } | |
| 148 | let status = tar.wait().map_err(|error| error.to_string())?; | |
| 149 | if !status.success() { | |
| 150 | return Err(format!("tar exited with {status}")); | |
| 151 | } | |
| 152 | Ok(()) | |
| 153 | } | |
| 154 | ||
| 155 | /// `tar --zstd -x` into `home`, reading from its stdin. | |
| 156 | fn untar(home: &Path) -> io::Result<Child> { | |
| 157 | Command::new("tar").args(["--zstd", "-x", "-f", "-", "-C"]).arg(home).stdin(Stdio::piped()).stdout(Stdio::null()).stderr(Stdio::inherit()).spawn() | |
| 158 | } | |
| 159 | ||
| 160 | /// Serves one connection: one request, one response. | |
| 161 | fn serve(client: TcpStream, server: &Server) -> io::Result<()> { | |
| 162 | let mut reader = BufReader::new(client.try_clone()?); | |
| 163 | let mut writer = client; | |
| 164 | let Some(head) = http::read_head(&mut reader)? else { return Ok(()) }; | |
| 165 | let body = http::request_body(&head); | |
| 166 | let (path, query) = split_target(head.target()); | |
| 167 | let authorised = match &server.token { | |
| 168 | None => true, | |
| 169 | Some(token) => head.header("authorization").is_some_and(|value| value.trim() == format!("Bearer {token}")), | |
| 170 | }; | |
| 171 | if !authorised { | |
| 172 | let _ = http::read_body(&mut reader, body); | |
| 173 | return respond(&mut writer, 401, "Unauthorized", &json!({ "message": "This computer takes requests from its own Durable Object only." })); | |
| 174 | } | |
| 175 | match (head.method(), path) { | |
| 176 | ("GET", "/health") => { | |
| 177 | let _ = http::read_body(&mut reader, body); | |
| 178 | respond(&mut writer, 200, "OK", &health(server)) | |
| 179 | } | |
| 180 | ("POST", "/exec") => { | |
| 181 | let raw = http::read_body(&mut reader, body)?; | |
| 182 | let request: Value = serde_json::from_slice(&raw).unwrap_or(Value::Null); | |
| 183 | let Some(cmd) = request.get("cmd").and_then(Value::as_str).filter(|c| !c.trim().is_empty()) else { | |
| 184 | return respond(&mut writer, 400, "Bad Request", &json!({ "message": "Give the command to run as cmd." })); | |
| 185 | }; | |
| 186 | let cwd = match within_home(request.get("cwd").and_then(Value::as_str).unwrap_or("")) { | |
| 187 | Some(cwd) => cwd, | |
| 188 | None => return respond(&mut writer, 400, "Bad Request", &json!({ "message": format!("cwd must be under {HOME}.") })), | |
| 189 | }; | |
| 190 | let timeout = request.get("timeout_seconds").and_then(Value::as_u64).unwrap_or(DEFAULT_TIMEOUT_SECONDS).clamp(1, MAX_TIMEOUT_SECONDS); | |
| 191 | let env = request.get("env").and_then(Value::as_object).map(|map| map.iter().filter_map(|(k, v)| v.as_str().map(|v| (k.clone(), v.to_owned()))).collect()).unwrap_or_default(); | |
| 192 | server.running.fetch_add(1, Ordering::SeqCst); | |
| 193 | let result = stream_exec(&mut writer, cmd, &cwd, Duration::from_secs(timeout), &env); | |
| 194 | server.running.fetch_sub(1, Ordering::SeqCst); | |
| 195 | result | |
| 196 | } | |
| 197 | ("GET", "/files") => { | |
| 198 | let _ = http::read_body(&mut reader, body); | |
| 199 | let Some(path) = query.get("path").and_then(|p| resolve_in_home(p)) else { | |
| 200 | return respond(&mut writer, 400, "Bad Request", &json!({ "message": format!("path must be under {HOME}.") })); | |
| 201 | }; | |
| 202 | let size = match std::fs::metadata(&path) { | |
| 203 | Ok(meta) if meta.is_file() => meta.len(), | |
| 204 | Ok(_) => return respond(&mut writer, 400, "Bad Request", &json!({ "message": "That is a directory, not a file." })), | |
| 205 | Err(_) => return respond(&mut writer, 404, "Not Found", &json!({ "message": "There is no such file." })), | |
| 206 | }; | |
| 207 | if size > MAX_FILE_BYTES { | |
| 208 | return respond(&mut writer, 413, "Payload Too Large", &json!({ "message": format!("The file is {size} bytes; files read this way are at most {MAX_FILE_BYTES}."), "bytes": size })); | |
| 209 | } | |
| 210 | let bytes = std::fs::read(&path)?; | |
| 211 | let head = Head { start: "HTTP/1.1 200 OK".into(), headers: vec![("Content-Type".into(), "application/octet-stream".into())] }; | |
| 212 | writer.write_all(&http::with_length(head, &bytes))?; | |
| 213 | writer.flush() | |
| 214 | } | |
| 215 | ("PUT", "/files") => { | |
| 216 | let Some(path) = query.get("path").and_then(|p| resolve_in_home(p)) else { | |
| 217 | let _ = http::read_body(&mut reader, body); | |
| 218 | return respond(&mut writer, 400, "Bad Request", &json!({ "message": format!("path must be under {HOME}.") })); | |
| 219 | }; | |
| 220 | if let Body::Length(n) = body { | |
| 221 | if n > MAX_FILE_BYTES { | |
| 222 | return respond(&mut writer, 413, "Payload Too Large", &json!({ "message": format!("Files written this way are at most {MAX_FILE_BYTES} bytes.") })); | |
| 223 | } | |
| 224 | } | |
| 225 | let bytes = http::read_body(&mut reader, body)?; | |
| 226 | if bytes.len() as u64 > MAX_FILE_BYTES { | |
| 227 | return respond(&mut writer, 413, "Payload Too Large", &json!({ "message": format!("Files written this way are at most {MAX_FILE_BYTES} bytes.") })); | |
| 228 | } | |
| 229 | if let Some(parent) = path.parent() { | |
| 230 | std::fs::create_dir_all(parent)?; | |
| 231 | } | |
| 232 | std::fs::write(&path, &bytes)?; | |
| 233 | respond(&mut writer, 200, "OK", &json!({ "ok": true, "path": path.to_string_lossy(), "bytes": bytes.len() })) | |
| 234 | } | |
| 235 | ("POST", "/restore") => { | |
| 236 | let mut tar = untar(Path::new(HOME))?; | |
| 237 | let copied = { | |
| 238 | let mut stdin = tar.stdin.take().expect("tar's stdin"); | |
| 239 | http::copy_body(&mut reader, &mut stdin, body) | |
| 240 | }; | |
| 241 | let status = tar.wait()?; | |
| 242 | match (copied, status.success()) { | |
| 243 | (Ok(()), true) => respond(&mut writer, 200, "OK", &json!({ "ok": true, "disk_used_bytes": dir_size(Path::new(HOME)) })), | |
| 244 | (Err(error), _) => respond(&mut writer, 400, "Bad Request", &json!({ "message": format!("The snapshot could not be read: {error}") })), | |
| 245 | (Ok(()), false) => respond(&mut writer, 500, "Internal Server Error", &json!({ "message": format!("tar exited with {status}") })), | |
| 246 | } | |
| 247 | } | |
| 248 | ("POST", "/snapshot") => { | |
| 249 | let _ = http::read_body(&mut reader, body); | |
| 250 | let home = Path::new(HOME); | |
| 251 | let used = dir_size(home); | |
| 252 | if used > DISK_CAP_BYTES { | |
| 253 | let largest: Vec<Value> = largest_top_level(home).into_iter().take(8).map(|(name, bytes)| json!({ "path": name, "bytes": bytes })).collect(); | |
| 254 | return respond( | |
| 255 | &mut writer, | |
| 256 | 413, | |
| 257 | "Payload Too Large", | |
| 258 | &json!({ "message": over_cap_message(used, &largest_top_level(home)), "disk_used_bytes": used, "disk_cap_bytes": DISK_CAP_BYTES, "largest": largest }), | |
| 259 | ); | |
| 260 | } | |
| 261 | stream_snapshot(&mut writer, home, used) | |
| 262 | } | |
| 263 | ("POST", "/stop") => { | |
| 264 | let _ = http::read_body(&mut reader, body); | |
| 265 | respond(&mut writer, 200, "OK", &json!({ "ok": true }))?; | |
| 266 | let _ = writer.shutdown(std::net::Shutdown::Both); | |
| 267 | eprintln!("g1t-runner: the computer is going to sleep"); | |
| 268 | std::process::exit(0); | |
| 269 | } | |
| 270 | _ => { | |
| 271 | let _ = http::read_body(&mut reader, body); | |
| 272 | respond(&mut writer, 404, "Not Found", &json!({ "message": "No such route." })) | |
| 273 | } | |
| 274 | } | |
| 275 | } | |
| 276 | ||
| 277 | fn respond(writer: &mut TcpStream, status: u16, reason: &str, body: &Value) -> io::Result<()> { | |
| 278 | writer.write_all(&http::json_response(status, reason, body))?; | |
| 279 | writer.flush() | |
| 280 | } | |
| 281 | ||
| 282 | fn health(server: &Server) -> Value { | |
| 283 | json!({ | |
| 284 | "uptime_seconds": server.started.elapsed().as_secs(), | |
| 285 | "disk_used_bytes": dir_size(Path::new(HOME)), | |
| 286 | "disk_cap_bytes": DISK_CAP_BYTES, | |
| 287 | "running_commands": server.running.load(Ordering::SeqCst), | |
| 288 | "home": HOME, | |
| 289 | }) | |
| 290 | } | |
| 291 | ||
| 292 | /// A request target as its path and its query, decoded. | |
| 293 | pub(crate) fn split_target(target: &str) -> (&str, BTreeMap<String, String>) { | |
| 294 | let (path, query) = target.split_once('?').unwrap_or((target, "")); | |
| 295 | let mut out = BTreeMap::new(); | |
| 296 | for pair in query.split('&').filter(|p| !p.is_empty()) { | |
| 297 | let (key, value) = pair.split_once('=').unwrap_or((pair, "")); | |
| 298 | out.insert(percent_decode(key), percent_decode(value)); | |
| 299 | } | |
| 300 | (path, out) | |
| 301 | } | |
| 302 | ||
| 303 | fn percent_decode(text: &str) -> String { | |
| 304 | let bytes = text.as_bytes(); | |
| 305 | let mut out = Vec::with_capacity(bytes.len()); | |
| 306 | let mut i = 0; | |
| 307 | while i < bytes.len() { | |
| 308 | match bytes[i] { | |
| 309 | b'%' if i + 2 < bytes.len() => { | |
| 310 | let hex = &text[i + 1..i + 3]; | |
| 311 | match u8::from_str_radix(hex, 16) { | |
| 312 | Ok(byte) => { | |
| 313 | out.push(byte); | |
| 314 | i += 3; | |
| 315 | } | |
| 316 | Err(_) => { | |
| 317 | out.push(b'%'); | |
| 318 | i += 1; | |
| 319 | } | |
| 320 | } | |
| 321 | } | |
| 322 | b'+' => { | |
| 323 | out.push(b' '); | |
| 324 | i += 1; | |
| 325 | } | |
| 326 | byte => { | |
| 327 | out.push(byte); | |
| 328 | i += 1; | |
| 329 | } | |
| 330 | } | |
| 331 | } | |
| 332 | String::from_utf8_lossy(&out).into_owned() | |
| 333 | } | |
| 334 | ||
| 335 | // ── Paths ───────────────────────────────────────────────────────────────── | |
| 336 | ||
| 337 | /// `given` as a path under `home`, or `None` when it would leave it: a | |
| 338 | /// relative path is taken from the home, `..` is resolved without touching | |
| 339 | /// the filesystem, and anything that climbs above the home is refused. The | |
| 340 | /// empty path is the home itself. Pure, so it is tested everywhere. | |
| 341 | pub(crate) fn within(home: &str, given: &str) -> Option<PathBuf> { | |
| 342 | if given.contains('\0') { | |
| 343 | return None; | |
| 344 | } | |
| 345 | let given = given.trim(); | |
| 346 | let home = home.trim_end_matches('/'); | |
| 347 | let mut parts: Vec<&str> = if given.starts_with('/') { Vec::new() } else { home.split('/').filter(|p| !p.is_empty()).collect() }; | |
| 348 | for part in given.split('/') { | |
| 349 | match part { | |
| 350 | "" | "." => {} | |
| 351 | ".." => { | |
| 352 | parts.pop()?; | |
| 353 | } | |
| 354 | name => parts.push(name), | |
| 355 | } | |
| 356 | } | |
| 357 | let joined = format!("/{}", parts.join("/")); | |
| 358 | if joined == home || joined.starts_with(&format!("{home}/")) { Some(PathBuf::from(joined)) } else { None } | |
| 359 | } | |
| 360 | ||
| 361 | /// `within` for the real home. | |
| 362 | pub(crate) fn within_home(given: &str) -> Option<PathBuf> { | |
| 363 | within(HOME, given) | |
| 364 | } | |
| 365 | ||
| 366 | /// `within_home`, and then the path as the filesystem resolves it: a | |
| 367 | /// symlink inside the home that points out of it is refused too. The | |
| 368 | /// deepest existing ancestor is what is resolved, so a file not made yet | |
| 369 | /// can still be written. | |
| 370 | fn resolve_in_home(given: &str) -> Option<PathBuf> { | |
| 371 | let path = within_home(given)?; | |
| 372 | let home = std::fs::canonicalize(HOME).unwrap_or_else(|_| PathBuf::from(HOME)); | |
| 373 | let mut probe = path.clone(); | |
| 374 | let mut tail: Vec<std::ffi::OsString> = Vec::new(); | |
| 375 | loop { | |
| 376 | if let Ok(real) = std::fs::canonicalize(&probe) { | |
| 377 | if real != home && !real.starts_with(&home) { | |
| 378 | return None; | |
| 379 | } | |
| 380 | let mut out = real; | |
| 381 | for part in tail.iter().rev() { | |
| 382 | out.push(part); | |
| 383 | } | |
| 384 | return Some(out); | |
| 385 | } | |
| 386 | let name = probe.file_name()?.to_owned(); | |
| 387 | tail.push(name); | |
| 388 | probe = probe.parent()?.to_path_buf(); | |
| 389 | } | |
| 390 | } | |
| 391 | ||
| 392 | // ── Running commands ────────────────────────────────────────────────────── | |
| 393 | ||
| 394 | /// One line of a command's output, as streamed. | |
| 395 | pub(crate) fn output_line(stream: &str, line: &str) -> Value { | |
| 396 | json!({ "stream": stream, "line": line }) | |
| 397 | } | |
| 398 | ||
| 399 | /// The closing line of a command: how it ended. | |
| 400 | pub(crate) fn final_line(exit_code: i32, duration: Duration, dropped_bytes: usize, timed_out: bool) -> Value { | |
| 401 | let mut line = json!({ "exit_code": exit_code, "duration_ms": duration.as_millis() as u64, "truncated": dropped_bytes > 0, "timed_out": timed_out }); | |
| 402 | if dropped_bytes > 0 { | |
| 403 | line["note"] = json!(format!("Output past {} was dropped: {} more bytes.", human_bytes(MAX_OUTPUT_BYTES as u64), dropped_bytes)); | |
| 404 | } | |
| 405 | if timed_out { | |
| 406 | line["note"] = json!(format!( | |
| 407 | "{}The command was still running at its timeout and was killed.", | |
| 408 | line.get("note").and_then(Value::as_str).map(|n| format!("{n} ")).unwrap_or_default() | |
| 409 | )); | |
| 410 | } | |
| 411 | line | |
| 412 | } | |
| 413 | ||
| 414 | /// Environment the agent may set for a command: plain names, and never the | |
| 415 | /// process's own credentials or anything that changes how programs load. | |
| 416 | pub(crate) fn safe_env(given: &BTreeMap<String, String>) -> BTreeMap<String, String> { | |
| 417 | given | |
| 418 | .iter() | |
| 419 | .filter(|(name, _)| { | |
| 420 | let plain = !name.is_empty() && name.chars().all(|c| c.is_ascii_alphanumeric() || c == '_') && !name.starts_with(|c: char| c.is_ascii_digit()); | |
| 421 | let risky = name.starts_with("G1T_") || name.starts_with("LD_") || matches!(name.as_str(), "BASH_ENV" | "ENV" | "PROMPT_COMMAND" | "HOME" | "SHELLOPTS" | "BASHOPTS"); | |
| 422 | plain && !risky | |
| 423 | }) | |
| 424 | .map(|(name, value)| (name.clone(), value.clone())) | |
| 425 | .collect() | |
| 426 | } | |
| 427 | ||
| 428 | /// Something a command's output produced. | |
| 429 | enum Produced { | |
| 430 | Line(&'static str, String), | |
| 431 | Closed, | |
| 432 | } | |
| 433 | ||
| 434 | /// Runs `cmd` and gives every line to `sink` as it comes, then the closing | |
| 435 | /// line. Returns what the closing line said. | |
| 436 | pub(crate) fn run_streaming(cmd: &str, cwd: &Path, timeout: Duration, env: &BTreeMap<String, String>, sink: &mut dyn FnMut(Value) -> io::Result<()>) -> io::Result<Value> { | |
| 437 | std::fs::create_dir_all(cwd)?; | |
| 438 | let started = Instant::now(); | |
| 439 | let mut command = Command::new("bash"); | |
| 440 | command.args(["-lc", cmd]).current_dir(cwd).stdin(Stdio::null()).stdout(Stdio::piped()).stderr(Stdio::piped()); | |
| 441 | command.env_remove("G1T_COMPUTER_TOKEN").env_remove("G1T_HOME_RESTORE").env_remove("G1T_HOME_RESTORE_TOKEN"); | |
| 442 | command.env("HOME", HOME).env("G1T_COMPUTER", "1"); | |
| 443 | for (name, value) in safe_env(env) { | |
| 444 | command.env(name, value); | |
| 445 | } | |
| 446 | let mut child = match command.spawn() { | |
| 447 | Ok(child) => child, | |
| 448 | Err(error) => { | |
| 449 | let line = json!({ "exit_code": 127, "duration_ms": 0, "truncated": false, "timed_out": false, "note": format!("bash could not be started: {error}") }); | |
| 450 | sink(line.clone())?; | |
| 451 | return Ok(line); | |
| 452 | } | |
| 453 | }; | |
| 454 | let (sender, receiver) = mpsc::channel::<Produced>(); | |
| 455 | let mut readers = Vec::new(); | |
| 456 | for (name, pipe) in [("stdout", child.stdout.take().map(|p| Box::new(p) as Box<dyn Read + Send>)), ("stderr", child.stderr.take().map(|p| Box::new(p) as Box<dyn Read + Send>))] { | |
| 457 | let Some(pipe) = pipe else { continue }; | |
| 458 | let sender = sender.clone(); | |
| 459 | readers.push(std::thread::spawn(move || { | |
| 460 | let mut reader = BufReader::new(pipe); | |
| 461 | let mut buffer = Vec::new(); | |
| 462 | loop { | |
| 463 | buffer.clear(); | |
| 464 | match reader.read_until(b'\n', &mut buffer) { | |
| 465 | Ok(0) | Err(_) => break, | |
| 466 | Ok(_) => { | |
| 467 | let text = String::from_utf8_lossy(&buffer).trim_end_matches(['\r', '\n']).to_owned(); | |
| 468 | if sender.send(Produced::Line(name, text)).is_err() { | |
| 469 | break; | |
| 470 | } | |
| 471 | } | |
| 472 | } | |
| 473 | } | |
| 474 | let _ = sender.send(Produced::Closed); | |
| 475 | })); | |
| 476 | } | |
| 477 | drop(sender); | |
| 478 | let mut open = readers.len(); | |
| 479 | let mut sent = 0usize; | |
| 480 | let mut dropped = 0usize; | |
| 481 | let mut timed_out = false; | |
| 482 | let deadline = started + timeout; | |
| 483 | while open > 0 { | |
| 484 | let now = Instant::now(); | |
| 485 | if now >= deadline { | |
| 486 | timed_out = true; | |
| 487 | let _ = child.kill(); | |
| 488 | break; | |
| 489 | } | |
| 490 | match receiver.recv_timeout(deadline - now) { | |
| 491 | Ok(Produced::Line(stream, text)) => { | |
| 492 | let bytes = text.len() + 1; | |
| 493 | if sent + bytes > MAX_OUTPUT_BYTES { | |
| 494 | dropped += bytes; | |
| 495 | continue; | |
| 496 | } | |
| 497 | sent += bytes; | |
| 498 | sink(output_line(stream, &text))?; | |
| 499 | } | |
| 500 | Ok(Produced::Closed) => open -= 1, | |
| 501 | Err(mpsc::RecvTimeoutError::Timeout) => { | |
| 502 | timed_out = true; | |
| 503 | let _ = child.kill(); | |
| 504 | break; | |
| 505 | } | |
| 506 | Err(mpsc::RecvTimeoutError::Disconnected) => break, | |
| 507 | } | |
| 508 | } | |
| 509 | let status = child.wait()?; | |
| 510 | for reader in readers { | |
| 511 | let _ = reader.join(); | |
| 512 | } | |
| 513 | let exit_code = if timed_out { TIMED_OUT_EXIT } else { status.code().unwrap_or(-1) }; | |
| 514 | let line = final_line(exit_code, started.elapsed(), dropped, timed_out); | |
| 515 | sink(line.clone())?; | |
| 516 | Ok(line) | |
| 517 | } | |
| 518 | ||
| 519 | /// `/exec`: the command's lines as a chunked NDJSON response. | |
| 520 | fn stream_exec(writer: &mut TcpStream, cmd: &str, cwd: &Path, timeout: Duration, env: &BTreeMap<String, String>) -> io::Result<()> { | |
| 521 | writer.write_all(b"HTTP/1.1 200 OK\r\nContent-Type: application/x-ndjson\r\nTransfer-Encoding: chunked\r\n\r\n")?; | |
| 522 | writer.flush()?; | |
| 523 | let mut sink = |line: Value| -> io::Result<()> { | |
| 524 | let mut text = line.to_string(); | |
| 525 | text.push('\n'); | |
| 526 | write_chunk(writer, text.as_bytes()) | |
| 527 | }; | |
| 528 | run_streaming(cmd, cwd, timeout, env, &mut sink)?; | |
| 529 | writer.write_all(b"0\r\n\r\n")?; | |
| 530 | writer.flush() | |
| 531 | } | |
| 532 | ||
| 533 | /// One chunk of a chunked body, flushed so it is seen as it is written. | |
| 534 | pub(crate) fn write_chunk<W: Write>(writer: &mut W, bytes: &[u8]) -> io::Result<()> { | |
| 535 | if bytes.is_empty() { | |
| 536 | return Ok(()); | |
| 537 | } | |
| 538 | write!(writer, "{:x}\r\n", bytes.len())?; | |
| 539 | writer.write_all(bytes)?; | |
| 540 | writer.write_all(b"\r\n")?; | |
| 541 | writer.flush() | |
| 542 | } | |
| 543 | ||
| 544 | // ── Snapshots ───────────────────────────────────────────────────────────── | |
| 545 | ||
| 546 | /// Every file's size under `dir`, symlinks not followed. | |
| 547 | pub(crate) fn dir_size(dir: &Path) -> u64 { | |
| 548 | let Ok(entries) = std::fs::read_dir(dir) else { return 0 }; | |
| 549 | let mut total = 0; | |
| 550 | for entry in entries.flatten() { | |
| 551 | let Ok(meta) = entry.metadata() else { continue }; | |
| 552 | if meta.is_dir() { | |
| 553 | total += dir_size(&entry.path()); | |
| 554 | } else if meta.is_file() { | |
| 555 | total += meta.len(); | |
| 556 | } | |
| 557 | } | |
| 558 | total | |
| 559 | } | |
| 560 | ||
| 561 | /// The home's top-level entries by size, largest first. | |
| 562 | pub(crate) fn largest_top_level(home: &Path) -> Vec<(String, u64)> { | |
| 563 | let Ok(entries) = std::fs::read_dir(home) else { return Vec::new() }; | |
| 564 | let mut out: Vec<(String, u64)> = entries | |
| 565 | .flatten() | |
| 566 | .map(|entry| { | |
| 567 | let size = entry.metadata().map(|meta| if meta.is_dir() { dir_size(&entry.path()) } else { meta.len() }).unwrap_or(0); | |
| 568 | (entry.file_name().to_string_lossy().into_owned(), size) | |
| 569 | }) | |
| 570 | .collect(); | |
| 571 | out.sort_by(|a, b| b.1.cmp(&a.1)); | |
| 572 | out | |
| 573 | } | |
| 574 | ||
| 575 | /// What the computer says when its home is over the cap. | |
| 576 | pub(crate) fn over_cap_message(used: u64, largest: &[(String, u64)]) -> String { | |
| 577 | let named: Vec<String> = largest.iter().take(5).map(|(name, bytes)| format!("{name} ({})", human_bytes(*bytes))).collect(); | |
| 578 | format!( | |
| 579 | "The home holds {}, over its {} cap, so it can't be saved. Largest at the top level: {}. Remove or trim some of it, or reset the computer.", | |
| 580 | human_bytes(used), | |
| 581 | human_bytes(DISK_CAP_BYTES), | |
| 582 | if named.is_empty() { "nothing".to_owned() } else { named.join(", ") } | |
| 583 | ) | |
| 584 | } | |
| 585 | ||
| 586 | pub(crate) fn human_bytes(bytes: u64) -> String { | |
| 587 | const UNITS: [&str; 5] = ["B", "KB", "MB", "GB", "TB"]; | |
| 588 | let mut value = bytes as f64; | |
| 589 | let mut unit = 0; | |
| 590 | while value >= 1000.0 && unit < UNITS.len() - 1 { | |
| 591 | value /= 1000.0; | |
| 592 | unit += 1; | |
| 593 | } | |
| 594 | if unit == 0 { format!("{bytes} B") } else { format!("{value:.1} {}", UNITS[unit]) } | |
| 595 | } | |
| 596 | ||
| 597 | /// `/snapshot`: the home as tar+zstd, chunked as tar writes it. A tar | |
| 598 | /// that fails ends the connection without the closing chunk, so the reader | |
| 599 | /// knows the archive is not whole. | |
| 600 | fn stream_snapshot(writer: &mut TcpStream, home: &Path, used: u64) -> io::Result<()> { | |
| 601 | let mut tar = Command::new("tar") | |
| 602 | .args(["--zstd", "-c", "-f", "-", "-C"]) | |
| 603 | .arg(home) | |
| 604 | .arg(".") | |
| 605 | .stdin(Stdio::null()) | |
| 606 | .stdout(Stdio::piped()) | |
| 607 | .stderr(Stdio::inherit()) | |
| 608 | .spawn()?; | |
| 609 | let mut stdout = tar.stdout.take().expect("tar's stdout"); | |
| 610 | writer.write_all(format!("HTTP/1.1 200 OK\r\nContent-Type: application/zstd\r\nTransfer-Encoding: chunked\r\nX-G1t-Disk-Used: {used}\r\n\r\n").as_bytes())?; | |
| 611 | let mut buffer = vec![0u8; 256 * 1024]; | |
| 612 | loop { | |
| 613 | let n = stdout.read(&mut buffer)?; | |
| 614 | if n == 0 { | |
| 615 | break; | |
| 616 | } | |
| 617 | write_chunk(writer, &buffer[..n])?; | |
| 618 | } | |
| 619 | let status = tar.wait()?; | |
| 620 | if !status.success() { | |
| 621 | let _ = writer.shutdown(std::net::Shutdown::Both); | |
| 622 | return Err(io::Error::other(format!("tar exited with {status}"))); | |
| 623 | } | |
| 624 | writer.write_all(b"0\r\n\r\n")?; | |
| 625 | writer.flush() | |
| 626 | } | |
| 627 | ||
| 628 | #[cfg(test)] | |
| 629 | mod tests { | |
| 630 | use super::*; | |
| 631 | ||
| 632 | #[test] | |
| 633 | fn paths_stay_under_the_home() { | |
| 634 | let home = "/home/agent"; | |
| 635 | assert_eq!(within(home, ""), Some(PathBuf::from("/home/agent"))); | |
| 636 | assert_eq!(within(home, "notes.md"), Some(PathBuf::from("/home/agent/notes.md"))); | |
| 637 | assert_eq!(within(home, "./sessions/asn_1/"), Some(PathBuf::from("/home/agent/sessions/asn_1"))); | |
| 638 | assert_eq!(within(home, "/home/agent/x/../y"), Some(PathBuf::from("/home/agent/y"))); | |
| 639 | assert_eq!(within(home, "a/b/../../c"), Some(PathBuf::from("/home/agent/c"))); | |
| 640 | assert_eq!(within(home, "/home/agent"), Some(PathBuf::from("/home/agent"))); | |
| 641 | } | |
| 642 | ||
| 643 | #[test] | |
| 644 | fn paths_that_leave_the_home_are_refused() { | |
| 645 | let home = "/home/agent"; | |
| 646 | assert_eq!(within(home, "../etc/passwd"), None); | |
| 647 | assert_eq!(within(home, "/etc/passwd"), None); | |
| 648 | assert_eq!(within(home, "/home/agentx/secret"), None); | |
| 649 | assert_eq!(within(home, "/home/agent/../node/.ssh"), None); | |
| 650 | assert_eq!(within(home, "a/../../.."), None); | |
| 651 | assert_eq!(within(home, "/"), None); | |
| 652 | assert_eq!(within(home, "ok\0bad"), None); | |
| 653 | } | |
| 654 | ||
| 655 | #[test] | |
| 656 | fn targets_split_into_path_and_query() { | |
| 657 | let (path, query) = split_target("/files?path=%2Fhome%2Fagent%2Fa+b.txt&x=1"); | |
| 658 | assert_eq!(path, "/files"); | |
| 659 | assert_eq!(query.get("path").map(String::as_str), Some("/home/agent/a b.txt")); | |
| 660 | assert_eq!(query.get("x").map(String::as_str), Some("1")); | |
| 661 | assert_eq!(split_target("/health").0, "/health"); | |
| 662 | } | |
| 663 | ||
| 664 | #[test] | |
| 665 | fn the_closing_line_says_how_it_ended() { | |
| 666 | let plain = final_line(0, Duration::from_millis(1234), 0, false); | |
| 667 | assert_eq!(plain["exit_code"], 0); | |
| 668 | assert_eq!(plain["duration_ms"], 1234); | |
| 669 | assert_eq!(plain["truncated"], false); | |
| 670 | assert_eq!(plain["timed_out"], false); | |
| 671 | assert!(plain.get("note").is_none()); | |
| 672 | let cut = final_line(TIMED_OUT_EXIT, Duration::from_secs(2), 10, true); | |
| 673 | assert_eq!(cut["truncated"], true); | |
| 674 | assert_eq!(cut["timed_out"], true); | |
| 675 | let note = cut["note"].as_str().unwrap(); | |
| 676 | assert!(note.contains("dropped") && note.contains("killed"), "{note}"); | |
| 677 | } | |
| 678 | ||
| 679 | #[test] | |
| 680 | fn the_environment_keeps_out_what_it_must() { | |
| 681 | let mut given = BTreeMap::new(); | |
| 682 | for (k, v) in [("FOO", "1"), ("G1T_TOKEN", "x"), ("LD_PRELOAD", "evil.so"), ("HOME", "/root"), ("bad-name", "1"), ("1ABC", "1"), ("CI", "true")] { | |
| 683 | given.insert(k.to_owned(), v.to_owned()); | |
| 684 | } | |
| 685 | let kept = safe_env(&given); | |
| 686 | assert_eq!(kept.keys().cloned().collect::<Vec<_>>(), vec!["CI".to_owned(), "FOO".to_owned()]); | |
| 687 | } | |
| 688 | ||
| 689 | #[test] | |
| 690 | fn chunks_are_framed() { | |
| 691 | let mut out = Vec::new(); | |
| 692 | write_chunk(&mut out, b"hello").unwrap(); | |
| 693 | write_chunk(&mut out, b"").unwrap(); | |
| 694 | assert_eq!(out, b"5\r\nhello\r\n"); | |
| 695 | } | |
| 696 | ||
| 697 | #[test] | |
| 698 | fn bytes_read_as_words() { | |
| 699 | assert_eq!(human_bytes(512), "512 B"); | |
| 700 | assert_eq!(human_bytes(5_000_000_000), "5.0 GB"); | |
| 701 | assert!(over_cap_message(6_000_000_000, &[("node_modules".into(), 4_000_000_000)]).contains("node_modules (4.0 GB)")); | |
| 702 | } | |
| 703 | ||
| 704 | #[cfg(unix)] | |
| 705 | #[test] | |
| 706 | fn a_command_streams_its_lines_then_how_it_ended() { | |
| 707 | let dir = std::env::temp_dir().join(format!("g1t-computer-{}", std::process::id())); | |
| 708 | let mut lines = Vec::new(); | |
| 709 | let last = run_streaming("echo one; echo two 1>&2; exit 3", &dir, Duration::from_secs(10), &BTreeMap::new(), &mut |line| { | |
| 710 | lines.push(line); | |
| 711 | Ok(()) | |
| 712 | }) | |
| 713 | .unwrap(); | |
| 714 | let _ = std::fs::remove_dir_all(&dir); | |
| 715 | assert_eq!(lines.len(), 3, "{lines:?}"); | |
| 716 | assert!(lines.iter().any(|l| l["stream"] == "stdout" && l["line"] == "one")); | |
| 717 | assert!(lines.iter().any(|l| l["stream"] == "stderr" && l["line"] == "two")); | |
| 718 | assert_eq!(last["exit_code"], 3); | |
| 719 | assert_eq!(last["timed_out"], false); | |
| 720 | assert_eq!(lines[2], last); | |
| 721 | } | |
| 722 | ||
| 723 | #[cfg(unix)] | |
| 724 | #[test] | |
| 725 | fn a_command_past_its_timeout_is_killed() { | |
| 726 | let dir = std::env::temp_dir().join(format!("g1t-computer-slow-{}", std::process::id())); | |
| 727 | let mut lines = Vec::new(); | |
| 728 | let last = run_streaming("echo start; sleep 30; echo never", &dir, Duration::from_millis(300), &BTreeMap::new(), &mut |line| { | |
| 729 | lines.push(line); | |
| 730 | Ok(()) | |
| 731 | }) | |
| 732 | .unwrap(); | |
| 733 | let _ = std::fs::remove_dir_all(&dir); | |
| 734 | assert_eq!(last["exit_code"], TIMED_OUT_EXIT); | |
| 735 | assert_eq!(last["timed_out"], true); | |
| 736 | assert!(!lines.iter().any(|l| l["line"] == "never")); | |
| 737 | } | |
| 738 | } |