diff --git a/rust_hft/prediction-markets/crates/ploy-market-data/src/feeds.rs b/rust_hft/prediction-markets/crates/ploy-market-data/src/feeds.rs index b9f622861..6c93fadcb 100644 --- a/rust_hft/prediction-markets/crates/ploy-market-data/src/feeds.rs +++ b/rust_hft/prediction-markets/crates/ploy-market-data/src/feeds.rs @@ -47,6 +47,10 @@ const POLYMARKET_CLOB_FAILURE_CAPTURE_ENV: &str = "MONDAY_POLYMARKET_CLOB_FAILUR const POLYMARKET_CLOB_CROSSED_BOOK_ERROR: &str = "Polymarket price-change batch produced a crossed book"; const MAX_POLYMARKET_CLOB_FAILURE_CAPTURE_BYTES: usize = 1_048_576; +// Deltas buffered per token between (re)subscription and the first book +// snapshot; exceeding the bound fails the token closed instead of growing +// without limit. +const MAX_POLYMARKET_CLOB_PENDING_CHANGES: usize = 256; const NEAR_DEPTH_PCT_RANGE: f64 = 0.001; const DB_POLYMARKET_SETTLEMENT_RETRY_LOOKBACK_SECS: i64 = 30 * 60; @@ -958,9 +962,11 @@ struct ClobBookState { snapshot_timestamp: Option, initialized: bool, // Set when the cached depth proved stale against a provider-reported BBA. - // A dirty token stays failed closed until a full book snapshot clears it, - // unlike a merely uninitialized new token which may quote entry BBA. + // A dirty token stays failed closed until a full book snapshot clears it. dirty: bool, + // Deltas buffered while the token awaits its first book snapshot after + // (re)subscription; replayed in arrival order once the snapshot lands. + pending: Vec<(i64, PriceChangeBatchEntry)>, } impl ClobBookState { @@ -1140,14 +1146,93 @@ impl ClobBookState { } } +fn clob_failed_closed_updates(token_id: &str, ts: DateTime) -> Vec { + let token_id: Arc = Arc::from(token_id); + vec![ + MarketUpdate::QuoteCollectionFailure { + token_id: Arc::clone(&token_id), + request_started_at: ts, + http_status: None, + error_kind: Arc::from("websocket_payload"), + ts, + }, + MarketUpdate::Quote { + token_id, + bid: None, + ask: None, + bid_size: None, + ask_size: None, + bid_levels: Vec::new(), + ask_levels: Vec::new(), + ts, + }, + ] +} + fn market_update_from_clob_book( book: &BookUpdate, state: &mut ClobBookState, -) -> Result { +) -> Result<(Vec, i64), String> { + let pending = std::mem::take(&mut state.pending); state.replace(book)?; - let ts = DateTime::from_timestamp_millis(book.timestamp) + // Replay buffered messages in arrival order. Messages strictly older than + // an already-applied timestamp are skipped (the live path likewise rejects + // regressing source time), each replayed message's final entry faces the + // same BBA reconciliation as a live batch, and the timestamp watermark is + // rebuilt from what was actually applied. + let mut watermark = book.timestamp; + let mut replay_dirty = false; + for (index, (timestamp, entry)) in pending.iter().enumerate() { + if *timestamp < watermark { + continue; + } + state.apply(entry, *timestamp); + watermark = *timestamp; + let message_ends = pending + .get(index + 1) + .is_none_or(|(next_timestamp, _)| next_timestamp != timestamp); + if message_ends + && (entry + .best_bid + .is_some_and(|price| pm_tradeable_price(price) && !state.bids.contains_key(&price)) + || entry.best_ask.is_some_and(|price| { + pm_tradeable_price(price) && !state.asks.contains_key(&price) + })) + { + replay_dirty = true; + break; + } + } + if replay_dirty { + *state = ClobBookState::default(); + state.dirty = true; + return Ok(( + clob_failed_closed_updates(&book.asset_id.to_string(), Utc::now()), + watermark, + )); + } + let ts = DateTime::from_timestamp_millis(watermark) .ok_or_else(|| "Polymarket book timestamp is out of range".to_string())?; - Ok(state.quote(book.asset_id.to_string(), ts, None)) + let quote = state.quote(book.asset_id.to_string(), ts, None); + // A replayed book can cross when buffered entries lack BBA fields, so the + // final quote faces the same crossed-book check as a live batch, with the + // same in-band isolation as the reconciliation above. + if let MarketUpdate::Quote { + bid: Some(bid), + ask: Some(ask), + .. + } = quote + { + if bid > ask { + *state = ClobBookState::default(); + state.dirty = true; + return Ok(( + clob_failed_closed_updates(&book.asset_id.to_string(), Utc::now()), + watermark, + )); + } + } + Ok((vec![quote], watermark)) } fn market_updates_from_price_change( @@ -1209,10 +1294,21 @@ fn market_updates_from_price_change( } for (token_id, entry) in applicable_entries { - books - .entry(token_id.clone()) - .or_default() - .apply(entry, change.timestamp); + let state = books.entry(token_id.clone()).or_default(); + if !state.initialized && !state.dirty { + // Syncing: no snapshot has arrived since (re)subscription. Buffer + // the delta for ordered replay after the first book instead of + // dropping it while still poisoning last_timestamp, and bound the + // buffer by failing the token closed on overflow. + if state.pending.len() < MAX_POLYMARKET_CLOB_PENDING_CHANGES { + state.pending.push((change.timestamp, entry.clone())); + continue; + } + *state = ClobBookState::default(); + state.dirty = true; + } else { + state.apply(entry, change.timestamp); + } last_timestamp.insert(token_id.clone(), change.timestamp); if last_entries.insert(token_id.clone(), entry).is_none() { updated_tokens.push(token_id); @@ -1245,26 +1341,7 @@ fn market_updates_from_price_change( .into_iter() .flat_map(|token_id| { if books[&token_id].dirty { - let token_id: Arc = Arc::from(token_id); - vec![ - MarketUpdate::QuoteCollectionFailure { - token_id: Arc::clone(&token_id), - request_started_at: failed_at, - http_status: None, - error_kind: Arc::from("websocket_payload"), - ts: failed_at, - }, - MarketUpdate::Quote { - token_id, - bid: None, - ask: None, - bid_size: None, - ask_size: None, - bid_levels: Vec::new(), - ask_levels: Vec::new(), - ts: failed_at, - }, - ] + clob_failed_closed_updates(&token_id, failed_at) } else { vec![books[&token_id].quote( token_id.clone(), @@ -1382,11 +1459,13 @@ fn forward_clob_ws_payload( { return Err("Polymarket book source time moved backwards".to_string()); } - last_timestamp.insert(token_id.clone(), book.timestamp); - let state = books_by_token.entry(token_id).or_default(); - let update = market_update_from_clob_book(&book, state)?; - if tx.send(update).is_err() { - return Ok(false); + let state = books_by_token.entry(token_id.clone()).or_default(); + let (updates, watermark) = market_update_from_clob_book(&book, state)?; + last_timestamp.insert(token_id, watermark); + for update in updates { + if tx.send(update).is_err() { + return Ok(false); + } } } WsMessage::PriceChange(change) => { @@ -2240,7 +2319,7 @@ mod tests { mark_db_event_expired_if_resolved, market_update_from_clob_book, market_updates_from_price_change, parse_agg_trade_msg, parse_equity_price_payload, rtds_market_data_ws_config, send_quote_collection_failure_and_empty, ClobBookState, - RestBook, POLYMARKET_CLOB_CROSSED_BOOK_ERROR, U256, + RestBook, MAX_POLYMARKET_CLOB_PENDING_CHANGES, POLYMARKET_CLOB_CROSSED_BOOK_ERROR, U256, }; use chrono::Utc; use ploy_market_contracts::MarketUpdate; @@ -2298,8 +2377,9 @@ mod tests { })) .expect("valid CLOB book update"); - let update = market_update_from_clob_book(&book, &mut ClobBookState::default()) + let (updates, _) = market_update_from_clob_book(&book, &mut ClobBookState::default()) .expect("tradeable quote"); + let [update] = updates.try_into().ok().expect("single quote update"); let MarketUpdate::Quote { token_id, bid, @@ -2325,7 +2405,7 @@ mod tests { } #[test] - fn clob_price_change_becomes_immediate_bba_tick() { + fn clob_price_change_before_first_snapshot_is_buffered_without_publishing() { let change = serde_json::from_value(json!({ "market": "0x0000000000000000000000000000000000000000000000000000000000000000", "timestamp": "1712205600456", @@ -2341,35 +2421,15 @@ mod tests { })) .expect("valid price change"); - let updates = market_updates_from_price_change( - &change, - &mut std::collections::HashMap::new(), - &mut std::collections::HashMap::new(), - ) - .expect("valid price change"); - assert_eq!(updates.len(), 1); - let MarketUpdate::Quote { - token_id, - bid, - ask, - bid_size, - ask_size, - bid_levels, - ask_levels, - ts, - } = &updates[0] - else { - panic!("expected quote update"); - }; - - assert_eq!(token_id.as_ref(), "7"); - assert_eq!(*bid, Some(dec!(0.51))); - assert_eq!(*ask, Some(dec!(0.53))); - assert_eq!(*bid_size, None); - assert_eq!(*ask_size, Some(dec!(4))); - assert!(bid_levels.is_empty()); - assert!(ask_levels.is_empty()); - assert_eq!(ts.timestamp_millis(), 1_712_205_600_456); + let mut books = std::collections::HashMap::new(); + let mut timestamps = std::collections::HashMap::new(); + let updates = market_updates_from_price_change(&change, &mut books, &mut timestamps) + .expect("valid price change"); + assert!(updates.is_empty()); + assert_eq!(books["7"].pending.len(), 1); + assert!(!books["7"].initialized); + assert!(!books["7"].dirty); + assert!(!timestamps.contains_key("7")); } #[test] @@ -2404,11 +2464,11 @@ mod tests { })) .expect("valid empty CLOB book update"); - let update = market_update_from_clob_book(&book, &mut ClobBookState::default()) + let (updates, _) = market_update_from_clob_book(&book, &mut ClobBookState::default()) .expect("empty book is still a state transition"); assert!(matches!( - update, - MarketUpdate::Quote { + updates.as_slice(), + [MarketUpdate::Quote { bid: None, ask: None, bid_size: None, @@ -2416,7 +2476,7 @@ mod tests { bid_levels, ask_levels, .. - } if bid_levels.is_empty() && ask_levels.is_empty() + }] if bid_levels.is_empty() && ask_levels.is_empty() )); } @@ -2447,9 +2507,11 @@ mod tests { })) .expect("valid terminal-only CLOB book update"); + let (updates, _) = market_update_from_clob_book(&book, &mut ClobBookState::default()) + .expect("valid terminal-only book"); assert!(matches!( - market_update_from_clob_book(&book, &mut ClobBookState::default()).unwrap(), - MarketUpdate::Quote { + updates.as_slice(), + [MarketUpdate::Quote { bid: None, ask: None, bid_size: None, @@ -2457,7 +2519,7 @@ mod tests { bid_levels, ask_levels, .. - } if bid_levels.len() == 1 && ask_levels.len() == 1 + }] if bid_levels.len() == 1 && ask_levels.len() == 1 )); } @@ -2908,7 +2970,7 @@ mod tests { } #[test] - fn clob_stale_book_self_heals_dirty_token_but_fails_closed_for_healthy_token() { + fn clob_stale_book_resyncs_buffering_token_but_fails_closed_for_healthy_token() { let (tx, mut rx) = tokio::sync::broadcast::channel(8); let mut books = std::collections::HashMap::new(); let mut timestamps = std::collections::HashMap::new(); @@ -2919,25 +2981,33 @@ mod tests { assert!( forward_clob_ws_payload(early_change, &tx, &mut books, &mut timestamps) - .expect("delta arriving before any snapshot is tolerated") + .expect("delta arriving before any snapshot is buffered for replay") ); assert!(!books["7"].initialized); assert!( !books["7"].dirty, - "a new token without any snapshot is uninitialized but not dirty" + "a new token without any snapshot is syncing, not dirty" + ); + assert_eq!(books["7"].pending.len(), 1); + assert!( + !timestamps.contains_key("7"), + "a buffered delta must not poison last_timestamp" + ); + assert!( + matches!( + rx.try_recv(), + Err(tokio::sync::broadcast::error::TryRecvError::Empty) + ), + "a syncing token publishes nothing until its first snapshot" ); - assert!(matches!( - rx.try_recv().unwrap(), - MarketUpdate::Quote { token_id, .. } if token_id.as_ref() == "7" - )); assert!( - forward_clob_ws_payload(stale_book, &tx, &mut books, &mut timestamps).expect( - "an uninitialized token has no valid state to protect from an older snapshot" - ) + forward_clob_ws_payload(stale_book, &tx, &mut books, &mut timestamps) + .expect("a syncing token has no valid state to protect from an older snapshot") ); assert!(books["7"].initialized); assert!(!books["7"].dirty); + assert!(books["7"].pending.is_empty()); assert!(matches!( rx.try_recv().unwrap(), MarketUpdate::Quote { @@ -3004,6 +3074,542 @@ mod tests { assert!(books["7"].dirty); } + #[test] + fn clob_price_change_before_first_snapshot_is_buffered_then_replayed() { + let (tx, mut rx) = tokio::sync::broadcast::channel(8); + let mut books = std::collections::HashMap::new(); + let mut timestamps = std::collections::HashMap::new(); + let early_change = br#"{"event_type":"price_change","market":"0x0000000000000000000000000000000000000000000000000000000000000000","timestamp":"1712205600200","price_changes":[{"asset_id":"7","price":"0.55","size":"3","side":"BUY","hash":null,"best_bid":"0.55","best_ask":"0.60"}]}"#; + let snapshot = br#"{"event_type":"book","asset_id":"7","market":"0x0000000000000000000000000000000000000000000000000000000000000000","timestamp":"1712205600100","bids":[{"price":"0.50","size":"7"}],"asks":[{"price":"0.60","size":"9"}]}"#; + let follow_up = br#"{"event_type":"price_change","market":"0x0000000000000000000000000000000000000000000000000000000000000000","timestamp":"1712205600300","price_changes":[{"asset_id":"7","price":"0.54","size":"2","side":"BUY","hash":null,"best_bid":"0.55","best_ask":"0.60"}]}"#; + + assert!( + forward_clob_ws_payload(early_change, &tx, &mut books, &mut timestamps) + .expect("buffering a pre-snapshot delta must not error") + ); + assert!(matches!( + rx.try_recv(), + Err(tokio::sync::broadcast::error::TryRecvError::Empty) + )); + assert!(!timestamps.contains_key("7")); + + assert!( + forward_clob_ws_payload(snapshot, &tx, &mut books, &mut timestamps) + .expect("the first snapshot resyncs the buffered delta") + ); + assert!(books["7"].initialized); + assert!(books["7"].pending.is_empty()); + assert_eq!(timestamps["7"], 1_712_205_600_200); + assert!(matches!( + rx.try_recv().unwrap(), + MarketUpdate::Quote { + token_id, + bid: Some(bid), + ask: Some(ask), + .. + } if token_id.as_ref() == "7" && bid == dec!(0.55) && ask == dec!(0.60) + )); + + assert!( + forward_clob_ws_payload(follow_up, &tx, &mut books, &mut timestamps) + .expect("post-resync delta publishes normally") + ); + assert!(matches!( + rx.try_recv().unwrap(), + MarketUpdate::Quote { + token_id, + bid: Some(bid), + ask: Some(ask), + .. + } if token_id.as_ref() == "7" && bid == dec!(0.55) && ask == dec!(0.60) + )); + assert!(!books["7"].dirty); + } + + #[test] + fn clob_replayed_delta_fills_snapshot_gap_so_bba_reconciles() { + let change = serde_json::from_value(json!({ + "market": "0x0000000000000000000000000000000000000000000000000000000000000000", + "timestamp": "1712205600150", + "price_changes": [{ + "asset_id": "7", "price": "0.55", "size": "3", "side": "BUY", + "hash": null, "best_bid": "0.55", "best_ask": "0.60" + }] + })) + .expect("valid pre-snapshot price change"); + let mut books = std::collections::HashMap::new(); + let mut timestamps = std::collections::HashMap::new(); + + let updates = market_updates_from_price_change(&change, &mut books, &mut timestamps) + .expect("pre-snapshot delta is buffered"); + assert!(updates.is_empty()); + assert_eq!(books["7"].pending.len(), 1); + + let book = serde_json::from_value(json!({ + "asset_id": "7", + "market": "0x0000000000000000000000000000000000000000000000000000000000000000", + "timestamp": "1712205600100", + "bids": [{"price": "0.50", "size": "7"}], + "asks": [{"price": "0.60", "size": "9"}], + "hash": null + })) + .expect("valid CLOB book update"); + let (updates, watermark) = + market_update_from_clob_book(&book, books.get_mut("7").expect("syncing token state")) + .expect("snapshot replays the buffered delta"); + assert_eq!(watermark, 1_712_205_600_150); + timestamps.insert("7".to_string(), watermark); + assert!(matches!( + updates.as_slice(), + [MarketUpdate::Quote { + bid: Some(bid), + ask: Some(ask), + .. + }] if *bid == dec!(0.55) && *ask == dec!(0.60) + )); + + let report = serde_json::from_value(json!({ + "market": "0x0000000000000000000000000000000000000000000000000000000000000000", + "timestamp": "1712205600200", + "price_changes": [{ + "asset_id": "7", "price": "0.55", "size": "5", "side": "BUY", + "hash": null, "best_bid": "0.55", "best_ask": "0.60" + }] + })) + .expect("valid post-resync price change"); + let updates = market_updates_from_price_change(&report, &mut books, &mut timestamps) + .expect("reported BBA reconciles against the replayed level"); + assert!(matches!( + updates.as_slice(), + [MarketUpdate::Quote { bid: Some(bid), .. }] if *bid == dec!(0.55) + )); + assert!(!books["7"].dirty); + } + + #[test] + fn clob_replay_drops_changes_older_than_the_snapshot() { + let stale_change = serde_json::from_value(json!({ + "market": "0x0000000000000000000000000000000000000000000000000000000000000000", + "timestamp": "1712205600050", + "price_changes": [{ + "asset_id": "7", "price": "0.70", "size": "1", "side": "BUY", + "hash": null, "best_bid": "0.70", "best_ask": "0.75" + }] + })) + .expect("valid pre-snapshot price change"); + let mut books = std::collections::HashMap::new(); + let mut timestamps = std::collections::HashMap::new(); + + let updates = market_updates_from_price_change(&stale_change, &mut books, &mut timestamps) + .expect("pre-snapshot delta is buffered"); + assert!(updates.is_empty()); + + let book = serde_json::from_value(json!({ + "asset_id": "7", + "market": "0x0000000000000000000000000000000000000000000000000000000000000000", + "timestamp": "1712205600100", + "bids": [{"price": "0.50", "size": "7"}], + "asks": [{"price": "0.60", "size": "9"}], + "hash": null + })) + .expect("valid CLOB book update"); + let (updates, watermark) = + market_update_from_clob_book(&book, books.get_mut("7").expect("syncing token state")) + .expect("snapshot applies"); + assert_eq!(watermark, 1_712_205_600_100); + assert!(matches!( + updates.as_slice(), + [MarketUpdate::Quote { + bid: Some(bid), + ask: Some(ask), + .. + }] if *bid == dec!(0.50) && *ask == dec!(0.60) + )); + assert!(!books["7"].bids.contains_key(&dec!(0.70))); + assert!(books["7"].pending.is_empty()); + } + + #[test] + fn clob_replay_applies_same_millisecond_deltas_like_the_live_path() { + let stale_change = serde_json::from_value(json!({ + "market": "0x0000000000000000000000000000000000000000000000000000000000000000", + "timestamp": "1712205600050", + "price_changes": [{ + "asset_id": "7", "price": "0.70", "size": "1", "side": "BUY", + "hash": null, "best_bid": "0.70", "best_ask": "0.75" + }] + })) + .expect("valid older pre-snapshot price change"); + let same_millisecond = serde_json::from_value(json!({ + "market": "0x0000000000000000000000000000000000000000000000000000000000000000", + "timestamp": "1712205600100", + "price_changes": [{ + "asset_id": "7", "price": "0.55", "size": "3", "side": "BUY", + "hash": null, "best_bid": "0.55", "best_ask": "0.60" + }] + })) + .expect("valid same-millisecond pre-snapshot price change"); + let mut books = std::collections::HashMap::new(); + let mut timestamps = std::collections::HashMap::new(); + + market_updates_from_price_change(&stale_change, &mut books, &mut timestamps) + .expect("older delta is buffered"); + market_updates_from_price_change(&same_millisecond, &mut books, &mut timestamps) + .expect("same-millisecond delta is buffered"); + + let book = serde_json::from_value(json!({ + "asset_id": "7", + "market": "0x0000000000000000000000000000000000000000000000000000000000000000", + "timestamp": "1712205600100", + "bids": [{"price": "0.50", "size": "7"}], + "asks": [{"price": "0.60", "size": "9"}], + "hash": null + })) + .expect("valid CLOB book update"); + let (updates, watermark) = + market_update_from_clob_book(&book, books.get_mut("7").expect("syncing token state")) + .expect("snapshot applies"); + assert_eq!(watermark, 1_712_205_600_100); + assert!( + matches!( + updates.as_slice(), + [MarketUpdate::Quote { + bid: Some(bid), + ask: Some(ask), + .. + }] if *bid == dec!(0.55) && *ask == dec!(0.60) + ), + "a delta sharing the snapshot millisecond must replay like the live path" + ); + assert!(books["7"].bids.contains_key(&dec!(0.55))); + assert!(!books["7"].bids.contains_key(&dec!(0.70))); + assert!(books["7"].pending.is_empty()); + } + + #[test] + fn clob_replay_crossed_book_isolates_token_instead_of_publishing() { + let book = serde_json::from_value(json!({ + "asset_id": "7", + "market": "0x0000000000000000000000000000000000000000000000000000000000000000", + "timestamp": "1712205600100", + "bids": [{"price": "0.50", "size": "7"}], + "asks": [{"price": "0.60", "size": "9"}], + "hash": null + })) + .expect("valid CLOB book update"); + // Without BBA fields apply() cannot prune the opposite side, so the + // replayed level crosses the snapshot ask. + let crossing = serde_json::from_value(json!({ + "market": "0x0000000000000000000000000000000000000000000000000000000000000000", + "timestamp": "1712205600150", + "price_changes": [{ + "asset_id": "7", "price": "0.65", "size": "3", "side": "BUY", "hash": null + }] + })) + .expect("valid crossing pre-snapshot price change"); + let mut books = std::collections::HashMap::new(); + let mut timestamps = std::collections::HashMap::new(); + market_updates_from_price_change(&crossing, &mut books, &mut timestamps) + .expect("crossing delta is buffered"); + + let (updates, _) = + market_update_from_clob_book(&book, books.get_mut("7").expect("syncing token state")) + .expect("a crossed replay must isolate the token, not error the batch"); + assert!( + matches!( + updates.as_slice(), + [MarketUpdate::QuoteCollectionFailure { + token_id, + error_kind, + .. + }, MarketUpdate::Quote { + token_id: quote_token_id, + bid: None, + ask: None, + .. + }] if token_id.as_ref() == "7" + && error_kind.as_ref() == "websocket_payload" + && quote_token_id.as_ref() == "7" + ), + "a crossed replayed book must fail closed instead of publishing" + ); + assert!(books["7"].dirty); + } + + #[test] + fn clob_replayed_quote_carries_the_replay_watermark_timestamp() { + let change = serde_json::from_value(json!({ + "market": "0x0000000000000000000000000000000000000000000000000000000000000000", + "timestamp": "1712205600150", + "price_changes": [{ + "asset_id": "7", "price": "0.55", "size": "3", "side": "BUY", + "hash": null, "best_bid": "0.55", "best_ask": "0.60" + }] + })) + .expect("valid pre-snapshot price change"); + let mut books = std::collections::HashMap::new(); + let mut timestamps = std::collections::HashMap::new(); + market_updates_from_price_change(&change, &mut books, &mut timestamps) + .expect("pre-snapshot delta is buffered"); + + let book = serde_json::from_value(json!({ + "asset_id": "7", + "market": "0x0000000000000000000000000000000000000000000000000000000000000000", + "timestamp": "1712205600100", + "bids": [{"price": "0.50", "size": "7"}], + "asks": [{"price": "0.60", "size": "9"}], + "hash": null + })) + .expect("valid CLOB book update"); + let (updates, _) = + market_update_from_clob_book(&book, books.get_mut("7").expect("syncing token state")) + .expect("snapshot replays the buffered delta"); + assert!( + matches!( + updates.as_slice(), + [MarketUpdate::Quote { ts, .. }] if ts.timestamp_millis() == 1_712_205_600_150 + ), + "a quote built from replayed deltas must carry the replay watermark, \ + not the older snapshot timestamp" + ); + } + + #[test] + fn clob_replay_reconciles_every_buffered_message_not_just_the_last() { + let book = serde_json::from_value(json!({ + "asset_id": "7", + "market": "0x0000000000000000000000000000000000000000000000000000000000000000", + "timestamp": "1712205600100", + "bids": [{"price": "0.50", "size": "7"}], + "asks": [{"price": "0.60", "size": "9"}], + "hash": null + })) + .expect("valid CLOB book update"); + let missing_bba = serde_json::from_value(json!({ + "market": "0x0000000000000000000000000000000000000000000000000000000000000000", + "timestamp": "1712205600150", + "price_changes": [{ + "asset_id": "7", "price": "0.41", "size": "10", "side": "SELL", + "hash": null, "best_bid": "0.40", "best_ask": "0.41" + }] + })) + .expect("valid earlier message whose reported bid is missing locally"); + let no_bba = serde_json::from_value(json!({ + "market": "0x0000000000000000000000000000000000000000000000000000000000000000", + "timestamp": "1712205600200", + "price_changes": [{ + "asset_id": "7", "price": "0.55", "size": "3", "side": "SELL", "hash": null + }] + })) + .expect("valid later message without BBA fields"); + let mut books = std::collections::HashMap::new(); + let mut timestamps = std::collections::HashMap::new(); + market_updates_from_price_change(&missing_bba, &mut books, &mut timestamps) + .expect("earlier message is buffered"); + market_updates_from_price_change(&no_bba, &mut books, &mut timestamps) + .expect("later message is buffered"); + + let (updates, _) = + market_update_from_clob_book(&book, books.get_mut("7").expect("syncing token state")) + .expect("reconciliation isolates the token in-band"); + assert!( + matches!( + updates.as_slice(), + [MarketUpdate::QuoteCollectionFailure { token_id, .. }, MarketUpdate::Quote { + token_id: quote_token_id, + bid: None, + ask: None, + .. + }] if token_id.as_ref() == "7" && quote_token_id.as_ref() == "7" + ), + "an earlier message's missing-BBA evidence must not be drowned by a later message" + ); + assert!(books["7"].dirty); + } + + #[test] + fn clob_replay_skips_regressing_messages_and_keeps_watermark() { + let book = serde_json::from_value(json!({ + "asset_id": "7", + "market": "0x0000000000000000000000000000000000000000000000000000000000000000", + "timestamp": "1712205600100", + "bids": [{"price": "0.50", "size": "7"}], + "asks": [{"price": "0.60", "size": "9"}], + "hash": null + })) + .expect("valid CLOB book update"); + let newer = serde_json::from_value(json!({ + "market": "0x0000000000000000000000000000000000000000000000000000000000000000", + "timestamp": "1712205600300", + "price_changes": [{ + "asset_id": "7", "price": "0.55", "size": "3", "side": "BUY", + "hash": null, "best_bid": "0.55", "best_ask": "0.60" + }] + })) + .expect("valid newer buffered message"); + let regressing = serde_json::from_value(json!({ + "market": "0x0000000000000000000000000000000000000000000000000000000000000000", + "timestamp": "1712205600200", + "price_changes": [{ + "asset_id": "7", "price": "0.55", "size": "99", "side": "BUY", + "hash": null, "best_bid": "0.55", "best_ask": "0.60" + }] + })) + .expect("valid regressing buffered message"); + let mut books = std::collections::HashMap::new(); + let mut timestamps = std::collections::HashMap::new(); + market_updates_from_price_change(&newer, &mut books, &mut timestamps) + .expect("newer message is buffered"); + market_updates_from_price_change(®ressing, &mut books, &mut timestamps) + .expect("regressing message is buffered"); + + let (updates, watermark) = + market_update_from_clob_book(&book, books.get_mut("7").expect("syncing token state")) + .expect("snapshot replays non-regressing messages"); + assert_eq!(watermark, 1_712_205_600_300); + assert!( + matches!( + updates.as_slice(), + [MarketUpdate::Quote { + bid: Some(bid), + bid_size: Some(size), + .. + }] if *bid == dec!(0.55) && *size == dec!(3) + ), + "a regressing buffered message must not overwrite newer replayed state" + ); + assert_eq!(books["7"].bids[&dec!(0.55)], dec!(3)); + } + + #[test] + fn clob_resubscription_rebuilds_sync_state_without_backwards_errors() { + let (tx, mut rx) = tokio::sync::broadcast::channel(8); + let mut books = std::collections::HashMap::new(); + let mut timestamps = std::collections::HashMap::new(); + let first_book = br#"{"event_type":"book","asset_id":"7","market":"0x0000000000000000000000000000000000000000000000000000000000000000","timestamp":"1712205600300","bids":[{"price":"0.40","size":"7"}],"asks":[{"price":"0.60","size":"9"}]}"#; + let early_change = br#"{"event_type":"price_change","market":"0x0000000000000000000000000000000000000000000000000000000000000000","timestamp":"1712205600150","price_changes":[{"asset_id":"7","price":"0.55","size":"3","side":"BUY","hash":null,"best_bid":"0.55","best_ask":"0.60"}]}"#; + let resync_book = br#"{"event_type":"book","asset_id":"7","market":"0x0000000000000000000000000000000000000000000000000000000000000000","timestamp":"1712205600100","bids":[{"price":"0.50","size":"7"}],"asks":[{"price":"0.60","size":"9"}]}"#; + let stale_delta = br#"{"event_type":"price_change","market":"0x0000000000000000000000000000000000000000000000000000000000000000","timestamp":"1712205600120","price_changes":[{"asset_id":"7","price":"0.55","size":"5","side":"BUY","hash":null,"best_bid":"0.55","best_ask":"0.60"}]}"#; + let fresh_delta = br#"{"event_type":"price_change","market":"0x0000000000000000000000000000000000000000000000000000000000000000","timestamp":"1712205600160","price_changes":[{"asset_id":"7","price":"0.54","size":"2","side":"BUY","hash":null,"best_bid":"0.55","best_ask":"0.60"}]}"#; + + assert!(forward_clob_ws_payload(first_book, &tx, &mut books, &mut timestamps).unwrap()); + assert!(matches!( + rx.try_recv().unwrap(), + MarketUpdate::Quote { token_id, .. } if token_id.as_ref() == "7" + )); + + // The spawn loop clears both maps on reconnect; resubscription then + // restarts the token in syncing state. + books.clear(); + timestamps.clear(); + + assert!( + forward_clob_ws_payload(early_change, &tx, &mut books, &mut timestamps) + .expect("post-reconnect delta is buffered") + ); + assert!(matches!( + rx.try_recv(), + Err(tokio::sync::broadcast::error::TryRecvError::Empty) + )); + + assert!( + forward_clob_ws_payload(resync_book, &tx, &mut books, &mut timestamps) + .expect("resubscription resyncs without a backwards error") + ); + assert_eq!(timestamps["7"], 1_712_205_600_150); + assert!(matches!( + rx.try_recv().unwrap(), + MarketUpdate::Quote { + token_id, + bid: Some(bid), + ask: Some(ask), + .. + } if token_id.as_ref() == "7" && bid == dec!(0.55) && ask == dec!(0.60) + )); + + assert!( + forward_clob_ws_payload(stale_delta, &tx, &mut books, &mut timestamps) + .expect("replayed watermark makes an older delta skippable, not fatal") + ); + assert!(matches!( + rx.try_recv(), + Err(tokio::sync::broadcast::error::TryRecvError::Empty) + )); + + assert!( + forward_clob_ws_payload(fresh_delta, &tx, &mut books, &mut timestamps) + .expect("fresh delta publishes after resync") + ); + assert!(matches!( + rx.try_recv().unwrap(), + MarketUpdate::Quote { + token_id, + bid: Some(bid), + .. + } if token_id.as_ref() == "7" && bid == dec!(0.55) + )); + } + + #[test] + fn clob_syncing_buffer_overflow_fails_closed() { + let mut books = std::collections::HashMap::new(); + let mut timestamps = std::collections::HashMap::new(); + for i in 0..MAX_POLYMARKET_CLOB_PENDING_CHANGES { + let change = serde_json::from_value(json!({ + "market": "0x0000000000000000000000000000000000000000000000000000000000000000", + "timestamp": (1_712_205_600_000i64 + i as i64).to_string(), + "price_changes": [{ + "asset_id": "7", "price": "0.55", "size": "3", "side": "BUY", + "hash": null, "best_bid": "0.55", "best_ask": "0.60" + }] + })) + .expect("valid buffered price change"); + let updates = market_updates_from_price_change(&change, &mut books, &mut timestamps) + .expect("buffered delta"); + assert!(updates.is_empty()); + } + assert_eq!( + books["7"].pending.len(), + MAX_POLYMARKET_CLOB_PENDING_CHANGES + ); + + let overflow = serde_json::from_value(json!({ + "market": "0x0000000000000000000000000000000000000000000000000000000000000000", + "timestamp": (1_712_205_600_000i64 + + MAX_POLYMARKET_CLOB_PENDING_CHANGES as i64) + .to_string(), + "price_changes": [{ + "asset_id": "7", "price": "0.55", "size": "3", "side": "BUY", + "hash": null, "best_bid": "0.55", "best_ask": "0.60" + }] + })) + .expect("valid overflowing price change"); + let updates = market_updates_from_price_change(&overflow, &mut books, &mut timestamps) + .expect("overflow fails the token closed instead of growing the buffer"); + assert!(matches!( + updates.as_slice(), + [MarketUpdate::QuoteCollectionFailure { + token_id, + error_kind, + .. + }, MarketUpdate::Quote { + token_id: quote_token_id, + bid: None, + ask: None, + .. + }] if token_id.as_ref() == "7" + && error_kind.as_ref() == "websocket_payload" + && quote_token_id.as_ref() == "7" + )); + assert!(books["7"].dirty); + assert!(books["7"].pending.is_empty()); + + let updates = market_updates_from_price_change(&overflow, &mut books, &mut timestamps) + .expect("a dirty token stays failed closed"); + assert_eq!(updates.len(), 2); + assert!(books["7"].pending.is_empty()); + } + #[test] fn clob_cancellation_bba_prunes_only_more_competitive_stale_depth() { let book = serde_json::from_value(json!({