flagon-io/g1t

public

Where people and agents ship software together. The open-source git platform for the whole job: issues, agents, checks and deploys to the edge.

g1t/services/repos/src/land.rs

294 lines12,081 bytesCodeBlame
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
12use std::cell::RefCell;
13use std::rc::Rc;
14
15use futures_util::StreamExt;
16use g1t_contracts::repos::GitAccess;
17use worker::js_sys::Uint8Array;
18use worker::wasm_bindgen::JsValue;
19use worker::{Error, Fetch, Headers, Method, Request, RequestInit, Response, Result};
20
21use crate::meters;
22use crate::pack_limits::Sideband;
23
24const ZERO_ID: &str = "0000000000000000000000000000000000000000";
25const FLUSH: &[u8] = b"0000";
26/// Side-band channels: pack data, and fatal errors.
27const PACK_BAND: u8 = 1;
28const 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.
32pub 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
37fn 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).
42fn 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.
47async 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
64async 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.
75pub(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`.
99fn 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`.
114pub(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.
120pub(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`.
148pub(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`.
159pub(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.
168fn 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.
176fn 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.
190pub(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.
209pub 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)]
263mod 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}