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 / lfs_roundtrip.rs
52 kB 1462 lines
1mod common; 2 3use std::collections::BTreeSet; 4use std::io::Write; 5use std::path::{Path, PathBuf}; 6use std::process::{Command, Stdio}; 7use std::sync::Arc; 8use std::time::Duration; 9 10use base64::Engine; 11use base64::engine::general_purpose::URL_SAFE_NO_PAD; 12use bytes::Bytes; 13use http::Method; 14use knot_atproto::Atproto; 15use knot_cob::{CobHome, CobStore}; 16use knot_cobs::{Registration, RegistryChange}; 17use knot_edge::RequiresFullHandshake; 18use knot_git::{Layout, Repo}; 19use knot_lfs::{FreeSpaceFloor, LfsHandle, LfsOid, LfsSize, LfsStore, LfsStorePath}; 20use knot_runtime::{ 21 FakeHttp, HttpResponse, K256Signer, ManualClock, OsEntropy, Signer, UnixMicros, 22}; 23use knot_secrets::{MasterKey, SealedStore}; 24use knot_types::{ 25 AccountDid, AdmissionPolicy, AuthorName, Email, KnotHostname, KnotId, OwnerDid, RepoDid, 26 RepoName, RepoRkey, UnixSeconds, 27}; 28use sha2::{Digest, Sha256}; 29use tempfile::TempDir; 30use tokio::net::TcpListener; 31use tower::ServiceExt; 32use url::Url; 33 34const REPO_DID: &str = "did:plc:squid"; 35const REPO_NAME: &str = "anemone"; 36const FORK_NAME: &str = "anemone-fork"; 37const OWNER_DID: &str = "did:plc:nel"; 38const PDS_HOST: &str = "pds.oyster.cafe"; 39const KNOT_DID: &str = "did:web:nel.pet"; 40const PINNED_DATE: &str = "2026-07-07T12:00:00+00:00"; 41 42fn require_git_lfs() -> bool { 43 let available = Command::new("git-lfs") 44 .arg("version") 45 .output() 46 .map(|out| out.status.success()) 47 .unwrap_or(false); 48 match (available, std::env::var("KNOT_LFS_ROUNDTRIP").as_deref()) { 49 (true, _) => true, 50 (false, Ok("skip")) => { 51 eprintln!( 52 "skipping lfs round trip gate: git-lfs unavailable and KNOT_LFS_ROUNDTRIP=skip" 53 ); 54 false 55 } 56 (false, _) => panic!( 57 "the lfs round trip gate found no working git-lfs on PATH. \ 58 Install git-lfs or set KNOT_LFS_ROUNDTRIP=skip to skip the gate." 59 ), 60 } 61} 62 63fn media_bytes() -> Vec<u8> { 64 (0..1_048_576u32) 65 .map(|n| (n.wrapping_mul(31) % 251) as u8) 66 .collect() 67} 68 69fn second_media_bytes() -> Vec<u8> { 70 (0..524_288u32) 71 .map(|n| (n.wrapping_mul(97).wrapping_add(13) % 253) as u8) 72 .collect() 73} 74 75fn require_scutiger() -> bool { 76 let available = Command::new("git-lfs-transfer") 77 .arg("--help") 78 .output() 79 .map(|out| out.status.success()) 80 .unwrap_or(false); 81 match (available, std::env::var("KNOT_LFS_CONFORMANCE").as_deref()) { 82 (true, _) => true, 83 (false, Ok("skip")) => { 84 eprintln!( 85 "skipping lfs conformance gate: git-lfs-transfer unavailable and \ 86 KNOT_LFS_CONFORMANCE=skip" 87 ); 88 false 89 } 90 (false, _) => panic!( 91 "the lfs conformance gate found no scutiger git-lfs-transfer on PATH. \ 92 Install it or set KNOT_LFS_CONFORMANCE=skip to skip the gate." 93 ), 94 } 95} 96 97fn git(cwd: &Path, env: &[(String, String)], args: &[&str]) -> (bool, String) { 98 let mut command = knot_fixtures::command_at(cwd, PINNED_DATE); 99 command.args(args); 100 env.iter().for_each(|(key, value)| { 101 command.env(key, value); 102 }); 103 let out = command.output().expect("git runs"); 104 ( 105 out.status.success(), 106 format!( 107 "{}{}", 108 String::from_utf8_lossy(&out.stdout), 109 String::from_utf8_lossy(&out.stderr) 110 ), 111 ) 112} 113 114fn keygen(dir: &Path) -> (String, String) { 115 let path = dir.join("client"); 116 let out = Command::new("ssh-keygen") 117 .args([ 118 "-t", 119 "ed25519", 120 "-N", 121 "", 122 "-C", 123 "nel@oyster.cafe", 124 "-f", 125 path.to_str().unwrap(), 126 ]) 127 .output() 128 .expect("ssh-keygen runs"); 129 assert!(out.status.success()); 130 let public_line = std::fs::read_to_string(dir.join("client.pub")) 131 .unwrap() 132 .trim() 133 .to_string(); 134 (path.to_str().unwrap().to_string(), public_line) 135} 136 137fn actor_signer() -> K256Signer { 138 K256Signer::from_slice(&[9u8; 32]).unwrap() 139} 140 141fn did_document(did: &str) -> Vec<u8> { 142 let multikey = knot_types::crypto::multikey(0xe7, actor_signer().public_key().as_bytes()); 143 serde_json::to_vec(&serde_json::json!({ 144 "id": did, 145 "alsoKnownAs": ["at://nel.pet"], 146 "verificationMethod": [{ 147 "id": format!("{did}#atproto"), 148 "type": "Multikey", 149 "controller": did, 150 "publicKeyMultibase": multikey 151 }], 152 "service": [{ 153 "id": "#atproto_pds", 154 "type": "AtprotoPersonalDataServer", 155 "serviceEndpoint": format!("https://{PDS_HOST}") 156 }] 157 })) 158 .unwrap() 159} 160 161fn list_records_body(public_line: &str) -> Vec<u8> { 162 serde_json::to_vec(&serde_json::json!({ 163 "records": [{ 164 "uri": format!("at://{OWNER_DID}/sh.tangled.publicKey/1"), 165 "value": { 166 "$type": "sh.tangled.publicKey", 167 "key": public_line, 168 "name": "laptop", 169 "createdAt": "2026-07-01T00:00:00Z" 170 } 171 }] 172 })) 173 .unwrap() 174} 175 176fn fake_http( 177 published_line: String, 178) -> FakeHttp< 179 impl Fn(&knot_runtime::HttpRequest) -> Result<HttpResponse, knot_runtime::NetworkError> 180 + Send 181 + Sync, 182> { 183 FakeHttp::new(move |request: &knot_runtime::HttpRequest| { 184 let host = request.url.host_str().unwrap_or_default().to_string(); 185 let path = request.url.path().to_string(); 186 let body = if host == PDS_HOST { 187 list_records_body(&published_line) 188 } else if host == "plc.directory" && request.method == http::Method::POST { 189 b"{}".to_vec() 190 } else if host == "plc.directory" && path.starts_with("/did:") { 191 did_document(path.trim_start_matches('/')) 192 } else { 193 return Ok(HttpResponse { 194 status: http::StatusCode::NOT_FOUND, 195 headers: http::HeaderMap::new(), 196 body: bytes::Bytes::new(), 197 }); 198 }; 199 Ok(HttpResponse { 200 status: http::StatusCode::OK, 201 headers: http::HeaderMap::new(), 202 body: bytes::Bytes::from(body), 203 }) 204 }) 205} 206 207fn service_jwt(nsid: &str, jti: &str) -> String { 208 service_jwt_as(OWNER_DID, nsid, jti) 209} 210 211fn service_jwt_as(iss: &str, nsid: &str, jti: &str) -> String { 212 let header = URL_SAFE_NO_PAD.encode(br#"{"alg":"ES256K","typ":"JWT"}"#); 213 let claims = serde_json::json!({ 214 "iss": iss, 215 "aud": KNOT_DID, 216 "exp": 1_001, 217 "iat": 999, 218 "jti": jti, 219 "lxm": nsid, 220 }); 221 let payload = URL_SAFE_NO_PAD.encode(serde_json::to_vec(&claims).unwrap()); 222 let signing_input = format!("{header}.{payload}"); 223 let signature = actor_signer().sign(signing_input.as_bytes()); 224 format!( 225 "{signing_input}.{}", 226 URL_SAFE_NO_PAD.encode(signature.as_bytes()) 227 ) 228} 229 230struct World { 231 _scan: TempDir, 232 lfs: LfsHandle, 233 ssh_port: u16, 234 http_base: String, 235 router: axum::Router, 236 layout: Layout, 237 h3: Option<common::Edge>, 238 _certdir: Option<TempDir>, 239} 240 241async fn spawn_world(published_line: String) -> World { 242 spawn(published_line, false).await 243} 244 245async fn spawn(published_line: String, with_h3: bool) -> World { 246 let scan = tempfile::tempdir().unwrap(); 247 let meta_path = scan.path().join("meta"); 248 Repo::create(&meta_path).unwrap(); 249 let layout = Layout::new(scan.path().join("repos")); 250 let repo_did = RepoDid::new(REPO_DID).unwrap(); 251 layout.create(&repo_did).unwrap(); 252 253 let knot = KnotId::new(KNOT_DID).unwrap(); 254 let secrets = Arc::new( 255 SealedStore::open( 256 scan.path().join("keys.sealed"), 257 &MasterKey::new([7u8; 32]).unwrap(), 258 Box::new(OsEntropy), 259 ) 260 .unwrap(), 261 ); 262 secrets.ensure(&knot).unwrap(); 263 let knot_signer = secrets.signer(&knot).unwrap(); 264 265 let meta = Repo::open(&meta_path).unwrap(); 266 CobStore::new(&meta) 267 .create( 268 &CobHome::from(&knot), 269 &RegistryChange::Register(Registration { 270 owner: OwnerDid::new(OWNER_DID).unwrap(), 271 rkey: RepoRkey::new(REPO_NAME).unwrap(), 272 name: RepoName::new(REPO_NAME).unwrap(), 273 repo: repo_did.clone(), 274 created_at: UnixSeconds::new(1), 275 }), 276 &knot_signer, 277 UnixSeconds::new(1), 278 ) 279 .unwrap(); 280 281 let index = Arc::new(knot_index::Index::new(meta_path.clone(), layout.clone())); 282 index.rebuild().unwrap(); 283 284 let atproto = Arc::new(Atproto::new( 285 fake_http(published_line), 286 ManualClock::new(UnixMicros::new(1_000_000_000)), 287 knot.clone(), 288 knot_atproto::PlcDirectory::new(Url::parse("https://plc.directory/").unwrap()).unwrap(), 289 )); 290 291 let lfs_store_dir = scan.path().join("lfs"); 292 std::fs::create_dir_all(&lfs_store_dir).unwrap(); 293 let lfs = LfsHandle::open( 294 LfsStorePath::new(&lfs_store_dir), 295 LfsSize::new(64 * 1024 * 1024), 296 FreeSpaceFloor::new(0), 297 ) 298 .unwrap(); 299 300 let key_dir = scan.path().join("hostkey"); 301 std::fs::create_dir_all(&key_dir).unwrap(); 302 let host_key = knot_ssh::load_or_create_host_key(&key_dir.join("host")).unwrap(); 303 let events = Arc::new(knot_events::EventLog::new( 304 ManualClock::new(UnixMicros::new(1_000_000_000)), 305 knot_events::ReplayBounds::new( 306 knot_events::ReplayEvents::new(64).unwrap(), 307 knot_events::ReplayBytes::new(16 << 20).unwrap(), 308 ), 309 )); 310 let ssh_state = Arc::new( 311 knot_ssh::SshState::new(knot_ssh::SshConfig { 312 layout: layout.clone(), 313 index: Arc::clone(&index), 314 atproto: Arc::clone(&atproto), 315 knot_actor: knot_types::ActorId::from_secp256k1(actor_signer().public_key().as_bytes()), 316 events: Arc::clone(&events), 317 hostname: KnotHostname::new("nel.pet").unwrap(), 318 appview: knot_types::AppviewEndpoint::new("https://tangled.test").unwrap(), 319 admins: BTreeSet::from([AccountDid::new(OWNER_DID).unwrap()]), 320 admission: AdmissionPolicy::Closed, 321 max_pack_bytes: knot_xrpc::MaxWireBytes::new(1 << 30), 322 archive_limit: knot_git::ArchiveLimit::default(), 323 languages_push_budget: knot_xrpc::LanguagesPushBudget::new(Duration::from_secs(2)), 324 ci_logs: None, 325 }) 326 .with_lfs(lfs.clone(), 16), 327 ); 328 let ssh_listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); 329 let ssh_port = ssh_listener.local_addr().unwrap().port(); 330 tokio::spawn(async move { 331 let _ = knot_ssh::serve_on_socket(ssh_listener, host_key, ssh_state).await; 332 }); 333 334 let http_listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); 335 let http_base = format!( 336 "http://127.0.0.1:{}", 337 http_listener.local_addr().unwrap().port() 338 ); 339 340 let xrpc_state = Arc::new(knot_xrpc::XrpcState { 341 ci_logs: None, 342 layout: layout.clone(), 343 index: Arc::clone(&index), 344 atproto, 345 secrets, 346 entropy: Arc::new(OsEntropy), 347 admins: BTreeSet::from([AccountDid::new(OWNER_DID).unwrap()]), 348 admission: AdmissionPolicy::Closed, 349 knot_did: knot, 350 knot_hostname: KnotHostname::new("nel.pet").unwrap(), 351 meta_path, 352 knot_service_url: knot_types::KnotServiceUrl::new(http_base.clone()).unwrap(), 353 limiter: Arc::new(knot_xrpc::PreAuthLimiter::default()), 354 cob_locks: Arc::new(knot_xrpc::CobLocks::default()), 355 reservations: Arc::new(knot_xrpc::Reservations::new( 356 knot_xrpc::ReservationTtl::new(1_000_000), 357 knot_xrpc::PerActorQuota::new(16), 358 knot_xrpc::GlobalQuota::new(16), 359 )), 360 proxy_trust: knot_types::ProxyTrust::default(), 361 committer: knot_xrpc::Committer { 362 name: AuthorName::new("Tangled"), 363 email: Email::new("noreply@tangled.sh"), 364 }, 365 byte_limits: knot_xrpc::ByteLimits { 366 pack: knot_xrpc::MaxWireBytes::new(1 << 30), 367 ..knot_xrpc::ByteLimits::default() 368 }, 369 budgets: knot_xrpc::Budgets::default(), 370 git_http: Arc::new(FakeHttp::new(|_request: &knot_runtime::HttpRequest| { 371 Err(knot_runtime::NetworkError::Connect( 372 "no remote upstream is served in this gate".to_string(), 373 )) 374 })), 375 pack_limits: knot_pack::PackLimits::default(), 376 service_owner: AccountDid::new(OWNER_DID).unwrap(), 377 events, 378 subscriber_gate: Arc::new(knot_events::SubscriberGate::new( 379 knot_events::GlobalSubscriberLimit::new(16), 380 knot_events::PerPeerSubscriberLimit::new(4), 381 )), 382 maintenance: knot_maintenance::MaintenanceHandle::disabled(), 383 appview: knot_types::AppviewEndpoint::new("https://tangled.test").unwrap(), 384 slots: knot_resource::Slots::testing(8), 385 lfs: Some(knot_xrpc::LfsWeb::new(lfs.clone(), 8)), 386 catalog: Arc::new(knot_messages::Catalog::defaults()), 387 }); 388 389 let resolver: Arc<dyn knot_pack::RepoResolver> = { 390 let index = Arc::clone(&index); 391 Arc::new(move |target: &knot_pack::RepoTarget| match target { 392 knot_pack::RepoTarget::Did(did) => match index.owner_of(did) { 393 knot_index::Resolved::Ready(Some(_)) => knot_pack::RepoLookup::Hosted(did.clone()), 394 knot_index::Resolved::Ready(None) => knot_pack::RepoLookup::Unhosted, 395 knot_index::Resolved::Warming => knot_pack::RepoLookup::Unavailable, 396 }, 397 knot_pack::RepoTarget::OwnerPath(owner, path) => { 398 match index.resolve_clone_path(owner, path) { 399 knot_index::Resolved::Ready(Some(found)) => { 400 knot_pack::RepoLookup::Hosted(found) 401 } 402 knot_index::Resolved::Ready(None) => knot_pack::RepoLookup::Unhosted, 403 knot_index::Resolved::Warming => knot_pack::RepoLookup::Unavailable, 404 } 405 } 406 }) 407 }; 408 let advertiser = knot_xrpc::receive_advertiser(Arc::clone(&xrpc_state)); 409 let (write_routes, advertisement) = knot_pack::edge_routes(knot_pack::EdgeConfig { 410 receive: Some(Arc::clone(&advertiser)), 411 pack_slots: knot_resource::PackSlots::new(4), 412 ..knot_pack::EdgeConfig::serving( 413 layout.clone(), 414 Arc::clone(&resolver), 415 Arc::new(knot_runtime::SystemClock), 416 ) 417 }); 418 let router = write_routes 419 .merge(advertisement.into_router()) 420 .merge(knot_xrpc::router(Arc::clone(&xrpc_state))); 421 let served = router.clone(); 422 tokio::spawn(async move { 423 let _ = axum::serve(http_listener, served).await; 424 }); 425 426 let (h3, certdir) = match with_h3 { 427 true => { 428 let certdir = tempfile::tempdir().unwrap(); 429 let edge = common::serve_edge(certdir.path(), || { 430 let (write_routes, advertisement) = knot_pack::edge_routes(knot_pack::EdgeConfig { 431 receive: Some(Arc::clone(&advertiser)), 432 pack_slots: knot_resource::PackSlots::new(4), 433 ..knot_pack::EdgeConfig::serving( 434 layout.clone(), 435 Arc::clone(&resolver), 436 Arc::new(knot_runtime::SystemClock), 437 ) 438 }); 439 let app = RequiresFullHandshake::new( 440 write_routes.merge(knot_xrpc::router(Arc::clone(&xrpc_state))), 441 ); 442 (app, advertisement) 443 }) 444 .await; 445 (Some(edge), Some(certdir)) 446 } 447 false => (None, None), 448 }; 449 450 World { 451 _scan: scan, 452 lfs, 453 ssh_port, 454 http_base, 455 router, 456 layout, 457 h3, 458 _certdir: certdir, 459 } 460} 461 462async fn in_git_blocking<T: Send + 'static>(task: impl FnOnce() -> T + Send + 'static) -> T { 463 tokio::task::spawn_blocking(task).await.unwrap() 464} 465 466fn seed_lfs_work(work: &Path, env: &[(String, String)], media: &[u8]) { 467 std::fs::create_dir_all(work).unwrap(); 468 let steps: [&[&str]; 2] = [ 469 &["init", "-q", "-b", "main"], 470 &["lfs", "install", "--local"], 471 ]; 472 steps.iter().for_each(|args| { 473 let (ok, out) = git(work, env, args); 474 assert!(ok, "{args:?} failed:\n{out}"); 475 }); 476 let (ok, out) = git(work, env, &["lfs", "track", "*.bin"]); 477 assert!(ok, "lfs track failed:\n{out}"); 478 let (ok, out) = git(work, env, &["config", "lfs.locksverify", "false"]); 479 assert!(ok, "config failed:\n{out}"); 480 std::fs::write(work.join("media.bin"), media).unwrap(); 481 std::fs::write(work.join("README.md"), "media lives in lfs\n").unwrap(); 482 let commit: [&[&str]; 2] = [&["add", "-A"], &["commit", "-q", "-m", "media"]]; 483 commit.iter().for_each(|args| { 484 let (ok, out) = git(work, env, args); 485 assert!(ok, "{args:?} failed:\n{out}"); 486 }); 487} 488 489fn clone_and_pull(base: &Path, url: &str, name: &str, env: &[(String, String)]) -> PathBuf { 490 let skip_smudge: Vec<(String, String)> = env 491 .iter() 492 .cloned() 493 .chain([("GIT_LFS_SKIP_SMUDGE".to_string(), "1".to_string())]) 494 .collect(); 495 let (ok, out) = git(base, &skip_smudge, &["clone", "-q", url, name]); 496 assert!(ok, "anonymous clone of {url} failed:\n{out}"); 497 let dst = base.join(name); 498 let pointer = std::fs::read_to_string(dst.join("media.bin")).unwrap(); 499 assert!( 500 pointer.contains("git-lfs.github.com/spec/v1"), 501 "clone must land the pointer before lfs pull, got:\n{pointer}" 502 ); 503 let (ok, out) = git(&dst, env, &["lfs", "install", "--local"]); 504 assert!(ok, "lfs install in {name} failed:\n{out}"); 505 let (ok, out) = git(&dst, env, &["lfs", "pull"]); 506 assert!(ok, "git lfs pull in {name} failed:\n{out}"); 507 dst 508} 509 510async fn create_fork(world: &World, jti: &str) -> (http::StatusCode, serde_json::Value) { 511 let token = service_jwt("sh.tangled.repo.create", jti); 512 let body = serde_json::json!({ 513 "rkey": FORK_NAME, 514 "name": FORK_NAME, 515 "source": format!("{}/{OWNER_DID}/{REPO_NAME}", world.http_base), 516 }); 517 let request = http::Request::builder() 518 .method("POST") 519 .uri("/xrpc/sh.tangled.repo.create") 520 .header(http::header::CONTENT_TYPE, "application/json") 521 .header(http::header::AUTHORIZATION, format!("Bearer {token}")) 522 .body(axum::body::Body::from(serde_json::to_vec(&body).unwrap())) 523 .unwrap(); 524 let response = world.router.clone().oneshot(request).await.unwrap(); 525 let status = response.status(); 526 let bytes = axum::body::to_bytes(response.into_body(), usize::MAX) 527 .await 528 .unwrap(); 529 let value = serde_json::from_slice(&bytes).unwrap_or(serde_json::Value::Null); 530 (status, value) 531} 532 533#[tokio::test(flavor = "multi_thread", worker_threads = 4)] 534async fn the_lfs_round_trip_gate_holds_over_both_transports_and_the_fork() { 535 if !require_git_lfs() { 536 return; 537 } 538 let scratch = tempfile::tempdir().unwrap(); 539 let (key_path, public_line) = keygen(scratch.path()); 540 let world = spawn_world(public_line).await; 541 542 let media = media_bytes(); 543 let media_oid = LfsOid::from_digest(Sha256::digest(&media).into()); 544 let ssh = format!( 545 "ssh -i {key_path} -o IdentitiesOnly=yes -o StrictHostKeyChecking=no \ 546 -o UserKnownHostsFile=/dev/null -o PreferredAuthentications=publickey -o BatchMode=yes" 547 ); 548 let path_env = std::env::var("PATH").unwrap_or_default(); 549 let home = scratch.path().to_str().unwrap().to_string(); 550 let env: Vec<(String, String)> = [ 551 ("GIT_SSH_COMMAND", &ssh), 552 ("PATH", &path_env), 553 ("HOME", &home), 554 ] 555 .map(|(key, value)| (key.to_string(), value.clone())) 556 .to_vec(); 557 558 let work = scratch.path().join("work"); 559 seed_lfs_work(&work, &env, &media); 560 561 let push_url = format!( 562 "ssh://git@127.0.0.1:{}/{OWNER_DID}/{REPO_NAME}", 563 world.ssh_port 564 ); 565 let (ok, out) = { 566 let work = work.clone(); 567 let env = env.clone(); 568 in_git_blocking(move || git(&work, &env, &["push", "-q", &push_url, "main"])).await 569 }; 570 assert!(ok, "lfs push over ssh failed:\n{out}"); 571 572 let source_repo = RepoDid::new(REPO_DID).unwrap(); 573 assert_eq!( 574 world.lfs.store.probe(&source_repo, &media_oid).unwrap(), 575 Some(LfsSize::new(media.len() as u64)), 576 "pushed media must be durable in the store" 577 ); 578 579 let clone_url = format!("{}/{OWNER_DID}/{REPO_NAME}", world.http_base); 580 let dst = { 581 let base = scratch.path().to_path_buf(); 582 let env = env.clone(); 583 in_git_blocking(move || clone_and_pull(&base, &clone_url, "reader", &env)).await 584 }; 585 assert_eq!( 586 std::fs::read(dst.join("media.bin")).unwrap(), 587 media, 588 "anonymous http reader must see byte-identical media" 589 ); 590 591 let (status, created) = create_fork(&world, "gate-fork-1").await; 592 assert_eq!( 593 status, 594 http::StatusCode::OK, 595 "fork create failed: {created}" 596 ); 597 assert!( 598 created.get("lfsMissing").is_none(), 599 "local fork must copy every object, got {created}" 600 ); 601 let fork_did = RepoDid::new(created["repoDid"].as_str().unwrap()).unwrap(); 602 assert_eq!( 603 world.lfs.store.probe(&fork_did, &media_oid).unwrap(), 604 Some(LfsSize::new(media.len() as u64)), 605 "fork prefix must hold its own copy of the media" 606 ); 607 608 let fork_url = format!("{}/{OWNER_DID}/{FORK_NAME}", world.http_base); 609 let fork_dst = { 610 let base = scratch.path().to_path_buf(); 611 let env = env.clone(); 612 in_git_blocking(move || clone_and_pull(&base, &fork_url, "fork-reader", &env)).await 613 }; 614 assert_eq!( 615 std::fs::read(fork_dst.join("media.bin")).unwrap(), 616 media, 617 "anonymous clone of the fork must see byte-identical media" 618 ); 619} 620 621fn seed_many_lfs(work: &Path, env: &[(String, String)], count: usize) { 622 std::fs::create_dir_all(work).unwrap(); 623 let steps: [&[&str]; 2] = [ 624 &["init", "-q", "-b", "main"], 625 &["lfs", "install", "--local"], 626 ]; 627 steps.iter().for_each(|args| { 628 let (ok, out) = git(work, env, args); 629 assert!(ok, "{args:?} failed:\n{out}"); 630 }); 631 let (ok, out) = git(work, env, &["lfs", "track", "*.bin"]); 632 assert!(ok, "lfs track failed:\n{out}"); 633 let (ok, out) = git(work, env, &["config", "lfs.locksverify", "false"]); 634 assert!(ok, "config failed:\n{out}"); 635 (0..count).for_each(|index| { 636 let size = 200 + index * 7; 637 let bytes: Vec<u8> = (0..size) 638 .map(|n| (n.wrapping_mul(31).wrapping_add(index) % 251) as u8) 639 .collect(); 640 std::fs::write(work.join(format!("object-{index}.bin")), bytes).unwrap(); 641 }); 642 let commit: [&[&str]; 2] = [&["add", "-A"], &["commit", "-q", "-m", "many media"]]; 643 commit.iter().for_each(|args| { 644 let (ok, out) = git(work, env, args); 645 assert!(ok, "{args:?} failed:\n{out}"); 646 }); 647} 648 649#[tokio::test(flavor = "multi_thread", worker_threads = 4)] 650async fn many_objects_ride_default_git_lfs_concurrency_over_both_transports() { 651 if !require_git_lfs() { 652 return; 653 } 654 let scratch = tempfile::tempdir().unwrap(); 655 let (key_path, public_line) = keygen(scratch.path()); 656 let world = spawn_world(public_line).await; 657 658 let ssh = format!( 659 "ssh -i {key_path} -o IdentitiesOnly=yes -o StrictHostKeyChecking=no \ 660 -o UserKnownHostsFile=/dev/null -o PreferredAuthentications=publickey -o BatchMode=yes" 661 ); 662 let path_env = std::env::var("PATH").unwrap_or_default(); 663 let home = scratch.path().to_str().unwrap().to_string(); 664 let env: Vec<(String, String)> = [ 665 ("GIT_SSH_COMMAND", &ssh), 666 ("PATH", &path_env), 667 ("HOME", &home), 668 ] 669 .map(|(key, value)| (key.to_string(), value.clone())) 670 .to_vec(); 671 672 let count = 25usize; 673 let work = scratch.path().join("work"); 674 seed_many_lfs(&work, &env, count); 675 676 let push_url = format!( 677 "ssh://git@127.0.0.1:{}/{OWNER_DID}/{REPO_NAME}", 678 world.ssh_port 679 ); 680 let (ok, out) = { 681 let work = work.clone(); 682 let env = env.clone(); 683 in_git_blocking(move || git(&work, &env, &["push", "-q", &push_url, "main"])).await 684 }; 685 assert!( 686 ok, 687 "git-lfs at its default concurrency must push {count} objects over ssh without tripping \ 688 the per-peer connection limit:\n{out}" 689 ); 690 691 let clone_url = format!("{}/{OWNER_DID}/{REPO_NAME}", world.http_base); 692 let dst = { 693 let base = scratch.path().to_path_buf(); 694 let env = env.clone(); 695 in_git_blocking(move || { 696 let skip_smudge: Vec<(String, String)> = env 697 .iter() 698 .cloned() 699 .chain([("GIT_LFS_SKIP_SMUDGE".to_string(), "1".to_string())]) 700 .collect(); 701 let (ok, out) = git(&base, &skip_smudge, &["clone", "-q", &clone_url, "reader"]); 702 assert!(ok, "anonymous clone failed:\n{out}"); 703 let dst = base.join("reader"); 704 let (ok, out) = git(&dst, &env, &["lfs", "install", "--local"]); 705 assert!(ok, "lfs install failed:\n{out}"); 706 let (ok, out) = git(&dst, &env, &["lfs", "pull"]); 707 assert!( 708 ok, 709 "anonymous http pull of {count} objects mustn't be throttled by the xrpc \ 710 pre-auth limiter:\n{out}" 711 ); 712 dst 713 }) 714 .await 715 }; 716 (0..count).for_each(|index| { 717 assert_eq!( 718 std::fs::read(work.join(format!("object-{index}.bin"))).unwrap(), 719 std::fs::read(dst.join(format!("object-{index}.bin"))).unwrap(), 720 "object-{index}.bin must be byte-identical over anonymous http" 721 ); 722 }); 723} 724 725fn write_shim(dir: &Path) -> String { 726 let shim = dir.join("local-ssh.sh"); 727 std::fs::write( 728 &shim, 729 "#!/bin/sh\n\ 730 while [ \"$#\" -gt 0 ]; do\n\ 731 case \"$1\" in\n\ 732 -o|-p) shift 2 ;;\n\ 733 -*) shift ;;\n\ 734 *) break ;;\n\ 735 esac\n\ 736 done\n\ 737 shift\n\ 738 eval exec \"$@\"\n", 739 ) 740 .unwrap(); 741 let mut permissions = std::fs::metadata(&shim).unwrap().permissions(); 742 std::os::unix::fs::PermissionsExt::set_mode(&mut permissions, 0o755); 743 std::fs::set_permissions(&shim, permissions).unwrap(); 744 shim.to_str().unwrap().to_string() 745} 746 747fn hex_object_files(root: &Path) -> Vec<(String, u64, PathBuf)> { 748 let entries = match std::fs::read_dir(root) { 749 Ok(entries) => entries, 750 Err(_) => return Vec::new(), 751 }; 752 entries 753 .filter_map(Result::ok) 754 .flat_map(|entry| { 755 let path = entry.path(); 756 if path.is_dir() { 757 return hex_object_files(&path); 758 } 759 path.file_name() 760 .and_then(|name| name.to_str()) 761 .filter(|name| LfsOid::new(*name).is_ok()) 762 .map(|name| { 763 let size = std::fs::metadata(&path).map(|meta| meta.len()).unwrap_or(0); 764 vec![(name.to_string(), size, path.clone())] 765 }) 766 .unwrap_or_default() 767 }) 768 .collect() 769} 770 771fn pull_verdict(base: &Path, url: &str, name: &str, env: &[(String, String)]) -> (bool, String) { 772 let skip_smudge: Vec<(String, String)> = env 773 .iter() 774 .cloned() 775 .chain([("GIT_LFS_SKIP_SMUDGE".to_string(), "1".to_string())]) 776 .collect(); 777 let (ok, out) = git(base, &skip_smudge, &["clone", "-q", url, name]); 778 assert!(ok, "clone of {url} failed:\n{out}"); 779 let dst = base.join(name); 780 let (ok, out) = git(&dst, env, &["lfs", "install", "--local"]); 781 assert!(ok, "lfs install in {name} failed:\n{out}"); 782 git(&dst, env, &["lfs", "pull"]) 783} 784 785#[tokio::test(flavor = "multi_thread", worker_threads = 4)] 786async fn the_lfs_stack_is_conformant_with_the_reference_server_and_client() { 787 if !require_git_lfs() || !require_scutiger() { 788 return; 789 } 790 let scratch = tempfile::tempdir().unwrap(); 791 let (key_path, public_line) = keygen(scratch.path()); 792 let world = spawn_world(public_line).await; 793 794 let media = media_bytes(); 795 let second = second_media_bytes(); 796 let media_oid = LfsOid::from_digest(Sha256::digest(&media).into()); 797 let second_oid = LfsOid::from_digest(Sha256::digest(&second).into()); 798 let expected: std::collections::BTreeSet<(String, u64)> = [ 799 (media_oid.as_str().to_string(), media.len() as u64), 800 (second_oid.as_str().to_string(), second.len() as u64), 801 ] 802 .into(); 803 804 let ssh = format!( 805 "ssh -i {key_path} -o IdentitiesOnly=yes -o StrictHostKeyChecking=no \ 806 -o UserKnownHostsFile=/dev/null -o PreferredAuthentications=publickey -o BatchMode=yes" 807 ); 808 let path_env = std::env::var("PATH").unwrap_or_default(); 809 let home = scratch.path().to_str().unwrap().to_string(); 810 let knot_env: Vec<(String, String)> = [ 811 ("GIT_SSH_COMMAND", &ssh), 812 ("PATH", &path_env), 813 ("HOME", &home), 814 ] 815 .map(|(key, value)| (key.to_string(), value.clone())) 816 .to_vec(); 817 let shim = write_shim(scratch.path()); 818 let reference_env: Vec<(String, String)> = [ 819 ("GIT_SSH_COMMAND", &shim), 820 ("PATH", &path_env), 821 ("HOME", &home), 822 ] 823 .map(|(key, value)| (key.to_string(), value.clone())) 824 .to_vec(); 825 826 let upstream = scratch.path().join("reference-upstream.git"); 827 let (ok, out) = git( 828 scratch.path(), 829 &reference_env, 830 &[ 831 "init", 832 "-q", 833 "--bare", 834 "-b", 835 "main", 836 upstream.to_str().unwrap(), 837 ], 838 ); 839 assert!(ok, "reference upstream init failed:\n{out}"); 840 841 let work = scratch.path().join("work"); 842 seed_lfs_work(&work, &knot_env, &media); 843 std::fs::write(work.join("extra.bin"), &second).unwrap(); 844 let commit: [&[&str]; 2] = [&["add", "-A"], &["commit", "-q", "-m", "extra media"]]; 845 commit.iter().for_each(|args| { 846 let (ok, out) = git(&work, &knot_env, args); 847 assert!(ok, "{args:?} failed:\n{out}"); 848 }); 849 850 let knot_push_url = format!( 851 "ssh://git@127.0.0.1:{}/{OWNER_DID}/{REPO_NAME}", 852 world.ssh_port 853 ); 854 let reference_push_url = format!("ssh://ref@localhost{}", upstream.display()); 855 let pushes = { 856 let work = work.clone(); 857 let knot_env = knot_env.clone(); 858 let reference_env = reference_env.clone(); 859 let knot_push_url = knot_push_url.clone(); 860 let reference_push_url = reference_push_url.clone(); 861 in_git_blocking(move || { 862 [ 863 git(&work, &knot_env, &["push", "-q", &knot_push_url, "main"]), 864 git( 865 &work, 866 &reference_env, 867 &["push", "-q", &reference_push_url, "main"], 868 ), 869 ] 870 }) 871 .await 872 }; 873 pushes.iter().for_each(|(ok, out)| { 874 assert!(ok, "push failed:\n{out}"); 875 }); 876 877 let source_repo = RepoDid::new(REPO_DID).unwrap(); 878 let knot_objects: std::collections::BTreeSet<(String, u64)> = world 879 .lfs 880 .store 881 .enumerate(&source_repo) 882 .unwrap() 883 .into_iter() 884 .map(|object| (object.oid.as_str().to_string(), object.size.get())) 885 .collect(); 886 let reference_objects: std::collections::BTreeSet<(String, u64)> = hex_object_files(&upstream) 887 .into_iter() 888 .map(|(name, size, _)| (name, size)) 889 .collect(); 890 assert_eq!( 891 knot_objects, expected, 892 "the knot store holds exactly the pushed object set" 893 ); 894 assert_eq!( 895 knot_objects, reference_objects, 896 "both servers hold identical object sets after the same push" 897 ); 898 899 let knot_clone_url = format!("{}/{OWNER_DID}/{REPO_NAME}", world.http_base); 900 let (knot_dst, reference_dst) = { 901 let base = scratch.path().to_path_buf(); 902 let knot_env = knot_env.clone(); 903 let reference_env = reference_env.clone(); 904 let knot_clone_url = knot_clone_url.clone(); 905 let reference_push_url = reference_push_url.clone(); 906 in_git_blocking(move || { 907 ( 908 clone_and_pull(&base, &knot_clone_url, "knot-reader", &knot_env), 909 clone_and_pull( 910 &base, 911 &reference_push_url, 912 "reference-reader", 913 &reference_env, 914 ), 915 ) 916 }) 917 .await 918 }; 919 ["media.bin", "extra.bin"].iter().for_each(|file| { 920 assert_eq!( 921 std::fs::read(knot_dst.join(file)).unwrap(), 922 std::fs::read(reference_dst.join(file)).unwrap(), 923 "{file}: both servers must check out identical media" 924 ); 925 }); 926 assert_eq!(std::fs::read(knot_dst.join("media.bin")).unwrap(), media); 927 assert_eq!(std::fs::read(knot_dst.join("extra.bin")).unwrap(), second); 928 929 let knot_removed = world 930 .lfs 931 .store 932 .object_file(&source_repo, &second_oid) 933 .unwrap() 934 .unwrap() 935 .1; 936 std::fs::remove_file(knot_removed).unwrap(); 937 let removed = hex_object_files(&upstream) 938 .into_iter() 939 .filter(|(name, _, _)| name == second_oid.as_str()) 940 .map(|(_, _, path)| std::fs::remove_file(path).unwrap()) 941 .count(); 942 assert!( 943 removed > 0, 944 "the reference server holds the object to remove" 945 ); 946 947 let verdicts = { 948 let base = scratch.path().to_path_buf(); 949 let knot_env = knot_env.clone(); 950 let reference_env = reference_env.clone(); 951 let knot_push_url = knot_push_url.clone(); 952 in_git_blocking(move || { 953 [ 954 pull_verdict(&base, &knot_clone_url, "knot-missing-http", &knot_env), 955 pull_verdict(&base, &knot_push_url, "knot-missing-ssh", &knot_env), 956 pull_verdict( 957 &base, 958 &reference_push_url, 959 "reference-missing", 960 &reference_env, 961 ), 962 ] 963 }) 964 .await 965 }; 966 let [(http_ok, http_out), (ssh_ok, ssh_out), (reference_ok, _)] = verdicts; 967 assert!( 968 !http_ok, 969 "an http pull of a missing object must fail loudly, never succeed silently:\n{http_out}" 970 ); 971 assert!( 972 !ssh_ok, 973 "an ssh pull of a missing object must fail loudly, never succeed silently:\n{ssh_out}" 974 ); 975 assert!( 976 reference_ok, 977 "scutiger 0.3.0 answers noop for a missing download and the client silently \ 978 succeeds. This pin is the recorded reason knot answers download instead, so its \ 979 get-object 404 turns the pull into a loud failure. If the reference starts failing \ 980 loudly too, the divergence note can be retired." 981 ); 982} 983 984async fn lfs_batch( 985 world: &World, 986 auth: Option<&str>, 987 op: &str, 988 oid: &LfsOid, 989 size: u64, 990) -> (http::StatusCode, serde_json::Value) { 991 let body = serde_json::json!({ 992 "operation": op, 993 "transfers": ["basic"], 994 "objects": [{ "oid": oid.as_str(), "size": size }], 995 "hash_algo": "sha256", 996 }); 997 let mut builder = http::Request::builder() 998 .method("POST") 999 .uri(format!("/{OWNER_DID}/{REPO_NAME}/info/lfs/objects/batch")) 1000 .header(http::header::CONTENT_TYPE, "application/vnd.git-lfs+json"); 1001 if let Some(auth) = auth { 1002 builder = builder.header(http::header::AUTHORIZATION, auth); 1003 } 1004 let request = builder 1005 .body(axum::body::Body::from(serde_json::to_vec(&body).unwrap())) 1006 .unwrap(); 1007 let response = world.router.clone().oneshot(request).await.unwrap(); 1008 let status = response.status(); 1009 let bytes = axum::body::to_bytes(response.into_body(), usize::MAX) 1010 .await 1011 .unwrap(); 1012 ( 1013 status, 1014 serde_json::from_slice(&bytes).unwrap_or(serde_json::Value::Null), 1015 ) 1016} 1017 1018async fn lfs_put( 1019 world: &World, 1020 auth: Option<&str>, 1021 oid: &LfsOid, 1022 bytes: Vec<u8>, 1023) -> http::StatusCode { 1024 let mut builder = http::Request::builder() 1025 .method("PUT") 1026 .uri(format!( 1027 "/{OWNER_DID}/{REPO_NAME}/info/lfs/objects/{}", 1028 oid.as_str() 1029 )) 1030 .header(http::header::CONTENT_LENGTH, bytes.len()); 1031 if let Some(auth) = auth { 1032 builder = builder.header(http::header::AUTHORIZATION, auth); 1033 } 1034 let request = builder.body(axum::body::Body::from(bytes)).unwrap(); 1035 world 1036 .router 1037 .clone() 1038 .oneshot(request) 1039 .await 1040 .unwrap() 1041 .status() 1042} 1043 1044fn basic_auth(token: &str) -> String { 1045 let raw = base64::engine::general_purpose::STANDARD.encode(format!("x-tangled-token:{token}")); 1046 format!("Basic {raw}") 1047} 1048 1049#[tokio::test(flavor = "multi_thread", worker_threads = 4)] 1050async fn lfs_http_push_stores_an_object_with_a_push_token() { 1051 let scratch = tempfile::tempdir().unwrap(); 1052 let (_key_path, public_line) = keygen(scratch.path()); 1053 let world = spawn_world(public_line).await; 1054 let repo = RepoDid::new(REPO_DID).unwrap(); 1055 1056 let payload: Vec<u8> = (0..4096u32) 1057 .map(|n| (n.wrapping_mul(17) % 251) as u8) 1058 .collect(); 1059 let oid = LfsOid::from_digest(Sha256::digest(&payload).into()); 1060 let size = payload.len() as u64; 1061 1062 let bearer = format!( 1063 "Bearer {}", 1064 service_jwt("sh.tangled.repo.push", "lfs-http-push-1") 1065 ); 1066 let (status, body) = lfs_batch(&world, Some(&bearer), "upload", &oid, size).await; 1067 assert_eq!( 1068 status, 1069 http::StatusCode::OK, 1070 "authenticated upload batch: {body}" 1071 ); 1072 let href = body["objects"][0]["actions"]["upload"]["href"] 1073 .as_str() 1074 .unwrap_or_else(|| panic!("expected an upload action, got {body}")); 1075 assert_eq!( 1076 href, 1077 format!( 1078 "{}/{OWNER_DID}/{REPO_NAME}/info/lfs/objects/{}", 1079 world.http_base, 1080 oid.as_str() 1081 ), 1082 "upload href points at the object route on this knot" 1083 ); 1084 assert_ne!( 1085 body["objects"][0]["authenticated"], 1086 serde_json::Value::Bool(true), 1087 "upload objects mustn't claim authenticated=true, else git-lfs sends the object put with no auth and loops on 401: {body}" 1088 ); 1089 assert!( 1090 world.lfs.store.probe(&repo, &oid).unwrap().is_none(), 1091 "object must be absent before the put" 1092 ); 1093 1094 let put = lfs_put(&world, Some(&bearer), &oid, payload.clone()).await; 1095 assert_eq!( 1096 put, 1097 http::StatusCode::OK, 1098 "the same push token must authorize both the batch and the object put" 1099 ); 1100 assert_eq!( 1101 world.lfs.store.probe(&repo, &oid).unwrap(), 1102 Some(LfsSize::new(size)), 1103 "the put object must be durable in the store" 1104 ); 1105 1106 let get = http::Request::builder() 1107 .method("GET") 1108 .uri(format!( 1109 "/{OWNER_DID}/{REPO_NAME}/info/lfs/objects/{}", 1110 oid.as_str() 1111 )) 1112 .body(axum::body::Body::empty()) 1113 .unwrap(); 1114 let response = world.router.clone().oneshot(get).await.unwrap(); 1115 assert_eq!( 1116 response.status(), 1117 http::StatusCode::OK, 1118 "anonymous download" 1119 ); 1120 let served = axum::body::to_bytes(response.into_body(), usize::MAX) 1121 .await 1122 .unwrap(); 1123 assert_eq!( 1124 served.as_ref(), 1125 payload.as_slice(), 1126 "an anonymous reader sees the byte-identical object a push stored" 1127 ); 1128 1129 let payload2: Vec<u8> = (0..2048u32) 1130 .map(|n| (n.wrapping_mul(29) % 251) as u8) 1131 .collect(); 1132 let oid2 = LfsOid::from_digest(Sha256::digest(&payload2).into()); 1133 let basic = basic_auth(&service_jwt("sh.tangled.repo.push", "lfs-http-push-2")); 1134 let (status, body) = 1135 lfs_batch(&world, Some(&basic), "upload", &oid2, payload2.len() as u64).await; 1136 assert_eq!( 1137 status, 1138 http::StatusCode::OK, 1139 "basic-auth upload batch: {body}" 1140 ); 1141 let put = lfs_put(&world, Some(&basic), &oid2, payload2.clone()).await; 1142 assert_eq!( 1143 put, 1144 http::StatusCode::OK, 1145 "a push token presented as the http basic password must authenticate the put" 1146 ); 1147 assert_eq!( 1148 world.lfs.store.probe(&repo, &oid2).unwrap(), 1149 Some(LfsSize::new(payload2.len() as u64)), 1150 "the basic-authenticated object must be durable too" 1151 ); 1152} 1153 1154#[tokio::test(flavor = "multi_thread", worker_threads = 4)] 1155async fn lfs_http_push_rejects_missing_and_mismatched_credentials() { 1156 let scratch = tempfile::tempdir().unwrap(); 1157 let (_key_path, public_line) = keygen(scratch.path()); 1158 let world = spawn_world(public_line).await; 1159 1160 let payload: Vec<u8> = (0..1024u32) 1161 .map(|n| (n.wrapping_mul(13) % 251) as u8) 1162 .collect(); 1163 let oid = LfsOid::from_digest(Sha256::digest(&payload).into()); 1164 let size = payload.len() as u64; 1165 1166 let (status, _) = lfs_batch(&world, None, "upload", &oid, size).await; 1167 assert_eq!( 1168 status, 1169 http::StatusCode::UNAUTHORIZED, 1170 "an unauthenticated upload batch is challenged" 1171 ); 1172 1173 let wrong_method = format!( 1174 "Bearer {}", 1175 service_jwt("sh.tangled.repo.create", "lfs-http-neg-method") 1176 ); 1177 let (status, _) = lfs_batch(&world, Some(&wrong_method), "upload", &oid, size).await; 1178 assert_eq!( 1179 status, 1180 http::StatusCode::UNAUTHORIZED, 1181 "a token bound to another method cannot authorize a push" 1182 ); 1183 1184 let stranger = format!( 1185 "Bearer {}", 1186 service_jwt_as(REPO_DID, "sh.tangled.repo.push", "lfs-http-neg-acl") 1187 ); 1188 let (status, _) = lfs_batch(&world, Some(&stranger), "upload", &oid, size).await; 1189 assert_eq!( 1190 status, 1191 http::StatusCode::FORBIDDEN, 1192 "a valid push token from a did that cannot push is refused by the acl" 1193 ); 1194 1195 let put = lfs_put(&world, None, &oid, payload).await; 1196 assert_eq!( 1197 put, 1198 http::StatusCode::UNAUTHORIZED, 1199 "an unauthenticated object put is challenged" 1200 ); 1201 assert!( 1202 world 1203 .lfs 1204 .store 1205 .probe(&RepoDid::new(REPO_DID).unwrap(), &oid) 1206 .unwrap() 1207 .is_none(), 1208 "no rejected request may leave an object behind" 1209 ); 1210} 1211 1212#[tokio::test(flavor = "multi_thread", worker_threads = 4)] 1213async fn git_push_over_http_authenticates_and_lands_the_ref() { 1214 let scratch = tempfile::tempdir().unwrap(); 1215 let (_key_path, public_line) = keygen(scratch.path()); 1216 let world = spawn_world(public_line).await; 1217 1218 let path_env = std::env::var("PATH").unwrap_or_default(); 1219 let home = scratch.path().to_str().unwrap().to_string(); 1220 let env: Vec<(String, String)> = [("PATH", &path_env), ("HOME", &home)] 1221 .map(|(key, value)| (key.to_string(), value.clone())) 1222 .to_vec(); 1223 1224 let work = scratch.path().join("work"); 1225 std::fs::create_dir_all(&work).unwrap(); 1226 let (ok, out) = git(&work, &env, &["init", "-q", "-b", "main"]); 1227 assert!(ok, "init failed:\n{out}"); 1228 std::fs::write(work.join("README.md"), "hello over http\n").unwrap(); 1229 let (ok, out) = git(&work, &env, &["add", "-A"]); 1230 assert!(ok, "add failed:\n{out}"); 1231 let (ok, out) = git(&work, &env, &["commit", "-q", "-m", "init over http"]); 1232 assert!(ok, "commit failed:\n{out}"); 1233 1234 let url = format!("{}/{OWNER_DID}/{REPO_NAME}", world.http_base); 1235 1236 let (ok, out) = { 1237 let work = work.clone(); 1238 let env = env.clone(); 1239 let url = url.clone(); 1240 in_git_blocking(move || git(&work, &env, &["push", "-q", &url, "main"])).await 1241 }; 1242 assert!( 1243 !ok, 1244 "an unauthenticated http push must be refused, git reported success:\n{out}" 1245 ); 1246 1247 let token = service_jwt("sh.tangled.repo.push", "git-http-push-1"); 1248 let header = format!( 1249 "http.extraHeader=Authorization: Basic {}", 1250 base64::engine::general_purpose::STANDARD.encode(format!("x-tangled-token:{token}")) 1251 ); 1252 let (ok, out) = { 1253 let work = work.clone(); 1254 let env = env.clone(); 1255 let url = url.clone(); 1256 let header = header.clone(); 1257 in_git_blocking(move || git(&work, &env, &["-c", &header, "push", "-q", &url, "main"])) 1258 .await 1259 }; 1260 assert!(ok, "authenticated http push failed:\n{out}"); 1261 1262 let (ok, refs) = { 1263 let scratch = scratch.path().to_path_buf(); 1264 let env = env.clone(); 1265 let url = url.clone(); 1266 in_git_blocking(move || git(&scratch, &env, &["ls-remote", &url])).await 1267 }; 1268 assert!( 1269 ok && refs.contains("refs/heads/main"), 1270 "the pushed ref must be advertised to an anonymous reader:\n{refs}" 1271 ); 1272 1273 let wrong = format!( 1274 "http.extraHeader=Authorization: Basic {}", 1275 base64::engine::general_purpose::STANDARD.encode(format!( 1276 "x-tangled-token:{}", 1277 service_jwt("sh.tangled.repo.create", "git-http-push-neg") 1278 )) 1279 ); 1280 std::fs::write(work.join("README.md"), "second write\n").unwrap(); 1281 let (ok, out) = git(&work, &env, &["commit", "-q", "-am", "second"]); 1282 assert!(ok, "second commit failed:\n{out}"); 1283 let (ok, out) = { 1284 let work = work.clone(); 1285 let env = env.clone(); 1286 let url = url.clone(); 1287 let wrong = wrong.clone(); 1288 in_git_blocking(move || git(&work, &env, &["-c", &wrong, "push", "-q", &url, "main"])).await 1289 }; 1290 assert!( 1291 !ok, 1292 "a token bound to another method mustn't authorize a push:\n{out}" 1293 ); 1294} 1295 1296fn build_pack(work: &Path, env: &[(String, String)], tip: &str) -> Vec<u8> { 1297 let mut command = knot_fixtures::command(work); 1298 command 1299 .args(["pack-objects", "--revs", "--stdout", "--delta-base-offset"]) 1300 .stdin(Stdio::piped()) 1301 .stdout(Stdio::piped()) 1302 .stderr(Stdio::piped()); 1303 env.iter().for_each(|(key, value)| { 1304 command.env(key, value); 1305 }); 1306 let mut child = command.spawn().expect("git pack-objects spawns"); 1307 child 1308 .stdin 1309 .take() 1310 .unwrap() 1311 .write_all(format!("{tip}\n").as_bytes()) 1312 .unwrap(); 1313 let out = child.wait_with_output().unwrap(); 1314 assert!( 1315 out.status.success(), 1316 "git pack-objects failed: {}", 1317 String::from_utf8_lossy(&out.stderr) 1318 ); 1319 out.stdout 1320} 1321 1322fn build_receive_body(tip: &str, pack: &[u8]) -> Bytes { 1323 let mut command = format!("{} {tip} refs/heads/main", "0".repeat(tip.len())).into_bytes(); 1324 command.push(0); 1325 command.extend_from_slice(b"report-status side-band-64k agent=knot-h3-test/0"); 1326 command.push(b'\n'); 1327 let mut body = common::pkt(&command); 1328 body.extend_from_slice(b"0000"); 1329 body.extend_from_slice(pack); 1330 Bytes::from(body) 1331} 1332 1333#[tokio::test(flavor = "multi_thread", worker_threads = 4)] 1334async fn git_push_over_h3_authenticates_and_lands_the_ref() { 1335 let scratch = tempfile::tempdir().unwrap(); 1336 let (_key_path, public_line) = keygen(scratch.path()); 1337 let world = spawn(public_line, true).await; 1338 let edge = world.h3.as_ref().expect("the h3 edge is stood up"); 1339 1340 let path_env = std::env::var("PATH").unwrap_or_default(); 1341 let home = scratch.path().to_str().unwrap().to_string(); 1342 let env: Vec<(String, String)> = [("PATH", &path_env), ("HOME", &home)] 1343 .map(|(key, value)| (key.to_string(), value.clone())) 1344 .to_vec(); 1345 1346 let work = scratch.path().join("work"); 1347 std::fs::create_dir_all(&work).unwrap(); 1348 let (ok, out) = git(&work, &env, &["init", "-q", "-b", "main"]); 1349 assert!(ok, "init failed:\n{out}"); 1350 std::fs::write(work.join("README.md"), "hello over http3\n").unwrap(); 1351 let (ok, out) = git(&work, &env, &["add", "-A"]); 1352 assert!(ok, "add failed:\n{out}"); 1353 let (ok, out) = git(&work, &env, &["commit", "-q", "-m", "init over http3"]); 1354 assert!(ok, "commit failed:\n{out}"); 1355 let (ok, tip) = git(&work, &env, &["rev-parse", "HEAD"]); 1356 assert!(ok, "rev-parse failed:\n{tip}"); 1357 let tip = tip.trim().to_string(); 1358 let pack = build_pack(&work, &env, &tip); 1359 1360 let advert_uri = 1361 format!("https://localhost/{OWNER_DID}/{REPO_NAME}/info/refs?service=git-receive-pack"); 1362 let receive_uri = format!("https://localhost/{OWNER_DID}/{REPO_NAME}/git-receive-pack"); 1363 let warmup = 1364 format!("https://localhost/{OWNER_DID}/{REPO_NAME}/info/refs?service=git-upload-pack"); 1365 const RECEIVE_CT: &str = "application/x-git-receive-pack-request"; 1366 1367 let (status, _) = 1368 common::h3_request(edge, Method::GET, advert_uri.clone(), &[], None, None).await; 1369 assert_eq!( 1370 status, 1371 http::StatusCode::UNAUTHORIZED, 1372 "an unauthenticated receive advertisement must be challenged over h3" 1373 ); 1374 1375 let token = basic_auth(&service_jwt("sh.tangled.repo.push", "git-h3-adv-1")); 1376 let (status, advert) = common::h3_request( 1377 edge, 1378 Method::GET, 1379 advert_uri, 1380 &[("authorization", token.as_str())], 1381 None, 1382 None, 1383 ) 1384 .await; 1385 assert_eq!( 1386 status, 1387 http::StatusCode::OK, 1388 "an authenticated receive advertisement is served over h3" 1389 ); 1390 assert!( 1391 String::from_utf8_lossy(&advert).contains("# service=git-receive-pack"), 1392 "the h3 receive advertisement includes the service banner" 1393 ); 1394 1395 let body = build_receive_body(&tip, &pack); 1396 let (status, _) = common::h3_request( 1397 edge, 1398 Method::POST, 1399 receive_uri.clone(), 1400 &[("content-type", RECEIVE_CT)], 1401 Some(body.clone()), 1402 Some(warmup.as_str()), 1403 ) 1404 .await; 1405 assert_eq!( 1406 status, 1407 http::StatusCode::UNAUTHORIZED, 1408 "an unauthenticated receive-pack post must be challenged over h3" 1409 ); 1410 1411 let wrong = basic_auth(&service_jwt("sh.tangled.repo.create", "git-h3-neg")); 1412 let (status, _) = common::h3_request( 1413 edge, 1414 Method::POST, 1415 receive_uri.clone(), 1416 &[ 1417 ("content-type", RECEIVE_CT), 1418 ("authorization", wrong.as_str()), 1419 ], 1420 Some(body.clone()), 1421 Some(warmup.as_str()), 1422 ) 1423 .await; 1424 assert_eq!( 1425 status, 1426 http::StatusCode::UNAUTHORIZED, 1427 "a token bound to another method cannot authorize a receive-pack over h3" 1428 ); 1429 1430 let good = basic_auth(&service_jwt("sh.tangled.repo.push", "git-h3-push-1")); 1431 let (status, report) = common::h3_request( 1432 edge, 1433 Method::POST, 1434 receive_uri, 1435 &[ 1436 ("content-type", RECEIVE_CT), 1437 ("authorization", good.as_str()), 1438 ], 1439 Some(body), 1440 Some(warmup.as_str()), 1441 ) 1442 .await; 1443 assert_eq!(status, http::StatusCode::OK, "the authenticated h3 push"); 1444 let report = String::from_utf8_lossy(&report); 1445 assert!( 1446 report.contains("unpack ok"), 1447 "the pack must unpack cleanly over h3:\n{report}" 1448 ); 1449 assert!( 1450 report.contains("ok refs/heads/main"), 1451 "the ref update must be accepted over h3:\n{report}" 1452 ); 1453 1454 let repo = world.layout.open(&RepoDid::new(REPO_DID).unwrap()).unwrap(); 1455 let landed = repo.references().unwrap().into_iter().any(|record| { 1456 record.name.as_str() == "refs/heads/main" && record.target.to_string() == tip 1457 }); 1458 assert!( 1459 landed, 1460 "the ref pushed over h3 must be durable in the bare repo" 1461 ); 1462}