This repository has no description
1use std::io::Read;
2use std::sync::Arc;
3use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
4use std::time::Instant;
5
6#[global_allocator]
7static GLOBAL: tikv_jemallocator::Jemalloc = tikv_jemallocator::Jemalloc;
8
9const PAGE: u64 = 4096;
10
11fn rss_bytes() -> u64 {
12 let statm = std::fs::read_to_string("/proc/self/statm").unwrap();
13 statm
14 .split_whitespace()
15 .nth(1)
16 .and_then(|pages| pages.parse::<u64>().ok())
17 .map(|pages| pages * PAGE)
18 .unwrap()
19}
20
21fn vm_hwm_bytes() -> u64 {
22 let status = std::fs::read_to_string("/proc/self/status").unwrap();
23 status
24 .lines()
25 .find_map(|line| line.strip_prefix("VmHWM:"))
26 .and_then(|rest| rest.split_whitespace().next())
27 .and_then(|kb| kb.parse::<u64>().ok())
28 .map(|kb| kb * 1024)
29 .unwrap()
30}
31
32fn mib(bytes: u64) -> u64 {
33 bytes / (1024 * 1024)
34}
35
36fn pktline(data: &[u8]) -> Vec<u8> {
37 let mut out = format!("{:04x}", 4 + data.len()).into_bytes();
38 out.extend_from_slice(data);
39 out
40}
41
42fn main() {
43 let max_threads = std::env::var("KNOT_MAX_THREADS")
44 .ok()
45 .and_then(|value| value.parse::<usize>().ok());
46 knot_resource::init(knot_resource::Ceilings {
47 max_threads: max_threads.map(knot_resource::ThreadCount::new),
48 max_memory: None,
49 });
50
51 let pack_path = std::path::PathBuf::from(
52 std::env::args()
53 .nth(1)
54 .expect("usage: receive_mem <pack-file>"),
55 );
56 let pack_size = std::fs::metadata(&pack_path).map(|m| m.len()).unwrap_or(0);
57 println!("pack: {} MiB on disk", mib(pack_size));
58
59 let mut preamble = pktline(
60 b"0000000000000000000000000000000000000000 \
61 1111111111111111111111111111111111111111 refs/heads/main\0report-status\n",
62 );
63 preamble.extend_from_slice(b"0000");
64
65 let dir = tempfile::tempdir().expect("tempdir");
66 let limit = knot_pack::MaxWireBytes::new(16 * 1024 * 1024 * 1024);
67 let mut receiver = knot_pack::PackReceiver::new(
68 dir.path(),
69 limit,
70 knot_pack::PackLimits::default(),
71 gix::hash::Kind::Sha1,
72 )
73 .expect("receiver");
74
75 let peak = Arc::new(AtomicU64::new(0));
76 let stop = Arc::new(AtomicBool::new(false));
77 let sampler = {
78 let peak = Arc::clone(&peak);
79 let stop = Arc::clone(&stop);
80 std::thread::spawn(move || {
81 while !stop.load(Ordering::Relaxed) {
82 peak.fetch_max(rss_bytes(), Ordering::Relaxed);
83 std::thread::sleep(std::time::Duration::from_millis(5));
84 }
85 })
86 };
87
88 let start = Instant::now();
89 receiver.write(&preamble).expect("preamble");
90 let mut file = std::fs::File::open(&pack_path).expect("open pack");
91 let mut buf = vec![0u8; 1 << 20];
92 loop {
93 let read = file.read(&mut buf).expect("read pack");
94 if read == 0 {
95 break;
96 }
97 receiver.write(&buf[..read]).expect("receive");
98 }
99 let received = receiver.finish().expect("finish");
100 let receipt_elapsed = start.elapsed();
101 let receipt_peak = peak.load(Ordering::Relaxed);
102 println!(
103 "receipt: {:.1}s, peak {} MiB (pack streamed to a temp file, never resident)",
104 receipt_elapsed.as_secs_f64(),
105 mib(receipt_peak)
106 );
107
108 let staged = received.open_pack().expect("open pack").expect("has pack");
109 let objects = tempfile::tempdir().expect("objects tempdir");
110 let result =
111 knot_pack::bench_admit_and_ingest(objects.path(), staged.path(), gix::hash::Kind::Sha1);
112 let total_elapsed = start.elapsed();
113 let total_peak = peak.load(Ordering::Relaxed);
114 stop.store(true, Ordering::Relaxed);
115 sampler.join().ok();
116
117 match result {
118 Ok(folded) => {
119 println!("ADMITTED, fold engaged: {folded}");
120 println!("receive + ingest: {:.1}s", total_elapsed.as_secs_f64());
121 println!("receive + ingest peak: {} MiB", mib(total_peak));
122 println!("VmHWM,kernel peak: {} MiB", mib(vm_hwm_bytes()));
123 }
124 Err(error) => {
125 println!("REFUSED cleanly before the fold: {error:?}");
126 println!("peak at refusal: {} MiB", mib(total_peak));
127 }
128 }
129}