Feature Elimination backend
This commit is contained in:
@@ -1,7 +1,7 @@
|
||||
"""Live view of the batches this process knows about.
|
||||
|
||||
Same pattern and the same documented trade-offs as `store_catalog_job_store.py`
|
||||
and its siblings: a process-local dict behind a lock, not shared across uvicorn
|
||||
Same pattern and the same documented trade-offs as `job_store.py` and its
|
||||
siblings: a process-local dict behind a lock, not shared across uvicorn
|
||||
workers. Adding a broker for this would be operational weight the project has
|
||||
already decided against (see `job_store.py`).
|
||||
|
||||
|
||||
@@ -1,183 +0,0 @@
|
||||
"""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,
|
||||
)
|
||||
@@ -1,93 +0,0 @@
|
||||
"""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()
|
||||
@@ -31,7 +31,7 @@ from app.api.routers import health, brands, search, suggest, chat, catalog, syst
|
||||
from app.api.routers import stores, discounts, analytics as store_analytics, trending, recommendations, store_admin
|
||||
from app.api.routers import nutrition, nutrition_admin, upload
|
||||
from app.api.routers import auth, user_products, admin_train, mcp_info
|
||||
from app.api.routers import store_catalog, batch_catalog, uploads
|
||||
from app.api.routers import batch_catalog, uploads
|
||||
from app.services.store_db import ensure_store_intelligence_schema
|
||||
from app.services.nutrition_db import ensure_nutrition_schema
|
||||
|
||||
@@ -305,7 +305,6 @@ app.include_router(store_admin.router, prefix="/api")
|
||||
app.include_router(nutrition.router, prefix="/api")
|
||||
app.include_router(nutrition_admin.router, prefix="/api")
|
||||
app.include_router(upload.router, prefix="/api")
|
||||
app.include_router(store_catalog.router, prefix="/api")
|
||||
app.include_router(batch_catalog.router, prefix="/api")
|
||||
app.include_router(uploads.router, prefix="/api")
|
||||
app.include_router(mcp_info.router, prefix="/api")
|
||||
|
||||
@@ -44,12 +44,19 @@ httpx>=0.27.2
|
||||
# you can skip everything below this line.
|
||||
# Retry/backoff for the barcode enrichment sources
|
||||
# (app/services/enrichment/barcode/retry.py). This is NOT an optional extra:
|
||||
# retry.py imports it at module scope, and that module is reached from
|
||||
# app/main.py's own import of the store_catalog router - so a missing tenacity
|
||||
# is not a degraded feature, it is the container exiting 1 on boot with
|
||||
# ModuleNotFoundError before uvicorn ever binds a socket. Swarm then restarts
|
||||
# it, and each restart re-runs ~21s of pandas/scipy/sklearn imports, which on a
|
||||
# 1-vCPU host reads as pinned-at-100% CPU rather than as a crash.
|
||||
# retry.py imports it at module scope, and that module is reached at BOOT, via
|
||||
# app/main.py -> routers/{uploads,batch_catalog} -> api/batch_common
|
||||
# -> core/store_catalog_pipeline -> enrichment/barcode/retry
|
||||
# so a missing tenacity is not a degraded feature, it is the container exiting 1
|
||||
# on boot with ModuleNotFoundError before uvicorn ever binds a socket. Swarm
|
||||
# then restarts it, and each restart re-runs ~21s of pandas/scipy/sklearn
|
||||
# imports, which on a 1-vCPU host reads as pinned-at-100% CPU rather than as a
|
||||
# crash.
|
||||
#
|
||||
# This note used to name the store_catalog router as the importer. That router
|
||||
# was deleted with the Store Catalog Ingestion tab; the chain above is the one
|
||||
# that survives, and it runs through every remaining ingestion route. Re-check
|
||||
# it here before ever concluding this dependency has become unused.
|
||||
tenacity>=8.2.3
|
||||
|
||||
beautifulsoup4>=4.12.3
|
||||
|
||||
@@ -2,8 +2,15 @@
|
||||
|
||||
Follows the pattern established by test_user_products_upload.py: monkeypatch
|
||||
the storage/network boundary *on the pipeline module object* (it imports those
|
||||
names directly), build real .xlsx fixtures in memory, and assert on what was
|
||||
captured rather than on the HTTP status alone.
|
||||
names directly), build real .xlsx fixtures in memory, and assert on the rows
|
||||
that were written rather than on a status code.
|
||||
|
||||
This file covers the PIPELINE only, and deliberately so. It once ended with a
|
||||
section of HTTP tests against /api/admin/store-catalog/*; that router was
|
||||
deleted along with the Store Catalog Ingestion admin tab, and the same eleven
|
||||
stages are now driven from POST /api/uploads/catalog and the admin batch routes
|
||||
- whose HTTP surfaces are covered in test_uploads_autorun.py, test_uploads_api.py
|
||||
and test_batch_catalog_ingest.py.
|
||||
|
||||
Every test runs with use_llm=False and fetch_images=False so nothing here
|
||||
touches Ollama or the open web.
|
||||
@@ -323,64 +330,3 @@ def test_the_same_pack_listed_twice_is_written_once(store):
|
||||
)
|
||||
_run(content)
|
||||
assert len(store) == 1
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# HTTP surface
|
||||
# ---------------------------------------------------------------------------
|
||||
def test_preview_reports_how_the_columns_were_understood(client, admin_headers):
|
||||
content = _sheet(["Item Name", "Segment", "Mystery"], [["Britannia 50-50", "Biscuits", "?"]])
|
||||
resp = client.post(
|
||||
"/api/admin/store-catalog/preview",
|
||||
files={"file": ("store.xlsx", content)},
|
||||
headers=admin_headers,
|
||||
)
|
||||
assert resp.status_code == 200
|
||||
body = resp.json()
|
||||
assert body["recognised_columns"]["product_name"] == "Item Name"
|
||||
assert body["unrecognised_columns"] == ["Mystery"]
|
||||
assert body["brand_column_present"] is False
|
||||
assert len(body["stages"]) == 11
|
||||
|
||||
|
||||
def test_preview_requires_admin(client):
|
||||
content = _sheet(["Item Name"], [["Britannia 50-50"]])
|
||||
resp = client.post("/api/admin/store-catalog/preview", files={"file": ("s.xlsx", content)})
|
||||
assert resp.status_code == 401
|
||||
|
||||
|
||||
def test_an_unparseable_file_is_rejected_before_a_job_is_created(client, admin_headers):
|
||||
resp = client.post(
|
||||
"/api/admin/store-catalog/ingest",
|
||||
files={"file": ("notes.txt", b"this is not a spreadsheet")},
|
||||
headers=admin_headers,
|
||||
)
|
||||
assert resp.status_code == 400
|
||||
|
||||
|
||||
def test_ingest_returns_a_job_id_that_can_be_polled(client, admin_headers, monkeypatch):
|
||||
# Keep the worker off the network and out of the database.
|
||||
monkeypatch.setattr(pipeline, "upsert_brand_products", lambda b, r, cleanup=False: len(r))
|
||||
monkeypatch.setattr(pipeline, "get_products_by_brand", lambda b, **kw: [])
|
||||
monkeypatch.setattr(pipeline, "embed_texts", lambda texts: [[0.0] * 384 for _ in texts])
|
||||
|
||||
content = _sheet(["Item Name", "Segment"], [["Britannia 50-50", "Biscuits"]])
|
||||
resp = client.post(
|
||||
"/api/admin/store-catalog/ingest",
|
||||
files={"file": ("store.xlsx", content)},
|
||||
headers=admin_headers,
|
||||
)
|
||||
assert resp.status_code == 202
|
||||
job_id = resp.json()["job_id"]
|
||||
assert resp.json()["rows_total"] == 1
|
||||
|
||||
poll = client.get(f"/api/admin/store-catalog/jobs/{job_id}", headers=admin_headers)
|
||||
assert poll.status_code == 200
|
||||
body = poll.json()
|
||||
assert body["status"] in {"pending", "running", "done", "failed"}
|
||||
assert body["total_stages"] == 11
|
||||
|
||||
|
||||
def test_polling_an_unknown_job_is_a_404(client, admin_headers):
|
||||
resp = client.get("/api/admin/store-catalog/jobs/does-not-exist", headers=admin_headers)
|
||||
assert resp.status_code == 404
|
||||
|
||||
@@ -303,7 +303,6 @@ def test_uploader_key_is_refused_on_every_admin_route(client, upload_headers):
|
||||
("get", "/api/admin/catalog-batch/batches/anything"),
|
||||
("post", "/api/admin/catalog-batch/batches/anything/resume"),
|
||||
("post", "/api/admin/catalog-batch/batches/anything/cancel"),
|
||||
("post", "/api/admin/store-catalog/ingest"),
|
||||
("post", "/api/system/init"),
|
||||
]
|
||||
for method, path in admin_routes:
|
||||
|
||||
Reference in New Issue
Block a user