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 storage,
16 users,
17};
18use futures_util::StreamExt;
19use tokio::io::AsyncWriteExt;
20
21use crate::container;
22
23/// Where the entrypoint tees the pane, matching `session-entrypoint.sh`.
24const TRANSCRIPT_IN_CONTAINER: &str = "/tmp/anvil-transcript";
25
26/// Resolve a repository and starting ref to `(bare repo path, commit sha)`.
27pub async fn resolve_base(
28 app: &App,
29 repo_id: i64,
30 base_ref: &str,
31) -> Result<(PathBuf, String), String> {
32 let repo = repos::find_by_id(&app.db, repo_id)
33 .await
34 .map_err(|e| e.to_string())?
35 .ok_or("repository not found")?;
36 let owner = users::find_by_id(&app.db, repo.owner_id)
37 .await
38 .map_err(|e| e.to_string())?
39 .ok_or("owner not found")?;
40 let repo_path = storage::repo_path(&app.config.repositories_dir(), &owner.username, &repo.name);
41
42 let commit = anvil_git::browse::resolve_commit(&repo_path, base_ref)
43 .map_err(|e| format!("resolving {base_ref}: {e}"))?;
44 Ok((repo_path, commit))
45}
46
47/// Create, seed and start the session's container, then hand it to a
48/// supervisor task.
49///
50/// Returns once the container is running; the caller sends the user to the
51/// session page, which attaches over a websocket.
52pub async fn launch(app: App, session: AgentSession, repo_path: PathBuf) -> Result<(), String> {
53 let cfg = &app.config.agent;
54 let docker = anvil_ci::docker::connect()?;
55 anvil_ci::docker::ensure_image(&docker, &cfg.image).await?;
56
57 // The agent CLI, run interactively under tmux. `--dangerously-skip-
58 // permissions` is defensible precisely because the container *is* the
59 // sandbox: all capabilities dropped, no socket, no mounts, nothing of
60 // anvil's on the filesystem. A permission prompt inside a container the
61 // agent already fully owns buys nothing.
62 let command = vec![
63 "claude".to_string(),
64 "--dangerously-skip-permissions".to_string(),
65 ];
66
67 let container_id = container::create(&docker, cfg, session.id, &[], &command).await?;
68
69 // Seed the workspace. M1 uploads the commit's files, exactly as CI does —
70 // no `.git`, so no credential is needed anywhere yet.
71 let files = anvil_git::browse::read_tree_files(&repo_path, &session.base_commit)
72 .map_err(|e| format!("reading tree {}: {e}", session.base_commit))?;
73 let entries: Vec<_> = files
74 .into_iter()
75 .map(|f| (f.path, f.content, f.executable))
76 .collect();
77 let workspace_tar = container::build_tar(container::WORKDIR.trim_start_matches('/'), &entries);
78 if let Err(e) = container::upload(&docker, &container_id, workspace_tar).await {
79 container::remove(&docker, &container_id).await;
80 return Err(e);
81 }
82
83 // Seed the agent's credentials the same way — uploaded, never mounted, so
84 // the no-bind-mounts invariant in docs/untrusted-mode.md still holds.
85 if let Err(e) = seed_credentials(&docker, &container_id, cfg).await {
86 container::remove(&docker, &container_id).await;
87 return Err(e);
88 }
89
90 if let Err(e) = container::start(&docker, &container_id).await {
91 container::remove(&docker, &container_id).await;
92 return Err(e);
93 }
94
95 let (shutdown, shutdown_rx) = tokio::sync::watch::channel(false);
96 let handle = Handle {
97 container_id: container_id.clone(),
98 shutdown,
99 last_activity: std::sync::Arc::new(std::sync::atomic::AtomicI64::new(agent::now_secs())),
100 attached: std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0)),
101 };
102 app.sessions.insert(session.id, handle);
103
104 agent::mark_running(&app.db, session.id, &container_id)
105 .await
106 .map_err(|e| e.to_string())?;
107
108 let prompt = session.prompt.clone();
109 tokio::spawn(supervise(
110 app,
111 session.id,
112 container_id,
113 shutdown_rx,
114 prompt,
115 ));
116 Ok(())
117}
118
119/// Upload the configured credentials directory to the session user's home.
120///
121/// Empty config means sessions start unauthenticated, which is a legitimate
122/// choice (and the default) rather than an error.
123async fn seed_credentials(
124 docker: &bollard::Docker,
125 container_id: &str,
126 cfg: &anvil_core::config::AgentConfig,
127) -> Result<(), String> {
128 if cfg.credentials_dir.as_os_str().is_empty() {
129 return Ok(());
130 }
131 if !cfg.credentials_dir.is_dir() {
132 return Err(format!(
133 "agent.credentials_dir {} is not a directory",
134 cfg.credentials_dir.display()
135 ));
136 }
137 let entries = container::read_dir_recursive(&cfg.credentials_dir)?;
138 if entries.is_empty() {
139 return Ok(());
140 }
141 let prefix = format!("{}/.claude", container::HOME.trim_start_matches('/'));
142 let tar = container::build_tar(&prefix, &entries);
143 container::upload(docker, container_id, tar).await
144}
145
146/// Own one session until it ends: keep the transcript, honour a shutdown
147/// request, and close the row out exactly once.
148async fn supervise(
149 app: App,
150 session_id: i64,
151 container_id: String,
152 mut shutdown: tokio::sync::watch::Receiver<bool>,
153 prompt: String,
154) {
155 let docker = match anvil_ci::docker::connect() {
156 Ok(docker) => docker,
157 Err(e) => {
158 finish(&app, session_id, &container_id, status::FAILED, 0, &e).await;
159 return;
160 }
161 };
162
163 // An autonomous session is the *interactive* agent with its prompt typed
164 // at it, rather than a headless run. That way a human can take over
165 // mid-flight just by attaching — there is no separate mode to bail out of.
166 if !prompt.is_empty() {
167 let docker = docker.clone();
168 let container_id = container_id.clone();
169 let prompt = prompt.clone();
170 tokio::spawn(async move {
171 // Let the CLI finish drawing before typing into it.
172 tokio::time::sleep(std::time::Duration::from_secs(5)).await;
173 if let Err(e) = container::send_keys(&docker, &container_id, &prompt).await {
174 tracing::warn!("agent: sending prompt to session {session_id}: {e}");
175 }
176 });
177 }
178
179 let transcript = tokio::spawn(pump_transcript(
180 docker.clone(),
181 app.clone(),
182 session_id,
183 container_id.clone(),
184 ));
185
186 // Wait for whichever comes first: the container exiting on its own, or a
187 // shutdown request from the sweep / an operator.
188 let mut wait = docker.wait_container(
189 &container_id,
190 None::<bollard::container::WaitContainerOptions<String>>,
191 );
192
193 let (final_status, exit_code) = tokio::select! {
194 result = wait.next() => match result {
195 Some(Ok(response)) => (status::EXITED, response.status_code),
196 Some(Err(e)) => {
197 tracing::warn!("agent: waiting on session {session_id}: {e}");
198 (status::FAILED, 0)
199 }
200 None => (status::EXITED, 0),
201 },
202 _ = shutdown.changed() => (status::REAPED, 0),
203 };
204
205 transcript.abort();
206 finish(&app, session_id, &container_id, final_status, exit_code, "").await;
207}
208
209/// Follow the in-container transcript and append it to the session's file on
210/// disk, so a finished session can be replayed without the container.
211///
212/// Reads the pipe-pane output rather than attaching a tmux client: an extra
213/// client would take part in tmux's window sizing and shrink everyone else's
214/// terminal.
215async fn pump_transcript(docker: bollard::Docker, app: App, session_id: i64, container_id: String) {
216 let dir = app.config.sessions_dir();
217 if let Err(e) = tokio::fs::create_dir_all(&dir).await {
218 tracing::warn!("agent: creating {}: {e}", dir.display());
219 return;
220 }
221 let path = storage::session_transcript_path(&dir, session_id);
222 let mut file = match tokio::fs::File::create(&path).await {
223 Ok(file) => file,
224 Err(e) => {
225 tracing::warn!("agent: creating {}: {e}", path.display());
226 return;
227 }
228 };
229
230 // `tail -F` from the start, so nothing emitted before this task got going
231 // is lost and a not-yet-created file is not fatal.
232 let exec = docker
233 .create_exec(
234 &container_id,
235 bollard::exec::CreateExecOptions {
236 attach_stdout: Some(true),
237 attach_stderr: Some(false),
238 tty: Some(false),
239 user: Some(container::RUN_AS.to_string()),
240 cmd: Some(vec![
241 "tail".to_string(),
242 "-n".to_string(),
243 "+1".to_string(),
244 "-F".to_string(),
245 TRANSCRIPT_IN_CONTAINER.to_string(),
246 ]),
247 ..Default::default()
248 },
249 )
250 .await;
251 let Ok(exec) = exec else {
252 tracing::warn!("agent: transcript exec for session {session_id} failed");
253 return;
254 };
255
256 let Ok(bollard::exec::StartExecResults::Attached { mut output, .. }) =
257 docker.start_exec(&exec.id, None).await
258 else {
259 return;
260 };
261
262 let handle = app.sessions.get(session_id);
263 while let Some(chunk) = output.next().await {
264 let Ok(log) = chunk else { break };
265 let bytes = log.into_bytes();
266 if bytes.is_empty() {
267 continue;
268 }
269 // Output counts as activity: an agent thinking out loud for an hour
270 // with nobody watching is working, not idle.
271 if let Some(handle) = &handle {
272 handle.touch();
273 }
274 if let Err(e) = file.write_all(bytes.as_ref()).await {
275 tracing::warn!("agent: writing transcript for session {session_id}: {e}");
276 break;
277 }
278 let _ = file.flush().await;
279 }
280}
281
282/// Close a session out: remove the container, drop the handle, stamp the row.
283async fn finish(
284 app: &App,
285 session_id: i64,
286 container_id: &str,
287 final_status: &str,
288 exit_code: i64,
289 error: &str,
290) {
291 if let Ok(docker) = anvil_ci::docker::connect() {
292 container::remove(&docker, container_id).await;
293 }
294 app.sessions.remove(session_id);
295 if let Err(e) = agent::finish(&app.db, session_id, final_status, exit_code, error).await {
296 tracing::error!("agent: closing out session {session_id}: {e}");
297 }
298 tracing::info!("agent: session {session_id} {final_status}");
299}
300
301/// Attach a browser to a session's tmux, returning the duplex terminal.
302///
303/// Bumps the attached count for the idle sweep; the caller must call
304/// [`detached`] when the websocket closes.
305pub async fn attach_stream(
306 app: &App,
307 session_id: i64,
308 cols: u16,
309 rows: u16,
310) -> Result<(container::Terminal, bollard::Docker), String> {
311 let handle = app
312 .sessions
313 .get(session_id)
314 .ok_or("session is not running")?;
315
316 let docker = anvil_ci::docker::connect()?;
317 let terminal = container::attach(&docker, &handle.container_id, cols, rows).await?;
318
319 handle
320 .attached
321 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
322 handle.touch();
323 let _ = agent::touch_attach(&app.db, session_id).await;
324
325 Ok((terminal, docker))
326}
327
328/// Note that a browser disconnected.
329pub fn detached(app: &App, session_id: i64) {
330 if let Some(handle) = app.sessions.get(session_id) {
331 handle
332 .attached
333 .fetch_sub(1, std::sync::atomic::Ordering::Relaxed);
334 handle.touch();
335 }
336}