1use alloc::borrow::ToOwned;
2use alloc::boxed::Box;
3use alloc::collections::{BTreeMap, BTreeSet};
4use alloc::string::{String, ToString};
5use alloc::vec::Vec;
6use core::error::Error;
7use core::pin::Pin;
8
9use miden_protocol::vm::FutureMaybeSend;
10
11type RpcFuture<T> = Pin<Box<dyn FutureMaybeSend<T>>>;
12
13use miden_objects::DecodeMessageExt;
14use miden_protocol::account::{
15 AccountCode,
16 AccountId,
17 AccountVaultPatch,
18 StorageMapPatchEntries,
19 StorageSlotName,
20};
21use miden_protocol::address::NetworkId;
22use miden_protocol::batch::{ProposedBatch, ProvenBatch};
23use miden_protocol::block::account_tree::AccountWitness;
24use miden_protocol::block::{BlockHeader, BlockNumber, SignedBlock};
25use miden_protocol::crypto::merkle::MerklePath;
26use miden_protocol::crypto::merkle::mmr::{Forest, MmrPath, MmrProof};
27use miden_protocol::note::{NoteId, NoteScript, NoteTag};
28use miden_protocol::transaction::ProvenTransaction;
29use miden_protocol::vm::ExecutionProof;
30use miden_protocol::{EMPTY_WORD, Word};
31use miden_tx::utils::sync::RwLock;
32use tonic::Status;
33use tracing::{info, warn};
34
35use super::domain::account::{
36 AccountProof,
37 AccountStorageRequirements,
38 GetAccountRequest,
39 StorageMapFetch,
40};
41use super::domain::note::{FetchedNote, SyncNotesBlock};
42use super::domain::nullifier::NullifierUpdate;
43use super::encryption::{
44 AttestedTransactionEncryptionKey,
45 NextTransactionEncryptionKey,
46 SealedTransactionInputs,
47 ValidatorAttestation,
48};
49use super::generated::rpc::GetAccountRequest as ProtoGetAccountRequest;
50use super::generated::rpc::get_account_request::AccountDetailRequest;
51use super::{Endpoint, NodeRpcClient, RpcEndpoint, RpcError, RpcStatusInfo};
52use crate::rpc::domain::account_vault::AccountVaultInfo;
53use crate::rpc::domain::limits::RpcLimits;
54use crate::rpc::domain::status::NetworkNoteStatusInfo;
55use crate::rpc::domain::storage_map::StorageMapInfo;
56use crate::rpc::domain::sync::{ChainMmrInfo, SyncTarget};
57use crate::rpc::domain::transaction::TransactionRecord;
58use crate::rpc::errors::node::{parse_node_error, parse_status_error};
59use crate::rpc::errors::{AcceptHeaderContext, AcceptHeaderError, GrpcError, RpcConversionError};
60use crate::rpc::generated::rpc::BlockRange;
61use crate::rpc::{AccountStateAt, generated as proto};
62
63mod api_client;
64mod retry;
65
66use api_client::api_client_wrapper::ApiClient;
67
68struct BlockPagination {
70 current_block_from: BlockNumber,
71 block_to: BlockNumber,
72 iterations: u32,
73}
74
75enum PaginationResult {
76 Continue,
77 Done {
78 chain_tip: BlockNumber,
79 block_num: BlockNumber,
80 },
81}
82
83impl BlockPagination {
84 const MAX_ITERATIONS: u32 = 1000;
89
90 fn new(block_from: BlockNumber, block_to: BlockNumber) -> Self {
91 Self {
92 current_block_from: block_from,
93 block_to,
94 iterations: 0,
95 }
96 }
97
98 fn current_block_from(&self) -> BlockNumber {
99 self.current_block_from
100 }
101
102 fn block_to(&self) -> BlockNumber {
103 self.block_to
104 }
105
106 fn advance(
107 &mut self,
108 block_num: BlockNumber,
109 chain_tip: BlockNumber,
110 ) -> Result<PaginationResult, RpcError> {
111 if self.iterations >= Self::MAX_ITERATIONS {
112 return Err(RpcError::PaginationError(
113 "too many pagination iterations, possible infinite loop".to_owned(),
114 ));
115 }
116 self.iterations += 1;
117
118 if block_num < self.current_block_from {
119 return Err(RpcError::PaginationError(
120 "invalid pagination: block_num went backwards".to_owned(),
121 ));
122 }
123
124 if block_num > self.block_to {
127 return Err(RpcError::PaginationError(format!(
128 "invalid pagination: block_num {block_num} is past the requested block_to {}",
129 self.block_to
130 )));
131 }
132
133 let target_block = self.block_to.min(chain_tip);
134
135 if block_num >= target_block {
136 return Ok(PaginationResult::Done { chain_tip, block_num });
137 }
138
139 self.current_block_from = BlockNumber::from(block_num.as_u32().saturating_add(1));
140
141 Ok(PaginationResult::Continue)
142 }
143}
144
145const DEFAULT_MAX_RESPONSE_SIZE_BYTES: usize = 4 * 1024 * 1024 * 115 / 100;
151
152pub struct GrpcClient {
162 client: RwLock<Option<ApiClient>>,
164 endpoint: String,
166 timeout_ms: u64,
168 genesis_commitment: RwLock<Option<Word>>,
170 limits: RwLock<Option<RpcLimits>>,
172 max_retries: u32,
174 retry_interval_ms: u64,
176 bearer_token: Option<String>,
180 max_decoding_message_size: usize,
183}
184
185impl GrpcClient {
186 pub fn new(endpoint: &Endpoint, timeout_ms: u64) -> GrpcClient {
189 GrpcClient {
190 client: RwLock::new(None),
191 endpoint: endpoint.to_string(),
192 timeout_ms,
193 genesis_commitment: RwLock::new(None),
194 limits: RwLock::new(None),
195 max_retries: retry::DEFAULT_MAX_RETRIES,
196 retry_interval_ms: retry::DEFAULT_RETRY_INTERVAL_MS,
197 bearer_token: None,
198 max_decoding_message_size: DEFAULT_MAX_RESPONSE_SIZE_BYTES,
199 }
200 }
201
202 #[must_use]
205 pub fn with_max_retries(mut self, max_retries: u32) -> Self {
206 self.max_retries = max_retries;
207 self
208 }
209
210 #[must_use]
213 pub fn with_retry_interval_ms(mut self, retry_interval_ms: u64) -> Self {
214 self.retry_interval_ms = retry_interval_ms;
215 self
216 }
217
218 #[must_use]
225 pub fn with_max_decoding_message_size(mut self, max_decoding_message_size: usize) -> Self {
226 self.max_decoding_message_size = max_decoding_message_size;
227 self
228 }
229
230 #[must_use]
254 pub fn with_bearer_auth(mut self, token: String) -> Self {
255 self.bearer_token = Some(token);
256 self
257 }
258
259 async fn ensure_connected(&self) -> Result<ApiClient, RpcError> {
262 if self.client.read().is_none() {
263 self.connect().await?;
264 }
265
266 Ok(self.client.read().as_ref().expect("rpc_api should be initialized").clone())
267 }
268
269 async fn connect(&self) -> Result<(), RpcError> {
272 let genesis_commitment = *self.genesis_commitment.read();
273 let new_client = ApiClient::new_client(
274 self.endpoint.clone(),
275 self.timeout_ms,
276 genesis_commitment,
277 self.bearer_token.clone(),
278 self.max_decoding_message_size,
279 )
280 .await?;
281 let mut client = self.client.write();
282 client.replace(new_client);
283
284 Ok(())
285 }
286
287 fn rpc_error_from_status(&self, endpoint: RpcEndpoint, status: Status) -> RpcError {
288 let genesis_commitment = self
289 .genesis_commitment
290 .read()
291 .as_ref()
292 .map_or_else(|| "none".to_string(), Word::to_hex);
293 let context = AcceptHeaderContext {
294 client_version: env!("CARGO_PKG_VERSION").to_string(),
295 genesis_commitment,
296 };
297 RpcError::from_grpc_error_with_context(endpoint, status, context)
298 }
299
300 async fn call_with_retry<T: Send + 'static>(
313 &self,
314 endpoint: RpcEndpoint,
315 mut call: impl FnMut(ApiClient) -> RpcFuture<Result<tonic::Response<T>, Status>>,
316 ) -> Result<tonic::Response<T>, RpcError> {
317 let mut retry_state =
318 retry::RetryState::new(endpoint, self.max_retries, self.retry_interval_ms);
319
320 loop {
321 let rpc_api = self.ensure_connected().await?;
322
323 match call(rpc_api).await {
324 Ok(response) => return Ok(response),
325 Err(status) if retry_state.should_retry(&status).await => {},
326 Err(status) => return Err(self.rpc_error_from_status(endpoint, status)),
327 }
328 }
329 }
330
331 pub async fn get_status_unversioned(&self) -> Result<RpcStatusInfo, RpcError> {
337 let mut rpc_api = ApiClient::new_client_without_accept_header(
338 self.endpoint.clone(),
339 self.timeout_ms,
340 self.bearer_token.clone(),
341 self.max_decoding_message_size,
342 )
343 .await?;
344 rpc_api
345 .status(proto::rpc::StatusRequest {})
346 .await
347 .map_err(|status| self.rpc_error_from_status(RpcEndpoint::Status, status))
348 .map(tonic::Response::into_inner)
349 .and_then(RpcStatusInfo::try_from)
350 }
351}
352
353#[cfg_attr(not(target_arch = "wasm32"), async_trait::async_trait)]
354#[cfg_attr(target_arch = "wasm32", async_trait::async_trait(?Send))]
355impl NodeRpcClient for GrpcClient {
356 fn has_genesis_commitment(&self) -> Option<Word> {
361 *self.genesis_commitment.read()
362 }
363
364 async fn set_genesis_commitment(&self, commitment: Word) -> Result<(), RpcError> {
365 if self.genesis_commitment.read().is_some() {
367 return Ok(());
369 }
370
371 self.genesis_commitment.write().replace(commitment);
373
374 let mut client_guard = self.client.write();
377 if let Some(client) = client_guard.as_mut() {
378 client.set_genesis_commitment(commitment);
379 }
380
381 Ok(())
382 }
383
384 async fn get_transaction_encryption_key(
385 &self,
386 ) -> Result<AttestedTransactionEncryptionKey, RpcError> {
387 let api_response = self
388 .call_with_retry(RpcEndpoint::GetTransactionEncryptionKey, |mut rpc_api| {
389 Box::pin(async move {
390 rpc_api
391 .get_transaction_encryption_key(
392 proto::rpc::GetTransactionEncryptionKeyRequest {},
393 )
394 .await
395 })
396 })
397 .await?;
398 let response = api_response
399 .into_inner()
400 .key
401 .ok_or(RpcError::ExpectedDataMissing("TransactionEncryptionKey".to_owned()))?;
402
403 let attestations = response
407 .attestations
408 .into_iter()
409 .filter_map(|attestation| {
410 let validator_key =
411 attestation.validator_public_key.and_then(|key| key.decode_and_verify().ok());
412 let signature =
413 attestation.signature.and_then(|signature| signature.decode_and_verify().ok());
414 let decoded = validator_key.zip(signature).map(|(validator_key, signature)| {
415 ValidatorAttestation { validator_key, signature }
416 });
417 if decoded.is_none() {
418 warn!(
419 "skipping a transaction encryption key attestation that failed to decode"
420 );
421 }
422 decoded
423 })
424 .collect::<Vec<_>>();
425
426 let wire_scheme = |scheme: i32| {
429 u32::try_from(scheme)
430 .map_err(|_| RpcError::InvalidResponse(format!("negative IES scheme '{scheme}'")))
431 };
432
433 let next_key = response
434 .next_key
435 .map(|next| {
436 Ok::<_, RpcError>(NextTransactionEncryptionKey {
437 scheme: wire_scheme(next.scheme)?,
438 key_id: next.key_id,
439 public_key: next.public_key,
440 rotation_block_num: next.rotation_block_num.into(),
441 })
442 })
443 .transpose()?;
444
445 Ok(AttestedTransactionEncryptionKey {
446 scheme: wire_scheme(response.scheme)?,
447 key_id: response.key_id,
448 public_key: response.public_key,
449 attestations,
450 next_key,
451 })
452 }
453
454 async fn submit_proven_transaction(
455 &self,
456 proven_transaction: &ProvenTransaction,
457 sealed_transaction_inputs: SealedTransactionInputs,
458 ) -> Result<BlockNumber, RpcError> {
459 let request = proto::submission::ProvenTransactionSubmission {
460 transaction: Some(proven_transaction.into()),
461 sealed_transaction_inputs: Some(sealed_transaction_inputs.into()),
462 };
463
464 let request = proto::rpc::SubmitProvenTxRequest { submission: Some(request) };
465 let api_response = self
466 .call_with_retry(RpcEndpoint::SubmitProvenTx, |mut rpc_api| {
467 let request = request.clone();
468 Box::pin(async move { rpc_api.submit_proven_tx(request).await })
469 })
470 .await?;
471
472 Ok(BlockNumber::from(api_response.into_inner().block_num))
473 }
474
475 async fn submit_proven_batch(
476 &self,
477 proven_batch: &ProvenBatch,
478 proposed_batch: &ProposedBatch,
479 sealed_transaction_inputs: Vec<SealedTransactionInputs>,
480 ) -> Result<BlockNumber, RpcError> {
481 let request = proto::submission::TransactionBatch {
482 batch: Some(proven_batch.into()),
483 proposed_batch: Some(proposed_batch.into()),
484 sealed_transaction_inputs: sealed_transaction_inputs
485 .into_iter()
486 .map(Into::into)
487 .collect(),
488 };
489
490 let request = proto::rpc::SubmitProvenTxBatchRequest { submission: Some(request) };
491 let api_response = self
492 .call_with_retry(RpcEndpoint::SubmitProvenBatch, |mut rpc_api| {
493 let request = request.clone();
494 Box::pin(async move { rpc_api.submit_proven_tx_batch(request).await })
495 })
496 .await?;
497
498 Ok(BlockNumber::from(api_response.into_inner().block_num))
499 }
500
501 async fn get_block_header_by_number(
502 &self,
503 block_num: Option<BlockNumber>,
504 include_mmr_proof: bool,
505 ) -> Result<(BlockHeader, Option<MmrProof>), RpcError> {
506 let request = proto::rpc::GetBlockHeaderByNumberRequest {
507 block_num: block_num.as_ref().map(BlockNumber::as_u32),
508 include_mmr_proof: Some(include_mmr_proof),
509 include_protocol_config: None,
510 };
511
512 info!("Calling GetBlockHeaderByNumber: {:?}", request);
513
514 let api_response = self
515 .call_with_retry(RpcEndpoint::GetBlockHeaderByNumber, |mut rpc_api| {
516 Box::pin(async move { rpc_api.get_block_header_by_number(request).await })
517 })
518 .await?;
519
520 let response = api_response.into_inner();
521
522 let block_header: BlockHeader = response
523 .block_header
524 .ok_or(RpcError::ExpectedDataMissing("BlockHeader".into()))?
525 .decode_and_build_unchecked()?;
526
527 let mmr_proof = if include_mmr_proof {
528 let forest = response
529 .chain_length
530 .ok_or(RpcError::ExpectedDataMissing("ChainLength".into()))?;
531 let merkle_path: MerklePath = response
532 .mmr_path
533 .ok_or(RpcError::ExpectedDataMissing("MmrPath".into()))?
534 .decode_and_verify()?;
535
536 let forest_size = usize::try_from(forest).expect("u64 should fit in usize");
537 let forest = Forest::new(forest_size).map_err(|_| {
538 RpcError::InvalidResponse(format!("invalid forest size: {forest_size}"))
539 })?;
540 Some(MmrProof::new(
541 MmrPath::new(forest, block_header.block_num().as_usize(), merkle_path),
542 block_header.commitment(),
543 ))
544 } else {
545 None
546 };
547
548 Ok((block_header, mmr_proof))
549 }
550
551 async fn get_notes_by_id(&self, note_ids: &[NoteId]) -> Result<Vec<FetchedNote>, RpcError> {
552 let limits = self.get_rpc_limits().await?;
553 let mut notes = Vec::with_capacity(note_ids.len());
554 for chunk in note_ids.chunks(limits.note_ids_limit as usize) {
555 let request = proto::rpc::GetNotesByIdRequest {
556 note_ids: chunk.iter().map(proto::note::NoteId::from).collect(),
557 };
558
559 let api_response = self
560 .call_with_retry(RpcEndpoint::GetNotesById, |mut rpc_api| {
561 let request = request.clone();
562 Box::pin(async move { rpc_api.get_notes_by_id(request).await })
563 })
564 .await?;
565
566 let response_notes = api_response
567 .into_inner()
568 .notes
569 .into_iter()
570 .map(FetchedNote::try_from)
571 .collect::<Result<Vec<FetchedNote>, RpcConversionError>>()?;
572
573 notes.extend(response_notes);
574 }
575 Ok(notes)
576 }
577
578 async fn sync_chain_mmr(
579 &self,
580 current_block_height: BlockNumber,
581 upper_bound: SyncTarget,
582 ) -> Result<ChainMmrInfo, RpcError> {
583 let finality_level: proto::rpc::FinalityLevel = upper_bound.into();
584
585 let request = proto::rpc::SyncChainMmrRequest {
586 current_client_block_height: current_block_height.as_u32(),
587 finality_level: finality_level.into(),
588 };
589
590 let response = self
591 .call_with_retry(RpcEndpoint::SyncChainMmr, |mut rpc_api| {
592 Box::pin(async move { rpc_api.sync_chain_mmr(request).await })
593 })
594 .await?;
595
596 response.into_inner().try_into()
597 }
598
599 async fn get_account(
611 &self,
612 account_id: AccountId,
613 request: GetAccountRequest,
614 ) -> Result<(BlockNumber, AccountProof), RpcError> {
615 let GetAccountRequest { storage, at, known_code, vault } = request;
616
617 let known_code_commitment = known_code.as_ref().map_or(EMPTY_WORD, AccountCode::commitment);
618 let mut known_codes_by_commitment: BTreeMap<Word, AccountCode> = BTreeMap::new();
619 if let Some(account_code) = known_code {
620 known_codes_by_commitment.insert(account_code.commitment(), account_code);
621 }
622
623 let requirements = match storage.clone() {
625 StorageMapFetch::Slots(reqs) => reqs,
626 StorageMapFetch::Skip | StorageMapFetch::All => AccountStorageRequirements::default(),
627 };
628
629 let account_details = if account_id.is_public() {
632 Some(AccountDetailRequest {
633 code_commitment: Some(known_code_commitment.into()),
634 asset_vault_commitment: vault.into(),
635 storage_request: storage.into(),
636 })
637 } else {
638 None
639 };
640
641 let block_num = match at {
642 AccountStateAt::Block(number) => Some(number.into()),
643 AccountStateAt::ChainTip => None,
644 };
645
646 let proto_request = ProtoGetAccountRequest {
647 account_id: Some(account_id.into()),
648 block_num,
649 details: account_details,
650 };
651
652 let response = self
653 .call_with_retry(RpcEndpoint::GetAccount, |mut rpc_api| {
654 let request = proto_request.clone();
655 Box::pin(async move { rpc_api.get_account(request).await })
656 })
657 .await?
658 .into_inner();
659
660 let account_witness: AccountWitness = response
661 .witness
662 .ok_or(RpcError::ExpectedDataMissing("AccountWitness".to_string()))?
663 .decode_and_verify()?;
664
665 let response_block_num: BlockNumber = response
666 .block_num
667 .ok_or(RpcError::ExpectedDataMissing("response block num".to_string()))?
668 .block_num
669 .into();
670
671 let headers = if account_witness.id().is_public() {
673 let details = response
674 .details
675 .ok_or(RpcError::ExpectedDataMissing("Account.Details".to_string()))?
676 .into_domain(&known_codes_by_commitment, &requirements)?;
677
678 Some(details)
679 } else {
680 None
681 };
682
683 let proof = AccountProof::new(account_witness, headers)
684 .map_err(|err| RpcError::InvalidResponse(err.to_string()))?;
685
686 Ok((response_block_num, proof))
687 }
688
689 async fn register_account(
690 &self,
691 invitation_code: &str,
692 account_id: AccountId,
693 ) -> Result<(), RpcError> {
694 let request = proto::rpc::RegisterAccountRequest {
696 invitation_code: invitation_code.to_string(),
697 account_id: Some(account_id.into()),
698 };
699
700 self.call_with_retry(RpcEndpoint::RegisterAccount, |mut rpc_api| {
701 let request = request.clone();
702 Box::pin(async move { rpc_api.register_account(request).await })
703 })
704 .await?;
705
706 Ok(())
707 }
708
709 async fn is_account_allowed(&self, account_id: AccountId) -> Result<bool, RpcError> {
710 let request = proto::rpc::IsAccountAllowedRequest { account_id: Some(account_id.into()) };
711
712 let response = self
713 .call_with_retry(RpcEndpoint::IsAccountAllowed, |mut rpc_api| {
714 Box::pin(async move { rpc_api.is_account_allowed(request).await })
715 })
716 .await?;
717
718 Ok(response.into_inner().allowed)
719 }
720
721 async fn is_invitation_code_valid(&self, invitation_code: &str) -> Result<bool, RpcError> {
722 let request =
723 proto::rpc::IsInvitationCodeValidRequest { invitation_code: invitation_code.into() };
724
725 let response = self
726 .call_with_retry(RpcEndpoint::IsInvitationCodeValid, |mut rpc_api| {
727 let request = request.clone();
728 Box::pin(async move { rpc_api.is_invitation_code_valid(request).await })
729 })
730 .await?;
731
732 Ok(response.into_inner().valid)
733 }
734
735 async fn sync_notes(
741 &self,
742 block_from: BlockNumber,
743 block_to: BlockNumber,
744 note_tags: &BTreeSet<NoteTag>,
745 ) -> Result<Vec<SyncNotesBlock>, RpcError> {
746 if note_tags.is_empty() {
747 return Ok(Vec::new());
748 }
749
750 let limits = self.get_rpc_limits().await?;
751 let tags: Vec<NoteTag> = note_tags.iter().copied().collect();
752
753 let mut merged_blocks: BTreeMap<BlockNumber, SyncNotesBlock> = BTreeMap::new();
756
757 for chunk in tags.chunks(limits.note_tags_limit as usize) {
758 let proto_tags: Vec<u32> = chunk.iter().map(|&t| t.into()).collect();
759 let mut pagination = BlockPagination::new(block_from, block_to);
760
761 loop {
762 let request = proto::rpc::SyncNotesRequest {
763 block_range: Some(BlockRange {
764 block_from: pagination.current_block_from().as_u32(),
765 block_to: block_to.as_u32(),
766 }),
767 note_tags: proto_tags.clone(),
768 };
769
770 let response = self
771 .call_with_retry(RpcEndpoint::SyncNotes, |mut rpc_api| {
772 let request = request.clone();
773 Box::pin(async move { rpc_api.sync_notes(request).await })
774 })
775 .await?
776 .into_inner();
777
778 let page = response.pagination_info.ok_or(RpcError::ExpectedDataMissing(
779 "SyncNotesResponse.pagination_info".to_owned(),
780 ))?;
781 let page_chain_tip = BlockNumber::from(page.chain_tip);
782 let page_block_to = BlockNumber::from(page.block_num);
783
784 for proto_block in response.blocks {
785 let block: SyncNotesBlock = proto_block.try_into()?;
786 let bn = block.block_header.block_num();
787 if let Some(existing) = merged_blocks.get_mut(&bn) {
788 for (id, note) in block.notes {
789 existing.notes.entry(id).or_insert(note);
790 }
791 } else {
792 merged_blocks.insert(bn, block);
793 }
794 }
795
796 match pagination.advance(page_block_to, page_chain_tip)? {
797 PaginationResult::Continue => {},
798 PaginationResult::Done { .. } => break,
799 }
800 }
801 }
802
803 Ok(merged_blocks.into_values().collect())
804 }
805
806 async fn sync_nullifiers(
807 &self,
808 prefixes: &[u16],
809 block_from: BlockNumber,
810 block_to: BlockNumber,
811 ) -> Result<Vec<NullifierUpdate>, RpcError> {
812 let limits = self.get_rpc_limits().await?;
813 let mut all_nullifiers = BTreeSet::new();
814
815 for chunk in prefixes.chunks(limits.nullifiers_limit as usize) {
818 let proto_prefixes: Vec<u32> = chunk.iter().map(|&x| u32::from(x)).collect();
819 let mut pagination = BlockPagination::new(block_from, block_to);
820
821 loop {
822 let request = proto::rpc::SyncNullifiersRequest {
823 nullifiers: proto_prefixes.clone(),
824 prefix_len: 16,
825 block_range: Some(BlockRange {
826 block_from: pagination.current_block_from().as_u32(),
827 block_to: pagination.block_to().as_u32(),
828 }),
829 };
830
831 let response = self
832 .call_with_retry(RpcEndpoint::SyncNullifiers, |mut rpc_api| {
833 let request = request.clone();
834 Box::pin(async move { rpc_api.sync_nullifiers(request).await })
835 })
836 .await?
837 .into_inner();
838
839 let batch_nullifiers = response
840 .nullifiers
841 .iter()
842 .map(TryFrom::try_from)
843 .collect::<Result<Vec<NullifierUpdate>, _>>()
844 .map_err(|err| RpcError::InvalidResponse(err.to_string()))?;
845
846 all_nullifiers.extend(batch_nullifiers);
847
848 let page = response.pagination_info.ok_or(RpcError::ExpectedDataMissing(
849 "SyncNullifiersResponse.pagination_info".to_owned(),
850 ))?;
851
852 match pagination.advance(page.block_num.into(), page.chain_tip.into())? {
853 PaginationResult::Continue => {},
854 PaginationResult::Done { .. } => break,
855 }
856 }
857 }
858 Ok(all_nullifiers.into_iter().collect::<Vec<_>>())
859 }
860
861 async fn get_block_by_number(
862 &self,
863 block_num: BlockNumber,
864 include_proof: bool,
865 ) -> Result<(SignedBlock, Option<ExecutionProof>), RpcError> {
866 let request = proto::rpc::GetBlockByNumberRequest {
867 block_num: block_num.as_u32(),
868 include_proof: Some(include_proof),
869 };
870
871 let response = self
872 .call_with_retry(RpcEndpoint::GetBlockByNumber, |mut rpc_api| {
873 Box::pin(async move { rpc_api.get_block_by_number(request).await })
874 })
875 .await?;
876
877 decode_block_response(response.into_inner())
878 }
879
880 async fn get_note_script_by_root(&self, root: Word) -> Result<Option<NoteScript>, RpcError> {
881 let request = proto::rpc::GetNoteScriptByRootRequest { root: Some(root.into()) };
882
883 let response = self
884 .call_with_retry(RpcEndpoint::GetNoteScriptByRoot, |mut rpc_api| {
885 let request = request.clone();
886 Box::pin(async move { rpc_api.get_note_script_by_root(request).await })
887 })
888 .await?;
889
890 let Some(script) = response.into_inner().script else {
892 return Ok(None);
893 };
894 let note_script: NoteScript = script.decode_and_verify()?;
895
896 Ok(Some(note_script))
897 }
898
899 async fn sync_storage_maps(
900 &self,
901 block_from: BlockNumber,
902 block_to: BlockNumber,
903 account_id: AccountId,
904 ) -> Result<StorageMapInfo, RpcError> {
905 let mut pagination = BlockPagination::new(block_from, block_to);
906 let mut map_entries: BTreeMap<StorageSlotName, StorageMapPatchEntries> = BTreeMap::new();
907
908 let (chain_tip, block_number) = loop {
909 let request = proto::rpc::SyncAccountStorageMapsRequest {
910 block_range: Some(BlockRange {
911 block_from: pagination.current_block_from().as_u32(),
912 block_to: block_to.as_u32(),
913 }),
914 account_id: Some(account_id.into()),
915 };
916 let response = self
917 .call_with_retry(RpcEndpoint::SyncStorageMaps, |mut rpc_api| {
918 Box::pin(async move { rpc_api.sync_account_storage_maps(request).await })
919 })
920 .await?;
921 let page = StorageMapInfo::try_from(response.into_inner())?;
922
923 for (slot_name, entries) in page.map_entries {
924 map_entries
925 .entry(slot_name)
926 .or_default()
927 .as_map_mut()
928 .extend(entries.into_map());
929 }
930
931 match pagination.advance(page.block_number, page.chain_tip)? {
932 PaginationResult::Continue => {},
933 PaginationResult::Done {
934 chain_tip: final_chain_tip,
935 block_num: final_block_num,
936 } => break (final_chain_tip, final_block_num),
937 }
938 };
939
940 Ok(StorageMapInfo { chain_tip, block_number, map_entries })
941 }
942
943 async fn sync_account_vault(
944 &self,
945 block_from: BlockNumber,
946 block_to: BlockNumber,
947 account_id: AccountId,
948 ) -> Result<AccountVaultInfo, RpcError> {
949 let mut pagination = BlockPagination::new(block_from, block_to);
950 let mut vault_patch = AccountVaultPatch::default();
951
952 let (chain_tip, block_number) = loop {
953 let request = proto::rpc::SyncAccountVaultRequest {
954 block_range: Some(BlockRange {
955 block_from: pagination.current_block_from().as_u32(),
956 block_to: block_to.as_u32(),
957 }),
958 account_id: Some(account_id.into()),
959 };
960 let response = self
961 .call_with_retry(RpcEndpoint::SyncAccountVault, |mut rpc_api| {
962 Box::pin(async move { rpc_api.sync_account_vault(request).await })
963 })
964 .await?;
965 let page = AccountVaultInfo::try_from(response.into_inner())?;
966
967 vault_patch.merge(page.vault_patch);
968
969 match pagination.advance(page.block_number, page.chain_tip)? {
970 PaginationResult::Continue => {},
971 PaginationResult::Done {
972 chain_tip: final_chain_tip,
973 block_num: final_block_num,
974 } => break (final_chain_tip, final_block_num),
975 }
976 };
977
978 Ok(AccountVaultInfo { chain_tip, block_number, vault_patch })
979 }
980
981 async fn sync_transactions(
987 &self,
988 block_from: BlockNumber,
989 block_to: BlockNumber,
990 account_ids: Vec<AccountId>,
991 ) -> Result<Vec<TransactionRecord>, RpcError> {
992 if account_ids.is_empty() {
993 return Ok(Vec::new());
994 }
995
996 let limits = self.get_rpc_limits().await?;
997 let mut transactions: Vec<TransactionRecord> = Vec::new();
998
999 for chunk in account_ids.chunks(limits.account_ids_limit as usize) {
1000 let proto_account_ids: Vec<_> = chunk.iter().map(|acc_id| (*acc_id).into()).collect();
1001 let mut pagination = BlockPagination::new(block_from, block_to);
1002
1003 loop {
1004 let request = proto::rpc::SyncTransactionsRequest {
1005 block_range: Some(BlockRange {
1006 block_from: pagination.current_block_from().as_u32(),
1007 block_to: block_to.as_u32(),
1008 }),
1009 account_ids: proto_account_ids.clone(),
1010 };
1011
1012 let response = self
1013 .call_with_retry(RpcEndpoint::SyncTransactions, |mut rpc_api| {
1014 let request = request.clone();
1015 Box::pin(async move { rpc_api.sync_transactions(request).await })
1016 })
1017 .await?
1018 .into_inner();
1019
1020 let page = response.pagination_info.ok_or(RpcError::ExpectedDataMissing(
1021 "SyncTransactionsResponse.pagination_info".to_owned(),
1022 ))?;
1023 let page_chain_tip = BlockNumber::from(page.chain_tip);
1024 let page_block_to = BlockNumber::from(page.block_num);
1025
1026 for proto_tx in response.transactions {
1027 transactions.push(TransactionRecord::try_from(proto_tx)?);
1028 }
1029
1030 match pagination.advance(page_block_to, page_chain_tip)? {
1031 PaginationResult::Continue => {},
1032 PaginationResult::Done { .. } => break,
1033 }
1034 }
1035 }
1036
1037 Ok(transactions)
1038 }
1039
1040 async fn get_network_id(&self) -> Result<NetworkId, RpcError> {
1041 let endpoint: Endpoint =
1042 Endpoint::try_from(self.endpoint.as_str()).map_err(RpcError::InvalidNodeEndpoint)?;
1043 Ok(endpoint.to_network_id())
1044 }
1045
1046 async fn get_rpc_limits(&self) -> Result<RpcLimits, RpcError> {
1047 if let Some(limits) = *self.limits.read() {
1048 return Ok(limits);
1049 }
1050
1051 let response = self
1052 .call_with_retry(RpcEndpoint::GetLimits, |mut rpc_api| {
1053 Box::pin(async move { rpc_api.get_limits(proto::rpc::GetLimitsRequest {}).await })
1054 })
1055 .await?;
1056 let limits = RpcLimits::try_from(response.into_inner()).map_err(RpcError::from)?;
1057
1058 self.limits.write().replace(limits);
1060 Ok(limits)
1061 }
1062
1063 fn has_rpc_limits(&self) -> Option<RpcLimits> {
1064 *self.limits.read()
1065 }
1066
1067 async fn set_rpc_limits(&self, limits: RpcLimits) {
1068 self.limits.write().replace(limits);
1069 }
1070
1071 async fn get_status_unversioned(&self) -> Result<RpcStatusInfo, RpcError> {
1072 GrpcClient::get_status_unversioned(self).await
1073 }
1074
1075 async fn get_network_note_status(
1076 &self,
1077 note_id: NoteId,
1078 ) -> Result<NetworkNoteStatusInfo, RpcError> {
1079 let request = proto::rpc::GetNetworkNoteStatusRequest {
1080 note_id: Some(proto::note::NoteId::from(¬e_id)),
1081 };
1082
1083 let response = self
1084 .call_with_retry(RpcEndpoint::GetNetworkNoteStatus, |mut rpc_api| {
1085 let request = request.clone();
1086 Box::pin(async move { rpc_api.get_network_note_status(request).await })
1087 })
1088 .await?;
1089
1090 response.into_inner().try_into()
1091 }
1092}
1093
1094impl RpcError {
1098 pub fn from_grpc_error_with_context(
1099 endpoint: RpcEndpoint,
1100 status: Status,
1101 context: AcceptHeaderContext,
1102 ) -> Self {
1103 if let Some(accept_error) =
1104 AcceptHeaderError::try_from_message_with_context(status.message(), context)
1105 {
1106 return Self::AcceptHeaderError(accept_error);
1107 }
1108
1109 let error_kind = GrpcError::from(&status);
1110
1111 let endpoint_error = parse_node_error(&endpoint, status.details(), status.message())
1113 .or_else(|| parse_status_error(&endpoint, &error_kind, status.message()));
1114
1115 let source = Box::new(status) as Box<dyn Error + Send + Sync + 'static>;
1116
1117 Self::RequestError {
1118 endpoint,
1119 error_kind,
1120 endpoint_error,
1121 source: Some(source),
1122 }
1123 }
1124}
1125
1126impl From<&Status> for GrpcError {
1127 fn from(status: &Status) -> Self {
1128 GrpcError::from_code(status.code() as i32, Some(status.message().to_string()))
1129 }
1130}
1131
1132fn decode_block_response(
1141 response: proto::rpc::GetBlockByNumberResponse,
1142) -> Result<(SignedBlock, Option<ExecutionProof>), RpcError> {
1143 let block: SignedBlock = response
1146 .block
1147 .ok_or(RpcError::ExpectedDataMissing("GetBlockByNumberResponse.block".to_string()))?
1148 .decode_and_build_unchecked()?;
1149
1150 let proof = response.proof.map(ExecutionProof::try_from).transpose()?;
1153
1154 Ok((block, proof))
1155}
1156
1157#[cfg(test)]
1158mod tests {
1159 use std::boxed::Box;
1160 use std::vec;
1161
1162 use miden_protocol::Word;
1163 use miden_protocol::block::{BlockNumber, SignedBlock};
1164 use miden_testing::MockChain;
1165
1166 use super::{
1167 BlockPagination,
1168 DEFAULT_MAX_RESPONSE_SIZE_BYTES,
1169 GrpcClient,
1170 PaginationResult,
1171 decode_block_response,
1172 proto,
1173 };
1174 use crate::alloc::string::ToString;
1175 use crate::rpc::{Endpoint, NodeRpcClient, RpcError};
1176
1177 fn assert_send_sync<T: Send + Sync>() {}
1178
1179 fn genesis_block_messages()
1181 -> (proto::blockchain::SignedBlock, proto::primitives::ExecutionProof) {
1182 let chain = MockChain::new();
1183 let block = chain.proven_blocks().first().expect("the chain has a genesis block").clone();
1184 let (header, body, signatures, proof) = block.into_parts();
1185
1186 (SignedBlock::new_unchecked(header, body, signatures).into(), proof.into())
1187 }
1188
1189 #[test]
1190 fn decode_block_response_reads_a_requested_proof() {
1191 let (block, proof_message) = genesis_block_messages();
1192 let response = proto::rpc::GetBlockByNumberResponse {
1193 block: Some(block),
1194 proof: Some(proof_message.clone()),
1195 };
1196
1197 let (_block, proof) = decode_block_response(response).unwrap();
1198
1199 let decoded: proto::primitives::ExecutionProof =
1200 proof.expect("the response carries a proof").into();
1201 assert_eq!(decoded, proof_message);
1202 }
1203
1204 #[test]
1205 fn decode_block_response_omits_an_absent_proof() {
1206 let (block, _) = genesis_block_messages();
1207 let response = proto::rpc::GetBlockByNumberResponse { block: Some(block), proof: None };
1208
1209 let (_block, proof) = decode_block_response(response).unwrap();
1210
1211 assert!(proof.is_none());
1212 }
1213
1214 #[test]
1215 fn decode_block_response_rejects_malformed_proof_bytes() {
1216 let (block, _) = genesis_block_messages();
1217 let response = proto::rpc::GetBlockByNumberResponse {
1218 block: Some(block),
1219 proof: Some(proto::primitives::ExecutionProof { encoded: vec![0xff; 32] }),
1220 };
1221
1222 let res = decode_block_response(response);
1223
1224 assert!(matches!(res, Err(RpcError::DeserializationError(_))));
1225 }
1226
1227 #[test]
1228 fn decode_block_response_rejects_an_absent_block() {
1229 let response = proto::rpc::GetBlockByNumberResponse { block: None, proof: None };
1230
1231 let res = decode_block_response(response);
1232
1233 assert!(matches!(res, Err(RpcError::ExpectedDataMissing(_))));
1234 }
1235
1236 #[test]
1237 fn is_send_sync() {
1238 assert_send_sync::<GrpcClient>();
1239 assert_send_sync::<Box<dyn NodeRpcClient>>();
1240 }
1241
1242 #[test]
1243 fn block_pagination_errors_when_block_num_goes_backwards() {
1244 let mut pagination = BlockPagination::new(10_u32.into(), 20_u32.into());
1245
1246 let res = pagination.advance(9_u32.into(), 20_u32.into());
1247 assert!(matches!(res, Err(RpcError::PaginationError(_))));
1248 }
1249
1250 #[test]
1251 fn block_pagination_errors_when_block_num_passes_block_to() {
1252 let mut pagination = BlockPagination::new(10_u32.into(), 20_u32.into());
1253
1254 let res = pagination.advance(21_u32.into(), 100_u32.into());
1255 assert!(matches!(res, Err(RpcError::PaginationError(_))));
1256 }
1257
1258 #[test]
1259 fn block_pagination_errors_after_max_iterations() {
1260 let mut pagination = BlockPagination::new(0_u32.into(), 10_000_u32.into());
1261 let chain_tip: BlockNumber = 10_000_u32.into();
1262
1263 for _ in 0..BlockPagination::MAX_ITERATIONS {
1264 let current = pagination.current_block_from();
1265 let res = pagination
1266 .advance(current, chain_tip)
1267 .expect("expected pagination to continue within iteration limit");
1268 assert!(matches!(res, PaginationResult::Continue));
1269 }
1270
1271 let res = pagination.advance(pagination.current_block_from(), chain_tip);
1272 assert!(matches!(res, Err(RpcError::PaginationError(_))));
1273 }
1274
1275 #[test]
1276 fn block_pagination_stops_at_min_of_block_to_and_chain_tip() {
1277 let mut pagination = BlockPagination::new(0_u32.into(), 50_u32.into());
1279
1280 let res = pagination
1281 .advance(30_u32.into(), 30_u32.into())
1282 .expect("expected pagination to succeed");
1283
1284 assert!(matches!(
1285 res,
1286 PaginationResult::Done {
1287 chain_tip,
1288 block_num
1289 } if chain_tip.as_u32() == 30 && block_num.as_u32() == 30
1290 ));
1291 }
1292
1293 #[test]
1294 fn block_pagination_advances_cursor_by_one() {
1295 let mut pagination = BlockPagination::new(5_u32.into(), 100_u32.into());
1296
1297 let res = pagination
1298 .advance(5_u32.into(), 100_u32.into())
1299 .expect("expected pagination to succeed");
1300 assert!(matches!(res, PaginationResult::Continue));
1301 assert_eq!(pagination.current_block_from().as_u32(), 6);
1302 }
1303
1304 async fn dyn_trait_send_fut(client: Box<dyn NodeRpcClient>) {
1306 let res = client.get_block_header_by_number(None, false).await;
1308 assert!(res.is_ok());
1309 }
1310
1311 #[tokio::test]
1312 async fn future_is_send() {
1313 let endpoint = &Endpoint::devnet();
1314 let client = GrpcClient::new(endpoint, 10000);
1315 let client: Box<GrpcClient> = client.into();
1316 tokio::task::spawn(async move { dyn_trait_send_fut(client).await });
1317 }
1318
1319 #[tokio::test]
1320 async fn set_genesis_commitment_sets_the_commitment_when_its_not_already_set() {
1321 let endpoint = &Endpoint::devnet();
1322 let client = GrpcClient::new(endpoint, 10000);
1323
1324 assert!(client.genesis_commitment.read().is_none());
1325
1326 let commitment = Word::default();
1327 client.set_genesis_commitment(commitment).await.unwrap();
1328
1329 assert_eq!(client.genesis_commitment.read().unwrap(), commitment);
1330 }
1331
1332 #[tokio::test]
1333 async fn set_genesis_commitment_does_nothing_if_the_commitment_is_already_set() {
1334 let endpoint = &Endpoint::devnet();
1335 let client = GrpcClient::new(endpoint, 10000);
1336
1337 let initial_commitment = Word::default();
1338 client.set_genesis_commitment(initial_commitment).await.unwrap();
1339
1340 let new_commitment = Word::from([1u32, 2, 3, 4]);
1341 client.set_genesis_commitment(new_commitment).await.unwrap();
1342
1343 assert_eq!(client.genesis_commitment.read().unwrap(), initial_commitment);
1344 }
1345
1346 #[tokio::test]
1347 async fn set_genesis_commitment_updates_the_client_if_already_connected() {
1348 let endpoint = &Endpoint::devnet();
1349 let client = GrpcClient::new(endpoint, 10000);
1350
1351 client.connect().await.unwrap();
1353
1354 let commitment = Word::default();
1355 client.set_genesis_commitment(commitment).await.unwrap();
1356
1357 assert_eq!(client.genesis_commitment.read().unwrap(), commitment);
1358 assert!(client.client.read().as_ref().is_some());
1359 }
1360
1361 #[test]
1362 fn with_bearer_auth_stores_token() {
1363 let endpoint = &Endpoint::devnet();
1364 let client = GrpcClient::new(endpoint, 10000).with_bearer_auth("token-one".to_string());
1365
1366 assert_eq!(client.bearer_token.as_deref(), Some("token-one"));
1367 }
1368
1369 #[test]
1370 fn with_bearer_auth_overwrites_on_repeat_call() {
1371 let endpoint = &Endpoint::devnet();
1372 let client = GrpcClient::new(endpoint, 10000)
1373 .with_bearer_auth("token-one".to_string())
1374 .with_bearer_auth("token-two".to_string());
1375
1376 assert_eq!(client.bearer_token.as_deref(), Some("token-two"));
1378 }
1379
1380 #[tokio::test]
1381 async fn with_bearer_auth_surfaces_invalid_ascii_value_at_connect_time() {
1382 let endpoint = &Endpoint::devnet();
1386 let client = GrpcClient::new(endpoint, 10000).with_bearer_auth("bad\nvalue".to_string());
1387
1388 let err = client.connect().await.expect_err("expected invalid token to fail connect");
1389 assert!(
1390 matches!(err, RpcError::ConnectionError(_)),
1391 "expected ConnectionError, got {err:?}",
1392 );
1393 }
1394
1395 #[tokio::test]
1396 async fn with_bearer_auth_is_preserved_across_set_genesis_commitment() {
1397 let endpoint = &Endpoint::devnet();
1398 let client = GrpcClient::new(endpoint, 10000).with_bearer_auth("token".to_string());
1399 client.connect().await.unwrap();
1400
1401 client.set_genesis_commitment(Word::default()).await.unwrap();
1402
1403 assert_eq!(client.bearer_token.as_deref(), Some("token"));
1405 assert!(client.client.read().as_ref().is_some());
1406 }
1407
1408 #[test]
1409 fn with_max_decoding_message_size_overrides_default() {
1410 let endpoint = &Endpoint::devnet();
1411
1412 let default_client = GrpcClient::new(endpoint, 10_000);
1414 assert_eq!(default_client.max_decoding_message_size, DEFAULT_MAX_RESPONSE_SIZE_BYTES);
1415
1416 let custom =
1418 GrpcClient::new(endpoint, 10_000).with_max_decoding_message_size(8 * 1024 * 1024);
1419 assert_eq!(custom.max_decoding_message_size, 8 * 1024 * 1024);
1420 }
1421
1422 #[tokio::test]
1433 #[ignore = "requires network access to public testnet"]
1434 async fn with_bearer_auth_does_not_break_real_rpc_against_testnet() {
1435 let endpoint = &Endpoint::testnet();
1436 let client = GrpcClient::new(endpoint, 10_000).with_bearer_auth("smoke-test".to_string());
1437
1438 let status = client
1439 .get_status_unversioned()
1440 .await
1441 .expect("testnet status with caller auth header must succeed");
1442 assert!(!status.version.is_empty(), "status must include a server version");
1443 }
1444}