This repository has no description
1mod common;
2
3use std::collections::BTreeSet;
4use std::io::Write;
5use std::path::Path;
6use std::process::Stdio;
7use std::sync::Arc;
8
9use bytes::Bytes;
10use common::Edge;
11use http::Method;
12use knot_edge::RequiresFullHandshake;
13use knot_git::Layout;
14use knot_pack::{CacheConfig, RepoLookup, RepoResolver, RepoTarget};
15use knot_types::{ObjectFormat, RepoDid};
16
17const PINNED_DATE: &str = "2026-06-20T12:00:00+00:00";
18
19fn git(cwd: &Path, args: &[&str]) -> String {
20 let out = knot_fixtures::command(cwd)
21 .args(args)
22 .env("GIT_AUTHOR_DATE", PINNED_DATE)
23 .env("GIT_COMMITTER_DATE", PINNED_DATE)
24 .output()
25 .expect("git runs");
26 assert!(
27 out.status.success(),
28 "git {args:?} failed: {}",
29 String::from_utf8_lossy(&out.stderr)
30 );
31 String::from_utf8(out.stdout)
32 .expect("git stdout is utf-8")
33 .trim()
34 .to_string()
35}
36
37fn object_set(dir: &Path) -> BTreeSet<String> {
38 git(
39 dir,
40 &[
41 "cat-file",
42 "--batch-all-objects",
43 "--batch-check=%(objectname)",
44 ],
45 )
46 .lines()
47 .map(|line| line.trim().to_string())
48 .filter(|line| !line.is_empty())
49 .collect()
50}
51
52fn seed_bare(work: &Path, bare: &Path, format: ObjectFormat) -> String {
53 let fmt = format!("--object-format={}", format.capability());
54 std::fs::create_dir_all(work).unwrap();
55 git(work, &["init", &fmt, "-q", "-b", "main"]);
56 std::fs::write(work.join("README.md"), "h3 over the simulated quic edge\n").unwrap();
57 git(work, &["add", "-A"]);
58 git(work, &["commit", "-q", "-m", "c1"]);
59 std::fs::write(work.join("extra.txt"), "a second object\n").unwrap();
60 git(work, &["add", "-A"]);
61 git(work, &["commit", "-q", "-m", "c2"]);
62 git(work, &["push", "-q", bare.to_str().unwrap(), "main"]);
63 git(bare, &["symbolic-ref", "HEAD", "refs/heads/main"]);
64 git(work, &["rev-parse", "HEAD"])
65}
66
67fn init(dir: &Path, format: ObjectFormat) {
68 std::fs::create_dir_all(dir).unwrap();
69 let fmt = format!("--object-format={}", format.capability());
70 git(dir, &["init", &fmt, "-q", dir.to_str().unwrap()]);
71}
72
73fn index_pack(repo: &Path, pack: &[u8]) {
74 let mut child = knot_fixtures::command(repo)
75 .args(["index-pack", "--stdin", "--fix-thin"])
76 .stdin(Stdio::piped())
77 .stdout(Stdio::piped())
78 .stderr(Stdio::piped())
79 .spawn()
80 .expect("spawn index-pack");
81 child.stdin.take().unwrap().write_all(pack).unwrap();
82 let output = child.wait_with_output().unwrap();
83 assert!(
84 output.status.success(),
85 "index-pack failed: {}",
86 String::from_utf8_lossy(&output.stderr)
87 );
88}
89
90fn fetch_body(tip: &str) -> Bytes {
91 let mut body = Vec::new();
92 body.extend_from_slice(&common::pkt(b"command=fetch\n"));
93 body.extend_from_slice(&common::pkt(b"agent=knot/0\n"));
94 body.extend_from_slice(b"0001");
95 body.extend_from_slice(&common::pkt(b"no-progress\n"));
96 body.extend_from_slice(&common::pkt(b"ofs-delta\n"));
97 body.extend_from_slice(&common::pkt(format!("want {tip}\n").as_bytes()));
98 body.extend_from_slice(&common::pkt(b"done\n"));
99 body.extend_from_slice(b"0000");
100 Bytes::from(body)
101}
102
103fn extract_pack(response: &[u8]) -> Vec<u8> {
104 let mut channel = Vec::new();
105 let mut pos = 0usize;
106 while pos + 4 <= response.len() {
107 let len = std::str::from_utf8(&response[pos..pos + 4])
108 .ok()
109 .and_then(|hex| usize::from_str_radix(hex, 16).ok())
110 .unwrap_or(0);
111 pos += 4;
112 if len < 4 {
113 continue;
114 }
115 let end = (pos + len - 4).min(response.len());
116 let payload = &response[pos..end];
117 pos = end;
118 if payload.first() == Some(&1) {
119 channel.extend_from_slice(&payload[1..]);
120 }
121 }
122 match channel.windows(4).position(|window| window == b"PACK") {
123 Some(start) => channel.split_off(start),
124 None => channel,
125 }
126}
127
128fn serve_dids() -> Arc<dyn RepoResolver> {
129 Arc::new(|target: &RepoTarget| match target {
130 RepoTarget::Did(did) => RepoLookup::Hosted(did.clone()),
131 RepoTarget::OwnerPath(_, _) => RepoLookup::Unhosted,
132 })
133}
134
135async fn clone_over_h3(edge: &Edge, did: &str, tip: &str) -> Vec<u8> {
136 let connection = edge
137 .client
138 .connect(edge.addr, "localhost")
139 .unwrap()
140 .await
141 .unwrap();
142 let quic = connection.clone();
143 let (mut driver, mut sender) = h3::client::new(h3_quinn::Connection::new(connection))
144 .await
145 .unwrap();
146 let drive = tokio::spawn(async move {
147 let _ = std::future::poll_fn(|cx| driver.poll_close(cx)).await;
148 });
149
150 let warmup = http::Request::get(format!(
151 "https://localhost/{did}/info/refs?service=git-upload-pack"
152 ))
153 .header("git-protocol", "version=2")
154 .body(())
155 .unwrap();
156 let mut warm = sender.send_request(warmup).await.unwrap();
157 common::finish_request(&mut warm).await;
158 assert!(
159 warm.recv_response().await.unwrap().status().is_success(),
160 "h3 info/refs advertisement must serve over QUIC"
161 );
162 common::drain(&mut warm).await;
163
164 let request = http::Request::builder()
165 .method(Method::POST)
166 .uri(format!("https://localhost/{did}/git-upload-pack"))
167 .header("content-type", "application/x-git-upload-pack-request")
168 .header("git-protocol", "version=2")
169 .body(())
170 .unwrap();
171 let mut stream = sender.send_request(request).await.unwrap();
172 stream.send_data(fetch_body(tip)).await.unwrap();
173 common::finish_request(&mut stream).await;
174 let response = stream.recv_response().await.unwrap();
175 assert!(
176 response.status().is_success(),
177 "h3 upload-pack returned {}",
178 response.status()
179 );
180 let out = common::drain(&mut stream).await;
181 quic.close(0u32.into(), b"done");
182 drive.abort();
183 out
184}
185
186async fn cloned_set(
187 layout: &Layout,
188 did: &str,
189 tip: &str,
190 format: ObjectFormat,
191) -> BTreeSet<String> {
192 let certdir = tempfile::tempdir().unwrap();
193 let clonedir = tempfile::tempdir().unwrap();
194 let edge = common::serve_edge(certdir.path(), || {
195 let (write_routes, advertisement) = knot_pack::edge_routes(
196 layout.clone(),
197 serve_dids(),
198 None,
199 None,
200 knot_resource::PackSlots::new(4),
201 CacheConfig::default(),
202 Arc::new(knot_messages::Catalog::defaults()),
203 knot_pack::default_hostname().clone(),
204 Arc::new(knot_runtime::SystemClock),
205 );
206 (RequiresFullHandshake::new(write_routes), advertisement)
207 })
208 .await;
209 let pack = extract_pack(&clone_over_h3(&edge, did, tip).await);
210 edge.shutdown.cancel();
211 edge.task.abort();
212 init(clonedir.path(), format);
213 index_pack(clonedir.path(), &pack);
214 object_set(clonedir.path())
215}
216
217async fn the_live_h3_transport_is_logically_reproducible(format: ObjectFormat, did_str: &str) {
218 let scan = tempfile::tempdir().unwrap();
219 let scratch = tempfile::tempdir().unwrap();
220
221 let did = RepoDid::new(did_str).unwrap();
222 let layout = Layout::new(scan.path()).with_object_format(format);
223 layout.create(&did).unwrap();
224 let bare = layout.repo_path(&did).unwrap();
225 let tip = seed_bare(&scratch.path().join("work"), &bare, format);
226 let truth = object_set(&bare);
227
228 let first = cloned_set(&layout, did_str, &tip, format).await;
229 let second = cloned_set(&layout, did_str, &tip, format).await;
230
231 assert_eq!(
232 first, second,
233 "{format:?}: two independent live h3 connections must deliver the same object set"
234 );
235 assert_eq!(
236 first, truth,
237 "{format:?}: the h3 clone's object set must equal the bare's full object set"
238 );
239}
240
241#[tokio::test(flavor = "multi_thread")]
242async fn the_new_h3_transport_replays_to_the_same_object_set_off_the_recorded_trace() {
243 the_live_h3_transport_is_logically_reproducible(ObjectFormat::SHA1, "did:plc:squid").await;
244 the_live_h3_transport_is_logically_reproducible(ObjectFormat::SHA256, "did:plc:cuttle").await;
245}