backend stores data file enrichment pipeline

This commit is contained in:
sriram
2026-08-18 16:58:36 +05:30
parent 691efef880
commit 7a4583372f
36 changed files with 4599 additions and 0 deletions

View File

@@ -0,0 +1,183 @@
"""Admin endpoints for turning a store spreadsheet into brand-catalog rows.
Three endpoints, matching how the other long-running admin jobs in this
project are exposed (see `nutrition_admin.py`):
POST /api/admin/store-catalog/preview - parse only, show the column map
POST /api/admin/store-catalog/ingest - 202 + job_id, runs in background
GET /api/admin/store-catalog/jobs/{id} - poll stage + row progress
`preview` exists because the mapping from a store's headers onto catalog
fields is a guess. Running an 11-stage scrape over 2000 rows only to discover
that "Item" was read as the description is expensive; showing the operator the
mapping first costs one parse.
"""
from __future__ import annotations
import logging
from typing import Optional
from fastapi import APIRouter, Depends, File, HTTPException, UploadFile, status
from pydantic import BaseModel
from app.api.background import run_in_background
from app.api.deps import require_admin
from app.api.store_catalog_job_store import store_catalog_job_store
from app.core import store_catalog_pipeline as pipeline
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/admin/store-catalog", tags=["admin", "catalog"])
# Same ceilings as the user-products upload path, for the same reasons.
MAX_UPLOAD_BYTES = 10 * 1024 * 1024
MAX_UPLOAD_ROWS = 2000
PREVIEW_ROWS = 10
class StoreCatalogJobOut(BaseModel):
job_id: str
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
result: Optional[dict] = None
async def _read_upload(file: UploadFile) -> bytes:
contents = await file.read()
if not contents:
raise HTTPException(status_code=400, detail="The uploaded file is empty.")
if len(contents) > MAX_UPLOAD_BYTES:
raise HTTPException(
status_code=413,
detail=f"File is larger than the {MAX_UPLOAD_BYTES // (1024 * 1024)}MB limit.",
)
return contents
@router.post("/preview", dependencies=[Depends(require_admin)])
async def preview_store_catalog(file: UploadFile = File(...)) -> dict:
"""Parse the sheet and report how its columns were understood."""
contents = await _read_upload(file)
try:
df, mapping = pipeline.parse_spreadsheet(file.filename or "upload.xlsx", contents)
except HTTPException:
raise
except ValueError as exc:
raise HTTPException(status_code=400, detail=str(exc)) from exc
except Exception as exc: # noqa: BLE001 - unreadable file is a user error
raise HTTPException(status_code=400, detail=f"Could not parse the file: {exc}") from exc
if df.empty:
raise HTTPException(status_code=400, detail="The file has no data rows.")
if len(df) > MAX_UPLOAD_ROWS:
raise HTTPException(
status_code=413,
detail=f"{len(df)} rows exceeds the {MAX_UPLOAD_ROWS}-row limit for one upload.",
)
sample = df.head(PREVIEW_ROWS).fillna("").astype(str).to_dict(orient="records")
return {
"filename": file.filename,
"rows_total": int(len(df)),
"recognised_columns": {f: str(c) for f, c in mapping.columns.items()},
"ignored_columns": mapping.ignored,
"unrecognised_columns": mapping.unrecognised,
"brand_column_present": "brand" in mapping.columns,
"preview": sample,
"stages": list(pipeline.STAGE_NAMES),
}
def _run_ingest_job(job_id: str, filename: str, contents: bytes,
use_llm: bool, fetch_images: bool) -> None:
store_catalog_job_store.update(job_id, status="running", detail="Starting pipeline")
def progress(stage_index: int, stage_name: str, done: int, total: int) -> None:
store_catalog_job_store.update(
job_id, stage_index=stage_index, stage_name=stage_name,
rows_done=done, rows_total=total,
)
try:
result = pipeline.run_pipeline(
filename, contents, progress=progress,
use_llm=use_llm, fetch_images=fetch_images,
)
body = result.as_dict()
if body.get("storage_error"):
# Rows were built but none reached the database. Reporting this as
# success would leave the operator believing the catalog changed.
store_catalog_job_store.update(
job_id, status="failed", result=body,
detail=f"Built {body['products_built']} row(s) but storing them failed: "
f"{body['storage_error']}",
)
return
store_catalog_job_store.update(
job_id, status="done", result=body,
detail=(f"{body['inserted']} inserted, {body['backfilled']} backfilled, "
f"{body['skipped_existing']} unchanged, {body['rejected']} rejected"),
)
except Exception as exc: # noqa: BLE001 - a daemon thread must not die silently
logger.exception("Store-catalog ingestion job %s failed", job_id)
store_catalog_job_store.update(job_id, status="failed", detail=str(exc))
@router.post("/ingest", status_code=status.HTTP_202_ACCEPTED,
dependencies=[Depends(require_admin)])
async def ingest_store_catalog(
file: UploadFile = File(...),
use_llm: bool = True,
fetch_images: bool = True,
) -> StoreCatalogJobOut:
"""Kick off the 11-stage pipeline and return a job id to poll.
The file is read here rather than in the worker: `UploadFile` is backed by
a temporary file tied to the request, so it is gone by the time a
background thread would get to it.
"""
contents = await _read_upload(file)
filename = file.filename or "upload.xlsx"
# Fail fast on an unparseable file so the caller gets a 400 now rather than
# a job that transitions straight to "failed".
try:
df, _mapping = pipeline.parse_spreadsheet(filename, contents)
except Exception as exc: # noqa: BLE001
raise HTTPException(status_code=400, detail=f"Could not parse the file: {exc}") from exc
if df.empty:
raise HTTPException(status_code=400, detail="The file has no data rows.")
if len(df) > MAX_UPLOAD_ROWS:
raise HTTPException(
status_code=413,
detail=f"{len(df)} rows exceeds the {MAX_UPLOAD_ROWS}-row limit for one upload.",
)
job = store_catalog_job_store.create(filename)
store_catalog_job_store.update(job.job_id, rows_total=int(len(df)))
run_in_background(
lambda: _run_ingest_job(job.job_id, filename, contents, use_llm, fetch_images),
name=f"store-catalog-{job.job_id[:8]}",
)
return StoreCatalogJobOut(
job_id=job.job_id, filename=filename, status=job.status,
rows_total=int(len(df)),
)
@router.get("/jobs/{job_id}", dependencies=[Depends(require_admin)])
def get_store_catalog_job(job_id: str) -> StoreCatalogJobOut:
job = store_catalog_job_store.get(job_id)
if not job:
raise HTTPException(status_code=404, detail="Job not found")
return StoreCatalogJobOut(
job_id=job.job_id, filename=job.filename, status=job.status, detail=job.detail,
stage_index=job.stage_index, stage_name=job.stage_name,
total_stages=job.total_stages, rows_done=job.rows_done,
rows_total=job.rows_total, result=job.result,
)

View File

@@ -0,0 +1,93 @@
"""Job tracking for store-spreadsheet catalog ingestion.
Same pattern and the same documented trade-offs as `job_store.py`,
`store_job_store.py` and `nutrition_job_store.py`: a process-local dict behind
a lock, lost on restart, not shared across uvicorn workers. Adding a broker for
this would be operational weight the project has already decided against (see
`job_store.py`).
It is its own module rather than a reuse of `nutrition_job_store` because this
job reports a different shape of progress: an 11-stage pipeline needs to say
*which stage* it is on as well as how many rows it has finished, so the UI can
show "Stage 6/11 - Image Search, 48/120 rows".
"""
from __future__ import annotations
import threading
import time
import uuid
from dataclasses import dataclass, field
from typing import Dict, Optional
@dataclass
class StoreCatalogJob:
job_id: str
filename: str
status: str = "pending" # pending -> running -> done | failed
detail: Optional[str] = None
result: Optional[dict] = None
stage_index: int = 0 # 1-based; 0 while still pending
stage_name: str = ""
total_stages: int = 11
rows_done: int = 0
rows_total: int = 0
created_at: float = field(default_factory=time.time)
updated_at: float = field(default_factory=time.time)
class StoreCatalogJobStore:
def __init__(self) -> None:
self._jobs: Dict[str, StoreCatalogJob] = {}
self._lock = threading.Lock()
def create(self, filename: str) -> StoreCatalogJob:
job = StoreCatalogJob(job_id=str(uuid.uuid4()), filename=filename)
with self._lock:
self._jobs[job.job_id] = job
return job
def update(
self,
job_id: str,
status: Optional[str] = None,
detail: Optional[str] = None,
result: Optional[dict] = None,
stage_index: Optional[int] = None,
stage_name: Optional[str] = None,
rows_done: Optional[int] = None,
rows_total: Optional[int] = None,
) -> None:
"""Every field is optional and None means "leave alone".
That matters: `job_store.update` sets `detail` unconditionally, so
marking a job running there wipes whatever detail it already had. A
progress callback firing many times per second must not erase state it
was not asked to change.
"""
with self._lock:
job = self._jobs.get(job_id)
if not job:
return
if status is not None:
job.status = status
if detail is not None:
job.detail = detail
if result is not None:
job.result = result
if stage_index is not None:
job.stage_index = stage_index
if stage_name is not None:
job.stage_name = stage_name
if rows_done is not None:
job.rows_done = rows_done
if rows_total is not None:
job.rows_total = rows_total
job.updated_at = time.time()
def get(self, job_id: str) -> Optional[StoreCatalogJob]:
with self._lock:
return self._jobs.get(job_id)
store_catalog_job_store = StoreCatalogJobStore()