Skip to main content

degenbot_cli_core/
pool.rs

1//! The `pool` command arms (ADR-051 D1;).
2//!
3//! Ports `cli/pool.py` arm for arm:
4//!
5//! - `pool update` — the chunk-loop hand-off to
6//!   [`run_pool_update`](degenbot_pool_updater::run_pool_update), with the
7//!   `--chunk` / `--to-block` / `--verify-chunk` / `--verify-all` /
8//!   `--verify-all-interval` flags 1:1. Progress stays a no-op sink here
9//!   (ADR-051 D9: the bin paints it); a cooperative cancel returns a friendly
10//!   cancelled report (exit 0), matching the Python `RuntimeError` guard.
11//! - `pool verify` — the read-only, ad-hoc sibling: fetch the COMMITTED
12//!   liquidity map, compare against on-chain truth at `--block`, and render
13//!   GREEN / the named divergence list.
14
15use 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
36/// The RPC retry budget the verify arm's spot reads use.
37const RPC_MAX_RETRIES: u32 = 5;
38
39/// The pool family `--family` selects.
40#[derive(Debug, Clone, Copy, PartialEq, Eq)]
41pub enum PoolFamily {
42    /// V3 (`ticks()`/`tickBitmap()`).
43    V3,
44    /// V4 (`PoolManager` `extsload`).
45    V4,
46}
47
48impl PoolFamily {
49    /// The wire spelling (`v3` / `v4`).
50    #[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/// The `pool` command group.
60#[derive(Debug, Clone, PartialEq, Eq)]
61pub enum PoolCommand {
62    /// Advance every active exchange's liquidity state to `to_block`.
63    Update {
64        /// Max blocks per chunk before committing.
65        chunk_size: u64,
66        /// The raw `--to-block` identifier (tag / tag:offset / integer).
67        to_block: String,
68        /// Run the pre-commit per-chunk on-chain-truth gate.
69        verify_chunk: bool,
70        /// Run the pre-commit market-wide verification at the interval + completion.
71        verify_all: bool,
72        /// The block interval for the `--verify-all` gate.
73        verify_all_interval: u64,
74    },
75    /// Verify a committed pool's liquidity map against on-chain truth.
76    Verify {
77        /// The HTTP RPC endpoint.
78        rpc_url: String,
79        /// The chain the pool lives on.
80        chain_id: i64,
81        /// The block number to read on-chain truth at.
82        block_number: u64,
83        /// V3 pool address, or V4 `PoolId` (bytes32 hex).
84        pool: String,
85        /// The pool family.
86        family: PoolFamily,
87        /// (V4 only) the `PoolManager` singleton.
88        pool_manager: Option<String>,
89    },
90}
91
92impl PoolCommand {
93    /// The arm's confirmation policy: neither pool arm prompts.
94    #[must_use]
95    pub const fn prompt_plan(&self, _ctx: &CliContext<'_>) -> PromptPlan {
96        PromptPlan::None
97    }
98}
99
100/// The resolved verify target for a family.
101enum VerifyTarget {
102    /// A V3 pool contract address.
103    V3(Address),
104    /// A V4 `PoolManager` + `PoolId`.
105    V4 { manager: Address, pool_id: B256 },
106}
107
108/// Execute a `pool` command.
109///
110/// # Errors
111///
112/// [`CliError::InvalidBlockTag`] / [`CliError::InvalidArgument`] /
113/// [`CliError::InvalidAddress`] for malformed inputs, [`CliError::Config`] for
114/// an unresolved driver-domain value, and [`CliError::PoolUpdate`] for a core
115/// failure.
116pub(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
157/// `pool update`.
158fn 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    // Self-serve registration: every supported exchange pair not found in the
171    // DB registers inactive, so the update never depends on prior CREATEs.
172    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    // The pre-run resume snapshot, for the failure line's requested-range
182    // bound: it mirrors the core's own `initial_start_block` (the minimum
183    // `last_update_block + 1` across the active exchanges). Read-only, so a
184    // failed diagnostic read must not fail the run — it degrades to the
185    // requested target.
186    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            // Post-failure resume snapshot: the committed per-chunk progress
211            // (a chunk-boundary commit advances the exchange cursor), read
212            // only — the outstanding work a rerun resumes from. `None` when
213            // the snapshot read itself failed; the rendering says so.
214            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
227/// The read-only per-exchange cursor snapshot for `chain_id` — the resume
228/// state a mid-run failure leaves committed. Plain `SELECT`, so it writes
229/// nothing; `None` cursors are the never-updated exchanges.
230fn 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
251/// The run's intended start block, mirroring the core's own
252/// `initial_start_block`: the minimum `last_update_block + 1` across the
253/// active exchanges (a never-updated exchange starts at 1). An empty
254/// registration set falls back to the requested target.
255fn 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
267/// `pool verify`.
268fn 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    // The arm verifies a row stored for `chain_id` against the node, so the
284    // provider is bound to that chain: an endpoint serving another chain is
285    // refused before any on-chain comparison reports phantom divergences.
286    let chain = u64::try_from(chain_id).map_err(|_| {
287        CliError::InvalidArgument(format!("chain id {chain_id} is not a valid chain"))
288    })?;
289    // ONE chain-verified provider for the whole arm, on the process-wide
290    // shared runtime — the pattern `run_pool_update` uses for its chunk
291    // loop — NOT a fresh runtime + transport built inside `block::block_on`
292    // per target family (a per-family, per-verification build with its
293    // connection tasks destroyed alongside the ad-hoc runtime is the churn
294    // that intermittently surfaces as alloy's
295    // `TransportErrorKind::BackendGone`). The construction is lazy (the
296    // node is dialed only after the committed row above is resolved), the
297    // chain binding still runs exactly once (`for_chain` refuses a
298    // foreign-chain endpoint before any comparison), and no retry layer is
299    // added: a transport failure surfaces as the typed run failure.
300    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
349/// Drive `fut` on the process-wide shared runtime — the `&'static`
350/// `degenbot_core::runtime::get_runtime()` singleton the updater cores
351/// drive — instead of building a throwaway `Builder` runtime per call.
352/// Same nesting constraint as [`crate::block::block_on`]: the arm must not
353/// run from inside another runtime.
354///
355/// # Errors
356///
357/// [`CliError::RuntimeNested`] when called from inside an existing runtime.
358fn 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
365/// The process-lifetime count of `pool verify` transport builds (the
366/// `for_chain` site above). The lifecycle contract to pin: ONE build per
367/// verify run — zero when the arm fails before the node is needed, and
368/// never one per target family. The companion counter lives at the chunk
369/// loop's own build site (`degenbot-pool-updater`'s
370/// `run_provider_build_count`).
371static 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/// The instrumented transport-build count — a test-visible seam so the
378/// one-build-per-run contract is pinnable offline (see the crate's
379/// `pool_verify_lifecycle` test).
380#[doc(hidden)]
381#[must_use]
382pub fn verify_provider_build_count() -> u64 {
383    VERIFY_PROVIDER_BUILDS.load(Ordering::Relaxed)
384}
385
386/// Reset the transport-build counter to zero, returning the previous value.
387#[doc(hidden)]
388#[must_use]
389pub fn reset_verify_provider_build_count() -> u64 {
390    VERIFY_PROVIDER_BUILDS.swap(0, Ordering::Relaxed)
391}
392
393/// Fetch the pool's committed liquidity map + resolve the verify target.
394fn 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}