Skip to main content

data_preprocess/
parquet_store.rs

1//! Parquet-based storage backend for tick and bar data.
2//!
3//! Uses Hive-style directory partitioning:
4//!   {root}/ticks/exchange={ex}/symbol={sym}/{date}.parquet
5//!   {root}/bars/exchange={ex}/symbol={sym}/timeframe={tf}/{date}.parquet
6
7use std::collections::HashMap;
8use std::ffi::OsString;
9use std::fs::{self, File, OpenOptions};
10use std::path::{Path, PathBuf};
11use std::sync::atomic::{AtomicU64, Ordering};
12
13use chrono::NaiveDateTime;
14use polars::prelude::*;
15
16use crate::convert::{
17    bars_to_dataframe, dataframe_to_bars, dataframe_to_price_bars, dataframe_to_stored_ticks,
18    dataframe_to_ticks, ndt_to_date_string, price_bars_to_dataframe, stored_ticks_to_dataframe,
19    ticks_to_dataframe,
20};
21use crate::error::{DataError, Result};
22use crate::models::{
23    Bar, BarQueryOpts, PriceBar, QueryOpts, SeriesDescriptor, StatRow, StoredTick, Tick,
24};
25use crate::scanner::{ParquetScanBounds, ParquetTickScan};
26
27/// Parquet-based storage backend for tick and bar data.
28pub struct ParquetStore {
29    root: PathBuf,
30}
31
32impl ParquetStore {
33    /// Open a Parquet data store rooted at the given directory, creating it if needed.
34    pub fn open(root: impl AsRef<Path>) -> Result<Self> {
35        let root = root.as_ref().to_path_buf();
36        fs::create_dir_all(&root)?;
37        Ok(Self { root })
38    }
39
40    // ── Import ──────────────────────────────────────────────────
41
42    /// Import ticks, deduplicating against existing data per date partition.
43    /// Returns the number of rows actually inserted (after dedup).
44    pub fn insert_ticks(&self, ticks: &[Tick]) -> Result<usize> {
45        if ticks.is_empty() {
46            return Ok(0);
47        }
48
49        // Group ticks by (exchange, symbol, date)
50        let mut groups: HashMap<(String, String, String), Vec<&Tick>> = HashMap::new();
51        for tick in ticks {
52            let date = ndt_to_date_string(&tick.ts);
53            let key = (tick.exchange.clone(), tick.symbol.clone(), date);
54            groups.entry(key).or_default().push(tick);
55        }
56
57        let mut total_inserted = 0usize;
58
59        for ((exchange, symbol, date), group_ticks) in &groups {
60            let dir = self.tick_dir(exchange, symbol);
61            fs::create_dir_all(&dir)?;
62            let file_path = dir.join(format!("{date}.parquet"));
63
64            let owned: Vec<Tick> = group_ticks.iter().map(|t| (*t).clone()).collect();
65            let new_df = ticks_to_dataframe(&owned)?;
66
67            if file_path.exists() {
68                let existing_df = read_parquet_file(&file_path)?;
69                let existing_count = existing_df.height();
70                let combined = concat_and_dedup_ticks(existing_df, new_df)?;
71                total_inserted += combined.height().saturating_sub(existing_count);
72                write_parquet_file(&file_path, &mut combined.clone())?;
73            } else {
74                let deduped = dedup_ticks(new_df)?;
75                total_inserted += deduped.height();
76                write_parquet_file(&file_path, &mut deduped.clone())?;
77            }
78        }
79
80        Ok(total_inserted)
81    }
82
83    /// Persist enhanced ticks without timestamp deduplication.
84    pub fn insert_stored_ticks(&self, ticks: &[StoredTick]) -> Result<usize> {
85        let mut groups: HashMap<(String, String, String), Vec<StoredTick>> = HashMap::new();
86        for tick in ticks {
87            tick.validate()?;
88            groups
89                .entry((
90                    tick.tick.exchange.clone(),
91                    tick.tick.symbol.clone(),
92                    ndt_to_date_string(&tick.tick.ts),
93                ))
94                .or_default()
95                .push(tick.clone());
96        }
97        let mut inserted = 0usize;
98        for ((exchange, symbol, date), mut incoming) in groups {
99            let dir = self.enhanced_tick_dir(&exchange, &symbol);
100            fs::create_dir_all(&dir)?;
101            let path = dir.join(format!("{date}.parquet"));
102            let mut rows = if path.exists() {
103                dataframe_to_stored_ticks(&read_parquet_file(&path)?)?
104            } else {
105                Vec::new()
106            };
107            let before = rows.len();
108            merge_stored_ticks(&mut rows, &mut incoming)?;
109            rows.sort_by_key(|row| (row.tick.ts, row.source_ordinal));
110            let mut frame = stored_ticks_to_dataframe(&rows)?;
111            write_parquet_file(&path, &mut frame)?;
112            inserted = inserted
113                .checked_add(rows.len() - before)
114                .ok_or_else(|| DataError::Other("insert count overflowed".into()))?;
115        }
116        Ok(inserted)
117    }
118
119    pub fn query_stored_ticks(&self, exchange: &str, symbol: &str) -> Result<Vec<StoredTick>> {
120        self.query_stored_ticks_bounded(exchange, symbol, usize::MAX, usize::MAX, || false)
121    }
122
123    pub fn query_stored_ticks_bounded<F>(
124        &self,
125        exchange: &str,
126        symbol: &str,
127        max_rows: usize,
128        max_resident_bytes: usize,
129        mut is_cancelled: F,
130    ) -> Result<Vec<StoredTick>>
131    where
132        F: FnMut() -> bool,
133    {
134        if max_rows == 0 || max_resident_bytes == 0 {
135            return Err(DataError::Other(
136                "enhanced tick query limits must be positive".into(),
137            ));
138        }
139        ensure_not_cancelled(&mut is_cancelled)?;
140        let dir = self.enhanced_tick_dir(exchange, symbol);
141        if !dir.exists() {
142            return Ok(Vec::new());
143        }
144        let mut rows = Vec::new();
145        for path in list_date_files_cancellable(&dir, None, None, &mut is_cancelled)? {
146            ensure_not_cancelled(&mut is_cancelled)?;
147            let frame = read_parquet_file(&path)?;
148            let admitted = rows
149                .len()
150                .checked_add(frame.height())
151                .ok_or_else(|| DataError::Other("enhanced tick row count overflowed".into()))?;
152            check_enhanced_query_bound::<StoredTick>(admitted, max_rows, max_resident_bytes)?;
153            let partition = dataframe_to_stored_ticks(&frame)?;
154            let bytes = stored_tick_bytes(&rows)?
155                .checked_add(stored_tick_bytes(&partition)?)
156                .ok_or_else(|| DataError::Other("enhanced tick byte count overflowed".into()))?;
157            if bytes > max_resident_bytes {
158                return Err(DataError::Other(
159                    "enhanced tick query exceeds resident-byte limit".into(),
160                ));
161            }
162            rows.extend(partition);
163            ensure_not_cancelled(&mut is_cancelled)?;
164        }
165        rows.sort_by_key(|row| (row.tick.ts, row.source_ordinal));
166        Ok(rows)
167    }
168
169    /// Import bars, deduplicating against existing data per date partition.
170    /// Returns the number of rows actually inserted (after dedup).
171    pub fn insert_bars(&self, bars: &[Bar]) -> Result<usize> {
172        if bars.is_empty() {
173            return Ok(0);
174        }
175
176        // Group bars by (exchange, symbol, timeframe, date)
177        let mut groups: HashMap<(String, String, String, String), Vec<&Bar>> = HashMap::new();
178        for bar in bars {
179            let date = ndt_to_date_string(&bar.ts);
180            let key = (
181                bar.exchange.clone(),
182                bar.symbol.clone(),
183                bar.timeframe.as_str().to_string(),
184                date,
185            );
186            groups.entry(key).or_default().push(bar);
187        }
188
189        let mut total_inserted = 0usize;
190
191        for ((exchange, symbol, timeframe, date), group_bars) in &groups {
192            let dir = self.bar_dir(exchange, symbol, timeframe);
193            fs::create_dir_all(&dir)?;
194            let file_path = dir.join(format!("{date}.parquet"));
195
196            let owned: Vec<Bar> = group_bars.iter().map(|b| (*b).clone()).collect();
197            let new_df = bars_to_dataframe(&owned)?;
198
199            if file_path.exists() {
200                let existing_df = read_parquet_file(&file_path)?;
201                let existing_count = existing_df.height();
202                let combined = concat_and_dedup_bars(existing_df, new_df)?;
203                total_inserted += combined.height().saturating_sub(existing_count);
204                write_parquet_file(&file_path, &mut combined.clone())?;
205            } else {
206                let deduped = dedup_bars(new_df)?;
207                total_inserted += deduped.height();
208                write_parquet_file(&file_path, &mut deduped.clone())?;
209            }
210        }
211
212        Ok(total_inserted)
213    }
214
215    pub fn insert_price_bars(
216        &self,
217        descriptor: &SeriesDescriptor,
218        bars: &[PriceBar],
219    ) -> Result<usize> {
220        descriptor.validate()?;
221        if !descriptor.verified {
222            return Err(DataError::Other(
223                "new price-only writes require a verified descriptor".into(),
224            ));
225        }
226        let dir = self.price_bar_dir(
227            &descriptor.exchange,
228            &descriptor.symbol,
229            descriptor.timeframe_seconds,
230        );
231        fs::create_dir_all(&dir)?;
232        persist_descriptor(&dir, descriptor)?;
233        let mut groups: HashMap<String, Vec<PriceBar>> = HashMap::new();
234        for bar in bars {
235            bar.validate()?;
236            groups
237                .entry(ndt_to_date_string(&bar.ts))
238                .or_default()
239                .push(bar.clone());
240        }
241        let mut inserted = 0usize;
242        for (date, mut incoming) in groups {
243            let path = dir.join(format!("{date}.parquet"));
244            let mut rows = if path.exists() {
245                dataframe_to_price_bars(&read_parquet_file(&path)?)?
246            } else {
247                Vec::new()
248            };
249            let before = rows.len();
250            for bar in incoming.drain(..) {
251                if let Some(existing) = rows.iter().find(|value| value.ts == bar.ts) {
252                    if existing != &bar {
253                        return Err(DataError::Other(
254                            "conflicting price-only bar at one bucket".into(),
255                        ));
256                    }
257                } else {
258                    rows.push(bar)
259                }
260            }
261            rows.sort_by_key(|bar| (bar.available_at, bar.ts));
262            let mut frame = price_bars_to_dataframe(&rows)?;
263            write_parquet_file(&path, &mut frame)?;
264            inserted += rows.len() - before;
265        }
266        Ok(inserted)
267    }
268
269    pub fn query_price_bars(&self, descriptor: &SeriesDescriptor) -> Result<Vec<PriceBar>> {
270        self.query_price_bars_bounded(descriptor, usize::MAX, usize::MAX, || false)
271    }
272
273    pub fn query_price_bars_bounded<F>(
274        &self,
275        descriptor: &SeriesDescriptor,
276        max_rows: usize,
277        max_resident_bytes: usize,
278        mut is_cancelled: F,
279    ) -> Result<Vec<PriceBar>>
280    where
281        F: FnMut() -> bool,
282    {
283        descriptor.validate()?;
284        if max_rows == 0 || max_resident_bytes == 0 {
285            return Err(DataError::Other(
286                "price-bar query limits must be positive".into(),
287            ));
288        }
289        ensure_not_cancelled(&mut is_cancelled)?;
290        let dir = self.price_bar_dir(
291            &descriptor.exchange,
292            &descriptor.symbol,
293            descriptor.timeframe_seconds,
294        );
295        if !dir.exists() {
296            return Ok(Vec::new());
297        }
298        verify_descriptor(&dir, descriptor)?;
299        let mut rows = Vec::new();
300        for path in list_date_files_cancellable(&dir, None, None, &mut is_cancelled)? {
301            ensure_not_cancelled(&mut is_cancelled)?;
302            let frame = read_parquet_file(&path)?;
303            let admitted = rows
304                .len()
305                .checked_add(frame.height())
306                .ok_or_else(|| DataError::Other("price-bar row count overflowed".into()))?;
307            check_enhanced_query_bound::<PriceBar>(admitted, max_rows, max_resident_bytes)?;
308            let partition = dataframe_to_price_bars(&frame)?;
309            let bytes = price_bar_bytes(&rows)?
310                .checked_add(price_bar_bytes(&partition)?)
311                .ok_or_else(|| DataError::Other("price-bar byte count overflowed".into()))?;
312            if bytes > max_resident_bytes {
313                return Err(DataError::Other(
314                    "price-bar query exceeds resident-byte limit".into(),
315                ));
316            }
317            rows.extend(partition);
318            ensure_not_cancelled(&mut is_cancelled)?;
319        }
320        rows.sort_by_key(|bar| bar.ts);
321        Ok(rows)
322    }
323
324    // ── Query ───────────────────────────────────────────────────
325
326    /// Query ticks for a given exchange+symbol, optionally filtered by date range.
327    /// Returns (ticks, total_count_matching_filters).
328    pub fn query_ticks(&self, opts: &QueryOpts) -> Result<(Vec<Tick>, u64)> {
329        self.query_ticks_cancellable(opts, || false)
330    }
331
332    /// Cancellable tick query.
333    ///
334    /// Cancellation is checked while traversing directory entries, before and
335    /// after every Parquet file, and between DataFrame processing stages. The
336    /// low-level Polars collect/read for one file remains atomic because Polars
337    /// does not expose an interruption hook for that operation.
338    pub fn query_ticks_cancellable<F>(
339        &self,
340        opts: &QueryOpts,
341        mut is_cancelled: F,
342    ) -> Result<(Vec<Tick>, u64)>
343    where
344        F: FnMut() -> bool,
345    {
346        ensure_not_cancelled(&mut is_cancelled)?;
347        let dir = self.tick_dir(&opts.exchange, &opts.symbol);
348        if !dir.exists() {
349            return Ok((Vec::new(), 0));
350        }
351
352        let files = list_date_files_cancellable(&dir, opts.from, opts.to, &mut is_cancelled)?;
353        if files.is_empty() {
354            return Ok((Vec::new(), 0));
355        }
356
357        let mut all_dfs: Vec<DataFrame> = Vec::with_capacity(files.len());
358        for file in &files {
359            ensure_not_cancelled(&mut is_cancelled)?;
360            let df = read_parquet_file(file)?;
361            ensure_not_cancelled(&mut is_cancelled)?;
362            all_dfs.push(df);
363        }
364        let mut combined = concat_dataframes(all_dfs)?;
365        ensure_not_cancelled(&mut is_cancelled)?;
366
367        combined = apply_ts_filter(combined, opts.from, opts.to)?;
368        ensure_not_cancelled(&mut is_cancelled)?;
369        combined = combined.sort(["ts"], SortMultipleOptions::default())?;
370        ensure_not_cancelled(&mut is_cancelled)?;
371
372        let total = combined.height() as u64;
373        combined = apply_pagination(combined, opts.limit, opts.tail, opts.descending)?;
374        ensure_not_cancelled(&mut is_cancelled)?;
375
376        let ticks = dataframe_to_ticks(&combined)?;
377        ensure_not_cancelled(&mut is_cancelled)?;
378        Ok((ticks, total))
379    }
380
381    /// Return the latest tick with a valid quote strictly before `before`.
382    pub fn latest_valid_tick_before(
383        &self,
384        exchange: &str,
385        symbol: &str,
386        before: NaiveDateTime,
387    ) -> Result<Option<Tick>> {
388        self.latest_valid_tick_before_cancellable(exchange, symbol, before, || false)
389    }
390
391    /// Cancellable strict-before lookup for the latest tick with a valid quote.
392    ///
393    /// Date partitions and their rows are searched newest-to-oldest.
394    /// Ticks at `before` are excluded, and ticks with missing, non-finite, non-positive, or crossed bid/ask prices are skipped.
395    pub fn latest_valid_tick_before_cancellable<F>(
396        &self,
397        exchange: &str,
398        symbol: &str,
399        before: NaiveDateTime,
400        mut is_cancelled: F,
401    ) -> Result<Option<Tick>>
402    where
403        F: FnMut() -> bool,
404    {
405        ensure_not_cancelled(&mut is_cancelled)?;
406        let scan = ParquetTickScan::describe_cancellable(
407            &self.root,
408            exchange,
409            symbol,
410            ParquetScanBounds::new(None, Some(before)),
411            &mut is_cancelled,
412        )?;
413        let latest = scan
414            .latest_valid_tick_before_cancellable(before, &mut is_cancelled)?
415            .map(|row| row.row);
416        ensure_not_cancelled(&mut is_cancelled)?;
417        Ok(latest)
418    }
419
420    /// Query bars for a given exchange+symbol+timeframe, optionally filtered by date range.
421    /// Returns (bars, total_count_matching_filters).
422    pub fn query_bars(&self, opts: &BarQueryOpts) -> Result<(Vec<Bar>, u64)> {
423        self.query_bars_cancellable(opts, || false)
424    }
425
426    /// Cancellable bar query with the same cooperative boundaries as
427    /// [`Self::query_ticks_cancellable`].
428    pub fn query_bars_cancellable<F>(
429        &self,
430        opts: &BarQueryOpts,
431        mut is_cancelled: F,
432    ) -> Result<(Vec<Bar>, u64)>
433    where
434        F: FnMut() -> bool,
435    {
436        ensure_not_cancelled(&mut is_cancelled)?;
437        let dir = self.bar_dir(&opts.exchange, &opts.symbol, &opts.timeframe);
438        if !dir.exists() {
439            return Ok((Vec::new(), 0));
440        }
441
442        let files = list_date_files_cancellable(&dir, opts.from, opts.to, &mut is_cancelled)?;
443        if files.is_empty() {
444            return Ok((Vec::new(), 0));
445        }
446
447        let mut all_dfs: Vec<DataFrame> = Vec::with_capacity(files.len());
448        for file in &files {
449            ensure_not_cancelled(&mut is_cancelled)?;
450            let df = read_parquet_file(file)?;
451            ensure_not_cancelled(&mut is_cancelled)?;
452            all_dfs.push(df);
453        }
454        let mut combined = concat_dataframes(all_dfs)?;
455        ensure_not_cancelled(&mut is_cancelled)?;
456
457        combined = apply_ts_filter(combined, opts.from, opts.to)?;
458        ensure_not_cancelled(&mut is_cancelled)?;
459        combined = combined.sort(["ts"], SortMultipleOptions::default())?;
460        ensure_not_cancelled(&mut is_cancelled)?;
461
462        let total = combined.height() as u64;
463        combined = apply_pagination(combined, opts.limit, opts.tail, opts.descending)?;
464        ensure_not_cancelled(&mut is_cancelled)?;
465
466        let bars = dataframe_to_bars(&combined)?;
467        ensure_not_cancelled(&mut is_cancelled)?;
468        Ok((bars, total))
469    }
470
471    // ── Delete ──────────────────────────────────────────────────
472
473    /// Delete ticks matching exchange+symbol, optionally within a date range.
474    pub fn delete_ticks(
475        &self,
476        exchange: &str,
477        symbol: &str,
478        from: Option<NaiveDateTime>,
479        to: Option<NaiveDateTime>,
480    ) -> Result<usize> {
481        let dir = self.tick_dir(exchange, symbol);
482        if !dir.exists() {
483            return Ok(0);
484        }
485        delete_from_partition(&dir, from, to)
486    }
487
488    /// Delete bars matching exchange+symbol+timeframe, optionally within a date range.
489    pub fn delete_bars(
490        &self,
491        exchange: &str,
492        symbol: &str,
493        timeframe: &str,
494        from: Option<NaiveDateTime>,
495        to: Option<NaiveDateTime>,
496    ) -> Result<usize> {
497        let dir = self.bar_dir(exchange, symbol, timeframe);
498        if !dir.exists() {
499            return Ok(0);
500        }
501        delete_from_partition(&dir, from, to)
502    }
503
504    /// Delete ALL data (ticks + bars) for an exchange+symbol pair.
505    pub fn delete_symbol(&self, exchange: &str, symbol: &str) -> Result<(usize, usize)> {
506        let tick_count = self.count_rows_in_dir(&self.tick_dir(exchange, symbol));
507        let bar_count = self.count_all_bars_for_symbol(exchange, symbol);
508
509        // Remove tick directory
510        let tick_dir = self.tick_dir(exchange, symbol);
511        if tick_dir.exists() {
512            fs::remove_dir_all(&tick_dir)?;
513        }
514
515        // Remove bar directories for all timeframes
516        let bar_sym_dir = self
517            .root
518            .join("bars")
519            .join(format!("exchange={exchange}"))
520            .join(format!("symbol={symbol}"));
521        if bar_sym_dir.exists() {
522            fs::remove_dir_all(&bar_sym_dir)?;
523        }
524
525        Ok((tick_count, bar_count))
526    }
527
528    /// Delete ALL data for an entire exchange.
529    pub fn delete_exchange(&self, exchange: &str) -> Result<(usize, usize)> {
530        let tick_ex_dir = self.root.join("ticks").join(format!("exchange={exchange}"));
531        let bar_ex_dir = self.root.join("bars").join(format!("exchange={exchange}"));
532
533        let tick_count = self.count_rows_recursive(&tick_ex_dir);
534        let bar_count = self.count_rows_recursive(&bar_ex_dir);
535
536        if tick_ex_dir.exists() {
537            fs::remove_dir_all(&tick_ex_dir)?;
538        }
539        if bar_ex_dir.exists() {
540            fs::remove_dir_all(&bar_ex_dir)?;
541        }
542
543        Ok((tick_count, bar_count))
544    }
545
546    // ── Stats ───────────────────────────────────────────────────
547
548    /// Summary statistics across all data, optionally filtered by exchange and/or symbol.
549    pub fn stats(&self, exchange: Option<&str>, symbol: Option<&str>) -> Result<Vec<StatRow>> {
550        let mut rows = Vec::new();
551
552        // Collect tick stats
553        self.collect_tick_stats(&mut rows, exchange, symbol)?;
554
555        // Collect bar stats
556        self.collect_bar_stats(&mut rows, exchange, symbol)?;
557
558        // Sort by exchange, symbol, data_type
559        rows.sort_by(|a, b| {
560            a.exchange
561                .cmp(&b.exchange)
562                .then(a.symbol.cmp(&b.symbol))
563                .then(a.data_type.cmp(&b.data_type))
564        });
565
566        Ok(rows)
567    }
568
569    /// Total size of all Parquet files under the data root (bytes).
570    pub fn total_size(&self) -> Option<u64> {
571        let mut total = 0u64;
572        for entry in walkdir(&self.root) {
573            if entry.extension().is_some_and(|e| e == "parquet")
574                && let Ok(meta) = fs::metadata(&entry)
575            {
576                total += meta.len();
577            }
578        }
579        if total == 0 { None } else { Some(total) }
580    }
581
582    // ── Private helpers ─────────────────────────────────────────
583
584    pub(crate) fn root_path(&self) -> &Path {
585        &self.root
586    }
587
588    pub(crate) fn verify_price_bar_series(&self, descriptor: &SeriesDescriptor) -> Result<()> {
589        descriptor.validate()?;
590        let directory = self.price_bar_dir(
591            &descriptor.exchange,
592            &descriptor.symbol,
593            descriptor.timeframe_seconds,
594        );
595        if directory.exists() {
596            verify_descriptor(&directory, descriptor)?;
597        }
598        Ok(())
599    }
600
601    /// Build tick directory path for a given exchange+symbol.
602    fn tick_dir(&self, exchange: &str, symbol: &str) -> PathBuf {
603        self.root
604            .join("ticks")
605            .join(format!("exchange={exchange}"))
606            .join(format!("symbol={symbol}"))
607    }
608
609    fn enhanced_tick_dir(&self, exchange: &str, symbol: &str) -> PathBuf {
610        self.root
611            .join("ordered_ticks")
612            .join(format!("exchange={exchange}"))
613            .join(format!("symbol={symbol}"))
614    }
615
616    fn price_bar_dir(&self, exchange: &str, symbol: &str, timeframe_seconds: u64) -> PathBuf {
617        self.root
618            .join("price_bars")
619            .join(format!("exchange={exchange}"))
620            .join(format!("symbol={symbol}"))
621            .join(format!("timeframe_seconds={timeframe_seconds}"))
622    }
623
624    /// Build bar directory path for a given exchange+symbol+timeframe.
625    fn bar_dir(&self, exchange: &str, symbol: &str, timeframe: &str) -> PathBuf {
626        self.root
627            .join("bars")
628            .join(format!("exchange={exchange}"))
629            .join(format!("symbol={symbol}"))
630            .join(format!("timeframe={timeframe}"))
631    }
632
633    /// Count total rows across all parquet files in a directory.
634    fn count_rows_in_dir(&self, dir: &Path) -> usize {
635        if !dir.exists() {
636            return 0;
637        }
638        let mut count = 0;
639        if let Ok(entries) = fs::read_dir(dir) {
640            for entry in entries.flatten() {
641                let path = entry.path();
642                if path.extension().is_some_and(|e| e == "parquet")
643                    && let Ok(df) = read_parquet_file(&path)
644                {
645                    count += df.height();
646                }
647            }
648        }
649        count
650    }
651
652    /// Count total rows recursively across all parquet files under a directory.
653    fn count_rows_recursive(&self, dir: &Path) -> usize {
654        if !dir.exists() {
655            return 0;
656        }
657        let mut count = 0;
658        for path in walkdir(dir) {
659            if path.extension().is_some_and(|e| e == "parquet")
660                && let Ok(df) = read_parquet_file(&path)
661            {
662                count += df.height();
663            }
664        }
665        count
666    }
667
668    /// Count all bar rows for a given exchange+symbol across all timeframes.
669    fn count_all_bars_for_symbol(&self, exchange: &str, symbol: &str) -> usize {
670        let bar_sym_dir = self
671            .root
672            .join("bars")
673            .join(format!("exchange={exchange}"))
674            .join(format!("symbol={symbol}"));
675        self.count_rows_recursive(&bar_sym_dir)
676    }
677
678    /// Collect tick stats from the directory tree.
679    fn collect_tick_stats(
680        &self,
681        rows: &mut Vec<StatRow>,
682        exchange_filter: Option<&str>,
683        symbol_filter: Option<&str>,
684    ) -> Result<()> {
685        let ticks_dir = self.root.join("ticks");
686        if !ticks_dir.exists() {
687            return Ok(());
688        }
689
690        for (exchange, symbol, dir) in self.iter_exchange_symbol_dirs(&ticks_dir)? {
691            if let Some(ef) = exchange_filter
692                && exchange != ef
693            {
694                continue;
695            }
696            if let Some(sf) = symbol_filter
697                && symbol != sf
698            {
699                continue;
700            }
701
702            let (count, ts_min, ts_max) = self.aggregate_parquet_stats(&dir)?;
703            if count > 0 {
704                rows.push(StatRow {
705                    exchange,
706                    symbol,
707                    data_type: "tick".to_string(),
708                    count,
709                    ts_min: ts_min.unwrap_or_default(),
710                    ts_max: ts_max.unwrap_or_default(),
711                });
712            }
713        }
714
715        Ok(())
716    }
717
718    /// Collect bar stats from the directory tree.
719    fn collect_bar_stats(
720        &self,
721        rows: &mut Vec<StatRow>,
722        exchange_filter: Option<&str>,
723        symbol_filter: Option<&str>,
724    ) -> Result<()> {
725        let bars_dir = self.root.join("bars");
726        if !bars_dir.exists() {
727            return Ok(());
728        }
729
730        for (exchange, symbol, timeframe, dir) in self.iter_exchange_symbol_tf_dirs(&bars_dir)? {
731            if let Some(ef) = exchange_filter
732                && exchange != ef
733            {
734                continue;
735            }
736            if let Some(sf) = symbol_filter
737                && symbol != sf
738            {
739                continue;
740            }
741
742            let (count, ts_min, ts_max) = self.aggregate_parquet_stats(&dir)?;
743            if count > 0 {
744                rows.push(StatRow {
745                    exchange,
746                    symbol,
747                    data_type: format!("bar ({timeframe})"),
748                    count,
749                    ts_min: ts_min.unwrap_or_default(),
750                    ts_max: ts_max.unwrap_or_default(),
751                });
752            }
753        }
754
755        Ok(())
756    }
757
758    /// Iterate over exchange/symbol directories under a top-level dir.
759    fn iter_exchange_symbol_dirs(&self, base: &Path) -> Result<Vec<(String, String, PathBuf)>> {
760        let mut result = Vec::new();
761        if !base.exists() {
762            return Ok(result);
763        }
764
765        for ex_entry in fs::read_dir(base)?.flatten() {
766            let ex_path = ex_entry.path();
767            if !ex_path.is_dir() {
768                continue;
769            }
770            let exchange =
771                parse_partition_value(ex_path.file_name().unwrap().to_str().unwrap_or(""));
772            if exchange.is_empty() {
773                continue;
774            }
775
776            for sym_entry in fs::read_dir(&ex_path)?.flatten() {
777                let sym_path = sym_entry.path();
778                if !sym_path.is_dir() {
779                    continue;
780                }
781                let symbol =
782                    parse_partition_value(sym_path.file_name().unwrap().to_str().unwrap_or(""));
783                if symbol.is_empty() {
784                    continue;
785                }
786                result.push((exchange.clone(), symbol, sym_path));
787            }
788        }
789
790        Ok(result)
791    }
792
793    /// Iterate over exchange/symbol/timeframe directories under a top-level dir.
794    fn iter_exchange_symbol_tf_dirs(
795        &self,
796        base: &Path,
797    ) -> Result<Vec<(String, String, String, PathBuf)>> {
798        let mut result = Vec::new();
799        if !base.exists() {
800            return Ok(result);
801        }
802
803        for ex_entry in fs::read_dir(base)?.flatten() {
804            let ex_path = ex_entry.path();
805            if !ex_path.is_dir() {
806                continue;
807            }
808            let exchange =
809                parse_partition_value(ex_path.file_name().unwrap().to_str().unwrap_or(""));
810            if exchange.is_empty() {
811                continue;
812            }
813
814            for sym_entry in fs::read_dir(&ex_path)?.flatten() {
815                let sym_path = sym_entry.path();
816                if !sym_path.is_dir() {
817                    continue;
818                }
819                let symbol =
820                    parse_partition_value(sym_path.file_name().unwrap().to_str().unwrap_or(""));
821                if symbol.is_empty() {
822                    continue;
823                }
824
825                for tf_entry in fs::read_dir(&sym_path)?.flatten() {
826                    let tf_path = tf_entry.path();
827                    if !tf_path.is_dir() {
828                        continue;
829                    }
830                    let timeframe =
831                        parse_partition_value(tf_path.file_name().unwrap().to_str().unwrap_or(""));
832                    if timeframe.is_empty() {
833                        continue;
834                    }
835                    result.push((exchange.clone(), symbol.clone(), timeframe, tf_path));
836                }
837            }
838        }
839
840        Ok(result)
841    }
842
843    /// Read all parquet files in a directory and aggregate row count + min/max ts.
844    fn aggregate_parquet_stats(
845        &self,
846        dir: &Path,
847    ) -> Result<(u64, Option<NaiveDateTime>, Option<NaiveDateTime>)> {
848        let mut total_count = 0u64;
849        let mut global_min: Option<i64> = None;
850        let mut global_max: Option<i64> = None;
851
852        if !dir.exists() {
853            return Ok((0, None, None));
854        }
855
856        for entry in fs::read_dir(dir)?.flatten() {
857            let path = entry.path();
858            if path.extension().is_some_and(|e| e == "parquet") {
859                let df = read_parquet_file(&path)?;
860                total_count += df.height() as u64;
861
862                if df.height() > 0 {
863                    let ts_col = df.column("ts").ok().and_then(|c| c.datetime().ok());
864                    if let Some(ts) = ts_col {
865                        if let Some(min_val) = ts.min() {
866                            global_min =
867                                Some(global_min.map_or(min_val, |cur: i64| cur.min(min_val)));
868                        }
869                        if let Some(max_val) = ts.max() {
870                            global_max =
871                                Some(global_max.map_or(max_val, |cur: i64| cur.max(max_val)));
872                        }
873                    }
874                }
875            }
876        }
877
878        let ts_min = global_min.map(micros_to_ndt);
879        let ts_max = global_max.map(micros_to_ndt);
880
881        Ok((total_count, ts_min, ts_max))
882    }
883}
884
885// ── Free functions ──────────────────────────────────────────────
886
887/// Parse a Hive partition value from a directory name like "exchange=ctrader".
888fn parse_partition_value(dir_name: &str) -> String {
889    dir_name
890        .split_once('=')
891        .map(|(_, v)| v.to_string())
892        .unwrap_or_default()
893}
894
895fn check_enhanced_query_bound<T>(
896    rows: usize,
897    max_rows: usize,
898    max_resident_bytes: usize,
899) -> Result<()> {
900    if rows > max_rows
901        || rows
902            .checked_mul(std::mem::size_of::<T>())
903            .is_none_or(|bytes| bytes > max_resident_bytes)
904    {
905        Err(DataError::Other(
906            "enhanced query exceeds its row or resident-byte limit".into(),
907        ))
908    } else {
909        Ok(())
910    }
911}
912
913fn stored_tick_bytes(rows: &[StoredTick]) -> Result<usize> {
914    rows.iter().try_fold(0usize, |bytes, row| {
915        bytes
916            .checked_add(std::mem::size_of::<StoredTick>())
917            .and_then(|value| value.checked_add(row.tick.exchange.len()))
918            .and_then(|value| value.checked_add(row.tick.symbol.len()))
919            .and_then(|value| {
920                value.checked_add(row.source_identity.as_ref().map_or(0, String::len))
921            })
922            .ok_or_else(|| DataError::Other("enhanced tick byte count overflowed".into()))
923    })
924}
925
926fn price_bar_bytes(rows: &[PriceBar]) -> Result<usize> {
927    rows.iter().try_fold(0usize, |bytes, row| {
928        bytes
929            .checked_add(std::mem::size_of::<PriceBar>())
930            .and_then(|value| value.checked_add(row.exchange.len()))
931            .and_then(|value| value.checked_add(row.symbol.len()))
932            .ok_or_else(|| DataError::Other("price-bar byte count overflowed".into()))
933    })
934}
935
936fn ensure_not_cancelled(is_cancelled: &mut dyn FnMut() -> bool) -> Result<()> {
937    if is_cancelled() {
938        Err(DataError::Cancelled)
939    } else {
940        Ok(())
941    }
942}
943
944/// List parquet files in a directory, optionally filtered by date range in filename.
945fn list_date_files(
946    dir: &Path,
947    from: Option<NaiveDateTime>,
948    to: Option<NaiveDateTime>,
949) -> Result<Vec<PathBuf>> {
950    list_date_files_cancellable(dir, from, to, &mut || false)
951}
952
953fn list_date_files_cancellable(
954    dir: &Path,
955    from: Option<NaiveDateTime>,
956    to: Option<NaiveDateTime>,
957    is_cancelled: &mut dyn FnMut() -> bool,
958) -> Result<Vec<PathBuf>> {
959    ensure_not_cancelled(is_cancelled)?;
960    let mut files = Vec::new();
961    let from_date = from.map(|d| d.format("%Y-%m-%d").to_string());
962    let to_date = to.map(|d| d.format("%Y-%m-%d").to_string());
963
964    for entry in fs::read_dir(dir)? {
965        ensure_not_cancelled(is_cancelled)?;
966        let path = entry?.path();
967        if path.extension().is_some_and(|e| e == "parquet") {
968            let stem = path.file_stem().and_then(|s| s.to_str()).unwrap_or("");
969
970            // Filename-level date pruning
971            let dominated_by_from = from_date.as_ref().is_some_and(|fd| stem < fd.as_str());
972            let past_to = to_date.as_ref().is_some_and(|td| stem > td.as_str());
973
974            if !dominated_by_from && !past_to {
975                files.push(path);
976            }
977        }
978    }
979
980    ensure_not_cancelled(is_cancelled)?;
981    files.sort();
982    Ok(files)
983}
984
985/// Read a single Parquet file into a DataFrame.
986fn read_parquet_file(path: &Path) -> Result<DataFrame> {
987    let file = std::fs::File::open(path)?;
988    let df = ParquetReader::new(file).finish()?;
989    Ok(df)
990}
991
992static NEXT_TEMP_FILE_ID: AtomicU64 = AtomicU64::new(0);
993
994/// Write a DataFrame to a temporary file and atomically replace the partition.
995fn write_parquet_file(path: &Path, df: &mut DataFrame) -> Result<()> {
996    let (temp_path, mut file) = create_partition_temp_file(path)?;
997    let write_result = (|| -> Result<()> {
998        ParquetWriter::new(&mut file)
999            .with_compression(ParquetCompression::Zstd(None))
1000            .finish(df)?;
1001        file.sync_all()?;
1002        Ok(())
1003    })();
1004    drop(file);
1005
1006    if let Err(error) = write_result {
1007        fs::remove_file(&temp_path).ok();
1008        return Err(error);
1009    }
1010    if let Err(error) = atomic_replace(&temp_path, path) {
1011        fs::remove_file(&temp_path).ok();
1012        return Err(error.into());
1013    }
1014    Ok(())
1015}
1016
1017fn create_partition_temp_file(path: &Path) -> Result<(PathBuf, File)> {
1018    let parent = path.parent().ok_or_else(|| {
1019        DataError::Other(format!("partition path has no parent: {}", path.display()))
1020    })?;
1021    let file_name = path.file_name().ok_or_else(|| {
1022        DataError::Other(format!(
1023            "partition path has no file name: {}",
1024            path.display()
1025        ))
1026    })?;
1027
1028    loop {
1029        let id = NEXT_TEMP_FILE_ID.fetch_add(1, Ordering::Relaxed);
1030        let mut temp_name = OsString::from(".");
1031        temp_name.push(file_name);
1032        temp_name.push(format!(".{}.{}.tmp", std::process::id(), id));
1033        let temp_path = parent.join(temp_name);
1034        match OpenOptions::new()
1035            .write(true)
1036            .create_new(true)
1037            .open(&temp_path)
1038        {
1039            Ok(file) => return Ok((temp_path, file)),
1040            Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => continue,
1041            Err(error) => return Err(error.into()),
1042        }
1043    }
1044}
1045
1046#[cfg(not(windows))]
1047fn atomic_replace(from: &Path, to: &Path) -> std::io::Result<()> {
1048    fs::rename(from, to)
1049}
1050
1051#[cfg(windows)]
1052fn atomic_replace(from: &Path, to: &Path) -> std::io::Result<()> {
1053    use std::os::windows::ffi::OsStrExt;
1054
1055    const MOVEFILE_REPLACE_EXISTING: u32 = 0x1;
1056    const MOVEFILE_WRITE_THROUGH: u32 = 0x8;
1057
1058    #[link(name = "kernel32")]
1059    unsafe extern "system" {
1060        fn MoveFileExW(
1061            existing_file_name: *const u16,
1062            new_file_name: *const u16,
1063            flags: u32,
1064        ) -> i32;
1065    }
1066
1067    let from = from
1068        .as_os_str()
1069        .encode_wide()
1070        .chain(Some(0))
1071        .collect::<Vec<_>>();
1072    let to = to
1073        .as_os_str()
1074        .encode_wide()
1075        .chain(Some(0))
1076        .collect::<Vec<_>>();
1077    let replaced = unsafe {
1078        MoveFileExW(
1079            from.as_ptr(),
1080            to.as_ptr(),
1081            MOVEFILE_REPLACE_EXISTING | MOVEFILE_WRITE_THROUGH,
1082        )
1083    };
1084    if replaced == 0 {
1085        Err(std::io::Error::last_os_error())
1086    } else {
1087        Ok(())
1088    }
1089}
1090
1091fn merge_stored_ticks(
1092    existing: &mut Vec<StoredTick>,
1093    incoming: &mut Vec<StoredTick>,
1094) -> Result<()> {
1095    let mut ordinals = existing
1096        .iter()
1097        .map(|row| row.source_ordinal)
1098        .collect::<std::collections::BTreeSet<_>>();
1099    for row in incoming.drain(..) {
1100        if !ordinals.insert(row.source_ordinal) {
1101            return Err(DataError::Other(
1102                "duplicate persisted source ordinal".into(),
1103            ));
1104        }
1105        if let (Some(source), Some(sequence)) = (&row.source_identity, row.provider_sequence)
1106            && let Some(previous) = existing.iter().find(|value| {
1107                value.source_identity.as_ref() == Some(source)
1108                    && value.provider_sequence == Some(sequence)
1109            })
1110        {
1111            if previous.tick == row.tick {
1112                continue;
1113            }
1114            return Err(DataError::Other(
1115                "conflicting payload for provider sequence".into(),
1116            ));
1117        }
1118        existing.push(row);
1119    }
1120    Ok(())
1121}
1122
1123fn persist_descriptor(dir: &Path, descriptor: &SeriesDescriptor) -> Result<()> {
1124    let path = dir.join("_descriptor.json");
1125    if path.exists() {
1126        return verify_descriptor(dir, descriptor);
1127    }
1128    fs::write(
1129        path,
1130        serde_json::to_vec_pretty(descriptor)
1131            .map_err(|error| DataError::Other(error.to_string()))?,
1132    )?;
1133    Ok(())
1134}
1135fn verify_descriptor(dir: &Path, descriptor: &SeriesDescriptor) -> Result<()> {
1136    let path = dir.join("_descriptor.json");
1137    if !path.exists() {
1138        if descriptor.verified {
1139            return Err(DataError::Other(
1140                "verified descriptor metadata is absent".into(),
1141            ));
1142        }
1143        return Ok(());
1144    }
1145    let stored: SeriesDescriptor = serde_json::from_slice(&fs::read(path)?)
1146        .map_err(|error| DataError::Other(error.to_string()))?;
1147    if &stored != descriptor {
1148        return Err(DataError::Other(
1149            "series descriptor conflicts with stored metadata".into(),
1150        ));
1151    }
1152    Ok(())
1153}
1154
1155/// Concat two tick DataFrames, dedup on (exchange, symbol, ts), sort by ts.
1156fn concat_and_dedup_ticks(existing: DataFrame, new: DataFrame) -> Result<DataFrame> {
1157    let combined = concat_dataframes(vec![existing, new])?;
1158    dedup_ticks(combined)
1159}
1160
1161/// Dedup a tick DataFrame on (exchange, symbol, ts) keeping first, sort by ts.
1162fn dedup_ticks(df: DataFrame) -> Result<DataFrame> {
1163    let cols: Vec<String> = vec!["exchange".into(), "symbol".into(), "ts".into()];
1164    let deduped = df
1165        .unique_stable(Some(&cols), UniqueKeepStrategy::First, None)?
1166        .sort(["ts"], SortMultipleOptions::default())?;
1167    Ok(deduped)
1168}
1169
1170/// Concat two bar DataFrames, dedup on (exchange, symbol, timeframe, ts), sort by ts.
1171fn concat_and_dedup_bars(existing: DataFrame, new: DataFrame) -> Result<DataFrame> {
1172    let combined = concat_dataframes(vec![existing, new])?;
1173    dedup_bars(combined)
1174}
1175
1176/// Dedup a bar DataFrame on (exchange, symbol, timeframe, ts) keeping first, sort by ts.
1177fn dedup_bars(df: DataFrame) -> Result<DataFrame> {
1178    let cols: Vec<String> = vec![
1179        "exchange".into(),
1180        "symbol".into(),
1181        "timeframe".into(),
1182        "ts".into(),
1183    ];
1184    let deduped = df
1185        .unique_stable(Some(&cols), UniqueKeepStrategy::First, None)?
1186        .sort(["ts"], SortMultipleOptions::default())?;
1187    Ok(deduped)
1188}
1189
1190/// Vertically concatenate multiple DataFrames.
1191fn concat_dataframes(dfs: Vec<DataFrame>) -> Result<DataFrame> {
1192    if dfs.is_empty() {
1193        return Err(DataError::Other("no dataframes to concat".into()));
1194    }
1195    if dfs.len() == 1 {
1196        return Ok(dfs.into_iter().next().unwrap());
1197    }
1198    let lazy_frames: Vec<LazyFrame> = dfs.into_iter().map(|df| df.lazy()).collect();
1199    let combined = polars::prelude::concat(lazy_frames, Default::default())?.collect()?;
1200    Ok(combined)
1201}
1202
1203/// Apply timestamp range filter to a DataFrame with a "ts" datetime column.
1204fn apply_ts_filter(
1205    df: DataFrame,
1206    from: Option<NaiveDateTime>,
1207    to: Option<NaiveDateTime>,
1208) -> Result<DataFrame> {
1209    if from.is_none() && to.is_none() {
1210        return Ok(df);
1211    }
1212
1213    let mut lf = df.lazy();
1214
1215    if let Some(f) = from {
1216        let from_micros = f.and_utc().timestamp_micros();
1217        lf = lf.filter(
1218            col("ts")
1219                .gt_eq(lit(from_micros).cast(DataType::Datetime(TimeUnit::Microseconds, None))),
1220        );
1221    }
1222    if let Some(t) = to {
1223        let to_micros = t.and_utc().timestamp_micros();
1224        lf = lf.filter(
1225            col("ts").lt_eq(lit(to_micros).cast(DataType::Datetime(TimeUnit::Microseconds, None))),
1226        );
1227    }
1228
1229    Ok(lf.collect()?)
1230}
1231
1232/// Apply limit, tail, and descending pagination to a sorted DataFrame.
1233fn apply_pagination(
1234    df: DataFrame,
1235    limit: usize,
1236    tail: bool,
1237    descending: bool,
1238) -> Result<DataFrame> {
1239    // limit == 0 means "no limit" — return all rows.
1240    let result = if tail && limit > 0 {
1241        // Take last N rows, then optionally reverse for descending
1242        let n = limit.min(df.height());
1243        let tailed = df.tail(Some(n));
1244        if descending {
1245            tailed.sort(
1246                ["ts"],
1247                SortMultipleOptions::default().with_order_descending(true),
1248            )?
1249        } else {
1250            tailed
1251        }
1252    } else if descending {
1253        let sorted = df.sort(
1254            ["ts"],
1255            SortMultipleOptions::default().with_order_descending(true),
1256        )?;
1257        if limit > 0 {
1258            sorted.head(Some(limit))
1259        } else {
1260            sorted
1261        }
1262    } else if limit > 0 {
1263        df.head(Some(limit))
1264    } else {
1265        df
1266    };
1267    Ok(result)
1268}
1269
1270/// Delete rows from a date-partitioned directory, optionally within a date range.
1271fn delete_from_partition(
1272    dir: &Path,
1273    from: Option<NaiveDateTime>,
1274    to: Option<NaiveDateTime>,
1275) -> Result<usize> {
1276    if from.is_none() && to.is_none() {
1277        // Delete everything in the directory
1278        let count = count_all_rows_in_dir(dir);
1279        // Remove all parquet files but keep the directory
1280        for entry in fs::read_dir(dir)?.flatten() {
1281            let path = entry.path();
1282            if path.extension().is_some_and(|e| e == "parquet") {
1283                fs::remove_file(&path)?;
1284            }
1285        }
1286        return Ok(count);
1287    }
1288
1289    let files = list_date_files(dir, from, to)?;
1290    let mut total_deleted = 0usize;
1291
1292    for file_path in &files {
1293        let df = read_parquet_file(file_path)?;
1294        let original_count = df.height();
1295
1296        // Filter to keep rows OUTSIDE the delete range
1297        let filtered = apply_ts_filter_inverted(df, from, to)?;
1298
1299        if filtered.height() == 0 {
1300            // All rows deleted — remove the file
1301            fs::remove_file(file_path)?;
1302            total_deleted += original_count;
1303        } else if filtered.height() < original_count {
1304            // Partial deletion — rewrite the file
1305            total_deleted += original_count - filtered.height();
1306            write_parquet_file(file_path, &mut filtered.clone())?;
1307        }
1308        // else: no rows matched the range in this file
1309    }
1310
1311    Ok(total_deleted)
1312}
1313
1314/// Filter to keep rows OUTSIDE a timestamp range (inverse of apply_ts_filter).
1315fn apply_ts_filter_inverted(
1316    df: DataFrame,
1317    from: Option<NaiveDateTime>,
1318    to: Option<NaiveDateTime>,
1319) -> Result<DataFrame> {
1320    let mut lf = df.lazy();
1321
1322    match (from, to) {
1323        (Some(f), Some(t)) => {
1324            let from_micros = f.and_utc().timestamp_micros();
1325            let to_micros = t.and_utc().timestamp_micros();
1326            let from_lit = lit(from_micros).cast(DataType::Datetime(TimeUnit::Microseconds, None));
1327            let to_lit = lit(to_micros).cast(DataType::Datetime(TimeUnit::Microseconds, None));
1328            // Keep rows where ts < from OR ts > to
1329            lf = lf.filter(col("ts").lt(from_lit).or(col("ts").gt(to_lit)));
1330        }
1331        (Some(f), None) => {
1332            let from_micros = f.and_utc().timestamp_micros();
1333            let from_lit = lit(from_micros).cast(DataType::Datetime(TimeUnit::Microseconds, None));
1334            lf = lf.filter(col("ts").lt(from_lit));
1335        }
1336        (None, Some(t)) => {
1337            let to_micros = t.and_utc().timestamp_micros();
1338            let to_lit = lit(to_micros).cast(DataType::Datetime(TimeUnit::Microseconds, None));
1339            lf = lf.filter(col("ts").gt(to_lit));
1340        }
1341        (None, None) => {}
1342    }
1343
1344    Ok(lf.collect()?)
1345}
1346
1347/// Count all rows across parquet files in a directory (non-recursive).
1348fn count_all_rows_in_dir(dir: &Path) -> usize {
1349    let mut count = 0;
1350    if let Ok(entries) = fs::read_dir(dir) {
1351        for entry in entries.flatten() {
1352            let path = entry.path();
1353            if path.extension().is_some_and(|e| e == "parquet")
1354                && let Ok(df) = read_parquet_file(&path)
1355            {
1356                count += df.height();
1357            }
1358        }
1359    }
1360    count
1361}
1362
1363/// Recursively walk a directory and collect all file paths.
1364fn walkdir(dir: &Path) -> Vec<PathBuf> {
1365    let mut result = Vec::new();
1366    if !dir.exists() {
1367        return result;
1368    }
1369    if let Ok(entries) = fs::read_dir(dir) {
1370        for entry in entries.flatten() {
1371            let path = entry.path();
1372            if path.is_dir() {
1373                result.extend(walkdir(&path));
1374            } else {
1375                result.push(path);
1376            }
1377        }
1378    }
1379    result
1380}
1381
1382/// Convert microsecond epoch to NaiveDateTime.
1383fn micros_to_ndt(micros: i64) -> NaiveDateTime {
1384    let secs = micros / 1_000_000;
1385    let nsecs = ((micros % 1_000_000) * 1_000) as u32;
1386    chrono::DateTime::from_timestamp(secs, nsecs)
1387        .map(|dt| dt.naive_utc())
1388        .unwrap_or_default()
1389}