Skip to main content

pagers_core/ops/
mod.rs

1mod evict;
2mod lock;
3mod lockall;
4mod process;
5mod query;
6mod touch;
7
8use std::fs::File;
9use std::path::Path;
10use std::sync::Arc;
11use std::sync::atomic::AtomicUsize;
12
13use memmap2::Mmap;
14
15use crate::Cancellation;
16use crate::mincore::{DefaultPageMap, PageMap};
17
18pub use evict::Evict;
19pub use lock::{Lock, LockedFile};
20pub use lockall::Lockall;
21pub use process::{CountsResult, FullResult, file_info};
22pub(crate) use process::{PreparedFile, prepare_file};
23pub use query::Query;
24pub use touch::Touch;
25
26#[derive(Debug, Clone, Copy, PartialEq, Eq)]
27#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
28pub enum ResidencyEffect {
29    Preserve,
30    Populate,
31    EvictAdvisory,
32}
33
34impl ResidencyEffect {
35    pub const fn action_sign(self) -> isize {
36        match self {
37            Self::Preserve => 0,
38            Self::Populate => 1,
39            Self::EvictAdvisory => -1,
40        }
41    }
42
43    pub const fn has_action(self) -> bool {
44        !matches!(self, Self::Preserve)
45    }
46
47    pub const fn progress_resident(self) -> Option<bool> {
48        match self {
49            Self::Preserve => None,
50            Self::Populate => Some(true),
51            Self::EvictAdvisory => Some(false),
52        }
53    }
54
55    pub fn action_pages(self, before: usize, after: usize) -> usize {
56        match self {
57            Self::Preserve => 0,
58            Self::Populate => after.saturating_sub(before),
59            Self::EvictAdvisory => before.saturating_sub(after),
60        }
61    }
62}
63
64pub trait Op: Sync {
65    const LABEL: &str;
66    const EFFECT: ResidencyEffect;
67
68    type Output: Send;
69    fn execute<PM: PageMap + Sync>(&self, ctx: &FileContext<'_, PM>)
70    -> crate::Result<Self::Output>;
71
72    fn finish(&self) -> crate::Result<()> {
73        Ok(())
74    }
75}
76
77pub struct FileContext<'a, PM: PageMap = DefaultPageMap> {
78    prepared: PreparedFile,
79    cancellation: Cancellation,
80    on_progress: Option<&'a (dyn Fn(usize, usize) + Sync)>,
81    residency: Option<&'a PM>,
82}
83
84impl<'a, PM: PageMap> FileContext<'a, PM> {
85    pub(crate) fn new(prepared: PreparedFile, cancellation: Cancellation) -> Self {
86        Self {
87            prepared,
88            cancellation,
89            on_progress: None,
90            residency: None,
91        }
92    }
93
94    pub(crate) fn with_progress(
95        mut self,
96        on_progress: Option<&'a (dyn Fn(usize, usize) + Sync)>,
97    ) -> Self {
98        self.on_progress = on_progress;
99        self
100    }
101
102    pub(crate) fn with_residency(mut self, residency: Option<&'a PM>) -> Self {
103        self.residency = residency;
104        self
105    }
106
107    pub fn file(&self) -> &File {
108        &self.prepared.file
109    }
110
111    pub fn path(&self) -> &Path {
112        &self.prepared.path
113    }
114
115    pub fn mmap(&self) -> &Mmap {
116        &self.prepared.mmap
117    }
118
119    pub fn mapping(&self) -> Arc<Mmap> {
120        Arc::clone(&self.prepared.mmap)
121    }
122
123    pub fn offset(&self) -> u64 {
124        self.prepared.offset()
125    }
126
127    pub fn len(&self) -> usize {
128        self.prepared.len()
129    }
130
131    pub fn is_empty(&self) -> bool {
132        self.prepared.mmap.is_empty()
133    }
134
135    pub fn total_pages(&self) -> usize {
136        self.prepared.total_pages()
137    }
138
139    pub fn residency(&self) -> Option<&PM> {
140        self.residency
141    }
142
143    pub fn report_progress(&self, pages_walked: usize, action_count: usize) {
144        if let Some(on_progress) = self.on_progress {
145            on_progress(pages_walked, action_count);
146        }
147    }
148
149    pub fn check_cancelled(&self) -> crate::Result<()> {
150        self.cancellation.check()
151    }
152}
153
154impl<PM: PageMap> std::fmt::Debug for FileContext<'_, PM> {
155    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
156        f.debug_struct("FileContext")
157            .field("path", &self.path())
158            .field("offset", &self.offset())
159            .field("len", &self.len())
160            .finish()
161    }
162}
163
164#[derive(Debug, Clone, Copy, PartialEq, Eq)]
165#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
166pub struct FileRange {
167    offset: u64,
168    max_len: Option<u64>,
169}
170
171impl FileRange {
172    pub fn full() -> Self {
173        Self {
174            offset: 0,
175            max_len: None,
176        }
177    }
178
179    pub fn new(offset: u64, max_len: Option<u64>) -> crate::Result<Self> {
180        let range = Self { offset, max_len };
181        range.validate()?;
182        Ok(range)
183    }
184
185    pub fn offset(&self) -> u64 {
186        self.offset
187    }
188
189    pub fn max_len(&self) -> Option<u64> {
190        self.max_len
191    }
192
193    pub(crate) fn validate(&self) -> crate::Result<()> {
194        let page_size = *crate::pagesize::PAGE_SIZE as u64;
195        if !self.offset.is_multiple_of(page_size) {
196            return Err(crate::Error::UnalignedRange {
197                offset: self.offset,
198                page_size,
199            });
200        }
201        if self.max_len == Some(0) {
202            return Err(crate::Error::EmptyRange);
203        }
204        if let Some(max_len) = self.max_len
205            && self.offset.checked_add(max_len).is_none()
206        {
207            return Err(crate::Error::RangeOverflow {
208                offset: self.offset,
209                max_len,
210            });
211        }
212        Ok(())
213    }
214}
215
216#[derive(Debug, Clone, PartialEq)]
217#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
218pub struct FileInfo<PM = DefaultPageMap> {
219    pub total_pages: usize,
220    pub residency: PM,
221}
222
223#[derive(Debug, Default)]
224pub struct Stats {
225    pub total_pages: AtomicUsize,
226    pub initial_pages_in_core: AtomicUsize,
227    pub action_pages: AtomicUsize,
228    pub total_files: AtomicUsize,
229    pub total_dirs: AtomicUsize,
230}
231
232impl Stats {
233    pub fn new() -> Self {
234        Self::default()
235    }
236}
237
238#[cfg(test)]
239mod tests {
240    use std::io::Write;
241    use std::sync::Arc;
242
243    use crate::Cancellation;
244
245    use super::{FileContext, FileRange, Op, ResidencyEffect, Touch, prepare_file};
246
247    #[test]
248    fn residency_effect_derives_sign_and_delta() {
249        assert_eq!(ResidencyEffect::Preserve.action_sign(), 0);
250        assert_eq!(ResidencyEffect::Preserve.action_pages(3, 7), 0);
251
252        assert_eq!(ResidencyEffect::Populate.action_sign(), 1);
253        assert_eq!(ResidencyEffect::Populate.action_pages(3, 7), 4);
254
255        assert_eq!(ResidencyEffect::EvictAdvisory.action_sign(), -1);
256        assert_eq!(ResidencyEffect::EvictAdvisory.action_pages(7, 3), 4);
257        assert_eq!(ResidencyEffect::EvictAdvisory.action_pages(3, 7), 0);
258    }
259
260    #[test]
261    fn file_context_from_prepared_file_preserves_mapping() {
262        let mut file = tempfile::NamedTempFile::new().unwrap();
263        file.write_all(&[0; 64]).unwrap();
264        let prepared = prepare_file(
265            file.path(),
266            &FileRange {
267                offset: 0,
268                max_len: None,
269            },
270        )
271        .unwrap()
272        .unwrap();
273        let mmap = Arc::clone(&prepared.mmap);
274
275        let context: FileContext<'_, Vec<bool>> = FileContext::new(prepared, Cancellation::new());
276
277        assert_eq!(context.len(), mmap.len());
278        assert_eq!(context.path(), file.path());
279        assert_eq!(context.total_pages(), 1);
280    }
281
282    #[test]
283    fn touch_stops_when_cancelled_during_page_walk() {
284        let mut file = tempfile::NamedTempFile::new().unwrap();
285        file.write_all(&vec![0; *crate::pagesize::PAGE_SIZE * 512])
286            .unwrap();
287        let prepared = prepare_file(file.path(), &FileRange::full())
288            .unwrap()
289            .unwrap();
290        let cancellation = Cancellation::new();
291        let context_cancellation = cancellation.clone();
292        let cancel_after_progress = |_: usize, _: usize| {
293            cancellation.cancel();
294        };
295        let context: FileContext<'_, Vec<bool>> = FileContext::new(prepared, context_cancellation)
296            .with_progress(Some(&cancel_after_progress));
297
298        let result = Touch.execute(&context);
299
300        assert!(matches!(result, Err(crate::Error::Cancelled)));
301    }
302}