"""`upload_source_file` emits `ingestion.job.*` log events (ADR-0011). Every failure branch funnels through `_mark_job_failed`, so this asserts the log event once per branch rather than re-testing the Postgres job-row behavior already covered in `test_upload_service.py`. Uses `structlog.testing.capture_logs()`, which captures events independent of whichever handlers/renderers happen to be configured in this process. """ from collections.abc import MutableMapping from typing import Any import pytest import structlog from anyio import CapacityLimiter, Semaphore from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker from src.application.auth.context import AuthContext from src.application.files.upload import upload_source_file from src.application.ingestion.errors import ( EmbedderError, IngestionTimeoutError, PointIndexingError, ) from src.config import ChunkingSettings, IngestionSettings, QdrantSettings from tests.fakes import ( FakeDenseEmbedder, FakeObjectStorage, FakePointStorage, FakeSparseEmbedder, ) from tests.support.factories import create_api_key, create_tenant, create_tenant_domain pytestmark = [ pytest.mark.integration, pytest.mark.postgres, pytest.mark.asyncio(loop_scope="session"), ] _CSV_BYTES = b"name,value\nfirst,1\n" async def _auth_for(db_session: AsyncSession, *, domain: str = "general") -> AuthContext: tenant = await create_tenant(db_session) api_key, _ = await create_api_key(db_session, tenant=tenant) await create_tenant_domain(db_session, tenant=tenant, domain=domain) await db_session.commit() return AuthContext( tenant_id=tenant.id, tenant_slug=tenant.slug, api_key_id=api_key.id, scopes=frozenset({"files:write"}), actor_type="backend", ) def _events_by_name( logs: list[MutableMapping[str, Any]], name: str ) -> list[MutableMapping[str, Any]]: return [entry for entry in logs if entry.get("event") == name] async def test_successful_upload_emits_started_and_completed_events( db_session: AsyncSession, db_sessionmaker: async_sessionmaker[AsyncSession] ) -> None: auth = await _auth_for(db_session) with structlog.testing.capture_logs() as logs: result = await upload_source_file( sessionmaker=db_sessionmaker, storage=FakeObjectStorage(), point_storage=FakePointStorage(), auth=auth, domain="general", filename="report.csv", data=_CSV_BYTES, ingestion_settings=IngestionSettings(), chunking_settings=ChunkingSettings(), qdrant_settings=QdrantSettings(), thread_limiter=CapacityLimiter(2), concurrency_limiter=Semaphore(2), dense_embedders=[ FakeDenseEmbedder(name="dense_nomic"), FakeDenseEmbedder(name="dense_openai"), ], sparse_embedder=FakeSparseEmbedder(), ) started = _events_by_name(logs, "ingestion.job.started") completed = _events_by_name(logs, "ingestion.job.completed") assert len(started) == 1 assert started[0]["tenant_id"] == str(auth.tenant_id) assert started[0]["ingestion_job_id"] == str(result.ingestion_job_id) assert len(completed) == 1 assert completed[0]["points_upserted"] == result.chunks_indexed assert _events_by_name(logs, "ingestion.job.failed") == [] async def test_storage_failure_emits_ingestion_job_failed( db_session: AsyncSession, db_sessionmaker: async_sessionmaker[AsyncSession] ) -> None: auth = await _auth_for(db_session) with structlog.testing.capture_logs() as logs, pytest.raises(OSError): await upload_source_file( sessionmaker=db_sessionmaker, storage=FakeObjectStorage(fail_next=True), point_storage=FakePointStorage(), auth=auth, domain="general", filename="report.csv", data=_CSV_BYTES, ingestion_settings=IngestionSettings(), chunking_settings=ChunkingSettings(), qdrant_settings=QdrantSettings(), thread_limiter=CapacityLimiter(2), concurrency_limiter=Semaphore(2), dense_embedders=[ FakeDenseEmbedder(name="dense_nomic"), FakeDenseEmbedder(name="dense_openai"), ], sparse_embedder=FakeSparseEmbedder(), ) failed = _events_by_name(logs, "ingestion.job.failed") assert len(failed) == 1 assert failed[0]["error_code"] == "storage_upload_failed" assert failed[0]["tenant_id"] == str(auth.tenant_id) async def test_embedding_failure_emits_ingestion_job_failed( db_session: AsyncSession, db_sessionmaker: async_sessionmaker[AsyncSession] ) -> None: """Previously silent: embedding_failed reached Postgres but never logged.""" auth = await _auth_for(db_session) failing_embedder = FakeDenseEmbedder(name="dense_nomic", fail_next=True) with structlog.testing.capture_logs() as logs, pytest.raises(EmbedderError): await upload_source_file( sessionmaker=db_sessionmaker, storage=FakeObjectStorage(), point_storage=FakePointStorage(), auth=auth, domain="general", filename="report.csv", data=_CSV_BYTES, ingestion_settings=IngestionSettings(), chunking_settings=ChunkingSettings(), qdrant_settings=QdrantSettings(), thread_limiter=CapacityLimiter(2), concurrency_limiter=Semaphore(2), dense_embedders=[failing_embedder, FakeDenseEmbedder(name="dense_openai")], sparse_embedder=FakeSparseEmbedder(), ) failed = _events_by_name(logs, "ingestion.job.failed") assert len(failed) == 1 assert failed[0]["error_code"] == "embedding_failed" async def test_index_failure_emits_ingestion_job_failed( db_session: AsyncSession, db_sessionmaker: async_sessionmaker[AsyncSession] ) -> None: """Previously silent: index_failed reached Postgres but never logged.""" auth = await _auth_for(db_session) with structlog.testing.capture_logs() as logs, pytest.raises(PointIndexingError): await upload_source_file( sessionmaker=db_sessionmaker, storage=FakeObjectStorage(), point_storage=FakePointStorage(fail_on_batch=0), auth=auth, domain="general", filename="report.csv", data=_CSV_BYTES, ingestion_settings=IngestionSettings(), chunking_settings=ChunkingSettings(), qdrant_settings=QdrantSettings(), thread_limiter=CapacityLimiter(2), concurrency_limiter=Semaphore(2), dense_embedders=[ FakeDenseEmbedder(name="dense_nomic"), FakeDenseEmbedder(name="dense_openai"), ], sparse_embedder=FakeSparseEmbedder(), ) failed = _events_by_name(logs, "ingestion.job.failed") assert len(failed) == 1 assert failed[0]["error_code"] == "index_failed" async def test_timeout_emits_ingestion_job_failed_not_a_duplicate_event( db_session: AsyncSession, db_sessionmaker: async_sessionmaker[AsyncSession] ) -> None: """The old ad-hoc `files.upload.timeout` log is gone -- `_mark_job_failed` is now the single place a failure is logged, so there is exactly one `ingestion.job.failed` event, not two events for one failure. """ auth = await _auth_for(db_session) slow_embedder = FakeDenseEmbedder(name="dense_nomic", delay_seconds=10) with structlog.testing.capture_logs() as logs, pytest.raises(IngestionTimeoutError): await upload_source_file( sessionmaker=db_sessionmaker, storage=FakeObjectStorage(), point_storage=FakePointStorage(), auth=auth, domain="general", filename="report.csv", data=_CSV_BYTES, ingestion_settings=IngestionSettings(timeout_seconds=0.05), chunking_settings=ChunkingSettings(), qdrant_settings=QdrantSettings(), thread_limiter=CapacityLimiter(2), concurrency_limiter=Semaphore(2), dense_embedders=[slow_embedder, FakeDenseEmbedder(name="dense_openai")], sparse_embedder=FakeSparseEmbedder(), ) failed = _events_by_name(logs, "ingestion.job.failed") assert len(failed) == 1 assert failed[0]["error_code"] == "timeout"