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>
568 lines
20 KiB
Python
568 lines
20 KiB
Python
"""
|
|
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
|