mirror of
https://github.com/bitcoinresearchkit/brk.git
synced 2026-07-19 15:08:12 -07:00
global: adding support for safe lengths
This commit is contained in:
@@ -1,9 +1,9 @@
|
||||
use std::ops::Range;
|
||||
|
||||
use brk_error::Result;
|
||||
use brk_indexer::Indexer;
|
||||
use brk_error::{Error, Result};
|
||||
use brk_indexer::{Indexer, Lengths};
|
||||
use brk_oracle::{Config, NUM_BINS, Oracle, START_HEIGHT, bin_to_cents, cents_to_bin};
|
||||
use brk_types::{Cents, Indexes, OutputType, Sats, TxIndex, TxOutIndex};
|
||||
use brk_types::{Cents, OutputType, Sats, TxIndex, TxOutIndex};
|
||||
use tracing::info;
|
||||
use vecdb::{AnyStoredVec, AnyVec, Exit, ReadableVec, StorageMode, VecIndex, WritableVec};
|
||||
|
||||
@@ -15,32 +15,33 @@ impl Vecs {
|
||||
&mut self,
|
||||
indexer: &Indexer,
|
||||
indexes: &indexes::Vecs,
|
||||
starting_indexes: &Indexes,
|
||||
exit: &Exit,
|
||||
) -> Result<()> {
|
||||
self.db.sync_bg_tasks()?;
|
||||
|
||||
self.compute_prices(indexer, starting_indexes, exit)?;
|
||||
let starting_lengths = indexer.safe_lengths();
|
||||
|
||||
self.compute_prices(indexer, exit)?;
|
||||
self.split.open.cents.compute_first(
|
||||
starting_indexes,
|
||||
&starting_lengths,
|
||||
&self.spot.cents.height,
|
||||
indexes,
|
||||
exit,
|
||||
)?;
|
||||
self.split.high.cents.compute_max(
|
||||
starting_indexes,
|
||||
&starting_lengths,
|
||||
&self.spot.cents.height,
|
||||
indexes,
|
||||
exit,
|
||||
)?;
|
||||
self.split.low.cents.compute_min(
|
||||
starting_indexes,
|
||||
&starting_lengths,
|
||||
&self.spot.cents.height,
|
||||
indexes,
|
||||
exit,
|
||||
)?;
|
||||
self.ohlc.cents.compute_from_split(
|
||||
starting_indexes,
|
||||
&starting_lengths,
|
||||
indexes,
|
||||
&self.split.open.cents,
|
||||
&self.split.high.cents,
|
||||
@@ -57,12 +58,9 @@ impl Vecs {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn compute_prices(
|
||||
&mut self,
|
||||
indexer: &Indexer,
|
||||
starting_indexes: &Indexes,
|
||||
exit: &Exit,
|
||||
) -> Result<()> {
|
||||
fn compute_prices(&mut self, indexer: &Indexer, exit: &Exit) -> Result<()> {
|
||||
let starting_height = indexer.safe_lengths().height;
|
||||
|
||||
let source_version =
|
||||
indexer.vecs.outputs.value.version() + indexer.vecs.outputs.output_type.version();
|
||||
self.spot
|
||||
@@ -77,13 +75,8 @@ impl Vecs {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
// Reorg: truncate to starting_indexes
|
||||
let truncate_to = self
|
||||
.spot
|
||||
.cents
|
||||
.height
|
||||
.len()
|
||||
.min(starting_indexes.height.to_usize());
|
||||
// Reorg: truncate to starting_lengths
|
||||
let truncate_to = self.spot.cents.height.len().min(starting_height.to_usize());
|
||||
self.spot
|
||||
.cents
|
||||
.height
|
||||
@@ -157,17 +150,50 @@ impl Vecs {
|
||||
}
|
||||
|
||||
/// Feed a range of blocks from the indexer into an Oracle (skipping coinbase),
|
||||
/// returning per-block ref_bin values.
|
||||
/// returning per-block ref_bin values. Uncapped: derives boundaries from
|
||||
/// raw indexer vec lengths. Use during compute, when the indexer is
|
||||
/// quiescent and `safe_lengths` is still pinned at the pre-pass value.
|
||||
fn feed_blocks<M: StorageMode>(
|
||||
oracle: &mut Oracle,
|
||||
indexer: &Indexer<M>,
|
||||
range: Range<usize>,
|
||||
) -> Vec<f64> {
|
||||
let total_txs = indexer.vecs.transactions.txid.len();
|
||||
let total_outputs = indexer.vecs.outputs.value.len();
|
||||
Self::feed_blocks_inner(oracle, indexer, range, None)
|
||||
}
|
||||
|
||||
/// Capped variant: derives boundaries from `cap` instead of raw vec
|
||||
/// lengths, so concurrent writer pushes past `cap` are invisible.
|
||||
/// Reader paths (live_oracle) use this with the current `safe_lengths`.
|
||||
fn feed_blocks_capped<M: StorageMode>(
|
||||
oracle: &mut Oracle,
|
||||
indexer: &Indexer<M>,
|
||||
range: Range<usize>,
|
||||
cap: &Lengths,
|
||||
) -> Vec<f64> {
|
||||
Self::feed_blocks_inner(oracle, indexer, range, Some(cap))
|
||||
}
|
||||
|
||||
fn feed_blocks_inner<M: StorageMode>(
|
||||
oracle: &mut Oracle,
|
||||
indexer: &Indexer<M>,
|
||||
range: Range<usize>,
|
||||
cap: Option<&Lengths>,
|
||||
) -> Vec<f64> {
|
||||
let (total_txs, total_outputs, height_len) = match cap {
|
||||
Some(c) => (
|
||||
c.tx_index.to_usize(),
|
||||
c.txout_index.to_usize(),
|
||||
c.height.to_usize(),
|
||||
),
|
||||
None => (
|
||||
indexer.vecs.transactions.txid.len(),
|
||||
indexer.vecs.outputs.value.len(),
|
||||
indexer.vecs.transactions.first_tx_index.len(),
|
||||
),
|
||||
};
|
||||
|
||||
// Pre-collect height-indexed data for the range (plus one extra for next-block lookups)
|
||||
let collect_end = (range.end + 1).min(indexer.vecs.transactions.first_tx_index.len());
|
||||
let collect_end = (range.end + 1).min(height_len);
|
||||
let first_tx_indexes: Vec<TxIndex> = indexer
|
||||
.vecs
|
||||
.transactions
|
||||
@@ -243,17 +269,34 @@ impl<M: StorageMode> Vecs<M> {
|
||||
/// window_size blocks already processed. Ready for additional blocks (e.g. mempool).
|
||||
pub fn live_oracle<IM: StorageMode>(&self, indexer: &Indexer<IM>) -> Result<Oracle> {
|
||||
let config = Config::default();
|
||||
let height = indexer.vecs.blocks.timestamp.len();
|
||||
let safe_lengths = indexer.safe_lengths();
|
||||
let height = safe_lengths.height.to_usize();
|
||||
let last_idx = self
|
||||
.spot
|
||||
.cents
|
||||
.height
|
||||
.len()
|
||||
.checked_sub(1)
|
||||
.ok_or(Error::NotFound(
|
||||
"oracle prices not yet computed".to_string(),
|
||||
))?;
|
||||
let last_cents = self
|
||||
.spot
|
||||
.cents
|
||||
.height
|
||||
.collect_one_at(self.spot.cents.height.len() - 1)
|
||||
.unwrap();
|
||||
.collect_one_at(last_idx)
|
||||
.ok_or(Error::NotFound(
|
||||
"oracle prices not yet computed".to_string(),
|
||||
))?;
|
||||
let seed_bin = cents_to_bin(last_cents.inner() as f64);
|
||||
let window_size = config.window_size;
|
||||
let oracle = Oracle::from_checkpoint(seed_bin, config, |o| {
|
||||
Vecs::feed_blocks(o, indexer, height.saturating_sub(window_size)..height);
|
||||
Vecs::feed_blocks_capped(
|
||||
o,
|
||||
indexer,
|
||||
height.saturating_sub(window_size)..height,
|
||||
&safe_lengths,
|
||||
);
|
||||
});
|
||||
|
||||
Ok(oracle)
|
||||
|
||||
@@ -1,8 +1,9 @@
|
||||
use brk_error::Result;
|
||||
use brk_indexer::Lengths;
|
||||
use brk_traversable::Traversable;
|
||||
use brk_types::{
|
||||
Cents, Close, Day1, Day3, Epoch, Halving, High, Hour1, Hour4, Hour12, Indexes, Low, Minute10,
|
||||
Minute30, Month1, Month3, Month6, OHLCCents, Open, Version, Week1, Year1, Year10,
|
||||
Cents, Close, Day1, Day3, Epoch, Halving, High, Hour1, Hour4, Hour12, Low, Minute10, Minute30,
|
||||
Month1, Month3, Month6, OHLCCents, Open, Version, Week1, Year1, Year10,
|
||||
};
|
||||
use derive_more::{Deref, DerefMut};
|
||||
use schemars::JsonSchema;
|
||||
@@ -81,7 +82,7 @@ impl OhlcVecs<OHLCCents> {
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
pub(crate) fn compute_from_split(
|
||||
&mut self,
|
||||
starting_indexes: &Indexes,
|
||||
starting_lengths: &Lengths,
|
||||
indexes: &indexes::Vecs,
|
||||
open: &EagerIndexes<Cents>,
|
||||
high: &EagerIndexes<Cents>,
|
||||
@@ -89,7 +90,7 @@ impl OhlcVecs<OHLCCents> {
|
||||
close: &Resolutions<Cents>,
|
||||
exit: &Exit,
|
||||
) -> Result<()> {
|
||||
let prev_height = starting_indexes.height.decremented().unwrap_or_default();
|
||||
let prev_height = starting_lengths.height.decremented().unwrap_or_default();
|
||||
|
||||
macro_rules! period {
|
||||
($field:ident) => {
|
||||
|
||||
Reference in New Issue
Block a user