feat(kb): requeue-failed endpoint
Moves every object under failed/<kind>/ back to inbox/<kind>/ for retry after a transient upstream failure (e.g. Voyage 401 before a key rotation). Refs Task Master #2 Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -84,6 +84,23 @@ def _derive_metadata(filename: str, kind: str) -> dict:
|
|||||||
return {"title": stem, "identifier": None}
|
return {"title": stem, "identifier": None}
|
||||||
|
|
||||||
|
|
||||||
|
@router.post("/requeue-failed")
|
||||||
|
async def requeue_failed(request: Request):
|
||||||
|
"""Move everything under failed/<kind>/ back to inbox/<kind>/ for retry."""
|
||||||
|
_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)
|
||||||
|
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")
|
@router.post("/scan-inbox")
|
||||||
async def scan_inbox(request: Request):
|
async def scan_inbox(request: Request):
|
||||||
_verify_admin(request)
|
_verify_admin(request)
|
||||||
|
|||||||
@@ -95,6 +95,25 @@ def mark_processed(src_key: str, filename: str, kind: str) -> str:
|
|||||||
return dst
|
return dst
|
||||||
|
|
||||||
|
|
||||||
|
def list_failed() -> list[dict]:
|
||||||
|
s3 = _client()
|
||||||
|
bucket = _bucket()
|
||||||
|
items: list[dict] = []
|
||||||
|
for kind in _KINDS:
|
||||||
|
prefix = f"failed/{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"]
|
||||||
|
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 mark_failed(src_key: str, filename: str, kind: str, reason: str) -> str:
|
def mark_failed(src_key: str, filename: str, kind: str, reason: str) -> str:
|
||||||
s3 = _client()
|
s3 = _client()
|
||||||
bucket = _bucket()
|
bucket = _bucket()
|
||||||
|
|||||||
Reference in New Issue
Block a user