feat: add volatile fetch endpoints for scheduler-driven prefetch

- 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 <noreply@anthropic.com>
This commit is contained in:
2025-12-26 13:40:40 +01:00
co-authored by Claude Opus 4.5
parent 2b8c229f53
commit ab892745fa
2 changed files with 550 additions and 1 deletions
+191 -1
View File
@@ -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,
+359
View File
@@ -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)
)