use super::*;
impl ConcurrentRadixTreeCompressed {
pub(super) fn find_in_subtree(
start: &SharedNode,
hash: ExternalSequenceBlockHash,
) -> Option<SharedNode> {
let mut queue = VecDeque::new();
start.push_children_into(&mut queue);
while let Some(node) = queue.pop_front() {
if node.contains_edge_hash(hash) {
return Some(node);
}
node.push_children_into(&mut queue);
}
None
}
pub(super) fn resolve_lookup(
&self,
lookup: &mut FxHashMap<WorkerWithDpRank, WorkerLookup>,
worker: WorkerWithDpRank,
hash: ExternalSequenceBlockHash,
direction: LookupRepairDirection,
) -> Option<SharedNode> {
let node = lookup.get(&worker)?.get(&hash)?.clone();
if node.contains_edge_hash(hash) {
return Some(node);
}
let resolved = Self::find_in_subtree(&node, hash)?;
#[cfg(feature = "bench")]
self.bench_metrics
.lookup_repair_scans
.fetch_add(1, Ordering::Relaxed);
self.repair_lookup_for_resolved_node(lookup, hash, &resolved, direction);
Some(resolved)
}
pub(super) fn repair_lookup_for_resolved_node(
&self,
lookup: &mut FxHashMap<WorkerWithDpRank, WorkerLookup>,
hash: ExternalSequenceBlockHash,
resolved: &SharedNode,
direction: LookupRepairDirection,
) {
#[cfg(feature = "bench")]
let mut changed_entries_total = 0u64;
for (&worker, worker_lookup) in lookup.iter_mut() {
let _changed_entries = update_arc_lookup_for_keys(
worker_lookup,
resolved.lookup_hashes_for_worker_repair(worker, hash, direction),
resolved,
);
#[cfg(feature = "bench")]
{
changed_entries_total += _changed_entries as u64;
}
}
#[cfg(feature = "bench")]
self.bench_metrics
.lookup_repair_entries
.fetch_add(changed_entries_total, Ordering::Relaxed);
}
}