Skip to content
341 linesCodeBlameRaw

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.

GitHub Actions on g1t, part two: running workflows1//! Telling g1t how a job is going: its steps, its log in batches, its
2//! annotations, and how it ended. Every report carries the job's token.
3
Merge Actions runs: summaries, attempts and re-runs, graceful cancel, log downloads, badges (actions 0009)4use std::sync::atomic::{AtomicBool, Ordering};
GitHub Actions on g1t, part two: running workflows5use std::time::{Duration, Instant};
6
7use anyhow::Result;
Merge branch 'worktree-agent-a3abfcce648e87dca'8use g1t_actions::mask;
GitHub Actions on g1t, part two: running workflows9use serde_json::{Value, json};
10
11/// How long log lines wait before they are sent.
12const FLUSH_EVERY: Duration = Duration::from_millis(1500);
13/// How much log is sent at once.
14const FLUSH_BYTES: usize = 64 * 1024;
Merge Actions runs: summaries, attempts and re-runs, graceful cancel, log downloads, badges (actions 0009)15/// The most one step may add to the job's summary, as on GitHub.
16pub(crate) const MAX_SUMMARY_BYTES: usize = 1024 * 1024;
17/// How long the runner goes without a report before it asks whether the
18/// job was cancelled.
19const PING_EVERY: Duration = Duration::from_secs(10);
20
21/// Whether the job was cancelled, as g1t's answers to its reports say.
22struct Cancel {
23 /// g1t said so.
24 said: AtomicBool,
25 /// The job has taken it in: the step it was on was stopped, and the
26 /// cleanup steps that follow are not stopped for it again.
27 taken: AtomicBool,
28}
29
30impl Cancel {
31 const fn new() -> Cancel {
32 Cancel { said: AtomicBool::new(false), taken: AtomicBool::new(false) }
33 }
34
35 /// Reads an answer to a report: `cancelled` when the run was cancelled.
36 fn hear(&self, answer: &Value) {
37 if answer["cancelled"].as_bool() == Some(true) {
38 self.said.store(true, Ordering::Relaxed);
39 }
40 }
41
42 fn cancelled(&self) -> bool {
43 self.said.load(Ordering::Relaxed)
44 }
GitHub Actions on g1t, part two: running workflows45
Merge Actions runs: summaries, attempts and re-runs, graceful cancel, log downloads, badges (actions 0009)46 fn interrupt(&self) -> bool {
47 self.cancelled() && !self.taken.load(Ordering::Relaxed)
48 }
49
50 fn take(&self) {
51 self.taken.store(true, Ordering::Relaxed);
52 }
53}
54
55/// This job's (one runs per process).
56static CANCEL: Cancel = Cancel::new();
57
58/// Whether g1t said the job was cancelled.
59pub(crate) fn cancelled() -> bool {
60 CANCEL.cancelled()
61}
62
63/// Whether a running step should be stopped: cancelled, and not yet taken in.
64pub(crate) fn interrupt() -> bool {
65 CANCEL.interrupt()
66}
67
68/// The job has seen the cancellation: what runs from here is its cleanup.
69pub(crate) fn take_cancel() {
70 CANCEL.take();
71}
72
GitHub Actions on g1t, part two: running workflows73pub(crate) struct Api {
74 pub(crate) base: String,
75 pub(crate) job: String,
76 pub(crate) token: String,
77}
78
79impl Api {
80 pub(crate) fn spec(&self) -> Result<Value> {
81 let response = ureq::post(&format!("{}/actions/jobs/{}/spec", self.base, self.job))
82 .timeout(Duration::from_secs(60))
83 .send_json(json!({ "token": self.token }))?;
84 Ok(response.into_json()?)
85 }
86
Merge Actions: cross-repo workflows and actions, release and deployment triggers, step timeouts87 /// Where to fetch another repository's action from: `{"source": "g1t",
88 /// "url", "ref", "token"}` or `{"source": "github"}`. Err holds why g1t
89 /// refused (`true`: the repository is private and may not be used
90 /// here), or that it could not be asked (`false`).
91 pub(crate) fn action(&self, repository: &str, git_ref: &str) -> std::result::Result<Value, (bool, String)> {
92 let sent = ureq::post(&format!("{}/actions/jobs/{}/action", self.base, self.job))
93 .timeout(Duration::from_secs(30))
94 .send_json(json!({ "token": self.token, "report": { "repository": repository, "ref": git_ref } }));
95 match sent {
96 Ok(response) => response.into_json().map_err(|error| (false, error.to_string())),
97 Err(ureq::Error::Status(code, response)) => {
98 let body: Value = response.into_json().unwrap_or(Value::Null);
99 let message = body["error"]["message"].as_str().unwrap_or("g1t did not answer.").to_owned();
100 Err((code == 403, message))
101 }
102 Err(error) => Err((false, error.to_string())),
103 }
104 }
105
GitHub Actions on g1t, part two: running workflows106 pub(crate) fn report(&self, report: Value) {
107 // A report that cannot be sent is tried a few times, then dropped:
108 // the job goes on, and g1t notices a silent job by itself.
109 for attempt in 0..3 {
110 let sent = ureq::post(&format!("{}/actions/jobs/{}", self.base, self.job))
111 .timeout(Duration::from_secs(30))
112 .send_json(json!({ "token": self.token, "report": report }));
113 match sent {
Merge Actions runs: summaries, attempts and re-runs, graceful cancel, log downloads, badges (actions 0009)114 Ok(response) => {
115 if let Ok(answer) = response.into_json::<Value>() {
116 CANCEL.hear(&answer);
117 }
118 return;
119 }
GitHub Actions on g1t, part two: running workflows120 // Refused: the job was cancelled or finished; nothing to retry.
121 Err(ureq::Error::Status(code, _)) if (400..500).contains(&code) => return,
122 Err(_) => std::thread::sleep(Duration::from_millis(500 * (attempt + 1))),
123 }
124 }
125 }
126}
127
Merge branch 'worktree-agent-a3abfcce648e87dca'128/// The job's log, masked, sent in batches. `masks` holds every form of
129/// every secret (`g1t_actions::mask`), longest first.
GitHub Actions on g1t, part two: running workflows130pub(crate) struct Log {
131 pub(crate) api: Api,
132 pub(crate) masks: Vec<String>,
133 step: u32,
134 buffer: String,
135 last: Instant,
Merge Actions runs: summaries, attempts and re-runs, graceful cancel, log downloads, badges (actions 0009)136 /// When anything was last sent, for the ping.
137 sent: Instant,
GitHub Actions on g1t, part two: running workflows138}
139
140impl Log {
Merge branch 'worktree-agent-a3abfcce648e87dca'141 pub(crate) fn new(api: Api, mut masks: Vec<String>) -> Log {
142 masks.retain(|mask| !mask.is_empty());
143 masks.sort_by_key(|mask| std::cmp::Reverse(mask.len()));
GitHub Actions on g1t, part two: running workflows144 Log {
145 api,
146 masks,
147 step: 0,
148 buffer: String::new(),
149 last: Instant::now(),
Merge Actions runs: summaries, attempts and re-runs, graceful cancel, log downloads, badges (actions 0009)150 sent: Instant::now(),
GitHub Actions on g1t, part two: running workflows151 }
152 }
153
Merge Actions runs: summaries, attempts and re-runs, graceful cancel, log downloads, badges (actions 0009)154 /// What waits to be sent.
155 #[cfg(all(test, unix))]
156 pub(crate) fn buffered(&self) -> String {
157 self.buffer.clone()
158 }
159
GitHub Actions on g1t, part two: running workflows160 /// Starts writing to step `number` (0 for the job's setup).
161 pub(crate) fn step(&mut self, number: u32) {
162 self.flush();
163 self.step = number;
164 }
165
166 pub(crate) fn mask(&self, text: &str) -> String {
Merge branch 'worktree-agent-a3abfcce648e87dca'167 mask::apply(text, &self.masks)
168 }
169
170 /// `::add-mask::`: masks `value` from here on, in every form it can
171 /// take (each line, base64, JSON-escaped), as a secret is.
172 pub(crate) fn add_mask(&mut self, value: &str) {
173 let mut added = false;
174 for variant in mask::variants(value) {
175 if !self.masks.contains(&variant) {
176 self.masks.push(variant);
177 added = true;
GitHub Actions on g1t, part two: running workflows178 }
179 }
Merge branch 'worktree-agent-a3abfcce648e87dca'180 if added {
181 self.masks.sort_by_key(|mask| std::cmp::Reverse(mask.len()));
182 }
GitHub Actions on g1t, part two: running workflows183 }
184
185 pub(crate) fn line(&mut self, text: &str) {
186 let masked = self.mask(text);
187 self.buffer.push_str(&masked);
188 self.buffer.push('\n');
189 if self.buffer.len() >= FLUSH_BYTES || self.last.elapsed() >= FLUSH_EVERY {
190 self.flush();
191 }
192 }
193
Merge Actions runs: summaries, attempts and re-runs, graceful cancel, log downloads, badges (actions 0009)194 /// Sends what is waiting if it has waited long enough, and asks after
195 /// the job when nothing has been sent for a while, so a cancellation
196 /// reaches a step that prints nothing.
GitHub Actions on g1t, part two: running workflows197 pub(crate) fn tick(&mut self) {
198 if !self.buffer.is_empty() && self.last.elapsed() >= FLUSH_EVERY {
199 self.flush();
Merge Actions runs: summaries, attempts and re-runs, graceful cancel, log downloads, badges (actions 0009)200 } else if self.sent.elapsed() >= PING_EVERY {
201 self.sent = Instant::now();
202 self.api.report(json!({ "kind": "ping" }));
GitHub Actions on g1t, part two: running workflows203 }
204 }
205
206 pub(crate) fn flush(&mut self) {
207 self.last = Instant::now();
208 if self.buffer.is_empty() {
209 return;
210 }
Merge Actions runs: summaries, attempts and re-runs, graceful cancel, log downloads, badges (actions 0009)211 self.sent = Instant::now();
GitHub Actions on g1t, part two: running workflows212 let text = std::mem::take(&mut self.buffer);
213 self.api.report(json!({ "kind": "log", "step": self.step, "text": text }));
214 }
215
216 pub(crate) fn steps(&self, names: &[String]) {
217 self.api.report(json!({ "kind": "steps", "steps": names }));
218 }
219
220 pub(crate) fn step_state(&mut self, number: u32, name: &str, status: &str, conclusion: Option<&str>) {
221 self.flush();
222 let name = self.mask(name);
223 self.api.report(json!({ "kind": "step", "number": number, "name": name, "status": status, "conclusion": conclusion }));
224 }
225
Merge Actions runs: summaries, attempts and re-runs, graceful cancel, log downloads, badges (actions 0009)226 /// What a step wrote to `$GITHUB_STEP_SUMMARY`, masked, for the run's
227 /// page. More than 1 MiB is refused, as on GitHub, with an error in the
228 /// log.
229 pub(crate) fn summary(&mut self, markdown: &str) {
230 if markdown.trim().is_empty() {
231 return;
232 }
233 if markdown.len() > MAX_SUMMARY_BYTES {
234 self.line(&format!(
235 "##[error]$GITHUB_STEP_SUMMARY upload aborted: a step's summary may be up to 1024k, and this one is {}k.",
236 markdown.len().div_ceil(1024)
237 ));
238 return;
239 }
240 let markdown = self.mask(markdown);
241 self.flush();
242 self.api.report(json!({ "kind": "summary", "step": self.step, "markdown": markdown }));
243 }
244
GitHub Actions on g1t, part two: running workflows245 pub(crate) fn annotation(&mut self, level: &str, message: &str, properties: &serde_json::Map<String, Value>) {
246 let message = self.mask(message);
Merge branch 'worktree-agent-a3abfcce648e87dca'247 let title = properties.get("title").and_then(Value::as_str).map(|title| self.mask(title));
GitHub Actions on g1t, part two: running workflows248 self.api.report(json!({
249 "kind": "annotation",
250 "level": level,
251 "message": message,
Merge branch 'worktree-agent-a3abfcce648e87dca'252 "title": title,
GitHub Actions on g1t, part two: running workflows253 "file": properties.get("file"),
254 "line": properties.get("line").and_then(|l| l.as_str()).and_then(|l| l.parse::<u32>().ok()),
255 }));
256 }
257
258 pub(crate) fn done(&mut self, conclusion: &str, outputs: &serde_json::Map<String, Value>, reason: Option<&str>) {
Merge branch 'worktree-agent-a3abfcce648e87dca'259 let outputs = self.withhold_secrets(outputs);
GitHub Actions on g1t, part two: running workflows260 self.flush();
261 self.api.report(json!({ "kind": "done", "conclusion": conclusion, "outputs": outputs, "reason": reason }));
262 }
Merge branch 'worktree-agent-a3abfcce648e87dca'263
264 /// The job's outputs without any that hold a secret, as on GitHub: an
265 /// output goes to other jobs and to the run's page, where no mask
266 /// reaches. Each one left out is warned of in the log.
267 pub(crate) fn withhold_secrets(&mut self, outputs: &serde_json::Map<String, Value>) -> serde_json::Map<String, Value> {
268 let mut kept = serde_json::Map::new();
269 for (name, value) in outputs {
270 let text = match value {
271 Value::String(text) => text.clone(),
272 other => other.to_string(),
273 };
274 if mask::reveals(&text, &self.masks) {
275 self.line(&format!("##[warning]The output `{name}` was left out: it holds a secret."));
276 continue;
277 }
278 kept.insert(name.clone(), value.clone());
279 }
280 kept
281 }
282}
283
284#[cfg(test)]
285mod tests {
286 use super::*;
287
288 fn log(secrets: &[&str]) -> Log {
289 let api = Api { base: "http://127.0.0.1:9".into(), job: "job_1".into(), token: "t".into() };
290 Log::new(api, mask::all_variants(secrets.iter().copied()))
291 }
292
293 #[test]
294 fn outputs_holding_a_secret_are_left_out() {
295 let mut log = log(&["s3cr3t-token"]);
296 let mut outputs = serde_json::Map::new();
297 outputs.insert("version".into(), json!("1.2.0"));
298 outputs.insert("leak".into(), json!("token=s3cr3t-token"));
299 outputs.insert("encoded".into(), json!("czNjcjN0LXRva2Vu"));
300 let kept = log.withhold_secrets(&outputs);
301 assert_eq!(kept.keys().collect::<Vec<_>>(), ["version"]);
302 assert!(log.buffer.contains("The output `leak` was left out"));
303 assert!(log.buffer.contains("The output `encoded` was left out"));
304 }
305
306 #[test]
Merge Actions runs: summaries, attempts and re-runs, graceful cancel, log downloads, badges (actions 0009)307 fn a_cancelled_answer_is_heard_once_and_taken_in() {
308 let cancel = Cancel::new();
309 cancel.hear(&json!({ "ok": true }));
310 assert!(!cancel.cancelled());
311 cancel.hear(&json!({ "ok": true, "cancelled": true }));
312 assert!(cancel.cancelled() && cancel.interrupt());
313 cancel.take();
314 assert!(cancel.cancelled() && !cancel.interrupt(), "cleanup steps are not stopped again");
315 cancel.hear(&json!({ "ok": true, "cancelled": false }));
316 assert!(cancel.cancelled(), "a cancelled job stays cancelled");
317 }
318
319 #[test]
320 fn a_summary_past_its_limit_is_refused_in_the_log() {
321 let mut log = log(&[]);
322 log.summary("
323");
324 assert!(log.buffer.is_empty(), "an empty summary is not sent");
325 log.summary(&"x".repeat(MAX_SUMMARY_BYTES + 1));
326 assert!(log.buffer.contains("$GITHUB_STEP_SUMMARY upload aborted"));
327 assert!(log.buffer.contains("1025k"));
328 }
329
330 #[test]
Merge branch 'worktree-agent-a3abfcce648e87dca'331 fn added_masks_cover_every_form() {
332 let mut log = log(&[]);
333 log.add_mask("line one\nline two");
334 assert_eq!(log.mask("first: line one"), "first: ***");
335 assert_eq!(log.mask("then line two"), "then ***");
336 // Longest first: the whole value, not its pieces.
337 log.add_mask("abc");
338 log.add_mask("abcdef");
339 assert_eq!(log.mask("abcdef"), "***");
340 }
GitHub Actions on g1t, part two: running workflows341}

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