221 lines
8.9 KiB
Python
221 lines
8.9 KiB
Python
import asyncio
|
|
import faulthandler
|
|
import logging
|
|
import signal
|
|
from datetime import datetime
|
|
|
|
from bars import BarManager
|
|
from config import config
|
|
from connection import ib_conn
|
|
from state import PositionTracker
|
|
from strategies import MAStockStrategy, ForexMAStrategy, ShortTermMAVWAPStrategy, MeanReversionStrategy
|
|
|
|
logging.basicConfig(
|
|
level=logging.INFO,
|
|
format="%(asctime)s [%(levelname)s] %(name)s: %(message)s",
|
|
datefmt="%H:%M:%S",
|
|
)
|
|
# silence ib_insync's per-message INFO spam (portfolio updates etc.)
|
|
logging.getLogger("ib_insync.wrapper").setLevel(logging.WARNING)
|
|
logging.getLogger("ib_insync.client").setLevel(logging.WARNING)
|
|
logger = logging.getLogger("main")
|
|
|
|
# dump stack traces to stderr on SIGUSR1 (diagnostics for a hung main loop)
|
|
faulthandler.register(signal.SIGUSR1, all_threads=True)
|
|
|
|
|
|
class TradingApp:
|
|
def __init__(self):
|
|
self.strategies: list = []
|
|
self._running = False
|
|
self._reconnect_task: asyncio.Task | None = None
|
|
self.tracker = PositionTracker(config.state_file)
|
|
self.bar_manager = BarManager(ib_conn.ib)
|
|
|
|
def setup_signal_handlers(self):
|
|
for sig in (signal.SIGINT, signal.SIGTERM):
|
|
signal.signal(sig, self._signal_handler)
|
|
|
|
def _signal_handler(self, sig, frame):
|
|
logger.info("Received signal %s, shutting down...", sig)
|
|
self._running = False
|
|
|
|
async def connect(self) -> bool:
|
|
retries = 0
|
|
while retries < config.max_retries:
|
|
if await ib_conn.ensure_connected():
|
|
logger.info(
|
|
"Connected to IB Gateway at %s:%s (%s)",
|
|
config.ib.host, config.ib.port, config.ib.port_label,
|
|
)
|
|
return True
|
|
retries += 1
|
|
if retries < config.max_retries:
|
|
logger.warning("Retry %d/%d in %ds...", retries, config.max_retries, config.retry_delay)
|
|
await asyncio.sleep(config.retry_delay)
|
|
logger.error("Failed to connect after %d retries", config.max_retries)
|
|
return False
|
|
|
|
def _account_position_keys(self) -> dict[str, float]:
|
|
"""Real account positions keyed like the tracker: 'META' or 'EUR.USD'."""
|
|
result: dict[str, float] = {}
|
|
for p in ib_conn.ib.positions():
|
|
c = p.contract
|
|
key = f"{c.symbol}.{c.currency}" if c.secType == "CASH" else c.symbol
|
|
result[key] = result.get(key, 0) + p.position
|
|
return result
|
|
|
|
async def _completed_stp_fills(self) -> dict[str, float]:
|
|
"""{orderRef: avgFillPrice} for today's filled STP orders, so reconcile can
|
|
backfill the ledger with real prices for fills the bot missed while offline."""
|
|
fills: dict[str, float] = {}
|
|
try:
|
|
for t in await ib_conn.ib.reqCompletedOrdersAsync(apiOnly=True):
|
|
ref = getattr(t.order, "orderRef", "") or ""
|
|
fill_price = getattr(t.orderStatus, "avgFillPrice", 0) or 0
|
|
if ref.endswith(":STP") and t.orderStatus.status == "Filled" and fill_price > 0:
|
|
fills[ref] = fill_price
|
|
# ib_insync inserts completed orders into wrapper.trades with a
|
|
# non-terminal OrderStatus (filled=0, status from orderState).
|
|
# Purge them so they don't masquerade as open orders: has_open_order
|
|
# would block stop placement forever for a key whose real STP is
|
|
# missing but whose completed order is still visible (08-04: NOK
|
|
# zombie STP appeared "PreSubmitted" in openTrades after restart).
|
|
ib_conn.ib.wrapper.trades.pop(t.order.permId, None)
|
|
except Exception as e:
|
|
logger.warning("reqCompletedOrders failed: %s", e)
|
|
return fills
|
|
|
|
async def run(self):
|
|
self.setup_signal_handlers()
|
|
self._running = True
|
|
|
|
if config.ib.is_paper:
|
|
logger.info("Starting Trading Bot | Account: %s | Mode: PAPER", config.ib.account)
|
|
else:
|
|
logger.warning("*" * 70)
|
|
logger.warning(
|
|
"LIVE TRADING MODE - real orders will be placed! Account=%s Gateway=%s:%s (%s)",
|
|
config.ib.account, config.ib.host, config.ib.port, config.ib.port_label,
|
|
)
|
|
logger.warning("*" * 70)
|
|
|
|
if not await self.connect():
|
|
logger.error("Initial connection failed, entering reconnect loop (Ctrl+C to abort)...")
|
|
await self._handle_disconnect()
|
|
if not self._running or not ib_conn.is_connected():
|
|
return
|
|
|
|
# registered exactly once; guarded against shutdown and duplicates
|
|
ib_conn.ib.disconnectedEvent += self._on_disconnected
|
|
|
|
if config.stock.enabled:
|
|
self.strategies.append(MAStockStrategy(ib_conn.ib, self.bar_manager, self.tracker))
|
|
if config.short_term.enabled:
|
|
self.strategies.append(ShortTermMAVWAPStrategy(ib_conn.ib, self.bar_manager, self.tracker))
|
|
if config.forex.enabled:
|
|
self.strategies.append(ForexMAStrategy(ib_conn.ib, self.bar_manager, self.tracker))
|
|
if config.mean_reversion.enabled:
|
|
self.strategies.append(MeanReversionStrategy(ib_conn.ib, self.bar_manager, self.tracker))
|
|
|
|
if not self.strategies:
|
|
logger.error("No strategies enabled")
|
|
return
|
|
|
|
# reconcile tracked ownership against real account positions before trading
|
|
external_fills = await self._completed_stp_fills()
|
|
self.tracker.reconcile(self._account_position_keys(), external_fills)
|
|
|
|
# populate open orders so strategies can adopt existing GTC stop orders
|
|
try:
|
|
await ib_conn.ib.reqOpenOrdersAsync()
|
|
except Exception as e:
|
|
logger.warning("reqOpenOrders failed: %s", e)
|
|
|
|
for s in self.strategies:
|
|
await s.on_start()
|
|
logger.info("Strategy loaded: %s", s.name)
|
|
|
|
logger.info("Bot is running. Press Ctrl+C to stop. Check interval: %.0fs", config.loop_interval)
|
|
|
|
try:
|
|
while self._running:
|
|
# skip cycles while the reconnect task is working
|
|
if not ib_conn.is_connected():
|
|
await asyncio.sleep(config.retry_delay)
|
|
continue
|
|
loop_start = datetime.now()
|
|
for s in self.strategies:
|
|
if not self._running:
|
|
break
|
|
try:
|
|
await s.on_bar()
|
|
except Exception as e:
|
|
logger.exception("Error in %s: %s", s.name, e)
|
|
|
|
if not self._running:
|
|
break
|
|
|
|
elapsed = (datetime.now() - loop_start).total_seconds()
|
|
await asyncio.sleep(max(0, config.loop_interval - elapsed))
|
|
except asyncio.CancelledError:
|
|
pass
|
|
finally:
|
|
await self.shutdown()
|
|
|
|
def _on_disconnected(self):
|
|
# ignore disconnects we triggered ourselves during shutdown
|
|
if not self._running:
|
|
return
|
|
# never run two reconnect loops concurrently
|
|
if self._reconnect_task and not self._reconnect_task.done():
|
|
return
|
|
logger.warning("IB Gateway disconnected! Attempting reconnect...")
|
|
self._reconnect_task = asyncio.create_task(self._handle_disconnect())
|
|
|
|
async def _handle_disconnect(self):
|
|
attempt = 0
|
|
while self._running:
|
|
attempt += 1
|
|
delay = min(config.retry_delay * attempt, 60)
|
|
logger.warning("Reconnect attempt %d, waiting %ds...", attempt, delay)
|
|
await asyncio.sleep(delay)
|
|
if not self._running:
|
|
return
|
|
if await ib_conn.ensure_connected():
|
|
logger.info("Reconnected successfully on attempt %d", attempt)
|
|
# live bar subscriptions died with the connection; re-subscribe
|
|
self.bar_manager.reset()
|
|
external_fills = await self._completed_stp_fills()
|
|
self.tracker.reconcile(self._account_position_keys(), external_fills)
|
|
try:
|
|
await ib_conn.ib.reqOpenOrdersAsync()
|
|
except Exception as e:
|
|
logger.warning("reqOpenOrders failed after reconnect: %s", e)
|
|
for s in self.strategies:
|
|
await s.on_start()
|
|
return
|
|
|
|
async def shutdown(self):
|
|
logger.info("Shutting down...")
|
|
self._running = False
|
|
if self._reconnect_task and not self._reconnect_task.done():
|
|
self._reconnect_task.cancel()
|
|
for s in self.strategies:
|
|
try:
|
|
await s.on_stop()
|
|
except Exception as e:
|
|
logger.exception("Error stopping %s: %s", s.name, e)
|
|
self.bar_manager.reset()
|
|
ib_conn.disconnect()
|
|
logger.info("Shutdown complete")
|
|
|
|
|
|
async def main():
|
|
app = TradingApp()
|
|
await app.run()
|
|
|
|
|
|
if __name__ == "__main__":
|
|
asyncio.run(main())
|