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 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
160pub 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}