Skip to content
331 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//! Running one process for a step: its output streamed to the log as it
2//! comes, with GitHub's workflow commands (`::error::`, `::group::`,
3//! `::add-mask::`…) read out of it.
4
5use std::collections::BTreeMap;
6use std::io::{BufRead, BufReader, Read};
7use std::process::{Command, Stdio};
8use std::sync::mpsc;
9use std::time::{Duration, Instant};
10
11use serde_json::{Map, Value};
12
13use super::report::Log;
14
15/// What a step's workflow commands left behind.
16#[derive(Default)]
17pub(crate) struct Commands {
18 /// `::set-output` (old, still honoured).
19 pub(crate) outputs: BTreeMap<String, String>,
20 /// `::save-state`, for the action's post step.
21 pub(crate) state: BTreeMap<String, String>,
22 /// Set by `::stop-commands::token` until `::token::`.
23 pub(crate) stopped: Option<String>,
24 /// Whether `::debug::` lines are shown (`ACTIONS_STEP_DEBUG`).
25 pub(crate) debug: bool,
26}
27
28/// `%25`, `%0D`, `%0A`, and in properties `%3A` and `%2C`, as the
29/// toolkit escapes them.
30fn unescape(text: &str, property: bool) -> String {
31 let mut out = text.replace("%0D", "\r").replace("%0A", "\n");
32 if property {
33 out = out.replace("%3A", ":").replace("%2C", ",");
34 }
35 out.replace("%25", "%")
36}
37
38/// `::name key=value,key=value::message`, if the line is a command.
39pub(crate) fn parse_command(line: &str) -> Option<(String, Map<String, Value>, String)> {
40 let rest = line.trim_start().strip_prefix("::")?;
41 let end = rest.find("::")?;
42 let (head, data) = (&rest[..end], &rest[end + 2..]);
43 let (name, properties) = match head.split_once(' ') {
44 Some((name, properties)) => (name, properties),
45 None => (head, ""),
46 };
47 if name.is_empty() || !name.chars().all(|c| c.is_ascii_alphanumeric() || c == '-' || c == '_') {
48 return None;
49 }
50 let mut map = Map::new();
51 for pair in properties.split(',').filter(|p| !p.trim().is_empty()) {
52 if let Some((key, value)) = pair.split_once('=') {
53 map.insert(key.trim().to_owned(), Value::String(unescape(value, true)));
54 }
55 }
56 Some((name.to_owned(), map, unescape(data, false)))
57}
58
59impl Commands {
60 /// Handles one line of output: a command is acted on, and the line to
61 /// show (if any) is returned.
62 pub(crate) fn handle(&mut self, line: &str, log: &mut Log) -> Option<String> {
63 if let Some(token) = &self.stopped {
64 if line.trim() == format!("::{token}::") {
65 self.stopped = None;
66 return None;
67 }
68 return Some(line.to_owned());
69 }
70 let Some((name, properties, data)) = parse_command(line) else {
71 return Some(line.to_owned());
72 };
73 match name.as_str() {
74 "add-mask" => {
75 if !data.trim().is_empty() {
Merge branch 'worktree-agent-a3abfcce648e87dca'76 log.add_mask(data.trim());
GitHub Actions on g1t, part two: running workflows77 }
78 None
79 }
80 "error" | "warning" | "notice" => {
81 log.annotation(&name, &data, &properties);
82 let label = match name.as_str() {
83 "error" => "Error",
84 "warning" => "Warning",
85 _ => "Notice",
86 };
87 Some(format!("##[{name}]{label}: {data}"))
88 }
89 "group" => Some(format!("##[group]{data}")),
90 "endgroup" => Some("##[endgroup]".to_owned()),
91 "debug" => self.debug.then(|| format!("##[debug]{data}")),
92 "set-output" => {
93 if let Some(name) = properties.get("name").and_then(Value::as_str) {
94 self.outputs.insert(name.to_owned(), data);
95 }
96 None
97 }
98 "save-state" => {
99 if let Some(name) = properties.get("name").and_then(Value::as_str) {
100 self.state.insert(name.to_owned(), data);
101 }
102 None
103 }
104 "stop-commands" => {
105 self.stopped = Some(data);
106 None
107 }
108 "echo" => None,
109 "add-path" | "set-env" => Some(format!(
110 "##[error]The `{name}` command is disabled, as on GitHub. Write to the file in $GITHUB_{} instead.",
111 if name == "add-path" { "PATH" } else { "ENV" }
112 )),
113 _ => Some(line.to_owned()),
114 }
115 }
116}
117
118/// How a process ended.
119pub(crate) enum Ended {
120 Exited(i32),
121 TimedOut,
Merge Actions runs: summaries, attempts and re-runs, graceful cancel, log downloads, badges (actions 0009)122 /// The run was cancelled: it was interrupted, then stopped.
123 Cancelled,
124}
125
126/// After a cancellation, how long a step has after SIGINT before SIGTERM,
127/// and after SIGTERM before it is killed, as GitHub's runner waits.
128const INTERRUPT_GRACE: Duration = Duration::from_millis(7500);
129const TERMINATE_GRACE: Duration = Duration::from_millis(2500);
130
131/// Sends `signal` to the process's group (it leads its own), so what the
132/// step started hears it too. Windows has no signals: it is left to `kill`.
Merge remote-tracking branch 'origin/main' into workspace-chat133/// Sent with kill(2) itself: a machine without the `kill` program (procps
134/// is not in every image) would otherwise give the step no SIGINT at all,
135/// and leave what it started running after the step was killed.
136#[cfg(unix)]
Merge Actions runs: summaries, attempts and re-runs, graceful cancel, log downloads, badges (actions 0009)137fn signal(child: &std::process::Child, signal: &str) {
Merge remote-tracking branch 'origin/main' into workspace-chat138 let number = match signal {
139 "INT" => libc::SIGINT,
140 "TERM" => libc::SIGTERM,
141 _ => libc::SIGKILL,
142 };
143 let Ok(group) = libc::pid_t::try_from(child.id()) else {
144 return;
145 };
146 // SAFETY: kill(2) with a negative pid signals that process group; it
147 // reads no memory of ours.
148 unsafe {
149 libc::kill(-group, number);
Merge Actions runs: summaries, attempts and re-runs, graceful cancel, log downloads, badges (actions 0009)150 }
GitHub Actions on g1t, part two: running workflows151}
152
Merge remote-tracking branch 'origin/main' into workspace-chat153#[cfg(not(unix))]
154fn signal(_child: &std::process::Child, _signal: &str) {}
155
GitHub Actions on g1t, part two: running workflows156/// Runs the command, sending its output (stdout and stderr together, a
157/// line at a time) through `commands` to the log, until it ends or
158/// `timeout` passes.
Merge Actions runs: summaries, attempts and re-runs, graceful cancel, log downloads, badges (actions 0009)159pub(crate) fn run(command: Command, timeout: Duration, log: &mut Log, commands: &mut Commands) -> std::io::Result<Ended> {
160 run_until(command, timeout, log, commands, &super::report::interrupt)
161}
162
163/// `run`, stopping the process gracefully once `interrupt` says so.
164fn run_until(
165 mut command: Command,
166 timeout: Duration,
167 log: &mut Log,
168 commands: &mut Commands,
169 interrupt: &dyn Fn() -> bool,
170) -> std::io::Result<Ended> {
GitHub Actions on g1t, part two: running workflows171 command.stdin(Stdio::null()).stdout(Stdio::piped()).stderr(Stdio::piped());
Merge Actions runs: summaries, attempts and re-runs, graceful cancel, log downloads, badges (actions 0009)172 // A group of its own, so a cancellation reaches what the step started.
173 #[cfg(unix)]
Merge remote-tracking branch 'origin/main' into workspace-chat174 {
175 use std::os::unix::process::CommandExt;
176 command.process_group(0);
177 // A shell cannot trap a signal it was started with ignored, and a
178 // runner started in the background (as a service, or under `&`)
179 // passes SIGINT on ignored. Back to the defaults, so a cancelled
180 // step hears SIGINT and its own trap runs.
181 // SAFETY: only signal(), which is async-signal-safe, between fork and exec.
182 unsafe {
183 command.pre_exec(|| {
184 for signal in [libc::SIGINT, libc::SIGTERM, libc::SIGQUIT] {
185 libc::signal(signal, libc::SIG_DFL);
186 }
187 Ok(())
188 });
189 }
190 }
GitHub Actions on g1t, part two: running workflows191 let mut child = command.spawn()?;
192 let (sender, lines) = mpsc::channel::<String>();
193 let mut readers = Vec::new();
194 let pipes: Vec<Box<dyn Read + Send>> = vec![
195 Box::new(child.stdout.take().expect("piped")),
196 Box::new(child.stderr.take().expect("piped")),
197 ];
198 for pipe in pipes {
199 let sender = sender.clone();
200 readers.push(std::thread::spawn(move || {
201 let mut reader = BufReader::new(pipe);
202 let mut buffer = Vec::new();
203 loop {
204 buffer.clear();
205 match reader.read_until(b'\n', &mut buffer) {
206 Ok(0) | Err(_) => break,
207 Ok(_) => {
208 let text = String::from_utf8_lossy(&buffer);
209 let text = text.trim_end_matches(['\n', '\r']);
210 // A progress bar redraws with \r; keep its last state.
211 let text = text.rsplit('\r').next().unwrap_or(text);
212 if sender.send(text.to_owned()).is_err() {
213 break;
214 }
215 }
216 }
217 }
218 }));
219 }
220 drop(sender);
221 let deadline = Instant::now() + timeout;
222 let mut timed_out = false;
Merge Actions runs: summaries, attempts and re-runs, graceful cancel, log downloads, badges (actions 0009)223 // When the cancellation reached it, and which signal it has had.
224 let mut interrupted: Option<(Instant, u8)> = None;
GitHub Actions on g1t, part two: running workflows225 loop {
226 match lines.recv_timeout(Duration::from_millis(250)) {
227 Ok(line) => {
228 if let Some(shown) = commands.handle(&line, log) {
229 log.line(&shown);
230 }
231 }
Merge branch 'main' into actions-toolkit-oidc-artifacts232 Err(mpsc::RecvTimeoutError::Timeout) => {
233 // The job's Docker Engine starting, from its own thread.
234 for note in crate::docker::take_notes() {
235 log.line(&note);
236 }
237 log.tick();
238 }
GitHub Actions on g1t, part two: running workflows239 Err(mpsc::RecvTimeoutError::Disconnected) => break,
240 }
241 if Instant::now() >= deadline {
242 timed_out = true;
Merge Actions runs: summaries, attempts and re-runs, graceful cancel, log downloads, badges (actions 0009)243 signal(&child, "KILL");
GitHub Actions on g1t, part two: running workflows244 let _ = child.kill();
245 break;
246 }
Merge Actions runs: summaries, attempts and re-runs, graceful cancel, log downloads, badges (actions 0009)247 // Cancelled: SIGINT, then SIGTERM, then killed, as on GitHub.
248 match interrupted {
249 None if interrupt() => {
250 log.line("##[error]The operation was canceled.");
251 signal(&child, "INT");
252 interrupted = Some((Instant::now(), 1));
253 }
254 Some((at, 1)) if at.elapsed() >= INTERRUPT_GRACE => {
255 signal(&child, "TERM");
256 interrupted = Some((Instant::now(), 2));
257 }
258 Some((at, 2)) if at.elapsed() >= TERMINATE_GRACE => {
259 signal(&child, "KILL");
260 let _ = child.kill();
261 break;
262 }
263 _ => {}
264 }
265 if interrupted.is_some() && matches!(child.try_wait(), Ok(Some(_))) {
266 // The step is gone; what it started may still hold its output.
267 signal(&child, "KILL");
268 break;
269 }
GitHub Actions on g1t, part two: running workflows270 }
271 let status = child.wait()?;
272 for reader in readers {
273 let _ = reader.join();
274 }
275 // Whatever arrived after the readers finished.
276 while let Ok(line) = lines.try_recv() {
277 if let Some(shown) = commands.handle(&line, log) {
278 log.line(&shown);
279 }
280 }
281 if timed_out {
282 return Ok(Ended::TimedOut);
283 }
Merge Actions runs: summaries, attempts and re-runs, graceful cancel, log downloads, badges (actions 0009)284 if interrupted.is_some() {
285 return Ok(Ended::Cancelled);
286 }
GitHub Actions on g1t, part two: running workflows287 Ok(Ended::Exited(status.code().unwrap_or(1)))
288}
289
290#[cfg(test)]
291mod tests {
292 use super::parse_command;
293
Merge Actions runs: summaries, attempts and re-runs, graceful cancel, log downloads, badges (actions 0009)294 /// A cancelled step hears SIGINT, and its own trap runs.
295 #[cfg(unix)]
296 #[test]
297 fn a_cancelled_step_is_interrupted_and_may_clean_up() {
298 use super::{Commands, Ended, run_until};
299 use crate::actions::report::{Api, Log};
300 use std::process::Command;
301 use std::time::{Duration, Instant};
302
303 let api = Api { base: "http://127.0.0.1:9".into(), job: "job_1".into(), token: "t".into() };
304 let mut log = Log::new(api, Vec::new());
305 let mut command = Command::new("sh");
306 command.args(["-c", "trap 'echo cleaned up; exit 3' INT; echo started; while true; do sleep 0.1; done"]);
307 let began = Instant::now();
308 let ask = move || began.elapsed() >= Duration::from_millis(600);
309 let ended = run_until(command, Duration::from_secs(60), &mut log, &mut Commands::default(), &ask).unwrap();
310 assert!(matches!(ended, Ended::Cancelled));
311 assert!(began.elapsed() < Duration::from_secs(8), "SIGINT ended it, not the kill after the grace period");
312 let text = log.buffered();
313 assert!(text.contains("The operation was canceled."), "{text}");
314 assert!(text.contains("cleaned up"), "{text}");
315 }
316
GitHub Actions on g1t, part two: running workflows317 #[test]
318 fn commands_are_read_with_their_properties() {
319 let (name, properties, data) = parse_command("::error file=app.js,line=10,title=Bad%3A thing::Something%0Abroke").unwrap();
320 assert_eq!(name, "error");
321 assert_eq!(properties["file"], "app.js");
322 assert_eq!(properties["line"], "10");
323 assert_eq!(properties["title"], "Bad: thing");
324 assert_eq!(data, "Something\nbroke");
325 let (name, properties, data) = parse_command("::group::Install").unwrap();
326 assert_eq!((name.as_str(), properties.len(), data.as_str()), ("group", 0, "Install"));
327 assert_eq!(parse_command("::set-output name=version::1.2.3").unwrap().1["name"], "version");
328 assert!(parse_command("plain text").is_none());
329 assert!(parse_command(":: not a command").is_none());
330 }
331}

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