use std::future::Future;
use std::path::Path;
use std::str::FromStr;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use alloy::primitives::{Address, B256};
use degenbot_core::runtime::get_runtime;
use degenbot_db::{ComputedLiquidityUpdate, DegenbotDb};
use degenbot_pool_updater::{
run_pool_update, verify_v3_liquidity_map_on_chain, verify_v4_liquidity_map_on_chain,
NoProgress, RunError,
};
use crate::block::{parse_to_block, resolve_to_block};
use crate::cancel::CancelHandle;
use crate::context::CliContext;
use crate::error::{CliError, ExchangeResumeState, PoolUpdateFailure};
use crate::prompt::{PromptPlan, Prompter};
use crate::report::PoolReport;
const RPC_MAX_RETRIES: u32 = 5;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PoolFamily {
V3,
V4,
}
impl PoolFamily {
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::V3 => "v3",
Self::V4 => "v4",
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum PoolCommand {
Update {
chunk_size: u64,
to_block: String,
verify_chunk: bool,
verify_all: bool,
verify_all_interval: u64,
},
Verify {
rpc_url: String,
chain_id: i64,
block_number: u64,
pool: String,
family: PoolFamily,
pool_manager: Option<String>,
},
}
impl PoolCommand {
#[must_use]
pub const fn prompt_plan(&self, _ctx: &CliContext<'_>) -> PromptPlan {
PromptPlan::None
}
}
enum VerifyTarget {
V3(Address),
V4 { manager: Address, pool_id: B256 },
}
pub(crate) fn execute(
command: &PoolCommand,
ctx: &CliContext<'_>,
_prompter: &dyn Prompter,
cancel: &CancelHandle,
) -> Result<PoolReport, CliError> {
match command {
PoolCommand::Update {
chunk_size,
to_block,
verify_chunk,
verify_all,
verify_all_interval,
} => update(
ctx,
cancel,
*chunk_size,
to_block,
*verify_chunk,
*verify_all,
*verify_all_interval,
),
PoolCommand::Verify {
rpc_url,
chain_id,
block_number,
pool,
family,
pool_manager,
} => verify(
ctx,
rpc_url,
*chain_id,
*block_number,
pool,
*family,
pool_manager.as_deref(),
),
}
}
fn update(
ctx: &CliContext<'_>,
cancel: &CancelHandle,
chunk_size: u64,
to_block: &str,
verify_chunk: bool,
verify_all: bool,
verify_all_interval: u64,
) -> Result<PoolReport, CliError> {
let database_path = ctx.database_path()?.value;
let chain_id = ctx.chain_id()?.value;
let rpc_url = ctx.node_request_uri()?.value;
crate::registrations::ensure_supported_registrations(&database_path)?;
let resolved = resolve_to_block(parse_to_block(to_block)?, &rpc_url)?;
let chain = i64::try_from(chain_id)
.map_err(|_| CliError::InvalidArgument(format!("chain id {chain_id} is out of range")))?;
let interval = if verify_all {
Some(verify_all_interval)
} else {
None
};
let cursors = read_exchange_cursors(database_path.as_path(), chain).unwrap_or_default();
let from_block = initial_run_start_block(&cursors, resolved);
let run = match shared_runtime_block_on(async {
degenbot_rpc::provider::AlloyProvider::new(&rpc_url, RPC_MAX_RETRIES).await
}) {
Ok(Ok(provider)) => run_pool_update(
&database_path,
chain,
resolved,
chunk_size,
provider,
cancel.flag(),
Arc::new(NoProgress),
verify_chunk,
interval,
verify_all,
),
Ok(Err(err)) => Err(RunError::from(err)),
Err(cli_err) => return Err(cli_err),
};
match run {
Ok(report) => Ok(PoolReport::Updated {
chain_id: report.chain_id,
from_block: report.from_block,
to_block: report.to_block,
chunks_committed: report.chunks_committed,
total_pools_written: report.total_pools_written,
total_liquidity_applies: report.total_liquidity_applies,
}),
Err(RunError::Cancelled) => Ok(PoolReport::UpdateCancelled { chain_id: chain }),
Err(err) => {
let resume = read_exchange_cursors(database_path.as_path(), chain).ok();
Err(CliError::PoolUpdate(Box::new(PoolUpdateFailure {
error: err,
rpc_url: rpc_url.clone(),
chain_id: chain,
from_block,
to_block: resolved,
resume,
})))
}
}
}
fn read_exchange_cursors(
database_path: &Path,
chain_id: i64,
) -> Result<Vec<ExchangeResumeState>, rusqlite::Error> {
let conn = rusqlite::Connection::open_with_flags(
database_path,
rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY,
)?;
let mut statement = conn.prepare(
"SELECT name, last_update_block FROM exchanges \
WHERE chain_id = ?1 AND active = 1 ORDER BY name",
)?;
let rows = statement.query_map([chain_id], |row| {
Ok(ExchangeResumeState {
name: row.get(0)?,
last_update_block: row.get(1)?,
})
})?;
rows.collect()
}
fn initial_run_start_block(cursors: &[ExchangeResumeState], to_block: Option<u64>) -> u64 {
cursors
.iter()
.map(|row| {
row.last_update_block
.map_or(1, |block| u64::try_from(block).unwrap_or(0) + 1)
})
.min()
.unwrap_or_else(|| to_block.unwrap_or(1))
.max(1)
}
fn verify(
ctx: &CliContext<'_>,
rpc_url: &str,
chain_id: i64,
block_number: u64,
pool: &str,
family: PoolFamily,
pool_manager: Option<&str>,
) -> Result<PoolReport, CliError> {
let database_path = ctx.database_path()?.value;
let (computed, target) = {
let (db, _state) = DegenbotDb::open(&database_path)?;
let conn = db.lock();
fetch_verify_state(&conn, chain_id, pool, family, pool_manager)?
};
let chain = u64::try_from(chain_id).map_err(|_| {
CliError::InvalidArgument(format!("chain id {chain_id} is not a valid chain"))
})?;
let divergences = shared_runtime_block_on(async {
note_verify_provider_build();
let provider =
degenbot_rpc::provider::AlloyProvider::for_chain(rpc_url, chain, RPC_MAX_RETRIES)
.await
.map_err(|err| CliError::BlockResolution(err.to_string()))?;
match target {
VerifyTarget::V3(address) => {
verify_v3_liquidity_map_on_chain(&provider, address, &computed, block_number)
.await
.map_err(|err| {
CliError::PoolUpdate(Box::new(PoolUpdateFailure {
error: err,
rpc_url: rpc_url.to_string(),
chain_id,
from_block: block_number,
to_block: Some(block_number),
resume: None,
}))
})
}
VerifyTarget::V4 { manager, pool_id } => verify_v4_liquidity_map_on_chain(
&provider,
manager,
pool_id,
&computed,
block_number,
)
.await
.map_err(|err| {
CliError::PoolUpdate(Box::new(PoolUpdateFailure {
error: err,
rpc_url: rpc_url.to_string(),
chain_id,
from_block: block_number,
to_block: Some(block_number),
resume: None,
}))
}),
}
})??;
Ok(PoolReport::Verified {
pool: pool.to_string(),
family,
block_number,
divergences,
})
}
pub(crate) fn shared_runtime_block_on<F: Future>(fut: F) -> Result<F::Output, CliError> {
if tokio::runtime::Handle::try_current().is_ok() {
return Err(CliError::RuntimeNested);
}
Ok(get_runtime().block_on(fut))
}
static VERIFY_PROVIDER_BUILDS: AtomicU64 = AtomicU64::new(0);
fn note_verify_provider_build() {
VERIFY_PROVIDER_BUILDS.fetch_add(1, Ordering::Relaxed);
}
#[doc(hidden)]
#[must_use]
pub fn verify_provider_build_count() -> u64 {
VERIFY_PROVIDER_BUILDS.load(Ordering::Relaxed)
}
#[doc(hidden)]
#[must_use]
pub fn reset_verify_provider_build_count() -> u64 {
VERIFY_PROVIDER_BUILDS.swap(0, Ordering::Relaxed)
}
fn fetch_verify_state(
conn: &rusqlite::Connection,
chain_id: i64,
pool: &str,
family: PoolFamily,
pool_manager: Option<&str>,
) -> Result<(ComputedLiquidityUpdate, VerifyTarget), CliError> {
match family {
PoolFamily::V3 => {
let address: Address = pool
.parse()
.map_err(|_| CliError::InvalidAddress(pool.to_string()))?;
let key = address.to_checksum(None);
let state = DegenbotDb::fetch_v3_pool_update_state_on_conn(conn, chain_id, &key)?
.ok_or_else(|| {
CliError::InvalidArgument(format!(
"v3 pool {key} not found on chain {chain_id}"
))
})?;
let (tick_bitmap, tick_data) =
DegenbotDb::fetch_v3_liquidity_map_on_conn(conn, state.pool_id)?;
Ok((
ComputedLiquidityUpdate {
pool_id: state.pool_id,
tick_spacing: state.tick_spacing,
tick_data,
tick_bitmap,
last_event: None,
},
VerifyTarget::V3(address),
))
}
PoolFamily::V4 => {
let manager = pool_manager.ok_or_else(|| {
CliError::InvalidArgument(
"--pool-manager is required for --family v4 (the PoolManager singleton)."
.to_string(),
)
})?;
let manager: Address = manager
.parse()
.map_err(|_| CliError::InvalidAddress(manager.to_string()))?;
let pool_id = B256::from_str(pool.strip_prefix("0x").unwrap_or(pool))
.map_err(|_| CliError::InvalidAddress(pool.to_string()))?;
let state = DegenbotDb::fetch_v4_pool_update_state_on_conn(conn, pool, chain_id)?
.ok_or_else(|| {
CliError::InvalidArgument(format!(
"v4 pool {pool} not found on chain {chain_id}"
))
})?;
let (tick_bitmap, tick_data) =
DegenbotDb::fetch_v4_liquidity_map_on_conn(conn, state.pool_id)?;
Ok((
ComputedLiquidityUpdate {
pool_id: state.pool_id,
tick_spacing: state.tick_spacing,
tick_data,
tick_bitmap,
last_event: None,
},
VerifyTarget::V4 { manager, pool_id },
))
}
}
}