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