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 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 let try_socket = TcpListener::bind(&addr).await;
320 let listener = try_socket.expect("Failed to bind");
321 info!("Listening on: {}", addr);
322
323 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}