From ab892745fa68711fcacdfd4b9bbf8e1003719d61 Mon Sep 17 00:00:00 2001 From: Jeroen Schweitzer Date: Fri, 26 Dec 2025 13:40:40 +0100 Subject: [PATCH] feat: add volatile fetch endpoints for scheduler-driven prefetch MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Add VolatileFetchService to orchestrate API fetch and cache storage - Add POST /volatile/fetch/weather/{city} endpoint - Add POST /volatile/fetch/news/{category} endpoint - Add POST /volatile/fetch/stock/{symbol} endpoint - Add POST /volatile/fetch/crypto/{symbol} endpoint Endpoints integrate with external API providers (OpenMeteo, NOS/BBC, AlphaVantage) and store results in volatile cache with configurable TTL. Designed for scheduler cron jobs to prefetch user-relevant data. 🤖 Generated with [Claude Code](https://claude.com/claude-code) Co-Authored-By: Claude Opus 4.5 --- src/routers/volatile.py | 192 ++++++++++++- src/services/volatile_fetch_service.py | 359 +++++++++++++++++++++++++ 2 files changed, 550 insertions(+), 1 deletion(-) create mode 100644 src/services/volatile_fetch_service.py diff --git a/src/routers/volatile.py b/src/routers/volatile.py index 1adb018..c36c4cb 100644 --- a/src/routers/volatile.py +++ b/src/routers/volatile.py @@ -19,7 +19,15 @@ from src.models.volatile import ( NAMESPACE_DEFAULT_TTL, ) from src.services.volatile_service import VolatileCacheService -from src.core.dependencies import verify_api_key, QdrantDep, OllamaDep +from src.services.volatile_fetch_service import VolatileFetchService +from src.core.dependencies import ( + verify_api_key, + QdrantDep, + OllamaDep, + get_weather_provider, + get_news_provider, + get_alphavantage_provider, +) from src.core.multi_tenancy import DEFAULT_USER from src.config import get_settings @@ -221,6 +229,188 @@ async def store_volatile( raise HTTPException(status_code=500, detail=f"Failed to store record: {str(e)}") +@router.post("/fetch/weather/{city}") +async def fetch_weather( + city: str, + user: str = Query(default=DEFAULT_USER, description="User identifier"), + ttl: int = Query(default=86400, ge=60, le=604800, description="TTL in seconds"), + qdrant: QdrantDep = None, + ollama: OllamaDep = None, + api_key: str = Depends(verify_api_key) +): + """ + Fetch current weather for a city and store in volatile cache. + + Called by scheduler for prefetch or on-demand. Geocodes city name + and fetches weather from Open-Meteo API. + + **Example:** + ``` + POST /volatile/fetch/weather/amsterdam?user=jpmschweitzer + ``` + """ + volatile_service = get_volatile_service(qdrant, ollama) + weather_provider = get_weather_provider() + + fetch_service = VolatileFetchService( + volatile_service=volatile_service, + weather_provider=weather_provider, + ) + + result = await fetch_service.fetch_weather(user, city, ttl=ttl) + + if not result.success: + raise HTTPException(status_code=500, detail=result.error) + + return { + "success": True, + "namespace": result.namespace, + "key": result.key, + "record": result.record, + } + + +@router.post("/fetch/news/{category}") +async def fetch_news( + category: str = "general", + user: str = Query(default=DEFAULT_USER, description="User identifier"), + limit: int = Query(default=10, ge=1, le=50, description="Max headlines"), + ttl: int = Query(default=7200, ge=60, le=86400, description="TTL in seconds"), + qdrant: QdrantDep = None, + ollama: OllamaDep = None, + api_key: str = Depends(verify_api_key) +): + """ + Fetch news headlines and store in volatile cache. + + Fetches from configured news sources (NOS, BBC) based on user settings. + Categories: general, world, tech, business, politics, etc. + + **Example:** + ``` + POST /volatile/fetch/news/tech?user=jpmschweitzer&limit=15 + ``` + """ + volatile_service = get_volatile_service(qdrant, ollama) + weather_provider = get_weather_provider() + news_provider = await get_news_provider() + + fetch_service = VolatileFetchService( + volatile_service=volatile_service, + weather_provider=weather_provider, + news_provider=news_provider, + ) + + result = await fetch_service.fetch_news(user, category, limit=limit, ttl=ttl) + + if not result.success: + raise HTTPException(status_code=500, detail=result.error) + + return { + "success": True, + "namespace": result.namespace, + "key": result.key, + "record": result.record, + } + + +@router.post("/fetch/stock/{symbol}") +async def fetch_stock( + symbol: str, + user: str = Query(default=DEFAULT_USER, description="User identifier"), + ttl: int = Query(default=300, ge=60, le=3600, description="TTL in seconds"), + qdrant: QdrantDep = None, + ollama: OllamaDep = None, + api_key: str = Depends(verify_api_key) +): + """ + Fetch stock quote and store in volatile cache. + + Fetches from Alpha Vantage API. Requires API key configured in settings. + + **Example:** + ``` + POST /volatile/fetch/stock/AAPL?user=jpmschweitzer + ``` + """ + volatile_service = get_volatile_service(qdrant, ollama) + weather_provider = get_weather_provider() + financial_provider = await get_alphavantage_provider() + + if not financial_provider: + raise HTTPException( + status_code=503, + detail="Financial provider not configured (Alpha Vantage API key missing)" + ) + + fetch_service = VolatileFetchService( + volatile_service=volatile_service, + weather_provider=weather_provider, + financial_provider=financial_provider, + ) + + result = await fetch_service.fetch_stock(user, symbol, ttl=ttl) + + if not result.success: + raise HTTPException(status_code=500, detail=result.error) + + return { + "success": True, + "namespace": result.namespace, + "key": result.key, + "record": result.record, + } + + +@router.post("/fetch/crypto/{symbol}") +async def fetch_crypto( + symbol: str, + market: str = Query(default="USD", description="Market currency"), + user: str = Query(default=DEFAULT_USER, description="User identifier"), + ttl: int = Query(default=300, ge=60, le=3600, description="TTL in seconds"), + qdrant: QdrantDep = None, + ollama: OllamaDep = None, + api_key: str = Depends(verify_api_key) +): + """ + Fetch cryptocurrency quote and store in volatile cache. + + Fetches from Alpha Vantage API. Requires API key configured in settings. + + **Example:** + ``` + POST /volatile/fetch/crypto/BTC?market=EUR&user=jpmschweitzer + ``` + """ + volatile_service = get_volatile_service(qdrant, ollama) + weather_provider = get_weather_provider() + financial_provider = await get_alphavantage_provider() + + if not financial_provider: + raise HTTPException( + status_code=503, + detail="Financial provider not configured (Alpha Vantage API key missing)" + ) + + fetch_service = VolatileFetchService( + volatile_service=volatile_service, + weather_provider=weather_provider, + financial_provider=financial_provider, + ) + + result = await fetch_service.fetch_crypto(user, symbol, market=market, ttl=ttl) + + if not result.success: + raise HTTPException(status_code=500, detail=result.error) + + return { + "success": True, + "namespace": result.namespace, + "key": result.key, + "record": result.record, + } + + @router.get("/{namespace}/{key}", response_model=VolatileRecordResponse) async def get_record( namespace: str, diff --git a/src/services/volatile_fetch_service.py b/src/services/volatile_fetch_service.py new file mode 100644 index 0000000..2112f1f --- /dev/null +++ b/src/services/volatile_fetch_service.py @@ -0,0 +1,359 @@ +""" +Volatile Fetch service for Library Desk. + +Orchestrates fetching data from external APIs and storing in volatile cache. +Called by scheduler for prefetch or by HybridRAG for reactive caching. +""" + +import logging +from typing import Optional +from dataclasses import dataclass + +from src.apis import ( + OpenMeteoProvider, + AggregatedNewsProvider, + AlphaVantageProvider, + CurrentWeather, + NewsFeed, + StockQuote, +) +from src.services.volatile_service import VolatileCacheService +from src.models.volatile import VolatileRecordResponse, VolatileNamespace + +logger = logging.getLogger(__name__) + + +@dataclass +class FetchResult: + """Result of a volatile fetch operation.""" + success: bool + namespace: str + key: str + record: Optional[VolatileRecordResponse] = None + error: Optional[str] = None + + +class VolatileFetchService: + """ + Service to fetch external data and store in volatile cache. + + Supports: + - Weather: Current conditions and forecast via Open-Meteo + - News: Headlines from configured sources (NOS, BBC) + - Financial: Stock/crypto quotes via Alpha Vantage + """ + + def __init__( + self, + volatile_service: VolatileCacheService, + weather_provider: OpenMeteoProvider, + news_provider: Optional[AggregatedNewsProvider] = None, + financial_provider: Optional[AlphaVantageProvider] = None, + ): + """ + Initialize volatile fetch service. + + Args: + volatile_service: Service for volatile cache storage + weather_provider: Open-Meteo weather provider + news_provider: Aggregated news provider (optional) + financial_provider: Alpha Vantage provider (optional) + """ + self.volatile = volatile_service + self.weather = weather_provider + self.news = news_provider + self.financial = financial_provider + + async def fetch_weather( + self, + user: str, + city: str, + ttl: int = 86400, # 24 hours + ) -> FetchResult: + """ + Fetch current weather for a city and store in volatile cache. + + Args: + user: User identifier + city: City name (will be geocoded) + ttl: Time-to-live in seconds + + Returns: + FetchResult with success status and stored record + """ + try: + # Geocode city and get weather + location = await self.weather.geocode(city) + if not location: + return FetchResult( + success=False, + namespace="weather", + key=city.lower(), + error=f"Could not geocode city: {city}" + ) + + weather = await self.weather.get_current(location) + + # Convert to storage format + data = { + "temperature": weather.temperature, + "feels_like": weather.feels_like, + "humidity": weather.humidity, + "wind_speed": weather.wind_speed, + "wind_direction": weather.wind_direction, + "conditions": weather.condition_text, + "condition_code": weather.condition.value, + "location": weather.location, + "text": weather.to_text(), + } + + # Store in volatile cache + record = await self.volatile.store( + user=user, + namespace=VolatileNamespace.WEATHER, + key=city.lower(), + data=data, + source="openmeteo", + ttl=ttl, + ) + + logger.info(f"Stored weather for {city} (user={user})") + return FetchResult( + success=True, + namespace="weather", + key=city.lower(), + record=record + ) + + except Exception as e: + logger.error(f"Failed to fetch weather for {city}: {e}") + return FetchResult( + success=False, + namespace="weather", + key=city.lower(), + error=str(e) + ) + + async def fetch_news( + self, + user: str, + category: str = "general", + limit: int = 10, + ttl: int = 7200, # 2 hours + ) -> FetchResult: + """ + Fetch news headlines and store in volatile cache. + + Args: + user: User identifier + category: News category (general, tech, world, etc.) + limit: Maximum headlines to fetch + ttl: Time-to-live in seconds + + Returns: + FetchResult with success status and stored record + """ + if not self.news: + return FetchResult( + success=False, + namespace="news", + key=category, + error="News provider not configured" + ) + + try: + feed = await self.news.get_feed(category, limit=limit) + + # Convert to storage format + headlines = [] + for item in feed.items: + headlines.append({ + "title": item.title, + "description": item.description, + "url": item.url, + "source": item.source, + "published": item.published.isoformat() if item.published else None, + }) + + data = { + "category": category, + "headlines": headlines, + "count": len(headlines), + "sources": list(set(h["source"] for h in headlines)), + "text": feed.to_text(), + } + + # Store in volatile cache + record = await self.volatile.store( + user=user, + namespace=VolatileNamespace.NEWS, + key=category, + data=data, + source="aggregated", + ttl=ttl, + ) + + logger.info(f"Stored {len(headlines)} headlines for {category} (user={user})") + return FetchResult( + success=True, + namespace="news", + key=category, + record=record + ) + + except Exception as e: + logger.error(f"Failed to fetch news for {category}: {e}") + return FetchResult( + success=False, + namespace="news", + key=category, + error=str(e) + ) + + async def fetch_stock( + self, + user: str, + symbol: str, + ttl: int = 300, # 5 minutes + ) -> FetchResult: + """ + Fetch stock quote and store in volatile cache. + + Args: + user: User identifier + symbol: Stock ticker symbol (e.g., "AAPL") + ttl: Time-to-live in seconds + + Returns: + FetchResult with success status and stored record + """ + if not self.financial: + return FetchResult( + success=False, + namespace="financial", + key=symbol.lower(), + error="Financial provider not configured" + ) + + try: + quote = await self.financial.get_quote(symbol) + if not quote: + return FetchResult( + success=False, + namespace="financial", + key=symbol.lower(), + error=f"No quote found for symbol: {symbol}" + ) + + # Convert to storage format + data = { + "symbol": quote.symbol, + "name": quote.name, + "price": quote.price, + "currency": quote.currency, + "change": quote.change, + "change_percent": quote.change_percent, + "text": quote.to_text(), + } + + # Store in volatile cache + record = await self.volatile.store( + user=user, + namespace=VolatileNamespace.FINANCIAL, + key=symbol.lower(), + data=data, + source="alphavantage", + ttl=ttl, + ) + + logger.info(f"Stored quote for {symbol} (user={user})") + return FetchResult( + success=True, + namespace="financial", + key=symbol.lower(), + record=record + ) + + except Exception as e: + logger.error(f"Failed to fetch quote for {symbol}: {e}") + return FetchResult( + success=False, + namespace="financial", + key=symbol.lower(), + error=str(e) + ) + + async def fetch_crypto( + self, + user: str, + symbol: str, + market: str = "USD", + ttl: int = 300, # 5 minutes + ) -> FetchResult: + """ + Fetch cryptocurrency quote and store in volatile cache. + + Args: + user: User identifier + symbol: Crypto symbol (e.g., "BTC", "ETH") + market: Market currency (default: USD) + ttl: Time-to-live in seconds + + Returns: + FetchResult with success status and stored record + """ + if not self.financial: + return FetchResult( + success=False, + namespace="financial", + key=f"{symbol.lower()}_{market.lower()}", + error="Financial provider not configured" + ) + + try: + quote = await self.financial.get_crypto_quote(symbol, market) + if not quote: + return FetchResult( + success=False, + namespace="financial", + key=f"{symbol.lower()}_{market.lower()}", + error=f"No quote found for crypto: {symbol}/{market}" + ) + + key = f"{symbol.lower()}_{market.lower()}" + + # Convert to storage format + data = { + "symbol": quote.symbol, + "name": quote.name, + "price": quote.price, + "currency": quote.currency, + "text": quote.to_text(), + } + + # Store in volatile cache + record = await self.volatile.store( + user=user, + namespace=VolatileNamespace.FINANCIAL, + key=key, + data=data, + source="alphavantage", + ttl=ttl, + ) + + logger.info(f"Stored crypto quote for {symbol}/{market} (user={user})") + return FetchResult( + success=True, + namespace="financial", + key=key, + record=record + ) + + except Exception as e: + logger.error(f"Failed to fetch crypto quote for {symbol}: {e}") + return FetchResult( + success=False, + namespace="financial", + key=f"{symbol.lower()}_{market.lower()}", + error=str(e) + )