251 lines
7.3 KiB
Python
251 lines
7.3 KiB
Python
"""
|
|
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
|
|
|
|
|
|
def _unique_dest(dest_dir: Path, name: str) -> Path:
|
|
target = dest_dir / name
|
|
if not target.exists():
|
|
return target
|
|
stem, suffix = Path(name).stem, Path(name).suffix
|
|
i = 2
|
|
while True:
|
|
cand = dest_dir / f"{stem}_{i}{suffix}"
|
|
if not cand.exists():
|
|
return cand
|
|
i += 1
|
|
|
|
|
|
async def _stage_uploads(uploads: list[UploadFile], input_dir: Path) -> tuple[int, str]:
|
|
"""Write one or more uploads into input_dir. Returns (file_count, label)."""
|
|
saved = 0
|
|
total = 0
|
|
labels: list[str] = []
|
|
max_bytes = 120 * 1024 * 1024
|
|
|
|
for up in uploads:
|
|
name = (up.filename or "").strip()
|
|
if not name:
|
|
continue
|
|
base = Path(name.replace("\\", "/")).name
|
|
if not base or base.startswith("."):
|
|
continue
|
|
ext = Path(base).suffix.lower()
|
|
if ext not in ALLOWED_EXT:
|
|
continue
|
|
|
|
data = await up.read()
|
|
if not data:
|
|
continue
|
|
total += len(data)
|
|
if total > max_bytes:
|
|
raise HTTPException(400, "Upload too large (max 120 MB)")
|
|
|
|
if ext == ".zip":
|
|
n = _safe_extract_zip(data, input_dir)
|
|
saved += n
|
|
else:
|
|
_unique_dest(input_dir, base).write_bytes(data)
|
|
saved += 1
|
|
labels.append(base)
|
|
|
|
if saved == 0:
|
|
raise HTTPException(
|
|
400,
|
|
f"No usable files. Allowed: {sorted(ALLOWED_EXT)}",
|
|
)
|
|
label = labels[0] if len(labels) == 1 else f"{len(labels)} files"
|
|
return saved, label
|
|
|
|
|
|
@router.post("/rescan")
|
|
async def rescan(
|
|
background_tasks: BackgroundTasks,
|
|
file: list[UploadFile] = File(..., description="ZIP, one help file, or several files from a folder"),
|
|
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, one .docx/.html/.pdf/.txt/.doc, or several such files
|
|
(folder pick). Process on server, write sections to OUTPUT_DIR and Postgres.
|
|
"""
|
|
_check_token(x_rescan_token or token)
|
|
|
|
staging = Path(tempfile.mkdtemp(prefix="rip_rescan_"))
|
|
input_dir = staging / "input"
|
|
input_dir.mkdir(parents=True, exist_ok=True)
|
|
|
|
try:
|
|
_count, label = await _stage_uploads(file, input_dir)
|
|
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=label,
|
|
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)
|