g1t/services/actions/src/lib.rs
| 1 | //! The actions service: a repository's GitHub Actions workflows, run on |
| 2 | //! g1t as they are. See `g1t_contracts::actions` for the methods and |
| 3 | //! `g1t_actions` for how workflows, expressions and filters are read. |
| 4 | //! |
| 5 | //! - [`sync`] reads workflow files, at the default branch for the list and |
| 6 | //! at an event's own commit for its runs. |
| 7 | //! - [`trigger`] turns events, schedules and manual runs into runs. |
| 8 | //! - [`plan`] moves a run's jobs along: each waits for the jobs it needs, |
| 9 | //! is skipped or expanded into its matrix, queued, started in a sandbox |
| 10 | //! when the workspace has room, and finished by the sandbox's report. |
| 11 | //! - [`payload`] builds the webhook-shaped `github.event`. |
| 12 | //! - [`settings`] keeps secrets and variables. |
| 13 | //! |
| 14 | //! The service acts as the repository's workspace: it reads what the |
| 15 | //! workspace can read, and a job's `GITHUB_TOKEN` is a short-lived token of |
| 16 | //! the workspace's. |
| 17 | |
| 18 | mod payload; |
| 19 | mod plan; |
| 20 | mod settings; |
| 21 | mod sync; |
| 22 | mod trigger; |
| 23 | mod views; |
| 24 | |
| 25 | use g1t_contracts::events::Event; |
| 26 | use g1t_contracts::identity::{SlugArgs, Workspace}; |
| 27 | use g1t_contracts::repos::{GetArgs, GetByIdArgs, Repo, RepoPath}; |
| 28 | use g1t_contracts::{FailureCode, Membership, Outcome, PrincipalKind, User, Viewer}; |
| 29 | use g1t_kit::{args, reply, rpc_method}; |
| 30 | use g1t_secrets::Sealer; |
| 31 | use serde::Deserialize; |
| 32 | use serde_json::Value; |
| 33 | use worker::wasm_bindgen::JsValue; |
| 34 | use worker::{Context, D1Database, Env, Fetcher, MessageBatch, MessageExt, Request, Response, Result, ScheduleContext, ScheduledEvent, event}; |
| 35 | |
| 36 | /// The most workflow files read from a repository. |
| 37 | pub const MAX_WORKFLOWS: usize = 50; |
| 38 | /// Jobs one workspace may have running at once; the rest wait their turn. |
| 39 | pub const RUNNING_PER_WORKSPACE: u32 = 4; |
| 40 | /// The longest a job may run, whatever its `timeout-minutes`. |
| 41 | pub const MAX_TIMEOUT_MINUTES: u32 = 60; |
| 42 | /// A running job that has said nothing for this long is taken as lost. |
| 43 | pub const SILENT_MS: u64 = 10 * 60 * 1000; |
| 44 | pub const SITE: &str = "https://g1t.sh"; |
| 45 | pub const API: &str = "https://api.g1t.sh"; |
| 46 | |
| 47 | #[derive(Deserialize)] |
| 48 | struct Count { |
| 49 | n: u32, |
| 50 | } |
| 51 | |
| 52 | pub fn optional(value: Option<&str>) -> JsValue { |
| 53 | value.map_or(JsValue::NULL, JsValue::from) |
| 54 | } |
| 55 | |
| 56 | pub fn fail<T>(code: FailureCode, message: impl Into<String>) -> Outcome<T> { |
| 57 | Outcome::fail(code, message) |
| 58 | } |
| 59 | |
| 60 | /// `owner/name` as a path. |
| 61 | pub fn repo_path(full_name: &str) -> RepoPath { |
| 62 | let (namespace, name) = full_name.split_once('/').unwrap_or((full_name, "")); |
| 63 | RepoPath { |
| 64 | namespace: namespace.to_owned(), |
| 65 | name: name.to_owned(), |
| 66 | } |
| 67 | } |
| 68 | |
| 69 | pub struct Actions { |
| 70 | db: D1Database, |
| 71 | repos: Fetcher, |
| 72 | work: Fetcher, |
| 73 | identity: Fetcher, |
| 74 | runner: Fetcher, |
| 75 | events: Fetcher, |
| 76 | /// Projects: a repository's secrets and variables belong to its project. |
| 77 | projects: Fetcher, |
| 78 | /// Seals secrets; absent until `ACTIONS_KEY` is set, when secrets |
| 79 | /// cannot be saved. |
| 80 | sealer: Option<Sealer>, |
| 81 | } |
| 82 | |
| 83 | impl Actions { |
| 84 | fn new(env: &Env) -> Result<Self> { |
| 85 | Ok(Actions { |
| 86 | db: env.d1("DB")?, |
| 87 | repos: env.service("REPOS")?, |
| 88 | work: env.service("WORK")?, |
| 89 | identity: env.service("IDENTITY")?, |
| 90 | runner: env.service("RUNNER")?, |
| 91 | events: env.service("EVENTS")?, |
| 92 | projects: env.service("PROJECTS")?, |
| 93 | sealer: env.secret("ACTIONS_KEY").ok().and_then(|key| Sealer::new(&key.to_string())), |
| 94 | }) |
| 95 | } |
| 96 | |
| 97 | /// The workspace itself, as the service acts. |
| 98 | async fn workspace_actor(&self, slug: &str) -> Result<Option<User>> { |
| 99 | let workspace: Option<Workspace> = g1t_kit::call(&self.identity, "get_workspace", &SlugArgs { slug: slug.to_owned() }).await?; |
| 100 | Ok(workspace.map(|workspace| User { |
| 101 | id: workspace.id, |
| 102 | username: workspace.slug.clone(), |
| 103 | kind: PrincipalKind::Workspace, |
| 104 | verified: true, |
| 105 | workspaces: vec![Membership::member(workspace.slug)], |
| 106 | ..User::default() |
| 107 | })) |
| 108 | } |
| 109 | |
| 110 | /// The repository, if the viewer may see it and it is not a pull |
| 111 | /// request's working copy. |
| 112 | async fn visible_repo(&self, path: &RepoPath, viewer: &Viewer) -> Result<Option<Repo>> { |
| 113 | let found: Outcome<Repo> = g1t_kit::call( |
| 114 | &self.repos, |
| 115 | "get", |
| 116 | &GetArgs { |
| 117 | path: path.clone(), |
| 118 | viewer: viewer.clone(), |
| 119 | }, |
| 120 | ) |
| 121 | .await?; |
| 122 | Ok(found.into_result().ok().filter(|repo| repo.fork_of.is_none())) |
| 123 | } |
| 124 | |
| 125 | /// A repository by id, as its workspace sees it. |
| 126 | async fn repo_by_id(&self, id: &str) -> Result<Option<(Repo, User)>> { |
| 127 | let path: Option<RepoPath> = g1t_kit::call(&self.repos, "path_by_id", &g1t_contracts::repos::PathByIdArgs { id: id.to_owned() }).await?; |
| 128 | let Some(path) = path else { return Ok(None) }; |
| 129 | let Some(actor) = self.workspace_actor(&path.namespace).await? else { |
| 130 | return Ok(None); |
| 131 | }; |
| 132 | let found: Outcome<Repo> = g1t_kit::call( |
| 133 | &self.repos, |
| 134 | "get_by_id", |
| 135 | &GetByIdArgs { |
| 136 | id: id.to_owned(), |
| 137 | viewer: Some(actor.clone()), |
| 138 | }, |
| 139 | ) |
| 140 | .await?; |
| 141 | Ok(found.into_result().ok().filter(|repo| repo.fork_of.is_none()).map(|repo| (repo, actor))) |
| 142 | } |
| 143 | |
| 144 | /// Refuses anyone but a member of the repository's workspace. |
| 145 | fn member(actor: &User, repo: &RepoPath) -> Option<Outcome<()>> { |
| 146 | (actor.kind == PrincipalKind::Agent || !actor.is_member(&repo.namespace.to_lowercase())) |
| 147 | .then(|| fail(FailureCode::Forbidden, format!("Only members of {} can do that.", repo.namespace))) |
| 148 | } |
| 149 | } |
| 150 | |
| 151 | /// Unwraps an `Outcome`, or returns its failure from the enclosing method. |
| 152 | #[macro_export] |
| 153 | macro_rules! check { |
| 154 | ($outcome:expr) => { |
| 155 | match $outcome { |
| 156 | g1t_contracts::Outcome::Ok(value) => value, |
| 157 | g1t_contracts::Outcome::Fail(refused) => return Ok(g1t_contracts::Outcome::Fail(refused)), |
| 158 | } |
| 159 | }; |
| 160 | } |
| 161 | |
| 162 | #[event(fetch)] |
| 163 | async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> { |
| 164 | let Some(method) = rpc_method(&request) else { |
| 165 | return Response::error("Not found", 404); |
| 166 | }; |
| 167 | let body: Value = request.json().await?; |
| 168 | let service = Actions::new(&env)?; |
| 169 | match method.as_str() { |
| 170 | "workflows" => reply(&service.workflows(args(body)?).await?), |
| 171 | "runs" => reply(&service.runs(args(body)?).await?), |
| 172 | "run" => reply(&service.run(args(body)?).await?), |
| 173 | "logs" => reply(&service.logs(args(body)?).await?), |
| 174 | "dispatch" => reply(&service.dispatch(args(body)?).await?), |
| 175 | "merge_group" => reply(&service.merge_group(args(body)?).await?), |
| 176 | "cancel" => reply(&service.cancel(args(body)?).await?), |
| 177 | "rerun" => reply(&service.rerun(args(body)?).await?), |
| 178 | "set_workflow_enabled" => reply(&service.set_workflow_enabled(args(body)?).await?), |
| 179 | "settings" => reply(&service.settings(args(body)?).await?), |
| 180 | "set_setting" => reply(&service.set_setting(args(body)?).await?), |
| 181 | "delete_setting" => reply(&service.delete_setting(args(body)?).await?), |
| 182 | "resolve_settings" => reply(&service.resolve_settings(args(body)?).await?), |
| 183 | "job_spec" => reply(&service.job_spec(args(body)?).await?), |
| 184 | "job_auth" => reply(&service.job_auth(args(body)?).await?), |
| 185 | "job_report" => reply(&service.job_report(args(body)?).await?), |
| 186 | _ => Response::error("Unknown method", 404), |
| 187 | } |
| 188 | } |
| 189 | |
| 190 | /// Events from the bus, on this service's own queue. |
| 191 | #[event(queue)] |
| 192 | async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> { |
| 193 | let service = Actions::new(&env)?; |
| 194 | for message in batch.messages()? { |
| 195 | if let Err(error) = service.on_event(message.body()).await { |
| 196 | worker::console_error!("actions: event {} failed: {error}", message.body().id); |
| 197 | message.retry(); |
| 198 | continue; |
| 199 | } |
| 200 | message.ack(); |
| 201 | } |
| 202 | Ok(()) |
| 203 | } |
| 204 | |
| 205 | /// Every minute: schedules that fire, jobs waiting for room, and jobs |
| 206 | /// whose sandbox went quiet. |
| 207 | #[event(scheduled)] |
| 208 | async fn scheduled(_event: ScheduledEvent, env: Env, _ctx: ScheduleContext) { |
| 209 | match Actions::new(&env) { |
| 210 | Ok(service) => { |
| 211 | if let Err(error) = service.on_minute(g1t_kit::now_ms()).await { |
| 212 | worker::console_error!("actions: the sweep failed: {error}"); |
| 213 | } |
| 214 | } |
| 215 | Err(error) => worker::console_error!("actions: could not start: {error}"), |
| 216 | } |
| 217 | } |