Skip to content
2,241 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//! A run's life: made, its jobs waiting on the jobs they need, each job
2//! skipped or expanded into its matrix and queued, started in a sandbox
3//! when its workspace has room, reporting its steps and logs as it goes,
4//! and finished; the run finishes with its last job.
5
6use g1t_actions::events::{RunInfo, runner_context};
7use g1t_actions::expr::{self, Scope, Status};
8use g1t_actions::matrix;
9use g1t_actions::workflow::{self, Workflow};
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look10use g1t_contracts::access::Capability;
Merge branch 'worktree-agent-a3abfcce648e87dca'11use g1t_contracts::actions::{JobCallArgs, RunActionArgs, RunApproval, StartJobArgs, WorkflowRun};
12use g1t_contracts::identity::{CreateJobTokenArgs, CreatedAccessToken, RevokeJobTokensArgs};
GitHub Actions on g1t, part two: running workflows13use g1t_contracts::repos::{Repo, RepoPath};
14use g1t_contracts::time::rfc3339;
15use g1t_contracts::{FailureCode, Outcome, new_id};
16use g1t_kit::now_ms;
17use g1t_secrets::{random_hex, same, sha256_hex};
18use serde::Deserialize;
19use serde_json::{Map, Value, json};
20use worker::Result;
21
Merge branch 'worktree-agent-a3abfcce648e87dca'22use crate::protection::Gate;
GitHub Actions on g1t, part two: running workflows23use crate::sync::WorkflowRow;
Fast pages, required checks on the branch, self-hosted runners, honest incidents24use crate::{Actions, Count, MAX_TIMEOUT_MINUTES, RUNNING_PER_WORKSPACE, SELF_HOSTED_MAX_TIMEOUT_MINUTES, SILENT_MS, SITE, check, fail, optional, repo_path};
25use g1t_contracts::runners::{Wanted, waiting_reason};
GitHub Actions on g1t, part two: running workflows26
27/// The most log one job keeps, in bytes; past it, the log says so and stops.
28const MAX_LOG_BYTES: usize = 4 * 1024 * 1024;
29/// The most a single log report may add.
30const MAX_CHUNK_BYTES: usize = 256 * 1024;
31const MAX_ANNOTATIONS: usize = 50;
32
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look33/// The repository whose unfinished runs `event` stops, by id: one deleted,
34/// or archived (not unarchived).
35pub fn stops_runs(event: &g1t_contracts::events::Event) -> Option<String> {
36 let stops = match event.kind.as_str() {
37 "repo.deleted" => true,
38 "repo.archived" => event.data["archived"].as_bool() != Some(false),
39 _ => false,
40 };
41 if !stops {
42 return None;
43 }
44 event.repo_id.clone().or_else(|| event.data["repoId"].as_str().map(str::to_owned))
45}
46
GitHub Actions on g1t, part two: running workflows47pub struct NewRun {
48 pub repo: Repo,
49 pub path: String,
50 pub source: String,
51 pub workflow: Workflow,
52 pub info: RunInfo,
53 pub action: Option<String>,
54 pub pull: Option<u32>,
55 pub title: String,
56 pub inputs: Map<String, Value>,
57 pub event_key: String,
58 pub actor_id: Option<String>,
59 pub actor: Option<String>,
60 pub trusted: bool,
Merge branch 'worktree-agent-a3abfcce648e87dca'61 /// Why the run waits for someone with the Write role to approve it
62 /// before anything starts: a pull request from outside (protection.rs).
63 pub approval: Option<String>,
GitHub Actions on g1t, part two: running workflows64}
65
66#[derive(Clone, Deserialize)]
67pub struct RunRow {
68 pub id: String,
69 pub workflow_id: String,
70 pub repo_id: String,
71 pub repo: String,
72 pub path: String,
73 pub name: String,
74 pub title: String,
75 pub number: u64,
76 pub attempt: u64,
77 pub event: String,
78 pub action: Option<String>,
79 pub git_ref: String,
80 pub sha: String,
81 pub pull: Option<u32>,
82 pub status: String,
83 pub conclusion: Option<String>,
84 pub error: Option<String>,
85 pub actor: Option<String>,
86 pub actor_id: Option<String>,
87 pub source: String,
88 pub info: String,
89 pub inputs: String,
90 pub trusted: u32,
91 pub concurrency_group: Option<String>,
92 pub created_at: String,
93 pub started_at: Option<String>,
94 pub finished_at: Option<String>,
Merge branch 'worktree-agent-a3abfcce648e87dca'95 /// JSON `RunApproval`, for a run that needed approval (migration 0006).
96 #[serde(default)]
97 pub approval: Option<String>,
98 /// Whether its concurrency group cancels what it replaces.
99 #[serde(default)]
100 pub cancel_in_progress: u32,
GitHub Actions on g1t, part two: running workflows101}
102
103impl RunRow {
Merge branch 'worktree-agent-a3abfcce648e87dca'104 pub fn approval(&self) -> Option<RunApproval> {
105 self.approval.as_deref().and_then(|text| serde_json::from_str(text).ok())
106 }
107
GitHub Actions on g1t, part two: running workflows108 pub fn info(&self) -> RunInfo {
109 let mut info: RunInfo = serde_json::from_str(&self.info).unwrap_or_default();
110 info.run_id = self.id.clone();
111 info.run_number = self.number;
112 info.run_attempt = self.attempt;
113 info.workflow_path = self.path.clone();
114 if info.workflow.is_empty() {
115 info.workflow = self.name.clone();
116 }
117 info
118 }
119
120 pub fn inputs(&self) -> Map<String, Value> {
121 serde_json::from_str(&self.inputs).unwrap_or_default()
122 }
123
124 pub fn summary(&self) -> WorkflowRun {
125 WorkflowRun {
126 id: self.id.clone(),
127 workflow_id: self.workflow_id.clone(),
128 path: self.path.clone(),
129 name: self.name.clone(),
130 title: self.title.clone(),
131 number: self.number,
132 attempt: self.attempt,
133 event: self.event.clone(),
134 git_ref: self.git_ref.clone(),
135 sha: self.sha.clone(),
136 pull: self.pull,
137 status: self.status.clone(),
138 conclusion: self.conclusion.clone(),
139 error: self.error.clone(),
140 actor: self.actor.clone(),
141 created_at: self.created_at.clone(),
142 started_at: self.started_at.clone(),
143 finished_at: self.finished_at.clone(),
144 }
145 }
146}
147
148#[derive(Clone, Deserialize)]
149pub struct JobRow {
150 pub id: String,
151 pub run_id: String,
152 pub repo_id: String,
153 pub namespace: String,
154 pub key: String,
155 pub ordinal: u32,
156 pub name: String,
157 pub needs: String,
158 pub matrix: Option<String>,
159 pub status: String,
160 pub conclusion: Option<String>,
161 pub steps: String,
162 pub annotations: String,
163 pub outputs: String,
164 pub reason: Option<String>,
165 pub token_hash: Option<String>,
166 pub timeout_minutes: u32,
167 pub continue_on_error: u32,
168 pub max_parallel: Option<u32>,
Actions: reusable workflows in the repository169 /// Set for a job that calls a reusable workflow, and for that
170 /// workflow's jobs (see migration 0002).
171 pub call: Option<String>,
GitHub Actions on g1t, part two: running workflows172 pub seen_at: Option<String>,
173 pub started_at: Option<String>,
174 pub finished_at: Option<String>,
Fast pages, required checks on the branch, self-hosted runners, honest incidents175 /// For a job whose `runs-on` names self-hosted runners: what it asks
176 /// for, as a JSON array (see `g1t_contracts::runners::Wanted`), when it
177 /// started waiting, and the runner that took it (migration 0004).
178 #[serde(default)]
179 pub labels: Option<String>,
180 #[serde(default)]
181 pub queued_at: Option<String>,
182 #[serde(default)]
183 pub runner_id: Option<String>,
184 #[serde(default)]
185 pub runner_name: Option<String>,
Merge branch 'worktree-agent-a3abfcce648e87dca'186 /// The environment it names, read when its needs were done; a job its
187 /// rules hold is `pending` (migration 0006).
188 #[serde(default)]
189 pub environment: Option<String>,
190 /// Its own `concurrency` group, and whether that cancels what it replaces.
191 #[serde(default)]
192 pub concurrency_group: Option<String>,
193 #[serde(default)]
194 pub cancel_in_progress: u32,
GitHub Actions on g1t, part two: running workflows195}
196
197impl JobRow {
198 pub fn needs(&self) -> Vec<String> {
199 serde_json::from_str(&self.needs).unwrap_or_default()
200 }
Actions: reusable workflows in the repository201
202 pub fn call(&self) -> Option<Value> {
203 self.call.as_deref().and_then(|call| serde_json::from_str(call).ok())
204 }
205
206 /// For a job of a called workflow: that workflow, the job's own id in
207 /// it, and the job.
208 pub fn callee(&self) -> Option<(Workflow, workflow::Job, Value)> {
209 let call = self.call().filter(|call| call["role"] == "callee")?;
210 let called = workflow::parse(call["source"].as_str()?).ok()?;
211 let job = called.jobs.iter().find(|job| call["job"].as_str() == Some(job.id.as_str()))?.clone();
212 Some((called, job, call))
213 }
GitHub Actions on g1t, part two: running workflows214}
215
Actions: reusable workflows in the repository216/// A job to decide on: its key, its definition, and what it needs, as
217/// (name in `needs`, key of the jobs).
218type Unit = (String, workflow::Job, Vec<(String, String)>);
219
220/// How deep reusable workflows may call one another, as on GitHub.
221const MAX_CALL_DEPTH: u64 = 4;
222
GitHub Actions on g1t, part two: running workflows223/// What the jobs of one key came to, for `needs.<key>`.
224fn key_result(rows: &[&JobRow]) -> &'static str {
225 let failed = |row: &&&JobRow| row.conclusion.as_deref() == Some("failure") && row.continue_on_error == 0;
226 if rows.iter().any(|row| failed(&row)) {
227 "failure"
228 } else if rows.iter().any(|row| row.conclusion.as_deref() == Some("cancelled")) {
229 "cancelled"
230 } else if rows.iter().all(|row| row.conclusion.as_deref() == Some("skipped")) {
231 "skipped"
232 } else {
233 "success"
234 }
235}
236
Fast pages, required checks on the branch, self-hosted runners, honest incidents237/// Whether any job before `key` failed: one it needs, or one those need,
238/// however far back. A job after a skipped one still sees the failure
239/// that skipped it, as GitHub's failure() does. `needs_of`: each key's
240/// needs, as keys; `failed`: whether a key's jobs came to a failure.
241fn ancestor_failed(needs_of: &std::collections::HashMap<&str, Vec<&str>>, key: &str, failed: impl Fn(&str) -> bool) -> bool {
242 let mut seen = std::collections::HashSet::new();
243 let mut stack: Vec<&str> = needs_of.get(key).cloned().unwrap_or_default();
244 while let Some(next) = stack.pop() {
245 if !seen.insert(next) {
246 continue;
247 }
248 if failed(next) {
249 return true;
250 }
251 stack.extend(needs_of.get(next).into_iter().flatten().copied());
252 }
253 false
254}
255
Merge branch 'main' into worktree-agent-a69aeabc4b0deeb97256/// A job's `environment:`, as deployments read it: its name, the address in
257/// `environment.url` once its expressions are filled in from the run, and
258/// whether the job deploys to it. A job with `deployment: false` only reads
259/// the environment's secrets and variables, and makes no deployment.
260#[derive(Clone, Debug, Default, PartialEq)]
261pub(crate) struct JobEnvironment {
262 pub(crate) name: String,
263 pub(crate) url: Option<String>,
264 pub(crate) deploys: bool,
265}
266
267/// The environment `environment:` names, as written in `raw` (a job), with
268/// `contexts` to fill in expressions: null when it names none, or only
269/// through an expression this cannot read before the job runs.
270pub(crate) fn environment_of(raw: &Value, contexts: &Map<String, Value>) -> Option<JobEnvironment> {
271 let scope = Scope { contexts, status: Status::Success, hash_files: None };
272 let plain = |value: &Value| -> Option<String> {
273 let text = match value {
274 Value::String(text) if text.contains("${{") => expr::interpolate_value(value, &scope).ok().map(|v| expr::to_text(&v))?,
275 Value::String(text) => text.clone(),
276 _ => return None,
277 };
278 let text = text.trim().to_owned();
279 (!text.is_empty()).then_some(text)
280 };
281 match raw.get("environment")? {
282 name @ Value::String(_) => Some(JobEnvironment { name: plain(name)?, url: None, deploys: true }),
283 Value::Object(env) => {
284 let name = plain(env.get("name")?)?;
285 let url = env
286 .get("url")
287 .and_then(plain)
288 .filter(|url| url.starts_with("https://") || url.starts_with("http://"));
289 let deploys = !matches!(env.get("deployment"), Some(Value::Bool(false)))
290 && !matches!(env.get("deployment"), Some(Value::String(text)) if text.trim() == "false");
291 Some(JobEnvironment { name, url, deploys })
292 }
293 _ => None,
294 }
295}
296
297/// The contexts a job's `runs-on` and `environment` are read with before it
298/// runs: the run's `github`, its inputs, and the job's matrix.
299fn start_contexts(run: &RunRow, job: &JobRow, key: &str) -> Map<String, Value> {
300 let mut contexts = Map::new();
301 contexts.insert("github".into(), run.info().context(key, "", run.action.as_deref()));
302 contexts.insert("inputs".into(), Value::Object(run.inputs()));
303 contexts.insert("matrix".into(), job.matrix.as_deref().and_then(|m| serde_json::from_str(m).ok()).unwrap_or_else(|| json!({})));
304 contexts
305}
306
307/// The job as its workflow (or the workflow it calls) defines it.
308fn job_spec(run: &RunRow, job: &JobRow) -> Option<workflow::Job> {
309 match job.callee() {
310 Some((_, spec, _)) => Some(spec),
311 None => workflow::parse(&run.source).ok().and_then(|w| w.jobs.into_iter().find(|j| j.id == job.key)),
312 }
313}
314
Merge branch 'worktree-agent-a3abfcce648e87dca'315/// The environment a job deploys to, if it deploys: by the name read when
316/// its needs were done (which may have needed their outputs), else as
317/// its `environment:` reads now.
Merge branch 'main' into worktree-agent-a69aeabc4b0deeb97318pub(crate) fn deploys_to(run: &RunRow, job: &JobRow) -> Option<JobEnvironment> {
319 let spec = job_spec(run, job)?;
Merge branch 'worktree-agent-a3abfcce648e87dca'320 let read = environment_of(&spec.raw, &start_contexts(run, job, &spec.id));
321 let env = match (read, &job.environment) {
322 (Some(env), Some(name)) => Some(JobEnvironment { name: name.clone(), ..env }),
323 (Some(env), None) => Some(env),
324 (None, Some(name)) => {
325 let deployment = spec.raw.get("environment").and_then(|env| env.get("deployment"));
326 let deploys = !matches!(deployment, Some(Value::Bool(false)))
327 && !matches!(deployment, Some(Value::String(text)) if text.trim() == "false");
328 Some(JobEnvironment { name: name.clone(), url: None, deploys })
329 }
330 (None, None) => None,
331 };
332 env.filter(|env| env.deploys)
Merge branch 'main' into worktree-agent-a69aeabc4b0deeb97333}
334
Fast pages, required checks on the branch, self-hosted runners, honest incidents335/// What the runner is told about a job it starts, from its workflow: the
336/// environment it names plainly, and the machine its `runs-on` asks for
337/// (`instance_for`; none for the standard one).
338#[derive(Default)]
339struct StartDetails {
340 environment: Option<String>,
341 instance: Option<String>,
342}
343
344fn start_details(run: &RunRow, job: &JobRow) -> StartDetails {
345 let spec = match job.callee() {
346 Some((_, spec, _)) => spec,
347 None => match workflow::parse(&run.source).ok().and_then(|w| w.jobs.into_iter().find(|j| j.id == job.key)) {
348 Some(spec) => spec,
349 None => return StartDetails::default(),
350 },
351 };
Merge branch 'worktree-agent-a3abfcce648e87dca'352 // As read when its needs were done, an expression's included.
353 let environment = job.environment.clone();
Fast pages, required checks on the branch, self-hosted runners, honest incidents354 // `runs-on` as the job was queued with: its matrix and the run's
355 // inputs. A label that needs more than those is the standard machine.
Merge branch 'main' into worktree-agent-a69aeabc4b0deeb97356 let contexts = start_contexts(run, job, &spec.id);
Fast pages, required checks on the branch, self-hosted runners, honest incidents357 let scope = Scope { contexts: &contexts, status: Status::Success, hash_files: None };
358 let labels: Vec<String> = match expr::interpolate_value(&spec.runs_on, &scope).unwrap_or(Value::Null) {
359 Value::String(label) => vec![label],
360 Value::Array(labels) => labels.iter().map(expr::to_text).collect(),
361 Value::Object(given) => match given.get("labels") {
362 Some(Value::Array(labels)) => labels.iter().map(expr::to_text).collect(),
363 Some(label) => vec![expr::to_text(label)],
364 None => Vec::new(),
365 },
366 _ => Vec::new(),
367 };
368 let instance = g1t_contracts::actions::instance_for(&labels);
369 StartDetails {
370 environment,
371 instance: (instance != g1t_contracts::actions::STANDARD_INSTANCE).then(|| instance.label.to_owned()),
372 }
373}
374
GitHub Actions on g1t, part two: running workflows375fn now() -> String {
376 rfc3339(now_ms())
377}
378
Merge branch 'main' into worktree-agent-a69aeabc4b0deeb97379/// How a finished run went for one environment, from the conclusions of
380/// the jobs that deploy to it: failed if any failed, an error if any was
381/// cancelled, else a success; nothing when none ran (all skipped).
382pub(crate) fn deployment_outcome(conclusions: &[Option<String>]) -> Option<&'static str> {
383 let ran: Vec<&str> = conclusions.iter().flatten().map(String::as_str).filter(|c| *c != "skipped").collect();
384 if ran.is_empty() {
385 return None;
386 }
387 if ran.iter().any(|c| *c == "failure" || *c == "timed_out") {
388 return Some("failure");
389 }
390 if ran.contains(&"cancelled") {
391 return Some("error");
392 }
393 Some("success")
394}
395
GitHub Actions on g1t, part two: running workflows396impl Actions {
397 pub async fn run_row(&self, id: &str) -> Result<Option<RunRow>> {
398 self.db.prepare("SELECT * FROM runs WHERE id = ?").bind(&[id.into()])?.first::<RunRow>(None).await
399 }
400
401 pub async fn job_rows(&self, run_id: &str) -> Result<Vec<JobRow>> {
402 self.db
403 .prepare("SELECT * FROM jobs WHERE run_id = ? ORDER BY rowid")
404 .bind(&[run_id.into()])?
405 .all()
406 .await?
407 .results::<JobRow>()
408 }
409
410 pub async fn run_summary(&self, id: &str) -> Result<Outcome<WorkflowRun>> {
411 Ok(match self.run_row(id).await? {
412 Some(row) => Outcome::Ok(row.summary()),
413 None => fail(FailureCode::NotFound, "No such run."),
414 })
415 }
416
417 /// The contexts every expression outside a job's steps may use.
418 fn base_contexts(run: &RunRow, vars: &Map<String, Value>, job: &str) -> Map<String, Value> {
419 let mut contexts = Map::new();
420 contexts.insert("github".into(), run.info().context(job, "", run.action.as_deref()));
421 contexts.insert("inputs".into(), Value::Object(run.inputs()));
422 contexts.insert("vars".into(), Value::Object(vars.clone()));
423 contexts.insert("needs".into(), json!({}));
424 contexts.insert("runner".into(), runner_context());
425 contexts
426 }
427
428 /// Makes a run and its jobs, and starts what can start. `None` when the
429 /// event already started this workflow.
430 pub async fn create_run(&self, new: NewRun) -> Result<Option<String>> {
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look431 // Nothing starts on an archived repository. A deleted one is never
432 // found to start anything on.
433 if new.repo.archived() {
434 return Ok(None);
435 }
GitHub Actions on g1t, part two: running workflows436 let workflow_row = self.workflow_row(&new.repo, &new.path, &new.workflow.display_name(&new.path), &new.source).await?;
437 let numbered = self
438 .db
439 .prepare("UPDATE workflows SET run_count = run_count + 1 WHERE id = ? RETURNING run_count AS n")
440 .bind(&[workflow_row.id.as_str().into()])?
441 .first::<Count>(None)
442 .await?
443 .map_or(1, |count| count.n);
444 let id = new_id("run", now_ms());
445 let mut info = new.info.clone();
446 info.workflow = new.workflow.display_name(&new.path);
447 info.workflow_path = new.path.clone();
448 info.run_id = id.clone();
449 info.run_number = u64::from(numbered);
Secrets and variables: one list, rows per environment, for workflows and deployments450 let vars = self
451 .variables_for(&new.repo.id, &format!("{}/{}", new.repo.namespace, new.repo.name), None, new.trusted)
452 .await?;
GitHub Actions on g1t, part two: running workflows453
454 // run-name and the concurrency group read github, inputs and vars.
455 let mut contexts = Map::new();
456 contexts.insert("github".into(), info.context("", "", new.action.as_deref()));
457 contexts.insert("inputs".into(), Value::Object(new.inputs.clone()));
458 contexts.insert("vars".into(), Value::Object(vars));
459 let scope = Scope {
460 contexts: &contexts,
461 status: Status::Success,
462 hash_files: None,
463 };
464 let title = new
465 .workflow
466 .run_name
467 .as_deref()
468 .and_then(|run_name| expr::interpolate(run_name, &scope).ok())
469 .filter(|title| !title.trim().is_empty())
470 .unwrap_or(new.title.clone());
471 let group = new.workflow.concurrency.as_ref().and_then(|c| expr::interpolate(&c.group, &scope).ok());
472 let cancel_in_progress = new
473 .workflow
474 .concurrency
475 .as_ref()
476 .and_then(|c| expr::interpolate_value(&c.cancel_in_progress, &scope).ok())
477 .is_some_and(|value| expr::truthy(&value));
478
Merge branch 'worktree-agent-a3abfcce648e87dca'479 let approval = new.approval.as_ref().map(|reason| RunApproval { state: "required".into(), reason: reason.clone(), approved_by: None });
GitHub Actions on g1t, part two: running workflows480 let inserted = self
481 .db
482 .prepare(
483 "INSERT OR IGNORE INTO runs (id, workflow_id, repo_id, repo, path, name, title, number, event, action, git_ref, sha,
Merge branch 'worktree-agent-a3abfcce648e87dca'484 pull, status, actor, actor_id, source, info, inputs, trusted, concurrency_group, event_key, created_at,
485 approval, cancel_in_progress)
486 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) RETURNING id",
GitHub Actions on g1t, part two: running workflows487 )
488 .bind(&[
489 id.as_str().into(),
490 workflow_row.id.as_str().into(),
491 new.repo.id.as_str().into(),
492 format!("{}/{}", new.repo.namespace, new.repo.name).into(),
493 new.path.as_str().into(),
494 info.workflow.as_str().into(),
495 title.as_str().into(),
496 numbered.into(),
497 info.event_name.as_str().into(),
498 optional(new.action.as_deref()),
499 info.git_ref.as_str().into(),
500 info.sha.as_str().into(),
501 new.pull.map_or(worker::wasm_bindgen::JsValue::NULL, Into::into),
Merge branch 'worktree-agent-a3abfcce648e87dca'502 if approval.is_some() { "action_required" } else { "queued" }.into(),
GitHub Actions on g1t, part two: running workflows503 optional(new.actor.as_deref()),
504 optional(new.actor_id.as_deref()),
505 new.source.as_str().into(),
506 serde_json::to_string(&info)?.into(),
507 serde_json::to_string(&new.inputs)?.into(),
508 u32::from(new.trusted).into(),
509 optional(group.as_deref()),
510 new.event_key.as_str().into(),
511 now().into(),
Merge branch 'worktree-agent-a3abfcce648e87dca'512 approval.as_ref().map(serde_json::to_string).transpose()?.as_deref().map_or(worker::wasm_bindgen::JsValue::NULL, Into::into),
513 u32::from(cancel_in_progress).into(),
GitHub Actions on g1t, part two: running workflows514 ])?
515 .first::<Value>(None)
516 .await?;
517 if inserted.is_none() {
518 return Ok(None);
519 }
520 if let Some(run) = self.run_row(&id).await? {
521 self.report_pending(&run).await?;
522 }
523
524 // Every job, waiting; each is expanded when the jobs it needs are done.
525 let mut statements = Vec::new();
526 for job in &new.workflow.jobs {
527 statements.push(
528 self.db
529 .prepare(
530 "INSERT INTO jobs (id, run_id, repo_id, namespace, key, name, needs, status) VALUES (?, ?, ?, ?, ?, ?, ?, 'waiting')",
531 )
532 .bind(&[
533 new_id("job", now_ms()).into(),
534 id.as_str().into(),
535 new.repo.id.as_str().into(),
536 new.repo.namespace.as_str().into(),
537 job.id.as_str().into(),
538 job.name.clone().filter(|n| !expr::has_expression(n)).unwrap_or(job.id.clone()).into(),
539 serde_json::to_string(&job.needs)?.into(),
540 ])?,
541 );
542 }
543 self.db.batch(statements).await?;
544
Merge branch 'worktree-agent-a3abfcce648e87dca'545 // A pull request's run from outside waits to be approved; it joins
546 // its concurrency group once it is (protection.rs, approve_run).
547 if approval.is_some() {
548 return Ok(Some(id));
549 }
550 if let Some(run) = self.run_row(&id).await? {
551 self.enter_group(&run).await?;
552 }
553 Ok(Some(id))
554 }
555
556 /// A run about to start: one run at a time per concurrency group, as on
557 /// GitHub. A newer run replaces one that waits in the group, and with
558 /// `cancel-in-progress` the one that runs; otherwise it waits as
559 /// `pending` for the one that runs. Then it moves along.
560 pub(crate) async fn enter_group(&self, run: &RunRow) -> Result<()> {
561 if let Some(group) = &run.concurrency_group {
562 let cancel_in_progress = run.cancel_in_progress != 0;
GitHub Actions on g1t, part two: running workflows563 let others = self
564 .db
Merge branch 'worktree-agent-a3abfcce648e87dca'565 .prepare(
566 "SELECT * FROM runs WHERE repo_id = ? AND concurrency_group = ? AND id != ? AND status NOT IN ('completed', 'action_required') ORDER BY id",
567 )
568 .bind(&[run.repo_id.as_str().into(), group.as_str().into(), run.id.as_str().into()])?
GitHub Actions on g1t, part two: running workflows569 .all()
570 .await?
571 .results::<RunRow>()?;
572 for other in &others {
573 if cancel_in_progress || other.status == "pending" {
574 // A newer run replaces a waiting one, as on GitHub.
575 self.cancel_run(other, "A newer run in the same concurrency group replaced it.").await?;
576 }
577 }
578 if !cancel_in_progress && others.iter().any(|other| other.status != "pending") {
579 self.db
580 .prepare("UPDATE runs SET status = 'pending' WHERE id = ?")
Merge branch 'worktree-agent-a3abfcce648e87dca'581 .bind(&[run.id.as_str().into()])?
GitHub Actions on g1t, part two: running workflows582 .run()
583 .await?;
Merge branch 'worktree-agent-a3abfcce648e87dca'584 return Ok(());
GitHub Actions on g1t, part two: running workflows585 }
586 }
Merge branch 'worktree-agent-a3abfcce648e87dca'587 self.advance(&run.id).await
GitHub Actions on g1t, part two: running workflows588 }
589
590 /// A run that could not start, such as for a workflow file that does not read.
591 #[allow(clippy::too_many_arguments)]
592 pub async fn record_failed_run(
593 &self,
594 row: &WorkflowRow,
595 git_ref: &str,
596 sha: &str,
597 event_key: &str,
598 actor_id: Option<&str>,
599 actor: &str,
600 problem: &str,
601 ) -> Result<()> {
602 let numbered = self
603 .db
604 .prepare("UPDATE workflows SET run_count = run_count + 1 WHERE id = ? RETURNING run_count AS n")
605 .bind(&[row.id.as_str().into()])?
606 .first::<Count>(None)
607 .await?
608 .map_or(1, |count| count.n);
609 let at = now();
610 self.db
611 .prepare(
612 "INSERT OR IGNORE INTO runs (id, workflow_id, repo_id, repo, path, name, title, number, event, git_ref, sha, status,
613 conclusion, error, actor, actor_id, source, info, event_key, created_at, finished_at)
614 VALUES (?, ?, ?, ?, ?, ?, ?, ?, 'push', ?, ?, 'completed', 'failure', ?, ?, ?, ?, '{}', ?, ?, ?)",
615 )
616 .bind(&[
617 new_id("run", now_ms()).into(),
618 row.id.as_str().into(),
619 row.repo_id.as_str().into(),
620 row.repo.as_str().into(),
621 row.path.as_str().into(),
622 row.path.as_str().into(),
623 "Invalid workflow file".into(),
624 numbered.into(),
625 git_ref.into(),
626 sha.into(),
627 problem.into(),
628 actor.into(),
629 optional(actor_id),
630 row.source.as_str().into(),
631 event_key.into(),
632 at.as_str().into(),
633 at.as_str().into(),
634 ])?
635 .run()
636 .await?;
637 Ok(())
638 }
639
640 /// Moves a run along: jobs whose needs are done are decided on, and
641 /// jobs that can start are started.
642 pub async fn advance(&self, run_id: &str) -> Result<()> {
643 // Each pass may finish jobs (skipped ones), which may free others.
644 for _ in 0..20 {
645 let Some(run) = self.run_row(run_id).await? else { return Ok(()) };
Merge branch 'worktree-agent-a3abfcce648e87dca'646 if matches!(run.status.as_str(), "completed" | "pending" | "action_required") {
GitHub Actions on g1t, part two: running workflows647 return Ok(());
648 }
649 let workflow = match workflow::parse(&run.source) {
650 Ok(workflow) => workflow,
651 Err(problem) => {
652 self.finish_run(&run, Some(&problem)).await?;
653 return Ok(());
654 }
655 };
656 let jobs = self.job_rows(run_id).await?;
657 let mut changed = false;
Actions: reusable workflows in the repository658 // What to decide on: the workflow's jobs, and the jobs of the
659 // workflows they call, each with its needs as (name, key).
660 let mut units: Vec<Unit> = workflow
661 .jobs
662 .iter()
663 .map(|job| (job.id.clone(), job.clone(), job.needs.iter().map(|n| (n.clone(), n.clone())).collect()))
664 .collect();
665 let mut seen = std::collections::HashSet::new();
666 for row in &jobs {
667 if !seen.insert(row.key.clone()) {
668 continue;
669 }
670 if let Some((_, job, call)) = row.callee() {
671 let parent = call["parent"].as_str().unwrap_or_default().to_owned();
672 let needs = job.needs.iter().map(|n| (n.clone(), format!("{parent}/{n}"))).collect();
673 units.push((row.key.clone(), job, needs));
674 }
675 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents676 // Each key's needs, as keys, to look back through for failure().
677 let needs_of: std::collections::HashMap<&str, Vec<&str>> =
678 units.iter().map(|(key, _, needs)| (key.as_str(), needs.iter().map(|(_, need)| need.as_str()).collect())).collect();
679 let key_failed = |key: &str| key_result(&jobs.iter().filter(|row| row.key == key).collect::<Vec<_>>()) == "failure";
Actions: reusable workflows in the repository680 for (key, job, needs) in &units {
681 let rows: Vec<&JobRow> = jobs.iter().filter(|row| &row.key == key).collect();
GitHub Actions on g1t, part two: running workflows682 if rows.is_empty() || !rows.iter().all(|row| row.status == "waiting") {
683 continue;
684 }
685 let needed: Vec<(&String, Vec<&JobRow>)> =
Actions: reusable workflows in the repository686 needs.iter().map(|(name, need)| (name, jobs.iter().filter(|row| &row.key == need).collect())).collect();
GitHub Actions on g1t, part two: running workflows687 if !needed.iter().all(|(_, rows)| rows.iter().all(|row| row.status == "completed")) {
688 continue;
689 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents690 let failed_before = ancestor_failed(&needs_of, key, key_failed);
691 self.decide(&run, job, rows[0], &needed, failed_before).await?;
GitHub Actions on g1t, part two: running workflows692 changed = true;
693 }
Actions: reusable workflows in the repository694 // A job that called a workflow finishes with that workflow's jobs.
695 for row in jobs.iter().filter(|row| row.status == "calling") {
696 let children: Vec<&JobRow> = jobs
697 .iter()
698 .filter(|child| child.call().is_some_and(|call| call["role"] == "callee" && call["parent"].as_str() == Some(row.key.as_str())))
699 .collect();
700 if !children.is_empty() && children.iter().all(|child| child.status == "completed") {
701 self.finish_call(row, &children).await?;
702 changed = true;
703 }
704 }
GitHub Actions on g1t, part two: running workflows705 if !changed {
706 break;
707 }
708 }
709 self.start_queued().await?;
710 self.finish_if_done(run_id).await
711 }
712
713 /// Decides on one job whose needs are done: skip it, fail it, or expand
Fast pages, required checks on the branch, self-hosted runners, honest incidents714 /// it into its matrix and queue it. `failed_before`: whether any job
715 /// before it failed, however far back (`ancestor_failed`).
716 async fn decide(&self, run: &RunRow, job: &workflow::Job, row: &JobRow, needed: &[(&String, Vec<&JobRow>)], failed_before: bool) -> Result<()> {
Secrets and variables: one list, rows per environment, for workflows and deployments717 let vars = self.variables_for(&run.repo_id, &run.repo, None, run.trusted != 0).await?;
GitHub Actions on g1t, part two: running workflows718 let mut contexts = Self::base_contexts(run, &vars, &job.id);
Actions: reusable workflows in the repository719 // A called workflow's jobs read the inputs they were called with.
720 let call = row.call();
721 let parent = call.as_ref().filter(|c| c["role"] == "callee").and_then(|c| c["parent"].as_str().map(str::to_owned));
722 if let Some(call) = call.as_ref().filter(|c| c["role"] == "callee") {
723 contexts.insert("inputs".into(), call["inputs"].clone());
724 }
GitHub Actions on g1t, part two: running workflows725 let mut needs = Map::new();
Fast pages, required checks on the branch, self-hosted runners, honest incidents726 let mut results = Vec::new();
GitHub Actions on g1t, part two: running workflows727 for (key, rows) in needed {
728 let result = key_result(rows);
729 let mut outputs = Map::new();
730 for row in rows {
731 if let Ok(Value::Object(more)) = serde_json::from_str::<Value>(&row.outputs) {
732 outputs.extend(more);
733 }
734 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents735 results.push(result);
GitHub Actions on g1t, part two: running workflows736 needs.insert((*key).clone(), json!({ "result": result, "outputs": outputs }));
737 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents738 // As on GitHub: a need that was skipped makes success() false but
739 // not failure(); failure() is a failure anywhere before the job.
740 let status = expr::job_status(results, failed_before, run.conclusion.as_deref() == Some("cancelled"));
GitHub Actions on g1t, part two: running workflows741 contexts.insert("needs".into(), Value::Object(needs));
742 let scope = Scope {
743 contexts: &contexts,
744 status,
745 hash_files: None,
746 };
747 let condition = job.condition.as_deref().unwrap_or_default();
748 match expr::condition(condition, &scope) {
749 Ok(true) => {}
750 Ok(false) => return self.skip_job(row, None).await,
751 Err(problem) => return self.fail_job(row, &format!("Its `if` does not read: {problem}")).await,
752 }
Actions: reusable workflows in the repository753 if let Some(uses) = &job.uses {
754 return self.call_workflow(run, job, row, uses, &scope).await;
GitHub Actions on g1t, part two: running workflows755 }
756
757 // Its matrix, which may come from a needed job's outputs.
758 let combinations = match &job.matrix {
759 None => vec![Map::new()],
760 Some(matrix) => {
761 let value = match expr::interpolate_value(matrix, &scope) {
762 Ok(value) => value,
763 Err(problem) => return self.fail_job(row, &format!("Its matrix does not read: {problem}")).await,
764 };
765 match matrix::expand(&value) {
766 Ok(combinations) if !combinations.is_empty() => combinations,
767 Ok(_) => return self.fail_job(row, "Its matrix makes no jobs.").await,
768 Err(problem) => return self.fail_job(row, &problem).await,
769 }
770 }
771 };
772 let total = combinations.len();
773 let raw = job.raw.as_object().cloned().unwrap_or_default();
Fast pages, required checks on the branch, self-hosted runners, honest incidents774 // Whether pull requests from forks may use self-hosted runners here,
775 // asked once, and only for a run that is not trusted.
776 let mut forks_allowed: Option<bool> = None;
Merge branch 'worktree-agent-a3abfcce648e87dca'777 // Each concurrency group its jobs join, and whether it cancels.
778 let mut groups: Vec<(String, bool)> = Vec::new();
GitHub Actions on g1t, part two: running workflows779 let mut statements = Vec::new();
780 for (index, combination) in combinations.iter().enumerate() {
781 let mut contexts = contexts.clone();
782 contexts.insert("matrix".into(), Value::Object(combination.clone()));
783 contexts.insert(
784 "strategy".into(),
785 json!({ "fail-fast": job.fail_fast, "job-index": index, "job-total": total, "max-parallel": job.max_parallel.unwrap_or(total as u32) }),
786 );
787 let scope = Scope {
788 contexts: &contexts,
789 status: Status::Success,
790 hash_files: None,
791 };
792 let base_name = job.name.clone().unwrap_or(job.id.clone());
Actions: reusable workflows in the repository793 // A called workflow's job is shown under the job that called it.
794 let base_name = match &parent {
795 Some(parent) => format!("{} / {base_name}", parent.replace('/', " / ")),
796 None => base_name,
797 };
GitHub Actions on g1t, part two: running workflows798 let name = if expr::has_expression(&base_name) {
799 expr::interpolate(&base_name, &scope).unwrap_or(base_name)
800 } else if job.matrix.is_some() {
801 matrix::job_name(&base_name, combination)
802 } else {
803 base_name
804 };
805 let runs_on = expr::interpolate_value(&job.runs_on, &scope).unwrap_or(Value::Null);
806 let labels = match &runs_on {
807 Value::String(label) => vec![label.clone()],
808 Value::Array(labels) => labels.iter().map(expr::to_text).collect(),
809 Value::Object(spec) => spec.get("labels").map(|l| match l {
810 Value::Array(labels) => labels.iter().map(expr::to_text).collect(),
811 other => vec![expr::to_text(other)],
812 }).unwrap_or_default(),
813 _ => Vec::new(),
814 };
Fast pages, required checks on the branch, self-hosted runners, honest incidents815 // `runs-on: self-hosted` (or a group): the workspace's own
816 // machines, which may run any OS. Otherwise g1t's Linux sandboxes.
817 let group = match &runs_on {
818 Value::Object(spec) => spec.get("group").map(expr::to_text),
819 _ => None,
820 };
821 let wanted = Wanted::of(&labels, group.as_deref());
822 let mut reason = if wanted.self_hosted {
823 None
824 } else {
825 labels
826 .iter()
827 .find(|label| {
828 let lower = label.to_ascii_lowercase();
829 lower.contains("windows") || lower.contains("macos")
830 })
831 .map(|label| {
832 let os = if label.to_ascii_lowercase().contains("windows") { "windows" } else { "macos" };
833 format!("`runs-on: {label}`: g1t's own runners are Linux. A self-hosted runner can run it: `runs-on: [self-hosted, {os}]`.")
834 })
835 };
836 // A pull request from a fork runs code anyone could write: never
837 // on the workspace's machines unless it said they may.
838 if wanted.self_hosted && run.trusted == 0 {
839 let allowed = match forks_allowed {
840 Some(allowed) => allowed,
841 None => {
842 let allowed = self.effective_runner_settings(&row.namespace, Some(&run.repo_id)).await?.fork_pull_requests;
843 forks_allowed = Some(allowed);
844 allowed
845 }
846 };
847 if !allowed {
848 reason = Some(
849 "Pull requests from forks do not run on self-hosted runners here. An admin can allow it under Settings, Runners.".to_owned(),
850 );
851 }
852 }
853 let max_minutes = if wanted.self_hosted { SELF_HOSTED_MAX_TIMEOUT_MINUTES } else { MAX_TIMEOUT_MINUTES };
GitHub Actions on g1t, part two: running workflows854 let timeout = raw
855 .get("timeout-minutes")
856 .and_then(|value| expr::interpolate_value(value, &scope).ok())
857 .and_then(|value| value.as_f64().or_else(|| expr::to_text(&value).parse().ok()))
Fast pages, required checks on the branch, self-hosted runners, honest incidents858 .map_or(MAX_TIMEOUT_MINUTES, |minutes| (minutes.ceil() as u32).clamp(1, max_minutes));
GitHub Actions on g1t, part two: running workflows859 let continue_on_error = raw
860 .get("continue-on-error")
861 .and_then(|value| expr::interpolate_value(value, &scope).ok())
862 .is_some_and(|value| expr::truthy(&value));
Merge branch 'worktree-agent-a3abfcce648e87dca'863 let (mut status, mut conclusion, mut finished) = match &reason {
GitHub Actions on g1t, part two: running workflows864 Some(_) => ("completed", Some("failure"), Some(now())),
865 None => ("queued", None, None),
866 };
Fast pages, required checks on the branch, self-hosted runners, honest incidents867 // A self-hosted job waits, saying for what, until a runner takes it.
868 let (labels_json, queued_at) = if wanted.self_hosted && reason.is_none() {
869 reason = Some(waiting_reason(&wanted));
870 (Some(serde_json::to_string(&wanted.stored())?), Some(now()))
871 } else {
872 (None, None)
873 };
Merge branch 'worktree-agent-a3abfcce648e87dca'874 // The environment it names, an expression read now, and what
875 // that environment's protection rules say: it starts, waits as
876 // `pending` (protection.rs), or may not deploy there.
877 let environment = environment_of(&job.raw, &contexts).map(|env| env.name);
878 if status == "queued"
879 && let Some(name) = &environment
880 {
881 match self.gate(run, name).await? {
882 Gate::Open => {}
883 Gate::Held(why) => {
884 status = "pending";
885 reason = Some(why);
886 }
887 Gate::Refused(why) => {
888 (status, conclusion, finished) = ("completed", Some("failure"), Some(now()));
889 reason = Some(why);
890 }
891 }
892 }
893 // Its own concurrency group.
894 let job_group = job
895 .concurrency
896 .as_ref()
897 .and_then(|c| expr::interpolate(&c.group, &scope).ok())
898 .map(|group| group.trim().to_owned())
899 .filter(|group| !group.is_empty());
900 let job_cancel = job
901 .concurrency
902 .as_ref()
903 .and_then(|c| expr::interpolate_value(&c.cancel_in_progress, &scope).ok())
904 .is_some_and(|value| expr::truthy(&value));
905 if status != "completed"
906 && let Some(group) = &job_group
907 && !groups.iter().any(|(known, _)| known == group)
908 {
909 groups.push((group.clone(), job_cancel));
910 }
GitHub Actions on g1t, part two: running workflows911 let values: Vec<worker::wasm_bindgen::JsValue> = vec![
912 name.into(),
913 serde_json::to_string(combination)?.into(),
914 status.into(),
915 optional(conclusion),
916 optional(reason.as_deref()),
917 timeout.into(),
918 u32::from(continue_on_error).into(),
919 job.max_parallel.map_or(worker::wasm_bindgen::JsValue::NULL, Into::into),
920 optional(finished.as_deref()),
Fast pages, required checks on the branch, self-hosted runners, honest incidents921 optional(labels_json.as_deref()),
922 optional(queued_at.as_deref()),
Merge branch 'worktree-agent-a3abfcce648e87dca'923 optional(environment.as_deref()),
924 optional(job_group.as_deref()),
925 u32::from(job_cancel).into(),
GitHub Actions on g1t, part two: running workflows926 ];
927 if index == 0 {
928 let mut bound = values;
929 bound.push(row.id.as_str().into());
930 statements.push(
931 self.db
932 .prepare(
933 "UPDATE jobs SET name = ?, matrix = ?, status = ?, conclusion = ?, reason = ?, timeout_minutes = ?,
Merge branch 'worktree-agent-a3abfcce648e87dca'934 continue_on_error = ?, max_parallel = ?, finished_at = ?, labels = ?, queued_at = ?, environment = ?,
935 concurrency_group = ?, cancel_in_progress = ? WHERE id = ?",
GitHub Actions on g1t, part two: running workflows936 )
937 .bind(&bound)?,
938 );
939 } else {
940 let mut bound: Vec<worker::wasm_bindgen::JsValue> = vec![
941 new_id("job", now_ms()).into(),
942 row.run_id.as_str().into(),
943 row.repo_id.as_str().into(),
944 row.namespace.as_str().into(),
945 row.key.as_str().into(),
946 (index as u32).into(),
947 row.needs.as_str().into(),
Actions: reusable workflows in the repository948 optional(row.call.as_deref()),
GitHub Actions on g1t, part two: running workflows949 ];
950 bound.extend(values);
951 statements.push(
952 self.db
953 .prepare(
Actions: reusable workflows in the repository954 "INSERT INTO jobs (id, run_id, repo_id, namespace, key, ordinal, needs, call, name, matrix, status, conclusion, reason,
Merge branch 'worktree-agent-a3abfcce648e87dca'955 timeout_minutes, continue_on_error, max_parallel, finished_at, labels, queued_at, environment,
956 concurrency_group, cancel_in_progress)
957 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
GitHub Actions on g1t, part two: running workflows958 )
959 .bind(&bound)?,
960 );
961 }
962 }
963 self.db.batch(statements).await?;
Merge branch 'worktree-agent-a3abfcce648e87dca'964 self.replace_in_groups(run, &row.key, &groups).await
965 }
966
967 /// A job queued in its own concurrency group: with `cancel-in-progress`
968 /// it cancels the group's other jobs; otherwise it replaces any that
969 /// are queued and not started, and waits for the one running
970 /// (`start_queued` starts one of a group at a time). As on GitHub.
971 async fn replace_in_groups(&self, run: &RunRow, key: &str, groups: &[(String, bool)]) -> Result<()> {
972 if groups.is_empty() {
973 return Ok(());
974 }
975 #[derive(Deserialize)]
976 struct Id {
977 id: String,
978 }
979 let mine: Vec<String> = self
980 .db
981 .prepare("SELECT id FROM jobs WHERE run_id = ? AND key = ?")
982 .bind(&[run.id.as_str().into(), key.into()])?
983 .all()
984 .await?
985 .results::<Id>()?
986 .into_iter()
987 .map(|row| row.id)
988 .collect();
989 let mut moved: Vec<String> = Vec::new();
990 for (group, cancel) in groups {
991 let others = self
992 .db
993 .prepare("SELECT * FROM jobs WHERE repo_id = ? AND concurrency_group = ? AND status IN ('queued', 'pending', 'in_progress')")
994 .bind(&[run.repo_id.as_str().into(), group.as_str().into()])?
995 .all()
996 .await?
997 .results::<JobRow>()?;
998 for other in others.iter().filter(|other| !mine.contains(&other.id)) {
999 if *cancel || other.status == "queued" {
1000 self.stop_job(other, "A newer job in the same concurrency group replaced it.").await?;
1001 if !moved.contains(&other.run_id) {
1002 moved.push(other.run_id.clone());
1003 }
1004 }
1005 }
1006 }
1007 for other_run in moved.iter().filter(|id| **id != run.id) {
1008 Box::pin(self.advance(other_run)).await?;
1009 }
GitHub Actions on g1t, part two: running workflows1010 Ok(())
1011 }
1012
Merge Actions: cross-repo workflows and actions, release and deployment triggers, step timeouts1013 /// A job that calls a reusable workflow, in the repository or another
1014 /// (reach.rs): that workflow's jobs join the run under it, with the
1015 /// inputs and secrets it passes.
Actions: reusable workflows in the repository1016 async fn call_workflow(&self, run: &RunRow, job: &workflow::Job, row: &JobRow, uses: &str, scope: &Scope<'_>) -> Result<()> {
1017 let depth = row.call().and_then(|c| c["depth"].as_u64()).unwrap_or(0) + 1;
1018 if depth > MAX_CALL_DEPTH {
1019 return self.fail_job(row, &format!("Reusable workflows call each other more than {MAX_CALL_DEPTH} deep.")).await;
1020 }
Merge Actions: cross-repo workflows and actions, release and deployment triggers, step timeouts1021 let (file, source, origin) = match self.called_workflow(run, row, uses).await? {
1022 Ok(found) => found,
1023 Err(why) => return self.fail_job(row, &why).await,
Actions: reusable workflows in the repository1024 };
1025 let called = match workflow::parse(&source) {
1026 Ok(called) => called,
1027 Err(problem) => return self.fail_job(row, &format!("`{file}` does not read: {problem}")).await,
1028 };
1029 let Some(trigger) = called.trigger("workflow_call") else {
1030 return self.fail_job(row, &format!("`{file}` cannot be called: it has no `on: workflow_call`.")).await;
1031 };
1032 // Inputs: what the caller passes, else the called workflow's defaults.
1033 let given = match job.raw.get("with") {
1034 Some(with) => match expr::interpolate_value(with, scope) {
1035 Ok(Value::Object(given)) => given,
1036 Ok(_) => Map::new(),
1037 Err(problem) => return self.fail_job(row, &format!("Its `with` does not read: {problem}")).await,
1038 },
1039 None => Map::new(),
1040 };
1041 let mut inputs = Map::new();
1042 for (name, spec) in &trigger.inputs {
1043 let value = given.get(name).cloned().or_else(|| spec.get("default").cloned()).unwrap_or(Value::Null);
1044 if value.is_null() && spec.get("required").and_then(Value::as_bool) == Some(true) {
1045 return self.fail_job(row, &format!("`{file}` needs the input `{name}`.")).await;
1046 }
1047 inputs.insert(name.clone(), value);
1048 }
1049 for (name, value) in given {
1050 inputs.entry(name).or_insert(value);
1051 }
Merge Actions: cross-repo workflows and actions, release and deployment triggers, step timeouts1052 // Secrets: none but the job's token unless `secrets:` passes them,
1053 // by name or with `inherit`; read when each job starts (job_spec).
1054 let outer = row.call().filter(|c| c["role"] == "callee").and_then(|c| c.get("secrets").cloned());
1055 let secrets = crate::reach::secrets_plan(&job.raw, scope.contexts, outer);
1056 if let Some(name) = crate::reach::missing_secrets(&called.raw, &secrets).first() {
1057 return self.fail_job(row, &format!("`{file}` needs the secret `{name}`: pass it under `secrets:`, or use `secrets: inherit`.")).await;
1058 }
Actions: reusable workflows in the repository1059 let mut statements = Vec::new();
1060 for called_job in &called.jobs {
1061 let needs: Vec<String> = called_job.needs.iter().map(|n| format!("{}/{n}", row.key)).collect();
1062 let call = json!({
1063 "role": "callee", "parent": row.key, "job": called_job.id, "path": file,
Merge Actions: cross-repo workflows and actions, release and deployment triggers, step timeouts1064 "source": source, "inputs": inputs, "depth": depth, "origin": origin, "secrets": secrets,
Actions: reusable workflows in the repository1065 });
1066 statements.push(
1067 self.db
1068 .prepare("INSERT INTO jobs (id, run_id, repo_id, namespace, key, name, needs, status, call) VALUES (?, ?, ?, ?, ?, ?, ?, 'waiting', ?)")
1069 .bind(&[
1070 new_id("job", now_ms()).into(),
1071 row.run_id.as_str().into(),
1072 row.repo_id.as_str().into(),
1073 row.namespace.as_str().into(),
1074 format!("{}/{}", row.key, called_job.id).into(),
1075 format!("{} / {}", row.name, called_job.name.clone().unwrap_or(called_job.id.clone())).into(),
1076 serde_json::to_string(&needs)?.into(),
1077 serde_json::to_string(&call)?.into(),
1078 ])?,
1079 );
1080 }
1081 statements.push(
1082 self.db
1083 .prepare("UPDATE jobs SET status = 'calling', call = ?, reason = ?, started_at = ? WHERE id = ?")
1084 .bind(&[
1085 serde_json::to_string(&json!({ "role": "caller", "path": file, "source": source }))?.into(),
1086 format!("Calls `{file}`.").into(),
1087 now().into(),
1088 row.id.as_str().into(),
1089 ])?,
1090 );
1091 self.db.batch(statements).await?;
1092 Ok(())
1093 }
1094
1095 /// A job that called a workflow, finished with its jobs: their result,
1096 /// and the outputs the workflow declares.
1097 async fn finish_call(&self, row: &JobRow, children: &[&JobRow]) -> Result<()> {
1098 let call = row.call().unwrap_or_default();
1099 let called = call["source"].as_str().and_then(|s| workflow::parse(s).ok());
1100 let mut jobs_context = Map::new();
1101 let mut by_key: std::collections::BTreeMap<String, Vec<&JobRow>> = std::collections::BTreeMap::new();
1102 for child in children {
1103 by_key.entry(child.key.clone()).or_default().push(child);
1104 }
1105 for (key, rows) in &by_key {
1106 let mut outputs = Map::new();
1107 for child in rows {
1108 if let Ok(Value::Object(more)) = serde_json::from_str::<Value>(&child.outputs) {
1109 outputs.extend(more);
1110 }
1111 }
1112 let id = key.rsplit('/').next().unwrap_or(key);
1113 jobs_context.insert(id.to_owned(), json!({ "result": key_result(rows), "outputs": outputs }));
1114 }
1115 let inputs = children.first().and_then(|c| c.call()).map(|c| c["inputs"].clone()).unwrap_or(json!({}));
1116 let mut contexts = Map::new();
1117 contexts.insert("jobs".into(), Value::Object(jobs_context));
1118 contexts.insert("inputs".into(), inputs);
1119 let scope = Scope { contexts: &contexts, status: Status::Success, hash_files: None };
1120 let mut outputs = Map::new();
1121 if let Some(Value::Object(declared)) = called.as_ref().map(|w| {
1122 let on = w.raw.get("on").or_else(|| w.raw.get("true")).cloned().unwrap_or(Value::Null);
1123 on.get("workflow_call").and_then(|c| c.get("outputs")).cloned().unwrap_or(Value::Null)
1124 }) {
1125 for (name, spec) in declared {
1126 if let Some(value) = spec.get("value") {
1127 let value = expr::interpolate_value(value, &scope).unwrap_or(Value::Null);
1128 outputs.insert(name, Value::String(expr::to_text(&value)));
1129 }
1130 }
1131 }
1132 let conclusion = key_result(children);
1133 self.db
1134 .prepare("UPDATE jobs SET status = 'completed', conclusion = ?, outputs = ?, finished_at = ? WHERE id = ? AND status = 'calling'")
1135 .bind(&[conclusion.into(), serde_json::to_string(&outputs)?.into(), now().into(), row.id.as_str().into()])?
1136 .run()
1137 .await?;
1138 Ok(())
1139 }
1140
1141 /// A file's text at a commit, if it is there.
Merge Actions: cross-repo workflows and actions, release and deployment triggers, step timeouts1142 pub(crate) async fn read_file(&self, path: &RepoPath, ws: &g1t_contracts::User, sha: &str, file: &str) -> Result<Option<String>> {
Actions: reusable workflows in the repository1143 let blob: Outcome<g1t_contracts::repos::BlobView> = g1t_kit::call(
1144 &self.repos,
1145 "blob",
1146 &g1t_contracts::repos::BlobArgs {
1147 path: path.clone(),
1148 viewer: Some(ws.clone()),
1149 git_ref: sha.to_owned(),
1150 file_path: file.to_owned(),
1151 },
1152 )
1153 .await?;
1154 Ok(match blob {
1155 Outcome::Ok(view) => view.text,
1156 Outcome::Fail(_) => None,
1157 })
1158 }
1159
GitHub Actions on g1t, part two: running workflows1160 async fn skip_job(&self, row: &JobRow, reason: Option<&str>) -> Result<()> {
1161 self.db
1162 .prepare("UPDATE jobs SET status = 'completed', conclusion = 'skipped', reason = ?, finished_at = ? WHERE id = ?")
1163 .bind(&[optional(reason), now().into(), row.id.as_str().into()])?
1164 .run()
1165 .await?;
1166 Ok(())
1167 }
1168
1169 async fn fail_job(&self, row: &JobRow, reason: &str) -> Result<()> {
1170 self.db
1171 .prepare("UPDATE jobs SET status = 'completed', conclusion = 'failure', reason = ?, finished_at = ? WHERE id = ? AND status != 'completed'")
1172 .bind(&[reason.into(), now().into(), row.id.as_str().into()])?
1173 .run()
1174 .await?;
1175 Ok(())
1176 }
1177
1178 /// Starts queued jobs, oldest first, while their workspace has room.
1179 pub async fn start_queued(&self) -> Result<()> {
1180 let queued = self
1181 .db
Fast pages, required checks on the branch, self-hosted runners, honest incidents1182 // Self-hosted jobs are taken by their runners (runners.rs).
1183 .prepare("SELECT * FROM jobs WHERE status = 'queued' AND labels IS NULL ORDER BY rowid LIMIT 50")
GitHub Actions on g1t, part two: running workflows1184 .all()
1185 .await?
1186 .results::<JobRow>()?;
1187 let mut running: std::collections::HashMap<String, u32> = std::collections::HashMap::new();
1188 for job in queued {
1189 let in_workspace = match running.get(&job.namespace) {
1190 Some(n) => *n,
1191 None => {
1192 let n = self
1193 .db
Fast pages, required checks on the branch, self-hosted runners, honest incidents1194 // Only g1t's own sandboxes count against the workspace's room.
1195 .prepare("SELECT COUNT(*) AS n FROM jobs WHERE status = 'in_progress' AND namespace = ? AND runner_id IS NULL")
GitHub Actions on g1t, part two: running workflows1196 .bind(&[job.namespace.as_str().into()])?
1197 .first::<Count>(None)
1198 .await?
1199 .map_or(0, |count| count.n);
1200 running.insert(job.namespace.clone(), n);
1201 n
1202 }
1203 };
1204 if in_workspace >= RUNNING_PER_WORKSPACE {
1205 continue;
1206 }
1207 if let Some(max) = job.max_parallel {
1208 let siblings = self
1209 .db
1210 .prepare("SELECT COUNT(*) AS n FROM jobs WHERE run_id = ? AND key = ? AND status = 'in_progress'")
1211 .bind(&[job.run_id.as_str().into(), job.key.as_str().into()])?
1212 .first::<Count>(None)
1213 .await?
1214 .map_or(0, |count| count.n);
1215 if siblings >= max {
1216 continue;
1217 }
1218 }
Merge branch 'worktree-agent-a3abfcce648e87dca'1219 // One job of a concurrency group runs at a time.
1220 if let Some(group) = &job.concurrency_group {
1221 let running = self
1222 .db
1223 .prepare("SELECT COUNT(*) AS n FROM jobs WHERE repo_id = ? AND concurrency_group = ? AND status = 'in_progress' AND id != ?")
1224 .bind(&[job.repo_id.as_str().into(), group.as_str().into(), job.id.as_str().into()])?
1225 .first::<Count>(None)
1226 .await?
1227 .map_or(0, |count| count.n);
1228 if running > 0 {
1229 continue;
1230 }
1231 }
GitHub Actions on g1t, part two: running workflows1232 let token = random_hex(24);
1233 let at = now();
1234 let claimed = self
1235 .db
1236 .prepare(
1237 "UPDATE jobs SET status = 'in_progress', token_hash = ?, started_at = ?, seen_at = ? WHERE id = ? AND status = 'queued' RETURNING id",
1238 )
1239 .bind(&[sha256_hex(&token).into(), at.as_str().into(), at.as_str().into(), job.id.as_str().into()])?
1240 .first::<Value>(None)
1241 .await?;
1242 if claimed.is_none() {
1243 continue;
1244 }
1245 running.insert(job.namespace.clone(), in_workspace + 1);
1246 self.db
1247 .prepare("UPDATE runs SET status = 'in_progress', started_at = COALESCE(started_at, ?) WHERE id = ? AND status = 'queued'")
1248 .bind(&[at.as_str().into(), job.run_id.as_str().into()])?
1249 .run()
1250 .await?;
1251 let run = self.run_row(&job.run_id).await?;
1252 let repo: RepoPath = run.as_ref().map(|run| repo_path(&run.repo)).unwrap_or(RepoPath {
1253 namespace: job.namespace.clone(),
1254 name: String::new(),
1255 });
Fast pages, required checks on the branch, self-hosted runners, honest incidents1256 // Its environment and the machine it asked for, for the runner.
1257 let details = run.as_ref().map(|run| start_details(run, &job)).unwrap_or_default();
GitHub Actions on g1t, part two: running workflows1258 let started: Outcome<Value> = g1t_kit::call(
1259 &self.runner,
1260 "start_actions_job",
1261 &StartJobArgs {
1262 job: job.id.clone(),
1263 token,
1264 repo,
1265 timeout_minutes: job.timeout_minutes,
Fast pages, required checks on the branch, self-hosted runners, honest incidents1266 workflow: run.as_ref().map(|run| run.path.clone()),
1267 environment: details.environment,
1268 trusted: run.as_ref().is_some_and(|run| run.trusted != 0),
1269 instance: details.instance,
GitHub Actions on g1t, part two: running workflows1270 },
1271 )
1272 .await
1273 .unwrap_or_else(|error| fail(FailureCode::Conflict, format!("The runner could not be reached: {error}")));
1274 if let Outcome::Fail(refused) = started {
1275 Box::pin(self.finish_job(&job.id, "failure", Some(&refused.message), None)).await?;
Merge branch 'main' into worktree-agent-a69aeabc4b0deeb971276 } else {
1277 self.job_started(&job.id).await?;
GitHub Actions on g1t, part two: running workflows1278 }
1279 }
1280 Ok(())
1281 }
1282
1283 /// Finishes a job and moves its run along.
1284 pub async fn finish_job(&self, job_id: &str, conclusion: &str, reason: Option<&str>, outputs: Option<&Map<String, Value>>) -> Result<()> {
1285 let finished = self
1286 .db
1287 .prepare(
1288 "UPDATE jobs SET status = 'completed', conclusion = ?, reason = COALESCE(?, reason), outputs = COALESCE(?, outputs),
1289 finished_at = ?, token_hash = NULL WHERE id = ? AND status != 'completed' RETURNING *",
1290 )
1291 .bind(&[
1292 conclusion.into(),
1293 optional(reason),
1294 outputs.map(|o| serde_json::to_string(o).unwrap_or_default()).as_deref().map_or(worker::wasm_bindgen::JsValue::NULL, Into::into),
1295 now().into(),
1296 job_id.into(),
1297 ])?
1298 .first::<JobRow>(None)
1299 .await?;
1300 let Some(job) = finished else { return Ok(()) };
Merge branch 'worktree-agent-a3abfcce648e87dca'1301 // Its G1T_TOKEN stops working with it.
1302 self.revoke_job_tokens(&job.id).await;
Merge branch 'main' into worktree-agent-a69aeabc4b0deeb971303 // A job that deploys failed: so did its run's deployment, now.
1304 if conclusion == "failure" && job.continue_on_error == 0 && job.started_at.is_some()
1305 && let Some(run) = self.run_row(&job.run_id).await?
1306 && let Some(env) = deploys_to(&run, &job)
1307 {
1308 self.report_deployment(&run, &env, "failure", false).await;
1309 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents1310 // A self-hosted runner's job: the runner is free again, and its time
1311 // is recorded, at nothing.
1312 if job.runner_id.is_some() {
1313 self.released(&job).await?;
1314 }
GitHub Actions on g1t, part two: running workflows1315 // Steps still marked as going are not going any more.
1316 let mut steps: Vec<Value> = serde_json::from_str(&job.steps).unwrap_or_default();
1317 let mut touched = false;
1318 for step in steps.iter_mut() {
1319 if step["status"] != "completed" {
1320 let was_running = step["status"] == "in_progress";
1321 step["status"] = json!("completed");
1322 step["conclusion"] = json!(if was_running { conclusion } else { "skipped" });
1323 touched = true;
1324 }
1325 }
1326 if touched {
1327 self.db
1328 .prepare("UPDATE jobs SET steps = ? WHERE id = ?")
1329 .bind(&[serde_json::to_string(&steps)?.into(), job.id.as_str().into()])?
1330 .run()
1331 .await?;
1332 }
1333 // fail-fast: one failed combination stops the rest of its matrix.
1334 if conclusion == "failure" && job.continue_on_error == 0 && job.matrix.as_deref().is_some_and(|m| m != "{}") {
1335 let run = self.run_row(&job.run_id).await?;
1336 let fail_fast = run
1337 .as_ref()
1338 .and_then(|run| workflow::parse(&run.source).ok())
1339 .and_then(|workflow| workflow.jobs.into_iter().find(|j| j.id == job.key))
1340 .is_none_or(|j| j.fail_fast);
1341 if fail_fast {
1342 let siblings = self
1343 .db
1344 .prepare("SELECT * FROM jobs WHERE run_id = ? AND key = ? AND status != 'completed'")
1345 .bind(&[job.run_id.as_str().into(), job.key.as_str().into()])?
1346 .all()
1347 .await?
1348 .results::<JobRow>()?;
1349 for sibling in siblings {
1350 self.stop_job(&sibling, "Another job of its matrix failed, and the matrix is fail-fast.").await?;
1351 }
1352 }
1353 }
1354 Box::pin(self.advance(&job.run_id)).await
1355 }
1356
1357 /// Cancels a job, stopping its sandbox if it has one.
1358 async fn stop_job(&self, job: &JobRow, reason: &str) -> Result<()> {
Fast pages, required checks on the branch, self-hosted runners, honest incidents1359 // A self-hosted runner hears it was cancelled on its next poll.
1360 if job.status == "in_progress" && job.runner_id.is_none() {
GitHub Actions on g1t, part two: running workflows1361 let _: Result<Value> = g1t_kit::call(&self.runner, "stop_actions_job", &json!({ "job": job.id })).await;
1362 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents1363 let stopped = self
1364 .db
1365 .prepare("UPDATE jobs SET status = 'completed', conclusion = 'cancelled', reason = ?, finished_at = ?, token_hash = NULL WHERE id = ? AND status != 'completed' RETURNING *")
GitHub Actions on g1t, part two: running workflows1366 .bind(&[reason.into(), now().into(), job.id.as_str().into()])?
Fast pages, required checks on the branch, self-hosted runners, honest incidents1367 .first::<JobRow>(None)
GitHub Actions on g1t, part two: running workflows1368 .await?;
Merge branch 'worktree-agent-a3abfcce648e87dca'1369 if stopped.as_ref().is_some_and(|row| row.started_at.is_some()) {
1370 self.revoke_job_tokens(&job.id).await;
1371 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents1372 if let Some(stopped) = stopped.filter(|row| row.runner_id.is_some()) {
1373 self.released(&stopped).await?;
1374 }
GitHub Actions on g1t, part two: running workflows1375 Ok(())
1376 }
1377
Merge branch 'worktree-agent-a3abfcce648e87dca'1378 /// Ends a job's tokens at once. A failure is logged: the token expires
1379 /// on its own soon after the job's time limit.
1380 async fn revoke_job_tokens(&self, job_id: &str) {
1381 let revoked: Result<bool> =
1382 g1t_kit::call(&self.identity, "revoke_job_tokens", &RevokeJobTokensArgs { job_id: job_id.to_owned() }).await;
1383 if let Err(error) = revoked {
1384 worker::console_error!("actions: the tokens of job {job_id} were not revoked: {error}");
1385 }
1386 }
1387
GitHub Actions on g1t, part two: running workflows1388 /// Finishes the run when every job has.
1389 async fn finish_if_done(&self, run_id: &str) -> Result<()> {
1390 let Some(run) = self.run_row(run_id).await? else { return Ok(()) };
Merge branch 'worktree-agent-a3abfcce648e87dca'1391 if matches!(run.status.as_str(), "completed" | "pending" | "action_required") {
GitHub Actions on g1t, part two: running workflows1392 return Ok(());
1393 }
1394 let jobs = self.job_rows(run_id).await?;
1395 if !jobs.iter().all(|job| job.status == "completed") {
1396 return Ok(());
1397 }
1398 self.finish_run(&run, None).await
1399 }
1400
1401 async fn finish_run(&self, run: &RunRow, error: Option<&str>) -> Result<()> {
1402 let jobs = self.job_rows(&run.id).await?;
1403 let rows: Vec<&JobRow> = jobs.iter().collect();
1404 let conclusion = if error.is_some() {
1405 "failure"
1406 } else if run.conclusion.as_deref() == Some("cancelled") {
1407 "cancelled"
1408 } else if rows.is_empty() {
1409 "skipped"
1410 } else {
1411 key_result(&rows)
1412 };
1413 let done = self
1414 .db
1415 .prepare("UPDATE runs SET status = 'completed', conclusion = ?, error = COALESCE(?, error), finished_at = ? WHERE id = ? AND status != 'completed' RETURNING id")
1416 .bind(&[conclusion.into(), optional(error), now().into(), run.id.as_str().into()])?
1417 .first::<Value>(None)
1418 .await?;
1419 if done.is_none() {
1420 return Ok(());
1421 }
1422 self.report_status(run, conclusion).await?;
Merge branch 'main' into worktree-agent-a69aeabc4b0deeb971423 self.settle_deployments(run, &jobs, conclusion == "cancelled").await;
Actions: workflow_run, workflow.completed, artifacts on the run page, Node 241424 let published: Result<()> = g1t_kit::call(
1425 &self.events,
1426 "publish",
1427 &g1t_contracts::events::Publish {
1428 events: vec![g1t_contracts::events::NewEvent {
1429 kind: "workflow.completed",
1430 source: "actions",
1431 repo_id: Some(run.repo_id.clone()),
1432 actor: run.actor_id.clone(),
1433 data: g1t_contracts::events::WorkflowEvent {
1434 run_id: run.id.clone(),
1435 repo_id: run.repo_id.clone(),
1436 workflow: run.name.clone(),
1437 path: run.path.clone(),
1438 number: run.number,
1439 event: run.event.clone(),
1440 conclusion: conclusion.to_owned(),
1441 git_ref: run.git_ref.clone(),
1442 sha: run.sha.clone(),
1443 pull: run.pull,
1444 },
1445 }],
1446 },
1447 )
1448 .await;
1449 if let Err(error) = published {
1450 worker::console_error!("actions: could not publish workflow.completed: {error}");
1451 }
GitHub Actions on g1t, part two: running workflows1452 // The next run waiting in its concurrency group.
1453 if let Some(group) = &run.concurrency_group {
1454 let next = self
1455 .db
1456 .prepare("SELECT * FROM runs WHERE repo_id = ? AND concurrency_group = ? AND status = 'pending' ORDER BY id LIMIT 1")
1457 .bind(&[run.repo_id.as_str().into(), group.as_str().into()])?
1458 .first::<RunRow>(None)
1459 .await?;
1460 if let Some(next) = next {
1461 self.db.prepare("UPDATE runs SET status = 'queued' WHERE id = ?").bind(&[next.id.as_str().into()])?.run().await?;
1462 Box::pin(self.advance(&next.id)).await?;
1463 }
1464 }
1465 Ok(())
1466 }
1467
Merge branch 'main' into worktree-agent-a69aeabc4b0deeb971468 /// Tells the deployments service how a run's deployment to `env` stands:
1469 /// made when its first job naming the environment starts, failed when one
1470 /// of them fails, and settled (`last`) when the run finishes. Never fails
1471 /// the run: a deployment that could not be recorded is logged.
1472 async fn report_deployment(&self, run: &RunRow, env: &JobEnvironment, state: &str, last: bool) {
1473 let path = repo_path(&run.repo);
1474 let reported: Result<Value> = g1t_kit::call(
1475 &self.deployments,
1476 "actions_deployment",
1477 &json!({
1478 "repoId": run.repo_id,
1479 "repo": { "namespace": path.namespace, "name": path.name },
1480 "runId": run.id,
1481 "attempt": run.attempt,
1482 "runUrl": format!("{SITE}/{}/actions/runs/{}", run.repo, run.id),
1483 "environment": env.name,
1484 "url": env.url,
1485 "ref": run.git_ref,
1486 "sha": run.sha,
1487 "state": state,
1488 "final": last,
1489 "creator": run.actor,
1490 "workflow": run.name,
1491 }),
1492 )
1493 .await;
1494 if let Err(error) = reported {
1495 worker::console_error!("actions: deployment to {} not recorded for run {}: {error}", env.name, run.id);
1496 }
1497 }
1498
1499 /// A job that deploys has started: its run's deployment to the
1500 /// environment is under way.
1501 pub(crate) async fn job_started(&self, job_id: &str) -> Result<()> {
1502 let Some(job) = self.db.prepare("SELECT * FROM jobs WHERE id = ?").bind(&[job_id.into()])?.first::<JobRow>(None).await? else {
1503 return Ok(());
1504 };
1505 let Some(run) = self.run_row(&job.run_id).await? else { return Ok(()) };
1506 if let Some(env) = deploys_to(&run, &job) {
1507 self.report_deployment(&run, &env, "in_progress", false).await;
1508 }
1509 Ok(())
1510 }
1511
1512 /// The run is over: each environment its jobs deployed to takes the
1513 /// outcome of those jobs (`deployment_outcome`).
1514 async fn settle_deployments(&self, run: &RunRow, jobs: &[JobRow], cancelled: bool) {
1515 let mut seen: Vec<(JobEnvironment, Vec<Option<String>>)> = Vec::new();
1516 for job in jobs {
1517 let Some(env) = deploys_to(run, job) else { continue };
1518 let conclusion = if cancelled && job.conclusion.as_deref() != Some("skipped") && job.started_at.is_some() {
1519 Some("cancelled".to_owned())
1520 } else if job.started_at.is_none() {
1521 Some("skipped".to_owned())
1522 } else {
1523 job.conclusion.clone()
1524 };
1525 match seen.iter_mut().find(|(known, _)| known.name.eq_ignore_ascii_case(&env.name)) {
1526 Some((known, conclusions)) => {
1527 if known.url.is_none() {
1528 known.url = env.url.clone();
1529 }
1530 conclusions.push(conclusion);
1531 }
1532 None => seen.push((env, vec![conclusion])),
1533 }
1534 }
1535 for (env, conclusions) in seen {
1536 if let Some(state) = deployment_outcome(&conclusions) {
1537 self.report_deployment(run, &env, state, true).await;
1538 }
1539 }
1540 }
1541
GitHub Actions on g1t, part two: running workflows1542 /// Tells the pull request (or commit) how the run went, as a status.
1543 async fn report_status(&self, run: &RunRow, conclusion: &str) -> Result<()> {
1544 let state = match conclusion {
1545 "success" | "skipped" => "success",
1546 "cancelled" => "error",
1547 _ => "failure",
1548 };
1549 let _: Result<Value> = g1t_kit::call(
1550 &self.work,
1551 "set_commit_status",
1552 &json!({
1553 "repoId": run.repo_id,
1554 "sha": run.sha,
1555 "context": format!("{} / {}", run.name, run.event),
1556 "state": state,
1557 "description": format!("{} {}", run.name, match conclusion {
1558 "success" => "passed",
1559 "skipped" => "was skipped",
1560 "cancelled" => "was cancelled",
1561 _ => "failed",
1562 }),
1563 "targetUrl": format!("{SITE}/{}/actions/runs/{}", run.repo, run.id),
Merge rulesets: branch and tag rules, agent-first, enforced on push and merge1564 "source": "actions",
GitHub Actions on g1t, part two: running workflows1565 }),
1566 )
1567 .await;
1568 Ok(())
1569 }
1570
1571 /// Tells the pull request a run has started on its head.
1572 pub async fn report_pending(&self, run: &RunRow) -> Result<()> {
1573 let _: Result<Value> = g1t_kit::call(
1574 &self.work,
1575 "set_commit_status",
1576 &json!({
1577 "repoId": run.repo_id,
1578 "sha": run.sha,
1579 "context": format!("{} / {}", run.name, run.event),
1580 "state": "pending",
Merge branch 'worktree-agent-a3abfcce648e87dca'1581 "description": if run.status == "action_required" {
1582 format!("{} is waiting for approval", run.name)
1583 } else {
1584 format!("{} is running", run.name)
1585 },
GitHub Actions on g1t, part two: running workflows1586 "targetUrl": format!("{SITE}/{}/actions/runs/{}", run.repo, run.id),
Merge rulesets: branch and tag rules, agent-first, enforced on push and merge1587 "source": "actions",
GitHub Actions on g1t, part two: running workflows1588 }),
1589 )
1590 .await;
1591 Ok(())
1592 }
1593
1594 /// Cancels a run: its waiting and queued jobs, and stops its running ones.
1595 pub async fn cancel_run(&self, run: &RunRow, reason: &str) -> Result<()> {
1596 self.db
1597 .prepare("UPDATE runs SET conclusion = 'cancelled' WHERE id = ? AND status != 'completed'")
1598 .bind(&[run.id.as_str().into()])?
1599 .run()
1600 .await?;
1601 for job in self.job_rows(&run.id).await?.iter().filter(|job| job.status != "completed") {
1602 self.stop_job(job, reason).await?;
1603 }
Merge branch 'worktree-agent-a3abfcce648e87dca'1604 if run.status == "pending" || run.status == "action_required" {
GitHub Actions on g1t, part two: running workflows1605 self.db.prepare("UPDATE runs SET status = 'queued' WHERE id = ?").bind(&[run.id.as_str().into()])?.run().await?;
1606 }
1607 self.advance(&run.id).await
1608 }
1609
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look1610 /// Cancels every run of the repository that has not finished, for
1611 /// `repo.deleted` and `repo.archived`.
1612 pub async fn stop_runs(&self, repo_id: &str) -> Result<()> {
1613 let runs = self
1614 .db
1615 .prepare("SELECT * FROM runs WHERE repo_id = ? AND status != 'completed'")
1616 .bind(&[repo_id.into()])?
1617 .all()
1618 .await?
1619 .results::<RunRow>()?;
1620 for run in runs {
1621 self.cancel_run(&run, "The repository was archived or deleted.").await?;
1622 }
1623 Ok(())
1624 }
1625
GitHub Actions on g1t, part two: running workflows1626 pub async fn cancel(&self, a: RunActionArgs) -> Result<Outcome<WorkflowRun>> {
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look1627 if let Outcome::Fail(refused) = self.may(&a.actor, &a.repo, Capability::Run).await? {
GitHub Actions on g1t, part two: running workflows1628 return Ok(Outcome::Fail(refused));
1629 }
1630 let run = check!(self.run_in(&a.repo, &a.id).await?);
1631 if run.status == "completed" {
1632 return Ok(fail(FailureCode::Conflict, "The run has already finished."));
1633 }
1634 self.cancel_run(&run, &format!("{} cancelled the run.", a.actor.username)).await?;
1635 self.run_summary(&run.id).await
1636 }
1637
1638 /// Runs again: every job, or with `failed_only` those that did not
1639 /// succeed and the jobs that need them.
1640 pub async fn rerun(&self, a: RunActionArgs) -> Result<Outcome<WorkflowRun>> {
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look1641 if let Outcome::Fail(refused) = self.may(&a.actor, &a.repo, Capability::Run).await? {
GitHub Actions on g1t, part two: running workflows1642 return Ok(Outcome::Fail(refused));
1643 }
1644 let run = check!(self.run_in(&a.repo, &a.id).await?);
1645 if run.status != "completed" {
1646 return Ok(fail(FailureCode::Conflict, "The run is still going: cancel it first."));
1647 }
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look1648 // Nothing starts again on an archived repository.
1649 match self.visible_repo(&a.repo, &Some(a.actor.clone())).await? {
1650 Some(repo) if repo.archived() => {
1651 return Ok(fail(FailureCode::Forbidden, g1t_contracts::repos::archived_message(&repo.namespace, &repo.name)));
1652 }
1653 Some(_) => {}
1654 None => return Ok(fail(FailureCode::NotFound, "There is no such repository.")),
1655 }
GitHub Actions on g1t, part two: running workflows1656 if run.error.is_some() {
1657 return Ok(fail(FailureCode::Conflict, "This run never started: fix the workflow file and push again."));
1658 }
1659 let jobs = self.job_rows(&run.id).await?;
1660 let workflow = workflow::parse(&run.source).ok();
1661 // Which keys run again: failed ones and, transitively, those needing them.
1662 let mut again: Vec<String> = Vec::new();
1663 for key in workflow.as_ref().map(|w| w.job_order()).unwrap_or_default() {
1664 let rows: Vec<&JobRow> = jobs.iter().filter(|j| j.key == key).collect();
1665 let failed = rows.iter().any(|row| row.conclusion.as_deref() != Some("success"));
1666 let needs_again = rows.first().is_some_and(|row| row.needs().iter().any(|need| again.contains(need)));
1667 if !a.failed_only || failed || needs_again {
1668 again.push(key.to_owned());
1669 }
1670 }
1671 if again.is_empty() {
1672 return Ok(fail(FailureCode::Conflict, "Every job succeeded: there is nothing to run again."));
1673 }
1674 let mut statements = Vec::new();
1675 for key in &again {
1676 statements.push(self.db.prepare("DELETE FROM logs WHERE job_id IN (SELECT id FROM jobs WHERE run_id = ? AND key = ?)").bind(&[run.id.as_str().into(), key.as_str().into()])?);
1677 statements.push(self.db.prepare("DELETE FROM jobs WHERE run_id = ? AND key = ? AND ordinal > 0").bind(&[run.id.as_str().into(), key.as_str().into()])?);
Actions: reusable workflows in the repository1678 // The jobs of a workflow it called are made again when it calls it again.
GitHub Actions on g1t, part two: running workflows1679 statements.push(
1680 self.db
Actions: reusable workflows in the repository1681 .prepare("DELETE FROM logs WHERE job_id IN (SELECT id FROM jobs WHERE run_id = ? AND key LIKE ?)")
1682 .bind(&[run.id.as_str().into(), format!("{key}/%").into()])?,
1683 );
1684 statements.push(
1685 self.db
1686 .prepare("DELETE FROM jobs WHERE run_id = ? AND key LIKE ?")
1687 .bind(&[run.id.as_str().into(), format!("{key}/%").into()])?,
1688 );
1689 statements.push(
1690 self.db
GitHub Actions on g1t, part two: running workflows1691 .prepare(
1692 "UPDATE jobs SET status = 'waiting', conclusion = NULL, steps = '[]', annotations = '[]', outputs = '{}', reason = NULL,
Fast pages, required checks on the branch, self-hosted runners, honest incidents1693 matrix = NULL, call = NULL, token_hash = NULL, seen_at = NULL, started_at = NULL, finished_at = NULL,
Merge branch 'worktree-agent-a3abfcce648e87dca'1694 labels = NULL, queued_at = NULL, runner_id = NULL, runner_name = NULL, environment = NULL,
1695 concurrency_group = NULL, cancel_in_progress = 0 WHERE run_id = ? AND key = ?",
GitHub Actions on g1t, part two: running workflows1696 )
1697 .bind(&[run.id.as_str().into(), key.as_str().into()])?,
1698 );
1699 }
1700 statements.push(
1701 self.db
1702 .prepare("UPDATE runs SET status = 'queued', conclusion = NULL, attempt = attempt + 1, started_at = NULL, finished_at = NULL WHERE id = ?")
1703 .bind(&[run.id.as_str().into()])?,
1704 );
1705 self.db.batch(statements).await?;
1706 if let Some(run) = self.run_row(&run.id).await? {
1707 self.report_pending(&run).await?;
1708 }
1709 self.advance(&run.id).await?;
1710 self.run_summary(&run.id).await
1711 }
1712
1713 pub async fn run_in(&self, repo: &RepoPath, id: &str) -> Result<Outcome<RunRow>> {
1714 let row = self
1715 .db
1716 .prepare("SELECT * FROM runs WHERE id = ? AND lower(repo) = lower(?)")
1717 .bind(&[id.into(), format!("{}/{}", repo.namespace, repo.name).into()])?
1718 .first::<RunRow>(None)
1719 .await?;
1720 Ok(row.map_or_else(|| fail(FailureCode::NotFound, "No such run."), Outcome::Ok))
1721 }
1722
1723 // --- The sandbox's side -----------------------------------------------------
1724
Fast pages, required checks on the branch, self-hosted runners, honest incidents1725 pub(crate) async fn job_for_token(&self, a: &JobCallArgs) -> Result<Outcome<JobRow>> {
GitHub Actions on g1t, part two: running workflows1726 let job = self.db.prepare("SELECT * FROM jobs WHERE id = ?").bind(&[a.job.as_str().into()])?.first::<JobRow>(None).await?;
1727 Ok(match job {
1728 Some(job) if job.status == "in_progress" && job.token_hash.as_deref().is_some_and(|hash| same(hash, &sha256_hex(&a.token))) => {
1729 Outcome::Ok(job)
1730 }
1731 _ => fail(FailureCode::Unauthenticated, "That job is not running, or the token is not its."),
1732 })
1733 }
1734
A repository has its own sidebar, as settings do1735 /// `job_auth`: which run and repository a running job's token is for,
1736 /// so the API can keep its artifacts and cache.
1737 pub async fn job_auth(&self, a: JobCallArgs) -> Result<Outcome<Value>> {
1738 let job = check!(self.job_for_token(&a).await?);
1739 Ok(Outcome::Ok(json!({ "run": job.run_id, "repoId": job.repo_id })))
1740 }
1741
GitHub Actions on g1t, part two: running workflows1742 /// `job_spec`: everything the sandbox needs to run the job.
1743 pub async fn job_spec(&self, a: JobCallArgs) -> Result<Outcome<Value>> {
1744 let job = check!(self.job_for_token(&a).await?);
1745 let Some(run) = self.run_row(&job.run_id).await? else {
1746 return Ok(fail(FailureCode::NotFound, "No such run."));
1747 };
Actions: reusable workflows in the repository1748 let Ok(caller) = workflow::parse(&run.source) else {
GitHub Actions on g1t, part two: running workflows1749 return Ok(fail(FailureCode::Invalid, "The workflow no longer reads."));
1750 };
Actions: reusable workflows in the repository1751 // A called workflow's job runs as that workflow defines it.
1752 let callee = job.callee();
1753 let (workflow, spec, call_inputs) = match callee {
1754 Some((called, spec, call)) => (called, spec, Some(call["inputs"].clone())),
1755 None => match caller.jobs.iter().find(|j| j.id == job.key) {
1756 Some(spec) => (caller.clone(), spec.clone(), None),
1757 None => return Ok(fail(FailureCode::NotFound, "The job is not in the workflow.")),
1758 },
GitHub Actions on g1t, part two: running workflows1759 };
Actions: reusable workflows in the repository1760 let spec = &spec;
GitHub Actions on g1t, part two: running workflows1761 let repo = repo_path(&run.repo);
1762 let trusted = run.trusted != 0;
Merge branch 'worktree-agent-a3abfcce648e87dca'1763 // The job's `environment:`, by name, as read when its needs were done
1764 // (an expression included), once the environment's protection rules
1765 // let it start: entries with a value for it give that value instead
1766 // of their default, as GitHub's environment secrets do.
1767 let environment: Option<String> = job.environment.clone();
1768 // What its token may do: its `permissions:` (a called workflow's
1769 // jobs no more than the job that calls it), else the repository's
1770 // default; read-only for a pull request from outside.
1771 // The repository's default, its workspace's taken in, and whether
1772 // its jobs may open and approve pull requests.
1773 let (default, pull_requests) = self.token_policy(&run.repo_id).await?;
1774 let mut permissions = spec.permissions(&workflow, default);
1775 if let Some(parent) = &job.call().filter(|c| c["role"] == "callee").and_then(|c| c["parent"].as_str().map(str::to_owned)) {
1776 let top = parent.split('/').next().unwrap_or(parent);
1777 if let Some(caller_job) = caller.jobs.iter().find(|j| j.id == top) {
1778 permissions = permissions.capped_by(&caller_job.permissions(&caller, default));
GitHub Actions on g1t, part two: running workflows1779 }
Merge branch 'worktree-agent-a3abfcce648e87dca'1780 }
1781 if !trusted {
1782 permissions = permissions.read_only();
1783 }
1784 // G1T_TOKEN, and GITHUB_TOKEN as its alias: a token of the
1785 // workspace's that reaches this repository only, with the scopes
1786 // its permissions give, until the job ends.
1787 let token = match self.workspace_actor(&repo.namespace).await? {
1788 Some(workspace) => {
1789 let created: CreatedAccessToken = g1t_kit::call(
1790 &self.identity,
1791 "create_job_token",
1792 &CreateJobTokenArgs {
1793 workspace,
1794 repo: repo.clone(),
1795 run_id: run.id.clone(),
1796 job_id: job.id.clone(),
1797 name: format!("G1T_TOKEN for {} run {}", run.repo, run.number),
1798 ttl_seconds: u64::from(job.timeout_minutes) * 60 + 600,
1799 scopes: permissions.scopes().into_iter().map(str::to_owned).collect(),
1800 pull_requests: pull_requests && trusted,
1801 },
1802 )
1803 .await?;
1804 created.token
1805 }
1806 None => String::new(),
GitHub Actions on g1t, part two: running workflows1807 };
Secrets and variables: one list, rows per environment, for workflows and deployments1808 // A run that is not trusted (a pull request from outside the
1809 // workspace) gets no secrets and an empty token.
Merge Actions: cross-repo workflows and actions, release and deployment triggers, step timeouts1810 let passed = job.call().filter(|c| c["role"] == "callee").and_then(|c| c.get("secrets").cloned());
1811 let mut secrets = match (trusted, passed) {
1812 (false, _) => Map::new(),
1813 (true, None) => self.secrets_for(&run.repo_id, &run.repo, environment.as_deref(), true).await?,
1814 // A called workflow's job: what its callers passed it (reach.rs),
1815 // and its own environment's secrets over them.
1816 (true, Some(plan)) => {
1817 let base = self.secrets_for(&run.repo_id, &run.repo, None, true).await?;
1818 let vars = self.variables_for(&run.repo_id, &run.repo, None, true).await?;
1819 let github = run.info().context(&job.key, "", run.action.as_deref());
1820 let mut passed = crate::reach::resolve_secrets(&plan, &base, &github, &vars);
1821 if let Some(name) = environment.as_deref() {
1822 let own = self.secrets_for(&run.repo_id, &run.repo, Some(name), true).await?;
1823 for (key, value) in own {
1824 if base.get(&key) != Some(&value) {
1825 passed.insert(key, value);
1826 }
1827 }
1828 }
1829 passed
1830 }
Secrets and variables: one list, rows per environment, for workflows and deployments1831 };
1832 secrets.insert("G1T_TOKEN".into(), Value::String(token.clone()));
GitHub Actions on g1t, part two: running workflows1833 secrets.insert("GITHUB_TOKEN".into(), Value::String(token.clone()));
Merge branch 'worktree-agent-a3abfcce648e87dca'1834 // Each secret as it is, a line at a time, base64 and JSON-escaped.
1835 let mut masks: Vec<String> = g1t_actions::mask::all_variants(secrets.values().filter_map(|v| v.as_str()));
Secrets and variables: one list, rows per environment, for workflows and deployments1836 let vars = self.variables_for(&run.repo_id, &run.repo, environment.as_deref(), trusted).await?;
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R21837 // The toolkit's runtime token (runtime.rs), for as long as the job
1838 // may run; the API puts it and the toolkit's addresses in the
1839 // job's variables.
1840 let runtime_token = crate::runtime::runtime_token(
1841 &job.id,
1842 &job.run_id,
1843 job.token_hash.as_deref().unwrap_or_default(),
1844 now_ms() / 1000,
1845 u64::from(job.timeout_minutes) * 60 + 600,
1846 );
1847 masks.push(runtime_token.clone());
1848 let retention_days = self.retention_setting(&run.repo_id).await?;
GitHub Actions on g1t, part two: running workflows1849
1850 let jobs = self.job_rows(&run.id).await?;
1851 let mut needs = Map::new();
Actions: reusable workflows in the repository1852 // In a called workflow, its jobs' keys sit under the job that called it.
1853 let parent = job.call().filter(|c| c["role"] == "callee").and_then(|c| c["parent"].as_str().map(str::to_owned));
GitHub Actions on g1t, part two: running workflows1854 for need in &spec.needs {
Actions: reusable workflows in the repository1855 let key = match &parent {
1856 Some(parent) => format!("{parent}/{need}"),
1857 None => need.clone(),
1858 };
1859 let rows: Vec<&JobRow> = jobs.iter().filter(|row| row.key == key).collect();
GitHub Actions on g1t, part two: running workflows1860 let mut outputs = Map::new();
1861 for row in &rows {
1862 if let Ok(Value::Object(more)) = serde_json::from_str::<Value>(&row.outputs) {
1863 outputs.extend(more);
1864 }
1865 }
1866 needs.insert(need.clone(), json!({ "result": key_result(&rows), "outputs": outputs }));
1867 }
1868 let siblings = jobs.iter().filter(|row| row.key == job.key).count();
1869 let matrix: Value = job.matrix.as_deref().and_then(|m| serde_json::from_str(m).ok()).unwrap_or(json!({}));
1870 let info = run.info();
1871 let mut github = info.context(&job.key, &token, run.action.as_deref());
1872 github["token"] = json!(token);
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R21873 github["retention_days"] = json!(retention_days);
Fast pages, required checks on the branch, self-hosted runners, honest incidents1874 // On a self-hosted runner, `runner` and `RUNNER_*` describe that
1875 // machine rather than g1t's sandbox.
1876 let mut variables = info.variables(&job.key);
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R21877 variables.insert("GITHUB_RETENTION_DAYS".into(), json!(retention_days.to_string()));
Fast pages, required checks on the branch, self-hosted runners, honest incidents1878 let runner = match &job.runner_id {
1879 Some(id) => self.runner_context_for(id, &mut variables).await?,
1880 None => runner_context(),
1881 };
GitHub Actions on g1t, part two: running workflows1882
1883 // Where to check out: a pull request's fork, or the repository.
1884 let clone_url = match run.pull {
1885 Some(number) if run.event.starts_with("pull_request") && run.event != "pull_request_target" => {
1886 let located: Outcome<g1t_contracts::work::PullDetail> = g1t_kit::call(
1887 &self.work,
1888 "get_pull",
1889 &g1t_contracts::work::ViewArgs {
1890 repo: repo.clone(),
1891 number,
1892 viewer: self.workspace_actor(&repo.namespace).await?,
1893 after_seq: 0,
1894 },
1895 )
1896 .await?;
1897 match located {
1898 Outcome::Ok(detail) => match detail.pull.fork {
1899 Some(fork) => format!("{SITE}/{}/{}.git", fork.namespace, fork.name),
1900 None => format!("{SITE}/{}.git", run.repo),
1901 },
1902 Outcome::Fail(_) => format!("{SITE}/{}.git", run.repo),
1903 }
1904 }
1905 _ => format!("{SITE}/{}.git", run.repo),
1906 };
1907
1908 Ok(Outcome::Ok(json!({
1909 "job": job.id,
1910 "run": run.id,
1911 "key": job.key,
1912 "name": job.name,
1913 "spec": spec.raw,
1914 "workflow": {
1915 "env": workflow.env,
1916 "defaults": workflow.raw.get("defaults").cloned().unwrap_or(Value::Null),
1917 },
1918 "github": github,
Fast pages, required checks on the branch, self-hosted runners, honest incidents1919 "variables": variables,
GitHub Actions on g1t, part two: running workflows1920 "event": info.event,
1921 "contexts": {
1922 "vars": vars,
1923 "secrets": secrets,
Actions: reusable workflows in the repository1924 "inputs": call_inputs.unwrap_or_else(|| Value::Object(run.inputs())),
GitHub Actions on g1t, part two: running workflows1925 "matrix": matrix,
1926 "needs": needs,
1927 "strategy": {
1928 "fail-fast": spec.fail_fast,
1929 "job-index": job.ordinal,
1930 "job-total": siblings,
1931 "max-parallel": spec.max_parallel.unwrap_or(siblings as u32),
1932 },
Fast pages, required checks on the branch, self-hosted runners, honest incidents1933 "runner": runner,
GitHub Actions on g1t, part two: running workflows1934 },
1935 "checkout": {
1936 "repository": run.repo,
1937 "url": clone_url,
1938 "sha": run.sha,
1939 "ref": run.git_ref,
1940 "token": token,
1941 },
1942 "timeoutMinutes": job.timeout_minutes,
1943 "masks": masks,
Merge branch 'worktree-agent-a3abfcce648e87dca'1944 // As the job's log lists them at its start.
1945 "permissions": permissions.listed().into_iter().map(|(name, access)| (name.to_owned(), json!(access.as_str()))).collect::<Map<String, Value>>(),
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R21946 // Whether the job may ask for an OIDC token decides whether it
1947 // is told where to.
1948 "runtime": {
1949 "token": runtime_token,
1950 "idToken": self.oidc_allowed(&run, &job),
1951 },
GitHub Actions on g1t, part two: running workflows1952 })))
1953 }
1954
1955 /// `job_report`: the sandbox telling how the job is going.
1956 pub async fn job_report(&self, a: JobCallArgs) -> Result<Outcome<Value>> {
1957 let job = check!(self.job_for_token(&a).await?);
1958 let report = &a.report;
1959 let at = now();
1960 match report["kind"].as_str().unwrap_or_default() {
1961 "steps" => {
1962 // The list can grow as the job goes (post steps), so steps
1963 // already reported keep where they stand.
1964 let known: Vec<Value> = serde_json::from_str(&job.steps).unwrap_or_default();
1965 let steps: Vec<Value> = report["steps"]
1966 .as_array()
1967 .map(|names| {
1968 names
1969 .iter()
1970 .enumerate()
1971 .map(|(i, name)| match known.get(i) {
1972 Some(step) if step["status"] != "queued" => step.clone(),
1973 _ => json!({ "number": i + 1, "name": expr::to_text(name), "status": "queued", "conclusion": null, "startedAt": null, "finishedAt": null }),
1974 })
1975 .collect()
1976 })
1977 .unwrap_or_default();
1978 self.db
1979 .prepare("UPDATE jobs SET steps = ?, seen_at = ? WHERE id = ?")
1980 .bind(&[serde_json::to_string(&steps)?.into(), at.as_str().into(), job.id.as_str().into()])?
1981 .run()
1982 .await?;
1983 }
1984 "step" => {
1985 let number = report["number"].as_u64().unwrap_or(0) as usize;
1986 let mut steps: Vec<Value> = serde_json::from_str(&job.steps).unwrap_or_default();
1987 if let Some(step) = number.checked_sub(1).and_then(|i| steps.get_mut(i)) {
1988 let status = report["status"].as_str().unwrap_or("in_progress");
1989 step["status"] = json!(status);
1990 if status == "in_progress" {
1991 step["startedAt"] = json!(at);
1992 }
1993 if status == "completed" {
1994 step["finishedAt"] = json!(at);
1995 step["conclusion"] = report["conclusion"].clone();
1996 }
1997 if let Some(name) = report["name"].as_str() {
1998 step["name"] = json!(name);
1999 }
2000 }
2001 self.db
2002 .prepare("UPDATE jobs SET steps = ?, seen_at = ? WHERE id = ?")
2003 .bind(&[serde_json::to_string(&steps)?.into(), at.as_str().into(), job.id.as_str().into()])?
2004 .run()
2005 .await?;
2006 }
2007 "log" => {
2008 let mut text = report["text"].as_str().unwrap_or_default().to_owned();
2009 if text.len() > MAX_CHUNK_BYTES {
2010 let mut cut = MAX_CHUNK_BYTES;
2011 while !text.is_char_boundary(cut) {
2012 cut -= 1;
2013 }
2014 text.truncate(cut);
2015 }
2016 #[derive(Deserialize)]
2017 struct Size {
GitHub Actions on g1t, part three: .g1t/workflows, the pages, the docs2018 n: Option<f64>,
2019 seq: Option<f64>,
GitHub Actions on g1t, part two: running workflows2020 }
2021 let size = self
2022 .db
2023 .prepare("SELECT SUM(LENGTH(text)) AS n, MAX(seq) AS seq FROM logs WHERE job_id = ?")
2024 .bind(&[job.id.as_str().into()])?
2025 .first::<Size>(None)
2026 .await?;
GitHub Actions on g1t, part three: .g1t/workflows, the pages, the docs2027 let (used, seq) = size.map_or((0, 0.0), |s| (s.n.unwrap_or(0.0) as usize, s.seq.unwrap_or(0.0)));
GitHub Actions on g1t, part two: running workflows2028 if used < MAX_LOG_BYTES {
2029 if used + text.len() >= MAX_LOG_BYTES {
2030 text.push_str("\n… The log reached its limit of 4 MB; the rest is not kept.\n");
2031 }
2032 self.db
2033 .prepare("INSERT INTO logs (job_id, seq, step, text) VALUES (?, ?, ?, ?)")
GitHub Actions on g1t, part three: .g1t/workflows, the pages, the docs2034 .bind(&[job.id.as_str().into(), // Numbers go to D1 as f64: a u64 would be a BigInt, which it refuses.
2035 (seq + 1.0).into(), (report["step"].as_u64().unwrap_or(0) as u32).into(), text.into()])?
GitHub Actions on g1t, part two: running workflows2036 .run()
2037 .await?;
2038 }
2039 self.db.prepare("UPDATE jobs SET seen_at = ? WHERE id = ?").bind(&[at.into(), job.id.as_str().into()])?.run().await?;
2040 }
2041 "annotation" => {
2042 let mut annotations: Vec<Value> = serde_json::from_str(&job.annotations).unwrap_or_default();
2043 if annotations.len() < MAX_ANNOTATIONS {
2044 annotations.push(json!({
2045 "level": report["level"].as_str().unwrap_or("notice"),
2046 "message": report["message"].as_str().unwrap_or_default().chars().take(4000).collect::<String>(),
2047 "title": report["title"],
2048 "file": report["file"],
2049 "line": report["line"],
2050 }));
2051 self.db
2052 .prepare("UPDATE jobs SET annotations = ?, seen_at = ? WHERE id = ?")
2053 .bind(&[serde_json::to_string(&annotations)?.into(), at.as_str().into(), job.id.as_str().into()])?
2054 .run()
2055 .await?;
2056 }
2057 }
2058 "done" => {
2059 let conclusion = report["conclusion"]
2060 .as_str()
2061 .filter(|c| matches!(*c, "success" | "failure" | "cancelled"))
2062 .unwrap_or("failure");
2063 let outputs = report["outputs"].as_object().cloned();
2064 Box::pin(self.finish_job(&job.id, conclusion, report["reason"].as_str(), outputs.as_ref())).await?;
2065 }
2066 other => return Ok(fail(FailureCode::Invalid, format!("There is no report called `{other}`."))),
2067 }
2068 Ok(Outcome::Ok(json!({ "ok": true })))
2069 }
2070
2071 // --- Every minute ---------------------------------------------------------------
2072
2073 pub async fn on_minute(&self, now_ms: u64) -> Result<()> {
2074 let minute = now_ms / 60_000 * 60_000;
2075 if let Err(error) = self.run_schedules(minute).await {
2076 worker::console_error!("actions: schedules failed: {error}");
2077 }
2078 // Jobs whose sandbox went quiet or ran past their time.
2079 let running = self.db.prepare("SELECT * FROM jobs WHERE status = 'in_progress'").all().await?.results::<JobRow>()?;
2080 for job in running {
2081 // Times in g1t's format compare as text.
2082 let before = |ms: u64| rfc3339(now_ms.saturating_sub(ms));
2083 let silent = job.seen_at.as_deref().is_some_and(|seen| seen < before(SILENT_MS).as_str());
2084 let limit = (u64::from(job.timeout_minutes) * 60 + 120) * 1000;
2085 let over = job.started_at.as_deref().is_some_and(|started| started < before(limit).as_str());
2086 if over {
2087 let reason = format!("It ran longer than its time limit of {} minutes.", job.timeout_minutes);
Fast pages, required checks on the branch, self-hosted runners, honest incidents2088 // A self-hosted runner is told to stop on its next poll.
2089 if job.runner_id.is_none() {
2090 let _: Result<Value> = g1t_kit::call(&self.runner, "stop_actions_job", &json!({ "job": job.id })).await;
2091 }
GitHub Actions on g1t, part two: running workflows2092 self.finish_job(&job.id, "failure", Some(&reason), None).await?;
2093 } else if silent {
Fast pages, required checks on the branch, self-hosted runners, honest incidents2094 let reason = match &job.runner_name {
2095 Some(name) => format!("The self-hosted runner {name} stopped answering."),
2096 None => "The runner stopped answering.".to_owned(),
2097 };
2098 self.finish_job(&job.id, "failure", Some(&reason), None).await?;
GitHub Actions on g1t, part two: running workflows2099 }
2100 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents2101 if let Err(error) = self.sweep_runners(now_ms).await {
2102 worker::console_error!("actions: the runners' sweep failed: {error}");
2103 }
Merge branch 'worktree-agent-a3abfcce648e87dca'2104 // Jobs held at an environment whose wait timer has run out.
2105 if let Err(error) = self.release_gates().await {
2106 worker::console_error!("actions: environments' gates failed: {error}");
2107 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents2108 // Once an hour: the cache's expired entries, and its storage.
2109 if (now_ms / 60_000) % 60 == 7
2110 && let Err(error) = self.sweep_cache(now_ms).await
2111 {
2112 worker::console_error!("actions: the cache's sweep failed: {error}");
2113 }
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R22114 // And artifacts past their time, and the toolkit's abandoned parts.
2115 if (now_ms / 60_000) % 60 == 37 {
2116 if let Err(error) = self.sweep_artifacts(now_ms).await {
2117 worker::console_error!("actions: the artifacts' sweep failed: {error}");
2118 }
2119 if let Err(error) = self.sweep_blob_parts(now_ms).await {
2120 worker::console_error!("actions: the blob parts' sweep failed: {error}");
2121 }
2122 }
GitHub Actions on g1t, part two: running workflows2123 self.start_queued().await
2124 }
2125}
2126
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look2127
2128#[cfg(test)]
2129mod stopping {
2130 use super::stops_runs;
2131 use g1t_contracts::events::Event;
2132 use serde_json::{Value, json};
2133
2134 fn event(kind: &str, data: Value) -> Event {
2135 Event {
2136 id: "evt_1".into(),
2137 kind: kind.into(),
2138 source: "repos".into(),
2139 time: "2026-10-05T00:00:00Z".into(),
2140 repo_id: Some("rep_1".into()),
2141 actor: None,
2142 data,
2143 }
2144 }
2145
2146 #[test]
2147 fn deleting_or_archiving_stops_runs() {
2148 assert_eq!(stops_runs(&event("repo.deleted", json!({ "repoId": "rep_1" }))).as_deref(), Some("rep_1"));
2149 assert_eq!(stops_runs(&event("repo.archived", json!({ "archived": true }))).as_deref(), Some("rep_1"));
2150 assert_eq!(stops_runs(&event("repo.unarchived", json!({ "archived": false }))), None);
2151 assert_eq!(stops_runs(&event("repo.restored", json!({}))), None);
2152 assert_eq!(stops_runs(&event("git.push", json!({}))), None);
2153 }
2154}
Fast pages, required checks on the branch, self-hosted runners, honest incidents2155
2156#[cfg(test)]
2157mod status_of_needs {
2158 use std::collections::HashMap;
2159
2160 use super::ancestor_failed;
2161
2162 /// check -> plan -> (migrate) -> core -> edge, as deploy.yml has them,
2163 /// and a job that needs only the last.
2164 fn graph() -> HashMap<&'static str, Vec<&'static str>> {
2165 HashMap::from([
2166 ("check", vec![]),
2167 ("plan", vec!["check"]),
2168 ("migrate", vec!["plan"]),
2169 ("core", vec!["plan", "migrate"]),
2170 ("edge", vec!["plan", "migrate", "core"]),
2171 ("notify", vec!["edge"]),
2172 ])
2173 }
2174
2175 #[test]
2176 fn a_failure_is_seen_however_far_back() {
2177 let needs = graph();
2178 let failed = |which: &'static str| move |key: &str| key == which;
2179 // check failed; plan, and everything after, was skipped for it.
2180 assert!(ancestor_failed(&needs, "notify", failed("check")));
2181 assert!(ancestor_failed(&needs, "core", failed("check")));
2182 assert!(ancestor_failed(&needs, "edge", failed("core")));
2183 // Nothing before a job failed: a skipped migrate is not a failure.
2184 assert!(!ancestor_failed(&needs, "edge", |_| false));
2185 assert!(!ancestor_failed(&needs, "core", failed("edge")));
2186 assert!(!ancestor_failed(&needs, "check", failed("check")));
2187 }
2188
2189 #[test]
2190 fn cycles_and_unknown_keys_end() {
2191 let needs = HashMap::from([("a", vec!["b"]), ("b", vec!["a"])]);
2192 assert!(!ancestor_failed(&needs, "a", |_| false));
2193 assert!(!ancestor_failed(&needs, "missing", |_| true));
2194 }
2195}
Merge branch 'main' into worktree-agent-a69aeabc4b0deeb972196
2197#[cfg(test)]
2198mod deployments {
2199 use serde_json::{Map, Value, json};
2200
2201 use super::{JobEnvironment, deployment_outcome, environment_of};
2202
2203 #[test]
2204 fn a_jobs_environment_is_read_for_deployments() {
2205 let contexts: Map<String, Value> = serde_json::from_value(json!({
2206 "github": { "ref_name": "main", "repository": "acme/web" },
2207 "inputs": { "target": "staging" },
2208 "matrix": {},
2209 }))
2210 .unwrap();
2211 let read = |raw: Value| environment_of(&raw, &contexts);
2212 assert_eq!(read(json!({})), None);
2213 assert_eq!(
2214 read(json!({ "environment": "production" })),
2215 Some(JobEnvironment { name: "production".into(), url: None, deploys: true })
2216 );
2217 assert_eq!(
2218 read(json!({ "environment": { "name": "production", "url": "https://g1t.sh" } })),
2219 Some(JobEnvironment { name: "production".into(), url: Some("https://g1t.sh".into()), deploys: true })
2220 );
2221 // Expressions are filled in from the run.
2222 assert_eq!(
2223 read(json!({ "environment": { "name": "${{ inputs.target }}", "url": "https://${{ github.ref_name }}.example.com" } })),
2224 Some(JobEnvironment { name: "staging".into(), url: Some("https://main.example.com".into()), deploys: true })
2225 );
2226 // Secrets only: no deployment.
2227 assert!(!read(json!({ "environment": { "name": "production", "deployment": false } })).unwrap().deploys);
2228 // Only http(s) addresses.
2229 assert_eq!(read(json!({ "environment": { "name": "production", "url": "javascript:alert(1)" } })).unwrap().url, None);
2230 }
2231
2232 #[test]
2233 fn a_runs_outcome_for_an_environment() {
2234 let of = |list: &[&str]| deployment_outcome(&list.iter().map(|c| Some((*c).to_owned())).collect::<Vec<_>>());
2235 assert_eq!(of(&["success", "skipped"]), Some("success"));
2236 assert_eq!(of(&["success", "failure"]), Some("failure"));
2237 assert_eq!(of(&["success", "cancelled"]), Some("error"));
2238 assert_eq!(of(&["skipped"]), None);
2239 assert_eq!(deployment_outcome(&[None]), None);
2240 }
2241}

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