Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
168 changes: 167 additions & 1 deletion crates/engine/src/chain.rs
Original file line number Diff line number Diff line change
Expand Up @@ -136,6 +136,16 @@ pub async fn fetch_chains(client: &Client, now: OffsetDateTime) -> Result<AssetC
cutoff_ms, "fetched latest ticks for chain build"
);

Ok(assemble_chains(rows, now))
}

/// Fold a flat list of [`ChainRow`]s into per-asset, per-expiry chains.
///
/// Pulled out of [`fetch_chains`] so the row-→-tree logic can be tested
/// without a `ClickHouse` connection. Unknown assets / kinds and
/// non-finite strikes are dropped silently (the warn-log path lives
/// here too, so the test surface stays honest).
fn assemble_chains(rows: Vec<ChainRow>, now: OffsetDateTime) -> AssetChains {
let mut out: AssetChains = HashMap::new();
for row in rows {
let Some(asset) = parse_asset(&row.asset) else {
Expand Down Expand Up @@ -191,7 +201,7 @@ pub async fn fetch_chains(client: &Client, now: OffsetDateTime) -> Result<AssetC
let _ = (row.bid, row.ask, row.underlying); // currently unused at chain level; reserved for future filters
}

Ok(out)
out
}

fn parse_asset(s: &str) -> Option<Asset> {
Expand Down Expand Up @@ -263,4 +273,160 @@ mod tests {
assert_eq!(parse_kind("put"), Some(OptionKind::Put));
assert_eq!(parse_kind("straddle"), None);
}

// ------- assemble_chains tests ------------------------------------

fn row(
asset: &str,
expiry: OffsetDateTime,
strike: f64,
kind: &str,
mid: Option<f64>,
iv: Option<f64>,
) -> ChainRow {
ChainRow {
asset: asset.to_string(),
expiry,
strike,
kind: kind.to_string(),
bid: None,
ask: None,
mid,
iv,
underlying: 100_000.0,
}
}

#[test]
fn assemble_drops_unknown_asset() {
let now = datetime!(2026-05-25 00:00:00 UTC);
let exp = now + time::Duration::days(7);
let rows = vec![
row("btc", exp, 100.0, "call", Some(1.0), Some(0.5)),
row("doge", exp, 100.0, "call", Some(1.0), Some(0.5)),
];
let out = assemble_chains(rows, now);
assert!(out.contains_key(&Asset::Btc));
assert_eq!(out.len(), 1);
}

#[test]
fn assemble_drops_unknown_kind() {
let now = datetime!(2026-05-25 00:00:00 UTC);
let exp = now + time::Duration::days(7);
let rows = vec![
row("btc", exp, 100.0, "call", Some(1.0), Some(0.5)),
row("btc", exp, 100.0, "straddle", Some(1.0), Some(0.5)),
];
let out = assemble_chains(rows, now);
let chain = &out[&Asset::Btc][&exp];
assert_eq!(chain.legs.len(), 1);
assert!(chain.legs[0].call_mid_usd.is_some());
assert!(chain.legs[0].put_mid_usd.is_none());
}

#[test]
fn assemble_drops_non_finite_or_non_positive_strike() {
let now = datetime!(2026-05-25 00:00:00 UTC);
let exp = now + time::Duration::days(7);
let rows = vec![
row("btc", exp, f64::NAN, "call", Some(1.0), Some(0.5)),
row("btc", exp, f64::INFINITY, "call", Some(1.0), Some(0.5)),
row("btc", exp, -50.0, "call", Some(1.0), Some(0.5)),
row("btc", exp, 0.0, "call", Some(1.0), Some(0.5)),
row("btc", exp, 100.0, "call", Some(1.0), Some(0.5)),
];
let out = assemble_chains(rows, now);
let chain = &out[&Asset::Btc][&exp];
assert_eq!(chain.legs.len(), 1, "only the K=100 row should survive");
assert_eq!(chain.legs[0].strike, 100.0);
}

#[test]
fn assemble_folds_call_and_put_into_one_leg() {
let now = datetime!(2026-05-25 00:00:00 UTC);
let exp = now + time::Duration::days(7);
let rows = vec![
row("btc", exp, 100.0, "call", Some(5.0), Some(0.5)),
row("btc", exp, 100.0, "put", Some(4.0), Some(0.55)),
];
let out = assemble_chains(rows, now);
let chain = &out[&Asset::Btc][&exp];
assert_eq!(chain.legs.len(), 1);
assert_eq!(chain.legs[0].strike, 100.0);
assert_eq!(chain.legs[0].call_mid_usd, Some(5.0));
assert_eq!(chain.legs[0].put_mid_usd, Some(4.0));
assert_eq!(chain.legs[0].call_iv, Some(0.5));
assert_eq!(chain.legs[0].put_iv, Some(0.55));
}

#[test]
fn assemble_isolates_btc_and_eth() {
let now = datetime!(2026-05-25 00:00:00 UTC);
let exp = now + time::Duration::days(7);
let rows = vec![
row("btc", exp, 100_000.0, "call", Some(5.0), Some(0.5)),
row("eth", exp, 3_000.0, "call", Some(2.0), Some(0.6)),
];
let out = assemble_chains(rows, now);
assert_eq!(out.len(), 2);
assert_eq!(out[&Asset::Btc][&exp].legs[0].strike, 100_000.0);
assert_eq!(out[&Asset::Eth][&exp].legs[0].strike, 3_000.0);
}

#[test]
fn assemble_isolates_expiries_within_an_asset() {
let now = datetime!(2026-05-25 00:00:00 UTC);
let near = now + time::Duration::days(7);
let far = now + time::Duration::days(40);
let rows = vec![
row("btc", near, 100.0, "call", Some(1.0), Some(0.5)),
row("btc", far, 100.0, "call", Some(2.0), Some(0.55)),
];
let out = assemble_chains(rows, now);
let by_expiry = &out[&Asset::Btc];
assert_eq!(by_expiry.len(), 2);
assert!((by_expiry[&near].time_to_expiry.0 - 7.0 / 365.0).abs() < 1e-12);
assert!((by_expiry[&far].time_to_expiry.0 - 40.0 / 365.0).abs() < 1e-12);
}

#[test]
fn assemble_filters_non_finite_iv_and_mid() {
let now = datetime!(2026-05-25 00:00:00 UTC);
let exp = now + time::Duration::days(7);
let rows = vec![
row(
"btc",
exp,
100.0,
"call",
Some(f64::NAN),
Some(f64::INFINITY),
),
row("btc", exp, 100.0, "put", Some(2.0), Some(0.5)),
];
let out = assemble_chains(rows, now);
let leg = &out[&Asset::Btc][&exp].legs[0];
assert_eq!(leg.call_mid_usd, None, "NaN mid filtered to None");
assert_eq!(leg.call_iv, None, "Inf iv filtered to None");
assert_eq!(leg.put_mid_usd, Some(2.0));
assert_eq!(leg.put_iv, Some(0.5));
}

#[test]
fn assemble_uses_passed_now_for_time_to_expiry() {
let now = datetime!(2026-01-01 00:00:00 UTC);
let exp = datetime!(2026-01-31 00:00:00 UTC);
let rows = vec![row("btc", exp, 100.0, "call", Some(1.0), Some(0.5))];
let out = assemble_chains(rows, now);
let tt = out[&Asset::Btc][&exp].time_to_expiry.0;
assert!((tt - 30.0 / 365.0).abs() < 1e-12);
}

#[test]
fn assemble_empty_input_yields_empty_map() {
let now = datetime!(2026-05-25 00:00:00 UTC);
let out = assemble_chains(Vec::new(), now);
assert!(out.is_empty());
}
}
103 changes: 103 additions & 0 deletions crates/engine/src/sinks.rs
Original file line number Diff line number Diff line change
Expand Up @@ -229,3 +229,106 @@ fn leg_envelope(s: &Strip) -> serde_json::Value {
"quotes": quotes,
})
}

#[cfg(test)]
#[allow(clippy::float_cmp)]
mod tests {
use super::*;
use time::macros::datetime;
use volx_shared_types::ids::IndexId;
use volx_shared_types::index::{IndexValue, StripHash};
use volx_shared_types::strip::StripQuote;
use volx_shared_types::units::Years;

#[allow(clippy::cast_precision_loss)] // tiny n_points in fixtures (≤ 32)
fn fixture_strip(forward: f64, k_zero: f64, t_y: f64, n_points: usize) -> Strip {
let mut quotes = Vec::with_capacity(n_points);
let step = forward * 0.01;
for i in 0..n_points {
let k = forward - step * (n_points as f64 / 2.0) + step * i as f64;
quotes.push(StripQuote {
strike: k,
q_usd: 1.0 + i as f64 * 0.01,
iv: 0.5 + i as f64 * 0.001,
});
}
Strip {
forward,
k_zero,
time_to_expiry: Years(t_y),
quotes,
}
}

fn fixture_index_value() -> IndexValue {
IndexValue {
index_id: IndexId::Bvol,
value: 42.5,
confidence: 0.95,
strip_hash: StripHash([7u8; 32]),
ts: datetime!(2026-05-27 12:00:00 UTC),
}
}

#[test]
fn leg_envelope_pins_field_names_and_quotes_triple_shape() {
let strip = fixture_strip(100.0, 99.5, 30.0 / 365.0, 4);
let env = leg_envelope(&strip);

assert_eq!(env["forward"], 100.0);
assert_eq!(env["k_zero"], 99.5);
assert!((env["time_to_expiry_y"].as_f64().unwrap() - 30.0 / 365.0).abs() < 1e-12);

let quotes = env["quotes"].as_array().unwrap();
assert_eq!(quotes.len(), 4);
// Each entry is a JSON array (not object) of length 3 in [K, Q, iv]
// order — the public-API contract.
for (i, q) in quotes.iter().enumerate() {
let arr = q.as_array().unwrap();
assert_eq!(arr.len(), 3, "entry {i} should be [K, Q, iv] triple");
assert_eq!(arr[0], strip.quotes[i].strike);
assert_eq!(arr[1], strip.quotes[i].q_usd);
assert_eq!(arr[2], strip.quotes[i].iv);
}
}

#[test]
fn strip_envelope_wraps_two_legs_with_top_level_id_and_ts() {
let near = fixture_strip(100.0, 99.5, 7.0 / 365.0, 3);
let next = fixture_strip(100.0, 99.5, 40.0 / 365.0, 3);
let iv = fixture_index_value();

let env = strip_envelope(&iv, &near, &next);

assert_eq!(env["index_id"], "BVOL");
// `ts` is currently serialized via `time::OffsetDateTime`'s default
// serializer (numeric array), not the IndexValue's `rfc3339`
// string form — tracked as #73. Document the current broken
// shape explicitly so the bug-fix PR can flip
// `is_array()` → `is_string()` in one line.
assert!(
env["ts"].is_array() || env["ts"].is_string(),
"ts must be numeric-array (current bug #73) or RFC 3339 string (post-#73); got {:?}",
env["ts"]
);

// Both legs present with leg_envelope shape inherited.
assert_eq!(env["near"]["forward"], 100.0);
assert_eq!(env["next"]["forward"], 100.0);
assert!((env["near"]["time_to_expiry_y"].as_f64().unwrap() - 7.0 / 365.0).abs() < 1e-12);
assert!((env["next"]["time_to_expiry_y"].as_f64().unwrap() - 40.0 / 365.0).abs() < 1e-12);

// Near + next are distinct envelopes (no aliasing).
assert_ne!(env["near"], env["next"]);
}

#[test]
fn strip_envelope_preserves_quote_count_per_leg() {
let near = fixture_strip(100.0, 99.5, 7.0 / 365.0, 5);
let next = fixture_strip(100.0, 99.5, 40.0 / 365.0, 7);
let iv = fixture_index_value();
let env = strip_envelope(&iv, &near, &next);
assert_eq!(env["near"]["quotes"].as_array().unwrap().len(), 5);
assert_eq!(env["next"]["quotes"].as_array().unwrap().len(), 7);
}
}
Loading
Loading