1use crate::error::GrpcError;
4use crate::server::MassaPublicGrpc;
5use crate::{EndorsementDraw, SlotDraw, SlotRange};
6
7use itertools::{izip, Itertools};
8use massa_execution_exports::mapping_grpc::{
9 to_event_filter, to_execution_query_response, to_querystate_filter,
10};
11use massa_execution_exports::{
12 ExecutionQueryRequest, ExecutionStackElement, ReadOnlyExecutionRequest, ReadOnlyExecutionTarget,
13};
14use massa_models::address::Address;
15use massa_models::amount::Amount;
16use massa_models::block::{Block, BlockGraphStatus};
17use massa_models::block_id::BlockId;
18use massa_models::config::CompactConfig;
19use massa_models::datastore::DatastoreDeserializer;
20use massa_models::endorsement::{EndorsementId, SecureShareEndorsement};
21use massa_models::operation::{OperationId, SecureShareOperation};
22use massa_models::prehash::{PreHashMap, PreHashSet};
23use massa_models::slot::Slot;
24use massa_models::timeslots::get_latest_block_slot_at_timestamp;
25use massa_proto_rs::massa::api::v1::{self as grpc_api};
26use massa_proto_rs::massa::model::v1::{self as grpc_model, read_only_execution_call};
27use massa_serialization::{DeserializeError, Deserializer};
28use massa_time::MassaTime;
29use massa_versioning::versioning_factory::{FactoryStrategy, VersioningFactory};
30use std::collections::HashSet;
31use std::str::FromStr;
32
33#[cfg(feature = "execution-trace")]
34use massa_execution_exports::types_trace_info::AbiTrace;
35#[cfg(feature = "execution-trace")]
36use massa_proto_rs::massa::api::v1::abi_call_stack_element_parent::CallStackElement;
37#[cfg(feature = "execution-trace")]
38use massa_proto_rs::massa::api::v1::{
39 AbiCallStack, AbiCallStackElement, AbiCallStackElementCall, AbiCallStackElementParent,
40 AscabiCallStack, DeferredCallAbiCallStack, GetOperationAbiCallStacksResponse,
41 GetSlotAbiCallStacksResponse, OperationAbiCallStack, SlotAbiCallStacks, TransferInfo,
42 TransferInfos,
43};
44
45pub(crate) fn execute_read_only_call(
47 grpc: &MassaPublicGrpc,
48 request: tonic::Request<grpc_api::ExecuteReadOnlyCallRequest>,
49) -> Result<grpc_api::ExecuteReadOnlyCallResponse, GrpcError> {
50 let call: grpc_model::ReadOnlyExecutionCall = request
51 .into_inner()
52 .call
53 .ok_or_else(|| GrpcError::InvalidArgument("no call provided".to_string()))?;
54
55 let caller_address = match call.caller_address {
56 Some(addr) => Address::from_str(&addr)?,
57 None => {
58 let now = MassaTime::now();
59 let keypair = grpc.keypair_factory.create(&(), FactoryStrategy::At(now))?;
60 Address::from_public_key(&keypair.get_public_key())
61 }
62 };
63
64 let mut call_stack = Vec::new();
65 let mut coins = None;
66 let target = if let Some(call_target) = call.target {
67 match call_target {
68 read_only_execution_call::Target::BytecodeCall(value) => {
69 let op_datastore = if value.operation_datastore.is_empty() {
70 None
71 } else {
72 let deserializer = DatastoreDeserializer::new(
73 grpc.grpc_config.max_op_datastore_entry_count,
74 grpc.grpc_config.max_op_datastore_key_length,
75 grpc.grpc_config.max_op_datastore_value_length,
76 );
77 match deserializer.deserialize::<DeserializeError>(&value.operation_datastore) {
78 Ok((_, deserialized)) => Some(deserialized),
79 Err(e) => {
80 return Err(GrpcError::InvalidArgument(format!(
81 "Datastore deserializing error: {}",
82 e
83 )))
84 }
85 }
86 };
87
88 call_stack.push(ExecutionStackElement {
89 address: caller_address,
90 coins: Default::default(),
91 owned_addresses: vec![caller_address],
92 operation_datastore: op_datastore,
93 });
94
95 ReadOnlyExecutionTarget::BytecodeExecution(value.bytecode)
96 }
97 read_only_execution_call::Target::FunctionCall(call) => {
98 let target_address = Address::from_str(&call.target_address)?;
99
100 coins = call
101 .coins
102 .map(|native_amount| {
103 Amount::from_mantissa_scale(native_amount.mantissa, native_amount.scale)
104 .map_err(|_| GrpcError::InvalidArgument("invalid amount".to_string()))
105 })
106 .transpose()?;
107
108 call_stack.push(ExecutionStackElement {
109 address: caller_address,
110 coins: Default::default(),
111 owned_addresses: vec![caller_address],
112 operation_datastore: None, });
114 call_stack.push(ExecutionStackElement {
117 address: target_address,
118 coins: coins.unwrap_or_default(),
119 owned_addresses: vec![target_address],
120 operation_datastore: None, });
122
123 ReadOnlyExecutionTarget::FunctionCall {
124 target_addr: Address::from_str(&call.target_address)?,
125 target_func: call.target_function,
126 parameter: call.parameter,
127 }
128 }
129 }
130 } else {
131 return Err(GrpcError::InvalidArgument(
132 "no call target provided".to_string(),
133 ));
134 };
135
136 let read_only_call = ReadOnlyExecutionRequest {
137 max_gas: call.max_gas,
138 call_stack,
139 target,
140 coins,
141 fee: call
142 .fee
143 .map(|native_amount| {
144 Amount::from_mantissa_scale(native_amount.mantissa, native_amount.scale)
145 .map_err(|_| GrpcError::InvalidArgument("invalid amount".to_string()))
146 })
147 .transpose()?,
148 };
149
150 if call.fee.is_some()
152 && read_only_call
153 .fee
154 .unwrap_or_default()
155 .checked_sub(grpc.grpc_config.minimal_fees)
156 .is_none()
157 {
158 return Err(GrpcError::InvalidArgument(format!(
159 "fee is too low provided: {} , minimal_fees required: {}",
160 read_only_call.fee.unwrap_or_default(),
161 grpc.grpc_config.minimal_fees
162 )));
163 }
164
165 let output = grpc
166 .execution_controller
167 .execute_readonly_request(read_only_call)?;
168
169 let result = grpc_model::ReadOnlyExecutionOutput {
170 out: Some(output.out.into()),
171 used_gas: output.gas_cost,
172 call_result: output.call_result,
173 };
174
175 Ok(grpc_api::ExecuteReadOnlyCallResponse {
176 output: Some(result),
177 })
178}
179
180pub(crate) fn get_blocks(
182 grpc: &MassaPublicGrpc,
183 request: tonic::Request<grpc_api::GetBlocksRequest>,
184) -> Result<grpc_api::GetBlocksResponse, GrpcError> {
185 let ids = request.into_inner().block_ids;
186
187 if ids.is_empty() {
188 return Err(GrpcError::InvalidArgument(
189 "no block id provided".to_string(),
190 ));
191 }
192
193 if ids.len() as u32 > grpc.grpc_config.max_block_ids_per_request {
194 return Err(GrpcError::InvalidArgument(format!(
195 "too many block ids received. Only a maximum of {} block ids are accepted per request",
196 grpc.grpc_config.max_block_ids_per_request
197 )));
198 }
199
200 let mut block_ids: Vec<BlockId> = ids
201 .into_iter()
202 .map(|id| {
203 BlockId::from_str(id.as_str())
204 .map_err(|_| GrpcError::InvalidArgument(format!("invalid block id: {}", id)))
205 })
206 .collect::<Result<_, _>>()?;
207
208 let mut blocks: Vec<Block> = Vec::with_capacity(block_ids.len());
209 {
210 let block_storage_lock = grpc.storage.read_blocks();
211 block_ids.retain(|id| {
212 if let Some(wrapped_block) = block_storage_lock.get(id) {
213 blocks.push(wrapped_block.content.clone());
214 return true;
215 };
216 false
217 });
218 }
219
220 let block_statuses = grpc.consensus_controller.get_block_statuses(&block_ids);
221
222 let result = blocks
223 .into_iter()
224 .zip(block_statuses)
225 .map(|(block, block_graph_status)| grpc_model::BlockWrapper {
226 block: Some(block.into()),
227 status: block_graph_status.into(),
228 })
229 .collect();
230
231 Ok(grpc_api::GetBlocksResponse {
232 wrapped_blocks: result,
233 })
234}
235
236pub(crate) fn get_datastore_entries(
238 grpc: &MassaPublicGrpc,
239 request: tonic::Request<grpc_api::GetDatastoreEntriesRequest>,
240) -> Result<grpc_api::GetDatastoreEntriesResponse, GrpcError> {
241 let inner_req = request.into_inner();
242
243 if inner_req.filters.is_empty() {
245 return Err(GrpcError::InvalidArgument("no filter provided".to_string()));
246 }
247
248 if inner_req.filters.len() as u64 > grpc.grpc_config.max_datastore_entries_per_request {
250 return Err(GrpcError::InvalidArgument(format!(
251 "too many datastore entries received. Only a maximum of {} datastore entries are accepted per request",
252 grpc.grpc_config.max_datastore_entries_per_request
253 )));
254 }
255
256 let filters: Vec<(Address, Vec<u8>)> = inner_req
257 .filters
258 .into_iter()
259 .map(|filter| {
260 let filter = filter
261 .filter
262 .ok_or_else(|| GrpcError::InvalidArgument("no filter provided".to_string()))?;
263 match filter {
264 grpc_api::get_datastore_entry_filter::Filter::AddressKey(addrs) => {
265 let add = Address::from_str(&addrs.address).map_err(|_| {
266 GrpcError::InvalidArgument(format!("invalid address: {}", addrs.address))
267 })?;
268 Ok((add, addrs.key))
269 }
270 }
271 })
272 .collect::<Result<Vec<(Address, Vec<u8>)>, GrpcError>>()?;
273
274 let entries = grpc
275 .execution_controller
276 .get_final_and_active_data_entry(filters)
277 .into_iter()
278 .map(|output| grpc_model::DatastoreEntry {
279 final_value: output.0.unwrap_or_default(),
280 candidate_value: output.1.unwrap_or_default(),
281 })
282 .collect();
283
284 Ok(grpc_api::GetDatastoreEntriesResponse {
285 datastore_entries: entries,
286 })
287}
288
289pub(crate) fn get_endorsements(
291 grpc: &MassaPublicGrpc,
292 request: tonic::Request<grpc_api::GetEndorsementsRequest>,
293) -> Result<grpc_api::GetEndorsementsResponse, GrpcError> {
294 let ids = request.into_inner().endorsement_ids;
295
296 if ids.is_empty() {
297 return Err(GrpcError::InvalidArgument(
298 "no endorsement id provided".to_string(),
299 ));
300 }
301
302 if ids.len() as u32 > grpc.grpc_config.max_endorsement_ids_per_request {
303 return Err(GrpcError::InvalidArgument(format!(
304 "too many endorsement ids received. Only a maximum of {} endorsement ids are accepted per request",
305 grpc.grpc_config.max_endorsement_ids_per_request
306 )));
307 }
308
309 let mut endorsement_ids: Vec<EndorsementId> = ids
310 .into_iter()
311 .map(|id| {
312 EndorsementId::from_str(id.as_str())
313 .map_err(|_| GrpcError::InvalidArgument(format!("invalid endorsement id: {}", id)))
314 })
315 .collect::<Result<_, _>>()?;
316
317 let mut secure_share_endorsements: Vec<SecureShareEndorsement> =
318 Vec::with_capacity(endorsement_ids.len());
319 {
320 let endorsement_storage_lock = grpc.storage.read_endorsements();
321 endorsement_ids.retain(|id| {
322 if let Some(wrapped_endorsement) = endorsement_storage_lock.get(id) {
323 secure_share_endorsements.push(wrapped_endorsement.clone());
324 return true;
325 };
326 false
327 });
328 }
329
330 let storage_info: Vec<(SecureShareEndorsement, PreHashSet<BlockId>)> = {
331 let read_blocks = grpc.storage.read_blocks();
332 secure_share_endorsements
333 .into_iter()
334 .map(|secure_share_operation| {
335 let ed_id = secure_share_operation.id;
336 (
337 secure_share_operation,
338 read_blocks
339 .get_blocks_by_endorsement(&ed_id)
340 .cloned()
341 .unwrap_or_default(),
342 )
343 })
344 .collect()
345 };
346
347 let in_pool = grpc
349 .pool_controller
350 .contains_endorsements(
351 &endorsement_ids,
352 Some(MassaTime::from_millis(
353 grpc.grpc_config.timeout.as_millis() as u64
354 )),
355 )
356 .unwrap_or_else(|_| vec![false; endorsement_ids.len()]);
357
358 let is_final: Vec<bool> = {
360 let involved_blocks: Vec<BlockId> = storage_info
361 .iter()
362 .flat_map(|(_ed, bs)| bs.iter())
363 .unique()
364 .cloned()
365 .collect();
366
367 let involved_block_statuses = grpc
368 .consensus_controller
369 .get_block_statuses(&involved_blocks);
370
371 let block_statuses: PreHashMap<BlockId, BlockGraphStatus> = involved_blocks
372 .into_iter()
373 .zip(involved_block_statuses)
374 .collect();
375
376 storage_info
377 .iter()
378 .map(|(_ed, bs)| {
379 bs.iter()
380 .any(|b| block_statuses.get(b) == Some(&BlockGraphStatus::Final))
381 })
382 .collect()
383 };
384
385 let mut result: Vec<grpc_model::EndorsementWrapper> = Vec::with_capacity(endorsement_ids.len());
387 let zipped_iterator = izip!(
388 storage_info.into_iter(),
389 in_pool.into_iter(),
390 is_final.into_iter()
391 );
392
393 for ((endorsement, in_blocks), in_pool, is_final) in zipped_iterator {
394 result.push(grpc_model::EndorsementWrapper {
395 in_pool,
396 is_final,
397 in_blocks: in_blocks
398 .into_iter()
399 .map(|block_id| block_id.to_string())
400 .collect(),
401 endorsement: Some(endorsement.into()),
402 });
403 }
404
405 Ok(grpc_api::GetEndorsementsResponse {
406 wrapped_endorsements: result,
407 })
408}
409
410pub(crate) fn get_stakers(
412 grpc: &MassaPublicGrpc,
413 request: tonic::Request<grpc_api::GetStakersRequest>,
414) -> Result<grpc_api::GetStakersResponse, GrpcError> {
415 let inner_req = request.into_inner();
416
417 let mut filter_opt = (None, None, None);
419
420 inner_req
422 .filters
423 .iter()
424 .for_each(|filter| match filter.filter {
425 Some(grpc_api::stakers_filter::Filter::MinRolls(min_rolls)) => {
426 filter_opt.0 = Some(min_rolls);
427 }
428 Some(grpc_api::stakers_filter::Filter::MaxRolls(max_rolls)) => {
429 filter_opt.1 = Some(max_rolls);
430 }
431 Some(grpc_api::stakers_filter::Filter::Limit(limit)) => {
432 filter_opt.2 = Some(limit);
433 }
434 None => {}
435 });
436
437 let now: MassaTime = MassaTime::now();
439
440 let latest_block_slot_at_timestamp_result = get_latest_block_slot_at_timestamp(
441 grpc.grpc_config.thread_count,
442 grpc.grpc_config.t0,
443 grpc.grpc_config.genesis_timestamp,
444 now,
445 );
446
447 let (cur_cycle, _cur_slot) = match latest_block_slot_at_timestamp_result {
448 Ok(Some(cur_slot)) if cur_slot.period <= grpc.grpc_config.last_start_period => (
449 Slot::new(grpc.grpc_config.last_start_period, 0)
450 .get_cycle(grpc.grpc_config.periods_per_cycle),
451 cur_slot,
452 ),
453 Ok(Some(cur_slot)) => (
454 cur_slot.get_cycle(grpc.grpc_config.periods_per_cycle),
455 cur_slot,
456 ),
457 Ok(None) => (0, Slot::new(0, 0)),
458 Err(e) => return Err(GrpcError::ModelsError(e)),
459 };
460
461 let mut staker_vec = grpc
463 .execution_controller
464 .get_cycle_active_rolls(cur_cycle)
465 .into_iter()
466 .filter_map(|(addr, rolls)| {
467 if let Some(min_rolls) = filter_opt.0 {
468 if rolls < min_rolls {
469 return None;
470 }
471 }
472 if let Some(max_rolls) = filter_opt.1 {
473 if rolls > max_rolls {
474 return None;
475 }
476 }
477 Some((addr.to_string(), rolls))
478 })
479 .collect::<Vec<(String, u64)>>();
480
481 staker_vec.sort_by_key(|&(_, roll_counts)| std::cmp::Reverse(roll_counts));
483
484 if let Some(limit) = filter_opt.2 {
485 staker_vec = staker_vec
486 .into_iter()
487 .take(limit as usize)
488 .collect::<Vec<(String, u64)>>();
489 }
490
491 let stakers = staker_vec
492 .into_iter()
493 .map(|(address, rolls)| grpc_model::StakerEntry { address, rolls })
494 .collect();
495
496 Ok(grpc_api::GetStakersResponse { stakers })
497}
498
499#[cfg(feature = "execution-trace")]
500pub fn into_element(abi_trace: &AbiTrace) -> AbiCallStackElementParent {
502 if abi_trace.sub_calls.is_none() {
503 AbiCallStackElementParent {
504 call_stack_element: Some(CallStackElement::Element(AbiCallStackElement {
505 name: abi_trace.name.clone(),
506 parameters: abi_trace
507 .parameters
508 .iter()
509 .map(|p| serde_json::to_string(p).unwrap_or_default())
510 .collect::<Vec<String>>(),
511 return_value: serde_json::to_string(&abi_trace.return_value).unwrap_or_default(),
512 })),
513 }
514 } else {
515 AbiCallStackElementParent {
516 call_stack_element: Some(CallStackElement::ElementCall(AbiCallStackElementCall {
517 name: abi_trace.name.clone(),
518 parameters: abi_trace
519 .parameters
520 .iter()
521 .map(|p| serde_json::to_string(p).unwrap_or_default())
522 .collect(),
523 return_value: serde_json::to_string(&abi_trace.return_value).unwrap_or_default(),
524 sub_calls: abi_trace
525 .sub_calls
526 .clone()
527 .unwrap_or_default()
528 .iter()
529 .map(into_element)
530 .collect(),
531 })),
532 }
533 }
534}
535
536#[cfg(feature = "execution-trace")]
537pub(crate) fn get_slot_transfers(
539 grpc: &MassaPublicGrpc,
540 request: tonic::Request<grpc_api::GetSlotTransfersRequest>,
541) -> Result<grpc_api::GetSlotTransfersResponse, GrpcError> {
542 use massa_proto_rs::massa::api::v1::GetSlotTransfersResponse;
543
544 let slots = request.into_inner().slots;
545
546 let mut transfer_each_slot: Vec<TransferInfos> = vec![];
547 for slot in slots {
548 let mut slot_transfers = TransferInfos {
549 slot: slot.clone().into(),
550 transfers: vec![],
551 };
552 let (abi_calls, direct_transfers) = grpc
556 .execution_controller
557 .get_slot_abi_call_stack_and_transfers(slot.clone().into());
558 if let Some(abi_calls) = abi_calls {
559 let abi_transfer_1 = "assembly_script_transfer_coins".to_string();
562 let abi_transfer_2 = "assembly_script_transfer_coins_for".to_string();
563 let abi_transfer_3 = "abi_transfer_coins".to_string();
564 let transfer_abi_names = vec![abi_transfer_1, abi_transfer_2, abi_transfer_3];
565 for (i, asc_call_stack) in abi_calls.asc_call_stacks.iter().enumerate() {
566 for abi_trace in asc_call_stack {
567 let only_transfer = abi_trace.flatten_filter(&transfer_abi_names);
568 for transfer in only_transfer {
569 let (t_from, t_to, t_amount) = transfer.parse_transfer();
570 slot_transfers.transfers.push(TransferInfo {
571 from: t_from.clone(),
572 to: t_to.clone(),
573 amount: t_amount,
574 operation_id_or_asc_index: Some(
575 grpc_api::transfer_info::OperationIdOrAscIndex::AscIndex(i as u64),
576 ),
577 });
578 }
579 }
580 }
581
582 for op_call_stack in abi_calls.operation_call_stacks {
583 let op_id = op_call_stack.0;
584 let op_call_stack = op_call_stack.1;
585 for abi_trace in op_call_stack {
586 let only_transfer = abi_trace.flatten_filter(&transfer_abi_names);
587 for transfer in only_transfer {
588 let (t_from, t_to, t_amount) = transfer.parse_transfer();
589 slot_transfers.transfers.push(TransferInfo {
590 from: t_from.clone(),
591 to: t_to.clone(),
592 amount: t_amount,
593 operation_id_or_asc_index: Some(
594 grpc_api::transfer_info::OperationIdOrAscIndex::OperationId(
595 op_id.to_string(),
596 ),
597 ),
598 });
599 }
600 }
601 }
602 }
603
604 if let Some(transfers) = direct_transfers {
605 for transfer in transfers {
606 slot_transfers.transfers.push(TransferInfo {
607 from: transfer.from.to_string(),
608 to: transfer.to.to_string(),
609 amount: transfer.amount.to_raw(),
610 operation_id_or_asc_index: Some(
611 grpc_api::transfer_info::OperationIdOrAscIndex::OperationId(
612 transfer.op_id.to_string(),
613 ),
614 ),
615 });
616 }
617 }
618
619 transfer_each_slot.push(slot_transfers);
620 }
621
622 Ok(GetSlotTransfersResponse { transfer_each_slot })
623}
624
625#[cfg(feature = "execution-trace")]
626pub(crate) fn get_operation_abi_call_stacks(
628 grpc: &MassaPublicGrpc,
629 request: tonic::Request<grpc_api::GetOperationAbiCallStacksRequest>,
630) -> Result<grpc_api::GetOperationAbiCallStacksResponse, GrpcError> {
631 let op_ids_ = request.into_inner().operation_ids;
632
633 let op_ids: Vec<OperationId> = op_ids_
634 .iter()
635 .map(|o| OperationId::from_str(o))
636 .collect::<Result<Vec<_>, _>>()?;
637
638 let mut elements = vec![];
639 for op_id in op_ids {
640 let abi_traces_ = grpc
641 .execution_controller
642 .get_operation_abi_call_stack(op_id);
643 if let Some(abi_traces) = abi_traces_ {
644 for abi_trace in abi_traces.iter() {
645 elements.push(into_element(abi_trace));
646 }
647 } else {
648 elements.push(AbiCallStackElementParent {
649 call_stack_element: None,
650 })
651 }
652 }
653
654 let resp = GetOperationAbiCallStacksResponse {
655 call_stacks: vec![AbiCallStack {
656 call_stack: elements,
657 }],
658 };
659
660 Ok(resp)
661}
662
663#[cfg(feature = "execution-trace")]
664pub(crate) fn get_slot_abi_call_stacks(
665 grpc: &MassaPublicGrpc,
666 request: tonic::Request<grpc_api::GetSlotAbiCallStacksRequest>,
667) -> Result<grpc_api::GetSlotAbiCallStacksResponse, GrpcError> {
668 let slots = request.into_inner().slots;
669
670 let mut slot_elements = vec![];
671 for slot in slots {
672 let call_stack_ = grpc
673 .execution_controller
674 .get_slot_abi_call_stack(slot.into());
675
676 let mut slot_abi_call_stacks = SlotAbiCallStacks {
677 asc_call_stacks: vec![],
678 deferred_call_stacks: vec![],
679 operation_call_stacks: vec![],
680 };
681
682 if let Some(call_stack) = call_stack_ {
683 for (call_id, deferred_call_stack) in call_stack.deferred_call_stacks {
684 slot_abi_call_stacks
685 .deferred_call_stacks
686 .push(DeferredCallAbiCallStack {
687 deferred_call_id: call_id.to_string(),
688 call_stack: deferred_call_stack.iter().map(into_element).collect(),
689 })
690 }
691 for (i, asc_call_stack) in call_stack.asc_call_stacks.into_iter().enumerate() {
692 slot_abi_call_stacks.asc_call_stacks.push(AscabiCallStack {
693 index: i as u64,
694 call_stack: asc_call_stack.iter().map(into_element).collect(),
695 })
696 }
697 for (op_id, op_call_stack) in call_stack.operation_call_stacks {
698 slot_abi_call_stacks
699 .operation_call_stacks
700 .push(OperationAbiCallStack {
701 operation_id: op_id.to_string(),
702 call_stack: op_call_stack.iter().map(into_element).collect(),
703 })
704 }
705 }
706 slot_elements.push(slot_abi_call_stacks);
707 }
708
709 let resp = GetSlotAbiCallStacksResponse {
710 slot_call_stacks: slot_elements,
711 };
712
713 Ok(resp)
714}
715
716pub(crate) fn get_next_block_best_parents(
718 grpc: &MassaPublicGrpc,
719 _request: tonic::Request<grpc_api::GetNextBlockBestParentsRequest>,
720) -> Result<grpc_api::GetNextBlockBestParentsResponse, GrpcError> {
721 let block_parents = grpc
722 .consensus_controller
723 .get_best_parents()
724 .into_iter()
725 .map(|p| grpc_model::BlockParent {
726 block_id: p.0.to_string(),
727 period: p.1,
728 })
729 .collect();
730 Ok(grpc_api::GetNextBlockBestParentsResponse { block_parents })
731}
732
733pub(crate) fn get_operations(
735 grpc: &MassaPublicGrpc,
736 request: tonic::Request<grpc_api::GetOperationsRequest>,
737) -> Result<grpc_api::GetOperationsResponse, GrpcError> {
738 let operation_ids = request.into_inner().operation_ids;
739
740 if operation_ids.is_empty() {
741 return Err(GrpcError::InvalidArgument(
742 "no operations ids specified".to_string(),
743 ));
744 }
745
746 if operation_ids.len() as u32 > grpc.grpc_config.max_operation_ids_per_request {
747 return Err(GrpcError::InvalidArgument(format!("too many operations received. Only a maximum of {} operations are accepted per request", grpc.grpc_config.max_operation_ids_per_request)));
748 }
749
750 let operation_ids: Vec<OperationId> = operation_ids
751 .into_iter()
752 .map(|id| {
753 OperationId::from_str(id.as_str())
754 .map_err(|_| GrpcError::InvalidArgument(format!("invalid operation id: {}", id)))
755 })
756 .collect::<Result<_, _>>()?;
757
758 let secure_share_operations: Vec<SecureShareOperation> = {
759 let read_ops = grpc.storage.read_operations();
760 operation_ids
761 .iter()
762 .filter_map(|id| read_ops.get(id).cloned())
763 .collect()
764 };
765
766 let storage_info: Vec<(SecureShareOperation, PreHashSet<BlockId>)> = {
767 let read_blocks = grpc.storage.read_blocks();
768 secure_share_operations
769 .into_iter()
770 .map(|secure_share_operation| {
771 let op_id = secure_share_operation.id;
772 (
773 secure_share_operation,
774 read_blocks
775 .get_blocks_by_operation(&op_id)
776 .cloned()
777 .unwrap_or_default(),
778 )
779 })
780 .collect()
781 };
782
783 let operations: Vec<grpc_model::OperationWrapper> = storage_info
784 .into_iter()
785 .map(|secure_share| {
786 let (secure_share_operation, block_ids) = secure_share;
787 grpc_model::OperationWrapper {
788 thread: secure_share_operation
789 .content_creator_address
790 .get_thread(grpc.grpc_config.thread_count) as u32,
791 operation: Some(secure_share_operation.into()),
792 block_ids: block_ids.into_iter().map(|id| id.to_string()).collect(),
793 }
794 })
795 .collect();
796
797 Ok(grpc_api::GetOperationsResponse {
798 wrapped_operations: operations,
799 })
800}
801
802pub(crate) fn get_sc_execution_events(
804 grpc: &MassaPublicGrpc,
805 request: tonic::Request<grpc_api::GetScExecutionEventsRequest>,
806) -> Result<grpc_api::GetScExecutionEventsResponse, GrpcError> {
807 let event_filter = to_event_filter(request.into_inner().filters)?;
808 let events: Vec<grpc_model::ScExecutionEvent> = grpc
809 .execution_controller
810 .get_filtered_sc_output_event(event_filter)
811 .into_iter()
812 .map(|event| event.into())
813 .collect();
814
815 Ok(grpc_api::GetScExecutionEventsResponse { events })
816}
817
818pub(crate) fn get_selector_draws(
820 grpc: &MassaPublicGrpc,
821 request: tonic::Request<grpc_api::GetSelectorDrawsRequest>,
822) -> Result<grpc_api::GetSelectorDrawsResponse, GrpcError> {
823 let inner_req = request.into_inner();
824 if inner_req.filters.len() as u32 > grpc.grpc_config.max_filters_per_request {
825 return Err(GrpcError::InvalidArgument(format!(
826 "too many filters received. Only a maximum of {} filters are accepted per request",
827 grpc.grpc_config.max_filters_per_request
828 )));
829 }
830
831 let mut addresses_filter: Option<PreHashSet<Address>> = None;
832 let mut slot_ranges_filter: Option<HashSet<SlotRange>> = None;
833 for query in inner_req.filters.into_iter() {
835 if let Some(filter) = query.filter {
836 match filter {
837 grpc_api::selector_draws_filter::Filter::Addresses(addrs) => {
838 if addrs.addresses.len() as u32 > grpc.grpc_config.max_addresses_per_request {
839 return Err(GrpcError::InvalidArgument(format!(
840 "too many addresses received. Only a maximum of {} addresses are accepted per request",
841 grpc.grpc_config.max_addresses_per_request
842 )));
843 }
844 let addresses = addresses_filter.get_or_insert_with(PreHashSet::default);
845 for address in addrs.addresses {
846 addresses.insert(Address::from_str(&address).map_err(|_| {
847 GrpcError::InvalidArgument(format!("invalid address: {}", address))
848 })?);
849 }
850 }
851 grpc_api::selector_draws_filter::Filter::SlotRange(s_range) => {
852 let slot_ranges = slot_ranges_filter.get_or_insert_with(HashSet::new);
853 if slot_ranges.len() as u32 >= grpc.grpc_config.max_slot_ranges_per_request {
854 return Err(GrpcError::InvalidArgument(format!(
855 "too many slot ranges received. Only a maximum of {} slot ranges are accepted per request",
856 grpc.grpc_config.max_slot_ranges_per_request
857 )));
858 }
859
860 let start_slot: Option<Slot> = s_range.start_slot.map(|s| s.into());
861 let end_slot: Option<Slot> = s_range.end_slot.map(|s| s.into());
862
863 let slot_range = SlotRange {
864 start_slot,
865 end_slot,
866 };
867 slot_range.check()?;
868 slot_ranges.insert(slot_range);
869 }
870 }
871 }
872 }
873
874 let selection_draws: HashSet<SlotDraw> = if let Some(slot_ranges) = slot_ranges_filter {
876 if slot_ranges.is_empty() {
877 return Err(GrpcError::InvalidArgument(
878 "at least, one slot range is required".to_string(),
879 ));
880 }
881
882 let mut start_slot = Slot::new(0, 0); let mut end_slot = Slot::new(u64::MAX, grpc.grpc_config.thread_count - 1); for slot_range in &slot_ranges {
885 start_slot = start_slot.max(slot_range.start_slot.unwrap_or_else(|| Slot::new(0, 0)));
886 end_slot = end_slot.min(
887 slot_range
888 .end_slot
889 .unwrap_or_else(|| Slot::new(u64::MAX, grpc.grpc_config.thread_count - 1)),
890 );
891 }
892 end_slot = end_slot.max(start_slot);
893
894 let mut restrict_to_addresses: Option<&PreHashSet<Address>> = None;
896 if let Some(addresses) = &addresses_filter {
897 if !addresses.is_empty() {
898 restrict_to_addresses = Some(addresses);
899 }
900 }
901
902 grpc.selector_controller
903 .get_available_selections_in_range(start_slot..=end_slot, restrict_to_addresses)
904 .unwrap_or_default()
905 .into_iter()
906 .map(|(v_slot, v_sel)| {
907 let endorsement_producers: Vec<EndorsementDraw> = v_sel
908 .endorsements
909 .into_iter()
910 .enumerate()
911 .map(|(index, endo_sel)| EndorsementDraw {
912 index: index as u64,
913 producer: endo_sel.to_string(),
914 })
915 .collect();
916
917 SlotDraw {
918 slot: Some(v_slot),
919 block_producer: Some(v_sel.producer.to_string()),
920 endorsement_draws: endorsement_producers,
921 }
922 })
923 .collect()
924 } else {
925 return Err(GrpcError::InvalidArgument(
926 "at least, one slot range is required".to_string(),
927 ));
928 };
929
930 Ok(grpc_api::GetSelectorDrawsResponse {
931 draws: selection_draws.into_iter().map(Into::into).collect(),
932 })
933}
934
935pub(crate) fn get_status(
937 grpc: &MassaPublicGrpc,
938 _request: tonic::Request<grpc_api::GetStatusRequest>,
939) -> Result<grpc_api::GetStatusResponse, GrpcError> {
940 let config = CompactConfig::default();
941 let now = MassaTime::now();
942 let last_slot = get_latest_block_slot_at_timestamp(
943 grpc.grpc_config.thread_count,
944 grpc.grpc_config.t0,
945 grpc.grpc_config.genesis_timestamp,
946 now,
947 )?;
948
949 let current_cycle = last_slot
950 .unwrap_or_else(|| Slot::new(0, 0))
951 .get_cycle(grpc.grpc_config.periods_per_cycle);
952 let cycle_duration = grpc
953 .grpc_config
954 .t0
955 .checked_mul(grpc.grpc_config.periods_per_cycle)?;
956 let current_cycle_time = if current_cycle == 0 {
957 grpc.grpc_config.genesis_timestamp
958 } else {
959 cycle_duration
960 .checked_mul(current_cycle)
961 .and_then(|elapsed_time_before_current_cycle| {
962 grpc.grpc_config
963 .genesis_timestamp
964 .checked_add(elapsed_time_before_current_cycle)
965 })?
966 };
967 let next_cycle_time = current_cycle_time.checked_add(cycle_duration)?;
968 let empty_request = ExecutionQueryRequest {
970 requests: vec![],
971 max_response_size: grpc.grpc_config.max_encoding_message_size,
972 max_event_count: None,
973 query_state_deadline_ms: None,
974 };
975 let state = grpc.execution_controller.query_state(empty_request);
976
977 let current_mip_version = grpc.keypair_factory.mip_store.get_network_version_current();
978
979 let status = grpc_model::PublicStatus {
980 node_id: grpc.node_id.to_string(),
981 version: grpc.version.to_string(),
982 current_time: Some(now.into()),
983 current_cycle,
984 current_cycle_time: Some(current_cycle_time.into()),
985 next_cycle_time: Some(next_cycle_time.into()),
986 last_executed_final_slot: Some(state.final_cursor.into()),
987 last_executed_speculative_slot: Some(state.candidate_cursor.into()),
988 final_state_fingerprint: state.final_state_fingerprint.to_string(),
989 config: Some(config.into()),
990 chain_id: grpc.grpc_config.chain_id,
991 minimal_fees: Some(grpc.grpc_config.minimal_fees.into()),
992 current_mip_version,
993 max_datastore_keys_query: grpc.grpc_config.max_datastore_keys_queries,
994 };
995
996 Ok(grpc_api::GetStatusResponse {
997 status: Some(status),
998 })
999}
1000
1001pub(crate) fn get_transactions_throughput(
1003 grpc: &MassaPublicGrpc,
1004 _request: tonic::Request<grpc_api::GetTransactionsThroughputRequest>,
1005) -> Result<grpc_api::GetTransactionsThroughputResponse, GrpcError> {
1006 let stats = grpc.execution_controller.get_stats();
1007 let nb_sec_range = stats
1008 .time_window_end
1009 .saturating_sub(stats.time_window_start)
1010 .to_duration()
1011 .as_secs();
1012
1013 let throughput = stats
1014 .final_executed_operations_count
1015 .checked_div(nb_sec_range as usize)
1016 .unwrap_or_default() as u32;
1017
1018 Ok(grpc_api::GetTransactionsThroughputResponse { throughput })
1019}
1020
1021pub(crate) fn query_state(
1023 grpc: &MassaPublicGrpc,
1024 request: tonic::Request<grpc_api::QueryStateRequest>,
1025) -> Result<grpc_api::QueryStateResponse, GrpcError> {
1026 let queries = request.into_inner().queries;
1027 if queries.is_empty() {
1028 return Err(GrpcError::InvalidArgument(
1029 "no query items specified".to_string(),
1030 ));
1031 }
1032 if queries.len() as u32 > grpc.grpc_config.max_query_items_per_request {
1033 return Err(GrpcError::InvalidArgument(format!("too many query items received. Only a maximum of {} operations are accepted per request", grpc.grpc_config.max_query_items_per_request)));
1034 }
1035
1036 let queries = queries
1037 .into_iter()
1038 .map(|q| {
1039 to_querystate_filter(
1040 q,
1041 grpc.grpc_config.max_datastore_keys_queries,
1042 grpc.grpc_config.max_datastore_key_length,
1043 )
1044 })
1045 .collect::<Result<Vec<_>, _>>()?;
1046
1047 let response = grpc
1048 .execution_controller
1049 .query_state(ExecutionQueryRequest {
1050 requests: queries,
1051 max_response_size: grpc.grpc_config.max_encoding_message_size,
1052 max_event_count: Some(grpc.grpc_config.max_event_per_query as usize),
1053 query_state_deadline_ms: grpc.grpc_config.query_state_deadline_ms,
1054 });
1055
1056 Ok(grpc_api::QueryStateResponse {
1057 final_cursor: Some(response.final_cursor.into()),
1058 candidate_cursor: Some(response.candidate_cursor.into()),
1059 final_state_fingerprint: response.final_state_fingerprint.to_string(),
1060 responses: response
1061 .responses
1062 .into_iter()
1063 .map(to_execution_query_response)
1064 .collect(),
1065 })
1066}
1067
1068pub(crate) fn search_blocks(
1070 grpc: &MassaPublicGrpc,
1071 request: tonic::Request<grpc_api::SearchBlocksRequest>,
1072) -> Result<grpc_api::SearchBlocksResponse, GrpcError> {
1073 let inner_req = request.into_inner();
1074 if inner_req.filters.len() as u32 > grpc.grpc_config.max_filters_per_request {
1075 return Err(GrpcError::InvalidArgument(format!(
1076 "too many filters received. Only a maximum of {} filters are accepted per request",
1077 grpc.grpc_config.max_filters_per_request
1078 )));
1079 }
1080
1081 let mut block_ids_filter: Option<PreHashSet<BlockId>> = None;
1082 let mut addresses_filter: Option<PreHashSet<Address>> = None;
1083 let mut slot_ranges_filter: Option<HashSet<SlotRange>> = None;
1084
1085 for query in inner_req.filters.into_iter() {
1087 if let Some(filter) = query.filter {
1088 match filter {
1089 grpc_api::search_blocks_filter::Filter::BlockIds(ids) => {
1090 if ids.block_ids.len() as u32 > grpc.grpc_config.max_block_ids_per_request {
1091 return Err(GrpcError::InvalidArgument(format!(
1092 "too many block ids received. Only a maximum of {} block ids are accepted per request",
1093 grpc.grpc_config.max_block_ids_per_request
1094 )));
1095 }
1096 let block_ids = block_ids_filter.get_or_insert_with(PreHashSet::default);
1097 for block_id in ids.block_ids {
1098 block_ids.insert(BlockId::from_str(&block_id).map_err(|_| {
1099 GrpcError::InvalidArgument(format!("invalid block id: {}", block_id))
1100 })?);
1101 }
1102 }
1103 grpc_api::search_blocks_filter::Filter::Addresses(addrs) => {
1104 if addrs.addresses.len() as u32 > grpc.grpc_config.max_addresses_per_request {
1105 return Err(GrpcError::InvalidArgument(format!(
1106 "too many addresses received. Only a maximum of {} addresses are accepted per request",
1107 grpc.grpc_config.max_addresses_per_request
1108 )));
1109 }
1110 let addresses = addresses_filter.get_or_insert_with(PreHashSet::default);
1111 for address in addrs.addresses {
1112 addresses.insert(Address::from_str(&address).map_err(|_| {
1113 GrpcError::InvalidArgument(format!("invalid address: {}", address))
1114 })?);
1115 }
1116 }
1117 grpc_api::search_blocks_filter::Filter::SlotRange(s_range) => {
1118 let slot_ranges = slot_ranges_filter.get_or_insert_with(HashSet::new);
1119 if slot_ranges.len() as u32 >= grpc.grpc_config.max_slot_ranges_per_request {
1120 return Err(GrpcError::InvalidArgument(format!(
1121 "too many slot ranges received. Only a maximum of {} slot ranges are accepted per request",
1122 grpc.grpc_config.max_slot_ranges_per_request
1123 )));
1124 }
1125
1126 let start_slot: Option<Slot> = s_range.start_slot.map(|s| s.into());
1127 let end_slot: Option<Slot> = s_range.end_slot.map(|s| s.into());
1128
1129 let slot_range = SlotRange {
1130 start_slot,
1131 end_slot,
1132 };
1133 slot_range.check()?;
1134 slot_ranges.insert(slot_range);
1135 }
1136 }
1137 }
1138 }
1139
1140 if block_ids_filter.is_none() && addresses_filter.is_none() && slot_ranges_filter.is_none() {
1142 return Err(GrpcError::InvalidArgument("no filter provided".to_string()));
1143 }
1144
1145 let mut res: Option<PreHashSet<BlockId>> = None;
1146
1147 if let Some(mut b_ids) = block_ids_filter {
1149 let read_lock = grpc.storage.read_blocks();
1150 b_ids.retain(|id: &BlockId| read_lock.contains(id));
1151
1152 res = Some(b_ids);
1153 }
1154
1155 if let Some(addrs) = addresses_filter {
1157 let b_ids: PreHashSet<BlockId> = {
1158 let read_lock = grpc.storage.read_blocks();
1159 let mut b_ids: PreHashSet<BlockId> = PreHashSet::default();
1160 for addr in addrs {
1161 if let Some(addr_b_ids) = read_lock.get_blocks_created_by(&addr) {
1162 b_ids.extend(addr_b_ids.clone());
1163 }
1164 }
1165
1166 b_ids
1167 };
1168 if let Some(block_ids) = res.as_mut() {
1169 block_ids.retain(|id: &BlockId| b_ids.contains(id));
1170 } else {
1171 res = Some(b_ids)
1172 }
1173 }
1174
1175 if let Some(slot_ranges) = slot_ranges_filter {
1177 let mut start_slot = Slot::new(0, 0); let mut end_slot = Slot::new(u64::MAX, grpc.grpc_config.thread_count - 1); for slot_range in &slot_ranges {
1180 start_slot = start_slot.max(slot_range.start_slot.unwrap_or_else(|| Slot::new(0, 0)));
1181 end_slot = end_slot.min(
1182 slot_range
1183 .end_slot
1184 .unwrap_or_else(|| Slot::new(u64::MAX, grpc.grpc_config.thread_count - 1)),
1185 );
1186 }
1187 end_slot = end_slot.max(start_slot);
1188
1189 let read_lock = grpc.storage.read_blocks();
1190 let b_ids: PreHashSet<BlockId> =
1191 read_lock.aggregate_blocks_by_slot_range(start_slot..end_slot);
1192
1193 if let Some(block_ids) = res.as_mut() {
1194 block_ids.retain(|id: &BlockId| b_ids.contains(id));
1195 } else {
1196 res = Some(b_ids)
1197 }
1198 }
1199
1200 let block_ids: Vec<BlockId> = res.unwrap_or_default().into_iter().collect();
1201
1202 if block_ids.is_empty() {
1203 return Ok(grpc_api::SearchBlocksResponse {
1204 block_infos: vec![],
1205 });
1206 }
1207
1208 let blocks_status = grpc.consensus_controller.get_block_statuses(&block_ids);
1209
1210 let result = block_ids
1211 .iter()
1212 .zip(blocks_status)
1213 .map(|(block_id, block_graph_status)| grpc_model::BlockInfo {
1214 block_id: block_id.to_string(),
1215 status: block_graph_status.into(),
1216 })
1217 .collect();
1218
1219 Ok(grpc_api::SearchBlocksResponse {
1220 block_infos: result,
1221 })
1222}
1223
1224pub(crate) fn search_endorsements(
1226 grpc: &MassaPublicGrpc,
1227 request: tonic::Request<grpc_api::SearchEndorsementsRequest>,
1228) -> Result<grpc_api::SearchEndorsementsResponse, GrpcError> {
1229 let inner_req = request.into_inner();
1230 if inner_req.filters.len() as u32 > grpc.grpc_config.max_filters_per_request {
1231 return Err(GrpcError::InvalidArgument(format!(
1232 "too many filters received. Only a maximum of {} filters are accepted per request",
1233 grpc.grpc_config.max_filters_per_request
1234 )));
1235 }
1236
1237 let mut endorsement_ids_filter: Option<PreHashSet<EndorsementId>> = None;
1238 let mut addresses_filter: Option<PreHashSet<Address>> = None;
1239 let mut block_ids_filter: Option<PreHashSet<BlockId>> = None;
1240
1241 for query in inner_req.filters.into_iter() {
1243 if let Some(filter) = query.filter {
1244 match filter {
1245 grpc_api::search_endorsements_filter::Filter::EndorsementIds(ids) => {
1246 if ids.endorsement_ids.len() as u32
1247 > grpc.grpc_config.max_endorsement_ids_per_request
1248 {
1249 return Err(GrpcError::InvalidArgument(format!(
1250 "too many endorsement ids received. Only a maximum of {} endorsement ids are accepted per request",
1251 grpc.grpc_config.max_endorsement_ids_per_request
1252 )));
1253 }
1254 let endorsement_ids =
1255 endorsement_ids_filter.get_or_insert_with(PreHashSet::default);
1256 for id in ids.endorsement_ids {
1257 endorsement_ids.insert(EndorsementId::from_str(&id).map_err(|_| {
1258 GrpcError::InvalidArgument(format!("invalid endorsement id: {}", id))
1259 })?);
1260 }
1261 }
1262 grpc_api::search_endorsements_filter::Filter::Addresses(addrs) => {
1263 if addrs.addresses.len() as u32 > grpc.grpc_config.max_addresses_per_request {
1264 return Err(GrpcError::InvalidArgument(format!(
1265 "too many addresses received. Only a maximum of {} addresses are accepted per request",
1266 grpc.grpc_config.max_addresses_per_request
1267 )));
1268 }
1269 let addresses = addresses_filter.get_or_insert_with(PreHashSet::default);
1270 for address in addrs.addresses {
1271 addresses.insert(Address::from_str(&address).map_err(|_| {
1272 GrpcError::InvalidArgument(format!("invalid address: {}", address))
1273 })?);
1274 }
1275 }
1276 grpc_api::search_endorsements_filter::Filter::BlockIds(ids) => {
1277 if ids.block_ids.len() as u32 > grpc.grpc_config.max_block_ids_per_request {
1278 return Err(GrpcError::InvalidArgument(format!(
1279 "too many block ids received. Only a maximum of {} block ids are accepted per request",
1280 grpc.grpc_config.max_block_ids_per_request
1281 )));
1282 }
1283 let block_ids = block_ids_filter.get_or_insert_with(PreHashSet::default);
1284 for block_id in ids.block_ids {
1285 block_ids.insert(BlockId::from_str(&block_id).map_err(|_| {
1286 GrpcError::InvalidArgument(format!("invalid block id: {}", block_id))
1287 })?);
1288 }
1289 }
1290 }
1291 }
1292 }
1293
1294 if endorsement_ids_filter.is_none() && addresses_filter.is_none() && block_ids_filter.is_none()
1296 {
1297 return Err(GrpcError::InvalidArgument("no filter provided".to_string()));
1298 }
1299
1300 let mut eds_ids: Option<PreHashSet<EndorsementId>> = None;
1301
1302 if let Some(mut e_ids) = endorsement_ids_filter {
1304 let read_lock = grpc.storage.read_endorsements();
1305 e_ids.retain(|id: &EndorsementId| read_lock.contains(id));
1306 eds_ids = Some(e_ids);
1307 }
1308
1309 if let Some(addrs) = addresses_filter {
1311 let e_ids: PreHashSet<EndorsementId> = {
1312 let mut e_ids: PreHashSet<EndorsementId> = PreHashSet::default();
1313 let read_lock = grpc.storage.read_endorsements();
1314 for addr in addrs {
1315 if let Some(addr_e_ids) = read_lock.get_endorsements_created_by(&addr) {
1316 e_ids.extend(addr_e_ids.clone());
1317 }
1318 }
1319
1320 e_ids
1321 };
1322 if let Some(endorsement_ids) = eds_ids.as_mut() {
1323 endorsement_ids.retain(|id: &EndorsementId| e_ids.contains(id));
1324 } else {
1325 eds_ids = Some(e_ids)
1326 }
1327 }
1328
1329 if let Some(b_ids) = block_ids_filter {
1331 let mut e_ids: PreHashSet<EndorsementId> = PreHashSet::default();
1332 let read_lock = grpc.storage.read_blocks();
1333 for block_id in b_ids {
1334 if let Some(wrapped_block) = read_lock.get(&block_id) {
1335 let b_endorsements: PreHashSet<EndorsementId> = wrapped_block
1336 .content
1337 .header
1338 .content
1339 .endorsements
1340 .iter()
1341 .map(|wrapped_endorsement| wrapped_endorsement.id)
1342 .collect();
1343 e_ids.extend(&b_endorsements);
1344 }
1345 }
1346
1347 if let Some(endorsement_ids) = eds_ids.as_mut() {
1348 endorsement_ids.retain(|id: &EndorsementId| e_ids.contains(id));
1349 } else {
1350 eds_ids = Some(e_ids)
1351 }
1352 }
1353
1354 let storage_info: Vec<(EndorsementId, PreHashSet<BlockId>)> = {
1355 let read_blocks_lock = grpc.storage.read_blocks();
1356 if let Some(endorsement_ids) = eds_ids {
1357 endorsement_ids
1358 .into_iter()
1359 .map(|id| {
1360 let block_ids = read_blocks_lock
1361 .get_blocks_by_endorsement(&id)
1362 .cloned()
1363 .unwrap_or_default();
1364
1365 (id, block_ids)
1366 })
1367 .collect()
1368 } else {
1369 return Ok(grpc_api::SearchEndorsementsResponse {
1370 endorsement_infos: Vec::new(),
1371 });
1372 }
1373 };
1374
1375 let e_ids: Vec<EndorsementId> = storage_info.iter().map(|(ed, _)| *ed).collect();
1377
1378 let in_pool = grpc
1380 .pool_controller
1381 .contains_endorsements(
1382 &e_ids,
1383 Some(MassaTime::from_millis(
1384 grpc.grpc_config.timeout.as_millis() as u64
1385 )),
1386 )
1387 .unwrap_or_else(|_| vec![false; e_ids.len()]);
1388
1389 let is_final: Vec<bool> = {
1391 let involved_blocks: Vec<BlockId> = storage_info
1392 .iter()
1393 .flat_map(|(_ed, bs)| bs.iter())
1394 .unique()
1395 .cloned()
1396 .collect();
1397
1398 let involved_block_statuses = grpc
1399 .consensus_controller
1400 .get_block_statuses(&involved_blocks);
1401
1402 let block_statuses: PreHashMap<BlockId, BlockGraphStatus> = involved_blocks
1403 .into_iter()
1404 .zip(involved_block_statuses)
1405 .collect();
1406 storage_info
1407 .iter()
1408 .map(|(_ed, bs)| {
1409 bs.iter()
1410 .any(|b| block_statuses.get(b) == Some(&BlockGraphStatus::Final))
1411 })
1412 .collect()
1413 };
1414
1415 let mut res: Vec<grpc_model::EndorsementInfo> = Vec::with_capacity(e_ids.len());
1417 let zipped_iterator = izip!(
1418 storage_info.into_iter(),
1419 in_pool.into_iter(),
1420 is_final.into_iter()
1421 );
1422
1423 for ((e_id, in_blocks), in_pool, is_final) in zipped_iterator {
1424 res.push(grpc_model::EndorsementInfo {
1425 in_pool,
1426 is_final,
1427 in_blocks: in_blocks
1428 .into_iter()
1429 .map(|block_id| block_id.to_string())
1430 .collect(),
1431 endorsement_id: e_id.to_string(),
1432 });
1433 }
1434
1435 Ok(grpc_api::SearchEndorsementsResponse {
1436 endorsement_infos: res,
1437 })
1438}
1439
1440pub(crate) fn search_operations(
1442 grpc: &MassaPublicGrpc,
1443 request: tonic::Request<grpc_api::SearchOperationsRequest>,
1444) -> Result<grpc_api::SearchOperationsResponse, GrpcError> {
1445 let inner_req: grpc_api::SearchOperationsRequest = request.into_inner();
1446 if inner_req.filters.len() as u32 > grpc.grpc_config.max_filters_per_request {
1447 return Err(GrpcError::InvalidArgument(format!(
1448 "too many filters received. Only a maximum of {} filters are accepted per request",
1449 grpc.grpc_config.max_filters_per_request
1450 )));
1451 }
1452 let mut operation_ids_filter: Option<PreHashSet<OperationId>> = None;
1453 let mut addresses_filter: Option<PreHashSet<Address>> = None;
1454
1455 for query in inner_req.filters.into_iter() {
1457 if let Some(filter) = query.filter {
1458 match filter {
1459 grpc_api::search_operations_filter::Filter::OperationIds(ids) => {
1460 if ids.operation_ids.len() as u32
1461 > grpc.grpc_config.max_operation_ids_per_request
1462 {
1463 return Err(GrpcError::InvalidArgument(format!(
1464 "too many operation ids received. Only a maximum of {} operation ids are accepted per request",
1465 grpc.grpc_config.max_operation_ids_per_request
1466 )));
1467 }
1468 let operation_ids =
1469 operation_ids_filter.get_or_insert_with(PreHashSet::default);
1470 for id in ids.operation_ids {
1471 operation_ids.insert(OperationId::from_str(&id).map_err(|_| {
1472 GrpcError::InvalidArgument(format!("invalid operation id: {}", id))
1473 })?);
1474 }
1475 }
1476 grpc_api::search_operations_filter::Filter::Addresses(addrs) => {
1477 if addrs.addresses.len() as u32 > grpc.grpc_config.max_addresses_per_request {
1478 return Err(GrpcError::InvalidArgument(format!(
1479 "too many addresses received. Only a maximum of {} addresses are accepted per request",
1480 grpc.grpc_config.max_addresses_per_request
1481 )));
1482 }
1483 let addresses = addresses_filter.get_or_insert_with(PreHashSet::default);
1484 for address in addrs.addresses {
1485 addresses.insert(Address::from_str(&address).map_err(|_| {
1486 GrpcError::InvalidArgument(format!("invalid address: {}", address))
1487 })?);
1488 }
1489 }
1490 }
1491 }
1492 }
1493
1494 if operation_ids_filter.is_none() && addresses_filter.is_none() {
1495 return Err(GrpcError::InvalidArgument("no filter provided".to_string()));
1496 }
1497
1498 let mut ops_ids: Option<PreHashSet<OperationId>> = None;
1499
1500 if let Some(mut o_ids) = operation_ids_filter {
1502 let read_lock = grpc.storage.read_operations();
1503 o_ids.retain(|id: &OperationId| read_lock.contains(id));
1504 ops_ids = Some(o_ids);
1505 }
1506
1507 if let Some(addrs) = addresses_filter {
1509 let o_ids: PreHashSet<OperationId> = {
1510 let read_lock = grpc.storage.read_operations();
1511 let mut o_ids: PreHashSet<OperationId> = PreHashSet::default();
1512 for addr in addrs {
1513 if let Some(addr_o_ids) = read_lock.get_operations_created_by(&addr) {
1514 o_ids.extend(addr_o_ids.clone());
1515 }
1516 }
1517
1518 o_ids
1519 };
1520 if let Some(operation_ids) = ops_ids.as_mut() {
1521 operation_ids.retain(|id: &OperationId| o_ids.contains(id));
1522 } else {
1523 ops_ids = Some(o_ids)
1524 }
1525 }
1526
1527 let operations: Vec<grpc_model::OperationInfo> = if let Some(operation_ids) = ops_ids {
1528 let secure_share_operations: Vec<SecureShareOperation> = {
1529 let read_ops = grpc.storage.read_operations();
1530 operation_ids
1531 .iter()
1532 .filter_map(|id| read_ops.get(id).cloned())
1533 .collect()
1534 };
1535
1536 let storage_info: Vec<(SecureShareOperation, PreHashSet<BlockId>)> = {
1537 let read_blocks = grpc.storage.read_blocks();
1538 secure_share_operations
1539 .into_iter()
1540 .map(|secure_share_operation| {
1541 let op_id = secure_share_operation.id;
1542 (
1543 secure_share_operation,
1544 read_blocks
1545 .get_blocks_by_operation(&op_id)
1546 .cloned()
1547 .unwrap_or_default(),
1548 )
1549 })
1550 .collect()
1551 };
1552
1553 storage_info
1554 .into_iter()
1555 .map(|secureshare| {
1556 let (secureshare_operation, block_ids) = secureshare;
1557 grpc_model::OperationInfo {
1558 id: secureshare_operation.id.to_string(),
1559 thread: secureshare_operation
1560 .content_creator_address
1561 .get_thread(grpc.grpc_config.thread_count)
1562 as u32,
1563 block_ids: block_ids.into_iter().map(|id| id.to_string()).collect(),
1564 }
1565 })
1566 .collect()
1567 } else {
1568 Vec::new()
1569 };
1570
1571 Ok(grpc_api::SearchOperationsResponse {
1572 operation_infos: operations,
1573 })
1574}