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

222 lines8,562 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};
8
9use anyhow::{Context, Result, bail};
10use serde_json::Value;
11
12use crate::report::{Entry, Reporter};
13
14const MAX_TURNS: &str = "80";
15
16/// A compact, readable rendering of a tool's input.
17fn describe_input(tool: &str, input: &Value) -> String {
18 let field = |name: &str| input.get(name).and_then(Value::as_str);
19 match tool {
20 "Bash" => field("command").unwrap_or_default().to_owned(),
21 "Read" | "Write" | "Edit" | "MultiEdit" | "NotebookEdit" => {
22 field("file_path").unwrap_or_default().to_owned()
23 }
24 "Glob" | "Grep" => field("pattern").unwrap_or_default().to_owned(),
25 _ => serde_json::to_string(input).unwrap_or_default(),
26 }
27}
28
29/// The text of a tool result, which is either a string or content blocks.
30fn result_text(content: &Value) -> String {
31 match content {
32 Value::String(text) => text.clone(),
33 Value::Array(blocks) => blocks
34 .iter()
35 .filter_map(|block| block.get("text").and_then(Value::as_str))
36 .collect::<Vec<_>>()
37 .join("\n"),
38 _ => String::new(),
39 }
40}
41
42/// Tells g1t what the run cost, so that the workspace it was for can be
43/// charged. The run's own token, given to this sandbox and to nothing
44/// else, is the credential. Does nothing where runs are not billed.
45fn report_cost(cost_usd: f64, turns: u64) {
46 let (Ok(api), Ok(run), Ok(token)) = (
47 std::env::var("G1T_API"),
48 std::env::var("BILLING_RUN"),
49 std::env::var("BILLING_TOKEN"),
50 ) else {
51 return;
52 };
53 let sent = ureq::post(&format!("{api}/runs/{run}/usage")).send_json(serde_json::json!({
54 "token": token,
55 "cost_usd": cost_usd,
56 "turns": turns,
57 }));
58 if let Err(error) = sent {
59 eprintln!("g1t-runner: could not report what the run cost: {error}");
60 }
61}
62
63/// Records one line of Claude Code's `stream-json` output. Returns the final
64/// result when the line is the one that ends the run.
65fn handle_event(event: &Value, reporter: &mut Reporter) -> Option<Result<String>> {
66 let blocks = || {
67 event["message"]["content"]
68 .as_array()
69 .map(Vec::as_slice)
70 .unwrap_or_default()
71 };
72 match event["type"].as_str()? {
73 "assistant" => {
74 for block in blocks() {
75 match block["type"].as_str() {
76 Some("text") => {
77 let text = block["text"].as_str().unwrap_or_default().trim();
78 if !text.is_empty() {
79 reporter.record(Entry::new("message", text));
80 }
81 }
82 Some("tool_use") => {
83 let tool = block["name"].as_str().unwrap_or("tool");
84 reporter.record(Entry::tool(
85 "tool_call",
86 tool,
87 &describe_input(tool, &block["input"]),
88 ));
89 }
90 _ => {}
91 }
92 }
93 None
94 }
95 "user" => {
96 for block in blocks() {
97 if block["type"] == "tool_result" {
98 let text = result_text(&block["content"]);
99 if !text.trim().is_empty() {
100 reporter.record(Entry::tool("tool_result", "result", text.trim()));
101 }
102 }
103 }
104 None
105 }
106 "result" => {
107 let text = event["result"].as_str().unwrap_or_default().to_owned();
108 // What the run cost, as the harness worked it out, kept with
109 // the session so that spend can be read per pull request.
110 if let Some(cost) = event["total_cost_usd"].as_f64() {
111 let turns = event["num_turns"].as_u64().unwrap_or_default();
112 report_cost(cost, turns);
113 reporter.record(Entry::new(
114 "note",
115 &format!("This run cost ${cost:.4} over {turns} turns."),
116 ));
117 }
118 Some(if event["is_error"].as_bool().unwrap_or(false) {
119 Err(anyhow::anyhow!("the agent reported an error: {text}"))
120 } else {
121 Ok(text)
122 })
123 }
124 _ => None,
125 }
126}
127
128/// Runs Claude Code on `prompt` in `workdir` and returns its closing
129/// summary.
130pub fn run_claude(workdir: &Path, prompt: &str, reporter: &mut Reporter) -> Result<String> {
131 // g1t's own tools, through a token that can do a few things in this one
132 // repository: read its issues, pull requests and merge queue, open an
133 // issue, and comment.
134 let mut tools = Vec::new();
135 if let Ok(token) = std::env::var("G1T_AGENT_TOKEN") {
136 let config = serde_json::json!({
137 "mcpServers": {
138 "g1t": {
139 "type": "http",
140 "url": std::env::var("G1T_MCP").unwrap_or_else(|_| "https://mcp.g1t.sh".to_owned()),
141 "headers": { "Authorization": format!("Bearer {token}") },
142 }
143 }
144 });
145 let path = "/work/g1t-mcp.json";
146 if std::fs::write(path, config.to_string()).is_ok() {
147 tools = vec!["--mcp-config".to_owned(), path.to_owned()];
148 }
149 // People can message the agent while it works: after each tool call
150 // a hook asks g1t for messages and hands any to the agent.
151 if let (Ok(repo), Ok(number)) = (std::env::var("G1T_REPO"), std::env::var("PULL_NUMBER")) {
152 let steer = serde_json::json!({
153 "api": std::env::var("G1T_API").unwrap_or_else(|_| "https://api.g1t.sh".to_owned()),
154 "token": token,
155 "repo": repo,
156 "number": number.parse::<u32>().unwrap_or_default(),
157 });
158 let hooks = serde_json::json!({
159 "hooks": {
160 "PostToolUse": [{
161 "matcher": "*",
162 "hooks": [{ "type": "command", "command": "MODE=steer /usr/local/bin/g1t-runner", "timeout": 15 }],
163 }],
164 "Stop": [{
165 "hooks": [{ "type": "command", "command": "MODE=steer G1T_HOOK=stop /usr/local/bin/g1t-runner", "timeout": 15 }],
166 }],
167 }
168 });
169 let home = std::env::var("HOME").unwrap_or_else(|_| "/home/node".to_owned());
170 let settings = format!("{home}/.claude/settings.json");
171 if std::fs::write(crate::steer::CONFIG, steer.to_string()).is_ok()
172 && std::fs::create_dir_all(format!("{home}/.claude")).is_ok()
173 {
174 let _ = std::fs::write(settings, hooks.to_string());
175 }
176 }
177 }
178 let mut child = Command::new("claude")
179 .current_dir(workdir)
180 .args(&tools)
181 .args([
182 "--print",
183 prompt,
184 "--output-format",
185 "stream-json",
186 "--verbose",
187 "--max-turns",
188 MAX_TURNS,
189 // The sandbox is the permission boundary: it holds one fork and
190 // one short-lived token, and nothing else.
191 "--dangerously-skip-permissions",
192 ])
193 // The agent needs the model key and nothing else of ours.
194 .env_remove("G1T_TOKEN")
195 .env_remove("REVIEW_TOKEN")
196 .env_remove("CHECK_TOKEN")
197 .env_remove("BILLING_TOKEN")
198 .env_remove("PLAN_TOKEN")
199 .env_remove("G1T_AGENT_TOKEN")
200 .stdin(Stdio::null())
201 .stdout(Stdio::piped())
202 .stderr(Stdio::inherit())
203 .spawn()
204 .context("could not start Claude Code")?;
205
206 let stdout = child.stdout.take().context("no output from Claude Code")?;
207 let mut outcome = None;
208 for line in BufReader::new(stdout).lines() {
209 let line = line?;
210 let Ok(event) = serde_json::from_str::<Value>(&line) else {
211 continue;
212 };
213 if let Some(result) = handle_event(&event, reporter) {
214 outcome = Some(result);
215 }
216 }
217 let status = child.wait()?;
218 match outcome {
219 Some(result) => result,
220 None => bail!("Claude Code exited ({status}) without a result"),
221 }
222}