massa_grpc/stream/
new_slot_transfers.rs

1use std::pin::Pin;
2
3#[cfg(feature = "execution-trace")]
4use crate::{error::GrpcError, server::MassaPublicGrpc};
5#[cfg(feature = "execution-trace")]
6use tonic::{Request, Streaming};
7
8/// Type declaration for NewSlotTransfers
9pub type NewSlotTransfersStreamType = Pin<
10    Box<
11        dyn futures_util::Stream<
12                Item = Result<
13                    massa_proto_rs::massa::api::v1::NewSlotTransfersResponse,
14                    tonic::Status,
15                >,
16            > + Send
17            + 'static,
18    >,
19>;
20
21#[cfg(feature = "execution-trace")]
22/// Creates a new stream of new slots transfers
23pub(crate) async fn new_slot_transfers(
24    grpc: &MassaPublicGrpc,
25    request: Request<Streaming<massa_proto_rs::massa::api::v1::NewSlotTransfersRequest>>,
26) -> Result<NewSlotTransfersStreamType, GrpcError> {
27    use crate::error::match_for_io_error;
28    use futures_util::StreamExt;
29    use massa_proto_rs::massa::api::v1::{self as grpc_api, FinalityLevel, TransferInfo};
30    use tokio::select;
31    use tracing::{error, warn};
32
33    let (tx, rx) = tokio::sync::mpsc::channel(grpc.grpc_config.max_channel_size);
34    // Extract the incoming stream of abi call stacks messages
35    let mut in_stream = request.into_inner();
36
37    // Subscribe to the new slot execution events channel
38    let mut subscriber = grpc
39        .execution_channels
40        .slot_execution_traces_sender
41        .subscribe();
42
43    tokio::spawn({
44        async move {
45            // Issue #5066: standardize startup with `new_slot_execution_outputs` by waiting
46            // for the client's first request before emitting anything, so the default filter
47            // is applied only once the client has had a chance to configure it and this stream
48            // never emits data under an implicit default.
49            let mut finality = match in_stream.next().await {
50                Some(Ok(first)) => first.finality_level(),
51                // No usable initial request (error or client disconnected): stream nothing.
52                _ => {
53                    error!("empty request");
54                    return;
55                }
56            };
57            loop {
58                select! {
59                    // Receive a new slot execution traces from the subscriber
60                    event = subscriber.recv() => {
61                        match event {
62                            Ok((massa_slot_execution_trace, slot_transfers, is_final)) => {
63                                // A `FinalityLevel::Unspecified` finality level (the default,
64                                // e.g. when the client sends no explicit level) accepts BOTH
65                                // candidate and final transfers, matching the default filter of
66                                // `FilterNewSlotExec` used by `new_slot_execution_outputs`.
67                                if (finality == FinalityLevel::Final && !is_final) ||
68                                    (finality == FinalityLevel::Candidate && is_final) {
69                                    continue;
70                                }
71                                let mut ret_transfers = Vec::new();
72                                // flatten & filter transfer trace in asc_call_stacks
73
74                                let abi_transfer_1 = "assembly_script_transfer_coins".to_string();
75                                let abi_transfer_2 = "assembly_script_transfer_coins_for".to_string();
76                                let abi_transfer_3 = "abi_transfer_coins".to_string();
77                                let transfer_abi_names = vec![abi_transfer_1, abi_transfer_2, abi_transfer_3];
78                                for (i, asc_call_stack) in massa_slot_execution_trace.asc_call_stacks.iter().enumerate() {
79                                    for abi_trace in asc_call_stack {
80                                        let only_transfer = abi_trace.flatten_filter(&transfer_abi_names);
81                                        for transfer in only_transfer {
82                                            let (t_from, t_to, t_amount) = transfer.parse_transfer();
83                                            ret_transfers.push(TransferInfo {
84                                                from: t_from.clone(),
85                                                to: t_to.clone(),
86                                                amount: t_amount,
87                                                operation_id_or_asc_index: Some(
88                                                    grpc_api::transfer_info::OperationIdOrAscIndex::AscIndex(i as u64),
89                                                ),
90                                            });
91                                        }
92                                    }
93                                }
94
95                                for deferred_call_call_stack in massa_slot_execution_trace.deferred_call_stacks {
96                                    let deferred_call_id = deferred_call_call_stack.0;
97                                    let deferred_call_call_stack = deferred_call_call_stack.1;
98                                    for abi_trace in deferred_call_call_stack {
99                                        let only_transfer = abi_trace.flatten_filter(&transfer_abi_names);
100                                        for transfer in only_transfer {
101                                            let (t_from, t_to, t_amount) = transfer.parse_transfer();
102                                            ret_transfers.push(TransferInfo {
103                                                from: t_from.clone(),
104                                                to: t_to.clone(),
105                                                amount: t_amount,
106                                                operation_id_or_asc_index: Some(
107                                                    grpc_api::transfer_info::OperationIdOrAscIndex::DeferredCallId(deferred_call_id.to_string()),
108                                                ),
109                                            });
110                                        }
111                                    }
112                                }
113
114                                for op_call_stack in massa_slot_execution_trace.operation_call_stacks {
115                                    let op_id = op_call_stack.0;
116                                    let op_call_stack = op_call_stack.1;
117                                    for abi_trace in op_call_stack {
118                                        let only_transfer = abi_trace.flatten_filter(&transfer_abi_names);
119                                        for transfer in only_transfer {
120                                            let (t_from, t_to, t_amount) = transfer.parse_transfer();
121                                            ret_transfers.push(TransferInfo {
122                                                from: t_from.clone(),
123                                                to: t_to.clone(),
124                                                amount: t_amount,
125                                                operation_id_or_asc_index: Some(
126                                                    grpc_api::transfer_info::OperationIdOrAscIndex::OperationId(op_id.to_string()),
127                                                ),
128                                            });
129                                        }
130                                    }
131                                }
132                                // Transfers are carried in the event itself, bound to the same
133                                // execution instance as the abi call stacks above. This avoids a
134                                // separate slot-keyed lookup, which could return transfers from a
135                                // re-execution of the same slot and yield a hybrid response.
136                                for transfer in slot_transfers {
137                                    ret_transfers.push(TransferInfo {
138                                        from: transfer.from.to_string(),
139                                        to: transfer.to.to_string(),
140                                        amount: transfer.amount.to_raw(),
141                                        operation_id_or_asc_index: Some(
142                                            grpc_api::transfer_info::OperationIdOrAscIndex::OperationId(
143                                                transfer.op_id.to_string(),
144                                            ),
145                                        ),
146                                    });
147                                }
148                                let ret = grpc_api::NewSlotTransfersResponse {
149                                    slot: Some(massa_slot_execution_trace.slot.into()),
150                                    transfers: ret_transfers,
151                                };
152
153                                if let Err(e) = tx.send(Ok(ret)).await {
154                                    error!("failed to send new slot execution trace: {}", e);
155                                    break;
156                                }
157                            }
158                            Err(e) => {
159                                error!("error on receive new slot execution trace : {}", e)
160                            }
161                        }
162                    }
163                    // Receive a new message from the in_stream
164                    res = in_stream.next() => {
165                        match res {
166                            Some(res) => {
167                                match res {
168                                    Ok(message) => {
169                                        finality = message.finality_level();
170                                    },
171                                    Err(e) => {
172                                        // Any io error -> break
173                                        if let Some(io_err) = match_for_io_error(&e) {
174                                            warn!("client disconnected, broken pipe: {}", io_err);
175                                            break;
176                                        }
177                                        error!("{}", e);
178                                        if let Err(e2) = tx.send(Err(e)).await {
179                                            error!("failed to send back error response: {}", e2);
180                                            break;
181                                        }
182                                    }
183                                }
184                            }
185                            None => {
186                                // the client has disconnected
187                                break;
188                            }
189                        }
190                    }
191                }
192            }
193        }
194    });
195    let out_stream = tokio_stream::wrappers::ReceiverStream::new(rx);
196    Ok(Box::pin(out_stream) as NewSlotTransfersStreamType)
197}