Skip to main content

perpl_sdk/stream/
trade.rs

1use std::{collections::HashMap, num::NonZeroU16};
2
3use alloy::{eips::BlockId, primitives::U256, providers::Provider};
4use futures::{Stream, StreamExt};
5
6use crate::{
7    Chain,
8    abi::dex::Exchange::{ExchangeEvents, ExchangeInstance},
9    error::DexError,
10    num, state, types,
11};
12
13pub type TradeEvent = types::EventContext<types::Trade>;
14pub type BlockTrades = types::BlockEvents<TradeEvent>;
15
16/// Returns stream of normalized trade events aggregated from the [`super::raw`]
17/// event stream, batched per block.
18///
19/// Listens to `MakerOrderFilledV2` and `TakerOrderFilledV2` events (and their
20/// V1 predecessors when replaying pre-v1.1.7.4 history), batches all maker
21/// fills per taker into unified `Trade` events, normalizes fixed-point values
22/// to decimals and recovers builder attribution from the order requests.
23///
24/// # Safety note
25///
26/// The returned stream is not cancellation-safe and should not be used within
27/// `select!`.
28///
29/// # Architecture
30///
31/// The module separates pure processing logic from async I/O:
32///
33/// - [`TradeProcessor`] - Pure, synchronous trade extraction from raw events
34/// - [`NormalizationConfig`] - Configuration fetched once at startup
35///
36/// # Data Model
37///
38/// Each [`TradeEvent`] represents a single taker order execution that may have
39/// matched against multiple maker orders. The `maker_fills` vector contains
40/// all individual [`types::MakerFill`]s that occurred as part of this trade.
41///
42/// # Example
43///
44/// ```ignore
45/// use perpl_sdk::{Chain, stream, types::StateInstant};
46///
47/// let chain = Chain::testnet();
48/// let provider = /* setup provider */;
49/// let from = StateInstant::new(latest_block, timestamp);
50///
51/// let raw_stream = stream::raw(
52///     &chain,
53///     provider.clone(),
54///     types::StateInstant::new(block_num, 0),
55///     tokio::time::sleep,
56/// );
57/// let mut trade_stream = pin!(stream::trade(&chain, provider, raw_stream).await.unwrap());
58///
59/// while let Some(Ok(block_trades)) = trade_stream.next().await {
60///     if !block_trades.events().is_empty() {
61///         println!(
62///             "Block {} - {} trade(s):",
63///             block_trades.instant().block_number(),
64///             block_trades.events().len()
65///         );
66///         for event in block_trades.events() {
67///             let trade = event.event();
68///             println!(
69///                 "  Taker {} {:?} {} @ {} on perp={} (fee: {})",
70///                 trade.taker_account_id,
71///                 trade.taker_side,
72///                 trade.total_size(),
73///                 trade.avg_price().unwrap_or_default(),
74///                 trade.perpetual_id,
75///                 trade.taker_fee,
76///             );
77///             for fill in &trade.maker_fills {
78///                 println!(
79///                     "    <- Maker {} order {} filled {} @ {} (fee: {})",
80///                     fill.maker_account_id, fill.maker_order_id, fill.size, fill.price, fill.fee,
81///                 );
82///             }
83///         }
84///     }
85/// }
86/// ```
87pub async fn trade<P>(
88    chain: &Chain,
89    provider: P,
90    raw_events: impl Stream<Item = Result<super::RawBlockEvents, DexError>>,
91) -> Result<impl Stream<Item = Result<BlockTrades, DexError>>, DexError>
92where
93    P: Provider + Clone,
94{
95    // Fetch normalization config
96    let config = NormalizationConfig::fetch(chain, &provider).await?;
97    // Setup trade processor
98    let mut processor = TradeProcessor::new(config);
99
100    let stream = raw_events.map(move |block_result| {
101        block_result.map(|block_events| processor.process_block(&block_events))
102    });
103
104    Ok(stream)
105}
106
107/// Configuration for normalization.
108#[derive(Clone)]
109pub struct NormalizationConfig {
110    collateral_converter: num::Converter,
111    perpetuals: HashMap<types::PerpetualId, PerpetualConverters>,
112}
113
114/// Converters for a single perpetual.
115#[derive(Clone, Copy)]
116struct PerpetualConverters {
117    price_converter: num::Converter,
118    size_converter: num::Converter,
119}
120
121/// Context for tracking order requests (reuses pattern from exchange.rs).
122struct OrderContext {
123    perpetual_id: types::PerpetualId,
124    account_id: types::AccountId,
125    request_id: types::RequestId,
126    side: types::OrderSide,
127    builder: Option<types::BuilderAttribution>,
128}
129
130/// Pending maker fill waiting for taker match.
131struct PendingMakerFill {
132    tx_hash: alloy::primitives::TxHash,
133    log_index: u64,
134    perpetual_id: types::PerpetualId,
135    maker_account_id: types::AccountId,
136    maker_order_id: types::OrderId,
137    maker_client_order_id: Option<types::RequestId>,
138    maker_builder: Option<types::BuilderAttribution>,
139    price: fastnum::UD64,
140    size: fastnum::UD64,
141    maker_fee: fastnum::UD64,
142    maker_builder_fee: fastnum::UD64,
143}
144
145/// Raw maker fill data, common to the V1 and V2 `MakerOrderFilled*` events.
146struct RawMakerFill {
147    perp_id: U256,
148    account_id: U256,
149    order_id: U256,
150    price_pns: U256,
151    lot_lns: U256,
152    fee_cns: U256,
153    builder_fee_cns: U256,
154}
155
156/// Trade processor - pure logic, no async.
157pub struct TradeProcessor {
158    config: NormalizationConfig,
159    order_context: Option<OrderContext>,
160    // Entries are retained after orders close and overwritten on ID reuse, so this can
161    // hold up to 65,535 entries per configured perpetual.
162    maker_orders: HashMap<(types::PerpetualId, types::OrderId), PlacedOrder>,
163    pending_maker_fills: Vec<PendingMakerFill>,
164    prev_tx_index: Option<u64>,
165}
166
167/// Attribution of an order observed being placed, recovered at fill time.
168#[derive(Clone, Copy)]
169struct PlacedOrder {
170    client_order_id: types::RequestId,
171    builder: Option<types::BuilderAttribution>,
172}
173
174impl TradeProcessor {
175    /// Create a new trade processor with the given normalization config.
176    pub fn new(config: NormalizationConfig) -> Self {
177        Self {
178            config,
179            order_context: None,
180            maker_orders: HashMap::new(),
181            pending_maker_fills: Vec::new(),
182            prev_tx_index: None,
183        }
184    }
185
186    /// Process a block of raw events and extract trades.
187    ///
188    /// This is pure logic - no async, no I/O.
189    pub fn process_block(&mut self, events: &super::RawBlockEvents) -> BlockTrades {
190        let mut trades = Vec::new();
191
192        for event in events.events() {
193            // Reset context at transaction boundary (pattern from exchange.rs)
194            if self.prev_tx_index.is_some_and(|idx| idx < event.tx_index()) {
195                self.order_context.take();
196                self.pending_maker_fills.clear();
197            }
198
199            if let Some(trade) = self.process_event(event) {
200                trades.push(trade);
201            }
202
203            self.prev_tx_index = Some(event.tx_index());
204        }
205
206        BlockTrades::new(events.instant(), trades)
207    }
208
209    /// Process a single event, potentially emitting a trade.
210    fn process_event(&mut self, event: &super::RawEvent) -> Option<TradeEvent> {
211        match event.event() {
212            // V1 order/fill events are never emitted by contract v1.1.7.4+, but
213            // stay handled to keep historical replay working
214            ExchangeEvents::OrderRequest(e) => {
215                self.track_order_request(e.perpId, e.accountId, e.orderDescId, e.orderType, None);
216                None
217            },
218            ExchangeEvents::OrderRequestV2(e) => {
219                self.track_order_request(
220                    e.perpId,
221                    e.accountId,
222                    e.orderDescId,
223                    e.orderType,
224                    types::BuilderAttribution::decode(&e.extension)
225                        .ok()
226                        .flatten(),
227                );
228                None
229            },
230            ExchangeEvents::OrderBatchCompleted(_) => {
231                self.order_context.take();
232                self.pending_maker_fills.clear();
233                None
234            },
235            ExchangeEvents::OrderPlaced(e) => {
236                if let Some(context) = self.order_context.as_ref()
237                    && self.config.perpetuals.contains_key(&context.perpetual_id)
238                    && let Some(order_id) = NonZeroU16::new(e.orderId.to())
239                {
240                    self.maker_orders.insert(
241                        (context.perpetual_id, order_id),
242                        PlacedOrder {
243                            client_order_id: context.request_id,
244                            builder: context.builder,
245                        },
246                    );
247                }
248                None
249            },
250            ExchangeEvents::MakerOrderFilled(e) => {
251                self.handle_maker_fill(
252                    event,
253                    RawMakerFill {
254                        perp_id: e.perpId,
255                        account_id: e.accountId,
256                        order_id: e.orderId,
257                        price_pns: e.pricePNS,
258                        lot_lns: e.lotLNS,
259                        fee_cns: e.feeCNS,
260                        builder_fee_cns: U256::ZERO,
261                    },
262                );
263                None
264            },
265            ExchangeEvents::MakerOrderFilledV2(e) => {
266                self.handle_maker_fill(
267                    event,
268                    RawMakerFill {
269                        perp_id: e.perpId,
270                        account_id: e.accountId,
271                        order_id: e.orderId,
272                        price_pns: e.pricePNS,
273                        lot_lns: e.lotLNS,
274                        fee_cns: e.feeCNS,
275                        builder_fee_cns: e.builderFeeCNS,
276                    },
277                );
278                None
279            },
280            ExchangeEvents::TakerOrderFilled(e) => {
281                self.handle_taker_fill(event, e.feeCNS, U256::ZERO)
282            },
283            ExchangeEvents::TakerOrderFilledV2(e) => {
284                self.handle_taker_fill(event, e.feeCNS, e.builderFeeCNS)
285            },
286            _ => None,
287        }
288    }
289
290    fn track_order_request(
291        &mut self,
292        perp_id: U256,
293        account_id: U256,
294        request_id: U256,
295        order_type: u8,
296        builder: Option<types::BuilderAttribution>,
297    ) {
298        let request_type: types::RequestType = order_type.into();
299        // Only track context for order types that can have fills
300        if let Some(side) = request_type.try_side() {
301            self.order_context = Some(OrderContext {
302                perpetual_id: perp_id.to(),
303                account_id: account_id.to(),
304                request_id: request_id.to(),
305                side,
306                builder,
307            });
308        }
309    }
310
311    fn handle_maker_fill(&mut self, event: &super::RawEvent, fill: RawMakerFill) {
312        let perp_id: types::PerpetualId = fill.perp_id.to();
313        let maker_order_id = NonZeroU16::new(fill.order_id.to()).expect("non-zero maker order ID");
314        if let Some(converters) = self.config.perpetuals.get(&perp_id) {
315            let maker_order = self.maker_orders.get(&(perp_id, maker_order_id)).copied();
316            self.pending_maker_fills.push(PendingMakerFill {
317                tx_hash: event.tx_hash(),
318                log_index: event.log_index(),
319                perpetual_id: perp_id,
320                maker_account_id: fill.account_id.to(),
321                maker_order_id,
322                maker_client_order_id: maker_order.map(|o| o.client_order_id),
323                maker_builder: maker_order.and_then(|o| o.builder),
324                price: converters.price_converter.from_unsigned(fill.price_pns),
325                size: converters.size_converter.from_unsigned(fill.lot_lns),
326                maker_fee: self.config.collateral_converter.from_unsigned(fill.fee_cns),
327                maker_builder_fee: self
328                    .config
329                    .collateral_converter
330                    .from_unsigned(fill.builder_fee_cns),
331            });
332        }
333    }
334
335    fn handle_taker_fill(
336        &mut self,
337        event: &super::RawEvent,
338        fee_cns: U256,
339        builder_fee_cns: U256,
340    ) -> Option<TradeEvent> {
341        let makers = std::mem::take(&mut self.pending_maker_fills);
342        if makers.is_empty() {
343            return None;
344        }
345
346        let ctx = self.order_context.as_ref()?;
347        let taker_tx_hash = event.tx_hash();
348
349        // Validate all maker fills have the same tx_hash as the taker fill
350        // This ensures proper correlation within the same transaction
351        if !makers.iter().all(|m| m.tx_hash == taker_tx_hash) {
352            // Data corruption: maker fills from different transaction
353            // Skip this trade to avoid incorrect correlations
354            return None;
355        }
356
357        // All makers should have the same perpetual_id (from the same order request)
358        let perpetual_id = makers.first()?.perpetual_id;
359
360        Some(
361            event.pass(types::Trade {
362                perpetual_id,
363                taker_account_id: ctx.account_id,
364                taker_request_id: ctx.request_id,
365                taker_side: ctx.side,
366                taker_fee: self.config.collateral_converter.from_unsigned(fee_cns),
367                taker_builder: ctx.builder,
368                taker_builder_fee: self
369                    .config
370                    .collateral_converter
371                    .from_unsigned(builder_fee_cns),
372                maker_fills: makers
373                    .into_iter()
374                    .map(|m| types::MakerFill {
375                        log_index: m.log_index,
376                        maker_account_id: m.maker_account_id,
377                        maker_order_id: m.maker_order_id,
378                        maker_client_order_id: m.maker_client_order_id,
379                        price: m.price,
380                        size: m.size,
381                        fee: m.maker_fee,
382                        builder: m.maker_builder,
383                        builder_fee: m.maker_builder_fee,
384                    })
385                    .collect(),
386            }),
387        )
388    }
389}
390
391impl NormalizationConfig {
392    /// Fetch normalization config from the chain.
393    ///
394    /// Tracks every perpetual listed on the exchange unless the chain
395    /// configuration names a subset, see [`Chain::perpetuals`].
396    pub async fn fetch<P: Provider + Clone>(chain: &Chain, provider: &P) -> Result<Self, DexError> {
397        let instance = ExchangeInstance::new(chain.exchange(), provider);
398
399        // Fetch exchange info for collateral decimals
400        let exchange_info = instance
401            .getExchangeInfo()
402            .call()
403            .await
404            .map_err(|err| DexError::Provider(err.into()))?;
405        let collateral_converter = num::Converter::new(exchange_info.collateralDecimals.to());
406
407        let perpetual_ids = if chain.perpetuals().is_empty() {
408            state::listed_perpetuals(chain, provider.clone(), BlockId::latest()).await?
409        } else {
410            chain.perpetuals().to_vec()
411        };
412
413        // Fetch perpetual info for each perpetual
414        let mut perpetuals = HashMap::new();
415        for perp_id in &perpetual_ids {
416            let perp_info = instance
417                .getPerpetualInfo(U256::from(*perp_id))
418                .call()
419                .await
420                .map_err(|err| DexError::Provider(err.into()))?;
421            perpetuals.insert(
422                *perp_id,
423                PerpetualConverters {
424                    price_converter: num::Converter::new(perp_info.priceDecimals.to()),
425                    size_converter: num::Converter::new(perp_info.lotDecimals.to()),
426                },
427            );
428        }
429
430        Ok(Self { collateral_converter, perpetuals })
431    }
432}
433
434#[cfg(test)]
435mod tests {
436    use std::time::Duration;
437
438    use alloy::{
439        primitives::I256, providers::ProviderBuilder, rpc::client::RpcClient,
440        transports::layers::RetryBackoffLayer,
441    };
442    use fastnum::udec64;
443    use futures::StreamExt;
444
445    use super::*;
446    use crate::{
447        Chain,
448        abi::dex::Exchange::{MakerOrderFilledV2, OrderPlaced, OrderRequestV2, TakerOrderFilledV2},
449        stream::RawEvent,
450    };
451
452    fn order_request(
453        perpetual_id: types::PerpetualId,
454        account_id: types::AccountId,
455        request_id: types::RequestId,
456        order_type: u8,
457        builder: Option<types::BuilderAttribution>,
458    ) -> ExchangeEvents {
459        ExchangeEvents::OrderRequestV2(OrderRequestV2 {
460            perpId: U256::from(perpetual_id),
461            accountId: U256::from(account_id),
462            orderDescId: U256::from(request_id),
463            orderId: U256::ZERO,
464            orderType: order_type,
465            pricePNS: U256::from(100),
466            lotLNS: U256::from(1),
467            expiryBlock: U256::ZERO,
468            postOnly: false,
469            fillOrKill: false,
470            immediateOrCancel: false,
471            maxMatches: U256::ZERO,
472            leverageHdths: U256::ZERO,
473            lastExecutionBlock: U256::ZERO,
474            amountCNS: U256::ZERO,
475            maxNegPnlCollatBPS: U256::ZERO,
476            gasLeft: U256::ZERO,
477            extension: builder
478                .map(|b| b.encode().expect("fee within range"))
479                .unwrap_or_default(),
480        })
481    }
482
483    fn normalization_config(perpetual_id: types::PerpetualId) -> NormalizationConfig {
484        let converter = num::Converter::new(0);
485        NormalizationConfig {
486            collateral_converter: converter,
487            perpetuals: HashMap::from([(
488                perpetual_id,
489                PerpetualConverters { price_converter: converter, size_converter: converter },
490            )]),
491        }
492    }
493
494    #[test]
495    fn client_order_ids_ignore_unconfigured_perpetuals() {
496        let mut processor = TradeProcessor::new(normalization_config(1));
497        processor.order_context = Some(OrderContext {
498            perpetual_id: 2,
499            account_id: 7,
500            request_id: 42,
501            side: types::OrderSide::Ask,
502            builder: None,
503        });
504        let event = RawEvent::empty(ExchangeEvents::OrderPlaced(OrderPlaced {
505            orderId: U256::from(9),
506            lotLNS: U256::from(1),
507            lockedBalanceCNS: U256::ZERO,
508            amountCNS: I256::ZERO,
509            balanceCNS: U256::ZERO,
510        }));
511
512        _ = processor.process_event(&event);
513
514        assert!(processor.maker_orders.is_empty());
515    }
516
517    /// Drives a single maker-vs-taker match through the processor, attributing
518    /// each side to the given builder, and returns the resulting trade.
519    fn one_match_trade(
520        maker_builder: Option<types::BuilderAttribution>,
521        maker_builder_fee: u64,
522        taker_builder: Option<types::BuilderAttribution>,
523        taker_builder_fee: u64,
524    ) -> types::Trade {
525        const PERPETUAL_ID: types::PerpetualId = 1;
526        const MAKER_ACCOUNT_ID: types::AccountId = 7;
527        const MAKER_CLIENT_ORDER_ID: types::RequestId = 42;
528        const TAKER_REQUEST_ID: types::RequestId = 84;
529        const MAKER_ORDER_ID: u16 = 9;
530
531        let mut processor = TradeProcessor::new(normalization_config(PERPETUAL_ID));
532        _ = processor.process_event(&RawEvent::empty(order_request(
533            PERPETUAL_ID,
534            MAKER_ACCOUNT_ID,
535            MAKER_CLIENT_ORDER_ID,
536            1,
537            maker_builder,
538        )));
539        _ = processor.process_event(&RawEvent::empty(ExchangeEvents::OrderPlaced(OrderPlaced {
540            orderId: U256::from(MAKER_ORDER_ID),
541            lotLNS: U256::from(1),
542            lockedBalanceCNS: U256::ZERO,
543            amountCNS: I256::ZERO,
544            balanceCNS: U256::ZERO,
545        })));
546        _ = processor.process_event(&RawEvent::empty(order_request(
547            PERPETUAL_ID,
548            8,
549            TAKER_REQUEST_ID,
550            0,
551            taker_builder,
552        )));
553        _ = processor.process_event(&RawEvent::empty(ExchangeEvents::MakerOrderFilledV2(
554            MakerOrderFilledV2 {
555                perpId: U256::from(PERPETUAL_ID),
556                accountId: U256::from(MAKER_ACCOUNT_ID),
557                orderId: U256::from(MAKER_ORDER_ID),
558                pricePNS: U256::from(100),
559                lotLNS: U256::from(1),
560                feeCNS: U256::from(2),
561                lockedBalanceCNS: U256::ZERO,
562                amountCNS: I256::ZERO,
563                balanceCNS: U256::ZERO,
564                builderId: U256::from(maker_builder.map(|b| b.builder_id()).unwrap_or_default()),
565                builderFeeCNS: U256::from(maker_builder_fee),
566            },
567        )));
568        processor
569            .process_event(&RawEvent::empty(ExchangeEvents::TakerOrderFilledV2(
570                TakerOrderFilledV2 {
571                    entryPricePNS: U256::from(100),
572                    collatPricePNS: U256::from(100),
573                    pnlPricePNS: U256::from(100),
574                    lotLNS: U256::from(1),
575                    feeCNS: U256::from(3),
576                    amountCNS: I256::ZERO,
577                    balanceCNS: U256::ZERO,
578                    builderId: U256::from(
579                        taker_builder.map(|b| b.builder_id()).unwrap_or_default(),
580                    ),
581                    builderFeeCNS: U256::from(taker_builder_fee),
582                },
583            )))
584            .expect("trade exists")
585            .event()
586            .clone()
587    }
588
589    #[test]
590    fn maker_fill_includes_observed_client_order_id() {
591        let trade = one_match_trade(None, 0, None, 0);
592        let maker_fill = trade.maker_fills.first().expect("maker fill exists");
593
594        assert_eq!(maker_fill.maker_client_order_id, Some(42));
595        assert_eq!(maker_fill.builder, None);
596        assert_eq!(trade.total_builder_fees(), udec64!(0));
597    }
598
599    #[test]
600    fn builder_attribution_recovered_from_order_requests() {
601        let maker_builder = types::BuilderAttribution::new(3, udec64!(0.0005));
602        let taker_builder = types::BuilderAttribution::new(4, udec64!(0.001));
603        let trade = one_match_trade(Some(maker_builder), 1, Some(taker_builder), 2);
604
605        // Attribution rides on the request that placed each order...
606        let maker_fill = trade.maker_fills.first().expect("maker fill exists");
607        assert_eq!(maker_fill.builder, Some(maker_builder));
608        assert_eq!(trade.taker_builder, Some(taker_builder));
609
610        // ...while the fee earned comes from the fill events, and is part of the
611        // reported fee rather than additional to it.
612        assert_eq!(maker_fill.builder_fee, udec64!(1));
613        assert_eq!(trade.taker_builder_fee, udec64!(2));
614        assert_eq!(trade.total_builder_fees(), udec64!(3));
615        assert_eq!(trade.builder_total(3), udec64!(1));
616        assert_eq!(trade.builder_total(4), udec64!(2));
617        assert_eq!(trade.builder_total(5), udec64!(0));
618        assert!(maker_fill.builder_fee < maker_fill.fee);
619        assert!(trade.taker_builder_fee < trade.taker_fee);
620    }
621
622    #[tokio::test]
623    async fn test_stream_recent_blocks() {
624        let client = RpcClient::builder()
625            .layer(RetryBackoffLayer::new(10, 100, 200))
626            .connect("https://testnet-rpc.monad.xyz")
627            .await
628            .unwrap();
629        client.set_poll_interval(Duration::from_millis(100));
630        let provider = ProviderBuilder::new().connect_client(client);
631
632        let testnet = Chain::testnet();
633        let block_num = provider.get_block_number().await.unwrap() + 1;
634        let raw_stream = crate::stream::raw(
635            &testnet,
636            provider.clone(),
637            types::StateInstant::new(block_num, 0),
638            tokio::time::sleep,
639        );
640
641        let trade_stream = trade(&testnet, provider, raw_stream).await.unwrap();
642        let block_trades = trade_stream.take(10).collect::<Vec<_>>().await;
643
644        for bt in &block_trades {
645            println!("block trades: {:?}", bt);
646        }
647    }
648}