1use 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}