Skip to main content

datui_lib/formats/
files.rs

1//! A directory of one spec's files, `{date}/{venue}/trades.bin`: one table, the parts
2//! of each file's path as columns.
3
4use super::{Bytes, Opened, Spec, SpecRecords};
5use polars::prelude::*;
6use std::path::{Path, PathBuf};
7use std::sync::Arc;
8
9/// The most files one tree is read from.
10const MAX_FILES: usize = 100_000;
11
12/// Each file's records and the values of its path's parts.
13pub struct SpecFiles {
14    parts: Vec<(Arc<dyn SpecRecords>, Vec<AnyValue<'static>>)>,
15    part_fields: Vec<(PlSmallStr, DataType)>,
16    /// The first row of each file.
17    starts: Vec<usize>,
18    rows: usize,
19    schema: SchemaRef,
20    sources: Vec<Arc<Bytes>>,
21}
22
23/// One part's pattern, compiled: its regex, and the part names it captures.
24fn component_regex(component: &str) -> Result<(regex::Regex, Vec<String>), String> {
25    let mut out = String::from("^");
26    let mut names = Vec::new();
27    let mut rest = component;
28    while let Some(open) = rest.find('{') {
29        out.push_str(&regex::escape(&rest[..open]));
30        let after = &rest[open + 1..];
31        let close = after.find('}').ok_or("a `{` without its `}`")?;
32        let inside = &after[..close];
33        let name = inside.split_once(':').map_or(inside, |(n, _)| n);
34        out.push_str("(.+?)");
35        names.push(name.to_string());
36        rest = &after[close + 1..];
37    }
38    out.push_str(&regex::escape(rest));
39    out.push('$');
40    Ok((regex::Regex::new(&out).map_err(|e| e.to_string())?, names))
41}
42
43/// A part's name and the text it matched in a path.
44type PartText = (String, String);
45
46/// The files under `dir` the pattern names, each with its parts' text.
47fn matching(dir: &Path, pattern: &str) -> Result<Vec<(PathBuf, Vec<PartText>)>, String> {
48    let components: Vec<(regex::Regex, Vec<String>)> = pattern
49        .split('/')
50        .map(component_regex)
51        .collect::<Result<_, _>>()?;
52    let mut found = Vec::new();
53    let mut stack = vec![(dir.to_path_buf(), 0usize, Vec::<PartText>::new())];
54    while let Some((at, depth, parts)) = stack.pop() {
55        let Ok(listing) = std::fs::read_dir(&at) else {
56            continue;
57        };
58        let (re, names) = &components[depth];
59        let last = depth + 1 == components.len();
60        for entry in listing.flatten() {
61            let name = entry.file_name().to_string_lossy().into_owned();
62            let Some(caps) = re.captures(&name) else {
63                continue;
64            };
65            let mut parts = parts.clone();
66            for (i, part) in names.iter().enumerate() {
67                parts.push((
68                    part.clone(),
69                    caps.get(i + 1).map_or("", |m| m.as_str()).to_string(),
70                ));
71            }
72            let path = entry.path();
73            if last {
74                if path.is_file() {
75                    found.push((path, parts));
76                    if found.len() > MAX_FILES {
77                        return Err(format!("more than {MAX_FILES} files match"));
78                    }
79                }
80            } else if path.is_dir() {
81                stack.push((path, depth + 1, parts));
82            }
83        }
84    }
85    found.sort_by(|a, b| a.0.cmp(&b.0));
86    Ok(found)
87}
88
89/// Open the tree of `spec`'s files under `dir`.
90pub fn open(spec: &Spec, dir: &Path) -> Result<Opened, String> {
91    let files = spec.files.as_ref().expect("a spec of files");
92    if !dir.is_dir() {
93        return Err(format!(
94            "{} reads a directory of its files ({}), and {} is not one",
95            spec.name,
96            files.pattern,
97            dir.display()
98        ));
99    }
100    let mut one = spec.clone();
101    one.files = None;
102    let mut parts = Vec::new();
103    let mut notes = Vec::new();
104    let mut header = None;
105    let mut sources = Vec::new();
106    let mut schema: Option<SchemaRef> = None;
107    let found = matching(dir, &files.pattern)?;
108    if found.is_empty() {
109        return Err(format!(
110            "no files under {} match {}",
111            dir.display(),
112            files.pattern
113        ));
114    }
115    for (path, texts) in found {
116        let shown = path
117            .strip_prefix(dir)
118            .unwrap_or(&path)
119            .display()
120            .to_string();
121        let mut values = Vec::new();
122        let mut ok = true;
123        for part in &files.parts {
124            let text = texts
125                .iter()
126                .find(|(n, _)| *n == part.name)
127                .map_or("", |(_, t)| t.as_str());
128            match &part.date {
129                Some(format) => match chrono::NaiveDate::parse_from_str(text, format) {
130                    Ok(date) => {
131                        let days = (date
132                            - chrono::NaiveDate::from_ymd_opt(1970, 1, 1).expect("a date"))
133                        .num_days();
134                        values.push(AnyValue::Date(days as i32));
135                    }
136                    Err(_) => {
137                        notes.push(format!(
138                            "{shown}: `{text}` is not a date as {format}; left out"
139                        ));
140                        ok = false;
141                        break;
142                    }
143                },
144                None => values.push(AnyValue::StringOwned(text.into())),
145            }
146        }
147        if !ok {
148            continue;
149        }
150        let opened = match one.open(&path, &shown) {
151            Ok(o) => o,
152            Err(e) => {
153                notes.push(format!("{shown}: {e}; left out"));
154                continue;
155            }
156        };
157        let theirs = opened.records.schema();
158        match &schema {
159            Some(s) if *s != theirs => {
160                notes.push(format!(
161                    "{shown}: its columns differ from the first file's; left out"
162                ));
163                continue;
164            }
165            Some(_) => {}
166            None => schema = Some(theirs),
167        }
168        notes.extend(opened.notes.into_iter().map(|n| format!("{shown}: {n}")));
169        sources.extend(opened.records.sources().iter().cloned());
170        if header.is_none() {
171            header = Some(opened.header);
172        }
173        parts.push((opened.records, values));
174    }
175    let Some(record_schema) = schema else {
176        return Err(format!(
177            "none of the files under {} could be read",
178            dir.display()
179        ));
180    };
181    let part_fields: Vec<(PlSmallStr, DataType)> = files
182        .parts
183        .iter()
184        .map(|p| {
185            (
186                PlSmallStr::from(p.name.as_str()),
187                if p.date.is_some() {
188                    DataType::Date
189                } else {
190                    DataType::String
191                },
192            )
193        })
194        .collect();
195    let mut fields: Vec<Field> = part_fields
196        .iter()
197        .map(|(n, d)| Field::new(n.clone(), d.clone()))
198        .collect();
199    fields.extend(record_schema.iter_fields());
200    let mut starts = Vec::with_capacity(parts.len());
201    let mut rows = 0usize;
202    for (records, _) in &parts {
203        starts.push(rows);
204        rows = rows.saturating_add(records.rows());
205    }
206    let records = SpecFiles {
207        parts,
208        part_fields,
209        starts,
210        rows: rows.min(IdxSize::MAX as usize),
211        schema: Arc::new(Schema::from_iter(fields)),
212        sources,
213    };
214    Ok(Opened {
215        records: Arc::new(records),
216        notes,
217        header: header.unwrap_or_default(),
218    })
219}
220
221impl std::fmt::Debug for SpecFiles {
222    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
223        f.debug_struct("SpecFiles")
224            .field("files", &self.parts.len())
225            .field("rows", &self.rows)
226            .finish()
227    }
228}
229
230impl SpecFiles {
231    /// `lf`, the rows of part `i`, with the part's path values in front.
232    fn dressed(&self, i: usize, lf: LazyFrame) -> LazyFrame {
233        let mut exprs: Vec<Expr> = self
234            .part_fields
235            .iter()
236            .zip(&self.parts[i].1)
237            .map(|((name, dtype), value)| {
238                let value = match value {
239                    AnyValue::Date(d) => lit(*d).cast(DataType::Date),
240                    AnyValue::StringOwned(s) => lit(s.as_str()),
241                    _ => lit(NULL).cast(dtype.clone()),
242                };
243                value.alias(name.clone())
244            })
245            .collect();
246        exprs.push(all().as_expr());
247        lf.select(exprs)
248    }
249}
250
251impl crate::formats::pushdown::Windowed for SpecFiles {
252    fn window(&self, start: usize, len: usize) -> PolarsResult<LazyFrame> {
253        let end = start.saturating_add(len).min(self.rows);
254        let mut frames = Vec::new();
255        let first = self
256            .starts
257            .partition_point(|s| *s <= start)
258            .saturating_sub(1);
259        for i in first..self.parts.len() {
260            let begin = self.starts[i];
261            if begin >= end {
262                break;
263            }
264            let rows = self.parts[i].0.rows();
265            let from = start.saturating_sub(begin).min(rows);
266            let take = (end - begin).min(rows) - from;
267            if take == 0 {
268                continue;
269            }
270            frames.push(self.dressed(i, self.parts[i].0.window(from, take)?));
271        }
272        if frames.is_empty() {
273            return Ok(DataFrame::empty_with_schema(&self.schema).lazy());
274        }
275        concat(frames, UnionArgs::default())
276    }
277}
278
279impl SpecRecords for SpecFiles {
280    fn rows(&self) -> usize {
281        self.rows
282    }
283    fn schema(&self) -> SchemaRef {
284        self.schema.clone()
285    }
286    fn into_lazy(self: Arc<Self>) -> PolarsResult<LazyFrame> {
287        let frames = (0..self.parts.len())
288            .map(|i| Ok(self.dressed(i, self.parts[i].0.clone().into_lazy()?)))
289            .collect::<PolarsResult<Vec<_>>>()?;
290        concat(frames, UnionArgs::default())
291    }
292    fn collect(&self, rows: usize) -> PolarsResult<DataFrame> {
293        crate::formats::pushdown::Windowed::window(self, 0, rows)?.collect()
294    }
295    fn sources(&self) -> &[Arc<Bytes>] {
296        &self.sources
297    }
298}