use std::collections::HashSet;
use std::sync::atomic::AtomicBool;
use std::time::Duration;
use anyhow::Result;
use bsv_wallet_toolbox::monitor::ArcadeEventsTask;
use bsv_wallet_toolbox::services::providers::arcade::statuses;
use bsv_wallet_toolbox::{
ArcadeSseClient, ArcadeStatusEvent, BroadcastMemory, BroadcastStatus, Chain, MonitorStorage,
PoisonOutcome, PoisonReport, Services, StorageSqlx, Wallet, BROADCAST_PROVIDER_NETWORK,
BROADCAST_SEEN_STALE_SECS, BROADCAST_STATUS_MINED, BROADCAST_STATUS_SEEN,
BROADCAST_STATUS_UNKNOWN, PROVIDER_ARCADE_V2,
};
use sqlx::Row;
use crate::broadcast_verify::{BroadcastVerifier, NetworkEvidence, PresenceReport};
pub const DEFAULT_ABSENCE_MINUTES: i64 = 30;
pub const DEFAULT_INTERVAL_SECS: u64 = 60;
pub const DEFAULT_MAX_PROBES: usize = 20;
pub const SERVE_MAX_AGE_HOURS: i64 = 24;
const PROBE_PACE: Duration = Duration::from_millis(350);
const SSE_BUDGET: Duration = Duration::from_secs(8);
const SSE_IDLE: Duration = Duration::from_secs(2);
pub type ServedWallet = Wallet<StorageSqlx, Services>;
#[derive(Debug, Clone)]
pub struct ArcadeSse {
pub url: String,
pub token: String,
}
#[derive(Debug, Clone)]
pub struct ReconcileOptions {
pub execute: bool,
pub max_probes: usize,
pub max_age_hours: Option<i64>,
pub absence_minutes: i64,
pub sse: Option<ArcadeSse>,
}
impl ReconcileOptions {
pub fn for_serve(sse: Option<ArcadeSse>) -> Self {
Self {
execute: true,
max_probes: env_parse("BROADCAST_RECONCILE_MAX_PROBES", DEFAULT_MAX_PROBES),
max_age_hours: Some(SERVE_MAX_AGE_HOURS),
absence_minutes: absence_minutes_from_env(),
sse,
}
}
pub fn for_command(execute: bool, max_probes: usize, sse: Option<ArcadeSse>) -> Self {
Self {
execute,
max_probes,
max_age_hours: None,
absence_minutes: absence_minutes_from_env(),
sse,
}
}
}
pub fn absence_minutes_from_env() -> i64 {
env_parse("BROADCAST_ABSENCE_MINUTES", DEFAULT_ABSENCE_MINUTES)
}
fn env_parse<T: std::str::FromStr>(key: &str, default: T) -> T {
std::env::var(key)
.ok()
.and_then(|v| v.trim().parse::<T>().ok())
.unwrap_or(default)
}
pub fn arcade_sse_for(db_path: &str) -> Option<ArcadeSse> {
crate::services_env::arcade_runtime(db_path)
.ok()
.flatten()
.map(|rt| ArcadeSse {
url: rt.url,
token: rt.callback_token,
})
}
#[derive(Debug, Default, Clone)]
pub struct ReconcileBroadcastsReport {
pub sse_events: u64,
pub sse_fatal: Vec<String>,
pub candidates: usize,
pub fresh: usize,
pub probed: usize,
pub seen: Vec<String>,
pub mined: Vec<String>,
pub held: Vec<String>,
pub inconclusive: Vec<String>,
pub absent: Vec<(String, i64)>,
pub fatal: Vec<String>,
pub retired: Vec<PoisonReport>,
}
impl ReconcileBroadcastsReport {
pub fn retired_txids(&self) -> Vec<String> {
self.retired
.iter()
.filter(|r| r.outcome == PoisonOutcome::Retired)
.flat_map(|r| r.retirable_txids())
.collect()
}
pub fn summary(&self, execute: bool) -> String {
let retired: Vec<&PoisonReport> = self
.retired
.iter()
.filter(|r| r.outcome == PoisonOutcome::Retired)
.collect();
let txs = self.retired_txids().len();
let restored: u32 = retired.iter().map(|r| r.restored).sum();
let restored_sats: i64 = retired.iter().map(|r| r.restored_sats).sum();
let invalidated: u32 = retired.iter().map(|r| r.invalidated).sum();
let invalidated_sats: i64 = retired.iter().map(|r| r.invalidated_sats).sum();
let internalized: usize = retired.iter().map(|r| r.internalized.len()).sum();
let alive = self
.retired
.iter()
.filter(|r| r.outcome == PoisonOutcome::Alive)
.count();
let refused = self
.retired
.iter()
.filter(|r| matches!(r.outcome, PoisonOutcome::Refused { .. }))
.count();
format!(
"reconcile-broadcasts{}: sse_events={} candidates={} fresh={} probed={} seen={} mined={} held={} inconclusive={} absent={} fatal={} retired_roots={} retired_txs={} restored={} ({} sats) invalidated={} ({} sats) internalized={} alive={} refused={}",
if execute { "" } else { " (dry run)" },
self.sse_events,
self.candidates,
self.fresh,
self.probed,
self.seen.len(),
self.mined.len(),
self.held.len(),
self.inconclusive.len(),
self.absent.len(),
self.fatal.len(),
retired.len(),
txs,
restored,
restored_sats,
invalidated,
invalidated_sats,
internalized,
alive,
refused,
)
}
}
pub async fn run_pass(
wallet: &ServedWallet,
verifier: &BroadcastVerifier,
opts: &ReconcileOptions,
) -> Result<ReconcileBroadcastsReport> {
let storage = wallet.storage();
let mut report = ReconcileBroadcastsReport::default();
if let Some(sse) = &opts.sse {
let (events, fatal) = drain_sse(storage, sse).await;
report.sse_events = events;
report.sse_fatal = fatal;
}
let candidates = select_candidates(storage, opts.max_age_hours).await?;
report.candidates = candidates.len();
let records = storage.broadcast_records(None, &candidates).await?;
let now = chrono::Utc::now();
let mut to_probe: Vec<String> = Vec::new();
let mut roots: Vec<(String, &'static str)> = Vec::new();
for txid in &candidates {
let mine: Vec<_> = records.iter().filter(|r| &r.txid == txid).collect();
if mine
.iter()
.any(|r| r.ladder_status() == Some(BroadcastStatus::Rejected))
{
roots.push((txid.clone(), "rejected memory row"));
continue;
}
let fresh = mine.iter().any(|r| {
r.ladder_status().is_some_and(|s| s.is_network_evidence())
&& (now - r.seen_at).num_seconds() <= BROADCAST_SEEN_STALE_SECS
});
if fresh {
report.fresh += 1;
continue;
}
to_probe.push(txid.clone());
}
for (index, txid) in to_probe.iter().take(opts.max_probes).enumerate() {
if index > 0 {
tokio::time::sleep(PROBE_PACE).await;
}
report.probed += 1;
let presence = verifier.verify_report(txid).await;
match classify(&presence) {
Probe::Evidence(evidence) => {
credit_presence(storage, txid, evidence, presence.evidence_provider).await;
match evidence {
NetworkEvidence::Seen => report.seen.push(txid.clone()),
NetworkEvidence::Mined => report.mined.push(txid.clone()),
}
}
Probe::Fatal => {
report.fatal.push(txid.clone());
roots.push((txid.clone(), "fatal verdict from the broadcaster"));
}
Probe::Absent => {
let minutes = absence_minutes(storage, txid).await;
report.absent.push((txid.clone(), minutes));
if minutes >= opts.absence_minutes {
roots.push((txid.clone(), "absent from every network source"));
}
}
Probe::Held => report.held.push(txid.clone()),
Probe::Inconclusive => report.inconclusive.push(txid.clone()),
}
}
for root in sweep_roots(storage).await? {
if !roots.iter().any(|(t, _)| t == &root) {
roots.push((root, "failed transaction with unproven descendants"));
}
}
let mut covered: HashSet<String> = HashSet::new();
for (root, reason) in roots {
if covered.contains(&root) {
continue;
}
let poison = storage
.retire_poisoned_chain(wallet.services(), &root, "invalid", opts.execute)
.await?;
match &poison.outcome {
PoisonOutcome::Retired => {
tracing::warn!(
root = %root,
reason,
txs = poison.chain.len(),
executed = poison.executed,
"reconcile-broadcasts: poisoned chain {}",
if opts.execute { "retired" } else { "would be retired" }
);
}
PoisonOutcome::Alive => {
if opts.execute {
if let Err(e) = storage
.record_broadcast_status(
&root,
BROADCAST_PROVIDER_NETWORK,
BROADCAST_STATUS_SEEN,
)
.await
{
tracing::warn!(root = %root, error = %e, "could not record the alive verdict");
}
}
tracing::info!(root = %root, reason, "reconcile-broadcasts: root is alive per the status service, kept");
}
PoisonOutcome::Refused { proven_txid } => {
tracing::error!(root = %root, proven = %proven_txid, reason, "reconcile-broadcasts: retire refused, a descendant is proven");
}
PoisonOutcome::NotFound => {}
}
for tx in &poison.chain {
covered.insert(tx.txid.clone());
}
report.retired.push(poison);
}
Ok(report)
}
pub async fn run_sweep(wallet: &ServedWallet, execute: bool) -> Result<Vec<PoisonReport>> {
let storage = wallet.storage();
let mut reports = Vec::new();
let mut covered: HashSet<String> = HashSet::new();
let mut roots = sweep_roots(storage).await?;
roots.extend(rejected_roots(storage).await?);
for root in roots {
if covered.contains(&root) {
continue;
}
let poison = storage
.retire_poisoned_chain(wallet.services(), &root, "invalid", execute)
.await?;
for tx in &poison.chain {
covered.insert(tx.txid.clone());
}
reports.push(poison);
}
Ok(reports)
}
pub fn spawn_serve_loop(
wallet: std::sync::Arc<ServedWallet>,
chain: Chain,
db_path: &str,
) -> Option<tokio::task::JoinHandle<()>> {
let enabled = std::env::var("BROADCAST_RECONCILE")
.map(|v| v != "0" && !v.eq_ignore_ascii_case("false"))
.unwrap_or(true);
if !enabled {
tracing::info!("broadcast reconcile loop disabled (BROADCAST_RECONCILE=0)");
return None;
}
let opts = ReconcileOptions::for_serve(arcade_sse_for(db_path));
let interval_secs =
env_parse("BROADCAST_RECONCILE_INTERVAL_SECS", DEFAULT_INTERVAL_SECS).max(5);
let verifier = BroadcastVerifier::single_pass(chain);
tracing::info!(
interval_secs,
max_probes = opts.max_probes,
absence_minutes = opts.absence_minutes,
sse = opts.sse.is_some(),
"broadcast reconcile loop started"
);
Some(tokio::spawn(async move {
let mut interval = tokio::time::interval(Duration::from_secs(interval_secs));
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
interval.tick().await;
loop {
interval.tick().await;
match run_pass(&wallet, &verifier, &opts).await {
Ok(report) => {
let quiet = report.probed == 0
&& report.sse_events == 0
&& report.retired.is_empty()
&& report.candidates == 0;
if quiet {
tracing::debug!("{}", report.summary(true));
} else {
tracing::info!("{}", report.summary(true));
}
}
Err(e) => tracing::warn!(error = %e, "broadcast reconcile pass failed"),
}
}
}))
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Probe {
Evidence(NetworkEvidence),
Fatal,
Absent,
Held,
Inconclusive,
}
fn classify(report: &PresenceReport) -> Probe {
if let Some(evidence) = report.evidence {
return Probe::Evidence(evidence);
}
if report.broadcaster_fatal {
return Probe::Fatal;
}
if report.network_absent {
return Probe::Absent;
}
match report.verification {
crate::broadcast_verify::BroadcastVerification::Confirmed => Probe::Held,
_ => Probe::Inconclusive,
}
}
async fn credit_presence(
storage: &StorageSqlx,
txid: &str,
evidence: NetworkEvidence,
provider: &'static str,
) {
if let Err(e) = storage
.mark_transaction_seen_on_network_by(txid, provider)
.await
{
tracing::warn!(txid = %txid, error = %e, "could not record the network presence");
}
if evidence == NetworkEvidence::Mined {
if let Err(e) = storage
.record_broadcast_status(txid, provider, BROADCAST_STATUS_MINED)
.await
{
tracing::warn!(txid = %txid, error = %e, "could not record the mined evidence");
}
}
match unproven_parents(storage, txid).await {
Ok(parents) if !parents.is_empty() => {
if let Err(e) = storage
.sqlx_broadcast_memory()
.record_broadcast_status_many(provider, BROADCAST_STATUS_SEEN, &parents)
.await
{
tracing::warn!(txid = %txid, error = %e, "could not credit the parents");
}
tracing::info!(
txid = %txid,
?evidence,
provider,
parents = parents.len(),
"network presence: the transaction and its unproven parents connected"
);
}
Ok(_) => {
tracing::info!(txid = %txid, ?evidence, provider, "network presence recorded");
}
Err(e) => tracing::warn!(txid = %txid, error = %e, "could not list the parents"),
}
}
async fn absence_minutes(storage: &StorageSqlx, txid: &str) -> i64 {
if let Err(e) = storage
.record_broadcast_status(txid, BROADCAST_PROVIDER_NETWORK, BROADCAST_STATUS_UNKNOWN)
.await
{
tracing::warn!(txid = %txid, error = %e, "could not record the absence");
return 0;
}
match storage
.broadcast_status_of(txid, BROADCAST_PROVIDER_NETWORK)
.await
{
Ok(Some(row)) if row.ladder_status() == Some(BroadcastStatus::Unknown) => {
let minutes = (chrono::Utc::now() - row.seen_at).num_minutes().max(0);
tracing::info!(
txid = %txid,
minutes,
"absent from every network source (absence clock running)"
);
minutes
}
_ => 0,
}
}
async fn select_candidates(
storage: &StorageSqlx,
max_age_hours: Option<i64>,
) -> Result<Vec<String>> {
let modifier = max_age_hours.map(|h| format!("-{} hours", h.max(0)));
let rows = sqlx::query(
"SELECT txid FROM transactions \
WHERE txid IS NOT NULL \
AND (status = 'unproven' \
OR (status = 'sending' AND datetime(created_at) <= datetime('now', '-600 seconds'))) \
AND (? IS NULL OR datetime(created_at) >= datetime('now', ?)) \
ORDER BY datetime(created_at) ASC, transaction_id ASC",
)
.bind(&modifier)
.bind(&modifier)
.fetch_all(storage.pool())
.await?;
let mut seen = HashSet::new();
Ok(rows
.iter()
.map(|r| r.get::<String, _>("txid"))
.filter(|t| seen.insert(t.clone()))
.collect())
}
async fn unproven_parents(storage: &StorageSqlx, txid: &str) -> Result<Vec<String>> {
let rows = sqlx::query(
"SELECT DISTINCT t.txid FROM outputs o \
JOIN transactions t ON o.transaction_id = t.transaction_id \
WHERE o.spent_by = (SELECT transaction_id FROM transactions WHERE txid = ? LIMIT 1) \
AND t.status IN ('unproven', 'sending') AND t.txid IS NOT NULL",
)
.bind(txid)
.fetch_all(storage.pool())
.await?;
Ok(rows.iter().map(|r| r.get::<String, _>("txid")).collect())
}
async fn sweep_roots(storage: &StorageSqlx) -> Result<Vec<String>> {
let rows = sqlx::query(
"SELECT DISTINCT p.txid FROM transactions p \
JOIN outputs o ON o.transaction_id = p.transaction_id \
JOIN transactions c ON c.transaction_id = o.spent_by \
WHERE p.status = 'failed' AND p.txid IS NOT NULL \
AND c.status IN ('unproven', 'sending', 'nosend') \
ORDER BY p.transaction_id ASC",
)
.fetch_all(storage.pool())
.await?;
Ok(rows.iter().map(|r| r.get::<String, _>("txid")).collect())
}
async fn rejected_roots(storage: &StorageSqlx) -> Result<Vec<String>> {
let rows = sqlx::query(
"SELECT DISTINCT t.txid FROM transactions t \
JOIN broadcast_seen b ON b.txid = t.txid \
WHERE t.status IN ('unproven', 'sending') AND b.status = 'rejected' \
ORDER BY t.transaction_id ASC",
)
.fetch_all(storage.pool())
.await?;
Ok(rows.iter().map(|r| r.get::<String, _>("txid")).collect())
}
async fn drain_sse(storage: &StorageSqlx, sse: &ArcadeSse) -> (u64, Vec<String>) {
let mut client = match ArcadeSseClient::new(&sse.url, &sse.token) {
Ok(c) => c,
Err(e) => {
tracing::warn!(error = %e, "arcade SSE client could not be built");
return (0, Vec::new());
}
};
let (tx, mut rx) = tokio::sync::mpsc::channel::<ArcadeStatusEvent>(256);
let stream = tokio::spawn(async move { client.stream_once(tx).await });
let deadline = tokio::time::Instant::now() + SSE_BUDGET;
let trigger = AtomicBool::new(false);
let mut applied = 0u64;
let mut fatal = Vec::new();
loop {
let remaining = deadline.saturating_duration_since(tokio::time::Instant::now());
if remaining.is_zero() {
break;
}
match tokio::time::timeout(SSE_IDLE.min(remaining), rx.recv()).await {
Ok(Some(ev)) => {
match ArcadeEventsTask::<StorageSqlx>::apply_event(storage, &ev, &trigger).await {
Ok(updated) => {
applied += 1;
tracing::debug!(txid = %ev.txid, status = %ev.tx_status, updated, "arcade SSE status");
}
Err(e) => {
tracing::warn!(txid = %ev.txid, status = %ev.tx_status, error = %e, "arcade SSE status not applied");
}
}
if ev.tx_status == statuses::MINED {
if let Err(e) = storage
.record_broadcast_status(
&ev.txid,
PROVIDER_ARCADE_V2,
BROADCAST_STATUS_MINED,
)
.await
{
tracing::warn!(txid = %ev.txid, error = %e, "could not record the mined verdict");
}
}
if bsv_wallet_toolbox::is_fatal_status(&ev.tx_status) {
fatal.push(ev.txid.clone());
}
}
Ok(None) => break,
Err(_) => break,
}
}
stream.abort();
(applied, fatal)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::broadcast_verify::BroadcastVerification;
#[test]
fn a_probe_report_classifies_in_evidence_order() {
let mut r = PresenceReport::from_verification(BroadcastVerification::Confirmed);
assert_eq!(classify(&r), Probe::Held);
r.network_absent = true;
assert_eq!(
classify(&r),
Probe::Absent,
"held by the broadcaster, absent from the network"
);
r.broadcaster_fatal = true;
assert_eq!(classify(&r), Probe::Fatal);
r.evidence = Some(NetworkEvidence::Seen);
assert_eq!(classify(&r), Probe::Evidence(NetworkEvidence::Seen));
let i = PresenceReport::from_verification(BroadcastVerification::Inconclusive);
assert_eq!(classify(&i), Probe::Inconclusive);
let mut rejected = PresenceReport::from_verification(BroadcastVerification::Rejected);
rejected.network_absent = true;
assert_eq!(classify(&rejected), Probe::Absent);
}
#[test]
fn options_read_the_env_knobs_with_defaults() {
let serve = ReconcileOptions::for_serve(None);
assert!(serve.execute);
assert_eq!(serve.max_age_hours, Some(SERVE_MAX_AGE_HOURS));
let command = ReconcileOptions::for_command(false, 7, None);
assert!(!command.execute);
assert_eq!(command.max_probes, 7);
assert_eq!(command.max_age_hours, None);
assert!(command.absence_minutes >= 1);
}
#[test]
fn the_summary_is_one_line() {
let report = ReconcileBroadcastsReport::default();
let line = report.summary(false);
assert!(line.starts_with("reconcile-broadcasts (dry run):"));
assert!(!line.contains('\n'));
assert!(report.retired_txids().is_empty());
}
async fn storage_with_chain() -> StorageSqlx {
use bsv_wallet_toolbox::WalletStorageWriter;
let storage = StorageSqlx::in_memory().await.unwrap();
storage
.migrate("reconcile-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();
let basket = storage
.find_or_create_default_basket(user.user_id)
.await
.unwrap()
.basket_id;
let now = chrono::Utc::now();
let mut ids = Vec::new();
for (txid, status, created) in [
("aa".repeat(32), "failed", "2020-01-01T00:00:00+00:00"),
("bb".repeat(32), "unproven", "2020-01-02T00:00:00+00:00"),
("cc".repeat(32), "sending", "2020-01-03 00:00:00"),
("dd".repeat(32), "unproven", ""),
] {
let created_at: String = if created.is_empty() {
now.to_rfc3339()
} else {
created.to_string()
};
let id = 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.user_id)
.bind(status)
.bind(&txid[..6])
.bind(&txid)
.bind(&created_at)
.bind(now)
.execute(storage.pool())
.await
.unwrap()
.last_insert_rowid();
ids.push(id);
}
let lock = hex::decode("76a914dbc0a7c84983c5bf199b7b2d41b3acf0408ee5aa88ac").unwrap();
for (tx_row, txid, spent_by) in [
(ids[0], "aa".repeat(32), Some(ids[1])),
(ids[1], "bb".repeat(32), Some(ids[3])),
(ids[3], "dd".repeat(32), None),
] {
sqlx::query(
"INSERT INTO outputs (user_id, transaction_id, basket_id, vout, satoshis, locking_script, txid, type, spendable, change, spent_by, provided_by, purpose, output_description, created_at, updated_at) \
VALUES (?, ?, ?, 0, 1000, ?, ?, 'P2PKH', ?, 1, ?, 'storage', 'change', 'c', ?, ?)",
)
.bind(user.user_id)
.bind(tx_row)
.bind(basket)
.bind(&lock)
.bind(&txid)
.bind(spent_by.is_none() as i64)
.bind(spent_by)
.bind(now)
.bind(now)
.execute(storage.pool())
.await
.unwrap();
}
storage
}
#[tokio::test]
async fn candidates_parents_and_sweep_roots_come_from_the_wallet_graph() {
let storage = storage_with_chain().await;
let all = select_candidates(&storage, None).await.unwrap();
assert_eq!(
all,
vec!["bb".repeat(32), "cc".repeat(32), "dd".repeat(32)],
"unproven and stale sending, oldest first"
);
let recent = select_candidates(&storage, Some(24)).await.unwrap();
assert_eq!(recent, vec!["dd".repeat(32)]);
let parents = unproven_parents(&storage, &"dd".repeat(32)).await.unwrap();
assert_eq!(parents, vec!["bb".repeat(32)]);
assert!(
unproven_parents(&storage, &"bb".repeat(32))
.await
.unwrap()
.is_empty(),
"a failed parent is not credited"
);
let roots = sweep_roots(&storage).await.unwrap();
assert_eq!(
roots,
vec!["aa".repeat(32)],
"failed with unproven children"
);
storage
.record_broadcast_status(&"cc".repeat(32), "ArcadeV2", "rejected")
.await
.unwrap();
assert_eq!(
rejected_roots(&storage).await.unwrap(),
vec!["cc".repeat(32)]
);
}
#[tokio::test]
async fn the_absence_clock_starts_on_the_first_absence_and_keeps_its_start() {
let storage = storage_with_chain().await;
let txid = "bb".repeat(32);
assert_eq!(absence_minutes(&storage, &txid).await, 0);
sqlx::query(
"UPDATE broadcast_seen SET seen_at = datetime('now', '-31 minutes') WHERE txid = ?",
)
.bind(&txid)
.execute(storage.pool())
.await
.unwrap();
let minutes = absence_minutes(&storage, &txid).await;
assert!((30..=32).contains(&minutes), "{}", minutes);
storage
.record_broadcast_status(&txid, BROADCAST_PROVIDER_NETWORK, BROADCAST_STATUS_SEEN)
.await
.unwrap();
assert_eq!(absence_minutes(&storage, &txid).await, 0);
}
}