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