riftenlabs-indexer/src/main.rs

418 lines
14 KiB
Rust
Raw Normal View History

2024-02-14 11:11:16 +01:00
// Copyright (C) 2024 Riften Labs AS
2024-02-05 16:32:42 +01:00
//
// 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 anyhow::{bail, Result};
use bcmr::wellknowndowloader::WellKnownDownloader;
use bitcoincash::{consensus::deserialize, Block, BlockHash};
2024-10-21 12:48:47 +02:00
use crc20::crc20fetcher::CRC20Fetcher;
use db::{DBPool, DB};
use electrum::electrum_get_tip;
use electrum_client_netagnostic::{Client, ElectrumApi, Param};
use log::{error, info, warn};
use rocket::{launch, routes};
2024-01-19 14:38:31 +01:00
use rocket_cors::{AllowedHeaders, AllowedOrigins};
use rpc::ResponseCache;
use rusqlite::OpenFlags;
use std::sync::atomic::AtomicBool;
2023-11-29 12:25:04 +01:00
use std::{
backtrace::Backtrace,
collections::HashMap,
2023-11-29 12:25:04 +01:00
panic,
path::Path,
process,
sync::{Arc, Mutex},
2025-02-14 18:35:12 +01:00
thread::{self, sleep},
time::{Duration, Instant},
2023-11-29 12:25:04 +01:00
};
2024-02-14 11:11:16 +01:00
use stderrlog::LogLevelNum;
2023-11-29 12:25:04 +01:00
use crate::bcmr::bcmrdownloader::BCMRDownloader;
2025-08-15 18:07:15 +02:00
use crate::db::cauldron::config::check_db_version;
use crate::db::cauldron::header::load_all_headers;
2025-09-05 13:53:10 +00:00
use crate::db::cauldron::tokenlist::db_utils::create_cached_token_metrics_table;
use crate::db::cauldron::tokenlist::metrics_cache::spawn_token_metrics_updater;
use crate::index::{index_blocks, update_mempool};
2023-11-29 12:25:04 +01:00
#[macro_use]
extern crate configure_me;
include_config!();
2024-02-14 11:11:16 +01:00
// The block where first cauldron contract was deployed. (Block 799870)
#[allow(dead_code)]
2024-02-14 11:11:16 +01:00
const RIFTEN_LABS_GENESIS_BLOCK: &str =
"000000000000000000ed24c811077f7268a21ecf25cb437655aaba33d8ff4997";
// Start parsing for BCMR data from this height
const CASHTOKEN_ACTIVATION_HEIGHT: &str =
"000000000000000002b678c471841c3e404ec7ae9ca9c32026fe27eb6e3a1ed1";
2023-11-29 12:25:04 +01:00
// Last indexed block height.
const KEY_LAST_INDEXED: &str = "last_indexed";
mod bcmr;
2024-04-03 21:40:33 +02:00
mod cashaddr;
2024-02-14 11:11:16 +01:00
mod chain;
2024-10-21 12:48:47 +02:00
mod crc20;
2023-11-29 12:25:04 +01:00
mod db;
2024-11-13 11:16:24 +01:00
mod def;
2024-02-14 11:11:16 +01:00
mod electrum;
mod index;
2024-03-04 16:40:50 +01:00
mod rpc;
2024-04-03 09:49:44 +02:00
mod timeutil;
2024-10-31 10:55:29 +00:00
mod utiltest;
2024-10-21 12:48:47 +02:00
mod utiltoken;
mod utiltx;
2023-11-29 12:25:04 +01:00
fn set_panic_hook() {
panic::set_hook(Box::new(|panic_info| {
2024-03-04 16:40:50 +01:00
error!("A thread panicked, terminating the program.");
if let Some(error) = panic_info.payload().downcast_ref::<anyhow::Error>() {
2025-07-16 09:31:14 +02:00
error!("Panic occurred: {error:?}");
2024-03-04 16:40:50 +01:00
error!("Anyhow backtrace:\n{}", error.backtrace());
let mut source = error.source();
while let Some(cause) = source {
2025-07-16 09:31:14 +02:00
error!("Caused by: {cause:?}");
2024-03-04 16:40:50 +01:00
source = cause.source();
}
} else if let Some(message) = panic_info.payload().downcast_ref::<&str>() {
2025-07-16 09:31:14 +02:00
error!("Panic occurred: {message}");
2024-10-21 12:48:47 +02:00
} else if let Some(message) = panic_info.payload().downcast_ref::<String>() {
2025-07-16 09:31:14 +02:00
error!("Panic occurred: {message}");
2024-02-14 11:11:16 +01:00
} else {
2025-07-16 09:31:14 +02:00
error!("Panic info: {panic_info:?}");
2023-11-29 12:25:04 +01:00
}
2024-03-04 16:40:50 +01:00
let backtrace = Backtrace::capture();
2025-07-16 09:31:14 +02:00
error!("Backtrace (if RUST_BACKTRACE=1):\n{backtrace}");
2023-11-29 12:25:04 +01:00
process::exit(1);
}));
}
fn db_initialize(c: &mut rusqlite::Connection) -> Result<(), rusqlite::Error> {
2025-02-14 18:35:12 +01:00
// Set busy timeout for normal queries.
c.busy_timeout(Duration::from_secs(30))?;
c.execute("PRAGMA foreign_keys=1;", [])?;
// PRAGMA journal_mode may not honor busy_timeout, so retry manually.
let start = Instant::now();
let timeout = Duration::from_secs(30);
loop {
match c.execute_batch("PRAGMA journal_mode=WAL;") {
Ok(_) => break,
Err(rusqlite::Error::SqliteFailure(err, _))
if err.code == rusqlite::ErrorCode::DatabaseBusy =>
{
if start.elapsed() >= timeout {
return Err(rusqlite::Error::SqliteFailure(
err,
Some("database locked after retries".into()),
));
}
sleep(Duration::from_millis(100));
}
Err(e) => return Err(e),
}
}
Ok(())
}
fn db_sanity_check(c: &rusqlite::Connection) -> rusqlite::Result<()> {
// Verify that foreign_keys are enabled
let mut stmt = c.prepare("PRAGMA foreign_keys;")?;
let foreign_keys_enabled: i32 = stmt.query_row([], |row| row.get(0))?;
if foreign_keys_enabled != 1 {
panic!(
2025-07-16 09:31:14 +02:00
"Foreign key enforcement is not enabled for the connection. (result: {foreign_keys_enabled})"
);
}
Ok(())
}
fn start_program(
config: Config,
) -> Result<(DB, BCMRDownloader, WellKnownDownloader, CRC20Fetcher)> {
let create_db_pool = |db_path| -> (bool, DBPool, DBPool) {
let db_exists = Path::new(db_path).exists();
2025-07-16 09:31:14 +02:00
info!("Initializing connection to {db_path}");
2025-02-14 18:35:12 +01:00
let write_manager = r2d2_sqlite::SqliteConnectionManager::file(db_path)
.with_flags(OpenFlags::SQLITE_OPEN_READ_WRITE | OpenFlags::SQLITE_OPEN_CREATE)
.with_init(db_initialize);
2024-03-14 12:14:13 +01:00
let write_pool =
Arc::new(r2d2::Pool::new(write_manager).expect("Failed to initialize database"));
let read_manager = r2d2_sqlite::SqliteConnectionManager::file(db_path)
.with_flags(OpenFlags::SQLITE_OPEN_READ_ONLY)
.with_init(db_initialize);
let read_pool =
Arc::new(r2d2::Pool::new(read_manager).expect("Failed to initialize database"));
db_sanity_check(&read_pool.get().unwrap()).expect("db sanity check failed");
db_sanity_check(&write_pool.get().unwrap()).expect("db sanity check failed");
(db_exists, write_pool, read_pool)
};
let (db_exists, cauldron_db_write, cauldron_db_read) = create_db_pool("cauldron.db");
if !db_exists {
db::cauldron::prepare_tables(
&cauldron_db_write
.get()
.expect("failed to get sqlite connection"),
);
2025-08-15 18:07:15 +02:00
} else {
// Check database version for existing database
check_db_version(&cauldron_db_read.get().unwrap()).expect("Database version check failed");
}
let (db_exists, bcmr_db_write, bcmr_db_read) = create_db_pool("bcmr.db");
2023-11-29 12:25:04 +01:00
if !db_exists {
db::bcmr::prepare_tables(
&bcmr_db_write
.get()
.expect("failed to create sql connection"),
);
2023-11-29 12:25:04 +01:00
}
2024-10-21 12:48:47 +02:00
let (db_exists, crc20_db_write, crc20_db_read) = create_db_pool("crc20.db");
if !db_exists {
db::crc20::prepare_tables(
&crc20_db_write
.get()
.expect("failed to create sqlite crc20 connection"),
);
}
let (db_exists, oracle_db_write, oracle_db_read) = create_db_pool("oracle.db");
if !db_exists {
db::oracle::prepare_tables(
&oracle_db_write
.get()
.expect("failed to create sqlite oracle connection"),
);
}
// Create a shared flag for indexing status
let indexing_in_progress = Arc::new(AtomicBool::new(false));
2024-02-14 11:11:16 +01:00
let client = Arc::new(Mutex::new(
match Client::new(&format!("tcp://{}", config.rostrum_addr)) {
Ok(server) => server,
Err(e) => {
error!(
"Failed to connect to {}: {}. See --help for setting a different server.",
config.rostrum_addr, e
);
bail!(e)
}
},
2024-02-14 11:11:16 +01:00
));
let genesis = client
.lock()
.unwrap()
.raw_call("blockchain.block.get", vec![Param::U32(0)])
.unwrap();
let genesis: Block = deserialize(&hex::decode(genesis.as_str().unwrap()).unwrap()).unwrap();
let chain = Arc::new(Mutex::new(chain::Chain::new(genesis.header)));
// initialize insert sequence for pool history
db::cauldron::pool::initialize_seq(&cauldron_db_read.get().unwrap());
2024-02-14 11:11:16 +01:00
info!("Loading block headers...");
let all_headers = load_all_headers(&cauldron_db_read.get().unwrap()).unwrap();
2024-02-14 11:11:16 +01:00
info!("Initializing {} headers...", all_headers.len());
chain.lock().unwrap().load(all_headers).unwrap();
info!("Headers loaded.");
let db = DB {
cauldron_w: cauldron_db_write,
cauldron_r: cauldron_db_read,
bcmr_w: bcmr_db_write,
bcmr_r: bcmr_db_read,
2024-10-21 12:48:47 +02:00
crc20_w: crc20_db_write,
crc20_r: crc20_db_read,
oracle_w: oracle_db_write,
oracle_r: oracle_db_read,
};
let db_cpy = db.clone();
2024-10-21 12:48:47 +02:00
let mut crc20fetcher = CRC20Fetcher::new();
crc20fetcher.start(
db.crc20_w.clone(),
client.clone(),
indexing_in_progress.clone(),
)?;
2023-11-29 12:25:04 +01:00
{
// clear oracle mempool
let conn = db.oracle_w.get().unwrap();
db::oracle::clear_mempool(&conn).unwrap();
}
let indexing_in_progress_clone = indexing_in_progress.clone();
2023-11-29 12:25:04 +01:00
thread::spawn(move || {
let db = db_cpy;
2024-02-14 11:11:16 +01:00
let mut tip: BlockHash = loop {
indexing_in_progress_clone.store(true, std::sync::atomic::Ordering::Relaxed);
break match index_blocks(chain.clone(), db.clone(), client.clone(), true) {
Ok(tip) => tip,
Err(e) => {
if e.to_string().contains("database is locked") {
2025-07-16 09:31:14 +02:00
warn!("initial index error, trying again: {e}");
continue;
}
panic!("Initial index failed: {}\n {}", e, e.backtrace());
}
};
2024-02-14 11:11:16 +01:00
};
indexing_in_progress_clone.store(false, std::sync::atomic::Ordering::Relaxed);
2023-11-29 12:25:04 +01:00
loop {
let new_tip = match electrum_get_tip(&client.lock().unwrap()) {
Ok(t) => t.0.block_hash(),
Err(e) => {
2025-07-16 09:31:14 +02:00
warn!("Failed to get block chain tip from electrum: {e}");
thread::sleep(Duration::from_secs(5));
continue;
}
};
2023-11-29 12:25:04 +01:00
if new_tip != tip {
indexing_in_progress_clone.store(true, std::sync::atomic::Ordering::Relaxed);
tip = match index_blocks(chain.clone(), db.clone(), client.clone(), true) {
Ok(t) => t,
Err(e) => {
warn!("Indexing block failed: {} {}", e, e.backtrace());
tip
}
};
indexing_in_progress_clone.store(false, std::sync::atomic::Ordering::Relaxed);
2023-11-29 12:25:04 +01:00
}
if let Err(e) =
update_mempool(db.cauldron_w.clone(), db.oracle_w.clone(), client.clone())
{
2025-07-16 09:31:14 +02:00
error!("Failed to update mempool: {e}");
2024-02-16 09:11:40 +01:00
}
thread::sleep(Duration::from_secs(5));
2023-11-29 12:25:04 +01:00
}
});
let mut bcmrdownloader = BCMRDownloader::new(db.bcmr_w.clone());
bcmrdownloader.start()?;
let mut wellknowndownloader = WellKnownDownloader::new(db.bcmr_w.clone());
wellknowndownloader.start()?;
2024-10-21 12:48:47 +02:00
Ok((db, bcmrdownloader, wellknowndownloader, crc20fetcher))
2024-02-14 11:11:16 +01:00
}
#[launch]
fn launch() -> _ {
stderrlog::new()
.verbosity(LogLevelNum::Info)
2024-02-14 11:11:16 +01:00
.init()
.unwrap();
set_panic_hook();
let (config, _extra) =
Config::including_optional_config_files(std::iter::empty::<std::ffi::OsString>())
.unwrap_or_exit();
let (dbpool, bcmrdownloader, wellknowndownloader, crc20fetcher) = match start_program(config) {
2024-02-14 11:11:16 +01:00
Ok(db) => db,
Err(e) => {
let backtrace = Backtrace::capture();
2025-07-16 09:31:14 +02:00
error!("Backtrace (if RUST_BACKTRACE=1):\n{backtrace}");
error!("Error: {e}");
2024-02-14 11:11:16 +01:00
panic!("Failed at program startup")
}
};
2024-01-19 14:38:31 +01:00
let allowed_origins = AllowedOrigins::all();
let cors = rocket_cors::CorsOptions {
allowed_origins,
allowed_methods: vec![rocket::http::Method::Get]
.into_iter()
.map(From::from)
.collect(),
allowed_headers: AllowedHeaders::some(&["Authorization", "Accept"]),
allow_credentials: true,
..Default::default()
}
.to_cors()
.unwrap();
let response_cache: ResponseCache = Arc::new(Mutex::new(HashMap::default()));
2025-09-05 13:53:10 +00:00
{
let conn = dbpool.cauldron_w.get().expect("get write conn");
// Ensure the table exists on both fresh and existing DBs
create_cached_token_metrics_table(&conn).expect("ensure cached_token_metrics exists");
}
spawn_token_metrics_updater(dbpool.clone());
2024-01-05 14:21:50 +01:00
rocket::build()
2024-03-14 12:14:13 +01:00
.manage(dbpool)
.manage(response_cache)
// give rocket ownership of downloader to ensure thread isn't dropped
.manage(bcmrdownloader)
.manage(wellknowndownloader)
2024-10-21 12:48:47 +02:00
.manage(crc20fetcher)
2024-03-04 16:40:50 +01:00
.mount(
"/cauldron/",
2024-03-18 12:34:46 +01:00
routes![
2024-04-03 09:49:44 +02:00
rpc::tvl::deprecated_tvl,
rpc::tvl::valuelocked_token,
rpc::tvl::valuelocked_all,
rpc::volume::volume_all,
rpc::volume::volume_token,
rpc::tokens::list_by_volume,
2024-10-31 10:55:29 +00:00
rpc::tokens::search_by_volume,
2025-09-05 13:53:10 +00:00
rpc::tokens::search_cached,
rpc::tokens::list_cached,
rpc::tokens::list_cached_by_ids,
2024-04-03 11:39:49 +02:00
rpc::price::price_history,
2025-06-17 13:22:46 +00:00
rpc::candlesticks::price_candlesticks,
2024-04-03 11:39:49 +02:00
rpc::price::price_current,
2024-10-14 11:18:50 +00:00
rpc::price::price_at,
2024-04-03 21:40:33 +02:00
rpc::pool::list_pools_by_apy,
rpc::pool::list_active_pools,
2024-11-13 11:16:24 +01:00
rpc::pool::pool_history,
rpc::apy::aggregate_apy,
2024-04-03 11:39:49 +02:00
rpc::contract::contract_count_token,
rpc::contract::contract_count_all,
rpc::contract::contract_volume,
2024-09-25 11:56:55 +02:00
rpc::user::unique_addresses,
rpc::tx::tx_latest,
2024-03-18 12:34:46 +01:00
],
2024-03-04 16:40:50 +01:00
)
.mount(
"/bcmr",
routes![rpc::bcmr::token_bcmr, rpc::bcmr::token_bcmr_all],
)
.mount(
"/oracle",
routes![
rpc::oracle::oracle_get_closest,
2025-06-17 14:04:26 +00:00
rpc::oracle::oracle_get_range,
rpc::oracle::oracle_get_history
],
)
2024-01-19 14:38:31 +01:00
.attach(cors)
2023-11-29 12:25:04 +01:00
}