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