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 / examples / receive_mem.rs
4.2 kB 129 lines
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}