use std::collections::BTreeMap;
use crate::shard::actor::decode_sequence_key;
use crate::store::NodeStore;
use crate::tree::{Hash, Node, TreeError};
use crate::wal::{Mutation, WalBuffer, WalError};
use super::handle::{ShardError, StreamSeq};
pub(super) fn scan_sequences<S>(
store: &S,
committed_root: Option<Hash>,
buffer: &WalBuffer,
) -> Result<Vec<StreamSeq>, ShardError>
where
S: NodeStore + ?Sized,
{
let mut merged: BTreeMap<Vec<u8>, Vec<u8>> = BTreeMap::new();
if let Some(root) = committed_root {
collect_tree(store, root, &mut merged)?;
}
let buffered = buffer.iter();
for mutation in buffered {
match mutation {
Mutation::Put { key, value } => {
merged.insert(key.clone(), value.clone());
}
Mutation::Delete { key } => {
merged.remove(key);
}
}
}
let mut streams = Vec::new();
for (key, value) in merged {
if let Some(stream_key) = decode_sequence_key(&key) {
let next_seq = decode_seq_value(&value)?;
streams.push((stream_key.to_vec(), next_seq));
}
}
Ok(streams)
}
fn collect_tree<S>(
store: &S,
root: Hash,
out: &mut BTreeMap<Vec<u8>, Vec<u8>>,
) -> Result<(), ShardError>
where
S: NodeStore + ?Sized,
{
let mut stack = vec![root];
while let Some(hash) = stack.pop() {
match load_node(store, hash)? {
Node::Leaf(leaf) => {
for (key, value) in leaf.entries() {
out.insert(key.clone(), value.clone());
}
}
Node::Internal(internal) => {
for (_separator, child) in internal.children() {
stack.push(*child);
}
}
}
}
Ok(())
}
fn load_node<S>(store: &S, hash: Hash) -> Result<Node, ShardError>
where
S: NodeStore + ?Sized,
{
store
.get(&hash)
.map_err(|error| ShardError::Wal(WalError::TreeError(error.to_string())))?
.ok_or(ShardError::Tree(TreeError::MissingNode { hash }))
}
fn decode_seq_value(value: &[u8]) -> Result<u64, ShardError> {
value
.try_into()
.map(u64::from_be_bytes)
.map_err(|_| ShardError::Wal(WalError::TreeError("invalid sequence metadata".to_owned())))
}