global: snapshot

This commit is contained in:
nym21
2026-03-04 23:21:56 +01:00
parent 9e23de4ba1
commit ef0b77baa8
51 changed files with 2109 additions and 5730 deletions
@@ -19,6 +19,8 @@ pub(crate) fn process_received(
empty_addr_count: &mut ByAddressType<u64>,
activity_counts: &mut AddressTypeToActivityCounts,
) {
let mut aggregated: FxHashMap<TypeIndex, (Sats, u32)> = FxHashMap::default();
for (output_type, vec) in received_data.unwrap().into_iter() {
if vec.is_empty() {
continue;
@@ -31,14 +33,13 @@ pub(crate) fn process_received(
// Aggregate receives by address - each address processed exactly once
// Track (total_value, output_count) for correct UTXO counting
let mut aggregated: FxHashMap<TypeIndex, (Sats, u32)> = FxHashMap::default();
for (type_index, value) in vec {
let entry = aggregated.entry(type_index).or_default();
entry.0 += value;
entry.1 += 1;
}
for (type_index, (total_value, output_count)) in aggregated {
for (type_index, (total_value, output_count)) in aggregated.drain() {
let (addr_data, status) = lookup.get_or_create_for_receive(output_type, type_index);
// Track receiving activity - each address in receive aggregation
@@ -40,9 +40,9 @@ pub(crate) fn process_sent(
height_to_timestamp: &[Timestamp],
current_height: Height,
current_timestamp: Timestamp,
seen_senders: &mut ByAddressType<FxHashSet<TypeIndex>>,
) -> Result<()> {
// Track unique senders per address type (simple set, no extra data needed)
let mut seen_senders: ByAddressType<FxHashSet<TypeIndex>> = ByAddressType::default();
seen_senders.values_mut().for_each(|set| set.clear());
for (receive_height, by_type) in sent_data.into_iter() {
let prev_price = height_to_price[receive_height.to_usize()];
@@ -62,15 +62,15 @@ pub(crate) fn process_inputs(
.map(|local_idx| -> Result<_> {
let txindex = txinindex_to_txindex[local_idx];
let prev_height = *txinindex_to_prev_height.get(local_idx).unwrap();
let value = *txinindex_to_value.get(local_idx).unwrap();
let input_type = *txinindex_to_outputtype.get(local_idx).unwrap();
let prev_height = txinindex_to_prev_height[local_idx];
let value = txinindex_to_value[local_idx];
let input_type = txinindex_to_outputtype[local_idx];
if input_type.is_not_address() {
return Ok((prev_height, value, input_type, None));
}
let typeindex = *txinindex_to_typeindex.get(local_idx).unwrap();
let typeindex = txinindex_to_typeindex[local_idx];
// Look up address data
let addr_data_opt = load_uncached_address_data(
@@ -1,6 +1,7 @@
use brk_cohort::ByAddressType;
use brk_error::Result;
use brk_types::{FundedAddressData, Sats, TxIndex, TypeIndex};
use rayon::prelude::*;
use smallvec::SmallVec;
use crate::distribution::{
@@ -47,7 +48,40 @@ pub(crate) fn process_outputs(
) -> Result<OutputsResult> {
let output_count = txoutdata_vec.len();
// Pre-allocate result structures
// Phase 1: Parallel address lookups (mmap reads)
let items: Vec<_> = (0..output_count)
.into_par_iter()
.map(|local_idx| -> Result<_> {
let txoutdata = &txoutdata_vec[local_idx];
let value = txoutdata.value;
let output_type = txoutdata.outputtype;
if output_type.is_not_address() {
return Ok((value, output_type, None));
}
let typeindex = txoutdata.typeindex;
let txindex = txoutindex_to_txindex[local_idx];
let addr_data_opt = load_uncached_address_data(
output_type,
typeindex,
first_addressindexes,
cache,
vr,
any_address_indexes,
addresses_data,
)?;
Ok((
value,
output_type,
Some((typeindex, txindex, value, addr_data_opt)),
))
})
.collect::<Result<Vec<_>>>()?;
// Phase 2: Sequential accumulation
let estimated_per_type = (output_count / 8).max(8);
let mut transacted = Transacted::default();
let mut received_data = AddressTypeToVec::with_capacity(estimated_per_type);
@@ -58,45 +92,26 @@ pub(crate) fn process_outputs(
let mut txindex_vecs =
AddressTypeToTypeIndexMap::<SmallVec<[TxIndex; 4]>>::with_capacity(estimated_per_type);
// Single pass: read from pre-collected vecs and accumulate
for (local_idx, txoutdata) in txoutdata_vec.iter().enumerate() {
let txindex = txoutindex_to_txindex[local_idx];
let value = txoutdata.value;
let output_type = txoutdata.outputtype;
for (value, output_type, addr_info) in items {
transacted.iterate(value, output_type);
if output_type.is_not_address() {
continue;
if let Some((typeindex, txindex, value, addr_data_opt)) = addr_info {
received_data
.get_mut(output_type)
.unwrap()
.push((typeindex, value));
if let Some(addr_data) = addr_data_opt {
address_data.insert_for_type(output_type, typeindex, addr_data);
}
txindex_vecs
.get_mut(output_type)
.unwrap()
.entry(typeindex)
.or_default()
.push(txindex);
}
let typeindex = txoutdata.typeindex;
received_data
.get_mut(output_type)
.unwrap()
.push((typeindex, value));
let addr_data_opt = load_uncached_address_data(
output_type,
typeindex,
first_addressindexes,
cache,
vr,
any_address_indexes,
addresses_data,
)?;
if let Some(addr_data) = addr_data_opt {
address_data.insert_for_type(output_type, typeindex, addr_data);
}
txindex_vecs
.get_mut(output_type)
.unwrap()
.entry(typeindex)
.or_default()
.push(txindex);
}
Ok(OutputsResult {