- New ff_calendar table (event_date, time, currency, impact, actual, forecast, previous) - New service ff_calendar.py: bulk CSV import (83K events 2007-2025) + live sync from faireconomy.media JSON endpoint (this week / next week) - New API endpoints: POST /api/eco/ff-import, POST /api/eco/ff-sync, GET /api/eco/calendar (period filter), GET /api/eco/ff-stats - CalendarPage.tsx full rewrite: period tabs (Recent/Today/Tomorrow/This Week…), currency flags filter, impact filter, unified date-grouped table with Time·Flag·Currency·Impact·Event·Actual·Forecast·Previous columns, green/red actual vs forecast, TODAY badge, auto-refresh 60s Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
439 lines
17 KiB
Python
439 lines
17 KiB
Python
"""
|
|
Economic calendar router — FRED-backed historical events + Forex Factory calendar.
|
|
Prefix: /api/eco
|
|
"""
|
|
import json
|
|
import logging
|
|
import os
|
|
from typing import Any, Dict, List, Optional
|
|
|
|
from fastapi import APIRouter, BackgroundTasks, HTTPException, Query
|
|
|
|
logger = logging.getLogger(__name__)
|
|
router = APIRouter(prefix="/api/eco", tags=["Economic Calendar"])
|
|
|
|
|
|
# ── Helpers ───────────────────────────────────────────────────────────────────
|
|
|
|
def _parse_row(row: Dict) -> Dict:
|
|
for field in ("assets_impacted",):
|
|
val = row.get(field)
|
|
if isinstance(val, str):
|
|
try:
|
|
row[field] = json.loads(val) if val else []
|
|
except Exception:
|
|
row[field] = []
|
|
elif val is None:
|
|
row[field] = []
|
|
return row
|
|
|
|
|
|
# ── Bootstrap ─────────────────────────────────────────────────────────────────
|
|
|
|
_bootstrap_status: Dict[str, Any] = {"running": False, "last_result": None}
|
|
|
|
|
|
def _run_bootstrap(from_date: str, series_ids: Optional[List[str]], force: bool):
|
|
global _bootstrap_status
|
|
_bootstrap_status["running"] = True
|
|
try:
|
|
from services.fred_bootstrap import bootstrap_fred
|
|
result = bootstrap_fred(from_date=from_date, series_ids=series_ids, force=force)
|
|
_bootstrap_status["last_result"] = result
|
|
total_inserted = sum(v.get("inserted", 0) for v in result.values() if isinstance(v, dict))
|
|
logger.info(f"[eco/bootstrap] Done — {total_inserted} rows inserted")
|
|
except Exception as e:
|
|
logger.error(f"[eco/bootstrap] Failed: {e}")
|
|
_bootstrap_status["last_result"] = {"error": str(e)}
|
|
finally:
|
|
_bootstrap_status["running"] = False
|
|
|
|
|
|
@router.get("/fred-key")
|
|
def get_fred_key() -> Dict[str, Any]:
|
|
"""Check if a FRED API key is configured."""
|
|
from services.database import get_config
|
|
key = get_config("fred_api_key") or ""
|
|
return {
|
|
"configured": bool(key),
|
|
"preview": (key[:6] + "…" + key[-4:]) if len(key) > 10 else ("***" if key else ""),
|
|
}
|
|
|
|
|
|
@router.post("/fred-key")
|
|
def save_fred_key(key: str = Query(..., description="FRED API key")) -> Dict[str, Any]:
|
|
"""Save FRED API key to DB config."""
|
|
from services.database import set_config
|
|
key = key.strip()
|
|
if not key:
|
|
raise HTTPException(400, "Empty key")
|
|
set_config("fred_api_key", key)
|
|
return {"status": "saved", "preview": key[:6] + "…" + key[-4:]}
|
|
|
|
|
|
@router.post("/bootstrap")
|
|
def bootstrap(
|
|
background_tasks: BackgroundTasks,
|
|
from_date: str = Query("2020-01-01", description="Start date for historical data"),
|
|
series: Optional[str] = Query(None, description="Comma-separated series IDs (blank = all)"),
|
|
force: bool = Query(False, description="Overwrite existing rows"),
|
|
) -> Dict[str, Any]:
|
|
"""
|
|
Trigger FRED data bootstrap (runs in background).
|
|
Fetches public CSV endpoint — no API key required.
|
|
"""
|
|
if _bootstrap_status["running"]:
|
|
raise HTTPException(409, "Bootstrap already running")
|
|
series_ids = [s.strip() for s in series.split(",") if s.strip()] if series else None
|
|
background_tasks.add_task(_run_bootstrap, from_date, series_ids, force)
|
|
return {"status": "started", "from_date": from_date, "series": series_ids or "all"}
|
|
|
|
|
|
@router.get("/bootstrap/status")
|
|
def bootstrap_status() -> Dict[str, Any]:
|
|
return _bootstrap_status
|
|
|
|
|
|
# ── Series catalog ────────────────────────────────────────────────────────────
|
|
|
|
@router.get("/series")
|
|
def list_series() -> List[Dict[str, Any]]:
|
|
"""Return available FRED series metadata."""
|
|
from services.fred_bootstrap import FRED_SERIES, CATEGORIES
|
|
return [
|
|
{
|
|
"id": sid,
|
|
"name": meta["name"],
|
|
"category": meta["category"],
|
|
"freq": meta["freq"],
|
|
"unit": meta["unit"],
|
|
"assets": meta["assets"],
|
|
}
|
|
for sid, meta in FRED_SERIES.items()
|
|
]
|
|
|
|
|
|
# ── Events list ───────────────────────────────────────────────────────────────
|
|
|
|
_SORT_COLS = {
|
|
"date": "ee.event_date",
|
|
"zscore": "ABS(COALESCE(ee.surprise_zscore, 0))",
|
|
"series": "ee.series_id",
|
|
"name": "ee.event_name",
|
|
}
|
|
|
|
|
|
@router.get("/events")
|
|
def list_eco_events(
|
|
date_from: Optional[str] = Query(None),
|
|
date_to: Optional[str] = Query(None),
|
|
series: Optional[str] = Query(None, description="Comma-separated series IDs"),
|
|
category: Optional[str] = Query(None, description="employment|inflation|growth|monetary|credit|rates"),
|
|
min_zscore: float = Query(0.0, ge=0.0, description="Minimum |z-score| filter"),
|
|
direction: Optional[str] = Query(None, description="bullish|bearish|neutral"),
|
|
sort_by: str = Query("date", description="date|zscore|series|name"),
|
|
sort_dir: str = Query("desc", description="asc|desc"),
|
|
limit: int = Query(200, ge=1, le=2000),
|
|
offset: int = Query(0, ge=0),
|
|
) -> Dict[str, Any]:
|
|
from services.database import get_conn
|
|
from services.fred_bootstrap import FRED_SERIES
|
|
|
|
conn = get_conn()
|
|
where_parts: List[str] = []
|
|
params: List[Any] = []
|
|
|
|
if date_from:
|
|
where_parts.append("ee.event_date >= ?")
|
|
params.append(date_from)
|
|
if date_to:
|
|
where_parts.append("ee.event_date <= ?")
|
|
params.append(date_to)
|
|
|
|
if series:
|
|
ids = [s.strip() for s in series.split(",") if s.strip()]
|
|
if ids:
|
|
placeholders = ",".join("?" * len(ids))
|
|
where_parts.append(f"ee.series_id IN ({placeholders})")
|
|
params.extend(ids)
|
|
|
|
if category:
|
|
# Map category → series IDs
|
|
matching = [sid for sid, m in FRED_SERIES.items() if m["category"] == category]
|
|
if matching:
|
|
placeholders = ",".join("?" * len(matching))
|
|
where_parts.append(f"ee.series_id IN ({placeholders})")
|
|
params.extend(matching)
|
|
else:
|
|
# No matching series — return empty
|
|
conn.close()
|
|
return {"total": 0, "offset": offset, "limit": limit, "events": []}
|
|
|
|
if min_zscore > 0:
|
|
where_parts.append("ABS(COALESCE(ee.surprise_zscore, 0)) >= ?")
|
|
params.append(min_zscore)
|
|
if direction:
|
|
where_parts.append("ee.surprise_direction = ?")
|
|
params.append(direction)
|
|
|
|
where_sql = ("WHERE " + " AND ".join(where_parts)) if where_parts else ""
|
|
sort_col = _SORT_COLS.get(sort_by, "ee.event_date")
|
|
order_sql = f"ORDER BY {sort_col} {'DESC' if sort_dir == 'desc' else 'ASC'}"
|
|
|
|
count_sql = f"SELECT COUNT(*) FROM economic_events ee {where_sql}"
|
|
data_sql = f"SELECT ee.* FROM economic_events ee {where_sql} {order_sql} LIMIT ? OFFSET ?"
|
|
|
|
try:
|
|
total = conn.execute(count_sql, params).fetchone()[0]
|
|
rows = conn.execute(data_sql, params + [limit, offset]).fetchall()
|
|
finally:
|
|
conn.close()
|
|
|
|
events = [_parse_row(dict(r)) for r in rows]
|
|
return {"total": total, "offset": offset, "limit": limit, "events": events}
|
|
|
|
|
|
# ── Upcoming releases ────────────────────────────────────────────────────────
|
|
|
|
# Typical publication lag (days) after the end of the reference period
|
|
_RELEASE_LAG: Dict[str, Dict[str, Any]] = {
|
|
"PAYEMS": {"freq_days": 30, "lag_days": 35, "note": "1er vendredi du mois"},
|
|
"UNRATE": {"freq_days": 30, "lag_days": 35, "note": "Avec NFP"},
|
|
"CPIAUCSL": {"freq_days": 30, "lag_days": 15, "note": "Mi-mois"},
|
|
"CPILFESL": {"freq_days": 30, "lag_days": 15, "note": "Avec CPI"},
|
|
"PCEPILFE": {"freq_days": 30, "lag_days": 30, "note": "Fin de mois"},
|
|
"FEDFUNDS": {"freq_days": 45, "lag_days": 45, "note": "Réunion FOMC (~8/an)"},
|
|
"ICSA": {"freq_days": 7, "lag_days": 4, "note": "Chaque jeudi"},
|
|
"GDPC1": {"freq_days": 90, "lag_days": 30, "note": "Estimation avancée BEA"},
|
|
"BAMLH0A0HYM2": {"freq_days": 7, "lag_days": 1, "note": "Hebdomadaire"},
|
|
"T10Y2Y": {"freq_days": 7, "lag_days": 1, "note": "Hebdomadaire"},
|
|
"T10Y3M": {"freq_days": 7, "lag_days": 1, "note": "Hebdomadaire"},
|
|
}
|
|
|
|
|
|
@router.get("/upcoming")
|
|
def upcoming_events() -> List[Dict[str, Any]]:
|
|
"""Estimate next release dates for each series based on last stored date + typical lag."""
|
|
from services.database import get_conn
|
|
from services.fred_bootstrap import FRED_SERIES
|
|
from datetime import date, timedelta
|
|
|
|
conn = get_conn()
|
|
today = date.today()
|
|
upcoming = []
|
|
|
|
for sid, lag_info in _RELEASE_LAG.items():
|
|
meta = FRED_SERIES.get(sid)
|
|
if not meta:
|
|
continue
|
|
try:
|
|
row = conn.execute(
|
|
"SELECT MAX(event_date), actual_value, actual_unit, surprise_direction "
|
|
"FROM economic_events WHERE series_id=?",
|
|
(sid,),
|
|
).fetchone()
|
|
if not row or not row[0]:
|
|
# No data in DB yet — still show as "not loaded"
|
|
upcoming.append({
|
|
"series_id": sid,
|
|
"name": meta["name"],
|
|
"note": lag_info["note"],
|
|
"category": meta["category"],
|
|
"last_release": None,
|
|
"last_value": None,
|
|
"last_unit": meta["unit"],
|
|
"last_direction": None,
|
|
"next_expected": None,
|
|
"days_until": None,
|
|
"status": "no_data",
|
|
})
|
|
continue
|
|
|
|
last_date = date.fromisoformat(row[0])
|
|
freq_days = lag_info["freq_days"]
|
|
lag_days = lag_info["lag_days"]
|
|
|
|
# Next expected = last ref period end + lag
|
|
next_release = last_date + timedelta(days=freq_days + lag_days)
|
|
|
|
days_until = (next_release - today).days
|
|
|
|
if days_until < -freq_days:
|
|
# Very overdue — likely bootstrap missing recent data
|
|
status = "overdue"
|
|
elif days_until < 0:
|
|
status = "due_soon" # probably released, not yet in DB
|
|
elif days_until <= 7:
|
|
status = "imminent"
|
|
elif days_until <= 30:
|
|
status = "upcoming"
|
|
else:
|
|
status = "scheduled"
|
|
|
|
upcoming.append({
|
|
"series_id": sid,
|
|
"name": meta["name"],
|
|
"note": lag_info["note"],
|
|
"category": meta["category"],
|
|
"last_release": str(last_date),
|
|
"last_value": round(float(row[1]), 3) if row[1] is not None else None,
|
|
"last_unit": row[2] or meta["unit"],
|
|
"last_direction": row[3],
|
|
"next_expected": str(next_release),
|
|
"days_until": days_until,
|
|
"status": status,
|
|
})
|
|
except Exception as e:
|
|
logger.debug(f"[eco/upcoming] {sid}: {e}")
|
|
|
|
conn.close()
|
|
upcoming.sort(key=lambda x: (x["next_expected"] or "9999", x["series_id"]))
|
|
return upcoming
|
|
|
|
|
|
# ── Forex Factory calendar ────────────────────────────────────────────────────
|
|
|
|
_ff_import_status: Dict[str, Any] = {"running": False, "last_result": None}
|
|
_ff_sync_status: Dict[str, Any] = {"running": False, "last_result": None}
|
|
|
|
_FF_CSV_PATH = os.path.join(
|
|
os.path.dirname(__file__), "..", "..", "forex_factory_cache.csv"
|
|
)
|
|
|
|
|
|
def _run_ff_import():
|
|
global _ff_import_status
|
|
_ff_import_status["running"] = True
|
|
try:
|
|
from services.ff_calendar import import_csv
|
|
csv_path = os.path.abspath(_FF_CSV_PATH)
|
|
result = import_csv(csv_path)
|
|
_ff_import_status["last_result"] = result
|
|
except Exception as e:
|
|
logger.error(f"[eco/ff-import] Failed: {e}")
|
|
_ff_import_status["last_result"] = {"error": str(e)}
|
|
finally:
|
|
_ff_import_status["running"] = False
|
|
|
|
|
|
def _run_ff_sync():
|
|
global _ff_sync_status
|
|
_ff_sync_status["running"] = True
|
|
try:
|
|
from services.ff_calendar import sync_live
|
|
result = sync_live()
|
|
_ff_sync_status["last_result"] = result
|
|
except Exception as e:
|
|
logger.error(f"[eco/ff-sync] Failed: {e}")
|
|
_ff_sync_status["last_result"] = {"error": str(e)}
|
|
finally:
|
|
_ff_sync_status["running"] = False
|
|
|
|
|
|
@router.post("/ff-import")
|
|
def ff_import(background_tasks: BackgroundTasks) -> Dict[str, Any]:
|
|
"""Import forex_factory_cache.csv into ff_calendar table (background)."""
|
|
if _ff_import_status["running"]:
|
|
raise HTTPException(409, "Import already running")
|
|
csv_path = os.path.abspath(_FF_CSV_PATH)
|
|
if not os.path.exists(csv_path):
|
|
raise HTTPException(404, f"CSV not found: {csv_path}")
|
|
background_tasks.add_task(_run_ff_import)
|
|
return {"status": "started", "csv": csv_path}
|
|
|
|
|
|
@router.get("/ff-import/status")
|
|
def ff_import_status() -> Dict[str, Any]:
|
|
return _ff_import_status
|
|
|
|
|
|
@router.post("/ff-sync")
|
|
def ff_sync(background_tasks: BackgroundTasks) -> Dict[str, Any]:
|
|
"""Fetch this week + next week from faireconomy.media and upsert."""
|
|
if _ff_sync_status["running"]:
|
|
raise HTTPException(409, "Sync already running")
|
|
background_tasks.add_task(_run_ff_sync)
|
|
return {"status": "started"}
|
|
|
|
|
|
@router.get("/ff-sync/status")
|
|
def ff_sync_status_ep() -> Dict[str, Any]:
|
|
return _ff_sync_status
|
|
|
|
|
|
@router.get("/calendar")
|
|
def ff_calendar(
|
|
period: str = Query("recent", description="recent|today|tomorrow|yesterday|this_week|next_week|previous_week|this_month|next_month|previous_month"),
|
|
currencies: Optional[str] = Query(None, description="Comma-separated: USD,EUR,GBP,..."),
|
|
impacts: Optional[str] = Query(None, description="Comma-separated: high,medium,low"),
|
|
limit: int = Query(300, ge=1, le=2000),
|
|
offset: int = Query(0, ge=0),
|
|
) -> Dict[str, Any]:
|
|
"""Unified Forex Factory calendar — past releases + upcoming events."""
|
|
from services.ff_calendar import get_calendar
|
|
cur_list = [c.strip().upper() for c in currencies.split(",") if c.strip()] if currencies else None
|
|
imp_list = [i.strip().lower() for i in impacts.split(",") if i.strip()] if impacts else None
|
|
return get_calendar(
|
|
period=period,
|
|
currencies=cur_list,
|
|
impacts=imp_list,
|
|
limit=limit,
|
|
offset=offset,
|
|
)
|
|
|
|
|
|
@router.get("/ff-stats")
|
|
def ff_stats() -> Dict[str, Any]:
|
|
"""Quick inventory of ff_calendar table."""
|
|
from services.database import get_conn
|
|
conn = get_conn()
|
|
try:
|
|
total = conn.execute("SELECT COUNT(*) FROM ff_calendar").fetchone()[0]
|
|
if total == 0:
|
|
return {"total": 0}
|
|
earliest = conn.execute("SELECT MIN(event_date) FROM ff_calendar").fetchone()[0]
|
|
latest = conn.execute("SELECT MAX(event_date) FROM ff_calendar").fetchone()[0]
|
|
by_currency = conn.execute(
|
|
"SELECT currency, COUNT(*) AS cnt FROM ff_calendar GROUP BY currency ORDER BY cnt DESC"
|
|
).fetchall()
|
|
by_impact = conn.execute(
|
|
"SELECT impact, COUNT(*) AS cnt FROM ff_calendar GROUP BY impact ORDER BY cnt DESC"
|
|
).fetchall()
|
|
return {
|
|
"total": total,
|
|
"earliest": earliest,
|
|
"latest": latest,
|
|
"by_currency": [dict(r) for r in by_currency],
|
|
"by_impact": [dict(r) for r in by_impact],
|
|
}
|
|
finally:
|
|
conn.close()
|
|
|
|
|
|
# ── DB count summary ──────────────────────────────────────────────────────────
|
|
|
|
@router.get("/status")
|
|
def eco_status() -> Dict[str, Any]:
|
|
from services.database import get_conn
|
|
conn = get_conn()
|
|
try:
|
|
total = conn.execute("SELECT COUNT(*) FROM economic_events").fetchone()[0]
|
|
latest = conn.execute(
|
|
"SELECT MAX(event_date) FROM economic_events"
|
|
).fetchone()[0]
|
|
earliest = conn.execute(
|
|
"SELECT MIN(event_date) FROM economic_events"
|
|
).fetchone()[0]
|
|
by_series = conn.execute(
|
|
"SELECT series_id, COUNT(*) as cnt, MAX(event_date) as latest "
|
|
"FROM economic_events GROUP BY series_id ORDER BY series_id"
|
|
).fetchall()
|
|
return {
|
|
"total": total,
|
|
"earliest": earliest,
|
|
"latest": latest,
|
|
"by_series": [dict(r) for r in by_series],
|
|
}
|
|
finally:
|
|
conn.close()
|