""" execution/portfolio.py — Position tracking, P&L, exposure, overnight flag. Reads from Alpaca API (live) or virtual_portfolio (DRY_RUN). """ from __future__ import annotations import datetime import logging from zoneinfo import ZoneInfo _ET = ZoneInfo("America/New_York") import config from execution import broker from execution.order_manager import get_virtual_portfolio logger = logging.getLogger("trading_system.portfolio") class Portfolio: """Unified portfolio interface for live and dry-run modes.""" def __init__(self): self._positions: dict[str, dict] = {} self._open_alpaca_orders: dict[str, dict] = {} self._reservations: dict[str, float] = {} # reservation_id -> notional self._daily_realized_pnl: float = 0.0 self._closed_trades: list[dict] = [] import threading self._lock = threading.RLock() def refresh(self): """Refresh positions and open orders from Alpaca or virtual portfolio.""" with self._lock: if config.DRY_RUN: self._positions = get_virtual_portfolio() from execution.order_manager import get_virtual_orders self._open_alpaca_orders = { oid: o for oid, o in get_virtual_orders().items() if o.get("status") not in ("filled", "canceled", "rejected", "expired") } else: try: raw_pos = broker.list_positions() self._positions = {} for p in raw_pos: sym = p.get("symbol", "") self._positions[sym] = { "symbol": sym, "qty": float(p.get("qty", 0)), "side": p.get("side", "long"), "entry_price": float(p.get("avg_entry_price", 0)), "current_price": float(p.get("current_price", 0)), "unrealized_pl": float(p.get("unrealized_pl", 0)), "unrealized_plpc": float(p.get("unrealized_plpc", 0)), "market_value": float(p.get("market_value", 0)), } raw_orders = broker.list_orders(status="open") self._open_alpaca_orders = {o.get("id"): o for o in raw_orders} # Also fetch recent closed/filled orders to clear stale reservations recent_closed = broker.list_orders(status="closed") closed_ids = {o.get("id") for o in recent_closed} # Reconcile reservations: if a reserved order is in closed_ids, release it (it's either filled or canceled) for oid in list(self._reservations.keys()): if oid in closed_ids: self._reservations.pop(oid, None) logger.info("Auto-released reservation for closed order %s", oid) except Exception as e: logger.error("Failed to refresh portfolio: %s", e) @property def positions(self) -> dict[str, dict]: return dict(self._positions) @property def open_count(self) -> int: return len(self._positions) @property def open_symbols(self) -> list[str]: return list(self._positions.keys()) def get_position(self, symbol: str) -> dict | None: return self._positions.get(symbol) def get_open_directions(self) -> dict[str, str]: """Get {symbol: "buy"|"sell"} for open positions.""" directions = {} for sym, pos in self._positions.items(): side = pos.get("side", "long") if side in ("long", "buy"): directions[sym] = "buy" else: directions[sym] = "sell" return directions def get_open_positions_with_dates(self) -> dict[str, dict]: """Get positions with open_date for PDT tracking.""" result = {} for sym, pos in self._positions.items(): result[sym] = { "side": "buy" if pos.get("side", "long") in ("long", "buy") else "sell", "open_date": pos.get("open_date", datetime.date.today().isoformat()), "qty": pos.get("qty", 0), } return result @property def total_unrealized_pnl(self) -> float: with self._lock: return sum(p.get("unrealized_pl", 0) for p in self._positions.values()) @property def total_market_value(self) -> float: with self._lock: return sum(abs(p.get("market_value", p.get("qty", 0) * p.get("current_price", 0))) for p in self._positions.values()) @property def total_exposure(self) -> float: """Calculates total $ exposure: Filled Positions + Open Alpaca Buy Orders + Local Uncommitted Reservations.""" with self._lock: pos_val = self.total_market_value alpaca_order_val = 0.0 for o in self._open_alpaca_orders.values(): if o.get("side", "") == "buy": qty = float(o.get("qty") or 0) limit_price = float(o.get("limit_price") or o.get("stop_price") or 0) notional = float(o.get("notional") or (qty * limit_price)) alpaca_order_val += notional local_val = 0.0 for oid, amt in self._reservations.items(): if oid not in self._open_alpaca_orders: local_val += amt return pos_val + alpaca_order_val + local_val def reserve_allocation(self, reservation_id: str, amount: float) -> bool: """Atomically reserve allocation for a new buy order. Idempotent.""" with self._lock: if reservation_id in self._reservations or reservation_id in self._open_alpaca_orders: return True if self.total_exposure + amount > config.ALLOCATED_CAPITAL: logger.warning( "Reservation rejected for %s: Exposure $%.2f + $%.2f > Limit $%.2f", reservation_id, self.total_exposure, amount, config.ALLOCATED_CAPITAL ) return False self._reservations[reservation_id] = amount logger.info("Reserved $%.2f for %s (Total Exposure: $%.2f)", amount, reservation_id, self.total_exposure) return True def commit_allocation(self, temp_id: str, alpaca_order_id: str): """Map a pre-submission temporary ID to a real Alpaca order ID.""" with self._lock: if temp_id in self._reservations: amount = self._reservations.pop(temp_id) self._reservations[alpaca_order_id] = amount logger.info("Committed reservation %s -> %s ($%.2f)", temp_id, alpaca_order_id, amount) def release_allocation(self, reservation_id: str): """Idempotently release a reservation (e.g., on rejection or cancellation).""" with self._lock: if reservation_id in self._reservations: amount = self._reservations.pop(reservation_id) logger.info("Released reservation %s ($%.2f)", reservation_id, amount) def record_close(self, symbol: str, realized_pnl: float, entry_price: float, exit_price: float, qty: float): """Record a closed position.""" self._daily_realized_pnl += realized_pnl self._closed_trades.append({ "symbol": symbol, "realized_pnl": realized_pnl, "entry_price": entry_price, "exit_price": exit_price, "qty": qty, "closed_at": datetime.datetime.now(datetime.timezone.utc).isoformat(), }) self._positions.pop(symbol, None) logger.info( "Closed %s: PnL=$%.2f (entry=%.2f, exit=%.2f, qty=%.2f)", symbol, realized_pnl, entry_price, exit_price, qty, ) @property def daily_realized_pnl(self) -> float: return self._daily_realized_pnl @property def total_pnl(self) -> float: return self._daily_realized_pnl + self.total_unrealized_pnl @property def closed_trades(self) -> list[dict]: return list(self._closed_trades) def has_overnight_risk(self) -> bool: """Check if any position is at risk of being held overnight.""" if not self._positions: return False now = datetime.datetime.now(_ET) return now.hour >= 15 and now.minute >= 30 def force_close_all(self, alert_callback=None) -> list[dict]: """Force close all positions (end-of-day or shutdown).""" closed = [] for symbol in list(self._positions.keys()): try: if config.DRY_RUN: pos = self._positions.pop(symbol, {}) pnl = (pos.get("current_price", 0) - pos.get("entry_price", 0)) * pos.get("qty", 0) closed.append({"symbol": symbol, "pnl": pnl, "reason": "end_of_day_forced_exit"}) logger.info("[DRY_RUN] Force closed %s: PnL=$%.2f", symbol, pnl) else: broker.close_position(symbol) closed.append({"symbol": symbol, "reason": "end_of_day_forced_exit"}) logger.info("Force closed %s: reason=end_of_day_forced_exit", symbol) except Exception as e: logger.error("Failed to force close %s: %s", symbol, e) if alert_callback and closed: summary = ", ".join(f"{c['symbol']}" for c in closed) alert_callback(f"📊 End-of-day close: {summary}") return closed def reset_daily(self): """Reset daily tracking.""" self._daily_realized_pnl = 0.0 self._closed_trades = []