1use std::collections::VecDeque;
4use std::fs::{self, File, Metadata};
5use std::io::{Seek, SeekFrom};
6#[cfg(unix)]
7use std::os::unix::fs::MetadataExt;
8#[cfg(windows)]
9use std::os::windows::fs::MetadataExt;
10use std::path::{Path, PathBuf};
11use std::time::SystemTime;
12
13use chrono::{NaiveDate, NaiveDateTime};
14use polars::prelude::*;
15
16use crate::convert::{
17 dataframe_to_bars, dataframe_to_price_bars, dataframe_to_stored_ticks, dataframe_to_ticks,
18};
19use crate::error::{DataError, Result};
20use crate::models::{Bar, PriceBar, SeriesDescriptor, StoredTick, Tick};
21use crate::parquet_store::ParquetStore;
22
23pub const DEFAULT_PARQUET_SCAN_ROWS: usize = 65_536;
25
26#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
28pub struct ParquetScanBounds {
29 pub from: Option<NaiveDateTime>,
30 pub to: Option<NaiveDateTime>,
31}
32
33impl ParquetScanBounds {
34 pub const fn new(from: Option<NaiveDateTime>, to: Option<NaiveDateTime>) -> Self {
35 Self { from, to }
36 }
37
38 fn validate(self) -> Result<Self> {
39 if let (Some(from), Some(to)) = (self.from, self.to)
40 && from > to
41 {
42 return Err(DataError::InvalidScanBounds { from, to });
43 }
44 Ok(self)
45 }
46
47 fn contains(self, ts: NaiveDateTime) -> bool {
48 self.from.is_none_or(|from| ts >= from) && self.to.is_none_or(|to| ts <= to)
49 }
50}
51
52#[derive(Debug, Clone, PartialEq)]
54pub struct ParquetScannedRow<T> {
55 pub row: T,
56 pub source_row_ordinal: u64,
57}
58
59trait Timestamped {
60 fn timestamp(&self) -> NaiveDateTime;
61}
62
63impl Timestamped for Tick {
64 fn timestamp(&self) -> NaiveDateTime {
65 self.ts
66 }
67}
68
69impl Timestamped for Bar {
70 fn timestamp(&self) -> NaiveDateTime {
71 self.ts
72 }
73}
74
75impl Timestamped for StoredTick {
76 fn timestamp(&self) -> NaiveDateTime {
77 self.tick.ts
78 }
79}
80
81impl Timestamped for PriceBar {
82 fn timestamp(&self) -> NaiveDateTime {
83 self.available_at
84 }
85}
86
87#[derive(Debug, Clone, PartialEq, Eq)]
88struct FileFingerprint {
89 len: u64,
90 modified: Option<SystemTime>,
91 created: Option<SystemTime>,
92 #[cfg(unix)]
93 device: u64,
94 #[cfg(unix)]
95 inode: u64,
96 #[cfg(windows)]
97 volume_serial_number: Option<u32>,
98 #[cfg(windows)]
99 file_index: Option<u64>,
100}
101
102impl FileFingerprint {
103 fn from_metadata(metadata: &Metadata) -> Self {
104 Self {
105 len: metadata.len(),
106 modified: metadata.modified().ok(),
107 created: metadata.created().ok(),
108 #[cfg(unix)]
109 device: metadata.dev(),
110 #[cfg(unix)]
111 inode: metadata.ino(),
112 #[cfg(windows)]
113 volume_serial_number: metadata.volume_serial_number(),
114 #[cfg(windows)]
115 file_index: metadata.file_index(),
116 }
117 }
118}
119
120#[derive(Debug, Clone)]
121struct PartitionDescriptor {
122 path: PathBuf,
123 fingerprint: FileFingerprint,
124 row_count: usize,
125 source_row_base: u64,
126}
127
128impl PartitionDescriptor {
129 fn describe(path: PathBuf, source_row_base: u64) -> Result<Self> {
130 let file = File::open(&path)?;
131 let fingerprint = FileFingerprint::from_metadata(&file.metadata()?);
132 ensure_path_matches(&path, &fingerprint)?;
133
134 let mut reader = ParquetReader::new(file.try_clone()?);
135 let row_count = reader.num_rows()?;
136 ensure_file_matches(&path, &file, &fingerprint)?;
137 ensure_path_matches(&path, &fingerprint)?;
138
139 Ok(Self {
140 path,
141 fingerprint,
142 row_count,
143 source_row_base,
144 })
145 }
146
147 fn validate(&self) -> Result<()> {
148 ensure_path_matches(&self.path, &self.fingerprint)
149 }
150
151 fn open_generation(&self) -> Result<OpenPartition> {
152 self.validate()?;
153 let file = match File::open(&self.path) {
154 Ok(file) => file,
155 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
156 return Err(partition_changed(&self.path));
157 }
158 Err(error) => return Err(error.into()),
159 };
160 ensure_file_matches(&self.path, &file, &self.fingerprint)?;
161 self.validate()?;
162 Ok(OpenPartition {
163 descriptor: self.clone(),
164 file,
165 })
166 }
167
168 fn ordinal(&self, row_offset: usize) -> Result<u64> {
169 self.source_row_base
170 .checked_add(
171 u64::try_from(row_offset).map_err(|_| {
172 DataError::Other("Parquet source row ordinal exceeds u64".into())
173 })?,
174 )
175 .ok_or_else(|| DataError::Other("Parquet source row ordinal overflow".into()))
176 }
177}
178
179struct OpenPartition {
180 descriptor: PartitionDescriptor,
181 file: File,
182}
183
184impl OpenPartition {
185 fn validate(&self) -> Result<()> {
186 ensure_file_matches(
187 &self.descriptor.path,
188 &self.file,
189 &self.descriptor.fingerprint,
190 )?;
191 self.descriptor.validate()
192 }
193
194 fn reader_file(&self) -> Result<File> {
195 self.validate()?;
196 let mut file = self.file.try_clone()?;
197 file.seek(SeekFrom::Start(0))?;
198 Ok(file)
199 }
200}
201
202fn partition_changed(path: &Path) -> DataError {
203 DataError::ParquetPartitionChanged {
204 path: path.display().to_string(),
205 }
206}
207
208fn ensure_file_matches(path: &Path, file: &File, expected: &FileFingerprint) -> Result<()> {
209 if FileFingerprint::from_metadata(&file.metadata()?) == *expected {
210 Ok(())
211 } else {
212 Err(partition_changed(path))
213 }
214}
215
216fn ensure_path_matches(path: &Path, expected: &FileFingerprint) -> Result<()> {
217 let metadata = match fs::metadata(path) {
218 Ok(metadata) => metadata,
219 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
220 return Err(partition_changed(path));
221 }
222 Err(error) => return Err(error.into()),
223 };
224 if FileFingerprint::from_metadata(&metadata) == *expected {
225 Ok(())
226 } else {
227 Err(partition_changed(path))
228 }
229}
230
231#[derive(Debug, Clone)]
232struct PartitionScan {
233 partitions: Vec<PartitionDescriptor>,
234 bounds: ParquetScanBounds,
235}
236
237impl PartitionScan {
238 fn describe(
239 directory: &Path,
240 bounds: ParquetScanBounds,
241 is_cancelled: &mut dyn FnMut() -> bool,
242 ) -> Result<Self> {
243 let bounds = bounds.validate()?;
244 let partitions = list_partitions(directory, bounds, is_cancelled)?;
245 Ok(Self { partitions, bounds })
246 }
247
248 fn validate(&self) -> Result<()> {
249 for partition in &self.partitions {
250 partition.validate()?;
251 }
252 Ok(())
253 }
254}
255
256type PartitionLoader<T> = fn(File, usize, usize) -> Result<Vec<T>>;
257
258struct PartitionCursor<T> {
259 partitions: VecDeque<PartitionDescriptor>,
260 current_partition: Option<OpenPartition>,
261 current_rows: std::vec::IntoIter<ParquetScannedRow<T>>,
262 bounds: ParquetScanBounds,
263 rows_per_read: usize,
264 partition_offset: usize,
265 finish_current_partition: bool,
266 last_read_ts: Option<NaiveDateTime>,
267 load_partition: PartitionLoader<T>,
268}
269
270impl<T: Timestamped> PartitionCursor<T> {
271 fn open(
272 scan: PartitionScan,
273 rows_per_read: usize,
274 load_partition: PartitionLoader<T>,
275 ) -> Result<Self> {
276 if rows_per_read == 0 {
277 return Err(DataError::InvalidScanReadSize);
278 }
279 scan.validate()?;
280 Ok(Self {
281 partitions: scan.partitions.into(),
282 current_partition: None,
283 current_rows: Vec::new().into_iter(),
284 bounds: scan.bounds,
285 rows_per_read,
286 partition_offset: 0,
287 finish_current_partition: false,
288 last_read_ts: None,
289 load_partition,
290 })
291 }
292
293 fn next_row(
294 &mut self,
295 is_cancelled: &mut dyn FnMut() -> bool,
296 ) -> Result<Option<ParquetScannedRow<T>>> {
297 loop {
298 ensure_not_cancelled(is_cancelled)?;
299 if !self.current_rows.as_slice().is_empty() {
300 self.current_partition
301 .as_ref()
302 .expect("buffered rows retain their open partition")
303 .validate()?;
304 return Ok(self.current_rows.next());
305 }
306
307 if self.finish_current_partition {
308 self.current_partition = None;
309 self.partition_offset = 0;
310 self.finish_current_partition = false;
311 }
312
313 if self.current_partition.is_none() {
314 let Some(descriptor) = self.partitions.pop_front() else {
315 return Ok(None);
316 };
317 self.current_partition = Some(descriptor.open_generation()?);
318 }
319
320 let partition = self
321 .current_partition
322 .as_ref()
323 .expect("current partition was opened");
324 if self.partition_offset >= partition.descriptor.row_count {
325 self.finish_current_partition = true;
326 continue;
327 }
328
329 ensure_not_cancelled(is_cancelled)?;
330 let rows_to_read = self
331 .rows_per_read
332 .min(partition.descriptor.row_count - self.partition_offset);
333 let reader = partition.reader_file()?;
334 let rows = (self.load_partition)(reader, self.partition_offset, rows_to_read)?;
335 ensure_not_cancelled(is_cancelled)?;
336 partition.validate()?;
337 if rows.len() != rows_to_read {
338 return Err(DataError::Other(format!(
339 "Parquet slice {} at offset {} returned {} rows instead of {}",
340 partition.descriptor.path.display(),
341 self.partition_offset,
342 rows.len(),
343 rows_to_read
344 )));
345 }
346
347 let read_offset = self.partition_offset;
348 let (keep_len, last_read_ts, exceeded_upper_bound) = validate_monotonic_prefix(
349 self.last_read_ts,
350 &rows,
351 &partition.descriptor.path,
352 self.bounds.to,
353 )?;
354 self.last_read_ts = last_read_ts;
355 self.partition_offset = self
356 .partition_offset
357 .checked_add(rows.len())
358 .ok_or_else(|| DataError::Other("Parquet scan row offset overflow".into()))?;
359 self.finish_current_partition =
360 exceeded_upper_bound || self.partition_offset >= partition.descriptor.row_count;
361 if exceeded_upper_bound {
362 self.partitions.clear();
363 }
364
365 let descriptor = &partition.descriptor;
366 let mut bounded_rows = Vec::with_capacity(keep_len);
367 for (index, row) in rows.into_iter().take(keep_len).enumerate() {
368 if self.bounds.contains(row.timestamp()) {
369 bounded_rows.push(ParquetScannedRow {
370 row,
371 source_row_ordinal: descriptor.ordinal(read_offset + index)?,
372 });
373 }
374 }
375 self.current_rows = bounded_rows.into_iter();
376 }
377 }
378
379 fn remaining_partitions(&self) -> usize {
380 self.partitions.len() + usize::from(self.current_partition.is_some())
381 }
382
383 fn rows_per_read(&self) -> usize {
384 self.rows_per_read
385 }
386}
387
388fn validate_monotonic_prefix<T: Timestamped>(
389 mut previous: Option<NaiveDateTime>,
390 rows: &[T],
391 path: &Path,
392 inclusive_to: Option<NaiveDateTime>,
393) -> Result<(usize, Option<NaiveDateTime>, bool)> {
394 for (index, row) in rows.iter().enumerate() {
395 let current = row.timestamp();
396 if let Some(previous) = previous
397 && current < previous
398 {
399 return Err(DataError::NonMonotonicParquetData {
400 path: path.display().to_string(),
401 previous,
402 current,
403 });
404 }
405 if inclusive_to.is_some_and(|to| current > to) {
406 return Ok((index, previous, true));
407 }
408 previous = Some(current);
409 }
410 Ok((rows.len(), previous, false))
411}
412
413fn find_latest_row<T, P>(
414 scan: &PartitionScan,
415 rows_per_read: usize,
416 load_partition: PartitionLoader<T>,
417 mut predicate: P,
418 is_cancelled: &mut dyn FnMut() -> bool,
419) -> Result<Option<ParquetScannedRow<T>>>
420where
421 T: Timestamped,
422 P: FnMut(&T) -> bool,
423{
424 if rows_per_read == 0 {
425 return Err(DataError::InvalidScanReadSize);
426 }
427 scan.validate()?;
428 let mut later_first_ts = None;
429
430 for descriptor in scan.partitions.iter().rev() {
431 ensure_not_cancelled(is_cancelled)?;
432 let partition = descriptor.open_generation()?;
433 let mut end = descriptor.row_count;
434 while end > 0 {
435 ensure_not_cancelled(is_cancelled)?;
436 let start = end.saturating_sub(rows_per_read);
437 let reader = partition.reader_file()?;
438 let rows = load_partition(reader, start, end - start)?;
439 ensure_not_cancelled(is_cancelled)?;
440 partition.validate()?;
441 if rows.len() != end - start {
442 return Err(DataError::Other(format!(
443 "Parquet reverse slice {} at offset {} returned {} rows instead of {}",
444 descriptor.path.display(),
445 start,
446 rows.len(),
447 end - start
448 )));
449 }
450 validate_monotonic_rows(None, &rows, &descriptor.path)?;
451 if let (Some(last), Some(later)) =
452 (rows.last().map(Timestamped::timestamp), later_first_ts)
453 && last > later
454 {
455 return Err(DataError::NonMonotonicParquetData {
456 path: descriptor.path.display().to_string(),
457 previous: last,
458 current: later,
459 });
460 }
461 if let Some(first) = rows.first().map(Timestamped::timestamp) {
462 later_first_ts = Some(first);
463 }
464
465 for (index, row) in rows.into_iter().enumerate().rev() {
466 let timestamp = row.timestamp();
467 if scan.bounds.from.is_some_and(|from| timestamp < from) {
468 return Ok(None);
469 }
470 if scan.bounds.contains(timestamp) && predicate(&row) {
471 return Ok(Some(ParquetScannedRow {
472 row,
473 source_row_ordinal: descriptor.ordinal(start + index)?,
474 }));
475 }
476 }
477 end = start;
478 }
479 }
480
481 ensure_not_cancelled(is_cancelled)?;
482 Ok(None)
483}
484
485#[derive(Debug, Clone)]
487pub struct ParquetTickScan {
488 inner: PartitionScan,
489}
490
491impl ParquetTickScan {
492 pub fn describe(
494 root: impl AsRef<Path>,
495 exchange: &str,
496 symbol: &str,
497 bounds: ParquetScanBounds,
498 ) -> Result<Self> {
499 Self::describe_cancellable(root, exchange, symbol, bounds, || false)
500 }
501
502 pub fn describe_cancellable<F>(
504 root: impl AsRef<Path>,
505 exchange: &str,
506 symbol: &str,
507 bounds: ParquetScanBounds,
508 mut is_cancelled: F,
509 ) -> Result<Self>
510 where
511 F: FnMut() -> bool,
512 {
513 let directory = root
514 .as_ref()
515 .join("ticks")
516 .join(format!("exchange={exchange}"))
517 .join(format!("symbol={symbol}"));
518 Ok(Self {
519 inner: PartitionScan::describe(&directory, bounds, &mut is_cancelled)?,
520 })
521 }
522
523 pub fn cursor(&self) -> Result<ParquetTickCursor> {
525 self.cursor_with_read_size(DEFAULT_PARQUET_SCAN_ROWS)
526 }
527
528 pub fn cursor_with_read_size(&self, rows_per_read: usize) -> Result<ParquetTickCursor> {
530 Ok(ParquetTickCursor {
531 inner: PartitionCursor::open(self.inner.clone(), rows_per_read, read_tick_partition)?,
532 })
533 }
534
535 pub fn latest_valid_tick_cancellable<F>(
537 &self,
538 mut is_cancelled: F,
539 ) -> Result<Option<ParquetScannedRow<Tick>>>
540 where
541 F: FnMut() -> bool,
542 {
543 find_latest_row(
544 &self.inner,
545 DEFAULT_PARQUET_SCAN_ROWS,
546 read_tick_partition,
547 tick_has_valid_quote,
548 &mut is_cancelled,
549 )
550 }
551
552 pub fn latest_valid_tick_before_cancellable<F>(
554 &self,
555 before: NaiveDateTime,
556 mut is_cancelled: F,
557 ) -> Result<Option<ParquetScannedRow<Tick>>>
558 where
559 F: FnMut() -> bool,
560 {
561 find_latest_row(
562 &self.inner,
563 DEFAULT_PARQUET_SCAN_ROWS,
564 read_tick_partition,
565 |tick| tick.ts < before && tick_has_valid_quote(tick),
566 &mut is_cancelled,
567 )
568 }
569}
570
571pub struct ParquetTickCursor {
573 inner: PartitionCursor<Tick>,
574}
575
576impl ParquetTickCursor {
577 pub fn open(
579 root: impl AsRef<Path>,
580 exchange: &str,
581 symbol: &str,
582 bounds: ParquetScanBounds,
583 ) -> Result<Self> {
584 Self::open_with_read_size(root, exchange, symbol, bounds, DEFAULT_PARQUET_SCAN_ROWS)
585 }
586
587 pub fn open_with_read_size(
589 root: impl AsRef<Path>,
590 exchange: &str,
591 symbol: &str,
592 bounds: ParquetScanBounds,
593 rows_per_read: usize,
594 ) -> Result<Self> {
595 Self::open_cancellable_with_read_size(root, exchange, symbol, bounds, rows_per_read, || {
596 false
597 })
598 }
599
600 pub fn open_cancellable<F>(
602 root: impl AsRef<Path>,
603 exchange: &str,
604 symbol: &str,
605 bounds: ParquetScanBounds,
606 is_cancelled: F,
607 ) -> Result<Self>
608 where
609 F: FnMut() -> bool,
610 {
611 Self::open_cancellable_with_read_size(
612 root,
613 exchange,
614 symbol,
615 bounds,
616 DEFAULT_PARQUET_SCAN_ROWS,
617 is_cancelled,
618 )
619 }
620
621 pub fn open_cancellable_with_read_size<F>(
623 root: impl AsRef<Path>,
624 exchange: &str,
625 symbol: &str,
626 bounds: ParquetScanBounds,
627 rows_per_read: usize,
628 is_cancelled: F,
629 ) -> Result<Self>
630 where
631 F: FnMut() -> bool,
632 {
633 ParquetTickScan::describe_cancellable(root, exchange, symbol, bounds, is_cancelled)?
634 .cursor_with_read_size(rows_per_read)
635 }
636
637 pub fn next_tick(&mut self) -> Result<Option<Tick>> {
639 self.next_tick_with_ordinal()
640 .map(|row| row.map(|row| row.row))
641 }
642
643 pub fn next_tick_with_ordinal(&mut self) -> Result<Option<ParquetScannedRow<Tick>>> {
645 self.next_tick_with_ordinal_cancellable(|| false)
646 }
647
648 pub fn next_tick_cancellable<F>(&mut self, is_cancelled: F) -> Result<Option<Tick>>
650 where
651 F: FnMut() -> bool,
652 {
653 self.next_tick_with_ordinal_cancellable(is_cancelled)
654 .map(|row| row.map(|row| row.row))
655 }
656
657 pub fn next_tick_with_ordinal_cancellable<F>(
659 &mut self,
660 mut is_cancelled: F,
661 ) -> Result<Option<ParquetScannedRow<Tick>>>
662 where
663 F: FnMut() -> bool,
664 {
665 self.inner.next_row(&mut is_cancelled)
666 }
667
668 pub fn remaining_partitions(&self) -> usize {
670 self.inner.remaining_partitions()
671 }
672
673 pub fn rows_per_read(&self) -> usize {
675 self.inner.rows_per_read()
676 }
677}
678
679#[derive(Debug, Clone)]
681pub struct ParquetBarScan {
682 inner: PartitionScan,
683}
684
685impl ParquetBarScan {
686 pub fn describe(
688 root: impl AsRef<Path>,
689 exchange: &str,
690 symbol: &str,
691 timeframe: &str,
692 bounds: ParquetScanBounds,
693 ) -> Result<Self> {
694 Self::describe_cancellable(root, exchange, symbol, timeframe, bounds, || false)
695 }
696
697 pub fn describe_cancellable<F>(
699 root: impl AsRef<Path>,
700 exchange: &str,
701 symbol: &str,
702 timeframe: &str,
703 bounds: ParquetScanBounds,
704 mut is_cancelled: F,
705 ) -> Result<Self>
706 where
707 F: FnMut() -> bool,
708 {
709 let directory = root
710 .as_ref()
711 .join("bars")
712 .join(format!("exchange={exchange}"))
713 .join(format!("symbol={symbol}"))
714 .join(format!("timeframe={timeframe}"));
715 Ok(Self {
716 inner: PartitionScan::describe(&directory, bounds, &mut is_cancelled)?,
717 })
718 }
719
720 pub fn cursor(&self) -> Result<ParquetBarCursor> {
722 self.cursor_with_read_size(DEFAULT_PARQUET_SCAN_ROWS)
723 }
724
725 pub fn cursor_with_read_size(&self, rows_per_read: usize) -> Result<ParquetBarCursor> {
727 Ok(ParquetBarCursor {
728 inner: PartitionCursor::open(self.inner.clone(), rows_per_read, read_bar_partition)?,
729 })
730 }
731
732 pub fn latest_valid_bar_cancellable<F>(
734 &self,
735 mut is_cancelled: F,
736 ) -> Result<Option<ParquetScannedRow<Bar>>>
737 where
738 F: FnMut() -> bool,
739 {
740 find_latest_row(
741 &self.inner,
742 DEFAULT_PARQUET_SCAN_ROWS,
743 read_bar_partition,
744 |bar| bar.close.is_finite() && bar.close > 0.0,
745 &mut is_cancelled,
746 )
747 }
748}
749
750pub struct ParquetBarCursor {
752 inner: PartitionCursor<Bar>,
753}
754
755impl ParquetBarCursor {
756 pub fn open(
758 root: impl AsRef<Path>,
759 exchange: &str,
760 symbol: &str,
761 timeframe: &str,
762 bounds: ParquetScanBounds,
763 ) -> Result<Self> {
764 Self::open_with_read_size(
765 root,
766 exchange,
767 symbol,
768 timeframe,
769 bounds,
770 DEFAULT_PARQUET_SCAN_ROWS,
771 )
772 }
773
774 pub fn open_with_read_size(
776 root: impl AsRef<Path>,
777 exchange: &str,
778 symbol: &str,
779 timeframe: &str,
780 bounds: ParquetScanBounds,
781 rows_per_read: usize,
782 ) -> Result<Self> {
783 Self::open_cancellable_with_read_size(
784 root,
785 exchange,
786 symbol,
787 timeframe,
788 bounds,
789 rows_per_read,
790 || false,
791 )
792 }
793
794 pub fn open_cancellable<F>(
796 root: impl AsRef<Path>,
797 exchange: &str,
798 symbol: &str,
799 timeframe: &str,
800 bounds: ParquetScanBounds,
801 is_cancelled: F,
802 ) -> Result<Self>
803 where
804 F: FnMut() -> bool,
805 {
806 Self::open_cancellable_with_read_size(
807 root,
808 exchange,
809 symbol,
810 timeframe,
811 bounds,
812 DEFAULT_PARQUET_SCAN_ROWS,
813 is_cancelled,
814 )
815 }
816
817 pub fn open_cancellable_with_read_size<F>(
819 root: impl AsRef<Path>,
820 exchange: &str,
821 symbol: &str,
822 timeframe: &str,
823 bounds: ParquetScanBounds,
824 rows_per_read: usize,
825 is_cancelled: F,
826 ) -> Result<Self>
827 where
828 F: FnMut() -> bool,
829 {
830 ParquetBarScan::describe_cancellable(
831 root,
832 exchange,
833 symbol,
834 timeframe,
835 bounds,
836 is_cancelled,
837 )?
838 .cursor_with_read_size(rows_per_read)
839 }
840
841 pub fn next_bar(&mut self) -> Result<Option<Bar>> {
843 self.next_bar_with_ordinal()
844 .map(|row| row.map(|row| row.row))
845 }
846
847 pub fn next_bar_with_ordinal(&mut self) -> Result<Option<ParquetScannedRow<Bar>>> {
849 self.next_bar_with_ordinal_cancellable(|| false)
850 }
851
852 pub fn next_bar_cancellable<F>(&mut self, is_cancelled: F) -> Result<Option<Bar>>
854 where
855 F: FnMut() -> bool,
856 {
857 self.next_bar_with_ordinal_cancellable(is_cancelled)
858 .map(|row| row.map(|row| row.row))
859 }
860
861 pub fn next_bar_with_ordinal_cancellable<F>(
863 &mut self,
864 mut is_cancelled: F,
865 ) -> Result<Option<ParquetScannedRow<Bar>>>
866 where
867 F: FnMut() -> bool,
868 {
869 self.inner.next_row(&mut is_cancelled)
870 }
871
872 pub fn remaining_partitions(&self) -> usize {
874 self.inner.remaining_partitions()
875 }
876
877 pub fn rows_per_read(&self) -> usize {
879 self.inner.rows_per_read()
880 }
881}
882
883#[derive(Debug, Clone)]
885pub struct ParquetStoredTickScan {
886 inner: PartitionScan,
887}
888
889impl ParquetStoredTickScan {
890 pub fn describe_cancellable<F>(
891 root: impl AsRef<Path>,
892 exchange: &str,
893 symbol: &str,
894 bounds: ParquetScanBounds,
895 mut is_cancelled: F,
896 ) -> Result<Self>
897 where
898 F: FnMut() -> bool,
899 {
900 let directory = root
901 .as_ref()
902 .join("ordered_ticks")
903 .join(format!("exchange={exchange}"))
904 .join(format!("symbol={symbol}"));
905 Ok(Self {
906 inner: PartitionScan::describe(&directory, bounds, &mut is_cancelled)?,
907 })
908 }
909
910 pub fn cursor_with_read_size(&self, rows_per_read: usize) -> Result<ParquetStoredTickCursor> {
911 Ok(ParquetStoredTickCursor {
912 inner: PartitionCursor::open(
913 self.inner.clone(),
914 rows_per_read,
915 read_stored_tick_partition,
916 )?,
917 })
918 }
919}
920
921pub struct ParquetStoredTickCursor {
923 inner: PartitionCursor<StoredTick>,
924}
925
926impl ParquetStoredTickCursor {
927 pub fn next_stored_tick_with_ordinal_cancellable<F>(
928 &mut self,
929 mut is_cancelled: F,
930 ) -> Result<Option<ParquetScannedRow<StoredTick>>>
931 where
932 F: FnMut() -> bool,
933 {
934 self.inner.next_row(&mut is_cancelled)
935 }
936
937 pub fn rows_per_read(&self) -> usize {
938 self.inner.rows_per_read()
939 }
940}
941
942#[derive(Debug, Clone)]
944pub struct ParquetPriceBarScan {
945 inner: PartitionScan,
946}
947
948impl ParquetPriceBarScan {
949 pub fn describe_cancellable<F>(
950 root: impl AsRef<Path>,
951 descriptor: &SeriesDescriptor,
952 bounds: ParquetScanBounds,
953 mut is_cancelled: F,
954 ) -> Result<Self>
955 where
956 F: FnMut() -> bool,
957 {
958 descriptor.validate()?;
959 let directory = root
960 .as_ref()
961 .join("price_bars")
962 .join(format!("exchange={}", descriptor.exchange))
963 .join(format!("symbol={}", descriptor.symbol))
964 .join(format!(
965 "timeframe_seconds={}",
966 descriptor.timeframe_seconds
967 ));
968 let discovery_bounds = ParquetScanBounds::default();
969 let mut inner = PartitionScan::describe(&directory, discovery_bounds, &mut is_cancelled)?;
970 inner.bounds = bounds.validate()?;
971 Ok(Self { inner })
972 }
973
974 pub fn cursor_with_read_size(&self, rows_per_read: usize) -> Result<ParquetPriceBarCursor> {
975 Ok(ParquetPriceBarCursor {
976 inner: PartitionCursor::open(
977 self.inner.clone(),
978 rows_per_read,
979 read_price_bar_partition,
980 )?,
981 })
982 }
983}
984
985pub struct ParquetPriceBarCursor {
987 inner: PartitionCursor<PriceBar>,
988}
989
990impl ParquetPriceBarCursor {
991 pub fn next_price_bar_with_ordinal_cancellable<F>(
992 &mut self,
993 mut is_cancelled: F,
994 ) -> Result<Option<ParquetScannedRow<PriceBar>>>
995 where
996 F: FnMut() -> bool,
997 {
998 self.inner.next_row(&mut is_cancelled)
999 }
1000
1001 pub fn rows_per_read(&self) -> usize {
1002 self.inner.rows_per_read()
1003 }
1004}
1005
1006impl ParquetStore {
1007 pub fn scan_ticks(
1009 &self,
1010 exchange: &str,
1011 symbol: &str,
1012 bounds: ParquetScanBounds,
1013 ) -> Result<ParquetTickCursor> {
1014 ParquetTickCursor::open(self.root_path(), exchange, symbol, bounds)
1015 }
1016
1017 pub fn scan_ticks_cancellable<F>(
1019 &self,
1020 exchange: &str,
1021 symbol: &str,
1022 bounds: ParquetScanBounds,
1023 is_cancelled: F,
1024 ) -> Result<ParquetTickCursor>
1025 where
1026 F: FnMut() -> bool,
1027 {
1028 ParquetTickCursor::open_cancellable(
1029 self.root_path(),
1030 exchange,
1031 symbol,
1032 bounds,
1033 is_cancelled,
1034 )
1035 }
1036
1037 pub fn scan_bars(
1039 &self,
1040 exchange: &str,
1041 symbol: &str,
1042 timeframe: &str,
1043 bounds: ParquetScanBounds,
1044 ) -> Result<ParquetBarCursor> {
1045 ParquetBarCursor::open(self.root_path(), exchange, symbol, timeframe, bounds)
1046 }
1047
1048 pub fn scan_bars_cancellable<F>(
1050 &self,
1051 exchange: &str,
1052 symbol: &str,
1053 timeframe: &str,
1054 bounds: ParquetScanBounds,
1055 is_cancelled: F,
1056 ) -> Result<ParquetBarCursor>
1057 where
1058 F: FnMut() -> bool,
1059 {
1060 ParquetBarCursor::open_cancellable(
1061 self.root_path(),
1062 exchange,
1063 symbol,
1064 timeframe,
1065 bounds,
1066 is_cancelled,
1067 )
1068 }
1069
1070 pub fn scan_stored_ticks_cancellable<F>(
1071 &self,
1072 exchange: &str,
1073 symbol: &str,
1074 bounds: ParquetScanBounds,
1075 rows_per_read: usize,
1076 is_cancelled: F,
1077 ) -> Result<ParquetStoredTickCursor>
1078 where
1079 F: FnMut() -> bool,
1080 {
1081 ParquetStoredTickScan::describe_cancellable(
1082 self.root_path(),
1083 exchange,
1084 symbol,
1085 bounds,
1086 is_cancelled,
1087 )?
1088 .cursor_with_read_size(rows_per_read)
1089 }
1090
1091 pub fn scan_price_bars_cancellable<F>(
1092 &self,
1093 descriptor: &SeriesDescriptor,
1094 bounds: ParquetScanBounds,
1095 rows_per_read: usize,
1096 is_cancelled: F,
1097 ) -> Result<ParquetPriceBarCursor>
1098 where
1099 F: FnMut() -> bool,
1100 {
1101 self.verify_price_bar_series(descriptor)?;
1102 ParquetPriceBarScan::describe_cancellable(
1103 self.root_path(),
1104 descriptor,
1105 bounds,
1106 is_cancelled,
1107 )?
1108 .cursor_with_read_size(rows_per_read)
1109 }
1110}
1111
1112fn list_partitions(
1113 directory: &Path,
1114 bounds: ParquetScanBounds,
1115 is_cancelled: &mut dyn FnMut() -> bool,
1116) -> Result<Vec<PartitionDescriptor>> {
1117 ensure_not_cancelled(is_cancelled)?;
1118 if !directory.exists() {
1119 return Ok(Vec::new());
1120 }
1121
1122 let from_date = bounds.from.map(|ts| ts.date());
1123 let to_date = bounds.to.map(|ts| ts.date());
1124 let mut paths = Vec::new();
1125 for entry in fs::read_dir(directory)? {
1126 ensure_not_cancelled(is_cancelled)?;
1127 let path = entry?.path();
1128 if path
1129 .extension()
1130 .is_none_or(|extension| extension != "parquet")
1131 {
1132 continue;
1133 }
1134
1135 let stem = path
1136 .file_stem()
1137 .and_then(|value| value.to_str())
1138 .ok_or_else(|| DataError::InvalidDatePartition(path.display().to_string()))?;
1139 let date = NaiveDate::parse_from_str(stem, "%Y-%m-%d")
1140 .map_err(|_| DataError::InvalidDatePartition(path.display().to_string()))?;
1141 if date.format("%Y-%m-%d").to_string() != stem {
1142 return Err(DataError::InvalidDatePartition(path.display().to_string()));
1143 }
1144 if from_date.is_some_and(|from| date < from) || to_date.is_some_and(|to| date > to) {
1145 continue;
1146 }
1147 paths.push((date, path));
1148 }
1149 paths.sort_by(|left, right| left.0.cmp(&right.0).then_with(|| left.1.cmp(&right.1)));
1150
1151 let mut partitions = Vec::with_capacity(paths.len());
1152 let mut source_row_base = 0u64;
1153 for (_, path) in paths {
1154 ensure_not_cancelled(is_cancelled)?;
1155 let descriptor = PartitionDescriptor::describe(path, source_row_base)?;
1156 source_row_base =
1157 source_row_base
1158 .checked_add(u64::try_from(descriptor.row_count).map_err(|_| {
1159 DataError::Other("Parquet partition row count exceeds u64".into())
1160 })?)
1161 .ok_or_else(|| DataError::Other("Parquet source row ordinal overflow".into()))?;
1162 partitions.push(descriptor);
1163 }
1164 ensure_not_cancelled(is_cancelled)?;
1165 Ok(partitions)
1166}
1167
1168fn validate_monotonic_rows<T: Timestamped>(
1169 mut previous: Option<NaiveDateTime>,
1170 rows: &[T],
1171 path: &Path,
1172) -> Result<()> {
1173 for row in rows {
1174 let current = row.timestamp();
1175 if let Some(previous) = previous
1176 && current < previous
1177 {
1178 return Err(DataError::NonMonotonicParquetData {
1179 path: path.display().to_string(),
1180 previous,
1181 current,
1182 });
1183 }
1184 previous = Some(current);
1185 }
1186 Ok(())
1187}
1188
1189fn ensure_not_cancelled(is_cancelled: &mut dyn FnMut() -> bool) -> Result<()> {
1190 if is_cancelled() {
1191 Err(DataError::Cancelled)
1192 } else {
1193 Ok(())
1194 }
1195}
1196
1197fn read_tick_partition(file: File, offset: usize, rows: usize) -> Result<Vec<Tick>> {
1198 let dataframe = ParquetReader::new(file)
1199 .with_slice(Some((offset, rows)))
1200 .finish()?;
1201 dataframe_to_ticks(&dataframe)
1202}
1203
1204fn read_bar_partition(file: File, offset: usize, rows: usize) -> Result<Vec<Bar>> {
1205 let dataframe = ParquetReader::new(file)
1206 .with_slice(Some((offset, rows)))
1207 .finish()?;
1208 dataframe_to_bars(&dataframe)
1209}
1210
1211fn read_stored_tick_partition(file: File, offset: usize, rows: usize) -> Result<Vec<StoredTick>> {
1212 let dataframe = ParquetReader::new(file)
1213 .with_slice(Some((offset, rows)))
1214 .finish()?;
1215 dataframe_to_stored_ticks(&dataframe)
1216}
1217
1218fn read_price_bar_partition(file: File, offset: usize, rows: usize) -> Result<Vec<PriceBar>> {
1219 let dataframe = ParquetReader::new(file)
1220 .with_slice(Some((offset, rows)))
1221 .finish()?;
1222 dataframe_to_price_bars(&dataframe)
1223}
1224
1225fn tick_has_valid_quote(tick: &Tick) -> bool {
1226 match (tick.bid, tick.ask) {
1227 (Some(bid), Some(ask)) => {
1228 bid.is_finite() && ask.is_finite() && bid > 0.0 && ask > 0.0 && bid <= ask
1229 }
1230 _ => false,
1231 }
1232}
1233
1234#[cfg(test)]
1235mod tests {
1236 use std::time::{SystemTime, UNIX_EPOCH};
1237
1238 use super::*;
1239 use crate::convert::ticks_to_dataframe;
1240 use crate::models::Timeframe;
1241
1242 fn ts(day: u32, hour: u32) -> NaiveDateTime {
1243 NaiveDate::from_ymd_opt(2026, 1, day)
1244 .unwrap()
1245 .and_hms_opt(hour, 0, 0)
1246 .unwrap()
1247 }
1248
1249 fn temp_root(name: &str) -> PathBuf {
1250 let nonce = SystemTime::now()
1251 .duration_since(UNIX_EPOCH)
1252 .unwrap()
1253 .as_nanos();
1254 std::env::temp_dir().join(format!(
1255 "qs-data-preprocess-{name}-{}-{nonce}",
1256 std::process::id()
1257 ))
1258 }
1259
1260 fn tick(at: NaiveDateTime) -> Tick {
1261 Tick {
1262 exchange: "test".into(),
1263 symbol: "EURUSD".into(),
1264 ts: at,
1265 bid: Some(1.0),
1266 ask: Some(1.1),
1267 last: None,
1268 volume: None,
1269 flags: None,
1270 }
1271 }
1272
1273 fn bar(at: NaiveDateTime) -> Bar {
1274 Bar {
1275 exchange: "test".into(),
1276 symbol: "EURUSD".into(),
1277 timeframe: Timeframe::M1,
1278 ts: at,
1279 open: 1.0,
1280 high: 1.2,
1281 low: 0.9,
1282 close: 1.1,
1283 tick_vol: 10,
1284 volume: 10,
1285 spread: 1,
1286 }
1287 }
1288
1289 #[test]
1290 fn tick_cursor_scans_partitions_in_ascending_bounded_order() {
1291 let root = temp_root("ticks");
1292 let store = ParquetStore::open(&root).unwrap();
1293 store
1294 .insert_ticks(&[tick(ts(2, 1)), tick(ts(1, 2)), tick(ts(1, 1))])
1295 .unwrap();
1296
1297 let mut cursor = store
1298 .scan_ticks(
1299 "test",
1300 "EURUSD",
1301 ParquetScanBounds::new(Some(ts(1, 2)), Some(ts(2, 1))),
1302 )
1303 .unwrap();
1304 assert_eq!(cursor.remaining_partitions(), 2);
1305 assert_eq!(cursor.next_tick().unwrap().unwrap().ts, ts(1, 2));
1306 assert_eq!(cursor.next_tick().unwrap().unwrap().ts, ts(2, 1));
1307 assert!(cursor.next_tick().unwrap().is_none());
1308
1309 fs::remove_dir_all(root).ok();
1310 }
1311
1312 #[test]
1313 fn bar_cursor_scans_ordered_date_partitions() {
1314 let root = temp_root("bars");
1315 let store = ParquetStore::open(&root).unwrap();
1316 store.insert_bars(&[bar(ts(2, 1)), bar(ts(1, 1))]).unwrap();
1317
1318 let mut cursor = store
1319 .scan_bars("test", "EURUSD", "1m", ParquetScanBounds::default())
1320 .unwrap();
1321 assert_eq!(cursor.next_bar().unwrap().unwrap().ts, ts(1, 1));
1322 assert_eq!(cursor.next_bar().unwrap().unwrap().ts, ts(2, 1));
1323 assert!(cursor.next_bar().unwrap().is_none());
1324
1325 fs::remove_dir_all(root).ok();
1326 }
1327
1328 #[test]
1329 fn cursor_reads_large_partitions_through_bounded_slices() {
1330 let root = temp_root("slices");
1331 let store = ParquetStore::open(&root).unwrap();
1332 store
1333 .insert_ticks(&[tick(ts(1, 1)), tick(ts(1, 2)), tick(ts(1, 3))])
1334 .unwrap();
1335 let mut cursor = ParquetTickCursor::open_with_read_size(
1336 &root,
1337 "test",
1338 "EURUSD",
1339 ParquetScanBounds::default(),
1340 1,
1341 )
1342 .unwrap();
1343
1344 assert_eq!(cursor.rows_per_read(), 1);
1345 assert_eq!(cursor.next_tick().unwrap().unwrap().ts, ts(1, 1));
1346 assert_eq!(cursor.next_tick().unwrap().unwrap().ts, ts(1, 2));
1347 assert_eq!(cursor.next_tick().unwrap().unwrap().ts, ts(1, 3));
1348 assert!(cursor.next_tick().unwrap().is_none());
1349 assert!(matches!(
1350 ParquetTickCursor::open_with_read_size(
1351 &root,
1352 "test",
1353 "EURUSD",
1354 ParquetScanBounds::default(),
1355 0,
1356 ),
1357 Err(DataError::InvalidScanReadSize)
1358 ));
1359
1360 fs::remove_dir_all(root).ok();
1361 }
1362
1363 #[test]
1364 fn cursor_cancellation_does_not_consume_the_partition() {
1365 let root = temp_root("cancel");
1366 let store = ParquetStore::open(&root).unwrap();
1367 store.insert_ticks(&[tick(ts(1, 1))]).unwrap();
1368 let mut cursor = store
1369 .scan_ticks("test", "EURUSD", ParquetScanBounds::default())
1370 .unwrap();
1371 let mut checks = 0;
1372
1373 let error = cursor
1374 .next_tick_cancellable(|| {
1375 checks += 1;
1376 checks == 3
1377 })
1378 .unwrap_err();
1379 assert!(matches!(error, DataError::Cancelled));
1380 assert_eq!(cursor.next_tick().unwrap().unwrap().ts, ts(1, 1));
1381
1382 fs::remove_dir_all(root).ok();
1383 }
1384
1385 #[test]
1386 fn cursor_rejects_non_monotonic_partition_rows_before_emitting_them() {
1387 let root = temp_root("monotonic");
1388 let path = root
1389 .join("ticks")
1390 .join("exchange=test")
1391 .join("symbol=EURUSD")
1392 .join("2026-01-01.parquet");
1393 fs::create_dir_all(path.parent().unwrap()).unwrap();
1394 let mut dataframe = ticks_to_dataframe(&[tick(ts(1, 2)), tick(ts(1, 1))]).unwrap();
1395 ParquetWriter::new(File::create(&path).unwrap())
1396 .finish(&mut dataframe)
1397 .unwrap();
1398
1399 let mut cursor =
1400 ParquetTickCursor::open(&root, "test", "EURUSD", ParquetScanBounds::default()).unwrap();
1401 assert!(matches!(
1402 cursor.next_tick(),
1403 Err(DataError::NonMonotonicParquetData {
1404 previous,
1405 current,
1406 ..
1407 }) if previous == ts(1, 2) && current == ts(1, 1)
1408 ));
1409
1410 fs::remove_dir_all(root).ok();
1411 }
1412
1413 #[test]
1414 fn cursor_reports_physical_ordinals_across_filtered_rows_and_partitions() {
1415 let root = temp_root("ordinals");
1416 let store = ParquetStore::open(&root).unwrap();
1417 store
1418 .insert_ticks(&[
1419 tick(ts(1, 1)),
1420 tick(ts(1, 2)),
1421 tick(ts(1, 3)),
1422 tick(ts(2, 1)),
1423 ])
1424 .unwrap();
1425 let mut cursor = ParquetTickCursor::open_with_read_size(
1426 &root,
1427 "test",
1428 "EURUSD",
1429 ParquetScanBounds::new(Some(ts(1, 2)), None),
1430 2,
1431 )
1432 .unwrap();
1433
1434 let rows = [
1435 cursor.next_tick_with_ordinal().unwrap().unwrap(),
1436 cursor.next_tick_with_ordinal().unwrap().unwrap(),
1437 cursor.next_tick_with_ordinal().unwrap().unwrap(),
1438 ];
1439 assert_eq!(
1440 rows.map(|row| (row.row.ts, row.source_row_ordinal)),
1441 [(ts(1, 2), 1), (ts(1, 3), 2), (ts(2, 1), 3)]
1442 );
1443
1444 fs::remove_dir_all(root).ok();
1445 }
1446
1447 #[test]
1448 fn running_cursor_fails_if_atomic_replacement_changes_its_partition() {
1449 let root = temp_root("replace-running");
1450 let store = ParquetStore::open(&root).unwrap();
1451 store
1452 .insert_ticks(&[tick(ts(1, 1)), tick(ts(1, 2)), tick(ts(1, 3))])
1453 .unwrap();
1454 let mut cursor = ParquetTickCursor::open_with_read_size(
1455 &root,
1456 "test",
1457 "EURUSD",
1458 ParquetScanBounds::default(),
1459 1,
1460 )
1461 .unwrap();
1462 assert_eq!(cursor.next_tick().unwrap().unwrap().ts, ts(1, 1));
1463 let partition_path = root
1464 .join("ticks")
1465 .join("exchange=test")
1466 .join("symbol=EURUSD")
1467 .join("2026-01-01.parquet");
1468 #[cfg(unix)]
1469 let original_inode = {
1470 use std::os::unix::fs::MetadataExt as _;
1471 fs::metadata(&partition_path).unwrap().ino()
1472 };
1473
1474 store.insert_ticks(&[tick(ts(1, 4))]).unwrap();
1475 #[cfg(unix)]
1476 {
1477 use std::os::unix::fs::MetadataExt as _;
1478 assert_ne!(original_inode, fs::metadata(&partition_path).unwrap().ino());
1479 }
1480 assert!(
1481 fs::read_dir(partition_path.parent().unwrap())
1482 .unwrap()
1483 .all(|entry| !entry
1484 .unwrap()
1485 .file_name()
1486 .to_string_lossy()
1487 .ends_with(".tmp"))
1488 );
1489 assert!(matches!(
1490 cursor.next_tick(),
1491 Err(DataError::ParquetPartitionChanged { .. })
1492 ));
1493
1494 fs::remove_dir_all(root).ok();
1495 }
1496
1497 #[test]
1498 fn described_scan_rejects_replacement_before_reopen() {
1499 let root = temp_root("replace-described");
1500 let store = ParquetStore::open(&root).unwrap();
1501 store.insert_ticks(&[tick(ts(1, 1))]).unwrap();
1502 let scan = ParquetTickScan::describe(&root, "test", "EURUSD", ParquetScanBounds::default())
1503 .unwrap();
1504
1505 store.insert_ticks(&[tick(ts(1, 2))]).unwrap();
1506 assert!(matches!(
1507 scan.cursor(),
1508 Err(DataError::ParquetPartitionChanged { .. })
1509 ));
1510
1511 fs::remove_dir_all(root).ok();
1512 }
1513
1514 #[test]
1515 fn upper_bound_stops_at_the_first_later_monotonic_row() {
1516 let root = temp_root("upper-bound");
1517 let store = ParquetStore::open(&root).unwrap();
1518 store
1519 .insert_ticks(&[
1520 tick(ts(1, 1)),
1521 tick(ts(1, 2)),
1522 tick(ts(1, 3)),
1523 tick(ts(1, 4)),
1524 ])
1525 .unwrap();
1526 let mut cursor = ParquetTickCursor::open_with_read_size(
1527 &root,
1528 "test",
1529 "EURUSD",
1530 ParquetScanBounds::new(None, Some(ts(1, 2))),
1531 3,
1532 )
1533 .unwrap();
1534
1535 assert_eq!(cursor.next_tick().unwrap().unwrap().ts, ts(1, 1));
1536 assert_eq!(cursor.next_tick().unwrap().unwrap().ts, ts(1, 2));
1537 assert!(cursor.next_tick().unwrap().is_none());
1538 assert_eq!(cursor.remaining_partitions(), 0);
1539
1540 fs::remove_dir_all(root).ok();
1541 }
1542
1543 #[test]
1544 fn reverse_chunks_find_the_latest_valid_tick_without_materializing_the_day() {
1545 let root = temp_root("latest-reverse");
1546 let store = ParquetStore::open(&root).unwrap();
1547 let mut invalid = tick(ts(1, 4));
1548 invalid.ask = None;
1549 store
1550 .insert_ticks(&[tick(ts(1, 1)), tick(ts(1, 2)), tick(ts(1, 3)), invalid])
1551 .unwrap();
1552 let scan = ParquetTickScan::describe(&root, "test", "EURUSD", ParquetScanBounds::default())
1553 .unwrap();
1554
1555 let latest = find_latest_row(
1556 &scan.inner,
1557 2,
1558 read_tick_partition,
1559 tick_has_valid_quote,
1560 &mut || false,
1561 )
1562 .unwrap()
1563 .unwrap();
1564 assert_eq!(latest.row.ts, ts(1, 3));
1565 assert_eq!(latest.source_row_ordinal, 2);
1566
1567 fs::remove_dir_all(root).ok();
1568 }
1569}