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