Skip to main content

data_preprocess/
scanner.rs

1//! Streaming cursors for ascending Parquet tick and bar partitions.
2
3use 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
23/// Default maximum number of rows decoded by one cursor read.
24pub const DEFAULT_PARQUET_SCAN_ROWS: usize = 65_536;
25
26/// Inclusive timestamp bounds for a Parquet cursor.
27#[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/// A decoded row paired with its physical ordinal in the described scan.
53#[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/// An immutable description of bounded tick partitions.
486#[derive(Debug, Clone)]
487pub struct ParquetTickScan {
488    inner: PartitionScan,
489}
490
491impl ParquetTickScan {
492    /// Describe tick partitions and bind each path to its current file fingerprint.
493    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    /// Describe tick partitions with cooperative cancellation.
503    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    /// Open a cursor against the file generations bound by this description.
524    pub fn cursor(&self) -> Result<ParquetTickCursor> {
525        self.cursor_with_read_size(DEFAULT_PARQUET_SCAN_ROWS)
526    }
527
528    /// Open a cursor with a maximum number of rows decoded per read.
529    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    /// Find the latest valid quote inside the inclusive scan bounds.
536    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    /// Find the latest valid quote strictly before `before` and inside the scan bounds.
553    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
571/// Partition-at-a-time cursor over ascending tick rows.
572pub struct ParquetTickCursor {
573    inner: PartitionCursor<Tick>,
574}
575
576impl ParquetTickCursor {
577    /// Open a tick cursor directly from a Parquet data root.
578    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    /// Open a tick cursor with a maximum number of rows decoded per read.
588    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    /// Open a tick cursor while checking cancellation during partition discovery.
601    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    /// Open a bounded tick cursor with cancellable partition discovery.
622    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    /// Read the next ascending tick.
638    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    /// Read the next ascending tick with its physical source-row ordinal.
644    pub fn next_tick_with_ordinal(&mut self) -> Result<Option<ParquetScannedRow<Tick>>> {
645        self.next_tick_with_ordinal_cancellable(|| false)
646    }
647
648    /// Read the next ascending tick with cooperative cancellation.
649    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    /// Read the next tick and source ordinal with cooperative cancellation.
658    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    /// Number of date partitions still active or pending.
669    pub fn remaining_partitions(&self) -> usize {
670        self.inner.remaining_partitions()
671    }
672
673    /// Maximum number of rows decoded by one Parquet read.
674    pub fn rows_per_read(&self) -> usize {
675        self.inner.rows_per_read()
676    }
677}
678
679/// An immutable description of bounded bar partitions.
680#[derive(Debug, Clone)]
681pub struct ParquetBarScan {
682    inner: PartitionScan,
683}
684
685impl ParquetBarScan {
686    /// Describe bar partitions and bind each path to its current file fingerprint.
687    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    /// Describe bar partitions with cooperative cancellation.
698    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    /// Open a cursor against the file generations bound by this description.
721    pub fn cursor(&self) -> Result<ParquetBarCursor> {
722        self.cursor_with_read_size(DEFAULT_PARQUET_SCAN_ROWS)
723    }
724
725    /// Open a cursor with a maximum number of rows decoded per read.
726    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    /// Find the latest bar with a valid close inside the inclusive scan bounds.
733    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
750/// Partition-at-a-time cursor over ascending bar rows.
751pub struct ParquetBarCursor {
752    inner: PartitionCursor<Bar>,
753}
754
755impl ParquetBarCursor {
756    /// Open a bar cursor directly from a Parquet data root.
757    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    /// Open a bar cursor with a maximum number of rows decoded per read.
775    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    /// Open a bar cursor while checking cancellation during partition discovery.
795    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    /// Open a bounded bar cursor with cancellable partition discovery.
818    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    /// Read the next ascending bar.
842    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    /// Read the next ascending bar with its physical source-row ordinal.
848    pub fn next_bar_with_ordinal(&mut self) -> Result<Option<ParquetScannedRow<Bar>>> {
849        self.next_bar_with_ordinal_cancellable(|| false)
850    }
851
852    /// Read the next ascending bar with cooperative cancellation.
853    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    /// Read the next bar and source ordinal with cooperative cancellation.
862    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    /// Number of date partitions still active or pending.
873    pub fn remaining_partitions(&self) -> usize {
874        self.inner.remaining_partitions()
875    }
876
877    /// Maximum number of rows decoded by one Parquet read.
878    pub fn rows_per_read(&self) -> usize {
879        self.inner.rows_per_read()
880    }
881}
882
883/// Immutable description of persisted ordered-tick partitions.
884#[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
921/// Slice-bounded cursor over persisted ordered ticks.
922pub 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/// Immutable description of availability-ordered price-bar partitions.
943#[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
985/// Slice-bounded cursor whose order is each price bar's actual availability.
986pub 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    /// Create an ascending tick cursor over this store.
1008    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    /// Create an ascending tick cursor with cancellable partition discovery.
1018    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    /// Create an ascending bar cursor over this store.
1038    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    /// Create an ascending bar cursor with cancellable partition discovery.
1049    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}