35
AWS-standalone migration: drop Worker relay, add Solana RPC/WS failover (#124)
115a637 unverified | """ | |
| modules/local_executor.py — Garden Angel Local On-Chain Execution v1.6 | |
| ────────────────────────────────────────────────────────────────────────────── | |
| v1.6 — modules/oracle.py (the Cloudflare Worker relay client) is deleted; | |
| LocalExecutor is now the only execute_trade() implementation, no more | |
| fallback. ExecutionResult moves to modules/execution_types.py (import | |
| updated below) so this file doesn't import a now-deleted module. | |
| v1.5 (2026-07-11) — FIX/DIAGNOSTIC: both legs' getAmountsOut() quotes are now | |
| pinned to the SAME explicit block_identifier (read once via | |
| w3.eth.block_number at the top of _execute_sync) instead of each | |
| independently resolving "latest". The block number is now embedded in | |
| both skip-reason messages. Prompted by production evidence that doesn't | |
| add up under the old code: two consecutive WBTC/USDT rejections 82 | |
| seconds apart logged byte-identical quotes ("expected back $1,398.00 on | |
| $1,464.74 in, fee $1.32") even though scanner.py's own | |
| reserve_block_ts tracking showed the sell-leg pool had genuinely traded | |
| in between — i.e. this method's "live" quote was not actually reflecting | |
| current chain state at the time it ran. Pinning the block closes one | |
| possible cause (the two legs racing against different blocks a few | |
| seconds apart) outright. It does not by itself rule out the other: | |
| BSC_RPC_URL is a public free-tier Ankr endpoint (see this deployment's | |
| config), and some free-tier RPC gateways cache eth_call responses | |
| server-side independent of the requested block tag. If a future skip | |
| message still shows the SAME block number across two attempts more than | |
| a few seconds apart on BSC (~3s block time), that's on-the-record proof | |
| the RPC endpoint itself is serving stale reads and BSC_RPC_URL needs to | |
| move off the free tier — not something this file can fix on its own. | |
| See scanner.py's matching v18.21 changelog for a related, independently | |
| confirmed loan-sizing bug found while chasing this (the loan-liquidity | |
| cap was applied against each pool's TOTAL two-sided liquidity instead of | |
| the single reserve side actually being traded against, sizing every | |
| loan at ~2x the intended percentage of real pool depth). | |
| v1.4 (2026-07-10) — NEW: optional `qwen=` constructor param | |
| (modules/qwen_client.py). When configured (QWEN_API_KEY set as a Space | |
| secret), a CONFIRMED on-chain execution's Telegram alert gets a short | |
| AI-generated context note appended — e.g. flagging anything unusual | |
| about the pair/spread. Explicitly advisory-only, confirmed with the | |
| operator: Qwen is only ever asked AFTER the trade has already executed | |
| on-chain, so it can never influence the BUY/HOLD decision, sizing, or | |
| whether execute_trade() runs — all unchanged, still 100% the existing | |
| spread/fee/liquidity math and the on-chain minProfit gate. Best-effort: | |
| any Qwen failure/timeout just means no commentary that cycle — the | |
| plain trade alert still sends exactly as it did before this existed. | |
| v1.3 (2026-07-10) — FIX: the pre-flight go/no-go check was priced as if the | |
| loan principal were at risk, when startArbitrage is an ATOMIC flash-loan | |
| round trip — a trade that lands below the on-chain minProfit gate reverts | |
| entirely, costing only gas (cents on BSC), never the loan. The old check | |
| demanded the DOUBLE worst-case slippage haircut still beat amount_in | |
| (an effective ~1%+ spread after DEX fees) before even attempting — a bar | |
| confirmed in production to reject essentially everything, including | |
| spreads a real arb bot would happily race for. Now: attempt whenever the | |
| EXPECTED live quotes clear the caller's profit floor after the 0.09% | |
| flash-loan fee; keep the worst-case haircuts on minIntermediateOut/ | |
| minFinalOut (execution protection); set on-chain minProfit to the | |
| caller's floor so the contract atomically guarantees "net at least the | |
| floor, or revert for pennies." Also: skip-reasons now log in USD instead | |
| of raw wei, so log lines are readable at a glance. | |
| v1.2 (2026-07-10) — NEW: optional `telegram=` constructor param. When a | |
| TelegramClient is passed in, a CONFIRMED on-chain execution fires an | |
| explicit "✅ LIVE TRADE EXECUTED" alert with the tx hash/BscScan link — | |
| deliberately separate from ReportChannel's every-scan reporting (most | |
| users leave that off, since it's noisy) so a real trade is never buried. | |
| Best-effort: a notification failure never affects the already-returned | |
| trade result. See bot.py's matching change (passes self._tg through). | |
| v1.1 (2026-07-10) — FIX: BSC's "WBTC" scan pair is actually BTCB (Binance-Peg | |
| Bitcoin), a different token than Ethereum's real WBTC — constants.py's | |
| BSC_TOKENS keys it "BTCB", not "WBTC". _resolve_token() aliased BNB→WBNB | |
| but had no equivalent WBTC→BTCB alias, so a live BSC WBTC BUY signal | |
| (confirmed in production: pancakeswap, net +$493.54) failed to resolve | |
| and was silently dropped — silently because execute_trade()'s | |
| `except ValueError` branch also didn't log anything. Both are fixed: | |
| WBTC now aliases to BTCB, and every resolution failure is logged at | |
| ERROR so this can never fail invisibly again. | |
| Purpose: this is "fix hunt". Until now, every BUY signal Scanner computed was | |
| handed to OracleClient, which POSTs to a third-party Cloudflare Worker | |
| (garden-angel-production.elghaly.workers.dev/execute) to actually execute the | |
| trade on-chain. That Worker's /execute route has never been shown to work in | |
| this deployment — gac-defi-bot's src/index.js /execute-signal handler (the | |
| Worker-side counterpart) fails closed with 403/501 (see that repo's | |
| SECURITY_AUDIT.md), so effectively every execute_trade() call against it has | |
| returned confirmed=False. Net effect: this bot has been "hunting" — computing | |
| real BUY signals — for as long as it's been deployed, while confirming | |
| approximately zero of them on-chain, because the only thing standing between | |
| "found an opportunity" and "executed it" was a Worker endpoint that was never | |
| wired up on the other end. | |
| LocalExecutor removes that dependency entirely. It implements the exact same | |
| interface OracleClient exposes — async execute_trade(...) -> ExecutionResult, | |
| importing the identical dataclass from modules/oracle.py so callers can't | |
| tell the two apart — but instead of an HTTP round trip to a third-party | |
| Worker, it calls ContractManager.start_arbitrage() directly against the | |
| FlashArbitrageV2 contract this deployment already owns and controls on BSC — | |
| the same contract the Wallet tab already reads from. | |
| Wired in from bot.py, opt-in only: | |
| LOCAL_EXECUTION_ENABLED=true — required. Unset/false leaves bot.py's | |
| existing OracleClient (Worker-based) | |
| path completely unchanged; nothing in | |
| this file runs until this is explicitly | |
| turned on. | |
| DRY_RUN=true — respected exactly like the rest of | |
| this deployment (PayoutManager, etc.): | |
| LIVE router quotes are still fetched | |
| so the numbers are real, but no | |
| transaction is ever built, signed, or | |
| broadcast. Returns confirmed=False with | |
| a clear "(dry run)" message so Scanner | |
| never credits a dry-run as a real | |
| execution. | |
| Safety properties: | |
| - EVM-only, hard-enforced. A `chain` kwarg other than "ARB" or "BSC" (or | |
| unset — every call site targets Arbitrum today, "BSC" kept only as a | |
| legacy-compatible value) is refused up front with confirmed=False. | |
| This never silently attempts an execution against the wrong chain's | |
| contract/router registry. (v1.6 FIX, 2026-07-15 — this note previously | |
| said "BSC-only", stale since the ARB migration; the code below has | |
| accepted ARB since that migration, only the docstring hadn't caught up.) | |
| - Every write goes through the SAME ContractManager/Web3Client the Wallet | |
| tab already uses for reads — one signing key (PAYOUT_PRIVATE_KEY), one | |
| verified RPC resolution path (see web3_client.py v1.2), no second, | |
| divergent execution path to keep in sync. | |
| - Slippage-protected quoting: minIntermediateOut/minFinalOut are computed | |
| from LIVE getAmountsOut() router quotes at execution time — not from | |
| Scanner's pre-trade spread estimate — with a SLIPPAGE_BPS haircut | |
| (constants.py CFG["SLIPPAGE_BPS"], default 50 = 0.50%) applied before | |
| being sent on-chain. | |
| - minProfit is the caller's own floor (min_profit_usd/net_profit_usd — | |
| never invented, never loosened to whatever quotes happen to show), | |
| enforced ATOMICALLY by the contract: the flash-loan round trip either | |
| nets at least the floor or the entire transaction reverts, costing only | |
| gas (a few cents on Arbitrum), never principal. The pre-flight check | |
| (v1.3) only skips trades whose EXPECTED live quotes can't clear the | |
| floor after the 0.05% flash-loan fee (AAVE_FLASHLOAN_FEE_BPS=5, see | |
| constants.py — this note previously said 0.09%, out of sync with that | |
| constant and with this file's own error-message text below) — a | |
| doomed-by-construction tx is never sent, but a genuine opportunity is | |
| attempted even though it might lose the race and revert for pennies. | |
| - Never fabricates a tx_hash or a realized_profit_usd. A transaction that | |
| reverts (receipt.status == 0) is reported confirmed=False. This module | |
| does not attempt to decode a realized-profit event — the deployed | |
| FlashArbitrageV2's exact event ABI has not been independently confirmed | |
| — so realized_profit_usd is left None on success; Scanner's | |
| _notify_payout() already falls back to its own pre-trade estimate | |
| whenever realized_profit_usd is None, so this is an existing, safe | |
| contract, not a new gap. | |
| - Every blocking web3.py call (getAmountsOut, build_transaction, sign, | |
| send_raw_transaction, wait_for_transaction_receipt) runs inside | |
| asyncio.to_thread() — execute_trade() is an async method called from | |
| Scanner's async _notify_payout(), and web3.py's HTTPProvider is fully | |
| synchronous; calling it directly would block the whole bot's event loop | |
| (Telegram polling, keep-alive pings, other in-flight scans) for the | |
| entire duration of an on-chain round trip. | |
| Required env vars (only exercised when LOCAL_EXECUTION_ENABLED=true): | |
| FLASH_ARBITRAGE_CONTRACT_ADDRESS, PAYOUT_PRIVATE_KEY, ARB_RPC_URL | |
| — all already required by ContractManager/Web3Client. LocalExecutor | |
| introduces no new required env vars beyond LOCAL_EXECUTION_ENABLED | |
| and the existing DRY_RUN flag documented above. | |
| """ | |
| from __future__ import annotations | |
| import asyncio | |
| import logging | |
| import time | |
| from typing import Any, Optional | |
| from web3 import Web3 | |
| from modules.contract_manager import ( | |
| ContractManager, get_contract_manager, resolve_router, | |
| get_contract_manager_v3, resolve_v3_router, | |
| UNISWAP_V3_QUOTER_V2, | |
| ) | |
| from modules.execution_types import ExecutionResult | |
| logger = logging.getLogger("garden_angel.local_executor") | |
| ROUTER_QUOTE_ABI = [ | |
| { | |
| "inputs": [ | |
| {"name": "amountIn", "type": "uint256"}, | |
| {"name": "path", "type": "address[]"}, | |
| ], | |
| "name": "getAmountsOut", | |
| "outputs": [{"name": "amounts", "type": "uint256[]"}], | |
| "stateMutability": "view", | |
| "type": "function", | |
| } | |
| ] | |
| # Uniswap V3 QuoterV2 — same simulate-via-eth_call pattern as | |
| # ROUTER_QUOTE_ABI above, used only for the pre-flight go/no-go check on a | |
| # V3-involving leg. See modules/scanner.py's identical ABI/usage for the | |
| # discovery-side counterpart. | |
| V3_QUOTER_ABI = [ | |
| { | |
| "inputs": [ | |
| { | |
| "components": [ | |
| {"name": "tokenIn", "type": "address"}, | |
| {"name": "tokenOut", "type": "address"}, | |
| {"name": "amountIn", "type": "uint256"}, | |
| {"name": "fee", "type": "uint24"}, | |
| {"name": "sqrtPriceLimitX96", "type": "uint160"}, | |
| ], | |
| "name": "params", | |
| "type": "tuple", | |
| } | |
| ], | |
| "name": "quoteExactInputSingle", | |
| "outputs": [ | |
| {"name": "amountOut", "type": "uint256"}, | |
| {"name": "sqrtPriceX96After", "type": "uint160"}, | |
| {"name": "initializedTicksCrossed", "type": "uint32"}, | |
| {"name": "gasEstimate", "type": "uint256"}, | |
| ], | |
| "stateMutability": "nonpayable", | |
| "type": "function", | |
| } | |
| ] | |
| _DEFAULT_SLIPPAGE_BPS = 50 # 0.50% — matches constants.py CFG["SLIPPAGE_BPS"] | |
| _RECEIPT_TIMEOUT_SECS = 180 # matches the rest of this deployment's fund-moving-call budget | |
| def _bps_haircut(amount: int, bps: int) -> int: | |
| """amount * (1 - bps/10000), integer-safe (no float rounding on-chain-bound values).""" | |
| return (amount * (10_000 - bps)) // 10_000 | |
| class _DryRunSkip(Exception): | |
| """Internal signal only — the dry-run path never reaches a real send/receipt.""" | |
| class LocalExecutor: | |
| """ | |
| Drop-in replacement for OracleClient. Same execute_trade() signature | |
| (plus additive optional kwargs) and the same ExecutionResult return | |
| type — no HTTP, no third party, calls ContractManager.start_arbitrage() | |
| directly. | |
| """ | |
| def __init__( | |
| self, | |
| contract_manager: Optional[ContractManager] = None, | |
| dry_run: bool = True, | |
| slippage_bps: int = _DEFAULT_SLIPPAGE_BPS, | |
| token_registry: Optional[dict] = None, | |
| telegram: Optional[Any] = None, | |
| qwen: Optional[Any] = None, | |
| ) -> None: | |
| self._cm = contract_manager or get_contract_manager() | |
| self._dry_run = bool(dry_run) | |
| self._slippage_bps = int(slippage_bps) | |
| # Optional TelegramClient (duck-typed via Any to avoid importing | |
| # modules.telegram_client here — this module has no other reason | |
| # to depend on it). When set, a CONFIRMED on-chain execution fires | |
| # an explicit alert — deliberately separate from ReportChannel's | |
| # every-scan reporting (which most users leave off, since it's | |
| # noisy) so "a real trade just executed" is never buried in scan | |
| # spam. Best-effort: a notification failure never affects the | |
| # trade result already returned to the caller. | |
| self._telegram = telegram | |
| # v1.4 — optional QwenClient (modules/qwen_client.py). Advisory | |
| # ONLY: called after a trade is already confirmed on-chain, purely | |
| # to append a short context note to the Telegram alert above. Can | |
| # never influence the trade that already happened. Inert/None-safe | |
| # if not passed or not configured (QWEN_API_KEY unset). | |
| self._qwen = qwen | |
| if token_registry is None: | |
| from constants import BSC_TOKENS # single source of truth for BSC token addresses | |
| token_registry = BSC_TOKENS | |
| self._tokens = { | |
| sym.upper(): Web3.to_checksum_address(info["address"]) | |
| for sym, info in token_registry.items() | |
| } | |
| logger.warning( | |
| "[LocalExecutor] ACTIVE — BUY signals will be executed DIRECTLY " | |
| "on-chain against %s via ContractManager, bypassing the " | |
| "Cloudflare Worker /execute path entirely. dry_run=%s " | |
| "slippage_bps=%d", | |
| self._cm.contract.address, self._dry_run, self._slippage_bps, | |
| ) | |
| def _resolve_token(self, symbol: str) -> str: | |
| sym = symbol.upper() | |
| # Migrated 2026-07 to Arbitrum One. Native gas token is ETH, which has | |
| # no ERC-20 address; the contract only ever moves wrapped WETH. WBTC on | |
| # Arbitrum is real WBTC (8 decimals), so no BTCB-style alias is needed. | |
| if sym in ("ETH", "WETH"): | |
| sym = "WETH" | |
| if sym not in self._tokens: | |
| raise ValueError( | |
| f"LocalExecutor: unknown token symbol '{symbol}' — not in " | |
| f"ARB_TOKENS. Known: {', '.join(sorted(self._tokens))}." | |
| ) | |
| return self._tokens[sym] | |
| async def execute_trade( | |
| self, | |
| base_asset: str, | |
| stable_asset: str, | |
| target_dex: str, | |
| amount: float, | |
| net_profit_usd: float, | |
| sell_dex: Optional[str] = None, | |
| chain: Optional[str] = None, | |
| min_profit_usd: Optional[float] = None, | |
| buy_fee_tier: int = 0, | |
| sell_fee_tier: int = 0, | |
| ) -> ExecutionResult: | |
| t0 = time.monotonic() | |
| if chain and chain.upper() not in ("ARB", "BSC"): | |
| return ExecutionResult( | |
| confirmed=False, | |
| error=f"LocalExecutor (EVM) runs on Arbitrum One, got chain={chain!r}. " | |
| f"Solana signals are handled by modules/solana_arb.py.", | |
| duration_ms=(time.monotonic() - t0) * 1000, | |
| ) | |
| if not sell_dex: | |
| return ExecutionResult( | |
| confirmed=False, | |
| error="LocalExecutor requires sell_dex (the exit-leg DEX) " | |
| "— caller did not provide one.", | |
| duration_ms=(time.monotonic() - t0) * 1000, | |
| ) | |
| # Uniswap V3 leg(s) — a SEPARATE contract/path (FlashArbitrageV3.sol, | |
| # see that file's own header for why FlashArbitrageV2 can't run | |
| # this). Dispatched here rather than woven into the V2-only path | |
| # below so a live-trading deployment's existing, already-proven V2 | |
| # execution is completely untouched by this addition. The common | |
| # case (buy_fee_tier == sell_fee_tier == 0, i.e. every trade before | |
| # this existed) falls straight through to the unchanged code below. | |
| if buy_fee_tier or sell_fee_tier: | |
| return await self._execute_trade_mixed( | |
| base_asset, stable_asset, target_dex, sell_dex, amount, | |
| net_profit_usd, min_profit_usd, buy_fee_tier, sell_fee_tier, t0, | |
| ) | |
| try: | |
| stable_addr = self._resolve_token(stable_asset) | |
| base_addr = self._resolve_token(base_asset) | |
| router_buy = resolve_router(target_dex) | |
| router_sell = resolve_router(sell_dex) | |
| except ValueError as exc: | |
| # FIX — this branch previously returned without logging anything, | |
| # so a token/router resolution failure (e.g. the WBTC/BTCB gap | |
| # this exact bug fixed) was completely silent: no error line in | |
| # the HF Space console, nothing to grep for. A real BUY signal | |
| # could vanish here with zero trace. Now logged at ERROR so a | |
| # resolution failure is always visible, same as every other | |
| # failure path in this method. | |
| logger.error( | |
| "[LocalExecutor] Token/router resolution failed for " | |
| "%s/%s on %s→%s: %s", | |
| base_asset, stable_asset, target_dex, sell_dex, exc, | |
| ) | |
| return ExecutionResult( | |
| confirmed=False, error=str(exc), | |
| duration_ms=(time.monotonic() - t0) * 1000, | |
| ) | |
| floor_usd = min_profit_usd if min_profit_usd is not None else max(net_profit_usd, 0.0) | |
| try: | |
| tx_hash, realized = await asyncio.to_thread( | |
| self._execute_sync, | |
| stable_addr, base_addr, router_buy, router_sell, | |
| amount, floor_usd, | |
| ) | |
| except _DryRunSkip as exc: | |
| logger.info("[LocalExecutor] %s", exc) | |
| return ExecutionResult( | |
| confirmed=False, error=f"(dry run) {exc}", | |
| duration_ms=(time.monotonic() - t0) * 1000, | |
| ) | |
| except Exception as exc: | |
| duration_ms = (time.monotonic() - t0) * 1000 | |
| logger.error( | |
| "[LocalExecutor] Execution failed for %s/%s on %s→%s: %s", | |
| base_asset, stable_asset, target_dex, sell_dex, exc, | |
| ) | |
| return ExecutionResult( | |
| confirmed=False, error=str(exc)[:300], | |
| duration_ms=duration_ms, | |
| ) | |
| duration_ms = (time.monotonic() - t0) * 1000 | |
| logger.info( | |
| "[LocalExecutor] Execution confirmed on-chain in %.0fms — tx=%s", | |
| duration_ms, tx_hash, | |
| ) | |
| await self._notify_confirmed( | |
| tx_hash, base_asset, stable_asset, target_dex, sell_dex, | |
| amount, net_profit_usd, realized, | |
| ) | |
| return ExecutionResult( | |
| confirmed=True, | |
| tx_hash=tx_hash, | |
| realized_profit_usd=realized, | |
| status_raw="confirmed", | |
| duration_ms=duration_ms, | |
| ) | |
| # ------------------------------------------------------------------ | |
| # Uniswap V3 mixed-leg execution — separate contract, opt-in only. | |
| # ------------------------------------------------------------------ | |
| async def _execute_trade_mixed( | |
| self, | |
| base_asset: str, stable_asset: str, | |
| target_dex: str, sell_dex: str, | |
| amount: float, net_profit_usd: float, min_profit_usd: Optional[float], | |
| buy_fee_tier: int, sell_fee_tier: int, | |
| t0: float, | |
| ) -> ExecutionResult: | |
| """Mirrors execute_trade()'s V2-only path above, but against | |
| FlashArbitrageV3.sol (a SEPARATE, independently-deployed contract — | |
| see that file's own header). Fails safe and explicitly: if the | |
| operator hasn't deployed+configured it, this returns confirmed=False | |
| with a clear reason, exactly like every other resolution failure in | |
| this module — it never falls back to silently skipping the trade | |
| without a trace, and it never attempts to force the V2-only path | |
| with mismatched router interfaces. | |
| """ | |
| import constants # local import — same lazy-import convention _resolve_token uses | |
| if not constants.v3_execution_enabled(): | |
| return ExecutionResult( | |
| confirmed=False, | |
| error="Uniswap V3 leg found but V3_EXECUTION_ENABLED is not " | |
| "'true' — reported only, not executed. Deploy " | |
| "Contrect/FlashArbitrageV3.sol and set " | |
| "FLASH_ARBITRAGE_V3_CONTRACT_ADDRESS + " | |
| "V3_EXECUTION_ENABLED=true to go live on this leg.", | |
| duration_ms=(time.monotonic() - t0) * 1000, | |
| ) | |
| cm_v3 = get_contract_manager_v3() | |
| if cm_v3 is None: | |
| return ExecutionResult( | |
| confirmed=False, | |
| error="Uniswap V3 leg found but FLASH_ARBITRAGE_V3_CONTRACT_ADDRESS " | |
| "is unset/unreachable — deploy Contrect/FlashArbitrageV3.sol " | |
| "yourself first (see that file's header).", | |
| duration_ms=(time.monotonic() - t0) * 1000, | |
| ) | |
| try: | |
| stable_addr = self._resolve_token(stable_asset) | |
| base_addr = self._resolve_token(base_asset) | |
| router_buy = resolve_v3_router(target_dex) if buy_fee_tier else resolve_router(target_dex) | |
| router_sell = resolve_v3_router(sell_dex) if sell_fee_tier else resolve_router(sell_dex) | |
| except ValueError as exc: | |
| logger.error( | |
| "[LocalExecutor] V3 token/router resolution failed for " | |
| "%s/%s on %s→%s: %s", | |
| base_asset, stable_asset, target_dex, sell_dex, exc, | |
| ) | |
| return ExecutionResult( | |
| confirmed=False, error=str(exc), | |
| duration_ms=(time.monotonic() - t0) * 1000, | |
| ) | |
| floor_usd = min_profit_usd if min_profit_usd is not None else max(net_profit_usd, 0.0) | |
| try: | |
| tx_hash, realized = await asyncio.to_thread( | |
| self._execute_sync_mixed, | |
| cm_v3, stable_addr, base_addr, | |
| router_buy, bool(buy_fee_tier), int(buy_fee_tier), | |
| router_sell, bool(sell_fee_tier), int(sell_fee_tier), | |
| amount, floor_usd, | |
| ) | |
| except _DryRunSkip as exc: | |
| logger.info("[LocalExecutor] %s", exc) | |
| return ExecutionResult( | |
| confirmed=False, error=f"(dry run) {exc}", | |
| duration_ms=(time.monotonic() - t0) * 1000, | |
| ) | |
| except Exception as exc: | |
| duration_ms = (time.monotonic() - t0) * 1000 | |
| logger.error( | |
| "[LocalExecutor] V3 execution failed for %s/%s on %s→%s: %s", | |
| base_asset, stable_asset, target_dex, sell_dex, exc, | |
| ) | |
| return ExecutionResult( | |
| confirmed=False, error=str(exc)[:300], | |
| duration_ms=duration_ms, | |
| ) | |
| duration_ms = (time.monotonic() - t0) * 1000 | |
| logger.info( | |
| "[LocalExecutor] V3 execution confirmed on-chain in %.0fms — tx=%s", | |
| duration_ms, tx_hash, | |
| ) | |
| await self._notify_confirmed( | |
| tx_hash, base_asset, stable_asset, target_dex, sell_dex, | |
| amount, net_profit_usd, realized, | |
| ) | |
| return ExecutionResult( | |
| confirmed=True, | |
| tx_hash=tx_hash, | |
| realized_profit_usd=realized, | |
| status_raw="confirmed", | |
| duration_ms=duration_ms, | |
| ) | |
| def _execute_sync_mixed( | |
| self, | |
| cm_v3, | |
| stable_addr: str, base_addr: str, | |
| router_buy: str, buy_is_v3: bool, buy_fee_tier: int, | |
| router_sell: str, sell_is_v3: bool, sell_fee_tier: int, | |
| amount: float, floor_usd: float, | |
| ) -> tuple[str, Optional[float]]: | |
| """Runs inside asyncio.to_thread(). Mirrors _execute_sync's pre- | |
| flight/slippage/on-chain-minProfit structure exactly, but quotes a | |
| V3 leg via QuoterV2.quoteExactInputSingle instead of getAmountsOut() | |
| (V3 has no two-address `path` — it's single-hop per fee tier).""" | |
| w3 = cm_v3.w3 | |
| amount_in = cm_v3.assets.scale_amount(stable_addr, amount) | |
| floor_wei = cm_v3.assets.scale_amount(stable_addr, max(floor_usd, 0.0)) | |
| block_number = w3.eth.block_number | |
| def _quote_leg(is_v3: bool, router: str, fee: int, token_in: str, token_out: str, amount_in_raw: int) -> int: | |
| if is_v3: | |
| quoter = w3.eth.contract( | |
| address=Web3.to_checksum_address(UNISWAP_V3_QUOTER_V2), abi=V3_QUOTER_ABI, | |
| ) | |
| params = ( | |
| Web3.to_checksum_address(token_in), Web3.to_checksum_address(token_out), | |
| amount_in_raw, fee, 0, | |
| ) | |
| # Pinned to the same explicit block as the V2 leg's | |
| # getAmountsOut() call below (see _execute_sync's v1.5 | |
| # note on why both legs must share one consistent block). | |
| out = quoter.functions.quoteExactInputSingle(params).call( | |
| block_identifier=block_number | |
| ) | |
| return int(out[0]) | |
| router_c = w3.eth.contract(address=router, abi=ROUTER_QUOTE_ABI) | |
| quote = router_c.functions.getAmountsOut( | |
| amount_in_raw, [token_in, token_out] | |
| ).call(block_identifier=block_number) | |
| return int(quote[-1]) | |
| expected_intermediate = _quote_leg(buy_is_v3, router_buy, buy_fee_tier, stable_addr, base_addr, amount_in) | |
| min_intermediate_out = _bps_haircut(expected_intermediate, self._slippage_bps) | |
| expected_final = _quote_leg(sell_is_v3, router_sell, sell_fee_tier, base_addr, stable_addr, expected_intermediate) | |
| min_final_out = _bps_haircut(expected_final, 2 * self._slippage_bps) | |
| from constants import AAVE_FLASHLOAN_FEE_BPS | |
| flash_fee_wei = (amount_in * AAVE_FLASHLOAN_FEE_BPS) // 10_000 | |
| expected_profit_wei = expected_final - amount_in - flash_fee_wei | |
| dec = cm_v3.assets.decimals(stable_addr) | |
| def _usd(wei: int) -> str: | |
| return f"${wei / 10 ** dec:,.2f}" | |
| if expected_profit_wei <= 0: | |
| raise RuntimeError( | |
| f"V3 live quotes unprofitable after the 0.05% flash-loan fee " | |
| f"(expected back {_usd(expected_final)} on {_usd(amount_in)} in, " | |
| f"fee {_usd(flash_fee_wei)}) — spread too thin, skipping." | |
| ) | |
| if floor_wei > 0 and expected_profit_wei < floor_wei: | |
| raise RuntimeError( | |
| f"V3 expected profit {_usd(expected_profit_wei)} is below the " | |
| f"configured floor {_usd(floor_wei)} — skipping rather than " | |
| "sending a trade that would revert on the contract's own " | |
| "minProfit gate." | |
| ) | |
| min_profit = floor_wei if floor_wei > 0 else 1 | |
| if self._dry_run: | |
| raise _DryRunSkip( | |
| f"would call startArbitrageMixed(asset={stable_addr}, " | |
| f"amount={amount_in}, routerBuy={router_buy} v3={buy_is_v3} " | |
| f"fee={buy_fee_tier}, routerSell={router_sell} v3={sell_is_v3} " | |
| f"fee={sell_fee_tier}, intermediateToken={base_addr}, " | |
| f"minIntermediateOut={min_intermediate_out}, " | |
| f"minFinalOut={min_final_out}, minProfit={min_profit}) — " | |
| f"DRY_RUN=true, nothing broadcast." | |
| ) | |
| tx_hash = cm_v3.start_arbitrage_mixed( | |
| asset=stable_addr, | |
| amount=amount_in, | |
| router_buy=router_buy, buy_is_v3=buy_is_v3, buy_fee_tier=buy_fee_tier, | |
| router_sell=router_sell, sell_is_v3=sell_is_v3, sell_fee_tier=sell_fee_tier, | |
| intermediate_token=base_addr, | |
| min_intermediate_out=min_intermediate_out, | |
| min_final_out=min_final_out, | |
| min_profit=min_profit, | |
| ) | |
| receipt = w3.eth.wait_for_transaction_receipt(tx_hash, timeout=_RECEIPT_TIMEOUT_SECS) | |
| if receipt.status != 1: | |
| raise RuntimeError(f"Transaction {tx_hash} reverted on-chain (status=0).") | |
| return tx_hash, None | |
| async def _notify_confirmed( | |
| self, | |
| tx_hash: str, base_asset: str, stable_asset: str, | |
| target_dex: str, sell_dex: str, | |
| amount: float, net_profit_usd: float, realized: Optional[float], | |
| ) -> None: | |
| """Best-effort Telegram alert for a CONFIRMED execution — see | |
| __init__'s self._telegram docstring note. Never raises: a | |
| notification failure must never look like the trade itself | |
| failed, since by this point it has already succeeded on-chain.""" | |
| if self._telegram is None: | |
| return | |
| profit_usd = realized if realized is not None else net_profit_usd | |
| profit_line = ( | |
| f"${realized:,.2f} (realized)" if realized is not None | |
| else f"~${net_profit_usd:,.2f} (pre-trade estimate)" | |
| ) | |
| # FIX (2026-07-15, operator went live — DRY_RUN=false) — this linked | |
| # to bscscan.com, a leftover from before the migration. ContractManager/ | |
| # Web3Client (self._cm) only ever operate against Arbitrum in this | |
| # deployment (see execute_trade's own chain gate above), so every real | |
| # tx_hash here is an Arbitrum One transaction — bscscan.com can never | |
| # find it. This is exactly the link an operator opens the moment their | |
| # very first live trade confirms, to verify funds actually moved. | |
| text = ( | |
| f"✅ LIVE TRADE EXECUTED\n" | |
| f"{base_asset}/{stable_asset} — {target_dex}→{sell_dex}\n" | |
| f"Loan: ${amount:,.2f} | Profit: {profit_line}\n" | |
| f"https://arbiscan.io/tx/{tx_hash}" | |
| ) | |
| # v1.4 — advisory-only Qwen commentary, appended AFTER the trade | |
| # already confirmed on-chain (see __init__'s self._qwen docstring | |
| # note). approx spread_pct here is just profit/loan — a rough | |
| # proxy for the commentary prompt, not the scanner's precise | |
| # pre-trade spread_pct (which isn't available at this call site). | |
| # Never allowed to delay or block the core alert: if this fails, | |
| # times out, or isn't configured, the plain alert above still | |
| # sends exactly as it always has. | |
| if self._qwen is not None and getattr(self._qwen, "enabled", False): | |
| try: | |
| approx_spread_pct = (profit_usd / amount * 100) if amount else 0.0 | |
| commentary = await self._qwen.analyze_opportunity( | |
| base_asset=base_asset, stable_asset=stable_asset, | |
| buy_dex=target_dex, sell_dex=sell_dex, | |
| spread_pct=approx_spread_pct, loan_amount=amount, | |
| net_profit_usd=profit_usd, | |
| ) | |
| if commentary: | |
| text += f"\n\n🤖 {commentary}" | |
| except Exception as exc: | |
| logger.warning("[LocalExecutor] Qwen commentary failed (non-fatal): %s", exc) | |
| try: | |
| await self._telegram.safe_send(text) | |
| except Exception as exc: | |
| logger.warning("[LocalExecutor] Confirmed-trade Telegram alert failed: %s", exc) | |
| def _execute_sync( | |
| self, | |
| stable_addr: str, base_addr: str, | |
| router_buy: str, router_sell: str, | |
| amount: float, floor_usd: float, | |
| ) -> tuple[str, Optional[float]]: | |
| """Runs inside asyncio.to_thread() — blocking web3.py calls only.""" | |
| w3 = self._cm.w3 | |
| amount_in = self._cm.assets.scale_amount(stable_addr, amount) | |
| floor_wei = self._cm.assets.scale_amount(stable_addr, max(floor_usd, 0.0)) | |
| # v1.5 FIX/DIAGNOSTIC (2026-07-11) — pin both legs' getAmountsOut() | |
| # calls to the SAME explicit block, and surface that block number | |
| # in every skip/failure message. Prompted by production evidence: | |
| # two consecutive WBTC/USDT rejections 82s apart logged | |
| # byte-identical numbers ("expected back $1,398.00 on $1,464.74 | |
| # in, fee $1.32") even though the sell-leg pool's own | |
| # reserve_block_ts (tracked by scanner.py) showed it had traded | |
| # in between — i.e. the quote this method used was NOT current | |
| # chain state at the time it ran. Two possible causes, both | |
| # addressed by this change: (1) the two legs racing against | |
| # different blocks if a new block lands between the calls — fixed | |
| # by pinning both to one explicit block_identifier instead of each | |
| # implicitly defaulting to "latest" independently; (2) the RPC | |
| # endpoint itself (BSC_RPC_URL — a public free-tier Ankr URL per | |
| # this deployment's config) caching eth_call responses server-side | |
| # regardless of block tag — NOT fixable from this file, but now | |
| # provable: if this block number stops advancing between | |
| # consecutive failures on the same pair, that's the smoking gun, | |
| # and BSC_RPC_URL needs to move off the free-tier endpoint. | |
| block_number = w3.eth.block_number | |
| buy_router = w3.eth.contract(address=router_buy, abi=ROUTER_QUOTE_ABI) | |
| sell_router = w3.eth.contract(address=router_sell, abi=ROUTER_QUOTE_ABI) | |
| buy_quote = buy_router.functions.getAmountsOut( | |
| amount_in, [stable_addr, base_addr] | |
| ).call(block_identifier=block_number) | |
| expected_intermediate = buy_quote[-1] | |
| min_intermediate_out = _bps_haircut(expected_intermediate, self._slippage_bps) | |
| # Quote the sell leg from the EXPECTED intermediate amount — the | |
| # profitability decision should be made on what the quotes say | |
| # will actually happen, not on the worst-case haircut chain. | |
| # Pinned to the SAME block as the buy-leg quote above (see v1.5 | |
| # note) so both legs are evaluated against one consistent chain | |
| # state instead of each independently resolving "latest". | |
| sell_quote = sell_router.functions.getAmountsOut( | |
| expected_intermediate, [base_addr, stable_addr] | |
| ).call(block_identifier=block_number) | |
| expected_final = sell_quote[-1] | |
| # minFinalOut stays worst-case (both legs slipping to their | |
| # haircut limits) — this is EXECUTION protection, enforced by the | |
| # router/contract on-chain, not the go/no-go decision. | |
| min_final_out = _bps_haircut(expected_final, 2 * self._slippage_bps) | |
| # v1.3 — the go/no-go decision. Previous versions demanded the | |
| # DOUBLE-HAIRCUT worst case (min_final_out) still beat amount_in, | |
| # i.e. an effective ~1%+ spread after DEX fees, before even | |
| # attempting — a bar so high it near-never fires on real BSC | |
| # blue-chip pairs. That caution was misplaced: startArbitrage is | |
| # an ATOMIC flash-loan round trip. If reality lands below the | |
| # on-chain minProfit gate at execution time, the ENTIRE | |
| # transaction reverts — the loan never completes, no principal is | |
| # ever at risk, and the only cost is the gas of a reverted tx | |
| # (cents on BSC). So the right economics are: attempt whenever | |
| # the EXPECTED quotes clear the caller's profit floor after the | |
| # flash-loan fee, and let the atomic on-chain minProfit revert be | |
| # the (cheap) loss-taker when the price races away. Losing that | |
| # race occasionally costs cents; never entering it costs every | |
| # real opportunity. | |
| from constants import AAVE_FLASHLOAN_FEE_BPS # single source of truth (5 bps on Arbitrum) | |
| flash_fee_wei = (amount_in * AAVE_FLASHLOAN_FEE_BPS) // 10_000 | |
| expected_profit_wei = expected_final - amount_in - flash_fee_wei | |
| dec = self._cm.assets.decimals(stable_addr) | |
| def _usd(wei: int) -> str: | |
| return f"${wei / 10 ** dec:,.2f}" | |
| if expected_profit_wei <= 0: | |
| raise RuntimeError( | |
| f"Live quotes unprofitable after the 0.05% flash-loan fee " | |
| f"(expected back {_usd(expected_final)} on {_usd(amount_in)} in, " | |
| f"fee {_usd(flash_fee_wei)}, quoted @ block {block_number}) — " | |
| f"spread too thin, skipping." | |
| ) | |
| if floor_wei > 0 and expected_profit_wei < floor_wei: | |
| raise RuntimeError( | |
| f"Expected profit {_usd(expected_profit_wei)} is below the " | |
| f"configured floor {_usd(floor_wei)} (quoted @ block " | |
| f"{block_number}) — skipping rather than sending a trade " | |
| f"that would revert on the contract's own minProfit gate." | |
| ) | |
| # On-chain minProfit = the caller's floor, enforced atomically by | |
| # the contract: it either nets at least this or the whole flash | |
| # loan reverts for pennies of gas. Never loosened to whatever the | |
| # quotes happen to show — the floor is the operator's risk knob. | |
| min_profit = floor_wei if floor_wei > 0 else 1 | |
| if self._dry_run: | |
| raise _DryRunSkip( | |
| f"would call startArbitrage(asset={stable_addr}, " | |
| f"amount={amount_in}, routerBuy={router_buy}, " | |
| f"routerSell={router_sell}, intermediateToken={base_addr}, " | |
| f"minIntermediateOut={min_intermediate_out}, " | |
| f"minFinalOut={min_final_out}, minProfit={min_profit}) — " | |
| f"DRY_RUN=true, nothing broadcast." | |
| ) | |
| tx_hash = self._cm.start_arbitrage( | |
| asset=stable_addr, | |
| amount=amount_in, | |
| router_buy=router_buy, | |
| router_sell=router_sell, | |
| intermediate_token=base_addr, | |
| min_intermediate_out=min_intermediate_out, | |
| min_final_out=min_final_out, | |
| min_profit=min_profit, | |
| ) | |
| receipt = w3.eth.wait_for_transaction_receipt(tx_hash, timeout=_RECEIPT_TIMEOUT_SECS) | |
| if receipt.status != 1: | |
| raise RuntimeError(f"Transaction {tx_hash} reverted on-chain (status=0).") | |
| # realized_profit_usd deliberately left None — see module docstring | |
| # (no independently-confirmed event ABI to decode it from). Scanner | |
| # falls back to its own pre-trade estimate for ledger crediting. | |
| return tx_hash, None | |
| _local_executor: Optional[LocalExecutor] = None | |
| def get_local_executor() -> LocalExecutor: | |
| global _local_executor | |
| if _local_executor is None: | |
| import os | |
| _local_executor = LocalExecutor( | |
| dry_run=os.environ.get("DRY_RUN", "false").strip().lower() == "true", | |
| ) | |
| return _local_executor | |