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 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, _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}