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