// Copyright (C) 2024-2026 Whiterun LLC // // This software is licensed under the GNU Affero General Public License (AGPL), version 3.0 or later. // A copy of the license can be found in the LICENSE file or at https://www.gnu.org/licenses/agpl-3.0.html use std::{ sync::{atomic::AtomicBool, atomic::Ordering, Arc}, thread::{self, JoinHandle}, time::Duration, }; use crate::bcmr::parsedbcmr::{parse_bcmr_json, SOURCE_ON_CHAIN}; use crate::db::bcmr::insert_bcmr_data; use crate::db::bcmr::update_bcmr_failure; use crate::{ bcmr::parse_bcmr_from_opreturn, db::{ bcmr::{get_entries_missing_bcmr_download, AuthChainEntry}, DBPool, }, }; use anyhow::*; use bitcoin_hashes::hex::ToHex; use bitcoincash::Script; use log::{info, warn}; use rand::thread_rng; use serde_json::Value; use rand::seq::SliceRandom; use rayon::prelude::*; use std::result::Result::Ok; use super::utilurl::get_url; // We are generous on timeout to allow for slow ipfs gateway const DOWNLOAD_TIMEOUT: Duration = Duration::from_secs(60); const IPFS_GATEWAYS: [&str; 2] = ["https://ipfs.io/ipfs/", "https://w3s.link/ipfs/"]; const MAX_BCMR_SIZE: usize = 1024 * 1024 * 100; // 100 MB const MAX_PARALLEL_DOWNLOADS: usize = 5; pub struct BCMRDownloader { db: DBPool, keep_running: Arc, download_thread: Option>, } const LOOP_SLEEP_TIME: Duration = Duration::from_secs(10); fn get_download_candidates(db: &DBPool) -> Vec { let conn = match db.get() { Ok(c) => c, Err(e) => { warn!("bcmr: Failed to get a db connection: {e}"); return Vec::default(); } }; match get_entries_missing_bcmr_download(&conn) { Ok(c) => c, Err(e) => { warn!("bcmr: Failed to fetch bcmr download candidates: {e}"); Vec::default() } } } fn fetch_bcmr( agent: &ureq::Agent, urls: &[String], expected_hash: &str, ) -> ( Option<(String, String)>, String, /* error */ bool, /* fatal error (don't try again) */ ) { let mut candidate: Option<(String, String)> = None; let mut errors: Vec = Vec::new(); for url in urls { let url = if let Some(stripped) = url.strip_prefix("ipfs://") { format!( "{}{}", IPFS_GATEWAYS.choose(&mut thread_rng()).unwrap(), stripped ) } else if !url.starts_with("https://") { format!("https://{url}") } else { url.to_string() }; let (contents, actual_hash) = match get_url(agent, &url, DOWNLOAD_TIMEOUT, MAX_BCMR_SIZE) { Ok(c) => c, Err(e) => { errors.push(e.to_string()); continue; } }; if actual_hash == expected_hash { return (Some((contents, actual_hash)), String::default(), false); } else { candidate = Some((contents, actual_hash)) } } // No URL's with matching expected hash if candidate.is_some() { (candidate, String::default(), false) } else if errors.is_empty() { (None, "No URLs to fetch BCMR from".to_string(), true) } else { ( None, format!("Errors fetching BCMR: {}", errors.join("; ")), false, ) } } impl BCMRDownloader { pub fn new(db: DBPool) -> Self { Self { db, keep_running: Arc::new(AtomicBool::new(true)), download_thread: None, } } pub fn start(&mut self) -> Result<()> { let db_cpy = self.db.clone(); let keep_running_cpy = self.keep_running.clone(); self.download_thread = Some( thread::Builder::new() .name("bcmr downloader".to_string()) .spawn(move || { // Create thread pool once, reuse across iterations let pool = rayon::ThreadPoolBuilder::new() .num_threads(MAX_PARALLEL_DOWNLOADS) .build() .unwrap(); // Create HTTP agent once, reuse for connection pooling let agent = ureq::Agent::new(); loop { if !keep_running_cpy.load(Ordering::Relaxed) { info!("Exiting bcmr download thread"); return; } let mut queue = get_download_candidates(&db_cpy); // Shuffle so we don't starve any single token if the list is large let mut rng = thread_rng(); queue.shuffle(&mut rng); if queue.is_empty() { thread::sleep(LOOP_SLEEP_TIME); continue; } info!("bcmr: {} tokens need BCMR download", queue.len()); pool.install(|| { queue.par_iter().for_each(|entry| { info!( "bcmr: Downloading BCMR for token {}, txid {}", entry.token_id.to_hex(), entry.txid.to_hex() ); // We may need to record a failure immediately; get a DB conn first. let conn = match db_cpy.get() { Ok(c) => c, Err(e) => { warn!("bcmr: Failed to get BCMR db connection: {e}"); return; } }; // Guard: bcmr_data must be present on the auth head let Some(raw) = entry.bcmr_data.as_ref() else { // Just warn and skip; likely a race. Do not mark fatal. warn!( "bcmr: Missing OP_RETURN BCMR for utxo {} (token {}, txid {}); skipping.", entry.utxo.to_hex(), entry.token_id.to_hex(), entry.txid.to_hex() ); return; }; // parse_bcmr_from_opreturn returns Option<_>; handle None gracefully let bcmr = match parse_bcmr_from_opreturn(&Script::from(raw.clone())) { Some(v) => v, None => { if let Err(e) = update_bcmr_failure( &conn, &entry.utxo, &entry.txid, &entry.token_id, "Invalid BCMR OP_RETURN in DB", true, ) { warn!("bcmr: Failed to set BCMR error: {e}"); } return; } }; let (json_str, error, is_fatal) = fetch_bcmr(&agent, &bcmr.uris, &bcmr.hash.to_hex()); let (json_str, actual_hash) = match json_str { Some(j) => j, None => { if let Err(e) = update_bcmr_failure( &conn, &entry.utxo, &entry.txid, &entry.token_id, &format!("Failed to fetch BCMR: {error}"), is_fatal, ) { warn!("bcmr: Failed to set BCMR error: {e}"); } return; } }; let json: Value = match serde_json::from_str(&json_str) { Ok(b) => b, Err(e) => { if let Err(e2) = update_bcmr_failure( &conn, &entry.utxo, &entry.txid, &entry.token_id, &format!("BCMR invalid JSON error: {e}"), true, ) { warn!("bcmr: Failed to set BCMR error: {e2}"); } return; } }; let bcmr_parsed = match parse_bcmr_json( &json, &entry.token_id, Some(actual_hash), Some(bcmr.hash.to_hex()), SOURCE_ON_CHAIN, ) { Ok(b) => b, Err(err) => { if let Err(e2) = update_bcmr_failure( &conn, &entry.utxo, &entry.txid, &entry.token_id, &format!("BCMR contents error: {err}"), true, ) { warn!("bcmr: Failed to set BCMR error: {e2}"); } return; } }; if let Err(err) = insert_bcmr_data(&conn, &entry.token_id, &entry.utxo, &bcmr_parsed) { warn!("bcmr: Failed to insert BCMR data {err}") } }); }); thread::sleep(LOOP_SLEEP_TIME) } }) .expect("failed to start bcmr download thread"), ); Ok(()) } } impl Drop for BCMRDownloader { fn drop(&mut self) { self.keep_running .store(false, std::sync::atomic::Ordering::SeqCst); if let Some(thread) = self.download_thread.take() { // Wake the thread in case it is sleeping thread.thread().unpark(); match thread.join() { Ok(_) => info!("bcmr download thread done"), Err(e) => warn!("Failed to join bcmr download thread: {e:?}"), } } } }