bsv_wallet_cli/
arc_ingest.rs1use anyhow::{anyhow, Result};
17use bsv_wallet_toolbox::monitor::ArcadeEventsTask;
18use bsv_wallet_toolbox::{MonitorStorage, ProofIngestOutcome, StorageSqlx};
19use std::sync::atomic::AtomicBool;
20
21#[derive(Debug, Clone, PartialEq, Eq)]
23pub enum IngestAction {
24 ProofIngested,
26 ProofRejected(String),
28 StatusApplied,
30 StatusIgnored,
32}
33
34pub async fn ingest_arc_payload(
36 storage: &StorageSqlx,
37 payload: &serde_json::Value,
38) -> Result<IngestAction> {
39 let txid = payload
40 .get("txid")
41 .and_then(|v| v.as_str())
42 .ok_or_else(|| anyhow!("payload missing txid"))?;
43 if txid.len() != 64 || !txid.chars().all(|c| c.is_ascii_hexdigit()) {
44 return Err(anyhow!("invalid txid"));
45 }
46 let tx_status = payload
47 .get("txStatus")
48 .and_then(|v| v.as_str())
49 .unwrap_or("");
50
51 let merkle_path_hex = payload
52 .get("merklePath")
53 .and_then(|v| v.as_str())
54 .filter(|s| !s.is_empty());
55
56 if let Some(mp_hex) = merkle_path_hex {
57 let merkle_path = hex::decode(mp_hex).map_err(|e| anyhow!("merklePath not hex: {}", e))?;
58 let block_height = payload
59 .get("blockHeight")
60 .and_then(|v| v.as_u64())
61 .ok_or_else(|| anyhow!("merklePath payload missing blockHeight"))?
62 as u32;
63 let block_hash = payload
64 .get("blockHash")
65 .and_then(|v| v.as_str())
66 .unwrap_or_default();
67
68 let _ = storage.mark_transaction_seen_on_network(txid).await;
71
72 match storage
73 .ingest_merkle_proof(txid, &merkle_path, block_height, block_hash, None)
74 .await
75 .map_err(|e| anyhow!("ingest_merkle_proof: {}", e))?
76 {
77 ProofIngestOutcome::Ingested(status) => {
78 tracing::info!(
79 txid = %txid,
80 block_height = ?status.block_height,
81 "arc-callback: merkle proof ingested"
82 );
83 Ok(IngestAction::ProofIngested)
84 }
85 ProofIngestOutcome::InvalidMerkleRoot { computed_root } => {
86 tracing::warn!(txid = %txid, computed_root = %computed_root, "arc-callback: proof rejected (invalid merkle root)");
87 Ok(IngestAction::ProofRejected("invalid merkle root".into()))
88 }
89 ProofIngestOutcome::InvalidProof(e) => {
90 tracing::warn!(txid = %txid, error = %e, "arc-callback: proof rejected (unparseable)");
91 Ok(IngestAction::ProofRejected(e))
92 }
93 ProofIngestOutcome::TrackerError(e) => {
94 tracing::warn!(txid = %txid, error = %e, "arc-callback: proof deferred (ChainTracker error) — polling sync will retry");
95 Ok(IngestAction::ProofRejected(format!("tracker error: {}", e)))
96 }
97 }
98 } else {
99 let trigger = AtomicBool::new(false);
106 let updated =
107 ArcadeEventsTask::<StorageSqlx>::apply_status_event(storage, txid, tx_status, &trigger)
108 .await
109 .map_err(|e| anyhow!("apply_status_event: {}", e))?;
110 tracing::info!(
111 txid = %txid,
112 status = %tx_status,
113 updated,
114 "arc-callback: status webhook received"
115 );
116 if updated {
117 Ok(IngestAction::StatusApplied)
118 } else {
119 Ok(IngestAction::StatusIgnored)
120 }
121 }
122}