anvilsign in

collin/anvil

1//! Per-session supervision: bring a container up, keep a durable transcript,
2//! and wind everything down exactly once.
3
4use std::path::PathBuf;
5
6use anvil_core::{
7 App,
8 agent::{
9 self,
10 Handle,
11 status,
12 },
13 models::AgentSession,
14 repos,
15 secrets,
16 storage,
17 users,
18};
19use futures_util::StreamExt;
20use tokio::io::AsyncWriteExt;
21
22use crate::container;
23
24/// Where the entrypoint tees the pane, matching `session-entrypoint.sh`.
25const TRANSCRIPT_IN_CONTAINER: &str = "/tmp/anvil-transcript";
26
27/// Resolve a repository and starting ref to `(bare repo path, commit sha)`.
28pub async fn resolve_base(
29 app: &App,
30 repo_id: i64,
31 base_ref: &str,
32) -> Result<(PathBuf, String), String> {
33 let repo = repos::find_by_id(&app.db, repo_id)
34 .await
35 .map_err(|e| e.to_string())?
36 .ok_or("repository not found")?;
37 let owner = users::find_by_id(&app.db, repo.owner_id)
38 .await
39 .map_err(|e| e.to_string())?
40 .ok_or("owner not found")?;
41 let repo_path = storage::repo_path(&app.config.repositories_dir(), &owner.username, &repo.name);
42
43 let commit = anvil_git::browse::resolve_commit(&repo_path, base_ref)
44 .map_err(|e| format!("resolving {base_ref}: {e}"))?;
45 Ok((repo_path, commit))
46}
47
48/// Create, seed and start the session's container, then hand it to a
49/// supervisor task.
50///
51/// Returns once the container is running; the caller sends the user to the
52/// session page, which attaches over a websocket.
53pub async fn launch(app: App, session: AgentSession, repo_path: PathBuf) -> Result<(), String> {
54 let cfg = &app.config.agent;
55 let docker = anvil_docker::connect()?;
56 anvil_docker::ensure_image(&docker, &cfg.image).await?;
57
58 // The user secrets (secrets::UserSecret) this session opted into at
59 // start time — an empty choice is the normal case, not an error. Resolved
60 // before anything is created, so a missing unlock fails fast rather than
61 // after a container already exists. See docs/secrets.md: this is opt-in
62 // per session rather than blanket, on purpose — a session can be steered
63 // by repo content it reads (prompt injection, docs/untrusted-mode.md),
64 // so it should only ever reach what it was deliberately given.
65 let (env_secrets, file_writes, json_merges) = resolve_secrets(&app, &session).await?;
66
67 // The agent CLI, run interactively under tmux. `--dangerously-skip-
68 // permissions` is defensible precisely because the container *is* the
69 // sandbox: all capabilities dropped, no socket, no mounts, nothing of
70 // anvil's on the filesystem. A permission prompt inside a container the
71 // agent already fully owns buys nothing.
72 let command = vec![
73 "claude".to_string(),
74 "--dangerously-skip-permissions".to_string(),
75 ];
76
77 let container_id = container::create(&docker, cfg, session.id, &env_secrets, &command).await?;
78
79 // Seed the workspace. M1 uploads the commit's files, exactly as CI does —
80 // no `.git`, so no credential is needed anywhere yet.
81 let files = anvil_git::browse::read_tree_files(&repo_path, &session.base_commit)
82 .map_err(|e| format!("reading tree {}: {e}", session.base_commit))?;
83 let entries: Vec<_> = files
84 .into_iter()
85 .map(|f| (f.path, f.content, f.executable))
86 .collect();
87 let workspace_tar = container::build_tar(container::WORKDIR.trim_start_matches('/'), &entries);
88 if let Err(e) = container::upload(&docker, &container_id, workspace_tar).await {
89 container::remove(&docker, &container_id).await;
90 return Err(e);
91 }
92
93 // Seed the agent's credentials the same way — uploaded, never mounted, so
94 // the no-bind-mounts invariant in docs/untrusted-mode.md still holds.
95 if let Err(e) = seed_credentials(&docker, &container_id, cfg).await {
96 container::remove(&docker, &container_id).await;
97 return Err(e);
98 }
99
100 // "file"/"json" secrets, written into $HOME — after credentials_dir, so a
101 // session's own explicit choice wins if the two ever target the same
102 // path. A "json" merge may download what credentials_dir (or the image,
103 // e.g. claude-onboarding.json) just put there, hence after both and not
104 // interleaved with them.
105 if let Err(e) = write_secret_files(&docker, &container_id, file_writes, json_merges).await {
106 container::remove(&docker, &container_id).await;
107 return Err(e);
108 }
109
110 if let Err(e) = container::start(&docker, &container_id).await {
111 container::remove(&docker, &container_id).await;
112 return Err(e);
113 }
114
115 let (shutdown, shutdown_rx) = tokio::sync::watch::channel(false);
116 let handle = Handle {
117 container_id: container_id.clone(),
118 shutdown,
119 last_activity: std::sync::Arc::new(std::sync::atomic::AtomicI64::new(agent::now_secs())),
120 attached: std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0)),
121 };
122 app.sessions.insert(session.id, handle);
123
124 agent::mark_running(&app.db, session.id, &container_id)
125 .await
126 .map_err(|e| e.to_string())?;
127
128 let prompt = session.prompt.clone();
129 tokio::spawn(supervise(
130 app,
131 session.id,
132 container_id,
133 shutdown_rx,
134 prompt,
135 ));
136 Ok(())
137}
138
139/// Resolve `session.secret_names` into what `launch` actually needs: pairs
140/// to inject as environment variables, whole-file writes, and json-field
141/// merges — see [`secrets::kind`]. Empty everywhere is the normal case for a
142/// session that opted into nothing.
143#[allow(clippy::type_complexity)]
144async fn resolve_secrets(
145 app: &App,
146 session: &AgentSession,
147) -> Result<
148 (
149 Vec<(String, String)>,
150 Vec<(String, String)>,
151 Vec<(String, String, String)>,
152 ),
153 String,
154> {
155 let names: Vec<String> = session
156 .secret_names
157 .split(',')
158 .map(str::trim)
159 .filter(|s| !s.is_empty())
160 .map(str::to_string)
161 .collect();
162 if names.is_empty() {
163 return Ok((Vec::new(), Vec::new(), Vec::new()));
164 }
165
166 let values = app
167 .user_vault
168 .take(session.user_id, &names)
169 .map_err(|missing| {
170 format!(
171 "session wants secret(s) {} but they are sealed, or were unlocked without them — \
172 run `anvild secret user unlock`",
173 missing.join(", ")
174 )
175 })?;
176
177 let mut env_secrets = Vec::new();
178 let mut file_writes = Vec::new();
179 let mut json_merges = Vec::new();
180 for (name, value) in values {
181 // Deleted between opting in and this session actually launching:
182 // skip it rather than fail the whole session over a secret that no
183 // longer exists.
184 let Some(secret) = secrets::find_for_user(&app.db, session.user_id, &name)
185 .await
186 .map_err(|e| e.to_string())?
187 else {
188 continue;
189 };
190 match secret.kind.as_str() {
191 secrets::kind::FILE => file_writes.push((secret.dest_path, value)),
192 secrets::kind::JSON => json_merges.push((secret.dest_path, secret.field, value)),
193 _ => env_secrets.push((name, value)),
194 }
195 }
196 Ok((env_secrets, file_writes, json_merges))
197}
198
199/// Write "file"/"json" user secrets into the session's `$HOME`. Multiple
200/// entries targeting the same path compose against one shared in-memory copy
201/// — a "file" secret provides the base if one names that path, otherwise a
202/// "json" merge downloads whatever is there already (from the image or
203/// `seed_credentials`) — rather than each overwriting the last.
204async fn write_secret_files(
205 docker: &bollard::Docker,
206 container_id: &str,
207 file_writes: Vec<(String, String)>,
208 json_merges: Vec<(String, String, String)>,
209) -> Result<(), String> {
210 if file_writes.is_empty() && json_merges.is_empty() {
211 return Ok(());
212 }
213
214 let mut contents: std::collections::BTreeMap<String, Vec<u8>> =
215 std::collections::BTreeMap::new();
216 for (dest_path, value) in file_writes {
217 contents.insert(dest_path, value.into_bytes());
218 }
219 for (dest_path, field, value) in json_merges {
220 let current = match contents.get(&dest_path) {
221 Some(bytes) => bytes.clone(),
222 None => {
223 let absolute = format!("{}/{dest_path}", container::HOME);
224 container::download_file(docker, container_id, &absolute)
225 .await
226 .unwrap_or_default()
227 }
228 };
229 let merged =
230 anvil_core::secrets::json_merge(&current, &field, &value).map_err(|e| e.to_string())?;
231 contents.insert(dest_path, merged);
232 }
233
234 let entries: Vec<_> = contents
235 .into_iter()
236 .map(|(path, bytes)| (path, bytes, false))
237 .collect();
238 let tar = container::build_tar(container::HOME.trim_start_matches('/'), &entries);
239 container::upload(docker, container_id, tar).await
240}
241
242/// Upload the configured credentials directory to the session user's home.
243///
244/// Empty config means sessions start unauthenticated, which is a legitimate
245/// choice (and the default) rather than an error.
246async fn seed_credentials(
247 docker: &bollard::Docker,
248 container_id: &str,
249 cfg: &anvil_core::config::AgentConfig,
250) -> Result<(), String> {
251 if cfg.credentials_dir.as_os_str().is_empty() {
252 return Ok(());
253 }
254 if !cfg.credentials_dir.is_dir() {
255 return Err(format!(
256 "agent.credentials_dir {} is not a directory",
257 cfg.credentials_dir.display()
258 ));
259 }
260 let entries = container::read_dir_recursive(&cfg.credentials_dir)?;
261 if entries.is_empty() {
262 return Ok(());
263 }
264 let prefix = format!("{}/.claude", container::HOME.trim_start_matches('/'));
265 let tar = container::build_tar(&prefix, &entries);
266 container::upload(docker, container_id, tar).await
267}
268
269/// Own one session until it ends: keep the transcript, honour a shutdown
270/// request, and close the row out exactly once.
271async fn supervise(
272 app: App,
273 session_id: i64,
274 container_id: String,
275 mut shutdown: tokio::sync::watch::Receiver<bool>,
276 prompt: String,
277) {
278 let docker = match anvil_docker::connect() {
279 Ok(docker) => docker,
280 Err(e) => {
281 finish(&app, session_id, &container_id, status::FAILED, 0, &e).await;
282 return;
283 }
284 };
285
286 // An autonomous session is the *interactive* agent with its prompt typed
287 // at it, rather than a headless run. That way a human can take over
288 // mid-flight just by attaching — there is no separate mode to bail out of.
289 if !prompt.is_empty() {
290 let docker = docker.clone();
291 let container_id = container_id.clone();
292 let prompt = prompt.clone();
293 tokio::spawn(async move {
294 // Let the CLI finish drawing before typing into it.
295 tokio::time::sleep(std::time::Duration::from_secs(5)).await;
296 if let Err(e) = container::send_keys(&docker, &container_id, &prompt).await {
297 tracing::warn!("agent: sending prompt to session {session_id}: {e}");
298 }
299 });
300 }
301
302 let transcript = tokio::spawn(pump_transcript(
303 docker.clone(),
304 app.clone(),
305 session_id,
306 container_id.clone(),
307 ));
308
309 // Wait for whichever comes first: the container exiting on its own, or a
310 // shutdown request from the sweep / an operator.
311 let mut wait = docker.wait_container(
312 &container_id,
313 None::<bollard::container::WaitContainerOptions<String>>,
314 );
315
316 let (final_status, exit_code) = tokio::select! {
317 result = wait.next() => match result {
318 Some(Ok(response)) => (status::EXITED, response.status_code),
319 Some(Err(e)) => {
320 tracing::warn!("agent: waiting on session {session_id}: {e}");
321 (status::FAILED, 0)
322 }
323 None => (status::EXITED, 0),
324 },
325 _ = shutdown.changed() => (status::REAPED, 0),
326 };
327
328 transcript.abort();
329 finish(&app, session_id, &container_id, final_status, exit_code, "").await;
330}
331
332/// Follow the in-container transcript and append it to the session's file on
333/// disk, so a finished session can be replayed without the container.
334///
335/// Reads the pipe-pane output rather than attaching a tmux client: an extra
336/// client would take part in tmux's window sizing and shrink everyone else's
337/// terminal.
338async fn pump_transcript(docker: bollard::Docker, app: App, session_id: i64, container_id: String) {
339 let dir = app.config.sessions_dir();
340 if let Err(e) = tokio::fs::create_dir_all(&dir).await {
341 tracing::warn!("agent: creating {}: {e}", dir.display());
342 return;
343 }
344 let path = storage::session_transcript_path(&dir, session_id);
345 let mut file = match tokio::fs::File::create(&path).await {
346 Ok(file) => file,
347 Err(e) => {
348 tracing::warn!("agent: creating {}: {e}", path.display());
349 return;
350 }
351 };
352
353 // `tail -F` from the start, so nothing emitted before this task got going
354 // is lost and a not-yet-created file is not fatal.
355 let exec = docker
356 .create_exec(
357 &container_id,
358 bollard::exec::CreateExecOptions {
359 attach_stdout: Some(true),
360 attach_stderr: Some(false),
361 tty: Some(false),
362 user: Some(container::RUN_AS.to_string()),
363 cmd: Some(vec![
364 "tail".to_string(),
365 "-n".to_string(),
366 "+1".to_string(),
367 "-F".to_string(),
368 TRANSCRIPT_IN_CONTAINER.to_string(),
369 ]),
370 ..Default::default()
371 },
372 )
373 .await;
374 let Ok(exec) = exec else {
375 tracing::warn!("agent: transcript exec for session {session_id} failed");
376 return;
377 };
378
379 let Ok(bollard::exec::StartExecResults::Attached { mut output, .. }) =
380 docker.start_exec(&exec.id, None).await
381 else {
382 return;
383 };
384
385 let handle = app.sessions.get(session_id);
386 while let Some(chunk) = output.next().await {
387 let Ok(log) = chunk else { break };
388 let bytes = log.into_bytes();
389 if bytes.is_empty() {
390 continue;
391 }
392 // Output counts as activity: an agent thinking out loud for an hour
393 // with nobody watching is working, not idle.
394 if let Some(handle) = &handle {
395 handle.touch();
396 }
397 if let Err(e) = file.write_all(bytes.as_ref()).await {
398 tracing::warn!("agent: writing transcript for session {session_id}: {e}");
399 break;
400 }
401 let _ = file.flush().await;
402 }
403}
404
405/// Close a session out: remove the container, drop the handle, stamp the row.
406async fn finish(
407 app: &App,
408 session_id: i64,
409 container_id: &str,
410 final_status: &str,
411 exit_code: i64,
412 error: &str,
413) {
414 if let Ok(docker) = anvil_docker::connect() {
415 container::remove(&docker, container_id).await;
416 }
417 app.sessions.remove(session_id);
418 if let Err(e) = agent::finish(&app.db, session_id, final_status, exit_code, error).await {
419 tracing::error!("agent: closing out session {session_id}: {e}");
420 }
421 tracing::info!("agent: session {session_id} {final_status}");
422}
423
424/// Attach a browser to a session's tmux, returning the duplex terminal.
425///
426/// Bumps the attached count for the idle sweep; the caller must call
427/// [`detached`] when the websocket closes.
428pub async fn attach_stream(
429 app: &App,
430 session_id: i64,
431 cols: u16,
432 rows: u16,
433) -> Result<(container::Terminal, bollard::Docker), String> {
434 let handle = app
435 .sessions
436 .get(session_id)
437 .ok_or("session is not running")?;
438
439 let docker = anvil_docker::connect()?;
440 let terminal = container::attach(&docker, &handle.container_id, cols, rows).await?;
441
442 handle
443 .attached
444 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
445 handle.touch();
446 let _ = agent::touch_attach(&app.db, session_id).await;
447
448 Ok((terminal, docker))
449}
450
451/// Note that a browser disconnected.
452pub fn detached(app: &App, session_id: i64) {
453 if let Some(handle) = app.sessions.get(session_id) {
454 handle
455 .attached
456 .fetch_sub(1, std::sync::atomic::Ordering::Relaxed);
457 handle.touch();
458 }
459}