massa_grpc/stream/
new_slot_abi_call_stacks.rs

1use futures_util::StreamExt;
2use massa_proto_rs::massa::api::v1::{self as grpc_api, FinalityLevel};
3use std::pin::Pin;
4use tokio::select;
5use tonic::{Request, Streaming};
6use tracing::{error, warn};
7
8use crate::error::match_for_io_error;
9use crate::{error::GrpcError, server::MassaPublicGrpc};
10
11#[cfg(feature = "execution-trace")]
12use crate::public::into_element;
13#[cfg(not(feature = "execution-trace"))]
14use massa_models::slot::Slot;
15#[cfg(feature = "execution-trace")]
16use massa_proto_rs::massa::api::v1::{
17    AscabiCallStack, DeferredCallAbiCallStack, OperationAbiCallStack,
18};
19#[cfg(not(feature = "execution-trace"))]
20#[derive(Clone)]
21struct SlotAbiCallStack {
22    /// Slot
23    pub slot: Slot,
24}
25
26/// Type declaration for NewSlotExecutionOutputs
27pub type NewSlotABICallStacksStreamType = Pin<
28    Box<
29        dyn futures_util::Stream<
30                Item = Result<grpc_api::NewSlotAbiCallStacksResponse, tonic::Status>,
31            > + Send
32            + 'static,
33    >,
34>;
35
36/// Creates a new stream of new slots abi call stacks
37pub(crate) async fn new_slot_abi_call_stacks(
38    grpc: &MassaPublicGrpc,
39    request: Request<Streaming<grpc_api::NewSlotAbiCallStacksRequest>>,
40) -> Result<NewSlotABICallStacksStreamType, GrpcError> {
41    let (tx, rx) = tokio::sync::mpsc::channel(grpc.grpc_config.max_channel_size);
42    // Extract the incoming stream of abi call stacks messages
43    let mut in_stream = request.into_inner();
44
45    // Clone the slot execution traces channel sender to subscribe from the spawned task
46    #[cfg(feature = "execution-trace")]
47    let traces_sender = grpc.execution_channels.slot_execution_traces_sender.clone();
48
49    tokio::spawn(async move {
50        // Wait for the first request to establish the finality level before
51        // subscribing, so that no traces are forwarded before the client's
52        // selection is known
53        let mut finality: FinalityLevel = match in_stream.next().await {
54            Some(Ok(message)) => message.finality_level(),
55            _ => {
56                error!("empty request");
57                return;
58            }
59        };
60
61        // Subscribe to the new slot execution events channel
62        #[cfg(feature = "execution-trace")]
63        let mut subscriber = traces_sender.subscribe();
64        #[cfg(not(feature = "execution-trace"))]
65        let (mut subscriber, _receiver) = {
66            let (subscriber_, receiver) =
67                tokio::sync::broadcast::channel::<(SlotAbiCallStack, (), bool)>(0);
68            (subscriber_.subscribe(), receiver)
69        };
70
71        loop {
72            select! {
73                // Receive a new slot execution traces from the subscriber
74                event = subscriber.recv() => {
75                    match event {
76                        Ok((massa_slot_execution_trace, _slot_transfers, received_finality)) => {
77                            if (finality == FinalityLevel::Final && !received_finality) ||
78                                (finality == FinalityLevel::Candidate && received_finality) {
79                                continue;
80                            }
81
82                            #[allow(unused_mut)]
83                            let mut ret = grpc_api::NewSlotAbiCallStacksResponse {
84                                slot: Some(massa_slot_execution_trace.slot.into()),
85                                asc_call_stacks: vec![],
86                                operation_call_stacks: vec![],
87                                deferred_call_stacks: vec![],
88                            };
89
90                            #[cfg(feature = "execution-trace")]
91                            {
92                                for (i, asc_call_stack) in massa_slot_execution_trace.asc_call_stacks.iter().enumerate() {
93                                    ret.asc_call_stacks.push(
94                                        AscabiCallStack {
95                                            index: i as u64,
96                                            call_stack: asc_call_stack.iter().map(into_element).collect()
97                                        }
98                                    )
99                                }
100                                for (op_id, op_call_stack) in massa_slot_execution_trace.operation_call_stacks.iter() {
101                                    ret.operation_call_stacks.push(
102                                        OperationAbiCallStack {
103                                            operation_id: op_id.to_string(),
104                                            call_stack: op_call_stack.iter().map(into_element).collect()
105                                        }
106                                    )
107                                }
108                                for (deferred_call_id, deferred_call_call_stack) in massa_slot_execution_trace.deferred_call_stacks.iter() {
109                                    ret.deferred_call_stacks.push(
110                                        DeferredCallAbiCallStack {
111                                            deferred_call_id: deferred_call_id.to_string(),
112                                            call_stack: deferred_call_call_stack.iter().map(into_element).collect()
113                                        }
114                                    )
115                                }
116
117                            }
118
119                            if let Err(e) = tx.send(Ok(ret)).await {
120                                error!("failed to send new slot execution trace: {}", e);
121                                break;
122                            }
123                        }
124                        Err(e) => {
125                            error!("error on receive new slot execution trace : {}", e)
126                        }
127                    }
128                }
129                // Receive a new message from the in_stream
130                res = in_stream.next() => {
131                    match res {
132                        Some(res) => {
133                            match res {
134                                Ok(message) => {
135                                    finality = message.finality_level();
136                                },
137                                Err(e) => {
138                                    // Any io error -> break
139                                    if let Some(io_err) = match_for_io_error(&e) {
140                                        warn!("client disconnected, broken pipe: {}", io_err);
141                                        break;
142                                    }
143                                    error!("{}", e);
144                                    if let Err(e2) = tx.send(Err(e)).await {
145                                        error!("failed to send back error response: {}", e2);
146                                        break;
147                                    }
148                                }
149                            }
150                        }
151                        None => {
152                            // the client has disconnected
153                            break;
154                        }
155                    }
156                }
157            }
158        }
159    });
160    let out_stream = tokio_stream::wrappers::ReceiverStream::new(rx);
161    Ok(Box::pin(out_stream) as NewSlotABICallStacksStreamType)
162}