use super::record::{ContractQuote, ContractSide, ExpirationRecord, QuoteRow, SnapshotRecord};
use crate::utils::ChainError;
use chrono::{DateTime, Utc};
use clickhouse::Row;
use positive::Positive;
use rust_decimal::Decimal;
use serde::{Deserialize, Serialize};
use uuid::Uuid;
pub(crate) const DECIMAL_SCALE: u32 = 28;
pub(crate) const SNAPSHOTS_TABLE: &str = "simulation_snapshots";
pub(crate) const QUOTES_TABLE: &str = "simulation_option_quotes";
#[derive(Debug, Clone, PartialEq, Row, Serialize, Deserialize)]
pub(crate) struct SnapshotMetaRow {
pub(crate) simulation_id: String,
pub(crate) simulation_generation: u64,
pub(crate) step: u64,
pub(crate) snapshot_id: String,
pub(crate) simulated_at: i64,
pub(crate) symbol: String,
pub(crate) underlying_price: i128,
pub(crate) base_volatility: i128,
pub(crate) quote_count: u64,
pub(crate) expiration_count: u64,
pub(crate) complete: bool,
pub(crate) inserted_at_ms: u64,
}
#[derive(Debug, Clone, PartialEq, Row, Serialize, Deserialize)]
pub(crate) struct OptionQuoteRow {
pub(crate) simulation_id: String,
pub(crate) simulation_generation: u64,
pub(crate) step: u64,
pub(crate) expires_at: i64,
pub(crate) strike: i128,
pub(crate) snapshot_id: String,
pub(crate) simulated_at: i64,
pub(crate) symbol: String,
pub(crate) days_to_expiration: i128,
pub(crate) labels: Vec<String>,
pub(crate) implied_volatility: i128,
pub(crate) call_bid: Option<i128>,
pub(crate) call_ask: Option<i128>,
pub(crate) call_mid: Option<i128>,
pub(crate) put_bid: Option<i128>,
pub(crate) put_ask: Option<i128>,
pub(crate) put_mid: Option<i128>,
pub(crate) delta_call: Option<i128>,
pub(crate) delta_put: Option<i128>,
pub(crate) gamma: Option<i128>,
pub(crate) inserted_at_ms: u64,
}
#[derive(Debug, Clone, PartialEq, Row, Deserialize)]
pub(crate) struct SnapshotMetaReadRow {
pub(crate) step: u64,
pub(crate) snapshot_id: String,
pub(crate) simulated_at: i64,
pub(crate) symbol: String,
pub(crate) underlying_price: i128,
pub(crate) base_volatility: i128,
pub(crate) quote_count: u64,
}
#[derive(Debug, Clone, PartialEq, Row, Deserialize)]
pub(crate) struct QuoteReadRow {
pub(crate) step: u64,
pub(crate) expires_at: i64,
pub(crate) days_to_expiration: i128,
pub(crate) labels: Vec<String>,
pub(crate) strike: i128,
pub(crate) implied_volatility: i128,
pub(crate) call_bid: Option<i128>,
pub(crate) call_ask: Option<i128>,
pub(crate) call_mid: Option<i128>,
pub(crate) put_bid: Option<i128>,
pub(crate) put_ask: Option<i128>,
pub(crate) put_mid: Option<i128>,
pub(crate) delta_call: Option<i128>,
pub(crate) delta_put: Option<i128>,
pub(crate) gamma: Option<i128>,
}
#[derive(Debug, Clone, PartialEq, Row, Deserialize)]
pub(crate) struct ContractReadRow {
pub(crate) step: u64,
pub(crate) simulated_at: i64,
pub(crate) expires_at: i64,
pub(crate) days_to_expiration: i128,
pub(crate) strike: i128,
pub(crate) implied_volatility: i128,
pub(crate) bid: Option<i128>,
pub(crate) ask: Option<i128>,
pub(crate) mid: Option<i128>,
pub(crate) delta: Option<i128>,
pub(crate) gamma: Option<i128>,
}
pub(crate) fn to_storage_decimal(value: Decimal, field: &str) -> Result<i128, ChainError> {
let shift = DECIMAL_SCALE.checked_sub(value.scale()).ok_or_else(|| {
ChainError::ClickHouseError(format!(
"{field} has scale {} beyond the storable {DECIMAL_SCALE}",
value.scale()
))
})?;
let factor = 10_i128.checked_pow(shift).ok_or_else(|| {
ChainError::ClickHouseError(format!("{field} needs an unrepresentable scale factor"))
})?;
value.mantissa().checked_mul(factor).ok_or_else(|| {
ChainError::ClickHouseError(format!(
"{field} value {value} does not fit a Decimal(38, {DECIMAL_SCALE}) column"
))
})
}
pub(crate) fn from_storage_decimal(raw: i128, field: &str) -> Result<Decimal, ChainError> {
let mut mantissa = raw;
let mut scale = DECIMAL_SCALE;
while scale > 0 && mantissa % 10 == 0 {
mantissa /= 10;
scale -= 1;
}
Decimal::try_from_i128_with_scale(mantissa, scale).map_err(|error| {
ChainError::ClickHouseError(format!("{field} holds an unreadable decimal: {error}"))
})
}
pub(crate) fn to_storage_positive(value: Positive, field: &str) -> Result<i128, ChainError> {
to_storage_decimal(value.to_dec(), field)
}
pub(crate) fn from_storage_positive(raw: i128, field: &str) -> Result<Positive, ChainError> {
let value = from_storage_decimal(raw, field)?;
Positive::new_decimal(value).map_err(|error| {
ChainError::ClickHouseError(format!("{field} holds a non-positive value: {error}"))
})
}
fn to_storage_optional(value: Option<Decimal>, field: &str) -> Result<Option<i128>, ChainError> {
value
.map(|value| to_storage_decimal(value, field))
.transpose()
}
fn to_storage_optional_positive(
value: Option<Positive>,
field: &str,
) -> Result<Option<i128>, ChainError> {
value
.map(|value| to_storage_positive(value, field))
.transpose()
}
fn from_storage_optional(raw: Option<i128>, field: &str) -> Result<Option<Decimal>, ChainError> {
raw.map(|raw| from_storage_decimal(raw, field)).transpose()
}
fn from_storage_optional_positive(
raw: Option<i128>,
field: &str,
) -> Result<Option<Positive>, ChainError> {
raw.map(|raw| from_storage_positive(raw, field)).transpose()
}
pub(crate) fn to_storage_instant(value: DateTime<Utc>, field: &str) -> Result<i64, ChainError> {
value.timestamp_nanos_opt().ok_or_else(|| {
ChainError::ClickHouseError(format!(
"{field} instant {value} is outside the storable 1678-2262 range"
))
})
}
#[must_use]
pub(crate) fn from_storage_instant(raw: i64) -> DateTime<Utc> {
DateTime::from_timestamp_nanos(raw)
}
pub(crate) fn meta_row(
record: &SnapshotRecord,
inserted_at_ms: u64,
) -> Result<SnapshotMetaRow, ChainError> {
Ok(SnapshotMetaRow {
simulation_id: record.simulation.to_string(),
simulation_generation: record.generation,
step: to_storage_count(record.step, "step")?,
snapshot_id: record.snapshot_id().to_string(),
simulated_at: to_storage_instant(record.simulated_at, "simulated_at")?,
symbol: record.symbol.clone(),
underlying_price: to_storage_positive(record.spot, "underlying_price")?,
base_volatility: to_storage_positive(record.base_volatility, "base_volatility")?,
quote_count: to_storage_count(record.quote_count(), "quote_count")?,
expiration_count: to_storage_count(record.expirations.len(), "expiration_count")?,
complete: true,
inserted_at_ms,
})
}
pub(crate) fn quote_rows(
record: &SnapshotRecord,
inserted_at_ms: u64,
) -> Result<Vec<OptionQuoteRow>, ChainError> {
let simulation_id = record.simulation.to_string();
let snapshot_id = record.snapshot_id().to_string();
let simulated_at = to_storage_instant(record.simulated_at, "simulated_at")?;
let step = to_storage_count(record.step, "step")?;
let mut rows = Vec::with_capacity(record.quote_count());
for expiration in &record.expirations {
let expires_at = to_storage_instant(expiration.expires_at, "expires_at")?;
let days_to_expiration =
to_storage_positive(expiration.days_to_expiration, "days_to_expiration")?;
for quote in &expiration.quotes {
rows.push(OptionQuoteRow {
simulation_id: simulation_id.clone(),
simulation_generation: record.generation,
step,
expires_at,
strike: to_storage_positive(quote.strike, "strike")?,
snapshot_id: snapshot_id.clone(),
simulated_at,
symbol: record.symbol.clone(),
days_to_expiration,
labels: expiration.labels.clone(),
implied_volatility: to_storage_positive(
quote.implied_volatility,
"implied_volatility",
)?,
call_bid: to_storage_optional_positive(quote.call_bid, "call_bid")?,
call_ask: to_storage_optional_positive(quote.call_ask, "call_ask")?,
call_mid: to_storage_optional_positive(quote.call_mid, "call_mid")?,
put_bid: to_storage_optional_positive(quote.put_bid, "put_bid")?,
put_ask: to_storage_optional_positive(quote.put_ask, "put_ask")?,
put_mid: to_storage_optional_positive(quote.put_mid, "put_mid")?,
delta_call: to_storage_optional(quote.delta_call, "delta_call")?,
delta_put: to_storage_optional(quote.delta_put, "delta_put")?,
gamma: to_storage_optional(quote.gamma, "gamma")?,
inserted_at_ms,
});
}
}
Ok(rows)
}
pub(crate) fn record_from_rows(
simulation: Uuid,
generation: u64,
meta: &SnapshotMetaReadRow,
quotes: &[QuoteReadRow],
) -> Result<SnapshotRecord, ChainError> {
let step = usize::try_from(meta.step).map_err(|_| {
ChainError::ClickHouseError(format!("step {} is not addressable here", meta.step))
})?;
let expected_id = super::record::snapshot_id(simulation, generation, step);
if meta.snapshot_id != expected_id.to_string() {
return Err(ChainError::ClickHouseError(format!(
"snapshot {simulation}/{generation}/{step} is stored under identity {} instead of \
{expected_id}",
meta.snapshot_id
)));
}
let mut expirations: Vec<ExpirationRecord> = Vec::new();
for row in quotes {
let expires_at = from_storage_instant(row.expires_at);
let quote = quote_from_row(row)?;
match expirations.last_mut() {
Some(current) if current.expires_at == expires_at => current.quotes.push(quote),
_ => expirations.push(ExpirationRecord {
expires_at,
days_to_expiration: from_storage_positive(
row.days_to_expiration,
"days_to_expiration",
)?,
labels: row.labels.clone(),
quotes: vec![quote],
}),
}
}
Ok(SnapshotRecord {
simulation,
generation,
step,
simulated_at: from_storage_instant(meta.simulated_at),
symbol: meta.symbol.clone(),
spot: from_storage_positive(meta.underlying_price, "underlying_price")?,
base_volatility: from_storage_positive(meta.base_volatility, "base_volatility")?,
expirations,
})
}
fn quote_from_row(row: &QuoteReadRow) -> Result<QuoteRow, ChainError> {
Ok(QuoteRow {
strike: from_storage_positive(row.strike, "strike")?,
implied_volatility: from_storage_positive(row.implied_volatility, "implied_volatility")?,
call_bid: from_storage_optional_positive(row.call_bid, "call_bid")?,
call_ask: from_storage_optional_positive(row.call_ask, "call_ask")?,
call_mid: from_storage_optional_positive(row.call_mid, "call_mid")?,
put_bid: from_storage_optional_positive(row.put_bid, "put_bid")?,
put_ask: from_storage_optional_positive(row.put_ask, "put_ask")?,
put_mid: from_storage_optional_positive(row.put_mid, "put_mid")?,
delta_call: from_storage_optional(row.delta_call, "delta_call")?,
delta_put: from_storage_optional(row.delta_put, "delta_put")?,
gamma: from_storage_optional(row.gamma, "gamma")?,
})
}
pub(crate) fn contract_quote_from_row(
row: &ContractReadRow,
side: ContractSide,
) -> Result<ContractQuote, ChainError> {
Ok(ContractQuote {
step: usize::try_from(row.step).map_err(|_| {
ChainError::ClickHouseError(format!("step {} is not addressable here", row.step))
})?,
simulated_at: from_storage_instant(row.simulated_at),
expires_at: from_storage_instant(row.expires_at),
days_to_expiration: from_storage_positive(row.days_to_expiration, "days_to_expiration")?,
strike: from_storage_positive(row.strike, "strike")?,
side,
implied_volatility: from_storage_positive(row.implied_volatility, "implied_volatility")?,
bid: from_storage_optional_positive(row.bid, "bid")?,
ask: from_storage_optional_positive(row.ask, "ask")?,
mid: from_storage_optional_positive(row.mid, "mid")?,
delta: from_storage_optional(row.delta, "delta")?,
gamma: from_storage_optional(row.gamma, "gamma")?,
})
}
fn to_storage_count(value: usize, field: &str) -> Result<u64, ChainError> {
u64::try_from(value).map_err(|_| ChainError::Validation {
field: field.to_string(),
reason: format!("{value} does not fit a UInt64 column"),
})
}
#[cfg(test)]
mod tests {
use super::*;
use chrono::{TimeZone, Timelike};
use positive::pos_or_panic;
use rust_decimal_macros::dec;
use std::str::FromStr;
fn instant(day: u32) -> DateTime<Utc> {
match Utc.with_ymd_and_hms(2026, 1, day, 14, 30, 0).single() {
Some(instant) => instant,
None => panic!("the test instant must be valid"),
}
}
fn quote(strike: f64) -> QuoteRow {
let long = match Decimal::from_str("1.2345678901234567890123456789") {
Ok(value) => value,
Err(error) => panic!("the fixture decimal must parse: {error}"),
};
let long_positive = match Positive::new_decimal(long) {
Ok(value) => value,
Err(error) => panic!("the fixture premium must be positive: {error}"),
};
QuoteRow::new(pos_or_panic!(strike), pos_or_panic!(0.185))
.with_call(
Some(long_positive),
Some(pos_or_panic!(1.3)),
Some(pos_or_panic!(1.25)),
Some(dec!(0.5123)),
)
.with_put(None, Some(pos_or_panic!(1.1)), None, Some(dec!(-0.4877)))
.with_gamma(Some(dec!(0.00312345)))
}
fn expiration(day: u32, strikes: &[f64]) -> ExpirationRecord {
ExpirationRecord::new(
instant(day),
pos_or_panic!(f64::from(day)),
vec!["weeklies".to_string(), "zero_dte".to_string()],
strikes.iter().copied().map(quote).collect(),
)
}
fn record() -> SnapshotRecord {
SnapshotRecord::new(
Uuid::from_u128(42),
3,
7,
instant(5),
"SPX".to_string(),
pos_or_panic!(5000.25),
pos_or_panic!(0.18),
vec![
expiration(6, &[4975.0, 5000.0, 5025.0]),
expiration(9, &[4975.0, 5000.0]),
],
)
}
fn read_rows(record: &SnapshotRecord) -> (SnapshotMetaReadRow, Vec<QuoteReadRow>) {
let meta = match meta_row(record, 1_700_000_000_000) {
Ok(row) => row,
Err(error) => panic!("the fixture must convert: {error}"),
};
let quotes = match quote_rows(record, 1_700_000_000_000) {
Ok(rows) => rows,
Err(error) => panic!("the fixture must convert: {error}"),
};
let meta = SnapshotMetaReadRow {
step: meta.step,
snapshot_id: meta.snapshot_id,
simulated_at: meta.simulated_at,
symbol: meta.symbol,
underlying_price: meta.underlying_price,
base_volatility: meta.base_volatility,
quote_count: meta.quote_count,
};
let quotes = quotes
.into_iter()
.map(|row| QuoteReadRow {
step: row.step,
expires_at: row.expires_at,
days_to_expiration: row.days_to_expiration,
labels: row.labels,
strike: row.strike,
implied_volatility: row.implied_volatility,
call_bid: row.call_bid,
call_ask: row.call_ask,
call_mid: row.call_mid,
put_bid: row.put_bid,
put_ask: row.put_ask,
put_mid: row.put_mid,
delta_call: row.delta_call,
delta_put: row.delta_put,
gamma: row.gamma,
})
.collect();
(meta, quotes)
}
#[test]
fn test_a_decimal_round_trips_exactly() {
let values = [
dec!(0),
dec!(5000.25),
dec!(-0.4877),
dec!(0.00000000000000000000000001),
match Decimal::from_str("1.2345678901234567890123456789") {
Ok(value) => value,
Err(error) => panic!("the fixture decimal must parse: {error}"),
},
];
for value in values {
let raw = match to_storage_decimal(value, "test") {
Ok(raw) => raw,
Err(error) => panic!("{value} must be storable: {error}"),
};
match from_storage_decimal(raw, "test") {
Ok(read) => assert_eq!(read, value, "{value} did not survive the column"),
Err(error) => panic!("{value} must be readable: {error}"),
}
}
}
#[test]
fn test_the_scale_matches_the_decimal_maximum() {
assert_eq!(DECIMAL_SCALE, 28);
}
#[test]
fn test_an_oversized_decimal_is_rejected() {
let huge = match Decimal::from_str("100000000000") {
Ok(value) => value,
Err(error) => panic!("the fixture decimal must parse: {error}"),
};
match to_storage_decimal(huge, "underlying_price") {
Err(ChainError::ClickHouseError(message)) => {
assert!(message.contains("underlying_price"), "{message}");
}
other => panic!("expected a ClickHouse error, got {other:?}"),
}
}
#[test]
fn test_an_instant_round_trips_to_the_nanosecond() {
let precise = match instant(5).with_nanosecond(123_456_789) {
Some(value) => value,
None => panic!("the fixture nanosecond must be valid"),
};
let raw = match to_storage_instant(precise, "simulated_at") {
Ok(raw) => raw,
Err(error) => panic!("the instant must be storable: {error}"),
};
assert_eq!(from_storage_instant(raw), precise);
}
#[test]
fn test_a_snapshot_flattens_into_its_quote_count() {
let record = record();
let rows = match quote_rows(&record, 1) {
Ok(rows) => rows,
Err(error) => panic!("the record must convert: {error}"),
};
assert_eq!(rows.len(), record.quote_count());
assert_eq!(rows.len(), 5);
for pair in rows.windows(2) {
if let [left, right] = pair {
assert!(
(left.expires_at, left.strike) < (right.expires_at, right.strike),
"rows must be written in the table's sorting order"
);
}
}
}
#[test]
fn test_every_quote_row_carries_the_snapshot_identity() {
let record = record();
let expected = record.snapshot_id().to_string();
let rows = match quote_rows(&record, 9) {
Ok(rows) => rows,
Err(error) => panic!("the record must convert: {error}"),
};
for row in &rows {
assert_eq!(row.snapshot_id, expected);
assert_eq!(row.simulation_id, record.simulation.to_string());
assert_eq!(row.simulation_generation, record.generation);
assert_eq!(row.inserted_at_ms, 9);
}
}
#[test]
fn test_the_marker_carries_the_expected_counts() {
let record = record();
let meta = match meta_row(&record, 5) {
Ok(row) => row,
Err(error) => panic!("the record must convert: {error}"),
};
assert!(meta.complete);
assert_eq!(meta.quote_count, 5);
assert_eq!(meta.expiration_count, 2);
assert_eq!(meta.snapshot_id, record.snapshot_id().to_string());
assert_eq!(meta.inserted_at_ms, 5);
}
#[test]
fn test_a_snapshot_round_trips_through_the_row_types() {
let original = record();
let (meta, quotes) = read_rows(&original);
match record_from_rows(original.simulation, original.generation, &meta, "es) {
Ok(reconstructed) => assert_eq!(reconstructed, original),
Err(error) => panic!("the snapshot must reconstruct: {error}"),
}
}
#[test]
fn test_an_empty_snapshot_round_trips() {
let mut original = record();
original.expirations.clear();
let (meta, quotes) = read_rows(&original);
assert!(quotes.is_empty());
match record_from_rows(original.simulation, original.generation, &meta, "es) {
Ok(reconstructed) => assert_eq!(reconstructed, original),
Err(error) => panic!("the snapshot must reconstruct: {error}"),
}
}
#[test]
fn test_a_foreign_snapshot_identity_is_refused() {
let original = record();
let (mut meta, quotes) = read_rows(&original);
meta.snapshot_id = Uuid::from_u128(1).to_string();
match record_from_rows(original.simulation, original.generation, &meta, "es) {
Err(ChainError::ClickHouseError(message)) => {
assert!(message.contains("stored under identity"), "{message}");
}
other => panic!("expected a ClickHouse error, got {other:?}"),
}
}
#[test]
fn test_a_contract_row_rebuilds_its_side() {
let row = ContractReadRow {
step: 7,
simulated_at: match to_storage_instant(instant(5), "simulated_at") {
Ok(raw) => raw,
Err(error) => panic!("the fixture instant must convert: {error}"),
},
expires_at: match to_storage_instant(instant(9), "expires_at") {
Ok(raw) => raw,
Err(error) => panic!("the fixture instant must convert: {error}"),
},
days_to_expiration: match to_storage_decimal(dec!(4.0), "days_to_expiration") {
Ok(raw) => raw,
Err(error) => panic!("the fixture decimal must convert: {error}"),
},
strike: match to_storage_decimal(dec!(5000), "strike") {
Ok(raw) => raw,
Err(error) => panic!("the fixture decimal must convert: {error}"),
},
implied_volatility: match to_storage_decimal(dec!(0.18), "implied_volatility") {
Ok(raw) => raw,
Err(error) => panic!("the fixture decimal must convert: {error}"),
},
bid: None,
ask: None,
mid: None,
delta: match to_storage_decimal(dec!(-0.4877), "delta") {
Ok(raw) => Some(raw),
Err(error) => panic!("the fixture decimal must convert: {error}"),
},
gamma: None,
};
match contract_quote_from_row(&row, ContractSide::Put) {
Ok(quote) => {
assert_eq!(quote.side, ContractSide::Put);
assert_eq!(quote.step, 7);
assert_eq!(quote.strike, pos_or_panic!(5000.0));
assert_eq!(quote.delta, Some(dec!(-0.4877)));
assert_eq!(quote.bid, None);
}
Err(error) => panic!("the contract row must rebuild: {error}"),
}
}
}