massa_grpc/stream/
trait_filters_impl.rs

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
18/// Trait implementation for filtering the output based on the request
19pub(crate) trait FilterGrpc<RequestType, FilterType, Data> {
20    /// Build the filter from the request
21    fn build_from_request(
22        request: RequestType,
23        grpc_config: &GrpcConfig,
24    ) -> Result<FilterType, GrpcError>;
25    /// Filter the output based on the filter
26    fn filter_output(&self, content: Data, grpc_config: &GrpcConfig) -> Option<Data>;
27}
28
29/// Type declaration for NewSlotExecutionOutputsFilter
30#[derive(Clone, Debug, Default)]
31pub(crate) struct FilterNewSlotExec {
32    // Execution output status to filter
33    status_filter: Option<i32>,
34    // Slot range to filter
35    slot_ranges_filter: Option<Vec<grpc_model::SlotRange>>,
36    // Async pool changes filter
37    async_pool_changes_filter: Option<Vec<grpc_api::async_pool_changes_filter::Filter>>,
38    // Executed denounciation filter
39    executed_denounciation_filter: Option<grpc_api::executed_denounciation_filter::Filter>,
40    // Execution event filter
41    execution_event_filter: Option<Vec<grpc_api::execution_event_filter::Filter>>,
42    // Executed ops changes filter
43    executed_ops_changes_filter: Option<Vec<grpc_api::executed_ops_changes_filter::Filter>>,
44    // Ledger changes filter
45    ledger_changes_filter: Option<Vec<grpc_api::ledger_changes_filter::Filter>>,
46}
47
48// Type declaration for NewOperationsFilter
49#[derive(Debug)]
50pub(crate) struct FilterNewOperations {
51    // Operation ids to filter
52    operation_ids: Option<HashSet<OperationId>>,
53    // Addresses to filter
54    addresses: Option<HashSet<Address>>,
55    // Operation types to filter
56    operation_types: Option<HashSet<i32>>,
57}
58
59// Type declaration for NewBlocksFilter
60#[derive(Clone, Debug)]
61pub(crate) struct FilterNewBlocks {
62    // Block ids to filter
63    block_ids: Option<HashSet<BlockId>>,
64    // Addresses to filter
65    addresses: Option<HashSet<Address>>,
66    // Slot range to filter
67    slot_ranges: Option<HashSet<SlotRange>>,
68}
69
70// Type declaration for NewFilledBlocksFilter
71#[derive(Clone, Debug)]
72pub(crate) struct FilterNewFilledBlocks {
73    // Block ids to filter
74    block_ids: Option<HashSet<BlockId>>,
75    // Addresses to filter
76    addresses: Option<HashSet<Address>>,
77    // Slot range to filter
78    slot_ranges: Option<HashSet<SlotRange>>,
79}
80
81// Type declaration for NewEndorsementsFilter
82#[derive(Debug)]
83pub(crate) struct NewEndorsementsFilter {
84    // Endorsement ids to filter
85    endorsement_ids: Option<HashSet<EndorsementId>>,
86    // Addresses to filter
87    addresses: Option<HashSet<Address>>,
88    // Block ids to filter
89    block_ids: Option<HashSet<BlockId>>,
90}
91
92// Filter for execution-info streams (only compiled with feature execution-info)
93#[cfg(feature = "execution-info")]
94pub(crate) struct NewExecutionInfoFilter {
95    // Address to filter
96    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
175// Return if the execution outputs should be send and remove the fields that are not needed
176fn 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    // Filter on status
183    if let Some(status) = filters.status_filter {
184        if status.ne(&exec_status) {
185            return None;
186        }
187    }
188
189    // Filter Slot Range
190    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    // Filter Exec event
204    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    // Filter async pool changes
234    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    // Filter executed ops id
330    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    // Filter ledger changes
349    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        // Get params filter from the request.
390        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); // inclusive
475            let mut end_slot = Slot::new(u64::MAX, grpc_config.thread_count - 1); // exclusive
476
477            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        // Get params filter from the request.
517        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                        // The length limited to the number of operation types in the enum
552                        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        // Get params filter from the request.
619        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); // inclusive
699            let mut end_slot = Slot::new(u64::MAX, grpc_config.thread_count - 1); // exclusive
700
701            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        // Get params filter from the request.
741        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}