use serde_json::Value;
use crate::actors::{
CandleData, DataMessage, ExchangeConnector, FundingData, OrderBookData, TickerData, TradeData,
TradeSide, WebSocketConfig,
};
use crate::bybit::BybitCategory;
use crate::error::Result;
const WS_BASE: &str = "wss://stream.bybit.com/v5/public";
const EXCHANGE_NAME: &str = "bybit";
const PING_INTERVAL_SECS: u64 = 20;
#[derive(Debug, Clone)]
pub struct BybitConnector {
pub url: String,
pub category: BybitCategory,
pub topics: Vec<String>,
}
impl BybitConnector {
#[must_use]
pub fn new(category: BybitCategory, topics: Vec<String>) -> Self {
Self {
url: format!("{WS_BASE}/{}", category.as_str()),
category,
topics,
}
}
#[must_use]
pub fn with_url(url: impl Into<String>, category: BybitCategory, topics: Vec<String>) -> Self {
Self {
url: url.into(),
category,
topics,
}
}
#[must_use]
pub fn trade_topic(symbol: &str) -> String {
format!("publicTrade.{symbol}")
}
#[must_use]
pub fn ticker_topic(symbol: &str) -> String {
format!("tickers.{symbol}")
}
#[must_use]
pub fn kline_topic(symbol: &str, interval: &str) -> String {
format!("kline.{interval}.{symbol}")
}
#[must_use]
pub fn orderbook_topic(symbol: &str, depth: u32) -> String {
format!("orderbook.{depth}.{symbol}")
}
}
impl ExchangeConnector for BybitConnector {
fn exchange_name(&self) -> &str {
EXCHANGE_NAME
}
fn ws_url(&self) -> &str {
&self.url
}
fn build_ws_config(&self, symbol: &str) -> WebSocketConfig {
WebSocketConfig {
url: self.url.clone(),
exchange: EXCHANGE_NAME.to_string(),
symbol: symbol.to_string(),
subscription_msg: self.subscription_message(symbol),
ping_interval_secs: PING_INTERVAL_SECS,
reconnect_delay_secs: 5,
max_reconnect_attempts: 5,
}
}
fn subscription_message(&self, _symbol: &str) -> Option<String> {
if self.topics.is_empty() {
return None;
}
serde_json::to_string(&serde_json::json!({
"op": "subscribe",
"args": &self.topics,
}))
.ok()
}
fn ping_message(&self) -> Option<String> {
Some(r#"{"op":"ping"}"#.to_string())
}
fn parse_message(&self, raw: &str) -> Result<Vec<DataMessage>> {
let json: Value = serde_json::from_str(raw)?;
if json.get("op").is_some() {
return Ok(vec![]);
}
let topic = json.get("topic").and_then(Value::as_str).unwrap_or("");
let Some(data) = json.get("data") else {
return Ok(vec![]);
};
let is_snapshot = json
.get("type")
.and_then(Value::as_str)
.unwrap_or("snapshot")
== "snapshot";
if topic.starts_with("publicTrade.") {
Ok(parse_trade_batch(data))
} else if topic.starts_with("tickers.") {
Ok(parse_ticker(data))
} else if topic.starts_with("kline.") {
Ok(parse_kline_batch(data))
} else if topic.starts_with("orderbook.") {
Ok(parse_orderbook(data, is_snapshot))
} else {
Ok(vec![])
}
}
}
fn now_ms() -> i64 {
chrono::Utc::now().timestamp_millis()
}
fn str_f64(v: &Value, key: &str) -> f64 {
v.get(key)
.and_then(Value::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(Value::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_trade_batch(data: &Value) -> Vec<DataMessage> {
let Some(arr) = data.as_array() else {
return vec![];
};
arr.iter()
.map(|t| {
let side = match t.get("S").and_then(Value::as_str).unwrap_or("Buy") {
s if s.eq_ignore_ascii_case("Sell") => TradeSide::Sell,
_ => TradeSide::Buy,
};
DataMessage::Trade(TradeData {
symbol: t["s"].as_str().unwrap_or("").to_string(),
exchange: EXCHANGE_NAME.to_string(),
side,
price: str_f64(t, "p"),
amount: str_f64(t, "v"),
exchange_ts: t["T"].as_i64().unwrap_or(0),
receipt_ts: now_ms(),
trade_id: t.get("i").and_then(Value::as_str).unwrap_or("").to_string(),
})
})
.collect()
}
fn parse_ticker(data: &Value) -> Vec<DataMessage> {
let symbol = data["symbol"].as_str().unwrap_or("").to_string();
let last_price = str_f64(data, "lastPrice");
let bid = opt_str_f64(data, "bid1Price").unwrap_or(0.0);
let ask = opt_str_f64(data, "ask1Price").unwrap_or(0.0);
let now = now_ms();
vec![DataMessage::Ticker(TickerData {
symbol,
exchange: EXCHANGE_NAME.to_string(),
price: last_price,
best_bid: bid,
best_ask: ask,
exchange_ts: now,
receipt_ts: now,
})]
}
fn parse_kline_batch(data: &Value) -> Vec<DataMessage> {
let Some(arr) = data.as_array() else {
return vec![];
};
arr.iter()
.map(|k| {
DataMessage::Candle(CandleData {
symbol: String::new(),
exchange: EXCHANGE_NAME.to_string(),
interval: k["interval"].as_str().unwrap_or("").to_string(),
open_ts: k["start"].as_i64().unwrap_or(0),
open: str_f64(k, "open"),
high: str_f64(k, "high"),
low: str_f64(k, "low"),
close: str_f64(k, "close"),
volume: str_f64(k, "volume"),
is_closed: k["confirm"].as_bool().unwrap_or(false),
receipt_ts: now_ms(),
})
})
.collect()
}
fn parse_orderbook(data: &Value, is_snapshot: bool) -> Vec<DataMessage> {
let symbol = data["s"].as_str().unwrap_or("").to_string();
let now = now_ms();
vec![DataMessage::OrderBook(OrderBookData {
symbol,
exchange: EXCHANGE_NAME.to_string(),
asks: parse_levels(&data["a"]),
bids: parse_levels(&data["b"]),
exchange_ts: now,
receipt_ts: now,
is_snapshot,
})]
}
#[allow(dead_code)]
fn parse_ticker_funding(data: &Value) -> Option<FundingData> {
let funding_rate = opt_str_f64(data, "fundingRate")?;
let next_funding = data
.get("nextFundingTime")
.and_then(Value::as_str)
.and_then(|s| s.parse::<i64>().ok())
.unwrap_or(0);
Some(FundingData {
symbol: data["symbol"].as_str().unwrap_or("").to_string(),
exchange: EXCHANGE_NAME.to_string(),
funding_rate,
next_funding_time: next_funding,
mark_price: opt_str_f64(data, "markPrice"),
index_price: opt_str_f64(data, "indexPrice"),
exchange_ts: now_ms(),
receipt_ts: now_ms(),
})
}
#[cfg(test)]
mod tests {
use super::*;
fn connector() -> BybitConnector {
BybitConnector::new(
BybitCategory::Linear,
vec!["publicTrade.BTCUSDT".into(), "kline.1.BTCUSDT".into()],
)
}
#[test]
fn subscription_message_packs_topics() {
let c = connector();
let sub = c.subscription_message("ignored").expect("sub");
let parsed: Value = serde_json::from_str(&sub).unwrap();
assert_eq!(parsed["op"], "subscribe");
let args = parsed["args"].as_array().unwrap();
assert_eq!(args.len(), 2);
assert_eq!(args[0], "publicTrade.BTCUSDT");
}
#[test]
fn empty_topics_returns_no_subscription() {
let c = BybitConnector::new(BybitCategory::Spot, vec![]);
assert!(c.subscription_message("BTCUSDT").is_none());
}
#[test]
fn ping_uses_bybit_op_format() {
assert_eq!(
connector().ping_message().as_deref(),
Some(r#"{"op":"ping"}"#)
);
}
#[test]
fn parse_op_ack_returns_empty() {
let raw = r#"{"op":"subscribe","conn_id":"x","ret_msg":"","success":true,"req_id":"y"}"#;
let msgs = connector().parse_message(raw).expect("parse");
assert!(msgs.is_empty());
}
#[test]
fn parse_pong_returns_empty() {
let raw = r#"{"op":"pong","ret_msg":"pong","success":true}"#;
let msgs = connector().parse_message(raw).expect("parse");
assert!(msgs.is_empty());
}
#[test]
fn parse_public_trade_emits_one_per_array_element() {
let raw = r#"{
"topic": "publicTrade.BTCUSDT",
"type": "snapshot",
"ts": 1700000000000,
"data": [
{"T":1700000000050,"s":"BTCUSDT","S":"Buy","v":"0.1","p":"96000.0","L":"PlusTick","i":"id-1","BT":false},
{"T":1700000000080,"s":"BTCUSDT","S":"Sell","v":"0.05","p":"96005.0","L":"MinusTick","i":"id-2","BT":false}
]
}"#;
let msgs = connector().parse_message(raw).expect("parse");
assert_eq!(msgs.len(), 2);
match &msgs[0] {
DataMessage::Trade(t) => {
assert_eq!(t.symbol, "BTCUSDT");
assert_eq!(t.exchange, "bybit");
assert_eq!(t.side, TradeSide::Buy);
assert!((t.price - 96_000.0).abs() < 1e-9);
assert_eq!(t.trade_id, "id-1");
}
other => panic!("expected Trade, got {other:?}"),
}
match &msgs[1] {
DataMessage::Trade(t) => assert_eq!(t.side, TradeSide::Sell),
_ => panic!("expected Trade variant"),
}
}
#[test]
fn parse_ticker_into_ticker_data() {
let raw = r#"{
"topic":"tickers.BTCUSDT",
"type":"snapshot",
"ts":1700000000000,
"data":{
"symbol":"BTCUSDT",
"lastPrice":"96000.0",
"bid1Price":"95999.0","bid1Size":"1.0",
"ask1Price":"96001.0","ask1Size":"1.5",
"markPrice":"96010.0","indexPrice":"96005.0",
"fundingRate":"0.0001","nextFundingTime":"1700028800000"
}
}"#;
let msgs = connector().parse_message(raw).expect("parse");
match &msgs[0] {
DataMessage::Ticker(t) => {
assert_eq!(t.symbol, "BTCUSDT");
assert!((t.price - 96_000.0).abs() < 1e-9);
assert!((t.best_bid - 95_999.0).abs() < 1e-9);
assert!((t.best_ask - 96_001.0).abs() < 1e-9);
}
other => panic!("expected Ticker, got {other:?}"),
}
}
#[test]
fn parse_kline_into_candle() {
let raw = r#"{
"topic":"kline.1.BTCUSDT",
"type":"snapshot",
"ts":1700000000050,
"data":[{
"start":1700000000000,"end":1700000059999,"interval":"1",
"open":"96000.0","close":"96100.0","high":"96200.0","low":"95900.0",
"volume":"10.0","turnover":"961000.0","confirm":true,"timestamp":1700000000050
}]
}"#;
let msgs = connector().parse_message(raw).expect("parse");
match &msgs[0] {
DataMessage::Candle(c) => {
assert_eq!(c.interval, "1");
assert_eq!(c.open_ts, 1_700_000_000_000);
assert!((c.open - 96_000.0).abs() < 1e-9);
assert!(c.is_closed);
}
other => panic!("expected Candle, got {other:?}"),
}
}
#[test]
fn parse_orderbook_snapshot_then_delta() {
let snap = r#"{
"topic":"orderbook.50.BTCUSDT",
"type":"snapshot",
"ts":1700000000000,
"data":{"s":"BTCUSDT","b":[["96000.0","1.5"]],"a":[["96001.0","0.5"]],"u":1,"seq":1}
}"#;
let delta = r#"{
"topic":"orderbook.50.BTCUSDT",
"type":"delta",
"ts":1700000000100,
"data":{"s":"BTCUSDT","b":[["95999.5","0.0"]],"a":[],"u":2,"seq":2}
}"#;
let c = connector();
match &c.parse_message(snap).expect("snap")[0] {
DataMessage::OrderBook(ob) => {
assert_eq!(ob.symbol, "BTCUSDT");
assert!(ob.is_snapshot);
assert_eq!(ob.asks.len(), 1);
}
_ => panic!("snapshot was not OrderBook"),
}
match &c.parse_message(delta).expect("delta")[0] {
DataMessage::OrderBook(ob) => {
assert!(!ob.is_snapshot);
assert!((ob.bids[0][1] - 0.0).abs() < 1e-12); }
_ => panic!("delta was not OrderBook"),
}
}
#[test]
fn topic_builders_format() {
assert_eq!(
BybitConnector::trade_topic("BTCUSDT"),
"publicTrade.BTCUSDT"
);
assert_eq!(BybitConnector::ticker_topic("BTCUSDT"), "tickers.BTCUSDT");
assert_eq!(
BybitConnector::kline_topic("BTCUSDT", "1"),
"kline.1.BTCUSDT"
);
assert_eq!(
BybitConnector::orderbook_topic("BTCUSDT", 50),
"orderbook.50.BTCUSDT"
);
}
#[test]
fn ws_url_per_category() {
assert!(
BybitConnector::new(BybitCategory::Spot, vec![])
.url
.ends_with("/v5/public/spot")
);
assert!(
BybitConnector::new(BybitCategory::Linear, vec![])
.url
.ends_with("/v5/public/linear")
);
assert!(
BybitConnector::new(BybitCategory::Inverse, vec![])
.url
.ends_with("/v5/public/inverse")
);
}
#[test]
fn ticker_funding_extractor_round_trip() {
let data = serde_json::json!({
"symbol":"BTCUSDT",
"markPrice":"96010.0","indexPrice":"96005.0",
"fundingRate":"0.0001","nextFundingTime":"1700028800000"
});
let f = parse_ticker_funding(&data).expect("funding");
assert_eq!(f.exchange, "bybit");
assert_eq!(f.next_funding_time, 1_700_028_800_000);
assert_eq!(f.mark_price, Some(96_010.0));
}
}