1use std::fs::File;
4use std::io::{self, Read, Seek, Write};
5use std::num::NonZeroU64;
6use std::sync::{Arc, Mutex};
7
8pub const DEFAULT_CACHE_RELEASE_WINDOW_BYTES: u64 = 1024 * 1024 * 1024;
10
11pub fn cache_release_window_for_streams(active_streams: usize) -> io::Result<NonZeroU64> {
17 let active_streams = u64::try_from(active_streams)
18 .map_err(|_| io::Error::other("cache-release stream count overflow"))?;
19 if active_streams == 0 {
20 return Err(io::Error::other("cache-release stream count is zero"));
21 }
22 let window = DEFAULT_CACHE_RELEASE_WINDOW_BYTES
23 .checked_div(active_streams)
24 .and_then(NonZeroU64::new)
25 .ok_or_else(|| io::Error::other("cache-release operation budget is exhausted"))?;
26 let aggregate = window
27 .get()
28 .checked_mul(active_streams)
29 .ok_or_else(|| io::Error::other("cache-release aggregate window overflow"))?;
30 if aggregate > DEFAULT_CACHE_RELEASE_WINDOW_BYTES {
31 return Err(io::Error::other(
32 "cache-release aggregate window exceeds operation budget",
33 ));
34 }
35 Ok(window)
36}
37
38pub fn validate_cache_release_operation_windows(windows: &[NonZeroU64]) -> io::Result<u64> {
45 if windows.is_empty() {
46 return Err(io::Error::other(
47 "cache-release operation has no active streams",
48 ));
49 }
50 let aggregate = windows.iter().try_fold(0_u64, |sum, window| {
51 sum.checked_add(window.get())
52 .ok_or_else(|| io::Error::other("cache-release aggregate window overflow"))
53 })?;
54 if aggregate > DEFAULT_CACHE_RELEASE_WINDOW_BYTES {
55 return Err(io::Error::other(
56 "cache-release aggregate window exceeds operation budget",
57 ));
58 }
59 Ok(aggregate)
60}
61
62#[derive(Debug, Clone, Copy, PartialEq, Eq)]
64pub enum FileCacheReleaseOutcome {
65 Released,
67 Unsupported,
69}
70
71#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
73pub struct FileCacheReleaseEvidence {
74 pub sync_operations: u64,
76 pub release_operations: u64,
78 pub unsupported_operations: u64,
80 pub released_bytes: u64,
82 pub peak_window_bytes: u64,
84}
85
86#[derive(Clone, Debug, Default)]
88pub struct FileCacheReleaseTracker {
89 state: Arc<Mutex<FileCacheReleaseTrackerState>>,
90}
91
92#[derive(Debug, Default)]
93struct FileCacheReleaseTrackerState {
94 evidence: FileCacheReleaseEvidence,
95 deferred_error: Option<(io::ErrorKind, String)>,
96}
97
98impl FileCacheReleaseTracker {
99 #[must_use]
101 pub fn evidence(&self) -> FileCacheReleaseEvidence {
102 self.state
103 .lock()
104 .unwrap_or_else(std::sync::PoisonError::into_inner)
105 .evidence
106 }
107
108 pub fn check_error(&self) -> io::Result<()> {
113 let mut state = self
114 .state
115 .lock()
116 .unwrap_or_else(std::sync::PoisonError::into_inner);
117 match state.deferred_error.take() {
118 Some((kind, message)) => Err(io::Error::new(kind, message)),
119 None => Ok(()),
120 }
121 }
122
123 fn account(&self, outcome: FileCacheReleaseOutcome, bytes: u64) {
124 let mut state = self
125 .state
126 .lock()
127 .unwrap_or_else(std::sync::PoisonError::into_inner);
128 state.evidence.peak_window_bytes = state.evidence.peak_window_bytes.max(bytes);
129 match outcome {
130 FileCacheReleaseOutcome::Released => {
131 state.evidence.release_operations =
132 state.evidence.release_operations.saturating_add(1);
133 state.evidence.released_bytes = state.evidence.released_bytes.saturating_add(bytes);
134 }
135 FileCacheReleaseOutcome::Unsupported => {
136 state.evidence.unsupported_operations =
137 state.evidence.unsupported_operations.saturating_add(1);
138 }
139 }
140 }
141
142 fn defer(&self, error: &io::Error) {
143 let mut state = self
144 .state
145 .lock()
146 .unwrap_or_else(std::sync::PoisonError::into_inner);
147 if state.deferred_error.is_none() {
148 state.deferred_error = Some((error.kind(), error.to_string()));
149 }
150 }
151}
152
153#[derive(Debug)]
155pub struct FileCacheReleasingReader {
156 file: File,
157 window_bytes: NonZeroU64,
158 pending_offset: u64,
159 pending_bytes: u64,
160 tracker: FileCacheReleaseTracker,
161 finished: bool,
162}
163
164impl FileCacheReleasingReader {
165 pub fn new(file: File) -> io::Result<Self> {
170 Self::with_window_bytes(
171 file,
172 NonZeroU64::new(DEFAULT_CACHE_RELEASE_WINDOW_BYTES)
173 .expect("default cache-release window is non-zero"),
174 FileCacheReleaseTracker::default(),
175 )
176 }
177
178 pub fn with_tracker(file: File, tracker: FileCacheReleaseTracker) -> io::Result<Self> {
183 Self::with_window_bytes(
184 file,
185 NonZeroU64::new(DEFAULT_CACHE_RELEASE_WINDOW_BYTES)
186 .expect("default cache-release window is non-zero"),
187 tracker,
188 )
189 }
190
191 pub fn with_window_bytes(
196 mut file: File,
197 window_bytes: NonZeroU64,
198 tracker: FileCacheReleaseTracker,
199 ) -> io::Result<Self> {
200 let pending_offset = file.stream_position()?;
201 Ok(Self {
202 file,
203 window_bytes,
204 pending_offset,
205 pending_bytes: 0,
206 tracker,
207 finished: false,
208 })
209 }
210
211 pub fn finish(&mut self) -> io::Result<FileCacheReleaseEvidence> {
216 self.release_pending()?;
217 self.finished = true;
218 self.tracker.check_error()?;
219 Ok(self.tracker.evidence())
220 }
221
222 #[must_use]
224 pub fn tracker(&self) -> FileCacheReleaseTracker {
225 self.tracker.clone()
226 }
227
228 #[must_use]
230 pub const fn window_bytes(&self) -> NonZeroU64 {
231 self.window_bytes
232 }
233
234 #[must_use]
236 pub const fn file(&self) -> &File {
237 &self.file
238 }
239
240 fn release_pending(&mut self) -> io::Result<()> {
241 if self.pending_bytes == 0 {
242 return Ok(());
243 }
244 let bytes = self.pending_bytes;
245 let end = self
246 .pending_offset
247 .checked_add(bytes)
248 .ok_or_else(|| io::Error::other("cache-release reader offset overflow"))?;
249 let outcome = release_file_cache(&self.file, self.pending_offset, bytes)?;
250 self.tracker.account(outcome, bytes);
251 self.pending_offset = end;
252 self.pending_bytes = 0;
253 Ok(())
254 }
255}
256
257impl Read for FileCacheReleasingReader {
258 fn read(&mut self, buffer: &mut [u8]) -> io::Result<usize> {
259 if buffer.is_empty() || self.finished {
260 return Ok(0);
261 }
262 if file_cache_release_supported() && self.pending_bytes == self.window_bytes.get() {
263 self.release_pending()?;
264 }
265 let limit = if file_cache_release_supported() {
266 usize::try_from(self.window_bytes.get() - self.pending_bytes)
267 .unwrap_or(usize::MAX)
268 .min(buffer.len())
269 } else {
270 buffer.len()
271 };
272 self.pending_offset
273 .checked_add(self.pending_bytes)
274 .and_then(|offset| offset.checked_add(u64::try_from(limit).ok()?))
275 .ok_or_else(|| io::Error::other("cache-release reader offset overflow"))?;
276 let read = self.file.read(&mut buffer[..limit])?;
277 self.pending_bytes = self
278 .pending_bytes
279 .checked_add(u64::try_from(read).map_err(io::Error::other)?)
280 .ok_or_else(|| io::Error::other("cache-release reader byte count overflow"))?;
281 if read == 0 {
282 self.release_pending()?;
283 self.finished = true;
284 }
285 Ok(read)
286 }
287}
288
289impl Seek for FileCacheReleasingReader {
290 fn seek(&mut self, position: io::SeekFrom) -> io::Result<u64> {
291 self.release_pending()?;
292 let offset = self.file.seek(position)?;
293 self.pending_offset = offset;
294 self.finished = false;
295 Ok(offset)
296 }
297}
298
299impl Drop for FileCacheReleasingReader {
300 fn drop(&mut self) {
301 if !self.finished
302 && let Err(error) = self.release_pending()
303 {
304 self.tracker.defer(&error);
305 }
306 }
307}
308
309#[derive(Debug)]
311pub struct DurableFileCacheWriter {
312 file: File,
313 window_bytes: NonZeroU64,
314 pending_offset: u64,
315 pending_bytes: u64,
316 evidence: FileCacheReleaseEvidence,
317}
318
319impl DurableFileCacheWriter {
320 pub fn new(file: File) -> io::Result<Self> {
325 Self::with_window_bytes(
326 file,
327 NonZeroU64::new(DEFAULT_CACHE_RELEASE_WINDOW_BYTES)
328 .expect("default cache-release window is non-zero"),
329 )
330 }
331
332 pub fn with_window_bytes(file: File, window_bytes: NonZeroU64) -> io::Result<Self> {
337 Self::with_window_bytes_checked(file, window_bytes, || Ok(()))
338 }
339
340 pub fn with_window_bytes_checked(
349 mut file: File,
350 window_bytes: NonZeroU64,
351 setup_check: impl FnOnce() -> io::Result<()>,
352 ) -> io::Result<Self> {
353 let pending_offset = file.stream_position()?;
354 setup_check()?;
355 Ok(Self {
356 file,
357 window_bytes,
358 pending_offset,
359 pending_bytes: 0,
360 evidence: FileCacheReleaseEvidence::default(),
361 })
362 }
363
364 pub fn sync_all_and_release(&mut self) -> io::Result<()> {
369 self.synchronize_pending(true)
370 }
371
372 #[must_use]
374 pub fn file(&self) -> &File {
375 &self.file
376 }
377
378 #[must_use]
380 pub const fn evidence(&self) -> FileCacheReleaseEvidence {
381 self.evidence
382 }
383
384 #[must_use]
386 pub const fn window_bytes(&self) -> NonZeroU64 {
387 self.window_bytes
388 }
389
390 #[must_use]
392 pub fn into_file(self) -> File {
393 self.file
394 }
395
396 fn synchronize_pending(&mut self, final_barrier: bool) -> io::Result<()> {
397 if self.pending_bytes == 0 && !final_barrier {
398 return Ok(());
399 }
400 if self.pending_bytes == 0 {
401 synchronize_file(&self.file)?;
402 self.evidence.sync_operations = self
403 .evidence
404 .sync_operations
405 .checked_add(1)
406 .ok_or_else(|| io::Error::other("cache-release sync count overflow"))?;
407 return Ok(());
408 }
409 let bytes = self.pending_bytes;
410 let end = self
411 .pending_offset
412 .checked_add(bytes)
413 .ok_or_else(|| io::Error::other("cache-release writer offset overflow"))?;
414 let outcome = synchronize_before_release(
415 || synchronize_file(&self.file),
416 || release_file_cache(&self.file, self.pending_offset, bytes),
417 )?;
418 self.evidence.sync_operations = self
419 .evidence
420 .sync_operations
421 .checked_add(1)
422 .ok_or_else(|| io::Error::other("cache-release sync count overflow"))?;
423 self.evidence.peak_window_bytes = self.evidence.peak_window_bytes.max(bytes);
424 match outcome {
425 FileCacheReleaseOutcome::Released => {
426 self.evidence.release_operations = self
427 .evidence
428 .release_operations
429 .checked_add(1)
430 .ok_or_else(|| io::Error::other("cache-release operation count overflow"))?;
431 self.evidence.released_bytes = self
432 .evidence
433 .released_bytes
434 .checked_add(bytes)
435 .ok_or_else(|| io::Error::other("cache-release byte count overflow"))?;
436 }
437 FileCacheReleaseOutcome::Unsupported => {
438 self.evidence.unsupported_operations = self
439 .evidence
440 .unsupported_operations
441 .checked_add(1)
442 .ok_or_else(|| io::Error::other("unsupported cache-release count overflow"))?;
443 }
444 }
445 self.pending_offset = end;
446 self.pending_bytes = 0;
447 Ok(())
448 }
449}
450
451fn synchronize_before_release<T>(
452 synchronize: impl FnOnce() -> io::Result<()>,
453 release: impl FnOnce() -> io::Result<T>,
454) -> io::Result<T> {
455 synchronize()?;
456 release()
457}
458
459impl Write for DurableFileCacheWriter {
460 fn write(&mut self, buffer: &[u8]) -> io::Result<usize> {
461 if file_cache_release_supported() && self.pending_bytes == self.window_bytes.get() {
462 self.synchronize_pending(false)?;
465 }
466 let limit = if file_cache_release_supported() {
467 let remaining = self.window_bytes.get() - self.pending_bytes;
468 buffer
469 .len()
470 .min(usize::try_from(remaining).unwrap_or(usize::MAX))
471 } else {
472 buffer.len()
473 };
474 self.pending_offset
475 .checked_add(self.pending_bytes)
476 .and_then(|offset| offset.checked_add(u64::try_from(limit).ok()?))
477 .ok_or_else(|| io::Error::other("cache-release writer offset overflow"))?;
478 let written = self.file.write(&buffer[..limit])?;
479 self.pending_bytes = self
480 .pending_bytes
481 .checked_add(u64::try_from(written).map_err(io::Error::other)?)
482 .ok_or_else(|| io::Error::other("cache-release writer byte count overflow"))?;
483 Ok(written)
484 }
485
486 fn flush(&mut self) -> io::Result<()> {
487 self.file.flush()
488 }
489}
490
491pub fn release_file_cache(
500 file: &File,
501 offset: u64,
502 bytes: u64,
503) -> io::Result<FileCacheReleaseOutcome> {
504 #[cfg(test)]
505 CACHE_RELEASE_FAILURE.with(|failure| {
506 if failure.replace(false) {
507 return Err(io::Error::other("injected cache-release failure"));
508 }
509 Ok(())
510 })?;
511 if bytes == 0 {
512 return Ok(if file_cache_release_supported() {
513 FileCacheReleaseOutcome::Released
514 } else {
515 FileCacheReleaseOutcome::Unsupported
516 });
517 }
518 release_file_cache_inner(file, offset, bytes)
519}
520
521fn synchronize_file(file: &File) -> io::Result<()> {
522 #[cfg(test)]
523 CACHE_SYNC_FAILURE.with(|failure| {
524 if failure.replace(false) {
525 return Err(io::Error::other("injected cache-sync failure"));
526 }
527 Ok(())
528 })?;
529 crate::ObservedSync::observed_sync_all(file)
530}
531
532#[cfg(test)]
533thread_local! {
534 static CACHE_RELEASE_FAILURE: std::cell::Cell<bool> = const {
535 std::cell::Cell::new(false)
536 };
537 static CACHE_SYNC_FAILURE: std::cell::Cell<bool> = const {
538 std::cell::Cell::new(false)
539 };
540}
541
542#[cfg(test)]
543fn inject_cache_release_failure() {
544 CACHE_RELEASE_FAILURE.with(|failure| failure.set(true));
545}
546
547#[cfg(test)]
548fn inject_cache_sync_failure() {
549 CACHE_SYNC_FAILURE.with(|failure| failure.set(true));
550}
551
552#[cfg(target_os = "linux")]
553const fn file_cache_release_supported() -> bool {
554 true
555}
556
557#[cfg(not(target_os = "linux"))]
558const fn file_cache_release_supported() -> bool {
559 false
560}
561
562#[cfg(target_os = "linux")]
563fn release_file_cache_inner(
564 file: &File,
565 offset: u64,
566 bytes: u64,
567) -> io::Result<FileCacheReleaseOutcome> {
568 rustix::fs::fadvise(
569 file,
570 offset,
571 NonZeroU64::new(bytes),
572 rustix::fs::Advice::DontNeed,
573 )
574 .map_err(io::Error::from)?;
575 Ok(FileCacheReleaseOutcome::Released)
576}
577
578#[cfg(not(target_os = "linux"))]
579#[expect(
580 clippy::unnecessary_wraps,
581 reason = "signature must match the fallible Linux implementation"
582)]
583fn release_file_cache_inner(
584 _file: &File,
585 _offset: u64,
586 _bytes: u64,
587) -> io::Result<FileCacheReleaseOutcome> {
588 Ok(FileCacheReleaseOutcome::Unsupported)
589}
590
591#[cfg(test)]
592mod tests;