Skip to content
738 linesCodeBlameRaw

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
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}