degenbot.runner¶

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):

  • BotRunner — the runtime driver facade (the start / build_paths / consume / dispatch seams).

  • ArbitrageConfig — the unified frozen config (build).

  • The build family (build_paths / PathRegistrationPipeline / ConstructionContext / resolve_directions) and the CLI arg parser (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¶

Package Contents¶

class degenbot.runner.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 _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.

cfg¶
property bot: 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).

property engine_registry: degenbot.arbitrage.engine_registry.EngineRegistry | None¶

the injected seam.

Type:

The session’s engine registry. Before start()

property async_w3: degenbot.provider.AsyncAlloyProvider | None¶

the injected seam.

Type:

The session’s dispatch-path provider. Before start()

property dispatcher: degenbot.dispatch.Dispatcher | None¶

None (no session yet).

Type:

The session’s dispatcher. Before start()

property session: _SessionState | None¶

The session owner (None only before start()).

property session_watch: degenbot.runner._session_watch.SessionWatch¶

The session watch the ritual attaches the run’s members to.

property consumer: Any¶

The injected consumer (None = the production block loop).

property readiness: degenbot.strategy.StrategyReadinessView | None¶

The readiness view resolved once at the start() boundary.

property settlement_active: bool¶

The resolved settlement-arm disposition (read by the ritual).

property path_builder: Any¶

The injected path builder (None = the real build_paths).

property scheduler: collections.abc.Callable[[collections.abc.Coroutine[Any, Any, None]], asyncio.Task[Any]]¶

The registration hand-off scheduler (data; production create_task).

property trim: collections.abc.Callable[..., None]¶

The python-state trim (the registration hand-off’s completion duty).

property pump_finished_watchdog: collections.abc.Callable[..., collections.abc.Coroutine[Any, Any, None]]¶

The pump-finished watchdog factory (the watch’s always-on member).

async start() → BotRunner¶

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 PhaseError.

Returns:

The started runner, stopped at Backfilled.

Raises:

RuntimeError – If the latest-block fetch fails at session start.

async run() → degenbot.runner._session_watch.SessionEndVerdict¶

Run the cockpit main loop until the consumer task ends.

Requires the Started phase — the session phase machine (PhaseError, delegating to the Rust host’s SessionPhase table) owns the lifecycle gate. The startup ordering itself is the run ritual’s (_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.

async enqueue_path(path_steps: Any, directions: list[bool] | None = None) → None¶

Add ONE specific path at any time (the operator surface).

Delegates to the session’s live PathRegistrationPipeline (created by the run ritual); path_steps + optional directions are the same shapes as 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).

async trigger_discovery(*, bound: int | None = None) → int¶

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).

async shutdown() → None¶

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.

class degenbot.runner.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).

bot: degenbot.Bot¶
chain_id: int¶
database_path: pathlib.Path¶
construction_route: degenbot.builders.request.ConstructionRoute¶
weth: Any¶
classmethod for_bot(bot: degenbot.Bot) → ConstructionContext¶

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.

class degenbot.runner.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 (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.

constr_ctx¶
constr_bot¶
constr_chain_id¶
constr_database_path¶
construction_route¶
weth¶
engine_registry¶
retry_policy_obj¶
discovery_batch_size¶
pool_types: list[degenbot.pathfinding.PoolKind] = []¶
pool_type_per_depth: list[set[degenbot.pathfinding.PoolKind] | None] | None = None¶
path_count = 0¶
cap_skip_count = 0¶
skip_count = 0¶
token_filter_count = 0¶
engine_reject_count = 0¶
dup_count = 0¶
register_fail_count = 0¶
v4_pool_count = 0¶
v4_hook_rejected = 0¶
v4_dynamic_fee_rejected = 0¶
other_exc_count = 0¶
capped = False¶
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.

async run_registration(*, producer: collections.abc.AsyncIterable[object]) → None¶

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 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<SnapshotDb> 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).

async enqueue_path(path_steps: Any, directions: list[bool] | None = None) → None¶

Add ONE specific path at any time (D1c operator surface).

async trigger_discovery(*, bound: int | None = None) → int¶

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.

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.

async degenbot.runner.build_paths(*, bot: degenbot.Bot, engine_registry: degenbot.arbitrage.engine_registry.EngineRegistry, options: BuildPathsOptions) → None¶

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 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.

degenbot.runner.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 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 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.

class degenbot.runner.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 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.

operator_address: str¶
operator_private_key: str¶
chain_id: int¶
node_http: str¶
node_ws: str¶
executor_address: str¶
executor_owner: str¶
inject_executor_code: bool¶
injected_address: str¶
permutation_filter: frozenset[str] | None¶
verification_retry_policy: degenbot.arbitrage.RetryPolicy¶
erc6909_profit: bool¶
reg_progress_secs: float¶
max_registered_paths: int¶
discovery_batch_size: int¶
contracts_dir: str¶
sim_exit_on_fail: bool¶
sim_exit_ignore_buckets: str¶
dry_run: bool¶
diag: degenbot.runner.diag.DiagConfig¶
executor_runtime: str | pathlib.Path | None = None¶
classmethod build(*, live: bool, permutation: str | None, rpc: RpcCascadeOverrides | None = None, values: degenbot.config.ConfigValues | None = None) → ArbitrageConfig¶

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 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 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 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.