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