Add server-side /api/rescan upload pipeline and Home UI.
Accept ZIP/single help files, process with help_processor into OUTPUT_DIR+Postgres, pollable jobs, optional RESCAN_TOKEN. Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
214
webapp/rescan.py
Normal file
214
webapp/rescan.py
Normal file
@@ -0,0 +1,214 @@
|
||||
"""
|
||||
Server-side rescan: upload help files (zip or single) → process → OUTPUT_DIR + Postgres.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import io
|
||||
import shutil
|
||||
import tempfile
|
||||
import threading
|
||||
import uuid
|
||||
import zipfile
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
from typing import Optional
|
||||
|
||||
from fastapi import APIRouter, BackgroundTasks, File, Form, Header, HTTPException, UploadFile
|
||||
from fastapi.responses import JSONResponse
|
||||
|
||||
router = APIRouter(prefix="/api", tags=["rescan"])
|
||||
|
||||
_jobs: dict[str, dict] = {}
|
||||
_lock = threading.Lock()
|
||||
|
||||
ALLOWED_EXT = {".html", ".htm", ".docx", ".doc", ".txt", ".pdf", ".zip"}
|
||||
|
||||
|
||||
def _cfg():
|
||||
import os
|
||||
|
||||
output = Path(
|
||||
os.getenv(
|
||||
"OUTPUT_DIR",
|
||||
"/mnt/mssql/share/RIP/RIP_Help_Source/Output",
|
||||
)
|
||||
)
|
||||
return {
|
||||
"conn": os.getenv("HELP_DB_CONN", ""),
|
||||
"api_key": os.getenv("ANTHROPIC_API_KEY", ""),
|
||||
"token": os.getenv("RESCAN_TOKEN", ""),
|
||||
"output_dir": output,
|
||||
"remote_root": os.getenv(
|
||||
"HELP_REMOTE_ROOT",
|
||||
str(output).replace("\\", "/"),
|
||||
),
|
||||
}
|
||||
|
||||
|
||||
def _check_token(provided: Optional[str]):
|
||||
expected = _cfg()["token"]
|
||||
if not expected:
|
||||
return
|
||||
if not provided or provided != expected:
|
||||
raise HTTPException(401, "Invalid or missing RESCAN_TOKEN")
|
||||
|
||||
|
||||
def _set_job(job_id: str, **fields):
|
||||
with _lock:
|
||||
job = _jobs.setdefault(job_id, {"id": job_id})
|
||||
job.update(fields)
|
||||
job["updated_at"] = datetime.now(timezone.utc).isoformat()
|
||||
|
||||
|
||||
def _safe_extract_zip(data: bytes, dest: Path) -> int:
|
||||
dest = dest.resolve()
|
||||
count = 0
|
||||
with zipfile.ZipFile(io.BytesIO(data)) as zf:
|
||||
for info in zf.infolist():
|
||||
if info.is_dir():
|
||||
continue
|
||||
name = Path(info.filename)
|
||||
# skip macOS junk / absolute paths
|
||||
if any(p.startswith("__MACOSX") or p.startswith(".") for p in name.parts):
|
||||
continue
|
||||
target = (dest / name.name).resolve()
|
||||
if not str(target).startswith(str(dest)):
|
||||
raise ValueError(f"unsafe zip path: {info.filename}")
|
||||
target.parent.mkdir(parents=True, exist_ok=True)
|
||||
with zf.open(info) as src, open(target, "wb") as out:
|
||||
shutil.copyfileobj(src, out)
|
||||
count += 1
|
||||
return count
|
||||
|
||||
|
||||
def _run_rescan(job_id: str, staging: Path, prefix: str, force: bool):
|
||||
cfg = _cfg()
|
||||
try:
|
||||
if not cfg["conn"]:
|
||||
raise RuntimeError("HELP_DB_CONN is not set on the server")
|
||||
if not cfg["api_key"]:
|
||||
raise RuntimeError("ANTHROPIC_API_KEY is not set on the server")
|
||||
|
||||
out = cfg["output_dir"]
|
||||
out.mkdir(parents=True, exist_ok=True)
|
||||
if not os_access_writable(out):
|
||||
raise RuntimeError(
|
||||
f"OUTPUT_DIR is not writable: {out} "
|
||||
"(Coolify Persistent Storage must be read-write)"
|
||||
)
|
||||
|
||||
_set_job(job_id, status="processing", message="running help_processor")
|
||||
|
||||
# Import here so webapp boots even if processor deps fail later
|
||||
from help_processor import process_directory
|
||||
|
||||
process_directory(
|
||||
input_dir=staging,
|
||||
output_dir=out,
|
||||
conn_str=cfg["conn"],
|
||||
api_key=cfg["api_key"],
|
||||
prefix=prefix,
|
||||
force=force,
|
||||
purge_missing=False,
|
||||
remote_root=cfg["remote_root"],
|
||||
)
|
||||
_set_job(
|
||||
job_id,
|
||||
status="done",
|
||||
message="ok",
|
||||
output_dir=str(out),
|
||||
prefix=prefix,
|
||||
)
|
||||
except Exception as e:
|
||||
_set_job(job_id, status="error", message=str(e))
|
||||
finally:
|
||||
shutil.rmtree(staging, ignore_errors=True)
|
||||
|
||||
|
||||
def os_access_writable(path: Path) -> bool:
|
||||
try:
|
||||
path.mkdir(parents=True, exist_ok=True)
|
||||
probe = path / f".write_probe_{uuid.uuid4().hex}"
|
||||
probe.write_text("ok", encoding="utf-8")
|
||||
probe.unlink()
|
||||
return True
|
||||
except Exception:
|
||||
return False
|
||||
|
||||
|
||||
@router.post("/rescan")
|
||||
async def rescan(
|
||||
background_tasks: BackgroundTasks,
|
||||
file: UploadFile = File(..., description="ZIP or single help file"),
|
||||
prefix: str = Form("RIP"),
|
||||
force: str = Form("false"),
|
||||
x_rescan_token: Optional[str] = Header(None, alias="X-Rescan-Token"),
|
||||
token: Optional[str] = Form(None),
|
||||
):
|
||||
"""
|
||||
Upload a .zip (or one .docx/.html/.pdf/.txt/.doc), process on server,
|
||||
write sections to OUTPUT_DIR and Postgres.
|
||||
"""
|
||||
_check_token(x_rescan_token or token)
|
||||
|
||||
name = (file.filename or "upload.bin").strip()
|
||||
ext = Path(name).suffix.lower()
|
||||
if ext not in ALLOWED_EXT:
|
||||
raise HTTPException(
|
||||
400,
|
||||
f"Unsupported type {ext}. Allowed: {sorted(ALLOWED_EXT)}",
|
||||
)
|
||||
|
||||
data = await file.read()
|
||||
if not data:
|
||||
raise HTTPException(400, "Empty upload")
|
||||
if len(data) > 120 * 1024 * 1024:
|
||||
raise HTTPException(400, "Upload too large (max 120 MB)")
|
||||
|
||||
staging = Path(tempfile.mkdtemp(prefix="rip_rescan_"))
|
||||
input_dir = staging / "input"
|
||||
input_dir.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
try:
|
||||
if ext == ".zip":
|
||||
n = _safe_extract_zip(data, input_dir)
|
||||
if n == 0:
|
||||
raise HTTPException(400, "ZIP has no usable files")
|
||||
else:
|
||||
(input_dir / Path(name).name).write_bytes(data)
|
||||
except HTTPException:
|
||||
shutil.rmtree(staging, ignore_errors=True)
|
||||
raise
|
||||
except Exception as e:
|
||||
shutil.rmtree(staging, ignore_errors=True)
|
||||
raise HTTPException(400, f"Failed to unpack upload: {e}") from e
|
||||
|
||||
job_id = uuid.uuid4().hex[:12]
|
||||
force_flag = str(force).lower() in ("1", "true", "yes", "on")
|
||||
_set_job(
|
||||
job_id,
|
||||
status="queued",
|
||||
message="accepted",
|
||||
filename=name,
|
||||
prefix=prefix,
|
||||
force=force_flag,
|
||||
created_at=datetime.now(timezone.utc).isoformat(),
|
||||
)
|
||||
background_tasks.add_task(_run_rescan, job_id, staging, prefix, force_flag)
|
||||
return JSONResponse(
|
||||
{
|
||||
"ok": True,
|
||||
"job_id": job_id,
|
||||
"status": "queued",
|
||||
"poll": f"/api/rescan/{job_id}",
|
||||
}
|
||||
)
|
||||
|
||||
|
||||
@router.get("/rescan/{job_id}")
|
||||
def rescan_status(job_id: str):
|
||||
with _lock:
|
||||
job = _jobs.get(job_id)
|
||||
if not job:
|
||||
raise HTTPException(404, "job not found")
|
||||
return JSONResponse(job)
|
||||
Reference in New Issue
Block a user