Skip to content
692 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
A host a sandbox was refused is a note on the run, never what its card says it is doing, and tools in a sandbox with only allowed hosts are told not to send usage reports home27/// How the runner starts a step that notes something about the sandbox
28/// rather than what the run is doing: a host it was refused (the runner's
29/// egress.ts `blockedStep`).
30const NOTE_PREFIXES: [&str; 1] = ["Blocked: "];
31
32/// Where a report leaves the run: its newest step that says what the run
33/// is doing. A note (a refused host) stays in the steps but never stands
34/// for the run, which carries on past it. `None` when there is only notes.
35fn current_step(steps: &[String]) -> Option<String> {
36 steps.iter().rev().find(|step| !NOTE_PREFIXES.iter().any(|prefix| step.starts_with(prefix))).cloned()
37}
38
Agents and memory, checks and conflicts, profiles, slug renames, custom domains39/// The most steps a run keeps; older ones fall off the start.
40const MAX_STEPS: u32 = 200;
41/// The most steps taken from one report.
42const MAX_STEPS_PER_REPORT: usize = 20;
43/// The longest a step is kept.
44const MAX_STEP_CHARS: usize = 240;
45/// A run that has said nothing for this long is taken to have died.
46const SILENT_HOURS: u64 = 3;
47const DEFAULT_LIST: u32 = 50;
48const MAX_LIST: u32 = 200;
49/// The most session entries one view returns.
50const SESSION_ENTRIES: u32 = 2000;
51
52/// Every column but the steps, which only a run's own page needs.
53const COLUMNS: &str = "id, workspace, repo_id, repo, number, pull_id, kind, agent, model, status,
54 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 API55 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 step56 COALESCE((SELECT title FROM pulls WHERE pulls.id = agent_runs.pull_id), agent_runs.title) AS title,
57 (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 domains58
59#[derive(Deserialize)]
60pub(crate) struct RunRow {
61 id: String,
62 #[allow(dead_code)]
63 workspace: String,
64 #[allow(dead_code)]
65 repo_id: String,
66 repo: String,
67 number: Option<u32>,
68 pub(crate) pull_id: Option<String>,
69 kind: String,
70 agent: String,
71 model: Option<String>,
72 status: String,
73 step: Option<String>,
74 #[serde(default)]
75 steps: Option<String>,
76 step_count: u32,
77 cost_usd: Option<f64>,
78 turns: Option<u32>,
79 sandbox: String,
80 token_hash: String,
81 started_by: Option<String>,
82 error: Option<String>,
83 created_at: String,
84 started_at: Option<String>,
85 finished_at: Option<String>,
86 updated_at: String,
87 title: Option<String>,
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API88 #[serde(default)]
89 budget_usd: Option<f64>,
90 #[serde(default)]
91 time_cap_minutes: Option<u32>,
92 #[serde(default)]
93 halted: Option<String>,
Thirteen MCP tools and classic token scopes; agents rate their confidence and can be put on an issue in one step94 /// JSON of how sure g1t was of the change as the run left it.
95 #[serde(default)]
96 confidence: Option<String>,
Agents and memory, checks and conflicts, profiles, slug renames, custom domains97}
98
99impl RunRow {
100 fn status(&self) -> RunStatus {
101 RunStatus::parse(&self.status).unwrap_or(RunStatus::Failed)
102 }
103
104 /// The run as `member` (or not) may see it: model and cost are the
105 /// workspace's business.
106 fn into_run(self, member: bool) -> AgentRun {
107 let (namespace, name) = self.repo.split_once('/').unwrap_or((&self.repo, ""));
108 AgentRun {
109 repo: RepoPath {
110 namespace: namespace.to_owned(),
111 name: name.to_owned(),
112 },
113 status: self.status(),
114 kind: RunKind::parse(&self.kind).unwrap_or(RunKind::Implement),
115 steps: self
116 .steps
117 .as_deref()
118 .and_then(|steps| serde_json::from_str(steps).ok())
119 .unwrap_or_default(),
120 id: self.id,
121 number: self.number.filter(|number| *number > 0),
122 title: self.title,
123 agent: self.agent,
124 model: self.model.filter(|_| member),
125 step: self.step,
126 step_count: self.step_count,
127 started_by: self.started_by,
128 error: self.error,
129 cost_usd: self.cost_usd.filter(|_| member),
130 turns: self.turns,
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API131 budget_usd: self.budget_usd.filter(|_| member),
132 time_cap_minutes: self.time_cap_minutes,
133 halted: self.halted,
Thirteen MCP tools and classic token scopes; agents rate their confidence and can be put on an issue in one step134 confidence: self
135 .confidence
136 .as_deref()
137 .and_then(|detail| serde_json::from_str(detail).ok()),
Agents and memory, checks and conflicts, profiles, slug renames, custom domains138 created_at: self.created_at,
139 started_at: self.started_at,
140 finished_at: self.finished_at,
141 updated_at: self.updated_at,
142 }
143 }
144}
145
146/// One line, short enough to read at a glance.
147pub(crate) fn one_line(text: &str, limit: usize) -> String {
148 let line = text.split_whitespace().collect::<Vec<_>>().join(" ");
149 if line.chars().count() <= limit {
150 return line;
151 }
152 let mut short: String = line.chars().take(limit.saturating_sub(1)).collect();
153 short.push('…');
154 short
155}
156
157/// 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 look158/// trusted service asks for. What it is given reads and nothing more.
Agents and memory, checks and conflicts, profiles, slug renames, custom domains159pub(crate) fn member_of(actor: &User, namespace: &str) -> Viewer {
160 let mut viewer = actor.clone();
161 let slug = namespace.to_lowercase();
162 if !viewer.is_member(&slug) {
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look163 viewer.workspaces.push(Membership {
164 base_permission: Some(g1t_contracts::access::BasePermission::Read),
165 ..Membership::member(slug)
166 });
Agents and memory, checks and conflicts, profiles, slug renames, custom domains167 }
168 Some(viewer)
169}
170
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look171/// Whether the viewer belongs to the workspace: what a run spent, its model
172/// and its budget are the workspace's business, not every reader's (an
173/// outside collaborator sees the work, not the bill).
Agents and memory, checks and conflicts, profiles, slug renames, custom domains174fn is_member(viewer: &Viewer, namespace: &str) -> bool {
175 viewer
176 .as_ref()
177 .is_some_and(|viewer| viewer.is_member(&namespace.to_lowercase()))
178}
179
180#[derive(Deserialize)]
181struct SessionListRow {
182 pull_id: String,
183 number: u32,
184 title: String,
185 status: PullStatus,
186 agent: String,
187 entries: u32,
188 tools: u32,
189 started_at: String,
190 last_at: String,
191 prompt: Option<String>,
192}
193
194#[derive(Deserialize)]
195struct RunSumRow {
196 pull_id: String,
197 kind: String,
198 status: String,
199 cost_usd: Option<f64>,
200}
201
202impl Work {
203 /// 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 look204 /// are a member of its workspace (which shows what runs cost).
Agents and memory, checks and conflicts, profiles, slug renames, custom domains205 async fn visible_repo(&self, path: &RepoPath, viewer: &Viewer) -> Result<Outcome<(Repo, bool)>> {
206 Ok(match self.repo(path, viewer).await? {
207 Outcome::Ok(repo) => {
208 let member = is_member(viewer, &repo.namespace);
209 Outcome::Ok((repo, member))
210 }
211 Outcome::Fail(failure) => Outcome::Fail(failure),
212 })
213 }
214
215 async fn run_row(&self, id: &str) -> Result<Option<RunRow>> {
216 self.db
217 .prepare(format!("SELECT {COLUMNS}, steps FROM agent_runs WHERE id = ?"))
218 .bind(&[id.into()])?
219 .first::<RunRow>(None)
220 .await
221 }
222
223 /// The run on pull request `number` that is at work now, if one is.
224 pub(crate) async fn active_run_on(&self, repo_id: &str, number: u32) -> Result<Option<String>> {
225 self.db
226 .prepare(
227 "SELECT id AS value FROM agent_runs
228 WHERE repo_id = ? AND number = ? AND status IN ('queued', 'running')
229 ORDER BY created_at DESC LIMIT 1",
230 )
231 .bind(&[repo_id.into(), number.into()])?
232 .first::<String>(Some("value"))
233 .await
234 }
235
236 pub(crate) async fn open_run(&self, a: OpenRunArgs) -> Result<Outcome<AgentRunTicket>> {
237 // 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.md238 // sandbox in, whoever the sandbox acts as. The run records where
239 // the repository is now, looked up by its id, not the path it was
240 // named by, which may be from before a transfer or a rename.
241 let (repo_id, path) = match &a.pull_id {
Agents and memory, checks and conflicts, profiles, slug renames, custom domains242 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.md243 Some(pull) => {
244 let now: Option<RepoPath> = g1t_kit::call(
245 &self.repos,
246 "path_by_id",
247 &g1t_contracts::repos::PathByIdArgs { id: pull.repo_id.clone() },
248 )
249 .await?;
250 (pull.repo_id, now.unwrap_or_else(|| a.repo.clone()))
251 }
Agents and memory, checks and conflicts, profiles, slug renames, custom domains252 None => return Ok(Outcome::fail(FailureCode::NotFound, "Pull request not found.")),
253 },
254 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.md255 Outcome::Ok(repo) => (repo.id, RepoPath { namespace: repo.namespace, name: repo.name }),
Agents and memory, checks and conflicts, profiles, slug renames, custom domains256 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
257 },
258 };
Merge demo dry run fixes: settled checks, other attempts, queue links, clone box, activity feed, DEMO.md259 let namespace = path.namespace.to_lowercase();
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look260 // Nothing new starts on an archived or deleted repository.
261 if !self.repo_active(&repo_id).await? {
262 return Ok(Outcome::fail(
263 FailureCode::Forbidden,
Merge demo dry run fixes: settled checks, other attempts, queue links, clone box, activity feed, DEMO.md264 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 look265 ));
266 }
Agents and memory, checks and conflicts, profiles, slug renames, custom domains267 let now = now_ms();
268 let id = new_id("arn", now);
269 let token = new_token();
270 let timestamp = rfc3339(now);
271 let agent = a
272 .agent
273 .clone()
g1t is one name: its agent's work, commits and comments show as @g1t, and nobody can claim g1t or g1t-agent274 .unwrap_or_else(|| g1t_contracts::identity::AGENT_NAME.to_owned());
Agents and memory, checks and conflicts, profiles, slug renames, custom domains275 self.db
276 .prepare(
277 "INSERT INTO agent_runs
278 (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 API279 sandbox, token_hash, started_by, created_at, updated_at, budget_usd, time_cap_minutes)
280 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 'queued', ?, ?, ?, ?, ?, ?, ?)",
Agents and memory, checks and conflicts, profiles, slug renames, custom domains281 )
282 .bind(&[
283 id.as_str().into(),
284 namespace.into(),
285 repo_id.into(),
Merge demo dry run fixes: settled checks, other attempts, queue links, clone box, activity feed, DEMO.md286 format!("{}/{}", path.namespace, path.name).into(),
Agents and memory, checks and conflicts, profiles, slug renames, custom domains287 a.number.filter(|n| *n > 0).map_or(JsValue::NULL, JsValue::from),
288 optional(&a.pull_id),
289 optional(&a.title.map(|title| one_line(&title, 200))),
290 a.kind.as_str().into(),
291 agent.into(),
292 optional(&a.model),
293 a.sandbox.into(),
294 hash(&token).into(),
295 optional(&a.started_by),
296 timestamp.as_str().into(),
297 timestamp.as_str().into(),
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API298 a.budget_usd.filter(|usd| usd.is_finite() && *usd > 0.0).map_or(JsValue::NULL, JsValue::from),
299 a.time_cap_minutes.map_or(JsValue::NULL, JsValue::from),
Agents and memory, checks and conflicts, profiles, slug renames, custom domains300 ])?
301 .run()
302 .await?;
303 Ok(Outcome::Ok(AgentRunTicket { run_id: id, token }))
304 }
305
306 /// 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 API307 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 domains308 self.db
309 .prepare(
310 "UPDATE agent_runs SET
311 steps = CASE WHEN json_array_length(steps) >= ?4
312 THEN json_insert(json_remove(steps, '$[0]'), '$[#]', json_object('at', ?1, 'text', ?2))
313 ELSE json_insert(steps, '$[#]', json_object('at', ?1, 'text', ?2)) END,
314 step_count = step_count + 1
315 WHERE id = ?3",
316 )
317 .bind(&[at.into(), text.into(), run_id.into(), MAX_STEPS.into()])
318 }
319
320 pub(crate) async fn report_run(&self, a: ReportRunArgs) -> Result<Outcome<RunStatus>> {
321 let Some(run) = self
322 .run_row(&a.run_id)
323 .await?
324 .filter(|run| run.token_hash == hash(&a.token))
325 else {
326 return Ok(Outcome::fail(FailureCode::NotFound, "Run not found."));
327 };
328 // Finished, or stopped by a person: nothing more is taken, and the
329 // sandbox learns why.
330 if run.finished_at.is_some() {
331 return Ok(Outcome::Ok(run.status()));
332 }
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API333 // It reached a cap of its guardrails (guardrails.rs).
334 if let Some(halt) = a.halt {
335 let pull_id = run.pull_id.as_deref();
336 return self
337 .halt_run(&run.id, &run.repo_id, pull_id, run.number, &run.kind, halt, a.error, a.cost_usd)
338 .await;
339 }
Agents and memory, checks and conflicts, profiles, slug renames, custom domains340 let now = rfc3339(now_ms());
341 let steps: Vec<String> = a
342 .steps
343 .iter()
344 .map(|step| one_line(step, MAX_STEP_CHARS))
345 .filter(|step| !step.is_empty())
346 .collect();
347 // A burst keeps its end: that is where the run is now.
348 let skip = steps.len().saturating_sub(MAX_STEPS_PER_REPORT);
349 let mut statements = Vec::new();
350 for step in &steps[skip..] {
351 statements.push(self.add_step(&run.id, &now, step)?);
352 }
353 let current = a
354 .step
355 .as_deref()
356 .map(|step| one_line(step, MAX_STEP_CHARS))
357 .filter(|step| !step.is_empty())
A host a sandbox was refused is a note on the run, never what its card says it is doing, and tools in a sandbox with only allowed hosts are told not to send usage reports home358 .or_else(|| current_step(&steps));
Agents and memory, checks and conflicts, profiles, slug renames, custom domains359 let outcome = a
360 .outcome
361 .filter(|outcome| matches!(outcome, RunStatus::Succeeded | RunStatus::Failed));
362 let cost = a.cost_usd.filter(|cost| cost.is_finite() && *cost >= 0.0);
363 statements.push(
364 self.db
365 .prepare(
366 "UPDATE agent_runs SET
367 status = COALESCE(?1, CASE WHEN status = 'queued' THEN 'running' ELSE status END),
368 started_at = COALESCE(started_at, ?2),
369 step = COALESCE(?3, step),
370 cost_usd = COALESCE(?4, cost_usd),
371 turns = COALESCE(?5, turns),
372 error = COALESCE(?6, error),
373 finished_at = CASE WHEN ?1 IS NULL THEN finished_at ELSE ?2 END,
374 updated_at = ?2
375 WHERE id = ?7 AND finished_at IS NULL",
376 )
377 .bind(&[
378 outcome.map_or(JsValue::NULL, |outcome| outcome.as_str().into()),
379 now.as_str().into(),
380 optional(&current),
381 cost.map_or(JsValue::NULL, JsValue::from),
382 a.turns.map_or(JsValue::NULL, JsValue::from),
383 optional(&a.error.map(|error| one_line(&error, 1000))),
384 run.id.as_str().into(),
385 ])?,
386 );
387 self.db.batch(statements).await?;
388 Ok(Outcome::Ok(outcome.unwrap_or(RunStatus::Running)))
389 }
390
391 pub(crate) async fn stop_run(&self, a: StopRunArgs) -> Result<Outcome<StoppedRun>> {
392 let viewer = Some(a.actor.clone());
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look393 let (repo, _) = match self.visible_repo(&a.repo, &viewer).await? {
Agents and memory, checks and conflicts, profiles, slug renames, custom domains394 Outcome::Ok(found) => found,
395 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
396 };
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look397 if !a.actor.verified {
398 return Ok(Outcome::fail(FailureCode::Forbidden, crate::UNVERIFIED));
399 }
400 if let Outcome::Fail(failure) =
401 crate::allowed(Some(&a.actor), &repo, g1t_contracts::access::Capability::Run)
402 {
403 return Ok(Outcome::Fail(failure));
Agents and memory, checks and conflicts, profiles, slug renames, custom domains404 }
405 let Some(run) = self
406 .run_row(&a.id)
407 .await?
408 .filter(|run| run.repo_id == repo.id)
409 else {
410 return Ok(Outcome::fail(FailureCode::NotFound, "Run not found."));
411 };
412 if !run.status().is_active() {
413 return Ok(Outcome::fail(
414 FailureCode::Conflict,
415 format!("This run has already {}.", match run.status() {
416 RunStatus::Stopped => "been stopped",
417 RunStatus::Succeeded => "finished",
418 _ => "ended",
419 }),
420 ));
421 }
422 let now = rfc3339(now_ms());
423 let said = format!("Stopped by {}.", a.actor.username);
424 let claimed = self
425 .db
426 .prepare(
427 "UPDATE agent_runs SET status = 'stopped', step = ?, finished_at = ?, updated_at = ?
428 WHERE id = ? AND status IN ('queued', 'running') RETURNING id AS value",
429 )
430 .bind(&[
431 said.as_str().into(),
432 now.as_str().into(),
433 now.as_str().into(),
434 run.id.as_str().into(),
435 ])?
436 .first::<String>(Some("value"))
437 .await?;
438 if claimed.is_none() {
439 return Ok(Outcome::fail(FailureCode::Conflict, "This run has already ended."));
440 }
441 self.add_step(&run.id, &now, &said)?.run().await?;
442 // g1t stops seeing the pull request through, so it does not start
443 // the same work again; a person decides what happens next.
444 if let (Some(pull_id), Some(number)) = (&run.pull_id, run.number) {
445 self.stall(StallArgs {
446 pull_id: pull_id.clone(),
447 reason: format!(
448 "{} stopped the agent's {} run. Ask for a review, a revision or a catch-up to start again.",
449 a.actor.username, run.kind
450 ),
Events: review requests, assignments, stops and deployments are published451 by: Some(a.actor.id.clone()),
Agents and memory, checks and conflicts, profiles, slug renames, custom domains452 })
453 .await?;
454 self.note(
455 &repo.id,
456 number,
457 (a.actor.id.as_str(), a.actor.username.as_str()),
458 &format!("stopped {}'s {} run", run.agent, run.kind),
459 )
460 .await?;
461 }
462 let sandbox = run.sandbox.clone();
463 let Some(row) = self.run_row(&run.id).await? else {
464 return Ok(Outcome::fail(FailureCode::NotFound, "Run not found."));
465 };
466 Ok(Outcome::Ok(StoppedRun {
467 run: row.into_run(true),
468 sandbox,
469 }))
470 }
471
472 /// Runs whose sandbox died without anyone noticing, marked failed.
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look473 pub(crate) async fn sweep_silent(&self) -> Result<()> {
Agents and memory, checks and conflicts, profiles, slug renames, custom domains474 let now = now_ms();
475 let cutoff = rfc3339(now.saturating_sub(SILENT_HOURS * 3_600_000));
476 self.db
477 .prepare(
478 "UPDATE agent_runs
479 SET status = 'failed', error = COALESCE(error, 'It stopped reporting.'),
480 finished_at = ?, updated_at = ?
481 WHERE status IN ('queued', 'running') AND updated_at < ?",
482 )
483 .bind(&[rfc3339(now).into(), rfc3339(now).into(), cutoff.into()])?
484 .run()
485 .await?;
486 Ok(())
487 }
488
489 pub(crate) async fn list_runs(&self, a: ListRunsArgs) -> Result<Outcome<Vec<AgentRun>>> {
490 let (column, key, member) = if let Some(workspace) = &a.workspace {
491 let slug = workspace.to_lowercase();
492 if !is_member(&a.viewer, &slug) {
493 return Ok(Outcome::fail(
494 FailureCode::Forbidden,
495 "Only members of the workspace can see its fleet.",
496 ));
497 }
498 ("workspace", slug, true)
499 } else if let Some(path) = &a.repo {
500 match self.visible_repo(path, &a.viewer).await? {
501 Outcome::Ok((repo, member)) => ("repo_id", repo.id, member),
502 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
503 }
504 } else {
505 return Ok(Outcome::fail(FailureCode::Invalid, "Name a repository or a workspace."));
506 };
507 self.sweep_silent().await?;
508 let limit = a.limit.unwrap_or(DEFAULT_LIST).clamp(1, MAX_LIST);
509 let rows = self
510 .db
511 .prepare(format!(
512 "SELECT {COLUMNS} FROM agent_runs
513 WHERE {column} = ?1
514 AND (?2 = 0 OR status IN ('queued', 'running'))
515 AND (?3 IS NULL OR kind = ?3)
516 AND (?4 IS NULL OR status = ?4)
517 AND (?5 IS NULL OR number = ?5)
518 ORDER BY created_at DESC LIMIT ?6"
519 ))
520 .bind(&[
521 key.into(),
522 u32::from(a.active).into(),
523 a.kind.map_or(JsValue::NULL, |kind| kind.as_str().into()),
524 a.status.map_or(JsValue::NULL, |status| status.as_str().into()),
525 a.number.map_or(JsValue::NULL, JsValue::from),
526 limit.into(),
527 ])?
528 .all()
529 .await?
530 .results::<RunRow>()?;
531 Ok(Outcome::Ok(rows.into_iter().map(|row| row.into_run(member)).collect()))
532 }
533
534 pub(crate) async fn get_run(&self, a: GetRunArgs) -> Result<Outcome<AgentRun>> {
535 let (repo, member) = match self.visible_repo(&a.repo, &a.viewer).await? {
536 Outcome::Ok(found) => found,
537 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
538 };
539 Ok(match self.run_row(&a.id).await?.filter(|run| run.repo_id == repo.id) {
540 Some(row) => Outcome::Ok(row.into_run(member)),
541 None => Outcome::fail(FailureCode::NotFound, "Run not found."),
542 })
543 }
544
545 // --- Sessions ----------------------------------------------------------
546
547 pub(crate) async fn list_sessions(&self, a: ListSessionsArgs) -> Result<Outcome<Vec<SessionSummary>>> {
548 let Some(path) = &a.repo else {
549 return Ok(Outcome::fail(FailureCode::Invalid, "Name a repository."));
550 };
551 let (repo, member) = match self.visible_repo(path, &a.viewer).await? {
552 Outcome::Ok(found) => found,
553 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
554 };
555 let rows = self
556 .db
557 .prepare(
558 "SELECT p.id AS pull_id, p.number, p.title, p.status, p.agent,
559 count(*) AS entries,
560 sum(CASE WHEN e.kind = 'tool_call' THEN 1 ELSE 0 END) AS tools,
561 min(e.at) AS started_at, max(e.at) AS last_at,
562 (SELECT substr(f.text, 1, 300) FROM session_entries f
563 WHERE f.pull_id = p.id AND f.kind = 'prompt' ORDER BY f.seq LIMIT 1) AS prompt
564 FROM pulls p JOIN session_entries e ON e.pull_id = p.id
565 WHERE p.repo_id = ?1
566 AND (?2 IS NULL OR p.status = ?2)
567 AND (?3 IS NULL OR p.number = ?3)
568 AND (?4 IS NULL OR EXISTS
569 (SELECT 1 FROM agent_runs r WHERE r.pull_id = p.id AND r.kind = ?4))
570 GROUP BY p.id ORDER BY last_at DESC LIMIT 100",
571 )
572 .bind(&[
573 repo.id.as_str().into(),
574 a.outcome.map_or(JsValue::NULL, |status| status.as_str().into()),
575 a.number.map_or(JsValue::NULL, JsValue::from),
576 a.kind.map_or(JsValue::NULL, |kind| kind.as_str().into()),
577 ])?
578 .all()
579 .await?
580 .results::<SessionListRow>()?;
581 let ids: Vec<&str> = rows.iter().map(|row| row.pull_id.as_str()).collect();
582 let runs = if ids.is_empty() {
583 Vec::new()
584 } else {
585 self.db
586 .prepare(
587 "SELECT pull_id, kind, status, cost_usd FROM agent_runs
588 WHERE repo_id = ? AND pull_id IN (SELECT value FROM json_each(?))",
589 )
590 .bind(&[repo.id.as_str().into(), serde_json::to_string(&ids)?.into()])?
591 .all()
592 .await?
593 .results::<RunSumRow>()?
594 };
595 Ok(Outcome::Ok(
596 rows.into_iter()
597 .map(|row| {
598 let mine: Vec<&RunSumRow> = runs.iter().filter(|run| run.pull_id == row.pull_id).collect();
599 let mut kinds: Vec<RunKind> = Vec::new();
600 for run in &mine {
601 if let Some(kind) = RunKind::parse(&run.kind)
602 && !kinds.contains(&kind)
603 {
604 kinds.push(kind);
605 }
606 }
607 let spent: f64 = mine.iter().filter_map(|run| run.cost_usd).sum();
608 SessionSummary {
609 number: row.number,
610 title: row.title,
611 status: row.status,
612 agent: row.agent,
613 entries: row.entries,
614 tools: row.tools,
615 prompt: row.prompt.map(|prompt| one_line(&prompt, 200)),
616 runs: mine.len() as u32,
617 cost_usd: (member && mine.iter().any(|run| run.cost_usd.is_some())).then_some(spent),
618 active: mine
619 .iter()
620 .any(|run| RunStatus::parse(&run.status).is_some_and(RunStatus::is_active)),
621 kinds,
622 started_at: row.started_at,
623 last_at: row.last_at,
624 }
625 })
626 .collect(),
627 ))
628 }
629
630 pub(crate) async fn get_session(&self, a: GetSessionArgs) -> Result<Outcome<SessionView>> {
631 let (repo, pull) = match self.pull_at(&a.repo, a.number, &a.viewer).await? {
632 Outcome::Ok(found) => found,
633 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
634 };
635 let member = is_member(&a.viewer, &repo.namespace);
636 let entries: Vec<SessionEntry> = self
637 .db
638 .prepare(
639 "SELECT seq, kind, text, tool, \"commit\", at FROM session_entries
640 WHERE pull_id = ? ORDER BY seq LIMIT ?",
641 )
642 .bind(&[pull.id.as_str().into(), SESSION_ENTRIES.into()])?
643 .all()
644 .await?
645 .results::<SessionRow>()?
646 .into_iter()
647 .map(SessionEntry::from)
648 .collect();
649 let runs: Vec<AgentRun> = self
650 .db
651 .prepare(format!(
652 "SELECT {COLUMNS}, steps FROM agent_runs WHERE pull_id = ? ORDER BY created_at DESC LIMIT 50"
653 ))
654 .bind(&[pull.id.as_str().into()])?
655 .all()
656 .await?
657 .results::<RunRow>()?
658 .into_iter()
659 .map(|row| row.into_run(member))
660 .collect();
661 let cost_usd = (member && runs.iter().any(|run| run.cost_usd.is_some()))
662 .then(|| runs.iter().filter_map(|run| run.cost_usd).sum());
663 Ok(Outcome::Ok(SessionView {
664 pull,
665 entries,
666 runs,
667 cost_usd,
668 }))
669 }
670}
671
672#[cfg(test)]
673mod tests {
674 use super::*;
675
676 #[test]
A host a sandbox was refused is a note on the run, never what its card says it is doing, and tools in a sandbox with only allowed hosts are told not to send usage reports home677 fn a_refused_host_is_noted_but_never_where_the_run_is() {
678 let steps = |list: &[&str]| list.iter().map(|step| (*step).to_owned()).collect::<Vec<String>>();
679 assert_eq!(current_step(&steps(&["Running npm ci", "Blocked: sparrow.cloudflare.com (not an allowed domain)"])).as_deref(), Some("Running npm ci"));
680 assert_eq!(current_step(&steps(&["Blocked: a.com (not an allowed domain)"])), None);
681 assert_eq!(current_step(&steps(&["Blocked: a.com (not an allowed domain)", "Typecheck"])).as_deref(), Some("Typecheck"));
682 assert_eq!(current_step(&[]), None);
683 }
684
685 #[test]
Agents and memory, checks and conflicts, profiles, slug renames, custom domains686 fn steps_are_one_short_line() {
687 assert_eq!(one_line(" Read\n src/lib.rs ", 40), "Read src/lib.rs");
688 let long = one_line(&"x".repeat(300), 10);
689 assert_eq!(long.chars().count(), 10);
690 assert!(long.ends_with('…'));
691 }
692}

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