Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
37f8e1819e | ||
|
|
1f848c4878 | ||
|
|
7297e6b9f1 | ||
|
|
2552bfd1f9 |
@@ -5,6 +5,70 @@ 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.4.3] - 2025-12-24
|
||||
|
||||
### Changed
|
||||
|
||||
- **Volatile Cache System Refactored to Vector Storage**
|
||||
- Backend migrated from Redis to Qdrant for semantic search capability
|
||||
- Data converted to natural language for embedding and semantic retrieval
|
||||
- Collection naming: `volatile_{user}` for per-user isolation
|
||||
- TTL implemented via `ttl_expiry` timestamp in vector payload
|
||||
- Simplified endpoints:
|
||||
- `GET /volatile/search?q=...` - Semantic search across volatile data
|
||||
- `POST /volatile/store?namespace=...&key=...` - Store with query params
|
||||
- `GET /volatile/{namespace}/{key}` - Get specific record
|
||||
- `DELETE /volatile/{namespace}/{key}` - Delete record
|
||||
- Removed namespace-specific URL patterns (simpler API for LLM tool use)
|
||||
|
||||
### Added
|
||||
|
||||
- **HybridRAG Volatile Integration** - Volatile cache now included in multi-source search
|
||||
- Volatile results get priority boost in RRF fusion (current data ranks higher)
|
||||
- New config options: `enable_volatile`, `volatile_limit` (default 1), `volatile_threshold`
|
||||
- Timing breakdown includes `volatile_ms`
|
||||
- **Volatile Cleanup Endpoint** - `POST /maintenance/cleanup/volatile`
|
||||
- Purges expired records across all `volatile_*` collections
|
||||
- Scheduler task for every 10 minutes recommended
|
||||
- Returns per-collection cleanup counts
|
||||
- **Natural Language Conversion** - Structured data converted for embedding
|
||||
- Template-based conversion for each namespace (weather, news, financial, etc.)
|
||||
- Fallback for custom namespaces
|
||||
|
||||
## [1.4.2] - 2025-12-24
|
||||
|
||||
### Added
|
||||
|
||||
- **Volatile Cache System** - Ephemeral data storage with TTL
|
||||
- `GET /volatile/{namespace}/{key}` - Retrieve cached record
|
||||
- `POST /volatile/{namespace}/{key}` - Store/update record with TTL
|
||||
- `DELETE /volatile/{namespace}/{key}` - Remove record
|
||||
- `GET /volatile/{namespace}` - List keys in namespace
|
||||
- `DELETE /volatile/{namespace}` - Clear all records in namespace
|
||||
- `GET /volatile/stats` - Cache statistics by namespace
|
||||
- `GET /volatile/scheduled` - Records needing refresh (for scheduler)
|
||||
- `GET /volatile/namespaces` - List available namespaces with default TTLs
|
||||
- **Volatile Namespaces** - Predefined categories with appropriate TTLs:
|
||||
- `weather` (30min) - Weather conditions and forecasts
|
||||
- `news` (1hr) - Headlines and breaking news
|
||||
- `financial` (5min) - Stock prices, exchange rates
|
||||
- `transit` (5min) - Train/bus schedules, delays
|
||||
- `traffic` (10min) - Commute times, road conditions
|
||||
- `air_quality` (1hr) - Pollution, pollen counts
|
||||
- `sports` (1min) - Live scores, matches
|
||||
- `social` (10min) - Social notifications
|
||||
- `system` (1min) - Service health status
|
||||
- `context` (1hr) - Session state
|
||||
- `custom` (1hr) - User-defined data
|
||||
- **Refresh Schedule Support** - Optional cron expressions for scheduler integration
|
||||
|
||||
## [1.4.1] - 2025-12-24
|
||||
|
||||
### Fixed
|
||||
|
||||
- Wiki.js API token now optional - GraphQL API works without authentication
|
||||
- Container startup failure when `WIKI_GRAPHQL_API` env var not set
|
||||
|
||||
## [1.4.0] - 2025-12-24
|
||||
|
||||
### Added
|
||||
|
||||
+1
-1
@@ -1,6 +1,6 @@
|
||||
[project]
|
||||
name = "library-desk"
|
||||
version = "1.4.0"
|
||||
version = "1.4.3"
|
||||
description = "Coordination service for The Library system - HybridRAG queries, document ingestion, entity extraction, and knowledge consolidation"
|
||||
readme = "README.md"
|
||||
requires-python = ">=3.12"
|
||||
|
||||
@@ -11,7 +11,7 @@ Provides async vector operations with:
|
||||
from qdrant_client import QdrantClient
|
||||
from qdrant_client.models import (
|
||||
Distance, VectorParams, PointStruct,
|
||||
Filter, FieldCondition, MatchValue
|
||||
Filter, FieldCondition, MatchValue, Range
|
||||
)
|
||||
from typing import List, Dict, Any, Optional
|
||||
import uuid
|
||||
@@ -676,4 +676,134 @@ class QdrantClientWrapper:
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Failed to list collections: {e}", exc_info=True)
|
||||
return []
|
||||
|
||||
# ========== Volatile Data Methods ==========
|
||||
|
||||
async def search_with_expiry_filter(
|
||||
self,
|
||||
collection_name: str,
|
||||
query_vector: List[float],
|
||||
current_timestamp: int,
|
||||
limit: int = 10,
|
||||
score_threshold: float = 0.7
|
||||
) -> List[Dict[str, Any]]:
|
||||
"""
|
||||
Search vectors filtering out expired records.
|
||||
|
||||
Args:
|
||||
collection_name: Collection name
|
||||
query_vector: Query embedding vector
|
||||
current_timestamp: Current time in milliseconds
|
||||
limit: Maximum results
|
||||
score_threshold: Minimum similarity score
|
||||
|
||||
Returns:
|
||||
List of non-expired search results
|
||||
"""
|
||||
# Filter: ttl_expiry > current_timestamp (not expired)
|
||||
expiry_filter = Filter(
|
||||
must=[
|
||||
FieldCondition(
|
||||
key="ttl_expiry",
|
||||
range=Range(gt=current_timestamp)
|
||||
)
|
||||
]
|
||||
)
|
||||
|
||||
try:
|
||||
response = self.client.query_points(
|
||||
collection_name=collection_name,
|
||||
query=query_vector,
|
||||
limit=limit,
|
||||
score_threshold=score_threshold,
|
||||
query_filter=expiry_filter,
|
||||
with_payload=True
|
||||
)
|
||||
|
||||
return [
|
||||
{
|
||||
"id": str(point.id),
|
||||
"score": point.score,
|
||||
"payload": dict(point.payload)
|
||||
}
|
||||
for point in response.points
|
||||
]
|
||||
except Exception as e:
|
||||
logger.error(f"Volatile search failed: {e}", exc_info=True)
|
||||
return []
|
||||
|
||||
async def delete_expired_vectors(
|
||||
self,
|
||||
collection_name: str,
|
||||
current_timestamp: int
|
||||
) -> int:
|
||||
"""
|
||||
Delete all vectors where ttl_expiry < current_timestamp.
|
||||
|
||||
Args:
|
||||
collection_name: Collection name
|
||||
current_timestamp: Current time in milliseconds
|
||||
|
||||
Returns:
|
||||
Number of points deleted (approximate)
|
||||
"""
|
||||
# Filter: ttl_expiry < current_timestamp (expired)
|
||||
expiry_filter = Filter(
|
||||
must=[
|
||||
FieldCondition(
|
||||
key="ttl_expiry",
|
||||
range=Range(lt=current_timestamp)
|
||||
)
|
||||
]
|
||||
)
|
||||
|
||||
try:
|
||||
# First count how many will be deleted (scroll to count)
|
||||
count = 0
|
||||
offset = None
|
||||
while True:
|
||||
points, next_offset = self.client.scroll(
|
||||
collection_name=collection_name,
|
||||
scroll_filter=expiry_filter,
|
||||
limit=100,
|
||||
offset=offset,
|
||||
with_payload=False
|
||||
)
|
||||
count += len(points)
|
||||
if next_offset is None:
|
||||
break
|
||||
offset = next_offset
|
||||
|
||||
if count == 0:
|
||||
return 0
|
||||
|
||||
# Delete expired points
|
||||
self.client.delete(
|
||||
collection_name=collection_name,
|
||||
points_selector=expiry_filter
|
||||
)
|
||||
|
||||
logger.info(f"Deleted {count} expired vectors from {collection_name}")
|
||||
return count
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Failed to delete expired vectors: {e}", exc_info=True)
|
||||
return 0
|
||||
|
||||
async def get_volatile_collections(self) -> List[str]:
|
||||
"""
|
||||
Get all volatile collections (prefixed with 'volatile_').
|
||||
|
||||
Returns:
|
||||
List of volatile collection names
|
||||
"""
|
||||
try:
|
||||
collections = self.client.get_collections()
|
||||
return [
|
||||
c.name for c in collections.collections
|
||||
if c.name.startswith("volatile_")
|
||||
]
|
||||
except Exception as e:
|
||||
logger.error(f"Failed to list volatile collections: {e}", exc_info=True)
|
||||
return []
|
||||
@@ -35,16 +35,19 @@ class WikiJSClient:
|
||||
self.graphql_url = f"{self.base_url}/graphql"
|
||||
self.api_token = api_token
|
||||
self.client = httpx.AsyncClient(timeout=30.0)
|
||||
logger.info(f"Initialized Wiki.js client: {base_url} (using API token)")
|
||||
auth_mode = "with API token" if api_token else "without auth (open API)"
|
||||
logger.info(f"Initialized Wiki.js client: {base_url} ({auth_mode})")
|
||||
|
||||
async def close(self):
|
||||
"""Close HTTP client"""
|
||||
await self.client.aclose()
|
||||
|
||||
def _ensure_authenticated(self):
|
||||
"""Verify API token is configured."""
|
||||
if not self.api_token:
|
||||
raise Exception("Wiki.js API token not configured")
|
||||
def _get_headers(self) -> Dict[str, str]:
|
||||
"""Get request headers, optionally including auth token."""
|
||||
headers = {"Content-Type": "application/json"}
|
||||
if self.api_token:
|
||||
headers["Authorization"] = f"Bearer {self.api_token}"
|
||||
return headers
|
||||
|
||||
async def _execute_query(
|
||||
self,
|
||||
@@ -64,18 +67,12 @@ class WikiJSClient:
|
||||
Raises:
|
||||
Exception: If query fails or returns errors
|
||||
"""
|
||||
# Ensure API token is configured
|
||||
self._ensure_authenticated()
|
||||
|
||||
payload = {
|
||||
"query": query,
|
||||
"variables": variables or {}
|
||||
}
|
||||
|
||||
headers = {
|
||||
"Authorization": f"Bearer {self.api_token}",
|
||||
"Content-Type": "application/json"
|
||||
}
|
||||
headers = self._get_headers()
|
||||
|
||||
try:
|
||||
response = await self.client.post(
|
||||
|
||||
+1
-1
@@ -43,7 +43,7 @@ class Settings(BaseSettings):
|
||||
|
||||
# Wiki.js Configuration
|
||||
wikijs_url: str = Field(default="http://wiki:3000", description="Wiki.js URL")
|
||||
wiki_graphql_api: str = Field(..., description="Wiki.js GraphQL API token (JWT)")
|
||||
wiki_graphql_api: str = Field(default="", description="Wiki.js GraphQL API token (optional - API may be open)")
|
||||
# Legacy auth fields - kept for backwards compatibility but deprecated
|
||||
wikijs_username: str = Field(default="", description="Wiki.js username (deprecated, use wiki_graphql_api)")
|
||||
wikijs_password: str = Field(default="", description="Wiki.js password (deprecated, use wiki_graphql_api)")
|
||||
|
||||
+2
-1
@@ -51,7 +51,7 @@ app.add_middleware(
|
||||
from src.routers import (
|
||||
wiki, tools, graph, vector, hybrid_rag, consolidation,
|
||||
ingestion, entity_linking, webhooks, rag_search, content,
|
||||
maintenance
|
||||
maintenance, volatile
|
||||
)
|
||||
|
||||
app.include_router(wiki.router)
|
||||
@@ -66,6 +66,7 @@ app.include_router(webhooks.router)
|
||||
app.include_router(rag_search.router)
|
||||
app.include_router(content.router)
|
||||
app.include_router(maintenance.router)
|
||||
app.include_router(volatile.router)
|
||||
|
||||
# Mount static files directory for Wiki.js integration scripts
|
||||
static_dir = Path(__file__).parent.parent / "static"
|
||||
|
||||
@@ -14,13 +14,16 @@ class HybridRAGConfig(BaseModel):
|
||||
vector_limit: int = Field(default=10, ge=1, le=50, description="Max vector results")
|
||||
graph_limit: int = Field(default=10, ge=1, le=50, description="Max graph results")
|
||||
web_limit: int = Field(default=5, ge=1, le=20, description="Max web results")
|
||||
volatile_limit: int = Field(default=1, ge=1, le=5, description="Max volatile results (typically 1)")
|
||||
enable_vector: bool = Field(default=True, description="Enable vector search")
|
||||
enable_graph: bool = Field(default=True, description="Enable graph search")
|
||||
enable_web: bool = Field(default=True, description="Enable web search")
|
||||
enable_volatile: bool = Field(default=True, description="Enable volatile cache search")
|
||||
enable_reranking: bool = Field(default=True, description="Enable LLM re-ranking")
|
||||
enable_enrichment: bool = Field(default=True, description="Enable graph enrichment")
|
||||
final_result_count: int = Field(default=10, ge=1, le=50, description="Final results to return")
|
||||
rrf_k: int = Field(default=60, ge=1, le=100, description="RRF constant")
|
||||
volatile_threshold: float = Field(default=0.8, ge=0.5, le=1.0, description="Volatile similarity threshold")
|
||||
|
||||
|
||||
class RelatedDossier(BaseModel):
|
||||
@@ -34,7 +37,7 @@ class RelatedDossier(BaseModel):
|
||||
|
||||
class HybridRAGResult(BaseModel):
|
||||
"""Single result from HybridRAG query."""
|
||||
source_type: str = Field(..., description="Source: 'vector', 'graph', 'web'")
|
||||
source_type: str = Field(..., description="Source: 'wiki', 'web', 'volatile'")
|
||||
title: str
|
||||
content: str
|
||||
url: Optional[str] = Field(None, description="URL for web results")
|
||||
@@ -53,6 +56,7 @@ class TimingBreakdown(BaseModel):
|
||||
vector_ms: float = Field(..., description="Phase 1: Vector search")
|
||||
graph_ms: float = Field(..., description="Phase 1: Graph search")
|
||||
web_ms: float = Field(..., description="Phase 1: Web search")
|
||||
volatile_ms: float = Field(default=0, description="Phase 1: Volatile cache search")
|
||||
fusion_ms: float = Field(..., description="Phase 2: RRF fusion")
|
||||
enrichment_ms: float = Field(..., description="Phase 3: Graph enrichment")
|
||||
reranking_ms: float = Field(..., description="Phase 4: LLM re-ranking")
|
||||
|
||||
@@ -0,0 +1,130 @@
|
||||
"""
|
||||
Volatile memory models for Library Desk.
|
||||
|
||||
Provides models for ephemeral cached data with TTL - weather, news, financial data,
|
||||
transit schedules, and other time-sensitive external information.
|
||||
"""
|
||||
|
||||
from pydantic import BaseModel, Field
|
||||
from typing import Dict, Any, Optional, List
|
||||
from datetime import datetime
|
||||
from enum import Enum
|
||||
|
||||
|
||||
class VolatileNamespace(str, Enum):
|
||||
"""
|
||||
Predefined namespaces for volatile data.
|
||||
|
||||
Each namespace can have different default TTLs and refresh schedules.
|
||||
"""
|
||||
# Real-time external data
|
||||
WEATHER = "weather" # Current conditions, forecasts
|
||||
NEWS = "news" # Headlines, breaking news
|
||||
FINANCIAL = "financial" # Stock prices, exchange rates, crypto
|
||||
TRANSIT = "transit" # Train/bus schedules, delays, disruptions
|
||||
TRAFFIC = "traffic" # Commute times, road conditions
|
||||
AIR_QUALITY = "air_quality" # Pollution levels, pollen counts
|
||||
SPORTS = "sports" # Live scores, upcoming matches
|
||||
|
||||
# System/integration data
|
||||
SOCIAL = "social" # Social media mentions, notifications
|
||||
SYSTEM = "system" # Service health, infrastructure status
|
||||
|
||||
# Ephemeral context
|
||||
CONTEXT = "context" # Conversation context, session state
|
||||
CUSTOM = "custom" # User-defined volatile data
|
||||
|
||||
|
||||
# Default TTLs per namespace (in seconds)
|
||||
NAMESPACE_DEFAULT_TTL: Dict[str, int] = {
|
||||
VolatileNamespace.WEATHER: 1800, # 30 min - weather changes slowly
|
||||
VolatileNamespace.NEWS: 3600, # 1 hour - news cycles
|
||||
VolatileNamespace.FINANCIAL: 300, # 5 min - markets move fast
|
||||
VolatileNamespace.TRANSIT: 300, # 5 min - schedules update frequently
|
||||
VolatileNamespace.TRAFFIC: 600, # 10 min - traffic patterns
|
||||
VolatileNamespace.AIR_QUALITY: 3600, # 1 hour - air quality stable
|
||||
VolatileNamespace.SPORTS: 60, # 1 min - live scores
|
||||
VolatileNamespace.SOCIAL: 600, # 10 min - social notifications
|
||||
VolatileNamespace.SYSTEM: 60, # 1 min - system health
|
||||
VolatileNamespace.CONTEXT: 3600, # 1 hour - session context
|
||||
VolatileNamespace.CUSTOM: 3600, # 1 hour - default for custom
|
||||
}
|
||||
|
||||
|
||||
class VolatileRecord(BaseModel):
|
||||
"""
|
||||
A volatile cache record with TTL.
|
||||
|
||||
Volatile records are ephemeral data stored in Redis with automatic expiration.
|
||||
Used for weather, news, financial data, and other time-sensitive information.
|
||||
"""
|
||||
key: str = Field(..., description="Record key (e.g., 'rotterdam', 'nos-headlines')")
|
||||
namespace: str = Field(..., description="Namespace (e.g., 'weather', 'news', 'financial')")
|
||||
data: Dict[str, Any] = Field(..., description="Actual content/payload")
|
||||
source: Optional[str] = Field(None, description="Origin API/service (e.g., 'openweathermap', 'nos.nl')")
|
||||
created_at: datetime = Field(default_factory=datetime.utcnow, description="When record was created")
|
||||
updated_at: datetime = Field(default_factory=datetime.utcnow, description="When record was last updated")
|
||||
ttl: int = Field(..., ge=60, le=604800, description="Time-to-live in seconds (max 7 days)")
|
||||
refresh_schedule: Optional[str] = Field(None, description="Cron expression for scheduled refresh")
|
||||
user: str = Field(..., description="User identifier for multi-tenancy")
|
||||
|
||||
|
||||
class VolatileRecordCreate(BaseModel):
|
||||
"""Request model for creating/updating a volatile record."""
|
||||
data: Dict[str, Any] = Field(..., description="Content to store")
|
||||
source: Optional[str] = Field(None, description="Origin API/service")
|
||||
ttl: Optional[int] = Field(None, ge=60, le=604800, description="TTL in seconds (uses namespace default if not set)")
|
||||
refresh_schedule: Optional[str] = Field(None, description="Cron expression for scheduled refresh")
|
||||
|
||||
|
||||
class VolatileRecordResponse(BaseModel):
|
||||
"""Response model for a volatile record."""
|
||||
key: str = Field(..., description="Record key")
|
||||
namespace: str = Field(..., description="Namespace")
|
||||
data: Dict[str, Any] = Field(..., description="Stored content")
|
||||
source: Optional[str] = Field(None, description="Origin API/service")
|
||||
created_at: datetime = Field(..., description="Creation timestamp")
|
||||
updated_at: datetime = Field(..., description="Last update timestamp")
|
||||
ttl: int = Field(..., description="TTL in seconds")
|
||||
ttl_remaining: int = Field(..., description="Seconds until expiration")
|
||||
refresh_schedule: Optional[str] = Field(None, description="Cron expression if scheduled")
|
||||
user: str = Field(..., description="User identifier")
|
||||
|
||||
|
||||
class VolatileListResponse(BaseModel):
|
||||
"""Response model for listing volatile records."""
|
||||
namespace: str = Field(..., description="Namespace queried")
|
||||
keys: List[str] = Field(..., description="List of keys in namespace")
|
||||
count: int = Field(..., description="Number of keys")
|
||||
user: str = Field(..., description="User identifier")
|
||||
|
||||
|
||||
class VolatileScheduledResponse(BaseModel):
|
||||
"""Response model for records needing refresh."""
|
||||
records: List[VolatileRecordResponse] = Field(..., description="Records with refresh schedules")
|
||||
count: int = Field(..., description="Number of scheduled records")
|
||||
user: str = Field(..., description="User identifier")
|
||||
|
||||
|
||||
class VolatileStatsResponse(BaseModel):
|
||||
"""Response model for volatile cache statistics."""
|
||||
total_records: int = Field(..., description="Total volatile records for user")
|
||||
by_namespace: Dict[str, int] = Field(..., description="Record count per namespace")
|
||||
scheduled_count: int = Field(..., description="Records with refresh schedules")
|
||||
total_memory_bytes: Optional[int] = Field(None, description="Approximate memory usage")
|
||||
user: str = Field(..., description="User identifier")
|
||||
|
||||
|
||||
class VolatileDeleteResponse(BaseModel):
|
||||
"""Response model for delete operation."""
|
||||
key: str = Field(..., description="Deleted key")
|
||||
namespace: str = Field(..., description="Namespace")
|
||||
deleted: bool = Field(..., description="Whether record was found and deleted")
|
||||
user: str = Field(..., description="User identifier")
|
||||
|
||||
|
||||
class VolatileBulkDeleteResponse(BaseModel):
|
||||
"""Response model for bulk delete operations."""
|
||||
namespace: Optional[str] = Field(None, description="Namespace if namespace-wide delete")
|
||||
deleted_count: int = Field(..., description="Number of records deleted")
|
||||
user: str = Field(..., description="User identifier")
|
||||
@@ -1,7 +1,7 @@
|
||||
"""
|
||||
HybridRAG router for multi-source search API.
|
||||
|
||||
Provides endpoint for combining vector, graph, and web search
|
||||
Provides endpoint for combining vector, graph, volatile cache, and web search
|
||||
with RRF fusion and LLM re-ranking.
|
||||
"""
|
||||
|
||||
@@ -34,10 +34,12 @@ def get_hybrid_rag_service(
|
||||
"""Get HybridRAG service instance with all dependencies."""
|
||||
from src.services.vector_service import VectorService
|
||||
from src.services.graph_service import GraphService
|
||||
from src.services.volatile_service import VolatileCacheService
|
||||
|
||||
# Create component services
|
||||
vector_service = VectorService(qdrant_client, wiki_client, ollama_client)
|
||||
graph_service = GraphService(neo4j_client, wiki_client)
|
||||
volatile_service = VolatileCacheService(qdrant_client, ollama_client, settings)
|
||||
|
||||
# Create HybridRAG service
|
||||
return HybridRAGService(
|
||||
@@ -46,7 +48,8 @@ def get_hybrid_rag_service(
|
||||
searxng_client=searxng_client,
|
||||
ollama_client=ollama_client,
|
||||
content_extractor=content_extractor,
|
||||
settings=settings
|
||||
settings=settings,
|
||||
volatile_service=volatile_service
|
||||
)
|
||||
|
||||
|
||||
@@ -58,26 +61,28 @@ async def hybrid_search(
|
||||
api_key: str = Depends(verify_api_key)
|
||||
):
|
||||
"""
|
||||
Execute HybridRAG query combining vector, graph, and web search.
|
||||
Execute HybridRAG query combining vector, graph, volatile cache, and web search.
|
||||
|
||||
**6-Phase Pipeline:**
|
||||
1. **Query Enhancement**: Extract keywords/synonyms with LLM
|
||||
2. **Parallel Retrieval**: Search vector (Qdrant), graph (Neo4j), web (SearXNG)
|
||||
3. **RRF Fusion**: Merge results with Reciprocal Rank Fusion
|
||||
2. **Parallel Retrieval**: Search vector (Qdrant), graph (Neo4j), volatile cache, web (SearXNG)
|
||||
3. **RRF Fusion**: Merge results with Reciprocal Rank Fusion (volatile gets priority boost)
|
||||
4. **Enrichment**: Add related documents via shared entities
|
||||
5. **LLM Re-ranking**: Re-rank with mistral-nemo for relevance
|
||||
5. **LLM Re-ranking**: Re-rank with configured model for relevance
|
||||
6. **Context Formatting**: Format for LLM consumption
|
||||
7. **Persistence**: Store for Librarian knowledge consolidation
|
||||
|
||||
**Example Request:**
|
||||
```json
|
||||
{
|
||||
"query": "How does Docker orchestration work with Kubernetes?",
|
||||
"query": "What's the weather in Rotterdam?",
|
||||
"user": "jpmschweitzer",
|
||||
"config": {
|
||||
"vector_limit": 10,
|
||||
"graph_limit": 10,
|
||||
"web_limit": 5,
|
||||
"volatile_limit": 5,
|
||||
"enable_volatile": true,
|
||||
"enable_reranking": true,
|
||||
"final_result_count": 10
|
||||
}
|
||||
@@ -85,7 +90,7 @@ async def hybrid_search(
|
||||
```
|
||||
|
||||
**Returns:**
|
||||
- Ranked results from all sources
|
||||
- Ranked results from all sources (wiki, volatile, web)
|
||||
- Extracted keywords/synonyms
|
||||
- Related dossiers (via graph)
|
||||
- Formatted context for LLM
|
||||
|
||||
@@ -16,10 +16,12 @@ import time
|
||||
|
||||
from src.services.vector_service import VectorService
|
||||
from src.services.graph_service import GraphService
|
||||
from src.services.volatile_service import VolatileCacheService
|
||||
from src.core.dependencies import (
|
||||
VectorServiceDep, GraphServiceDep, WikiJSDep, RedisDep,
|
||||
verify_api_key
|
||||
QdrantDep, OllamaDep, verify_api_key
|
||||
)
|
||||
from src.config import get_settings
|
||||
from datetime import datetime, timezone
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
@@ -209,6 +211,15 @@ class ReconcileIndexResponse(BaseModel):
|
||||
total_duration_ms: float
|
||||
|
||||
|
||||
class VolatileCleanupResponse(BaseModel):
|
||||
"""Response from volatile cache cleanup operation."""
|
||||
success: bool
|
||||
collections_processed: int
|
||||
total_expired_purged: int
|
||||
by_collection: Dict[str, int] = Field(default_factory=dict)
|
||||
duration_ms: float
|
||||
|
||||
|
||||
# ========== Endpoints ==========
|
||||
|
||||
@router.post("/cleanup/vectors", response_model=VectorCleanupResponse)
|
||||
@@ -490,6 +501,61 @@ async def cleanup_all(
|
||||
raise HTTPException(status_code=500, detail=str(e))
|
||||
|
||||
|
||||
@router.post("/cleanup/volatile", response_model=VolatileCleanupResponse)
|
||||
async def cleanup_volatile(
|
||||
qdrant: QdrantDep = None,
|
||||
ollama: OllamaDep = None,
|
||||
api_key: str = Depends(verify_api_key)
|
||||
):
|
||||
"""
|
||||
Purge expired volatile cache records across all users.
|
||||
|
||||
Loops through all volatile_* collections and removes records where
|
||||
ttl_expiry < current_timestamp.
|
||||
|
||||
**Scheduler Task** - Recommended to run every 10 minutes.
|
||||
|
||||
**Scheduler Integration:**
|
||||
```json
|
||||
{
|
||||
"task_name": "volatile_cleanup",
|
||||
"schedule": "*/10 * * * *",
|
||||
"endpoint": "POST /maintenance/cleanup/volatile",
|
||||
"description": "Purge expired volatile cache records"
|
||||
}
|
||||
```
|
||||
"""
|
||||
start_time = time.time()
|
||||
|
||||
try:
|
||||
settings = get_settings()
|
||||
service = VolatileCacheService(
|
||||
qdrant_client=qdrant,
|
||||
ollama_client=ollama,
|
||||
settings=settings
|
||||
)
|
||||
|
||||
# Purge expired from all volatile collections
|
||||
results = await service.purge_all_expired()
|
||||
|
||||
total_purged = sum(results.values())
|
||||
duration_ms = (time.time() - start_time) * 1000
|
||||
|
||||
logger.info(f"Volatile cleanup complete: {total_purged} expired records purged from {len(results)} collections")
|
||||
|
||||
return VolatileCleanupResponse(
|
||||
success=True,
|
||||
collections_processed=len(results),
|
||||
total_expired_purged=total_purged,
|
||||
by_collection=results,
|
||||
duration_ms=duration_ms
|
||||
)
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Volatile cleanup failed: {e}", exc_info=True)
|
||||
raise HTTPException(status_code=500, detail=str(e))
|
||||
|
||||
|
||||
@router.get("/health", response_model=HealthCheckResponse)
|
||||
async def maintenance_health(
|
||||
user: str = Query(..., description="User identifier"),
|
||||
|
||||
@@ -0,0 +1,273 @@
|
||||
"""
|
||||
Volatile cache router for Library Desk API.
|
||||
|
||||
Endpoints for ephemeral cached data with TTL - weather, news, financial, etc.
|
||||
Data is stored as vectors in Qdrant for semantic search retrieval.
|
||||
"""
|
||||
|
||||
from fastapi import APIRouter, HTTPException, Depends, Query
|
||||
import logging
|
||||
|
||||
from src.models.volatile import (
|
||||
VolatileRecordCreate,
|
||||
VolatileRecordResponse,
|
||||
VolatileListResponse,
|
||||
VolatileScheduledResponse,
|
||||
VolatileStatsResponse,
|
||||
VolatileDeleteResponse,
|
||||
VolatileNamespace,
|
||||
NAMESPACE_DEFAULT_TTL,
|
||||
)
|
||||
from src.services.volatile_service import VolatileCacheService
|
||||
from src.core.dependencies import verify_api_key, QdrantDep, OllamaDep
|
||||
from src.core.multi_tenancy import DEFAULT_USER
|
||||
from src.config import get_settings
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
router = APIRouter(prefix="/volatile", tags=["Volatile Cache"])
|
||||
|
||||
|
||||
def get_volatile_service(qdrant: QdrantDep, ollama: OllamaDep) -> VolatileCacheService:
|
||||
"""Get volatile cache service instance."""
|
||||
settings = get_settings()
|
||||
return VolatileCacheService(
|
||||
qdrant_client=qdrant,
|
||||
ollama_client=ollama,
|
||||
settings=settings
|
||||
)
|
||||
|
||||
|
||||
@router.get("/stats", response_model=VolatileStatsResponse)
|
||||
async def get_stats(
|
||||
user: str = Query(default=DEFAULT_USER, description="User identifier"),
|
||||
qdrant: QdrantDep = None,
|
||||
ollama: OllamaDep = None,
|
||||
api_key: str = Depends(verify_api_key)
|
||||
):
|
||||
"""
|
||||
Get volatile cache statistics.
|
||||
|
||||
Returns counts of records by namespace and scheduled refresh info.
|
||||
"""
|
||||
service = get_volatile_service(qdrant, ollama)
|
||||
stats = await service.get_stats(user)
|
||||
|
||||
return VolatileStatsResponse(
|
||||
total_records=stats["total_records"],
|
||||
by_namespace=stats["by_namespace"],
|
||||
scheduled_count=stats["scheduled_count"],
|
||||
total_memory_bytes=None,
|
||||
user=user,
|
||||
)
|
||||
|
||||
|
||||
@router.get("/scheduled", response_model=VolatileScheduledResponse)
|
||||
async def get_scheduled(
|
||||
user: str = Query(default=DEFAULT_USER, description="User identifier"),
|
||||
qdrant: QdrantDep = None,
|
||||
ollama: OllamaDep = None,
|
||||
api_key: str = Depends(verify_api_key)
|
||||
):
|
||||
"""
|
||||
Get records with refresh schedules.
|
||||
|
||||
Used by scheduler to determine what volatile data needs refreshing.
|
||||
Returns all records that have a refresh_schedule cron expression set.
|
||||
"""
|
||||
service = get_volatile_service(qdrant, ollama)
|
||||
records = await service.get_scheduled(user)
|
||||
|
||||
return VolatileScheduledResponse(
|
||||
records=records,
|
||||
count=len(records),
|
||||
user=user,
|
||||
)
|
||||
|
||||
|
||||
@router.get("/namespaces")
|
||||
async def list_namespaces(
|
||||
api_key: str = Depends(verify_api_key)
|
||||
):
|
||||
"""
|
||||
List available namespaces and their default TTLs.
|
||||
|
||||
Returns predefined namespaces with their default TTL values.
|
||||
"""
|
||||
return {
|
||||
"namespaces": [
|
||||
{
|
||||
"name": ns.value,
|
||||
"default_ttl": NAMESPACE_DEFAULT_TTL.get(ns, 3600),
|
||||
"description": _get_namespace_description(ns),
|
||||
}
|
||||
for ns in VolatileNamespace
|
||||
]
|
||||
}
|
||||
|
||||
|
||||
def _get_namespace_description(ns: VolatileNamespace) -> str:
|
||||
"""Get human-readable description for namespace."""
|
||||
descriptions = {
|
||||
VolatileNamespace.WEATHER: "Weather conditions and forecasts",
|
||||
VolatileNamespace.NEWS: "Headlines and breaking news",
|
||||
VolatileNamespace.FINANCIAL: "Stock prices, exchange rates, crypto",
|
||||
VolatileNamespace.TRANSIT: "Train/bus schedules, delays",
|
||||
VolatileNamespace.TRAFFIC: "Commute times, road conditions",
|
||||
VolatileNamespace.AIR_QUALITY: "Pollution levels, pollen counts",
|
||||
VolatileNamespace.SPORTS: "Live scores, upcoming matches",
|
||||
VolatileNamespace.SOCIAL: "Social media mentions, notifications",
|
||||
VolatileNamespace.SYSTEM: "Service health, infrastructure status",
|
||||
VolatileNamespace.CONTEXT: "Conversation context, session state",
|
||||
VolatileNamespace.CUSTOM: "User-defined volatile data",
|
||||
}
|
||||
return descriptions.get(ns, "Custom namespace")
|
||||
|
||||
|
||||
@router.get("/search")
|
||||
async def search_volatile(
|
||||
q: str = Query(..., min_length=1, description="Search query"),
|
||||
user: str = Query(default=DEFAULT_USER, description="User identifier"),
|
||||
limit: int = Query(default=5, ge=1, le=20, description="Maximum results"),
|
||||
threshold: float = Query(default=0.75, ge=0.5, le=1.0, description="Minimum similarity score"),
|
||||
qdrant: QdrantDep = None,
|
||||
ollama: OllamaDep = None,
|
||||
api_key: str = Depends(verify_api_key)
|
||||
):
|
||||
"""
|
||||
Semantic search across volatile data.
|
||||
|
||||
Searches all volatile data for semantically similar content.
|
||||
Higher threshold = stricter matching.
|
||||
|
||||
**Example:**
|
||||
```
|
||||
GET /volatile/search?q=weather%20rotterdam&user=jpmschweitzer
|
||||
```
|
||||
"""
|
||||
service = get_volatile_service(qdrant, ollama)
|
||||
results = await service.search(user, q, limit=limit, score_threshold=threshold)
|
||||
|
||||
return {
|
||||
"query": q,
|
||||
"results": results,
|
||||
"count": len(results),
|
||||
"user": user,
|
||||
}
|
||||
|
||||
|
||||
@router.post("/store", response_model=VolatileRecordResponse)
|
||||
async def store_volatile(
|
||||
namespace: str = Query(..., description="Data namespace (weather, news, etc.)"),
|
||||
key: str = Query(..., description="Record key (e.g., 'rotterdam', 'nos-headlines')"),
|
||||
request: VolatileRecordCreate = None,
|
||||
user: str = Query(default=DEFAULT_USER, description="User identifier"),
|
||||
qdrant: QdrantDep = None,
|
||||
ollama: OllamaDep = None,
|
||||
api_key: str = Depends(verify_api_key)
|
||||
):
|
||||
"""
|
||||
Store volatile data.
|
||||
|
||||
Data is converted to natural language and embedded for semantic search.
|
||||
If the same namespace+key already exists, it will be updated.
|
||||
|
||||
**Example Request:**
|
||||
```json
|
||||
POST /volatile/store?namespace=weather&key=rotterdam
|
||||
{
|
||||
"data": {
|
||||
"temperature": 8,
|
||||
"conditions": "Cloudy",
|
||||
"humidity": 85
|
||||
},
|
||||
"source": "openweathermap",
|
||||
"ttl": 1800,
|
||||
"refresh_schedule": "0 * * * *"
|
||||
}
|
||||
```
|
||||
|
||||
**Refresh Schedule:**
|
||||
Optional cron expression for automatic refresh. The scheduler
|
||||
will query `/volatile/scheduled` and trigger refreshes.
|
||||
"""
|
||||
# Validate namespace if not custom
|
||||
if namespace != VolatileNamespace.CUSTOM:
|
||||
try:
|
||||
VolatileNamespace(namespace)
|
||||
except ValueError:
|
||||
valid = [ns.value for ns in VolatileNamespace]
|
||||
raise HTTPException(
|
||||
status_code=400,
|
||||
detail=f"Invalid namespace '{namespace}'. Valid: {valid}"
|
||||
)
|
||||
|
||||
service = get_volatile_service(qdrant, ollama)
|
||||
|
||||
try:
|
||||
record = await service.store(
|
||||
user=user,
|
||||
namespace=namespace,
|
||||
key=key,
|
||||
data=request.data,
|
||||
source=request.source,
|
||||
ttl=request.ttl,
|
||||
refresh_schedule=request.refresh_schedule,
|
||||
)
|
||||
return record
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Failed to store volatile record: {e}")
|
||||
raise HTTPException(status_code=500, detail=f"Failed to store record: {str(e)}")
|
||||
|
||||
|
||||
@router.get("/{namespace}/{key}", response_model=VolatileRecordResponse)
|
||||
async def get_record(
|
||||
namespace: str,
|
||||
key: str,
|
||||
user: str = Query(default=DEFAULT_USER, description="User identifier"),
|
||||
qdrant: QdrantDep = None,
|
||||
ollama: OllamaDep = None,
|
||||
api_key: str = Depends(verify_api_key)
|
||||
):
|
||||
"""
|
||||
Get a specific volatile record by namespace and key.
|
||||
|
||||
**Example:**
|
||||
```
|
||||
GET /volatile/weather/rotterdam?user=jpmschweitzer
|
||||
```
|
||||
"""
|
||||
service = get_volatile_service(qdrant, ollama)
|
||||
record = await service.get(user, namespace, key)
|
||||
|
||||
if not record:
|
||||
raise HTTPException(
|
||||
status_code=404,
|
||||
detail=f"Record '{key}' not found in namespace '{namespace}'"
|
||||
)
|
||||
|
||||
return record
|
||||
|
||||
|
||||
@router.delete("/{namespace}/{key}", response_model=VolatileDeleteResponse)
|
||||
async def delete_record(
|
||||
namespace: str,
|
||||
key: str,
|
||||
user: str = Query(default=DEFAULT_USER, description="User identifier"),
|
||||
qdrant: QdrantDep = None,
|
||||
ollama: OllamaDep = None,
|
||||
api_key: str = Depends(verify_api_key)
|
||||
):
|
||||
"""
|
||||
Delete a specific volatile record.
|
||||
"""
|
||||
service = get_volatile_service(qdrant, ollama)
|
||||
deleted = await service.delete(user, namespace, key)
|
||||
|
||||
return VolatileDeleteResponse(
|
||||
key=key,
|
||||
namespace=namespace,
|
||||
deleted=deleted,
|
||||
user=user,
|
||||
)
|
||||
@@ -20,6 +20,7 @@ import logging
|
||||
|
||||
from src.services.vector_service import VectorService
|
||||
from src.services.graph_service import GraphService
|
||||
from src.services.volatile_service import VolatileCacheService
|
||||
from src.clients.searxng_client import SearXNGClient
|
||||
from src.clients.ollama_client import OllamaClient
|
||||
from src.clients.content_extractor import ContentExtractor
|
||||
@@ -46,7 +47,8 @@ class HybridRAGService:
|
||||
searxng_client: SearXNGClient,
|
||||
ollama_client: OllamaClient,
|
||||
content_extractor: ContentExtractor,
|
||||
settings: Settings
|
||||
settings: Settings,
|
||||
volatile_service: Optional[VolatileCacheService] = None
|
||||
):
|
||||
"""
|
||||
Initialize HybridRAG service.
|
||||
@@ -58,6 +60,7 @@ class HybridRAGService:
|
||||
ollama_client: Client for LLM (keyword extraction, re-ranking)
|
||||
content_extractor: Client for extracting full content from URLs
|
||||
settings: Application settings
|
||||
volatile_service: Service for volatile cache search (optional)
|
||||
"""
|
||||
self.vector = vector_service
|
||||
self.graph = graph_service
|
||||
@@ -65,6 +68,7 @@ class HybridRAGService:
|
||||
self.ollama = ollama_client
|
||||
self.content_extractor = content_extractor
|
||||
self.settings = settings
|
||||
self.volatile = volatile_service
|
||||
self.reranker_model = settings.ollama_model
|
||||
|
||||
async def search(
|
||||
@@ -104,8 +108,9 @@ class HybridRAGService:
|
||||
timing["vector_ms"] = raw_results.get("timing", {}).get("vector_ms", 0)
|
||||
timing["graph_ms"] = raw_results.get("timing", {}).get("graph_ms", 0)
|
||||
timing["web_ms"] = raw_results.get("timing", {}).get("web_ms", 0)
|
||||
timing["volatile_ms"] = raw_results.get("timing", {}).get("volatile_ms", 0)
|
||||
|
||||
# Phase 2: Two-Stage RRF Fusion
|
||||
# Phase 2: Three-Source RRF Fusion
|
||||
phase2_start = time.time()
|
||||
|
||||
# Stage 1: Merge wiki sources (vector + graph) into single ranking
|
||||
@@ -115,10 +120,12 @@ class HybridRAGService:
|
||||
k=config.rrf_k
|
||||
)
|
||||
|
||||
# Stage 2: Final RRF between wiki and web (equal footing)
|
||||
# Stage 2: Final RRF between wiki, volatile, and web
|
||||
# Volatile gets priority boost (smaller k = higher contribution per rank)
|
||||
fused_results = self._reciprocal_rank_fusion(
|
||||
wiki_results=wiki_merged,
|
||||
web_results=raw_results.get("web", []),
|
||||
volatile_results=raw_results.get("volatile", []),
|
||||
k=config.rrf_k
|
||||
)
|
||||
timing["fusion_ms"] = (time.time() - phase2_start) * 1000
|
||||
@@ -389,6 +396,37 @@ JSON:"""
|
||||
|
||||
tasks["web"] = web_search()
|
||||
|
||||
# Volatile cache search
|
||||
if config.enable_volatile and self.volatile:
|
||||
async def volatile_search():
|
||||
start = time.time()
|
||||
try:
|
||||
results = await self.volatile.search(
|
||||
user=user,
|
||||
query=query,
|
||||
limit=config.volatile_limit,
|
||||
score_threshold=config.volatile_threshold
|
||||
)
|
||||
formatted = [
|
||||
{
|
||||
"key": r.key,
|
||||
"namespace": r.namespace,
|
||||
"title": f"{r.namespace}: {r.key}",
|
||||
"content": r.data.get("text", "") if isinstance(r.data, dict) else str(r.data),
|
||||
"raw_data": r.data,
|
||||
"source_api": r.source,
|
||||
"ttl_remaining": r.ttl_remaining,
|
||||
"source": "volatile"
|
||||
}
|
||||
for r in results
|
||||
]
|
||||
return formatted, (time.time() - start) * 1000
|
||||
except Exception as e:
|
||||
logger.error(f"Volatile search failed: {e}", exc_info=True)
|
||||
return [], (time.time() - start) * 1000
|
||||
|
||||
tasks["volatile"] = volatile_search()
|
||||
|
||||
# Execute all searches in parallel
|
||||
results_dict = await asyncio.gather(*tasks.values())
|
||||
|
||||
@@ -401,7 +439,8 @@ JSON:"""
|
||||
|
||||
logger.info(
|
||||
f"Parallel retrieval: vector={len(output.get('vector', []))}, "
|
||||
f"graph={len(output.get('graph', []))}, web={len(output.get('web', []))}"
|
||||
f"graph={len(output.get('graph', []))}, web={len(output.get('web', []))}, "
|
||||
f"volatile={len(output.get('volatile', []))}"
|
||||
)
|
||||
|
||||
return output
|
||||
@@ -491,23 +530,42 @@ JSON:"""
|
||||
self,
|
||||
wiki_results: List[Dict],
|
||||
web_results: List[Dict],
|
||||
volatile_results: Optional[List[Dict]] = None,
|
||||
k: int = 60
|
||||
) -> List[Dict[str, Any]]:
|
||||
"""
|
||||
Stage 2: Final RRF between wiki (single source) and web.
|
||||
Stage 2: Final RRF between wiki, volatile, and web.
|
||||
|
||||
Wiki results are pre-merged from vector+graph, so wiki and web
|
||||
now compete on equal footing.
|
||||
Wiki results are pre-merged from vector+graph. Volatile results
|
||||
get a priority boost (smaller effective k) since they represent
|
||||
current, time-sensitive information.
|
||||
|
||||
Args:
|
||||
wiki_results: Pre-merged wiki results from _merge_wiki_sources()
|
||||
web_results: Results from web search
|
||||
volatile_results: Results from volatile cache (fresh data)
|
||||
k: RRF constant (default 60)
|
||||
|
||||
Returns:
|
||||
Final merged and sorted results
|
||||
"""
|
||||
rrf_scores = {}
|
||||
volatile_results = volatile_results or []
|
||||
|
||||
# Volatile results get priority boost (k/2 = stronger score per rank)
|
||||
volatile_k = k // 2
|
||||
for rank, result in enumerate(volatile_results, start=1):
|
||||
key = result.get("key")
|
||||
namespace = result.get("namespace", "unknown")
|
||||
if not key:
|
||||
continue
|
||||
result_id = f"volatile_{namespace}_{key}"
|
||||
rrf_scores[result_id] = {
|
||||
"result": result,
|
||||
"rrf_score": 1 / (volatile_k + rank), # Priority boost
|
||||
"sources": ["volatile"],
|
||||
"source_type": "volatile"
|
||||
}
|
||||
|
||||
# Wiki results (single source, already merged)
|
||||
for rank, result in enumerate(wiki_results, start=1):
|
||||
@@ -542,7 +600,8 @@ JSON:"""
|
||||
reverse=True
|
||||
)
|
||||
|
||||
logger.info(f"Final RRF: {len(sorted_results)} results (wiki + web)")
|
||||
volatile_count = len([r for r in sorted_results if r["source_type"] == "volatile"])
|
||||
logger.info(f"Final RRF: {len(sorted_results)} results (wiki + volatile[{volatile_count}] + web)")
|
||||
|
||||
return sorted_results
|
||||
|
||||
|
||||
@@ -0,0 +1,564 @@
|
||||
"""
|
||||
Volatile Cache service for Library Desk.
|
||||
|
||||
Provides ephemeral data storage with TTL using Qdrant vectors:
|
||||
- Weather, news, financial data
|
||||
- Transit schedules, traffic conditions
|
||||
- System status, social notifications
|
||||
|
||||
Data is stored as embedded vectors for semantic search retrieval.
|
||||
"""
|
||||
|
||||
import hashlib
|
||||
import logging
|
||||
import time
|
||||
from datetime import datetime
|
||||
from typing import List, Optional, Dict, Any
|
||||
|
||||
from src.clients.qdrant_client import QdrantClientWrapper
|
||||
from src.clients.ollama_client import OllamaClient
|
||||
from src.config import Settings
|
||||
from src.models.volatile import (
|
||||
VolatileRecordResponse,
|
||||
VolatileNamespace,
|
||||
NAMESPACE_DEFAULT_TTL,
|
||||
)
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class VolatileCacheService:
|
||||
"""
|
||||
Service for volatile data with TTL stored in Qdrant.
|
||||
|
||||
Stores ephemeral data as vectors for semantic search retrieval.
|
||||
Each user has an isolated volatile collection.
|
||||
"""
|
||||
|
||||
COLLECTION_PREFIX = "volatile_"
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
qdrant_client: QdrantClientWrapper,
|
||||
ollama_client: OllamaClient,
|
||||
settings: Settings
|
||||
):
|
||||
"""
|
||||
Initialize volatile cache service.
|
||||
|
||||
Args:
|
||||
qdrant_client: Qdrant client for vector storage
|
||||
ollama_client: Ollama client for embeddings
|
||||
settings: Application settings
|
||||
"""
|
||||
self.qdrant = qdrant_client
|
||||
self.ollama = ollama_client
|
||||
self.settings = settings
|
||||
|
||||
logger.info("Initialized VolatileCacheService (Qdrant backend)")
|
||||
|
||||
def _collection_name(self, user: str) -> str:
|
||||
"""Get volatile collection name for user."""
|
||||
return f"{self.COLLECTION_PREFIX}{user}"
|
||||
|
||||
def _make_vector_id(self, namespace: str, key: str) -> str:
|
||||
"""
|
||||
Generate deterministic vector ID for namespace/key.
|
||||
|
||||
Same namespace+key always produces same ID for upsert behavior.
|
||||
"""
|
||||
combined = f"{namespace}:{key}"
|
||||
return hashlib.md5(combined.encode()).hexdigest()
|
||||
|
||||
def _get_default_ttl(self, namespace: str) -> int:
|
||||
"""Get default TTL for a namespace."""
|
||||
try:
|
||||
ns = VolatileNamespace(namespace)
|
||||
return NAMESPACE_DEFAULT_TTL.get(ns, self.settings.volatile_default_ttl)
|
||||
except ValueError:
|
||||
return self.settings.volatile_default_ttl
|
||||
|
||||
def _current_timestamp_ms(self) -> int:
|
||||
"""Get current timestamp in milliseconds."""
|
||||
return int(time.time() * 1000)
|
||||
|
||||
def _to_natural_language(
|
||||
self,
|
||||
namespace: str,
|
||||
key: str,
|
||||
data: Dict[str, Any]
|
||||
) -> str:
|
||||
"""
|
||||
Convert structured data to natural language for embedding.
|
||||
|
||||
This creates a text representation that embeds well semantically.
|
||||
"""
|
||||
# Template-based conversion for known namespaces
|
||||
if namespace == VolatileNamespace.WEATHER:
|
||||
temp = data.get("temperature", data.get("temp", "unknown"))
|
||||
conditions = data.get("conditions", data.get("weather", ""))
|
||||
humidity = data.get("humidity", "")
|
||||
text = f"Current weather in {key}: {temp}°C"
|
||||
if conditions:
|
||||
text += f", {conditions}"
|
||||
if humidity:
|
||||
text += f", humidity {humidity}%"
|
||||
return text
|
||||
|
||||
elif namespace == VolatileNamespace.NEWS:
|
||||
title = data.get("title", data.get("headline", ""))
|
||||
summary = data.get("summary", data.get("description", ""))
|
||||
source = data.get("source", "")
|
||||
text = f"News: {title}"
|
||||
if summary:
|
||||
text += f". {summary}"
|
||||
if source:
|
||||
text += f" (Source: {source})"
|
||||
return text
|
||||
|
||||
elif namespace == VolatileNamespace.FINANCIAL:
|
||||
symbol = data.get("symbol", key)
|
||||
price = data.get("price", "")
|
||||
change = data.get("change", data.get("change_percent", ""))
|
||||
text = f"Financial data for {symbol}"
|
||||
if price:
|
||||
text += f": price {price}"
|
||||
if change:
|
||||
text += f", change {change}%"
|
||||
return text
|
||||
|
||||
elif namespace == VolatileNamespace.TRANSIT:
|
||||
route = data.get("route", data.get("line", key))
|
||||
status = data.get("status", "")
|
||||
delay = data.get("delay", data.get("delay_minutes", ""))
|
||||
text = f"Transit {route}"
|
||||
if status:
|
||||
text += f": {status}"
|
||||
if delay:
|
||||
text += f", delay {delay} minutes"
|
||||
return text
|
||||
|
||||
elif namespace == VolatileNamespace.TRAFFIC:
|
||||
location = data.get("location", key)
|
||||
duration = data.get("duration", data.get("travel_time", ""))
|
||||
congestion = data.get("congestion", "")
|
||||
text = f"Traffic for {location}"
|
||||
if duration:
|
||||
text += f": {duration} minutes"
|
||||
if congestion:
|
||||
text += f", congestion level {congestion}"
|
||||
return text
|
||||
|
||||
elif namespace == VolatileNamespace.AIR_QUALITY:
|
||||
location = data.get("location", key)
|
||||
aqi = data.get("aqi", data.get("index", ""))
|
||||
quality = data.get("quality", "")
|
||||
text = f"Air quality in {location}"
|
||||
if aqi:
|
||||
text += f": AQI {aqi}"
|
||||
if quality:
|
||||
text += f" ({quality})"
|
||||
return text
|
||||
|
||||
elif namespace == VolatileNamespace.SPORTS:
|
||||
event = data.get("event", data.get("match", key))
|
||||
score = data.get("score", "")
|
||||
status = data.get("status", "")
|
||||
text = f"Sports: {event}"
|
||||
if score:
|
||||
text += f" - Score: {score}"
|
||||
if status:
|
||||
text += f" ({status})"
|
||||
return text
|
||||
|
||||
elif namespace == VolatileNamespace.SYSTEM:
|
||||
service = data.get("service", key)
|
||||
status = data.get("status", "unknown")
|
||||
message = data.get("message", "")
|
||||
text = f"System status for {service}: {status}"
|
||||
if message:
|
||||
text += f". {message}"
|
||||
return text
|
||||
|
||||
# Fallback: serialize key fields
|
||||
text_parts = [f"{namespace} data for {key}:"]
|
||||
for k, v in data.items():
|
||||
if isinstance(v, (str, int, float, bool)):
|
||||
text_parts.append(f"{k}: {v}")
|
||||
return " ".join(text_parts)
|
||||
|
||||
async def store(
|
||||
self,
|
||||
user: str,
|
||||
namespace: str,
|
||||
key: str,
|
||||
data: Dict[str, Any],
|
||||
source: Optional[str] = None,
|
||||
ttl: Optional[int] = None,
|
||||
refresh_schedule: Optional[str] = None
|
||||
) -> VolatileRecordResponse:
|
||||
"""
|
||||
Store volatile data as an embedded vector.
|
||||
|
||||
Args:
|
||||
user: User identifier
|
||||
namespace: Data namespace (from controlled list)
|
||||
key: Record key (normalized slug)
|
||||
data: Structured data to store
|
||||
source: Origin API/service
|
||||
ttl: TTL in seconds (uses namespace default if not set)
|
||||
refresh_schedule: Optional cron expression for refresh
|
||||
|
||||
Returns:
|
||||
The stored record
|
||||
"""
|
||||
collection = self._collection_name(user)
|
||||
|
||||
# Ensure collection exists
|
||||
await self.qdrant.ensure_collection(collection)
|
||||
|
||||
# Calculate TTL and expiry
|
||||
effective_ttl = ttl if ttl is not None else self._get_default_ttl(namespace)
|
||||
now_ms = self._current_timestamp_ms()
|
||||
expiry_ms = now_ms + (effective_ttl * 1000)
|
||||
|
||||
# Convert to natural language for embedding
|
||||
text = self._to_natural_language(namespace, key, data)
|
||||
|
||||
# Generate embedding
|
||||
embedding = await self.ollama.embed(text)
|
||||
if not embedding:
|
||||
raise ValueError("Failed to generate embedding for volatile data")
|
||||
|
||||
# Build payload
|
||||
now = datetime.utcnow()
|
||||
payload = {
|
||||
"doc_type": "volatile",
|
||||
"namespace": namespace,
|
||||
"key": key,
|
||||
"text": text,
|
||||
"raw_data": data,
|
||||
"source": source,
|
||||
"created_at": now.isoformat(),
|
||||
"updated_at": now.isoformat(),
|
||||
"ttl": effective_ttl,
|
||||
"ttl_expiry": expiry_ms,
|
||||
"refresh_schedule": refresh_schedule,
|
||||
"user": user,
|
||||
}
|
||||
|
||||
# Upsert vector (same namespace+key = same ID = update)
|
||||
vector_id = self._make_vector_id(namespace, key)
|
||||
success = await self.qdrant.upsert_vector(
|
||||
collection_name=collection,
|
||||
vector_id=vector_id,
|
||||
vector=embedding,
|
||||
payload=payload
|
||||
)
|
||||
|
||||
if not success:
|
||||
raise ValueError("Failed to store volatile vector")
|
||||
|
||||
logger.debug(f"Stored volatile {namespace}:{key} with TTL {effective_ttl}s")
|
||||
|
||||
return VolatileRecordResponse(
|
||||
key=key,
|
||||
namespace=namespace,
|
||||
data=data,
|
||||
source=source,
|
||||
created_at=now,
|
||||
updated_at=now,
|
||||
ttl=effective_ttl,
|
||||
ttl_remaining=effective_ttl,
|
||||
refresh_schedule=refresh_schedule,
|
||||
user=user,
|
||||
)
|
||||
|
||||
async def search(
|
||||
self,
|
||||
user: str,
|
||||
query: str,
|
||||
limit: int = 5,
|
||||
score_threshold: float = 0.75
|
||||
) -> List[VolatileRecordResponse]:
|
||||
"""
|
||||
Semantic search across volatile data.
|
||||
|
||||
Args:
|
||||
user: User identifier
|
||||
query: Search query
|
||||
limit: Maximum results
|
||||
score_threshold: Minimum similarity score (higher = stricter)
|
||||
|
||||
Returns:
|
||||
List of matching volatile records
|
||||
"""
|
||||
collection = self._collection_name(user)
|
||||
|
||||
# Check if collection exists
|
||||
if not await self.qdrant.collection_exists(collection):
|
||||
return []
|
||||
|
||||
# Generate query embedding
|
||||
query_embedding = await self.ollama.embed(query)
|
||||
if not query_embedding:
|
||||
logger.error("Failed to embed query for volatile search")
|
||||
return []
|
||||
|
||||
# Search with expiry filter
|
||||
now_ms = self._current_timestamp_ms()
|
||||
results = await self.qdrant.search_with_expiry_filter(
|
||||
collection_name=collection,
|
||||
query_vector=query_embedding,
|
||||
current_timestamp=now_ms,
|
||||
limit=limit,
|
||||
score_threshold=score_threshold
|
||||
)
|
||||
|
||||
# Convert to response models
|
||||
responses = []
|
||||
for result in results:
|
||||
payload = result["payload"]
|
||||
ttl_expiry = payload.get("ttl_expiry", 0)
|
||||
ttl_remaining = max(0, (ttl_expiry - now_ms) // 1000)
|
||||
|
||||
responses.append(VolatileRecordResponse(
|
||||
key=payload["key"],
|
||||
namespace=payload["namespace"],
|
||||
data=payload.get("raw_data", {}),
|
||||
source=payload.get("source"),
|
||||
created_at=datetime.fromisoformat(payload["created_at"]),
|
||||
updated_at=datetime.fromisoformat(payload["updated_at"]),
|
||||
ttl=payload.get("ttl", 0),
|
||||
ttl_remaining=ttl_remaining,
|
||||
refresh_schedule=payload.get("refresh_schedule"),
|
||||
user=payload["user"],
|
||||
))
|
||||
|
||||
return responses
|
||||
|
||||
async def get(
|
||||
self,
|
||||
user: str,
|
||||
namespace: str,
|
||||
key: str
|
||||
) -> Optional[VolatileRecordResponse]:
|
||||
"""
|
||||
Get a specific volatile record by namespace and key.
|
||||
|
||||
Args:
|
||||
user: User identifier
|
||||
namespace: Data namespace
|
||||
key: Record key
|
||||
|
||||
Returns:
|
||||
Record if found and not expired, None otherwise
|
||||
"""
|
||||
# Use search with high threshold to find exact match
|
||||
query = self._to_natural_language(namespace, key, {"key": key})
|
||||
results = await self.search(user, query, limit=10, score_threshold=0.5)
|
||||
|
||||
# Find exact namespace+key match
|
||||
for result in results:
|
||||
if result.namespace == namespace and result.key == key:
|
||||
return result
|
||||
|
||||
return None
|
||||
|
||||
async def delete(
|
||||
self,
|
||||
user: str,
|
||||
namespace: str,
|
||||
key: str
|
||||
) -> bool:
|
||||
"""
|
||||
Delete a specific volatile record.
|
||||
|
||||
Args:
|
||||
user: User identifier
|
||||
namespace: Data namespace
|
||||
key: Record key
|
||||
|
||||
Returns:
|
||||
True if deleted, False if not found
|
||||
"""
|
||||
collection = self._collection_name(user)
|
||||
|
||||
if not await self.qdrant.collection_exists(collection):
|
||||
return False
|
||||
|
||||
vector_id = self._make_vector_id(namespace, key)
|
||||
|
||||
try:
|
||||
deleted = await self.qdrant.delete_by_ids(
|
||||
collection_name=collection,
|
||||
point_ids=[vector_id]
|
||||
)
|
||||
return deleted > 0
|
||||
except Exception as e:
|
||||
logger.error(f"Failed to delete volatile {namespace}:{key}: {e}")
|
||||
return False
|
||||
|
||||
async def get_scheduled(
|
||||
self,
|
||||
user: str
|
||||
) -> List[VolatileRecordResponse]:
|
||||
"""
|
||||
Get all records with refresh schedules.
|
||||
|
||||
Used by scheduler to determine what needs refreshing.
|
||||
|
||||
Args:
|
||||
user: User identifier
|
||||
|
||||
Returns:
|
||||
List of records with refresh_schedule set
|
||||
"""
|
||||
collection = self._collection_name(user)
|
||||
|
||||
if not await self.qdrant.collection_exists(collection):
|
||||
return []
|
||||
|
||||
now_ms = self._current_timestamp_ms()
|
||||
scheduled = []
|
||||
|
||||
# Scroll through all non-expired records
|
||||
try:
|
||||
all_points = await self.qdrant.scroll_all_points(
|
||||
collection_name=collection,
|
||||
with_payload=True
|
||||
)
|
||||
|
||||
for point in all_points:
|
||||
payload = point.get("payload", {})
|
||||
ttl_expiry = payload.get("ttl_expiry", 0)
|
||||
|
||||
# Skip expired
|
||||
if ttl_expiry <= now_ms:
|
||||
continue
|
||||
|
||||
# Only include if has refresh schedule
|
||||
if payload.get("refresh_schedule"):
|
||||
ttl_remaining = max(0, (ttl_expiry - now_ms) // 1000)
|
||||
scheduled.append(VolatileRecordResponse(
|
||||
key=payload["key"],
|
||||
namespace=payload["namespace"],
|
||||
data=payload.get("raw_data", {}),
|
||||
source=payload.get("source"),
|
||||
created_at=datetime.fromisoformat(payload["created_at"]),
|
||||
updated_at=datetime.fromisoformat(payload["updated_at"]),
|
||||
ttl=payload.get("ttl", 0),
|
||||
ttl_remaining=ttl_remaining,
|
||||
refresh_schedule=payload["refresh_schedule"],
|
||||
user=payload["user"],
|
||||
))
|
||||
|
||||
return scheduled
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Failed to get scheduled volatile records: {e}")
|
||||
return []
|
||||
|
||||
async def get_stats(
|
||||
self,
|
||||
user: str
|
||||
) -> Dict[str, Any]:
|
||||
"""
|
||||
Get cache statistics for user.
|
||||
|
||||
Args:
|
||||
user: User identifier
|
||||
|
||||
Returns:
|
||||
Statistics dict
|
||||
"""
|
||||
collection = self._collection_name(user)
|
||||
|
||||
if not await self.qdrant.collection_exists(collection):
|
||||
return {
|
||||
"total_records": 0,
|
||||
"by_namespace": {},
|
||||
"scheduled_count": 0,
|
||||
"expired_count": 0,
|
||||
}
|
||||
|
||||
now_ms = self._current_timestamp_ms()
|
||||
by_namespace: Dict[str, int] = {}
|
||||
total = 0
|
||||
scheduled = 0
|
||||
expired = 0
|
||||
|
||||
try:
|
||||
all_points = await self.qdrant.scroll_all_points(
|
||||
collection_name=collection,
|
||||
with_payload=True
|
||||
)
|
||||
|
||||
for point in all_points:
|
||||
payload = point.get("payload", {})
|
||||
namespace = payload.get("namespace", "unknown")
|
||||
ttl_expiry = payload.get("ttl_expiry", 0)
|
||||
|
||||
if ttl_expiry <= now_ms:
|
||||
expired += 1
|
||||
else:
|
||||
total += 1
|
||||
by_namespace[namespace] = by_namespace.get(namespace, 0) + 1
|
||||
if payload.get("refresh_schedule"):
|
||||
scheduled += 1
|
||||
|
||||
return {
|
||||
"total_records": total,
|
||||
"by_namespace": by_namespace,
|
||||
"scheduled_count": scheduled,
|
||||
"expired_count": expired,
|
||||
}
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Failed to get volatile stats: {e}")
|
||||
return {
|
||||
"total_records": 0,
|
||||
"by_namespace": {},
|
||||
"scheduled_count": 0,
|
||||
"expired_count": 0,
|
||||
}
|
||||
|
||||
async def purge_expired(
|
||||
self,
|
||||
user: str
|
||||
) -> int:
|
||||
"""
|
||||
Purge all expired volatile records for user.
|
||||
|
||||
Args:
|
||||
user: User identifier
|
||||
|
||||
Returns:
|
||||
Number of records purged
|
||||
"""
|
||||
collection = self._collection_name(user)
|
||||
|
||||
if not await self.qdrant.collection_exists(collection):
|
||||
return 0
|
||||
|
||||
now_ms = self._current_timestamp_ms()
|
||||
return await self.qdrant.delete_expired_vectors(collection, now_ms)
|
||||
|
||||
async def purge_all_expired(self) -> Dict[str, int]:
|
||||
"""
|
||||
Purge expired records from all volatile collections.
|
||||
|
||||
Returns:
|
||||
Dict of collection -> purged count
|
||||
"""
|
||||
collections = await self.qdrant.get_volatile_collections()
|
||||
results = {}
|
||||
now_ms = self._current_timestamp_ms()
|
||||
|
||||
for collection in collections:
|
||||
purged = await self.qdrant.delete_expired_vectors(collection, now_ms)
|
||||
if purged > 0:
|
||||
results[collection] = purged
|
||||
logger.info(f"Purged {purged} expired from {collection}")
|
||||
|
||||
return results
|
||||
@@ -0,0 +1,567 @@
|
||||
"""
|
||||
Tests for volatile cache router and service (Qdrant backend).
|
||||
|
||||
Tests:
|
||||
- Volatile record CRUD operations
|
||||
- Namespace listing and management
|
||||
- Scheduled record retrieval
|
||||
- TTL behavior and expiry filtering
|
||||
- Semantic search
|
||||
- Natural language conversion
|
||||
"""
|
||||
|
||||
import pytest
|
||||
from datetime import datetime
|
||||
from unittest.mock import AsyncMock, MagicMock, patch
|
||||
|
||||
from src.models.volatile import (
|
||||
VolatileRecord,
|
||||
VolatileRecordCreate,
|
||||
VolatileRecordResponse,
|
||||
VolatileListResponse,
|
||||
VolatileScheduledResponse,
|
||||
VolatileStatsResponse,
|
||||
VolatileDeleteResponse,
|
||||
VolatileBulkDeleteResponse,
|
||||
VolatileNamespace,
|
||||
NAMESPACE_DEFAULT_TTL,
|
||||
)
|
||||
|
||||
|
||||
class TestVolatileModels:
|
||||
"""Test volatile data models."""
|
||||
|
||||
def test_volatile_record_creation(self):
|
||||
"""Test VolatileRecord model creation."""
|
||||
record = VolatileRecord(
|
||||
key="rotterdam",
|
||||
namespace="weather",
|
||||
data={"temperature": 18, "conditions": "Cloudy"},
|
||||
source="openweathermap",
|
||||
ttl=1800,
|
||||
user="jpmschweitzer",
|
||||
)
|
||||
assert record.key == "rotterdam"
|
||||
assert record.namespace == "weather"
|
||||
assert record.data["temperature"] == 18
|
||||
assert record.ttl == 1800
|
||||
assert record.refresh_schedule is None
|
||||
|
||||
def test_volatile_record_with_schedule(self):
|
||||
"""Test VolatileRecord with refresh schedule."""
|
||||
record = VolatileRecord(
|
||||
key="nos-headlines",
|
||||
namespace="news",
|
||||
data={"headlines": ["Test headline"]},
|
||||
source="nos.nl",
|
||||
ttl=3600,
|
||||
refresh_schedule="0 * * * *",
|
||||
user="jpmschweitzer",
|
||||
)
|
||||
assert record.refresh_schedule == "0 * * * *"
|
||||
|
||||
def test_volatile_record_create(self):
|
||||
"""Test VolatileRecordCreate model."""
|
||||
create = VolatileRecordCreate(
|
||||
data={"price": 150.50, "change": 2.3},
|
||||
source="alpha_vantage",
|
||||
ttl=300,
|
||||
)
|
||||
assert create.data["price"] == 150.50
|
||||
assert create.ttl == 300
|
||||
|
||||
def test_volatile_record_response(self):
|
||||
"""Test VolatileRecordResponse model."""
|
||||
response = VolatileRecordResponse(
|
||||
key="rotterdam",
|
||||
namespace="weather",
|
||||
data={"temperature": 18},
|
||||
source="openweathermap",
|
||||
created_at=datetime.utcnow(),
|
||||
updated_at=datetime.utcnow(),
|
||||
ttl=1800,
|
||||
ttl_remaining=1500,
|
||||
user="jpmschweitzer",
|
||||
)
|
||||
assert response.ttl_remaining == 1500
|
||||
assert response.ttl == 1800
|
||||
|
||||
|
||||
class TestVolatileNamespaces:
|
||||
"""Test volatile namespaces and defaults."""
|
||||
|
||||
def test_all_namespaces_have_default_ttl(self):
|
||||
"""Verify all namespaces have default TTLs defined."""
|
||||
for ns in VolatileNamespace:
|
||||
assert ns in NAMESPACE_DEFAULT_TTL, f"Missing TTL for {ns}"
|
||||
assert NAMESPACE_DEFAULT_TTL[ns] > 0
|
||||
|
||||
def test_weather_default_ttl(self):
|
||||
"""Test weather namespace default TTL."""
|
||||
assert NAMESPACE_DEFAULT_TTL[VolatileNamespace.WEATHER] == 1800 # 30 min
|
||||
|
||||
def test_financial_default_ttl(self):
|
||||
"""Test financial namespace default TTL."""
|
||||
assert NAMESPACE_DEFAULT_TTL[VolatileNamespace.FINANCIAL] == 300 # 5 min
|
||||
|
||||
def test_sports_default_ttl(self):
|
||||
"""Test sports namespace default TTL (fast updates)."""
|
||||
assert NAMESPACE_DEFAULT_TTL[VolatileNamespace.SPORTS] == 60 # 1 min
|
||||
|
||||
def test_namespace_count(self):
|
||||
"""Test we have the expected number of namespaces."""
|
||||
assert len(VolatileNamespace) == 11
|
||||
|
||||
|
||||
class TestVolatileListResponse:
|
||||
"""Test list response models."""
|
||||
|
||||
def test_list_response(self):
|
||||
"""Test VolatileListResponse model."""
|
||||
response = VolatileListResponse(
|
||||
namespace="weather",
|
||||
keys=["rotterdam", "amsterdam", "utrecht"],
|
||||
count=3,
|
||||
user="jpmschweitzer",
|
||||
)
|
||||
assert response.count == 3
|
||||
assert "rotterdam" in response.keys
|
||||
|
||||
|
||||
class TestVolatileScheduledResponse:
|
||||
"""Test scheduled records response."""
|
||||
|
||||
def test_scheduled_response_empty(self):
|
||||
"""Test empty scheduled response."""
|
||||
response = VolatileScheduledResponse(
|
||||
records=[],
|
||||
count=0,
|
||||
user="jpmschweitzer",
|
||||
)
|
||||
assert response.count == 0
|
||||
assert response.records == []
|
||||
|
||||
def test_scheduled_response_with_records(self):
|
||||
"""Test scheduled response with records."""
|
||||
record = VolatileRecordResponse(
|
||||
key="nos-headlines",
|
||||
namespace="news",
|
||||
data={"headlines": []},
|
||||
source="nos.nl",
|
||||
created_at=datetime.utcnow(),
|
||||
updated_at=datetime.utcnow(),
|
||||
ttl=3600,
|
||||
ttl_remaining=3000,
|
||||
refresh_schedule="0 */6 * * *",
|
||||
user="jpmschweitzer",
|
||||
)
|
||||
response = VolatileScheduledResponse(
|
||||
records=[record],
|
||||
count=1,
|
||||
user="jpmschweitzer",
|
||||
)
|
||||
assert response.count == 1
|
||||
assert response.records[0].refresh_schedule == "0 */6 * * *"
|
||||
|
||||
|
||||
class TestVolatileStatsResponse:
|
||||
"""Test stats response model."""
|
||||
|
||||
def test_stats_response(self):
|
||||
"""Test VolatileStatsResponse model."""
|
||||
response = VolatileStatsResponse(
|
||||
total_records=15,
|
||||
by_namespace={"weather": 3, "news": 5, "financial": 7},
|
||||
scheduled_count=2,
|
||||
total_memory_bytes=None,
|
||||
user="jpmschweitzer",
|
||||
)
|
||||
assert response.total_records == 15
|
||||
assert response.by_namespace["weather"] == 3
|
||||
assert response.scheduled_count == 2
|
||||
|
||||
|
||||
class TestVolatileDeleteResponses:
|
||||
"""Test delete response models."""
|
||||
|
||||
def test_delete_response(self):
|
||||
"""Test VolatileDeleteResponse model."""
|
||||
response = VolatileDeleteResponse(
|
||||
key="rotterdam",
|
||||
namespace="weather",
|
||||
deleted=True,
|
||||
user="jpmschweitzer",
|
||||
)
|
||||
assert response.deleted is True
|
||||
|
||||
def test_delete_not_found(self):
|
||||
"""Test delete response when record not found."""
|
||||
response = VolatileDeleteResponse(
|
||||
key="nonexistent",
|
||||
namespace="weather",
|
||||
deleted=False,
|
||||
user="jpmschweitzer",
|
||||
)
|
||||
assert response.deleted is False
|
||||
|
||||
def test_bulk_delete_response(self):
|
||||
"""Test VolatileBulkDeleteResponse model."""
|
||||
response = VolatileBulkDeleteResponse(
|
||||
namespace="weather",
|
||||
deleted_count=5,
|
||||
user="jpmschweitzer",
|
||||
)
|
||||
assert response.deleted_count == 5
|
||||
assert response.namespace == "weather"
|
||||
|
||||
|
||||
class TestVolatileService:
|
||||
"""Test VolatileCacheService functionality (Qdrant backend)."""
|
||||
|
||||
@pytest.fixture
|
||||
def mock_qdrant(self):
|
||||
"""Create mock Qdrant client."""
|
||||
qdrant = AsyncMock()
|
||||
qdrant.ensure_collection = AsyncMock()
|
||||
qdrant.collection_exists = AsyncMock(return_value=True)
|
||||
qdrant.upsert_vector = AsyncMock(return_value=True)
|
||||
qdrant.delete_by_ids = AsyncMock(return_value=1)
|
||||
qdrant.search_with_expiry_filter = AsyncMock(return_value=[])
|
||||
qdrant.scroll_all_points = AsyncMock(return_value=[])
|
||||
qdrant.delete_expired_vectors = AsyncMock(return_value=0)
|
||||
qdrant.get_volatile_collections = AsyncMock(return_value=[])
|
||||
return qdrant
|
||||
|
||||
@pytest.fixture
|
||||
def mock_ollama(self):
|
||||
"""Create mock Ollama client."""
|
||||
ollama = AsyncMock()
|
||||
ollama.embed = AsyncMock(return_value=[0.1] * 768) # Return 768-dim embedding
|
||||
return ollama
|
||||
|
||||
@pytest.fixture
|
||||
def mock_settings(self):
|
||||
"""Create mock settings."""
|
||||
settings = MagicMock()
|
||||
settings.volatile_default_ttl = 3600
|
||||
return settings
|
||||
|
||||
@pytest.fixture
|
||||
def volatile_service(self, mock_qdrant, mock_ollama, mock_settings):
|
||||
"""Create VolatileCacheService with mocks."""
|
||||
from src.services.volatile_service import VolatileCacheService
|
||||
return VolatileCacheService(
|
||||
qdrant_client=mock_qdrant,
|
||||
ollama_client=mock_ollama,
|
||||
settings=mock_settings
|
||||
)
|
||||
|
||||
def test_collection_name(self, volatile_service):
|
||||
"""Test collection naming pattern."""
|
||||
name = volatile_service._collection_name("jpmschweitzer")
|
||||
assert name == "volatile_jpmschweitzer"
|
||||
|
||||
def test_make_vector_id(self, volatile_service):
|
||||
"""Test deterministic vector ID generation."""
|
||||
id1 = volatile_service._make_vector_id("weather", "rotterdam")
|
||||
id2 = volatile_service._make_vector_id("weather", "rotterdam")
|
||||
id3 = volatile_service._make_vector_id("weather", "amsterdam")
|
||||
|
||||
assert id1 == id2 # Same namespace+key = same ID
|
||||
assert id1 != id3 # Different key = different ID
|
||||
assert len(id1) == 32 # MD5 hex length
|
||||
|
||||
def test_get_default_ttl_known_namespace(self, volatile_service):
|
||||
"""Test default TTL for known namespace."""
|
||||
ttl = volatile_service._get_default_ttl("weather")
|
||||
assert ttl == 1800 # Weather namespace default
|
||||
|
||||
def test_get_default_ttl_unknown_namespace(self, volatile_service):
|
||||
"""Test default TTL for unknown namespace."""
|
||||
ttl = volatile_service._get_default_ttl("unknown_namespace")
|
||||
assert ttl == 3600 # Falls back to settings default
|
||||
|
||||
def test_to_natural_language_weather(self, volatile_service):
|
||||
"""Test natural language conversion for weather data."""
|
||||
text = volatile_service._to_natural_language(
|
||||
namespace="weather",
|
||||
key="rotterdam",
|
||||
data={"temperature": 18, "conditions": "Cloudy", "humidity": 75}
|
||||
)
|
||||
assert "rotterdam" in text.lower()
|
||||
assert "18" in text
|
||||
assert "Cloudy" in text
|
||||
assert "75" in text
|
||||
|
||||
def test_to_natural_language_news(self, volatile_service):
|
||||
"""Test natural language conversion for news data."""
|
||||
text = volatile_service._to_natural_language(
|
||||
namespace="news",
|
||||
key="nos-headlines",
|
||||
data={"title": "Breaking News", "summary": "Something happened", "source": "NOS"}
|
||||
)
|
||||
assert "Breaking News" in text
|
||||
assert "Something happened" in text
|
||||
assert "NOS" in text
|
||||
|
||||
def test_to_natural_language_financial(self, volatile_service):
|
||||
"""Test natural language conversion for financial data."""
|
||||
text = volatile_service._to_natural_language(
|
||||
namespace="financial",
|
||||
key="AAPL",
|
||||
data={"symbol": "AAPL", "price": 150.50, "change": 2.3}
|
||||
)
|
||||
assert "AAPL" in text
|
||||
assert "price" in text.lower()
|
||||
assert "change" in text.lower()
|
||||
|
||||
def test_to_natural_language_transit(self, volatile_service):
|
||||
"""Test natural language conversion for transit data."""
|
||||
text = volatile_service._to_natural_language(
|
||||
namespace="transit",
|
||||
key="ns-intercity",
|
||||
data={"route": "Amsterdam-Rotterdam", "status": "On time", "delay": 0}
|
||||
)
|
||||
assert "Amsterdam-Rotterdam" in text or "ns-intercity" in text.lower()
|
||||
assert "On time" in text
|
||||
|
||||
def test_to_natural_language_fallback(self, volatile_service):
|
||||
"""Test natural language fallback for unknown namespace."""
|
||||
text = volatile_service._to_natural_language(
|
||||
namespace="custom",
|
||||
key="test-key",
|
||||
data={"foo": "bar", "count": 42}
|
||||
)
|
||||
assert "custom" in text.lower()
|
||||
assert "foo" in text or "bar" in text
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_store_success(self, volatile_service, mock_qdrant, mock_ollama):
|
||||
"""Test successful store operation."""
|
||||
result = await volatile_service.store(
|
||||
user="jpmschweitzer",
|
||||
namespace="weather",
|
||||
key="rotterdam",
|
||||
data={"temperature": 18, "conditions": "Sunny"},
|
||||
source="openweathermap",
|
||||
ttl=1800
|
||||
)
|
||||
|
||||
assert result.key == "rotterdam"
|
||||
assert result.namespace == "weather"
|
||||
assert result.ttl == 1800
|
||||
mock_qdrant.ensure_collection.assert_called_once()
|
||||
mock_ollama.embed.assert_called_once()
|
||||
mock_qdrant.upsert_vector.assert_called_once()
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_store_uses_namespace_default_ttl(self, volatile_service, mock_qdrant, mock_ollama):
|
||||
"""Test store uses namespace default TTL when not specified."""
|
||||
result = await volatile_service.store(
|
||||
user="jpmschweitzer",
|
||||
namespace="weather",
|
||||
key="amsterdam",
|
||||
data={"temperature": 16},
|
||||
source="openweathermap",
|
||||
ttl=None # Not specified
|
||||
)
|
||||
|
||||
assert result.ttl == 1800 # Weather default
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_search_empty_collection(self, volatile_service, mock_qdrant, mock_ollama):
|
||||
"""Test search when collection doesn't exist."""
|
||||
mock_qdrant.collection_exists.return_value = False
|
||||
|
||||
results = await volatile_service.search(
|
||||
user="jpmschweitzer",
|
||||
query="weather rotterdam"
|
||||
)
|
||||
|
||||
assert results == []
|
||||
mock_ollama.embed.assert_not_called()
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_search_with_results(self, volatile_service, mock_qdrant, mock_ollama):
|
||||
"""Test search returns results."""
|
||||
import time
|
||||
now_ms = int(time.time() * 1000)
|
||||
|
||||
mock_qdrant.search_with_expiry_filter.return_value = [
|
||||
{
|
||||
"score": 0.95,
|
||||
"payload": {
|
||||
"key": "rotterdam",
|
||||
"namespace": "weather",
|
||||
"raw_data": {"temperature": 18},
|
||||
"source": "openweathermap",
|
||||
"created_at": datetime.utcnow().isoformat(),
|
||||
"updated_at": datetime.utcnow().isoformat(),
|
||||
"ttl": 1800,
|
||||
"ttl_expiry": now_ms + 900000, # 15 min remaining
|
||||
"refresh_schedule": None,
|
||||
"user": "jpmschweitzer"
|
||||
}
|
||||
}
|
||||
]
|
||||
|
||||
results = await volatile_service.search(
|
||||
user="jpmschweitzer",
|
||||
query="weather rotterdam"
|
||||
)
|
||||
|
||||
assert len(results) == 1
|
||||
assert results[0].key == "rotterdam"
|
||||
assert results[0].namespace == "weather"
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_delete_success(self, volatile_service, mock_qdrant):
|
||||
"""Test successful delete."""
|
||||
mock_qdrant.delete_by_ids.return_value = 1
|
||||
|
||||
result = await volatile_service.delete("jpmschweitzer", "weather", "rotterdam")
|
||||
|
||||
assert result is True
|
||||
mock_qdrant.delete_by_ids.assert_called_once()
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_delete_not_found(self, volatile_service, mock_qdrant):
|
||||
"""Test delete when record not found."""
|
||||
mock_qdrant.delete_by_ids.return_value = 0
|
||||
|
||||
result = await volatile_service.delete("jpmschweitzer", "weather", "nonexistent")
|
||||
|
||||
assert result is False
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_get_stats_empty(self, volatile_service, mock_qdrant):
|
||||
"""Test stats with no records."""
|
||||
mock_qdrant.collection_exists.return_value = False
|
||||
|
||||
stats = await volatile_service.get_stats("jpmschweitzer")
|
||||
|
||||
assert stats["total_records"] == 0
|
||||
assert stats["by_namespace"] == {}
|
||||
assert stats["scheduled_count"] == 0
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_get_stats_with_records(self, volatile_service, mock_qdrant):
|
||||
"""Test stats with records."""
|
||||
import time
|
||||
now_ms = int(time.time() * 1000)
|
||||
|
||||
mock_qdrant.scroll_all_points.return_value = [
|
||||
{"payload": {"namespace": "weather", "ttl_expiry": now_ms + 100000}},
|
||||
{"payload": {"namespace": "weather", "ttl_expiry": now_ms + 100000, "refresh_schedule": "0 * * * *"}},
|
||||
{"payload": {"namespace": "news", "ttl_expiry": now_ms + 100000}},
|
||||
{"payload": {"namespace": "weather", "ttl_expiry": now_ms - 100000}}, # Expired
|
||||
]
|
||||
|
||||
stats = await volatile_service.get_stats("jpmschweitzer")
|
||||
|
||||
assert stats["total_records"] == 3 # Excludes expired
|
||||
assert stats["by_namespace"]["weather"] == 2
|
||||
assert stats["by_namespace"]["news"] == 1
|
||||
assert stats["scheduled_count"] == 1
|
||||
assert stats["expired_count"] == 1
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_purge_expired(self, volatile_service, mock_qdrant):
|
||||
"""Test purging expired records."""
|
||||
mock_qdrant.delete_expired_vectors.return_value = 5
|
||||
|
||||
result = await volatile_service.purge_expired("jpmschweitzer")
|
||||
|
||||
assert result == 5
|
||||
mock_qdrant.delete_expired_vectors.assert_called_once()
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_purge_all_expired(self, volatile_service, mock_qdrant):
|
||||
"""Test purging expired from all collections."""
|
||||
mock_qdrant.get_volatile_collections.return_value = [
|
||||
"volatile_user1",
|
||||
"volatile_user2"
|
||||
]
|
||||
mock_qdrant.delete_expired_vectors.side_effect = [3, 2]
|
||||
|
||||
results = await volatile_service.purge_all_expired()
|
||||
|
||||
assert results["volatile_user1"] == 3
|
||||
assert results["volatile_user2"] == 2
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_get_scheduled(self, volatile_service, mock_qdrant):
|
||||
"""Test getting scheduled records."""
|
||||
import time
|
||||
now_ms = int(time.time() * 1000)
|
||||
|
||||
mock_qdrant.scroll_all_points.return_value = [
|
||||
{
|
||||
"payload": {
|
||||
"key": "nos-headlines",
|
||||
"namespace": "news",
|
||||
"raw_data": {"headlines": []},
|
||||
"source": "nos.nl",
|
||||
"created_at": datetime.utcnow().isoformat(),
|
||||
"updated_at": datetime.utcnow().isoformat(),
|
||||
"ttl": 3600,
|
||||
"ttl_expiry": now_ms + 1800000,
|
||||
"refresh_schedule": "0 */6 * * *",
|
||||
"user": "jpmschweitzer"
|
||||
}
|
||||
},
|
||||
{
|
||||
"payload": {
|
||||
"key": "rotterdam",
|
||||
"namespace": "weather",
|
||||
"raw_data": {"temperature": 18},
|
||||
"source": "openweathermap",
|
||||
"created_at": datetime.utcnow().isoformat(),
|
||||
"updated_at": datetime.utcnow().isoformat(),
|
||||
"ttl": 1800,
|
||||
"ttl_expiry": now_ms + 900000,
|
||||
"refresh_schedule": None, # Not scheduled
|
||||
"user": "jpmschweitzer"
|
||||
}
|
||||
}
|
||||
]
|
||||
|
||||
scheduled = await volatile_service.get_scheduled("jpmschweitzer")
|
||||
|
||||
assert len(scheduled) == 1
|
||||
assert scheduled[0].key == "nos-headlines"
|
||||
assert scheduled[0].refresh_schedule == "0 */6 * * *"
|
||||
|
||||
|
||||
class TestVolatileCleanupEndpoint:
|
||||
"""Test volatile cleanup in maintenance router."""
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_cleanup_volatile(self):
|
||||
"""Test volatile cleanup endpoint."""
|
||||
from src.routers.maintenance import cleanup_volatile, VolatileCleanupResponse
|
||||
|
||||
mock_qdrant = AsyncMock()
|
||||
mock_qdrant.get_volatile_collections = AsyncMock(return_value=[
|
||||
"volatile_user1",
|
||||
"volatile_user2"
|
||||
])
|
||||
mock_qdrant.delete_expired_vectors = AsyncMock(side_effect=[3, 2])
|
||||
|
||||
mock_ollama = AsyncMock()
|
||||
|
||||
mock_settings = MagicMock()
|
||||
mock_settings.volatile_default_ttl = 3600
|
||||
|
||||
with patch('src.routers.maintenance.get_settings', return_value=mock_settings):
|
||||
result = await cleanup_volatile(
|
||||
qdrant=mock_qdrant,
|
||||
ollama=mock_ollama,
|
||||
api_key="test"
|
||||
)
|
||||
|
||||
assert result.success is True
|
||||
assert result.collections_processed == 2
|
||||
assert result.total_expired_purged == 5
|
||||
assert result.by_collection["volatile_user1"] == 3
|
||||
assert result.by_collection["volatile_user2"] == 2
|
||||
Reference in New Issue
Block a user