| """
|
| execution/order_manager.py β Order submission with limit/bracket, fractional handling,
|
| spread checks, timeout monitoring. Implements LinkedOrderGroup for fractional orders.
|
| """
|
| from __future__ import annotations
|
|
|
| import datetime
|
| import json
|
| import logging
|
| import random
|
| import threading
|
| import time
|
| import uuid
|
|
|
| import config
|
| from contracts import OrderRequest, LinkedOrderGroup
|
| from execution import broker
|
| from data import storage
|
|
|
| logger = logging.getLogger("trading_system.order_manager")
|
|
|
|
|
|
|
|
|
| _virtual_portfolio: dict[str, dict] = {}
|
| _virtual_orders: dict[str, dict] = {}
|
| _virtual_order_counter = 0
|
| _vp_lock = threading.Lock()
|
|
|
|
|
| def _next_virtual_id() -> str:
|
| global _virtual_order_counter
|
| _virtual_order_counter += 1
|
| return f"dry_run_{uuid.uuid4().hex[:12]}"
|
|
|
|
|
| def _simulate_fill(order: OrderRequest) -> dict:
|
| """Simulate an order fill for DRY_RUN mode."""
|
| delay = random.uniform(2, 15)
|
| time.sleep(min(delay, 2))
|
| order_id = _next_virtual_id()
|
| fill = {
|
| "id": order_id,
|
| "symbol": order.symbol,
|
| "side": order.side,
|
| "qty": str(order.qty),
|
| "filled_qty": str(order.qty),
|
| "filled_avg_price": str(order.limit_price),
|
| "status": "filled",
|
| "type": "limit",
|
| "created_at": datetime.datetime.now(datetime.timezone.utc).isoformat(),
|
| }
|
| logger.info("[DRY_RUN] Simulated fill: %s", json.dumps(fill))
|
|
|
| with _vp_lock:
|
| _virtual_orders[order_id] = fill
|
| if order.side == "buy":
|
| _virtual_portfolio[order.symbol] = {
|
| "symbol": order.symbol,
|
| "qty": order.qty,
|
| "side": "buy",
|
| "entry_price": order.limit_price,
|
| "current_price": order.limit_price,
|
| "unrealized_pl": 0.0,
|
| "open_date": datetime.datetime.now(datetime.timezone.utc),
|
| }
|
| elif order.symbol in _virtual_portfolio:
|
| del _virtual_portfolio[order.symbol]
|
|
|
| return fill
|
|
|
|
|
|
|
|
|
| def check_spread(symbol: str) -> tuple[bool, float, float, float]:
|
| """Check if the current bid-ask spread is acceptable.
|
|
|
| Returns: (passes, bid, ask, spread_pct)
|
| For paper trading with stale quotes, falls back to last trade price.
|
| """
|
| if config.DRY_RUN:
|
| return True, 100.0, 100.01, 0.01
|
|
|
| quote_data = broker.get_latest_quote(symbol)
|
| quote = quote_data.get("quote", quote_data)
|
| bid = float(quote.get("bp", quote.get("bid_price", 0)))
|
| ask = float(quote.get("ap", quote.get("ask_price", 0)))
|
|
|
| if ask <= 0:
|
| return False, bid, ask, 100.0
|
|
|
| spread_pct = ((ask - bid) / ask) * 100
|
| passes = spread_pct <= config.MAX_SPREAD_PCT
|
|
|
|
|
|
|
| if not passes and config.TRADING_MODE == "paper":
|
| try:
|
| trade_data = broker.get_latest_trade(symbol)
|
| trade = trade_data.get("trade", trade_data)
|
| last_price = float(trade.get("p", trade.get("price", 0)))
|
| if last_price > 0 and bid < last_price < ask:
|
|
|
| logger.info(
|
| "%s: quote spread %.2f%% is stale (bid=%.2f ask=%.2f), "
|
| "using last trade $%.2f as reference",
|
| symbol, spread_pct, bid, ask, last_price,
|
| )
|
|
|
| bid = round(last_price - 0.01, 2)
|
| ask = round(last_price + 0.01, 2)
|
| spread_pct = 0.01
|
| passes = True
|
| except Exception as e:
|
| logger.debug("Could not fetch last trade for %s: %s", symbol, e)
|
|
|
| if not passes:
|
| logger.warning(
|
| "%s spread too wide: %.2f%% (max %.2f%%)",
|
| symbol, spread_pct, config.MAX_SPREAD_PCT,
|
| )
|
|
|
| return passes, bid, ask, spread_pct
|
|
|
|
|
| def compute_limit_price(side: str, bid: float, ask: float) -> float:
|
| """Compute limit price based on AGGRESSIVE_ENTRY setting."""
|
| if side == "buy":
|
| if config.AGGRESSIVE_ENTRY:
|
| return ask
|
| return round(ask - 0.01, 2)
|
| else:
|
| if config.AGGRESSIVE_ENTRY:
|
| return bid
|
| return round(bid + 0.01, 2)
|
|
|
|
|
|
|
|
|
| def _submit_fractional_order(order: OrderRequest, alert_callback=None) -> dict | None:
|
| """Handle fractional order: 3-step (entry + stop + TP) with orphan detection.
|
|
|
| Step 1: Create LinkedOrderGroup in DB BEFORE any order.
|
| Step 2: Submit entry limit order.
|
| Step 3: On fill β submit stop-loss, then take-profit.
|
| """
|
| entry_order_id = _next_virtual_id() if config.DRY_RUN else None
|
|
|
|
|
| group = LinkedOrderGroup(
|
| symbol=order.symbol,
|
| entry_order_id=entry_order_id or "pending",
|
| is_fractional=True,
|
| orphaned=False,
|
| )
|
|
|
|
|
| if config.DRY_RUN:
|
| fill = _simulate_fill(order)
|
| entry_order_id = fill["id"]
|
| else:
|
| try:
|
| if getattr(config, "ALLOW_MARKET_ORDERS", False):
|
| result = broker.submit_market_order(
|
| symbol=order.symbol,
|
| side=order.side,
|
| qty=order.qty,
|
| )
|
| else:
|
| result = broker.submit_limit_order(
|
| symbol=order.symbol,
|
| side=order.side,
|
| qty=order.qty,
|
| limit_price=order.limit_price,
|
| )
|
| entry_order_id = result.get("id", "")
|
| except Exception as e:
|
| logger.error("Entry order failed for %s: %s", order.symbol, e)
|
| return None
|
|
|
|
|
| group.entry_order_id = entry_order_id
|
| storage.insert_linked_order_group(group.model_dump())
|
| logger.info("LinkedOrderGroup created: %s β %s", order.symbol, entry_order_id)
|
|
|
|
|
| if config.DRY_RUN:
|
| storage.update_linked_order_group(entry_order_id, {
|
| "entry_filled": True,
|
| "entry_price": order.limit_price,
|
| })
|
| _submit_fractional_legs(
|
| entry_order_id, order, order.limit_price, alert_callback
|
| )
|
|
|
| return {"entry_order_id": entry_order_id, "status": "submitted"}
|
|
|
|
|
| def _submit_fractional_legs(
|
| entry_order_id: str,
|
| order: OrderRequest,
|
| fill_price: float,
|
| alert_callback=None,
|
| ):
|
| """Submit stop-loss and take-profit orders after entry is filled."""
|
| close_side = "sell" if order.side == "buy" else "buy"
|
|
|
|
|
| stop_order_id = None
|
| if order.stop_price:
|
| try:
|
| if config.DRY_RUN:
|
| stop_order_id = _next_virtual_id()
|
| logger.info("[DRY_RUN] Simulated stop order: %s", stop_order_id)
|
| else:
|
| if getattr(order, "trail_percent", None):
|
| result = broker.submit_trailing_stop_order(
|
| symbol=order.symbol,
|
| side=close_side,
|
| qty=order.qty,
|
| trail_percent=order.trail_percent,
|
| )
|
| else:
|
| result = broker.submit_stop_order(
|
| symbol=order.symbol,
|
| side=close_side,
|
| qty=order.qty,
|
| stop_price=order.stop_price,
|
| )
|
| stop_order_id = result.get("id")
|
|
|
| storage.update_linked_order_group(entry_order_id, {
|
| "stop_submitted": True,
|
| "stop_order_id": stop_order_id,
|
| })
|
| except Exception as e:
|
| logger.critical(
|
| "Stop-loss submission FAILED for %s (entry %s): %s. ORPHANED.",
|
| order.symbol, entry_order_id, e,
|
| )
|
| storage.update_linked_order_group(entry_order_id, {"orphaned": True})
|
| if alert_callback:
|
| alert_callback(
|
| f"π¨ ORPHANED ORDER: {order.symbol} entry filled but stop-loss failed. "
|
| f"Attempting market close."
|
| )
|
| _emergency_close(order.symbol, order.qty, close_side)
|
| return
|
|
|
|
|
| tp_order_id = None
|
| if order.tp_price:
|
| try:
|
| if config.DRY_RUN:
|
| tp_order_id = _next_virtual_id()
|
| logger.info("[DRY_RUN] Simulated TP order: %s", tp_order_id)
|
| else:
|
| result = broker.submit_limit_tp_order(
|
| symbol=order.symbol,
|
| side=close_side,
|
| qty=order.qty,
|
| limit_price=order.tp_price,
|
| )
|
| tp_order_id = result.get("id")
|
|
|
| storage.update_linked_order_group(entry_order_id, {
|
| "tp_submitted": True,
|
| "tp_order_id": tp_order_id,
|
| })
|
| except Exception as e:
|
| logger.critical(
|
| "TP submission FAILED for %s (entry %s): %s. ORPHANED.",
|
| order.symbol, entry_order_id, e,
|
| )
|
| storage.update_linked_order_group(entry_order_id, {"orphaned": True})
|
| if alert_callback:
|
| alert_callback(
|
| f"π¨ ORPHANED ORDER: {order.symbol} entry filled but TP failed. "
|
| f"Attempting market close."
|
| )
|
|
|
| if stop_order_id and not config.DRY_RUN:
|
| try:
|
| broker.cancel_order(stop_order_id)
|
| except Exception:
|
| pass
|
| _emergency_close(order.symbol, order.qty, close_side)
|
| return
|
|
|
|
|
| def _emergency_close(symbol: str, qty: float, close_side: str):
|
| """Emergency market close when stop/TP legs fail."""
|
| if config.DRY_RUN:
|
| logger.warning("[DRY_RUN] Emergency close simulated: %s", symbol)
|
| with _vp_lock:
|
| _virtual_portfolio.pop(symbol, None)
|
| return
|
|
|
| try:
|
| broker.close_position(symbol)
|
| logger.warning("Emergency position close executed: %s", symbol)
|
| except Exception as e:
|
| logger.critical("EMERGENCY CLOSE FAILED for %s: %s", symbol, e)
|
|
|
|
|
|
|
|
|
| def _submit_bracket_order(order: OrderRequest) -> dict | None:
|
| """Submit a bracket order for whole-share quantities."""
|
| if config.DRY_RUN:
|
| fill = _simulate_fill(order)
|
| return fill
|
|
|
| is_market = getattr(config, "ALLOW_MARKET_ORDERS", False)
|
|
|
| if order.stop_price and order.tp_price:
|
| result = broker.submit_bracket_order(
|
| symbol=order.symbol,
|
| side=order.side,
|
| qty=int(order.qty),
|
| limit_price=order.limit_price,
|
| stop_price=order.stop_price,
|
| tp_price=order.tp_price,
|
| is_market=is_market,
|
| )
|
| else:
|
| if is_market:
|
| result = broker.submit_market_order(
|
| symbol=order.symbol,
|
| side=order.side,
|
| qty=order.qty,
|
| )
|
| else:
|
| result = broker.submit_limit_order(
|
| symbol=order.symbol,
|
| side=order.side,
|
| qty=order.qty,
|
| limit_price=order.limit_price,
|
| )
|
|
|
| return result
|
|
|
|
|
|
|
|
|
| def submit_order(order: OrderRequest, alert_callback=None) -> dict | None:
|
| """Submit an order, routing to fractional or bracket flow.
|
|
|
| Args:
|
| order: OrderRequest contract
|
| alert_callback: Optional callable for critical alerts
|
|
|
| Returns: Order result dict or None on failure
|
| """
|
| prefix = "[DRY_RUN] " if config.DRY_RUN else ""
|
| logger.info(
|
| "%sSubmitting order: %s %s %.4f shares @ $%.2f (stop=$%s, tp=$%s)",
|
| prefix, order.side, order.symbol, order.qty, order.limit_price,
|
| order.stop_price, order.tp_price,
|
| )
|
|
|
| try:
|
| if order.is_fractional or getattr(order, "trail_percent", None):
|
| result = _submit_fractional_order(order, alert_callback)
|
| else:
|
| result = _submit_bracket_order(order)
|
|
|
| if result:
|
|
|
| oid = result.get("id") or result.get("entry_order_id")
|
| if oid and getattr(order, "type", "limit") == "limit":
|
| track_pending_order(oid, order)
|
|
|
| elif oid:
|
| track_pending_order(oid, order)
|
|
|
| return result
|
| except Exception as e:
|
| logger.error("Order submission failed: %s", e)
|
| return None
|
|
|
|
|
|
|
|
|
| _pending_orders: dict[str, dict] = {}
|
| _pending_lock = threading.Lock()
|
|
|
|
|
| def track_pending_order(order_id: str, order: OrderRequest):
|
| """Start tracking an order for timeout."""
|
| with _pending_lock:
|
| _pending_orders[order_id] = {
|
| "submitted_at": time.monotonic(),
|
| "order": order,
|
| }
|
|
|
|
|
| def check_timeouts(portfolio) -> list[str]:
|
| """Cancel timed-out limit orders and release reserved capital."""
|
| cancelled = []
|
| timeout = config.LIMIT_ORDER_TIMEOUT_SEC
|
| now = time.monotonic()
|
|
|
| with _pending_lock:
|
| expired = [
|
| oid for oid, info in _pending_orders.items()
|
| if now - info["submitted_at"] > timeout
|
| ]
|
|
|
| for oid in expired:
|
| try:
|
| if not config.DRY_RUN:
|
| order_status = broker.get_order(oid)
|
| if order_status.get("status") in ("filled", "canceled", "rejected", "expired"):
|
| with _pending_lock:
|
| _pending_orders.pop(oid, None)
|
| continue
|
|
|
| broker.cancel_order(oid)
|
|
|
|
|
|
|
| if portfolio:
|
| portfolio.release_allocation(oid)
|
|
|
| cancelled.append(oid)
|
| logger.info("Timed out order cancelled and allocation released: %s", oid)
|
| except Exception as e:
|
| logger.error("Failed to cancel timed-out order %s: %s", oid, e)
|
|
|
| with _pending_lock:
|
| _pending_orders.pop(oid, None)
|
|
|
| return cancelled
|
|
|
|
|
|
|
|
|
| def check_orphans(alert_callback=None) -> int:
|
| """Check for orphaned linked order groups.
|
|
|
| Runs every 60 seconds. Attempts to resubmit missing legs once.
|
| If resubmission fails: mark orphaned, market-close, alert CRITICAL.
|
|
|
| Returns: count of orphans processed
|
| """
|
| candidates = storage.get_orphan_candidates()
|
| processed = 0
|
|
|
| for group in candidates:
|
| entry_id = group["entry_order_id"]
|
| symbol = group["symbol"]
|
| entry_price = group.get("entry_price", 0)
|
|
|
| logger.warning("Orphan candidate found: %s (entry %s)", symbol, entry_id)
|
|
|
|
|
| storage.update_linked_order_group(entry_id, {"orphaned": True})
|
| if alert_callback:
|
| alert_callback(
|
| f"π¨ ORPHANED: {symbol} entry {entry_id} missing stop/TP legs. Closing."
|
| )
|
| if not config.DRY_RUN:
|
| try:
|
| broker.close_position(symbol)
|
| except Exception as e:
|
| logger.critical("Failed to close orphan %s: %s", symbol, e)
|
| else:
|
| with _vp_lock:
|
| _virtual_portfolio.pop(symbol, None)
|
|
|
| processed += 1
|
|
|
| return processed
|
|
|
|
|
|
|
|
|
| def get_virtual_portfolio() -> dict:
|
| """Get the DRY_RUN virtual portfolio."""
|
| with _vp_lock:
|
| return dict(_virtual_portfolio)
|
|
|
|
|
| def close_virtual_position(symbol: str):
|
| """Close a virtual position (DRY_RUN mode)."""
|
| with _vp_lock:
|
| _virtual_portfolio.pop(symbol, None)
|
| logger.info("[DRY_RUN] Virtual position closed: %s", symbol)
|
|
|
|
|
| def get_virtual_orders() -> dict:
|
| with _vp_lock:
|
| return dict(_virtual_orders)
|
|
|