Skip to main content

drive_abci/query/
service.rs

1use crate::error::query::QueryError;
2use crate::error::Error;
3use crate::metrics::{abci_response_code_metric_label, query_duration_metric};
4use crate::platform_types::platform::Platform;
5use crate::platform_types::platform_state::PlatformState;
6use crate::platform_types::platform_state::PlatformStateV0Methods;
7use crate::query::QueryValidationResult;
8use crate::rpc::core::DefaultCoreRPC;
9use crate::utils::spawn_blocking_task_with_name_if_supported;
10use async_trait::async_trait;
11use dapi_grpc::drive::v0::drive_internal_server::DriveInternal;
12use dapi_grpc::drive::v0::{GetProofsRequest, GetProofsResponse};
13use dapi_grpc::platform::v0::get_path_elements_request;
14use dapi_grpc::platform::v0::platform_server::Platform as PlatformService;
15use dapi_grpc::platform::v0::{
16    BroadcastStateTransitionRequest, BroadcastStateTransitionResponse, GetAddressInfoRequest,
17    GetAddressInfoResponse, GetAddressesBranchStateRequest, GetAddressesBranchStateResponse,
18    GetAddressesInfosRequest, GetAddressesInfosResponse, GetAddressesTrunkStateRequest,
19    GetAddressesTrunkStateResponse, GetConsensusParamsRequest, GetConsensusParamsResponse,
20    GetContestedResourceIdentityVotesRequest, GetContestedResourceIdentityVotesResponse,
21    GetContestedResourceVoteStateRequest, GetContestedResourceVoteStateResponse,
22    GetContestedResourceVotersForIdentityRequest, GetContestedResourceVotersForIdentityResponse,
23    GetContestedResourcesRequest, GetContestedResourcesResponse, GetCurrentQuorumsInfoRequest,
24    GetCurrentQuorumsInfoResponse, GetDataContractHistoryRequest, GetDataContractHistoryResponse,
25    GetDataContractRequest, GetDataContractResponse, GetDataContractsRequest,
26    GetDataContractsResponse, GetDocumentHistoryRequest, GetDocumentHistoryResponse,
27    GetDocumentsRequest, GetDocumentsResponse, GetEpochsInfoRequest, GetEpochsInfoResponse,
28    GetEvonodesProposedEpochBlocksByIdsRequest, GetEvonodesProposedEpochBlocksByRangeRequest,
29    GetEvonodesProposedEpochBlocksResponse, GetFinalizedEpochInfosRequest,
30    GetFinalizedEpochInfosResponse, GetGroupActionSignersRequest, GetGroupActionSignersResponse,
31    GetGroupActionsRequest, GetGroupActionsResponse, GetGroupInfoRequest, GetGroupInfoResponse,
32    GetGroupInfosRequest, GetGroupInfosResponse, GetIdentitiesBalancesRequest,
33    GetIdentitiesBalancesResponse, GetIdentitiesContractKeysRequest,
34    GetIdentitiesContractKeysResponse, GetIdentitiesTokenBalancesRequest,
35    GetIdentitiesTokenBalancesResponse, GetIdentitiesTokenInfosRequest,
36    GetIdentitiesTokenInfosResponse, GetIdentityBalanceAndRevisionRequest,
37    GetIdentityBalanceAndRevisionResponse, GetIdentityBalanceRequest, GetIdentityBalanceResponse,
38    GetIdentityByNonUniquePublicKeyHashRequest, GetIdentityByNonUniquePublicKeyHashResponse,
39    GetIdentityByPublicKeyHashRequest, GetIdentityByPublicKeyHashResponse,
40    GetIdentityContractNonceRequest, GetIdentityContractNonceResponse, GetIdentityKeysRequest,
41    GetIdentityKeysResponse, GetIdentityNonceRequest, GetIdentityNonceResponse, GetIdentityRequest,
42    GetIdentityResponse, GetIdentityTokenBalancesRequest, GetIdentityTokenBalancesResponse,
43    GetIdentityTokenInfosRequest, GetIdentityTokenInfosResponse,
44    GetMostRecentShieldedAnchorRequest, GetMostRecentShieldedAnchorResponse,
45    GetPathElementsRequest, GetPathElementsResponse, GetPrefundedSpecializedBalanceRequest,
46    GetPrefundedSpecializedBalanceResponse, GetProtocolVersionUpgradeStateRequest,
47    GetProtocolVersionUpgradeStateResponse, GetProtocolVersionUpgradeVoteStatusRequest,
48    GetProtocolVersionUpgradeVoteStatusResponse, GetRecentAddressBalanceChangesRequest,
49    GetRecentAddressBalanceChangesResponse, GetRecentCompactedAddressBalanceChangesRequest,
50    GetRecentCompactedAddressBalanceChangesResponse, GetShieldedAnchorsRequest,
51    GetShieldedAnchorsResponse, GetShieldedEncryptedNotesRequest,
52    GetShieldedEncryptedNotesResponse, GetShieldedNotesCountRequest, GetShieldedNotesCountResponse,
53    GetShieldedNullifiersRequest, GetShieldedNullifiersResponse, GetShieldedPoolStateRequest,
54    GetShieldedPoolStateResponse, GetStatusRequest, GetStatusResponse, GetTokenContractInfoRequest,
55    GetTokenContractInfoResponse, GetTokenDirectPurchasePricesRequest,
56    GetTokenDirectPurchasePricesResponse, GetTokenPerpetualDistributionLastClaimRequest,
57    GetTokenPerpetualDistributionLastClaimResponse, GetTokenPreProgrammedDistributionsRequest,
58    GetTokenPreProgrammedDistributionsResponse, GetTokenStatusesRequest, GetTokenStatusesResponse,
59    GetTokenTotalSupplyRequest, GetTokenTotalSupplyResponse, GetTotalCreditsInPlatformRequest,
60    GetTotalCreditsInPlatformResponse, GetVotePollsByEndDateRequest, GetVotePollsByEndDateResponse,
61    WaitForStateTransitionResultRequest, WaitForStateTransitionResultResponse,
62};
63use dapi_grpc::tonic::{Code, Request, Response, Status};
64use dpp::version::PlatformVersion;
65use std::sync::atomic::Ordering;
66use std::sync::Arc;
67use std::thread::sleep;
68use std::time::Duration;
69use tracing::Instrument;
70
71const MAX_PATH_COMPONENTS: usize = 256;
72const MAX_GROVEDB_KEY_BYTES: usize = 255;
73const MAX_PATH_QUERY_BYTES: usize = 64 * 1024;
74
75/// Service to handle platform queries
76pub struct QueryService {
77    platform: Arc<Platform<DefaultCoreRPC>>,
78}
79
80type QueryMethod<RQ, RS> = fn(
81    &Platform<DefaultCoreRPC>,
82    RQ,
83    &PlatformState,
84    &PlatformVersion,
85) -> Result<QueryValidationResult<RS>, Error>;
86
87impl QueryService {
88    /// Creates new QueryService
89    pub fn new(platform: Arc<Platform<DefaultCoreRPC>>) -> Self {
90        Self { platform }
91    }
92
93    async fn handle_blocking_query<RQ, RS>(
94        &self,
95        request: Request<RQ>,
96        query_method: QueryMethod<RQ, RS>,
97        endpoint_name: &str,
98    ) -> Result<Response<RS>, Status>
99    where
100        RS: Clone + Send + 'static,
101        RQ: Send + Clone + 'static,
102    {
103        let mut response_duration_metric = query_duration_metric(endpoint_name);
104
105        let platform = Arc::clone(&self.platform);
106
107        let result = spawn_blocking_task_with_name_if_supported("query", move || {
108            let mut result;
109
110            let query_request = request.into_inner();
111
112            let mut query_counter = 0;
113
114            loop {
115                let platform_state = platform.state.load();
116
117                let platform_version = platform_state
118                    .current_platform_version()
119                    .map_err(|_| Status::unavailable("platform is not initialized"))?;
120
121                // Query is using Platform execution state and Drive state to during the execution.
122                // They are updating every block in finalize block ABCI handler.
123                // The problem is that these two operations aren't atomic and some latency between
124                // them could lead to data races. `committed_block_height_guard` counter that represents
125                // the latest the height of latest committed Drive state and logic bellow ensures
126                // that query is executed only after/before both states are updated.
127                let mut needs_restart = false;
128
129                loop {
130                    let committed_block_height_guard = platform
131                        .committed_block_height_guard
132                        .load(Ordering::Relaxed);
133                    let mut counter = 0;
134                    if platform_state.last_committed_block_height() == committed_block_height_guard
135                    {
136                        break;
137                    } else {
138                        counter += 1;
139                        sleep(Duration::from_millis(10))
140                    }
141
142                    // We try for up to 1 second
143                    if counter >= 100 {
144                        query_counter += 1;
145                        needs_restart = true;
146                        break;
147                    }
148                }
149
150                if query_counter > 3 {
151                    return Err(query_error_into_status(QueryError::NotServiceable(
152                        "platform is saturated (did not attempt query)".to_string(),
153                    )));
154                }
155
156                if needs_restart {
157                    continue;
158                }
159
160                result = query_method(
161                    &platform,
162                    query_request.clone(),
163                    &platform_state,
164                    platform_version,
165                );
166
167                let committed_block_height_guard = platform
168                    .committed_block_height_guard
169                    .load(Ordering::Relaxed);
170
171                if platform_state.last_committed_block_height() == committed_block_height_guard {
172                    // in this case the query almost certainly executed correctly
173                    break;
174                } else {
175                    query_counter += 1;
176
177                    if query_counter > 2 {
178                        // This should never be possible
179                        return Err(query_error_into_status(QueryError::NotServiceable(
180                            "platform is saturated".to_string(),
181                        )));
182                    }
183                }
184            }
185
186            let mut query_result = result.map_err(error_into_status)?;
187
188            if query_result.is_valid() {
189                let response = query_result
190                    .into_data()
191                    .map_err(|error| error_into_status(error.into()))?;
192
193                Ok(Response::new(response))
194            } else {
195                let error = query_result.errors.swap_remove(0);
196
197                Err(query_error_into_status(error))
198            }
199        })?
200        .instrument(tracing::trace_span!("query", endpoint_name))
201        .await
202        .map_err(|error| Status::internal(format!("query thread failed: {}", error)))?;
203
204        // Query logging and metrics
205        let code = match &result {
206            Ok(_) => Code::Ok,
207            Err(status) => status.code(),
208        };
209
210        let code_label = format!("{:?}", code).to_lowercase();
211
212        // Add code to response duration metric
213        let label = abci_response_code_metric_label(code);
214        response_duration_metric.add_label(label);
215
216        match code {
217            // User errors
218            Code::Ok
219            | Code::InvalidArgument
220            | Code::NotFound
221            | Code::AlreadyExists
222            | Code::ResourceExhausted
223            | Code::PermissionDenied
224            | Code::Unavailable
225            | Code::Aborted
226            | Code::FailedPrecondition
227            | Code::OutOfRange
228            | Code::Cancelled
229            | Code::DeadlineExceeded
230            | Code::Unauthenticated => {
231                let elapsed_time = response_duration_metric.elapsed().as_secs_f64();
232
233                tracing::trace!(
234                    elapsed_time,
235                    endpoint_name,
236                    code = code_label,
237                    "query '{}' executed with code {:?} in {} secs",
238                    endpoint_name,
239                    code,
240                    elapsed_time
241                );
242            }
243            // System errors
244            Code::Unknown | Code::Unimplemented | Code::Internal | Code::DataLoss => {
245                tracing::error!(
246                    endpoint_name,
247                    code = code_label,
248                    "query '{}' execution failed with code {:?}",
249                    endpoint_name,
250                    code
251                );
252            }
253        }
254
255        result
256    }
257}
258
259fn respond_with_unimplemented<RS>(name: &str) -> Result<Response<RS>, Status> {
260    tracing::error!("{} endpoint is called but it's not supported", name);
261
262    Err(Status::unimplemented("the endpoint is not supported"))
263}
264
265#[async_trait]
266impl PlatformService for QueryService {
267    async fn broadcast_state_transition(
268        &self,
269        _request: Request<BroadcastStateTransitionRequest>,
270    ) -> Result<Response<BroadcastStateTransitionResponse>, Status> {
271        respond_with_unimplemented("broadcast_state_transition")
272    }
273
274    async fn get_identity(
275        &self,
276        request: Request<GetIdentityRequest>,
277    ) -> Result<Response<GetIdentityResponse>, Status> {
278        self.handle_blocking_query(
279            request,
280            Platform::<DefaultCoreRPC>::query_identity,
281            "get_identity",
282        )
283        .await
284    }
285
286    async fn get_identities_contract_keys(
287        &self,
288        request: Request<GetIdentitiesContractKeysRequest>,
289    ) -> Result<Response<GetIdentitiesContractKeysResponse>, Status> {
290        self.handle_blocking_query(
291            request,
292            Platform::<DefaultCoreRPC>::query_identities_contract_keys,
293            "get_identities_contract_keys",
294        )
295        .await
296    }
297
298    async fn get_identity_keys(
299        &self,
300        request: Request<GetIdentityKeysRequest>,
301    ) -> Result<Response<GetIdentityKeysResponse>, Status> {
302        self.handle_blocking_query(
303            request,
304            Platform::<DefaultCoreRPC>::query_keys,
305            "get_identity_keys",
306        )
307        .await
308    }
309
310    async fn get_identity_nonce(
311        &self,
312        request: Request<GetIdentityNonceRequest>,
313    ) -> Result<Response<GetIdentityNonceResponse>, Status> {
314        self.handle_blocking_query(
315            request,
316            Platform::<DefaultCoreRPC>::query_identity_nonce,
317            "get_identity_nonce",
318        )
319        .await
320    }
321
322    async fn get_identity_contract_nonce(
323        &self,
324        request: Request<GetIdentityContractNonceRequest>,
325    ) -> Result<Response<GetIdentityContractNonceResponse>, Status> {
326        self.handle_blocking_query(
327            request,
328            Platform::<DefaultCoreRPC>::query_identity_contract_nonce,
329            "get_identity_contract_nonce",
330        )
331        .await
332    }
333
334    async fn get_identity_balance(
335        &self,
336        request: Request<GetIdentityBalanceRequest>,
337    ) -> Result<Response<GetIdentityBalanceResponse>, Status> {
338        self.handle_blocking_query(
339            request,
340            Platform::<DefaultCoreRPC>::query_balance,
341            "get_identity_balance",
342        )
343        .await
344    }
345
346    async fn get_identity_balance_and_revision(
347        &self,
348        request: Request<GetIdentityBalanceAndRevisionRequest>,
349    ) -> Result<Response<GetIdentityBalanceAndRevisionResponse>, Status> {
350        self.handle_blocking_query(
351            request,
352            Platform::<DefaultCoreRPC>::query_balance_and_revision,
353            "get_identity_balance_and_revision",
354        )
355        .await
356    }
357
358    async fn get_data_contract(
359        &self,
360        request: Request<GetDataContractRequest>,
361    ) -> Result<Response<GetDataContractResponse>, Status> {
362        self.handle_blocking_query(
363            request,
364            Platform::<DefaultCoreRPC>::query_data_contract,
365            "get_data_contract",
366        )
367        .await
368    }
369
370    async fn get_data_contract_history(
371        &self,
372        request: Request<GetDataContractHistoryRequest>,
373    ) -> Result<Response<GetDataContractHistoryResponse>, Status> {
374        self.handle_blocking_query(
375            request,
376            Platform::<DefaultCoreRPC>::query_data_contract_history,
377            "get_data_contract_history",
378        )
379        .await
380    }
381
382    async fn get_data_contracts(
383        &self,
384        request: Request<GetDataContractsRequest>,
385    ) -> Result<Response<GetDataContractsResponse>, Status> {
386        self.handle_blocking_query(
387            request,
388            Platform::<DefaultCoreRPC>::query_data_contracts,
389            "get_data_contracts",
390        )
391        .await
392    }
393
394    async fn get_document_history(
395        &self,
396        request: Request<GetDocumentHistoryRequest>,
397    ) -> Result<Response<GetDocumentHistoryResponse>, Status> {
398        self.handle_blocking_query(
399            request,
400            Platform::<DefaultCoreRPC>::query_document_history,
401            "get_document_history",
402        )
403        .await
404    }
405
406    async fn get_documents(
407        &self,
408        request: Request<GetDocumentsRequest>,
409    ) -> Result<Response<GetDocumentsResponse>, Status> {
410        self.handle_blocking_query(
411            request,
412            Platform::<DefaultCoreRPC>::query_documents,
413            "get_documents",
414        )
415        .await
416    }
417
418    async fn get_identity_by_public_key_hash(
419        &self,
420        request: Request<GetIdentityByPublicKeyHashRequest>,
421    ) -> Result<Response<GetIdentityByPublicKeyHashResponse>, Status> {
422        self.handle_blocking_query(
423            request,
424            Platform::<DefaultCoreRPC>::query_identity_by_public_key_hash,
425            "get_identity_by_public_key_hash",
426        )
427        .await
428    }
429
430    async fn get_identity_by_non_unique_public_key_hash(
431        &self,
432        request: Request<GetIdentityByNonUniquePublicKeyHashRequest>,
433    ) -> Result<Response<GetIdentityByNonUniquePublicKeyHashResponse>, Status> {
434        self.handle_blocking_query(
435            request,
436            Platform::<DefaultCoreRPC>::query_identity_by_non_unique_public_key_hash,
437            "get_identity_by_non_unique_public_key_hash",
438        )
439        .await
440    }
441
442    async fn wait_for_state_transition_result(
443        &self,
444        _request: Request<WaitForStateTransitionResultRequest>,
445    ) -> Result<Response<WaitForStateTransitionResultResponse>, Status> {
446        respond_with_unimplemented("wait_for_state_transition_result")
447    }
448
449    async fn get_consensus_params(
450        &self,
451        _request: Request<GetConsensusParamsRequest>,
452    ) -> Result<Response<GetConsensusParamsResponse>, Status> {
453        respond_with_unimplemented("get_consensus_params")
454    }
455
456    async fn get_protocol_version_upgrade_state(
457        &self,
458        request: Request<GetProtocolVersionUpgradeStateRequest>,
459    ) -> Result<Response<GetProtocolVersionUpgradeStateResponse>, Status> {
460        self.handle_blocking_query(
461            request,
462            Platform::<DefaultCoreRPC>::query_version_upgrade_state,
463            "get_protocol_version_upgrade_state",
464        )
465        .await
466    }
467
468    async fn get_protocol_version_upgrade_vote_status(
469        &self,
470        request: Request<GetProtocolVersionUpgradeVoteStatusRequest>,
471    ) -> Result<Response<GetProtocolVersionUpgradeVoteStatusResponse>, Status> {
472        self.handle_blocking_query(
473            request,
474            Platform::<DefaultCoreRPC>::query_version_upgrade_vote_status,
475            "get_protocol_version_upgrade_vote_status",
476        )
477        .await
478    }
479
480    async fn get_epochs_info(
481        &self,
482        request: Request<GetEpochsInfoRequest>,
483    ) -> Result<Response<GetEpochsInfoResponse>, Status> {
484        self.handle_blocking_query(
485            request,
486            Platform::<DefaultCoreRPC>::query_epoch_infos,
487            "get_epochs_info",
488        )
489        .await
490    }
491
492    async fn get_path_elements(
493        &self,
494        request: Request<GetPathElementsRequest>,
495    ) -> Result<Response<GetPathElementsResponse>, Status> {
496        validate_path_elements_request(request.get_ref())?;
497        self.handle_blocking_query(
498            request,
499            Platform::<DefaultCoreRPC>::query_path_elements,
500            "get_path_elements",
501        )
502        .await
503    }
504
505    async fn get_contested_resources(
506        &self,
507        request: Request<GetContestedResourcesRequest>,
508    ) -> Result<Response<GetContestedResourcesResponse>, Status> {
509        self.handle_blocking_query(
510            request,
511            Platform::<DefaultCoreRPC>::query_contested_resources,
512            "get_contested_resources",
513        )
514        .await
515    }
516
517    async fn get_contested_resource_vote_state(
518        &self,
519        request: Request<GetContestedResourceVoteStateRequest>,
520    ) -> Result<Response<GetContestedResourceVoteStateResponse>, Status> {
521        self.handle_blocking_query(
522            request,
523            Platform::<DefaultCoreRPC>::query_contested_resource_vote_state,
524            "get_contested_resource_vote_state",
525        )
526        .await
527    }
528
529    async fn get_contested_resource_voters_for_identity(
530        &self,
531        request: Request<GetContestedResourceVotersForIdentityRequest>,
532    ) -> Result<Response<GetContestedResourceVotersForIdentityResponse>, Status> {
533        self.handle_blocking_query(
534            request,
535            Platform::<DefaultCoreRPC>::query_contested_resource_voters_for_identity,
536            "get_contested_resource_voters_for_identity",
537        )
538        .await
539    }
540
541    async fn get_contested_resource_identity_votes(
542        &self,
543        request: Request<GetContestedResourceIdentityVotesRequest>,
544    ) -> Result<Response<GetContestedResourceIdentityVotesResponse>, Status> {
545        self.handle_blocking_query(
546            request,
547            Platform::<DefaultCoreRPC>::query_contested_resource_identity_votes,
548            "get_contested_resource_identity_votes",
549        )
550        .await
551    }
552
553    async fn get_vote_polls_by_end_date(
554        &self,
555        request: Request<GetVotePollsByEndDateRequest>,
556    ) -> Result<Response<GetVotePollsByEndDateResponse>, Status> {
557        self.handle_blocking_query(
558            request,
559            Platform::<DefaultCoreRPC>::query_vote_polls_by_end_date_query,
560            "get_vote_polls_by_end_date",
561        )
562        .await
563    }
564
565    async fn get_prefunded_specialized_balance(
566        &self,
567        request: Request<GetPrefundedSpecializedBalanceRequest>,
568    ) -> Result<Response<GetPrefundedSpecializedBalanceResponse>, Status> {
569        self.handle_blocking_query(
570            request,
571            Platform::<DefaultCoreRPC>::query_prefunded_specialized_balance,
572            "get_prefunded_specialized_balance",
573        )
574        .await
575    }
576
577    async fn get_total_credits_in_platform(
578        &self,
579        request: Request<GetTotalCreditsInPlatformRequest>,
580    ) -> Result<Response<GetTotalCreditsInPlatformResponse>, Status> {
581        self.handle_blocking_query(
582            request,
583            Platform::<DefaultCoreRPC>::query_total_credits_in_platform,
584            "get_total_credits_in_platform",
585        )
586        .await
587    }
588
589    async fn get_identities_balances(
590        &self,
591        request: Request<GetIdentitiesBalancesRequest>,
592    ) -> Result<Response<GetIdentitiesBalancesResponse>, Status> {
593        self.handle_blocking_query(
594            request,
595            Platform::<DefaultCoreRPC>::query_identities_balances,
596            "get_identities_balances",
597        )
598        .await
599    }
600
601    async fn get_status(
602        &self,
603        request: Request<GetStatusRequest>,
604    ) -> Result<Response<GetStatusResponse>, Status> {
605        self.handle_blocking_query(
606            request,
607            Platform::<DefaultCoreRPC>::query_partial_status,
608            "query_partial_status",
609        )
610        .await
611    }
612
613    async fn get_evonodes_proposed_epoch_blocks_by_ids(
614        &self,
615        request: Request<GetEvonodesProposedEpochBlocksByIdsRequest>,
616    ) -> Result<Response<GetEvonodesProposedEpochBlocksResponse>, Status> {
617        self.handle_blocking_query(
618            request,
619            Platform::<DefaultCoreRPC>::query_proposed_block_counts_by_evonode_ids,
620            "query_proposed_block_counts_by_evonode_ids",
621        )
622        .await
623    }
624
625    async fn get_evonodes_proposed_epoch_blocks_by_range(
626        &self,
627        request: Request<GetEvonodesProposedEpochBlocksByRangeRequest>,
628    ) -> Result<Response<GetEvonodesProposedEpochBlocksResponse>, Status> {
629        self.handle_blocking_query(
630            request,
631            Platform::<DefaultCoreRPC>::query_proposed_block_counts_by_range,
632            "query_proposed_block_counts_by_range",
633        )
634        .await
635    }
636
637    async fn get_current_quorums_info(
638        &self,
639        request: Request<GetCurrentQuorumsInfoRequest>,
640    ) -> Result<Response<GetCurrentQuorumsInfoResponse>, Status> {
641        self.handle_blocking_query(
642            request,
643            Platform::<DefaultCoreRPC>::query_current_quorums_info,
644            "query_current_quorums_info",
645        )
646        .await
647    }
648
649    async fn get_identity_token_balances(
650        &self,
651        request: Request<GetIdentityTokenBalancesRequest>,
652    ) -> Result<Response<GetIdentityTokenBalancesResponse>, Status> {
653        self.handle_blocking_query(
654            request,
655            Platform::<DefaultCoreRPC>::query_identity_token_balances,
656            "query_identity_token_balances",
657        )
658        .await
659    }
660
661    async fn get_identities_token_balances(
662        &self,
663        request: Request<GetIdentitiesTokenBalancesRequest>,
664    ) -> Result<Response<GetIdentitiesTokenBalancesResponse>, Status> {
665        self.handle_blocking_query(
666            request,
667            Platform::<DefaultCoreRPC>::query_identities_token_balances,
668            "query_identities_token_balances",
669        )
670        .await
671    }
672
673    async fn get_identity_token_infos(
674        &self,
675        request: Request<GetIdentityTokenInfosRequest>,
676    ) -> Result<Response<GetIdentityTokenInfosResponse>, Status> {
677        self.handle_blocking_query(
678            request,
679            Platform::<DefaultCoreRPC>::query_identity_token_infos,
680            "query_identity_token_infos",
681        )
682        .await
683    }
684
685    async fn get_identities_token_infos(
686        &self,
687        request: Request<GetIdentitiesTokenInfosRequest>,
688    ) -> Result<Response<GetIdentitiesTokenInfosResponse>, Status> {
689        self.handle_blocking_query(
690            request,
691            Platform::<DefaultCoreRPC>::query_identities_token_infos,
692            "query_identities_token_infos",
693        )
694        .await
695    }
696
697    async fn get_token_statuses(
698        &self,
699        request: Request<GetTokenStatusesRequest>,
700    ) -> Result<Response<GetTokenStatusesResponse>, Status> {
701        self.handle_blocking_query(
702            request,
703            Platform::<DefaultCoreRPC>::query_token_statuses,
704            "get_token_statuses",
705        )
706        .await
707    }
708
709    async fn get_token_pre_programmed_distributions(
710        &self,
711        request: Request<GetTokenPreProgrammedDistributionsRequest>,
712    ) -> Result<Response<GetTokenPreProgrammedDistributionsResponse>, Status> {
713        self.handle_blocking_query(
714            request,
715            Platform::<DefaultCoreRPC>::query_token_pre_programmed_distributions,
716            "get_token_pre_programmed_distributions",
717        )
718        .await
719    }
720
721    async fn get_token_total_supply(
722        &self,
723        request: Request<GetTokenTotalSupplyRequest>,
724    ) -> Result<Response<GetTokenTotalSupplyResponse>, Status> {
725        self.handle_blocking_query(
726            request,
727            Platform::<DefaultCoreRPC>::query_token_total_supply,
728            "get_token_total_supply",
729        )
730        .await
731    }
732
733    async fn get_group_info(
734        &self,
735        request: Request<GetGroupInfoRequest>,
736    ) -> Result<Response<GetGroupInfoResponse>, Status> {
737        self.handle_blocking_query(
738            request,
739            Platform::<DefaultCoreRPC>::query_group_info,
740            "get_group_info",
741        )
742        .await
743    }
744
745    async fn get_group_infos(
746        &self,
747        request: Request<GetGroupInfosRequest>,
748    ) -> Result<Response<GetGroupInfosResponse>, Status> {
749        self.handle_blocking_query(
750            request,
751            Platform::<DefaultCoreRPC>::query_group_infos,
752            "get_group_infos",
753        )
754        .await
755    }
756
757    async fn get_group_actions(
758        &self,
759        request: Request<GetGroupActionsRequest>,
760    ) -> Result<Response<GetGroupActionsResponse>, Status> {
761        self.handle_blocking_query(
762            request,
763            Platform::<DefaultCoreRPC>::query_group_actions,
764            "get_group_actions",
765        )
766        .await
767    }
768
769    async fn get_group_action_signers(
770        &self,
771        request: Request<GetGroupActionSignersRequest>,
772    ) -> Result<Response<GetGroupActionSignersResponse>, Status> {
773        self.handle_blocking_query(
774            request,
775            Platform::<DefaultCoreRPC>::query_group_action_signers,
776            "get_group_action_signers",
777        )
778        .await
779    }
780
781    async fn get_token_direct_purchase_prices(
782        &self,
783        request: Request<GetTokenDirectPurchasePricesRequest>,
784    ) -> Result<Response<GetTokenDirectPurchasePricesResponse>, Status> {
785        self.handle_blocking_query(
786            request,
787            Platform::<DefaultCoreRPC>::query_token_direct_purchase_prices,
788            "get_token_direct_purchase_prices",
789        )
790        .await
791    }
792
793    async fn get_token_contract_info(
794        &self,
795        request: Request<GetTokenContractInfoRequest>,
796    ) -> Result<Response<GetTokenContractInfoResponse>, Status> {
797        self.handle_blocking_query(
798            request,
799            Platform::<DefaultCoreRPC>::query_token_contract_info,
800            "get_token_contract_info",
801        )
802        .await
803    }
804
805    async fn get_token_perpetual_distribution_last_claim(
806        &self,
807        request: Request<GetTokenPerpetualDistributionLastClaimRequest>,
808    ) -> Result<Response<GetTokenPerpetualDistributionLastClaimResponse>, Status> {
809        self.handle_blocking_query(
810            request,
811            Platform::<DefaultCoreRPC>::query_token_perpetual_distribution_last_claim,
812            "get_token_perpetual_distribution_last_claim",
813        )
814        .await
815    }
816
817    async fn get_finalized_epoch_infos(
818        &self,
819        request: Request<GetFinalizedEpochInfosRequest>,
820    ) -> Result<Response<GetFinalizedEpochInfosResponse>, Status> {
821        self.handle_blocking_query(
822            request,
823            Platform::<DefaultCoreRPC>::query_finalized_epoch_infos,
824            "get_finalized_epoch_infos",
825        )
826        .await
827    }
828
829    async fn get_address_info(
830        &self,
831        request: Request<GetAddressInfoRequest>,
832    ) -> Result<Response<GetAddressInfoResponse>, Status> {
833        self.handle_blocking_query(
834            request,
835            Platform::<DefaultCoreRPC>::query_address_info,
836            "get_address_info",
837        )
838        .await
839    }
840
841    async fn get_addresses_infos(
842        &self,
843        request: Request<GetAddressesInfosRequest>,
844    ) -> Result<Response<GetAddressesInfosResponse>, Status> {
845        self.handle_blocking_query(
846            request,
847            Platform::<DefaultCoreRPC>::query_addresses_infos,
848            "get_addresses_infos",
849        )
850        .await
851    }
852
853    async fn get_addresses_trunk_state(
854        &self,
855        request: Request<GetAddressesTrunkStateRequest>,
856    ) -> Result<Response<GetAddressesTrunkStateResponse>, Status> {
857        self.handle_blocking_query(
858            request,
859            Platform::<DefaultCoreRPC>::query_addresses_trunk_state,
860            "get_addresses_trunk_state",
861        )
862        .await
863    }
864
865    async fn get_addresses_branch_state(
866        &self,
867        request: Request<GetAddressesBranchStateRequest>,
868    ) -> Result<Response<GetAddressesBranchStateResponse>, Status> {
869        self.handle_blocking_query(
870            request,
871            Platform::<DefaultCoreRPC>::query_addresses_branch_state,
872            "get_addresses_branch_state",
873        )
874        .await
875    }
876
877    async fn get_recent_address_balance_changes(
878        &self,
879        request: Request<GetRecentAddressBalanceChangesRequest>,
880    ) -> Result<Response<GetRecentAddressBalanceChangesResponse>, Status> {
881        self.handle_blocking_query(
882            request,
883            Platform::<DefaultCoreRPC>::query_recent_address_balance_changes,
884            "get_recent_address_balance_changes",
885        )
886        .await
887    }
888
889    async fn get_recent_compacted_address_balance_changes(
890        &self,
891        request: Request<GetRecentCompactedAddressBalanceChangesRequest>,
892    ) -> Result<Response<GetRecentCompactedAddressBalanceChangesResponse>, Status> {
893        self.handle_blocking_query(
894            request,
895            Platform::<DefaultCoreRPC>::query_recent_compacted_address_balance_changes,
896            "get_recent_compacted_address_balance_changes",
897        )
898        .await
899    }
900
901    async fn get_shielded_encrypted_notes(
902        &self,
903        request: Request<GetShieldedEncryptedNotesRequest>,
904    ) -> Result<Response<GetShieldedEncryptedNotesResponse>, Status> {
905        self.handle_blocking_query(
906            request,
907            Platform::<DefaultCoreRPC>::query_shielded_encrypted_notes,
908            "get_shielded_encrypted_notes",
909        )
910        .await
911    }
912
913    async fn get_shielded_anchors(
914        &self,
915        request: Request<GetShieldedAnchorsRequest>,
916    ) -> Result<Response<GetShieldedAnchorsResponse>, Status> {
917        self.handle_blocking_query(
918            request,
919            Platform::<DefaultCoreRPC>::query_shielded_anchors,
920            "get_shielded_anchors",
921        )
922        .await
923    }
924
925    async fn get_most_recent_shielded_anchor(
926        &self,
927        request: Request<GetMostRecentShieldedAnchorRequest>,
928    ) -> Result<Response<GetMostRecentShieldedAnchorResponse>, Status> {
929        self.handle_blocking_query(
930            request,
931            Platform::<DefaultCoreRPC>::query_most_recent_shielded_anchor,
932            "get_most_recent_shielded_anchor",
933        )
934        .await
935    }
936
937    async fn get_shielded_pool_state(
938        &self,
939        request: Request<GetShieldedPoolStateRequest>,
940    ) -> Result<Response<GetShieldedPoolStateResponse>, Status> {
941        self.handle_blocking_query(
942            request,
943            Platform::<DefaultCoreRPC>::query_shielded_pool_state,
944            "get_shielded_pool_state",
945        )
946        .await
947    }
948
949    async fn get_shielded_notes_count(
950        &self,
951        request: Request<GetShieldedNotesCountRequest>,
952    ) -> Result<Response<GetShieldedNotesCountResponse>, Status> {
953        self.handle_blocking_query(
954            request,
955            Platform::<DefaultCoreRPC>::query_shielded_notes_count,
956            "get_shielded_notes_count",
957        )
958        .await
959    }
960
961    async fn get_shielded_nullifiers(
962        &self,
963        request: Request<GetShieldedNullifiersRequest>,
964    ) -> Result<Response<GetShieldedNullifiersResponse>, Status> {
965        self.handle_blocking_query(
966            request,
967            Platform::<DefaultCoreRPC>::query_shielded_nullifiers,
968            "get_shielded_nullifiers",
969        )
970        .await
971    }
972}
973
974#[async_trait]
975impl DriveInternal for QueryService {
976    async fn get_proofs(
977        &self,
978        request: Request<GetProofsRequest>,
979    ) -> Result<Response<GetProofsResponse>, Status> {
980        self.handle_blocking_query(
981            request,
982            Platform::<DefaultCoreRPC>::query_proofs,
983            "get_proofs",
984        )
985        .await
986    }
987}
988
989fn query_error_into_status(error: QueryError) -> Status {
990    match error {
991        QueryError::NotFound(message) => Status::not_found(message),
992        QueryError::InvalidArgument(message) => Status::invalid_argument(message),
993        QueryError::Query(error) => Status::invalid_argument(error.to_string()),
994        QueryError::ResourceExhausted(message) => Status::resource_exhausted(message),
995        _ => {
996            tracing::error!("unexpected query error: {:?}", error);
997
998            Status::unknown(error.to_string())
999        }
1000    }
1001}
1002
1003fn error_into_status(error: Error) -> Status {
1004    Status::internal(format!("query: {}", error))
1005}
1006
1007fn validate_path_elements_request(request: &GetPathElementsRequest) -> Result<(), Status> {
1008    let v0 = match request.version.as_ref() {
1009        Some(get_path_elements_request::Version::V0(v0)) => v0,
1010        None => return Err(Status::invalid_argument("missing request version")),
1011    };
1012    let max_keys = PlatformVersion::latest()
1013        .drive_abci
1014        .query
1015        .max_returned_elements as usize;
1016    if v0.path.len() > MAX_PATH_COMPONENTS || v0.keys.len() > max_keys {
1017        return Err(Status::resource_exhausted(
1018            "too many path or key components",
1019        ));
1020    }
1021
1022    let total_bytes = v0
1023        .path
1024        .iter()
1025        .chain(&v0.keys)
1026        .try_fold(0usize, |total, component| {
1027            if component.len() > MAX_GROVEDB_KEY_BYTES {
1028                return Err(Status::resource_exhausted(
1029                    "path or key component is too large",
1030                ));
1031            }
1032            total
1033                .checked_add(component.len())
1034                .ok_or_else(|| Status::resource_exhausted("path query size overflow"))
1035        })?;
1036
1037    if total_bytes > MAX_PATH_QUERY_BYTES {
1038        return Err(Status::resource_exhausted(
1039            "aggregate path query bytes exceed limit",
1040        ));
1041    }
1042
1043    Ok(())
1044}
1045
1046#[cfg(test)]
1047mod tests {
1048    use super::*;
1049    use dapi_grpc::platform::v0::get_path_elements_request::GetPathElementsRequestV0;
1050
1051    #[test]
1052    fn path_elements_request_is_bounded_before_debug_formatting() {
1053        let request = GetPathElementsRequest {
1054            version: Some(get_path_elements_request::Version::V0(
1055                GetPathElementsRequestV0 {
1056                    path: vec![vec![]; MAX_PATH_COMPONENTS + 1],
1057                    keys: vec![],
1058                    prove: false,
1059                },
1060            )),
1061        };
1062
1063        let status = validate_path_elements_request(&request)
1064            .expect_err("expected excessive path depth to be rejected");
1065        assert_eq!(status.code(), Code::ResourceExhausted);
1066    }
1067}
1068
1069#[cfg(test)]
1070mod query_error_status_tests {
1071    use super::*;
1072
1073    #[test]
1074    fn resource_exhausted_query_error_maps_to_retryable_grpc_status() {
1075        let status = query_error_into_status(QueryError::ResourceExhausted(
1076            "server-side retained state is over capacity".to_string(),
1077        ));
1078
1079        assert_eq!(status.code(), Code::ResourceExhausted);
1080    }
1081}