async_kraken_ws 0.1.2

Minimal wrapper for the Kraken exchange WebSockets using async-std
Documentation
use async_kraken::client::KrakenClient;
use async_std::channel::{SendError, Sender};
use futures::{pin_mut, StreamExt};
use serde_json::{json, Value};
use tungstenite::Message;

use crate::config::{API_URL_PRIVATE, API_URL_PUBLIC};

pub struct KrakenWSHandler {
    token: Option<String>,
    sx: Sender<Value>,
}

impl KrakenWSHandler {
    pub async fn ping(&self) -> Result<(), SendError<Value>> {
        self.sx.send(json!({"event": "ping"})).await
    }

    async fn subscribe(&self, pairs: Vec<&str>, payload: Value) -> Result<(), SendError<Value>> {
        let json = if pairs.len() > 0 {
            json!({ "event": "subscribe", "pair" : pairs, "subscription": payload })
        } else {
            json!({ "event": "subscribe", "subscription": payload })
        };
        println!("{}", json.to_string());
        self.sx.send(json).await
    }

    pub async fn subscribe_ticker(&self, pairs: Vec<&str>) -> Result<(), SendError<Value>> {
        self.subscribe(pairs, json!({ "name": "ticker"})).await
    }

    pub async fn subscribe_ohlc(
        &self,
        pairs: Vec<&str>,
        interval: usize,
    ) -> Result<(), SendError<Value>> {
        self.subscribe(pairs, json!({ "name": "ohlc", "interval": interval}))
            .await
    }

    pub async fn subscribe_trade(&self, pairs: Vec<&str>) -> Result<(), SendError<Value>> {
        self.subscribe(pairs, json!({ "name": "trade"})).await
    }

    pub async fn subscribe_spread(&self, pairs: Vec<&str>) -> Result<(), SendError<Value>> {
        self.subscribe(pairs, json!({ "name": "spread"})).await
    }

    pub async fn subscribe_book(
        &self,
        pairs: Vec<&str>,
        depth: usize,
    ) -> Result<(), SendError<Value>> {
        self.subscribe(pairs, json!({ "name": "book", "depth":depth}))
            .await
    }

    pub async fn subscribe_own_trades(&self) -> Result<(), SendError<Value>> {
        self.subscribe(vec![], json!({ "name": "ownTrades", "token": self.token}))
            .await
    }

    pub async fn subscribe_open_orders(&self) -> Result<(), SendError<Value>> {
        self.subscribe(vec![], json!({ "name": "openOrders", "token": self.token}))
            .await
    }

    pub async fn send_private(&self, payload: Value) -> Result<(), SendError<Value>> {
        let mut payload = payload.clone();
        if let Some(token) = &self.token {
            let m = payload.as_object_mut().unwrap();
            m.insert(String::from("token"), json!(token));
            payload = serde_json::to_value(m).unwrap();
        }
        self.sx.send(payload).await
    }

    pub async fn send_public(&self, payload: Value) -> Result<(), SendError<Value>> {
        self.sx.send(payload).await
    }
}

pub struct KrakenWS {}
impl KrakenWS {
    pub async fn run_public(event_callback: fn(Value)) -> KrakenWSHandler {
        let sx = Self::connect(event_callback, API_URL_PUBLIC).await;
        KrakenWSHandler { token: None, sx }
    }

    pub async fn run_private(
        event_callback: fn(Value),
        api_key: String,
        api_secret: String,
    ) -> KrakenWSHandler {
        let result = match KrakenClient::with_credentials(api_key, api_secret)
            .api_request("GetWebSocketsToken", json!({}))
            .await
        {
            Ok(r) => r,
            Err(e) => panic!("{}", e),
        };
        let token = result["token"].as_str().unwrap().to_string();
        let sx = Self::connect(event_callback, API_URL_PRIVATE).await;

        KrakenWSHandler {
            token: Some(token),
            sx,
        }
    }

    async fn connect(event_callback: fn(Value), uri: &str) -> Sender<Value> {
        let (client, _) = async_tungstenite::async_std::connect_async(uri)
            .await
            .expect("Can't connect to websocket server");
        println!("Connected to the Kraken WS server ({})", uri);

        let (sx, rx) = async_std::channel::unbounded::<Value>();
        async_std::task::spawn(async move {
            let (writer, reader) = client.split();
            let send_ws = rx.map(|x| Ok(Message::text(x.to_string()))).forward(writer);
            let ws_callback = reader.for_each(|msg| async {
                let json: Value = serde_json::from_str(&msg.unwrap().to_string()).unwrap();
                event_callback(json);
            });
            pin_mut!(send_ws, ws_callback);
            futures::future::select(send_ws, ws_callback).await;
        });
        sx
    }
}