1use 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
27pub struct ParquetStore {
29 root: PathBuf,
30}
31
32impl ParquetStore {
33 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 pub fn insert_ticks(&self, ticks: &[Tick]) -> Result<usize> {
45 if ticks.is_empty() {
46 return Ok(0);
47 }
48
49 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 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 pub fn insert_bars(&self, bars: &[Bar]) -> Result<usize> {
172 if bars.is_empty() {
173 return Ok(0);
174 }
175
176 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 pub fn query_ticks(&self, opts: &QueryOpts) -> Result<(Vec<Tick>, u64)> {
329 self.query_ticks_cancellable(opts, || false)
330 }
331
332 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 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 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 pub fn query_bars(&self, opts: &BarQueryOpts) -> Result<(Vec<Bar>, u64)> {
423 self.query_bars_cancellable(opts, || false)
424 }
425
426 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 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 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 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 let tick_dir = self.tick_dir(exchange, symbol);
511 if tick_dir.exists() {
512 fs::remove_dir_all(&tick_dir)?;
513 }
514
515 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 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 pub fn stats(&self, exchange: Option<&str>, symbol: Option<&str>) -> Result<Vec<StatRow>> {
550 let mut rows = Vec::new();
551
552 self.collect_tick_stats(&mut rows, exchange, symbol)?;
554
555 self.collect_bar_stats(&mut rows, exchange, symbol)?;
557
558 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 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 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 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 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 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 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 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 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 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 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 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 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
885fn 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
944fn 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 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
985fn 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
994fn 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
1155fn concat_and_dedup_ticks(existing: DataFrame, new: DataFrame) -> Result<DataFrame> {
1157 let combined = concat_dataframes(vec![existing, new])?;
1158 dedup_ticks(combined)
1159}
1160
1161fn 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
1170fn concat_and_dedup_bars(existing: DataFrame, new: DataFrame) -> Result<DataFrame> {
1172 let combined = concat_dataframes(vec![existing, new])?;
1173 dedup_bars(combined)
1174}
1175
1176fn 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
1190fn 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
1203fn 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
1232fn apply_pagination(
1234 df: DataFrame,
1235 limit: usize,
1236 tail: bool,
1237 descending: bool,
1238) -> Result<DataFrame> {
1239 let result = if tail && limit > 0 {
1241 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
1270fn 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 let count = count_all_rows_in_dir(dir);
1279 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 let filtered = apply_ts_filter_inverted(df, from, to)?;
1298
1299 if filtered.height() == 0 {
1300 fs::remove_file(file_path)?;
1302 total_deleted += original_count;
1303 } else if filtered.height() < original_count {
1304 total_deleted += original_count - filtered.height();
1306 write_parquet_file(file_path, &mut filtered.clone())?;
1307 }
1308 }
1310
1311 Ok(total_deleted)
1312}
1313
1314fn 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 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
1347fn 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
1363fn 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
1382fn 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}