qs-data-preprocess 0.4.0

Historical market data storage and preprocessing CLI
Documentation
# data-preprocess

Historical market data storage and preprocessing CLI. Imports tick and OHLCV bar data from CSV files into Parquet (default) or DuckDB storage, with support for multiple exchanges, deduplication, querying, management, and bounded ascending Parquet cursors.

## Storage Backends

| Backend | Feature Flag | Default | Build Time | Description |
|---------|-------------|---------|------------|-------------|
| **Parquet + Polars** | `parquet` | ✅ Yes | ~30s | Hive-partitioned Parquet files. No C++ compilation. zstd compressed. |
| **DuckDB** | `duckdb-backend` | No | Native build, environment dependent | Embedded columnar database. Opt-in for SQL exploration. |

### Parquet Directory Layout (Hive-Style Partitioning)

```
{data_dir}/
├── ticks/
│   └── exchange={exchange}/
│       └── symbol={symbol}/
│           ├── 2026-01-15.parquet
│           └── 2026-01-16.parquet
└── bars/
    └── exchange={exchange}/
        └── symbol={symbol}/
            └── timeframe={timeframe}/
                ├── 2026-01-15.parquet
                └── 2026-01-16.parquet
```

Each legacy file covers one date for one exchange+symbol (or exchange+symbol+timeframe for bars). Files are sorted by timestamp ascending and compressed with zstd.

The additive library-only enhanced paths use `ordered_ticks/` and `price_bars/`. `StoredTick` persists a source ordinal and preserves equal timestamps; provider sequence is treated as duplicate identity only together with a source identity, and conflicting payloads reject. `PriceBar` stores nullable `tick_count`, actual `available_at`, and a verified `SeriesDescriptor`; unknown count is never represented by zero, one, volume, or another sentinel. Enhanced paths support both bounded materialization and fingerprinted slice-bounded cursors. `ParquetStoredTickCursor` and `ParquetPriceBarCursor` decode at most the admitted rows per read, validate partition generations and monotonic order between slices, and check cancellation without materializing the full dataset. Materialized compatibility queries still check row and resident-byte limits while accumulating partitions and include owned string storage in retained-byte accounting. Replay metadata keeps `available_at` separate from the nominal bucket timestamp, so delayed bars are ordered by actual publication and cannot replay their old range as a new executable bar. The legacy `Tick` and `Bar { tick_vol: i64 }` APIs and partitions remain unchanged. New enhanced writes require verified descriptors, while a caller may read metadata-free legacy data only as an explicitly unverified assertion.

## Quick Start

```bash
# Import tick data (symbol extracted from filename, UTC+2 default)
cargo run -p qs-data-preprocess -- \
  input tick --exchange ctrader BTCUSD_202602161900_202602210954.csv

# Import bar data (timeframe required)
cargo run -p qs-data-preprocess -- \
  input bar --exchange ctrader --timeframe 1m BTCUSD_M1_202602210045_202602211009.csv

# View statistics
cargo run -p qs-data-preprocess -- stats

# Query ticks
cargo run -p qs-data-preprocess -- \
  view tick --exchange ctrader --symbol BTCUSD --limit 20 --tail

# Query bars
cargo run -p qs-data-preprocess -- \
  view bar --exchange ctrader --symbol BTCUSD --timeframe 1m --limit 20

# Remove data
cargo run -p qs-data-preprocess -- \
  remove tick --exchange ctrader --symbol BTCUSD --from 2026-02-16 --to 2026-02-18
cargo run -p qs-data-preprocess -- remove symbol --exchange ctrader BTCUSD
cargo run -p qs-data-preprocess -- remove exchange binance

# Run tests (Parquet only, default)
cargo test -p qs-data-preprocess

# Run tests (both backends)
cargo test -p qs-data-preprocess --features duckdb-backend

# Use DuckDB backend at runtime
cargo run -p qs-data-preprocess --features duckdb-backend -- \
  --backend duckdb --db market_data.duckdb stats
```

## CLI Reference

```
data-preprocess [OPTIONS] <COMMAND>

Global options:
  --backend <parquet|duckdb>   Storage backend [default: parquet]
  --data-dir <PATH>            Root directory for Parquet files [default: market_data]
                               Also reads DATA_PREPROCESS_DIR env var
  --db <PATH>                  Path to DuckDB file (duckdb backend only) [default: market_data.duckdb]
                               Also reads DATA_PREPROCESS_DB env var

Commands:
  input          Import market data from CSV file(s)
  resample       Build stored bars from stored ticks
  remove         Remove data by exchange / symbol / type / date range
  stats          Show summary statistics
  view           Query and display stored data
```

### `resample`

```
data-preprocess resample [OPTIONS]

  -e, --exchange <EX>              Exchange name (REQUIRED)
  -s, --symbol <SYM>               Symbol to resample (REQUIRED)
  -t, --timeframe <TF>             Target timeframe: 1m, 3m, 5m, 15m, 30m, 1h, 4h, 1d, 1w (REQUIRED)
      --price-basis <BASIS>        bid | ask | mid [default: mid]
      --align-offset-seconds <N>   Bucket alignment offset [default: 0]
      --digits <N>                 Price decimal digits, used for the spread in points [default: 5]
      --from <DATETIME>            Inclusive start (YYYY-MM-DD or YYYY-MM-DDTHH:MM:SS)
      --to <DATETIME>              Inclusive end
      --flush-partial              Also write the final, still-incomplete bucket
```

Reads stored ticks and writes stored bars using the same bucket arithmetic, quote-acceptance rule, and accumulation the replay engine applies when it builds bars in memory. Bars produced here are therefore identical to the bars a tick-driven replay would derive from the same input, which lets a parameter search run over bars and a confirmation run over ticks without disagreeing.

Each bar records the average bid/ask spread observed while it formed, expressed in points for the given digit count, so bar-driven replay can execute against a realistic two-sided quote instead of a zero spread.

Replay treats a bar's close as the midpoint and rebuilds the quote as `close ± spread / 2`. Bars stored on the `mid` basis therefore reconstruct the original two-sided quote; bars stored on `bid` or `ask` reconstruct a quote whose centre sits half a spread away from the real midpoint, which shifts every stop, target, and pending-order trigger by that amount. Use `bid` or `ask` only to match a strategy series that analyses that side, and confirm on ticks.

Bars count accepted ticks in `tick_vol` and leave `volume` at zero, matching the tick-count volume the replay engine projects. Ticks must arrive in chronological order; any tick belonging to an already-closed bucket is dropped and reported at the end of the run.

An interval with no accepted tick produces no bar, so market closures simply have no bars. The final incomplete bucket is dropped unless `--flush-partial` is given, because the replay engine never emits an incomplete bar. Monthly bars have no fixed duration and are rejected. Weekly buckets align to the Unix epoch week unless an alignment offset moves them.

```bash
# Hourly bars for one symbol over a date range
data-preprocess --data-dir /data/forex resample \
  --exchange icmarkets --symbol EURUSD --timeframe 1h \
  --digits 5 --from 2025-01-01 --to 2026-09-04

# Daily bars that close at 22:00 UTC
data-preprocess --data-dir /data/forex resample \
  --exchange icmarkets --symbol EURUSD --timeframe 1d \
  --align-offset-seconds 79200 --digits 5
```

Resampling reads the Parquet tick store; it is unavailable on the DuckDB backend.

### `input tick`

```
data-preprocess input tick [OPTIONS] <FILES>...

  -e, --exchange <EX>      Exchange name (REQUIRED)
      --symbol <SYM>       Override symbol (default: from filename)
      --tz-offset <TZ>     Source timezone offset [default: +02:00]
```

### `input bar`

```
data-preprocess input bar [OPTIONS] <FILES>...

  -e, --exchange <EX>      Exchange name (REQUIRED)
  -t, --timeframe <TF>     Timeframe: 1m, 3m, 5m, 15m, 30m, 1h, 4h, 1d, 1w, 1M (REQUIRED)
      --symbol <SYM>       Override symbol (default: from filename)
      --tz-offset <TZ>     Source timezone offset [default: +02:00]
```

### `stats`

```
data-preprocess stats [--exchange <EX>] [--symbol <SYM>]
```

### `view tick` / `view bar`

```
data-preprocess view tick -e <EX> --symbol <SYM> [--from <DT>] [--to <DT>] [--limit N] [--tail] [--desc]
data-preprocess view bar  -e <EX> --symbol <SYM> -t <TF> [--from <DT>] [--to <DT>] [--limit N] [--tail] [--desc]
```

### `remove`

```
data-preprocess remove tick     -e <EX> --symbol <SYM> [--from <DT>] [--to <DT>]
data-preprocess remove bar      -e <EX> --symbol <SYM> -t <TF> [--from <DT>] [--to <DT>]
data-preprocess remove symbol   -e <EX> <SYMBOL>
data-preprocess remove exchange <EXCHANGE>
```

## Input CSV Formats

### Tick CSV

Tab-delimited, with header. Filename convention: `{SYMBOL}_*.csv`

```
<DATE>	<TIME>	<BID>	<ASK>	<LAST>	<VOLUME>	<FLAGS>
2026.02.16	19:00:00.083	67849.69	67861.69			6
```

### Bar CSV

Tab-delimited, with header. Filename convention: `{SYMBOL}_*.csv`

```
<DATE>	<TIME>	<OPEN>	<HIGH>	<LOW>	<CLOSE>	<TICKVOL>	<VOL>	<SPREAD>
2026.02.21	00:45:00	67932.44	67934.19	67888.89	67910.24	184	0	1200
```

## Data Conventions

- **Timeframes** use canonical storage labels; `1m` is one minute and the case-sensitive `1M` label is one month (`MN1` is also accepted as input)
- **Exchanges** are always stored lowercase (`ctrader`, `binance`)
- **Symbols** are always stored uppercase (`BTCUSD`, `EURUSD`)
- **Timestamps** are stored in UTC — source timezone is converted on import
- **Legacy deduplication** uses `(exchange, symbol, ts)` for ticks and `(exchange, symbol, timeframe, ts)` for bars
  - Parquet: read-merge-write per date partition file (bounded to one file per dedup operation)
  - DuckDB: `INSERT OR IGNORE` with UNIQUE constraints
  - `parse_tick_csv_with_audit` reports parsed/distinct/simultaneous row counts and the greatest observed fractional-second precision. The legacy CSV has no provider sequence and the legacy write path persists no source ordinal, so its audit always marks exact quote-path capability false. More than six fractional digits also warns that the Parquet microsecond representation may lose timestamp precision.
  - Existing legacy partitions remain readable but may already have discarded simultaneous rows; conversion cannot reconstruct them and must not mark them complete.
- **Enhanced ordered ticks** use persisted `(timestamp, source_ordinal)` order without timestamp deduplication. A provider sequence is duplicate identity only when paired with a nonempty source identity; without it, even identical simultaneous payloads remain distinct.

## Library Usage

### Parquet backend (default)

```rust
use data_preprocess::{ParquetScanBounds, ParquetStore, models::{QueryOpts, BarQueryOpts}};

let store = ParquetStore::open("market_data")?;

// Import ticks
let inserted = store.insert_ticks(&ticks)?;

// Query ticks. The from/to bounds remain inclusive.
let (ticks, total) = store.query_ticks(&QueryOpts {
    exchange: "ctrader".into(),
    symbol: "BTCUSD".into(),
    from: None,
    to: None,
    limit: 1000,
    tail: false,
    descending: false,
})?;

// Find the newest executable quote strictly before a timestamp.
let prior_tick = store.latest_valid_tick_before(
    "ctrader",
    "BTCUSD",
    decision_ts,
)?;

// The cancellable variant checks between partitions, file reads, and rows.
let prior_tick = store.latest_valid_tick_before_cancellable(
    "ctrader",
    "BTCUSD",
    decision_ts,
    || cancellation.is_cancelled(),
)?;

// Stream inclusive, ascending rows without materializing the full range.
let bounds = ParquetScanBounds::new(from, to);
let mut cursor = store.scan_ticks_cancellable(
    "ctrader",
    "BTCUSD",
    bounds,
    || cancellation.is_cancelled(),
)?;
while let Some(tick) = cursor.next_tick_cancellable(|| cancellation.is_cancelled())? {
    consume(tick);
}

// Stats
let stats = store.stats(None, None)?;

// Delete
let deleted = store.delete_ticks("ctrader", "BTCUSD", None, None)?;
let (tick_count, bar_count) = store.delete_symbol("ctrader", "BTCUSD")?;
```

### DuckDB backend (opt-in)

Requires `features = ["duckdb-backend"]` in your `Cargo.toml`.

```rust
use data_preprocess::{Database, models::{QueryOpts, BarQueryOpts}};

let db = Database::open("market_data.duckdb".as_ref())?;
let (ticks, total) = db.query_ticks(&QueryOpts {
    exchange: "ctrader".into(),
    symbol: "BTCUSD".into(),
    from: None,
    to: None,
    limit: 1000,
    tail: false,
    descending: false,
})?;
```

### Strict-before quote lookup

`ParquetStore::latest_valid_tick_before` and `latest_valid_tick_before_cancellable` search date partitions newest-to-oldest and return at most one tick. The cutoff is strict: a tick whose timestamp equals the requested timestamp is never returned. Missing, non-finite, non-positive, and crossed bid/ask quotes are skipped. This is separate from `query_ticks`, whose `from` and `to` range bounds remain inclusive.

### Bounded Parquet cursors

`ParquetTickCursor` and `ParquetBarCursor` enumerate date partitions in ascending order and decode bounded row slices, using 65,536 rows per read by default. Bounds are inclusive. Cursor reads preserve physical row order for equal timestamps and reject timestamp regressions instead of sorting or buffering a complete series.

Cancellable opens and reads check cancellation during partition discovery and before and after each bounded Polars read. One individual Polars read remains atomic because the backend has no interruption hook.

### Deterministic backtest feed ordering

`qs-backtest` keeps `MarketEvent` unchanged and attaches non-persisted `EventMetadata` in `FeedEvent`. Metadata contains `SeriesRoles`, `series_rank`, and `row_sequence`; deterministic feeds order events by `(timestamp, series_rank, row_sequence)`.

- `SeriesRoles::PRIMARY_AND_CONVERSION` represents one shared tick once with both roles.
- `EventBatchFeed` groups each ascending cursor into complete timestamp batches with stable source row sequence.
- `KWayMergeFeed` retains one batch of lookahead per series and emits complete merged timestamp batches in deterministic order.
- A primary bar and conversion tick can remain separate `FeedEvent`s in the same batch.
- Source failures, cancellation, or timestamp regressions fail the stream rather than emitting a partial timestamp batch.
- Materialized `VecFeed`, `ticks_to_feed_with_metadata`, `bars_to_feed_with_metadata`, and `merge_feeds` remain available for compatibility and tests.

### Consumer crates (models only, no backend)

For crates that only need the type definitions (`Tick`, `Bar`, `Timeframe`, `QueryOpts`), disable default features to avoid pulling in Polars:

```toml
[dependencies]
qs-data-preprocess = { path = "../data-preprocess", default-features = false }
```

## Feature Flags

| Feature | Dependencies | Use Case |
|---------|-------------|----------|
| `parquet` (default) | `polars` | Parquet read/write with Polars |
| `duckdb-backend` | `duckdb` | DuckDB embedded database |
| _(none)_ | — | Model types + parsers only (fastest build) |

## API Parity

Core storage operations have matching logical behavior and return types. Parquet-only streaming extensions are listed explicitly:

| Operation | `ParquetStore` | `Database` |
|-----------|---------------|-----------|
| Open | `open(root_dir)` | `open(file_path)` |
| Insert ticks | `insert_ticks(&[Tick]) -> usize` | `insert_ticks(&[Tick]) -> usize` |
| Insert bars | `insert_bars(&[Bar]) -> usize` | `insert_bars(&[Bar]) -> usize` |
| Query ticks | `query_ticks(&QueryOpts) -> (Vec<Tick>, u64)` | `query_ticks(&QueryOpts) -> (Vec<Tick>, u64)` |
| Latest valid tick strictly before timestamp | `latest_valid_tick_before(ex, sym, before) -> Option<Tick>` | Not available |
| Ascending bounded tick cursor | `scan_ticks(ex, sym, bounds) -> ParquetTickCursor` | Not available |
| Ascending bounded bar cursor | `scan_bars(ex, sym, tf, bounds) -> ParquetBarCursor` | Not available |
| Query bars | `query_bars(&BarQueryOpts) -> (Vec<Bar>, u64)` | `query_bars(&BarQueryOpts) -> (Vec<Bar>, u64)` |
| Delete ticks | `delete_ticks(ex, sym, from, to) -> usize` | `delete_ticks(ex, sym, from, to) -> usize` |
| Delete bars | `delete_bars(ex, sym, tf, from, to) -> usize` | `delete_bars(ex, sym, tf, from, to) -> usize` |
| Delete symbol | `delete_symbol(ex, sym) -> (usize, usize)` | `delete_symbol(ex, sym) -> (usize, usize)` |
| Delete exchange | `delete_exchange(ex) -> (usize, usize)` | `delete_exchange(ex) -> (usize, usize)` |
| Stats | `stats(ex?, sym?) -> Vec<StatRow>` | `stats(ex?, sym?) -> Vec<StatRow>` |
| Size | `total_size() -> Option<u64>` | `file_size() -> Option<u64>` |

## Used By

| Crate | How |
|-------|-----|
| `qs-backtest` | `default-features = false` — imports `Tick`, `Bar`, `Timeframe`, `QueryOpts`, `BarQueryOpts` model types only (no Polars/DuckDB). Uses `ticks_to_feed()` / `bars_to_feed()` converters. |
| `qs-backtest-server` | Full `parquet` feature - opens bounded tick/bar cursors for streaming FutureQuote replay, retains materialized Legacy queries, and uses `stats()` for the `list_symbols` RPC. |

## License

Licensed under either of

- Apache License, Version 2.0 ([LICENSE-APACHE](../../LICENSE-APACHE) or <http://www.apache.org/licenses/LICENSE-2.0>)
- MIT License ([LICENSE-MIT](../../LICENSE-MIT) or <http://opensource.org/licenses/MIT>)