Trading-Bot-M20 / execution /portfolio.py
raghava4u's picture
Upload folder using huggingface_hub
d53dc44 verified
Raw History Blame Contribute Delete
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)
@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 = []