This repository has no description
0

Configure Feed

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

core / knot2 / crates / knot-lfs / tests / soak.rs
5.3 kB 160 lines
1mod common; 2 3use std::io::Write; 4 5use common::{ 6 GROWTH_SLACK, PEAK_CEILING, download_script, incompressible, oid_of, repo, rss_bytes, 7 upload_script, 8}; 9use knot_lfs::{ 10 ClaimedSize, DiskStore, FreeSpaceFloor, LfsOid, LfsSize, LfsStore, LfsStorePath, 11 StoreAdmission, TransferOp, serve_transfer, 12}; 13 14const OBJECT_BYTES: usize = 8 * 1024 * 1024; 15const SEEDS: usize = 4; 16const WRITERS: u64 = 4; 17const READERS: usize = 4; 18const ROUNDS: u64 = 3; 19 20struct CountingSink(u64); 21 22impl Write for CountingSink { 23 fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> { 24 self.0 += buf.len() as u64; 25 Ok(buf.len()) 26 } 27 28 fn flush(&mut self) -> std::io::Result<()> { 29 Ok(()) 30 } 31} 32 33#[test] 34fn sustained_concurrent_transfers_stay_bounded_and_leak_nothing() { 35 let dir = tempfile::tempdir().unwrap(); 36 let store = DiskStore::open(LfsStorePath::new(dir.path())).unwrap(); 37 let admission = StoreAdmission::new( 38 LfsStorePath::new(dir.path()), 39 LfsSize::new(u64::MAX), 40 FreeSpaceFloor::new(0), 41 ); 42 43 let seeded: Vec<LfsOid> = (0..SEEDS) 44 .map(|seed| { 45 let body = incompressible(OBJECT_BYTES, 0x5eed_0000 + seed as u64); 46 let oid = oid_of(&body); 47 store 48 .put( 49 &repo(), 50 &oid, 51 ClaimedSize::new(body.len() as u64), 52 &mut &body[..], 53 ) 54 .unwrap(); 55 oid 56 }) 57 .collect(); 58 59 let storm = |round: u64| { 60 std::thread::scope(|scope| { 61 let writers: Vec<_> = (0..WRITERS) 62 .map(|writer| { 63 let store = &store; 64 let admission = &admission; 65 scope.spawn(move || { 66 let body = 67 incompressible(OBJECT_BYTES, 0xfeed_0000 + round * WRITERS + writer); 68 let (oid, script) = upload_script(&body); 69 let mut out = Vec::new(); 70 serve_transfer( 71 store, 72 admission, 73 &repo(), 74 TransferOp::Upload, 75 &knot_messages::default_catalog().lfs, 76 &script[..], 77 &mut out, 78 ) 79 .unwrap(); 80 let replies = String::from_utf8_lossy(&out); 81 assert!( 82 !replies.contains("status 4") && !replies.contains("status 5"), 83 "round {round} writer {writer}: upload session failed:\n{replies}" 84 ); 85 oid 86 }) 87 }) 88 .collect(); 89 let readers: Vec<_> = (0..READERS) 90 .map(|reader| { 91 let store = &store; 92 let admission = &admission; 93 let oid = seeded[reader % SEEDS].clone(); 94 scope.spawn(move || { 95 let script = download_script(&oid); 96 let mut sink = CountingSink(0); 97 serve_transfer( 98 store, 99 admission, 100 &repo(), 101 TransferOp::Download, 102 &knot_messages::default_catalog().lfs, 103 &script[..], 104 &mut sink, 105 ) 106 .unwrap(); 107 assert!( 108 sink.0 >= OBJECT_BYTES as u64, 109 "round {round} reader {reader}: streamed {} bytes", 110 sink.0 111 ); 112 }) 113 }) 114 .collect(); 115 let uploaded: Vec<LfsOid> = writers 116 .into_iter() 117 .map(|writer| writer.join().unwrap()) 118 .collect(); 119 readers.into_iter().for_each(|reader| { 120 reader.join().unwrap(); 121 }); 122 uploaded 123 }) 124 }; 125 126 let mut uploaded = storm(0); 127 let settled = rss_bytes(); 128 129 let peaks: Vec<u64> = (1..ROUNDS) 130 .map(|round| { 131 uploaded.extend(storm(round)); 132 rss_bytes() 133 }) 134 .collect(); 135 136 let peak = peaks.iter().copied().max().unwrap_or(settled); 137 assert!( 138 peak < PEAK_CEILING, 139 "concurrent transfers peaked at {peak} bytes, ceiling {PEAK_CEILING}" 140 ); 141 let last = *peaks.last().unwrap_or(&settled); 142 assert!( 143 last <= settled + GROWTH_SLACK, 144 "rss grew from {settled} to {last} across rounds, transfers are leaking" 145 ); 146 147 uploaded.iter().chain(seeded.iter()).for_each(|oid| { 148 assert!( 149 store.probe(&repo(), oid).unwrap().is_some(), 150 "object {oid} must be readable after the concurrent rounds" 151 ); 152 }); 153 let leftover: Vec<_> = std::fs::read_dir(dir.path().join(".incoming")) 154 .unwrap() 155 .collect(); 156 assert!( 157 leftover.is_empty(), 158 "temporary upload files leaked: {leftover:?}" 159 ); 160}