g1t/crates/runner/src/report.rs
| 1 | //! Reports a pull request's progress to g1t through its public API. |
| 2 | |
| 3 | use std::time::{Duration, Instant}; |
| 4 | |
| 5 | use anyhow::{Context, Result}; |
| 6 | use serde::Serialize; |
| 7 | |
| 8 | /// Entries are sent in batches at most this often. |
| 9 | const FLUSH_INTERVAL: Duration = Duration::from_millis(1500); |
| 10 | const FLUSH_SIZE: usize = 50; |
| 11 | /// Tool output can be enormous; the session keeps the start of it. |
| 12 | const MAX_ENTRY_CHARS: usize = 8000; |
| 13 | |
| 14 | /// One step of the session, as the API accepts it. |
| 15 | #[derive(Debug, Serialize)] |
| 16 | pub struct Entry { |
| 17 | kind: &'static str, |
| 18 | text: String, |
| 19 | #[serde(skip_serializing_if = "Option::is_none")] |
| 20 | tool: Option<String>, |
| 21 | } |
| 22 | |
| 23 | impl Entry { |
| 24 | pub fn new(kind: &'static str, text: &str) -> Self { |
| 25 | let mut chars = text.chars(); |
| 26 | let mut text: String = chars.by_ref().take(MAX_ENTRY_CHARS).collect(); |
| 27 | if chars.next().is_some() { |
| 28 | text.push_str("\n… (truncated)"); |
| 29 | } |
| 30 | Entry { |
| 31 | kind, |
| 32 | text, |
| 33 | tool: None, |
| 34 | } |
| 35 | } |
| 36 | |
| 37 | pub fn tool(kind: &'static str, tool: &str, text: &str) -> Self { |
| 38 | Entry { |
| 39 | tool: Some(tool.to_owned()), |
| 40 | ..Entry::new(kind, text) |
| 41 | } |
| 42 | } |
| 43 | } |
| 44 | |
| 45 | /// Replaces every occurrence of a secret with a marker. An agent can print |
| 46 | /// its environment or a git remote; whatever it prints is recorded. |
| 47 | fn redact(text: &str, secrets: &[String]) -> String { |
| 48 | secrets.iter().fold(text.to_owned(), |text, secret| { |
| 49 | text.replace(secret, "[redacted]") |
| 50 | }) |
| 51 | } |
| 52 | |
| 53 | pub struct Reporter { |
| 54 | api: String, |
| 55 | token: String, |
| 56 | /// The pull request's path in the API: `repos/<owner>/<name>/pulls/<number>`. |
| 57 | pull: String, |
| 58 | /// Values that must never reach a session, which is as public as the |
| 59 | /// repository: the credentials this process was started with. |
| 60 | secrets: Vec<String>, |
| 61 | pending: Vec<Entry>, |
| 62 | last_flush: Instant, |
| 63 | /// False for work that has no session: entries are dropped. |
| 64 | recording: bool, |
| 65 | } |
| 66 | |
| 67 | impl Reporter { |
| 68 | /// A reporter that records nothing, for a run with no session to keep. |
| 69 | pub fn silent() -> Self { |
| 70 | Reporter { |
| 71 | api: String::new(), |
| 72 | token: String::new(), |
| 73 | pull: String::new(), |
| 74 | secrets: Vec::new(), |
| 75 | pending: Vec::new(), |
| 76 | last_flush: Instant::now(), |
| 77 | recording: false, |
| 78 | } |
| 79 | } |
| 80 | |
| 81 | pub fn from_env() -> Result<Self> { |
| 82 | let var = |name: &str| std::env::var(name).with_context(|| format!("{name} is not set")); |
| 83 | let token = var("G1T_TOKEN")?; |
| 84 | let secrets = std::iter::once(token.clone()) |
| 85 | .chain(std::env::var("ANTHROPIC_API_KEY")) |
| 86 | .chain(std::env::var("AI_GATEWAY_TOKEN")) |
| 87 | .chain(std::env::var("BILLING_TOKEN")) |
| 88 | .filter(|secret| !secret.is_empty()) |
| 89 | .collect(); |
| 90 | Ok(Reporter { |
| 91 | api: var("G1T_API")?, |
| 92 | token, |
| 93 | pull: format!("repos/{}/pulls/{}", var("G1T_REPO")?, var("PULL_NUMBER")?), |
| 94 | secrets, |
| 95 | pending: Vec::new(), |
| 96 | last_flush: Instant::now(), |
| 97 | recording: true, |
| 98 | }) |
| 99 | } |
| 100 | |
| 101 | fn post(&self, action: &str, body: serde_json::Value) -> Result<()> { |
| 102 | ureq::post(&format!("{}/{}/{action}", self.api, self.pull)) |
| 103 | .set("authorization", &format!("Bearer {}", self.token)) |
| 104 | .send_json(body) |
| 105 | .with_context(|| format!("{action} request failed"))?; |
| 106 | Ok(()) |
| 107 | } |
| 108 | |
| 109 | /// Queues an entry, sending the batch if it is due. |
| 110 | pub fn record(&mut self, mut entry: Entry) { |
| 111 | if !self.recording { |
| 112 | return; |
| 113 | } |
| 114 | entry.text = redact(&entry.text, &self.secrets); |
| 115 | self.pending.push(entry); |
| 116 | if self.pending.len() >= FLUSH_SIZE || self.last_flush.elapsed() >= FLUSH_INTERVAL { |
| 117 | self.flush(); |
| 118 | } |
| 119 | } |
| 120 | |
| 121 | /// Sends everything queued. A failed send is logged and the entries |
| 122 | /// kept, so one bad request does not lose the session or stop the run. |
| 123 | pub fn flush(&mut self) { |
| 124 | self.last_flush = Instant::now(); |
| 125 | if self.pending.is_empty() { |
| 126 | return; |
| 127 | } |
| 128 | let body = serde_json::json!({ "entries": self.pending }); |
| 129 | match self.post("session", body) { |
| 130 | Ok(()) => self.pending.clear(), |
| 131 | Err(error) => eprintln!("g1t-runner: {error:#}"), |
| 132 | } |
| 133 | } |
| 134 | |
| 135 | /// Marks the pull request ready for review, with `summary` as its |
| 136 | /// description. |
| 137 | pub fn ready(&self, summary: &str) -> Result<()> { |
| 138 | self.post("ready", serde_json::json!({ "summary": summary })) |
| 139 | } |
| 140 | |
| 141 | /// Closes the pull request without merging. |
| 142 | pub fn close(&self) -> Result<()> { |
| 143 | self.post("close", serde_json::json!({})) |
| 144 | } |
| 145 | } |
| 146 | |
| 147 | #[cfg(test)] |
| 148 | mod tests { |
| 149 | use super::*; |
| 150 | |
| 151 | #[test] |
| 152 | fn secrets_are_removed_from_entries() { |
| 153 | let secrets = vec!["g1t_abc123".to_owned(), "sk-ant-xyz".to_owned()]; |
| 154 | let text = "origin https://me:g1t_abc123@g1t.sh/a.git |
| 155 | KEY=sk-ant-xyz g1t_abc123"; |
| 156 | let clean = redact(text, &secrets); |
| 157 | assert!(!clean.contains("g1t_abc123")); |
| 158 | assert!(!clean.contains("sk-ant-xyz")); |
| 159 | assert_eq!(clean.matches("[redacted]").count(), 3); |
| 160 | } |
| 161 | |
| 162 | #[test] |
| 163 | fn long_entries_are_truncated() { |
| 164 | let entry = Entry::new("note", &"x".repeat(MAX_ENTRY_CHARS + 10)); |
| 165 | assert!(entry.text.ends_with("(truncated)")); |
| 166 | assert!(Entry::new("note", "short").text == "short"); |
| 167 | } |
| 168 | } |