""" Specialist Desks — per asset-class fundamental configs + report catalogue. Prefix: /api/specialist-desks """ import uuid import logging from typing import Optional, List from fastapi import APIRouter, HTTPException from pydantic import BaseModel logger = logging.getLogger(__name__) router = APIRouter(prefix="/api/specialist-desks", tags=["specialist-desks"]) # ── Pydantic models ──────────────────────────────────────────────────────────── class MacroSensitivity(BaseModel): regime: str effect: str class PriceDeltaThresholds(BaseModel): significant_day: float = 1.0 extreme_day: float = 3.0 significant_week: float = 2.0 extreme_week: float = 5.0 class DeskConfigUpdate(BaseModel): display_name: str icon: str = "" fundamentals: str = "" macro_sensitivity: List[MacroSensitivity] = [] price_delta_thresholds: PriceDeltaThresholds = PriceDeltaThresholds() notes: str = "" class ReportUpsert(BaseModel): name: str source: str = "" cadence: str = "monthly" next_date: Optional[str] = None last_date: Optional[str] = None importance: int = 2 notes: str = "" url: str = "" desks: List[str] = [] class ReportResultUpdate(BaseModel): actual_value: Optional[float] = None consensus_estimate: Optional[float] = None text_sentiment_score: Optional[float] = None text_sentiment_label: Optional[str] = None text_sentiment_summary: Optional[str] = None class ScoreTextRequest(BaseModel): text: str report_name: str = "" desk: str = "forex" # ── Configs ──────────────────────────────────────────────────────────────────── @router.get("/configs") def list_configs(): from services.database import get_all_asset_class_configs return get_all_asset_class_configs() @router.get("/configs/{asset_class}") def get_config(asset_class: str): from services.database import get_asset_class_config, get_desk_reports cfg = get_asset_class_config(asset_class) if not cfg: raise HTTPException(404, f"Desk '{asset_class}' not found") cfg["reports"] = get_desk_reports(asset_class) return cfg @router.put("/configs/{asset_class}") def update_config(asset_class: str, body: DeskConfigUpdate): from services.database import upsert_asset_class_config upsert_asset_class_config( asset_class=asset_class, display_name=body.display_name, icon=body.icon, fundamentals=body.fundamentals, macro_sensitivity=[s.dict() for s in body.macro_sensitivity], price_delta_thresholds=body.price_delta_thresholds.dict(), notes=body.notes, ) return {"saved": asset_class} # ── Reports ──────────────────────────────────────────────────────────────────── @router.get("/reports") def list_reports(): from services.database import get_all_specialist_reports return get_all_specialist_reports() @router.post("/reports") def create_report(body: ReportUpsert): from services.database import upsert_specialist_report, link_report_desk report_id = "RPT_" + uuid.uuid4().hex[:8].upper() upsert_specialist_report( report_id=report_id, name=body.name, source=body.source, cadence=body.cadence, next_date=body.next_date, last_date=body.last_date, importance=body.importance, notes=body.notes, url=body.url, ) for ac in body.desks: link_report_desk(report_id, ac) return {"id": report_id} @router.put("/reports/{report_id}") def update_report(report_id: str, body: ReportUpsert): from services.database import upsert_specialist_report, link_report_desk, unlink_report_desk, get_all_specialist_reports # Check exists all_reports = get_all_specialist_reports() if not any(r["id"] == report_id for r in all_reports): raise HTTPException(404, "Report not found") upsert_specialist_report( report_id=report_id, name=body.name, source=body.source, cadence=body.cadence, next_date=body.next_date, last_date=body.last_date, importance=body.importance, notes=body.notes, url=body.url, ) # Sync desk links: remove all, re-add existing = next(r for r in all_reports if r["id"] == report_id) for old_ac in existing.get("desks", []): unlink_report_desk(report_id, old_ac) for new_ac in body.desks: link_report_desk(report_id, new_ac) return {"saved": report_id} @router.delete("/reports/{report_id}") def delete_report(report_id: str): from services.database import delete_specialist_report delete_specialist_report(report_id) return {"deleted": report_id} @router.post("/reports/{report_id}/link/{asset_class}") def link_to_desk(report_id: str, asset_class: str): from services.database import link_report_desk link_report_desk(report_id, asset_class) return {"linked": report_id, "desk": asset_class} @router.delete("/reports/{report_id}/link/{asset_class}") def unlink_from_desk(report_id: str, asset_class: str): from services.database import unlink_report_desk unlink_report_desk(report_id, asset_class) return {"unlinked": report_id, "desk": asset_class} # ── Surprise Index ───────────────────────────────────────────────────────────── @router.put("/reports/{report_id}/result") def update_report_result(report_id: str, body: ReportResultUpdate): from services.database import update_report_result as _upd _upd( report_id=report_id, actual_value=body.actual_value, consensus_estimate=body.consensus_estimate, text_sentiment_score=body.text_sentiment_score, text_sentiment_label=body.text_sentiment_label, text_sentiment_summary=body.text_sentiment_summary, ) return {"saved": report_id} # ── COT Data ─────────────────────────────────────────────────────────────────── @router.get("/cot") def get_cot_data(): from services.database import get_latest_cot_data return get_latest_cot_data() @router.post("/cot/fetch") def fetch_cot(): from services.cot_fetcher import fetch_all_cot from services.database import save_cot_data try: data = fetch_all_cot() if data: saved = save_cot_data(data) return {"fetched": len(data), "saved": saved} return {"fetched": 0, "saved": 0, "warning": "No COT data returned"} except Exception as e: logger.error(f"COT fetch error: {e}") raise HTTPException(500, str(e)) # ── Forward Curves ───────────────────────────────────────────────────────────── @router.get("/forward-curves") def get_forward_curves(): from services.database import get_latest_forward_curves return get_latest_forward_curves() @router.post("/forward-curves/fetch") def fetch_forward_curves_endpoint(): from services.forward_curve import fetch_forward_curves from services.database import save_forward_curves try: data = fetch_forward_curves() if data: saved = save_forward_curves(data) return {"fetched": len(data), "saved": saved} return {"fetched": 0, "saved": 0, "warning": "No curve data returned"} except Exception as e: logger.error(f"Forward curve fetch error: {e}") raise HTTPException(500, str(e)) # ── Hawk / Dove Text Scorer ──────────────────────────────────────────────────── @router.post("/score-text") def score_text(body: ScoreTextRequest): """Score a central bank / commodity report excerpt via GPT-4o-mini.""" try: from services.ai_analyzer import score_report_text result = score_report_text(body.text, body.report_name, body.desk) return result except Exception as e: logger.error(f"Text scoring error: {e}") raise HTTPException(500, str(e)) # ── Bank Forecasts ───────────────────────────────────────────────────────────── class BankSourceUpsert(BaseModel): name: str url: str = "" active: bool = True notes: str = "" @router.get("/bank-forecasts/sources") def list_bank_sources(): from services.database import get_conn conn = get_conn() try: rows = conn.execute( "SELECT id, name, url, active, last_scraped, notes FROM bank_forecast_sources ORDER BY name" ).fetchall() return [dict(r) for r in rows] finally: conn.close() @router.put("/bank-forecasts/sources/{source_id}") def upsert_bank_source(source_id: str, body: BankSourceUpsert): from services.database import get_conn conn = get_conn() try: conn.execute( """INSERT INTO bank_forecast_sources (id, name, url, active, notes) VALUES (?,?,?,?,?) ON CONFLICT(id) DO UPDATE SET name=excluded.name, url=excluded.url, active=excluded.active, notes=excluded.notes""", (source_id, body.name, body.url, int(body.active), body.notes) ) conn.commit() return {"saved": source_id} finally: conn.close() @router.get("/bank-forecasts") def list_bank_forecasts(series_id: str = "", event_date: str = ""): from services.database import get_conn conn = get_conn() try: wheres = ["1=1"] params = [] if series_id: wheres.append("bf.series_id = ?"); params.append(series_id) if event_date: wheres.append("bf.event_date = ?"); params.append(event_date) rows = conn.execute( f"""SELECT bf.*, bs.name as bank_name FROM bank_forecasts bf JOIN bank_forecast_sources bs ON bs.id = bf.source_id WHERE {' AND '.join(wheres)} ORDER BY bf.event_date DESC, bf.extracted_at DESC LIMIT 500""", params ).fetchall() return [dict(r) for r in rows] finally: conn.close() @router.get("/bank-forecasts/consensus") def get_consensus(): """Return one consensus row per (series_id, event_date) = average of bank forecasts.""" from services.database import get_conn conn = get_conn() try: rows = conn.execute( """SELECT series_id, event_name, event_date, ROUND(AVG(forecast_value), 4) as consensus, COUNT(*) as bank_count, MIN(forecast_value) as min_forecast, MAX(forecast_value) as max_forecast FROM bank_forecasts WHERE forecast_value IS NOT NULL GROUP BY series_id, event_date ORDER BY event_date DESC""" ).fetchall() return [dict(r) for r in rows] finally: conn.close() @router.post("/bank-forecasts/scrape") def trigger_scrape(source_id: str = ""): """Scrape one source (source_id) or all active sources (empty).""" from services.database import get_conn from services.bank_forecast_scraper import scrape_source, scrape_all_active conn = get_conn() try: if source_id: s = conn.execute("SELECT * FROM bank_forecast_sources WHERE id=?", (source_id,)).fetchone() if not s: raise HTTPException(404, "Source not found") result = [scrape_source(conn, dict(s))] else: result = scrape_all_active(conn) return {"results": result, "total_sources": len(result)} finally: conn.close() @router.post("/bank-forecasts/push-consensus") def push_consensus_endpoint(series_id: str = "", event_date: str = ""): """Push bank consensus to macro_series_log. Filters optional.""" from services.database import get_conn from services.bank_forecast_scraper import push_consensus_to_log, push_all_consensus conn = get_conn() try: if series_id and event_date: val = push_consensus_to_log(conn, series_id, event_date) return {"pushed": [{"series_id": series_id, "event_date": event_date, "consensus": val}]} else: results = push_all_consensus(conn) return {"pushed": results} finally: conn.close()