ingestion updates

This commit is contained in:
sriram
2026-08-28 08:57:10 +05:30
parent a54bd43f8b
commit 0e75d32f61
9 changed files with 1147 additions and 4 deletions

5
.gitignore vendored
View File

@@ -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/

View File

@@ -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))

212
app/api/routers/uploads.py Normal file
View File

@@ -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,
)

View File

@@ -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,
)

348
app/core/inbox.py Normal file
View File

@@ -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

View File

@@ -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)

View File

@@ -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.<domain>. 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

View File

@@ -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,

390
tests/test_inbox_upload.py Normal file
View File

@@ -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")