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

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