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