1use 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
27const CONTINUATION: [u8; 4] = [0xff; 4];
30
31const MAX_SCHEMA: usize = 16 << 20;
34
35pub 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
61const SCHEMA_PREFIX: usize = 64 << 10;
64
65pub fn is_stream_file(path: &Path) -> bool {
67 File::open(path).is_ok_and(is_stream)
68}
69
70fn is_stream(mut source: impl Read) -> bool {
73 let mut head = Vec::new();
74 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
104fn 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
112pub fn starts_with_stream(paths: &[PathBuf]) -> bool {
117 paths.first().is_some_and(|p| is_stream_file(p))
118}
119
120pub fn any_stream(paths: &[PathBuf]) -> bool {
122 paths.iter().any(|p| is_stream_file(p))
123}
124
125#[derive(Debug, Clone, PartialEq, Eq)]
127pub enum Part {
128 InPlace(PathBuf),
130 Converted {
132 source: PathBuf,
133 offset: u64,
134 rows: u64,
135 },
136}
137
138#[derive(Debug)]
140pub(crate) struct Converted {
141 pub file: TempDownload,
142 pub parts: Vec<Part>,
143}
144
145pub(crate) fn is_ipc_file_head(head: &[u8]) -> bool {
147 head.starts_with(b"ARROW1")
148}
149
150pub(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
188pub(crate) struct Merge<'a> {
192 out: Option<(FileWriter<BufWriter<File>>, PathBuf, Columns)>,
195 file: tempfile::NamedTempFile,
197 claim: crate::unfinished::Claim,
198 dir: PathBuf,
199 writer: &'a Writer,
200 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 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 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 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 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 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 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 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 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 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
376const ELSEWHERE: &str = "Choose another place with --temp-dir or the temp_dir setting.";
378
379pub(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
388fn 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
402fn 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
419struct Counting<'a, R> {
421 inner: R,
422 at: u64,
423 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
437struct 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 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 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 #[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 #[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 #[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 #[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 #[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 #[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 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 #[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 #[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}