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

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