Trading Engine

Trading Controller Lifecycle

Detailed analysis of the initialization, state management, position lifecycle, event processing loop, and graceful shutdown inside the core TradingController.

⏱️ 12 min read📊 Level: Advanced

The TradingController (ot_module/controller.py, ~13,800 LOC) is the central nervous system of the DepthSight bot execution engine. It orchestrates incoming market data, runs strategy instances, validates every signal through the RiskManager, and dispatches orders to exchange executors — all within a single asynchronous event loop.


Critical Structural Components

When initialized, the constructor (init, lines 440–530) establishes these critical structures to coordinate asynchronous processing and ensure thread safety across hundreds of concurrently running strategies:

ComponentTypePurpose
event_queuesyncio.Queue(maxsize=1000)Central bounded event bus for incoming market data and order updates
_event_handler_semaphoresyncio.Semaphore(64)Limits concurrent event handlers to prevent thread exhaustion
_positions_dict_locksyncio.LockProtects structural integrity of the ActivePositionMap dictionary
_symbol_locksDict[str, asyncio.Lock]Per-symbol locks ensuring sequential order processing on the same coin
_active_positionsActivePositionMapThread-safe, market-aware runtime storage for active positions
running_strategy_instancesDict[str, Tuple[BaseStrategy, dict]]Maps running strategy config_id to its strategy instance and runtime config
executorsDict[str, Executor]References to "live" (BinanceExecutor) and "paper" (PaperTradingExecutor)
redis_clientredis.asyncio.RedisClient for streaming state changes, runtime variables, and logs

ActivePositionMap — Multi-Market Position Storage

The ActivePositionMap is a custom dict subclass that stores active positions using composite keys of the form market_type:SYMBOL (e.g., utures_usdtm:BTCUSDT). It preserves backward compatibility by allowing symbol-only queries when the lookup is unambiguous, and falls back to matching logic when multiple market types share the same ticker symbol. This design is critical for multi-market support (linear futures, spot, coin-margined) on the same asset.


Startup Orchestration ( start() )

The start() method (lines 975–1090) executes a carefully ordered sequence of startup tasks that must complete before any trading activity can occur:

Rendering diagram...

Step 1 — Configuration Loading (lines 980–1020)

Fetches the user's AppConfig from PostgreSQL and applies runtime settings:

Sources:

Key operations:

  • Risk Limits: Applies max_concurrent_trades, daily_loss_limit, leverage constraints to the RiskManager.
  • Notifications: Loads custom Telegram chat IDs for real-time trade and alert delivery.
  • Blacklist: Syncs the static blacklist from the database into the controller's runtime memory.

Step 2 — Symbol Selection Config (lines 1022–1040)

Loads the SymbolSelectionConfig which dictates the bot's market scope:

ModeBehaviorUse Case
STATICTrades only a fixed whitelist of symbols from strategy configConservative, manually curated portfolios
DYNAMICSelects symbols based on NATR (Normalized ATR) volatility rankingAdaptive market making across top movers
ORACLEAI-driven symbol selection using ML Oracle regime parametersExperimental / high-alpha strategies

The controller iterates through selected symbols and initializes a DataConsumer subscription for each, ensuring market data flows before strategies begin evaluation.

Step 3 — Runtime State Recovery (lines 1042–1060)

After a server restart or process migration, the controller must reconnect to its prior state without losing track of open trades:

Sources:

Redis keys used for persistence:

  • depthsight:state:positions:{user_id} — serialized list of active positions
  • depthsight:state:portfolio:{user_id} — portfolio-level stats (daily PnL, balance)
  • depthsight:state:strategies:{user_id} — running strategy instances and their params

Step 4 — Exchange Reconciliation (lines 1062–1090)

Automatically cross-references local DB records with actual exchange positions via CCXTExecutor.fetch_positions(). This critical step handles:

  • Orphaned Positions: If the exchange reports an open position that the controller has no record of (e.g., due to a prior crash), it imports the position with a "RECOVERED" flag.
  • Stale Entries: If the controller has a position in _active_positions but the exchange reports it as closed, it marks the trade completed and writes the closing record to the database.
  • Quantity Discrepancies: If partial fills occurred during downtime, the controller adjusts internal quantity tracking to match exchange reality.

Step 5 — Data Consumer Registration (lines 1092–1110)

The controller registers event listeners with the DataConsumer. For each active symbol:

  1. Candle close events (type: "CANDLE_CLOSE") — triggers strategy evaluation on each new closed candle.
  2. Tick events (type: "TICK") — used for intra-candle stop-loss monitoring and tape-reading strategies.
  3. Order book snapshots — for market-impact-aware execution and paper-trading fill simulation.

Event Processing Loop (_event_loop())

Once started, the controller enters its main asynchronous event loop (lines 1120–1250). This is the heart of all trading activity:

Rendering diagram...

Event Types and Dispatch Logic

Event TypeSourceAction
CANDLE_CLOSEDataConsumer — kline streamEvaluates strategy entry conditions, triggers signal generation
TICKDataConsumer — aggTrade streamUpdates stop-loss/take-profit tracking, trailing stops
ORDER_UPDATEExchange user data streamMatches fill confirmation, updates position state
SIGNALStrategy instance (internal)Dispatches to RiskManager for assessment and execution
MANUAL_EXITREST API / user commandLiquidates specified position
REGIME_CHANGEOracle ML engineTriggers risk re-evaluation, possible mass liquidation

Concurrency Model

The controller uses a triple-lock strategy to prevent race conditions:

  1. _positions_dict_lock (asyncio.Lock): Guards structural mutations to _active_positions dictionary (adding/removing positions).
  2. _symbol_locks[symbol] (asyncio.Lock): Per-symbol ordering guarantee. Ensures events for BTCUSDT are processed sequentially, preventing duplicate orders or conflicting state updates on the same asset.
  3. _event_handler_semaphore (Semaphore 64): Global throttle capping concurrent event handler tasks, preventing memory exhaustion during volatile market conditions when thousands of events arrive per second.

Position Lifecycle (Open — Manage — Close)

Each position traverses a well-defined lifecycle within the controller:

Rendering diagram...

1. Position Entry (_open_position(), lines ~2500–2650)

When a signal passes the RiskManager assessment:

Sources:

2. Position Management (_manage_positions(), lines ~3200–3400)

On each CANDLE_CLOSE or TICK event, the controller iterates active positions:

  • Stop Loss Check: if low <= position.stop_loss (longs) or if high >= position.stop_loss (shorts).
  • Take Profit Check: Same logic applied with TP threshold.
  • Partial Targets: Each configured level is checked independently; on hit, (target_qty / total_qty) * 100% of the position is closed.
  • Trailing Stop: SL ratchets up (long) or down (short) by configurable % of MFE (Maximum Favorable Excursion).
  • Breakeven: After price moves reakeven_activation_pct favorably, SL moves to entry + buffer.
  • DCA / Grid: Unfilled grid orders are monitored and filled as price steps through predefined levels.

3. Position Exit (_close_position(), lines ~3800–4000)

Sources:

Order Submission Flow (Signal to Order)

The path from strategy signal to live order traverses four distinct layers:

Rendering diagram...

RiskManager Gate

The ssess_signal() call (ot_module/risk_manager.py, lines 1284–1793) performs 11 sequential stages of validation before any capital is committed:

StageCheckRejection Code
0Symbol blacklistSYMBOL_BLACKLISTED
1Balance fetchZERO_BALANCE
2Risk limit check (daily loss, max trades)GLOBAL_RISK_LIMIT
3Dynamic strategy/symbol multiplierZERO_RISK
4Entry price validationINVALID_PRICE
5Stop-loss validationSL_WRONG_SIDE, ZERO_SL_DISTANCE
6Reward/Risk ratioLOW_RR
7Exchange lot filters (stepSize, minQty)MIN_QTY_VIOLATION
8Min notional checkMIN_NOTIONAL_VIOLATION
9Dollar-based R/R checkLOW_DOLLAR_RR

Exchange Dispatch

Once approved, the controller selects the correct executor and calls place_order(). The CCXTExecutor.place_order() method (ot_module/exchanges/ccxt_executor.py, lines 605–889) handles:

  1. Symbol normalization: BTCUSDT to BTC/USDT:USDT (futures) or BTC/USDT (spot).
  2. Order type mapping: Binance STOP_MARKET to CCXT "market" with stopLossPrice param.
  3. Exchange-specific parameters: OKX uses posSide, Gate.io uses settle: "usdt", Bitget requires hedged: True.
  4. Rate-limit compliance: Delegated to CCXT's built-in enableRateLimit token bucket.

State Persistence & Redis Publishing

The controller continuously publishes its runtime state to Redis, enabling:

  1. Real-Time Frontend Updates: The WebSocket server subscribes to depthsight:events:positions:{user_id} and forwards updates to the browser dashboard.
  2. Crash Recovery: A new controller instance can reconstruct its active position map from Redis snapshots.
  3. Multi-Process Coordination: The ot_runner.py daemon monitors Redis state keys to decide when to spawn or terminate controller processes.
Sources:

Shutdown & Cleanup Sequence (stop())

When a stop command is received (via Redis command bus or API), the stop() method (lines ~1400–1550) executes a graceful teardown:

Rendering diagram...

Shutdown Steps

  1. _running = False: The event loop terminates on its next iteration.
  2. Queue Drain: Pending events in event_queue are consumed but discarded.
  3. Position Liquidation (configurable): If liquidate_on_stop=True, sends market orders to close all active positions. Otherwise, positions are left open and their state preserved in Redis.
  4. Order Cancellation: Calls Executor.cancel_all_open_orders() for pending orders.
  5. Data Consumer Unsubscription: Removes symbol subscriptions to stop data flow.
  6. State Finalization: Writes final trade records to PostgreSQL and publishes a terminal state snapshot to Redis.

Error Handling & Recovery

The controller implements multiple layers of error resilience:

Exchange API Failures

When place_order() raises an exception (network error, rate limit, exchange maintenance):

Sources:

Position Reconciliation Timer

A background task (lines ~1600–1650) runs every config.POSITION_RECONCILE_INTERVAL (default: 300 seconds) and:

  1. Fetches open positions from the exchange via CCXTExecutor.get_open_positions().
  2. Cross-references with _active_positions.
  3. Repairs any discrepancies — missing positions, quantity mismatches, stale entries.

Graceful Degradation

If Redis becomes unavailable, the controller continues running in a degraded mode:

  • State persistence is disabled (no snapshots).
  • Trading operations continue using in-memory state only.
  • When Redis recovers, a full state sync is triggered automatically.