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

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,
Thirteen MCP tools and classic token scopes; agents rate their confidence and can be put on an issue in one step44 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";
Agents and memory, checks and conflicts, profiles, slug renames, custom domains46
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>,
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API76 #[serde(default)]
77 budget_usd: Option<f64>,
78 #[serde(default)]
79 time_cap_minutes: Option<u32>,
80 #[serde(default)]
81 halted: Option<String>,
Thirteen MCP tools and classic token scopes; agents rate their confidence and can be put on an issue in one step82 /// JSON of how sure g1t was of the change as the run left it.
83 #[serde(default)]
84 confidence: Option<String>,
Agents and memory, checks and conflicts, profiles, slug renames, custom domains85}
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,
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API119 budget_usd: self.budget_usd.filter(|_| member),
120 time_cap_minutes: self.time_cap_minutes,
121 halted: self.halted,
Thirteen MCP tools and classic token scopes; agents rate their confidence and can be put on an issue in one step122 confidence: self
123 .confidence
124 .as_deref()
125 .and_then(|detail| serde_json::from_str(detail).ok()),
Agents and memory, checks and conflicts, profiles, slug renames, custom domains126 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
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look146/// trusted service asks for. What it is given reads and nothing more.
Agents and memory, checks and conflicts, profiles, slug renames, custom domains147pub(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) {
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look151 viewer.workspaces.push(Membership {
152 base_permission: Some(g1t_contracts::access::BasePermission::Read),
153 ..Membership::member(slug)
154 });
Agents and memory, checks and conflicts, profiles, slug renames, custom domains155 }
156 Some(viewer)
157}
158
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look159/// 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).
Agents and memory, checks and conflicts, profiles, slug renames, custom domains162fn 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
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look192 /// are a member of its workspace (which shows what runs cost).
Agents and memory, checks and conflicts, profiles, slug renames, custom domains193 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 };
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look237 // 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 }
Agents and memory, checks and conflicts, profiles, slug renames, custom domains244 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,
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API256 sandbox, token_hash, started_by, created_at, updated_at, budget_usd, time_cap_minutes)
257 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 'queued', ?, ?, ?, ?, ?, ?, ?)",
Agents and memory, checks and conflicts, profiles, slug renames, custom domains258 )
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(),
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API275 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),
Agents and memory, checks and conflicts, profiles, slug renames, custom domains277 ])?
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`.
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API284 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 domains285 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 }
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API310 // 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 }
Agents and memory, checks and conflicts, profiles, slug renames, custom domains317 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());
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look370 let (repo, _) = match self.visible_repo(&a.repo, &viewer).await? {
Agents and memory, checks and conflicts, profiles, slug renames, custom domains371 Outcome::Ok(found) => found,
372 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
373 };
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look374 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));
Agents and memory, checks and conflicts, profiles, slug renames, custom domains381 }
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.
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look449 pub(crate) async fn sweep_silent(&self) -> Result<()> {
Agents and memory, checks and conflicts, profiles, slug renames, custom domains450 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}