Skip to content
738 linesCodeBlameRaw
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
34use std::collections::BTreeMap;
35use std::io::{self, BufRead, BufReader, Read, Write};
36use std::net::{TcpListener, TcpStream};
37use std::path::{Path, PathBuf};
38use std::process::{Child, Command, Stdio};
39use std::sync::atomic::{AtomicUsize, Ordering};
40use std::sync::mpsc;
41use std::sync::Arc;
42use std::time::{Duration, Instant};
43
44use serde_json::{Value, json};
45
46use crate::docker::http::{self, Body, Head};
47
48/// Where the agent lives: its files, clones, caches and notes.
49pub(crate) const HOME: &str = "/home/agent";
50/// Where a session's work goes by default, under the home.
51pub(crate) const SESSIONS_DIR: &str = "sessions";
52/// The port the Durable Object reaches the computer on.
53pub(crate) const PORT: u16 = 8787;
54/// How much the home may hold: 5 GB, a hard cap in v1 (no disk charge).
55pub(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.
57pub(crate) const MAX_OUTPUT_BYTES: usize = 1024 * 1024;
58/// The most a file read or written through `/files` may be.
59pub(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.
61pub(crate) const DEFAULT_TIMEOUT_SECONDS: u64 = 120;
62pub(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.
64pub(crate) const TIMED_OUT_EXIT: i32 = 124;
65
66/// What the process knows while it serves.
67struct 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`.
75pub 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.
116fn 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
128fn 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.
136fn 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.
156fn 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.
161fn 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
277fn 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
282fn 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.
293pub(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
303fn 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.
341pub(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.
362pub(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.
370fn 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.
395pub(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.
400pub(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.
416pub(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.
429enum 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.
436pub(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.
520fn 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.
534pub(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.
547pub(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.
562pub(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.
576pub(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
586pub(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.
600fn 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)]
629mod 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}