diff --git a/.gitignore b/.gitignore index 14b8968..08dfb06 100644 --- a/.gitignore +++ b/.gitignore @@ -58,3 +58,8 @@ orchestration/.dagster_home/.telemetry/ # would put customer files in the repository permanently. # BATCH_UPLOAD_DIR overrides the location; this covers the default. data/batch_uploads/ + +# Spreadsheets a colleague dropped through POST /api/uploads/catalog, waiting +# for review. Same reasoning as above, and more so: these arrive from outside +# the team and nobody here has looked at them yet. +data/inbox/ diff --git a/app/api/routers/batch_catalog.py b/app/api/routers/batch_catalog.py index ec20eec..dac52fb 100644 --- a/app/api/routers/batch_catalog.py +++ b/app/api/routers/batch_catalog.py @@ -33,7 +33,7 @@ from pydantic import BaseModel from app.api.batch_job_store import batch_job_store from app.api.deps import require_admin -from app.core import batch_ingest, batch_worker +from app.core import batch_ingest, batch_worker, inbox from app.core import store_catalog_pipeline as pipeline from app.infrastructure.settings import ( BATCH_MAX_FILES, @@ -77,6 +77,7 @@ class BatchOut(BaseModel): batch_id: str status: str detail: Optional[str] = None + submitted_by: Optional[str] = None created_at: float updated_at: float files_total: int @@ -308,6 +309,139 @@ async def ingest_catalog_batch( return _to_out(manifest) +# --------------------------------------------------------------------------- +# The review inbox: files a colleague dropped, waiting for a decision +# --------------------------------------------------------------------------- +# These are the admin half of the two-actor flow. The uploader half lives in +# app/api/routers/uploads.py and can reach none of this. + + +class InboxSelection(BaseModel): + file_ids: List[str] + use_llm: bool = False + fetch_images: bool = False + + +class InboxDismissal(BaseModel): + file_ids: List[str] + + +@router.get("/inbox", dependencies=[Depends(require_admin)]) +def list_inbox() -> dict: + """Files awaiting review, grouped by the drop they arrived in. + + `pending_count` is what the badge renders, and it is computed here rather + than by summing the response client-side so the two can never disagree. + """ + submissions = inbox.list_pending() + return { + "pending_count": sum(len(s.pending_files) for s in submissions), + "submissions": [ + { + "submission_id": s.submission_id, + "submitted_by": s.submitted_by, + "created_at": s.created_at, + "files": [ + { + "file_id": f.file_id, + "filename": f.filename, + "rows_total": f.rows_total, + "size_bytes": f.size_bytes, + } + for f in s.pending_files + ], + } + for s in submissions + ], + } + + +@router.post("/from-inbox", status_code=status.HTTP_202_ACCEPTED, + dependencies=[Depends(require_admin)]) +def start_batch_from_inbox(selection: InboxSelection) -> BatchOut: + """Compose a batch out of the selected inbox files and start it. + + The files may come from different drops on different days; that is the + point of selecting per file rather than per submission. From here on this + is an ordinary batch and every existing path - progress, resume, cancel - + applies unchanged. + """ + if not selection.file_ids: + raise HTTPException(status_code=400, detail="No files were selected.") + + try: + uploads, submitters = inbox.collect_for_batch(selection.file_ids) + except KeyError as exc: + raise HTTPException( + status_code=404, detail=f"No such file in the inbox: {exc.args[0]}" + ) from exc + except ValueError as exc: + # Two admin tabs open on the same inbox. Tell the second one what + # happened rather than silently running the file a second time. + raise HTTPException(status_code=409, detail=str(exc)) from exc + + # Re-parse rather than trusting the row counts recorded at upload: the + # ceilings are a property of the batch about to run, not of the drops it + # was assembled from, and a selection can span any number of drops. + read = [(name, contents) for name, contents in uploads] + valid, invalid, _rows = _parse_all(read) + if not valid: + detail = "; ".join(f"{name}: {reason}" for name, reason in invalid) + raise HTTPException( + status_code=400, + detail=f"None of the selected files could be ingested. {detail}", + ) + + manifest = batch_ingest.stage_batch( + [(name, contents) for name, contents, _n in valid], + use_llm=selection.use_llm, + fetch_images=selection.fetch_images, + invalid=invalid, + ) + for entry, (_name, _contents, rows) in zip(manifest.files, valid): + entry.rows_total = rows + manifest.submitted_by = ", ".join(submitters) or None + batch_ingest.write_manifest(manifest) + batch_job_store.put(manifest) + + try: + batch_worker.submit(manifest.batch_id) + except queue.Full: + manifest.detail = ( + "The ingestion queue is full. This batch is staged and can be started " + "with Resume once the running batches finish." + ) + batch_ingest.write_manifest(manifest) + batch_job_store.put(manifest) + raise HTTPException( + status_code=429, + detail=( + "Too many batches are already queued. This selection has been staged - " + "press Resume on it once the current batch finishes." + ), + ) + + # Only now, once the batch exists AND is queued. Marking first would drop + # the files out of the inbox with nothing left to retry from if staging + # had then failed. + inbox.mark_consumed(selection.file_ids, manifest.batch_id) + return _to_out(manifest) + + +@router.post("/inbox/dismiss", dependencies=[Depends(require_admin)]) +def dismiss_inbox_files(dismissal: InboxDismissal) -> dict: + """Mark files as never-to-run, so the badge can reach zero.""" + if not dismissal.file_ids: + raise HTTPException(status_code=400, detail="No files were selected.") + changed = inbox.dismiss(dismissal.file_ids) + if not changed: + raise HTTPException( + status_code=409, + detail="None of those files were still awaiting review.", + ) + return {"dismissed": changed, "pending_count": inbox.pending_count()} + + @router.get("/batches", dependencies=[Depends(require_admin)]) def list_catalog_batches(limit: int = 20) -> dict: limit = max(1, min(limit, 100)) diff --git a/app/api/routers/uploads.py b/app/api/routers/uploads.py new file mode 100644 index 0000000..7a79ef3 --- /dev/null +++ b/app/api/routers/uploads.py @@ -0,0 +1,212 @@ +"""The one endpoint an outside contributor may call. + + POST /api/uploads/catalog - drop spreadsheets into the review inbox + +This is deliberately its own router, with its own prefix and its own guard, so +that the difference between it and everything else in the app is visible in one +screen rather than inferred from a decorator halfway down a 400-line admin +module. Everything else that touches catalog data is `require_admin`; this is +`require_permission("upload_catalog")`, and that single line is the whole +security boundary of the feature. + +WHAT THIS ENDPOINT CANNOT DO +---------------------------- +Start work. Files land in the inbox and wait for an admin to select them +(`app/core/inbox.py`). No worker is touched, no queue is entered, no thread is +started. That is what makes it safe to hand a credential to someone outside the +team: the worst a leaked uploader key costs is bounded disk, never CPU on a +one-vCPU host that is also serving the API. + +It also cannot read anything. There is no GET here on purpose - the decision was +that this is a one-way drop, so the credential grants no visibility into the +catalog, into other submissions, or even into the submitter's own past uploads. + +WHY A BAD SHEET IS REJECTED HERE AND NOT LATER +---------------------------------------------- +The file is parsed during the request, while the colleague is still watching. +Telling them "row 1 has no product name column" in the 202 is worth far more +than discovering it days later in an admin panel they cannot see, with no way to +ask them for a corrected file except out of band. +""" +from __future__ import annotations + +import logging +from typing import List, Optional + +from fastapi import APIRouter, Depends, File, HTTPException, UploadFile, status +from pydantic import BaseModel + +from app.api.deps import require_permission +from app.core import inbox +from app.core import store_catalog_pipeline as pipeline +from app.infrastructure.security import Principal +from app.infrastructure.settings import ( + BATCH_MAX_FILES, + BATCH_MAX_TOTAL_BYTES, + BATCH_MAX_TOTAL_ROWS, + INBOX_MAX_PENDING_FILES, +) + +logger = logging.getLogger(__name__) +router = APIRouter(prefix="/uploads", tags=["uploads"]) + +# Identical to the admin batch path. A file that is too large for an admin to +# upload is not somehow acceptable because a colleague sent it. +MAX_UPLOAD_BYTES = 10 * 1024 * 1024 +MAX_UPLOAD_ROWS = 2000 + + +class SubmittedFileOut(BaseModel): + filename: str + accepted: bool + rows_total: int = 0 + error: Optional[str] = None + + +class SubmissionOut(BaseModel): + submission_id: Optional[str] + submitted_by: str + files_accepted: int + files_rejected: int + files: List[SubmittedFileOut] + message: str + + +@router.post("/catalog", status_code=status.HTTP_202_ACCEPTED) +async def submit_catalog_files( + files: List[UploadFile] = File(...), + principal: Principal = Depends(require_permission("upload_catalog")), +) -> SubmissionOut: + """Accept spreadsheets into the review inbox. Starts nothing. + + Returns 202 with a per-file verdict. A drop where some files parse and some + do not is a partial success, not a failure: the good ones are kept and the + caller is told precisely which sheet to fix. + """ + if not files: + raise HTTPException(status_code=400, detail="No files were uploaded.") + if len(files) > BATCH_MAX_FILES: + raise HTTPException( + status_code=413, + detail=( + f"{len(files)} files exceeds the {BATCH_MAX_FILES}-file limit for one " + f"upload. Send them in smaller drops." + ), + ) + + # Refuse before reading a byte if the inbox is already backed up. This is + # the ceiling that stops an unattended key filling the disk one perfectly + # valid file at a time. + already_waiting = inbox.pending_count() + if already_waiting + len(files) > INBOX_MAX_PENDING_FILES: + raise HTTPException( + status_code=429, + detail=( + f"The review inbox already holds {already_waiting} file(s) awaiting " + f"review, and the limit is {INBOX_MAX_PENDING_FILES}. Please wait until " + f"some have been processed." + ), + ) + + accepted: list = [] + rejected: list = [] + total_bytes = 0 + total_rows = 0 + + for upload in files: + name = upload.filename or "upload.xlsx" + contents = await upload.read() + + if not contents: + rejected.append((name, "The file is empty.")) + continue + if len(contents) > MAX_UPLOAD_BYTES: + rejected.append(( + name, + f"Larger than the {MAX_UPLOAD_BYTES // (1024 * 1024)}MB per-file limit.", + )) + continue + + total_bytes += len(contents) + if total_bytes > BATCH_MAX_TOTAL_BYTES: + raise HTTPException( + status_code=413, + detail=( + f"This drop is larger than the " + f"{BATCH_MAX_TOTAL_BYTES // (1024 * 1024)}MB total limit." + ), + ) + + # Parse now, while the sender is still here to be told. + try: + frame, _mapping = pipeline.parse_spreadsheet(name, contents) + except HTTPException as exc: + # read_products_dataframe raises this for an unsupported extension + # or a missing Excel reader; its message already names the file and + # says what to do. + rejected.append((name, str(exc.detail))) + continue + except Exception as exc: # noqa: BLE001 - an unreadable sheet is caller error + rejected.append((name, f"Could not parse the file: {exc}")) + continue + + if frame.empty: + rejected.append((name, "The file has no data rows.")) + continue + rows = int(len(frame)) + if rows > MAX_UPLOAD_ROWS: + rejected.append(( + name, f"{rows} rows exceeds the {MAX_UPLOAD_ROWS}-row per-file limit." + )) + continue + + total_rows += rows + if total_rows > BATCH_MAX_TOTAL_ROWS: + raise HTTPException( + status_code=413, + detail=( + f"This drop totals more than {BATCH_MAX_TOTAL_ROWS} rows. " + f"Send it in two smaller drops." + ), + ) + + accepted.append((name, contents, rows)) + + submission = None + if accepted: + submission = inbox.stage_submission( + accepted, + # The key's NAME, never its secret. Principal.username is the name + # half of the API_KEYS entry (see principal_for_api_key). + submitted_by=principal.username, + ) + logger.info( + "Inbox: %d file(s) from %s awaiting review (submission %s)", + len(accepted), principal.username, submission.submission_id, + ) + + out = [ + SubmittedFileOut(filename=n, accepted=True, rows_total=r) + for n, _c, r in accepted + ] + [ + SubmittedFileOut(filename=n, accepted=False, error=e) for n, e in rejected + ] + + if accepted and rejected: + message = ( + f"{len(accepted)} file(s) received and awaiting review. " + f"{len(rejected)} could not be read - see the errors below and resend those." + ) + elif accepted: + message = f"{len(accepted)} file(s) received and awaiting review." + else: + message = "No file could be read. Nothing was received - see the errors below." + + return SubmissionOut( + submission_id=submission.submission_id if submission else None, + submitted_by=principal.username, + files_accepted=len(accepted), + files_rejected=len(rejected), + files=out, + message=message, + ) diff --git a/app/core/batch_ingest.py b/app/core/batch_ingest.py index ca02f52..89e5401 100644 --- a/app/core/batch_ingest.py +++ b/app/core/batch_ingest.py @@ -124,6 +124,10 @@ class BatchManifest: created_at: float = field(default_factory=time.time) updated_at: float = field(default_factory=time.time) detail: Optional[str] = None + # Who sent the files, when they came through the review inbox. None for a + # batch an admin uploaded directly. Carried so the Batch tab can say where + # a run came from instead of leaving it to be guessed from filenames. + submitted_by: Optional[str] = None files: List[BatchFile] = field(default_factory=list) # -- derived, recomputed rather than stored, so they cannot drift --------- @@ -188,6 +192,7 @@ class BatchManifest: "created_at": self.created_at, "updated_at": self.updated_at, "detail": self.detail, + "submitted_by": self.submitted_by, "files_total": self.files_total, "files_done": self.files_done, "files_failed": self.files_failed, @@ -212,6 +217,7 @@ class BatchManifest: created_at=float(raw.get("created_at") or time.time()), updated_at=float(raw.get("updated_at") or time.time()), detail=raw.get("detail"), + submitted_by=raw.get("submitted_by"), files=files, ) diff --git a/app/core/inbox.py b/app/core/inbox.py new file mode 100644 index 0000000..e6e2987 --- /dev/null +++ b/app/core/inbox.py @@ -0,0 +1,348 @@ +""" +A review inbox: a third party drops spreadsheets, an admin decides what runs. + +WHY THIS IS NOT PART OF `batch_ingest` +-------------------------------------- +A batch is a RUN. It has a worker, a state machine +(`queued -> running -> done | failed | partial | cancelled | interrupted`), a +resume path, and thirty-five tests built around exactly that. A submission is +not a run and never becomes one: it is a pile of files sitting still, with no +worker and no run state, waiting for a human. + +Folding "awaiting review" into `BatchManifest` would have put review logic +inside the thing that executes pipelines, and every one of those tests would +have needed re-reasoning to answer "can this state reach the worker?". Two +small state machines are easier to be sure about than one that means two +things. + +So the flow is: + + colleague admin batch_ingest + --------- ----- ------------ + POST /api/uploads/catalog + | + v + SUBMISSION (inert) ---> select file ids ---> stage_batch() ---> worker + +`claim_files()` is the seam. It copies the chosen files into a fresh batch +directory and hands back what `stage_batch` needs, so everything downstream is +the machinery that already exists and is already tested. + +WHY COPY RATHER THAN MOVE +------------------------- +A batch directory has to be self-contained: that is what lets a batch be +resumed after a restart without caring whether anything else still exists. If a +batch referenced files living in the inbox, dismissing or purging an inbox entry +would break the resume of a run that had already started. The cost is a +duplicate of a file that is at most 10MB, until retention clears the inbox copy. + +NOTHING HERE EXECUTES ANYTHING +------------------------------ +No import of `batch_worker`, no thread, no queue. That is the security property +the uploader role depends on: a credential that can reach only this module can +consume bounded disk and can never consume CPU. +""" +from __future__ import annotations + +import json +import logging +import os +import re +import shutil +import time +import uuid +from dataclasses import asdict, dataclass, field +from pathlib import Path +from typing import Any, Dict, List, Optional, Tuple + +from app.core.batch_ingest import _safe_name +from app.infrastructure.settings import ( + INBOX_RETENTION_DAYS, + INBOX_UPLOAD_DIR, +) + +logger = logging.getLogger(__name__) + +RECORD_NAME = "submission.json" + +# File states inside the inbox. Deliberately disjoint from the batch states - +# an inbox file is never "running"; the COPY of it inside a batch is. +PENDING = "pending" # waiting for an admin to look at it +CONSUMED = "consumed" # copied into a batch and started +DISMISSED = "dismissed" # an admin decided it will never run + +CLOSED_STATES = {CONSUMED, DISMISSED} + + +@dataclass +class InboxFile: + """One spreadsheet in a drop.""" + + file_id: str + filename: str # what the colleague called it + stored_name: str # what it is called on disk + size_bytes: int = 0 + rows_total: int = 0 + status: str = PENDING + detail: Optional[str] = None + batch_id: Optional[str] = None # set when consumed, so the trail is followable + decided_at: Optional[float] = None + + +@dataclass +class Submission: + """One POST from one colleague. Serialised to submission.json verbatim.""" + + submission_id: str + submitted_by: str # the API key's NAME, never its secret + created_at: float = field(default_factory=time.time) + updated_at: float = field(default_factory=time.time) + files: List[InboxFile] = field(default_factory=list) + + @property + def pending_files(self) -> List[InboxFile]: + return [f for f in self.files if f.status == PENDING] + + def to_dict(self) -> Dict[str, Any]: + return { + "submission_id": self.submission_id, + "submitted_by": self.submitted_by, + "created_at": self.created_at, + "updated_at": self.updated_at, + "files": [asdict(f) for f in self.files], + } + + @classmethod + def from_dict(cls, raw: Dict[str, Any]) -> "Submission": + allowed = set(InboxFile.__dataclass_fields__) + return cls( + submission_id=raw["submission_id"], + submitted_by=raw.get("submitted_by") or "unknown", + created_at=float(raw.get("created_at") or time.time()), + updated_at=float(raw.get("updated_at") or time.time()), + files=[ + InboxFile(**{k: v for k, v in entry.items() if k in allowed}) + for entry in (raw.get("files") or []) + ], + ) + + +# --------------------------------------------------------------------------- +# Disk layout +# --------------------------------------------------------------------------- +def inbox_root() -> Path: + """Read at call time, not import time, so tests can repoint the directory.""" + return Path(INBOX_UPLOAD_DIR) + + +def submission_dir(submission_id: str) -> Path: + # Ids are generated here (uuid4 hex), never taken from a request body, but + # this is the function that turns one into a path - so it validates anyway. + if not re.fullmatch(r"[A-Za-z0-9_-]{1,64}", submission_id or ""): + raise ValueError("Invalid submission id: {!r}".format(submission_id)) + return inbox_root() / submission_id + + +def record_path(submission_id: str) -> Path: + return submission_dir(submission_id) / RECORD_NAME + + +def write_record(submission: Submission) -> None: + """Temp file then os.replace, for the same reason batch_ingest does it: + a half-written record is indistinguishable from a corrupt one, and the + listing path reads every record it finds.""" + target = record_path(submission.submission_id) + target.parent.mkdir(parents=True, exist_ok=True) + tmp = target.with_name(RECORD_NAME + ".tmp") + tmp.write_text(json.dumps(submission.to_dict(), indent=2), encoding="utf-8") + os.replace(tmp, target) + + +def read_record(submission_id: str) -> Optional[Submission]: + try: + path = record_path(submission_id) + except ValueError: + return None + if not path.exists(): + return None + try: + return Submission.from_dict(json.loads(path.read_text(encoding="utf-8"))) + except Exception as exc: # noqa: BLE001 - one bad record must not break the inbox + logger.warning("Ignoring unreadable submission record %s: %s", path, exc) + return None + + +def list_submissions() -> List[Submission]: + """Every readable submission on disk, newest first.""" + root = inbox_root() + if not root.exists(): + return [] + found = [] + for child in sorted(root.iterdir()): + if not child.is_dir(): + continue + record = read_record(child.name) + if record: + found.append(record) + return sorted(found, key=lambda s: s.created_at, reverse=True) + + +def list_pending() -> List[Submission]: + """Submissions that still have at least one file awaiting a decision.""" + return [s for s in list_submissions() if s.pending_files] + + +def pending_count() -> int: + """What the badge shows.""" + return sum(len(s.pending_files) for s in list_submissions()) + + +# --------------------------------------------------------------------------- +# Intake +# --------------------------------------------------------------------------- +def stage_submission( + uploads: List[Tuple[str, bytes, int]], + *, + submitted_by: str, + rejected: Optional[List[Tuple[str, str]]] = None, +) -> Submission: + """Write a colleague's drop to disk. Starts nothing. + + `uploads` is (filename, content, rows_total) for files that already parsed + cleanly. `rejected` carries the ones that did not; they are recorded so the + colleague's 202 can tell them which sheet to fix, but they are NOT written + to disk and never appear in the inbox - there is nothing an admin could + usefully do with a file that cannot be read. + """ + submission_id = uuid.uuid4().hex + directory = submission_dir(submission_id) + directory.mkdir(parents=True, exist_ok=True) + + submission = Submission(submission_id=submission_id, submitted_by=submitted_by) + + for position, (filename, content, rows) in enumerate(uploads): + stored = "{:02d}_{}".format(position, _safe_name(filename)) + (directory / stored).write_bytes(content) + submission.files.append( + InboxFile( + file_id=uuid.uuid4().hex, + filename=filename or stored, + stored_name=stored, + size_bytes=len(content), + rows_total=rows, + ) + ) + + write_record(submission) + return submission + + +# --------------------------------------------------------------------------- +# The seam into batch_ingest +# --------------------------------------------------------------------------- +def find_file(file_id: str) -> Optional[Tuple[Submission, InboxFile]]: + for submission in list_submissions(): + for entry in submission.files: + if entry.file_id == file_id: + return submission, entry + return None + + +def collect_for_batch(file_ids: List[str]) -> Tuple[List[Tuple[str, bytes]], List[str]]: + """Read the chosen pending files into memory, in the order given. + + Returns (uploads, submitters). Raises KeyError naming the first id that is + unknown, and ValueError naming the first that is no longer pending - two + admin tabs open on the same inbox must not both start the same file, and + the second one gets told why rather than silently double-running it. + """ + uploads: List[Tuple[str, bytes]] = [] + submitters: List[str] = [] + + for file_id in file_ids: + found = find_file(file_id) + if found is None: + raise KeyError(file_id) + submission, entry = found + if entry.status != PENDING: + raise ValueError( + "'{}' was already {} and cannot be started again.".format( + entry.filename, entry.status + ) + ) + path = submission_dir(submission.submission_id) / entry.stored_name + uploads.append((entry.filename, path.read_bytes())) + if submission.submitted_by not in submitters: + submitters.append(submission.submitted_by) + + return uploads, submitters + + +def mark_consumed(file_ids: List[str], batch_id: str) -> None: + """Called only after the batch has actually been created and queued. + + Ordering matters: marking first and staging second would lose the files + from the inbox if staging then failed, with nothing left to retry from. + """ + _decide(file_ids, CONSUMED, batch_id=batch_id, + detail="Started as part of batch {}.".format(batch_id)) + + +def dismiss(file_ids: List[str]) -> int: + """An admin has decided these will never run. Returns how many changed. + + This exists so the badge can reach zero. Without it a file nobody intends + to process sits in the inbox forever, and a count that never clears is a + count people stop reading. + """ + return _decide(file_ids, DISMISSED, detail="Dismissed without running.") + + +def _decide(file_ids: List[str], status: str, *, batch_id: Optional[str] = None, + detail: Optional[str] = None) -> int: + wanted = set(file_ids) + changed = 0 + for submission in list_submissions(): + touched = False + for entry in submission.files: + if entry.file_id in wanted and entry.status == PENDING: + entry.status = status + entry.detail = detail + entry.batch_id = batch_id + entry.decided_at = time.time() + touched = True + changed += 1 + if touched: + submission.updated_at = time.time() + write_record(submission) + return changed + + +# --------------------------------------------------------------------------- +# Retention +# --------------------------------------------------------------------------- +def purge_expired(now: Optional[float] = None) -> List[str]: + """Delete submissions whose files are all decided and past retention. + + A submission with anything still PENDING is never purged, however old. + Deleting a file nobody has looked at yet would lose a colleague's work + silently, which is worse than the disk it occupies. + """ + if INBOX_RETENTION_DAYS <= 0: + return [] + cutoff = (now if now is not None else time.time()) - INBOX_RETENTION_DAYS * 86400 + removed: List[str] = [] + for submission in list_submissions(): + if submission.created_at >= cutoff: + continue + if any(f.status == PENDING for f in submission.files): + continue + try: + shutil.rmtree(submission_dir(submission.submission_id)) + removed.append(submission.submission_id) + except OSError as exc: + logger.warning("Could not purge submission %s: %s", + submission.submission_id, exc) + if removed: + logger.info("Purged %d expired inbox submission(s).", len(removed)) + return removed diff --git a/app/infrastructure/security.py b/app/infrastructure/security.py index 1b2e01f..f8df7e9 100644 --- a/app/infrastructure/security.py +++ b/app/infrastructure/security.py @@ -72,6 +72,22 @@ ROLE_PERMISSIONS: Dict[str, List[str]] = { "view_nutrition_insights", "optimize_profits", ], + # A third party who may drop spreadsheets into the review inbox and do + # NOTHING else. One permission, deliberately. + # + # This role exists because API keys carry no per-key scoping: + # principal_for_api_key() derives permissions entirely from the role, so + # "upload-only" can only be expressed as a role. Reusing `user` would have + # been less code and would also have handed an outside contributor + # add_product, upload_batch_products and upload_store_inventory - real + # write access to the catalog - to solve a problem that needed one verb. + # + # Nothing this role can do starts work: an uploaded file waits in the inbox + # until an admin selects it. So the worst an leaked uploader key costs is + # bounded disk (INBOX_MAX_PENDING_FILES), never CPU on a one-vCPU host. + "uploader": [ + "upload_catalog", + ], } VALID_ROLES = frozenset(ROLE_PERMISSIONS) diff --git a/app/infrastructure/settings.py b/app/infrastructure/settings.py index ddacfa8..945d312 100644 --- a/app/infrastructure/settings.py +++ b/app/infrastructure/settings.py @@ -169,6 +169,24 @@ BATCH_RETENTION_DAYS = int(os.getenv("BATCH_RETENTION_DAYS", "7")) # which is precisely how a slow start turns into an unrecoverable spiral. BATCH_AUTO_RESUME = _bool("BATCH_AUTO_RESUME", "false") +# --- Review inbox (a third party drops files; an admin decides) ------------- +# Files arrive here from POST /api/uploads/catalog and WAIT. Nothing in this +# directory is ever executed until an admin selects it, which is the whole +# security property of the feature: an uploader credential can consume bounded +# disk but can never consume CPU on a one-vCPU host. +INBOX_UPLOAD_DIR = _dir("INBOX_UPLOAD_DIR", DATA_DIR / "inbox") + +# Per-file ceilings are the batch ones (10MB / 2000 rows). This bounds the +# QUEUE: how many unreviewed files may accumulate before uploads are refused +# with 429. Without it an unattended key fills the disk one valid file at a +# time, and every one of them looks legitimate. +INBOX_MAX_PENDING_FILES = int(os.getenv("INBOX_MAX_PENDING_FILES", "200")) + +# Consumed and dismissed files are deleted this many days after upload. +# Pending files are never purged - deleting something nobody has looked at yet +# would lose a colleague's work silently. +INBOX_RETENTION_DAYS = int(os.getenv("INBOX_RETENTION_DAYS", "7")) + # Pristine copies of the bundled seed catalogs and pre-trained models, placed # here by the Dockerfile at a path that is never itself mounted over. # @@ -457,9 +475,15 @@ def _parse_api_keys(raw: str) -> dict: f"comma-separated between entries." ) name, role, secret = (p.strip() for p in parts) - if role not in {"admin", "user"}: + # MUST stay in step with ROLE_PERMISSIONS in app/infrastructure/security.py, + # which is the source of truth. It is duplicated rather than imported + # because security.py imports THIS module, so importing it back here + # would be a cycle. A role added there but not here is rejected at boot + # with the message below - loud, and before any request is served. + if role not in {"admin", "user", "uploader"}: raise RuntimeError( - f"API_KEYS entry {name!r} has role {role!r}; expected 'admin' or 'user'." + f"API_KEYS entry {name!r} has role {role!r}; expected 'admin', 'user' " + f"or 'uploader'." ) if not secret: raise RuntimeError(f"API_KEYS entry {name!r} has an empty secret.") @@ -476,6 +500,13 @@ def _parse_api_keys(raw: str) -> dict: # Machine consumers of api.. Empty by default - browser sessions go # through /api/auth/login instead, and a key that nobody needs is only risk. +# +# NAME THE KEY FOR ITS FUNCTION, NOT THE PERSON HOLDING IT. +# /api/health is public and reports {name, role, fingerprint} for every +# configured key (describe_api_keys in security.py). The secret is never +# exposed, but the NAME is - so `catalog-drop:uploader:...` is right and +# `priya-laptop:uploader:...` publishes a colleague's name to anyone who +# curls the health endpoint. API_KEYS = _parse_api_keys(os.getenv("API_KEYS", "")) # Default RAG behaviour diff --git a/app/main.py b/app/main.py index 98c3573..eaf8c45 100644 --- a/app/main.py +++ b/app/main.py @@ -31,7 +31,7 @@ from app.api.routers import health, brands, search, suggest, chat, catalog, syst from app.api.routers import stores, discounts, analytics as store_analytics, trending, recommendations, store_admin from app.api.routers import nutrition, nutrition_admin, upload from app.api.routers import auth, user_products, admin_train, mcp_info -from app.api.routers import store_catalog, batch_catalog +from app.api.routers import store_catalog, batch_catalog, uploads from app.services.store_db import ensure_store_intelligence_schema from app.services.nutrition_db import ensure_nutrition_schema @@ -293,6 +293,7 @@ app.include_router(nutrition_admin.router, prefix="/api") app.include_router(upload.router, prefix="/api") app.include_router(store_catalog.router, prefix="/api") app.include_router(batch_catalog.router, prefix="/api") +app.include_router(uploads.router, prefix="/api") app.include_router(mcp_info.router, prefix="/api") # MCP lives outside /api on purpose: it is a protocol endpoint for AI clients, diff --git a/tests/test_inbox_upload.py b/tests/test_inbox_upload.py new file mode 100644 index 0000000..1aca704 --- /dev/null +++ b/tests/test_inbox_upload.py @@ -0,0 +1,390 @@ +"""Two-actor flow: a colleague drops files, an admin decides what runs. + +Follows test_batch_catalog_ingest.py: point the directories at tmp_path, +monkeypatch the storage/embedding boundary, and stub the background worker by +default. Nothing here loads sentence-transformers or torch, nothing reaches the +network, and nothing writes into the repository's data directory. + +The most important tests in this file are the ones asserting what an `uploader` +credential CANNOT do. That role exists specifically so an outside contributor +does not get the `user` role's catalog write access, and an assertion is the +only thing that keeps it true as endpoints are added. +""" +from __future__ import annotations + +import io + +import pytest + +from app.core import batch_ingest, inbox +from app.core import store_catalog_pipeline as pipeline + +openpyxl = pytest.importorskip("openpyxl") + +HEADERS = ["Product Name", "Category", "Brand"] +ROWS = [["Amul Butter 100g", "Butter", "Amul"]] + +UPLOAD_KEY = "k" * 43 +UPLOAD = "/api/uploads/catalog" +INBOX = "/api/admin/catalog-batch/inbox" +FROM_INBOX = "/api/admin/catalog-batch/from-inbox" +DISMISS = "/api/admin/catalog-batch/inbox/dismiss" + + +def _csv(rows=ROWS) -> bytes: + return ("\n".join([",".join(HEADERS)] + [",".join(r) for r in rows])).encode() + + +def _sheet(rows=ROWS) -> bytes: + wb = openpyxl.Workbook() + ws = wb.active + ws.append(HEADERS) + for row in rows: + ws.append(row) + buf = io.BytesIO() + wb.save(buf) + return buf.getvalue() + + +def _files(*pairs): + return [("files", (n, io.BytesIO(c), "application/octet-stream")) + for n, c in pairs] + + +@pytest.fixture(autouse=True) +def _isolate_sku_counter(tmp_path, monkeypatch): + from app.services import sku_service + monkeypatch.setattr(sku_service, "_data_dir", tmp_path / "sku_sequences") + + +@pytest.fixture(autouse=True) +def dirs(tmp_path, monkeypatch): + """Both staging directories under tmp_path. Read at call time by design.""" + monkeypatch.setattr(inbox, "INBOX_UPLOAD_DIR", tmp_path / "inbox") + monkeypatch.setattr(batch_ingest, "BATCH_UPLOAD_DIR", tmp_path / "batch_uploads") + return tmp_path + + +@pytest.fixture(autouse=True) +def no_background_worker(monkeypatch): + """Same trap as in test_batch_catalog_ingest: a worker still running after + teardown resolves the directory setting again and writes into the real + data/ directory. Stub it; the tests here are about routing and state.""" + from app.core import batch_worker + + submitted: list = [] + monkeypatch.setattr(batch_worker, "submit", submitted.append) + return submitted + + +@pytest.fixture(autouse=True) +def upload_key(monkeypatch): + """One uploader API key. + + Patched on `security`, not on `settings`: security.py does + `from ...settings import API_KEYS` at import, binding the dict object, so + rebinding the name in `settings` leaves `principal_for_api_key` still + looking at the original and every request comes back 401. + """ + from app.infrastructure import security + + monkeypatch.setattr( + security, "API_KEYS", {UPLOAD_KEY: ("catalog-drop", "uploader")} + ) + return {"X-API-Key": UPLOAD_KEY} + + +@pytest.fixture +def store(monkeypatch): + table: dict = {} + + def fake_upsert(brand, rows, cleanup=False): + assert cleanup is False + for row in rows: + table[row["image_id"]] = dict(row) + return len(rows) + + monkeypatch.setattr(pipeline, "upsert_brand_products", fake_upsert) + monkeypatch.setattr(pipeline, "get_products_by_brand", lambda b, **kw: list(table.values())) + monkeypatch.setattr(pipeline, "embed_texts", lambda texts: [[0.0] * 384 for _ in texts]) + return table + + +def _submit(client, upload_key, *pairs): + return client.post(UPLOAD, files=_files(*pairs), headers=upload_key) + + +# --------------------------------------------------------------------------- +# The security boundary - the reason the `uploader` role exists +# --------------------------------------------------------------------------- +def test_uploader_key_can_submit(client, upload_key): + r = _submit(client, upload_key, ("a.csv", _csv())) + assert r.status_code == 202 + body = r.json() + assert body["files_accepted"] == 1 + assert body["submitted_by"] == "catalog-drop" + + +@pytest.mark.parametrize("method,path", [ + ("get", "/api/admin/catalog-batch/inbox"), + ("post", "/api/admin/catalog-batch/from-inbox"), + ("post", "/api/admin/catalog-batch/inbox/dismiss"), + ("post", "/api/admin/catalog-batch/preview"), + ("post", "/api/admin/catalog-batch/ingest"), + ("get", "/api/admin/catalog-batch/batches"), + ("get", "/api/admin/catalog-batch/batches/abc"), + ("post", "/api/admin/catalog-batch/batches/abc/resume"), + ("post", "/api/admin/catalog-batch/batches/abc/cancel"), + ("post", "/api/admin/store-catalog/preview"), + ("post", "/api/admin/store-catalog/ingest"), + ("get", "/api/admin/store-catalog/jobs/abc"), +]) +def test_uploader_key_is_refused_on_every_admin_route(client, upload_key, method, path): + """An uploader credential must reach nothing but its own endpoint.""" + r = getattr(client, method)(path, headers=upload_key) + assert r.status_code == 403, f"{method.upper()} {path} returned {r.status_code}" + + +@pytest.mark.parametrize("path", [ + "/api/user/products/add", + "/api/user/products/upload-file", + "/api/upload/stores", +]) +def test_uploader_key_has_none_of_the_user_roles_write_access(client, upload_key, path): + """This is the specific over-grant avoided by NOT reusing the `user` role. + + A `user` key would carry add_product / upload_batch_products / + upload_store_inventory - real catalog writes - to solve a problem that + needed one verb. + """ + r = client.post(path, headers=upload_key) + assert r.status_code == 403, f"{path} returned {r.status_code}" + + +def test_upload_endpoint_rejects_anonymous(client): + r = client.post(UPLOAD, files=_files(("a.csv", _csv()))) + assert r.status_code == 401 + + +def test_admin_can_still_use_the_upload_endpoint(client, admin_headers): + """admin bypasses every permission check (Principal.has_permission).""" + r = client.post(UPLOAD, files=_files(("a.csv", _csv())), headers=admin_headers) + assert r.status_code == 202 + + +def test_role_allow_lists_stay_in_step(client): + """settings._parse_api_keys duplicates the role names from security.py + because importing back would be a cycle. Pin the duplication.""" + from app.infrastructure import settings + from app.infrastructure.security import VALID_ROLES + + for role in VALID_ROLES: + parsed = settings._parse_api_keys(f"n:{role}:{'x' * 43}") + assert parsed[("x" * 43)] == ("n", role) + + with pytest.raises(RuntimeError): + settings._parse_api_keys(f"n:wizard:{'x' * 43}") + + +# --------------------------------------------------------------------------- +# Intake +# --------------------------------------------------------------------------- +def test_a_broken_sheet_is_refused_at_upload_and_never_enters_the_inbox(client, upload_key): + r = _submit(client, upload_key, ("good.csv", _csv()), ("broken.pdf", b"%PDF-1.4")) + assert r.status_code == 202 + body = r.json() + assert body["files_accepted"] == 1 + assert body["files_rejected"] == 1 + bad = [f for f in body["files"] if f["filename"] == "broken.pdf"][0] + assert bad["accepted"] is False and bad["error"] + # Only the good one is waiting. + assert inbox.pending_count() == 1 + + +def test_a_drop_where_nothing_parses_stages_no_submission(client, upload_key): + r = _submit(client, upload_key, ("broken.pdf", b"%PDF-1.4")) + assert r.status_code == 202 + assert r.json()["submission_id"] is None + assert inbox.pending_count() == 0 + assert inbox.list_submissions() == [] + + +def test_uploaded_filenames_cannot_escape_the_inbox(client, upload_key, dirs): + r = _submit(client, upload_key, ("../../evil.csv", _csv())) + assert r.status_code == 202 + submission = inbox.list_submissions()[0] + directory = inbox.submission_dir(submission.submission_id) + for entry in submission.files: + written = directory / entry.stored_name + assert written.resolve().parent == directory.resolve() + assert not (dirs / "evil.csv").exists() + + +def test_a_full_inbox_refuses_further_uploads(client, upload_key, monkeypatch): + from app.api.routers import uploads + + monkeypatch.setattr(uploads, "INBOX_MAX_PENDING_FILES", 1) + assert _submit(client, upload_key, ("a.csv", _csv())).status_code == 202 + r = _submit(client, upload_key, ("b.csv", _csv())) + assert r.status_code == 429 + assert "awaiting review" in r.json()["detail"] + + +def test_too_many_files_in_one_drop_is_413(client, upload_key, monkeypatch): + from app.api.routers import uploads + + monkeypatch.setattr(uploads, "BATCH_MAX_FILES", 2) + r = _submit(client, upload_key, + ("a.csv", _csv()), ("b.csv", _csv()), ("c.csv", _csv())) + assert r.status_code == 413 + + +def test_uploading_starts_no_work(client, upload_key, no_background_worker): + """The whole safety property: an uploader can spend disk, never CPU.""" + _submit(client, upload_key, ("a.csv", _csv())) + assert no_background_worker == [] + + +# --------------------------------------------------------------------------- +# The admin side +# --------------------------------------------------------------------------- +def test_inbox_lists_pending_grouped_by_drop(client, upload_key, admin_headers): + _submit(client, upload_key, ("a.csv", _csv()), ("b.csv", _csv())) + _submit(client, upload_key, ("c.csv", _csv())) + + r = client.get(INBOX, headers=admin_headers) + assert r.status_code == 200 + body = r.json() + assert body["pending_count"] == 3 + assert len(body["submissions"]) == 2 + assert all(s["submitted_by"] == "catalog-drop" for s in body["submissions"]) + assert body["submissions"][0]["files"][0]["rows_total"] == 1 + + +def test_selecting_files_across_two_drops_makes_one_batch( + client, upload_key, admin_headers, store, no_background_worker +): + """The requirement that ruled out reusing batch-resume: files from + different drops, composed into a single run.""" + _submit(client, upload_key, ("a.csv", _csv()), ("b.csv", _csv())) + _submit(client, upload_key, ("c.csv", _csv())) + + listing = client.get(INBOX, headers=admin_headers).json() + first = listing["submissions"][0]["files"][0]["file_id"] + second = listing["submissions"][1]["files"][0]["file_id"] + + r = client.post(FROM_INBOX, json={"file_ids": [first, second]}, headers=admin_headers) + assert r.status_code == 202 + batch = r.json() + assert batch["files_total"] == 2 + assert batch["submitted_by"] == "catalog-drop" + assert no_background_worker == [batch["batch_id"]] + + # The unselected file is untouched. + after = client.get(INBOX, headers=admin_headers).json() + assert after["pending_count"] == 1 + + +def test_a_started_file_cannot_be_started_twice(client, upload_key, admin_headers, store): + """Two admin tabs on the same inbox must not double-run a file.""" + _submit(client, upload_key, ("a.csv", _csv())) + file_id = client.get(INBOX, headers=admin_headers).json()["submissions"][0]["files"][0]["file_id"] + + assert client.post(FROM_INBOX, json={"file_ids": [file_id]}, headers=admin_headers).status_code == 202 + again = client.post(FROM_INBOX, json={"file_ids": [file_id]}, headers=admin_headers) + assert again.status_code == 409 + assert "already consumed" in again.json()["detail"] + + +def test_selecting_an_unknown_file_is_404(client, admin_headers): + r = client.post(FROM_INBOX, json={"file_ids": ["nope"]}, headers=admin_headers) + assert r.status_code == 404 + + +def test_selecting_nothing_is_400(client, admin_headers): + assert client.post(FROM_INBOX, json={"file_ids": []}, headers=admin_headers).status_code == 400 + + +def test_dismiss_clears_the_badge_without_running_anything( + client, upload_key, admin_headers, no_background_worker +): + _submit(client, upload_key, ("a.csv", _csv())) + file_id = client.get(INBOX, headers=admin_headers).json()["submissions"][0]["files"][0]["file_id"] + + r = client.post(DISMISS, json={"file_ids": [file_id]}, headers=admin_headers) + assert r.status_code == 200 + assert r.json() == {"dismissed": 1, "pending_count": 0} + assert no_background_worker == [], "dismiss must not start work" + assert client.get(INBOX, headers=admin_headers).json()["pending_count"] == 0 + + +def test_dismissing_an_already_decided_file_is_409(client, upload_key, admin_headers): + _submit(client, upload_key, ("a.csv", _csv())) + file_id = client.get(INBOX, headers=admin_headers).json()["submissions"][0]["files"][0]["file_id"] + client.post(DISMISS, json={"file_ids": [file_id]}, headers=admin_headers) + again = client.post(DISMISS, json={"file_ids": [file_id]}, headers=admin_headers) + assert again.status_code == 409 + + +# --------------------------------------------------------------------------- +# End to end, and retention +# --------------------------------------------------------------------------- +def test_selected_files_actually_reach_the_catalog(client, upload_key, admin_headers, store): + """Run the composed batch for real (worker stubbed, pipeline not).""" + _submit(client, upload_key, ("a.csv", _csv()), ("b.xlsx", _sheet())) + listing = client.get(INBOX, headers=admin_headers).json() + ids = [f["file_id"] for f in listing["submissions"][0]["files"]] + + started = client.post(FROM_INBOX, json={"file_ids": ids}, headers=admin_headers).json() + result = batch_ingest.run_batch(started["batch_id"]) + + assert result.status == "done" + assert result.files_done == 2 + assert result.submitted_by == "catalog-drop" + assert store, "no rows reached the fake brand table" + + +def test_purge_removes_decided_submissions_and_spares_pending(client, upload_key, admin_headers, monkeypatch): + import time + + monkeypatch.setattr(inbox, "INBOX_RETENTION_DAYS", 7) + + _submit(client, upload_key, ("old.csv", _csv())) + old = inbox.list_submissions()[0] + file_id = old.files[0].file_id + client.post(DISMISS, json={"file_ids": [file_id]}, headers=admin_headers) + aged = inbox.read_record(old.submission_id) + aged.created_at = time.time() - 30 * 86400 + inbox.write_record(aged) + + _submit(client, upload_key, ("pending.csv", _csv())) + still_pending = [s for s in inbox.list_submissions() if s.pending_files][0] + stale_but_unreviewed = inbox.read_record(still_pending.submission_id) + stale_but_unreviewed.created_at = time.time() - 30 * 86400 + inbox.write_record(stale_but_unreviewed) + + removed = inbox.purge_expired() + + assert removed == [old.submission_id] + assert not inbox.submission_dir(old.submission_id).exists() + assert inbox.submission_dir(still_pending.submission_id).exists(), \ + "a file nobody has reviewed must never be purged" + + +def test_purge_is_disabled_when_retention_is_zero(client, upload_key, monkeypatch): + monkeypatch.setattr(inbox, "INBOX_RETENTION_DAYS", 0) + _submit(client, upload_key, ("a.csv", _csv())) + assert inbox.purge_expired() == [] + + +def test_unreadable_record_is_skipped_not_raised(client, upload_key): + _submit(client, upload_key, ("a.csv", _csv())) + sid = inbox.list_submissions()[0].submission_id + inbox.record_path(sid).write_text("{ not json", encoding="utf-8") + assert inbox.read_record(sid) is None + assert inbox.list_submissions() == [] + + +def test_submission_dir_rejects_a_traversing_id(): + with pytest.raises(ValueError): + inbox.submission_dir("../escape")