use std::fs::File;
use std::io::{BufReader, Read};
use std::path::Path;
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use color_eyre::Result;
use color_eyre::eyre::eyre;
use polars::prelude::{DataFrame, LazyFrame, PolarsResult};
use crate::cloud::download::TempDownload;
use crate::formats::model_files::MetaValue;
use crate::formats::readers::{ConvertIn, ConvertOut};
use crate::formats::segments::{Converted, Segments};
use crate::loading::unfinished::Writer;
use crate::notes::Note;
use crate::numfmt::group_chrome;
use crate::{CompressionFormat, OpenOptions};
const CHUNK: usize = 1 << 16;
#[derive(Debug, Clone, Default, PartialEq)]
pub struct Detail {
pub tab: &'static str,
pub lines: Vec<String>,
pub warnings: Vec<String>,
pub list_title: &'static str,
pub list: Vec<(String, MetaValue)>,
pub first: bool,
pub own_columns: bool,
pub tables: Vec<String>,
pub table: Option<String>,
}
pub const fn tab(format: crate::FileFormat) -> &'static str {
match format.descriptor().summary {
crate::Summary::Tab(tab) => tab,
crate::Summary::None(_) => format.descriptor().title,
}
}
pub trait BatchReader {
fn push(&mut self, piece: &[u8]) -> Result<()>;
fn take_batch(&mut self) -> PolarsResult<Option<DataFrame>>;
fn finish(&mut self) -> Result<DataFrame>;
}
struct Counted<'a, R> {
inner: R,
read: &'a AtomicU64,
}
impl<R: Read> Read for Counted<'_, R> {
fn read(&mut self, buf: &mut [u8]) -> std::io::Result<usize> {
let n = self.inner.read(buf)?;
self.read.fetch_add(n as u64, Ordering::Relaxed);
Ok(n)
}
}
pub(crate) fn open_reader<'a>(
file: &Path,
options: &OpenOptions,
read: &'a AtomicU64,
) -> Result<Box<dyn Read + 'a>> {
let f = BufReader::new(Counted {
inner: File::open(file)?,
read,
});
Ok(
match options
.compression
.or_else(|| CompressionFormat::from_extension(file))
{
None => Box::new(f),
Some(CompressionFormat::Gzip) => Box::new(flate2::read::MultiGzDecoder::new(f)),
Some(CompressionFormat::Zstd) => Box::new(zstd::Decoder::new(f)?),
Some(CompressionFormat::Bzip2) => Box::new(bzip2::read::BzDecoder::new(f)),
Some(CompressionFormat::Xz) => Box::new(xz2::read::XzDecoder::new(f)),
},
)
}
pub(crate) fn read_through<R: BatchReader>(
file: &Path,
options: &OpenOptions,
writer: &Writer,
read: &AtomicU64,
reader: &mut R,
mut each: impl FnMut(&DataFrame),
) -> Result<(LazyFrame, Vec<TempDownload>)> {
let mut source = open_reader(file, options, read)?;
let mut segments = Segments::new(options, writer);
let mut chunk = vec![0u8; CHUNK];
loop {
if writer.stopped() {
return Err(eyre!("Reading was stopped."));
}
let n = match source.read(&mut chunk) {
Ok(0) => break,
Ok(n) => n,
Err(e) if e.kind() == std::io::ErrorKind::Interrupted => continue,
Err(e) => return Err(e.into()),
};
reader.push(&chunk[..n])?;
if let Some(df) = reader.take_batch()? {
each(&df);
segments.write(&df)?;
}
}
let last = reader.finish()?;
each(&last);
segments.write(&last)?;
segments.finish()
}
pub(crate) fn convert_with<R: BatchReader>(
input: &ConvertIn<'_>,
mut reader: R,
finish: impl FnOnce(&R, LazyFrame) -> Result<(LazyFrame, Vec<Note>, Detail)>,
) -> ConvertOut {
let [file] = input.files else {
return Err(eyre!("Open {} files one at a time.", input.format.name()));
};
let (lf, files) = read_through(
file,
input.options,
input.writer,
input.read,
&mut reader,
|_| {},
)?;
let (lf, notes, detail) = finish(&reader, lf)?;
let converted = Converted {
lf,
files,
notes,
other_tables: Vec::new(),
};
Ok((converted, Some(Arc::new(detail))))
}
pub(crate) fn note(summary: String, scope: String) -> Note {
Note {
summary,
scope,
read_as_text: None,
passed_over: None,
}
}
pub(crate) fn count(n: u64, one: &str, many: &str) -> String {
let n = usize::try_from(n).unwrap_or(usize::MAX);
format!("{} {}", group_chrome(n), if n == 1 { one } else { many })
}
pub(crate) fn capped_list(
rows: impl Iterator<Item = (String, MetaValue)>,
total: usize,
) -> Vec<(String, MetaValue)> {
let mut list: Vec<_> = rows.take(crate::limits::get().detail_rows).collect();
if total > list.len() {
let more = total - list.len();
list.push((
crate::glyphs::get().ellipsis.to_string(),
MetaValue::Text(format!(
"{} more {} limits.detail_rows raises it",
group_chrome(more),
crate::glyphs::get().middot
)),
));
}
list
}