use serde_json::{Value, json};
use crate::actors::{
CandleData, DataMessage, ExchangeConnector, OrderBookData, TickerData, TradeData, TradeSide,
WebSocketConfig,
};
use crate::error::Result;
const WS_PUBLIC_URL: &str = "wss://stream.crypto.com/exchange/v1/market";
const WS_PRIVATE_URL: &str = "wss://stream.crypto.com/exchange/v1/user";
const EXCHANGE_NAME: &str = "cryptocom";
#[derive(Debug, Clone)]
pub struct CryptocomConnector {
pub url: String,
}
impl CryptocomConnector {
#[must_use]
pub fn public() -> Self {
Self {
url: WS_PUBLIC_URL.to_string(),
}
}
#[must_use]
pub fn private() -> Self {
Self {
url: WS_PRIVATE_URL.to_string(),
}
}
#[must_use]
pub fn with_url(url: impl Into<String>) -> Self {
Self { url: url.into() }
}
#[must_use]
pub fn trade_channel(instrument: &str) -> String {
format!("trade.{instrument}")
}
#[must_use]
pub fn ticker_channel(instrument: &str) -> String {
format!("ticker.{instrument}")
}
#[must_use]
pub fn candlestick_channel(instrument: &str, timeframe: &str) -> String {
format!("candlestick.{timeframe}.{instrument}")
}
#[must_use]
pub fn book_channel(instrument: &str, depth: u32) -> String {
format!("book.{instrument}.{depth}")
}
#[must_use]
pub fn subscribe_frame(id: i64, channels: &[String]) -> String {
let channels_ref: Vec<&str> = channels.iter().map(String::as_str).collect();
json!({
"id": id,
"method": "subscribe",
"params": {"channels": channels_ref},
})
.to_string()
}
}
impl ExchangeConnector for CryptocomConnector {
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: None,
ping_interval_secs: 30,
reconnect_delay_secs: 5,
max_reconnect_attempts: 5,
}
}
fn subscription_message(&self, _symbol: &str) -> Option<String> {
None
}
fn ping_message(&self) -> Option<String> {
None
}
fn response_for(&self, raw: &str) -> Option<String> {
if !raw.contains("public/heartbeat") {
return None;
}
let json: Value = serde_json::from_str(raw).ok()?;
if json.get("method").and_then(Value::as_str)? != "public/heartbeat" {
return None;
}
let id = json.get("id").and_then(Value::as_i64)?;
Some(json!({"id": id, "method": "public/respond-heartbeat"}).to_string())
}
fn parse_message(&self, raw: &str) -> Result<Vec<DataMessage>> {
let json: Value = serde_json::from_str(raw)?;
if json.get("method").and_then(Value::as_str) == Some("public/heartbeat") {
return Ok(vec![]);
}
let Some(result) = json.get("result") else {
return Ok(vec![]);
};
let channel = result.get("channel").and_then(Value::as_str).unwrap_or("");
let Some(data) = result.get("data").and_then(Value::as_array) else {
return Ok(vec![]);
};
let instrument = result
.get("instrument_name")
.and_then(Value::as_str)
.unwrap_or("")
.to_string();
let is_snapshot = result
.get("type")
.and_then(Value::as_str)
.unwrap_or("update")
== "snapshot";
match channel {
"trade" => Ok(parse_trades(data, &instrument)),
"ticker" => Ok(parse_tickers(data, &instrument)),
"candlestick" => Ok(parse_candles(data, &instrument)),
"book" => Ok(parse_books(data, &instrument, is_snapshot)),
_ => Ok(vec![]),
}
}
}
fn now_ms() -> i64 {
chrono::Utc::now().timestamp_millis()
}
fn flexible_f64(v: &Value, key: &str) -> f64 {
match v.get(key) {
Some(Value::String(s)) => s.parse().unwrap_or(0.0),
Some(Value::Number(n)) => n.as_f64().unwrap_or(0.0),
_ => 0.0,
}
}
fn parse_trades(data: &[Value], instrument_fallback: &str) -> Vec<DataMessage> {
data.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,
};
let symbol = t
.get("i")
.and_then(Value::as_str)
.unwrap_or(instrument_fallback)
.to_string();
DataMessage::Trade(TradeData {
symbol,
exchange: EXCHANGE_NAME.to_string(),
side,
price: flexible_f64(t, "p"),
amount: flexible_f64(t, "q"),
exchange_ts: t.get("t").and_then(Value::as_i64).unwrap_or(0),
receipt_ts: now_ms(),
trade_id: t.get("d").and_then(Value::as_str).unwrap_or("").to_string(),
})
})
.collect()
}
fn parse_tickers(data: &[Value], instrument_fallback: &str) -> Vec<DataMessage> {
let now = now_ms();
data.iter()
.map(|t| {
let symbol = t
.get("i")
.and_then(Value::as_str)
.unwrap_or(instrument_fallback)
.to_string();
DataMessage::Ticker(TickerData {
symbol,
exchange: EXCHANGE_NAME.to_string(),
price: flexible_f64(t, "a"),
best_bid: flexible_f64(t, "b"),
best_ask: flexible_f64(t, "k"),
exchange_ts: t.get("t").and_then(Value::as_i64).unwrap_or(now),
receipt_ts: now,
})
})
.collect()
}
fn parse_candles(data: &[Value], instrument_fallback: &str) -> Vec<DataMessage> {
let now = now_ms();
data.iter()
.map(|c| {
let symbol = c
.get("i")
.and_then(Value::as_str)
.unwrap_or(instrument_fallback)
.to_string();
let interval = c
.get("interval")
.and_then(Value::as_str)
.unwrap_or("")
.to_string();
DataMessage::Candle(CandleData {
symbol,
exchange: EXCHANGE_NAME.to_string(),
interval,
open_ts: c.get("t").and_then(Value::as_i64).unwrap_or(now),
open: flexible_f64(c, "o"),
high: flexible_f64(c, "h"),
low: flexible_f64(c, "l"),
close: flexible_f64(c, "c"),
volume: flexible_f64(c, "v"),
is_closed: false,
receipt_ts: now,
})
})
.collect()
}
fn parse_books(data: &[Value], instrument_fallback: &str, is_snapshot: bool) -> Vec<DataMessage> {
let now = now_ms();
data.iter()
.map(|b| {
DataMessage::OrderBook(OrderBookData {
symbol: instrument_fallback.to_string(),
exchange: EXCHANGE_NAME.to_string(),
asks: parse_levels(b.get("asks")),
bids: parse_levels(b.get("bids")),
exchange_ts: b.get("t").and_then(Value::as_i64).unwrap_or(now),
receipt_ts: now,
is_snapshot,
})
})
.collect()
}
fn parse_levels(v: Option<&Value>) -> Vec<[f64; 2]> {
v.and_then(Value::as_array)
.map(|arr| {
arr.iter()
.filter_map(|lvl| {
let p: f64 = lvl.get(0)?.as_str()?.parse().ok()?;
let q: f64 = lvl.get(1)?.as_str()?.parse().ok()?;
Some([p, q])
})
.collect()
})
.unwrap_or_default()
}
#[cfg(test)]
mod tests {
use super::*;
fn connector() -> CryptocomConnector {
CryptocomConnector::public()
}
#[test]
fn ws_url_picks_public_or_private() {
assert!(CryptocomConnector::public().url.ends_with("/market"));
assert!(CryptocomConnector::private().url.ends_with("/user"));
}
#[test]
fn channel_builders_format() {
assert_eq!(
CryptocomConnector::trade_channel("BTC_USDT"),
"trade.BTC_USDT"
);
assert_eq!(
CryptocomConnector::ticker_channel("BTC_USDT"),
"ticker.BTC_USDT"
);
assert_eq!(
CryptocomConnector::candlestick_channel("BTC_USDT", "1m"),
"candlestick.1m.BTC_USDT"
);
assert_eq!(
CryptocomConnector::book_channel("BTC_USDT", 10),
"book.BTC_USDT.10"
);
}
#[test]
fn subscribe_frame_carries_all_channels() {
let frame = CryptocomConnector::subscribe_frame(
7,
&[
CryptocomConnector::trade_channel("BTC_USDT"),
CryptocomConnector::book_channel("BTC_USDT", 10),
],
);
let v: Value = serde_json::from_str(&frame).unwrap();
assert_eq!(v["id"], 7);
assert_eq!(v["method"], "subscribe");
assert_eq!(v["params"]["channels"][0], "trade.BTC_USDT");
assert_eq!(v["params"]["channels"][1], "book.BTC_USDT.10");
}
#[test]
fn ping_message_is_none() {
assert!(connector().ping_message().is_none());
}
#[test]
fn response_for_heartbeat_echoes_id() {
let raw = r#"{"id":1234,"method":"public/heartbeat"}"#;
let resp = connector().response_for(raw).expect("response");
let v: Value = serde_json::from_str(&resp).unwrap();
assert_eq!(v["id"], 1234);
assert_eq!(v["method"], "public/respond-heartbeat");
}
#[test]
fn response_for_non_heartbeat_is_none() {
assert!(
connector()
.response_for(r#"{"id":1,"method":"subscribe","code":0}"#)
.is_none()
);
assert!(
connector()
.response_for(r#"{"result":{"channel":"trade","data":[]}}"#)
.is_none()
);
}
#[test]
fn response_for_short_circuits_on_non_heartbeat_text() {
let raw = r#"{"result":{"channel":"book","data":[]}}"#;
assert!(connector().response_for(raw).is_none());
}
#[test]
fn parse_subscribe_ack_returns_empty() {
let raw = r#"{"id":1,"method":"subscribe","code":0}"#;
assert!(connector().parse_message(raw).expect("parse").is_empty());
}
#[test]
fn parse_heartbeat_returns_empty() {
let raw = r#"{"id":1,"method":"public/heartbeat"}"#;
assert!(connector().parse_message(raw).expect("parse").is_empty());
}
#[test]
fn parse_trade_uses_inner_symbol_and_emits_one_per_element() {
let raw = r#"{
"id":-1,"method":"subscribe","code":0,
"result":{
"instrument_name":"BTC_USDT","channel":"trade","subscription":"trade.BTC_USDT",
"data":[
{"i":"BTC_USDT","s":"buy","p":"96000","q":"0.05","t":1700000000000,"d":"id-1"},
{"i":"BTC_USDT","s":"sell","p":"96005","q":"0.10","t":1700000000500,"d":"id-2"}
]
}
}"#;
let msgs = connector().parse_message(raw).expect("parse");
assert_eq!(msgs.len(), 2);
match &msgs[0] {
DataMessage::Trade(t) => {
assert_eq!(t.symbol, "BTC_USDT");
assert_eq!(t.exchange, "cryptocom");
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"),
}
}
#[test]
fn parse_ticker_routes_letter_fields() {
let raw = r#"{
"result":{
"instrument_name":"BTC_USDT","channel":"ticker","subscription":"ticker.BTC_USDT",
"data":[{
"i":"BTC_USDT","a":"96000.0","b":"95999.0","k":"96001.0",
"h":"96500.0","l":"95500.0","v":"100.5","vv":"9650000",
"c":"0.005","t":1700000000000
}]
}
}"#;
match &connector().parse_message(raw).expect("parse")[0] {
DataMessage::Ticker(t) => {
assert_eq!(t.symbol, "BTC_USDT");
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_candlestick_emits_one_per_element() {
let raw = r#"{
"result":{
"instrument_name":"BTC_USDT","channel":"candlestick","subscription":"candlestick.1m.BTC_USDT",
"data":[
{"o":"96000","h":"96100","l":"95900","c":"96050","v":"10.5","t":1700000000000}
]
}
}"#;
match &connector().parse_message(raw).expect("parse")[0] {
DataMessage::Candle(c) => {
assert_eq!(c.symbol, "BTC_USDT");
assert!((c.open - 96_000.0).abs() < 1e-9);
assert!((c.close - 96_050.0).abs() < 1e-9);
assert!(!c.is_closed);
}
other => panic!("expected Candle, got {other:?}"),
}
}
#[test]
fn parse_book_snapshot_then_delta() {
let snap = r#"{
"result":{
"instrument_name":"BTC_USDT","channel":"book","subscription":"book.BTC_USDT.10","type":"snapshot",
"data":[{
"asks":[["96001","0.5","1"]],
"bids":[["96000","1.5","2"],["95999","2.0","3"]],
"t":1700000000000,"s":1
}]
}
}"#;
let delta = r#"{
"result":{
"instrument_name":"BTC_USDT","channel":"book","subscription":"book.BTC_USDT.10","type":"update",
"data":[{
"asks":[],
"bids":[["95998","0","0"]],
"t":1700000000100,"s":2
}]
}
}"#;
let c = connector();
match &c.parse_message(snap).expect("snap")[0] {
DataMessage::OrderBook(ob) => {
assert_eq!(ob.symbol, "BTC_USDT");
assert!(ob.is_snapshot);
assert_eq!(ob.bids.len(), 2);
}
_ => 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 parse_unknown_channel_returns_empty() {
let raw = r#"{"result":{"channel":"user.order","data":[{"x":"y"}],"instrument_name":"BTC_USDT"}}"#;
assert!(connector().parse_message(raw).expect("parse").is_empty());
}
}