g1t/services/actions/src/trigger.rs

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