Skip to main content

data_preprocess/parser/
tick_csv.rs

1use std::path::Path;
2
3use chrono::FixedOffset;
4
5use crate::error::Result;
6use crate::models::{Tick, TickImportAudit};
7use crate::parser::parse_datetime_to_utc;
8
9/// Parse optional float field; empty or whitespace-only → None.
10fn parse_optional_f64(field: Option<&str>) -> Option<f64> {
11    field
12        .map(str::trim)
13        .filter(|s| !s.is_empty())
14        .and_then(|s| s.parse::<f64>().ok())
15}
16
17/// Parse optional integer field; empty or whitespace-only → None.
18fn parse_optional_i32(field: Option<&str>) -> Option<i32> {
19    field
20        .map(str::trim)
21        .filter(|s| !s.is_empty())
22        .and_then(|s| s.parse::<i32>().ok())
23}
24
25/// Parse a tab-delimited tick CSV into `Vec<Tick>`.
26///
27/// Malformed rows produce warnings instead of hard errors so partial
28/// files can still be imported.
29pub fn parse_tick_csv(
30    path: &Path,
31    exchange: &str,
32    symbol: &str,
33    source_offset: &FixedOffset,
34) -> Result<(Vec<Tick>, Vec<String>)> {
35    let (ticks, warnings, _) = parse_tick_csv_with_audit(path, exchange, symbol, source_offset)?;
36    Ok((ticks, warnings))
37}
38
39/// Parse legacy tick CSV and report the capabilities that survive the current import path.
40pub fn parse_tick_csv_with_audit(
41    path: &Path,
42    exchange: &str,
43    symbol: &str,
44    source_offset: &FixedOffset,
45) -> Result<(Vec<Tick>, Vec<String>, TickImportAudit)> {
46    let mut reader = csv::ReaderBuilder::new()
47        .delimiter(b'\t')
48        .has_headers(true)
49        .flexible(true)
50        .from_path(path)?;
51
52    let mut ticks = Vec::new();
53    let mut warnings = Vec::new();
54    let mut maximum_fractional_digits = 0_u8;
55
56    for (line_idx, record) in reader.records().enumerate() {
57        let record = match record {
58            Ok(r) => r,
59            Err(e) => {
60                warnings.push(format!("line {}: {}", line_idx + 2, e));
61                continue;
62            }
63        };
64
65        let date_str = record.get(0).unwrap_or("").trim();
66        let time_str = record.get(1).unwrap_or("").trim();
67        let bid = parse_optional_f64(record.get(2));
68        let ask = parse_optional_f64(record.get(3));
69        let last = parse_optional_f64(record.get(4));
70        let vol = parse_optional_f64(record.get(5));
71        let flags = parse_optional_i32(record.get(6));
72
73        let ts = match parse_datetime_to_utc(date_str, time_str, source_offset) {
74            Ok(t) => t,
75            Err(e) => {
76                warnings.push(format!("line {}: {}", line_idx + 2, e));
77                continue;
78            }
79        };
80
81        maximum_fractional_digits = maximum_fractional_digits.max(
82            time_str
83                .split_once('.')
84                .map_or(0, |(_, fraction)| fraction.len().min(9) as u8),
85        );
86        ticks.push(Tick {
87            exchange: exchange.to_string(),
88            symbol: symbol.to_string(),
89            ts,
90            bid,
91            ask,
92            last,
93            volume: vol,
94            flags,
95        });
96    }
97
98    let distinct_timestamps = ticks
99        .iter()
100        .map(|tick| tick.ts)
101        .collect::<std::collections::BTreeSet<_>>()
102        .len();
103    let audit = TickImportAudit {
104        parsed_rows: ticks.len(),
105        distinct_timestamps,
106        simultaneous_rows: ticks.len().saturating_sub(distinct_timestamps),
107        maximum_fractional_digits,
108        provider_sequence_available: false,
109        stable_import_ordinal_persisted: false,
110        legacy_timestamp_dedup: true,
111        exact_quote_path_capable: false,
112    };
113    Ok((ticks, warnings, audit))
114}