feat(kb): MinIO inbox scanner — list/scan/move endpoints
Adds a MinIO-backed ingest flow: upload to s3://insurance-kb/inbox/<kind>/ via MinIO Console or any S3 client, then POST /admin/kb/scan-inbox to process all pending files. Successful ingests move to processed/<kind>/, failures move to failed/<kind>/ with the error stored as S3 metadata. - api/services/kb/s3.py: boto3 client, list_inbox/fetch/move helpers. - api/routes/admin_kb.py: GET /admin/kb/inbox, POST /admin/kb/scan-inbox. - The kind is derived from the folder path; title/identifier default to the filename stem and can be refined by re-ingesting with metadata. Refs Task Master #2 Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -9,6 +9,7 @@ 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")
|
||||
|
||||
@@ -64,3 +65,61 @@ async def ingest(
|
||||
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
|
||||
# Title is just the stem; identifier left empty unless obvious.
|
||||
return {"title": stem, "identifier": None}
|
||||
|
||||
|
||||
@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"]
|
||||
meta = _derive_metadata(filename, kind)
|
||||
try:
|
||||
data = kb_s3.fetch(src_key)
|
||||
result = await kb_ingest.ingest_source(
|
||||
kind=kind,
|
||||
title=meta["title"],
|
||||
identifier=meta["identifier"],
|
||||
content=data,
|
||||
filename=filename,
|
||||
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}
|
||||
|
||||
@@ -0,0 +1,115 @@
|
||||
"""MinIO/S3 helpers for the insurance-kb bucket.
|
||||
|
||||
Layout in the bucket:
|
||||
inbox/{law,regulation,circular}/<file>
|
||||
processed/{law,regulation,circular}/<file>
|
||||
failed/{law,regulation,circular}/<file>
|
||||
|
||||
After a successful ingest, the file is moved from inbox/* to processed/*.
|
||||
On failure it is moved to failed/*.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import os
|
||||
from typing import Iterable
|
||||
|
||||
import boto3
|
||||
from botocore.client import Config as BotoConfig
|
||||
|
||||
logger = logging.getLogger("shira.kb.s3")
|
||||
|
||||
_KINDS = ("law", "regulation", "circular")
|
||||
|
||||
|
||||
def _client():
|
||||
endpoint = os.environ.get("KB_S3_ENDPOINT")
|
||||
access = os.environ.get("KB_S3_ACCESS_KEY")
|
||||
secret = os.environ.get("KB_S3_SECRET_KEY")
|
||||
region = os.environ.get("KB_S3_REGION", "us-east-1")
|
||||
if not (endpoint and access and secret):
|
||||
raise RuntimeError("KB_S3_* environment variables are not set")
|
||||
return boto3.client(
|
||||
"s3",
|
||||
endpoint_url=endpoint,
|
||||
aws_access_key_id=access,
|
||||
aws_secret_access_key=secret,
|
||||
region_name=region,
|
||||
config=BotoConfig(signature_version="s3v4"),
|
||||
)
|
||||
|
||||
|
||||
def _bucket() -> str:
|
||||
return os.environ.get("KB_S3_BUCKET", "insurance-kb")
|
||||
|
||||
|
||||
def list_inbox() -> list[dict]:
|
||||
"""Return inbox objects across all kinds, excluding folder markers."""
|
||||
s3 = _client()
|
||||
bucket = _bucket()
|
||||
items: list[dict] = []
|
||||
for kind in _KINDS:
|
||||
prefix = f"inbox/{kind}/"
|
||||
paginator = s3.get_paginator("list_objects_v2")
|
||||
for page in paginator.paginate(Bucket=bucket, Prefix=prefix):
|
||||
for obj in page.get("Contents", []) or []:
|
||||
key = obj["Key"]
|
||||
# Skip folder markers (MinIO creates zero-byte objects with trailing /)
|
||||
if key.endswith("/"):
|
||||
continue
|
||||
filename = key[len(prefix):]
|
||||
if not filename:
|
||||
continue
|
||||
items.append({
|
||||
"key": key,
|
||||
"kind": kind,
|
||||
"filename": filename,
|
||||
"size": obj["Size"],
|
||||
})
|
||||
return items
|
||||
|
||||
|
||||
def fetch(key: str) -> bytes:
|
||||
s3 = _client()
|
||||
resp = s3.get_object(Bucket=_bucket(), Key=key)
|
||||
return resp["Body"].read()
|
||||
|
||||
|
||||
def move(src_key: str, dst_key: str) -> None:
|
||||
s3 = _client()
|
||||
bucket = _bucket()
|
||||
s3.copy_object(
|
||||
Bucket=bucket,
|
||||
Key=dst_key,
|
||||
CopySource={"Bucket": bucket, "Key": src_key},
|
||||
MetadataDirective="COPY",
|
||||
)
|
||||
s3.delete_object(Bucket=bucket, Key=src_key)
|
||||
logger.info("[kb.s3] moved %s -> %s", src_key, dst_key)
|
||||
|
||||
|
||||
def mark_processed(src_key: str, filename: str, kind: str) -> str:
|
||||
dst = f"processed/{kind}/{filename}"
|
||||
move(src_key, dst)
|
||||
return dst
|
||||
|
||||
|
||||
def mark_failed(src_key: str, filename: str, kind: str, reason: str) -> str:
|
||||
s3 = _client()
|
||||
bucket = _bucket()
|
||||
dst = f"failed/{kind}/{filename}"
|
||||
s3.copy_object(
|
||||
Bucket=bucket,
|
||||
Key=dst,
|
||||
CopySource={"Bucket": bucket, "Key": src_key},
|
||||
Metadata={"ingest_error": reason[:800]},
|
||||
MetadataDirective="REPLACE",
|
||||
)
|
||||
s3.delete_object(Bucket=bucket, Key=src_key)
|
||||
logger.info("[kb.s3] failed %s -> %s (%s)", src_key, dst, reason[:80])
|
||||
return dst
|
||||
|
||||
|
||||
def kinds() -> Iterable[str]:
|
||||
return _KINDS
|
||||
Reference in New Issue
Block a user