Skip to main content

nautilus_event_store/capture/
builtins.rs

1// -------------------------------------------------------------------------------------------------
2//  Copyright (C) 2015-2026 Nautech Systems Pty Ltd. All rights reserved.
3//  https://nautechsystems.io
4//
5//  Licensed under the GNU Lesser General Public License Version 3.0 (the "License");
6//  You may not use this file except in compliance with the License.
7//  You may obtain a copy of the License at https://www.gnu.org/licenses/lgpl-3.0.en.html
8//
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14// -------------------------------------------------------------------------------------------------
15
16//! Representative encoders for the SPEC's allow-listed message surface.
17//!
18//! Phase 6 shipped a sample triple (`SubmitOrder` command, `OrderFilled` generated event,
19//! `OrderStatusReport` raw venue report) so the bus capture adapter had a working
20//! allow-list end-to-end. Phase 7 adds envelope-aware dispatchers for the
21//! wrapper enums production code actually pushes through `send_trading_command`,
22//! `publish_order_event`, `send_execution_report`, and `publish_position_event`
23//! ([`TradingCommand`], [`OrderEventAny`], [`ExecutionReport`], [`PositionEvent`]).
24//! The same pattern covers `send_data_command` and `send_data_response` ([`DataCommand`],
25//! [`DataResponse`]). These reach the bus tap as their wrapper [`std::any::TypeId`] and
26//! the bare-type registrations would miss them. Each dispatcher unwraps its variant,
27//! runs the inner-typed encode, and stamps the inner-variant's canonical `payload_type`
28//! tag so forensics scans see entries identical to the bare-type capture path.
29//!
30//! The payload serialization format is MessagePack via `rmp-serde`. The on-disk envelope
31//! uses the positional codec; MessagePack inside the payload handles the upstream Nautilus
32//! types that carry `#[serde(tag = "type")]` internal tagging, which a non-self-describing
33//! format cannot round-trip.
34
35use 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
78/// The canonical `payload_type` tag for [`SubmitOrder`].
79pub const PAYLOAD_TYPE_SUBMIT_ORDER: &str = "SubmitOrder";
80/// The canonical `payload_type` tag for [`SubmitOrderList`].
81pub const PAYLOAD_TYPE_SUBMIT_ORDER_LIST: &str = "SubmitOrderList";
82/// The canonical `payload_type` tag for [`ModifyOrder`].
83pub const PAYLOAD_TYPE_MODIFY_ORDER: &str = "ModifyOrder";
84/// The canonical `payload_type` tag for [`BatchModifyOrders`].
85pub const PAYLOAD_TYPE_BATCH_MODIFY_ORDERS: &str = "BatchModifyOrders";
86/// The canonical `payload_type` tag for [`CancelOrder`].
87pub const PAYLOAD_TYPE_CANCEL_ORDER: &str = "CancelOrder";
88/// The canonical `payload_type` tag for [`CancelAllOrders`].
89pub const PAYLOAD_TYPE_CANCEL_ALL_ORDERS: &str = "CancelAllOrders";
90/// The canonical `payload_type` tag for [`BatchCancelOrders`].
91pub const PAYLOAD_TYPE_BATCH_CANCEL_ORDERS: &str = "BatchCancelOrders";
92/// The canonical `payload_type` tag for [`QueryOrder`].
93pub const PAYLOAD_TYPE_QUERY_ORDER: &str = "QueryOrder";
94/// The canonical `payload_type` tag for [`QueryAccount`].
95pub const PAYLOAD_TYPE_QUERY_ACCOUNT: &str = "QueryAccount";
96
97/// The canonical `payload_type` tag for [`OrderInitialized`].
98pub const PAYLOAD_TYPE_ORDER_INITIALIZED: &str = "OrderInitialized";
99/// The canonical `payload_type` tag for [`OrderDenied`].
100pub const PAYLOAD_TYPE_ORDER_DENIED: &str = "OrderDenied";
101/// The canonical `payload_type` tag for [`OrderEmulated`].
102pub const PAYLOAD_TYPE_ORDER_EMULATED: &str = "OrderEmulated";
103/// The canonical `payload_type` tag for [`OrderReleased`].
104pub const PAYLOAD_TYPE_ORDER_RELEASED: &str = "OrderReleased";
105/// The canonical `payload_type` tag for [`OrderSubmitted`].
106pub const PAYLOAD_TYPE_ORDER_SUBMITTED: &str = "OrderSubmitted";
107/// The canonical `payload_type` tag for [`OrderAccepted`].
108pub const PAYLOAD_TYPE_ORDER_ACCEPTED: &str = "OrderAccepted";
109/// The canonical `payload_type` tag for [`OrderRejected`].
110pub const PAYLOAD_TYPE_ORDER_REJECTED: &str = "OrderRejected";
111/// The canonical `payload_type` tag for [`OrderCanceled`].
112pub const PAYLOAD_TYPE_ORDER_CANCELED: &str = "OrderCanceled";
113/// The canonical `payload_type` tag for [`OrderExpired`].
114pub const PAYLOAD_TYPE_ORDER_EXPIRED: &str = "OrderExpired";
115/// The canonical `payload_type` tag for [`OrderTriggered`].
116pub const PAYLOAD_TYPE_ORDER_TRIGGERED: &str = "OrderTriggered";
117/// The canonical `payload_type` tag for [`OrderPendingUpdate`].
118pub const PAYLOAD_TYPE_ORDER_PENDING_UPDATE: &str = "OrderPendingUpdate";
119/// The canonical `payload_type` tag for [`OrderPendingCancel`].
120pub const PAYLOAD_TYPE_ORDER_PENDING_CANCEL: &str = "OrderPendingCancel";
121/// The canonical `payload_type` tag for [`OrderModifyRejected`].
122pub const PAYLOAD_TYPE_ORDER_MODIFY_REJECTED: &str = "OrderModifyRejected";
123/// The canonical `payload_type` tag for [`OrderCancelRejected`].
124pub const PAYLOAD_TYPE_ORDER_CANCEL_REJECTED: &str = "OrderCancelRejected";
125/// The canonical `payload_type` tag for [`OrderUpdated`].
126pub const PAYLOAD_TYPE_ORDER_UPDATED: &str = "OrderUpdated";
127/// The canonical `payload_type` tag for [`OrderFilled`].
128pub const PAYLOAD_TYPE_ORDER_FILLED: &str = "OrderFilled";
129/// The canonical `payload_type` tag for [`OrderFillVoided`].
130pub const PAYLOAD_TYPE_ORDER_FILL_VOIDED: &str = "OrderFillVoided";
131/// The canonical `payload_type` tag for [`OrderStatusReport`].
132pub const PAYLOAD_TYPE_ORDER_STATUS_REPORT: &str = "OrderStatusReport";
133/// The canonical `payload_type` tag for [`FillReport`].
134pub const PAYLOAD_TYPE_FILL_REPORT: &str = "FillReport";
135/// The canonical `payload_type` tag for the [`ExecutionReport::OrderWithFills`] bundle.
136pub const PAYLOAD_TYPE_ORDER_WITH_FILLS: &str = "OrderWithFills";
137/// The canonical `payload_type` tag for [`PositionStatusReport`].
138pub const PAYLOAD_TYPE_POSITION_STATUS_REPORT: &str = "PositionStatusReport";
139/// The canonical `payload_type` tag for [`ExecutionMassStatus`].
140pub const PAYLOAD_TYPE_EXECUTION_MASS_STATUS: &str = "ExecutionMassStatus";
141/// The canonical `payload_type` tag for [`PositionOpened`].
142pub const PAYLOAD_TYPE_POSITION_OPENED: &str = "PositionOpened";
143/// The canonical `payload_type` tag for [`PositionChanged`].
144pub const PAYLOAD_TYPE_POSITION_CHANGED: &str = "PositionChanged";
145/// The canonical `payload_type` tag for [`PositionClosed`].
146pub const PAYLOAD_TYPE_POSITION_CLOSED: &str = "PositionClosed";
147/// The canonical `payload_type` tag for [`PositionAdjusted`].
148pub const PAYLOAD_TYPE_POSITION_ADJUSTED: &str = "PositionAdjusted";
149/// The canonical `payload_type` tag for [`AccountState`].
150pub const PAYLOAD_TYPE_ACCOUNT_STATE: &str = "AccountState";
151/// The canonical `payload_type` tag for [`TimeEvent`].
152pub const PAYLOAD_TYPE_TIME_EVENT: &str = "TimeEvent";
153
154/// The canonical `payload_type` tag for `RequestCommand`.
155pub const PAYLOAD_TYPE_REQUEST_COMMAND: &str = "RequestCommand";
156/// The canonical `payload_type` tag for `SubscribeCommand`.
157pub const PAYLOAD_TYPE_SUBSCRIBE_COMMAND: &str = "SubscribeCommand";
158/// The canonical `payload_type` tag for `UnsubscribeCommand`.
159pub const PAYLOAD_TYPE_UNSUBSCRIBE_COMMAND: &str = "UnsubscribeCommand";
160#[cfg(feature = "defi")]
161/// The canonical `payload_type` tag for `DefiRequestCommand`.
162pub const PAYLOAD_TYPE_DEFI_REQUEST_COMMAND: &str = "DefiRequestCommand";
163#[cfg(feature = "defi")]
164/// The canonical `payload_type` tag for `DefiSubscribeCommand`.
165pub const PAYLOAD_TYPE_DEFI_SUBSCRIBE_COMMAND: &str = "DefiSubscribeCommand";
166#[cfg(feature = "defi")]
167/// The canonical `payload_type` tag for `DefiUnsubscribeCommand`.
168pub const PAYLOAD_TYPE_DEFI_UNSUBSCRIBE_COMMAND: &str = "DefiUnsubscribeCommand";
169
170/// The canonical `payload_type` tag for [`CustomDataResponse`].
171pub const PAYLOAD_TYPE_CUSTOM_DATA_RESPONSE: &str = "CustomDataResponse";
172/// The canonical `payload_type` tag for [`InstrumentResponse`].
173pub const PAYLOAD_TYPE_INSTRUMENT_RESPONSE: &str = "InstrumentResponse";
174/// The canonical `payload_type` tag for [`InstrumentsResponse`].
175pub const PAYLOAD_TYPE_INSTRUMENTS_RESPONSE: &str = "InstrumentsResponse";
176/// The canonical `payload_type` tag for [`BookResponse`].
177pub const PAYLOAD_TYPE_BOOK_RESPONSE: &str = "BookResponse";
178/// The canonical `payload_type` tag for [`BookDeltasResponse`].
179pub const PAYLOAD_TYPE_BOOK_DELTAS_RESPONSE: &str = "BookDeltasResponse";
180/// The canonical `payload_type` tag for [`BookDepthResponse`].
181pub const PAYLOAD_TYPE_BOOK_DEPTH_RESPONSE: &str = "BookDepthResponse";
182/// The canonical `payload_type` tag for [`QuotesResponse`].
183pub const PAYLOAD_TYPE_QUOTES_RESPONSE: &str = "QuotesResponse";
184/// The canonical `payload_type` tag for [`TradesResponse`].
185pub const PAYLOAD_TYPE_TRADES_RESPONSE: &str = "TradesResponse";
186/// The canonical `payload_type` tag for [`FundingRatesResponse`].
187pub const PAYLOAD_TYPE_FUNDING_RATES_RESPONSE: &str = "FundingRatesResponse";
188/// The canonical `payload_type` tag for [`ForwardPricesResponse`].
189pub const PAYLOAD_TYPE_FORWARD_PRICES_RESPONSE: &str = "ForwardPricesResponse";
190/// The canonical `payload_type` tag for [`BarsResponse`].
191pub const PAYLOAD_TYPE_BARS_RESPONSE: &str = "BarsResponse";
192
193// Wrapper-level fallback tag reached only when a dispatcher returns an
194// `EncodedPayload` without an override. Every current variant stamps its own
195// inner tag, so this is a sentinel for a future variant that forgets the
196// override rather than a tag the writer is expected to commit.
197const 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/// Returns an [`EncoderRegistry`] preloaded with the default encoders.
271///
272/// Callers can extend the returned registry with additional encoders before constructing
273/// the [`crate::capture::BusCaptureAdapter`].
274#[must_use]
275pub fn default_registry() -> EncoderRegistry {
276    let mut registry = EncoderRegistry::new();
277    register_default(&mut registry);
278    registry
279}
280
281/// Adds the default encoders to `registry`.
282///
283/// The bare-type registrations remain for capture sites that already submit the inner
284/// type directly (the kernel's `RunStarted` path and a few internal tests). The envelope
285/// registrations are what production bus traffic actually hits: `send_trading_command`
286/// reaches the tap as [`TradingCommand`], `publish_order_event` reaches it as
287/// [`OrderEventAny`], `send_execution_report` reaches it as [`ExecutionReport`],
288/// `publish_position_event` reaches it as [`PositionEvent`], and `send_data_response`
289/// reaches it as [`DataResponse`]. Without these wrapper-aware dispatchers the tap looks
290/// up the wrapper's [`std::any::TypeId`], finds no encoder, and silently drops the
291/// capture.
292///
293/// [`AccountState`] is registered as a bare type: `publish_account_state` and
294/// `send_account_state` both reach the tap as the same `AccountState` `TypeId`, so a
295/// single registration covers both dispatch paths.
296///
297/// [`OrderStatusReport`], [`FillReport`], and [`PositionStatusReport`] are registered
298/// as bare types because the execution engine publishes raw venue reports through
299/// `publish_any` on `reconciliation.raw.*` topics before any state mutation. The
300/// bare-type registration is what captures those raw inputs for forensic replay.
301pub 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
347/// Attaches header extractors for every type that carries `correlation_id` or
348/// `causation_id` today.
349///
350/// Header propagation lands incrementally per the SPEC's workstream A: extractors only
351/// exist for types whose underlying struct has grown the fields. Other types fall back to
352/// the registry's no-op extractor, which yields [`Headers::empty`]; capture still works
353/// for them, the entry just carries no correlation metadata until the field set arrives.
354fn 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
369/// Attaches identity extractors for the types production dispatch pushes through more
370/// than one tap-visible boundary (portfolio endpoint send plus strategy topic publish,
371/// command hops through risk to execution, account states on both dispatch paths, data
372/// commands through the queue endpoint and the drained execute endpoint), so the
373/// adapter captures each logical message exactly once. The venue report types
374/// deliberately carry no extractor: the raw `reconciliation.raw.*` publish and the
375/// engine-bound dispatch are distinct capture boundaries.
376fn 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        // `DataCommand` is `#[non_exhaustive]`; future variants capture per dispatch
397        _ => 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
480// `send_trading_command` reaches the bus tap with the wrapper's `TypeId`, so the
481// extractor must mirror the encoder's variant dispatch to surface the inner command's
482// correlation metadata on the captured entry.
483fn 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
497// `send_data_command` reaches the bus tap with the wrapper's `TypeId`. The data engine
498// keys RPC request/response pairs by the request's `request_id`: the response's
499// `correlation_id` echoes that uuid back. Surfacing `request_id` as the captured entry's
500// `correlation_id` therefore lines a request entry up with its eventual response entry
501// under the same chain key. Subscribe / Unsubscribe variants carry an explicit
502// `correlation_id` field, which we forward as-is. DeFi variants are not yet wired through
503// header propagation.
504fn 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        // `DataCommand` is `#[non_exhaustive]` and the defi variants do not yet carry
510        // header propagation; future variants drop through this arm with empty headers
511        // until their correlation field shape lands.
512        _ => Headers::empty(),
513    }
514}
515
516// Every `DataResponse` variant carries a required `correlation_id` that pairs the
517// response with its originating request; the captured entry mirrors that value.
518fn 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
532/// Encodes a [`SubmitOrder`] command into canonical bytes plus its `client_order_id` index.
533///
534/// # Errors
535///
536/// Returns [`EncodeError::Serialize`] when MessagePack rejects the payload (a malformed
537/// value the type system should make unrepresentable; surfaced rather than swallowed
538/// because the audit contract refuses to drop captured commands).
539pub 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
548/// Encodes an [`OrderFilled`] event into canonical bytes plus its `client_order_id` and
549/// `venue_order_id` indices.
550///
551/// # Errors
552///
553/// Returns [`EncodeError::Serialize`] when MessagePack rejects the payload.
554pub 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
566/// Encodes a [`TradingCommand`] envelope by dispatching on the variant.
567///
568/// The captured entry's `payload_type` matches the inner-variant tag (e.g. `SubmitOrder`
569/// rather than `TradingCommand`) so forensics scans pair with the same decoder as the
570/// bare-type capture path. The serialized payload is the inner variant; the wrapper enum
571/// is never written to disk.
572///
573/// # Errors
574///
575/// Returns the inner encoder's [`EncodeError`] for the [`TradingCommand::SubmitOrder`]
576/// variant; other variants return [`EncodeError::Serialize`] when MessagePack rejects the
577/// inner payload.
578pub 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
594/// Encodes an [`OrderEventAny`] envelope by dispatching on the inner variant.
595///
596/// The captured entry's `payload_type` matches the inner-variant tag (e.g. `OrderFilled`
597/// rather than `OrderEventAny`); the serialized payload is the inner variant.
598///
599/// # Errors
600///
601/// Returns the inner encoder's [`EncodeError`] for the [`OrderEventAny::Filled`] variant;
602/// other variants return [`EncodeError::Serialize`] when MessagePack rejects the inner
603/// payload.
604pub 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
626/// Encodes an [`ExecutionReport`] envelope by dispatching on the variant.
627///
628/// `send_execution_report` hands the bus tap an [`ExecutionReport`] wrapper, so the tap
629/// dispatches by the wrapper's [`std::any::TypeId`] and the inner variants never reach
630/// their bare-type encoders. The dispatcher unwraps each variant, encodes the inner type
631/// with its own index keys, and stamps the inner-variant tag so forensics scans see
632/// entries identical to a bare capture path.
633///
634/// The [`ExecutionReport::Order`] arm reuses [`encode_order_status_report`] because the
635/// bare-type encoder already exists; the remaining variants delegate to private inner
636/// encoders that index the report's identifiers individually.
637///
638/// # Errors
639///
640/// Returns the inner encoder's [`EncodeError`] when MessagePack rejects the inner
641/// payload.
642pub 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
655/// Encodes a [`FillReport`] into canonical bytes plus its `venue_order_id` index and,
656/// when present, its `client_order_id` index.
657///
658/// # Errors
659///
660/// Returns [`EncodeError::Serialize`] when MessagePack rejects the payload.
661pub 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
683/// Encodes a [`PositionStatusReport`] into canonical bytes with no sidecar indices.
684///
685/// `PositionStatusReport` carries only `AccountId`, `InstrumentId`, and `PositionId`;
686/// none of those have a matching [`IndexKind`] variant today. Capture with no sidecar
687/// indices so the entry is forensics-discoverable by sequential scan rather than
688/// synthesising an index against an identifier the reader cannot query.
689///
690/// # Errors
691///
692/// Returns [`EncodeError::Serialize`] when MessagePack rejects the payload.
693pub 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    // PositionStatusReport identifiers are not indexable today, see
810    // `encode_position_status_report`.
811    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
829/// Encodes a [`PositionEvent`] envelope by dispatching on the variant.
830///
831/// `publish_position_event` hands the bus tap a [`PositionEvent`] wrapper, so the tap
832/// dispatches by the wrapper's [`std::any::TypeId`] and the inner variants never reach
833/// their bare-type encoders. The dispatcher unwraps each variant, encodes the inner
834/// struct, and stamps the inner-variant tag so forensics scans see entries identical
835/// to the bare-type capture path.
836///
837/// # Errors
838///
839/// Returns the inner encoder's [`EncodeError`] when MessagePack rejects the inner
840/// payload.
841pub 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    // Opening and closing client_order_ids are distinct in normal operation; dedup
889    // guards the rare case where a single order both opens and closes the position.
890    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    // PositionAdjusted carries no client_order_id; identifiers are PositionId,
908    // AccountId, and InstrumentId, none of which have a matching IndexKind today.
909    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
1185/// Encodes an [`OrderStatusReport`] into canonical bytes plus its `venue_order_id` index
1186/// and, when present, its `client_order_id` index.
1187///
1188/// External orders observed only at the venue may not carry a `client_order_id`; the
1189/// index is omitted in that case so the secondary index never records an empty key.
1190///
1191/// # Errors
1192///
1193/// Returns [`EncodeError::Serialize`] when MessagePack rejects the payload.
1194pub 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
1213/// Encodes an [`AccountState`] into canonical bytes with no sidecar indices.
1214///
1215/// `AccountState` carries `AccountId` and `event_id` (UUID4); neither matches an
1216/// [`IndexKind`] variant today, so the encoder emits no sidecar keys and forensics
1217/// scans rely on sequential range over `seq`. This mirrors the [`PositionStatusReport`]
1218/// precedent.
1219///
1220/// # Errors
1221///
1222/// Returns [`EncodeError::Serialize`] when MessagePack rejects the payload.
1223pub 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
1236/// Encodes a fired [`TimeEvent`] into canonical bytes with no sidecar indices.
1237///
1238/// Time events carry a callback boundary rather than a cache-state key. The event store
1239/// captures them for forensic ordering and deterministic replay inputs, while cache
1240/// replay leaves clock re-arming to the later clock lifecycle event workstream.
1241///
1242/// # Errors
1243///
1244/// Returns [`EncodeError::Serialize`] when MessagePack rejects the payload.
1245pub 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
1255/// Encodes a [`DataCommand`] envelope by dispatching on its command category.
1256///
1257/// `send_data_command` hands the bus tap a [`DataCommand`] wrapper, so the tap
1258/// dispatches by the wrapper's [`std::any::TypeId`] and the inner command category
1259/// never reaches a bare-type encoder. The dispatcher unwraps the category, encodes the
1260/// serializable inner enum (`RequestCommand`, `SubscribeCommand`, or
1261/// `UnsubscribeCommand`), and stamps that category's canonical `payload_type` tag.
1262///
1263/// Request IDs, command IDs, and correlation IDs do not have a matching [`IndexKind`]
1264/// today, so data commands emit no sidecar indices. Correlation is recovered from the
1265/// captured payload and, once header propagation lands, propagated headers.
1266///
1267/// # Errors
1268///
1269/// Returns [`EncodeError::Serialize`] when MessagePack rejects the inner payload, or
1270/// when a future non-exhaustive [`DataCommand`] variant has no encoder yet.
1271pub 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
1312/// Encodes a [`DataResponse`] envelope by dispatching on the variant.
1313///
1314/// `send_data_response` hands the bus tap a [`DataResponse`] wrapper, so the tap
1315/// dispatches by the wrapper's [`std::any::TypeId`] and the inner variants never reach
1316/// their bare-type encoders. The dispatcher unwraps each variant, encodes the inner
1317/// struct, and stamps the inner-variant tag so forensics scans see entries identical
1318/// to a bare-type capture path.
1319///
1320/// Each variant carries a `correlation_id` (UUID4) pairing the response with the
1321/// originating `RequestCommand::request_id`. [`IndexKind`] has no matching variant
1322/// today, so every variant emits zero sidecar indices, mirroring the
1323/// [`PositionStatusReport`] and [`AccountState`] precedents.
1324///
1325/// The [`DataResponse::Data`] and [`DataResponse::Book`] variants carry payloads that
1326/// are not directly serializable: [`CustomDataResponse`] holds an `Arc<dyn Any>` and
1327/// [`BookResponse`] holds a [`nautilus_model::orderbook::OrderBook`] without serde
1328/// derives. The dispatcher serializes the audit-relevant metadata for those two
1329/// variants via local borrowed wrapper structs (the `encode_order_with_fills`
1330/// precedent) and omits the opaque payload. `BookResponse` is state-affecting on
1331/// the data engine path (`handle_book_response` clones the book into the cache); a
1332/// follow-up that adds serde to `OrderBook`/`BookLadder` can replace the metadata
1333/// wrapper with full payload capture without changing the dispatcher contract.
1334///
1335/// # Errors
1336///
1337/// Returns the inner encoder's [`EncodeError`] when MessagePack rejects the inner
1338/// payload.
1339pub 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    // `data: Arc<dyn Any + Send + Sync>` is type-erased at the dispatcher; the
1359    // payload is captured via a metadata-only wrapper so the audit entry pairs
1360    // with the originating request without depending on per-registration
1361    // serializers for the inner Any payload.
1362    #[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    // `data: OrderBook` is not serde-derived today (BookLadder/BookLevel chain), and
1415    // the full book state is rarely the audit value at this level. Capture
1416    // response-level metadata via a borrowed wrapper so the entry pairs with the
1417    // originating request.
1418    #[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, // correlation_id
1597        )
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        // Production pushes every queued data command through two tapped sends
1764        // (queue, then drained execute); the identity must key both hops to the
1765        // same command.
1766        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, // correlation_id
1842        )
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, // correlation_id
1854        )
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, // correlation_id
1882        )
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, // correlation_id
1896        )
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, // correlation_id
1910        )
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, // correlation_id
1925        )
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, // correlation_id
1939        )
1940    }
1941
1942    fn make_submit_order_list(client_order_ids: Vec<ClientOrderId>) -> SubmitOrderList {
1943        // OrderList::new asserts that order_inits' client_order_ids match the list's,
1944        // so we mint one OrderInitialized per id with the matching client_order_id.
1945        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, // correlation_id
1969        )
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        // TradingCommand reaches the bus tap as the wrapper TypeId; the dispatcher must
2103        // unwrap to SubmitOrder, produce the same bytes and indices as the bare-type
2104        // encoder, and stamp the inner payload_type so forensics scans pair the entry
2105        // with the SubmitOrder decoder.
2106        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        // QueryAccount carries no client_order_id or venue_order_id; the dispatcher
2149        // must not invent empty index keys.
2150        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    // Walks every TradingCommand variant and asserts the dispatcher stamps the inner-variant
2198    // payload_type tag and emits the expected number of index keys. Catches a swapped match
2199    // arm or a forgotten `with_payload_type` override that would otherwise fall back to the
2200    // wrapper sentinel tag.
2201    #[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    // Walks every OrderEventAny variant. Builds each variant from `Default::default()`
2269    // (gated by the `test-support` feature on `nautilus-model`) with the test client_order_id
2270    // patched in, so the assertion can verify the index value alongside the tag.
2271    #[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    // Optional venue_order_id branch coverage for variants that carry one. None must skip
2341    // the VenueOrderId index entirely; Some must push it second.
2342    #[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        // Two cancels: one with a venue_order_id (contributes 2 indices), one without
2398        // (contributes 1). The dispatcher must preserve both children's identifiers in
2399        // the same order they appear in the batch.
2400        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        // SubmitOrderList carries N orders; the dispatcher must emit a ClientOrderId
2464        // index per child so forensics can resolve any of the list's intents to the
2465        // captured seq.
2466        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        // ExecutionReport::Order maps onto the existing OrderStatusReport bare-type
2536        // encoder; the dispatcher must produce identical bytes and indices and stamp
2537        // the OrderStatusReport tag so forensics scans pair the entry with the same
2538        // decoder as the bare-type capture path.
2539        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        // PositionStatusReport identifiers (AccountId, InstrumentId, PositionId) have
2590        // no matching IndexKind today; the dispatcher must not invent sidecar indices
2591        // pointing at an identifier the reader cannot query.
2592        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        // The order and its fills semantically carry the same venue_order_id and
2610        // client_order_id, so the dispatcher must dedupe rather than emit a duplicate
2611        // (kind, key) pair the backend would silently drop.
2612        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        // A bundled OrderWithFills can carry a fill whose client_order_id differs
2634        // from the order's (rare, but real: external orders observed via a fill
2635        // before the venue confirms the canonical id). The dispatcher must index
2636        // both client_order_ids so forensics can resolve either to the same seq.
2637        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        // One order with venue+client ids, one fill sharing both ids, one position
2680        // (unindexable). The dispatcher must dedupe the shared ids.
2681        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        // Two distinct orders + a fill for a third venue_order_id with its own
2694        // client_order_id. The dispatcher must record an index for each unique id
2695        // so forensics can resolve any child of the status report.
2696        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        // Three venue_order_ids + three client_order_ids = 6 distinct keys
2721        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    // Walks every ExecutionReport variant and asserts the dispatcher stamps the
2743    // inner-variant tag. Catches a swapped match arm or a forgotten override that
2744    // would otherwise fall back to the wrapper sentinel tag.
2745    #[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        // Rare but real: a single order both opens and closes the position (e.g.,
2977        // reduce-only fills against a stale position). The dispatcher must dedupe
2978        // rather than insert the same (kind, key) twice.
2979        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        // PositionAdjusted has no ClientOrderId field; PositionId/AccountId/
2991        // InstrumentId have no matching IndexKind today, so the dispatcher must
2992        // not invent sidecar indices pointing at an identifier the reader cannot
2993        // query.
2994        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        // AccountState carries AccountId and event_id (UUID4); neither matches an
3060        // IndexKind variant today. The encoder must capture the payload without
3061        // synthesising sidecar indices pointing at identifiers the reader cannot
3062        // query, mirroring the PositionStatusReport precedent.
3063        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        // The default registry must dispatch AccountState through encode_account_state
3086        // and stamp PAYLOAD_TYPE_ACCOUNT_STATE so both `publish_account_state` and
3087        // `send_account_state` capture under the same canonical tag.
3088        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        // DataCommand reaches the bus tap as the wrapper TypeId. The dispatcher must
3152        // unwrap one level, encode RequestCommand, and stamp the category tag so the
3153        // reader pairs the bytes with the RequestCommand decoder.
3154        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        // The default registry must dispatch DataCommand through encode_data_command and
3221        // stamp the inner category tag, not the wrapper sentinel tag.
3222        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        // CustomDataResponse holds `data: Arc<dyn Any>`; the dispatcher captures
3551        // metadata via a borrowed wrapper. Correlation_id has no matching IndexKind
3552        // today, so the encoder emits no sidecar keys.
3553        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        // The CustomDataResponseRef wrapper omits `data` (Arc<dyn Any>); the audit
3568        // entry captures correlation_id, client_id, venue, data_type, timing, and
3569        // params so forensics can pair the response with its request.
3570        #[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        // BookResponseRef omits `data: OrderBook` (not serde-derived). The audit
3607        // entry captures the correlation and addressing metadata.
3608        #[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    // Walks every DataResponse variant and asserts the dispatcher stamps the
3707    // inner-variant tag and emits zero sidecar indices. Catches a swapped match
3708    // arm or a forgotten `with_payload_type` override that would otherwise fall
3709    // back to the wrapper sentinel tag.
3710    #[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        // The default registry must dispatch DataResponse through encode_data_response
3770        // and stamp the inner-variant tag, so `send_data_response` captures under the
3771        // same canonical tag as a bare-type capture path.
3772        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        // The data engine pairs RPC requests and responses by correlation_id (see
3786        // `crates/data/src/engine/mod.rs` `send_response`). Captured DataResponse entries
3787        // must carry that uuid in `Headers::correlation_id` so forensics can join a
3788        // captured response to its request.
3789        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        // The request_id of an outbound RequestCommand IS the chain root: the eventual
3804        // DataResponse echoes the same uuid back as its correlation_id. Surfacing
3805        // request_id in `Headers::correlation_id` lines the request entry up with its
3806        // response entry under one chain key.
3807        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        // Subscribe variants carry an optional correlation_id field; the extractor
3821        // must forward whatever the inner command reports so captured subscribe
3822        // traffic joins its acknowledgements under one chain key.
3823        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        // The TradingCommand envelope dispatch must route every variant to the matching
3861        // per-type extractor and forward both correlation_id and causation_id intact.
3862        // A swap of args inside any extract_*_headers helper or a misrouted wrapper arm
3863        // is caught by exercising each variant with distinct populated values.
3864        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        // Every DataResponse variant carries a required correlation_id paired with its
3971        // originating request; the extractor must forward that uuid intact regardless
3972        // of which variant is captured.
3973        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}