1use std::collections::HashSet;
2use std::str::FromStr;
3
4use crate::SlotRange;
5use crate::{config::GrpcConfig, error::GrpcError};
6#[cfg(feature = "execution-info")]
7use massa_execution_exports::execution_info::ExecutionInfoForSlot;
8use massa_execution_exports::{ExecutionOutput, SlotExecutionOutput};
9use massa_models::address::Address;
10use massa_models::block::{FilledBlock, SecureShareBlock};
11use massa_models::block_id::BlockId;
12use massa_models::endorsement::{EndorsementId, SecureShareEndorsement};
13use massa_models::operation::{OperationId, SecureShareOperation};
14use massa_models::slot::Slot;
15use massa_proto_rs::massa::api::v1::{self as grpc_api};
16use massa_proto_rs::massa::model::v1::{self as grpc_model};
17
18pub(crate) trait FilterGrpc<RequestType, FilterType, Data> {
20 fn build_from_request(
22 request: RequestType,
23 grpc_config: &GrpcConfig,
24 ) -> Result<FilterType, GrpcError>;
25 fn filter_output(&self, content: Data, grpc_config: &GrpcConfig) -> Option<Data>;
27}
28
29#[derive(Clone, Debug, Default)]
31pub(crate) struct FilterNewSlotExec {
32 status_filter: Option<i32>,
34 slot_ranges_filter: Option<Vec<grpc_model::SlotRange>>,
36 async_pool_changes_filter: Option<Vec<grpc_api::async_pool_changes_filter::Filter>>,
38 executed_denounciation_filter: Option<grpc_api::executed_denounciation_filter::Filter>,
40 execution_event_filter: Option<Vec<grpc_api::execution_event_filter::Filter>>,
42 executed_ops_changes_filter: Option<Vec<grpc_api::executed_ops_changes_filter::Filter>>,
44 ledger_changes_filter: Option<Vec<grpc_api::ledger_changes_filter::Filter>>,
46}
47
48#[derive(Debug)]
50pub(crate) struct FilterNewOperations {
51 operation_ids: Option<HashSet<OperationId>>,
53 addresses: Option<HashSet<Address>>,
55 operation_types: Option<HashSet<i32>>,
57}
58
59#[derive(Clone, Debug)]
61pub(crate) struct FilterNewBlocks {
62 block_ids: Option<HashSet<BlockId>>,
64 addresses: Option<HashSet<Address>>,
66 slot_ranges: Option<HashSet<SlotRange>>,
68}
69
70#[derive(Clone, Debug)]
72pub(crate) struct FilterNewFilledBlocks {
73 block_ids: Option<HashSet<BlockId>>,
75 addresses: Option<HashSet<Address>>,
77 slot_ranges: Option<HashSet<SlotRange>>,
79}
80
81#[derive(Debug)]
83pub(crate) struct NewEndorsementsFilter {
84 endorsement_ids: Option<HashSet<EndorsementId>>,
86 addresses: Option<HashSet<Address>>,
88 block_ids: Option<HashSet<BlockId>>,
90}
91
92#[cfg(feature = "execution-info")]
94pub(crate) struct NewExecutionInfoFilter {
95 address: Option<Address>,
97}
98
99impl
100 FilterGrpc<Vec<grpc_api::NewSlotExecutionOutputsFilter>, FilterNewSlotExec, SlotExecutionOutput>
101 for FilterNewSlotExec
102{
103 fn build_from_request(
104 filters: Vec<grpc_api::NewSlotExecutionOutputsFilter>,
105 grpc_config: &GrpcConfig,
106 ) -> Result<FilterNewSlotExec, GrpcError> {
107 if filters.len() as u32 > grpc_config.max_filters_per_request {
108 return Err(GrpcError::InvalidArgument(format!(
109 "too many filters received. Only a maximum of {} filters are accepted per request",
110 grpc_config.max_filters_per_request
111 )));
112 }
113
114 let mut result = FilterNewSlotExec::default();
115
116 for query in filters.into_iter() {
117 if let Some(filter) = query.filter {
118 match filter {
119 grpc_api::new_slot_execution_outputs_filter::Filter::Status(status) => result.status_filter = Some(status),
120 grpc_api::new_slot_execution_outputs_filter::Filter::SlotRange(s_range) => {
121 result.slot_ranges_filter.get_or_insert(Vec::new()).push(s_range);
122 },
123 grpc_api::new_slot_execution_outputs_filter::Filter::AsyncPoolChangesFilter(filter) => {
124 if let Some(request_f) = filter.filter {
125 result.async_pool_changes_filter.get_or_insert(Vec::new()).push(request_f);
126 }
127 },
128 grpc_api::new_slot_execution_outputs_filter::Filter::ExecutedDenounciationFilter(filter) => result.executed_denounciation_filter = filter.filter,
129 grpc_api::new_slot_execution_outputs_filter::Filter::EventFilter(filter) => {
130 if let Some(request_f) = filter.filter {
131 result.execution_event_filter.get_or_insert(Vec::new()).push(request_f);
132 }
133 },
134 grpc_api::new_slot_execution_outputs_filter::Filter::ExecutedOpsChangesFilter(filter) => {
135 if let Some(request_f) = filter.filter {
136 result.executed_ops_changes_filter.get_or_insert(Vec::new()).push(request_f);
137 }
138 },
139 grpc_api::new_slot_execution_outputs_filter::Filter::LedgerChangesFilter(filter) => {
140 if let Some(request_f) = filter.filter {
141 result.ledger_changes_filter.get_or_insert(Vec::new()).push(request_f);
142 }
143 },
144 }
145 }
146 }
147
148 Ok(result)
149 }
150
151 fn filter_output(
152 &self,
153 content: SlotExecutionOutput,
154 grpc_config: &GrpcConfig,
155 ) -> Option<SlotExecutionOutput> {
156 match content {
157 SlotExecutionOutput::ExecutedSlot(e_output) => filter_map_exec_output_inner(
158 e_output,
159 self,
160 grpc_config,
161 grpc_model::ExecutionOutputStatus::Candidate as i32,
162 )
163 .map(SlotExecutionOutput::ExecutedSlot),
164 SlotExecutionOutput::FinalizedSlot(e_output) => filter_map_exec_output_inner(
165 e_output,
166 self,
167 grpc_config,
168 grpc_model::ExecutionOutputStatus::Final as i32,
169 )
170 .map(SlotExecutionOutput::FinalizedSlot),
171 }
172 }
173}
174
175fn filter_map_exec_output_inner(
177 mut exec_output: ExecutionOutput,
178 filters: &FilterNewSlotExec,
179 _grpc_config: &GrpcConfig,
180 exec_status: i32,
181) -> Option<ExecutionOutput> {
182 if let Some(status) = filters.status_filter {
184 if status.ne(&exec_status) {
185 return None;
186 }
187 }
188
189 if let Some(slot_ranges) = &filters.slot_ranges_filter {
191 if slot_ranges.iter().any(|slot_range| {
192 slot_range
193 .start_slot
194 .is_some_and(|start| exec_output.slot < start.into())
195 || slot_range
196 .end_slot
197 .is_some_and(|end| exec_output.slot >= end.into())
198 }) {
199 return None;
200 }
201 }
202
203 if let Some(execution_event_filter) = filters.execution_event_filter.as_ref() {
205 exec_output.events.0.retain(|event| {
206 execution_event_filter.iter().all(|filter| match filter {
207 grpc_api::execution_event_filter::Filter::None(_) => false,
208 grpc_api::execution_event_filter::Filter::CallerAddress(addr) => event
209 .context
210 .call_stack
211 .front()
212 .is_some_and(|call| call.to_string().eq(addr)),
213 grpc_api::execution_event_filter::Filter::EmitterAddress(addr) => event
214 .context
215 .call_stack
216 .back()
217 .is_some_and(|emit| emit.to_string().eq(addr)),
218 grpc_api::execution_event_filter::Filter::OriginalOperationId(ope_id) => event
219 .context
220 .origin_operation_id
221 .is_some_and(|ope| ope.to_string().eq(ope_id)),
222 grpc_api::execution_event_filter::Filter::IsFailure(b) => {
223 event.context.is_error.eq(b)
224 }
225 })
226 });
227
228 if exec_output.events.0.is_empty() {
229 return None;
230 }
231 }
232
233 if let Some(async_pool_changes_filter) = &filters.async_pool_changes_filter {
235 exec_output.state_changes.async_pool_changes.0.retain(
236 |(_msg_id, _slot, _emission_index), changes| {
237 async_pool_changes_filter.iter().all(|filter| match filter {
238 grpc_api::async_pool_changes_filter::Filter::None(_empty) => false,
239 grpc_api::async_pool_changes_filter::Filter::Type(filter_type) => match changes
240 {
241 massa_models::types::SetUpdateOrDelete::Set(_) => {
242 (grpc_model::AsyncPoolChangeType::Set as i32).eq(filter_type)
243 }
244 massa_models::types::SetUpdateOrDelete::Update(_) => {
245 (grpc_model::AsyncPoolChangeType::Update as i32).eq(filter_type)
246 }
247 massa_models::types::SetUpdateOrDelete::Delete => {
248 (grpc_model::AsyncPoolChangeType::Delete as i32).eq(filter_type)
249 }
250 },
251 grpc_api::async_pool_changes_filter::Filter::Handler(handler) => {
252 match changes {
253 massa_models::types::SetUpdateOrDelete::Set(msg) => {
254 msg.function.eq(handler)
255 }
256 massa_models::types::SetUpdateOrDelete::Update(msg) => {
257 match &msg.function {
258 massa_models::types::SetOrKeep::Set(func) => func.eq(handler),
259 massa_models::types::SetOrKeep::Keep => false,
260 }
261 }
262 massa_models::types::SetUpdateOrDelete::Delete => false,
263 }
264 }
265 grpc_api::async_pool_changes_filter::Filter::DestinationAddress(
266 filter_dest_addr,
267 ) => match changes {
268 massa_models::types::SetUpdateOrDelete::Set(msg) => {
269 msg.destination.to_string().eq(filter_dest_addr)
270 }
271 massa_models::types::SetUpdateOrDelete::Update(msg) => {
272 match msg.destination {
273 massa_models::types::SetOrKeep::Set(dest) => {
274 dest.to_string().eq(filter_dest_addr)
275 }
276 massa_models::types::SetOrKeep::Keep => false,
277 }
278 }
279 massa_models::types::SetUpdateOrDelete::Delete => false,
280 },
281 grpc_api::async_pool_changes_filter::Filter::EmitterAddress(
282 filter_emit_addr,
283 ) => match changes {
284 massa_models::types::SetUpdateOrDelete::Set(msg) => {
285 msg.sender.to_string().eq(filter_emit_addr)
286 }
287 massa_models::types::SetUpdateOrDelete::Update(msg) => match msg.sender {
288 massa_models::types::SetOrKeep::Set(addr) => {
289 addr.to_string().eq(filter_emit_addr)
290 }
291 massa_models::types::SetOrKeep::Keep => false,
292 },
293 massa_models::types::SetUpdateOrDelete::Delete => false,
294 },
295 grpc_api::async_pool_changes_filter::Filter::CanBeExecuted(filter_exec) => {
296 match changes {
297 massa_models::types::SetUpdateOrDelete::Set(msg) => {
298 msg.can_be_executed.eq(filter_exec)
299 }
300 massa_models::types::SetUpdateOrDelete::Update(msg) => {
301 match msg.can_be_executed {
302 massa_models::types::SetOrKeep::Set(b) => b.eq(filter_exec),
303 massa_models::types::SetOrKeep::Keep => false,
304 }
305 }
306 massa_models::types::SetUpdateOrDelete::Delete => false,
307 }
308 }
309 })
310 },
311 );
312
313 if exec_output.state_changes.async_pool_changes.0.is_empty() {
314 return None;
315 }
316 }
317
318 if let Some(executed_denounciation_filter) = &filters.executed_denounciation_filter {
319 match executed_denounciation_filter {
320 grpc_api::executed_denounciation_filter::Filter::None(_empty) => {
321 exec_output
322 .state_changes
323 .executed_denunciations_changes
324 .clear();
325 }
326 }
327 }
328
329 if let Some(executed_ops_changes_filter) = &filters.executed_ops_changes_filter {
331 exec_output
332 .state_changes
333 .executed_ops_changes
334 .retain(|op, (_success, _slot)| {
335 executed_ops_changes_filter.iter().all(|f| match f {
336 grpc_api::executed_ops_changes_filter::Filter::None(_empty) => false,
337 grpc_api::executed_ops_changes_filter::Filter::OperationId(filter_ope_id) => {
338 op.to_string().eq(filter_ope_id)
339 }
340 })
341 });
342
343 if exec_output.state_changes.executed_ops_changes.is_empty() {
344 return None;
345 }
346 }
347
348 if let Some(ledger_changes_filter) = &filters.ledger_changes_filter {
350 exec_output
351 .state_changes
352 .ledger_changes
353 .0
354 .retain(|addr, _ledger_change| {
355 ledger_changes_filter.iter().all(|filter| match filter {
356 grpc_api::ledger_changes_filter::Filter::None(_empty) => false,
357 grpc_api::ledger_changes_filter::Filter::Address(filter_addr) => {
358 addr.to_string().eq(filter_addr)
359 }
360 })
361 });
362
363 if exec_output.state_changes.ledger_changes.0.is_empty() {
364 return None;
365 }
366 }
367
368 Some(exec_output)
369}
370
371impl FilterGrpc<Vec<grpc_api::NewBlocksFilter>, FilterNewBlocks, SecureShareBlock>
372 for FilterNewBlocks
373{
374 fn build_from_request(
375 filters: Vec<grpc_api::NewBlocksFilter>,
376 grpc_config: &GrpcConfig,
377 ) -> Result<FilterNewBlocks, GrpcError> {
378 if filters.len() as u32 > grpc_config.max_filters_per_request {
379 return Err(GrpcError::InvalidArgument(format!(
380 "too many filters received. Only a maximum of {} filters are accepted per request",
381 grpc_config.max_filters_per_request
382 )));
383 }
384
385 let mut block_ids_filter: Option<HashSet<BlockId>> = None;
386 let mut addresses_filter: Option<HashSet<Address>> = None;
387 let mut slot_ranges_filter: Option<HashSet<SlotRange>> = None;
388
389 for query in filters.into_iter() {
391 if let Some(filter) = query.filter {
392 match filter {
393 grpc_api::new_blocks_filter::Filter::BlockIds(ids) => {
394 if ids.block_ids.len() as u32 > grpc_config.max_block_ids_per_request {
395 return Err(GrpcError::InvalidArgument(format!(
396 "too many block ids received. Only a maximum of {} block ids are accepted per request",
397 grpc_config.max_block_ids_per_request
398 )));
399 }
400
401 let block_ids = block_ids_filter.get_or_insert_with(HashSet::new);
402 for block_id in ids.block_ids {
403 block_ids.insert(BlockId::from_str(&block_id).map_err(|_| {
404 GrpcError::InvalidArgument(format!(
405 "invalid block id: {}",
406 block_id
407 ))
408 })?);
409 }
410 }
411 grpc_api::new_blocks_filter::Filter::Addresses(addrs) => {
412 if addrs.addresses.len() as u32 > grpc_config.max_addresses_per_request {
413 return Err(GrpcError::InvalidArgument(format!(
414 "too many addresses received. Only a maximum of {} addresses are accepted per request",
415 grpc_config.max_addresses_per_request
416 )));
417 }
418
419 let addresses = addresses_filter.get_or_insert_with(HashSet::new);
420 for address in addrs.addresses {
421 addresses.insert(Address::from_str(&address).map_err(|_| {
422 GrpcError::InvalidArgument(format!("invalid address: {}", address))
423 })?);
424 }
425 }
426 grpc_api::new_blocks_filter::Filter::SlotRange(s_range) => {
427 let slot_ranges = slot_ranges_filter.get_or_insert_with(HashSet::new);
428 if slot_ranges.len() as u32 >= grpc_config.max_slot_ranges_per_request {
429 return Err(GrpcError::InvalidArgument(format!(
430 "too many slot ranges received. Only a maximum of {} slot ranges are accepted per request",
431 grpc_config.max_slot_ranges_per_request
432 )));
433 }
434
435 let start_slot = s_range.start_slot.map(|s| s.into());
436 let end_slot = s_range.end_slot.map(|s| s.into());
437
438 let slot_range = SlotRange {
439 start_slot,
440 end_slot,
441 };
442 slot_range.check()?;
443 slot_ranges.insert(slot_range);
444 }
445 }
446 }
447 }
448
449 Ok(FilterNewBlocks {
450 block_ids: block_ids_filter,
451 addresses: addresses_filter,
452 slot_ranges: slot_ranges_filter,
453 })
454 }
455
456 fn filter_output(
457 &self,
458 content: SecureShareBlock,
459 grpc_config: &GrpcConfig,
460 ) -> Option<SecureShareBlock> {
461 if let Some(block_ids) = &self.block_ids {
462 if !block_ids.contains(&content.id) {
463 return None;
464 }
465 }
466
467 if let Some(addresses) = &self.addresses {
468 if !addresses.contains(&content.content_creator_address) {
469 return None;
470 }
471 }
472
473 if let Some(slot_ranges) = &self.slot_ranges {
474 let mut start_slot = Slot::new(0, 0); let mut end_slot = Slot::new(u64::MAX, grpc_config.thread_count - 1); for slot_range in slot_ranges {
478 start_slot =
479 start_slot.max(slot_range.start_slot.unwrap_or_else(|| Slot::new(0, 0)));
480 end_slot = end_slot.min(
481 slot_range
482 .end_slot
483 .unwrap_or_else(|| Slot::new(u64::MAX, grpc_config.thread_count - 1)),
484 );
485 }
486 end_slot = end_slot.max(start_slot);
487 let current_slot = content.content.header.content.slot;
488
489 if !(current_slot >= start_slot && current_slot < end_slot) {
490 return None;
491 }
492 }
493
494 Some(content)
495 }
496}
497
498impl FilterGrpc<Vec<grpc_api::NewOperationsFilter>, FilterNewOperations, SecureShareOperation>
499 for FilterNewOperations
500{
501 fn build_from_request(
502 filters: Vec<grpc_api::NewOperationsFilter>,
503 grpc_config: &GrpcConfig,
504 ) -> Result<FilterNewOperations, GrpcError> {
505 if filters.len() as u32 > grpc_config.max_filters_per_request {
506 return Err(GrpcError::InvalidArgument(format!(
507 "too many filters received. Only a maximum of {} filters are accepted per request",
508 grpc_config.max_filters_per_request
509 )));
510 }
511
512 let mut operation_ids_filter: Option<HashSet<OperationId>> = None;
513 let mut addresses_filter: Option<HashSet<Address>> = None;
514 let mut operation_types_filter: Option<HashSet<i32>> = None;
515
516 for query in filters.into_iter() {
518 if let Some(filter) = query.filter {
519 match filter {
520 grpc_api::new_operations_filter::Filter::OperationIds(ids) => {
521 if ids.operation_ids.len() as u32
522 > grpc_config.max_operation_ids_per_request
523 {
524 return Err(GrpcError::InvalidArgument(format!(
525 "too many operation ids received. Only a maximum of {} operation ids are accepted per request",
526 grpc_config.max_operation_ids_per_request
527 )));
528 }
529 let operation_ids = operation_ids_filter.get_or_insert_with(HashSet::new);
530 for id in ids.operation_ids {
531 operation_ids.insert(OperationId::from_str(&id).map_err(|_| {
532 GrpcError::InvalidArgument(format!("invalid operation id: {}", id))
533 })?);
534 }
535 }
536 grpc_api::new_operations_filter::Filter::Addresses(addrs) => {
537 if addrs.addresses.len() as u32 > grpc_config.max_addresses_per_request {
538 return Err(GrpcError::InvalidArgument(format!(
539 "too many addresses received. Only a maximum of {} addresses are accepted per request",
540 grpc_config.max_addresses_per_request
541 )));
542 }
543 let addresses = addresses_filter.get_or_insert_with(HashSet::new);
544 for address in addrs.addresses {
545 addresses.insert(Address::from_str(&address).map_err(|_| {
546 GrpcError::InvalidArgument(format!("invalid address: {}", address))
547 })?);
548 }
549 }
550 grpc_api::new_operations_filter::Filter::OperationTypes(ope_types) => {
551 if ope_types.op_types.len() as u64 > 6 {
553 return Err(GrpcError::InvalidArgument(
554 "too many operation types received. Only a maximum of 6 operation types are accepted per request".to_string()
555 ));
556 }
557 let operation_types =
558 operation_types_filter.get_or_insert_with(HashSet::new);
559 operation_types.extend(&ope_types.op_types);
560 }
561 }
562 }
563 }
564
565 Ok(FilterNewOperations {
566 operation_ids: operation_ids_filter,
567 addresses: addresses_filter,
568 operation_types: operation_types_filter,
569 })
570 }
571
572 fn filter_output(
573 &self,
574 content: SecureShareOperation,
575 _grpc_config: &GrpcConfig,
576 ) -> Option<SecureShareOperation> {
577 if let Some(operation_ids) = &self.operation_ids {
578 if !operation_ids.contains(&content.id) {
579 return None;
580 }
581 }
582
583 if let Some(addresses) = &self.addresses {
584 if !addresses.contains(&content.content_creator_address) {
585 return None;
586 }
587 }
588
589 if let Some(operation_types) = &self.operation_types {
590 let op_type = grpc_model::OpType::from(content.content.op.clone()) as i32;
591 if !operation_types.contains(&op_type) {
592 return None;
593 }
594 }
595
596 Some(content)
597 }
598}
599
600impl FilterGrpc<Vec<grpc_api::NewBlocksFilter>, FilterNewFilledBlocks, FilledBlock>
601 for FilterNewFilledBlocks
602{
603 fn build_from_request(
604 filters: Vec<grpc_api::NewBlocksFilter>,
605 grpc_config: &GrpcConfig,
606 ) -> Result<FilterNewFilledBlocks, GrpcError> {
607 if filters.len() as u32 > grpc_config.max_filters_per_request {
608 return Err(GrpcError::InvalidArgument(format!(
609 "too many filters received. Only a maximum of {} filters are accepted per request",
610 grpc_config.max_filters_per_request
611 )));
612 }
613
614 let mut block_ids_filter: Option<HashSet<BlockId>> = None;
615 let mut addresses_filter: Option<HashSet<Address>> = None;
616 let mut slot_ranges_filter: Option<HashSet<SlotRange>> = None;
617
618 for query in filters.into_iter() {
620 if let Some(filter) = query.filter {
621 match filter {
622 grpc_api::new_blocks_filter::Filter::BlockIds(ids) => {
623 if ids.block_ids.len() as u32 > grpc_config.max_block_ids_per_request {
624 return Err(GrpcError::InvalidArgument(format!(
625 "too many block ids received. Only a maximum of {} block ids are accepted per request",
626 grpc_config.max_block_ids_per_request
627 )));
628 }
629 let block_ids = block_ids_filter.get_or_insert_with(HashSet::new);
630 for block_id in ids.block_ids {
631 block_ids.insert(BlockId::from_str(&block_id).map_err(|_| {
632 GrpcError::InvalidArgument(format!(
633 "invalid block id: {}",
634 block_id
635 ))
636 })?);
637 }
638 }
639 grpc_api::new_blocks_filter::Filter::Addresses(addrs) => {
640 if addrs.addresses.len() as u32 > grpc_config.max_addresses_per_request {
641 return Err(GrpcError::InvalidArgument(format!(
642 "too many addresses received. Only a maximum of {} addresses are accepted per request",
643 grpc_config.max_addresses_per_request
644 )));
645 }
646
647 let addresses = addresses_filter.get_or_insert_with(HashSet::new);
648 for address in addrs.addresses {
649 addresses.insert(Address::from_str(&address).map_err(|_| {
650 GrpcError::InvalidArgument(format!("invalid address: {}", address))
651 })?);
652 }
653 }
654 grpc_api::new_blocks_filter::Filter::SlotRange(s_range) => {
655 let slot_ranges = slot_ranges_filter.get_or_insert_with(HashSet::new);
656 if slot_ranges.len() as u32 >= grpc_config.max_slot_ranges_per_request {
657 return Err(GrpcError::InvalidArgument(format!(
658 "too many slot ranges received. Only a maximum of {} slot ranges are accepted per request",
659 grpc_config.max_slot_ranges_per_request
660 )));
661 }
662
663 let start_slot = s_range.start_slot.map(|s| s.into());
664 let end_slot = s_range.end_slot.map(|s| s.into());
665
666 let slot_range = SlotRange {
667 start_slot,
668 end_slot,
669 };
670 slot_range.check()?;
671 slot_ranges.insert(slot_range);
672 }
673 }
674 }
675 }
676
677 Ok(FilterNewFilledBlocks {
678 block_ids: block_ids_filter,
679 addresses: addresses_filter,
680 slot_ranges: slot_ranges_filter,
681 })
682 }
683
684 fn filter_output(&self, content: FilledBlock, grpc_config: &GrpcConfig) -> Option<FilledBlock> {
685 if let Some(block_ids) = &self.block_ids {
686 if !block_ids.contains(&content.header.id) {
687 return None;
688 }
689 }
690
691 if let Some(addresses) = &self.addresses {
692 if !addresses.contains(&content.header.content_creator_address) {
693 return None;
694 }
695 }
696
697 if let Some(slot_ranges) = &self.slot_ranges {
698 let mut start_slot = Slot::new(0, 0); let mut end_slot = Slot::new(u64::MAX, grpc_config.thread_count - 1); for slot_range in slot_ranges {
702 start_slot =
703 start_slot.max(slot_range.start_slot.unwrap_or_else(|| Slot::new(0, 0)));
704 end_slot = end_slot.min(
705 slot_range
706 .end_slot
707 .unwrap_or_else(|| Slot::new(u64::MAX, grpc_config.thread_count - 1)),
708 );
709 }
710 end_slot = end_slot.max(start_slot);
711 let current_slot = content.header.content.slot;
712
713 if !(current_slot >= start_slot && current_slot < end_slot) {
714 return None;
715 }
716 }
717
718 Some(content)
719 }
720}
721
722impl FilterGrpc<Vec<grpc_api::NewEndorsementsFilter>, NewEndorsementsFilter, SecureShareEndorsement>
723 for NewEndorsementsFilter
724{
725 fn build_from_request(
726 filters: Vec<grpc_api::NewEndorsementsFilter>,
727 grpc_config: &GrpcConfig,
728 ) -> Result<NewEndorsementsFilter, GrpcError> {
729 if filters.len() as u32 > grpc_config.max_filters_per_request {
730 return Err(GrpcError::InvalidArgument(format!(
731 "too many filters received. Only a maximum of {} filters are accepted per request",
732 grpc_config.max_filters_per_request
733 )));
734 }
735
736 let mut endorsement_ids_filter: Option<HashSet<EndorsementId>> = None;
737 let mut addresses_filter: Option<HashSet<Address>> = None;
738 let mut block_ids_filter: Option<HashSet<BlockId>> = None;
739
740 for query in filters.into_iter() {
742 if let Some(filter) = query.filter {
743 match filter {
744 grpc_api::new_endorsements_filter::Filter::EndorsementIds(ids) => {
745 if ids.endorsement_ids.len() as u32
746 > grpc_config.max_endorsement_ids_per_request
747 {
748 return Err(GrpcError::InvalidArgument(format!(
749 "too many endorsement ids received. Only a maximum of {} endorsement ids are accepted per request",
750 grpc_config.max_endorsement_ids_per_request
751 )));
752 }
753 let endorsement_ids =
754 endorsement_ids_filter.get_or_insert_with(HashSet::new);
755 for id in ids.endorsement_ids {
756 endorsement_ids.insert(EndorsementId::from_str(&id).map_err(|_| {
757 GrpcError::InvalidArgument(format!(
758 "invalid endorsement id: {}",
759 id
760 ))
761 })?);
762 }
763 }
764 grpc_api::new_endorsements_filter::Filter::Addresses(addrs) => {
765 if addrs.addresses.len() as u32 > grpc_config.max_addresses_per_request {
766 return Err(GrpcError::InvalidArgument(format!(
767 "too many addresses received. Only a maximum of {} addresses are accepted per request",
768 grpc_config.max_addresses_per_request
769 )));
770 }
771 let addresses = addresses_filter.get_or_insert_with(HashSet::new);
772 for address in addrs.addresses {
773 addresses.insert(Address::from_str(&address).map_err(|_| {
774 GrpcError::InvalidArgument(format!("invalid address: {}", address))
775 })?);
776 }
777 }
778 grpc_api::new_endorsements_filter::Filter::BlockIds(ids) => {
779 if ids.block_ids.len() as u32 > grpc_config.max_block_ids_per_request {
780 return Err(GrpcError::InvalidArgument(format!(
781 "too many block ids received. Only a maximum of {} block ids are accepted per request",
782 grpc_config.max_block_ids_per_request
783 )));
784 }
785 let block_ids = block_ids_filter.get_or_insert_with(HashSet::new);
786 for block_id in ids.block_ids {
787 block_ids.insert(BlockId::from_str(&block_id).map_err(|_| {
788 GrpcError::InvalidArgument(format!(
789 "invalid block id: {}",
790 block_id
791 ))
792 })?);
793 }
794 }
795 }
796 }
797 }
798
799 Ok(NewEndorsementsFilter {
800 endorsement_ids: endorsement_ids_filter,
801 addresses: addresses_filter,
802 block_ids: block_ids_filter,
803 })
804 }
805
806 fn filter_output(
807 &self,
808 content: SecureShareEndorsement,
809 _grpc_config: &GrpcConfig,
810 ) -> Option<SecureShareEndorsement> {
811 if let Some(endorsement_ids) = &self.endorsement_ids {
812 if !endorsement_ids.contains(&content.id) {
813 return None;
814 }
815 }
816
817 if let Some(addresses) = &self.addresses {
818 if !addresses.contains(&content.content_creator_address) {
819 return None;
820 }
821 }
822
823 if let Some(block_ids) = &self.block_ids {
824 if !block_ids.contains(&content.content.endorsed_block) {
825 return None;
826 }
827 }
828
829 Some(content)
830 }
831}
832
833#[cfg(feature = "execution-info")]
834impl FilterGrpc<Option<String>, NewExecutionInfoFilter, ExecutionInfoForSlot>
835 for NewExecutionInfoFilter
836{
837 fn build_from_request(
838 filter: Option<String>,
839 _grpc_config: &GrpcConfig,
840 ) -> Result<NewExecutionInfoFilter, GrpcError> {
841 let address = filter
842 .map(|addr| {
843 Address::from_str(&addr)
844 .map_err(|_| GrpcError::InvalidArgument(format!("invalid address: {}", addr)))
845 })
846 .transpose()?;
847 Ok(NewExecutionInfoFilter { address })
848 }
849
850 fn filter_output(
851 &self,
852 mut content: ExecutionInfoForSlot,
853 _grpc_config: &GrpcConfig,
854 ) -> Option<ExecutionInfoForSlot> {
855 if let Some(address_filter) = &self.address {
856 content.transfers.retain(|transfer| {
857 transfer.from.is_some_and(|from| from.eq(address_filter))
858 || transfer.to.is_some_and(|to| to.eq(address_filter))
859 });
860 }
861
862 if content.transfers.is_empty() {
863 None
864 } else {
865 Some(content)
866 }
867 }
868}