diff --git a/src/book.rs b/src/book.rs index 37bdf35..2040bd8 100644 --- a/src/book.rs +++ b/src/book.rs @@ -413,8 +413,8 @@ impl OrderBook { } self.sequence = timestamp; - self.timestamp = - chrono::DateTime::::from_timestamp(timestamp as i64, 0).unwrap_or_else(Utc::now); + self.timestamp = chrono::DateTime::::from_timestamp_millis(timestamp as i64) + .unwrap_or_else(Utc::now); Ok(true) } @@ -443,20 +443,26 @@ impl OrderBook { Ok(()) } - /// Finish applying a WS `book` update. - pub(crate) fn finish_ws_book_update(&mut self) { + /// Finish applying a WS `book` snapshot. + pub(crate) fn finish_ws_book_update( + &mut self, + mut has_bid: impl FnMut(Price) -> bool, + mut has_ask: impl FnMut(Price) -> bool, + ) { + self.bids.retain(|price_ticks, _| has_bid(*price_ticks)); + self.asks.retain(|price_ticks, _| has_ask(*price_ticks)); self.trim_depth(); } /// Apply a WebSocket `book` update for this token. /// - /// The official Polymarket CLOB WebSocket `book` event contains batches of - /// price levels for both sides. Unlike `apply_delta_fast`, this method can - /// apply many levels that share the same message timestamp. + /// The official Polymarket CLOB WebSocket `book` event is a full snapshot. + /// Unlike `apply_delta_fast`, this method can apply many levels that share + /// the same message timestamp. /// /// Notes: - /// - This performs upserts (update/insert/remove) for the provided levels. - /// - It does **not** infer removals for levels omitted from the message. + /// - This replaces the snapshot for both sides. + /// - Levels omitted from the message are removed. /// - Insertions of *new* price levels may allocate (BTreeMap node growth). pub fn apply_book_update(&mut self, update: &BookUpdate) -> Result<()> { if update.asset_id != self.token_id { @@ -471,7 +477,7 @@ impl OrderBook { } self.sequence = update.timestamp; - self.timestamp = chrono::DateTime::::from_timestamp(update.timestamp as i64, 0) + self.timestamp = chrono::DateTime::::from_timestamp_millis(update.timestamp as i64) .unwrap_or_else(Utc::now); // Apply bids (BUY) and asks (SELL) as level upserts. @@ -513,6 +519,10 @@ impl OrderBook { } } + self.bids + .retain(|price_ticks, _| book_update_has_level(&update.bids, *price_ticks)); + self.asks + .retain(|price_ticks, _| book_update_has_level(&update.asks, *price_ticks)); self.trim_depth(); Ok(()) } @@ -759,6 +769,19 @@ impl OrderBook { } } +fn book_update_has_level(levels: &[OrderSummary], price_ticks: Price) -> bool { + levels.iter().any(|level| { + let Ok(level_price_ticks) = decimal_to_price(level.price) else { + return false; + }; + let Ok(size_units) = decimal_to_qty(level.size) else { + return false; + }; + + size_units != 0 && level_price_ticks == price_ticks + }) +} + /// Market impact calculation result /// This tells you what would happen if you executed a large order #[derive(Debug, Clone)] @@ -1021,6 +1044,115 @@ mod tests { assert_eq!(book.best_bid().unwrap().size, dec!(100)); // Should be our size } + #[test] + fn test_book_update_replaces_snapshot_and_uses_millis_timestamp() { + let mut book = OrderBook::new("test_token".to_string(), 10); + let timestamp = 1_757_908_892_351; + + book.apply_book_update(&BookUpdate { + asset_id: "test_token".to_string(), + market: "0xabc".to_string(), + timestamp, + bids: vec![ + OrderSummary { + price: dec!(0.48), + size: dec!(10), + }, + OrderSummary { + price: dec!(0.49), + size: dec!(20), + }, + ], + asks: vec![ + OrderSummary { + price: dec!(0.52), + size: dec!(30), + }, + OrderSummary { + price: dec!(0.53), + size: dec!(40), + }, + ], + hash: None, + }) + .unwrap(); + + book.apply_book_update(&BookUpdate { + asset_id: "test_token".to_string(), + market: "0xabc".to_string(), + timestamp: timestamp + 1, + bids: vec![OrderSummary { + price: dec!(0.49), + size: dec!(25), + }], + asks: vec![OrderSummary { + price: dec!(0.53), + size: dec!(45), + }], + hash: None, + }) + .unwrap(); + + assert_eq!(book.timestamp.timestamp_millis(), timestamp as i64 + 1); + assert_eq!(book.bids(None).len(), 1); + assert_eq!(book.asks(None).len(), 1); + assert_eq!(book.best_bid().unwrap().price, dec!(0.49)); + assert_eq!(book.best_bid().unwrap().size, dec!(25)); + assert_eq!(book.best_ask().unwrap().price, dec!(0.53)); + assert_eq!(book.best_ask().unwrap().size, dec!(45)); + } + + #[test] + fn test_book_update_depth_keeps_best_prices_independent_of_payload_order() { + let mut book = OrderBook::new("test_token".to_string(), 2); + + book.apply_book_update(&BookUpdate { + asset_id: "test_token".to_string(), + market: "0xabc".to_string(), + timestamp: 1_757_908_892_351, + bids: vec![ + OrderSummary { + price: dec!(0.49), + size: dec!(10), + }, + OrderSummary { + price: dec!(0.50), + size: dec!(20), + }, + OrderSummary { + price: dec!(0.48), + size: dec!(30), + }, + ], + asks: vec![ + OrderSummary { + price: dec!(0.52), + size: dec!(40), + }, + OrderSummary { + price: dec!(0.54), + size: dec!(50), + }, + OrderSummary { + price: dec!(0.53), + size: dec!(60), + }, + ], + hash: None, + }) + .unwrap(); + + let bids = book.bids(None); + assert_eq!(bids.len(), 2); + assert_eq!(bids[0].price, dec!(0.50)); + assert_eq!(bids[1].price, dec!(0.49)); + + let asks = book.asks(None); + assert_eq!(asks.len(), 2); + assert_eq!(asks[0].price, dec!(0.52)); + assert_eq!(asks[1].price, dec!(0.53)); + } + #[test] fn test_spread_calculation() { // Test that we can calculate the spread between bid and ask diff --git a/src/stream.rs b/src/stream.rs index a3e2613..5452b1e 100644 --- a/src/stream.rs +++ b/src/stream.rs @@ -762,4 +762,31 @@ mod tests { assert_eq!(snapshot.asks[0].price, Decimal::from_str("0.76").unwrap()); assert_eq!(snapshot.asks[0].size, Decimal::from_str("5").unwrap()); } + + #[test] + fn test_websocket_book_applier_replaces_snapshot() { + let books = crate::book::OrderBookManager::new(64); + let _ = books.get_or_create_book("12345").unwrap(); + + let processor = WsBookUpdateProcessor::new(1024); + let stream = WebSocketStream::new("wss://example.com/ws"); + let mut applier = stream.into_book_applier(&books, processor); + + let first = r#"{"event_type":"book","asset_id":"12345","timestamp":"1757908892351","bids":[{"price":"0.74","size":"8"},{"price":"0.75","size":"10"}],"asks":[{"price":"0.76","size":"5"},{"price":"0.77","size":"4"}]}"#.to_string(); + applier.apply_text_message(first).unwrap(); + + let second = r#"{"event_type":"book","asset_id":"12345","timestamp":"1757908892352","bids":[{"price":"0.75","size":"11"}],"asks":[{"price":"0.76","size":"6"}]}"#.to_string(); + let stats = applier.apply_text_message(second).unwrap(); + assert_eq!(stats.book_messages, 1); + assert_eq!(stats.book_levels_applied, 2); + + let snapshot = books.get_book("12345").unwrap(); + assert_eq!(snapshot.timestamp.timestamp_millis(), 1_757_908_892_352); + assert_eq!(snapshot.bids.len(), 1); + assert_eq!(snapshot.asks.len(), 1); + assert_eq!(snapshot.bids[0].price, Decimal::from_str("0.75").unwrap()); + assert_eq!(snapshot.bids[0].size, Decimal::from_str("11").unwrap()); + assert_eq!(snapshot.asks[0].price, Decimal::from_str("0.76").unwrap()); + assert_eq!(snapshot.asks[0].size, Decimal::from_str("6").unwrap()); + } } diff --git a/src/ws_hot_path.rs b/src/ws_hot_path.rs index 7eaa98c..a3a9233 100644 --- a/src/ws_hot_path.rs +++ b/src/ws_hot_path.rs @@ -9,7 +9,7 @@ use crate::book::OrderBookManager; use crate::errors::{PolyfillError, Result}; -use crate::types::{decimal_to_price, decimal_to_qty, Side}; +use crate::types::{decimal_to_price, decimal_to_qty, Price, Side}; use rust_decimal::Decimal; use simd_json::prelude::*; use std::str::FromStr; @@ -141,7 +141,10 @@ fn process_stream_object<'tape, 'input>( applied += apply_levels(book, Side::SELL, asks)?; } - book.finish_ws_book_update(); + book.finish_ws_book_update( + |price_ticks| ws_levels_contain_price(bids, price_ticks), + |price_ticks| ws_levels_contain_price(asks, price_ticks), + ); Ok(applied) })?; @@ -193,3 +196,38 @@ fn apply_levels<'tape, 'input>( Ok(applied) } + +fn ws_levels_contain_price<'tape, 'input>( + levels: Option>, + price_ticks: Price, +) -> bool { + let Some(levels) = levels else { + return false; + }; + + levels.iter().any(|level| { + let Some(obj) = level.as_object() else { + return false; + }; + let Some(price_str) = obj.get("price").and_then(|v| v.into_string()) else { + return false; + }; + let Some(size_str) = obj.get("size").and_then(|v| v.into_string()) else { + return false; + }; + let Ok(price_decimal) = Decimal::from_str(price_str) else { + return false; + }; + let Ok(size_decimal) = Decimal::from_str(size_str) else { + return false; + }; + let Ok(level_price_ticks) = decimal_to_price(price_decimal) else { + return false; + }; + let Ok(size_units) = decimal_to_qty(size_decimal) else { + return false; + }; + + size_units != 0 && level_price_ticks == price_ticks + }) +}