flagon-io/g1t

public

Where people and agents ship software together. The open-source git platform for the whole job: issues, agents, checks and deploys to the edge.

g1t/services/work/src/runs.rs

651 lines26,225 bytesCodeBlame
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/// 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,
43 started_at, finished_at, updated_at, budget_usd, time_cap_minutes, halted,
44 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>,
75 #[serde(default)]
76 budget_usd: Option<f64>,
77 #[serde(default)]
78 time_cap_minutes: Option<u32>,
79 #[serde(default)]
80 halted: Option<String>,
81}
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,
115 budget_usd: self.budget_usd.filter(|_| member),
116 time_cap_minutes: self.time_cap_minutes,
117 halted: self.halted,
118 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
138/// trusted service asks for. What it is given reads and nothing more.
139pub(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) {
143 viewer.workspaces.push(Membership {
144 base_permission: Some(g1t_contracts::access::BasePermission::Read),
145 ..Membership::member(slug)
146 });
147 }
148 Some(viewer)
149}
150
151/// 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).
154fn 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
184 /// are a member of its workspace (which shows what runs cost).
185 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 };
229 // 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 }
236 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,
248 sandbox, token_hash, started_by, created_at, updated_at, budget_usd, time_cap_minutes)
249 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 'queued', ?, ?, ?, ?, ?, ?, ?)",
250 )
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(),
267 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),
269 ])?
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`.
276 pub(crate) fn add_step(&self, run_id: &str, at: &str, text: &str) -> Result<worker::D1PreparedStatement> {
277 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 }
302 // 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 }
309 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());
362 let (repo, _) = match self.visible_repo(&a.repo, &viewer).await? {
363 Outcome::Ok(found) => found,
364 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
365 };
366 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));
373 }
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.
441 pub(crate) async fn sweep_silent(&self) -> Result<()> {
442 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}