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

318 lines13,013 bytesCodeBlame

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.

Account dropdown, llms.txt onboarding, hosted agent runner (not yet deployed)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};
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API8use std::sync::Arc;
9use std::sync::atomic::{AtomicBool, Ordering};
10use std::time::Duration;
Account dropdown, llms.txt onboarding, hosted agent runner (not yet deployed)11
12use anyhow::{Context, Result, bail};
13use serde_json::Value;
14
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API15use crate::guard::{self, Halted, Policy};
Agents and memory, checks and conflicts, profiles, slug renames, custom domains16use crate::progress::{self, Progress};
Account dropdown, llms.txt onboarding, hosted agent runner (not yet deployed)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
Agents as a team: lifecycle, merge queue, billing and a new shell47/// 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
Account dropdown, llms.txt onboarding, hosted agent runner (not yet deployed)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.
Agents and memory, checks and conflicts, profiles, slug renames, custom domains70fn handle_event(
71 event: &Value,
72 reporter: &mut Reporter,
73 progress: &mut Option<Progress>,
74) -> Option<Result<String>> {
Account dropdown, llms.txt onboarding, hosted agent runner (not yet deployed)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));
Agents and memory, checks and conflicts, profiles, slug renames, custom domains89 if let Some(progress) = progress {
90 progress.step(&format!("Said: {}", text.lines().next().unwrap_or_default()));
91 }
Account dropdown, llms.txt onboarding, hosted agent runner (not yet deployed)92 }
93 }
94 Some("tool_use") => {
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look95 crate::abuse::touch();
Account dropdown, llms.txt onboarding, hosted agent runner (not yet deployed)96 let tool = block["name"].as_str().unwrap_or("tool");
Agents and memory, checks and conflicts, profiles, slug renames, custom domains97 if let Some(progress) = progress {
98 progress.step(&progress::describe_tool(tool, &block["input"]));
99 }
Account dropdown, llms.txt onboarding, hosted agent runner (not yet deployed)100 reporter.record(Entry::tool(
101 "tool_call",
102 tool,
103 &describe_input(tool, &block["input"]),
104 ));
105 }
106 _ => {}
107 }
108 }
109 None
110 }
111 "user" => {
112 for block in blocks() {
113 if block["type"] == "tool_result" {
114 let text = result_text(&block["content"]);
115 if !text.trim().is_empty() {
116 reporter.record(Entry::tool("tool_result", "result", text.trim()));
117 }
118 }
119 }
120 None
121 }
122 "result" => {
123 let text = event["result"].as_str().unwrap_or_default().to_owned();
Agents as a team: lifecycle, merge queue, billing and a new shell124 // What the run cost, as the harness worked it out, kept with
125 // the session so that spend can be read per pull request.
126 if let Some(cost) = event["total_cost_usd"].as_f64() {
127 let turns = event["num_turns"].as_u64().unwrap_or_default();
128 report_cost(cost, turns);
Agents and memory, checks and conflicts, profiles, slug renames, custom domains129 if let Some(progress) = progress {
130 progress.cost(cost, turns);
131 }
Agents as a team: lifecycle, merge queue, billing and a new shell132 reporter.record(Entry::new(
133 "note",
134 &format!("This run cost ${cost:.4} over {turns} turns."),
135 ));
136 }
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API137 // Claude Code stopped at the run's cost cap (`--max-budget-usd`).
138 if event["subtype"] == "error_max_budget_usd" {
139 return Some(Err(anyhow::Error::new(Halted::Budget)));
140 }
Account dropdown, llms.txt onboarding, hosted agent runner (not yet deployed)141 Some(if event["is_error"].as_bool().unwrap_or(false) {
142 Err(anyhow::anyhow!("the agent reported an error: {text}"))
143 } else {
144 Ok(text)
145 })
146 }
147 _ => None,
148 }
149}
150
151/// Runs Claude Code on `prompt` in `workdir` and returns its closing
152/// summary.
153pub fn run_claude(workdir: &Path, prompt: &str, reporter: &mut Reporter) -> Result<String> {
Agents as a team: lifecycle, merge queue, billing and a new shell154 // g1t's own tools, through a token that can do a few things in this one
155 // repository: read its issues, pull requests and merge queue, open an
156 // issue, and comment.
157 let mut tools = Vec::new();
158 if let Ok(token) = std::env::var("G1T_AGENT_TOKEN") {
159 let config = serde_json::json!({
160 "mcpServers": {
161 "g1t": {
162 "type": "http",
163 "url": std::env::var("G1T_MCP").unwrap_or_else(|_| "https://mcp.g1t.sh".to_owned()),
164 "headers": { "Authorization": format!("Bearer {token}") },
165 }
166 }
167 });
168 let path = "/work/g1t-mcp.json";
169 if std::fs::write(path, config.to_string()).is_ok() {
170 tools = vec!["--mcp-config".to_owned(), path.to_owned()];
171 }
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request172 // People can message the agent while it works: after each tool call
173 // a hook asks g1t for messages and hands any to the agent.
174 if let (Ok(repo), Ok(number)) = (std::env::var("G1T_REPO"), std::env::var("PULL_NUMBER")) {
175 let steer = serde_json::json!({
176 "api": std::env::var("G1T_API").unwrap_or_else(|_| "https://api.g1t.sh".to_owned()),
177 "token": token,
178 "repo": repo,
179 "number": number.parse::<u32>().unwrap_or_default(),
180 });
181 let hooks = serde_json::json!({
182 "hooks": {
183 "PostToolUse": [{
184 "matcher": "*",
185 "hooks": [{ "type": "command", "command": "MODE=steer /usr/local/bin/g1t-runner", "timeout": 15 }],
186 }],
Messages reach the agent even as it finishes187 "Stop": [{
188 "hooks": [{ "type": "command", "command": "MODE=steer G1T_HOOK=stop /usr/local/bin/g1t-runner", "timeout": 15 }],
189 }],
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request190 }
191 });
192 let home = std::env::var("HOME").unwrap_or_else(|_| "/home/node".to_owned());
193 let settings = format!("{home}/.claude/settings.json");
194 if std::fs::write(crate::steer::CONFIG, steer.to_string()).is_ok()
195 && std::fs::create_dir_all(format!("{home}/.claude")).is_ok()
196 {
197 let _ = std::fs::write(settings, hooks.to_string());
198 }
199 }
Agents as a team: lifecycle, merge queue, billing and a new shell200 }
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API201 // The run's guardrails (guard.rs): the hook before every tool call, the
202 // permission rules, and the cost cap, which Claude Code enforces itself.
203 let policy = guard::Policy::from_env();
204 if let Some(policy) = &policy {
205 let installed = guard::install(policy, &std::env::var("GUARDRAILS").unwrap_or_default());
206 if let Some(file) = installed.settings_file {
207 tools.extend(["--settings".to_owned(), file]);
208 }
209 if let Some(budget) = policy.budget_usd.filter(|budget| *budget > 0.0) {
210 tools.extend(["--max-budget-usd".to_owned(), format!("{budget:.2}")]);
211 }
212 }
213 // A fork's checkout is anyone's: none of its CLAUDE.md, .claude
214 // settings, hooks, MCP servers or commands are loaded (guard.rs).
215 let trusted = guard::checkout_trusted_from_env();
216 if !trusted {
217 tools.extend(guard::UNTRUSTED_FLAGS.iter().map(|flag| (*flag).to_owned()));
218 reporter.record(Entry::new(
219 "note",
220 "This checkout is not the repository's own branch, so its CLAUDE.md and .claude settings, hooks, MCP servers and commands were not loaded.",
221 ));
222 }
223 let mut denials = guard::Denials::default();
Agents and memory, checks and conflicts, profiles, slug renames, custom domains224 // How the run goes, step by step, for people watching it live.
225 let mut progress = Progress::from_env();
226 if let Some(progress) = &mut progress {
227 progress.step("Started the agent");
228 }
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API229 let mut command = Command::new("claude");
230 if !trusted {
231 command.env(guard::UNTRUSTED_ENV.0, guard::UNTRUSTED_ENV.1);
232 }
233 let mut child = command
Account dropdown, llms.txt onboarding, hosted agent runner (not yet deployed)234 .current_dir(workdir)
Agents as a team: lifecycle, merge queue, billing and a new shell235 .args(&tools)
Account dropdown, llms.txt onboarding, hosted agent runner (not yet deployed)236 .args([
237 "--print",
238 prompt,
239 "--output-format",
240 "stream-json",
241 "--verbose",
242 "--max-turns",
243 MAX_TURNS,
244 // The sandbox is the permission boundary: it holds one fork and
245 // one short-lived token, and nothing else.
246 "--dangerously-skip-permissions",
247 ])
248 // The agent needs the model key and nothing else of ours.
249 .env_remove("G1T_TOKEN")
Agents as a team: lifecycle, merge queue, billing and a new shell250 .env_remove("REVIEW_TOKEN")
251 .env_remove("CHECK_TOKEN")
252 .env_remove("BILLING_TOKEN")
253 .env_remove("PLAN_TOKEN")
254 .env_remove("G1T_AGENT_TOKEN")
Agents and memory, checks and conflicts, profiles, slug renames, custom domains255 .env_remove("AGENT_RUN_TOKEN")
Account dropdown, llms.txt onboarding, hosted agent runner (not yet deployed)256 .stdin(Stdio::null())
257 .stdout(Stdio::piped())
258 .stderr(Stdio::inherit())
259 .spawn()
260 .context("could not start Claude Code")?;
261
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API262 // The time cap: the agent is stopped when it passes, and the run with it.
263 let timed_out = Arc::new(AtomicBool::new(false));
264 let finished = Arc::new(AtomicBool::new(false));
265 if let Some(minutes) = policy.as_ref().and_then(|policy| policy.minutes) {
266 let (pid, timed_out, finished) = (child.id(), timed_out.clone(), finished.clone());
267 std::thread::spawn(move || {
268 std::thread::sleep(Duration::from_secs(u64::from(minutes) * 60));
269 if !finished.load(Ordering::SeqCst) {
270 timed_out.store(true, Ordering::SeqCst);
271 let _ = Command::new("kill").arg(pid.to_string()).status();
272 }
273 });
274 }
275
Account dropdown, llms.txt onboarding, hosted agent runner (not yet deployed)276 let stdout = child.stdout.take().context("no output from Claude Code")?;
277 let mut outcome = None;
278 for line in BufReader::new(stdout).lines() {
279 let line = line?;
280 let Ok(event) = serde_json::from_str::<Value>(&line) else {
281 continue;
282 };
Agents and memory, checks and conflicts, profiles, slug renames, custom domains283 if let Some(result) = handle_event(&event, reporter, &mut progress) {
Account dropdown, llms.txt onboarding, hosted agent runner (not yet deployed)284 outcome = Some(result);
285 }
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API286 // What the guardrails refused, as steps people can see.
287 for denied in denials.take() {
288 reporter.record(Entry::new("note", &denied));
289 if let Some(progress) = &mut progress {
290 progress.step(&denied);
291 }
292 }
Account dropdown, llms.txt onboarding, hosted agent runner (not yet deployed)293 }
294 let status = child.wait()?;
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API295 finished.store(true, Ordering::SeqCst);
296 if timed_out.load(Ordering::SeqCst) {
297 outcome = Some(Err(anyhow::Error::new(Halted::Time)));
298 }
Agents and memory, checks and conflicts, profiles, slug renames, custom domains299 if let Some(progress) = &mut progress {
300 progress.flush();
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API301 if let Some(halted) = outcome.as_ref().and_then(|o| o.as_ref().err()).and_then(|e| e.downcast_ref::<Halted>()) {
302 let message = match (halted, &policy) {
303 (Halted::Budget, Some(Policy { budget_usd: Some(budget), .. })) => {
304 format!("Stopped: it reached its cost cap of ${budget:.2}.")
305 }
306 (Halted::Time, Some(Policy { minutes: Some(minutes), .. })) => {
307 format!("Stopped: it reached its time cap of {minutes} minutes.")
308 }
309 _ => format!("Stopped: {halted}."),
310 };
311 progress.halt(halted.reason(), &message);
312 }
Agents and memory, checks and conflicts, profiles, slug renames, custom domains313 }
Account dropdown, llms.txt onboarding, hosted agent runner (not yet deployed)314 match outcome {
315 Some(result) => result,
316 None => bail!("Claude Code exited ({status}) without a result"),
317 }
318}