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}