Skip to main content

openlimits_coinbase/model/websocket/
mod.rs

1use openlimits_exchange::errors::OpenLimitsError;
2use openlimits_exchange::model::websocket::OpenLimitsWebSocketMessage;
3use openlimits_exchange::model::websocket::WebSocketResponse;
4use openlimits_exchange::model::AskBid;
5use openlimits_exchange::model::OrderBookResponse;
6use openlimits_exchange::shared::Result;
7use std::convert::{TryFrom, TryInto};
8
9use super::OrderSide;
10
11mod activate;
12mod auth;
13mod change;
14mod channel_type;
15mod channel;
16mod coinbase_subscription;
17mod coinbase_websocket_message;
18mod done;
19mod full;
20mod input_message;
21mod level2_snapshot_record;
22mod level2_update_record;
23mod level2;
24mod match_;
25mod open;
26mod reason;
27mod received;
28mod stop_type;
29mod subscribe_cmd;
30mod subscribe;
31mod ticker;
32
33pub use activate::Activate;
34pub use auth::Auth;
35pub use change::Change;
36pub use channel_type::ChannelType;
37pub use channel::Channel;
38pub use coinbase_subscription::CoinbaseSubscription;
39pub use coinbase_websocket_message::CoinbaseWebsocketMessage;
40pub use done::Done;
41pub use full::Full;
42pub use input_message::InputMessage;
43pub use level2_snapshot_record::Level2SnapshotRecord;
44pub use level2_update_record::Level2UpdateRecord;
45pub use level2::Level2;
46pub use match_::Match;
47pub use open::Open;
48pub use reason::Reason;
49pub use received::Received;
50pub use stop_type::StopType;
51pub use subscribe_cmd::SubscribeCmd;
52pub use subscribe::Subscribe;
53pub use ticker::Ticker;
54pub use super::shared;
55use openlimits_exchange::model::Trade;
56
57impl TryFrom<CoinbaseWebsocketMessage> for WebSocketResponse<CoinbaseWebsocketMessage> {
58    type Error = OpenLimitsError;
59
60    fn try_from(value: CoinbaseWebsocketMessage) -> Result<Self> {
61        match value {
62            CoinbaseWebsocketMessage::Level2(level2) => {
63                Ok(WebSocketResponse::Generic(level2.try_into()?))
64            },
65            CoinbaseWebsocketMessage::Match(match_) => {
66                Ok(WebSocketResponse::Generic(match_.into()))
67            },
68            CoinbaseWebsocketMessage::Full(full) => {
69                Ok(WebSocketResponse::Generic(full.into()))
70            },
71            _ => Ok(WebSocketResponse::Raw(value))
72        }
73    }
74}
75
76impl From<Full> for OpenLimitsWebSocketMessage {
77    fn from(from: Full) -> Self {
78        match from {
79            Full::Match(match_) => match_.into(),
80            _ => todo!("Full is not fully implemented :)")
81        }
82    }
83}
84
85impl From<Match> for OpenLimitsWebSocketMessage {
86    fn from(match_: Match) -> Self {
87        let market_pair = match_.product_id;
88        let price = match_.price;
89        let qty = match_.size;
90        let id = format!("{}", match_.trade_id);
91        let buyer_order_id = Some(match_.taker_order_id);
92        let seller_order_id = Some(match_.maker_order_id);
93        let created_at = match_.time;
94        let fees = None;
95        let liquidity = None;
96        let side = match_.side.into();
97        let trade = Trade { market_pair, price, qty, id, buyer_order_id, created_at, fees, liquidity, seller_order_id, side };
98        let trades = vec![trade];
99        Self::Trades(trades)
100    }
101}
102
103impl TryFrom<Level2> for OpenLimitsWebSocketMessage {
104    type Error = OpenLimitsError;
105
106    fn try_from(level2: Level2) -> std::result::Result<Self, Self::Error> {
107        // FIXME: How can we get the update id?
108        let last_update_id = None;
109        let update_id = None;
110        Ok(match level2 {
111            Level2::Snapshot { asks, bids, .. } => {
112                let bids = bids.iter().map(|bid| bid.into()).collect();
113                let asks = asks.iter().map(|ask| ask.into()).collect();
114                let order_book_response = OrderBookResponse {
115                    bids,
116                    asks,
117                    update_id,
118                    last_update_id,
119                };
120                OpenLimitsWebSocketMessage::OrderBook(order_book_response)
121            }
122            Level2::L2update { changes, .. } => {
123                let bids = changes
124                    .iter()
125                    .filter(|change| change.side == OrderSide::Buy)
126                    .map(|change| change.into())
127                    .collect();
128                let asks = changes
129                    .iter()
130                    .filter(|change| change.side == OrderSide::Sell)
131                    .map(|change| change.into())
132                    .collect();
133                let order_book_response = OrderBookResponse {
134                    bids,
135                    asks,
136                    update_id,
137                    last_update_id,
138                };
139                OpenLimitsWebSocketMessage::OrderBook(order_book_response)
140            }
141        })
142    }
143}
144
145impl From<&Level2SnapshotRecord> for AskBid {
146    fn from(record: &Level2SnapshotRecord) -> Self {
147        let price = record.price;
148        let qty = record.size;
149        Self { price, qty }
150    }
151}
152
153impl From<&Level2UpdateRecord> for AskBid {
154    fn from(record: &Level2UpdateRecord) -> Self {
155        let price = record.price;
156        let qty = record.size;
157        Self { price, qty }
158    }
159}