Skip to main content

snarkos_node_rest/
lib.rs

1// Copyright (c) 2019-2026 Provable Inc.
2// This file is part of the snarkOS library.
3
4// Licensed under the Apache License, Version 2.0 (the "License");
5// you may not use this file except in compliance with the License.
6// You may obtain a copy of the License at:
7
8// http://www.apache.org/licenses/LICENSE-2.0
9
10// Unless required by applicable law or agreed to in writing, software
11// distributed under the License is distributed on an "AS IS" BASIS,
12// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13// See the License for the specific language governing permissions and
14// limitations under the License.
15
16#![forbid(unsafe_code)]
17
18#[macro_use]
19extern crate tracing;
20
21mod helpers;
22// Imports custom `Path` type, to be used instead of `axum`'s.
23pub use helpers::*;
24
25mod history_compat;
26use history_compat::*;
27
28mod routes;
29
30mod version;
31
32use snarkos_node_cdn::CdnBlockSync;
33use snarkos_node_consensus::Consensus;
34use snarkos_node_router::{
35    Routing,
36    messages::{Message, UnconfirmedTransaction},
37};
38use snarkos_node_sync::BlockSync;
39use snarkvm::{
40    console::{program::ProgramID, types::Field},
41    ledger::narwhal::Data,
42    prelude::{Ledger, Network, VM, cfg_into_iter, store::ConsensusStorage},
43};
44
45use anyhow::{Context, Result};
46use axum::{
47    body::Body,
48    extract::{ConnectInfo, DefaultBodyLimit, Query, State},
49    http::{Method, Request, StatusCode, header::CONTENT_TYPE},
50    middleware,
51    response::Response,
52    routing::{get, post},
53};
54use axum_extra::response::ErasedJson;
55#[cfg(feature = "locktick")]
56use locktick::parking_lot::Mutex;
57use lru::LruCache;
58#[cfg(not(feature = "locktick"))]
59use parking_lot::Mutex;
60use std::{net::SocketAddr, num::NonZeroUsize, sync::Arc, time::Duration};
61use tokio::{net::TcpListener, sync::Semaphore, task::JoinHandle};
62use tower_governor::{GovernorLayer, governor::GovernorConfigBuilder};
63use tower_http::{
64    cors::{Any, CorsLayer},
65    trace::TraceLayer,
66};
67use tracing::Span;
68
69/// The default port used for the REST API
70pub const DEFAULT_REST_PORT: u16 = 3030;
71
72/// The API version prefixes.
73pub const API_VERSION_V1: &str = "v1";
74pub const API_VERSION_V2: &str = "v2";
75
76/// The capacity of the LRU holding recently requested blocks.
77const BLOCK_CACHE_SIZE: usize = 128;
78
79/// A REST API server for the ledger.
80#[derive(Clone)]
81pub struct Rest<N: Network, C: ConsensusStorage<N>, R: Routing<N>> {
82    /// CDN sync (only if node is using the CDN to sync).
83    cdn_sync: Option<Arc<CdnBlockSync>>,
84    /// The consensus module.
85    consensus: Option<Consensus<N>>,
86    /// The ledger.
87    ledger: Ledger<N, C>,
88    /// The node (routing).
89    routing: Arc<R>,
90    /// The server handles.
91    handles: Arc<Mutex<Vec<JoinHandle<()>>>>,
92    /// A reference to BlockSync,
93    block_sync: Arc<BlockSync<N>>,
94    /// The number of ongoing deploy transaction verifications via REST.
95    num_verifying_deploys: Arc<Semaphore>,
96    /// The number of ongoing execute transaction verifications via REST.
97    num_verifying_executions: Arc<Semaphore>,
98    /// The number of ongoing solution verifications via REST.
99    num_verifying_solutions: Arc<Semaphore>,
100    /// A cache containing recently requested blocks.
101    block_cache: Arc<Mutex<LruCache<N::BlockHash, ErasedJson>>>,
102    /// The upstream for the routes of the removed `history` feature, if `--history-compat-mode` is set.
103    history_compat: Option<Arc<HistoryCompat>>,
104}
105
106impl<N: Network, C: 'static + ConsensusStorage<N>, R: Routing<N>> Rest<N, C, R> {
107    /// Initializes a new instance of the server.
108    #[allow(clippy::too_many_arguments)]
109    pub async fn start(
110        rest_ip: SocketAddr,
111        rest_rps: u32,
112        history_api_url: Option<String>,
113        consensus: Option<Consensus<N>>,
114        ledger: Ledger<N, C>,
115        routing: Arc<R>,
116        cdn_sync: Option<Arc<CdnBlockSync>>,
117        block_sync: Arc<BlockSync<N>>,
118    ) -> Result<Self> {
119        // Initialize the history compatibility upstream, if requested.
120        let history_compat = match history_api_url {
121            Some(url) => Some(Arc::new(HistoryCompat::new(&url, N::SHORT_NAME)?)),
122            None => None,
123        };
124        // Initialize the server.
125        let mut server = Self {
126            consensus,
127            ledger,
128            routing,
129            cdn_sync,
130            block_sync,
131            handles: Default::default(),
132            num_verifying_deploys: Arc::new(Semaphore::new(VM::<N, C>::MAX_PARALLEL_DEPLOY_VERIFICATIONS)),
133            num_verifying_executions: Arc::new(Semaphore::new(VM::<N, C>::MAX_PARALLEL_EXECUTE_VERIFICATIONS)),
134            num_verifying_solutions: Arc::new(Semaphore::new(N::MAX_SOLUTIONS)),
135            block_cache: Arc::new(Mutex::new(LruCache::new(NonZeroUsize::new(BLOCK_CACHE_SIZE).unwrap()))),
136            history_compat,
137        };
138        // Spawn the server.
139        server.spawn_server(rest_ip, rest_rps).await?;
140        // Return the server.
141        Ok(server)
142    }
143}
144
145impl<N: Network, C: ConsensusStorage<N>, R: Routing<N>> Rest<N, C, R> {
146    /// Returns the ledger.
147    pub const fn ledger(&self) -> &Ledger<N, C> {
148        &self.ledger
149    }
150
151    /// Returns the handles.
152    pub const fn handles(&self) -> &Arc<Mutex<Vec<JoinHandle<()>>>> {
153        &self.handles
154    }
155
156    /// Shuts down the REST instance.
157    pub fn shut_down(&self) {
158        self.handles.lock().iter().for_each(|handle| handle.abort());
159    }
160}
161
162impl<N: Network, C: ConsensusStorage<N>, R: Routing<N>> Rest<N, C, R> {
163    fn build_routes(&self, rest_rps: u32) -> axum::Router {
164        let cors = CorsLayer::new()
165            .allow_origin(Any)
166            .allow_methods([Method::GET, Method::POST, Method::DELETE, Method::OPTIONS])
167            .allow_headers([CONTENT_TYPE]);
168
169        // Prepare the rate limiting setup.
170        let governor_config = Box::new(
171            GovernorConfigBuilder::default()
172                .per_nanosecond((1_000_000_000 / rest_rps) as u64)
173                .burst_size(rest_rps)
174                .error_handler(|error| {
175                    // Properly return a 429 Too Many Requests error
176                    let error_message = error.to_string();
177                    let mut response = Response::new(error_message.clone().into());
178                    *response.status_mut() = StatusCode::INTERNAL_SERVER_ERROR;
179                    if error_message.contains("Too Many Requests") {
180                        *response.status_mut() = StatusCode::TOO_MANY_REQUESTS;
181                    }
182                    response
183                })
184                .finish()
185                .expect("Couldn't set up rate limiting for the REST server!"),
186        );
187
188        // Build the JWT auth-protected endpoints. #[cfg] cannot appear inside a method chain, so we
189        // build this router as a named binding and conditionally extend it before applying the layer.
190        let auth_routes = axum::Router::new()
191            .route("/node/address", get(Self::get_node_address))
192            .route("/program/{id}/mapping/{name}", get(Self::get_mapping_values))
193            .route("/db_backup", post(Self::db_backup));
194
195        // Slipstream plugin management endpoints require auth.
196        #[cfg(feature = "slipstream-plugins")]
197        let auth_routes = auth_routes
198            .route("/slipstream/plugins", get(Self::slipstream_list_plugins).post(Self::slipstream_load_plugin))
199            .route(
200                "/slipstream/plugins/{name}",
201                // TODO: PUT (reload) is not yet implemented.
202                axum::routing::delete(Self::slipstream_unload_plugin),
203            );
204
205        let routes = axum::Router::new()
206            .merge(auth_routes.route_layer(middleware::from_fn(auth_middleware)))
207
208            // All endpoints declared after here are not protected
209
210             // Get ../consensus_version
211            .route("/consensus_version", get(Self::get_consensus_version))
212
213            // GET ../block/..
214            .route("/block/height/latest", get(Self::get_block_height_latest))
215            .route("/block/hash/latest", get(Self::get_block_hash_latest))
216            .route("/block/latest", get(Self::get_block_latest))
217            .route("/block/{height_or_hash}", get(Self::get_block))
218            // The path param here is actually only the height, but the name must match the route
219            // above, otherwise there'll be a conflict at runtime.
220            .route("/block/{height_or_hash}/header", get(Self::get_block_header))
221            .route("/block/{height_or_hash}/transactions", get(Self::get_block_transactions))
222
223            // GET and POST ../transaction/..
224            .route("/transaction/{id}", get(Self::get_transaction))
225            .route("/transaction/confirmed/{id}", get(Self::get_confirmed_transaction))
226            .route("/transaction/unconfirmed/{id}", get(Self::get_unconfirmed_transaction))
227            .route("/transaction/rejected/{id}/reason", get(Self::get_transaction_rejection_reason))
228            .route("/transaction/broadcast", post(Self::transaction_broadcast))
229
230            // GET and POST ../solution/..
231            .route("/solution/limits/{prover_address}", get(Self::get_solution_limits_for_prover))
232            .route("/solution/broadcast", post(Self::solution_broadcast))
233
234            // GET ../find/..
235            .route("/find/blockHash/{tx_id}", get(Self::find_block_hash))
236            .route("/find/blockHeight/{state_root}", get(Self::find_block_height_from_state_root))
237            .route("/find/transactionID/deployment/{program_id}", get(Self::find_latest_transaction_id_from_program_id))
238            .route("/find/transactionID/deployment/{program_id}/{edition}", get(Self::find_latest_transaction_id_from_program_id_and_edition))
239            .route("/find/transactionID/deployment/{program_id}/{edition}/original", get(Self::find_original_deployment_transaction_id))
240            .route("/find/transactionID/deployment/{program_id}/{edition}/{amendment}", get(Self::find_transaction_id_from_program_id_edition_and_amendment))
241            .route("/find/transactionID/{transition_id}", get(Self::find_transaction_id_from_transition_id))
242            .route("/find/transitionID/{input_or_output_id}", get(Self::find_transition_id))
243
244            // GET ../connections/p2p/.. (with ../peers/.. aliases)
245            .route("/peers/count", get(Self::get_peers_count))
246            .route("/peers/all", get(Self::get_peers_all))
247            .route("/peers/all/metrics", get(Self::get_peers_all_metrics))
248            .route("/connections/p2p/count", get(Self::get_peers_count))
249            .route("/connections/p2p/all", get(Self::get_peers_all))
250            .route("/connections/p2p/all/metrics", get(Self::get_peers_all_metrics))
251
252            // GET ../program/..
253            .route("/program/{id}", get(Self::get_program))
254            .route("/program/{id}/latest_edition", get(Self::get_latest_program_edition))
255            .route("/program/{id}/{edition}", get(Self::get_program_for_edition))
256            .route("/program/{id}/mappings", get(Self::get_mapping_names))
257            .route("/program/{id}/mapping/{name}/{key}", get(Self::get_mapping_value))
258            .route("/program/{id}/amendment_count", get(Self::get_program_amendment_count))
259            .route("/program/{id}/{edition}/amendment_count", get(Self::get_program_amendment_count_for_edition))
260
261            // GET ../sync/..
262            // Note: keeping ../sync_status for compatibility
263            .route("/sync_status", get(Self::get_sync_status))
264            .route("/sync/status", get(Self::get_sync_status))
265            .route("/sync/peers", get(Self::get_sync_peers))
266            .route("/sync/requests", get(Self::get_sync_requests_summary))
267            .route("/sync/requests/list", get(Self::get_sync_requests_list))
268
269            // GET misc endpoints.
270            .route("/version", get(Self::get_version))
271            .route("/blocks", get(Self::get_blocks))
272            .route("/height/{hash}", get(Self::get_height))
273            .route("/memoryPool/transmissions", get(Self::get_memory_pool_transmissions))
274            .route("/memoryPool/solutions", get(Self::get_memory_pool_solutions))
275            .route("/memoryPool/transactions", get(Self::get_memory_pool_transactions))
276            .route("/statePath/{commitment}", get(Self::get_state_path_for_commitment))
277            .route("/statePaths", get(Self::get_state_paths_for_commitments))
278            .route("/stateRoot/latest", get(Self::get_state_root_latest))
279            .route("/stateRoot/{height}", get(Self::get_state_root))
280            .route("/committee/latest", get(Self::get_committee_latest))
281            .route("/committee/{height}", get(Self::get_committee))
282            .route("/delegators/{validator}", get(Self::get_delegators_for_validator));
283
284        // If the node is a validator, enable the BFT connections endpoints.
285        let routes = match self.consensus {
286            Some(_) => routes
287                .route("/connections/bft/count", get(Self::get_bft_connections_count))
288                .route("/connections/bft/all", get(Self::get_bft_connections_all)),
289            None => routes,
290        };
291
292        // If the node is a validator and `telemetry` features is enabled, enable the additional endpoint.
293        #[cfg(feature = "metrics")]
294        let routes = match self.consensus {
295            Some(_) => routes.route("/validators/participation", get(Self::get_validator_participation_scores)),
296            None => routes,
297        };
298
299        // Register the view-at-latest-height endpoint (always available, no history required).
300        let routes = routes.route("/program/{id}/view/{function}", post(Self::evaluate_view_latest));
301
302        // In history compatibility mode, serve the routes of the removed `history` feature from the
303        // upstream historical API (see `history_compat`).
304        let routes = if self.history_compat.is_some() {
305            routes
306                .route("/program/{id}/mapping/{name}/{key}/history/{height}", get(Self::get_history_compat))
307                .route("/program/{id}/mapping/{name}/history/{height}", get(Self::get_history_batch_compat))
308                .route("/program/{id}/view/{function}/{height}", post(Self::evaluate_view_at_height_compat))
309                .route("/staking/rewards/{address}/{height}", get(Self::get_staking_reward_compat))
310        } else {
311            routes
312        };
313
314        // If the `history-staking-rewards` feature is enabled, enable the additional endpoint (unless
315        // compatibility mode already serves it).
316        #[cfg(feature = "history-staking-rewards")]
317        let routes = if self.history_compat.is_some() {
318            routes
319        } else {
320            routes.route("/staking/rewards/{address}/{height}", get(Self::get_staking_reward))
321        };
322
323        let trace_layer = TraceLayer::new_for_http()
324            .make_span_with(|request: &Request<_>| {
325                let addr = request
326                    .extensions()
327                    .get::<ConnectInfo<SocketAddr>>()
328                    .map(|ConnectInfo(addr)| addr.to_string())
329                    .unwrap_or_else(|| "unknown".to_string());
330
331                // Create a span that includes method, path, and our extracted IP
332                tracing::info_span!(
333                    "REST",
334                    method = %request.method(),
335                    uri = %request.uri().path(),
336                    addr = %addr,
337                )
338            })
339            .on_request(|_request: &Request<_>, _span: &Span| {
340                info!("Received a request");
341            })
342            .on_response(|_response: &Response<_>, latency: Duration, _span: &Span| {
343                info!("Finished request in {:?}", latency);
344            });
345
346        routes
347            // Pass in `Rest` to make things convenient.
348            .with_state(self.clone())
349            // Cap the request body size at 1.5MiB.
350            .layer(DefaultBodyLimit::max(2 * 768 * 1024))
351            .layer(GovernorLayer {
352                config: governor_config.into(),
353            })
354            // Enable CORS.
355            .layer(cors)
356            // Enable tower-http tracing.
357            .layer(trace_layer)
358    }
359
360    async fn spawn_server(&mut self, rest_ip: SocketAddr, rest_rps: u32) -> Result<()> {
361        // Log the REST rate limit per IP.
362        debug!("REST rate limit per IP - {rest_rps} RPS");
363
364        // Add the v1 API as default and under "/v1".
365        let default_router = axum::Router::new().nest(
366            &format!("/{}", N::SHORT_NAME),
367            self.build_routes(rest_rps).layer(middleware::map_response(v1_error_middleware)),
368        );
369        let v1_router = axum::Router::new().nest(
370            &format!("/{API_VERSION_V1}/{}", N::SHORT_NAME),
371            self.build_routes(rest_rps).layer(middleware::map_response(v1_error_middleware)),
372        );
373
374        // Add the v2 API under "/v2".
375        let v2_router =
376            axum::Router::new().nest(&format!("/{API_VERSION_V2}/{}", N::SHORT_NAME), self.build_routes(rest_rps));
377
378        // Combine all routes.
379        let router = default_router.merge(v1_router).merge(v2_router);
380
381        let rest_listener =
382            TcpListener::bind(rest_ip).await.with_context(|| "Failed to bind TCP port for REST endpoints")?;
383
384        let handle = tokio::spawn(async move {
385            axum::serve(rest_listener, router.into_make_service_with_connect_info::<SocketAddr>())
386                .await
387                .expect("couldn't start rest server");
388        });
389
390        self.handles.lock().push(handle);
391        Ok(())
392    }
393}
394
395/// Converts errors to the old style for the v1 API.
396/// The error code will always be 500 and the content a simple string.
397async fn v1_error_middleware(response: Response) -> Response {
398    // The status code used by all v1 errors
399    const V1_STATUS_CODE: StatusCode = StatusCode::INTERNAL_SERVER_ERROR;
400
401    if response.status().is_success() {
402        return response;
403    }
404
405    // Returns a opaque error instead of panicking.
406    let fallback = || {
407        let mut response = Response::new(Body::from("Failed to convert error"));
408        *response.status_mut() = V1_STATUS_CODE;
409        response
410    };
411
412    let Ok(bytes) = axum::body::to_bytes(response.into_body(), usize::MAX).await else {
413        return fallback();
414    };
415
416    // Deserialize REST error so we can convert it to a string
417    let Ok(json_err) = serde_json::from_slice::<SerializedRestError>(&bytes) else {
418        return fallback();
419    };
420
421    let mut message = json_err.message;
422    for next in json_err.chain.into_iter() {
423        message = format!("{message} — {next}");
424    }
425
426    let mut response = Response::new(Body::from(message));
427
428    *response.status_mut() = V1_STATUS_CODE;
429
430    response
431}
432
433/// Formats an ID into a truncated identifier (for logging purposes).
434pub fn fmt_id(id: impl ToString) -> String {
435    let id = id.to_string();
436    let mut formatted_id = id.chars().take(16).collect::<String>();
437    if id.chars().count() > 16 {
438        formatted_id.push_str("..");
439    }
440    formatted_id
441}
442
443#[cfg(test)]
444mod tests {
445    use super::*;
446    use anyhow::anyhow;
447    use axum::{
448        Router,
449        body::Body,
450        http::{Request, StatusCode},
451        middleware,
452        routing::get,
453    };
454    use tower::ServiceExt; // for `oneshot`
455
456    fn test_app() -> Router {
457        let build_routes = || {
458            Router::new()
459                .route("/not_found", get(|| async { Err::<(), RestError>(RestError::not_found(anyhow!("missing"))) }))
460                .route("/bad_request", get(|| async { Err::<(), RestError>(RestError::bad_request(anyhow!("bad"))) }))
461                .route(
462                    "/service_unavailable",
463                    get(|| async { Err::<(), RestError>(RestError::service_unavailable(anyhow!("gone"))) }),
464                )
465        };
466        let router_v1 = build_routes().route_layer(middleware::map_response(v1_error_middleware));
467        let router_v2 = Router::new().nest(&format!("/{API_VERSION_V2}"), build_routes());
468        router_v1.merge(router_v2)
469    }
470
471    #[tokio::test]
472    async fn v1_routes_force_internal_server_error() {
473        let app = test_app();
474
475        let res = app.clone().oneshot(Request::builder().uri("/not_found").body(Body::empty()).unwrap()).await.unwrap();
476        assert_eq!(res.status(), StatusCode::INTERNAL_SERVER_ERROR);
477
478        let res =
479            app.clone().oneshot(Request::builder().uri("/bad_request").body(Body::empty()).unwrap()).await.unwrap();
480        assert_eq!(res.status(), StatusCode::INTERNAL_SERVER_ERROR);
481
482        let res =
483            app.oneshot(Request::builder().uri("/service_unavailable").body(Body::empty()).unwrap()).await.unwrap();
484        assert_eq!(res.status(), StatusCode::INTERNAL_SERVER_ERROR);
485    }
486
487    #[tokio::test]
488    async fn v2_routes_return_specific_errors() {
489        let app = test_app();
490
491        let res =
492            app.clone().oneshot(Request::builder().uri("/v2/not_found").body(Body::empty()).unwrap()).await.unwrap();
493        assert_eq!(res.status(), StatusCode::NOT_FOUND);
494
495        let res =
496            app.clone().oneshot(Request::builder().uri("/v2/bad_request").body(Body::empty()).unwrap()).await.unwrap();
497        assert_eq!(res.status(), StatusCode::BAD_REQUEST);
498
499        let res =
500            app.oneshot(Request::builder().uri("/v2/service_unavailable").body(Body::empty()).unwrap()).await.unwrap();
501        assert_eq!(res.status(), StatusCode::SERVICE_UNAVAILABLE);
502    }
503}
504
505#[cfg(test)]
506mod route_tests {
507    use super::*;
508    use snarkos_node_bft_ledger_service::MockLedgerService;
509    use snarkos_node_network::ConnectionMode;
510    use snarkos_node_router::test_helpers::{TestRouter, client, sample_genesis_block};
511    use snarkvm::{
512        ledger::{committee::test_helpers::sample_committee, store::helpers::memory::ConsensusMemory},
513        prelude::MainnetV0,
514        utilities::TestRng,
515    };
516
517    use aleo_std::StorageMode;
518    use axum::body::to_bytes;
519    use tower::ServiceExt; // for `oneshot`
520
521    type CurrentNetwork = MainnetV0;
522    type CurrentRest = Rest<CurrentNetwork, ConsensusMemory<CurrentNetwork>, TestRouter<CurrentNetwork>>;
523
524    /// The rate limit given to the router under test. The governor layer is applied by
525    /// `build_routes`, so this is set high enough that a test making several requests in quick
526    /// succession is never the thing that trips it.
527    const TEST_RPS: u32 = 1_000;
528
529    /// Builds a `Rest` over an in-memory ledger containing only the genesis block.
530    ///
531    /// This constructs the struct directly rather than calling `Rest::start`, which would bind a
532    /// port and spawn a server. None of the routes exercised here touch `consensus`, `cdn_sync`,
533    /// `routing` or `block_sync`; those fields exist only to satisfy the type.
534    async fn sample_rest() -> CurrentRest {
535        let rng = &mut TestRng::default();
536
537        // `Ledger::load` reaches snarkVM's sequential-operation thread and blocks on the reply,
538        // which panics if called from an async context. Production always drives these from a
539        // blocking task, so do the same here.
540        let ledger = tokio::task::spawn_blocking(|| {
541            Ledger::<CurrentNetwork, ConsensusMemory<CurrentNetwork>>::load(
542                sample_genesis_block::<CurrentNetwork>(),
543                StorageMode::new_test(None),
544            )
545        })
546        .await
547        .expect("the ledger task panicked")
548        .expect("couldn't load the test ledger");
549
550        let ledger_service = Arc::new(MockLedgerService::new(sample_committee(rng)));
551
552        Rest {
553            cdn_sync: None,
554            consensus: None,
555            ledger,
556            routing: Arc::new(client(0, 10, rng).await),
557            handles: Default::default(),
558            block_sync: Arc::new(BlockSync::new(ledger_service, ConnectionMode::Router)),
559            num_verifying_deploys: Arc::new(Semaphore::new(1)),
560            num_verifying_executions: Arc::new(Semaphore::new(1)),
561            num_verifying_solutions: Arc::new(Semaphore::new(1)),
562            block_cache: Arc::new(Mutex::new(LruCache::new(NonZeroUsize::new(BLOCK_CACHE_SIZE).unwrap()))),
563            history_compat: None,
564        }
565    }
566
567    /// Issues a GET request against the routes, without the network prefix that `spawn_server`
568    /// nests them under.
569    async fn get(rest: &CurrentRest, uri: &str) -> (StatusCode, String) {
570        request(rest, Method::GET, uri).await
571    }
572
573    /// Issues a request against the routes, without the network prefix that `spawn_server` nests
574    /// them under.
575    ///
576    /// The governor layer keys on the peer IP taken from `ConnectInfo`, which a request built by
577    /// hand does not carry, so this attaches one; without it every request fails the rate limiter's
578    /// key extractor rather than reaching a handler.
579    async fn request(rest: &CurrentRest, method: Method, uri: &str) -> (StatusCode, String) {
580        let mut request = Request::builder().method(method).uri(uri).body(Body::empty()).unwrap();
581        request.extensions_mut().insert(ConnectInfo(SocketAddr::from(([127, 0, 0, 1], 4130))));
582
583        let response = rest.build_routes(TEST_RPS).oneshot(request).await.unwrap();
584        let status = response.status();
585        let body = to_bytes(response.into_body(), usize::MAX).await.unwrap();
586
587        (status, String::from_utf8(body.to_vec()).unwrap())
588    }
589
590    #[tokio::test]
591    async fn latest_block_hash_is_the_genesis_hash() {
592        let rest = sample_rest().await;
593
594        // The test ledger holds only the genesis block.
595        let (status, body) = get(&rest, "/block/hash/latest").await;
596        assert_eq!(status, StatusCode::OK);
597
598        let hash: <CurrentNetwork as Network>::BlockHash = serde_json::from_str(&body).unwrap();
599        assert_eq!(hash, sample_genesis_block::<CurrentNetwork>().hash());
600    }
601
602    /// The routes of the removed `history` feature, in history compatibility mode.
603    mod history_compat {
604        use super::*;
605        use crate::history_compat::fixtures;
606
607        /// A stand-in for the upstream historical API: serves the block-1,000,000 fixtures at every
608        /// height for the mappings it has, and for `withdraw` the 500 that the real upstream answers
609        /// for a height it has no snapshot of.
610        async fn spawn_upstream() -> String {
611            async fn snapshot(Path((_height, mapping)): Path<(u32, String)>) -> (StatusCode, &'static str) {
612                match mapping.as_str() {
613                    "unbonding" => (StatusCode::OK, fixtures::UNBONDING),
614                    "bonded" => (StatusCode::OK, fixtures::BONDED),
615                    "stakingrewards" => (StatusCode::OK, fixtures::STAKING_REWARDS),
616                    "metadata" => (StatusCode::OK, fixtures::METADATA),
617                    _ => (StatusCode::INTERNAL_SERVER_ERROR, fixtures::MISSING),
618                }
619            }
620            let app =
621                axum::Router::new().route("/mainnet/block/{height}/history/{mapping}", axum::routing::get(snapshot));
622            let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
623            let address = listener.local_addr().unwrap();
624            tokio::spawn(async move { axum::serve(listener, app).await.unwrap() });
625            format!("http://{address}")
626        }
627
628        /// A `Rest` in history compatibility mode against the stub upstream.
629        async fn sample_compat_rest() -> CurrentRest {
630            let mut rest = sample_rest().await;
631            rest.history_compat = Some(Arc::new(HistoryCompat::new(&spawn_upstream().await, "mainnet").unwrap()));
632            rest
633        }
634
635        const UNBONDING_STAKER: &str = "aleo1sdjqhlcm9qltpu74ek0vxewt52zsdmn6swmpjn6m0tp9xf57dvpq740r8j";
636        const BONDED_STAKER: &str = "aleo1qy4qufq03wcph05fdf5aj09ez67vcmmlrzqf0zza352qwaq43gyqt3wdf6";
637        const VALIDATOR: &str = "aleo1vfukg8ky2mhfprw63s0k0hl4vvd8573s6fkn8cv9y0ca6q27eq8qwdnxls";
638
639        #[tokio::test]
640        async fn history_serves_the_value_from_the_upstream_snapshot() {
641            let rest = sample_compat_rest().await;
642            // The test ledger is at height 0, so that is the one height the upstream is asked for.
643            let (status, body) =
644                get(&rest, &format!("/program/credits.aleo/mapping/unbonding/{UNBONDING_STAKER}/history/0")).await;
645            assert_eq!(status, StatusCode::OK, "{body}");
646            // The body is the value's plaintext string, as the removed feature returned it.
647            let value: Option<String> = serde_json::from_str(&body).unwrap();
648            assert_eq!(value.as_deref(), Some("{\n  microcredits: 10113730488u64,\n  height: 621255u32\n}"));
649        }
650
651        #[tokio::test]
652        async fn history_answers_null_for_an_absent_key() {
653            let rest = sample_compat_rest().await;
654            // A key absent from the snapshot: e.g. an unbond that was claimed.
655            let (status, body) =
656                get(&rest, &format!("/program/credits.aleo/mapping/unbonding/{BONDED_STAKER}/history/0")).await;
657            assert_eq!(status, StatusCode::OK, "{body}");
658            assert_eq!(body.trim(), "null");
659        }
660
661        #[tokio::test]
662        async fn history_is_served_regardless_of_the_node_height() {
663            // The test ledger is at height 0; the upstream is the source of truth, so a height far
664            // above the node's is answered from it all the same.
665            let rest = sample_compat_rest().await;
666            let (status, body) =
667                get(&rest, &format!("/program/credits.aleo/mapping/unbonding/{UNBONDING_STAKER}/history/1000000"))
668                    .await;
669            assert_eq!(status, StatusCode::OK, "{body}");
670            let value: Option<String> = serde_json::from_str(&body).unwrap();
671            assert!(value.is_some());
672            let (status, body) = get(&rest, &format!("/staking/rewards/{BONDED_STAKER}/1000000")).await;
673            assert_eq!(status, StatusCode::OK, "{body}");
674            assert_ne!(body.trim(), "null");
675        }
676
677        #[tokio::test]
678        async fn history_batch_serves_every_key_from_one_snapshot() {
679            let rest = sample_compat_rest().await;
680            let (status, body) = get(
681                &rest,
682                &format!("/program/credits.aleo/mapping/unbonding/history/0?keys={UNBONDING_STAKER},{BONDED_STAKER}"),
683            )
684            .await;
685            assert_eq!(status, StatusCode::OK, "{body}");
686            let values: Vec<serde_json::Value> = serde_json::from_str(&body).unwrap();
687            assert_eq!(values.len(), 2);
688            assert_eq!(values[0]["key"], UNBONDING_STAKER);
689            assert_eq!(values[0]["value"], "{\n  microcredits: 10113730488u64,\n  height: 621255u32\n}");
690            assert_eq!(values[1]["key"], BONDED_STAKER);
691            assert_eq!(values[1]["value"], serde_json::Value::Null);
692        }
693
694        #[tokio::test]
695        async fn history_rejects_what_the_upstream_does_not_record() {
696            let rest = sample_compat_rest().await;
697            // A `credits.aleo` mapping the upstream has no snapshot of.
698            let (status, body) =
699                get(&rest, &format!("/program/credits.aleo/mapping/committee/{VALIDATOR}/history/0")).await;
700            assert_eq!(status, StatusCode::NOT_FOUND, "{body}");
701            assert!(body.contains("credits.aleo/committee"), "{body}");
702            // Another program.
703            let (status, body) = get(&rest, "/program/other.aleo/mapping/bonded/1field/history/0").await;
704            assert_eq!(status, StatusCode::NOT_FOUND, "{body}");
705            assert!(body.contains("other.aleo/bonded"), "{body}");
706            // A supported mapping for which the upstream has no snapshot at that height (which it
707            // reports as a 500, not a 404).
708            let (status, body) =
709                get(&rest, &format!("/program/credits.aleo/mapping/withdraw/{VALIDATOR}/history/0")).await;
710            assert_eq!(status, StatusCode::NOT_FOUND, "{body}");
711            assert!(body.contains("No snapshot of 'withdraw'"), "{body}");
712            // A view at a past height.
713            let (status, body) = request(&rest, Method::POST, "/program/credits.aleo/view/anything/0").await;
714            assert_eq!(status, StatusCode::NOT_FOUND, "{body}");
715            assert!(body.contains("latest height"), "{body}");
716        }
717
718        #[tokio::test]
719        async fn staking_reward_joins_the_rewards_and_bonded_snapshots() {
720            let rest = sample_compat_rest().await;
721            let (status, body) = get(&rest, &format!("/staking/rewards/{BONDED_STAKER}/0")).await;
722            assert_eq!(status, StatusCode::OK, "{body}");
723            // `[validator, reward, new_stake]`, as the `history-staking-rewards` feature returned it.
724            let reward: (String, u64, u64) = serde_json::from_str(&body).unwrap();
725            assert_eq!(reward, (VALIDATOR.to_string(), 6477, 141_347_021_440));
726            // A staker with no reward at that height.
727            let (status, body) = get(&rest, &format!("/staking/rewards/{UNBONDING_STAKER}/0")).await;
728            assert_eq!(status, StatusCode::OK, "{body}");
729            assert_eq!(body.trim(), "null");
730        }
731
732        #[tokio::test]
733        async fn history_routes_are_absent_without_compatibility_mode() {
734            let rest = sample_rest().await;
735            let (status, _) =
736                get(&rest, &format!("/program/credits.aleo/mapping/unbonding/{UNBONDING_STAKER}/history/0")).await;
737            assert_eq!(status, StatusCode::NOT_FOUND);
738        }
739    }
740}