g1t/crates/runner/src/progress.rs
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 domains | 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 | let line = one_line(&self.clean(text)); | |
| 122 | if line.is_empty() { | |
| 123 | return; | |
| 124 | } | |
| 125 | self.pending.push(line); | |
| 126 | if self.pending.len() > MAX_PENDING { | |
| 127 | self.pending.remove(0); | |
| 128 | } | |
| 129 | if self.last.elapsed() >= INTERVAL { | |
| 130 | self.flush(); | |
| 131 | } | |
| 132 | } | |
| 133 | ||
| 134 | /// Sends the steps gathered so far. | |
| 135 | pub fn flush(&mut self) { | |
| 136 | self.last = Instant::now(); | |
| 137 | if self.pending.is_empty() { | |
| 138 | return; | |
| 139 | } | |
| 140 | let steps = std::mem::take(&mut self.pending); | |
| 141 | self.send(json!({ "steps": steps })); | |
| 142 | } | |
| 143 | ||
| 144 | /// What the run cost, as the harness worked it out. | |
| 145 | pub fn cost(&mut self, cost_usd: f64, turns: u64) { | |
| 146 | self.flush(); | |
| 147 | self.send(json!({ "costUsd": cost_usd, "turns": turns })); | |
| 148 | } | |
| 149 | ||
| 150 | /// Whether a person stopped the run. | |
| 151 | #[allow(dead_code)] | |
| 152 | pub fn stopped(&self) -> bool { | |
| 153 | self.stopped | |
| 154 | } | |
| 155 | } | |
| 156 | ||
| 157 | #[cfg(test)] | |
| 158 | mod tests { | |
| 159 | use super::*; | |
| 160 | ||
| 161 | #[test] | |
| 162 | fn tool_calls_read_as_steps() { | |
| 163 | assert_eq!( | |
| 164 | describe_tool("Edit", &json!({ "file_path": "/work/repo/src/lib.rs" })), | |
| 165 | "Edited src/lib.rs" | |
| 166 | ); | |
| 167 | assert_eq!(describe_tool("Bash", &json!({ "command": "cargo test\n -q" })), "Ran cargo test -q"); | |
| 168 | assert_eq!(describe_tool("mcp__g1t__create_issue", &json!({})), "Used g1t: create issue"); | |
| 169 | } | |
| 170 | ||
| 171 | #[test] | |
| 172 | fn steps_are_one_short_line() { | |
| 173 | let long = one_line(&"word ".repeat(100)); | |
| 174 | assert_eq!(long.chars().count(), MAX_STEP_CHARS); | |
| 175 | assert!(long.ends_with('…')); | |
| 176 | } | |
| 177 | } |