g1t/services/repos/src/land.rs

397 lines16,843 bytesCodeBlame

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

Issues and pull requests replace intents and attempts1//! Landing a pull request: moving a repository's branch forward to a commit
Pull requests from branches2//! from a fork, or from another of its own branches.
Rust repos service with shipping; pull requests kept in the model3//!
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
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily7//! 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`).
Rust repos service with shipping; pull requests kept in the model11
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily12use std::cell::RefCell;
13use std::rc::Rc;
14
15use futures_util::StreamExt;
Rust repos service with shipping; pull requests kept in the model16use g1t_contracts::repos::GitAccess;
17use worker::js_sys::Uint8Array;
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily18use worker::wasm_bindgen::JsValue;
19use worker::{Error, Fetch, Headers, Method, Request, RequestInit, Response, Result};
20
21use crate::meters;
22use crate::pack_limits::Sideband;
Rust repos service with shipping; pull requests kept in the model23
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
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily30/// 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
Rust repos service with shipping; pull requests kept in the model37fn pkt_line(payload: &str) -> Vec<u8> {
38 format!("{:04x}{payload}", payload.len() + 4).into_bytes()
39}
40
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily41/// 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> {
Rust repos service with shipping; pull requests kept in the model48 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();
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily53 init.with_method(Method::Post).with_headers(headers).with_body(Some(body));
Rust repos service with shipping; pull requests kept in the model54 let request = Request::new_with_init(&format!("{}/{service}", access.remote), &init)?;
55 let mut response = Fetch::Request(request).send().await?;
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily56 meters::record_remote(meter(service), &access.remote, sent, 0);
Rust repos service with shipping; pull requests kept in the model57 if response.status_code() != 200 {
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily58 let text = response.text().await.unwrap_or_default();
59 return Err(Error::RustError(format!("{service} returned {}: {text}", response.status_code())));
Rust repos service with shipping; pull requests kept in the model60 }
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily61 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 }
Rust repos service with shipping; pull requests kept in the model71 Ok(bytes)
72}
73
74/// The payloads of the pkt-lines in `bytes`, and the offset where they stop.
Pull requests from branches75pub(crate) fn read_pkt_lines(bytes: &[u8]) -> (Vec<&[u8]>, usize) {
Rust repos service with shipping; pull requests kept in the model76 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
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily97/// 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> {
Rust repos service with shipping; pull requests kept in the model100 // 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"));
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily109 body
110}
Rust repos service with shipping; pull requests kept in the model111
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily112/// 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?;
Agents as a team: lifecycle, merge queue, billing and a new shell116 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);
Rust repos service with shipping; pull requests kept in the model122 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!(
Pull requests from branches128 "the source refused the fetch: {}",
Rust repos service with shipping; pull requests kept in the model129 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!(
Pull requests from branches138 "the source did not send a pack: {}",
Agents as a team: lifecycle, merge queue, billing and a new shell139 String::from_utf8_lossy(response)
Rust repos service with shipping; pull requests kept in the model140 )));
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`.
Agents as a team: lifecycle, merge queue, billing and a new shell148pub(crate) async fn push_pack(
Rust repos service with shipping; pull requests kept in the model149 target: &GitAccess,
150 branch: &str,
151 old: Option<&str>,
152 new: &str,
153 pack: Vec<u8>,
154) -> Result<std::result::Result<(), String>> {
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily155 push_ref(target, &format!("refs/heads/{branch}"), old, new, Some(pack)).await
Merge queue: tested states are deleted once their entry leaves156}
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>> {
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily164 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
Merge queue: tested states are deleted once their entry leaves173}
174
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily175/// 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(
Merge queue: tested states are deleted once their entry leaves191 target: &GitAccess,
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily192 reference: &str,
Merge queue: tested states are deleted once their entry leaves193 old: Option<&str>,
194 new: &str,
195 pack: Option<Vec<u8>>,
196) -> Result<std::result::Result<(), String>> {
197 let sends_pack = pack.is_some();
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily198 let mut body = command(reference, old, new, sends_pack);
Merge queue: tested states are deleted once their entry leaves199 if let Some(pack) = pack {
200 body.extend(pack);
201 }
Rust repos service with shipping; pull requests kept in the model202 let response = post(target, "git-receive-pack", body).await?;
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily203 Ok(reported(&response, reference, sends_pack))
Rust repos service with shipping; pull requests kept in the model204}
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`.
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily208/// The pack goes from the source's answer into the push as it arrives.
Rust repos service with shipping; pull requests kept in the model209pub 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>> {
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily216 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
Merge branch 'worktree-agent-a2013627e5ea4ab13'262/// The commands at the start of a receive-pack request for many refs:
263/// `(name, old, new)`, `new` the zero id to delete.
264fn commands_block(commands: &[crate::mirror::Command]) -> Vec<u8> {
265 let deletes = commands.iter().any(|(_, _, new)| new == ZERO_ID);
266 let capabilities = if deletes { "report-status delete-refs" } else { "report-status" };
267 let mut body = Vec::new();
268 for (index, (name, old, new)) in commands.iter().enumerate() {
269 let old = old.as_deref().unwrap_or(ZERO_ID);
270 let line = if index == 0 { format!("{old} {new} {name}\0 {capabilities}\n") } else { format!("{old} {new} {name}\n") };
271 body.extend(pkt_line(&line));
272 }
273 body.extend_from_slice(FLUSH);
274 body
275}
276
277/// Makes `target`'s refs match `commands`, fetching from `source` the
278/// objects `wants` names that `haves` (commits the target holds) do not
279/// reach. The pack streams from one into the other, never held whole, so
280/// a whole repository can be copied this way (moves.rs). `Err` in the
281/// inner result says which refs the target refused, and why.
282pub(crate) async fn copy_refs(
283 source: &GitAccess,
284 target: &GitAccess,
285 commands: &[crate::mirror::Command],
286 wants: &[String],
287 haves: &[String],
288) -> Result<std::result::Result<(), String>> {
289 if commands.is_empty() {
290 return Ok(Ok(()));
291 }
292 let head = commands_block(commands);
293 let sends_pack = commands.iter().any(|(_, _, new)| new != ZERO_ID);
294 let report = if wants.is_empty() {
295 // Every object is there already: only refs move.
296 let mut body = head;
297 if sends_pack {
298 body.extend_from_slice(EMPTY_PACK);
299 }
300 post(target, "git-receive-pack", body).await?
301 } else {
302 let request = crate::mirror::upload_request(wants, haves);
303 let sent = request.len() as u64;
304 let mut fetched = send(source, "git-upload-pack", Uint8Array::from(request.as_slice()).into(), sent).await?;
305 let failure: Rc<RefCell<Option<String>>> = Rc::default();
306 let demux = Rc::new(RefCell::new(Sideband::default()));
307 let pack = {
308 let failure = failure.clone();
309 let demux = demux.clone();
310 fetched.stream()?.map(move |chunk| {
311 let chunk = chunk?;
312 demux.borrow_mut().feed(&chunk).map_err(|why| {
313 *failure.borrow_mut() = Some(why.clone());
314 Error::RustError(why)
315 })
316 })
317 };
318 let end = {
319 let failure = failure.clone();
320 let demux = demux.clone();
321 futures_util::stream::once(async move {
322 demux.borrow().finish().map(|()| Vec::new()).map_err(|why| {
323 *failure.borrow_mut() = Some(why.clone());
324 Error::RustError(why)
325 })
326 })
327 };
328 let body = futures_util::stream::once(async move { Ok::<Vec<u8>, Error>(head) })
329 .chain(pack)
330 .chain(end)
331 .filter(|chunk| futures_util::future::ready(!matches!(chunk, Ok(bytes) if bytes.is_empty())));
332 let pushed = send(target, "git-receive-pack", crate::git_http::stream_body(body)?, 0).await;
333 if let Some(why) = failure.borrow_mut().take() {
334 return Err(Error::RustError(why));
335 }
336 let report = pushed?.bytes().await?;
337 let moved = demux.borrow().pack_bytes;
338 if let (Some(from), Some(to)) = (crate::store::key_from_remote(&source.remote), crate::store::key_from_remote(&target.remote)) {
339 meters::record_bytes("internal.git.fetch", &from, 0, moved);
340 meters::record_bytes("internal.git.receive_pack", &to, moved, 0);
341 }
342 report
343 };
344 let problems = crate::mirror::refused(&report, commands);
345 Ok(if problems.is_empty() { Ok(()) } else { Err(problems.join("; ")) })
346}
347
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily348#[cfg(test)]
349mod tests {
350 use super::*;
351
352 #[test]
Merge branch 'worktree-agent-a2013627e5ea4ab13'353 fn many_refs_are_moved_by_one_block_of_commands() {
354 let a = "c71546fcd893ef8b0f57388b65e620d759705dda".to_owned();
355 let b = "4807077b296e6edbf410d55e72749d3e1170c291".to_owned();
356 let block = String::from_utf8(commands_block(&[
357 ("refs/heads/main".to_owned(), None, a.clone()),
358 ("refs/tags/v1".to_owned(), Some(a.clone()), b.clone()),
359 ]))
360 .unwrap();
361 // Capabilities ride on the first command only.
362 assert!(block.contains(&format!("{ZERO_ID} {a} refs/heads/main\0 report-status\n")));
363 assert!(block.contains(&format!("{a} {b} refs/tags/v1\n")));
364 assert!(block.ends_with("0000"));
365 let deleting = String::from_utf8(commands_block(&[("refs/heads/old".to_owned(), Some(a), ZERO_ID.to_owned())])).unwrap();
366 assert!(deleting.contains("delete-refs"));
367 }
368
369 #[test]
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily370 fn one_ref_is_moved_by_one_command() {
371 let new = "4807077b296e6edbf410d55e72749d3e1170c291";
372 let body = String::from_utf8(command("refs/pull/pul_1/head", None, new, true)).unwrap();
373 assert!(body.starts_with(&format!("{:04x}", body.len() - 4)));
374 assert!(body.contains(&format!("{ZERO_ID} {new} refs/pull/pul_1/head\0 report-status\n")));
375 assert!(body.ends_with("0000"));
376 let deletion = String::from_utf8(command("refs/heads/x", Some(new), ZERO_ID, false)).unwrap();
377 assert!(deletion.contains("delete-refs"));
378 }
379
380 #[test]
381 fn a_report_says_whether_the_ref_moved() {
382 let ok = [pkt_line("unpack ok\n"), pkt_line("ok refs/heads/main\n"), FLUSH.to_vec()].concat();
383 assert_eq!(reported(&ok, "refs/heads/main", true), Ok(()));
384 let refused = [pkt_line("unpack ok\n"), pkt_line("ng refs/heads/main non-fast-forward\n"), FLUSH.to_vec()].concat();
385 assert!(reported(&refused, "refs/heads/main", true).unwrap_err().contains("non-fast-forward"));
386 // Nothing unpacked: a server may not say so.
387 let bare = [pkt_line("ok refs/heads/main\n"), FLUSH.to_vec()].concat();
388 assert_eq!(reported(&bare, "refs/heads/main", false), Ok(()));
389 assert!(reported(&bare, "refs/heads/main", true).is_err());
390 }
391
392 #[test]
393 fn the_empty_pack_is_a_pack_of_nothing() {
394 let pack = g1t_scan::pack::write_pack(&[]);
395 assert_eq!(pack, EMPTY_PACK);
396 }
Rust repos service with shipping; pull requests kept in the model397}

This file's history is long; its oldest lines are credited to the oldest commit read.