Skip to main content

pagers_core/
mode.rs

1use std::path::Path;
2use std::sync::atomic::Ordering;
3use std::sync::mpsc::Sender;
4
5use crate::Cancellation;
6use crate::events::Event;
7use crate::mincore::{DefaultPageMap, PageMap};
8use crate::ops::{
9    self, FileContext, FileRange, Op, PreparedFile, ResidencyEffect, Stats, prepare_file,
10};
11
12pub trait DisplayMode<PM: PageMap = DefaultPageMap>: Sync {
13    fn process_one<O: Op>(
14        &self,
15        op: &O,
16        path: &Path,
17        range: &FileRange,
18        stats: &Stats,
19        cancellation: &Cancellation,
20    ) -> crate::Result<Option<O::Output>>;
21
22    fn finish(&self) {}
23}
24
25pub struct Tui<PM: PageMap = DefaultPageMap> {
26    sink: Sender<Event<PM>>,
27    next_id: std::sync::atomic::AtomicUsize,
28}
29
30impl<PM: PageMap> Tui<PM> {
31    pub fn new(sender: Sender<Event<PM>>) -> Self {
32        Self {
33            sink: sender,
34            next_id: std::sync::atomic::AtomicUsize::new(0),
35        }
36    }
37
38    fn send_residency(&self, id: usize, page_offset: usize, residency: &PM) {
39        let mut start = 0;
40        while start < residency.len() {
41            let resident = residency.is_set(start);
42            let mut end = start + 1;
43            while end < residency.len() && residency.is_set(end) == resident {
44                end += 1;
45            }
46            let _ = self.sink.send(Event::FileProgress {
47                id,
48                page_offset: page_offset + start,
49                pages_walked: end - start,
50                resident,
51            });
52            start = end;
53        }
54    }
55}
56
57impl<PM: PageMap + Clone + Send + Sync> DisplayMode<PM> for Tui<PM> {
58    fn process_one<O: Op>(
59        &self,
60        op: &O,
61        path: &Path,
62        range: &FileRange,
63        stats: &Stats,
64        cancellation: &Cancellation,
65    ) -> crate::Result<Option<O::Output>> {
66        let id = self.next_id.fetch_add(1, Ordering::Relaxed);
67        let path_str: std::sync::Arc<str> = path.display().to_string().into();
68
69        let pf = match prepare_file(path, range)? {
70            Some(pf) => pf,
71            None => return Ok(None),
72        };
73        let residency: PM = crate::mincore::residency(&pf.mmap, pf.len())?;
74        let pages_in_core = residency.count_filled();
75        let total_pages = pf.total_pages();
76
77        stats.total_files.fetch_add(1, Ordering::Relaxed);
78        stats.total_pages.fetch_add(total_pages, Ordering::Relaxed);
79        stats
80            .initial_pages_in_core
81            .fetch_add(pages_in_core, Ordering::Relaxed);
82        let _ = self.sink.send(Event::FileStart {
83            id,
84            path: path_str.clone(),
85            total_pages,
86            residency: residency.clone(),
87        });
88
89        let reported_action = std::sync::atomic::AtomicUsize::new(0);
90        let on_progress = |pages_walked: usize, action_count: usize| {
91            let action = action_count;
92            let delta = action.saturating_sub(reported_action.swap(action, Ordering::Relaxed));
93            stats.action_pages.fetch_add(delta, Ordering::Relaxed);
94            if let Some(resident) = O::EFFECT.progress_resident() {
95                let _ = self.sink.send(Event::FileProgress {
96                    id,
97                    page_offset: 0,
98                    pages_walked,
99                    resident,
100                });
101            }
102        };
103
104        let prepared = Some((pf, residency, pages_in_core));
105        let processed =
106            full_process_file::<O, PM>(op, path, range, Some(&on_progress), prepared, cancellation);
107        let result = match processed {
108            Ok(Some(result)) => result,
109            Ok(None) => {
110                let _ = self.sink.send(Event::FileDone { id });
111                return Ok(None);
112            }
113            Err(error) => {
114                let _ = self.sink.send(Event::FileDone { id });
115                return Err(error);
116            }
117        };
118
119        // Flush remaining action_pages not covered by the progress callback.
120        let reported = reported_action.load(Ordering::Relaxed);
121        let total_action =
122            O::EFFECT.action_pages(result.pages_in_core_before, result.pages_in_core_after);
123        stats
124            .action_pages
125            .fetch_add(total_action.saturating_sub(reported), Ordering::Relaxed);
126
127        if O::EFFECT.has_action()
128            && let Some(residency_after) = result.residency_after.as_ref()
129        {
130            self.send_residency(id, 0, residency_after);
131        }
132
133        let _ = self.sink.send(Event::FileDone { id });
134
135        Ok(Some(result.output))
136    }
137
138    fn finish(&self) {
139        let _ = self.sink.send(Event::AllDone);
140    }
141}
142
143pub struct Cli;
144
145impl<PM: PageMap + Send + Sync> DisplayMode<PM> for Cli {
146    fn process_one<O: Op>(
147        &self,
148        op: &O,
149        path: &Path,
150        range: &FileRange,
151        stats: &Stats,
152        cancellation: &Cancellation,
153    ) -> crate::Result<Option<O::Output>> {
154        let result = match counts_process_file::<O, PM>(op, path, range, cancellation)? {
155            Some(result) => result,
156            None => return Ok(None),
157        };
158        cli_record_stats::<O>(&result, stats);
159        tracing::info!(
160            "{}: {}/{} pages resident",
161            path.display(),
162            result.pages_in_core_after,
163            result.total_pages,
164        );
165        Ok(Some(result.output))
166    }
167}
168
169pub(crate) fn full_process_file<O: Op, PM: PageMap + Sync>(
170    op: &O,
171    path: &Path,
172    range: &FileRange,
173    on_progress: Option<&(dyn Fn(usize, usize) + Sync)>,
174    prepared: Option<(PreparedFile, PM, usize)>,
175    cancellation: &Cancellation,
176) -> crate::Result<Option<ops::FullResult<O::Output, PM>>> {
177    cancellation.check()?;
178    let (pf, residency_before, pages_in_core_before) = match prepared {
179        Some(tuple) => tuple,
180        None => {
181            let Some(pf) = prepare_file(path, range)? else {
182                return Ok(None);
183            };
184            let residency_before: PM = crate::mincore::residency(&pf.mmap, pf.len())?;
185            let pages_in_core_before = residency_before.count_filled();
186            (pf, residency_before, pages_in_core_before)
187        }
188    };
189
190    let ctx = FileContext::new(pf, cancellation.clone())
191        .with_progress(on_progress)
192        .with_residency(Some(&residency_before));
193
194    let output = op.execute(&ctx)?;
195    let total_pages = ctx.total_pages();
196
197    let (pages_in_core_after, residency_after) = match O::EFFECT {
198        ResidencyEffect::Preserve => (pages_in_core_before, None),
199        ResidencyEffect::Populate => {
200            let after = PM::from_bools(std::iter::repeat_n(true, total_pages));
201            (total_pages, Some(after))
202        }
203        ResidencyEffect::EvictAdvisory => {
204            let after: PM = crate::mincore::residency(ctx.mmap(), ctx.len())?;
205            let count = after.count_filled();
206            (count, Some(after))
207        }
208    };
209    drop(ctx);
210    let (residency_before, residency_after) = match O::EFFECT {
211        ResidencyEffect::Preserve => (None, Some(residency_before)),
212        ResidencyEffect::Populate | ResidencyEffect::EvictAdvisory => {
213            (Some(residency_before), residency_after)
214        }
215    };
216
217    Ok(Some(ops::FullResult {
218        output,
219        total_pages,
220        pages_in_core_before,
221        pages_in_core_after,
222        residency_before,
223        residency_after,
224    }))
225}
226
227pub(crate) fn counts_process_file<O: Op, PM: PageMap + Sync>(
228    op: &O,
229    path: &Path,
230    range: &FileRange,
231    cancellation: &Cancellation,
232) -> crate::Result<Option<ops::CountsResult<O::Output>>> {
233    cancellation.check()?;
234    let Some(pf) = prepare_file(path, range)? else {
235        return Ok(None);
236    };
237
238    let residency: Option<PM> = match O::EFFECT {
239        ResidencyEffect::Populate => Some(crate::mincore::residency(&pf.mmap, pf.len())?),
240        ResidencyEffect::Preserve | ResidencyEffect::EvictAdvisory => None,
241    };
242    let pages_in_core_before = match residency.as_ref() {
243        Some(residency) => residency.count_filled(),
244        None => counts_page_count::<PM>(&pf.file, &pf.mmap, pf.offset(), pf.len())?,
245    };
246
247    let ctx = FileContext::new(pf, cancellation.clone()).with_residency(residency.as_ref());
248
249    let output = op.execute(&ctx)?;
250    let total_pages = ctx.total_pages();
251
252    let pages_in_core_after = match O::EFFECT {
253        ResidencyEffect::Preserve => pages_in_core_before,
254        ResidencyEffect::Populate => total_pages,
255        ResidencyEffect::EvictAdvisory => {
256            counts_page_count::<PM>(ctx.file(), ctx.mmap(), ctx.offset(), ctx.len())?
257        }
258    };
259
260    Ok(Some(ops::CountsResult {
261        output,
262        total_pages,
263        pages_in_core_before,
264        pages_in_core_after,
265    }))
266}
267
268fn cli_record_stats<O: Op>(result: &ops::CountsResult<O::Output>, stats: &Stats) {
269    let action = O::EFFECT.action_pages(result.pages_in_core_before, result.pages_in_core_after);
270    stats
271        .total_pages
272        .fetch_add(result.total_pages, Ordering::Relaxed);
273    stats
274        .initial_pages_in_core
275        .fetch_add(result.pages_in_core_before, Ordering::Relaxed);
276    stats.action_pages.fetch_add(action, Ordering::Relaxed);
277    stats.total_files.fetch_add(1, Ordering::Relaxed);
278}
279
280#[allow(unused_variables)]
281fn counts_page_count<PM: PageMap>(
282    file: &std::fs::File,
283    mmap: &memmap2::Mmap,
284    offset: u64,
285    len: usize,
286) -> crate::Result<usize> {
287    #[cfg(target_os = "linux")]
288    if *crate::cachestat::SUPPORTED {
289        use std::os::unix::io::AsFd;
290        return Ok(crate::cachestat::cached_pages(file.as_fd(), offset, len as u64)?.try_into()?);
291    }
292    let residency: PM = crate::mincore::residency(mmap, len)?;
293    Ok(residency.count_filled())
294}
295
296#[cfg(test)]
297mod tests {
298    use super::*;
299
300    #[test]
301    fn tui_reports_only_the_requested_range() {
302        let file = tempfile::NamedTempFile::new().unwrap();
303        let page_size = *crate::pagesize::PAGE_SIZE;
304        file.as_file().set_len((page_size * 3) as u64).unwrap();
305        let range = FileRange::new(page_size as u64, Some(page_size as u64)).unwrap();
306        let stats = Stats::new();
307        let cancellation = Cancellation::new();
308        let (sender, receiver) = std::sync::mpsc::channel();
309        let tui = Tui::<DefaultPageMap>::new(sender);
310
311        tui.process_one(&ops::Query, file.path(), &range, &stats, &cancellation)
312            .unwrap();
313
314        let Event::FileStart {
315            total_pages,
316            residency,
317            ..
318        } = receiver.recv().unwrap()
319        else {
320            panic!("expected file start event");
321        };
322        assert_eq!(total_pages, 1);
323        assert_eq!(residency.len(), 1);
324        assert_eq!(stats.total_pages.load(Ordering::Relaxed), 1);
325    }
326}