"""Fundamental data service. Stores fundamental data (P/E, revenue growth, earnings surprise, market cap) and marks the fundamental dimension score as stale on new data. """ from __future__ import annotations import json import logging from datetime import datetime, timezone from sqlalchemy import select, update from sqlalchemy.ext.asyncio import AsyncSession from app.database import insert_for_session from app.exceptions import NotFoundError from app.models.fundamental import FundamentalData from app.models.score import DimensionScore from app.models.ticker import Ticker logger = logging.getLogger(__name__) async def _get_ticker(db: AsyncSession, symbol: str) -> Ticker: """Look up a ticker by symbol.""" normalised = symbol.strip().upper() result = await db.execute(select(Ticker).where(Ticker.symbol == normalised)) ticker = result.scalar_one_or_none() if ticker is None: raise NotFoundError(f"Ticker not found: {normalised}") return ticker async def store_fundamental( db: AsyncSession, symbol: str, pe_ratio: float | None = None, revenue_growth: float | None = None, earnings_surprise: float | None = None, market_cap: float | None = None, next_earnings_date=None, unavailable_fields: dict[str, str] | None = None, ) -> FundamentalData: """Store or update fundamental data for a ticker. Keeps a single latest snapshot per ticker. On new data, marks the fundamental dimension score as stale (if one exists). """ ticker = await _get_ticker(db, symbol) now = datetime.now(timezone.utc) unavailable_fields_json = json.dumps(unavailable_fields or {}) stmt = insert_for_session(db, FundamentalData).values( ticker_id=ticker.id, pe_ratio=pe_ratio, revenue_growth=revenue_growth, earnings_surprise=earnings_surprise, market_cap=market_cap, next_earnings_date=next_earnings_date, fetched_at=now, unavailable_fields_json=unavailable_fields_json, ) stmt = stmt.on_conflict_do_update( index_elements=["ticker_id"], set_={ "pe_ratio": stmt.excluded.pe_ratio, "revenue_growth": stmt.excluded.revenue_growth, "earnings_surprise": stmt.excluded.earnings_surprise, "market_cap": stmt.excluded.market_cap, "next_earnings_date": stmt.excluded.next_earnings_date, "fetched_at": stmt.excluded.fetched_at, "unavailable_fields_json": stmt.excluded.unavailable_fields_json, }, ).returning(FundamentalData) record = (await db.execute(stmt)).scalar_one() # Mark fundamental dimension score as stale if it exists # TODO: Use DimensionScore service when built await db.execute( update(DimensionScore) .where( DimensionScore.ticker_id == ticker.id, DimensionScore.dimension == "fundamental", ) .values(is_stale=True) ) await db.commit() return record async def get_fundamental( db: AsyncSession, symbol: str, ) -> FundamentalData | None: """Get the latest fundamental data for a ticker.""" ticker = await _get_ticker(db, symbol) result = await db.execute( select(FundamentalData).where(FundamentalData.ticker_id == ticker.id) ) return result.scalar_one_or_none()