use super::{state::*, wire::BLOCK_SYNC_MESSAGE_TYPE_BYTES, *};
#[derive(Clone, Debug)]
pub(crate) struct ReorderBuffer {
blocks: BTreeMap<block::Height, BufferedBlock>,
buffered_bytes: u64,
decoded_attributed_memory_bytes: u64,
}
impl ReorderBuffer {
pub(super) fn new() -> Self {
Self {
blocks: BTreeMap::new(),
buffered_bytes: 0,
decoded_attributed_memory_bytes: 0,
}
}
pub(super) fn buffered_bytes(&self) -> u64 {
self.buffered_bytes
}
pub(super) fn decoded_attributed_memory_bytes(&self) -> u64 {
self.decoded_attributed_memory_bytes
}
#[cfg(test)]
pub(super) fn decoded_attributed_memory_bytes_scanned(&self) -> u64 {
self.blocks
.values()
.map(|buffered| buffered.body.decoded_attributed_memory_size_bytes())
.fold(0u64, u64::saturating_add)
}
pub(super) fn len(&self) -> usize {
self.blocks.len()
}
pub(super) fn contains(&self, height: block::Height) -> bool {
self.blocks.contains_key(&height)
}
pub(super) fn contains_at_or_above(&self, height: block::Height) -> bool {
self.blocks.range(height..).next().is_some()
}
pub(super) fn hash(&self, height: block::Height) -> Option<block::Hash> {
self.blocks.get(&height).map(|buffered| buffered.hash)
}
#[cfg(test)]
pub(super) fn insert(
&mut self,
height: block::Height,
block: Arc<block::Block>,
bytes: u64,
source_peer: ZakuraPeerId,
) -> ReorderInsertResult {
let previous_block_hash = block.header.previous_block_hash;
self.insert_body(
height,
block.hash(),
previous_block_hash,
BufferedBlockBody::from_decoded_block(block, None),
bytes,
source_peer,
)
}
pub(super) fn insert_body(
&mut self,
height: block::Height,
hash: block::Hash,
previous_block_hash: block::Hash,
body: BufferedBlockBody,
bytes: u64,
source_peer: ZakuraPeerId,
) -> ReorderInsertResult {
if self.blocks.contains_key(&height) {
return ReorderInsertResult::Duplicate;
}
let decoded_attributed_memory_size_bytes = body.decoded_attributed_memory_size_bytes();
self.blocks.insert(
height,
BufferedBlock {
hash,
previous_block_hash,
body,
bytes,
source_peer,
},
);
self.buffered_bytes = self.buffered_bytes.saturating_add(bytes);
self.decoded_attributed_memory_bytes = self
.decoded_attributed_memory_bytes
.saturating_add(decoded_attributed_memory_size_bytes);
ReorderInsertResult::Inserted
}
pub(super) fn drain_contiguous_prefix(
&mut self,
verified_block_tip: block::Height,
) -> Vec<DrainedBlock> {
let mut released = Vec::new();
let mut next = match next_height(verified_block_tip) {
Some(next) => next,
None => return released,
};
while let Some(buffered) = self.blocks.remove(&next) {
self.buffered_bytes = self.buffered_bytes.saturating_sub(buffered.bytes);
self.decoded_attributed_memory_bytes = self
.decoded_attributed_memory_bytes
.saturating_sub(buffered.body.decoded_attributed_memory_size_bytes());
released.push(DrainedBlock {
height: next,
hash: buffered.hash,
previous_block_hash: buffered.previous_block_hash,
body: buffered.body,
bytes: buffered.bytes,
source_peer: buffered.source_peer,
});
let Some(after) = next_height(next) else {
break;
};
next = after;
}
released
}
pub(crate) fn clear(&mut self) -> u64 {
self.drop_from(block::Height::MIN)
}
pub(crate) fn drop_through(&mut self, through: block::Height) -> u64 {
let heights: Vec<_> = self
.blocks
.range(..=through)
.map(|(height, _)| *height)
.collect();
let mut released = 0u64;
for height in heights {
if let Some(buffered) = self.blocks.remove(&height) {
self.buffered_bytes = self.buffered_bytes.saturating_sub(buffered.bytes);
self.decoded_attributed_memory_bytes = self
.decoded_attributed_memory_bytes
.saturating_sub(buffered.body.decoded_attributed_memory_size_bytes());
released = released.saturating_add(buffered.bytes);
}
}
released
}
pub(crate) fn drop_from(&mut self, from: block::Height) -> u64 {
let heights: Vec<_> = self
.blocks
.range(from..)
.map(|(height, _)| *height)
.collect();
let mut released = 0u64;
for height in heights {
if let Some(buffered) = self.blocks.remove(&height) {
self.buffered_bytes = self.buffered_bytes.saturating_sub(buffered.bytes);
self.decoded_attributed_memory_bytes = self
.decoded_attributed_memory_bytes
.saturating_sub(buffered.body.decoded_attributed_memory_size_bytes());
released = released.saturating_add(buffered.bytes);
}
}
released
}
}
#[derive(Copy, Clone, Debug, Eq, PartialEq)]
pub(super) enum ReorderInsertResult {
Inserted,
Duplicate,
}
#[derive(Clone, Debug)]
pub(super) struct DrainedBlock {
pub(super) height: block::Height,
pub(super) hash: block::Hash,
pub(super) previous_block_hash: block::Hash,
pub(super) body: BufferedBlockBody,
pub(super) bytes: u64,
pub(super) source_peer: ZakuraPeerId,
}
#[derive(Clone, Debug)]
struct BufferedBlock {
hash: block::Hash,
previous_block_hash: block::Hash,
body: BufferedBlockBody,
bytes: u64,
source_peer: ZakuraPeerId,
}
#[derive(Clone, Debug)]
pub(super) enum BufferedBlockBody {
RawFramePayload(Arc<[u8]>),
Decoded {
block: Arc<block::Block>,
decoded_attributed_memory_size_bytes: u64,
},
DecodedWithRawFramePayload {
block: Arc<block::Block>,
raw_frame_payload: Arc<[u8]>,
decoded_attributed_memory_size_bytes: u64,
},
}
impl BufferedBlockBody {
#[cfg(any(test, feature = "internal-bench"))]
pub(super) fn from_decoded_block(
block: Arc<block::Block>,
raw_frame_payload: Option<Arc<[u8]>>,
) -> Self {
let decoded_attributed_memory_size_bytes = block.attributed_memory_size_bytes();
Self::from_measured_decoded_block(
block,
raw_frame_payload,
decoded_attributed_memory_size_bytes,
)
}
pub(super) fn from_measured_decoded_block(
block: Arc<block::Block>,
raw_frame_payload: Option<Arc<[u8]>>,
decoded_attributed_memory_size_bytes: u64,
) -> Self {
match raw_frame_payload {
Some(raw_frame_payload) => BufferedBlockBody::DecodedWithRawFramePayload {
block,
raw_frame_payload,
decoded_attributed_memory_size_bytes,
},
None => BufferedBlockBody::Decoded {
block,
decoded_attributed_memory_size_bytes,
},
}
}
pub(super) fn decoded_attributed_memory_size_bytes(&self) -> u64 {
match self {
BufferedBlockBody::RawFramePayload(_) => 0,
BufferedBlockBody::Decoded {
decoded_attributed_memory_size_bytes,
..
}
| BufferedBlockBody::DecodedWithRawFramePayload {
decoded_attributed_memory_size_bytes,
..
} => *decoded_attributed_memory_size_bytes,
}
}
pub(super) fn retain_for_backlog(self) -> Self {
match self {
BufferedBlockBody::DecodedWithRawFramePayload {
raw_frame_payload, ..
} => BufferedBlockBody::RawFramePayload(raw_frame_payload),
body => body,
}
}
pub(super) fn retain_for_backlog_in_place(&mut self) {
if let BufferedBlockBody::DecodedWithRawFramePayload {
raw_frame_payload, ..
} = self
{
*self = BufferedBlockBody::RawFramePayload(raw_frame_payload.clone());
}
}
#[cfg(test)]
pub(super) fn is_decoded(&self) -> bool {
!matches!(self, BufferedBlockBody::RawFramePayload(_))
}
pub(super) fn decoded_block(&mut self) -> Arc<block::Block> {
match self {
BufferedBlockBody::Decoded { block, .. }
| BufferedBlockBody::DecodedWithRawFramePayload { block, .. } => block.clone(),
BufferedBlockBody::RawFramePayload(payload) => {
let block = decode_raw_frame_payload(payload);
let decoded_attributed_memory_size_bytes = block.attributed_memory_size_bytes();
let serialized_bytes = payload.len().saturating_sub(BLOCK_SYNC_MESSAGE_TYPE_BYTES);
metrics::histogram!(
"sync.block.body.decoded.attributed_memory_size_bytes",
"stage" => "reorder"
)
.record(decoded_attributed_memory_size_bytes as f64);
if serialized_bytes > 0 {
metrics::histogram!(
"sync.block.body.decoded.to_serialized_ratio",
"stage" => "reorder"
)
.record(decoded_attributed_memory_size_bytes as f64 / serialized_bytes as f64);
}
*self = BufferedBlockBody::DecodedWithRawFramePayload {
block: block.clone(),
raw_frame_payload: payload.clone(),
decoded_attributed_memory_size_bytes,
};
block
}
}
}
}
fn decode_raw_frame_payload(payload: &Arc<[u8]>) -> Arc<block::Block> {
let mut reader = Cursor::new(&payload[BLOCK_SYNC_MESSAGE_TYPE_BYTES..]);
Arc::new(block::Block::zcash_deserialize(&mut reader).expect(
"raw block bytes deserialize because the peer routine decoded them before buffering",
))
}