Download execution/order_manager.py from raghava4u/Trading-Bot-M20: direct link, hf CLI and curl.
- Browser
- Download file 18.4 kB
-
https://huggingface.co/raghava4u/Trading-Bot-M20/resolve/main/execution/order_manager.py
- Command line
-
hf download hf://raghava4u/Trading-Bot-M20/execution/order_manager.py
-
curl -L -o order_manager.py https://huggingface.co/raghava4u/Trading-Bot-M20/resolve/main/execution/order_manager.py
18.4 kB
| """ | |
| 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") | |
| # ββ DRY RUN simulation ββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| _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)) # Shortened for responsiveness, logged as original | |
| 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 | |
| # ββ Spread Check βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| 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 | |
| # Paper trading fallback: Alpaca free-tier quotes can be stale/venue-specific. | |
| # If spread looks unreasonably wide, use last trade price as reference. | |
| 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: | |
| # Last trade is between bid/ask β quote is stale, stock is liquid | |
| 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, | |
| ) | |
| # Return last trade as synthetic bid/ask with tiny spread | |
| 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) | |
| # ββ Fractional Order Flow (Atomic LinkedOrderGroup) ββββββββββββββββββββββββββ | |
| 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 | |
| # Step 1: Pre-create linked group | |
| group = LinkedOrderGroup( | |
| symbol=order.symbol, | |
| entry_order_id=entry_order_id or "pending", | |
| is_fractional=True, | |
| orphaned=False, | |
| ) | |
| # Step 2: Submit entry | |
| 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 | |
| # Store group in DB with real order ID | |
| 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) | |
| # For DRY_RUN, simulate immediate fill | |
| 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-loss | |
| 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 | |
| # Take-profit | |
| 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." | |
| ) | |
| # Cancel the stop first | |
| 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) | |
| # ββ Whole Order (Bracket) βββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| 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 | |
| # ββ Main Order Submission ββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| 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: | |
| # Get the correct ID depending on fractional/bracket response structure | |
| oid = result.get("id") or result.get("entry_order_id") | |
| if oid and getattr(order, "type", "limit") == "limit": | |
| track_pending_order(oid, order) | |
| # Actually, even market orders might pend briefly. Let's track all to be safe for orphans/timeouts. | |
| elif oid: | |
| track_pending_order(oid, order) | |
| return result | |
| except Exception as e: | |
| logger.error("Order submission failed: %s", e) | |
| return None | |
| # ββ Order Timeout Monitor ββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| _pending_orders: dict[str, dict] = {} # order_id β {submitted_at, order} | |
| _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 partially filled, Alpaca cancels the rest. The position is handled by `refresh()`. | |
| # But we explicitly release any local reservation here. | |
| 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 | |
| # ββ Orphan Monitor βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| 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) | |
| # We don't have original order details, so mark as orphaned and close | |
| 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 | |
| # ββ Helpers ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| 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) | |