Skip to main content

datui_lib/
stdin.rs

1//! Data piped in: `cmd | datui` and `datui -`.
2//!
3//! Standard input is read once, to a temporary file, as a phase of the open
4//! ([`crate::loading`]); the scan of that file then stays lazy, as for any file. The
5//! format is read off the first bytes, since a pipe has no extension to go by, unless
6//! `--format` or `--compression` says. Keys come from the terminal meanwhile: Crossterm
7//! reads `/dev/tty` on Unix when standard input is not one, and `CONIN$` on Windows.
8
9use std::io::Read;
10use std::path::{Path, PathBuf};
11use std::sync::atomic::AtomicU64;
12
13use crate::download::{Opened, StreamError, TempDownload};
14use crate::unfinished::Writer;
15use crate::{CompressionFormat, FileFormat, OpenOptions};
16
17/// The path that names standard input on the command line.
18pub const PATH: &str = "-";
19
20/// What data read from standard input is called on screen and in messages.
21pub const NAME: &str = "stdin";
22
23/// Whether `path` names standard input.
24pub fn is_stdin(path: &Path) -> bool {
25    path.as_os_str() == PATH
26}
27
28/// The name `path` goes by: `stdin` for standard input, itself otherwise.
29pub fn named(path: &Path) -> PathBuf {
30    if is_stdin(path) {
31        PathBuf::from(NAME)
32    } else {
33        path.to_path_buf()
34    }
35}
36
37/// Whether standard input carries data: a pipe or a file. A terminal there is the
38/// user, and a device such as `/dev/null` or Windows' `NUL`, where a launcher points
39/// it, holds nothing.
40pub fn piped() -> bool {
41    use std::io::IsTerminal;
42    let stdin = std::io::stdin();
43    carries_data(stdin.is_terminal(), is_device(&stdin))
44}
45
46/// What [`piped`] decides from what standard input is: neither a terminal nor a
47/// device.
48pub fn carries_data(terminal: bool, device: bool) -> bool {
49    !terminal && !device
50}
51
52/// Whether standard input is a character device. One that cannot be asked is taken
53/// for one, so nothing is read from it.
54#[cfg(unix)]
55fn is_device(stdin: &std::io::Stdin) -> bool {
56    use std::os::fd::AsFd;
57    use std::os::unix::fs::FileTypeExt;
58    stdin
59        .as_fd()
60        .try_clone_to_owned()
61        .and_then(|fd| std::fs::File::from(fd).metadata())
62        .ok()
63        .is_none_or(|meta| meta.file_type().is_char_device())
64}
65
66/// Whether standard input is anything but a file or a pipe: `NUL`, a console, which
67/// `is_terminal` has already answered for, or no handle at all, as a process started
68/// without one has; none of these is read.
69#[cfg(windows)]
70fn is_device(stdin: &std::io::Stdin) -> bool {
71    use std::os::windows::io::AsRawHandle;
72    use windows_sys::Win32::Storage::FileSystem::{FILE_TYPE_DISK, FILE_TYPE_PIPE, GetFileType};
73    // SAFETY: the handle is standard input's, open for the life of the process; a null
74    // or invalid one makes `GetFileType` answer `FILE_TYPE_UNKNOWN`.
75    let kind = unsafe { GetFileType(stdin.as_raw_handle()) };
76    !matches!(kind, FILE_TYPE_DISK | FILE_TYPE_PIPE)
77}
78
79#[cfg(not(any(unix, windows)))]
80fn is_device(_stdin: &std::io::Stdin) -> bool {
81    false
82}
83
84/// `path` as a host other than the command line means it: `-` is a file of that name,
85/// since only the command line reads standard input.
86pub fn as_file(path: PathBuf) -> PathBuf {
87    if is_stdin(&path) {
88        Path::new(".").join(path)
89    } else {
90        path
91    }
92}
93
94/// The paths to open: those named, or standard input when none are and something is
95/// piped in.
96pub fn paths_or_stdin(paths: Vec<PathBuf>, piped: bool) -> Vec<PathBuf> {
97    if paths.is_empty() && piped {
98        vec![PathBuf::from(PATH)]
99    } else {
100        paths
101    }
102}
103
104/// Why standard input cannot be read as asked, before anything is read: nothing is
105/// piped in, or it is named with other paths.
106pub fn refuse(paths: &[PathBuf], piped: bool) -> Option<&'static str> {
107    if !paths.iter().any(|path| is_stdin(path)) {
108        return None;
109    }
110    if paths.len() > 1 {
111        return Some("Standard input (-) is read on its own. Name it without other paths.");
112    }
113    (!piped)
114        .then_some("Nothing is piped to standard input. Pipe data in, as in: cat data.csv | datui")
115}
116
117/// The format and compression the first bytes of a file say it is: compression by its
118/// magic numbers, then whatever a format's signature says ([`crate::readers::sniff`]),
119/// and text no format claims as JSON, CSV or TSV on evidence, lines otherwise
120/// ([`crate::lines::guess`]). `head` is all there is when it is shorter than [`HEAD`].
121pub fn sniff(head: &[u8]) -> (FileFormat, Option<CompressionFormat>) {
122    let (format, compression, _) = sniffed(head);
123    (format, compression)
124}
125
126/// [`sniff`] as `options` ask: a guess of lines is CSV with `--delimiter`. Also says
127/// whether the format was guessed rather than said by a signature.
128pub(crate) fn sniff_for(
129    head: &[u8],
130    options: &OpenOptions,
131) -> (FileFormat, Option<CompressionFormat>, bool) {
132    let (format, compression, guessed) = sniffed(head);
133    match guessed {
134        true => (crate::lines::as_asked(format, options), compression, true),
135        false => (format, compression, false),
136    }
137}
138
139/// [`sniff`], and whether the format was guessed rather than said by a signature.
140fn sniffed(head: &[u8]) -> (FileFormat, Option<CompressionFormat>, bool) {
141    const COMPRESSED: [(&[u8], CompressionFormat); 4] = [
142        (b"\x1f\x8b", CompressionFormat::Gzip),
143        (b"\x28\xb5\x2f\xfd", CompressionFormat::Zstd),
144        (b"BZh", CompressionFormat::Bzip2),
145        (b"\xfd7zXZ\x00", CompressionFormat::Xz),
146    ];
147    if let Some((_, compression)) = COMPRESSED.iter().find(|(magic, _)| head.starts_with(magic)) {
148        // What is inside is looked at once it is on disk ([`inside`]).
149        return (FileFormat::TEXT, Some(*compression), false);
150    }
151    if let Some(format) = crate::readers::sniff(head, None, crate::readers::Asked::Pipe, |_| true) {
152        return (format, None, false);
153    }
154    let format = crate::lines::guess(head, head.len() < HEAD).unwrap_or(FileFormat::TEXT);
155    (format, None, true)
156}
157
158/// What compressed data piped in holds, by its first bytes once decompressed: a
159/// format read through its compression by its signature, else delimited text on
160/// evidence, else lines.
161fn inside(file: &Path, compression: CompressionFormat) -> (FileFormat, bool) {
162    let Some(head) = crate::formats::head_of(file, Some(compression), HEAD as u64) else {
163        return (FileFormat::TEXT, false);
164    };
165    if let Some(format) = crate::readers::sniff(
166        &head,
167        None,
168        crate::readers::Asked::Pipe,
169        FileFormat::reads_into,
170    ) {
171        return (format, false);
172    }
173    let format = crate::lines::guess(&head, head.len() < HEAD)
174        .filter(|f| f.decompressed_once())
175        .unwrap_or(FileFormat::TEXT);
176    (format, true)
177}
178
179/// Bytes [`sniff`] looks at.
180const HEAD: usize = crate::readers::HEAD;
181
182/// Where what comes in on standard input is copied: `--temp-dir`, else `spool` in the
183/// cache directory, on disk, rather than the system's temp directory, which is memory
184/// on many Linux machines. The system's when the cache directory cannot be made.
185pub(crate) fn spool_dir(options: &OpenOptions) -> Option<PathBuf> {
186    if let Some(dir) = &options.temp_dir {
187        return Some(dir.clone());
188    }
189    let dir = crate::cache::CacheManager::new(crate::APP_NAME)
190        .ok()?
191        .cache_dir()
192        .join("spool");
193    std::fs::create_dir_all(&dir).ok()?;
194    forget_old_spools(&dir);
195    Some(dir)
196}
197
198/// How old a spool left behind is before it goes: one a session that crashed did not
199/// remove. Long enough that no session still reading its own is near it.
200const OLD_SPOOL: std::time::Duration = std::time::Duration::from_secs(7 * 24 * 60 * 60);
201
202/// Remove the spools in `dir` older than [`OLD_SPOOL`]. Best effort: the cache is
203/// not the system temp directory, which a reboot empties.
204fn forget_old_spools(dir: &Path) {
205    let Ok(entries) = std::fs::read_dir(dir) else {
206        return;
207    };
208    for entry in entries.flatten() {
209        let old = entry
210            .metadata()
211            .and_then(|m| m.modified())
212            .ok()
213            .and_then(|at| at.elapsed().ok())
214            .is_some_and(|age| age > OLD_SPOOL);
215        // Not one a live datui still holds, however long it has been reading.
216        #[cfg(unix)]
217        let free = !crate::download::held_elsewhere(&entry.path());
218        // Elsewhere a file held open cannot be removed, and the removal fails.
219        #[cfg(not(unix))]
220        let free = true;
221        if old && free && entry.path().extension().is_some_and(|e| e == "tmp") {
222            let _ = std::fs::remove_file(entry.path());
223        }
224    }
225}
226
227/// Read what `open` answers with into a temporary file in the spool directory
228/// ([`spool_dir`]), counting the bytes into `read`, and say what it holds: `options`
229/// with the format and compression the first bytes say, where the user did not.
230/// Stops, removing the file, once `writer`'s open is stopped.
231pub(crate) fn spool<R: Read>(
232    open: impl FnOnce() -> Opened<R> + Send + 'static,
233    options: OpenOptions,
234    writer: &Writer,
235    read: &AtomicU64,
236) -> Result<(TempDownload, OpenOptions), String> {
237    let file = crate::download::spool_to_temp(spool_dir(&options).as_deref(), open, writer, read)
238        .map_err(|error| match error {
239        StreamError::Open(e) | StreamError::Read(e) => {
240            format!("Could not read standard input: {e}")
241        }
242        StreamError::Write(report) => crate::error_display::user_message_from_report(&report, None),
243        StreamError::Short { .. } | StreamError::Cut => {
244            "Reading standard input was stopped.".to_string()
245        }
246    })?;
247    let options = described(file.path(), options)?;
248    Ok((file, options))
249}
250
251/// What all of standard input, copied to `file`, holds: `options` with the format and
252/// compression its first bytes say, where the user did not.
253pub(crate) fn described(file: &Path, options: OpenOptions) -> Result<OpenOptions, String> {
254    let mut head = Vec::with_capacity(HEAD);
255    std::fs::File::open(file)
256        .and_then(|f| f.take(HEAD as u64).read_to_end(&mut head))
257        .map_err(|e| format!("Could not read standard input back: {e}"))?;
258    if head.is_empty() {
259        return Err("Nothing came in on standard input.".to_string());
260    }
261    let (mut format, compression, mut guessed) = sniffed(&head);
262    if options.format.is_none()
263        && let Some(compression) = options.compression.or(compression)
264    {
265        (format, guessed) = inside(file, compression);
266    }
267    if guessed {
268        format = crate::lines::as_asked(format, &options);
269    }
270    Ok(match (options.format, options.compression) {
271        // A delimited format named and compression not: the bytes say whether it is
272        // compressed, as a file's extension would.
273        (Some(named), None) if named.separator().is_some() => OpenOptions {
274            compression,
275            ..options
276        },
277        // Named by the user: theirs, compression and all.
278        (Some(_), _) => options,
279        // Compression named, and what it holds read through it.
280        (None, Some(_)) => OpenOptions {
281            format: Some(format),
282            format_guessed: guessed,
283            ..options
284        },
285        (None, None) => OpenOptions {
286            format: Some(format),
287            compression,
288            format_guessed: guessed,
289            ..options
290        },
291    })
292}
293
294/// Whether standard input opened with `options` may be shown as it arrives: nothing
295/// named rules it out (a format read once it is finished, or compression).
296pub(crate) fn may_read_as_it_arrives(options: &OpenOptions) -> bool {
297    options.compression.is_none()
298        && options
299            .format
300            .is_none_or(|format| format.follows() || format == FileFormat::Arrow)
301}
302
303#[cfg(test)]
304mod tests {
305    use super::*;
306    use std::io::Write;
307    use std::sync::Arc;
308    use std::sync::atomic::{AtomicBool, Ordering};
309
310    fn files_in(dir: &Path) -> usize {
311        std::fs::read_dir(dir).unwrap().count()
312    }
313
314    fn options_in(dir: &Path) -> OpenOptions {
315        OpenOptions {
316            temp_dir: Some(dir.to_path_buf()),
317            ..Default::default()
318        }
319    }
320
321    /// Each format by its first bytes; whitespace and a byte-order mark before JSON
322    /// are not the data's first character.
323    #[test]
324    fn the_first_bytes_say_the_format() {
325        let cases: [(&[u8], FileFormat, Option<CompressionFormat>); 37] = [
326            (b"PAR1\x15\x04", FileFormat::Parquet, None),
327            (b"\x93NUMPY\x01\x00", FileFormat::Numpy, None),
328            (b"\x7fELF\x02\x01\x01", FileFormat::Elf, None),
329            (b"ULog\x01\x12\x35\x01", FileFormat::Ulog, None),
330            (b"(1.000100) can0 123#DEADBEEF\n", FileFormat::Candump, None),
331            (b"\xa3\x95\x80\x80\x59FMT\0", FileFormat::Dataflash, None),
332            (b"SQLite format 3\0\x10\x00", FileFormat::Sqlite, None),
333            (b"8=FIX.4.4|9=5|35=0|10=000|\n", FileFormat::Fix, None),
334            (b"$version Verilator $end\n", FileFormat::Vcd, None),
335            (
336                b"aspirin\n\n\n  1  0  0  0  0  0  0  0  0  0999 V2000\n",
337                FileFormat::Sdf,
338                None,
339            ),
340            (b"GGUF\x03\x00\x00\x00", FileFormat::Gguf, None),
341            (b"MThd\0\0\0\x06\0\x01", FileFormat::Midi, None),
342            (
343                b"\x02\x00\x00\x00\x00\x00\x00\x00{}",
344                FileFormat::Safetensors,
345                None,
346            ),
347            (b"$GPGGA,123519,4807.038,N", FileFormat::Nmea, None),
348            // A CSV header with the shape of a sentence.
349            (b"$USD,$EUR\n1,2\n", FileFormat::Csv, None),
350            (
351                b"<?xml version=\"1.0\"?>\n<gpx version=\"1.1\">",
352                FileFormat::Gpx,
353                None,
354            ),
355            (b"RIFF\x24\x00\x00\x00WAVEfmt ", FileFormat::Audio, None),
356            (b"FORM\x00\x00\x00\x2eAIFFCOMM", FileFormat::Audio, None),
357            (b"ARROW1\x00\x00", FileFormat::Arrow, None),
358            // An Arrow IPC stream, cut short of its schema message.
359            (b"\xff\xff\xff\xff\x10\x01\x00\x00", FileFormat::Arrow, None),
360            (b"Obj\x01\x04", FileFormat::Avro, None),
361            (
362                b"\x1f\x8b\x08\x00",
363                FileFormat::Text,
364                Some(CompressionFormat::Gzip),
365            ),
366            (
367                b"\x28\xb5\x2f\xfd\x04",
368                FileFormat::Text,
369                Some(CompressionFormat::Zstd),
370            ),
371            (b"BZh91AY", FileFormat::Text, Some(CompressionFormat::Bzip2)),
372            (
373                b"\xfd7zXZ\x00\x00",
374                FileFormat::Text,
375                Some(CompressionFormat::Xz),
376            ),
377            (b"  \n[{\"a\": 1}]", FileFormat::Json, None),
378            (b"\xef\xbb\xbf{\"a\": 1}\n", FileFormat::Jsonl, None),
379            (b"a,b\n1,2\n", FileFormat::Csv, None),
380            // One field a line, a log, a header alone: lines.
381            (b"1\n2\n3\n", FileFormat::Text, None),
382            (
383                b"Oct  3 12:00:01 host sshd[1]: Accepted\n",
384                FileFormat::Text,
385                None,
386            ),
387            (b"id,name\n", FileFormat::Text, None),
388            // A pretty-printed object, as `curl` gets from an API, is not one per line.
389            (b"{\n  \"a\": 1\n}\n", FileFormat::Json, None),
390            (b"{\"a\": 1}  \r\n{\"a\": 2}", FileFormat::Jsonl, None),
391            (b"{\"a\": 1}", FileFormat::Jsonl, None),
392            (b"id\tname\n1\tx\n", FileFormat::Tsv, None),
393            (b"id\tname,first\n", FileFormat::Text, None),
394            (b"id,name\n1,a\tb\n", FileFormat::Csv, None),
395        ];
396        for (head, format, compression) in cases {
397            assert_eq!(sniff(head), (format, compression), "{head:?}");
398        }
399    }
400
401    /// A terminal is the user and a device (`/dev/null`, `NUL`) holds nothing; only
402    /// a pipe or a file is read (#567).
403    #[test]
404    fn only_a_pipe_or_a_file_carries_data() {
405        assert!(carries_data(false, false), "a pipe or a file");
406        assert!(!carries_data(true, false), "a terminal");
407        assert!(!carries_data(false, true), "/dev/null or NUL");
408        assert!(!carries_data(true, true), "a console");
409    }
410
411    /// `-` is standard input wherever it is named, alone; no paths and a pipe on
412    /// standard input select it, and a terminal there does not.
413    #[test]
414    fn stdin_is_chosen_by_a_dash_or_a_pipe() {
415        assert_eq!(paths_or_stdin(Vec::new(), true), vec![PathBuf::from("-")]);
416        assert!(paths_or_stdin(Vec::new(), false).is_empty());
417        let listed = vec![PathBuf::from("a.csv")];
418        assert_eq!(paths_or_stdin(listed.clone(), true), listed);
419
420        assert_eq!(refuse(&listed, false), None);
421        assert_eq!(refuse(&[PathBuf::from("-")], true), None);
422        assert!(refuse(&[PathBuf::from("-")], false).is_some());
423        assert!(refuse(&[PathBuf::from("-"), PathBuf::from("a.csv")], true).is_some());
424        assert_eq!(named(Path::new("-")), PathBuf::from("stdin"));
425        assert_eq!(named(Path::new("./-")), PathBuf::from("./-"));
426        assert!(!is_stdin(&as_file(PathBuf::from("-"))));
427        assert_eq!(as_file(PathBuf::from("a.csv")), PathBuf::from("a.csv"));
428    }
429
430    /// Spooled whole, counted, and its format read off the file; what the user named
431    /// wins over the bytes.
432    #[test]
433    fn a_reader_is_spooled_whole_and_its_format_read() {
434        let dir = tempfile::tempdir().unwrap();
435        let read = AtomicU64::new(0);
436        let body = b"[{\"a\": 1}]".to_vec();
437        let (file, options) = spool(
438            move || Ok((std::io::Cursor::new(body), None)),
439            options_in(dir.path()),
440            &Writer::default(),
441            &read,
442        )
443        .unwrap();
444        assert_eq!(std::fs::read(file.path()).unwrap(), b"[{\"a\": 1}]");
445        assert_eq!(read.load(Ordering::Relaxed), 10);
446        assert_eq!(options.format, Some(FileFormat::Json));
447        assert_eq!(options.compression, None);
448
449        let named = OpenOptions {
450            compression: Some(CompressionFormat::Zstd),
451            ..options_in(dir.path())
452        };
453        let (_, options) = spool(
454            || Ok((std::io::Cursor::new(b"\x1f\x8b".to_vec()), None)),
455            named,
456            &Writer::default(),
457            &AtomicU64::new(0),
458        )
459        .unwrap();
460        // Not zstd after all: nothing inside says more than lines.
461        assert_eq!(options.format, Some(FileFormat::Text));
462        assert_eq!(options.compression, Some(CompressionFormat::Zstd));
463
464        // A delimited format named alone: its compression still comes from the bytes.
465        let named = OpenOptions {
466            format: Some(FileFormat::Tsv),
467            ..options_in(dir.path())
468        };
469        let (_, options) = spool(
470            || Ok((std::io::Cursor::new(b"\x1f\x8b".to_vec()), None)),
471            named,
472            &Writer::default(),
473            &AtomicU64::new(0),
474        )
475        .unwrap();
476        assert_eq!(options.format, Some(FileFormat::Tsv));
477        assert_eq!(options.compression, Some(CompressionFormat::Gzip));
478
479        // Compressed and unnamed: what is inside says CSV on evidence, lines otherwise.
480        for (body, format) in [
481            (&b"id,name\n1,a\n2,b\n"[..], FileFormat::Csv),
482            (b"started\n\nstopped, after 2s\n", FileFormat::Text),
483        ] {
484            let mut gz = flate2::write::GzEncoder::new(Vec::new(), Default::default());
485            gz.write_all(body).unwrap();
486            let gz = gz.finish().unwrap();
487            let (_, options) = spool(
488                move || Ok((std::io::Cursor::new(gz), None)),
489                options_in(dir.path()),
490                &Writer::default(),
491                &AtomicU64::new(0),
492            )
493            .unwrap();
494            assert_eq!(options.format, Some(format));
495            assert_eq!(options.compression, Some(CompressionFormat::Gzip));
496            assert!(options.format_guessed);
497        }
498
499        let empty = spool(
500            || Ok((std::io::empty(), None)),
501            options_in(dir.path()),
502            &Writer::default(),
503            &AtomicU64::new(0),
504        );
505        assert!(empty.is_err(), "nothing piped in is said");
506    }
507
508    /// A producer that goes quiet mid-stream: what came is counted, and stopping the
509    /// open ends the read and removes the partial file, though the pipe stays open.
510    #[test]
511    fn a_stop_mid_spool_removes_the_partial_file() {
512        let dir = tempfile::tempdir().unwrap();
513        let (reader, mut pipe) = std::io::pipe().unwrap();
514        let stop = Arc::new(AtomicBool::new(false));
515        let writer = crate::unfinished::Unfinished::default().writer(stop.clone());
516        let read = Arc::new(AtomicU64::new(0));
517        let worker = {
518            let (writer, read) = (writer.clone(), read.clone());
519            let options = options_in(dir.path());
520            std::thread::spawn(move || spool(move || Ok((reader, None)), options, &writer, &read))
521        };
522        pipe.write_all(b"a,b\n1,2\n").unwrap();
523        let deadline = std::time::Instant::now() + std::time::Duration::from_secs(60);
524        while read.load(Ordering::Relaxed) < 8 {
525            assert!(
526                std::time::Instant::now() < deadline,
527                "the bytes never landed"
528            );
529            std::thread::sleep(std::time::Duration::from_millis(5));
530        }
531        assert_eq!(files_in(dir.path()), 1, "the partial file is there");
532        stop.store(true, Ordering::Relaxed);
533        let answer = worker.join().unwrap();
534        assert!(answer.is_err(), "a stopped read is not a dataset");
535        assert_eq!(files_in(dir.path()), 0, "and its file is gone");
536        drop(pipe);
537    }
538}