Skip to content
Open
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
7 changes: 7 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

19 changes: 19 additions & 0 deletions crates/networking/manager/src/p2p_sender.rs
Original file line number Diff line number Diff line change
Expand Up @@ -90,4 +90,23 @@ impl P2PSender {
warn!("Failed to send error response: {err}");
}
}

pub fn send_invalid_request(
&self,
peer_id: PeerId,
connection_id: ConnectionId,
stream_id: u64,
error: &str,
) {
if let Err(err) = self.0.send(P2PMessage::Response(P2PResponse {
peer_id,
connection_id,
stream_id,
message: Box::new(RespMessage::Error(ReqRespError::InvalidData(
error.to_string(),
))),
})) {
warn!("Failed to send invalid-request response: {err}");
}
}
}
148 changes: 119 additions & 29 deletions crates/networking/manager/src/req_resp.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,11 +3,18 @@ use std::sync::Arc;
use libp2p::{PeerId, swarm::ConnectionId};
use ream_consensus_beacon::{blob_sidecar::BlobIdentifier, data_column_sidecar::ColumnIdentifier};
use ream_p2p::network::beacon::network_state::NetworkState;
use ream_req_resp::beacon::messages::{
BeaconRequestMessage, BeaconResponseMessage,
blob_sidecars::{BlobSidecarsByRangeV1Request, BlobSidecarsByRootV1Request},
blocks::{BeaconBlocksByRangeV2Request, BeaconBlocksByRootV2Request},
data_column_sidecars::{DataColumnSidecarsByRangeV1Request, DataColumnSidecarsByRootV1Request},
use ream_req_resp::{
beacon::messages::{
BeaconRequestMessage, BeaconResponseMessage,
blob_sidecars::{BlobSidecarsByRangeV1Request, BlobSidecarsByRootV1Request},
blocks::{BeaconBlocksByRangeV2Request, BeaconBlocksByRootV2Request},
data_column_sidecars::{
DataColumnSidecarsByRangeV1Request, DataColumnSidecarsByRootV1Request,
},
},
constants::{
MAX_REQUEST_BLOCKS, MAX_REQUEST_BLOCKS_DENEB, MAX_REQUEST_DATA_COLUMN_SIDECARS_PER_COLUMN,
},
};
use ream_storage::{
db::beacon::BeaconDB,
Expand Down Expand Up @@ -50,16 +57,41 @@ pub async fn handle_req_resp_message(
count,
..
}) => {
for slot in start_slot..start_slot + count {
let Ok(Some(block_root)) = ream_db.slot_index_provider().get(slot) else {
trace!("No block root found for slot {slot}");
p2p_sender.send_error_response(
peer_id,
connection_id,
stream_id,
&format!("No block root found for slot {slot}"),
);
return;
if count > MAX_REQUEST_BLOCKS {
p2p_sender.send_invalid_request(

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Same here, let's read from the beacon state: beacon_network_spec().max_request_blocks

peer_id,
connection_id,
stream_id,
&format!("Requested count {count} exceeds MAX_REQUEST_BLOCKS"),
);
return;
}
let Some(end_slot_exclusive) = start_slot.checked_add(count) else {
p2p_sender.send_invalid_request(
peer_id,
connection_id,
stream_id,
&format!("start_slot {start_slot} + count {count} overflows"),
);
return;
};

for slot in start_slot..end_slot_exclusive {
let block_root = match ream_db.slot_index_provider().get(slot) {
Ok(Some(block_root)) => block_root,
Ok(None) => {
trace!("No block root found for slot {slot}");
continue;
}
Err(err) => {
p2p_sender.send_error_response(
peer_id,
connection_id,
stream_id,
&format!("Failed to read slot index for slot {slot}: {err:?}"),
);
return;
}
};
let Ok(Some(block)) = ream_db.block_provider().get(block_root) else {
trace!("No block found for root {block_root}");
Expand Down Expand Up @@ -109,16 +141,41 @@ pub async fn handle_req_resp_message(
start_slot,
count,
}) => {
for slot in start_slot..start_slot + count {
let Ok(Some(block_root)) = ream_db.slot_index_provider().get(slot) else {
trace!("No block root found for slot {slot}");
p2p_sender.send_error_response(
peer_id,
connection_id,
stream_id,
&format!("No block root found for slot {slot}"),
);
return;
if count > MAX_REQUEST_BLOCKS_DENEB {
Comment on lines 143 to +144

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

lets avoid using hardcoded const when we can read from the beacon state: beacon_network_spec().max_request_blocks_deneb

p2p_sender.send_invalid_request(
peer_id,
connection_id,
stream_id,
&format!("Requested count {count} exceeds MAX_REQUEST_BLOCKS_DENEB"),
);
return;
}
let Some(end_slot_exclusive) = start_slot.checked_add(count) else {
p2p_sender.send_invalid_request(
peer_id,
connection_id,
stream_id,
&format!("start_slot {start_slot} + count {count} overflows"),
);
return;
};

for slot in start_slot..end_slot_exclusive {
let block_root = match ream_db.slot_index_provider().get(slot) {
Ok(Some(block_root)) => block_root,
Ok(None) => {
trace!("No block root found for slot {slot}");
continue;
}
Err(err) => {
p2p_sender.send_error_response(
peer_id,
connection_id,
stream_id,
&format!("Failed to read slot index for slot {slot}: {err:?}"),
);
return;
}
};
let Ok(Some(block)) = ream_db.block_provider().get(block_root) else {
trace!("No block found for root {block_root}");
Expand Down Expand Up @@ -232,10 +289,43 @@ pub async fn handle_req_resp_message(
count,
columns,
}) => {
for slot in start_slot..start_slot + count {
let Ok(Some(block_root)) = ream_db.slot_index_provider().get(slot) else {
trace!("No block root found for slot {slot}");
continue;
if count > MAX_REQUEST_DATA_COLUMN_SIDECARS_PER_COLUMN {
p2p_sender.send_invalid_request(
peer_id,
connection_id,
stream_id,
&format!(
"Requested count {count} exceeds MAX_REQUEST_DATA_COLUMN_SIDECARS_PER_COLUMN"
),
);
return;
}
let Some(end_slot_exclusive) = start_slot.checked_add(count) else {
p2p_sender.send_invalid_request(
peer_id,
connection_id,
stream_id,
&format!("start_slot {start_slot} + count {count} overflows"),
);
return;
};

for slot in start_slot..end_slot_exclusive {
let block_root = match ream_db.slot_index_provider().get(slot) {
Ok(Some(block_root)) => block_root,
Ok(None) => {
trace!("No block root found for slot {slot}");
continue;
}
Err(err) => {
p2p_sender.send_error_response(
peer_id,
connection_id,
stream_id,
&format!("Failed to read slot index for slot {slot}: {err:?}"),
);
return;
}
};

for &column_index in &columns {
Expand Down
70 changes: 47 additions & 23 deletions crates/networking/manager/src/service.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,8 @@ use alloy_primitives::B256;
use libp2p::PeerId;
use ream_chain_beacon::beacon_chain::BeaconChain;
use ream_consensus_misc::{
constants::beacon::NUM_CUSTODY_GROUPS, misc::compute_start_slot_at_epoch,
constants::beacon::{NUM_CUSTODY_GROUPS, SLOTS_PER_EPOCH},
misc::compute_start_slot_at_epoch,
};
use ream_discv5::{
config::DiscoveryConfig,
Expand Down Expand Up @@ -37,7 +38,10 @@ use ream_syncer::{
block_range::BlockRangeSyncer,
unknown_parent_lookups::{MAX_LOOKUPS, UnknownBlockMeta, UnknownParentLookupCoordinator},
};
use tokio::{sync::mpsc, time::interval};
use tokio::{
sync::mpsc,
time::{interval, sleep},
};
use tracing::{error, info, warn};
use tree_hash::TreeHash;

Expand Down Expand Up @@ -107,6 +111,7 @@ pub struct NetworkManagerService {
pub ream_db: BeaconDB,
pub cached_db: Arc<BeaconCacheDB>,
pub sync_committee_pool: Arc<SyncCommitteePool>,
pub executor: ReamExecutor,
}

struct ReconciledBlockLookupState {
Expand Down Expand Up @@ -274,6 +279,7 @@ impl NetworkManagerService {
ream_db,
cached_db,
sync_committee_pool,
executor,
})
}

Expand All @@ -290,9 +296,13 @@ impl NetworkManagerService {
cached_db,
network_state,
block_range_syncer,
executor,
..
} = self;

let mut resync_check_interval = interval(Duration::from_secs(
SLOTS_PER_EPOCH * beacon_network_spec().seconds_per_slot(),
));
let mut interval = interval(Duration::from_secs(
beacon_network_spec().seconds_per_slot(),
));
Expand All @@ -314,6 +324,8 @@ impl NetworkManagerService {
let mut syncer_handle = block_range_syncer.start();
// Avoid polling a completed JoinHandle after the syncer has caught up.
let mut syncer_active = true;
// Kept alive (not dropped) once synced, so the periodic re-check can restart it later.
let mut idle_syncer: Option<BlockRangeSyncer> = None;
loop {
tokio::select! {
// Drive unknown-parent lookup actions and results.
Expand Down Expand Up @@ -527,31 +539,43 @@ impl NetworkManagerService {
// Restart range sync until the finalized target is reached.
result = &mut syncer_handle, if syncer_active => {
syncer_active = false;
let joined_result = match result {
Ok(joined_result) => joined_result,
Err(err) => {
error!("Block range syncer failed to join task: {err}");
continue;
match result {
Ok(Ok((mut block_range_syncer, sync_result))) => {
if let Err(err) = sync_result {
warn!("Block range sync segment failed: {err:?}");
}
if block_range_syncer.is_synced_to_head_slot().await {
idle_syncer = Some(block_range_syncer);
} else {
syncer_handle = block_range_syncer.start();
syncer_active = true;
}
}
};

let thread_result = match joined_result {
Ok(result) => result,
Err(err) => {
error!("Block range syncer thread failed: {err}");
continue;
Ok(Err(err)) => {
// Executor shutdown cancelled the task; do not re-arm.
error!("Block range syncer task cancelled: {err}");
}
};

let block_range_syncer = match thread_result {
Ok(syncer) => syncer,
Err(err) => {
error!("Block range syncer failed to start: {err}");
continue;
error!("Block range syncer task panicked, reconstructing: {err}");
sleep(Duration::from_secs(5)).await;
let block_range_syncer = BlockRangeSyncer::new(
beacon_chain.clone(),
p2p_sender.0.clone(),
network_state.clone(),
executor.clone(),
);
syncer_handle = block_range_syncer.start();
syncer_active = true;
}
};

if !block_range_syncer.is_synced_to_finalized_slot().await {
}
}
_ = resync_check_interval.tick(), if !syncer_active && idle_syncer.is_some() => {
let mut block_range_syncer = idle_syncer
.take()
.expect("checked by the select guard above");
if block_range_syncer.is_synced_to_head_slot().await {
idle_syncer = Some(block_range_syncer);
} else {
syncer_handle = block_range_syncer.start();
syncer_active = true;
}
Expand Down
5 changes: 3 additions & 2 deletions crates/networking/manager/src/unknown_parent_lookup.rs
Original file line number Diff line number Diff line change
Expand Up @@ -337,6 +337,7 @@ mod tests {
use std::sync::Arc;

use anyhow::anyhow;
use ream_p2p::network::beacon::channel::P2PCallbackError;
use tokio::sync::mpsc::{self, UnboundedReceiver, error::TryRecvError};

use super::*;
Expand All @@ -350,7 +351,7 @@ mod tests {
receiver: &mut UnboundedReceiver<P2PMessage>,
expected_peer: PeerId,
expected_root: B256,
) -> mpsc::Sender<anyhow::Result<P2PCallbackResponse>> {
) -> mpsc::Sender<Result<P2PCallbackResponse, P2PCallbackError>> {
let message = receiver.recv().await.expect("request should be sent");
let P2PMessage::Request(P2PRequest::BlockRoots {
peer_id,
Expand Down Expand Up @@ -490,7 +491,7 @@ mod tests {

let callback = intercept_request(&mut receiver, peer_id, expected_root).await;
callback
.send(Err(anyhow!("transport failed")))
.send(Err(anyhow!("transport failed").into()))
.await
.expect("callback error should be delivered");

Expand Down
Loading
Loading