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 sync_notes(
727 &self,
728 block_from: BlockNumber,
729 block_to: BlockNumber,
730 note_tags: &BTreeSet<NoteTag>,
731 ) -> Result<Vec<SyncNotesBlock>, RpcError> {
732 if note_tags.is_empty() {
733 return Ok(Vec::new());
734 }
735
736 let limits = self.get_rpc_limits().await?;
737 let tags: Vec<NoteTag> = note_tags.iter().copied().collect();
738
739 let mut merged_blocks: BTreeMap<BlockNumber, SyncNotesBlock> = BTreeMap::new();
742
743 for chunk in tags.chunks(limits.note_tags_limit as usize) {
744 let proto_tags: Vec<u32> = chunk.iter().map(|&t| t.into()).collect();
745 let mut pagination = BlockPagination::new(block_from, block_to);
746
747 loop {
748 let request = proto::rpc::SyncNotesRequest {
749 block_range: Some(BlockRange {
750 block_from: pagination.current_block_from().as_u32(),
751 block_to: block_to.as_u32(),
752 }),
753 note_tags: proto_tags.clone(),
754 };
755
756 let response = self
757 .call_with_retry(RpcEndpoint::SyncNotes, |mut rpc_api| {
758 let request = request.clone();
759 Box::pin(async move { rpc_api.sync_notes(request).await })
760 })
761 .await?
762 .into_inner();
763
764 let page = response.pagination_info.ok_or(RpcError::ExpectedDataMissing(
765 "SyncNotesResponse.pagination_info".to_owned(),
766 ))?;
767 let page_chain_tip = BlockNumber::from(page.chain_tip);
768 let page_block_to = BlockNumber::from(page.block_num);
769
770 for proto_block in response.blocks {
771 let block: SyncNotesBlock = proto_block.try_into()?;
772 let bn = block.block_header.block_num();
773 if let Some(existing) = merged_blocks.get_mut(&bn) {
774 for (id, note) in block.notes {
775 existing.notes.entry(id).or_insert(note);
776 }
777 } else {
778 merged_blocks.insert(bn, block);
779 }
780 }
781
782 match pagination.advance(page_block_to, page_chain_tip)? {
783 PaginationResult::Continue => {},
784 PaginationResult::Done { .. } => break,
785 }
786 }
787 }
788
789 Ok(merged_blocks.into_values().collect())
790 }
791
792 async fn sync_nullifiers(
793 &self,
794 prefixes: &[u16],
795 block_from: BlockNumber,
796 block_to: BlockNumber,
797 ) -> Result<Vec<NullifierUpdate>, RpcError> {
798 let limits = self.get_rpc_limits().await?;
799 let mut all_nullifiers = BTreeSet::new();
800
801 for chunk in prefixes.chunks(limits.nullifiers_limit as usize) {
804 let proto_prefixes: Vec<u32> = chunk.iter().map(|&x| u32::from(x)).collect();
805 let mut pagination = BlockPagination::new(block_from, block_to);
806
807 loop {
808 let request = proto::rpc::SyncNullifiersRequest {
809 nullifiers: proto_prefixes.clone(),
810 prefix_len: 16,
811 block_range: Some(BlockRange {
812 block_from: pagination.current_block_from().as_u32(),
813 block_to: pagination.block_to().as_u32(),
814 }),
815 };
816
817 let response = self
818 .call_with_retry(RpcEndpoint::SyncNullifiers, |mut rpc_api| {
819 let request = request.clone();
820 Box::pin(async move { rpc_api.sync_nullifiers(request).await })
821 })
822 .await?
823 .into_inner();
824
825 let batch_nullifiers = response
826 .nullifiers
827 .iter()
828 .map(TryFrom::try_from)
829 .collect::<Result<Vec<NullifierUpdate>, _>>()
830 .map_err(|err| RpcError::InvalidResponse(err.to_string()))?;
831
832 all_nullifiers.extend(batch_nullifiers);
833
834 let page = response.pagination_info.ok_or(RpcError::ExpectedDataMissing(
835 "SyncNullifiersResponse.pagination_info".to_owned(),
836 ))?;
837
838 match pagination.advance(page.block_num.into(), page.chain_tip.into())? {
839 PaginationResult::Continue => {},
840 PaginationResult::Done { .. } => break,
841 }
842 }
843 }
844 Ok(all_nullifiers.into_iter().collect::<Vec<_>>())
845 }
846
847 async fn get_block_by_number(
848 &self,
849 block_num: BlockNumber,
850 include_proof: bool,
851 ) -> Result<(SignedBlock, Option<ExecutionProof>), RpcError> {
852 let request = proto::rpc::GetBlockByNumberRequest {
853 block_num: block_num.as_u32(),
854 include_proof: Some(include_proof),
855 };
856
857 let response = self
858 .call_with_retry(RpcEndpoint::GetBlockByNumber, |mut rpc_api| {
859 Box::pin(async move { rpc_api.get_block_by_number(request).await })
860 })
861 .await?;
862
863 decode_block_response(response.into_inner())
864 }
865
866 async fn get_note_script_by_root(&self, root: Word) -> Result<Option<NoteScript>, RpcError> {
867 let request = proto::rpc::GetNoteScriptByRootRequest { root: Some(root.into()) };
868
869 let response = self
870 .call_with_retry(RpcEndpoint::GetNoteScriptByRoot, |mut rpc_api| {
871 let request = request.clone();
872 Box::pin(async move { rpc_api.get_note_script_by_root(request).await })
873 })
874 .await?;
875
876 let Some(script) = response.into_inner().script else {
878 return Ok(None);
879 };
880 let note_script: NoteScript = script.decode_and_verify()?;
881
882 Ok(Some(note_script))
883 }
884
885 async fn sync_storage_maps(
886 &self,
887 block_from: BlockNumber,
888 block_to: BlockNumber,
889 account_id: AccountId,
890 ) -> Result<StorageMapInfo, RpcError> {
891 let mut pagination = BlockPagination::new(block_from, block_to);
892 let mut map_entries: BTreeMap<StorageSlotName, StorageMapPatchEntries> = BTreeMap::new();
893
894 let (chain_tip, block_number) = loop {
895 let request = proto::rpc::SyncAccountStorageMapsRequest {
896 block_range: Some(BlockRange {
897 block_from: pagination.current_block_from().as_u32(),
898 block_to: block_to.as_u32(),
899 }),
900 account_id: Some(account_id.into()),
901 };
902 let response = self
903 .call_with_retry(RpcEndpoint::SyncStorageMaps, |mut rpc_api| {
904 Box::pin(async move { rpc_api.sync_account_storage_maps(request).await })
905 })
906 .await?;
907 let page = StorageMapInfo::try_from(response.into_inner())?;
908
909 for (slot_name, entries) in page.map_entries {
910 map_entries
911 .entry(slot_name)
912 .or_default()
913 .as_map_mut()
914 .extend(entries.into_map());
915 }
916
917 match pagination.advance(page.block_number, page.chain_tip)? {
918 PaginationResult::Continue => {},
919 PaginationResult::Done {
920 chain_tip: final_chain_tip,
921 block_num: final_block_num,
922 } => break (final_chain_tip, final_block_num),
923 }
924 };
925
926 Ok(StorageMapInfo { chain_tip, block_number, map_entries })
927 }
928
929 async fn sync_account_vault(
930 &self,
931 block_from: BlockNumber,
932 block_to: BlockNumber,
933 account_id: AccountId,
934 ) -> Result<AccountVaultInfo, RpcError> {
935 let mut pagination = BlockPagination::new(block_from, block_to);
936 let mut vault_patch = AccountVaultPatch::default();
937
938 let (chain_tip, block_number) = loop {
939 let request = proto::rpc::SyncAccountVaultRequest {
940 block_range: Some(BlockRange {
941 block_from: pagination.current_block_from().as_u32(),
942 block_to: block_to.as_u32(),
943 }),
944 account_id: Some(account_id.into()),
945 };
946 let response = self
947 .call_with_retry(RpcEndpoint::SyncAccountVault, |mut rpc_api| {
948 Box::pin(async move { rpc_api.sync_account_vault(request).await })
949 })
950 .await?;
951 let page = AccountVaultInfo::try_from(response.into_inner())?;
952
953 vault_patch.merge(page.vault_patch);
954
955 match pagination.advance(page.block_number, page.chain_tip)? {
956 PaginationResult::Continue => {},
957 PaginationResult::Done {
958 chain_tip: final_chain_tip,
959 block_num: final_block_num,
960 } => break (final_chain_tip, final_block_num),
961 }
962 };
963
964 Ok(AccountVaultInfo { chain_tip, block_number, vault_patch })
965 }
966
967 async fn sync_transactions(
973 &self,
974 block_from: BlockNumber,
975 block_to: BlockNumber,
976 account_ids: Vec<AccountId>,
977 ) -> Result<Vec<TransactionRecord>, RpcError> {
978 if account_ids.is_empty() {
979 return Ok(Vec::new());
980 }
981
982 let limits = self.get_rpc_limits().await?;
983 let mut transactions: Vec<TransactionRecord> = Vec::new();
984
985 for chunk in account_ids.chunks(limits.account_ids_limit as usize) {
986 let proto_account_ids: Vec<_> = chunk.iter().map(|acc_id| (*acc_id).into()).collect();
987 let mut pagination = BlockPagination::new(block_from, block_to);
988
989 loop {
990 let request = proto::rpc::SyncTransactionsRequest {
991 block_range: Some(BlockRange {
992 block_from: pagination.current_block_from().as_u32(),
993 block_to: block_to.as_u32(),
994 }),
995 account_ids: proto_account_ids.clone(),
996 };
997
998 let response = self
999 .call_with_retry(RpcEndpoint::SyncTransactions, |mut rpc_api| {
1000 let request = request.clone();
1001 Box::pin(async move { rpc_api.sync_transactions(request).await })
1002 })
1003 .await?
1004 .into_inner();
1005
1006 let page = response.pagination_info.ok_or(RpcError::ExpectedDataMissing(
1007 "SyncTransactionsResponse.pagination_info".to_owned(),
1008 ))?;
1009 let page_chain_tip = BlockNumber::from(page.chain_tip);
1010 let page_block_to = BlockNumber::from(page.block_num);
1011
1012 for proto_tx in response.transactions {
1013 transactions.push(TransactionRecord::try_from(proto_tx)?);
1014 }
1015
1016 match pagination.advance(page_block_to, page_chain_tip)? {
1017 PaginationResult::Continue => {},
1018 PaginationResult::Done { .. } => break,
1019 }
1020 }
1021 }
1022
1023 Ok(transactions)
1024 }
1025
1026 async fn get_network_id(&self) -> Result<NetworkId, RpcError> {
1027 let endpoint: Endpoint =
1028 Endpoint::try_from(self.endpoint.as_str()).map_err(RpcError::InvalidNodeEndpoint)?;
1029 Ok(endpoint.to_network_id())
1030 }
1031
1032 async fn get_rpc_limits(&self) -> Result<RpcLimits, RpcError> {
1033 if let Some(limits) = *self.limits.read() {
1034 return Ok(limits);
1035 }
1036
1037 let response = self
1038 .call_with_retry(RpcEndpoint::GetLimits, |mut rpc_api| {
1039 Box::pin(async move { rpc_api.get_limits(proto::rpc::GetLimitsRequest {}).await })
1040 })
1041 .await?;
1042 let limits = RpcLimits::try_from(response.into_inner()).map_err(RpcError::from)?;
1043
1044 self.limits.write().replace(limits);
1046 Ok(limits)
1047 }
1048
1049 fn has_rpc_limits(&self) -> Option<RpcLimits> {
1050 *self.limits.read()
1051 }
1052
1053 async fn set_rpc_limits(&self, limits: RpcLimits) {
1054 self.limits.write().replace(limits);
1055 }
1056
1057 async fn get_status_unversioned(&self) -> Result<RpcStatusInfo, RpcError> {
1058 GrpcClient::get_status_unversioned(self).await
1059 }
1060
1061 async fn get_network_note_status(
1062 &self,
1063 note_id: NoteId,
1064 ) -> Result<NetworkNoteStatusInfo, RpcError> {
1065 let request = proto::rpc::GetNetworkNoteStatusRequest {
1066 note_id: Some(proto::note::NoteId::from(¬e_id)),
1067 };
1068
1069 let response = self
1070 .call_with_retry(RpcEndpoint::GetNetworkNoteStatus, |mut rpc_api| {
1071 let request = request.clone();
1072 Box::pin(async move { rpc_api.get_network_note_status(request).await })
1073 })
1074 .await?;
1075
1076 response.into_inner().try_into()
1077 }
1078}
1079
1080impl RpcError {
1084 pub fn from_grpc_error_with_context(
1085 endpoint: RpcEndpoint,
1086 status: Status,
1087 context: AcceptHeaderContext,
1088 ) -> Self {
1089 if let Some(accept_error) =
1090 AcceptHeaderError::try_from_message_with_context(status.message(), context)
1091 {
1092 return Self::AcceptHeaderError(accept_error);
1093 }
1094
1095 let error_kind = GrpcError::from(&status);
1096
1097 let endpoint_error = parse_node_error(&endpoint, status.details(), status.message())
1099 .or_else(|| parse_status_error(&endpoint, &error_kind, status.message()));
1100
1101 let source = Box::new(status) as Box<dyn Error + Send + Sync + 'static>;
1102
1103 Self::RequestError {
1104 endpoint,
1105 error_kind,
1106 endpoint_error,
1107 source: Some(source),
1108 }
1109 }
1110}
1111
1112impl From<&Status> for GrpcError {
1113 fn from(status: &Status) -> Self {
1114 GrpcError::from_code(status.code() as i32, Some(status.message().to_string()))
1115 }
1116}
1117
1118fn decode_block_response(
1127 response: proto::rpc::GetBlockByNumberResponse,
1128) -> Result<(SignedBlock, Option<ExecutionProof>), RpcError> {
1129 let block: SignedBlock = response
1132 .block
1133 .ok_or(RpcError::ExpectedDataMissing("GetBlockByNumberResponse.block".to_string()))?
1134 .decode_and_build_unchecked()?;
1135
1136 let proof = response.proof.map(ExecutionProof::try_from).transpose()?;
1139
1140 Ok((block, proof))
1141}
1142
1143#[cfg(test)]
1144mod tests {
1145 use std::boxed::Box;
1146 use std::vec;
1147
1148 use miden_protocol::Word;
1149 use miden_protocol::block::{BlockNumber, SignedBlock};
1150 use miden_testing::MockChain;
1151
1152 use super::{
1153 BlockPagination,
1154 DEFAULT_MAX_RESPONSE_SIZE_BYTES,
1155 GrpcClient,
1156 PaginationResult,
1157 decode_block_response,
1158 proto,
1159 };
1160 use crate::alloc::string::ToString;
1161 use crate::rpc::{Endpoint, NodeRpcClient, RpcError};
1162
1163 fn assert_send_sync<T: Send + Sync>() {}
1164
1165 fn genesis_block_messages()
1167 -> (proto::blockchain::SignedBlock, proto::primitives::ExecutionProof) {
1168 let chain = MockChain::new();
1169 let block = chain.proven_blocks().first().expect("the chain has a genesis block").clone();
1170 let (header, body, signatures, proof) = block.into_parts();
1171
1172 (SignedBlock::new_unchecked(header, body, signatures).into(), proof.into())
1173 }
1174
1175 #[test]
1176 fn decode_block_response_reads_a_requested_proof() {
1177 let (block, proof_message) = genesis_block_messages();
1178 let response = proto::rpc::GetBlockByNumberResponse {
1179 block: Some(block),
1180 proof: Some(proof_message.clone()),
1181 };
1182
1183 let (_block, proof) = decode_block_response(response).unwrap();
1184
1185 let decoded: proto::primitives::ExecutionProof =
1186 proof.expect("the response carries a proof").into();
1187 assert_eq!(decoded, proof_message);
1188 }
1189
1190 #[test]
1191 fn decode_block_response_omits_an_absent_proof() {
1192 let (block, _) = genesis_block_messages();
1193 let response = proto::rpc::GetBlockByNumberResponse { block: Some(block), proof: None };
1194
1195 let (_block, proof) = decode_block_response(response).unwrap();
1196
1197 assert!(proof.is_none());
1198 }
1199
1200 #[test]
1201 fn decode_block_response_rejects_malformed_proof_bytes() {
1202 let (block, _) = genesis_block_messages();
1203 let response = proto::rpc::GetBlockByNumberResponse {
1204 block: Some(block),
1205 proof: Some(proto::primitives::ExecutionProof { encoded: vec![0xff; 32] }),
1206 };
1207
1208 let res = decode_block_response(response);
1209
1210 assert!(matches!(res, Err(RpcError::DeserializationError(_))));
1211 }
1212
1213 #[test]
1214 fn decode_block_response_rejects_an_absent_block() {
1215 let response = proto::rpc::GetBlockByNumberResponse { block: None, proof: None };
1216
1217 let res = decode_block_response(response);
1218
1219 assert!(matches!(res, Err(RpcError::ExpectedDataMissing(_))));
1220 }
1221
1222 #[test]
1223 fn is_send_sync() {
1224 assert_send_sync::<GrpcClient>();
1225 assert_send_sync::<Box<dyn NodeRpcClient>>();
1226 }
1227
1228 #[test]
1229 fn block_pagination_errors_when_block_num_goes_backwards() {
1230 let mut pagination = BlockPagination::new(10_u32.into(), 20_u32.into());
1231
1232 let res = pagination.advance(9_u32.into(), 20_u32.into());
1233 assert!(matches!(res, Err(RpcError::PaginationError(_))));
1234 }
1235
1236 #[test]
1237 fn block_pagination_errors_when_block_num_passes_block_to() {
1238 let mut pagination = BlockPagination::new(10_u32.into(), 20_u32.into());
1239
1240 let res = pagination.advance(21_u32.into(), 100_u32.into());
1241 assert!(matches!(res, Err(RpcError::PaginationError(_))));
1242 }
1243
1244 #[test]
1245 fn block_pagination_errors_after_max_iterations() {
1246 let mut pagination = BlockPagination::new(0_u32.into(), 10_000_u32.into());
1247 let chain_tip: BlockNumber = 10_000_u32.into();
1248
1249 for _ in 0..BlockPagination::MAX_ITERATIONS {
1250 let current = pagination.current_block_from();
1251 let res = pagination
1252 .advance(current, chain_tip)
1253 .expect("expected pagination to continue within iteration limit");
1254 assert!(matches!(res, PaginationResult::Continue));
1255 }
1256
1257 let res = pagination.advance(pagination.current_block_from(), chain_tip);
1258 assert!(matches!(res, Err(RpcError::PaginationError(_))));
1259 }
1260
1261 #[test]
1262 fn block_pagination_stops_at_min_of_block_to_and_chain_tip() {
1263 let mut pagination = BlockPagination::new(0_u32.into(), 50_u32.into());
1265
1266 let res = pagination
1267 .advance(30_u32.into(), 30_u32.into())
1268 .expect("expected pagination to succeed");
1269
1270 assert!(matches!(
1271 res,
1272 PaginationResult::Done {
1273 chain_tip,
1274 block_num
1275 } if chain_tip.as_u32() == 30 && block_num.as_u32() == 30
1276 ));
1277 }
1278
1279 #[test]
1280 fn block_pagination_advances_cursor_by_one() {
1281 let mut pagination = BlockPagination::new(5_u32.into(), 100_u32.into());
1282
1283 let res = pagination
1284 .advance(5_u32.into(), 100_u32.into())
1285 .expect("expected pagination to succeed");
1286 assert!(matches!(res, PaginationResult::Continue));
1287 assert_eq!(pagination.current_block_from().as_u32(), 6);
1288 }
1289
1290 async fn dyn_trait_send_fut(client: Box<dyn NodeRpcClient>) {
1292 let res = client.get_block_header_by_number(None, false).await;
1294 assert!(res.is_ok());
1295 }
1296
1297 #[tokio::test]
1298 async fn future_is_send() {
1299 let endpoint = &Endpoint::devnet();
1300 let client = GrpcClient::new(endpoint, 10000);
1301 let client: Box<GrpcClient> = client.into();
1302 tokio::task::spawn(async move { dyn_trait_send_fut(client).await });
1303 }
1304
1305 #[tokio::test]
1306 async fn set_genesis_commitment_sets_the_commitment_when_its_not_already_set() {
1307 let endpoint = &Endpoint::devnet();
1308 let client = GrpcClient::new(endpoint, 10000);
1309
1310 assert!(client.genesis_commitment.read().is_none());
1311
1312 let commitment = Word::default();
1313 client.set_genesis_commitment(commitment).await.unwrap();
1314
1315 assert_eq!(client.genesis_commitment.read().unwrap(), commitment);
1316 }
1317
1318 #[tokio::test]
1319 async fn set_genesis_commitment_does_nothing_if_the_commitment_is_already_set() {
1320 let endpoint = &Endpoint::devnet();
1321 let client = GrpcClient::new(endpoint, 10000);
1322
1323 let initial_commitment = Word::default();
1324 client.set_genesis_commitment(initial_commitment).await.unwrap();
1325
1326 let new_commitment = Word::from([1u32, 2, 3, 4]);
1327 client.set_genesis_commitment(new_commitment).await.unwrap();
1328
1329 assert_eq!(client.genesis_commitment.read().unwrap(), initial_commitment);
1330 }
1331
1332 #[tokio::test]
1333 async fn set_genesis_commitment_updates_the_client_if_already_connected() {
1334 let endpoint = &Endpoint::devnet();
1335 let client = GrpcClient::new(endpoint, 10000);
1336
1337 client.connect().await.unwrap();
1339
1340 let commitment = Word::default();
1341 client.set_genesis_commitment(commitment).await.unwrap();
1342
1343 assert_eq!(client.genesis_commitment.read().unwrap(), commitment);
1344 assert!(client.client.read().as_ref().is_some());
1345 }
1346
1347 #[test]
1348 fn with_bearer_auth_stores_token() {
1349 let endpoint = &Endpoint::devnet();
1350 let client = GrpcClient::new(endpoint, 10000).with_bearer_auth("token-one".to_string());
1351
1352 assert_eq!(client.bearer_token.as_deref(), Some("token-one"));
1353 }
1354
1355 #[test]
1356 fn with_bearer_auth_overwrites_on_repeat_call() {
1357 let endpoint = &Endpoint::devnet();
1358 let client = GrpcClient::new(endpoint, 10000)
1359 .with_bearer_auth("token-one".to_string())
1360 .with_bearer_auth("token-two".to_string());
1361
1362 assert_eq!(client.bearer_token.as_deref(), Some("token-two"));
1364 }
1365
1366 #[tokio::test]
1367 async fn with_bearer_auth_surfaces_invalid_ascii_value_at_connect_time() {
1368 let endpoint = &Endpoint::devnet();
1372 let client = GrpcClient::new(endpoint, 10000).with_bearer_auth("bad\nvalue".to_string());
1373
1374 let err = client.connect().await.expect_err("expected invalid token to fail connect");
1375 assert!(
1376 matches!(err, RpcError::ConnectionError(_)),
1377 "expected ConnectionError, got {err:?}",
1378 );
1379 }
1380
1381 #[tokio::test]
1382 async fn with_bearer_auth_is_preserved_across_set_genesis_commitment() {
1383 let endpoint = &Endpoint::devnet();
1384 let client = GrpcClient::new(endpoint, 10000).with_bearer_auth("token".to_string());
1385 client.connect().await.unwrap();
1386
1387 client.set_genesis_commitment(Word::default()).await.unwrap();
1388
1389 assert_eq!(client.bearer_token.as_deref(), Some("token"));
1391 assert!(client.client.read().as_ref().is_some());
1392 }
1393
1394 #[test]
1395 fn with_max_decoding_message_size_overrides_default() {
1396 let endpoint = &Endpoint::devnet();
1397
1398 let default_client = GrpcClient::new(endpoint, 10_000);
1400 assert_eq!(default_client.max_decoding_message_size, DEFAULT_MAX_RESPONSE_SIZE_BYTES);
1401
1402 let custom =
1404 GrpcClient::new(endpoint, 10_000).with_max_decoding_message_size(8 * 1024 * 1024);
1405 assert_eq!(custom.max_decoding_message_size, 8 * 1024 * 1024);
1406 }
1407
1408 #[tokio::test]
1419 #[ignore = "requires network access to public testnet"]
1420 async fn with_bearer_auth_does_not_break_real_rpc_against_testnet() {
1421 let endpoint = &Endpoint::testnet();
1422 let client = GrpcClient::new(endpoint, 10_000).with_bearer_auth("smoke-test".to_string());
1423
1424 let status = client
1425 .get_status_unversioned()
1426 .await
1427 .expect("testnet status with caller auth header must succeed");
1428 assert!(!status.version.is_empty(), "status must include a server version");
1429 }
1430}