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
Original file line number Diff line number Diff line change
@@ -0,0 +1,69 @@
use crate::inclusion_list_verification::InclusionListVerificationError;
use crate::{BeaconChain, BeaconChainTypes};
use tracing::debug;
use types::{ChainSpec, SignedInclusionList};

pub struct GossipVerificationContext<'a, T: BeaconChainTypes> {
// TODO(heze): complete while implementing the gossip verification of inclusion lists
pub slot_clock: &'a T::SlotClock,
pub spec: &'a ChainSpec,
}

pub struct GossipVerifiedInclusionList {
pub signed_inclusion_list: SignedInclusionList,
pub is_timely: bool,
}

impl GossipVerifiedInclusionList {
pub fn new<T: BeaconChainTypes>(
signed_inclusion_list: SignedInclusionList,
_ctx: &GossipVerificationContext<'_, T>,
) -> Result<Self, InclusionListVerificationError> {
// TODO(heze): implement gossip verification for inclusion lists
Ok(Self {
signed_inclusion_list,
is_timely: true,
})
}
}

impl<T: BeaconChainTypes> BeaconChain<T> {
pub fn inclusion_list_gossip_verification_context(&self) -> GossipVerificationContext<'_, T> {
GossipVerificationContext {
slot_clock: &self.slot_clock,
spec: &self.spec,
}
}

pub fn verify_inclusion_list_for_gossip(
&self,
signed_inclusion_list: SignedInclusionList,
) -> Result<GossipVerifiedInclusionList, InclusionListVerificationError> {
let slot = signed_inclusion_list.message.slot;
let validator_index = signed_inclusion_list.message.validator_index;

let ctx = self.inclusion_list_gossip_verification_context();
match GossipVerifiedInclusionList::new(signed_inclusion_list, &ctx) {
Ok(verified) => {
debug!(
%slot,
%validator_index,
"Successfully verified gossip inclusion list"
);

// TODO(heze): emit the inclusion_list SSE event

Ok(verified)
}
Err(e) => {
debug!(
error = ?e,
%slot,
%validator_index,
"Rejected gossip inclusion list"
);
Err(e)
}
}
}
}
35 changes: 35 additions & 0 deletions beacon_node/beacon_chain/src/inclusion_list_verification/mod.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,35 @@
use crate::BeaconChainError;
use std::sync::Arc;
use types::{BeaconStateError, Slot};

pub mod gossip_verified_inclusion_list;

#[derive(Debug)]
pub enum InclusionListVerificationError {
/// Two valid inclusion lists were already seen from this validator for this slot.
AlreadySeenTwice { validator_index: u64, slot: Slot },
/// The slot clock cannot read.
UnableToReadSlot,
/// Beacon Chain error
BeaconChainError(Arc<BeaconChainError>),
/// Beacon State error
BeaconStateError(BeaconStateError),
}

impl std::fmt::Display for InclusionListVerificationError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{:?}", self)
}
}

impl From<BeaconChainError> for InclusionListVerificationError {
fn from(e: BeaconChainError) -> Self {
InclusionListVerificationError::BeaconChainError(Arc::new(e))
}
}

impl From<BeaconStateError> for InclusionListVerificationError {
fn from(e: BeaconStateError) -> Self {
InclusionListVerificationError::BeaconStateError(e)
}
}
1 change: 1 addition & 0 deletions beacon_node/beacon_chain/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@ pub mod fork_choice_signal;
pub mod graffiti_calculator;
pub mod historical_blocks;
pub mod historical_data_columns;
pub mod inclusion_list_verification;
pub mod invariants;
pub mod kzg_utils;
pub mod light_client_finality_update_verification;
Expand Down
25 changes: 24 additions & 1 deletion beacon_node/http_api/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1906,6 +1906,27 @@ pub async fn serve<T: BeaconChainTypes>(
},
);

/*
* inclusion lists
*/

// POST validator/inclusion_list (JSON)
let post_validator_inclusion_list = post_validator_inclusion_list(
eth_v1.clone(),
not_while_syncing_filter.clone(),
task_spawner_filter.clone(),
chain_filter.clone(),
network_tx_filter.clone(),
);

// POST validator/inclusion_list (SSZ)
let post_validator_inclusion_list_ssz = post_validator_inclusion_list_ssz(
eth_v1.clone(),
not_while_syncing_filter.clone(),
task_spawner_filter.clone(),
chain_filter.clone(),
network_tx_filter.clone(),
);
/*
* config
*/
Expand Down Expand Up @@ -3464,7 +3485,8 @@ pub async fn serve<T: BeaconChainTypes>(
.uor(post_beacon_execution_payload_envelopes_ssz)
.uor(post_beacon_execution_payload_bids_ssz)
.uor(post_beacon_pool_payload_attestations_ssz)
.uor(post_validator_proposer_preferences_ssz),
.uor(post_validator_proposer_preferences_ssz)
.uor(post_validator_inclusion_list_ssz),
)
.uor(post_beacon_blocks)
.uor(post_beacon_blinded_blocks)
Expand All @@ -3478,6 +3500,7 @@ pub async fn serve<T: BeaconChainTypes>(
.uor(post_beacon_pool_payload_attestations)
.uor(post_beacon_pool_bls_to_execution_changes)
.uor(post_validator_proposer_preferences)
.uor(post_validator_inclusion_list)
.uor(post_beacon_execution_payload_envelopes)
.uor(post_beacon_execution_payload_bids)
.uor(post_beacon_state_validators)
Expand Down
153 changes: 149 additions & 4 deletions beacon_node/http_api/src/validator/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ use crate::utils::{
use crate::version::{V1, V2, V3, V4, add_ssz_content_type_header, unsupported_version_rejection};
use crate::{StateId, attester_duties, proposer_duties, ptc_duties, sync_committees};
use beacon_chain::attestation_verification::VerifiedAttestation;
use beacon_chain::inclusion_list_verification::InclusionListVerificationError;
use beacon_chain::proposer_preferences_verification::ProposerPreferencesError;
use beacon_chain::{AttestationError, BeaconChain, BeaconChainError, BeaconChainTypes};
use bls::PublicKeyBytes;
Expand All @@ -31,8 +32,8 @@ use tokio::sync::oneshot;
use tracing::{debug, error, info, warn};
use types::{
BeaconState, Epoch, EthSpec, ForkName, ProposerPreparationData, SignedAggregateAndProof,
SignedContributionAndProof, SignedProposerPreferences, SignedValidatorRegistrationData, Slot,
SyncContributionData, ValidatorSubscription,
SignedContributionAndProof, SignedInclusionList, SignedProposerPreferences,
SignedValidatorRegistrationData, Slot, SyncContributionData, ValidatorSubscription,
};
use warp::{Filter, Rejection, Reply, http::response::Builder};
use warp_utils::reject::convert_rejection;
Expand Down Expand Up @@ -71,7 +72,7 @@ pub fn get_validator_sync_committee_contribution<T: BeaconChainTypes>(
.and(warp::path("sync_committee_contribution"))
.and(warp::path::end())
.and(warp::query::<SyncContributionData>())
.and(not_while_syncing_filter.clone())
.and(not_while_syncing_filter)
.and(task_spawner_filter.clone())
.and(chain_filter.clone())
.then(
Expand Down Expand Up @@ -118,7 +119,7 @@ pub fn post_validator_duties_sync<T: BeaconChainTypes>(
))
}))
.and(warp::path::end())
.and(not_while_syncing_filter.clone())
.and(not_while_syncing_filter)
.and(warp_utils::json::json())
.and(task_spawner_filter.clone())
.and(chain_filter.clone())
Expand Down Expand Up @@ -1297,3 +1298,147 @@ fn publish_proposer_preferences<T: BeaconChainTypes>(
))
}
}

/// POST validator/inclusion_list (JSON)
pub fn post_validator_inclusion_list<T: BeaconChainTypes>(
eth_v1: EthV1Filter,
not_while_syncing_filter: NotWhileSyncingFilter,
task_spawner_filter: TaskSpawnerFilter<T>,
chain_filter: ChainFilter<T>,
network_tx_filter: NetworkTxFilter<T>,
) -> ResponseFilter {
eth_v1
.and(warp::path("validator"))
.and(warp::path("inclusion_list"))
.and(warp::path::end())
.and(warp_utils::json::json())
.and(warp::header::<ForkName>(CONSENSUS_VERSION_HEADER))
.and(not_while_syncing_filter.clone())
.and(task_spawner_filter)
.and(chain_filter)
.and(network_tx_filter)
.then(
|request_body: GenericResponse<SignedInclusionList>,
fork_name: ForkName,
not_synced_filter: Result<(), Rejection>,
task_spawner: TaskSpawner<T::EthSpec>,
chain: Arc<BeaconChain<T>>,
network_tx: UnboundedSender<NetworkMessage<T::EthSpec>>| {
task_spawner.blocking_response_task(Priority::P0, move || {
not_synced_filter?;
ensure_heze_consensus_version(fork_name)?;
publish_inclusion_list(&chain, &network_tx, request_body.data)?;
Ok(warp::reply())
})
},
)
.boxed()
}

/// POST validator/inclusion_list (SSZ)
pub fn post_validator_inclusion_list_ssz<T: BeaconChainTypes>(
eth_v1: EthV1Filter,
not_while_syncing_filter: NotWhileSyncingFilter,
task_spawner_filter: TaskSpawnerFilter<T>,
chain_filter: ChainFilter<T>,
network_tx_filter: NetworkTxFilter<T>,
) -> ResponseFilter {
eth_v1
.and(warp::path("validator"))
.and(warp::path("inclusion_list"))
.and(warp::path::end())
.and(warp::body::bytes())
.and(warp::header::<ForkName>(CONSENSUS_VERSION_HEADER))
.and(not_while_syncing_filter)
.and(task_spawner_filter)
.and(chain_filter)
.and(network_tx_filter)
.then(
|body_bytes: Bytes,
fork_name: ForkName,
not_synced_filter: Result<(), Rejection>,
task_spawner: TaskSpawner<T::EthSpec>,
chain: Arc<BeaconChain<T>>,
network_tx: UnboundedSender<NetworkMessage<T::EthSpec>>| {
task_spawner.blocking_response_task(Priority::P0, move || {
not_synced_filter?;
ensure_heze_consensus_version(fork_name)?;
let signed_inclusion_list = SignedInclusionList::from_ssz_bytes(&body_bytes)
.map_err(|e| {
warp_utils::reject::custom_bad_request(format!("invalid SSZ: {e:?}"))
})?;
publish_inclusion_list(&chain, &network_tx, signed_inclusion_list)?;
Ok(warp::reply())
})
},
)
.boxed()
}

fn ensure_heze_consensus_version(fork_name: ForkName) -> Result<(), Rejection> {
if !fork_name.heze_enabled() {
return Err(warp_utils::reject::custom_bad_request(format!(
"Eth-Consensus-Version {fork_name} is not supported for inclusion lists"
)));
}
Ok(())
}

fn publish_inclusion_list<T: BeaconChainTypes>(
chain: &BeaconChain<T>,
network_tx: &UnboundedSender<NetworkMessage<T::EthSpec>>,
signed_inclusion_list: SignedInclusionList,
) -> Result<(), warp::Rejection> {
let slot = signed_inclusion_list.message.slot;
let validator_index = signed_inclusion_list.message.validator_index;
let fork_name = chain.spec.fork_name_at_slot::<T::EthSpec>(slot);

if !fork_name.heze_enabled() {
return Err(warp_utils::reject::custom_bad_request(
"Inclusion lists publishing is not supported before the Heze fork".into(),
));
}

match chain.verify_inclusion_list_for_gossip(signed_inclusion_list) {
Ok(verified_inclusion_list) => {
crate::utils::publish_pubsub_message(
network_tx,
PubsubMessage::InclusionList(Box::new(
verified_inclusion_list.signed_inclusion_list,
)),
)?;
Ok(())
}
Err(InclusionListVerificationError::AlreadySeenTwice { .. }) => {
debug!(
%slot,
%validator_index,
"Two valid inclusion lists were already seen"
);
Ok(())
}
Err(
e @ (InclusionListVerificationError::BeaconChainError(_)
| InclusionListVerificationError::BeaconStateError(_)
| InclusionListVerificationError::UnableToReadSlot),
) => {
error!(%slot, error = ?e, "Internal error verifying inclusion list");
Err(warp_utils::reject::custom_server_error(format!(
"internal error verifying inclusion list: {e}"
)))
}
// TODO(heze): remove once the IL gossip verification errors are added to InclusionListVerificationError
#[allow(unreachable_patterns)]
Err(e) => {
warn!(
%slot,
%validator_index,
error = ?e,
"Inclusion list failed gossip verification"
);
Err(warp_utils::reject::custom_bad_request(format!(
"inclusion list failed gossip verification: {e}"
)))
}
}
}
Loading
Loading