use std::collections::{BTreeMap, HashMap, VecDeque};
use std::convert::Infallible;
use anyhow::Result;
use bsv_sdk::transaction::MerklePath;
use bsv_sdk::wallet::{InternalizeActionArgs, WalletInterface};
use bsv_tracker::{
CheckError, Clock, Evidence, Header, HeaderError, Headers, Height, Hint, HintSource,
HintStatus, HostAction, Input, Params, Proof, ProofFetcher, Reask, State, Timestamp, TxId,
Word,
};
use bsv_wallet_toolbox::services::PostBeefResult;
use bsv_wallet_toolbox::{BroadcastStatus, StorageSqlx, WalletServices};
use serde::{Deserialize, Serialize};
use sqlx::Row as _;
pub const DEFAULT_AGE_SECS: u64 = 600;
pub const DEFAULT_MAX_ASKS: usize = 20;
const MAX_BACKOFF_DOUBLINGS: u32 = 6;
const MAX_ASKS_PER_TX: u32 = 16;
const TIP_RING: usize = 12;
#[derive(Debug, Clone)]
pub struct HeaderSnapshot {
tip: std::result::Result<(Height, String), String>,
headers: BTreeMap<Height, std::result::Result<Header, String>>,
}
impl HeaderSnapshot {
pub async fn take<V: WalletServices + ?Sized>(
services: &V,
heights: impl IntoIterator<Item = Height>,
) -> Self {
let tip = services
.get_chain_tip_header()
.await
.map(|h| (h.height, h.hash.to_ascii_lowercase()))
.map_err(|e| e.to_string());
let mut snapshot = Self {
tip,
headers: BTreeMap::new(),
};
snapshot.extend(services, heights).await;
snapshot
}
pub fn unread() -> Self {
Self {
tip: Err("the header service was not read".to_string()),
headers: BTreeMap::new(),
}
}
pub async fn extend<V: WalletServices + ?Sized>(
&mut self,
services: &V,
heights: impl IntoIterator<Item = Height>,
) {
for height in heights {
if self.headers.contains_key(&height) {
continue;
}
let answer = match services.get_header_for_height(height).await {
Ok(bytes) => header_of(&bytes),
Err(e) => Err(e.to_string()),
};
self.headers.insert(height, answer);
}
}
pub fn tip(&self) -> Option<(Height, &str)> {
self.tip.as_ref().ok().map(|(h, hash)| (*h, hash.as_str()))
}
}
impl Headers for HeaderSnapshot {
fn header_at(&self, height: Height) -> std::result::Result<Option<Header>, HeaderError> {
match self.headers.get(&height) {
Some(Ok(header)) => Ok(Some(header.clone())),
Some(Err(fault)) => Err(HeaderError(fault.clone())),
None => Ok(None),
}
}
fn tip_height(&self) -> std::result::Result<Height, HeaderError> {
self.tip
.as_ref()
.map(|(height, _)| *height)
.map_err(|fault| HeaderError(fault.clone()))
}
}
fn header_of(bytes: &[u8]) -> std::result::Result<Header, String> {
if bytes.len() != 80 {
return Err(format!("a header of {} bytes", bytes.len()));
}
let display = |raw: &[u8]| {
let mut reversed = raw.to_vec();
reversed.reverse();
hex::encode(reversed)
};
Ok(Header {
hash: display(&bsv_sdk::primitives::hash::sha256d(bytes)),
merkle_root: display(&bytes[36..68]),
})
}
#[derive(Debug, Clone, Copy, Default)]
pub struct SystemClock;
impl Clock for SystemClock {
fn now(&self) -> Timestamp {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0)
}
}
#[derive(Debug, Default)]
pub struct MemoryHints {
queue: VecDeque<(TxId, Hint)>,
}
impl MemoryHints {
async fn read(storage: &StorageSqlx, rows: &[TrackedRow]) -> Result<Self> {
let txids: Vec<String> = rows.iter().map(|r| r.txid.clone()).collect();
if txids.is_empty() {
return Ok(Self::default());
}
let mut records = storage.broadcast_records(None, &txids).await?;
records.sort_by_key(|r| r.seen_at);
let mut queue = VecDeque::new();
for record in records {
let Some(status) = record.ladder_status() else {
continue;
};
let txid = record.txid.to_ascii_lowercase();
let Some(row) = rows.iter().find(|r| r.txid == txid) else {
continue;
};
let (label, status) = match status {
BroadcastStatus::Accepted => ("accepted", HintStatus::Accepted),
BroadcastStatus::Seen => ("seen", HintStatus::Seen),
BroadcastStatus::Mined => ("mined", HintStatus::Unknown),
BroadcastStatus::Rejected => (
"rejected",
HintStatus::Rejected {
reason: record.status.clone(),
},
),
BroadcastStatus::Unknown => ("unknown", HintStatus::Unknown),
};
let source = format!("{}|{}", record.provider, label);
if row.hints.iter().any(|h| h.source == source) {
continue;
}
let observed = record.seen_at.timestamp().max(0) as Timestamp;
queue.push_back((txid, Hint::new(source, status, observed)));
}
Ok(Self { queue })
}
}
impl HintSource for MemoryHints {
type Error = Infallible;
fn next_hint(&mut self) -> std::result::Result<Option<(TxId, Hint)>, Self::Error> {
Ok(self.queue.pop_front())
}
}
#[derive(Debug, Default)]
pub struct Courier {
carried: HashMap<TxId, std::result::Result<Option<Proof>, String>>,
}
impl Courier {
pub async fn carry<V: WalletServices + ?Sized>(services: &V, txids: &[TxId]) -> Self {
let mut carried = HashMap::new();
for txid in txids {
let answer = match services.get_merkle_path(txid, false).await {
Ok(found) => match found.merkle_path {
Some(hex) => MerklePath::from_hex(&hex)
.map_err(|e| format!("not a merkle path: {e}"))
.and_then(|path| Proof::new(txid.clone(), path).map_err(|e| e.to_string()))
.map(Some),
None => Ok(None),
},
Err(e) => Err(e.to_string()),
};
carried.insert(txid.clone(), answer);
}
Self { carried }
}
fn heights(&self) -> Vec<Height> {
self.carried
.values()
.filter_map(|answer| answer.as_ref().ok()?.as_ref().map(Proof::height))
.collect()
}
}
impl ProofFetcher for Courier {
type Error = String;
fn fetch(
&mut self,
txid: &str,
_reask: &Reask,
) -> std::result::Result<Option<Proof>, Self::Error> {
self.carried.remove(txid).unwrap_or(Ok(None))
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
struct StoredHint {
source: String,
status: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
height: Option<Height>,
#[serde(default, skip_serializing_if = "Option::is_none")]
reason: Option<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
competitors: Vec<TxId>,
observed: Timestamp,
}
impl StoredHint {
fn of(hint: &Hint) -> Self {
let (status, height, reason, competitors) = match &hint.status {
HintStatus::Accepted => ("accepted", None, None, vec![]),
HintStatus::Seen => ("seen", None, None, vec![]),
HintStatus::Mined { height } => ("mined", Some(*height), None, vec![]),
HintStatus::StaleBlock => ("stale_block", None, None, vec![]),
HintStatus::Rejected { reason } => ("rejected", None, Some(reason.clone()), vec![]),
HintStatus::DoubleSpend { competitors } => {
("double_spend", None, None, competitors.clone())
}
HintStatus::OrphanMempool => ("orphan_mempool", None, None, vec![]),
HintStatus::Unknown => ("unknown", None, None, vec![]),
};
Self {
source: hint.source.clone(),
status: status.to_string(),
height,
reason,
competitors,
observed: hint.observed,
}
}
fn hint(&self) -> Hint {
let status = match self.status.as_str() {
"accepted" => HintStatus::Accepted,
"seen" => HintStatus::Seen,
"mined" => self
.height
.map_or(HintStatus::Unknown, |height| HintStatus::Mined { height }),
"stale_block" => HintStatus::StaleBlock,
"rejected" => HintStatus::Rejected {
reason: self.reason.clone().unwrap_or_default(),
},
"double_spend" => HintStatus::DoubleSpend {
competitors: self.competitors.clone(),
},
"orphan_mempool" => HintStatus::OrphanMempool,
_ => HintStatus::Unknown,
};
Hint::new(self.source.clone(), status, self.observed)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct TrackedRow {
txid: TxId,
built: bool,
hints: Vec<StoredHint>,
bump: Option<String>,
asks: u32,
next_ask_at: Option<Timestamp>,
reannounce_beef: Option<Vec<u8>>,
reannounce_at: Option<Timestamp>,
on_accept: Option<String>,
}
impl TrackedRow {
fn new(txid: &str) -> Self {
Self {
txid: txid.to_ascii_lowercase(),
built: true,
hints: vec![],
bump: None,
asks: 0,
next_ask_at: None,
reannounce_beef: None,
reannounce_at: None,
on_accept: None,
}
}
fn proof(&self) -> Option<std::result::Result<Proof, CheckError>> {
let hex = self.bump.as_deref()?;
Some(
MerklePath::from_hex(hex)
.map_err(|e| CheckError::InvalidProof(e.to_string()))
.and_then(|path| Proof::new(self.txid.clone(), path)),
)
}
fn proof_height(&self) -> Option<Height> {
self.proof()?.ok().map(|p| p.height())
}
fn replay<H: Headers>(&self, params: &Params, headers: &H) -> (State, Option<CheckError>) {
let mut state = State::new(self.txid.clone());
if self.built {
let _ = state.step(params, headers, Input::Host(HostAction::Build));
}
for hint in &self.hints {
let _ = state.step(params, headers, Input::Hint(hint.hint()));
}
let fault = match self.proof() {
Some(Ok(proof)) => state
.step(params, headers, Input::Evidence(Evidence::Proof(proof)))
.err(),
Some(Err(fault)) => Some(fault),
None => None,
};
(state, fault)
}
}
pub fn word_label(word: &Word) -> &'static str {
match word {
Word::Unknown => "unknown",
Word::Built => "built",
Word::Announced => "announced",
Word::Seen => "seen",
Word::Mined(_) => "mined",
Word::Stale => "stale",
Word::Rejected { .. } => "rejected",
Word::Conflicted { .. } => "conflicted",
Word::Abandoned => "abandoned",
}
}
fn reask_label(reask: &Reask) -> String {
match reask {
Reask::Reorg { fork_height } => format!("reorg@{fork_height}"),
Reask::Fork { height } => format!("fork@{height}"),
Reask::Hint(hint) => format!("hint:{}", hint.source),
Reask::ProofFailed { height } => format!("proof_failed@{height}"),
Reask::Recheck => "recheck".to_string(),
Reask::Age => "age".to_string(),
Reask::Spend => "spend".to_string(),
}
}
async fn ensure_tables(storage: &StorageSqlx) -> Result<()> {
for sql in [
"CREATE TABLE IF NOT EXISTS tracker_states (\
txid TEXT PRIMARY KEY, \
built INTEGER NOT NULL DEFAULT 1, \
hints TEXT NOT NULL DEFAULT '[]', \
bump TEXT, \
word TEXT NOT NULL DEFAULT 'unknown', \
height INTEGER, \
header_hash TEXT, \
reask TEXT, \
asks INTEGER NOT NULL DEFAULT 0, \
next_ask_at INTEGER, \
reannounce_beef BLOB, \
reannounce_at INTEGER, \
on_accept TEXT, \
updated_at INTEGER NOT NULL DEFAULT 0)",
"CREATE INDEX IF NOT EXISTS tracker_states_word ON tracker_states (word)",
"CREATE INDEX IF NOT EXISTS tracker_states_height ON tracker_states (height)",
"CREATE TABLE IF NOT EXISTS tracker_meta (key TEXT PRIMARY KEY, value TEXT NOT NULL)",
] {
sqlx::query(sql).execute(storage.pool()).await?;
}
Ok(())
}
const ROW_COLUMNS: &str =
"txid, built, hints, bump, asks, next_ask_at, reannounce_beef, reannounce_at, on_accept";
fn row_of(row: &sqlx::sqlite::SqliteRow) -> TrackedRow {
let hints: String = row.get("hints");
TrackedRow {
txid: row.get("txid"),
built: row.get::<i64, _>("built") != 0,
hints: serde_json::from_str(&hints).unwrap_or_default(),
bump: row.get("bump"),
asks: row.get::<i64, _>("asks").max(0) as u32,
next_ask_at: row
.get::<Option<i64>, _>("next_ask_at")
.map(|t| t.max(0) as Timestamp),
reannounce_beef: row.get("reannounce_beef"),
reannounce_at: row
.get::<Option<i64>, _>("reannounce_at")
.map(|t| t.max(0) as Timestamp),
on_accept: row.get("on_accept"),
}
}
async fn load_row(storage: &StorageSqlx, txid: &str) -> Result<Option<TrackedRow>> {
let row = sqlx::query(&format!(
"SELECT {ROW_COLUMNS} FROM tracker_states WHERE txid = ?"
))
.bind(txid.to_ascii_lowercase())
.fetch_optional(storage.pool())
.await?;
Ok(row.as_ref().map(row_of))
}
async fn store_row(
storage: &StorageSqlx,
row: &TrackedRow,
state: &State,
now: Timestamp,
) -> Result<()> {
let (height, header_hash) = match state.word() {
Word::Mined(mined) => (
Some(mined.height() as i64),
Some(mined.checked().header().hash.clone()),
),
_ => (None, None),
};
sqlx::query(
"INSERT INTO tracker_states \
(txid, built, hints, bump, word, height, header_hash, reask, asks, next_ask_at, \
reannounce_beef, reannounce_at, on_accept, updated_at) \
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) \
ON CONFLICT(txid) DO UPDATE SET built = excluded.built, hints = excluded.hints, \
bump = excluded.bump, word = excluded.word, height = excluded.height, \
header_hash = excluded.header_hash, reask = excluded.reask, asks = excluded.asks, \
next_ask_at = excluded.next_ask_at, reannounce_beef = excluded.reannounce_beef, \
reannounce_at = excluded.reannounce_at, on_accept = excluded.on_accept, \
updated_at = excluded.updated_at",
)
.bind(&row.txid)
.bind(row.built as i64)
.bind(serde_json::to_string(&row.hints)?)
.bind(&row.bump)
.bind(word_label(state.word()))
.bind(height)
.bind(header_hash)
.bind(state.reask().map(reask_label))
.bind(row.asks as i64)
.bind(row.next_ask_at.map(|t| t as i64))
.bind(&row.reannounce_beef)
.bind(row.reannounce_at.map(|t| t as i64))
.bind(&row.on_accept)
.bind(now as i64)
.execute(storage.pool())
.await?;
Ok(())
}
pub async fn stored_word(
storage: &StorageSqlx,
txid: &str,
) -> Result<Option<(String, Option<Height>, Option<String>)>> {
ensure_tables(storage).await?;
let row = sqlx::query("SELECT word, height, reask FROM tracker_states WHERE txid = ?")
.bind(txid.to_ascii_lowercase())
.fetch_optional(storage.pool())
.await?;
Ok(row.map(|r| {
(
r.get("word"),
r.get::<Option<i64>, _>("height").map(|h| h as Height),
r.get("reask"),
)
}))
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum PostAnswer {
Accepted { provider: String },
NotFinal { provider: String, said: String },
Refused { said: String },
}
pub fn says_not_final(text: &str) -> bool {
let lower = text.to_ascii_lowercase();
["non-final", "nonfinal", "non final", "not final"]
.iter()
.any(|word| lower.contains(word))
|| lower
.split(|c: char| !c.is_ascii_alphanumeric())
.any(|token| token == "476")
}
pub fn read_post(results: &[PostBeefResult]) -> PostAnswer {
if let Some(accepted) = results.iter().find(|r| r.is_success()) {
return PostAnswer::Accepted {
provider: accepted.name.clone(),
};
}
let said: Vec<(String, String)> = results
.iter()
.map(|r| {
let detail = r
.error
.clone()
.or_else(|| r.txid_results.iter().find_map(|t| t.data.clone()))
.unwrap_or_else(|| r.status.clone());
(r.name.clone(), detail)
})
.collect();
if let Some((provider, detail)) = said.iter().find(|(_, detail)| says_not_final(detail)) {
return PostAnswer::NotFinal {
provider: provider.clone(),
said: detail.clone(),
};
}
PostAnswer::Refused {
said: if said.is_empty() {
"none answered".to_string()
} else {
said.iter()
.map(|(name, detail)| format!("{name}: {detail}"))
.collect::<Vec<_>>()
.join("; ")
},
}
}
#[derive(Debug, Clone)]
pub struct Reannounce {
pub beef: Vec<u8>,
pub at: Timestamp,
pub on_accept: Option<InternalizeActionArgs>,
}
pub async fn heard<C: Clock>(
storage: &StorageSqlx,
clock: &C,
opts: &TickOptions,
txid: &str,
hint: Hint,
not_before: Option<Timestamp>,
reannounce: Option<Reannounce>,
) -> Result<(String, Option<String>)> {
ensure_tables(storage).await?;
let mut row = match load_row(storage, txid).await? {
Some(row) if row.bump.is_some() => {
let stored = stored_word(storage, txid).await?.unwrap_or_default();
return Ok((stored.0, stored.2));
}
Some(row) => row,
None => TrackedRow::new(txid),
};
let params = opts.params();
let headers = HeaderSnapshot::unread();
let (mut state, _) = row.replay(¶ms, &headers);
if !row.hints.iter().any(|h| h.source == hint.source) {
row.hints.push(StoredHint::of(&hint));
let _ = state.step(¶ms, &headers, Input::Hint(hint));
}
row.asks = 0;
row.next_ask_at = not_before;
match reannounce {
Some(again) => {
row.reannounce_beef = Some(again.beef);
row.reannounce_at = Some(again.at);
row.on_accept = match again.on_accept {
Some(args) => Some(serde_json::to_string(&args)?),
None => None,
};
}
None => {
row.reannounce_beef = None;
row.reannounce_at = None;
row.on_accept = None;
}
}
store_row(storage, &row, &state, clock.now()).await?;
Ok((
word_label(state.word()).to_string(),
state.reask().map(reask_label),
))
}
#[derive(Debug, Clone)]
pub struct TickOptions {
pub age_threshold: Timestamp,
pub max_asks: usize,
}
impl TickOptions {
pub fn from_env() -> Self {
fn env<T: std::str::FromStr>(key: &str, default: T) -> T {
std::env::var(key)
.ok()
.and_then(|v| v.trim().parse().ok())
.unwrap_or(default)
}
Self {
age_threshold: env("TRACKER_AGE_SECS", DEFAULT_AGE_SECS).max(1),
max_asks: env("TRACKER_MAX_ASKS", DEFAULT_MAX_ASKS),
}
}
fn params(&self) -> Params {
Params {
age_threshold: self.age_threshold,
}
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)]
pub struct TickReport {
pub tracked: usize,
pub adopted: usize,
pub hints: usize,
pub asked: usize,
pub mined: Vec<String>,
pub no_proof: Vec<String>,
pub moved: Vec<String>,
pub moved_from: Option<Height>,
pub reannounced: Vec<String>,
pub not_final: Vec<String>,
#[serde(skip)]
pub on_accept: Vec<(String, String)>,
pub faults: Vec<String>,
pub stored: Vec<String>,
}
impl TickReport {
pub fn is_quiet(&self) -> bool {
self.adopted == 0
&& self.hints == 0
&& self.asked == 0
&& self.reannounced.is_empty()
&& self.not_final.is_empty()
&& self.moved.is_empty()
&& self.faults.is_empty()
}
pub fn summary(&self) -> String {
format!(
"tracker tick: {} tracked ({} new), {} hint(s), {} announced again ({} still not final), {} re-ask(s): {} mined, {} without a proof, {} moved, {} fault(s)",
self.tracked,
self.adopted,
self.hints,
self.reannounced.len(),
self.not_final.len(),
self.asked,
self.mined.len(),
self.no_proof.len(),
self.moved.len(),
self.faults.len(),
)
}
}
async fn read_ring(storage: &StorageSqlx) -> Vec<(Height, String)> {
sqlx::query("SELECT value FROM tracker_meta WHERE key = 'tips'")
.fetch_optional(storage.pool())
.await
.ok()
.flatten()
.and_then(|r| serde_json::from_str(&r.get::<String, _>("value")).ok())
.unwrap_or_default()
}
async fn write_ring(storage: &StorageSqlx, ring: &[(Height, String)]) -> Result<()> {
sqlx::query(
"INSERT INTO tracker_meta (key, value) VALUES ('tips', ?) \
ON CONFLICT(key) DO UPDATE SET value = excluded.value",
)
.bind(serde_json::to_string(ring)?)
.execute(storage.pool())
.await?;
Ok(())
}
async fn moved_since<V: WalletServices + ?Sized>(
services: &V,
snapshot: &mut HeaderSnapshot,
mut ring: Vec<(Height, String)>,
) -> (Option<Height>, Vec<(Height, String)>) {
let Some((tip_height, tip_hash)) = snapshot.tip().map(|(h, hash)| (h, hash.to_string())) else {
return (None, ring);
};
let mut moved = None;
while let Some((height, hash)) = ring.last().cloned() {
let active = if height == tip_height {
Some(tip_hash.clone())
} else if height > tip_height {
None
} else {
snapshot.extend(services, [height]).await;
match snapshot.header_at(height) {
Ok(Some(header)) => Some(header.hash),
_ => return (moved, ring),
}
};
if active.as_deref() == Some(hash.as_str()) {
break;
}
moved = Some(height);
ring.pop();
}
if ring.last().map(|(h, _)| *h) != Some(tip_height) {
ring.push((tip_height, tip_hash));
}
if ring.len() > TIP_RING {
let extra = ring.len() - TIP_RING;
ring.drain(..extra);
}
(moved, ring)
}
pub async fn tick<V: WalletServices + ?Sized, C: Clock>(
storage: &StorageSqlx,
services: &V,
clock: &C,
opts: &TickOptions,
) -> Result<TickReport> {
ensure_tables(storage).await?;
let params = opts.params();
let now = clock.now();
let mut report = TickReport::default();
let adopted = sqlx::query(
"INSERT OR IGNORE INTO tracker_states (txid, word, updated_at) \
SELECT DISTINCT lower(txid), 'built', ? FROM transactions \
WHERE txid IS NOT NULL AND status IN ('unproven', 'sending')",
)
.bind(now as i64)
.execute(storage.pool())
.await?;
report.adopted = adopted.rows_affected() as usize;
let mut snapshot = HeaderSnapshot::take(services, []).await;
let (moved_from, ring) = moved_since(services, &mut snapshot, read_ring(storage).await).await;
report.moved_from = moved_from;
let mut rows: Vec<TrackedRow> = sqlx::query(&format!(
"SELECT {ROW_COLUMNS} FROM tracker_states \
WHERE word IN ('unknown', 'built', 'announced', 'seen', 'stale') \
OR (word = 'mined' AND ? IS NOT NULL AND height >= ?) \
ORDER BY txid"
))
.bind(moved_from.map(|h| h as i64))
.bind(moved_from.map(|h| h as i64))
.fetch_all(storage.pool())
.await?
.iter()
.map(row_of)
.collect();
report.tracked = rows.len();
snapshot
.extend(services, rows.iter().filter_map(TrackedRow::proof_height))
.await;
let mut hints = MemoryHints::read(storage, &rows).await?;
let mut states: HashMap<TxId, State> = HashMap::new();
for row in rows.iter_mut() {
let (state, fault) = row.replay(¶ms, &snapshot);
if let Some(fault) = fault {
if matches!(fault, CheckError::RootMismatch(_)) {
row.bump = None;
report.moved.push(row.txid.clone());
tracing::warn!(
marker = "tracker_stored_proof_moved",
txid = %row.txid,
"a stored merkle path no longer meets the active header; asking again"
);
} else {
report.faults.push(format!("{}: {fault}", row.txid));
}
row.asks = 0;
row.next_ask_at = None;
}
states.insert(row.txid.clone(), state);
}
while let Ok(Some((txid, hint))) = hints.next_hint() {
let (Some(state), Some(row)) = (
states.get_mut(&txid),
rows.iter_mut().find(|r| r.txid == txid),
) else {
continue;
};
row.hints.push(StoredHint::of(&hint));
let _ = state.step(¶ms, &snapshot, Input::Hint(hint));
report.hints += 1;
if matches!(state.reask(), Some(Reask::Hint(_))) {
row.asks = 0;
row.next_ask_at = None;
}
}
for row in rows.iter_mut() {
let (Some(beef), Some(at)) = (row.reannounce_beef.clone(), row.reannounce_at) else {
continue;
};
let Some(state) = states.get_mut(&row.txid) else {
continue;
};
if at > now || matches!(state.word(), Word::Mined(_)) {
continue;
}
let pause = Some(now.saturating_add(opts.age_threshold));
let answer = services
.post_beef(&beef, std::slice::from_ref(&row.txid))
.await
.map(|results| read_post(&results));
match answer {
Ok(PostAnswer::Accepted { provider }) => {
let hint = Hint::new(format!("{provider}|accepted"), HintStatus::Accepted, now);
row.hints.push(StoredHint::of(&hint));
let _ = state.step(¶ms, &snapshot, Input::Hint(hint));
row.reannounce_beef = None;
row.reannounce_at = None;
row.asks = 0;
row.next_ask_at = pause;
report.reannounced.push(row.txid.clone());
if let Some(args) = row.on_accept.take() {
report.on_accept.push((row.txid.clone(), args));
}
}
Ok(PostAnswer::NotFinal { .. }) => {
row.asks = row.asks.saturating_add(1);
row.next_ask_at = pause;
if row.asks >= MAX_ASKS_PER_TX {
row.reannounce_at = None;
report.faults.push(format!(
"{}: still not final after {MAX_ASKS_PER_TX} announces; not announced again",
row.txid
));
} else {
row.reannounce_at = pause;
report.not_final.push(row.txid.clone());
}
}
Ok(PostAnswer::Refused { said }) => {
row.reannounce_at = None;
row.next_ask_at = pause;
report.faults.push(format!(
"{}: refused when announced again: {said}",
row.txid
));
}
Err(e) => {
row.reannounce_at = pause;
report
.faults
.push(format!("{}: could not announce again: {e}", row.txid));
}
}
}
for state in states.values_mut() {
let _ = state.step(¶ms, &snapshot, Input::Host(HostAction::Tick(now)));
}
let due: Vec<TxId> = rows
.iter()
.filter(|row| {
states[&row.txid].reask().is_some()
&& row.asks < MAX_ASKS_PER_TX
&& row.next_ask_at.is_none_or(|at| at <= now)
})
.take(opts.max_asks)
.map(|row| row.txid.clone())
.collect();
let mut courier = Courier::carry(services, &due).await;
snapshot.extend(services, courier.heights()).await;
for txid in &due {
let (Some(state), Some(row)) = (
states.get_mut(txid),
rows.iter_mut().find(|r| &r.txid == txid),
) else {
continue;
};
report.asked += 1;
let reask = state.reask().cloned().unwrap_or(Reask::Age);
let fetched = courier.fetch(txid, &reask);
let proven = match fetched {
Ok(Some(proof)) => {
let hex = proof.path().to_hex();
match state.step(¶ms, &snapshot, Input::Evidence(Evidence::Proof(proof))) {
Ok(()) => {
row.bump = Some(hex);
true
}
Err(fault) => {
report.faults.push(format!("{txid}: {fault}"));
false
}
}
}
Ok(None) => {
report.no_proof.push(txid.clone());
false
}
Err(fault) => {
report
.faults
.push(format!("{txid}: could not ask for a proof: {fault}"));
false
}
};
if proven {
row.asks = 0;
row.next_ask_at = None;
report.mined.push(txid.clone());
} else {
let pause = opts
.age_threshold
.saturating_mul(1 << row.asks.min(MAX_BACKOFF_DOUBLINGS));
row.asks = row.asks.saturating_add(1);
row.next_ask_at = Some(now.saturating_add(pause));
}
}
for row in &rows {
let state = &states[&row.txid];
store_row(storage, row, state, now).await?;
if let (true, Word::Mined(mined)) = (report.mined.contains(&row.txid), state.word()) {
let checked = mined.checked();
let outcome = storage
.ingest_merkle_proof(
&row.txid,
&checked.proof().path().to_binary(),
checked.height(),
&checked.header().hash,
Some(checked.root()),
)
.await;
report.stored.push(match outcome {
Ok(outcome) => format!("{}: {outcome:?}", row.txid),
Err(e) => format!("{}: not stored: {e}", row.txid),
});
}
}
write_ring(storage, &ring).await?;
Ok(report)
}
pub async fn record_accepted<W: WalletInterface>(wallet: &W, report: &mut TickReport) {
for (txid, args) in std::mem::take(&mut report.on_accept) {
let line = match serde_json::from_str::<InternalizeActionArgs>(&args) {
Ok(args) => match wallet.internalize_action(args, "bsv-wallet-cli").await {
Ok(result) if result.accepted => format!("{txid}: recorded in the wallet"),
Ok(_) => format!("{txid}: the wallet did not accept its own transaction"),
Err(e) => format!(
"{txid}: not recorded in the wallet ({e}); `bsv-wallet receive {txid}` records it once it is mined"
),
},
Err(e) => format!("{txid}: not recorded in the wallet (stored arguments: {e})"),
};
if !line.ends_with("recorded in the wallet") {
tracing::warn!(marker = "own_transaction_not_recorded", "{line}");
}
report.stored.push(line);
}
}
fn readme_spend_guard<H: Headers>(
state: &mut State,
params: &Params,
headers: &H,
) -> std::result::Result<bool, CheckError> {
state.step(params, headers, Input::Host(HostAction::SpendAttempt))?;
state.step(params, headers, Input::Evidence(Evidence::Recheck))?;
Ok(matches!(state.word(), Word::Mined(_)) && !state.suspect())
}
pub async fn spend_guard<V: WalletServices + ?Sized, C: Clock>(
storage: &StorageSqlx,
services: &V,
clock: &C,
opts: &TickOptions,
txid: &str,
) -> Result<std::result::Result<bool, CheckError>> {
ensure_tables(storage).await?;
let Some(row) = load_row(storage, txid).await? else {
return Ok(Ok(false));
};
let params = opts.params();
let snapshot = HeaderSnapshot::take(services, row.proof_height()).await;
let (mut state, fault) = row.replay(¶ms, &snapshot);
match fault {
Some(fault) if !matches!(fault, CheckError::RootMismatch(_)) => return Ok(Err(fault)),
_ => {}
}
let verdict = readme_spend_guard(&mut state, ¶ms, &snapshot);
store_row(storage, &row, &state, clock.now()).await?;
Ok(verdict)
}
#[cfg(test)]
pub(crate) mod tests {
use super::*;
use bsv_wallet_toolbox::services::mock::{MockResponse, MockWalletServices};
use bsv_wallet_toolbox::services::{BlockHeader, GetMerklePathResult};
use bsv_wallet_toolbox::{
WalletStorageWriter, BROADCAST_STATUS_ACCEPTED, BROADCAST_STATUS_MINED,
};
use chrono::Utc;
pub(crate) struct FixedClock(pub Timestamp);
impl Clock for FixedClock {
fn now(&self) -> Timestamp {
self.0
}
}
pub(crate) fn opts() -> TickOptions {
TickOptions {
age_threshold: 600,
max_asks: 20,
}
}
pub(crate) async fn wallet_storage() -> (StorageSqlx, i64) {
let storage = StorageSqlx::in_memory().await.unwrap();
storage
.migrate("tracker-tests", &("02".to_string() + &"ab".repeat(32)))
.await
.unwrap();
storage.make_available().await.unwrap();
let (user, _) = storage
.find_or_insert_user(&("02".to_string() + &"cd".repeat(32)))
.await
.unwrap();
(storage, user.user_id)
}
pub(crate) async fn insert_tx(storage: &StorageSqlx, user_id: i64, txid: &str, status: &str) {
sqlx::query(
"INSERT INTO transactions (user_id, status, reference, is_outgoing, satoshis, version, lock_time, description, txid, raw_tx, created_at, updated_at) \
VALUES (?, ?, ?, 1, 0, 1, 0, 'd', ?, X'01000000', ?, ?)",
)
.bind(user_id)
.bind(status)
.bind(&txid[..8])
.bind(txid)
.bind(Utc::now())
.bind(Utc::now())
.execute(storage.pool())
.await
.unwrap();
}
pub(crate) fn header(height: u32, root: &str, nonce: u32) -> BlockHeader {
BlockHeader {
version: 1,
previous_hash: "00".repeat(32),
merkle_root: root.to_string(),
time: 1_700_000_000,
bits: 0x1d00_ffff,
nonce,
hash: String::new(),
height,
}
}
pub(crate) fn path_of(txid: &str, height: u32) -> String {
MerklePath::from_coinbase_txid(txid, height).to_hex()
}
pub(crate) fn proof_answer(path: Option<String>) -> MockResponse<GetMerklePathResult> {
MockResponse::Success(GetMerklePathResult {
name: Some("Services".to_string()),
merkle_path: path,
header: None,
error: None,
notes: vec![],
})
}
#[tokio::test]
async fn a_pass_tracks_announces_and_proves_one_named_transaction() {
let (storage, user) = wallet_storage().await;
let txid = "a1".repeat(32);
insert_tx(&storage, user, &txid, "unproven").await;
storage
.record_broadcast_status(&txid, "arcade", BROADCAST_STATUS_ACCEPTED)
.await
.unwrap();
let services = MockWalletServices::builder()
.get_merkle_path_response(proof_answer(Some(path_of(&txid, 900))))
.build();
services.set_header_for_height(header(900, &txid, 0));
let now = Utc::now().timestamp() as u64;
let report = tick(&storage, &services, &FixedClock(now), &opts())
.await
.unwrap();
assert_eq!((report.adopted, report.hints, report.asked), (1, 1, 0));
assert_eq!(
stored_word(&storage, &txid).await.unwrap(),
Some(("announced".to_string(), None, None))
);
assert_eq!(services.call_count("get_merkle_path"), 0);
let report = tick(&storage, &services, &FixedClock(now + 601), &opts())
.await
.unwrap();
assert_eq!(report.asked, 1);
assert_eq!(report.mined, vec![txid.clone()]);
assert_eq!(
stored_word(&storage, &txid).await.unwrap(),
Some(("mined".to_string(), Some(900), None))
);
assert_eq!(services.call_count("get_merkle_path"), 1);
let report = tick(&storage, &services, &FixedClock(now + 1300), &opts())
.await
.unwrap();
assert_eq!((report.tracked, report.asked), (0, 0));
assert_eq!(services.call_count("get_merkle_path"), 1);
}
#[tokio::test]
async fn a_broadcasters_mined_is_a_hint_that_asks_for_the_proof() {
let (storage, user) = wallet_storage().await;
let txid = "b2".repeat(32);
insert_tx(&storage, user, &txid, "unproven").await;
storage
.record_broadcast_status(&txid, "arcade", BROADCAST_STATUS_MINED)
.await
.unwrap();
let services = MockWalletServices::builder()
.get_merkle_path_response(proof_answer(None))
.build();
let now = Utc::now().timestamp() as u64;
let report = tick(&storage, &services, &FixedClock(now), &opts())
.await
.unwrap();
assert_eq!(report.asked, 1);
assert_eq!(report.no_proof, vec![txid.clone()]);
assert!(report.mined.is_empty());
let (word, height, reask) = stored_word(&storage, &txid).await.unwrap().unwrap();
assert_eq!((word.as_str(), height), ("built", None));
assert_eq!(reask.as_deref(), Some("hint:arcade|mined"));
let report = tick(&storage, &services, &FixedClock(now + 60), &opts())
.await
.unwrap();
assert_eq!(report.asked, 0);
let report = tick(&storage, &services, &FixedClock(now + 600), &opts())
.await
.unwrap();
assert_eq!(report.asked, 1);
let report = tick(&storage, &services, &FixedClock(now + 1200), &opts())
.await
.unwrap();
assert_eq!(report.asked, 0);
assert_eq!(services.call_count("get_merkle_path"), 2);
}
#[tokio::test]
async fn a_moved_header_reasks_the_rows_at_or_above_it_and_no_other() {
let (storage, user) = wallet_storage().await;
let low = "c3".repeat(32);
let high = "d4".repeat(32);
let services = MockWalletServices::builder()
.get_merkle_path_response(MockResponse::Sequence(vec![
proof_answer(Some(path_of(&low, 100))),
proof_answer(Some(path_of(&high, 200))),
proof_answer(Some(path_of(&high, 201))),
]))
.build();
services.set_header_for_height(header(100, &low, 0));
services.set_header_for_height(header(200, &high, 0));
services.set_tip_header(header_with_hash(200, &high, 0));
let now = Utc::now().timestamp() as u64;
for txid in [&low, &high] {
insert_tx(&storage, user, txid, "unproven").await;
storage
.record_broadcast_status(txid, "arcade", BROADCAST_STATUS_MINED)
.await
.unwrap();
}
let report = tick(&storage, &services, &FixedClock(now), &opts())
.await
.unwrap();
assert_eq!(report.mined.len(), 2, "{report:?}");
services.set_header_for_height(header(200, &"ee".repeat(32), 7));
services.set_header_for_height(header(201, &high, 0));
services.set_tip_header(header_with_hash(201, &high, 0));
let report = tick(&storage, &services, &FixedClock(now + 60), &opts())
.await
.unwrap();
assert_eq!(report.moved_from, Some(200));
assert_eq!(
report.tracked, 1,
"only the row at or above the moved header"
);
assert_eq!(report.moved, vec![high.clone()]);
assert_eq!(report.mined, vec![high.clone()]);
assert_eq!(
stored_word(&storage, &high).await.unwrap().unwrap().1,
Some(201)
);
assert_eq!(
stored_word(&storage, &low).await.unwrap().unwrap().1,
Some(100)
);
}
#[tokio::test]
async fn a_not_final_transaction_is_announced_again_at_its_lock_time() {
use bsv_wallet_toolbox::services::mock::{
error_post_beef_result, success_post_beef_result,
};
let (storage, _user) = wallet_storage().await;
let txid = "c7".repeat(32);
let lock = 2_000_000_000u64;
let services = MockWalletServices::builder()
.post_beef_response(MockResponse::Sequence(vec![
MockResponse::Success(vec![error_post_beef_result("Arcade", "476 non-final")]),
MockResponse::Success(vec![success_post_beef_result("Arcade", &[&txid])]),
]))
.get_merkle_path_response(proof_answer(None))
.build();
let hint = Hint::new(
"Arcade|rejected",
HintStatus::Rejected {
reason: "476 non-final".to_string(),
},
lock - 3600,
);
let again = Reannounce {
beef: vec![1, 2, 3],
at: lock,
on_accept: None,
};
let word = heard(
&storage,
&FixedClock(lock - 3600),
&opts(),
&txid,
hint,
Some(lock),
Some(again),
)
.await
.unwrap();
assert_eq!(
word,
(
"built".to_string(),
Some("hint:Arcade|rejected".to_string())
)
);
let report = tick(&storage, &services, &FixedClock(lock - 1), &opts())
.await
.unwrap();
assert_eq!((report.tracked, report.asked), (1, 0), "{report:?}");
assert_eq!(services.call_count("post_beef"), 0);
let report = tick(&storage, &services, &FixedClock(lock), &opts())
.await
.unwrap();
assert_eq!(report.not_final, vec![txid.clone()]);
assert_eq!(services.call_count("post_beef"), 1);
let report = tick(&storage, &services, &FixedClock(lock + 60), &opts())
.await
.unwrap();
assert!(report.not_final.is_empty() && report.reannounced.is_empty());
assert_eq!(services.call_count("post_beef"), 1, "one pause between two");
let report = tick(&storage, &services, &FixedClock(lock + 600), &opts())
.await
.unwrap();
assert_eq!(report.reannounced, vec![txid.clone()]);
assert_eq!(services.call_count("post_beef"), 2);
assert_eq!(
stored_word(&storage, &txid).await.unwrap().unwrap().0,
"announced"
);
assert_eq!(services.call_count("get_merkle_path"), 0);
let report = tick(&storage, &services, &FixedClock(lock + 1300), &opts())
.await
.unwrap();
assert_eq!(report.asked, 1);
assert_eq!(services.call_count("post_beef"), 2);
assert_eq!(services.call_count("get_merkle_path"), 1);
}
#[test]
fn not_final_is_read_from_the_refusal() {
for said in [
"476 non-final",
"bad-txns-nonfinal",
"tx is not final",
"code 476",
] {
assert!(says_not_final(said), "{said}");
}
for said in [
"fee too low",
"DOUBLE_SPEND_ATTEMPTED",
"txid 4476aa",
"HTTP 500",
] {
assert!(!says_not_final(said), "{said}");
}
}
fn header_with_hash(height: u32, root: &str, nonce: u32) -> BlockHeader {
let mut h = header(height, root, nonce);
h.hash = header_of(&h.to_binary()).unwrap().hash;
h
}
#[tokio::test]
async fn the_readme_spend_guard_runs_as_the_host() {
let (storage, user) = wallet_storage().await;
let txid = "e5".repeat(32);
insert_tx(&storage, user, &txid, "unproven").await;
storage
.record_broadcast_status(&txid, "arcade", BROADCAST_STATUS_MINED)
.await
.unwrap();
let services = MockWalletServices::builder()
.get_merkle_path_response(proof_answer(Some(path_of(&txid, 900))))
.build();
services.set_header_for_height(header(900, &txid, 0));
let clock = FixedClock(Utc::now().timestamp() as u64);
tick(&storage, &services, &clock, &opts()).await.unwrap();
let guard = spend_guard(&storage, &services, &clock, &opts(), &txid)
.await
.unwrap();
println!(
"spend guard, header in place: {guard:?}, stored {:?}",
stored_word(&storage, &txid).await.unwrap().unwrap()
);
assert_eq!(guard, Ok(true));
services.set_tip_unavailable();
let guard = spend_guard(&storage, &services, &clock, &opts(), &txid)
.await
.unwrap();
println!(
"spend guard, no header service: {guard:?}, stored {:?}",
stored_word(&storage, &txid).await.unwrap().unwrap()
);
assert!(matches!(guard, Err(CheckError::Headers(_))), "{guard:?}");
assert_eq!(
stored_word(&storage, &txid).await.unwrap().unwrap().0,
"mined",
"could not look changes no word"
);
services.set_tip_header(header_with_hash(901, &"00".repeat(32), 1));
services.set_header_for_height(header(900, &"ee".repeat(32), 7));
let guard = spend_guard(&storage, &services, &clock, &opts(), &txid)
.await
.unwrap();
println!(
"spend guard, header moved: {guard:?}, stored {:?}",
stored_word(&storage, &txid).await.unwrap().unwrap()
);
assert_eq!(guard, Ok(false));
let guard = spend_guard(&storage, &services, &clock, &opts(), &"f6".repeat(32))
.await
.unwrap();
assert_eq!(guard, Ok(false));
let params = opts().params();
services.set_tip_header(header_with_hash(900, &txid, 0));
services.set_header_for_height(header(900, &txid, 0));
let held = HeaderSnapshot::take(&services, [900]).await;
let row = load_row(&storage, &txid).await.unwrap().unwrap();
let (mut state, fault) = row.replay(¶ms, &held);
assert_eq!(fault, None);
let lost = HeaderSnapshot::take(&services, []).await;
let guard = readme_spend_guard(&mut state, ¶ms, &lost);
println!("spend guard, height unanswered: {guard:?}");
assert_eq!(guard, Err(CheckError::Unavailable(900)));
}
}