Compare commits

...
5 Commits
Author SHA1 Message Date
jpmschweitzerandClaude Opus 4.5 e6e65d6d78 feat: add test data cleanup endpoint
Build and Push / build (release) Successful in 28s
Add POST /maintenance/cleanup/test-data endpoint to purge LLM test data
from wiki, graph, and vectors. Security-restricted to test user namespace
only (users/llm-tester/*, users/llm_tester/*).

- Supports dry_run=true (default) to preview before deleting
- Cleans vectors, graph nodes, and wiki pages
- Scheduler task configured for weekly cleanup (Sunday 3:00 AM)

🤖 Generated with [Claude Code](https://claude.com/claude-code)

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
2025-12-24 21:12:16 +01:00
jpmschweitzerandClaude Opus 4.5 37f8e1819e feat: refactor volatile cache to vector storage with HybridRAG integration
Build and Push / build (release) Successful in 28s
- Migrate volatile backend from Redis to Qdrant for semantic search
- Add natural language conversion for structured data embedding
- Simplify API: /volatile/search, /volatile/store, /{namespace}/{key}
- Integrate volatile into HybridRAG with priority boost in RRF fusion
- Add POST /maintenance/cleanup/volatile for expiry purging
- Update tests for new Qdrant-based architecture (37/37 pass)

🤖 Generated with [Claude Code](https://claude.com/claude-code)

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
2025-12-24 20:03:37 +01:00
jpmschweitzerandClaude Opus 4.5 1f848c4878 release: v1.4.2 - volatile cache system
Build and Push / build (release) Successful in 28s
Phase 2 of Memory Management System complete.
See CHANGELOG.md for details.

🤖 Generated with [Claude Code](https://claude.com/claude-code)

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
2025-12-24 17:32:17 +01:00
jpmschweitzerandClaude Opus 4.5 7297e6b9f1 feat: add volatile cache system for ephemeral data
Phase 2 of Memory Management System - volatile memory tier:

- VolatileCacheService: Redis-backed TTL storage
- Volatile router with full CRUD operations
- Predefined namespaces: weather, news, financial, transit, traffic,
  air_quality, sports, social, system, context, custom
- Each namespace has appropriate default TTL (1min to 1hr)
- Refresh schedule support via cron expressions
- Scheduler integration endpoint: GET /volatile/scheduled

Endpoints:
- GET/POST/DELETE /volatile/{namespace}/{key}
- GET/DELETE /volatile/{namespace}
- GET /volatile/stats
- GET /volatile/scheduled
- GET /volatile/namespaces

🤖 Generated with [Claude Code](https://claude.com/claude-code)

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
2025-12-24 17:31:45 +01:00
jpmschweitzerandClaude Opus 4.5 2552bfd1f9 fix: make Wiki.js API token optional for open GraphQL endpoints
Build and Push / build (release) Successful in 33s
The Wiki.js GraphQL API is accessible without authentication.
Make WIKI_GRAPHQL_API env var optional with empty default to fix
container startup failures.

🤖 Generated with [Claude Code](https://claude.com/claude-code)

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
2025-12-24 16:59:41 +01:00
15 changed files with 2161 additions and 123 deletions
+74
View File
@@ -5,6 +5,80 @@ 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.4] - 2025-12-24
### Added
- **Test Data Cleanup Endpoint** - `POST /maintenance/cleanup/test-data`
- Purges LLM test data from wiki, graph, and vectors
- Security-restricted to test user namespace only (`users/llm-tester/*`, `users/llm_tester/*`)
- Supports `dry_run=true` (default) to preview before deleting
- Scheduler task configured for weekly cleanup (Sunday 3:00 AM)
## [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
+138 -89
View File
@@ -6,15 +6,24 @@ A three-tier memory architecture for Library Desk with intelligent orchestration
| Tier | Storage | Purpose | TTL |
|------|---------|---------|-----|
| **Volatile** | Redis | Weather, news, financial, ephemeral context | 5min - 2hr |
| **Volatile** | Qdrant (vectors) | Weather, news, financial, ephemeral context | 5min - 2hr |
| **Documents** | TBD (research) | Git mirrors, PDFs, video, images | Permanent |
| **Knowledge** | Wiki + Neo4j | Personal dossiers, research, summaries | Permanent |
**Implementation Priority**: Cleanup → Volatile → Documents
**Implementation Priority**: Cleanup → Volatile → Documents → Test Data Cleanup
### Phase Status
| Phase | Status | Version |
|-------|--------|---------|
| Phase 1: Cleanup System | ✅ Complete | v1.4.0 |
| Phase 2: Volatile Memory | ✅ Complete | v1.4.3 |
| Phase 3: Document Storage | ⏳ Pending | - |
| Phase 4: Test Data Cleanup | ✅ Complete | v1.4.4 |
---
## Phase 1: Cleanup System Completion
## Phase 1: Cleanup System Completion
### Current State
- **COMPLETE** - All Phase 1 tasks implemented
@@ -57,9 +66,9 @@ A three-tier memory architecture for Library Desk with intelligent orchestration
---
## Phase 2: Volatile Memory System
## Phase 2: Volatile Memory System
### Architecture
### Architecture (Final Implementation)
```
┌─────────────────┐ ┌──────────────┐ ┌─────────────────┐
@@ -71,88 +80,46 @@ A three-tier memory architecture for Library Desk with intelligent orchestration
┌─────────────────┐
Redis
(DB 4, TTL)
Qdrant
(volatile_{user})
└─────────────────┘
```
### Data Model
**Key design decisions:**
- Vector storage in Qdrant (not Redis) for semantic search
- Collection per user: `volatile_{user}`
- TTL via `ttl_expiry` timestamp in payload
- Natural language conversion for embedding structured data
- Integrated into HybridRAG with priority boost
```python
class VolatileRecord(BaseModel):
key: str # e.g., "weather:rotterdam"
namespace: str # e.g., "weather", "news", "financial"
data: dict # Actual content
source: Optional[str] # Origin API/service
created_at: datetime
updated_at: datetime
ttl: int # Seconds until expiration
refresh_schedule: Optional[str] # Cron expression, if repeating
user: str # Multi-tenant isolation
```
**Key pattern**: `{user}:volatile:{namespace}:{key_hash}`
### Implementation Order: Integration-First
1. **Start with Consolidation Hook** - Understand data flow through existing system
2. **Build Service Layer** - VolatileCacheService with Redis operations
3. **Add API Endpoints** - REST interface for volatile data
4. **Biographer Integration** - Query user preferences for relevance
### Tasks
#### 2.1 Integrate with Consolidation (FIRST)
**New file**: `src/services/volatile_service.py`
```python
class VolatileCacheService:
async def get(user, namespace, key) -> Optional[VolatileRecord]
async def set(user, namespace, key, data, ttl, refresh_schedule=None)
async def delete(user, namespace, key)
async def list_namespace(user, namespace) -> List[str]
async def get_scheduled(user) -> List[VolatileRecord] # For scheduler
```
#### 2.2 Create Volatile API Router
**New file**: `src/routers/volatile.py`
### Endpoints (Implemented)
| Endpoint | Method | Purpose |
|----------|--------|---------|
| `/volatile/{namespace}/{key}` | GET | Retrieve record |
| `/volatile/{namespace}/{key}` | POST | Store/update record |
| `/volatile/search?q=...` | GET | Semantic search across volatile data |
| `/volatile/store?namespace=...&key=...` | POST | Store/update record |
| `/volatile/{namespace}/{key}` | GET | Retrieve specific record |
| `/volatile/{namespace}/{key}` | DELETE | Remove record |
| `/volatile/{namespace}` | GET | List keys in namespace |
| `/volatile/scheduled` | GET | List records needing refresh |
| `/volatile/stats` | GET | Cache statistics |
| `/volatile/scheduled` | GET | Records needing refresh |
| `/volatile/namespaces` | GET | List available namespaces |
| `/maintenance/cleanup/volatile` | POST | Purge expired records |
#### 2.3 Integrate with Consolidation
**File**: `src/services/consolidation_service.py`
### Namespaces
Add relevance trigger detection:
1. During consolidation, analyze search results for location/interest patterns
2. Query tatlock's Biographer collection for user preferences
3. If match found, create/update volatile refresh schedule
#### 2.4 Biographer Integration
**File**: `src/core/dependencies.py`
```python
def get_biographer_qdrant() -> QdrantClientWrapper:
"""Direct access to tatlock's Biographer collection."""
# Configure to connect to tatlock's Qdrant
```
#### 2.5 Scheduler-Side Configuration
Document required scheduler tasks:
```json
{
"task_name": "volatile_refresh",
"schedule": "*/15 * * * *",
"endpoint": "GET /volatile/scheduled",
"follow_up": "For each record, call refresh endpoint with record.refresh_schedule"
}
```
| Namespace | Default TTL | Use Case |
|-----------|-------------|----------|
| weather | 30 min | Current conditions, forecasts |
| news | 1 hour | Headlines, breaking news |
| financial | 5 min | Stock prices, exchange rates |
| transit | 5 min | Train/bus schedules, delays |
| traffic | 10 min | Commute times, road conditions |
| air_quality | 1 hour | Pollution, pollen counts |
| sports | 1 min | Live scores, matches |
| social | 10 min | Social notifications |
| system | 1 min | Service health status |
| context | 1 hour | Session state |
| custom | 1 hour | User-defined data |
---
@@ -219,34 +186,116 @@ Add LLM-powered category descriptor generation:
---
## Files to Modify/Create
## Phase 4: LLM Tester Data Cleanup ✅
### Phase 1 (Cleanup)
- `src/routers/maintenance.py` - Add timestamp tracking
### Problem
LLM testing creates accumulated cruft across the system:
- Wiki.js pages under `llm-tester/` and `llm_tester/` paths
- Graph nodes (Document, Entity) linked to test pages
- Vector chunks in Qdrant for test content
This data accumulates over time and clutters Wiki.js visually (no separate tenant scope for tests).
### Solution
Add a maintenance endpoint to purge all LLM tester artifacts across wiki, graph, and vectors.
### Tasks
#### 4.1 Identify Test Data Patterns ✅
**Patterns matched** (security-restricted to test user namespace):
- `users/llm-tester/*`
- `users/llm_tester/*`
#### 4.2 Add Cleanup Endpoint ✅
**File**: `src/routers/maintenance.py`
```python
@router.post("/cleanup/test-data")
async def cleanup_test_data(
dry_run: bool = Query(default=True),
wiki: WikiJSDep = None,
vector_service: VectorServiceDep = None,
graph_service: GraphServiceDep = None,
api_key: str = Depends(verify_api_key)
):
"""
Purge LLM tester data from wiki, graph, and vectors.
**Security**: Only deletes pages in the test user namespace:
- users/llm-tester/*
- users/llm_tester/*
Use dry_run=true to preview what would be deleted.
"""
```
#### 4.3 Implementation Steps ✅
1. **Wiki cleanup**: Delete pages via GraphQL mutation
2. **Graph cleanup**: Delete Document nodes using `delete_page()` method
3. **Vector cleanup**: Delete chunks using `delete_page_chunks()` method
#### 4.4 Scheduler Integration ✅
**Recommended schedule**: Weekly (Sunday 3:00 AM)
```json
{
"task_name": "test_data_cleanup",
"schedule": "0 3 * * 0",
"endpoint": "POST /maintenance/cleanup/test-data?dry_run=false",
"description": "Weekly cleanup of LLM test data"
}
```
### Files to Modify
- `src/routers/maintenance.py` - Add cleanup endpoint
- `src/services/wiki_service.py` - Add bulk delete by path pattern (if needed)
- `src/services/graph_service.py` - May need pattern-based node deletion
- `src/services/vector_service.py` - Add pattern-based chunk deletion
---
## Files Modified/Created
### Phase 1 (Cleanup) ✅
- `src/routers/maintenance.py` - Timestamp tracking, cleanup endpoints
- `src/services/graph_service.py` - Bidirectional validation
- `src/services/vector_service.py` - Cross-reference checks
- `LIBRARIAN_INTEGRATION.md` - Scheduler config docs
### Phase 2 (Volatile)
- `src/services/volatile_service.py` - **NEW**
- `src/routers/volatile.py` - **NEW**
- `src/models/volatile.py` - **NEW**
- `src/core/dependencies.py` - Add Biographer client
- `src/services/consolidation_service.py` - Relevance triggers
- `tests/test_volatile.py` - **NEW**
### Phase 2 (Volatile)
- `src/services/volatile_service.py` - Qdrant-based volatile cache
- `src/routers/volatile.py` - Simplified endpoints
- `src/models/volatile.py` - Namespaces and models
- `src/models/hybrid_rag.py` - Volatile config options
- `src/services/hybrid_rag_service.py` - Volatile integration
- `src/clients/qdrant_client.py` - Expiry filter methods
- `tests/test_volatile.py` - 37 tests
### Phase 3 (Documents)
- `docs/DOCUMENT_STORAGE_RESEARCH.md` - **NEW**
- `src/services/document_store_service.py` - **NEW** (post-research)
- `src/routers/documents.py` - **NEW** (post-research)
### Phase 4 (Test Data Cleanup)
- `src/routers/maintenance.py` - Add cleanup endpoint
- `src/services/wiki_service.py` - Bulk delete by path pattern
- `src/services/graph_service.py` - Pattern-based node deletion
- `src/services/vector_service.py` - Pattern-based chunk deletion
---
## Resolved Design Decisions
1. **Biographer Qdrant**: Same Qdrant instance, different collection. Library-Desk queries directly.
2. **Scheduler API**: Has REST API for task registration. Library-Desk can programmatically create refresh schedules.
3. **External API calls**: Library-Desk routes through SearXNG for web search. Consider dedicated API integrations for high-value volatiles (weather, financial) for consistent quality.
1. **Volatile Storage**: Qdrant vectors (not Redis) for semantic search capability
2. **Collection Naming**: `volatile_{user}` for per-user isolation
3. **TTL Mechanism**: `ttl_expiry` timestamp in payload, background cleanup job
4. **HybridRAG Integration**: Volatile as third source with RRF priority boost
5. **Biographer Qdrant**: Same Qdrant instance, different collection
6. **Scheduler API**: Has REST API for task registration
---
+1 -1
View File
@@ -1,6 +1,6 @@
[project]
name = "library-desk"
version = "1.4.0"
version = "1.4.4"
description = "Coordination service for The Library system - HybridRAG queries, document ingestion, entity extraction, and knowledge consolidation"
readme = "README.md"
requires-python = ">=3.12"
+131 -1
View File
@@ -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 []
+9 -12
View File
@@ -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
View File
@@ -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
View File
@@ -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"
+5 -1
View File
@@ -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")
+130
View File
@@ -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")
+13 -8
View File
@@ -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
+186 -1
View File
@@ -16,10 +16,13 @@ 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 src.core.multi_tenancy import DEFAULT_USER
from datetime import datetime, timezone
logger = logging.getLogger(__name__)
@@ -209,6 +212,44 @@ 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
class TestDataCleanupResponse(BaseModel):
"""Response from test data cleanup operation."""
success: bool
dry_run: bool
wiki_pages_deleted: int
graph_nodes_deleted: int
vector_chunks_deleted: int
pages_found: List[Dict[str, Any]] = Field(default_factory=list)
duration_ms: float
# Test data path patterns - restricted to test user namespace only
# These are the only paths that can be cleaned up for safety
TEST_USER_PATH_PREFIXES = [
"users/llm-tester/",
"users/llm_tester/",
]
def _matches_test_user_path(path: str) -> bool:
"""Check if a path is in the test user namespace.
Only matches paths that START with test user prefixes for safety.
This prevents accidental deletion of non-test data.
"""
path_lower = path.lower()
return any(path_lower.startswith(prefix) for prefix in TEST_USER_PATH_PREFIXES)
# ========== Endpoints ==========
@router.post("/cleanup/vectors", response_model=VectorCleanupResponse)
@@ -490,6 +531,150 @@ 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.post("/cleanup/test-data", response_model=TestDataCleanupResponse)
async def cleanup_test_data(
dry_run: bool = Query(default=True, description="Preview only, don't delete"),
wiki: WikiJSDep = None,
vector_service: VectorServiceDep = None,
graph_service: GraphServiceDep = None,
api_key: str = Depends(verify_api_key)
):
"""
Purge LLM tester data from wiki, graph, and vectors.
**Security**: Only deletes pages in the test user namespace:
- users/llm-tester/*
- users/llm_tester/*
This endpoint cannot delete data outside these paths.
**Use dry_run=true (default) to preview what would be deleted.**
**Scheduler Integration:**
```json
{
"task_name": "test_data_cleanup",
"schedule": "0 3 * * 0",
"endpoint": "POST /maintenance/cleanup/test-data?dry_run=false",
"description": "Weekly cleanup of LLM test data"
}
```
"""
start_time = time.time()
try:
# List all wiki pages
all_pages = await wiki.list_all_pages(batch_size=500)
# Filter for test user paths only (security: restricted to test namespace)
test_pages = [
{"id": p["id"], "path": p["path"], "title": p.get("title", "")}
for p in all_pages
if _matches_test_user_path(p.get("path", ""))
]
logger.info(f"Found {len(test_pages)} test pages matching patterns: {TEST_USER_PATH_PREFIXES}")
wiki_deleted = 0
graph_deleted = 0
vector_deleted = 0
if not dry_run and test_pages:
for page in test_pages:
page_id = page["id"]
page_path = page["path"]
try:
# Delete vector chunks for this page (using DEFAULT_USER collection)
chunks_removed = await vector_service.delete_page_chunks(page_id, DEFAULT_USER)
vector_deleted += chunks_removed
# Delete graph node for this page (returns count, may be 0 if no node)
graph_removed = await graph_service.delete_page(page_id, DEFAULT_USER)
graph_deleted += graph_removed
# Delete wiki page (raises exception on failure, returns None on success)
await wiki.delete_page(page_id)
wiki_deleted += 1
logger.info(f"Deleted test page: {page_path} (id={page_id})")
except Exception as e:
logger.error(f"Failed to delete page {page_path}: {e}")
continue
duration_ms = (time.time() - start_time) * 1000
return TestDataCleanupResponse(
success=True,
dry_run=dry_run,
wiki_pages_deleted=wiki_deleted,
graph_nodes_deleted=graph_deleted,
vector_chunks_deleted=vector_deleted,
pages_found=test_pages,
duration_ms=duration_ms
)
except Exception as e:
logger.error(f"Test data 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"),
+273
View File
@@ -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,
)
+67 -8
View File
@@ -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
+564
View File
@@ -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
+567
View File
@@ -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