mirror of
https://github.com/pezkuwichain/pezkuwi-subxt.git
synced 2026-07-21 12:05:42 +00:00
[BEEFY] Avoid missing voting sessions during node restart (#3074)
Related to https://github.com/paritytech/polkadot-sdk/issues/3003 and https://github.com/paritytech/polkadot-sdk/issues/2842 --------- Co-authored-by: Adrian Catangiu <adrian@parity.io>
This commit is contained in:
@@ -17,18 +17,20 @@
|
||||
// along with this program. If not, see <https://www.gnu.org/licenses/>.
|
||||
|
||||
use crate::{
|
||||
aux_schema,
|
||||
communication::{
|
||||
gossip::{proofs_topic, votes_topic, GossipFilterCfg, GossipMessage, GossipValidator},
|
||||
peers::PeerReport,
|
||||
request_response::outgoing_requests_engine::{OnDemandJustificationsEngine, ResponseInfo},
|
||||
},
|
||||
error::Error,
|
||||
expect_validator_set,
|
||||
justification::BeefyVersionedFinalityProof,
|
||||
keystore::{BeefyKeystore, BeefySignatureHasher},
|
||||
metric_inc, metric_set,
|
||||
metrics::VoterMetrics,
|
||||
round::{Rounds, VoteImportResult},
|
||||
BeefyVoterLinks, LOG_TARGET,
|
||||
wait_for_parent_header, BeefyVoterLinks, HEADER_SYNC_DELAY, LOG_TARGET,
|
||||
};
|
||||
use codec::{Codec, Decode, DecodeAll, Encode};
|
||||
use futures::{stream::Fuse, FutureExt, StreamExt};
|
||||
@@ -38,6 +40,7 @@ use sc_network_gossip::GossipEngine;
|
||||
use sc_utils::{mpsc::TracingUnboundedReceiver, notification::NotificationReceiver};
|
||||
use sp_api::ProvideRuntimeApi;
|
||||
use sp_arithmetic::traits::{AtLeast32Bit, Saturating};
|
||||
use sp_blockchain::Backend as BlockchainBackend;
|
||||
use sp_consensus::SyncOracle;
|
||||
use sp_consensus_beefy::{
|
||||
check_equivocation_proof,
|
||||
@@ -53,6 +56,7 @@ use sp_runtime::{
|
||||
use std::{
|
||||
collections::{BTreeMap, BTreeSet, VecDeque},
|
||||
fmt::Debug,
|
||||
marker::PhantomData,
|
||||
sync::Arc,
|
||||
};
|
||||
|
||||
@@ -176,6 +180,13 @@ impl<B: Block> VoterOracle<B> {
|
||||
}
|
||||
}
|
||||
|
||||
// Check if an observed session can be added to the Oracle.
|
||||
fn can_add_session(&self, session_start: NumberFor<B>) -> bool {
|
||||
let latest_known_session_start =
|
||||
self.sessions.back().map(|session| session.session_start());
|
||||
Some(session_start) > latest_known_session_start
|
||||
}
|
||||
|
||||
/// Add new observed session to the Oracle.
|
||||
pub fn add_session(&mut self, rounds: Rounds<B>) {
|
||||
self.sessions.push_back(rounds);
|
||||
@@ -236,12 +247,10 @@ impl<B: Block> VoterOracle<B> {
|
||||
/// Return `Some(number)` if we should be voting on block `number`,
|
||||
/// return `None` if there is no block we should vote on.
|
||||
pub fn voting_target(&self) -> Option<NumberFor<B>> {
|
||||
let rounds = if let Some(r) = self.sessions.front() {
|
||||
r
|
||||
} else {
|
||||
let rounds = self.sessions.front().or_else(|| {
|
||||
debug!(target: LOG_TARGET, "🥩 No voting round started");
|
||||
return None
|
||||
};
|
||||
None
|
||||
})?;
|
||||
let best_grandpa = *self.best_grandpa_block_header.number();
|
||||
let best_beefy = self.best_beefy_block;
|
||||
|
||||
@@ -327,50 +336,171 @@ pub(crate) struct BeefyComms<B: Block> {
|
||||
pub on_demand_justifications: OnDemandJustificationsEngine<B>,
|
||||
}
|
||||
|
||||
/// A BEEFY worker plays the BEEFY protocol
|
||||
pub(crate) struct BeefyWorker<B: Block, BE, P, RuntimeApi, S> {
|
||||
pub(crate) struct BeefyWorkerBase<B: Block, BE, RuntimeApi> {
|
||||
// utilities
|
||||
pub backend: Arc<BE>,
|
||||
pub payload_provider: P,
|
||||
pub runtime: Arc<RuntimeApi>,
|
||||
pub sync: Arc<S>,
|
||||
pub key_store: BeefyKeystore,
|
||||
|
||||
// communication (created once, but returned and reused if worker is restarted/reinitialized)
|
||||
pub comms: BeefyComms<B>,
|
||||
|
||||
// channels
|
||||
/// Links between the block importer, the background voter and the RPC layer.
|
||||
pub links: BeefyVoterLinks<B>,
|
||||
|
||||
// voter state
|
||||
/// BEEFY client metrics.
|
||||
pub metrics: Option<VoterMetrics>,
|
||||
/// Buffer holding justifications for future processing.
|
||||
pub pending_justifications: BTreeMap<NumberFor<B>, BeefyVersionedFinalityProof<B>>,
|
||||
/// Persisted voter state.
|
||||
pub persisted_state: PersistedState<B>,
|
||||
|
||||
pub _phantom: PhantomData<B>,
|
||||
}
|
||||
|
||||
impl<B, BE, P, R, S> BeefyWorker<B, BE, P, R, S>
|
||||
impl<B, BE, R> BeefyWorkerBase<B, BE, R>
|
||||
where
|
||||
B: Block + Codec,
|
||||
BE: Backend<B>,
|
||||
P: PayloadProvider<B>,
|
||||
S: SyncOracle,
|
||||
R: ProvideRuntimeApi<B>,
|
||||
R::Api: BeefyApi<B, AuthorityId>,
|
||||
{
|
||||
fn best_grandpa_block(&self) -> NumberFor<B> {
|
||||
*self.persisted_state.voting_oracle.best_grandpa_block_header.number()
|
||||
// If no persisted state present, walk back the chain from first GRANDPA notification to either:
|
||||
// - latest BEEFY finalized block, or if none found on the way,
|
||||
// - BEEFY pallet genesis;
|
||||
// Enqueue any BEEFY mandatory blocks (session boundaries) found on the way, for voter to
|
||||
// finalize.
|
||||
async fn init_state(
|
||||
&self,
|
||||
beefy_genesis: NumberFor<B>,
|
||||
best_grandpa: <B as Block>::Header,
|
||||
min_block_delta: u32,
|
||||
) -> Result<PersistedState<B>, Error> {
|
||||
let blockchain = self.backend.blockchain();
|
||||
|
||||
let beefy_genesis = self
|
||||
.runtime
|
||||
.runtime_api()
|
||||
.beefy_genesis(best_grandpa.hash())
|
||||
.ok()
|
||||
.flatten()
|
||||
.filter(|genesis| *genesis == beefy_genesis)
|
||||
.ok_or_else(|| Error::Backend("BEEFY pallet expected to be active.".into()))?;
|
||||
// Walk back the imported blocks and initialize voter either, at the last block with
|
||||
// a BEEFY justification, or at pallet genesis block; voter will resume from there.
|
||||
let mut sessions = VecDeque::new();
|
||||
let mut header = best_grandpa.clone();
|
||||
let state = loop {
|
||||
if let Some(true) = blockchain
|
||||
.justifications(header.hash())
|
||||
.ok()
|
||||
.flatten()
|
||||
.map(|justifs| justifs.get(BEEFY_ENGINE_ID).is_some())
|
||||
{
|
||||
info!(
|
||||
target: LOG_TARGET,
|
||||
"🥩 Initialize BEEFY voter at last BEEFY finalized block: {:?}.",
|
||||
*header.number()
|
||||
);
|
||||
let best_beefy = *header.number();
|
||||
// If no session boundaries detected so far, just initialize new rounds here.
|
||||
if sessions.is_empty() {
|
||||
let active_set =
|
||||
expect_validator_set(self.runtime.as_ref(), self.backend.as_ref(), &header)
|
||||
.await?;
|
||||
let mut rounds = Rounds::new(best_beefy, active_set);
|
||||
// Mark the round as already finalized.
|
||||
rounds.conclude(best_beefy);
|
||||
sessions.push_front(rounds);
|
||||
}
|
||||
let state = PersistedState::checked_new(
|
||||
best_grandpa,
|
||||
best_beefy,
|
||||
sessions,
|
||||
min_block_delta,
|
||||
beefy_genesis,
|
||||
)
|
||||
.ok_or_else(|| Error::Backend("Invalid BEEFY chain".into()))?;
|
||||
break state
|
||||
}
|
||||
|
||||
if *header.number() == beefy_genesis {
|
||||
// We've reached BEEFY genesis, initialize voter here.
|
||||
let genesis_set =
|
||||
expect_validator_set(self.runtime.as_ref(), self.backend.as_ref(), &header)
|
||||
.await?;
|
||||
info!(
|
||||
target: LOG_TARGET,
|
||||
"🥩 Loading BEEFY voter state from genesis on what appears to be first startup. \
|
||||
Starting voting rounds at block {:?}, genesis validator set {:?}.",
|
||||
beefy_genesis,
|
||||
genesis_set,
|
||||
);
|
||||
|
||||
sessions.push_front(Rounds::new(beefy_genesis, genesis_set));
|
||||
break PersistedState::checked_new(
|
||||
best_grandpa,
|
||||
Zero::zero(),
|
||||
sessions,
|
||||
min_block_delta,
|
||||
beefy_genesis,
|
||||
)
|
||||
.ok_or_else(|| Error::Backend("Invalid BEEFY chain".into()))?
|
||||
}
|
||||
|
||||
if let Some(active) = find_authorities_change::<B>(&header) {
|
||||
info!(
|
||||
target: LOG_TARGET,
|
||||
"🥩 Marking block {:?} as BEEFY Mandatory.",
|
||||
*header.number()
|
||||
);
|
||||
sessions.push_front(Rounds::new(*header.number(), active));
|
||||
}
|
||||
|
||||
// Move up the chain.
|
||||
header = wait_for_parent_header(blockchain, header, HEADER_SYNC_DELAY).await?;
|
||||
};
|
||||
|
||||
aux_schema::write_current_version(self.backend.as_ref())?;
|
||||
aux_schema::write_voter_state(self.backend.as_ref(), &state)?;
|
||||
Ok(state)
|
||||
}
|
||||
|
||||
fn voting_oracle(&self) -> &VoterOracle<B> {
|
||||
&self.persisted_state.voting_oracle
|
||||
}
|
||||
pub async fn load_or_init_state(
|
||||
&mut self,
|
||||
beefy_genesis: NumberFor<B>,
|
||||
best_grandpa: <B as Block>::Header,
|
||||
min_block_delta: u32,
|
||||
) -> Result<PersistedState<B>, Error> {
|
||||
// Initialize voter state from AUX DB if compatible.
|
||||
if let Some(mut state) = crate::aux_schema::load_persistent(self.backend.as_ref())?
|
||||
// Verify state pallet genesis matches runtime.
|
||||
.filter(|state| state.pallet_genesis() == beefy_genesis)
|
||||
{
|
||||
// Overwrite persisted state with current best GRANDPA block.
|
||||
state.set_best_grandpa(best_grandpa.clone());
|
||||
// Overwrite persisted data with newly provided `min_block_delta`.
|
||||
state.set_min_block_delta(min_block_delta);
|
||||
info!(target: LOG_TARGET, "🥩 Loading BEEFY voter state from db: {:?}.", state);
|
||||
|
||||
fn active_rounds(&mut self) -> Result<&Rounds<B>, Error> {
|
||||
self.persisted_state.voting_oracle.active_rounds()
|
||||
// Make sure that all the headers that we need have been synced.
|
||||
let mut new_sessions = vec![];
|
||||
let mut header = best_grandpa.clone();
|
||||
while *header.number() > state.best_beefy() {
|
||||
if state.voting_oracle.can_add_session(*header.number()) {
|
||||
if let Some(active) = find_authorities_change::<B>(&header) {
|
||||
new_sessions.push((active, *header.number()));
|
||||
}
|
||||
}
|
||||
header =
|
||||
wait_for_parent_header(self.backend.blockchain(), header, HEADER_SYNC_DELAY)
|
||||
.await?;
|
||||
}
|
||||
|
||||
// Make sure we didn't miss any sessions during node restart.
|
||||
for (validator_set, new_session_start) in new_sessions.drain(..).rev() {
|
||||
info!(
|
||||
target: LOG_TARGET,
|
||||
"🥩 Handling missed BEEFY session after node restart: {:?}.",
|
||||
new_session_start
|
||||
);
|
||||
self.init_session_at(&mut state, validator_set, new_session_start);
|
||||
}
|
||||
return Ok(state)
|
||||
}
|
||||
|
||||
// No valid voter-state persisted, re-initialize from pallet genesis.
|
||||
self.init_state(beefy_genesis, best_grandpa, min_block_delta).await
|
||||
}
|
||||
|
||||
/// Verify `active` validator set for `block` against the key store
|
||||
@@ -394,7 +524,7 @@ where
|
||||
if store.intersection(&active).count() == 0 {
|
||||
let msg = "no authority public key found in store".to_string();
|
||||
debug!(target: LOG_TARGET, "🥩 for block {:?} {}", block, msg);
|
||||
metric_inc!(self, beefy_no_authority_found_in_store);
|
||||
metric_inc!(self.metrics, beefy_no_authority_found_in_store);
|
||||
Err(Error::Keystore(msg))
|
||||
} else {
|
||||
Ok(())
|
||||
@@ -404,13 +534,14 @@ where
|
||||
/// Handle session changes by starting new voting round for mandatory blocks.
|
||||
fn init_session_at(
|
||||
&mut self,
|
||||
persisted_state: &mut PersistedState<B>,
|
||||
validator_set: ValidatorSet<AuthorityId>,
|
||||
new_session_start: NumberFor<B>,
|
||||
) {
|
||||
debug!(target: LOG_TARGET, "🥩 New active validator set: {:?}", validator_set);
|
||||
|
||||
// BEEFY should finalize a mandatory block during each session.
|
||||
if let Ok(active_session) = self.active_rounds() {
|
||||
if let Ok(active_session) = persisted_state.voting_oracle.active_rounds() {
|
||||
if !active_session.mandatory_done() {
|
||||
debug!(
|
||||
target: LOG_TARGET,
|
||||
@@ -418,7 +549,7 @@ where
|
||||
validator_set.id(),
|
||||
active_session.validator_set_id(),
|
||||
);
|
||||
metric_inc!(self, beefy_lagging_sessions);
|
||||
metric_inc!(self.metrics, beefy_lagging_sessions);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -428,10 +559,10 @@ where
|
||||
}
|
||||
|
||||
let id = validator_set.id();
|
||||
self.persisted_state
|
||||
persisted_state
|
||||
.voting_oracle
|
||||
.add_session(Rounds::new(new_session_start, validator_set));
|
||||
metric_set!(self, beefy_validator_set_id, id);
|
||||
metric_set!(self.metrics, beefy_validator_set_id, id);
|
||||
info!(
|
||||
target: LOG_TARGET,
|
||||
"🥩 New Rounds for validator set id: {:?} with session_start {:?}",
|
||||
@@ -439,6 +570,61 @@ where
|
||||
new_session_start
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
/// A BEEFY worker/voter that follows the BEEFY protocol
|
||||
pub(crate) struct BeefyWorker<B: Block, BE, P, RuntimeApi, S> {
|
||||
pub base: BeefyWorkerBase<B, BE, RuntimeApi>,
|
||||
|
||||
// utils
|
||||
pub payload_provider: P,
|
||||
pub sync: Arc<S>,
|
||||
|
||||
// communication (created once, but returned and reused if worker is restarted/reinitialized)
|
||||
pub comms: BeefyComms<B>,
|
||||
|
||||
// channels
|
||||
/// Links between the block importer, the background voter and the RPC layer.
|
||||
pub links: BeefyVoterLinks<B>,
|
||||
|
||||
// voter state
|
||||
/// Buffer holding justifications for future processing.
|
||||
pub pending_justifications: BTreeMap<NumberFor<B>, BeefyVersionedFinalityProof<B>>,
|
||||
/// Persisted voter state.
|
||||
pub persisted_state: PersistedState<B>,
|
||||
}
|
||||
|
||||
impl<B, BE, P, R, S> BeefyWorker<B, BE, P, R, S>
|
||||
where
|
||||
B: Block + Codec,
|
||||
BE: Backend<B>,
|
||||
P: PayloadProvider<B>,
|
||||
S: SyncOracle,
|
||||
R: ProvideRuntimeApi<B>,
|
||||
R::Api: BeefyApi<B, AuthorityId>,
|
||||
{
|
||||
fn best_grandpa_block(&self) -> NumberFor<B> {
|
||||
*self.persisted_state.voting_oracle.best_grandpa_block_header.number()
|
||||
}
|
||||
|
||||
fn voting_oracle(&self) -> &VoterOracle<B> {
|
||||
&self.persisted_state.voting_oracle
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
fn active_rounds(&mut self) -> Result<&Rounds<B>, Error> {
|
||||
self.persisted_state.voting_oracle.active_rounds()
|
||||
}
|
||||
|
||||
/// Handle session changes by starting new voting round for mandatory blocks.
|
||||
fn init_session_at(
|
||||
&mut self,
|
||||
validator_set: ValidatorSet<AuthorityId>,
|
||||
new_session_start: NumberFor<B>,
|
||||
) {
|
||||
self.base
|
||||
.init_session_at(&mut self.persisted_state, validator_set, new_session_start);
|
||||
}
|
||||
|
||||
fn handle_finality_notification(
|
||||
&mut self,
|
||||
@@ -452,7 +638,8 @@ where
|
||||
);
|
||||
let header = ¬ification.header;
|
||||
|
||||
self.runtime
|
||||
self.base
|
||||
.runtime
|
||||
.runtime_api()
|
||||
.beefy_genesis(header.hash())
|
||||
.ok()
|
||||
@@ -466,7 +653,7 @@ where
|
||||
self.persisted_state.set_best_grandpa(header.clone());
|
||||
|
||||
// Check all (newly) finalized blocks for new session(s).
|
||||
let backend = self.backend.clone();
|
||||
let backend = self.base.backend.clone();
|
||||
for header in notification
|
||||
.tree_route
|
||||
.iter()
|
||||
@@ -485,7 +672,7 @@ where
|
||||
}
|
||||
|
||||
if new_session_added {
|
||||
crate::aux_schema::write_voter_state(&*self.backend, &self.persisted_state)
|
||||
crate::aux_schema::write_voter_state(&*self.base.backend, &self.persisted_state)
|
||||
.map_err(|e| Error::Backend(e.to_string()))?;
|
||||
}
|
||||
|
||||
@@ -519,7 +706,7 @@ where
|
||||
true,
|
||||
);
|
||||
},
|
||||
RoundAction::Drop => metric_inc!(self, beefy_stale_votes),
|
||||
RoundAction::Drop => metric_inc!(self.base.metrics, beefy_stale_votes),
|
||||
RoundAction::Enqueue => error!(target: LOG_TARGET, "🥩 unexpected vote: {:?}.", vote),
|
||||
};
|
||||
Ok(())
|
||||
@@ -539,23 +726,23 @@ where
|
||||
match self.voting_oracle().triage_round(block_num)? {
|
||||
RoundAction::Process => {
|
||||
debug!(target: LOG_TARGET, "🥩 Process justification for round: {:?}.", block_num);
|
||||
metric_inc!(self, beefy_imported_justifications);
|
||||
metric_inc!(self.base.metrics, beefy_imported_justifications);
|
||||
self.finalize(justification)?
|
||||
},
|
||||
RoundAction::Enqueue => {
|
||||
debug!(target: LOG_TARGET, "🥩 Buffer justification for round: {:?}.", block_num);
|
||||
if self.pending_justifications.len() < MAX_BUFFERED_JUSTIFICATIONS {
|
||||
self.pending_justifications.entry(block_num).or_insert(justification);
|
||||
metric_inc!(self, beefy_buffered_justifications);
|
||||
metric_inc!(self.base.metrics, beefy_buffered_justifications);
|
||||
} else {
|
||||
metric_inc!(self, beefy_buffered_justifications_dropped);
|
||||
metric_inc!(self.base.metrics, beefy_buffered_justifications_dropped);
|
||||
warn!(
|
||||
target: LOG_TARGET,
|
||||
"🥩 Buffer justification dropped for round: {:?}.", block_num
|
||||
);
|
||||
}
|
||||
},
|
||||
RoundAction::Drop => metric_inc!(self, beefy_stale_justifications),
|
||||
RoundAction::Drop => metric_inc!(self.base.metrics, beefy_stale_justifications),
|
||||
};
|
||||
Ok(())
|
||||
}
|
||||
@@ -577,7 +764,7 @@ where
|
||||
// We created the `finality_proof` and know to be valid.
|
||||
// New state is persisted after finalization.
|
||||
self.finalize(finality_proof.clone())?;
|
||||
metric_inc!(self, beefy_good_votes_processed);
|
||||
metric_inc!(self.base.metrics, beefy_good_votes_processed);
|
||||
return Ok(Some(finality_proof))
|
||||
},
|
||||
VoteImportResult::Ok => {
|
||||
@@ -588,17 +775,20 @@ where
|
||||
.map(|(mandatory_num, _)| mandatory_num == block_number)
|
||||
.unwrap_or(false)
|
||||
{
|
||||
crate::aux_schema::write_voter_state(&*self.backend, &self.persisted_state)
|
||||
.map_err(|e| Error::Backend(e.to_string()))?;
|
||||
crate::aux_schema::write_voter_state(
|
||||
&*self.base.backend,
|
||||
&self.persisted_state,
|
||||
)
|
||||
.map_err(|e| Error::Backend(e.to_string()))?;
|
||||
}
|
||||
metric_inc!(self, beefy_good_votes_processed);
|
||||
metric_inc!(self.base.metrics, beefy_good_votes_processed);
|
||||
},
|
||||
VoteImportResult::Equivocation(proof) => {
|
||||
metric_inc!(self, beefy_equivocation_votes);
|
||||
metric_inc!(self.base.metrics, beefy_equivocation_votes);
|
||||
self.report_equivocation(proof)?;
|
||||
},
|
||||
VoteImportResult::Invalid => metric_inc!(self, beefy_invalid_votes),
|
||||
VoteImportResult::Stale => metric_inc!(self, beefy_stale_votes),
|
||||
VoteImportResult::Invalid => metric_inc!(self.base.metrics, beefy_invalid_votes),
|
||||
VoteImportResult::Stale => metric_inc!(self.base.metrics, beefy_stale_votes),
|
||||
};
|
||||
Ok(None)
|
||||
}
|
||||
@@ -625,14 +815,15 @@ where
|
||||
|
||||
// Set new best BEEFY block number.
|
||||
self.persisted_state.set_best_beefy(block_num);
|
||||
crate::aux_schema::write_voter_state(&*self.backend, &self.persisted_state)
|
||||
crate::aux_schema::write_voter_state(&*self.base.backend, &self.persisted_state)
|
||||
.map_err(|e| Error::Backend(e.to_string()))?;
|
||||
|
||||
metric_set!(self, beefy_best_block, block_num);
|
||||
metric_set!(self.base.metrics, beefy_best_block, block_num);
|
||||
|
||||
self.comms.on_demand_justifications.cancel_requests_older_than(block_num);
|
||||
|
||||
if let Err(e) = self
|
||||
.base
|
||||
.backend
|
||||
.blockchain()
|
||||
.expect_block_hash_from_id(&BlockId::Number(block_num))
|
||||
@@ -642,7 +833,8 @@ where
|
||||
.notify(|| Ok::<_, ()>(hash))
|
||||
.expect("forwards closure result; the closure always returns Ok; qed.");
|
||||
|
||||
self.backend
|
||||
self.base
|
||||
.backend
|
||||
.append_justification(hash, (BEEFY_ENGINE_ID, finality_proof.encode()))
|
||||
}) {
|
||||
debug!(
|
||||
@@ -679,12 +871,16 @@ where
|
||||
|
||||
for (num, justification) in justifs_to_process.into_iter() {
|
||||
debug!(target: LOG_TARGET, "🥩 Handle buffered justification for: {:?}.", num);
|
||||
metric_inc!(self, beefy_imported_justifications);
|
||||
metric_inc!(self.base.metrics, beefy_imported_justifications);
|
||||
if let Err(err) = self.finalize(justification) {
|
||||
error!(target: LOG_TARGET, "🥩 Error finalizing block: {}", err);
|
||||
}
|
||||
}
|
||||
metric_set!(self, beefy_buffered_justifications, self.pending_justifications.len());
|
||||
metric_set!(
|
||||
self.base.metrics,
|
||||
beefy_buffered_justifications,
|
||||
self.pending_justifications.len()
|
||||
);
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
@@ -693,7 +889,7 @@ where
|
||||
fn try_to_vote(&mut self) -> Result<(), Error> {
|
||||
// Vote if there's now a new vote target.
|
||||
if let Some(target) = self.voting_oracle().voting_target() {
|
||||
metric_set!(self, beefy_should_vote_on, target);
|
||||
metric_set!(self.base.metrics, beefy_should_vote_on, target);
|
||||
if target > self.persisted_state.best_voted {
|
||||
self.do_vote(target)?;
|
||||
}
|
||||
@@ -713,6 +909,7 @@ where
|
||||
self.persisted_state.voting_oracle.best_grandpa_block_header.clone()
|
||||
} else {
|
||||
let hash = self
|
||||
.base
|
||||
.backend
|
||||
.blockchain()
|
||||
.expect_block_hash_from_id(&BlockId::Number(target_number))
|
||||
@@ -724,7 +921,7 @@ where
|
||||
Error::Backend(err_msg)
|
||||
})?;
|
||||
|
||||
self.backend.blockchain().expect_header(hash).map_err(|err| {
|
||||
self.base.backend.blockchain().expect_header(hash).map_err(|err| {
|
||||
let err_msg = format!(
|
||||
"Couldn't get header for block #{:?} ({:?}) (error: {:?}), skipping vote..",
|
||||
target_number, hash, err
|
||||
@@ -744,7 +941,7 @@ where
|
||||
let rounds = self.persisted_state.voting_oracle.active_rounds_mut()?;
|
||||
let (validators, validator_set_id) = (rounds.validators(), rounds.validator_set_id());
|
||||
|
||||
let authority_id = if let Some(id) = self.key_store.authority_id(validators) {
|
||||
let authority_id = if let Some(id) = self.base.key_store.authority_id(validators) {
|
||||
debug!(target: LOG_TARGET, "🥩 Local authority id: {:?}", id);
|
||||
id
|
||||
} else {
|
||||
@@ -758,7 +955,7 @@ where
|
||||
let commitment = Commitment { payload, block_number: target_number, validator_set_id };
|
||||
let encoded_commitment = commitment.encode();
|
||||
|
||||
let signature = match self.key_store.sign(&authority_id, &encoded_commitment) {
|
||||
let signature = match self.base.key_store.sign(&authority_id, &encoded_commitment) {
|
||||
Ok(sig) => sig,
|
||||
Err(err) => {
|
||||
warn!(target: LOG_TARGET, "🥩 Error signing commitment: {:?}", err);
|
||||
@@ -783,7 +980,7 @@ where
|
||||
.gossip_engine
|
||||
.gossip_message(proofs_topic::<B>(), encoded_proof, true);
|
||||
} else {
|
||||
metric_inc!(self, beefy_votes_sent);
|
||||
metric_inc!(self.base.metrics, beefy_votes_sent);
|
||||
debug!(target: LOG_TARGET, "🥩 Sent vote message: {:?}", vote);
|
||||
let encoded_vote = GossipMessage::<B>::Vote(vote).encode();
|
||||
self.comms.gossip_engine.gossip_message(votes_topic::<B>(), encoded_vote, false);
|
||||
@@ -791,8 +988,8 @@ where
|
||||
|
||||
// Persist state after vote to avoid double voting in case of voter restarts.
|
||||
self.persisted_state.best_voted = target_number;
|
||||
metric_set!(self, beefy_best_voted, target_number);
|
||||
crate::aux_schema::write_voter_state(&*self.backend, &self.persisted_state)
|
||||
metric_set!(self.base.metrics, beefy_best_voted, target_number);
|
||||
crate::aux_schema::write_voter_state(&*self.base.backend, &self.persisted_state)
|
||||
.map_err(|e| Error::Backend(e.to_string()))
|
||||
}
|
||||
|
||||
@@ -966,7 +1163,7 @@ where
|
||||
if !check_equivocation_proof::<_, _, BeefySignatureHasher>(&proof) {
|
||||
debug!(target: LOG_TARGET, "🥩 Skip report for bad equivocation {:?}", proof);
|
||||
return Ok(())
|
||||
} else if let Some(local_id) = self.key_store.authority_id(validators) {
|
||||
} else if let Some(local_id) = self.base.key_store.authority_id(validators) {
|
||||
if offender_id == local_id {
|
||||
debug!(target: LOG_TARGET, "🥩 Skip equivocation report for own equivocation");
|
||||
return Ok(())
|
||||
@@ -975,6 +1172,7 @@ where
|
||||
|
||||
let number = *proof.round_number();
|
||||
let hash = self
|
||||
.base
|
||||
.backend
|
||||
.blockchain()
|
||||
.expect_block_hash_from_id(&BlockId::Number(number))
|
||||
@@ -985,7 +1183,7 @@ where
|
||||
);
|
||||
Error::Backend(err_msg)
|
||||
})?;
|
||||
let runtime_api = self.runtime.runtime_api();
|
||||
let runtime_api = self.base.runtime.runtime_api();
|
||||
// generate key ownership proof at that block
|
||||
let key_owner_proof = match runtime_api
|
||||
.generate_key_ownership_proof(hash, validator_set_id, offender_id)
|
||||
@@ -1002,7 +1200,7 @@ where
|
||||
};
|
||||
|
||||
// submit equivocation report at **best** block
|
||||
let best_block_hash = self.backend.blockchain().info().best_hash;
|
||||
let best_block_hash = self.base.backend.blockchain().info().best_hash;
|
||||
runtime_api
|
||||
.submit_report_equivocation_unsigned_extrinsic(best_block_hash, proof, key_owner_proof)
|
||||
.map_err(Error::RuntimeApi)?;
|
||||
@@ -1028,7 +1226,7 @@ where
|
||||
|
||||
/// Calculate next block number to vote on.
|
||||
///
|
||||
/// Return `None` if there is no voteable target yet.
|
||||
/// Return `None` if there is no votable target yet.
|
||||
fn vote_target<N>(best_grandpa: N, best_beefy: N, session_start: N, min_delta: u32) -> Option<N>
|
||||
where
|
||||
N: AtLeast32Bit + Copy + Debug,
|
||||
@@ -1189,14 +1387,17 @@ pub(crate) mod tests {
|
||||
on_demand_justifications,
|
||||
};
|
||||
BeefyWorker {
|
||||
backend,
|
||||
base: BeefyWorkerBase {
|
||||
backend,
|
||||
runtime: api,
|
||||
key_store: Some(keystore).into(),
|
||||
metrics,
|
||||
_phantom: Default::default(),
|
||||
},
|
||||
payload_provider,
|
||||
runtime: api,
|
||||
key_store: Some(keystore).into(),
|
||||
sync: Arc::new(sync),
|
||||
links,
|
||||
comms,
|
||||
metrics,
|
||||
sync: Arc::new(sync),
|
||||
pending_justifications: BTreeMap::new(),
|
||||
persisted_state,
|
||||
}
|
||||
@@ -1470,19 +1671,19 @@ pub(crate) mod tests {
|
||||
let mut worker = create_beefy_worker(net.peer(0), &keys[0], 1, validator_set.clone());
|
||||
|
||||
// keystore doesn't contain other keys than validators'
|
||||
assert_eq!(worker.verify_validator_set(&1, &validator_set), Ok(()));
|
||||
assert_eq!(worker.base.verify_validator_set(&1, &validator_set), Ok(()));
|
||||
|
||||
// unknown `Bob` key
|
||||
let keys = &[Keyring::Bob];
|
||||
let validator_set = ValidatorSet::new(make_beefy_ids(keys), 0).unwrap();
|
||||
let err_msg = "no authority public key found in store".to_string();
|
||||
let expected = Err(Error::Keystore(err_msg));
|
||||
assert_eq!(worker.verify_validator_set(&1, &validator_set), expected);
|
||||
assert_eq!(worker.base.verify_validator_set(&1, &validator_set), expected);
|
||||
|
||||
// worker has no keystore
|
||||
worker.key_store = None.into();
|
||||
worker.base.key_store = None.into();
|
||||
let expected_err = Err(Error::Keystore("no Keystore".into()));
|
||||
assert_eq!(worker.verify_validator_set(&1, &validator_set), expected_err);
|
||||
assert_eq!(worker.base.verify_validator_set(&1, &validator_set), expected_err);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
@@ -1634,7 +1835,7 @@ pub(crate) mod tests {
|
||||
|
||||
let mut net = BeefyTestNet::new(1);
|
||||
let mut worker = create_beefy_worker(net.peer(0), &keys[0], 1, validator_set.clone());
|
||||
worker.runtime = api_alice.clone();
|
||||
worker.base.runtime = api_alice.clone();
|
||||
|
||||
// let there be a block with num = 1:
|
||||
let _ = net.peer(0).push_blocks(1, false);
|
||||
|
||||
Reference in New Issue
Block a user