mirror of
https://github.com/floor-licker/polyfill-rs.git
synced 2026-08-25 18:38:09 +00:00
Merge pull request #62 from floor-licker/fix/separate-book-update-clocks
fix: separate book delta and snapshot clocks
This commit is contained in:
+111
-12
@@ -165,10 +165,23 @@ pub struct OrderBook {
|
|||||||
/// Hash of token_id for fast lookups (avoids string comparisons in hot path)
|
/// Hash of token_id for fast lookups (avoids string comparisons in hot path)
|
||||||
pub token_id_hash: u64,
|
pub token_id_hash: u64,
|
||||||
|
|
||||||
/// Current sequence number for ordering updates
|
/// Legacy delta sequence number.
|
||||||
/// This helps us ignore old/duplicate updates that arrive out of order
|
///
|
||||||
|
/// 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,
|
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
|
/// Last update timestamp - when we last got new data for this book
|
||||||
pub timestamp: chrono::DateTime<Utc>,
|
pub timestamp: chrono::DateTime<Utc>,
|
||||||
|
|
||||||
@@ -235,7 +248,9 @@ impl OrderBook {
|
|||||||
Self {
|
Self {
|
||||||
token_id,
|
token_id,
|
||||||
token_id_hash,
|
token_id_hash,
|
||||||
sequence: 0, // Start at 0, will increment as we get updates
|
sequence: 0, // Compatibility alias for last_delta_sequence
|
||||||
|
last_delta_sequence: 0,
|
||||||
|
last_snapshot_timestamp_ms: 0,
|
||||||
timestamp: Utc::now(),
|
timestamp: Utc::now(),
|
||||||
bids: BookSide::new(BookSideKind::Bid, max_depth), // Empty to start
|
bids: BookSide::new(BookSideKind::Bid, max_depth), // Empty to start
|
||||||
asks: BookSide::new(BookSideKind::Ask, max_depth), // Empty to start
|
asks: BookSide::new(BookSideKind::Ask, max_depth), // Empty to start
|
||||||
@@ -444,7 +459,7 @@ impl OrderBook {
|
|||||||
timestamp: self.timestamp,
|
timestamp: self.timestamp,
|
||||||
bids: self.bids(None), // Get all bids (up to max_depth)
|
bids: self.bids(None), // Get all bids (up to max_depth)
|
||||||
asks: self.asks(None), // Get all asks (up to max_depth)
|
asks: self.asks(None), // Get all asks (up to max_depth)
|
||||||
sequence: self.sequence,
|
sequence: self.last_delta_sequence,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -475,11 +490,11 @@ impl OrderBook {
|
|||||||
pub fn apply_delta_fast(&mut self, delta: FastOrderDelta) -> Result<()> {
|
pub fn apply_delta_fast(&mut self, delta: FastOrderDelta) -> Result<()> {
|
||||||
// Validate sequence ordering - ignore old updates that arrive late
|
// Validate sequence ordering - ignore old updates that arrive late
|
||||||
// This is crucial for maintaining data integrity in real-time systems
|
// This is crucial for maintaining data integrity in real-time systems
|
||||||
if delta.sequence <= self.sequence {
|
if delta.sequence <= self.last_delta_sequence {
|
||||||
trace!(
|
trace!(
|
||||||
"Ignoring stale delta: {} <= {}",
|
"Ignoring stale delta: {} <= {}",
|
||||||
delta.sequence,
|
delta.sequence,
|
||||||
self.sequence
|
self.last_delta_sequence
|
||||||
);
|
);
|
||||||
return Ok(());
|
return Ok(());
|
||||||
}
|
}
|
||||||
@@ -510,6 +525,7 @@ impl OrderBook {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Update our tracking info
|
// Update our tracking info
|
||||||
|
self.last_delta_sequence = delta.sequence;
|
||||||
self.sequence = delta.sequence;
|
self.sequence = delta.sequence;
|
||||||
self.timestamp = delta.timestamp;
|
self.timestamp = delta.timestamp;
|
||||||
|
|
||||||
@@ -583,9 +599,10 @@ impl OrderBook {
|
|||||||
return Err(PolyfillError::validation("Token ID mismatch"));
|
return Err(PolyfillError::validation("Token ID mismatch"));
|
||||||
}
|
}
|
||||||
|
|
||||||
// Use the exchange-provided timestamp as our monotonic snapshot marker.
|
// Use the exchange-provided timestamp as the monotonic snapshot marker.
|
||||||
// If the feed emits two snapshots in the same millisecond, a different hash is treated as
|
// If the feed emits two snapshots in the same millisecond, a different hash is treated as
|
||||||
// a distinct newer state. Equal timestamp without a hash is considered duplicate/stale.
|
// a distinct newer state. Equal timestamp without a hash is considered duplicate/stale.
|
||||||
|
// Snapshot timestamps are intentionally separate from legacy delta sequence numbers.
|
||||||
if !self.should_apply_snapshot(update.timestamp, update.hash.as_deref()) {
|
if !self.should_apply_snapshot(update.timestamp, update.hash.as_deref()) {
|
||||||
return Ok(());
|
return Ok(());
|
||||||
}
|
}
|
||||||
@@ -602,7 +619,7 @@ impl OrderBook {
|
|||||||
self.validate_snapshot_level(parsed)?;
|
self.validate_snapshot_level(parsed)?;
|
||||||
}
|
}
|
||||||
|
|
||||||
self.sequence = update.timestamp;
|
self.last_snapshot_timestamp_ms = update.timestamp;
|
||||||
self.last_snapshot_hash = update.hash.clone();
|
self.last_snapshot_hash = update.hash.clone();
|
||||||
self.timestamp = chrono::DateTime::<Utc>::from_timestamp_millis(update.timestamp as i64)
|
self.timestamp = chrono::DateTime::<Utc>::from_timestamp_millis(update.timestamp as i64)
|
||||||
.unwrap_or_else(Utc::now);
|
.unwrap_or_else(Utc::now);
|
||||||
@@ -706,10 +723,10 @@ impl OrderBook {
|
|||||||
|
|
||||||
#[inline]
|
#[inline]
|
||||||
fn should_apply_snapshot(&self, timestamp: u64, hash: Option<&str>) -> bool {
|
fn should_apply_snapshot(&self, timestamp: u64, hash: Option<&str>) -> bool {
|
||||||
if timestamp > self.sequence {
|
if timestamp > self.last_snapshot_timestamp_ms {
|
||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
if timestamp < self.sequence {
|
if timestamp < self.last_snapshot_timestamp_ms {
|
||||||
return false;
|
return false;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -751,7 +768,7 @@ impl OrderBook {
|
|||||||
self.validate_snapshot_level(level)?;
|
self.validate_snapshot_level(level)?;
|
||||||
}
|
}
|
||||||
|
|
||||||
self.sequence = timestamp;
|
self.last_snapshot_timestamp_ms = timestamp;
|
||||||
self.last_snapshot_hash = hash.map(str::to_owned);
|
self.last_snapshot_hash = hash.map(str::to_owned);
|
||||||
self.timestamp = chrono::DateTime::<Utc>::from_timestamp_millis(timestamp as i64)
|
self.timestamp = chrono::DateTime::<Utc>::from_timestamp_millis(timestamp as i64)
|
||||||
.unwrap_or_else(Utc::now);
|
.unwrap_or_else(Utc::now);
|
||||||
@@ -1280,6 +1297,8 @@ mod tests {
|
|||||||
|
|
||||||
book.apply_delta(delta).unwrap();
|
book.apply_delta(delta).unwrap();
|
||||||
assert_eq!(book.sequence, 1); // Sequence should update
|
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().price, dec!(0.5)); // Should be our bid
|
||||||
assert_eq!(book.best_bid().unwrap().size, dec!(100)); // Should be our size
|
assert_eq!(book.best_bid().unwrap().size, dec!(100)); // Should be our size
|
||||||
}
|
}
|
||||||
@@ -1406,6 +1425,84 @@ mod tests {
|
|||||||
assert_eq!(book.best_ask().unwrap().size, dec!(45));
|
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;
|
||||||
|
let levels = [
|
||||||
|
ParsedBookLevel {
|
||||||
|
side: Side::BUY,
|
||||||
|
price_ticks: 5_000,
|
||||||
|
size_units: 100_000,
|
||||||
|
},
|
||||||
|
ParsedBookLevel {
|
||||||
|
side: Side::SELL,
|
||||||
|
price_ticks: 6_000,
|
||||||
|
size_units: 100_000,
|
||||||
|
},
|
||||||
|
];
|
||||||
|
|
||||||
|
assert!(book
|
||||||
|
.apply_ws_book_snapshot_fast("test_token", snapshot_timestamp_ms, None, &levels)
|
||||||
|
.unwrap());
|
||||||
|
|
||||||
|
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]
|
#[test]
|
||||||
fn test_book_update_allows_same_timestamp_with_different_hash() {
|
fn test_book_update_allows_same_timestamp_with_different_hash() {
|
||||||
let mut book = OrderBook::new("test_token".to_string(), 10);
|
let mut book = OrderBook::new("test_token".to_string(), 10);
|
||||||
@@ -1554,7 +1651,9 @@ mod tests {
|
|||||||
.unwrap_err();
|
.unwrap_err();
|
||||||
|
|
||||||
assert!(err.to_string().contains("Price not aligned"));
|
assert!(err.to_string().contains("Price not aligned"));
|
||||||
assert_eq!(book.sequence, 100);
|
assert_eq!(book.sequence, 0);
|
||||||
|
assert_eq!(book.last_delta_sequence, 0);
|
||||||
|
assert_eq!(book.last_snapshot_timestamp_ms, 100);
|
||||||
assert_eq!(book.bids(None).len(), 1);
|
assert_eq!(book.bids(None).len(), 1);
|
||||||
assert_eq!(book.asks(None).len(), 1);
|
assert_eq!(book.asks(None).len(), 1);
|
||||||
assert_eq!(book.best_bid().unwrap().price, dec!(0.50));
|
assert_eq!(book.best_bid().unwrap().price, dec!(0.50));
|
||||||
|
|||||||
+6
-3
@@ -370,7 +370,8 @@ mod tests {
|
|||||||
.is_err());
|
.is_err());
|
||||||
|
|
||||||
let snapshot = books.get_book("test_asset_id").unwrap();
|
let snapshot = books.get_book("test_asset_id").unwrap();
|
||||||
assert_eq!(snapshot.sequence, 1000);
|
assert_eq!(snapshot.sequence, 0);
|
||||||
|
assert_eq!(snapshot.timestamp.timestamp_millis(), 1000);
|
||||||
assert_eq!(snapshot.bids.len(), 1);
|
assert_eq!(snapshot.bids.len(), 1);
|
||||||
assert_eq!(snapshot.asks.len(), 1);
|
assert_eq!(snapshot.asks.len(), 1);
|
||||||
assert_eq!(snapshot.bids[0].price, dec!(0.50));
|
assert_eq!(snapshot.bids[0].price, dec!(0.50));
|
||||||
@@ -407,7 +408,8 @@ mod tests {
|
|||||||
.is_err());
|
.is_err());
|
||||||
|
|
||||||
let snapshot = books.get_book("test_asset_id").unwrap();
|
let snapshot = books.get_book("test_asset_id").unwrap();
|
||||||
assert_eq!(snapshot.sequence, 1000);
|
assert_eq!(snapshot.sequence, 0);
|
||||||
|
assert_eq!(snapshot.timestamp.timestamp_millis(), 1000);
|
||||||
assert_eq!(snapshot.bids.len(), 1);
|
assert_eq!(snapshot.bids.len(), 1);
|
||||||
assert_eq!(snapshot.asks.len(), 1);
|
assert_eq!(snapshot.asks.len(), 1);
|
||||||
assert_eq!(snapshot.bids[0].price, dec!(0.50));
|
assert_eq!(snapshot.bids[0].price, dec!(0.50));
|
||||||
@@ -433,7 +435,8 @@ mod tests {
|
|||||||
assert_eq!(second_stats.book_levels_applied, 2);
|
assert_eq!(second_stats.book_levels_applied, 2);
|
||||||
|
|
||||||
let snapshot = books.get_book("test_asset_id").unwrap();
|
let snapshot = books.get_book("test_asset_id").unwrap();
|
||||||
assert_eq!(snapshot.sequence, 1000);
|
assert_eq!(snapshot.sequence, 0);
|
||||||
|
assert_eq!(snapshot.timestamp.timestamp_millis(), 1000);
|
||||||
assert_eq!(snapshot.bids[0].price, dec!(0.51));
|
assert_eq!(snapshot.bids[0].price, dec!(0.51));
|
||||||
assert_eq!(snapshot.bids[0].size, dec!(11));
|
assert_eq!(snapshot.bids[0].size, dec!(11));
|
||||||
assert_eq!(snapshot.asks[0].price, dec!(0.61));
|
assert_eq!(snapshot.asks[0].price, dec!(0.61));
|
||||||
|
|||||||
Reference in New Issue
Block a user