diff --git a/tests/test_retention.py b/tests/test_retention.py new file mode 100644 index 0000000..95056d7 --- /dev/null +++ b/tests/test_retention.py @@ -0,0 +1,79 @@ +"""Tests for content retention cleanup of the posts.content column.""" + +from __future__ import annotations + +from unittest.mock import AsyncMock, MagicMock, patch + +import pytest + +import voucherbot.services.retention as retention + + +def _fake_session(rowcount: int = 0) -> MagicMock: + session = MagicMock() + session.execute = AsyncMock(return_value=MagicMock(rowcount=rowcount)) + session.commit = AsyncMock() + return session + + +async def _run_purge(session: MagicMock) -> int: + with patch("voucherbot.services.retention.session_scope") as scope: + scope.return_value.__aenter__ = AsyncMock(return_value=session) + scope.return_value.__aexit__ = AsyncMock(return_value=False) + return await retention.purge_expired_post_content() + + +@pytest.mark.asyncio +async def test_purge_only_targets_content_of_old_posts() -> None: + session = _fake_session(rowcount=3) + rows = await _run_purge(session) + + assert rows == 3 + session.commit.assert_awaited_once() + + stmt = session.execute.await_args.args[0] + sql = str(stmt) + assert "UPDATE posts" in sql + assert "SET content = NULL" in sql + assert "created_at" in sql + assert "cutoff" in str(stmt) + + +@pytest.mark.asyncio +async def test_purge_never_touches_other_columns() -> None: + """The UPDATE's SET clause must contain exactly one column: content.""" + session = _fake_session(rowcount=1) + await _run_purge(session) + + stmt = session.execute.await_args.args[0] + sql = str(stmt) + set_clause = sql.split("SET", 1)[1].split("WHERE", 1)[0] + assert set_clause.strip() == "content = NULL" + assert "updated_at" not in set_clause + assert "summary" not in set_clause + + +@pytest.mark.asyncio +async def test_purge_skips_already_null_content() -> None: + session = _fake_session(rowcount=0) + rows = await _run_purge(session) + + assert rows == 0 + stmt = session.execute.await_args.args[0] + assert "IS NOT NULL" in str(stmt) + + +@pytest.mark.asyncio +async def test_purge_logs_and_swallows_errors() -> None: + session = MagicMock() + session.execute = AsyncMock(side_effect=RuntimeError("db down")) + with ( + patch("voucherbot.services.retention.session_scope") as scope, + patch("voucherbot.services.retention.logger.warning") as warn, + ): + scope.return_value.__aenter__ = AsyncMock(return_value=session) + scope.return_value.__aexit__ = AsyncMock(return_value=False) + rows = await retention.purge_expired_post_content() + + assert rows == 0 + warn.assert_called_once() diff --git a/voucherbot/config/settings.py b/voucherbot/config/settings.py index eec3a68..d9964ea 100644 --- a/voucherbot/config/settings.py +++ b/voucherbot/config/settings.py @@ -107,6 +107,10 @@ class Settings(BaseSettings): source_backoff_base_minutes: int = 5 source_backoff_max_minutes: int = 360 + # Content retention: posts older than this many days have their content + # column nulled out on each scheduler sweep (all other columns untouched). + content_retention_days: int = 7 + # AI providers gemini_api_key: Optional[str] = None groq_api_key: Optional[str] = None diff --git a/voucherbot/services/retention.py b/voucherbot/services/retention.py new file mode 100644 index 0000000..4c238c6 --- /dev/null +++ b/voucherbot/services/retention.py @@ -0,0 +1,63 @@ +"""Content retention: periodically null out the content of old posts. + +Posts accumulate large scraped bodies over time. Since only the most recent +posts are ever surfaced to users, this module keeps the ``posts.content`` +column bounded by nulling it for every post whose ``created_at`` is older than +``settings.content_retention_days``. Only the ``content`` column is updated — +every other column is left untouched. +""" + +from __future__ import annotations + +from datetime import datetime, timedelta, timezone +from typing import Any, cast + +import structlog +from sqlalchemy import CursorResult, text + +from voucherbot.config.settings import settings +from voucherbot.database.connection import session_scope + +logger = structlog.get_logger(__name__) + + +async def purge_expired_post_content() -> int: + """Null out ``posts.content`` for every post older than the retention window. + + Returns the number of rows updated. Never raises — failures are logged so + the scheduler loop stays healthy. + """ + cutoff = datetime.now(timezone.utc) - timedelta( + days=settings.content_retention_days + ) + try: + async with session_scope() as session: + # Raw UPDATE on purpose: a Core/ORM ``update()`` would also write + # ``updated_at = now()`` via the model's ``onupdate`` hook, but this + # retention job must touch *only* the content column. + result = await session.execute( + text( + """ + UPDATE posts + SET content = NULL + WHERE created_at < :cutoff + AND content IS NOT NULL + """ + ), + {"cutoff": cutoff}, + ) + await session.commit() + rows = cast(CursorResult[Any], result).rowcount or 0 + if rows: + logger.info( + "retention: purged expired post content", + rows=rows, + cutoff=cutoff.isoformat(), + ) + return rows + except Exception as exc: + logger.warning( + "retention: content purge sweep failed", + error=str(exc)[:160], + ) + return 0 diff --git a/voucherbot/services/scheduler.py b/voucherbot/services/scheduler.py index b106ff3..65b51cc 100644 --- a/voucherbot/services/scheduler.py +++ b/voucherbot/services/scheduler.py @@ -27,6 +27,7 @@ from voucherbot.providers.training_provider.collector import TrainingProviderCollector from voucherbot.services.dispatcher import dispatch_tick from voucherbot.services.email.notifications import retry_pending_notifications +from voucherbot.services.retention import purge_expired_post_content logger = structlog.get_logger(__name__) @@ -146,6 +147,8 @@ async def _run_loop() -> None: # Retry undelivered voucher alerts (idempotency keyed) even when # no source ran; failed rows are re-attempted on later sweeps. await retry_pending_notifications() + # Null out content of posts older than the retention window. + await purge_expired_post_content() sleep_seconds = await _seconds_until_next_due() logger.info( "scheduler: sleeping until next sweep",