g1t/services/repos/src/land.rs
| 1 | //! Landing a pull request: moving a repository's branch forward to a commit |
| 2 | //! from a fork, or from another of its own branches. |
| 3 | //! |
| 4 | //! The Artifacts binding cannot write, so this speaks git's smart HTTP |
| 5 | //! protocol directly. It asks the fork for a pack holding exactly the |
| 6 | //! objects the target is missing, then pushes that pack to the target |
| 7 | //! unchanged. No object is parsed or rebuilt along the way. When landing, |
| 8 | //! the pack streams from one to the other as it arrives |
| 9 | //! ([`fast_forward`]): it is never held whole, however large the pull |
| 10 | //! request (pack_limits.rs `Sideband`). |
| 11 | |
| 12 | use std::cell::RefCell; |
| 13 | use std::rc::Rc; |
| 14 | |
| 15 | use futures_util::StreamExt; |
| 16 | use g1t_contracts::repos::GitAccess; |
| 17 | use worker::js_sys::Uint8Array; |
| 18 | use worker::wasm_bindgen::JsValue; |
| 19 | use worker::{Error, Fetch, Headers, Method, Request, RequestInit, Response, Result}; |
| 20 | |
| 21 | use crate::meters; |
| 22 | use crate::pack_limits::Sideband; |
| 23 | |
| 24 | const ZERO_ID: &str = "0000000000000000000000000000000000000000"; |
| 25 | const FLUSH: &[u8] = b"0000"; |
| 26 | /// Side-band channels: pack data, and fatal errors. |
| 27 | const PACK_BAND: u8 = 1; |
| 28 | const ERROR_BAND: u8 = 3; |
| 29 | |
| 30 | /// A pack with no objects: what a push that only points a ref at a commit |
| 31 | /// the repository already has sends. |
| 32 | pub const EMPTY_PACK: &[u8] = &[ |
| 33 | b'P', b'A', b'C', b'K', 0, 0, 0, 2, 0, 0, 0, 0, 0x02, 0x9d, 0x08, 0x82, 0x3b, 0xd8, 0xa8, |
| 34 | 0xea, 0xb5, 0x10, 0xad, 0x6a, 0xc7, 0x5c, 0x82, 0x3c, 0xfd, 0x3e, 0xd3, 0x1e, |
| 35 | ]; |
| 36 | |
| 37 | fn pkt_line(payload: &str) -> Vec<u8> { |
| 38 | format!("{:04x}{payload}", payload.len() + 4).into_bytes() |
| 39 | } |
| 40 | |
| 41 | /// The meter for g1t's own request to `service` (meters.rs). |
| 42 | fn meter(service: &str) -> &'static str { |
| 43 | if service == "git-receive-pack" { "internal.git.receive_pack" } else { "internal.git.fetch" } |
| 44 | } |
| 45 | |
| 46 | /// Sends `body` to `service` on the store, answered in full or an error. |
| 47 | async fn send(access: &GitAccess, service: &str, body: JsValue, sent: u64) -> Result<Response> { |
| 48 | let headers = Headers::new(); |
| 49 | headers.set("authorization", &format!("Bearer {}", access.token))?; |
| 50 | headers.set("content-type", &format!("application/x-{service}-request"))?; |
| 51 | headers.set("accept", &format!("application/x-{service}-result"))?; |
| 52 | let mut init = RequestInit::new(); |
| 53 | init.with_method(Method::Post).with_headers(headers).with_body(Some(body)); |
| 54 | let request = Request::new_with_init(&format!("{}/{service}", access.remote), &init)?; |
| 55 | let mut response = Fetch::Request(request).send().await?; |
| 56 | meters::record_remote(meter(service), &access.remote, sent, 0); |
| 57 | if response.status_code() != 200 { |
| 58 | let text = response.text().await.unwrap_or_default(); |
| 59 | return Err(Error::RustError(format!("{service} returned {}: {text}", response.status_code()))); |
| 60 | } |
| 61 | Ok(response) |
| 62 | } |
| 63 | |
| 64 | async fn post(access: &GitAccess, service: &str, body: Vec<u8>) -> Result<Vec<u8>> { |
| 65 | let sent = body.len() as u64; |
| 66 | let mut response = send(access, service, Uint8Array::from(body.as_slice()).into(), sent).await?; |
| 67 | let bytes = response.bytes().await?; |
| 68 | if let Some(key) = crate::store::key_from_remote(&access.remote) { |
| 69 | meters::record_bytes(meter(service), &key, 0, bytes.len() as u64); |
| 70 | } |
| 71 | Ok(bytes) |
| 72 | } |
| 73 | |
| 74 | /// The payloads of the pkt-lines in `bytes`, and the offset where they stop. |
| 75 | pub(crate) fn read_pkt_lines(bytes: &[u8]) -> (Vec<&[u8]>, usize) { |
| 76 | let mut lines = Vec::new(); |
| 77 | let mut position = 0; |
| 78 | while position + 4 <= bytes.len() { |
| 79 | let Some(length) = std::str::from_utf8(&bytes[position..position + 4]) |
| 80 | .ok() |
| 81 | .and_then(|hex| usize::from_str_radix(hex, 16).ok()) |
| 82 | else { |
| 83 | break; |
| 84 | }; |
| 85 | if length < 4 { |
| 86 | // Flush, delimiter and response-end packets carry no payload. |
| 87 | position += 4; |
| 88 | continue; |
| 89 | } |
| 90 | let end = (position + length).min(bytes.len()); |
| 91 | lines.push(&bytes[position + 4..end]); |
| 92 | position = end; |
| 93 | } |
| 94 | (lines, position) |
| 95 | } |
| 96 | |
| 97 | /// The upload-pack request for everything reachable from `want` that is |
| 98 | /// not reachable from `have`. |
| 99 | fn want_request(want: &str, have: Option<&str>) -> Vec<u8> { |
| 100 | // Side-band framing puts the pack in its own channel, so its exact bytes |
| 101 | // can be recovered. Without it the response ends in a stray flush packet |
| 102 | // that a receiver rejects as junk after the pack. |
| 103 | let mut body = pkt_line(&format!("want {want} side-band-64k\n")); |
| 104 | body.extend_from_slice(FLUSH); |
| 105 | if let Some(have) = have { |
| 106 | body.extend(pkt_line(&format!("have {have}\n"))); |
| 107 | } |
| 108 | body.extend(pkt_line("done\n")); |
| 109 | body |
| 110 | } |
| 111 | |
| 112 | /// A pack holding everything reachable from `want` that is not reachable |
| 113 | /// from `have`. |
| 114 | pub(crate) async fn fetch_pack(source: &GitAccess, want: &str, have: Option<&str>) -> Result<Vec<u8>> { |
| 115 | let response = post(source, "git-upload-pack", want_request(want, have)).await?; |
| 116 | unpack_sideband(&response) |
| 117 | } |
| 118 | |
| 119 | /// The pack in an upload-pack response that used side-band framing. |
| 120 | pub(crate) fn unpack_sideband(response: &[u8]) -> Result<Vec<u8>> { |
| 121 | let (lines, _) = read_pkt_lines(response); |
| 122 | let mut pack = Vec::new(); |
| 123 | for line in lines { |
| 124 | match line.first() { |
| 125 | Some(&PACK_BAND) => pack.extend_from_slice(&line[1..]), |
| 126 | Some(&ERROR_BAND) => { |
| 127 | return Err(Error::RustError(format!( |
| 128 | "the source refused the fetch: {}", |
| 129 | String::from_utf8_lossy(&line[1..]) |
| 130 | ))); |
| 131 | } |
| 132 | // ACK and NAK lines, and progress messages. |
| 133 | _ => {} |
| 134 | } |
| 135 | } |
| 136 | if !pack.starts_with(b"PACK") { |
| 137 | return Err(Error::RustError(format!( |
| 138 | "the source did not send a pack: {}", |
| 139 | String::from_utf8_lossy(response) |
| 140 | ))); |
| 141 | } |
| 142 | Ok(pack) |
| 143 | } |
| 144 | |
| 145 | /// Updates `branch` on the target from `old` to `new`, sending `pack`. |
| 146 | /// `Err(reason)` in the inner result means git refused the update, for |
| 147 | /// example because the branch is no longer at `old`. |
| 148 | pub(crate) async fn push_pack( |
| 149 | target: &GitAccess, |
| 150 | branch: &str, |
| 151 | old: Option<&str>, |
| 152 | new: &str, |
| 153 | pack: Vec<u8>, |
| 154 | ) -> Result<std::result::Result<(), String>> { |
| 155 | push_ref(target, &format!("refs/heads/{branch}"), old, new, Some(pack)).await |
| 156 | } |
| 157 | |
| 158 | /// Removes `branch` from the target, if it is still at `old`. |
| 159 | pub(crate) async fn delete_ref( |
| 160 | target: &GitAccess, |
| 161 | branch: &str, |
| 162 | old: &str, |
| 163 | ) -> Result<std::result::Result<(), String>> { |
| 164 | push_ref(target, &format!("refs/heads/{branch}"), Some(old), ZERO_ID, None).await |
| 165 | } |
| 166 | |
| 167 | /// The commands at the start of a receive-pack request for one ref. |
| 168 | fn command(reference: &str, old: Option<&str>, new: &str, sends_pack: bool) -> Vec<u8> { |
| 169 | let capabilities = if sends_pack { "report-status" } else { "report-status delete-refs" }; |
| 170 | let mut body = pkt_line(&format!("{} {new} {reference}\0 {capabilities}\n", old.unwrap_or(ZERO_ID))); |
| 171 | body.extend_from_slice(FLUSH); |
| 172 | body |
| 173 | } |
| 174 | |
| 175 | /// Whether a report-status says `reference` moved. |
| 176 | fn reported(report: &[u8], reference: &str, sent_pack: bool) -> std::result::Result<(), String> { |
| 177 | let (lines, _) = read_pkt_lines(report); |
| 178 | let lines: Vec<String> = lines |
| 179 | .into_iter() |
| 180 | .map(|line| String::from_utf8_lossy(line).trim_end().to_owned()) |
| 181 | .collect(); |
| 182 | // With nothing to unpack a server may not say so. |
| 183 | let unpacked = !sent_pack || lines.iter().any(|line| line == "unpack ok"); |
| 184 | let updated = lines.iter().any(|line| *line == format!("ok {reference}")); |
| 185 | if unpacked && updated { Ok(()) } else { Err(lines.join("; ")) } |
| 186 | } |
| 187 | |
| 188 | /// Moves one ref (`refs/heads/main`, `refs/pull/<id>/head`) on the target |
| 189 | /// from `old` to `new`; a deletion sends no pack. |
| 190 | pub(crate) async fn push_ref( |
| 191 | target: &GitAccess, |
| 192 | reference: &str, |
| 193 | old: Option<&str>, |
| 194 | new: &str, |
| 195 | pack: Option<Vec<u8>>, |
| 196 | ) -> Result<std::result::Result<(), String>> { |
| 197 | let sends_pack = pack.is_some(); |
| 198 | let mut body = command(reference, old, new, sends_pack); |
| 199 | if let Some(pack) = pack { |
| 200 | body.extend(pack); |
| 201 | } |
| 202 | let response = post(target, "git-receive-pack", body).await?; |
| 203 | Ok(reported(&response, reference, sends_pack)) |
| 204 | } |
| 205 | |
| 206 | /// Moves `branch` on `target` from `old` to `new`, a commit that exists in |
| 207 | /// `source`. The caller must have checked that `new` descends from `old`. |
| 208 | /// The pack goes from the source's answer into the push as it arrives. |
| 209 | pub async fn fast_forward( |
| 210 | source: &GitAccess, |
| 211 | target: &GitAccess, |
| 212 | branch: &str, |
| 213 | old: Option<&str>, |
| 214 | new: &str, |
| 215 | ) -> Result<std::result::Result<(), String>> { |
| 216 | let request = want_request(new, old); |
| 217 | let sent = request.len() as u64; |
| 218 | let mut fetched = send(source, "git-upload-pack", Uint8Array::from(request.as_slice()).into(), sent).await?; |
| 219 | let reference = format!("refs/heads/{branch}"); |
| 220 | let failure: Rc<RefCell<Option<String>>> = Rc::default(); |
| 221 | let demux = Rc::new(RefCell::new(Sideband::default())); |
| 222 | let pack = { |
| 223 | let failure = failure.clone(); |
| 224 | let demux = demux.clone(); |
| 225 | fetched.stream()?.map(move |chunk| { |
| 226 | let chunk = chunk?; |
| 227 | demux.borrow_mut().feed(&chunk).map_err(|why| { |
| 228 | *failure.borrow_mut() = Some(why.clone()); |
| 229 | Error::RustError(why) |
| 230 | }) |
| 231 | }) |
| 232 | }; |
| 233 | // The source said all it had: a pack must have come. |
| 234 | let end = { |
| 235 | let failure = failure.clone(); |
| 236 | let demux = demux.clone(); |
| 237 | futures_util::stream::once(async move { |
| 238 | demux.borrow().finish().map(|()| Vec::new()).map_err(|why| { |
| 239 | *failure.borrow_mut() = Some(why.clone()); |
| 240 | Error::RustError(why) |
| 241 | }) |
| 242 | }) |
| 243 | }; |
| 244 | let head = command(&reference, old, new, true); |
| 245 | let body = futures_util::stream::once(async move { Ok::<Vec<u8>, Error>(head) }) |
| 246 | .chain(pack) |
| 247 | .chain(end) |
| 248 | .filter(|chunk| futures_util::future::ready(!matches!(chunk, Ok(bytes) if bytes.is_empty()))); |
| 249 | let pushed = send(target, "git-receive-pack", crate::git_http::stream_body(body)?, 0).await; |
| 250 | if let Some(why) = failure.borrow_mut().take() { |
| 251 | return Err(Error::RustError(why)); |
| 252 | } |
| 253 | let report = pushed?.bytes().await?; |
| 254 | let moved = demux.borrow().pack_bytes; |
| 255 | if let (Some(from), Some(to)) = (crate::store::key_from_remote(&source.remote), crate::store::key_from_remote(&target.remote)) { |
| 256 | meters::record_bytes("internal.git.fetch", &from, 0, moved); |
| 257 | meters::record_bytes("internal.git.receive_pack", &to, moved, 0); |
| 258 | } |
| 259 | Ok(reported(&report, &reference, true)) |
| 260 | } |
| 261 | |
| 262 | #[cfg(test)] |
| 263 | mod tests { |
| 264 | use super::*; |
| 265 | |
| 266 | #[test] |
| 267 | fn one_ref_is_moved_by_one_command() { |
| 268 | let new = "4807077b296e6edbf410d55e72749d3e1170c291"; |
| 269 | let body = String::from_utf8(command("refs/pull/pul_1/head", None, new, true)).unwrap(); |
| 270 | assert!(body.starts_with(&format!("{:04x}", body.len() - 4))); |
| 271 | assert!(body.contains(&format!("{ZERO_ID} {new} refs/pull/pul_1/head\0 report-status\n"))); |
| 272 | assert!(body.ends_with("0000")); |
| 273 | let deletion = String::from_utf8(command("refs/heads/x", Some(new), ZERO_ID, false)).unwrap(); |
| 274 | assert!(deletion.contains("delete-refs")); |
| 275 | } |
| 276 | |
| 277 | #[test] |
| 278 | fn a_report_says_whether_the_ref_moved() { |
| 279 | let ok = [pkt_line("unpack ok\n"), pkt_line("ok refs/heads/main\n"), FLUSH.to_vec()].concat(); |
| 280 | assert_eq!(reported(&ok, "refs/heads/main", true), Ok(())); |
| 281 | let refused = [pkt_line("unpack ok\n"), pkt_line("ng refs/heads/main non-fast-forward\n"), FLUSH.to_vec()].concat(); |
| 282 | assert!(reported(&refused, "refs/heads/main", true).unwrap_err().contains("non-fast-forward")); |
| 283 | // Nothing unpacked: a server may not say so. |
| 284 | let bare = [pkt_line("ok refs/heads/main\n"), FLUSH.to_vec()].concat(); |
| 285 | assert_eq!(reported(&bare, "refs/heads/main", false), Ok(())); |
| 286 | assert!(reported(&bare, "refs/heads/main", true).is_err()); |
| 287 | } |
| 288 | |
| 289 | #[test] |
| 290 | fn the_empty_pack_is_a_pack_of_nothing() { |
| 291 | let pack = g1t_scan::pack::write_pack(&[]); |
| 292 | assert_eq!(pack, EMPTY_PACK); |
| 293 | } |
| 294 | } |