flagon-io/g1t

public

Git for AI scale: a forge for thousands of agents working on the same code at once.

g1t/services/work/src/runs.rs

651 lines26,225 bytesCodeBlame

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

Agents and memory, checks and conflicts, profiles, slug renames, custom domains1//! Agent runs: every sandbox g1t starts for an agent, and for checks and
2//! the merge queue, as a record people can watch, stop and look back on.
3//!
4//! The runner service opens a run as it starts a sandbox (`open_run`) and
5//! the sandbox reports its steps through the API with the run's one-time
6//! token (`report_run`), as checks do. The runner closes it when the
7//! sandbox stops. Members stop a run with `stop_run`; the runner then
8//! destroys its sandbox.
9//!
10//! Sessions are the read side of what agents record on a pull request:
11//! `list_sessions` and `get_session`.
12
13use g1t_contracts::agents::*;
14use g1t_contracts::repos::{Repo, RepoPath};
15use g1t_contracts::time::rfc3339;
16use g1t_contracts::work::{PullStatus, SessionEntry, StallArgs};
17use g1t_contracts::{FailureCode, Membership, Outcome, User, Viewer, new_id};
18use g1t_kit::now_ms;
19use serde::Deserialize;
20use worker::Result;
21use worker::wasm_bindgen::JsValue;
22
23use crate::checks::{hash, new_token};
24use crate::rows::SessionRow;
25use crate::{Work, optional};
26
27/// The most steps a run keeps; older ones fall off the start.
28const MAX_STEPS: u32 = 200;
29/// The most steps taken from one report.
30const MAX_STEPS_PER_REPORT: usize = 20;
31/// The longest a step is kept.
32const MAX_STEP_CHARS: usize = 240;
33/// A run that has said nothing for this long is taken to have died.
34const SILENT_HOURS: u64 = 3;
35const DEFAULT_LIST: u32 = 50;
36const MAX_LIST: u32 = 200;
37/// The most session entries one view returns.
38const SESSION_ENTRIES: u32 = 2000;
39
40/// Every column but the steps, which only a run's own page needs.
41const COLUMNS: &str = "id, workspace, repo_id, repo, number, pull_id, kind, agent, model, status,
42 step, step_count, cost_usd, turns, sandbox, token_hash, started_by, error, created_at,
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API43 started_at, finished_at, updated_at, budget_usd, time_cap_minutes, halted,
Agents and memory, checks and conflicts, profiles, slug renames, custom domains44 COALESCE((SELECT title FROM pulls WHERE pulls.id = agent_runs.pull_id), agent_runs.title) AS title";
45
46#[derive(Deserialize)]
47pub(crate) struct RunRow {
48 id: String,
49 #[allow(dead_code)]
50 workspace: String,
51 #[allow(dead_code)]
52 repo_id: String,
53 repo: String,
54 number: Option<u32>,
55 pub(crate) pull_id: Option<String>,
56 kind: String,
57 agent: String,
58 model: Option<String>,
59 status: String,
60 step: Option<String>,
61 #[serde(default)]
62 steps: Option<String>,
63 step_count: u32,
64 cost_usd: Option<f64>,
65 turns: Option<u32>,
66 sandbox: String,
67 token_hash: String,
68 started_by: Option<String>,
69 error: Option<String>,
70 created_at: String,
71 started_at: Option<String>,
72 finished_at: Option<String>,
73 updated_at: String,
74 title: Option<String>,
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API75 #[serde(default)]
76 budget_usd: Option<f64>,
77 #[serde(default)]
78 time_cap_minutes: Option<u32>,
79 #[serde(default)]
80 halted: Option<String>,
Agents and memory, checks and conflicts, profiles, slug renames, custom domains81}
82
83impl RunRow {
84 fn status(&self) -> RunStatus {
85 RunStatus::parse(&self.status).unwrap_or(RunStatus::Failed)
86 }
87
88 /// The run as `member` (or not) may see it: model and cost are the
89 /// workspace's business.
90 fn into_run(self, member: bool) -> AgentRun {
91 let (namespace, name) = self.repo.split_once('/').unwrap_or((&self.repo, ""));
92 AgentRun {
93 repo: RepoPath {
94 namespace: namespace.to_owned(),
95 name: name.to_owned(),
96 },
97 status: self.status(),
98 kind: RunKind::parse(&self.kind).unwrap_or(RunKind::Implement),
99 steps: self
100 .steps
101 .as_deref()
102 .and_then(|steps| serde_json::from_str(steps).ok())
103 .unwrap_or_default(),
104 id: self.id,
105 number: self.number.filter(|number| *number > 0),
106 title: self.title,
107 agent: self.agent,
108 model: self.model.filter(|_| member),
109 step: self.step,
110 step_count: self.step_count,
111 started_by: self.started_by,
112 error: self.error,
113 cost_usd: self.cost_usd.filter(|_| member),
114 turns: self.turns,
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API115 budget_usd: self.budget_usd.filter(|_| member),
116 time_cap_minutes: self.time_cap_minutes,
117 halted: self.halted,
Agents and memory, checks and conflicts, profiles, slug renames, custom domains118 created_at: self.created_at,
119 started_at: self.started_at,
120 finished_at: self.finished_at,
121 updated_at: self.updated_at,
122 }
123 }
124}
125
126/// One line, short enough to read at a glance.
127pub(crate) fn one_line(text: &str, limit: usize) -> String {
128 let line = text.split_whitespace().collect::<Vec<_>>().join(" ");
129 if line.chars().count() <= limit {
130 return line;
131 }
132 let mut short: String = line.chars().take(limit.saturating_sub(1)).collect();
133 short.push('…');
134 short
135}
136
137/// A principal that can read any repository of `namespace`, for lookups a
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look138/// trusted service asks for. What it is given reads and nothing more.
Agents and memory, checks and conflicts, profiles, slug renames, custom domains139pub(crate) fn member_of(actor: &User, namespace: &str) -> Viewer {
140 let mut viewer = actor.clone();
141 let slug = namespace.to_lowercase();
142 if !viewer.is_member(&slug) {
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look143 viewer.workspaces.push(Membership {
144 base_permission: Some(g1t_contracts::access::BasePermission::Read),
145 ..Membership::member(slug)
146 });
Agents and memory, checks and conflicts, profiles, slug renames, custom domains147 }
148 Some(viewer)
149}
150
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look151/// Whether the viewer belongs to the workspace: what a run spent, its model
152/// and its budget are the workspace's business, not every reader's (an
153/// outside collaborator sees the work, not the bill).
Agents and memory, checks and conflicts, profiles, slug renames, custom domains154fn is_member(viewer: &Viewer, namespace: &str) -> bool {
155 viewer
156 .as_ref()
157 .is_some_and(|viewer| viewer.is_member(&namespace.to_lowercase()))
158}
159
160#[derive(Deserialize)]
161struct SessionListRow {
162 pull_id: String,
163 number: u32,
164 title: String,
165 status: PullStatus,
166 agent: String,
167 entries: u32,
168 tools: u32,
169 started_at: String,
170 last_at: String,
171 prompt: Option<String>,
172}
173
174#[derive(Deserialize)]
175struct RunSumRow {
176 pull_id: String,
177 kind: String,
178 status: String,
179 cost_usd: Option<f64>,
180}
181
182impl Work {
183 /// The repository at `path`, if `viewer` may see it, and whether they
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look184 /// are a member of its workspace (which shows what runs cost).
Agents and memory, checks and conflicts, profiles, slug renames, custom domains185 async fn visible_repo(&self, path: &RepoPath, viewer: &Viewer) -> Result<Outcome<(Repo, bool)>> {
186 Ok(match self.repo(path, viewer).await? {
187 Outcome::Ok(repo) => {
188 let member = is_member(viewer, &repo.namespace);
189 Outcome::Ok((repo, member))
190 }
191 Outcome::Fail(failure) => Outcome::Fail(failure),
192 })
193 }
194
195 async fn run_row(&self, id: &str) -> Result<Option<RunRow>> {
196 self.db
197 .prepare(format!("SELECT {COLUMNS}, steps FROM agent_runs WHERE id = ?"))
198 .bind(&[id.into()])?
199 .first::<RunRow>(None)
200 .await
201 }
202
203 /// The run on pull request `number` that is at work now, if one is.
204 pub(crate) async fn active_run_on(&self, repo_id: &str, number: u32) -> Result<Option<String>> {
205 self.db
206 .prepare(
207 "SELECT id AS value FROM agent_runs
208 WHERE repo_id = ? AND number = ? AND status IN ('queued', 'running')
209 ORDER BY created_at DESC LIMIT 1",
210 )
211 .bind(&[repo_id.into(), number.into()])?
212 .first::<String>(Some("value"))
213 .await
214 }
215
216 pub(crate) async fn open_run(&self, a: OpenRunArgs) -> Result<Outcome<AgentRunTicket>> {
217 // The runner is trusted: it names the repository it is starting a
218 // sandbox in, whoever the sandbox acts as.
219 let (repo_id, namespace) = match &a.pull_id {
220 Some(pull_id) => match self.pull_by_id(pull_id).await? {
221 Some(pull) => (pull.repo_id, a.repo.namespace.to_lowercase()),
222 None => return Ok(Outcome::fail(FailureCode::NotFound, "Pull request not found.")),
223 },
224 None => match self.repo(&a.repo, &member_of(&a.actor, &a.repo.namespace)).await? {
225 Outcome::Ok(repo) => (repo.id, repo.namespace.to_lowercase()),
226 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
227 },
228 };
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look229 // Nothing new starts on an archived or deleted repository.
230 if !self.repo_active(&repo_id).await? {
231 return Ok(Outcome::fail(
232 FailureCode::Forbidden,
233 format!("{}/{} is archived or deleted, so nothing new starts on it.", a.repo.namespace, a.repo.name),
234 ));
235 }
Agents and memory, checks and conflicts, profiles, slug renames, custom domains236 let now = now_ms();
237 let id = new_id("arn", now);
238 let token = new_token();
239 let timestamp = rfc3339(now);
240 let agent = a
241 .agent
242 .clone()
243 .unwrap_or_else(|| if a.kind.is_agent() { "g1t-agent" } else { "g1t" }.to_owned());
244 self.db
245 .prepare(
246 "INSERT INTO agent_runs
247 (id, workspace, repo_id, repo, number, pull_id, title, kind, agent, model, status,
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API248 sandbox, token_hash, started_by, created_at, updated_at, budget_usd, time_cap_minutes)
249 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 'queued', ?, ?, ?, ?, ?, ?, ?)",
Agents and memory, checks and conflicts, profiles, slug renames, custom domains250 )
251 .bind(&[
252 id.as_str().into(),
253 namespace.into(),
254 repo_id.into(),
255 format!("{}/{}", a.repo.namespace, a.repo.name).into(),
256 a.number.filter(|n| *n > 0).map_or(JsValue::NULL, JsValue::from),
257 optional(&a.pull_id),
258 optional(&a.title.map(|title| one_line(&title, 200))),
259 a.kind.as_str().into(),
260 agent.into(),
261 optional(&a.model),
262 a.sandbox.into(),
263 hash(&token).into(),
264 optional(&a.started_by),
265 timestamp.as_str().into(),
266 timestamp.as_str().into(),
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API267 a.budget_usd.filter(|usd| usd.is_finite() && *usd > 0.0).map_or(JsValue::NULL, JsValue::from),
268 a.time_cap_minutes.map_or(JsValue::NULL, JsValue::from),
Agents and memory, checks and conflicts, profiles, slug renames, custom domains269 ])?
270 .run()
271 .await?;
272 Ok(Outcome::Ok(AgentRunTicket { run_id: id, token }))
273 }
274
275 /// Adds `text` to a run's steps, keeping the latest `MAX_STEPS`.
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API276 pub(crate) fn add_step(&self, run_id: &str, at: &str, text: &str) -> Result<worker::D1PreparedStatement> {
Agents and memory, checks and conflicts, profiles, slug renames, custom domains277 self.db
278 .prepare(
279 "UPDATE agent_runs SET
280 steps = CASE WHEN json_array_length(steps) >= ?4
281 THEN json_insert(json_remove(steps, '$[0]'), '$[#]', json_object('at', ?1, 'text', ?2))
282 ELSE json_insert(steps, '$[#]', json_object('at', ?1, 'text', ?2)) END,
283 step_count = step_count + 1
284 WHERE id = ?3",
285 )
286 .bind(&[at.into(), text.into(), run_id.into(), MAX_STEPS.into()])
287 }
288
289 pub(crate) async fn report_run(&self, a: ReportRunArgs) -> Result<Outcome<RunStatus>> {
290 let Some(run) = self
291 .run_row(&a.run_id)
292 .await?
293 .filter(|run| run.token_hash == hash(&a.token))
294 else {
295 return Ok(Outcome::fail(FailureCode::NotFound, "Run not found."));
296 };
297 // Finished, or stopped by a person: nothing more is taken, and the
298 // sandbox learns why.
299 if run.finished_at.is_some() {
300 return Ok(Outcome::Ok(run.status()));
301 }
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API302 // It reached a cap of its guardrails (guardrails.rs).
303 if let Some(halt) = a.halt {
304 let pull_id = run.pull_id.as_deref();
305 return self
306 .halt_run(&run.id, &run.repo_id, pull_id, run.number, &run.kind, halt, a.error, a.cost_usd)
307 .await;
308 }
Agents and memory, checks and conflicts, profiles, slug renames, custom domains309 let now = rfc3339(now_ms());
310 let steps: Vec<String> = a
311 .steps
312 .iter()
313 .map(|step| one_line(step, MAX_STEP_CHARS))
314 .filter(|step| !step.is_empty())
315 .collect();
316 // A burst keeps its end: that is where the run is now.
317 let skip = steps.len().saturating_sub(MAX_STEPS_PER_REPORT);
318 let mut statements = Vec::new();
319 for step in &steps[skip..] {
320 statements.push(self.add_step(&run.id, &now, step)?);
321 }
322 let current = a
323 .step
324 .as_deref()
325 .map(|step| one_line(step, MAX_STEP_CHARS))
326 .filter(|step| !step.is_empty())
327 .or_else(|| steps.last().cloned());
328 let outcome = a
329 .outcome
330 .filter(|outcome| matches!(outcome, RunStatus::Succeeded | RunStatus::Failed));
331 let cost = a.cost_usd.filter(|cost| cost.is_finite() && *cost >= 0.0);
332 statements.push(
333 self.db
334 .prepare(
335 "UPDATE agent_runs SET
336 status = COALESCE(?1, CASE WHEN status = 'queued' THEN 'running' ELSE status END),
337 started_at = COALESCE(started_at, ?2),
338 step = COALESCE(?3, step),
339 cost_usd = COALESCE(?4, cost_usd),
340 turns = COALESCE(?5, turns),
341 error = COALESCE(?6, error),
342 finished_at = CASE WHEN ?1 IS NULL THEN finished_at ELSE ?2 END,
343 updated_at = ?2
344 WHERE id = ?7 AND finished_at IS NULL",
345 )
346 .bind(&[
347 outcome.map_or(JsValue::NULL, |outcome| outcome.as_str().into()),
348 now.as_str().into(),
349 optional(&current),
350 cost.map_or(JsValue::NULL, JsValue::from),
351 a.turns.map_or(JsValue::NULL, JsValue::from),
352 optional(&a.error.map(|error| one_line(&error, 1000))),
353 run.id.as_str().into(),
354 ])?,
355 );
356 self.db.batch(statements).await?;
357 Ok(Outcome::Ok(outcome.unwrap_or(RunStatus::Running)))
358 }
359
360 pub(crate) async fn stop_run(&self, a: StopRunArgs) -> Result<Outcome<StoppedRun>> {
361 let viewer = Some(a.actor.clone());
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look362 let (repo, _) = match self.visible_repo(&a.repo, &viewer).await? {
Agents and memory, checks and conflicts, profiles, slug renames, custom domains363 Outcome::Ok(found) => found,
364 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
365 };
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look366 if !a.actor.verified {
367 return Ok(Outcome::fail(FailureCode::Forbidden, crate::UNVERIFIED));
368 }
369 if let Outcome::Fail(failure) =
370 crate::allowed(Some(&a.actor), &repo, g1t_contracts::access::Capability::Run)
371 {
372 return Ok(Outcome::Fail(failure));
Agents and memory, checks and conflicts, profiles, slug renames, custom domains373 }
374 let Some(run) = self
375 .run_row(&a.id)
376 .await?
377 .filter(|run| run.repo_id == repo.id)
378 else {
379 return Ok(Outcome::fail(FailureCode::NotFound, "Run not found."));
380 };
381 if !run.status().is_active() {
382 return Ok(Outcome::fail(
383 FailureCode::Conflict,
384 format!("This run has already {}.", match run.status() {
385 RunStatus::Stopped => "been stopped",
386 RunStatus::Succeeded => "finished",
387 _ => "ended",
388 }),
389 ));
390 }
391 let now = rfc3339(now_ms());
392 let said = format!("Stopped by {}.", a.actor.username);
393 let claimed = self
394 .db
395 .prepare(
396 "UPDATE agent_runs SET status = 'stopped', step = ?, finished_at = ?, updated_at = ?
397 WHERE id = ? AND status IN ('queued', 'running') RETURNING id AS value",
398 )
399 .bind(&[
400 said.as_str().into(),
401 now.as_str().into(),
402 now.as_str().into(),
403 run.id.as_str().into(),
404 ])?
405 .first::<String>(Some("value"))
406 .await?;
407 if claimed.is_none() {
408 return Ok(Outcome::fail(FailureCode::Conflict, "This run has already ended."));
409 }
410 self.add_step(&run.id, &now, &said)?.run().await?;
411 // g1t stops seeing the pull request through, so it does not start
412 // the same work again; a person decides what happens next.
413 if let (Some(pull_id), Some(number)) = (&run.pull_id, run.number) {
414 self.stall(StallArgs {
415 pull_id: pull_id.clone(),
416 reason: format!(
417 "{} stopped the agent's {} run. Ask for a review, a revision or a catch-up to start again.",
418 a.actor.username, run.kind
419 ),
420 })
421 .await?;
422 self.note(
423 &repo.id,
424 number,
425 (a.actor.id.as_str(), a.actor.username.as_str()),
426 &format!("stopped {}'s {} run", run.agent, run.kind),
427 )
428 .await?;
429 }
430 let sandbox = run.sandbox.clone();
431 let Some(row) = self.run_row(&run.id).await? else {
432 return Ok(Outcome::fail(FailureCode::NotFound, "Run not found."));
433 };
434 Ok(Outcome::Ok(StoppedRun {
435 run: row.into_run(true),
436 sandbox,
437 }))
438 }
439
440 /// Runs whose sandbox died without anyone noticing, marked failed.
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look441 pub(crate) async fn sweep_silent(&self) -> Result<()> {
Agents and memory, checks and conflicts, profiles, slug renames, custom domains442 let now = now_ms();
443 let cutoff = rfc3339(now.saturating_sub(SILENT_HOURS * 3_600_000));
444 self.db
445 .prepare(
446 "UPDATE agent_runs
447 SET status = 'failed', error = COALESCE(error, 'It stopped reporting.'),
448 finished_at = ?, updated_at = ?
449 WHERE status IN ('queued', 'running') AND updated_at < ?",
450 )
451 .bind(&[rfc3339(now).into(), rfc3339(now).into(), cutoff.into()])?
452 .run()
453 .await?;
454 Ok(())
455 }
456
457 pub(crate) async fn list_runs(&self, a: ListRunsArgs) -> Result<Outcome<Vec<AgentRun>>> {
458 let (column, key, member) = if let Some(workspace) = &a.workspace {
459 let slug = workspace.to_lowercase();
460 if !is_member(&a.viewer, &slug) {
461 return Ok(Outcome::fail(
462 FailureCode::Forbidden,
463 "Only members of the workspace can see its fleet.",
464 ));
465 }
466 ("workspace", slug, true)
467 } else if let Some(path) = &a.repo {
468 match self.visible_repo(path, &a.viewer).await? {
469 Outcome::Ok((repo, member)) => ("repo_id", repo.id, member),
470 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
471 }
472 } else {
473 return Ok(Outcome::fail(FailureCode::Invalid, "Name a repository or a workspace."));
474 };
475 self.sweep_silent().await?;
476 let limit = a.limit.unwrap_or(DEFAULT_LIST).clamp(1, MAX_LIST);
477 let rows = self
478 .db
479 .prepare(format!(
480 "SELECT {COLUMNS} FROM agent_runs
481 WHERE {column} = ?1
482 AND (?2 = 0 OR status IN ('queued', 'running'))
483 AND (?3 IS NULL OR kind = ?3)
484 AND (?4 IS NULL OR status = ?4)
485 AND (?5 IS NULL OR number = ?5)
486 ORDER BY created_at DESC LIMIT ?6"
487 ))
488 .bind(&[
489 key.into(),
490 u32::from(a.active).into(),
491 a.kind.map_or(JsValue::NULL, |kind| kind.as_str().into()),
492 a.status.map_or(JsValue::NULL, |status| status.as_str().into()),
493 a.number.map_or(JsValue::NULL, JsValue::from),
494 limit.into(),
495 ])?
496 .all()
497 .await?
498 .results::<RunRow>()?;
499 Ok(Outcome::Ok(rows.into_iter().map(|row| row.into_run(member)).collect()))
500 }
501
502 pub(crate) async fn get_run(&self, a: GetRunArgs) -> Result<Outcome<AgentRun>> {
503 let (repo, member) = match self.visible_repo(&a.repo, &a.viewer).await? {
504 Outcome::Ok(found) => found,
505 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
506 };
507 Ok(match self.run_row(&a.id).await?.filter(|run| run.repo_id == repo.id) {
508 Some(row) => Outcome::Ok(row.into_run(member)),
509 None => Outcome::fail(FailureCode::NotFound, "Run not found."),
510 })
511 }
512
513 // --- Sessions ----------------------------------------------------------
514
515 pub(crate) async fn list_sessions(&self, a: ListSessionsArgs) -> Result<Outcome<Vec<SessionSummary>>> {
516 let Some(path) = &a.repo else {
517 return Ok(Outcome::fail(FailureCode::Invalid, "Name a repository."));
518 };
519 let (repo, member) = match self.visible_repo(path, &a.viewer).await? {
520 Outcome::Ok(found) => found,
521 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
522 };
523 let rows = self
524 .db
525 .prepare(
526 "SELECT p.id AS pull_id, p.number, p.title, p.status, p.agent,
527 count(*) AS entries,
528 sum(CASE WHEN e.kind = 'tool_call' THEN 1 ELSE 0 END) AS tools,
529 min(e.at) AS started_at, max(e.at) AS last_at,
530 (SELECT substr(f.text, 1, 300) FROM session_entries f
531 WHERE f.pull_id = p.id AND f.kind = 'prompt' ORDER BY f.seq LIMIT 1) AS prompt
532 FROM pulls p JOIN session_entries e ON e.pull_id = p.id
533 WHERE p.repo_id = ?1
534 AND (?2 IS NULL OR p.status = ?2)
535 AND (?3 IS NULL OR p.number = ?3)
536 AND (?4 IS NULL OR EXISTS
537 (SELECT 1 FROM agent_runs r WHERE r.pull_id = p.id AND r.kind = ?4))
538 GROUP BY p.id ORDER BY last_at DESC LIMIT 100",
539 )
540 .bind(&[
541 repo.id.as_str().into(),
542 a.outcome.map_or(JsValue::NULL, |status| status.as_str().into()),
543 a.number.map_or(JsValue::NULL, JsValue::from),
544 a.kind.map_or(JsValue::NULL, |kind| kind.as_str().into()),
545 ])?
546 .all()
547 .await?
548 .results::<SessionListRow>()?;
549 let ids: Vec<&str> = rows.iter().map(|row| row.pull_id.as_str()).collect();
550 let runs = if ids.is_empty() {
551 Vec::new()
552 } else {
553 self.db
554 .prepare(
555 "SELECT pull_id, kind, status, cost_usd FROM agent_runs
556 WHERE repo_id = ? AND pull_id IN (SELECT value FROM json_each(?))",
557 )
558 .bind(&[repo.id.as_str().into(), serde_json::to_string(&ids)?.into()])?
559 .all()
560 .await?
561 .results::<RunSumRow>()?
562 };
563 Ok(Outcome::Ok(
564 rows.into_iter()
565 .map(|row| {
566 let mine: Vec<&RunSumRow> = runs.iter().filter(|run| run.pull_id == row.pull_id).collect();
567 let mut kinds: Vec<RunKind> = Vec::new();
568 for run in &mine {
569 if let Some(kind) = RunKind::parse(&run.kind)
570 && !kinds.contains(&kind)
571 {
572 kinds.push(kind);
573 }
574 }
575 let spent: f64 = mine.iter().filter_map(|run| run.cost_usd).sum();
576 SessionSummary {
577 number: row.number,
578 title: row.title,
579 status: row.status,
580 agent: row.agent,
581 entries: row.entries,
582 tools: row.tools,
583 prompt: row.prompt.map(|prompt| one_line(&prompt, 200)),
584 runs: mine.len() as u32,
585 cost_usd: (member && mine.iter().any(|run| run.cost_usd.is_some())).then_some(spent),
586 active: mine
587 .iter()
588 .any(|run| RunStatus::parse(&run.status).is_some_and(RunStatus::is_active)),
589 kinds,
590 started_at: row.started_at,
591 last_at: row.last_at,
592 }
593 })
594 .collect(),
595 ))
596 }
597
598 pub(crate) async fn get_session(&self, a: GetSessionArgs) -> Result<Outcome<SessionView>> {
599 let (repo, pull) = match self.pull_at(&a.repo, a.number, &a.viewer).await? {
600 Outcome::Ok(found) => found,
601 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
602 };
603 let member = is_member(&a.viewer, &repo.namespace);
604 let entries: Vec<SessionEntry> = self
605 .db
606 .prepare(
607 "SELECT seq, kind, text, tool, \"commit\", at FROM session_entries
608 WHERE pull_id = ? ORDER BY seq LIMIT ?",
609 )
610 .bind(&[pull.id.as_str().into(), SESSION_ENTRIES.into()])?
611 .all()
612 .await?
613 .results::<SessionRow>()?
614 .into_iter()
615 .map(SessionEntry::from)
616 .collect();
617 let runs: Vec<AgentRun> = self
618 .db
619 .prepare(format!(
620 "SELECT {COLUMNS}, steps FROM agent_runs WHERE pull_id = ? ORDER BY created_at DESC LIMIT 50"
621 ))
622 .bind(&[pull.id.as_str().into()])?
623 .all()
624 .await?
625 .results::<RunRow>()?
626 .into_iter()
627 .map(|row| row.into_run(member))
628 .collect();
629 let cost_usd = (member && runs.iter().any(|run| run.cost_usd.is_some()))
630 .then(|| runs.iter().filter_map(|run| run.cost_usd).sum());
631 Ok(Outcome::Ok(SessionView {
632 pull,
633 entries,
634 runs,
635 cost_usd,
636 }))
637 }
638}
639
640#[cfg(test)]
641mod tests {
642 use super::*;
643
644 #[test]
645 fn steps_are_one_short_line() {
646 assert_eq!(one_line(" Read\n src/lib.rs ", 40), "Read src/lib.rs");
647 let long = one_line(&"x".repeat(300), 10);
648 assert_eq!(long.chars().count(), 10);
649 assert!(long.ends_with('…'));
650 }
651}