// 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::{ collections::HashSet, sync::{mpsc::sync_channel, Arc, Mutex}, time::Duration, }; use bitcoin_hashes::{ hex::{FromHex, ToHex}, Hash, }; use bitcoincash::{consensus::deserialize, Block, BlockHash, Network, Transaction, Txid}; use electrum_client_netagnostic::{Client, ElectrumApi, Param}; use log::{debug, info, warn}; use riftenlabs_defi::cauldron::{parse_cauldrons_from_tx, ParsedContract}; use crate::{ bcmr::index_bcmr, chain::{get_new_headers, Chain, StoreBlockUndoer}, crc20::index_crc20, db::{ self, cauldron::{ config::{config_get, config_set}, header::{db_get_header, store_headers}, tx::insert_block_tx, user::insert_user_action, utxo_funding::insert_utxo_funding, utxo_spending::insert_utxo_spending, }, oracle::index_oracle, DBPool, DB, }, electrum::{electrum_fetch_mempool, electrum_get_tip, electrum_get_tx}, signal::shutdown_requested, timeutil::time_now, utiltx::ttor_sorted_kahn, CASHTOKEN_ACTIVATION_HEIGHT, CHIPNET_START_BLOCK, KEY_LAST_INDEXED, }; use anyhow::{bail, Context, Result}; pub fn update_mempool( cauldron_db: DBPool, oracle_db: DBPool, electrum: Arc>, ) -> Result<()> { let our_mempool_txs: HashSet = db::cauldron::mempool::load_mempool(&cauldron_db.get().unwrap())?; let (cauldron_txs, oracle_txs) = electrum_fetch_mempool(&electrum.lock().unwrap())?; let txs_to_delete = our_mempool_txs.difference(&cauldron_txs); let txs_to_add = cauldron_txs.difference(&our_mempool_txs); let txs_to_add: Vec = txs_to_add .into_iter() .filter_map( |txid| match electrum_get_tx(&electrum.lock().unwrap(), txid) { Ok(tx) => Some(tx), Err(e) => { info!("Failed to get mempool tx {txid}: {e}"); None } }, ) .collect(); let txs_to_add = ttor_sorted_kahn(txs_to_add); let mut db_conn = cauldron_db.get().unwrap(); let db_tx = db_conn.transaction()?; db::cauldron::mempool::delete_mempool_txs(&db_tx, txs_to_delete)?; let mut all_cauldrons = vec![]; let current_timestamp = time_now() as u64; for tx in txs_to_add { let cauldrons: Vec = parse_cauldrons_from_tx(&tx); if cauldrons.is_empty() { info!("ignoring non-cauldron tx {}", tx.txid()); continue; } let txid = tx.txid(); db::cauldron::tx::insert_mempool_tx(&db_tx, &txid, current_timestamp)?; info!( "mempool add {} with {} cauldrons", txid.to_hex(), cauldrons.len() ); insert_utxo_funding(&db_tx, &cauldrons, &txid, false)?; insert_utxo_spending(&db_tx, &cauldrons, &txid, false)?; insert_user_action(&db_tx, &cauldrons, &tx, true)?; all_cauldrons.extend(cauldrons); } db::cauldron::pool::update_pool_history(&db_tx, all_cauldrons, None, Some(current_timestamp))?; db_tx.commit()?; { // oracle updates let mut conn = oracle_db.get().context("failed to get oracle db")?; let tx = conn.transaction()?; // filter entries we have let oracle_txs = oracle_txs.into_iter().filter(|txid| { !db::oracle::has_entry(&tx, txid).expect("Failed to check if entry exists") }); let txs_to_add: Vec = oracle_txs .into_iter() .filter_map( |txid| match electrum_get_tx(&electrum.lock().unwrap(), &txid) { Ok(tx) => Some(tx), Err(e) => { info!("Failed to get mempool tx {txid}: {e}"); None } }, ) .collect(); index_oracle(&tx, &txs_to_add, &BlockHash::all_zeros())?; tx.commit()?; } Ok(()) } pub fn index_blocks( chain: Arc>, db: DB, client: Arc>, bcmr_enabled: bool, network: Option, ) -> Result { let (tip_header, _) = electrum_get_tip(&client.lock().unwrap())?; // Validate network before proceeding let start_block_hash = match network { Some(Network::Chipnet) => BlockHash::from_hex(CHIPNET_START_BLOCK) .map_err(|e| anyhow::anyhow!("Invalid CHIPNET_START_BLOCK: {}", e))?, Some(Network::Bitcoin) => BlockHash::from_hex(CASHTOKEN_ACTIVATION_HEIGHT) .map_err(|e| anyhow::anyhow!("Invalid CASHTOKEN_ACTIVATION_HEIGHT: {}", e))?, None => BlockHash::from_hex(CASHTOKEN_ACTIVATION_HEIGHT) .map_err(|e| anyhow::anyhow!("Invalid CASHTOKEN_ACTIVATION_HEIGHT: {}", e))?, Some(net) => bail!("Unknown network: {:?}", net), }; let (block_send, block_recv) = sync_channel::>(10); // Update header chain (and undo any blocks that may have reorged away) { let chain = chain.lock().unwrap(); if tip_header.block_hash() != chain.tip_hash() { debug!( "Updating header chain from {} to {}", chain.tip_hash().to_hex(), tip_header.block_hash().to_hex() ); let new_headers = get_new_headers(&client.lock().unwrap(), &chain, &tip_header.block_hash())?; debug!("Storing headers"); { let mut db_conn = db.cauldron_w.get().unwrap(); for chunk in new_headers.chunks(100000) { let db_tx = db_conn.transaction()?; store_headers(&db_tx, chunk)?; db_tx.commit()?; } } let undoer = StoreBlockUndoer::new(db.clone())?; chain.update(undoer, new_headers, None)?; debug!("Header update done"); } } let db_cpy = db.clone(); std::thread::spawn(move || { let db = db_cpy; let last_indexed = config_get(&db.cauldron_r.get().unwrap(), KEY_LAST_INDEXED).unwrap(); let mut last_indexed = if let Some(last) = last_indexed { BlockHash::from_hex(&last).unwrap() } else { // No last_indexed found, use network-specific starting point start_block_hash }; // Check if last indexed has been orphaned loop { if chain.lock().unwrap().contains(&last_indexed) { break; } info!( "Last indexed block ({}) has been orphaned", last_indexed.to_hex() ); let header = db_get_header(&db.cauldron_r.get().unwrap(), &last_indexed) .context("Failed to get header for last indexed") .unwrap(); last_indexed = header.prev_blockhash; info!( "Last indexed block rolled back to {}", last_indexed.to_hex() ); } loop { if tip_header.block_hash() == last_indexed { if let Err(e) = block_send.send(None) { warn!("Failed to end EOL to block reader: {e}"); } break; } let next_height = chain .lock() .unwrap() .get_block_height(&last_indexed) .expect("last_indexed height not found in main chain") + 1; let res = client .lock() .unwrap() .raw_call("blockchain.block.get", vec![Param::U32(next_height as u32)]) .unwrap(); let block_hex: String = serde_json::from_str(&res.to_string()).unwrap(); let block: Block = deserialize(&hex::decode(&block_hex).unwrap()).unwrap(); let block_hash = block.block_hash(); if let Err(e) = block_send.send(Some(( next_height, chain.lock().unwrap().get_mtp(next_height).unwrap(), block, ))) { warn!("Failed to send block to reader: {e}"); break; } last_indexed = block_hash; } }); loop { // Use recv_timeout to allow checking for shutdown let recv_result = block_recv.recv_timeout(Duration::from_secs(1)); let (block_height, mtp, block) = match recv_result { Ok(Some(data)) => data, Ok(None) => break, // End of blocks signal Err(std::sync::mpsc::RecvTimeoutError::Timeout) => { if shutdown_requested() { info!("Shutdown requested, exiting block indexing"); break; } continue; } Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => { warn!("Block sender disconnected"); break; } }; let mut db_conn = db.cauldron_w.get().unwrap(); let db_tx = db_conn.transaction()?; let blockhash = block.block_hash(); let mut total_cauldrons = 0; let mut all_cauldrons = vec![]; let sorted_txs = ttor_sorted_kahn(block.txdata); for tx in &sorted_txs { let cauldrons = parse_cauldrons_from_tx(tx); if cauldrons.is_empty() { continue; } let txid = tx.txid(); insert_block_tx(&db_tx, &txid, &blockhash, mtp as i64).context("inserting block tx")?; insert_utxo_funding(&db_tx, &cauldrons, &txid, true) .context("inserting funding utxos")?; insert_utxo_spending(&db_tx, &cauldrons, &txid, true) .context("inserting spending utxos")?; insert_user_action(&db_tx, &cauldrons, tx, true).context("inserting user actions")?; total_cauldrons += cauldrons.len(); all_cauldrons.extend(cauldrons); } // Figuring out initial utxo needs to be done on all cauldrons in a block. db::cauldron::pool::update_pool_history(&db_tx, all_cauldrons, Some(mtp), None)?; config_set(&db_tx, KEY_LAST_INDEXED, &blockhash.to_hex()); { // crc20 let mut conn = db.crc20_w.get().context("failed to get crc20 db")?; let tx = conn.transaction()?; index_crc20(&tx, &sorted_txs)?; tx.commit()?; } { // oracle updates let mut conn = db.oracle_w.get().context("failed to get oracle db")?; let tx = conn.transaction()?; index_oracle(&tx, &sorted_txs, &blockhash)?; tx.commit()?; } let autheader_updates = if bcmr_enabled { let mut conn = db.bcmr_w.get().context("failed to get bcmr db")?; let bcmr_db_tx = conn.transaction()?; let updates = index_bcmr(&bcmr_db_tx, &blockhash, sorted_txs)?; bcmr_db_tx.commit()?; updates as i64 } else { -1 }; // cauldron commit needs to come at last; as it tracks what the last successful block indexed was db_tx.commit()?; info!( "Indexed {}; mtp: {}, height {}, {} trades, {} autheader updates.", blockhash.to_hex(), mtp, block_height, total_cauldrons, autheader_updates, ); } Ok(tip_header.block_hash()) }