use std::collections::HashMap;
use super::events::DepthLadder;
use super::identity::InstrumentKey;
pub const MAX_DEPTH_BOOKS: usize = 1024;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
#[repr(u8)]
pub enum DepthStatus {
#[default]
Fresh,
ResyncNeeded,
}
#[derive(Debug, Clone)]
pub struct DepthBook {
ladder: DepthLadder,
status: DepthStatus,
}
impl DepthBook {
#[must_use]
pub fn ladder(&self) -> &DepthLadder {
&self.ladder
}
#[must_use]
pub fn status(&self) -> DepthStatus {
self.status
}
#[must_use]
pub fn needs_resync(&self) -> bool {
matches!(self.status, DepthStatus::ResyncNeeded)
}
#[must_use]
pub fn level_count(&self) -> usize {
self.ladder.bids.len() + self.ladder.asks.len()
}
}
#[derive(Debug, Clone, Default)]
pub struct DepthStore {
books: HashMap<InstrumentKey, DepthBook>,
dropped_capacity: u64,
}
impl DepthStore {
#[must_use]
pub fn new() -> Self {
Self::default()
}
pub fn apply(&mut self, ladder: DepthLadder) -> bool {
let key = ladder.instrument.key.clone();
match self.books.get_mut(&key) {
Some(book) => {
book.status = if depth_continues(book.ladder.change_id, ladder.change_id) {
DepthStatus::Fresh
} else {
DepthStatus::ResyncNeeded
};
book.ladder = ladder;
true
}
None => {
if self.books.len() >= MAX_DEPTH_BOOKS {
self.dropped_capacity = self
.dropped_capacity
.checked_add(1)
.unwrap_or(self.dropped_capacity);
return false;
}
let _ = self.books.insert(
key,
DepthBook {
ladder,
status: DepthStatus::Fresh,
},
);
true
}
}
}
#[must_use]
pub fn book(&self, key: &InstrumentKey) -> Option<&DepthBook> {
self.books.get(key)
}
#[must_use]
pub fn len(&self) -> usize {
self.books.len()
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.books.is_empty()
}
#[must_use]
pub fn dropped_capacity(&self) -> u64 {
self.dropped_capacity
}
}
#[must_use]
pub fn depth_continues(prev: Option<u64>, next: Option<u64>) -> bool {
match (prev, next) {
(Some(p), Some(n)) => n >= p,
(Some(_), None) => false,
(None, _) => true,
}
}
#[cfg(test)]
mod tests {
use chrono::{DateTime, Utc};
use optionstratlib::OptionStyle;
use optionstratlib::prelude::Positive;
use super::*;
use crate::chain::events::DepthLevel;
use crate::chain::identity::{
ContractSpecFingerprint, ExerciseStyle, Instrument, InstrumentKey, ProviderId,
SettlementStyle,
};
#[track_caller]
fn pid(id: &str) -> ProviderId {
match ProviderId::new(id) {
Ok(p) => p,
Err(e) => panic!("invalid provider id `{id}`: {e}"),
}
}
#[track_caller]
fn utc(secs: i64) -> DateTime<Utc> {
match DateTime::<Utc>::from_timestamp(secs, 0) {
Some(t) => t,
None => panic!("invalid test timestamp: {secs}"),
}
}
#[track_caller]
fn pos(value: f64) -> Positive {
match Positive::new(value) {
Ok(p) => p,
Err(e) => panic!("invalid test positive `{value}`: {e}"),
}
}
fn key(strike: f64, style: OptionStyle) -> InstrumentKey {
InstrumentKey {
underlying: "BTC".to_owned(),
expiration_utc: utc(1_700_000_000),
strike: pos(strike),
style,
}
}
fn instrument(strike: f64, style: OptionStyle) -> Instrument {
Instrument {
key: key(strike, style),
provider: pid("deribit"),
native_symbol: "BTC-27JUN25-60000-C".to_owned(),
stream_symbol: None,
spec: ContractSpecFingerprint {
contract_multiplier: 1,
settlement: SettlementStyle::Cash,
exercise: ExerciseStyle::European,
quote_currency: "USD".to_owned(),
venue_product_code: "BTC".to_owned(),
},
}
}
fn ladder(strike: f64, style: OptionStyle, change_id: Option<u64>) -> DepthLadder {
DepthLadder {
instrument: instrument(strike, style),
bids: vec![
DepthLevel {
price: pos(60_000.0),
size: pos(2.0),
},
DepthLevel {
price: pos(59_990.0),
size: pos(5.0),
},
],
asks: vec![DepthLevel {
price: pos(60_010.0),
size: pos(1.0),
}],
event_time: Some(utc(1_700_000_099)),
received_time: utc(1_700_000_100),
change_id,
}
}
#[test]
fn test_depth_continues_advance_and_repeat_are_continuous() {
assert!(depth_continues(Some(10), Some(11)), "advance continues");
assert!(
depth_continues(Some(10), Some(50)),
"a forward skip (coalesced drop) still continues — no false resync",
);
assert!(depth_continues(Some(10), Some(10)), "a repeat continues");
}
#[test]
fn test_depth_continues_regression_and_lost_sequence_are_gaps() {
assert!(
!depth_continues(Some(10), Some(9)),
"a regression is a resync boundary",
);
assert!(
!depth_continues(Some(10), None),
"losing the sequence is a resync boundary",
);
}
#[test]
fn test_depth_continues_seed_and_sequenceless_never_flag() {
assert!(depth_continues(None, Some(1)), "the first ladder seeds");
assert!(
depth_continues(None, None),
"a feed with no sequence never flags a resync",
);
}
#[test]
fn test_apply_seeds_a_fresh_book() {
let mut store = DepthStore::new();
assert!(store.is_empty());
assert!(store.apply(ladder(60_000.0, OptionStyle::Call, Some(1))));
let book = match store.book(&key(60_000.0, OptionStyle::Call)) {
Some(b) => b,
None => panic!("the seeded book must be retrievable by its key"),
};
assert_eq!(book.status(), DepthStatus::Fresh, "a seed is fresh");
assert!(!book.needs_resync());
assert_eq!(book.level_count(), 3, "two bids + one ask");
assert_eq!(store.len(), 1);
}
#[test]
fn test_apply_overwrite_keeps_fresh_on_advancing_change_id() {
let mut store = DepthStore::new();
let _ = store.apply(ladder(60_000.0, OptionStyle::Call, Some(1)));
let _ = store.apply(ladder(60_000.0, OptionStyle::Call, Some(5)));
let book = store.book(&key(60_000.0, OptionStyle::Call));
match book {
Some(b) => {
assert_eq!(b.status(), DepthStatus::Fresh, "advancing stays fresh");
assert_eq!(b.ladder().change_id, Some(5), "the latest ladder wins");
}
None => panic!("book present"),
}
assert_eq!(store.len(), 1, "an overwrite does not grow the store");
}
#[test]
fn test_apply_regression_flags_resync_then_clears_on_resume() {
let mut store = DepthStore::new();
let _ = store.apply(ladder(60_000.0, OptionStyle::Call, Some(10)));
let _ = store.apply(ladder(60_000.0, OptionStyle::Call, Some(3)));
match store.book(&key(60_000.0, OptionStyle::Call)) {
Some(b) => assert!(b.needs_resync(), "a change_id regression flags resync"),
None => panic!("book present"),
}
let _ = store.apply(ladder(60_000.0, OptionStyle::Call, Some(4)));
match store.book(&key(60_000.0, OptionStyle::Call)) {
Some(b) => assert_eq!(
b.status(),
DepthStatus::Fresh,
"an advancing resumed sequence clears the resync badge",
),
None => panic!("book present"),
}
}
#[test]
fn test_apply_tracks_distinct_instruments_separately() {
let mut store = DepthStore::new();
let _ = store.apply(ladder(60_000.0, OptionStyle::Call, Some(1)));
let _ = store.apply(ladder(60_000.0, OptionStyle::Put, Some(1)));
assert_eq!(
store.len(),
2,
"call and put at one strike are distinct books"
);
assert!(store.book(&key(60_000.0, OptionStyle::Call)).is_some());
assert!(store.book(&key(60_000.0, OptionStyle::Put)).is_some());
assert!(
store.book(&key(62_000.0, OptionStyle::Call)).is_none(),
"an unseen instrument has no book",
);
}
#[test]
fn test_apply_drops_new_instrument_at_capacity() {
let mut store = DepthStore::new();
for i in 0..MAX_DEPTH_BOOKS {
let strike = 1_000.0 + i as f64;
assert!(store.apply(ladder(strike, OptionStyle::Call, Some(1))));
}
assert_eq!(store.len(), MAX_DEPTH_BOOKS);
assert!(
!store.apply(ladder(999_999.0, OptionStyle::Call, Some(1))),
"a new instrument at capacity is dropped",
);
assert_eq!(store.len(), MAX_DEPTH_BOOKS, "the cap holds");
assert_eq!(store.dropped_capacity(), 1);
assert!(
store.apply(ladder(1_000.0, OptionStyle::Call, Some(2))),
"an existing instrument still updates at capacity",
);
assert_eq!(store.len(), MAX_DEPTH_BOOKS);
}
}