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/crates/runner/src/progress.rs

185 lines6,390 bytesCodeBlame
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
9use std::time::{Duration, Instant};
10
11use serde_json::{Value, json};
12
13/// Steps are sent at most this often.
14const INTERVAL: Duration = Duration::from_secs(3);
15/// The most steps one report carries; a burst keeps its end.
16const MAX_PENDING: usize = 20;
17const MAX_STEP_CHARS: usize = 200;
18
19pub 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`.
31pub 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.
42pub 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
68impl 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)]
166mod 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}