This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / knot2 / crates / knot-sim / tests / h3.rs
8.2 kB 245 lines
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}