sheet upload fix

This commit is contained in:
Suriyakumarvijayanayagam
2026-08-28 11:22:53 +05:30
parent 0e75d32f61
commit 52f5d3be1d
41 changed files with 2056 additions and 1507 deletions

268
app/api/batch_common.py Normal file
View File

@@ -0,0 +1,268 @@
"""Shared machinery for the two routes that start a catalog batch.
app/api/routers/batch_catalog.py POST /api/admin/catalog-batch/ingest
app/api/routers/uploads.py POST /api/uploads/catalog
Both accept spreadsheets, both run the same 11 stages over them, and both hand
back a batch id to poll. What differs is only who may call them and what the
caller is allowed to see afterwards - so everything between "read the upload"
and "queue the batch" lives here instead of being written twice and drifting.
WHY THE LIMITS ARE ARGUMENTS RATHER THAN IMPORTS
------------------------------------------------
`read_uploads` and `parse_all` take an `UploadLimits` instead of reading the
settings themselves. Each router builds one from ITS OWN module globals, at
call time, which is what keeps
monkeypatch.setattr(batch_catalog, "BATCH_MAX_FILES", 2)
working - the idiom the existing suite is written in. Had this module read the
settings directly, those patches would become silently inert and a limit test
that no longer exercises its limit would still pass.
"""
from __future__ import annotations
import logging
import queue
from dataclasses import dataclass
from typing import List, Optional, Tuple
from fastapi import HTTPException, UploadFile
from pydantic import BaseModel
from app.api.batch_job_store import batch_job_store
from app.core import batch_ingest, batch_worker
from app.core import store_catalog_pipeline as pipeline
logger = logging.getLogger(__name__)
# The worker cannot import the API layer without a cycle, so the wiring is done
# once, here, at import. This module is imported by every router that can start
# a batch, which is why it is the right place: whichever of them loads first,
# the worker is configured before anything can be submitted to it.
batch_worker.configure(
on_change=batch_job_store.put,
should_cancel=batch_job_store.is_cancelled,
)
@dataclass(frozen=True)
class UploadLimits:
"""Ceilings for one submission. Per-file first, then per-batch."""
max_files: int
max_file_bytes: int
max_file_rows: int
max_total_bytes: int
max_total_rows: int
# ---------------------------------------------------------------------------
# Reading and parsing
# ---------------------------------------------------------------------------
async def read_uploads(
files: List[UploadFile], limits: UploadLimits
) -> List[Tuple[str, bytes]]:
"""Read every upload into memory, enforcing the count and size ceilings.
Read here rather than in the worker for the reason `store_catalog.py` gives:
`UploadFile` is backed by a temporary file tied to the request, and it is
gone before a background thread would reach it.
"""
if not files:
raise HTTPException(status_code=400, detail="No files were uploaded.")
if len(files) > limits.max_files:
raise HTTPException(
status_code=413,
detail=(
f"{len(files)} files exceeds the {limits.max_files}-file limit for one "
f"batch. Split the drop and send it in two."
),
)
read: List[Tuple[str, bytes]] = []
total = 0
for upload in files:
contents = await upload.read()
name = upload.filename or "upload.xlsx"
if not contents:
# Recorded rather than raised - an empty file among nine good ones
# is a fact about that file, not a reason to reject the drop.
read.append((name, b""))
continue
if len(contents) > limits.max_file_bytes:
raise HTTPException(
status_code=413,
detail=(
f"'{name}' is larger than the "
f"{limits.max_file_bytes // (1024 * 1024)}MB per-file limit."
),
)
total += len(contents)
if total > limits.max_total_bytes:
raise HTTPException(
status_code=413,
detail=(
f"The batch is larger than the "
f"{limits.max_total_bytes // (1024 * 1024)}MB total limit."
),
)
read.append((name, contents))
return read
def parse_all(read: List[Tuple[str, bytes]], limits: UploadLimits):
"""Split the uploads into (valid, invalid) by trying to parse each one.
Failing here is what stops a batch transitioning straight to "failed" a
second after it started - the same reasoning as `store_catalog.py`, applied
per file so one bad sheet does not condemn the others.
"""
valid: List[Tuple[str, bytes, int]] = []
invalid: List[Tuple[str, str]] = []
rows_total = 0
for name, contents in read:
if not contents:
invalid.append((name, "The file is empty."))
continue
try:
df, mapping = pipeline.parse_spreadsheet(name, contents)
except HTTPException as exc:
# read_products_dataframe raises HTTPException for an unsupported
# extension or a missing Excel reader; its message already names the
# file and says what to do about it.
invalid.append((name, str(exc.detail)))
continue
except Exception as exc: # noqa: BLE001 - an unreadable sheet is caller error
invalid.append((name, f"Could not parse the file: {exc}"))
continue
if df.empty:
invalid.append((name, "The file has no data rows."))
continue
# A sheet whose headers carry no product name is not a catalog, and
# without this it is accepted with a 202 and then ingests nothing - the
# worst possible answer, because it looks like success from every angle
# the caller can see. `user_products.py` has always made this check on
# its own upload path; the catalog paths did not, and an API client
# sending the wrong export is the likeliest mistake there is.
if "product_name" not in mapping.columns:
recognised = ", ".join(sorted(mapping.columns)) or "none"
invalid.append((
name,
f"No product name column was found. Headers read: "
f"{', '.join(str(c) for c in df.columns)}. Recognised fields: "
f"{recognised}.",
))
continue
if len(df) > limits.max_file_rows:
invalid.append((
name,
f"{len(df)} rows exceeds the {limits.max_file_rows}-row per-file limit.",
))
continue
rows_total += int(len(df))
if rows_total > limits.max_total_rows:
raise HTTPException(
status_code=413,
detail=(
f"The batch totals more than {limits.max_total_rows} rows. "
f"Split it and send it in two."
),
)
valid.append((name, contents, int(len(df))))
return valid, invalid, rows_total
# ---------------------------------------------------------------------------
# Response shape
# ---------------------------------------------------------------------------
class BatchFileOut(BaseModel):
index: int
filename: str
status: str
detail: Optional[str] = None
stage_index: int = 0
stage_name: str = ""
total_stages: int = pipeline.TOTAL_STAGES
rows_done: int = 0
rows_total: int = 0
size_bytes: int = 0
result: Optional[dict] = None
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
files_done: int
files_failed: int
current_file: Optional[str] = None
use_llm: bool
fetch_images: bool
totals: dict
brands: List[str]
files: List[BatchFileOut]
def to_out(manifest: batch_ingest.BatchManifest) -> BatchOut:
body = manifest.to_dict()
body["files"] = [BatchFileOut(**{
key: entry[key] for key in BatchFileOut.model_fields if key in entry
}) for entry in body["files"]]
return BatchOut(**{k: v for k, v in body.items() if k in BatchOut.model_fields})
# ---------------------------------------------------------------------------
# Staging and queueing
# ---------------------------------------------------------------------------
def stage_and_queue(
valid: List[Tuple[str, bytes, int]],
invalid: List[Tuple[str, str]],
*,
use_llm: bool,
fetch_images: bool,
submitted_by: Optional[str] = None,
) -> Tuple[batch_ingest.BatchManifest, bool]:
"""Write the files down, publish the batch, and try to start it.
Returns `(manifest, started)`. `started` is False only when the worker queue
was full: the batch is staged and durable either way, and the caller decides
what to say about it - an admin has a Resume button, an API client does not,
and the two deserve different words for the same 429.
"""
manifest = batch_ingest.stage_batch(
[(name, contents) for name, contents, _n in valid],
use_llm=use_llm,
fetch_images=fetch_images,
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
batch_ingest.write_manifest(manifest)
batch_job_store.put(manifest)
try:
batch_worker.submit(manifest.batch_id)
except queue.Full:
manifest.status = batch_ingest.QUEUED
manifest.detail = (
"The ingestion queue was full when this batch arrived. It is staged "
"and can be started with Resume once the running batches finish."
)
batch_ingest.write_manifest(manifest)
batch_job_store.put(manifest)
return manifest, False
return manifest, True

View File

@@ -3,7 +3,7 @@ from __future__ import annotations
import io
import logging
from typing import Any, Dict, List, Optional
from typing import List, Optional
import pandas as pd
from pydantic import BaseModel, Field
from fastapi import APIRouter, Depends, File, HTTPException, UploadFile
@@ -12,7 +12,6 @@ from app.api.deps import require_permission
from app.infrastructure.settings import S3_BUCKET
from app.services.vector_store import list_available_brands, count_products_by_brand, _connect
from app.services.s3_service import s3_service
from app.services import store_db
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/admin/training", tags=["admin_train"])

View File

@@ -21,19 +21,26 @@ image stage, a Playwright subprocess that can burn three minutes on its own -
on a single-vCPU container that is also serving the API. So a batch opts IN to
those stages; it does not opt out. `USE_OLLAMA` is false in production anyway,
which makes `use_llm` a no-op there and the honest default obvious.
THE OTHER WAY INTO THE SAME PIPELINE
------------------------------------
`app/api/routers/uploads.py` exposes ingestion to outside API clients under
`upload_catalog` rather than `require_admin`. It stages and queues through the
identical helpers (`app/api/batch_common.py`) and produces ordinary batches,
so everything here - the list, resume, cancel - applies to those too. The only
difference is that its reads are filtered to the caller's own submissions.
"""
from __future__ import annotations
import logging
import queue
from typing import List, Optional
from typing import List
from fastapi import APIRouter, Depends, File, HTTPException, UploadFile, status
from pydantic import BaseModel
from app.api import batch_common
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, inbox
from app.core import batch_ingest, batch_worker
from app.core import store_catalog_pipeline as pipeline
from app.infrastructure.settings import (
BATCH_MAX_FILES,
@@ -41,7 +48,6 @@ from app.infrastructure.settings import (
BATCH_MAX_TOTAL_ROWS,
)
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/admin/catalog-batch", tags=["admin", "catalog"])
# Per-file ceilings match store_catalog.py exactly. A file that is too big for
@@ -51,152 +57,23 @@ MAX_UPLOAD_BYTES = 10 * 1024 * 1024
MAX_UPLOAD_ROWS = 2000
PREVIEW_ROWS = 10
# The worker cannot import the API layer without a cycle, so the wiring is done
# here, at import, once.
batch_worker.configure(
on_change=batch_job_store.put,
should_cancel=batch_job_store.is_cancelled,
)
# The response shape, the upload readers and the staging helper are shared with
# app/api/routers/uploads.py - see app/api/batch_common.py, which also wires the
# worker to the job store at import.
BatchFileOut = batch_common.BatchFileOut
BatchOut = batch_common.BatchOut
_to_out = batch_common.to_out
class BatchFileOut(BaseModel):
index: int
filename: str
status: str
detail: Optional[str] = None
stage_index: int = 0
stage_name: str = ""
total_stages: int = pipeline.TOTAL_STAGES
rows_done: int = 0
rows_total: int = 0
size_bytes: int = 0
result: Optional[dict] = None
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
files_done: int
files_failed: int
current_file: Optional[str] = None
use_llm: bool
fetch_images: bool
totals: dict
brands: List[str]
files: List[BatchFileOut]
def _to_out(manifest: batch_ingest.BatchManifest) -> BatchOut:
body = manifest.to_dict()
body["files"] = [BatchFileOut(**{
key: entry[key] for key in BatchFileOut.model_fields if key in entry
}) for entry in body["files"]]
return BatchOut(**{k: v for k, v in body.items() if k in BatchOut.model_fields})
async def _read_uploads(files: List[UploadFile]) -> List[tuple]:
"""Read every upload into memory, enforcing the count and size ceilings.
Read here rather than in the worker for the same reason `store_catalog.py`
gives: `UploadFile` is backed by a temporary file tied to the request, and
it is gone before a background thread would reach it.
"""
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"batch. Split the drop and upload it in two batches."
),
)
read: List[tuple] = []
total = 0
for upload in files:
contents = await upload.read()
name = upload.filename or "upload.xlsx"
if not contents:
# Recorded rather than raised - an empty file among nine good ones
# is a fact about that file, not a reason to reject the drop.
read.append((name, b""))
continue
if len(contents) > MAX_UPLOAD_BYTES:
raise HTTPException(
status_code=413,
detail=(
f"'{name}' is larger than the "
f"{MAX_UPLOAD_BYTES // (1024 * 1024)}MB per-file limit."
),
)
total += len(contents)
if total > BATCH_MAX_TOTAL_BYTES:
raise HTTPException(
status_code=413,
detail=(
f"The batch is larger than the "
f"{BATCH_MAX_TOTAL_BYTES // (1024 * 1024)}MB total limit."
),
)
read.append((name, contents))
return read
def _parse_all(read: List[tuple]):
"""Split the uploads into (valid, invalid) by trying to parse each one.
Failing fast here is what stops a batch transitioning straight to "failed"
a second after it started - the same reasoning as `store_catalog.py:147`,
applied per file so that one bad sheet does not condemn the others.
"""
valid: List[tuple] = []
invalid: List[tuple] = []
rows_total = 0
for name, contents in read:
if not contents:
invalid.append((name, "The file is empty."))
continue
try:
df, _mapping = pipeline.parse_spreadsheet(name, contents)
except HTTPException as exc:
# read_products_dataframe raises HTTPException for an unsupported
# extension or a missing Excel reader; its message already names
# the file and what to do about it.
invalid.append((name, str(exc.detail)))
continue
except Exception as exc: # noqa: BLE001 - an unreadable sheet is user error
invalid.append((name, f"Could not parse the file: {exc}"))
continue
if df.empty:
invalid.append((name, "The file has no data rows."))
continue
if len(df) > MAX_UPLOAD_ROWS:
invalid.append((
name,
f"{len(df)} rows exceeds the {MAX_UPLOAD_ROWS}-row per-file limit.",
))
continue
rows_total += int(len(df))
if rows_total > BATCH_MAX_TOTAL_ROWS:
raise HTTPException(
status_code=413,
detail=(
f"The batch totals more than {BATCH_MAX_TOTAL_ROWS} rows. "
f"Split it and upload in two batches."
),
)
valid.append((name, contents, int(len(df))))
return valid, invalid, rows_total
def _limits() -> batch_common.UploadLimits:
"""Read at call time, from THIS module's globals - see batch_common."""
return batch_common.UploadLimits(
max_files=BATCH_MAX_FILES,
max_file_bytes=MAX_UPLOAD_BYTES,
max_file_rows=MAX_UPLOAD_ROWS,
max_total_bytes=BATCH_MAX_TOTAL_BYTES,
max_total_rows=BATCH_MAX_TOTAL_ROWS,
)
@router.post("/preview", dependencies=[Depends(require_admin)])
@@ -207,7 +84,7 @@ async def preview_catalog_batch(files: List[UploadFile] = File(...)) -> dict:
headers onto catalog fields is a guess, and finding out that "Item" was read
as the description after twenty files have been scraped is expensive.
"""
read = await _read_uploads(files)
read = await batch_common.read_uploads(files, _limits())
out = []
for name, contents in read:
if not contents:
@@ -267,8 +144,9 @@ async def ingest_catalog_batch(
on a different host to the frontend, so anything that sat on the request
path would be racing an idle timeout nobody here controls.
"""
read = await _read_uploads(files)
valid, invalid, _rows = _parse_all(read)
limits = _limits()
read = await batch_common.read_uploads(files, limits)
valid, invalid, _rows = batch_common.parse_all(read, limits)
if not valid:
detail = "; ".join(f"{name}: {reason}" for name, reason in invalid)
@@ -277,27 +155,10 @@ async def ingest_catalog_batch(
detail=f"None of the uploaded files could be ingested. {detail}",
)
manifest = batch_ingest.stage_batch(
[(name, contents) for name, contents, _n in valid],
use_llm=use_llm,
fetch_images=fetch_images,
invalid=invalid,
manifest, started = batch_common.stage_and_queue(
valid, invalid, use_llm=use_llm, fetch_images=fetch_images,
)
for entry, (_name, _contents, rows) in zip(manifest.files, valid):
entry.rows_total = rows
batch_ingest.write_manifest(manifest)
batch_job_store.put(manifest)
try:
batch_worker.submit(manifest.batch_id)
except queue.Full:
manifest.status = batch_ingest.QUEUED
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)
if not started:
raise HTTPException(
status_code=429,
detail=(
@@ -309,139 +170,6 @@ 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))

View File

@@ -3,7 +3,6 @@ from __future__ import annotations
import logging
import os
import subprocess
import threading
from pathlib import Path
from typing import Any, Dict

View File

@@ -1,7 +1,43 @@
"""
Operator spreadsheet uploads: store inventory, sales history, nutrition facts.
POST /api/upload/stores (+ /stores/upload)
POST /api/upload/analytics (+ /analytics/upload)
POST /api/upload/nutrition (+ /nutrition/upload)
GET /api/upload/template/{tab_type}
These are the endpoints the Admin UI's upload tabs call, through
`/api/upload/${tabType}`.
WHY NOTHING HERE INVENTS A VALUE
--------------------------------
Every reader in this module returns `None` for an absent column, and every
writer either stores NULL or skips the row. That is a deliberate correction of
how this file used to work: it filled a missing column with a plausible
constant, so a spreadsheet whose headers did not match silently produced rows
attributed to brand "amul" in store "store_mumbai_1" at MRP 100 - and, worse,
nutrition rows carrying invented calories and `allergens = ['None']` stamped
`data_status = 'verified'`.
That last one is the reason this rule is absolute rather than a preference.
`app/services/nutrition_db.py` states the contract these tables are built on:
Every numeric column in `nutrition_facts` is nullable and stays NULL
unless a value was actually returned by a trusted source. Nothing in
this module ever writes an estimated, interpolated, or LLM-guessed
number into these columns.
A missing allergens column means "we were not told", never "this product is
allergen-free" - and the two are indistinguishable once a default has been
written. So a row that lacks the fields identifying it is REPORTED BACK to the
uploader, per row, rather than repaired into something importable.
Arithmetic on values that were supplied is not invention: `line_total` may be
derived from `quantity * unit_price`, because both were given.
"""
from __future__ import annotations
import io
import json
import uuid
import logging
import pandas as pd
@@ -18,6 +54,10 @@ from app.services import store_db, nutrition_db
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/upload", tags=["upload"])
# A caller who sent a 2000-row sheet with the wrong headers does not need 2000
# identical messages to understand what went wrong.
MAX_REPORTED_ERRORS = 50
def _normalize_col(col: str) -> str:
"""Normalize dataframe column names (lower, strip, replace spaces/hyphens with underscore)."""
@@ -36,39 +76,114 @@ def read_df_from_upload(filename: str, contents: bytes) -> pd.DataFrame:
df = pd.read_csv(io.BytesIO(contents))
except Exception:
df = pd.read_csv(io.BytesIO(contents), sep=None, engine='python')
# Rename columns to normalized format
df.columns = [_normalize_col(c) for c in df.columns]
return df
def _get_str(row: dict, keys: List[str], default: str = "") -> str:
# ---------------------------------------------------------------------------
# Cell readers
# ---------------------------------------------------------------------------
# All three return None for "this row did not carry the value", which every
# caller below is required to handle explicitly. There is deliberately no
# `default=` parameter: that parameter is what made a missing column look like
# data, and adding it back would reintroduce the bug this module documents.
def _opt_str(row: dict, keys: List[str]) -> Optional[str]:
for k in keys:
if k in row and pd.notna(row[k]):
val = str(row[k]).strip()
if val:
return val
return default
return None
def _get_float(row: dict, keys: List[str], default: float = 0.0) -> float:
def _opt_float(row: dict, keys: List[str]) -> Optional[float]:
for k in keys:
if k in row and pd.notna(row[k]):
try:
return float(row[k])
except (ValueError, TypeError):
pass
return default
return None
def _get_int(row: dict, keys: List[str], default: int = 0) -> int:
def _opt_int(row: dict, keys: List[str]) -> Optional[int]:
for k in keys:
if k in row and pd.notna(row[k]):
try:
return int(float(row[k]))
except (ValueError, TypeError):
pass
return default
return None
def _slug_id(brand: str, product_name: str) -> str:
"""A deterministic image_id for a row that did not carry one.
Derived entirely from values the sheet supplied, so it is a formatting
decision rather than an invented fact, and the same product re-uploaded
lands on the same key instead of duplicating.
"""
return f"{brand}_{product_name.lower().replace(' ', '_')}"
def _record(errors: List[Dict[str, Any]], index: Any, message: str) -> None:
"""Note why a row was skipped. Row numbers are as the uploader sees them:
1-based, counting the header, which is what their spreadsheet shows."""
if len(errors) < MAX_REPORTED_ERRORS:
try:
row_no = int(index) + 2
except (TypeError, ValueError):
row_no = -1
errors.append({"row": row_no, "error": message})
def _finish(
filename: Optional[str],
rows_total: int,
imported: int,
errors: List[Dict[str, Any]],
skipped: int,
noun: str,
**extra: Any,
) -> Dict[str, Any]:
"""Build the response, or refuse the upload if nothing at all landed.
Nothing imported is a 422 rather than a 200 with `rows_imported: 0`. The
old shape reported success for a file that stored not one row, which is how
a header mismatch went unnoticed for as long as it did.
"""
body = {
"status": "success" if skipped == 0 else "partial",
"filename": filename,
"rows_total": rows_total,
"rows_imported": imported,
"rows_skipped": skipped,
"errors": errors,
**extra,
}
if imported == 0:
raise HTTPException(
status_code=422,
detail={
"message": (
f"No {noun} could be imported from '{filename}'. "
f"Check the column headers against the sample template "
f"(GET /api/upload/template/...)."
),
"rows_total": rows_total,
"errors": errors,
},
)
if skipped:
body["message"] = (
f"Imported {imported} {noun}; skipped {skipped} row(s) that were "
f"missing required fields - see 'errors'."
)
else:
body["message"] = f"Successfully imported {imported} {noun}."
return body
# ---------------------------------------------------------------------------
@@ -79,69 +194,106 @@ def _get_int(row: dict, keys: List[str], default: int = 0) -> int:
@router.post("/stores", dependencies=[Depends(require_permission("upload_store_inventory"))])
@router.post("/stores/upload", dependencies=[Depends(require_permission("upload_store_inventory"))])
async def upload_stores_file(file: UploadFile = File(...)) -> Dict[str, Any]:
"""Import store inventory, and prices where the sheet carries them.
`store_id`, `brand` and `product_name` identify the row and are required;
a row without them is skipped and reported rather than filed under a
default store and brand.
Stock levels fall back to 0 - the column default the schema itself
declares - because `store_inventory` cannot hold NULL there. Prices are
different: `store_prices` requires all three of mrp/cost_price/selling_price
NOT NULL, and cost price cannot be derived from anything else on the row,
so a sheet without it gets its inventory imported and its prices left
alone, counted in `prices_skipped`. That is the honest outcome; the
alternative is a margin computed from a cost nobody supplied.
"""
if not file.filename:
raise HTTPException(status_code=400, detail="No file uploaded")
contents = await file.read()
try:
df = read_df_from_upload(file.filename, contents)
except Exception as e:
raise HTTPException(status_code=400, detail=f"Could not parse Excel/CSV file: {e}")
if df.empty:
raise HTTPException(status_code=400, detail="Uploaded file contains no data rows")
conn = _connect()
if not conn:
raise HTTPException(status_code=500, detail="Database connection failed")
imported_count = 0
prices_written = 0
prices_skipped = 0
stores_created = set()
errors: List[Dict[str, Any]] = []
skipped = 0
try:
with conn, conn.cursor() as cur:
# Ensure tables exist
store_db.ensure_store_intelligence_schema()
for _, r in df.iterrows():
for index, r in df.iterrows():
row = r.to_dict()
store_id = _get_str(row, ['store_id', 'store'], 'store_mumbai_1')
brand = _get_str(row, ['brand', 'brand_name'], 'amul').lower()
product_name = _get_str(row, ['product_name', 'title', 'name', 'item'], 'Product Item')
image_id = _get_str(row, ['image_id', 'sku', 'product_sku', 'item_id'], '')
if not image_id:
image_id = f"{brand}_{product_name.lower().replace(' ', '_')}"
category = _get_str(row, ['category', 'cat'], 'Dairy')
avail_stock = _get_int(row, ['available_stock', 'stock', 'qty', 'quantity'], 50)
reserved_stock = _get_int(row, ['reserved_stock', 'reserved'], 0)
reorder_lvl = _get_int(row, ['reorder_level', 'reorder'], 15)
safety_stk = _get_int(row, ['safety_stock', 'safety'], 10)
mrp = _get_float(row, ['mrp', 'price'], 100.0)
cost_price = _get_float(row, ['cost_price', 'cost'], 70.0)
selling_price = _get_float(row, ['selling_price', 'sell_price'], mrp * 0.9 if mrp else 90.0)
# 1. Ensure store exists
store_id = _opt_str(row, ['store_id', 'store'])
brand = _opt_str(row, ['brand', 'brand_name'])
product_name = _opt_str(row, ['product_name', 'title', 'name', 'item'])
missing = [
label for label, value in (
("store_id", store_id), ("brand", brand), ("product_name", product_name),
) if not value
]
if missing:
skipped += 1
_record(errors, index, f"missing required field(s): {', '.join(missing)}")
continue
brand = brand.lower()
image_id = _opt_str(row, ['image_id', 'sku', 'product_sku', 'item_id']) \
or _slug_id(brand, product_name)
category = _opt_str(row, ['category', 'cat'])
city = _opt_str(row, ['city', 'store_city'])
store_name = _opt_str(row, ['store_name']) or store_id.replace('_', ' ').title()
# NOT NULL with a schema default of 0. Absent means "not
# counted", which 0 represents as faithfully as anything can.
avail_stock = _opt_int(row, ['available_stock', 'stock', 'qty', 'quantity']) or 0
reserved_stock = _opt_int(row, ['reserved_stock', 'reserved']) or 0
reorder_lvl = _opt_int(row, ['reorder_level', 'reorder']) or 0
safety_stk = _opt_int(row, ['safety_stock', 'safety']) or 0
mrp = _opt_float(row, ['mrp', 'price'])
cost_price = _opt_float(row, ['cost_price', 'cost'])
selling_price = _opt_float(row, ['selling_price', 'sell_price'])
# 1. Ensure store exists. Only the id and a display name are
# asserted; city/tier/footfall stay at their schema defaults
# unless the sheet said otherwise.
cur.execute(
"""
INSERT INTO stores (store_id, store_name, city, tier, footfall_index)
VALUES (%s, %s, %s, %s, %s)
ON CONFLICT (store_id) DO NOTHING
INSERT INTO stores (store_id, store_name, city)
VALUES (%s, %s, %s)
ON CONFLICT (store_id) DO UPDATE SET
city = COALESCE(EXCLUDED.city, stores.city)
""",
(store_id, store_id.replace('_', ' ').title(), 'Mumbai', 'standard', 25.0)
(store_id, store_name, city)
)
stores_created.add(store_id)
# 2. Upsert store_inventory
cur.execute(
"""
INSERT INTO store_inventory
INSERT INTO store_inventory
(store_id, brand, image_id, title, category, available_stock, reserved_stock, reorder_level, safety_stock)
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s)
ON CONFLICT (store_id, brand, image_id) DO UPDATE SET
title = EXCLUDED.title,
category = EXCLUDED.category,
category = COALESCE(EXCLUDED.category, store_inventory.category),
available_stock = EXCLUDED.available_stock,
reserved_stock = EXCLUDED.reserved_stock,
reorder_level = EXCLUDED.reorder_level,
@@ -150,36 +302,41 @@ async def upload_stores_file(file: UploadFile = File(...)) -> Dict[str, Any]:
""",
(store_id, brand, image_id, product_name, category, avail_stock, reserved_stock, reorder_lvl, safety_stk)
)
# 3. Upsert store_prices
cur.execute(
"""
INSERT INTO store_prices (store_id, brand, image_id, mrp, cost_price, selling_price)
VALUES (%s, %s, %s, %s, %s, %s)
ON CONFLICT (store_id, brand, image_id) DO UPDATE SET
mrp = EXCLUDED.mrp,
cost_price = EXCLUDED.cost_price,
selling_price = EXCLUDED.selling_price,
updated_at = CURRENT_TIMESTAMP
""",
(store_id, brand, image_id, mrp, cost_price, selling_price)
)
# 3. Upsert store_prices - only with a complete price triple.
if mrp is not None and cost_price is not None and selling_price is not None:
cur.execute(
"""
INSERT INTO store_prices (store_id, brand, image_id, mrp, cost_price, selling_price)
VALUES (%s, %s, %s, %s, %s, %s)
ON CONFLICT (store_id, brand, image_id) DO UPDATE SET
mrp = EXCLUDED.mrp,
cost_price = EXCLUDED.cost_price,
selling_price = EXCLUDED.selling_price,
updated_at = CURRENT_TIMESTAMP
""",
(store_id, brand, image_id, mrp, cost_price, selling_price)
)
prices_written += 1
else:
prices_skipped += 1
imported_count += 1
except HTTPException:
raise
except Exception as e:
logger.error("Stores upload failed: %s", e)
raise HTTPException(status_code=500, detail=f"Database import failed: {e}")
finally:
conn.close()
return {
"status": "success",
"filename": file.filename,
"rows_total": len(df),
"rows_imported": imported_count,
"stores_affected": list(stores_created),
"message": f"Successfully imported {imported_count} store inventory items across {len(stores_created)} store(s)."
}
return _finish(
file.filename, len(df), imported_count, errors, skipped, "store inventory items",
stores_affected=sorted(stores_created),
prices_written=prices_written,
prices_skipped=prices_skipped,
)
# ---------------------------------------------------------------------------
@@ -188,244 +345,407 @@ async def upload_stores_file(file: UploadFile = File(...)) -> Dict[str, Any]:
@router.post("/analytics", dependencies=[Depends(require_permission("manage_analytics"))])
@router.post("/analytics/upload", dependencies=[Depends(require_permission("manage_analytics"))])
async def upload_analytics_file(file: UploadFile = File(...)) -> Dict[str, Any]:
"""Import sales transactions into `orders` / `order_items`.
THE COLUMN NAMES HERE WERE WRONG AND THE ENDPOINT NEVER WORKED.
The insert named `total_price`, which is not a column on `order_items`
(it is `line_total` - see store_db.py, which writes the same table), and it
omitted `store_id`, which is NOT NULL. Every call therefore raised, was
swallowed by the except below, and came back as a flat
500 "Database import failed". The input column may still be spelled
`total_price` in a customer's sheet - that alias is kept - but it is stored
in `line_total`, which is the column every analytics reader actually sums
(`intelligence/analytics.py`, `trending_model.py`, `features.py`).
Re-uploading the same file appends its lines again: `order_items` has a
surrogate key and no natural uniqueness to conflict on. Import each file
once, or give the rows stable `order_id`s and clear them first.
`customer_id` is required rather than defaulted. It used to fall back to a
single shared 'cust_imported', which silently merges every buyer in the
file into one customer and corrupts exactly the per-customer models -
purchase propensity, engagement - that this table exists to feed.
"""
if not file.filename:
raise HTTPException(status_code=400, detail="No file uploaded")
contents = await file.read()
try:
df = read_df_from_upload(file.filename, contents)
except Exception as e:
raise HTTPException(status_code=400, detail=f"Could not parse Excel/CSV file: {e}")
if df.empty:
raise HTTPException(status_code=400, detail="Uploaded file contains no data rows")
conn = _connect()
if not conn:
raise HTTPException(status_code=500, detail="Database connection failed")
imported_orders = 0
total_revenue = 0.0
errors: List[Dict[str, Any]] = []
skipped = 0
dates_defaulted = 0
touched_orders: List[str] = []
try:
with conn, conn.cursor() as cur:
store_db.ensure_store_intelligence_schema()
for _, r in df.iterrows():
for index, r in df.iterrows():
row = r.to_dict()
store_id = _get_str(row, ['store_id', 'store'], 'store_mumbai_1')
brand = _get_str(row, ['brand', 'brand_name'], 'amul').lower()
image_id = _get_str(row, ['image_id', 'sku', 'product_sku'], '')
product_name = _get_str(row, ['product_name', 'title', 'item'], 'Analytics Item')
if not image_id:
image_id = f"{brand}_{product_name.lower().replace(' ', '_')}"
order_id = _get_str(row, ['order_id', 'transaction_id'], f"ord_up_{uuid.uuid4().hex[:8]}")
customer_id = _get_str(row, ['customer_id', 'user_id', 'customer'], 'cust_imported')
raw_date = _get_str(row, ['order_date', 'date', 'timestamp'], '')
order_date = datetime.now()
store_id = _opt_str(row, ['store_id', 'store'])
brand = _opt_str(row, ['brand', 'brand_name'])
product_name = _opt_str(row, ['product_name', 'title', 'item'])
image_id = _opt_str(row, ['image_id', 'sku', 'product_sku'])
customer_id = _opt_str(row, ['customer_id', 'user_id', 'customer'])
qty = _opt_int(row, ['quantity', 'units_sold', 'qty', 'count'])
unit_price = _opt_float(row, ['unit_price', 'selling_price', 'price'])
missing = [
label for label, value in (
("store_id", store_id),
("brand", brand),
("customer_id", customer_id),
("quantity", qty),
("unit_price", unit_price),
) if value is None or value == ""
]
if not image_id and not product_name:
missing.append("image_id or product_name")
if missing:
skipped += 1
_record(errors, index, f"missing required field(s): {', '.join(missing)}")
continue
brand = brand.lower()
image_id = image_id or _slug_id(brand, product_name)
# A surrogate key, not a fact about the sale: a sheet without
# order ids is one row per order, which is what this generates.
order_id = _opt_str(row, ['order_id', 'transaction_id']) \
or f"ord_up_{uuid.uuid4().hex[:8]}"
raw_date = _opt_str(row, ['order_date', 'date', 'timestamp'])
order_date = None
if raw_date:
try:
order_date = pd.to_datetime(raw_date).to_pydatetime()
except Exception:
pass
qty = _get_int(row, ['quantity', 'units_sold', 'qty', 'count'], 1)
unit_price = _get_float(row, ['unit_price', 'selling_price', 'price'], 100.0)
tot_price = _get_float(row, ['total_price', 'revenue', 'total'], qty * unit_price)
# Ensure store exists
order_date = None
if order_date is None:
# orders.order_date is NOT NULL. Import time is the only
# defensible stand-in, and it is counted so the response can
# say how much of the file is not really dated.
order_date = datetime.now()
dates_defaulted += 1
# Arithmetic over supplied values, not invention.
line_total = _opt_float(row, ['total_price', 'line_total', 'revenue', 'total'])
if line_total is None:
line_total = qty * unit_price
payment_method = _opt_str(row, ['payment_method', 'payment'])
delivery_status = _opt_str(row, ['delivery_status', 'status'])
cur.execute(
"INSERT INTO stores (store_id, store_name, city, tier, footfall_index) VALUES (%s, %s, %s, %s, %s) ON CONFLICT (store_id) DO NOTHING",
(store_id, store_id.replace('_', ' ').title(), 'Mumbai', 'standard', 25.0)
"""
INSERT INTO stores (store_id, store_name)
VALUES (%s, %s)
ON CONFLICT (store_id) DO NOTHING
""",
(store_id, store_id.replace('_', ' ').title())
)
# Insert order header
cur.execute(
"""
INSERT INTO orders (order_id, customer_id, store_id, order_date, payment_method, order_value, delivery_status)
VALUES (%s, %s, %s, %s, %s, %s, %s)
ON CONFLICT (order_id) DO UPDATE SET order_value = EXCLUDED.order_value
ON CONFLICT (order_id) DO NOTHING
""",
(order_id, customer_id, store_id, order_date, 'upi', tot_price, 'delivered')
(order_id, customer_id, store_id, order_date, payment_method, 0, delivery_status)
)
# Insert order item
cur.execute(
"""
INSERT INTO order_items (order_id, brand, image_id, quantity, unit_price, total_price)
VALUES (%s, %s, %s, %s, %s, %s)
INSERT INTO order_items
(order_id, store_id, brand, image_id, quantity, unit_price, line_total)
VALUES (%s, %s, %s, %s, %s, %s, %s)
""",
(order_id, brand, image_id, qty, unit_price, tot_price)
(order_id, store_id, brand, image_id, qty, unit_price, line_total)
)
touched_orders.append(order_id)
imported_orders += 1
total_revenue += tot_price
total_revenue += line_total
# order_value is the sum of the order's lines, so it is recomputed
# from order_items rather than accumulated per row. A multi-line
# order previously ended up carrying only its last line's value.
if touched_orders:
cur.execute(
"""
UPDATE orders o
SET order_value = s.total
FROM (
SELECT order_id, SUM(line_total) AS total
FROM order_items
WHERE order_id = ANY(%s)
GROUP BY order_id
) s
WHERE o.order_id = s.order_id
""",
(list(set(touched_orders)),)
)
except HTTPException:
raise
except Exception as e:
logger.error("Analytics upload failed: %s", e)
raise HTTPException(status_code=500, detail=f"Database import failed: {e}")
finally:
conn.close()
return {
"status": "success",
"filename": file.filename,
"rows_total": len(df),
"rows_imported": imported_orders,
"total_revenue": round(total_revenue, 2),
"message": f"Successfully imported {imported_orders} sales transactions (Total Revenue: ₹{total_revenue:,.2f})."
}
return _finish(
file.filename, len(df), imported_orders, errors, skipped, "sales transactions",
total_revenue=round(total_revenue, 2),
dates_defaulted_to_now=dates_defaulted,
)
# ---------------------------------------------------------------------------
# Nutrition Intelligence Excel / CSV Upload
# ---------------------------------------------------------------------------
# The per-100g fields the "high protein" / "low sugar" endpoints sort on. Which
# of these a row carries decides its data_status, so the list is named once
# rather than repeated in the status logic.
_CORE_NUTRIENTS = (
"calories_kcal", "protein_g", "carbohydrates_g", "total_sugar_g",
"dietary_fiber_g", "total_fat_g", "sodium_mg",
)
@router.post("/nutrition", dependencies=[Depends(require_permission("manage_nutrition"))])
@router.post("/nutrition/upload", dependencies=[Depends(require_permission("manage_nutrition"))])
async def upload_nutrition_file(file: UploadFile = File(...)) -> Dict[str, Any]:
"""Import nutrition facts, storing NULL for everything the sheet omitted.
THIS IS THE ENDPOINT THE DATA-INTEGRITY RULE EXISTS FOR.
It used to substitute a constant for every absent column - 150 kcal, 5g
protein, health_score 78, `diet_tags = ['High Protein', 'Gluten Free']`,
`allergens = ['None']` - and write the result with
`data_status = 'verified'`. A sheet with no allergens column therefore
asserted, verified, that every product in it was allergen-free, and the
public /api/nutrition endpoints served that.
Now: absent means NULL, `data_status` reflects what the row actually
carried ('verified' only with the full core set, else 'partial', else
'unavailable'), and `allergen_source` records 'upload' or 'unavailable' so
a reader can tell "no allergens" from "not told". The upserts COALESCE, so
a later sheet that omits a column can never blank a value an earlier
trusted source established, and a row already 'verified' is never
downgraded by a thinner upload.
"""
if not file.filename:
raise HTTPException(status_code=400, detail="No file uploaded")
contents = await file.read()
try:
df = read_df_from_upload(file.filename, contents)
except Exception as e:
raise HTTPException(status_code=400, detail=f"Could not parse Excel/CSV file: {e}")
if df.empty:
raise HTTPException(status_code=400, detail="Uploaded file contains no data rows")
conn = _connect()
if not conn:
raise HTTPException(status_code=500, detail="Database connection failed")
imported_count = 0
errors: List[Dict[str, Any]] = []
skipped = 0
status_counts = {"verified": 0, "partial": 0, "unavailable": 0}
try:
with conn, conn.cursor() as cur:
nutrition_db.ensure_nutrition_schema()
for _, r in df.iterrows():
for index, r in df.iterrows():
row = r.to_dict()
brand = _get_str(row, ['brand', 'brand_name'], 'amul').lower()
product_name = _get_str(row, ['product_name', 'title', 'item', 'name'], 'Nutrition Item')
image_id = _get_str(row, ['image_id', 'sku', 'id'], '')
if not image_id:
image_id = f"{brand}_{product_name.lower().replace(' ', '_')}"
category = _get_str(row, ['category', 'cat'], 'Food')
calories = _get_float(row, ['calories', 'calories_kcal', 'energy'], 150.0)
protein = _get_float(row, ['protein', 'protein_g'], 5.0)
carbs = _get_float(row, ['carbohydrates', 'carbs', 'carbohydrates_g'], 20.0)
sugar = _get_float(row, ['sugar', 'total_sugar_g', 'sugars'], 4.0)
fiber = _get_float(row, ['fiber', 'dietary_fiber_g'], 2.0)
fat = _get_float(row, ['fat', 'total_fat_g'], 6.0)
sodium = _get_float(row, ['sodium', 'sodium_mg'], 120.0)
calcium = _get_float(row, ['calcium', 'calcium_mg'], 80.0)
iron = _get_float(row, ['iron', 'iron_mg'], 1.5)
vitamin_c = _get_float(row, ['vitamin_c', 'vitamin_c_mg'], 5.0)
health_score = _get_float(row, ['health_score', 'nutrition_score', 'score'], 78.0)
diet_tags_raw = _get_str(row, ['diet_tags', 'tags', 'diet'], 'High Protein, Gluten Free')
allergens_raw = _get_str(row, ['allergens', 'allergen'], 'None')
diet_tags = [t.strip() for t in diet_tags_raw.split(',') if t.strip()]
allergens = [a.strip() for a in allergens_raw.split(',') if a.strip()]
# 1. Upsert nutrition_facts
brand = _opt_str(row, ['brand', 'brand_name'])
product_name = _opt_str(row, ['product_name', 'title', 'item', 'name'])
image_id = _opt_str(row, ['image_id', 'sku', 'id'])
missing = [
label for label, value in (("brand", brand),) if not value
]
if not image_id and not product_name:
missing.append("image_id or product_name")
if missing:
skipped += 1
_record(errors, index, f"missing required field(s): {', '.join(missing)}")
continue
brand = brand.lower()
image_id = image_id or _slug_id(brand, product_name)
category = _opt_str(row, ['category', 'cat'])
values = {
"calories_kcal": _opt_float(row, ['calories', 'calories_kcal', 'energy']),
"protein_g": _opt_float(row, ['protein', 'protein_g']),
"carbohydrates_g": _opt_float(row, ['carbohydrates', 'carbs', 'carbohydrates_g']),
"total_sugar_g": _opt_float(row, ['sugar', 'total_sugar_g', 'sugars']),
"dietary_fiber_g": _opt_float(row, ['fiber', 'dietary_fiber_g']),
"total_fat_g": _opt_float(row, ['fat', 'total_fat_g']),
"sodium_mg": _opt_float(row, ['sodium', 'sodium_mg']),
"calcium_mg": _opt_float(row, ['calcium', 'calcium_mg']),
"iron_mg": _opt_float(row, ['iron', 'iron_mg']),
"vitamin_c_mg": _opt_float(row, ['vitamin_c', 'vitamin_c_mg']),
}
present_core = [k for k in _CORE_NUTRIENTS if values.get(k) is not None]
if len(present_core) == len(_CORE_NUTRIENTS):
data_status = "verified"
elif present_core:
data_status = "partial"
else:
data_status = "unavailable"
status_counts[data_status] += 1
cur.execute(
"""
INSERT INTO nutrition_facts
(brand, image_id, product_name, category, data_status, data_source,
calories_kcal, protein_g, carbohydrates_g, total_sugar_g, dietary_fiber_g,
total_fat_g, sodium_mg, calcium_mg, iron_mg, vitamin_c_mg)
VALUES (%s, %s, %s, %s, 'verified', 'excel_upload', %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
total_fat_g, sodium_mg, calcium_mg, iron_mg, vitamin_c_mg, updated_at)
VALUES (%s, %s, %s, %s, %s, 'excel_upload',
%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, CURRENT_TIMESTAMP)
ON CONFLICT (brand, image_id) DO UPDATE SET
product_name = EXCLUDED.product_name,
category = EXCLUDED.category,
data_status = 'verified',
calories_kcal = EXCLUDED.calories_kcal,
protein_g = EXCLUDED.protein_g,
carbohydrates_g = EXCLUDED.carbohydrates_g,
total_sugar_g = EXCLUDED.total_sugar_g,
dietary_fiber_g = EXCLUDED.dietary_fiber_g,
total_fat_g = EXCLUDED.total_fat_g,
sodium_mg = EXCLUDED.sodium_mg,
calcium_mg = EXCLUDED.calcium_mg,
iron_mg = EXCLUDED.iron_mg,
vitamin_c_mg = EXCLUDED.vitamin_c_mg
product_name = COALESCE(EXCLUDED.product_name, nutrition_facts.product_name),
category = COALESCE(EXCLUDED.category, nutrition_facts.category),
data_source = 'excel_upload',
-- Never downgrade a row a trusted source already verified.
data_status = CASE WHEN nutrition_facts.data_status = 'verified'
THEN 'verified' ELSE EXCLUDED.data_status END,
calories_kcal = COALESCE(EXCLUDED.calories_kcal, nutrition_facts.calories_kcal),
protein_g = COALESCE(EXCLUDED.protein_g, nutrition_facts.protein_g),
carbohydrates_g = COALESCE(EXCLUDED.carbohydrates_g, nutrition_facts.carbohydrates_g),
total_sugar_g = COALESCE(EXCLUDED.total_sugar_g, nutrition_facts.total_sugar_g),
dietary_fiber_g = COALESCE(EXCLUDED.dietary_fiber_g, nutrition_facts.dietary_fiber_g),
total_fat_g = COALESCE(EXCLUDED.total_fat_g, nutrition_facts.total_fat_g),
sodium_mg = COALESCE(EXCLUDED.sodium_mg, nutrition_facts.sodium_mg),
calcium_mg = COALESCE(EXCLUDED.calcium_mg, nutrition_facts.calcium_mg),
iron_mg = COALESCE(EXCLUDED.iron_mg, nutrition_facts.iron_mg),
vitamin_c_mg = COALESCE(EXCLUDED.vitamin_c_mg, nutrition_facts.vitamin_c_mg),
updated_at = CURRENT_TIMESTAMP
""",
(brand, image_id, product_name, category, calories, protein, carbs, sugar, fiber, fat, sodium, calcium, iron, vitamin_c)
(brand, image_id, product_name, category, data_status,
values["calories_kcal"], values["protein_g"], values["carbohydrates_g"],
values["total_sugar_g"], values["dietary_fiber_g"], values["total_fat_g"],
values["sodium_mg"], values["calcium_mg"], values["iron_mg"],
values["vitamin_c_mg"])
)
# 2. Upsert nutrition_insights
insights_json = json.dumps({
"brand": brand,
"image_id": image_id,
"data_status": "verified",
"nutrition_score": health_score,
"health_score": health_score,
"positive_insights": [f"Contains {protein}g protein per 100g", f"Provides {fiber}g dietary fiber"],
"nutritional_cautions": [f"{sugar}g sugar per 100g"],
"diet_tags": diet_tags,
"allergens": allergens
})
cur.execute(
"""
INSERT INTO nutrition_insights
(brand, image_id, data_status, nutrition_score, health_score, score_breakdown,
positive_insights, nutritional_cautions, diet_tags, allergens)
VALUES (%s, %s, 'verified', %s, %s, %s, %s, %s, %s, %s)
ON CONFLICT (brand, image_id) DO UPDATE SET
data_status = 'verified',
nutrition_score = EXCLUDED.nutrition_score,
health_score = EXCLUDED.health_score,
positive_insights = EXCLUDED.positive_insights,
nutritional_cautions = EXCLUDED.nutritional_cautions,
diet_tags = EXCLUDED.diet_tags,
allergens = EXCLUDED.allergens
""",
(brand, image_id, health_score, health_score, json.dumps({"protein": 85, "fiber": 80}),
[f"Contains {protein}g protein per 100g"], [f"{sugar}g sugar per 100g"], diet_tags, allergens)
)
# -- insights -------------------------------------------------
# Only ever built from values this row actually carried. A row
# that carried none produces no insight row at all, rather than
# a confident-looking one full of defaults.
health_score = _opt_float(row, ['health_score', 'nutrition_score', 'score'])
diet_tags_raw = _opt_str(row, ['diet_tags', 'tags', 'diet'])
allergens_raw = _opt_str(row, ['allergens', 'allergen'])
diet_tags = [t.strip() for t in diet_tags_raw.split(',') if t.strip()] \
if diet_tags_raw is not None else None
allergens = [a.strip() for a in allergens_raw.split(',') if a.strip()] \
if allergens_raw is not None else None
# 'unavailable' is the schema's own vocabulary for "we were not
# told", and it is what stops an empty list reading as "none".
allergen_source = "upload" if allergens is not None else "unavailable"
positives = []
if values["protein_g"] is not None:
positives.append(f"Contains {values['protein_g']}g protein per 100g")
if values["dietary_fiber_g"] is not None:
positives.append(f"Provides {values['dietary_fiber_g']}g dietary fiber")
cautions = []
if values["total_sugar_g"] is not None:
cautions.append(f"{values['total_sugar_g']}g sugar per 100g")
if health_score is not None or diet_tags or allergens or positives or cautions:
cur.execute(
"""
INSERT INTO nutrition_insights
(brand, image_id, data_status, nutrition_score, health_score,
scoring_version, positive_insights, nutritional_cautions,
diet_tags, allergens, allergen_source, generated_at)
VALUES (%s, %s, %s, %s, %s, 'excel_upload', %s, %s, %s, %s, %s, CURRENT_TIMESTAMP)
ON CONFLICT (brand, image_id) DO UPDATE SET
data_status = CASE WHEN nutrition_insights.data_status = 'verified'
THEN 'verified' ELSE EXCLUDED.data_status END,
nutrition_score = COALESCE(EXCLUDED.nutrition_score, nutrition_insights.nutrition_score),
health_score = COALESCE(EXCLUDED.health_score, nutrition_insights.health_score),
scoring_version = 'excel_upload',
positive_insights = COALESCE(EXCLUDED.positive_insights, nutrition_insights.positive_insights),
nutritional_cautions = COALESCE(EXCLUDED.nutritional_cautions, nutrition_insights.nutritional_cautions),
diet_tags = COALESCE(EXCLUDED.diet_tags, nutrition_insights.diet_tags),
allergens = COALESCE(EXCLUDED.allergens, nutrition_insights.allergens),
allergen_source = CASE WHEN EXCLUDED.allergens IS NULL
THEN nutrition_insights.allergen_source
ELSE EXCLUDED.allergen_source END,
generated_at = CURRENT_TIMESTAMP
""",
(brand, image_id, data_status, health_score, health_score,
positives or None, cautions or None,
diet_tags, allergens, allergen_source)
)
imported_count += 1
except HTTPException:
raise
except Exception as e:
logger.error("Nutrition upload failed: %s", e)
raise HTTPException(status_code=500, detail=f"Database import failed: {e}")
finally:
conn.close()
return {
"status": "success",
"filename": file.filename,
"rows_total": len(df),
"rows_imported": imported_count,
"message": f"Successfully imported {imported_count} nutritional intelligence items."
}
return _finish(
file.filename, len(df), imported_count, errors, skipped,
"nutritional intelligence items",
data_status_counts=status_counts,
)
# ---------------------------------------------------------------------------
# Template Downloads
# ---------------------------------------------------------------------------
# The required columns below are the ones the handlers refuse a row without.
# Everything else is genuinely optional and is stored as NULL when omitted -
# these endpoints no longer fill a gap with a plausible-looking constant, so a
# template that omits a column now produces an honest blank rather than a
# confident wrong number.
@router.get("/template/{tab_type}")
def get_sample_template(tab_type: str) -> Response:
tab_type = tab_type.lower()
if tab_type == 'stores':
# Required: store_id, brand, product_name.
# Prices are written only when mrp, cost_price and selling_price are
# all present - store_prices declares all three NOT NULL.
content = (
"store_id,brand,image_id,product_name,category,available_stock,reserved_stock,mrp,cost_price,selling_price,reorder_level,safety_stock\n"
"store_mumbai_1,amul,amul_amul_butter_500ml,Amul Butter 500ml,Dairy,120,5,250.00,200.00,235.00,20,10\n"
"store_mumbai_1,amul,amul_amul_ghee_1l,Amul Ghee 1L,Dairy,85,2,650.00,520.00,610.00,15,5\n"
"store_delhi_2,nestle,nestle_everyday_1kg,Everyday Milk Powder 1kg,Dairy,45,0,420.00,340.00,399.00,10,5\n"
"store_id,brand,image_id,product_name,category,available_stock,reserved_stock,mrp,cost_price,selling_price,reorder_level,safety_stock,city\n"
"store_mumbai_1,amul,amul_amul_butter_500ml,Amul Butter 500ml,Dairy,120,5,250.00,200.00,235.00,20,10,Mumbai\n"
"store_mumbai_1,amul,amul_amul_ghee_1l,Amul Ghee 1L,Dairy,85,2,650.00,520.00,610.00,15,5,Mumbai\n"
"store_delhi_2,nestle,nestle_everyday_1kg,Everyday Milk Powder 1kg,Dairy,45,0,420.00,340.00,399.00,10,5,Delhi\n"
)
filename = "sample_stores_inventory_template.csv"
elif tab_type == 'analytics':
# Required: store_id, brand, customer_id, quantity, unit_price, and one
# of image_id / product_name. `total_price` is optional - it is derived
# from quantity * unit_price when absent - and is stored in the
# `line_total` column every analytics reader sums.
content = (
"order_id,store_id,brand,image_id,product_name,order_date,customer_id,quantity,unit_price,total_price\n"
"ORD_9001,store_mumbai_1,amul,amul_amul_butter_500ml,Amul Butter 500ml,2026-08-01 10:30:00,cust_101,2,235.00,470.00\n"
@@ -434,16 +754,21 @@ def get_sample_template(tab_type: str) -> Response:
)
filename = "sample_analytics_sales_template.csv"
elif tab_type == 'nutrition':
# Required: brand, and one of image_id / product_name. Every nutrient
# column is optional and stays NULL when omitted; a row is marked
# 'verified' only when the full core set is present. Leave `allergens`
# out entirely rather than writing "None" - an empty cell records "not
# told", which is not the same claim as "contains no allergens".
content = (
"brand,image_id,product_name,category,calories_kcal,protein_g,carbohydrates_g,total_sugar_g,dietary_fiber_g,total_fat_g,sodium_mg,health_score,diet_tags,allergens\n"
"amul,amul_amul_butter_500ml,Amul Butter 500ml,Dairy,717,0.8,0.1,0.0,0.0,81.0,650,75,Vegetarian,Dairy\n"
"amul,amul_amul_ghee_1l,Amul Ghee 1L,Dairy,898,0.0,0.0,0.0,0.0,99.8,0,82,Vegetarian,Keto Friendly\n"
"nestle,nestle_everyday_1kg,Everyday Milk Powder 1kg,Dairy,496,25.5,38.0,38.0,0.0,27.0,350,88,High Protein,Dairy\n"
"amul,amul_amul_butter_500ml,Amul Butter 500ml,Dairy,717,0.8,0.1,0.0,0.0,81.0,650,75,Vegetarian,Milk\n"
"amul,amul_amul_ghee_1l,Amul Ghee 1L,Dairy,898,0.0,0.0,0.0,0.0,99.8,0,82,Vegetarian,Milk\n"
"nestle,nestle_everyday_1kg,Everyday Milk Powder 1kg,Dairy,496,25.5,38.0,38.0,0.0,27.0,350,88,High Protein,Milk\n"
)
filename = "sample_nutrition_intelligence_template.csv"
else:
raise HTTPException(status_code=400, detail=f"Unknown template type '{tab_type}'. Use stores, analytics, or nutrition.")
return PlainTextResponse(
content=content,
media_type="text/csv",

View File

@@ -1,6 +1,8 @@
"""The one endpoint an outside contributor may call.
"""The catalog ingestion API given to outside API users.
POST /api/uploads/catalog - drop spreadsheets into the review inbox
POST /api/uploads/catalog - send spreadsheets, the pipeline runs
GET /api/uploads/catalog - the batches this caller has sent
GET /api/uploads/catalog/{batch_id} - progress and result of one of them
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
@@ -9,42 +11,55 @@ 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.
WHAT HAPPENS WHEN A FILE ARRIVES
--------------------------------
It is parsed during the request - while the caller is still on the phone - so
an unusable sheet comes back as a 400 naming the problem rather than as a job
that fails a minute later into a void. Then the bytes are staged to disk, a
batch is queued, and the same 11-stage pipeline the admin routes use runs over
them: `app/core/batch_ingest.py` -> `store_catalog_pipeline.run_pipeline`.
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.
The response is a `batch_id`. Ingestion is far too slow to finish inside a
request - it is thousands of rows through eleven stages - so the caller polls
GET /api/uploads/catalog/{batch_id} until `status` leaves `queued`/`running`.
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.
THIS CREDENTIAL NOW COSTS CPU, AND THAT IS THE POINT
----------------------------------------------------
An earlier version of this endpoint parked files in a review inbox and started
nothing, so that a leaked key could cost only disk. That is not the product:
an API user sends a file in order for it to be ingested, and a queue that needs
an admin to press a button is not an API.
So the bound is no longer "this role cannot start work" but "all work, from
every source, goes through one worker". `batch_worker` runs a single batch at a
time behind a queue of `BATCH_QUEUE_MAX`; past that this endpoint answers 429.
An uploader key can therefore occupy the ingestion worker, which is what it is
for - it cannot multiply it, which is what matters on a one-vCPU host that is
also serving the API and its healthcheck.
WHAT THIS ENDPOINT STILL CANNOT DO
----------------------------------
See anyone else's data. Every read here is filtered by `submitted_by`, so a key
sees the batches it sent and nothing else - not the catalog, not other callers'
submissions, not the admin batch list. Nothing here can cancel, resume, or
delete; those stay on the admin router.
"""
from __future__ import annotations
import logging
from typing import List, Optional
from typing import List
from fastapi import APIRouter, Depends, File, HTTPException, UploadFile, status
from pydantic import BaseModel
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.core import inbox
from app.core import store_catalog_pipeline as pipeline
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_FILES,
)
logger = logging.getLogger(__name__)
@@ -56,157 +71,138 @@ 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
def _limits() -> batch_common.UploadLimits:
"""Read at call time, from THIS module's globals - see batch_common."""
return batch_common.UploadLimits(
max_files=BATCH_MAX_FILES,
max_file_bytes=MAX_UPLOAD_BYTES,
max_file_rows=MAX_UPLOAD_ROWS,
max_total_bytes=BATCH_MAX_TOTAL_BYTES,
max_total_rows=BATCH_MAX_TOTAL_ROWS,
)
class SubmissionOut(BaseModel):
submission_id: Optional[str]
submitted_by: str
files_accepted: int
files_rejected: int
files: List[SubmittedFileOut]
class CatalogUploadOut(batch_common.BatchOut):
"""A started batch, plus a sentence a human can read without a schema."""
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.
def _owns(manifest: batch_ingest.BatchManifest, principal: Principal) -> bool:
"""May this caller see this batch?
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.
An admin sees everything - they already have the whole batch router. Anyone
else sees only what their own credential sent, matched on the credential
NAME, which is what `principal.username` is for an API key (see
security.principal_for_api_key).
"""
if not files:
raise HTTPException(status_code=400, detail="No files were uploaded.")
if len(files) > BATCH_MAX_FILES:
if principal.role == "admin":
return True
return bool(manifest.submitted_by) and manifest.submitted_by == principal.username
@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")),
) -> CatalogUploadOut:
"""Accept spreadsheets and run the catalog pipeline over them.
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.
`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.
"""
limits = _limits()
read = await batch_common.read_uploads(files, limits)
valid, invalid, _rows = batch_common.parse_all(read, limits)
if not valid:
detail = "; ".join(f"{name}: {reason}" for name, reason in invalid)
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."
),
status_code=400,
detail=f"None of the uploaded files could be ingested. {detail}",
)
# 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:
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,
)
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"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."
f"Too many batches are already queued. Batch {manifest.batch_id} has "
f"been saved but not started; retry this upload shortly."
),
)
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,
logger.info(
"Catalog batch %s queued: %d file(s), %d rejected, from %s",
manifest.batch_id, len(valid), len(invalid), principal.username,
)
body = batch_common.to_out(manifest).model_dump()
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."
)
return CatalogUploadOut(**body, message=message)
@router.get("/catalog")
def list_my_catalog_batches(
limit: int = 20,
principal: Principal = Depends(require_permission("upload_catalog")),
) -> dict:
"""The batches this credential has sent, newest first."""
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
# caller's own - or none at all while another key is busy.
mine = [
m for m in batch_job_store.recent(limit * 10) if _owns(m, principal)
][:limit]
return {"batches": [batch_common.to_out(m) for m in mine]}
@router.get("/catalog/{batch_id}")
def get_my_catalog_batch(
batch_id: str,
principal: Principal = Depends(require_permission("upload_catalog")),
) -> batch_common.BatchOut:
"""Progress and result for one batch this credential sent.
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.
"""
manifest = batch_job_store.get(batch_id)
if not manifest or not _owns(manifest, principal):
raise HTTPException(status_code=404, detail="Batch not found")
return batch_common.to_out(manifest)

View File

@@ -293,6 +293,7 @@ def stage_batch(
use_llm: bool = False,
fetch_images: bool = False,
invalid: Optional[List[Tuple[str, str]]] = None,
submitted_by: Optional[str] = None,
) -> BatchManifest:
"""Write the uploads to disk and return the manifest describing them.
@@ -300,12 +301,23 @@ def stage_batch(
They are recorded as failed members of the batch rather than dropped: an
operator who selected six files and sees five must be told what happened to
the sixth, and the batch page is the only place they will look.
`submitted_by` is the credential name the files arrived under. It is what
scopes an API client's view to its own batches, so it is set at staging
time rather than patched on afterwards - a batch that existed for even a
moment without an owner is a batch the ownership filter would hide from
the only person entitled to see it.
"""
batch_id = uuid.uuid4().hex
directory = batch_dir(batch_id)
directory.mkdir(parents=True, exist_ok=True)
manifest = BatchManifest(batch_id=batch_id, use_llm=use_llm, fetch_images=fetch_images)
manifest = BatchManifest(
batch_id=batch_id,
use_llm=use_llm,
fetch_images=fetch_images,
submitted_by=submitted_by,
)
position = 0
for filename, content in uploads:

View File

@@ -7,11 +7,10 @@ resort). The Node.js/Crawlee scraping layer that used to live here has
been removed - see app/services/image_search.py and
app/services/playwright_image_fallback.py for details.
"""
import asyncio
import json
import sys
from pathlib import Path
from typing import Dict, List, Optional, Any
from typing import Dict, List, Any
import logging
# Add app directory to path
@@ -41,10 +40,8 @@ def generate_product_highlights(product: Dict[str, Any], brand: str) -> List[str
product = {}
highlights = []
title = str(product.get('title', '')).strip()
description = str(product.get('description', '')).strip()
category = str(product.get('category', '')).strip()
price_range = str(product.get('price_range', '')).strip()
size_variants = product.get('size_variants', [])
# Ensure size_variants is a list
@@ -820,7 +817,7 @@ class ProductCatalogEngine:
price_str = variant.split(' - ₹')[1]
price_num = float(price_str)
prices.append(price_num)
except:
except (ValueError, TypeError):
pass
if prices:
min_price = min(prices)
@@ -918,7 +915,7 @@ class ProductCatalogEngine:
'products': enhanced_products
}
logger.info(f"🎉 Catalog generation complete!")
logger.info("🎉 Catalog generation complete!")
logger.info(f"📊 Products: {catalog['total_products']}")
logger.info(f"🖼️ Total images: {catalog['total_images']}")

View File

@@ -1,348 +0,0 @@
"""
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

@@ -53,7 +53,6 @@ from app.infrastructure.settings import (
ENABLE_BARCODE_LOOKUP,
ENABLE_HSN_GST_ENRICHMENT,
ENABLE_PRODUCT_VALIDATION,
ENABLE_SKU_WEB_LOOKUP,
MAX_VARIANTS_PER_PRODUCT,
USE_EMBEDDINGS,
)
@@ -64,7 +63,6 @@ from app.infrastructure.settings import (
# and tests/test_user_products_upload.py monkeypatches these names on that
# module object.
from app.api.routers.user_products import (
AddProductRequest,
_slugify,
_text,
map_spreadsheet_columns,

View File

@@ -72,8 +72,8 @@ 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.
# An outside API client that may send spreadsheets for catalog ingestion 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
@@ -82,9 +82,17 @@ ROLE_PERMISSIONS: Dict[str, List[str]] = {
# 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.
# WHAT A LEAKED UPLOADER KEY COSTS. Real CPU: this permission starts the
# 11-stage pipeline, which is the point of the endpoint. The bound is not
# "this role cannot work" but "all ingestion, from every source, shares one
# worker" - batch_worker runs a single batch at a time behind a queue of
# BATCH_QUEUE_MAX, past which POST /api/uploads/catalog answers 429. So a
# key can occupy the ingestion worker; it cannot multiply it, and it cannot
# touch the request path the healthcheck reads.
#
# What it still cannot do: read the catalog, read another caller's
# submissions (every read on that router is filtered by submitted_by), or
# cancel, resume or delete anything.
"uploader": [
"upload_catalog",
],

View File

@@ -169,24 +169,6 @@ 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.
#

View File

@@ -17,7 +17,7 @@ business-rule signal.
from __future__ import annotations
from dataclasses import dataclass
from typing import Dict, List
from typing import Dict
import numpy as np
import pandas as pd

View File

@@ -9,10 +9,9 @@ ML feature logic without a live Postgres instance.
from __future__ import annotations
import math
from datetime import date, datetime
from typing import Dict, Iterable, List, Optional
from datetime import date
from typing import Dict, Iterable, Optional
import numpy as np
import pandas as pd
# ---------------------------------------------------------------------------

View File

@@ -26,7 +26,7 @@ individually; only the training data is pooled).
"""
from __future__ import annotations
from datetime import date, timedelta
from datetime import date
from typing import List
import numpy as np

View File

@@ -20,7 +20,6 @@ from __future__ import annotations
import logging
from typing import Any, Dict, List, Tuple
import numpy as np
import pandas as pd
from app.intelligence.model_utils import ModelBundle, load_bundle, save_bundle

View File

@@ -22,7 +22,6 @@ from __future__ import annotations
import logging
from typing import Any, Dict, List
import numpy as np
import pandas as pd
from app.intelligence.model_utils import ModelBundle, load_bundle, save_bundle

View File

@@ -22,9 +22,9 @@ repeatable ML training and for demos.
"""
from __future__ import annotations
from dataclasses import dataclass, field
from dataclasses import dataclass
from datetime import date, timedelta
from typing import Dict, List, Optional, Sequence
from typing import Dict, List, Sequence
import hashlib
import numpy as np

View File

@@ -18,10 +18,9 @@ propensity score in the analytics UI).
"""
from __future__ import annotations
from datetime import date, timedelta
from datetime import date
from typing import List
import numpy as np
import pandas as pd
from app.intelligence import features as F

View File

@@ -15,7 +15,6 @@ heavy tuning.
"""
from __future__ import annotations
from typing import List
import numpy as np
import pandas as pd

View File

@@ -17,8 +17,8 @@ is `predicted_score.sort_values(ascending=False)` on live order data.
"""
from __future__ import annotations
from datetime import date, timedelta
from typing import Dict, List, Literal, Optional
from datetime import date
from typing import Dict, List, Literal
import numpy as np
import pandas as pd

View File

@@ -1,6 +1,6 @@
from __future__ import annotations
from typing import Dict, List, Optional
from typing import Dict, List
import pandas as pd
@@ -50,7 +50,6 @@ def product_dashboard(brand: str, image_id: str) -> Dict:
popularity = None
if engagement_row:
from app.intelligence import features as F
feat_row = pd.DataFrame([{
"views_norm": min(engagement_row["views"] / 10, 100),
"wishlist_norm": min(engagement_row["wishlist_count"] / 2, 100),

View File

@@ -10,7 +10,7 @@ ranked by health_score improvement rather than by distance.
"""
from __future__ import annotations
from typing import Any, Dict, List, Optional
from typing import Any, Dict, List
from app.services import nutrition_db

View File

@@ -39,7 +39,6 @@ treating NULL as "zero".
"""
from __future__ import annotations
import json
import logging
from typing import Any, Dict, List, Optional

View File

@@ -46,6 +46,11 @@ from typing import Any, Optional
from pydantic import BaseModel, Field
from app.infrastructure.settings import (
VALIDATION_REJECT_THRESHOLD,
VALIDATION_REVIEW_THRESHOLD,
)
from app.services.title_validator import find_category_conflicts
from app.services import price_estimator
from app.services import category_units as cu
@@ -119,10 +124,20 @@ class ValidationIssue(BaseModel):
class ValidationConfig(BaseModel):
"""Tunable thresholds - see docs/VALIDATION_PIPELINE.md for defaults
and how to change them via settings.py without editing this file."""
and how to change them via settings.py without editing this file.
reject_below: float = 0.35
review_below: float = 0.70
The two settings really are read now. They were declared in settings.py
with these same numbers as their defaults, and nothing ever imported them:
the values below were hardcoded copies, so setting VALIDATION_REJECT_THRESHOLD
in the environment did exactly nothing while appearing to work.
`default_factory` rather than a plain default so the setting is read when a
config is constructed, not once at import - which is what lets a test move
the threshold without reloading the module.
"""
reject_below: float = Field(default_factory=lambda: VALIDATION_REJECT_THRESHOLD)
review_below: float = Field(default_factory=lambda: VALIDATION_REVIEW_THRESHOLD)
class ValidationReport(BaseModel):

View File

@@ -6,13 +6,11 @@ functions need, then persists/reads the result cache via `store_db.py`.
from __future__ import annotations
import logging
from typing import Dict, List, Optional
from typing import Dict, List
import numpy as np
import pandas as pd
from app.intelligence import recommendation_engine as RE
from app.intelligence.popularity_model import popularity_scorer
from app.services import store_db
logger = logging.getLogger(__name__)

View File

@@ -34,7 +34,7 @@ from __future__ import annotations
import json
import logging
from datetime import date, datetime
from datetime import date
from typing import Any, Dict, List, Optional
import pandas as pd

View File

@@ -39,7 +39,7 @@ from __future__ import annotations
import re
import logging
from typing import Iterable, Optional
from typing import Optional
logger = logging.getLogger(__name__)

View File

@@ -1,8 +1,7 @@
from __future__ import annotations
from typing import List, Optional, Dict, Any
from typing import List, Optional, Dict, Any, Tuple
import json
import logging
import os
import re
@@ -10,7 +9,7 @@ import time
import psycopg
from app.infrastructure.settings import (
DATABASE_URL, USE_PGVECTOR, DB_HOST, DB_PORT, DB_NAME, DB_USER, DB_PASSWORD,
USE_PGVECTOR, DB_HOST, DB_PORT, DB_NAME, DB_USER, DB_PASSWORD,
DB_CONNECT_TIMEOUT_SECONDS,
)
from app.services.brand_registry import BRAND_ALIASES, resolve_parent_brand
@@ -161,8 +160,8 @@ def _ensure_columns(cur, table_name: str) -> None:
"embedding": "vector(384)",
}
cur.execute(
f"SELECT column_name, is_nullable, column_default FROM information_schema.columns "
f"WHERE table_schema = 'public' AND table_name = %s",
"SELECT column_name, is_nullable, column_default FROM information_schema.columns "
"WHERE table_schema = 'public' AND table_name = %s",
(table_name,),
)
col_info = cur.fetchall()