3d5bc38401
Each file in the inbox can now be accompanied by a `<filename>.meta.json` sidecar that supplies identifier, published_at, effective_at, and source_url. The scanner applies filename-derived defaults first and lets the sidecar override individual fields. - list_inbox / list_failed skip `.meta.json` so sidecars aren't ingested as standalone documents. - mark_processed, mark_failed, and requeue-failed move the sidecar alongside its parent (best-effort). - mark_failed now sanitizes the error string to ASCII before writing it to S3 object metadata (HTTP header encoding). - search SQL selects published_at/effective_at; the formatter shows the full heading_path plus publication date so citations are unambiguous. - scripts/kb_sidecar_template.json documents the expected shape. Refs Task Master #2 Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
185 lines
6.1 KiB
Python
185 lines
6.1 KiB
Python
"""Admin endpoints for managing the Israeli National Insurance KB."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
import os
|
|
from typing import Literal
|
|
|
|
from fastapi import APIRouter, File, Form, HTTPException, Request, UploadFile
|
|
|
|
from api.services.kb import ingest as kb_ingest
|
|
from api.services.kb import s3 as kb_s3
|
|
|
|
logger = logging.getLogger("shira.admin.kb")
|
|
|
|
router = APIRouter(prefix="/admin/kb", tags=["admin-kb"])
|
|
|
|
|
|
def _verify_admin(request: Request) -> None:
|
|
expected = os.environ.get("ADMIN_API_KEY") or os.environ.get("API_KEY", "")
|
|
if not expected:
|
|
raise HTTPException(status_code=503, detail="Admin API key not configured")
|
|
provided = request.headers.get("X-Admin-Key") or request.headers.get("X-Api-Key") or ""
|
|
if provided != expected:
|
|
raise HTTPException(status_code=401, detail="Unauthorized")
|
|
|
|
|
|
@router.post("/ingest")
|
|
async def ingest(
|
|
request: Request,
|
|
kind: Literal["law", "regulation", "circular"] = Form(...),
|
|
title: str = Form(...),
|
|
identifier: str | None = Form(None),
|
|
published_at: str | None = Form(None),
|
|
effective_at: str | None = Form(None),
|
|
source_url: str | None = Form(None),
|
|
original_path: str | None = Form(None),
|
|
file: UploadFile = File(...),
|
|
):
|
|
_verify_admin(request)
|
|
data = await file.read()
|
|
if not data:
|
|
raise HTTPException(status_code=400, detail="Empty file")
|
|
try:
|
|
result = await kb_ingest.ingest_source(
|
|
kind=kind,
|
|
title=title,
|
|
identifier=identifier,
|
|
content=data,
|
|
filename=file.filename or "upload",
|
|
published_at=published_at,
|
|
effective_at=effective_at,
|
|
source_url=source_url,
|
|
original_path=original_path,
|
|
)
|
|
except kb_ingest.IngestError as e:
|
|
raise HTTPException(status_code=422, detail=str(e))
|
|
except Exception as e:
|
|
logger.exception("[admin.kb] ingest failed")
|
|
raise HTTPException(status_code=500, detail=str(e))
|
|
return result
|
|
|
|
|
|
@router.get("/stats")
|
|
async def stats(request: Request):
|
|
_verify_admin(request)
|
|
return await kb_ingest.stats()
|
|
|
|
|
|
@router.get("/inbox")
|
|
async def list_inbox(request: Request):
|
|
_verify_admin(request)
|
|
return {"items": kb_s3.list_inbox()}
|
|
|
|
|
|
def _derive_metadata(filename: str, kind: str) -> dict:
|
|
"""Best-effort metadata from the filename stem. User can override later."""
|
|
stem = filename.rsplit("/", 1)[-1]
|
|
for ext in (".pdf", ".docx", ".txt", ".md"):
|
|
if stem.lower().endswith(ext):
|
|
stem = stem[: -len(ext)]
|
|
break
|
|
return {
|
|
"title": stem,
|
|
"identifier": None,
|
|
"published_at": None,
|
|
"effective_at": None,
|
|
"source_url": None,
|
|
}
|
|
|
|
|
|
def _sidecar_metadata(item_key: str) -> dict:
|
|
"""Fetch the optional `<key>.meta.json` sidecar and return parsed fields.
|
|
|
|
The sidecar is expected to be a JSON object with any of:
|
|
{title, identifier, published_at, effective_at, source_url}
|
|
|
|
Missing or unreadable sidecars return {}.
|
|
"""
|
|
import json
|
|
sidecar_key = item_key + ".meta.json"
|
|
try:
|
|
data = kb_s3.fetch(sidecar_key)
|
|
except Exception:
|
|
return {}
|
|
try:
|
|
parsed = json.loads(data.decode("utf-8"))
|
|
except Exception as e:
|
|
logger.warning("[admin.kb] bad sidecar %s: %s", sidecar_key, e)
|
|
return {}
|
|
# Accept only whitelisted keys — ignore typos / extra fields.
|
|
allowed = {"title", "identifier", "published_at", "effective_at", "source_url"}
|
|
return {k: v for k, v in parsed.items() if k in allowed and v}
|
|
|
|
|
|
@router.post("/requeue-failed")
|
|
async def requeue_failed(request: Request):
|
|
"""Move everything under failed/<kind>/ back to inbox/<kind>/ for retry.
|
|
|
|
The sidecar `<name>.meta.json` (if present) is moved alongside.
|
|
"""
|
|
_verify_admin(request)
|
|
items = kb_s3.list_failed()
|
|
moved = []
|
|
for item in items:
|
|
src = item["key"]
|
|
dst = f"inbox/{item['kind']}/{item['filename']}"
|
|
try:
|
|
kb_s3.move(src, dst)
|
|
# Best-effort sidecar move.
|
|
try:
|
|
kb_s3.move(src + ".meta.json", dst + ".meta.json")
|
|
except Exception:
|
|
pass
|
|
moved.append({"from": src, "to": dst})
|
|
except Exception as e:
|
|
moved.append({"from": src, "error": str(e)})
|
|
return {"requeued": len([m for m in moved if "to" in m]), "items": moved}
|
|
|
|
|
|
@router.post("/scan-inbox")
|
|
async def scan_inbox(request: Request):
|
|
_verify_admin(request)
|
|
items = kb_s3.list_inbox()
|
|
results = []
|
|
for item in items:
|
|
src_key = item["key"]
|
|
kind = item["kind"]
|
|
filename = item["filename"]
|
|
# Merge: filename-derived defaults < sidecar overrides.
|
|
meta = _derive_metadata(filename, kind)
|
|
meta.update(_sidecar_metadata(src_key))
|
|
try:
|
|
data = kb_s3.fetch(src_key)
|
|
result = await kb_ingest.ingest_source(
|
|
kind=kind,
|
|
title=meta["title"],
|
|
identifier=meta.get("identifier"),
|
|
content=data,
|
|
filename=filename,
|
|
published_at=meta.get("published_at"),
|
|
effective_at=meta.get("effective_at"),
|
|
source_url=meta.get("source_url"),
|
|
original_path=f"s3://{kb_s3._bucket()}/processed/{kind}/{filename}",
|
|
)
|
|
dst = kb_s3.mark_processed(src_key, filename, kind)
|
|
results.append({
|
|
"key": src_key,
|
|
"status": "processed",
|
|
"destination": dst,
|
|
**result,
|
|
})
|
|
except Exception as e:
|
|
logger.exception("[scan-inbox] failed for %s", src_key)
|
|
dst = kb_s3.mark_failed(src_key, filename, kind, str(e))
|
|
results.append({
|
|
"key": src_key,
|
|
"status": "failed",
|
|
"destination": dst,
|
|
"error": str(e),
|
|
})
|
|
return {"processed": sum(1 for r in results if r["status"] == "processed"),
|
|
"failed": sum(1 for r in results if r["status"] == "failed"),
|
|
"items": results}
|