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::events::{Event, EventSink};
6use crate::mincore::{DefaultPageMap, PageMap};
7use crate::ops::{
8    self, FileContext, FileProcessed, FileRange, Op, PreparedFile, Stats, prepare_file,
9};
10
11pub trait DisplayMode<PM: PageMap = DefaultPageMap>: Sync {
12    fn process_one<O: Op>(
13        &self,
14        op: &O,
15        path: &Path,
16        range: &FileRange,
17        stats: &Stats,
18    ) -> Option<O::Output>;
19
20    fn finish(&self) {}
21}
22
23pub struct Tui<PM: PageMap = DefaultPageMap> {
24    sink: EventSink<PM>,
25}
26
27impl<PM: PageMap> Tui<PM> {
28    pub fn new(sender: Sender<Event<PM>>) -> Self {
29        Self {
30            sink: EventSink::new(sender),
31        }
32    }
33}
34
35impl<PM: PageMap + Clone + Send + Sync> DisplayMode<PM> for Tui<PM> {
36    fn process_one<O: Op>(
37        &self,
38        op: &O,
39        path: &Path,
40        range: &FileRange,
41        stats: &Stats,
42    ) -> Option<O::Output> {
43        let path_str: std::sync::Arc<str> = path.display().to_string().into();
44        let full_file = FileRange {
45            offset: 0,
46            max_len: None,
47        };
48
49        let pf = match prepare_file(path, &full_file) {
50            Ok(Some(pf)) => pf,
51            Ok(None) => return None,
52            Err(e) => {
53                tracing::warn!("{}: {e}", path.display());
54                return None;
55            }
56        };
57        let residency: PM = match crate::mincore::residency(&pf.mmap, pf.len) {
58            Ok(r) => r,
59            Err(e) => {
60                tracing::warn!("{}: {e}", path.display());
61                return None;
62            }
63        };
64        let pages_in_core = residency.count_filled();
65        let total_pages = pf.total_pages;
66
67        stats.total_files.fetch_add(1, Ordering::Relaxed);
68        stats.total_pages.fetch_add(total_pages, Ordering::Relaxed);
69        stats
70            .initial_pages_in_core
71            .fetch_add(pages_in_core, Ordering::Relaxed);
72        self.sink.send(Event::FileStart {
73            path: path_str.clone(),
74            total_pages,
75            residency: residency.clone(),
76        });
77
78        let prepared_pf = if range.is_full_file() { Some(pf) } else { None };
79
80        if O::SKIP_RESIDENCY {
81            let result = match skip_process_file::<O, PM>(op, path, range, prepared_pf) {
82                Ok(Some(r)) => r,
83                Ok(None) => return None,
84                Err(e) => {
85                    tracing::warn!("{e}");
86                    return None;
87                }
88            };
89
90            let total_action = O::action_pages(
91                result.output_ref(),
92                result.total_pages(),
93                result.pages_in_core_before(),
94                result.pages_in_core_after(),
95            );
96            stats
97                .action_pages
98                .fetch_add(total_action, Ordering::Relaxed);
99
100            self.sink.send(Event::FileProgress {
101                path: path_str.clone(),
102                page_offset: 0,
103                pages_walked: total_pages,
104                resident: O::ACTION_SIGN >= 0,
105            });
106            self.sink.send(Event::FileDone { path: path_str });
107
108            return Some(result.into_output());
109        }
110
111        let page_offset = range.offset as usize / *crate::pagesize::PAGE_SIZE;
112        let reported_action = std::sync::atomic::AtomicUsize::new(0);
113        let on_progress = |pages_walked: usize, action_count: usize| {
114            let action = action_count;
115            let delta = action - reported_action.swap(action, Ordering::Relaxed);
116            stats.action_pages.fetch_add(delta, Ordering::Relaxed);
117            self.sink.send(Event::FileProgress {
118                path: path_str.clone(),
119                page_offset,
120                pages_walked,
121                resident: O::ACTION_SIGN >= 0,
122            });
123        };
124
125        let prepared_full = prepared_pf.map(|pf| (pf, residency, pages_in_core));
126        let result =
127            match full_process_file::<O, PM>(op, path, range, Some(&on_progress), prepared_full) {
128                Ok(Some(r)) => r,
129                Ok(None) => return None,
130                Err(e) => {
131                    tracing::warn!("{e}");
132                    return None;
133                }
134            };
135
136        // Flush remaining action_pages not covered by the progress callback.
137        let reported = reported_action.load(Ordering::Relaxed);
138        let total_action = O::action_pages(
139            result.output_ref(),
140            result.total_pages(),
141            result.pages_in_core_before(),
142            result.pages_in_core_after(),
143        );
144        stats
145            .action_pages
146            .fetch_add(total_action - reported, Ordering::Relaxed);
147
148        self.sink.send(Event::FileDone { path: path_str });
149
150        Some(result.into_output())
151    }
152
153    fn finish(&self) {
154        self.sink.send(Event::AllDone);
155    }
156}
157
158pub struct Cli;
159
160// Marker ZSTs for run-mode dispatch
161pub struct TuiMode;
162pub struct CliMode;
163pub struct Daemon;
164pub struct NoDaemon;
165
166impl<PM: PageMap + Send + Sync> DisplayMode<PM> for Cli {
167    fn process_one<O: Op>(
168        &self,
169        op: &O,
170        path: &Path,
171        range: &FileRange,
172        stats: &Stats,
173    ) -> Option<O::Output> {
174        if O::SKIP_RESIDENCY {
175            let result = match skip_process_file::<O, PM>(op, path, range, None) {
176                Ok(Some(r)) => r,
177                Ok(None) => return None,
178                Err(e) => {
179                    tracing::warn!("{e}");
180                    return None;
181                }
182            };
183            cli_record_stats::<O>(&result, stats);
184            tracing::info!("{}: {} pages", path.display(), result.total_pages());
185            return Some(result.into_output());
186        }
187
188        let result = match counts_process_file::<O, PM>(op, path, range) {
189            Ok(Some(r)) => r,
190            Ok(None) => return None,
191            Err(e) => {
192                tracing::warn!("{e}");
193                return None;
194            }
195        };
196        cli_record_stats::<O>(&result, stats);
197        tracing::info!(
198            "{}: {}/{} pages resident",
199            path.display(),
200            result.pages_in_core_after(),
201            result.total_pages(),
202        );
203        Some(result.into_output())
204    }
205}
206
207pub(crate) fn full_process_file<O: Op, PM: PageMap + Sync>(
208    op: &O,
209    path: &Path,
210    range: &FileRange,
211    on_progress: Option<&(dyn Fn(usize, usize) + Sync)>,
212    prepared: Option<(PreparedFile, PM, usize)>,
213) -> crate::Result<Option<ops::FullResult<O::Output, PM>>> {
214    let (pf, residency_before, pages_in_core_before) = match prepared {
215        Some(tuple) => tuple,
216        None => {
217            let Some(pf) = prepare_file(path, range)? else {
218                return Ok(None);
219            };
220            let residency_before: PM = crate::mincore::residency(&pf.mmap, pf.len)?;
221            let pages_in_core_before = residency_before.count_filled();
222            (pf, residency_before, pages_in_core_before)
223        }
224    };
225
226    let ctx = FileContext {
227        file: &pf.file,
228        path,
229        mmap: std::sync::Arc::clone(&pf.mmap),
230        offset: pf.offset,
231        len: pf.len,
232        on_progress,
233        residency: Some(&residency_before),
234    };
235
236    let output = op.execute(&ctx)?;
237
238    let (pages_in_core_after, residency_before, residency_after) = if O::MUTATES_RESIDENCY {
239        let fill = O::ACTION_SIGN >= 0;
240        let after = PM::from_bools(std::iter::repeat_n(fill, pf.total_pages));
241        let count = if fill { pf.total_pages } else { 0 };
242        (count, Some(residency_before), Some(after))
243    } else {
244        (pages_in_core_before, None, Some(residency_before))
245    };
246
247    Ok(Some(ops::FullResult {
248        output,
249        total_pages: pf.total_pages,
250        pages_in_core_before,
251        pages_in_core_after,
252        residency_before,
253        residency_after,
254    }))
255}
256
257pub(crate) fn counts_process_file<O: Op, PM: PageMap + Sync>(
258    op: &O,
259    path: &Path,
260    range: &FileRange,
261) -> crate::Result<Option<ops::CountsResult<O::Output>>> {
262    let Some(pf) = prepare_file(path, range)? else {
263        return Ok(None);
264    };
265
266    let residency: Option<PM> = if O::MUTATES_RESIDENCY {
267        Some(crate::mincore::residency(&pf.mmap, pf.len)?)
268    } else {
269        None
270    };
271
272    let ctx = FileContext {
273        file: &pf.file,
274        path,
275        mmap: std::sync::Arc::clone(&pf.mmap),
276        offset: pf.offset,
277        len: pf.len,
278        on_progress: None,
279        residency: residency.as_ref(),
280    };
281
282    let output = op.execute(&ctx)?;
283
284    let pages_in_core_after = if O::MUTATES_RESIDENCY {
285        if O::ACTION_SIGN >= 0 {
286            pf.total_pages
287        } else {
288            0
289        }
290    } else {
291        counts_page_count::<PM>(&pf.file, &pf.mmap, pf.offset, pf.len)?
292    };
293
294    Ok(Some(ops::CountsResult {
295        output,
296        total_pages: pf.total_pages,
297        pages_in_core_after,
298    }))
299}
300
301pub(crate) fn skip_process_file<O: Op, PM: PageMap + Sync>(
302    op: &O,
303    path: &Path,
304    range: &FileRange,
305    prepared: Option<PreparedFile>,
306) -> crate::Result<Option<ops::SkipResult<O::Output>>> {
307    let pf = match prepared {
308        Some(pf) => pf,
309        None => {
310            let Some(pf) = prepare_file(path, range)? else {
311                return Ok(None);
312            };
313            pf
314        }
315    };
316
317    let ctx = FileContext {
318        file: &pf.file,
319        path,
320        mmap: std::sync::Arc::clone(&pf.mmap),
321        offset: pf.offset,
322        len: pf.len,
323        on_progress: None,
324        residency: None::<&PM>,
325    };
326
327    let output = op.execute(&ctx)?;
328
329    Ok(Some(ops::SkipResult {
330        output,
331        total_pages: pf.total_pages,
332    }))
333}
334
335fn cli_record_stats<O: Op>(result: &impl FileProcessed<Output = O::Output>, stats: &Stats) {
336    let action = O::action_pages(
337        result.output_ref(),
338        result.total_pages(),
339        result.pages_in_core_before(),
340        result.pages_in_core_after(),
341    );
342    let signed_action = action as isize * O::ACTION_SIGN;
343    let initial = (result.pages_in_core_after() as isize - signed_action).max(0) as usize;
344    stats
345        .total_pages
346        .fetch_add(result.total_pages(), Ordering::Relaxed);
347    stats
348        .initial_pages_in_core
349        .fetch_add(initial, Ordering::Relaxed);
350    stats.action_pages.fetch_add(action, Ordering::Relaxed);
351    stats.total_files.fetch_add(1, Ordering::Relaxed);
352}
353
354#[allow(unused_variables)]
355fn counts_page_count<PM: PageMap>(
356    file: &std::fs::File,
357    mmap: &memmap2::Mmap,
358    offset: u64,
359    len: usize,
360) -> crate::Result<usize> {
361    #[cfg(target_os = "linux")]
362    if *crate::cachestat::SUPPORTED {
363        use std::os::unix::io::AsFd;
364        return Ok(crate::cachestat::cached_pages(file.as_fd(), offset, len as u64)?.try_into()?);
365    }
366    let residency: PM = crate::mincore::residency(mmap, len)?;
367    Ok(residency.count_filled())
368}