Skip to main content

bulk_client/msgs/
md.rs

1//! Incoming market-data message types.
2//!
3//! These structs are deserialized from the WebSocket JSON feed and correspond
4//! to the Python definitions in `md.py`.
5
6use crate::common::side::Side;
7use crate::transaction::ActionMeta;
8use serde::{Deserialize, Deserializer, Serialize};
9
10// ============================================================================
11// Matrix MD
12// ============================================================================
13
14/// Matrix with named columns and rows
15#[derive(Clone, Debug, Deserialize, Serialize)]
16pub struct Matrix {
17    /// named labels for the matrix
18    pub index: Vec<String>,
19    pub matrix: Vec<Vec<f64>>,
20
21    #[serde(skip)]
22    pub meta: ActionMeta,
23}
24
25// ============================================================================
26// Summary-level Data
27// ============================================================================
28
29/// Market ticker data.
30///
31/// Received on the `"ticker"` channel.
32#[derive(Debug, Clone, Deserialize)]
33#[serde(rename_all = "camelCase")]
34#[allow(unused)]
35pub struct Ticker {
36    pub symbol: String,
37    #[serde(deserialize_with = "f64_or_nan")]
38    pub last_price: f64,
39    #[serde(deserialize_with = "f64_or_nan")]
40    pub mark_price: f64,
41    #[serde(deserialize_with = "f64_or_nan")]
42    pub oracle_price: f64,
43    #[serde(deserialize_with = "f64_or_nan")]
44    pub price_change: f64,
45    #[serde(deserialize_with = "f64_or_nan")]
46    pub price_change_percent: f64,
47    #[serde(deserialize_with = "f64_or_nan")]
48    pub high_price: f64,
49    #[serde(deserialize_with = "f64_or_nan")]
50    pub low_price: f64,
51    #[serde(deserialize_with = "f64_or_nan")]
52    pub volume: f64,
53    #[serde(deserialize_with = "f64_or_nan")]
54    pub quote_volume: f64,
55    #[serde(deserialize_with = "f64_or_nan")]
56    pub open_interest: f64,
57    #[serde(deserialize_with = "f64_or_nan")]
58    pub funding_rate: f64,
59}
60
61// ============================================================================
62// Candles
63// ============================================================================
64
65/// OHLCV candlestick data.
66///
67/// Received on the `"candle"` channel.
68/// The `symbol` and `interval` are populated from the subscription topic,
69/// not from the candle payload itself.
70#[derive(Debug, Clone, Deserialize)]
71#[allow(unused)]
72pub struct Candle {
73    #[serde(rename = "t")]
74    pub open_time: u64,
75    #[serde(rename = "T")]
76    pub close_time: u64,
77    #[serde(skip)]
78    pub symbol: String,
79    #[serde(skip)]
80    pub interval: String,
81    #[serde(rename = "o")]
82    pub open: f64,
83    #[serde(rename = "h")]
84    pub high: f64,
85    #[serde(rename = "l")]
86    pub low: f64,
87    #[serde(rename = "c")]
88    pub close: f64,
89    #[serde(rename = "v")]
90    pub volume: f64,
91    #[serde(rename = "n")]
92    pub num_trades: u64,
93}
94
95// ============================================================================
96// Trades
97// ============================================================================
98
99/// A single public trade.
100///
101/// Received on the `"trades"` channel (as a list).
102#[derive(Debug, Clone, Deserialize)]
103#[allow(unused)]
104pub struct Trade {
105    #[serde(rename = "time")]
106    pub timestamp: u64,
107    #[serde(rename = "s")]
108    pub symbol: String,
109    #[serde(rename = "b")]
110    pub side: Side,
111    #[serde(rename = "sz")]
112    pub size: f64,
113    #[serde(rename = "px")]
114    pub price: f64,
115    pub maker: String,
116    pub taker: String,
117}
118
119// ============================================================================
120// Order Book
121// ============================================================================
122
123/// Single price level in the order book.
124#[derive(Debug, Clone, PartialEq, Deserialize)]
125#[allow(unused)]
126pub struct OrderBookLevel {
127    /// Price
128    #[serde(rename = "px")]
129    pub price: f64,
130    /// Size (aggregate quantity at this price)
131    #[serde(rename = "sz")]
132    pub size: f64,
133    /// Number of orders at this level
134    #[serde(rename = "n")]
135    pub num_orders: u32,
136}
137
138impl std::fmt::Display for OrderBookLevel {
139    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
140        write!(f, "{} @ {}", self.size, self.price)
141    }
142}
143
144/// Full L2 order-book snapshot.
145///
146/// Received on the `"l2Snapshot"` channel.
147#[derive(Debug, Clone, Deserialize)]
148#[allow(unused)]
149pub struct L2Snapshot {
150    pub timestamp: u64,
151    pub symbol: String,
152    /// `[[bid_levels], [ask_levels]]`
153    pub levels: (Vec<OrderBookLevel>, Vec<OrderBookLevel>),
154}
155
156// ============================================================================
157// Helpers
158// ============================================================================
159
160/// Deserialize an `f64` that may be `null` in JSON, mapping `null` → `NaN`.
161pub fn f64_or_nan<'de, D>(deserializer: D) -> Result<f64, D::Error>
162where
163    D: Deserializer<'de>,
164{
165    Ok(Option::<f64>::deserialize(deserializer)?.unwrap_or(f64::NAN))
166}
167
168// ============================================================================
169// Optional SDK integration
170// ============================================================================
171
172#[cfg(feature = "with-sdk")]
173use bulk_sdk_core::{
174    data::L2Delta as SDKL2Delta, data::L2Snapshot as SDKL2Snapshot,
175    data::PriceLevel as SDKPriceLevel, markets::MktId,
176};
177
178#[cfg(feature = "with-sdk")]
179impl From<&L2Snapshot> for SDKL2Snapshot {
180    fn from(snap: &L2Snapshot) -> Self {
181        let (bids, asks) = &snap.levels;
182
183        let instrument = MktId::new(snap.symbol.as_str()).unwrap();
184        let newbids = bids
185            .iter()
186            .map(|x| SDKPriceLevel {
187                amount: x.size,
188                price: x.price,
189                num_orders: x.num_orders,
190                cum_vwap: 0.0,
191                cum_amount: 0.0,
192            })
193            .collect();
194        let newasks = asks
195            .iter()
196            .map(|x| SDKPriceLevel {
197                amount: x.size,
198                price: x.price,
199                num_orders: x.num_orders,
200                cum_vwap: 0.0,
201                cum_amount: 0.0,
202            })
203            .collect();
204
205        SDKL2Snapshot {
206            stamp: snap.timestamp,
207            instrument,
208            bids: newbids,
209            asks: newasks,
210            trackable_id: Default::default(),
211        }
212    }
213}
214
215#[cfg(feature = "with-sdk")]
216impl From<&OrderBookLevel> for SDKL2Delta {
217    fn from(level: &OrderBookLevel) -> Self {
218        SDKL2Delta {
219            stamp: 0,
220            instrument: Default::default(),
221            side: Default::default(),
222            amount: level.size,
223            price: level.price,
224        }
225    }
226}
227
228//
229// Unit Tests
230//
231
232#[cfg(test)]
233mod tests {
234    use super::*;
235
236    #[test]
237    fn test_ticker_deserialize_with_nulls() {
238        let json = r#"{
239            "symbol": "BTC-USD",
240            "priceChange": 0.0,
241            "priceChangePercent": 0.0,
242            "lastPrice": 100000.1824189499,
243            "highPrice": 100000.52095557356,
244            "lowPrice": 99999.62809703613,
245            "volume": 0.0,
246            "quoteVolume": 0.0,
247            "markPrice": null,
248            "oraclePrice": null,
249            "openInterest": 0.0,
250            "fundingRate": 0.0
251        }"#;
252
253        let ticker: Ticker = serde_json::from_str(json).unwrap();
254
255        assert_eq!(ticker.symbol, "BTC-USD");
256        assert!((ticker.last_price - 100000.1824189499).abs() < 1e-6);
257        assert!((ticker.high_price - 100000.52095557356).abs() < 1e-6);
258        assert!((ticker.low_price - 99999.62809703613).abs() < 1e-6);
259        assert_eq!(ticker.price_change, 0.0);
260        assert_eq!(ticker.price_change_percent, 0.0);
261        assert_eq!(ticker.volume, 0.0);
262        assert_eq!(ticker.quote_volume, 0.0);
263        assert_eq!(ticker.open_interest, 0.0);
264        assert_eq!(ticker.funding_rate, 0.0);
265
266        // null → NaN
267        assert!(ticker.mark_price.is_nan());
268        assert!(ticker.oracle_price.is_nan());
269    }
270
271    #[test]
272    fn test_ticker_deserialize_with_values() {
273        let json = r#"{
274            "symbol": "ETH-USD",
275            "priceChange": 10.5,
276            "priceChangePercent": 0.33,
277            "lastPrice": 3200.0,
278            "highPrice": 3250.0,
279            "lowPrice": 3150.0,
280            "volume": 1234.56,
281            "quoteVolume": 3950000.0,
282            "markPrice": 3201.5,
283            "oraclePrice": 3200.8,
284            "openInterest": 50000.0,
285            "fundingRate": 0.0001
286        }"#;
287
288        let ticker: Ticker = serde_json::from_str(json).unwrap();
289
290        assert_eq!(ticker.symbol, "ETH-USD");
291        assert!((ticker.mark_price - 3201.5).abs() < 1e-6);
292        assert!((ticker.oracle_price - 3200.8).abs() < 1e-6);
293        assert!(!ticker.mark_price.is_nan());
294        assert!(!ticker.oracle_price.is_nan());
295    }
296
297    #[test]
298    fn test_l2_snapshot_deserialize() {
299        let json = serde_json::json!({
300            "timestamp": 1770906894450242133u64,
301            "symbol": "ETH-USD",
302            "updateType": "snapshot",
303            "levels": [
304                [
305                    {"px": 1988.36, "sz": 50.0000001, "n": 5},
306                    {"px": 1988.35, "sz": 71.3240001, "n": 5},
307                    {"px": 1988.34, "sz": 40.00000006, "n": 4},
308                    {"px": 1988.32, "sz": 102.8540001, "n": 5},
309                    {"px": 1988.31, "sz": 30.00000003, "n": 3},
310                    {"px": 1988.30, "sz": 30.00000003, "n": 3},
311                    {"px": 1988.29, "sz": 50.8950001, "n": 5},
312                    {"px": 1988.28, "sz": 64.4280001, "n": 5},
313                    {"px": 1988.27, "sz": 53.7410001, "n": 5},
314                    {"px": 1988.25, "sz": 121.9040001, "n": 5}
315                ],
316                [
317                    {"px": 1988.47, "sz": 10.0, "n": 1},
318                    {"px": 1988.75, "sz": 70.0, "n": 7},
319                    {"px": 1989.00, "sz": 230.00000006, "n": 23},
320                    {"px": 1989.25, "sz": 140.0, "n": 14},
321                    {"px": 1989.50, "sz": 440.25680003, "n": 44},
322                    {"px": 1989.75, "sz": 230.00000007, "n": 23},
323                    {"px": 1990.00, "sz": 200.0, "n": 20},
324                    {"px": 1990.25, "sz": 230.00000009, "n": 23},
325                    {"px": 1990.50, "sz": 340.0, "n": 34},
326                    {"px": 1990.75, "sz": 250.0, "n": 25}
327                ]
328            ]
329        });
330
331        let snap: L2Snapshot = serde_json::from_value(json).expect("deserialize L2Snapshot");
332
333        assert_eq!(snap.symbol, "ETH-USD");
334        assert_eq!(snap.timestamp, 1770906894450242133);
335
336        let (bids, asks) = &snap.levels;
337        assert_eq!(bids.len(), 10);
338        assert_eq!(asks.len(), 10);
339
340        // Best bid
341        assert_eq!(bids[0].price, 1988.36);
342        assert_eq!(bids[0].size, 50.0000001);
343        assert_eq!(bids[0].num_orders, 5);
344
345        // Best ask
346        assert_eq!(asks[0].price, 1988.47);
347        assert_eq!(asks[0].size, 10.0);
348        assert_eq!(asks[0].num_orders, 1);
349
350        // Last bid
351        assert_eq!(bids[9].price, 1988.25);
352        assert_eq!(bids[9].size, 121.9040001);
353        assert_eq!(bids[9].num_orders, 5);
354
355        // Last ask
356        assert_eq!(asks[9].price, 1990.75);
357        assert_eq!(asks[9].size, 250.0);
358        assert_eq!(asks[9].num_orders, 25);
359    }
360}