feat: add environment endpoint for weather, forecast, and sun data
Add GET /tools/environment endpoint that fetches weather, forecast, sun times, and air quality data from user's volatile Qdrant collection. - Add qdrant-client dependency - Create QdrantReadClient wrapper for read-only queries - Add environment schemas and service in tools domain - Parse weather, forecast, sun times, and air quality from Qdrant payloads - Support user-specific collections via preferred_username from OIDC - Add comprehensive service tests (13 tests) 🤖 Generated with [Claude Code](https://claude.com/claude-code) Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 4.5
parent
b8fca7060f
commit
6045c6ac6a
@@ -7,6 +7,7 @@ from src.shared.clients.portainer_client import PortainerClient, get_portainer_c
|
||||
from src.shared.clients.npm_client import NPMClient, get_npm_client
|
||||
from src.shared.clients.homeassistant_client import HomeAssistantClient, get_homeassistant_client
|
||||
from src.shared.clients.authentik_client import AuthentikClient, get_authentik_client
|
||||
from src.shared.clients.qdrant_client import QdrantReadClient, get_qdrant_client
|
||||
|
||||
__all__ = [
|
||||
"PortainerClient",
|
||||
@@ -17,4 +18,6 @@ __all__ = [
|
||||
"get_homeassistant_client",
|
||||
"AuthentikClient",
|
||||
"get_authentik_client",
|
||||
"QdrantReadClient",
|
||||
"get_qdrant_client",
|
||||
]
|
||||
|
||||
@@ -0,0 +1,258 @@
|
||||
"""
|
||||
Qdrant Vector Database Client (Read-Only)
|
||||
|
||||
Provides read-only access to Qdrant collections for querying volatile data.
|
||||
Used to fetch weather, forecast, and sun times from the volatile_{user} collection.
|
||||
"""
|
||||
import time
|
||||
from typing import List, Dict, Any, Optional
|
||||
from qdrant_client import QdrantClient
|
||||
from qdrant_client.models import Filter, FieldCondition, MatchValue, Range
|
||||
|
||||
from src.shared.logging import get_logger
|
||||
from src.shared.config import get_settings
|
||||
|
||||
logger = get_logger(__name__)
|
||||
settings = get_settings()
|
||||
|
||||
|
||||
class QdrantReadClient:
|
||||
"""
|
||||
Read-only Qdrant client for accessing volatile data.
|
||||
|
||||
Connects to Qdrant and provides methods to query collections
|
||||
with filtering by namespace and TTL expiry.
|
||||
"""
|
||||
|
||||
VOLATILE_COLLECTION_PREFIX = "volatile_"
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
host: Optional[str] = None,
|
||||
port: Optional[int] = None,
|
||||
):
|
||||
"""
|
||||
Initialize Qdrant read client.
|
||||
|
||||
Args:
|
||||
host: Qdrant server host (default from settings)
|
||||
port: Qdrant server port (default from settings)
|
||||
"""
|
||||
self.host = host or settings.qdrant_host
|
||||
self.port = port or settings.qdrant_port
|
||||
self._client: Optional[QdrantClient] = None
|
||||
|
||||
logger.info(f"Initialized QdrantReadClient: {self.host}:{self.port}")
|
||||
|
||||
@property
|
||||
def client(self) -> QdrantClient:
|
||||
"""Lazy-load Qdrant client connection."""
|
||||
if self._client is None:
|
||||
self._client = QdrantClient(
|
||||
host=self.host,
|
||||
port=self.port,
|
||||
)
|
||||
return self._client
|
||||
|
||||
def _get_volatile_collection(self, user: str) -> str:
|
||||
"""Get volatile collection name for user."""
|
||||
return f"{self.VOLATILE_COLLECTION_PREFIX}{user}"
|
||||
|
||||
def _current_timestamp_ms(self) -> int:
|
||||
"""Get current timestamp in milliseconds."""
|
||||
return int(time.time() * 1000)
|
||||
|
||||
async def collection_exists(self, collection_name: str) -> bool:
|
||||
"""
|
||||
Check if a collection exists.
|
||||
|
||||
Args:
|
||||
collection_name: Name of collection to check
|
||||
|
||||
Returns:
|
||||
True if collection exists
|
||||
"""
|
||||
try:
|
||||
collections = self.client.get_collections()
|
||||
existing = [c.name for c in collections.collections]
|
||||
return collection_name in existing
|
||||
except Exception as e:
|
||||
logger.error(f"Error checking collection existence: {e}")
|
||||
return False
|
||||
|
||||
async def get_by_namespace(
|
||||
self,
|
||||
user: str,
|
||||
namespace: str,
|
||||
include_expired: bool = False
|
||||
) -> List[Dict[str, Any]]:
|
||||
"""
|
||||
Get all records for a specific namespace from user's volatile collection.
|
||||
|
||||
Args:
|
||||
user: User identifier (e.g., 'jpmschweitzer' or 'default')
|
||||
namespace: Namespace to filter (e.g., 'weather', 'forecast', 'sun')
|
||||
include_expired: Whether to include expired records (default False)
|
||||
|
||||
Returns:
|
||||
List of records with payload data
|
||||
"""
|
||||
collection_name = self._get_volatile_collection(user)
|
||||
|
||||
if not await self.collection_exists(collection_name):
|
||||
logger.debug(f"Collection {collection_name} does not exist")
|
||||
return []
|
||||
|
||||
# Build filter conditions
|
||||
conditions = [
|
||||
FieldCondition(
|
||||
key="namespace",
|
||||
match=MatchValue(value=namespace)
|
||||
)
|
||||
]
|
||||
|
||||
# Add TTL expiry filter unless including expired
|
||||
if not include_expired:
|
||||
now_ms = self._current_timestamp_ms()
|
||||
conditions.append(
|
||||
FieldCondition(
|
||||
key="ttl_expiry",
|
||||
range=Range(gt=now_ms)
|
||||
)
|
||||
)
|
||||
|
||||
query_filter = Filter(must=conditions)
|
||||
|
||||
try:
|
||||
# Scroll through matching records
|
||||
points, _ = self.client.scroll(
|
||||
collection_name=collection_name,
|
||||
scroll_filter=query_filter,
|
||||
limit=100,
|
||||
with_payload=True,
|
||||
with_vectors=False
|
||||
)
|
||||
|
||||
results = []
|
||||
for point in points:
|
||||
payload = dict(point.payload) if point.payload else {}
|
||||
results.append({
|
||||
"id": str(point.id),
|
||||
"namespace": payload.get("namespace"),
|
||||
"key": payload.get("key"),
|
||||
"raw_data": payload.get("raw_data", {}),
|
||||
"source": payload.get("source"),
|
||||
"ttl_expiry": payload.get("ttl_expiry"),
|
||||
"updated_at": payload.get("updated_at"),
|
||||
})
|
||||
|
||||
logger.debug(
|
||||
f"Found {len(results)} records in {collection_name}/{namespace}"
|
||||
)
|
||||
return results
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error fetching from {collection_name}/{namespace}: {e}")
|
||||
return []
|
||||
|
||||
async def get_environment_data(
|
||||
self,
|
||||
user: str
|
||||
) -> Dict[str, Any]:
|
||||
"""
|
||||
Get all environment data (weather, forecast, sun times) for a user.
|
||||
|
||||
Convenience method that fetches all environment-related namespaces
|
||||
in a single call.
|
||||
|
||||
Args:
|
||||
user: User identifier
|
||||
|
||||
Returns:
|
||||
Dict with 'weather', 'forecast', 'sun_times', 'air_quality' keys
|
||||
(each may be None if no data found)
|
||||
"""
|
||||
result = {
|
||||
"weather": None,
|
||||
"forecast": None,
|
||||
"sun_times": None,
|
||||
"air_quality": None,
|
||||
}
|
||||
|
||||
# Fetch weather data
|
||||
weather_records = await self.get_by_namespace(user, "weather")
|
||||
if weather_records:
|
||||
# Get the first/most recent weather record
|
||||
result["weather"] = weather_records[0].get("raw_data")
|
||||
# Check if air quality is embedded in weather data
|
||||
if result["weather"]:
|
||||
aqi = result["weather"].get("aqi") or result["weather"].get("air_quality")
|
||||
if aqi:
|
||||
result["air_quality"] = aqi if isinstance(aqi, dict) else {"aqi": aqi}
|
||||
|
||||
# Fetch forecast data
|
||||
forecast_records = await self.get_by_namespace(user, "forecast")
|
||||
if forecast_records:
|
||||
# Forecast might be a single record with list or multiple records
|
||||
first_record = forecast_records[0].get("raw_data")
|
||||
if isinstance(first_record, list):
|
||||
result["forecast"] = first_record
|
||||
elif isinstance(first_record, dict):
|
||||
# Could be a dict with 'days' or 'forecast' key
|
||||
result["forecast"] = first_record.get(
|
||||
"days",
|
||||
first_record.get("forecast", [first_record])
|
||||
)
|
||||
|
||||
# Fetch sun times data
|
||||
sun_records = await self.get_by_namespace(user, "sun")
|
||||
if sun_records:
|
||||
result["sun_times"] = sun_records[0].get("raw_data")
|
||||
|
||||
# Check for separate air quality namespace if not embedded
|
||||
if result["air_quality"] is None:
|
||||
aq_records = await self.get_by_namespace(user, "air_quality")
|
||||
if aq_records:
|
||||
result["air_quality"] = aq_records[0].get("raw_data")
|
||||
|
||||
return result
|
||||
|
||||
async def health_check(self) -> Dict[str, Any]:
|
||||
"""
|
||||
Check Qdrant connectivity.
|
||||
|
||||
Returns:
|
||||
Dict with connection status and info
|
||||
"""
|
||||
try:
|
||||
collections = self.client.get_collections()
|
||||
volatile_collections = [
|
||||
c.name for c in collections.collections
|
||||
if c.name.startswith(self.VOLATILE_COLLECTION_PREFIX)
|
||||
]
|
||||
return {
|
||||
"status": "healthy",
|
||||
"connected": True,
|
||||
"host": f"{self.host}:{self.port}",
|
||||
"volatile_collections": volatile_collections,
|
||||
}
|
||||
except Exception as e:
|
||||
logger.error(f"Qdrant health check failed: {e}")
|
||||
return {
|
||||
"status": "unhealthy",
|
||||
"connected": False,
|
||||
"host": f"{self.host}:{self.port}",
|
||||
"error": str(e),
|
||||
}
|
||||
|
||||
|
||||
# Singleton instance for reuse
|
||||
_qdrant_client: Optional[QdrantReadClient] = None
|
||||
|
||||
|
||||
def get_qdrant_client() -> QdrantReadClient:
|
||||
"""Get or create singleton Qdrant client instance."""
|
||||
global _qdrant_client
|
||||
if _qdrant_client is None:
|
||||
_qdrant_client = QdrantReadClient()
|
||||
return _qdrant_client
|
||||
Reference in New Issue
Block a user