1use 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
17pub const PATH: &str = "-";
19
20pub const NAME: &str = "stdin";
22
23pub fn is_stdin(path: &Path) -> bool {
25 path.as_os_str() == PATH
26}
27
28pub fn named(path: &Path) -> PathBuf {
30 if is_stdin(path) {
31 PathBuf::from(NAME)
32 } else {
33 path.to_path_buf()
34 }
35}
36
37pub 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
46pub fn carries_data(terminal: bool, device: bool) -> bool {
49 !terminal && !device
50}
51
52#[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#[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 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
84pub fn as_file(path: PathBuf) -> PathBuf {
87 if is_stdin(&path) {
88 Path::new(".").join(path)
89 } else {
90 path
91 }
92}
93
94pub 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
104pub 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
117pub fn sniff(head: &[u8]) -> (FileFormat, Option<CompressionFormat>) {
122 let (format, compression, _) = sniffed(head);
123 (format, compression)
124}
125
126pub(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
139fn 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 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
158fn 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
179const HEAD: usize = crate::readers::HEAD;
181
182pub(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
198const OLD_SPOOL: std::time::Duration = std::time::Duration::from_secs(7 * 24 * 60 * 60);
201
202fn 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 #[cfg(unix)]
217 let free = !crate::download::held_elsewhere(&entry.path());
218 #[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
227pub(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
251pub(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 (Some(named), None) if named.separator().is_some() => OpenOptions {
274 compression,
275 ..options
276 },
277 (Some(_), _) => options,
279 (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
294pub(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 #[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 (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 (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 (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 (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 #[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 #[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 #[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 assert_eq!(options.format, Some(FileFormat::Text));
462 assert_eq!(options.compression, Some(CompressionFormat::Zstd));
463
464 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 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 #[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}