Brand Ingestion

This commit is contained in:
sriram
2026-09-07 17:52:45 +05:30
parent 75dd3eb3ce
commit 2749bee1a3
11 changed files with 5069 additions and 10 deletions

View File

@@ -0,0 +1,253 @@
"""Admin routes for brand discovery: a brand NAME into the 11-stage pipeline.
Two steps on purpose, mirroring `batch_catalog.py`'s preview/ingest split.
`/preview` discovers and returns; it writes nothing, anywhere. `/ingest` takes
the rows the admin kept, renders them as a CSV, and hands the bytes to the same
`batch_common.stage_and_queue` an uploaded spreadsheet goes through - so the
batch manifest, the stage timeline, Resume, Cancel and the nutrition
auto-enrichment that follows a batch all work here without a line of new code.
The gap between the two steps is the point. Discovery's language-model half can
invent a product that nothing downstream is able to catch: a well-formed
fiction resolves a category, gets a price band and an internal SKU, and clears
`product_validator`'s "verified" threshold comfortably. `product_validator` was
built to reject MALFORMED rows, not false ones. A person looking at the list is
the check, so the list is shown before anything is written.
"""
from __future__ import annotations
import logging
from typing import Any, Dict, List, Optional
from fastapi import APIRouter, Depends, HTTPException, status
from pydantic import BaseModel, Field
from starlette.concurrency import run_in_threadpool
from app.api import batch_common
from app.api.deps import require_admin
from app.core import store_catalog_pipeline as pipeline
from app.infrastructure.settings import (
BATCH_MAX_FILES,
BATCH_MAX_TOTAL_BYTES,
BATCH_MAX_TOTAL_ROWS,
BRAND_DISCOVERY_DEADLINE_SECONDS,
BRAND_DISCOVERY_MAX_PRODUCTS,
)
from app.services import active_brands, brand_discovery
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/admin/brand-discovery", tags=["admin", "catalog"])
# The same per-file ceilings the admin batch routes apply. Discovery emits one
# CSV row per product and pack-size explosion happens later, inside stage 4, so
# 200 products is 200 rows here - three orders of magnitude inside the limit.
MAX_UPLOAD_BYTES = 10 * 1024 * 1024
MAX_UPLOAD_ROWS = 2000
def _limits() -> batch_common.UploadLimits:
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,
)
# ---------------------------------------------------------------------------
# Request bodies
# ---------------------------------------------------------------------------
class DiscoveryPreviewRequest(BaseModel):
brand: str
max_products: int = Field(default=BRAND_DISCOVERY_MAX_PRODUCTS, ge=1, le=2000)
use_openfacts: bool = True
use_llm: bool = True
# Ungrounded language-model rows are dropped rather than shown by default.
# Turning this off is how an admin sees them - they arrive unticked.
require_evidence: bool = True
# Shorter than the service default: somebody is watching a spinner.
deadline_seconds: float = Field(default=90.0, ge=0.0,
le=BRAND_DISCOVERY_DEADLINE_SECONDS)
refresh_corpus: bool = False
class DiscoveredProductIn(BaseModel):
"""One row the admin kept. Mirrors `DiscoveredProduct`'s written fields.
Sent back rather than re-discovered so that what is ingested is exactly what
was reviewed - a second discovery pass could legitimately return something
different, and then the approval would have been of a different list.
"""
product_name: str
title: Optional[str] = None
category: Optional[str] = None
description: Optional[str] = None
size_variants: List[str] = Field(default_factory=list)
providers: List[str] = Field(default_factory=list)
highlights: List[str] = Field(default_factory=list)
nutrients: List[str] = Field(default_factory=list)
fssai_license: Optional[str] = None
barcode: Optional[str] = None
image_url: Optional[str] = None
class DiscoveryIngestRequest(BaseModel):
brand: str
products: List[DiscoveredProductIn]
# Stage 2 only calls Ollama for a row whose description is blank, and
# discovery leaves most of them blank on purpose (see brand_discovery's note
# on the boilerplate generator). On an unreachable Ollama each such row
# costs up to OLLAMA_TIMEOUT_SECONDS, so this is worth being able to turn
# off for a large brand.
use_llm: bool = True
fetch_images: bool = True
# Ingesting a brand outside ACTIVE_BRANDS writes rows nothing can read.
# Refused unless the caller says they mean it - see the 409 below.
acknowledge_inactive: bool = False
# ---------------------------------------------------------------------------
# Routes
# ---------------------------------------------------------------------------
@router.post("/preview", dependencies=[Depends(require_admin)])
async def preview_brand_discovery(payload: DiscoveryPreviewRequest) -> Dict[str, Any]:
"""Discover a brand's products and return them. Writes nothing.
Run on a worker thread: discovery does blocking HTTP to Open Food Facts and,
when the language model is enabled, a series of blocking Ollama calls. On
the event loop that would stall every other request for the duration.
"""
try:
result = await run_in_threadpool(
brand_discovery.discover_brand_products,
payload.brand,
max_products=payload.max_products,
deadline_seconds=payload.deadline_seconds,
use_openfacts=payload.use_openfacts,
use_llm=payload.use_llm,
require_evidence=payload.require_evidence,
refresh_corpus=payload.refresh_corpus,
)
except ValueError as exc:
raise HTTPException(status_code=400, detail=str(exc)) from exc
except Exception as exc: # noqa: BLE001 - report the failure, do not 500
logger.exception("Brand discovery failed for %r", payload.brand)
raise HTTPException(
status_code=502,
detail=f"Discovery failed for {payload.brand!r}: {exc}",
) from exc
body = result.as_dict()
# Served rather than duplicated in the frontend, exactly as BatchOut does,
# so the two cannot drift when a stage is added.
body["stages"] = list(pipeline.STAGE_NAMES)
if result.filtering_enabled and not result.brand_active:
body["warnings"] = list(body.get("warnings") or []) + [_inactive_message(result)]
return body
@router.post("/ingest", status_code=status.HTTP_202_ACCEPTED,
dependencies=[Depends(require_admin)])
async def ingest_brand_discovery(payload: DiscoveryIngestRequest) -> batch_common.BatchOut:
"""Stage the reviewed products as a catalog batch and return an id to poll."""
brand = (payload.brand or "").strip()
if not brand:
raise HTTPException(status_code=400, detail="A brand name is required.")
if not payload.products:
raise HTTPException(
status_code=400,
detail="No products were selected, so there is nothing to ingest.",
)
# A green run over an unreadable catalog is the failure this project refuses
# to ship. The rows WOULD be written and fully enriched - nutrition
# auto-enrichment runs with include_inactive=True - but /api/brands, search,
# suggest and the category listing all filter the brand out, so the result
# looks like nothing happened.
if active_brands.filtering_enabled() and not active_brands.is_active_brand(brand):
if not payload.acknowledge_inactive:
raise HTTPException(status_code=409, detail=_inactive_detail(brand))
products = [
brand_discovery.DiscoveredProduct(
brand=brand,
product_name=item.product_name,
title=item.title or item.product_name,
category=item.category or "",
category_hint="",
description=item.description or "",
size_variants=list(item.size_variants),
providers=list(item.providers),
highlights=list(item.highlights),
nutrients=list(item.nutrients),
fssai_license=item.fssai_license,
barcode=item.barcode,
image_url=item.image_url,
)
for item in payload.products
]
filename = brand_discovery.synthetic_filename(brand)
contents = brand_discovery.rows_to_csv_bytes(products)
# ONE FILE, NEVER CHUNKED. stage 11 groups by brand per file and reads the
# existing catalog per file, and its intra-file de-duplication
# ({image_id: row}) is per file too - so the same image_id split across two
# chunks would not be caught. A single file makes that de-duplication total.
valid, invalid, _rows = batch_common.parse_all([(filename, contents)], _limits())
if not valid:
detail = "; ".join(f"{name}: {reason}" for name, reason in invalid)
raise HTTPException(
status_code=400,
detail=f"The discovered products could not be staged. {detail}",
)
manifest, started = batch_common.stage_and_queue(
valid, invalid,
use_llm=payload.use_llm,
fetch_images=payload.fetch_images,
submitted_by=f"brand-discovery: {brand}",
)
if not started:
raise HTTPException(
status_code=429,
detail=(
"Too many batches are already queued. This one has been saved - "
"press Resume on it once the current batch finishes."
),
)
return batch_common.to_out(manifest)
# ---------------------------------------------------------------------------
# The ACTIVE_BRANDS message, in one place
# ---------------------------------------------------------------------------
def _env_line(brand: str) -> str:
names = active_brands.active_display_names()
return "ACTIVE_BRANDS=" + ",".join(list(names) + [brand])
def _inactive_message(result: brand_discovery.DiscoveryResult) -> str:
return (
f"{result.brand} is not in ACTIVE_BRANDS, so these products would be "
f"written to {result.table} and then filtered out of /api/brands, "
f"search, suggest and the category listing. The rows would be complete "
f"and correct, just unreadable. To make them visible, set "
f"'{_env_line(result.brand)}' in backend/.env and restart the API - "
f"settings are read once at import, so a restart is required."
)
def _inactive_detail(brand: str) -> str:
return (
f"{brand} is not in ACTIVE_BRANDS. Ingesting it now would write a "
f"complete catalog that no endpoint can read. Either set "
f"'{_env_line(brand)}' in backend/.env and restart the API first, or "
f"re-send with acknowledge_inactive=true to stage the data anyway - "
f"adding the brand later is a config change and a restart, not a "
f"re-ingest."
)

View File

@@ -48,6 +48,20 @@ def generate_catalog(payload: CatalogGenerateRequest) -> CatalogJobOut:
"""Kick off brand catalog ingestion (discovery -> images -> embeddings ->
pgvector) as a background daemon thread and return immediately with a job id.
PREFER /api/admin/brand-discovery/* FOR NEW WORK. This route is unchanged
and still supported, but it does not run the eleven stages in
`app/core/store_catalog_pipeline.py` - no title validation, no pack-size
explosion, no SKU resolution, no barcode, no HSN/GST, no validation gate -
and it mints `image_id` with `s3_service.generate_image_id()`, which appends
a random uuid4. Nothing it writes can ever match an existing row, so running
it twice for one brand produces two catalogs. It also calls
`upsert_brand_products(cleanup=True)`, which deletes every row not in the
batch it just built.
The discovery routes do run all eleven stages, use a deterministic
`image_id`, write with `cleanup=False`, and show the products for approval
before anything is stored. See docs/BRAND_DISCOVERY.md.
NOTE: on an 8GB RAM / CPU-only machine, running ingestion (which loads
the embeddings model and calls Ollama repeatedly) at the same time as
heavy chat traffic will be slow. This is intended as an occasional

View File

@@ -8,6 +8,18 @@ product discovery (Ollama), per-product image search + S3 upload,
pricing/description enrichment, embedding generation, and the pgvector
upsert. This module exists only to give that pipeline one clear, reusable
entry point and a consistent result shape for callers.
NOT THE ELEVEN-STAGE PIPELINE, and prefer `app/services/brand_discovery.py`
for new work. "Full pipeline" above means this module's own sequence, not the
eleven stages in `app/core/store_catalog_pipeline.py`: there is no title
validation, no pack-size explosion, no SKU service, no barcode, no HSN/GST and
no validation gate here, and `image_id` carries a random uuid4 suffix, so a
second run for the same brand cannot match the first and duplicates it.
This path is unchanged and still works. It also rewrites the brand's seed
catalog via `brand_sync.export_brand_to_seed_file` (below), which the discovery
path deliberately does not - so the two are not drop-in replacements for each
other. See docs/BRAND_DISCOVERY.md.
"""
from __future__ import annotations

View File

@@ -316,6 +316,37 @@ BRAND_SYNC_INTERVAL_SECONDS = int(os.getenv("BRAND_SYNC_INTERVAL_SECONDS", "300"
# See app/services/active_brands.py.
ACTIVE_BRANDS = os.getenv("ACTIVE_BRANDS", "")
# ---------------------------------------------------------------------------
# Brand discovery (brand name -> the 11-stage pipeline)
# ---------------------------------------------------------------------------
# Discovery turns a brand NAME into rows the ordinary catalog pipeline ingests.
# See app/services/brand_discovery.py. Every value below has a working default,
# so the feature needs no configuration to run.
#
# Open Food Facts is the primary source and the language model is the
# supplement, not the reverse: OFF returns real products carrying a real GTIN,
# while the default OLLAMA_MODEL_NAME (qwen2.5:1.5b) invents plausible ones that
# nothing downstream can catch. Turning BRAND_DISCOVERY_USE_OFF off leaves the
# result resting on the model alone.
BRAND_DISCOVERY_USE_OFF = _bool("BRAND_DISCOVERY_USE_OFF", "true")
BRAND_DISCOVERY_USE_LLM = _bool("BRAND_DISCOVERY_USE_LLM", "true")
# Products per discovery run. One CSV row per product; pack-size explosion
# happens later in stage 4, so this is well inside the 2000-row per-file cap.
BRAND_DISCOVERY_MAX_PRODUCTS = int(os.getenv("BRAND_DISCOVERY_MAX_PRODUCTS", "200"))
# Pack sizes kept per product when only the language model offers any. Stage 6
# runs an image search per exploded row, so this multiplies the slowest part of
# the run; 3 keeps a large brand inside a sane wall-clock.
BRAND_DISCOVERY_MAX_SIZES = int(os.getenv("BRAND_DISCOVERY_MAX_SIZES", "3"))
# Wall-clock ceiling on the LLM half of a run, checked between prompts. Open
# Food Facts runs first and is never subject to it, so a run that hits this
# still returns the evidence-backed products.
BRAND_DISCOVERY_DEADLINE_SECONDS = float(
os.getenv("BRAND_DISCOVERY_DEADLINE_SECONDS", "300")
)
# ---------------------------------------------------------------------------
# S3 / DigitalOcean Spaces (product image storage) - optional
# ---------------------------------------------------------------------------

View File

@@ -31,7 +31,7 @@ from app.api.routers import health, brands, search, suggest, chat, catalog, syst
from app.api.routers import stores, discounts, analytics as store_analytics, trending, recommendations, store_admin
from app.api.routers import nutrition, nutrition_admin, upload
from app.api.routers import auth, user_products, admin_train, mcp_info
from app.api.routers import batch_catalog, uploads
from app.api.routers import batch_catalog, uploads, brand_discovery
from app.services.store_db import ensure_store_intelligence_schema
from app.services.nutrition_db import ensure_nutrition_schema
@@ -307,6 +307,7 @@ app.include_router(nutrition_admin.router, prefix="/api")
app.include_router(upload.router, prefix="/api")
app.include_router(batch_catalog.router, prefix="/api")
app.include_router(uploads.router, prefix="/api")
app.include_router(brand_discovery.router, prefix="/api")
app.include_router(mcp_info.router, prefix="/api")
# MCP lives outside /api on purpose: it is a protocol endpoint for AI clients,

File diff suppressed because it is too large Load Diff