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/harness.rs

317 lines12,966 bytesCodeBlame
1//! The agent harness. g1t does not implement its own agent loop: it runs
2//! Claude Code headless and translates its event stream into session
3//! entries.
4
5use std::io::{BufRead, BufReader};
6use std::path::Path;
7use std::process::{Command, Stdio};
8use std::sync::Arc;
9use std::sync::atomic::{AtomicBool, Ordering};
10use std::time::Duration;
11
12use anyhow::{Context, Result, bail};
13use serde_json::Value;
14
15use crate::guard::{self, Halted, Policy};
16use crate::progress::{self, Progress};
17use crate::report::{Entry, Reporter};
18
19const MAX_TURNS: &str = "80";
20
21/// A compact, readable rendering of a tool's input.
22fn describe_input(tool: &str, input: &Value) -> String {
23 let field = |name: &str| input.get(name).and_then(Value::as_str);
24 match tool {
25 "Bash" => field("command").unwrap_or_default().to_owned(),
26 "Read" | "Write" | "Edit" | "MultiEdit" | "NotebookEdit" => {
27 field("file_path").unwrap_or_default().to_owned()
28 }
29 "Glob" | "Grep" => field("pattern").unwrap_or_default().to_owned(),
30 _ => serde_json::to_string(input).unwrap_or_default(),
31 }
32}
33
34/// The text of a tool result, which is either a string or content blocks.
35fn result_text(content: &Value) -> String {
36 match content {
37 Value::String(text) => text.clone(),
38 Value::Array(blocks) => blocks
39 .iter()
40 .filter_map(|block| block.get("text").and_then(Value::as_str))
41 .collect::<Vec<_>>()
42 .join("\n"),
43 _ => String::new(),
44 }
45}
46
47/// Tells g1t what the run cost, so that the workspace it was for can be
48/// charged. The run's own token, given to this sandbox and to nothing
49/// else, is the credential. Does nothing where runs are not billed.
50fn report_cost(cost_usd: f64, turns: u64) {
51 let (Ok(api), Ok(run), Ok(token)) = (
52 std::env::var("G1T_API"),
53 std::env::var("BILLING_RUN"),
54 std::env::var("BILLING_TOKEN"),
55 ) else {
56 return;
57 };
58 let sent = ureq::post(&format!("{api}/runs/{run}/usage")).send_json(serde_json::json!({
59 "token": token,
60 "cost_usd": cost_usd,
61 "turns": turns,
62 }));
63 if let Err(error) = sent {
64 eprintln!("g1t-runner: could not report what the run cost: {error}");
65 }
66}
67
68/// Records one line of Claude Code's `stream-json` output. Returns the final
69/// result when the line is the one that ends the run.
70fn handle_event(
71 event: &Value,
72 reporter: &mut Reporter,
73 progress: &mut Option<Progress>,
74) -> Option<Result<String>> {
75 let blocks = || {
76 event["message"]["content"]
77 .as_array()
78 .map(Vec::as_slice)
79 .unwrap_or_default()
80 };
81 match event["type"].as_str()? {
82 "assistant" => {
83 for block in blocks() {
84 match block["type"].as_str() {
85 Some("text") => {
86 let text = block["text"].as_str().unwrap_or_default().trim();
87 if !text.is_empty() {
88 reporter.record(Entry::new("message", text));
89 if let Some(progress) = progress {
90 progress.step(&format!("Said: {}", text.lines().next().unwrap_or_default()));
91 }
92 }
93 }
94 Some("tool_use") => {
95 let tool = block["name"].as_str().unwrap_or("tool");
96 if let Some(progress) = progress {
97 progress.step(&progress::describe_tool(tool, &block["input"]));
98 }
99 reporter.record(Entry::tool(
100 "tool_call",
101 tool,
102 &describe_input(tool, &block["input"]),
103 ));
104 }
105 _ => {}
106 }
107 }
108 None
109 }
110 "user" => {
111 for block in blocks() {
112 if block["type"] == "tool_result" {
113 let text = result_text(&block["content"]);
114 if !text.trim().is_empty() {
115 reporter.record(Entry::tool("tool_result", "result", text.trim()));
116 }
117 }
118 }
119 None
120 }
121 "result" => {
122 let text = event["result"].as_str().unwrap_or_default().to_owned();
123 // What the run cost, as the harness worked it out, kept with
124 // the session so that spend can be read per pull request.
125 if let Some(cost) = event["total_cost_usd"].as_f64() {
126 let turns = event["num_turns"].as_u64().unwrap_or_default();
127 report_cost(cost, turns);
128 if let Some(progress) = progress {
129 progress.cost(cost, turns);
130 }
131 reporter.record(Entry::new(
132 "note",
133 &format!("This run cost ${cost:.4} over {turns} turns."),
134 ));
135 }
136 // Claude Code stopped at the run's cost cap (`--max-budget-usd`).
137 if event["subtype"] == "error_max_budget_usd" {
138 return Some(Err(anyhow::Error::new(Halted::Budget)));
139 }
140 Some(if event["is_error"].as_bool().unwrap_or(false) {
141 Err(anyhow::anyhow!("the agent reported an error: {text}"))
142 } else {
143 Ok(text)
144 })
145 }
146 _ => None,
147 }
148}
149
150/// Runs Claude Code on `prompt` in `workdir` and returns its closing
151/// summary.
152pub fn run_claude(workdir: &Path, prompt: &str, reporter: &mut Reporter) -> Result<String> {
153 // g1t's own tools, through a token that can do a few things in this one
154 // repository: read its issues, pull requests and merge queue, open an
155 // issue, and comment.
156 let mut tools = Vec::new();
157 if let Ok(token) = std::env::var("G1T_AGENT_TOKEN") {
158 let config = serde_json::json!({
159 "mcpServers": {
160 "g1t": {
161 "type": "http",
162 "url": std::env::var("G1T_MCP").unwrap_or_else(|_| "https://mcp.g1t.sh".to_owned()),
163 "headers": { "Authorization": format!("Bearer {token}") },
164 }
165 }
166 });
167 let path = "/work/g1t-mcp.json";
168 if std::fs::write(path, config.to_string()).is_ok() {
169 tools = vec!["--mcp-config".to_owned(), path.to_owned()];
170 }
171 // People can message the agent while it works: after each tool call
172 // a hook asks g1t for messages and hands any to the agent.
173 if let (Ok(repo), Ok(number)) = (std::env::var("G1T_REPO"), std::env::var("PULL_NUMBER")) {
174 let steer = serde_json::json!({
175 "api": std::env::var("G1T_API").unwrap_or_else(|_| "https://api.g1t.sh".to_owned()),
176 "token": token,
177 "repo": repo,
178 "number": number.parse::<u32>().unwrap_or_default(),
179 });
180 let hooks = serde_json::json!({
181 "hooks": {
182 "PostToolUse": [{
183 "matcher": "*",
184 "hooks": [{ "type": "command", "command": "MODE=steer /usr/local/bin/g1t-runner", "timeout": 15 }],
185 }],
186 "Stop": [{
187 "hooks": [{ "type": "command", "command": "MODE=steer G1T_HOOK=stop /usr/local/bin/g1t-runner", "timeout": 15 }],
188 }],
189 }
190 });
191 let home = std::env::var("HOME").unwrap_or_else(|_| "/home/node".to_owned());
192 let settings = format!("{home}/.claude/settings.json");
193 if std::fs::write(crate::steer::CONFIG, steer.to_string()).is_ok()
194 && std::fs::create_dir_all(format!("{home}/.claude")).is_ok()
195 {
196 let _ = std::fs::write(settings, hooks.to_string());
197 }
198 }
199 }
200 // The run's guardrails (guard.rs): the hook before every tool call, the
201 // permission rules, and the cost cap, which Claude Code enforces itself.
202 let policy = guard::Policy::from_env();
203 if let Some(policy) = &policy {
204 let installed = guard::install(policy, &std::env::var("GUARDRAILS").unwrap_or_default());
205 if let Some(file) = installed.settings_file {
206 tools.extend(["--settings".to_owned(), file]);
207 }
208 if let Some(budget) = policy.budget_usd.filter(|budget| *budget > 0.0) {
209 tools.extend(["--max-budget-usd".to_owned(), format!("{budget:.2}")]);
210 }
211 }
212 // A fork's checkout is anyone's: none of its CLAUDE.md, .claude
213 // settings, hooks, MCP servers or commands are loaded (guard.rs).
214 let trusted = guard::checkout_trusted_from_env();
215 if !trusted {
216 tools.extend(guard::UNTRUSTED_FLAGS.iter().map(|flag| (*flag).to_owned()));
217 reporter.record(Entry::new(
218 "note",
219 "This checkout is not the repository's own branch, so its CLAUDE.md and .claude settings, hooks, MCP servers and commands were not loaded.",
220 ));
221 }
222 let mut denials = guard::Denials::default();
223 // How the run goes, step by step, for people watching it live.
224 let mut progress = Progress::from_env();
225 if let Some(progress) = &mut progress {
226 progress.step("Started the agent");
227 }
228 let mut command = Command::new("claude");
229 if !trusted {
230 command.env(guard::UNTRUSTED_ENV.0, guard::UNTRUSTED_ENV.1);
231 }
232 let mut child = command
233 .current_dir(workdir)
234 .args(&tools)
235 .args([
236 "--print",
237 prompt,
238 "--output-format",
239 "stream-json",
240 "--verbose",
241 "--max-turns",
242 MAX_TURNS,
243 // The sandbox is the permission boundary: it holds one fork and
244 // one short-lived token, and nothing else.
245 "--dangerously-skip-permissions",
246 ])
247 // The agent needs the model key and nothing else of ours.
248 .env_remove("G1T_TOKEN")
249 .env_remove("REVIEW_TOKEN")
250 .env_remove("CHECK_TOKEN")
251 .env_remove("BILLING_TOKEN")
252 .env_remove("PLAN_TOKEN")
253 .env_remove("G1T_AGENT_TOKEN")
254 .env_remove("AGENT_RUN_TOKEN")
255 .stdin(Stdio::null())
256 .stdout(Stdio::piped())
257 .stderr(Stdio::inherit())
258 .spawn()
259 .context("could not start Claude Code")?;
260
261 // The time cap: the agent is stopped when it passes, and the run with it.
262 let timed_out = Arc::new(AtomicBool::new(false));
263 let finished = Arc::new(AtomicBool::new(false));
264 if let Some(minutes) = policy.as_ref().and_then(|policy| policy.minutes) {
265 let (pid, timed_out, finished) = (child.id(), timed_out.clone(), finished.clone());
266 std::thread::spawn(move || {
267 std::thread::sleep(Duration::from_secs(u64::from(minutes) * 60));
268 if !finished.load(Ordering::SeqCst) {
269 timed_out.store(true, Ordering::SeqCst);
270 let _ = Command::new("kill").arg(pid.to_string()).status();
271 }
272 });
273 }
274
275 let stdout = child.stdout.take().context("no output from Claude Code")?;
276 let mut outcome = None;
277 for line in BufReader::new(stdout).lines() {
278 let line = line?;
279 let Ok(event) = serde_json::from_str::<Value>(&line) else {
280 continue;
281 };
282 if let Some(result) = handle_event(&event, reporter, &mut progress) {
283 outcome = Some(result);
284 }
285 // What the guardrails refused, as steps people can see.
286 for denied in denials.take() {
287 reporter.record(Entry::new("note", &denied));
288 if let Some(progress) = &mut progress {
289 progress.step(&denied);
290 }
291 }
292 }
293 let status = child.wait()?;
294 finished.store(true, Ordering::SeqCst);
295 if timed_out.load(Ordering::SeqCst) {
296 outcome = Some(Err(anyhow::Error::new(Halted::Time)));
297 }
298 if let Some(progress) = &mut progress {
299 progress.flush();
300 if let Some(halted) = outcome.as_ref().and_then(|o| o.as_ref().err()).and_then(|e| e.downcast_ref::<Halted>()) {
301 let message = match (halted, &policy) {
302 (Halted::Budget, Some(Policy { budget_usd: Some(budget), .. })) => {
303 format!("Stopped: it reached its cost cap of ${budget:.2}.")
304 }
305 (Halted::Time, Some(Policy { minutes: Some(minutes), .. })) => {
306 format!("Stopped: it reached its time cap of {minutes} minutes.")
307 }
308 _ => format!("Stopped: {halted}."),
309 };
310 progress.halt(halted.reason(), &message);
311 }
312 }
313 match outcome {
314 Some(result) => result,
315 None => bail!("Claude Code exited ({status}) without a result"),
316 }
317}