anvilsign in

collin/anvil

1//! Agent sessions: a long-lived container per session running tmux plus an
2//! agent CLI against one repository, attachable from the browser.
3//!
4//! Deliberately a sibling of `anvil-ci` rather than part of it. The CI runner
5//! is a single task processing one job at a time; a session lives for hours and
6//! must never sit in that queue. Sharing the runner image and the Docker
7//! plumbing (`anvil_docker`) is the extent of the overlap.
8//!
9//! The shape:
10//!
11//! ```text
12//! browser <--ws--> anvil-web <--broadcast--,
13//! |
14//! supervisor task --docker exec (tty)--> tmux client
15//! | |
16//! '--> transcript on disk the agent
17//! ```
18//!
19//! One supervisor per session owns the container and a broadcast channel.
20//! Attached browsers are subscribers, and so is the transcript writer — which
21//! is why output is durable whether or not anyone is watching.
22
23pub mod container;
24mod supervisor;
25
26use anvil_core::{
27 App,
28 agent::{
29 self,
30 status,
31 },
32};
33pub use supervisor::{
34 attach_stream,
35 detached,
36};
37
38/// Why a session could not be started.
39#[derive(Debug)]
40pub enum StartError {
41 /// `agent.enabled` is false.
42 Disabled,
43 /// `agent.max_concurrent` sessions are already running.
44 AtCapacity(usize),
45 /// Anything else, with detail.
46 Failed(String),
47}
48
49impl std::fmt::Display for StartError {
50 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
51 match self {
52 Self::Disabled => write!(f, "agent sessions are disabled (set agent.enabled)"),
53 Self::AtCapacity(n) => {
54 write!(f, "already running {n} sessions (agent.max_concurrent)")
55 }
56 Self::Failed(e) => write!(f, "{e}"),
57 }
58 }
59}
60
61/// Start a session against `repo_id` at `base_ref`, returning its id.
62///
63/// `secret_names` are the user secrets (`anvil_core::secrets::list_for_user`)
64/// this session opted into at start time — see `supervisor::launch` for how
65/// they actually get into the container. Empty is a normal choice, not an
66/// error: most sessions want none.
67///
68/// Creates the row first so a failure part-way through is still visible in the
69/// UI, then hands off to a supervisor task. Returns as soon as the container is
70/// up — the caller redirects to the session page and attaches from there.
71pub async fn start(
72 app: &App,
73 repo_id: i64,
74 user_id: i64,
75 kind: &str,
76 base_ref: &str,
77 prompt: &str,
78 secret_names: &[String],
79) -> Result<i64, StartError> {
80 let cfg = &app.config.agent;
81 if !cfg.enabled {
82 return Err(StartError::Disabled);
83 }
84 let live = app.sessions.len();
85 if cfg.max_concurrent > 0 && live >= cfg.max_concurrent {
86 return Err(StartError::AtCapacity(live));
87 }
88
89 let (repo_path, base_commit) = supervisor::resolve_base(app, repo_id, base_ref)
90 .await
91 .map_err(StartError::Failed)?;
92
93 let session = agent::create(
94 &app.db,
95 repo_id,
96 user_id,
97 kind,
98 base_ref,
99 &base_commit,
100 prompt,
101 &cfg.image,
102 &secret_names.join(","),
103 )
104 .await
105 .map_err(|e| StartError::Failed(e.to_string()))?;
106
107 let id = session.id;
108 match supervisor::launch(app.clone(), session, repo_path).await {
109 Ok(()) => Ok(id),
110 Err(e) => {
111 let _ = agent::finish(&app.db, id, status::FAILED, 0, &e).await;
112 Err(StartError::Failed(e))
113 }
114 }
115}
116
117/// Ask a session to wind up. Idempotent, and safe for a session this process
118/// does not have a handle for (it just marks the row).
119pub async fn stop(app: &App, session_id: i64) {
120 if let Some(handle) = app.sessions.get(session_id) {
121 let _ = handle.shutdown.send(true);
122 return;
123 }
124 // No live handle: either it already ended, or it belongs to a previous
125 // process. Either way the row should not claim to be running.
126 if let Ok(Some(session)) = agent::get(&app.db, session_id).await
127 && status::is_live(&session.status)
128 {
129 let _ = agent::finish(&app.db, session_id, status::REAPED, 0, "").await;
130 }
131}
132
133/// Reconcile `agent_sessions` rows against the containers Docker still has.
134///
135/// The registry is in-memory, so after a restart this process has a handle for
136/// nothing: every labelled container is an orphan and every row claiming to be
137/// live is stale. Both are reaped, which mirrors what `ci::requeue_interrupted`
138/// does for interrupted runs — except a terminal session cannot be resumed, so
139/// it ends rather than re-queues.
140pub async fn reconcile(app: &App) {
141 let docker = match anvil_docker::connect() {
142 Ok(docker) => docker,
143 Err(e) => {
144 tracing::warn!("agent: skipping reconcile, {e}");
145 return;
146 }
147 };
148
149 match container::list_sessions(&docker).await {
150 Ok(found) => {
151 for (container_id, session_id) in found {
152 tracing::info!("agent: reaping orphaned container for session {session_id}");
153 container::remove(&docker, &container_id).await;
154 }
155 }
156 Err(e) => tracing::warn!("agent: listing session containers: {e}"),
157 }
158
159 match agent::live(&app.db).await {
160 Ok(stale) => {
161 for session in stale {
162 let _ = agent::finish(
163 &app.db,
164 session.id,
165 status::REAPED,
166 0,
167 "anvil restarted while this session was running",
168 )
169 .await;
170 }
171 }
172 Err(e) => tracing::error!("agent: reconciling session rows: {e}"),
173 }
174}
175
176/// Periodic job enforcing the idle and wall-clock caps.
177///
178/// A session with an attached viewer is never idle, however quiet the agent is
179/// — otherwise watching a long think would kill it.
180pub struct SweepJob;
181
182#[async_trait::async_trait]
183impl anvil_core::periodic::PeriodicJob for SweepJob {
184 fn name(&self) -> &str {
185 "agent session sweep"
186 }
187
188 async fn run(&self, app: &App) -> anvil_core::Result<()> {
189 let cfg = &app.config.agent;
190 let now = anvil_core::agent::now_secs();
191
192 for session in agent::live(&app.db).await? {
193 let Some(handle) = app.sessions.get(session.id) else {
194 // Live row, no handle: a leftover this process cannot drive.
195 let _ = agent::finish(&app.db, session.id, status::REAPED, 0, "").await;
196 continue;
197 };
198
199 let lifetime = now - session.started_at.max(session.created_at);
200 if cfg.max_lifetime_secs > 0 && lifetime as u64 > cfg.max_lifetime_secs {
201 tracing::info!("agent: session {} hit the lifetime cap", session.id);
202 let _ = handle.shutdown.send(true);
203 continue;
204 }
205
206 if cfg.idle_timeout_secs > 0
207 && !handle.has_viewers()
208 && handle.idle_secs() as u64 > cfg.idle_timeout_secs
209 {
210 tracing::info!("agent: session {} idle, reaping", session.id);
211 let _ = handle.shutdown.send(true);
212 }
213 }
214 Ok(())
215 }
216}