Skip to content
986 linesCodeBlameRaw
1//! Copying every branch and tag from one git server to another: importing a
2//! repository with a credential, keeping a mirror in step with the host it
3//! mirrors, pushing a repository's refs out to a host it is mirrored to,
4//! and moving chosen refs either way when a takeover is handed back (see
5//! `g1t_contracts::mirrors`).
6//!
7//! Like landing (see `land.rs`), this speaks git's smart HTTP protocol and
8//! relays the pack it receives unchanged. Both sides are asked for their
9//! refs; the source is asked for one pack holding what the target lacks;
10//! the target is sent one push that moves every ref that differs.
11//!
12//! Nothing a mirror held is lost silently: when catching up moves a branch
13//! somewhere its old commit is not part of (a force-push or deletion on the
14//! remote), the old commit is kept under `refs/g1t/replaced/`, for at
15//! least [`REPLACED_KEPT_DAYS`] days.
16
17use std::collections::{BTreeMap, HashSet};
18use std::fmt;
19
20use futures_util::StreamExt;
21use worker::js_sys::Uint8Array;
22use worker::{Fetch, Headers, Method, Request, RequestInit, Result};
23
24use g1t_contracts::repos::{
25 GitAccess, MirrorApplied, MirrorApplyArgs, MirrorArgs, MirrorDirection, MirrorRefs, MirrorRefsArgs, Mirrored, RefMoved,
26 Repo, SetMirrorArgs,
27};
28use g1t_contracts::{FailureCode, Outcome};
29use g1t_kit::now_ms;
30
31use crate::land::{push_ref, read_pkt_lines, unpack_sideband};
32use crate::registry::store_key;
33use crate::store::{GitRepo, GitStore, Scope};
34use crate::{Repos, descends_from, import, not_found};
35
36const ZERO_ID: &str = "0000000000000000000000000000000000000000";
37/// The most a pack may hold. A Worker holds it in memory while relaying it.
38pub const MAX_PACK_BYTES: usize = 40 * 1024 * 1024;
39/// How many of the target's commits are named to the source as already
40/// had, so a sync only carries what is new.
41const MAX_HAVES: usize = 256;
42const USER_AGENT: &str = "git/2.45.0 (g1t mirror)";
43/// Where a mirror keeps a commit the remote stopped pointing at.
44pub const REPLACED_PREFIX: &str = "refs/g1t/replaced/";
45/// How long a kept commit stays, at least.
46pub const REPLACED_KEPT_DAYS: u64 = 30;
47/// How far back a moved branch is searched for its old commit.
48const REPLACED_HISTORY: u32 = 200;
49/// A pack with no objects: for a push whose refs all name objects the
50/// target already holds. The last 20 bytes are the SHA-1 of the first 12.
51const EMPTY_PACK: [u8; 32] = [
52 b'P', b'A', b'C', b'K', 0, 0, 0, 2, 0, 0, 0, 0, 0x02, 0x9d, 0x08, 0x82, 0x3b, 0xd8, 0xa8, 0xea,
53 0xb5, 0x10, 0xad, 0x6a, 0xc7, 0x5c, 0x82, 0x3c, 0xfd, 0x3e, 0xd3, 0x1e,
54];
55
56/// One side of a copy: a repository's smart HTTP address and the
57/// `authorization` header that opens it. The header is never logged.
58pub struct Endpoint {
59 pub url: String,
60 pub authorization: String,
61}
62
63impl Endpoint {
64 /// A g1t repository, with a credential from the git store.
65 pub fn bearer(url: &str, token: &str) -> Self {
66 Endpoint {
67 url: url.trim_end_matches('/').to_owned(),
68 authorization: format!("Bearer {token}"),
69 }
70 }
71
72 /// A public repository anywhere, read with no credential: importing
73 /// one copies every branch and tag, as with a credential.
74 pub fn anonymous(url: &str) -> Self {
75 Endpoint {
76 url: url.trim_end_matches('/').to_owned(),
77 authorization: String::new(),
78 }
79 }
80
81 /// A GitHub repository, with an installation access token. GitHub takes
82 /// one as the password of the user `x-access-token`. The token is used
83 /// as given: its length and shape are GitHub's to change.
84 pub fn github(url: &str, token: &str) -> Self {
85 Self::basic(url, GITHUB_USER, token)
86 }
87
88 /// Any https git host, with a token sent as `username`'s password.
89 pub fn basic(url: &str, username: &str, token: &str) -> Self {
90 Endpoint {
91 url: url.trim_end_matches('/').to_owned(),
92 authorization: format!("Basic {}", base64(&format!("{username}:{token}"))),
93 }
94 }
95
96 /// A remote named in a service's arguments: GitHub's user unless
97 /// another is given.
98 pub fn remote(url: &str, username: Option<&str>, token: &str) -> Self {
99 Self::basic(url, username.filter(|name| !name.trim().is_empty()).unwrap_or(GITHUB_USER), token)
100 }
101}
102
103const GITHUB_USER: &str = "x-access-token";
104
105/// Why a copy did not happen.
106#[derive(Clone, Debug, PartialEq, Eq)]
107pub enum Problem {
108 /// The other host did not answer, or answered with a server error: it
109 /// may be down.
110 Unreachable(String),
111 /// It answered, and refused or could not be used.
112 Refused(String),
113}
114
115impl Problem {
116 pub fn code(&self) -> FailureCode {
117 match self {
118 Problem::Unreachable(_) => FailureCode::Unavailable,
119 Problem::Refused(_) => FailureCode::Conflict,
120 }
121 }
122
123 pub fn message(&self) -> &str {
124 match self {
125 Problem::Unreachable(message) | Problem::Refused(message) => message,
126 }
127 }
128
129 fn map(self, wrap: impl Fn(&str) -> String) -> Self {
130 match self {
131 Problem::Unreachable(message) => Problem::Unreachable(wrap(&message)),
132 Problem::Refused(message) => Problem::Refused(wrap(&message)),
133 }
134 }
135}
136
137impl fmt::Display for Problem {
138 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
139 f.write_str(self.message())
140 }
141}
142
143/// What an answer's status says about the host: a server error, a timeout
144/// or a rate limit means it may be down; anything else, that it answered.
145pub fn problem_for_status(status: u16, what: impl Into<String>) -> Problem {
146 if status >= 500 || status == 408 || status == 429 {
147 Problem::Unreachable(what.into())
148 } else {
149 Problem::Refused(what.into())
150 }
151}
152
153/// What a copy changed.
154#[derive(Debug, Default, PartialEq, Eq)]
155pub struct Copied {
156 /// Refs created or moved, each with where it was and where it is now.
157 pub updated: Vec<(String, Option<String>, String)>,
158 pub deleted: Vec<String>,
159 /// What each deleted ref pointed at.
160 pub deleted_from: BTreeMap<String, String>,
161 /// The branch the source's HEAD names, when it said.
162 pub head: Option<String>,
163}
164
165/// Which refs a copy moves.
166#[derive(Clone, Copy, Debug, PartialEq, Eq)]
167pub enum Prune {
168 /// Refs the source no longer has are deleted from the target: the
169 /// target is a mirror of the source.
170 Yes,
171 /// Refs only the target has are left alone: pushing out to a host that
172 /// may have branches of its own.
173 No,
174}
175
176fn base64(text: &str) -> String {
177 const ALPHABET: &[u8; 64] = b"ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789+/";
178 let bytes = text.as_bytes();
179 let mut out = String::with_capacity(bytes.len().div_ceil(3) * 4);
180 for chunk in bytes.chunks(3) {
181 let n = (chunk[0] as u32) << 16
182 | (*chunk.get(1).unwrap_or(&0) as u32) << 8
183 | *chunk.get(2).unwrap_or(&0) as u32;
184 for (i, shift) in [18, 12, 6, 0].into_iter().enumerate() {
185 if i <= chunk.len() {
186 out.push(ALPHABET[(n >> shift) as usize & 63] as char);
187 } else {
188 out.push('=');
189 }
190 }
191 }
192 out
193}
194
195fn pkt_line(payload: &str) -> Vec<u8> {
196 format!("{:04x}{payload}", payload.len() + 4).into_bytes()
197}
198
199/// A ref advertisement: branch and tag refs by name, and the branch HEAD
200/// points to when the server said.
201#[derive(Debug, Default, PartialEq, Eq)]
202pub struct Advertised {
203 pub refs: BTreeMap<String, String>,
204 pub head: Option<String>,
205}
206
207impl Advertised {
208 /// The branch to make a copy's default: the one HEAD names, else
209 /// `main`, else the first branch. `None` for a repository with none.
210 pub fn default_branch(&self) -> Option<String> {
211 let branches: Vec<&str> = self.refs.keys().filter_map(|name| name.strip_prefix("refs/heads/")).collect();
212 self.head
213 .clone()
214 .filter(|head| branches.contains(&head.as_str()))
215 .or_else(|| branches.iter().find(|name| **name == "main").map(|name| name.to_string()))
216 .or_else(|| branches.first().map(|name| name.to_string()))
217 }
218}
219
220/// Asks a source for its refs before anything is made from it, so an
221/// address or credential that does not work leaves nothing behind.
222pub async fn probe(source: &Endpoint) -> Result<std::result::Result<Advertised, String>> {
223 Ok(advertise(source, "git-upload-pack").await?.map_err(|problem| problem.to_string()))
224}
225
226/// Reads `info/refs`. Peeled tags (`^{}`) and the placeholder an empty
227/// repository sends are skipped; only branches and tags are kept.
228pub fn parse_advertisement(bytes: &[u8]) -> Advertised {
229 let (lines, _) = read_pkt_lines(bytes);
230 let mut advertised = Advertised::default();
231 for line in lines {
232 let mut parts = line.splitn(2, |byte| *byte == 0);
233 let Some(reference) = parts.next().and_then(|bytes| std::str::from_utf8(bytes).ok()) else {
234 continue;
235 };
236 if let Some(capabilities) = parts.next().and_then(|bytes| std::str::from_utf8(bytes).ok()) {
237 advertised.head = capabilities
238 .split(' ')
239 .find_map(|capability| capability.trim().strip_prefix("symref=HEAD:refs/heads/"))
240 .map(str::to_owned)
241 .or(advertised.head);
242 }
243 let Some((hash, name)) = reference.trim_end().split_once(' ') else {
244 continue;
245 };
246 let kept = (name.starts_with("refs/heads/") || name.starts_with("refs/tags/")) && !name.ends_with("^{}");
247 if kept && hash.len() == 40 && hash != ZERO_ID {
248 advertised.refs.insert(name.to_owned(), hash.to_owned());
249 }
250 }
251 advertised
252}
253
254/// One ref command for a push: `(name, old, new)`, `new` the zero id to
255/// delete.
256pub type Command = (String, Option<String>, String);
257
258/// What has to change on `target` so its refs match `source`'s.
259pub fn plan(source: &Advertised, target: &Advertised, prune: Prune) -> Vec<Command> {
260 let mut commands = Vec::new();
261 for (name, hash) in &source.refs {
262 let old = target.refs.get(name);
263 if old != Some(hash) {
264 commands.push((name.clone(), old.cloned(), hash.clone()));
265 }
266 }
267 if prune == Prune::Yes {
268 for (name, hash) in &target.refs {
269 if !source.refs.contains_key(name) {
270 commands.push((name.clone(), Some(hash.clone()), ZERO_ID.to_owned()));
271 }
272 }
273 }
274 commands
275}
276
277/// The body of an upload-pack request for `wants`, naming `haves`.
278pub fn upload_request(wants: &[String], haves: &[String]) -> Vec<u8> {
279 let mut body = Vec::new();
280 for (index, want) in wants.iter().enumerate() {
281 // Capabilities ride on the first want only.
282 let line = if index == 0 { format!("want {want} side-band-64k\n") } else { format!("want {want}\n") };
283 body.extend(pkt_line(&line));
284 }
285 body.extend_from_slice(b"0000");
286 for have in haves.iter().take(MAX_HAVES) {
287 body.extend(pkt_line(&format!("have {have}\n")));
288 }
289 body.extend(pkt_line("done\n"));
290 body
291}
292
293/// The body of a receive-pack request: the commands, then the pack when
294/// anything but deletions is sent.
295pub fn receive_request(commands: &[Command], pack: Option<Vec<u8>>) -> Vec<u8> {
296 let deletes = commands.iter().any(|(_, _, new)| new == ZERO_ID);
297 let capabilities = if deletes { "report-status delete-refs" } else { "report-status" };
298 let mut body = Vec::new();
299 for (index, (name, old, new)) in commands.iter().enumerate() {
300 let old = old.as_deref().unwrap_or(ZERO_ID);
301 let line = if index == 0 { format!("{old} {new} {name}\0 {capabilities}\n") } else { format!("{old} {new} {name}\n") };
302 body.extend(pkt_line(&line));
303 }
304 body.extend_from_slice(b"0000");
305 if commands.iter().any(|(_, _, new)| new != ZERO_ID) {
306 body.extend(pack.unwrap_or_else(|| EMPTY_PACK.to_vec()));
307 }
308 body
309}
310
311/// Which commands a receive-pack report says were refused, with why.
312pub fn refused(report: &[u8], commands: &[Command]) -> Vec<String> {
313 let (lines, _) = read_pkt_lines(report);
314 let lines: Vec<String> = lines
315 .into_iter()
316 .map(|line| {
317 // A report may come inside side-band channel 1.
318 let line = if line.first() == Some(&1) { &line[1..] } else { line };
319 String::from_utf8_lossy(line).trim_end().to_owned()
320 })
321 .collect();
322 let mut problems: Vec<String> = lines
323 .iter()
324 .filter(|line| line.starts_with("unpack ") && *line != "unpack ok")
325 .cloned()
326 .collect();
327 for (name, _, _) in commands {
328 if !lines.iter().any(|line| *line == format!("ok {name}")) {
329 let reason = lines
330 .iter()
331 .find_map(|line| line.strip_prefix(&format!("ng {name} ")))
332 .unwrap_or("not reported");
333 problems.push(format!("{name}: {reason}"));
334 }
335 }
336 problems
337}
338
339fn request(method: Method, url: &str, endpoint: &Endpoint, body: Option<(&str, Vec<u8>)>) -> Result<Request> {
340 // Requests to g1t's own store are metered (meters.rs); GitHub's are not.
341 if endpoint.authorization.starts_with("Bearer ") {
342 let sent = body.as_ref().map_or(0, |(_, body)| body.len() as u64);
343 let meter = match body.as_ref().map(|(service, _)| *service) {
344 Some("git-receive-pack") => "internal.git.receive_pack",
345 Some(_) => "internal.git.fetch",
346 None => "internal.git.info_refs",
347 };
348 crate::meters::record_remote(meter, &endpoint.url, sent, 0);
349 }
350 let headers = Headers::new();
351 headers.set("user-agent", USER_AGENT)?;
352 if !endpoint.authorization.is_empty() {
353 headers.set("authorization", &endpoint.authorization)?;
354 }
355 let mut init = RequestInit::new();
356 if let Some((service, body)) = body {
357 headers.set("content-type", &format!("application/x-{service}-request"))?;
358 headers.set("accept", &format!("application/x-{service}-result"))?;
359 init.with_body(Some(Uint8Array::from(body.as_slice()).into()));
360 }
361 init.with_method(method).with_headers(headers);
362 Request::new_with_init(url, &init)
363}
364
365/// Asks a server for its refs, as the given service would see them.
366async fn advertise(endpoint: &Endpoint, service: &str) -> Result<std::result::Result<Advertised, Problem>> {
367 Ok(advertise_raw(endpoint, service).await?.map(|bytes| parse_advertisement(&bytes)))
368}
369
370/// `info/refs` as sent.
371async fn advertise_raw(endpoint: &Endpoint, service: &str) -> Result<std::result::Result<Vec<u8>, Problem>> {
372 let url = format!("{}/info/refs?service={service}", endpoint.url);
373 let mut response = match Fetch::Request(request(Method::Get, &url, endpoint, None)?).send().await {
374 Ok(response) => response,
375 Err(_) => return Ok(Err(Problem::Unreachable("The repository could not be reached.".to_owned()))),
376 };
377 match response.status_code() {
378 200 => Ok(Ok(response.bytes().await?)),
379 401 | 403 => Ok(Err(Problem::Refused("The repository refused g1t's credential.".to_owned()))),
380 404 => Ok(Err(Problem::Refused("The repository was not found, or g1t may not read it.".to_owned()))),
381 status => Ok(Err(problem_for_status(status, format!("The repository answered {status}.")))),
382 }
383}
384
385/// The kept commits in an advertisement: `refs/g1t/replaced/...` names
386/// with what they point at.
387pub fn replaced_refs(bytes: &[u8]) -> Vec<(String, String)> {
388 let (lines, _) = read_pkt_lines(bytes);
389 lines
390 .into_iter()
391 .filter_map(|line| {
392 let line = std::str::from_utf8(line.split(|byte| *byte == 0).next()?).ok()?;
393 let (hash, name) = line.trim_end().split_once(' ')?;
394 (name.starts_with(REPLACED_PREFIX) && hash.len() == 40).then(|| (name.to_owned(), hash.to_owned()))
395 })
396 .collect()
397}
398
399/// The ref that keeps `old` when `name` stops pointing at it, stamped with
400/// the time in milliseconds so expired ones are found from the name alone.
401pub fn replaced_name(name: &str, now_ms: u64) -> String {
402 format!("{REPLACED_PREFIX}{}/{now_ms}", name.trim_start_matches("refs/"))
403}
404
405/// The ref a kept commit was replaced on: `refs/heads/main` for
406/// `refs/g1t/replaced/heads/main/<ms>`.
407pub fn replaced_from(kept: &str) -> Option<String> {
408 let rest = kept.strip_prefix(REPLACED_PREFIX)?;
409 let (name, _) = rest.rsplit_once('/')?;
410 Some(format!("refs/{name}"))
411}
412
413/// Whether a kept commit's ref is older than [`REPLACED_KEPT_DAYS`].
414pub fn replaced_expired(name: &str, now_ms: u64) -> bool {
415 name.rsplit('/')
416 .next()
417 .and_then(|stamp| stamp.parse::<u64>().ok())
418 .is_some_and(|stamp| now_ms.saturating_sub(stamp) > REPLACED_KEPT_DAYS * 86_400_000)
419}
420
421/// Fetches one pack holding `wants`, reading in pieces so that one too
422/// large is noticed before it is all held.
423async fn fetch_pack(source: &Endpoint, wants: &[String], haves: &[String]) -> Result<std::result::Result<Vec<u8>, Problem>> {
424 let body = upload_request(wants, haves);
425 let url = format!("{}/git-upload-pack", source.url);
426 let Ok(mut response) = Fetch::Request(request(Method::Post, &url, source, Some(("git-upload-pack", body)))?)
427 .send()
428 .await
429 else {
430 return Ok(Err(Problem::Unreachable("The repository stopped answering while sending.".to_owned())));
431 };
432 if response.status_code() != 200 {
433 let status = response.status_code();
434 return Ok(Err(problem_for_status(status, format!("The source refused to send the repository ({status})."))));
435 }
436 let mut received = Vec::new();
437 let mut stream = response.stream()?;
438 while let Some(chunk) = stream.next().await {
439 received.extend_from_slice(&chunk?);
440 if received.len() > MAX_PACK_BYTES {
441 return Ok(Err(Problem::Refused(format!(
442 "The change is larger than {} MB, the most g1t copies at once. Push it with git instead.",
443 MAX_PACK_BYTES / 1024 / 1024
444 ))));
445 }
446 }
447 Ok(unpack_sideband(&received).map_err(|_| Problem::Refused("The source did not send a usable pack.".to_owned())))
448}
449
450/// Makes `target`'s branches and tags match `source`'s. With
451/// `Prune::No`, refs only the target has are kept. `Err` in the inner
452/// result is a reason to show a person.
453pub async fn copy(source: &Endpoint, target: &Endpoint, prune: Prune) -> Result<std::result::Result<Copied, Problem>> {
454 let theirs = match advertise(source, "git-upload-pack").await? {
455 Ok(refs) => refs,
456 Err(problem) => return Ok(Err(problem)),
457 };
458 let ours = match advertise(target, "git-receive-pack").await? {
459 Ok(refs) => refs,
460 Err(problem) => return Ok(Err(problem.map(|reason| format!("Writing the copy failed: {reason}")))),
461 };
462 let commands = plan(&theirs, &ours, prune);
463 let mut copied = Copied {
464 head: theirs.head.clone(),
465 ..Copied::default()
466 };
467 if commands.is_empty() {
468 return Ok(Ok(copied));
469 }
470 let results = match transfer(source, target, &ours, &commands).await? {
471 Ok(results) => results,
472 Err(problem) => return Ok(Err(problem)),
473 };
474 let problems: Vec<String> = results
475 .iter()
476 .filter_map(|(name, problem)| problem.as_ref().map(|problem| format!("{name}: {problem}")))
477 .collect();
478 if !problems.is_empty() {
479 return Ok(Err(Problem::Refused(format!("Some refs were not updated: {}", problems.join("; ")))));
480 }
481 for (name, old, new) in commands {
482 if new == ZERO_ID {
483 if let Some(old) = old {
484 copied.deleted_from.insert(name.clone(), old);
485 }
486 copied.deleted.push(name);
487 } else {
488 copied.updated.push((name, old, new));
489 }
490 }
491 Ok(Ok(copied))
492}
493
494/// Sends `commands` to `target` in one push, with a pack from `source`
495/// holding what the target lacks (`held`: what its refs already name).
496/// Each command comes back with why it was refused, if it was.
497async fn transfer(
498 source: &Endpoint,
499 target: &Endpoint,
500 held: &Advertised,
501 commands: &[Command],
502) -> Result<std::result::Result<Vec<(String, Option<String>)>, Problem>> {
503 let have: HashSet<&String> = held.refs.values().collect();
504 let mut wants: Vec<String> = commands
505 .iter()
506 .map(|(_, _, new)| new)
507 .filter(|new| *new != ZERO_ID && !have.contains(new))
508 .cloned()
509 .collect();
510 wants.sort();
511 wants.dedup();
512 let pack = if wants.is_empty() {
513 None
514 } else {
515 let haves: Vec<String> = held.refs.values().cloned().collect::<HashSet<_>>().into_iter().collect();
516 match fetch_pack(source, &wants, &haves).await? {
517 Ok(pack) => Some(pack),
518 Err(problem) => return Ok(Err(problem)),
519 }
520 };
521 let body = receive_request(commands, pack);
522 let url = format!("{}/git-receive-pack", target.url);
523 let Ok(mut response) = Fetch::Request(request(Method::Post, &url, target, Some(("git-receive-pack", body)))?)
524 .send()
525 .await
526 else {
527 return Ok(Err(Problem::Unreachable("The repository stopped answering while receiving.".to_owned())));
528 };
529 let report = response.bytes().await?;
530 if response.status_code() != 200 {
531 let status = response.status_code();
532 return Ok(Err(problem_for_status(status, format!("The push was refused ({status})."))));
533 }
534 Ok(Ok(outcomes(&report, commands)))
535}
536
537/// Each command with why it was refused, if the report says it was.
538pub fn outcomes(report: &[u8], commands: &[Command]) -> Vec<(String, Option<String>)> {
539 let refused = refused(report, commands);
540 commands
541 .iter()
542 .map(|(name, _, _)| {
543 let prefix = format!("{name}: ");
544 (name.clone(), refused.iter().find_map(|line| line.strip_prefix(&prefix).map(str::to_owned)))
545 })
546 .collect()
547}
548
549/// The pushes an import announces, one for each ref it made: the default
550/// branch first (what follows a repository, such as its Composer package,
551/// starts from it), then the other branches, then the tags.
552pub fn import_pushes(updated: &[(String, Option<String>, String)], default_branch: &str) -> Vec<(String, String)> {
553 let default = format!("refs/heads/{default_branch}");
554 let rank = |name: &str| {
555 if name == default {
556 0
557 } else if name.starts_with("refs/heads/") {
558 1
559 } else {
560 2
561 }
562 };
563 let mut pushes: Vec<(String, String)> = updated
564 .iter()
565 .filter(|(name, _, new)| (name.starts_with("refs/heads/") || name.starts_with("refs/tags/")) && new != ZERO_ID)
566 .map(|(name, _, new)| (name.clone(), new.clone()))
567 .collect();
568 pushes.sort_by(|a, b| rank(&a.0).cmp(&rank(&b.0)).then_with(|| a.0.cmp(&b.0)));
569 pushes
570}
571
572impl<S: GitStore> Repos<S> {
573 /// `mirror`: a mirror catching up with the host it mirrors, or a
574 /// repository pushing its refs out to one. Each branch moved on g1t is
575 /// announced as a push copied in, so what follows a mirror's state
576 /// (workflows, deployments) can tell.
577 pub(crate) async fn mirror(&self, a: MirrorArgs) -> Result<Outcome<Mirrored>> {
578 let Some(repo) = self.registry.by_id(&a.repo_id).await? else {
579 return Ok(not_found());
580 };
581 let Some(url) = import::clean_url(&a.url) else {
582 return Ok(Outcome::fail(FailureCode::Invalid, "That is not an https repository address."));
583 };
584 // Catching up writes: it waits for a move between namespaces (moves.rs).
585 let repo = if a.direction == MirrorDirection::Pull {
586 match self.unpaused(repo).await? {
587 Ok(repo) => repo,
588 Err((code, message)) => return Ok(Outcome::fail(code, message)),
589 }
590 } else {
591 repo
592 };
593 let scope = match a.direction {
594 MirrorDirection::Pull => Scope::Write,
595 MirrorDirection::Push => Scope::Read,
596 };
597 let access = self.store.open(&store_key(&repo)).await?.access(scope).await?;
598 let ours = Endpoint::bearer(&access.remote, &access.token);
599 let theirs = Endpoint::remote(&url, a.username.as_deref(), &a.token);
600 let copied = match a.direction {
601 MirrorDirection::Pull => {
602 let copied = copy(&theirs, &ours, Prune::Yes).await?;
603 self.refs_moved(&repo.id).await;
604 copied
605 }
606 MirrorDirection::Push => copy(&ours, &theirs, Prune::No).await?,
607 };
608 let copied = match copied {
609 Ok(copied) => copied,
610 Err(problem) => return Ok(Outcome::fail(problem.code(), problem.to_string())),
611 };
612 let mut replaced = Vec::new();
613 if a.direction == MirrorDirection::Pull {
614 let moved: Vec<(String, Option<String>, String)> = copied
615 .updated
616 .iter()
617 .cloned()
618 .chain(
619 copied
620 .deleted_from
621 .iter()
622 .map(|(name, old)| (name.clone(), Some(old.clone()), ZERO_ID.to_owned())),
623 )
624 .collect();
625 replaced = self.keep_replaced(&repo, &access, &moved).await?;
626 for (name, old, new) in &copied.updated {
627 if name.starts_with("refs/heads/") {
628 self.publish_mirrored_push(&repo, name, old.as_deref(), new).await?;
629 }
630 }
631 }
632 Ok(Outcome::Ok(Mirrored {
633 updated: copied.updated.into_iter().map(|(name, _, _)| name).collect(),
634 deleted: copied.deleted,
635 replaced,
636 }))
637 }
638
639 /// Keeps each commit a pull stopped pointing at, when the ref's new
640 /// commit does not contain it: a force-push or deletion on the remote.
641 /// Returns the refs that keep them. Expired ones are let go at the same
642 /// time.
643 async fn keep_replaced(&self, repo: &Repo, access: &GitAccess, moved: &[(String, Option<String>, String)]) -> Result<Vec<String>> {
644 let git = self.store.open(&store_key(repo)).await?;
645 let now = now_ms();
646 let mut kept = Vec::new();
647 for (name, old, new) in moved {
648 let Some(old) = old.as_deref() else { continue };
649 let lost = if new == ZERO_ID || name.starts_with("refs/tags/") {
650 true
651 } else {
652 match git.log(name, REPLACED_HISTORY).await {
653 Ok(history) => !descends_from(&git, &history, old).await.unwrap_or(true),
654 Err(_) => false,
655 }
656 };
657 if !lost {
658 continue;
659 }
660 let keep = replaced_name(name, now);
661 if let Ok(Ok(())) = push_ref(access, &keep, None, old, Some(EMPTY_PACK.to_vec())).await {
662 kept.push(keep);
663 }
664 }
665 if !kept.is_empty() {
666 let ours = Endpoint::bearer(&access.remote, &access.token);
667 if let Ok(Ok(bytes)) = advertise_raw(&ours, "git-receive-pack").await {
668 for (name, hash) in replaced_refs(&bytes).into_iter().filter(|(name, _)| replaced_expired(name, now)) {
669 let _ = push_ref(access, &name, Some(&hash), ZERO_ID, None).await;
670 }
671 }
672 self.refs_moved(&repo.id).await;
673 }
674 Ok(kept)
675 }
676
677 /// `mirror_refs`: both sides' branches and tags.
678 pub(crate) async fn mirror_refs(&self, a: MirrorRefsArgs) -> Result<Outcome<MirrorRefs>> {
679 let Some(repo) = self.registry.by_id(&a.repo_id).await? else {
680 return Ok(not_found());
681 };
682 let access = self.store.open(&store_key(&repo)).await?.access(Scope::Read).await?;
683 let ours = match advertise(&Endpoint::bearer(&access.remote, &access.token), "git-upload-pack").await? {
684 Ok(ours) => ours,
685 Err(problem) => return Ok(Outcome::fail(FailureCode::Conflict, problem.to_string())),
686 };
687 // No address: only g1t's side was asked for.
688 if a.url.is_empty() {
689 return Ok(Outcome::Ok(MirrorRefs { ours: ours.refs, theirs: None, unreachable: None }));
690 }
691 let Some(url) = import::clean_url(&a.url) else {
692 return Ok(Outcome::fail(FailureCode::Invalid, "That is not an https repository address."));
693 };
694 let (theirs, unreachable) = match advertise(&Endpoint::remote(&url, a.username.as_deref(), &a.token), "git-upload-pack").await? {
695 Ok(theirs) => (Some(theirs.refs), None),
696 Err(Problem::Unreachable(reason)) => (None, Some(reason)),
697 Err(problem) => return Ok(Outcome::fail(problem.code(), problem.to_string())),
698 };
699 Ok(Outcome::Ok(MirrorRefs {
700 ours: ours.refs,
701 theirs,
702 unreachable,
703 }))
704 }
705
706 /// `mirror_apply`: moves chosen refs either way, each only from the
707 /// commit the caller saw. Pushes go out in one request, pulls come in
708 /// in another.
709 pub(crate) async fn mirror_apply(&self, a: MirrorApplyArgs) -> Result<Outcome<MirrorApplied>> {
710 let Some(repo) = self.registry.by_id(&a.repo_id).await? else {
711 return Ok(not_found());
712 };
713 let Some(url) = import::clean_url(&a.url) else {
714 return Ok(Outcome::fail(FailureCode::Invalid, "That is not an https repository address."));
715 };
716 let repo = match self.unpaused(repo).await? {
717 Ok(repo) => repo,
718 Err((code, message)) => return Ok(Outcome::fail(code, message)),
719 };
720 let access = self.store.open(&store_key(&repo)).await?.access(Scope::Write).await?;
721 let ours = Endpoint::bearer(&access.remote, &access.token);
722 let theirs = Endpoint::remote(&url, a.username.as_deref(), &a.token);
723 let mut moved = Vec::new();
724 for direction in [MirrorDirection::Push, MirrorDirection::Pull] {
725 let moves: Vec<_> = a.moves.iter().filter(|m| m.direction == direction).collect();
726 if moves.is_empty() {
727 continue;
728 }
729 let commands: Vec<Command> = moves
730 .iter()
731 .map(|m| {
732 (
733 m.to.clone().unwrap_or_else(|| m.git_ref.clone()),
734 m.old.clone(),
735 m.new.clone().unwrap_or_else(|| ZERO_ID.to_owned()),
736 )
737 })
738 .collect();
739 let (source, target) = match direction {
740 MirrorDirection::Push => (&ours, &theirs),
741 MirrorDirection::Pull => (&theirs, &ours),
742 };
743 let held = match advertise(target, "git-receive-pack").await? {
744 Ok(held) => held,
745 Err(problem) => return Ok(Outcome::fail(problem.code(), problem.to_string())),
746 };
747 let results = match transfer(source, target, &held, &commands).await? {
748 Ok(results) => results,
749 Err(problem) => return Ok(Outcome::fail(problem.code(), problem.to_string())),
750 };
751 let mut landed = Vec::new();
752 for ((name, problem), command) in results.into_iter().zip(&commands) {
753 if problem.is_none() {
754 landed.push(command.clone());
755 }
756 moved.push(RefMoved {
757 git_ref: name,
758 direction,
759 problem,
760 replaced: None,
761 });
762 }
763 if direction == MirrorDirection::Pull && !landed.is_empty() {
764 self.refs_moved(&repo.id).await;
765 for keep in self.keep_replaced(&repo, &access, &landed).await? {
766 let from = replaced_from(&keep);
767 if let Some(entry) = moved
768 .iter_mut()
769 .find(|m| m.direction == direction && Some(&m.git_ref) == from.as_ref())
770 {
771 entry.replaced = Some(keep);
772 }
773 }
774 for (name, old, new) in &landed {
775 if name.starts_with("refs/heads/") && new != ZERO_ID {
776 self.publish_mirrored_push(&repo, name, old.as_deref(), new).await?;
777 }
778 }
779 }
780 }
781 Ok(Outcome::Ok(MirrorApplied { moved }))
782 }
783
784 /// `set_mirror`: what the integrations service decided about a
785 /// repository's mirror, kept on its row.
786 pub(crate) async fn set_mirror(&self, a: SetMirrorArgs) -> Result<Outcome<Repo>> {
787 let Some(mut repo) = self.registry.by_id(&a.repo_id).await? else {
788 return Ok(not_found());
789 };
790 self.registry.set_mirror(&repo.id, a.mirror.as_ref()).await?;
791 repo.mirror = a.mirror;
792 Ok(Outcome::Ok(repo))
793 }
794}
795
796#[cfg(test)]
797mod tests {
798 use super::*;
799
800 const A: &str = "c71546fcd893ef8b0f57388b65e620d759705dda";
801 const B: &str = "4807077b296e6edbf410d55e72749d3e1170c291";
802 const C: &str = "1111111111111111111111111111111111111111";
803
804 fn advertised(refs: &[(&str, &str)]) -> Advertised {
805 Advertised {
806 refs: refs.iter().map(|(name, hash)| (name.to_string(), hash.to_string())).collect(),
807 head: None,
808 }
809 }
810
811 #[test]
812 fn the_basic_header_carries_the_token_whole() {
813 assert_eq!(base64("x-access-token:abc"), "eC1hY2Nlc3MtdG9rZW46YWJj");
814 assert_eq!(base64("a"), "YQ==");
815 assert_eq!(base64("ab"), "YWI=");
816 // GitHub's stateless installation tokens run to hundreds of
817 // characters; nothing here assumes a length.
818 let token = format!("ghs_{}", "x".repeat(516));
819 assert_eq!(token.len(), 520);
820 let endpoint = Endpoint::github("https://github.com/o/r.git/", &token);
821 assert_eq!(endpoint.url, "https://github.com/o/r.git");
822 let encoded = endpoint.authorization.strip_prefix("Basic ").unwrap();
823 assert_eq!(encoded, base64(&format!("x-access-token:{token}")));
824 assert_eq!(encoded.len(), (15 + 520usize).div_ceil(3) * 4);
825 }
826
827 #[test]
828 fn the_advertisement_keeps_branches_and_tags() {
829 let bytes = [
830 pkt_line("# service=git-upload-pack\n"),
831 b"0000".to_vec(),
832 pkt_line(&format!("{A} HEAD\0multi_ack side-band-64k symref=HEAD:refs/heads/trunk agent=git/x\n")),
833 pkt_line(&format!("{A} refs/heads/trunk\n")),
834 pkt_line(&format!("{B} refs/tags/v1\n")),
835 pkt_line(&format!("{A} refs/tags/v1^{{}}\n")),
836 pkt_line(&format!("{C} refs/pull/1/head\n")),
837 b"0000".to_vec(),
838 ]
839 .concat();
840 let parsed = parse_advertisement(&bytes);
841 assert_eq!(parsed.head.as_deref(), Some("trunk"));
842 assert_eq!(parsed, Advertised {
843 refs: advertised(&[("refs/heads/trunk", A), ("refs/tags/v1", B)]).refs,
844 head: Some("trunk".to_owned()),
845 });
846 assert_eq!(parsed.default_branch().as_deref(), Some("trunk"));
847 assert_eq!(advertised(&[("refs/heads/a", A), ("refs/heads/main", A)]).default_branch().as_deref(), Some("main"));
848 assert_eq!(advertised(&[("refs/tags/v1", A)]).default_branch(), None);
849 let empty = [pkt_line(&format!("{ZERO_ID} capabilities^{{}}\0report-status\n")), b"0000".to_vec()].concat();
850 assert!(parse_advertisement(&empty).refs.is_empty());
851 }
852
853 #[test]
854 fn a_public_import_copies_every_branch_and_tag_annotated_ones_whole() {
855 // An empty repository being filled from a public source: every
856 // branch and tag is made, an annotated tag as its tag object (so it
857 // stays annotated), and the source's pull request refs and peeled
858 // lines are left out.
859 let bytes = [
860 pkt_line("# service=git-upload-pack\n"),
861 b"0000".to_vec(),
862 pkt_line(&format!("{A} HEAD\0multi_ack symref=HEAD:refs/heads/master\n")),
863 pkt_line(&format!("{A} refs/heads/master\n")),
864 pkt_line(&format!("{B} refs/heads/next\n")),
865 pkt_line(&format!("{C} refs/pull/7/head\n")),
866 pkt_line(&format!("{C} refs/tags/1.0.0\n")),
867 pkt_line(&format!("{A} refs/tags/1.0.0^{{}}\n")),
868 pkt_line(&format!("{B} refs/tags/2.0.0\n")),
869 b"0000".to_vec(),
870 ]
871 .concat();
872 let source = parse_advertisement(&bytes);
873 assert_eq!(source.default_branch().as_deref(), Some("master"), "HEAD stays the default branch");
874 let commands = plan(&source, &Advertised::default(), Prune::Yes);
875 let names: Vec<(&str, &str)> = commands.iter().map(|(name, _, new)| (name.as_str(), new.as_str())).collect();
876 assert_eq!(
877 names,
878 [("refs/heads/master", A), ("refs/heads/next", B), ("refs/tags/1.0.0", C), ("refs/tags/2.0.0", B)]
879 );
880 let pushes = import_pushes(&commands, "master");
881 let order: Vec<&str> = pushes.iter().map(|(name, _)| name.as_str()).collect();
882 assert_eq!(order, ["refs/heads/master", "refs/heads/next", "refs/tags/1.0.0", "refs/tags/2.0.0"]);
883 let pushes = import_pushes(&commands, "next");
884 assert_eq!(pushes[0].0, "refs/heads/next", "the default branch is announced first");
885 assert_eq!(Endpoint::anonymous("https://github.com/php-fig/log.git/").url, "https://github.com/php-fig/log.git");
886 assert!(Endpoint::anonymous("https://x").authorization.is_empty(), "no credential is sent");
887 }
888
889 #[test]
890 fn a_mirror_prunes_and_a_push_out_does_not() {
891 let source = advertised(&[("refs/heads/main", A), ("refs/tags/v1", B)]);
892 let target = advertised(&[("refs/heads/main", B), ("refs/heads/old", C)]);
893 assert_eq!(plan(&source, &target, Prune::Yes), vec![
894 ("refs/heads/main".to_owned(), Some(B.to_owned()), A.to_owned()),
895 ("refs/tags/v1".to_owned(), None, B.to_owned()),
896 ("refs/heads/old".to_owned(), Some(C.to_owned()), ZERO_ID.to_owned()),
897 ]);
898 assert_eq!(plan(&source, &target, Prune::No).len(), 2);
899 assert!(plan(&source, &source, Prune::Yes).is_empty());
900 }
901
902 #[test]
903 fn a_push_of_known_objects_sends_an_empty_pack() {
904 let commands = vec![("refs/tags/v1".to_owned(), None, A.to_owned())];
905 let body = receive_request(&commands, None);
906 assert!(body.ends_with(&EMPTY_PACK));
907 let deletes = vec![("refs/heads/x".to_owned(), Some(A.to_owned()), ZERO_ID.to_owned())];
908 let body = receive_request(&deletes, None);
909 assert!(body.ends_with(b"0000"));
910 assert!(String::from_utf8_lossy(&body).contains("delete-refs"));
911 }
912
913 #[test]
914 fn wants_carry_capabilities_once() {
915 let body = String::from_utf8(upload_request(&[A.to_owned(), B.to_owned()], &[C.to_owned()])).unwrap();
916 assert_eq!(body.matches("side-band-64k").count(), 1);
917 assert!(body.contains(&format!("want {B}\n")));
918 assert!(body.contains(&format!("have {C}\n")));
919 assert!(body.ends_with("0009done\n"));
920 }
921
922 #[test]
923 fn replaced_commits_are_kept_by_name_and_time() {
924 let kept = replaced_name("refs/heads/feature/x", 1_000);
925 assert_eq!(kept, "refs/g1t/replaced/heads/feature/x/1000");
926 assert_eq!(replaced_from(&kept).as_deref(), Some("refs/heads/feature/x"));
927 assert!(!replaced_expired(&kept, 1_000 + REPLACED_KEPT_DAYS * 86_400_000));
928 assert!(replaced_expired(&kept, 1_001 + REPLACED_KEPT_DAYS * 86_400_000));
929 assert!(!replaced_expired("refs/g1t/replaced/heads/main/not-a-time", u64::MAX));
930 let bytes = [
931 pkt_line(&format!("{A} refs/heads/main\0report-status\n")),
932 pkt_line(&format!("{B} {kept}\n")),
933 b"0000".to_vec(),
934 ]
935 .concat();
936 assert_eq!(replaced_refs(&bytes), vec![(kept.clone(), B.to_owned())]);
937 assert!(parse_advertisement(&bytes).refs.get(&kept).is_none(), "kept commits are not branches");
938 }
939
940 #[test]
941 fn a_server_error_means_the_host_may_be_down() {
942 assert_eq!(problem_for_status(503, "x").code(), FailureCode::Unavailable);
943 assert_eq!(problem_for_status(429, "x").code(), FailureCode::Unavailable);
944 assert_eq!(problem_for_status(403, "x").code(), FailureCode::Conflict);
945 let endpoint = Endpoint::remote("https://git.example/a.git", Some("chase"), "t");
946 assert_eq!(endpoint.authorization, format!("Basic {}", base64("chase:t")));
947 assert_eq!(Endpoint::remote("https://x", Some(" "), "t").authorization, Endpoint::github("https://x", "t").authorization);
948 }
949
950 #[test]
951 fn each_command_comes_back_with_its_refusal() {
952 let commands = vec![
953 ("refs/heads/main".to_owned(), None, A.to_owned()),
954 ("refs/heads/x".to_owned(), None, B.to_owned()),
955 ];
956 let report = [
957 pkt_line("unpack ok\n"),
958 pkt_line("ok refs/heads/main\n"),
959 pkt_line("ng refs/heads/x protected branch hook declined\n"),
960 b"0000".to_vec(),
961 ]
962 .concat();
963 assert_eq!(outcomes(&report, &commands), vec![
964 ("refs/heads/main".to_owned(), None),
965 ("refs/heads/x".to_owned(), Some("protected branch hook declined".to_owned())),
966 ]);
967 }
968
969 #[test]
970 fn refusals_are_reported_by_ref() {
971 let commands = vec![
972 ("refs/heads/main".to_owned(), None, A.to_owned()),
973 ("refs/heads/x".to_owned(), None, B.to_owned()),
974 ];
975 let report = [
976 pkt_line("unpack ok\n"),
977 pkt_line("ok refs/heads/main\n"),
978 pkt_line("ng refs/heads/x protected branch\n"),
979 b"0000".to_vec(),
980 ]
981 .concat();
982 assert_eq!(refused(&report, &commands), vec!["refs/heads/x: protected branch".to_owned()]);
983 let fine = [pkt_line("unpack ok\n"), pkt_line("ok refs/heads/main\n"), pkt_line("ok refs/heads/x\n")].concat();
984 assert!(refused(&fine, &commands).is_empty());
985 }
986}