63 lines
2.2 KiB
Python
63 lines
2.2 KiB
Python
"""ARQ background jobs for the knowledge plugin.
|
|
|
|
Registered via ``register_job()`` at import time; the worker discovers this
|
|
module through ``KnowledgePlugin.get_job_modules()`` — no core imports of
|
|
plugin models needed.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
from datetime import UTC, datetime, timedelta
|
|
from typing import Any
|
|
|
|
from app.core.job_registry import register_job
|
|
from app.plugins.builtins.knowledge.models import KnowledgeExtraction
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
_KNOWLEDGE_RETENTION_DAYS = 90
|
|
|
|
|
|
async def cleanup_knowledge_job(ctx: dict[str, Any]) -> None:
|
|
"""Delete old knowledge extractions (rejected or auto_created) older than 90 days.
|
|
|
|
Runs daily. Keeps approved extractions indefinitely.
|
|
Iterates per-tenant for RLS compliance.
|
|
"""
|
|
from sqlalchemy import text as sa_text, delete as sa_delete
|
|
|
|
from app.core.db import get_worker_session_factory
|
|
|
|
factory = get_worker_session_factory()
|
|
async with factory() as db:
|
|
try:
|
|
tenant_result = await db.execute(sa_text("SELECT id FROM tenants"))
|
|
tenant_ids = [row[0] for row in tenant_result]
|
|
|
|
cutoff = datetime.now(UTC).replace(tzinfo=None) - timedelta(days=_KNOWLEDGE_RETENTION_DAYS)
|
|
total_deleted = 0
|
|
for tenant_id in tenant_ids:
|
|
await db.execute(
|
|
sa_text("SELECT set_config('app.current_tenant_id', :tid, true)"),
|
|
{"tid": str(tenant_id)},
|
|
)
|
|
# Delete rejected and auto_created extractions older than retention window
|
|
result = await db.execute(
|
|
sa_delete(KnowledgeExtraction).where(
|
|
KnowledgeExtraction.status.in_(["rejected", "auto_created"]),
|
|
KnowledgeExtraction.created_at < cutoff,
|
|
)
|
|
)
|
|
total_deleted += result.rowcount
|
|
await db.commit()
|
|
|
|
if total_deleted:
|
|
logger.info("Knowledge retention: cleaned up %d old extractions", total_deleted)
|
|
except Exception:
|
|
logger.error("Knowledge retention cleanup failed", exc_info=True)
|
|
await db.rollback()
|
|
|
|
|
|
register_job("cleanup_knowledge", cleanup_knowledge_job)
|