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 (thestart / build_paths / consume / dispatchseams).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, theDispatcher, the block clock, the sim context, the pipelines) lives on ONE_SessionStatebuilt instart()— 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 onlyengine_registry+async_w3+ dispatcher once trimmed — the Python pool/token caches are scaffolding once the Rust engine owns canonical state.Testability seams (mirrors
EngineRegistry’sengine=seam):bot,engine_registry,async_w3,snapshots,path_builder,consumer, and the registrationschedulerare injectable. When injected,start()/run()orchestrate the fakes and the phase ordering is verifiable offline; whenNone(production), the actors are built fromcfgand the real module functions are called.- cfg¶
- property bot: degenbot.Bot | None¶
The session’s Python-companion bot (
Noneonce 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_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 path_builder: Any¶
The injected path builder (
None= the realbuild_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— BEFOREresume(). Zero result batches emit during this window (the pump isn’t running), sorun()can attach the consumer in the gap beforeresume()without a stale-backlog window. Idempotent via the phase alone: re-entry once Started is a no-op; Running/Closed re-entry raisesPhaseError.- 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
Startedphase — the session phase machine (PhaseError, delegating to the Rust host’sSessionPhasetable) 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+ optionaldirectionsare the same shapes asPathRegistrationPipeline.enqueue_path(). The path is built via the retainedConstructionContext(RustPoolBuilder), registered + verified, released toLive, 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_registrymay beNone), afterrun()exited, or from aSIGINT/KeyboardInterrupthandler. Mirrors the Ruststop()contract: idempotent, sets the shutdown flag + aborts the pump task so the WS stream’scombined.next().awaitunblocks 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, whichasyncio.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_pathsneeds 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¶
- 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_pathspreviously ran inline: construction (through the retainedConstructionContext— the RustPoolBuilder), engine registration + verification, direction resolution, registered-path dedup, per-path release, and the summary counters.It is LONG-LIVED by design: it keeps the
ConstructionContextAND theengine_registryfor 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 afterrun()trims the main-loop bot. The context survives the trim (Sub-A seam), so these methods never need the dropped Pythonbot. The pipeline never awaits the pump, so adds/discovery cannot block update/solve/dispatch.max_pathsis the registered-path cap (0= uncapped) and is REQUIRED, because it is a configuration value: the caller that resolved it (max_registered_pathsin 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/VerificationRpcErroris 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¶
- 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] Progressline only fires whenpath_countreaches 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 onforce(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
PoolStateUpdaterintake 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 mostREG_INTAKE_WINDOWreceipts 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 beforebuild_pathsreturns, 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_asyncis 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 (transientVerificationRpcErroris retried;VerificationMismatchErroris 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 beforeregister_vN_pool.optionsis 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 asDirectionResolutionError— 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 ontomain()lives in theBotRunnerorchestration. The operational stances resolve through the core verdict, and the operator/executor identity is read from the process environment. Live defaults reproducemain()’s behavior exactly — no new defaults invented.- verification_retry_policy: degenbot.arbitrage.RetryPolicy¶
- 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 everyDEGENBOT_*operational stance is a declared schema key that arrives fromdegenbot.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 (explicitnode> OS envDEGENBOT_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 raisesRpcNotConfiguredError.executor: zero address is a fatal
ValueError(a factory cannot return early likemain()’sreturn).inject code: when
simulation.inject_executor_coderesolves true, the executor address is overridden toINJECTED_EXECUTOR_ADDRESS.permutation: a CLI string becomes a singleton frozenset;
NonestaysNone.
- Returns:
A frozen
ArbitrageConfigwith cascade-resolvednode_http/node_ws.- Raises:
ValueError – missing or placeholder operator in live mode, zero-address executor, a retired knob present, or
RpcNotConfiguredError(aValueErrorsubclass) when no RPC endpoint is configured forchain_idin any cascade layer.