Skip to content
1,057 linesCodeBlameRaw
1//! What starts a run: an event on the bus, a schedule, or someone running a
2//! workflow by hand. Each finds the workflows that want it, at the commit
3//! the event is about, and checks their filters.
4
5use g1t_actions::events::{RunInfo, github_events};
6use g1t_actions::workflow::{self, Trigger, Workflow};
7use g1t_contracts::access::{self, Capability};
8use g1t_contracts::actions::{DispatchArgs, RepositoryDispatchArgs, WorkflowRun};
9use g1t_contracts::events::{Event, caused_by_job};
10use g1t_contracts::identity::{AGENT_ID, AGENT_NAME, UsernamesArgs};
11use g1t_contracts::repos::{Commit, CompareArgs, Comparison, LogArgs, Repo, RepoPath};
12use g1t_contracts::work::{IssueDetail, PullDetail, ViewArgs};
13use g1t_contracts::{FailureCode, Outcome, User, new_id};
14use g1t_kit::now_ms;
15use serde_json::{Map, Value, json};
16use worker::Result;
17
18use crate::plan::NewRun;
19use crate::sync::{Read, WorkflowRow};
20use crate::{API, Actions, SITE, check, fail, payload};
21
22/// What an event is about, worked out once for every workflow it starts.
23struct Subject {
24 /// Where the workflow files are read, and at which commit.
25 source: RepoPath,
26 source_ref: Option<String>,
27 git_ref: String,
28 sha: String,
29 head_ref: Option<String>,
30 base_ref: Option<String>,
31 pull: Option<u32>,
32 /// The branch or tag for `branches`/`tags` filters; for pull requests,
33 /// the branch they merge into.
34 filter_ref: String,
35 /// The files it changes, for `paths` filters; `None` until needed.
36 paths: Option<Vec<String>>,
37 /// For a push, what to compare to find the files.
38 compare: Option<(Option<String>, String)>,
39 payload: Value,
40 title: String,
41 trusted: bool,
42 /// Why its runs wait for approval first: a pull request from outside,
43 /// by the repository's approval policy (protection.rs).
44 approval: Option<String>,
45}
46
47/// Whether whoever a pull request is for is trusted without asking
48/// identity: g1t's agent in work nobody asked it for, or someone whose
49/// role here is known to allow pushing.
50fn trusted_outright(owner: &User, repo: &Repo) -> bool {
51 owner.id == AGENT_ID || access::can(Some(owner), repo, Capability::Push)
52}
53
54impl Actions {
55 async fn username(&self, id: Option<&str>) -> Result<Option<String>> {
56 let Some(id) = id else { return Ok(None) };
57 if id == AGENT_ID {
58 return Ok(Some(AGENT_NAME.to_owned()));
59 }
60 let names: std::collections::HashMap<String, String> =
61 g1t_kit::call(&self.identity, "usernames", &UsernamesArgs { ids: vec![id.to_owned()] }).await?;
62 Ok(names.get(id).cloned())
63 }
64
65 /// Whether whoever a pull request is for (Pull::owner: whoever asked
66 /// g1t for it, or its author) could push to the repository, so its
67 /// runs get the secrets and a token. Anyone else's, a reader's included
68 /// (who may open one on a private repository too), runs without them.
69 /// A change g1t made for someone is trusted as they are.
70 async fn insider(&self, owner: &User, repo: &Repo, ws: &User) -> Result<bool> {
71 if trusted_outright(owner, repo) {
72 return Ok(true);
73 }
74 // Stored authors carry no memberships or grants: ask identity, as
75 // the workspace (which may see anyone's permission).
76 let permission: Outcome<access::PermissionInfo> = g1t_kit::call(
77 &self.identity,
78 "collaborator_permission",
79 &access::CollaboratorPermissionArgs {
80 viewer: Some(ws.clone()),
81 path: RepoPath { namespace: repo.namespace.clone(), name: repo.name.clone() },
82 username: owner.username.clone(),
83 },
84 )
85 .await?;
86 Ok(permission
87 .into_result()
88 .ok()
89 .and_then(|info| info.role)
90 .is_some_and(|role| access::allows(role, Capability::Push)))
91 }
92
93 async fn commits(&self, repo: &Repo, actor: &User, after: &str, before: Option<&str>) -> Result<Vec<Commit>> {
94 let log: Outcome<Vec<Commit>> = g1t_kit::call(
95 &self.repos,
96 "log",
97 &LogArgs {
98 path: RepoPath {
99 namespace: repo.namespace.clone(),
100 name: repo.name.clone(),
101 },
102 viewer: Some(actor.clone()),
103 git_ref: Some(after.to_owned()),
104 limit: 20,
105 },
106 )
107 .await?;
108 let mut commits: Vec<Commit> = log.into_result().unwrap_or_default();
109 if let Some(before) = before
110 && let Some(at) = commits.iter().position(|commit| commit.hash == before)
111 {
112 commits.truncate(at);
113 }
114 // GitHub lists them oldest first, with the head commit last.
115 commits.reverse();
116 Ok(commits)
117 }
118
119 async fn changed_paths(&self, repo: &Repo, actor: &User, base: Option<String>, head: String) -> Result<Vec<String>> {
120 let compared: Outcome<Comparison> = g1t_kit::call(
121 &self.repos,
122 "compare",
123 &CompareArgs {
124 repo_id: repo.id.clone(),
125 viewer: Some(actor.clone()),
126 base,
127 head: Some(head),
128 base_branch: None,
129 },
130 )
131 .await?;
132 Ok(compared.into_result().map(|c| c.files.into_iter().map(|f| f.path).collect()).unwrap_or_default())
133 }
134
135 async fn default_head(&self, repo: &Repo) -> Result<Option<String>> {
136 g1t_kit::call(
137 &self.repos,
138 "head",
139 &g1t_contracts::repos::HeadArgs {
140 repo_id: repo.id.clone(),
141 branch: repo.default_branch.clone(),
142 },
143 )
144 .await
145 }
146
147 fn repo_path(repo: &Repo) -> RepoPath {
148 RepoPath {
149 namespace: repo.namespace.clone(),
150 name: repo.name.clone(),
151 }
152 }
153
154 /// The subject of an event of `kind`, as GitHub's `event_name`.
155 async fn subject(&self, event: &Event, event_name: &str, action: Option<&str>, repo: &Repo, ws: &User, sender: &str) -> Result<Option<Subject>> {
156 let path = Self::repo_path(repo);
157 let data = &event.data;
158 let on_default = |sha: String, payload: Value, title: String, pull: Option<u32>| Subject {
159 source: path.clone(),
160 source_ref: None,
161 git_ref: format!("refs/heads/{}", repo.default_branch),
162 sha,
163 head_ref: None,
164 base_ref: None,
165 pull,
166 filter_ref: format!("refs/heads/{}", repo.default_branch),
167 paths: None,
168 compare: None,
169 payload,
170 title,
171 trusted: true,
172 approval: None,
173 };
174 let view = |number: u32| ViewArgs {
175 repo: path.clone(),
176 number,
177 viewer: Some(ws.clone()),
178 after_seq: 0,
179 };
180 Ok(match event_name {
181 "push" => {
182 let (Some(git_ref), Some(after)) = (data["ref"].as_str(), data["after"].as_str()) else {
183 return Ok(None);
184 };
185 // The merge queue's states run merge_group workflows, not push ones.
186 if git_ref.starts_with("refs/heads/g1t-queue/") {
187 return Ok(None);
188 }
189 let before = data["before"].as_str();
190 let commits = self.commits(repo, ws, after, before).await?;
191 let title = commits.last().map(|c| c.message.lines().next().unwrap_or_default().to_owned()).unwrap_or_default();
192 let mut payload = payload::push(repo, git_ref, before, after, &commits, sender);
193 if let Some(head) = commits.last() {
194 payload["head_commit"] = payload::commit(repo, head);
195 }
196 Some(Subject {
197 source: path.clone(),
198 source_ref: Some(after.to_owned()),
199 git_ref: git_ref.to_owned(),
200 sha: after.to_owned(),
201 head_ref: None,
202 base_ref: None,
203 pull: None,
204 filter_ref: git_ref.to_owned(),
205 paths: None,
206 compare: Some((before.map(str::to_owned), after.to_owned())),
207 payload,
208 title,
209 trusted: true,
210 approval: None,
211 })
212 }
213 // A branch or tag was made: its own commit, as on GitHub.
214 "create" => {
215 let (Some(git_ref), Some(after)) = (data["ref"].as_str(), data["after"].as_str()) else {
216 return Ok(None);
217 };
218 if git_ref.starts_with("refs/heads/g1t-queue/") || !data["before"].is_null() {
219 return Ok(None);
220 }
221 let (ref_type, name) = match (git_ref.strip_prefix("refs/heads/"), git_ref.strip_prefix("refs/tags/")) {
222 (Some(branch), _) => ("branch", branch),
223 (_, Some(tag)) => ("tag", tag),
224 _ => return Ok(None),
225 };
226 let payload = json!({
227 "ref": name,
228 "ref_type": ref_type,
229 "master_branch": repo.default_branch,
230 "description": repo.description,
231 "pusher_type": "user",
232 "repository": payload::repository(repo),
233 "sender": payload::user(sender),
234 });
235 Some(Subject {
236 source: path.clone(),
237 source_ref: Some(after.to_owned()),
238 git_ref: git_ref.to_owned(),
239 sha: after.to_owned(),
240 head_ref: None,
241 base_ref: None,
242 pull: None,
243 filter_ref: git_ref.to_owned(),
244 paths: None,
245 compare: None,
246 payload,
247 title: format!("Created {ref_type} {name}"),
248 trusted: true,
249 approval: None,
250 })
251 }
252 "pull_request" | "pull_request_target" | "pull_request_review" => {
253 let Some(number) = data["number"].as_u64().map(|n| n as u32) else { return Ok(None) };
254 let detail: Outcome<PullDetail> = g1t_kit::call(&self.work, "get_pull", &view(number)).await?;
255 let Outcome::Ok(detail) = detail else { return Ok(None) };
256 let pull = &detail.pull;
257 let base_ref = pull.base_branch(&repo.default_branch).to_owned();
258 let mut payload = json!({
259 "action": action,
260 "number": pull.number,
261 "pull_request": payload::pull(repo, pull),
262 "repository": payload::repository(repo),
263 "sender": payload::user(sender),
264 });
265 payload::changed(&mut payload, data);
266 if event_name == "pull_request_review" {
267 let review = detail.comments.iter().rev().find(|c| c.verdict.is_some());
268 payload["review"] = json!({
269 "state": review.and_then(|r| r.verdict).map(|v| format!("{v:?}").to_lowercase()),
270 "body": review.map(|r| r.body.clone()),
271 "user": review.map(|r| payload::user(&r.author.username)),
272 });
273 }
274 let trusted = self.insider(pull.owner(), repo, ws).await?;
275 // A pull request from outside may wait for approval before
276 // its head's code runs. `pull_request_target` runs the
277 // base's code, and a merged one's run the commit it landed
278 // as, so neither waits.
279 let approval = if event_name != "pull_request_target" && action != Some("closed") {
280 self.approval_needed(repo, pull.owner(), ws).await?
281 } else {
282 None
283 };
284 let head_ref = payload::head_ref(pull);
285 if event_name == "pull_request_target" {
286 // In the base's context: its workflows, its head.
287 let Some(sha) = self.default_head(repo).await? else { return Ok(None) };
288 let mut subject = on_default(sha, payload, pull.title.clone(), Some(pull.number));
289 subject.head_ref = Some(head_ref);
290 subject.base_ref = Some(base_ref.clone());
291 subject.filter_ref = format!("refs/heads/{base_ref}");
292 subject.paths = Some(pull.files.iter().map(|f| f.path.clone()).collect());
293 return Ok(Some(subject));
294 }
295 // A merged pull request's run is on the commit it landed as,
296 // in the repository; otherwise on its head, where that is.
297 let landed = match (action, data["commit"].as_str()) {
298 (Some("closed"), Some(commit)) => Some(commit.to_owned()),
299 _ => None,
300 };
301 let sha = match (&landed, data["commit"].as_str(), &pull.head_commit) {
302 (Some(commit), _, _) => commit.clone(),
303 (None, Some(commit), _) => commit.to_owned(),
304 (None, None, Some(head)) => head.clone(),
305 (None, None, None) => return Ok(None),
306 };
307 let source = match landed {
308 Some(_) => path.clone(),
309 None => pull.fork.clone().unwrap_or_else(|| path.clone()),
310 };
311 Some(Subject {
312 source,
313 source_ref: Some(sha.clone()),
314 git_ref: format!("refs/pull/{}/merge", pull.number),
315 sha,
316 head_ref: Some(head_ref),
317 base_ref: Some(base_ref.clone()),
318 pull: Some(pull.number),
319 // `branches` filters on pull requests name the base.
320 filter_ref: format!("refs/heads/{base_ref}"),
321 paths: Some(pull.files.iter().map(|f| f.path.clone()).collect()),
322 compare: None,
323 payload,
324 title: pull.title.clone(),
325 trusted,
326 approval,
327 })
328 }
329 "workflow_run" => {
330 // A run of a workflow_run workflow does not start another,
331 // so two such workflows cannot set each other off.
332 if data["event"].as_str() == Some("workflow_run") {
333 return Ok(None);
334 }
335 let Some(sha) = self.default_head(repo).await? else { return Ok(None) };
336 let head_branch = data["ref"].as_str().unwrap_or_default().trim_start_matches("refs/heads/").to_owned();
337 let name = data["workflow"].as_str().unwrap_or_default();
338 let payload = json!({
339 "action": "completed",
340 "workflow_run": {
341 "id": data["runId"],
342 "name": name,
343 "path": data["path"],
344 "event": data["event"],
345 "status": "completed",
346 "conclusion": data["conclusion"],
347 "head_sha": data["sha"],
348 "head_branch": head_branch,
349 "run_number": data["number"],
350 "html_url": format!("{SITE}/{}/{}/actions/runs/{}", repo.namespace, repo.name, data["runId"].as_str().unwrap_or_default()),
351 "pull_requests": data["pull"].as_u64().map(|n| vec![json!({ "number": n })]).unwrap_or_default(),
352 },
353 "workflow": { "name": name, "path": data["path"] },
354 "repository": payload::repository(repo),
355 "sender": payload::user(sender),
356 });
357 let mut subject = on_default(sha, payload, format!("After {name}"), None);
358 // Branch filters apply to the branch the followed run was on.
359 subject.filter_ref = format!("refs/heads/{head_branch}");
360 Some(subject)
361 }
362 "issues" | "issue_comment" => {
363 let Some(number) = data["number"].as_u64().map(|n| n as u32) else { return Ok(None) };
364 let Some(sha) = self.default_head(repo).await? else { return Ok(None) };
365 let issue: Outcome<IssueDetail> = g1t_kit::call(&self.work, "get_issue", &view(number)).await?;
366 let (issue_json, comments, title, on_pull) = match issue {
367 Outcome::Ok(detail) => (payload::issue(repo, &detail.issue), detail.comments, detail.issue.title.clone(), false),
368 Outcome::Fail(_) => {
369 let pull: Outcome<PullDetail> = g1t_kit::call(&self.work, "get_pull", &view(number)).await?;
370 let Outcome::Ok(detail) = pull else { return Ok(None) };
371 (payload::pull_as_issue(repo, &detail.pull), detail.comments, detail.pull.title.clone(), true)
372 }
373 };
374 let mut payload = json!({
375 "action": action,
376 "issue": issue_json,
377 "repository": payload::repository(repo),
378 "sender": payload::user(sender),
379 });
380 payload::changed(&mut payload, data);
381 if event_name == "issue_comment" {
382 let comment_id = data["commentId"].as_str();
383 let comment = comments.iter().find(|c| Some(c.id.as_str()) == comment_id).or(comments.last());
384 match comment {
385 Some(comment) => payload["comment"] = payload::comment(repo, number, comment, on_pull),
386 None => return Ok(None),
387 }
388 }
389 Some(on_default(sha, payload, title, on_pull.then_some(number)))
390 }
391 _ => None,
392 })
393 }
394
395 pub async fn on_event(&self, event: &Event) -> Result<()> {
396 let Some(repo_id) = event.repo_id.as_deref() else { return Ok(()) };
397 let mut mapped = github_events(&event.kind);
398 // A new branch or tag is also `create`.
399 if event.kind == "git.push" && event.data["before"].is_null() {
400 mapped.push(("create", None));
401 }
402 let pushed_default = event.kind == "git.push" && event.data["defaultBranch"].as_bool() == Some(true);
403 if mapped.is_empty() && !pushed_default {
404 return Ok(());
405 }
406 let Some((repo, ws)) = self.repo_by_id(repo_id).await? else { return Ok(()) };
407 if pushed_default {
408 self.sync(&repo, &ws).await?;
409 }
410 // What a workflow job's own token did starts no workflows, as on
411 // GitHub, so a workflow cannot set itself off; only
412 // `workflow_dispatch` and `repository_dispatch` do.
413 if let Some(run) = caused_by_job(&event.data) {
414 worker::console_log!("actions: {} {} came from run {run}'s token; no workflows start for it", event.kind, event.id);
415 return Ok(());
416 }
417 let sender = self.username(event.actor.as_deref()).await?.unwrap_or_else(|| repo.namespace.clone());
418 for (event_name, action) in mapped {
419 // Issues and comments start the default branch's workflows,
420 // which the synced table lists: when none listens, nothing is
421 // read from git. Agents make many of these events.
422 if matches!(event_name, "issues" | "issue_comment") && self.listens(repo_id, event_name).await? == Some(false) {
423 continue;
424 }
425 let Some(mut subject) = self.subject(event, event_name, action, &repo, &ws, &sender).await? else {
426 continue;
427 };
428 let read = self.read_workflows(&subject.source, &ws, subject.source_ref.as_deref()).await?;
429 // A pull request's head runs each workflow once, however many
430 // events say it is there (marked ready, and pushed).
431 let key = match subject.pull {
432 Some(number) if event_name.starts_with("pull_request") && event_name != "pull_request_review" => {
433 let phase = if action == Some("closed") { "closed" } else { "open" };
434 format!("{event_name}:{number}:{}:{phase}", subject.sha)
435 }
436 _ if event_name == "create" => format!("{}:create", event.id),
437 _ => event.id.clone(),
438 };
439 self.start_matching(&repo, &ws, read, &mut subject, event_name, action, &key, event.actor.as_deref(), &sender)
440 .await?;
441 }
442 Ok(())
443 }
444
445 #[allow(clippy::too_many_arguments)]
446 async fn start_matching(
447 &self,
448 repo: &Repo,
449 ws: &User,
450 read: Read,
451 subject: &mut Subject,
452 event_name: &str,
453 action: Option<&str>,
454 event_key: &str,
455 actor_id: Option<&str>,
456 sender: &str,
457 ) -> Result<()> {
458 for file in read.files {
459 let parsed = workflow::parse(&file.source);
460 let workflow = match parsed {
461 Ok(workflow) => workflow,
462 Err(problem) => {
463 // A push shows a broken workflow as a failed run, as GitHub does.
464 if event_name == "push" && file.source.contains("on") {
465 self.record_invalid(repo, &file.path, &file.source, subject, event_key, actor_id, sender, &problem)
466 .await?;
467 }
468 continue;
469 }
470 };
471 let Some(trigger) = workflow.trigger(event_name) else { continue };
472 // workflow_run follows the workflows it names.
473 if event_name == "workflow_run" {
474 let followed = subject.payload["workflow_run"]["name"].as_str().unwrap_or_default();
475 if !trigger.workflows.iter().any(|name| name == followed) {
476 continue;
477 }
478 }
479 if !trigger.wants_type(action) || !self.passes(repo, ws, trigger, subject, event_name).await? {
480 continue;
481 }
482 if self.disabled(&repo.id, &file.path).await? {
483 continue;
484 }
485 self.create_run(NewRun {
486 repo: repo.clone(),
487 path: file.path,
488 source: file.source,
489 info: self.run_info(repo, &workflow, event_name, subject, sender, actor_id),
490 workflow,
491 action: action.map(str::to_owned),
492 pull: subject.pull,
493 title: subject.title.clone(),
494 inputs: Map::new(),
495 event_key: event_key.to_owned(),
496 actor_id: actor_id.map(str::to_owned),
497 actor: Some(sender.to_owned()),
498 trusted: subject.trusted,
499 approval: subject.approval.clone(),
500 })
501 .await?;
502 }
503 Ok(())
504 }
505
506 /// Whether the branch, tag and path filters let the event through.
507 async fn passes(&self, repo: &Repo, ws: &User, trigger: &Trigger, subject: &mut Subject, event_name: &str) -> Result<bool> {
508 let git_ref = subject.filter_ref.as_str();
509 if let Some(tag) = git_ref.strip_prefix("refs/tags/") {
510 // A tag push runs a workflow that filters tags, or filters nothing.
511 if trigger.tags.is_set() {
512 if !trigger.tags.allows(tag) {
513 return Ok(false);
514 }
515 } else if trigger.branches.is_set() {
516 return Ok(false);
517 }
518 // Paths are not checked for tags, as on GitHub.
519 return Ok(true);
520 }
521 let branch = git_ref.strip_prefix("refs/heads/").unwrap_or(git_ref);
522 if trigger.branches.is_set() {
523 if !trigger.branches.allows(branch) {
524 return Ok(false);
525 }
526 } else if event_name == "push" && trigger.tags.is_set() {
527 return Ok(false);
528 }
529 if trigger.paths.is_set() {
530 if subject.paths.is_none() {
531 let (base, head) = subject.compare.clone().unwrap_or((None, subject.sha.clone()));
532 subject.paths = Some(self.changed_paths(repo, ws, base, head).await?);
533 }
534 if !trigger.paths.allows_paths(subject.paths.as_deref().unwrap_or_default()) {
535 return Ok(false);
536 }
537 }
538 Ok(true)
539 }
540
541 async fn disabled(&self, repo_id: &str, path: &str) -> Result<bool> {
542 let row = self
543 .db
544 .prepare("SELECT * FROM workflows WHERE repo_id = ? AND path = ?")
545 .bind(&[repo_id.into(), path.into()])?
546 .first::<WorkflowRow>(None)
547 .await?;
548 Ok(row.is_some_and(|row| row.state == "disabled"))
549 }
550
551 fn run_info(&self, repo: &Repo, workflow: &Workflow, event_name: &str, subject: &Subject, sender: &str, actor_id: Option<&str>) -> RunInfo {
552 RunInfo {
553 repository: format!("{}/{}", repo.namespace, repo.name),
554 repository_id: repo.id.clone(),
555 default_branch: repo.default_branch.clone(),
556 event_name: event_name.to_owned(),
557 event: subject.payload.clone(),
558 git_ref: subject.git_ref.clone(),
559 sha: subject.sha.clone(),
560 head_ref: subject.head_ref.clone(),
561 base_ref: subject.base_ref.clone(),
562 actor: sender.to_owned(),
563 actor_id: actor_id.unwrap_or_default().to_owned(),
564 triggering_actor: sender.to_owned(),
565 run_id: String::new(),
566 run_number: 0,
567 run_attempt: 1,
568 workflow: workflow.name.clone().unwrap_or_default(),
569 workflow_path: String::new(),
570 server_url: SITE.to_owned(),
571 api_url: API.to_owned(),
572 }
573 }
574
575 #[allow(clippy::too_many_arguments)]
576 async fn record_invalid(
577 &self,
578 repo: &Repo,
579 path: &str,
580 source: &str,
581 subject: &Subject,
582 event_key: &str,
583 actor_id: Option<&str>,
584 sender: &str,
585 problem: &str,
586 ) -> Result<()> {
587 let row = self.workflow_row(repo, path, path, source).await?;
588 self.record_failed_run(&row, subject.git_ref.as_str(), &subject.sha, event_key, actor_id, sender, problem).await
589 }
590
591 /// Scheduled workflows whose cron fires this minute, on the default branch.
592 pub async fn run_schedules(&self, minute: u64) -> Result<()> {
593 let rows = self
594 .db
595 .prepare("SELECT * FROM workflows WHERE state = 'active' AND crons != '[]' AND error IS NULL")
596 .all()
597 .await?
598 .results::<WorkflowRow>()?;
599 for row in rows {
600 let crons: Vec<String> = serde_json::from_str(&row.crons).unwrap_or_default();
601 let Some(cron) = crons.iter().find(|cron| g1t_actions::cron::Schedule::parse(cron).is_ok_and(|s| s.fires_at(minute))) else {
602 continue;
603 };
604 let Ok(workflow) = workflow::parse(&row.source) else { continue };
605 // Schedules wait while a repository is archived; a deleted one is not found.
606 let Some((repo, _ws)) = self.repo_by_id(&row.repo_id).await?.filter(|(repo, _)| !repo.archived()) else { continue };
607 let Some(sha) = self.default_head(&repo).await? else { continue };
608 let payload = json!({ "schedule": cron, "repository": payload::repository(&repo), "workflow": row.path });
609 let mut subject = Subject {
610 source: Self::repo_path(&repo),
611 source_ref: None,
612 git_ref: format!("refs/heads/{}", repo.default_branch),
613 sha,
614 head_ref: None,
615 base_ref: None,
616 pull: None,
617 filter_ref: String::new(),
618 paths: None,
619 compare: None,
620 payload,
621 title: format!("Scheduled: {cron}"),
622 trusted: true,
623 approval: None,
624 };
625 subject.filter_ref = subject.git_ref.clone();
626 let info = self.run_info(&repo, &workflow, "schedule", &subject, &repo.namespace, None);
627 self.create_run(NewRun {
628 repo: repo.clone(),
629 path: row.path.clone(),
630 source: row.source.clone(),
631 workflow,
632 info,
633 action: None,
634 pull: None,
635 title: subject.title.clone(),
636 inputs: Map::new(),
637 event_key: format!("schedule:{minute}"),
638 actor_id: None,
639 actor: None,
640 trusted: true,
641 approval: None,
642 })
643 .await?;
644 }
645 Ok(())
646 }
647
648 /// `dispatch`: someone with the Write role runs a workflow that has
649 /// `workflow_dispatch`.
650 pub async fn dispatch(&self, a: DispatchArgs) -> Result<Outcome<WorkflowRun>> {
651 let repo = check!(self.may(&a.actor, &a.repo, Capability::Run).await?);
652 if repo.archived() {
653 return Ok(fail(FailureCode::Forbidden, g1t_contracts::repos::archived_message(&repo.namespace, &repo.name)));
654 }
655 let Some(ws) = self.workspace_actor(&repo.namespace).await? else {
656 return Ok(fail(FailureCode::NotFound, "There is no such workspace."));
657 };
658 let git_ref = a.git_ref.clone().unwrap_or_else(|| repo.default_branch.clone());
659 let full_ref = if git_ref.starts_with("refs/") {
660 git_ref.clone()
661 } else {
662 // A branch if there is one by that name, otherwise a tag.
663 let branches: Outcome<Vec<g1t_contracts::repos::Branch>> = g1t_kit::call(
664 &self.repos,
665 "branches",
666 &g1t_contracts::repos::BranchesArgs {
667 path: Self::repo_path(&repo),
668 viewer: Some(ws.clone()),
669 },
670 )
671 .await?;
672 let is_branch = branches.into_result().unwrap_or_default().iter().any(|branch| branch.name == git_ref);
673 format!("refs/{}/{git_ref}", if is_branch { "heads" } else { "tags" })
674 };
675 let short = full_ref.trim_start_matches("refs/heads/").trim_start_matches("refs/tags/").to_owned();
676 let read = self.read_workflows(&Self::repo_path(&repo), &ws, Some(&short)).await?;
677 let Some(sha) = read.head.clone() else {
678 return Ok(fail(FailureCode::NotFound, format!("There is no branch or tag called {short}.")));
679 };
680 // A workflow is named by its file (`build.yml`), its path, or its id
681 // (`wfl_…`), which stands for the path it was read from.
682 let by_id = if a.workflow.starts_with("wfl_") {
683 self.db
684 .prepare("SELECT * FROM workflows WHERE repo_id = ? AND id = ?")
685 .bind(&[repo.id.as_str().into(), a.workflow.as_str().into()])?
686 .first::<WorkflowRow>(None)
687 .await?
688 .map(|row| row.path)
689 } else {
690 None
691 };
692 let named = by_id.as_deref().unwrap_or(&a.workflow);
693 let wanted = named.trim_start_matches(".g1t/workflows/");
694 let Some(file) = read.files.iter().find(|file| {
695 file.path.rsplit('/').next() == Some(wanted) || file.path == named
696 }) else {
697 return Ok(fail(FailureCode::NotFound, format!("There is no workflow {wanted} on {short}.")));
698 };
699 let workflow = match workflow::parse(&file.source) {
700 Ok(workflow) => workflow,
701 Err(problem) => return Ok(fail(FailureCode::Invalid, format!("The workflow does not read: {problem}"))),
702 };
703 let Some(trigger) = workflow.trigger("workflow_dispatch") else {
704 return Ok(fail(FailureCode::Invalid, "That workflow cannot be run by hand: it has no `workflow_dispatch` trigger."));
705 };
706 let inputs = check!(dispatch_inputs(trigger, &a.inputs));
707 let payload = json!({
708 "inputs": inputs,
709 "ref": full_ref,
710 "repository": payload::repository(&repo),
711 "sender": payload::user(&a.actor.username),
712 "workflow": file.path,
713 });
714 let subject = Subject {
715 source: Self::repo_path(&repo),
716 source_ref: Some(sha.clone()),
717 git_ref: full_ref.clone(),
718 sha,
719 head_ref: None,
720 base_ref: None,
721 pull: None,
722 filter_ref: full_ref,
723 paths: None,
724 compare: None,
725 payload,
726 title: format!("{} run by {}", workflow.display_name(&file.path), a.actor.username),
727 trusted: true,
728 approval: None,
729 };
730 let info = self.run_info(&repo, &workflow, "workflow_dispatch", &subject, &a.actor.username, Some(&a.actor.id));
731 let created = self
732 .create_run(NewRun {
733 repo: repo.clone(),
734 path: file.path.clone(),
735 source: file.source.clone(),
736 workflow,
737 info,
738 action: None,
739 pull: None,
740 title: subject.title.clone(),
741 inputs,
742 event_key: format!("dispatch:{}", new_id("dsp", now_ms())),
743 actor_id: Some(a.actor.id.clone()),
744 actor: Some(a.actor.username.clone()),
745 trusted: true,
746 approval: None,
747 })
748 .await?;
749 match created {
750 Some(id) => self.run_summary(&id).await,
751 None => Ok(fail(FailureCode::Conflict, "It did not start.")),
752 }
753 }
754
755 /// `repository_dispatch`: an outside event, by name, starts the default
756 /// branch's workflows that run `on: repository_dispatch` with that type
757 /// (or with no `types`). A workflow job's token may send one: this,
758 /// with `workflow_dispatch`, is how a workflow starts another.
759 pub async fn repository_dispatch(&self, a: RepositoryDispatchArgs) -> Result<Outcome<u32>> {
760 let repo = check!(self.may(&a.actor, &a.repo, Capability::Push).await?);
761 if repo.archived() {
762 return Ok(fail(FailureCode::Forbidden, g1t_contracts::repos::archived_message(&repo.namespace, &repo.name)));
763 }
764 let event_type = a.event_type.trim().to_owned();
765 if event_type.is_empty() || event_type.chars().count() > 100 {
766 return Ok(fail(FailureCode::Invalid, "event_type is 1 to 100 characters."));
767 }
768 let client_payload = match a.client_payload {
769 Value::Null => json!({}),
770 Value::Object(map) if map.len() <= 10 => Value::Object(map),
771 Value::Object(_) => return Ok(fail(FailureCode::Invalid, "client_payload has at most 10 top-level properties.")),
772 _ => return Ok(fail(FailureCode::Invalid, "client_payload is a JSON object.")),
773 };
774 if serde_json::to_string(&client_payload).map_or(0, |text| text.len()) > 64 * 1024 {
775 return Ok(fail(FailureCode::Invalid, "client_payload is at most 64 KB."));
776 }
777 let Some(ws) = self.workspace_actor(&repo.namespace).await? else {
778 return Ok(fail(FailureCode::NotFound, "There is no such workspace."));
779 };
780 let read = self.read_workflows(&Self::repo_path(&repo), &ws, Some(&repo.default_branch)).await?;
781 let Some(sha) = read.head.clone() else {
782 return Ok(fail(FailureCode::NotFound, "The repository has no default branch to run on yet."));
783 };
784 let git_ref = format!("refs/heads/{}", repo.default_branch);
785 let payload = json!({
786 "action": event_type,
787 "branch": repo.default_branch,
788 "client_payload": client_payload,
789 "repository": payload::repository(&repo),
790 "sender": payload::user(&a.actor.username),
791 });
792 let key = format!("repository_dispatch:{}", new_id("dsp", now_ms()));
793 let mut started = 0u32;
794 for file in read.files {
795 let Ok(workflow) = workflow::parse(&file.source) else { continue };
796 let Some(trigger) = workflow.trigger("repository_dispatch") else { continue };
797 if !trigger.wants_type(Some(&event_type)) || self.disabled(&repo.id, &file.path).await? {
798 continue;
799 }
800 let subject = Subject {
801 source: Self::repo_path(&repo),
802 source_ref: Some(sha.clone()),
803 git_ref: git_ref.clone(),
804 sha: sha.clone(),
805 head_ref: None,
806 base_ref: None,
807 pull: None,
808 filter_ref: git_ref.clone(),
809 paths: None,
810 compare: None,
811 payload: payload.clone(),
812 title: event_type.clone(),
813 trusted: true,
814 approval: None,
815 };
816 let info = self.run_info(&repo, &workflow, "repository_dispatch", &subject, &a.actor.username, Some(&a.actor.id));
817 let created = self
818 .create_run(NewRun {
819 repo: repo.clone(),
820 path: file.path.clone(),
821 source: file.source.clone(),
822 workflow,
823 info,
824 action: Some(event_type.clone()),
825 pull: None,
826 title: subject.title.clone(),
827 inputs: Map::new(),
828 event_key: key.clone(),
829 actor_id: Some(a.actor.id.clone()),
830 actor: Some(a.actor.username.clone()),
831 trusted: true,
832 approval: None,
833 })
834 .await?;
835 if created.is_some() {
836 started += 1;
837 }
838 }
839 Ok(Outcome::Ok(started))
840 }
841}
842
843#[derive(serde::Deserialize)]
844#[serde(rename_all = "camelCase")]
845pub struct MergeGroupArgs {
846 pub repo_id: String,
847 pub entry: String,
848 pub sha: String,
849 pub head_ref: String,
850 #[serde(default)]
851 pub base_sha: Option<String>,
852 pub number: u32,
853 #[serde(default)]
854 pub ahead: Vec<u32>,
855}
856
857impl Actions {
858 /// `merge_group`: the merge queue built a state and its checks passed.
859 /// Starts the workflows that run `on: merge_group` on it, as GitHub's
860 /// queue does, and says how many started; the queue waits for their
861 /// statuses on that commit.
862 pub async fn merge_group(&self, a: MergeGroupArgs) -> Result<Outcome<Value>> {
863 let Some((repo, ws)) = self.repo_by_id(&a.repo_id).await? else {
864 return Ok(Outcome::Ok(json!({ "runs": 0 })));
865 };
866 let read = self.read_workflows(&Self::repo_path(&repo), &ws, Some(&a.sha)).await?;
867 let head_commit = self.commits(&repo, &ws, &a.sha, None).await?.pop();
868 let payload = json!({
869 "action": "checks_requested",
870 "merge_group": {
871 "head_sha": a.sha,
872 "head_ref": a.head_ref,
873 "base_sha": a.base_sha,
874 "base_ref": format!("refs/heads/{}", repo.default_branch),
875 "head_commit": head_commit.as_ref().map(|c| payload::commit(&repo, c)),
876 },
877 "repository": payload::repository(&repo),
878 "sender": payload::user(&repo.namespace),
879 });
880 let mut started = 0u32;
881 for file in read.files {
882 let Ok(workflow) = workflow::parse(&file.source) else { continue };
883 let Some(trigger) = workflow.trigger("merge_group") else { continue };
884 if !trigger.wants_type(Some("checks_requested")) || self.disabled(&repo.id, &file.path).await? {
885 continue;
886 }
887 // Branch filters on merge_group name the branch it merges into.
888 if trigger.branches.is_set() && !trigger.branches.allows(&repo.default_branch) {
889 continue;
890 }
891 let ahead = if a.ahead.is_empty() {
892 String::new()
893 } else {
894 format!(" after {}", a.ahead.iter().map(|n| format!("#{n}")).collect::<Vec<_>>().join(", "))
895 };
896 let subject = Subject {
897 source: Self::repo_path(&repo),
898 source_ref: Some(a.sha.clone()),
899 git_ref: a.head_ref.clone(),
900 sha: a.sha.clone(),
901 head_ref: None,
902 base_ref: Some(repo.default_branch.clone()),
903 pull: Some(a.number),
904 filter_ref: format!("refs/heads/{}", repo.default_branch),
905 paths: None,
906 compare: None,
907 payload: payload.clone(),
908 title: format!("Merge queue: #{}{ahead}", a.number),
909 trusted: true,
910 approval: None,
911 };
912 let info = self.run_info(&repo, &workflow, "merge_group", &subject, &repo.namespace, None);
913 let created = self
914 .create_run(NewRun {
915 repo: repo.clone(),
916 path: file.path.clone(),
917 source: file.source.clone(),
918 workflow,
919 info,
920 action: Some("checks_requested".to_owned()),
921 pull: Some(a.number),
922 title: subject.title.clone(),
923 inputs: Map::new(),
924 event_key: format!("merge_group:{}:{}", a.entry, a.sha),
925 actor_id: None,
926 actor: None,
927 trusted: true,
928 approval: None,
929 })
930 .await?;
931 if created.is_some() {
932 started += 1;
933 }
934 }
935 Ok(Outcome::Ok(json!({ "runs": started })))
936 }
937}
938
939/// The inputs of a manual run: what was given, checked against the
940/// workflow's declared inputs, with their defaults filled in.
941fn dispatch_inputs(trigger: &Trigger, given: &Map<String, Value>) -> Outcome<Map<String, Value>> {
942 let mut inputs = Map::new();
943 for (name, spec) in &trigger.inputs {
944 let kind = spec.get("type").and_then(Value::as_str).unwrap_or("string");
945 let value = given.get(name).cloned().or_else(|| spec.get("default").cloned());
946 let required = spec.get("required").and_then(Value::as_bool).unwrap_or(false);
947 let value = match value {
948 Some(Value::Null) | None if required => return fail(FailureCode::Invalid, format!("The input `{name}` is required.")),
949 Some(Value::Null) | None => match kind {
950 "boolean" => Value::Bool(false),
951 _ => Value::String(String::new()),
952 },
953 Some(value) => match kind {
954 "boolean" => Value::Bool(match &value {
955 Value::Bool(flag) => *flag,
956 Value::String(text) => text == "true",
957 _ => false,
958 }),
959 "number" => match &value {
960 Value::Number(_) => value,
961 Value::String(text) => match text.parse::<f64>().ok().and_then(serde_json::Number::from_f64) {
962 Some(number) => Value::Number(number),
963 None => return fail(FailureCode::Invalid, format!("The input `{name}` is a number.")),
964 },
965 _ => return fail(FailureCode::Invalid, format!("The input `{name}` is a number.")),
966 },
967 "choice" => {
968 let text = g1t_actions::expr::to_text(&value);
969 let options: Vec<String> =
970 spec.get("options").and_then(Value::as_array).map(|o| o.iter().map(g1t_actions::expr::to_text).collect()).unwrap_or_default();
971 if !options.is_empty() && !options.contains(&text) {
972 return fail(FailureCode::Invalid, format!("The input `{name}` is one of {}.", options.join(", ")));
973 }
974 Value::String(text)
975 }
976 _ => Value::String(g1t_actions::expr::to_text(&value)),
977 },
978 };
979 inputs.insert(name.clone(), value);
980 }
981 Outcome::Ok(inputs)
982}
983
984#[cfg(test)]
985mod tests {
986 use super::*;
987 use g1t_contracts::work::{Pull, g1t_author};
988 use g1t_contracts::{Membership, PrincipalKind};
989
990 fn repo() -> Repo {
991 serde_json::from_value(json!({
992 "id": "rep_1", "namespace": "acme", "name": "web", "description": null, "isPrivate": true,
993 "ownerId": "ws_1", "defaultBranch": "main", "forkOf": null, "createdAt": ""
994 }))
995 .unwrap()
996 }
997
998 fn person(id: &str, username: &str) -> User {
999 User { id: id.into(), username: username.into(), kind: PrincipalKind::User, ..User::default() }
1000 }
1001
1002 fn made_for(asker: User) -> Pull {
1003 serde_json::from_value(json!({
1004 "id": "pr_1", "repoId": "rep_1", "number": 14, "issue": 12, "title": "Fix it", "body": null,
1005 "agent": "g1t", "runtime": "hosted", "status": "open",
1006 "fork": { "namespace": "pulls", "name": "pr_1" }, "forkRepoId": "rep_f",
1007 "branch": null, "headCommit": "abc", "mergeBase": null, "mergedBy": null, "mergedAt": null,
1008 "supersededBy": null, "checkStatus": null,
1009 "author": g1t_author(), "requestedBy": asker,
1010 "createdAt": "", "updatedAt": ""
1011 }))
1012 .unwrap()
1013 }
1014
1015 #[test]
1016 fn g1t_s_change_for_someone_is_trusted_as_they_are() {
1017 // Stored people carry no memberships, so identity is asked about
1018 // them; being g1t's change gives it nothing more.
1019 let pull = made_for(person("usr_2", "ana"));
1020 assert!(!trusted_outright(pull.owner(), &repo()));
1021 // Someone known to be able to push is trusted at once.
1022 let mut member = person("usr_1", "syntaqx");
1023 member.workspaces.push(Membership::member("acme"));
1024 let pull = made_for(member);
1025 assert!(trusted_outright(pull.owner(), &repo()));
1026 }
1027
1028 #[test]
1029 fn the_payload_names_g1t_as_its_user_and_who_asked_for_it() {
1030 let pull = made_for(person("usr_1", "syntaqx"));
1031 let event = payload::pull(&repo(), &pull);
1032 assert_eq!(event["user"]["login"], "g1t");
1033 assert_eq!(event["user"]["type"], "Bot");
1034 assert_eq!(event["requested_by"]["login"], "syntaqx");
1035 assert_eq!(event["requested_by"]["type"], "User");
1036 let as_issue = payload::pull_as_issue(&repo(), &pull);
1037 assert_eq!(as_issue["user"]["login"], "g1t");
1038 assert_eq!(as_issue["requested_by"]["login"], "syntaqx");
1039 }
1040
1041 #[test]
1042 fn a_pull_request_names_its_own_base_labels_and_milestone() {
1043 let mut pull = made_for(person("usr_1", "syntaqx"));
1044 let event = payload::pull(&repo(), &pull);
1045 assert_eq!(event["base"]["ref"], repo().default_branch);
1046 pull.base = Some("release/1.x".into());
1047 pull.labels = vec!["bug".into()];
1048 pull.milestone = Some(g1t_contracts::work::MilestoneRef { number: 2, title: "1.1".into() });
1049 let event = payload::pull(&repo(), &pull);
1050 assert_eq!(event["base"]["ref"], "release/1.x");
1051 assert_eq!(event["labels"], serde_json::json!([{ "name": "bug" }]));
1052 assert_eq!(event["milestone"]["title"], "1.1");
1053 let mut labeled = serde_json::json!({ "action": "labeled" });
1054 payload::changed(&mut labeled, &serde_json::json!({ "label": { "name": "bug", "color": "d73a4a" } }));
1055 assert_eq!(labeled["label"]["color"], "d73a4a");
1056 }
1057}