Download execution/portfolio.py from raghava4u/Trading-Bot-M20: direct link, hf CLI and curl.
- Browser
- Download file 10.3 kB
-
https://huggingface.co/raghava4u/Trading-Bot-M20/resolve/main/execution/portfolio.py
- Command line
-
hf download hf://raghava4u/Trading-Bot-M20/execution/portfolio.py
-
curl -L -o portfolio.py https://huggingface.co/raghava4u/Trading-Bot-M20/resolve/main/execution/portfolio.py
10.3 kB
| """ | |
| 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) | |
| def positions(self) -> dict[str, dict]: | |
| return dict(self._positions) | |
| def open_count(self) -> int: | |
| return len(self._positions) | |
| 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 | |
| def total_unrealized_pnl(self) -> float: | |
| with self._lock: | |
| return sum(p.get("unrealized_pl", 0) for p in self._positions.values()) | |
| 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()) | |
| 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, | |
| ) | |
| def daily_realized_pnl(self) -> float: | |
| return self._daily_realized_pnl | |
| def total_pnl(self) -> float: | |
| return self._daily_realized_pnl + self.total_unrealized_pnl | |
| 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 = [] | |