Skip to main content

iota_core/
subscription_handler.rs

1// Copyright (c) Mysten Labs, Inc.
2// Modifications Copyright (c) 2024 IOTA Stiftung
3// SPDX-License-Identifier: Apache-2.0
4
5use std::sync::Arc;
6
7use iota_json_rpc_types::{
8    EffectsWithInput, EventFilter, IotaEvent, IotaTransactionBlockEffects,
9    IotaTransactionBlockEffectsAPI, IotaTransactionBlockEvents, TransactionFilter,
10};
11use iota_sdk_types::Transaction;
12use iota_types::error::IotaResult;
13use prometheus_filtered::{
14    IntCounterVec, IntGaugeVec, Registry, register_int_counter_vec_with_registry,
15    register_int_gauge_vec_with_registry,
16};
17use tokio_stream::Stream;
18use tracing::{error, instrument, trace};
19
20use crate::streamer::Streamer;
21
22#[cfg(test)]
23#[path = "unit_tests/subscription_handler_tests.rs"]
24mod subscription_handler_tests;
25
26pub const EVENT_DISPATCH_BUFFER_SIZE: usize = 1000;
27
28pub struct SubscriptionMetrics {
29    pub streaming_success: IntCounterVec,
30    pub streaming_failure: IntCounterVec,
31    pub streaming_active_subscriber_number: IntGaugeVec,
32    pub dropped_submissions: IntCounterVec,
33}
34
35impl SubscriptionMetrics {
36    pub fn new(registry: &Registry) -> Self {
37        Self {
38            streaming_success: register_int_counter_vec_with_registry!(
39                "streaming_success",
40                "Total number of items that are streamed successfully",
41                &["type"],
42                registry,
43            )
44            .unwrap(),
45            streaming_failure: register_int_counter_vec_with_registry!(
46                "streaming_failure",
47                "Total number of items that fail to be streamed",
48                &["type"],
49                registry,
50            )
51            .unwrap(),
52            streaming_active_subscriber_number: register_int_gauge_vec_with_registry!(
53                "streaming_active_subscriber_number",
54                "Current number of active subscribers",
55                &["type"],
56                registry,
57            )
58            .unwrap(),
59            dropped_submissions: register_int_counter_vec_with_registry!(
60                "streaming_dropped_submissions",
61                "Total number of submissions that are dropped",
62                &["type"],
63                registry,
64            )
65            .unwrap(),
66        }
67    }
68}
69
70pub struct SubscriptionHandler {
71    event_streamer: Streamer<IotaEvent, IotaEvent, EventFilter>,
72    transaction_streamer:
73        Streamer<EffectsWithInput, IotaTransactionBlockEffects, TransactionFilter>,
74}
75
76impl SubscriptionHandler {
77    pub fn new(registry: &Registry) -> Self {
78        let metrics = Arc::new(SubscriptionMetrics::new(registry));
79        Self {
80            event_streamer: Streamer::spawn(EVENT_DISPATCH_BUFFER_SIZE, metrics.clone(), "event"),
81            transaction_streamer: Streamer::spawn(EVENT_DISPATCH_BUFFER_SIZE, metrics, "tx"),
82        }
83    }
84}
85
86impl SubscriptionHandler {
87    #[instrument(level = "trace", skip_all, fields(tx_digest =? effects.transaction_digest()), err)]
88    pub fn process_tx(
89        &self,
90        input: &Transaction,
91        effects: &IotaTransactionBlockEffects,
92        events: &IotaTransactionBlockEvents,
93    ) -> IotaResult {
94        trace!(
95            num_events = events.data.len(),
96            tx_digest =? effects.transaction_digest(),
97            "Processing tx/event subscription"
98        );
99
100        if let Err(e) = self.transaction_streamer.try_send(EffectsWithInput {
101            input: input.clone(),
102            effects: effects.clone(),
103        }) {
104            error!(error =? e, "Failed to send transaction to dispatch");
105        }
106
107        // serially dispatch event processing to honor events' orders.
108        for event in events.data.clone() {
109            // Send to unified event streamer (serves both JSON-RPC and gRPC subscribers)
110            if let Err(e) = self.event_streamer.try_send(event) {
111                error!(error =? e, "Failed to send event to dispatch");
112            }
113        }
114        Ok(())
115    }
116
117    pub fn subscribe_events(&self, filter: EventFilter) -> impl Stream<Item = IotaEvent> {
118        self.event_streamer.subscribe(filter)
119    }
120
121    pub fn subscribe_transactions(
122        &self,
123        filter: TransactionFilter,
124    ) -> impl Stream<Item = IotaTransactionBlockEffects> {
125        self.transaction_streamer.subscribe(filter)
126    }
127}