| 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 | |
| 5 | use 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. |
| 12 | pub fn now_secs() -> i64 { |
| 13 | crate::now() |
| 14 | } |
| 15 | |
| 16 | /// Status values stored in [`AgentSession::status`]. |
| 17 | pub 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`]. |
| 37 | pub 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)] |
| 47 | pub 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. |
| 85 | pub 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`. |
| 95 | pub 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 | /// Delete every session row for a repository, returning the deleted ids so the |
| 109 | /// caller can remove their transcripts. Refuses nothing: callers check for live |
| 110 | /// sessions first — see [`repos::delete`](crate::repos::delete), which will not |
| 111 | /// delete a repository whose containers are still up. |
| 112 | pub async fn delete_for_repo(db: &toasty::Db, repo_id: i64) -> Result<Vec<i64>> { |
| 113 | let mut conn = db.clone(); |
| 114 | let sessions = AgentSession::filter(AgentSession::fields().repo_id().eq(repo_id)) |
| 115 | .exec(&mut conn) |
| 116 | .await?; |
| 117 | let mut ids = Vec::with_capacity(sessions.len()); |
| 118 | for session in sessions { |
| 119 | ids.push(session.id); |
| 120 | let mut conn = db.clone(); |
| 121 | session.delete().exec(&mut conn).await?; |
| 122 | } |
| 123 | Ok(ids) |
| 124 | } |
| 125 | |
| 126 | /// A repository's sessions that still believe they have a container. |
| 127 | pub async fn live_for_repo(db: &toasty::Db, repo_id: i64) -> Result<Vec<AgentSession>> { |
| 128 | Ok(live(db) |
| 129 | .await? |
| 130 | .into_iter() |
| 131 | .filter(|s| s.repo_id == repo_id) |
| 132 | .collect()) |
| 133 | } |
| 134 | |
| 135 | /// Every session that believes it still has a container — used by the sweep |
| 136 | /// and by the startup reconcile. |
| 137 | pub async fn live(db: &toasty::Db) -> Result<Vec<AgentSession>> { |
| 138 | let mut conn = db.clone(); |
| 139 | let starting = AgentSession::filter(AgentSession::fields().status().eq(status::STARTING)) |
| 140 | .exec(&mut conn) |
| 141 | .await?; |
| 142 | let mut conn = db.clone(); |
| 143 | let running = AgentSession::filter(AgentSession::fields().status().eq(status::RUNNING)) |
| 144 | .exec(&mut conn) |
| 145 | .await?; |
| 146 | let mut all: Vec<_> = starting.into_iter().chain(running).collect(); |
| 147 | all.sort_by_key(|s| s.id); |
| 148 | Ok(all) |
| 149 | } |
| 150 | |
| 151 | /// How many sessions are currently live, for the concurrency cap. |
| 152 | pub async fn live_count(db: &toasty::Db) -> Result<usize> { |
| 153 | Ok(live(db).await?.len()) |
| 154 | } |
| 155 | |
| 156 | /// Record the container backing a session and move it to |
| 157 | /// [`status::RUNNING`]. |
| 158 | pub async fn mark_running(db: &toasty::Db, id: i64, container_id: &str) -> Result<()> { |
| 159 | let Some(mut session) = get(db, id).await? else { |
| 160 | return Ok(()); |
| 161 | }; |
| 162 | let now = crate::now(); |
| 163 | let mut conn = db.clone(); |
| 164 | session |
| 165 | .update() |
| 166 | .status(status::RUNNING) |
| 167 | .container_id(container_id) |
| 168 | .started_at(now) |
| 169 | .last_attach_at(now) |
| 170 | .exec(&mut conn) |
| 171 | .await?; |
| 172 | Ok(()) |
| 173 | } |
| 174 | |
| 175 | /// Note that a viewer attached, which is what holds off the idle sweep. |
| 176 | pub async fn touch_attach(db: &toasty::Db, id: i64) -> Result<()> { |
| 177 | let Some(mut session) = get(db, id).await? else { |
| 178 | return Ok(()); |
| 179 | }; |
| 180 | let mut conn = db.clone(); |
| 181 | session |
| 182 | .update() |
| 183 | .last_attach_at(crate::now()) |
| 184 | .exec(&mut conn) |
| 185 | .await?; |
| 186 | Ok(()) |
| 187 | } |
| 188 | |
| 189 | /// Record the branch the agent's work is on, once there is one. |
| 190 | pub async fn set_branch(db: &toasty::Db, id: i64, branch: &str) -> Result<()> { |
| 191 | let Some(mut session) = get(db, id).await? else { |
| 192 | return Ok(()); |
| 193 | }; |
| 194 | let mut conn = db.clone(); |
| 195 | session.update().branch(branch).exec(&mut conn).await?; |
| 196 | Ok(()) |
| 197 | } |
| 198 | |
| 199 | /// Close a session out with a terminal status. `error` is empty unless the |
| 200 | /// status is [`status::FAILED`]. |
| 201 | pub async fn finish( |
| 202 | db: &toasty::Db, |
| 203 | id: i64, |
| 204 | status_value: &str, |
| 205 | exit_code: i64, |
| 206 | error: &str, |
| 207 | ) -> Result<()> { |
| 208 | let Some(mut session) = get(db, id).await? else { |
| 209 | return Ok(()); |
| 210 | }; |
| 211 | let mut conn = db.clone(); |
| 212 | session |
| 213 | .update() |
| 214 | .status(status_value) |
| 215 | .exit_code(exit_code) |
| 216 | .error(error) |
| 217 | .finished_at(crate::now()) |
| 218 | .exec(&mut conn) |
| 219 | .await?; |
| 220 | Ok(()) |
| 221 | } |
| 222 | |
| 223 | // --- the live-session registry --------------------------------------------- |
| 224 | |
| 225 | /// Handles to the sessions running right now, keyed by session id. |
| 226 | /// |
| 227 | /// Mirrors [`secrets::Vault`](crate::secrets::Vault): in-memory only, lives on |
| 228 | /// [`App`](crate::App), and empty after a restart — which is exactly why |
| 229 | /// startup reconciles the `agent_sessions` rows against the containers Docker |
| 230 | /// still has (see `anvil_agent::reconcile`). |
| 231 | /// |
| 232 | /// Deliberately free of any Docker type so it can live in this crate: the |
| 233 | /// registry stores a container *id*, and `anvil-agent` owns everything that |
| 234 | /// knows what to do with one. |
| 235 | #[derive(Clone, Default)] |
| 236 | pub struct Registry { |
| 237 | inner: std::sync::Arc<std::sync::Mutex<std::collections::HashMap<i64, Handle>>>, |
| 238 | } |
| 239 | |
| 240 | /// One live session's shared state. |
| 241 | /// |
| 242 | /// Note what is *not* here: a fan-out channel for terminal output. tmux already |
| 243 | /// multiplexes — every attached browser opens its own tmux client, and tmux |
| 244 | /// repaints each one — so a broadcast in front of it would duplicate work the |
| 245 | /// terminal multiplexer exists to do. The durable transcript comes from |
| 246 | /// `tmux pipe-pane` inside the container instead, which keeps recording |
| 247 | /// whether or not anyone is attached. |
| 248 | #[derive(Clone)] |
| 249 | pub struct Handle { |
| 250 | /// Docker container backing the session. |
| 251 | pub container_id: String, |
| 252 | /// Flipped to `true` to ask the supervisor to wind the session up. |
| 253 | pub shutdown: tokio::sync::watch::Sender<bool>, |
| 254 | /// Unix seconds of the last attach or output byte, driving the idle sweep. |
| 255 | pub last_activity: std::sync::Arc<std::sync::atomic::AtomicI64>, |
| 256 | /// How many browsers are attached right now. A session with viewers is |
| 257 | /// never idle, however quiet the agent is. |
| 258 | pub attached: std::sync::Arc<std::sync::atomic::AtomicUsize>, |
| 259 | } |
| 260 | |
| 261 | impl Handle { |
| 262 | /// Note that something happened, holding off the idle sweep. |
| 263 | pub fn touch(&self) { |
| 264 | self.last_activity |
| 265 | .store(crate::now(), std::sync::atomic::Ordering::Relaxed); |
| 266 | } |
| 267 | |
| 268 | /// Unix seconds since the last attach or output byte. |
| 269 | pub fn idle_secs(&self) -> i64 { |
| 270 | crate::now() |
| 271 | - self |
| 272 | .last_activity |
| 273 | .load(std::sync::atomic::Ordering::Relaxed) |
| 274 | } |
| 275 | |
| 276 | /// Whether a browser is currently attached. |
| 277 | pub fn has_viewers(&self) -> bool { |
| 278 | self.attached.load(std::sync::atomic::Ordering::Relaxed) > 0 |
| 279 | } |
| 280 | } |
| 281 | |
| 282 | impl Registry { |
| 283 | /// Register a live session. |
| 284 | pub fn insert(&self, session_id: i64, handle: Handle) { |
| 285 | if let Ok(mut map) = self.inner.lock() { |
| 286 | map.insert(session_id, handle); |
| 287 | } |
| 288 | } |
| 289 | |
| 290 | /// The handle for a session, if it is live in *this* process. |
| 291 | pub fn get(&self, session_id: i64) -> Option<Handle> { |
| 292 | self.inner.lock().ok()?.get(&session_id).cloned() |
| 293 | } |
| 294 | |
| 295 | /// Drop a session's handle, e.g. once its container is gone. |
| 296 | pub fn remove(&self, session_id: i64) -> Option<Handle> { |
| 297 | self.inner.lock().ok()?.remove(&session_id) |
| 298 | } |
| 299 | |
| 300 | /// Every live session id. |
| 301 | pub fn ids(&self) -> Vec<i64> { |
| 302 | self.inner |
| 303 | .lock() |
| 304 | .map(|m| m.keys().copied().collect()) |
| 305 | .unwrap_or_default() |
| 306 | } |
| 307 | |
| 308 | /// How many sessions are live, for the concurrency cap. Counted from the |
| 309 | /// registry rather than the database because it is the containers, not the |
| 310 | /// rows, that consume the host. |
| 311 | pub fn len(&self) -> usize { |
| 312 | self.inner.lock().map(|m| m.len()).unwrap_or(0) |
| 313 | } |
| 314 | |
| 315 | /// Whether no session is live. |
| 316 | pub fn is_empty(&self) -> bool { |
| 317 | self.len() == 0 |
| 318 | } |
| 319 | } |