fix: separate book delta and snapshot clocks

This commit is contained in:
floor-licker
2026-06-22 19:50:13 -04:00
parent cb1838e77d
commit 39ac56356e
+101 -13
View File
@@ -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<Utc>,
@@ -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::<Utc>::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::<Utc>::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);