Skip to main content

nxd_core/progress/
target.rs

1use std::collections::BTreeMap;
2use std::fs::File;
3use std::io::Write;
4use std::sync::{Arc, Mutex};
5use std::time::{Duration, Instant};
6
7use crate::progress::color::ColorMode;
8use crate::progress::log::{ProgressEvent, StatusLevel, normalize_multiline_message};
9use crate::progress::redaction::{ProgressRedactor, redact_sensitive_text};
10
11pub enum LogTarget {
12	Terminal {
13		file: Option<File>,
14		bulk: crate::progress::bulk::TerminalBulkState,
15	},
16	File(File),
17	Batch {
18		prefix: String,
19		file: File,
20		started_at: Option<Instant>,
21		last_event_at: Option<Instant>,
22		outcome: Option<BatchOutcome>,
23		/// A single-target run echoes subprocess output to the terminal so the
24		/// operator sees the same narrative the log records. Concurrent targets
25		/// would interleave line by line, so a fleet run leaves this off and keeps
26		/// the terminal to milestones.
27		echo: Option<crate::progress::bulk::TerminalBulkState>,
28	},
29	Routed {
30		fallback: Logger,
31		hosts: BTreeMap<String, Logger>,
32	},
33	Silent,
34}
35
36impl LogTarget {
37	pub fn log_event(&mut self, level: StatusLevel, message: &str) {
38		self.write_event(level, message, message);
39	}
40
41	pub fn log_status(&mut self, message: &str) {
42		self.write_event(StatusLevel::Info, message, message);
43	}
44
45	pub fn log_debug(&mut self, message: &str) {
46		self.write_event(StatusLevel::Debug, message, message);
47	}
48
49	fn write_event(&mut self, level: StatusLevel, message: &str, file_message: &str) {
50		let message = normalize_multiline_message(message);
51		let message = redact_sensitive_text(&message);
52		let file_message = normalize_multiline_message(file_message);
53		let file_message = redact_sensitive_text(&file_message);
54		let message_lines: Vec<&str> =
55			if message.is_empty() { vec![""] } else { message.lines().collect() };
56		let file_lines: Vec<&str> =
57			if file_message.is_empty() { vec![""] } else { file_message.lines().collect() };
58
59		match self {
60			LogTarget::Terminal { file, bulk } => {
61				if let Some(f) = file {
62					for line in file_lines {
63						let _ = writeln!(f, "[{}] {}", level.label(), line);
64					}
65				}
66				bulk.clear_active_counter();
67				bulk.flush_summary();
68				for line in message_lines {
69					let event = ProgressEvent::new(level, None, line);
70					event.log(ColorMode::Auto);
71				}
72			}
73			LogTarget::Batch { prefix, file, started_at, last_event_at, outcome, echo } => {
74				let now = Instant::now();
75				let started = *started_at.get_or_insert(now);
76				let elapsed = now.saturating_duration_since(started);
77				*last_event_at = Some(now);
78				match level {
79					StatusLevel::Success if *outcome != Some(BatchOutcome::Failed) => {
80						*outcome = Some(BatchOutcome::Succeeded);
81					}
82					StatusLevel::Failure | StatusLevel::Error => *outcome = Some(BatchOutcome::Failed),
83					_ => {}
84				}
85				// A milestone must not be garbled by an in-place counter still on screen.
86				if let Some(bulk) = echo.as_mut() {
87					bulk.clear_active_counter();
88					bulk.flush_summary();
89				}
90				let terminal_message = message.strip_prefix(&format!("{prefix}: ")).unwrap_or(&message);
91				let terminal_message = timed_batch_message(level, terminal_message, elapsed);
92				for line in terminal_message.lines() {
93					let event = ProgressEvent::new(level, Some(prefix), line);
94					event.log(ColorMode::Auto);
95				}
96				for line in file_lines {
97					let line = timed_batch_message(level, line, elapsed);
98					let _ = writeln!(file, "[{}] {}", level.label(), line);
99				}
100			}
101			LogTarget::File(file) => {
102				for line in file_lines {
103					let _ = writeln!(file, "[{}] {}", level.label(), line);
104				}
105			}
106			LogTarget::Routed { .. } => {}
107			LogTarget::Silent => {}
108		}
109	}
110}
111
112fn timed_batch_message(level: StatusLevel, message: &str, elapsed: Duration) -> String {
113	let elapsed = crate::progress::log::format_elapsed(elapsed);
114	match level {
115		StatusLevel::Success => format!("{message} in {elapsed}"),
116		StatusLevel::Failure | StatusLevel::Error => format!("{message} after {elapsed}"),
117		_ => message.to_string(),
118	}
119}
120
121#[derive(Clone, Copy, Debug, PartialEq, Eq)]
122pub enum BatchOutcome {
123	Succeeded,
124	Failed,
125}
126
127#[derive(Clone)]
128pub struct Logger {
129	pub target: Arc<Mutex<LogTarget>>,
130	redactor: Arc<Mutex<ProgressRedactor>>,
131}
132
133impl Logger {
134	pub fn new(target: Arc<Mutex<LogTarget>>) -> Self {
135		Self { target, redactor: Arc::new(Mutex::new(ProgressRedactor::default())) }
136	}
137
138	pub fn with_sensitive_values<'a>(&self, values: impl IntoIterator<Item = &'a [u8]>) -> Self {
139		Self {
140			target: self.target.clone(),
141			redactor: Arc::new(Mutex::new(ProgressRedactor::with_sensitive_values(values))),
142		}
143	}
144
145	pub fn sanitize(&self, message: &str) -> String {
146		self
147			.redactor
148			.lock()
149			.map(|mut redactor| redactor.redact(message))
150			.unwrap_or_else(|_| redact_sensitive_text(message))
151	}
152
153	pub fn lock(&self) -> std::sync::LockResult<std::sync::MutexGuard<'_, LogTarget>> {
154		self.target.lock()
155	}
156
157	pub fn silent() -> Self {
158		Self::new(Arc::new(Mutex::new(LogTarget::Silent)))
159	}
160
161	pub fn terminal() -> Self {
162		Self::new(Arc::new(Mutex::new(LogTarget::Terminal {
163			file: None,
164			bulk: crate::progress::bulk::TerminalBulkState::default(),
165		})))
166	}
167
168	pub fn terminal_with_file(file: File) -> Self {
169		Self::new(Arc::new(Mutex::new(LogTarget::Terminal {
170			file: Some(file),
171			bulk: crate::progress::bulk::TerminalBulkState::default(),
172		})))
173	}
174
175	pub fn file(file: File) -> Self {
176		Self::new(Arc::new(Mutex::new(LogTarget::File(file))))
177	}
178
179	pub fn batch(prefix: String, file: File) -> Self {
180		Self::batch_with_echo(prefix, file, false)
181	}
182
183	/// A batch logger that also echoes subprocess output to the terminal. Used for
184	/// a single-target run, where nothing else is writing to the terminal
185	/// concurrently, so the operator can watch the same narrative the log records.
186	pub fn batch_with_echo(prefix: String, file: File, echo: bool) -> Self {
187		Self::new(Arc::new(Mutex::new(LogTarget::Batch {
188			prefix,
189			file,
190			started_at: None,
191			last_event_at: None,
192			outcome: None,
193			echo: echo.then(crate::progress::bulk::TerminalBulkState::default),
194		})))
195	}
196
197	pub fn routed(fallback: Logger, hosts: BTreeMap<String, Logger>) -> Self {
198		Self::new(Arc::new(Mutex::new(LogTarget::Routed { fallback, hosts })))
199	}
200
201	pub fn for_host(&self, hostname: &str) -> Self {
202		self
203			.target
204			.lock()
205			.ok()
206			.and_then(|target| match &*target {
207				LogTarget::Routed { fallback, hosts } => {
208					Some(hosts.get(hostname).unwrap_or(fallback).clone())
209				}
210				_ => None,
211			})
212			.unwrap_or_else(|| self.clone())
213	}
214
215	fn routed_fallback(&self) -> Option<Logger> {
216		self.target.lock().ok().and_then(|target| match &*target {
217			LogTarget::Routed { fallback, .. } => Some(fallback.clone()),
218			_ => None,
219		})
220	}
221
222	pub fn dump_suppressed_tail(&self) {
223		if let Some(fallback) = self.routed_fallback() {
224			fallback.dump_suppressed_tail();
225			return;
226		}
227		let mut target = match self.target.lock() {
228			Ok(guard) => guard,
229			Err(_) => return,
230		};
231		let bulk = match &mut *target {
232			LogTarget::Terminal { bulk, .. } => bulk,
233			_ => return,
234		};
235		bulk.clear_active_counter();
236		bulk.flush_summary();
237		let tail = bulk.take_suppressed_tail();
238		if !tail.is_empty() {
239			eprintln!("[suppressed output tail on failure]:");
240			for line in tail {
241				eprintln!("{}", crate::progress::stream::format_output_line(&line));
242			}
243		}
244	}
245
246	pub fn event(&self, level: StatusLevel, message: &str) {
247		if level == StatusLevel::Debug && !crate::config::get_runtime_options().debug {
248			return;
249		}
250		let message = self.sanitize(message);
251		if let Some(fallback) = self.routed_fallback() {
252			fallback.event(level, &message);
253			return;
254		}
255		if let Ok(mut lock) = self.target.lock() {
256			lock.log_event(level, &message);
257		}
258	}
259
260	/// Record command detail. Single-host terminal runs retain it; batch runs
261	/// write it only to the host log so concurrent subprocess output cannot
262	/// overwhelm the interactive fleet summary.
263	pub fn detail(&self, message: &str) {
264		let message = self.sanitize(message);
265		if let Some(fallback) = self.routed_fallback() {
266			fallback.detail(&message);
267			return;
268		}
269		if let Ok(mut target) = self.target.lock() {
270			match &mut *target {
271				LogTarget::Terminal { file, bulk } => {
272					let normalized = normalize_multiline_message(&message);
273					if let Some(f) = file {
274						for line in normalized.lines() {
275							let _ = writeln!(f, "{}", line);
276						}
277					}
278					for line in normalized.lines() {
279						if crate::config::get_runtime_options().verbose {
280							bulk.clear_active_counter();
281							bulk.flush_summary();
282							println!("{}", crate::progress::stream::format_output_line(line));
283						} else {
284							bulk.process_line(line);
285						}
286					}
287				}
288				LogTarget::Batch { file, echo, .. } => {
289					for line in normalize_multiline_message(&message).lines() {
290						let _ = writeln!(file, "{}", line);
291						if let Some(bulk) = echo.as_mut() {
292							if crate::config::get_runtime_options().verbose {
293								bulk.clear_active_counter();
294								bulk.flush_summary();
295								println!("{}", crate::progress::stream::format_output_line(line));
296							} else {
297								bulk.process_line(line);
298							}
299						}
300					}
301				}
302				LogTarget::File(file) => {
303					for line in normalize_multiline_message(&message).lines() {
304						let _ = writeln!(file, "{}", line);
305					}
306				}
307				LogTarget::Routed { .. } | LogTarget::Silent => {}
308			}
309		}
310	}
311
312	/// Elapsed time between this host's first and most recent milestone.
313	/// Batch callers use this after apply so a fast host is not charged for
314	/// time spent waiting for another host in the fleet.
315	pub fn elapsed(&self) -> Option<Duration> {
316		self.target.lock().ok().and_then(|target| match &*target {
317			LogTarget::Batch { started_at: Some(started), last_event_at: Some(last), .. } => {
318				Some(last.saturating_duration_since(*started))
319			}
320			_ => None,
321		})
322	}
323
324	pub fn batch_outcome(&self) -> Option<BatchOutcome> {
325		self.target.lock().ok().and_then(|target| match &*target {
326			LogTarget::Batch { outcome, .. } => *outcome,
327			_ => None,
328		})
329	}
330
331	pub fn output_line(&self, line: &str, stderr: bool) {
332		let line = self.sanitize(line);
333		if let Some(fallback) = self.routed_fallback() {
334			fallback.output_line(&line, stderr);
335			return;
336		}
337		if let Ok(mut target) = self.target.lock() {
338			match &mut *target {
339				LogTarget::Terminal { file, bulk } => {
340					if let Some(f) = file {
341						let _ = writeln!(f, "{}", line);
342					}
343					if stderr || crate::config::get_runtime_options().verbose {
344						bulk.clear_active_counter();
345						bulk.flush_summary();
346						let formatted = crate::progress::stream::format_output_line(&line);
347						if stderr {
348							eprintln!("{}", formatted);
349						} else {
350							println!("{}", formatted);
351						}
352					} else {
353						bulk.process_line(&line);
354					}
355				}
356				LogTarget::Batch { file, echo, .. } => {
357					let _ = writeln!(file, "{}", line);
358					if let Some(bulk) = echo.as_mut() {
359						if stderr || crate::config::get_runtime_options().verbose {
360							bulk.clear_active_counter();
361							bulk.flush_summary();
362							let rendered = crate::progress::stream::format_output_line(&line);
363							if stderr {
364								eprintln!("{rendered}");
365							} else {
366								println!("{rendered}");
367							}
368						} else {
369							bulk.process_line(&line);
370						}
371					}
372				}
373				LogTarget::File(file) => {
374					let _ = writeln!(file, "{}", line);
375				}
376				LogTarget::Routed { .. } | LogTarget::Silent => {}
377			}
378		}
379	}
380
381	pub fn info(&self, message: &str) {
382		self.event(StatusLevel::Info, message);
383	}
384
385	pub fn success(&self, message: &str) {
386		self.event(StatusLevel::Success, message);
387	}
388
389	pub fn warn(&self, message: &str) {
390		self.event(StatusLevel::Warning, message);
391	}
392
393	pub fn error(&self, message: &str) {
394		self.event(StatusLevel::Error, message);
395	}
396
397	pub fn failure(&self, message: &str) {
398		self.event(StatusLevel::Failure, message);
399	}
400
401	pub fn debug(&self, message: &str) {
402		self.event(StatusLevel::Debug, message);
403	}
404}
405
406#[macro_export]
407macro_rules! info {
408    ($logger:expr, $fmt:literal) => { $logger.info(&format!($fmt)) };
409    ($logger:expr, $msg:expr) => { $logger.info($msg) };
410    ($logger:expr, $fmt:literal, $($arg:tt)*) => { $logger.info(&format!($fmt, $($arg)*)) };
411}
412
413#[macro_export]
414macro_rules! success {
415    ($logger:expr, $fmt:literal) => { $logger.success(&format!($fmt)) };
416    ($logger:expr, $msg:expr) => { $logger.success($msg) };
417    ($logger:expr, $fmt:literal, $($arg:tt)*) => { $logger.success(&format!($fmt, $($arg)*)) };
418}
419
420#[macro_export]
421macro_rules! warn {
422    ($logger:expr, $fmt:literal) => { $logger.warn(&format!($fmt)) };
423    ($logger:expr, $msg:expr) => { $logger.warn($msg) };
424    ($logger:expr, $fmt:literal, $($arg:tt)*) => { $logger.warn(&format!($fmt, $($arg)*)) };
425}
426
427#[macro_export]
428macro_rules! error {
429    ($logger:expr, $fmt:literal) => { $logger.error(&format!($fmt)) };
430    ($logger:expr, $msg:expr) => { $logger.error($msg) };
431    ($logger:expr, $fmt:literal, $($arg:tt)*) => { $logger.error(&format!($fmt, $($arg)*)) };
432}
433
434#[macro_export]
435macro_rules! failure {
436    ($logger:expr, $fmt:literal) => { $logger.failure(&format!($fmt)) };
437    ($logger:expr, $msg:expr) => { $logger.failure($msg) };
438    ($logger:expr, $fmt:literal, $($arg:tt)*) => { $logger.failure(&format!($fmt, $($arg)*)) };
439}
440
441#[macro_export]
442macro_rules! debug {
443    ($logger:expr, $fmt:literal) => { $logger.debug(&format!($fmt)) };
444    ($logger:expr, $msg:expr) => { $logger.debug($msg) };
445    ($logger:expr, $fmt:literal, $($arg:tt)*) => { $logger.debug(&format!($fmt, $($arg)*)) };
446}
447
448#[cfg(test)]
449#[path = "target_tests.rs"]
450mod tests;