Skip to main content

datui_lib/
ipc_stream.rs

1//! Arrow IPC streams: the format Hugging Face `datasets` writes its cache in.
2//!
3//! A stream is the IPC file format without the `ARROW1` magic and the footer that says
4//! where each record batch is, so Polars cannot scan it. An open converts the stream,
5//! or every stream shard of a directory, once into one IPC file in the temp directory
6//! and scans that; IPC files among the shards are scanned where they are, not copied
7//! ([`Part`]). The conversion holds one record batch at a time: 0.2 GiB at its peak
8//! for a 3.0 GiB stream. Read eagerly instead, that stream held 3.6 GiB for as long as
9//! it was open, and the same rows with ZSTD buffers, 1.3 GiB on disk, the same 3.6 GiB.
10
11use std::fs::File;
12use std::io::{BufRead, BufReader, BufWriter, Read, Seek, SeekFrom, Write};
13use std::path::{Path, PathBuf};
14use std::sync::Arc;
15use std::sync::atomic::{AtomicU64, Ordering};
16
17use color_eyre::{Result, eyre::eyre};
18use polars_arrow::io::ipc::format::ipc::planus::ReadAsRoot;
19use polars_arrow::io::ipc::format::ipc::{MessageHeaderRef, MessageRef};
20use polars_arrow::io::ipc::read::{StreamReader, StreamState, read_stream_metadata};
21use polars_arrow::io::ipc::write::{FileWriter, WriteOptions};
22
23use crate::download::TempDownload;
24use crate::error_display::{FileError, user_message_from_io};
25use crate::unfinished::Writer;
26
27/// What a stream's messages start with since Arrow 0.15. Older streams start with the
28/// schema message's length.
29const CONTINUATION: [u8; 4] = [0xff; 4];
30
31/// The longest schema message read to tell a stream from other bytes. A schema with
32/// thousands of columns and its metadata fits well inside.
33const MAX_SCHEMA: usize = 16 << 20;
34
35/// Whether `head`, the first bytes of a file, begins an Arrow IPC stream: a schema
36/// message, after the continuation marker or, in a stream older than it, without.
37///
38/// Where `head` holds the whole message it has to be a schema message. Where it is cut
39/// short, the marker followed by a length is taken as a stream; an older stream, with
40/// no marker to go on, is not. The message's fields are not read: Polars panics on a
41/// column type it has not implemented, and this runs on any file being opened.
42pub fn is_stream_head(head: &[u8]) -> bool {
43    let marked = head.starts_with(&CONTINUATION);
44    let rest = if marked { &head[4..] } else { head };
45    let Some(length) = rest.get(..4) else {
46        return false;
47    };
48    let length = i32::from_le_bytes([length[0], length[1], length[2], length[3]]);
49    let Ok(length) = usize::try_from(length) else {
50        return false;
51    };
52    if length == 0 || length > MAX_SCHEMA {
53        return false;
54    }
55    let Some(message) = rest.get(4..4 + length) else {
56        return marked;
57    };
58    begins_schema(message)
59}
60
61/// How much of a long schema message is read to see that it begins like one before the
62/// rest is: the schema table is near the front, its fields and metadata after it.
63const SCHEMA_PREFIX: usize = 64 << 10;
64
65/// Whether the file at `path` is an Arrow IPC stream, by its contents.
66pub fn is_stream_file(path: &Path) -> bool {
67    File::open(path).is_ok_and(is_stream)
68}
69
70/// Whether `source` holds an Arrow IPC stream. A file that only happens to start with
71/// a small number, as many binary files do, is read no further than [`SCHEMA_PREFIX`].
72fn is_stream(mut source: impl Read) -> bool {
73    let mut head = Vec::new();
74    // The marker and the length, then the message the length names.
75    if (&mut source).take(8).read_to_end(&mut head).is_err() {
76        return false;
77    }
78    let at = if head.starts_with(&CONTINUATION) {
79        4
80    } else {
81        0
82    };
83    let Some(length) = head
84        .get(at..at + 4)
85        .map(|b| i32::from_le_bytes([b[0], b[1], b[2], b[3]]))
86        .and_then(|l| usize::try_from(l).ok())
87        .filter(|l| (1..=MAX_SCHEMA).contains(l))
88    else {
89        return false;
90    };
91    let mut read_to = |end: usize, head: &mut Vec<u8>| {
92        let more = end.saturating_sub(head.len()) as u64;
93        (&mut source).take(more).read_to_end(head).is_ok()
94    };
95    let start = at + 4;
96    if length > SCHEMA_PREFIX
97        && !(read_to(start + SCHEMA_PREFIX, &mut head) && begins_schema(&head[start..]))
98    {
99        return false;
100    }
101    read_to(start + length, &mut head) && is_stream_head(&head)
102}
103
104/// Whether `message`, all or the front of one, is a schema message as far as it goes.
105fn begins_schema(message: &[u8]) -> bool {
106    matches!(
107        MessageRef::read_as_root(message).and_then(|m| m.header()),
108        Ok(Some(MessageHeaderRef::Schema(_)))
109    )
110}
111
112/// Whether the Arrow `paths` are read by converting their streams: when the first is a
113/// stream. Only the first is opened. The conversion opens the rest, and leaves the IPC
114/// files among them where they are; a stream behind an IPC file is found by
115/// [`any_stream`] once the scan, which reads every IPC file's footer, has failed on it.
116pub fn starts_with_stream(paths: &[PathBuf]) -> bool {
117    paths.first().is_some_and(|p| is_stream_file(p))
118}
119
120/// Whether any of `paths` is a stream: asked once a scan of them as IPC files failed.
121pub fn any_stream(paths: &[PathBuf]) -> bool {
122    paths.iter().any(|p| is_stream_file(p))
123}
124
125/// Where one input's rows are read from once its streams are converted.
126#[derive(Debug, Clone, PartialEq, Eq)]
127pub enum Part {
128    /// An IPC file, scanned where it is: a local path or an object's URL.
129    InPlace(PathBuf),
130    /// A stream, `source`, converted: `rows` rows of the converted file from `offset`.
131    Converted {
132        source: PathBuf,
133        offset: u64,
134        rows: u64,
135    },
136}
137
138/// What a conversion wrote: the IPC file of the streams, and each input's place.
139#[derive(Debug)]
140pub(crate) struct Converted {
141    pub file: TempDownload,
142    pub parts: Vec<Part>,
143}
144
145/// Whether an Arrow file starts as an IPC file does, with its footer at the end.
146pub(crate) fn is_ipc_file_head(head: &[u8]) -> bool {
147    head.starts_with(b"ARROW1")
148}
149
150/// Convert the Arrow streams at `paths`, in order, into one IPC file in `temp_dir` (the
151/// system temp directory when `None`), created and claimed through `writer` and
152/// removed if the open stops or this fails. The IPC files among them are not copied:
153/// they are scanned where they are, in their place among the streams. `read` counts
154/// the bytes looked at so far.
155///
156/// One record batch is in memory at a time. Buffers compressed with LZ4 or ZSTD are
157/// written out uncompressed, so the file maps and scans like any other.
158pub(crate) fn convert(
159    paths: &[PathBuf],
160    temp_dir: Option<&Path>,
161    writer: &Writer,
162    read: &AtomicU64,
163) -> Result<Converted> {
164    let mut merge = Merge::create(temp_dir, writer)?;
165    let mut parts = Vec::with_capacity(paths.len());
166    let mut before = 0;
167    for path in paths {
168        let mut source = File::open(path)?;
169        let size = source.metadata()?.len();
170        let mut head = Vec::with_capacity(6);
171        (&mut source).take(6).read_to_end(&mut head)?;
172        if is_ipc_file_head(&head) {
173            parts.push(Part::InPlace(path.clone()));
174        } else {
175            source.seek(SeekFrom::Start(0))?;
176            has_room(size, temp_dir)?;
177            parts.push(merge.append(source, path, before, read)?);
178        }
179        before += size;
180        read.store(before, Ordering::Relaxed);
181    }
182    Ok(Converted {
183        file: merge.finish()?,
184        parts,
185    })
186}
187
188/// One IPC file being written from the batches of Arrow streams, appended one at a
189/// time: what a conversion writes, and what a bucket's streams are read into as they
190/// download. Dropped unfinished, the file goes.
191pub(crate) struct Merge<'a> {
192    /// The writer, the first input's name and its columns, once one is appended. Its
193    /// handle on the file is let go before the file is removed.
194    out: Option<(FileWriter<BufWriter<File>>, PathBuf, Columns)>,
195    /// The file, then its claim: dropped in that order.
196    file: tempfile::NamedTempFile,
197    claim: crate::unfinished::Claim,
198    dir: PathBuf,
199    writer: &'a Writer,
200    /// The rows written so far.
201    rows: u64,
202}
203
204type Columns = Vec<(
205    polars::prelude::PlSmallStr,
206    polars_arrow::datatypes::ArrowDataType,
207)>;
208
209fn columns(schema: &polars_arrow::datatypes::ArrowSchema) -> Columns {
210    schema
211        .iter_values()
212        .map(|f| (f.name.clone(), f.dtype.clone()))
213        .collect()
214}
215
216impl<'a> Merge<'a> {
217    /// The empty file, in `temp_dir`, claimed through `writer`.
218    pub(crate) fn create(temp_dir: Option<&Path>, writer: &'a Writer) -> Result<Self> {
219        let Some((file, claim)) =
220            writer.create(|| TempDownload::create(temp_dir, Some("arrow")))?
221        else {
222            return Err(stopped());
223        };
224        Ok(Self {
225            file,
226            claim,
227            dir: temp_dir
228                .map(Path::to_path_buf)
229                .unwrap_or_else(std::env::temp_dir),
230            writer,
231            out: None,
232            rows: 0,
233        })
234    }
235
236    /// Append the batches of the stream `source` reads, which errors call `name`,
237    /// counting its bytes into `read` after the `before` read ahead of it. Its place
238    /// in the file is returned.
239    pub(crate) fn append(
240        &mut self,
241        source: impl Read,
242        name: &Path,
243        before: u64,
244        read: &AtomicU64,
245    ) -> Result<Part> {
246        let offset = self.rows;
247        match self.batches(source, name, before, read) {
248            Ok(true) => Ok(Part::Converted {
249                source: name.to_path_buf(),
250                offset,
251                rows: self.rows - offset,
252            }),
253            Ok(false) => Err(stopped()),
254            Err(e) => Err(out_of_room(e, &self.dir)),
255        }
256    }
257
258    /// The finished file, held with its claim.
259    pub(crate) fn finish(self) -> Result<TempDownload> {
260        let Some((mut out, _, _)) = self.out else {
261            return Err(eyre!("No Arrow IPC stream to read."));
262        };
263        let finished = out
264            .finish()
265            .map_err(color_eyre::Report::from)
266            .and_then(|()| Ok(out.into_inner().flush()?));
267        if let Err(e) = finished {
268            return Err(out_of_room(e, &self.dir));
269        }
270        Ok(TempDownload::held(self.file, Some(self.claim)))
271    }
272
273    /// `false` when the open was stopped first.
274    fn batches(
275        &mut self,
276        source: impl Read,
277        name: &Path,
278        before: u64,
279        read: &AtomicU64,
280    ) -> Result<bool> {
281        let mut reader = Forward {
282            inner: BufReader::with_capacity(
283                1 << 20,
284                Counting {
285                    inner: source,
286                    at: 0,
287                    before,
288                    read,
289                },
290            ),
291            at: 0,
292        };
293        let unreadable = move |e: &dyn std::fmt::Display| -> color_eyre::Report {
294            FileError::new(name, format!("not a readable Arrow IPC stream: {e}")).into()
295        };
296        // A read that failed under the stream (a download cut off) says so, not that
297        // the stream is damaged.
298        let failed = move |e: polars::prelude::PolarsError| match e {
299            polars::prelude::PolarsError::IO { error, .. } => {
300                FileError::new(name, user_message_from_io(&error, None)).into()
301            }
302            e => unreadable(&e),
303        };
304        // Polars panics on a column type it has not implemented, such as run-end
305        // encoding: that is a stream it cannot read, not a crash.
306        let metadata = crate::logging::catch_panic(|| read_stream_metadata(&mut reader))
307            .map_err(|_| unreadable(&"it has a column type Polars cannot read"))?
308            .map_err(failed)?;
309        self.start(
310            name,
311            &metadata.schema,
312            &metadata.ipc_schema.fields,
313            metadata.custom_schema_metadata.as_ref(),
314        )?;
315        let mut batches = StreamReader::new(reader, metadata, None);
316        let (out, _, _) = self.out.as_mut().expect("started just above");
317        loop {
318            if self.writer.stopped() {
319                return Ok(false);
320            }
321            // Polars also panics on some malformed record batches, rather than erring.
322            let next = crate::logging::catch_panic(|| batches.next())
323                .map_err(|_| unreadable(&"a record batch in it is damaged"))?;
324            let batch = match next {
325                Some(Ok(StreamState::Some(batch))) => batch,
326                // The end of a stream written without its end-of-stream marker.
327                Some(Ok(StreamState::Waiting)) | None => break,
328                Some(Err(e)) => return Err(failed(e)),
329            };
330            self.rows += batch.len() as u64;
331            out.write(&batch, None)?;
332        }
333        Ok(true)
334    }
335
336    /// Start the file with the first input's schema, or check a later one has the same
337    /// columns.
338    fn start(
339        &mut self,
340        name: &Path,
341        schema: &polars_arrow::datatypes::ArrowSchema,
342        fields: &[polars_arrow::io::ipc::IpcField],
343        custom: Option<&polars_arrow::datatypes::Metadata>,
344    ) -> Result<()> {
345        match &self.out {
346            None => {
347                let mut out = FileWriter::try_new(
348                    BufWriter::with_capacity(1 << 20, self.file.as_file().try_clone()?),
349                    Arc::new(schema.clone()),
350                    Some(fields.to_vec()),
351                    WriteOptions { compression: None },
352                )?;
353                if let Some(custom) = custom {
354                    out.set_custom_schema_metadata(Arc::new(custom.clone()));
355                }
356                self.out = Some((out, name.to_path_buf(), columns(schema)));
357            }
358            Some((_, first, first_columns)) => {
359                if columns(schema) != *first_columns {
360                    return Err(eyre!(
361                        "{} has different columns from {}, so they cannot be read as one table.",
362                        name.display(),
363                        first.display()
364                    ));
365                }
366            }
367        }
368        Ok(())
369    }
370}
371
372fn stopped() -> color_eyre::Report {
373    eyre!("Converting the Arrow stream was stopped.")
374}
375
376/// What to do about a temp directory too small for the copy.
377const ELSEWHERE: &str = "Choose another place with --temp-dir or the temp_dir setting.";
378
379/// Whether `temp_dir` (the system temp directory when `None`) has room for a copy of
380/// `needs` bytes of Arrow; see [`room`].
381pub(crate) fn has_room(needs: u64, temp_dir: Option<&Path>) -> Result<()> {
382    let dir = temp_dir
383        .map(Path::to_path_buf)
384        .unwrap_or_else(std::env::temp_dir);
385    room(needs, crate::local_copy::free_space(&dir), &dir)
386}
387
388/// The copy is about the size of the streams, larger where their buffers are
389/// compressed: refused before any of it is written where `dir` has less free.
390fn room(needs: u64, free: Option<u64>, dir: &Path) -> Result<()> {
391    match free {
392        Some(free) if free < needs => Err(eyre!(
393            "Converting the Arrow stream needs {} free in {}, which has {}. {ELSEWHERE}",
394            crate::discover::format_size(needs),
395            dir.display(),
396            crate::discover::format_size(free),
397        )),
398        _ => Ok(()),
399    }
400}
401
402/// A write that ran out of space, said so with where and what to do.
403fn out_of_room(error: color_eyre::Report, dir: &Path) -> color_eyre::Report {
404    let full = error.chain().any(|cause| {
405        cause
406            .downcast_ref::<std::io::Error>()
407            .is_some_and(|e| e.kind() == std::io::ErrorKind::StorageFull)
408    });
409    if full {
410        eyre!(
411            "{} ran out of space for the converted Arrow stream, which is written uncompressed. {ELSEWHERE}",
412            dir.display()
413        )
414    } else {
415        error
416    }
417}
418
419/// A stream being read, counting its bytes into the open's progress.
420struct Counting<'a, R> {
421    inner: R,
422    at: u64,
423    /// The bytes of the inputs before this one.
424    before: u64,
425    read: &'a AtomicU64,
426}
427
428impl<R: Read> Read for Counting<'_, R> {
429    fn read(&mut self, buf: &mut [u8]) -> std::io::Result<usize> {
430        let n = self.inner.read(buf)?;
431        self.at += n as u64;
432        self.read.store(self.before + self.at, Ordering::Relaxed);
433        Ok(n)
434    }
435}
436
437/// A stream read front to back, by the stream reader that seeks: only forward, over
438/// the padding after a message, so a download can be read as it arrives.
439struct Forward<R> {
440    inner: R,
441    at: u64,
442}
443
444impl<R: BufRead> Read for Forward<R> {
445    fn read(&mut self, buf: &mut [u8]) -> std::io::Result<usize> {
446        let n = self.inner.read(buf)?;
447        self.at += n as u64;
448        Ok(n)
449    }
450}
451
452impl<R: BufRead> Seek for Forward<R> {
453    fn seek(&mut self, pos: SeekFrom) -> std::io::Result<u64> {
454        let to = match pos {
455            SeekFrom::Start(to) => Some(to),
456            SeekFrom::Current(by) => self.at.checked_add_signed(by),
457            SeekFrom::End(_) => None,
458        };
459        let skip = to
460            .and_then(|to| to.checked_sub(self.at))
461            .ok_or_else(|| std::io::Error::other("an Arrow stream is read front to back"))?;
462        let skipped = std::io::copy(&mut (&mut self.inner).take(skip), &mut std::io::sink())?;
463        self.at += skipped;
464        Ok(self.at)
465    }
466}
467
468#[cfg(test)]
469pub(crate) mod tests {
470    use super::*;
471    use crate::unfinished::Unfinished;
472    use polars::prelude::*;
473    use polars_arrow::io::ipc::write::{Compression, StreamWriter};
474    use std::sync::atomic::AtomicBool;
475
476    fn frame(from: i64, n: i64) -> DataFrame {
477        df!(
478            "id" => (from..from + n).collect::<Vec<_>>(),
479            "text" => (from..from + n).map(|i| format!("row {i}")).collect::<Vec<_>>(),
480        )
481        .unwrap()
482    }
483
484    /// `df` as a stream of batches of three rows, as pyarrow and Hugging Face write it.
485    pub(crate) fn stream(
486        df: &DataFrame,
487        compression: Option<Compression>,
488        legacy: bool,
489    ) -> Vec<u8> {
490        let mut out = Vec::new();
491        let mut writer = StreamWriter::new(&mut out, WriteOptions { compression });
492        writer
493            .start(&df.schema().to_arrow(CompatLevel::newest()), None)
494            .unwrap();
495        for at in (0..df.height()).step_by(3) {
496            let part = df.slice(at as i64, 3);
497            for batch in part.iter_chunks(CompatLevel::newest(), false) {
498                writer.write(&batch, None).unwrap();
499            }
500        }
501        writer.finish().unwrap();
502        if legacy { strip_markers(&out) } else { out }
503    }
504
505    /// The same stream without the continuation markers, as Arrow wrote before 0.15.
506    fn strip_markers(stream: &[u8]) -> Vec<u8> {
507        let mut out = Vec::new();
508        let mut at = 0;
509        while at < stream.len() {
510            assert_eq!(stream[at..at + 4], CONTINUATION);
511            let length = i32::from_le_bytes(stream[at + 4..at + 8].try_into().unwrap());
512            out.extend_from_slice(&stream[at + 4..at + 8]);
513            at += 8;
514            if length == 0 {
515                break;
516            }
517            let message = &stream[at..at + length as usize];
518            let body = MessageRef::read_as_root(message)
519                .unwrap()
520                .body_length()
521                .unwrap() as usize;
522            out.extend_from_slice(&stream[at..at + length as usize + body]);
523            at += length as usize + body;
524        }
525        out
526    }
527
528    fn writer(stop: bool) -> Writer {
529        Unfinished::default().writer(Arc::new(AtomicBool::new(stop)))
530    }
531
532    fn rows(path: &Path) -> DataFrame {
533        LazyFrame::scan_ipc(
534            PlRefPath::try_from_path(path).unwrap(),
535            Default::default(),
536            Default::default(),
537        )
538        .unwrap()
539        .collect()
540        .unwrap()
541    }
542
543    #[test]
544    fn a_stream_is_told_by_its_first_message() {
545        let df = frame(0, 5);
546        let marked = stream(&df, None, false);
547        let legacy = stream(&df, None, true);
548        assert!(is_stream_head(&marked));
549        assert!(is_stream_head(&legacy));
550        assert!(is_stream_head(&marked[..8]), "the marker and a length");
551        assert!(
552            !is_stream_head(&legacy[..8]),
553            "no marker, and too little to parse"
554        );
555        let mut file = Vec::new();
556        polars::io::ipc::IpcWriter::new(&mut file)
557            .finish(&mut df.clone())
558            .unwrap();
559        for not in [
560            &file[..],
561            b"id,text\n0,a\n",
562            b"\xff\xff\xff\xff\x00\x00\x00\x00",
563            b"\xff\xff\xff\xff\xff\xff\xff\x7f",
564            b"\x10\x00\x00\x00garbage garbage garbage",
565            b"",
566        ] {
567            assert!(!is_stream_head(not), "{not:?}");
568        }
569    }
570
571    /// A schema message too long to read whole at a glance is still told by its bytes,
572    /// with or without the marker; a file that only starts with a small number, as
573    /// many binary files do, is read no further than its front.
574    #[test]
575    fn a_long_schema_is_read_but_a_lookalike_is_not() {
576        let names: Vec<String> = (0..3000)
577            .map(|i| format!("a_long_column_name_{i:05}"))
578            .collect();
579        let df = DataFrame::new(
580            1,
581            names
582                .iter()
583                .map(|n| Column::new(n.as_str().into(), [1i32]))
584                .collect(),
585        )
586        .unwrap();
587        for legacy in [false, true] {
588            let bytes = stream(&df, None, legacy);
589            let at = if legacy { 0 } else { 4 };
590            let length = i32::from_le_bytes(bytes[at..at + 4].try_into().unwrap()) as usize;
591            assert!(length > SCHEMA_PREFIX, "{length}");
592            assert!(is_stream(&bytes[..]), "legacy: {legacy}");
593        }
594
595        struct Counted<'a>(&'a [u8], usize);
596        impl Read for Counted<'_> {
597            fn read(&mut self, buf: &mut [u8]) -> std::io::Result<usize> {
598                let n = self.0.read(buf)?;
599                self.1 += n;
600                Ok(n)
601            }
602        }
603        let mut lookalike = vec![7u8; 12 << 20];
604        lookalike[..4].copy_from_slice(&(10i32 << 20).to_le_bytes());
605        let mut source = Counted(&lookalike, 0);
606        assert!(!is_stream(&mut source));
607        assert!(source.1 <= 4 + SCHEMA_PREFIX, "read {}", source.1);
608    }
609
610    /// Streams of every kind, compressed or not, with or without the markers, become
611    /// one IPC file holding their rows in order; the bytes read are counted.
612    #[test]
613    fn streams_convert_to_one_ipc_file() {
614        let dir = tempfile::tempdir().unwrap();
615        let kinds = [
616            (None, false),
617            (Some(Compression::LZ4), false),
618            (Some(Compression::ZSTD(Default::default())), false),
619            (None, true),
620        ];
621        let mut paths = Vec::new();
622        for (i, (compression, legacy)) in kinds.into_iter().enumerate() {
623            let path = dir.path().join(format!("data-{i:05}-of-00004.arrow"));
624            std::fs::write(
625                &path,
626                stream(&frame(i as i64 * 10, 10), compression, legacy),
627            )
628            .unwrap();
629            assert!(is_stream_file(&path), "{path:?}");
630            paths.push(path);
631        }
632        let out = tempfile::tempdir().unwrap();
633        let read = AtomicU64::new(0);
634        let converted = convert(&paths, Some(out.path()), &writer(false), &read).unwrap();
635        let file = converted.file;
636        assert!(!is_stream_file(file.path()), "an IPC file now");
637        assert_eq!(
638            converted.parts[1],
639            Part::Converted {
640                source: paths[1].clone(),
641                offset: 10,
642                rows: 10
643            }
644        );
645        let df = rows(file.path());
646        assert_eq!(df.height(), 40);
647        assert_eq!(
648            df.column("id").unwrap().i64().unwrap().get(25),
649            Some(25),
650            "in order"
651        );
652        let total: u64 = paths
653            .iter()
654            .map(|p| std::fs::metadata(p).unwrap().len())
655            .sum();
656        assert_eq!(read.load(Ordering::Relaxed), total);
657        let path = file.path().to_path_buf();
658        drop(file);
659        assert!(!path.exists());
660    }
661
662    /// A stopped open writes nothing and leaves nothing; a shard with other columns,
663    /// and a stream that is not one, are refused by name.
664    #[test]
665    fn a_stop_or_a_bad_stream_leaves_no_file() {
666        let dir = tempfile::tempdir().unwrap();
667        let good = dir.path().join("a.arrow");
668        std::fs::write(&good, stream(&frame(0, 9), None, false)).unwrap();
669        let out = tempfile::tempdir().unwrap();
670        let empty = || std::fs::read_dir(out.path()).unwrap().next().is_none();
671        let read = AtomicU64::new(0);
672
673        let stopped = convert(
674            std::slice::from_ref(&good),
675            Some(out.path()),
676            &writer(true),
677            &read,
678        );
679        assert!(stopped.is_err());
680        assert!(empty());
681
682        let other = dir.path().join("b.arrow");
683        let df = df!("x" => [1.5f64]).unwrap();
684        std::fs::write(&other, stream(&df, None, false)).unwrap();
685        let error = convert(
686            &[good.clone(), other.clone()],
687            Some(out.path()),
688            &writer(false),
689            &read,
690        )
691        .unwrap_err()
692        .to_string();
693        assert!(
694            error.contains("b.arrow has different columns from"),
695            "{error}"
696        );
697        assert!(empty());
698
699        let cut = dir.path().join("cut.arrow");
700        let bytes = stream(&frame(0, 9), None, false);
701        std::fs::write(&cut, &bytes[..bytes.len() / 2]).unwrap();
702        let error = convert(&[cut], Some(out.path()), &writer(false), &read)
703            .unwrap_err()
704            .to_string();
705        assert!(
706            error.contains("cut.arrow\": Not a readable Arrow IPC stream"),
707            "{error}"
708        );
709        assert!(empty());
710    }
711
712    /// A stream with no name to go on is Arrow by its bytes, to discovery and to a pipe.
713    #[test]
714    fn discovery_and_a_pipe_know_a_stream_by_its_bytes() {
715        let dir = tempfile::tempdir().unwrap();
716        for (name, legacy) in [("part-0", false), ("part-1", true)] {
717            let bytes = stream(&frame(0, 4), None, legacy);
718            let path = dir.path().join(name);
719            std::fs::write(&path, &bytes).unwrap();
720            assert_eq!(
721                crate::discover::sniff_format(&path),
722                Some(crate::FileFormat::Arrow),
723                "{name}"
724            );
725            assert_eq!(
726                crate::stdin::sniff(&bytes),
727                (crate::FileFormat::Arrow, None),
728                "{name}"
729            );
730        }
731    }
732
733    /// A column type Polars has not implemented makes a stream it cannot read, told
734    /// as one by its bytes and refused by the conversion rather than panicking it.
735    #[test]
736    fn a_stream_polars_cannot_read_is_refused() {
737        let bytes =
738            include_bytes!("../../../fuzz/corpus/ipc_stream_head/regression-run-end-encoded");
739        assert!(is_stream_head(bytes));
740        let dir = tempfile::tempdir().unwrap();
741        let path = dir.path().join("ree.arrow");
742        std::fs::write(&path, bytes).unwrap();
743        let error = convert(
744            &[path],
745            Some(dir.path()),
746            &writer(false),
747            &AtomicU64::new(0),
748        )
749        .unwrap_err()
750        .to_string();
751        assert!(
752            error.contains("a column type Polars cannot read"),
753            "{error}"
754        );
755        assert_eq!(std::fs::read_dir(dir.path()).unwrap().count(), 1);
756    }
757
758    /// A record batch Polars panics on, rather than erring, is a stream it cannot read,
759    /// refused by name with nothing left behind.
760    #[test]
761    fn a_damaged_record_batch_is_refused() {
762        let df = df!(
763            "id" => (0..50i64).collect::<Vec<_>>(),
764            "t" => (0..50).map(|i| format!("r{i}")).collect::<Vec<_>>(),
765        )
766        .unwrap();
767        let mut bytes = Vec::new();
768        let mut out = StreamWriter::new(&mut bytes, WriteOptions { compression: None });
769        out.start(&df.schema().to_arrow(CompatLevel::newest()), None)
770            .unwrap();
771        for batch in df.iter_chunks(CompatLevel::newest(), false) {
772            out.write(&batch, None).unwrap();
773        }
774        out.finish().unwrap();
775        // In the record batch's message, so that Polars cannot read its length, which
776        // it unwraps (found by mutating this stream).
777        assert_eq!(bytes[251], 0);
778        bytes[251] = 0x21;
779        let dir = tempfile::tempdir().unwrap();
780        let path = dir.path().join("damaged.arrow");
781        std::fs::write(&path, &bytes).unwrap();
782        let out = tempfile::tempdir().unwrap();
783        let error = convert(
784            &[path],
785            Some(out.path()),
786            &writer(false),
787            &AtomicU64::new(0),
788        )
789        .unwrap_err()
790        .to_string();
791        assert!(
792            error.contains("damaged.arrow\": Not a readable Arrow IPC stream"),
793            "{error}"
794        );
795        assert!(std::fs::read_dir(out.path()).unwrap().next().is_none());
796    }
797
798    /// A temp directory with less free than the streams refuses the conversion before
799    /// writing, and one that fills up says where, both with what to do.
800    #[test]
801    fn a_full_temp_directory_says_so() {
802        let dir = Path::new("/scratch");
803        assert!(room(10, Some(10), dir).is_ok());
804        assert!(room(10, None, dir).is_ok(), "free space unknown");
805        let error = room(2 << 30, Some(1 << 30), dir).unwrap_err().to_string();
806        assert!(
807            error.contains("needs 2.0 GB free in /scratch, which has 1.0 GB"),
808            "{error}"
809        );
810        assert!(error.contains("--temp-dir"), "{error}");
811
812        let full =
813            polars::error::PolarsError::from(std::io::Error::from(std::io::ErrorKind::StorageFull));
814        let error = out_of_room(full.into(), dir).to_string();
815        assert!(error.starts_with("/scratch ran out of space"), "{error}");
816        assert!(error.contains("--temp-dir"), "{error}");
817        let other = out_of_room(eyre!("something else"), dir).to_string();
818        assert_eq!(other, "something else");
819    }
820
821    /// Only the first file says whether a list of Arrow files is converted. The
822    /// conversion copies only the streams; an IPC file among them keeps its place and
823    /// is read where it is.
824    #[test]
825    fn only_the_streams_among_ipc_files_are_converted() {
826        let dir = tempfile::tempdir().unwrap();
827        let streamed = dir.path().join("s.arrow");
828        std::fs::write(&streamed, stream(&frame(0, 3), None, false)).unwrap();
829        let file = dir.path().join("f.arrow");
830        polars::io::ipc::IpcWriter::new(std::fs::File::create(&file).unwrap())
831            .finish(&mut frame(3, 4))
832            .unwrap();
833        assert!(!starts_with_stream(std::slice::from_ref(&file)));
834        assert!(starts_with_stream(&[streamed.clone(), file.clone()]));
835        assert!(!starts_with_stream(&[file.clone(), streamed.clone()]));
836        assert!(any_stream(&[file.clone(), streamed.clone()]));
837        assert!(!any_stream(std::slice::from_ref(&file)));
838        let out = tempfile::tempdir().unwrap();
839        let converted_part = Part::Converted {
840            source: streamed.clone(),
841            offset: 0,
842            rows: 3,
843        };
844        for (paths, parts) in [
845            (
846                [streamed.clone(), file.clone()],
847                [converted_part.clone(), Part::InPlace(file.clone())],
848            ),
849            (
850                [file.clone(), streamed.clone()],
851                [Part::InPlace(file.clone()), converted_part.clone()],
852            ),
853        ] {
854            let read = AtomicU64::new(0);
855            let converted = convert(&paths, Some(out.path()), &writer(false), &read).unwrap();
856            assert_eq!(converted.parts, parts);
857            let df = rows(converted.file.path());
858            assert_eq!(df.height(), 3, "the stream's rows only");
859            let total: u64 = paths
860                .iter()
861                .map(|p| std::fs::metadata(p).unwrap().len())
862                .sum();
863            assert_eq!(read.load(Ordering::Relaxed), total);
864        }
865    }
866}