1use std::future::Future;
16use std::path::Path;
17use std::str::FromStr;
18use std::sync::atomic::{AtomicU64, Ordering};
19use std::sync::Arc;
20
21use alloy::primitives::{Address, B256};
22use degenbot_core::runtime::get_runtime;
23use degenbot_db::{ComputedLiquidityUpdate, DegenbotDb};
24use degenbot_pool_updater::{
25 run_pool_update, verify_v3_liquidity_map_on_chain, verify_v4_liquidity_map_on_chain,
26 NoProgress, RunError,
27};
28
29use crate::block::{parse_to_block, resolve_to_block};
30use crate::cancel::CancelHandle;
31use crate::context::CliContext;
32use crate::error::{CliError, ExchangeResumeState, PoolUpdateFailure};
33use crate::prompt::{PromptPlan, Prompter};
34use crate::report::PoolReport;
35
36const RPC_MAX_RETRIES: u32 = 5;
38
39#[derive(Debug, Clone, Copy, PartialEq, Eq)]
41pub enum PoolFamily {
42 V3,
44 V4,
46}
47
48impl PoolFamily {
49 #[must_use]
51 pub const fn as_str(self) -> &'static str {
52 match self {
53 Self::V3 => "v3",
54 Self::V4 => "v4",
55 }
56 }
57}
58
59#[derive(Debug, Clone, PartialEq, Eq)]
61pub enum PoolCommand {
62 Update {
64 chunk_size: u64,
66 to_block: String,
68 verify_chunk: bool,
70 verify_all: bool,
72 verify_all_interval: u64,
74 },
75 Verify {
77 rpc_url: String,
79 chain_id: i64,
81 block_number: u64,
83 pool: String,
85 family: PoolFamily,
87 pool_manager: Option<String>,
89 },
90}
91
92impl PoolCommand {
93 #[must_use]
95 pub const fn prompt_plan(&self, _ctx: &CliContext<'_>) -> PromptPlan {
96 PromptPlan::None
97 }
98}
99
100enum VerifyTarget {
102 V3(Address),
104 V4 { manager: Address, pool_id: B256 },
106}
107
108pub(crate) fn execute(
117 command: &PoolCommand,
118 ctx: &CliContext<'_>,
119 _prompter: &dyn Prompter,
120 cancel: &CancelHandle,
121) -> Result<PoolReport, CliError> {
122 match command {
123 PoolCommand::Update {
124 chunk_size,
125 to_block,
126 verify_chunk,
127 verify_all,
128 verify_all_interval,
129 } => update(
130 ctx,
131 cancel,
132 *chunk_size,
133 to_block,
134 *verify_chunk,
135 *verify_all,
136 *verify_all_interval,
137 ),
138 PoolCommand::Verify {
139 rpc_url,
140 chain_id,
141 block_number,
142 pool,
143 family,
144 pool_manager,
145 } => verify(
146 ctx,
147 rpc_url,
148 *chain_id,
149 *block_number,
150 pool,
151 *family,
152 pool_manager.as_deref(),
153 ),
154 }
155}
156
157fn update(
159 ctx: &CliContext<'_>,
160 cancel: &CancelHandle,
161 chunk_size: u64,
162 to_block: &str,
163 verify_chunk: bool,
164 verify_all: bool,
165 verify_all_interval: u64,
166) -> Result<PoolReport, CliError> {
167 let database_path = ctx.database_path()?.value;
168 let chain_id = ctx.chain_id()?.value;
169 let rpc_url = ctx.node_request_uri()?.value;
170 crate::registrations::ensure_supported_registrations(&database_path)?;
173 let resolved = resolve_to_block(parse_to_block(to_block)?, &rpc_url)?;
174 let chain = i64::try_from(chain_id)
175 .map_err(|_| CliError::InvalidArgument(format!("chain id {chain_id} is out of range")))?;
176 let interval = if verify_all {
177 Some(verify_all_interval)
178 } else {
179 None
180 };
181 let cursors = read_exchange_cursors(database_path.as_path(), chain).unwrap_or_default();
187 let from_block = initial_run_start_block(&cursors, resolved);
188 match run_pool_update(
189 &database_path,
190 chain,
191 resolved,
192 chunk_size,
193 &rpc_url,
194 cancel.flag(),
195 Arc::new(NoProgress),
196 verify_chunk,
197 interval,
198 verify_all,
199 ) {
200 Ok(report) => Ok(PoolReport::Updated {
201 chain_id: report.chain_id,
202 from_block: report.from_block,
203 to_block: report.to_block,
204 chunks_committed: report.chunks_committed,
205 total_pools_written: report.total_pools_written,
206 total_liquidity_applies: report.total_liquidity_applies,
207 }),
208 Err(RunError::Cancelled) => Ok(PoolReport::UpdateCancelled { chain_id: chain }),
209 Err(err) => {
210 let resume = read_exchange_cursors(database_path.as_path(), chain).ok();
215 Err(CliError::PoolUpdate(Box::new(PoolUpdateFailure {
216 error: err,
217 rpc_url: rpc_url.clone(),
218 chain_id: chain,
219 from_block,
220 to_block: resolved,
221 resume,
222 })))
223 }
224 }
225}
226
227fn read_exchange_cursors(
231 database_path: &Path,
232 chain_id: i64,
233) -> Result<Vec<ExchangeResumeState>, rusqlite::Error> {
234 let conn = rusqlite::Connection::open_with_flags(
235 database_path,
236 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY,
237 )?;
238 let mut statement = conn.prepare(
239 "SELECT name, last_update_block FROM exchanges \
240 WHERE chain_id = ?1 AND active = 1 ORDER BY name",
241 )?;
242 let rows = statement.query_map([chain_id], |row| {
243 Ok(ExchangeResumeState {
244 name: row.get(0)?,
245 last_update_block: row.get(1)?,
246 })
247 })?;
248 rows.collect()
249}
250
251fn initial_run_start_block(cursors: &[ExchangeResumeState], to_block: Option<u64>) -> u64 {
256 cursors
257 .iter()
258 .map(|row| {
259 row.last_update_block
260 .map_or(1, |block| u64::try_from(block).unwrap_or(0) + 1)
261 })
262 .min()
263 .unwrap_or_else(|| to_block.unwrap_or(1))
264 .max(1)
265}
266
267fn verify(
269 ctx: &CliContext<'_>,
270 rpc_url: &str,
271 chain_id: i64,
272 block_number: u64,
273 pool: &str,
274 family: PoolFamily,
275 pool_manager: Option<&str>,
276) -> Result<PoolReport, CliError> {
277 let database_path = ctx.database_path()?.value;
278 let (computed, target) = {
279 let (db, _state) = DegenbotDb::open(&database_path)?;
280 let conn = db.lock();
281 fetch_verify_state(&conn, chain_id, pool, family, pool_manager)?
282 };
283 let chain = u64::try_from(chain_id).map_err(|_| {
287 CliError::InvalidArgument(format!("chain id {chain_id} is not a valid chain"))
288 })?;
289 let divergences = shared_runtime_block_on(async {
301 note_verify_provider_build();
302 let provider =
303 degenbot_rpc::provider::AlloyProvider::for_chain(rpc_url, chain, RPC_MAX_RETRIES)
304 .await
305 .map_err(|err| CliError::BlockResolution(err.to_string()))?;
306 match target {
307 VerifyTarget::V3(address) => {
308 verify_v3_liquidity_map_on_chain(&provider, address, &computed, block_number)
309 .await
310 .map_err(|err| {
311 CliError::PoolUpdate(Box::new(PoolUpdateFailure {
312 error: err,
313 rpc_url: rpc_url.to_string(),
314 chain_id,
315 from_block: block_number,
316 to_block: Some(block_number),
317 resume: None,
318 }))
319 })
320 }
321 VerifyTarget::V4 { manager, pool_id } => verify_v4_liquidity_map_on_chain(
322 &provider,
323 manager,
324 pool_id,
325 &computed,
326 block_number,
327 )
328 .await
329 .map_err(|err| {
330 CliError::PoolUpdate(Box::new(PoolUpdateFailure {
331 error: err,
332 rpc_url: rpc_url.to_string(),
333 chain_id,
334 from_block: block_number,
335 to_block: Some(block_number),
336 resume: None,
337 }))
338 }),
339 }
340 })??;
341 Ok(PoolReport::Verified {
342 pool: pool.to_string(),
343 family,
344 block_number,
345 divergences,
346 })
347}
348
349fn shared_runtime_block_on<F: Future>(fut: F) -> Result<F::Output, CliError> {
359 if tokio::runtime::Handle::try_current().is_ok() {
360 return Err(CliError::RuntimeNested);
361 }
362 Ok(get_runtime().block_on(fut))
363}
364
365static VERIFY_PROVIDER_BUILDS: AtomicU64 = AtomicU64::new(0);
372
373fn note_verify_provider_build() {
374 VERIFY_PROVIDER_BUILDS.fetch_add(1, Ordering::Relaxed);
375}
376
377#[doc(hidden)]
381#[must_use]
382pub fn verify_provider_build_count() -> u64 {
383 VERIFY_PROVIDER_BUILDS.load(Ordering::Relaxed)
384}
385
386#[doc(hidden)]
388#[must_use]
389pub fn reset_verify_provider_build_count() -> u64 {
390 VERIFY_PROVIDER_BUILDS.swap(0, Ordering::Relaxed)
391}
392
393fn fetch_verify_state(
395 conn: &rusqlite::Connection,
396 chain_id: i64,
397 pool: &str,
398 family: PoolFamily,
399 pool_manager: Option<&str>,
400) -> Result<(ComputedLiquidityUpdate, VerifyTarget), CliError> {
401 match family {
402 PoolFamily::V3 => {
403 let address: Address = pool
404 .parse()
405 .map_err(|_| CliError::InvalidAddress(pool.to_string()))?;
406 let key = address.to_checksum(None);
407 let state = DegenbotDb::fetch_v3_pool_update_state_on_conn(conn, chain_id, &key)?
408 .ok_or_else(|| {
409 CliError::InvalidArgument(format!(
410 "v3 pool {key} not found on chain {chain_id}"
411 ))
412 })?;
413 let (tick_bitmap, tick_data) =
414 DegenbotDb::fetch_v3_liquidity_map_on_conn(conn, state.pool_id)?;
415 Ok((
416 ComputedLiquidityUpdate {
417 pool_id: state.pool_id,
418 tick_spacing: state.tick_spacing,
419 tick_data,
420 tick_bitmap,
421 last_event: None,
422 },
423 VerifyTarget::V3(address),
424 ))
425 }
426 PoolFamily::V4 => {
427 let manager = pool_manager.ok_or_else(|| {
428 CliError::InvalidArgument(
429 "--pool-manager is required for --family v4 (the PoolManager singleton)."
430 .to_string(),
431 )
432 })?;
433 let manager: Address = manager
434 .parse()
435 .map_err(|_| CliError::InvalidAddress(manager.to_string()))?;
436 let pool_id = B256::from_str(pool.strip_prefix("0x").unwrap_or(pool))
437 .map_err(|_| CliError::InvalidAddress(pool.to_string()))?;
438 let state = DegenbotDb::fetch_v4_pool_update_state_on_conn(conn, pool, chain_id)?
439 .ok_or_else(|| {
440 CliError::InvalidArgument(format!(
441 "v4 pool {pool} not found on chain {chain_id}"
442 ))
443 })?;
444 let (tick_bitmap, tick_data) =
445 DegenbotDb::fetch_v4_liquidity_map_on_conn(conn, state.pool_id)?;
446 Ok((
447 ComputedLiquidityUpdate {
448 pool_id: state.pool_id,
449 tick_spacing: state.tick_spacing,
450 tick_data,
451 tick_bitmap,
452 last_event: None,
453 },
454 VerifyTarget::V4 { manager, pool_id },
455 ))
456 }
457 }
458}