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