From 5ddacbe57387681ad4498ecb8decac15316671e2 Mon Sep 17 00:00:00 2001 From: Chaim Date: Tue, 21 Apr 2026 14:45:28 +0000 Subject: [PATCH] =?UTF-8?q?feat(kb):=20MinIO=20inbox=20scanner=20=E2=80=94?= =?UTF-8?q?=20list/scan/move=20endpoints?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Adds a MinIO-backed ingest flow: upload to s3://insurance-kb/inbox// via MinIO Console or any S3 client, then POST /admin/kb/scan-inbox to process all pending files. Successful ingests move to processed//, failures move to failed// 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) --- api/routes/admin_kb.py | 59 +++++++++++++++++++++ api/services/kb/s3.py | 115 +++++++++++++++++++++++++++++++++++++++++ 2 files changed, 174 insertions(+) create mode 100644 api/services/kb/s3.py diff --git a/api/routes/admin_kb.py b/api/routes/admin_kb.py index 54724c3..d370746 100644 --- a/api/routes/admin_kb.py +++ b/api/routes/admin_kb.py @@ -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} diff --git a/api/services/kb/s3.py b/api/services/kb/s3.py new file mode 100644 index 0000000..0b7dc5f --- /dev/null +++ b/api/services/kb/s3.py @@ -0,0 +1,115 @@ +"""MinIO/S3 helpers for the insurance-kb bucket. + +Layout in the bucket: + inbox/{law,regulation,circular}/ + processed/{law,regulation,circular}/ + failed/{law,regulation,circular}/ + +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