diff --git a/rust_hft/data-pipelines/core/src/binance_market_tape.rs b/rust_hft/data-pipelines/core/src/binance_market_tape.rs index 40d7c8282..5e8c447c4 100644 --- a/rust_hft/data-pipelines/core/src/binance_market_tape.rs +++ b/rust_hft/data-pipelines/core/src/binance_market_tape.rs @@ -1,6 +1,6 @@ //! Shared contract for the immutable Binance market tape. -use std::collections::{btree_map::Entry, BTreeMap, HashMap}; +use std::collections::{btree_map::Entry, BTreeMap, BTreeSet, HashMap}; use anyhow::{Context, Result}; use rust_decimal::Decimal; @@ -10,9 +10,279 @@ use serde_json::{Map, Value}; pub const LEGACY_LOB_TAPE_SCHEMA: &str = "binance.lob_tape.v2"; pub const MARKET_TAPE_SCHEMA: &str = "binance.market_tape.v1"; pub const AGGREGATE_TRADE_SUMMARY_CONTRACT: &str = "binance.aggregate_trade_summary.v1"; +pub const LOB_CONTINUITY_SUMMARY_CONTRACT: &str = "binance.lob_continuity.v1"; pub const MAX_SOURCE_LEAD_MS: u64 = 1_000; pub const MAX_SOURCE_DELAY_MS: u64 = 30_000; +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct LobContinuitySummary { + pub contract: String, + pub capture_session_id: String, + pub reconnect_boundary: bool, + pub sequence_gaps: u64, + pub source_time_rollbacks: u64, + pub declared_symbol_count: u64, + pub covered_symbol_count: u64, + pub missing_symbols: Vec, + pub symbols: BTreeMap, +} + +#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct SymbolLobContinuitySummary { + pub snapshot_seed_count: u64, + pub diff_count: u64, + pub checkpoint_count: u64, + pub first_update_id: Option, + pub last_update_id: Option, + pub first_source_time_ms: Option, + pub last_source_time_ms: Option, + pub first_received_at_ns: Option, + pub last_received_at_ns: Option, + pub min_source_latency_ms: Option, + pub max_source_latency_ms: Option, + pub min_bid_levels: Option, + pub max_bid_levels: Option, + pub min_ask_levels: Option, + pub max_ask_levels: Option, +} + +#[derive(Debug)] +pub struct LobContinuitySummaryBuilder { + declared_symbols: BTreeSet, + capture_session_id: Option, + reconnect_boundary: bool, + sequence_gaps: u64, + source_time_rollbacks: u64, + symbols: BTreeMap, +} + +impl LobContinuitySummaryBuilder { + pub fn new(symbols: impl IntoIterator) -> Result { + let declared_symbols = symbols.into_iter().collect::>(); + if declared_symbols.is_empty() + || declared_symbols + .iter() + .any(|symbol| symbol.is_empty() || symbol != &symbol.to_ascii_uppercase()) + { + anyhow::bail!("LOB continuity summary requires uppercase declared symbols"); + } + Ok(Self { + symbols: declared_symbols + .iter() + .cloned() + .map(|symbol| (symbol, SymbolLobContinuitySummary::default())) + .collect(), + declared_symbols, + capture_session_id: None, + reconnect_boundary: false, + sequence_gaps: 0, + source_time_rollbacks: 0, + }) + } + + pub fn observe(&mut self, raw: &Map) -> Result<()> { + let event_type = raw + .get("type") + .and_then(Value::as_str) + .context("LOB continuity row is missing type")?; + if let Some(session_id) = raw.get("session_id").and_then(Value::as_str) { + if session_id.is_empty() { + anyhow::bail!("LOB continuity row has an empty session_id"); + } + if self + .capture_session_id + .as_deref() + .is_some_and(|expected| expected != session_id) + { + anyhow::bail!("LOB continuity segment contains mixed session_id values"); + } + self.capture_session_id + .get_or_insert_with(|| session_id.to_owned()); + } + let received_at_ns = raw + .get("received_at_ns") + .and_then(Value::as_u64) + .context("LOB continuity row is missing received_at_ns")?; + match event_type { + "session_start" => self.reconnect_boundary = true, + "snapshot" => { + let symbol = self.row_symbol(raw)?; + let snapshot = raw + .get("snapshot") + .and_then(Value::as_object) + .context("LOB continuity snapshot has no payload")?; + let (bid_levels, ask_levels) = book_depth(snapshot)?; + let summary = self.symbol_mut(&symbol)?; + summary.snapshot_seed_count = summary + .snapshot_seed_count + .checked_add(1) + .context("LOB snapshot seed count overflow")?; + summary.observe_received_at(received_at_ns); + summary.observe_depth(bid_levels, ask_levels); + } + "diff" => { + let clock = DepthSourceClock::from_archived_event(raw, received_at_ns)?; + let summary = self.symbol_mut(&clock.symbol)?; + summary.diff_count = summary + .diff_count + .checked_add(1) + .context("LOB diff count overflow")?; + summary.first_update_id.get_or_insert(clock.first_update_id); + summary.last_update_id = Some(clock.final_update_id); + summary + .first_source_time_ms + .get_or_insert(clock.event_time_ms); + summary.last_source_time_ms = Some(clock.event_time_ms); + summary.observe_received_at(received_at_ns); + summary.observe_latency(source_latency_ms(received_at_ns, clock.event_time_ms)?); + } + "checkpoint" => { + let symbol = self.row_symbol(raw)?; + let (bid_levels, ask_levels) = book_depth(raw)?; + let is_seed = raw.get("reason").and_then(Value::as_str) == Some("segment_open"); + let summary = self.symbol_mut(&symbol)?; + if is_seed { + summary.snapshot_seed_count = summary + .snapshot_seed_count + .checked_add(1) + .context("LOB snapshot seed count overflow")?; + } + summary.checkpoint_count = summary + .checkpoint_count + .checked_add(1) + .context("LOB checkpoint count overflow")?; + summary.observe_received_at(received_at_ns); + summary.observe_depth(bid_levels, ask_levels); + } + "sequence_gap" => { + self.sequence_gaps = self + .sequence_gaps + .checked_add(1) + .context("LOB sequence gap count overflow")?; + if raw + .get("error") + .and_then(Value::as_str) + .is_some_and(|error| error.contains("source-time rollback")) + { + self.source_time_rollbacks = self + .source_time_rollbacks + .checked_add(1) + .context("LOB source-time rollback count overflow")?; + } + } + _ => {} + } + Ok(()) + } + + pub fn finish(self) -> Result { + let capture_session_id = self + .capture_session_id + .context("LOB continuity segment has no capture session")?; + let missing_symbols = self + .symbols + .iter() + .filter(|(_, summary)| { + summary.snapshot_seed_count == 0 + || summary.diff_count == 0 + || summary.checkpoint_count == 0 + }) + .map(|(symbol, _)| symbol.clone()) + .collect::>(); + let covered_symbol_count = + self.declared_symbols + .len() + .checked_sub(missing_symbols.len()) + .context("LOB covered symbol count underflow")? as u64; + Ok(LobContinuitySummary { + contract: LOB_CONTINUITY_SUMMARY_CONTRACT.to_owned(), + capture_session_id, + reconnect_boundary: self.reconnect_boundary, + sequence_gaps: self.sequence_gaps, + source_time_rollbacks: self.source_time_rollbacks, + declared_symbol_count: self.declared_symbols.len() as u64, + covered_symbol_count, + missing_symbols, + symbols: self.symbols, + }) + } + + fn row_symbol(&self, raw: &Map) -> Result { + let symbol = raw + .get("symbol") + .and_then(Value::as_str) + .context("LOB continuity row is missing symbol")? + .to_ascii_uppercase(); + if !self.declared_symbols.contains(&symbol) { + anyhow::bail!("LOB continuity row symbol is outside its declared scope"); + } + Ok(symbol) + } + + fn symbol_mut(&mut self, symbol: &str) -> Result<&mut SymbolLobContinuitySummary> { + self.symbols + .get_mut(symbol) + .context("LOB continuity row symbol is outside its declared scope") + } +} + +impl SymbolLobContinuitySummary { + fn observe_received_at(&mut self, received_at_ns: u64) { + self.first_received_at_ns.get_or_insert(received_at_ns); + self.last_received_at_ns = Some(received_at_ns); + } + + fn observe_latency(&mut self, latency_ms: i64) { + self.min_source_latency_ms = Some( + self.min_source_latency_ms + .map_or(latency_ms, |current| current.min(latency_ms)), + ); + self.max_source_latency_ms = Some( + self.max_source_latency_ms + .map_or(latency_ms, |current| current.max(latency_ms)), + ); + } + + fn observe_depth(&mut self, bid_levels: u64, ask_levels: u64) { + self.min_bid_levels = Some( + self.min_bid_levels + .map_or(bid_levels, |current| current.min(bid_levels)), + ); + self.max_bid_levels = Some( + self.max_bid_levels + .map_or(bid_levels, |current| current.max(bid_levels)), + ); + self.min_ask_levels = Some( + self.min_ask_levels + .map_or(ask_levels, |current| current.min(ask_levels)), + ); + self.max_ask_levels = Some( + self.max_ask_levels + .map_or(ask_levels, |current| current.max(ask_levels)), + ); + } +} + +fn book_depth(raw: &Map) -> Result<(u64, u64)> { + let bids = raw + .get("bids") + .and_then(Value::as_array) + .context("LOB book is missing bids")?; + let asks = raw + .get("asks") + .and_then(Value::as_array) + .context("LOB book is missing asks")?; + Ok((bids.len() as u64, asks.len() as u64)) +} + +fn source_latency_ms(received_at_ns: u64, source_time_ms: u64) -> Result { + let received_at_ms = received_at_ns / 1_000_000; + let latency = i128::from(received_at_ms) - i128::from(source_time_ms); + i64::try_from(latency).context("LOB source latency exceeds its numeric range") +} + pub fn supported_schema(schema: &str) -> bool { matches!(schema, LEGACY_LOB_TAPE_SCHEMA | MARKET_TAPE_SCHEMA) } diff --git a/rust_hft/data-pipelines/core/src/binance_market_tape_artifact.rs b/rust_hft/data-pipelines/core/src/binance_market_tape_artifact.rs index 008b45f93..75513c85f 100644 --- a/rust_hft/data-pipelines/core/src/binance_market_tape_artifact.rs +++ b/rust_hft/data-pipelines/core/src/binance_market_tape_artifact.rs @@ -20,7 +20,8 @@ use crate::binance_lob_replay::{ use crate::binance_market_tape::{ event_type_allowed, AggregateTrade, AggregateTradeSequenceValidator, AggregateTradeSummary, AggregateTradeSummaryBuilder, DepthSourceClock, DepthSourceClockSequenceValidator, - AGGREGATE_TRADE_SUMMARY_CONTRACT, MARKET_TAPE_SCHEMA, + LobContinuitySummary, LobContinuitySummaryBuilder, AGGREGATE_TRADE_SUMMARY_CONTRACT, + LOB_CONTINUITY_SUMMARY_CONTRACT, MARKET_TAPE_SCHEMA, }; const REPLAY_SCOPE: &str = @@ -179,18 +180,31 @@ pub fn seal_binance_market_tape_triplet( pub fn verify_binance_market_tape( sealed: Vec, ) -> Result { - verify_binance_market_tape_with_summary_requirement(sealed, false) + verify_binance_market_tape_with_requirements(sealed, false, false) } pub fn verify_binance_market_tape_with_required_trade_summaries( sealed: Vec, ) -> Result { - verify_binance_market_tape_with_summary_requirement(sealed, true) + verify_binance_market_tape_with_requirements(sealed, true, false) } -fn verify_binance_market_tape_with_summary_requirement( +pub fn verify_binance_market_tape_with_required_lob_continuity( + sealed: Vec, +) -> Result { + verify_binance_market_tape_with_requirements(sealed, false, true) +} + +pub fn verify_binance_market_tape_with_required_trade_and_lob_summaries( + sealed: Vec, +) -> Result { + verify_binance_market_tape_with_requirements(sealed, true, true) +} + +fn verify_binance_market_tape_with_requirements( mut sealed: Vec, require_trade_summaries: bool, + require_lob_continuity: bool, ) -> Result { if sealed.is_empty() { bail!("market-tape segment set is empty"); @@ -268,6 +282,7 @@ fn verify_binance_market_tape_with_summary_requirement( previous_segment_end_received_at_ns = Some(segment.manifest.end_received_at_ns); let mut counts = BTreeMap::::new(); let mut trade_summaries = AggregateTradeSummaryBuilder::default(); + let mut lob_continuity = LobContinuitySummaryBuilder::new(symbols.iter().cloned())?; let mut checkpoints = BTreeSet::new(); let mut snapshot_seeds = BTreeSet::new(); for (index, range) in segment.rows.iter().enumerate() { @@ -278,6 +293,7 @@ fn verify_binance_market_tape_with_summary_requirement( .ok_or_else(|| anyhow!("market-tape row must be an object"))?; let (event_type, row_session_id, received_at_ns) = validate_row(raw, &segment.manifest)?; + lob_continuity.observe(raw)?; if event_type == "session_start" && highest_received_at_ns.is_some_and(|last| received_at_ns < last) { @@ -382,6 +398,16 @@ fn verify_binance_market_tape_with_summary_requirement( if checkpoints != symbols { bail!("market-tape segment is missing a replay-safe checkpoint"); } + let lob_continuity = lob_continuity.finish()?; + match segment.manifest.lob_continuity.as_ref() { + Some(manifest) if manifest != &lob_continuity => { + bail!("market-tape LOB continuity summary does not match raw rows") + } + None if require_lob_continuity => { + bail!("market-tape segment is missing the LOB continuity contract") + } + _ => {} + } identities.push(segment.identity(market, trade_summaries)); } let aggregate_trade_symbols = aggregate_trades @@ -492,6 +518,8 @@ struct TapeManifest { price_surface_derivation: String, trade_summary_contract: Option, trade_summaries: Option>, + #[serde(default)] + lob_continuity: Option, } fn parse_manifest(bytes: &[u8]) -> Result { @@ -536,6 +564,10 @@ fn validate_manifest_identity( .trade_summary_contract .as_deref() .is_some_and(|contract| contract != AGGREGATE_TRADE_SUMMARY_CONTRACT) + || manifest + .lob_continuity + .as_ref() + .is_some_and(|summary| summary.contract != LOB_CONTINUITY_SUMMARY_CONTRACT) || manifest.venue_depth_complete || manifest.file != data_name || manifest.bytes != data_bytes as u64 @@ -869,6 +901,23 @@ mod tests { }) } + fn add_lob_continuity( + triplet: &BinanceMarketTapeTriplet, + rows: &[Value], + symbols: &[&str], + ) -> BinanceMarketTapeTrustAnchor { + let mut summary = + LobContinuitySummaryBuilder::new(symbols.iter().map(|symbol| (*symbol).to_owned())) + .unwrap(); + for row in rows { + summary.observe(row.as_object().unwrap()).unwrap(); + } + let summary = summary.finish().unwrap(); + rewrite_manifest(triplet, |manifest| { + manifest["lob_continuity"] = serde_json::to_value(summary).unwrap(); + }) + } + fn one_trade_summary(base_volume: &str) -> Value { json!({ "BTCUSDT":{ @@ -992,6 +1041,36 @@ mod tests { assert!(error.to_string().contains("summary contract")); } + #[test] + fn strict_lob_verifier_rejects_manifest_without_continuity_contract() { + let root = tempdir(); + let (triplet, anchor) = write_triplet(root.path(), &valid_rows()); + let sealed = seal_binance_market_tape_triplet(&triplet, &anchor).unwrap(); + + let error = + verify_binance_market_tape_with_required_lob_continuity(vec![sealed]).unwrap_err(); + assert!(error.to_string().contains("LOB continuity contract")); + } + + #[test] + fn strict_lob_verifier_rejects_manifest_latency_not_derived_from_raw_rows() { + let root = tempdir(); + let rows = valid_rows(); + let (triplet, _) = write_triplet(root.path(), &rows); + let valid_anchor = add_lob_continuity(&triplet, &rows, &["BTCUSDT"]); + let sealed = seal_binance_market_tape_triplet(&triplet, &valid_anchor).unwrap(); + verify_binance_market_tape_with_required_lob_continuity(vec![sealed]).unwrap(); + + let anchor = rewrite_manifest(&triplet, |manifest| { + manifest["lob_continuity"]["symbols"]["BTCUSDT"]["max_source_latency_ms"] = json!(999); + }); + let sealed = seal_binance_market_tape_triplet(&triplet, &anchor).unwrap(); + + let error = + verify_binance_market_tape_with_required_lob_continuity(vec![sealed]).unwrap_err(); + assert!(error.to_string().contains("does not match raw rows")); + } + #[test] fn declared_summary_contract_requires_manifest_summaries() { let root = tempdir(); diff --git a/rust_hft/tools/collector/src/lob_archiver.rs b/rust_hft/tools/collector/src/lob_archiver.rs index bf391e02d..82ecd591f 100644 --- a/rust_hft/tools/collector/src/lob_archiver.rs +++ b/rust_hft/tools/collector/src/lob_archiver.rs @@ -1,10 +1,11 @@ +use anyhow::Context; pub use data::binance_lob_replay::{ source_revision, Market, ReplaySequenceEvent, ReplaySequenceValidator, }; use data::binance_market_tape::{ event_type_allowed, supported_schema, AggregateTrade, AggregateTradeSequenceValidator, - AggregateTradeSummary, AggregateTradeSummaryBuilder, AGGREGATE_TRADE_SUMMARY_CONTRACT, - LEGACY_LOB_TAPE_SCHEMA, + AggregateTradeSummary, AggregateTradeSummaryBuilder, LobContinuitySummary, + LobContinuitySummaryBuilder, AGGREGATE_TRADE_SUMMARY_CONTRACT, LEGACY_LOB_TAPE_SCHEMA, }; use engine::binance_md::{parse_fixed_6, BookSync, SequenceDecision, UpdateMeta}; use rand::random; @@ -14,7 +15,7 @@ use serde_json::{json, Value}; use sha2::{Digest, Sha256}; use std::collections::{BTreeMap, BTreeSet, HashMap}; use std::fs::{self, File, OpenOptions}; -use std::io::{BufWriter, Read, Write}; +use std::io::{BufRead, BufReader, BufWriter, Read, Write}; use std::path::{Path, PathBuf}; use std::process::{Command, ExitStatus}; use std::str::FromStr; @@ -570,6 +571,11 @@ fn finalize_segment( RAW_SCHEMA => "captured_aggregate_trades_plus_snapshot_seed_plus_sequence_checked_diffs", _ => anyhow::bail!("unsupported recovered tape schema {schema}"), }; + let lob_continuity = if schema == RAW_SCHEMA { + Some(summarize_lob_continuity(path, config.symbols.clone())?) + } else { + None + }; let data = path.with_file_name(format!("part-{start_ns}.jsonl.zst")); let temporary_data = data.with_extension("zst.tmp"); let mut command = Command::new("zstd"); @@ -635,6 +641,10 @@ fn finalize_segment( "trade_summary_contract".to_owned(), AGGREGATE_TRADE_SUMMARY_CONTRACT.into(), ); + metadata.insert( + "lob_continuity".to_owned(), + serde_json::to_value(lob_continuity.context("market tape has no LOB summary")?)?, + ); } let manifest = data.with_file_name(format!( "{}.manifest.json", @@ -658,6 +668,23 @@ fn finalize_segment( }) } +fn summarize_lob_continuity( + path: &Path, + symbols: Vec, +) -> anyhow::Result { + let mut summary = LobContinuitySummaryBuilder::new(symbols)?; + for (index, line) in BufReader::new(File::open(path)?).lines().enumerate() { + let line = line?; + let raw: Value = serde_json::from_str(&line) + .with_context(|| format!("parse LOB continuity row {}", index + 1))?; + summary.observe( + raw.as_object() + .context("LOB continuity row must be an object")?, + )?; + } + summary.finish() +} + pub fn write_success_marker(data: &Path, digest: &str) -> anyhow::Result { let success = data.with_file_name(format!( "{}._SUCCESS", @@ -1342,18 +1369,59 @@ mod tests { snapshot_limit: 100, zstd_timeout: Duration::from_secs(30), }; - let mut segment = Segment::create(config, 1_700_000_000_000_000_000).unwrap(); + let start_ns = 1_700_000_000_000_000_000; + let mut segment = Segment::create(config, start_ns).unwrap(); segment - .write("snapshot", json!({"symbol":"BTCUSDT"}), segment.start_ns) + .write( + "session_start", + json!({ + "session_id":"session-1", + "market":"spot", + "symbols":1, + "websocket_shards":1 + }), + segment.start_ns, + ) + .unwrap(); + segment + .write( + "snapshot", + json!({ + "session_id":"session-1", + "symbol":"BTCUSDT", + "request_started_at_ns":segment.start_ns + 50_000_000, + "snapshot":{ + "lastUpdateId":100, + "bids":[["100","1"]], + "asks":[["101","1"]] + } + }), + segment.start_ns + 100_000_000, + ) .unwrap(); + let diff_received_at_ns = segment.start_ns + 200_000_000; segment .write( "diff", - json!({"symbol":"BTCUSDT","schema":"binance.market_tape.v999"}), - segment.start_ns + 1, + json!({ + "session_id":"session-1", + "frame":{ + "stream":"btcusdt@depth@100ms", + "data":{ + "e":"depthUpdate", + "E":diff_received_at_ns / 1_000_000, + "s":"BTCUSDT", + "U":101, + "u":101, + "b":[["100","2"]], + "a":[] + } + } + }), + diff_received_at_ns, ) .unwrap(); - let first_trade_received_at_ns = segment.start_ns + 200_000_000; + let first_trade_received_at_ns = segment.start_ns + 300_000_000; segment .write( "agg_trade", @@ -1378,7 +1446,7 @@ mod tests { first_trade_received_at_ns, ) .unwrap(); - let last_trade_received_at_ns = segment.start_ns + 300_000_000; + let last_trade_received_at_ns = segment.start_ns + 400_000_000; segment .write( "agg_trade", @@ -1406,8 +1474,18 @@ mod tests { segment .write( "checkpoint", - json!({"symbol":"BTCUSDT", "synced":true, "bridged":true}), - segment.start_ns + 400_000_000, + json!({ + "session_id":"session-1", + "symbol":"BTCUSDT", + "last_update_id":101, + "synced":true, + "bridged":true, + "bids":[["100","2"]], + "asks":[["101","1"]], + "reason":"test", + "replay_safe":true + }), + segment.start_ns + 500_000_000, ) .unwrap(); let artifacts = segment.close().unwrap().unwrap(); @@ -1419,7 +1497,39 @@ mod tests { ); assert_eq!( manifest["event_types"], - json!({"agg_trade":2,"checkpoint":1,"diff":1,"snapshot":1}) + json!({"agg_trade":2,"checkpoint":1,"diff":1,"session_start":1,"snapshot":1}) + ); + assert_eq!( + manifest["lob_continuity"], + json!({ + "contract":"binance.lob_continuity.v1", + "capture_session_id":"session-1", + "reconnect_boundary":true, + "sequence_gaps":0, + "source_time_rollbacks":0, + "declared_symbol_count":1, + "covered_symbol_count":1, + "missing_symbols":[], + "symbols":{ + "BTCUSDT":{ + "snapshot_seed_count":1, + "diff_count":1, + "checkpoint_count":1, + "first_update_id":101, + "last_update_id":101, + "first_source_time_ms":diff_received_at_ns / 1_000_000, + "last_source_time_ms":diff_received_at_ns / 1_000_000, + "first_received_at_ns":start_ns + 100_000_000, + "last_received_at_ns":start_ns + 500_000_000, + "min_source_latency_ms":0, + "max_source_latency_ms":0, + "min_bid_levels":1, + "max_bid_levels":1, + "min_ask_levels":1, + "max_ask_levels":1 + } + } + }) ); assert_eq!( manifest["trade_summaries"]["BTCUSDT"], @@ -1475,7 +1585,13 @@ mod tests { rows.iter() .map(|row| row["type"].as_str().unwrap()) .collect::>(), - BTreeSet::from(["agg_trade", "checkpoint", "diff", "snapshot"]) + BTreeSet::from([ + "agg_trade", + "checkpoint", + "diff", + "session_start", + "snapshot" + ]) ); assert!(!artifacts.success.exists()); write_success_marker(&artifacts.data, &artifacts.sha256).unwrap();