1use std::collections::HashSet;
36
37use bytes::Bytes;
38use nautilus_common::{
39 messages::{
40 data::{
41 BarsResponse, BookDeltasResponse, BookDepthResponse, BookResponse, CustomDataResponse,
42 DataCommand, DataResponse, ForwardPricesResponse, FundingRatesResponse,
43 InstrumentResponse, InstrumentsResponse, QuotesResponse, TradesResponse,
44 },
45 execution::{
46 BatchCancelOrders, BatchModifyOrders, CancelAllOrders, CancelOrder, ExecutionReport,
47 ModifyOrder, QueryAccount, QueryOrder, SubmitOrder, SubmitOrderList, TradingCommand,
48 },
49 },
50 timer::TimeEvent,
51};
52use nautilus_core::{Params, UUID4, UnixNanos};
53use nautilus_model::{
54 data::DataType,
55 events::{
56 AccountState, OrderAccepted, OrderCancelRejected, OrderCanceled, OrderDenied,
57 OrderEmulated, OrderEventAny, OrderExpired, OrderFillVoided, OrderFilled, OrderInitialized,
58 OrderModifyRejected, OrderPendingCancel, OrderPendingUpdate, OrderRejected, OrderReleased,
59 OrderSubmitted, OrderTriggered, OrderUpdated, PositionAdjusted, PositionChanged,
60 PositionClosed, PositionEvent, PositionOpened,
61 },
62 identifiers::{ClientId, InstrumentId, Venue},
63 reports::{ExecutionMassStatus, FillReport, OrderStatusReport, PositionStatusReport},
64};
65use serde::Serialize;
66use ustr::Ustr;
67
68use crate::{
69 backend::{IndexKey, IndexKind},
70 capture::{
71 encoder::{EncodeError, EncodedPayload},
72 registry::EncoderRegistry,
73 },
74 entry::PayloadType,
75 headers::Headers,
76};
77
78pub const PAYLOAD_TYPE_SUBMIT_ORDER: &str = "SubmitOrder";
80pub const PAYLOAD_TYPE_SUBMIT_ORDER_LIST: &str = "SubmitOrderList";
82pub const PAYLOAD_TYPE_MODIFY_ORDER: &str = "ModifyOrder";
84pub const PAYLOAD_TYPE_BATCH_MODIFY_ORDERS: &str = "BatchModifyOrders";
86pub const PAYLOAD_TYPE_CANCEL_ORDER: &str = "CancelOrder";
88pub const PAYLOAD_TYPE_CANCEL_ALL_ORDERS: &str = "CancelAllOrders";
90pub const PAYLOAD_TYPE_BATCH_CANCEL_ORDERS: &str = "BatchCancelOrders";
92pub const PAYLOAD_TYPE_QUERY_ORDER: &str = "QueryOrder";
94pub const PAYLOAD_TYPE_QUERY_ACCOUNT: &str = "QueryAccount";
96
97pub const PAYLOAD_TYPE_ORDER_INITIALIZED: &str = "OrderInitialized";
99pub const PAYLOAD_TYPE_ORDER_DENIED: &str = "OrderDenied";
101pub const PAYLOAD_TYPE_ORDER_EMULATED: &str = "OrderEmulated";
103pub const PAYLOAD_TYPE_ORDER_RELEASED: &str = "OrderReleased";
105pub const PAYLOAD_TYPE_ORDER_SUBMITTED: &str = "OrderSubmitted";
107pub const PAYLOAD_TYPE_ORDER_ACCEPTED: &str = "OrderAccepted";
109pub const PAYLOAD_TYPE_ORDER_REJECTED: &str = "OrderRejected";
111pub const PAYLOAD_TYPE_ORDER_CANCELED: &str = "OrderCanceled";
113pub const PAYLOAD_TYPE_ORDER_EXPIRED: &str = "OrderExpired";
115pub const PAYLOAD_TYPE_ORDER_TRIGGERED: &str = "OrderTriggered";
117pub const PAYLOAD_TYPE_ORDER_PENDING_UPDATE: &str = "OrderPendingUpdate";
119pub const PAYLOAD_TYPE_ORDER_PENDING_CANCEL: &str = "OrderPendingCancel";
121pub const PAYLOAD_TYPE_ORDER_MODIFY_REJECTED: &str = "OrderModifyRejected";
123pub const PAYLOAD_TYPE_ORDER_CANCEL_REJECTED: &str = "OrderCancelRejected";
125pub const PAYLOAD_TYPE_ORDER_UPDATED: &str = "OrderUpdated";
127pub const PAYLOAD_TYPE_ORDER_FILLED: &str = "OrderFilled";
129pub const PAYLOAD_TYPE_ORDER_FILL_VOIDED: &str = "OrderFillVoided";
131pub const PAYLOAD_TYPE_ORDER_STATUS_REPORT: &str = "OrderStatusReport";
133pub const PAYLOAD_TYPE_FILL_REPORT: &str = "FillReport";
135pub const PAYLOAD_TYPE_ORDER_WITH_FILLS: &str = "OrderWithFills";
137pub const PAYLOAD_TYPE_POSITION_STATUS_REPORT: &str = "PositionStatusReport";
139pub const PAYLOAD_TYPE_EXECUTION_MASS_STATUS: &str = "ExecutionMassStatus";
141pub const PAYLOAD_TYPE_POSITION_OPENED: &str = "PositionOpened";
143pub const PAYLOAD_TYPE_POSITION_CHANGED: &str = "PositionChanged";
145pub const PAYLOAD_TYPE_POSITION_CLOSED: &str = "PositionClosed";
147pub const PAYLOAD_TYPE_POSITION_ADJUSTED: &str = "PositionAdjusted";
149pub const PAYLOAD_TYPE_ACCOUNT_STATE: &str = "AccountState";
151pub const PAYLOAD_TYPE_TIME_EVENT: &str = "TimeEvent";
153
154pub const PAYLOAD_TYPE_REQUEST_COMMAND: &str = "RequestCommand";
156pub const PAYLOAD_TYPE_SUBSCRIBE_COMMAND: &str = "SubscribeCommand";
158pub const PAYLOAD_TYPE_UNSUBSCRIBE_COMMAND: &str = "UnsubscribeCommand";
160#[cfg(feature = "defi")]
161pub const PAYLOAD_TYPE_DEFI_REQUEST_COMMAND: &str = "DefiRequestCommand";
163#[cfg(feature = "defi")]
164pub const PAYLOAD_TYPE_DEFI_SUBSCRIBE_COMMAND: &str = "DefiSubscribeCommand";
166#[cfg(feature = "defi")]
167pub const PAYLOAD_TYPE_DEFI_UNSUBSCRIBE_COMMAND: &str = "DefiUnsubscribeCommand";
169
170pub const PAYLOAD_TYPE_CUSTOM_DATA_RESPONSE: &str = "CustomDataResponse";
172pub const PAYLOAD_TYPE_INSTRUMENT_RESPONSE: &str = "InstrumentResponse";
174pub const PAYLOAD_TYPE_INSTRUMENTS_RESPONSE: &str = "InstrumentsResponse";
176pub const PAYLOAD_TYPE_BOOK_RESPONSE: &str = "BookResponse";
178pub const PAYLOAD_TYPE_BOOK_DELTAS_RESPONSE: &str = "BookDeltasResponse";
180pub const PAYLOAD_TYPE_BOOK_DEPTH_RESPONSE: &str = "BookDepthResponse";
182pub const PAYLOAD_TYPE_QUOTES_RESPONSE: &str = "QuotesResponse";
184pub const PAYLOAD_TYPE_TRADES_RESPONSE: &str = "TradesResponse";
186pub const PAYLOAD_TYPE_FUNDING_RATES_RESPONSE: &str = "FundingRatesResponse";
188pub const PAYLOAD_TYPE_FORWARD_PRICES_RESPONSE: &str = "ForwardPricesResponse";
190pub const PAYLOAD_TYPE_BARS_RESPONSE: &str = "BarsResponse";
192
193const PAYLOAD_TYPE_TRADING_COMMAND: &str = "TradingCommand";
198
199const PAYLOAD_TYPE_ORDER_EVENT_ANY: &str = "OrderEventAny";
200
201const PAYLOAD_TYPE_EXECUTION_REPORT: &str = "ExecutionReport";
202
203const PAYLOAD_TYPE_POSITION_EVENT: &str = "PositionEvent";
204
205const PAYLOAD_TYPE_DATA_COMMAND: &str = "DataCommand";
206
207const PAYLOAD_TYPE_DATA_RESPONSE: &str = "DataResponse";
208
209#[cfg(test)]
210pub(crate) const DEFAULT_CAPTURE_PAYLOAD_TYPES: &[&str] = &[
211 PAYLOAD_TYPE_SUBMIT_ORDER,
212 PAYLOAD_TYPE_SUBMIT_ORDER_LIST,
213 PAYLOAD_TYPE_MODIFY_ORDER,
214 PAYLOAD_TYPE_BATCH_MODIFY_ORDERS,
215 PAYLOAD_TYPE_CANCEL_ORDER,
216 PAYLOAD_TYPE_CANCEL_ALL_ORDERS,
217 PAYLOAD_TYPE_BATCH_CANCEL_ORDERS,
218 PAYLOAD_TYPE_QUERY_ORDER,
219 PAYLOAD_TYPE_QUERY_ACCOUNT,
220 PAYLOAD_TYPE_ORDER_INITIALIZED,
221 PAYLOAD_TYPE_ORDER_DENIED,
222 PAYLOAD_TYPE_ORDER_EMULATED,
223 PAYLOAD_TYPE_ORDER_RELEASED,
224 PAYLOAD_TYPE_ORDER_SUBMITTED,
225 PAYLOAD_TYPE_ORDER_ACCEPTED,
226 PAYLOAD_TYPE_ORDER_REJECTED,
227 PAYLOAD_TYPE_ORDER_CANCELED,
228 PAYLOAD_TYPE_ORDER_EXPIRED,
229 PAYLOAD_TYPE_ORDER_TRIGGERED,
230 PAYLOAD_TYPE_ORDER_PENDING_UPDATE,
231 PAYLOAD_TYPE_ORDER_PENDING_CANCEL,
232 PAYLOAD_TYPE_ORDER_MODIFY_REJECTED,
233 PAYLOAD_TYPE_ORDER_CANCEL_REJECTED,
234 PAYLOAD_TYPE_ORDER_UPDATED,
235 PAYLOAD_TYPE_ORDER_FILLED,
236 PAYLOAD_TYPE_ORDER_FILL_VOIDED,
237 PAYLOAD_TYPE_ORDER_STATUS_REPORT,
238 PAYLOAD_TYPE_FILL_REPORT,
239 PAYLOAD_TYPE_ORDER_WITH_FILLS,
240 PAYLOAD_TYPE_POSITION_STATUS_REPORT,
241 PAYLOAD_TYPE_EXECUTION_MASS_STATUS,
242 PAYLOAD_TYPE_POSITION_OPENED,
243 PAYLOAD_TYPE_POSITION_CHANGED,
244 PAYLOAD_TYPE_POSITION_CLOSED,
245 PAYLOAD_TYPE_POSITION_ADJUSTED,
246 PAYLOAD_TYPE_ACCOUNT_STATE,
247 PAYLOAD_TYPE_TIME_EVENT,
248 PAYLOAD_TYPE_REQUEST_COMMAND,
249 PAYLOAD_TYPE_SUBSCRIBE_COMMAND,
250 PAYLOAD_TYPE_UNSUBSCRIBE_COMMAND,
251 #[cfg(feature = "defi")]
252 PAYLOAD_TYPE_DEFI_REQUEST_COMMAND,
253 #[cfg(feature = "defi")]
254 PAYLOAD_TYPE_DEFI_SUBSCRIBE_COMMAND,
255 #[cfg(feature = "defi")]
256 PAYLOAD_TYPE_DEFI_UNSUBSCRIBE_COMMAND,
257 PAYLOAD_TYPE_CUSTOM_DATA_RESPONSE,
258 PAYLOAD_TYPE_INSTRUMENT_RESPONSE,
259 PAYLOAD_TYPE_INSTRUMENTS_RESPONSE,
260 PAYLOAD_TYPE_BOOK_RESPONSE,
261 PAYLOAD_TYPE_BOOK_DELTAS_RESPONSE,
262 PAYLOAD_TYPE_BOOK_DEPTH_RESPONSE,
263 PAYLOAD_TYPE_QUOTES_RESPONSE,
264 PAYLOAD_TYPE_TRADES_RESPONSE,
265 PAYLOAD_TYPE_FUNDING_RATES_RESPONSE,
266 PAYLOAD_TYPE_FORWARD_PRICES_RESPONSE,
267 PAYLOAD_TYPE_BARS_RESPONSE,
268];
269
270#[must_use]
275pub fn default_registry() -> EncoderRegistry {
276 let mut registry = EncoderRegistry::new();
277 register_default(&mut registry);
278 registry
279}
280
281pub fn register_default(registry: &mut EncoderRegistry) {
302 registry
303 .register::<SubmitOrder, _>(payload_type(PAYLOAD_TYPE_SUBMIT_ORDER), encode_submit_order);
304 registry
305 .register::<OrderFilled, _>(payload_type(PAYLOAD_TYPE_ORDER_FILLED), encode_order_filled);
306 registry.register::<OrderStatusReport, _>(
307 payload_type(PAYLOAD_TYPE_ORDER_STATUS_REPORT),
308 encode_order_status_report,
309 );
310 registry.register::<FillReport, _>(payload_type(PAYLOAD_TYPE_FILL_REPORT), encode_fill_report);
311 registry.register::<PositionStatusReport, _>(
312 payload_type(PAYLOAD_TYPE_POSITION_STATUS_REPORT),
313 encode_position_status_report,
314 );
315 registry.register::<TradingCommand, _>(
316 payload_type(PAYLOAD_TYPE_TRADING_COMMAND),
317 encode_trading_command,
318 );
319 registry.register::<OrderEventAny, _>(
320 payload_type(PAYLOAD_TYPE_ORDER_EVENT_ANY),
321 encode_order_event_any,
322 );
323 registry.register::<ExecutionReport, _>(
324 payload_type(PAYLOAD_TYPE_EXECUTION_REPORT),
325 encode_execution_report,
326 );
327 registry.register::<PositionEvent, _>(
328 payload_type(PAYLOAD_TYPE_POSITION_EVENT),
329 encode_position_event,
330 );
331 registry.register::<AccountState, _>(
332 payload_type(PAYLOAD_TYPE_ACCOUNT_STATE),
333 encode_account_state,
334 );
335 registry.register::<TimeEvent, _>(payload_type(PAYLOAD_TYPE_TIME_EVENT), encode_time_event);
336 registry
337 .register::<DataCommand, _>(payload_type(PAYLOAD_TYPE_DATA_COMMAND), encode_data_command);
338 registry.register::<DataResponse, _>(
339 payload_type(PAYLOAD_TYPE_DATA_RESPONSE),
340 encode_data_response,
341 );
342
343 register_default_headers(registry);
344 register_default_identities(registry);
345}
346
347fn register_default_headers(registry: &mut EncoderRegistry) {
355 registry.register_headers::<SubmitOrder, _>(extract_submit_order_headers);
356 registry.register_headers::<SubmitOrderList, _>(extract_submit_order_list_headers);
357 registry.register_headers::<ModifyOrder, _>(extract_modify_order_headers);
358 registry.register_headers::<BatchModifyOrders, _>(extract_batch_modify_orders_headers);
359 registry.register_headers::<CancelOrder, _>(extract_cancel_order_headers);
360 registry.register_headers::<CancelAllOrders, _>(extract_cancel_all_orders_headers);
361 registry.register_headers::<BatchCancelOrders, _>(extract_batch_cancel_orders_headers);
362 registry.register_headers::<QueryOrder, _>(extract_query_order_headers);
363 registry.register_headers::<QueryAccount, _>(extract_query_account_headers);
364 registry.register_headers::<TradingCommand, _>(extract_trading_command_headers);
365 registry.register_headers::<DataCommand, _>(extract_data_command_headers);
366 registry.register_headers::<DataResponse, _>(extract_data_response_headers);
367}
368
369fn register_default_identities(registry: &mut EncoderRegistry) {
377 registry.register_identity::<SubmitOrder, _>(|command| Some(command.command_id));
378 registry.register_identity::<OrderFilled, _>(|event| Some(event.event_id));
379 registry.register_identity::<TradingCommand, _>(|c| Some(extract_trading_command_identity(c)));
380 registry.register_identity::<OrderEventAny, _>(|e| Some(extract_order_event_any_identity(e)));
381 registry.register_identity::<AccountState, _>(|state| Some(state.event_id));
382 registry.register_identity::<DataCommand, _>(extract_data_command_identity);
383}
384
385fn extract_data_command_identity(command: &DataCommand) -> Option<UUID4> {
386 match command {
387 DataCommand::Request(cmd) => Some(*cmd.request_id()),
388 DataCommand::Subscribe(cmd) => Some(cmd.command_id()),
389 DataCommand::Unsubscribe(cmd) => Some(cmd.command_id()),
390 #[cfg(feature = "defi")]
391 DataCommand::DefiRequest(cmd) => Some(*cmd.request_id()),
392 #[cfg(feature = "defi")]
393 DataCommand::DefiSubscribe(cmd) => Some(cmd.command_id()),
394 #[cfg(feature = "defi")]
395 DataCommand::DefiUnsubscribe(cmd) => Some(cmd.command_id()),
396 _ => None,
398 }
399}
400
401fn extract_trading_command_identity(command: &TradingCommand) -> UUID4 {
402 match command {
403 TradingCommand::SubmitOrder(c) => c.command_id,
404 TradingCommand::SubmitOrderList(c) => c.command_id,
405 TradingCommand::ModifyOrder(c) => c.command_id,
406 TradingCommand::ModifyOrders(c) => c.command_id,
407 TradingCommand::CancelOrder(c) => c.command_id,
408 TradingCommand::CancelOrders(c) => c.command_id,
409 TradingCommand::CancelAllOrders(c) => c.command_id,
410 TradingCommand::QueryOrder(c) => c.command_id,
411 TradingCommand::QueryAccount(c) => c.command_id,
412 }
413}
414
415fn extract_order_event_any_identity(event: &OrderEventAny) -> UUID4 {
416 match event {
417 OrderEventAny::Initialized(e) => e.event_id,
418 OrderEventAny::Denied(e) => e.event_id,
419 OrderEventAny::Emulated(e) => e.event_id,
420 OrderEventAny::Released(e) => e.event_id,
421 OrderEventAny::Submitted(e) => e.event_id,
422 OrderEventAny::Accepted(e) => e.event_id,
423 OrderEventAny::Rejected(e) => e.event_id,
424 OrderEventAny::Canceled(e) => e.event_id,
425 OrderEventAny::Expired(e) => e.event_id,
426 OrderEventAny::Triggered(e) => e.event_id,
427 OrderEventAny::PendingUpdate(e) => e.event_id,
428 OrderEventAny::PendingCancel(e) => e.event_id,
429 OrderEventAny::ModifyRejected(e) => e.event_id,
430 OrderEventAny::CancelRejected(e) => e.event_id,
431 OrderEventAny::Updated(e) => e.event_id,
432 OrderEventAny::Filled(e) => e.event_id,
433 OrderEventAny::FillVoided(e) => e.event_id,
434 }
435}
436
437fn headers_from_fields(correlation_id: Option<UUID4>, causation_id: Option<UUID4>) -> Headers {
438 Headers {
439 correlation_id,
440 causation_id,
441 }
442}
443
444fn extract_submit_order_headers(cmd: &SubmitOrder) -> Headers {
445 headers_from_fields(cmd.correlation_id, cmd.causation_id)
446}
447
448fn extract_submit_order_list_headers(cmd: &SubmitOrderList) -> Headers {
449 headers_from_fields(cmd.correlation_id, cmd.causation_id)
450}
451
452fn extract_modify_order_headers(cmd: &ModifyOrder) -> Headers {
453 headers_from_fields(cmd.correlation_id, cmd.causation_id)
454}
455
456fn extract_batch_modify_orders_headers(cmd: &BatchModifyOrders) -> Headers {
457 headers_from_fields(cmd.correlation_id, cmd.causation_id)
458}
459
460fn extract_cancel_order_headers(cmd: &CancelOrder) -> Headers {
461 headers_from_fields(cmd.correlation_id, cmd.causation_id)
462}
463
464fn extract_cancel_all_orders_headers(cmd: &CancelAllOrders) -> Headers {
465 headers_from_fields(cmd.correlation_id, cmd.causation_id)
466}
467
468fn extract_batch_cancel_orders_headers(cmd: &BatchCancelOrders) -> Headers {
469 headers_from_fields(cmd.correlation_id, cmd.causation_id)
470}
471
472fn extract_query_order_headers(cmd: &QueryOrder) -> Headers {
473 headers_from_fields(cmd.correlation_id, cmd.causation_id)
474}
475
476fn extract_query_account_headers(cmd: &QueryAccount) -> Headers {
477 headers_from_fields(cmd.correlation_id, cmd.causation_id)
478}
479
480fn extract_trading_command_headers(command: &TradingCommand) -> Headers {
484 match command {
485 TradingCommand::SubmitOrder(cmd) => extract_submit_order_headers(cmd),
486 TradingCommand::SubmitOrderList(cmd) => extract_submit_order_list_headers(cmd),
487 TradingCommand::ModifyOrder(cmd) => extract_modify_order_headers(cmd),
488 TradingCommand::ModifyOrders(cmd) => extract_batch_modify_orders_headers(cmd),
489 TradingCommand::CancelOrder(cmd) => extract_cancel_order_headers(cmd),
490 TradingCommand::CancelOrders(cmd) => extract_batch_cancel_orders_headers(cmd),
491 TradingCommand::CancelAllOrders(cmd) => extract_cancel_all_orders_headers(cmd),
492 TradingCommand::QueryOrder(cmd) => extract_query_order_headers(cmd),
493 TradingCommand::QueryAccount(cmd) => extract_query_account_headers(cmd),
494 }
495}
496
497fn extract_data_command_headers(command: &DataCommand) -> Headers {
505 match command {
506 DataCommand::Request(cmd) => headers_from_fields(Some(*cmd.request_id()), None),
507 DataCommand::Subscribe(cmd) => headers_from_fields(cmd.correlation_id(), None),
508 DataCommand::Unsubscribe(cmd) => headers_from_fields(cmd.correlation_id(), None),
509 _ => Headers::empty(),
513 }
514}
515
516fn extract_data_response_headers(response: &DataResponse) -> Headers {
519 headers_from_fields(Some(*response.correlation_id()), None)
520}
521
522fn payload_type(tag: &str) -> PayloadType {
523 Ustr::from(tag)
524}
525
526fn encode_serde<T: Serialize>(value: &T) -> Result<Bytes, EncodeError> {
527 rmp_serde::to_vec_named(value)
528 .map(Bytes::from)
529 .map_err(|e| EncodeError::Serialize(e.to_string()))
530}
531
532pub fn encode_submit_order(message: &SubmitOrder) -> Result<EncodedPayload, EncodeError> {
540 let payload = encode_serde(message)?;
541 let index_keys = vec![IndexKey::new(
542 IndexKind::ClientOrderId,
543 message.client_order_id.to_string(),
544 )];
545 Ok(EncodedPayload::new(payload, index_keys))
546}
547
548pub fn encode_order_filled(message: &OrderFilled) -> Result<EncodedPayload, EncodeError> {
555 let payload = encode_serde(message)?;
556 let index_keys = vec![
557 IndexKey::new(
558 IndexKind::ClientOrderId,
559 message.client_order_id.to_string(),
560 ),
561 IndexKey::new(IndexKind::VenueOrderId, message.venue_order_id.to_string()),
562 ];
563 Ok(EncodedPayload::new(payload, index_keys))
564}
565
566pub fn encode_trading_command(command: &TradingCommand) -> Result<EncodedPayload, EncodeError> {
579 match command {
580 TradingCommand::SubmitOrder(cmd) => {
581 Ok(retag(encode_submit_order(cmd)?, PAYLOAD_TYPE_SUBMIT_ORDER))
582 }
583 TradingCommand::SubmitOrderList(cmd) => encode_submit_order_list(cmd),
584 TradingCommand::ModifyOrder(cmd) => encode_modify_order(cmd),
585 TradingCommand::ModifyOrders(cmd) => encode_batch_modify_orders(cmd),
586 TradingCommand::CancelOrder(cmd) => encode_cancel_order(cmd),
587 TradingCommand::CancelOrders(cmd) => encode_batch_cancel_orders(cmd),
588 TradingCommand::CancelAllOrders(cmd) => encode_cancel_all_orders(cmd),
589 TradingCommand::QueryOrder(cmd) => encode_query_order(cmd),
590 TradingCommand::QueryAccount(cmd) => encode_query_account(cmd),
591 }
592}
593
594pub fn encode_order_event_any(event: &OrderEventAny) -> Result<EncodedPayload, EncodeError> {
605 match event {
606 OrderEventAny::Initialized(e) => encode_order_initialized(e),
607 OrderEventAny::Denied(e) => encode_order_denied(e),
608 OrderEventAny::Emulated(e) => encode_order_emulated(e),
609 OrderEventAny::Released(e) => encode_order_released(e),
610 OrderEventAny::Submitted(e) => encode_order_submitted(e),
611 OrderEventAny::Accepted(e) => encode_order_accepted(e),
612 OrderEventAny::Rejected(e) => encode_order_rejected(e),
613 OrderEventAny::Canceled(e) => encode_order_canceled(e),
614 OrderEventAny::Expired(e) => encode_order_expired(e),
615 OrderEventAny::Triggered(e) => encode_order_triggered(e),
616 OrderEventAny::PendingUpdate(e) => encode_order_pending_update(e),
617 OrderEventAny::PendingCancel(e) => encode_order_pending_cancel(e),
618 OrderEventAny::ModifyRejected(e) => encode_order_modify_rejected(e),
619 OrderEventAny::CancelRejected(e) => encode_order_cancel_rejected(e),
620 OrderEventAny::Updated(e) => encode_order_updated(e),
621 OrderEventAny::Filled(e) => Ok(retag(encode_order_filled(e)?, PAYLOAD_TYPE_ORDER_FILLED)),
622 OrderEventAny::FillVoided(e) => encode_order_fill_voided(e),
623 }
624}
625
626pub fn encode_execution_report(report: &ExecutionReport) -> Result<EncodedPayload, EncodeError> {
643 match report {
644 ExecutionReport::Order(r) => Ok(retag(
645 encode_order_status_report(r)?,
646 PAYLOAD_TYPE_ORDER_STATUS_REPORT,
647 )),
648 ExecutionReport::Fill(r) => encode_fill_report(r),
649 ExecutionReport::OrderWithFills(order, fills) => encode_order_with_fills(order, fills),
650 ExecutionReport::Position(r) => encode_position_status_report(r),
651 ExecutionReport::MassStatus(s) => encode_execution_mass_status(s),
652 }
653}
654
655pub fn encode_fill_report(report: &FillReport) -> Result<EncodedPayload, EncodeError> {
662 let payload = encode_serde(report)?;
663 let mut index_keys = Vec::with_capacity(2);
664 index_keys.push(IndexKey::new(
665 IndexKind::VenueOrderId,
666 report.venue_order_id.to_string(),
667 ));
668
669 if let Some(client_order_id) = &report.client_order_id {
670 index_keys.push(IndexKey::new(
671 IndexKind::ClientOrderId,
672 client_order_id.to_string(),
673 ));
674 }
675
676 Ok(EncodedPayload::with_payload_type(
677 payload_type(PAYLOAD_TYPE_FILL_REPORT),
678 payload,
679 index_keys,
680 ))
681}
682
683pub fn encode_position_status_report(
694 report: &PositionStatusReport,
695) -> Result<EncodedPayload, EncodeError> {
696 let payload = encode_serde(report)?;
697 Ok(EncodedPayload::with_payload_type(
698 payload_type(PAYLOAD_TYPE_POSITION_STATUS_REPORT),
699 payload,
700 Vec::new(),
701 ))
702}
703
704fn encode_order_with_fills(
705 order: &OrderStatusReport,
706 fills: &[FillReport],
707) -> Result<EncodedPayload, EncodeError> {
708 #[derive(Serialize)]
709 struct OrderWithFillsRef<'a> {
710 order_report: &'a OrderStatusReport,
711 fill_reports: &'a [FillReport],
712 }
713
714 let payload = encode_serde(&OrderWithFillsRef {
715 order_report: order,
716 fill_reports: fills,
717 })?;
718 let mut index_keys = Vec::new();
719 let mut seen = HashSet::new();
720 push_unique_index_key(
721 &mut index_keys,
722 &mut seen,
723 IndexKind::VenueOrderId,
724 order.venue_order_id.to_string(),
725 );
726
727 if let Some(client_order_id) = &order.client_order_id {
728 push_unique_index_key(
729 &mut index_keys,
730 &mut seen,
731 IndexKind::ClientOrderId,
732 client_order_id.to_string(),
733 );
734 }
735
736 for fill in fills {
737 push_unique_index_key(
738 &mut index_keys,
739 &mut seen,
740 IndexKind::VenueOrderId,
741 fill.venue_order_id.to_string(),
742 );
743
744 if let Some(client_order_id) = &fill.client_order_id {
745 push_unique_index_key(
746 &mut index_keys,
747 &mut seen,
748 IndexKind::ClientOrderId,
749 client_order_id.to_string(),
750 );
751 }
752 }
753
754 Ok(EncodedPayload::with_payload_type(
755 payload_type(PAYLOAD_TYPE_ORDER_WITH_FILLS),
756 payload,
757 index_keys,
758 ))
759}
760
761fn encode_execution_mass_status(
762 status: &ExecutionMassStatus,
763) -> Result<EncodedPayload, EncodeError> {
764 let payload = encode_serde(status)?;
765 let mut index_keys = Vec::new();
766 let mut seen = HashSet::new();
767
768 let order_reports = status.order_reports();
769
770 for (venue_order_id, report) in &order_reports {
771 push_unique_index_key(
772 &mut index_keys,
773 &mut seen,
774 IndexKind::VenueOrderId,
775 venue_order_id.to_string(),
776 );
777
778 if let Some(client_order_id) = &report.client_order_id {
779 push_unique_index_key(
780 &mut index_keys,
781 &mut seen,
782 IndexKind::ClientOrderId,
783 client_order_id.to_string(),
784 );
785 }
786 }
787
788 let fill_reports = status.fill_reports();
789
790 for (venue_order_id, fills) in &fill_reports {
791 push_unique_index_key(
792 &mut index_keys,
793 &mut seen,
794 IndexKind::VenueOrderId,
795 venue_order_id.to_string(),
796 );
797
798 for fill in fills {
799 if let Some(client_order_id) = &fill.client_order_id {
800 push_unique_index_key(
801 &mut index_keys,
802 &mut seen,
803 IndexKind::ClientOrderId,
804 client_order_id.to_string(),
805 );
806 }
807 }
808 }
809 Ok(EncodedPayload::with_payload_type(
812 payload_type(PAYLOAD_TYPE_EXECUTION_MASS_STATUS),
813 payload,
814 index_keys,
815 ))
816}
817
818fn push_unique_index_key(
819 index_keys: &mut Vec<IndexKey>,
820 seen: &mut HashSet<(IndexKind, String)>,
821 kind: IndexKind,
822 key: String,
823) {
824 if seen.insert((kind, key.clone())) {
825 index_keys.push(IndexKey::new(kind, key));
826 }
827}
828
829pub fn encode_position_event(event: &PositionEvent) -> Result<EncodedPayload, EncodeError> {
842 match event {
843 PositionEvent::PositionOpened(e) => encode_position_opened(e),
844 PositionEvent::PositionChanged(e) => encode_position_changed(e),
845 PositionEvent::PositionClosed(e) => encode_position_closed(e),
846 PositionEvent::PositionAdjusted(e) => encode_position_adjusted(e),
847 }
848}
849
850fn encode_position_opened(event: &PositionOpened) -> Result<EncodedPayload, EncodeError> {
851 let payload = encode_serde(event)?;
852 let index_keys = vec![IndexKey::new(
853 IndexKind::ClientOrderId,
854 event.opening_order_id.to_string(),
855 )];
856 Ok(EncodedPayload::with_payload_type(
857 payload_type(PAYLOAD_TYPE_POSITION_OPENED),
858 payload,
859 index_keys,
860 ))
861}
862
863fn encode_position_changed(event: &PositionChanged) -> Result<EncodedPayload, EncodeError> {
864 let payload = encode_serde(event)?;
865 let index_keys = vec![IndexKey::new(
866 IndexKind::ClientOrderId,
867 event.opening_order_id.to_string(),
868 )];
869 Ok(EncodedPayload::with_payload_type(
870 payload_type(PAYLOAD_TYPE_POSITION_CHANGED),
871 payload,
872 index_keys,
873 ))
874}
875
876fn encode_position_closed(event: &PositionClosed) -> Result<EncodedPayload, EncodeError> {
877 let payload = encode_serde(event)?;
878 let mut index_keys = Vec::new();
879 let mut seen = HashSet::new();
880
881 push_unique_index_key(
882 &mut index_keys,
883 &mut seen,
884 IndexKind::ClientOrderId,
885 event.opening_order_id.to_string(),
886 );
887
888 if let Some(closing_order_id) = &event.closing_order_id {
891 push_unique_index_key(
892 &mut index_keys,
893 &mut seen,
894 IndexKind::ClientOrderId,
895 closing_order_id.to_string(),
896 );
897 }
898
899 Ok(EncodedPayload::with_payload_type(
900 payload_type(PAYLOAD_TYPE_POSITION_CLOSED),
901 payload,
902 index_keys,
903 ))
904}
905
906fn encode_position_adjusted(event: &PositionAdjusted) -> Result<EncodedPayload, EncodeError> {
907 let payload = encode_serde(event)?;
910 Ok(EncodedPayload::with_payload_type(
911 payload_type(PAYLOAD_TYPE_POSITION_ADJUSTED),
912 payload,
913 Vec::new(),
914 ))
915}
916
917fn encode_submit_order_list(cmd: &SubmitOrderList) -> Result<EncodedPayload, EncodeError> {
918 let payload = encode_serde(cmd)?;
919 let index_keys = cmd
920 .order_list
921 .client_order_ids
922 .iter()
923 .map(|cid| IndexKey::new(IndexKind::ClientOrderId, cid.to_string()))
924 .collect();
925 Ok(EncodedPayload::with_payload_type(
926 payload_type(PAYLOAD_TYPE_SUBMIT_ORDER_LIST),
927 payload,
928 index_keys,
929 ))
930}
931
932fn encode_modify_order(cmd: &ModifyOrder) -> Result<EncodedPayload, EncodeError> {
933 encode_with_order_ids(
934 cmd,
935 PAYLOAD_TYPE_MODIFY_ORDER,
936 cmd.client_order_id.to_string(),
937 cmd.venue_order_id.map(|v| v.to_string()),
938 )
939}
940
941fn encode_batch_modify_orders(cmd: &BatchModifyOrders) -> Result<EncodedPayload, EncodeError> {
942 let payload = encode_serde(cmd)?;
943 let mut index_keys = Vec::with_capacity(cmd.modifies.len() * 2);
944 for c in &cmd.modifies {
945 index_keys.push(IndexKey::new(
946 IndexKind::ClientOrderId,
947 c.client_order_id.to_string(),
948 ));
949
950 if let Some(venue) = c.venue_order_id {
951 index_keys.push(IndexKey::new(IndexKind::VenueOrderId, venue.to_string()));
952 }
953 }
954 Ok(EncodedPayload::with_payload_type(
955 payload_type(PAYLOAD_TYPE_BATCH_MODIFY_ORDERS),
956 payload,
957 index_keys,
958 ))
959}
960
961fn encode_cancel_order(cmd: &CancelOrder) -> Result<EncodedPayload, EncodeError> {
962 encode_with_order_ids(
963 cmd,
964 PAYLOAD_TYPE_CANCEL_ORDER,
965 cmd.client_order_id.to_string(),
966 cmd.venue_order_id.map(|v| v.to_string()),
967 )
968}
969
970fn encode_cancel_all_orders(cmd: &CancelAllOrders) -> Result<EncodedPayload, EncodeError> {
971 let payload = encode_serde(cmd)?;
972 Ok(EncodedPayload::with_payload_type(
973 payload_type(PAYLOAD_TYPE_CANCEL_ALL_ORDERS),
974 payload,
975 Vec::new(),
976 ))
977}
978
979fn encode_batch_cancel_orders(cmd: &BatchCancelOrders) -> Result<EncodedPayload, EncodeError> {
980 let payload = encode_serde(cmd)?;
981 let mut index_keys = Vec::with_capacity(cmd.cancels.len() * 2);
982 for c in &cmd.cancels {
983 index_keys.push(IndexKey::new(
984 IndexKind::ClientOrderId,
985 c.client_order_id.to_string(),
986 ));
987
988 if let Some(venue) = c.venue_order_id {
989 index_keys.push(IndexKey::new(IndexKind::VenueOrderId, venue.to_string()));
990 }
991 }
992 Ok(EncodedPayload::with_payload_type(
993 payload_type(PAYLOAD_TYPE_BATCH_CANCEL_ORDERS),
994 payload,
995 index_keys,
996 ))
997}
998
999fn encode_query_order(cmd: &QueryOrder) -> Result<EncodedPayload, EncodeError> {
1000 encode_with_order_ids(
1001 cmd,
1002 PAYLOAD_TYPE_QUERY_ORDER,
1003 cmd.client_order_id.to_string(),
1004 cmd.venue_order_id.map(|v| v.to_string()),
1005 )
1006}
1007
1008fn encode_query_account(cmd: &QueryAccount) -> Result<EncodedPayload, EncodeError> {
1009 let payload = encode_serde(cmd)?;
1010 Ok(EncodedPayload::with_payload_type(
1011 payload_type(PAYLOAD_TYPE_QUERY_ACCOUNT),
1012 payload,
1013 Vec::new(),
1014 ))
1015}
1016
1017fn encode_order_initialized(e: &OrderInitialized) -> Result<EncodedPayload, EncodeError> {
1018 encode_with_order_ids(
1019 e,
1020 PAYLOAD_TYPE_ORDER_INITIALIZED,
1021 e.client_order_id.to_string(),
1022 None,
1023 )
1024}
1025
1026fn encode_order_denied(e: &OrderDenied) -> Result<EncodedPayload, EncodeError> {
1027 encode_with_order_ids(
1028 e,
1029 PAYLOAD_TYPE_ORDER_DENIED,
1030 e.client_order_id.to_string(),
1031 None,
1032 )
1033}
1034
1035fn encode_order_emulated(e: &OrderEmulated) -> Result<EncodedPayload, EncodeError> {
1036 encode_with_order_ids(
1037 e,
1038 PAYLOAD_TYPE_ORDER_EMULATED,
1039 e.client_order_id.to_string(),
1040 None,
1041 )
1042}
1043
1044fn encode_order_released(e: &OrderReleased) -> Result<EncodedPayload, EncodeError> {
1045 encode_with_order_ids(
1046 e,
1047 PAYLOAD_TYPE_ORDER_RELEASED,
1048 e.client_order_id.to_string(),
1049 None,
1050 )
1051}
1052
1053fn encode_order_submitted(e: &OrderSubmitted) -> Result<EncodedPayload, EncodeError> {
1054 encode_with_order_ids(
1055 e,
1056 PAYLOAD_TYPE_ORDER_SUBMITTED,
1057 e.client_order_id.to_string(),
1058 None,
1059 )
1060}
1061
1062fn encode_order_accepted(e: &OrderAccepted) -> Result<EncodedPayload, EncodeError> {
1063 encode_with_order_ids(
1064 e,
1065 PAYLOAD_TYPE_ORDER_ACCEPTED,
1066 e.client_order_id.to_string(),
1067 Some(e.venue_order_id.to_string()),
1068 )
1069}
1070
1071fn encode_order_rejected(e: &OrderRejected) -> Result<EncodedPayload, EncodeError> {
1072 encode_with_order_ids(
1073 e,
1074 PAYLOAD_TYPE_ORDER_REJECTED,
1075 e.client_order_id.to_string(),
1076 None,
1077 )
1078}
1079
1080fn encode_order_canceled(e: &OrderCanceled) -> Result<EncodedPayload, EncodeError> {
1081 encode_with_order_ids(
1082 e,
1083 PAYLOAD_TYPE_ORDER_CANCELED,
1084 e.client_order_id.to_string(),
1085 e.venue_order_id.map(|v| v.to_string()),
1086 )
1087}
1088
1089fn encode_order_expired(e: &OrderExpired) -> Result<EncodedPayload, EncodeError> {
1090 encode_with_order_ids(
1091 e,
1092 PAYLOAD_TYPE_ORDER_EXPIRED,
1093 e.client_order_id.to_string(),
1094 e.venue_order_id.map(|v| v.to_string()),
1095 )
1096}
1097
1098fn encode_order_triggered(e: &OrderTriggered) -> Result<EncodedPayload, EncodeError> {
1099 encode_with_order_ids(
1100 e,
1101 PAYLOAD_TYPE_ORDER_TRIGGERED,
1102 e.client_order_id.to_string(),
1103 e.venue_order_id.map(|v| v.to_string()),
1104 )
1105}
1106
1107fn encode_order_pending_update(e: &OrderPendingUpdate) -> Result<EncodedPayload, EncodeError> {
1108 encode_with_order_ids(
1109 e,
1110 PAYLOAD_TYPE_ORDER_PENDING_UPDATE,
1111 e.client_order_id.to_string(),
1112 e.venue_order_id.map(|v| v.to_string()),
1113 )
1114}
1115
1116fn encode_order_pending_cancel(e: &OrderPendingCancel) -> Result<EncodedPayload, EncodeError> {
1117 encode_with_order_ids(
1118 e,
1119 PAYLOAD_TYPE_ORDER_PENDING_CANCEL,
1120 e.client_order_id.to_string(),
1121 e.venue_order_id.map(|v| v.to_string()),
1122 )
1123}
1124
1125fn encode_order_modify_rejected(e: &OrderModifyRejected) -> Result<EncodedPayload, EncodeError> {
1126 encode_with_order_ids(
1127 e,
1128 PAYLOAD_TYPE_ORDER_MODIFY_REJECTED,
1129 e.client_order_id.to_string(),
1130 e.venue_order_id.map(|v| v.to_string()),
1131 )
1132}
1133
1134fn encode_order_cancel_rejected(e: &OrderCancelRejected) -> Result<EncodedPayload, EncodeError> {
1135 encode_with_order_ids(
1136 e,
1137 PAYLOAD_TYPE_ORDER_CANCEL_REJECTED,
1138 e.client_order_id.to_string(),
1139 e.venue_order_id.map(|v| v.to_string()),
1140 )
1141}
1142
1143fn encode_order_updated(e: &OrderUpdated) -> Result<EncodedPayload, EncodeError> {
1144 encode_with_order_ids(
1145 e,
1146 PAYLOAD_TYPE_ORDER_UPDATED,
1147 e.client_order_id.to_string(),
1148 e.venue_order_id.map(|v| v.to_string()),
1149 )
1150}
1151
1152fn encode_order_fill_voided(e: &OrderFillVoided) -> Result<EncodedPayload, EncodeError> {
1153 encode_with_order_ids(
1154 e,
1155 PAYLOAD_TYPE_ORDER_FILL_VOIDED,
1156 e.client_order_id.to_string(),
1157 Some(e.venue_order_id.to_string()),
1158 )
1159}
1160
1161fn encode_with_order_ids<T: Serialize>(
1162 value: &T,
1163 tag: &str,
1164 client_order_id: String,
1165 venue_order_id: Option<String>,
1166) -> Result<EncodedPayload, EncodeError> {
1167 let payload = encode_serde(value)?;
1168 let mut index_keys = Vec::with_capacity(2);
1169 index_keys.push(IndexKey::new(IndexKind::ClientOrderId, client_order_id));
1170 if let Some(venue) = venue_order_id {
1171 index_keys.push(IndexKey::new(IndexKind::VenueOrderId, venue));
1172 }
1173 Ok(EncodedPayload::with_payload_type(
1174 payload_type(tag),
1175 payload,
1176 index_keys,
1177 ))
1178}
1179
1180fn retag(mut encoded: EncodedPayload, tag: &str) -> EncodedPayload {
1181 encoded.payload_type = Some(payload_type(tag));
1182 encoded
1183}
1184
1185pub fn encode_order_status_report(
1195 message: &OrderStatusReport,
1196) -> Result<EncodedPayload, EncodeError> {
1197 let payload = encode_serde(message)?;
1198 let mut index_keys = Vec::with_capacity(2);
1199 index_keys.push(IndexKey::new(
1200 IndexKind::VenueOrderId,
1201 message.venue_order_id.to_string(),
1202 ));
1203
1204 if let Some(client_order_id) = &message.client_order_id {
1205 index_keys.push(IndexKey::new(
1206 IndexKind::ClientOrderId,
1207 client_order_id.to_string(),
1208 ));
1209 }
1210 Ok(EncodedPayload::new(payload, index_keys))
1211}
1212
1213pub fn encode_account_state(message: &AccountState) -> Result<EncodedPayload, EncodeError> {
1224 let payload = encode_serde(message)?;
1225 Ok(EncodedPayload::new(payload, Vec::new()))
1226}
1227
1228#[derive(Serialize)]
1229struct TimeEventPayload<'a> {
1230 name: &'a str,
1231 event_id: UUID4,
1232 ts_event: UnixNanos,
1233 ts_init: UnixNanos,
1234}
1235
1236pub fn encode_time_event(event: &TimeEvent) -> Result<EncodedPayload, EncodeError> {
1246 let payload = TimeEventPayload {
1247 name: event.name.as_str(),
1248 event_id: event.event_id,
1249 ts_event: event.ts_event,
1250 ts_init: event.ts_init,
1251 };
1252 Ok(EncodedPayload::new(encode_serde(&payload)?, Vec::new()))
1253}
1254
1255pub fn encode_data_command(command: &DataCommand) -> Result<EncodedPayload, EncodeError> {
1272 match command {
1273 DataCommand::Request(cmd) => {
1274 encode_data_command_category(cmd, PAYLOAD_TYPE_REQUEST_COMMAND)
1275 }
1276 DataCommand::Subscribe(cmd) => {
1277 encode_data_command_category(cmd, PAYLOAD_TYPE_SUBSCRIBE_COMMAND)
1278 }
1279 DataCommand::Unsubscribe(cmd) => {
1280 encode_data_command_category(cmd, PAYLOAD_TYPE_UNSUBSCRIBE_COMMAND)
1281 }
1282 #[cfg(feature = "defi")]
1283 DataCommand::DefiRequest(cmd) => {
1284 encode_data_command_category(cmd, PAYLOAD_TYPE_DEFI_REQUEST_COMMAND)
1285 }
1286 #[cfg(feature = "defi")]
1287 DataCommand::DefiSubscribe(cmd) => {
1288 encode_data_command_category(cmd, PAYLOAD_TYPE_DEFI_SUBSCRIBE_COMMAND)
1289 }
1290 #[cfg(feature = "defi")]
1291 DataCommand::DefiUnsubscribe(cmd) => {
1292 encode_data_command_category(cmd, PAYLOAD_TYPE_DEFI_UNSUBSCRIBE_COMMAND)
1293 }
1294 _ => Err(EncodeError::Serialize(
1295 "unsupported DataCommand variant".to_string(),
1296 )),
1297 }
1298}
1299
1300fn encode_data_command_category<T: Serialize>(
1301 command: &T,
1302 tag: &str,
1303) -> Result<EncodedPayload, EncodeError> {
1304 let payload = encode_serde(command)?;
1305 Ok(EncodedPayload::with_payload_type(
1306 payload_type(tag),
1307 payload,
1308 Vec::new(),
1309 ))
1310}
1311
1312pub fn encode_data_response(response: &DataResponse) -> Result<EncodedPayload, EncodeError> {
1340 match response {
1341 DataResponse::Data(resp) => encode_custom_data_response(resp),
1342 DataResponse::Instrument(resp) => encode_instrument_response(resp),
1343 DataResponse::Instruments(resp) => encode_instruments_response(resp),
1344 DataResponse::Book(resp) => encode_book_response(resp),
1345 DataResponse::BookDeltas(resp) => encode_book_deltas_response(resp),
1346 DataResponse::BookDepth(resp) => encode_book_depth_response(resp),
1347 DataResponse::Quotes(resp) => encode_quotes_response(resp),
1348 DataResponse::Trades(resp) => encode_trades_response(resp),
1349 DataResponse::FundingRates(resp) => encode_funding_rates_response(resp),
1350 DataResponse::ForwardPrices(resp) => encode_forward_prices_response(resp),
1351 DataResponse::Bars(resp) => encode_bars_response(resp),
1352 }
1353}
1354
1355fn encode_custom_data_response(
1356 response: &CustomDataResponse,
1357) -> Result<EncodedPayload, EncodeError> {
1358 #[derive(Serialize)]
1363 struct CustomDataResponseRef<'a> {
1364 correlation_id: &'a UUID4,
1365 client_id: &'a ClientId,
1366 venue: &'a Option<Venue>,
1367 data_type: &'a DataType,
1368 start: &'a Option<UnixNanos>,
1369 end: &'a Option<UnixNanos>,
1370 ts_init: &'a UnixNanos,
1371 params: &'a Option<Params>,
1372 }
1373
1374 let payload = encode_serde(&CustomDataResponseRef {
1375 correlation_id: &response.correlation_id,
1376 client_id: &response.client_id,
1377 venue: &response.venue,
1378 data_type: &response.data_type,
1379 start: &response.start,
1380 end: &response.end,
1381 ts_init: &response.ts_init,
1382 params: &response.params,
1383 })?;
1384 Ok(EncodedPayload::with_payload_type(
1385 payload_type(PAYLOAD_TYPE_CUSTOM_DATA_RESPONSE),
1386 payload,
1387 Vec::new(),
1388 ))
1389}
1390
1391fn encode_instrument_response(
1392 response: &InstrumentResponse,
1393) -> Result<EncodedPayload, EncodeError> {
1394 let payload = encode_serde(response)?;
1395 Ok(EncodedPayload::with_payload_type(
1396 payload_type(PAYLOAD_TYPE_INSTRUMENT_RESPONSE),
1397 payload,
1398 Vec::new(),
1399 ))
1400}
1401
1402fn encode_instruments_response(
1403 response: &InstrumentsResponse,
1404) -> Result<EncodedPayload, EncodeError> {
1405 let payload = encode_serde(response)?;
1406 Ok(EncodedPayload::with_payload_type(
1407 payload_type(PAYLOAD_TYPE_INSTRUMENTS_RESPONSE),
1408 payload,
1409 Vec::new(),
1410 ))
1411}
1412
1413fn encode_book_response(response: &BookResponse) -> Result<EncodedPayload, EncodeError> {
1414 #[derive(Serialize)]
1419 struct BookResponseRef<'a> {
1420 correlation_id: &'a UUID4,
1421 client_id: &'a ClientId,
1422 instrument_id: &'a InstrumentId,
1423 start: &'a Option<UnixNanos>,
1424 end: &'a Option<UnixNanos>,
1425 ts_init: &'a UnixNanos,
1426 params: &'a Option<Params>,
1427 }
1428
1429 let payload = encode_serde(&BookResponseRef {
1430 correlation_id: &response.correlation_id,
1431 client_id: &response.client_id,
1432 instrument_id: &response.instrument_id,
1433 start: &response.start,
1434 end: &response.end,
1435 ts_init: &response.ts_init,
1436 params: &response.params,
1437 })?;
1438 Ok(EncodedPayload::with_payload_type(
1439 payload_type(PAYLOAD_TYPE_BOOK_RESPONSE),
1440 payload,
1441 Vec::new(),
1442 ))
1443}
1444
1445fn encode_quotes_response(response: &QuotesResponse) -> Result<EncodedPayload, EncodeError> {
1446 let payload = encode_serde(response)?;
1447 Ok(EncodedPayload::with_payload_type(
1448 payload_type(PAYLOAD_TYPE_QUOTES_RESPONSE),
1449 payload,
1450 Vec::new(),
1451 ))
1452}
1453
1454fn encode_book_deltas_response(
1455 response: &BookDeltasResponse,
1456) -> Result<EncodedPayload, EncodeError> {
1457 let payload = encode_serde(response)?;
1458 Ok(EncodedPayload::with_payload_type(
1459 payload_type(PAYLOAD_TYPE_BOOK_DELTAS_RESPONSE),
1460 payload,
1461 Vec::new(),
1462 ))
1463}
1464
1465fn encode_book_depth_response(response: &BookDepthResponse) -> Result<EncodedPayload, EncodeError> {
1466 let payload = encode_serde(response)?;
1467 Ok(EncodedPayload::with_payload_type(
1468 payload_type(PAYLOAD_TYPE_BOOK_DEPTH_RESPONSE),
1469 payload,
1470 Vec::new(),
1471 ))
1472}
1473
1474fn encode_trades_response(response: &TradesResponse) -> Result<EncodedPayload, EncodeError> {
1475 let payload = encode_serde(response)?;
1476 Ok(EncodedPayload::with_payload_type(
1477 payload_type(PAYLOAD_TYPE_TRADES_RESPONSE),
1478 payload,
1479 Vec::new(),
1480 ))
1481}
1482
1483fn encode_funding_rates_response(
1484 response: &FundingRatesResponse,
1485) -> Result<EncodedPayload, EncodeError> {
1486 let payload = encode_serde(response)?;
1487 Ok(EncodedPayload::with_payload_type(
1488 payload_type(PAYLOAD_TYPE_FUNDING_RATES_RESPONSE),
1489 payload,
1490 Vec::new(),
1491 ))
1492}
1493
1494fn encode_forward_prices_response(
1495 response: &ForwardPricesResponse,
1496) -> Result<EncodedPayload, EncodeError> {
1497 let payload = encode_serde(response)?;
1498 Ok(EncodedPayload::with_payload_type(
1499 payload_type(PAYLOAD_TYPE_FORWARD_PRICES_RESPONSE),
1500 payload,
1501 Vec::new(),
1502 ))
1503}
1504
1505fn encode_bars_response(response: &BarsResponse) -> Result<EncodedPayload, EncodeError> {
1506 let payload = encode_serde(response)?;
1507 Ok(EncodedPayload::with_payload_type(
1508 payload_type(PAYLOAD_TYPE_BARS_RESPONSE),
1509 payload,
1510 Vec::new(),
1511 ))
1512}
1513
1514#[cfg(test)]
1515mod tests {
1516 use nautilus_common::messages::data::{
1517 RequestCommand, RequestQuotes, SubscribeCommand, SubscribeQuotes, UnsubscribeCommand,
1518 UnsubscribeQuotes,
1519 };
1520 #[cfg(feature = "defi")]
1521 use nautilus_common::messages::defi::{
1522 DefiRequestCommand, DefiSubscribeCommand, DefiUnsubscribeCommand, RequestPoolSnapshot,
1523 SubscribeBlocks, UnsubscribeBlocks,
1524 };
1525 use nautilus_core::{UUID4, UnixNanos};
1526 #[cfg(feature = "defi")]
1527 use nautilus_model::defi::Blockchain;
1528 use nautilus_model::{
1529 data::{Bar, BarType, stubs::stub_depth10},
1530 enums::{
1531 AccountType, BookType, LiquiditySide, OrderSide, OrderStatus, OrderType,
1532 PositionAdjustmentType, PositionSide, TimeInForce,
1533 },
1534 events::{
1535 PositionAdjusted, PositionChanged, PositionClosed, PositionOpened,
1536 order::spec::{
1537 OrderFillVoidedSpec, OrderFilledSpec, OrderInitializedSpec, OrderSubmittedSpec,
1538 },
1539 },
1540 identifiers::{
1541 AccountId, ClientId, ClientOrderId, InstrumentId, OrderListId, PositionId, StrategyId,
1542 TradeId, TraderId, Venue, VenueOrderId,
1543 },
1544 instruments::{InstrumentAny, stubs::currency_pair_ethusdt},
1545 orderbook::OrderBook,
1546 orders::OrderList,
1547 reports::{ExecutionMassStatus, FillReport, PositionStatusReport},
1548 types::{AccountBalance, Currency, Money, Price, Quantity},
1549 };
1550 use rstest::rstest;
1551 use serde::Deserialize;
1552
1553 use super::*;
1554
1555 fn trader_id() -> TraderId {
1556 TraderId::from("TRADER-001")
1557 }
1558
1559 fn strategy_id() -> StrategyId {
1560 StrategyId::from("S-001")
1561 }
1562
1563 fn instrument_id() -> InstrumentId {
1564 InstrumentId::from("ETHUSDT-PERP.BINANCE")
1565 }
1566
1567 fn client_order_id() -> ClientOrderId {
1568 ClientOrderId::from("O-20260510-000001")
1569 }
1570
1571 fn venue_order_id() -> VenueOrderId {
1572 VenueOrderId::from("V-12345")
1573 }
1574
1575 fn make_submit_order() -> SubmitOrder {
1576 let order_init = OrderInitializedSpec::builder()
1577 .instrument_id(instrument_id())
1578 .client_order_id(client_order_id())
1579 .quantity(Quantity::from("1"))
1580 .time_in_force(TimeInForce::Gtc)
1581 .ts_event(UnixNanos::from(1))
1582 .ts_init(UnixNanos::from(2))
1583 .build();
1584 SubmitOrder::new(
1585 trader_id(),
1586 Some(ClientId::from("BINANCE")),
1587 strategy_id(),
1588 instrument_id(),
1589 client_order_id(),
1590 order_init,
1591 None,
1592 None,
1593 None,
1594 UUID4::new(),
1595 UnixNanos::from(3),
1596 None, )
1598 }
1599
1600 fn make_order_filled() -> OrderFilled {
1601 OrderFilledSpec::builder()
1602 .instrument_id(instrument_id())
1603 .client_order_id(client_order_id())
1604 .venue_order_id(venue_order_id())
1605 .account_id(AccountId::from("BINANCE-001"))
1606 .trade_id(TradeId::from("T-9999"))
1607 .last_qty(Quantity::from("1"))
1608 .last_px(Price::from("100.00"))
1609 .currency(Currency::USDT())
1610 .ts_event(UnixNanos::from(10))
1611 .ts_init(UnixNanos::from(11))
1612 .commission(Money::new(0.10, Currency::USDT()))
1613 .build()
1614 }
1615
1616 fn make_order_status_report() -> OrderStatusReport {
1617 OrderStatusReport::new(
1618 AccountId::from("BINANCE-001"),
1619 instrument_id(),
1620 Some(client_order_id()),
1621 venue_order_id(),
1622 OrderSide::Buy.into(),
1623 OrderType::Market,
1624 TimeInForce::Gtc,
1625 OrderStatus::Filled,
1626 Quantity::from("1"),
1627 Quantity::from("1"),
1628 UnixNanos::from(20),
1629 UnixNanos::from(21),
1630 UnixNanos::from(22),
1631 Some(UUID4::new()),
1632 )
1633 }
1634
1635 #[rstest]
1636 fn submit_order_encoder_emits_client_order_id_index() {
1637 let cmd = make_submit_order();
1638 let encoded = encode_submit_order(&cmd).expect("encode");
1639
1640 assert!(!encoded.payload.is_empty());
1641 assert_eq!(encoded.index_keys.len(), 1);
1642 assert_eq!(encoded.index_keys[0].kind, IndexKind::ClientOrderId);
1643 assert_eq!(encoded.index_keys[0].key, cmd.client_order_id.to_string());
1644 }
1645
1646 #[rstest]
1647 fn order_filled_encoder_emits_client_and_venue_order_id_indices() {
1648 let event = make_order_filled();
1649 let encoded = encode_order_filled(&event).expect("encode");
1650
1651 assert!(!encoded.payload.is_empty());
1652 assert_eq!(encoded.index_keys.len(), 2);
1653 assert_eq!(encoded.index_keys[0].kind, IndexKind::ClientOrderId);
1654 assert_eq!(encoded.index_keys[0].key, event.client_order_id.to_string());
1655 assert_eq!(encoded.index_keys[1].kind, IndexKind::VenueOrderId);
1656 assert_eq!(encoded.index_keys[1].key, event.venue_order_id.to_string());
1657 }
1658
1659 #[rstest]
1660 fn order_status_report_encoder_includes_client_order_id_when_present() {
1661 let report = make_order_status_report();
1662 let encoded = encode_order_status_report(&report).expect("encode");
1663
1664 assert_eq!(encoded.index_keys.len(), 2);
1665 assert_eq!(encoded.index_keys[0].kind, IndexKind::VenueOrderId);
1666 assert_eq!(encoded.index_keys[1].kind, IndexKind::ClientOrderId);
1667 }
1668
1669 #[rstest]
1670 fn order_status_report_encoder_omits_client_order_id_when_absent() {
1671 let mut report = make_order_status_report();
1672 report.client_order_id = None;
1673 let encoded = encode_order_status_report(&report).expect("encode");
1674
1675 assert_eq!(encoded.index_keys.len(), 1);
1676 assert_eq!(encoded.index_keys[0].kind, IndexKind::VenueOrderId);
1677 }
1678
1679 #[rstest]
1680 fn default_registry_covers_published_state_affecting_surface() {
1681 let registry = default_registry();
1682 let expected = [
1683 (
1684 "send_any_value(SubmitOrder) / bare SubmitOrder",
1685 registry.contains::<SubmitOrder>(),
1686 ),
1687 (
1688 "publish_order_event(OrderFilled) / bare OrderFilled",
1689 registry.contains::<OrderFilled>(),
1690 ),
1691 (
1692 "reconciliation.raw.order_status / OrderStatusReport",
1693 registry.contains::<OrderStatusReport>(),
1694 ),
1695 (
1696 "reconciliation.raw.fill / FillReport",
1697 registry.contains::<FillReport>(),
1698 ),
1699 (
1700 "reconciliation.raw.position / PositionStatusReport",
1701 registry.contains::<PositionStatusReport>(),
1702 ),
1703 (
1704 "send_trading_command / TradingCommand",
1705 registry.contains::<TradingCommand>(),
1706 ),
1707 (
1708 "publish_order_event / OrderEventAny",
1709 registry.contains::<OrderEventAny>(),
1710 ),
1711 (
1712 "send_execution_report / ExecutionReport",
1713 registry.contains::<ExecutionReport>(),
1714 ),
1715 (
1716 "publish_position_event / PositionEvent",
1717 registry.contains::<PositionEvent>(),
1718 ),
1719 (
1720 "publish_account_state and send_account_state / AccountState",
1721 registry.contains::<AccountState>(),
1722 ),
1723 (
1724 "time event handler firing / TimeEvent",
1725 registry.contains::<TimeEvent>(),
1726 ),
1727 (
1728 "send_data_command / DataCommand",
1729 registry.contains::<DataCommand>(),
1730 ),
1731 (
1732 "send_data_response / DataResponse",
1733 registry.contains::<DataResponse>(),
1734 ),
1735 ];
1736 let missing: Vec<&str> = expected
1737 .iter()
1738 .filter_map(|(name, registered)| (!*registered).then_some(*name))
1739 .collect();
1740
1741 assert!(
1742 missing.is_empty(),
1743 "missing default event-store encoder registrations for {missing:?}",
1744 );
1745 assert_eq!(
1746 registry.len(),
1747 expected.len(),
1748 "default registry must match the audited state-affecting surface",
1749 );
1750 }
1751
1752 #[rstest]
1753 fn submit_order_payload_round_trips_through_msgpack() {
1754 let cmd = make_submit_order();
1755 let encoded = encode_submit_order(&cmd).expect("encode");
1756
1757 let decoded: SubmitOrder = rmp_serde::from_slice(&encoded.payload).expect("decode");
1758 assert_eq!(decoded, cmd);
1759 }
1760
1761 #[rstest]
1762 fn default_registry_data_command_identity_dedupes_dispatch_hops() {
1763 let registry = default_registry();
1767
1768 let request = make_request_command();
1769 let expected_request = *request.request_id();
1770 let subscribe = make_subscribe_command();
1771 let expected_subscribe = subscribe.command_id();
1772 let unsubscribe = make_unsubscribe_command();
1773 let expected_unsubscribe = unsubscribe.command_id();
1774
1775 let cases = [
1776 (DataCommand::Request(request), expected_request),
1777 (DataCommand::Subscribe(subscribe), expected_subscribe),
1778 (DataCommand::Unsubscribe(unsubscribe), expected_unsubscribe),
1779 ];
1780
1781 for (command, expected) in cases {
1782 assert_eq!(registry.identity_for_any(&command), Some(expected));
1783 }
1784 }
1785
1786 #[cfg(feature = "defi")]
1787 #[rstest]
1788 fn default_registry_defi_data_command_identity_dedupes_dispatch_hops() {
1789 let registry = default_registry();
1790
1791 let request = make_defi_request_command();
1792 let expected_request = *request.request_id();
1793 let subscribe = make_defi_subscribe_command();
1794 let expected_subscribe = subscribe.command_id();
1795 let unsubscribe = make_defi_unsubscribe_command();
1796 let expected_unsubscribe = unsubscribe.command_id();
1797
1798 let cases = [
1799 (DataCommand::DefiRequest(request), expected_request),
1800 (DataCommand::DefiSubscribe(subscribe), expected_subscribe),
1801 (
1802 DataCommand::DefiUnsubscribe(unsubscribe),
1803 expected_unsubscribe,
1804 ),
1805 ];
1806
1807 for (command, expected) in cases {
1808 assert_eq!(registry.identity_for_any(&command), Some(expected));
1809 }
1810 }
1811
1812 #[rstest]
1813 fn order_filled_payload_round_trips_through_msgpack() {
1814 let event = make_order_filled();
1815 let encoded = encode_order_filled(&event).expect("encode");
1816
1817 let decoded: OrderFilled = rmp_serde::from_slice(&encoded.payload).expect("decode");
1818 assert_eq!(decoded, event);
1819 }
1820
1821 #[rstest]
1822 fn order_status_report_payload_round_trips_through_msgpack() {
1823 let report = make_order_status_report();
1824 let encoded = encode_order_status_report(&report).expect("encode");
1825
1826 let decoded: OrderStatusReport = rmp_serde::from_slice(&encoded.payload).expect("decode");
1827 assert_eq!(decoded, report);
1828 }
1829
1830 fn make_cancel_order() -> CancelOrder {
1831 CancelOrder::new(
1832 trader_id(),
1833 Some(ClientId::from("BINANCE")),
1834 strategy_id(),
1835 instrument_id(),
1836 client_order_id(),
1837 Some(venue_order_id()),
1838 UUID4::new(),
1839 UnixNanos::from(4),
1840 None,
1841 None, )
1843 }
1844
1845 fn make_query_account() -> QueryAccount {
1846 QueryAccount::new(
1847 trader_id(),
1848 Some(ClientId::from("BINANCE")),
1849 AccountId::from("BINANCE-001"),
1850 UUID4::new(),
1851 UnixNanos::from(5),
1852 None,
1853 None, )
1855 }
1856
1857 fn make_order_submitted() -> OrderSubmitted {
1858 OrderSubmittedSpec::builder()
1859 .instrument_id(instrument_id())
1860 .client_order_id(client_order_id())
1861 .account_id(AccountId::from("BINANCE-001"))
1862 .ts_event(UnixNanos::from(30))
1863 .ts_init(UnixNanos::from(31))
1864 .build()
1865 }
1866
1867 fn make_modify_order(venue: Option<VenueOrderId>) -> ModifyOrder {
1868 ModifyOrder::new(
1869 trader_id(),
1870 Some(ClientId::from("BINANCE")),
1871 strategy_id(),
1872 instrument_id(),
1873 client_order_id(),
1874 venue,
1875 Some(Quantity::from("2")),
1876 Some(Price::from("100.00")),
1877 None,
1878 UUID4::new(),
1879 UnixNanos::from(6),
1880 None,
1881 None, )
1883 }
1884
1885 fn make_batch_modify_orders(modifies: Vec<ModifyOrder>) -> BatchModifyOrders {
1886 BatchModifyOrders::new(
1887 trader_id(),
1888 Some(ClientId::from("BINANCE")),
1889 strategy_id(),
1890 instrument_id(),
1891 modifies,
1892 UUID4::new(),
1893 UnixNanos::from(7),
1894 None,
1895 None, )
1897 }
1898
1899 fn make_cancel_all_orders() -> CancelAllOrders {
1900 CancelAllOrders::new(
1901 trader_id(),
1902 Some(ClientId::from("BINANCE")),
1903 strategy_id(),
1904 instrument_id(),
1905 Some(OrderSide::Buy),
1906 UUID4::new(),
1907 UnixNanos::from(7),
1908 None,
1909 None, )
1911 }
1912
1913 fn make_query_order(venue: Option<VenueOrderId>) -> QueryOrder {
1914 QueryOrder::new(
1915 trader_id(),
1916 Some(ClientId::from("BINANCE")),
1917 strategy_id(),
1918 instrument_id(),
1919 client_order_id(),
1920 venue,
1921 UUID4::new(),
1922 UnixNanos::from(8),
1923 None,
1924 None, )
1926 }
1927
1928 fn make_batch_cancel_orders(cancels: Vec<CancelOrder>) -> BatchCancelOrders {
1929 BatchCancelOrders::new(
1930 trader_id(),
1931 Some(ClientId::from("BINANCE")),
1932 strategy_id(),
1933 instrument_id(),
1934 cancels,
1935 UUID4::new(),
1936 UnixNanos::from(9),
1937 None,
1938 None, )
1940 }
1941
1942 fn make_submit_order_list(client_order_ids: Vec<ClientOrderId>) -> SubmitOrderList {
1943 let order_inits: Vec<OrderInitialized> = client_order_ids
1946 .iter()
1947 .copied()
1948 .map(make_order_initialized_with_id)
1949 .collect();
1950 let order_list = OrderList::new(
1951 OrderListId::from("OL-1"),
1952 instrument_id(),
1953 strategy_id(),
1954 client_order_ids,
1955 UnixNanos::from(10),
1956 );
1957 SubmitOrderList::new(
1958 trader_id(),
1959 Some(ClientId::from("BINANCE")),
1960 strategy_id(),
1961 order_list,
1962 order_inits,
1963 None,
1964 None,
1965 None,
1966 UUID4::new(),
1967 UnixNanos::from(11),
1968 None, )
1970 }
1971
1972 fn make_order_initialized_with_id(client_order_id: ClientOrderId) -> OrderInitialized {
1973 OrderInitialized {
1974 client_order_id,
1975 ..OrderInitialized::default()
1976 }
1977 }
1978
1979 fn ev_initialized() -> OrderEventAny {
1980 OrderEventAny::Initialized(make_order_initialized_with_id(client_order_id()))
1981 }
1982
1983 fn ev_denied() -> OrderEventAny {
1984 OrderEventAny::Denied(OrderDenied {
1985 client_order_id: client_order_id(),
1986 ..Default::default()
1987 })
1988 }
1989
1990 fn ev_emulated() -> OrderEventAny {
1991 OrderEventAny::Emulated(OrderEmulated {
1992 client_order_id: client_order_id(),
1993 ..Default::default()
1994 })
1995 }
1996
1997 fn ev_released() -> OrderEventAny {
1998 OrderEventAny::Released(OrderReleased {
1999 client_order_id: client_order_id(),
2000 ..Default::default()
2001 })
2002 }
2003
2004 fn ev_submitted() -> OrderEventAny {
2005 OrderEventAny::Submitted(make_order_submitted())
2006 }
2007
2008 fn ev_accepted_with_venue(venue: VenueOrderId) -> OrderEventAny {
2009 OrderEventAny::Accepted(OrderAccepted {
2010 client_order_id: client_order_id(),
2011 venue_order_id: venue,
2012 ..Default::default()
2013 })
2014 }
2015
2016 fn ev_rejected() -> OrderEventAny {
2017 OrderEventAny::Rejected(OrderRejected {
2018 client_order_id: client_order_id(),
2019 ..Default::default()
2020 })
2021 }
2022
2023 fn ev_canceled(venue: Option<VenueOrderId>) -> OrderEventAny {
2024 OrderEventAny::Canceled(OrderCanceled {
2025 client_order_id: client_order_id(),
2026 venue_order_id: venue,
2027 ..Default::default()
2028 })
2029 }
2030
2031 fn ev_expired(venue: Option<VenueOrderId>) -> OrderEventAny {
2032 OrderEventAny::Expired(OrderExpired {
2033 client_order_id: client_order_id(),
2034 venue_order_id: venue,
2035 ..Default::default()
2036 })
2037 }
2038
2039 fn ev_triggered(venue: Option<VenueOrderId>) -> OrderEventAny {
2040 OrderEventAny::Triggered(OrderTriggered {
2041 client_order_id: client_order_id(),
2042 venue_order_id: venue,
2043 ..Default::default()
2044 })
2045 }
2046
2047 fn ev_pending_update(venue: Option<VenueOrderId>) -> OrderEventAny {
2048 OrderEventAny::PendingUpdate(OrderPendingUpdate {
2049 client_order_id: client_order_id(),
2050 venue_order_id: venue,
2051 ..Default::default()
2052 })
2053 }
2054
2055 fn ev_pending_cancel(venue: Option<VenueOrderId>) -> OrderEventAny {
2056 OrderEventAny::PendingCancel(OrderPendingCancel {
2057 client_order_id: client_order_id(),
2058 venue_order_id: venue,
2059 ..Default::default()
2060 })
2061 }
2062
2063 fn ev_modify_rejected(venue: Option<VenueOrderId>) -> OrderEventAny {
2064 OrderEventAny::ModifyRejected(OrderModifyRejected {
2065 client_order_id: client_order_id(),
2066 venue_order_id: venue,
2067 ..Default::default()
2068 })
2069 }
2070
2071 fn ev_cancel_rejected(venue: Option<VenueOrderId>) -> OrderEventAny {
2072 OrderEventAny::CancelRejected(OrderCancelRejected {
2073 client_order_id: client_order_id(),
2074 venue_order_id: venue,
2075 ..Default::default()
2076 })
2077 }
2078
2079 fn ev_updated(venue: Option<VenueOrderId>) -> OrderEventAny {
2080 OrderEventAny::Updated(OrderUpdated {
2081 client_order_id: client_order_id(),
2082 venue_order_id: venue,
2083 ..Default::default()
2084 })
2085 }
2086
2087 fn ev_filled() -> OrderEventAny {
2088 OrderEventAny::Filled(make_order_filled())
2089 }
2090
2091 fn ev_fill_voided() -> OrderEventAny {
2092 OrderEventAny::FillVoided(
2093 OrderFillVoidedSpec::builder()
2094 .client_order_id(client_order_id())
2095 .venue_order_id(venue_order_id())
2096 .build(),
2097 )
2098 }
2099
2100 #[rstest]
2101 fn trading_command_envelope_stamps_inner_submit_order_payload_type() {
2102 let cmd = make_submit_order();
2107 let bare = encode_submit_order(&cmd).expect("bare");
2108
2109 let envelope = TradingCommand::SubmitOrder(cmd);
2110 let wrapped = encode_trading_command(&envelope).expect("envelope");
2111
2112 assert_eq!(wrapped.payload, bare.payload);
2113 assert_eq!(wrapped.index_keys, bare.index_keys);
2114 assert_eq!(
2115 wrapped.payload_type.expect("override").as_str(),
2116 PAYLOAD_TYPE_SUBMIT_ORDER,
2117 );
2118 }
2119
2120 #[rstest]
2121 fn trading_command_cancel_order_envelope_emits_client_and_venue_indices() {
2122 let cancel = make_cancel_order();
2123 let envelope = TradingCommand::CancelOrder(cancel.clone());
2124 let wrapped = encode_trading_command(&envelope).expect("envelope");
2125
2126 assert_eq!(
2127 wrapped.payload_type.expect("override").as_str(),
2128 PAYLOAD_TYPE_CANCEL_ORDER,
2129 );
2130 assert_eq!(wrapped.index_keys.len(), 2);
2131 assert_eq!(wrapped.index_keys[0].kind, IndexKind::ClientOrderId);
2132 assert_eq!(
2133 wrapped.index_keys[0].key,
2134 cancel.client_order_id.to_string(),
2135 );
2136 assert_eq!(wrapped.index_keys[1].kind, IndexKind::VenueOrderId);
2137 assert_eq!(
2138 wrapped.index_keys[1].key,
2139 cancel.venue_order_id.expect("set").to_string(),
2140 );
2141
2142 let decoded: CancelOrder = rmp_serde::from_slice(&wrapped.payload).expect("decode");
2143 assert_eq!(decoded, cancel);
2144 }
2145
2146 #[rstest]
2147 fn trading_command_query_account_envelope_records_no_order_indices() {
2148 let envelope = TradingCommand::QueryAccount(make_query_account());
2151 let wrapped = encode_trading_command(&envelope).expect("envelope");
2152
2153 assert_eq!(
2154 wrapped.payload_type.expect("override").as_str(),
2155 PAYLOAD_TYPE_QUERY_ACCOUNT,
2156 );
2157 assert!(wrapped.index_keys.is_empty());
2158 }
2159
2160 #[rstest]
2161 fn order_event_any_envelope_stamps_inner_filled_payload_type() {
2162 let filled = make_order_filled();
2163 let bare = encode_order_filled(&filled).expect("bare");
2164
2165 let envelope = OrderEventAny::Filled(filled);
2166 let wrapped = encode_order_event_any(&envelope).expect("envelope");
2167
2168 assert_eq!(wrapped.payload, bare.payload);
2169 assert_eq!(wrapped.index_keys, bare.index_keys);
2170 assert_eq!(
2171 wrapped.payload_type.expect("override").as_str(),
2172 PAYLOAD_TYPE_ORDER_FILLED,
2173 );
2174 }
2175
2176 #[rstest]
2177 fn order_event_any_submitted_envelope_emits_client_order_id_index() {
2178 let submitted = make_order_submitted();
2179 let envelope = OrderEventAny::Submitted(submitted);
2180 let wrapped = encode_order_event_any(&envelope).expect("envelope");
2181
2182 assert_eq!(
2183 wrapped.payload_type.expect("override").as_str(),
2184 PAYLOAD_TYPE_ORDER_SUBMITTED,
2185 );
2186 assert_eq!(wrapped.index_keys.len(), 1);
2187 assert_eq!(wrapped.index_keys[0].kind, IndexKind::ClientOrderId);
2188 assert_eq!(
2189 wrapped.index_keys[0].key,
2190 submitted.client_order_id.to_string(),
2191 );
2192
2193 let decoded: OrderSubmitted = rmp_serde::from_slice(&wrapped.payload).expect("decode");
2194 assert_eq!(decoded, submitted);
2195 }
2196
2197 #[rstest]
2202 #[case::submit_order(
2203 TradingCommand::SubmitOrder(make_submit_order()),
2204 PAYLOAD_TYPE_SUBMIT_ORDER,
2205 1
2206 )]
2207 #[case::submit_order_list(
2208 TradingCommand::SubmitOrderList(make_submit_order_list(vec![
2209 ClientOrderId::from("O-A"),
2210 ClientOrderId::from("O-B"),
2211 ])),
2212 PAYLOAD_TYPE_SUBMIT_ORDER_LIST,
2213 2,
2214 )]
2215 #[case::modify_order(
2216 TradingCommand::ModifyOrder(make_modify_order(Some(venue_order_id()))),
2217 PAYLOAD_TYPE_MODIFY_ORDER,
2218 2
2219 )]
2220 #[case::batch_modify_orders(
2221 TradingCommand::ModifyOrders(make_batch_modify_orders(vec![
2222 make_modify_order(Some(venue_order_id())),
2223 ])),
2224 PAYLOAD_TYPE_BATCH_MODIFY_ORDERS,
2225 2,
2226 )]
2227 #[case::cancel_order(
2228 TradingCommand::CancelOrder(make_cancel_order()),
2229 PAYLOAD_TYPE_CANCEL_ORDER,
2230 2
2231 )]
2232 #[case::batch_cancel_orders(
2233 TradingCommand::CancelOrders(make_batch_cancel_orders(vec![make_cancel_order()])),
2234 PAYLOAD_TYPE_BATCH_CANCEL_ORDERS,
2235 2,
2236 )]
2237 #[case::cancel_all_orders(
2238 TradingCommand::CancelAllOrders(make_cancel_all_orders()),
2239 PAYLOAD_TYPE_CANCEL_ALL_ORDERS,
2240 0
2241 )]
2242 #[case::query_order(
2243 TradingCommand::QueryOrder(make_query_order(Some(venue_order_id()))),
2244 PAYLOAD_TYPE_QUERY_ORDER,
2245 2
2246 )]
2247 #[case::query_account(
2248 TradingCommand::QueryAccount(make_query_account()),
2249 PAYLOAD_TYPE_QUERY_ACCOUNT,
2250 0
2251 )]
2252 fn trading_command_envelope_stamps_inner_tag_for_every_variant(
2253 #[case] command: TradingCommand,
2254 #[case] expected_tag: &str,
2255 #[case] expected_index_count: usize,
2256 ) {
2257 let encoded = encode_trading_command(&command).expect("encode");
2258 let tag = encoded.payload_type.expect("override").as_str().to_string();
2259
2260 assert_eq!(tag, expected_tag);
2261 assert_ne!(
2262 tag, PAYLOAD_TYPE_TRADING_COMMAND,
2263 "wrapper fallback tag must never reach the writer",
2264 );
2265 assert_eq!(encoded.index_keys.len(), expected_index_count);
2266 }
2267
2268 #[rstest]
2272 #[case::initialized(ev_initialized(), PAYLOAD_TYPE_ORDER_INITIALIZED, false)]
2273 #[case::denied(ev_denied(), PAYLOAD_TYPE_ORDER_DENIED, false)]
2274 #[case::emulated(ev_emulated(), PAYLOAD_TYPE_ORDER_EMULATED, false)]
2275 #[case::released(ev_released(), PAYLOAD_TYPE_ORDER_RELEASED, false)]
2276 #[case::submitted(ev_submitted(), PAYLOAD_TYPE_ORDER_SUBMITTED, false)]
2277 #[case::accepted(
2278 ev_accepted_with_venue(venue_order_id()),
2279 PAYLOAD_TYPE_ORDER_ACCEPTED,
2280 true
2281 )]
2282 #[case::rejected(ev_rejected(), PAYLOAD_TYPE_ORDER_REJECTED, false)]
2283 #[case::canceled(ev_canceled(Some(venue_order_id())), PAYLOAD_TYPE_ORDER_CANCELED, true)]
2284 #[case::expired(ev_expired(Some(venue_order_id())), PAYLOAD_TYPE_ORDER_EXPIRED, true)]
2285 #[case::triggered(
2286 ev_triggered(Some(venue_order_id())),
2287 PAYLOAD_TYPE_ORDER_TRIGGERED,
2288 true
2289 )]
2290 #[case::pending_update(
2291 ev_pending_update(Some(venue_order_id())),
2292 PAYLOAD_TYPE_ORDER_PENDING_UPDATE,
2293 true
2294 )]
2295 #[case::pending_cancel(
2296 ev_pending_cancel(Some(venue_order_id())),
2297 PAYLOAD_TYPE_ORDER_PENDING_CANCEL,
2298 true
2299 )]
2300 #[case::modify_rejected(
2301 ev_modify_rejected(Some(venue_order_id())),
2302 PAYLOAD_TYPE_ORDER_MODIFY_REJECTED,
2303 true
2304 )]
2305 #[case::cancel_rejected(
2306 ev_cancel_rejected(Some(venue_order_id())),
2307 PAYLOAD_TYPE_ORDER_CANCEL_REJECTED,
2308 true
2309 )]
2310 #[case::updated(ev_updated(Some(venue_order_id())), PAYLOAD_TYPE_ORDER_UPDATED, true)]
2311 #[case::filled(ev_filled(), PAYLOAD_TYPE_ORDER_FILLED, true)]
2312 #[case::fill_voided(ev_fill_voided(), PAYLOAD_TYPE_ORDER_FILL_VOIDED, true)]
2313 fn order_event_any_envelope_stamps_inner_tag_for_every_variant(
2314 #[case] event: OrderEventAny,
2315 #[case] expected_tag: &str,
2316 #[case] expects_venue_index: bool,
2317 ) {
2318 let encoded = encode_order_event_any(&event).expect("encode");
2319 let tag = encoded.payload_type.expect("override").as_str().to_string();
2320
2321 assert_eq!(tag, expected_tag);
2322 assert_ne!(
2323 tag, PAYLOAD_TYPE_ORDER_EVENT_ANY,
2324 "wrapper fallback tag must never reach the writer",
2325 );
2326 assert_eq!(
2327 encoded.index_keys[0].kind,
2328 IndexKind::ClientOrderId,
2329 "first index must always be ClientOrderId for every order event",
2330 );
2331 assert_eq!(encoded.index_keys[0].key, client_order_id().to_string(),);
2332 if expects_venue_index {
2333 assert_eq!(encoded.index_keys.len(), 2);
2334 assert_eq!(encoded.index_keys[1].kind, IndexKind::VenueOrderId);
2335 } else {
2336 assert_eq!(encoded.index_keys.len(), 1);
2337 }
2338 }
2339
2340 #[rstest]
2343 #[case::cancel_order_some(TradingCommand::CancelOrder(make_cancel_order()), 2)]
2344 #[case::cancel_order_none(
2345 TradingCommand::CancelOrder(CancelOrder {
2346 venue_order_id: None,
2347 ..make_cancel_order()
2348 }),
2349 1,
2350 )]
2351 #[case::modify_order_some(
2352 TradingCommand::ModifyOrder(make_modify_order(Some(venue_order_id()))),
2353 2
2354 )]
2355 #[case::modify_order_none(TradingCommand::ModifyOrder(make_modify_order(None)), 1)]
2356 #[case::batch_modify_orders(
2357 TradingCommand::ModifyOrders(make_batch_modify_orders(vec![
2358 make_modify_order(Some(venue_order_id())),
2359 make_modify_order(None),
2360 ])),
2361 3,
2362 )]
2363 #[case::query_order_some(
2364 TradingCommand::QueryOrder(make_query_order(Some(venue_order_id()))),
2365 2
2366 )]
2367 #[case::query_order_none(TradingCommand::QueryOrder(make_query_order(None)), 1)]
2368 fn trading_command_envelope_index_count_matches_venue_optionality(
2369 #[case] command: TradingCommand,
2370 #[case] expected_index_count: usize,
2371 ) {
2372 let encoded = encode_trading_command(&command).expect("encode");
2373 assert_eq!(encoded.index_keys.len(), expected_index_count);
2374 assert_eq!(encoded.index_keys[0].kind, IndexKind::ClientOrderId);
2375 if expected_index_count == 2 {
2376 assert_eq!(encoded.index_keys[1].kind, IndexKind::VenueOrderId);
2377 }
2378 }
2379
2380 #[rstest]
2381 #[case::canceled_some(ev_canceled(Some(venue_order_id())), 2)]
2382 #[case::canceled_none(ev_canceled(None), 1)]
2383 #[case::updated_some(ev_updated(Some(venue_order_id())), 2)]
2384 #[case::updated_none(ev_updated(None), 1)]
2385 #[case::pending_update_some(ev_pending_update(Some(venue_order_id())), 2)]
2386 #[case::pending_update_none(ev_pending_update(None), 1)]
2387 fn order_event_any_envelope_index_count_matches_venue_optionality(
2388 #[case] event: OrderEventAny,
2389 #[case] expected_index_count: usize,
2390 ) {
2391 let encoded = encode_order_event_any(&event).expect("encode");
2392 assert_eq!(encoded.index_keys.len(), expected_index_count);
2393 }
2394
2395 #[rstest]
2396 fn batch_cancel_orders_envelope_indexes_each_child_with_optional_venue() {
2397 let with_venue = make_cancel_order();
2401 let mut without_venue = make_cancel_order();
2402 without_venue.venue_order_id = None;
2403 without_venue.client_order_id = ClientOrderId::from("O-NOVENUE");
2404 let batch = make_batch_cancel_orders(vec![with_venue.clone(), without_venue.clone()]);
2405
2406 let encoded = encode_trading_command(&TradingCommand::CancelOrders(batch)).expect("encode");
2407
2408 assert_eq!(
2409 encoded.payload_type.expect("override").as_str(),
2410 PAYLOAD_TYPE_BATCH_CANCEL_ORDERS,
2411 );
2412 assert_eq!(encoded.index_keys.len(), 3);
2413 assert_eq!(encoded.index_keys[0].kind, IndexKind::ClientOrderId);
2414 assert_eq!(
2415 encoded.index_keys[0].key,
2416 with_venue.client_order_id.to_string(),
2417 );
2418 assert_eq!(encoded.index_keys[1].kind, IndexKind::VenueOrderId);
2419 assert_eq!(
2420 encoded.index_keys[1].key,
2421 with_venue.venue_order_id.expect("set").to_string(),
2422 );
2423 assert_eq!(encoded.index_keys[2].kind, IndexKind::ClientOrderId);
2424 assert_eq!(
2425 encoded.index_keys[2].key,
2426 without_venue.client_order_id.to_string(),
2427 );
2428 }
2429
2430 #[rstest]
2431 fn batch_modify_orders_envelope_indexes_each_child_with_optional_venue() {
2432 let with_venue = make_modify_order(Some(venue_order_id()));
2433 let mut without_venue = make_modify_order(None);
2434 without_venue.client_order_id = ClientOrderId::from("O-NOVENUE");
2435 let batch = make_batch_modify_orders(vec![with_venue.clone(), without_venue.clone()]);
2436
2437 let encoded = encode_trading_command(&TradingCommand::ModifyOrders(batch)).expect("encode");
2438
2439 assert_eq!(
2440 encoded.payload_type.expect("override").as_str(),
2441 PAYLOAD_TYPE_BATCH_MODIFY_ORDERS,
2442 );
2443 assert_eq!(encoded.index_keys.len(), 3);
2444 assert_eq!(encoded.index_keys[0].kind, IndexKind::ClientOrderId);
2445 assert_eq!(
2446 encoded.index_keys[0].key,
2447 with_venue.client_order_id.to_string(),
2448 );
2449 assert_eq!(encoded.index_keys[1].kind, IndexKind::VenueOrderId);
2450 assert_eq!(
2451 encoded.index_keys[1].key,
2452 with_venue.venue_order_id.expect("set").to_string(),
2453 );
2454 assert_eq!(encoded.index_keys[2].kind, IndexKind::ClientOrderId);
2455 assert_eq!(
2456 encoded.index_keys[2].key,
2457 without_venue.client_order_id.to_string(),
2458 );
2459 }
2460
2461 #[rstest]
2462 fn submit_order_list_envelope_indexes_each_client_order_id() {
2463 let ids = vec![
2467 ClientOrderId::from("O-LIST-1"),
2468 ClientOrderId::from("O-LIST-2"),
2469 ClientOrderId::from("O-LIST-3"),
2470 ];
2471 let cmd = make_submit_order_list(ids.clone());
2472 let encoded =
2473 encode_trading_command(&TradingCommand::SubmitOrderList(cmd)).expect("encode");
2474
2475 assert_eq!(
2476 encoded.payload_type.expect("override").as_str(),
2477 PAYLOAD_TYPE_SUBMIT_ORDER_LIST,
2478 );
2479 assert_eq!(encoded.index_keys.len(), ids.len());
2480 for (idx, expected_id) in ids.iter().enumerate() {
2481 assert_eq!(encoded.index_keys[idx].kind, IndexKind::ClientOrderId);
2482 assert_eq!(encoded.index_keys[idx].key, expected_id.to_string());
2483 }
2484 }
2485
2486 fn make_fill_report() -> FillReport {
2487 FillReport::new(
2488 AccountId::from("BINANCE-001"),
2489 instrument_id(),
2490 venue_order_id(),
2491 TradeId::from("T-1111"),
2492 OrderSide::Buy,
2493 Quantity::from("1"),
2494 Price::from("100.00"),
2495 Money::new(0.10, Currency::USDT()),
2496 LiquiditySide::Taker,
2497 Some(client_order_id()),
2498 None,
2499 UnixNanos::from(40),
2500 UnixNanos::from(41),
2501 None,
2502 )
2503 }
2504
2505 fn make_position_status_report() -> PositionStatusReport {
2506 PositionStatusReport::new(
2507 AccountId::from("BINANCE-001"),
2508 instrument_id(),
2509 PositionSide::Long,
2510 Quantity::from("1"),
2511 UnixNanos::from(50),
2512 UnixNanos::from(51),
2513 None,
2514 Some(PositionId::from("P-001")),
2515 None,
2516 )
2517 }
2518
2519 fn make_execution_mass_status_with_reports() -> ExecutionMassStatus {
2520 let mut status = ExecutionMassStatus::new(
2521 ClientId::from("BINANCE"),
2522 AccountId::from("BINANCE-001"),
2523 Venue::from("BINANCE"),
2524 UnixNanos::from(60),
2525 None,
2526 );
2527 status.add_order_reports(vec![make_order_status_report()]);
2528 status.add_fill_reports(vec![make_fill_report()]);
2529 status.add_position_reports(vec![make_position_status_report()]);
2530 status
2531 }
2532
2533 #[rstest]
2534 fn execution_report_order_envelope_reuses_bare_status_encoder() {
2535 let report = make_order_status_report();
2540 let bare = encode_order_status_report(&report).expect("bare");
2541
2542 let envelope = ExecutionReport::Order(Box::new(report));
2543 let wrapped = encode_execution_report(&envelope).expect("envelope");
2544
2545 assert_eq!(wrapped.payload, bare.payload);
2546 assert_eq!(wrapped.index_keys, bare.index_keys);
2547 assert_eq!(
2548 wrapped.payload_type.expect("override").as_str(),
2549 PAYLOAD_TYPE_ORDER_STATUS_REPORT,
2550 );
2551 }
2552
2553 #[rstest]
2554 fn execution_report_fill_envelope_emits_venue_and_client_order_id_indices() {
2555 let fill = make_fill_report();
2556 let envelope = ExecutionReport::Fill(Box::new(fill.clone()));
2557 let encoded = encode_execution_report(&envelope).expect("encode");
2558
2559 assert_eq!(
2560 encoded.payload_type.expect("override").as_str(),
2561 PAYLOAD_TYPE_FILL_REPORT,
2562 );
2563 assert_eq!(encoded.index_keys.len(), 2);
2564 assert_eq!(encoded.index_keys[0].kind, IndexKind::VenueOrderId);
2565 assert_eq!(encoded.index_keys[0].key, fill.venue_order_id.to_string());
2566 assert_eq!(encoded.index_keys[1].kind, IndexKind::ClientOrderId);
2567 assert_eq!(
2568 encoded.index_keys[1].key,
2569 fill.client_order_id.expect("set").to_string(),
2570 );
2571
2572 let decoded: FillReport = rmp_serde::from_slice(&encoded.payload).expect("decode");
2573 assert_eq!(decoded, fill);
2574 }
2575
2576 #[rstest]
2577 fn execution_report_fill_envelope_omits_client_order_id_when_absent() {
2578 let mut fill = make_fill_report();
2579 fill.client_order_id = None;
2580 let envelope = ExecutionReport::Fill(Box::new(fill));
2581 let encoded = encode_execution_report(&envelope).expect("encode");
2582
2583 assert_eq!(encoded.index_keys.len(), 1);
2584 assert_eq!(encoded.index_keys[0].kind, IndexKind::VenueOrderId);
2585 }
2586
2587 #[rstest]
2588 fn execution_report_position_envelope_records_no_indices() {
2589 let position = make_position_status_report();
2593 let envelope = ExecutionReport::Position(Box::new(position.clone()));
2594 let encoded = encode_execution_report(&envelope).expect("encode");
2595
2596 assert_eq!(
2597 encoded.payload_type.expect("override").as_str(),
2598 PAYLOAD_TYPE_POSITION_STATUS_REPORT,
2599 );
2600 assert!(encoded.index_keys.is_empty());
2601
2602 let decoded: PositionStatusReport =
2603 rmp_serde::from_slice(&encoded.payload).expect("decode");
2604 assert_eq!(decoded, position);
2605 }
2606
2607 #[rstest]
2608 fn execution_report_order_with_fills_envelope_dedupes_shared_order_ids() {
2609 let order = make_order_status_report();
2613 let fills = vec![make_fill_report()];
2614 let envelope = ExecutionReport::OrderWithFills(Box::new(order.clone()), fills);
2615 let encoded = encode_execution_report(&envelope).expect("encode");
2616
2617 assert_eq!(
2618 encoded.payload_type.expect("override").as_str(),
2619 PAYLOAD_TYPE_ORDER_WITH_FILLS,
2620 );
2621 assert_eq!(encoded.index_keys.len(), 2);
2622 assert_eq!(encoded.index_keys[0].kind, IndexKind::VenueOrderId);
2623 assert_eq!(encoded.index_keys[0].key, order.venue_order_id.to_string());
2624 assert_eq!(encoded.index_keys[1].kind, IndexKind::ClientOrderId);
2625 assert_eq!(
2626 encoded.index_keys[1].key,
2627 order.client_order_id.expect("set").to_string(),
2628 );
2629 }
2630
2631 #[rstest]
2632 fn execution_report_order_with_fills_envelope_indexes_distinct_fill_ids() {
2633 let order = make_order_status_report();
2638 let mut fill = make_fill_report();
2639 fill.client_order_id = Some(ClientOrderId::from("O-EXTRA-001"));
2640 let envelope = ExecutionReport::OrderWithFills(Box::new(order), vec![fill.clone()]);
2641 let encoded = encode_execution_report(&envelope).expect("encode");
2642
2643 assert_eq!(encoded.index_keys.len(), 3);
2644 assert_eq!(encoded.index_keys[2].kind, IndexKind::ClientOrderId);
2645 assert_eq!(
2646 encoded.index_keys[2].key,
2647 fill.client_order_id.expect("set").to_string(),
2648 );
2649 }
2650
2651 #[rstest]
2652 fn execution_report_order_with_fills_payload_round_trips() {
2653 #[derive(serde::Deserialize)]
2654 struct OrderWithFillsOwned {
2655 order_report: OrderStatusReport,
2656 fill_reports: Vec<FillReport>,
2657 }
2658
2659 let order = make_order_status_report();
2660 let fills = vec![make_fill_report()];
2661 let envelope = ExecutionReport::OrderWithFills(Box::new(order.clone()), fills.clone());
2662 let encoded = encode_execution_report(&envelope).expect("encode");
2663
2664 let decoded: OrderWithFillsOwned = rmp_serde::from_slice(&encoded.payload).expect("decode");
2665 assert_eq!(decoded.order_report, order);
2666 assert_eq!(decoded.fill_reports, fills);
2667 }
2668
2669 #[rstest]
2670 fn execution_report_mass_status_envelope_indexes_orders_and_fills() {
2671 let status = make_execution_mass_status_with_reports();
2672 let envelope = ExecutionReport::MassStatus(Box::new(status.clone()));
2673 let encoded = encode_execution_report(&envelope).expect("encode");
2674
2675 assert_eq!(
2676 encoded.payload_type.expect("override").as_str(),
2677 PAYLOAD_TYPE_EXECUTION_MASS_STATUS,
2678 );
2679 assert_eq!(encoded.index_keys.len(), 2);
2682 assert_eq!(encoded.index_keys[0].kind, IndexKind::VenueOrderId);
2683 assert_eq!(encoded.index_keys[0].key, venue_order_id().to_string());
2684 assert_eq!(encoded.index_keys[1].kind, IndexKind::ClientOrderId);
2685 assert_eq!(encoded.index_keys[1].key, client_order_id().to_string());
2686
2687 let decoded: ExecutionMassStatus = rmp_serde::from_slice(&encoded.payload).expect("decode");
2688 assert_eq!(decoded, status);
2689 }
2690
2691 #[rstest]
2692 fn execution_report_mass_status_envelope_indexes_distinct_children() {
2693 let mut status = ExecutionMassStatus::new(
2697 ClientId::from("BINANCE"),
2698 AccountId::from("BINANCE-001"),
2699 Venue::from("BINANCE"),
2700 UnixNanos::from(60),
2701 None,
2702 );
2703 let order_a = make_order_status_report();
2704 let order_b = OrderStatusReport {
2705 client_order_id: Some(ClientOrderId::from("O-B")),
2706 venue_order_id: VenueOrderId::from("V-B"),
2707 ..make_order_status_report()
2708 };
2709 let fill_c = FillReport {
2710 venue_order_id: VenueOrderId::from("V-C"),
2711 client_order_id: Some(ClientOrderId::from("O-C")),
2712 ..make_fill_report()
2713 };
2714 status.add_order_reports(vec![order_a, order_b]);
2715 status.add_fill_reports(vec![fill_c]);
2716
2717 let encoded = encode_execution_report(&ExecutionReport::MassStatus(Box::new(status)))
2718 .expect("encode");
2719
2720 assert_eq!(encoded.index_keys.len(), 6);
2722 let venue_keys: Vec<&str> = encoded
2723 .index_keys
2724 .iter()
2725 .filter(|k| k.kind == IndexKind::VenueOrderId)
2726 .map(|k| k.key.as_str())
2727 .collect();
2728 let client_keys: Vec<&str> = encoded
2729 .index_keys
2730 .iter()
2731 .filter(|k| k.kind == IndexKind::ClientOrderId)
2732 .map(|k| k.key.as_str())
2733 .collect();
2734 assert!(venue_keys.contains(&venue_order_id().to_string().as_str()));
2735 assert!(venue_keys.contains(&"V-B"));
2736 assert!(venue_keys.contains(&"V-C"));
2737 assert!(client_keys.contains(&client_order_id().to_string().as_str()));
2738 assert!(client_keys.contains(&"O-B"));
2739 assert!(client_keys.contains(&"O-C"));
2740 }
2741
2742 #[rstest]
2746 #[case::order(
2747 ExecutionReport::Order(Box::new(make_order_status_report())),
2748 PAYLOAD_TYPE_ORDER_STATUS_REPORT
2749 )]
2750 #[case::fill(
2751 ExecutionReport::Fill(Box::new(make_fill_report())),
2752 PAYLOAD_TYPE_FILL_REPORT
2753 )]
2754 #[case::order_with_fills(
2755 ExecutionReport::OrderWithFills(
2756 Box::new(make_order_status_report()),
2757 vec![make_fill_report()],
2758 ),
2759 PAYLOAD_TYPE_ORDER_WITH_FILLS,
2760 )]
2761 #[case::position(
2762 ExecutionReport::Position(Box::new(make_position_status_report())),
2763 PAYLOAD_TYPE_POSITION_STATUS_REPORT
2764 )]
2765 #[case::mass_status(
2766 ExecutionReport::MassStatus(Box::new(make_execution_mass_status_with_reports())),
2767 PAYLOAD_TYPE_EXECUTION_MASS_STATUS
2768 )]
2769 fn execution_report_envelope_stamps_inner_tag_for_every_variant(
2770 #[case] report: ExecutionReport,
2771 #[case] expected_tag: &str,
2772 ) {
2773 let encoded = encode_execution_report(&report).expect("encode");
2774 let tag = encoded.payload_type.expect("override").as_str().to_string();
2775
2776 assert_eq!(tag, expected_tag);
2777 assert_ne!(
2778 tag, PAYLOAD_TYPE_EXECUTION_REPORT,
2779 "wrapper fallback tag must never reach the writer",
2780 );
2781 }
2782
2783 fn opening_order_id() -> ClientOrderId {
2784 ClientOrderId::from("O-OPEN-001")
2785 }
2786
2787 fn closing_order_id() -> ClientOrderId {
2788 ClientOrderId::from("O-CLOSE-001")
2789 }
2790
2791 fn position_id() -> PositionId {
2792 PositionId::from("P-001")
2793 }
2794
2795 fn make_position_opened() -> PositionOpened {
2796 PositionOpened {
2797 trader_id: trader_id(),
2798 strategy_id: strategy_id(),
2799 instrument_id: instrument_id(),
2800 position_id: position_id(),
2801 account_id: AccountId::from("BINANCE-001"),
2802 opening_order_id: opening_order_id(),
2803 entry: OrderSide::Buy,
2804 side: PositionSide::Long,
2805 signed_qty: 1.0,
2806 quantity: Quantity::from("1"),
2807 last_qty: Quantity::from("1"),
2808 last_px: Price::from("100.00"),
2809 currency: Currency::USDT(),
2810 avg_px_open: 100.0,
2811 realized_pnl: Some(Money::new(-0.1, Currency::USDT())),
2812 event_id: UUID4::new(),
2813 ts_event: UnixNanos::from(70),
2814 ts_init: UnixNanos::from(71),
2815 }
2816 }
2817
2818 fn make_position_changed() -> PositionChanged {
2819 PositionChanged {
2820 trader_id: trader_id(),
2821 strategy_id: strategy_id(),
2822 instrument_id: instrument_id(),
2823 position_id: position_id(),
2824 account_id: AccountId::from("BINANCE-001"),
2825 opening_order_id: opening_order_id(),
2826 entry: OrderSide::Buy,
2827 side: PositionSide::Long,
2828 signed_qty: 2.0,
2829 quantity: Quantity::from("2"),
2830 peak_quantity: Quantity::from("2"),
2831 last_qty: Quantity::from("1"),
2832 last_px: Price::from("101.00"),
2833 currency: Currency::USDT(),
2834 avg_px_open: 100.5,
2835 avg_px_close: None,
2836 realized_return: 0.0,
2837 realized_pnl: None,
2838 unrealized_pnl: Money::new(1.0, Currency::USDT()),
2839 event_id: UUID4::new(),
2840 ts_opened: UnixNanos::from(70),
2841 ts_event: UnixNanos::from(80),
2842 ts_init: UnixNanos::from(81),
2843 }
2844 }
2845
2846 fn make_position_closed() -> PositionClosed {
2847 PositionClosed {
2848 trader_id: trader_id(),
2849 strategy_id: strategy_id(),
2850 instrument_id: instrument_id(),
2851 position_id: position_id(),
2852 account_id: AccountId::from("BINANCE-001"),
2853 opening_order_id: opening_order_id(),
2854 closing_order_id: Some(closing_order_id()),
2855 entry: OrderSide::Buy,
2856 side: PositionSide::Flat,
2857 signed_qty: 0.0,
2858 quantity: Quantity::from("0"),
2859 peak_quantity: Quantity::from("2"),
2860 last_qty: Quantity::from("2"),
2861 last_px: Price::from("102.00"),
2862 currency: Currency::USDT(),
2863 avg_px_open: 100.5,
2864 avg_px_close: Some(102.0),
2865 realized_return: 0.015,
2866 realized_pnl: Some(Money::new(3.0, Currency::USDT())),
2867 unrealized_pnl: Money::new(0.0, Currency::USDT()),
2868 duration: 3_600_000_000_000,
2869 event_id: UUID4::new(),
2870 ts_opened: UnixNanos::from(70),
2871 ts_closed: Some(UnixNanos::from(90)),
2872 ts_event: UnixNanos::from(90),
2873 ts_init: UnixNanos::from(91),
2874 }
2875 }
2876
2877 fn make_position_adjusted() -> PositionAdjusted {
2878 PositionAdjusted::new(
2879 trader_id(),
2880 strategy_id(),
2881 instrument_id(),
2882 position_id(),
2883 AccountId::from("BINANCE-001"),
2884 PositionAdjustmentType::Commission,
2885 None,
2886 None,
2887 None,
2888 UUID4::new(),
2889 UnixNanos::from(100),
2890 UnixNanos::from(101),
2891 )
2892 }
2893
2894 #[rstest]
2895 fn position_event_opened_envelope_emits_opening_order_id_index() {
2896 let opened = make_position_opened();
2897 let envelope = PositionEvent::PositionOpened(opened.clone());
2898 let encoded = encode_position_event(&envelope).expect("encode");
2899
2900 assert_eq!(
2901 encoded.payload_type.expect("override").as_str(),
2902 PAYLOAD_TYPE_POSITION_OPENED,
2903 );
2904 assert_eq!(encoded.index_keys.len(), 1);
2905 assert_eq!(encoded.index_keys[0].kind, IndexKind::ClientOrderId);
2906 assert_eq!(
2907 encoded.index_keys[0].key,
2908 opened.opening_order_id.to_string(),
2909 );
2910
2911 let decoded: PositionOpened = rmp_serde::from_slice(&encoded.payload).expect("decode");
2912 assert_eq!(decoded, opened);
2913 }
2914
2915 #[rstest]
2916 fn position_event_changed_envelope_emits_opening_order_id_index() {
2917 let changed = make_position_changed();
2918 let envelope = PositionEvent::PositionChanged(changed.clone());
2919 let encoded = encode_position_event(&envelope).expect("encode");
2920
2921 assert_eq!(
2922 encoded.payload_type.expect("override").as_str(),
2923 PAYLOAD_TYPE_POSITION_CHANGED,
2924 );
2925 assert_eq!(encoded.index_keys.len(), 1);
2926 assert_eq!(encoded.index_keys[0].kind, IndexKind::ClientOrderId);
2927 assert_eq!(
2928 encoded.index_keys[0].key,
2929 changed.opening_order_id.to_string(),
2930 );
2931
2932 let decoded: PositionChanged = rmp_serde::from_slice(&encoded.payload).expect("decode");
2933 assert_eq!(decoded, changed);
2934 }
2935
2936 #[rstest]
2937 fn position_event_closed_envelope_indexes_both_opening_and_closing_order_ids() {
2938 let closed = make_position_closed();
2939 let envelope = PositionEvent::PositionClosed(closed.clone());
2940 let encoded = encode_position_event(&envelope).expect("encode");
2941
2942 assert_eq!(
2943 encoded.payload_type.expect("override").as_str(),
2944 PAYLOAD_TYPE_POSITION_CLOSED,
2945 );
2946 assert_eq!(encoded.index_keys.len(), 2);
2947 assert_eq!(encoded.index_keys[0].kind, IndexKind::ClientOrderId);
2948 assert_eq!(
2949 encoded.index_keys[0].key,
2950 closed.opening_order_id.to_string(),
2951 );
2952 assert_eq!(encoded.index_keys[1].kind, IndexKind::ClientOrderId);
2953 assert_eq!(
2954 encoded.index_keys[1].key,
2955 closed.closing_order_id.expect("set").to_string(),
2956 );
2957
2958 let decoded: PositionClosed = rmp_serde::from_slice(&encoded.payload).expect("decode");
2959 assert_eq!(decoded, closed);
2960 }
2961
2962 #[rstest]
2963 fn position_event_closed_envelope_omits_closing_order_id_when_absent() {
2964 let mut closed = make_position_closed();
2965 closed.closing_order_id = None;
2966 let envelope = PositionEvent::PositionClosed(closed);
2967 let encoded = encode_position_event(&envelope).expect("encode");
2968
2969 assert_eq!(encoded.index_keys.len(), 1);
2970 assert_eq!(encoded.index_keys[0].kind, IndexKind::ClientOrderId);
2971 assert_eq!(encoded.index_keys[0].key, opening_order_id().to_string());
2972 }
2973
2974 #[rstest]
2975 fn position_event_closed_envelope_dedupes_when_open_and_close_match() {
2976 let mut closed = make_position_closed();
2980 closed.closing_order_id = Some(closed.opening_order_id);
2981 let envelope = PositionEvent::PositionClosed(closed);
2982 let encoded = encode_position_event(&envelope).expect("encode");
2983
2984 assert_eq!(encoded.index_keys.len(), 1);
2985 assert_eq!(encoded.index_keys[0].key, opening_order_id().to_string());
2986 }
2987
2988 #[rstest]
2989 fn position_event_adjusted_envelope_records_no_indices() {
2990 let adjusted = make_position_adjusted();
2995 let envelope = PositionEvent::PositionAdjusted(adjusted);
2996 let encoded = encode_position_event(&envelope).expect("encode");
2997
2998 assert_eq!(
2999 encoded.payload_type.expect("override").as_str(),
3000 PAYLOAD_TYPE_POSITION_ADJUSTED,
3001 );
3002 assert!(encoded.index_keys.is_empty());
3003
3004 let decoded: PositionAdjusted = rmp_serde::from_slice(&encoded.payload).expect("decode");
3005 assert_eq!(decoded, adjusted);
3006 }
3007
3008 #[rstest]
3009 #[case::opened(
3010 PositionEvent::PositionOpened(make_position_opened()),
3011 PAYLOAD_TYPE_POSITION_OPENED
3012 )]
3013 #[case::changed(
3014 PositionEvent::PositionChanged(make_position_changed()),
3015 PAYLOAD_TYPE_POSITION_CHANGED
3016 )]
3017 #[case::closed(
3018 PositionEvent::PositionClosed(make_position_closed()),
3019 PAYLOAD_TYPE_POSITION_CLOSED
3020 )]
3021 #[case::adjusted(
3022 PositionEvent::PositionAdjusted(make_position_adjusted()),
3023 PAYLOAD_TYPE_POSITION_ADJUSTED
3024 )]
3025 fn position_event_envelope_stamps_inner_tag_for_every_variant(
3026 #[case] event: PositionEvent,
3027 #[case] expected_tag: &str,
3028 ) {
3029 let encoded = encode_position_event(&event).expect("encode");
3030 let tag = encoded.payload_type.expect("override").as_str().to_string();
3031
3032 assert_eq!(tag, expected_tag);
3033 assert_ne!(
3034 tag, PAYLOAD_TYPE_POSITION_EVENT,
3035 "wrapper fallback tag must never reach the writer",
3036 );
3037 }
3038
3039 fn make_account_state() -> AccountState {
3040 AccountState::new(
3041 AccountId::from("BINANCE-001"),
3042 AccountType::Cash,
3043 vec![AccountBalance::new(
3044 Money::from("1000000 USD"),
3045 Money::from("0 USD"),
3046 Money::from("1000000 USD"),
3047 )],
3048 vec![],
3049 true,
3050 UUID4::new(),
3051 UnixNanos::from(110),
3052 UnixNanos::from(111),
3053 Some(Currency::USD()),
3054 )
3055 }
3056
3057 #[rstest]
3058 fn account_state_encoder_records_no_indices() {
3059 let state = make_account_state();
3064 let encoded = encode_account_state(&state).expect("encode");
3065
3066 assert!(!encoded.payload.is_empty());
3067 assert!(encoded.index_keys.is_empty());
3068 assert!(
3069 encoded.payload_type.is_none(),
3070 "bare-type encoders inherit the registry's registered tag",
3071 );
3072 }
3073
3074 #[rstest]
3075 fn account_state_payload_round_trips_through_msgpack() {
3076 let state = make_account_state();
3077 let encoded = encode_account_state(&state).expect("encode");
3078
3079 let decoded: AccountState = rmp_serde::from_slice(&encoded.payload).expect("decode");
3080 assert_eq!(decoded, state);
3081 }
3082
3083 #[rstest]
3084 fn account_state_registered_under_canonical_payload_type() {
3085 let registry = default_registry();
3089 let state = make_account_state();
3090 let (tag, encoded) = registry
3091 .encode(&state)
3092 .expect("encode")
3093 .expect("registered");
3094
3095 assert_eq!(tag.as_str(), PAYLOAD_TYPE_ACCOUNT_STATE);
3096 assert!(encoded.index_keys.is_empty());
3097 }
3098
3099 #[derive(Debug, Deserialize, PartialEq, Eq)]
3100 struct DecodedTimeEventPayload {
3101 name: String,
3102 event_id: UUID4,
3103 ts_event: UnixNanos,
3104 ts_init: UnixNanos,
3105 }
3106
3107 #[rstest]
3108 fn time_event_payload_round_trips_through_msgpack() {
3109 let event = TimeEvent::new(
3110 Ustr::from("heartbeat"),
3111 UUID4::new(),
3112 UnixNanos::from(100),
3113 UnixNanos::from(99),
3114 );
3115 let encoded = encode_time_event(&event).expect("encode");
3116
3117 let decoded: DecodedTimeEventPayload =
3118 rmp_serde::from_slice(&encoded.payload).expect("decode");
3119 assert_eq!(
3120 decoded,
3121 DecodedTimeEventPayload {
3122 name: event.name.to_string(),
3123 event_id: event.event_id,
3124 ts_event: event.ts_event,
3125 ts_init: event.ts_init,
3126 },
3127 );
3128 assert!(encoded.index_keys.is_empty());
3129 }
3130
3131 #[rstest]
3132 fn time_event_registered_under_canonical_payload_type() {
3133 let registry = default_registry();
3134 let event = TimeEvent::new(
3135 Ustr::from("heartbeat"),
3136 UUID4::new(),
3137 UnixNanos::from(100),
3138 UnixNanos::from(99),
3139 );
3140 let (tag, encoded) = registry
3141 .encode(&event)
3142 .expect("encode")
3143 .expect("registered");
3144
3145 assert_eq!(tag.as_str(), PAYLOAD_TYPE_TIME_EVENT);
3146 assert!(encoded.index_keys.is_empty());
3147 }
3148
3149 #[rstest]
3150 fn data_command_request_envelope_stamps_request_command_payload_type() {
3151 let request = make_request_command();
3155 let envelope = DataCommand::Request(request.clone());
3156 let encoded = encode_data_command(&envelope).expect("encode");
3157
3158 assert_eq!(
3159 encoded.payload_type.expect("override").as_str(),
3160 PAYLOAD_TYPE_REQUEST_COMMAND,
3161 );
3162 assert!(encoded.index_keys.is_empty());
3163
3164 let decoded: RequestCommand = rmp_serde::from_slice(&encoded.payload).expect("decode");
3165 match (decoded, request) {
3166 (RequestCommand::Quotes(decoded), RequestCommand::Quotes(expected)) => {
3167 assert_eq!(decoded.request_id, expected.request_id);
3168 assert_eq!(decoded.instrument_id, expected.instrument_id);
3169 }
3170 other => panic!("expected RequestCommand::Quotes round trip, was {other:?}"),
3171 }
3172 }
3173
3174 #[rstest]
3175 fn data_command_subscribe_envelope_stamps_subscribe_command_payload_type() {
3176 let subscribe = make_subscribe_command();
3177 let envelope = DataCommand::Subscribe(subscribe.clone());
3178 let encoded = encode_data_command(&envelope).expect("encode");
3179
3180 assert_eq!(
3181 encoded.payload_type.expect("override").as_str(),
3182 PAYLOAD_TYPE_SUBSCRIBE_COMMAND,
3183 );
3184 assert!(encoded.index_keys.is_empty());
3185
3186 let decoded: SubscribeCommand = rmp_serde::from_slice(&encoded.payload).expect("decode");
3187 match (decoded, subscribe) {
3188 (SubscribeCommand::Quotes(decoded), SubscribeCommand::Quotes(expected)) => {
3189 assert_eq!(decoded.command_id, expected.command_id);
3190 assert_eq!(decoded.instrument_id, expected.instrument_id);
3191 }
3192 other => panic!("expected SubscribeCommand::Quotes round trip, was {other:?}"),
3193 }
3194 }
3195
3196 #[rstest]
3197 fn data_command_unsubscribe_envelope_stamps_unsubscribe_command_payload_type() {
3198 let unsubscribe = make_unsubscribe_command();
3199 let envelope = DataCommand::Unsubscribe(unsubscribe.clone());
3200 let encoded = encode_data_command(&envelope).expect("encode");
3201
3202 assert_eq!(
3203 encoded.payload_type.expect("override").as_str(),
3204 PAYLOAD_TYPE_UNSUBSCRIBE_COMMAND,
3205 );
3206 assert!(encoded.index_keys.is_empty());
3207
3208 let decoded: UnsubscribeCommand = rmp_serde::from_slice(&encoded.payload).expect("decode");
3209 match (decoded, unsubscribe) {
3210 (UnsubscribeCommand::Quotes(decoded), UnsubscribeCommand::Quotes(expected)) => {
3211 assert_eq!(decoded.command_id, expected.command_id);
3212 assert_eq!(decoded.instrument_id, expected.instrument_id);
3213 }
3214 other => panic!("expected UnsubscribeCommand::Quotes round trip, was {other:?}"),
3215 }
3216 }
3217
3218 #[rstest]
3219 fn data_command_registered_under_category_payload_type() {
3220 let registry = default_registry();
3223 let envelope = DataCommand::Subscribe(make_subscribe_command());
3224 let (tag, encoded) = registry
3225 .encode(&envelope)
3226 .expect("encode")
3227 .expect("registered");
3228
3229 assert_eq!(tag.as_str(), PAYLOAD_TYPE_SUBSCRIBE_COMMAND);
3230 assert!(encoded.index_keys.is_empty());
3231 }
3232
3233 #[cfg(feature = "defi")]
3234 #[rstest]
3235 fn data_command_defi_request_envelope_stamps_defi_request_command_payload_type() {
3236 let request = make_defi_request_command();
3237 let envelope = DataCommand::DefiRequest(request.clone());
3238 let encoded = encode_data_command(&envelope).expect("encode");
3239
3240 assert_eq!(
3241 encoded.payload_type.expect("override").as_str(),
3242 PAYLOAD_TYPE_DEFI_REQUEST_COMMAND,
3243 );
3244 assert!(encoded.index_keys.is_empty());
3245
3246 let decoded: DefiRequestCommand = rmp_serde::from_slice(&encoded.payload).expect("decode");
3247 match (decoded, request) {
3248 (
3249 DefiRequestCommand::PoolSnapshot(decoded),
3250 DefiRequestCommand::PoolSnapshot(expected),
3251 ) => {
3252 assert_eq!(decoded.request_id, expected.request_id);
3253 assert_eq!(decoded.instrument_id, expected.instrument_id);
3254 assert_eq!(decoded.client_id, expected.client_id);
3255 assert_eq!(decoded.ts_init, expected.ts_init);
3256 }
3257 }
3258 }
3259
3260 #[cfg(feature = "defi")]
3261 #[rstest]
3262 fn data_command_defi_subscribe_envelope_stamps_defi_subscribe_command_payload_type() {
3263 let subscribe = make_defi_subscribe_command();
3264 let envelope = DataCommand::DefiSubscribe(subscribe.clone());
3265 let encoded = encode_data_command(&envelope).expect("encode");
3266
3267 assert_eq!(
3268 encoded.payload_type.expect("override").as_str(),
3269 PAYLOAD_TYPE_DEFI_SUBSCRIBE_COMMAND,
3270 );
3271 assert!(encoded.index_keys.is_empty());
3272
3273 let decoded: DefiSubscribeCommand =
3274 rmp_serde::from_slice(&encoded.payload).expect("decode");
3275
3276 match (decoded, subscribe) {
3277 (DefiSubscribeCommand::Blocks(decoded), DefiSubscribeCommand::Blocks(expected)) => {
3278 assert_eq!(decoded.command_id, expected.command_id);
3279 assert_eq!(decoded.chain, expected.chain);
3280 assert_eq!(decoded.client_id, expected.client_id);
3281 assert_eq!(decoded.ts_init, expected.ts_init);
3282 }
3283 other => panic!("expected DefiSubscribeCommand::Blocks round trip, was {other:?}"),
3284 }
3285 }
3286
3287 #[cfg(feature = "defi")]
3288 #[rstest]
3289 fn data_command_defi_unsubscribe_envelope_stamps_defi_unsubscribe_command_payload_type() {
3290 let unsubscribe = make_defi_unsubscribe_command();
3291 let envelope = DataCommand::DefiUnsubscribe(unsubscribe.clone());
3292 let encoded = encode_data_command(&envelope).expect("encode");
3293
3294 assert_eq!(
3295 encoded.payload_type.expect("override").as_str(),
3296 PAYLOAD_TYPE_DEFI_UNSUBSCRIBE_COMMAND,
3297 );
3298 assert!(encoded.index_keys.is_empty());
3299
3300 let decoded: DefiUnsubscribeCommand =
3301 rmp_serde::from_slice(&encoded.payload).expect("decode");
3302
3303 match (decoded, unsubscribe) {
3304 (DefiUnsubscribeCommand::Blocks(decoded), DefiUnsubscribeCommand::Blocks(expected)) => {
3305 assert_eq!(decoded.command_id, expected.command_id);
3306 assert_eq!(decoded.chain, expected.chain);
3307 assert_eq!(decoded.client_id, expected.client_id);
3308 assert_eq!(decoded.ts_init, expected.ts_init);
3309 }
3310 other => panic!("expected DefiUnsubscribeCommand::Blocks round trip, was {other:?}"),
3311 }
3312 }
3313
3314 fn client_id() -> ClientId {
3315 ClientId::from("BINANCE")
3316 }
3317
3318 fn venue() -> Venue {
3319 Venue::from("BINANCE")
3320 }
3321
3322 fn correlation_id() -> UUID4 {
3323 UUID4::new()
3324 }
3325
3326 fn bar_type() -> BarType {
3327 BarType::from("ETHUSDT-PERP.BINANCE-1-MINUTE-LAST-EXTERNAL")
3328 }
3329
3330 fn data_type() -> DataType {
3331 DataType::new("Bar", None, None)
3332 }
3333
3334 fn make_request_command() -> RequestCommand {
3335 RequestCommand::Quotes(RequestQuotes::new(
3336 instrument_id(),
3337 None,
3338 None,
3339 None,
3340 Some(client_id()),
3341 correlation_id(),
3342 UnixNanos::from(197),
3343 None,
3344 ))
3345 }
3346
3347 fn make_subscribe_command() -> SubscribeCommand {
3348 SubscribeCommand::Quotes(SubscribeQuotes::new(
3349 instrument_id(),
3350 Some(client_id()),
3351 Some(venue()),
3352 correlation_id(),
3353 UnixNanos::from(198),
3354 Some(correlation_id()),
3355 None,
3356 ))
3357 }
3358
3359 fn make_unsubscribe_command() -> UnsubscribeCommand {
3360 UnsubscribeCommand::Quotes(UnsubscribeQuotes::new(
3361 instrument_id(),
3362 Some(client_id()),
3363 Some(venue()),
3364 correlation_id(),
3365 UnixNanos::from(199),
3366 Some(correlation_id()),
3367 None,
3368 ))
3369 }
3370
3371 #[cfg(feature = "defi")]
3372 fn make_defi_request_command() -> DefiRequestCommand {
3373 DefiRequestCommand::PoolSnapshot(RequestPoolSnapshot::new(
3374 instrument_id(),
3375 Some(client_id()),
3376 correlation_id(),
3377 UnixNanos::from(196),
3378 None,
3379 ))
3380 }
3381
3382 #[cfg(feature = "defi")]
3383 fn make_defi_subscribe_command() -> DefiSubscribeCommand {
3384 DefiSubscribeCommand::Blocks(SubscribeBlocks::new(
3385 Blockchain::Ethereum,
3386 Some(client_id()),
3387 correlation_id(),
3388 UnixNanos::from(195),
3389 None,
3390 ))
3391 }
3392
3393 #[cfg(feature = "defi")]
3394 fn make_defi_unsubscribe_command() -> DefiUnsubscribeCommand {
3395 DefiUnsubscribeCommand::Blocks(UnsubscribeBlocks::new(
3396 Blockchain::Ethereum,
3397 Some(client_id()),
3398 correlation_id(),
3399 UnixNanos::from(194),
3400 None,
3401 ))
3402 }
3403
3404 fn make_custom_data_response() -> CustomDataResponse {
3405 CustomDataResponse::new(
3406 correlation_id(),
3407 client_id(),
3408 Some(venue()),
3409 data_type(),
3410 (),
3411 None,
3412 None,
3413 UnixNanos::from(200),
3414 None,
3415 )
3416 }
3417
3418 fn make_instrument_response() -> InstrumentResponse {
3419 InstrumentResponse::new(
3420 correlation_id(),
3421 client_id(),
3422 instrument_id(),
3423 InstrumentAny::CurrencyPair(currency_pair_ethusdt()),
3424 None,
3425 None,
3426 UnixNanos::from(201),
3427 None,
3428 )
3429 }
3430
3431 fn make_instruments_response() -> InstrumentsResponse {
3432 InstrumentsResponse::new(
3433 correlation_id(),
3434 client_id(),
3435 venue(),
3436 vec![InstrumentAny::CurrencyPair(currency_pair_ethusdt())],
3437 None,
3438 None,
3439 UnixNanos::from(202),
3440 None,
3441 )
3442 }
3443
3444 fn make_book_response() -> BookResponse {
3445 BookResponse::new(
3446 correlation_id(),
3447 client_id(),
3448 instrument_id(),
3449 OrderBook::new(instrument_id(), BookType::L2_MBP),
3450 None,
3451 None,
3452 UnixNanos::from(203),
3453 None,
3454 )
3455 }
3456
3457 fn make_book_deltas_response() -> BookDeltasResponse {
3458 BookDeltasResponse::new(
3459 correlation_id(),
3460 client_id(),
3461 instrument_id(),
3462 Vec::new(),
3463 None,
3464 None,
3465 UnixNanos::from(204),
3466 None,
3467 )
3468 }
3469
3470 fn make_book_depth_response() -> BookDepthResponse {
3471 let mut depth = stub_depth10();
3472 depth.instrument_id = instrument_id();
3473 BookDepthResponse::new(
3474 correlation_id(),
3475 client_id(),
3476 instrument_id(),
3477 vec![depth],
3478 None,
3479 None,
3480 UnixNanos::from(205),
3481 None,
3482 )
3483 }
3484
3485 fn make_quotes_response() -> QuotesResponse {
3486 QuotesResponse::new(
3487 correlation_id(),
3488 client_id(),
3489 instrument_id(),
3490 Vec::new(),
3491 None,
3492 None,
3493 UnixNanos::from(206),
3494 None,
3495 )
3496 }
3497
3498 fn make_trades_response() -> TradesResponse {
3499 TradesResponse::new(
3500 correlation_id(),
3501 client_id(),
3502 instrument_id(),
3503 Vec::new(),
3504 None,
3505 None,
3506 UnixNanos::from(207),
3507 None,
3508 )
3509 }
3510
3511 fn make_funding_rates_response() -> FundingRatesResponse {
3512 FundingRatesResponse::new(
3513 correlation_id(),
3514 client_id(),
3515 instrument_id(),
3516 Vec::new(),
3517 None,
3518 None,
3519 UnixNanos::from(208),
3520 None,
3521 )
3522 }
3523
3524 fn make_forward_prices_response() -> ForwardPricesResponse {
3525 ForwardPricesResponse::new(
3526 correlation_id(),
3527 client_id(),
3528 venue(),
3529 Vec::new(),
3530 UnixNanos::from(209),
3531 None,
3532 )
3533 }
3534
3535 fn make_bars_response() -> BarsResponse {
3536 BarsResponse::new(
3537 correlation_id(),
3538 client_id(),
3539 bar_type(),
3540 Vec::<Bar>::new(),
3541 None,
3542 None,
3543 UnixNanos::from(210),
3544 None,
3545 )
3546 }
3547
3548 #[rstest]
3549 fn data_response_custom_data_envelope_stamps_inner_tag_with_no_indices() {
3550 let response = make_custom_data_response();
3554 let envelope = DataResponse::Data(response);
3555 let encoded = encode_data_response(&envelope).expect("encode");
3556
3557 assert_eq!(
3558 encoded.payload_type.expect("override").as_str(),
3559 PAYLOAD_TYPE_CUSTOM_DATA_RESPONSE,
3560 );
3561 assert!(encoded.index_keys.is_empty());
3562 assert!(!encoded.payload.is_empty());
3563 }
3564
3565 #[rstest]
3566 fn data_response_custom_data_payload_round_trips_metadata() {
3567 #[derive(serde::Deserialize)]
3571 struct CustomDataResponseOwned {
3572 correlation_id: UUID4,
3573 client_id: ClientId,
3574 venue: Option<Venue>,
3575 ts_init: UnixNanos,
3576 }
3577
3578 let response = make_custom_data_response();
3579 let envelope = DataResponse::Data(response.clone());
3580 let encoded = encode_data_response(&envelope).expect("encode");
3581
3582 let decoded: CustomDataResponseOwned =
3583 rmp_serde::from_slice(&encoded.payload).expect("decode");
3584 assert_eq!(decoded.correlation_id, response.correlation_id);
3585 assert_eq!(decoded.client_id, response.client_id);
3586 assert_eq!(decoded.venue, response.venue);
3587 assert_eq!(decoded.ts_init, response.ts_init);
3588 }
3589
3590 #[rstest]
3591 fn data_response_book_envelope_stamps_inner_tag_with_no_indices() {
3592 let response = make_book_response();
3593 let envelope = DataResponse::Book(response);
3594 let encoded = encode_data_response(&envelope).expect("encode");
3595
3596 assert_eq!(
3597 encoded.payload_type.expect("override").as_str(),
3598 PAYLOAD_TYPE_BOOK_RESPONSE,
3599 );
3600 assert!(encoded.index_keys.is_empty());
3601 assert!(!encoded.payload.is_empty());
3602 }
3603
3604 #[rstest]
3605 fn data_response_book_payload_round_trips_metadata() {
3606 #[derive(serde::Deserialize)]
3609 struct BookResponseOwned {
3610 correlation_id: UUID4,
3611 instrument_id: InstrumentId,
3612 }
3613
3614 let response = make_book_response();
3615 let envelope = DataResponse::Book(response.clone());
3616 let encoded = encode_data_response(&envelope).expect("encode");
3617
3618 let decoded: BookResponseOwned = rmp_serde::from_slice(&encoded.payload).expect("decode");
3619 assert_eq!(decoded.correlation_id, response.correlation_id);
3620 assert_eq!(decoded.instrument_id, response.instrument_id);
3621 }
3622
3623 #[rstest]
3624 fn data_response_book_depth_payload_round_trips() {
3625 let response = make_book_depth_response();
3626 let envelope = DataResponse::BookDepth(response.clone());
3627 let encoded = encode_data_response(&envelope).expect("encode");
3628
3629 assert_eq!(
3630 encoded.payload_type.expect("override").as_str(),
3631 PAYLOAD_TYPE_BOOK_DEPTH_RESPONSE,
3632 );
3633 assert!(encoded.index_keys.is_empty());
3634
3635 let decoded: BookDepthResponse = rmp_serde::from_slice(&encoded.payload).expect("decode");
3636 assert_eq!(decoded.correlation_id, response.correlation_id);
3637 assert_eq!(decoded.instrument_id, response.instrument_id);
3638 assert_eq!(decoded.data, response.data);
3639 }
3640
3641 #[rstest]
3642 fn data_response_quotes_payload_round_trips() {
3643 let response = make_quotes_response();
3644 let envelope = DataResponse::Quotes(response.clone());
3645 let encoded = encode_data_response(&envelope).expect("encode");
3646
3647 assert_eq!(
3648 encoded.payload_type.expect("override").as_str(),
3649 PAYLOAD_TYPE_QUOTES_RESPONSE,
3650 );
3651 assert!(encoded.index_keys.is_empty());
3652
3653 let decoded: QuotesResponse = rmp_serde::from_slice(&encoded.payload).expect("decode");
3654 assert_eq!(decoded.correlation_id, response.correlation_id);
3655 assert_eq!(decoded.instrument_id, response.instrument_id);
3656 assert_eq!(decoded.data, response.data);
3657 }
3658
3659 #[rstest]
3660 fn data_response_trades_payload_round_trips() {
3661 let response = make_trades_response();
3662 let envelope = DataResponse::Trades(response.clone());
3663 let encoded = encode_data_response(&envelope).expect("encode");
3664
3665 assert_eq!(
3666 encoded.payload_type.expect("override").as_str(),
3667 PAYLOAD_TYPE_TRADES_RESPONSE,
3668 );
3669
3670 let decoded: TradesResponse = rmp_serde::from_slice(&encoded.payload).expect("decode");
3671 assert_eq!(decoded.correlation_id, response.correlation_id);
3672 }
3673
3674 #[rstest]
3675 fn data_response_bars_payload_round_trips() {
3676 let response = make_bars_response();
3677 let envelope = DataResponse::Bars(response.clone());
3678 let encoded = encode_data_response(&envelope).expect("encode");
3679
3680 assert_eq!(
3681 encoded.payload_type.expect("override").as_str(),
3682 PAYLOAD_TYPE_BARS_RESPONSE,
3683 );
3684
3685 let decoded: BarsResponse = rmp_serde::from_slice(&encoded.payload).expect("decode");
3686 assert_eq!(decoded.correlation_id, response.correlation_id);
3687 assert_eq!(decoded.bar_type, response.bar_type);
3688 }
3689
3690 #[rstest]
3691 fn data_response_instrument_payload_round_trips() {
3692 let response = make_instrument_response();
3693 let envelope = DataResponse::Instrument(Box::new(response.clone()));
3694 let encoded = encode_data_response(&envelope).expect("encode");
3695
3696 assert_eq!(
3697 encoded.payload_type.expect("override").as_str(),
3698 PAYLOAD_TYPE_INSTRUMENT_RESPONSE,
3699 );
3700
3701 let decoded: InstrumentResponse = rmp_serde::from_slice(&encoded.payload).expect("decode");
3702 assert_eq!(decoded.correlation_id, response.correlation_id);
3703 assert_eq!(decoded.instrument_id, response.instrument_id);
3704 }
3705
3706 #[rstest]
3711 #[case::data(
3712 DataResponse::Data(make_custom_data_response()),
3713 PAYLOAD_TYPE_CUSTOM_DATA_RESPONSE
3714 )]
3715 #[case::instrument(
3716 DataResponse::Instrument(Box::new(make_instrument_response())),
3717 PAYLOAD_TYPE_INSTRUMENT_RESPONSE
3718 )]
3719 #[case::instruments(
3720 DataResponse::Instruments(make_instruments_response()),
3721 PAYLOAD_TYPE_INSTRUMENTS_RESPONSE
3722 )]
3723 #[case::book(DataResponse::Book(make_book_response()), PAYLOAD_TYPE_BOOK_RESPONSE)]
3724 #[case::book_deltas(
3725 DataResponse::BookDeltas(make_book_deltas_response()),
3726 PAYLOAD_TYPE_BOOK_DELTAS_RESPONSE
3727 )]
3728 #[case::book_depth(
3729 DataResponse::BookDepth(make_book_depth_response()),
3730 PAYLOAD_TYPE_BOOK_DEPTH_RESPONSE
3731 )]
3732 #[case::quotes(
3733 DataResponse::Quotes(make_quotes_response()),
3734 PAYLOAD_TYPE_QUOTES_RESPONSE
3735 )]
3736 #[case::trades(
3737 DataResponse::Trades(make_trades_response()),
3738 PAYLOAD_TYPE_TRADES_RESPONSE
3739 )]
3740 #[case::funding_rates(
3741 DataResponse::FundingRates(make_funding_rates_response()),
3742 PAYLOAD_TYPE_FUNDING_RATES_RESPONSE
3743 )]
3744 #[case::forward_prices(
3745 DataResponse::ForwardPrices(make_forward_prices_response()),
3746 PAYLOAD_TYPE_FORWARD_PRICES_RESPONSE
3747 )]
3748 #[case::bars(DataResponse::Bars(make_bars_response()), PAYLOAD_TYPE_BARS_RESPONSE)]
3749 fn data_response_envelope_stamps_inner_tag_for_every_variant(
3750 #[case] response: DataResponse,
3751 #[case] expected_tag: &str,
3752 ) {
3753 let encoded = encode_data_response(&response).expect("encode");
3754 let tag = encoded.payload_type.expect("override").as_str().to_string();
3755
3756 assert_eq!(tag, expected_tag);
3757 assert_ne!(
3758 tag, PAYLOAD_TYPE_DATA_RESPONSE,
3759 "wrapper fallback tag must never reach the writer",
3760 );
3761 assert!(
3762 encoded.index_keys.is_empty(),
3763 "DataResponse correlation_id has no matching IndexKind today",
3764 );
3765 }
3766
3767 #[rstest]
3768 fn data_response_registered_under_canonical_payload_type() {
3769 let registry = default_registry();
3773 let envelope = DataResponse::Quotes(make_quotes_response());
3774 let (tag, encoded) = registry
3775 .encode(&envelope)
3776 .expect("encode")
3777 .expect("registered");
3778
3779 assert_eq!(tag.as_str(), PAYLOAD_TYPE_QUOTES_RESPONSE);
3780 assert!(encoded.index_keys.is_empty());
3781 }
3782
3783 #[rstest]
3784 fn data_response_headers_extractor_surfaces_correlation_id() {
3785 let registry = default_registry();
3790 let response = make_quotes_response();
3791 let expected = response.correlation_id;
3792 let envelope = DataResponse::Quotes(response);
3793
3794 let headers = registry
3795 .headers_for_any(&envelope as &dyn std::any::Any)
3796 .expect("registered");
3797 assert_eq!(headers.correlation_id, Some(expected));
3798 assert_eq!(headers.causation_id, None);
3799 }
3800
3801 #[rstest]
3802 fn data_command_request_headers_use_request_id_as_correlation() {
3803 let registry = default_registry();
3808 let request = make_quotes_request();
3809 let expected = request.request_id;
3810 let envelope = DataCommand::Request(RequestCommand::Quotes(request));
3811
3812 let headers = registry
3813 .headers_for_any(&envelope as &dyn std::any::Any)
3814 .expect("registered");
3815 assert_eq!(headers.correlation_id, Some(expected));
3816 }
3817
3818 #[rstest]
3819 fn data_command_subscribe_headers_surface_correlation_id() {
3820 let registry = default_registry();
3824 let subscribe = make_subscribe_command();
3825 let expected = subscribe.correlation_id();
3826 let envelope = DataCommand::Subscribe(subscribe);
3827
3828 let headers = registry
3829 .headers_for_any(&envelope as &dyn std::any::Any)
3830 .expect("registered");
3831 assert_eq!(headers.correlation_id, expected);
3832 }
3833
3834 #[rstest]
3835 fn data_command_unsubscribe_headers_surface_correlation_id() {
3836 let registry = default_registry();
3837 let unsubscribe = make_unsubscribe_command();
3838 let expected = unsubscribe.correlation_id();
3839 let envelope = DataCommand::Unsubscribe(unsubscribe);
3840
3841 let headers = registry
3842 .headers_for_any(&envelope as &dyn std::any::Any)
3843 .expect("registered");
3844 assert_eq!(headers.correlation_id, expected);
3845 }
3846
3847 #[rstest]
3848 #[case::submit_order(trading_command_submit_order)]
3849 #[case::submit_order_list(trading_command_submit_order_list)]
3850 #[case::modify_order(trading_command_modify_order)]
3851 #[case::batch_modify_orders(trading_command_batch_modify_orders)]
3852 #[case::cancel_order(trading_command_cancel_order)]
3853 #[case::cancel_all_orders(trading_command_cancel_all_orders)]
3854 #[case::batch_cancel_orders(trading_command_batch_cancel_orders)]
3855 #[case::query_order(trading_command_query_order)]
3856 #[case::query_account(trading_command_query_account)]
3857 fn trading_command_extractor_surfaces_both_headers(
3858 #[case] builder: fn() -> (TradingCommand, UUID4, UUID4),
3859 ) {
3860 let (envelope, corr, caus) = builder();
3865 let registry = default_registry();
3866
3867 let headers = registry
3868 .headers_for_any(&envelope as &dyn std::any::Any)
3869 .expect("registered");
3870 assert_eq!(headers.correlation_id, Some(corr));
3871 assert_eq!(headers.causation_id, Some(caus));
3872 }
3873
3874 fn trading_command_submit_order() -> (TradingCommand, UUID4, UUID4) {
3875 let corr = UUID4::new();
3876 let caus = UUID4::new();
3877 let mut cmd = make_submit_order();
3878 cmd.correlation_id = Some(corr);
3879 cmd.causation_id = Some(caus);
3880 (TradingCommand::SubmitOrder(cmd), corr, caus)
3881 }
3882
3883 fn trading_command_submit_order_list() -> (TradingCommand, UUID4, UUID4) {
3884 let corr = UUID4::new();
3885 let caus = UUID4::new();
3886 let mut cmd = make_submit_order_list(vec![client_order_id()]);
3887 cmd.correlation_id = Some(corr);
3888 cmd.causation_id = Some(caus);
3889 (TradingCommand::SubmitOrderList(cmd), corr, caus)
3890 }
3891
3892 fn trading_command_modify_order() -> (TradingCommand, UUID4, UUID4) {
3893 let corr = UUID4::new();
3894 let caus = UUID4::new();
3895 let mut cmd = make_modify_order(Some(venue_order_id()));
3896 cmd.correlation_id = Some(corr);
3897 cmd.causation_id = Some(caus);
3898 (TradingCommand::ModifyOrder(cmd), corr, caus)
3899 }
3900
3901 fn trading_command_batch_modify_orders() -> (TradingCommand, UUID4, UUID4) {
3902 let corr = UUID4::new();
3903 let caus = UUID4::new();
3904 let mut cmd = make_batch_modify_orders(vec![make_modify_order(Some(venue_order_id()))]);
3905 cmd.correlation_id = Some(corr);
3906 cmd.causation_id = Some(caus);
3907 (TradingCommand::ModifyOrders(cmd), corr, caus)
3908 }
3909
3910 fn trading_command_cancel_order() -> (TradingCommand, UUID4, UUID4) {
3911 let corr = UUID4::new();
3912 let caus = UUID4::new();
3913 let mut cmd = make_cancel_order();
3914 cmd.correlation_id = Some(corr);
3915 cmd.causation_id = Some(caus);
3916 (TradingCommand::CancelOrder(cmd), corr, caus)
3917 }
3918
3919 fn trading_command_batch_cancel_orders() -> (TradingCommand, UUID4, UUID4) {
3920 let corr = UUID4::new();
3921 let caus = UUID4::new();
3922 let mut cmd = make_batch_cancel_orders(vec![make_cancel_order()]);
3923 cmd.correlation_id = Some(corr);
3924 cmd.causation_id = Some(caus);
3925 (TradingCommand::CancelOrders(cmd), corr, caus)
3926 }
3927
3928 fn trading_command_cancel_all_orders() -> (TradingCommand, UUID4, UUID4) {
3929 let corr = UUID4::new();
3930 let caus = UUID4::new();
3931 let mut cmd = make_cancel_all_orders();
3932 cmd.correlation_id = Some(corr);
3933 cmd.causation_id = Some(caus);
3934 (TradingCommand::CancelAllOrders(cmd), corr, caus)
3935 }
3936
3937 fn trading_command_query_order() -> (TradingCommand, UUID4, UUID4) {
3938 let corr = UUID4::new();
3939 let caus = UUID4::new();
3940 let mut cmd = make_query_order(Some(venue_order_id()));
3941 cmd.correlation_id = Some(corr);
3942 cmd.causation_id = Some(caus);
3943 (TradingCommand::QueryOrder(cmd), corr, caus)
3944 }
3945
3946 fn trading_command_query_account() -> (TradingCommand, UUID4, UUID4) {
3947 let corr = UUID4::new();
3948 let caus = UUID4::new();
3949 let mut cmd = make_query_account();
3950 cmd.correlation_id = Some(corr);
3951 cmd.causation_id = Some(caus);
3952 (TradingCommand::QueryAccount(cmd), corr, caus)
3953 }
3954
3955 #[rstest]
3956 #[case::data(data_response_data())]
3957 #[case::instrument(data_response_instrument())]
3958 #[case::instruments(data_response_instruments())]
3959 #[case::book(data_response_book())]
3960 #[case::book_deltas(data_response_book_deltas())]
3961 #[case::book_depth(data_response_book_depth())]
3962 #[case::quotes(data_response_quotes())]
3963 #[case::trades(data_response_trades())]
3964 #[case::funding_rates(data_response_funding_rates())]
3965 #[case::forward_prices(data_response_forward_prices())]
3966 #[case::bars(data_response_bars())]
3967 fn data_response_extractor_surfaces_correlation_id_for_every_variant(
3968 #[case] envelope_with_expected: (DataResponse, UUID4),
3969 ) {
3970 let (envelope, expected) = envelope_with_expected;
3974 let registry = default_registry();
3975
3976 let headers = registry
3977 .headers_for_any(&envelope as &dyn std::any::Any)
3978 .expect("registered");
3979 assert_eq!(headers.correlation_id, Some(expected));
3980 assert_eq!(headers.causation_id, None);
3981 }
3982
3983 fn data_response_data() -> (DataResponse, UUID4) {
3984 let resp = make_custom_data_response();
3985 let expected = resp.correlation_id;
3986 (DataResponse::Data(resp), expected)
3987 }
3988
3989 fn data_response_instrument() -> (DataResponse, UUID4) {
3990 let resp = make_instrument_response();
3991 let expected = resp.correlation_id;
3992 (DataResponse::Instrument(Box::new(resp)), expected)
3993 }
3994
3995 fn data_response_instruments() -> (DataResponse, UUID4) {
3996 let resp = make_instruments_response();
3997 let expected = resp.correlation_id;
3998 (DataResponse::Instruments(resp), expected)
3999 }
4000
4001 fn data_response_book() -> (DataResponse, UUID4) {
4002 let resp = make_book_response();
4003 let expected = resp.correlation_id;
4004 (DataResponse::Book(resp), expected)
4005 }
4006
4007 fn data_response_book_deltas() -> (DataResponse, UUID4) {
4008 let resp = make_book_deltas_response();
4009 let expected = resp.correlation_id;
4010 (DataResponse::BookDeltas(resp), expected)
4011 }
4012
4013 fn data_response_book_depth() -> (DataResponse, UUID4) {
4014 let resp = make_book_depth_response();
4015 let expected = resp.correlation_id;
4016 (DataResponse::BookDepth(resp), expected)
4017 }
4018
4019 fn data_response_quotes() -> (DataResponse, UUID4) {
4020 let resp = make_quotes_response();
4021 let expected = resp.correlation_id;
4022 (DataResponse::Quotes(resp), expected)
4023 }
4024
4025 fn data_response_trades() -> (DataResponse, UUID4) {
4026 let resp = make_trades_response();
4027 let expected = resp.correlation_id;
4028 (DataResponse::Trades(resp), expected)
4029 }
4030
4031 fn data_response_funding_rates() -> (DataResponse, UUID4) {
4032 let resp = make_funding_rates_response();
4033 let expected = resp.correlation_id;
4034 (DataResponse::FundingRates(resp), expected)
4035 }
4036
4037 fn data_response_forward_prices() -> (DataResponse, UUID4) {
4038 let resp = make_forward_prices_response();
4039 let expected = resp.correlation_id;
4040 (DataResponse::ForwardPrices(resp), expected)
4041 }
4042
4043 fn data_response_bars() -> (DataResponse, UUID4) {
4044 let resp = make_bars_response();
4045 let expected = resp.correlation_id;
4046 (DataResponse::Bars(resp), expected)
4047 }
4048
4049 fn make_quotes_request() -> RequestQuotes {
4050 RequestQuotes {
4051 instrument_id: InstrumentId::from("EUR/USD.SIM"),
4052 start: None,
4053 end: None,
4054 limit: None,
4055 client_id: None,
4056 request_id: UUID4::new(),
4057 ts_init: UnixNanos::default(),
4058 params: None,
4059 }
4060 }
4061}