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