Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
38 changes: 36 additions & 2 deletions lean_client/http_api/src/handlers.rs
Original file line number Diff line number Diff line change
@@ -1,23 +1,26 @@
use std::sync::Arc;
use std::{collections::HashMap, sync::Arc};

use axum::{
Json,
extract::{FromRef, State},
http::StatusCode,
response::{IntoResponse, Response},
};
use containers::SignedBlock;
use fork_choice::store::Store;
use parking_lot::RwLock;
use serde_json::{Value, json};
use ssz::SszWrite;
use ssz::{H256, SszWrite};

use crate::aggregator_controller::SharedController;

pub type SharedStore = Arc<RwLock<Store>>;
pub type SharedSignedBlocks = Arc<RwLock<HashMap<H256, SignedBlock>>>;

#[derive(Clone)]
pub struct AppState {
pub store: SharedStore,
pub signed_blocks: SharedSignedBlocks,
pub controller: SharedController,
}

Expand All @@ -27,6 +30,12 @@ impl FromRef<AppState> for SharedStore {
}
}

impl FromRef<AppState> for SharedSignedBlocks {
fn from_ref(app_state: &AppState) -> Self {
app_state.signed_blocks.clone()
}
}

impl FromRef<AppState> for SharedController {
fn from_ref(app_state: &AppState) -> Self {
app_state.controller.clone()
Expand Down Expand Up @@ -62,6 +71,31 @@ pub async fn states_finalized(State(store): State<SharedStore>) -> Result<Respon
.into_response())
}

pub async fn blocks_finalized(
State(store): State<SharedStore>,
State(signed_blocks): State<SharedSignedBlocks>,
) -> Result<Response, StatusCode> {
let store = store.read();
let signed_blocks = signed_blocks.read();

let finalized_root = store.latest_finalized.root;

let block = signed_blocks
.get(&finalized_root)
.ok_or(StatusCode::NOT_FOUND)?;

let ssz_bytes = block
.to_ssz()
.map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?;

Ok((
StatusCode::OK,
[(axum::http::header::CONTENT_TYPE, "application/octet-stream")],
ssz_bytes,
)
.into_response())
}

pub async fn checkpoints_justified(State(store): State<SharedStore>) -> impl IntoResponse {
let store = store.read();

Expand Down
2 changes: 1 addition & 1 deletion lean_client/http_api/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@ mod test_driver;

pub use aggregator_controller::AggregatorController;
pub use config::HttpServerConfig;
pub use handlers::SharedStore;
pub use handlers::{SharedSignedBlocks, SharedStore};
pub use routing::normal_routes;
pub use server::{run_server, run_test_driver_server};
pub use test_driver::{TestDriverState, test_driver_routes};
10 changes: 8 additions & 2 deletions lean_client/http_api/src/routing.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,19 +4,25 @@ use crate::{
aggregator_controller::SharedController,
aggregator_handlers,
config::HttpServerConfig,
handlers::{self, AppState, SharedStore},
handlers::{self, AppState, SharedSignedBlocks, SharedStore},
};

pub fn normal_routes(
_config: &HttpServerConfig,
store: SharedStore,
signed_blocks: SharedSignedBlocks,
controller: SharedController,
) -> Router {
let app_state = AppState { store, controller };
let app_state = AppState {
store,
signed_blocks,
controller,
};

Router::new()
.route("/lean/v0/health", get(handlers::health))
.route("/lean/v0/states/finalized", get(handlers::states_finalized))
.route("/lean/v0/blocks/finalized", get(handlers::blocks_finalized))
.route(
"/lean/v0/checkpoints/justified",
get(handlers::checkpoints_justified),
Expand Down
8 changes: 5 additions & 3 deletions lean_client/http_api/src/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,17 +7,18 @@ use tracing::info;
use crate::{
aggregator_controller::SharedController,
config::HttpServerConfig,
handlers::SharedStore,
handlers::{SharedSignedBlocks, SharedStore},
routing::normal_routes,
test_driver::{TestDriverState, test_driver_routes},
};

pub async fn run_server(
config: HttpServerConfig,
store: SharedStore,
signed_blocks: SharedSignedBlocks,
aggregator_controller: SharedController,
) -> Result<()> {
let router = normal_routes(&config, store, aggregator_controller);
let router = normal_routes(&config, store, signed_blocks, aggregator_controller);
serve(config, router).await
}

Expand All @@ -32,10 +33,11 @@ pub async fn run_server(
pub async fn run_test_driver_server(
config: HttpServerConfig,
store: SharedStore,
signed_blocks: SharedSignedBlocks,
aggregator_controller: SharedController,
) -> Result<()> {
let driver_state = TestDriverState::new(store.clone());
let router = normal_routes(&config, store, aggregator_controller)
let router = normal_routes(&config, store, signed_blocks, aggregator_controller)
.merge(test_driver_routes(driver_state));
serve(config, router).await
}
Expand Down
109 changes: 107 additions & 2 deletions lean_client/http_api/tests/api_endpoint.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,12 +12,16 @@ use axum::{
body::Body,
http::{Request, header::CONTENT_TYPE},
};
use containers::{Block, BlockBody, Checkpoint, MultiMessageAggregate, SignedBlock, Slot};
use fork_choice::store::Store;
use http_api::{AggregatorController, HttpServerConfig, SharedStore, normal_routes};
use http_api::{
AggregatorController, HttpServerConfig, SharedSignedBlocks, SharedStore, normal_routes,
};
use http_body_util::BodyExt;
use parking_lot::RwLock;
use serde::Deserialize;
use serde_json::Value;
use ssz::{H256, SszHash, SszReadDefault as _};
use test_generator::test_resources;
use tower::ServiceExt;

Expand Down Expand Up @@ -73,10 +77,12 @@ fn api_endpoint(spec_file: &str) {
..Default::default()
}));

let signed_blocks: SharedSignedBlocks = Arc::new(RwLock::new(HashMap::new()));

let controller = Some(Arc::new(AggregatorController::new(store.clone(), None)));

let config = HttpServerConfig::default();
let router = normal_routes(&config, store, controller);
let router = normal_routes(&config, store, signed_blocks, controller);

let request = match (case.method.as_str(), case.request_body.as_ref()) {
("POST", Some(body)) => {
Expand Down Expand Up @@ -138,3 +144,102 @@ fn api_endpoint(spec_file: &str) {
}
});
}

fn sample_signed_block() -> SignedBlock {
SignedBlock {
block: Block {
slot: Slot(13),
proposer_index: 0,
parent_root: H256::default(),
state_root: H256::default(),
body: BlockBody {
attestations: Default::default(),
},
},
proof: MultiMessageAggregate::default(),
}
}

fn blocks_finalized_request() -> Request<Body> {
Request::builder()
.method("GET")
.uri("/lean/v0/blocks/finalized")
.body(Body::empty())
.expect("build request")
}

#[test]
fn blocks_finalized_returns_signed_block_ssz() {
let rt = tokio::runtime::Runtime::new().expect("tokio runtime");
rt.block_on(async move {
let signed_block = sample_signed_block();
let root = signed_block.block.hash_tree_root();

let store: SharedStore = Arc::new(RwLock::new(Store {
latest_finalized: Checkpoint {
root,
slot: signed_block.block.slot,
},
..Default::default()
}));
let signed_blocks: SharedSignedBlocks =
Arc::new(RwLock::new(HashMap::from([(root, signed_block.clone())])));
let controller = Some(Arc::new(AggregatorController::new(store.clone(), None)));

let config = HttpServerConfig::default();
let router = normal_routes(&config, store, signed_blocks, controller);

let response = router
.oneshot(blocks_finalized_request())
.await
.expect("router oneshot");

assert_eq!(response.status().as_u16(), 200);
assert_eq!(
response
.headers()
.get(CONTENT_TYPE)
.and_then(|value| value.to_str().ok()),
Some("application/octet-stream"),
);

let body_bytes = response
.into_body()
.collect()
.await
.expect("collect response body")
.to_bytes();

let decoded = SignedBlock::from_ssz_default(&body_bytes).expect("decode SignedBlock SSZ");
assert_eq!(decoded.block.hash_tree_root(), root);
assert_eq!(decoded.block.slot, signed_block.block.slot);
});
}

#[test]
fn blocks_finalized_returns_404_when_signed_block_missing() {
let rt = tokio::runtime::Runtime::new().expect("tokio runtime");
rt.block_on(async move {
let signed_block = sample_signed_block();

let store: SharedStore = Arc::new(RwLock::new(Store {
latest_finalized: Checkpoint {
root: signed_block.block.hash_tree_root(),
slot: signed_block.block.slot,
},
..Default::default()
}));
let signed_blocks: SharedSignedBlocks = Arc::new(RwLock::new(HashMap::new()));
let controller = Some(Arc::new(AggregatorController::new(store.clone(), None)));

let config = HttpServerConfig::default();
let router = normal_routes(&config, store, signed_blocks, controller);

let response = router
.oneshot(blocks_finalized_request())
.await
.expect("router oneshot");

assert_eq!(response.status().as_u16(), 404);
});
}
11 changes: 10 additions & 1 deletion lean_client/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1100,6 +1100,7 @@ async fn main() -> Result<()> {
let chain_outbound_sender = outbound_p2p_sender.clone();

let http_store = store.clone();
let http_signed_blocks = signed_block_provider.clone();
let aggregator_controller =
Arc::new(AggregatorController::new(store.clone(), vs_for_controller));
// The hive `spec-assets-*` test suites drive the client through the
Expand All @@ -1114,17 +1115,25 @@ async fn main() -> Result<()> {
.map(str::trim),
Some("1") | Some("true") | Some("TRUE") | Some("yes")
);

task::spawn(async move {
let result = if test_driver_enabled {
info!("HTTP server starting in test-driver mode (HIVE_LEAN_TEST_DRIVER=1)");
http_api::run_test_driver_server(
args.http_config,
http_store,
http_signed_blocks,
Some(aggregator_controller),
)
.await
} else {
http_api::run_server(args.http_config, http_store, Some(aggregator_controller)).await
http_api::run_server(
args.http_config,
http_store,
http_signed_blocks,
Some(aggregator_controller),
)
.await
};
if let Err(err) = result {
error!("HTTP Server failed with error: {err:?}");
Expand Down
Loading