This repository has no description
0

Configure Feed

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

core / knot2 / crates / knot-pack / tests / quarantine.rs
18 kB 583 lines
1use std::path::Path; 2use std::time::Instant; 3 4use knot_cob::{CobHome, CobStore}; 5use knot_cobs::{Grant, MembersChange}; 6use knot_git::{Layout, Repo}; 7use knot_pack::{ 8 MaxWireBytes, PackLimits, PackReceiver, ReceiveCommand, ReceiveFramer, ReceiveGuard, 9 RefDecision, receive_pack_guarded, receive_request_complete, sweep_incoming, 10}; 11use knot_runtime::{K256Signer, SeededEntropy, Signer}; 12use knot_types::{AccountDid, ActorId, ObjectFormat, Oid, RepoDid, UnixSeconds}; 13 14const SHA1: ObjectFormat = ObjectFormat::SHA1; 15 16mod common; 17use common::{commit, must, pack_objects, receive_request}; 18 19fn limits() -> PackLimits { 20 PackLimits::default() 21} 22 23fn seeded_pack(layout: &Layout, did: &RepoDid) -> (Repo, String, Vec<u8>) { 24 let bare = layout.create(did).unwrap(); 25 let work_dir = tempfile::tempdir().unwrap(); 26 let work = work_dir.path(); 27 must(work, &["init", "-q", "-b", "main"]); 28 commit(work, "a.txt", "x\n", "c1"); 29 let c1 = must(work, &["rev-parse", "HEAD"]); 30 let oids: Vec<String> = must(work, &["rev-list", "--objects", &c1]) 31 .lines() 32 .map(|line| line.split_whitespace().next().unwrap().to_string()) 33 .collect(); 34 let pack = pack_objects(work, &oids); 35 (bare, c1, pack) 36} 37 38fn report_text(report: &[u8]) -> String { 39 String::from_utf8_lossy(report).replace('\0', "") 40} 41 42fn dir_count(path: &Path) -> usize { 43 std::fs::read_dir(path) 44 .map(|entries| entries.filter_map(Result::ok).count()) 45 .unwrap_or(0) 46} 47 48fn live_object_count(repo: &Repo) -> usize { 49 let objects = repo.objects_dir(); 50 let loose: usize = std::fs::read_dir(&objects) 51 .map(|entries| { 52 entries 53 .filter_map(Result::ok) 54 .filter(|entry| { 55 entry 56 .file_name() 57 .to_str() 58 .is_some_and(|name| name.len() == 2) 59 && entry.path().is_dir() 60 }) 61 .map(|shard| dir_count(&shard.path())) 62 .sum() 63 }) 64 .unwrap_or(0); 65 let packs = std::fs::read_dir(objects.join("pack")) 66 .map(|entries| { 67 entries 68 .filter_map(Result::ok) 69 .filter(|entry| entry.path().extension().is_some_and(|ext| ext == "pack")) 70 .count() 71 }) 72 .unwrap_or(0); 73 loose + packs 74} 75 76struct AllowPublic; 77 78impl ReceiveGuard for AllowPublic { 79 fn authorize(&self, _staged: &Repo, commands: &[ReceiveCommand]) -> Vec<RefDecision> { 80 commands 81 .iter() 82 .map(|command| { 83 if command.name().is_some_and(knot_git::is_public_ref) { 84 RefDecision::Allow 85 } else { 86 RefDecision::Reject("not public ref".to_string()) 87 } 88 }) 89 .collect() 90 } 91} 92 93struct DenyAll; 94 95impl ReceiveGuard for DenyAll { 96 fn authorize(&self, _staged: &Repo, commands: &[ReceiveCommand]) -> Vec<RefDecision> { 97 commands 98 .iter() 99 .map(|_| RefDecision::Reject("unauthorized".to_string())) 100 .collect() 101 } 102} 103 104struct AllowAll; 105 106impl ReceiveGuard for AllowAll { 107 fn authorize(&self, _staged: &Repo, commands: &[ReceiveCommand]) -> Vec<RefDecision> { 108 commands.iter().map(|_| RefDecision::Allow).collect() 109 } 110} 111 112struct CobVerify { 113 owner: ActorId, 114 home: CobHome, 115} 116 117impl ReceiveGuard for CobVerify { 118 fn authorize(&self, staged: &Repo, commands: &[ReceiveCommand]) -> Vec<RefDecision> { 119 commands 120 .iter() 121 .map(|command| { 122 let store = CobStore::new(staged); 123 let Ok(name) = knot_types::RefName::new(command.refname()) else { 124 return RefDecision::Reject("invalid ref name".to_string()); 125 }; 126 match knot_cobs::verify_cob_ref(&store, &self.home, &name, &self.owner) { 127 Ok(_) => RefDecision::Allow, 128 Err(error) => RefDecision::Reject(error.to_string()), 129 } 130 }) 131 .collect() 132 } 133} 134 135#[test] 136fn a_denied_push_leaves_no_objects_in_the_live_odb() { 137 let scan = tempfile::tempdir().unwrap(); 138 let layout = Layout::new(scan.path()); 139 let did = RepoDid::new("did:plc:squid").unwrap(); 140 let (bare, c1, pack) = seeded_pack(&layout, &did); 141 142 assert_eq!( 143 live_object_count(&bare), 144 0, 145 "fresh bare repo has no objects" 146 ); 147 let request = receive_request("refs/heads/main", &Oid::null().to_hex(), &c1, &pack); 148 let report = receive_pack_guarded( 149 &bare, 150 &request, 151 &limits(), 152 &DenyAll, 153 &|_| {}, 154 &knot_pack::default_catalog().reject, 155 ) 156 .unwrap() 157 .report; 158 let report = report_text(&report); 159 assert!( 160 report.contains("ng refs/heads/main unauthorized"), 161 "{report}" 162 ); 163 164 assert!( 165 bare.references().unwrap().is_empty(), 166 "denied push mustn't create the ref" 167 ); 168 assert_eq!( 169 live_object_count(&bare), 170 0, 171 "denied push must leave no objects in the live odb" 172 ); 173 assert!(!bare.contains(Oid::from_hex(&c1).unwrap())); 174} 175 176#[test] 177fn an_authorized_push_migrates_objects_then_a_stale_push_is_rejected() { 178 let scan = tempfile::tempdir().unwrap(); 179 let layout = Layout::new(scan.path()); 180 let did = RepoDid::new("did:plc:squid").unwrap(); 181 let (bare, c1, pack) = seeded_pack(&layout, &did); 182 183 let report = receive_pack_guarded( 184 &bare, 185 &receive_request("refs/heads/main", &Oid::null().to_hex(), &c1, &pack), 186 &limits(), 187 &AllowPublic, 188 &|_| {}, 189 &knot_pack::default_catalog().reject, 190 ) 191 .unwrap() 192 .report; 193 let report = report_text(&report); 194 assert!(report.contains("unpack ok"), "{report}"); 195 assert!(report.contains("ok refs/heads/main"), "{report}"); 196 let refs = bare.references().unwrap(); 197 assert_eq!(refs.len(), 1); 198 assert_eq!(refs[0].name.as_str(), "refs/heads/main"); 199 assert_eq!(refs[0].target, Oid::from_hex(&c1).unwrap()); 200 assert!(bare.contains(Oid::from_hex(&c1).unwrap())); 201 let after_first = live_object_count(&bare); 202 203 let wrong_old = "1".repeat(40); 204 let fresh = "2".repeat(40); 205 let report = receive_pack_guarded( 206 &bare, 207 &receive_request("refs/heads/main", &wrong_old, &fresh, b""), 208 &limits(), 209 &AllowPublic, 210 &|_| {}, 211 &knot_pack::default_catalog().reject, 212 ) 213 .unwrap() 214 .report; 215 let report = report_text(&report); 216 assert!(report.contains("ng refs/heads/main"), "{report}"); 217 assert_eq!( 218 bare.find_ref(&knot_types::RefName::new("refs/heads/main").unwrap()) 219 .unwrap(), 220 Some(Oid::from_hex(&c1).unwrap()), 221 "stale push mustn't move the ref" 222 ); 223 assert_eq!( 224 live_object_count(&bare), 225 after_first, 226 "stale push mustn't add objects to the live odb" 227 ); 228} 229 230#[test] 231fn a_cob_ref_verifies_against_the_owner_key_and_is_refused_for_a_stranger() { 232 let scan = tempfile::tempdir().unwrap(); 233 let layout = Layout::new(scan.path()); 234 let signer = K256Signer::generate(&SeededEntropy::new(7)); 235 236 let source_did = RepoDid::new("did:plc:source").unwrap(); 237 let home = CobHome::from(&source_did); 238 let source = layout.create(&source_did).unwrap(); 239 let store = CobStore::new(&source); 240 let grant = Grant { 241 subject: AccountDid::new("did:plc:nel").unwrap(), 242 added_by: AccountDid::new("did:plc:nel").unwrap(), 243 created_at: UnixSeconds::new(1), 244 }; 245 let created = store 246 .create( 247 &home, 248 &MembersChange::Add(grant), 249 &signer, 250 UnixSeconds::new(1), 251 ) 252 .unwrap(); 253 let tip = created.tip.oid().to_hex(); 254 let refname = format!( 255 "refs/cobs/sh.tangled.knot.member/{}", 256 created.object.oid().to_hex() 257 ); 258 259 let reachable: Vec<String> = must(source.path(), &["rev-list", "--objects", &tip]) 260 .lines() 261 .map(|line| line.split_whitespace().next().unwrap().to_string()) 262 .collect(); 263 let pack = pack_objects(source.path(), &reachable); 264 let owner = ActorId::from_secp256k1(signer.public_key().as_bytes()); 265 266 let dest = layout 267 .create(&RepoDid::new("did:plc:dest").unwrap()) 268 .unwrap(); 269 let report = receive_pack_guarded( 270 &dest, 271 &receive_request(&refname, &Oid::null().to_hex(), &tip, &pack), 272 &limits(), 273 &CobVerify { 274 owner: owner.clone(), 275 home: home.clone(), 276 }, 277 &|_| {}, 278 &knot_pack::default_catalog().reject, 279 ) 280 .unwrap() 281 .report; 282 assert!( 283 report_text(&report).contains(&format!("ok {refname}")), 284 "owner-signed COB ref must be accepted at the receive boundary: {}", 285 report_text(&report) 286 ); 287 288 let stranger_key = K256Signer::generate(&SeededEntropy::new(9)); 289 let stranger = ActorId::from_secp256k1(stranger_key.public_key().as_bytes()); 290 let dest2 = layout 291 .create(&RepoDid::new("did:plc:dest2").unwrap()) 292 .unwrap(); 293 let report = receive_pack_guarded( 294 &dest2, 295 &receive_request(&refname, &Oid::null().to_hex(), &tip, &pack), 296 &limits(), 297 &CobVerify { 298 owner: stranger, 299 home: home.clone(), 300 }, 301 &|_| {}, 302 &knot_pack::default_catalog().reject, 303 ) 304 .unwrap() 305 .report; 306 assert!( 307 report_text(&report).contains(&format!("ng {refname}")), 308 "COB ref not signed by the resolved owner key must be refused: {}", 309 report_text(&report) 310 ); 311 assert_eq!( 312 live_object_count(&dest2), 313 0, 314 "refused COB ref push leaves no objects behind" 315 ); 316} 317 318#[test] 319fn a_reserved_ref_update_is_refused_even_when_the_guard_allows_it() { 320 let scan = tempfile::tempdir().unwrap(); 321 let layout = Layout::new(scan.path()); 322 let did = RepoDid::new("did:plc:squid").unwrap(); 323 let bare = layout.create(&did).unwrap(); 324 325 let cob_ref = format!("refs/cobs/sh.tangled.knot.member/{}", "a".repeat(40)); 326 let old = "1".repeat(40); 327 let new = "2".repeat(40); 328 let report = receive_pack_guarded( 329 &bare, 330 &receive_request(&cob_ref, &old, &new, b""), 331 &limits(), 332 &AllowAll, 333 &|_| {}, 334 &knot_pack::default_catalog().reject, 335 ) 336 .unwrap() 337 .report; 338 let report = report_text(&report); 339 340 assert!( 341 report.contains(&format!("ng {cob_ref}")), 342 "non-create update to a reserved ref must be refused even under an allow-all guard: {report}" 343 ); 344 assert!( 345 report.contains("cannot be modified"), 346 "refusal must name the create-only rule, not connectivity or compare-and-swap: {report}" 347 ); 348 assert!( 349 bare.references().unwrap().is_empty(), 350 "refused reserved-ref update must land nothing" 351 ); 352 assert_eq!( 353 live_object_count(&bare), 354 0, 355 "refused reserved-ref update must migrate no objects" 356 ); 357} 358 359fn pseudo_random(seed: u64, len: usize) -> Vec<u8> { 360 let mut state = seed.wrapping_mul(0x9E37_79B9_7F4A_7C15).wrapping_add(1); 361 (0..len) 362 .map(|_| { 363 state ^= state << 13; 364 state ^= state >> 7; 365 state ^= state << 17; 366 (state >> 24) as u8 367 }) 368 .collect() 369} 370 371fn incompressible_pack(megabytes: usize) -> Vec<u8> { 372 let work_dir = tempfile::tempdir().unwrap(); 373 let work = work_dir.path(); 374 must(work, &["init", "-q", "-b", "main"]); 375 (0..megabytes).for_each(|i| { 376 let blob = pseudo_random(i as u64 + 1, 1024 * 1024); 377 std::fs::write(work.join(format!("blob-{i:04}.bin")), blob).unwrap(); 378 }); 379 must(work, &["add", "-A"]); 380 must(work, &["commit", "-q", "-m", "bulk"]); 381 let head = must(work, &["rev-parse", "HEAD"]); 382 let oids: Vec<String> = must(work, &["rev-list", "--objects", &head]) 383 .lines() 384 .map(|line| line.split_whitespace().next().unwrap().to_string()) 385 .collect(); 386 pack_objects(work, &oids) 387} 388 389#[test] 390fn the_receive_framer_scans_incrementally_in_linear_time() { 391 const READ_CHUNK: usize = 64 * 1024; 392 let pack = incompressible_pack(24); 393 let body = receive_request( 394 "refs/heads/main", 395 &Oid::null().to_hex(), 396 &"1".repeat(40), 397 &pack, 398 ); 399 assert!( 400 pack.len() > 8 * 1024 * 1024, 401 "pack must be large enough to expose any quadratic scaling: {} bytes", 402 pack.len() 403 ); 404 405 let single_start = Instant::now(); 406 let complete = ReceiveFramer::new(limits(), SHA1.kind()) 407 .advance_bytes(&body) 408 .unwrap(); 409 let single = single_start.elapsed(); 410 assert_eq!( 411 complete, 412 Some(body.len()), 413 "one advance over the full body frames whole request" 414 ); 415 416 let chunks = body.len().div_ceil(READ_CHUNK); 417 let chunked_start = Instant::now(); 418 let mut framer = ReceiveFramer::new(limits(), SHA1.kind()); 419 let mut framed_at = None; 420 (1..=chunks).for_each(|n| { 421 let end = (n * READ_CHUNK).min(body.len()); 422 if framed_at.is_none() 423 && let Some(total) = framer.advance_bytes(&body[..end]).unwrap() 424 { 425 framed_at = Some(total); 426 } 427 }); 428 let chunked = chunked_start.elapsed(); 429 assert_eq!( 430 framed_at, 431 Some(body.len()), 432 "feeding the same body in 64KB chunks frames it at identical length" 433 ); 434 435 let ratio = chunked.as_secs_f64() / single.as_secs_f64().max(1e-6); 436 assert!( 437 ratio < 5.0, 438 "resumable framer never re-inflates settled objects, so the {chunks}-chunk feed \ 439 must stay within a small constant of one pass; got {ratio:.2}x" 440 ); 441} 442 443#[test] 444fn receive_request_completion_framing_and_the_per_object_limit() { 445 let scan = tempfile::tempdir().unwrap(); 446 let layout = Layout::new(scan.path()); 447 let did = RepoDid::new("did:plc:squid").unwrap(); 448 let (_bare, c1, pack) = seeded_pack(&layout, &did); 449 let request = receive_request("refs/heads/main", &Oid::null().to_hex(), &c1, &pack); 450 451 assert_eq!( 452 receive_request_complete(&request[..request.len() / 2], &limits(), SHA1.kind()).unwrap(), 453 None, 454 "half-delivered request isn't yet complete" 455 ); 456 assert_eq!( 457 receive_request_complete(&request, &limits(), SHA1.kind()).unwrap(), 458 Some(request.len()), 459 "whole request reports its exact length" 460 ); 461 let mut trailing = request.clone(); 462 trailing.extend_from_slice(b"junk-after-the-pack"); 463 assert_eq!( 464 receive_request_complete(&trailing, &limits(), SHA1.kind()).unwrap(), 465 Some(request.len()), 466 "framer stops at the pack trailer, ignoring trailing bytes" 467 ); 468 469 let tight = PackLimits { 470 max_object_bytes: knot_pack::MaxObjectBytes::new(1), 471 ..PackLimits::default() 472 }; 473 assert!( 474 receive_request_complete(&request, &tight, SHA1.kind()).is_err(), 475 "object beyond the per-object limit is refused by the framer" 476 ); 477} 478 479#[test] 480fn sweep_incoming_removes_abandoned_staging_directories() { 481 let scan = tempfile::tempdir().unwrap(); 482 let layout = Layout::new(scan.path()); 483 let did = RepoDid::new("did:plc:squid").unwrap(); 484 let bare = layout.create(&did).unwrap(); 485 486 let staging = bare.path().join(".knot-incoming.4242.0"); 487 std::fs::create_dir_all(staging.join("objects")).unwrap(); 488 std::fs::write(staging.join("objects").join("leftover"), b"x").unwrap(); 489 assert!(staging.exists()); 490 491 let swept = sweep_incoming(scan.path()); 492 assert_eq!(swept, 1, "exactly one staging directory is swept"); 493 assert!(!staging.exists(), "abandoned staging directory is gone"); 494 assert!( 495 bare.path().join("objects").exists(), 496 "live repository is untouched" 497 ); 498} 499 500#[test] 501fn streamed_receipt_reassembles_a_request_byte_for_byte() { 502 let scan = tempfile::tempdir().unwrap(); 503 let layout = Layout::new(scan.path()); 504 let did = RepoDid::new("did:plc:squid").unwrap(); 505 let (_bare, c1, pack) = seeded_pack(&layout, &did); 506 let request = receive_request("refs/heads/main", &Oid::null().to_hex(), &c1, &pack); 507 508 let dir = tempfile::tempdir().unwrap(); 509 let mut receiver = PackReceiver::new( 510 dir.path(), 511 MaxWireBytes::new(1 << 30), 512 limits(), 513 SHA1.kind(), 514 ) 515 .unwrap(); 516 let completed = request 517 .chunks(7) 518 .try_fold(false, |_, chunk| receiver.write(chunk)) 519 .unwrap(); 520 assert!( 521 completed, 522 "the framer must detect completion within the streamed request" 523 ); 524 525 let body = receiver.finish().unwrap(); 526 let mut reconstructed = body.preamble().to_vec(); 527 if let Some(pack) = body.open_pack().unwrap() { 528 reconstructed.extend_from_slice(&std::fs::read(pack.path()).unwrap()); 529 } 530 assert_eq!( 531 &reconstructed[..], 532 &request[..], 533 "the streamed body must be byte-identical to the buffered request" 534 ); 535 assert_eq!( 536 receive_request_complete(&request, &limits(), SHA1.kind()).unwrap(), 537 Some(request.len()), 538 "streamed total must agree with the buffered framer" 539 ); 540} 541 542#[test] 543fn streamed_receipt_handles_empty_and_delete_only_requests() { 544 let scan = tempfile::tempdir().unwrap(); 545 let layout = Layout::new(scan.path()); 546 let did = RepoDid::new("did:plc:squid").unwrap(); 547 let (_bare, c1, _pack) = seeded_pack(&layout, &did); 548 let dir = tempfile::tempdir().unwrap(); 549 550 let empty = PackReceiver::new( 551 dir.path(), 552 MaxWireBytes::new(1 << 30), 553 limits(), 554 SHA1.kind(), 555 ) 556 .unwrap() 557 .finish() 558 .unwrap(); 559 assert!( 560 empty.is_empty(), 561 "a stream with no bytes yields an empty body" 562 ); 563 564 let delete = receive_request("refs/heads/main", &c1, &Oid::null().to_hex(), &[]); 565 let mut receiver = PackReceiver::new( 566 dir.path(), 567 MaxWireBytes::new(1 << 30), 568 limits(), 569 SHA1.kind(), 570 ) 571 .unwrap(); 572 let completed = delete 573 .chunks(5) 574 .try_fold(false, |_, chunk| receiver.write(chunk)) 575 .unwrap(); 576 assert!(completed, "a delete-only push completes without a pack"); 577 let body = receiver.finish().unwrap(); 578 assert_eq!(body.preamble(), &delete[..]); 579 assert!( 580 body.open_pack().unwrap().is_none(), 581 "a delete-only push has no pack" 582 ); 583}