| 1 | //! Docker specifics for an agent session: creating the container, seeding it, |
| 2 | //! attaching a terminal to the tmux session inside, and tearing it down. |
| 3 | //! |
| 4 | //! The containment story is deliberately the CI one (see |
| 5 | //! `docs/untrusted-mode.md` §1): all capabilities dropped, `no-new-privileges`, |
| 6 | //! pids/memory/cpu caps, **no Docker socket, no bind mounts, no volumes**. |
| 7 | //! Everything the container needs — the checkout and the agent's credentials — |
| 8 | //! is uploaded as a tar through the Docker API, so nothing on the anvil host is |
| 9 | //! ever exposed to it. |
| 10 | |
| 11 | use anvil_core::config::AgentConfig; |
| 12 | use bollard::{ |
| 13 | Docker, |
| 14 | container::{ |
| 15 | Config, |
| 16 | CreateContainerOptions, |
| 17 | RemoveContainerOptions, |
| 18 | StartContainerOptions, |
| 19 | UploadToContainerOptions, |
| 20 | }, |
| 21 | exec::{ |
| 22 | CreateExecOptions, |
| 23 | ResizeExecOptions, |
| 24 | StartExecOptions, |
| 25 | StartExecResults, |
| 26 | }, |
| 27 | models::HostConfig, |
| 28 | }; |
| 29 | |
| 30 | /// Label carrying the session id, so containers can be found again after a |
| 31 | /// restart — the registry is in-memory and does not survive one. |
| 32 | pub const LABEL_SESSION: &str = "anvil.session"; |
| 33 | |
| 34 | /// Working directory inside the container, matching the CI runner's. |
| 35 | pub const WORKDIR: &str = "/workspace"; |
| 36 | |
| 37 | /// The unprivileged user `deploy/runner/Dockerfile` creates. Sessions run as |
| 38 | /// this rather than root; CI keeps the image's default (root) because plenty of |
| 39 | /// pipelines expect to `apt-get`. |
| 40 | pub const RUN_AS: &str = "agent"; |
| 41 | |
| 42 | /// That user's uid/gid, needed when building tars so the uploaded files are |
| 43 | /// owned by the account that has to write them. |
| 44 | const RUN_AS_UID: u64 = 1000; |
| 45 | |
| 46 | /// Home directory of [`RUN_AS`]; the agent CLI's config lives under it. |
| 47 | pub const HOME: &str = "/home/agent"; |
| 48 | |
| 49 | /// tmux session name, matching `session-entrypoint.sh`. |
| 50 | pub const TMUX_SESSION: &str = "agent"; |
| 51 | |
| 52 | /// Create (but do not start) a session container. |
| 53 | pub async fn create( |
| 54 | docker: &Docker, |
| 55 | cfg: &AgentConfig, |
| 56 | session_id: i64, |
| 57 | env: &[(String, String)], |
| 58 | command: &[String], |
| 59 | ) -> Result<String, String> { |
| 60 | // The same sandbox CI uses. Limits of 0 mean "unlimited" and omit the cap. |
| 61 | // |
| 62 | // Note there is no `network_mode` opt-out: unlike a CI job, a session is |
| 63 | // useless without network — it has to reach the model API, and from M2 it |
| 64 | // clones and pushes back to anvil. |
| 65 | let host_config = HostConfig { |
| 66 | cap_drop: Some(vec!["ALL".to_string()]), |
| 67 | security_opt: Some(vec!["no-new-privileges:true".to_string()]), |
| 68 | pids_limit: (cfg.pids_limit > 0).then_some(cfg.pids_limit), |
| 69 | memory: (cfg.memory_mb > 0).then(|| cfg.memory_mb * 1024 * 1024), |
| 70 | memory_swap: (cfg.memory_mb > 0).then(|| cfg.memory_mb * 1024 * 1024), |
| 71 | nano_cpus: (cfg.cpus > 0.0).then_some((cfg.cpus * 1e9) as i64), |
| 72 | ..Default::default() |
| 73 | }; |
| 74 | |
| 75 | let mut cmd = vec!["anvil-session".to_string()]; |
| 76 | cmd.extend(command.iter().cloned()); |
| 77 | |
| 78 | let config = Config { |
| 79 | image: Some(cfg.image.clone()), |
| 80 | cmd: Some(cmd), |
| 81 | env: Some( |
| 82 | env.iter() |
| 83 | .map(|(k, v)| format!("{k}={v}")) |
| 84 | .chain([format!("HOME={HOME}"), "TERM=xterm-256color".to_string()]) |
| 85 | .collect(), |
| 86 | ), |
| 87 | working_dir: Some(WORKDIR.to_string()), |
| 88 | user: Some(RUN_AS.to_string()), |
| 89 | labels: Some( |
| 90 | [(LABEL_SESSION.to_string(), session_id.to_string())] |
| 91 | .into_iter() |
| 92 | .collect(), |
| 93 | ), |
| 94 | // A terminal program needs a tty even before anyone attaches, or tmux |
| 95 | // starts with a 80x24 dumb terminal and never recovers. |
| 96 | tty: Some(true), |
| 97 | open_stdin: Some(true), |
| 98 | host_config: Some(host_config), |
| 99 | ..Default::default() |
| 100 | }; |
| 101 | |
| 102 | let created = docker |
| 103 | .create_container(None::<CreateContainerOptions<String>>, config) |
| 104 | .await |
| 105 | .map_err(|e| format!("create session container: {e}"))?; |
| 106 | Ok(created.id) |
| 107 | } |
| 108 | |
| 109 | /// Upload a tar into the container, rooted at `/`. |
| 110 | pub async fn upload(docker: &Docker, id: &str, tar: Vec<u8>) -> Result<(), String> { |
| 111 | docker |
| 112 | .upload_to_container( |
| 113 | id, |
| 114 | Some(UploadToContainerOptions { |
| 115 | path: "/".to_string(), |
| 116 | ..Default::default() |
| 117 | }), |
| 118 | tar.into(), |
| 119 | ) |
| 120 | .await |
| 121 | .map_err(|e| format!("upload to session container: {e}")) |
| 122 | } |
| 123 | |
| 124 | /// Build a tar of `(path, contents, executable)` entries, rooted at `prefix` |
| 125 | /// (no leading slash) and owned by the session user. |
| 126 | pub fn build_tar(prefix: &str, files: &[(String, Vec<u8>, bool)]) -> Vec<u8> { |
| 127 | let mut builder = tar::Builder::new(Vec::new()); |
| 128 | for (path, content, executable) in files { |
| 129 | let mut header = tar::Header::new_gnu(); |
| 130 | header.set_size(content.len() as u64); |
| 131 | header.set_mode(if *executable { 0o755 } else { 0o644 }); |
| 132 | // Ownership matters: the container runs as `agent`, and a checkout it |
| 133 | // cannot write is not a workspace. |
| 134 | header.set_uid(RUN_AS_UID); |
| 135 | header.set_gid(RUN_AS_UID); |
| 136 | header.set_cksum(); |
| 137 | let full = format!("{prefix}/{path}"); |
| 138 | if builder |
| 139 | .append_data(&mut header, &full, content.as_slice()) |
| 140 | .is_err() |
| 141 | { |
| 142 | continue; |
| 143 | } |
| 144 | } |
| 145 | builder.into_inner().unwrap_or_default() |
| 146 | } |
| 147 | |
| 148 | /// Read a host directory into `(relative path, bytes, executable)` entries. |
| 149 | /// |
| 150 | /// Used for the agent CLI's credentials directory, which is uploaded rather |
| 151 | /// than bind-mounted so the no-mounts invariant survives. |
| 152 | pub fn read_dir_recursive(root: &std::path::Path) -> Result<Vec<(String, Vec<u8>, bool)>, String> { |
| 153 | fn walk( |
| 154 | base: &std::path::Path, |
| 155 | dir: &std::path::Path, |
| 156 | out: &mut Vec<(String, Vec<u8>, bool)>, |
| 157 | ) -> std::io::Result<()> { |
| 158 | for entry in std::fs::read_dir(dir)? { |
| 159 | let entry = entry?; |
| 160 | let path = entry.path(); |
| 161 | let meta = entry.metadata()?; |
| 162 | if meta.is_dir() { |
| 163 | walk(base, &path, out)?; |
| 164 | } else if meta.is_file() { |
| 165 | let Ok(rel) = path.strip_prefix(base) else { |
| 166 | continue; |
| 167 | }; |
| 168 | let executable = { |
| 169 | #[cfg(unix)] |
| 170 | { |
| 171 | use std::os::unix::fs::PermissionsExt; |
| 172 | meta.permissions().mode() & 0o111 != 0 |
| 173 | } |
| 174 | #[cfg(not(unix))] |
| 175 | { |
| 176 | false |
| 177 | } |
| 178 | }; |
| 179 | out.push(( |
| 180 | rel.to_string_lossy().replace('\\', "/"), |
| 181 | std::fs::read(&path)?, |
| 182 | executable, |
| 183 | )); |
| 184 | } |
| 185 | } |
| 186 | Ok(()) |
| 187 | } |
| 188 | |
| 189 | let mut out = Vec::new(); |
| 190 | walk(root, root, &mut out).map_err(|e| format!("reading {}: {e}", root.display()))?; |
| 191 | Ok(out) |
| 192 | } |
| 193 | |
| 194 | /// Start a created container. |
| 195 | pub async fn start(docker: &Docker, id: &str) -> Result<(), String> { |
| 196 | docker |
| 197 | .start_container(id, None::<StartContainerOptions<String>>) |
| 198 | .await |
| 199 | .map_err(|e| format!("start session container: {e}")) |
| 200 | } |
| 201 | |
| 202 | /// A live terminal attached to the container's tmux session. |
| 203 | pub struct Terminal { |
| 204 | /// Docker exec id, needed to resize the pty. |
| 205 | pub exec_id: String, |
| 206 | /// Bytes the terminal produces. |
| 207 | pub output: std::pin::Pin< |
| 208 | Box< |
| 209 | dyn futures_util::Stream< |
| 210 | Item = Result<bollard::container::LogOutput, bollard::errors::Error>, |
| 211 | > + Send, |
| 212 | >, |
| 213 | >, |
| 214 | /// Keystrokes go here. |
| 215 | pub input: std::pin::Pin<Box<dyn tokio::io::AsyncWrite + Send>>, |
| 216 | } |
| 217 | |
| 218 | /// Attach a new tmux client to the container's session. |
| 219 | /// |
| 220 | /// Each caller gets its own `docker exec`, which is the point: a dropped |
| 221 | /// websocket kills that client only, never the agent, because the agent is a |
| 222 | /// process inside tmux rather than a child of the exec. |
| 223 | pub async fn attach(docker: &Docker, id: &str, cols: u16, rows: u16) -> Result<Terminal, String> { |
| 224 | let exec = docker |
| 225 | .create_exec( |
| 226 | id, |
| 227 | CreateExecOptions { |
| 228 | attach_stdin: Some(true), |
| 229 | attach_stdout: Some(true), |
| 230 | attach_stderr: Some(true), |
| 231 | tty: Some(true), |
| 232 | user: Some(RUN_AS.to_string()), |
| 233 | env: Some(vec![ |
| 234 | "TERM=xterm-256color".to_string(), |
| 235 | format!("HOME={HOME}"), |
| 236 | ]), |
| 237 | cmd: Some(vec![ |
| 238 | "tmux".to_string(), |
| 239 | "-f".to_string(), |
| 240 | "/etc/anvil/tmux.conf".to_string(), |
| 241 | "attach-session".to_string(), |
| 242 | "-t".to_string(), |
| 243 | TMUX_SESSION.to_string(), |
| 244 | ]), |
| 245 | ..Default::default() |
| 246 | }, |
| 247 | ) |
| 248 | .await |
| 249 | .map_err(|e| format!("create attach exec: {e}"))?; |
| 250 | |
| 251 | let started = docker |
| 252 | .start_exec( |
| 253 | &exec.id, |
| 254 | Some(StartExecOptions { |
| 255 | detach: false, |
| 256 | tty: true, |
| 257 | ..Default::default() |
| 258 | }), |
| 259 | ) |
| 260 | .await |
| 261 | .map_err(|e| format!("start attach exec: {e}"))?; |
| 262 | |
| 263 | let StartExecResults::Attached { output, input } = started else { |
| 264 | return Err("attach exec detached unexpectedly".to_string()); |
| 265 | }; |
| 266 | |
| 267 | // Size the pty before the first byte, so the TUI lays out correctly rather |
| 268 | // than redrawing from an 80x24 assumption. |
| 269 | let terminal = Terminal { |
| 270 | exec_id: exec.id, |
| 271 | output, |
| 272 | input, |
| 273 | }; |
| 274 | resize(docker, &terminal.exec_id, cols, rows).await; |
| 275 | Ok(terminal) |
| 276 | } |
| 277 | |
| 278 | /// Resize an attached terminal. Best-effort: a resize racing the exec's exit is |
| 279 | /// routine and not worth failing an attach over. |
| 280 | pub async fn resize(docker: &Docker, exec_id: &str, cols: u16, rows: u16) { |
| 281 | if cols == 0 || rows == 0 { |
| 282 | return; |
| 283 | } |
| 284 | if let Err(e) = docker |
| 285 | .resize_exec( |
| 286 | exec_id, |
| 287 | ResizeExecOptions { |
| 288 | height: rows, |
| 289 | width: cols, |
| 290 | }, |
| 291 | ) |
| 292 | .await |
| 293 | { |
| 294 | tracing::debug!("resize exec {exec_id}: {e}"); |
| 295 | } |
| 296 | } |
| 297 | |
| 298 | /// Run a one-shot command in the container and collect its stdout. |
| 299 | /// |
| 300 | /// This is the `capture-pane` / `send-keys` path: short, non-interactive tmux |
| 301 | /// control commands rather than a terminal. |
| 302 | pub async fn exec_capture(docker: &Docker, id: &str, cmd: &[&str]) -> Result<Vec<u8>, String> { |
| 303 | use futures_util::StreamExt; |
| 304 | |
| 305 | let exec = docker |
| 306 | .create_exec( |
| 307 | id, |
| 308 | CreateExecOptions { |
| 309 | attach_stdout: Some(true), |
| 310 | attach_stderr: Some(true), |
| 311 | tty: Some(false), |
| 312 | user: Some(RUN_AS.to_string()), |
| 313 | env: Some(vec![format!("HOME={HOME}")]), |
| 314 | cmd: Some(cmd.iter().map(|s| s.to_string()).collect()), |
| 315 | ..Default::default() |
| 316 | }, |
| 317 | ) |
| 318 | .await |
| 319 | .map_err(|e| format!("create exec {cmd:?}: {e}"))?; |
| 320 | |
| 321 | let started = docker |
| 322 | .start_exec(&exec.id, None) |
| 323 | .await |
| 324 | .map_err(|e| format!("start exec {cmd:?}: {e}"))?; |
| 325 | |
| 326 | let StartExecResults::Attached { mut output, .. } = started else { |
| 327 | return Ok(Vec::new()); |
| 328 | }; |
| 329 | |
| 330 | let mut buf = Vec::new(); |
| 331 | while let Some(chunk) = output.next().await { |
| 332 | match chunk { |
| 333 | Ok(log) => buf.extend_from_slice(log.into_bytes().as_ref()), |
| 334 | Err(e) => return Err(format!("exec {cmd:?}: {e}")), |
| 335 | } |
| 336 | } |
| 337 | Ok(buf) |
| 338 | } |
| 339 | |
| 340 | /// The tmux pane's current screen plus scrollback, with escape sequences kept. |
| 341 | /// |
| 342 | /// This is what a reconnecting browser is sent before the live stream, so it |
| 343 | /// resumes with the screen it left rather than a blank one. |
| 344 | pub async fn capture_pane(docker: &Docker, id: &str) -> Result<Vec<u8>, String> { |
| 345 | exec_capture( |
| 346 | docker, |
| 347 | id, |
| 348 | &[ |
| 349 | "tmux", |
| 350 | "capture-pane", |
| 351 | "-p", |
| 352 | "-e", |
| 353 | "-S", |
| 354 | "-", |
| 355 | "-t", |
| 356 | TMUX_SESSION, |
| 357 | ], |
| 358 | ) |
| 359 | .await |
| 360 | } |
| 361 | |
| 362 | /// Type a line into the tmux session, as if the user had. |
| 363 | /// |
| 364 | /// This is how an autonomous session is driven: rather than a headless |
| 365 | /// `-p` invocation, the same interactive agent gets its prompt typed at it, so |
| 366 | /// a human can take over mid-run just by attaching. |
| 367 | pub async fn send_keys(docker: &Docker, id: &str, text: &str) -> Result<(), String> { |
| 368 | exec_capture( |
| 369 | docker, |
| 370 | id, |
| 371 | &["tmux", "send-keys", "-t", TMUX_SESSION, text, "Enter"], |
| 372 | ) |
| 373 | .await |
| 374 | .map(|_| ()) |
| 375 | } |
| 376 | |
| 377 | /// Force-remove the container. Best-effort; a container already gone is a |
| 378 | /// success as far as callers are concerned. |
| 379 | pub async fn remove(docker: &Docker, id: &str) { |
| 380 | if let Err(e) = docker |
| 381 | .remove_container( |
| 382 | id, |
| 383 | Some(RemoveContainerOptions { |
| 384 | force: true, |
| 385 | ..Default::default() |
| 386 | }), |
| 387 | ) |
| 388 | .await |
| 389 | { |
| 390 | tracing::debug!("removing session container {id}: {e}"); |
| 391 | } |
| 392 | } |
| 393 | |
| 394 | /// Session containers Docker still knows about, as `(container id, session id)`. |
| 395 | /// |
| 396 | /// Startup uses this to reconcile rows against reality: the registry is |
| 397 | /// in-memory, so after a restart every one of these is an orphan. |
| 398 | pub async fn list_sessions(docker: &Docker) -> Result<Vec<(String, i64)>, String> { |
| 399 | use bollard::container::ListContainersOptions; |
| 400 | |
| 401 | let containers = docker |
| 402 | .list_containers(Some(ListContainersOptions::<String> { |
| 403 | all: true, |
| 404 | filters: [("label".to_string(), vec![LABEL_SESSION.to_string()])] |
| 405 | .into_iter() |
| 406 | .collect(), |
| 407 | ..Default::default() |
| 408 | })) |
| 409 | .await |
| 410 | .map_err(|e| format!("listing session containers: {e}"))?; |
| 411 | |
| 412 | Ok(containers |
| 413 | .into_iter() |
| 414 | .filter_map(|c| { |
| 415 | let id = c.id?; |
| 416 | let session_id = c.labels?.get(LABEL_SESSION)?.parse().ok()?; |
| 417 | Some((id, session_id)) |
| 418 | }) |
| 419 | .collect()) |
| 420 | } |