iota_core/
subscription_handler.rs1use 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 for event in events.data.clone() {
109 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}