Skip to main content

bsv_wallet_cli/
tracker_host.rs

1//! The CLI as a host of the transaction tracker (`bsv-tracker`).
2//!
3//! The tracker holds one word per transaction, the evidence behind it, the
4//! hints it has heard and the next re-ask. Only headers and proofs change a
5//! word; every broadcaster word is a hint that schedules a re-ask. The crate
6//! fetches nothing, reads no clock and stores nothing: the host supplies the
7//! four things below, runs the re-asks, and keeps the state.
8//!
9//! # The four host things
10//!
11//! | the tracker's trait | here | what it is |
12//! |---|---|---|
13//! | `Headers` | [`HeaderSnapshot`] | the header service's tip and the headers at the heights a pass needs, read once per pass through the wallet's services; a height not read is "no answer", a fault is a fault |
14//! | `ProofFetcher` | [`Courier`] | `get_merkle_path` for one named transaction; what it carries is a raw path, checked by the tracker against the snapshot |
15//! | `HintSource` | [`MemoryHints`] | the broadcast memory rows of the wallet's storage (what a broadcaster answered, what Arcade pushed over SSE or the webhook), each word fed once |
16//! | `Clock` | [`SystemClock`] | seconds since the epoch |
17//!
18//! The tracker's traits are synchronous and the wallet's transports are
19//! not, so a pass reads first (headers, memory rows, proofs for the re-asks
20//! that are due) and steps second.
21//!
22//! # The state, in the CLI's storage
23//!
24//! One row per tracked transaction in `tracker_states`, in the wallet's own
25//! database. A row holds inputs, never a trusted word: whether the host
26//! built it, the hints in the order heard, and the raw merkle path. Loading
27//! a row replays those inputs, so a stored path is checked again against
28//! the active header every time it is read (the tracker has no
29//! deserializer for a checked proof, by design). The `word`, `height` and
30//! `reask` columns are what the tracker said at the last step, kept for
31//! reading and for the indexes; nothing is decided from them.
32//!
33//! # What a pass asks, and what it never asks
34//!
35//! A pass ([`tick`]) names transactions; it never walks the chain and
36//! never asks a chain index. It loads the rows that are not yet mined, feeds
37//! the new hints, ticks the clock, and runs the re-asks that are due: one
38//! proof request for one transaction, with a doubling pause between two
39//! requests for the same transaction. A mined row is read again only when
40//! the header at or below its height moved (the last twelve tips are kept
41//! to see that) or when a spend names it ([`spend_guard`]).
42//!
43//! One re-ask is an announce, not a proof request: a transaction a
44//! broadcaster called not final (a forced `gift-claim` before its lock
45//! time) is kept with its BEEF and handed to the broadcasters again at the
46//! lock time ([`heard`], and the re-announce step of [`tick`]).
47
48use std::collections::{BTreeMap, HashMap, VecDeque};
49use std::convert::Infallible;
50
51use anyhow::Result;
52use bsv_sdk::transaction::MerklePath;
53use bsv_sdk::wallet::{InternalizeActionArgs, WalletInterface};
54use bsv_tracker::{
55    CheckError, Clock, Evidence, Header, HeaderError, Headers, Height, Hint, HintSource,
56    HintStatus, HostAction, Input, Params, Proof, ProofFetcher, Reask, State, Timestamp, TxId,
57    Word,
58};
59use bsv_wallet_toolbox::services::PostBeefResult;
60use bsv_wallet_toolbox::{BroadcastStatus, StorageSqlx, WalletServices};
61use serde::{Deserialize, Serialize};
62use sqlx::Row as _;
63
64/// Default seconds without a new hint before an announced or seen
65/// transaction is asked for its proof (`TRACKER_AGE_SECS`): one block
66/// interval.
67pub const DEFAULT_AGE_SECS: u64 = 600;
68/// Default re-asks one pass runs (`TRACKER_MAX_ASKS`).
69pub const DEFAULT_MAX_ASKS: usize = 20;
70/// The pause between two re-asks of one transaction doubles up to this
71/// many times (600 s, then 1200 s, up to 64 times the age threshold).
72const MAX_BACKOFF_DOUBLINGS: u32 = 6;
73/// A transaction asked this many times in a row with no proof is not asked
74/// again until a new hint names it (about five days at the default age).
75const MAX_ASKS_PER_TX: u32 = 16;
76/// Recent tips kept to see a moved header.
77const TIP_RING: usize = 12;
78
79// ---------------------------------------------------------------------------
80// Headers
81// ---------------------------------------------------------------------------
82
83/// One pass's view of the active chain: the tip and the headers at the
84/// heights the pass named, read from the header service through the
85/// wallet's services. Fails closed: a height that was not read has no
86/// answer, a read that failed is a fault, and with no tip nothing checks.
87#[derive(Debug, Clone)]
88pub struct HeaderSnapshot {
89    tip: std::result::Result<(Height, String), String>,
90    headers: BTreeMap<Height, std::result::Result<Header, String>>,
91}
92
93impl HeaderSnapshot {
94    /// Read the tip and the headers at `heights`.
95    pub async fn take<V: WalletServices + ?Sized>(
96        services: &V,
97        heights: impl IntoIterator<Item = Height>,
98    ) -> Self {
99        let tip = services
100            .get_chain_tip_header()
101            .await
102            .map(|h| (h.height, h.hash.to_ascii_lowercase()))
103            .map_err(|e| e.to_string());
104        let mut snapshot = Self {
105            tip,
106            headers: BTreeMap::new(),
107        };
108        snapshot.extend(services, heights).await;
109        snapshot
110    }
111
112    /// A snapshot of nothing: no tip, no header. Every check fails closed.
113    pub fn unread() -> Self {
114        Self {
115            tip: Err("the header service was not read".to_string()),
116            headers: BTreeMap::new(),
117        }
118    }
119
120    /// Read the headers at `heights` that this snapshot does not hold yet.
121    pub async fn extend<V: WalletServices + ?Sized>(
122        &mut self,
123        services: &V,
124        heights: impl IntoIterator<Item = Height>,
125    ) {
126        for height in heights {
127            if self.headers.contains_key(&height) {
128                continue;
129            }
130            let answer = match services.get_header_for_height(height).await {
131                Ok(bytes) => header_of(&bytes),
132                Err(e) => Err(e.to_string()),
133            };
134            self.headers.insert(height, answer);
135        }
136    }
137
138    /// The tip's height and hash, when the header service gave one.
139    pub fn tip(&self) -> Option<(Height, &str)> {
140        self.tip.as_ref().ok().map(|(h, hash)| (*h, hash.as_str()))
141    }
142}
143
144impl Headers for HeaderSnapshot {
145    fn header_at(&self, height: Height) -> std::result::Result<Option<Header>, HeaderError> {
146        match self.headers.get(&height) {
147            Some(Ok(header)) => Ok(Some(header.clone())),
148            Some(Err(fault)) => Err(HeaderError(fault.clone())),
149            None => Ok(None),
150        }
151    }
152
153    fn tip_height(&self) -> std::result::Result<Height, HeaderError> {
154        self.tip
155            .as_ref()
156            .map(|(height, _)| *height)
157            .map_err(|fault| HeaderError(fault.clone()))
158    }
159}
160
161/// The tracker's header projection of an 80-byte block header: its hash
162/// and its merkle root, both as displayed.
163fn header_of(bytes: &[u8]) -> std::result::Result<Header, String> {
164    if bytes.len() != 80 {
165        return Err(format!("a header of {} bytes", bytes.len()));
166    }
167    let display = |raw: &[u8]| {
168        let mut reversed = raw.to_vec();
169        reversed.reverse();
170        hex::encode(reversed)
171    };
172    Ok(Header {
173        hash: display(&bsv_sdk::primitives::hash::sha256d(bytes)),
174        merkle_root: display(&bytes[36..68]),
175    })
176}
177
178// ---------------------------------------------------------------------------
179// Clock
180// ---------------------------------------------------------------------------
181
182/// The host's clock: seconds since the epoch.
183#[derive(Debug, Clone, Copy, Default)]
184pub struct SystemClock;
185
186impl Clock for SystemClock {
187    fn now(&self) -> Timestamp {
188        std::time::SystemTime::now()
189            .duration_since(std::time::UNIX_EPOCH)
190            .map(|d| d.as_secs())
191            .unwrap_or(0)
192    }
193}
194
195// ---------------------------------------------------------------------------
196// Hints
197// ---------------------------------------------------------------------------
198
199/// The hints of one pass: the broadcast memory rows of the tracked
200/// transactions whose `(provider, word)` the row has not heard yet. A word
201/// a broadcaster repeats every minute is one hint, so a refreshed `seen`
202/// never resets the age clock.
203#[derive(Debug, Default)]
204pub struct MemoryHints {
205    queue: VecDeque<(TxId, Hint)>,
206}
207
208impl MemoryHints {
209    async fn read(storage: &StorageSqlx, rows: &[TrackedRow]) -> Result<Self> {
210        let txids: Vec<String> = rows.iter().map(|r| r.txid.clone()).collect();
211        if txids.is_empty() {
212            return Ok(Self::default());
213        }
214        let mut records = storage.broadcast_records(None, &txids).await?;
215        records.sort_by_key(|r| r.seen_at);
216        let mut queue = VecDeque::new();
217        for record in records {
218            let Some(status) = record.ladder_status() else {
219                continue;
220            };
221            let txid = record.txid.to_ascii_lowercase();
222            let Some(row) = rows.iter().find(|r| r.txid == txid) else {
223                continue;
224            };
225            // A broadcaster's "mined" carries no height here and is no
226            // proof: it reads as a word that disagrees and asks for one.
227            let (label, status) = match status {
228                BroadcastStatus::Accepted => ("accepted", HintStatus::Accepted),
229                BroadcastStatus::Seen => ("seen", HintStatus::Seen),
230                BroadcastStatus::Mined => ("mined", HintStatus::Unknown),
231                BroadcastStatus::Rejected => (
232                    "rejected",
233                    HintStatus::Rejected {
234                        reason: record.status.clone(),
235                    },
236                ),
237                BroadcastStatus::Unknown => ("unknown", HintStatus::Unknown),
238            };
239            let source = format!("{}|{}", record.provider, label);
240            if row.hints.iter().any(|h| h.source == source) {
241                continue;
242            }
243            let observed = record.seen_at.timestamp().max(0) as Timestamp;
244            queue.push_back((txid, Hint::new(source, status, observed)));
245        }
246        Ok(Self { queue })
247    }
248}
249
250impl HintSource for MemoryHints {
251    type Error = Infallible;
252
253    fn next_hint(&mut self) -> std::result::Result<Option<(TxId, Hint)>, Self::Error> {
254        Ok(self.queue.pop_front())
255    }
256}
257
258// ---------------------------------------------------------------------------
259// Proofs
260// ---------------------------------------------------------------------------
261
262/// The host's proof transport: `get_merkle_path` for the transactions whose
263/// re-ask is due, read before the tracker steps. What it hands over is a
264/// raw path bound to its txid; the tracker checks it.
265#[derive(Debug, Default)]
266pub struct Courier {
267    carried: HashMap<TxId, std::result::Result<Option<Proof>, String>>,
268}
269
270impl Courier {
271    /// Ask for the proof of each of `txids`, once.
272    pub async fn carry<V: WalletServices + ?Sized>(services: &V, txids: &[TxId]) -> Self {
273        let mut carried = HashMap::new();
274        for txid in txids {
275            let answer = match services.get_merkle_path(txid, false).await {
276                Ok(found) => match found.merkle_path {
277                    Some(hex) => MerklePath::from_hex(&hex)
278                        .map_err(|e| format!("not a merkle path: {e}"))
279                        .and_then(|path| Proof::new(txid.clone(), path).map_err(|e| e.to_string()))
280                        .map(Some),
281                    None => Ok(None),
282                },
283                Err(e) => Err(e.to_string()),
284            };
285            carried.insert(txid.clone(), answer);
286        }
287        Self { carried }
288    }
289
290    fn heights(&self) -> Vec<Height> {
291        self.carried
292            .values()
293            .filter_map(|answer| answer.as_ref().ok()?.as_ref().map(Proof::height))
294            .collect()
295    }
296}
297
298impl ProofFetcher for Courier {
299    type Error = String;
300
301    fn fetch(
302        &mut self,
303        txid: &str,
304        _reask: &Reask,
305    ) -> std::result::Result<Option<Proof>, Self::Error> {
306        self.carried.remove(txid).unwrap_or(Ok(None))
307    }
308}
309
310// ---------------------------------------------------------------------------
311// The stored row
312// ---------------------------------------------------------------------------
313
314#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
315struct StoredHint {
316    source: String,
317    status: String,
318    #[serde(default, skip_serializing_if = "Option::is_none")]
319    height: Option<Height>,
320    #[serde(default, skip_serializing_if = "Option::is_none")]
321    reason: Option<String>,
322    #[serde(default, skip_serializing_if = "Vec::is_empty")]
323    competitors: Vec<TxId>,
324    observed: Timestamp,
325}
326
327impl StoredHint {
328    fn of(hint: &Hint) -> Self {
329        let (status, height, reason, competitors) = match &hint.status {
330            HintStatus::Accepted => ("accepted", None, None, vec![]),
331            HintStatus::Seen => ("seen", None, None, vec![]),
332            HintStatus::Mined { height } => ("mined", Some(*height), None, vec![]),
333            HintStatus::StaleBlock => ("stale_block", None, None, vec![]),
334            HintStatus::Rejected { reason } => ("rejected", None, Some(reason.clone()), vec![]),
335            HintStatus::DoubleSpend { competitors } => {
336                ("double_spend", None, None, competitors.clone())
337            }
338            HintStatus::OrphanMempool => ("orphan_mempool", None, None, vec![]),
339            HintStatus::Unknown => ("unknown", None, None, vec![]),
340        };
341        Self {
342            source: hint.source.clone(),
343            status: status.to_string(),
344            height,
345            reason,
346            competitors,
347            observed: hint.observed,
348        }
349    }
350
351    fn hint(&self) -> Hint {
352        let status = match self.status.as_str() {
353            "accepted" => HintStatus::Accepted,
354            "seen" => HintStatus::Seen,
355            "mined" => self
356                .height
357                .map_or(HintStatus::Unknown, |height| HintStatus::Mined { height }),
358            "stale_block" => HintStatus::StaleBlock,
359            "rejected" => HintStatus::Rejected {
360                reason: self.reason.clone().unwrap_or_default(),
361            },
362            "double_spend" => HintStatus::DoubleSpend {
363                competitors: self.competitors.clone(),
364            },
365            "orphan_mempool" => HintStatus::OrphanMempool,
366            _ => HintStatus::Unknown,
367        };
368        Hint::new(self.source.clone(), status, self.observed)
369    }
370}
371
372/// One tracked transaction as the CLI's storage holds it: inputs to replay.
373#[derive(Debug, Clone, PartialEq, Eq)]
374struct TrackedRow {
375    txid: TxId,
376    built: bool,
377    /// Oldest first.
378    hints: Vec<StoredHint>,
379    /// The raw merkle path, hex. Checked again on every load.
380    bump: Option<String>,
381    asks: u32,
382    next_ask_at: Option<Timestamp>,
383    /// The BEEF to hand to the broadcasters again, and when.
384    reannounce_beef: Option<Vec<u8>>,
385    reannounce_at: Option<Timestamp>,
386    /// What the wallet records once a broadcaster takes it: the
387    /// `internalizeAction` arguments, JSON.
388    on_accept: Option<String>,
389}
390
391impl TrackedRow {
392    fn new(txid: &str) -> Self {
393        Self {
394            txid: txid.to_ascii_lowercase(),
395            built: true,
396            hints: vec![],
397            bump: None,
398            asks: 0,
399            next_ask_at: None,
400            reannounce_beef: None,
401            reannounce_at: None,
402            on_accept: None,
403        }
404    }
405
406    fn proof(&self) -> Option<std::result::Result<Proof, CheckError>> {
407        let hex = self.bump.as_deref()?;
408        Some(
409            MerklePath::from_hex(hex)
410                .map_err(|e| CheckError::InvalidProof(e.to_string()))
411                .and_then(|path| Proof::new(self.txid.clone(), path)),
412        )
413    }
414
415    fn proof_height(&self) -> Option<Height> {
416        self.proof()?.ok().map(|p| p.height())
417    }
418
419    /// Replay the stored inputs. The stored path is checked against
420    /// `headers` like any other proof; a path that no longer checks leaves
421    /// the hint-tier word and a re-ask, and is returned as the fault.
422    fn replay<H: Headers>(&self, params: &Params, headers: &H) -> (State, Option<CheckError>) {
423        let mut state = State::new(self.txid.clone());
424        if self.built {
425            let _ = state.step(params, headers, Input::Host(HostAction::Build));
426        }
427        for hint in &self.hints {
428            let _ = state.step(params, headers, Input::Hint(hint.hint()));
429        }
430        let fault = match self.proof() {
431            Some(Ok(proof)) => state
432                .step(params, headers, Input::Evidence(Evidence::Proof(proof)))
433                .err(),
434            Some(Err(fault)) => Some(fault),
435            None => None,
436        };
437        (state, fault)
438    }
439}
440
441/// The tracker's word, as one lowercase label.
442pub fn word_label(word: &Word) -> &'static str {
443    match word {
444        Word::Unknown => "unknown",
445        Word::Built => "built",
446        Word::Announced => "announced",
447        Word::Seen => "seen",
448        Word::Mined(_) => "mined",
449        Word::Stale => "stale",
450        Word::Rejected { .. } => "rejected",
451        Word::Conflicted { .. } => "conflicted",
452        Word::Abandoned => "abandoned",
453    }
454}
455
456fn reask_label(reask: &Reask) -> String {
457    match reask {
458        Reask::Reorg { fork_height } => format!("reorg@{fork_height}"),
459        Reask::Fork { height } => format!("fork@{height}"),
460        Reask::Hint(hint) => format!("hint:{}", hint.source),
461        Reask::ProofFailed { height } => format!("proof_failed@{height}"),
462        Reask::Recheck => "recheck".to_string(),
463        Reask::Age => "age".to_string(),
464        Reask::Spend => "spend".to_string(),
465    }
466}
467
468async fn ensure_tables(storage: &StorageSqlx) -> Result<()> {
469    for sql in [
470        "CREATE TABLE IF NOT EXISTS tracker_states (\
471            txid TEXT PRIMARY KEY, \
472            built INTEGER NOT NULL DEFAULT 1, \
473            hints TEXT NOT NULL DEFAULT '[]', \
474            bump TEXT, \
475            word TEXT NOT NULL DEFAULT 'unknown', \
476            height INTEGER, \
477            header_hash TEXT, \
478            reask TEXT, \
479            asks INTEGER NOT NULL DEFAULT 0, \
480            next_ask_at INTEGER, \
481            reannounce_beef BLOB, \
482            reannounce_at INTEGER, \
483            on_accept TEXT, \
484            updated_at INTEGER NOT NULL DEFAULT 0)",
485        "CREATE INDEX IF NOT EXISTS tracker_states_word ON tracker_states (word)",
486        "CREATE INDEX IF NOT EXISTS tracker_states_height ON tracker_states (height)",
487        "CREATE TABLE IF NOT EXISTS tracker_meta (key TEXT PRIMARY KEY, value TEXT NOT NULL)",
488    ] {
489        sqlx::query(sql).execute(storage.pool()).await?;
490    }
491    Ok(())
492}
493
494const ROW_COLUMNS: &str =
495    "txid, built, hints, bump, asks, next_ask_at, reannounce_beef, reannounce_at, on_accept";
496
497fn row_of(row: &sqlx::sqlite::SqliteRow) -> TrackedRow {
498    let hints: String = row.get("hints");
499    TrackedRow {
500        txid: row.get("txid"),
501        built: row.get::<i64, _>("built") != 0,
502        hints: serde_json::from_str(&hints).unwrap_or_default(),
503        bump: row.get("bump"),
504        asks: row.get::<i64, _>("asks").max(0) as u32,
505        next_ask_at: row
506            .get::<Option<i64>, _>("next_ask_at")
507            .map(|t| t.max(0) as Timestamp),
508        reannounce_beef: row.get("reannounce_beef"),
509        reannounce_at: row
510            .get::<Option<i64>, _>("reannounce_at")
511            .map(|t| t.max(0) as Timestamp),
512        on_accept: row.get("on_accept"),
513    }
514}
515
516async fn load_row(storage: &StorageSqlx, txid: &str) -> Result<Option<TrackedRow>> {
517    let row = sqlx::query(&format!(
518        "SELECT {ROW_COLUMNS} FROM tracker_states WHERE txid = ?"
519    ))
520    .bind(txid.to_ascii_lowercase())
521    .fetch_optional(storage.pool())
522    .await?;
523    Ok(row.as_ref().map(row_of))
524}
525
526async fn store_row(
527    storage: &StorageSqlx,
528    row: &TrackedRow,
529    state: &State,
530    now: Timestamp,
531) -> Result<()> {
532    let (height, header_hash) = match state.word() {
533        Word::Mined(mined) => (
534            Some(mined.height() as i64),
535            Some(mined.checked().header().hash.clone()),
536        ),
537        _ => (None, None),
538    };
539    sqlx::query(
540        "INSERT INTO tracker_states \
541            (txid, built, hints, bump, word, height, header_hash, reask, asks, next_ask_at, \
542             reannounce_beef, reannounce_at, on_accept, updated_at) \
543         VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) \
544         ON CONFLICT(txid) DO UPDATE SET built = excluded.built, hints = excluded.hints, \
545            bump = excluded.bump, word = excluded.word, height = excluded.height, \
546            header_hash = excluded.header_hash, reask = excluded.reask, asks = excluded.asks, \
547            next_ask_at = excluded.next_ask_at, reannounce_beef = excluded.reannounce_beef, \
548            reannounce_at = excluded.reannounce_at, on_accept = excluded.on_accept, \
549            updated_at = excluded.updated_at",
550    )
551    .bind(&row.txid)
552    .bind(row.built as i64)
553    .bind(serde_json::to_string(&row.hints)?)
554    .bind(&row.bump)
555    .bind(word_label(state.word()))
556    .bind(height)
557    .bind(header_hash)
558    .bind(state.reask().map(reask_label))
559    .bind(row.asks as i64)
560    .bind(row.next_ask_at.map(|t| t as i64))
561    .bind(&row.reannounce_beef)
562    .bind(row.reannounce_at.map(|t| t as i64))
563    .bind(&row.on_accept)
564    .bind(now as i64)
565    .execute(storage.pool())
566    .await?;
567    Ok(())
568}
569
570/// What the storage holds for `txid`: the word, the mined height and the
571/// pending re-ask the tracker gave at its last step. A view for a person;
572/// nothing is decided from it.
573pub async fn stored_word(
574    storage: &StorageSqlx,
575    txid: &str,
576) -> Result<Option<(String, Option<Height>, Option<String>)>> {
577    ensure_tables(storage).await?;
578    let row = sqlx::query("SELECT word, height, reask FROM tracker_states WHERE txid = ?")
579        .bind(txid.to_ascii_lowercase())
580        .fetch_optional(storage.pool())
581        .await?;
582    Ok(row.map(|r| {
583        (
584            r.get("word"),
585            r.get::<Option<i64>, _>("height").map(|h| h as Height),
586            r.get("reask"),
587        )
588    }))
589}
590
591// ---------------------------------------------------------------------------
592// What a broadcaster answered to a transaction the host just handed over
593// ---------------------------------------------------------------------------
594
595/// The broadcasters' answer to one post, read as a hint.
596#[derive(Debug, Clone, PartialEq, Eq)]
597pub enum PostAnswer {
598    /// A broadcaster took it.
599    Accepted { provider: String },
600    /// None took it and one called it not final: its lock time has not
601    /// passed for the node behind that broadcaster. Transient.
602    NotFinal { provider: String, said: String },
603    /// None took it, for some other reason (or none answered).
604    Refused { said: String },
605}
606
607/// Does a broadcaster's refusal say "not final"? The node's word is
608/// `non-final`; ARC and Arcade carry it under code 476, which the toolbox
609/// reads as retryable (bsv-wallet-toolbox-rs `ARCADE_CODE_NON_FINAL`).
610pub fn says_not_final(text: &str) -> bool {
611    let lower = text.to_ascii_lowercase();
612    ["non-final", "nonfinal", "non final", "not final"]
613        .iter()
614        .any(|word| lower.contains(word))
615        || lower
616            .split(|c: char| !c.is_ascii_alphanumeric())
617            .any(|token| token == "476")
618}
619
620/// Read the results of one `post_beef`.
621pub fn read_post(results: &[PostBeefResult]) -> PostAnswer {
622    if let Some(accepted) = results.iter().find(|r| r.is_success()) {
623        return PostAnswer::Accepted {
624            provider: accepted.name.clone(),
625        };
626    }
627    let said: Vec<(String, String)> = results
628        .iter()
629        .map(|r| {
630            let detail = r
631                .error
632                .clone()
633                .or_else(|| r.txid_results.iter().find_map(|t| t.data.clone()))
634                .unwrap_or_else(|| r.status.clone());
635            (r.name.clone(), detail)
636        })
637        .collect();
638    if let Some((provider, detail)) = said.iter().find(|(_, detail)| says_not_final(detail)) {
639        return PostAnswer::NotFinal {
640            provider: provider.clone(),
641            said: detail.clone(),
642        };
643    }
644    PostAnswer::Refused {
645        said: if said.is_empty() {
646            "none answered".to_string()
647        } else {
648            said.iter()
649                .map(|(name, detail)| format!("{name}: {detail}"))
650                .collect::<Vec<_>>()
651                .join("; ")
652        },
653    }
654}
655
656/// A BEEF to hand to the broadcasters again, and when.
657#[derive(Debug, Clone)]
658pub struct Reannounce {
659    /// The BEEF, as it was posted.
660    pub beef: Vec<u8>,
661    /// Not before this time (seconds since the epoch).
662    pub at: Timestamp,
663    /// What the wallet records when a broadcaster takes it then.
664    pub on_accept: Option<InternalizeActionArgs>,
665}
666
667/// Track a transaction this wallet built and just handed to a broadcaster,
668/// and feed what the broadcaster said, as a hint. `not_before` holds every
669/// re-ask until that time (a transaction that cannot be mined before its
670/// lock time has no proof to ask for before it); `reannounce` keeps the
671/// BEEF to hand over again then. Returns the tracker's word and its
672/// pending re-ask.
673pub async fn heard<C: Clock>(
674    storage: &StorageSqlx,
675    clock: &C,
676    opts: &TickOptions,
677    txid: &str,
678    hint: Hint,
679    not_before: Option<Timestamp>,
680    reannounce: Option<Reannounce>,
681) -> Result<(String, Option<String>)> {
682    ensure_tables(storage).await?;
683    let mut row = match load_row(storage, txid).await? {
684        // A row with a stored path is past this: its word is read, never
685        // rewritten from a hint.
686        Some(row) if row.bump.is_some() => {
687            let stored = stored_word(storage, txid).await?.unwrap_or_default();
688            return Ok((stored.0, stored.2));
689        }
690        Some(row) => row,
691        None => TrackedRow::new(txid),
692    };
693    let params = opts.params();
694    let headers = HeaderSnapshot::unread();
695    let (mut state, _) = row.replay(&params, &headers);
696    if !row.hints.iter().any(|h| h.source == hint.source) {
697        row.hints.push(StoredHint::of(&hint));
698        let _ = state.step(&params, &headers, Input::Hint(hint));
699    }
700    row.asks = 0;
701    row.next_ask_at = not_before;
702    match reannounce {
703        Some(again) => {
704            row.reannounce_beef = Some(again.beef);
705            row.reannounce_at = Some(again.at);
706            row.on_accept = match again.on_accept {
707                Some(args) => Some(serde_json::to_string(&args)?),
708                None => None,
709            };
710        }
711        None => {
712            row.reannounce_beef = None;
713            row.reannounce_at = None;
714            row.on_accept = None;
715        }
716    }
717    store_row(storage, &row, &state, clock.now()).await?;
718    Ok((
719        word_label(state.word()).to_string(),
720        state.reask().map(reask_label),
721    ))
722}
723
724// ---------------------------------------------------------------------------
725// The pass
726// ---------------------------------------------------------------------------
727
728/// The host's parameters of a pass.
729#[derive(Debug, Clone)]
730pub struct TickOptions {
731    /// Seconds without a new hint before an announced or seen transaction
732    /// is asked for its proof (the tracker's `Params::age_threshold`).
733    pub age_threshold: Timestamp,
734    /// Re-asks one pass runs.
735    pub max_asks: usize,
736}
737
738impl TickOptions {
739    /// `TRACKER_AGE_SECS` and `TRACKER_MAX_ASKS`, with their defaults.
740    pub fn from_env() -> Self {
741        fn env<T: std::str::FromStr>(key: &str, default: T) -> T {
742            std::env::var(key)
743                .ok()
744                .and_then(|v| v.trim().parse().ok())
745                .unwrap_or(default)
746        }
747        Self {
748            age_threshold: env("TRACKER_AGE_SECS", DEFAULT_AGE_SECS).max(1),
749            max_asks: env("TRACKER_MAX_ASKS", DEFAULT_MAX_ASKS),
750        }
751    }
752
753    fn params(&self) -> Params {
754        Params {
755            age_threshold: self.age_threshold,
756        }
757    }
758}
759
760/// What one pass did.
761#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)]
762pub struct TickReport {
763    /// Rows read this pass (never the mined ones, unless a header moved).
764    pub tracked: usize,
765    /// The wallet's own unproven transactions tracked for the first time.
766    pub adopted: usize,
767    /// New hints fed.
768    pub hints: usize,
769    /// Re-asks run: one proof request each.
770    pub asked: usize,
771    /// Transactions whose proof checked this pass: the word is `mined`.
772    pub mined: Vec<String>,
773    /// Re-asks that brought no proof; each is asked again after its pause.
774    pub no_proof: Vec<String>,
775    /// Mined rows whose stored path no longer meets the active header.
776    pub moved: Vec<String>,
777    /// The lowest height whose header moved since the last pass, if any.
778    pub moved_from: Option<Height>,
779    /// Transactions handed to the broadcasters again and accepted.
780    pub reannounced: Vec<String>,
781    /// Transactions handed over again and still called not final; each is
782    /// handed over again after a pause.
783    pub not_final: Vec<String>,
784    /// What the wallet is to record for each transaction accepted on a
785    /// re-announce (`txid`, the `internalizeAction` arguments): the caller
786    /// that holds the wallet runs [`record_accepted`].
787    #[serde(skip)]
788    pub on_accept: Vec<(String, String)>,
789    /// Faults the host reports: `txid: what failed`.
790    pub faults: Vec<String>,
791    /// What the wallet's storage did with each checked proof.
792    pub stored: Vec<String>,
793}
794
795impl TickReport {
796    /// Nothing moved and nothing was asked.
797    pub fn is_quiet(&self) -> bool {
798        self.adopted == 0
799            && self.hints == 0
800            && self.asked == 0
801            && self.reannounced.is_empty()
802            && self.not_final.is_empty()
803            && self.moved.is_empty()
804            && self.faults.is_empty()
805    }
806
807    /// One line.
808    pub fn summary(&self) -> String {
809        format!(
810            "tracker tick: {} tracked ({} new), {} hint(s), {} announced again ({} still not final), {} re-ask(s): {} mined, {} without a proof, {} moved, {} fault(s)",
811            self.tracked,
812            self.adopted,
813            self.hints,
814            self.reannounced.len(),
815            self.not_final.len(),
816            self.asked,
817            self.mined.len(),
818            self.no_proof.len(),
819            self.moved.len(),
820            self.faults.len(),
821        )
822    }
823}
824
825/// The recent tips, newest last, read from and written to `tracker_meta`.
826async fn read_ring(storage: &StorageSqlx) -> Vec<(Height, String)> {
827    sqlx::query("SELECT value FROM tracker_meta WHERE key = 'tips'")
828        .fetch_optional(storage.pool())
829        .await
830        .ok()
831        .flatten()
832        .and_then(|r| serde_json::from_str(&r.get::<String, _>("value")).ok())
833        .unwrap_or_default()
834}
835
836async fn write_ring(storage: &StorageSqlx, ring: &[(Height, String)]) -> Result<()> {
837    sqlx::query(
838        "INSERT INTO tracker_meta (key, value) VALUES ('tips', ?) \
839         ON CONFLICT(key) DO UPDATE SET value = excluded.value",
840    )
841    .bind(serde_json::to_string(ring)?)
842    .execute(storage.pool())
843    .await?;
844    Ok(())
845}
846
847/// The lowest height whose header is no longer the one a past pass saw,
848/// and the ring with the moved entries dropped and the tip appended. One
849/// header read when the tip advanced, none when it did not, and one more
850/// per replaced block. Deeper than the ring, the ring's lowest height is
851/// the answer and the host says so.
852async fn moved_since<V: WalletServices + ?Sized>(
853    services: &V,
854    snapshot: &mut HeaderSnapshot,
855    mut ring: Vec<(Height, String)>,
856) -> (Option<Height>, Vec<(Height, String)>) {
857    let Some((tip_height, tip_hash)) = snapshot.tip().map(|(h, hash)| (h, hash.to_string())) else {
858        return (None, ring);
859    };
860    let mut moved = None;
861    while let Some((height, hash)) = ring.last().cloned() {
862        let active = if height == tip_height {
863            Some(tip_hash.clone())
864        } else if height > tip_height {
865            None
866        } else {
867            snapshot.extend(services, [height]).await;
868            match snapshot.header_at(height) {
869                Ok(Some(header)) => Some(header.hash),
870                // Could not look: nothing is concluded this pass.
871                _ => return (moved, ring),
872            }
873        };
874        if active.as_deref() == Some(hash.as_str()) {
875            break;
876        }
877        moved = Some(height);
878        ring.pop();
879    }
880    if ring.last().map(|(h, _)| *h) != Some(tip_height) {
881        ring.push((tip_height, tip_hash));
882    }
883    if ring.len() > TIP_RING {
884        let extra = ring.len() - TIP_RING;
885        ring.drain(..extra);
886    }
887    (moved, ring)
888}
889
890/// One pass of the host. See the module docs.
891pub async fn tick<V: WalletServices + ?Sized, C: Clock>(
892    storage: &StorageSqlx,
893    services: &V,
894    clock: &C,
895    opts: &TickOptions,
896) -> Result<TickReport> {
897    ensure_tables(storage).await?;
898    let params = opts.params();
899    let now = clock.now();
900    let mut report = TickReport::default();
901
902    // The wallet's own transactions that are out and unproven, tracked
903    // from the first pass that sees them. Never anyone else's.
904    let adopted = sqlx::query(
905        "INSERT OR IGNORE INTO tracker_states (txid, word, updated_at) \
906         SELECT DISTINCT lower(txid), 'built', ? FROM transactions \
907         WHERE txid IS NOT NULL AND status IN ('unproven', 'sending')",
908    )
909    .bind(now as i64)
910    .execute(storage.pool())
911    .await?;
912    report.adopted = adopted.rows_affected() as usize;
913
914    // Did a header we stood on move? Then the mined rows at or above it
915    // are read again, and no other mined row is.
916    let mut snapshot = HeaderSnapshot::take(services, []).await;
917    let (moved_from, ring) = moved_since(services, &mut snapshot, read_ring(storage).await).await;
918    report.moved_from = moved_from;
919
920    let mut rows: Vec<TrackedRow> = sqlx::query(&format!(
921        "SELECT {ROW_COLUMNS} FROM tracker_states \
922         WHERE word IN ('unknown', 'built', 'announced', 'seen', 'stale') \
923            OR (word = 'mined' AND ? IS NOT NULL AND height >= ?) \
924         ORDER BY txid"
925    ))
926    .bind(moved_from.map(|h| h as i64))
927    .bind(moved_from.map(|h| h as i64))
928    .fetch_all(storage.pool())
929    .await?
930    .iter()
931    .map(row_of)
932    .collect();
933    report.tracked = rows.len();
934
935    // Read: the headers the stored paths name, and the new hints.
936    snapshot
937        .extend(services, rows.iter().filter_map(TrackedRow::proof_height))
938        .await;
939    let mut hints = MemoryHints::read(storage, &rows).await?;
940
941    // Step: replay each row, feed its hints, tick its clock.
942    let mut states: HashMap<TxId, State> = HashMap::new();
943    for row in rows.iter_mut() {
944        let (state, fault) = row.replay(&params, &snapshot);
945        if let Some(fault) = fault {
946            // The stored path does not check against the active header:
947            // the row is back at its hint-tier word with a re-ask, due now.
948            if matches!(fault, CheckError::RootMismatch(_)) {
949                // The path proves inclusion in a block that is no longer
950                // active: it is dropped, and a new one is asked for.
951                row.bump = None;
952                report.moved.push(row.txid.clone());
953                tracing::warn!(
954                    marker = "tracker_stored_proof_moved",
955                    txid = %row.txid,
956                    "a stored merkle path no longer meets the active header; asking again"
957                );
958            } else {
959                report.faults.push(format!("{}: {fault}", row.txid));
960            }
961            row.asks = 0;
962            row.next_ask_at = None;
963        }
964        states.insert(row.txid.clone(), state);
965    }
966    while let Ok(Some((txid, hint))) = hints.next_hint() {
967        let (Some(state), Some(row)) = (
968            states.get_mut(&txid),
969            rows.iter_mut().find(|r| r.txid == txid),
970        ) else {
971            continue;
972        };
973        row.hints.push(StoredHint::of(&hint));
974        let _ = state.step(&params, &snapshot, Input::Hint(hint));
975        report.hints += 1;
976        // A word that disagrees is asked about now, not after the pause
977        // an earlier re-ask earned.
978        if matches!(state.reask(), Some(Reask::Hint(_))) {
979            row.asks = 0;
980            row.next_ask_at = None;
981        }
982    }
983    // Announce again what a broadcaster called not final, at its time.
984    // The lock time passing on our clock is not the node's rule passing
985    // (the node reads the median time past, about an hour behind; the rule
986    // is the sibling's, bsv-script-lean `docs/BLOCK-CONSENSUS.md` section
987    // 2.2), so "not final" again is expected and waits one more pause.
988    for row in rows.iter_mut() {
989        let (Some(beef), Some(at)) = (row.reannounce_beef.clone(), row.reannounce_at) else {
990            continue;
991        };
992        let Some(state) = states.get_mut(&row.txid) else {
993            continue;
994        };
995        if at > now || matches!(state.word(), Word::Mined(_)) {
996            continue;
997        }
998        let pause = Some(now.saturating_add(opts.age_threshold));
999        let answer = services
1000            .post_beef(&beef, std::slice::from_ref(&row.txid))
1001            .await
1002            .map(|results| read_post(&results));
1003        match answer {
1004            Ok(PostAnswer::Accepted { provider }) => {
1005                let hint = Hint::new(format!("{provider}|accepted"), HintStatus::Accepted, now);
1006                row.hints.push(StoredHint::of(&hint));
1007                let _ = state.step(&params, &snapshot, Input::Hint(hint));
1008                row.reannounce_beef = None;
1009                row.reannounce_at = None;
1010                row.asks = 0;
1011                row.next_ask_at = pause;
1012                report.reannounced.push(row.txid.clone());
1013                if let Some(args) = row.on_accept.take() {
1014                    report.on_accept.push((row.txid.clone(), args));
1015                }
1016            }
1017            Ok(PostAnswer::NotFinal { .. }) => {
1018                row.asks = row.asks.saturating_add(1);
1019                row.next_ask_at = pause;
1020                if row.asks >= MAX_ASKS_PER_TX {
1021                    row.reannounce_at = None;
1022                    report.faults.push(format!(
1023                        "{}: still not final after {MAX_ASKS_PER_TX} announces; not announced again",
1024                        row.txid
1025                    ));
1026                } else {
1027                    row.reannounce_at = pause;
1028                    report.not_final.push(row.txid.clone());
1029                }
1030            }
1031            Ok(PostAnswer::Refused { said }) => {
1032                row.reannounce_at = None;
1033                row.next_ask_at = pause;
1034                report.faults.push(format!(
1035                    "{}: refused when announced again: {said}",
1036                    row.txid
1037                ));
1038            }
1039            Err(e) => {
1040                // Could not look: the same time again, one pause on.
1041                row.reannounce_at = pause;
1042                report
1043                    .faults
1044                    .push(format!("{}: could not announce again: {e}", row.txid));
1045            }
1046        }
1047    }
1048    for state in states.values_mut() {
1049        let _ = state.step(&params, &snapshot, Input::Host(HostAction::Tick(now)));
1050    }
1051
1052    // Read again: a proof for each re-ask that is due, and its header.
1053    let due: Vec<TxId> = rows
1054        .iter()
1055        .filter(|row| {
1056            states[&row.txid].reask().is_some()
1057                && row.asks < MAX_ASKS_PER_TX
1058                && row.next_ask_at.is_none_or(|at| at <= now)
1059        })
1060        .take(opts.max_asks)
1061        .map(|row| row.txid.clone())
1062        .collect();
1063    let mut courier = Courier::carry(services, &due).await;
1064    snapshot.extend(services, courier.heights()).await;
1065
1066    // Step again: each proof through the tracker's check.
1067    for txid in &due {
1068        let (Some(state), Some(row)) = (
1069            states.get_mut(txid),
1070            rows.iter_mut().find(|r| &r.txid == txid),
1071        ) else {
1072            continue;
1073        };
1074        report.asked += 1;
1075        let reask = state.reask().cloned().unwrap_or(Reask::Age);
1076        let fetched = courier.fetch(txid, &reask);
1077        let proven = match fetched {
1078            Ok(Some(proof)) => {
1079                let hex = proof.path().to_hex();
1080                match state.step(&params, &snapshot, Input::Evidence(Evidence::Proof(proof))) {
1081                    Ok(()) => {
1082                        row.bump = Some(hex);
1083                        true
1084                    }
1085                    Err(fault) => {
1086                        report.faults.push(format!("{txid}: {fault}"));
1087                        false
1088                    }
1089                }
1090            }
1091            Ok(None) => {
1092                report.no_proof.push(txid.clone());
1093                false
1094            }
1095            Err(fault) => {
1096                report
1097                    .faults
1098                    .push(format!("{txid}: could not ask for a proof: {fault}"));
1099                false
1100            }
1101        };
1102        if proven {
1103            row.asks = 0;
1104            row.next_ask_at = None;
1105            report.mined.push(txid.clone());
1106        } else {
1107            let pause = opts
1108                .age_threshold
1109                .saturating_mul(1 << row.asks.min(MAX_BACKOFF_DOUBLINGS));
1110            row.asks = row.asks.saturating_add(1);
1111            row.next_ask_at = Some(now.saturating_add(pause));
1112        }
1113    }
1114
1115    // Keep the state, and hand each checked proof to the wallet's own
1116    // storage through its one proof funnel.
1117    for row in &rows {
1118        let state = &states[&row.txid];
1119        store_row(storage, row, state, now).await?;
1120        if let (true, Word::Mined(mined)) = (report.mined.contains(&row.txid), state.word()) {
1121            let checked = mined.checked();
1122            let outcome = storage
1123                .ingest_merkle_proof(
1124                    &row.txid,
1125                    &checked.proof().path().to_binary(),
1126                    checked.height(),
1127                    &checked.header().hash,
1128                    Some(checked.root()),
1129                )
1130                .await;
1131            report.stored.push(match outcome {
1132                Ok(outcome) => format!("{}: {outcome:?}", row.txid),
1133                Err(e) => format!("{}: not stored: {e}", row.txid),
1134            });
1135        }
1136    }
1137    write_ring(storage, &ring).await?;
1138    Ok(report)
1139}
1140
1141/// Record in the wallet what a pass's re-announces got accepted (E4): the
1142/// wallet's storage is the verdict for its own actions, so a claim the
1143/// wallet made is internalized from the BEEF it already holds, and nothing
1144/// scans the chain for it. Each outcome is added to `report.stored`.
1145pub async fn record_accepted<W: WalletInterface>(wallet: &W, report: &mut TickReport) {
1146    for (txid, args) in std::mem::take(&mut report.on_accept) {
1147        let line = match serde_json::from_str::<InternalizeActionArgs>(&args) {
1148            Ok(args) => match wallet.internalize_action(args, "bsv-wallet-cli").await {
1149                Ok(result) if result.accepted => format!("{txid}: recorded in the wallet"),
1150                Ok(_) => format!("{txid}: the wallet did not accept its own transaction"),
1151                Err(e) => format!(
1152                    "{txid}: not recorded in the wallet ({e}); `bsv-wallet receive {txid}` records it once it is mined"
1153                ),
1154            },
1155            Err(e) => format!("{txid}: not recorded in the wallet (stored arguments: {e})"),
1156        };
1157        if !line.ends_with("recorded in the wallet") {
1158            tracing::warn!(marker = "own_transaction_not_recorded", "{line}");
1159        }
1160        report.stored.push(line);
1161    }
1162}
1163
1164// ---------------------------------------------------------------------------
1165// The spend guard
1166// ---------------------------------------------------------------------------
1167
1168/// The wallet spend guard, as the tracker's README gives it (bsv-tracker
1169/// `README.md`, "Wallet: spend guard"), unchanged: the scheduled recheck
1170/// runs before the coin is used, and an unavailable header is an error that
1171/// prevents the spend.
1172fn readme_spend_guard<H: Headers>(
1173    state: &mut State,
1174    params: &Params,
1175    headers: &H,
1176) -> std::result::Result<bool, CheckError> {
1177    state.step(params, headers, Input::Host(HostAction::SpendAttempt))?;
1178    state.step(params, headers, Input::Evidence(Evidence::Recheck))?;
1179    Ok(matches!(state.word(), Word::Mined(_)) && !state.suspect())
1180}
1181
1182/// Run the spend guard for a coin of the tracked transaction `txid`, as
1183/// the host: the row is loaded (its stored path checked against a fresh
1184/// header), the README's guard runs on it, and the row is stored again.
1185/// `Ok(true)` only for a checked, non-suspect `mined`; `Err` when the
1186/// header service could not answer. A transaction the host does not track
1187/// is not mined.
1188pub async fn spend_guard<V: WalletServices + ?Sized, C: Clock>(
1189    storage: &StorageSqlx,
1190    services: &V,
1191    clock: &C,
1192    opts: &TickOptions,
1193    txid: &str,
1194) -> Result<std::result::Result<bool, CheckError>> {
1195    ensure_tables(storage).await?;
1196    let Some(row) = load_row(storage, txid).await? else {
1197        return Ok(Ok(false));
1198    };
1199    let params = opts.params();
1200    let snapshot = HeaderSnapshot::take(services, row.proof_height()).await;
1201    let (mut state, fault) = row.replay(&params, &snapshot);
1202    match fault {
1203        // Could not look (no tip, no header at the height, a path that
1204        // does not parse): an error, no spend, and the row is left as it
1205        // was. "Could not look" changes no word.
1206        Some(fault) if !matches!(fault, CheckError::RootMismatch(_)) => return Ok(Err(fault)),
1207        _ => {}
1208    }
1209    let verdict = readme_spend_guard(&mut state, &params, &snapshot);
1210    store_row(storage, &row, &state, clock.now()).await?;
1211    Ok(verdict)
1212}
1213
1214#[cfg(test)]
1215pub(crate) mod tests {
1216    use super::*;
1217    use bsv_wallet_toolbox::services::mock::{MockResponse, MockWalletServices};
1218    use bsv_wallet_toolbox::services::{BlockHeader, GetMerklePathResult};
1219    use bsv_wallet_toolbox::{
1220        WalletStorageWriter, BROADCAST_STATUS_ACCEPTED, BROADCAST_STATUS_MINED,
1221    };
1222    use chrono::Utc;
1223
1224    /// A clock the test sets.
1225    pub(crate) struct FixedClock(pub Timestamp);
1226    impl Clock for FixedClock {
1227        fn now(&self) -> Timestamp {
1228            self.0
1229        }
1230    }
1231
1232    pub(crate) fn opts() -> TickOptions {
1233        TickOptions {
1234            age_threshold: 600,
1235            max_asks: 20,
1236        }
1237    }
1238
1239    /// A migrated in-memory wallet with one user.
1240    pub(crate) async fn wallet_storage() -> (StorageSqlx, i64) {
1241        let storage = StorageSqlx::in_memory().await.unwrap();
1242        storage
1243            .migrate("tracker-tests", &("02".to_string() + &"ab".repeat(32)))
1244            .await
1245            .unwrap();
1246        storage.make_available().await.unwrap();
1247        let (user, _) = storage
1248            .find_or_insert_user(&("02".to_string() + &"cd".repeat(32)))
1249            .await
1250            .unwrap();
1251        (storage, user.user_id)
1252    }
1253
1254    pub(crate) async fn insert_tx(storage: &StorageSqlx, user_id: i64, txid: &str, status: &str) {
1255        sqlx::query(
1256            "INSERT INTO transactions (user_id, status, reference, is_outgoing, satoshis, version, lock_time, description, txid, raw_tx, created_at, updated_at) \
1257             VALUES (?, ?, ?, 1, 0, 1, 0, 'd', ?, X'01000000', ?, ?)",
1258        )
1259        .bind(user_id)
1260        .bind(status)
1261        .bind(&txid[..8])
1262        .bind(txid)
1263        .bind(Utc::now())
1264        .bind(Utc::now())
1265        .execute(storage.pool())
1266        .await
1267        .unwrap();
1268    }
1269
1270    /// A header at `height` whose merkle root is `root`.
1271    pub(crate) fn header(height: u32, root: &str, nonce: u32) -> BlockHeader {
1272        BlockHeader {
1273            version: 1,
1274            previous_hash: "00".repeat(32),
1275            merkle_root: root.to_string(),
1276            time: 1_700_000_000,
1277            bits: 0x1d00_ffff,
1278            nonce,
1279            hash: String::new(),
1280            height,
1281        }
1282    }
1283
1284    /// A path proving `txid` alone in a block at `height`: its root is the
1285    /// txid.
1286    pub(crate) fn path_of(txid: &str, height: u32) -> String {
1287        MerklePath::from_coinbase_txid(txid, height).to_hex()
1288    }
1289
1290    pub(crate) fn proof_answer(path: Option<String>) -> MockResponse<GetMerklePathResult> {
1291        MockResponse::Success(GetMerklePathResult {
1292            name: Some("Services".to_string()),
1293            merkle_path: path,
1294            header: None,
1295            error: None,
1296            notes: vec![],
1297        })
1298    }
1299
1300    /// E3 (the rulings of 2026-10-09). The wallet's own unproven
1301    /// transaction is tracked, a broadcaster's acceptance announces it, and
1302    /// its age asks for a proof: one request, for that transaction. The
1303    /// proof checks against the header service's header and the word is
1304    /// `mined`. The row is in the wallet's database, and the next pass
1305    /// reads no mined row and asks nothing.
1306    #[tokio::test]
1307    async fn a_pass_tracks_announces_and_proves_one_named_transaction() {
1308        let (storage, user) = wallet_storage().await;
1309        let txid = "a1".repeat(32);
1310        insert_tx(&storage, user, &txid, "unproven").await;
1311        storage
1312            .record_broadcast_status(&txid, "arcade", BROADCAST_STATUS_ACCEPTED)
1313            .await
1314            .unwrap();
1315        let services = MockWalletServices::builder()
1316            .get_merkle_path_response(proof_answer(Some(path_of(&txid, 900))))
1317            .build();
1318        services.set_header_for_height(header(900, &txid, 0));
1319        let now = Utc::now().timestamp() as u64;
1320
1321        // Young: announced on the hint, nothing asked.
1322        let report = tick(&storage, &services, &FixedClock(now), &opts())
1323            .await
1324            .unwrap();
1325        assert_eq!((report.adopted, report.hints, report.asked), (1, 1, 0));
1326        assert_eq!(
1327            stored_word(&storage, &txid).await.unwrap(),
1328            Some(("announced".to_string(), None, None))
1329        );
1330        assert_eq!(services.call_count("get_merkle_path"), 0);
1331
1332        // Past the age threshold: the re-ask runs, the proof checks.
1333        let report = tick(&storage, &services, &FixedClock(now + 601), &opts())
1334            .await
1335            .unwrap();
1336        assert_eq!(report.asked, 1);
1337        assert_eq!(report.mined, vec![txid.clone()]);
1338        assert_eq!(
1339            stored_word(&storage, &txid).await.unwrap(),
1340            Some(("mined".to_string(), Some(900), None))
1341        );
1342        assert_eq!(services.call_count("get_merkle_path"), 1);
1343
1344        // A mined row is not read again while no header moved.
1345        let report = tick(&storage, &services, &FixedClock(now + 1300), &opts())
1346            .await
1347            .unwrap();
1348        assert_eq!((report.tracked, report.asked), (0, 0));
1349        assert_eq!(services.call_count("get_merkle_path"), 1);
1350    }
1351
1352    /// The evidence rule in the host: a broadcaster's "mined" is a hint. It
1353    /// asks for the proof at once; with no proof the word does not move,
1354    /// and the same transaction is not asked again until its pause is over.
1355    #[tokio::test]
1356    async fn a_broadcasters_mined_is_a_hint_that_asks_for_the_proof() {
1357        let (storage, user) = wallet_storage().await;
1358        let txid = "b2".repeat(32);
1359        insert_tx(&storage, user, &txid, "unproven").await;
1360        storage
1361            .record_broadcast_status(&txid, "arcade", BROADCAST_STATUS_MINED)
1362            .await
1363            .unwrap();
1364        let services = MockWalletServices::builder()
1365            .get_merkle_path_response(proof_answer(None))
1366            .build();
1367        let now = Utc::now().timestamp() as u64;
1368
1369        let report = tick(&storage, &services, &FixedClock(now), &opts())
1370            .await
1371            .unwrap();
1372        assert_eq!(report.asked, 1);
1373        assert_eq!(report.no_proof, vec![txid.clone()]);
1374        assert!(report.mined.is_empty());
1375        let (word, height, reask) = stored_word(&storage, &txid).await.unwrap().unwrap();
1376        assert_eq!((word.as_str(), height), ("built", None));
1377        assert_eq!(reask.as_deref(), Some("hint:arcade|mined"));
1378
1379        // One minute later: the pause is not over, nothing is asked.
1380        let report = tick(&storage, &services, &FixedClock(now + 60), &opts())
1381            .await
1382            .unwrap();
1383        assert_eq!(report.asked, 0);
1384        // After the pause it is asked again, and the pause doubles.
1385        let report = tick(&storage, &services, &FixedClock(now + 600), &opts())
1386            .await
1387            .unwrap();
1388        assert_eq!(report.asked, 1);
1389        let report = tick(&storage, &services, &FixedClock(now + 1200), &opts())
1390            .await
1391            .unwrap();
1392        assert_eq!(report.asked, 0);
1393        assert_eq!(services.call_count("get_merkle_path"), 2);
1394    }
1395
1396    /// A header we stood on moved: the mined rows at or above it are read
1397    /// again and asked again, and a mined row below it is not touched.
1398    #[tokio::test]
1399    async fn a_moved_header_reasks_the_rows_at_or_above_it_and_no_other() {
1400        let (storage, user) = wallet_storage().await;
1401        let low = "c3".repeat(32);
1402        let high = "d4".repeat(32);
1403        let services = MockWalletServices::builder()
1404            .get_merkle_path_response(MockResponse::Sequence(vec![
1405                proof_answer(Some(path_of(&low, 100))),
1406                proof_answer(Some(path_of(&high, 200))),
1407                proof_answer(Some(path_of(&high, 201))),
1408            ]))
1409            .build();
1410        services.set_header_for_height(header(100, &low, 0));
1411        services.set_header_for_height(header(200, &high, 0));
1412        services.set_tip_header(header_with_hash(200, &high, 0));
1413        let now = Utc::now().timestamp() as u64;
1414        for txid in [&low, &high] {
1415            insert_tx(&storage, user, txid, "unproven").await;
1416            storage
1417                .record_broadcast_status(txid, "arcade", BROADCAST_STATUS_MINED)
1418                .await
1419                .unwrap();
1420        }
1421        let report = tick(&storage, &services, &FixedClock(now), &opts())
1422            .await
1423            .unwrap();
1424        assert_eq!(report.mined.len(), 2, "{report:?}");
1425
1426        // Height 200 is replaced by a block without `high`; it is mined
1427        // again at 201.
1428        services.set_header_for_height(header(200, &"ee".repeat(32), 7));
1429        services.set_header_for_height(header(201, &high, 0));
1430        services.set_tip_header(header_with_hash(201, &high, 0));
1431        let report = tick(&storage, &services, &FixedClock(now + 60), &opts())
1432            .await
1433            .unwrap();
1434        assert_eq!(report.moved_from, Some(200));
1435        assert_eq!(
1436            report.tracked, 1,
1437            "only the row at or above the moved header"
1438        );
1439        assert_eq!(report.moved, vec![high.clone()]);
1440        assert_eq!(report.mined, vec![high.clone()]);
1441        assert_eq!(
1442            stored_word(&storage, &high).await.unwrap().unwrap().1,
1443            Some(201)
1444        );
1445        assert_eq!(
1446            stored_word(&storage, &low).await.unwrap().unwrap().1,
1447            Some(100)
1448        );
1449    }
1450
1451    /// E2 (the rulings of 2026-10-09). A transaction a broadcaster called
1452    /// not final is held with its BEEF: nothing is asked or announced
1453    /// before its lock time; at the lock time it is announced again; "not
1454    /// final" again (the node's clock is behind ours) waits one pause; an
1455    /// acceptance announces it, and from there its age asks for the proof.
1456    #[tokio::test]
1457    async fn a_not_final_transaction_is_announced_again_at_its_lock_time() {
1458        use bsv_wallet_toolbox::services::mock::{
1459            error_post_beef_result, success_post_beef_result,
1460        };
1461        let (storage, _user) = wallet_storage().await;
1462        let txid = "c7".repeat(32);
1463        let lock = 2_000_000_000u64;
1464        let services = MockWalletServices::builder()
1465            .post_beef_response(MockResponse::Sequence(vec![
1466                MockResponse::Success(vec![error_post_beef_result("Arcade", "476 non-final")]),
1467                MockResponse::Success(vec![success_post_beef_result("Arcade", &[&txid])]),
1468            ]))
1469            .get_merkle_path_response(proof_answer(None))
1470            .build();
1471
1472        let hint = Hint::new(
1473            "Arcade|rejected",
1474            HintStatus::Rejected {
1475                reason: "476 non-final".to_string(),
1476            },
1477            lock - 3600,
1478        );
1479        let again = Reannounce {
1480            beef: vec![1, 2, 3],
1481            at: lock,
1482            on_accept: None,
1483        };
1484        let word = heard(
1485            &storage,
1486            &FixedClock(lock - 3600),
1487            &opts(),
1488            &txid,
1489            hint,
1490            Some(lock),
1491            Some(again),
1492        )
1493        .await
1494        .unwrap();
1495        assert_eq!(
1496            word,
1497            (
1498                "built".to_string(),
1499                Some("hint:Arcade|rejected".to_string())
1500            )
1501        );
1502
1503        // Before the lock time: nothing announced, no proof asked.
1504        let report = tick(&storage, &services, &FixedClock(lock - 1), &opts())
1505            .await
1506            .unwrap();
1507        assert_eq!((report.tracked, report.asked), (1, 0), "{report:?}");
1508        assert_eq!(services.call_count("post_beef"), 0);
1509
1510        // At the lock time: announced again; the node still says not final.
1511        let report = tick(&storage, &services, &FixedClock(lock), &opts())
1512            .await
1513            .unwrap();
1514        assert_eq!(report.not_final, vec![txid.clone()]);
1515        assert_eq!(services.call_count("post_beef"), 1);
1516        let report = tick(&storage, &services, &FixedClock(lock + 60), &opts())
1517            .await
1518            .unwrap();
1519        assert!(report.not_final.is_empty() && report.reannounced.is_empty());
1520        assert_eq!(services.call_count("post_beef"), 1, "one pause between two");
1521
1522        // One pause on: accepted. The word is the tracker's `announced`.
1523        let report = tick(&storage, &services, &FixedClock(lock + 600), &opts())
1524            .await
1525            .unwrap();
1526        assert_eq!(report.reannounced, vec![txid.clone()]);
1527        assert_eq!(services.call_count("post_beef"), 2);
1528        assert_eq!(
1529            stored_word(&storage, &txid).await.unwrap().unwrap().0,
1530            "announced"
1531        );
1532        assert_eq!(services.call_count("get_merkle_path"), 0);
1533
1534        // It is not announced a third time; its age asks for the proof.
1535        let report = tick(&storage, &services, &FixedClock(lock + 1300), &opts())
1536            .await
1537            .unwrap();
1538        assert_eq!(report.asked, 1);
1539        assert_eq!(services.call_count("post_beef"), 2);
1540        assert_eq!(services.call_count("get_merkle_path"), 1);
1541    }
1542
1543    /// The words that mean "not final", and one that does not.
1544    #[test]
1545    fn not_final_is_read_from_the_refusal() {
1546        for said in [
1547            "476 non-final",
1548            "bad-txns-nonfinal",
1549            "tx is not final",
1550            "code 476",
1551        ] {
1552            assert!(says_not_final(said), "{said}");
1553        }
1554        for said in [
1555            "fee too low",
1556            "DOUBLE_SPEND_ATTEMPTED",
1557            "txid 4476aa",
1558            "HTTP 500",
1559        ] {
1560            assert!(!says_not_final(said), "{said}");
1561        }
1562    }
1563
1564    /// A tip header carrying the hash its bytes have.
1565    fn header_with_hash(height: u32, root: &str, nonce: u32) -> BlockHeader {
1566        let mut h = header(height, root, nonce);
1567        h.hash = header_of(&h.to_binary()).unwrap().hash;
1568        h
1569    }
1570
1571    /// The README's wallet spend guard, run by the host over a stored row
1572    /// (bsv-tracker `README.md`, "Wallet: spend guard"). Three runs: the
1573    /// header is the one the proof was checked against; the header service
1574    /// cannot answer; the header moved.
1575    #[tokio::test]
1576    async fn the_readme_spend_guard_runs_as_the_host() {
1577        let (storage, user) = wallet_storage().await;
1578        let txid = "e5".repeat(32);
1579        insert_tx(&storage, user, &txid, "unproven").await;
1580        storage
1581            .record_broadcast_status(&txid, "arcade", BROADCAST_STATUS_MINED)
1582            .await
1583            .unwrap();
1584        let services = MockWalletServices::builder()
1585            .get_merkle_path_response(proof_answer(Some(path_of(&txid, 900))))
1586            .build();
1587        services.set_header_for_height(header(900, &txid, 0));
1588        let clock = FixedClock(Utc::now().timestamp() as u64);
1589        tick(&storage, &services, &clock, &opts()).await.unwrap();
1590
1591        // 1. Mined, checked, not suspect: the coin may be spent.
1592        let guard = spend_guard(&storage, &services, &clock, &opts(), &txid)
1593            .await
1594            .unwrap();
1595        println!(
1596            "spend guard, header in place:   {guard:?}, stored {:?}",
1597            stored_word(&storage, &txid).await.unwrap().unwrap()
1598        );
1599        assert_eq!(guard, Ok(true));
1600
1601        // 2. No tip from the header service: an error, and no spend.
1602        services.set_tip_unavailable();
1603        let guard = spend_guard(&storage, &services, &clock, &opts(), &txid)
1604            .await
1605            .unwrap();
1606        println!(
1607            "spend guard, no header service: {guard:?}, stored {:?}",
1608            stored_word(&storage, &txid).await.unwrap().unwrap()
1609        );
1610        assert!(matches!(guard, Err(CheckError::Headers(_))), "{guard:?}");
1611        assert_eq!(
1612            stored_word(&storage, &txid).await.unwrap().unwrap().0,
1613            "mined",
1614            "could not look changes no word"
1615        );
1616
1617        // 3. The header at the height moved: not mined, a re-ask pending.
1618        services.set_tip_header(header_with_hash(901, &"00".repeat(32), 1));
1619        services.set_header_for_height(header(900, &"ee".repeat(32), 7));
1620        let guard = spend_guard(&storage, &services, &clock, &opts(), &txid)
1621            .await
1622            .unwrap();
1623        println!(
1624            "spend guard, header moved:      {guard:?}, stored {:?}",
1625            stored_word(&storage, &txid).await.unwrap().unwrap()
1626        );
1627        assert_eq!(guard, Ok(false));
1628
1629        // A transaction the host does not track is not mined.
1630        let guard = spend_guard(&storage, &services, &clock, &opts(), &"f6".repeat(32))
1631            .await
1632            .unwrap();
1633        assert_eq!(guard, Ok(false));
1634
1635        // The guard's own error path, on a state held in memory: mined,
1636        // then the header service loses the height.
1637        let params = opts().params();
1638        services.set_tip_header(header_with_hash(900, &txid, 0));
1639        services.set_header_for_height(header(900, &txid, 0));
1640        let held = HeaderSnapshot::take(&services, [900]).await;
1641        let row = load_row(&storage, &txid).await.unwrap().unwrap();
1642        let (mut state, fault) = row.replay(&params, &held);
1643        assert_eq!(fault, None);
1644        let lost = HeaderSnapshot::take(&services, []).await;
1645        let guard = readme_spend_guard(&mut state, &params, &lost);
1646        println!("spend guard, height unanswered: {guard:?}");
1647        assert_eq!(guard, Err(CheckError::Unavailable(900)));
1648    }
1649}