Approval voting overlay db (#3366)

* node/approval-voting: Introduce Backend trait and Overlaybackend

This commit introduces a Backend trait and attempts to move away
from the Action model via an OverlayBackend as in the ChainSelection
subsystem.

* node/approval-voting: Add WriteOps for StoredBlockRange and BlocksAtHeight

* node/approval-voting: Add load_all_blocks to overlay

* node/approval-voting: Get all module tests to pass.

This commit modifies all tests to ensure tests are passing.

* node/approval-voting: Address oversights in the previous commit

This commit addresses some oversights in the prior commit.

1. Inner errors in backend.write were swallowed
2. One-off write functions removed to avoid useless abstraction
3. Touch-ups in general

* node/approval-voting: Move from TestDB to dyn KeyValueDB

This commit removes the TestDB from tests.rs and replaces it with
an in-memory kvdb.

* node/approval-voting: Address feedback

* node/approval-voting: Add license to ops.rs

* node/approval-voting: Address second-pass feedback

* Add TODO

* node/approval-voting: Bump spec_version

* node/approval-voting: Address final comments.
This commit is contained in:
Lldenaurois
2021-07-08 11:00:57 -04:00
committed by GitHub
parent de445adb6d
commit 7313e485d0
12 changed files with 1445 additions and 1145 deletions
+75 -141
View File
@@ -72,12 +72,15 @@ use time::{slot_number_to_tick, Tick, Clock, ClockExt, SystemClock};
mod approval_checking;
mod approval_db;
mod backend;
mod criteria;
mod import;
mod ops;
mod time;
mod persisted_entries;
use crate::approval_db::v1::Config as DatabaseConfig;
use crate::approval_db::v1::{DbBackend, Config as DatabaseConfig};
use crate::backend::{Backend, OverlayedBackend};
#[cfg(test)]
mod tests;
@@ -330,11 +333,13 @@ impl<C> Subsystem<C> for ApprovalVotingSubsystem
where C: SubsystemContext<Message = ApprovalVotingMessage>
{
fn start(self, ctx: C) -> SpawnedSubsystem {
let future = run::<C>(
let backend = DbBackend::new(self.db.clone(), self.db_config);
let future = run::<DbBackend, C>(
ctx,
self,
Box::new(SystemClock),
Box::new(RealAssignmentCriteria),
backend,
)
.map_err(|e| SubsystemError::with_origin("approval-voting", e))
.boxed();
@@ -461,72 +466,6 @@ impl Wakeups {
}
}
/// A read-only handle to a database.
trait DBReader {
fn load_block_entry(
&self,
block_hash: &Hash,
) -> SubsystemResult<Option<BlockEntry>>;
fn load_candidate_entry(
&self,
candidate_hash: &CandidateHash,
) -> SubsystemResult<Option<CandidateEntry>>;
fn load_all_blocks(&self) -> SubsystemResult<Vec<Hash>>;
}
// This is a submodule to enforce opacity of the inner DB type.
mod approval_db_v1_reader {
use super::{
DBReader, KeyValueDB, Hash, CandidateHash, BlockEntry, CandidateEntry,
SubsystemResult, SubsystemError, DatabaseConfig, approval_db,
};
/// A DB reader that uses the approval-db V1 under the hood.
pub(super) struct ApprovalDBV1Reader<T> {
inner: T,
config: DatabaseConfig,
}
impl<T> ApprovalDBV1Reader<T> {
pub(super) fn new(inner: T, config: DatabaseConfig) -> Self {
ApprovalDBV1Reader {
inner,
config,
}
}
}
impl<'a, T: 'a> DBReader for ApprovalDBV1Reader<T>
where T: std::ops::Deref<Target=(dyn KeyValueDB + 'a)>
{
fn load_block_entry(
&self,
block_hash: &Hash,
) -> SubsystemResult<Option<BlockEntry>> {
approval_db::v1::load_block_entry(&*self.inner, &self.config, block_hash)
.map(|e| e.map(Into::into))
.map_err(|e| SubsystemError::with_origin("approval-voting", e))
}
fn load_candidate_entry(
&self,
candidate_hash: &CandidateHash,
) -> SubsystemResult<Option<CandidateEntry>> {
approval_db::v1::load_candidate_entry(&*self.inner, &self.config, candidate_hash)
.map(|e| e.map(Into::into))
.map_err(|e| SubsystemError::with_origin("approval-voting", e))
}
fn load_all_blocks(&self) -> SubsystemResult<Vec<Hash>> {
approval_db::v1::load_all_blocks(&*self.inner, &self.config)
.map_err(|e| SubsystemError::with_origin("approval-voting", e))
}
}
}
use approval_db_v1_reader::ApprovalDBV1Reader;
struct ApprovalStatus {
required_tranches: RequiredTranches,
tranche_now: DelayTranche,
@@ -636,16 +575,15 @@ impl CurrentlyCheckingSet {
}
}
struct State<T> {
struct State {
session_window: RollingSessionWindow,
keystore: Arc<LocalKeystore>,
slot_duration_millis: u64,
db: T,
clock: Box<dyn Clock + Send + Sync>,
assignment_criteria: Box<dyn AssignmentCriteria + Send + Sync>,
}
impl<T> State<T> {
impl State {
fn session_info(&self, i: SessionIndex) -> Option<&SessionInfo> {
self.session_window.session_info(i)
}
@@ -705,8 +643,6 @@ enum Action {
candidate_hash: CandidateHash,
tick: Tick,
},
WriteBlockEntry(BlockEntry),
WriteCandidateEntry(CandidateHash, CandidateEntry),
LaunchApproval {
candidate_hash: CandidateHash,
indirect_cert: IndirectAssignmentCert,
@@ -723,19 +659,21 @@ enum Action {
Conclude,
}
async fn run<C>(
async fn run<B, C>(
mut ctx: C,
mut subsystem: ApprovalVotingSubsystem,
clock: Box<dyn Clock + Send + Sync>,
assignment_criteria: Box<dyn AssignmentCriteria + Send + Sync>,
mut backend: B,
) -> SubsystemResult<()>
where C: SubsystemContext<Message = ApprovalVotingMessage>
where
C: SubsystemContext<Message = ApprovalVotingMessage>,
B: Backend,
{
let mut state = State {
session_window: RollingSessionWindow::new(APPROVAL_SESSIONS),
keystore: subsystem.keystore,
slot_duration_millis: subsystem.slot_duration_millis,
db: ApprovalDBV1Reader::new(subsystem.db.clone(), subsystem.db_config.clone()),
clock,
assignment_criteria,
};
@@ -746,14 +684,14 @@ async fn run<C>(
let mut last_finalized_height: Option<BlockNumber> = None;
let db_writer = &*subsystem.db;
loop {
let mut overlayed_db = OverlayedBackend::new(&backend);
let actions = futures::select! {
(tick, woken_block, woken_candidate) = wakeups.next(&*state.clock).fuse() => {
subsystem.metrics.on_wakeup();
process_wakeup(
&mut state,
&mut overlayed_db,
woken_block,
woken_candidate,
tick,
@@ -763,9 +701,8 @@ async fn run<C>(
let mut actions = handle_from_overseer(
&mut ctx,
&mut state,
&mut overlayed_db,
&subsystem.metrics,
db_writer,
subsystem.db_config,
next_msg?,
&mut last_finalized_height,
&mut wakeups,
@@ -814,17 +751,23 @@ async fn run<C>(
if handle_actions(
&mut ctx,
&mut state,
&mut overlayed_db,
&subsystem.metrics,
&mut wakeups,
&mut currently_checking_set,
&mut approvals_cache,
db_writer,
subsystem.db_config,
&mut subsystem.mode,
actions,
).await? {
break;
}
if !overlayed_db.is_empty() {
let _timer = subsystem.metrics.time_db_transaction();
let ops = overlayed_db.into_write_ops();
backend.write(ops)?;
}
}
Ok(())
@@ -851,17 +794,15 @@ async fn run<C>(
// returns `true` if any of the actions was a `Conclude` command.
async fn handle_actions(
ctx: &mut impl SubsystemContext,
state: &mut State<impl DBReader>,
state: &mut State,
overlayed_db: &mut OverlayedBackend<'_, impl Backend>,
metrics: &Metrics,
wakeups: &mut Wakeups,
currently_checking_set: &mut CurrentlyCheckingSet,
approvals_cache: &mut lru::LruCache<CandidateHash, ApprovalOutcome>,
db: &dyn KeyValueDB,
db_config: DatabaseConfig,
mode: &mut Mode,
actions: Vec<Action>,
) -> SubsystemResult<bool> {
let mut transaction = approval_db::v1::Transaction::new(db_config);
let mut conclude = false;
let mut actions_iter = actions.into_iter();
@@ -872,15 +813,7 @@ async fn handle_actions(
block_number,
candidate_hash,
tick,
} => {
wakeups.schedule(block_hash, block_number, candidate_hash, tick)
}
Action::WriteBlockEntry(block_entry) => {
transaction.put_block_entry(block_entry.into());
}
Action::WriteCandidateEntry(candidate_hash, candidate_entry) => {
transaction.put_candidate_entry(candidate_hash, candidate_entry.into());
}
} => wakeups.schedule(block_hash, block_number, candidate_hash, tick),
Action::IssueApproval(candidate_hash, approval_request) => {
let mut sender = ctx.sender().clone();
// Note that the IssueApproval action will create additional
@@ -897,6 +830,7 @@ async fn handle_actions(
let next_actions: Vec<Action> = issue_approval(
&mut sender,
state,
overlayed_db,
metrics,
candidate_hash,
approval_request,
@@ -969,9 +903,7 @@ async fn handle_actions(
Action::BecomeActive => {
*mode = Mode::Active;
let messages = distribution_messages_for_activation(
ApprovalDBV1Reader::new(db, db_config)
)?;
let messages = distribution_messages_for_activation(overlayed_db)?;
ctx.send_messages(messages.into_iter().map(Into::into)).await;
}
@@ -979,20 +911,13 @@ async fn handle_actions(
}
}
if !transaction.is_empty() {
let _timer = metrics.time_db_transaction();
transaction.write(db)
.map_err(|e| SubsystemError::with_origin("approval-voting", e))?;
}
Ok(conclude)
}
fn distribution_messages_for_activation<'a>(
db: impl DBReader + 'a,
fn distribution_messages_for_activation(
db: &OverlayedBackend<'_, impl Backend>,
) -> SubsystemResult<Vec<ApprovalDistributionMessage>> {
let all_blocks = db.load_all_blocks()?;
let all_blocks: Vec<Hash> = db.load_all_blocks()?;
let mut approval_meta = Vec::with_capacity(all_blocks.len());
let mut messages = Vec::new();
@@ -1089,10 +1014,9 @@ fn distribution_messages_for_activation<'a>(
// Handle an incoming signal from the overseer. Returns true if execution should conclude.
async fn handle_from_overseer(
ctx: &mut impl SubsystemContext,
state: &mut State<impl DBReader>,
state: &mut State,
db: &mut OverlayedBackend<'_, impl Backend>,
metrics: &Metrics,
db_writer: &dyn KeyValueDB,
db_config: DatabaseConfig,
x: FromOverseer<ApprovalVotingMessage>,
last_finalized_height: &mut Option<BlockNumber>,
wakeups: &mut Wakeups,
@@ -1107,8 +1031,7 @@ async fn handle_from_overseer(
match import::handle_new_head(
ctx,
state,
db_writer,
db_config,
db,
head,
&*last_finalized_height,
).await {
@@ -1162,7 +1085,7 @@ async fn handle_from_overseer(
FromOverseer::Signal(OverseerSignal::BlockFinalized(block_hash, block_number)) => {
*last_finalized_height = Some(block_number);
approval_db::v1::canonicalize(db_writer, &db_config, block_number, block_hash)
crate::ops::canonicalize(db, block_number, block_hash)
.map_err(|e| SubsystemError::with_origin("db", e))?;
wakeups.prune_finalized_wakeups(block_number);
@@ -1175,15 +1098,16 @@ async fn handle_from_overseer(
FromOverseer::Communication { msg } => match msg {
ApprovalVotingMessage::CheckAndImportAssignment(a, claimed_core, res) => {
let (check_outcome, actions)
= check_and_import_assignment(state, a, claimed_core)?;
= check_and_import_assignment(state, db, a, claimed_core)?;
let _ = res.send(check_outcome);
actions
}
ApprovalVotingMessage::CheckAndImportApproval(a, res) => {
check_and_import_approval(state, metrics, a, |r| { let _ = res.send(r); })?.0
check_and_import_approval(state, db, metrics, a, |r| { let _ = res.send(r); })?.0
}
ApprovalVotingMessage::ApprovedAncestor(target, lower_bound, res ) => {
match handle_approved_ancestor(ctx, &state.db, target, lower_bound, wakeups).await {
match handle_approved_ancestor(ctx, db, target, lower_bound, wakeups).await {
Ok(v) => {
let _ = res.send(v);
}
@@ -1203,7 +1127,7 @@ async fn handle_from_overseer(
async fn handle_approved_ancestor(
ctx: &mut impl SubsystemContext,
db: &impl DBReader,
db: &OverlayedBackend<'_, impl Backend>,
target: Hash,
lower_bound: BlockNumber,
wakeups: &Wakeups,
@@ -1491,14 +1415,15 @@ fn schedule_wakeup_action(
}
fn check_and_import_assignment(
state: &State<impl DBReader>,
state: &State,
db: &mut OverlayedBackend<'_, impl Backend>,
assignment: IndirectAssignmentCert,
candidate_index: CandidateIndex,
) -> SubsystemResult<(AssignmentCheckResult, Vec<Action>)> {
const TICK_TOO_FAR_IN_FUTURE: Tick = 20; // 10 seconds.
let tick_now = state.clock.tick_now();
let block_entry = match state.db.load_block_entry(&assignment.block_hash)? {
let block_entry = match db.load_block_entry(&assignment.block_hash)? {
Some(b) => b,
None => return Ok((AssignmentCheckResult::Bad(
AssignmentCheckError::UnknownBlock(assignment.block_hash),
@@ -1523,7 +1448,7 @@ fn check_and_import_assignment(
), Vec::new())), // no candidate at core.
};
let mut candidate_entry = match state.db.load_candidate_entry(&assigned_candidate_hash)? {
let mut candidate_entry = match db.load_candidate_entry(&assigned_candidate_hash)? {
Some(c) => c,
None => {
return Ok((AssignmentCheckResult::Bad(
@@ -1605,13 +1530,14 @@ fn check_and_import_assignment(
}
// We also write the candidate entry as it now contains the new candidate.
actions.push(Action::WriteCandidateEntry(assigned_candidate_hash, candidate_entry));
db.write_candidate_entry(candidate_entry.into());
Ok((res, actions))
}
fn check_and_import_approval<T>(
state: &State<impl DBReader>,
state: &State,
db: &mut OverlayedBackend<'_, impl Backend>,
metrics: &Metrics,
approval: IndirectSignedApprovalVote,
with_response: impl FnOnce(ApprovalCheckResult) -> T,
@@ -1623,7 +1549,7 @@ fn check_and_import_approval<T>(
} }
}
let block_entry = match state.db.load_block_entry(&approval.block_hash)? {
let block_entry = match db.load_block_entry(&approval.block_hash)? {
Some(b) => b,
None => {
respond_early!(ApprovalCheckResult::Bad(
@@ -1666,7 +1592,7 @@ fn check_and_import_approval<T>(
))
}
let candidate_entry = match state.db.load_candidate_entry(&approved_candidate_hash)? {
let candidate_entry = match db.load_candidate_entry(&approved_candidate_hash)? {
Some(c) => c,
None => {
respond_early!(ApprovalCheckResult::Bad(
@@ -1704,6 +1630,7 @@ fn check_and_import_approval<T>(
let actions = import_checked_approval(
state,
db,
&metrics,
block_entry,
approved_candidate_hash,
@@ -1738,7 +1665,8 @@ impl ApprovalSource {
// validator on the candidate and block. This updates the block entry and candidate entry as
// necessary and schedules any further wakeups.
fn import_checked_approval(
state: &State<impl DBReader>,
state: &State,
db: &mut OverlayedBackend<'_, impl Backend>,
metrics: &Metrics,
mut block_entry: BlockEntry,
candidate_hash: CandidateHash,
@@ -1812,7 +1740,7 @@ fn import_checked_approval(
actions.push(Action::NoteApprovedInChainSelection(block_hash));
}
actions.push(Action::WriteBlockEntry(block_entry));
db.write_block_entry(block_entry.into());
}
(is_approved, status)
@@ -1856,11 +1784,11 @@ fn import_checked_approval(
//
// 1. The source is remote, as we don't store anything new in the approval entry.
// 2. The candidate is not newly approved, as we haven't altered the approval entry's
// approved flag with `mark_approved` above.
// approved flag with `mark_approved` above.
// 3. The source had already approved the candidate, as we haven't altered the bitfield.
if !source.is_remote() || newly_approved || !already_approved_by {
// In all other cases, we need to write the candidate entry.
actions.push(Action::WriteCandidateEntry(candidate_hash, candidate_entry));
db.write_candidate_entry(candidate_entry);
}
}
@@ -1904,7 +1832,8 @@ fn should_trigger_assignment(
}
fn process_wakeup(
state: &State<impl DBReader>,
state: &State,
db: &mut OverlayedBackend<'_, impl Backend>,
relay_block: Hash,
candidate_hash: CandidateHash,
expected_tick: Tick,
@@ -1917,8 +1846,8 @@ fn process_wakeup(
.with_candidate(candidate_hash)
.with_stage(jaeger::Stage::ApprovalChecking);
let block_entry = state.db.load_block_entry(&relay_block)?;
let candidate_entry = state.db.load_candidate_entry(&candidate_hash)?;
let block_entry = db.load_block_entry(&relay_block)?;
let candidate_entry = db.load_candidate_entry(&candidate_hash)?;
// If either is not present, we have nothing to wakeup. Might have lost a race with finality
let (block_entry, mut candidate_entry) = match (block_entry, candidate_entry) {
@@ -1981,7 +1910,10 @@ fn process_wakeup(
(should_trigger, approval_entry.backing_group())
};
let (mut actions, maybe_cert) = if should_trigger {
let mut actions = Vec::new();
let candidate_receipt = candidate_entry.candidate_receipt().clone();
let maybe_cert = if should_trigger {
let maybe_cert = {
let approval_entry = candidate_entry.approval_entry_mut(&relay_block)
.expect("should_trigger only true if this fetched earlier; qed");
@@ -1989,11 +1921,11 @@ fn process_wakeup(
approval_entry.trigger_our_assignment(state.clock.tick_now())
};
let actions = vec![Action::WriteCandidateEntry(candidate_hash, candidate_entry.clone())];
db.write_candidate_entry(candidate_entry.clone());
(actions, maybe_cert)
maybe_cert
} else {
(Vec::new(), None)
None
};
if let Some((cert, val_index, tranche)) = maybe_cert {
@@ -2010,7 +1942,7 @@ fn process_wakeup(
tracing::trace!(
target: LOG_TARGET,
?candidate_hash,
para_id = ?candidate_entry.candidate_receipt().descriptor.para_id,
para_id = ?candidate_receipt.descriptor.para_id,
block_hash = ?relay_block,
"Launching approval work.",
);
@@ -2023,7 +1955,7 @@ fn process_wakeup(
relay_block_hash: relay_block,
candidate_index: i as _,
session: block_entry.session(),
candidate: candidate_entry.candidate_receipt().clone(),
candidate: candidate_receipt,
backing_group,
});
}
@@ -2274,12 +2206,13 @@ async fn launch_approval(
// have been done.
fn issue_approval(
ctx: &mut impl SubsystemSender,
state: &State<impl DBReader>,
state: &mut State,
db: &mut OverlayedBackend<'_, impl Backend>,
metrics: &Metrics,
candidate_hash: CandidateHash,
ApprovalVoteRequest { validator_index, block_hash }: ApprovalVoteRequest,
) -> SubsystemResult<Vec<Action>> {
let block_entry = match state.db.load_block_entry(&block_hash)? {
let block_entry = match db.load_block_entry(&block_hash)? {
Some(b) => b,
None => {
// not a cause for alarm - just lost a race with pruning, most likely.
@@ -2337,7 +2270,7 @@ fn issue_approval(
}
};
let candidate_entry = match state.db.load_candidate_entry(&candidate_hash)? {
let candidate_entry = match db.load_candidate_entry(&candidate_hash)? {
Some(c) => c,
None => {
tracing::warn!(
@@ -2397,6 +2330,7 @@ fn issue_approval(
let actions = import_checked_approval(
state,
db,
metrics,
block_entry,
candidate_hash,