Phase 1: Real-Time Market Monitoring System

COMPLETE: Real-time unusual activity detection for congressional tickers

New Database Model:
- MarketAlert: Stores unusual market activity alerts
  * Tracks volume spikes, price movements, volatility
  * JSON details field for flexible data storage
  * Severity scoring (1-10 scale)
  * Indexed for efficient queries by ticker/timestamp

New Modules:
- src/pote/monitoring/market_monitor.py: Core monitoring engine
  * get_congressional_watchlist(): Top 50 most-traded tickers
  * check_ticker(): Analyze single stock for unusual activity
  * scan_watchlist(): Batch analysis of multiple tickers
  * Detection logic:
    - Unusual volume (3x average)
    - Price spikes/drops (>5%)
    - High volatility (2x normal)
  * save_alerts(): Persist to database
  * get_recent_alerts(): Query historical alerts

- src/pote/monitoring/alert_manager.py: Alert formatting & filtering
  * format_alert_text(): Human-readable output
  * format_alert_html(): HTML email format
  * filter_alerts(): By severity, ticker, type
  * generate_summary_report(): Text/HTML reports

Scripts:
- scripts/monitor_market.py: CLI monitoring tool
  * Continuous monitoring mode (--interval)
  * One-time scan (--once)
  * Custom ticker lists or auto-detect congressional watchlist
  * Severity filtering (--min-severity)
  * Report generation and saving

Migrations:
- alembic/versions/f44014715b40_add_market_alerts_table.py

Documentation:
- docs/11_live_market_monitoring.md: Complete explanation
  * Why you can't track WHO is trading
  * What IS possible (timing analysis)
  * How hybrid monitoring works
  * Data sources and APIs

Usage:
  # Monitor congressional tickers (one-time scan)
  python scripts/monitor_market.py --once

  # Continuous monitoring (every 5 minutes)
  python scripts/monitor_market.py --interval 300

  # Monitor specific tickers
  python scripts/monitor_market.py --tickers NVDA,MSFT,AAPL --once

Next Steps (Phase 2):
- Disclosure correlation engine
- Timing advantage calculator
- Suspicious trade flagging
This commit is contained in:
ilia
2025-12-15 15:10:49 -05:00
parent 8ba9d7ffdd
commit cfaf38b0be
8 changed files with 1191 additions and 0 deletions
+48
View File
@@ -13,6 +13,7 @@ from sqlalchemy import (
ForeignKey,
Index,
Integer,
JSON,
String,
Text,
UniqueConstraint,
@@ -218,3 +219,50 @@ class MetricTrade(Base):
__table_args__ = (
UniqueConstraint("trade_id", "calc_date", "calc_version", name="uq_metrics_trade"),
)
class MarketAlert(Base):
"""
Real-time market activity alerts.
Tracks unusual volume, price movements, and other anomalies.
"""
__tablename__ = "market_alerts"
id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
ticker: Mapped[str] = mapped_column(String(20), nullable=False, index=True)
alert_type: Mapped[str] = mapped_column(
String(50), nullable=False
) # 'unusual_volume', 'price_spike', 'options_flow', etc.
timestamp: Mapped[datetime] = mapped_column(DateTime, nullable=False, index=True)
# Alert details (stored as JSON)
details: Mapped[dict | None] = mapped_column(JSON)
# Metrics at time of alert
price: Mapped[Decimal | None] = mapped_column(DECIMAL(15, 4))
volume: Mapped[int | None] = mapped_column(Integer)
change_pct: Mapped[Decimal | None] = mapped_column(
DECIMAL(10, 4)
) # Price change %
# Severity scoring
severity: Mapped[int | None] = mapped_column(Integer) # 1-10 scale
# Metadata
source: Mapped[str] = mapped_column(String(50), default="market_monitor")
created_at: Mapped[datetime] = mapped_column(
DateTime, default=lambda: datetime.now(timezone.utc)
)
# Indexes for efficient queries
__table_args__ = (
Index("ix_market_alerts_ticker_timestamp", "ticker", "timestamp"),
Index("ix_market_alerts_alert_type", "alert_type"),
)
def __repr__(self) -> str:
return (
f"<MarketAlert(ticker='{self.ticker}', type='{self.alert_type}', "
f"timestamp={self.timestamp}, severity={self.severity})>"
)
+10
View File
@@ -0,0 +1,10 @@
"""
Market monitoring module.
Real-time tracking of unusual market activity.
"""
from .market_monitor import MarketMonitor
from .alert_manager import AlertManager
__all__ = ["MarketMonitor", "AlertManager"]
+244
View File
@@ -0,0 +1,244 @@
"""
Alert management and notification system.
Handles alert filtering, formatting, and delivery.
"""
import logging
from datetime import datetime, timezone
from typing import Any
from sqlalchemy.orm import Session
from pote.db.models import MarketAlert
logger = logging.getLogger(__name__)
class AlertManager:
"""Manage and deliver market alerts."""
def __init__(self, session: Session):
"""Initialize alert manager."""
self.session = session
def format_alert_text(self, alert: MarketAlert) -> str:
"""
Format alert as human-readable text.
Args:
alert: MarketAlert object
Returns:
Formatted alert string
"""
emoji_map = {
"unusual_volume": "📊",
"price_spike": "🚀",
"price_drop": "📉",
"high_volatility": "",
"options_flow": "💰",
}
emoji = emoji_map.get(alert.alert_type, "🔔")
severity_stars = "" * min(alert.severity or 1, 5)
lines = [
f"{emoji} {alert.ticker} - {alert.alert_type.upper().replace('_', ' ')}",
f" Severity: {severity_stars} ({alert.severity}/10)",
f" Time: {alert.timestamp.strftime('%Y-%m-%d %H:%M:%S')}",
f" Price: ${float(alert.price):.2f}" if alert.price else "",
f" Volume: {alert.volume:,}" if alert.volume else "",
f" Change: {float(alert.change_pct):+.2f}%" if alert.change_pct else "",
]
# Add details
if alert.details:
lines.append(" Details:")
for key, value in alert.details.items():
if isinstance(value, (int, float)):
if "pct" in key.lower() or "change" in key.lower():
lines.append(f" {key}: {value:+.2f}%")
else:
lines.append(f" {key}: {value:,.2f}")
else:
lines.append(f" {key}: {value}")
return "\n".join(line for line in lines if line)
def format_alert_html(self, alert: MarketAlert) -> str:
"""
Format alert as HTML.
Args:
alert: MarketAlert object
Returns:
HTML formatted alert
"""
severity_class = "high" if (alert.severity or 0) >= 7 else "medium" if (alert.severity or 0) >= 4 else "low"
html = f"""
<div class="alert {severity_class}">
<h3>{alert.ticker} - {alert.alert_type.replace('_', ' ').title()}</h3>
<p class="timestamp">{alert.timestamp.strftime('%Y-%m-%d %H:%M:%S')}</p>
<p class="severity">Severity: {alert.severity}/10</p>
<div class="metrics">
<span>Price: ${float(alert.price):.2f}</span>
<span>Volume: {alert.volume:,}</span>
<span>Change: {float(alert.change_pct):+.2f}%</span>
</div>
</div>
"""
return html
def filter_alerts(
self,
alerts: list[MarketAlert],
min_severity: int = 5,
tickers: list[str] | None = None,
alert_types: list[str] | None = None,
) -> list[MarketAlert]:
"""
Filter alerts by criteria.
Args:
alerts: List of alerts
min_severity: Minimum severity threshold
tickers: Only include these tickers (None = all)
alert_types: Only include these types (None = all)
Returns:
Filtered list of alerts
"""
filtered = alerts
# Filter by severity
filtered = [a for a in filtered if (a.severity or 0) >= min_severity]
# Filter by ticker
if tickers:
ticker_set = set(t.upper() for t in tickers)
filtered = [a for a in filtered if a.ticker.upper() in ticker_set]
# Filter by alert type
if alert_types:
type_set = set(alert_types)
filtered = [a for a in filtered if a.alert_type in type_set]
return filtered
def generate_summary_report(
self, alerts: list[MarketAlert], format: str = "text"
) -> str:
"""
Generate summary report of alerts.
Args:
alerts: List of alerts
format: Output format ('text' or 'html')
Returns:
Formatted summary report
"""
if format == "html":
return self._generate_html_summary(alerts)
else:
return self._generate_text_summary(alerts)
def _generate_text_summary(self, alerts: list[MarketAlert]) -> str:
"""Generate text summary report."""
if not alerts:
return "📭 No alerts to report."
lines = [
"=" * 80,
f" MARKET ACTIVITY ALERTS - {datetime.now(timezone.utc).strftime('%Y-%m-%d %H:%M:%S')} UTC",
f" {len(alerts)} Alerts",
"=" * 80,
"",
]
# Group by ticker
by_ticker: dict[str, list[MarketAlert]] = {}
for alert in alerts:
if alert.ticker not in by_ticker:
by_ticker[alert.ticker] = []
by_ticker[alert.ticker].append(alert)
# Sort tickers by max severity
sorted_tickers = sorted(
by_ticker.keys(),
key=lambda t: max((a.severity or 0) for a in by_ticker[t]),
reverse=True,
)
for ticker in sorted_tickers:
ticker_alerts = by_ticker[ticker]
max_sev = max((a.severity or 0) for a in ticker_alerts)
lines.append("" * 80)
lines.append(f"🎯 {ticker} - {len(ticker_alerts)} alerts (Max Severity: {max_sev}/10)")
lines.append("" * 80)
for alert in sorted(
ticker_alerts, key=lambda a: a.severity or 0, reverse=True
):
lines.append("")
lines.append(self.format_alert_text(alert))
lines.append("")
# Summary statistics
lines.append("=" * 80)
lines.append("📊 SUMMARY")
lines.append("=" * 80)
lines.append("")
lines.append(f"Total Alerts: {len(alerts)}")
lines.append(f"Unique Tickers: {len(by_ticker)}")
# Alert type breakdown
type_counts: dict[str, int] = {}
for alert in alerts:
type_counts[alert.alert_type] = type_counts.get(alert.alert_type, 0) + 1
lines.append("\nAlert Types:")
for alert_type, count in sorted(
type_counts.items(), key=lambda x: x[1], reverse=True
):
lines.append(f" {alert_type.replace('_', ' ').title():20s}: {count}")
# Top severity alerts
lines.append("\nTop 5 Highest Severity:")
top_alerts = sorted(alerts, key=lambda a: a.severity or 0, reverse=True)[:5]
for alert in top_alerts:
lines.append(
f" {alert.ticker:6s} - {alert.alert_type:20s} (Severity: {alert.severity}/10)"
)
lines.append("")
lines.append("=" * 80)
return "\n".join(lines)
def _generate_html_summary(self, alerts: list[MarketAlert]) -> str:
"""Generate HTML summary report."""
html_parts = [
"<html><head><style>",
"body { font-family: Arial, sans-serif; }",
".alert { border: 1px solid #ddd; padding: 15px; margin: 10px 0; border-radius: 5px; }",
".alert.high { background-color: #ffebee; border-color: #f44336; }",
".alert.medium { background-color: #fff3e0; border-color: #ff9800; }",
".alert.low { background-color: #e8f5e9; border-color: #4caf50; }",
".timestamp { color: #666; font-size: 0.9em; }",
".metrics span { margin-right: 20px; }",
"</style></head><body>",
f"<h1>Market Activity Alerts</h1>",
f"<p><strong>{len(alerts)} Alerts</strong> | {datetime.now(timezone.utc).strftime('%Y-%m-%d %H:%M:%S')} UTC</p>",
]
for alert in sorted(alerts, key=lambda a: a.severity or 0, reverse=True):
html_parts.append(self.format_alert_html(alert))
html_parts.append("</body></html>")
return "\n".join(html_parts)
+281
View File
@@ -0,0 +1,281 @@
"""
Real-time market monitoring for congressional tickers.
Detects unusual activity: volume spikes, price movements, volatility.
"""
import logging
from datetime import datetime, timedelta, timezone
from decimal import Decimal
from typing import Any
import yfinance as yf
from sqlalchemy.orm import Session
from pote.db.models import MarketAlert, Security, Trade
logger = logging.getLogger(__name__)
class MarketMonitor:
"""Monitor stocks for unusual market activity."""
def __init__(self, session: Session):
"""Initialize market monitor."""
self.session = session
def get_congressional_watchlist(self, limit: int = 50) -> list[str]:
"""
Get list of most-traded tickers by Congress.
Args:
limit: Maximum number of tickers to return
Returns:
List of ticker symbols
"""
from sqlalchemy import func
result = (
self.session.query(Security.ticker, func.count(Trade.id).label("count"))
.join(Trade)
.group_by(Security.ticker)
.order_by(func.count(Trade.id).desc())
.limit(limit)
.all()
)
tickers = [r[0] for r in result]
logger.info(f"Built watchlist of {len(tickers)} tickers from congressional trades")
return tickers
def check_ticker(self, ticker: str, lookback_days: int = 5) -> list[dict[str, Any]]:
"""
Check a single ticker for unusual activity.
Args:
ticker: Stock ticker symbol
lookback_days: Days of history to analyze
Returns:
List of alerts detected
"""
alerts = []
try:
stock = yf.Ticker(ticker)
# Get recent history
hist = stock.history(period=f"{lookback_days}d", interval="1d")
if len(hist) < 2:
logger.warning(f"Insufficient data for {ticker}")
return alerts
# Calculate baseline metrics
avg_volume = hist["Volume"].mean()
avg_price_change = hist["Close"].pct_change().abs().mean()
# Get latest data
latest = hist.iloc[-1]
prev = hist.iloc[-2]
current_volume = latest["Volume"]
current_price = latest["Close"]
price_change = (current_price - prev["Close"]) / prev["Close"]
# Check for unusual volume (3x average)
if current_volume > avg_volume * 3 and avg_volume > 0:
severity = min(10, int((current_volume / avg_volume) - 2))
alerts.append(
{
"ticker": ticker,
"alert_type": "unusual_volume",
"timestamp": datetime.now(timezone.utc),
"details": {
"current_volume": int(current_volume),
"avg_volume": int(avg_volume),
"multiplier": round(current_volume / avg_volume, 2),
},
"price": Decimal(str(current_price)),
"volume": int(current_volume),
"change_pct": Decimal(str(price_change * 100)),
"severity": severity,
}
)
# Check for significant price movement (>5%)
if abs(price_change) > 0.05:
severity = min(10, int(abs(price_change) * 100 / 2))
alerts.append(
{
"ticker": ticker,
"alert_type": "price_spike"
if price_change > 0
else "price_drop",
"timestamp": datetime.now(timezone.utc),
"details": {
"current_price": float(current_price),
"prev_price": float(prev["Close"]),
"change_pct": round(price_change * 100, 2),
},
"price": Decimal(str(current_price)),
"volume": int(current_volume),
"change_pct": Decimal(str(price_change * 100)),
"severity": severity,
}
)
# Check for unusual volatility (price swings)
if len(hist) >= 5:
recent_volatility = hist["Close"].iloc[-5:].pct_change().abs().mean()
if recent_volatility > avg_price_change * 2 and avg_price_change > 0:
severity = min(
10, int((recent_volatility / avg_price_change) - 1)
)
alerts.append(
{
"ticker": ticker,
"alert_type": "high_volatility",
"timestamp": datetime.now(timezone.utc),
"details": {
"recent_volatility": round(recent_volatility * 100, 2),
"avg_volatility": round(avg_price_change * 100, 2),
"multiplier": round(recent_volatility / avg_price_change, 2),
},
"price": Decimal(str(current_price)),
"volume": int(current_volume),
"change_pct": Decimal(str(price_change * 100)),
"severity": severity,
}
)
except Exception as e:
logger.error(f"Error checking {ticker}: {e}")
return alerts
def scan_watchlist(
self, tickers: list[str] | None = None, lookback_days: int = 5
) -> list[dict[str, Any]]:
"""
Scan multiple tickers for unusual activity.
Args:
tickers: List of tickers to scan (None = use congressional watchlist)
lookback_days: Days of history to analyze
Returns:
List of all alerts detected
"""
if tickers is None:
tickers = self.get_congressional_watchlist()
all_alerts = []
logger.info(f"Scanning {len(tickers)} tickers for unusual activity...")
for ticker in tickers:
alerts = self.check_ticker(ticker, lookback_days=lookback_days)
all_alerts.extend(alerts)
if alerts:
logger.info(
f"🔔 {ticker}: {len(alerts)} alerts - "
+ ", ".join(a["alert_type"] for a in alerts)
)
logger.info(f"Scan complete. Found {len(all_alerts)} total alerts.")
return all_alerts
def save_alerts(self, alerts: list[dict[str, Any]]) -> int:
"""
Save alerts to database.
Args:
alerts: List of alert dictionaries
Returns:
Number of alerts saved
"""
saved = 0
for alert_data in alerts:
alert = MarketAlert(**alert_data)
self.session.add(alert)
saved += 1
self.session.commit()
logger.info(f"Saved {saved} alerts to database")
return saved
def get_recent_alerts(
self,
ticker: str | None = None,
days: int = 7,
alert_type: str | None = None,
min_severity: int = 0,
) -> list[MarketAlert]:
"""
Query recent alerts from database.
Args:
ticker: Filter by ticker (None = all)
days: Look back this many days
alert_type: Filter by alert type (None = all)
min_severity: Minimum severity level
Returns:
List of MarketAlert objects
"""
since = datetime.now(timezone.utc) - timedelta(days=days)
query = self.session.query(MarketAlert).filter(MarketAlert.timestamp >= since)
if ticker:
query = query.filter(MarketAlert.ticker == ticker)
if alert_type:
query = query.filter(MarketAlert.alert_type == alert_type)
if min_severity > 0:
query = query.filter(MarketAlert.severity >= min_severity)
return query.order_by(MarketAlert.timestamp.desc()).all()
def get_ticker_alert_summary(self, days: int = 30) -> dict[str, dict]:
"""
Get summary of alerts by ticker.
Args:
days: Look back this many days
Returns:
Dict mapping ticker to alert summary
"""
since = datetime.now(timezone.utc) - timedelta(days=days)
from sqlalchemy import func
results = (
self.session.query(
MarketAlert.ticker,
func.count(MarketAlert.id).label("alert_count"),
func.avg(MarketAlert.severity).label("avg_severity"),
func.max(MarketAlert.severity).label("max_severity"),
)
.filter(MarketAlert.timestamp >= since)
.group_by(MarketAlert.ticker)
.order_by(func.count(MarketAlert.id).desc())
.all()
)
summary = {}
for r in results:
summary[r[0]] = {
"alert_count": r[1],
"avg_severity": round(float(r[2]), 2) if r[2] else 0,
"max_severity": r[3],
}
return summary