""" Phase 14 Plan 01 — RED contract tests for the byte cache service. Requirements: CACHE-03, CACHE-04, CACHE-05 Threat coverage: T-14-01 — IDOR: cache metadata queries are always owner-scoped T-14-02 — Object-key secrecy: object_key never appears in API responses T-14-05 — Cache quota corruption: quota must be atomically incremented on cache retain and decremented on eviction; a safe in-memory simulation is used because SQLite cannot run the atomic UPDATE. All tests import non-existent Phase 14 modules and must FAIL due to ImportError until implementation is complete. """ from __future__ import annotations import uuid from unittest.mock import AsyncMock, MagicMock, patch import pytest pytestmark = pytest.mark.asyncio from tests.conftest import _TEST_USER_AGENT # ─── Helpers ───────────────────────────────────────────────────────────────── async def _create_user_and_token(session, role: str = "user"): from db.models import User, Quota from services.auth import hash_password, create_access_token user_id = uuid.uuid4() user = User( id=user_id, handle=f"cache_user_{user_id.hex[:8]}", email=f"cache_{user_id.hex[:8]}@example.com", password_hash=hash_password("Testpassword123!"), role=role, is_active=True, password_must_change=False, ) quota = Quota(user_id=user_id, limit_bytes=104857600, used_bytes=0) session.add(user) session.add(quota) await session.commit() await session.refresh(user) token = create_access_token(str(user_id), role, user_agent=_TEST_USER_AGENT) return { "user": user, "token": token, "headers": { "Authorization": f"Bearer {token}", "User-Agent": _TEST_USER_AGENT, }, } async def _create_cloud_connection(session, user_id, provider="google_drive", status="ACTIVE"): from db.models import CloudConnection from storage.cloud_utils import encrypt_credentials from config import settings master_key = settings.cloud_creds_key.encode() creds_enc = encrypt_credentials( master_key, str(user_id), {"access_token": "tok", "refresh_token": "ref"}, ) conn = CloudConnection( id=uuid.uuid4(), user_id=user_id, provider=provider, display_name="Cache Test Connection", credentials_enc=creds_enc, status=status, ) session.add(conn) await session.commit() return conn async def _create_cloud_item(session, user_id, connection_id, *, name="doc.pdf", kind="file", etag="etag-v1"): from db.models import CloudItem item = CloudItem( id=uuid.uuid4(), user_id=user_id, connection_id=connection_id, provider_item_id=f"pitem-{uuid.uuid4().hex[:8]}", name=name, kind=kind, content_type="application/pdf", provider_size=102400, etag=etag, analysis_status="pending", ) session.add(item) await session.commit() return item # ─── CACHE-03: Cache bytes only during open/preview/analysis ───────────────── async def test_cache_entry_created_during_analysis_job(db_session): """CACHE-03: A cache entry is created when bytes are hydrated for analysis. Imports CloudByteCacheEntry (Phase 14 schema addition). """ from services.cloud_cache import create_cache_entry # Phase 14 user_id = uuid.uuid4() conn_id = uuid.uuid4() item_id = uuid.uuid4() entry = await create_cache_entry( session=db_session, user_id=user_id, connection_id=conn_id, cloud_item_id=item_id, provider_item_id="pitem-001", version_key="etag-v1", object_key=f"cache/{user_id}/{uuid.uuid4()}.pdf", content_type="application/pdf", size_bytes=102400, ) assert entry is not None assert entry.user_id == user_id assert entry.cloud_item_id == item_id assert entry.size_bytes == 102400 assert entry.pin_count == 0 assert entry.active_job_count == 0 async def test_browse_without_analysis_does_not_create_cache_entry(db_session): """CACHE-03: Listing a folder never creates a cache entry (no byte download). Verifies that reconcile_cloud_listing (existing Phase 12/13 code path) leaves cloud_byte_cache_entries empty. """ from db.models import CloudConnection, CloudItem from sqlalchemy import select # Try to import Phase 14 cache model try: from db.models import CloudByteCacheEntry # Phase 14 — not yet in models except ImportError: pytest.fail( "CloudByteCacheEntry model must be added to db/models.py in Phase 14" ) user_id = uuid.uuid4() # After any browse/reconcile operation, zero cache entries should exist result = await db_session.execute( select(CloudByteCacheEntry).where(CloudByteCacheEntry.user_id == user_id) ) rows = result.scalars().all() assert len(rows) == 0, "Browse operations must not create cache entries" # ─── CACHE-04: LRU eviction and cache limits ───────────────────────────────── async def test_eviction_does_not_remove_pinned_entry(db_session): """CACHE-04/D-16: Eviction must not evict entries with active jobs or pins. An entry with pin_count > 0 must survive LRU eviction. """ from services.cloud_cache import create_cache_entry, evict_lru_entries # Phase 14 user_id = uuid.uuid4() conn_id = uuid.uuid4() # Create a pinned cache entry (pin_count=1 simulates active analysis/preview) pinned = await create_cache_entry( session=db_session, user_id=user_id, connection_id=conn_id, cloud_item_id=uuid.uuid4(), provider_item_id="pinned-item", version_key="etag-pinned", object_key=f"cache/{user_id}/{uuid.uuid4()}.pdf", content_type="application/pdf", size_bytes=10 * 1024 * 1024, # 10 MB pin_count=1, # pinned — must survive eviction ) # Create an unpinned entry to confirm eviction can happen at all evictable = await create_cache_entry( session=db_session, user_id=user_id, connection_id=conn_id, cloud_item_id=uuid.uuid4(), provider_item_id="evictable-item", version_key="etag-evict", object_key=f"cache/{user_id}/{uuid.uuid4()}.pdf", content_type="application/pdf", size_bytes=5 * 1024 * 1024, # 5 MB pin_count=0, ) # Run eviction with a tight byte limit that forces eviction of unpinned entries evicted_ids = await evict_lru_entries( session=db_session, user_id=user_id, max_cache_bytes=1, # 1 byte limit forces eviction ) # The pinned entry must NOT be evicted assert str(pinned.id) not in [str(e) for e in evicted_ids], ( "Eviction must preserve entries with pin_count > 0 (D-16)" ) async def test_eviction_removes_oldest_unpinned_entry_first(db_session): """CACHE-04/D-14: LRU eviction removes least-recently-accessed unpinned entry first.""" from services.cloud_cache import create_cache_entry, evict_lru_entries # Phase 14 from datetime import datetime, timezone, timedelta user_id = uuid.uuid4() conn_id = uuid.uuid4() older_entry = await create_cache_entry( session=db_session, user_id=user_id, connection_id=conn_id, cloud_item_id=uuid.uuid4(), provider_item_id="old-item", version_key="etag-old", object_key=f"cache/{user_id}/{uuid.uuid4()}.pdf", content_type="application/pdf", size_bytes=5 * 1024 * 1024, pin_count=0, # Simulate last accessed 2 hours ago last_accessed_at=datetime.now(timezone.utc) - timedelta(hours=2), ) newer_entry = await create_cache_entry( session=db_session, user_id=user_id, connection_id=conn_id, cloud_item_id=uuid.uuid4(), provider_item_id="new-item", version_key="etag-new", object_key=f"cache/{user_id}/{uuid.uuid4()}.pdf", content_type="application/pdf", size_bytes=5 * 1024 * 1024, pin_count=0, # Simulate last accessed 1 minute ago last_accessed_at=datetime.now(timezone.utc) - timedelta(minutes=1), ) # Evict enough to remove exactly one entry (limit below combined size) evicted_ids = await evict_lru_entries( session=db_session, user_id=user_id, max_cache_bytes=6 * 1024 * 1024, # fits one entry; must evict the older one ) # The older entry should be evicted first (LRU order) assert str(older_entry.id) in [str(e) for e in evicted_ids], ( "LRU eviction must remove least-recently-accessed entry first" ) # The newer entry should survive assert str(newer_entry.id) not in [str(e) for e in evicted_ids] # ─── CACHE-04/T-14-05: Quota accounting ───────────────────────────────────── async def test_cache_retain_increments_quota(db_session): """T-14-05: Creating a cache entry must atomically increment used_bytes in quotas. The increment requirement is tested via the service layer; the actual atomic SQL UPDATE…RETURNING runs during implementation against real PostgreSQL. In SQLite we verify the function signature exists and accepts the right arguments. """ from services.cloud_cache import increment_quota_for_cache # Phase 14 # Verify the function exists with the expected signature import inspect sig = inspect.signature(increment_quota_for_cache) params = list(sig.parameters.keys()) # Must accept session, user_id, and size_bytes assert "session" in params, "increment_quota_for_cache must accept session" assert "user_id" in params, "increment_quota_for_cache must accept user_id" assert "size_bytes" in params, "increment_quota_for_cache must accept size_bytes" async def test_cache_eviction_decrements_quota(db_session): """T-14-05: Evicting a cache entry must decrement used_bytes in quotas. Verifies the decrement function signature exists (full atomic SQL in PostgreSQL). """ from services.cloud_cache import decrement_quota_for_cache # Phase 14 import inspect sig = inspect.signature(decrement_quota_for_cache) params = list(sig.parameters.keys()) assert "session" in params, "decrement_quota_for_cache must accept session" assert "user_id" in params, "decrement_quota_for_cache must accept user_id" assert "size_bytes" in params, "decrement_quota_for_cache must accept size_bytes" # ─── CACHE-05: Owner isolation ─────────────────────────────────────────────── async def test_cache_entry_query_is_owner_scoped(db_session): """CACHE-05/T-14-01: Cache metadata queries are scoped to owner. Foreign user must not see another user's cache entries. """ from services.cloud_cache import create_cache_entry, list_cache_entries # Phase 14 owner_id = uuid.uuid4() foreign_id = uuid.uuid4() conn_id = uuid.uuid4() # Create entry for owner await create_cache_entry( session=db_session, user_id=owner_id, connection_id=conn_id, cloud_item_id=uuid.uuid4(), provider_item_id="owner-item", version_key="etag-v1", object_key=f"cache/{owner_id}/{uuid.uuid4()}.pdf", content_type="application/pdf", size_bytes=1024, ) # Foreign user should see zero entries foreign_entries = await list_cache_entries(session=db_session, user_id=foreign_id) assert len(foreign_entries) == 0, ( "Cache entries must be invisible to non-owner users (CACHE-05)" ) # Owner should see their own entries owner_entries = await list_cache_entries(session=db_session, user_id=owner_id) assert len(owner_entries) == 1 async def test_cache_api_excludes_object_key_from_response(async_client, db_session): """T-14-02: Cache settings/status API response never includes the MinIO object key. Object keys are internal and must not appear in any API response. """ from api.cloud.analysis import router as analysis_router # Phase 14 ctx = await _create_user_and_token(db_session) # GET cache settings/status endpoint resp = await async_client.get( "/api/cloud/analysis/cache", headers=ctx["headers"], ) # May return 200 (empty state) or 404/422 if the endpoint exists assert resp.status_code != 500, "Cache status endpoint must not 500" if resp.status_code == 200: body_text = resp.text # T-14-02: object_key must not appear in any response assert "object_key" not in body_text, ( "MinIO object_key must not appear in cache API responses (T-14-02)" ) async def test_admin_cannot_read_user_cache_status(async_client, db_session): """T-14-01/CACHE-05: Admin cannot access user's cache metadata via the API.""" from api.cloud.analysis import router as analysis_router # Phase 14 user_ctx = await _create_user_and_token(db_session, role="user") admin_ctx = await _create_user_and_token(db_session, role="admin") # Admin tries to read the user's cache resp = await async_client.get( "/api/cloud/analysis/cache", headers=admin_ctx["headers"], ) # Admin is blocked from cloud analysis routes entirely (same as browse/mutations) assert resp.status_code in (403, 404), ( f"Admin must not read cache metadata — got {resp.status_code}" ) # ─── Version-key helper ─────────────────────────────────────────────────────── async def test_version_key_prefers_version_over_etag(): """ANALYZE-06/D-19: Version key computation prefers 'version' field over etag. Precedence: version > etag > metadata fingerprint > content hash. Content hash must never be a pre-download requirement. """ from services.cloud_cache import compute_version_key # Phase 14 # version takes priority key_with_version = compute_version_key( provider_item_id="item-001", version="v3", etag="etag-stale", size=1024, modified_at="2026-06-01T00:00:00Z", content_type="application/pdf", ) key_from_etag_only = compute_version_key( provider_item_id="item-001", version=None, etag="etag-stable", size=1024, modified_at="2026-06-01T00:00:00Z", content_type="application/pdf", ) key_from_fingerprint = compute_version_key( provider_item_id="item-001", version=None, etag=None, size=1024, modified_at="2026-06-01T00:00:00Z", content_type="application/pdf", ) # Keys derived from different inputs must differ assert key_with_version != key_from_etag_only, ( "version-derived key must differ from etag-only key" ) assert key_from_etag_only != key_from_fingerprint, ( "etag-derived key must differ from metadata-fingerprint key" ) # Same inputs must produce stable keys (deterministic) key_repeat = compute_version_key( provider_item_id="item-001", version="v3", etag="etag-stale", size=1024, modified_at="2026-06-01T00:00:00Z", content_type="application/pdf", ) assert key_with_version == key_repeat, "compute_version_key must be deterministic" async def test_version_key_fallback_to_fingerprint_without_etag(): """ANALYZE-06/D-19: Version key falls back to metadata fingerprint when etag absent.""" from services.cloud_cache import compute_version_key # Phase 14 key = compute_version_key( provider_item_id="item-fingerprint", version=None, etag=None, size=2048, modified_at="2026-06-15T10:00:00Z", content_type="text/plain", ) assert key is not None assert isinstance(key, str) assert len(key) > 0 # Same metadata must produce same fingerprint key key_repeat = compute_version_key( provider_item_id="item-fingerprint", version=None, etag=None, size=2048, modified_at="2026-06-15T10:00:00Z", content_type="text/plain", ) assert key == key_repeat async def test_content_hash_not_required_during_version_key_computation(): """ANALYZE-06/D-20: Content hash is not computed during version key resolution. Passing content_hash=None must not cause a byte-fetch fallback. """ from services.cloud_cache import compute_version_key # Phase 14 # Providing no content_hash still produces a valid key via metadata fingerprint key = compute_version_key( provider_item_id="item-no-hash", version=None, etag=None, size=512, modified_at="2026-06-20T12:00:00Z", content_type="application/pdf", content_hash=None, # Must be accepted without triggering any download ) assert key is not None assert isinstance(key, str) # ─── retain_or_reuse_cache_entry ───────────────────────────────────────────── async def test_retain_or_reuse_returns_existing_entry(db_session): """CACHE-03: retain_or_reuse_cache_entry returns an existing non-evicted entry.""" from services.cloud_cache import create_cache_entry, retain_or_reuse_cache_entry user_id = uuid.uuid4() conn_id = uuid.uuid4() item_id = uuid.uuid4() version_key = "etag-shared-v1" # Create the first entry as the existing cached bytes first = await create_cache_entry( session=db_session, user_id=user_id, connection_id=conn_id, cloud_item_id=item_id, provider_item_id="pitem-retain", version_key=version_key, object_key=f"cache/{user_id}/{uuid.uuid4()}.pdf", content_type="application/pdf", size_bytes=1024, ) # retain_or_reuse must return the existing entry, not create a new one reused, created_new = await retain_or_reuse_cache_entry( session=db_session, user_id=user_id, connection_id=conn_id, cloud_item_id=item_id, provider_item_id="pitem-retain", version_key=version_key, object_key=f"cache/{user_id}/{uuid.uuid4()}-new.pdf", # different key — should not be used size_bytes=1024, ) assert not created_new, "Should reuse existing entry, not create a new one" assert str(reused.id) == str(first.id), "Must return the same DB row" async def test_retain_or_reuse_creates_new_entry_when_absent(db_session): """CACHE-03: retain_or_reuse_cache_entry creates a new entry if none exists.""" from services.cloud_cache import retain_or_reuse_cache_entry user_id = uuid.uuid4() conn_id = uuid.uuid4() item_id = uuid.uuid4() new_key = f"cache/{user_id}/{uuid.uuid4()}.pdf" entry, created_new = await retain_or_reuse_cache_entry( session=db_session, user_id=user_id, connection_id=conn_id, cloud_item_id=item_id, provider_item_id="pitem-new", version_key="etag-brand-new", object_key=new_key, size_bytes=2048, ) assert created_new, "No existing entry — must create a new one" assert entry.size_bytes == 2048 assert entry.user_id == user_id async def test_retain_or_reuse_reactivates_evicted_entry(db_session): """CACHE-03: An evicted entry for the same version key is reactivated, not duplicated. The unique constraint on (user_id, connection_id, cloud_item_id, version_key) means inserting a new row is impossible. retain_or_reuse_cache_entry reactivates the evicted row (clears evicted_at) and updates the object_key to point to fresh bytes. """ from services.cloud_cache import create_cache_entry, retain_or_reuse_cache_entry from datetime import datetime, timezone user_id = uuid.uuid4() conn_id = uuid.uuid4() item_id = uuid.uuid4() version_key = "etag-evicted" # Create and evict an entry evicted = await create_cache_entry( session=db_session, user_id=user_id, connection_id=conn_id, cloud_item_id=item_id, provider_item_id="pitem-evicted", version_key=version_key, object_key=f"cache/{user_id}/{uuid.uuid4()}.pdf", size_bytes=512, ) evicted.evicted_at = datetime.now(timezone.utc) await db_session.flush() fresh_key = f"cache/{user_id}/{uuid.uuid4()}-fresh.pdf" reactivated, created_new = await retain_or_reuse_cache_entry( session=db_session, user_id=user_id, connection_id=conn_id, cloud_item_id=item_id, provider_item_id="pitem-evicted", version_key=version_key, object_key=fresh_key, size_bytes=512, ) # created_new=True because we needed to rehydrate bytes (evicted entry was unavailable) assert created_new, "Reactivated evicted entry must report created_new=True" # Same DB row is reused (no duplicate due to unique constraint) assert str(reactivated.id) == str(evicted.id) # evicted_at must be cleared assert reactivated.evicted_at is None, "Reactivated entry must have evicted_at=None" # Object key updated to fresh bytes assert reactivated.object_key == fresh_key async def test_retain_or_reuse_isolates_foreign_user(db_session): """T-14-06: retain_or_reuse_cache_entry never returns another user's entry.""" from services.cloud_cache import create_cache_entry, retain_or_reuse_cache_entry owner_id = uuid.uuid4() foreign_id = uuid.uuid4() conn_id = uuid.uuid4() item_id = uuid.uuid4() version_key = "etag-shared" # Create entry for owner owner_entry = await create_cache_entry( session=db_session, user_id=owner_id, connection_id=conn_id, cloud_item_id=item_id, provider_item_id="pitem-owner", version_key=version_key, object_key=f"cache/{owner_id}/{uuid.uuid4()}.pdf", size_bytes=4096, ) # Foreign user must get a new entry, not the owner's foreign_entry, created_new = await retain_or_reuse_cache_entry( session=db_session, user_id=foreign_id, connection_id=conn_id, cloud_item_id=item_id, provider_item_id="pitem-owner", version_key=version_key, object_key=f"cache/{foreign_id}/{uuid.uuid4()}.pdf", size_bytes=4096, ) assert created_new, "Foreign user must not reuse another user's cache entry (T-14-06)" assert str(foreign_entry.id) != str(owner_entry.id) # ─── pin_cache_entry / release_cache_entry ──────────────────────────────────── async def test_pin_increments_pin_count(db_session): """CACHE-04: pin_cache_entry increments pin_count to protect from eviction.""" from services.cloud_cache import create_cache_entry, pin_cache_entry user_id = uuid.uuid4() conn_id = uuid.uuid4() item_id = uuid.uuid4() entry = await create_cache_entry( session=db_session, user_id=user_id, connection_id=conn_id, cloud_item_id=item_id, provider_item_id="pitem-pin", version_key="etag-pin-v1", object_key=f"cache/{user_id}/{uuid.uuid4()}.pdf", size_bytes=1024, pin_count=0, ) assert entry.pin_count == 0 pinned = await pin_cache_entry(session=db_session, entry_id=entry.id, user_id=user_id) assert pinned.pin_count == 1, "pin_cache_entry must increment pin_count" async def test_release_decrements_pin_count(db_session): """CACHE-04: release_cache_entry decrements pin_count back to 0.""" from services.cloud_cache import create_cache_entry, pin_cache_entry, release_cache_entry user_id = uuid.uuid4() conn_id = uuid.uuid4() item_id = uuid.uuid4() entry = await create_cache_entry( session=db_session, user_id=user_id, connection_id=conn_id, cloud_item_id=item_id, provider_item_id="pitem-release", version_key="etag-release-v1", object_key=f"cache/{user_id}/{uuid.uuid4()}.pdf", size_bytes=1024, ) await pin_cache_entry(session=db_session, entry_id=entry.id, user_id=user_id) released = await release_cache_entry(session=db_session, entry_id=entry.id, user_id=user_id) assert released.pin_count == 0, "release_cache_entry must decrement pin_count" async def test_release_clamps_to_zero(db_session): """CACHE-04: release_cache_entry on an already-zero pin_count stays at 0 (defensive).""" from services.cloud_cache import create_cache_entry, release_cache_entry user_id = uuid.uuid4() conn_id = uuid.uuid4() item_id = uuid.uuid4() entry = await create_cache_entry( session=db_session, user_id=user_id, connection_id=conn_id, cloud_item_id=item_id, provider_item_id="pitem-clamp", version_key="etag-clamp", object_key=f"cache/{user_id}/{uuid.uuid4()}.pdf", size_bytes=512, pin_count=0, ) released = await release_cache_entry(session=db_session, entry_id=entry.id, user_id=user_id) assert released.pin_count == 0, "Double-release must not produce negative pin_count" async def test_pin_raises_on_foreign_entry(db_session): """T-14-06: pin_cache_entry raises ValueError on foreign-user entry.""" from services.cloud_cache import create_cache_entry, pin_cache_entry owner_id = uuid.uuid4() foreign_id = uuid.uuid4() conn_id = uuid.uuid4() entry = await create_cache_entry( session=db_session, user_id=owner_id, connection_id=conn_id, cloud_item_id=uuid.uuid4(), provider_item_id="pitem-foreign-pin", version_key="etag-foreign", object_key=f"cache/{owner_id}/{uuid.uuid4()}.pdf", size_bytes=512, ) with pytest.raises(ValueError, match="not found"): await pin_cache_entry(session=db_session, entry_id=entry.id, user_id=foreign_id) # ─── update_cache_access ───────────────────────────────────────────────────── # ─── hydrate_and_cache_bytes (Plan 06) ─────────────────────────────────────── class _FakeMinIO: """Minimal MinIO stub for hydrate_and_cache_bytes tests.""" def __init__(self, bytes_to_return: bytes = b"fake-content"): self._store: dict = {} self._bytes_to_return = bytes_to_return self.put_call_count = 0 self.get_call_count = 0 self.delete_call_count = 0 async def get_object(self, key: str) -> bytes: self.get_call_count += 1 if key in self._store: return self._store[key] raise Exception(f"Key not found: {key}") async def put_object(self, key: str, data: bytes, content_type: str = None) -> None: self.put_call_count += 1 self._store[key] = data async def delete_object(self, key: str) -> None: self.delete_call_count += 1 self._store.pop(key, None) async def test_hydrate_cache_miss_fetches_from_provider(db_session): """CACHE-03/Plan 06: Cache miss calls the provider fetch_fn and stores bytes in MinIO.""" from services.cloud_cache import hydrate_and_cache_bytes, list_cache_entries from unittest.mock import AsyncMock, patch user_id = uuid.uuid4() conn_id = uuid.uuid4() item_id = uuid.uuid4() version_key = "etag-miss-v1" expected_bytes = b"provider-content-hello" fetch_call_count = 0 async def fetch_fn(): nonlocal fetch_call_count fetch_call_count += 1 return expected_bytes minio = _FakeMinIO() # Patch increment_quota_for_cache to bypass quota check (no quota row in test DB) with patch("services.cloud_cache.increment_quota_for_cache", new=AsyncMock()): result_bytes, cache_entry_id = await hydrate_and_cache_bytes( db_session, user_id=user_id, connection_id=conn_id, cloud_item_id=item_id, provider_item_id="pitem-miss", version_key=version_key, content_type="application/pdf", fetch_fn=fetch_fn, minio_client=minio, cache_limit_bytes=512 * 1024 * 1024, ) assert result_bytes == expected_bytes, "Returned bytes must match provider bytes" assert fetch_call_count == 1, "Cache miss must call fetch_fn exactly once" assert minio.put_call_count == 1, "Cache miss must store bytes in MinIO" # A cache entry must have been created for the user. entries = await list_cache_entries(db_session, user_id=user_id) assert len(entries) == 1, "A cache entry must be created on cache miss" assert entries[0].version_key == version_key async def test_hydrate_cache_hit_skips_provider(db_session): """CACHE-03/Plan 06: Cache hit reads from MinIO and does NOT call fetch_fn (provider). The fetch_fn must not be called when a non-evicted entry exists for the version_key. """ from services.cloud_cache import ( create_cache_entry, hydrate_and_cache_bytes, ) user_id = uuid.uuid4() conn_id = uuid.uuid4() item_id = uuid.uuid4() version_key = "etag-hit-v1" cached_bytes = b"minio-cached-bytes" object_key = f"cache/{user_id}/{uuid.uuid4()}.pdf" # Pre-populate a cache entry entry = await create_cache_entry( session=db_session, user_id=user_id, connection_id=conn_id, cloud_item_id=item_id, provider_item_id="pitem-hit", version_key=version_key, object_key=object_key, content_type="application/pdf", size_bytes=len(cached_bytes), ) await db_session.commit() minio = _FakeMinIO() minio._store[object_key] = cached_bytes # Pre-load the fake MinIO store fetch_call_count = 0 async def fetch_fn(): nonlocal fetch_call_count fetch_call_count += 1 return b"should-not-be-called" result_bytes, cache_entry_id = await hydrate_and_cache_bytes( db_session, user_id=user_id, connection_id=conn_id, cloud_item_id=item_id, provider_item_id="pitem-hit", version_key=version_key, content_type="application/pdf", fetch_fn=fetch_fn, minio_client=minio, cache_limit_bytes=512 * 1024 * 1024, ) assert result_bytes == cached_bytes, "Cache hit must return MinIO bytes, not provider bytes" assert fetch_call_count == 0, "Cache hit must NOT call fetch_fn (provider avoided)" assert minio.put_call_count == 0, "Cache hit must NOT write to MinIO again" async def test_hydrate_pin_released_after_cache_hit(db_session): """CACHE-04/Plan 06: Pin is released after hydrate_and_cache_bytes returns on cache hit.""" from services.cloud_cache import ( create_cache_entry, hydrate_and_cache_bytes, list_cache_entries, ) user_id = uuid.uuid4() conn_id = uuid.uuid4() item_id = uuid.uuid4() version_key = "etag-pin-release" cached_bytes = b"pin-release-test-bytes" object_key = f"cache/{user_id}/{uuid.uuid4()}.pdf" entry = await create_cache_entry( session=db_session, user_id=user_id, connection_id=conn_id, cloud_item_id=item_id, provider_item_id="pitem-pin-rel", version_key=version_key, object_key=object_key, content_type="application/pdf", size_bytes=len(cached_bytes), ) await db_session.commit() minio = _FakeMinIO() minio._store[object_key] = cached_bytes async def fetch_fn(): return b"fallback" await hydrate_and_cache_bytes( db_session, user_id=user_id, connection_id=conn_id, cloud_item_id=item_id, provider_item_id="pitem-pin-rel", version_key=version_key, content_type="application/pdf", fetch_fn=fetch_fn, minio_client=minio, cache_limit_bytes=512 * 1024 * 1024, ) # After hydration completes, pin_count must be back to 0. await db_session.refresh(entry) assert entry.pin_count == 0, ( "Pin must be released after hydrate_and_cache_bytes returns (PATTERNS.md step 8)" ) async def test_hydrate_cache_miss_pin_released(db_session): """CACHE-04/Plan 06: Pin is released after cache-miss hydration completes.""" from services.cloud_cache import hydrate_and_cache_bytes, list_cache_entries from unittest.mock import AsyncMock, patch user_id = uuid.uuid4() conn_id = uuid.uuid4() item_id = uuid.uuid4() version_key = "etag-miss-pin-release" expected_bytes = b"miss-pin-release-bytes" async def fetch_fn(): return expected_bytes minio = _FakeMinIO() with patch("services.cloud_cache.increment_quota_for_cache", new=AsyncMock()): await hydrate_and_cache_bytes( db_session, user_id=user_id, connection_id=conn_id, cloud_item_id=item_id, provider_item_id="pitem-miss-pin", version_key=version_key, content_type="application/pdf", fetch_fn=fetch_fn, minio_client=minio, cache_limit_bytes=512 * 1024 * 1024, ) entries = await list_cache_entries(db_session, user_id=user_id) assert all(e.pin_count == 0 for e in entries), ( "All cache entries must have pin_count == 0 after hydration (no pin leak)" ) async def test_hydrate_foreign_user_gets_separate_cache_entry(db_session): """T-14-06/Plan 06: hydrate_and_cache_bytes never reuses another user's cache entry.""" from services.cloud_cache import ( create_cache_entry, hydrate_and_cache_bytes, list_cache_entries, ) from unittest.mock import AsyncMock, patch owner_id = uuid.uuid4() foreign_id = uuid.uuid4() conn_id = uuid.uuid4() item_id = uuid.uuid4() version_key = "etag-shared-user-isolation" cached_bytes = b"owner-only-bytes" object_key = f"cache/{owner_id}/{uuid.uuid4()}.pdf" # Create owner's cache entry await create_cache_entry( session=db_session, user_id=owner_id, connection_id=conn_id, cloud_item_id=item_id, provider_item_id="pitem-owner-iso", version_key=version_key, object_key=object_key, size_bytes=len(cached_bytes), ) await db_session.commit() minio = _FakeMinIO(bytes_to_return=b"foreign-bytes") fetch_call_count = 0 async def fetch_fn(): nonlocal fetch_call_count fetch_call_count += 1 return b"foreign-bytes" # Patch increment_quota_for_cache to bypass quota check (no quota row in test DB) with patch("services.cloud_cache.increment_quota_for_cache", new=AsyncMock()): await hydrate_and_cache_bytes( db_session, user_id=foreign_id, connection_id=conn_id, cloud_item_id=item_id, provider_item_id="pitem-owner-iso", version_key=version_key, content_type="application/pdf", fetch_fn=fetch_fn, minio_client=minio, cache_limit_bytes=512 * 1024 * 1024, ) # Foreign user must have had a cache miss (fetched from provider) assert fetch_call_count == 1, ( "Foreign user must not reuse owner's cache entry (T-14-06)" ) # Foreign user now has their own entry foreign_entries = await list_cache_entries(db_session, user_id=foreign_id) assert len(foreign_entries) == 1 assert str(foreign_entries[0].user_id) == str(foreign_id) async def test_update_cache_access_touches_last_accessed(db_session): """CACHE-04: update_cache_access refreshes last_accessed_at for LRU ordering.""" from services.cloud_cache import create_cache_entry, update_cache_access from datetime import datetime, timezone, timedelta user_id = uuid.uuid4() conn_id = uuid.uuid4() old_ts = datetime.now(timezone.utc) - timedelta(hours=3) entry = await create_cache_entry( session=db_session, user_id=user_id, connection_id=conn_id, cloud_item_id=uuid.uuid4(), provider_item_id="pitem-access", version_key="etag-access", object_key=f"cache/{user_id}/{uuid.uuid4()}.pdf", size_bytes=256, last_accessed_at=old_ts, ) updated = await update_cache_access( session=db_session, entry_id=entry.id, user_id=user_id ) assert updated.last_accessed_at > old_ts, ( "update_cache_access must refresh last_accessed_at beyond the old timestamp" )