A push is checked once and side by side: its pack is read and its bases fetched once for the rules, the workflow gate and the secret scan, which run together, other services are asked while the pack is read, and cache writes and rule records finish after git has its answer
The trees, histories and blobs each check walks are read a level, a commit or a pair at a time instead of one by one; a thin pack's base is asked for as what the pack's own trees say it is, or as a blob and a tree at once. What a push is refused for, and in which order, is unchanged.
10 files+658−2970/10 viewed
| 289 | 289 | | `mint` | Only when no credential was kept: the git store making one for the request | | |
| 290 | 290 | | `store` | The git store's answer to a clone or fetch | | |
| 291 | 291 | | `recv` | Only for a push: receiving it from git | | |
| 292 | − | | `rules` | Only for a push: checking it against the rules of the branches and tags it changes, and for workflow files a token may not change | | |
| 293 | − | | `scan` | Only for a push: checking it for secrets and private email addresses | | |
| 292 | + | | `checks` | Only for a push: checking it against the rules of the branches and tags it changes, for workflow files a token may not change, and for secrets and private email addresses, all at once | | |
| 294 | 293 | | `upload` | Only for a push: handing it to the git store and its answer | | |
| 295 | 294 | | `refs` | Only for a push: recording that the repository's refs changed | | |
| 296 | 295 | | `total` | Everything g1t did | | |
| 296 | + | ||
| 297 | + | A push's checks run side by side, so each also has its own entry, after | |
| 298 | + | the steps and not counted in the total: `read` (reading the push's objects | |
| 299 | + | and fetching what they build on from the repository), `rules`, `scan` | |
| 300 | + | (secrets and email addresses) and, for an access token, `gate` (workflow | |
| 301 | + | files). | |
| 297 | 302 | | `repos` | The same, measured where your request arrived | | |
| 298 | 303 | ||
| 299 | 304 | Two entries say how a step went rather than how long it took: |
| 396 | 396 | _ => None, | |
| 397 | 397 | } | |
| 398 | 398 | } | |
| 399 | + | ||
| 400 | + | /// What the pack's own commits and trees say the objects they name | |
| 401 | + | /// are: a commit's tree and parents, a tree's entries (submodules | |
| 402 | + | /// left out). How a missing base can be asked for as what it is. | |
| 403 | + | pub fn named_kinds(&self) -> HashMap<String, ObjectKind> { | |
| 404 | + | let mut kinds = HashMap::new(); | |
| 405 | + | for (kind, data) in self.objects.values() { | |
| 406 | + | match kind { | |
| 407 | + | ObjectKind::Commit => { | |
| 408 | + | let commit = parse_commit(data); | |
| 409 | + | kinds.insert(commit.tree, ObjectKind::Tree); | |
| 410 | + | kinds.extend(commit.parents.into_iter().map(|parent| (parent, ObjectKind::Commit))); | |
| 411 | + | } | |
| 412 | + | ObjectKind::Tree => { | |
| 413 | + | for item in parse_tree(data) { | |
| 414 | + | if item.is_tree() { | |
| 415 | + | kinds.insert(item.id, ObjectKind::Tree); | |
| 416 | + | } else if item.mode != "160000" { | |
| 417 | + | kinds.insert(item.id, ObjectKind::Blob); | |
| 418 | + | } | |
| 419 | + | } | |
| 420 | + | } | |
| 421 | + | _ => {} | |
| 422 | + | } | |
| 423 | + | } | |
| 424 | + | kinds | |
| 425 | + | } | |
| 399 | 426 | } | |
| 400 | 427 | ||
| 401 | 428 | /// What a commit says about its place in history, and whose it is. | |
| ⋯ | |||
| 571 | 598 | } | |
| 572 | 599 | ||
| 573 | 600 | #[test] | |
| 601 | + | fn a_pack_says_what_the_objects_it_names_are() { | |
| 602 | + | let blob_id = object_id(ObjectKind::Blob, b"x"); | |
| 603 | + | let sub_id = "1".repeat(40); | |
| 604 | + | let module_id = "2".repeat(40); | |
| 605 | + | let tree = encode_tree(&[ | |
| 606 | + | TreeItem { mode: "100644".into(), name: "a.txt".into(), id: blob_id.clone() }, | |
| 607 | + | TreeItem { mode: "40000".into(), name: "src".into(), id: sub_id.clone() }, | |
| 608 | + | TreeItem { mode: "160000".into(), name: "vendor".into(), id: module_id.clone() }, | |
| 609 | + | ]); | |
| 610 | + | let tree_id = object_id(ObjectKind::Tree, &tree); | |
| 611 | + | let parent = "3".repeat(40); | |
| 612 | + | let commit = format!("tree {tree_id}\nparent {parent}\nauthor A <a@x> 1 +0000\n\nm\n"); | |
| 613 | + | let pack = Pack::parse(&build_pack(&[(ObjectKind::Tree, tree), (ObjectKind::Commit, commit.into_bytes())], &[])).unwrap(); | |
| 614 | + | let kinds = pack.named_kinds(); | |
| 615 | + | assert_eq!(kinds.get(&blob_id), Some(&ObjectKind::Blob)); | |
| 616 | + | assert_eq!(kinds.get(&sub_id), Some(&ObjectKind::Tree)); | |
| 617 | + | assert_eq!(kinds.get(&tree_id), Some(&ObjectKind::Tree)); | |
| 618 | + | assert_eq!(kinds.get(&parent), Some(&ObjectKind::Commit)); | |
| 619 | + | // A submodule's commit lives in another repository. | |
| 620 | + | assert_eq!(kinds.get(&module_id), None); | |
| 621 | + | } | |
| 622 | + | ||
| 623 | + | #[test] | |
| 574 | 624 | fn commits_and_trees_are_parsed_and_trees_rebuilt() { | |
| 575 | 625 | let blob_id = object_id(ObjectKind::Blob, b"x"); | |
| 576 | 626 | let tree = encode_tree(&[ | |
| 14 | 14 | ||
| 15 | 15 | use crate::meters; | |
| 16 | 16 | use crate::pack_limits::{PackSizer, Violation}; | |
| 17 | + | use crate::push_checks::{Checks, Verdict}; | |
| 17 | 18 | use crate::resilience::{self, Busy, Failure}; | |
| 18 | 19 | ||
| 19 | 20 | const ENDPOINTS: [&str; 3] = ["info/refs", "git-upload-pack", "git-receive-pack"]; | |
| ⋯ | |||
| 87 | 88 | /// Step names and whole milliseconds only. The Workers clock moves only | |
| 88 | 89 | /// while a request waits on something, so each step is the time spent | |
| 89 | 90 | /// waiting on the database, another service or the git store. A note | |
| 90 | − | /// says how a step went without a duration: `refs;desc=hit-colo`. | |
| 91 | + | /// says how a step went without a duration: `refs;desc=hit-colo`. A part | |
| 92 | + | /// is one of several things a step did at once, with its own duration: a | |
| 93 | + | /// push's `checks` step is its `read`, `rules` and `scan` parts side by side. | |
| 91 | 94 | pub struct Timing { | |
| 92 | 95 | started: u64, | |
| 93 | 96 | last: u64, | |
| 94 | 97 | steps: Vec<(&'static str, u64)>, | |
| 98 | + | parts: Vec<(&'static str, u64)>, | |
| 95 | 99 | notes: Vec<(&'static str, &'static str)>, | |
| 96 | 100 | } | |
| 97 | 101 | ||
| ⋯ | |||
| 102 | 106 | started: now, | |
| 103 | 107 | last: now, | |
| 104 | 108 | steps: Vec::new(), | |
| 109 | + | parts: Vec::new(), | |
| 105 | 110 | notes: Vec::new(), | |
| 106 | 111 | } | |
| 107 | 112 | } | |
| 108 | 113 | ||
| 114 | + | /// Says how long `part` of the step under way took. Parts overlap, so | |
| 115 | + | /// they are listed after the steps and are not part of the total. | |
| 116 | + | pub fn part(&mut self, part: &'static str, ms: u64) { | |
| 117 | + | self.parts.push((part, ms)); | |
| 118 | + | } | |
| 119 | + | ||
| 109 | 120 | /// Says how `name` went: where an answer or a credential came from. | |
| 110 | 121 | pub fn note(&mut self, name: &'static str, description: &'static str) { | |
| 111 | 122 | self.notes.push((name, description)); | |
| ⋯ | |||
| 122 | 133 | pub fn apply(&self, response: Response) -> Result<Response> { | |
| 123 | 134 | let total = g1t_kit::now_ms().saturating_sub(self.started); | |
| 124 | 135 | let headers = response.headers().clone(); | |
| 125 | − | headers.set("server-timing", &server_timing(&self.steps, &self.notes, total))?; | |
| 136 | + | headers.set("server-timing", &server_timing(&self.steps, &self.parts, &self.notes, total))?; | |
| 126 | 137 | Ok(response.with_headers(headers)) | |
| 127 | 138 | } | |
| 128 | 139 | } | |
| 129 | 140 | ||
| 130 | − | /// A `Server-Timing` value: each step with its duration, the notes, then | |
| 131 | − | /// the total. | |
| 132 | − | fn server_timing(steps: &[(&str, u64)], notes: &[(&str, &str)], total: u64) -> String { | |
| 141 | + | /// A `Server-Timing` value: each step with its duration, then each part, | |
| 142 | + | /// the notes, and the total. | |
| 143 | + | fn server_timing(steps: &[(&str, u64)], parts: &[(&str, u64)], notes: &[(&str, &str)], total: u64) -> String { | |
| 133 | 144 | steps | |
| 134 | 145 | .iter() | |
| 146 | + | .chain(parts) | |
| 135 | 147 | .map(|(step, ms)| format!("{step};dur={ms}")) | |
| 136 | 148 | .chain(notes.iter().map(|(name, description)| format!("{name};desc={description}"))) | |
| 137 | 149 | .chain(std::iter::once(format!("total;dur={total}"))) | |
| ⋯ | |||
| 844 | 856 | } | |
| 845 | 857 | ||
| 846 | 858 | /// Sends the request on to the git store and returns its response as is, | |
| 847 | − | /// unless it is a push the rules refuse (`rules`, rules.rs), one the | |
| 848 | − | /// store could not hold (`limits`, pack_limits.rs), or one that `scan` | |
| 849 | − | /// (push protection) answers itself. A fetch's ref listing has its `HEAD` | |
| 850 | − | /// pointed at `default_branch` (see [`with_head`]). A POST's body is | |
| 851 | − | /// `read` when the caller has read it already. | |
| 859 | + | /// unless it is a push `checks` refuse (the workflow gate, the rules, push | |
| 860 | + | /// protection: push_checks.rs), or one the store could not hold (`limits`, | |
| 861 | + | /// pack_limits.rs). `checks` is given as much of the push as was read, | |
| 862 | + | /// whether that is all of it, and whether to scan it for secrets. A | |
| 863 | + | /// fetch's ref listing has its `HEAD` pointed at `default_branch` (see | |
| 864 | + | /// [`with_head`]). A POST's body is `read` when the caller has read it | |
| 865 | + | /// already. | |
| 852 | 866 | /// | |
| 853 | 867 | /// A push is read as it arrives: up to `limits.scan_cap` is kept, to be | |
| 854 | 868 | /// scanned and sent on whole; past it, the push is declined, or streamed | |
| ⋯ | |||
| 862 | 876 | read: Option<Vec<u8>>, | |
| 863 | 877 | git: &GitRequest, | |
| 864 | 878 | access: &GitAccess, | |
| 865 | − | rules: impl AsyncFnOnce(&[u8], bool) -> Result<Option<Response>>, | |
| 879 | + | checks: impl AsyncFnOnce(&[u8], bool, bool) -> Result<Checks>, | |
| 866 | 880 | default_branch: Option<&str>, | |
| 867 | 881 | limits: PushLimits, | |
| 868 | − | scan: impl AsyncFnOnce(&[u8]) -> Result<Option<Response>>, | |
| 869 | 882 | timing: &mut Timing, | |
| 870 | 883 | ) -> Result<Push> { | |
| 871 | 884 | let headers = Headers::new(); | |
| ⋯ | |||
| 882 | 895 | let namespace = crate::store::health_namespace(&access.remote); | |
| 883 | 896 | ||
| 884 | 897 | if method == Method::Post && git.endpoint == "git-receive-pack" { | |
| 885 | − | return push(request, &url, headers, rules, limits, scan, &namespace, timing).await; | |
| 898 | + | return push(request, &url, headers, checks, limits, &namespace, timing).await; | |
| 886 | 899 | } | |
| 887 | 900 | ||
| 888 | 901 | // A read: the ref advertisement, `ls-refs`, or a fetch of objects. | |
| ⋯ | |||
| 964 | 977 | } | |
| 965 | 978 | ||
| 966 | 979 | /// A receive-pack request; see [`forward`]. Its steps: `recv` (the push | |
| 967 | − | /// read), `rules`, `scan` (push protection), `upload` (the store's answer). | |
| 980 | + | /// read), `checks`, `upload` (the store's answer). | |
| 968 | 981 | #[allow(clippy::too_many_arguments)] | |
| 969 | 982 | async fn push( | |
| 970 | 983 | mut request: Request, | |
| 971 | 984 | url: &str, | |
| 972 | 985 | headers: Headers, | |
| 973 | − | rules: impl AsyncFnOnce(&[u8], bool) -> Result<Option<Response>>, | |
| 986 | + | checks: impl AsyncFnOnce(&[u8], bool, bool) -> Result<Checks>, | |
| 974 | 987 | limits: PushLimits, | |
| 975 | − | scan: impl AsyncFnOnce(&[u8]) -> Result<Option<Response>>, | |
| 976 | 988 | namespace: &str, | |
| 977 | 989 | timing: &mut Timing, | |
| 978 | 990 | ) -> Result<Push> { | |
| ⋯ | |||
| 997 | 1009 | } | |
| 998 | 1010 | } | |
| 999 | 1011 | timing.mark("recv"); | |
| 1000 | − | // The rules of the branches and tags it changes, first: what they | |
| 1001 | − | // refuse is refused whatever else is wrong with it. | |
| 1002 | − | let ruled = rules(&head, ended).await?; | |
| 1003 | − | timing.mark("rules"); | |
| 1004 | − | if let Some(response) = ruled { | |
| 1005 | − | if !ended { | |
| 1006 | − | drain(&mut stream).await?; | |
| 1007 | − | } | |
| 1008 | − | return Ok(Push::Refused(response)); | |
| 1012 | + | // The workflow gate and the rules of the branches and tags it changes, | |
| 1013 | + | // with push protection alongside for a push read whole that is within | |
| 1014 | + | // the size limits (push_checks.rs). What the rules refuse is refused | |
| 1015 | + | // whatever else is wrong with it. | |
| 1016 | + | let checked = checks(&head, ended, ended && violation.is_none()).await?; | |
| 1017 | + | for (part, ms) in checked.spans { | |
| 1018 | + | timing.part(part, ms); | |
| 1009 | 1019 | } | |
| 1020 | + | timing.mark("checks"); | |
| 1021 | + | let blocked = match checked.verdict { | |
| 1022 | + | Verdict::Refused(response) => { | |
| 1023 | + | if !ended { | |
| 1024 | + | drain(&mut stream).await?; | |
| 1025 | + | } | |
| 1026 | + | return Ok(Push::Refused(response)); | |
| 1027 | + | } | |
| 1028 | + | Verdict::Blocked(response) => Some(response), | |
| 1029 | + | Verdict::Clear => None, | |
| 1030 | + | }; | |
| 1010 | 1031 | if !ended && limits.large == LargePushes::Refuse && violation.is_none() { | |
| 1011 | 1032 | let size = head.len() as u64 + drain(&mut stream).await?; | |
| 1012 | 1033 | violation = Some(SizeViolation::Unscannable { size, cap: limits.scan_cap }); | |
| ⋯ | |||
| 1024 | 1045 | init.with_method(Method::Post).with_headers(headers); | |
| 1025 | 1046 | let started = g1t_kit::now_ms(); | |
| 1026 | 1047 | let (answered, pack_bytes, sent, unscanned) = if ended { | |
| 1027 | − | let scanned = scan(&head).await?; | |
| 1028 | − | timing.mark("scan"); | |
| 1029 | − | if let Some(response) = scanned { | |
| 1048 | + | if let Some(response) = blocked { | |
| 1030 | 1049 | return Ok(Push::Blocked(response)); | |
| 1031 | 1050 | } | |
| 1032 | 1051 | let pack = pack_bytes(&head); | |
| ⋯ | |||
| 1113 | 1132 | #[test] | |
| 1114 | 1133 | fn server_timing_names_each_step_and_the_total() { | |
| 1115 | 1134 | assert_eq!( | |
| 1116 | − | server_timing(&[("repo", 12), ("token", 0), ("store", 140)], &[], 153), | |
| 1135 | + | server_timing(&[("repo", 12), ("token", 0), ("store", 140)], &[], &[], 153), | |
| 1117 | 1136 | "repo;dur=12, token;dur=0, store;dur=140, total;dur=153" | |
| 1118 | 1137 | ); | |
| 1119 | − | assert_eq!(server_timing(&[], &[], 3), "total;dur=3"); | |
| 1138 | + | assert_eq!(server_timing(&[], &[], &[], 3), "total;dur=3"); | |
| 1139 | + | // A push's checks run side by side: their parts come after the steps. | |
| 1140 | + | assert_eq!( | |
| 1141 | + | server_timing(&[("recv", 30), ("checks", 120), ("upload", 300)], &[("read", 40), ("rules", 110), ("scan", 90)], &[], 450), | |
| 1142 | + | "recv;dur=30, checks;dur=120, upload;dur=300, read;dur=40, rules;dur=110, scan;dur=90, total;dur=450" | |
| 1143 | + | ); | |
| 1120 | 1144 | assert_eq!( | |
| 1121 | − | server_timing(&[("repo", 1), ("cache", 2)], &[("refs", "hit-colo")], 4), | |
| 1145 | + | server_timing(&[("repo", 1), ("cache", 2)], &[], &[("refs", "hit-colo")], 4), | |
| 1122 | 1146 | "repo;dur=1, cache;dur=2, refs;desc=hit-colo, total;dur=4" | |
| 1123 | 1147 | ); | |
| 1124 | 1148 | } | |
| 32 | 32 | mod namespaces; | |
| 33 | 33 | mod pack_cache; | |
| 34 | 34 | mod pack_limits; | |
| 35 | + | mod push_checks; | |
| 35 | 36 | mod push_commits; | |
| 36 | 37 | mod refs; | |
| 37 | 38 | mod refs_cache; | |
| ⋯ | |||
| 225 | 226 | shared: Option<Rc<shared::Shared>>, | |
| 226 | 227 | /// Packs for fresh clones (pack_cache.rs); `None` without the bucket. | |
| 227 | 228 | packs: Option<Rc<pack_cache::Packs>>, | |
| 229 | + | /// What this request started that need not hold up its answer | |
| 230 | + | /// (store.rs), handed to its `waitUntil` once it has answered. | |
| 231 | + | deferred: Rc<store::Deferred>, | |
| 228 | 232 | } | |
| 229 | 233 | ||
| 230 | 234 | impl<S: GitStore> Repos<S> { | |
| ⋯ | |||
| 1867 | 1871 | .zip(refs_cache::usable(registry::refs_state(&repo.id), now_ms())) | |
| 1868 | 1872 | .map(|(normalized, version)| pack_cache::Key::new(&repo.id, version, &normalized)); | |
| 1869 | 1873 | // A kept answer and the free workspace limits, with a kept | |
| 1870 | − | // credential looked up alongside. A kept answer goes back without | |
| 1871 | − | // waiting for the credential, which it does not need. | |
| 1872 | − | let ((answer, pack, limited), kept_access) = { | |
| 1874 | + | // credential and, for a push, what the repository holds looked up | |
| 1875 | + | // alongside. A kept answer goes back without waiting for the | |
| 1876 | + | // credential, which it does not need. | |
| 1877 | + | let pushing = write && !get; | |
| 1878 | + | let ((answer, pack, (limited, held)), kept_access) = { | |
| 1873 | 1879 | let shared = self.shared.as_deref(); | |
| 1874 | 1880 | let answer_and_limits = std::pin::pin!(futures_util::future::join3( | |
| 1875 | 1881 | async { | |
| ⋯ | |||
| 1884 | 1890 | _ => None, | |
| 1885 | 1891 | } | |
| 1886 | 1892 | }, | |
| 1887 | − | self.git_limits(call, git, &repo, env), | |
| 1893 | + | futures_util::future::join(self.git_limits(call, git, &repo, env), async { | |
| 1894 | + | if pushing { Some(self.held(&repo).await) } else { None } | |
| 1895 | + | }), | |
| 1888 | 1896 | )); | |
| 1889 | 1897 | let kept_access = std::pin::pin!(self.store.kept_access(&key, scope)); | |
| 1890 | 1898 | match futures_util::future::select(answer_and_limits, kept_access).await { | |
| 1891 | 1899 | futures_util::future::Either::Left((first, kept_access)) => { | |
| 1892 | − | let answered = first.0.is_some() || first.1.is_some() || matches!(first.2, Ok(Some(_)) | Err(_)); | |
| 1900 | + | let answered = first.0.is_some() || first.1.is_some() || matches!(first.2.0, Ok(Some(_)) | Err(_)); | |
| 1893 | 1901 | (first, if answered { None } else { kept_access.await }) | |
| 1894 | 1902 | } | |
| 1895 | 1903 | futures_util::future::Either::Right((kept_access, first)) => (first.await, kept_access), | |
| ⋯ | |||
| 1958 | 1966 | // request is tried again with a new one; the requests after it then | |
| 1959 | 1967 | // have that one too. | |
| 1960 | 1968 | let again = if get { Some(request.clone()?) } else { None }; | |
| 1961 | − | // Push protection: a push that adds a secret is refused. See secret_scan.rs. | |
| 1962 | − | let scan = async |body: &[u8]| self.protect(&repo, viewer.as_ref(), body).await; | |
| 1963 | − | // Rulesets: what the rules of the branches and tags it changes | |
| 1964 | − | // refuse is declined, saying which rule and why (rules.rs). | |
| 1965 | − | let rules = async |head: &[u8], whole: bool| self.check_push(&repo, viewer.as_ref(), head, whole).await; | |
| 1969 | + | // A push's checks: workflow files, the rules of the branches and | |
| 1970 | + | // tags it changes, saying which rule and why (rules.rs), and push | |
| 1971 | + | // protection, which refuses a push that adds a secret | |
| 1972 | + | // (secret_scan.rs). See push_checks.rs. | |
| 1973 | + | let checks = async |head: &[u8], whole: bool, scan: bool| self.check_receive(&repo, viewer.as_ref(), head, whole, scan).await; | |
| 1966 | 1974 | // What a push may bring (pack_limits.rs): the repository's size is | |
| 1967 | 1975 | // its own and its pull requests' working copies'. | |
| 1968 | − | let limits = if write && !get { | |
| 1976 | + | let limits = if let Some(held) = held { | |
| 1969 | 1977 | git_http::PushLimits { | |
| 1970 | − | held: self.held(&repo).await, | |
| 1978 | + | held, | |
| 1971 | 1979 | repo_limit: self.repo_limit, | |
| 1972 | 1980 | large: self.large_pushes, | |
| 1973 | 1981 | ..git_http::PushLimits::default() | |
| ⋯ | |||
| 1980 | 1988 | body, | |
| 1981 | 1989 | git, | |
| 1982 | 1990 | &access, | |
| 1983 | − | rules, | |
| 1991 | + | checks, | |
| 1984 | 1992 | default_branch.as_deref(), | |
| 1985 | 1993 | limits, | |
| 1986 | − | scan, | |
| 1987 | 1994 | timing, | |
| 1988 | 1995 | ) | |
| 1989 | 1996 | .await?; | |
| ⋯ | |||
| 1995 | 2002 | self.store.forget_access(&key).await; | |
| 1996 | 2003 | if let Some(again) = again { | |
| 1997 | 2004 | let access = self.store.mint_access(&key, scope).await?; | |
| 1998 | − | let nothing = async |_: &[u8]| Ok(None); | |
| 1999 | 2005 | outcome = git_http::forward( | |
| 2000 | 2006 | again, | |
| 2001 | 2007 | None, | |
| 2002 | 2008 | git, | |
| 2003 | 2009 | &access, | |
| 2004 | − | async |_: &[u8], _: bool| Ok(None), | |
| 2010 | + | async |_: &[u8], _: bool, _: bool| Ok(push_checks::Checks::clear()), | |
| 2005 | 2011 | default_branch.as_deref(), | |
| 2006 | 2012 | git_http::PushLimits::default(), | |
| 2007 | − | nothing, | |
| 2008 | 2013 | timing, | |
| 2009 | 2014 | ) | |
| 2010 | 2015 | .await?; | |
| ⋯ | |||
| 2341 | 2346 | { | |
| 2342 | 2347 | worker::console_error!("push not recorded: {error}"); | |
| 2343 | 2348 | } | |
| 2349 | + | repos.deferred.settle().await; | |
| 2344 | 2350 | }); | |
| 2345 | 2351 | } | |
| 2346 | 2352 | } | |
| 2347 | 2353 | ||
| 2348 | 2354 | fn service(env: &Env) -> Result<Repos<ArtifactsStore>> { | |
| 2349 | 2355 | let shared = shared::Shared::from_env(env).map(Rc::new); | |
| 2356 | + | let deferred = Rc::new(store::Deferred::default()); | |
| 2350 | 2357 | Ok(Repos { | |
| 2351 | 2358 | registry: Registry { db: env.d1("DB")? }, | |
| 2352 | − | store: ArtifactsStore::new(env, shared.clone())?, | |
| 2359 | + | store: ArtifactsStore::new(env, shared.clone(), deferred.clone())?, | |
| 2360 | + | deferred, | |
| 2353 | 2361 | shared, | |
| 2354 | 2362 | packs: pack_cache::Packs::from_env(env).map(Rc::new), | |
| 2355 | 2363 | events: env.service("EVENTS")?, | |
| ⋯ | |||
| 2435 | 2443 | && let Some((job_id, number)) = backup_part_path(&request.path()) | |
| 2436 | 2444 | { | |
| 2437 | 2445 | let answered = backup_part(&mut request, &env, &repos, job_id, number).await; | |
| 2446 | + | repos.deferred.hand_over(&ctx); | |
| 2438 | 2447 | flush_later(&env, &ctx); | |
| 2439 | 2448 | return answered; | |
| 2440 | 2449 | } | |
| 2441 | 2450 | let Some(method) = rpc_method(&request) else { | |
| 2442 | 2451 | let answered = repos.git_http(request, &env, &ctx).await; | |
| 2452 | + | repos.deferred.hand_over(&ctx); | |
| 2443 | 2453 | flush_later(&env, &ctx); | |
| 2444 | 2454 | return answered; | |
| 2445 | 2455 | }; | |
| ⋯ | |||
| 2686 | 2696 | }, | |
| 2687 | 2697 | answered => answered, | |
| 2688 | 2698 | }; | |
| 2699 | + | repos.deferred.hand_over(&ctx); | |
| 2689 | 2700 | flush_later(&env, &ctx); | |
| 2690 | 2701 | served.finish(answered) | |
| 2691 | 2702 | } | |
Binary or large file; its contents are not shown.
| 7 | 7 | use std::cell::Cell; | |
| 8 | 8 | use std::collections::{HashMap, HashSet, VecDeque}; | |
| 9 | 9 | ||
| 10 | + | use futures_util::future::{join_all, try_join_all}; | |
| 10 | 11 | use g1t_contracts::rules::{CommitFacts, FileChange, Signature}; | |
| 11 | 12 | use g1t_scan::pack::{ObjectKind, Pack}; | |
| 12 | 13 | use worker::Result; | |
| ⋯ | |||
| 89 | 90 | ||
| 90 | 91 | /// The files that differ between two trees, deletions included, with the | |
| 91 | 92 | /// size of each new blob the pack holds. Whether the list is complete. | |
| 93 | + | /// Each level of the trees is read at once. | |
| 92 | 94 | async fn changed<R: GitRepo>(objects: &Objects<'_, R>, old_root: Option<String>, new_root: Option<String>) -> Result<(Vec<FileChange>, bool)> { | |
| 93 | 95 | let mut files = Vec::new(); | |
| 94 | 96 | let mut level: Vec<(String, Option<String>, Option<String>)> = vec![(String::new(), old_root, new_root)]; | |
| 95 | 97 | while !level.is_empty() { | |
| 98 | + | let read = try_join_all(level.iter().map(|(_, old, new)| async move { | |
| 99 | + | let tree = async |id: &Option<String>| match id { | |
| 100 | + | Some(id) => objects.tree(id).await, | |
| 101 | + | None => Ok(Vec::new()), | |
| 102 | + | }; | |
| 103 | + | futures_util::future::try_join(tree(old), tree(new)).await | |
| 104 | + | })) | |
| 105 | + | .await?; | |
| 96 | 106 | let mut next = Vec::new(); | |
| 97 | − | for (prefix, old, new) in level { | |
| 98 | − | let old_items = match &old { | |
| 99 | − | Some(id) => objects.tree(id).await?, | |
| 100 | − | None => Vec::new(), | |
| 101 | − | }; | |
| 102 | − | let new_items = match &new { | |
| 103 | − | Some(id) => objects.tree(id).await?, | |
| 104 | − | None => Vec::new(), | |
| 105 | − | }; | |
| 107 | + | for ((prefix, _, _), (old_items, new_items)) in level.into_iter().zip(read) { | |
| 106 | 108 | for item in &new_items { | |
| 107 | 109 | let before = old_items.iter().find(|entry| entry.name == item.name); | |
| 108 | 110 | if before.is_some_and(|before| before.id == item.id && before.mode == item.mode) { | |
| ⋯ | |||
| 172 | 174 | } | |
| 173 | 175 | ||
| 174 | 176 | /// Whether `old` is in the history of `new`: a fast-forward. Walks the | |
| 175 | − | /// pack's commits, then the repository's history from where it leaves it. | |
| 177 | + | /// pack's commits, then the repository's history from where it leaves it, | |
| 178 | + | /// reading the histories from each place it leaves at once. | |
| 176 | 179 | pub async fn contains<R: GitRepo>(pack: &Pack, repo: &R, new: &str, old: &str, depth: u32) -> Result<bool> { | |
| 177 | 180 | if new == old { | |
| 178 | 181 | return Ok(true); | |
| ⋯ | |||
| 192 | 195 | _ => boundary.push(id), | |
| 193 | 196 | } | |
| 194 | 197 | } | |
| 195 | − | for start in boundary.iter().take(20) { | |
| 196 | − | let history = repo.log(start, depth).await?; | |
| 197 | − | if history.iter().any(|commit| commit.hash == old) { | |
| 198 | − | return Ok(true); | |
| 199 | − | } | |
| 200 | − | // Merges: the history is first-parent only, so look along the | |
| 201 | − | // second parents it names too, a step at a time. | |
| 202 | − | for commit in history.iter().filter(|commit| commit.parents.len() > 1).take(10) { | |
| 203 | − | for parent in commit.parents.iter().skip(1) { | |
| 204 | − | if parent == old || repo.log(parent, depth).await?.iter().any(|commit| commit.hash == old) { | |
| 205 | − | return Ok(true); | |
| 206 | − | } | |
| 207 | − | } | |
| 208 | − | } | |
| 198 | + | let histories = join_all(boundary.iter().take(20).map(|start| repo.log(start, depth))).await; | |
| 199 | + | if found_in(&histories, old) { | |
| 200 | + | return Ok(true); | |
| 201 | + | } | |
| 202 | + | // Merges: the history is first-parent only, so look along the second | |
| 203 | + | // parents it names too. | |
| 204 | + | let seconds: Vec<&String> = histories | |
| 205 | + | .iter() | |
| 206 | + | .flatten() | |
| 207 | + | .flat_map(|history| history.iter().filter(|commit| commit.parents.len() > 1).take(10)) | |
| 208 | + | .flat_map(|commit| commit.parents.iter().skip(1)) | |
| 209 | + | .collect(); | |
| 210 | + | if seconds.iter().any(|parent| *parent == old) { | |
| 211 | + | return Ok(true); | |
| 212 | + | } | |
| 213 | + | let further = join_all(seconds.iter().map(|parent| repo.log(parent, depth))).await; | |
| 214 | + | if found_in(&further, old) { | |
| 215 | + | return Ok(true); | |
| 216 | + | } | |
| 217 | + | // Not found: a history that could not be read may have held it. | |
| 218 | + | for history in histories.into_iter().chain(further) { | |
| 219 | + | history?; | |
| 209 | 220 | } | |
| 210 | 221 | Ok(false) | |
| 211 | 222 | } | |
| 212 | 223 | ||
| 224 | + | /// Whether any history read holds `old`. | |
| 225 | + | fn found_in(histories: &[Result<Vec<g1t_contracts::repos::Commit>>], old: &str) -> bool { | |
| 226 | + | histories.iter().flatten().any(|history| history.iter().any(|commit| commit.hash == old)) | |
| 227 | + | } | |
| 228 | + | ||
| 213 | 229 | /// The signature fingerprints and committer addresses of commits, for | |
| 214 | 230 | /// looking up who owns them. | |
| 215 | 231 | pub fn signing_facts(pack: &Pack, ids: &[String]) -> (Vec<String>, Vec<String>) { | |
| ⋯ | |||
| 232 | 248 | } | |
| 233 | 249 | ||
| 234 | 250 | /// Every commit of `ids` read, with its signature decided against who | |
| 235 | − | /// owns the keys and addresses (`owners`, when signatures matter). | |
| 251 | + | /// owns the keys and addresses (`owners`, when signatures matter). The | |
| 252 | + | /// commits are read at once. | |
| 236 | 253 | pub async fn read_all<R: GitRepo>( | |
| 237 | 254 | pack: &Pack, | |
| 238 | 255 | repo: &R, | |
| ⋯ | |||
| 240 | 257 | owners: Option<&Owners>, | |
| 241 | 258 | ) -> Result<Vec<CommitFacts>> { | |
| 242 | 259 | let objects = Objects { pack, repo, reads: Cell::new(0) }; | |
| 243 | − | let mut out = Vec::new(); | |
| 244 | − | for id in ids { | |
| 260 | + | let objects = &objects; | |
| 261 | + | let read = try_join_all(ids.iter().map(|id| async move { | |
| 245 | 262 | let signature = owners.and_then(|(keys, emails)| { | |
| 246 | 263 | let (_, data) = pack.get(id)?; | |
| 247 | 264 | let committer = read_commit(data).committer_email; | |
| 248 | 265 | Some(crate::signatures::decide(data, committer.as_deref(), keys, emails)) | |
| 249 | 266 | }); | |
| 250 | − | if let Some(facts) = facts(&objects, id, signature).await? { | |
| 251 | − | out.push(facts); | |
| 252 | − | } | |
| 253 | − | } | |
| 254 | − | Ok(out) | |
| 267 | + | facts(objects, id, signature).await | |
| 268 | + | })) | |
| 269 | + | .await?; | |
| 270 | + | Ok(read.into_iter().flatten().collect()) | |
| 255 | 271 | } | |
| 256 | 272 | ||
| 257 | 273 | #[cfg(test)] | |
| ⋯ | |||
| 288 | 304 | assert_eq!(added(&pack, &second_id, 1), None, "past the limit"); | |
| 289 | 305 | assert_eq!(added(&pack, &"0".repeat(40), 10), Some(Vec::new()), "a tip the pack does not hold adds nothing"); | |
| 290 | 306 | } | |
| 307 | + | ||
| 308 | + | /// Histories by where they start; a start missing from it fails. | |
| 309 | + | struct Histories(HashMap<String, Vec<g1t_contracts::repos::Commit>>); | |
| 310 | + | ||
| 311 | + | impl GitRepo for Histories { | |
| 312 | + | async fn access(&self, _scope: crate::store::Scope) -> Result<g1t_contracts::repos::GitAccess> { | |
| 313 | + | unimplemented!() | |
| 314 | + | } | |
| 315 | + | async fn branches(&self) -> Result<Vec<g1t_contracts::repos::Branch>> { | |
| 316 | + | Ok(Vec::new()) | |
| 317 | + | } | |
| 318 | + | async fn log(&self, git_ref: &str, _limit: u32) -> Result<Vec<g1t_contracts::repos::Commit>> { | |
| 319 | + | self.0.get(git_ref).cloned().ok_or_else(|| worker::Error::RustError(format!("no history from {git_ref}"))) | |
| 320 | + | } | |
| 321 | + | async fn parents(&self, _commit_hash: &str) -> Result<Option<Vec<String>>> { | |
| 322 | + | Ok(None) | |
| 323 | + | } | |
| 324 | + | async fn read_tree(&self, _tree_hash: &str) -> Result<Option<Vec<g1t_contracts::repos::TreeEntry>>> { | |
| 325 | + | Ok(None) | |
| 326 | + | } | |
| 327 | + | async fn read_blob(&self, _blob_hash: &str) -> Result<Option<Vec<u8>>> { | |
| 328 | + | Ok(None) | |
| 329 | + | } | |
| 330 | + | async fn read_file(&self, _git_ref: &str, _path: &str) -> Result<Option<Vec<u8>>> { | |
| 331 | + | Ok(None) | |
| 332 | + | } | |
| 333 | + | async fn fork(&self, _target_key: &str) -> Result<()> { | |
| 334 | + | Ok(()) | |
| 335 | + | } | |
| 336 | + | } | |
| 337 | + | ||
| 338 | + | fn run<F: std::future::Future>(future: F) -> F::Output { | |
| 339 | + | let waker = std::task::Waker::noop(); | |
| 340 | + | match std::pin::pin!(future).as_mut().poll(&mut std::task::Context::from_waker(waker)) { | |
| 341 | + | std::task::Poll::Ready(output) => output, | |
| 342 | + | std::task::Poll::Pending => panic!("the fake store never waits"), | |
| 343 | + | } | |
| 344 | + | } | |
| 345 | + | ||
| 346 | + | fn at(hash: &str, parents: &[&str]) -> g1t_contracts::repos::Commit { | |
| 347 | + | g1t_contracts::repos::Commit { | |
| 348 | + | hash: hash.to_owned(), | |
| 349 | + | tree_hash: String::new(), | |
| 350 | + | message: String::new(), | |
| 351 | + | author: g1t_contracts::repos::Signature { name: "A".into(), email: "a@x".into() }, | |
| 352 | + | parents: parents.iter().map(|parent| (*parent).to_owned()).collect(), | |
| 353 | + | authored_at: String::new(), | |
| 354 | + | } | |
| 355 | + | } | |
| 356 | + | ||
| 357 | + | #[test] | |
| 358 | + | fn a_fast_forward_is_found_along_the_history_and_its_merges() { | |
| 359 | + | use g1t_scan::pack::write_pack; | |
| 360 | + | // The push brings one commit on top of `base`; `base` merged `side`, | |
| 361 | + | // whose history holds `old`. | |
| 362 | + | let (base, side, old) = ("b".repeat(40), "5".repeat(40), "0".repeat(40)); | |
| 363 | + | let tip = format!("tree t | |
| 364 | + | parent {base} | |
| 365 | + | author A <a@x> 1 +0000 | |
| 366 | + | committer A <a@x> 1 +0000 | |
| 367 | + | ||
| 368 | + | tip | |
| 369 | + | ").into_bytes(); | |
| 370 | + | let tip_id = g1t_scan::pack::object_id(ObjectKind::Commit, &tip); | |
| 371 | + | let pack = Pack::parse(&write_pack(&[(ObjectKind::Commit, tip)])).unwrap(); | |
| 372 | + | let mut histories = HashMap::new(); | |
| 373 | + | histories.insert(base.clone(), vec![at(&base, &["1".repeat(40).as_str(), side.as_str()])]); | |
| 374 | + | histories.insert(side.clone(), vec![at(&side, &[]), at(&old, &[])]); | |
| 375 | + | let repo = Histories(histories); | |
| 376 | + | assert!(run(contains(&pack, &repo, &tip_id, &old, 100)).unwrap()); | |
| 377 | + | assert!(run(contains(&pack, &repo, &tip_id, &side, 100)).unwrap(), "a second parent itself"); | |
| 378 | + | assert!(!run(contains(&pack, &repo, &tip_id, &"9".repeat(40), 100)).unwrap()); | |
| 379 | + | // A history that cannot be read does not hide one found elsewhere, | |
| 380 | + | // and is an error only when nothing was found. | |
| 381 | + | let mut histories = repo.0; | |
| 382 | + | histories.remove(&side); | |
| 383 | + | let repo = Histories(histories); | |
| 384 | + | assert!(run(contains(&pack, &repo, &tip_id, &side, 100)).unwrap()); | |
| 385 | + | assert!(run(contains(&pack, &repo, &tip_id, &old, 100)).is_err()); | |
| 386 | + | } | |
| 291 | 387 | } | |
| 19 | 19 | use g1t_contracts::{FailureCode, Outcome, User}; | |
| 20 | 20 | use g1t_rules::push::{RefChange, judge}; | |
| 21 | 21 | use g1t_rules::{ActorFacts, Judged, Who, content, outcome, report}; | |
| 22 | − | use g1t_scan::pack::{Pack, pack_start}; | |
| 22 | + | use g1t_scan::pack::Pack; | |
| 23 | 23 | use worker::{Response, Result}; | |
| 24 | 24 | ||
| 25 | 25 | use crate::registry::store_key; | |
| ⋯ | |||
| 30 | 30 | /// Where people read the rules of a branch. | |
| 31 | 31 | const SITE: &str = "https://g1t.sh"; | |
| 32 | 32 | ||
| 33 | + | /// What the rules ask of a push, before its pack is read. | |
| 34 | + | pub(crate) enum PushRules { | |
| 35 | + | /// Nothing to judge: a pull request's working copy, a push that changes | |
| 36 | + | /// no branch or tag, or one no ruleset holds for. | |
| 37 | + | Nothing, | |
| 38 | + | /// Answered without the pack: by the old protection, where rulesets | |
| 39 | + | /// cannot be read, or declined because they could not be read just now. | |
| 40 | + | Answered(Result<Option<Response>>), | |
| 41 | + | /// The rulesets to judge it by, and the branches and tags it changes. | |
| 42 | + | Judge { rules: RefRules, updates: Vec<(String, Option<String>, Option<String>)> }, | |
| 43 | + | /// Judged, once its pack was read (push_checks.rs). | |
| 44 | + | Judged(PushJudged), | |
| 45 | + | } | |
| 46 | + | ||
| 47 | + | /// How the rules judged a push, and the response declining it, if they did. | |
| 48 | + | pub(crate) struct PushJudged { | |
| 49 | + | facts: ActorFacts, | |
| 50 | + | judged: Vec<Judged>, | |
| 51 | + | sha: Option<String>, | |
| 52 | + | pub(crate) refusal: Option<Response>, | |
| 53 | + | } | |
| 54 | + | ||
| 33 | 55 | /// What the rules said about a change. | |
| 34 | 56 | pub(crate) enum Ruled { | |
| 35 | 57 | /// No ruleset holds, or none refuses it. | |
| ⋯ | |||
| 163 | 185 | Ok(Ruled::Refused { message }) | |
| 164 | 186 | } | |
| 165 | 187 | ||
| 166 | − | /// Rules for a push: the response declining it, or `None` to let it | |
| 167 | − | /// on. `body` is as much of the push as was read; `whole` says whether | |
| 168 | − | /// that is all of it, so that its commits can be read. | |
| 169 | − | pub(crate) async fn check_push(&self, repo: &Repo, pusher: Option<&User>, body: &[u8], whole: bool) -> Result<Option<Response>> { | |
| 170 | − | // Workflow files need their own scope from a token, on a pull | |
| 171 | − | // request's working copy too (workflow_gate.rs). | |
| 172 | − | if let Some(response) = self.workflow_gate(repo, pusher, body, whole).await? { | |
| 173 | − | return Ok(Some(response)); | |
| 174 | − | } | |
| 188 | + | /// What the rules ask of a push, found before its pack is read: the | |
| 189 | + | /// rulesets that hold for the branches and tags it changes. | |
| 190 | + | pub(crate) async fn push_rules(&self, repo: &Repo, pusher: Option<&User>, body: &[u8]) -> PushRules { | |
| 175 | 191 | if repo.fork_of.is_some() { | |
| 176 | − | return Ok(None); | |
| 192 | + | return PushRules::Nothing; | |
| 177 | 193 | } | |
| 178 | 194 | let updates = crate::git_http::ref_updates(body); | |
| 179 | 195 | if updates.is_empty() { | |
| 180 | − | return Ok(None); | |
| 196 | + | return PushRules::Nothing; | |
| 181 | 197 | } | |
| 182 | 198 | let refs: Vec<String> = updates.iter().map(|(name, _, _)| name.clone()).collect(); | |
| 183 | − | let rules = match self.ref_rules(repo, pusher, refs).await { | |
| 184 | − | Ok(Some(rules)) => rules, | |
| 199 | + | match self.ref_rules(repo, pusher, refs).await { | |
| 200 | + | Ok(Some(rules)) if rules.rulesets.is_empty() => PushRules::Nothing, | |
| 201 | + | Ok(Some(rules)) => PushRules::Judge { rules, updates }, | |
| 185 | 202 | // Without the work service, the old protection holds. | |
| 186 | 203 | Ok(None) => { | |
| 187 | 204 | let protected = repo.protected.then(|| repo.default_branch.clone()); | |
| 188 | − | return protected | |
| 189 | − | .and_then(|branch| crate::git_http::refusal(body, &branch)) | |
| 190 | − | .map(crate::git_http::report_response) | |
| 191 | − | .transpose(); | |
| 205 | + | PushRules::Answered( | |
| 206 | + | protected | |
| 207 | + | .and_then(|branch| crate::git_http::refusal(body, &branch)) | |
| 208 | + | .map(crate::git_http::report_response) | |
| 209 | + | .transpose(), | |
| 210 | + | ) | |
| 192 | 211 | } | |
| 193 | 212 | Err(error) => { | |
| 194 | 213 | worker::console_error!("ref_rules failed during a push: {error}"); | |
| 195 | − | return Ok(Some(crate::git_http::declined( | |
| 196 | − | body, | |
| 197 | − | "rules could not be checked", | |
| 198 | − | &["The rules for this repository could not be checked just now. Push again in a moment.".to_owned()], | |
| 199 | − | )?)); | |
| 214 | + | PushRules::Answered( | |
| 215 | + | crate::git_http::declined( | |
| 216 | + | body, | |
| 217 | + | "rules could not be checked", | |
| 218 | + | &["The rules for this repository could not be checked just now. Push again in a moment.".to_owned()], | |
| 219 | + | ) | |
| 220 | + | .map(Some), | |
| 221 | + | ) | |
| 200 | 222 | } | |
| 201 | − | }; | |
| 202 | − | if rules.rulesets.is_empty() { | |
| 203 | − | return Ok(None); | |
| 204 | 223 | } | |
| 224 | + | } | |
| 225 | + | ||
| 226 | + | /// Judges a push by `rules` (see [`Repos::push_rules`]). `pack` is the | |
| 227 | + | /// push's, with its bases supplied, when it was read whole; its | |
| 228 | + | /// commits are read only when a rule needs them. | |
| 229 | + | #[allow(clippy::too_many_arguments)] | |
| 230 | + | pub(crate) async fn judge_push<R: GitRepo>( | |
| 231 | + | &self, | |
| 232 | + | repo: &Repo, | |
| 233 | + | pusher: Option<&User>, | |
| 234 | + | body: &[u8], | |
| 235 | + | rules: &RefRules, | |
| 236 | + | updates: Vec<(String, Option<String>, Option<String>)>, | |
| 237 | + | pack: Option<&Pack>, | |
| 238 | + | git: &R, | |
| 239 | + | ) -> Result<PushJudged> { | |
| 205 | 240 | let facts = who(pusher, repo); | |
| 206 | 241 | let content = needs_commits(&rules.rulesets, facts.kind); | |
| 207 | 242 | let ancestry = needs_ancestry(&rules.rulesets); | |
| 208 | − | let git = self.store.open(&store_key(repo)).await?; | |
| 209 | − | // The pack, read when a rule needs what it holds. | |
| 210 | − | let pack = if whole && (content || ancestry) { | |
| 211 | − | match pack_start(body).map(|start| Pack::parse(&body[start..])) { | |
| 212 | − | Some(Ok(mut pack)) => { | |
| 213 | − | crate::secret_scan::supply_bases(&mut pack, &git).await?; | |
| 214 | − | Some(pack) | |
| 215 | − | } | |
| 216 | − | Some(Err(problem)) => { | |
| 217 | − | worker::console_error!("a push's pack could not be read for rules: {problem}"); | |
| 218 | − | None | |
| 219 | − | } | |
| 220 | − | // Nothing but deletions, or pointing refs at commits the | |
| 221 | − | // repository has: an empty pack. | |
| 222 | − | None => Pack::parse(crate::land::EMPTY_PACK).ok(), | |
| 223 | − | } | |
| 224 | − | } else { | |
| 225 | − | None | |
| 226 | − | }; | |
| 243 | + | // The pack, used when a rule needs what it holds. | |
| 244 | + | let pack = pack.filter(|_| content || ancestry); | |
| 227 | 245 | let signatures = rules | |
| 228 | 246 | .rulesets | |
| 229 | 247 | .iter() | |
| 230 | 248 | .any(|ruleset| ruleset.rules.iter().any(|entry| matches!(entry.rule, Rule::RequiredSignatures(_)))); | |
| 231 | − | let mut changes = Vec::new(); | |
| 232 | − | for (git_ref, old, new) in updates { | |
| 249 | + | // Each ref's facts are read at once. | |
| 250 | + | let changes = futures_util::future::try_join_all(updates.into_iter().map(|(git_ref, old, new)| async move { | |
| 233 | 251 | let mut change = RefChange { git_ref, old: old.clone(), new: new.clone(), fast_forward: None, commits: Vec::new(), complete: false }; | |
| 234 | − | if let (Some(pack), Some(new)) = (&pack, &new) { | |
| 235 | − | if let Some(old) = &old | |
| 236 | − | && ancestry | |
| 237 | − | { | |
| 238 | − | change.fast_forward = Some(rule_facts::contains(pack, &git, new, old, MAX_ANCESTRY).await?); | |
| 239 | − | } | |
| 240 | − | if content { | |
| 241 | − | match rule_facts::added(pack, new, MAX_COMMITS) { | |
| 242 | − | Some(ids) => { | |
| 243 | − | let owners = if signatures { | |
| 244 | − | let (fingerprints, emails) = rule_facts::signing_facts(pack, &ids); | |
| 245 | − | Some(self.signing_owners(&fingerprints, &emails).await) | |
| 246 | − | } else { | |
| 247 | − | None | |
| 248 | − | }; | |
| 249 | − | change.commits = rule_facts::read_all(pack, &git, &ids, owners.as_ref()).await?; | |
| 250 | − | change.complete = change.commits.iter().all(|commit| commit.files_complete); | |
| 251 | − | } | |
| 252 | − | None => change.complete = false, | |
| 252 | + | if let (Some(pack), Some(new)) = (pack, &new) { | |
| 253 | + | let fast_forward = async { | |
| 254 | + | match &old { | |
| 255 | + | Some(old) if ancestry => Ok::<_, worker::Error>(Some(rule_facts::contains(pack, git, new, old, MAX_ANCESTRY).await?)), | |
| 256 | + | _ => Ok(None), | |
| 257 | + | } | |
| 258 | + | }; | |
| 259 | + | let commits = async { | |
| 260 | + | if !content { | |
| 261 | + | return Ok::<_, worker::Error>((Vec::new(), true)); | |
| 253 | 262 | } | |
| 254 | − | } else { | |
| 255 | − | change.complete = true; | |
| 256 | − | } | |
| 263 | + | let Some(ids) = rule_facts::added(pack, new, MAX_COMMITS) else { | |
| 264 | + | return Ok((Vec::new(), false)); | |
| 265 | + | }; | |
| 266 | + | let owners = if signatures { | |
| 267 | + | let (fingerprints, emails) = rule_facts::signing_facts(pack, &ids); | |
| 268 | + | Some(self.signing_owners(&fingerprints, &emails).await) | |
| 269 | + | } else { | |
| 270 | + | None | |
| 271 | + | }; | |
| 272 | + | let commits = rule_facts::read_all(pack, git, &ids, owners.as_ref()).await?; | |
| 273 | + | let complete = commits.iter().all(|commit| commit.files_complete); | |
| 274 | + | Ok((commits, complete)) | |
| 275 | + | }; | |
| 276 | + | let (fast_forward, commits) = futures_util::future::try_join(fast_forward, commits).await?; | |
| 277 | + | change.fast_forward = fast_forward; | |
| 278 | + | (change.commits, change.complete) = commits; | |
| 257 | 279 | } else if new.is_none() || !content { | |
| 258 | 280 | change.complete = true; | |
| 259 | 281 | } | |
| 260 | − | changes.push(change); | |
| 261 | − | } | |
| 282 | + | Ok::<_, worker::Error>(change) | |
| 283 | + | })) | |
| 284 | + | .await?; | |
| 262 | 285 | let judged: Vec<Judged> = changes | |
| 263 | 286 | .iter() | |
| 264 | 287 | .flat_map(|change| judge(&rules.rulesets, &rules.default_branch, facts.kind, change)) | |
| 265 | 288 | .collect(); | |
| 266 | 289 | let sha = changes.iter().find_map(|change| change.new.clone()); | |
| 267 | − | self.record_judged(repo, &judged, Action::Push, &facts, sha.as_deref()).await; | |
| 268 | − | if !outcome::refused(&judged) { | |
| 269 | − | return Ok(None); | |
| 290 | + | let refusal = if outcome::refused(&judged) { | |
| 291 | + | let refused_ref = judged.iter().find(|one| one.blocks()).map(|one| one.git_ref.clone()).unwrap_or_default(); | |
| 292 | + | let lines = report::remote_lines(&refused_ref, &judged, &rules_url(repo, &refused_ref)); | |
| 293 | + | Some(crate::git_http::declined(body, &report::ng_reason(&judged), &lines)?) | |
| 294 | + | } else { | |
| 295 | + | None | |
| 296 | + | }; | |
| 297 | + | Ok(PushJudged { facts, judged, sha, refusal }) | |
| 298 | + | } | |
| 299 | + | ||
| 300 | + | /// Records how the rules judged a push: before the answer when they | |
| 301 | + | /// refused it, after it otherwise, since a push let through does not | |
| 302 | + | /// wait on the record. | |
| 303 | + | pub(crate) async fn record_push_judged(&self, repo: &Repo, judged: &PushJudged) { | |
| 304 | + | if judged.refusal.is_some() { | |
| 305 | + | self.record_judged(repo, &judged.judged, Action::Push, &judged.facts, judged.sha.as_deref()).await; | |
| 306 | + | return; | |
| 307 | + | } | |
| 308 | + | let Some(work) = self.work.clone() else { return }; | |
| 309 | + | if judged.judged.is_empty() { | |
| 310 | + | return; | |
| 270 | 311 | } | |
| 271 | − | let refused_ref = judged.iter().find(|one| one.blocks()).map(|one| one.git_ref.clone()).unwrap_or_default(); | |
| 272 | − | let lines = report::remote_lines(&refused_ref, &judged, &rules_url(repo, &refused_ref)); | |
| 273 | − | Ok(Some(crate::git_http::declined(body, &report::ng_reason(&judged), &lines)?)) | |
| 312 | + | let evaluations = g1t_rules::evaluations(&judged.judged, &repo.id, &repo.namespace, Action::Push, &judged.facts, None, judged.sha.as_deref()); | |
| 313 | + | self.deferred.spawn(async move { | |
| 314 | + | let recorded: Result<u32> = g1t_kit::call(&work, "record_evaluations", &RecordEvaluationsArgs { evaluations }).await; | |
| 315 | + | if let Err(error) = recorded { | |
| 316 | + | worker::console_error!("rule evaluations not recorded: {error}"); | |
| 317 | + | } | |
| 318 | + | }); | |
| 274 | 319 | } | |
| 275 | 320 | ||
| 276 | 321 | /// Services only: the commits a pull request would land, read as rules | |
| 28 | 28 | use g1t_contracts::security_suite::{CheckSecretArgs, MatchPatternArgs, PatternMatch, PatternMatches, PatternSpec, PatternsForArgs, SecretValidity}; | |
| 29 | 29 | use g1t_scan::custom::{self, Compiled}; | |
| 30 | 30 | use g1t_scan::lockfiles::Lockfile; | |
| 31 | − | use g1t_scan::pack::{ObjectKind, Pack, TreeItem, encode_tree, pack_start}; | |
| 31 | + | use g1t_scan::pack::{ObjectKind, Pack, TreeItem, encode_tree}; | |
| 32 | 32 | use g1t_scan::protection::{self, Blocked}; | |
| 33 | − | use worker::{Response, Result}; | |
| 33 | + | use worker::Result; | |
| 34 | 34 | ||
| 35 | 35 | use crate::registry::store_key; | |
| 36 | 36 | use crate::store::{GitRepo, GitStore}; | |
| ⋯ | |||
| 224 | 224 | .filter(|change| !g1t_scan::secrets::skipped_path(&change.path)) | |
| 225 | 225 | .collect(); | |
| 226 | 226 | for batch in changes.chunks(READS_AT_ONCE) { | |
| 227 | + | // A change's old and new contents at once: the new is nearly always | |
| 228 | + | // in the pack, the old in the repository. | |
| 227 | 229 | let read = try_join_all(batch.iter().map(|change| async move { | |
| 228 | − | let new = objects.blob(&change.new).await?; | |
| 229 | − | let old = match (&change.old, &new) { | |
| 230 | − | (Some(old), Some(_)) => objects.blob(old).await?, | |
| 231 | − | _ => None, | |
| 230 | + | let old = async { | |
| 231 | + | match &change.old { | |
| 232 | + | Some(old) => objects.blob(old).await, | |
| 233 | + | None => Ok(None), | |
| 234 | + | } | |
| 232 | 235 | }; | |
| 236 | + | let (new, old) = futures_util::future::try_join(objects.blob(&change.new), old).await?; | |
| 233 | 237 | Ok::<_, worker::Error>((new, old)) | |
| 234 | 238 | })) | |
| 235 | 239 | .await?; | |
| ⋯ | |||
| 256 | 260 | Ok(found) | |
| 257 | 261 | } | |
| 258 | 262 | ||
| 259 | − | /// Fetches what a thin pack's deltas are based on from the repository. | |
| 263 | + | /// Fetches what a thin pack's deltas are based on from the repository, | |
| 264 | + | /// all at once. A base the pack's own trees name is read as what they say | |
| 265 | + | /// it is; any other is asked for as a blob and as a tree together, and | |
| 266 | + | /// whichever it is answers. | |
| 260 | 267 | pub(crate) async fn supply_bases<R: GitRepo>(pack: &mut Pack, repo: &R) -> Result<()> { | |
| 261 | 268 | for _ in 0..3 { | |
| 262 | 269 | let missing = pack.missing_bases(); | |
| 263 | 270 | if missing.is_empty() { | |
| 264 | 271 | return Ok(()); | |
| 265 | 272 | } | |
| 266 | − | let found = try_join_all(missing.iter().take(MAX_BASES).map(|id| async move { | |
| 267 | − | // A base is nearly always a blob; failing that, a tree. | |
| 268 | − | if let Ok(Some(bytes)) = repo.read_blob(id).await { | |
| 269 | − | return Ok::<_, worker::Error>(Some((ObjectKind::Blob, bytes))); | |
| 270 | − | } | |
| 271 | − | Ok(repo.read_tree(id).await.ok().flatten().map(|entries| { | |
| 273 | + | let named = pack.named_kinds(); | |
| 274 | + | let blob = async |id: &str| repo.read_blob(id).await.ok().flatten().map(|bytes| (ObjectKind::Blob, bytes)); | |
| 275 | + | let tree = async |id: &str| { | |
| 276 | + | repo.read_tree(id).await.ok().flatten().map(|entries| { | |
| 272 | 277 | let items: Vec<TreeItem> = entries | |
| 273 | 278 | .into_iter() | |
| 274 | 279 | .map(|entry| TreeItem { mode: mode(entry.kind).to_owned(), name: entry.name, id: entry.hash }) | |
| 275 | 280 | .collect(); | |
| 276 | 281 | (ObjectKind::Tree, encode_tree(&items)) | |
| 277 | − | })) | |
| 282 | + | }) | |
| 283 | + | }; | |
| 284 | + | let found = futures_util::future::join_all(missing.iter().take(MAX_BASES).map(|id| async { | |
| 285 | + | match named.get(id.as_str()) { | |
| 286 | + | Some(ObjectKind::Tree) => tree(id).await, | |
| 287 | + | Some(_) => blob(id).await, | |
| 288 | + | None => { | |
| 289 | + | let (as_blob, as_tree) = futures_util::future::join(blob(id), tree(id)).await; | |
| 290 | + | as_blob.or(as_tree) | |
| 291 | + | } | |
| 292 | + | } | |
| 278 | 293 | })) | |
| 279 | − | .await?; | |
| 294 | + | .await; | |
| 280 | 295 | let mut progress = false; | |
| 281 | 296 | for (id, object) in missing.iter().zip(found) { | |
| 282 | 297 | if let Some((kind, data)) = object { | |
| ⋯ | |||
| 295 | 310 | /// large to read is an error ([`unscannable`]): it is declined, never let | |
| 296 | 311 | /// through unread. A pack that cannot be read for another reason is let | |
| 297 | 312 | /// through, and said so in the logs; the store will judge it. | |
| 313 | + | #[cfg(test)] | |
| 298 | 314 | pub async fn scan_push<R: GitRepo>(repo: &R, body: &[u8], patterns: &[Compiled]) -> Result<Vec<NewSecret>> { | |
| 299 | 315 | if body.len() > MAX_SCANNED_PUSH { | |
| 300 | 316 | return Err(worker::Error::RustError(format!("{UNSCANNABLE} {} bytes", body.len()))); | |
| 301 | 317 | } | |
| 302 | − | let Some(start) = pack_start(body) else { | |
| 303 | − | return Ok(Vec::new()); | |
| 304 | − | }; | |
| 305 | − | let mut pack = match Pack::parse(&body[start..]) { | |
| 318 | + | let mut pack = crate::push_checks::read_pack(body); | |
| 319 | + | if let Ok(pack) = &mut pack { | |
| 320 | + | supply_bases(pack, repo).await?; | |
| 321 | + | } | |
| 322 | + | scan_pack(repo, body.len(), &pack, patterns).await | |
| 323 | + | } | |
| 324 | + | ||
| 325 | + | /// [`scan_push`] for a push of `size` bytes whose pack was read already, | |
| 326 | + | /// with its bases supplied (push_checks.rs). | |
| 327 | + | pub(crate) async fn scan_pack<R: GitRepo>(repo: &R, size: usize, pack: &std::result::Result<Pack, String>, patterns: &[Compiled]) -> Result<Vec<NewSecret>> { | |
| 328 | + | if size > MAX_SCANNED_PUSH { | |
| 329 | + | return Err(worker::Error::RustError(format!("{UNSCANNABLE} {size} bytes"))); | |
| 330 | + | } | |
| 331 | + | let pack = match pack { | |
| 306 | 332 | Ok(pack) => pack, | |
| 307 | 333 | Err(problem) if problem.contains("too large") => { | |
| 308 | 334 | return Err(worker::Error::RustError(format!("{UNSCANNABLE} {problem}"))); | |
| ⋯ | |||
| 312 | 338 | return Ok(Vec::new()); | |
| 313 | 339 | } | |
| 314 | 340 | }; | |
| 315 | − | supply_bases(&mut pack, repo).await?; | |
| 316 | 341 | if pack.unresolved() > 0 { | |
| 317 | 342 | worker::console_error!("{} objects of a push could not be resolved for scanning", pack.unresolved()); | |
| 318 | 343 | } | |
| 319 | − | let objects = Objects { pack: &pack, repo, reads: Cell::new(0) }; | |
| 320 | − | let commits: Vec<String> = pack.commits().iter().take(MAX_PUSH_COMMITS).cloned().collect(); | |
| 344 | + | let objects = Objects { pack, repo, reads: Cell::new(0) }; | |
| 345 | + | let objects = &objects; | |
| 346 | + | let commits: Vec<(String, g1t_scan::pack::CommitInfo)> = pack | |
| 347 | + | .commits() | |
| 348 | + | .iter() | |
| 349 | + | .take(MAX_PUSH_COMMITS) | |
| 350 | + | .filter_map(|id| Some((id.clone(), pack.commit(id)?))) | |
| 351 | + | .collect(); | |
| 352 | + | // The trees of the parents the pack does not hold, all read up front: | |
| 353 | + | // usually the one commit the push builds on. | |
| 354 | + | let mut outside: Vec<&String> = commits | |
| 355 | + | .iter() | |
| 356 | + | .filter_map(|(_, commit)| commit.parents.first()) | |
| 357 | + | .filter(|parent| !pack.contains(parent)) | |
| 358 | + | .collect(); | |
| 359 | + | outside.sort(); | |
| 360 | + | outside.dedup(); | |
| 361 | + | let outside_trees: std::collections::HashMap<&String, Option<String>> = outside | |
| 362 | + | .iter() | |
| 363 | + | .copied() | |
| 364 | + | .zip(try_join_all(outside.iter().map(|parent| objects.commit_tree(parent))).await?) | |
| 365 | + | .collect(); | |
| 366 | + | let outside_trees = &outside_trees; | |
| 367 | + | // What each commit changes, all read at once. | |
| 368 | + | let changed = try_join_all(commits.iter().map(|(_, commit)| async move { | |
| 369 | + | let old_tree = match commit.parents.first() { | |
| 370 | + | Some(parent) => match outside_trees.get(parent) { | |
| 371 | + | Some(tree) => tree.clone(), | |
| 372 | + | None => objects.commit_tree(parent).await?, | |
| 373 | + | }, | |
| 374 | + | None => None, | |
| 375 | + | }; | |
| 376 | + | changed_files(objects, old_tree, commit.tree.clone()).await | |
| 377 | + | })) | |
| 378 | + | .await?; | |
| 321 | 379 | let mut found = Vec::new(); | |
| 322 | 380 | let mut seen_blobs = HashSet::new(); | |
| 323 | 381 | let mut seen_secrets = HashSet::new(); | |
| 324 | − | for id in commits { | |
| 325 | − | let Some(commit) = pack.commit(&id) else { continue }; | |
| 326 | − | let old_tree = match commit.parents.first() { | |
| 327 | − | Some(parent) => objects.commit_tree(parent).await?, | |
| 328 | − | None => None, | |
| 329 | − | }; | |
| 382 | + | for ((id, _), changes) in commits.iter().zip(changed) { | |
| 330 | 383 | // Only content the push brings is new; a blob the repository has | |
| 331 | 384 | // was looked at when it arrived. | |
| 332 | − | let changes: Vec<Change> = changed_files(&objects, old_tree, commit.tree) | |
| 333 | − | .await? | |
| 385 | + | let changes: Vec<Change> = changes | |
| 334 | 386 | .into_iter() | |
| 335 | 387 | .filter(|change| pack.contains(&change.new) && seen_blobs.insert((change.path.clone(), change.new.clone()))) | |
| 336 | 388 | .collect(); | |
| 337 | − | for secret in scan_changes(&objects, &id, changes, patterns).await? { | |
| 389 | + | for secret in scan_changes(objects, id, changes, patterns).await? { | |
| 338 | 390 | if seen_secrets.insert(secret.fingerprint.clone()) { | |
| 339 | 391 | found.push(secret); | |
| 340 | 392 | } | |
| ⋯ | |||
| 347 | 399 | /// addresses while they keep it private: its id and the address. Only the | |
| 348 | 400 | /// commits the push adds are read; anyone else's address is no concern | |
| 349 | 401 | /// here. A pack that cannot be read is let through. | |
| 402 | + | #[cfg(test)] | |
| 350 | 403 | pub fn exposed_address(body: &[u8], guard: &PushEmailGuard) -> Option<(String, String)> { | |
| 351 | 404 | if body.len() > MAX_SCANNED_PUSH { | |
| 352 | 405 | return None; | |
| 353 | 406 | } | |
| 354 | − | let pack = Pack::parse(&body[pack_start(body)?..]).ok()?; | |
| 407 | + | exposed_in(&Pack::parse(&body[g1t_scan::pack::pack_start(body)?..]).ok()?, guard) | |
| 408 | + | } | |
| 409 | + | ||
| 410 | + | /// [`exposed_address`] for a pack read already. | |
| 411 | + | pub(crate) fn exposed_in(pack: &Pack, guard: &PushEmailGuard) -> Option<(String, String)> { | |
| 355 | 412 | pack.commits().iter().find_map(|id| { | |
| 356 | 413 | let commit = pack.commit(id)?; | |
| 357 | 414 | [commit.author_email, commit.committer_email] | |
| ⋯ | |||
| 379 | 436 | /// What a push by `pusher` must not publish: their own addresses, when | |
| 380 | 437 | /// they keep them private and block such pushes. An agent's push is | |
| 381 | 438 | /// its person's. `None` when nothing is guarded, or identity cannot say. | |
| 382 | − | async fn push_email_guard(&self, pusher: Option<&User>) -> Option<PushEmailGuard> { | |
| 439 | + | pub(crate) async fn push_email_guard(&self, pusher: Option<&User>) -> Option<PushEmailGuard> { | |
| 383 | 440 | let pusher = pusher?; | |
| 384 | 441 | let person = pusher.acting.as_ref().map_or(pusher.id.clone(), |acting| acting.on_behalf_of.id.clone()); | |
| 385 | 442 | let identity = self.identity.as_ref()?; | |
| ⋯ | |||
| 389 | 446 | worker::console_error!("push_email_guard failed: {error}"); | |
| 390 | 447 | None | |
| 391 | 448 | }) | |
| 392 | − | } | |
| 393 | − | ||
| 394 | − | /// Push protection: the response refusing a push that adds secrets | |
| 395 | − | /// nobody has allowed, or that would publish the pusher's private | |
| 396 | − | /// address, or `None` to let it through. | |
| 397 | − | /// `repo` is the repository pushed to, as the request read it. | |
| 398 | − | pub(crate) async fn protect(&self, repo: &Repo, pusher: Option<&User>, body: &[u8]) -> Result<Option<Response>> { | |
| 399 | − | // A pull request's findings belong to the repository it was made from. | |
| 400 | − | let owner = match &repo.fork_of { | |
| 401 | − | Some(id) => self.registry.by_id(id).await?.unwrap_or(repo.clone()), | |
| 402 | − | None => repo.clone(), | |
| 403 | − | }; | |
| 404 | − | // Asking identity about the pusher's address and scanning the push | |
| 405 | − | // do not depend on each other, so they happen at once. | |
| 406 | − | let scan = async { | |
| 407 | − | let patterns = compiled(&self.patterns_for(&owner).await); | |
| 408 | − | let git = self.store.open(&store_key(repo)).await?; | |
| 409 | − | scan_push(&git, body, &patterns).await | |
| 410 | − | }; | |
| 411 | − | let (guard, found) = futures_util::future::join(self.push_email_guard(pusher), scan).await; | |
| 412 | − | if let Some(guard) = guard | |
| 413 | − | && let Some((commit, email)) = exposed_address(body, &guard) | |
| 414 | − | { | |
| 415 | − | return Ok(Some(crate::git_http::declined( | |
| 416 | − | body, | |
| 417 | − | "push would publish a private email", | |
| 418 | − | &exposed_message(&commit, &email, &guard.noreply), | |
| 419 | − | )?)); | |
| 420 | − | } | |
| 421 | − | let found = match found { | |
| 422 | − | Err(error) if unscannable(&error) => { | |
| 423 | − | let (reason, messages) = crate::git_http::size_refusal(&crate::git_http::SizeViolation::Unscannable { | |
| 424 | − | size: body.len() as u64, | |
| 425 | − | cap: MAX_SCANNED_PUSH, | |
| 426 | − | }); | |
| 427 | − | return Ok(Some(crate::git_http::declined(body, &reason, &messages)?)); | |
| 428 | − | } | |
| 429 | − | found => found?, | |
| 430 | − | }; | |
| 431 | − | if found.is_empty() { | |
| 432 | − | return Ok(None); | |
| 433 | − | } | |
| 434 | − | let blocked = self.blocked(&owner, pusher, found).await; | |
| 435 | − | if blocked.is_empty() { | |
| 436 | − | return Ok(None); | |
| 437 | − | } | |
| 438 | − | Ok(Some(crate::git_http::declined( | |
| 439 | − | body, | |
| 440 | − | &protection::reason(&blocked), | |
| 441 | − | &protection::explain(&blocked), | |
| 442 | − | )?)) | |
| 443 | 449 | } | |
| 444 | 450 | ||
| 445 | 451 | /// The custom patterns the security service says `repo` is scanned | |
| ⋯ | |||
| 724 | 730 | ||
| 725 | 731 | #[cfg(test)] | |
| 726 | 732 | mod tests { | |
| 733 | + | use std::cell::Cell; | |
| 727 | 734 | use std::collections::HashMap; | |
| 728 | 735 | use std::future::Future; | |
| 729 | 736 | use std::pin::pin; | |
| ⋯ | |||
| 749 | 756 | blobs: HashMap<String, Vec<u8>>, | |
| 750 | 757 | trees: HashMap<String, Vec<TreeEntry>>, | |
| 751 | 758 | commits: HashMap<String, Commit>, | |
| 759 | + | /// How often each kind of read was asked for. | |
| 760 | + | blob_reads: Cell<u32>, | |
| 761 | + | tree_reads: Cell<u32>, | |
| 762 | + | log_reads: Cell<u32>, | |
| 752 | 763 | } | |
| 753 | 764 | ||
| 754 | 765 | impl GitRepo for FakeRepo { | |
| ⋯ | |||
| 759 | 770 | Ok(Vec::new()) | |
| 760 | 771 | } | |
| 761 | 772 | async fn log(&self, git_ref: &str, _limit: u32) -> Result<Vec<Commit>> { | |
| 773 | + | self.log_reads.set(self.log_reads.get() + 1); | |
| 762 | 774 | Ok(self.commits.get(git_ref).cloned().into_iter().collect()) | |
| 763 | 775 | } | |
| 764 | 776 | async fn parents(&self, commit_hash: &str) -> Result<Option<Vec<String>>> { | |
| 765 | 777 | Ok(self.commits.get(commit_hash).map(|commit| commit.parents.clone())) | |
| 766 | 778 | } | |
| 767 | 779 | async fn read_tree(&self, tree_hash: &str) -> Result<Option<Vec<TreeEntry>>> { | |
| 780 | + | self.tree_reads.set(self.tree_reads.get() + 1); | |
| 768 | 781 | Ok(self.trees.get(tree_hash).cloned()) | |
| 769 | 782 | } | |
| 770 | 783 | async fn read_blob(&self, blob_hash: &str) -> Result<Option<Vec<u8>>> { | |
| 784 | + | self.blob_reads.set(self.blob_reads.get() + 1); | |
| 771 | 785 | Ok(self.blobs.get(blob_hash).cloned()) | |
| 772 | 786 | } | |
| 773 | 787 | async fn read_file(&self, _git_ref: &str, _path: &str) -> Result<Option<Vec<u8>>> { | |
| ⋯ | |||
| 920 | 934 | } | |
| 921 | 935 | ||
| 922 | 936 | #[test] | |
| 937 | + | fn a_base_is_asked_for_as_what_the_pack_names_it_or_both_ways_at_once() { | |
| 938 | + | // A file's old version, which no tree in the pack names: asked for | |
| 939 | + | // as a blob and as a tree together, and found as a blob. | |
| 940 | + | let old = b"one | |
| 941 | + | ".to_vec(); | |
| 942 | + | let old_id = object_id(ObjectKind::Blob, &old); | |
| 943 | + | let mut repo = FakeRepo::default(); | |
| 944 | + | repo.blobs.insert(old_id.clone(), old.clone()); | |
| 945 | + | let mut delta = vec![old.len() as u8, (old.len() + 4) as u8, 0x80 | 0x10, old.len() as u8, 4]; | |
| 946 | + | delta.extend_from_slice(b"two | |
| 947 | + | "); | |
| 948 | + | let body = push(&[Entry::Delta(old_id.clone(), delta)]); | |
| 949 | + | let mut pack = crate::push_checks::read_pack(&body).unwrap(); | |
| 950 | + | run(supply_bases(&mut pack, &repo)).unwrap(); | |
| 951 | + | assert_eq!(pack.unresolved(), 0); | |
| 952 | + | assert_eq!((repo.blob_reads.get(), repo.tree_reads.get()), (1, 1)); | |
| 953 | + | ||
| 954 | + | // A directory the pack's root tree names as a tree: read as one. | |
| 955 | + | let file = object_id(ObjectKind::Blob, b"x"); | |
| 956 | + | let listed = [TreeItem { mode: "100644".into(), name: "a".into(), id: file.clone() }]; | |
| 957 | + | let sub = encode_tree(&listed); | |
| 958 | + | let sub_id = object_id(ObjectKind::Tree, &sub); | |
| 959 | + | let mut repo = FakeRepo::default(); | |
| 960 | + | repo.trees.insert(sub_id.clone(), vec![TreeEntry { name: "a".into(), hash: file, kind: EntryKind::Blob }]); | |
| 961 | + | let root = encode_tree(&[TreeItem { mode: "40000".into(), name: "src".into(), id: sub_id.clone() }]); | |
| 962 | + | let delta = vec![sub.len() as u8, sub.len() as u8, 0x80 | 0x10, sub.len() as u8]; | |
| 963 | + | let body = push(&[Entry::Whole(ObjectKind::Tree, root), Entry::Delta(sub_id, delta)]); | |
| 964 | + | let mut pack = crate::push_checks::read_pack(&body).unwrap(); | |
| 965 | + | run(supply_bases(&mut pack, &repo)).unwrap(); | |
| 966 | + | assert_eq!(pack.unresolved(), 0); | |
| 967 | + | assert_eq!((repo.blob_reads.get(), repo.tree_reads.get()), (0, 1)); | |
| 968 | + | } | |
| 969 | + | ||
| 970 | + | #[test] | |
| 971 | + | fn the_commit_a_push_builds_on_is_read_once_however_many_commits_build_on_it() { | |
| 972 | + | // Two branches pushed at once, each one commit on the same parent. | |
| 973 | + | let parent_id = "c71546fcd893ef8b0f57388b65e620d759705dda".to_owned(); | |
| 974 | + | let base_tree_id = object_id(ObjectKind::Tree, &[]); | |
| 975 | + | let mut repo = FakeRepo::default(); | |
| 976 | + | repo.trees.insert(base_tree_id.clone(), Vec::new()); | |
| 977 | + | repo.commits.insert( | |
| 978 | + | parent_id.clone(), | |
| 979 | + | Commit { | |
| 980 | + | hash: parent_id.clone(), | |
| 981 | + | tree_hash: base_tree_id, | |
| 982 | + | message: String::new(), | |
| 983 | + | author: Signature { name: "A".into(), email: "a@example.com".into() }, | |
| 984 | + | parents: Vec::new(), | |
| 985 | + | authored_at: String::new(), | |
| 986 | + | }, | |
| 987 | + | ); | |
| 988 | + | let mut entries = Vec::new(); | |
| 989 | + | for (name, line) in [("a.env", format!("KEY={} | |
| 990 | + | ", key())), ("b.txt", "nothing here | |
| 991 | + | ".to_owned())] { | |
| 992 | + | let blob = line.into_bytes(); | |
| 993 | + | let tree = encode_tree(&[TreeItem { mode: "100644".into(), name: name.into(), id: object_id(ObjectKind::Blob, &blob) }]); | |
| 994 | + | entries.push(Entry::Whole(ObjectKind::Commit, commit(&object_id(ObjectKind::Tree, &tree), Some(&parent_id)))); | |
| 995 | + | entries.push(Entry::Whole(ObjectKind::Tree, tree)); | |
| 996 | + | entries.push(Entry::Whole(ObjectKind::Blob, blob)); | |
| 997 | + | } | |
| 998 | + | let found = run(scan_push(&repo, &push(&entries), &[])).unwrap(); | |
| 999 | + | assert_eq!(found.len(), 1); | |
| 1000 | + | assert_eq!(found[0].path, "a.env"); | |
| 1001 | + | assert_eq!(repo.log_reads.get(), 1); | |
| 1002 | + | } | |
| 1003 | + | ||
| 1004 | + | #[test] | |
| 923 | 1005 | fn a_push_too_large_to_read_is_never_let_through_unread() { | |
| 924 | 1006 | let body = vec![0u8; MAX_SCANNED_PUSH + 1]; | |
| 925 | 1007 | let error = run(scan_push(&FakeRepo::default(), &body, &[])).unwrap_err(); | |
| 104 | 104 | } | |
| 105 | 105 | } | |
| 106 | 106 | ||
| 107 | + | /// Work a request starts that need not hold up its answer, such as keeping | |
| 108 | + | /// an object in the Cache API: begun at once, and handed to the request's | |
| 109 | + | /// `waitUntil` once it has answered ([`Deferred::hand_over`]), so it is | |
| 110 | + | /// neither awaited on the way nor cut short after. Each request has its own. | |
| 111 | + | #[derive(Default)] | |
| 112 | + | pub struct Deferred { | |
| 113 | + | started: RefCell<Vec<worker::js_sys::Promise>>, | |
| 114 | + | } | |
| 115 | + | ||
| 116 | + | impl Deferred { | |
| 117 | + | /// Starts `work` now. | |
| 118 | + | pub fn spawn(&self, work: impl std::future::Future<Output = ()> + 'static) { | |
| 119 | + | let promise = worker::wasm_bindgen_futures::future_to_promise(async move { | |
| 120 | + | work.await; | |
| 121 | + | Ok(JsValue::UNDEFINED) | |
| 122 | + | }); | |
| 123 | + | self.started.borrow_mut().push(promise); | |
| 124 | + | } | |
| 125 | + | ||
| 126 | + | /// Hands what was started to `ctx`, to finish after the answer. | |
| 127 | + | pub fn hand_over(&self, ctx: &worker::Context) { | |
| 128 | + | let started = std::mem::take(&mut *self.started.borrow_mut()); | |
| 129 | + | if started.is_empty() { | |
| 130 | + | return; | |
| 131 | + | } | |
| 132 | + | ctx.wait_until(async move { | |
| 133 | + | futures_util::future::join_all(started.into_iter().map(worker::wasm_bindgen_futures::JsFuture::from)).await; | |
| 134 | + | }); | |
| 135 | + | } | |
| 136 | + | ||
| 137 | + | /// Waits for what was started: for work already running after an | |
| 138 | + | /// answer, in a `waitUntil` of its own. | |
| 139 | + | pub async fn settle(&self) { | |
| 140 | + | let started = std::mem::take(&mut *self.started.borrow_mut()); | |
| 141 | + | futures_util::future::join_all(started.into_iter().map(worker::wasm_bindgen_futures::JsFuture::from)).await; | |
| 142 | + | } | |
| 143 | + | } | |
| 144 | + | ||
| 107 | 145 | /// A place repositories live. `key` is the store's own name for a repo. | |
| 108 | 146 | #[allow(async_fn_in_trait)] | |
| 109 | 147 | pub trait GitStore { | |
| ⋯ | |||
| 271 | 309 | namespaces: Rc<Vec<Namespace>>, | |
| 272 | 310 | /// Where isolates share the credentials they make; see shared.rs. | |
| 273 | 311 | shared: Option<Rc<crate::shared::Shared>>, | |
| 312 | + | /// The request's work that need not hold up its answer: objects kept | |
| 313 | + | /// for next time. | |
| 314 | + | deferred: Rc<Deferred>, | |
| 274 | 315 | } | |
| 275 | 316 | ||
| 276 | 317 | impl ArtifactsStore { | |
| 277 | − | pub fn new(env: &Env, shared: Option<Rc<crate::shared::Shared>>) -> Result<Self> { | |
| 318 | + | pub fn new(env: &Env, shared: Option<Rc<crate::shared::Shared>>, deferred: Rc<Deferred>) -> Result<Self> { | |
| 278 | 319 | let config = env.var("ARTIFACTS_NAMESPACES").ok().map(|value| value.to_string()); | |
| 279 | 320 | let text = |name: &str| env.var(name).ok().map(|value| value.to_string()); | |
| 280 | 321 | let secret = env.secret("GIT_FALLBACK_SECRET").ok().map(|value| value.to_string()); | |
| ⋯ | |||
| 318 | 359 | }); | |
| 319 | 360 | } | |
| 320 | 361 | } | |
| 321 | − | Ok(Self { namespaces: Rc::new(namespaces), shared }) | |
| 362 | + | Ok(Self { namespaces: Rc::new(namespaces), shared, deferred }) | |
| 322 | 363 | } | |
| 323 | 364 | ||
| 324 | 365 | fn binding(&self, namespace: &str) -> Result<&Target> { | |
| ⋯ | |||
| 675 | 716 | namespace, | |
| 676 | 717 | shared: self.shared.clone(), | |
| 677 | 718 | refs_version: None, | |
| 719 | + | deferred: self.deferred.clone(), | |
| 678 | 720 | }) | |
| 679 | 721 | } | |
| 680 | 722 | ||
| ⋯ | |||
| 715 | 757 | shared: Option<Rc<crate::shared::Shared>>, | |
| 716 | 758 | /// See [`GitRepo::at_refs_version`]. | |
| 717 | 759 | refs_version: Option<u64>, | |
| 760 | + | deferred: Rc<Deferred>, | |
| 718 | 761 | } | |
| 719 | 762 | ||
| 720 | 763 | /// Where cached git objects live. Trees and blobs are named by their | |
| ⋯ | |||
| 943 | 986 | self.cached_at(&format!("{kind}/{hash}"), true).await | |
| 944 | 987 | } | |
| 945 | 988 | ||
| 946 | − | async fn keep_at(&self, path: &str, bytes: Vec<u8>, max_age: &str) { | |
| 989 | + | /// Keeps an answer at `path`: in the isolate at once, and in the Cache | |
| 990 | + | /// API by a write that starts now and finishes in the request's | |
| 991 | + | /// `waitUntil`. A read never waits on the write, which only saves a | |
| 992 | + | /// later read. | |
| 993 | + | fn keep_at(&self, path: &str, bytes: Vec<u8>, max_age: &'static str) { | |
| 994 | + | let url = self.cache_url(path); | |
| 947 | 995 | if max_age == OBJECT_MAX_AGE { | |
| 948 | − | MEMORY.with(|memory| memory.borrow_mut().put(self.cache_url(path), &bytes)); | |
| 996 | + | MEMORY.with(|memory| memory.borrow_mut().put(url.clone(), &bytes)); | |
| 949 | 997 | } | |
| 950 | − | let Ok(mut response) = worker::Response::from_bytes(bytes) else { | |
| 951 | − | return; | |
| 952 | − | }; | |
| 953 | − | let _ = response.headers_mut().set("cache-control", max_age); | |
| 954 | − | let _ = worker::Cache::default().put(self.cache_url(path), response).await; | |
| 998 | + | self.deferred.spawn(async move { | |
| 999 | + | let Ok(mut response) = worker::Response::from_bytes(bytes) else { | |
| 1000 | + | return; | |
| 1001 | + | }; | |
| 1002 | + | let _ = response.headers_mut().set("cache-control", max_age); | |
| 1003 | + | let _ = worker::Cache::default().put(url, response).await; | |
| 1004 | + | }); | |
| 955 | 1005 | } | |
| 956 | 1006 | ||
| 957 | 1007 | /// Whether `read_file`'s key was found not to be a file a little while | |
| ⋯ | |||
| 967 | 1017 | } | |
| 968 | 1018 | ||
| 969 | 1019 | /// Keeps an object for next time. A failure only costs a later read. | |
| 970 | − | async fn keep(&self, kind: &str, hash: &str, bytes: Vec<u8>) { | |
| 971 | − | self.keep_at(&format!("{kind}/{hash}"), bytes, OBJECT_MAX_AGE).await; | |
| 1020 | + | fn keep(&self, kind: &str, hash: &str, bytes: Vec<u8>) { | |
| 1021 | + | self.keep_at(&format!("{kind}/{hash}"), bytes, OBJECT_MAX_AGE); | |
| 972 | 1022 | } | |
| 973 | 1023 | ||
| 974 | 1024 | async fn get_key(&self, key: &CacheKey) -> Option<Vec<u8>> { | |
| ⋯ | |||
| 978 | 1028 | } | |
| 979 | 1029 | } | |
| 980 | 1030 | ||
| 981 | − | async fn put_key(&self, key: &CacheKey, bytes: Vec<u8>) { | |
| 1031 | + | fn put_key(&self, key: &CacheKey, bytes: Vec<u8>) { | |
| 982 | 1032 | match key { | |
| 983 | − | CacheKey::Forever(path) => self.keep_at(path, bytes, OBJECT_MAX_AGE).await, | |
| 984 | − | CacheKey::Versioned(path) => self.keep_at(path, bytes, VERSIONED_MAX_AGE).await, | |
| 1033 | + | CacheKey::Forever(path) => self.keep_at(path, bytes, OBJECT_MAX_AGE), | |
| 1034 | + | CacheKey::Versioned(path) => self.keep_at(path, bytes, VERSIONED_MAX_AGE), | |
| 985 | 1035 | } | |
| 986 | 1036 | } | |
| 987 | 1037 | ||
| ⋯ | |||
| 1113 | 1163 | } | |
| 1114 | 1164 | let branches = crate::refs::branches(&self.access(Scope::Read).await?).await?; | |
| 1115 | 1165 | if let (Some(key), Ok(bytes)) = (&key, serde_json::to_vec(&branches)) { | |
| 1116 | − | self.put_key(key, bytes).await; | |
| 1166 | + | self.put_key(key, bytes); | |
| 1117 | 1167 | } | |
| 1118 | 1168 | Ok(branches) | |
| 1119 | 1169 | } | |
| ⋯ | |||
| 1138 | 1188 | && let Ok(bytes) = serde_json::to_vec(&commits) | |
| 1139 | 1189 | { | |
| 1140 | 1190 | if let Some(key) = &key { | |
| 1141 | − | self.put_key(key, bytes.clone()).await; | |
| 1191 | + | self.put_key(key, bytes.clone()); | |
| 1142 | 1192 | } | |
| 1143 | 1193 | // The same history, by the commit the name led to. | |
| 1144 | 1194 | if !is_commit_hash(git_ref) | |
| 1145 | 1195 | && let Some(CacheKey::Forever(path)) = log_key(&commits[0].hash, limit, None) | |
| 1146 | 1196 | { | |
| 1147 | − | self.keep_at(&path, bytes, OBJECT_MAX_AGE).await; | |
| 1197 | + | self.keep_at(&path, bytes, OBJECT_MAX_AGE); | |
| 1148 | 1198 | } | |
| 1149 | 1199 | } | |
| 1150 | 1200 | Ok(commits) | |
| ⋯ | |||
| 1162 | 1212 | if let Some(parents) = &parents | |
| 1163 | 1213 | && let Ok(bytes) = serde_json::to_vec(parents) | |
| 1164 | 1214 | { | |
| 1165 | − | self.keep("commit", commit_hash, bytes).await; | |
| 1215 | + | self.keep("commit", commit_hash, bytes); | |
| 1166 | 1216 | } | |
| 1167 | 1217 | Ok(parents) | |
| 1168 | 1218 | } | |
| ⋯ | |||
| 1187 | 1237 | if let Some(entries) = &entries | |
| 1188 | 1238 | && let Ok(bytes) = serde_json::to_vec(entries) | |
| 1189 | 1239 | { | |
| 1190 | − | self.keep("tree", tree_hash, bytes).await; | |
| 1240 | + | self.keep("tree", tree_hash, bytes); | |
| 1191 | 1241 | } | |
| 1192 | 1242 | Ok(entries) | |
| 1193 | 1243 | } | |
| ⋯ | |||
| 1201 | 1251 | meters::record_bytes("binding.read_blob", &self.key, 0, bytes.len() as u64); | |
| 1202 | 1252 | } | |
| 1203 | 1253 | if let Some(bytes) = bytes.as_ref().filter(|bytes| bytes.len() <= MAX_CACHED_BLOB) { | |
| 1204 | − | self.keep("blob", blob_hash, bytes.clone()).await; | |
| 1254 | + | self.keep("blob", blob_hash, bytes.clone()); | |
| 1205 | 1255 | } | |
| 1206 | 1256 | Ok(bytes) | |
| 1207 | 1257 | } | |
| ⋯ | |||
| 1225 | 1275 | }; | |
| 1226 | 1276 | if let Some(size) = size { | |
| 1227 | 1277 | meters::record_bytes("binding.read_blob", &self.key, 0, size); | |
| 1228 | − | self.keep_at(&path, size.to_string().into_bytes(), OBJECT_MAX_AGE).await; | |
| 1278 | + | self.keep_at(&path, size.to_string().into_bytes(), OBJECT_MAX_AGE); | |
| 1229 | 1279 | } | |
| 1230 | 1280 | Ok(size) | |
| 1231 | 1281 | } | |
| ⋯ | |||
| 1260 | 1310 | if let Some(bytes) = &bytes { | |
| 1261 | 1311 | meters::record_bytes("binding.read_file", &self.key, 0, bytes.len() as u64); | |
| 1262 | 1312 | } else if remember_absent && let Some(key) = &key { | |
| 1263 | − | self.keep_at(&absent_path(key), vec![1], ABSENT_MAX_AGE).await; | |
| 1313 | + | self.keep_at(&absent_path(key), vec![1], ABSENT_MAX_AGE); | |
| 1264 | 1314 | } | |
| 1265 | 1315 | if let (Some(key), Some(bytes)) = (&key, bytes.as_ref().filter(|bytes| bytes.len() <= MAX_CACHED_BLOB)) { | |
| 1266 | − | self.put_key(key, bytes.clone()).await; | |
| 1316 | + | self.put_key(key, bytes.clone()); | |
| 1267 | 1317 | } | |
| 1268 | 1318 | Ok(bytes) | |
| 1269 | 1319 | } | |
| 17 | 17 | ||
| 18 | 18 | use g1t_contracts::User; | |
| 19 | 19 | use g1t_contracts::scopes::{TokenAccess, WORKFLOW_DIRS, decide_workflow_files}; | |
| 20 | − | use g1t_scan::pack::{ObjectKind, Pack, TreeItem, pack_start}; | |
| 21 | − | use worker::{Response, Result}; | |
| 20 | + | use g1t_scan::pack::{ObjectKind, Pack, TreeItem}; | |
| 21 | + | use worker::Result; | |
| 22 | 22 | ||
| 23 | − | use crate::registry::store_key; | |
| 24 | 23 | use crate::rule_facts::{self, MAX_COMMITS}; | |
| 25 | 24 | use crate::secret_scan::Objects; | |
| 26 | − | use crate::store::{GitRepo, GitStore}; | |
| 27 | − | use crate::Repos; | |
| 25 | + | use crate::store::GitRepo; | |
| 28 | 26 | ||
| 29 | 27 | /// The token behind a push or an edit, when it is one this gate checks: | |
| 30 | 28 | /// any token without `workflow_files:write`, and every job's token. | |
| ⋯ | |||
| 104 | 102 | /// Why a push is refused, as the reason git shows beside each ref and the | |
| 105 | 103 | /// lines it prints, when `token` may not change workflow files and the | |
| 106 | 104 | /// push (`body`, read `whole` or not) does or cannot be checked. | |
| 105 | + | #[cfg(test)] | |
| 107 | 106 | pub(crate) async fn judge_push<R: GitRepo>(token: &TokenAccess, body: &[u8], whole: bool, git: &R) -> Result<Option<(&'static str, Vec<String>)>> { | |
| 107 | + | let mut pack = whole.then(|| crate::push_checks::read_pack(body)); | |
| 108 | + | if let Some(Ok(pack)) = &mut pack { | |
| 109 | + | crate::secret_scan::supply_bases(pack, git).await?; | |
| 110 | + | } | |
| 111 | + | judge_pack(token, body, pack.as_ref(), git).await | |
| 112 | + | } | |
| 113 | + | ||
| 114 | + | /// [`judge_push`] for a push whose pack was read already, with its bases | |
| 115 | + | /// supplied (push_checks.rs); `None` when it was not read whole. | |
| 116 | + | pub(crate) async fn judge_pack<R: GitRepo>( | |
| 117 | + | token: &TokenAccess, | |
| 118 | + | body: &[u8], | |
| 119 | + | pack: Option<&std::result::Result<Pack, String>>, | |
| 120 | + | git: &R, | |
| 121 | + | ) -> Result<Option<(&'static str, Vec<String>)>> { | |
| 108 | 122 | let updates = crate::git_http::ref_updates(body); | |
| 109 | 123 | if updates.iter().all(|(_, _, new)| new.is_none()) { | |
| 110 | 124 | return Ok(None); | |
| 111 | 125 | } | |
| 112 | 126 | let too_large = |line: String| Ok(Some(("push too large to check for workflow files", vec![line, "Push in smaller parts, or with a token that has the workflow_files:write scope.".to_owned()]))); | |
| 113 | − | if !whole { | |
| 127 | + | let Some(pack) = pack else { | |
| 114 | 128 | return too_large("This push is too large for g1t to check whether it changes workflow files, and this token may not change them.".to_owned()); | |
| 115 | − | } | |
| 116 | − | let mut pack = match pack_start(body).map(|start| Pack::parse(&body[start..])) { | |
| 117 | − | Some(Ok(pack)) => pack, | |
| 118 | − | // Nothing but ref moves to commits the repository has. | |
| 119 | − | None => return Ok(None), | |
| 120 | − | Some(Err(problem)) => { | |
| 129 | + | }; | |
| 130 | + | // A push with no pack moves refs to commits the repository has: an | |
| 131 | + | // empty pack, which changes nothing. | |
| 132 | + | let pack = match pack { | |
| 133 | + | Ok(pack) => pack, | |
| 134 | + | Err(problem) => { | |
| 121 | 135 | worker::console_error!("a push's pack could not be read for workflow files: {problem}"); | |
| 122 | 136 | return Ok(Some(("push could not be checked for workflow files", vec!["g1t could not read this push to check it for workflow files. Push again.".to_owned()]))); | |
| 123 | 137 | } | |
| 124 | 138 | }; | |
| 125 | − | crate::secret_scan::supply_bases(&mut pack, git).await?; | |
| 126 | 139 | for (_, _, new) in updates { | |
| 127 | 140 | let Some(new) = new else { continue }; | |
| 128 | − | let Some(ids) = rule_facts::added(&pack, &new, MAX_COMMITS) else { | |
| 141 | + | let Some(ids) = rule_facts::added(pack, &new, MAX_COMMITS) else { | |
| 129 | 142 | return too_large(format!("This push adds more than {MAX_COMMITS} commits to one ref, too many for g1t to check for workflow files, and this token may not change them.")); | |
| 130 | 143 | }; | |
| 131 | − | if let Some(path) = changed_in_commits(&pack, git, &ids).await? { | |
| 144 | + | if let Some(path) = changed_in_commits(pack, git, &ids).await? { | |
| 132 | 145 | let reason = decide_workflow_files(Some(token), [path.as_str()]) | |
| 133 | 146 | .and_then(|decision| decision.reason) | |
| 134 | 147 | .unwrap_or_else(|| format!("This access token cannot change the workflow file {path}.")); | |
| ⋯ | |||
| 139 | 152 | } | |
| 140 | 153 | } | |
| 141 | 154 | Ok(None) | |
| 142 | − | } | |
| 143 | − | ||
| 144 | − | impl<S: GitStore> Repos<S> { | |
| 145 | − | /// For a push by a token this gate checks: the response declining it | |
| 146 | − | /// when it changes a workflow file, or when it cannot be checked. | |
| 147 | − | pub(crate) async fn workflow_gate(&self, repo: &g1t_contracts::repos::Repo, pusher: Option<&User>, body: &[u8], whole: bool) -> Result<Option<Response>> { | |
| 148 | − | let Some(token) = gated(pusher) else { | |
| 149 | − | return Ok(None); | |
| 150 | − | }; | |
| 151 | − | let git = self.store.open(&store_key(repo)).await?; | |
| 152 | − | match judge_push(token, body, whole, &git).await? { | |
| 153 | − | Some((reason, lines)) => Ok(Some(crate::git_http::declined(body, reason, &lines)?)), | |
| 154 | − | None => Ok(None), | |
| 155 | − | } | |
| 156 | − | } | |
| 157 | 155 | } | |
| 158 | 156 | ||
| 159 | 157 | #[cfg(test)] | |