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

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

Agents and memory, checks and conflicts, profiles, slug renames, custom domains1//! 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,
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API43 started_at, finished_at, updated_at, budget_usd, time_cap_minutes, halted,
Agents and memory, checks and conflicts, profiles, slug renames, custom domains44 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>,
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API75 #[serde(default)]
76 budget_usd: Option<f64>,
77 #[serde(default)]
78 time_cap_minutes: Option<u32>,
79 #[serde(default)]
80 halted: Option<String>,
Agents and memory, checks and conflicts, profiles, slug renames, custom domains81}
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,
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API115 budget_usd: self.budget_usd.filter(|_| member),
116 time_cap_minutes: self.time_cap_minutes,
117 halted: self.halted,
Agents and memory, checks and conflicts, profiles, slug renames, custom domains118 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,
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API235 sandbox, token_hash, started_by, created_at, updated_at, budget_usd, time_cap_minutes)
236 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 'queued', ?, ?, ?, ?, ?, ?, ?)",
Agents and memory, checks and conflicts, profiles, slug renames, custom domains237 )
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(),
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API254 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),
Agents and memory, checks and conflicts, profiles, slug renames, custom domains256 ])?
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`.
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API263 pub(crate) fn add_step(&self, run_id: &str, at: &str, text: &str) -> Result<worker::D1PreparedStatement> {
Agents and memory, checks and conflicts, profiles, slug renames, custom domains264 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 }
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API289 // 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 }
Agents and memory, checks and conflicts, profiles, slug renames, custom domains296 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}