Skip to main content

miden_client/rpc/tonic_client/
mod.rs

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
68/// Tracks the pagination state for block-driven endpoints.
69struct 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    /// Maximum number of pagination iterations for a single request.
85    ///
86    /// Protects against nodes returning inconsistent pagination data that could otherwise trigger
87    /// an infinite loop.
88    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        // The node must not answer with a page that ends past the requested window. The cursor is
125        // trusted downstream, so an out-of-window cursor is rejected here.
126        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
145// GRPC CLIENT
146// ================================================================================================
147
148/// Default maximum size (in bytes) of a decoded gRPC response the client will accept: 15% above
149/// tonic's built-in 4 MiB receive limit. See [`GrpcClient::with_max_decoding_message_size`].
150const DEFAULT_MAX_RESPONSE_SIZE_BYTES: usize = 4 * 1024 * 1024 * 115 / 100;
151
152/// Client for the Node RPC API using gRPC.
153///
154/// If the `tonic` feature is enabled, this client will use a `tonic::transport::Channel` to
155/// communicate with the node. In this case the connection will be established lazily when the first
156/// request is made. If the `web-tonic` feature is enabled, this client will use a
157/// `tonic_web_wasm_client::Client` to communicate with the node.
158///
159/// In both cases, the [`GrpcClient`] depends on the types inside the `generated` module, which are
160/// generated by the build script and also depend on the target architecture.
161pub struct GrpcClient {
162    /// The underlying gRPC client, lazily initialized on first request.
163    client: RwLock<Option<ApiClient>>,
164    /// The node endpoint URL to connect to.
165    endpoint: String,
166    /// Request timeout in milliseconds.
167    timeout_ms: u64,
168    /// The genesis block commitment, used for request validation by the node.
169    genesis_commitment: RwLock<Option<Word>>,
170    /// Cached RPC limits fetched from the node.
171    limits: RwLock<Option<RpcLimits>>,
172    /// Maximum number of retry attempts for rate-limited or transiently unavailable requests.
173    max_retries: u32,
174    /// Fallback retry interval in milliseconds when no `retry-after` header is present.
175    retry_interval_ms: u64,
176    /// Optional bearer token injected as `authorization: Bearer <token>` on every outbound gRPC
177    /// call, alongside the standard `accept` header. Used when talking to an authenticating gateway
178    /// in front of the node.
179    bearer_token: Option<String>,
180    /// Maximum size (in bytes) of a decoded gRPC response the client will accept. Defaults to
181    /// [`DEFAULT_MAX_RESPONSE_SIZE_BYTES`].
182    max_decoding_message_size: usize,
183}
184
185impl GrpcClient {
186    /// Returns a new instance of [`GrpcClient`] that'll do calls to the provided [`Endpoint`] with
187    /// the given timeout in milliseconds.
188    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    /// Sets the maximum number of retry attempts for rate-limited or transiently unavailable
203    /// requests. Defaults to `4`.
204    #[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    /// Sets the fallback retry interval in milliseconds, used when the server does not provide a
211    /// `retry-after` header. Defaults to `100` ms.
212    #[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    /// Sets the maximum size (in bytes) of a decoded gRPC response the client will accept.
219    ///
220    /// Defaults to 15% above [tonic's built-in 4 MiB receive limit][tonic-decode], leaving headroom
221    /// for responses that land slightly over 4 MiB.
222    ///
223    /// [tonic-decode]: https://github.com/hyperium/tonic/blob/6cb6056b5a748bc5a29bd48f4602dbc4e552bb7d/tonic/src/codec/decode.rs#L192-L218
224    #[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    /// Attaches a `authorization: Bearer <token>` header to every outbound gRPC call made by this
231    /// client, alongside the standard `accept` header.
232    ///
233    /// Intended for connecting to a Miden node through an authenticating gateway (e.g.
234    /// `miden-testnet.eu-central-8.gateway.fm`) that rate-limits unauthenticated traffic. Without
235    /// an auth mechanism on the client side, callers would have no way to supply the token the
236    /// gateway requires.
237    ///
238    /// Calling this method twice overwrites the earlier token.
239    ///
240    /// Validation of the token against [`AsciiMetadataValue`](tonic::metadata::AsciiMetadataValue)
241    /// is deferred to connection time (printable ASCII plus tab only — `HeaderValue::from_str`
242    /// semantics): invalid tokens surface as
243    /// [`RpcError::ConnectionError`](crate::rpc::RpcError::ConnectionError) on the first request,
244    /// so CR/LF header-injection attempts are rejected.
245    ///
246    /// # Example
247    ///
248    /// ```no_run
249    /// # use miden_client::rpc::{Endpoint, GrpcClient};
250    /// let endpoint = Endpoint::new("https".into(), "node.example".into(), Some(443));
251    /// let client = GrpcClient::new(&endpoint, 10_000).with_bearer_auth("<api-key>".into());
252    /// ```
253    #[must_use]
254    pub fn with_bearer_auth(mut self, token: String) -> Self {
255        self.bearer_token = Some(token);
256        self
257    }
258
259    /// Takes care of establishing the RPC connection if not connected yet. It ensures that the
260    /// `rpc_api` field is initialized and returns a write guard to it.
261    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    /// Connects to the Miden node, setting the client API with the provided URL, timeout and
270    /// genesis commitment.
271    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    /// Executes an RPC call and automatically retries transient failures.
301    ///
302    /// The provided closure is invoked with a freshly connected [`ApiClient`] on each attempt.
303    /// Retries are delegated to [`retry::RetryState`], which handles gRPC
304    /// [`tonic::Code::ResourceExhausted`] responses on any endpoint and
305    /// [`tonic::Code::Unavailable`] only where repeating the call is safe (see
306    /// [`RpcEndpoint::is_idempotent`]), including honoring cooldown delays when the node provides
307    /// them.
308    ///
309    /// Returns the first successful gRPC response. If the call keeps failing after retries are
310    /// exhausted, or if the error is not retryable, this returns the corresponding [`RpcError`] for
311    /// the provided [`RpcEndpoint`].
312    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    /// Fetches RPC status without injecting an Accept header.
332    ///
333    /// This instantiates a separate API client without the Accept interceptor, so it does not reuse
334    /// the primary gRPC client. Any caller-supplied [`with_bearer_auth`](Self::with_bearer_auth)
335    /// token is still forwarded so gateway authentication keeps working.
336    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    /// Sets the genesis commitment for the client. If the client is already connected, it will be
357    /// updated to use the new commitment on subsequent requests. If the client is not connected,
358    /// the commitment will be stored and used when the client connects. If the genesis commitment
359    /// is already set, this method does nothing.
360    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        // Check if already set before doing anything else
366        if self.genesis_commitment.read().is_some() {
367            // Genesis commitment is already set, ignoring the new value.
368            return Ok(());
369        }
370
371        // Store the commitment for future connections
372        self.genesis_commitment.write().replace(commitment);
373
374        // If a client is already connected, update it to use the new genesis commitment. If not
375        // connected, the commitment will be used when connect() is called.
376        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        // An undecodable attestation is skipped rather than failing the whole response, so that one
404        // junk entry served by the relaying operator cannot hide a valid attestation behind it.
405        // Verification requires one that both decodes and verifies, so dropping the rest is safe.
406        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        // A negative scheme is a malformed response, not a scheme this client happens to not
427        // support, so it is rejected here rather than aliased onto a valid identifier.
428        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    /// Sends a `GetAccount` request to the Miden node, and extracts the [`AccountProof`] from the
600    /// response, as well as the block number that it was retrieved for.
601    ///
602    /// # Errors
603    ///
604    /// This function will return an error if:
605    ///
606    /// - The requested Account isn't returned by the node.
607    /// - There was an error sending the request to the node.
608    /// - The answer had a `None` for one of the expected fields.
609    /// - There is an error during storage deserialization.
610    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        // We need the requested slots to interpret the node response.
624        let requirements = match storage.clone() {
625            StorageMapFetch::Slots(reqs) => reqs,
626            StorageMapFetch::Skip | StorageMapFetch::All => AccountStorageRequirements::default(),
627        };
628
629        // Only request details for accounts with public state (Public or Network), passing the
630        // known code commitment so the node can skip re-sending code we already hold.
631        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        // For accounts with public state, details should be present when requested
672        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        // The invitation code is a secret. Keep it out of logs and out of error messages.
695        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    /// Sends one or more `SyncNoteRequest`s to the node and merges the responses into a list of
722    /// [`SyncNotesBlock`]s.
723    ///
724    /// Chunks `note_tags` by [`RpcLimits::note_tags_limit`] and paginates each chunk across the
725    /// requested block range.
726    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        // Merge blocks across tag-chunks: a single block can hold notes whose tags fall into
740        // different chunks, so the same block can appear in multiple chunks' responses.
741        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        // If the prefixes are too many, we need to chunk them into smaller groups to avoid
802        // violating the RPC limit.
803        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        // The node returns an empty payload when it has no script registered for the root.
877        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    /// Sends one or more `SyncTransactions` requests to the node and concatenates the responses
968    /// into a flat list of [`TransactionRecord`]s.
969    ///
970    /// Chunks `account_ids` by [`RpcLimits::account_ids_limit`] and paginates each chunk across the
971    /// requested block range.
972    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        // Cache fetched values
1045        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(&note_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
1080// ERRORS
1081// ================================================================================================
1082
1083impl 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        // Parse the application-level error from the status details
1098        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
1118// HELPERS
1119// ================================================================================================
1120
1121/// Decodes the response of `get_block_by_number`.
1122///
1123/// The response carries the signed block and its proof in separate fields, so the block bytes
1124/// decode as a [`SignedBlock`] and never as a `ProvenBlock`. The node omits the proof when it is
1125/// not requested, and also when the block is not proven yet, so an absent proof is not an error.
1126fn decode_block_response(
1127    response: proto::rpc::GetBlockByNumberResponse,
1128) -> Result<(SignedBlock, Option<ExecutionProof>), RpcError> {
1129    // The response carries the block and its proof in separate fields, so the block message holds a
1130    // signed block and never a proven one.
1131    let block: SignedBlock = response
1132        .block
1133        .ok_or(RpcError::ExpectedDataMissing("GetBlockByNumberResponse.block".to_string()))?
1134        .decode_and_build_unchecked()?;
1135
1136    // The node omits the proof when it is not requested, and also when the block is not proven yet,
1137    // so an absent proof is not an error.
1138    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    /// Returns the signed block and proof messages of the mock chain's genesis block.
1166    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        // block_to is beyond chain tip, so target should be chain_tip.
1264        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    // Function that returns a `Send` future from a dynamic trait that must be `Sync`.
1291    async fn dyn_trait_send_fut(client: Box<dyn NodeRpcClient>) {
1292        // This won't compile if `get_block_header_by_number` doesn't return a `Send+Sync` future.
1293        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        // "Connect" the client
1338        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        // Second call replaces the first.
1363        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        // Tokens containing control characters are rejected by `AsciiMetadataValue`. The fluent
1369        // builder defers the check to connection time, so the error must surface as a
1370        // `ConnectionError` on the first request — preventing CR/LF header-injection.
1371        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        // Rebuilding the interceptor after a genesis update must not drop the caller token.
1390        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        // A fresh client uses the default decode ceiling.
1399        let default_client = GrpcClient::new(endpoint, 10_000);
1400        assert_eq!(default_client.max_decoding_message_size, DEFAULT_MAX_RESPONSE_SIZE_BYTES);
1401
1402        // The knob overrides it for callers that hit responses above the default.
1403        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    /// Real-network smoke test: hitting the public testnet with a caller-supplied bearer token must
1409    /// return a real [`RpcStatusInfo`], proving the header is a valid
1410    /// [`AsciiMetadataValue`](tonic::metadata::AsciiMetadataValue) on the wire and that an
1411    /// unauthenticated node ignores it cleanly.
1412    ///
1413    /// `#[ignore]`d by default so offline CI doesn't fail; run with `cargo test -- --ignored
1414    /// with_bearer_auth_does_not_break_real_rpc_against_testnet` when validating against the real
1415    /// network. The interceptor-level test
1416    /// (`api_client::tests::interceptor_injects_bearer_token_onto_request`) already proves the
1417    /// header reaches outbound request metadata without needing the network.
1418    #[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}