anvilsign in

collin/anvil

1//! Agent sessions: persistence and lifecycle for a tmux-hosted agent CLI
2//! running against a repository. Execution (Docker, tmux, the terminal I/O
3//! pump) lives in `anvil-agent`; this module owns only the rows.
4
5use crate::{
6 error::Result,
7 models::AgentSession,
8};
9
10/// Current Unix time in seconds, so `anvil-agent` measures idle and lifetime
11/// against the same clock the rows are stamped with.
12pub fn now_secs() -> i64 {
13 crate::now()
14}
15
16/// Status values stored in [`AgentSession::status`].
17pub mod status {
18 /// The row exists and the container is being created.
19 pub const STARTING: &str = "starting";
20 /// The container is up and the tmux session is live.
21 pub const RUNNING: &str = "running";
22 /// The agent finished on its own; see `exit_code`.
23 pub const EXITED: &str = "exited";
24 /// The supervisor could not start or keep the session; see `error`.
25 pub const FAILED: &str = "failed";
26 /// A timeout, an operator stop, or a server restart took it.
27 pub const REAPED: &str = "reaped";
28
29 /// Whether `status` is one a session can still leave (i.e. it should have
30 /// a live container behind it).
31 pub fn is_live(status: &str) -> bool {
32 status == STARTING || status == RUNNING
33 }
34}
35
36/// Session kinds stored in [`AgentSession::kind`].
37pub mod kind {
38 /// A human drives the terminal.
39 pub const INTERACTIVE: &str = "interactive";
40 /// Started with a prompt and left to run.
41 pub const AUTONOMOUS: &str = "autonomous";
42}
43
44/// Create a session row in [`status::STARTING`]. The container does not exist
45/// yet — the supervisor fills in `container_id` once it does.
46#[allow(clippy::too_many_arguments)]
47pub async fn create(
48 db: &toasty::Db,
49 repo_id: i64,
50 user_id: i64,
51 kind: &str,
52 base_ref: &str,
53 base_commit: &str,
54 prompt: &str,
55 image: &str,
56 secret_names: &str,
57) -> Result<AgentSession> {
58 let mut conn = db.clone();
59 let now = crate::now();
60 let session = toasty::create!(AgentSession {
61 repo_id: repo_id,
62 user_id: user_id,
63 status: status::STARTING,
64 kind: kind,
65 base_ref: base_ref,
66 base_commit: base_commit,
67 branch: "",
68 prompt: prompt,
69 container_id: "",
70 image: image,
71 created_at: now,
72 started_at: 0,
73 finished_at: 0,
74 last_attach_at: 0,
75 exit_code: 0,
76 error: "",
77 secret_names: secret_names,
78 })
79 .exec(&mut conn)
80 .await?;
81 Ok(session)
82}
83
84/// One session by id.
85pub async fn get(db: &toasty::Db, id: i64) -> Result<Option<AgentSession>> {
86 let mut conn = db.clone();
87 let session = AgentSession::filter(AgentSession::fields().id().eq(id))
88 .first()
89 .exec(&mut conn)
90 .await?;
91 Ok(session)
92}
93
94/// A repository's sessions, newest first, capped at `limit`.
95pub async fn list_by_repo(
96 db: &toasty::Db,
97 repo_id: i64,
98 limit: usize,
99) -> Result<Vec<AgentSession>> {
100 let mut conn = db.clone();
101 let sessions = AgentSession::filter(AgentSession::fields().repo_id().eq(repo_id))
102 .order_by(AgentSession::fields().id().desc())
103 .exec(&mut conn)
104 .await?;
105 Ok(sessions.into_iter().take(limit).collect())
106}
107
108/// Every session that believes it still has a container — used by the sweep
109/// and by the startup reconcile.
110pub async fn live(db: &toasty::Db) -> Result<Vec<AgentSession>> {
111 let mut conn = db.clone();
112 let starting = AgentSession::filter(AgentSession::fields().status().eq(status::STARTING))
113 .exec(&mut conn)
114 .await?;
115 let mut conn = db.clone();
116 let running = AgentSession::filter(AgentSession::fields().status().eq(status::RUNNING))
117 .exec(&mut conn)
118 .await?;
119 let mut all: Vec<_> = starting.into_iter().chain(running).collect();
120 all.sort_by_key(|s| s.id);
121 Ok(all)
122}
123
124/// How many sessions are currently live, for the concurrency cap.
125pub async fn live_count(db: &toasty::Db) -> Result<usize> {
126 Ok(live(db).await?.len())
127}
128
129/// Record the container backing a session and move it to
130/// [`status::RUNNING`].
131pub async fn mark_running(db: &toasty::Db, id: i64, container_id: &str) -> Result<()> {
132 let Some(mut session) = get(db, id).await? else {
133 return Ok(());
134 };
135 let now = crate::now();
136 let mut conn = db.clone();
137 session
138 .update()
139 .status(status::RUNNING)
140 .container_id(container_id)
141 .started_at(now)
142 .last_attach_at(now)
143 .exec(&mut conn)
144 .await?;
145 Ok(())
146}
147
148/// Note that a viewer attached, which is what holds off the idle sweep.
149pub async fn touch_attach(db: &toasty::Db, id: i64) -> Result<()> {
150 let Some(mut session) = get(db, id).await? else {
151 return Ok(());
152 };
153 let mut conn = db.clone();
154 session
155 .update()
156 .last_attach_at(crate::now())
157 .exec(&mut conn)
158 .await?;
159 Ok(())
160}
161
162/// Record the branch the agent's work is on, once there is one.
163pub async fn set_branch(db: &toasty::Db, id: i64, branch: &str) -> Result<()> {
164 let Some(mut session) = get(db, id).await? else {
165 return Ok(());
166 };
167 let mut conn = db.clone();
168 session.update().branch(branch).exec(&mut conn).await?;
169 Ok(())
170}
171
172/// Close a session out with a terminal status. `error` is empty unless the
173/// status is [`status::FAILED`].
174pub async fn finish(
175 db: &toasty::Db,
176 id: i64,
177 status_value: &str,
178 exit_code: i64,
179 error: &str,
180) -> Result<()> {
181 let Some(mut session) = get(db, id).await? else {
182 return Ok(());
183 };
184 let mut conn = db.clone();
185 session
186 .update()
187 .status(status_value)
188 .exit_code(exit_code)
189 .error(error)
190 .finished_at(crate::now())
191 .exec(&mut conn)
192 .await?;
193 Ok(())
194}
195
196// --- the live-session registry ---------------------------------------------
197
198/// Handles to the sessions running right now, keyed by session id.
199///
200/// Mirrors [`secrets::Vault`](crate::secrets::Vault): in-memory only, lives on
201/// [`App`](crate::App), and empty after a restart — which is exactly why
202/// startup reconciles the `agent_sessions` rows against the containers Docker
203/// still has (see `anvil_agent::reconcile`).
204///
205/// Deliberately free of any Docker type so it can live in this crate: the
206/// registry stores a container *id*, and `anvil-agent` owns everything that
207/// knows what to do with one.
208#[derive(Clone, Default)]
209pub struct Registry {
210 inner: std::sync::Arc<std::sync::Mutex<std::collections::HashMap<i64, Handle>>>,
211}
212
213/// One live session's shared state.
214///
215/// Note what is *not* here: a fan-out channel for terminal output. tmux already
216/// multiplexes — every attached browser opens its own tmux client, and tmux
217/// repaints each one — so a broadcast in front of it would duplicate work the
218/// terminal multiplexer exists to do. The durable transcript comes from
219/// `tmux pipe-pane` inside the container instead, which keeps recording
220/// whether or not anyone is attached.
221#[derive(Clone)]
222pub struct Handle {
223 /// Docker container backing the session.
224 pub container_id: String,
225 /// Flipped to `true` to ask the supervisor to wind the session up.
226 pub shutdown: tokio::sync::watch::Sender<bool>,
227 /// Unix seconds of the last attach or output byte, driving the idle sweep.
228 pub last_activity: std::sync::Arc<std::sync::atomic::AtomicI64>,
229 /// How many browsers are attached right now. A session with viewers is
230 /// never idle, however quiet the agent is.
231 pub attached: std::sync::Arc<std::sync::atomic::AtomicUsize>,
232}
233
234impl Handle {
235 /// Note that something happened, holding off the idle sweep.
236 pub fn touch(&self) {
237 self.last_activity
238 .store(crate::now(), std::sync::atomic::Ordering::Relaxed);
239 }
240
241 /// Unix seconds since the last attach or output byte.
242 pub fn idle_secs(&self) -> i64 {
243 crate::now()
244 - self
245 .last_activity
246 .load(std::sync::atomic::Ordering::Relaxed)
247 }
248
249 /// Whether a browser is currently attached.
250 pub fn has_viewers(&self) -> bool {
251 self.attached.load(std::sync::atomic::Ordering::Relaxed) > 0
252 }
253}
254
255impl Registry {
256 /// Register a live session.
257 pub fn insert(&self, session_id: i64, handle: Handle) {
258 if let Ok(mut map) = self.inner.lock() {
259 map.insert(session_id, handle);
260 }
261 }
262
263 /// The handle for a session, if it is live in *this* process.
264 pub fn get(&self, session_id: i64) -> Option<Handle> {
265 self.inner.lock().ok()?.get(&session_id).cloned()
266 }
267
268 /// Drop a session's handle, e.g. once its container is gone.
269 pub fn remove(&self, session_id: i64) -> Option<Handle> {
270 self.inner.lock().ok()?.remove(&session_id)
271 }
272
273 /// Every live session id.
274 pub fn ids(&self) -> Vec<i64> {
275 self.inner
276 .lock()
277 .map(|m| m.keys().copied().collect())
278 .unwrap_or_default()
279 }
280
281 /// How many sessions are live, for the concurrency cap. Counted from the
282 /// registry rather than the database because it is the containers, not the
283 /// rows, that consume the host.
284 pub fn len(&self) -> usize {
285 self.inner.lock().map(|m| m.len()).unwrap_or(0)
286 }
287
288 /// Whether no session is live.
289 pub fn is_empty(&self) -> bool {
290 self.len() == 0
291 }
292}