This repository has no description
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}