387 lines
13 KiB
Rust
387 lines
13 KiB
Rust
// 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::{atomic::Ordering, 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,
|
|
DB,
|
|
},
|
|
electrum::{electrum_fetch_mempool, electrum_get_tip, electrum_get_tx},
|
|
signal::shutdown_requested,
|
|
timeutil::time_now,
|
|
utiltx::ttor_sorted_kahn,
|
|
IbdState, CASHTOKEN_ACTIVATION_HEIGHT, CHIPNET_START_BLOCK, KEY_LAST_INDEXED,
|
|
};
|
|
use anyhow::{bail, Context, Result};
|
|
|
|
pub async fn update_mempool(db: &DB, electrum: Arc<Mutex<Client>>) -> Result<()> {
|
|
let our_mempool_txs: HashSet<Txid> =
|
|
db::cauldron::mempool::load_mempool(&db.cauldron_w).await?;
|
|
|
|
let electrum_clone = electrum.clone();
|
|
let (cauldron_txs, oracle_txs) = tokio::task::spawn_blocking(move || {
|
|
electrum_fetch_mempool(&electrum_clone.lock().unwrap())
|
|
})
|
|
.await??;
|
|
|
|
let txs_to_delete: Vec<Txid> = our_mempool_txs.difference(&cauldron_txs).cloned().collect();
|
|
let txs_to_add: Vec<&Txid> = cauldron_txs.difference(&our_mempool_txs).collect();
|
|
|
|
let electrum_for_fetch = electrum.clone();
|
|
let txids_to_fetch: Vec<Txid> = txs_to_add.iter().map(|t| **t).collect();
|
|
let txs_to_add: Vec<Transaction> = tokio::task::spawn_blocking(move || {
|
|
txids_to_fetch
|
|
.iter()
|
|
.filter_map(
|
|
|txid| match electrum_get_tx(&electrum_for_fetch.lock().unwrap(), txid) {
|
|
Ok(tx) => Some(tx),
|
|
Err(e) => {
|
|
info!("Failed to get mempool tx {txid}: {e}");
|
|
None
|
|
}
|
|
},
|
|
)
|
|
.collect()
|
|
})
|
|
.await?;
|
|
|
|
let txs_to_add = ttor_sorted_kahn(txs_to_add);
|
|
|
|
{
|
|
let mut db_tx = db.cauldron_w.begin().await?;
|
|
|
|
db::cauldron::mempool::delete_mempool_txs(&mut db_tx, txs_to_delete.iter()).await?;
|
|
|
|
let mut all_cauldrons = vec![];
|
|
let current_timestamp = time_now() as u64;
|
|
|
|
for btx in txs_to_add {
|
|
let cauldrons: Vec<ParsedContract> = parse_cauldrons_from_tx(&btx);
|
|
|
|
if cauldrons.is_empty() {
|
|
info!("ignoring non-cauldron tx {}", btx.txid());
|
|
continue;
|
|
}
|
|
let txid = btx.txid();
|
|
db::cauldron::tx::insert_mempool_tx(&mut db_tx, &txid, current_timestamp).await?;
|
|
info!(
|
|
"mempool add {} with {} cauldrons",
|
|
txid.to_hex(),
|
|
cauldrons.len()
|
|
);
|
|
|
|
insert_utxo_funding(&mut db_tx, &cauldrons, &txid, false).await?;
|
|
insert_utxo_spending(&mut db_tx, &cauldrons, &txid, false).await?;
|
|
insert_user_action(&mut db_tx, &cauldrons, &btx, true).await?;
|
|
|
|
all_cauldrons.extend(cauldrons);
|
|
}
|
|
|
|
db::cauldron::pool::update_pool_history(
|
|
&mut db_tx,
|
|
all_cauldrons,
|
|
None,
|
|
Some(current_timestamp),
|
|
)
|
|
.await?;
|
|
db_tx.commit().await?;
|
|
}
|
|
|
|
// oracle updates
|
|
let oracle_electrum = electrum.clone();
|
|
|
|
// filter entries we have
|
|
let mut oracle_to_add = Vec::new();
|
|
for txid in oracle_txs {
|
|
if !db::oracle::has_entry(&db.oracle_w, &txid).await? {
|
|
oracle_to_add.push(txid);
|
|
}
|
|
}
|
|
|
|
let txs_to_add: Vec<Transaction> = tokio::task::spawn_blocking(move || {
|
|
oracle_to_add
|
|
.into_iter()
|
|
.filter_map(
|
|
|txid| match electrum_get_tx(&oracle_electrum.lock().unwrap(), &txid) {
|
|
Ok(tx) => Some(tx),
|
|
Err(e) => {
|
|
info!("Failed to get mempool tx {txid}: {e}");
|
|
None
|
|
}
|
|
},
|
|
)
|
|
.collect()
|
|
})
|
|
.await?;
|
|
|
|
index_oracle(&db.oracle_w, &txs_to_add, &BlockHash::all_zeros()).await?;
|
|
|
|
Ok(())
|
|
}
|
|
|
|
pub async fn index_blocks(
|
|
chain: Arc<Mutex<Chain>>,
|
|
db: DB,
|
|
client: Arc<Mutex<Client>>,
|
|
bcmr_enabled: bool,
|
|
network: Option<Network>,
|
|
ibd_state: Option<Arc<IbdState>>,
|
|
) -> Result<BlockHash> {
|
|
let electrum_clone = client.clone();
|
|
let (tip_header, tip_height) =
|
|
tokio::task::spawn_blocking(move || electrum_get_tip(&electrum_clone.lock().unwrap()))
|
|
.await??;
|
|
|
|
// Update target height for IBD progress tracking
|
|
if let Some(ref state) = ibd_state {
|
|
state.target_height.store(tip_height, Ordering::Relaxed);
|
|
}
|
|
|
|
// 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, mut block_recv) = tokio::sync::mpsc::channel::<Option<(u64, u64, Block)>>(10);
|
|
|
|
// Update header chain (and undo any blocks that may have reorged away)
|
|
{
|
|
let needs_update = {
|
|
let chain_guard = chain.lock().unwrap();
|
|
tip_header.block_hash() != chain_guard.tip_hash()
|
|
};
|
|
if needs_update {
|
|
{
|
|
let chain_guard = chain.lock().unwrap();
|
|
debug!(
|
|
"Updating header chain from {} to {}",
|
|
chain_guard.tip_hash().to_hex(),
|
|
tip_header.block_hash().to_hex()
|
|
);
|
|
}
|
|
let client_for_headers = client.clone();
|
|
let chain_for_headers = chain.clone();
|
|
let tip_hash = tip_header.block_hash();
|
|
let new_headers = tokio::task::spawn_blocking(move || {
|
|
let chain_guard = chain_for_headers.lock().unwrap();
|
|
get_new_headers(&client_for_headers.lock().unwrap(), &chain_guard, &tip_hash)
|
|
})
|
|
.await??;
|
|
debug!("Storing headers");
|
|
for chunk in new_headers.chunks(100000) {
|
|
let mut db_tx = db.cauldron_w.begin().await?;
|
|
store_headers(&mut db_tx, chunk).await?;
|
|
db_tx.commit().await?;
|
|
}
|
|
// Run chain update in spawn_blocking because StoreBlockUndoer::undo_block
|
|
// uses Handle::current().block_on() which panics on Tokio worker threads.
|
|
let chain_for_update = chain.clone();
|
|
let undoer_db = db.clone();
|
|
tokio::task::spawn_blocking(move || {
|
|
let chain_guard = chain_for_update.lock().unwrap();
|
|
let undoer = StoreBlockUndoer::new(undoer_db)?;
|
|
chain_guard.update(undoer, new_headers, None)
|
|
})
|
|
.await??;
|
|
debug!("Header update done");
|
|
}
|
|
}
|
|
|
|
let db_cpy = db.clone();
|
|
let chain_cpy = chain.clone();
|
|
let client_cpy = client.clone();
|
|
|
|
tokio::task::spawn_blocking(move || {
|
|
let db = db_cpy;
|
|
let handle = tokio::runtime::Handle::current();
|
|
|
|
let last_indexed = handle
|
|
.block_on(config_get(&db.cauldron_r, KEY_LAST_INDEXED))
|
|
.unwrap();
|
|
let mut last_indexed = if let Some(last) = last_indexed {
|
|
BlockHash::from_hex(&last).unwrap()
|
|
} else {
|
|
start_block_hash
|
|
};
|
|
|
|
// Check if last indexed has been orphaned
|
|
loop {
|
|
if chain_cpy.lock().unwrap().contains(&last_indexed) {
|
|
break;
|
|
}
|
|
info!(
|
|
"Last indexed block ({}) has been orphaned",
|
|
last_indexed.to_hex()
|
|
);
|
|
let header = handle
|
|
.block_on(db_get_header(&db.cauldron_r, &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.blocking_send(None) {
|
|
warn!("Failed to send EOL to block reader: {e}");
|
|
}
|
|
break;
|
|
}
|
|
|
|
let next_height = chain_cpy
|
|
.lock()
|
|
.unwrap()
|
|
.get_block_height(&last_indexed)
|
|
.expect("last_indexed height not found in main chain")
|
|
+ 1;
|
|
|
|
let res = client_cpy
|
|
.lock()
|
|
.unwrap()
|
|
.raw_call("blockchain.block.get", vec![Param::U32(next_height as u32)])
|
|
.unwrap_or_else(|e| panic!("Failed to fetch block at height {next_height} from electrum: {e}"));
|
|
|
|
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.blocking_send(Some((
|
|
next_height,
|
|
chain_cpy.lock().unwrap().get_mtp(next_height).unwrap(),
|
|
block,
|
|
))) {
|
|
warn!("Failed to send block to reader: {e}");
|
|
break;
|
|
}
|
|
|
|
last_indexed = block_hash;
|
|
}
|
|
});
|
|
|
|
loop {
|
|
let recv_result = tokio::time::timeout(Duration::from_secs(1), block_recv.recv()).await;
|
|
let (block_height, mtp, block) = match recv_result {
|
|
Ok(Some(Some(data))) => data,
|
|
Ok(Some(None)) | Ok(None) => break, // End of blocks signal or channel closed
|
|
Err(_) => {
|
|
// Timeout
|
|
if shutdown_requested() {
|
|
info!("Shutdown requested, exiting block indexing");
|
|
break;
|
|
}
|
|
continue;
|
|
}
|
|
};
|
|
|
|
let mut db_tx = db.cauldron_w.begin().await?;
|
|
|
|
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(&mut db_tx, &txid, &blockhash, mtp as i64)
|
|
.await
|
|
.context("inserting block tx")?;
|
|
insert_utxo_funding(&mut db_tx, &cauldrons, &txid, true)
|
|
.await
|
|
.context("inserting funding utxos")?;
|
|
insert_utxo_spending(&mut db_tx, &cauldrons, &txid, true)
|
|
.await
|
|
.context("inserting spending utxos")?;
|
|
insert_user_action(&mut db_tx, &cauldrons, tx, true)
|
|
.await
|
|
.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(&mut db_tx, all_cauldrons, Some(mtp), None).await?;
|
|
// config_set must be inside the cauldron transaction, as it tracks what the
|
|
// last successful block indexed was. It must commit atomically with block data.
|
|
config_set(&mut *db_tx, KEY_LAST_INDEXED, &blockhash.to_hex()).await;
|
|
|
|
// crc20
|
|
index_crc20(&db.crc20_w, &sorted_txs).await?;
|
|
|
|
// oracle updates
|
|
index_oracle(&db.oracle_w, &sorted_txs, &blockhash).await?;
|
|
|
|
let autheader_updates = if bcmr_enabled {
|
|
let updates = index_bcmr(&db.bcmr_w, &blockhash, sorted_txs).await?;
|
|
updates as i64
|
|
} else {
|
|
-1
|
|
};
|
|
|
|
db_tx.commit().await?;
|
|
|
|
// Update IBD progress
|
|
if let Some(ref state) = ibd_state {
|
|
state.current_height.store(block_height, Ordering::Relaxed);
|
|
}
|
|
|
|
info!(
|
|
"Indexed {}; mtp: {}, height {}, {} trades, {} autheader updates.",
|
|
blockhash.to_hex(),
|
|
mtp,
|
|
block_height,
|
|
total_cauldrons,
|
|
autheader_updates,
|
|
);
|
|
}
|
|
Ok(tip_header.block_hash())
|
|
}
|