Compare commits
10 commits
283c0c3e2c
...
f0a1b2cf13
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
f0a1b2cf13 | ||
|
|
0eef3ecd7c | ||
|
|
d961a5f762 | ||
|
|
8e45d5f216 | ||
|
|
e733bad6ad | ||
|
|
d3fe840596 | ||
|
|
e01232892d | ||
|
|
a9f3b4678e | ||
|
|
29276574e2 | ||
|
|
2ba379828d |
7 changed files with 1622 additions and 201 deletions
|
|
@ -10,6 +10,16 @@ use bitcoincash::TokenID;
|
|||
use serde::Serialize;
|
||||
use sqlx::{Row, SqlitePool};
|
||||
|
||||
pub use self::policy::{GuardParams, Leg, Policy, PoolReserves};
|
||||
|
||||
/// Interpret a database blob as a 32-byte hash, erroring rather than panicking on a
|
||||
/// short or oversized value.
|
||||
pub(crate) fn to_hash32(bytes: &[u8]) -> Result<[u8; 32]> {
|
||||
bytes
|
||||
.try_into()
|
||||
.map_err(|_| anyhow::anyhow!("expected a 32-byte hash, got {} bytes", bytes.len()))
|
||||
}
|
||||
|
||||
#[derive(Debug, Serialize)]
|
||||
pub struct CandlestickData {
|
||||
pub time: i64, // start of the interval
|
||||
|
|
@ -71,14 +81,16 @@ impl PriceInterval {
|
|||
}
|
||||
|
||||
fn aggregate_raw_trades(
|
||||
all_trades: &[(i64, i64, i64)],
|
||||
all_legs: &[Leg],
|
||||
intervals: Vec<PriceInterval>,
|
||||
step_size: i64,
|
||||
mut found_first_trade: bool,
|
||||
mut last_close_price: Option<f64>,
|
||||
policy: &mut Policy,
|
||||
) -> (Vec<CandlestickData>, bool, Option<f64>) {
|
||||
let mut result = Vec::with_capacity(intervals.len());
|
||||
let mut trade_index = 0;
|
||||
let mut leg_index = 0;
|
||||
let mut txid_set = std::collections::HashSet::new();
|
||||
|
||||
for interval in intervals {
|
||||
let interval_start = interval.start;
|
||||
|
|
@ -86,38 +98,51 @@ fn aggregate_raw_trades(
|
|||
let mut pi = PriceInterval::new(interval_start, step_size);
|
||||
let mut first_trade_in_interval = true;
|
||||
|
||||
while trade_index < all_trades.len() {
|
||||
let (ts, vol_sats, vol_tokens) = all_trades[trade_index];
|
||||
if ts < interval_start {
|
||||
trade_index += 1;
|
||||
while leg_index < all_legs.len() {
|
||||
let leg = &all_legs[leg_index];
|
||||
if leg.ts < interval_start {
|
||||
leg_index += 1;
|
||||
continue;
|
||||
}
|
||||
if ts >= interval_end {
|
||||
if leg.ts >= interval_end {
|
||||
break;
|
||||
}
|
||||
|
||||
if vol_tokens != 0 {
|
||||
let price = vol_sats as f64 / vol_tokens as f64;
|
||||
if first_trade_in_interval {
|
||||
pi.open = Some(price);
|
||||
pi.high = price;
|
||||
pi.low = price;
|
||||
first_trade_in_interval = false;
|
||||
}
|
||||
pi.close = Some(price);
|
||||
if price.is_finite() {
|
||||
pi.high = pi.high.max(price);
|
||||
pi.low = pi.low.min(price);
|
||||
let judge = policy.judge(leg);
|
||||
let sats_delta_abs = leg.sats_delta.unsigned_abs() as i64;
|
||||
let token_delta_abs = leg.token_delta.unsigned_abs() as i64;
|
||||
|
||||
// Always count volume, regardless of acceptance.
|
||||
pi.volume_sats += sats_delta_abs;
|
||||
pi.volume_tokens += token_delta_abs;
|
||||
txid_set.insert(leg.txid);
|
||||
|
||||
// Update OHLC only if accepted and priceable.
|
||||
if judge.accepted() {
|
||||
if let Some(price) = judge.price {
|
||||
if first_trade_in_interval {
|
||||
pi.open = Some(price);
|
||||
pi.high = price;
|
||||
pi.low = price;
|
||||
first_trade_in_interval = false;
|
||||
}
|
||||
pi.close = Some(price);
|
||||
if price.is_finite() {
|
||||
pi.high = pi.high.max(price);
|
||||
pi.low = pi.low.min(price);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pi.volume_sats += vol_sats;
|
||||
pi.volume_tokens += vol_tokens;
|
||||
pi.transaction_count += 1;
|
||||
trade_index += 1;
|
||||
policy.apply(leg);
|
||||
leg_index += 1;
|
||||
}
|
||||
|
||||
// Carry forward last close when volume exists but net tokens are zero.
|
||||
// Set tx_count to unique txids in this interval.
|
||||
pi.transaction_count = txid_set.len() as i64;
|
||||
txid_set.clear();
|
||||
|
||||
// Carry forward last close when no accepted prints in this interval.
|
||||
if pi.transaction_count > 0 && (pi.open.is_none() || pi.close.is_none()) {
|
||||
if let Some(prev) = last_close_price {
|
||||
if pi.open.is_none() {
|
||||
|
|
@ -206,33 +231,29 @@ fn fill_ohlcv_candles(
|
|||
(result, found_first_trade, last_close)
|
||||
}
|
||||
|
||||
/// Returns one row per transaction: `(effective_timestamp, volume_sats, volume_tokens)`.
|
||||
///
|
||||
/// Volumes are gross sums of the absolute per-leg deltas, so the derived price
|
||||
/// `volume_sats / volume_tokens` is the volume-weighted average of the prices actually
|
||||
/// executed by that transaction's legs, and is therefore always bounded by the cheapest
|
||||
/// and dearest leg. Summing the *signed* deltas instead lets a multi-pool arbitrage
|
||||
/// transaction — which buys from one pool and sells into others — cancel almost all of
|
||||
/// its token movement and divide real satoshis by a near-zero remainder, fabricating a
|
||||
/// price no leg ever traded at.
|
||||
async fn fetch_raw_trades(
|
||||
/// Returns per-leg data for policy evaluation.
|
||||
/// Each row is: (txid, effective_timestamp, pool, sats_delta, token_delta, sats, token_amount, sequence)
|
||||
async fn fetch_raw_legs(
|
||||
pool: &SqlitePool,
|
||||
token_blob: &[u8],
|
||||
timestamp_start: i64,
|
||||
timestamp_end: i64,
|
||||
) -> Result<Vec<(i64, i64, i64)>> {
|
||||
) -> Result<Vec<Leg>> {
|
||||
let sql = r#"
|
||||
SELECT
|
||||
phe.txid,
|
||||
phe.effective_timestamp,
|
||||
SUM(ABS(phe.sats_delta)) AS volume_sats,
|
||||
SUM(ABS(phe.token_delta)) AS volume_tokens,
|
||||
MIN(phe.sequence) AS min_sequence
|
||||
phe.pool,
|
||||
phe.sats_delta,
|
||||
phe.token_delta,
|
||||
phe.sats,
|
||||
phe.token_amount,
|
||||
phe.sequence
|
||||
FROM pool_history_entry AS phe
|
||||
WHERE phe.token_id = ?
|
||||
AND phe.effective_timestamp >= ?
|
||||
AND phe.effective_timestamp < ?
|
||||
GROUP BY phe.txid, phe.effective_timestamp
|
||||
ORDER BY phe.effective_timestamp ASC, min_sequence ASC;
|
||||
ORDER BY phe.effective_timestamp ASC, phe.sequence ASC;
|
||||
"#;
|
||||
let rows = sqlx::query(sql)
|
||||
.bind(token_blob)
|
||||
|
|
@ -241,38 +262,191 @@ ORDER BY phe.effective_timestamp ASC, min_sequence ASC;
|
|||
.fetch_all(pool)
|
||||
.await?;
|
||||
|
||||
Ok(rows
|
||||
.into_iter()
|
||||
.map(|r| (r.get(0), r.get(1), r.get(2)))
|
||||
.collect())
|
||||
rows.into_iter()
|
||||
.map(|r| {
|
||||
let txid: Vec<u8> = r.get(0);
|
||||
let pool_bytes: Vec<u8> = r.get(2);
|
||||
Ok(Leg {
|
||||
txid: to_hash32(&txid)?,
|
||||
pool: to_hash32(&pool_bytes)?,
|
||||
ts: r.get(1),
|
||||
sequence: r.get(7),
|
||||
sats_delta: r.get(3),
|
||||
token_delta: r.get(4),
|
||||
sats: r.get(5),
|
||||
token_amount: r.get(6),
|
||||
})
|
||||
})
|
||||
.collect()
|
||||
}
|
||||
|
||||
/// Returns the close price of the most recent priceable trade strictly before
|
||||
/// `timestamp_end`, using the same per-tx aggregation as `fetch_raw_trades`.
|
||||
/// Returns `None` when no prior trade exists (new token, no history).
|
||||
/// Pool reserves for `token_blob` as they stood immediately before `timestamp`.
|
||||
///
|
||||
/// Seeding a [`Policy`] with this is what makes its verdicts independent of the caller's
|
||||
/// window: the raw path, a rebuild batch, and a 1-week versus 1-month request all start
|
||||
/// from the same chain state rather than from whichever legs happened to be in range.
|
||||
///
|
||||
/// Enumerates the token's pools and seeks each one's latest history row, rather than
|
||||
/// scanning every row the token ever produced — a busy token has hundreds of thousands,
|
||||
/// and this runs once per query.
|
||||
///
|
||||
/// Pools already withdrawn at `timestamp` are excluded. A withdrawal leaves no history
|
||||
/// row behind, so a drained pool's last entry still shows its full pre-withdrawal
|
||||
/// reserves; counting those would let a pool that no longer holds anything go on voting
|
||||
/// the reference. Pools withdrawn *after* `timestamp` are kept, since they were live then
|
||||
/// — the same rule `poolvisitor` applies.
|
||||
///
|
||||
/// `confirmed_only` must match whatever the caller's leg query does. `rebuild_range` folds
|
||||
/// only confirmed legs, so seeding it from reserves that include mempool rows would make
|
||||
/// materialised output a function of mempool contents at rebuild time — and `INSERT OR
|
||||
/// IGNORE` would then freeze whichever run happened first. The live query paths do not
|
||||
/// filter their legs, so they pass `false`.
|
||||
pub(crate) async fn fetch_reserve_snapshot(
|
||||
pool: &SqlitePool,
|
||||
token_blob: &[u8],
|
||||
timestamp: i64,
|
||||
confirmed_only: bool,
|
||||
) -> Result<Vec<PoolReserves>> {
|
||||
// Interpolated, not bound: a constant fragment keeps the query planner able to use
|
||||
// idx_pool_history_entry_pool_timestamp instead of evaluating a flag per row.
|
||||
let confirmed_join = if confirmed_only {
|
||||
"JOIN tx AS prior_tx ON prior_tx.txid = prior.txid AND prior_tx.blockhash IS NOT NULL"
|
||||
} else {
|
||||
""
|
||||
};
|
||||
|
||||
let sql = format!(
|
||||
r#"
|
||||
SELECT phe.pool, phe.sequence, phe.sats, phe.token_amount
|
||||
FROM pool AS p
|
||||
JOIN pool_history_entry AS phe ON phe.utxo = (
|
||||
SELECT prior.utxo
|
||||
FROM pool_history_entry AS prior
|
||||
{confirmed_join}
|
||||
WHERE prior.pool = p.creation_utxo
|
||||
AND prior.effective_timestamp < ?
|
||||
ORDER BY prior.effective_timestamp DESC, prior.sequence DESC
|
||||
LIMIT 1
|
||||
)
|
||||
WHERE p.token_id = ?
|
||||
AND (p.withdrawn_in_utxo IS NULL OR (
|
||||
SELECT t.effective_timestamp
|
||||
FROM utxo_spending AS us
|
||||
JOIN tx AS t ON us.txid = t.txid
|
||||
WHERE us.spent_utxo_hash = p.withdrawn_in_utxo
|
||||
) >= ?);
|
||||
"#
|
||||
);
|
||||
let rows = sqlx::query(&sql)
|
||||
.bind(timestamp)
|
||||
.bind(token_blob)
|
||||
.bind(timestamp)
|
||||
.fetch_all(pool)
|
||||
.await?;
|
||||
|
||||
rows.into_iter()
|
||||
.map(|r| {
|
||||
let pool_bytes: Vec<u8> = r.get(0);
|
||||
Ok(PoolReserves {
|
||||
pool: to_hash32(&pool_bytes)?,
|
||||
sequence: r.get(1),
|
||||
sats: r.get(2),
|
||||
token_amount: r.get(3),
|
||||
})
|
||||
})
|
||||
.collect()
|
||||
}
|
||||
|
||||
/// A price to carry into a window that opens on silence, so that requesting 1W and 1M
|
||||
/// yields the same candles over the period they share.
|
||||
///
|
||||
/// Prefers the most recent accepted print strictly before `timestamp_end`. Falls back to
|
||||
/// the reference implied by pool reserves when the token's last activity was entirely
|
||||
/// muted or unpriceable — reserves still say what the market was, and a carried price
|
||||
/// beats a hole in the chart.
|
||||
///
|
||||
/// Returns `None` only when the token has no history at all before `timestamp_end`.
|
||||
///
|
||||
/// The lookback is deliberately unbounded. An earlier revision capped the scan at 24
|
||||
/// hours, which silently dropped the seed for any token trading less often than daily.
|
||||
async fn fetch_last_close_before(
|
||||
pool: &SqlitePool,
|
||||
token_blob: &[u8],
|
||||
timestamp_end: i64,
|
||||
) -> Result<Option<f64>> {
|
||||
let sql = r#"
|
||||
SELECT
|
||||
CAST(SUM(ABS(phe.sats_delta)) AS REAL) / CAST(SUM(ABS(phe.token_delta)) AS REAL) AS close_price
|
||||
FROM pool_history_entry AS phe
|
||||
WHERE phe.token_id = ?
|
||||
AND phe.effective_timestamp < ?
|
||||
GROUP BY phe.txid, phe.effective_timestamp
|
||||
HAVING SUM(ABS(phe.token_delta)) != 0
|
||||
ORDER BY phe.effective_timestamp DESC, MIN(phe.sequence) DESC
|
||||
LIMIT 1
|
||||
"#;
|
||||
let row = sqlx::query(sql)
|
||||
.bind(token_blob)
|
||||
.bind(timestamp_end)
|
||||
.fetch_optional(pool)
|
||||
.await?;
|
||||
// Locate the token's most recent activity first. Trusting `ohlcv_1h` before knowing
|
||||
// this would let a stale bucket shadow newer legs: the table is structurally at least
|
||||
// three hours behind the tip, and is empty for the whole of a post-version-bump
|
||||
// rebuild, so "newest materialised bucket" is routinely far older than the real last
|
||||
// print.
|
||||
let last_activity: Option<i64> = sqlx::query_scalar(
|
||||
"SELECT MAX(effective_timestamp) FROM pool_history_entry
|
||||
WHERE token_id = ? AND effective_timestamp < ?",
|
||||
)
|
||||
.bind(token_blob)
|
||||
.bind(timestamp_end)
|
||||
.fetch_optional(pool)
|
||||
.await?
|
||||
.flatten();
|
||||
|
||||
Ok(row.map(|r| r.get::<f64, _>(0)))
|
||||
let Some(last_activity) = last_activity else {
|
||||
return Ok(None);
|
||||
};
|
||||
|
||||
let scan_start = (last_activity / 3600) * 3600;
|
||||
|
||||
// The last accepted print lives in the hour containing that activity. If that exact
|
||||
// hour is materialised and closes at or before the cutoff, its close is the answer.
|
||||
//
|
||||
// The `scan_start + 3600 <= timestamp_end` test matters: a bucket straddling the cutoff
|
||||
// closes on a print that happened *after* the instant being asked about, which an
|
||||
// earlier revision leaked for any non-hour-aligned `timestamp_end`.
|
||||
if scan_start + 3600 <= timestamp_end {
|
||||
let exact: Option<f64> =
|
||||
sqlx::query_scalar("SELECT close FROM ohlcv_1h WHERE token_id = ? AND bucket_ts = ?")
|
||||
.bind(token_blob)
|
||||
.bind(scan_start)
|
||||
.fetch_optional(pool)
|
||||
.await?;
|
||||
|
||||
if exact.is_some() {
|
||||
return Ok(exact);
|
||||
}
|
||||
}
|
||||
|
||||
// Replay that hour through the guard. Only one hour is needed, because the policy is
|
||||
// seeded from reserves rather than rebuilt by folding all of history.
|
||||
let snapshot = fetch_reserve_snapshot(pool, token_blob, scan_start, false).await?;
|
||||
let mut policy = Policy::seeded(GuardParams::default(), &snapshot);
|
||||
|
||||
let legs = fetch_raw_legs(pool, token_blob, scan_start, timestamp_end).await?;
|
||||
let mut last_accepted_price = None;
|
||||
for leg in &legs {
|
||||
let judge = policy.judge(leg);
|
||||
if judge.accepted() {
|
||||
last_accepted_price = judge.price;
|
||||
}
|
||||
policy.apply(leg);
|
||||
}
|
||||
if last_accepted_price.is_some() {
|
||||
return Ok(last_accepted_price);
|
||||
}
|
||||
|
||||
// That hour was entirely muted or unpriceable. Reach further back through the
|
||||
// materialised table — unbounded, but only now that it cannot shadow newer legs.
|
||||
let earlier: Option<f64> = sqlx::query_scalar(
|
||||
"SELECT close FROM ohlcv_1h
|
||||
WHERE token_id = ? AND bucket_ts + 3600 <= ?
|
||||
ORDER BY bucket_ts DESC LIMIT 1",
|
||||
)
|
||||
.bind(token_blob)
|
||||
.bind(timestamp_end)
|
||||
.fetch_optional(pool)
|
||||
.await?;
|
||||
|
||||
// Last resort: the price implied by current reserves. Not a executed print, but a
|
||||
// carried price beats a hole in the chart.
|
||||
Ok(earlier.or_else(|| policy.reference()))
|
||||
}
|
||||
|
||||
/// `ohlcv_materialized_end`: exclusive upper bound of what is in `ohlcv_1h`.
|
||||
|
|
@ -290,6 +464,7 @@ pub async fn candlesticks(
|
|||
}
|
||||
|
||||
let token_blob = display_hex_to_blob::<TokenID>(token_id)?;
|
||||
let params = GuardParams::default();
|
||||
|
||||
// Seed gap-fill with the last known close price before this window so that
|
||||
// switching between timeframes (e.g. 1W vs 1M) produces consistent prices
|
||||
|
|
@ -319,7 +494,14 @@ pub async fn candlesticks(
|
|||
|
||||
if ohlcv_end < timestamp_end {
|
||||
// Tail: query raw for [ohlcv_end, timestamp_end) and append.
|
||||
let raw_trades = fetch_raw_trades(pool, &token_blob, ohlcv_end, timestamp_end).await?;
|
||||
//
|
||||
// The policy is seeded at `ohlcv_end`, not at `timestamp_start`: the buckets
|
||||
// before it came from `ohlcv_1h` and never passed through this policy, so a
|
||||
// policy built from the window's start would enter the tail believing the
|
||||
// token had no history and wave its first leg through unjudged.
|
||||
let snapshot = fetch_reserve_snapshot(pool, &token_blob, ohlcv_end, false).await?;
|
||||
let mut policy = Policy::seeded(params, &snapshot);
|
||||
let raw_legs = fetch_raw_legs(pool, &token_blob, ohlcv_end, timestamp_end).await?;
|
||||
|
||||
let mut tail_intervals = Vec::new();
|
||||
let mut t = ohlcv_end;
|
||||
|
|
@ -329,11 +511,12 @@ pub async fn candlesticks(
|
|||
}
|
||||
|
||||
let (tail, _, _) = aggregate_raw_trades(
|
||||
&raw_trades,
|
||||
&raw_legs,
|
||||
tail_intervals,
|
||||
step_size,
|
||||
found_first,
|
||||
last_close,
|
||||
&mut policy,
|
||||
);
|
||||
result.extend(tail);
|
||||
}
|
||||
|
|
@ -349,12 +532,22 @@ pub async fn candlesticks(
|
|||
current_start += step_size;
|
||||
}
|
||||
|
||||
let all_trades = fetch_raw_trades(pool, &token_blob, timestamp_start, timestamp_end).await?;
|
||||
let snapshot = fetch_reserve_snapshot(pool, &token_blob, timestamp_start, false).await?;
|
||||
let mut policy = Policy::seeded(params, &snapshot);
|
||||
let all_legs = fetch_raw_legs(pool, &token_blob, timestamp_start, timestamp_end).await?;
|
||||
|
||||
let (result, _, _) =
|
||||
aggregate_raw_trades(&all_trades, intervals, step_size, seed_found, seed_close);
|
||||
let (result, _, _) = aggregate_raw_trades(
|
||||
&all_legs,
|
||||
intervals,
|
||||
step_size,
|
||||
seed_found,
|
||||
seed_close,
|
||||
&mut policy,
|
||||
);
|
||||
Ok(result)
|
||||
}
|
||||
|
||||
pub mod policy;
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests;
|
||||
|
|
|
|||
412
src/db/cauldron/candlestick/policy.rs
Normal file
412
src/db/cauldron/candlestick/policy.rs
Normal file
|
|
@ -0,0 +1,412 @@
|
|||
// Copyright (C) 2024-2026 Whiterun LLC
|
||||
//
|
||||
// This software is licensed under the GNU Affero General Public License (AGPL), version 3.0 or later.
|
||||
// A copy of the license can be found in the LICENSE file or at https://www.gnu.org/licenses/agpl-3.0.html
|
||||
|
||||
use std::cmp::Ordering;
|
||||
use std::collections::HashMap;
|
||||
|
||||
/// Pools holding fewer tokens than this are dust and do not vote on the reference.
|
||||
const MIN_TOKEN_RESERVE: u64 = 10;
|
||||
|
||||
/// A leg of a transaction: one pool's change in a single transaction.
|
||||
#[derive(Clone, Debug)]
|
||||
pub struct Leg {
|
||||
pub txid: [u8; 32],
|
||||
pub pool: [u8; 32],
|
||||
pub ts: i64,
|
||||
pub sequence: i64,
|
||||
pub sats_delta: i64, // sats change at this pool
|
||||
pub token_delta: i64, // token change at this pool
|
||||
pub sats: u64, // post-change sats (reserves)
|
||||
pub token_amount: u64, // post-change token_amount (reserves)
|
||||
}
|
||||
|
||||
/// One pool's reserves at an instant, used to seed a [`Policy`].
|
||||
///
|
||||
/// Seeding is what makes the policy's verdicts a function of the chain rather than of
|
||||
/// where the caller happened to start reading: a request for one hour and a request for
|
||||
/// one month both begin from the same reserve state and so agree on the overlap.
|
||||
#[derive(Clone, Debug)]
|
||||
pub struct PoolReserves {
|
||||
pub pool: [u8; 32],
|
||||
pub sequence: i64,
|
||||
pub sats: u64,
|
||||
pub token_amount: u64,
|
||||
}
|
||||
|
||||
/// Why a leg was accepted or muted.
|
||||
///
|
||||
/// Deliberately allocation-free: `judge` runs once per leg, and a full `ohlcv_1h` rebuild
|
||||
/// walks millions of them.
|
||||
#[derive(Clone, Copy, Debug, PartialEq)]
|
||||
pub enum Verdict {
|
||||
/// No token moved, or no sats moved, so the leg has no price.
|
||||
Unpriceable,
|
||||
/// No pool held usable reserves, so there is nothing to deviate from.
|
||||
NoReference,
|
||||
/// Within `max_deviation_factor` of the reference.
|
||||
WithinBand { dev: f64 },
|
||||
/// Outside the band, but the pool holds enough qualification credit to print anyway.
|
||||
CreditExempt {
|
||||
dev: f64,
|
||||
credit: u64,
|
||||
largest_credit: u64,
|
||||
},
|
||||
/// Outside the band with too little credit: counted in volume, excluded from OHLC.
|
||||
Muted {
|
||||
dev: f64,
|
||||
credit: u64,
|
||||
largest_credit: u64,
|
||||
},
|
||||
}
|
||||
|
||||
impl Verdict {
|
||||
pub fn accepted(self) -> bool {
|
||||
matches!(
|
||||
self,
|
||||
Verdict::NoReference | Verdict::WithinBand { .. } | Verdict::CreditExempt { .. }
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
/// Outcome of judging a leg.
|
||||
#[derive(Clone, Copy, Debug)]
|
||||
pub struct JudgeResult {
|
||||
/// The leg's executed price, or `None` when it has none.
|
||||
pub price: Option<f64>,
|
||||
pub verdict: Verdict,
|
||||
}
|
||||
|
||||
impl JudgeResult {
|
||||
pub fn accepted(&self) -> bool {
|
||||
self.verdict.accepted()
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Debug)]
|
||||
pub struct GuardParams {
|
||||
/// `F`: a leg further than this factor from the reference needs credit to print.
|
||||
pub max_deviation_factor: f64,
|
||||
/// `q`: credit needed to qualify, as a fraction of the largest pool's credit.
|
||||
pub min_share_fraction: f64,
|
||||
}
|
||||
|
||||
impl Default for GuardParams {
|
||||
/// Calibrated against a production replay: `F = 5` sits well above the worst real
|
||||
/// non-qualifying print and well below the weakest observed attack print.
|
||||
///
|
||||
/// Single source of truth for every call site until these are threaded through
|
||||
/// `configure_me`; changing them requires an `OHLCV_VERSION` bump.
|
||||
fn default() -> Self {
|
||||
Self {
|
||||
max_deviation_factor: 5.0,
|
||||
min_share_fraction: 0.05,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Per-pool state: reserves only.
|
||||
///
|
||||
/// Qualification credit is *derived* from these on demand (see [`credit_from`]) rather than
|
||||
/// stored, which is what keeps a policy's verdicts a function of the reserves it holds
|
||||
/// rather than of the order it learned them in.
|
||||
#[derive(Clone, Debug)]
|
||||
struct PoolState {
|
||||
sequence: i64,
|
||||
sats: u64,
|
||||
token_amount: u64,
|
||||
}
|
||||
|
||||
/// The guard: judges each leg against a reference price derived from pool reserves.
|
||||
///
|
||||
/// Reserves are advanced by [`Policy::apply`] as legs stream past, but the starting
|
||||
/// state comes from a snapshot (see [`Policy::seeded`]) rather than from whatever legs
|
||||
/// happen to fall inside the caller's window. That is what keeps a bucket's verdict
|
||||
/// identical whether it was reached through the raw path, through a rebuild batch, or
|
||||
/// through a request for a different timeframe.
|
||||
pub struct Policy {
|
||||
pools: HashMap<[u8; 32], PoolState>,
|
||||
reference: Option<f64>,
|
||||
params: GuardParams,
|
||||
/// Reused across `recompute_reference` calls so the hot path does not allocate.
|
||||
scratch: Vec<(f64, u64, u64)>,
|
||||
}
|
||||
|
||||
impl Policy {
|
||||
/// A policy with no prior knowledge. Only correct when the caller genuinely starts
|
||||
/// at the token's first trade; otherwise use [`Policy::seeded`].
|
||||
pub fn new(params: GuardParams) -> Self {
|
||||
Self::seeded(params, &[])
|
||||
}
|
||||
|
||||
/// A policy primed with pool reserves as of some instant.
|
||||
///
|
||||
/// Nothing beyond reserves needs priming: credit is derived from them (see
|
||||
/// [`credit_from`]), so a seeded policy and one that folded its way to the same
|
||||
/// reserves are indistinguishable by construction.
|
||||
pub fn seeded(params: GuardParams, snapshot: &[PoolReserves]) -> Self {
|
||||
let mut policy = Self {
|
||||
pools: HashMap::with_capacity(snapshot.len()),
|
||||
reference: None,
|
||||
params,
|
||||
scratch: Vec::with_capacity(snapshot.len()),
|
||||
};
|
||||
|
||||
for r in snapshot {
|
||||
policy.pools.insert(
|
||||
r.pool,
|
||||
PoolState {
|
||||
sequence: r.sequence,
|
||||
sats: r.sats,
|
||||
token_amount: r.token_amount,
|
||||
},
|
||||
);
|
||||
}
|
||||
policy.recompute_reference();
|
||||
policy
|
||||
}
|
||||
|
||||
/// Judge a leg. Pure: call [`Policy::apply`] afterwards to advance state.
|
||||
pub fn judge(&self, leg: &Leg) -> JudgeResult {
|
||||
// A leg that moved no tokens, or bought them for nothing, has no price to print.
|
||||
if leg.token_delta == 0 || leg.sats_delta == 0 {
|
||||
return JudgeResult {
|
||||
price: None,
|
||||
verdict: Verdict::Unpriceable,
|
||||
};
|
||||
}
|
||||
|
||||
let price = leg.sats_delta.unsigned_abs() as f64 / leg.token_delta.unsigned_abs() as f64;
|
||||
|
||||
let Some(reference) = self.reference else {
|
||||
return JudgeResult {
|
||||
price: Some(price),
|
||||
verdict: Verdict::NoReference,
|
||||
};
|
||||
};
|
||||
|
||||
let dev = deviation(price, reference);
|
||||
if dev <= self.params.max_deviation_factor {
|
||||
return JudgeResult {
|
||||
price: Some(price),
|
||||
verdict: Verdict::WithinBand { dev },
|
||||
};
|
||||
}
|
||||
|
||||
// Both sides derived from present reserves, so this scan is the whole of tier 2's
|
||||
// state. It is O(pools), but only reached by off-band legs — `WithinBand` returns
|
||||
// above — so the common path stays a single comparison.
|
||||
let f = self.params.max_deviation_factor;
|
||||
let credit = self
|
||||
.pools
|
||||
.get(&leg.pool)
|
||||
.map_or(0, |s| credit_from(s.sats, s.token_amount, reference, f));
|
||||
let largest_credit = self
|
||||
.pools
|
||||
.values()
|
||||
.map(|s| credit_from(s.sats, s.token_amount, reference, f))
|
||||
.max()
|
||||
.unwrap_or(0);
|
||||
// `ceil().max(1)` keeps a pool with no usable depth (credit 0) from qualifying just
|
||||
// because the largest credit is small.
|
||||
let threshold =
|
||||
((largest_credit as f64 * self.params.min_share_fraction).ceil() as u64).max(1);
|
||||
|
||||
let verdict = if credit >= threshold {
|
||||
Verdict::CreditExempt {
|
||||
dev,
|
||||
credit,
|
||||
largest_credit,
|
||||
}
|
||||
} else {
|
||||
Verdict::Muted {
|
||||
dev,
|
||||
credit,
|
||||
largest_credit,
|
||||
}
|
||||
};
|
||||
|
||||
JudgeResult {
|
||||
price: Some(price),
|
||||
verdict,
|
||||
}
|
||||
}
|
||||
|
||||
/// Advance state past a leg. Must be called for every leg, accepted or not: a muted
|
||||
/// leg still moved real reserves.
|
||||
///
|
||||
/// Takes no verdict, because none is needed — credit follows the reserves this records.
|
||||
pub fn apply(&mut self, leg: &Leg) {
|
||||
// `i64::MIN` so a pool's very first leg clears the monotonicity guard below.
|
||||
let state = self.pools.entry(leg.pool).or_insert(PoolState {
|
||||
sequence: i64::MIN,
|
||||
sats: 0,
|
||||
token_amount: 0,
|
||||
});
|
||||
|
||||
// Strict `>`: the `i64::MIN` sentinel above is what lets a pool's first leg through,
|
||||
// so there is no reason to also relax the monotonicity rule for later ones.
|
||||
if leg.sequence > state.sequence {
|
||||
state.sequence = leg.sequence;
|
||||
state.sats = leg.sats;
|
||||
state.token_amount = leg.token_amount;
|
||||
}
|
||||
|
||||
self.recompute_reference();
|
||||
}
|
||||
|
||||
/// Current reference price.
|
||||
pub fn reference(&self) -> Option<f64> {
|
||||
self.reference
|
||||
}
|
||||
|
||||
/// Qualification credit held by a pool. Exposed for tests and diagnostics.
|
||||
#[allow(dead_code)]
|
||||
pub fn credit_of(&self, pool: &[u8; 32]) -> u64 {
|
||||
let Some(reference) = self.reference else {
|
||||
return 0;
|
||||
};
|
||||
self.pools.get(pool).map_or(0, |s| {
|
||||
credit_from(
|
||||
s.sats,
|
||||
s.token_amount,
|
||||
reference,
|
||||
self.params.max_deviation_factor,
|
||||
)
|
||||
})
|
||||
}
|
||||
|
||||
fn recompute_reference(&mut self) {
|
||||
// Detach the scratch buffer so `self.pools` can be borrowed while filling it.
|
||||
let mut scratch = std::mem::take(&mut self.scratch);
|
||||
scratch.clear();
|
||||
scratch.extend(self.pools.values().filter_map(|s| {
|
||||
spot_ratio(s.sats, s.token_amount).map(|spot| (spot, s.sats, s.token_amount))
|
||||
}));
|
||||
self.reference = weighted_median_reference(&mut scratch);
|
||||
self.scratch = scratch;
|
||||
}
|
||||
}
|
||||
|
||||
/// Qualification credit: a pool's min-depth, but only while its own spot sits within `F` of
|
||||
/// the reference.
|
||||
///
|
||||
/// Stateless by construction, and that is the point. An earlier revision stored "min-depth
|
||||
/// at the pool's last accepted print", which made credit path-dependent — a reserve
|
||||
/// snapshot cannot recover it, so [`Policy::seeded`] had to approximate, and a seeded policy
|
||||
/// returned `CreditExempt` where a folded one returned `Muted` on the very same leg. That
|
||||
/// broke the invariant the whole seeding design exists to hold, and because `ohlcv_1h` is
|
||||
/// written with `INSERT OR IGNORE`, whichever path ran first would freeze its answer in
|
||||
/// permanently.
|
||||
///
|
||||
/// Little security was given up. The path-dependent rule was a speed bump rather than a
|
||||
/// wall: an attacker needed one cheap in-band trade to convert a fresh pool into a credited
|
||||
/// one. What remains is a pure depth gate — printing far off-reference requires a pool
|
||||
/// holding a real share of the token's liquidity, which is the part that actually costs an
|
||||
/// attacker money. A pool that has been walked off-market is still excluded outright, so a
|
||||
/// snapshot cannot launder an attacker's pool into qualification.
|
||||
fn credit_from(sats: u64, token_amount: u64, reference: f64, f: f64) -> u64 {
|
||||
let Some(spot) = spot_ratio(sats, token_amount) else {
|
||||
return 0;
|
||||
};
|
||||
if deviation(spot, reference) > f {
|
||||
return 0;
|
||||
}
|
||||
min_depth(sats, token_amount, reference)
|
||||
}
|
||||
|
||||
/// A pool's spot price, or `None` if it is too small or one-sided to quote one.
|
||||
fn spot_ratio(sats: u64, token_amount: u64) -> Option<f64> {
|
||||
if token_amount < MIN_TOKEN_RESERVE || sats == 0 {
|
||||
return None;
|
||||
}
|
||||
Some(sats as f64 / token_amount as f64)
|
||||
}
|
||||
|
||||
/// How far apart two prices are, as a factor >= 1 in whichever direction.
|
||||
fn deviation(price: f64, reference: f64) -> f64 {
|
||||
// Finiteness is tested first so the `<=` comparisons below are total: NaN fails
|
||||
// `is_finite` and returns here, which is why plain `<=` is safe rather than needing
|
||||
// `!(price > 0.0)` to catch it.
|
||||
if !price.is_finite() || !reference.is_finite() || price <= 0.0 || reference <= 0.0 {
|
||||
return f64::INFINITY;
|
||||
}
|
||||
(price / reference).max(reference / price)
|
||||
}
|
||||
|
||||
/// `d = min(S, T * R)`: the side of the pool an attacker would have to actually fund.
|
||||
///
|
||||
/// Valuing the token side at the reference is what stops a pool stuffed with worthless
|
||||
/// tokens from out-voting a pool holding real sats.
|
||||
fn min_depth(sats: u64, token_amount: u64, reference: f64) -> u64 {
|
||||
let valued = (token_amount as f64 * reference).ceil();
|
||||
if !valued.is_finite() {
|
||||
return sats;
|
||||
}
|
||||
// Rust saturates on out-of-range float-to-int casts, so a huge reference clamps to
|
||||
// u64::MAX and `min` still picks the sats side.
|
||||
sats.min(valued as u64)
|
||||
}
|
||||
|
||||
/// Min-depth-weighted median of pool spot ratios.
|
||||
///
|
||||
/// `d_i = min(S_i, T_i * R)` needs an `R` to value the token side, and `R` is what we are
|
||||
/// solving for. Carrying the previous leg's `R` forward would resolve that circularity,
|
||||
/// but it would also make the answer depend on where the caller started reading. Instead
|
||||
/// seed with the unweighted median and reweight to a fixed point, which depends on
|
||||
/// nothing but the reserves themselves.
|
||||
fn weighted_median_reference(entries: &mut [(f64, u64, u64)]) -> Option<f64> {
|
||||
if entries.is_empty() {
|
||||
return None;
|
||||
}
|
||||
entries.sort_unstable_by(|a, b| a.0.partial_cmp(&b.0).unwrap_or(Ordering::Equal));
|
||||
|
||||
let mut reference = median_by(entries, |_| 1)?;
|
||||
|
||||
// Raising `reference` can only raise each `min(S, T * reference)` weight, and only for
|
||||
// pools not yet capped by their sats side — which are exactly the high-spot tail. So
|
||||
// the chosen index moves monotonically upwards and never revisits a position, giving a
|
||||
// fixed point within `entries.len()` passes. A fixed pass count would instead return a
|
||||
// truncated iterate: on a set of ten pools spread over many decades, stopping at four
|
||||
// passes yields 411_155 where the fixed point is 501_964, and wider spreads diverge by
|
||||
// more than an order of magnitude.
|
||||
for _ in 0..entries.len() {
|
||||
let Some(next) = median_by(entries, |e| min_depth(e.1, e.2, reference)) else {
|
||||
break;
|
||||
};
|
||||
if next == reference {
|
||||
break;
|
||||
}
|
||||
reference = next;
|
||||
}
|
||||
Some(reference)
|
||||
}
|
||||
|
||||
/// Lower weighted median of pre-sorted `entries`: the first value whose cumulative weight
|
||||
/// reaches half the total.
|
||||
///
|
||||
/// Compares `2 * cumulative >= total` rather than `cumulative >= total / 2` so integer
|
||||
/// division cannot bias the pick downwards. When every weight is zero it falls back to
|
||||
/// the positional median, so an all-dust set does not collapse to its smallest ratio.
|
||||
fn median_by(sorted: &[(f64, u64, u64)], weight: impl Fn(&(f64, u64, u64)) -> u64) -> Option<f64> {
|
||||
let total: u128 = sorted.iter().map(|e| weight(e) as u128).sum();
|
||||
if total == 0 {
|
||||
// Defensive only: `spot_ratio` rejects `sats == 0`, so every live weight is >= 1.
|
||||
// Index chosen to match what the loop below returns for equal weights, rather than
|
||||
// `len / 2`, which is the upper median for an even count.
|
||||
return sorted.get(sorted.len().saturating_sub(1) / 2).map(|e| e.0);
|
||||
}
|
||||
|
||||
let mut cumulative: u128 = 0;
|
||||
for entry in sorted {
|
||||
cumulative += weight(entry) as u128;
|
||||
if cumulative * 2 >= total {
|
||||
return Some(entry.0);
|
||||
}
|
||||
}
|
||||
sorted.last().map(|e| e.0)
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests;
|
||||
444
src/db/cauldron/candlestick/policy/tests.rs
Normal file
444
src/db/cauldron/candlestick/policy/tests.rs
Normal file
|
|
@ -0,0 +1,444 @@
|
|||
// Copyright (C) 2024-2026 Whiterun LLC
|
||||
//
|
||||
// This software is licensed under the GNU Affero General Public License (AGPL), version 3.0 or later.
|
||||
// A copy of the license can be found in the LICENSE file or at https://www.gnu.org/licenses/agpl-3.0.html
|
||||
|
||||
use super::*;
|
||||
|
||||
fn params() -> GuardParams {
|
||||
GuardParams {
|
||||
max_deviation_factor: 5.0,
|
||||
min_share_fraction: 0.05,
|
||||
}
|
||||
}
|
||||
|
||||
/// A leg at `pool`, leaving that pool holding `sats`/`tokens` afterwards.
|
||||
fn leg(
|
||||
pool: u8,
|
||||
sequence: i64,
|
||||
sats_delta: i64,
|
||||
token_delta: i64,
|
||||
sats: u64,
|
||||
token_amount: u64,
|
||||
) -> Leg {
|
||||
Leg {
|
||||
txid: [pool; 32],
|
||||
pool: [pool; 32],
|
||||
ts: 1000 + sequence,
|
||||
sequence,
|
||||
sats_delta,
|
||||
token_delta,
|
||||
sats,
|
||||
token_amount,
|
||||
}
|
||||
}
|
||||
|
||||
fn reserves(pool: u8, sequence: i64, sats: u64, token_amount: u64) -> PoolReserves {
|
||||
PoolReserves {
|
||||
pool: [pool; 32],
|
||||
sequence,
|
||||
sats,
|
||||
token_amount,
|
||||
}
|
||||
}
|
||||
|
||||
/// Feed a leg through judge-then-apply, returning the verdict.
|
||||
fn step(policy: &mut Policy, l: &Leg) -> JudgeResult {
|
||||
let judged = policy.judge(l);
|
||||
policy.apply(l);
|
||||
judged
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_leg_with_no_token_movement_is_unpriceable() {
|
||||
let policy = Policy::new(params());
|
||||
let judged = policy.judge(&leg(1, 1, 100, 0, 1_000_000, 10_000));
|
||||
assert_eq!(judged.verdict, Verdict::Unpriceable);
|
||||
assert!(judged.price.is_none());
|
||||
assert!(!judged.accepted());
|
||||
}
|
||||
|
||||
/// A leg that moved tokens for zero sats has no meaningful price; printing it as 0.0
|
||||
/// would drag a candle's low to the floor.
|
||||
#[test]
|
||||
fn test_leg_with_no_sats_movement_is_unpriceable() {
|
||||
let policy = Policy::new(params());
|
||||
let judged = policy.judge(&leg(1, 1, 0, -1_000, 1_000_000, 10_000));
|
||||
assert_eq!(judged.verdict, Verdict::Unpriceable);
|
||||
assert!(judged.price.is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_first_leg_of_a_new_token_prints_and_sets_the_reference() {
|
||||
let mut policy = Policy::new(params());
|
||||
let judged = step(&mut policy, &leg(1, 1, 100_000, -1_000, 1_000_000, 10_000));
|
||||
|
||||
assert_eq!(judged.verdict, Verdict::NoReference);
|
||||
assert!(judged.accepted());
|
||||
assert_eq!(judged.price, Some(100.0));
|
||||
assert_eq!(policy.reference(), Some(100.0));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_leg_near_the_reference_prints() {
|
||||
let mut policy = Policy::new(params());
|
||||
step(&mut policy, &leg(1, 1, 100_000, -1_000, 1_000_000, 10_000));
|
||||
|
||||
let judged = step(&mut policy, &leg(2, 2, 190_000, -2_000, 2_000_000, 20_000));
|
||||
assert!(matches!(judged.verdict, Verdict::WithinBand { .. }));
|
||||
assert!(judged.accepted());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_leg_far_from_the_reference_without_credit_is_muted() {
|
||||
let mut policy = Policy::new(params());
|
||||
step(&mut policy, &leg(1, 1, 100_000, -1_000, 1_000_000, 10_000));
|
||||
|
||||
// A fresh pool printing 100x off market, with no accepted print behind it.
|
||||
let judged = step(
|
||||
&mut policy,
|
||||
&leg(9, 2, 10_000_000, -1_000, 10_001_000, 1_010),
|
||||
);
|
||||
assert!(
|
||||
matches!(judged.verdict, Verdict::Muted { .. }),
|
||||
"{:?}",
|
||||
judged.verdict
|
||||
);
|
||||
assert!(!judged.accepted());
|
||||
// Muted, but the price was still computed — the caller still counts its volume.
|
||||
assert!(judged.price.is_some());
|
||||
}
|
||||
|
||||
/// Regression: `apply` inserted new pools with `sequence: leg.sequence` and then gated the
|
||||
/// reserve update behind `leg.sequence > state.sequence`, which is false on that very first
|
||||
/// insert — so a pool's opening leg left it recorded with zero reserves, and therefore zero
|
||||
/// credit. The `i64::MIN` sentinel is what fixes it.
|
||||
#[test]
|
||||
fn test_pool_holds_credit_after_its_very_first_leg() {
|
||||
let mut policy = Policy::new(params());
|
||||
step(&mut policy, &leg(1, 1, 100_000, -1_000, 1_000_000, 10_000));
|
||||
|
||||
let second = leg(2, 2, 200_000, -2_000, 2_000_000, 20_000);
|
||||
let judged = step(&mut policy, &second);
|
||||
assert!(judged.accepted(), "{:?}", judged.verdict);
|
||||
|
||||
assert!(
|
||||
policy.credit_of(&[2; 32]) > 0,
|
||||
"a pool's opening leg must record its reserves, not leave it at zero"
|
||||
);
|
||||
}
|
||||
|
||||
/// A wild print leaves the pool that made it sitting off-market, and an off-market pool
|
||||
/// holds no credit — so the pool cannot use one absurd leg to qualify the next.
|
||||
#[test]
|
||||
fn test_pool_left_off_market_by_its_own_print_holds_no_credit() {
|
||||
let mut policy = Policy::new(params());
|
||||
step(&mut policy, &leg(1, 1, 100_000, -1_000, 1_000_000, 10_000));
|
||||
step(
|
||||
&mut policy,
|
||||
&leg(9, 2, 10_000_000, -1_000, 10_001_000, 1_010),
|
||||
);
|
||||
|
||||
assert_eq!(
|
||||
policy.credit_of(&[9; 32]),
|
||||
0,
|
||||
"a pool 99x off the reference must not qualify for tier 2"
|
||||
);
|
||||
}
|
||||
|
||||
/// Regression for the seeded/folded credit divergence — the reason credit is derived from
|
||||
/// reserves rather than remembered.
|
||||
///
|
||||
/// Pool 2's only leg is a one-token trade at 100x the market. It is muted, but it is far too
|
||||
/// small to move the pool's own spot, so the pool is left deep and in-band. Under the old
|
||||
/// "min-depth at the last accepted print" rule the folded policy held pool 2 at zero credit
|
||||
/// while a policy seeded from the resulting reserves granted it full credit — and the two
|
||||
/// then returned `Muted` and `CreditExempt` for the identical next leg. Since `ohlcv_1h` is
|
||||
/// written with `INSERT OR IGNORE`, whichever path ran first would have frozen its answer in.
|
||||
#[test]
|
||||
fn test_seeded_and_folded_agree_on_credit_after_a_muted_print() {
|
||||
let mut folded = Policy::new(params());
|
||||
step(&mut folded, &leg(1, 1, 100_000, -1_000, 1_000_000, 10_000));
|
||||
// Pool 2's opening leg: absurd price, negligible size, so it ends in-band.
|
||||
step(&mut folded, &leg(2, 2, 10_000, -1, 2_010_000, 19_999));
|
||||
|
||||
let seeded = Policy::seeded(
|
||||
params(),
|
||||
&[
|
||||
reserves(1, 1, 1_000_000, 10_000),
|
||||
reserves(2, 2, 2_010_000, 19_999),
|
||||
],
|
||||
);
|
||||
|
||||
assert_eq!(folded.reference(), seeded.reference());
|
||||
assert_eq!(
|
||||
folded.credit_of(&[2; 32]),
|
||||
seeded.credit_of(&[2; 32]),
|
||||
"credit must not depend on whether the policy folded history or was seeded from it"
|
||||
);
|
||||
assert!(
|
||||
folded.credit_of(&[2; 32]) > 0,
|
||||
"the test is vacuous unless the pool actually ends up in-band and credited"
|
||||
);
|
||||
|
||||
// The consequence that made this worth fixing: identical verdicts on the next leg.
|
||||
let next = leg(2, 3, 500_000, -1, 2_510_000, 19_998);
|
||||
assert_eq!(folded.judge(&next).verdict, seeded.judge(&next).verdict);
|
||||
}
|
||||
|
||||
/// The property the whole seeded design exists for: the reference depends only on pool
|
||||
/// reserves, never on how many legs the caller happened to read to arrive at them.
|
||||
#[test]
|
||||
fn test_reference_is_a_function_of_reserves_not_of_history_read() {
|
||||
let mut folded = Policy::new(params());
|
||||
step(&mut folded, &leg(1, 1, 100_000, -1_000, 1_000_000, 10_000));
|
||||
step(&mut folded, &leg(2, 2, 190_000, -2_000, 2_000_000, 20_000));
|
||||
step(&mut folded, &leg(1, 3, 105_000, -1_000, 1_105_000, 9_000));
|
||||
|
||||
// Same end state, reached by seeding instead of folding.
|
||||
let seeded = Policy::seeded(
|
||||
params(),
|
||||
&[
|
||||
reserves(1, 3, 1_105_000, 9_000),
|
||||
reserves(2, 2, 2_000_000, 20_000),
|
||||
],
|
||||
);
|
||||
|
||||
assert_eq!(
|
||||
folded.reference(),
|
||||
seeded.reference(),
|
||||
"a seeded policy must agree with a folded one on identical reserves"
|
||||
);
|
||||
}
|
||||
|
||||
/// Corollary: two callers whose windows start at different points must reach the same
|
||||
/// verdict on a leg they both see.
|
||||
#[test]
|
||||
fn test_seeded_and_folded_policies_agree_on_a_shared_leg() {
|
||||
let mut folded = Policy::new(params());
|
||||
step(&mut folded, &leg(1, 1, 100_000, -1_000, 1_000_000, 10_000));
|
||||
step(&mut folded, &leg(2, 2, 190_000, -2_000, 2_000_000, 20_000));
|
||||
|
||||
let seeded = Policy::seeded(
|
||||
params(),
|
||||
&[
|
||||
reserves(1, 1, 1_000_000, 10_000),
|
||||
reserves(2, 2, 2_000_000, 20_000),
|
||||
],
|
||||
);
|
||||
|
||||
let shared = leg(1, 3, 300_000, -1_000, 1_300_000, 9_000);
|
||||
assert_eq!(
|
||||
folded.judge(&shared).accepted(),
|
||||
seeded.judge(&shared).accepted(),
|
||||
"verdict must not depend on where the caller started reading"
|
||||
);
|
||||
}
|
||||
|
||||
/// Seeding hands credit to pools already trading at market, but withholds it from one
|
||||
/// that has been walked far away — otherwise a snapshot would launder an attacker's pool
|
||||
/// straight into qualification.
|
||||
#[test]
|
||||
fn test_seeding_withholds_credit_from_off_market_pools() {
|
||||
let policy = Policy::seeded(
|
||||
params(),
|
||||
&[
|
||||
reserves(1, 1, 1_000_000, 10_000), // spot 100, at market
|
||||
reserves(2, 2, 2_000_000, 20_000), // spot 100, at market
|
||||
reserves(9, 3, 10_000_000, 1_000), // spot 10_000, 100x off
|
||||
],
|
||||
);
|
||||
|
||||
assert_eq!(policy.reference(), Some(100.0));
|
||||
assert!(policy.credit_of(&[1; 32]) > 0);
|
||||
assert!(policy.credit_of(&[2; 32]) > 0);
|
||||
assert_eq!(
|
||||
policy.credit_of(&[9; 32]),
|
||||
0,
|
||||
"a pool sitting 100x off market must not be seeded with credit"
|
||||
);
|
||||
}
|
||||
|
||||
/// Min-depth weighting is what stops shallow pools from voting the reference away from
|
||||
/// where the real liquidity sits.
|
||||
///
|
||||
/// The dust pools deliberately OUTNUMBER the real one, so an unweighted median lands on
|
||||
/// them (10_000) and only the weighting pulls the answer back to 100. A test where all
|
||||
/// pools quote the same price would pass no matter what the weighting did.
|
||||
#[test]
|
||||
fn test_dust_pools_do_not_outvote_one_deep_pool() {
|
||||
let policy = Policy::seeded(
|
||||
params(),
|
||||
&[
|
||||
reserves(1, 1, 100_000_000, 1_000_000), // spot 100, deep
|
||||
reserves(8, 2, 100_000, 10), // spot 10_000, shallow
|
||||
reserves(9, 3, 100_000, 10), // spot 10_000, shallow
|
||||
],
|
||||
);
|
||||
|
||||
assert_eq!(
|
||||
policy.reference(),
|
||||
Some(100.0),
|
||||
"two shallow pools must not outvote one holding a thousand times the depth"
|
||||
);
|
||||
}
|
||||
|
||||
/// The dust floor is a guard against degenerate arithmetic, not a weighting mechanism —
|
||||
/// min-depth already suppresses thin pools. What it must do is keep a pool that cannot
|
||||
/// quote a ratio from producing an infinity.
|
||||
#[test]
|
||||
fn test_pools_that_cannot_quote_a_ratio_are_excluded() {
|
||||
let policy = Policy::seeded(
|
||||
params(),
|
||||
&[
|
||||
reserves(1, 1, 1_000_000, 10_000), // spot 100
|
||||
reserves(8, 2, 1_000_000, 0), // no tokens: ratio would divide by zero
|
||||
reserves(9, 3, 0, 10_000), // no sats: ratio would be 0.0
|
||||
],
|
||||
);
|
||||
|
||||
let reference = policy
|
||||
.reference()
|
||||
.expect("the one quotable pool sets the reference");
|
||||
assert!(
|
||||
reference.is_finite() && reference == 100.0,
|
||||
"got {reference}"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_policy_with_no_usable_pools_has_no_reference() {
|
||||
let policy = Policy::seeded(params(), &[reserves(1, 1, 0, 0)]);
|
||||
assert_eq!(policy.reference(), None);
|
||||
}
|
||||
|
||||
/// `cumulative >= total / 2` truncates, which tips the pick towards the lower ratio.
|
||||
///
|
||||
/// Weights 2 and 3 are the minimal case that separates the two forms: the old test is
|
||||
/// `2 >= 5/2 == 2`, true, so it stops on the first entry; the correct `2*2 >= 5` is false,
|
||||
/// so it advances to the second. Driving `median_by` directly is what makes this
|
||||
/// falsifiable — going through `weighted_median_reference` cannot pin the weights.
|
||||
#[test]
|
||||
fn test_weighted_median_is_not_biased_by_integer_division() {
|
||||
let sorted = [(100.0, 2u64, 0u64), (200.0, 3u64, 0u64)];
|
||||
assert_eq!(
|
||||
median_by(&sorted, |e| e.1),
|
||||
Some(200.0),
|
||||
"weight 3 holds the majority, so the median is 200, not 100"
|
||||
);
|
||||
|
||||
// An odd, clearly-weighted set must land on the true weighted median.
|
||||
let mut skewed = vec![(10.0, 1, 1), (100.0, 1_000_000, 1_000_000), (1000.0, 1, 1)];
|
||||
assert_eq!(weighted_median_reference(&mut skewed), Some(100.0));
|
||||
}
|
||||
|
||||
/// Regression: the reweighting loop ran a fixed four passes and returned whatever iterate
|
||||
/// it had reached. Because raising the reference only ever raises weights, the pick climbs
|
||||
/// one position per pass, so a set spread over many decades stops short — this one settled
|
||||
/// on 411_155 against a true fixed point of 501_964, and wider spreads diverged by more
|
||||
/// than an order of magnitude.
|
||||
#[test]
|
||||
fn test_reference_reaches_its_fixed_point_not_a_truncated_iterate() {
|
||||
let pools: [(u64, u64); 10] = [
|
||||
(962, 249_939_293_128_076),
|
||||
(3_198, 71_909_139_611_364),
|
||||
(455_109_248_120, 243_581_757_429_499),
|
||||
(279, 9_934),
|
||||
(432_371, 2_608),
|
||||
(208_541_212_143, 728_829_157),
|
||||
(925_517_097_739, 737_748_546),
|
||||
(34_778_861_631_604, 885_976_135),
|
||||
(89_743_615_643_767, 218_272_037),
|
||||
(499_655_065_278_927, 995_401_018),
|
||||
];
|
||||
let mut entries: Vec<(f64, u64, u64)> = pools
|
||||
.iter()
|
||||
.map(|(sats, tokens)| (*sats as f64 / *tokens as f64, *sats, *tokens))
|
||||
.collect();
|
||||
|
||||
let reference = weighted_median_reference(&mut entries).unwrap();
|
||||
|
||||
// `entries` is left sorted, so one more pass is a direct fixed-point check.
|
||||
let next = median_by(&entries, |e| min_depth(e.1, e.2, reference)).unwrap();
|
||||
assert_eq!(
|
||||
next, reference,
|
||||
"reweighting must converge; {reference} moves to {next} on the next pass"
|
||||
);
|
||||
}
|
||||
|
||||
/// All-zero weights used to make the cumulative test pass on the first entry, returning
|
||||
/// the smallest ratio rather than the middle one.
|
||||
#[test]
|
||||
fn test_all_zero_weights_fall_back_to_the_positional_median() {
|
||||
let mut entries = vec![(10.0, 0, 0), (100.0, 0, 0), (1000.0, 0, 0)];
|
||||
assert_eq!(weighted_median_reference(&mut entries), Some(100.0));
|
||||
}
|
||||
|
||||
/// Documents a KNOWN GAP in tier-2 rather than endorsing it: the threshold is a share of
|
||||
/// `largest_credit`, and the pool holding that largest credit trivially clears a share of
|
||||
/// itself. So the deepest pool — and every pool of a single-pool token — can print any
|
||||
/// price it likes. Left as-is pending a decision on the tier-2 rule; this test exists so
|
||||
/// the behaviour is visible and any change to it is deliberate.
|
||||
#[test]
|
||||
fn test_known_gap_deepest_pool_self_qualifies_for_tier_two() {
|
||||
let mut policy = Policy::seeded(params(), &[reserves(1, 1, 10_000_000, 100_000)]);
|
||||
assert_eq!(policy.reference(), Some(100.0));
|
||||
assert!(policy.credit_of(&[1; 32]) > 0);
|
||||
|
||||
// The same pool now prints ten million times off its own spot.
|
||||
let absurd = leg(1, 2, 1_000_000_000, -1, 11_000_000_000, 99_999);
|
||||
let judged = policy.judge(&absurd);
|
||||
|
||||
assert!(
|
||||
matches!(judged.verdict, Verdict::CreditExempt { .. }),
|
||||
"expected the known self-qualification gap, got {:?}",
|
||||
judged.verdict
|
||||
);
|
||||
assert!(judged.accepted());
|
||||
let _ = &mut policy;
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_min_depth_values_the_token_side_at_the_reference() {
|
||||
// Token side is the thinner one: 10 tokens * 100 = 1_000 < 50_000 sats.
|
||||
assert_eq!(min_depth(50_000, 10, 100.0), 1_000);
|
||||
// Sats side is the thinner one.
|
||||
assert_eq!(min_depth(500, 10, 100.0), 500);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_min_depth_saturates_instead_of_overflowing() {
|
||||
assert_eq!(min_depth(42, u64::MAX, f64::MAX), 42);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_deviation_is_symmetric_and_rejects_degenerate_input() {
|
||||
assert_eq!(deviation(200.0, 100.0), 2.0);
|
||||
assert_eq!(deviation(50.0, 100.0), 2.0);
|
||||
assert_eq!(deviation(0.0, 100.0), f64::INFINITY);
|
||||
assert_eq!(deviation(100.0, 0.0), f64::INFINITY);
|
||||
// NaN and the infinities must land on INFINITY (i.e. "out of band"), never sneak
|
||||
// through as a small deviation. `<=` alone is false for NaN, so the finiteness test
|
||||
// has to run first — this is what pins that ordering.
|
||||
assert_eq!(deviation(f64::NAN, 100.0), f64::INFINITY);
|
||||
assert_eq!(deviation(100.0, f64::NAN), f64::INFINITY);
|
||||
assert_eq!(deviation(f64::INFINITY, 100.0), f64::INFINITY);
|
||||
assert_eq!(deviation(100.0, f64::INFINITY), f64::INFINITY);
|
||||
assert_eq!(deviation(-100.0, 100.0), f64::INFINITY);
|
||||
assert_eq!(deviation(100.0, -100.0), f64::INFINITY);
|
||||
}
|
||||
|
||||
/// A leg replayed at or below a pool's known sequence must not rewind its reserves.
|
||||
#[test]
|
||||
fn test_out_of_order_leg_does_not_rewind_reserves() {
|
||||
let mut policy = Policy::new(params());
|
||||
step(&mut policy, &leg(1, 10, 100_000, -1_000, 1_000_000, 10_000));
|
||||
let after_current = policy.reference();
|
||||
|
||||
step(&mut policy, &leg(1, 5, 100_000, -1_000, 7, 7));
|
||||
assert_eq!(
|
||||
policy.reference(),
|
||||
after_current,
|
||||
"a stale leg must not overwrite newer reserves"
|
||||
);
|
||||
}
|
||||
|
|
@ -9,6 +9,7 @@ use crate::db::cauldron::{
|
|||
pool::{self, dummy_init_seq},
|
||||
tx::{self, insert_block_tx, insert_mempool_tx},
|
||||
utxo_funding::{self, insert_utxo_funding},
|
||||
utxo_spending,
|
||||
};
|
||||
use crate::utiltest::mock_db_pool;
|
||||
use bitcoin_hashes::Hash;
|
||||
|
|
@ -38,12 +39,48 @@ fn dummy_cauldron(
|
|||
|
||||
async fn setup_db(pool: sqlx::SqlitePool) {
|
||||
utxo_funding::create_table(&pool).await;
|
||||
utxo_spending::create_table(&pool).await;
|
||||
tx::create_table(&pool).await;
|
||||
pool::create_table(&pool).await;
|
||||
ohlcv::create_table(&pool).await;
|
||||
dummy_init_seq();
|
||||
}
|
||||
|
||||
/// Register a pool row so the reserve snapshot can find it, optionally marking it
|
||||
/// withdrawn by the transaction `withdrawn_by` at `withdrawn_at`.
|
||||
async fn register_pool(
|
||||
conn: &mut sqlx::pool::PoolConnection<sqlx::Sqlite>,
|
||||
pool_hash: &OutPointHash,
|
||||
token: &TokenID,
|
||||
withdrawal: Option<(Txid, i64)>,
|
||||
) {
|
||||
let withdrawn_utxo = withdrawal.map(|(txid, ts)| {
|
||||
let spent = OutPointHash::from_byte_array(*txid.as_byte_array());
|
||||
(spent, txid, ts)
|
||||
});
|
||||
|
||||
sqlx::query("INSERT OR REPLACE INTO pool (creation_utxo, owner_pkh, token_id, withdrawn_in_utxo) VALUES (?, ?, ?, ?)")
|
||||
.bind(pool_hash.to_blob())
|
||||
.bind(PubkeyHash::all_zeros().to_blob())
|
||||
.bind(token.to_blob())
|
||||
.bind(withdrawn_utxo.as_ref().map(|(spent, _, _)| spent.to_blob()))
|
||||
.execute(&mut **conn)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
if let Some((spent, txid, ts)) = withdrawn_utxo {
|
||||
insert_mempool_tx(&mut **conn, &txid, ts as u64)
|
||||
.await
|
||||
.unwrap();
|
||||
sqlx::query("INSERT OR REPLACE INTO utxo_spending (spent_utxo_hash, txid) VALUES (?, ?)")
|
||||
.bind(spent.to_blob())
|
||||
.bind(txid.to_blob())
|
||||
.execute(&mut **conn)
|
||||
.await
|
||||
.unwrap();
|
||||
}
|
||||
}
|
||||
|
||||
async fn insert_trade_at(
|
||||
conn: &mut sqlx::pool::PoolConnection<sqlx::Sqlite>,
|
||||
token: &TokenID,
|
||||
|
|
@ -144,8 +181,11 @@ async fn insert_multileg_trade_at(
|
|||
/// Dividing the signed sums produced a price no leg traded at — on mainnet token NWB
|
||||
/// (tx 1E84F4E9…1916, 27 legs) that printed 30,792,599.5 sats/unit against legs that
|
||||
/// actually executed between 0.288 and 0.335.
|
||||
///
|
||||
/// The durable invariant is that a printed price must be one a leg actually executed at.
|
||||
/// Under per-leg pricing the close is specifically the last accepted leg's price.
|
||||
#[tokio::test]
|
||||
async fn test_multipool_arb_priced_by_gross_volume_not_net() {
|
||||
async fn test_multipool_arb_prices_from_its_legs_not_their_net() {
|
||||
let db = mock_db_pool(setup_db).await;
|
||||
let token = TokenID::from_byte_array([0xC1; 32]);
|
||||
let token_blob = token.to_blob();
|
||||
|
|
@ -161,6 +201,8 @@ async fn test_multipool_arb_priced_by_gross_volume_not_net() {
|
|||
.unwrap()
|
||||
.expect("arb transaction must price");
|
||||
|
||||
// Per-leg pricing: each leg executes at its own price, not a net ratio.
|
||||
// The last leg is the buy leg: 384_906_040 / 1_334_527_067 = 0.288421
|
||||
let low_leg = 384_906_040.0 / 1_334_527_067.0;
|
||||
let high_leg = 446_491_239.0 / 1_334_527_069.0;
|
||||
assert!(
|
||||
|
|
@ -168,25 +210,24 @@ async fn test_multipool_arb_priced_by_gross_volume_not_net() {
|
|||
"price {price} must lie within the executed leg range [{low_leg}, {high_leg}]"
|
||||
);
|
||||
|
||||
// Net-ratio pricing would divide 61,585,199 sats by 2 token units.
|
||||
let net_ratio = 61_585_199.0 / 2.0;
|
||||
// The close should be the last leg's price (buy leg).
|
||||
let expected = 384_906_040.0 / 1_334_527_067.0;
|
||||
assert!(
|
||||
price < net_ratio / 1000.0,
|
||||
"price {price} must not resemble the netting artifact {net_ratio}"
|
||||
(price - expected).abs() < 1e-9,
|
||||
"expected {expected} (last leg price), got {price}"
|
||||
);
|
||||
|
||||
let expected = 831_397_279.0 / 2_669_054_136.0;
|
||||
assert!((price - expected).abs() < 1e-9, "expected {expected}");
|
||||
}
|
||||
|
||||
/// Every leg pointing the same way is the ordinary case: gross and net agree exactly,
|
||||
/// so 99.76% of mainnet prints — including the OLA supply-shock crash — are untouched.
|
||||
/// Every leg pointing the same way is the ordinary case — the one covering 99.76% of
|
||||
/// mainnet prints, including the OLA supply-shock crash. Nothing here should be muted or
|
||||
/// distorted: each leg executed at a genuine price and the candle closes on the last.
|
||||
#[tokio::test]
|
||||
async fn test_single_direction_multileg_price_matches_net_ratio() {
|
||||
async fn test_single_direction_multileg_closes_on_its_last_leg() {
|
||||
let db = mock_db_pool(setup_db).await;
|
||||
let token = TokenID::from_byte_array([0xC2; 32]);
|
||||
let token_blob = token.to_blob();
|
||||
|
||||
// Leg prices: 50.0, 50.0, 49.5 — a tight, ordinary spread.
|
||||
let legs = [(150_000i64, -3_000i64), (50_000, -1_000), (99_000, -2_000)];
|
||||
let mut conn = db.cauldron_w.acquire().await.unwrap();
|
||||
insert_multileg_trade_at(&mut conn, &token, 0x22, 1000, &legs).await;
|
||||
|
|
@ -196,13 +237,157 @@ async fn test_single_direction_multileg_price_matches_net_ratio() {
|
|||
.unwrap()
|
||||
.expect("router transaction must price");
|
||||
|
||||
let signed_sats: i64 = legs.iter().map(|l| l.0).sum();
|
||||
let signed_tokens: i64 = legs.iter().map(|l| l.1).sum();
|
||||
let net_ratio = (signed_sats as f64 / signed_tokens as f64).abs();
|
||||
let leg_prices: Vec<f64> = legs
|
||||
.iter()
|
||||
.map(|(s, t)| *s as f64 / t.unsigned_abs() as f64)
|
||||
.collect();
|
||||
let lowest = leg_prices.iter().cloned().fold(f64::MAX, f64::min);
|
||||
let highest = leg_prices.iter().cloned().fold(f64::MIN, f64::max);
|
||||
assert!(
|
||||
(price - net_ratio).abs() < f64::EPSILON,
|
||||
"single-direction transactions must be unaffected: {price} vs {net_ratio}"
|
||||
price >= lowest && price <= highest,
|
||||
"price {price} must be one the legs actually executed at, within [{lowest}, {highest}]"
|
||||
);
|
||||
|
||||
let last_leg_price = *leg_prices.last().unwrap();
|
||||
assert!(
|
||||
(price - last_leg_price).abs() < f64::EPSILON,
|
||||
"close should be the last leg's price: {price} vs {last_leg_price}"
|
||||
);
|
||||
}
|
||||
|
||||
/// A pool that has been drained leaves its last history row showing full pre-withdrawal
|
||||
/// reserves, because a withdrawal writes no new row. Seeding the policy from those would
|
||||
/// let a pool holding nothing keep voting the reference — the OLA case, where a 10.9B-sat
|
||||
/// pool was withdrawn shortly before a crash and a lingering ghost would have muted it.
|
||||
#[tokio::test]
|
||||
async fn test_reserve_snapshot_drops_pools_withdrawn_before_the_instant() {
|
||||
let db = mock_db_pool(setup_db).await;
|
||||
let token = TokenID::from_byte_array([0xC6; 32]);
|
||||
let token_blob = token.to_blob();
|
||||
|
||||
let mut conn = db.cauldron_w.acquire().await.unwrap();
|
||||
insert_trade_at(&mut conn, &token, 0x41, 1000, 100_000, 2_000).await;
|
||||
|
||||
let pool_hash = OutPointHash::from_byte_array([0x41u8.wrapping_add(0x80); 32]);
|
||||
register_pool(&mut conn, &pool_hash, &token, None).await;
|
||||
|
||||
let live = super::fetch_reserve_snapshot(&db.cauldron_r, &token_blob, 5000, false)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(live.len(), 1, "a live pool must be in the snapshot");
|
||||
|
||||
// Withdraw it at ts=2000 and ask again afterwards.
|
||||
let withdrawal_tx = Txid::from_byte_array([0x42; 32]);
|
||||
register_pool(&mut conn, &pool_hash, &token, Some((withdrawal_tx, 2000))).await;
|
||||
|
||||
let after = super::fetch_reserve_snapshot(&db.cauldron_r, &token_blob, 5000, false)
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(
|
||||
after.is_empty(),
|
||||
"a pool withdrawn at 2000 must not vote on the reference at 5000: {after:?}"
|
||||
);
|
||||
|
||||
// But it was live at ts=1500, so a historical query must still see it.
|
||||
let before = super::fetch_reserve_snapshot(&db.cauldron_r, &token_blob, 1500, false)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
before.len(),
|
||||
1,
|
||||
"a pool withdrawn later was still live earlier and must remain visible"
|
||||
);
|
||||
}
|
||||
|
||||
/// Regression: the seed consulted `ohlcv_1h` first and returned on any hit. Since that
|
||||
/// table is structurally at least three hours behind the tip — and completely empty for
|
||||
/// the whole of a post-version-bump rebuild — a stale bucket shadowed every newer leg.
|
||||
#[tokio::test]
|
||||
async fn test_seed_prefers_recent_legs_over_a_stale_materialised_bucket() {
|
||||
let db = mock_db_pool(setup_db).await;
|
||||
let token = TokenID::from_byte_array([0xC7; 32]);
|
||||
let token_blob = token.to_blob();
|
||||
let mut conn = db.cauldron_w.acquire().await.unwrap();
|
||||
|
||||
sqlx::query(
|
||||
"INSERT INTO ohlcv_1h
|
||||
(token_id, bucket_ts, open, high, low, close, volume_sats, volume_tokens, tx_count)
|
||||
VALUES (?, 0, 10.0, 10.0, 10.0, 10.0, 1, 1, 1)",
|
||||
)
|
||||
.bind(&token_blob)
|
||||
.execute(&mut *conn)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
// A real print at 50.0 two hours later that the table has not caught up with.
|
||||
insert_trade_at(&mut conn, &token, 0x51, 7200, 100_000, 2_000).await;
|
||||
|
||||
let price = super::fetch_last_close_before(&db.cauldron_r, &token_blob, 10800)
|
||||
.await
|
||||
.unwrap()
|
||||
.expect("a seed must exist");
|
||||
|
||||
assert!(
|
||||
(price - 50.0).abs() < f64::EPSILON,
|
||||
"stale bucket must not shadow the newer print: got {price}"
|
||||
);
|
||||
}
|
||||
|
||||
/// Regression: the materialised lookup took any bucket with `bucket_ts < timestamp_end`,
|
||||
/// so for a non-hour-aligned cutoff it could return a bucket straddling that cutoff —
|
||||
/// closing on a print from *after* the instant being asked about.
|
||||
#[tokio::test]
|
||||
async fn test_seed_does_not_leak_a_price_from_after_the_cutoff() {
|
||||
let db = mock_db_pool(setup_db).await;
|
||||
let token = TokenID::from_byte_array([0xC8; 32]);
|
||||
let token_blob = token.to_blob();
|
||||
let mut conn = db.cauldron_w.acquire().await.unwrap();
|
||||
|
||||
// Bucket [7200, 10800) closes at 80.0 — a price set late in that hour.
|
||||
sqlx::query(
|
||||
"INSERT INTO ohlcv_1h
|
||||
(token_id, bucket_ts, open, high, low, close, volume_sats, volume_tokens, tx_count)
|
||||
VALUES (?, 7200, 50.0, 80.0, 50.0, 80.0, 1, 1, 1)",
|
||||
)
|
||||
.bind(&token_blob)
|
||||
.execute(&mut *conn)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
// What actually happened by ts=8000 was a single print at 50.0.
|
||||
insert_trade_at(&mut conn, &token, 0x52, 7300, 100_000, 2_000).await;
|
||||
|
||||
let price = super::fetch_last_close_before(&db.cauldron_r, &token_blob, 8000)
|
||||
.await
|
||||
.unwrap()
|
||||
.expect("a seed must exist");
|
||||
|
||||
assert!(
|
||||
(price - 50.0).abs() < f64::EPSILON,
|
||||
"must not return the 80.0 close of a bucket that ends after the cutoff: got {price}"
|
||||
);
|
||||
}
|
||||
|
||||
/// Regression: the seed lookback was briefly capped at 24 hours, which silently dropped
|
||||
/// the carry-forward price for any token trading less often than daily and left a hole
|
||||
/// where the leading candles should be.
|
||||
#[tokio::test]
|
||||
async fn test_seed_lookback_reaches_past_a_long_silence() {
|
||||
let db = mock_db_pool(setup_db).await;
|
||||
let token = TokenID::from_byte_array([0xC5; 32]);
|
||||
let token_blob = token.to_blob();
|
||||
|
||||
let mut conn = db.cauldron_w.acquire().await.unwrap();
|
||||
insert_trade_at(&mut conn, &token, 0x31, 1000, 100_000, 2_000).await;
|
||||
|
||||
// Ask 60 days later. The only trade is far outside any 24-hour window.
|
||||
let sixty_days = 1000 + 60 * 24 * 3600;
|
||||
let price = super::fetch_last_close_before(&db.cauldron_r, &token_blob, sixty_days)
|
||||
.await
|
||||
.unwrap()
|
||||
.expect("a dormant token must still carry its last known price");
|
||||
|
||||
assert!((price - 50.0).abs() < f64::EPSILON, "got {price}");
|
||||
}
|
||||
|
||||
/// A transaction whose legs cancel exactly used to print nothing at all (the chart
|
||||
|
|
@ -228,8 +413,12 @@ async fn test_exactly_cancelling_legs_still_price() {
|
|||
.unwrap()
|
||||
.expect("zero-net transaction must still price from its legs");
|
||||
|
||||
// 20_000 gross sats over 2_000 gross tokens, between the 9 and 11 leg prices.
|
||||
assert!((price - 10.0).abs() < f64::EPSILON, "got {price}");
|
||||
// Per-leg pricing: the close is the last leg's price.
|
||||
// Last leg: (11_000, -1_000) = 11,000 / 1,000 = 11.0
|
||||
assert!(
|
||||
(price - 11.0).abs() < f64::EPSILON,
|
||||
"close should be last leg's price: got {price}"
|
||||
);
|
||||
}
|
||||
|
||||
/// Legs that move no tokens cannot produce a price (division by zero volume).
|
||||
|
|
|
|||
|
|
@ -12,7 +12,20 @@ use sqlx::{Row, SqlitePool};
|
|||
/// `INSERT OR IGNORE`, so existing rows are never corrected in place.
|
||||
///
|
||||
/// 2: price switched from the signed net ratio to the gross volume ratio.
|
||||
pub const OHLCV_VERSION: u32 = 2;
|
||||
/// 3: pricing switched from per-transaction to per-leg with policy-based acceptance.
|
||||
/// 4: policy reference seeded from reserves rather than folded from the batch start;
|
||||
/// f64 sentinels no longer written; credit earned on a pool's first accepted print;
|
||||
/// zero-sats legs unpriceable; weighted-median rounding and zero-weight fixes;
|
||||
/// pools withdrawn before the snapshot instant excluded from the reference.
|
||||
/// 5: reweighting iterates to its fixed point instead of stopping at four passes, so the
|
||||
/// reference moves on multi-decade pool spreads; buckets with no accepted price are
|
||||
/// written flat at the carried close instead of dropped, restoring their volume; the
|
||||
/// rebuild seeds from confirmed-only reserves, so output no longer depends on mempool
|
||||
/// contents at rebuild time.
|
||||
/// 6: qualification credit derived from present reserves instead of remembered from a
|
||||
/// pool's last accepted print, so a seeded policy and a folded one can no longer
|
||||
/// disagree on the same leg. Tier 2 is now a pure depth gate.
|
||||
pub const OHLCV_VERSION: u32 = 6;
|
||||
const OHLCV_VERSION_KEY: &str = "ohlcv_version";
|
||||
|
||||
pub async fn create_table(pool: &SqlitePool) {
|
||||
|
|
@ -84,8 +97,11 @@ pub async fn get_min_trade_bucket_ts(pool: &SqlitePool) -> Result<Option<i64>> {
|
|||
/// Materialise all 1-hour OHLCV buckets for confirmed trades whose effective timestamp falls
|
||||
/// in `[since_ts, until_ts)`.
|
||||
///
|
||||
/// Two-phase approach: the slow aggregation SELECT runs against `read_pool` (no write lock),
|
||||
/// then the pre-computed rows are bulk-inserted via `write_pool` (write lock held briefly).
|
||||
/// Three-phase approach:
|
||||
/// 1. Fetch per-leg data for all tokens in the range (no write lock)
|
||||
/// 2. Fold through each token's legs using the Policy to compute buckets (no write lock)
|
||||
/// 3. Insert pre-computed buckets in one transaction (brief write lock)
|
||||
///
|
||||
/// Uses INSERT OR IGNORE so existing rows are never overwritten.
|
||||
/// Returns the number of rows inserted.
|
||||
pub async fn rebuild_range(
|
||||
|
|
@ -94,98 +110,36 @@ pub async fn rebuild_range(
|
|||
since_ts: i64,
|
||||
until_ts: i64,
|
||||
) -> Result<u64> {
|
||||
use crate::db::cauldron::candlestick::{
|
||||
fetch_reserve_snapshot, to_hash32, GuardParams, Leg, Policy,
|
||||
};
|
||||
use std::collections::HashMap;
|
||||
|
||||
if since_ts >= until_ts {
|
||||
return Ok(0);
|
||||
}
|
||||
|
||||
// Phase 1: aggregate using the read pool — no write lock held during the slow CTE.
|
||||
let select_sql = r#"
|
||||
WITH per_pool_tx_raw AS (
|
||||
SELECT
|
||||
phe.token_id,
|
||||
phe.txid,
|
||||
phe.effective_timestamp AS ts,
|
||||
phe.utxo,
|
||||
phe.sats_delta,
|
||||
phe.token_delta,
|
||||
phe.sequence
|
||||
FROM pool_history_entry AS phe
|
||||
JOIN tx ON tx.txid = phe.txid
|
||||
WHERE tx.blockhash IS NOT NULL
|
||||
AND phe.effective_timestamp >= ?
|
||||
AND phe.effective_timestamp < ?
|
||||
),
|
||||
per_pool_tx AS (
|
||||
SELECT
|
||||
token_id,
|
||||
txid,
|
||||
ts,
|
||||
(ts / 3600) * 3600 AS bucket_ts,
|
||||
utxo,
|
||||
MIN(sequence) AS min_sequence,
|
||||
SUM(ABS(sats_delta)) AS vol_sats,
|
||||
SUM(ABS(token_delta)) AS vol_tokens
|
||||
FROM per_pool_tx_raw
|
||||
GROUP BY token_id, txid, ts, utxo
|
||||
),
|
||||
tx_trades AS (
|
||||
SELECT
|
||||
token_id,
|
||||
txid,
|
||||
ts,
|
||||
bucket_ts,
|
||||
MIN(min_sequence) AS min_sequence,
|
||||
SUM(vol_sats) AS vol_sats,
|
||||
SUM(vol_tokens) AS vol_tokens
|
||||
FROM per_pool_tx
|
||||
GROUP BY token_id, txid, ts
|
||||
),
|
||||
priceable AS (
|
||||
SELECT
|
||||
token_id,
|
||||
bucket_ts,
|
||||
CAST(vol_sats AS REAL) / CAST(vol_tokens AS REAL) AS price,
|
||||
ROW_NUMBER() OVER (PARTITION BY token_id, bucket_ts ORDER BY ts ASC, min_sequence ASC) AS rn_asc,
|
||||
ROW_NUMBER() OVER (PARTITION BY token_id, bucket_ts ORDER BY ts DESC, min_sequence DESC) AS rn_desc
|
||||
FROM tx_trades
|
||||
WHERE vol_tokens != 0
|
||||
),
|
||||
ohlc AS (
|
||||
SELECT
|
||||
token_id,
|
||||
bucket_ts,
|
||||
MAX(CASE WHEN rn_asc = 1 THEN price END) AS open,
|
||||
MAX(CASE WHEN rn_desc = 1 THEN price END) AS close,
|
||||
MAX(price) AS high,
|
||||
MIN(price) AS low
|
||||
FROM priceable
|
||||
GROUP BY token_id, bucket_ts
|
||||
),
|
||||
vol AS (
|
||||
SELECT
|
||||
token_id,
|
||||
bucket_ts,
|
||||
SUM(vol_sats) AS volume_sats,
|
||||
SUM(vol_tokens) AS volume_tokens,
|
||||
COUNT(*) AS tx_count
|
||||
FROM tx_trades
|
||||
GROUP BY token_id, bucket_ts
|
||||
)
|
||||
// Phase 1: fetch all confirmed legs for the time range.
|
||||
let sql = r#"
|
||||
SELECT
|
||||
ohlc.token_id,
|
||||
ohlc.bucket_ts,
|
||||
ohlc.open,
|
||||
ohlc.high,
|
||||
ohlc.low,
|
||||
ohlc.close,
|
||||
vol.volume_sats,
|
||||
vol.volume_tokens,
|
||||
vol.tx_count
|
||||
FROM ohlc
|
||||
JOIN vol ON ohlc.token_id = vol.token_id AND ohlc.bucket_ts = vol.bucket_ts
|
||||
phe.token_id,
|
||||
phe.txid,
|
||||
phe.pool,
|
||||
phe.sats_delta,
|
||||
phe.token_delta,
|
||||
phe.sats,
|
||||
phe.token_amount,
|
||||
phe.sequence,
|
||||
phe.effective_timestamp
|
||||
FROM pool_history_entry AS phe
|
||||
JOIN tx ON tx.txid = phe.txid
|
||||
WHERE tx.blockhash IS NOT NULL
|
||||
AND phe.effective_timestamp >= ?
|
||||
AND phe.effective_timestamp < ?
|
||||
ORDER BY phe.token_id ASC, phe.effective_timestamp ASC, phe.sequence ASC;
|
||||
"#;
|
||||
|
||||
let rows = sqlx::query(select_sql)
|
||||
let rows = sqlx::query(sql)
|
||||
.bind(since_ts)
|
||||
.bind(until_ts)
|
||||
.fetch_all(read_pool)
|
||||
|
|
@ -195,35 +149,132 @@ JOIN vol ON ohlc.token_id = vol.token_id AND ohlc.bucket_ts = vol.bucket_ts
|
|||
return Ok(0);
|
||||
}
|
||||
|
||||
// Phase 2: insert pre-computed rows inside a single transaction.
|
||||
// The write lock is held only for these fast INSERTs, not during aggregation.
|
||||
let mut tx = write_pool.begin().await?;
|
||||
let mut inserted = 0u64;
|
||||
// Phase 2: fold through legs per-token using the policy core to compute buckets.
|
||||
let params = GuardParams::default();
|
||||
|
||||
let mut buckets: HashMap<(Vec<u8>, i64), OhlcvBucket> = HashMap::new();
|
||||
|
||||
let mut current_token: Option<Vec<u8>> = None;
|
||||
let mut policy = Policy::new(params);
|
||||
// Reserve-implied price for each token as the batch opened, used to price a leading
|
||||
// bucket whose legs were all muted so its volume is still recorded.
|
||||
let mut entry_price: HashMap<Vec<u8>, Option<f64>> = HashMap::new();
|
||||
|
||||
for row in &rows {
|
||||
let token_id: Vec<u8> = row.get(0);
|
||||
let bucket_ts: i64 = row.get(1);
|
||||
let open: f64 = row.get(2);
|
||||
let high: f64 = row.get(3);
|
||||
let low: f64 = row.get(4);
|
||||
let close: f64 = row.get(5);
|
||||
let volume_sats: i64 = row.get(6);
|
||||
let volume_tokens: i64 = row.get(7);
|
||||
let tx_count: i64 = row.get(8);
|
||||
let txid: Vec<u8> = row.get(1);
|
||||
let pool_bytes: Vec<u8> = row.get(2);
|
||||
let ts: i64 = row.get(8);
|
||||
let bucket_ts = (ts / 3600) * 3600;
|
||||
|
||||
// A new token starts a new policy — seeded from that token's reserves as they
|
||||
// stood when this batch opened. Callers materialise history in 24-hour batches,
|
||||
// so an unseeded policy would forget everything at every batch boundary and wave
|
||||
// through the first leg of every day of history unjudged.
|
||||
if current_token.as_ref() != Some(&token_id) {
|
||||
current_token = Some(token_id.clone());
|
||||
let snapshot = fetch_reserve_snapshot(read_pool, &token_id, since_ts, true).await?;
|
||||
policy = Policy::seeded(params, &snapshot);
|
||||
entry_price.insert(token_id.clone(), policy.reference());
|
||||
}
|
||||
|
||||
let leg = Leg {
|
||||
txid: to_hash32(&txid)?,
|
||||
pool: to_hash32(&pool_bytes)?,
|
||||
ts,
|
||||
sequence: row.get(7),
|
||||
sats_delta: row.get(3),
|
||||
token_delta: row.get(4),
|
||||
sats: row.get(5),
|
||||
token_amount: row.get(6),
|
||||
};
|
||||
|
||||
// Evaluate the leg through the policy.
|
||||
let judge = policy.judge(&leg);
|
||||
|
||||
let sats_delta_abs = leg.sats_delta.unsigned_abs() as i64;
|
||||
let token_delta_abs = leg.token_delta.unsigned_abs() as i64;
|
||||
|
||||
// Get or create the bucket for this (token, bucket_ts).
|
||||
let bucket_key = (token_id.clone(), bucket_ts);
|
||||
let bucket = buckets.entry(bucket_key).or_insert_with(OhlcvBucket::new);
|
||||
|
||||
// Always accumulate volume.
|
||||
bucket.volume_sats += sats_delta_abs;
|
||||
bucket.volume_tokens += token_delta_abs;
|
||||
bucket.txids.insert(leg.txid);
|
||||
|
||||
// Update OHLC only if accepted and priceable.
|
||||
if judge.accepted() {
|
||||
if let Some(price) = judge.price {
|
||||
if bucket.first_accepted_price.is_none() {
|
||||
bucket.first_accepted_price = Some(price);
|
||||
}
|
||||
bucket.last_accepted_price = Some(price);
|
||||
bucket.high = bucket.high.max(price);
|
||||
bucket.low = bucket.low.min(price);
|
||||
}
|
||||
}
|
||||
|
||||
policy.apply(&leg);
|
||||
}
|
||||
|
||||
// Phase 3: insert all computed buckets in one transaction.
|
||||
//
|
||||
// Ordered by (token, bucket_ts) so a token's buckets can carry their close forward.
|
||||
// A bucket whose legs were all muted or unpriceable has no price of its own, but it
|
||||
// does have real volume, and the guard makes such buckets ordinary rather than
|
||||
// exotic — a manipulation burst filling an hour produces exactly one. Writing it with
|
||||
// the carried close keeps that volume, where writing the raw accumulators would put
|
||||
// the f64::MIN/MAX sentinels into the table and skipping the row entirely would make
|
||||
// the materialised path report zero volume for an hour the raw path reports in full.
|
||||
let mut ordered: Vec<((Vec<u8>, i64), OhlcvBucket)> = buckets.into_iter().collect();
|
||||
ordered.sort_unstable_by(|a, b| a.0.cmp(&b.0));
|
||||
|
||||
let mut tx = write_pool.begin().await?;
|
||||
let mut inserted = 0u64;
|
||||
let mut carry_token: Option<Vec<u8>> = None;
|
||||
let mut carry: Option<f64> = None;
|
||||
|
||||
for ((token_id, bucket_ts), bucket) in ordered {
|
||||
if carry_token.as_ref() != Some(&token_id) {
|
||||
carry_token = Some(token_id.clone());
|
||||
carry = entry_price.get(&token_id).copied().flatten();
|
||||
}
|
||||
|
||||
let (open, close, high, low) = match bucket.last_accepted_price {
|
||||
Some(last) => {
|
||||
carry = Some(last);
|
||||
(
|
||||
bucket.first_accepted_price.unwrap_or(last),
|
||||
last,
|
||||
bucket.high,
|
||||
bucket.low,
|
||||
)
|
||||
}
|
||||
// No price anywhere in this bucket. Flat at the carried close if one exists;
|
||||
// if the token has no known price yet at all, there is nothing to plot and the
|
||||
// row is skipped, as the pre-policy SQL did by inner-joining prices to volume.
|
||||
None => match carry {
|
||||
Some(prev) => (prev, prev, prev, prev),
|
||||
None => continue,
|
||||
},
|
||||
};
|
||||
|
||||
inserted += sqlx::query(
|
||||
"INSERT OR IGNORE INTO ohlcv_1h
|
||||
(token_id, bucket_ts, open, high, low, close, volume_sats, volume_tokens, tx_count)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)",
|
||||
)
|
||||
.bind(token_id)
|
||||
.bind(&token_id)
|
||||
.bind(bucket_ts)
|
||||
.bind(open)
|
||||
.bind(high)
|
||||
.bind(low)
|
||||
.bind(close)
|
||||
.bind(volume_sats)
|
||||
.bind(volume_tokens)
|
||||
.bind(tx_count)
|
||||
.bind(bucket.volume_sats)
|
||||
.bind(bucket.volume_tokens)
|
||||
.bind(bucket.txids.len() as i64)
|
||||
.execute(&mut *tx)
|
||||
.await?
|
||||
.rows_affected();
|
||||
|
|
@ -233,6 +284,31 @@ JOIN vol ON ohlc.token_id = vol.token_id AND ohlc.bucket_ts = vol.bucket_ts
|
|||
Ok(inserted)
|
||||
}
|
||||
|
||||
/// Temporary bucket structure for accumulating OHLCV data per (token, bucket_ts).
|
||||
struct OhlcvBucket {
|
||||
first_accepted_price: Option<f64>,
|
||||
last_accepted_price: Option<f64>,
|
||||
high: f64,
|
||||
low: f64,
|
||||
volume_sats: i64,
|
||||
volume_tokens: i64,
|
||||
txids: std::collections::HashSet<[u8; 32]>,
|
||||
}
|
||||
|
||||
impl OhlcvBucket {
|
||||
fn new() -> Self {
|
||||
Self {
|
||||
first_accepted_price: None,
|
||||
last_accepted_price: None,
|
||||
high: f64::MIN,
|
||||
low: f64::MAX,
|
||||
volume_sats: 0,
|
||||
volume_tokens: 0,
|
||||
txids: std::collections::HashSet::new(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub struct OhlcvRow {
|
||||
pub bucket_ts: i64,
|
||||
pub open: f64,
|
||||
|
|
@ -282,7 +358,7 @@ pub async fn get_active_candles(
|
|||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::db::cauldron::{pool as cauldron_pool, tx, utxo_funding};
|
||||
use crate::db::cauldron::{pool as cauldron_pool, tx, utxo_funding, utxo_spending};
|
||||
use sqlx::sqlite::{SqliteConnectOptions, SqlitePoolOptions};
|
||||
use std::sync::atomic::{AtomicU64, Ordering};
|
||||
|
||||
|
|
@ -300,6 +376,7 @@ mod tests {
|
|||
async fn setup_db(pool: &SqlitePool) {
|
||||
tx::create_table(pool).await;
|
||||
utxo_funding::create_table(pool).await;
|
||||
utxo_spending::create_table(pool).await;
|
||||
cauldron_pool::create_table(pool).await;
|
||||
create_table(pool).await; // ohlcv_1h + idx_phe_txid
|
||||
}
|
||||
|
|
@ -476,6 +553,74 @@ mod tests {
|
|||
assert_eq!(n, 0, "mempool trades must not be materialised");
|
||||
}
|
||||
|
||||
/// Regression: a bucket whose legs are all unpriceable used to be written out with its
|
||||
/// `high`/`low` accumulators still at the `f64::MIN`/`f64::MAX` sentinels, putting
|
||||
/// ±1.8e308 into the table and from there straight onto the chart. Such a bucket has
|
||||
/// no OHLC to report and must simply not be materialised.
|
||||
#[tokio::test]
|
||||
async fn test_rebuild_range_skips_buckets_with_no_priceable_leg() {
|
||||
let pool = test_pool().await;
|
||||
setup_db(&pool).await;
|
||||
|
||||
// A confirmed leg that moves sats but no tokens: real volume, no price.
|
||||
insert_confirmed_trade(
|
||||
&pool, [0x01; 32], [0x02; 32], [0x03; 32], 1727963400, -1000, 0,
|
||||
)
|
||||
.await;
|
||||
|
||||
let n = rebuild_range(&pool, &pool, 1727960400, 1727967600)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(n, 0, "a bucket with no priceable leg must not be written");
|
||||
|
||||
let rows: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM ohlcv_1h")
|
||||
.fetch_one(&pool)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
rows, 0,
|
||||
"no row should exist, let alone one carrying f64 sentinels"
|
||||
);
|
||||
}
|
||||
|
||||
/// Once a token has a known price, an hour whose legs are all unpriceable or muted
|
||||
/// still carries real volume. Dropping the row made the materialised path report zero
|
||||
/// volume for an hour the raw path reports in full — and the guard makes such hours
|
||||
/// ordinary, since one manipulation burst can fill an entire bucket.
|
||||
#[tokio::test]
|
||||
async fn test_rebuild_range_keeps_volume_for_a_bucket_with_no_accepted_price() {
|
||||
let pool = test_pool().await;
|
||||
setup_db(&pool).await;
|
||||
let token = [0x03_u8; 32];
|
||||
|
||||
// Hour 1727960400: an ordinary priced trade at 1000/25 = 40.
|
||||
insert_confirmed_trade(&pool, [0x01; 32], [0x02; 32], token, 1727960500, -1000, 25).await;
|
||||
// Hour 1727964000: tokens move but no sats, so the leg has no price.
|
||||
insert_confirmed_trade(&pool, [0x11; 32], [0x12; 32], token, 1727964100, 0, 25).await;
|
||||
|
||||
let n = rebuild_range(&pool, &pool, 1727960400, 1727967600)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(n, 2, "both hours must be materialised");
|
||||
|
||||
let (volume, open, close, high, low): (i64, f64, f64, f64, f64) = sqlx::query_as(
|
||||
"SELECT volume_tokens, open, close, high, low FROM ohlcv_1h WHERE bucket_ts = 1727964000",
|
||||
)
|
||||
.fetch_one(&pool)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(
|
||||
volume, 25,
|
||||
"volume must survive a bucket with no accepted price"
|
||||
);
|
||||
assert_eq!(
|
||||
(open, close, high, low),
|
||||
(40.0, 40.0, 40.0, 40.0),
|
||||
"the bucket must sit flat at the carried close, not at a sentinel"
|
||||
);
|
||||
}
|
||||
|
||||
/// A multi-pool arbitrage transaction whose legs nearly cancel must materialise the
|
||||
/// price its legs executed at, not the signed-net ratio.
|
||||
#[tokio::test]
|
||||
|
|
@ -511,16 +656,31 @@ mod tests {
|
|||
.await
|
||||
.unwrap();
|
||||
|
||||
let close: f64 = sqlx::query_scalar("SELECT close FROM ohlcv_1h WHERE token_id = ?")
|
||||
.bind(token.as_slice())
|
||||
.fetch_one(&pool)
|
||||
.await
|
||||
.unwrap();
|
||||
let (close, open, high, low): (f64, f64, f64, f64) =
|
||||
sqlx::query_as("SELECT close, open, high, low FROM ohlcv_1h WHERE token_id = ?")
|
||||
.bind(token.as_slice())
|
||||
.fetch_one(&pool)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
// Per-leg pricing: each leg executes at its own price.
|
||||
// Leg 1 (buy): 446,491,239 / 1,334,527,069 = 0.334569
|
||||
// Leg 2 (sell): 384,906,040 / 1,334,527,067 = 0.288421
|
||||
// Close is the last leg, open is the first, high/low are the extremes.
|
||||
let leg1_price = 446_491_239.0 / 1_334_527_069.0;
|
||||
let leg2_price = 384_906_040.0 / 1_334_527_067.0;
|
||||
|
||||
let expected = 831_397_279.0 / 2_669_054_136.0;
|
||||
assert!(
|
||||
(close - expected).abs() < 1e-9,
|
||||
"materialised close {close} should be the gross ratio {expected}"
|
||||
(close - leg2_price).abs() < 1e-9,
|
||||
"materialised close {close} should be leg2 price {leg2_price}"
|
||||
);
|
||||
assert!(
|
||||
(open - leg1_price).abs() < 1e-9,
|
||||
"materialised open {open} should be leg1 price {leg1_price}"
|
||||
);
|
||||
assert!(
|
||||
(high - leg1_price).abs() < 1e-9 && (low - leg2_price).abs() < 1e-9,
|
||||
"high {high} and low {low} should span the leg prices"
|
||||
);
|
||||
}
|
||||
|
||||
|
|
|
|||
17
src/main.rs
17
src/main.rs
|
|
@ -515,18 +515,33 @@ async fn launch() -> _ {
|
|||
}
|
||||
|
||||
// Background task: incrementally materialise new 1-hour OHLCV buckets as blocks arrive.
|
||||
// The full historical backfill above already ran; this task only handles the tail.
|
||||
// After a version wipe this is also what repopulates history, since the synchronous
|
||||
// backfill above is skipped in that case.
|
||||
// Only processes buckets older than 3 hours (well beyond BCH reorg depth).
|
||||
{
|
||||
let ohlcv_write = dbpool.cauldron_w.clone();
|
||||
let ohlcv_read = dbpool.cauldron_r.clone();
|
||||
let ohlcv_state_bg = ohlcv_state.clone();
|
||||
let ohlcv_ibd = ibd_state.clone();
|
||||
tokio::spawn(async move {
|
||||
// Batch size: 1 day per SQL call to keep each write short.
|
||||
const BATCH_SECS: i64 = 24 * 3600;
|
||||
// Safety margin: only materialise buckets older than this many seconds.
|
||||
const SAFETY_SECS: i64 = 3 * 3600;
|
||||
|
||||
// Wait for IBD before materialising anything.
|
||||
//
|
||||
// The synchronous backfill above waits too, but it is skipped whenever
|
||||
// `ohlcv_wiped` is set — which includes every fresh database, since a missing
|
||||
// version key reads as stale. Without this the task would sweep from the
|
||||
// first trade all the way to `now - 3h` while indexing is still years behind,
|
||||
// writing nothing, contending with block writes for the cauldron write lock,
|
||||
// and advancing `materialized_end` to roughly now against an empty table — at
|
||||
// which point `candlesticks()` would take the fast path over nothing.
|
||||
while !ohlcv_ibd.initial_sync_complete.load(Ordering::Relaxed) {
|
||||
tokio::time::sleep(Duration::from_secs(1)).await;
|
||||
}
|
||||
|
||||
loop {
|
||||
let now = crate::timeutil::time_now();
|
||||
// Floor to 1-hour boundary, 3 hours ago.
|
||||
|
|
|
|||
|
|
@ -8,6 +8,7 @@ use crate::db::cauldron::{
|
|||
pool::{self, dummy_init_seq, insert_new_pool},
|
||||
tx::{self, insert_block_tx, insert_mempool_tx},
|
||||
utxo_funding::{self, insert_utxo_funding},
|
||||
utxo_spending,
|
||||
};
|
||||
use crate::utiltest::mock_db_pool;
|
||||
use crate::OhlcvState;
|
||||
|
|
@ -64,8 +65,10 @@ fn dummy_cauldron(
|
|||
|
||||
async fn setup_mock_db(pool: sqlx::SqlitePool) {
|
||||
utxo_funding::create_table(&pool).await;
|
||||
utxo_spending::create_table(&pool).await;
|
||||
tx::create_table(&pool).await;
|
||||
pool::create_table(&pool).await;
|
||||
ohlcv::create_table(&pool).await;
|
||||
dummy_init_seq();
|
||||
|
||||
let mut conn = pool.acquire().await.unwrap();
|
||||
|
|
@ -200,6 +203,7 @@ async fn setup_mock_db(pool: sqlx::SqlitePool) {
|
|||
|
||||
async fn setup_seed_db(pool: sqlx::SqlitePool) {
|
||||
utxo_funding::create_table(&pool).await;
|
||||
utxo_spending::create_table(&pool).await;
|
||||
tx::create_table(&pool).await;
|
||||
pool::create_table(&pool).await;
|
||||
ohlcv::create_table(&pool).await;
|
||||
|
|
@ -378,8 +382,10 @@ async fn test_multiple_candlesticks_endpoint() {
|
|||
async fn test_single_swap_multiple_pools() {
|
||||
let mock_db = mock_db_pool(|pool: sqlx::SqlitePool| async move {
|
||||
utxo_funding::create_table(&pool).await;
|
||||
utxo_spending::create_table(&pool).await;
|
||||
tx::create_table(&pool).await;
|
||||
pool::create_table(&pool).await;
|
||||
ohlcv::create_table(&pool).await;
|
||||
dummy_init_seq();
|
||||
|
||||
let mut conn = pool.acquire().await.unwrap();
|
||||
|
|
@ -479,8 +485,10 @@ async fn test_candle_carries_price_when_net_zero_tokens_but_has_volume() {
|
|||
let mock_db = mock_db_pool(|pool: sqlx::SqlitePool| async move {
|
||||
// --- boilerplate setup ---
|
||||
utxo_funding::create_table(&pool).await;
|
||||
utxo_spending::create_table(&pool).await;
|
||||
tx::create_table(&pool).await;
|
||||
pool::create_table(&pool).await;
|
||||
ohlcv::create_table(&pool).await;
|
||||
dummy_init_seq();
|
||||
|
||||
let mut conn = pool.acquire().await.unwrap();
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue