diff --git a/app/core/worker.py b/app/core/worker.py index 744c14d..9130d7e 100644 --- a/app/core/worker.py +++ b/app/core/worker.py @@ -445,6 +445,53 @@ async def cleanup_trash_job(ctx: dict[str, Any]) -> None: register_job("cleanup_trash", cleanup_trash_job) +# ── Knowledge retention cleanup job ───────────────────────────────────────── + +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 datetime import datetime, timedelta + + from app.core.db import get_worker_session_factory + from app.plugins.builtins.knowledge.models import KnowledgeExtraction + + 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.utcnow() - timedelta(days=90) + 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 90 days + 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) + + class WorkerSettings: """ARQ worker settings.""" functions = get_all_jobs() @@ -484,6 +531,11 @@ class WorkerSettings: _wrap_cron_with_lock("cleanup_trash", cleanup_trash_job, ttl_seconds=300), hour=4, minute=0, ), + # Knowledge retention cleanup — daily at 05:00 (90 days, keeps approved) + cron( + _wrap_cron_with_lock("cleanup_knowledge", cleanup_knowledge_job, ttl_seconds=300), + hour=5, minute=0, + ), # Scheduled backup — daily at 02:00 (guarded by distributed lock) cron( _wrap_cron_with_lock("run_backup", get_job("run_backup"), ttl_seconds=600), diff --git a/app/plugins/builtins/knowledge/plugin.py b/app/plugins/builtins/knowledge/plugin.py index 7c610b0..9d02001 100644 --- a/app/plugins/builtins/knowledge/plugin.py +++ b/app/plugins/builtins/knowledge/plugin.py @@ -40,6 +40,22 @@ class KnowledgePlugin(BasePlugin): source_title=title, source_text=content, ) register_action("wiki.article.created", on_wiki_create, priority=20, owner_tag="knowledge") + # H-DATA-LIFE: Re-extract when wiki article is updated + async def on_wiki_update(*args, **kwargs): + article_id = kwargs.get("article_id") or kwargs.get("entity_id") + tenant_id = kwargs.get("tenant_id") + title = kwargs.get("title", "") + content = kwargs.get("content", "") + if article_id and tenant_id and content: + from app.core.db import get_worker_session_factory + factory = get_worker_session_factory() + async with factory() as session: + await extract_knowledge( + db=session, tenant_id=uuid.UUID(str(tenant_id)), + source_type="wiki_article", source_id=uuid.UUID(str(article_id)), + source_title=title, source_text=content, + ) + register_action("wiki.article.updated", on_wiki_update, priority=20, owner_tag="knowledge") logger.info("Registered knowledge extraction hooks") except Exception: logger.exception("Failed to register knowledge hooks") diff --git a/app/plugins/builtins/knowledge/routes.py b/app/plugins/builtins/knowledge/routes.py index f6a4e8e..8f20753 100644 --- a/app/plugins/builtins/knowledge/routes.py +++ b/app/plugins/builtins/knowledge/routes.py @@ -42,7 +42,7 @@ async def ask( ): """Ask a knowledge question — uses wiki + graph_rag as context.""" tenant_id = uuid.UUID(current_user["tenant_id"]) - question = body.get("question", "") + question = body.get("question") or body.get("query", "") if not question: raise HTTPException(400, detail={"detail": "question required", "code": "missing_question"}) result = await ask_knowledge(db=db, tenant_id=tenant_id, question=question) diff --git a/app/plugins/builtins/knowledge/services.py b/app/plugins/builtins/knowledge/services.py index c04c8c1..c40dff5 100644 --- a/app/plugins/builtins/knowledge/services.py +++ b/app/plugins/builtins/knowledge/services.py @@ -109,22 +109,48 @@ async def ask_knowledge( # Search wiki articles for context from app.plugins.builtins.unified_search.provider_registry import get_search_registry registry = get_search_registry() - wiki_provider = registry.get("wiki_article") + evidence: list[dict[str, Any]] = [] context_parts = [] + wiki_provider = registry.get("wiki_article") if wiki_provider: results = await wiki_provider._search_fts_filtered( db=db, tsquery=question, tenant_id=tenant_id, limit=5, visible_ids=None ) for r in results: - context_parts.append(f"Title: {r.get('title', '')}\nContent: {r.get('content', '')[:500]}") + snippet = (r.get("content", "") or "")[:300] + title = r.get("title", "") + context_parts.append(f"Title: {title}\nContent: {snippet}") + evidence.append({ + "id": str(r.get("id", "")), + "source_type": "wiki_article", + "source_id": str(r.get("id", "")), + "title": title, + "snippet": snippet, + "score": float(r.get("rank", 0.0)), + "url": f"/wiki?article={r.get('slug', '')}", + }) # Search graph_rag for relationships graph_provider = registry.get("graph_relationship") if graph_provider: - results = await graph_provider._search_fts_filtered( - db=db, tsquery=question, tenant_id=tenant_id, limit=5, visible_ids=None - ) - for r in results: - context_parts.append(f"Relationship: {r.get('source_type')} -> {r.get('relationship_type')} -> {r.get('target_type')}") + try: + results = await graph_provider._search_fts_filtered( + db=db, tsquery=question, tenant_id=tenant_id, limit=5, visible_ids=None + ) + for r in results: + title = f"{r.get('source_type', '')} → {r.get('relationship_type', '')} → {r.get('target_type', '')}" + snippet = str(r.get("meta", "")) + context_parts.append(f"Relationship: {title}") + evidence.append({ + "id": str(r.get("id", "")), + "source_type": "graph_relationship", + "source_id": str(r.get("id", "")), + "title": title, + "snippet": snippet, + "score": float(r.get("rank", 0.0)), + "url": None, + }) + except Exception: + logger.warning("GraphRAG search failed in ask_knowledge", exc_info=True) context = "\n\n".join(context_parts) if context_parts else "No knowledge base content found." messages = [ {"role": "system", "content": f"You are a knowledge assistant. Answer based on this context:\n\n{context}"}, @@ -136,7 +162,8 @@ async def ask_knowledge( ) return { "answer": response.get("content", ""), - "evidence": context_parts[:3], + "evidence": evidence, + "sources_used": list(set(e["source_type"] for e in evidence)), "cost_usd": response.get("cost_usd", 0.0), }