flagon-io/g1t

public

Where people and agents ship software together. The open-source git platform for the whole job: issues, agents, checks and deploys to the edge.

g1t/services/actions/src/trigger.rs

788 lines35,508 bytesCodeBlame
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::actions::{DispatchArgs, WorkflowRun};
8use g1t_contracts::events::Event;
9use g1t_contracts::identity::{AGENT_ID, AGENT_NAME, UsernamesArgs};
10use g1t_contracts::repos::{Commit, CompareArgs, Comparison, LogArgs, Repo, RepoPath};
11use g1t_contracts::work::{IssueDetail, PullDetail, ViewArgs};
12use g1t_contracts::{FailureCode, Outcome, User, new_id};
13use g1t_kit::now_ms;
14use serde_json::{Map, Value, json};
15use worker::Result;
16
17use crate::plan::NewRun;
18use crate::sync::{Read, WorkflowRow};
19use crate::{API, Actions, SITE, check, fail, payload};
20
21/// What an event is about, worked out once for every workflow it starts.
22struct Subject {
23 /// Where the workflow files are read, and at which commit.
24 source: RepoPath,
25 source_ref: Option<String>,
26 git_ref: String,
27 sha: String,
28 head_ref: Option<String>,
29 base_ref: Option<String>,
30 pull: Option<u32>,
31 /// The branch or tag for `branches`/`tags` filters; for pull requests,
32 /// the branch they merge into.
33 filter_ref: String,
34 /// The files it changes, for `paths` filters; `None` until needed.
35 paths: Option<Vec<String>>,
36 /// For a push, what to compare to find the files.
37 compare: Option<(Option<String>, String)>,
38 payload: Value,
39 title: String,
40 trusted: bool,
41}
42
43impl Actions {
44 async fn username(&self, id: Option<&str>) -> Result<Option<String>> {
45 let Some(id) = id else { return Ok(None) };
46 if id == AGENT_ID {
47 return Ok(Some(AGENT_NAME.to_owned()));
48 }
49 let names: std::collections::HashMap<String, String> =
50 g1t_kit::call(&self.identity, "usernames", &UsernamesArgs { ids: vec![id.to_owned()] }).await?;
51 Ok(names.get(id).cloned())
52 }
53
54 /// Whether a pull request's author belongs to the workspace, so its
55 /// runs get the secrets and a token. On a private repository only
56 /// members can open one at all.
57 async fn insider(&self, author: &User, repo: &Repo, ws: &User) -> Result<bool> {
58 let slug = repo.namespace.to_lowercase();
59 if repo.is_private || author.id == AGENT_ID || author.is_member(&slug) {
60 return Ok(true);
61 }
62 // Stored authors carry no memberships: ask the workspace.
63 let members: Outcome<Vec<g1t_contracts::identity::Member>> = g1t_kit::call(
64 &self.identity,
65 "list_members",
66 &g1t_contracts::identity::ListMembersArgs { slug, viewer: Some(ws.clone()) },
67 )
68 .await?;
69 Ok(members.into_result().unwrap_or_default().iter().any(|m| m.username.eq_ignore_ascii_case(&author.username)))
70 }
71
72 async fn commits(&self, repo: &Repo, actor: &User, after: &str, before: Option<&str>) -> Result<Vec<Commit>> {
73 let log: Outcome<Vec<Commit>> = g1t_kit::call(
74 &self.repos,
75 "log",
76 &LogArgs {
77 path: RepoPath {
78 namespace: repo.namespace.clone(),
79 name: repo.name.clone(),
80 },
81 viewer: Some(actor.clone()),
82 git_ref: Some(after.to_owned()),
83 limit: 20,
84 },
85 )
86 .await?;
87 let mut commits: Vec<Commit> = log.into_result().unwrap_or_default();
88 if let Some(before) = before
89 && let Some(at) = commits.iter().position(|commit| commit.hash == before)
90 {
91 commits.truncate(at);
92 }
93 // GitHub lists them oldest first, with the head commit last.
94 commits.reverse();
95 Ok(commits)
96 }
97
98 async fn changed_paths(&self, repo: &Repo, actor: &User, base: Option<String>, head: String) -> Result<Vec<String>> {
99 let compared: Outcome<Comparison> = g1t_kit::call(
100 &self.repos,
101 "compare",
102 &CompareArgs {
103 repo_id: repo.id.clone(),
104 viewer: Some(actor.clone()),
105 base,
106 head: Some(head),
107 },
108 )
109 .await?;
110 Ok(compared.into_result().map(|c| c.files.into_iter().map(|f| f.path).collect()).unwrap_or_default())
111 }
112
113 async fn default_head(&self, repo: &Repo) -> Result<Option<String>> {
114 g1t_kit::call(
115 &self.repos,
116 "head",
117 &g1t_contracts::repos::HeadArgs {
118 repo_id: repo.id.clone(),
119 branch: repo.default_branch.clone(),
120 },
121 )
122 .await
123 }
124
125 fn repo_path(repo: &Repo) -> RepoPath {
126 RepoPath {
127 namespace: repo.namespace.clone(),
128 name: repo.name.clone(),
129 }
130 }
131
132 /// The subject of an event of `kind`, as GitHub's `event_name`.
133 async fn subject(&self, event: &Event, event_name: &str, action: Option<&str>, repo: &Repo, ws: &User, sender: &str) -> Result<Option<Subject>> {
134 let path = Self::repo_path(repo);
135 let data = &event.data;
136 let on_default = |sha: String, payload: Value, title: String, pull: Option<u32>| Subject {
137 source: path.clone(),
138 source_ref: None,
139 git_ref: format!("refs/heads/{}", repo.default_branch),
140 sha,
141 head_ref: None,
142 base_ref: None,
143 pull,
144 filter_ref: format!("refs/heads/{}", repo.default_branch),
145 paths: None,
146 compare: None,
147 payload,
148 title,
149 trusted: true,
150 };
151 let view = |number: u32| ViewArgs {
152 repo: path.clone(),
153 number,
154 viewer: Some(ws.clone()),
155 after_seq: 0,
156 };
157 Ok(match event_name {
158 "push" => {
159 let (Some(git_ref), Some(after)) = (data["ref"].as_str(), data["after"].as_str()) else {
160 return Ok(None);
161 };
162 // The merge queue's states run merge_group workflows, not push ones.
163 if git_ref.starts_with("refs/heads/g1t-queue/") {
164 return Ok(None);
165 }
166 let before = data["before"].as_str();
167 let commits = self.commits(repo, ws, after, before).await?;
168 let title = commits.last().map(|c| c.message.lines().next().unwrap_or_default().to_owned()).unwrap_or_default();
169 let mut payload = payload::push(repo, git_ref, before, after, &commits, sender);
170 if let Some(head) = commits.last() {
171 payload["head_commit"] = payload::commit(repo, head);
172 }
173 Some(Subject {
174 source: path.clone(),
175 source_ref: Some(after.to_owned()),
176 git_ref: git_ref.to_owned(),
177 sha: after.to_owned(),
178 head_ref: None,
179 base_ref: None,
180 pull: None,
181 filter_ref: git_ref.to_owned(),
182 paths: None,
183 compare: Some((before.map(str::to_owned), after.to_owned())),
184 payload,
185 title,
186 trusted: true,
187 })
188 }
189 "pull_request" | "pull_request_target" | "pull_request_review" => {
190 let Some(number) = data["number"].as_u64().map(|n| n as u32) else { return Ok(None) };
191 let detail: Outcome<PullDetail> = g1t_kit::call(&self.work, "get_pull", &view(number)).await?;
192 let Outcome::Ok(detail) = detail else { return Ok(None) };
193 let pull = &detail.pull;
194 let labels = detail.issue.as_ref().map(|i| i.labels.clone()).unwrap_or_default();
195 let mut payload = json!({
196 "action": action,
197 "number": pull.number,
198 "pull_request": payload::pull(repo, pull, &labels),
199 "repository": payload::repository(repo),
200 "sender": payload::user(sender),
201 });
202 if event_name == "pull_request_review" {
203 let review = detail.comments.iter().rev().find(|c| c.verdict.is_some());
204 payload["review"] = json!({
205 "state": review.and_then(|r| r.verdict).map(|v| format!("{v:?}").to_lowercase()),
206 "body": review.map(|r| r.body.clone()),
207 "user": review.map(|r| payload::user(&r.author.username)),
208 });
209 }
210 let trusted = self.insider(&pull.author, repo, ws).await?;
211 let head_ref = payload::head_ref(pull);
212 if event_name == "pull_request_target" {
213 // In the base's context: its workflows, its head.
214 let Some(sha) = self.default_head(repo).await? else { return Ok(None) };
215 let mut subject = on_default(sha, payload, pull.title.clone(), Some(pull.number));
216 subject.head_ref = Some(head_ref);
217 subject.base_ref = Some(repo.default_branch.clone());
218 subject.paths = Some(pull.files.iter().map(|f| f.path.clone()).collect());
219 return Ok(Some(subject));
220 }
221 // A merged pull request's run is on the commit it landed as,
222 // in the repository; otherwise on its head, where that is.
223 let landed = match (action, data["commit"].as_str()) {
224 (Some("closed"), Some(commit)) => Some(commit.to_owned()),
225 _ => None,
226 };
227 let sha = match (&landed, data["commit"].as_str(), &pull.head_commit) {
228 (Some(commit), _, _) => commit.clone(),
229 (None, Some(commit), _) => commit.to_owned(),
230 (None, None, Some(head)) => head.clone(),
231 (None, None, None) => return Ok(None),
232 };
233 let source = match landed {
234 Some(_) => path.clone(),
235 None => pull.fork.clone().unwrap_or_else(|| path.clone()),
236 };
237 Some(Subject {
238 source,
239 source_ref: Some(sha.clone()),
240 git_ref: format!("refs/pull/{}/merge", pull.number),
241 sha,
242 head_ref: Some(head_ref),
243 base_ref: Some(repo.default_branch.clone()),
244 pull: Some(pull.number),
245 filter_ref: format!("refs/heads/{}", repo.default_branch),
246 paths: Some(pull.files.iter().map(|f| f.path.clone()).collect()),
247 compare: None,
248 payload,
249 title: pull.title.clone(),
250 trusted,
251 })
252 }
253 "workflow_run" => {
254 // A run of a workflow_run workflow does not start another,
255 // so two such workflows cannot set each other off.
256 if data["event"].as_str() == Some("workflow_run") {
257 return Ok(None);
258 }
259 let Some(sha) = self.default_head(repo).await? else { return Ok(None) };
260 let head_branch = data["ref"].as_str().unwrap_or_default().trim_start_matches("refs/heads/").to_owned();
261 let name = data["workflow"].as_str().unwrap_or_default();
262 let payload = json!({
263 "action": "completed",
264 "workflow_run": {
265 "id": data["runId"],
266 "name": name,
267 "path": data["path"],
268 "event": data["event"],
269 "status": "completed",
270 "conclusion": data["conclusion"],
271 "head_sha": data["sha"],
272 "head_branch": head_branch,
273 "run_number": data["number"],
274 "html_url": format!("{SITE}/{}/{}/actions/runs/{}", repo.namespace, repo.name, data["runId"].as_str().unwrap_or_default()),
275 "pull_requests": data["pull"].as_u64().map(|n| vec![json!({ "number": n })]).unwrap_or_default(),
276 },
277 "workflow": { "name": name, "path": data["path"] },
278 "repository": payload::repository(repo),
279 "sender": payload::user(sender),
280 });
281 let mut subject = on_default(sha, payload, format!("After {name}"), None);
282 // Branch filters apply to the branch the followed run was on.
283 subject.filter_ref = format!("refs/heads/{head_branch}");
284 Some(subject)
285 }
286 "issues" | "issue_comment" => {
287 let Some(number) = data["number"].as_u64().map(|n| n as u32) else { return Ok(None) };
288 let Some(sha) = self.default_head(repo).await? else { return Ok(None) };
289 let issue: Outcome<IssueDetail> = g1t_kit::call(&self.work, "get_issue", &view(number)).await?;
290 let (issue_json, comments, title, on_pull) = match issue {
291 Outcome::Ok(detail) => (payload::issue(repo, &detail.issue), detail.comments, detail.issue.title.clone(), false),
292 Outcome::Fail(_) => {
293 let pull: Outcome<PullDetail> = g1t_kit::call(&self.work, "get_pull", &view(number)).await?;
294 let Outcome::Ok(detail) = pull else { return Ok(None) };
295 let labels = detail.issue.as_ref().map(|i| i.labels.clone()).unwrap_or_default();
296 (payload::pull_as_issue(repo, &detail.pull, &labels), detail.comments, detail.pull.title.clone(), true)
297 }
298 };
299 let mut payload = json!({
300 "action": action,
301 "issue": issue_json,
302 "repository": payload::repository(repo),
303 "sender": payload::user(sender),
304 });
305 if event_name == "issue_comment" {
306 let comment_id = data["commentId"].as_str();
307 let comment = comments.iter().find(|c| Some(c.id.as_str()) == comment_id).or(comments.last());
308 match comment {
309 Some(comment) => payload["comment"] = payload::comment(repo, number, comment, on_pull),
310 None => return Ok(None),
311 }
312 }
313 Some(on_default(sha, payload, title, on_pull.then_some(number)))
314 }
315 _ => None,
316 })
317 }
318
319 pub async fn on_event(&self, event: &Event) -> Result<()> {
320 let Some(repo_id) = event.repo_id.as_deref() else { return Ok(()) };
321 let mapped = github_events(&event.kind);
322 let pushed_default = event.kind == "git.push" && event.data["defaultBranch"].as_bool() == Some(true);
323 if mapped.is_empty() && !pushed_default {
324 return Ok(());
325 }
326 let Some((repo, ws)) = self.repo_by_id(repo_id).await? else { return Ok(()) };
327 if pushed_default {
328 self.sync(&repo, &ws).await?;
329 }
330 let sender = self.username(event.actor.as_deref()).await?.unwrap_or_else(|| repo.namespace.clone());
331 for (event_name, action) in mapped {
332 let Some(mut subject) = self.subject(event, event_name, action, &repo, &ws, &sender).await? else {
333 continue;
334 };
335 let read = self.read_workflows(&subject.source, &ws, subject.source_ref.as_deref()).await?;
336 // A pull request's head runs each workflow once, however many
337 // events say it is there (marked ready, and pushed).
338 let key = match subject.pull {
339 Some(number) if event_name.starts_with("pull_request") && event_name != "pull_request_review" => {
340 let phase = if action == Some("closed") { "closed" } else { "open" };
341 format!("{event_name}:{number}:{}:{phase}", subject.sha)
342 }
343 _ => event.id.clone(),
344 };
345 self.start_matching(&repo, &ws, read, &mut subject, event_name, action, &key, event.actor.as_deref(), &sender)
346 .await?;
347 }
348 Ok(())
349 }
350
351 #[allow(clippy::too_many_arguments)]
352 async fn start_matching(
353 &self,
354 repo: &Repo,
355 ws: &User,
356 read: Read,
357 subject: &mut Subject,
358 event_name: &str,
359 action: Option<&str>,
360 event_key: &str,
361 actor_id: Option<&str>,
362 sender: &str,
363 ) -> Result<()> {
364 for file in read.files {
365 let parsed = workflow::parse(&file.source);
366 let workflow = match parsed {
367 Ok(workflow) => workflow,
368 Err(problem) => {
369 // A push shows a broken workflow as a failed run, as GitHub does.
370 if event_name == "push" && file.source.contains("on") {
371 self.record_invalid(repo, &file.path, &file.source, subject, event_key, actor_id, sender, &problem)
372 .await?;
373 }
374 continue;
375 }
376 };
377 let Some(trigger) = workflow.trigger(event_name) else { continue };
378 // workflow_run follows the workflows it names.
379 if event_name == "workflow_run" {
380 let followed = subject.payload["workflow_run"]["name"].as_str().unwrap_or_default();
381 if !trigger.workflows.iter().any(|name| name == followed) {
382 continue;
383 }
384 }
385 if !trigger.wants_type(action) || !self.passes(repo, ws, trigger, subject, event_name).await? {
386 continue;
387 }
388 if self.disabled(&repo.id, &file.path).await? {
389 continue;
390 }
391 self.create_run(NewRun {
392 repo: repo.clone(),
393 path: file.path,
394 source: file.source,
395 info: self.run_info(repo, &workflow, event_name, subject, sender, actor_id),
396 workflow,
397 action: action.map(str::to_owned),
398 pull: subject.pull,
399 title: subject.title.clone(),
400 inputs: Map::new(),
401 event_key: event_key.to_owned(),
402 actor_id: actor_id.map(str::to_owned),
403 actor: Some(sender.to_owned()),
404 trusted: subject.trusted,
405 })
406 .await?;
407 }
408 Ok(())
409 }
410
411 /// Whether the branch, tag and path filters let the event through.
412 async fn passes(&self, repo: &Repo, ws: &User, trigger: &Trigger, subject: &mut Subject, event_name: &str) -> Result<bool> {
413 let git_ref = subject.filter_ref.as_str();
414 if let Some(tag) = git_ref.strip_prefix("refs/tags/") {
415 // A tag push runs a workflow that filters tags, or filters nothing.
416 if trigger.tags.is_set() {
417 if !trigger.tags.allows(tag) {
418 return Ok(false);
419 }
420 } else if trigger.branches.is_set() {
421 return Ok(false);
422 }
423 // Paths are not checked for tags, as on GitHub.
424 return Ok(true);
425 }
426 let branch = git_ref.strip_prefix("refs/heads/").unwrap_or(git_ref);
427 if trigger.branches.is_set() {
428 if !trigger.branches.allows(branch) {
429 return Ok(false);
430 }
431 } else if event_name == "push" && trigger.tags.is_set() {
432 return Ok(false);
433 }
434 if trigger.paths.is_set() {
435 if subject.paths.is_none() {
436 let (base, head) = subject.compare.clone().unwrap_or((None, subject.sha.clone()));
437 subject.paths = Some(self.changed_paths(repo, ws, base, head).await?);
438 }
439 if !trigger.paths.allows_paths(subject.paths.as_deref().unwrap_or_default()) {
440 return Ok(false);
441 }
442 }
443 Ok(true)
444 }
445
446 async fn disabled(&self, repo_id: &str, path: &str) -> Result<bool> {
447 let row = self
448 .db
449 .prepare("SELECT * FROM workflows WHERE repo_id = ? AND path = ?")
450 .bind(&[repo_id.into(), path.into()])?
451 .first::<WorkflowRow>(None)
452 .await?;
453 Ok(row.is_some_and(|row| row.state == "disabled"))
454 }
455
456 fn run_info(&self, repo: &Repo, workflow: &Workflow, event_name: &str, subject: &Subject, sender: &str, actor_id: Option<&str>) -> RunInfo {
457 RunInfo {
458 repository: format!("{}/{}", repo.namespace, repo.name),
459 repository_id: repo.id.clone(),
460 default_branch: repo.default_branch.clone(),
461 event_name: event_name.to_owned(),
462 event: subject.payload.clone(),
463 git_ref: subject.git_ref.clone(),
464 sha: subject.sha.clone(),
465 head_ref: subject.head_ref.clone(),
466 base_ref: subject.base_ref.clone(),
467 actor: sender.to_owned(),
468 actor_id: actor_id.unwrap_or_default().to_owned(),
469 triggering_actor: sender.to_owned(),
470 run_id: String::new(),
471 run_number: 0,
472 run_attempt: 1,
473 workflow: workflow.name.clone().unwrap_or_default(),
474 workflow_path: String::new(),
475 server_url: SITE.to_owned(),
476 api_url: API.to_owned(),
477 }
478 }
479
480 #[allow(clippy::too_many_arguments)]
481 async fn record_invalid(
482 &self,
483 repo: &Repo,
484 path: &str,
485 source: &str,
486 subject: &Subject,
487 event_key: &str,
488 actor_id: Option<&str>,
489 sender: &str,
490 problem: &str,
491 ) -> Result<()> {
492 let row = self.workflow_row(repo, path, path, source).await?;
493 self.record_failed_run(&row, subject.git_ref.as_str(), &subject.sha, event_key, actor_id, sender, problem).await
494 }
495
496 /// Scheduled workflows whose cron fires this minute, on the default branch.
497 pub async fn run_schedules(&self, minute: u64) -> Result<()> {
498 let rows = self
499 .db
500 .prepare("SELECT * FROM workflows WHERE state = 'active' AND crons != '[]' AND error IS NULL")
501 .all()
502 .await?
503 .results::<WorkflowRow>()?;
504 for row in rows {
505 let crons: Vec<String> = serde_json::from_str(&row.crons).unwrap_or_default();
506 let Some(cron) = crons.iter().find(|cron| g1t_actions::cron::Schedule::parse(cron).is_ok_and(|s| s.fires_at(minute))) else {
507 continue;
508 };
509 let Ok(workflow) = workflow::parse(&row.source) else { continue };
510 let Some((repo, _ws)) = self.repo_by_id(&row.repo_id).await? else { continue };
511 let Some(sha) = self.default_head(&repo).await? else { continue };
512 let payload = json!({ "schedule": cron, "repository": payload::repository(&repo), "workflow": row.path });
513 let mut subject = Subject {
514 source: Self::repo_path(&repo),
515 source_ref: None,
516 git_ref: format!("refs/heads/{}", repo.default_branch),
517 sha,
518 head_ref: None,
519 base_ref: None,
520 pull: None,
521 filter_ref: String::new(),
522 paths: None,
523 compare: None,
524 payload,
525 title: format!("Scheduled: {cron}"),
526 trusted: true,
527 };
528 subject.filter_ref = subject.git_ref.clone();
529 let info = self.run_info(&repo, &workflow, "schedule", &subject, &repo.namespace, None);
530 self.create_run(NewRun {
531 repo: repo.clone(),
532 path: row.path.clone(),
533 source: row.source.clone(),
534 workflow,
535 info,
536 action: None,
537 pull: None,
538 title: subject.title.clone(),
539 inputs: Map::new(),
540 event_key: format!("schedule:{minute}"),
541 actor_id: None,
542 actor: None,
543 trusted: true,
544 })
545 .await?;
546 }
547 Ok(())
548 }
549
550 /// `dispatch`: a member runs a workflow that has `workflow_dispatch`.
551 pub async fn dispatch(&self, a: DispatchArgs) -> Result<Outcome<WorkflowRun>> {
552 if let Some(refused) = Self::member(&a.actor, &a.repo) {
553 return Ok(check_refusal(refused));
554 }
555 let Some(repo) = self.visible_repo(&a.repo, &Some(a.actor.clone())).await? else {
556 return Ok(fail(FailureCode::NotFound, "There is no such repository."));
557 };
558 let Some(ws) = self.workspace_actor(&repo.namespace).await? else {
559 return Ok(fail(FailureCode::NotFound, "There is no such workspace."));
560 };
561 let git_ref = a.git_ref.clone().unwrap_or_else(|| repo.default_branch.clone());
562 let full_ref = if git_ref.starts_with("refs/") {
563 git_ref.clone()
564 } else {
565 // A branch if there is one by that name, otherwise a tag.
566 let branches: Outcome<Vec<g1t_contracts::repos::Branch>> = g1t_kit::call(
567 &self.repos,
568 "branches",
569 &g1t_contracts::repos::BranchesArgs {
570 path: Self::repo_path(&repo),
571 viewer: Some(ws.clone()),
572 },
573 )
574 .await?;
575 let is_branch = branches.into_result().unwrap_or_default().iter().any(|branch| branch.name == git_ref);
576 format!("refs/{}/{git_ref}", if is_branch { "heads" } else { "tags" })
577 };
578 let short = full_ref.trim_start_matches("refs/heads/").trim_start_matches("refs/tags/").to_owned();
579 let read = self.read_workflows(&Self::repo_path(&repo), &ws, Some(&short)).await?;
580 let Some(sha) = read.head.clone() else {
581 return Ok(fail(FailureCode::NotFound, format!("There is no branch or tag called {short}.")));
582 };
583 let wanted = a.workflow.trim_start_matches(".g1t/workflows/");
584 let Some(file) = read.files.iter().find(|file| {
585 file.path.rsplit('/').next() == Some(wanted) || file.path == a.workflow
586 }) else {
587 return Ok(fail(FailureCode::NotFound, format!("There is no workflow {wanted} on {short}.")));
588 };
589 let workflow = match workflow::parse(&file.source) {
590 Ok(workflow) => workflow,
591 Err(problem) => return Ok(fail(FailureCode::Invalid, format!("The workflow does not read: {problem}"))),
592 };
593 let Some(trigger) = workflow.trigger("workflow_dispatch") else {
594 return Ok(fail(FailureCode::Invalid, "That workflow cannot be run by hand: it has no `workflow_dispatch` trigger."));
595 };
596 let inputs = check!(dispatch_inputs(trigger, &a.inputs));
597 let payload = json!({
598 "inputs": inputs,
599 "ref": full_ref,
600 "repository": payload::repository(&repo),
601 "sender": payload::user(&a.actor.username),
602 "workflow": file.path,
603 });
604 let subject = Subject {
605 source: Self::repo_path(&repo),
606 source_ref: Some(sha.clone()),
607 git_ref: full_ref.clone(),
608 sha,
609 head_ref: None,
610 base_ref: None,
611 pull: None,
612 filter_ref: full_ref,
613 paths: None,
614 compare: None,
615 payload,
616 title: format!("{} run by {}", workflow.display_name(&file.path), a.actor.username),
617 trusted: true,
618 };
619 let info = self.run_info(&repo, &workflow, "workflow_dispatch", &subject, &a.actor.username, Some(&a.actor.id));
620 let created = self
621 .create_run(NewRun {
622 repo: repo.clone(),
623 path: file.path.clone(),
624 source: file.source.clone(),
625 workflow,
626 info,
627 action: None,
628 pull: None,
629 title: subject.title.clone(),
630 inputs,
631 event_key: format!("dispatch:{}", new_id("dsp", now_ms())),
632 actor_id: Some(a.actor.id.clone()),
633 actor: Some(a.actor.username.clone()),
634 trusted: true,
635 })
636 .await?;
637 match created {
638 Some(id) => self.run_summary(&id).await,
639 None => Ok(fail(FailureCode::Conflict, "It did not start.")),
640 }
641 }
642}
643
644#[derive(serde::Deserialize)]
645#[serde(rename_all = "camelCase")]
646pub struct MergeGroupArgs {
647 pub repo_id: String,
648 pub entry: String,
649 pub sha: String,
650 pub head_ref: String,
651 #[serde(default)]
652 pub base_sha: Option<String>,
653 pub number: u32,
654 #[serde(default)]
655 pub ahead: Vec<u32>,
656}
657
658impl Actions {
659 /// `merge_group`: the merge queue built a state and its checks passed.
660 /// Starts the workflows that run `on: merge_group` on it, as GitHub's
661 /// queue does, and says how many started; the queue waits for their
662 /// statuses on that commit.
663 pub async fn merge_group(&self, a: MergeGroupArgs) -> Result<Outcome<Value>> {
664 let Some((repo, ws)) = self.repo_by_id(&a.repo_id).await? else {
665 return Ok(Outcome::Ok(json!({ "runs": 0 })));
666 };
667 let read = self.read_workflows(&Self::repo_path(&repo), &ws, Some(&a.sha)).await?;
668 let head_commit = self.commits(&repo, &ws, &a.sha, None).await?.pop();
669 let payload = json!({
670 "action": "checks_requested",
671 "merge_group": {
672 "head_sha": a.sha,
673 "head_ref": a.head_ref,
674 "base_sha": a.base_sha,
675 "base_ref": format!("refs/heads/{}", repo.default_branch),
676 "head_commit": head_commit.as_ref().map(|c| payload::commit(&repo, c)),
677 },
678 "repository": payload::repository(&repo),
679 "sender": payload::user(&repo.namespace),
680 });
681 let mut started = 0u32;
682 for file in read.files {
683 let Ok(workflow) = workflow::parse(&file.source) else { continue };
684 let Some(trigger) = workflow.trigger("merge_group") else { continue };
685 if !trigger.wants_type(Some("checks_requested")) || self.disabled(&repo.id, &file.path).await? {
686 continue;
687 }
688 // Branch filters on merge_group name the branch it merges into.
689 if trigger.branches.is_set() && !trigger.branches.allows(&repo.default_branch) {
690 continue;
691 }
692 let ahead = if a.ahead.is_empty() {
693 String::new()
694 } else {
695 format!(" after {}", a.ahead.iter().map(|n| format!("#{n}")).collect::<Vec<_>>().join(", "))
696 };
697 let subject = Subject {
698 source: Self::repo_path(&repo),
699 source_ref: Some(a.sha.clone()),
700 git_ref: a.head_ref.clone(),
701 sha: a.sha.clone(),
702 head_ref: None,
703 base_ref: Some(repo.default_branch.clone()),
704 pull: Some(a.number),
705 filter_ref: format!("refs/heads/{}", repo.default_branch),
706 paths: None,
707 compare: None,
708 payload: payload.clone(),
709 title: format!("Merge queue: #{}{ahead}", a.number),
710 trusted: true,
711 };
712 let info = self.run_info(&repo, &workflow, "merge_group", &subject, &repo.namespace, None);
713 let created = self
714 .create_run(NewRun {
715 repo: repo.clone(),
716 path: file.path.clone(),
717 source: file.source.clone(),
718 workflow,
719 info,
720 action: Some("checks_requested".to_owned()),
721 pull: Some(a.number),
722 title: subject.title.clone(),
723 inputs: Map::new(),
724 event_key: format!("merge_group:{}:{}", a.entry, a.sha),
725 actor_id: None,
726 actor: None,
727 trusted: true,
728 })
729 .await?;
730 if created.is_some() {
731 started += 1;
732 }
733 }
734 Ok(Outcome::Ok(json!({ "runs": started })))
735 }
736}
737
738fn check_refusal<T>(refused: Outcome<()>) -> Outcome<T> {
739 match refused {
740 Outcome::Fail(failure) => Outcome::Fail(failure),
741 Outcome::Ok(()) => fail(FailureCode::Forbidden, "Not allowed."),
742 }
743}
744
745/// The inputs of a manual run: what was given, checked against the
746/// workflow's declared inputs, with their defaults filled in.
747fn dispatch_inputs(trigger: &Trigger, given: &Map<String, Value>) -> Outcome<Map<String, Value>> {
748 let mut inputs = Map::new();
749 for (name, spec) in &trigger.inputs {
750 let kind = spec.get("type").and_then(Value::as_str).unwrap_or("string");
751 let value = given.get(name).cloned().or_else(|| spec.get("default").cloned());
752 let required = spec.get("required").and_then(Value::as_bool).unwrap_or(false);
753 let value = match value {
754 Some(Value::Null) | None if required => return fail(FailureCode::Invalid, format!("The input `{name}` is required.")),
755 Some(Value::Null) | None => match kind {
756 "boolean" => Value::Bool(false),
757 _ => Value::String(String::new()),
758 },
759 Some(value) => match kind {
760 "boolean" => Value::Bool(match &value {
761 Value::Bool(flag) => *flag,
762 Value::String(text) => text == "true",
763 _ => false,
764 }),
765 "number" => match &value {
766 Value::Number(_) => value,
767 Value::String(text) => match text.parse::<f64>().ok().and_then(serde_json::Number::from_f64) {
768 Some(number) => Value::Number(number),
769 None => return fail(FailureCode::Invalid, format!("The input `{name}` is a number.")),
770 },
771 _ => return fail(FailureCode::Invalid, format!("The input `{name}` is a number.")),
772 },
773 "choice" => {
774 let text = g1t_actions::expr::to_text(&value);
775 let options: Vec<String> =
776 spec.get("options").and_then(Value::as_array).map(|o| o.iter().map(g1t_actions::expr::to_text).collect()).unwrap_or_default();
777 if !options.is_empty() && !options.contains(&text) {
778 return fail(FailureCode::Invalid, format!("The input `{name}` is one of {}.", options.join(", ")));
779 }
780 Value::String(text)
781 }
782 _ => Value::String(g1t_actions::expr::to_text(&value)),
783 },
784 };
785 inputs.insert(name.clone(), value);
786 }
787 Outcome::Ok(inputs)
788}