From 72f515bf61bf691eb09a72bf0b6e015346433985 Mon Sep 17 00:00:00 2001 From: Jeroen Schweitzer Date: Wed, 7 Jan 2026 11:53:46 +0100 Subject: [PATCH] feat: add combined environment endpoint for concurrent weather + air quality fetch MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - POST /volatile/fetch/environment/{city} fetches both in parallel - Single geocode lookup shared between API calls - Uses asyncio.gather() for concurrent external requests - Fix scheduler executor name (rest_api → rest_api_executor) 🤖 Generated with [Claude Code](https://claude.com/claude-code) Co-Authored-By: Claude Opus 4.5 --- CHANGELOG.md | 14 +++ pyproject.toml | 2 +- src/clients/scheduler_client.py | 4 +- src/routers/volatile.py | 58 ++++++++++ src/services/volatile_fetch_service.py | 140 ++++++++++++++++++++++++- 5 files changed, 214 insertions(+), 4 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 500caab..843c241 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,20 @@ All notable changes to Library Desk will be documented in this file. The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/), and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). +## [1.7.0] - 2026-01-07 + +### Added + +- **Combined Environment Endpoint** - `POST /volatile/fetch/environment/{city}` + - Fetches weather and air quality concurrently with `asyncio.gather()` + - Single geocode lookup shared between both API calls + - More efficient than calling weather and air_quality separately + - Reduces wall-clock time and eliminates redundant geocoding + +### Fixed + +- **Scheduler executor name** - Fixed `rest_api` → `rest_api_executor` in SchedulerTask model and register_volatile_fetch() to prevent "Executor module not found" errors + ## [1.6.2] - 2025-12-30 ### Added diff --git a/pyproject.toml b/pyproject.toml index 97ca48c..b333d6a 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1,6 +1,6 @@ [project] name = "library-desk" -version = "1.6.2" +version = "1.7.0" description = "Coordination service for The Library system - HybridRAG queries, document ingestion, entity extraction, and knowledge consolidation" readme = "README.md" requires-python = ">=3.12" diff --git a/src/clients/scheduler_client.py b/src/clients/scheduler_client.py index 90496e5..21a36c2 100644 --- a/src/clients/scheduler_client.py +++ b/src/clients/scheduler_client.py @@ -18,7 +18,7 @@ class SchedulerTask(BaseModel): task_name: str = Field(..., description="Unique task identifier") service: str = Field(default="library-desk", description="Service that owns this task") - executor: str = Field(default="rest_api", description="Executor type") + executor: str = Field(default="rest_api_executor", description="Executor type") priority: int = Field(default=50, ge=1, le=100, description="Priority (lower = higher)") description: Optional[str] = Field(None, description="Human-readable description") enabled: bool = Field(default=True, description="Whether task is enabled") @@ -291,7 +291,7 @@ class SchedulerClient: task = SchedulerTask( task_name=task_name, service="library-desk", - executor="rest_api", + executor="rest_api_executor", priority=60, # Background maintenance priority description=description or f"Prefetch {namespace}/{key} for {user}", minute=schedule.get("minute", -1), diff --git a/src/routers/volatile.py b/src/routers/volatile.py index ea94de3..e27680b 100644 --- a/src/routers/volatile.py +++ b/src/routers/volatile.py @@ -545,6 +545,64 @@ async def fetch_air_quality( } +@router.post("/fetch/environment/{city}") +async def fetch_environment( + city: str, + user: str = Query(default=DEFAULT_USER, description="User identifier"), + weather_ttl: int = Query(default=3600, ge=60, le=86400, description="Weather TTL in seconds"), + air_quality_ttl: int = Query(default=3600, ge=60, le=86400, description="Air quality TTL in seconds"), + qdrant: QdrantDep = None, + ollama: OllamaDep = None, + api_key: str = Depends(verify_api_key) +): + """ + Fetch weather and air quality concurrently for a city. + + Performs a single geocode lookup and fetches both weather and air quality + data in parallel, storing both in volatile cache. More efficient than + calling /fetch/weather and /fetch/air_quality separately. + + **Example:** + ``` + POST /volatile/fetch/environment/rotterdam?user=jpmschweitzer + ``` + + **Response includes:** + - weather: Current conditions (temperature, humidity, wind, UV) + - air_quality: AQI indices, pollutants, pollen data + """ + 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_environment( + user, city, weather_ttl=weather_ttl, air_quality_ttl=air_quality_ttl + ) + + if not result.success: + raise HTTPException(status_code=500, detail="; ".join(result.errors)) + + return { + "success": True, + "key": result.key, + "weather": { + "success": result.weather.success if result.weather else False, + "record": result.weather.record if result.weather else None, + "error": result.weather.error if result.weather else None, + }, + "air_quality": { + "success": result.air_quality.success if result.air_quality else False, + "record": result.air_quality.record if result.air_quality else None, + "error": result.air_quality.error if result.air_quality else None, + }, + "errors": result.errors, + } + + @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 index 5c6d0bf..32085d5 100644 --- a/src/services/volatile_fetch_service.py +++ b/src/services/volatile_fetch_service.py @@ -5,9 +5,10 @@ Orchestrates fetching data from external APIs and storing in volatile cache. Called by scheduler for prefetch or by HybridRAG for reactive caching. """ +import asyncio import logging from typing import Optional -from dataclasses import dataclass +from dataclasses import dataclass, field from src.apis import ( OpenMeteoProvider, @@ -36,6 +37,16 @@ class FetchResult: error: Optional[str] = None +@dataclass +class EnvironmentFetchResult: + """Result of combined environment fetch (weather + air quality).""" + success: bool + key: str + weather: Optional[FetchResult] = None + air_quality: Optional[FetchResult] = None + errors: list[str] = field(default_factory=list) + + class VolatileFetchService: """ Service to fetch external data and store in volatile cache. @@ -596,3 +607,130 @@ class VolatileFetchService: key=city.lower(), error=str(e) ) + + async def fetch_environment( + self, + user: str, + city: str, + weather_ttl: int = 3600, + air_quality_ttl: int = 3600, + ) -> EnvironmentFetchResult: + """ + Fetch weather and air quality concurrently for a city. + + Performs a single geocode lookup and fetches both weather and air quality + data in parallel, storing both in volatile cache. + + Args: + user: User identifier + city: City name (will be geocoded once) + weather_ttl: TTL for weather data (default 1 hour) + air_quality_ttl: TTL for air quality data (default 1 hour) + + Returns: + EnvironmentFetchResult with both weather and air quality results + """ + errors: list[str] = [] + key = city.lower() + + # Single geocode lookup (shared by both fetches) + try: + location = await self.weather.geocode(city) + if not location: + return EnvironmentFetchResult( + success=False, + key=key, + errors=[f"Could not geocode city: {city}"] + ) + except Exception as e: + return EnvironmentFetchResult( + success=False, + key=key, + errors=[f"Geocoding failed: {e}"] + ) + + # Fetch weather and air quality concurrently + async def fetch_weather_data() -> FetchResult: + try: + current = await self.weather.get_current(location) + text = current.to_text() + data = { + "temperature": current.temperature, + "feels_like": current.feels_like, + "humidity": current.humidity, + "wind_speed": current.wind_speed, + "wind_direction": current.wind_direction, + "conditions": current.condition_text, + "condition_code": current.condition.value, + "uv_index": current.uv_index, + "location": current.location, + "text": text, + } + record = await self.volatile.store( + user=user, + namespace=VolatileNamespace.WEATHER, + key=key, + data=data, + source="openmeteo", + ttl=weather_ttl, + ) + return FetchResult(success=True, namespace="weather", key=key, record=record) + except Exception as e: + return FetchResult(success=False, namespace="weather", key=key, error=str(e)) + + async def fetch_air_quality_data() -> FetchResult: + try: + air_quality = await self.weather.get_air_quality(location) + data = { + "location": air_quality.location, + "aqi_european": air_quality.aqi_european, + "aqi_us": air_quality.aqi_us, + "pm2_5": air_quality.pm2_5, + "pm10": air_quality.pm10, + "ozone": air_quality.ozone, + "nitrogen_dioxide": air_quality.nitrogen_dioxide, + "sulphur_dioxide": air_quality.sulphur_dioxide, + "carbon_monoxide": air_quality.carbon_monoxide, + "pollen_grass": air_quality.pollen_grass, + "pollen_birch": air_quality.pollen_birch, + "pollen_alder": air_quality.pollen_alder, + "text": air_quality.to_text(), + } + record = await self.volatile.store( + user=user, + namespace=VolatileNamespace.AIR_QUALITY, + key=key, + data=data, + source="openmeteo", + ttl=air_quality_ttl, + ) + return FetchResult(success=True, namespace="air_quality", key=key, record=record) + except Exception as e: + return FetchResult(success=False, namespace="air_quality", key=key, error=str(e)) + + # Run both fetches concurrently + weather_result, air_quality_result = await asyncio.gather( + fetch_weather_data(), + fetch_air_quality_data(), + ) + + # Collect any errors + if not weather_result.success: + errors.append(f"Weather: {weather_result.error}") + if not air_quality_result.success: + errors.append(f"Air quality: {air_quality_result.error}") + + success = weather_result.success or air_quality_result.success + logger.info( + f"Environment fetch for {city} (user={user}): " + f"weather={'ok' if weather_result.success else 'failed'}, " + f"air_quality={'ok' if air_quality_result.success else 'failed'}" + ) + + return EnvironmentFetchResult( + success=success, + key=key, + weather=weather_result, + air_quality=air_quality_result, + errors=errors, + )