use serde_json::Value;
use crate::actors::{
CandleData, DataMessage, ExchangeConnector, FundingData, OrderBookData, TickerData, TradeData,
TradeSide, WebSocketConfig,
};
use crate::error::Result;
const SPOT_WS_BASE: &str = "wss://stream.binance.com:9443";
const FUTURES_WS_BASE: &str = "wss://fstream.binance.com";
const EXCHANGE_SPOT: &str = "binance";
const EXCHANGE_FUTURES: &str = "binance-futures";
#[derive(Debug, Clone)]
pub struct BinanceConnector {
pub url: String,
pub exchange: &'static str,
}
impl BinanceConnector {
#[must_use]
pub fn spot(streams: &[&str]) -> Self {
Self {
url: build_combined_url(SPOT_WS_BASE, streams),
exchange: EXCHANGE_SPOT,
}
}
#[must_use]
pub fn futures(streams: &[&str]) -> Self {
Self {
url: build_combined_url(FUTURES_WS_BASE, streams),
exchange: EXCHANGE_FUTURES,
}
}
#[must_use]
pub fn with_url(url: impl Into<String>, exchange: &'static str) -> Self {
Self {
url: url.into(),
exchange,
}
}
#[must_use]
pub fn trade_stream(symbol: &str) -> String {
format!("{}@aggTrade", symbol.to_lowercase())
}
#[must_use]
pub fn ticker_stream(symbol: &str) -> String {
format!("{}@bookTicker", symbol.to_lowercase())
}
#[must_use]
pub fn kline_stream(symbol: &str, interval: &str) -> String {
format!("{}@kline_{interval}", symbol.to_lowercase())
}
#[must_use]
pub fn depth_stream(symbol: &str) -> String {
format!("{}@depth@100ms", symbol.to_lowercase())
}
#[must_use]
pub fn depth_snapshot_stream(symbol: &str, levels: u8) -> String {
let n = match levels {
0..=5 => 5,
6..=10 => 10,
_ => 20,
};
format!("{}@depth{n}@100ms", symbol.to_lowercase())
}
#[must_use]
pub fn mark_price_stream(symbol: &str) -> String {
format!("{}@markPrice@1s", symbol.to_lowercase())
}
}
fn build_combined_url(base: &str, streams: &[&str]) -> String {
let joined = streams.join("/");
format!("{base}/stream?streams={joined}")
}
impl ExchangeConnector for BinanceConnector {
fn exchange_name(&self) -> &str {
self.exchange
}
fn ws_url(&self) -> &str {
&self.url
}
fn build_ws_config(&self, symbol: &str) -> WebSocketConfig {
WebSocketConfig {
url: self.url.clone(),
exchange: self.exchange.to_string(),
symbol: symbol.to_string(),
subscription_msg: None,
ping_interval_secs: 180,
reconnect_delay_secs: 5,
max_reconnect_attempts: 5,
}
}
fn subscription_message(&self, _symbol: &str) -> Option<String> {
None
}
fn parse_message(&self, raw: &str) -> Result<Vec<DataMessage>> {
let json: Value = serde_json::from_str(raw)?;
let (stream_name, inner) = if let Some(stream) = json.get("stream").and_then(|v| v.as_str())
{
let data = json.get("data").cloned().unwrap_or(Value::Null);
(Some(stream.to_string()), data)
} else {
(None, json)
};
let event_type = inner.get("e").and_then(|v| v.as_str()).unwrap_or("");
match event_type {
"aggTrade" => Ok(parse_agg_trade(self.exchange, &inner)),
"kline" => Ok(parse_kline(self.exchange, &inner)),
"markPriceUpdate" => Ok(parse_mark_price(self.exchange, &inner)),
"depthUpdate" => Ok(parse_depth_update(self.exchange, &inner)),
"" => {
if inner.get("u").is_some() && inner.get("s").is_some() && inner.get("b").is_some()
{
Ok(parse_book_ticker(self.exchange, &inner))
} else if inner.get("lastUpdateId").is_some() {
let symbol = stream_name
.as_ref()
.and_then(|s| s.split('@').next())
.unwrap_or("")
.to_uppercase();
Ok(parse_depth_snapshot(self.exchange, &symbol, &inner))
} else {
Ok(vec![])
}
}
_ => Ok(vec![]), }
}
}
fn now_ms() -> i64 {
chrono::Utc::now().timestamp_millis()
}
fn str_f64(v: &Value, key: &str) -> f64 {
v.get(key)
.and_then(|x| x.as_str())
.and_then(|s| s.parse().ok())
.unwrap_or(0.0)
}
fn opt_str_f64(v: &Value, key: &str) -> Option<f64> {
v.get(key)
.and_then(|x| x.as_str())
.and_then(|s| s.parse().ok())
}
fn parse_levels(v: &Value) -> Vec<[f64; 2]> {
v.as_array()
.map(|arr| {
arr.iter()
.filter_map(|row| {
let p: f64 = row.get(0)?.as_str()?.parse().ok()?;
let q: f64 = row.get(1)?.as_str()?.parse().ok()?;
Some([p, q])
})
.collect()
})
.unwrap_or_default()
}
fn parse_agg_trade(exchange: &str, data: &Value) -> Vec<DataMessage> {
let symbol = data["s"].as_str().unwrap_or("").to_string();
let is_buyer_maker = data["m"].as_bool().unwrap_or(false);
let side = if is_buyer_maker {
TradeSide::Sell
} else {
TradeSide::Buy
};
vec![DataMessage::Trade(TradeData {
symbol,
exchange: exchange.to_string(),
side,
price: str_f64(data, "p"),
amount: str_f64(data, "q"),
exchange_ts: data["T"].as_i64().unwrap_or(0),
receipt_ts: now_ms(),
trade_id: data["a"].as_u64().unwrap_or(0).to_string(),
})]
}
fn parse_kline(exchange: &str, data: &Value) -> Vec<DataMessage> {
let Some(k) = data.get("k") else {
return vec![];
};
vec![DataMessage::Candle(CandleData {
symbol: k["s"].as_str().unwrap_or("").to_string(),
exchange: exchange.to_string(),
interval: k["i"].as_str().unwrap_or("").to_string(),
open_ts: k["t"].as_i64().unwrap_or(0),
open: str_f64(k, "o"),
high: str_f64(k, "h"),
low: str_f64(k, "l"),
close: str_f64(k, "c"),
volume: str_f64(k, "v"),
is_closed: k["x"].as_bool().unwrap_or(false),
receipt_ts: now_ms(),
})]
}
fn parse_mark_price(exchange: &str, data: &Value) -> Vec<DataMessage> {
vec![DataMessage::FundingRate(FundingData {
symbol: data["s"].as_str().unwrap_or("").to_string(),
exchange: exchange.to_string(),
funding_rate: str_f64(data, "r"),
next_funding_time: data["T"].as_i64().unwrap_or(0),
mark_price: opt_str_f64(data, "p"),
index_price: opt_str_f64(data, "i"),
exchange_ts: data["E"].as_i64().unwrap_or(0),
receipt_ts: now_ms(),
})]
}
fn parse_depth_update(exchange: &str, data: &Value) -> Vec<DataMessage> {
vec![DataMessage::OrderBook(OrderBookData {
symbol: data["s"].as_str().unwrap_or("").to_string(),
exchange: exchange.to_string(),
asks: parse_levels(&data["a"]),
bids: parse_levels(&data["b"]),
exchange_ts: data["E"].as_i64().unwrap_or(0),
receipt_ts: now_ms(),
is_snapshot: false,
})]
}
fn parse_book_ticker(exchange: &str, data: &Value) -> Vec<DataMessage> {
let now = now_ms();
vec![DataMessage::Ticker(TickerData {
symbol: data["s"].as_str().unwrap_or("").to_string(),
exchange: exchange.to_string(),
price: 0.0,
best_bid: str_f64(data, "b"),
best_ask: str_f64(data, "a"),
exchange_ts: now,
receipt_ts: now,
})]
}
fn parse_depth_snapshot(exchange: &str, symbol: &str, data: &Value) -> Vec<DataMessage> {
let now = now_ms();
vec![DataMessage::OrderBook(OrderBookData {
symbol: symbol.to_string(),
exchange: exchange.to_string(),
asks: parse_levels(&data["asks"]),
bids: parse_levels(&data["bids"]),
exchange_ts: now,
receipt_ts: now,
is_snapshot: true,
})]
}
#[cfg(test)]
mod tests {
use super::*;
fn connector() -> BinanceConnector {
BinanceConnector::spot(&["btcusdt@aggTrade"])
}
#[test]
fn parse_combined_stream_aggtrade() {
let raw = r#"{
"stream": "btcusdt@aggTrade",
"data": {
"e": "aggTrade",
"E": 1700000000000,
"s": "BTCUSDT",
"a": 12345,
"p": "96000.50",
"q": "0.05",
"f": 100,
"l": 102,
"T": 1700000000050,
"m": true,
"M": true
}
}"#;
let msgs = connector().parse_message(raw).expect("parse");
assert_eq!(msgs.len(), 1);
match &msgs[0] {
DataMessage::Trade(t) => {
assert_eq!(t.symbol, "BTCUSDT");
assert_eq!(t.exchange, "binance");
assert_eq!(t.side, TradeSide::Sell); assert!((t.price - 96_000.5).abs() < 1e-9);
assert!((t.amount - 0.05).abs() < 1e-12);
assert_eq!(t.exchange_ts, 1_700_000_000_050);
assert_eq!(t.trade_id, "12345");
}
other => panic!("expected Trade, got {other:?}"),
}
}
#[test]
fn parse_bookticker_into_ticker() {
let raw = r#"{
"stream": "btcusdt@bookTicker",
"data": {
"u": 400900217,
"s": "BTCUSDT",
"b": "96000.10",
"B": "1.5",
"a": "96001.00",
"A": "0.8"
}
}"#;
let msgs = connector().parse_message(raw).expect("parse");
match &msgs[0] {
DataMessage::Ticker(t) => {
assert_eq!(t.symbol, "BTCUSDT");
assert!((t.best_bid - 96_000.1).abs() < 1e-9);
assert!((t.best_ask - 96_001.0).abs() < 1e-9);
assert!((t.price - 0.0).abs() < 1e-12);
}
other => panic!("expected Ticker, got {other:?}"),
}
}
#[test]
fn parse_kline_into_candle() {
let raw = r#"{
"stream": "btcusdt@kline_1m",
"data": {
"e": "kline",
"E": 1700000000000,
"s": "BTCUSDT",
"k": {
"t": 1700000000000,
"T": 1700000059999,
"s": "BTCUSDT",
"i": "1m",
"f": 100,
"L": 200,
"o": "96000.00",
"c": "96100.00",
"h": "96200.00",
"l": "95900.00",
"v": "100.5",
"n": 250,
"x": true,
"q": "9650000.0",
"V": "50.0",
"Q": "4800000.0",
"B": "0"
}
}
}"#;
let msgs = connector().parse_message(raw).expect("parse");
match &msgs[0] {
DataMessage::Candle(c) => {
assert_eq!(c.symbol, "BTCUSDT");
assert_eq!(c.interval, "1m");
assert!((c.open - 96_000.0).abs() < 1e-9);
assert!((c.close - 96_100.0).abs() < 1e-9);
assert!(c.is_closed);
}
other => panic!("expected Candle, got {other:?}"),
}
}
#[test]
fn parse_depth_update_into_orderbook_delta() {
let raw = r#"{
"stream": "btcusdt@depth@100ms",
"data": {
"e": "depthUpdate",
"E": 1700000000000,
"s": "BTCUSDT",
"U": 157,
"u": 160,
"b": [["96000.00", "1.5"], ["95999.50", "0.0"]],
"a": [["96001.00", "0.5"]]
}
}"#;
let msgs = connector().parse_message(raw).expect("parse");
match &msgs[0] {
DataMessage::OrderBook(ob) => {
assert_eq!(ob.symbol, "BTCUSDT");
assert!(!ob.is_snapshot);
assert_eq!(ob.bids.len(), 2);
assert!((ob.bids[1][1] - 0.0).abs() < 1e-12); assert_eq!(ob.asks.len(), 1);
}
other => panic!("expected OrderBook, got {other:?}"),
}
}
#[test]
fn parse_depth_snapshot_uses_stream_name_for_symbol() {
let raw = r#"{
"stream": "btcusdt@depth5@100ms",
"data": {
"lastUpdateId": 999,
"bids": [["96000.00", "1.0"]],
"asks": [["96001.00", "0.5"]]
}
}"#;
let msgs = connector().parse_message(raw).expect("parse");
match &msgs[0] {
DataMessage::OrderBook(ob) => {
assert_eq!(ob.symbol, "BTCUSDT");
assert!(ob.is_snapshot);
assert_eq!(ob.bids.len(), 1);
assert_eq!(ob.asks.len(), 1);
}
other => panic!("expected OrderBook, got {other:?}"),
}
}
#[test]
fn parse_mark_price_into_funding_rate() {
let raw = r#"{
"stream": "btcusdt@markPrice@1s",
"data": {
"e": "markPriceUpdate",
"E": 1700000000000,
"s": "BTCUSDT",
"p": "96010.0",
"i": "96005.0",
"P": "96012.0",
"r": "0.0001",
"T": 1700028800000
}
}"#;
let msgs = connector().parse_message(raw).expect("parse");
match &msgs[0] {
DataMessage::FundingRate(f) => {
assert_eq!(f.symbol, "BTCUSDT");
assert!((f.funding_rate - 0.0001).abs() < 1e-9);
assert_eq!(f.next_funding_time, 1_700_028_800_000);
assert_eq!(f.mark_price, Some(96_010.0));
assert_eq!(f.index_price, Some(96_005.0));
}
other => panic!("expected FundingRate, got {other:?}"),
}
}
#[test]
fn unknown_event_returns_empty_vec() {
let raw = r#"{"e": "someFutureEvent", "s": "BTCUSDT"}"#;
let msgs = connector().parse_message(raw).expect("parse");
assert!(
msgs.is_empty(),
"unknown event should yield no DataMessages"
);
}
#[test]
fn stream_name_builders_lowercase_symbols() {
assert_eq!(
BinanceConnector::trade_stream("BTCUSDT"),
"btcusdt@aggTrade"
);
assert_eq!(
BinanceConnector::ticker_stream("ETHUSDT"),
"ethusdt@bookTicker"
);
assert_eq!(
BinanceConnector::kline_stream("BTCUSDT", "5m"),
"btcusdt@kline_5m"
);
assert_eq!(
BinanceConnector::depth_stream("BTCUSDT"),
"btcusdt@depth@100ms"
);
assert_eq!(
BinanceConnector::mark_price_stream("BTCUSDT"),
"btcusdt@markPrice@1s"
);
}
#[test]
fn depth_snapshot_stream_clamps_levels() {
assert!(BinanceConnector::depth_snapshot_stream("BTCUSDT", 0).ends_with("@depth5@100ms"));
assert!(BinanceConnector::depth_snapshot_stream("BTCUSDT", 5).ends_with("@depth5@100ms"));
assert!(BinanceConnector::depth_snapshot_stream("BTCUSDT", 8).ends_with("@depth10@100ms"));
assert!(BinanceConnector::depth_snapshot_stream("BTCUSDT", 20).ends_with("@depth20@100ms"));
assert!(
BinanceConnector::depth_snapshot_stream("BTCUSDT", 100).ends_with("@depth20@100ms")
);
}
#[test]
fn spot_url_combines_streams() {
let c = BinanceConnector::spot(&["btcusdt@aggTrade", "btcusdt@bookTicker"]);
assert!(
c.url
.starts_with("wss://stream.binance.com:9443/stream?streams=")
);
assert!(c.url.contains("btcusdt@aggTrade"));
assert!(c.url.contains("btcusdt@bookTicker"));
assert_eq!(c.exchange, "binance");
}
#[test]
fn futures_url_uses_fstream_host() {
let c = BinanceConnector::futures(&["btcusdt@markPrice@1s"]);
assert!(
c.url
.starts_with("wss://fstream.binance.com/stream?streams=")
);
assert_eq!(c.exchange, "binance-futures");
}
}