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 = path.display().to_string();
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, _pages_in_core_before) = if O::MUTATES_RESIDENCY {
267        let r: PM = crate::mincore::residency(&pf.mmap, pf.len)?;
268        let count = r.count_filled();
269        (Some(r), count)
270    } else {
271        (None, 0)
272    };
273
274    let ctx = FileContext {
275        file: &pf.file,
276        path,
277        mmap: std::sync::Arc::clone(&pf.mmap),
278        offset: pf.offset,
279        len: pf.len,
280        on_progress: None,
281        residency: residency.as_ref(),
282    };
283
284    let output = op.execute(&ctx)?;
285
286    let pages_in_core_after = if O::MUTATES_RESIDENCY {
287        if O::ACTION_SIGN >= 0 {
288            pf.total_pages
289        } else {
290            0
291        }
292    } else {
293        counts_page_count::<PM>(&pf.file, &pf.mmap, pf.offset, pf.len)?
294    };
295
296    Ok(Some(ops::CountsResult {
297        output,
298        total_pages: pf.total_pages,
299        pages_in_core_after,
300    }))
301}
302
303pub(crate) fn skip_process_file<O: Op, PM: PageMap + Sync>(
304    op: &O,
305    path: &Path,
306    range: &FileRange,
307    prepared: Option<PreparedFile>,
308) -> crate::Result<Option<ops::SkipResult<O::Output>>> {
309    let pf = match prepared {
310        Some(pf) => pf,
311        None => {
312            let Some(pf) = prepare_file(path, range)? else {
313                return Ok(None);
314            };
315            pf
316        }
317    };
318
319    let ctx = FileContext {
320        file: &pf.file,
321        path,
322        mmap: std::sync::Arc::clone(&pf.mmap),
323        offset: pf.offset,
324        len: pf.len,
325        on_progress: None,
326        residency: None::<&PM>,
327    };
328
329    let output = op.execute(&ctx)?;
330
331    Ok(Some(ops::SkipResult {
332        output,
333        total_pages: pf.total_pages,
334    }))
335}
336
337fn cli_record_stats<O: Op>(result: &impl FileProcessed<Output = O::Output>, stats: &Stats) {
338    let action = O::action_pages(
339        result.output_ref(),
340        result.total_pages(),
341        result.pages_in_core_before(),
342        result.pages_in_core_after(),
343    );
344    let signed_action = action as isize * O::ACTION_SIGN;
345    let initial = (result.pages_in_core_after() as isize - signed_action).max(0) as usize;
346    stats
347        .total_pages
348        .fetch_add(result.total_pages(), Ordering::Relaxed);
349    stats
350        .initial_pages_in_core
351        .fetch_add(initial, Ordering::Relaxed);
352    stats.action_pages.fetch_add(action, Ordering::Relaxed);
353    stats.total_files.fetch_add(1, Ordering::Relaxed);
354}
355
356#[allow(unused_variables)]
357fn counts_page_count<PM: PageMap>(
358    file: &std::fs::File,
359    mmap: &memmap2::Mmap,
360    offset: u64,
361    len: usize,
362) -> crate::Result<usize> {
363    #[cfg(target_os = "linux")]
364    if *crate::cachestat::SUPPORTED {
365        use std::os::unix::io::AsFd;
366        return Ok(crate::cachestat::cached_pages(file.as_fd(), offset, len as u64)?.try_into()?);
367    }
368    let residency: PM = crate::mincore::residency(mmap, len)?;
369    Ok(residency.count_filled())
370}