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_ci::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/// Creates the row first so a failure part-way through is still visible in the
64/// UI, then hands off to a supervisor task. Returns as soon as the container is
65/// up — the caller redirects to the session page and attaches from there.
66pub async fn start(
67 app: &App,
68 repo_id: i64,
69 user_id: i64,
70 kind: &str,
71 base_ref: &str,
72 prompt: &str,
73) -> Result<i64, StartError> {
74 let cfg = &app.config.agent;
75 if !cfg.enabled {
76 return Err(StartError::Disabled);
77 }
78 let live = app.sessions.len();
79 if cfg.max_concurrent > 0 && live >= cfg.max_concurrent {
80 return Err(StartError::AtCapacity(live));
81 }
82
83 let (repo_path, base_commit) = supervisor::resolve_base(app, repo_id, base_ref)
84 .await
85 .map_err(StartError::Failed)?;
86
87 let session = agent::create(
88 &app.db,
89 repo_id,
90 user_id,
91 kind,
92 base_ref,
93 &base_commit,
94 prompt,
95 &cfg.image,
96 )
97 .await
98 .map_err(|e| StartError::Failed(e.to_string()))?;
99
100 let id = session.id;
101 match supervisor::launch(app.clone(), session, repo_path).await {
102 Ok(()) => Ok(id),
103 Err(e) => {
104 let _ = agent::finish(&app.db, id, status::FAILED, 0, &e).await;
105 Err(StartError::Failed(e))
106 }
107 }
108}
109
110/// Ask a session to wind up. Idempotent, and safe for a session this process
111/// does not have a handle for (it just marks the row).
112pub async fn stop(app: &App, session_id: i64) {
113 if let Some(handle) = app.sessions.get(session_id) {
114 let _ = handle.shutdown.send(true);
115 return;
116 }
117 // No live handle: either it already ended, or it belongs to a previous
118 // process. Either way the row should not claim to be running.
119 if let Ok(Some(session)) = agent::get(&app.db, session_id).await
120 && status::is_live(&session.status)
121 {
122 let _ = agent::finish(&app.db, session_id, status::REAPED, 0, "").await;
123 }
124}
125
126/// Reconcile `agent_sessions` rows against the containers Docker still has.
127///
128/// The registry is in-memory, so after a restart this process has a handle for
129/// nothing: every labelled container is an orphan and every row claiming to be
130/// live is stale. Both are reaped, which mirrors what `ci::requeue_interrupted`
131/// does for interrupted runs — except a terminal session cannot be resumed, so
132/// it ends rather than re-queues.
133pub async fn reconcile(app: &App) {
134 let docker = match anvil_ci::docker::connect() {
135 Ok(docker) => docker,
136 Err(e) => {
137 tracing::warn!("agent: skipping reconcile, {e}");
138 return;
139 }
140 };
141
142 match container::list_sessions(&docker).await {
143 Ok(found) => {
144 for (container_id, session_id) in found {
145 tracing::info!("agent: reaping orphaned container for session {session_id}");
146 container::remove(&docker, &container_id).await;
147 }
148 }
149 Err(e) => tracing::warn!("agent: listing session containers: {e}"),
150 }
151
152 match agent::live(&app.db).await {
153 Ok(stale) => {
154 for session in stale {
155 let _ = agent::finish(
156 &app.db,
157 session.id,
158 status::REAPED,
159 0,
160 "anvil restarted while this session was running",
161 )
162 .await;
163 }
164 }
165 Err(e) => tracing::error!("agent: reconciling session rows: {e}"),
166 }
167}
168
169/// Periodic job enforcing the idle and wall-clock caps.
170///
171/// A session with an attached viewer is never idle, however quiet the agent is
172/// — otherwise watching a long think would kill it.
173pub struct SweepJob;
174
175#[async_trait::async_trait]
176impl anvil_core::periodic::PeriodicJob for SweepJob {
177 fn name(&self) -> &str {
178 "agent session sweep"
179 }
180
181 async fn run(&self, app: &App) -> anvil_core::Result<()> {
182 let cfg = &app.config.agent;
183 let now = anvil_core::agent::now_secs();
184
185 for session in agent::live(&app.db).await? {
186 let Some(handle) = app.sessions.get(session.id) else {
187 // Live row, no handle: a leftover this process cannot drive.
188 let _ = agent::finish(&app.db, session.id, status::REAPED, 0, "").await;
189 continue;
190 };
191
192 let lifetime = now - session.started_at.max(session.created_at);
193 if cfg.max_lifetime_secs > 0 && lifetime as u64 > cfg.max_lifetime_secs {
194 tracing::info!("agent: session {} hit the lifetime cap", session.id);
195 let _ = handle.shutdown.send(true);
196 continue;
197 }
198
199 if cfg.idle_timeout_secs > 0
200 && !handle.has_viewers()
201 && handle.idle_secs() as u64 > cfg.idle_timeout_secs
202 {
203 tracing::info!("agent: session {} idle, reaping", session.id);
204 let _ = handle.shutdown.send(true);
205 }
206 }
207 Ok(())
208 }
209}