feat: H-CITE + H-RET + H-DATA-LIFE — evidence references in ask_knowledge, knowledge retention ARQ cron job (daily 05:00), re-extraction hook on wiki.article.updated
Check Cross-Plugin Imports / check (push) Has been cancelled
Check Cross-Plugin Imports / check (push) Has been cancelled
This commit is contained in:
@@ -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),
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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),
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user