feat: add read-only nightly integrity-check maintenance endpoint
POST /maintenance/integrity-check {user} reports per tenant, without
ever fixing anything:
- wiki pages with ZERO vectors in Qdrant (silent-skip reindex victims)
- orphaned vectors whose wiki page no longer exists
- unexpected Qdrant collections vs known tenant patterns (test-tenant
residue and unknown namespaces flagged; foreign services counted)
- Neo4j Document nodes without wiki counterparts
- counts + duration_ms
The latest report is cached in Redis (library:integrity:latest:{user},
30-day TTL) so the weekly quality report can fold it in. Explicit user
required per Phase B. Offline tests assert the report contents, the
collection classification rules, and that no destructive client method
is ever invoked.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01QbFZyDvYksazX6nYQYZ67L
This commit is contained in:
@@ -21,6 +21,7 @@ from src.core.dependencies import (
|
||||
VectorServiceDep, GraphServiceDep, WikiJSDep, RedisDep,
|
||||
QdrantDep, OllamaDep, PaperlessDep, verify_api_key
|
||||
)
|
||||
from src.core.multi_tenancy import RequiredUser, sanitize_user_id
|
||||
from src.config import get_settings
|
||||
from datetime import datetime, timezone
|
||||
|
||||
@@ -1098,3 +1099,246 @@ async def reconcile_index(
|
||||
except Exception as e:
|
||||
logger.error(f"Reconcile-index failed: {e}", exc_info=True)
|
||||
raise HTTPException(status_code=500, detail=str(e))
|
||||
|
||||
|
||||
# ========== Integrity check (nightly, read-only) ==========
|
||||
|
||||
# Redis key holding the latest integrity report per tenant (folded into the
|
||||
# weekly quality report).
|
||||
INTEGRITY_LATEST_KEY = "library:integrity:latest:{user}"
|
||||
INTEGRITY_LATEST_TTL = 86400 * 30 # 30 days
|
||||
|
||||
#: Qdrant collection prefixes owned by library-desk.
|
||||
LIBRARY_COLLECTION_PREFIXES = ("library_desk_", "volatile_")
|
||||
|
||||
|
||||
def _looks_like_test_tenant(name: str) -> bool:
|
||||
"""Heuristic for test/probe residue in collection or tenant names."""
|
||||
lowered = name.lower()
|
||||
return (
|
||||
"llm_tester" in lowered
|
||||
or "llm-tester" in lowered
|
||||
or "test" in lowered
|
||||
or lowered.startswith("verify_probe")
|
||||
or lowered.startswith("verify-probe")
|
||||
)
|
||||
|
||||
|
||||
def classify_collection(name: str, known_tenants: set[str]) -> str:
|
||||
"""
|
||||
Classify a Qdrant collection against known tenant patterns.
|
||||
|
||||
Returns one of:
|
||||
- ``expected``: library-desk collection for a tenant with a wiki namespace
|
||||
- ``test_residue``: library-desk collection for a test/probe tenant
|
||||
- ``unknown_tenant``: library-desk collection for a tenant with no wiki
|
||||
namespace (orphaned or mis-scoped)
|
||||
- ``foreign_test_residue``: another service's collection that looks like
|
||||
test residue (reported, but owned elsewhere)
|
||||
- ``foreign``: another service's collection (informational only)
|
||||
"""
|
||||
for prefix in LIBRARY_COLLECTION_PREFIXES:
|
||||
if name.startswith(prefix):
|
||||
tenant = name[len(prefix):]
|
||||
if _looks_like_test_tenant(tenant):
|
||||
return "test_residue"
|
||||
if tenant in known_tenants:
|
||||
return "expected"
|
||||
return "unknown_tenant"
|
||||
if _looks_like_test_tenant(name):
|
||||
return "foreign_test_residue"
|
||||
return "foreign"
|
||||
|
||||
|
||||
class IntegrityCheckRequest(BaseModel):
|
||||
"""Request body for /maintenance/integrity-check."""
|
||||
user: RequiredUser = Field(
|
||||
...,
|
||||
description="User identifier (tenant). Required — the report is scoped to this tenant."
|
||||
)
|
||||
|
||||
|
||||
class IntegrityCheckResponse(BaseModel):
|
||||
"""Read-only integrity report for one tenant."""
|
||||
success: bool
|
||||
user: str
|
||||
generated_at: str
|
||||
pages_without_vectors: List[Dict[str, Any]] = Field(
|
||||
default_factory=list,
|
||||
description="Wiki pages with ZERO vectors in Qdrant (silent-skip reindex victims)"
|
||||
)
|
||||
orphaned_vector_chunks: int = Field(
|
||||
default=0, description="Vector chunks whose wiki page no longer exists"
|
||||
)
|
||||
orphaned_vector_page_ids: List[int] = Field(
|
||||
default_factory=list, description="Distinct stale page ids referenced by orphaned chunks"
|
||||
)
|
||||
unexpected_collections: List[Dict[str, str]] = Field(
|
||||
default_factory=list,
|
||||
description="Qdrant collections flagged as test residue or unknown tenants"
|
||||
)
|
||||
foreign_collections: int = Field(
|
||||
default=0, description="Collections owned by other services (informational)"
|
||||
)
|
||||
documents_without_wiki: List[Dict[str, Any]] = Field(
|
||||
default_factory=list,
|
||||
description="Neo4j Document nodes whose wiki page no longer exists"
|
||||
)
|
||||
counts: Dict[str, int] = Field(default_factory=dict)
|
||||
duration_ms: float = 0.0
|
||||
|
||||
|
||||
async def run_integrity_check(
|
||||
user: str,
|
||||
vector_service: VectorService,
|
||||
graph_service: GraphService,
|
||||
wiki_client,
|
||||
qdrant
|
||||
) -> IntegrityCheckResponse:
|
||||
"""
|
||||
Run the read-only integrity check for one tenant.
|
||||
|
||||
Reports (never fixes):
|
||||
1. Wiki pages with zero vectors in the tenant's Qdrant collection
|
||||
2. Orphaned vectors whose wiki page no longer exists
|
||||
3. Unexpected Qdrant collections (test residue / unknown tenants)
|
||||
4. Neo4j Document nodes without wiki counterparts
|
||||
"""
|
||||
start_time = time.time()
|
||||
tenant_prefix = f"users/{sanitize_user_id(user)}"
|
||||
|
||||
# One unfiltered listing serves both the tenant scan and the
|
||||
# known-tenant derivation for collection classification.
|
||||
all_pages = await wiki_client.list_all_pages()
|
||||
tenant_pages = [
|
||||
p for p in all_pages
|
||||
if ("/" + str(p.get("path", "")).lstrip("/")).startswith("/" + tenant_prefix)
|
||||
]
|
||||
|
||||
known_tenants = set()
|
||||
for p in all_pages:
|
||||
parts = str(p.get("path", "")).lstrip("/").split("/")
|
||||
if len(parts) >= 2 and parts[0] == "users":
|
||||
known_tenants.add(sanitize_user_id(parts[1]))
|
||||
|
||||
tenant_page_ids = {p["id"] for p in tenant_pages if p.get("id")}
|
||||
|
||||
# Vector side (tenant collection only)
|
||||
chunk_refs = await vector_service.get_all_chunk_references(user)
|
||||
wiki_chunk_refs = [r for r in chunk_refs if r.get("doc_type", "wiki") == "wiki"]
|
||||
vectorized_page_ids = {r["page_id"] for r in wiki_chunk_refs if r.get("page_id")}
|
||||
|
||||
pages_without_vectors = [
|
||||
{"page_id": p["id"], "path": p.get("path", ""), "title": p.get("title", "")}
|
||||
for p in tenant_pages
|
||||
if p.get("id") and p["id"] not in vectorized_page_ids
|
||||
]
|
||||
|
||||
orphaned_chunks = [
|
||||
r for r in wiki_chunk_refs
|
||||
if r.get("page_id") and r["page_id"] not in tenant_page_ids
|
||||
]
|
||||
orphaned_page_ids = sorted({r["page_id"] for r in orphaned_chunks})
|
||||
|
||||
# Collection audit (global listing, read-only)
|
||||
collections = await qdrant.list_collections()
|
||||
unexpected = []
|
||||
foreign_count = 0
|
||||
for coll in collections:
|
||||
category = classify_collection(coll["name"], known_tenants)
|
||||
if category in ("test_residue", "unknown_tenant", "foreign_test_residue"):
|
||||
unexpected.append({"name": coll["name"], "category": category})
|
||||
elif category == "foreign":
|
||||
foreign_count += 1
|
||||
|
||||
# Graph side (tenant labels only)
|
||||
graph_docs = await graph_service.get_all_document_references(user)
|
||||
documents_without_wiki = [
|
||||
{"page_id": d.get("page_id"), "path": d.get("path", ""), "title": d.get("title", "")}
|
||||
for d in graph_docs
|
||||
if d.get("doc_type") == "wiki"
|
||||
and d.get("page_id")
|
||||
and d["page_id"] not in tenant_page_ids
|
||||
]
|
||||
|
||||
duration_ms = (time.time() - start_time) * 1000
|
||||
|
||||
return IntegrityCheckResponse(
|
||||
success=True,
|
||||
user=user,
|
||||
generated_at=datetime.now(timezone.utc).isoformat(),
|
||||
pages_without_vectors=pages_without_vectors,
|
||||
orphaned_vector_chunks=len(orphaned_chunks),
|
||||
orphaned_vector_page_ids=orphaned_page_ids,
|
||||
unexpected_collections=unexpected,
|
||||
foreign_collections=foreign_count,
|
||||
documents_without_wiki=documents_without_wiki,
|
||||
counts={
|
||||
"tenant_wiki_pages": len(tenant_pages),
|
||||
"tenant_vector_chunks": len(wiki_chunk_refs),
|
||||
"tenant_graph_documents": len(graph_docs),
|
||||
"pages_without_vectors": len(pages_without_vectors),
|
||||
"orphaned_vector_chunks": len(orphaned_chunks),
|
||||
"unexpected_collections": len(unexpected),
|
||||
"documents_without_wiki": len(documents_without_wiki),
|
||||
},
|
||||
duration_ms=duration_ms
|
||||
)
|
||||
|
||||
|
||||
@router.post("/integrity-check", response_model=IntegrityCheckResponse)
|
||||
async def integrity_check(
|
||||
request: IntegrityCheckRequest,
|
||||
vector_service: VectorServiceDep = None,
|
||||
graph_service: GraphServiceDep = None,
|
||||
wiki_client: WikiJSDep = None,
|
||||
qdrant: QdrantDep = None,
|
||||
redis: RedisDep = None,
|
||||
api_key: str = Depends(verify_api_key)
|
||||
):
|
||||
"""
|
||||
Nightly integrity check (READ-ONLY: reports, never auto-fixes).
|
||||
|
||||
Reports per tenant:
|
||||
- Wiki pages with ZERO vectors in Qdrant (silent-skip reindex victims)
|
||||
- Orphaned vectors whose wiki page no longer exists
|
||||
- Unexpected Qdrant collections: anything not matching known tenant
|
||||
patterns — flags test-tenant residue and unknown namespaces
|
||||
- Neo4j Document nodes without wiki counterparts
|
||||
- Counts and duration
|
||||
|
||||
The latest report is cached in Redis (30 days) so the weekly quality
|
||||
report can fold it in without re-running the scan.
|
||||
|
||||
**Scheduler Task** — nightly at 04:30, see docs/scheduler-tasks.md.
|
||||
"""
|
||||
try:
|
||||
report = await run_integrity_check(
|
||||
user=request.user,
|
||||
vector_service=vector_service,
|
||||
graph_service=graph_service,
|
||||
wiki_client=wiki_client,
|
||||
qdrant=qdrant
|
||||
)
|
||||
|
||||
# Cache the latest report for the quality report (best-effort)
|
||||
if redis:
|
||||
try:
|
||||
import json as _json
|
||||
await redis.setex(
|
||||
INTEGRITY_LATEST_KEY.format(user=request.user),
|
||||
INTEGRITY_LATEST_TTL,
|
||||
_json.dumps(report.model_dump(mode="json"))
|
||||
)
|
||||
except Exception as e:
|
||||
logger.warning(f"Failed to cache integrity report: {e}")
|
||||
|
||||
logger.info(
|
||||
f"Integrity check for {request.user}: {report.counts} "
|
||||
f"in {report.duration_ms:.0f}ms"
|
||||
)
|
||||
return report
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Integrity check failed for {request.user}: {e}", exc_info=True)
|
||||
raise HTTPException(status_code=500, detail="Integrity check failed")
|
||||
|
||||
Reference in New Issue
Block a user