Skip to main content

acuity_index_substrate/
websockets.rs

1use crate::shared::*;
2use futures::{SinkExt, StreamExt};
3use sled::Tree;
4use std::net::SocketAddr;
5use subxt::backend::legacy::LegacyRpcMethods;
6use subxt::metadata::types::Metadata;
7use tokio::net::{TcpListener, TcpStream};
8use tokio::sync::{
9    mpsc::{UnboundedSender, unbounded_channel},
10    watch::Receiver,
11};
12use tokio_tungstenite::tungstenite;
13use tracing::{error, info};
14use zerocopy::AsBytes;
15use zerocopy::{BigEndian, FromBytes, byteorder::U32};
16
17pub fn process_msg_status<R: RuntimeIndexer>(span_db: &Tree) -> ResponseMessage<R::ChainKey> {
18    let mut spans = vec![];
19    for (key, value) in span_db.into_iter().flatten() {
20        let span_value = SpanDbValue::read_from(&value).unwrap();
21        let start: u32 = span_value.start.into();
22        let end: u32 = u32::from_be_bytes(key.as_ref().try_into().unwrap());
23        let span = Span { start, end };
24        spans.push(span);
25    }
26    ResponseMessage::Status(spans)
27}
28
29pub fn process_msg_subscribe_status<R: RuntimeIndexer>(
30    sub_tx: &UnboundedSender<SubscriptionMessage<R::ChainKey>>,
31    sub_response_tx: &UnboundedSender<ResponseMessage<R::ChainKey>>,
32) -> ResponseMessage<R::ChainKey> {
33    let msg = SubscriptionMessage::SubscribeStatus {
34        sub_response_tx: sub_response_tx.clone(),
35    };
36    sub_tx.send(msg).unwrap();
37    ResponseMessage::Subscribed
38}
39
40pub fn process_msg_unsubscribe_status<R: RuntimeIndexer>(
41    sub_tx: &UnboundedSender<SubscriptionMessage<R::ChainKey>>,
42    sub_response_tx: &UnboundedSender<ResponseMessage<R::ChainKey>>,
43) -> ResponseMessage<R::ChainKey> {
44    let msg = SubscriptionMessage::UnsubscribeStatus {
45        sub_response_tx: sub_response_tx.clone(),
46    };
47    sub_tx.send(msg).unwrap();
48    ResponseMessage::Unsubscribed
49}
50
51pub async fn process_msg_variants<R: RuntimeIndexer>(
52    rpc: &LegacyRpcMethods<R::RuntimeConfig>,
53) -> Result<ResponseMessage<R::ChainKey>, IndexError> {
54    let metadata: Metadata = rpc
55        .state_get_metadata(None)
56        .await?
57        .to_frame_metadata()?
58        .try_into()?;
59    let mut pallets = Vec::new();
60
61    for pallet in metadata.pallets() {
62        let mut pallet_meta = PalletMeta {
63            index: pallet.index(),
64            name: pallet.name().to_owned(),
65            events: Vec::new(),
66        };
67
68        if let Some(variants) = pallet.event_variants() {
69            for variant in variants {
70                pallet_meta.events.push(EventMeta {
71                    index: variant.index,
72                    name: variant.name.clone(),
73                })
74            }
75            pallets.push(pallet_meta);
76        }
77    }
78    Ok(ResponseMessage::Variants(pallets))
79}
80
81pub fn get_events_variant(tree: &Tree, pallet_id: u8, variant_id: u8) -> Vec<Event> {
82    let mut events = Vec::new();
83    let mut iter = tree.scan_prefix([pallet_id, variant_id]).keys();
84
85    while let Some(Ok(key)) = iter.next_back() {
86        let key = VariantKey::read_from(&key).unwrap();
87
88        events.push(Event {
89            block_number: key.block_number.into(),
90            event_index: key.event_index.into(),
91        });
92
93        if events.len() == 100 {
94            break;
95        }
96    }
97    events
98}
99
100pub fn get_events_bytes32(tree: &Tree, key: &Bytes32) -> Vec<Event> {
101    let mut events = Vec::new();
102    let mut iter = tree.scan_prefix(key).keys();
103
104    while let Some(Ok(key)) = iter.next_back() {
105        let key = Bytes32Key::read_from(&key).unwrap();
106
107        events.push(Event {
108            block_number: key.block_number.into(),
109            event_index: key.event_index.into(),
110        });
111
112        if events.len() == 100 {
113            break;
114        }
115    }
116    events
117}
118
119pub fn get_events_u32(tree: &Tree, key: u32) -> Vec<Event> {
120    let mut events = Vec::new();
121    let mut iter = tree.scan_prefix(key.to_be_bytes()).keys();
122
123    while let Some(Ok(key)) = iter.next_back() {
124        let key = U32Key::read_from(&key).unwrap();
125
126        events.push(Event {
127            block_number: key.block_number.into(),
128            event_index: key.event_index.into(),
129        });
130
131        if events.len() == 100 {
132            break;
133        }
134    }
135    events
136}
137
138pub fn process_msg_get_events_substrate<R: RuntimeIndexer>(
139    trees: &Trees<<R::ChainKey as IndexKey>::ChainTrees>,
140    key: &SubstrateKey,
141) -> Vec<Event> {
142    match key {
143        SubstrateKey::AccountId(account_id) => {
144            get_events_bytes32(&trees.substrate.account_id, account_id)
145        }
146        SubstrateKey::AccountIndex(account_index) => {
147            get_events_u32(&trees.substrate.account_index, *account_index)
148        }
149        SubstrateKey::BountyIndex(bounty_index) => {
150            get_events_u32(&trees.substrate.bounty_index, *bounty_index)
151        }
152        SubstrateKey::EraIndex(era_index) => get_events_u32(&trees.substrate.era_index, *era_index),
153        SubstrateKey::MessageId(message_id) => {
154            get_events_bytes32(&trees.substrate.message_id, message_id)
155        }
156        SubstrateKey::PoolId(pool_id) => get_events_u32(&trees.substrate.pool_id, *pool_id),
157        SubstrateKey::PreimageHash(preimage_hash) => {
158            get_events_bytes32(&trees.substrate.preimage_hash, preimage_hash)
159        }
160        SubstrateKey::ProposalHash(proposal_hash) => {
161            get_events_bytes32(&trees.substrate.proposal_hash, proposal_hash)
162        }
163        SubstrateKey::ProposalIndex(proposal_index) => {
164            get_events_u32(&trees.substrate.proposal_index, *proposal_index)
165        }
166        SubstrateKey::RefIndex(ref_index) => get_events_u32(&trees.substrate.ref_index, *ref_index),
167        SubstrateKey::RegistrarIndex(registrar_index) => {
168            get_events_u32(&trees.substrate.registrar_index, *registrar_index)
169        }
170        SubstrateKey::SessionIndex(session_index) => {
171            get_events_u32(&trees.substrate.session_index, *session_index)
172        }
173        SubstrateKey::TipHash(tip_hash) => get_events_bytes32(&trees.substrate.tip_hash, tip_hash),
174        SubstrateKey::SpendIndex(spend_index) => {
175            get_events_u32(&trees.substrate.spend_index, *spend_index)
176        }
177    }
178}
179
180pub fn process_msg_get_events<R: RuntimeIndexer>(
181    trees: &Trees<<R::ChainKey as IndexKey>::ChainTrees>,
182    key: Key<R::ChainKey>,
183) -> ResponseMessage<R::ChainKey> {
184    let events = match key {
185        Key::Variant(pallet_id, variant_id) => {
186            get_events_variant(&trees.variant, pallet_id, variant_id)
187        }
188        Key::Substrate(ref key) => process_msg_get_events_substrate::<R>(trees, key),
189        Key::Chain(ref key) => key.get_key_events(&trees.chain),
190    };
191
192    let mut block_numbers = events
193        .iter()
194        .map(|event| event.block_number)
195        .collect::<Vec<u32>>();
196    block_numbers.sort();
197    block_numbers.dedup();
198
199    let mut block_events = Vec::new();
200
201    for block_number in block_numbers.iter() {
202        let key: U32<BigEndian> = (*block_number).into();
203        if let Ok(Some(event_bytes)) = trees.block_events.get(key.as_bytes()) {
204            block_events.push(Block {
205                block_number: *block_number,
206                bytes: event_bytes.to_vec(),
207            });
208        };
209    }
210
211    ResponseMessage::Events {
212        key,
213        events,
214        block_events,
215    }
216}
217
218pub fn process_msg_subscribe_events<R: RuntimeIndexer>(
219    key: Key<R::ChainKey>,
220    sub_tx: &UnboundedSender<SubscriptionMessage<R::ChainKey>>,
221    sub_response_tx: &UnboundedSender<ResponseMessage<R::ChainKey>>,
222) -> ResponseMessage<R::ChainKey> {
223    let msg = SubscriptionMessage::SubscribeEvents {
224        key,
225        sub_response_tx: sub_response_tx.clone(),
226    };
227    sub_tx.send(msg).unwrap();
228    ResponseMessage::Subscribed
229}
230
231pub fn process_msg_unsubscribe_events<R: RuntimeIndexer>(
232    key: Key<R::ChainKey>,
233    sub_tx: &UnboundedSender<SubscriptionMessage<R::ChainKey>>,
234    sub_response_tx: &UnboundedSender<ResponseMessage<R::ChainKey>>,
235) -> ResponseMessage<R::ChainKey> {
236    let msg = SubscriptionMessage::UnsubscribeEvents {
237        key,
238        sub_response_tx: sub_response_tx.clone(),
239    };
240    sub_tx.send(msg).unwrap();
241    ResponseMessage::Unsubscribed
242}
243
244pub async fn process_msg<R: RuntimeIndexer>(
245    rpc: &LegacyRpcMethods<R::RuntimeConfig>,
246    trees: &Trees<<R::ChainKey as IndexKey>::ChainTrees>,
247    msg: RequestMessage<R::ChainKey>,
248    sub_tx: &UnboundedSender<SubscriptionMessage<R::ChainKey>>,
249    sub_response_tx: &UnboundedSender<ResponseMessage<R::ChainKey>>,
250) -> Result<ResponseMessage<R::ChainKey>, IndexError> {
251    Ok(match msg {
252        RequestMessage::Status => process_msg_status::<R>(&trees.span),
253        RequestMessage::SubscribeStatus => {
254            process_msg_subscribe_status::<R>(sub_tx, sub_response_tx)
255        }
256        RequestMessage::UnsubscribeStatus => {
257            process_msg_unsubscribe_status::<R>(sub_tx, sub_response_tx)
258        }
259        RequestMessage::Variants => process_msg_variants::<R>(rpc).await?,
260        RequestMessage::GetEvents { key } => process_msg_get_events::<R>(trees, key),
261        RequestMessage::SubscribeEvents { key } => {
262            process_msg_subscribe_events::<R>(key, sub_tx, sub_response_tx)
263        }
264        RequestMessage::UnsubscribeEvents { key } => {
265            process_msg_unsubscribe_events::<R>(key, sub_tx, sub_response_tx)
266        }
267        RequestMessage::SizeOnDisk => ResponseMessage::SizeOnDisk(trees.root.size_on_disk()?),
268    })
269}
270
271async fn handle_connection<R: RuntimeIndexer>(
272    rpc: LegacyRpcMethods<R::RuntimeConfig>,
273    raw_stream: TcpStream,
274    addr: SocketAddr,
275    trees: Trees<<R::ChainKey as IndexKey>::ChainTrees>,
276    sub_tx: UnboundedSender<SubscriptionMessage<R::ChainKey>>,
277) -> Result<(), IndexError> {
278    info!("Incoming TCP connection from: {}", addr);
279    let ws_stream = tokio_tungstenite::accept_async(raw_stream).await?;
280    info!("WebSocket connection established: {}", addr);
281
282    let (mut ws_sender, mut ws_receiver) = ws_stream.split();
283    // Create the channel for the substrate thread to send event messages to this thread.
284    let (sub_events_tx, mut sub_events_rx) = unbounded_channel();
285
286    loop {
287        tokio::select! {
288            Some(Ok(msg)) = ws_receiver.next() => {
289                if msg.is_text() || msg.is_binary() {
290                    match serde_json::from_str(msg.to_text()?) {
291                        Ok(request_json) => {
292                            let response_msg = process_msg::<R>(&rpc, &trees, request_json, &sub_tx, &sub_events_tx).await?;
293                            let response_json = serde_json::to_string(&response_msg).unwrap();
294                            ws_sender.send(tungstenite::Message::Text(response_json)).await?;
295                        },
296                        Err(error) => error!("{}", error),
297                    }
298                }
299            },
300            Some(msg) = sub_events_rx.recv() => {
301                let response_json = serde_json::to_string(&msg).unwrap();
302                ws_sender.send(tungstenite::Message::Text(response_json)).await?;
303            },
304        }
305    }
306}
307
308pub async fn websockets_listen<R: RuntimeIndexer + 'static>(
309    trees: Trees<<R::ChainKey as IndexKey>::ChainTrees>,
310    rpc: LegacyRpcMethods<R::RuntimeConfig>,
311    port: u16,
312    mut exit_rx: Receiver<bool>,
313    sub_tx: UnboundedSender<SubscriptionMessage<R::ChainKey>>,
314) {
315    let mut addr = "0.0.0.0:".to_string();
316    addr.push_str(&port.to_string());
317
318    // Create the event loop and TCP listener we'll accept connections on.
319    let try_socket = TcpListener::bind(&addr).await;
320    let listener = try_socket.expect("Failed to bind");
321    info!("Listening on: {}", addr);
322
323    // Let's spawn the handling of each connection in a separate task.
324    loop {
325        tokio::select! {
326            biased;
327
328            _ = exit_rx.changed() => {
329                break;
330            }
331            Ok((stream, addr)) = listener.accept() => {
332                tokio::spawn(handle_connection::<R>(
333                    rpc.clone(),
334                    stream,
335                    addr,
336                    trees.clone(),
337                    sub_tx.clone(),
338                ));
339            }
340        }
341    }
342}