Skip to content

g1t/crates/runner/src/harness.rs

362 lines14,633 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
Merge branch 'model-routing'47/// The tokens a run used, from Claude Code's closing `usage`: input,
48/// output, and prompt-cache reads and writes. On a workspace's own model
49/// key this is what g1t charges its agent rate on (the model proxy's own
50/// count, when it has more, wins), so it is reported with the cost.
51fn run_tokens(event: &Value) -> Value {
52 let usage = &event["usage"];
53 let count = |field: &str| usage[field].as_u64().unwrap_or_default();
54 serde_json::json!({
55 "input": count("input_tokens"),
56 "output": count("output_tokens"),
57 "cache_read": count("cache_read_input_tokens"),
58 "cache_write": count("cache_creation_input_tokens"),
59 })
60}
61
Agents as a team: lifecycle, merge queue, billing and a new shell62/// Tells g1t what the run cost, so that the workspace it was for can be
63/// charged. The run's own token, given to this sandbox and to nothing
64/// else, is the credential. Does nothing where runs are not billed.
Merge branch 'model-routing'65fn report_cost(cost_usd: f64, turns: u64, tokens: Value) {
Agents as a team: lifecycle, merge queue, billing and a new shell66 let (Ok(api), Ok(run), Ok(token)) = (
67 std::env::var("G1T_API"),
68 std::env::var("BILLING_RUN"),
69 std::env::var("BILLING_TOKEN"),
70 ) else {
71 return;
72 };
73 let sent = ureq::post(&format!("{api}/runs/{run}/usage")).send_json(serde_json::json!({
74 "token": token,
75 "cost_usd": cost_usd,
76 "turns": turns,
Merge branch 'model-routing'77 "tokens": tokens,
Agents as a team: lifecycle, merge queue, billing and a new shell78 }));
79 if let Err(error) = sent {
80 eprintln!("g1t-runner: could not report what the run cost: {error}");
81 }
82}
83
Account dropdown, llms.txt onboarding, hosted agent runner (not yet deployed)84/// Records one line of Claude Code's `stream-json` output. Returns the final
85/// result when the line is the one that ends the run.
Agents and memory, checks and conflicts, profiles, slug renames, custom domains86fn handle_event(
87 event: &Value,
88 reporter: &mut Reporter,
89 progress: &mut Option<Progress>,
90) -> Option<Result<String>> {
Account dropdown, llms.txt onboarding, hosted agent runner (not yet deployed)91 let blocks = || {
92 event["message"]["content"]
93 .as_array()
94 .map(Vec::as_slice)
95 .unwrap_or_default()
96 };
97 match event["type"].as_str()? {
98 "assistant" => {
99 for block in blocks() {
100 match block["type"].as_str() {
101 Some("text") => {
102 let text = block["text"].as_str().unwrap_or_default().trim();
103 if !text.is_empty() {
104 reporter.record(Entry::new("message", text));
Agents and memory, checks and conflicts, profiles, slug renames, custom domains105 if let Some(progress) = progress {
106 progress.step(&format!("Said: {}", text.lines().next().unwrap_or_default()));
107 }
Account dropdown, llms.txt onboarding, hosted agent runner (not yet deployed)108 }
109 }
110 Some("tool_use") => {
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look111 crate::abuse::touch();
Account dropdown, llms.txt onboarding, hosted agent runner (not yet deployed)112 let tool = block["name"].as_str().unwrap_or("tool");
Agents and memory, checks and conflicts, profiles, slug renames, custom domains113 if let Some(progress) = progress {
114 progress.step(&progress::describe_tool(tool, &block["input"]));
115 }
Account dropdown, llms.txt onboarding, hosted agent runner (not yet deployed)116 reporter.record(Entry::tool(
117 "tool_call",
118 tool,
119 &describe_input(tool, &block["input"]),
120 ));
121 }
122 _ => {}
123 }
124 }
125 None
126 }
127 "user" => {
128 for block in blocks() {
129 if block["type"] == "tool_result" {
130 let text = result_text(&block["content"]);
131 if !text.trim().is_empty() {
132 reporter.record(Entry::tool("tool_result", "result", text.trim()));
133 }
134 }
135 }
136 None
137 }
138 "result" => {
139 let text = event["result"].as_str().unwrap_or_default().to_owned();
Agents as a team: lifecycle, merge queue, billing and a new shell140 // What the run cost, as the harness worked it out, kept with
141 // the session so that spend can be read per pull request.
142 if let Some(cost) = event["total_cost_usd"].as_f64() {
143 let turns = event["num_turns"].as_u64().unwrap_or_default();
Merge branch 'model-routing'144 report_cost(cost, turns, run_tokens(event));
Agents and memory, checks and conflicts, profiles, slug renames, custom domains145 if let Some(progress) = progress {
146 progress.cost(cost, turns);
147 }
Agents as a team: lifecycle, merge queue, billing and a new shell148 reporter.record(Entry::new(
149 "note",
150 &format!("This run cost ${cost:.4} over {turns} turns."),
151 ));
152 }
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API153 // Claude Code stopped at the run's cost cap (`--max-budget-usd`).
154 if event["subtype"] == "error_max_budget_usd" {
155 return Some(Err(anyhow::Error::new(Halted::Budget)));
156 }
Account dropdown, llms.txt onboarding, hosted agent runner (not yet deployed)157 Some(if event["is_error"].as_bool().unwrap_or(false) {
158 Err(anyhow::anyhow!("the agent reported an error: {text}"))
159 } else {
160 Ok(text)
161 })
162 }
163 _ => None,
164 }
165}
166
167/// Runs Claude Code on `prompt` in `workdir` and returns its closing
168/// summary.
169pub fn run_claude(workdir: &Path, prompt: &str, reporter: &mut Reporter) -> Result<String> {
Agents as a team: lifecycle, merge queue, billing and a new shell170 // g1t's own tools, through a token that can do a few things in this one
171 // repository: read its issues, pull requests and merge queue, open an
172 // issue, and comment.
173 let mut tools = Vec::new();
174 if let Ok(token) = std::env::var("G1T_AGENT_TOKEN") {
175 let config = serde_json::json!({
176 "mcpServers": {
177 "g1t": {
178 "type": "http",
179 "url": std::env::var("G1T_MCP").unwrap_or_else(|_| "https://mcp.g1t.sh".to_owned()),
180 "headers": { "Authorization": format!("Bearer {token}") },
181 }
182 }
183 });
184 let path = "/work/g1t-mcp.json";
185 if std::fs::write(path, config.to_string()).is_ok() {
186 tools = vec!["--mcp-config".to_owned(), path.to_owned()];
187 }
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request188 // People can message the agent while it works: after each tool call
189 // a hook asks g1t for messages and hands any to the agent.
190 if let (Ok(repo), Ok(number)) = (std::env::var("G1T_REPO"), std::env::var("PULL_NUMBER")) {
191 let steer = serde_json::json!({
192 "api": std::env::var("G1T_API").unwrap_or_else(|_| "https://api.g1t.sh".to_owned()),
193 "token": token,
194 "repo": repo,
195 "number": number.parse::<u32>().unwrap_or_default(),
196 });
197 let hooks = serde_json::json!({
198 "hooks": {
199 "PostToolUse": [{
200 "matcher": "*",
201 "hooks": [{ "type": "command", "command": "MODE=steer /usr/local/bin/g1t-runner", "timeout": 15 }],
202 }],
Messages reach the agent even as it finishes203 "Stop": [{
204 "hooks": [{ "type": "command", "command": "MODE=steer G1T_HOOK=stop /usr/local/bin/g1t-runner", "timeout": 15 }],
205 }],
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request206 }
207 });
208 let home = std::env::var("HOME").unwrap_or_else(|_| "/home/node".to_owned());
209 let settings = format!("{home}/.claude/settings.json");
210 if std::fs::write(crate::steer::CONFIG, steer.to_string()).is_ok()
211 && std::fs::create_dir_all(format!("{home}/.claude")).is_ok()
212 {
213 let _ = std::fs::write(settings, hooks.to_string());
214 }
215 }
Agents as a team: lifecycle, merge queue, billing and a new shell216 }
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API217 // The run's guardrails (guard.rs): the hook before every tool call, the
218 // permission rules, and the cost cap, which Claude Code enforces itself.
219 let policy = guard::Policy::from_env();
220 if let Some(policy) = &policy {
221 let installed = guard::install(policy, &std::env::var("GUARDRAILS").unwrap_or_default());
222 if let Some(file) = installed.settings_file {
223 tools.extend(["--settings".to_owned(), file]);
224 }
225 if let Some(budget) = policy.budget_usd.filter(|budget| *budget > 0.0) {
226 tools.extend(["--max-budget-usd".to_owned(), format!("{budget:.2}")]);
227 }
228 }
229 // A fork's checkout is anyone's: none of its CLAUDE.md, .claude
230 // settings, hooks, MCP servers or commands are loaded (guard.rs).
231 let trusted = guard::checkout_trusted_from_env();
232 if !trusted {
233 tools.extend(guard::UNTRUSTED_FLAGS.iter().map(|flag| (*flag).to_owned()));
234 reporter.record(Entry::new(
235 "note",
236 "This checkout is not the repository's own branch, so its CLAUDE.md and .claude settings, hooks, MCP servers and commands were not loaded.",
237 ));
238 }
239 let mut denials = guard::Denials::default();
Agents and memory, checks and conflicts, profiles, slug renames, custom domains240 // How the run goes, step by step, for people watching it live.
241 let mut progress = Progress::from_env();
242 if let Some(progress) = &mut progress {
243 progress.step("Started the agent");
244 }
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API245 let mut command = Command::new("claude");
246 if !trusted {
247 command.env(guard::UNTRUSTED_ENV.0, guard::UNTRUSTED_ENV.1);
248 }
249 let mut child = command
Account dropdown, llms.txt onboarding, hosted agent runner (not yet deployed)250 .current_dir(workdir)
Agents as a team: lifecycle, merge queue, billing and a new shell251 .args(&tools)
Account dropdown, llms.txt onboarding, hosted agent runner (not yet deployed)252 .args([
253 "--print",
254 prompt,
255 "--output-format",
256 "stream-json",
257 "--verbose",
258 "--max-turns",
259 MAX_TURNS,
260 // The sandbox is the permission boundary: it holds one fork and
261 // one short-lived token, and nothing else.
262 "--dangerously-skip-permissions",
263 ])
264 // The agent needs the model key and nothing else of ours.
265 .env_remove("G1T_TOKEN")
Agents as a team: lifecycle, merge queue, billing and a new shell266 .env_remove("REVIEW_TOKEN")
267 .env_remove("CHECK_TOKEN")
268 .env_remove("BILLING_TOKEN")
269 .env_remove("PLAN_TOKEN")
270 .env_remove("G1T_AGENT_TOKEN")
Agents and memory, checks and conflicts, profiles, slug renames, custom domains271 .env_remove("AGENT_RUN_TOKEN")
Account dropdown, llms.txt onboarding, hosted agent runner (not yet deployed)272 .stdin(Stdio::null())
273 .stdout(Stdio::piped())
274 .stderr(Stdio::inherit())
275 .spawn()
276 .context("could not start Claude Code")?;
277
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API278 // The time cap: the agent is stopped when it passes, and the run with it.
279 let timed_out = Arc::new(AtomicBool::new(false));
280 let finished = Arc::new(AtomicBool::new(false));
281 if let Some(minutes) = policy.as_ref().and_then(|policy| policy.minutes) {
282 let (pid, timed_out, finished) = (child.id(), timed_out.clone(), finished.clone());
283 std::thread::spawn(move || {
284 std::thread::sleep(Duration::from_secs(u64::from(minutes) * 60));
285 if !finished.load(Ordering::SeqCst) {
286 timed_out.store(true, Ordering::SeqCst);
287 let _ = Command::new("kill").arg(pid.to_string()).status();
288 }
289 });
290 }
291
Account dropdown, llms.txt onboarding, hosted agent runner (not yet deployed)292 let stdout = child.stdout.take().context("no output from Claude Code")?;
293 let mut outcome = None;
294 for line in BufReader::new(stdout).lines() {
295 let line = line?;
296 let Ok(event) = serde_json::from_str::<Value>(&line) else {
297 continue;
298 };
Agents and memory, checks and conflicts, profiles, slug renames, custom domains299 if let Some(result) = handle_event(&event, reporter, &mut progress) {
Account dropdown, llms.txt onboarding, hosted agent runner (not yet deployed)300 outcome = Some(result);
301 }
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API302 // What the guardrails refused, as steps people can see.
303 for denied in denials.take() {
304 reporter.record(Entry::new("note", &denied));
305 if let Some(progress) = &mut progress {
306 progress.step(&denied);
307 }
308 }
Account dropdown, llms.txt onboarding, hosted agent runner (not yet deployed)309 }
310 let status = child.wait()?;
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API311 finished.store(true, Ordering::SeqCst);
312 if timed_out.load(Ordering::SeqCst) {
313 outcome = Some(Err(anyhow::Error::new(Halted::Time)));
314 }
Agents and memory, checks and conflicts, profiles, slug renames, custom domains315 if let Some(progress) = &mut progress {
316 progress.flush();
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API317 if let Some(halted) = outcome.as_ref().and_then(|o| o.as_ref().err()).and_then(|e| e.downcast_ref::<Halted>()) {
318 let message = match (halted, &policy) {
319 (Halted::Budget, Some(Policy { budget_usd: Some(budget), .. })) => {
320 format!("Stopped: it reached its cost cap of ${budget:.2}.")
321 }
322 (Halted::Time, Some(Policy { minutes: Some(minutes), .. })) => {
323 format!("Stopped: it reached its time cap of {minutes} minutes.")
324 }
325 _ => format!("Stopped: {halted}."),
326 };
327 progress.halt(halted.reason(), &message);
328 }
Agents and memory, checks and conflicts, profiles, slug renames, custom domains329 }
Account dropdown, llms.txt onboarding, hosted agent runner (not yet deployed)330 match outcome {
331 Some(result) => result,
332 None => bail!("Claude Code exited ({status}) without a result"),
333 }
334}
Merge branch 'model-routing'335
336#[cfg(test)]
337mod tests {
338 use super::*;
339
340 #[test]
341 fn a_runs_tokens_come_from_the_closing_usage_by_kind() {
342 let event = serde_json::json!({
343 "type": "result",
344 "total_cost_usd": 0.42,
345 "usage": {
346 "input_tokens": 1200,
347 "output_tokens": 340,
348 "cache_read_input_tokens": 90000,
349 "cache_creation_input_tokens": 5000,
350 },
351 });
352 assert_eq!(
353 run_tokens(&event),
354 serde_json::json!({ "input": 1200, "output": 340, "cache_read": 90000, "cache_write": 5000 })
355 );
356 // An older harness without usage reports nothing, not an error.
357 assert_eq!(
358 run_tokens(&serde_json::json!({ "type": "result" })),
359 serde_json::json!({ "input": 0, "output": 0, "cache_read": 0, "cache_write": 0 })
360 );
361 }
362}

This file's history is long; its oldest lines are credited to the oldest commit read.