diff --git a/.env.example b/.env.example index dd4d9c2..80ef243 100644 --- a/.env.example +++ b/.env.example @@ -339,6 +339,22 @@ ENABLE_SERVER_OCR=true #IMAGE_IDENTIFY_MIN_IMAGE_SCORE=0.70 #IMAGE_IDENTIFY_MIN_TEXT_SCORE=0.60 +# Capture-to-catalog: when identify cannot confirm a photo, read the label and - +# for a brand already in the catalog - add the product through the 11-stage +# pipeline in the background (validation_status=needs_review). The response +# carries a provisional card and a job id to poll at +# GET /api/search/identify/jobs/{id}. OFF by default: it is the only way a +# public route writes to a brand table. See app/services/capture_discovery.py. +ENABLE_CAPTURE_DISCOVERY=false +#CAPTURE_DIR=/app/data/captures +#CAPTURE_QUEUE_MAX=8 +#CAPTURE_MAX_PER_CLIENT_PER_HOUR=20 +# The API's public base URL. Set it to let the colleague's photo become the +# product image when the web search finds none; blank = never. +#CAPTURE_PUBLIC_BASE_URL=https://api.example.com +#CAPTURE_RETAIL_CHECK=true +#CAPTURE_USE_LLM=true + # Product SKU: try a live web search for a real marketplace product ID # (Amazon ASIN, Flipkart PID, etc.) before falling back to an internal SKU. # Set to false to always generate internal SKUs only (faster, offline-safe). diff --git a/app/api/routers/search.py b/app/api/routers/search.py index 248a1e5..211ad45 100644 --- a/app/api/routers/search.py +++ b/app/api/routers/search.py @@ -2,19 +2,23 @@ from __future__ import annotations from typing import Optional -from fastapi import APIRouter, File, Form, HTTPException, Query, UploadFile -from fastapi.responses import JSONResponse +from fastapi import APIRouter, Depends, File, Form, HTTPException, Query, Request, UploadFile +from fastapi.responses import FileResponse, JSONResponse from starlette.concurrency import run_in_threadpool +from app.api.deps import require_admin from app.api.routers.brands import _row_to_product_out from app.api.schemas import ( + CaptureJobOut, IdentifyOut, ImageMatchOut, ImageSearchOut, ImageVectorSearchRequest, + ProvisionalProductOut, SearchOut, SourceProductOut, ) +from app.infrastructure import settings from app.infrastructure.settings import ( IMAGE_SEARCH_DEFAULT_MIN_SCORE, IMAGE_SEARCH_DEFAULT_TOP_K, @@ -23,7 +27,7 @@ from app.infrastructure.settings import ( SEARCH_DEFAULT_TOP_K, SEARCH_MAX_TOP_K, ) -from app.services import image_embedder +from app.services import capture_discovery, image_embedder from app.services.catalog_search import search_catalog from app.services.image_match import ( ImageSearchResult, @@ -241,6 +245,7 @@ async def image_search_endpoint( @router.post("/search/identify", response_model=IdentifyOut) async def identify_endpoint( + request: Request, file: UploadFile = File(..., description="The product photo (JPEG/PNG/WebP), ideally cropped to the pack"), text: Optional[str] = Form(None, max_length=500, description="OCR text read off the label; when absent the server reads it"), @@ -290,4 +295,91 @@ async def identify_endpoint( "deployment. Send the label as `text`, or embed the photo client-side and POST " "the vector to /api/search/image-vector.", ) - return _to_identify_out(result) + out = _to_identify_out(result) + if settings.ENABLE_CAPTURE_DISCOVERY and not capture_discovery.is_confirmed( + result.matched_by, result.fallback_reason + ): + outcome = await run_in_threadpool( + capture_discovery.handle_miss, + label_text=result.ocr_text, image_bytes=content, vector=vector, + brand=brand, category=category, + client=request.client.host if request.client else "unknown", + ) + _apply_capture_outcome(out, outcome) + return out + + +def _apply_capture_outcome(out: IdentifyOut, outcome: capture_discovery.CaptureOutcome) -> None: + """Fold a capture-to-catalog decision into the identify response. + + The low-confidence rows the ladder returned are dropped whenever discovery + has an answer of its own: they are OTHER products (another Godrej line, a + rival detergent), and showing them beside "adding it now" invites the + colleague to pick the wrong one. + """ + out.discovery_status = outcome.status + out.discovery_job_id = outcome.job_id + out.discovery_message = outcome.message + out.provisional = ProvisionalProductOut(**outcome.provisional) if outcome.provisional else None + if outcome.status == capture_discovery.EXISTS and outcome.existing: + row = outcome.existing + card = _row_to_product_out(row, row.get("brand") or "") + out.results = [ImageMatchOut(**card.model_dump(), score=1.0, text_overlap=1.0)] + out.total = 1 + out.matched_by = capture_discovery.MATCHED_BY_LABEL_EXACT + elif outcome.status == capture_discovery.PENDING: + out.results = [] + out.total = 0 + out.matched_by = capture_discovery.MATCHED_BY_DISCOVERY + + +@router.get("/search/identify/jobs/{job_id}", response_model=CaptureJobOut) +def capture_job_endpoint(job_id: str) -> CaptureJobOut: + """A capture-to-catalog job started by /search/identify. + + Poll until `status` is terminal (done, rejected, failed, interrupted). On + done, `product` is the stored catalog row - validation_status needs_review. + """ + job = capture_discovery.get_job(job_id) + if job is None: + raise HTTPException(status_code=404, detail="No capture job with that id.") + return _to_capture_job_out(job) + + +def _to_capture_job_out(job: capture_discovery.CaptureJob) -> CaptureJobOut: + from app.services.vector_store import get_product_by_image_id + + product = None + validation_status = None + if job.status == capture_discovery.DONE and job.image_id: + row = get_product_by_image_id(job.parent, job.image_id) + if row: + product = _row_to_product_out(row, row.get("brand") or job.parent) + validation_status = row.get("validation_status") + return CaptureJobOut( + job_id=job.job_id, status=job.status, created_at=job.created_at, + updated_at=job.updated_at, provisional=job.provisional, product=product, + image_id=job.image_id, disposition=job.disposition, + validation_status=validation_status, retail_presence=job.retail_presence, + photo_used_as_image=job.photo_used_as_image, detail=job.detail, + warnings=job.warnings, + ) + + +@router.get("/admin/captures", response_model=list[CaptureJobOut], + dependencies=[Depends(require_admin)]) +def list_capture_jobs_endpoint(limit: int = Query(50, ge=1, le=200)) -> list[CaptureJobOut]: + """Recent capture-to-catalog jobs, newest first - the review list for + products colleagues added from the field.""" + return [_to_capture_job_out(job) for job in capture_discovery.list_jobs(limit)] + + +@router.get("/search/captures/{name}", include_in_schema=False) +def capture_photo_endpoint(name: str): + """A colleague's capture photo - served only because a product with no web + image uses it as its image (CAPTURE_PUBLIC_BASE_URL).""" + found = capture_discovery.photo_path(name) + if found is None: + raise HTTPException(status_code=404, detail="No such photo.") + path, media_type = found + return FileResponse(path, media_type=media_type) diff --git a/app/api/schemas.py b/app/api/schemas.py index 9901514..9b3e486 100644 --- a/app/api/schemas.py +++ b/app/api/schemas.py @@ -313,11 +313,48 @@ class IdentifyOut(ImageSearchOut): "image_vector" with no `fallback_reason`; otherwise `fallback_reason` says why the best effort shown is unconfirmed (see product_identify.py). """ - matched_by: str = "none" # "image_vector" | "text" | "none" + matched_by: str = "none" # "image_vector" | "text" | "label_exact" + # | "discovery_pending" | "none" ocr_text: Optional[str] = None # the label text the ladder used ocr_source: Optional[str] = None # "client" | "server" image_top_score: Optional[float] = None fallback_reason: Optional[str] = None + # Capture-to-catalog (ENABLE_CAPTURE_DISCOVERY). All None when the answer + # was confirmed or the feature is off. See app/services/capture_discovery.py. + discovery_status: Optional[str] = None # "pending" | "exists" | "needs_input" | "busy" + discovery_job_id: Optional[str] = None # poll GET /api/search/identify/jobs/{id} + discovery_message: Optional[str] = None # one line to show the colleague + provisional: Optional["ProvisionalProductOut"] = None + + +class ProvisionalProductOut(BaseModel): + """What the label says, before the pipeline has run. Not a catalog row.""" + brand: str + product_name: str + size: Optional[str] = None + category: Optional[str] = None + hsn_code: Optional[str] = None + gst_percent: Optional[float] = None + hsn_gst_needs_review: Optional[bool] = None + visible_in_search: bool = True # False: brand is not in ACTIVE_BRANDS + source: str = "label" + + +class CaptureJobOut(BaseModel): + """GET /search/identify/jobs/{job_id}: one capture-to-catalog job.""" + job_id: str + status: str # queued | running | done | rejected | failed | interrupted + created_at: float + updated_at: float + provisional: ProvisionalProductOut + product: Optional[ProductOut] = None # the stored row, once done + image_id: Optional[str] = None + disposition: Optional[str] = None # inserted | backfilled | unchanged + validation_status: Optional[str] = None + retail_presence: Optional[dict] = None + photo_used_as_image: bool = False + detail: Optional[str] = None + warnings: List[str] = Field(default_factory=list) # --------------------------------------------------------------------------- diff --git a/app/infrastructure/settings.py b/app/infrastructure/settings.py index a188189..d587214 100644 --- a/app/infrastructure/settings.py +++ b/app/infrastructure/settings.py @@ -469,6 +469,34 @@ OCR_NUM_THREADS = int(os.getenv("OCR_NUM_THREADS", "2")) # bounded alike. OCR_MAX_CHARS = int(os.getenv("OCR_MAX_CHARS", "500")) +# --------------------------------------------------------------------------- +# Capture-to-catalog - what /search/identify does when the photo is a product +# the catalog does not have (app/services/capture_discovery.py) +# --------------------------------------------------------------------------- +# OFF by default. When on, an unconfirmed identify reads the label, and - for a +# brand we already know - queues the product through the 11-stage pipeline and +# answers at once with a provisional card and a job id to poll. The row lands +# with validation_status=needs_review. This is the only path by which a public +# route WRITES to a brand table, which is why it is a switch. +ENABLE_CAPTURE_DISCOVERY = _bool("ENABLE_CAPTURE_DISCOVERY", "false") +# Photos and job records. Under DATA_DIR so a container keeps them on the volume. +CAPTURE_DIR = _dir("CAPTURE_DIR", DATA_DIR / "captures") +# Jobs waiting behind the one capture worker. Each runs all 11 stages - web +# image search and a live retail lookup included - so a minute or more apiece. +CAPTURE_QUEUE_MAX = int(os.getenv("CAPTURE_QUEUE_MAX", "8")) +# New discovery jobs one client (by IP) may start per hour. The route is public. +CAPTURE_MAX_PER_CLIENT_PER_HOUR = int(os.getenv("CAPTURE_MAX_PER_CLIENT_PER_HOUR", "20")) +# Absolute base of THIS API as the outside world reaches it, e.g. +# https://api.example.com. Needed only to show the colleague's own photo as the +# product image when the web search found none: image URLs must be absolute +# (vector_store._usable_image_url drops a bare path). Blank = never use it. +CAPTURE_PUBLIC_BASE_URL = os.getenv("CAPTURE_PUBLIC_BASE_URL", "").strip() +# A live retail-presence lookup (retail_presence.check_listing, live=True) for +# each discovered product: the pack-size check that catches a misread label. +CAPTURE_RETAIL_CHECK = _bool("CAPTURE_RETAIL_CHECK", "true") +# Stage 2's LLM description. Off falls back to the factual template. +CAPTURE_USE_LLM = _bool("CAPTURE_USE_LLM", "true") + # --------------------------------------------------------------------------- # USDA FoodData Central - nutrition for loose, unbranded commodities # --------------------------------------------------------------------------- diff --git a/app/services/brand_registry.py b/app/services/brand_registry.py index b8a466e..03da3a7 100644 --- a/app/services/brand_registry.py +++ b/app/services/brand_registry.py @@ -364,9 +364,36 @@ def resolve_parent_brand(brand: str) -> str: for alias, parent in BRAND_ALIASES.items(): if _contains_word(key, alias) or _contains_word(alias, key): return parent + prefix = _known_leading_brand(key) + if prefix: + return resolve_parent_brand(prefix) return brand +def _known_leading_brand(key: str) -> Optional[str]: + """The longest leading run of words in `key` that is itself a known brand. + + "Godrej Fab" matched nothing above: no alias contains it, and it contains + no alias ("godrej no.1" and friends are all longer than "godrej"). So it + fell through to the identity and a field capture would have built a + brand_godrej_fab table beside brand_godrej. The same held for "Amul Taaza" + and "Hindustan Unilever Rin" - a parent name followed by a product line + nobody had registered. + + Only reached when every rule above failed, so no input that resolves today + changes its answer. The prefix must be a PARENT or an exact ALIAS KEY, + never a fuzzy hit: letting "fab" through would send "Fab Detergent" to + Parle via the alias "parle fab". + """ + words = key.split() + known = set(BRAND_ALIASES) | set(BRAND_ALIASES.values()) + for n in range(len(words) - 1, 0, -1): + candidate = " ".join(words[:n]) + if candidate in known: + return candidate + return None + + def get_known_sub_brands(brand: str) -> list[str]: """Return known sub-brand/product names for a given brand from BRAND_ALIASES. diff --git a/app/services/capture_discovery.py b/app/services/capture_discovery.py new file mode 100644 index 0000000..7c8fa9f --- /dev/null +++ b/app/services/capture_discovery.py @@ -0,0 +1,748 @@ +"""Capture-to-catalog: a photographed product the catalog does not have. + +THE FLOW +-------- +/search/identify runs its read-only ladder (product_identify.py). When that +answer is NOT confirmed and ENABLE_CAPTURE_DISCOVERY is on, the router hands +the result here: + + label text ─► parse_label ─► brand, product name, pack size, category + │ unreadable / unknown brand ─► status "needs_input" (never "not found") + ▼ + dedup: the same image_id already stored ─► status "exists" + that product + the same product already queued ─► that job's id + ▼ + one-row CSV ─► capture worker ─► store_catalog_pipeline.run_pipeline + (all 11 stages: LLM description, category, pack size, + pricing, web images, SKU, barcode, HSN/GST, validation, + embed + upsert) + ▼ + post-pass on the INSERTED row only: + validation_status = needs_review (a rejection is left alone) + field_sources.capture = {origin, capture id, retail presence verdict} + no web image found ─► the colleague's photo becomes the image + (only with CAPTURE_PUBLIC_BASE_URL set) + +The caller gets a provisional card (built from the label, instantly) and a job +id; GET /api/search/identify/jobs/{id} returns the stored product when done. + +WHY THE PHOTO IS TRUSTED, AND WHAT IS NOT +----------------------------------------- +Most of this catalogue is LLM output nobody grounded. A photo taken in a shop +is the opposite: first-hand evidence the product exists. What can still be +wrong is the READING - OCR misses a letter, or the pack size. That is why the +live retail check matters (retail_presence matches the pack size, not only the +title) and why the row is needs_review rather than verified. + +WHY A WORKER OF ITS OWN, NOT batch_worker +----------------------------------------- +The post-pass has to run after the pipeline stores the row, and the batch +worker has no hook for that. The pipeline itself is reused unchanged - +`run_pipeline` is documented as "called on a background daemon thread". One +thread, one bounded queue, the same shape as image_vector's worker. + +The row is written ONLY for a brand we already know (a registry parent or an +existing brand table): an OCR misread must not be able to mint a brand table. +""" +from __future__ import annotations + +import json +import logging +import queue +import re +import threading +import time +import uuid +from dataclasses import asdict, dataclass, field +from pathlib import Path +from typing import Any, Callable, Dict, List, Optional, Tuple + +from app.infrastructure import settings + +logger = logging.getLogger(__name__) + +# Job states. +QUEUED = "queued" +RUNNING = "running" +DONE = "done" +REJECTED = "rejected" # the validation gate refused the product +FAILED = "failed" +INTERRUPTED = "interrupted" # the process restarted while it was queued/running +TERMINAL = {DONE, REJECTED, FAILED, INTERRUPTED} + +# What the identify response says about discovery (IdentifyOut.discovery_status). +PENDING = "pending" # a job was queued (or one already was) +EXISTS = "exists" # the product is stored; the ladder just missed it +NEEDS_INPUT = "needs_input" # the label could not be turned into a product +BUSY = "busy" # rate limit or full queue; provisional card only + +ORIGIN = "field_capture" +MATCHED_BY_DISCOVERY = "discovery_pending" +MATCHED_BY_LABEL_EXACT = "label_exact" + +_JOB_ID = re.compile(r"^[0-9a-f]{32}$") +_PHOTO_TYPES = { + "jpg": ("image/jpeg", b"\xff\xd8\xff"), + "png": ("image/png", b"\x89PNG"), + "webp": ("image/webp", b"RIFF"), +} +# Words that finish a product name when they follow a category keyword: +# "Detergent Powder", "Detergent Bar", "Dishwash Liquid". +_FORM_WORDS = {"powder", "bar", "liquid", "gel", "cake", "pods", "matic", "paste", "spray", "soap"} +_MAX_NAME_WORDS = 6 + + +# --------------------------------------------------------------------------- +# Label -> product +# --------------------------------------------------------------------------- +@dataclass +class LabelParse: + brand: str = "" # the display form written to the CSV ("Godrej") + parent: str = "" # resolve_parent_brand(brand) - the table key + product_name: str = "" # without the pack size ("Godrej Fab Detergent Powder") + size: str = "" # compact, normalised ("1kg"); "" when not on the label + category: str = "" # "" leaves stage 3 to decide + hsn_code: Optional[str] = None + gst_percent: Optional[float] = None + hsn_gst_needs_review: Optional[bool] = None + label_text: str = "" + problem: Optional[str] = None # set when the label is unusable; the message to show + + @property + def usable(self) -> bool: + return self.problem is None + + +def _compact_size(size: str) -> str: + return re.sub(r"(\d)\s+(?=[a-z])", r"\1", size.strip().lower()) + + +def _known_brand_names() -> List[str]: + """Every name that identifies a brand on its own: registry parents plus the + exact alias keys, longest first so "hindustan unilever" beats "hindustan".""" + from app.services.brand_registry import BRAND_ALIASES + + names = set(BRAND_ALIASES.values()) | set(BRAND_ALIASES) + return sorted(names, key=len, reverse=True) + + +def _find_brand(text: str) -> Optional[Tuple[str, int]]: + """(brand name as written on the label, its word position) or None. + + Registry parents are tried before `extract_brand_mention` because that + function's brand map is built from ACTIVE brands - Godrej is inactive in + this deployment, so "Godrej fab" would come back empty from it. + """ + lower = text.lower() + for name in _known_brand_names(): + match = re.search(r"(? List[str]: + """Up to _MAX_NAME_WORDS words from `start`, stopping at the first number.""" + out: List[str] = [] + for word in words[start:]: + if any(ch.isdigit() for ch in word) or len(out) >= _MAX_NAME_WORDS: + break + out.append(word) + return out + + +def _cut_after_category(window: List[str], brand_words: int) -> Tuple[List[str], Optional[str]]: + """Trim marketing copy: end the name at its category keyword (plus one form + word), so "Godrej fab Detergent Powder Superior cleaning" stops at "Powder". + + Never inside the brand: "Dairy Milk" names Cadbury's chocolate, and its + "milk" must not end "Dairy Milk Silk" at the brand as a Dairy product.""" + from app.services.category_registry import detect_category_from_text + + for end in range(brand_words + 1, len(window) + 1): + category = detect_category_from_text(" ".join(window[:end]), exact_only=True) + if category: + if end < len(window) and window[end].lower() in _FORM_WORDS: + end += 1 + return window[:end], category + return window, None + + +def parse_label(text: Optional[str], *, brand_hint: Optional[str] = None, + category_hint: Optional[str] = None) -> LabelParse: + """Turn OCR'd label text into a product. Pure apart from registry lookups.""" + from app.core.store_catalog_pipeline import _SIZE_IN_TITLE + from app.services.brand_registry import resolve_parent_brand + from app.services.category_registry import detect_category_from_text + from app.services.category_units import fix_or_reject_size + from app.services.enrichment.hsn_gst.models import resolve_hsn_gst + from app.services.label_match import clean_label + + raw = (text or "").strip() + parse = LabelParse(label_text=raw) + if not raw: + parse.problem = ("The label could not be read. Retake the photo closer to the " + "front of the pack, or type the product name.") + return parse + + cleaned = clean_label(raw) or raw + words = cleaned.split() + + if brand_hint and brand_hint.strip(): + brand = brand_hint.strip() + found = _find_brand(cleaned) + position = found[1] if found and resolve_parent_brand(found[0]) == resolve_parent_brand(brand) else None + else: + found = _find_brand(cleaned) + if not found: + parse.problem = ("No brand we know was found on the label. Type the brand and " + "product name so it can be added.") + return parse + brand, position = found + + if position is None: + # The caller named the brand but it is not on the label text: the name + # is the label's leading words, prefixed with the brand. + tail = _name_window(words, 0) + window = brand.split() + [w for w in tail if w.lower() not in brand.lower().split()] + else: + window = _name_window(words, position) + brand_words = len(brand.split()) + window, category = _cut_after_category(window, brand_words) + + if len(window) <= brand_words: + parse.problem = (f"Only the brand ({brand}) could be read. Type the product name " + "so it can be added.") + parse.brand, parse.parent = brand, resolve_parent_brand(brand) + return parse + + name = " ".join(w if w.isupper() and len(w) > 3 else w.capitalize() for w in window) + # The brand as the label wrote it, but in title case - "godrej" -> "Godrej". + display_brand = " ".join(w.capitalize() if w.islower() else w for w in brand.split()) + + category = category or detect_category_from_text(cleaned, exact_only=True) or (category_hint or "") + + size = "" + size_match = _SIZE_IN_TITLE.search(cleaned) + if size_match: + fixed, _changed, _reason = fix_or_reject_size(_compact_size(size_match.group(0)), category, name) + size = _compact_size(fixed) if fixed else "" + + parse.brand = display_brand + parse.parent = resolve_parent_brand(display_brand) + parse.product_name = name + parse.size = size + parse.category = category + if category: + info = resolve_hsn_gst(category, name) + if info is not None: + parse.hsn_code = info.hsn_code + parse.gst_percent = info.gst_percent + parse.hsn_gst_needs_review = info.hsn_gst_needs_review + return parse + + +def is_known_brand(parent: str, table_exists: Optional[Callable[[str], bool]] = None) -> bool: + """D4: only a brand we already know may receive a row - a registry parent, + or a brand table that already exists (active or not).""" + from app.services.brand_registry import BRAND_ALIASES + + key = (parent or "").strip().lower() + if not key: + return False + if key in set(BRAND_ALIASES.values()): + return True + return bool((table_exists or _brand_table_exists)(parent)) + + +def _brand_table_exists(parent: str) -> bool: + from app.services.vector_store import _connect, _table_exists, _table_name + + conn = _connect() + if conn is None: + return False + try: + with conn.cursor() as cur: + return bool(_table_exists(cur, _table_name(parent))) + except Exception: # noqa: BLE001 - "unknown" is the safe answer + return False + finally: + conn.close() + + +def expected_image_id(parse: LabelParse) -> str: + from app.core.store_catalog_pipeline import build_image_id + + return build_image_id(parse.parent, parse.product_name, parse.size) + + +def provisional_card(parse: LabelParse) -> Dict[str, Any]: + """What the colleague sees immediately, from the label alone.""" + from app.services import active_brands + + return { + "brand": parse.brand, + "product_name": parse.product_name, + "size": parse.size or None, + "category": parse.category or None, + "hsn_code": parse.hsn_code, + "gst_percent": parse.gst_percent, + "hsn_gst_needs_review": parse.hsn_gst_needs_review, + # False means: the row will be written, but search filters the brand + # out until it is added to ACTIVE_BRANDS. + "visible_in_search": (not active_brands.filtering_enabled()) + or active_brands.is_active_brand(parse.parent), + "source": "label", + } + + +# --------------------------------------------------------------------------- +# Jobs: durable JSON under CAPTURE_DIR/jobs +# --------------------------------------------------------------------------- +@dataclass +class CaptureJob: + job_id: str + status: str + created_at: float + updated_at: float + brand: str + parent: str + product_name: str + size: str + category: str + label_text: str + provisional: Dict[str, Any] + photo_ext: Optional[str] = None + submitted_by: Optional[str] = None + image_id: Optional[str] = None + disposition: Optional[str] = None # inserted | backfilled | unchanged + retail_presence: Optional[Dict[str, Any]] = None + photo_used_as_image: bool = False + detail: Optional[str] = None + warnings: List[str] = field(default_factory=list) + + def to_dict(self) -> Dict[str, Any]: + return asdict(self) + + @classmethod + def from_dict(cls, data: Dict[str, Any]) -> "CaptureJob": + known = {k: data[k] for k in cls.__dataclass_fields__ if k in data} + return cls(**known) + + +def _jobs_dir() -> Path: + path = Path(settings.CAPTURE_DIR) / "jobs" + path.mkdir(parents=True, exist_ok=True) + return path + + +def _photos_dir() -> Path: + path = Path(settings.CAPTURE_DIR) / "photos" + path.mkdir(parents=True, exist_ok=True) + return path + + +_state_lock = threading.Lock() +_jobs: Dict[str, CaptureJob] = {} +_live_ids: set = set() # queued or running in THIS process +_inflight: Dict[str, str] = {} # expected image_id -> job_id + + +def _save(job: CaptureJob) -> None: + job.updated_at = time.time() + path = _jobs_dir() / f"{job.job_id}.json" + tmp = path.with_suffix(".tmp") + tmp.write_text(json.dumps(job.to_dict(), ensure_ascii=False), encoding="utf-8") + tmp.replace(path) + with _state_lock: + _jobs[job.job_id] = job + + +def get_job(job_id: str) -> Optional[CaptureJob]: + if not _JOB_ID.match(job_id or ""): + return None + with _state_lock: + job = _jobs.get(job_id) + live = job_id in _live_ids + if job is None: + path = _jobs_dir() / f"{job_id}.json" + if not path.exists(): + return None + try: + job = CaptureJob.from_dict(json.loads(path.read_text(encoding="utf-8"))) + except Exception: # noqa: BLE001 - a corrupt record is a missing one + logger.warning("capture: unreadable job record %s", path) + return None + if job.status not in TERMINAL and not live: + # Queued or running according to disk, but no worker here owns it: the + # process restarted. Say so rather than letting a client poll forever. + job.status = INTERRUPTED + job.detail = "The server restarted before this job finished. Capture the product again." + _save(job) + return job + + +def list_jobs(limit: int = 50) -> List[CaptureJob]: + paths = sorted(_jobs_dir().glob("*.json"), key=lambda p: p.stat().st_mtime, reverse=True) + out = [] + for path in paths[:limit]: + job = get_job(path.stem) + if job is not None: + out.append(job) + return out + + +def photo_path(name: str) -> Optional[Tuple[Path, str]]: + """(file, media type) for a stored capture photo named ".".""" + stem, _, ext = (name or "").partition(".") + if not _JOB_ID.match(stem) or ext not in _PHOTO_TYPES: + return None + path = _photos_dir() / f"{stem}.{ext}" + return (path, _PHOTO_TYPES[ext][0]) if path.exists() else None + + +def _photo_ext(image_bytes: bytes) -> Optional[str]: + for ext, (_media, magic) in _PHOTO_TYPES.items(): + if image_bytes.startswith(magic): + if ext == "webp" and image_bytes[8:12] != b"WEBP": + continue + return ext + return None + + +def _photo_public_url(job: CaptureJob) -> Optional[str]: + base = settings.CAPTURE_PUBLIC_BASE_URL + if not base or not job.photo_ext: + return None + return f"{base.rstrip('/')}/api/search/captures/{job.job_id}.{job.photo_ext}" + + +# --------------------------------------------------------------------------- +# Rate limit (the route is public) +# --------------------------------------------------------------------------- +_starts: Dict[str, List[float]] = {} + + +def _allow_start(client: str, now: Optional[float] = None) -> bool: + now = time.time() if now is None else now + limit = settings.CAPTURE_MAX_PER_CLIENT_PER_HOUR + with _state_lock: + recent = [t for t in _starts.get(client, []) if now - t < 3600] + if len(recent) >= limit: + _starts[client] = recent + return False + recent.append(now) + _starts[client] = recent + return True + + +# --------------------------------------------------------------------------- +# The miss handler the router calls +# --------------------------------------------------------------------------- +@dataclass +class CaptureOutcome: + status: str + message: Optional[str] = None + job_id: Optional[str] = None + provisional: Optional[Dict[str, Any]] = None + existing: Optional[Dict[str, Any]] = None # the stored row, for EXISTS + + +def is_confirmed(matched_by: str, fallback_reason: Optional[str]) -> bool: + """The rule product_identify.py documents for clients.""" + return matched_by == "text" or (matched_by == "image_vector" and not fallback_reason) + + +def handle_miss( + *, + label_text: Optional[str], + image_bytes: Optional[bytes], + vector: Optional[List[float]], + brand: Optional[str], + category: Optional[str], + client: str, +) -> CaptureOutcome: + """Decide what an unconfirmed identify turns into. Never raises.""" + try: + return _handle_miss(label_text=label_text, image_bytes=image_bytes, vector=vector, + brand=brand, category=category, client=client) + except Exception: # noqa: BLE001 - discovery must never break identify + logger.exception("capture: discovery failed; identify answered without it") + return CaptureOutcome(status=NEEDS_INPUT, + message="The product could not be added automatically. " + "Type the brand and product name.") + + +def _handle_miss(*, label_text, image_bytes, vector, brand, category, client) -> CaptureOutcome: + from app.services.vector_store import get_product_by_image_id + + parse = parse_label(label_text, brand_hint=brand, category_hint=category) + if not parse.usable: + return CaptureOutcome(status=NEEDS_INPUT, message=parse.problem) + card = provisional_card(parse) + if not is_known_brand(parse.parent): + return CaptureOutcome( + status=NEEDS_INPUT, provisional=card, + message=(f"'{parse.brand}' is not a brand in the catalog yet, so it was not " + "added automatically. Ask an admin to add the brand, or correct the " + "brand name if the label was misread."), + ) + + image_id = expected_image_id(parse) + existing = get_product_by_image_id(parse.parent, image_id) + if existing: + return CaptureOutcome(status=EXISTS, provisional=card, existing=existing, + message="This product is already in the catalog.") + + with _state_lock: + running = _inflight.get(image_id) + if running: + return CaptureOutcome(status=PENDING, job_id=running, provisional=card, + message="This product is already being added.") + + if not _allow_start(client): + return CaptureOutcome(status=BUSY, provisional=card, + message="Too many new products from this device in the last " + "hour. The details above are read from the label.") + + now = time.time() + job = CaptureJob( + job_id=uuid.uuid4().hex, status=QUEUED, created_at=now, updated_at=now, + brand=parse.brand, parent=parse.parent, product_name=parse.product_name, + size=parse.size, category=parse.category, label_text=parse.label_text, + provisional=card, submitted_by=client, image_id=image_id, + ) + if image_bytes: + ext = _photo_ext(image_bytes) + if ext: + (_photos_dir() / f"{job.job_id}.{ext}").write_bytes(image_bytes) + job.photo_ext = ext + _save(job) + + with _state_lock: + _live_ids.add(job.job_id) + _inflight[image_id] = job.job_id + try: + _ensure_worker() + _queue.put_nowait((job.job_id, list(vector) if vector is not None else None)) + except queue.Full: + with _state_lock: + _live_ids.discard(job.job_id) + _inflight.pop(image_id, None) + job.status = FAILED + job.detail = "The capture queue was full. Capture the product again in a few minutes." + _save(job) + return CaptureOutcome(status=BUSY, provisional=card, + message="Many products are being added right now. The details " + "above are read from the label; try again shortly.") + return CaptureOutcome(status=PENDING, job_id=job.job_id, provisional=card, + message="Not in the catalog yet - adding it now. Details will " + "fill in within a few minutes.") + + +# --------------------------------------------------------------------------- +# The worker +# --------------------------------------------------------------------------- +_queue: "queue.Queue[Tuple[str, Optional[List[float]]]]" = queue.Queue( + maxsize=max(1, settings.CAPTURE_QUEUE_MAX)) +_worker: Optional[threading.Thread] = None +_worker_lock = threading.Lock() + + +def _ensure_worker() -> None: + global _worker + with _worker_lock: + if _worker is None or not _worker.is_alive(): + _worker = threading.Thread(target=_loop, name="capture-discovery", daemon=True) + _worker.start() + + +def _loop() -> None: + while True: + job_id, vector = _queue.get() + try: + job = get_job(job_id) + if job is not None: + run_job(job, vector) + except Exception: # noqa: BLE001 - one bad job must not kill the worker + logger.exception("capture: job %s crashed", job_id) + finally: + with _state_lock: + _live_ids.discard(job_id) + for key, value in list(_inflight.items()): + if value == job_id: + del _inflight[key] + _queue.task_done() + + +def build_csv(job: CaptureJob) -> Tuple[str, bytes]: + """The one-row sheet the pipeline reads, through the same writer the + brand-discovery ingest uses.""" + from app.services.brand_discovery import DiscoveredProduct, rows_to_csv_bytes + from app.services.vector_store import _sanitize_name + + product = DiscoveredProduct( + brand=job.brand, product_name=job.product_name, title=job.product_name, + category=job.category or "", category_hint="", description="", + size_variants=[job.size] if job.size else [], + ) + name = f"capture-{_sanitize_name(job.parent) or 'brand'}-{job.job_id[:8]}.csv" + return name, rows_to_csv_bytes([product]) + + +def run_job(job: CaptureJob, vector: Optional[List[float]] = None) -> CaptureJob: + """Run the pipeline for one capture and mark what it stored. Synchronous.""" + from app.core import store_catalog_pipeline as pipeline + + job.status = RUNNING + _save(job) + try: + filename, content = build_csv(job) + result = pipeline.run_pipeline(filename, content, use_llm=settings.CAPTURE_USE_LLM, + fetch_images=True) + except Exception as exc: # noqa: BLE001 + logger.exception("capture: pipeline failed for %s", job.product_name) + job.status, job.detail = FAILED, f"The catalog pipeline failed: {exc}" + _save(job) + return job + + job.warnings = list(result.warnings)[:20] + if result.storage_error: + job.status, job.detail = FAILED, f"Storing the product failed: {result.storage_error}" + _save(job) + return job + if not result.products: + reasons = "; ".join(r.get("reason", "") for r in result.rejections) or \ + "; ".join(e.error for e in result.errors) + job.status = REJECTED + job.detail = f"The product was not added: {reasons or 'it failed validation'}" + _save(job) + return job + + stored = result.products[0] + job.image_id = stored.get("image_id") or job.image_id + job.disposition = stored.get("disposition") + + if job.disposition == "inserted": + job.retail_presence = _retail_check(job) + fallback = _photo_public_url(job) + applied_photo = mark_captured_row(job, fallback) + job.photo_used_as_image = applied_photo + if applied_photo and vector is not None: + _store_photo_vector(job, vector, fallback) + else: + # The row existed after all (a race, or the label named a product the + # ladder missed). It keeps whatever verdict it had - a field capture + # never downgrades a row it did not create. + job.detail = "The product was already in the catalog; nothing was changed." + + job.status = DONE + _save(job) + return job + + +def _retail_check(job: CaptureJob) -> Optional[Dict[str, Any]]: + if not settings.CAPTURE_RETAIL_CHECK: + return None + try: + from app.services.retail_presence import check_listing + + evidence = check_listing(job.parent, job.product_name, job.size, live=True) + return { + "status": evidence.status, + "retailer": evidence.retailer, + "url": evidence.url, + "matched_title": evidence.matched_title, + "matched_size": evidence.matched_size, + } + except Exception as exc: # noqa: BLE001 - corroboration is best-effort + logger.warning("capture: retail check failed for %s: %s", job.product_name, exc) + return {"status": "unknown", "error": str(exc)} + + +def mark_captured_row(job: CaptureJob, fallback_image_url: Optional[str]) -> bool: + """needs_review + capture provenance on the row this job INSERTED; the photo + as the image only where the pipeline stored none. True when the photo was + applied. A rejected verdict is never lifted to needs_review.""" + from psycopg.types.json import Json + + from app.services.vector_store import _connect, _table_name + + provenance = {"capture": { + "origin": ORIGIN, + "capture_id": job.job_id, + "captured_at": job.created_at, + "label_text": job.label_text[:500], + "retail_presence": job.retail_presence, + }} + conn = _connect() + if conn is None: + job.warnings.append("could not mark the row needs_review: no database connection") + return False + table = _table_name(job.parent) + try: + with conn.cursor() as cur: + cur.execute( + f"UPDATE {table} SET " + f"validation_status = CASE WHEN validation_status = 'rejected' " + f"THEN validation_status ELSE 'needs_review' END, " + f"field_sources = COALESCE(field_sources, '{{}}'::jsonb) || %s " + f"WHERE image_id = %s", + (Json(provenance), job.image_id), + ) + applied = False + if fallback_image_url: + cur.execute( + f"UPDATE {table} SET image_url = %s, image_urls = ARRAY[%s]::text[] " + f"WHERE image_id = %s AND COALESCE(image_url, '') = '' " + f"AND COALESCE(cardinality(image_urls), 0) = 0", + (fallback_image_url, fallback_image_url, job.image_id), + ) + applied = cur.rowcount > 0 + return applied + except Exception as exc: # noqa: BLE001 + logger.exception("capture: marking %s failed", job.image_id) + job.warnings.append(f"could not mark the row needs_review: {exc}") + return False + finally: + conn.close() + + +def _store_photo_vector(job: CaptureJob, vector: List[float], src: str) -> None: + """The photo IS the primary image now, so its vector is the row's current + one - img_vector_src names the URL those exact bytes are served from.""" + from app.services import image_vector + from app.services.vector_store import _connect, _table_name + + conn = _connect() + if conn is None: + return + try: + with conn.cursor() as cur: + if image_vector.table_has_columns(cur, _table_name(job.parent)): + image_vector.store_vector(cur, _table_name(job.parent), job.image_id, vector, src) + except Exception: # noqa: BLE001 - the worker will compute it from the URL later + logger.exception("capture: storing the photo vector for %s failed", job.image_id) + finally: + conn.close() diff --git a/app/services/category_registry.py b/app/services/category_registry.py index 96b7825..da18ed1 100644 --- a/app/services/category_registry.py +++ b/app/services/category_registry.py @@ -108,7 +108,14 @@ CATEGORY_REGISTRY: List[Dict[str, object]] = [ {"category": "Hair Care", "keywords": ["shampoo", "shampooo", "conditioner", "hair oil"], "generic_term": "hair care product"}, {"category": "Bath Soap", "keywords": ["bath soap", "soap bar", "soap", "soaps"], "generic_term": "soap"}, {"category": "Skin & Bath Care", "keywords": ["face wash", "body lotion", "skin cream", "moisturizer", "body wash", "cream", "lotion"], "generic_term": "skin care product"}, - {"category": "Household Cleaning", "keywords": ["detergent", "laundry", "dishwash", "floor cleaner", "handwash", "cleaner"], "generic_term": "cleaning product"}, + # Laundry has its own category because the rest of the system already + # names it: title_validator maps "detergent" to it, and HSN_GST_TABLE gives + # it 3402 at 18% with no review flag. Folding it into "Household Cleaning" + # made every detergent land on that table's needs_review row instead. + # "fab" and bare "wash" are deliberately NOT keywords - Parle Fab is a + # biscuit, and "wash" is inside face wash and body wash. + {"category": "Detergents & Fabric Care", "keywords": ["detergent", "detergents", "washing powder", "detergent powder", "detergent bar", "laundry", "fabric wash", "fabric care"], "generic_term": "detergent"}, + {"category": "Household Cleaning", "keywords": ["dishwash", "floor cleaner", "handwash", "cleaner"], "generic_term": "cleaning product"}, {"category": "Fragrance & Deodorants", "keywords": ["deodorant", "deo spray", "perfume", "fragrance", "body spray", "deo"], "generic_term": "fragrance product"}, {"category": "Household - Agarbatti", "keywords": ["agarbatti", "incense sticks", "incense stick"], "generic_term": "agarbatti"}, {"category": "Household - Lamp Oil", "keywords": ["lamp oil"], "generic_term": "lamp oil"}, diff --git a/docs/CAPTURE_TO_CATALOG_WORKFLOW.md b/docs/CAPTURE_TO_CATALOG_WORKFLOW.md new file mode 100644 index 0000000..4853f81 --- /dev/null +++ b/docs/CAPTURE_TO_CATALOG_WORKFLOW.md @@ -0,0 +1,219 @@ +# Capture-to-Catalog: workflow + +**What this covers:** what happens when a colleague photographs a product, +including one the catalogue does not have yet (for example *Godrej Fab*, +Detergents & Fabric Care, missing from `brand_godrej`). + +**The requirement:** the colleague must never be told "No such product found". +The system reads the label, researches the product, adds a proper record to the +brand table, and returns the details. + +**Status (2026-09-24):** built and tested. It is **off by default** +(`ENABLE_CAPTURE_DISCOVERY=false`) until the rollout steps below are done. + +--- + +## 1. Before and after + +| | Before | After (flag on) | +|---|---|---| +| Product in catalogue | Returned by image or label match | Same, unchanged | +| Product NOT in catalogue | Returned the *closest other* products (e.g. a rival detergent) marked unconfirmed. Nothing was added. | Instant card from the label + background job that adds the product. The colleague gets the full record when it's ready. | +| Label unreadable | Best-effort list, unconfirmed | A clear ask: "retake closer, or type the product name" | +| Brand never seen before | Best-effort list | A clear ask; **no brand table is created from a photo** | + +**Was it a cache problem? No.** Identify results are never cached. The +real reasons a product could stay invisible are fixed or explained in §6. + +--- + +## 2. The workflow, step by step + +``` + Colleague takes photo (Nearle app) + │ + ▼ + POST /api/search/identify + │ + ┌────────┴──────────────────────────────────────────────┐ + │ STEP 1 Identify (read-only, as before) │ + │ a. Image match: photo vector vs img_vector (≥ 0.70) │ + │ b. Label match: OCR text → catalogue text search │ + └────────┬──────────────────────────────────────────────┘ + │ + confirmed? ── yes ──► return the product. DONE. + │ no + ▼ + ┌───────────────────────────────────────────────────────┐ + │ STEP 2 Read the label (capture_discovery.parse_label)│ + │ brand "Godrej" (from the brand word) │ + │ product "Godrej Fab Detergent Powder" │ + │ pack size "1kg" │ + │ category "Detergents & Fabric Care" │ + │ HSN / GST 3402 / 18% │ + └────────┬──────────────────────────────────────────────┘ + │ + unreadable, or brand not in catalogue ──► needs_input + │ ("type the brand / name") + ▼ + ┌───────────────────────────────────────────────────────┐ + │ STEP 3 Duplicate check │ + │ same image_id already stored ──► exists: return it │ + │ same product already queued ──► that job's id │ + │ device over hourly limit ──► busy + card │ + └────────┬──────────────────────────────────────────────┘ + ▼ + ┌───────────────────────────────────────────────────────┐ + │ STEP 4 Reply at once (colleague is not kept waiting) │ + │ matched_by "discovery_pending" │ + │ provisional card (from step 2) + discovery_job_id │ + └────────┬──────────────────────────────────────────────┘ + │ background worker, one job at a time + ▼ + ┌───────────────────────────────────────────────────────┐ + │ STEP 5 The existing 11-stage pipeline, on one row │ + │ 1 brand + FSSAI 2 LLM description │ + │ 3 title/category 4 pack size │ + │ 5 price band 6 web images (identity-checked) │ + │ 7 SKU 8-9 barcode, HSN/GST, content │ + │ 10 validation gate 11 embed + write to brand_godrej │ + └────────┬──────────────────────────────────────────────┘ + ▼ + ┌───────────────────────────────────────────────────────┐ + │ STEP 6 Mark the new row (only if it was INSERTED) │ + │ live retail check: is this exact pack sold online? │ + │ validation_status = needs_review │ + │ field_sources.capture = origin, job, retail verdict │ + │ no web image? → colleague's photo becomes the image │ + └────────┬──────────────────────────────────────────────┘ + ▼ + App polls GET /api/search/identify/jobs/{job_id} + → status "done" + the full stored product +``` + +### Timing the colleague sees +- **~1–3 s:** the identify response with the provisional card + (brand, name, size, category, HSN) read from the label. +- **About a minute or more:** the full record: description, images, price + band, SKU and tax fields. It is slower when jobs queue up, because one runs at + a time. +- **Next photo of the same product:** found straight away through the label + match. It also matches on the photo itself once the image vector exists. + +--- + +## 3. What the app receives + +**Identify response, new fields** (additive; existing fields unchanged): + +| Field | Meaning | +|---|---| +| `discovery_status` | `pending` / `exists` / `needs_input` / `busy`; `null` when the answer was already confirmed or the feature is off | +| `discovery_job_id` | Id to poll (only for `pending`) | +| `discovery_message` | One line to show the colleague | +| `provisional` | The card read from the label, including `visible_in_search` | + +A **confirmed** answer is now: `matched_by` = `text`, `label_exact`, or +`image_vector` with no `fallback_reason`. + +**Job poll:** `GET /api/search/identify/jobs/{id}` returns +`queued → running → done | rejected | failed | interrupted`. When the status is +`done`, `product` holds the stored catalogue card. `retail_presence` shows +whether a retailer lists that exact pack. + +**Admin review list:** `GET /api/admin/captures` (admin only) shows recent +capture jobs, newest first. + +Full contract: `backend/docs/IMAGE_SEARCH_API.md`, section *"When the product +is not in the catalogue"*. + +--- + +## 4. Decisions taken + +| # | Decision | Why | +|---|---|---| +| D1 | **Instant card + background job**, not one long request | The full pipeline takes minutes. A mobile request would time out. | +| D2 | New rows are **visible immediately, marked `needs_review`** | The photo proves the product exists, but OCR can misread a letter or the pack size. A human confirms later. | +| D3 | The **colleague's photo is used as the image only when the web finds none** | Clean retailer images look better. A product with no image at all is worse than a shelf photo. | +| D4 | **Only brands already in the catalogue** | An OCR misread must never create a junk brand table. Brands outside `ACTIVE_BRANDS` are still written, but stay hidden from search until switched on. | + +**Safety rules built in:** +- A capture never overwrites or downgrades an existing row. The review mark is + applied only to a row the job itself inserted. A `rejected` verdict is never + lifted. +- A product-line word alone ("Fab") is never treated as a brand, because in the + registry "Fab" is Parle's biscuit. +- Per-device hourly limit (`CAPTURE_MAX_PER_CLIENT_PER_HOUR`, default 20) and + a bounded queue (`CAPTURE_QUEUE_MAX`, default 8), since the route is public. +- If discovery errors, identify still answers. It falls back to `needs_input` + and never fails with a 500. + +--- + +## 5. Where the code is + +| Piece | File | +|---|---| +| Label parsing, duplicate check, jobs, worker, row marking | `backend/app/services/capture_discovery.py` (new) | +| Hook into identify, job/photo/admin routes | `backend/app/api/routers/search.py` | +| Response models | `backend/app/api/schemas.py` (`IdentifyOut` fields, `ProvisionalProductOut`, `CaptureJobOut`) | +| Settings | `backend/app/infrastructure/settings.py` (`ENABLE_CAPTURE_DISCOVERY`, `CAPTURE_*`) and `.env.example` | +| Brand fix ("Godrej Fab" → godrej) | `backend/app/services/brand_registry.py` (`_known_leading_brand`) | +| Category fix (detergents → "Detergents & Fabric Care") | `backend/app/services/category_registry.py` | +| Tests | `tests/test_capture_discovery.py`, `tests/test_brand_registry.py`, `tests/test_category_registry.py` | + +The 11-stage pipeline itself (`store_catalog_pipeline.run_pipeline`) is +**reused unchanged**. + +--- + +## 6. Things that hide a product, and their status + +| Issue | Effect | Status | +|---|---|---| +| "Godrej Fab" did not resolve to Godrej | Would have created a separate `brand_godrej_fab` table | **Fixed.** A known brand followed by a product line now joins that brand. | +| Detergents were categorised as "Household Cleaning" | Every detergent got its tax code flagged for review | **Fixed.** They now use "Detergents & Fabric Care" (HSN 3402, 18%, no flag). | +| `ACTIVE_BRANDS` excludes Godrej | Godrej Fab is stored but hidden from search | **Needs a decision:** add Godrej to `ACTIVE_BRANDS` (or leave it blank) and restart | +| Image vector is filled after the write | The next photo matches on the label first, not the image | By design. It catches up within seconds to minutes. | + +--- + +## 7. Known limits + +1. **Label reading is rule-based.** It needs a known brand word on the label. + Heavily stylised or partial labels get `needs_input`, and the colleague types + the name (sent as `text` / `brand`). An LLM fallback for messy labels is a + possible later addition. +2. **One vector per product.** When the web finds an image, the image vector + is built from that image, not the colleague's photo. Re-captures are then + matched by the label, which works because the row now exists. Matching on the + colleague's own photo in every case would need a second vector column. +3. **Non-food has no Open Food Facts source.** For detergents, the live retail + check is the only outside confirmation. That is why rows stay `needs_review`. +4. **The photo is shown as the image only when `CAPTURE_PUBLIC_BASE_URL` is set**, + because image URLs must be absolute. It is not uploaded to S3. +5. **Only `POST /search/identify`** (photo upload) triggers discovery. The + vector-only route `/search/image-vector` has no photo and does not. + +--- + +## 8. Rollout steps + +1. **Check production for split brand tables:** + `SELECT table_name FROM information_schema.tables WHERE table_name LIKE 'brand\_%';` + Look for any table named after a known brand plus a product line. This is + the one thing the brand fix could re-route. +2. Decide `ACTIVE_BRANDS`: add Godrej and any other brands colleagues will capture. +3. Set `CAPTURE_PUBLIC_BASE_URL` if the shelf photo should be used as the image + for products the web has no image for. +4. **Smoke test against the local database, not production.** `backend/.env` + currently points at **production**. A throwaway brand cannot be used, because + D4 refuses unknown brands on purpose. So a smoke test always writes a real + `needs_review` row into a real brand table: run it locally, or delete the row + afterwards (its `image_id` is in the job record). +5. Turn it on: `ENABLE_CAPTURE_DISCOVERY=true`, then restart the API. +6. Tell the Nearle app team about the new fields (§3). The app should show + `provisional` + `discovery_message` and poll the job. +7. Review new rows regularly: `GET /api/admin/captures`, or query + `validation_status = 'needs_review'` with `field_sources ? 'capture'`. diff --git a/docs/DEPLOYMENT_IMPACT_CAPTURE_TO_CATALOG.md b/docs/DEPLOYMENT_IMPACT_CAPTURE_TO_CATALOG.md new file mode 100644 index 0000000..7cca8e8 --- /dev/null +++ b/docs/DEPLOYMENT_IMPACT_CAPTURE_TO_CATALOG.md @@ -0,0 +1,239 @@ +# After deployment: what changes + +**Scope:** the changes made on 2026-09-24: +- brand-alias fix +- detergent category fix +- capture-to-catalog flow + +How the flow works is described in `CAPTURE_TO_CATALOG_WORKFLOW.md`. This +document covers what the running system does differently once these changes +are deployed. + +--- + +## 1. Summary + +The changes fall into two groups: + +| Group | When it takes effect | Needs any action? | +|---|---|---| +| **A. Live on deploy.** Category detection, brand-name resolution, extra response fields, new routes. | As soon as the new build starts | No. Verify with the checks in §6. | +| **B. Dormant.** Capture-to-catalog (adding products from photos). | Only after `ENABLE_CAPTURE_DISCOVERY=true` and a restart | Yes. Decisions and steps in §4. | + +**What the deploy does not change:** +- no database migration or new columns; +- no new container or Python dependency (rapidocr is already in the image); +- no change to existing rows. + +With the flag off, no route writes anything new. + +--- + +## 2. Live on deploy (flag-independent) + +### 2.1 Detergent searches return detergents +**What changed:** +- `detergent`, `laundry`, `washing powder`, `detergent powder`, `detergent bar`, + `fabric wash` and `fabric care` now map to **"Detergents & Fabric Care"**. +- They used to map to "Household Cleaning". + +**Effect on search and chat** (`catalog_search`, `rag_service`): +- A query like "detergent" or "Surf Excel detergent" is limited to the + detected category. +- **Before:** it was limited to "Household Cleaning". The last production + snapshot has about **38 rows** stored as "Detergents & Fabric Care" and + **2** as "Household Cleaning", so the real detergents were **filtered out**. +- **After:** those detergent rows are returned. +- "Household Cleaning" keeps dishwash, floor cleaner, handwash and cleaner. + +**Effect on new uploads:** +- A detergent row gets category "Detergents & Fabric Care" and HSN **3402 at + 18%** with **`hsn_gst_needs_review = false`**. Before, it was flagged for review. +- Stored rows keep their category. Nothing is recategorised. + +**Risk:** if either of the 2 "Household Cleaning" rows is actually a detergent, +a "detergent" search no longer returns it. It is still found by its name or brand. + +### 2.2 "Brand + product line" joins the parent brand +**What changed:** `resolve_parent_brand` now has one more rule, used only +when every older rule found nothing. A name that starts with a known brand +resolves to that brand: + +| Input | Before | After | +|---|---|---| +| Godrej Fab | own table `brand_godrej_fab` | `brand_godrej` | +| Amul Taaza | own table `brand_amul_taaza` | `brand_amul` | +| Hindustan Unilever Rin | own table | `brand_hindustan_unilever` | +| Fab Detergent Powder | unchanged | unchanged (a line name alone is never treated as a brand) | + +Anything that resolved before resolves the same way. All seed catalogs are +pinned by tests. + +**Where it applies:** +- uploads (stage 1 of the pipeline); +- per-brand reads such as `/api/brands/{brand}/products`; +- the S3 folder and the seed-file choice. + +**Risk (not verified, because production could not be read from here):** +production may already have a table created under such a name, for example a +`brand_amul_taaza` from an earlier upload. If so, requests for that brand name +now read and write the parent table instead, and the old table becomes +orphaned. Its rows are not lost, only unreachable by that name. Check this +before deploying (§6, check 1). + +### 2.3 Identify response has four extra keys +`POST /api/search/identify` (and `/search/image-vector` with +`text_fallback`) now always include these keys: + +`discovery_status`, `discovery_job_id`, `discovery_message`, `provisional` + +With the flag off they are all `null`. No existing key changed. The Nearle +app is unaffected unless it rejects unknown JSON keys. + +### 2.4 New routes (visible in the OpenAPI docs) + +| Route | With the flag off | +|---|---| +| `GET /api/search/identify/jobs/{job_id}` | always 404 (no jobs exist) | +| `GET /api/admin/captures` | admin only; returns an empty list | +| `GET /api/search/captures/{name}` | always 404 (no photos exist); hidden from the docs | + +--- + +## 3. Which environment file the server uses + +- **The image:** `backend/Dockerfile` copies `.env.production` into the image + as `.env`. +- **Compose:** `docker-compose.yml` also loads `./backend/.env` from the + server's checkout into the backend container. Values from compose win. + +| Setting | `.env.production` | local `backend/.env` | +|---|---|---| +| `ENABLE_CAPTURE_DISCOVERY` | not set, so **off** | not set, so **off** | +| `ACTIVE_BRANDS` | not set, so **all brands visible** | `Amul, Cadbury, Hindustan Unilever, Own Products` | + +The capture flow is off in production however it is deployed. + +Whether Godrej products appear in search depends on which file the server +actually loads. Confirm with `GET /api/brands` after deploy (§6, check 4). + +--- + +## 4. After `ENABLE_CAPTURE_DISCOVERY=true` + +### 4.1 What happens on each request +Each identify request that **cannot be confirmed** now also does the following: +1. Reads the label into brand, product name, pack size and category. This adds + a few milliseconds. +2. Chooses what to return: + - known brand, new product: queues a job and returns `pending` with a + provisional card. The unconfirmed lookalikes are no longer in `results`. + - product already stored: returns it as a confirmed match (`label_exact`). + - unreadable label or unknown brand: returns `needs_input` with a message. + - rate limit or full queue: returns `busy` with the card. + +Confirmed identifies behave exactly as before. + +### 4.2 What a job does in the background +- **One worker thread** in the backend process runs one job at a time. +- Each job runs the full 11-stage pipeline for one product, which makes + outbound calls: + - an Ollama description, up to 120 s; + - web image search, several providers, validated one by one; + - a live retail-presence search; + - an SKU lookup if `ENABLE_SKU_WEB_LOOKUP` is on. +- **Expect about a minute or more per job.** Later jobs wait in the queue + (up to `CAPTURE_QUEUE_MAX = 8`). +- It then writes **one row** to the brand table: + - `validation_status = needs_review`; + - `field_sources.capture` records the origin, job id, label text and retail verdict. +- The image vector is computed by the existing background worker. If the + colleague's photo became the product image, its vector is stored straight away. + +### 4.3 Resources +| Resource | Impact | +|---|---| +| CPU / RAM | One pipeline run at a time, inside the backend's existing 2.5 GB limit. It can run alongside an admin batch upload; both use Ollama, so each gets slower. | +| Network | Outbound calls to image search, retail search and Ollama for each job | +| Disk | `/app/data/captures/`, on the `catalog_rag_backend_data` volume, so it survives redeploys. One photo (roughly 0.5–3 MB) plus a small JSON record per job. **There is no automatic cleanup yet.** | +| Database | One row per new product, plus two small targeted updates | + +### 4.4 Restarts +- A queued or running job does not survive a restart. When it is next polled + it reports `interrupted`, and the colleague captures again. +- Rows that were already written stay in the table. +- The per-device rate-limit counters reset. + +### 4.5 Safety rules +- **Only known brands.** A registry parent or an existing brand table. + A photo can never create a new brand table. +- **Never overwrites.** The `needs_review` mark and provenance are applied only + to a row the job itself inserted. An existing row is left exactly as it was, + and a `rejected` verdict is never lifted. +- **Public route limits.** Each device (by IP) can start 20 jobs per hour + (`CAPTURE_MAX_PER_CLIENT_PER_HOUR`), and the queue is bounded. Behind a + proxy, several devices can share one IP. +- **Fails safe.** Any error inside discovery turns into `needs_input`. + Identify never returns a 500 because of it. + +--- + +## 5. How to enable it + +1. **Check production for tables split by product line** (see §6, check 1). +2. **Decide `ACTIVE_BRANDS`** for the server. Blank means every brand is + visible. Otherwise list every brand colleagues will photograph, Godrej included. +3. **Optional:** set `CAPTURE_PUBLIC_BASE_URL=https://api.`. This lets + the shelf photo be used as the image for products the web has no image for. +4. Set `ENABLE_CAPTURE_DISCOVERY=true` in the env file the server actually loads + (§3), then restart the backend. +5. **Tell the Nearle app team** about the new response fields and the job poll + (`IMAGE_SEARCH_API.md`, "When the product is not in the catalogue"). +6. **Smoke test with one real product** and delete the row afterwards if it + was only a test. A made-up brand can't be used, because the brand gate + refuses it on purpose. + +--- + +## 6. Post-deploy checks + +| # | Check | Expected | +|---|---|---| +| 1 | *(before deploy)* `SELECT table_name FROM information_schema.tables WHERE table_name LIKE 'brand\_%';` | No table named after a known brand plus a product line (e.g. `brand_amul_taaza`). If one exists, decide whether to merge it first. | +| 2 | `GET /api/health` | Healthy. `ocr` is available (the capture flow needs server OCR when the app sends no `text`). | +| 3 | `GET /api/search?q=detergent` | Detergent products are returned (category "Detergents & Fabric Care"). | +| 4 | `GET /api/brands` | The brand list matches your `ACTIVE_BRANDS` decision. | +| 5 | `POST /api/search/identify` with any photo | The same answer as before, plus four `null` discovery keys. | +| 6 | `GET /api/admin/captures` (admin token) | `[]` | +| 7 | *(after enabling)* Photograph a product that is not in the catalogue | `discovery_status: "pending"`, then the job reaches `done` and `product.validation_status` is `needs_review`. | + +--- + +## 7. Rollback + +| What | How | What remains | +|---|---|---| +| Stop adding products | `ENABLE_CAPTURE_DISCOVERY=false`, then restart | Rows already added, which can be found by the query below | +| Remove the whole change | Redeploy the previous build | The same rows. Nothing needs to be migrated back. | +| Remove the rows it added | `DELETE FROM brand_ WHERE field_sources ? 'capture' AND validation_status = 'needs_review';` (review first with a `SELECT`) | Photos under `/app/data/captures/`, which can be deleted | + +Reverting the category or brand fix affects **new** uploads only. Rows +written while the fix was live keep the values they were given. + +--- + +## 8. Open items + +1. **Photo and job retention.** There is no cleanup job yet, so photos build + up on the data volume. A purge after N days is worth adding before heavy use. +2. **The photo route is public.** It is reachable by anyone who has the + 32-character job id. The ids can't be guessed, but anyone with a job id can + view that photo. +3. **One image vector per product.** When the web supplies the image, the + colleague's photo is not used for image matching. Repeat photos of that + product are matched by reading the label instead. +4. **Rule-based label reading.** It needs a known brand word printed on the + pack. Otherwise the colleague is asked to type the name. +5. **The pipeline itself was not run end to end in tests.** It is tested with + the pipeline, database and web lookups replaced by fakes. The first real run + should be the smoke test in §5, step 6. diff --git a/docs/IMAGE_SEARCH_API.md b/docs/IMAGE_SEARCH_API.md index 1954e5b..a05c28d 100644 --- a/docs/IMAGE_SEARCH_API.md +++ b/docs/IMAGE_SEARCH_API.md @@ -244,6 +244,59 @@ has always been. On, the same ladder runs with the app's vector and `text` (no photo, so never server OCR), and the response is the identify shape above. Not on the GET variant. +## When the product is not in the catalogue (capture-to-catalog) + +With `ENABLE_CAPTURE_DISCOVERY=true` an unconfirmed identify no longer ends +at a best effort. The label is read into brand, product name, pack size and +category, and one of four things comes back in `discovery_status`: + +| `discovery_status` | `matched_by` | What happened | What to show | +|---|---|---|---| +| `pending` | `discovery_pending` | Queued for the 11-stage pipeline. `results` is empty. | `provisional` + `discovery_message`; poll the job | +| `exists` | `label_exact` | The product is stored; the ladder just missed it. `results` holds it. | the product (confirmed) | +| `needs_input` | unchanged | The label could not be read, or the brand is not in the catalogue. | `discovery_message`: ask for the name, resend as `text` / `brand` | +| `busy` | unchanged | Rate limit (per device per hour) or full queue. | `provisional` + `discovery_message` | + +`label_exact` is a confirmed answer. `discovery_status` is `null` when the +answer was already confirmed, or the feature is off. + +```jsonc +{ + "matched_by": "discovery_pending", + "results": [], + "discovery_status": "pending", + "discovery_job_id": "3f2a…", // 32 hex + "discovery_message": "Not in the catalog yet - adding it now. Details will fill in within a few minutes.", + "provisional": { + "brand": "Godrej", "product_name": "Godrej Fab Detergent Powder", "size": "1kg", + "category": "Detergents & Fabric Care", "hsn_code": "3402", "gst_percent": 18, + "hsn_gst_needs_review": false, + "visible_in_search": false, // brand not in ACTIVE_BRANDS: stored, but search hides it + "source": "label" + } +} +``` + +**Polling.** `GET /api/search/identify/jobs/{job_id}` returns `status`: +`queued`, `running`, then one of `done`, `rejected` (the validation gate refused +it; `detail` says why), `failed`, `interrupted` (the server restarted; capture +again). On `done`, `product` is the stored catalogue card with +`validation_status` `needs_review`, and `retail_presence` is the live check +of whether a retailer lists that exact pack. One job runs at a time, and each takes +roughly a minute or more (web images, retail lookup, LLM description). + +**Only known brands.** A brand that is neither in the brand registry nor an +existing brand table is never created from a label, so an OCR misread cannot +mint a brand. A product line on its own ("Fab") is not treated as a brand. + +**The photo.** It is kept under `CAPTURE_DIR`. It becomes the product image +only when the web search found no acceptable image *and* +`CAPTURE_PUBLIC_BASE_URL` is set; it is then served at +`/api/search/captures/{job_id}.jpg` and its vector is stored at once, so the +next photo of that pack matches on the image rung. + +Admins can list recent capture jobs at `GET /api/admin/captures`. + ## Errors | Code | Cause | What to do | diff --git a/tests/test_brand_registry.py b/tests/test_brand_registry.py index 968a121..83d0ea2 100644 --- a/tests/test_brand_registry.py +++ b/tests/test_brand_registry.py @@ -219,6 +219,28 @@ def test_unknown_brand_passes_through_unchanged() -> None: assert resolve_parent_brand("Totally Made Up Brand") == "Totally Made Up Brand" +@pytest.mark.parametrize("brand,parent", [ + ("Godrej Fab", "godrej"), + ("GODREJ FAB", "godrej"), + ("Godrej Ezee", "godrej"), + ("Amul Taaza", "amul"), + ("Hindustan Unilever Rin", "hindustan unilever"), + # The prefix inherits the parent's own answer, quirks included. + ("Sunfeast Farmlite", "itc"), +]) +def test_parent_followed_by_an_unregistered_line_joins_the_parent(brand: str, parent: str) -> None: + """The field-capture case: a label reading "Godrej Fab" must write into + brand_godrej, not build a brand_godrej_fab table beside it.""" + assert resolve_parent_brand(brand) == parent + + +def test_leading_brand_rule_never_uses_a_fuzzy_word() -> None: + """"fab" is only known as part of the alias "parle fab". Accepting it as a + leading brand would send another company's detergent to Parle.""" + assert resolve_parent_brand("Fab Detergent Powder") == "Fab Detergent Powder" + assert resolve_parent_brand("Sun Pharma") == "Sun Pharma" + + def test_new_brand_round_trips_to_its_own_table() -> None: """ITC is the worked example: it must own brand_itc, not merge elsewhere.""" assert resolve_parent_brand("ITC") == "itc" diff --git a/tests/test_capture_discovery.py b/tests/test_capture_discovery.py new file mode 100644 index 0000000..2439b1d --- /dev/null +++ b/tests/test_capture_discovery.py @@ -0,0 +1,331 @@ +"""Capture-to-catalog: a photographed product the catalog does not have. + +No database, no network, no worker thread: the DB reads, the pipeline and the +post-pass are patched, and jobs are written under a tmp CAPTURE_DIR. The +pipeline's own stages are covered by the store-catalog tests; here the +assertions are about the decisions this module makes around it. +""" +from __future__ import annotations + +import io +import math +from typing import Any, Dict, List + +import pytest + +from app.api.routers import search as search_router +from app.core.store_catalog_pipeline import PipelineResult +from app.services import capture_discovery as cd +from app.services.image_match import ImageSearchResult +from app.services.product_identify import IdentifyResult + +GODREJ_FAB_LABEL = ("NEW Godrej fab Detergent Powder Superior cleaning 1kg MRP Rs.110 " + "(incl of all taxes) www.godrejfab.com Mfd by Godrej Consumer Products Ltd") + + +# --------------------------------------------------------------------------- +# Fixtures +# --------------------------------------------------------------------------- +@pytest.fixture(autouse=True) +def isolated(tmp_path, monkeypatch): + """Jobs on tmp, no worker thread, clean in-process state, DB reads patched.""" + monkeypatch.setattr(cd.settings, "CAPTURE_DIR", tmp_path / "captures") + monkeypatch.setattr(cd.settings, "CAPTURE_MAX_PER_CLIENT_PER_HOUR", 20) + monkeypatch.setattr(cd.settings, "CAPTURE_PUBLIC_BASE_URL", "") + monkeypatch.setattr(cd, "_ensure_worker", lambda: None) + monkeypatch.setattr(cd, "_jobs", {}) + monkeypatch.setattr(cd, "_live_ids", set()) + monkeypatch.setattr(cd, "_inflight", {}) + monkeypatch.setattr(cd, "_starts", {}) + import queue as _queue + monkeypatch.setattr(cd, "_queue", _queue.Queue(maxsize=8)) + monkeypatch.setattr(cd, "_brand_table_exists", lambda parent: False) + import app.services.vector_store as vs + monkeypatch.setattr(vs, "get_product_by_image_id", lambda brand, image_id: None) + + +def _miss(label: str = GODREJ_FAB_LABEL, **kw) -> cd.CaptureOutcome: + args = dict(label_text=label, image_bytes=b"\xff\xd8\xff\xe0fakejpeg", vector=None, + brand=None, category=None, client="10.0.0.1") + args.update(kw) + return cd.handle_miss(**args) + + +# --------------------------------------------------------------------------- +# parse_label +# --------------------------------------------------------------------------- +def test_godrej_fab_label_becomes_a_godrej_detergent() -> None: + p = cd.parse_label(GODREJ_FAB_LABEL) + assert p.usable + assert (p.brand, p.parent) == ("Godrej", "godrej") + assert p.product_name == "Godrej Fab Detergent Powder" # marketing copy trimmed + assert p.size == "1kg" + assert p.category == "Detergents & Fabric Care" + assert (p.hsn_code, p.gst_percent, p.hsn_gst_needs_review) == ("3402", 18, False) + + +def test_sub_brand_printed_alone_reaches_its_parent() -> None: + p = cd.parse_label("Surf excel easy wash detergent powder 1 kg") + assert p.parent == "hindustan unilever" + assert p.brand == "Surf Excel" + assert p.size == "1kg" + + +def test_a_brand_that_is_also_a_category_word_is_not_cut_short() -> None: + p = cd.parse_label("Dairy Milk Silk 60g") + assert p.product_name == "Dairy Milk Silk" + assert p.parent == "cadbury" + assert p.category == "Chocolates" + + +def test_bare_line_name_is_never_given_a_brand() -> None: + """"Fab" alone is Parle's biscuit line in the registry; a detergent label + without its brand must ask, not guess.""" + p = cd.parse_label("fab detergent powder 500 g") + assert not p.usable + assert "brand" in p.problem.lower() + + +def test_brand_hint_completes_a_label_without_the_brand() -> None: + p = cd.parse_label("fab detergent powder 500 g", brand_hint="Godrej") + assert p.usable + assert p.product_name == "Godrej Fab Detergent Powder" + assert p.size == "500g" + + +@pytest.mark.parametrize("text", ["", " ", None]) +def test_unreadable_label_asks_for_input(text) -> None: + p = cd.parse_label(text) + assert not p.usable + assert "type the product name" in p.problem.lower() + + +def test_brand_only_label_asks_for_the_product_name() -> None: + p = cd.parse_label("Godrej 1kg") + assert not p.usable + assert "product name" in p.problem.lower() + + +def test_expected_image_id_matches_the_pipeline_key() -> None: + p = cd.parse_label(GODREJ_FAB_LABEL) + assert cd.expected_image_id(p) == "godrej_godrej_fab_detergent_powder_1kg" + + +# --------------------------------------------------------------------------- +# Confirmation rule and brand gate +# --------------------------------------------------------------------------- +@pytest.mark.parametrize("matched_by,reason,confirmed", [ + ("text", "image_below_threshold", True), + ("image_vector", None, True), + ("image_vector", "text_no_match", False), + ("none", "ocr_empty", False), +]) +def test_is_confirmed_follows_the_documented_client_rule(matched_by, reason, confirmed) -> None: + assert cd.is_confirmed(matched_by, reason) is confirmed + + +def test_known_brand_is_a_registry_parent_or_an_existing_table() -> None: + assert cd.is_known_brand("godrej", table_exists=lambda p: False) + assert cd.is_known_brand("idhayam", table_exists=lambda p: True) + assert not cd.is_known_brand("Xyzzy Foods", table_exists=lambda p: False) + assert not cd.is_known_brand("", table_exists=lambda p: True) + + +# --------------------------------------------------------------------------- +# handle_miss +# --------------------------------------------------------------------------- +def test_miss_queues_a_job_with_a_provisional_card() -> None: + out = _miss() + assert out.status == cd.PENDING + assert out.job_id and len(out.job_id) == 32 + assert out.provisional["product_name"] == "Godrej Fab Detergent Powder" + assert out.provisional["source"] == "label" + job = cd.get_job(out.job_id) + assert job.status == cd.QUEUED + assert job.photo_ext == "jpg" + assert cd.photo_path(f"{job.job_id}.jpg") is not None + + +def test_same_product_twice_reuses_the_running_job() -> None: + first = _miss() + second = _miss(client="10.0.0.2") + assert second.status == cd.PENDING + assert second.job_id == first.job_id + + +def test_stored_product_is_returned_not_rediscovered(monkeypatch) -> None: + import app.services.vector_store as vs + row = {"image_id": "godrej_godrej_fab_detergent_powder_1kg", "product_name": "Godrej Fab 1kg"} + monkeypatch.setattr(vs, "get_product_by_image_id", + lambda brand, image_id: row if image_id == row["image_id"] else None) + out = _miss() + assert out.status == cd.EXISTS + assert out.existing is row + assert out.job_id is None + + +def test_unknown_brand_never_mints_a_table() -> None: + out = _miss(label="Xyzzy Soap 100g", brand="Xyzzy") + assert out.status == cd.NEEDS_INPUT + assert out.job_id is None + assert "not a brand in the catalog" in out.message + + +def test_unreadable_label_is_needs_input_never_not_found() -> None: + out = _miss(label=None) + assert out.status == cd.NEEDS_INPUT + assert "not found" not in (out.message or "").lower() + + +def test_rate_limit_answers_busy_with_the_card(monkeypatch) -> None: + monkeypatch.setattr(cd.settings, "CAPTURE_MAX_PER_CLIENT_PER_HOUR", 1) + assert _miss().status == cd.PENDING + out = _miss(label="Amul Taaza Toned Milk 500 ml") + assert out.status == cd.BUSY + assert out.provisional["brand"] == "Amul" + + +def test_discovery_failure_never_breaks_identify(monkeypatch) -> None: + monkeypatch.setattr(cd, "parse_label", lambda *a, **k: 1 / 0) + out = _miss() + assert out.status == cd.NEEDS_INPUT + + +def test_restart_marks_an_orphaned_job_interrupted(monkeypatch) -> None: + out = _miss() + monkeypatch.setattr(cd, "_live_ids", set()) + monkeypatch.setattr(cd, "_jobs", {}) + assert cd.get_job(out.job_id).status == cd.INTERRUPTED + + +@pytest.mark.parametrize("name", ["../etc/passwd", "abc.jpg", "0" * 32 + ".exe", ""]) +def test_photo_path_rejects_anything_but_a_job_photo(name) -> None: + assert cd.photo_path(name) is None + + +# --------------------------------------------------------------------------- +# run_job +# --------------------------------------------------------------------------- +def _queued_job() -> cd.CaptureJob: + return cd.get_job(_miss().job_id) + + +def _result(disposition: str = "inserted") -> PipelineResult: + r = PipelineResult() + r.products = [{"image_id": "godrej_godrej_fab_detergent_powder_1kg", "disposition": disposition}] + return r + + +@pytest.fixture +def pipeline(monkeypatch): + from app.core import store_catalog_pipeline as p + seen: Dict[str, Any] = {"marked": []} + + def fake_run(filename, content, **kw): + seen["csv"] = content.decode("utf-8") + seen["kw"] = kw + return seen["result"] + + seen["result"] = _result() + monkeypatch.setattr(p, "run_pipeline", fake_run) + monkeypatch.setattr(cd, "_retail_check", lambda job: {"status": "found"}) + monkeypatch.setattr(cd, "mark_captured_row", + lambda job, url: seen["marked"].append((job.image_id, url)) or False) + return seen + + +def test_inserted_product_is_marked_for_review(pipeline) -> None: + job = cd.run_job(_queued_job()) + assert job.status == cd.DONE + assert job.disposition == "inserted" + assert job.retail_presence == {"status": "found"} + assert pipeline["marked"] == [("godrej_godrej_fab_detergent_powder_1kg", None)] + assert "Godrej Fab Detergent Powder" in pipeline["csv"] + assert "Detergents & Fabric Care" in pipeline["csv"] + assert pipeline["kw"]["fetch_images"] is True + + +def test_existing_row_is_never_downgraded(pipeline) -> None: + pipeline["result"] = _result("backfilled") + job = cd.run_job(_queued_job()) + assert job.status == cd.DONE + assert pipeline["marked"] == [] + + +def test_rejected_product_reports_why(pipeline) -> None: + r = PipelineResult() + r.rejections = [{"reason": "title does not name a product"}] + pipeline["result"] = r + job = cd.run_job(_queued_job()) + assert job.status == cd.REJECTED + assert "title does not name a product" in job.detail + + +def test_storage_error_fails_the_job(pipeline) -> None: + r = _result() + r.storage_error = "connection refused" + pipeline["result"] = r + job = cd.run_job(_queued_job()) + assert job.status == cd.FAILED + + +def test_photo_url_offered_only_with_a_public_base(pipeline, monkeypatch) -> None: + monkeypatch.setattr(cd.settings, "CAPTURE_PUBLIC_BASE_URL", "https://api.example.com/") + job = cd.run_job(_queued_job()) + url = pipeline["marked"][0][1] + assert url == f"https://api.example.com/api/search/captures/{job.job_id}.jpg" + + +# --------------------------------------------------------------------------- +# HTTP +# --------------------------------------------------------------------------- +def _unit(): + v = [math.cos(i / 7.0) for i in range(1024)] + n = math.sqrt(sum(x * x for x in v)) + return [x / n for x in v] + + +@pytest.fixture +def unconfirmed(monkeypatch): + monkeypatch.setattr(search_router.image_embedder, "available", lambda: True) + monkeypatch.setattr(search_router.image_embedder, "embedding_for_bytes", lambda data: _unit()) + rival = {"image_id": "hul_surf_excel_1kg", "product_name": "Surf Excel 1kg", "brand": "Hindustan Unilever", + "category": "Detergents & Fabric Care", "score": 0.52, "text_overlap": 0.0} + result = IdentifyResult(search=ImageSearchResult(rows=[rival], top_k=10), matched_by="image_vector", + ocr_text=GODREJ_FAB_LABEL, ocr_source="server", image_top_score=0.52, + fallback_reason="text_no_match") + monkeypatch.setattr(search_router, "identify_product", lambda **kw: result) + + +def _post(client): + return client.post("/api/search/identify", + files={"file": ("p.jpg", io.BytesIO(b"\xff\xd8\xff\xe0x"), "image/jpeg")}) + + +def test_identify_with_discovery_off_is_unchanged(client, unconfirmed, monkeypatch) -> None: + monkeypatch.setattr(search_router.settings, "ENABLE_CAPTURE_DISCOVERY", False) + body = _post(client).json() + assert body["discovery_status"] is None + assert body["matched_by"] == "image_vector" + assert body["total"] == 1 + + +def test_identify_miss_starts_discovery(client, unconfirmed, monkeypatch) -> None: + monkeypatch.setattr(search_router.settings, "ENABLE_CAPTURE_DISCOVERY", True) + body = _post(client).json() + assert body["discovery_status"] == "pending" + assert body["matched_by"] == "discovery_pending" + assert body["results"] == [] # the rival detergent is not offered + assert body["provisional"]["product_name"] == "Godrej Fab Detergent Powder" + job = client.get(f"/api/search/identify/jobs/{body['discovery_job_id']}").json() + assert job["status"] == "queued" + assert job["provisional"]["brand"] == "Godrej" + + +def test_job_endpoint_404s_for_an_unknown_id(client) -> None: + assert client.get("/api/search/identify/jobs/" + "0" * 32).status_code == 404 + assert client.get("/api/search/identify/jobs/not-an-id").status_code == 404 + + +def test_admin_capture_list_requires_admin(client) -> None: + assert client.get("/api/admin/captures").status_code in (401, 403) diff --git a/tests/test_category_registry.py b/tests/test_category_registry.py new file mode 100644 index 0000000..87c193c --- /dev/null +++ b/tests/test_category_registry.py @@ -0,0 +1,49 @@ +""" +Laundry products resolve to "Detergents & Fabric Care", the name the rest of +the system already uses (title_validator, HSN_GST_TABLE, consumability). + +The registry used to put detergents under "Household Cleaning", whose HSN row +carries needs_review=True, so every detergent was flagged for review over a +naming mismatch rather than any real doubt about its tax code. + +No database, no network. +""" +from __future__ import annotations + +import pytest + +from app.services.category_registry import ALL_CATEGORIES, detect_category_from_text +from app.services.enrichment.hsn_gst.models import resolve_hsn_gst + + +@pytest.mark.parametrize("title", [ + "Godrej Fab Detergent Powder 1kg", + "Surf Excel Matic Liquid Laundry Detergent", + "Tide Washing Powder", + "Rin Detergent Bar", +]) +def test_laundry_resolves_to_detergents_with_a_clean_hsn(title: str) -> None: + category = detect_category_from_text(title) + assert category == "Detergents & Fabric Care" + info = resolve_hsn_gst(category, title) + assert (info.hsn_code, info.gst_percent, info.hsn_gst_needs_review) == ("3402", 18, False) + + +@pytest.mark.parametrize("title", ["Vim Dishwash Bar", "Lizol Floor Cleaner", "Dettol Handwash"]) +def test_other_cleaners_stay_household_cleaning(title: str) -> None: + assert detect_category_from_text(title) == "Household Cleaning" + + +def test_fab_is_not_a_laundry_keyword() -> None: + """Parle Fab is a biscuit. A bare "Godrej Fab" gets its category from the + label text, not from the product-line name.""" + assert detect_category_from_text("Parle Fab Biscuits") == "Biscuits & Cookies" + assert detect_category_from_text("Godrej Fab") is None + + +def test_face_wash_is_not_laundry() -> None: + assert detect_category_from_text("Nivea Face Wash") == "Skin & Bath Care" + + +def test_category_names_are_unique() -> None: + assert len(ALL_CATEGORIES) == len(set(ALL_CATEGORIES)) diff --git a/tests/test_identify_api.py b/tests/test_identify_api.py index 445c98d..640221b 100644 --- a/tests/test_identify_api.py +++ b/tests/test_identify_api.py @@ -46,6 +46,9 @@ def _row(score: float = 0.631, image_id: str = "britannia_marie_gold_300g") -> D IMAGE_SEARCH_KEYS = {"results", "total", "detected_brand", "scoped_to_brand", "scope_fallback", "min_score", "top_k", "query_text"} IDENTIFY_KEYS = IMAGE_SEARCH_KEYS | {"matched_by", "ocr_text", "ocr_source", "image_top_score", "fallback_reason"} +# Capture-to-catalog fields: always present, None unless discovery ran (see +# tests/test_capture_discovery.py). Additive - no existing key changed. +IDENTIFY_KEYS |= {"discovery_status", "discovery_job_id", "discovery_message", "provisional"} def _post_identify(client, data: bytes, **form):