g1t/crates/runner/src/progress.rs
| 1 | //! Tells g1t what the agent is doing, step by step, so people can watch a |
| 2 | //! run live: each tool call and each thing the agent says, as one short |
| 3 | //! line. Sent at most every few seconds, with the run's own token, and |
| 4 | //! never allowed to slow down or fail the run. |
| 5 | //! |
| 6 | //! The runner service gives the sandbox `AGENT_RUN` and `AGENT_RUN_TOKEN` |
| 7 | //! when it records the run. Without them nothing is sent. |
| 8 | |
| 9 | use std::time::{Duration, Instant}; |
| 10 | |
| 11 | use serde_json::{Value, json}; |
| 12 | |
| 13 | /// Steps are sent at most this often. |
| 14 | const INTERVAL: Duration = Duration::from_secs(3); |
| 15 | /// The most steps one report carries; a burst keeps its end. |
| 16 | const MAX_PENDING: usize = 20; |
| 17 | const MAX_STEP_CHARS: usize = 200; |
| 18 | |
| 19 | pub struct Progress { |
| 20 | url: String, |
| 21 | token: String, |
| 22 | pending: Vec<String>, |
| 23 | last: Instant, |
| 24 | /// Values that must never be sent: the credentials this process has. |
| 25 | secrets: Vec<String>, |
| 26 | /// Set once g1t says the run was stopped. |
| 27 | stopped: bool, |
| 28 | } |
| 29 | |
| 30 | /// One line, at most `MAX_STEP_CHARS`. |
| 31 | pub fn one_line(text: &str) -> String { |
| 32 | let line = text.split_whitespace().collect::<Vec<_>>().join(" "); |
| 33 | if line.chars().count() <= MAX_STEP_CHARS { |
| 34 | return line; |
| 35 | } |
| 36 | let mut short: String = line.chars().take(MAX_STEP_CHARS - 1).collect(); |
| 37 | short.push('…'); |
| 38 | short |
| 39 | } |
| 40 | |
| 41 | /// A tool call as a step: what it did, to what. |
| 42 | pub fn describe_tool(tool: &str, input: &Value) -> String { |
| 43 | let field = |name: &str| input.get(name).and_then(Value::as_str).unwrap_or_default(); |
| 44 | let file = |name: &str| { |
| 45 | let path = field(name); |
| 46 | path.strip_prefix("/work/repo/").unwrap_or(path).to_owned() |
| 47 | }; |
| 48 | let line = match tool { |
| 49 | "Bash" => format!("Ran {}", field("command")), |
| 50 | "Read" => format!("Read {}", file("file_path")), |
| 51 | "Write" => format!("Wrote {}", file("file_path")), |
| 52 | "Edit" | "MultiEdit" => format!("Edited {}", file("file_path")), |
| 53 | "NotebookEdit" => format!("Edited {}", file("notebook_path")), |
| 54 | "Glob" => format!("Looked for {}", field("pattern")), |
| 55 | "Grep" => format!("Searched for {}", field("pattern")), |
| 56 | "WebFetch" => format!("Fetched {}", field("url")), |
| 57 | "WebSearch" => format!("Searched the web for {}", field("query")), |
| 58 | "TodoWrite" => "Updated its plan".to_owned(), |
| 59 | "Task" | "Agent" => format!("Asked a helper: {}", field("description")), |
| 60 | other => match other.strip_prefix("mcp__g1t__") { |
| 61 | Some(operation) => format!("Used g1t: {}", operation.replace('_', " ")), |
| 62 | None => format!("Used {other}"), |
| 63 | }, |
| 64 | }; |
| 65 | one_line(&line) |
| 66 | } |
| 67 | |
| 68 | impl Progress { |
| 69 | pub fn from_env() -> Option<Progress> { |
| 70 | let run = std::env::var("AGENT_RUN").ok().filter(|run| !run.is_empty())?; |
| 71 | let token = std::env::var("AGENT_RUN_TOKEN").ok().filter(|token| !token.is_empty())?; |
| 72 | let api = std::env::var("G1T_API").unwrap_or_else(|_| "https://api.g1t.sh".to_owned()); |
| 73 | let secrets = [ |
| 74 | "G1T_TOKEN", |
| 75 | "G1T_AGENT_TOKEN", |
| 76 | "ANTHROPIC_API_KEY", |
| 77 | "AI_GATEWAY_TOKEN", |
| 78 | "BILLING_TOKEN", |
| 79 | "AGENT_RUN_TOKEN", |
| 80 | ] |
| 81 | .into_iter() |
| 82 | .filter_map(|name| std::env::var(name).ok()) |
| 83 | .filter(|secret| secret.len() >= 8) |
| 84 | .collect(); |
| 85 | Some(Progress { |
| 86 | url: format!("{}/agent-runs/{run}/report", api.trim_end_matches('/')), |
| 87 | token, |
| 88 | pending: Vec::new(), |
| 89 | // So the first step goes at once. |
| 90 | last: Instant::now() - INTERVAL, |
| 91 | secrets, |
| 92 | stopped: false, |
| 93 | }) |
| 94 | } |
| 95 | |
| 96 | fn clean(&self, text: &str) -> String { |
| 97 | self.secrets |
| 98 | .iter() |
| 99 | .fold(text.to_owned(), |text, secret| text.replace(secret, "[redacted]")) |
| 100 | } |
| 101 | |
| 102 | fn send(&mut self, body: Value) { |
| 103 | let mut body = body; |
| 104 | body["token"] = json!(self.token); |
| 105 | let sent = ureq::post(&self.url) |
| 106 | .timeout(Duration::from_secs(10)) |
| 107 | .send_json(body); |
| 108 | match sent { |
| 109 | Ok(response) => { |
| 110 | let answer: Value = response.into_json().unwrap_or(Value::Null); |
| 111 | if answer["status"] == "stopped" { |
| 112 | self.stopped = true; |
| 113 | } |
| 114 | } |
| 115 | Err(error) => eprintln!("g1t-runner: could not report the run's progress: {error}"), |
| 116 | } |
| 117 | } |
| 118 | |
| 119 | /// Records a step, sending what has gathered if it is time. |
| 120 | pub fn step(&mut self, text: &str) { |
| 121 | crate::abuse::touch(); |
| 122 | let line = one_line(&self.clean(text)); |
| 123 | if line.is_empty() { |
| 124 | return; |
| 125 | } |
| 126 | self.pending.push(line); |
| 127 | if self.pending.len() > MAX_PENDING { |
| 128 | self.pending.remove(0); |
| 129 | } |
| 130 | if self.last.elapsed() >= INTERVAL { |
| 131 | self.flush(); |
| 132 | } |
| 133 | } |
| 134 | |
| 135 | /// Sends the steps gathered so far. |
| 136 | pub fn flush(&mut self) { |
| 137 | self.last = Instant::now(); |
| 138 | if self.pending.is_empty() { |
| 139 | return; |
| 140 | } |
| 141 | let steps = std::mem::take(&mut self.pending); |
| 142 | self.send(json!({ "steps": steps })); |
| 143 | } |
| 144 | |
| 145 | /// What the run cost, as the harness worked it out. |
| 146 | pub fn cost(&mut self, cost_usd: f64, turns: u64) { |
| 147 | self.flush(); |
| 148 | self.send(json!({ "costUsd": cost_usd, "turns": turns })); |
| 149 | } |
| 150 | |
| 151 | /// Ends the run as stopped because it reached a cap of its guardrails |
| 152 | /// (guard.rs): `budget` or `time`. |
| 153 | pub fn halt(&mut self, reason: &str, message: &str) { |
| 154 | self.flush(); |
| 155 | self.send(json!({ "halt": reason, "error": message })); |
| 156 | } |
| 157 | |
| 158 | /// Whether a person stopped the run. |
| 159 | #[allow(dead_code)] |
| 160 | pub fn stopped(&self) -> bool { |
| 161 | self.stopped |
| 162 | } |
| 163 | } |
| 164 | |
| 165 | #[cfg(test)] |
| 166 | mod tests { |
| 167 | use super::*; |
| 168 | |
| 169 | #[test] |
| 170 | fn tool_calls_read_as_steps() { |
| 171 | assert_eq!( |
| 172 | describe_tool("Edit", &json!({ "file_path": "/work/repo/src/lib.rs" })), |
| 173 | "Edited src/lib.rs" |
| 174 | ); |
| 175 | assert_eq!(describe_tool("Bash", &json!({ "command": "cargo test\n -q" })), "Ran cargo test -q"); |
| 176 | assert_eq!(describe_tool("mcp__g1t__create_issue", &json!({})), "Used g1t: create issue"); |
| 177 | } |
| 178 | |
| 179 | #[test] |
| 180 | fn steps_are_one_short_line() { |
| 181 | let long = one_line(&"word ".repeat(100)); |
| 182 | assert_eq!(long.chars().count(), MAX_STEP_CHARS); |
| 183 | assert!(long.ends_with('…')); |
| 184 | } |
| 185 | } |