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

636 lines25,387 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.
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::member(slug));
144 }
145 Some(viewer)
146}
147
148fn is_member(viewer: &Viewer, namespace: &str) -> bool {
149 viewer
150 .as_ref()
151 .is_some_and(|viewer| viewer.is_member(&namespace.to_lowercase()))
152}
153
154#[derive(Deserialize)]
155struct SessionListRow {
156 pull_id: String,
157 number: u32,
158 title: String,
159 status: PullStatus,
160 agent: String,
161 entries: u32,
162 tools: u32,
163 started_at: String,
164 last_at: String,
165 prompt: Option<String>,
166}
167
168#[derive(Deserialize)]
169struct RunSumRow {
170 pull_id: String,
171 kind: String,
172 status: String,
173 cost_usd: Option<f64>,
174}
175
176impl Work {
177 /// The repository at `path`, if `viewer` may see it, and whether they
178 /// are a member of its workspace.
179 async fn visible_repo(&self, path: &RepoPath, viewer: &Viewer) -> Result<Outcome<(Repo, bool)>> {
180 Ok(match self.repo(path, viewer).await? {
181 Outcome::Ok(repo) => {
182 let member = is_member(viewer, &repo.namespace);
183 Outcome::Ok((repo, member))
184 }
185 Outcome::Fail(failure) => Outcome::Fail(failure),
186 })
187 }
188
189 async fn run_row(&self, id: &str) -> Result<Option<RunRow>> {
190 self.db
191 .prepare(format!("SELECT {COLUMNS}, steps FROM agent_runs WHERE id = ?"))
192 .bind(&[id.into()])?
193 .first::<RunRow>(None)
194 .await
195 }
196
197 /// The run on pull request `number` that is at work now, if one is.
198 pub(crate) async fn active_run_on(&self, repo_id: &str, number: u32) -> Result<Option<String>> {
199 self.db
200 .prepare(
201 "SELECT id AS value FROM agent_runs
202 WHERE repo_id = ? AND number = ? AND status IN ('queued', 'running')
203 ORDER BY created_at DESC LIMIT 1",
204 )
205 .bind(&[repo_id.into(), number.into()])?
206 .first::<String>(Some("value"))
207 .await
208 }
209
210 pub(crate) async fn open_run(&self, a: OpenRunArgs) -> Result<Outcome<AgentRunTicket>> {
211 // The runner is trusted: it names the repository it is starting a
212 // sandbox in, whoever the sandbox acts as.
213 let (repo_id, namespace) = match &a.pull_id {
214 Some(pull_id) => match self.pull_by_id(pull_id).await? {
215 Some(pull) => (pull.repo_id, a.repo.namespace.to_lowercase()),
216 None => return Ok(Outcome::fail(FailureCode::NotFound, "Pull request not found.")),
217 },
218 None => match self.repo(&a.repo, &member_of(&a.actor, &a.repo.namespace)).await? {
219 Outcome::Ok(repo) => (repo.id, repo.namespace.to_lowercase()),
220 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
221 },
222 };
223 let now = now_ms();
224 let id = new_id("arn", now);
225 let token = new_token();
226 let timestamp = rfc3339(now);
227 let agent = a
228 .agent
229 .clone()
230 .unwrap_or_else(|| if a.kind.is_agent() { "g1t-agent" } else { "g1t" }.to_owned());
231 self.db
232 .prepare(
233 "INSERT INTO agent_runs
234 (id, workspace, repo_id, repo, number, pull_id, title, kind, agent, model, status,
235 sandbox, token_hash, started_by, created_at, updated_at, budget_usd, time_cap_minutes)
236 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 'queued', ?, ?, ?, ?, ?, ?, ?)",
237 )
238 .bind(&[
239 id.as_str().into(),
240 namespace.into(),
241 repo_id.into(),
242 format!("{}/{}", a.repo.namespace, a.repo.name).into(),
243 a.number.filter(|n| *n > 0).map_or(JsValue::NULL, JsValue::from),
244 optional(&a.pull_id),
245 optional(&a.title.map(|title| one_line(&title, 200))),
246 a.kind.as_str().into(),
247 agent.into(),
248 optional(&a.model),
249 a.sandbox.into(),
250 hash(&token).into(),
251 optional(&a.started_by),
252 timestamp.as_str().into(),
253 timestamp.as_str().into(),
254 a.budget_usd.filter(|usd| usd.is_finite() && *usd > 0.0).map_or(JsValue::NULL, JsValue::from),
255 a.time_cap_minutes.map_or(JsValue::NULL, JsValue::from),
256 ])?
257 .run()
258 .await?;
259 Ok(Outcome::Ok(AgentRunTicket { run_id: id, token }))
260 }
261
262 /// Adds `text` to a run's steps, keeping the latest `MAX_STEPS`.
263 pub(crate) fn add_step(&self, run_id: &str, at: &str, text: &str) -> Result<worker::D1PreparedStatement> {
264 self.db
265 .prepare(
266 "UPDATE agent_runs SET
267 steps = CASE WHEN json_array_length(steps) >= ?4
268 THEN json_insert(json_remove(steps, '$[0]'), '$[#]', json_object('at', ?1, 'text', ?2))
269 ELSE json_insert(steps, '$[#]', json_object('at', ?1, 'text', ?2)) END,
270 step_count = step_count + 1
271 WHERE id = ?3",
272 )
273 .bind(&[at.into(), text.into(), run_id.into(), MAX_STEPS.into()])
274 }
275
276 pub(crate) async fn report_run(&self, a: ReportRunArgs) -> Result<Outcome<RunStatus>> {
277 let Some(run) = self
278 .run_row(&a.run_id)
279 .await?
280 .filter(|run| run.token_hash == hash(&a.token))
281 else {
282 return Ok(Outcome::fail(FailureCode::NotFound, "Run not found."));
283 };
284 // Finished, or stopped by a person: nothing more is taken, and the
285 // sandbox learns why.
286 if run.finished_at.is_some() {
287 return Ok(Outcome::Ok(run.status()));
288 }
289 // It reached a cap of its guardrails (guardrails.rs).
290 if let Some(halt) = a.halt {
291 let pull_id = run.pull_id.as_deref();
292 return self
293 .halt_run(&run.id, &run.repo_id, pull_id, run.number, &run.kind, halt, a.error, a.cost_usd)
294 .await;
295 }
296 let now = rfc3339(now_ms());
297 let steps: Vec<String> = a
298 .steps
299 .iter()
300 .map(|step| one_line(step, MAX_STEP_CHARS))
301 .filter(|step| !step.is_empty())
302 .collect();
303 // A burst keeps its end: that is where the run is now.
304 let skip = steps.len().saturating_sub(MAX_STEPS_PER_REPORT);
305 let mut statements = Vec::new();
306 for step in &steps[skip..] {
307 statements.push(self.add_step(&run.id, &now, step)?);
308 }
309 let current = a
310 .step
311 .as_deref()
312 .map(|step| one_line(step, MAX_STEP_CHARS))
313 .filter(|step| !step.is_empty())
314 .or_else(|| steps.last().cloned());
315 let outcome = a
316 .outcome
317 .filter(|outcome| matches!(outcome, RunStatus::Succeeded | RunStatus::Failed));
318 let cost = a.cost_usd.filter(|cost| cost.is_finite() && *cost >= 0.0);
319 statements.push(
320 self.db
321 .prepare(
322 "UPDATE agent_runs SET
323 status = COALESCE(?1, CASE WHEN status = 'queued' THEN 'running' ELSE status END),
324 started_at = COALESCE(started_at, ?2),
325 step = COALESCE(?3, step),
326 cost_usd = COALESCE(?4, cost_usd),
327 turns = COALESCE(?5, turns),
328 error = COALESCE(?6, error),
329 finished_at = CASE WHEN ?1 IS NULL THEN finished_at ELSE ?2 END,
330 updated_at = ?2
331 WHERE id = ?7 AND finished_at IS NULL",
332 )
333 .bind(&[
334 outcome.map_or(JsValue::NULL, |outcome| outcome.as_str().into()),
335 now.as_str().into(),
336 optional(&current),
337 cost.map_or(JsValue::NULL, JsValue::from),
338 a.turns.map_or(JsValue::NULL, JsValue::from),
339 optional(&a.error.map(|error| one_line(&error, 1000))),
340 run.id.as_str().into(),
341 ])?,
342 );
343 self.db.batch(statements).await?;
344 Ok(Outcome::Ok(outcome.unwrap_or(RunStatus::Running)))
345 }
346
347 pub(crate) async fn stop_run(&self, a: StopRunArgs) -> Result<Outcome<StoppedRun>> {
348 let viewer = Some(a.actor.clone());
349 let (repo, member) = match self.visible_repo(&a.repo, &viewer).await? {
350 Outcome::Ok(found) => found,
351 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
352 };
353 if !member || !a.actor.verified {
354 return Ok(Outcome::fail(
355 FailureCode::Forbidden,
356 "Only members of the workspace can stop its agents.",
357 ));
358 }
359 let Some(run) = self
360 .run_row(&a.id)
361 .await?
362 .filter(|run| run.repo_id == repo.id)
363 else {
364 return Ok(Outcome::fail(FailureCode::NotFound, "Run not found."));
365 };
366 if !run.status().is_active() {
367 return Ok(Outcome::fail(
368 FailureCode::Conflict,
369 format!("This run has already {}.", match run.status() {
370 RunStatus::Stopped => "been stopped",
371 RunStatus::Succeeded => "finished",
372 _ => "ended",
373 }),
374 ));
375 }
376 let now = rfc3339(now_ms());
377 let said = format!("Stopped by {}.", a.actor.username);
378 let claimed = self
379 .db
380 .prepare(
381 "UPDATE agent_runs SET status = 'stopped', step = ?, finished_at = ?, updated_at = ?
382 WHERE id = ? AND status IN ('queued', 'running') RETURNING id AS value",
383 )
384 .bind(&[
385 said.as_str().into(),
386 now.as_str().into(),
387 now.as_str().into(),
388 run.id.as_str().into(),
389 ])?
390 .first::<String>(Some("value"))
391 .await?;
392 if claimed.is_none() {
393 return Ok(Outcome::fail(FailureCode::Conflict, "This run has already ended."));
394 }
395 self.add_step(&run.id, &now, &said)?.run().await?;
396 // g1t stops seeing the pull request through, so it does not start
397 // the same work again; a person decides what happens next.
398 if let (Some(pull_id), Some(number)) = (&run.pull_id, run.number) {
399 self.stall(StallArgs {
400 pull_id: pull_id.clone(),
401 reason: format!(
402 "{} stopped the agent's {} run. Ask for a review, a revision or a catch-up to start again.",
403 a.actor.username, run.kind
404 ),
405 })
406 .await?;
407 self.note(
408 &repo.id,
409 number,
410 (a.actor.id.as_str(), a.actor.username.as_str()),
411 &format!("stopped {}'s {} run", run.agent, run.kind),
412 )
413 .await?;
414 }
415 let sandbox = run.sandbox.clone();
416 let Some(row) = self.run_row(&run.id).await? else {
417 return Ok(Outcome::fail(FailureCode::NotFound, "Run not found."));
418 };
419 Ok(Outcome::Ok(StoppedRun {
420 run: row.into_run(true),
421 sandbox,
422 }))
423 }
424
425 /// Runs whose sandbox died without anyone noticing, marked failed.
426 async fn sweep_silent(&self) -> Result<()> {
427 let now = now_ms();
428 let cutoff = rfc3339(now.saturating_sub(SILENT_HOURS * 3_600_000));
429 self.db
430 .prepare(
431 "UPDATE agent_runs
432 SET status = 'failed', error = COALESCE(error, 'It stopped reporting.'),
433 finished_at = ?, updated_at = ?
434 WHERE status IN ('queued', 'running') AND updated_at < ?",
435 )
436 .bind(&[rfc3339(now).into(), rfc3339(now).into(), cutoff.into()])?
437 .run()
438 .await?;
439 Ok(())
440 }
441
442 pub(crate) async fn list_runs(&self, a: ListRunsArgs) -> Result<Outcome<Vec<AgentRun>>> {
443 let (column, key, member) = if let Some(workspace) = &a.workspace {
444 let slug = workspace.to_lowercase();
445 if !is_member(&a.viewer, &slug) {
446 return Ok(Outcome::fail(
447 FailureCode::Forbidden,
448 "Only members of the workspace can see its fleet.",
449 ));
450 }
451 ("workspace", slug, true)
452 } else if let Some(path) = &a.repo {
453 match self.visible_repo(path, &a.viewer).await? {
454 Outcome::Ok((repo, member)) => ("repo_id", repo.id, member),
455 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
456 }
457 } else {
458 return Ok(Outcome::fail(FailureCode::Invalid, "Name a repository or a workspace."));
459 };
460 self.sweep_silent().await?;
461 let limit = a.limit.unwrap_or(DEFAULT_LIST).clamp(1, MAX_LIST);
462 let rows = self
463 .db
464 .prepare(format!(
465 "SELECT {COLUMNS} FROM agent_runs
466 WHERE {column} = ?1
467 AND (?2 = 0 OR status IN ('queued', 'running'))
468 AND (?3 IS NULL OR kind = ?3)
469 AND (?4 IS NULL OR status = ?4)
470 AND (?5 IS NULL OR number = ?5)
471 ORDER BY created_at DESC LIMIT ?6"
472 ))
473 .bind(&[
474 key.into(),
475 u32::from(a.active).into(),
476 a.kind.map_or(JsValue::NULL, |kind| kind.as_str().into()),
477 a.status.map_or(JsValue::NULL, |status| status.as_str().into()),
478 a.number.map_or(JsValue::NULL, JsValue::from),
479 limit.into(),
480 ])?
481 .all()
482 .await?
483 .results::<RunRow>()?;
484 Ok(Outcome::Ok(rows.into_iter().map(|row| row.into_run(member)).collect()))
485 }
486
487 pub(crate) async fn get_run(&self, a: GetRunArgs) -> Result<Outcome<AgentRun>> {
488 let (repo, member) = match self.visible_repo(&a.repo, &a.viewer).await? {
489 Outcome::Ok(found) => found,
490 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
491 };
492 Ok(match self.run_row(&a.id).await?.filter(|run| run.repo_id == repo.id) {
493 Some(row) => Outcome::Ok(row.into_run(member)),
494 None => Outcome::fail(FailureCode::NotFound, "Run not found."),
495 })
496 }
497
498 // --- Sessions ----------------------------------------------------------
499
500 pub(crate) async fn list_sessions(&self, a: ListSessionsArgs) -> Result<Outcome<Vec<SessionSummary>>> {
501 let Some(path) = &a.repo else {
502 return Ok(Outcome::fail(FailureCode::Invalid, "Name a repository."));
503 };
504 let (repo, member) = match self.visible_repo(path, &a.viewer).await? {
505 Outcome::Ok(found) => found,
506 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
507 };
508 let rows = self
509 .db
510 .prepare(
511 "SELECT p.id AS pull_id, p.number, p.title, p.status, p.agent,
512 count(*) AS entries,
513 sum(CASE WHEN e.kind = 'tool_call' THEN 1 ELSE 0 END) AS tools,
514 min(e.at) AS started_at, max(e.at) AS last_at,
515 (SELECT substr(f.text, 1, 300) FROM session_entries f
516 WHERE f.pull_id = p.id AND f.kind = 'prompt' ORDER BY f.seq LIMIT 1) AS prompt
517 FROM pulls p JOIN session_entries e ON e.pull_id = p.id
518 WHERE p.repo_id = ?1
519 AND (?2 IS NULL OR p.status = ?2)
520 AND (?3 IS NULL OR p.number = ?3)
521 AND (?4 IS NULL OR EXISTS
522 (SELECT 1 FROM agent_runs r WHERE r.pull_id = p.id AND r.kind = ?4))
523 GROUP BY p.id ORDER BY last_at DESC LIMIT 100",
524 )
525 .bind(&[
526 repo.id.as_str().into(),
527 a.outcome.map_or(JsValue::NULL, |status| status.as_str().into()),
528 a.number.map_or(JsValue::NULL, JsValue::from),
529 a.kind.map_or(JsValue::NULL, |kind| kind.as_str().into()),
530 ])?
531 .all()
532 .await?
533 .results::<SessionListRow>()?;
534 let ids: Vec<&str> = rows.iter().map(|row| row.pull_id.as_str()).collect();
535 let runs = if ids.is_empty() {
536 Vec::new()
537 } else {
538 self.db
539 .prepare(
540 "SELECT pull_id, kind, status, cost_usd FROM agent_runs
541 WHERE repo_id = ? AND pull_id IN (SELECT value FROM json_each(?))",
542 )
543 .bind(&[repo.id.as_str().into(), serde_json::to_string(&ids)?.into()])?
544 .all()
545 .await?
546 .results::<RunSumRow>()?
547 };
548 Ok(Outcome::Ok(
549 rows.into_iter()
550 .map(|row| {
551 let mine: Vec<&RunSumRow> = runs.iter().filter(|run| run.pull_id == row.pull_id).collect();
552 let mut kinds: Vec<RunKind> = Vec::new();
553 for run in &mine {
554 if let Some(kind) = RunKind::parse(&run.kind)
555 && !kinds.contains(&kind)
556 {
557 kinds.push(kind);
558 }
559 }
560 let spent: f64 = mine.iter().filter_map(|run| run.cost_usd).sum();
561 SessionSummary {
562 number: row.number,
563 title: row.title,
564 status: row.status,
565 agent: row.agent,
566 entries: row.entries,
567 tools: row.tools,
568 prompt: row.prompt.map(|prompt| one_line(&prompt, 200)),
569 runs: mine.len() as u32,
570 cost_usd: (member && mine.iter().any(|run| run.cost_usd.is_some())).then_some(spent),
571 active: mine
572 .iter()
573 .any(|run| RunStatus::parse(&run.status).is_some_and(RunStatus::is_active)),
574 kinds,
575 started_at: row.started_at,
576 last_at: row.last_at,
577 }
578 })
579 .collect(),
580 ))
581 }
582
583 pub(crate) async fn get_session(&self, a: GetSessionArgs) -> Result<Outcome<SessionView>> {
584 let (repo, pull) = match self.pull_at(&a.repo, a.number, &a.viewer).await? {
585 Outcome::Ok(found) => found,
586 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
587 };
588 let member = is_member(&a.viewer, &repo.namespace);
589 let entries: Vec<SessionEntry> = self
590 .db
591 .prepare(
592 "SELECT seq, kind, text, tool, \"commit\", at FROM session_entries
593 WHERE pull_id = ? ORDER BY seq LIMIT ?",
594 )
595 .bind(&[pull.id.as_str().into(), SESSION_ENTRIES.into()])?
596 .all()
597 .await?
598 .results::<SessionRow>()?
599 .into_iter()
600 .map(SessionEntry::from)
601 .collect();
602 let runs: Vec<AgentRun> = self
603 .db
604 .prepare(format!(
605 "SELECT {COLUMNS}, steps FROM agent_runs WHERE pull_id = ? ORDER BY created_at DESC LIMIT 50"
606 ))
607 .bind(&[pull.id.as_str().into()])?
608 .all()
609 .await?
610 .results::<RunRow>()?
611 .into_iter()
612 .map(|row| row.into_run(member))
613 .collect();
614 let cost_usd = (member && runs.iter().any(|run| run.cost_usd.is_some()))
615 .then(|| runs.iter().filter_map(|run| run.cost_usd).sum());
616 Ok(Outcome::Ok(SessionView {
617 pull,
618 entries,
619 runs,
620 cost_usd,
621 }))
622 }
623}
624
625#[cfg(test)]
626mod tests {
627 use super::*;
628
629 #[test]
630 fn steps_are_one_short_line() {
631 assert_eq!(one_line(" Read\n src/lib.rs ", 40), "Read src/lib.rs");
632 let long = one_line(&"x".repeat(300), 10);
633 assert_eq!(long.chars().count(), 10);
634 assert!(long.ends_with('…'));
635 }
636}