Skip to content
322 linesCodeBlameRaw
1//! 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
4use std::sync::atomic::{AtomicBool, Ordering};
5use std::time::{Duration, Instant};
6
7use anyhow::Result;
8use g1t_actions::mask;
9use 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;
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 }
45
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
73pub(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
87 pub(crate) fn report(&self, report: Value) {
88 // A report that cannot be sent is tried a few times, then dropped:
89 // the job goes on, and g1t notices a silent job by itself.
90 for attempt in 0..3 {
91 let sent = ureq::post(&format!("{}/actions/jobs/{}", self.base, self.job))
92 .timeout(Duration::from_secs(30))
93 .send_json(json!({ "token": self.token, "report": report }));
94 match sent {
95 Ok(response) => {
96 if let Ok(answer) = response.into_json::<Value>() {
97 CANCEL.hear(&answer);
98 }
99 return;
100 }
101 // Refused: the job was cancelled or finished; nothing to retry.
102 Err(ureq::Error::Status(code, _)) if (400..500).contains(&code) => return,
103 Err(_) => std::thread::sleep(Duration::from_millis(500 * (attempt + 1))),
104 }
105 }
106 }
107}
108
109/// The job's log, masked, sent in batches. `masks` holds every form of
110/// every secret (`g1t_actions::mask`), longest first.
111pub(crate) struct Log {
112 pub(crate) api: Api,
113 pub(crate) masks: Vec<String>,
114 step: u32,
115 buffer: String,
116 last: Instant,
117 /// When anything was last sent, for the ping.
118 sent: Instant,
119}
120
121impl Log {
122 pub(crate) fn new(api: Api, mut masks: Vec<String>) -> Log {
123 masks.retain(|mask| !mask.is_empty());
124 masks.sort_by_key(|mask| std::cmp::Reverse(mask.len()));
125 Log {
126 api,
127 masks,
128 step: 0,
129 buffer: String::new(),
130 last: Instant::now(),
131 sent: Instant::now(),
132 }
133 }
134
135 /// What waits to be sent.
136 #[cfg(all(test, unix))]
137 pub(crate) fn buffered(&self) -> String {
138 self.buffer.clone()
139 }
140
141 /// Starts writing to step `number` (0 for the job's setup).
142 pub(crate) fn step(&mut self, number: u32) {
143 self.flush();
144 self.step = number;
145 }
146
147 pub(crate) fn mask(&self, text: &str) -> String {
148 mask::apply(text, &self.masks)
149 }
150
151 /// `::add-mask::`: masks `value` from here on, in every form it can
152 /// take (each line, base64, JSON-escaped), as a secret is.
153 pub(crate) fn add_mask(&mut self, value: &str) {
154 let mut added = false;
155 for variant in mask::variants(value) {
156 if !self.masks.contains(&variant) {
157 self.masks.push(variant);
158 added = true;
159 }
160 }
161 if added {
162 self.masks.sort_by_key(|mask| std::cmp::Reverse(mask.len()));
163 }
164 }
165
166 pub(crate) fn line(&mut self, text: &str) {
167 let masked = self.mask(text);
168 self.buffer.push_str(&masked);
169 self.buffer.push('\n');
170 if self.buffer.len() >= FLUSH_BYTES || self.last.elapsed() >= FLUSH_EVERY {
171 self.flush();
172 }
173 }
174
175 /// Sends what is waiting if it has waited long enough, and asks after
176 /// the job when nothing has been sent for a while, so a cancellation
177 /// reaches a step that prints nothing.
178 pub(crate) fn tick(&mut self) {
179 if !self.buffer.is_empty() && self.last.elapsed() >= FLUSH_EVERY {
180 self.flush();
181 } else if self.sent.elapsed() >= PING_EVERY {
182 self.sent = Instant::now();
183 self.api.report(json!({ "kind": "ping" }));
184 }
185 }
186
187 pub(crate) fn flush(&mut self) {
188 self.last = Instant::now();
189 if self.buffer.is_empty() {
190 return;
191 }
192 self.sent = Instant::now();
193 let text = std::mem::take(&mut self.buffer);
194 self.api.report(json!({ "kind": "log", "step": self.step, "text": text }));
195 }
196
197 pub(crate) fn steps(&self, names: &[String]) {
198 self.api.report(json!({ "kind": "steps", "steps": names }));
199 }
200
201 pub(crate) fn step_state(&mut self, number: u32, name: &str, status: &str, conclusion: Option<&str>) {
202 self.flush();
203 let name = self.mask(name);
204 self.api.report(json!({ "kind": "step", "number": number, "name": name, "status": status, "conclusion": conclusion }));
205 }
206
207 /// What a step wrote to `$GITHUB_STEP_SUMMARY`, masked, for the run's
208 /// page. More than 1 MiB is refused, as on GitHub, with an error in the
209 /// log.
210 pub(crate) fn summary(&mut self, markdown: &str) {
211 if markdown.trim().is_empty() {
212 return;
213 }
214 if markdown.len() > MAX_SUMMARY_BYTES {
215 self.line(&format!(
216 "##[error]$GITHUB_STEP_SUMMARY upload aborted: a step's summary may be up to 1024k, and this one is {}k.",
217 markdown.len().div_ceil(1024)
218 ));
219 return;
220 }
221 let markdown = self.mask(markdown);
222 self.flush();
223 self.api.report(json!({ "kind": "summary", "step": self.step, "markdown": markdown }));
224 }
225
226 pub(crate) fn annotation(&mut self, level: &str, message: &str, properties: &serde_json::Map<String, Value>) {
227 let message = self.mask(message);
228 let title = properties.get("title").and_then(Value::as_str).map(|title| self.mask(title));
229 self.api.report(json!({
230 "kind": "annotation",
231 "level": level,
232 "message": message,
233 "title": title,
234 "file": properties.get("file"),
235 "line": properties.get("line").and_then(|l| l.as_str()).and_then(|l| l.parse::<u32>().ok()),
236 }));
237 }
238
239 pub(crate) fn done(&mut self, conclusion: &str, outputs: &serde_json::Map<String, Value>, reason: Option<&str>) {
240 let outputs = self.withhold_secrets(outputs);
241 self.flush();
242 self.api.report(json!({ "kind": "done", "conclusion": conclusion, "outputs": outputs, "reason": reason }));
243 }
244
245 /// The job's outputs without any that hold a secret, as on GitHub: an
246 /// output goes to other jobs and to the run's page, where no mask
247 /// reaches. Each one left out is warned of in the log.
248 pub(crate) fn withhold_secrets(&mut self, outputs: &serde_json::Map<String, Value>) -> serde_json::Map<String, Value> {
249 let mut kept = serde_json::Map::new();
250 for (name, value) in outputs {
251 let text = match value {
252 Value::String(text) => text.clone(),
253 other => other.to_string(),
254 };
255 if mask::reveals(&text, &self.masks) {
256 self.line(&format!("##[warning]The output `{name}` was left out: it holds a secret."));
257 continue;
258 }
259 kept.insert(name.clone(), value.clone());
260 }
261 kept
262 }
263}
264
265#[cfg(test)]
266mod tests {
267 use super::*;
268
269 fn log(secrets: &[&str]) -> Log {
270 let api = Api { base: "http://127.0.0.1:9".into(), job: "job_1".into(), token: "t".into() };
271 Log::new(api, mask::all_variants(secrets.iter().copied()))
272 }
273
274 #[test]
275 fn outputs_holding_a_secret_are_left_out() {
276 let mut log = log(&["s3cr3t-token"]);
277 let mut outputs = serde_json::Map::new();
278 outputs.insert("version".into(), json!("1.2.0"));
279 outputs.insert("leak".into(), json!("token=s3cr3t-token"));
280 outputs.insert("encoded".into(), json!("czNjcjN0LXRva2Vu"));
281 let kept = log.withhold_secrets(&outputs);
282 assert_eq!(kept.keys().collect::<Vec<_>>(), ["version"]);
283 assert!(log.buffer.contains("The output `leak` was left out"));
284 assert!(log.buffer.contains("The output `encoded` was left out"));
285 }
286
287 #[test]
288 fn a_cancelled_answer_is_heard_once_and_taken_in() {
289 let cancel = Cancel::new();
290 cancel.hear(&json!({ "ok": true }));
291 assert!(!cancel.cancelled());
292 cancel.hear(&json!({ "ok": true, "cancelled": true }));
293 assert!(cancel.cancelled() && cancel.interrupt());
294 cancel.take();
295 assert!(cancel.cancelled() && !cancel.interrupt(), "cleanup steps are not stopped again");
296 cancel.hear(&json!({ "ok": true, "cancelled": false }));
297 assert!(cancel.cancelled(), "a cancelled job stays cancelled");
298 }
299
300 #[test]
301 fn a_summary_past_its_limit_is_refused_in_the_log() {
302 let mut log = log(&[]);
303 log.summary("
304");
305 assert!(log.buffer.is_empty(), "an empty summary is not sent");
306 log.summary(&"x".repeat(MAX_SUMMARY_BYTES + 1));
307 assert!(log.buffer.contains("$GITHUB_STEP_SUMMARY upload aborted"));
308 assert!(log.buffer.contains("1025k"));
309 }
310
311 #[test]
312 fn added_masks_cover_every_form() {
313 let mut log = log(&[]);
314 log.add_mask("line one\nline two");
315 assert_eq!(log.mask("first: line one"), "first: ***");
316 assert_eq!(log.mask("then line two"), "then ***");
317 // Longest first: the whole value, not its pieces.
318 log.add_mask("abc");
319 log.add_mask("abcdef");
320 assert_eq!(log.mask("abcdef"), "***");
321 }
322}