Skip to main content

datui_lib/
logging.rs

1//! The log file, and keeping stray stderr off the screen: while the TUI runs, fd 2
2//! points at the log, Polars warnings go through `log`, and non-fatal errors (cache,
3//! history) land here.
4
5use std::cell::Cell;
6use std::collections::{HashSet, VecDeque};
7use std::fs::{File, OpenOptions};
8use std::io::Write;
9use std::path::{Path, PathBuf};
10use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
11use std::sync::{LazyLock, Mutex};
12
13pub use log::LevelFilter;
14use log::{Log, Metadata, Record};
15
16/// The log's name in the cache directory.
17pub const LOG_FILE_NAME: &str = "datui.log";
18
19/// Past this the log is moved aside to `<name>.1`, replacing the previous one.
20const MAX_BYTES: u64 = 1024 * 1024;
21
22/// The level when `DATUI_LOG` is unset.
23const DEFAULT_LEVEL: LevelFilter = LevelFilter::Warn;
24
25/// Where the log goes and how much of it.
26#[derive(Debug, Clone, PartialEq, Eq)]
27pub struct LogSettings {
28    /// `None` when logging is off; stray stderr is then discarded.
29    pub path: Option<PathBuf>,
30    pub level: LevelFilter,
31    /// `DATUI_LOG` held something that is not a level; said in the log once it opens.
32    pub unknown_level: Option<String>,
33}
34
35impl LogSettings {
36    /// Resolve the file from `[log] file` (which `--log-file` overrides) or the
37    /// cache directory, and the level from `DATUI_LOG`.
38    pub fn resolve(
39        configured: Option<&str>,
40        level: Option<&str>,
41        cache_dir: Option<&Path>,
42    ) -> Self {
43        let (level, unknown_level) = match level.map(str::trim).filter(|l| !l.is_empty()) {
44            None => (DEFAULT_LEVEL, None),
45            Some(text) => match text.parse::<LevelFilter>() {
46                Ok(level) => (level, None),
47                Err(_) => (DEFAULT_LEVEL, Some(text.to_string())),
48            },
49        };
50        let path = match configured.map(str::trim).filter(|p| !p.is_empty()) {
51            Some(path) => Some(crate::config::expand_config_path(path)),
52            None => cache_dir.map(|dir| dir.join(LOG_FILE_NAME)),
53        };
54        Self {
55            path: path.filter(|_| level != LevelFilter::Off),
56            level,
57            unknown_level,
58        }
59    }
60}
61
62/// The rotated copy of `path`: `datui.log` → `datui.log.1`.
63pub fn rotated_path(path: &Path) -> PathBuf {
64    let mut name = path.as_os_str().to_owned();
65    name.push(".1");
66    PathBuf::from(name)
67}
68
69/// A size-capped log file with one rotation.
70pub struct FileLog {
71    path: PathBuf,
72    file: File,
73    cap: u64,
74}
75
76impl FileLog {
77    /// Open (creating) `path` for appending, rotating it first when already over `cap`.
78    pub fn open(path: &Path, cap: u64) -> std::io::Result<Self> {
79        if let Some(dir) = path.parent().filter(|d| !d.as_os_str().is_empty()) {
80            std::fs::create_dir_all(dir)?;
81        }
82        let file = Self::append(path)?;
83        let mut log = Self {
84            path: path.to_path_buf(),
85            file,
86            cap,
87        };
88        log.rotate_if_full()?;
89        Ok(log)
90    }
91
92    fn append(path: &Path) -> std::io::Result<File> {
93        // Append mode, so two sessions sharing the file interleave lines rather than
94        // overwrite each other.
95        OpenOptions::new().create(true).append(true).open(path)
96    }
97
98    /// Write one line, then rotate if that filled the file. Returns whether it rotated.
99    pub fn write_line(&mut self, line: &str) -> std::io::Result<bool> {
100        self.file.write_all(line.as_bytes())?;
101        if !line.ends_with('\n') {
102            self.file.write_all(b"\n")?;
103        }
104        self.rotate_if_full()
105    }
106
107    /// Measured on the file, since another session may share it. The rename happens under
108    /// a lock and only when the file at the path is full (another session may already have
109    /// rotated, leaving this one writing `<name>.1`); the path is reopened either way.
110    fn rotate_if_full(&mut self) -> std::io::Result<bool> {
111        if self.file.metadata()?.len() <= self.cap {
112            return Ok(false);
113        }
114        let mut lock_name = self.path.as_os_str().to_owned();
115        lock_name.push(".lock");
116        let Some(_lock) =
117            crate::cache::lock_file(Path::new(&lock_name), std::time::Duration::from_millis(200))?
118        else {
119            // Busy: the next line tries again.
120            return Ok(false);
121        };
122        let full = std::fs::metadata(&self.path).is_ok_and(|m| m.len() > self.cap);
123        if full {
124            std::fs::rename(&self.path, rotated_path(&self.path))?;
125        }
126        self.file = Self::append(&self.path)?;
127        Ok(full)
128    }
129}
130
131struct State {
132    level: LevelFilter,
133    file: Option<FileLog>,
134    /// Whether this session has written its header line yet.
135    started: bool,
136    /// Values never written to the log as they are: keys, tokens, passwords.
137    secrets: Vec<String>,
138}
139
140static STATE: Mutex<State> = Mutex::new(State {
141    level: LevelFilter::Off,
142    file: None,
143    started: false,
144    secrets: Vec::new(),
145});
146
147fn state() -> std::sync::MutexGuard<'static, State> {
148    STATE.lock().unwrap_or_else(|e| e.into_inner())
149}
150
151struct Logger;
152
153static LOGGER: Logger = Logger;
154
155impl Log for Logger {
156    fn enabled(&self, metadata: &Metadata) -> bool {
157        metadata.level() <= state().level
158    }
159
160    fn log(&self, record: &Record) {
161        if self.enabled(record.metadata()) {
162            write_record(record.level().as_str(), record.target(), record.args());
163        }
164    }
165
166    fn flush(&self) {
167        if let Some(file) = state().file.as_mut() {
168            let _ = file.file.flush();
169        }
170    }
171}
172
173thread_local! {
174    /// Set while this thread writes a record. A panic in the middle of one runs the
175    /// hook, which logs; without this it would wait on the lock the thread holds.
176    static WRITING: Cell<bool> = const { Cell::new(false) };
177}
178
179/// Clears [`WRITING`] however the write ends.
180struct Writing;
181
182impl Writing {
183    fn enter() -> Option<Self> {
184        (!WRITING.with(|w| w.replace(true))).then_some(Self)
185    }
186}
187
188impl Drop for Writing {
189    fn drop(&mut self) {
190        WRITING.with(|w| w.set(false));
191    }
192}
193
194/// Write one record whatever the level, when the log is open.
195fn write_record(level: &str, target: &str, message: impl std::fmt::Display) {
196    let Some(_writing) = Writing::enter() else {
197        return;
198    };
199    // Formatted before the lock, in case formatting logs something itself.
200    let message = message.to_string().replace('\n', "\n    ");
201    let now = chrono::Local::now()
202        .format("%Y-%m-%d %H:%M:%S%.3f")
203        .to_string();
204    let mut state = state();
205    if state.file.is_none() {
206        return;
207    }
208    let mut text = String::new();
209    if !state.started {
210        state.started = true;
211        text.push_str(&format!(
212            "{now} ----- datui {} (pid {})\n",
213            env!("CARGO_PKG_VERSION"),
214            std::process::id()
215        ));
216    }
217    text.push_str(&format!(
218        "{now} {level:<5} {target}: {}",
219        redact(&message, &state.secrets)
220    ));
221    if let Some(file) = state.file.as_mut() {
222        // A log that cannot be written has nowhere to say so; stderr is the screen.
223        let _ = file.write_line(&text);
224    }
225}
226
227/// Open the log and install the logger and Polars warning hook; repeatable (the Python
228/// binding runs the TUI per `view`), latest settings winning. Returns the message for
229/// the user if the log cannot open, unprinted since stderr may be that log.
230pub fn init(settings: &LogSettings) -> Option<String> {
231    static INSTALLED: std::sync::Once = std::sync::Once::new();
232    INSTALLED.call_once(|| {
233        let _ = log::set_logger(&LOGGER);
234        polars_error::set_warning_function(polars_warning);
235    });
236    let mut note = None;
237    let file = settings
238        .path
239        .as_deref()
240        .and_then(|path| match FileLog::open(path, MAX_BYTES) {
241            Ok(file) => Some(file),
242            Err(e) => {
243                note = Some(format!("cannot write the log {}: {e}", path.display()));
244                None
245            }
246        });
247    let level = if file.is_some() {
248        settings.level
249    } else {
250        LevelFilter::Off
251    };
252    {
253        let mut state = state();
254        state.file = file;
255        state.level = level;
256        state.started = false;
257    }
258    log::set_max_level(level);
259    keep_out_of_log_from_env();
260    if let Some(text) = &settings.unknown_level {
261        log::warn!(target: "datui", "DATUI_LOG={text} is not a level; using warn");
262    }
263    note
264}
265
266/// The file the log is writing to, if it is open.
267pub fn current_path() -> Option<PathBuf> {
268    state().file.as_ref().map(|f| f.path.clone())
269}
270
271/// Never write `secret` to the log as it is.
272pub fn keep_out_of_log(secret: &str) {
273    // Short values would mask ordinary words.
274    if secret.len() < 8 {
275        return;
276    }
277    let mut state = state();
278    if !state.secrets.iter().any(|s| s == secret) {
279        state.secrets.push(secret.to_string());
280    }
281}
282
283/// Values of variables whose names say they hold a credential, the `[cloud] env_files`
284/// ones included.
285fn keep_out_of_log_from_env() {
286    for (name, value) in crate::cloud::cloud_env::vars() {
287        if holds_a_credential(&name, &value) {
288            keep_out_of_log(&value);
289        }
290    }
291}
292
293/// Whether an environment variable's value is a credential to mask. The name decides,
294/// except that a file or directory is not the secret itself: masking the path in
295/// `AWS_WEB_IDENTITY_TOKEN_FILE` would hide every mention of it.
296fn holds_a_credential(name: &str, value: &str) -> bool {
297    let name = name.to_ascii_uppercase();
298    let named = [
299        "SECRET",
300        "TOKEN",
301        "PASSWORD",
302        "PASSWD",
303        "ACCESS_KEY",
304        "ACCOUNT_KEY",
305        "CONNECTION_STRING",
306    ]
307    .iter()
308    .any(|word| name.contains(word))
309        || name.ends_with("_KEY");
310    let names_a_place = [
311        "_FILE",
312        "_PATH",
313        "_DIR",
314        "_DIRECTORY",
315        "_HOME",
316        "_URL",
317        "_URI",
318    ]
319    .iter()
320    .any(|suffix| name.ends_with(suffix));
321    let path = Path::new(value);
322    named && !names_a_place && !(path.is_absolute() && path.exists())
323}
324
325/// Log a failure not worth stopping for (cache, history) instead of dropping it.
326pub trait LogFailure {
327    fn or_log(self, what: &str);
328}
329
330impl<T, E: std::fmt::Display> LogFailure for Result<T, E> {
331    fn or_log(self, what: &str) {
332        if let Err(e) = self {
333            log::warn!(target: "datui", "{what}: {e:#}");
334        }
335    }
336}
337
338/// Mask credentials in a log line: known secret values, the user and password in a
339/// URL, signed-URL and SAS query parameters, and authorization headers.
340pub fn redact(text: &str, secrets: &[String]) -> String {
341    static PATTERNS: LazyLock<Vec<regex::Regex>> = LazyLock::new(|| {
342        [
343            // scheme://user:password@host
344            r"(?i)(\b[a-z][a-z0-9+.-]*://)[^/\s:@]+:[^/\s@]+@",
345            // ?X-Amz-Signature=..., &sig=..., &token=..., or a SAS token on its own
346            r"(?i)((?:^|[?&;])(?:x-amz-signature|x-amz-credential|x-amz-security-token|x-goog-signature|x-goog-credential|signature|sig|token|access_token|api_key|apikey|key|password|secret)=)[^&\s]+",
347            // Authorization: Bearer ..., "authorization": "..."
348            r#"(?i)(authorization"?\s*[:=]\s*"?)[^"\r\n]+"#,
349            // Case-sensitive Basic, or "basic statistics" would lose its noun.
350            r"(\b(?i:bearer)\s+|\bBasic\s+)[A-Za-z0-9._~+/=-]{8,}",
351            // secret_access_key = ..., "session_token": "...", AccountKey=...
352            r#"(?i)((?:secret[_-]?access[_-]?key|session[_-]?token|account[_-]?key|client[_-]?secret|sas[_-]?token|password)"?\s*[:=]\s*"?)[^\s",;}]+"#,
353        ]
354        .iter()
355        .filter_map(|p| regex::Regex::new(p).ok())
356        .collect()
357    });
358    let mut out = text.to_string();
359    for secret in secrets {
360        if out.contains(secret.as_str()) {
361            out = out.replace(secret.as_str(), "***");
362        }
363    }
364    for pattern in PATTERNS.iter() {
365        out = pattern.replace_all(&out, "${1}***").into_owned();
366    }
367    out
368}
369
370/// Polars warnings already seen this session, and user warnings not yet flashed.
371struct PolarsWarnings {
372    seen: HashSet<String>,
373    unshown: VecDeque<String>,
374}
375
376static POLARS: Mutex<Option<PolarsWarnings>> = Mutex::new(None);
377
378/// A warning repeats for every chunk and every collect; past this many distinct
379/// ones, the rest are dropped rather than let them fill the log.
380const MAX_POLARS_WARNINGS: usize = 256;
381
382/// Where `polars_warn!` goes: each distinct warning logged once per session; user
383/// warnings (not deprecations) are also queued for the footer.
384fn polars_warning(message: &str, kind: polars_error::PolarsWarning) {
385    use polars_error::PolarsWarning as W;
386    let text = message.split_whitespace().collect::<Vec<_>>().join(" ");
387    {
388        let mut polars = POLARS.lock().unwrap_or_else(|e| e.into_inner());
389        let polars = polars.get_or_insert_with(|| PolarsWarnings {
390            seen: HashSet::new(),
391            unshown: VecDeque::new(),
392        });
393        if polars.seen.len() >= MAX_POLARS_WARNINGS || !polars.seen.insert(text.clone()) {
394            return;
395        }
396        if matches!(kind, W::UserWarning | W::CategoricalRemappingWarning) {
397            polars.unshown.push_back(text.clone());
398            tell_the_loop();
399        }
400    }
401    log::warn!(target: "polars", "{kind:?}: {text}");
402}
403
404/// What wakes the run loop when there is news it has to come and look for: a warning
405/// queued for the footer, a background panic nothing reported. The loop only
406/// wakes for events and deadlines, so without this either would wait for a key.
407static NEWS: Mutex<Option<Box<dyn Fn() + Send + Sync>>> = Mutex::new(None);
408
409fn tell_the_loop() {
410    if let Some(wake) = NEWS.lock().unwrap_or_else(|e| e.into_inner()).as_ref() {
411        wake();
412    }
413}
414
415/// The next Polars user warning not yet shown, for the footer.
416pub fn next_polars_warning() -> Option<String> {
417    POLARS
418        .lock()
419        .unwrap_or_else(|e| e.into_inner())
420        .as_mut()
421        .and_then(|p| p.unshown.pop_front())
422}
423
424/// Forget which Polars warnings were seen, so a new session logs them again.
425fn reset_polars_warnings() {
426    *POLARS.lock().unwrap_or_else(|e| e.into_inner()) = None;
427}
428
429/// Set while a TUI session owns the terminal.
430static TUI_ACTIVE: AtomicBool = AtomicBool::new(false);
431
432/// The last panic on a background thread, to print if it takes the TUI thread down.
433static BACKGROUND_PANIC: Mutex<Option<String>> = Mutex::new(None);
434
435/// Background panics nothing has told the user about yet: those on threads that do
436/// not report their own (see [`catch_panic`]).
437static UNREPORTED_PANICS: AtomicUsize = AtomicUsize::new(0);
438
439thread_local! {
440    /// Set while [`catch_panic`] runs work on this thread, which then reports its own.
441    static REPORTS_ITS_PANICS: Cell<bool> = const { Cell::new(false) };
442}
443
444/// Run `work`, turning a panic into a message for the user. The hook has already
445/// logged the panic with its backtrace; the message says where. A panic on a Polars
446/// thread that this work was waiting on is resumed here, so it is reported here too.
447pub fn catch_panic<T>(work: impl FnOnce() -> T) -> Result<T, String> {
448    let reported = REPORTS_ITS_PANICS.with(|r| r.replace(true));
449    let unreported = UNREPORTED_PANICS.load(Ordering::SeqCst);
450    let caught = std::panic::catch_unwind(std::panic::AssertUnwindSafe(work));
451    REPORTS_ITS_PANICS.with(|r| r.set(reported));
452    caught.map_err(|payload| {
453        // Those counted meanwhile were the Polars threads this work waited on.
454        UNREPORTED_PANICS.fetch_min(unreported, Ordering::SeqCst);
455        let what = payload
456            .downcast_ref::<&str>()
457            .map(|s| s.to_string())
458            .or_else(|| payload.downcast_ref::<String>().cloned())
459            .unwrap_or_else(|| "panic".to_string());
460        match current_path() {
461            Some(path) => format!("Internal error: {what}\n\nDetails: {}", path.display()),
462            None => format!("Internal error: {what}"),
463        }
464    })
465}
466
467/// What to flash when a background thread panicked with nothing to say so: a raw
468/// thread's result simply never arrives. `None` when none did since the last call.
469pub fn take_unreported_panic() -> Option<String> {
470    if UNREPORTED_PANICS.swap(0, Ordering::SeqCst) == 0 {
471        return None;
472    }
473    Some(match current_path() {
474        Some(path) => format!(
475            "A background task failed; see {}",
476            path.file_name().unwrap_or_default().to_string_lossy()
477        ),
478        None => "A background task failed".to_string(),
479    })
480}
481
482/// The span during which the TUI owns the terminal: stderr goes to the log (Unix),
483/// and a panic on a background thread is logged instead of printed over the screen.
484/// Dropping it, on every exit path, hands stderr back.
485pub struct TuiSession {
486    restore_terminal: fn(),
487}
488
489impl TuiSession {
490    /// Begin right after the terminal is taken. `restore_terminal` hands the screen
491    /// back if the TUI thread unwinds from a panic the hook never saw.
492    pub fn begin(restore_terminal: fn()) -> Self {
493        reset_polars_warnings();
494        *BACKGROUND_PANIC.lock().unwrap_or_else(|e| e.into_inner()) = None;
495        UNREPORTED_PANICS.store(0, Ordering::SeqCst);
496        #[cfg(unix)]
497        {
498            // Through the log's pipe even with no log open yet: the settings that open
499            // it are read after the session begins. Lines with no log are dropped.
500            stderr::redirect(true);
501        }
502        TUI_ACTIVE.store(true, Ordering::SeqCst);
503        install_panic_hook();
504        Self { restore_terminal }
505    }
506
507    /// Call `wake` whenever a background panic or a Polars warning is waiting for
508    /// [`take_unreported_panic`] or [`next_polars_warning`], for the session's length.
509    pub fn wake_with(&self, wake: impl Fn() + Send + Sync + 'static) {
510        *NEWS.lock().unwrap_or_else(|e| e.into_inner()) = Some(Box::new(wake));
511    }
512}
513
514impl Drop for TuiSession {
515    fn drop(&mut self) {
516        NEWS.lock().unwrap_or_else(|e| e.into_inner()).take();
517        // Already inactive when the hook saw this thread panic: it restored stderr, and
518        // the hooks below it the terminal. Restoring again would pop the shell's
519        // keyboard flags rather than ours.
520        let hook_handled_it = !TUI_ACTIVE.swap(false, Ordering::SeqCst);
521        #[cfg(unix)]
522        stderr::restore();
523        // A panic resumed from another thread (a Polars worker's, say) unwinds this one
524        // without running the hook, so the terminal and the message are ours to handle.
525        if std::thread::panicking() && !hook_handled_it {
526            (self.restore_terminal)();
527            if let Some(message) = BACKGROUND_PANIC
528                .lock()
529                .unwrap_or_else(|e| e.into_inner())
530                .take()
531            {
532                eprintln!("{message}");
533            }
534        }
535    }
536}
537
538/// Wraps whatever hook is installed (color-eyre's, under ratatui's). Installed per
539/// session, above the hook `ratatui::try_init` adds each time.
540fn install_panic_hook() {
541    let tui_thread = std::thread::current().id();
542    let previous = std::panic::take_hook();
543    std::panic::set_hook(Box::new(move |info| {
544        if !TUI_ACTIVE.load(Ordering::SeqCst) {
545            previous(info);
546            return;
547        }
548        if std::thread::current().id() != tui_thread {
549            // A worker's panic is caught (tokio's blocking pool, or a catch_unwind)
550            // and the app carries on, so it must not tear down the screen.
551            let message = format!(
552                "a background thread panicked at {}: {}\n{}",
553                info.location()
554                    .map(|l| l.to_string())
555                    .unwrap_or_else(|| "?".into()),
556                panic_payload(info),
557                std::backtrace::Backtrace::force_capture()
558            );
559            log::error!(target: "datui::panic", "{message}");
560            *BACKGROUND_PANIC.lock().unwrap_or_else(|e| e.into_inner()) = Some(message);
561            if !REPORTS_ITS_PANICS.with(Cell::get) {
562                UNREPORTED_PANICS.fetch_add(1, Ordering::SeqCst);
563                tell_the_loop();
564            }
565            return;
566        }
567        // The TUI thread: hand stderr back so the report reaches the terminal. Inactive
568        // from here, so an earlier session's hook further down the chain (the Python
569        // binding, run from another thread) passes the report on instead of keeping it.
570        TUI_ACTIVE.store(false, Ordering::SeqCst);
571        // The hooks below hand back the screen but not the mouse, whose reporting
572        // outlives the alternate screen: the shell would read every click as text.
573        // Unlike popping the keyboard flags, this is safe to repeat.
574        let _ = crossterm::execute!(std::io::stdout(), crossterm::event::DisableMouseCapture);
575        BACKGROUND_PANIC
576            .lock()
577            .unwrap_or_else(|e| e.into_inner())
578            .take();
579        #[cfg(unix)]
580        stderr::restore();
581        previous(info);
582    }));
583}
584
585fn panic_payload(info: &std::panic::PanicHookInfo<'_>) -> String {
586    let payload = info.payload();
587    payload
588        .downcast_ref::<&str>()
589        .map(|s| s.to_string())
590        .or_else(|| payload.downcast_ref::<String>().cloned())
591        .unwrap_or_else(|| "panic".to_string())
592}
593
594/// One line someone wrote to stderr while the TUI was up.
595#[cfg(unix)]
596fn write_stray_line(line: &[u8]) {
597    let text = String::from_utf8_lossy(line);
598    let text = text.trim_end_matches(['\n', '\r']);
599    if !text.trim().is_empty() {
600        write_record("WARN", "stderr", text);
601    }
602}
603
604/// Pointing fd 2 away and back (not on Windows, where the Polars hook suffices). With the
605/// log open, fd 2 is a pipe a thread copies into the log line by line, counted against
606/// the cap.
607#[cfg(unix)]
608mod stderr {
609    use std::io::BufRead;
610    use std::os::fd::{AsRawFd, FromRawFd, OwnedFd, RawFd};
611    use std::sync::Mutex;
612    use std::sync::mpsc::{Receiver, channel};
613    use std::time::Duration;
614
615    struct Redirect {
616        /// A duplicate of the real stderr.
617        terminal: OwnedFd,
618        /// Answers once the reader has read the pipe to its end; `None` for /dev/null.
619        drained: Option<Receiver<()>>,
620    }
621
622    static REDIRECT: Mutex<Option<Redirect>> = Mutex::new(None);
623
624    fn current() -> std::sync::MutexGuard<'static, Option<Redirect>> {
625        REDIRECT.lock().unwrap_or_else(|e| e.into_inner())
626    }
627
628    /// Point fd 2 at the log, or at `/dev/null` with logging off.
629    pub(super) fn redirect(logging: bool) {
630        let mut current = current();
631        if current.is_some() {
632            return;
633        }
634        // SAFETY: fcntl(F_DUPFD_CLOEXEC) reads no memory; on failure it returns -1.
635        // Close-on-exec, so a child process never inherits the terminal through it.
636        let copy = unsafe { libc::fcntl(libc::STDERR_FILENO, libc::F_DUPFD_CLOEXEC, 0) };
637        if copy < 0 {
638            return;
639        }
640        // SAFETY: `copy` is a new descriptor that nothing else owns.
641        let terminal = unsafe { OwnedFd::from_raw_fd(copy) };
642        let target = logging.then(into_the_log).flatten().or_else(|| {
643            std::fs::OpenOptions::new()
644                .write(true)
645                .open("/dev/null")
646                .ok()
647                .map(|null| (OwnedFd::from(null), None))
648        });
649        // `target` closes at the end of this function, which leaves fd 2 as the pipe's
650        // only writer: once fd 2 points back at the terminal, the reader sees the end.
651        if let Some((target, drained)) = target
652            && point(target.as_raw_fd())
653        {
654            *current = Some(Redirect { terminal, drained });
655        }
656    }
657
658    /// A pipe whose other end a thread copies into the log, line by line.
659    fn into_the_log() -> Option<(OwnedFd, Option<Receiver<()>>)> {
660        let (reader, writer) = std::io::pipe().ok()?;
661        let (done, drained) = channel();
662        std::thread::Builder::new()
663            .name("datui-stderr".into())
664            .spawn(move || {
665                let mut reader = std::io::BufReader::new(reader);
666                let mut line = Vec::new();
667                loop {
668                    line.clear();
669                    match reader.read_until(b'\n', &mut line) {
670                        Ok(0) => break,
671                        // Caught, because a reader that died would close the pipe and
672                        // turn every later `eprintln!` into a panic of its own.
673                        Ok(_) => {
674                            let _ = std::panic::catch_unwind(|| super::write_stray_line(&line));
675                        }
676                        Err(e) if e.kind() == std::io::ErrorKind::Interrupted => {}
677                        Err(_) => break,
678                    }
679                }
680                let _ = done.send(());
681            })
682            .ok()?;
683        Some((writer.into(), Some(drained)))
684    }
685
686    /// Hand the real stderr back. A no-op when it is not redirected.
687    pub(super) fn restore() {
688        let Some(redirect) = current().take() else {
689            return;
690        };
691        point(redirect.terminal.as_raw_fd());
692        // Let the reader finish what was written just before, so a panic report or a
693        // last warning is in the log. Bounded: a child process that inherited fd 2
694        // holds the pipe open for as long as it runs.
695        if let Some(drained) = redirect.drained {
696            let _ = drained.recv_timeout(Duration::from_millis(500));
697        }
698    }
699
700    fn point(fd: RawFd) -> bool {
701        // SAFETY: dup2 reads no memory; `fd` is open, as the caller holds it. fd 2 is
702        // replaced atomically, so a concurrent write lands on one file or the other.
703        unsafe { libc::dup2(fd, libc::STDERR_FILENO) >= 0 }
704    }
705}
706
707#[cfg(test)]
708mod tests {
709    use super::*;
710
711    #[test]
712    fn the_log_goes_to_the_cache_unless_configured() {
713        let cache = Path::new("/cache/datui");
714        let default = LogSettings::resolve(None, None, Some(cache));
715        assert_eq!(default.path, Some(cache.join("datui.log")));
716        assert_eq!(default.level, LevelFilter::Warn);
717
718        let chosen = LogSettings::resolve(Some("/elsewhere/x.log"), Some("debug"), Some(cache));
719        assert_eq!(chosen.path, Some(PathBuf::from("/elsewhere/x.log")));
720        assert_eq!(chosen.level, LevelFilter::Debug);
721
722        let blank = LogSettings::resolve(Some("  "), Some(" "), Some(cache));
723        assert_eq!(blank.path, Some(cache.join("datui.log")));
724        assert_eq!(blank.level, LevelFilter::Warn);
725    }
726
727    #[test]
728    fn off_means_no_file_and_a_typo_means_the_default() {
729        let off = LogSettings::resolve(None, Some("OFF"), Some(Path::new("/c")));
730        assert_eq!(off.path, None);
731
732        let typo = LogSettings::resolve(None, Some("loud"), Some(Path::new("/c")));
733        assert_eq!(typo.level, LevelFilter::Warn);
734        assert_eq!(typo.unknown_level.as_deref(), Some("loud"));
735        assert!(typo.path.is_some());
736    }
737
738    #[test]
739    fn a_full_log_moves_aside_once() {
740        let dir = tempfile::tempdir().unwrap();
741        let path = dir.path().join("sub").join("datui.log");
742        let mut log = FileLog::open(&path, 100).unwrap();
743        let line = "x".repeat(60);
744        assert!(!log.write_line(&line).unwrap());
745        assert!(log.write_line(&line).unwrap(), "past the cap it rotates");
746        assert!(rotated_path(&path).exists());
747        assert_eq!(std::fs::metadata(&path).unwrap().len(), 0);
748        // A second rotation replaces the first: one old file, never more.
749        log.write_line(&"y".repeat(120)).unwrap();
750        let old = std::fs::read_to_string(rotated_path(&path)).unwrap();
751        assert!(old.starts_with('y'), "{old:?}");
752        assert!(!dir.path().join("sub").join("datui.log.1.1").exists());
753    }
754
755    /// Two sessions share the log. One moves it aside; the other, still writing into
756    /// the moved file, must not then move the first one's fresh log over it.
757    #[test]
758    fn two_sessions_never_rotate_each_others_lines_away() {
759        let dir = tempfile::tempdir().unwrap();
760        let path = dir.path().join("datui.log");
761        let mut a = FileLog::open(&path, 100).unwrap();
762        let mut b = FileLog::open(&path, 100).unwrap();
763        let mut written = Vec::new();
764        // Under two caps' worth in all: one rotation, so nothing may be lost.
765        for n in 0..2 {
766            for (who, log) in [("a", &mut a), ("b", &mut b)] {
767                let line = format!("{who}{n} {}", "x".repeat(40));
768                log.write_line(&line).unwrap();
769                written.push(line);
770            }
771        }
772        let all = std::fs::read_to_string(rotated_path(&path)).unwrap_or_default()
773            + &std::fs::read_to_string(&path).unwrap();
774        for line in &written {
775            assert!(all.contains(line.as_str()), "lost {line:?} from {all:?}");
776        }
777    }
778
779    #[test]
780    fn credentials_are_masked() {
781        let secrets = vec!["wJalrXUtnFEMI/K7MDENG".to_string()];
782        let cases = [
783            ("key wJalrXUtnFEMI/K7MDENG refused", "wJalrXUtnFEMI"),
784            ("GET https://alice:hunter22@host/x failed", "hunter22"),
785            (
786                "403 for https://b.s3.amazonaws.com/k?X-Amz-Credential=AKIA123%2F&X-Amz-Signature=abcdef12",
787                "abcdef12",
788            ),
789            (
790                "https://acct.blob.core.windows.net/c?sv=2022&sig=Zm9vYmFy",
791                "Zm9vYmFy",
792            ),
793            ("Authorization: Bearer eyJhbGciOiJIUzI1NiJ9.x.y", "eyJhbGci"),
794            (
795                "DefaultEndpointsProtocol=https;AccountKey=c2VjcmV0a2V5;",
796                "c2VjcmV0a2V5",
797            ),
798        ];
799        for (line, secret) in cases {
800            let masked = redact(line, &secrets);
801            assert!(!masked.contains(secret), "{line} -> {masked}");
802            assert!(masked.contains("***"), "{masked}");
803        }
804        assert_eq!(
805            redact("s3://bucket/key.parquet: not found", &secrets),
806            "s3://bucket/key.parquet: not found"
807        );
808        for plain in [
809            "basic statistics failed for column_with_long_name",
810            "abfss://container@account.dfs.core.windows.net/data.parquet",
811        ] {
812            assert_eq!(redact(plain, &secrets), plain);
813        }
814        assert!(!redact("Basic YWxpY2U6aHVudGVyMg==", &secrets).contains("YWxpY2U6"));
815        assert!(!redact("sig=Zm9vYmFyYmF6&sv=2022", &secrets).contains("Zm9vYmFy"));
816    }
817
818    #[test]
819    fn a_variable_is_masked_by_its_name_but_not_when_it_names_a_place() {
820        let key = "wJalrXUtnFEMI/K7MDENG/bPxRfiCYEXAMPLEKEY";
821        for name in [
822            "AWS_SECRET_ACCESS_KEY",
823            "AWS_SESSION_TOKEN",
824            "AZURE_STORAGE_ACCOUNT_KEY",
825            "AZURE_STORAGE_CONNECTION_STRING",
826            "MINIO_ROOT_PASSWORD",
827            "hf_token",
828            "OPENAI_API_KEY",
829        ] {
830            assert!(holds_a_credential(name, key), "{name}");
831        }
832        for name in [
833            "AWS_WEB_IDENTITY_TOKEN_FILE",
834            "GITHUB_TOKEN_PATH",
835            "PASSWORD_STORE_DIR",
836            "VAULT_TOKEN_URL",
837        ] {
838            assert!(!holds_a_credential(name, key), "{name}");
839        }
840        assert!(!holds_a_credential("AWS_REGION", key));
841        let dir = tempfile::tempdir().unwrap();
842        let place = dir.path().to_string_lossy();
843        assert!(
844            !holds_a_credential("SOME_SECRET", &place),
845            "a path that exists is not the secret"
846        );
847    }
848
849    #[test]
850    fn a_caught_panic_becomes_a_message() {
851        assert_eq!(catch_panic(|| 7), Ok(7));
852        let message = catch_panic::<()>(|| panic!("worker died")).unwrap_err();
853        assert!(
854            message.starts_with("Internal error: worker died"),
855            "{message}"
856        );
857        let formatted = catch_panic::<()>(|| panic!("row {} of {}", 3, 9)).unwrap_err();
858        assert!(formatted.contains("row 3 of 9"), "{formatted}");
859    }
860
861    /// The only test in this binary that sets the global logger and the Polars hook, so
862    /// nothing else races it for the file.
863    #[test]
864    fn a_polars_warning_lands_in_the_log_once_and_a_user_warning_is_queued() {
865        let dir = tempfile::tempdir().unwrap();
866        let path = dir.path().join("datui.log");
867        init(&LogSettings {
868            path: Some(path.clone()),
869            level: LevelFilter::Warn,
870            unknown_level: None,
871        });
872        assert_eq!(current_path(), Some(path.clone()));
873
874        for _ in 0..3 {
875            polars_error::polars_warn!(
876                Deprecation,
877                "casting in test {} is deprecated.\nUse something else.",
878                std::process::id()
879            );
880        }
881        polars_error::polars_warn!(UserWarning, "remapped categories in test {}", 7);
882        log::info!("below the level");
883        log::logger().flush();
884
885        let text = std::fs::read_to_string(&path).unwrap();
886        let deprecation = format!(
887            "Deprecation: casting in test {} is deprecated. Use something else.",
888            std::process::id()
889        );
890        assert_eq!(text.matches(&deprecation).count(), 1, "{text}");
891        assert!(
892            text.contains("UserWarning: remapped categories in test 7"),
893            "{text}"
894        );
895        assert!(!text.contains("below the level"), "{text}");
896
897        // The user warning reaches the footer, once, when the bar is free.
898        let (tx, _rx) = std::sync::mpsc::channel();
899        let mut app = crate::App::new(tx, crate::tests::test_runtime());
900        app.busy = true;
901        assert!(!app.flash_polars_warning(), "a busy message outranks it");
902        app.busy = false;
903        app.error_modal.active = true;
904        assert!(!app.flash_polars_warning(), "a modal would hide it");
905        app.error_modal.active = false;
906        let mut flashed = Vec::new();
907        while app.flash_polars_warning() {
908            flashed.extend(app.flash_message().map(str::to_string));
909            app.flash = None;
910        }
911        assert!(
912            flashed.contains(&"Polars: remapped categories in test 7".to_string()),
913            "{flashed:?}"
914        );
915        assert!(!flashed.iter().any(|w| w.contains("casting in test")));
916        polars_error::polars_warn!(UserWarning, "remapped categories in test {}", 7);
917        assert!(
918            std::iter::from_fn(next_polars_warning).all(|w| !w.contains("in test 7")),
919            "a user warning is flashed once"
920        );
921    }
922}