Skip to main content

polyester/codecs/decode/
realtime.rs

1//! Realtime protobuf publication decoders (Go `codecs/decode/realtime.go` parity).
2
3use buffa::Message;
4
5use super::{
6    account_identity_from_proto, address_book_invalidation_from_proto, api_key_from_proto,
7    api_policy_from_proto, asset_balance_from_proto, candle_point_from_proto,
8    flow_summary_message_from_proto, market_overview_batch_from_proto, market_trade_from_proto,
9    order_from_proto, subaccount_from_proto, subaccount_policy_from_proto, transfer_row_from_proto,
10    trigger_event_from_proto, trigger_from_proto, user_trade_from_proto,
11    zipped_asset_supply_batch_from_proto,
12};
13use crate::errors::{Error, Result};
14use crate::models::{
15    AccountIdentity, AddressBookViewInvalidation, ApiData, ApiKeySummary, ApiPolicy, AssetBalance,
16    Candle, LedgerTransfer, LifecycleFlowSummary, MarketOverviewList, MarketTrade, Order,
17    OrderBookDeltaUpdate, PriceQtyPair, SubAccount, SubaccountPolicy, Trigger, TriggerEvent,
18    UserTrade, ZippedAssetSupplyBatch,
19};
20use crate::proto::auth::v1::{
21    AccountIdentity as ProtoAccountIdentity, AddressBookViewInvalidated, ApiKey as ProtoApiKey,
22    ApiPolicyView, Subaccount, SubaccountPolicyView,
23};
24use crate::proto::chain::lifecycle::v1::{FlowDetailView, FlowSummaryView};
25use crate::proto::chain::zipper::v1::ZippedAssetSupplyBatch as ProtoZippedAssetSupplyBatch;
26use crate::proto::ledger::read::v1::{AssetBalance as ProtoAssetBalance, TransferRow};
27use crate::proto::marketdata::v1::{
28    CandlePoint, HeatmapLiveBucket, MarketTrade as ProtoMarketTrade,
29};
30use crate::proto::marketoverview::v1::MarketOverviewBatch;
31use crate::proto::orderbook::v1::OrderBookDelta;
32use crate::proto::orders::v1::{Order as ProtoOrder, UserTrade as ProtoUserTrade};
33use crate::proto::triggers::v1::{Trigger as ProtoTrigger, TriggerEvent as ProtoTriggerEvent};
34
35fn decode_proto<M: Message + Default>(payload: &[u8]) -> Result<M> {
36    if payload.is_empty() {
37        return Err(Error::realtime(
38            "proto decode: empty publication payload".to_owned(),
39        ));
40    }
41    if payload.len() > crate::realtime::MAX_REALTIME_MESSAGE_BYTES {
42        return Err(Error::realtime(format!(
43            "proto decode: publication exceeds {} bytes",
44            crate::realtime::MAX_REALTIME_MESSAGE_BYTES
45        )));
46    }
47    M::decode_from_slice(payload).map_err(|e| Error::realtime(format!("proto decode: {e}")))
48}
49
50pub fn order_from_bytes(payload: &[u8]) -> Result<Order> {
51    let msg = decode_proto::<ProtoOrder>(payload)?;
52    Ok(order_from_proto(&msg))
53}
54
55pub fn user_trade_from_bytes(payload: &[u8]) -> Result<UserTrade> {
56    let msg = decode_proto::<ProtoUserTrade>(payload)?;
57    Ok(user_trade_from_proto(&msg))
58}
59
60pub fn asset_balance_from_bytes(payload: &[u8]) -> Result<AssetBalance> {
61    let msg = decode_proto::<ProtoAssetBalance>(payload)?;
62    Ok(asset_balance_from_proto(&msg))
63}
64
65pub fn ledger_transfer_from_bytes(payload: &[u8]) -> Result<LedgerTransfer> {
66    let msg = decode_proto::<TransferRow>(payload)?;
67    Ok(transfer_row_from_proto(&msg))
68}
69
70pub fn trigger_from_bytes(payload: &[u8]) -> Result<Trigger> {
71    let msg = decode_proto::<ProtoTrigger>(payload)?;
72    Ok(trigger_from_proto(&msg))
73}
74
75pub fn trigger_event_from_bytes(payload: &[u8]) -> Result<TriggerEvent> {
76    let msg = decode_proto::<ProtoTriggerEvent>(payload)?;
77    Ok(trigger_event_from_proto(&msg))
78}
79
80pub fn market_trade_from_bytes(
81    quantity_scale: u32,
82) -> impl Fn(&[u8]) -> Result<MarketTrade> + Send + Sync + 'static {
83    move |payload: &[u8]| {
84        let msg = decode_proto::<ProtoMarketTrade>(payload)?;
85        Ok(market_trade_from_proto(&msg, quantity_scale))
86    }
87}
88
89pub fn orderbook_delta_from_bytes(payload: &[u8]) -> Result<OrderBookDeltaUpdate> {
90    let msg = decode_proto::<OrderBookDelta>(payload)?;
91    Ok(OrderBookDeltaUpdate {
92        symbol_id: msg.symbol_id,
93        book_seq_start: msg.book_seq_start,
94        book_seq_end: msg.book_seq_end,
95        reset: msg.reset,
96        bids: msg
97            .bids
98            .iter()
99            .map(|l| PriceQtyPair {
100                price_ticks: l.price_ticks,
101                qty_scaled: l.qty_scaled,
102            })
103            .collect(),
104        asks: msg
105            .asks
106            .iter()
107            .map(|l| PriceQtyPair {
108                price_ticks: l.price_ticks,
109                qty_scaled: l.qty_scaled,
110            })
111            .collect(),
112    })
113}
114
115pub fn flow_summary_from_bytes(payload: &[u8]) -> Result<LifecycleFlowSummary> {
116    let msg = decode_proto::<FlowSummaryView>(payload)?;
117    Ok(flow_summary_message_from_proto(&msg))
118}
119
120pub fn flow_detail_from_bytes(payload: &[u8]) -> Result<LifecycleFlowSummary> {
121    let msg = decode_proto::<FlowDetailView>(payload)?;
122    msg.summary
123        .as_option()
124        .map(flow_summary_message_from_proto)
125        .ok_or_else(|| Error::realtime("proto decode: flow detail is missing summary".to_owned()))
126}
127
128pub fn account_identity_from_bytes(payload: &[u8]) -> Result<AccountIdentity> {
129    let msg = decode_proto::<ProtoAccountIdentity>(payload)?;
130    Ok(account_identity_from_proto(&msg))
131}
132
133pub fn candle_point_from_bytes(
134    symbol_id: u32,
135    timeframe: String,
136    volume_scale: u32,
137) -> impl Fn(&[u8]) -> Result<Candle> + Send + Sync + 'static {
138    move |payload: &[u8]| {
139        let point = decode_proto::<CandlePoint>(payload)?;
140        candle_point_from_proto(&point, volume_scale, symbol_id, &timeframe)
141    }
142}
143
144pub fn heatmap_live_bucket_from_bytes(payload: &[u8]) -> Result<ApiData> {
145    let msg = decode_proto::<HeatmapLiveBucket>(payload)?;
146    Ok(super::api_data_from_proto(&msg))
147}
148
149pub fn market_overview_batch_from_bytes(payload: &[u8]) -> Result<MarketOverviewList> {
150    let msg = decode_proto::<MarketOverviewBatch>(payload)?;
151    Ok(market_overview_batch_from_proto(&msg))
152}
153
154pub fn zipped_asset_supply_batch_from_bytes(
155    payload: &[u8],
156    scale_fn: impl Fn(u32) -> Option<u32>,
157) -> Result<ZippedAssetSupplyBatch> {
158    let msg = decode_proto::<ProtoZippedAssetSupplyBatch>(payload)?;
159    zipped_asset_supply_batch_from_proto(&msg, scale_fn)
160}
161
162pub fn api_key_from_bytes(payload: &[u8]) -> Result<ApiKeySummary> {
163    let msg = decode_proto::<ProtoApiKey>(payload)?;
164    Ok(api_key_from_proto(&msg))
165}
166
167pub fn subaccount_from_bytes(payload: &[u8]) -> Result<SubAccount> {
168    let msg = decode_proto::<Subaccount>(payload)?;
169    Ok(subaccount_from_proto(&msg))
170}
171
172pub fn subaccount_policy_from_bytes(payload: &[u8]) -> Result<SubaccountPolicy> {
173    let msg = decode_proto::<SubaccountPolicyView>(payload)?;
174    Ok(subaccount_policy_from_proto(&msg))
175}
176
177pub fn api_policy_from_bytes(payload: &[u8]) -> Result<ApiPolicy> {
178    let msg = decode_proto::<ApiPolicyView>(payload)?;
179    Ok(api_policy_from_proto(&msg))
180}
181
182pub fn address_book_invalidation_from_bytes(payload: &[u8]) -> Result<AddressBookViewInvalidation> {
183    let msg = decode_proto::<AddressBookViewInvalidated>(payload)?;
184    Ok(address_book_invalidation_from_proto(&msg))
185}
186
187#[cfg(test)]
188mod tests {
189    use super::*;
190    use crate::codecs::scalars::format_uint64_id;
191    use crate::proto::orders::v1::{OrderStatus, OrderType, Side, TimeInForce};
192
193    #[test]
194    fn order_from_bytes_round_trip() {
195        let msg = ProtoOrder {
196            order_id: 42,
197            symbol_id: 3,
198            client_order_id: "coid".into(),
199            side: Side::Buy.into(),
200            status: OrderStatus::Working.into(),
201            order_type: OrderType::Limit.into(),
202            time_in_force: TimeInForce::Gtc.into(),
203            ..Default::default()
204        };
205        let bytes = msg.encode_to_vec();
206        let order = order_from_bytes(&bytes).expect("decode");
207        assert_eq!(order.order_id, format_uint64_id(42));
208        assert_eq!(order.side, "buy");
209        assert_eq!(order.status, "working");
210    }
211
212    #[test]
213    fn api_policy_from_bytes_round_trip() {
214        use crate::proto::auth::v1::ApiPolicyView;
215
216        let msg = ApiPolicyView {
217            id: 9,
218            name: "bots".into(),
219            description: "api key policy".into(),
220            revision: 4,
221            ..Default::default()
222        };
223        let bytes = msg.encode_to_vec();
224        let policy = api_policy_from_bytes(&bytes).expect("decode");
225        assert_eq!(policy.policy_id, format_uint64_id(9));
226        assert_eq!(policy.name, "bots");
227        assert_eq!(policy.revision, 4);
228    }
229
230    #[test]
231    fn subaccount_policy_from_bytes_round_trip() {
232        use crate::proto::auth::v1::SubaccountPolicyView;
233
234        let msg = SubaccountPolicyView {
235            id: 7,
236            name: "trader".into(),
237            revision: 3,
238            ..Default::default()
239        };
240        let bytes = msg.encode_to_vec();
241        let policy = subaccount_policy_from_bytes(&bytes).expect("decode");
242        assert_eq!(policy.policy_id, format_uint64_id(7));
243        assert_eq!(policy.name, "trader");
244        assert_eq!(policy.revision, 3);
245    }
246
247    #[test]
248    fn market_trade_from_bytes_maps_side() {
249        let msg = ProtoMarketTrade {
250            symbol_id: 1,
251            match_id: 99,
252            is_buy: true,
253            price_ticks: 100,
254            qty_scaled: 5,
255            ts_ns: 123,
256            ..Default::default()
257        };
258        let bytes = msg.encode_to_vec();
259        let trade = market_trade_from_bytes(6)(&bytes).expect("decode");
260        assert_eq!(trade.match_id, "99");
261        assert_eq!(trade.side, "buy");
262        assert_eq!(
263            trade.qty.as_ref().unwrap().format(None).unwrap(),
264            "0.000005"
265        );
266    }
267
268    #[test]
269    fn flow_detail_without_required_summary_fails_closed() {
270        let msg = FlowDetailView {
271            from_live_state: true,
272            ..Default::default()
273        };
274        let error = flow_detail_from_bytes(&msg.encode_to_vec())
275            .expect_err("missing flow summary must not become an empty success");
276        assert!(error.to_string().contains("missing summary"));
277    }
278}