degenbot.runner =============== .. py:module:: degenbot.runner .. autoapi-nested-parse:: Settlement-arbitrage runtime driver (``BotRunner``) companion package. Extracted from ``examples/eth_backrun_v2_v3_v4_rust.py`` / ``eth_backrun_helpers.py``. This is the Python-companion ``stays-python`` cockpit over the Rust-owned engine: it owns config, discovery/registration, result consumption, and dispatch orchestration — never pool/engine state (ADR-003: ``Bot`` is the single Rust state owner; ADR-006: ``Bot`` is the per-chain orchestrator, this package is its deployment cockpit). The package presents one face: the driver cockpit. Public surface (what this module re-exports): - :class:`BotRunner` — the runtime driver facade (the ``start / build_paths / consume / dispatch`` seams). - :class:`ArbitrageConfig` — the unified frozen config (``build``). - The build family (``build_paths`` / ``PathRegistrationPipeline`` / ``ConstructionContext`` / ``resolve_directions``) and the CLI arg parser (:mod:`degenbot.runner.cli`). PRG-5: the bounded crawl shell retired — the crawl is the fleet-hosted intake now. Everything else is private by name (``_consume`` / ``_dispatch`` / ``_render`` / ``identity``) and is imported directly by name from its private module — nothing is smuggled in via the package root. Submodules ---------- .. toctree:: :maxdepth: 1 /autoapi/degenbot/runner/bot_runner/index /autoapi/degenbot/runner/build_paths/index /autoapi/degenbot/runner/cli/index /autoapi/degenbot/runner/config/index /autoapi/degenbot/runner/diag/index /autoapi/degenbot/runner/identity/index Package Contents ---------------- .. py:class:: BotRunner(cfg: degenbot.runner.config.ArbitrageConfig, *, actors: InjectedActors | None = None, install_sigint: bool = True) Orchestrator that collapses the settlement-arbitrage startup ritual behind one facade. Owns the config and the lifecycle; the coordination state itself (the three actors ``bot``/``engine_registry``/``async_w3``, the ``Dispatcher``, the block clock, the sim context, the pipelines) lives on ONE :class:`_SessionState` built in ``start()`` — the runner's same-named attributes are a facade over that owner, not mirrors. The runner is the ONE place that enforces the phase ordering the engine's state machine requires: start(): subscribe → stream snapshots → backfill → verify config (``EngineRegistry.start``, stops at Backfilled, pre-resume) run(): the run ritual (``_run_ritual.RunRitual``) owns the startup ordering — attach consumer → watch → ``resume()`` → registration → main loop — as a state machine; the cross-task fail-fast channel surfaces a fatal registration error. Usage (production):: cfg = ArbitrageConfig.build(live=not dry_run, permutation=args.permutation) async with BotRunner(cfg) as session: await session.run() In production ``run()`` schedules discovery+registration as a background task (through the injected scheduler) and enters the main loop immediately; the state-trim runs on registration completion (inside the scheduled hand-off), not on the main-loop entry path. A fatal verification error still crashes loudly through the cross-task channel. The hot loop keeps only ``engine_registry`` + ``async_w3`` + dispatcher once trimmed — the Python pool/token caches are scaffolding once the Rust engine owns canonical state. Testability seams (mirrors ``EngineRegistry``'s ``engine=`` seam): ``bot``, ``engine_registry``, ``async_w3``, ``snapshots``, ``path_builder``, ``consumer``, and the registration ``scheduler`` are injectable. When injected, ``start()``/``run()`` orchestrate the fakes and the phase ordering is verifiable offline; when ``None`` (production), the actors are built from ``cfg`` and the real module functions are called. .. py:attribute:: cfg .. py:property:: bot :type: degenbot.Bot | None The session's Python-companion bot (``None`` once the trim dropped it). Before ``start()``: the injected seam (a write is an injection). .. py:property:: engine_registry :type: degenbot.arbitrage.engine_registry.EngineRegistry | None the injected seam. :type: The session's engine registry. Before ``start()`` .. py:property:: async_w3 :type: degenbot.provider.AsyncAlloyProvider | None the injected seam. :type: The session's dispatch-path provider. Before ``start()`` .. py:property:: dispatcher :type: degenbot.dispatch.Dispatcher | None ``None`` (no session yet). :type: The session's dispatcher. Before ``start()`` .. py:property:: session :type: _SessionState | None The session owner (``None`` only before ``start()``). .. py:property:: session_watch :type: degenbot.runner._session_watch.SessionWatch The session watch the ritual attaches the run's members to. .. py:property:: consumer :type: Any The injected consumer (``None`` = the production block loop). .. py:property:: readiness :type: degenbot.strategy.StrategyReadinessView | None The readiness view resolved once at the ``start()`` boundary. .. py:property:: settlement_active :type: bool The resolved settlement-arm disposition (read by the ritual). .. py:property:: path_builder :type: Any The injected path builder (``None`` = the real ``build_paths``). .. py:property:: scheduler :type: collections.abc.Callable[[collections.abc.Coroutine[Any, Any, None]], asyncio.Task[Any]] The registration hand-off scheduler (data; production ``create_task``). .. py:property:: trim :type: collections.abc.Callable[..., None] The python-state trim (the registration hand-off's completion duty). .. py:property:: pump_finished_watchdog :type: collections.abc.Callable[..., collections.abc.Coroutine[Any, Any, None]] The pump-finished watchdog factory (the watch's always-on member). .. py:method:: start() -> BotRunner :async: Build the actors, fetch block state, load snapshots, run ``engine_registry.start()``. Stops at ``Backfilled`` — BEFORE ``resume()``. Zero result batches emit during this window (the pump isn't running), so ``run()`` can attach the consumer in the gap before ``resume()`` without a stale-backlog window. Idempotent via the phase alone: re-entry once Started is a no-op; Running/Closed re-entry raises :class:`PhaseError`. :returns: The started runner, stopped at ``Backfilled``. :raises RuntimeError: If the latest-block fetch fails at session start. .. py:method:: run() -> degenbot.runner._session_watch.SessionEndVerdict :async: Run the cockpit main loop until the consumer task ends. Requires the ``Started`` phase — the session phase machine (:class:`PhaseError`, delegating to the Rust host's ``SessionPhase`` table) owns the lifecycle gate. The startup ordering itself is the run ritual's (:mod:`~degenbot.runner._run_ritual`): this method is the thin phase-gated drive over that machine. :returns: The session's end verdict (the session watch's ranking over how the run ended). Consuming it is optional — callers that ignore it behave exactly as before this return existed. .. py:method:: enqueue_path(path_steps: Any, directions: list[bool] | None = None) -> None :async: Add ONE specific path at any time (the operator surface). Delegates to the session's live :class:`PathRegistrationPipeline` (created by the run ritual); ``path_steps`` + optional ``directions`` are the same shapes as :meth:`PathRegistrationPipeline.enqueue_path`. The path is built via the retained ``ConstructionContext`` (Rust ``PoolBuilder``), registered + verified, released to ``Live``, and registered — without disturbing the pump's update/solve/dispatch. :raises RuntimeError: if no live pipeline exists (injected fake builders have no construction surface, or ``run()`` has not run). .. py:method:: trigger_discovery(*, bound: int | None = None) -> int :async: Trigger a bounded one-shot discovery sweep (on-demand trigger). Delegates to the session's live pipeline. :returns: The number of paths processed by the sweep. :raises RuntimeError: if no live pipeline exists (injected fake builders, or ``run()`` has not run). .. py:method:: shutdown() -> None :async: Signal the Rust core to stop the pump (best-effort). Safe to call at any point in the lifecycle — before ``start()`` finished (``engine_registry`` may be ``None``), after ``run()`` exited, or from a ``SIGINT``/``KeyboardInterrupt`` handler. Mirrors the Rust ``stop()`` contract: idempotent, sets the shutdown flag + aborts the pump task so the WS stream's ``combined.next().await`` unblocks immediately (60s cold-shutdown otherwise). Any exception is swallowed and logged so a partial-startup teardown can't mask the original in-flight exception. This is the one place that closes the Rust core's pump — the ``KeyboardInterrupt``-exits-slowly bug was the pump task (spawned on the shared tokio runtime, decoupled from the asyncio loop) blocking on a silent WS subscription, which ``asyncio.run``'s teardown did not reach until the OS closed the socket. .. py:class:: ConstructionContext Registration-owned construction resources, kept out of run()'s trim. Bundles everything ``build_paths`` needs to construct and register pools, so the registration task owns them as a single self-contained context for its lifetime. ``BotRunner.run()`` trims *main-loop* state (``release_python_state()`` + ``self.bot = None``); the context is a *separate* identity that a background registration task holds and that the trim never severs — the decoupling seam for Sub-B (background registration on the pump runtime). The context holds RESOLVED POLICY VALUES only: the construction route (CONTEXT.md, Construction route — the ordered factory rungs + the generic builder rung) and the WETH token, built once here. The core route entry (``pool_builder::route``) owns the walk — route order, the DB two-step identity, get-or-register into session state — and classifies every failure on the build-refusal taxonomy, so no tracker or snapshot object lives here (the retired three-tracker fallback chain was the bare except-and-continue bug this replaces). .. py:attribute:: bot :type: degenbot.Bot .. py:attribute:: chain_id :type: int .. py:attribute:: database_path :type: pathlib.Path .. py:attribute:: construction_route :type: degenbot.builders.request.ConstructionRoute .. py:attribute:: weth :type: Any .. py:method:: for_bot(bot: degenbot.Bot) -> ConstructionContext :classmethod: Build the construction context for a bot. Resolves the route policy (the mainnet V3 fork factories in policy order, generic builder rung armed) and builds WETH once. The core route entry owns everything else about construction. :returns: The construction context. .. py:class:: PathRegistrationPipeline(*, context: ConstructionContext, engine_registry: degenbot.arbitrage.engine_registry.EngineRegistry, retry_policy: degenbot.arbitrage.RetryPolicy | None = None, max_paths: int, discovery_batch_size: int, progress_interval_secs: float | None = None) Reusable, pump-concurrent registration pipeline (D1c). Owns the per-path registration work that ``build_paths`` previously ran inline: construction (through the retained ``ConstructionContext`` — the Rust ``PoolBuilder``), engine registration + verification, direction resolution, registered-path dedup, per-path release, and the summary counters. It is LONG-LIVED by design: it keeps the ``ConstructionContext`` AND the ``engine_registry`` for the session's lifetime, so an operator can add a specific path (``enqueue_path``) or trigger a bounded on-demand discovery (``trigger_discovery``) at ANY time — including after ``run()`` trims the main-loop bot. The context survives the trim (Sub-A seam), so these methods never need the dropped Python ``bot``. The pipeline never awaits the pump, so adds/discovery cannot block update/solve/dispatch. ``max_paths`` is the registered-path cap (``0`` = uncapped) and is REQUIRED, because it is a configuration value: the caller that resolved it (:attr:`~degenbot.runner.config.ArbitrageConfig.max_registered_paths` in production) states it, and a default here would be a second authority that no config layer can reach. Registration stops accepting new paths once the engine path registry reaches the cap; the engine then reaches steady state with a bounded path universe, so solve performance is observable without ongoing registration load. The pipeline announces the cap it was given at startup, since the code default and the running environment may differ. The fail-fast tripwire is preserved: a fatal ``VerificationMismatchError`` / ``VerificationRpcError`` is NOT swallowed here — it propagates out of the worker and must abort the pipeline loudly. .. py:attribute:: constr_ctx .. py:attribute:: constr_bot .. py:attribute:: constr_chain_id .. py:attribute:: constr_database_path .. py:attribute:: construction_route .. py:attribute:: weth .. py:attribute:: engine_registry .. py:attribute:: retry_policy_obj .. py:attribute:: discovery_batch_size .. py:attribute:: pool_types :type: list[degenbot.pathfinding.PoolKind] :value: [] .. py:attribute:: pool_type_per_depth :type: list[set[degenbot.pathfinding.PoolKind] | None] | None :value: None .. py:attribute:: path_count :value: 0 .. py:attribute:: cap_skip_count :value: 0 .. py:attribute:: skip_count :value: 0 .. py:attribute:: token_filter_count :value: 0 .. py:attribute:: engine_reject_count :value: 0 .. py:attribute:: dup_count :value: 0 .. py:attribute:: register_fail_count :value: 0 .. py:attribute:: v4_pool_count :value: 0 .. py:attribute:: v4_hook_rejected :value: 0 .. py:attribute:: v4_dynamic_fee_rejected :value: 0 .. py:attribute:: other_exc_count :value: 0 .. py:attribute:: capped :value: False .. py:method:: emit_registration_progress(*, force: bool = False) -> None Log the registration counters + top skip-reason breakdown. The legacy ``[build_paths] Progress`` line only fires when ``path_count`` reaches a multiple of 1000. During a discovery-heavy crawl that registers few paths it never fires, so the skip/dup/reject counts (and their reasons) stay invisible. This is the same summary emitted on ``force`` (a wall-clock cadence) so the cause is always observable mid-crawl. .. py:method:: run_registration(*, producer: collections.abc.AsyncIterable[object]) -> None :async: Run the crawl: submit each discovered path as ONE fleet unit. PRG-5: the bounded producer/consumer queue retired with the crawl shell — discovery iterates directly and every path leaves as a single ``PoolStateUpdater`` intake unit (build + verify lifecycles + path registration inside the Rust core). The concurrency is the fleet's (duty-counted seats + its own bounded queue), and the driver-side backpressure is the submission window: at most :data:`REG_INTAKE_WINDOW` receipts are outstanding, so discovery can never outrun registration by more than the window. Units are resolved in FIFO submission order (the retired workers' ordering guarantee), so the Progress summary's counter drift and the 1000-boundary log lines keep their retired shapes exactly. Returns with EVERY submitted receipt resolved (the completion clause that replaced the retired executor-drain: all cloned ``Arc`` handles acquired inside units are dropped before ``build_paths`` returns, keeping the close_snapshot_tx() Arc::try_unwrap canary quiet). The unit's fatal exception (VerificationMismatchError / VerificationRpcError / DirectionResolutionError) propagates through the receipt and aborts the crawl loudly — the "shut down" contract, unchanged. On a fatal the crawl stops submitting immediately (outstanding units still finish — fleet units are never cancelled, the Deferrable cordon class). .. py:method:: enqueue_path(path_steps: Any, directions: list[bool] | None = None) -> None :async: Add ONE specific path at any time (D1c operator surface). .. py:method:: trigger_discovery(*, bound: int | None = None) -> int :async: Trigger a bounded one-shot discovery sweep (D1c). A sweep that runs to NATURAL completion (not bound-truncated, not capped) latches the structural graph edition; a later trigger over the same edition stops immediately and returns 0 — the unchanged structure can only re-yield paths the pipeline already processed. The latch re-arms itself when the edition changes (pool added or removed), when the probe is unavailable, and after any truncated sweep. :returns: The number of paths processed by the sweep. .. py:method:: discovery_sweep(*, find_paths_async: collections.abc.Callable[..., collections.abc.AsyncGenerator[object, None]] = find_paths_async) -> collections.abc.AsyncGenerator[object, None] Run a single discovery sweep over the DB subgraph (V2/V3/V4 DFS). ``find_paths_async`` is the discovery producer seam (tests inject a recording producer to observe the forwarded batch size); the default is the production adapter. :returns: The async generator of discovered candidate paths. .. py:function:: build_paths(*, bot: degenbot.Bot, engine_registry: degenbot.arbitrage.engine_registry.EngineRegistry, options: BuildPathsOptions) -> None :async: Discover V2/V3/V4 arb paths, build Python pools, register with Rust engine. V4 pools are discovered via find_paths_async and built through ``bot.build_managed_pool()``. V4 pool admission (amount-modifying hooks / dynamic fees) is enforced by the Rust core at registration time, surfacing as typed HookedPoolRejectedError / DynamicFeePoolRejectedError. Each per-pool verify lifecycle runs through the core-owned bounded retry dance with the resolved policy injected (transient ``VerificationRpcError`` is retried; ``VerificationMismatchError`` is never retried and crashes loudly). Discovery is a single pass over the DB subgraph driven through a reusable :class:`PathRegistrationPipeline`; after it completes the orphan sweep releases Tracked pools whose path was skipped before ``register_vN_pool``. ``options`` is required because it carries the registered-path cap: the caller that resolved the cap states it, and this function never invents one for a pipeline it builds itself. .. py:function:: resolve_directions(pools: list[degenbot.UniswapV2Pool | degenbot.UniswapV3Pool | degenbot.UniswapV4Pool], input_token_address: str) -> list[bool] Determine zero_for_one for each hop so the cycle closes. The cycle: input_token → hop_0 → intermediate → hop_1 → ... → input_token. Returns a list of zfo values (one per hop). The mechanic lives in the core (`degenbot_pathfinding::directions::resolve_directions`); this adapter flattens the constructed pools into the typed seam's hop tuples and re-raises the core's refusal as :class:`DirectionResolutionError`. Resolution is kind-blind. V4 pools use NATIVE_CURRENCY_ADDRESS (address(0)) for ETH, which the core treats as equivalent to WETH — the profit token is always WETH. The core's refusal surfaces as :class:`DirectionResolutionError` — a hop carries neither tracked token, or the cycle does not close: an invariant violation (wrong pool built, stale subgraph, or a builder bug), never a skip. :returns: One zero-for-one value per hop, in hop order. .. py:class:: ArbitrageConfig Unified settlement-arbitrage configuration — one object for the ~20 tunables `main()` reads. Replaces the scattered config sources (the example's dotenv file, module-top constants, and CLI flags) with a single frozen value object. Construct via :meth:`build`; the bridge onto ``main()`` lives in the ``BotRunner`` orchestration. The operational stances resolve through the core verdict, and the operator/executor identity is read from the process environment. Live defaults reproduce ``main()``'s behavior exactly — no new defaults invented. .. py:attribute:: operator_address :type: str .. py:attribute:: operator_private_key :type: str .. py:attribute:: chain_id :type: int .. py:attribute:: node_http :type: str .. py:attribute:: node_ws :type: str .. py:attribute:: executor_address :type: str .. py:attribute:: executor_owner :type: str .. py:attribute:: inject_executor_code :type: bool .. py:attribute:: injected_address :type: str .. py:attribute:: permutation_filter :type: frozenset[str] | None .. py:attribute:: verification_retry_policy :type: degenbot.arbitrage.RetryPolicy .. py:attribute:: erc6909_profit :type: bool .. py:attribute:: reg_progress_secs :type: float .. py:attribute:: max_registered_paths :type: int .. py:attribute:: discovery_batch_size :type: int .. py:attribute:: contracts_dir :type: str .. py:attribute:: sim_exit_on_fail :type: bool .. py:attribute:: sim_exit_ignore_buckets :type: str .. py:attribute:: dry_run :type: bool .. py:attribute:: diag :type: degenbot.runner.diag.DiagConfig .. py:attribute:: executor_runtime :type: str | pathlib.Path | None :value: None .. py:method:: build(*, live: bool, permutation: str | None, rpc: RpcCascadeOverrides | None = None, values: degenbot.config.ConfigValues | None = None) -> ArbitrageConfig :classmethod: Build an ArbitrageConfig from the resolved verdict + CLI flags + process identity. The dotenv file is no longer a source: operator/executor identity is read from the process environment (the launch shell exports it from ``bot.env``), and every ``DEGENBOT_*`` operational stance is a declared schema key that arrives from :func:`degenbot.config.resolved_config` and reaches the operator file as well as the environment; nothing here re-resolves a layer. Behavior: - operator: live mode requires both OPERATOR_ADDRESS/OPERATOR_PRIVATE_KEY from the process environment and refuses a known placeholder key (raises ValueError); dry-run defaults to a valid throwaway key + its derived address. - nodes: delegated to :func:`degenbot.config.resolve_rpc_uris`, so the four-layer cascade (explicit ``node`` > OS env ``DEGENBOT_RPC_{IPC,WS,HTTP}_CHAINID_{cid}`` > the operator file's ``[nodes.*]`` tables > a declared default) applies, per capability. There is **no ``localhost`` default** — a chain with no configured endpoint in any layer raises :class:`RpcNotConfiguredError`. - executor: zero address is a fatal ``ValueError`` (a factory cannot return early like ``main()``'s ``return``). - inject code: when ``simulation.inject_executor_code`` resolves true, the executor address is overridden to ``INJECTED_EXECUTOR_ADDRESS``. - permutation: a CLI string becomes a singleton frozenset; ``None`` stays ``None``. :returns: A frozen ``ArbitrageConfig`` with cascade-resolved ``node_http``/``node_ws``. :raises ValueError: missing or placeholder operator in live mode, zero-address executor, a retired knob present, or ``RpcNotConfiguredError`` (a ``ValueError`` subclass) when no RPC endpoint is configured for ``chain_id`` in any cascade layer.