import uuid from datetime import UTC, datetime, timedelta import pytest from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncEngine, AsyncSession, async_sessionmaker from app.ingestion.handlers import RETENTION_CLEANUP, ensure_retention_scheduled from app.ingestion.queue import ( MAX_ATTEMPTS, enqueue, job_handler, process_one, ) from app.models import ( AuthSession, Conversation, ConversationMode, Department, Job, JobStatus, User, ) @pytest.fixture def session_factory(db_engine: AsyncEngine) -> async_sessionmaker[AsyncSession]: return async_sessionmaker(db_engine, expire_on_commit=False) async def _get_job(db: AsyncSession, job_id: uuid.UUID) -> Job: db.expire_all() job = await db.get(Job, job_id) assert job is not None return job async def test_successful_job_is_marked_done( db: AsyncSession, session_factory: async_sessionmaker[AsyncSession] ) -> None: seen: list[dict] = [] @job_handler("t_ok") async def handle(handler_db: AsyncSession, job: Job) -> None: seen.append(job.payload) job = await enqueue(db, "t_ok", {"n": 1}) await db.commit() assert await process_one(session_factory) is True assert seen == [{"n": 1}] refreshed = await _get_job(db, job.id) assert refreshed.status == JobStatus.done assert refreshed.attempts == 1 # Nothing left to do. assert await process_one(session_factory) is False async def test_failing_job_retries_with_backoff( db: AsyncSession, session_factory: async_sessionmaker[AsyncSession] ) -> None: @job_handler("t_fail") async def handle(handler_db: AsyncSession, job: Job) -> None: raise ValueError("boom") job = await enqueue(db, "t_fail") await db.commit() assert await process_one(session_factory) is True refreshed = await _get_job(db, job.id) assert refreshed.status == JobStatus.pending assert refreshed.attempts == 1 assert refreshed.last_error is not None assert "ValueError: boom" in refreshed.last_error assert refreshed.run_after > datetime.now(UTC) + timedelta(seconds=10) # Backed off into the future: not claimable right now. assert await process_one(session_factory) is False async def test_job_fails_permanently_after_max_attempts( db: AsyncSession, session_factory: async_sessionmaker[AsyncSession] ) -> None: @job_handler("t_exhaust") async def handle(handler_db: AsyncSession, job: Job) -> None: raise RuntimeError("always broken") job = await enqueue(db, "t_exhaust") await db.commit() for _ in range(MAX_ATTEMPTS): refreshed = await _get_job(db, job.id) refreshed.run_after = datetime.now(UTC) - timedelta(seconds=1) await db.commit() assert await process_one(session_factory) is True refreshed = await _get_job(db, job.id) assert refreshed.status == JobStatus.failed assert refreshed.attempts == MAX_ATTEMPTS async def test_unknown_job_type_records_error( db: AsyncSession, session_factory: async_sessionmaker[AsyncSession] ) -> None: job = await enqueue(db, "t_nobody_home") await db.commit() assert await process_one(session_factory) is True refreshed = await _get_job(db, job.id) assert refreshed.status == JobStatus.pending assert refreshed.last_error is not None assert "LookupError" in refreshed.last_error async def test_locked_job_is_skipped_and_claimable_after_rollback( db: AsyncSession, session_factory: async_sessionmaker[AsyncSession] ) -> None: """SKIP LOCKED + crash-safety: a claim held by a dying worker (open tx) is invisible to others and becomes claimable again on rollback.""" @job_handler("t_locked") async def handle(handler_db: AsyncSession, job: Job) -> None: pass job = await enqueue(db, "t_locked") await db.commit() async with session_factory() as other: claimed = ( await other.execute( select(Job).where(Job.id == job.id).with_for_update(skip_locked=True) ) ).scalar_one() assert claimed.id == job.id # Row is locked by "another worker": nothing to process. assert await process_one(session_factory) is False await other.rollback() # the worker "crashes" # After the rollback the job is claimable again. assert await process_one(session_factory) is True refreshed = await _get_job(db, job.id) assert refreshed.status == JobStatus.done async def test_handler_writes_roll_back_atomically_on_failure( db: AsyncSession, session_factory: async_sessionmaker[AsyncSession] ) -> None: @job_handler("t_atomic") async def handle(handler_db: AsyncSession, job: Job) -> None: handler_db.add(Department(name="Ghost Department")) await handler_db.flush() raise RuntimeError("after write") await enqueue(db, "t_atomic") await db.commit() assert await process_one(session_factory) is True ghost = ( await db.execute( select(Department).where(Department.name == "Ghost Department") ) ).scalar_one_or_none() assert ghost is None async def test_retention_cleanup( db: AsyncSession, session_factory: async_sessionmaker[AsyncSession], seeded_user: User, ) -> None: now = datetime.now(UTC) old = now - timedelta(days=120) old_query = Conversation( mode=ConversationMode.query, user_id=seeded_user.id, updated_at=old ) fresh_query = Conversation(mode=ConversationMode.query, user_id=seeded_user.id) # A non-query mode (EE insight) must survive retention: only ephemeral # query threads are cleaned up. old_insight = Conversation( mode=ConversationMode.insight, user_id=seeded_user.id, updated_at=old ) expired_session = AuthSession( user_id=seeded_user.id, expires_at=now - timedelta(days=1) ) valid_session = AuthSession( user_id=seeded_user.id, expires_at=now + timedelta(days=1) ) db.add_all([old_query, fresh_query, old_insight, expired_session, valid_session]) await enqueue(db, RETENTION_CLEANUP) await db.commit() assert await process_one(session_factory) is True db.expire_all() remaining_conversations = { c.id for c in (await db.execute(select(Conversation))).scalars() } assert remaining_conversations == {fresh_query.id, old_insight.id} remaining_sessions = { s.id for s in (await db.execute(select(AuthSession))).scalars() } assert remaining_sessions == {valid_session.id} # Rescheduled itself for tomorrow. next_job = ( await db.execute( select(Job).where( Job.type == RETENTION_CLEANUP, Job.status == JobStatus.pending ) ) ).scalar_one() assert next_job.run_after > now + timedelta(hours=23) async def test_ensure_retention_scheduled_is_idempotent( db: AsyncSession, session_factory: async_sessionmaker[AsyncSession] ) -> None: async with session_factory() as first: await ensure_retention_scheduled(first) async with session_factory() as second: await ensure_retention_scheduled(second) jobs = ( (await db.execute(select(Job).where(Job.type == RETENTION_CLEANUP))) .scalars() .all() ) assert len(jobs) == 1 async def test_retention_respects_configured_days( db: AsyncSession, session_factory: async_sessionmaker[AsyncSession], seeded_user: User, monkeypatch: pytest.MonkeyPatch, ) -> None: """End-to-end with a short retention window: 1 day keeps yesterday's conversation out of scope for deletion at 20h but purges a 30h one.""" from app.config import get_settings monkeypatch.setenv("PABLAN_QUERY_RETENTION_DAYS", "1") get_settings.cache_clear() try: now = datetime.now(UTC) too_old = Conversation( mode=ConversationMode.query, user_id=seeded_user.id, updated_at=now - timedelta(hours=30), ) still_fresh = Conversation( mode=ConversationMode.query, user_id=seeded_user.id, updated_at=now - timedelta(hours=20), ) db.add_all([too_old, still_fresh]) await enqueue(db, RETENTION_CLEANUP) await db.commit() assert await process_one(session_factory) is True db.expire_all() remaining = {c.id for c in (await db.execute(select(Conversation))).scalars()} assert remaining == {still_fresh.id} finally: get_settings.cache_clear()