diff --git a/app/api/batch_common.py b/app/api/batch_common.py index 15b1885..c8c8346 100644 --- a/app/api/batch_common.py +++ b/app/api/batch_common.py @@ -266,3 +266,50 @@ def stage_and_queue( return manifest, False return manifest, True + + +def stage_pending( + valid: List[Tuple[str, bytes, int]], + invalid: List[Tuple[str, str]], + *, + submitted_by: Optional[str] = None, +) -> batch_ingest.BatchManifest: + """Write the files down and stop. The review inbox half of the flow. + + Deliberately a separate function rather than a `queue=False` flag on + `stage_and_queue`: the admin routes must queue unconditionally, and a shared + boolean is exactly the kind of default that gets inverted in a later edit and + silently starts running work nobody approved. + + `use_llm` / `fetch_images` are not taken here. They are run-time choices, and + the person who makes them is the admin pressing Start in the inbox - not the + colleague who dropped the file. They are supplied to `stage_and_queue` when + the selected files become a real batch. + """ + manifest = batch_ingest.stage_batch( + [(name, contents) for name, contents, _n in valid], + invalid=invalid, + submitted_by=submitted_by, + ) + # stage_batch records size but not row counts; it never parsed the files. + for entry, (_name, _contents, rows) in zip(manifest.files, valid): + entry.rows_total = rows + + manifest.status = batch_ingest.PENDING + manifest.detail = "Waiting for review. Nothing runs until an admin starts it." + batch_ingest.write_manifest(manifest) + batch_job_store.put(manifest) + return manifest + + +def pending_submissions() -> List[batch_ingest.BatchManifest]: + """Every drop awaiting review, newest first. + + Read from disk rather than the in-memory store: the store is populated by + whatever this process has seen, and an inbox that empties itself on restart + would look exactly like a colleague's files having been processed. + """ + return [ + m for m in batch_ingest.list_manifests() + if m.status == batch_ingest.PENDING and m.files + ] diff --git a/app/api/routers/batch_catalog.py b/app/api/routers/batch_catalog.py index 3780b5d..25283b6 100644 --- a/app/api/routers/batch_catalog.py +++ b/app/api/routers/batch_catalog.py @@ -6,6 +6,9 @@ GET /api/admin/catalog-batch/batches/{id} - poll one batch POST /api/admin/catalog-batch/batches/{id}/resume - after a restart POST /api/admin/catalog-batch/batches/{id}/cancel - stop the rest + GET /api/admin/catalog-batch/inbox - drops awaiting review + POST /api/admin/catalog-batch/from-inbox - run selected files + POST /api/admin/catalog-batch/inbox/dismiss - discard selected files This is the multi-file sibling of `store_catalog.py`, and it deliberately does not replace it: the single-file endpoints are untouched and still work. What is @@ -33,9 +36,10 @@ difference is that its reads are filtered to the caller's own submissions. from __future__ import annotations import queue -from typing import List +from typing import Dict, List, Optional from fastapi import APIRouter, Depends, File, HTTPException, UploadFile, status +from pydantic import BaseModel, Field from app.api import batch_common from app.api.batch_job_store import batch_job_store @@ -172,8 +176,20 @@ async def ingest_catalog_batch( @router.get("/batches", dependencies=[Depends(require_admin)]) def list_catalog_batches(limit: int = 20) -> dict: + """Runs, newest first. Inbox drops are NOT runs and are excluded. + + A `pending` submission has never been near the pipeline. Listing it here + would put a row with no progress and no result in the Batch tab, next to + real runs, and the Resume button beside it would be a lie. + """ limit = max(1, min(limit, 100)) - return {"batches": [_to_out(m) for m in batch_job_store.recent(limit)]} + # Over-fetch, then drop the pending ones, so filtering cannot return fewer + # than `limit` runs just because the inbox happens to be busy. + runs = [ + m for m in batch_job_store.recent(limit * 10) + if m.status != batch_ingest.PENDING + ][:limit] + return {"batches": [_to_out(m) for m in runs]} @router.get("/batches/{batch_id}", dependencies=[Depends(require_admin)]) @@ -245,3 +261,216 @@ def cancel_catalog_batch(batch_id: str) -> BatchOut: manifest.detail = "Cancelling - the file currently running will finish first." batch_job_store.put(manifest) return _to_out(manifest) + + +# --------------------------------------------------------------------------- +# The review inbox +# --------------------------------------------------------------------------- +# The other half of the flow in app/api/routers/uploads.py: a colleague posts +# spreadsheets there with no credential, they land as `pending`, and nothing +# runs until somebody here says so. These three routes are what the Inbox tab +# in the admin UI calls (frontend/src/pages/InboxPanel.jsx). +# +# Selection is per FILE and crosses submissions on purpose. "Two sheets from +# Monday's drop plus one from today, as one batch" is the request an operator +# actually has, and it has no expression in a model where the unit is the drop. + + +class InboxFileOut(BaseModel): + """One spreadsheet waiting for review.""" + + # "{batch_id}:{index}". Addressed by a compound id rather than a bare index + # because the UI holds one flat selection set spanning every submission, and + # an index alone is not unique across two of them. + file_id: str + filename: str + rows_total: int = 0 + size_bytes: int = 0 + + +class InboxSubmissionOut(BaseModel): + """One drop: the files that arrived together, and who sent them.""" + + submission_id: str + submitted_by: Optional[str] = None + created_at: float + files: List[InboxFileOut] + + +class InboxOut(BaseModel): + pending_count: int + submissions: List[InboxSubmissionOut] + + +class InboxSelection(BaseModel): + """The files an admin ticked.""" + + file_ids: List[str] = Field(default_factory=list) + + +class InboxStartRequest(InboxSelection): + # Chosen HERE, not by the sender - see the note on the POST handler in + # uploads.py. These commit the host to outbound work, so the decision + # belongs to the person who can see what the machine is already doing. + use_llm: bool = False + fetch_images: bool = False + + +class InboxDismissOut(BaseModel): + dismissed: int + + +def _parse_file_ids(file_ids: List[str]) -> Dict[str, List[int]]: + """Group "{batch_id}:{index}" into {batch_id: [index, ...]}. + + Malformed ids are dropped rather than raising: the UI polls every five + seconds and prunes its selection against what came back, so a tick can + legitimately refer to a file another tab started a moment ago. Failing the + whole request would let one stale checkbox block the rest. + """ + grouped: Dict[str, List[int]] = {} + for raw in file_ids or []: + batch_id, _, index = str(raw).partition(":") + if not batch_id or not index.isdigit(): + continue + grouped.setdefault(batch_id, []).append(int(index)) + return grouped + + +@router.get("/inbox", dependencies=[Depends(require_admin)]) +def list_inbox() -> InboxOut: + """Every drop awaiting review, newest first.""" + submissions = [] + pending_count = 0 + for manifest in batch_common.pending_submissions(): + files = [ + InboxFileOut( + file_id=f"{manifest.batch_id}:{entry.index}", + filename=entry.filename, + rows_total=entry.rows_total, + size_bytes=entry.size_bytes, + ) + # A file the sender's own upload already rejected is recorded on the + # manifest so they can be told about it, but it has no bytes on disk + # and cannot be run - so it is not offered for selection. + for entry in manifest.files + if entry.status != batch_ingest.FAILED and entry.stored_name + ] + if not files: + continue + pending_count += len(files) + submissions.append( + InboxSubmissionOut( + submission_id=manifest.batch_id, + submitted_by=manifest.submitted_by, + created_at=manifest.created_at, + files=files, + ) + ) + return InboxOut(pending_count=pending_count, submissions=submissions) + + +@router.post("/from-inbox", status_code=status.HTTP_202_ACCEPTED, + dependencies=[Depends(require_admin)]) +def start_batch_from_inbox(request: InboxStartRequest) -> BatchOut: + """Take the selected files out of the inbox and run them as one batch. + + The bytes are COPIED into a fresh batch rather than the pending manifest + being promoted in place. Two reasons: a selection can span submissions, and + there is no such thing as promoting two manifests into one; and the run gets + its own id, so it has an identity distinct from the drop it came from - + which is what the Batch tab lists and what Resume acts on. + + The originals are removed afterwards, so the same sheet cannot be started + twice from a stale checkbox in another tab. + """ + grouped = _parse_file_ids(request.file_ids) + if not grouped: + raise HTTPException(status_code=400, detail="No files were selected.") + + # (filename, bytes, rows) - the shape stage_and_queue takes, which is also + # what parse_all returns, so nothing needs reparsing here. + picked: List[tuple] = [] + senders: List[str] = [] + for batch_id, indices in grouped.items(): + manifest = batch_ingest.read_manifest(batch_id) + if not manifest or manifest.status != batch_ingest.PENDING: + continue + directory = batch_ingest.batch_dir(batch_id) + for entry in manifest.files: + if entry.index not in indices or not entry.stored_name: + continue + try: + contents = (directory / entry.stored_name).read_bytes() + except OSError: + # Retention or a concurrent dismiss got there first. Skipping is + # right: the file is genuinely gone, and the poll that follows + # will show it has left the inbox. + continue + picked.append((entry.filename, contents, entry.rows_total)) + if manifest.submitted_by: + senders.append(manifest.submitted_by) + + if not picked: + raise HTTPException( + status_code=409, + detail=( + "None of those files are still waiting - they may have been " + "started or dismissed already. Refresh the inbox." + ), + ) + + # Preserved so the Batch tab can say where a run came from: one name when a + # drop came from one colleague, a joined list when a batch was assembled + # from several, which is exactly when the question gets asked. + unique_senders = sorted(set(senders)) + submitted_by = ", ".join(unique_senders)[:120] if unique_senders else None + + manifest, started = batch_common.stage_and_queue( + picked, + [], + use_llm=request.use_llm, + fetch_images=request.fetch_images, + submitted_by=submitted_by, + ) + if not started: + raise HTTPException( + status_code=429, + detail=( + "Too many batches are already queued. These files have been " + "taken out of the inbox and saved as batch " + f"{manifest.batch_id} - press Resume on it once the current " + "batch finishes." + ), + ) + + # Only now, once the bytes are safely staged under a new id. Removing them + # first would lose the files outright if staging then failed. + for batch_id, indices in grouped.items(): + batch_ingest.remove_files(batch_id, indices) + + return _to_out(manifest) + + +@router.post("/inbox/dismiss", dependencies=[Depends(require_admin)]) +def dismiss_inbox_files(request: InboxSelection) -> InboxDismissOut: + """Discard the selected files. The bytes go with them. + + Deliberately irreversible and deliberately unceremonious: this is the + disposal path for a drop nobody wants, and on an endpoint anyone can post to + it is the control that keeps the volume from filling with rejected sheets. + """ + grouped = _parse_file_ids(request.file_ids) + if not grouped: + raise HTTPException(status_code=400, detail="No files were selected.") + + dismissed = 0 + for batch_id, indices in grouped.items(): + manifest = batch_ingest.read_manifest(batch_id) + # Only ever the inbox. Without this check a crafted id would delete + # files out of a batch that was mid-run. + if not manifest or manifest.status != batch_ingest.PENDING: + continue + dismissed += batch_ingest.remove_files(batch_id, indices) + + return InboxDismissOut(dismissed=dismissed) diff --git a/app/api/routers/uploads.py b/app/api/routers/uploads.py index 8c3e19a..0b2c2e5 100644 --- a/app/api/routers/uploads.py +++ b/app/api/routers/uploads.py @@ -47,19 +47,21 @@ delete; those stay on the admin router. from __future__ import annotations import logging -from typing import List +from typing import List, Optional -from fastapi import APIRouter, Depends, File, HTTPException, UploadFile, status +from fastapi import APIRouter, Depends, File, Form, HTTPException, UploadFile, status from app.api import batch_common from app.api.batch_job_store import batch_job_store -from app.api.deps import require_permission +from app.api.deps import get_optional_principal, require_permission from app.core import batch_ingest 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_BYTES, + INBOX_MAX_PENDING_FILES, ) logger = logging.getLogger(__name__) @@ -101,25 +103,71 @@ def _owns(manifest: batch_ingest.BatchManifest, principal: Principal) -> bool: return bool(manifest.submitted_by) and manifest.submitted_by == principal.username +def _inbox_capacity_or_429(incoming_files: int, incoming_bytes: int) -> None: + """Refuse a drop that would push the review inbox past its ceiling. + + This endpoint takes files from anyone, and it queues nothing, so + BATCH_QUEUE_MAX - the bound that made an authenticated uploader safe - does + not apply here. Unreviewed submissions accumulate on the volume until + somebody acts on them, which on this host is the scarcest resource there is. + + Counted over drops still awaiting review only, so starting or dismissing one + frees its share at once. + """ + pending = batch_common.pending_submissions() + files_now = sum(len(m.files) for m in pending) + bytes_now = sum(f.size_bytes for m in pending for f in m.files) + + if files_now + incoming_files > INBOX_MAX_PENDING_FILES: + raise HTTPException( + status_code=429, + detail=( + f"The review inbox is full ({files_now} file(s) awaiting review, " + f"limit {INBOX_MAX_PENDING_FILES}). Nothing was stored. Ask an " + f"admin to clear the inbox, then resend." + ), + ) + if bytes_now + incoming_bytes > INBOX_MAX_PENDING_BYTES: + raise HTTPException( + status_code=429, + detail=( + f"The review inbox is full " + f"({bytes_now // (1024 * 1024)}MB awaiting review, limit " + f"{INBOX_MAX_PENDING_BYTES // (1024 * 1024)}MB). Nothing was " + f"stored. Ask an admin to clear the inbox, then resend." + ), + ) + + @router.post("/catalog", status_code=status.HTTP_202_ACCEPTED) async def ingest_catalog_files( files: List[UploadFile] = File(...), - use_llm: bool = False, - fetch_images: bool = False, - principal: Principal = Depends(require_permission("upload_catalog")), + sender: Optional[str] = Form(None), + principal: Optional[Principal] = Depends(get_optional_principal), ) -> CatalogUploadOut: - """Accept spreadsheets and run the catalog pipeline over them. + """Accept spreadsheets into the review inbox. No credential required. Returns 202 and a `batch_id` to poll. A drop where some files parse and some - do not is a partial success, not a failure: the good ones are ingested and - the bad ones come back in `files` as `status: "failed"` with the reason, so - the caller knows exactly which sheet to fix and resend. + do not is a partial success, not a failure: the good ones are stored and the + bad ones come back in `files` as `status: "failed"` with the reason, so the + caller knows exactly which sheet to fix and resend. - `use_llm` and `fetch_images` default OFF, the opposite of the single-file - admin route. Both are network stages, and image search in particular spawns - a Playwright subprocess that can spend minutes per batch on a host with one - vCPU. An API client that genuinely wants them can ask; an API client that - does not think about it gets the cheap, predictable path. + NOTHING SENT HERE RUNS ON ARRIVAL. + The files are staged and the batch is left `pending`. An admin sees it in the + review inbox, ticks the sheets they want and presses Start; only then does + anything reach the pipeline or the catalog. That gate is what makes an + endpoint anybody can post to acceptable: the cost of an unwanted drop is + disk until someone declines it, not products in the live catalog. + + `use_llm` and `fetch_images` are NOT accepted here, though they used to be. + They decide how a run behaves, and the person who decides that is now the + admin pressing Start - not the sender. Leaving them on this endpoint would + let an anonymous caller commit the host to Playwright image search. + + `sender` is a free-text label, not identity - it is whatever the caller + typed. It exists because the inbox groups drops by who sent them, and three + colleagues all showing as "anonymous" is an inbox nobody can triage. A real + credential, if one is presented, wins over it. """ limits = _limits() read = await batch_common.read_uploads(files, limits) @@ -132,45 +180,43 @@ async def ingest_catalog_files( detail=f"None of the uploaded files could be ingested. {detail}", ) - manifest, started = batch_common.stage_and_queue( - valid, - invalid, - use_llm=use_llm, - fetch_images=fetch_images, - # 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, + _inbox_capacity_or_429( + incoming_files=len(valid), + incoming_bytes=sum(len(contents) for _n, contents, _r in valid), ) - if not started: - # The batch is staged and durable, but this caller has no Resume button - # - that lives on the admin router - so the honest instruction is to - # send it again shortly. The id is included so an admin can find and - # resume this one instead if the caller reports it. - raise HTTPException( - status_code=429, - detail=( - f"Too many batches are already queued. Batch {manifest.batch_id} has " - f"been saved but not started; retry this upload shortly." - ), - ) + # A presented credential still names the sender - `get_optional_principal` + # returns None only when NO credential was sent, and still raises on one + # that is present and wrong. `sender` is trusted for a label and nothing + # else; it is truncated because it is rendered in the admin UI. + submitted_by = ( + principal.username if principal + else ((sender or "").strip()[:60] or "anonymous") + ) + + manifest = batch_common.stage_pending( + valid, + invalid, + submitted_by=submitted_by, + ) logger.info( - "Catalog batch %s queued: %d file(s), %d rejected, from %s", - manifest.batch_id, len(valid), len(invalid), principal.username, + "Catalog drop %s received for review: %d file(s), %d rejected, from %s%s", + manifest.batch_id, len(valid), len(invalid), submitted_by, + "" if principal else " (no credential)", ) body = batch_common.to_out(manifest).model_dump() + message = ( + f"{len(valid)} file(s) received and waiting for review. Nothing runs " + f"until an admin starts them. " + f"Poll GET /api/uploads/catalog/{manifest.batch_id} for status." + ) if invalid: message = ( - f"{len(valid)} file(s) accepted and queued for ingestion. " - f"{len(invalid)} could not be read - see 'files' for the reason on each, " - f"and resend those." - ) - else: - message = ( - f"{len(valid)} file(s) accepted and queued for ingestion. " - f"Poll GET /api/uploads/catalog/{manifest.batch_id} for progress." + f"{len(valid)} file(s) received and waiting for review. " + f"{len(invalid)} could not be read - see 'files' for the reason on " + f"each, and resend those." ) return CatalogUploadOut(**body, message=message) @@ -180,7 +226,12 @@ def list_my_catalog_batches( limit: int = 20, principal: Principal = Depends(require_permission("upload_catalog")), ) -> dict: - """The batches this credential has sent, newest first.""" + """The batches this credential has sent, newest first. + + Still credentialed, unlike the single-batch read below. An anonymous LIST + would hand any caller every other sender's drops in one request, which is + a different thing entirely from letting someone check the id they hold. + """ limit = max(1, min(limit, 100)) # Over-fetch before filtering: `recent` orders by creation across every # caller, so taking `limit` first would return fewer than `limit` of this @@ -194,15 +245,26 @@ def list_my_catalog_batches( @router.get("/catalog/{batch_id}") def get_my_catalog_batch( batch_id: str, - principal: Principal = Depends(require_permission("upload_catalog")), + principal: Optional[Principal] = Depends(get_optional_principal), ) -> batch_common.BatchOut: - """Progress and result for one batch this credential sent. + """Progress and result for one drop, addressed by its id. - A batch belonging to someone else is a 404, not a 403: whether a given id - exists is not this caller's business, and the two answers must therefore be - indistinguishable. + THE ID IS THE CREDENTIAL HERE, and it has to be: the sender needed no + credential to post, so requiring one to read the result would leave them + unable to find out what happened to their own file. `batch_id` is a + `uuid4().hex` handed only to whoever submitted the drop - 128 bits, not + enumerable - so holding it is the proof of having sent it. + + A caller who DID present a credential is held to it, and sees only their own + submissions. That is stricter than anonymous access to the same row, which + is the right way round: a named key should not become a way to browse. + + Either way an id you may not see is a 404, never a 403 - whether it exists + is not the caller's business, so the two answers must be indistinguishable. """ manifest = batch_job_store.get(batch_id) - if not manifest or not _owns(manifest, principal): + if not manifest: + raise HTTPException(status_code=404, detail="Batch not found") + if principal is not None and not _owns(manifest, principal): raise HTTPException(status_code=404, detail="Batch not found") return batch_common.to_out(manifest) diff --git a/app/core/batch_ingest.py b/app/core/batch_ingest.py index 0e59a04..724a941 100644 --- a/app/core/batch_ingest.py +++ b/app/core/batch_ingest.py @@ -73,6 +73,12 @@ CANCELLED = "cancelled" INTERRUPTED = "interrupted" PARTIAL = "partial" +# A drop sitting in the review inbox: staged on disk, deliberately NOT queued. +# It is a batch state rather than a store of its own so that retention, the +# manifest format and the job store all apply to it unchanged. Its FILES stay +# `queued`, which is why `settle()` has to leave this status alone - see there. +PENDING = "pending" + TERMINAL_BATCH_STATES = {DONE, FAILED, PARTIAL, CANCELLED} # `..`, separators and drive letters all stripped. UploadFile.filename is @@ -167,7 +173,15 @@ class BatchManifest: return sorted(seen) def settle(self) -> str: - """Recompute the batch status from its files. Returns the new status.""" + """Recompute the batch status from its files. Returns the new status. + + A `pending` submission is exempt. Its files are `queued` - they are + genuinely waiting - so the rule below would promote the batch to + `queued` and it would read as work already accepted, which is the one + thing the review inbox exists to prevent. + """ + if self.status == PENDING: + return self.status states = {f.status for f in self.files} if states & {QUEUED, RUNNING}: self.status = RUNNING if RUNNING in states else QUEUED @@ -284,6 +298,57 @@ def list_manifests() -> List[BatchManifest]: return sorted(found, key=lambda m: m.created_at, reverse=True) +def remove_files(batch_id: str, indices) -> int: + """Drop these files from a staged batch, bytes and all. Returns the count. + + Used by both halves of the review inbox: dismissing files, and taking them + into a run (`from-inbox` copies the bytes into a fresh batch first, so the + original submission must not keep serving them a second time). + + Surviving files KEEP their original `index`. The inbox addresses a file as + "{batch_id}:{index}", and the admin UI holds those ids in a selection set + across polls, so renumbering here would silently repoint a tick from the + file someone chose to whichever file slid into its place. + + A submission with nothing left is deleted outright rather than left as an + empty drop the inbox would render as a heading with no rows. + """ + manifest = read_manifest(batch_id) + if not manifest: + return 0 + + wanted = {int(i) for i in indices} + keep, dropped = [], 0 + for entry in manifest.files: + if entry.index not in wanted: + keep.append(entry) + continue + if entry.stored_name: + try: + (batch_dir(batch_id) / entry.stored_name).unlink(missing_ok=True) + except OSError as exc: + logger.warning( + "Could not delete %s from batch %s: %s", + entry.stored_name, batch_id, exc, + ) + dropped += 1 + + if not dropped: + return 0 + + if not keep: + try: + shutil.rmtree(batch_dir(batch_id)) + except OSError as exc: + logger.warning("Could not remove empty batch %s: %s", batch_id, exc) + return dropped + + manifest.files = keep + manifest.updated_at = time.time() + write_manifest(manifest) + return dropped + + # --------------------------------------------------------------------------- # Staging # --------------------------------------------------------------------------- @@ -490,6 +555,12 @@ def scan_interrupted() -> List[str]: for manifest in list_manifests(): if manifest.status in TERMINAL_BATCH_STATES or manifest.status == INTERRUPTED: continue + # A pending drop was never running, so a restart did not cut it short. + # This sweep takes every non-terminal manifest, so without this line + # every restart would relabel the whole review inbox "interrupted" and + # offer an admin a Resume button for work nobody had started. + if manifest.status == PENDING: + continue for entry in manifest.files: if entry.status == RUNNING: entry.status = QUEUED @@ -524,9 +595,14 @@ def purge_expired(now: Optional[float] = None) -> List[str]: for manifest in list_manifests(): if manifest.created_at >= cutoff: continue - if manifest.status not in TERMINAL_BATCH_STATES: + if manifest.status not in TERMINAL_BATCH_STATES and manifest.status != PENDING: # An old batch still queued is a bug somewhere, but deleting the # only copy of its input is not the way to find out. + # + # A PENDING drop is the exception, and the reason retention matters + # now: /api/uploads/catalog takes files from anyone, and nothing + # about an unreviewed submission ever reaches a terminal state. Left + # out of this sweep it would occupy the volume permanently. continue try: shutil.rmtree(batch_dir(manifest.batch_id)) diff --git a/app/infrastructure/settings.py b/app/infrastructure/settings.py index b915a1e..3386092 100644 --- a/app/infrastructure/settings.py +++ b/app/infrastructure/settings.py @@ -163,6 +163,22 @@ BATCH_QUEUE_MAX = int(os.getenv("BATCH_QUEUE_MAX", "4")) # resource it has. BATCH_RETENTION_DAYS = int(os.getenv("BATCH_RETENTION_DAYS", "7")) +# --- Review inbox ---------------------------------------------------------- +# POST /api/uploads/catalog accepts files with NO credential, so that colleagues +# can send spreadsheets without one being issued to them. Nothing it accepts is +# queued - files wait in the admin review inbox - which removes BATCH_QUEUE_MAX +# as the bound on that endpoint and leaves the volume as the only thing an +# anonymous sender can exhaust. These are that bound; past either, the endpoint +# answers 429 and stages nothing. +# +# Both count only files still AWAITING review. Starting or dismissing a drop +# releases its share immediately, and BATCH_RETENTION_DAYS reclaims whatever +# nobody ever looks at. +INBOX_MAX_PENDING_FILES = int(os.getenv("INBOX_MAX_PENDING_FILES", "200")) +INBOX_MAX_PENDING_BYTES = int( + os.getenv("INBOX_MAX_PENDING_BYTES", str(200 * 1024 * 1024)) +) + # Deliberately false. A batch interrupted by a restart is marked "interrupted" # and waits for someone to press Resume. Auto-resuming would mean a container # stuck in a restart loop re-runs the heaviest work in the app on every boot, diff --git a/app/main.py b/app/main.py index eaf8c45..5a67fd0 100644 --- a/app/main.py +++ b/app/main.py @@ -94,6 +94,20 @@ async def lifespan(_app: FastAPI): interrupted = scan_interrupted() if interrupted: logger.info("Marked %d interrupted catalog batch(es).", len(interrupted)) + + # Retention, which otherwise only runs when the ingestion worker goes + # idle. That was sufficient while every upload was queued work; it is + # not now that /api/uploads/catalog stages drops for review and queues + # nothing. An inbox nobody acts on would never start the worker, so the + # sweep that reclaims abandoned uploads would never run either. + try: + from app.core.batch_ingest import purge_expired + + purged = purge_expired() + if purged: + logger.info("Purged %d expired batch upload(s) at startup.", len(purged)) + except Exception as exc: # noqa: BLE001 - housekeeping must not block boot + logger.warning("Could not purge expired batch uploads: %s", exc) if interrupted and BATCH_AUTO_RESUME: # Opt-in only. Importing the worker here rather than at module # scope keeps the queue and its thread out of a boot that never diff --git a/docs/INGESTION_API.md b/docs/INGESTION_API.md new file mode 100644 index 0000000..3e6934e --- /dev/null +++ b/docs/INGESTION_API.md @@ -0,0 +1,426 @@ +# Sending spreadsheets to the catalogue + +**Base:** `https://mcp.nearle.ai.in` · **Formats:** `.xlsx` `.xls` `.csv` `.tsv` + +Four endpoints accept a spreadsheet. One is **open — no credential needed** — and +stages files for review; three are operator imports that require a credential and +write directly to the store, sales and nutrition tables. + +> **Deployment status.** The review-inbox behaviour in §1 is committed but **not yet +> deployed**. Until it is, `POST /api/uploads/catalog` still answers `401` without a +> credential. Confirm before integrating — see +> [§7 Checking what is deployed](#7-checking-what-is-deployed). + +--- + +## Read this first + +These endpoints used to fill a missing column with a plausible constant. A sheet whose headers didn't match still returned `200 OK` — having written rows attributed to brand `amul` in store `store_mumbai_1` at MRP 100. Nutrition rows were worse: absent nutrients became invented numbers and an absent allergens column became `allergens = ['None']`, stored as `data_status: "verified"`. + +The rule now: **nothing is invented.** A missing value is stored as NULL. A row missing the fields that identify it is skipped and reported back with its row number. An upload where nothing could be imported is a `422`, not a `200` with `rows_imported: 0`. + +A blank allergens cell records *"we were not told"* — which is not the claim *"contains no allergens"*. Leave the column out rather than writing `None`. + +**Arithmetic is not invention.** `line_total` is still derived from `quantity × unit_price` when your sheet omits it, because both were supplied. That distinction is the whole design: compute from given values, never substitute for absent ones. + +--- + +## Authentication + +**The catalogue drop (§1) needs no credential.** Anyone who can reach the host can +send it a spreadsheet. That is only safe because nothing sent there runs on arrival: +files wait in an admin review inbox, so the cost of an unwanted drop is disk until +somebody declines it — never products in the live catalogue. + +The three operator imports (§2) write straight to the database and stay credentialed. +Two credential types are accepted there: + +``` +X-API-Key: # machine clients +Authorization: Bearer # from POST /api/auth/login, 12h lifetime +``` + +Send one or the other. Sending both is not rejected — the bearer token is tried first +and the API key ignored — but relying on that is a way to spend an hour debugging the +wrong credential. + +| Endpoint | Permission | Who holds it | +| --- | --- | --- | +| `POST /api/uploads/catalog` | **none** | anyone | +| `GET /api/uploads/catalog/{batch_id}` | **none** — the id is the credential | anyone holding that id | +| `GET /api/uploads/catalog` (list) | `upload_catalog` | `uploader`, `admin` | +| `/api/upload/stores` | `upload_store_inventory` | `user`, `admin` | +| `/api/upload/analytics` | `manage_analytics` | `admin` | +| `/api/upload/nutrition` | `manage_nutrition` | `admin` | + +`admin` is a superuser: `Principal.has_permission` grants it every permission rather +than requiring each to be listed. On the credentialed endpoints, no credential → `401` +and a valid credential without the permission → `403`. + +A credential that is **present but wrong** is a `401` everywhere, including on the open +drop. Anonymous is allowed; mistyped is not — a typo must not silently downgrade to an +anonymous submission the sender then cannot find. + +--- + +## 1. The catalogue drop + +This is the one to give a colleague. No credential, up to 20 sheets per request, +answers immediately with a `batch_id`. + +**Nothing you send here runs on arrival.** The files are stored and an admin sees them +in the review inbox; only when they select the sheets and press Start does anything +reach the 11-stage pipeline or the catalogue. + +### `POST /api/uploads/catalog` → `202` + +```bash +curl -X POST https://mcp.nearle.ai.in/api/uploads/catalog \ + -F 'files=@store-catalog.xlsx' \ + -F 'files=@second-store.xlsx' \ + -F 'sender=priya' +``` + +| Form field | | Purpose | +| --- | --- | --- | +| `files` | **required** | The spreadsheets. Repeat the field for more than one; up to 20. | +| `sender` | optional | A label for the inbox, so the admin can see who sent what. Free text, trimmed to 60 chars. Defaults to `anonymous`. | + +`use_llm` and `fetch_images` are **no longer accepted here.** They decide how a run +behaves and commit the host to outbound work, so the choice belongs to the admin +pressing Start — not to the sender. Passing them is inert. + +### Which columns are read + +Headers match case-insensitively by keyword, so `Product Name`, `ITEM` and `variant` all resolve to the same field. **A product-name column is the only hard requirement** — a sheet without one is rejected during the request. + +| Field | | Headers that match | +| --- | --- | --- | +| `product_name` | **required** | `product` · `item` · `variant` · `name` | +| `brand` | optional | `brand` · `manufacturer` · `company` | +| `category` | optional | `category` · `segment` | +| `size_variants` | optional | `size` · `pack` · `weight` · `volume` · `net qty` · `quantity` | +| `final_selling_price` | optional | `final price` · `selling price` · `price` · `mrp` · `rate` · `cost` | +| `barcode` | optional | `barcode` · `bar code` · `gtin` · `ean` · `upc` | +| `fssai_license` | optional | any header containing `fssai` | +| `hsn_code` | optional | any header containing `hsn` | +| `image_url` | optional | `image` · `photo` · `picture` · `url` · `link` | +| `description` | optional | `description` · `desc` · `detail` | +| `title` | optional | `title` | +| `product_sku` | optional | `sku` | +| `price_range` | optional | `price range` · `range` | +| `providers` | optional | `provider` · `platform` · `marketplace` · `available at` | +| `highlights` | optional | `highlight` · `feature` · `benefit` | +| `nutrients` | optional | `nutrient` · `nutrition` | + +Rules are evaluated in the order above, so a header matching two of them takes the +earlier field. Where two *columns* map to the same field, the **leftmost wins** and the +other is reported as ignored. Unrecognised columns are listed back, not an error. + +### Response — `202`, one good file and one bad + +The bad file is kept as a failed member rather than dropped, so a sender who submitted two files and sees one knows what happened to the other. + +```jsonc +{ + "batch_id": "49a82536866a483a9189954d3c749243", + "status": "pending", + "detail": "Waiting for review. Nothing runs until an admin starts it.", + "submitted_by": "priya", + "created_at": 1756370000.0, + "updated_at": 1756370000.0, + "files_total": 2, + "files_done": 0, + "files_failed": 1, + "current_file": null, + "use_llm": false, + "fetch_images": false, + "totals": { "rows_total": 0, "products_built": 0, "inserted": 0, + "backfilled": 0, "skipped_existing": 0, "rejected": 0 }, + "brands": [], + "files": [ + { "index": 0, "filename": "catalog.csv", "status": "queued", + "total_stages": 11, "rows_total": 1, "size_bytes": 57 }, + { "index": 1, "filename": "notes.txt", "status": "failed", + "detail": "The file has no data rows." } + ], + "message": "1 file(s) received and waiting for review. 1 could not be read - see 'files' for the reason on each, and resend those." +} +``` + +Note the two levels of status. The **batch** is `pending` — waiting on a person. An +individual **file** inside it reads `queued`, meaning it is intact and eligible to be +run; it is waiting behind a decision, not behind the worker. + +### Polling + +`GET /api/uploads/catalog/{batch_id}` — the same body, live. No credential: the +`batch_id` is a `uuid4` handed only to whoever sent the drop, so holding it is the +proof of having sent it. An id nobody issued is a `404`. + +Expect `pending` until an admin acts. After that the batch you sent stops existing +under that id — the files move into a run with an id of its own, and your drop +reports `404`. **A `404` after `pending` means your files were accepted and started,** +not that they were lost. Ask the admin for the run id if you need to follow it. + +`GET /api/uploads/catalog?limit=20` — your recent submissions, newest first, under +`{"batches": [...]}`. This one **does** need a credential: letting an anonymous caller +enumerate every sender's drops is a different thing from letting one check the id they +hold. + +### Status values + +| Status | Meaning | +| --- | --- | +| `pending` | In the review inbox. Nothing has run. | +| `queued` | Approved and waiting for the worker | +| `running` | In the pipeline now | +| `done` | Every file completed | +| `partial` | Some landed, some failed. Deliberately not `done` — four of five succeeding must not read as flat success | +| `failed` | No file completed | +| `interrupted` | A restart cut a run short. An admin resumes it; it never auto-restarts | +| `cancelled` | Stopped by an admin before the remaining files began | + +### The eleven stages + +1. Brand Resolution & FSSAI Licence Mapping +2. Row Intake & Normalisation +3. Title & Category Consistency Guard +4. Pack-Size Variant Explosion & Unit Safety +5. Pricing Band Estimation +6. Image Search & Contamination Filtering +7. Marketplace & Internal SKU Resolution +8. Barcode Retrieval & Enrichment +9. HSN / GST Tax Enrichment +10. Deterministic Product Validation Gate +11. Vector Embedding & Storage + +### Limits + +| Limit | Value | Env var | Exceeded | +| --- | --- | --- | --- | +| Files per request | 20 | `BATCH_MAX_FILES` | `413` | +| Bytes per file | 10 MB | — | `413` | +| Rows per file | 2,000 | — | file marked `failed`, others accepted | +| Bytes per request | 50 MB | `BATCH_MAX_TOTAL_BYTES` | `413` | +| Rows per request | 20,000 | `BATCH_MAX_TOTAL_ROWS` | `413` | +| Files awaiting review | 200 | `INBOX_MAX_PENDING_FILES` | `429` — nothing stored | +| Bytes awaiting review | 200 MB | `INBOX_MAX_PENDING_BYTES` | `429` — nothing stored | + +The last two are the ceiling on the inbox as a whole, across every sender. They count +only what is still *awaiting review*, so starting or dismissing a drop frees its share +immediately. A `429` here stores nothing — resend once an admin has cleared space. + +Drops nobody acts on are deleted after `BATCH_RETENTION_DAYS` (7). + +--- + +## 2. Operator imports + +Three endpoints that write straight to the store, sales and nutrition tables — no pipeline, no review, no polling. Synchronous, with a per-row report. All three share one response shape. **Form field name is `file`** (singular). + +Each also answers at an `/upload` alias — `/api/upload/stores/upload`, `/api/upload/analytics/upload`, `/api/upload/nutrition/upload` — carrying the same guard. Prefer the short form. + +```jsonc +// success +{ + "status": "success", + "filename": "sheet.csv", + "rows_total": 1, + "rows_imported": 1, + "rows_skipped": 0, + "errors": [], + "message": "Successfully imported 1 store inventory items." +} + +// partial — row 3 had no brand +{ + "status": "partial", + "rows_total": 2, + "rows_imported": 1, + "rows_skipped": 1, + "errors": [ + { "row": 3, "error": "missing required field(s): brand" } + ], + "message": "Imported 1 store inventory items; skipped 1 row(s) that were missing required fields - see 'errors'." +} +``` + +`row` is the number your spreadsheet shows — 1-based, counting the header — so row 3 is the second data row. At most **50** errors are returned; the rest are still skipped. + +### `POST /api/upload/stores` + +| Column | | When absent | +| --- | --- | --- | +| `store_id` | **required** | Row skipped and reported | +| `brand` | **required** | Row skipped and reported (lowercased on write) | +| `product_name` | **required** | Row skipped and reported | +| `image_id` | optional | Derived from brand + product name, so a re-upload lands on the same key | +| `store_name` | optional | Falls back to the `store_id`, underscores replaced and title-cased | +| `category`, `city` | optional | Stored NULL | +| `available_stock`, `reserved_stock`, `reorder_level`, `safety_stock` | optional | Stored `0` — the schema's own default; columns are NOT NULL | +| `mrp`, `cost_price`, `selling_price` | **all three** | Prices written only when all three present. Inventory still imports; counted in `prices_skipped` | + +Extra response fields: `stores_affected`, `prices_written`, `prices_skipped`. + +### `POST /api/upload/analytics` + +> **This endpoint never worked before.** Its insert named a `total_price` column that doesn't exist on `order_items` (the column is `line_total`) and omitted the NOT NULL `store_id`. Every call raised and returned `500 Database import failed`. If you have an integration that gave up on it, retry. + +| Column | | When absent | +| --- | --- | --- | +| `store_id`, `brand`, `customer_id`, `quantity`, `unit_price` | **required** | Row skipped and reported | +| `image_id` **or** `product_name` | one of | Row skipped. `image_id` derived from the name when only the name is given | +| `order_id` | optional | Generated. Each row then becomes its own order | +| `total_price` | optional | Derived as `quantity × unit_price`. Accepted under this name; stored in `line_total` | +| `order_date` | optional | Import time, counted in `dates_defaulted_to_now` so you can see how much of the file isn't really dated | +| `payment_method`, `delivery_status` | optional | Stored NULL | + +Extra response fields: `total_revenue`, `dates_defaulted_to_now`. + +`customer_id` is required rather than defaulted: it used to fall back to a shared `cust_imported`, which merges every buyer in a file into one customer and corrupts exactly the per-customer models this table feeds. Rows sharing an `order_id` become one order, whose `order_value` is recomputed as the **sum of its lines**. + +**Import each file once.** `order_items` has a surrogate key and nothing to conflict on, so re-uploading appends the lines again. + +### `POST /api/upload/nutrition` + +Only `brand`, plus one of `image_id` / `product_name`, is required. Every nutrient column is optional and stays NULL when omitted. Recognised: `calories_kcal`, `protein_g`, `carbohydrates_g`, `total_sugar_g`, `dietary_fiber_g`, `total_fat_g`, `sodium_mg`, `calcium_mg`, `iron_mg`, `vitamin_c_mg`, `health_score`, `diet_tags`, `allergens` (plus `category`). + +What you send decides the row's `data_status`, which every downstream reader must check before treating a NULL as zero: + +| Status | Awarded when | +| --- | --- | +| `verified` | All seven core nutrients present — calories, protein, carbohydrates, sugar, fibre, fat, sodium | +| `partial` | Some of the seven present | +| `unavailable` | None present. Row still created, carrying only its identity | + +```jsonc +// two of seven nutrients +{ "rows_imported": 1, + "data_status_counts": { "verified": 0, "partial": 1, "unavailable": 0 } } + +// full core set +{ "rows_imported": 1, + "data_status_counts": { "verified": 1, "partial": 0, "unavailable": 0 } } +``` + +> **Never write `"None"` in the allergens column.** An omitted allergens column stores NULL with `allergen_source: "unavailable"`. A literal `None` stores the string as an allergen name. The two must stay distinguishable: one says nobody told us, the other claims the product is allergen-free — and that claim is served by the public `/api/nutrition` endpoints. + +Upserts use `COALESCE` throughout, so a later, thinner sheet can never blank a value an earlier trusted source established, and a row already `verified` is never downgraded by a `partial` upload. + +--- + +## 3. Errors + +| Code | Cause | What to do | +| --- | --- | --- | +| `400` | File unparseable, no data rows, or (catalogue drop) no product-name column in any file | Read `detail`; it names the headers it found | +| `401` | Credential present but invalid; or absent on an endpoint that needs one | Check the table in Authentication | +| `403` | Valid credential, wrong permission | Check the role table | +| `413` | Over a file-count, byte or row ceiling | Split the drop | +| `422` | **Operator imports only:** not one row could be imported | `detail.errors` lists row numbers and what each was missing | +| `429` | **Catalogue drop:** the review inbox is full | Nothing was stored. Ask an admin to clear it, then resend | + +```jsonc +// 400 — wrong columns +{ "detail": "None of the uploaded files could be ingested. wrong.csv: No product name column was found. Headers read: colour, size. Recognised fields: size_variants." } + +// 422 — nothing importable +{ "detail": { + "message": "No store inventory items could be imported from 'sheet.csv'. Check the column headers against the sample template (GET /api/upload/template/...).", + "rows_total": 1, + "errors": [ { "row": 2, "error": "missing required field(s): store_id, brand" } ] } } +``` + +Note the shape difference: on `400`, `413` and `429`, `detail` is a **string**; on `422` it is an **object**. A client that renders `detail` directly will print `[object Object]` for a 422 unless it handles both. + +Sample sheets with the correct headers: `GET /api/upload/template/stores`, `/analytics`, `/nutrition` — no credential needed. + +--- + +## 4. If you already integrated + +| Was | Is now | +| --- | --- | +| `POST /api/uploads/catalog` needed an `X-API-Key` | No credential. Send the file; drop the header | +| Files ran on arrival | Files wait for review. Expect `pending`, not `queued` | +| `?use_llm` / `?fetch_images` on the drop | Ignored. The admin chooses at Start | +| `429` meant the worker queue was full | `429` now means the review inbox is full | +| `200` with `rows_imported: 0` | `422` with per-row reasons. Handle as a client error, not a server one | +| Missing `customer_id` became `cust_imported` | Row is skipped. Supply a real customer id | +| Missing `cost_price` invented as `mrp × 0.7` | Prices skipped for that row; inventory still lands. Send all three | +| Missing nutrients and allergens got plausible defaults | Stored NULL, and the row is `partial`/`unavailable` rather than `verified` | + +--- + +## 5. API keys are now optional + +Nobody needs a key to send catalogue spreadsheets. Issue one only for a machine client +that wants the credentialed reads (`GET /api/uploads/catalog`) or the operator imports. + +```bash +python scripts/make_auth_secrets.py --api-key catalog-drop:uploader +``` + +Put the resulting `name:role:secret` triple in the **Dokploy Environment tab**, not in +`.env.production`, and restart the service. Two reasons, both recorded in +`.env.production`'s own comments: + +1. `.env.production` is committed. A per-consumer key is the one credential that gets + issued and revoked often, and it does not belong in git. +2. `backend/Dockerfile` does `COPY .env.production .env`, so a value there is baked at + **build** time — issuing or revoking would mean rebuilding an image that installs CPU + torch, a build that has already failed once on disk space. `settings.py` calls + `load_dotenv()` **without** `override=True`, so the process environment wins. + +Container environment is fixed at creation, so the service must be **recreated**, not +merely restarted. `/api/health` will then report `api_keys_source: "process-env"`. + +Constraints enforced at boot, before any request is served: + +- Format `name:role:secret`, comma-separated between entries. +- `role` is `admin`, `user` or `uploader`. Keys never expire — treat one as a long-lived + secret and rotate it deliberately. Two entries can be live at once, which is how you + rotate without a cutover window. +- The secret must be at least **32 characters**, because `/api/health` publishes a digest. +- **Name the key for its function, not the person holding it.** `/api/health` is public + and reports `{name, role, fingerprint}` for every configured key. + +--- + +## 6. What an open drop costs + +Disk, and nothing else, until somebody looks at it. The endpoint queues nothing, so it +cannot occupy the ingestion worker and cannot reach the catalogue on its own — the +review gate is what makes accepting files from anyone acceptable. + +`INBOX_MAX_PENDING_FILES` and `INBOX_MAX_PENDING_BYTES` bound what unreviewed +submissions can occupy; past either the endpoint answers `429` and stores nothing. +Dismissing a drop deletes its bytes immediately, and retention reclaims anything nobody +ever looks at. + +If the host is reachable from the open internet, consider an IP allow-list at the proxy +as defence in depth. The controls above bound the damage; they do not stop a stranger +from filling the inbox with sheets an admin then has to decline. + +--- + +## 7. Checking what is deployed + +`GET /api/health` is public and answers `200` even when dependencies are down. + +```bash +curl -s https://mcp.nearle.ai.in/api/health +``` + +- **`POST /api/uploads/catalog` with no credential and no file returns `400`/`422`** → + the open drop is live. A `401` means the old, credentialed build is still running. +- **`auth.api_keys_count` / `api_keys` / `api_keys_source` present** → the build + includes the auth diagnostics. Absent → the deployment predates them. +- **`auth.api_keys[].fingerprint`** answers *"is my key on this deployment?"* without + anyone sending the secret — the question a `401` cannot answer, since an undeployed + key and a wrong key fail identically. Compare against + `scripts/make_auth_secrets.py --fingerprint`. + +`status: "degraded"` with `ollama: false` is expected in production: `USE_OLLAMA` is +`false` there. diff --git a/tests/test_review_inbox.py b/tests/test_review_inbox.py new file mode 100644 index 0000000..67191ca --- /dev/null +++ b/tests/test_review_inbox.py @@ -0,0 +1,339 @@ +"""The admin half of the review inbox. + + GET /api/admin/catalog-batch/inbox - drops awaiting review + POST /api/admin/catalog-batch/from-inbox - run selected files + POST /api/admin/catalog-batch/inbox/dismiss - discard selected files + +The sender half lives in test_uploads_api.py. What matters here is the gate +between them: a drop arrives with no credential and runs only when an admin +says so, and every assertion below is ultimately about one of two failures - + +* something running that nobody approved, and +* a file being lost, or run twice, on its way out of the inbox. + +The response SHAPES are asserted literally rather than loosely, because +frontend/src/pages/InboxPanel.jsx was written against this contract before the +backend existed. `submission_id`, `file_id`, `pending_count` and the `dismissed` +count are read by name there; a rename that only this file catches is cheap, and +one that nothing catches is a blank admin tab with no error. +""" +from __future__ import annotations + +import io + +import pytest + +from app.api.batch_job_store import batch_job_store +from app.core import batch_ingest + +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" + +HEADERS = ["Product Name", "Category", "Brand"] +ROWS = [["Amul Butter 100g", "Butter", "Amul"]] + + +def _csv(rows=ROWS) -> bytes: + return ("\n".join([",".join(HEADERS)] + [",".join(r) for r in rows])).encode() + + +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 batch_root(tmp_path, monkeypatch): + root = tmp_path / "batch_uploads" + monkeypatch.setattr(batch_ingest, "BATCH_UPLOAD_DIR", root) + return root + + +@pytest.fixture(autouse=True) +def submitted(monkeypatch): + """Stub the worker and record what reached it. + + Turns "did anything run?" into a direct assertion. A real worker would also + resolve BATCH_UPLOAD_DIR again after teardown and write into the repo's + data/ directory - see test_batch_catalog_ingest. + """ + from app.core import batch_worker + + seen: list = [] + monkeypatch.setattr(batch_worker, "submit", seen.append) + return seen + + +@pytest.fixture(autouse=True) +def _clean_job_store(): + batch_job_store._batches.clear() + batch_job_store._cancelled.clear() + yield + batch_job_store._batches.clear() + batch_job_store._cancelled.clear() + + +@pytest.fixture +def upload_headers(monkeypatch) -> dict: + """An `uploader` key, for asserting what it may NOT do here. + + Patched on `security`, not on `settings`: security.py binds the dict object + at import, so rebinding the name in `settings` leaves principal_for_api_key + looking at the original and every request comes back 401. + """ + from app.infrastructure import security + + secret = "k" * 43 + monkeypatch.setattr(security, "API_KEYS", {secret: ("catalog-drop", "uploader")}) + return {"X-API-Key": secret} + + +def _drop(client, sender: str, *names: str) -> str: + """Post a drop the way a colleague would - no credential - and return its id.""" + response = client.post( + UPLOAD, + files=_files(*[(n, _csv()) for n in (names or ("a.csv",))]), + data={"sender": sender}, + ) + assert response.status_code == 202, response.text + return response.json()["batch_id"] + + +# --------------------------------------------------------------------------- +# Listing +# --------------------------------------------------------------------------- +def test_a_drop_appears_in_the_inbox(client, admin_headers): + batch_id = _drop(client, "priya", "catalog.csv") + + body = client.get(INBOX, headers=admin_headers).json() + + assert body["pending_count"] == 1 + assert len(body["submissions"]) == 1 + submission = body["submissions"][0] + assert submission["submission_id"] == batch_id + assert submission["submitted_by"] == "priya" + assert submission["created_at"] > 0 + assert submission["files"][0]["file_id"] == f"{batch_id}:0" + assert submission["files"][0]["filename"] == "catalog.csv" + assert submission["files"][0]["rows_total"] == 1 + assert submission["files"][0]["size_bytes"] > 0 + + +def test_an_empty_inbox_is_empty_not_an_error(client, admin_headers): + body = client.get(INBOX, headers=admin_headers).json() + assert body == {"pending_count": 0, "submissions": []} + + +def test_the_inbox_survives_a_restart(client, admin_headers): + """Read from disk, not from the in-process cache. + + The job store is a module singleton populated by what this process has seen. + An inbox backed by it would empty itself on every deploy, which looks exactly + like a colleague's files having been dealt with. + """ + _drop(client, "priya") + batch_job_store._batches.clear() + + assert client.get(INBOX, headers=admin_headers).json()["pending_count"] == 1 + + +def test_a_rejected_file_is_not_offered_for_selection(client, admin_headers): + """A file the upload already refused is on the manifest so the sender can be + told, but it has no bytes and cannot be run.""" + response = client.post( + UPLOAD, files=_files(("good.csv", _csv()), ("bad.txt", b"nope")), + ) + assert response.status_code == 202, response.text + + body = client.get(INBOX, headers=admin_headers).json() + names = [f["filename"] for f in body["submissions"][0]["files"]] + assert names == ["good.csv"] + assert body["pending_count"] == 1 + + +def test_a_restart_does_not_relabel_the_inbox_interrupted(client, admin_headers): + """scan_interrupted() sweeps every non-terminal manifest. + + Without its PENDING exemption, every boot would mark the whole inbox + "interrupted" and offer an admin a Resume button for work nobody started. + """ + _drop(client, "priya") + + assert batch_ingest.scan_interrupted() == [] + assert client.get(INBOX, headers=admin_headers).json()["pending_count"] == 1 + + +def test_pending_drops_are_not_listed_as_runs(client, admin_headers): + """The Batch tab lists runs. A drop has never been near the pipeline, and a + row there with no progress and a Resume button would be a lie.""" + _drop(client, "priya") + + batches = client.get("/api/admin/catalog-batch/batches", + headers=admin_headers).json()["batches"] + assert batches == [] + + +# --------------------------------------------------------------------------- +# Starting +# --------------------------------------------------------------------------- +def test_starting_a_file_queues_it_and_empties_the_inbox(client, admin_headers, submitted): + batch_id = _drop(client, "priya", "catalog.csv") + + started = client.post(FROM_INBOX, json={"file_ids": [f"{batch_id}:0"]}, + headers=admin_headers) + + assert started.status_code == 202, started.text + run = started.json() + assert run["files_total"] == 1 + assert run["submitted_by"] == "priya" + # A NEW id: the run is a distinct thing from the drop it came from. + assert run["batch_id"] != batch_id + assert submitted == [run["batch_id"]] + + assert client.get(INBOX, headers=admin_headers).json()["pending_count"] == 0 + + +def test_files_from_two_drops_become_one_batch(client, admin_headers, submitted): + """The reason selection is per file rather than per drop.""" + first = _drop(client, "priya", "monday.csv") + second = _drop(client, "arun", "today.csv") + + started = client.post( + FROM_INBOX, + json={"file_ids": [f"{first}:0", f"{second}:0"]}, + headers=admin_headers, + ) + + assert started.status_code == 202, started.text + run = started.json() + assert run["files_total"] == 2 + assert sorted(f["filename"] for f in run["files"]) == ["monday.csv", "today.csv"] + # Both senders are preserved - the question gets asked precisely when a + # batch was assembled out of several drops. + assert run["submitted_by"] == "arun, priya" + assert len(submitted) == 1 + assert client.get(INBOX, headers=admin_headers).json()["pending_count"] == 0 + + +def test_starting_one_file_leaves_its_siblings_in_the_inbox(client, admin_headers): + batch_id = _drop(client, "priya", "one.csv", "two.csv") + + client.post(FROM_INBOX, json={"file_ids": [f"{batch_id}:0"]}, headers=admin_headers) + + body = client.get(INBOX, headers=admin_headers).json() + assert body["pending_count"] == 1 + remaining = body["submissions"][0]["files"][0] + assert remaining["filename"] == "two.csv" + # Its id is unchanged. The UI holds selections across polls, so renumbering + # here would silently repoint a tick at a different file. + assert remaining["file_id"] == f"{batch_id}:1" + + +def test_starting_the_same_file_twice_is_refused(client, admin_headers, submitted): + """A stale checkbox in a second tab must not run a sheet again.""" + batch_id = _drop(client, "priya") + payload = {"file_ids": [f"{batch_id}:0"]} + + assert client.post(FROM_INBOX, json=payload, headers=admin_headers).status_code == 202 + second = client.post(FROM_INBOX, json=payload, headers=admin_headers) + + assert second.status_code == 409 + assert "already" in second.json()["detail"].lower() + assert len(submitted) == 1 + + +def test_the_admin_chooses_the_network_stages(client, admin_headers): + """Not the sender - see the sender-side assertion in test_uploads_api.""" + batch_id = _drop(client, "priya") + + started = client.post( + FROM_INBOX, + json={"file_ids": [f"{batch_id}:0"], "fetch_images": True, "use_llm": True}, + headers=admin_headers, + ) + + assert started.json()["fetch_images"] is True + assert started.json()["use_llm"] is True + + +def test_starting_nothing_is_refused(client, admin_headers, submitted): + assert client.post(FROM_INBOX, json={"file_ids": []}, + headers=admin_headers).status_code == 400 + assert submitted == [] + + +# --------------------------------------------------------------------------- +# Dismissing +# --------------------------------------------------------------------------- +def test_dismissing_deletes_the_file_and_its_bytes(client, admin_headers, batch_root): + batch_id = _drop(client, "priya") + assert list((batch_root / batch_id).glob("*.csv")) + + result = client.post(DISMISS, json={"file_ids": [f"{batch_id}:0"]}, + headers=admin_headers) + + assert result.status_code == 200, result.text + assert result.json() == {"dismissed": 1} + assert client.get(INBOX, headers=admin_headers).json()["pending_count"] == 0 + # The disposal path on an endpoint anyone can post to. If the bytes survived + # a dismiss, the volume would fill with sheets that were already refused. + assert not (batch_root / batch_id).exists() + + +def test_dismissing_cannot_touch_a_running_batch(client, admin_headers, monkeypatch): + """A crafted id must not delete files out of a batch that is mid-run.""" + from app.api import batch_common + + manifest, _started = batch_common.stage_and_queue( + [("live.csv", _csv(), 1)], [], use_llm=False, fetch_images=False, + ) + assert manifest.status != batch_ingest.PENDING + + result = client.post(DISMISS, json={"file_ids": [f"{manifest.batch_id}:0"]}, + headers=admin_headers) + + assert result.json() == {"dismissed": 0} + assert batch_ingest.read_manifest(manifest.batch_id).files + + +def test_a_malformed_id_does_not_fail_the_whole_request(client, admin_headers): + """The UI polls every five seconds and prunes selections against what came + back, so one stale tick must not block the files that are still there.""" + batch_id = _drop(client, "priya") + + result = client.post( + DISMISS, + json={"file_ids": ["nonsense", f"{batch_id}:0", "also:bad:x"]}, + headers=admin_headers, + ) + + assert result.status_code == 200, result.text + assert result.json() == {"dismissed": 1} + + +# --------------------------------------------------------------------------- +# Who may reach it +# --------------------------------------------------------------------------- +def test_the_inbox_is_admin_only(client, upload_headers): + """An uploader key sends files; it does not get to decide what runs.""" + for method, path, kwargs in ( + ("get", INBOX, {}), + ("post", FROM_INBOX, {"json": {"file_ids": []}}), + ("post", DISMISS, {"json": {"file_ids": []}}), + ): + response = getattr(client, method)(path, headers=upload_headers, **kwargs) + assert response.status_code == 403, f"{method.upper()} {path} -> {response.status_code}" + + +def test_the_inbox_is_closed_to_anonymous(client): + assert client.get(INBOX).status_code == 401 + assert client.post(FROM_INBOX, json={"file_ids": []}).status_code == 401 + assert client.post(DISMISS, json={"file_ids": []}).status_code == 401 diff --git a/tests/test_uploads_api.py b/tests/test_uploads_api.py index a75b871..7e89e65 100644 --- a/tests/test_uploads_api.py +++ b/tests/test_uploads_api.py @@ -153,11 +153,14 @@ def store(monkeypatch): # --------------------------------------------------------------------------- # The requirement: a file sent here is ingested # --------------------------------------------------------------------------- -def test_an_uploaded_spreadsheet_is_queued_for_ingestion(client, upload_headers, submitted): +def test_an_uploaded_spreadsheet_waits_for_review(client, upload_headers, submitted): """THE regression test for this endpoint. - An earlier version parked the file and started nothing, which is - indistinguishable from working if you only assert on the status code. + It must accept the file AND NOT RUN IT. Asserting only on the 202 cannot + tell the two apart, and both directions have been wrong here before: the + endpoint once parked files and started nothing while reporting success, and + it later ran everything the moment it arrived - which is unacceptable now + that anybody can post to it. """ response = client.post(UPLOAD, files=_files(("catalog.xlsx", _sheet())), headers=upload_headers) @@ -165,10 +168,11 @@ def test_an_uploaded_spreadsheet_is_queued_for_ingestion(client, upload_headers, assert response.status_code == 202, response.text body = response.json() assert body["batch_id"] - assert body["status"] == "queued" + assert body["status"] == "pending" assert body["files_total"] == 1 - # The batch reached the worker. This is the assertion that was missing. - assert submitted == [body["batch_id"]] + # Nothing reached the worker. This is the assertion the review gate rests on. + assert submitted == [] + assert "review" in body["message"].lower() def test_an_excel_file_reaches_the_catalog(client, upload_headers, store): @@ -211,24 +215,74 @@ def test_network_stages_are_off_unless_asked_for(client, upload_headers): assert response.json()["use_llm"] is False -def test_network_stages_can_be_opted_into(client, upload_headers): - response = client.post(UPLOAD + "?fetch_images=true", files=_files(("a.csv", _csv())), - headers=upload_headers) - assert response.json()["fetch_images"] is True +def test_a_sender_cannot_opt_into_the_network_stages(client, upload_headers): + """The sender does not get to commit the host to Playwright image search. + + `fetch_images` and `use_llm` moved to the admin who presses Start. Passing + them here must be inert rather than honoured - on an endpoint open to + anonymous callers, an accepted query parameter is an invitation. + """ + response = client.post(UPLOAD + "?fetch_images=true&use_llm=true", + files=_files(("a.csv", _csv())), headers=upload_headers) + assert response.status_code == 202, response.text + assert response.json()["fetch_images"] is False + assert response.json()["use_llm"] is False # --------------------------------------------------------------------------- # Who may call it # --------------------------------------------------------------------------- -def test_upload_endpoint_rejects_anonymous(client): - assert client.post(UPLOAD, files=_files(("a.csv", _csv()))).status_code == 401 +def test_upload_endpoint_accepts_anonymous(client, submitted): + """The point of the endpoint: a colleague needs no credential to send a file. + + Safe only because of the assertion below it - the drop is staged for review, + not run. If that ever changes, this test is the one that has to be argued + with first. + """ + response = client.post(UPLOAD, files=_files(("a.csv", _csv()))) + assert response.status_code == 202, response.text + body = response.json() + assert body["status"] == "pending" + assert body["submitted_by"] == "anonymous" + assert submitted == [] + + +def test_a_sender_label_is_recorded_when_given(client): + """The inbox groups by sender; three colleagues all reading "anonymous" is + an inbox nobody can triage.""" + response = client.post(UPLOAD, files=_files(("a.csv", _csv())), + data={"sender": "priya"}) + assert response.status_code == 202, response.text + assert response.json()["submitted_by"] == "priya" + + +def test_a_credential_outranks_the_sender_label(client, upload_headers): + """`sender` is free text. A real credential is not, so it wins.""" + response = client.post(UPLOAD, files=_files(("a.csv", _csv())), + data={"sender": "someone-else"}, headers=upload_headers) + assert response.status_code == 202, response.text + assert response.json()["submitted_by"] == "catalog-drop" + + +def test_a_wrong_credential_is_still_refused(client): + """Anonymous is allowed; WRONG is not. A typo in a key must not silently + downgrade to an anonymous drop that the sender then cannot find.""" + response = client.post(UPLOAD, files=_files(("a.csv", _csv())), + headers={"X-API-Key": "n" * 43}) + assert response.status_code == 401 def test_admin_can_also_send_files(client, admin_headers, submitted): - """`admin` passes every permission check, so this path stays open to them.""" + """`admin` reaches it too - and lands in the same review queue. + + An admin who wants files to run immediately has + POST /api/admin/catalog-batch/ingest. This endpoint has one behaviour for + everyone, so the gate cannot be stepped around by whoever holds a key. + """ response = client.post(UPLOAD, files=_files(("a.csv", _csv())), headers=admin_headers) assert response.status_code == 202, response.text - assert submitted == [response.json()["batch_id"]] + assert response.json()["status"] == "pending" + assert submitted == [] def test_uploader_key_is_refused_on_every_admin_route(client, upload_headers): @@ -300,9 +354,12 @@ def test_a_bad_file_among_good_ones_is_reported_not_fatal(client, upload_headers ) assert response.status_code == 202, response.text body = response.json() - assert submitted == [body["batch_id"]] + assert submitted == [] by_name = {f["filename"]: f for f in body["files"]} + # The FILE is queued inside a batch that is pending: it is waiting its turn + # behind a decision, not behind the worker. `settle()` is what keeps the two + # apart - without its PENDING guard this batch would report itself queued. assert by_name["good.xlsx"]["status"] == "queued" assert by_name["bad.txt"]["status"] == "failed" assert by_name["bad.txt"]["detail"] @@ -367,27 +424,48 @@ def test_uploaded_filenames_cannot_escape_the_batch_directory(client, upload_hea assert not (batch_root.parent / "evil.csv").exists() -def test_a_full_queue_is_refused_but_the_batch_is_saved(client, upload_headers, monkeypatch): - """The caller has no Resume button, so they are told to retry - but the - files are on disk and an admin can start them.""" - import queue +def test_a_full_inbox_is_refused_and_stores_nothing(client, upload_headers, monkeypatch): + """The bound that replaced BATCH_QUEUE_MAX on this endpoint. - from app.core import batch_worker + Nothing here is queued any more, so the worker queue no longer limits what + an anonymous caller can leave on the volume. The inbox ceiling does, and a + refusal has to store NOTHING - a 429 that still wrote the file would be no + limit at all. + """ + from app.api.routers import uploads - def full(_batch_id): - raise queue.Full() + monkeypatch.setattr(uploads, "INBOX_MAX_PENDING_FILES", 1) - monkeypatch.setattr(batch_worker, "submit", full) - response = client.post(UPLOAD, files=_files(("a.csv", _csv())), headers=upload_headers) + first = client.post(UPLOAD, files=_files(("a.csv", _csv())), headers=upload_headers) + assert first.status_code == 202, first.text - assert response.status_code == 429 - detail = response.json()["detail"] - assert "retry" in detail.lower() - # The id is in the message precisely so an admin can find it. - staged = batch_ingest.list_manifests() - assert len(staged) == 1 - assert staged[0].batch_id in detail - assert staged[0].status == "queued" + second = client.post(UPLOAD, files=_files(("b.csv", _csv())), headers=upload_headers) + assert second.status_code == 429 + assert "full" in second.json()["detail"].lower() + + # Still exactly the one drop. The refused file was not written. + assert len(batch_ingest.list_manifests()) == 1 + + +def test_dismissing_a_drop_frees_inbox_capacity(client, upload_headers, admin_headers, + monkeypatch): + """The ceiling counts what is AWAITING review, so triage releases it.""" + from app.api.routers import uploads + + monkeypatch.setattr(uploads, "INBOX_MAX_PENDING_FILES", 1) + + first = client.post(UPLOAD, files=_files(("a.csv", _csv())), headers=upload_headers) + batch_id = first.json()["batch_id"] + assert client.post(UPLOAD, files=_files(("b.csv", _csv())), + headers=upload_headers).status_code == 429 + + dismissed = client.post("/api/admin/catalog-batch/inbox/dismiss", + json={"file_ids": [f"{batch_id}:0"]}, headers=admin_headers) + assert dismissed.status_code == 200, dismissed.text + assert dismissed.json()["dismissed"] == 1 + + assert client.post(UPLOAD, files=_files(("b.csv", _csv())), + headers=upload_headers).status_code == 202 # --------------------------------------------------------------------------- @@ -465,6 +543,25 @@ def test_a_traversing_batch_id_is_not_found(client, upload_headers): assert client.get(f"{UPLOAD}/..%2F..%2Fetc", headers=upload_headers).status_code == 404 -def test_reads_require_a_credential(client): +def test_the_batch_list_requires_a_credential(client): + """The LIST stays closed even though the drop is open. + + Letting an anonymous caller enumerate every sender's submissions is a + different thing entirely from letting one check the id they were handed. + """ assert client.get(UPLOAD).status_code == 401 - assert client.get(f"{UPLOAD}/anything").status_code == 401 + + +def test_an_anonymous_sender_can_poll_the_id_they_were_given(client): + """The id IS the credential here - the sender needed none to post, so + requiring one to read the outcome would strand them.""" + posted = client.post(UPLOAD, files=_files(("a.csv", _csv())), data={"sender": "priya"}) + batch_id = posted.json()["batch_id"] + + polled = client.get(f"{UPLOAD}/{batch_id}") + assert polled.status_code == 200, polled.text + assert polled.json()["status"] == "pending" + + # An id nobody issued is a 404, not a 401: whether it exists is not the + # caller's business, so the two answers must be indistinguishable. + assert client.get(f"{UPLOAD}/anything").status_code == 404