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://ws.kraken.com/v2";
const WS_PRIVATE_URL: &str = "wss://ws-auth.kraken.com/v2";
const EXCHANGE_NAME: &str = "kraken";
const PING_INTERVAL_SECS: u64 = 30;
#[derive(Debug, Clone)]
pub struct KrakenConnector {
pub url: String,
}
impl KrakenConnector {
#[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_subscription(pairs: &[&str]) -> String {
json!({
"method": "subscribe",
"params": {"channel": "trade", "symbol": pairs},
})
.to_string()
}
#[must_use]
pub fn ticker_subscription(pairs: &[&str]) -> String {
json!({
"method": "subscribe",
"params": {"channel": "ticker", "symbol": pairs},
})
.to_string()
}
#[must_use]
pub fn ohlc_subscription(pairs: &[&str], interval_mins: u32) -> String {
json!({
"method": "subscribe",
"params": {"channel": "ohlc", "symbol": pairs, "interval": interval_mins},
})
.to_string()
}
#[must_use]
pub fn book_subscription(pairs: &[&str], depth: u32) -> String {
json!({
"method": "subscribe",
"params": {"channel": "book", "symbol": pairs, "depth": depth},
})
.to_string()
}
}
impl ExchangeConnector for KrakenConnector {
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: PING_INTERVAL_SECS,
reconnect_delay_secs: 5,
max_reconnect_attempts: 5,
}
}
fn subscription_message(&self, _symbol: &str) -> Option<String> {
None
}
fn ping_message(&self) -> Option<String> {
Some(r#"{"method":"ping"}"#.to_string())
}
fn parse_message(&self, raw: &str) -> Result<Vec<DataMessage>> {
let json: Value = serde_json::from_str(raw)?;
if json.get("method").is_some() {
return Ok(vec![]);
}
let channel = json.get("channel").and_then(Value::as_str).unwrap_or("");
let Some(data) = json.get("data").and_then(Value::as_array) else {
return Ok(vec![]);
};
let is_snapshot =
json.get("type").and_then(Value::as_str).unwrap_or("update") == "snapshot";
match channel {
"trade" => Ok(parse_trades(data)),
"ticker" => Ok(parse_tickers(data)),
"ohlc" => Ok(parse_klines(data)),
"book" => Ok(parse_books(data, is_snapshot)),
_ => Ok(vec![]),
}
}
}
fn now_ms() -> i64 {
chrono::Utc::now().timestamp_millis()
}
fn iso_to_ms(s: &str) -> i64 {
chrono::DateTime::parse_from_rfc3339(s).map_or_else(|_| now_ms(), |dt| dt.timestamp_millis())
}
fn f64_field(v: &Value, key: &str) -> f64 {
v.get(key).and_then(Value::as_f64).unwrap_or(0.0)
}
fn parse_trades(data: &[Value]) -> Vec<DataMessage> {
data.iter()
.map(|t| {
let side = match t.get("side").and_then(Value::as_str).unwrap_or("buy") {
s if s.eq_ignore_ascii_case("sell") => TradeSide::Sell,
_ => TradeSide::Buy,
};
let ts = t
.get("timestamp")
.and_then(Value::as_str)
.map_or_else(now_ms, iso_to_ms);
let trade_id = t
.get("trade_id")
.and_then(Value::as_u64)
.map(|n| n.to_string())
.unwrap_or_default();
DataMessage::Trade(TradeData {
symbol: t["symbol"].as_str().unwrap_or("").to_string(),
exchange: EXCHANGE_NAME.to_string(),
side,
price: f64_field(t, "price"),
amount: f64_field(t, "qty"),
exchange_ts: ts,
receipt_ts: now_ms(),
trade_id,
})
})
.collect()
}
fn parse_tickers(data: &[Value]) -> Vec<DataMessage> {
let now = now_ms();
data.iter()
.map(|t| {
DataMessage::Ticker(TickerData {
symbol: t["symbol"].as_str().unwrap_or("").to_string(),
exchange: EXCHANGE_NAME.to_string(),
price: f64_field(t, "last"),
best_bid: f64_field(t, "bid"),
best_ask: f64_field(t, "ask"),
exchange_ts: now, receipt_ts: now,
})
})
.collect()
}
fn parse_klines(data: &[Value]) -> Vec<DataMessage> {
data.iter()
.map(|c| {
let interval = c
.get("interval")
.and_then(Value::as_u64)
.map(|n| n.to_string())
.unwrap_or_default();
DataMessage::Candle(CandleData {
symbol: c["symbol"].as_str().unwrap_or("").to_string(),
exchange: EXCHANGE_NAME.to_string(),
interval,
open_ts: c
.get("interval_begin")
.and_then(Value::as_str)
.map_or_else(now_ms, iso_to_ms),
open: f64_field(c, "open"),
high: f64_field(c, "high"),
low: f64_field(c, "low"),
close: f64_field(c, "close"),
volume: f64_field(c, "volume"),
is_closed: false,
receipt_ts: now_ms(),
})
})
.collect()
}
fn parse_books(data: &[Value], is_snapshot: bool) -> Vec<DataMessage> {
let now = now_ms();
data.iter()
.map(|b| {
let bids = parse_level_objects(b.get("bids"));
let asks = parse_level_objects(b.get("asks"));
DataMessage::OrderBook(OrderBookData {
symbol: b["symbol"].as_str().unwrap_or("").to_string(),
exchange: EXCHANGE_NAME.to_string(),
asks,
bids,
exchange_ts: now, receipt_ts: now,
is_snapshot,
})
})
.collect()
}
fn parse_level_objects(v: Option<&Value>) -> Vec<[f64; 2]> {
v.and_then(Value::as_array)
.map(|arr| {
arr.iter()
.filter_map(|lvl| {
let p = lvl.get("price")?.as_f64()?;
let q = lvl.get("qty")?.as_f64()?;
Some([p, q])
})
.collect()
})
.unwrap_or_default()
}
#[cfg(test)]
mod tests {
use super::*;
fn connector() -> KrakenConnector {
KrakenConnector::public()
}
#[test]
fn ws_url_picks_public_or_private() {
assert!(KrakenConnector::public().url.ends_with("ws.kraken.com/v2"));
assert!(
KrakenConnector::private()
.url
.ends_with("ws-auth.kraken.com/v2")
);
}
#[test]
fn ping_uses_kraken_method_format() {
assert_eq!(
connector().ping_message().as_deref(),
Some(r#"{"method":"ping"}"#)
);
}
#[test]
fn subscription_builders_emit_canonical_shape() {
let sub = KrakenConnector::trade_subscription(&["BTC/USD", "ETH/USD"]);
let v: Value = serde_json::from_str(&sub).unwrap();
assert_eq!(v["method"], "subscribe");
assert_eq!(v["params"]["channel"], "trade");
assert_eq!(v["params"]["symbol"][0], "BTC/USD");
assert_eq!(v["params"]["symbol"][1], "ETH/USD");
let ohlc = KrakenConnector::ohlc_subscription(&["BTC/USD"], 5);
let v: Value = serde_json::from_str(&ohlc).unwrap();
assert_eq!(v["params"]["channel"], "ohlc");
assert_eq!(v["params"]["interval"], 5);
let book = KrakenConnector::book_subscription(&["BTC/USD"], 100);
let v: Value = serde_json::from_str(&book).unwrap();
assert_eq!(v["params"]["channel"], "book");
assert_eq!(v["params"]["depth"], 100);
}
#[test]
fn parse_method_response_returns_empty() {
let raw = r#"{"method":"subscribe","success":true,"result":{"channel":"trade","symbol":"BTC/USD"}}"#;
assert!(connector().parse_message(raw).expect("parse").is_empty());
}
#[test]
fn parse_pong_returns_empty() {
let raw = r#"{"method":"pong"}"#;
assert!(connector().parse_message(raw).expect("parse").is_empty());
}
#[test]
fn parse_heartbeat_returns_empty() {
let raw = r#"{"channel":"heartbeat"}"#;
assert!(connector().parse_message(raw).expect("parse").is_empty());
}
#[test]
fn parse_trade_emits_one_per_array_element() {
let raw = r#"{
"channel":"trade","type":"snapshot",
"data":[
{"symbol":"BTC/USD","side":"buy","qty":0.1,"price":96000.0,"ord_type":"market","trade_id":1,"timestamp":"2026-05-25T12:00:00.000000Z"},
{"symbol":"BTC/USD","side":"sell","qty":0.05,"price":96005.0,"ord_type":"limit","trade_id":2,"timestamp":"2026-05-25T12:00:00.500000Z"}
]
}"#;
let msgs = connector().parse_message(raw).expect("parse");
assert_eq!(msgs.len(), 2);
match &msgs[0] {
DataMessage::Trade(t) => {
assert_eq!(t.symbol, "BTC/USD");
assert_eq!(t.exchange, "kraken");
assert_eq!(t.side, TradeSide::Buy);
assert!((t.price - 96_000.0).abs() < 1e-9);
assert!((t.amount - 0.1).abs() < 1e-12);
assert_eq!(t.trade_id, "1");
assert!(t.exchange_ts > 1_700_000_000_000); }
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_into_ticker_data() {
let raw = r#"{
"channel":"ticker","type":"snapshot",
"data":[{
"symbol":"BTC/USD","bid":95999.0,"ask":96001.0,
"bid_qty":1.0,"ask_qty":1.5,"last":96000.0,
"volume":100.5,"high":96500.0,"low":95500.0,
"vwap":95800.0,"change":250.0,"change_pct":0.26
}]
}"#;
let msgs = connector().parse_message(raw).expect("parse");
match &msgs[0] {
DataMessage::Ticker(t) => {
assert_eq!(t.symbol, "BTC/USD");
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_ohlc_into_candle() {
let raw = r#"{
"channel":"ohlc","type":"snapshot",
"data":[{
"symbol":"BTC/USD","interval":1,
"open":96000.0,"high":96100.0,"low":95900.0,"close":96050.0,
"trades":100,"volume":10.5,"vwap":96025.0,
"interval_begin":"2026-05-25T12:00:00.000000Z"
}]
}"#;
let msgs = connector().parse_message(raw).expect("parse");
match &msgs[0] {
DataMessage::Candle(c) => {
assert_eq!(c.symbol, "BTC/USD");
assert_eq!(c.interval, "1");
assert!((c.open - 96_000.0).abs() < 1e-9);
assert!((c.close - 96_050.0).abs() < 1e-9);
assert!(!c.is_closed);
assert!(c.open_ts > 1_700_000_000_000);
}
other => panic!("expected Candle, got {other:?}"),
}
}
#[test]
fn parse_book_snapshot_and_delta() {
let snap = r#"{
"channel":"book","type":"snapshot",
"data":[{
"symbol":"BTC/USD",
"bids":[{"price":96000.0,"qty":1.5},{"price":95999.0,"qty":2.0}],
"asks":[{"price":96001.0,"qty":0.5}],
"checksum":12345
}]
}"#;
let delta = r#"{
"channel":"book","type":"update",
"data":[{
"symbol":"BTC/USD",
"bids":[{"price":95998.0,"qty":0.0}],
"asks":[],
"checksum":12346
}]
}"#;
let c = connector();
match &c.parse_message(snap).expect("snap")[0] {
DataMessage::OrderBook(ob) => {
assert!(ob.is_snapshot);
assert_eq!(ob.bids.len(), 2);
assert!((ob.asks[0][0] - 96_001.0).abs() < 1e-9);
}
_ => 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#"{"channel":"executions","data":[{"x":"y"}]}"#;
let msgs = connector().parse_message(raw).expect("parse");
assert!(msgs.is_empty());
}
}