Skip to content
2,233 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
Actions: reusable workflows in the repository1013 /// A job that calls a reusable workflow in the repository: that
1014 /// workflow's jobs join the run under it, with the inputs it passes.
1015 async fn call_workflow(&self, run: &RunRow, job: &workflow::Job, row: &JobRow, uses: &str, scope: &Scope<'_>) -> Result<()> {
1016 let Some(local) = uses.strip_prefix("./") else {
1017 return self
1018 .fail_job(row, "Reusable workflows from other repositories are not called on g1t yet; ones in this repository (`./.g1t/workflows/…`) are.")
1019 .await;
1020 };
1021 let depth = row.call().and_then(|c| c["depth"].as_u64()).unwrap_or(0) + 1;
1022 if depth > MAX_CALL_DEPTH {
1023 return self.fail_job(row, &format!("Reusable workflows call each other more than {MAX_CALL_DEPTH} deep.")).await;
1024 }
1025 let local = local.split('@').next().unwrap_or(local).to_owned();
1026 let path = repo_path(&run.repo);
1027 let Some(ws) = self.workspace_actor(&path.namespace).await? else {
1028 return self.fail_job(row, "The workspace is gone.").await;
1029 };
1030 // A repository moved from GitHub keeps saying `.github/…`.
1031 let mut found = self.read_file(&path, &ws, &run.sha, &local).await?.map(|text| (local.clone(), text));
1032 if found.is_none()
1033 && let Some(rest) = local.strip_prefix(".github/")
1034 {
1035 let moved = format!(".g1t/{rest}");
1036 found = self.read_file(&path, &ws, &run.sha, &moved).await?.map(|text| (moved, text));
1037 }
1038 let Some((file, source)) = found else {
1039 return self.fail_job(row, &format!("`{uses}` is not in the repository at this commit.")).await;
1040 };
1041 let called = match workflow::parse(&source) {
1042 Ok(called) => called,
1043 Err(problem) => return self.fail_job(row, &format!("`{file}` does not read: {problem}")).await,
1044 };
1045 let Some(trigger) = called.trigger("workflow_call") else {
1046 return self.fail_job(row, &format!("`{file}` cannot be called: it has no `on: workflow_call`.")).await;
1047 };
1048 // Inputs: what the caller passes, else the called workflow's defaults.
1049 let given = match job.raw.get("with") {
1050 Some(with) => match expr::interpolate_value(with, scope) {
1051 Ok(Value::Object(given)) => given,
1052 Ok(_) => Map::new(),
1053 Err(problem) => return self.fail_job(row, &format!("Its `with` does not read: {problem}")).await,
1054 },
1055 None => Map::new(),
1056 };
1057 let mut inputs = Map::new();
1058 for (name, spec) in &trigger.inputs {
1059 let value = given.get(name).cloned().or_else(|| spec.get("default").cloned()).unwrap_or(Value::Null);
1060 if value.is_null() && spec.get("required").and_then(Value::as_bool) == Some(true) {
1061 return self.fail_job(row, &format!("`{file}` needs the input `{name}`.")).await;
1062 }
1063 inputs.insert(name.clone(), value);
1064 }
1065 for (name, value) in given {
1066 inputs.entry(name).or_insert(value);
1067 }
1068 let mut statements = Vec::new();
1069 for called_job in &called.jobs {
1070 let needs: Vec<String> = called_job.needs.iter().map(|n| format!("{}/{n}", row.key)).collect();
1071 let call = json!({
1072 "role": "callee", "parent": row.key, "job": called_job.id, "path": file,
1073 "source": source, "inputs": inputs, "depth": depth,
1074 });
1075 statements.push(
1076 self.db
1077 .prepare("INSERT INTO jobs (id, run_id, repo_id, namespace, key, name, needs, status, call) VALUES (?, ?, ?, ?, ?, ?, ?, 'waiting', ?)")
1078 .bind(&[
1079 new_id("job", now_ms()).into(),
1080 row.run_id.as_str().into(),
1081 row.repo_id.as_str().into(),
1082 row.namespace.as_str().into(),
1083 format!("{}/{}", row.key, called_job.id).into(),
1084 format!("{} / {}", row.name, called_job.name.clone().unwrap_or(called_job.id.clone())).into(),
1085 serde_json::to_string(&needs)?.into(),
1086 serde_json::to_string(&call)?.into(),
1087 ])?,
1088 );
1089 }
1090 statements.push(
1091 self.db
1092 .prepare("UPDATE jobs SET status = 'calling', call = ?, reason = ?, started_at = ? WHERE id = ?")
1093 .bind(&[
1094 serde_json::to_string(&json!({ "role": "caller", "path": file, "source": source }))?.into(),
1095 format!("Calls `{file}`.").into(),
1096 now().into(),
1097 row.id.as_str().into(),
1098 ])?,
1099 );
1100 self.db.batch(statements).await?;
1101 Ok(())
1102 }
1103
1104 /// A job that called a workflow, finished with its jobs: their result,
1105 /// and the outputs the workflow declares.
1106 async fn finish_call(&self, row: &JobRow, children: &[&JobRow]) -> Result<()> {
1107 let call = row.call().unwrap_or_default();
1108 let called = call["source"].as_str().and_then(|s| workflow::parse(s).ok());
1109 let mut jobs_context = Map::new();
1110 let mut by_key: std::collections::BTreeMap<String, Vec<&JobRow>> = std::collections::BTreeMap::new();
1111 for child in children {
1112 by_key.entry(child.key.clone()).or_default().push(child);
1113 }
1114 for (key, rows) in &by_key {
1115 let mut outputs = Map::new();
1116 for child in rows {
1117 if let Ok(Value::Object(more)) = serde_json::from_str::<Value>(&child.outputs) {
1118 outputs.extend(more);
1119 }
1120 }
1121 let id = key.rsplit('/').next().unwrap_or(key);
1122 jobs_context.insert(id.to_owned(), json!({ "result": key_result(rows), "outputs": outputs }));
1123 }
1124 let inputs = children.first().and_then(|c| c.call()).map(|c| c["inputs"].clone()).unwrap_or(json!({}));
1125 let mut contexts = Map::new();
1126 contexts.insert("jobs".into(), Value::Object(jobs_context));
1127 contexts.insert("inputs".into(), inputs);
1128 let scope = Scope { contexts: &contexts, status: Status::Success, hash_files: None };
1129 let mut outputs = Map::new();
1130 if let Some(Value::Object(declared)) = called.as_ref().map(|w| {
1131 let on = w.raw.get("on").or_else(|| w.raw.get("true")).cloned().unwrap_or(Value::Null);
1132 on.get("workflow_call").and_then(|c| c.get("outputs")).cloned().unwrap_or(Value::Null)
1133 }) {
1134 for (name, spec) in declared {
1135 if let Some(value) = spec.get("value") {
1136 let value = expr::interpolate_value(value, &scope).unwrap_or(Value::Null);
1137 outputs.insert(name, Value::String(expr::to_text(&value)));
1138 }
1139 }
1140 }
1141 let conclusion = key_result(children);
1142 self.db
1143 .prepare("UPDATE jobs SET status = 'completed', conclusion = ?, outputs = ?, finished_at = ? WHERE id = ? AND status = 'calling'")
1144 .bind(&[conclusion.into(), serde_json::to_string(&outputs)?.into(), now().into(), row.id.as_str().into()])?
1145 .run()
1146 .await?;
1147 Ok(())
1148 }
1149
1150 /// A file's text at a commit, if it is there.
1151 async fn read_file(&self, path: &RepoPath, ws: &g1t_contracts::User, sha: &str, file: &str) -> Result<Option<String>> {
1152 let blob: Outcome<g1t_contracts::repos::BlobView> = g1t_kit::call(
1153 &self.repos,
1154 "blob",
1155 &g1t_contracts::repos::BlobArgs {
1156 path: path.clone(),
1157 viewer: Some(ws.clone()),
1158 git_ref: sha.to_owned(),
1159 file_path: file.to_owned(),
1160 },
1161 )
1162 .await?;
1163 Ok(match blob {
1164 Outcome::Ok(view) => view.text,
1165 Outcome::Fail(_) => None,
1166 })
1167 }
1168
GitHub Actions on g1t, part two: running workflows1169 async fn skip_job(&self, row: &JobRow, reason: Option<&str>) -> Result<()> {
1170 self.db
1171 .prepare("UPDATE jobs SET status = 'completed', conclusion = 'skipped', reason = ?, finished_at = ? WHERE id = ?")
1172 .bind(&[optional(reason), now().into(), row.id.as_str().into()])?
1173 .run()
1174 .await?;
1175 Ok(())
1176 }
1177
1178 async fn fail_job(&self, row: &JobRow, reason: &str) -> Result<()> {
1179 self.db
1180 .prepare("UPDATE jobs SET status = 'completed', conclusion = 'failure', reason = ?, finished_at = ? WHERE id = ? AND status != 'completed'")
1181 .bind(&[reason.into(), now().into(), row.id.as_str().into()])?
1182 .run()
1183 .await?;
1184 Ok(())
1185 }
1186
1187 /// Starts queued jobs, oldest first, while their workspace has room.
1188 pub async fn start_queued(&self) -> Result<()> {
1189 let queued = self
1190 .db
Fast pages, required checks on the branch, self-hosted runners, honest incidents1191 // Self-hosted jobs are taken by their runners (runners.rs).
1192 .prepare("SELECT * FROM jobs WHERE status = 'queued' AND labels IS NULL ORDER BY rowid LIMIT 50")
GitHub Actions on g1t, part two: running workflows1193 .all()
1194 .await?
1195 .results::<JobRow>()?;
1196 let mut running: std::collections::HashMap<String, u32> = std::collections::HashMap::new();
1197 for job in queued {
1198 let in_workspace = match running.get(&job.namespace) {
1199 Some(n) => *n,
1200 None => {
1201 let n = self
1202 .db
Fast pages, required checks on the branch, self-hosted runners, honest incidents1203 // Only g1t's own sandboxes count against the workspace's room.
1204 .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 workflows1205 .bind(&[job.namespace.as_str().into()])?
1206 .first::<Count>(None)
1207 .await?
1208 .map_or(0, |count| count.n);
1209 running.insert(job.namespace.clone(), n);
1210 n
1211 }
1212 };
1213 if in_workspace >= RUNNING_PER_WORKSPACE {
1214 continue;
1215 }
1216 if let Some(max) = job.max_parallel {
1217 let siblings = self
1218 .db
1219 .prepare("SELECT COUNT(*) AS n FROM jobs WHERE run_id = ? AND key = ? AND status = 'in_progress'")
1220 .bind(&[job.run_id.as_str().into(), job.key.as_str().into()])?
1221 .first::<Count>(None)
1222 .await?
1223 .map_or(0, |count| count.n);
1224 if siblings >= max {
1225 continue;
1226 }
1227 }
Merge branch 'worktree-agent-a3abfcce648e87dca'1228 // One job of a concurrency group runs at a time.
1229 if let Some(group) = &job.concurrency_group {
1230 let running = self
1231 .db
1232 .prepare("SELECT COUNT(*) AS n FROM jobs WHERE repo_id = ? AND concurrency_group = ? AND status = 'in_progress' AND id != ?")
1233 .bind(&[job.repo_id.as_str().into(), group.as_str().into(), job.id.as_str().into()])?
1234 .first::<Count>(None)
1235 .await?
1236 .map_or(0, |count| count.n);
1237 if running > 0 {
1238 continue;
1239 }
1240 }
GitHub Actions on g1t, part two: running workflows1241 let token = random_hex(24);
1242 let at = now();
1243 let claimed = self
1244 .db
1245 .prepare(
1246 "UPDATE jobs SET status = 'in_progress', token_hash = ?, started_at = ?, seen_at = ? WHERE id = ? AND status = 'queued' RETURNING id",
1247 )
1248 .bind(&[sha256_hex(&token).into(), at.as_str().into(), at.as_str().into(), job.id.as_str().into()])?
1249 .first::<Value>(None)
1250 .await?;
1251 if claimed.is_none() {
1252 continue;
1253 }
1254 running.insert(job.namespace.clone(), in_workspace + 1);
1255 self.db
1256 .prepare("UPDATE runs SET status = 'in_progress', started_at = COALESCE(started_at, ?) WHERE id = ? AND status = 'queued'")
1257 .bind(&[at.as_str().into(), job.run_id.as_str().into()])?
1258 .run()
1259 .await?;
1260 let run = self.run_row(&job.run_id).await?;
1261 let repo: RepoPath = run.as_ref().map(|run| repo_path(&run.repo)).unwrap_or(RepoPath {
1262 namespace: job.namespace.clone(),
1263 name: String::new(),
1264 });
Fast pages, required checks on the branch, self-hosted runners, honest incidents1265 // Its environment and the machine it asked for, for the runner.
1266 let details = run.as_ref().map(|run| start_details(run, &job)).unwrap_or_default();
GitHub Actions on g1t, part two: running workflows1267 let started: Outcome<Value> = g1t_kit::call(
1268 &self.runner,
1269 "start_actions_job",
1270 &StartJobArgs {
1271 job: job.id.clone(),
1272 token,
1273 repo,
1274 timeout_minutes: job.timeout_minutes,
Fast pages, required checks on the branch, self-hosted runners, honest incidents1275 workflow: run.as_ref().map(|run| run.path.clone()),
1276 environment: details.environment,
1277 trusted: run.as_ref().is_some_and(|run| run.trusted != 0),
1278 instance: details.instance,
GitHub Actions on g1t, part two: running workflows1279 },
1280 )
1281 .await
1282 .unwrap_or_else(|error| fail(FailureCode::Conflict, format!("The runner could not be reached: {error}")));
1283 if let Outcome::Fail(refused) = started {
1284 Box::pin(self.finish_job(&job.id, "failure", Some(&refused.message), None)).await?;
Merge branch 'main' into worktree-agent-a69aeabc4b0deeb971285 } else {
1286 self.job_started(&job.id).await?;
GitHub Actions on g1t, part two: running workflows1287 }
1288 }
1289 Ok(())
1290 }
1291
1292 /// Finishes a job and moves its run along.
1293 pub async fn finish_job(&self, job_id: &str, conclusion: &str, reason: Option<&str>, outputs: Option<&Map<String, Value>>) -> Result<()> {
1294 let finished = self
1295 .db
1296 .prepare(
1297 "UPDATE jobs SET status = 'completed', conclusion = ?, reason = COALESCE(?, reason), outputs = COALESCE(?, outputs),
1298 finished_at = ?, token_hash = NULL WHERE id = ? AND status != 'completed' RETURNING *",
1299 )
1300 .bind(&[
1301 conclusion.into(),
1302 optional(reason),
1303 outputs.map(|o| serde_json::to_string(o).unwrap_or_default()).as_deref().map_or(worker::wasm_bindgen::JsValue::NULL, Into::into),
1304 now().into(),
1305 job_id.into(),
1306 ])?
1307 .first::<JobRow>(None)
1308 .await?;
1309 let Some(job) = finished else { return Ok(()) };
Merge branch 'worktree-agent-a3abfcce648e87dca'1310 // Its G1T_TOKEN stops working with it.
1311 self.revoke_job_tokens(&job.id).await;
Merge branch 'main' into worktree-agent-a69aeabc4b0deeb971312 // A job that deploys failed: so did its run's deployment, now.
1313 if conclusion == "failure" && job.continue_on_error == 0 && job.started_at.is_some()
1314 && let Some(run) = self.run_row(&job.run_id).await?
1315 && let Some(env) = deploys_to(&run, &job)
1316 {
1317 self.report_deployment(&run, &env, "failure", false).await;
1318 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents1319 // A self-hosted runner's job: the runner is free again, and its time
1320 // is recorded, at nothing.
1321 if job.runner_id.is_some() {
1322 self.released(&job).await?;
1323 }
GitHub Actions on g1t, part two: running workflows1324 // Steps still marked as going are not going any more.
1325 let mut steps: Vec<Value> = serde_json::from_str(&job.steps).unwrap_or_default();
1326 let mut touched = false;
1327 for step in steps.iter_mut() {
1328 if step["status"] != "completed" {
1329 let was_running = step["status"] == "in_progress";
1330 step["status"] = json!("completed");
1331 step["conclusion"] = json!(if was_running { conclusion } else { "skipped" });
1332 touched = true;
1333 }
1334 }
1335 if touched {
1336 self.db
1337 .prepare("UPDATE jobs SET steps = ? WHERE id = ?")
1338 .bind(&[serde_json::to_string(&steps)?.into(), job.id.as_str().into()])?
1339 .run()
1340 .await?;
1341 }
1342 // fail-fast: one failed combination stops the rest of its matrix.
1343 if conclusion == "failure" && job.continue_on_error == 0 && job.matrix.as_deref().is_some_and(|m| m != "{}") {
1344 let run = self.run_row(&job.run_id).await?;
1345 let fail_fast = run
1346 .as_ref()
1347 .and_then(|run| workflow::parse(&run.source).ok())
1348 .and_then(|workflow| workflow.jobs.into_iter().find(|j| j.id == job.key))
1349 .is_none_or(|j| j.fail_fast);
1350 if fail_fast {
1351 let siblings = self
1352 .db
1353 .prepare("SELECT * FROM jobs WHERE run_id = ? AND key = ? AND status != 'completed'")
1354 .bind(&[job.run_id.as_str().into(), job.key.as_str().into()])?
1355 .all()
1356 .await?
1357 .results::<JobRow>()?;
1358 for sibling in siblings {
1359 self.stop_job(&sibling, "Another job of its matrix failed, and the matrix is fail-fast.").await?;
1360 }
1361 }
1362 }
1363 Box::pin(self.advance(&job.run_id)).await
1364 }
1365
1366 /// Cancels a job, stopping its sandbox if it has one.
1367 async fn stop_job(&self, job: &JobRow, reason: &str) -> Result<()> {
Fast pages, required checks on the branch, self-hosted runners, honest incidents1368 // A self-hosted runner hears it was cancelled on its next poll.
1369 if job.status == "in_progress" && job.runner_id.is_none() {
GitHub Actions on g1t, part two: running workflows1370 let _: Result<Value> = g1t_kit::call(&self.runner, "stop_actions_job", &json!({ "job": job.id })).await;
1371 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents1372 let stopped = self
1373 .db
1374 .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 workflows1375 .bind(&[reason.into(), now().into(), job.id.as_str().into()])?
Fast pages, required checks on the branch, self-hosted runners, honest incidents1376 .first::<JobRow>(None)
GitHub Actions on g1t, part two: running workflows1377 .await?;
Merge branch 'worktree-agent-a3abfcce648e87dca'1378 if stopped.as_ref().is_some_and(|row| row.started_at.is_some()) {
1379 self.revoke_job_tokens(&job.id).await;
1380 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents1381 if let Some(stopped) = stopped.filter(|row| row.runner_id.is_some()) {
1382 self.released(&stopped).await?;
1383 }
GitHub Actions on g1t, part two: running workflows1384 Ok(())
1385 }
1386
Merge branch 'worktree-agent-a3abfcce648e87dca'1387 /// Ends a job's tokens at once. A failure is logged: the token expires
1388 /// on its own soon after the job's time limit.
1389 async fn revoke_job_tokens(&self, job_id: &str) {
1390 let revoked: Result<bool> =
1391 g1t_kit::call(&self.identity, "revoke_job_tokens", &RevokeJobTokensArgs { job_id: job_id.to_owned() }).await;
1392 if let Err(error) = revoked {
1393 worker::console_error!("actions: the tokens of job {job_id} were not revoked: {error}");
1394 }
1395 }
1396
GitHub Actions on g1t, part two: running workflows1397 /// Finishes the run when every job has.
1398 async fn finish_if_done(&self, run_id: &str) -> Result<()> {
1399 let Some(run) = self.run_row(run_id).await? else { return Ok(()) };
Merge branch 'worktree-agent-a3abfcce648e87dca'1400 if matches!(run.status.as_str(), "completed" | "pending" | "action_required") {
GitHub Actions on g1t, part two: running workflows1401 return Ok(());
1402 }
1403 let jobs = self.job_rows(run_id).await?;
1404 if !jobs.iter().all(|job| job.status == "completed") {
1405 return Ok(());
1406 }
1407 self.finish_run(&run, None).await
1408 }
1409
1410 async fn finish_run(&self, run: &RunRow, error: Option<&str>) -> Result<()> {
1411 let jobs = self.job_rows(&run.id).await?;
1412 let rows: Vec<&JobRow> = jobs.iter().collect();
1413 let conclusion = if error.is_some() {
1414 "failure"
1415 } else if run.conclusion.as_deref() == Some("cancelled") {
1416 "cancelled"
1417 } else if rows.is_empty() {
1418 "skipped"
1419 } else {
1420 key_result(&rows)
1421 };
1422 let done = self
1423 .db
1424 .prepare("UPDATE runs SET status = 'completed', conclusion = ?, error = COALESCE(?, error), finished_at = ? WHERE id = ? AND status != 'completed' RETURNING id")
1425 .bind(&[conclusion.into(), optional(error), now().into(), run.id.as_str().into()])?
1426 .first::<Value>(None)
1427 .await?;
1428 if done.is_none() {
1429 return Ok(());
1430 }
1431 self.report_status(run, conclusion).await?;
Merge branch 'main' into worktree-agent-a69aeabc4b0deeb971432 self.settle_deployments(run, &jobs, conclusion == "cancelled").await;
Actions: workflow_run, workflow.completed, artifacts on the run page, Node 241433 let published: Result<()> = g1t_kit::call(
1434 &self.events,
1435 "publish",
1436 &g1t_contracts::events::Publish {
1437 events: vec![g1t_contracts::events::NewEvent {
1438 kind: "workflow.completed",
1439 source: "actions",
1440 repo_id: Some(run.repo_id.clone()),
1441 actor: run.actor_id.clone(),
1442 data: g1t_contracts::events::WorkflowEvent {
1443 run_id: run.id.clone(),
1444 repo_id: run.repo_id.clone(),
1445 workflow: run.name.clone(),
1446 path: run.path.clone(),
1447 number: run.number,
1448 event: run.event.clone(),
1449 conclusion: conclusion.to_owned(),
1450 git_ref: run.git_ref.clone(),
1451 sha: run.sha.clone(),
1452 pull: run.pull,
1453 },
1454 }],
1455 },
1456 )
1457 .await;
1458 if let Err(error) = published {
1459 worker::console_error!("actions: could not publish workflow.completed: {error}");
1460 }
GitHub Actions on g1t, part two: running workflows1461 // The next run waiting in its concurrency group.
1462 if let Some(group) = &run.concurrency_group {
1463 let next = self
1464 .db
1465 .prepare("SELECT * FROM runs WHERE repo_id = ? AND concurrency_group = ? AND status = 'pending' ORDER BY id LIMIT 1")
1466 .bind(&[run.repo_id.as_str().into(), group.as_str().into()])?
1467 .first::<RunRow>(None)
1468 .await?;
1469 if let Some(next) = next {
1470 self.db.prepare("UPDATE runs SET status = 'queued' WHERE id = ?").bind(&[next.id.as_str().into()])?.run().await?;
1471 Box::pin(self.advance(&next.id)).await?;
1472 }
1473 }
1474 Ok(())
1475 }
1476
Merge branch 'main' into worktree-agent-a69aeabc4b0deeb971477 /// Tells the deployments service how a run's deployment to `env` stands:
1478 /// made when its first job naming the environment starts, failed when one
1479 /// of them fails, and settled (`last`) when the run finishes. Never fails
1480 /// the run: a deployment that could not be recorded is logged.
1481 async fn report_deployment(&self, run: &RunRow, env: &JobEnvironment, state: &str, last: bool) {
1482 let path = repo_path(&run.repo);
1483 let reported: Result<Value> = g1t_kit::call(
1484 &self.deployments,
1485 "actions_deployment",
1486 &json!({
1487 "repoId": run.repo_id,
1488 "repo": { "namespace": path.namespace, "name": path.name },
1489 "runId": run.id,
1490 "attempt": run.attempt,
1491 "runUrl": format!("{SITE}/{}/actions/runs/{}", run.repo, run.id),
1492 "environment": env.name,
1493 "url": env.url,
1494 "ref": run.git_ref,
1495 "sha": run.sha,
1496 "state": state,
1497 "final": last,
1498 "creator": run.actor,
1499 "workflow": run.name,
1500 }),
1501 )
1502 .await;
1503 if let Err(error) = reported {
1504 worker::console_error!("actions: deployment to {} not recorded for run {}: {error}", env.name, run.id);
1505 }
1506 }
1507
1508 /// A job that deploys has started: its run's deployment to the
1509 /// environment is under way.
1510 pub(crate) async fn job_started(&self, job_id: &str) -> Result<()> {
1511 let Some(job) = self.db.prepare("SELECT * FROM jobs WHERE id = ?").bind(&[job_id.into()])?.first::<JobRow>(None).await? else {
1512 return Ok(());
1513 };
1514 let Some(run) = self.run_row(&job.run_id).await? else { return Ok(()) };
1515 if let Some(env) = deploys_to(&run, &job) {
1516 self.report_deployment(&run, &env, "in_progress", false).await;
1517 }
1518 Ok(())
1519 }
1520
1521 /// The run is over: each environment its jobs deployed to takes the
1522 /// outcome of those jobs (`deployment_outcome`).
1523 async fn settle_deployments(&self, run: &RunRow, jobs: &[JobRow], cancelled: bool) {
1524 let mut seen: Vec<(JobEnvironment, Vec<Option<String>>)> = Vec::new();
1525 for job in jobs {
1526 let Some(env) = deploys_to(run, job) else { continue };
1527 let conclusion = if cancelled && job.conclusion.as_deref() != Some("skipped") && job.started_at.is_some() {
1528 Some("cancelled".to_owned())
1529 } else if job.started_at.is_none() {
1530 Some("skipped".to_owned())
1531 } else {
1532 job.conclusion.clone()
1533 };
1534 match seen.iter_mut().find(|(known, _)| known.name.eq_ignore_ascii_case(&env.name)) {
1535 Some((known, conclusions)) => {
1536 if known.url.is_none() {
1537 known.url = env.url.clone();
1538 }
1539 conclusions.push(conclusion);
1540 }
1541 None => seen.push((env, vec![conclusion])),
1542 }
1543 }
1544 for (env, conclusions) in seen {
1545 if let Some(state) = deployment_outcome(&conclusions) {
1546 self.report_deployment(run, &env, state, true).await;
1547 }
1548 }
1549 }
1550
GitHub Actions on g1t, part two: running workflows1551 /// Tells the pull request (or commit) how the run went, as a status.
1552 async fn report_status(&self, run: &RunRow, conclusion: &str) -> Result<()> {
1553 let state = match conclusion {
1554 "success" | "skipped" => "success",
1555 "cancelled" => "error",
1556 _ => "failure",
1557 };
1558 let _: Result<Value> = g1t_kit::call(
1559 &self.work,
1560 "set_commit_status",
1561 &json!({
1562 "repoId": run.repo_id,
1563 "sha": run.sha,
1564 "context": format!("{} / {}", run.name, run.event),
1565 "state": state,
1566 "description": format!("{} {}", run.name, match conclusion {
1567 "success" => "passed",
1568 "skipped" => "was skipped",
1569 "cancelled" => "was cancelled",
1570 _ => "failed",
1571 }),
1572 "targetUrl": format!("{SITE}/{}/actions/runs/{}", run.repo, run.id),
Merge rulesets: branch and tag rules, agent-first, enforced on push and merge1573 "source": "actions",
GitHub Actions on g1t, part two: running workflows1574 }),
1575 )
1576 .await;
1577 Ok(())
1578 }
1579
1580 /// Tells the pull request a run has started on its head.
1581 pub async fn report_pending(&self, run: &RunRow) -> Result<()> {
1582 let _: Result<Value> = g1t_kit::call(
1583 &self.work,
1584 "set_commit_status",
1585 &json!({
1586 "repoId": run.repo_id,
1587 "sha": run.sha,
1588 "context": format!("{} / {}", run.name, run.event),
1589 "state": "pending",
Merge branch 'worktree-agent-a3abfcce648e87dca'1590 "description": if run.status == "action_required" {
1591 format!("{} is waiting for approval", run.name)
1592 } else {
1593 format!("{} is running", run.name)
1594 },
GitHub Actions on g1t, part two: running workflows1595 "targetUrl": format!("{SITE}/{}/actions/runs/{}", run.repo, run.id),
Merge rulesets: branch and tag rules, agent-first, enforced on push and merge1596 "source": "actions",
GitHub Actions on g1t, part two: running workflows1597 }),
1598 )
1599 .await;
1600 Ok(())
1601 }
1602
1603 /// Cancels a run: its waiting and queued jobs, and stops its running ones.
1604 pub async fn cancel_run(&self, run: &RunRow, reason: &str) -> Result<()> {
1605 self.db
1606 .prepare("UPDATE runs SET conclusion = 'cancelled' WHERE id = ? AND status != 'completed'")
1607 .bind(&[run.id.as_str().into()])?
1608 .run()
1609 .await?;
1610 for job in self.job_rows(&run.id).await?.iter().filter(|job| job.status != "completed") {
1611 self.stop_job(job, reason).await?;
1612 }
Merge branch 'worktree-agent-a3abfcce648e87dca'1613 if run.status == "pending" || run.status == "action_required" {
GitHub Actions on g1t, part two: running workflows1614 self.db.prepare("UPDATE runs SET status = 'queued' WHERE id = ?").bind(&[run.id.as_str().into()])?.run().await?;
1615 }
1616 self.advance(&run.id).await
1617 }
1618
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look1619 /// Cancels every run of the repository that has not finished, for
1620 /// `repo.deleted` and `repo.archived`.
1621 pub async fn stop_runs(&self, repo_id: &str) -> Result<()> {
1622 let runs = self
1623 .db
1624 .prepare("SELECT * FROM runs WHERE repo_id = ? AND status != 'completed'")
1625 .bind(&[repo_id.into()])?
1626 .all()
1627 .await?
1628 .results::<RunRow>()?;
1629 for run in runs {
1630 self.cancel_run(&run, "The repository was archived or deleted.").await?;
1631 }
1632 Ok(())
1633 }
1634
GitHub Actions on g1t, part two: running workflows1635 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 look1636 if let Outcome::Fail(refused) = self.may(&a.actor, &a.repo, Capability::Run).await? {
GitHub Actions on g1t, part two: running workflows1637 return Ok(Outcome::Fail(refused));
1638 }
1639 let run = check!(self.run_in(&a.repo, &a.id).await?);
1640 if run.status == "completed" {
1641 return Ok(fail(FailureCode::Conflict, "The run has already finished."));
1642 }
1643 self.cancel_run(&run, &format!("{} cancelled the run.", a.actor.username)).await?;
1644 self.run_summary(&run.id).await
1645 }
1646
1647 /// Runs again: every job, or with `failed_only` those that did not
1648 /// succeed and the jobs that need them.
1649 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 look1650 if let Outcome::Fail(refused) = self.may(&a.actor, &a.repo, Capability::Run).await? {
GitHub Actions on g1t, part two: running workflows1651 return Ok(Outcome::Fail(refused));
1652 }
1653 let run = check!(self.run_in(&a.repo, &a.id).await?);
1654 if run.status != "completed" {
1655 return Ok(fail(FailureCode::Conflict, "The run is still going: cancel it first."));
1656 }
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look1657 // Nothing starts again on an archived repository.
1658 match self.visible_repo(&a.repo, &Some(a.actor.clone())).await? {
1659 Some(repo) if repo.archived() => {
1660 return Ok(fail(FailureCode::Forbidden, g1t_contracts::repos::archived_message(&repo.namespace, &repo.name)));
1661 }
1662 Some(_) => {}
1663 None => return Ok(fail(FailureCode::NotFound, "There is no such repository.")),
1664 }
GitHub Actions on g1t, part two: running workflows1665 if run.error.is_some() {
1666 return Ok(fail(FailureCode::Conflict, "This run never started: fix the workflow file and push again."));
1667 }
1668 let jobs = self.job_rows(&run.id).await?;
1669 let workflow = workflow::parse(&run.source).ok();
1670 // Which keys run again: failed ones and, transitively, those needing them.
1671 let mut again: Vec<String> = Vec::new();
1672 for key in workflow.as_ref().map(|w| w.job_order()).unwrap_or_default() {
1673 let rows: Vec<&JobRow> = jobs.iter().filter(|j| j.key == key).collect();
1674 let failed = rows.iter().any(|row| row.conclusion.as_deref() != Some("success"));
1675 let needs_again = rows.first().is_some_and(|row| row.needs().iter().any(|need| again.contains(need)));
1676 if !a.failed_only || failed || needs_again {
1677 again.push(key.to_owned());
1678 }
1679 }
1680 if again.is_empty() {
1681 return Ok(fail(FailureCode::Conflict, "Every job succeeded: there is nothing to run again."));
1682 }
1683 let mut statements = Vec::new();
1684 for key in &again {
1685 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()])?);
1686 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 repository1687 // The jobs of a workflow it called are made again when it calls it again.
GitHub Actions on g1t, part two: running workflows1688 statements.push(
1689 self.db
Actions: reusable workflows in the repository1690 .prepare("DELETE FROM logs WHERE job_id IN (SELECT id FROM jobs WHERE run_id = ? AND key LIKE ?)")
1691 .bind(&[run.id.as_str().into(), format!("{key}/%").into()])?,
1692 );
1693 statements.push(
1694 self.db
1695 .prepare("DELETE FROM jobs WHERE run_id = ? AND key LIKE ?")
1696 .bind(&[run.id.as_str().into(), format!("{key}/%").into()])?,
1697 );
1698 statements.push(
1699 self.db
GitHub Actions on g1t, part two: running workflows1700 .prepare(
1701 "UPDATE jobs SET status = 'waiting', conclusion = NULL, steps = '[]', annotations = '[]', outputs = '{}', reason = NULL,
Fast pages, required checks on the branch, self-hosted runners, honest incidents1702 matrix = NULL, call = NULL, token_hash = NULL, seen_at = NULL, started_at = NULL, finished_at = NULL,
Merge branch 'worktree-agent-a3abfcce648e87dca'1703 labels = NULL, queued_at = NULL, runner_id = NULL, runner_name = NULL, environment = NULL,
1704 concurrency_group = NULL, cancel_in_progress = 0 WHERE run_id = ? AND key = ?",
GitHub Actions on g1t, part two: running workflows1705 )
1706 .bind(&[run.id.as_str().into(), key.as_str().into()])?,
1707 );
1708 }
1709 statements.push(
1710 self.db
1711 .prepare("UPDATE runs SET status = 'queued', conclusion = NULL, attempt = attempt + 1, started_at = NULL, finished_at = NULL WHERE id = ?")
1712 .bind(&[run.id.as_str().into()])?,
1713 );
1714 self.db.batch(statements).await?;
1715 if let Some(run) = self.run_row(&run.id).await? {
1716 self.report_pending(&run).await?;
1717 }
1718 self.advance(&run.id).await?;
1719 self.run_summary(&run.id).await
1720 }
1721
1722 pub async fn run_in(&self, repo: &RepoPath, id: &str) -> Result<Outcome<RunRow>> {
1723 let row = self
1724 .db
1725 .prepare("SELECT * FROM runs WHERE id = ? AND lower(repo) = lower(?)")
1726 .bind(&[id.into(), format!("{}/{}", repo.namespace, repo.name).into()])?
1727 .first::<RunRow>(None)
1728 .await?;
1729 Ok(row.map_or_else(|| fail(FailureCode::NotFound, "No such run."), Outcome::Ok))
1730 }
1731
1732 // --- The sandbox's side -----------------------------------------------------
1733
Fast pages, required checks on the branch, self-hosted runners, honest incidents1734 pub(crate) async fn job_for_token(&self, a: &JobCallArgs) -> Result<Outcome<JobRow>> {
GitHub Actions on g1t, part two: running workflows1735 let job = self.db.prepare("SELECT * FROM jobs WHERE id = ?").bind(&[a.job.as_str().into()])?.first::<JobRow>(None).await?;
1736 Ok(match job {
1737 Some(job) if job.status == "in_progress" && job.token_hash.as_deref().is_some_and(|hash| same(hash, &sha256_hex(&a.token))) => {
1738 Outcome::Ok(job)
1739 }
1740 _ => fail(FailureCode::Unauthenticated, "That job is not running, or the token is not its."),
1741 })
1742 }
1743
A repository has its own sidebar, as settings do1744 /// `job_auth`: which run and repository a running job's token is for,
1745 /// so the API can keep its artifacts and cache.
1746 pub async fn job_auth(&self, a: JobCallArgs) -> Result<Outcome<Value>> {
1747 let job = check!(self.job_for_token(&a).await?);
1748 Ok(Outcome::Ok(json!({ "run": job.run_id, "repoId": job.repo_id })))
1749 }
1750
GitHub Actions on g1t, part two: running workflows1751 /// `job_spec`: everything the sandbox needs to run the job.
1752 pub async fn job_spec(&self, a: JobCallArgs) -> Result<Outcome<Value>> {
1753 let job = check!(self.job_for_token(&a).await?);
1754 let Some(run) = self.run_row(&job.run_id).await? else {
1755 return Ok(fail(FailureCode::NotFound, "No such run."));
1756 };
Actions: reusable workflows in the repository1757 let Ok(caller) = workflow::parse(&run.source) else {
GitHub Actions on g1t, part two: running workflows1758 return Ok(fail(FailureCode::Invalid, "The workflow no longer reads."));
1759 };
Actions: reusable workflows in the repository1760 // A called workflow's job runs as that workflow defines it.
1761 let callee = job.callee();
1762 let (workflow, spec, call_inputs) = match callee {
1763 Some((called, spec, call)) => (called, spec, Some(call["inputs"].clone())),
1764 None => match caller.jobs.iter().find(|j| j.id == job.key) {
1765 Some(spec) => (caller.clone(), spec.clone(), None),
1766 None => return Ok(fail(FailureCode::NotFound, "The job is not in the workflow.")),
1767 },
GitHub Actions on g1t, part two: running workflows1768 };
Actions: reusable workflows in the repository1769 let spec = &spec;
GitHub Actions on g1t, part two: running workflows1770 let repo = repo_path(&run.repo);
1771 let trusted = run.trusted != 0;
Merge branch 'worktree-agent-a3abfcce648e87dca'1772 // The job's `environment:`, by name, as read when its needs were done
1773 // (an expression included), once the environment's protection rules
1774 // let it start: entries with a value for it give that value instead
1775 // of their default, as GitHub's environment secrets do.
1776 let environment: Option<String> = job.environment.clone();
1777 // What its token may do: its `permissions:` (a called workflow's
1778 // jobs no more than the job that calls it), else the repository's
1779 // default; read-only for a pull request from outside.
1780 // The repository's default, its workspace's taken in, and whether
1781 // its jobs may open and approve pull requests.
1782 let (default, pull_requests) = self.token_policy(&run.repo_id).await?;
1783 let mut permissions = spec.permissions(&workflow, default);
1784 if let Some(parent) = &job.call().filter(|c| c["role"] == "callee").and_then(|c| c["parent"].as_str().map(str::to_owned)) {
1785 let top = parent.split('/').next().unwrap_or(parent);
1786 if let Some(caller_job) = caller.jobs.iter().find(|j| j.id == top) {
1787 permissions = permissions.capped_by(&caller_job.permissions(&caller, default));
GitHub Actions on g1t, part two: running workflows1788 }
Merge branch 'worktree-agent-a3abfcce648e87dca'1789 }
1790 if !trusted {
1791 permissions = permissions.read_only();
1792 }
1793 // G1T_TOKEN, and GITHUB_TOKEN as its alias: a token of the
1794 // workspace's that reaches this repository only, with the scopes
1795 // its permissions give, until the job ends.
1796 let token = match self.workspace_actor(&repo.namespace).await? {
1797 Some(workspace) => {
1798 let created: CreatedAccessToken = g1t_kit::call(
1799 &self.identity,
1800 "create_job_token",
1801 &CreateJobTokenArgs {
1802 workspace,
1803 repo: repo.clone(),
1804 run_id: run.id.clone(),
1805 job_id: job.id.clone(),
1806 name: format!("G1T_TOKEN for {} run {}", run.repo, run.number),
1807 ttl_seconds: u64::from(job.timeout_minutes) * 60 + 600,
1808 scopes: permissions.scopes().into_iter().map(str::to_owned).collect(),
1809 pull_requests: pull_requests && trusted,
1810 },
1811 )
1812 .await?;
1813 created.token
1814 }
1815 None => String::new(),
GitHub Actions on g1t, part two: running workflows1816 };
Secrets and variables: one list, rows per environment, for workflows and deployments1817 // A run that is not trusted (a pull request from outside the
1818 // workspace) gets no secrets and an empty token.
1819 let mut secrets = if trusted {
1820 self.secrets_for(&run.repo_id, &run.repo, environment.as_deref(), true).await?
1821 } else {
1822 Map::new()
1823 };
1824 secrets.insert("G1T_TOKEN".into(), Value::String(token.clone()));
GitHub Actions on g1t, part two: running workflows1825 secrets.insert("GITHUB_TOKEN".into(), Value::String(token.clone()));
Merge branch 'worktree-agent-a3abfcce648e87dca'1826 // Each secret as it is, a line at a time, base64 and JSON-escaped.
1827 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 deployments1828 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 R21829 // The toolkit's runtime token (runtime.rs), for as long as the job
1830 // may run; the API puts it and the toolkit's addresses in the
1831 // job's variables.
1832 let runtime_token = crate::runtime::runtime_token(
1833 &job.id,
1834 &job.run_id,
1835 job.token_hash.as_deref().unwrap_or_default(),
1836 now_ms() / 1000,
1837 u64::from(job.timeout_minutes) * 60 + 600,
1838 );
1839 masks.push(runtime_token.clone());
1840 let retention_days = self.retention_setting(&run.repo_id).await?;
GitHub Actions on g1t, part two: running workflows1841
1842 let jobs = self.job_rows(&run.id).await?;
1843 let mut needs = Map::new();
Actions: reusable workflows in the repository1844 // In a called workflow, its jobs' keys sit under the job that called it.
1845 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 workflows1846 for need in &spec.needs {
Actions: reusable workflows in the repository1847 let key = match &parent {
1848 Some(parent) => format!("{parent}/{need}"),
1849 None => need.clone(),
1850 };
1851 let rows: Vec<&JobRow> = jobs.iter().filter(|row| row.key == key).collect();
GitHub Actions on g1t, part two: running workflows1852 let mut outputs = Map::new();
1853 for row in &rows {
1854 if let Ok(Value::Object(more)) = serde_json::from_str::<Value>(&row.outputs) {
1855 outputs.extend(more);
1856 }
1857 }
1858 needs.insert(need.clone(), json!({ "result": key_result(&rows), "outputs": outputs }));
1859 }
1860 let siblings = jobs.iter().filter(|row| row.key == job.key).count();
1861 let matrix: Value = job.matrix.as_deref().and_then(|m| serde_json::from_str(m).ok()).unwrap_or(json!({}));
1862 let info = run.info();
1863 let mut github = info.context(&job.key, &token, run.action.as_deref());
1864 github["token"] = json!(token);
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R21865 github["retention_days"] = json!(retention_days);
Fast pages, required checks on the branch, self-hosted runners, honest incidents1866 // On a self-hosted runner, `runner` and `RUNNER_*` describe that
1867 // machine rather than g1t's sandbox.
1868 let mut variables = info.variables(&job.key);
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R21869 variables.insert("GITHUB_RETENTION_DAYS".into(), json!(retention_days.to_string()));
Fast pages, required checks on the branch, self-hosted runners, honest incidents1870 let runner = match &job.runner_id {
1871 Some(id) => self.runner_context_for(id, &mut variables).await?,
1872 None => runner_context(),
1873 };
GitHub Actions on g1t, part two: running workflows1874
1875 // Where to check out: a pull request's fork, or the repository.
1876 let clone_url = match run.pull {
1877 Some(number) if run.event.starts_with("pull_request") && run.event != "pull_request_target" => {
1878 let located: Outcome<g1t_contracts::work::PullDetail> = g1t_kit::call(
1879 &self.work,
1880 "get_pull",
1881 &g1t_contracts::work::ViewArgs {
1882 repo: repo.clone(),
1883 number,
1884 viewer: self.workspace_actor(&repo.namespace).await?,
1885 after_seq: 0,
1886 },
1887 )
1888 .await?;
1889 match located {
1890 Outcome::Ok(detail) => match detail.pull.fork {
1891 Some(fork) => format!("{SITE}/{}/{}.git", fork.namespace, fork.name),
1892 None => format!("{SITE}/{}.git", run.repo),
1893 },
1894 Outcome::Fail(_) => format!("{SITE}/{}.git", run.repo),
1895 }
1896 }
1897 _ => format!("{SITE}/{}.git", run.repo),
1898 };
1899
1900 Ok(Outcome::Ok(json!({
1901 "job": job.id,
1902 "run": run.id,
1903 "key": job.key,
1904 "name": job.name,
1905 "spec": spec.raw,
1906 "workflow": {
1907 "env": workflow.env,
1908 "defaults": workflow.raw.get("defaults").cloned().unwrap_or(Value::Null),
1909 },
1910 "github": github,
Fast pages, required checks on the branch, self-hosted runners, honest incidents1911 "variables": variables,
GitHub Actions on g1t, part two: running workflows1912 "event": info.event,
1913 "contexts": {
1914 "vars": vars,
1915 "secrets": secrets,
Actions: reusable workflows in the repository1916 "inputs": call_inputs.unwrap_or_else(|| Value::Object(run.inputs())),
GitHub Actions on g1t, part two: running workflows1917 "matrix": matrix,
1918 "needs": needs,
1919 "strategy": {
1920 "fail-fast": spec.fail_fast,
1921 "job-index": job.ordinal,
1922 "job-total": siblings,
1923 "max-parallel": spec.max_parallel.unwrap_or(siblings as u32),
1924 },
Fast pages, required checks on the branch, self-hosted runners, honest incidents1925 "runner": runner,
GitHub Actions on g1t, part two: running workflows1926 },
1927 "checkout": {
1928 "repository": run.repo,
1929 "url": clone_url,
1930 "sha": run.sha,
1931 "ref": run.git_ref,
1932 "token": token,
1933 },
1934 "timeoutMinutes": job.timeout_minutes,
1935 "masks": masks,
Merge branch 'worktree-agent-a3abfcce648e87dca'1936 // As the job's log lists them at its start.
1937 "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 R21938 // Whether the job may ask for an OIDC token decides whether it
1939 // is told where to.
1940 "runtime": {
1941 "token": runtime_token,
1942 "idToken": self.oidc_allowed(&run, &job),
1943 },
GitHub Actions on g1t, part two: running workflows1944 })))
1945 }
1946
1947 /// `job_report`: the sandbox telling how the job is going.
1948 pub async fn job_report(&self, a: JobCallArgs) -> Result<Outcome<Value>> {
1949 let job = check!(self.job_for_token(&a).await?);
1950 let report = &a.report;
1951 let at = now();
1952 match report["kind"].as_str().unwrap_or_default() {
1953 "steps" => {
1954 // The list can grow as the job goes (post steps), so steps
1955 // already reported keep where they stand.
1956 let known: Vec<Value> = serde_json::from_str(&job.steps).unwrap_or_default();
1957 let steps: Vec<Value> = report["steps"]
1958 .as_array()
1959 .map(|names| {
1960 names
1961 .iter()
1962 .enumerate()
1963 .map(|(i, name)| match known.get(i) {
1964 Some(step) if step["status"] != "queued" => step.clone(),
1965 _ => json!({ "number": i + 1, "name": expr::to_text(name), "status": "queued", "conclusion": null, "startedAt": null, "finishedAt": null }),
1966 })
1967 .collect()
1968 })
1969 .unwrap_or_default();
1970 self.db
1971 .prepare("UPDATE jobs SET steps = ?, seen_at = ? WHERE id = ?")
1972 .bind(&[serde_json::to_string(&steps)?.into(), at.as_str().into(), job.id.as_str().into()])?
1973 .run()
1974 .await?;
1975 }
1976 "step" => {
1977 let number = report["number"].as_u64().unwrap_or(0) as usize;
1978 let mut steps: Vec<Value> = serde_json::from_str(&job.steps).unwrap_or_default();
1979 if let Some(step) = number.checked_sub(1).and_then(|i| steps.get_mut(i)) {
1980 let status = report["status"].as_str().unwrap_or("in_progress");
1981 step["status"] = json!(status);
1982 if status == "in_progress" {
1983 step["startedAt"] = json!(at);
1984 }
1985 if status == "completed" {
1986 step["finishedAt"] = json!(at);
1987 step["conclusion"] = report["conclusion"].clone();
1988 }
1989 if let Some(name) = report["name"].as_str() {
1990 step["name"] = json!(name);
1991 }
1992 }
1993 self.db
1994 .prepare("UPDATE jobs SET steps = ?, seen_at = ? WHERE id = ?")
1995 .bind(&[serde_json::to_string(&steps)?.into(), at.as_str().into(), job.id.as_str().into()])?
1996 .run()
1997 .await?;
1998 }
1999 "log" => {
2000 let mut text = report["text"].as_str().unwrap_or_default().to_owned();
2001 if text.len() > MAX_CHUNK_BYTES {
2002 let mut cut = MAX_CHUNK_BYTES;
2003 while !text.is_char_boundary(cut) {
2004 cut -= 1;
2005 }
2006 text.truncate(cut);
2007 }
2008 #[derive(Deserialize)]
2009 struct Size {
GitHub Actions on g1t, part three: .g1t/workflows, the pages, the docs2010 n: Option<f64>,
2011 seq: Option<f64>,
GitHub Actions on g1t, part two: running workflows2012 }
2013 let size = self
2014 .db
2015 .prepare("SELECT SUM(LENGTH(text)) AS n, MAX(seq) AS seq FROM logs WHERE job_id = ?")
2016 .bind(&[job.id.as_str().into()])?
2017 .first::<Size>(None)
2018 .await?;
GitHub Actions on g1t, part three: .g1t/workflows, the pages, the docs2019 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 workflows2020 if used < MAX_LOG_BYTES {
2021 if used + text.len() >= MAX_LOG_BYTES {
2022 text.push_str("\n… The log reached its limit of 4 MB; the rest is not kept.\n");
2023 }
2024 self.db
2025 .prepare("INSERT INTO logs (job_id, seq, step, text) VALUES (?, ?, ?, ?)")
GitHub Actions on g1t, part three: .g1t/workflows, the pages, the docs2026 .bind(&[job.id.as_str().into(), // Numbers go to D1 as f64: a u64 would be a BigInt, which it refuses.
2027 (seq + 1.0).into(), (report["step"].as_u64().unwrap_or(0) as u32).into(), text.into()])?
GitHub Actions on g1t, part two: running workflows2028 .run()
2029 .await?;
2030 }
2031 self.db.prepare("UPDATE jobs SET seen_at = ? WHERE id = ?").bind(&[at.into(), job.id.as_str().into()])?.run().await?;
2032 }
2033 "annotation" => {
2034 let mut annotations: Vec<Value> = serde_json::from_str(&job.annotations).unwrap_or_default();
2035 if annotations.len() < MAX_ANNOTATIONS {
2036 annotations.push(json!({
2037 "level": report["level"].as_str().unwrap_or("notice"),
2038 "message": report["message"].as_str().unwrap_or_default().chars().take(4000).collect::<String>(),
2039 "title": report["title"],
2040 "file": report["file"],
2041 "line": report["line"],
2042 }));
2043 self.db
2044 .prepare("UPDATE jobs SET annotations = ?, seen_at = ? WHERE id = ?")
2045 .bind(&[serde_json::to_string(&annotations)?.into(), at.as_str().into(), job.id.as_str().into()])?
2046 .run()
2047 .await?;
2048 }
2049 }
2050 "done" => {
2051 let conclusion = report["conclusion"]
2052 .as_str()
2053 .filter(|c| matches!(*c, "success" | "failure" | "cancelled"))
2054 .unwrap_or("failure");
2055 let outputs = report["outputs"].as_object().cloned();
2056 Box::pin(self.finish_job(&job.id, conclusion, report["reason"].as_str(), outputs.as_ref())).await?;
2057 }
2058 other => return Ok(fail(FailureCode::Invalid, format!("There is no report called `{other}`."))),
2059 }
2060 Ok(Outcome::Ok(json!({ "ok": true })))
2061 }
2062
2063 // --- Every minute ---------------------------------------------------------------
2064
2065 pub async fn on_minute(&self, now_ms: u64) -> Result<()> {
2066 let minute = now_ms / 60_000 * 60_000;
2067 if let Err(error) = self.run_schedules(minute).await {
2068 worker::console_error!("actions: schedules failed: {error}");
2069 }
2070 // Jobs whose sandbox went quiet or ran past their time.
2071 let running = self.db.prepare("SELECT * FROM jobs WHERE status = 'in_progress'").all().await?.results::<JobRow>()?;
2072 for job in running {
2073 // Times in g1t's format compare as text.
2074 let before = |ms: u64| rfc3339(now_ms.saturating_sub(ms));
2075 let silent = job.seen_at.as_deref().is_some_and(|seen| seen < before(SILENT_MS).as_str());
2076 let limit = (u64::from(job.timeout_minutes) * 60 + 120) * 1000;
2077 let over = job.started_at.as_deref().is_some_and(|started| started < before(limit).as_str());
2078 if over {
2079 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 incidents2080 // A self-hosted runner is told to stop on its next poll.
2081 if job.runner_id.is_none() {
2082 let _: Result<Value> = g1t_kit::call(&self.runner, "stop_actions_job", &json!({ "job": job.id })).await;
2083 }
GitHub Actions on g1t, part two: running workflows2084 self.finish_job(&job.id, "failure", Some(&reason), None).await?;
2085 } else if silent {
Fast pages, required checks on the branch, self-hosted runners, honest incidents2086 let reason = match &job.runner_name {
2087 Some(name) => format!("The self-hosted runner {name} stopped answering."),
2088 None => "The runner stopped answering.".to_owned(),
2089 };
2090 self.finish_job(&job.id, "failure", Some(&reason), None).await?;
GitHub Actions on g1t, part two: running workflows2091 }
2092 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents2093 if let Err(error) = self.sweep_runners(now_ms).await {
2094 worker::console_error!("actions: the runners' sweep failed: {error}");
2095 }
Merge branch 'worktree-agent-a3abfcce648e87dca'2096 // Jobs held at an environment whose wait timer has run out.
2097 if let Err(error) = self.release_gates().await {
2098 worker::console_error!("actions: environments' gates failed: {error}");
2099 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents2100 // Once an hour: the cache's expired entries, and its storage.
2101 if (now_ms / 60_000) % 60 == 7
2102 && let Err(error) = self.sweep_cache(now_ms).await
2103 {
2104 worker::console_error!("actions: the cache's sweep failed: {error}");
2105 }
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R22106 // And artifacts past their time, and the toolkit's abandoned parts.
2107 if (now_ms / 60_000) % 60 == 37 {
2108 if let Err(error) = self.sweep_artifacts(now_ms).await {
2109 worker::console_error!("actions: the artifacts' sweep failed: {error}");
2110 }
2111 if let Err(error) = self.sweep_blob_parts(now_ms).await {
2112 worker::console_error!("actions: the blob parts' sweep failed: {error}");
2113 }
2114 }
GitHub Actions on g1t, part two: running workflows2115 self.start_queued().await
2116 }
2117}
2118
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look2119
2120#[cfg(test)]
2121mod stopping {
2122 use super::stops_runs;
2123 use g1t_contracts::events::Event;
2124 use serde_json::{Value, json};
2125
2126 fn event(kind: &str, data: Value) -> Event {
2127 Event {
2128 id: "evt_1".into(),
2129 kind: kind.into(),
2130 source: "repos".into(),
2131 time: "2026-10-05T00:00:00Z".into(),
2132 repo_id: Some("rep_1".into()),
2133 actor: None,
2134 data,
2135 }
2136 }
2137
2138 #[test]
2139 fn deleting_or_archiving_stops_runs() {
2140 assert_eq!(stops_runs(&event("repo.deleted", json!({ "repoId": "rep_1" }))).as_deref(), Some("rep_1"));
2141 assert_eq!(stops_runs(&event("repo.archived", json!({ "archived": true }))).as_deref(), Some("rep_1"));
2142 assert_eq!(stops_runs(&event("repo.unarchived", json!({ "archived": false }))), None);
2143 assert_eq!(stops_runs(&event("repo.restored", json!({}))), None);
2144 assert_eq!(stops_runs(&event("git.push", json!({}))), None);
2145 }
2146}
Fast pages, required checks on the branch, self-hosted runners, honest incidents2147
2148#[cfg(test)]
2149mod status_of_needs {
2150 use std::collections::HashMap;
2151
2152 use super::ancestor_failed;
2153
2154 /// check -> plan -> (migrate) -> core -> edge, as deploy.yml has them,
2155 /// and a job that needs only the last.
2156 fn graph() -> HashMap<&'static str, Vec<&'static str>> {
2157 HashMap::from([
2158 ("check", vec![]),
2159 ("plan", vec!["check"]),
2160 ("migrate", vec!["plan"]),
2161 ("core", vec!["plan", "migrate"]),
2162 ("edge", vec!["plan", "migrate", "core"]),
2163 ("notify", vec!["edge"]),
2164 ])
2165 }
2166
2167 #[test]
2168 fn a_failure_is_seen_however_far_back() {
2169 let needs = graph();
2170 let failed = |which: &'static str| move |key: &str| key == which;
2171 // check failed; plan, and everything after, was skipped for it.
2172 assert!(ancestor_failed(&needs, "notify", failed("check")));
2173 assert!(ancestor_failed(&needs, "core", failed("check")));
2174 assert!(ancestor_failed(&needs, "edge", failed("core")));
2175 // Nothing before a job failed: a skipped migrate is not a failure.
2176 assert!(!ancestor_failed(&needs, "edge", |_| false));
2177 assert!(!ancestor_failed(&needs, "core", failed("edge")));
2178 assert!(!ancestor_failed(&needs, "check", failed("check")));
2179 }
2180
2181 #[test]
2182 fn cycles_and_unknown_keys_end() {
2183 let needs = HashMap::from([("a", vec!["b"]), ("b", vec!["a"])]);
2184 assert!(!ancestor_failed(&needs, "a", |_| false));
2185 assert!(!ancestor_failed(&needs, "missing", |_| true));
2186 }
2187}
Merge branch 'main' into worktree-agent-a69aeabc4b0deeb972188
2189#[cfg(test)]
2190mod deployments {
2191 use serde_json::{Map, Value, json};
2192
2193 use super::{JobEnvironment, deployment_outcome, environment_of};
2194
2195 #[test]
2196 fn a_jobs_environment_is_read_for_deployments() {
2197 let contexts: Map<String, Value> = serde_json::from_value(json!({
2198 "github": { "ref_name": "main", "repository": "acme/web" },
2199 "inputs": { "target": "staging" },
2200 "matrix": {},
2201 }))
2202 .unwrap();
2203 let read = |raw: Value| environment_of(&raw, &contexts);
2204 assert_eq!(read(json!({})), None);
2205 assert_eq!(
2206 read(json!({ "environment": "production" })),
2207 Some(JobEnvironment { name: "production".into(), url: None, deploys: true })
2208 );
2209 assert_eq!(
2210 read(json!({ "environment": { "name": "production", "url": "https://g1t.sh" } })),
2211 Some(JobEnvironment { name: "production".into(), url: Some("https://g1t.sh".into()), deploys: true })
2212 );
2213 // Expressions are filled in from the run.
2214 assert_eq!(
2215 read(json!({ "environment": { "name": "${{ inputs.target }}", "url": "https://${{ github.ref_name }}.example.com" } })),
2216 Some(JobEnvironment { name: "staging".into(), url: Some("https://main.example.com".into()), deploys: true })
2217 );
2218 // Secrets only: no deployment.
2219 assert!(!read(json!({ "environment": { "name": "production", "deployment": false } })).unwrap().deploys);
2220 // Only http(s) addresses.
2221 assert_eq!(read(json!({ "environment": { "name": "production", "url": "javascript:alert(1)" } })).unwrap().url, None);
2222 }
2223
2224 #[test]
2225 fn a_runs_outcome_for_an_environment() {
2226 let of = |list: &[&str]| deployment_outcome(&list.iter().map(|c| Some((*c).to_owned())).collect::<Vec<_>>());
2227 assert_eq!(of(&["success", "skipped"]), Some("success"));
2228 assert_eq!(of(&["success", "failure"]), Some("failure"));
2229 assert_eq!(of(&["success", "cancelled"]), Some("error"));
2230 assert_eq!(of(&["skipped"]), None);
2231 assert_eq!(deployment_outcome(&[None]), None);
2232 }
2233}

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