Skip to content
671 linesCodeBlameRaw

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
Merge demo dry run fixes: settled checks, other attempts, queue links, clone box, activity feed, DEMO.md226 // 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 {
Agents and memory, checks and conflicts, profiles, slug renames, custom domains230 Some(pull_id) => match self.pull_by_id(pull_id).await? {
Merge demo dry run fixes: settled checks, other attempts, queue links, clone box, activity feed, DEMO.md231 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 }
Agents and memory, checks and conflicts, profiles, slug renames, custom domains240 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? {
Merge demo dry run fixes: settled checks, other attempts, queue links, clone box, activity feed, DEMO.md243 Outcome::Ok(repo) => (repo.id, RepoPath { namespace: repo.namespace, name: repo.name }),
Agents and memory, checks and conflicts, profiles, slug renames, custom domains244 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
245 },
246 };
Merge demo dry run fixes: settled checks, other attempts, queue links, clone box, activity feed, DEMO.md247 let namespace = path.namespace.to_lowercase();
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look248 // 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,
Merge demo dry run fixes: settled checks, other attempts, queue links, clone box, activity feed, DEMO.md252 format!("{}/{} is archived or deleted, so nothing new starts on it.", path.namespace, path.name),
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look253 ));
254 }
Agents and memory, checks and conflicts, profiles, slug renames, custom domains255 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()
g1t is one name: its agent's work, commits and comments show as @g1t, and nobody can claim g1t or g1t-agent262 .unwrap_or_else(|| g1t_contracts::identity::AGENT_NAME.to_owned());
Agents and memory, checks and conflicts, profiles, slug renames, custom domains263 self.db
264 .prepare(
265 "INSERT INTO agent_runs
266 (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 API267 sandbox, token_hash, started_by, created_at, updated_at, budget_usd, time_cap_minutes)
268 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 'queued', ?, ?, ?, ?, ?, ?, ?)",
Agents and memory, checks and conflicts, profiles, slug renames, custom domains269 )
270 .bind(&[
271 id.as_str().into(),
272 namespace.into(),
273 repo_id.into(),
Merge demo dry run fixes: settled checks, other attempts, queue links, clone box, activity feed, DEMO.md274 format!("{}/{}", path.namespace, path.name).into(),
Agents and memory, checks and conflicts, profiles, slug renames, custom domains275 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(),
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API286 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),
Agents and memory, checks and conflicts, profiles, slug renames, custom domains288 ])?
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`.
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API295 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 domains296 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 }
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API321 // 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 }
Agents and memory, checks and conflicts, profiles, slug renames, custom domains328 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());
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look381 let (repo, _) = match self.visible_repo(&a.repo, &viewer).await? {
Agents and memory, checks and conflicts, profiles, slug renames, custom domains382 Outcome::Ok(found) => found,
383 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
384 };
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look385 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));
Agents and memory, checks and conflicts, profiles, slug renames, custom domains392 }
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 ),
Events: review requests, assignments, stops and deployments are published439 by: Some(a.actor.id.clone()),
Agents and memory, checks and conflicts, profiles, slug renames, custom domains440 })
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.
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look461 pub(crate) async fn sweep_silent(&self) -> Result<()> {
Agents and memory, checks and conflicts, profiles, slug renames, custom domains462 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}

This file's history is long; its oldest lines are credited to the oldest commit read.