From 39ac56356ebd5eeef9f6444f12a72eade56480a6 Mon Sep 17 00:00:00 2001 From: floor-licker Date: Mon, 22 Jun 2026 19:50:13 -0400 Subject: [PATCH] fix: separate book delta and snapshot clocks --- src/book.rs | 114 ++++++++++++++++++++++++++++++++++++++++++++++------ 1 file changed, 101 insertions(+), 13 deletions(-) diff --git a/src/book.rs b/src/book.rs index a866f71..73dbeb0 100644 --- a/src/book.rs +++ b/src/book.rs @@ -158,10 +158,23 @@ pub struct OrderBook { /// Hash of token_id for fast lookups (avoids string comparisons in hot path) pub token_id_hash: u64, - /// Current sequence number for ordering updates - /// This helps us ignore old/duplicate updates that arrive out of order + /// Legacy delta sequence number. + /// + /// Kept as a compatibility alias for [`Self::last_delta_sequence`]. WebSocket + /// snapshot timestamps are tracked separately in [`Self::last_snapshot_timestamp_ms`]. pub sequence: u64, + /// Last accepted legacy/incremental delta sequence number. + /// + /// Delta sequences are a different clock from WebSocket snapshot timestamps and + /// must not be compared against millisecond timestamps. + pub last_delta_sequence: u64, + + /// Last accepted full-book snapshot timestamp in milliseconds. + /// + /// Used for WebSocket `book` snapshots and serde `BookUpdate` snapshots. + pub last_snapshot_timestamp_ms: u64, + /// Last update timestamp - when we last got new data for this book pub timestamp: chrono::DateTime, @@ -222,7 +235,9 @@ impl OrderBook { Self { token_id, token_id_hash, - sequence: 0, // Start at 0, will increment as we get updates + sequence: 0, // Legacy alias for last_delta_sequence + last_delta_sequence: 0, + last_snapshot_timestamp_ms: 0, timestamp: Utc::now(), bids: BookSide::new(BookSideKind::Bid, max_depth), // Empty to start asks: BookSide::new(BookSideKind::Ask, max_depth), // Empty to start @@ -430,7 +445,7 @@ impl OrderBook { timestamp: self.timestamp, bids: self.bids(None), // Get all bids (up to max_depth) asks: self.asks(None), // Get all asks (up to max_depth) - sequence: self.sequence, + sequence: self.last_delta_sequence, } } @@ -461,11 +476,11 @@ impl OrderBook { pub fn apply_delta_fast(&mut self, delta: FastOrderDelta) -> Result<()> { // Validate sequence ordering - ignore old updates that arrive late // This is crucial for maintaining data integrity in real-time systems - if delta.sequence <= self.sequence { + if delta.sequence <= self.last_delta_sequence { trace!( "Ignoring stale delta: {} <= {}", delta.sequence, - self.sequence + self.last_delta_sequence ); return Ok(()); } @@ -496,6 +511,7 @@ impl OrderBook { } // Update our tracking info + self.last_delta_sequence = delta.sequence; self.sequence = delta.sequence; self.timestamp = delta.timestamp; @@ -531,11 +547,11 @@ impl OrderBook { return Err(PolyfillError::validation("Token ID mismatch")); } - if timestamp <= self.sequence { + if timestamp <= self.last_snapshot_timestamp_ms { return Ok(false); } - self.sequence = timestamp; + self.last_snapshot_timestamp_ms = timestamp; self.timestamp = chrono::DateTime::::from_timestamp_millis(timestamp as i64) .unwrap_or_else(Utc::now); self.begin_snapshot(); @@ -585,14 +601,13 @@ impl OrderBook { return Err(PolyfillError::validation("Token ID mismatch")); } - // Use the exchange-provided timestamp as our monotonic sequence marker. - // This is less strict than the REST/legacy delta sequence but works for - // ignoring obviously stale book snapshots. - if update.timestamp <= self.sequence { + // Use the exchange-provided timestamp as the monotonic snapshot marker. + // Snapshot timestamps are separate from legacy/incremental delta sequences. + if update.timestamp <= self.last_snapshot_timestamp_ms { return Ok(()); } - self.sequence = update.timestamp; + self.last_snapshot_timestamp_ms = update.timestamp; self.timestamp = chrono::DateTime::::from_timestamp_millis(update.timestamp as i64) .unwrap_or_else(Utc::now); self.begin_snapshot(); @@ -1220,6 +1235,8 @@ mod tests { book.apply_delta(delta).unwrap(); assert_eq!(book.sequence, 1); // Sequence should update + assert_eq!(book.last_delta_sequence, 1); + assert_eq!(book.last_snapshot_timestamp_ms, 0); assert_eq!(book.best_bid().unwrap().price, dec!(0.5)); // Should be our bid assert_eq!(book.best_bid().unwrap().size, dec!(100)); // Should be our size } @@ -1346,6 +1363,77 @@ mod tests { assert_eq!(book.best_ask().unwrap().size, dec!(45)); } + #[test] + fn test_ws_snapshot_timestamp_does_not_block_delta_sequence() { + let mut book = OrderBook::new("test_token".to_string(), 10); + let snapshot_timestamp_ms = 1_757_908_892_351; + + assert!(book + .begin_ws_book_update("test_token", snapshot_timestamp_ms) + .unwrap()); + book.apply_ws_book_level_fast(Side::BUY, 5_000, 100_000) + .unwrap(); + book.apply_ws_book_level_fast(Side::SELL, 6_000, 100_000) + .unwrap(); + book.finish_ws_book_update(); + + assert_eq!(book.sequence, 0); + assert_eq!(book.last_delta_sequence, 0); + assert_eq!(book.last_snapshot_timestamp_ms, snapshot_timestamp_ms); + + book.apply_delta(OrderDelta { + token_id: "test_token".to_string(), + timestamp: Utc::now(), + side: Side::BUY, + price: dec!(0.51), + size: dec!(11), + sequence: 1, + }) + .unwrap(); + + assert_eq!(book.sequence, 1); + assert_eq!(book.last_delta_sequence, 1); + assert_eq!(book.last_snapshot_timestamp_ms, snapshot_timestamp_ms); + assert_eq!(book.best_bid().unwrap().price, dec!(0.51)); + } + + #[test] + fn test_delta_sequence_does_not_block_snapshot_timestamp() { + let mut book = OrderBook::new("test_token".to_string(), 10); + + book.apply_delta(OrderDelta { + token_id: "test_token".to_string(), + timestamp: Utc::now(), + side: Side::BUY, + price: dec!(0.40), + size: dec!(10), + sequence: 10_000, + }) + .unwrap(); + + book.apply_book_update(&BookUpdate { + asset_id: "test_token".to_string(), + market: "0xabc".to_string(), + timestamp: 1_000, + bids: vec![OrderSummary { + price: dec!(0.50), + size: dec!(20), + }], + asks: vec![OrderSummary { + price: dec!(0.60), + size: dec!(30), + }], + hash: None, + }) + .unwrap(); + + assert_eq!(book.sequence, 10_000); + assert_eq!(book.last_delta_sequence, 10_000); + assert_eq!(book.last_snapshot_timestamp_ms, 1_000); + assert_eq!(book.best_bid().unwrap().price, dec!(0.50)); + assert_eq!(book.best_ask().unwrap().price, dec!(0.60)); + } + #[test] fn test_book_update_depth_keeps_best_prices_independent_of_payload_order() { let mut book = OrderBook::new("test_token".to_string(), 2);