Image capturing Flow updates
This commit is contained in:
@@ -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.
|
||||
|
||||
|
||||
748
app/services/capture_discovery.py
Normal file
748
app/services/capture_discovery.py
Normal file
@@ -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"(?<!\w)" + re.escape(name) + r"(?!\w)", lower)
|
||||
if match:
|
||||
return text[match.start():match.end()], len(lower[:match.start()].split())
|
||||
from app.services.query_intent import extract_brand_mention
|
||||
|
||||
mention = extract_brand_mention(text)
|
||||
if mention:
|
||||
match = re.search(r"(?<!\w)" + re.escape(mention.lower()) + r"(?!\w)", lower)
|
||||
position = len(lower[:match.start()].split()) if match else 0
|
||||
return mention, position
|
||||
# A sub-brand printed on its own: "Surf excel" is only known through the
|
||||
# alias "hul surf excel". Two-word runs only - a single word would let
|
||||
# "Fab" reach Parle through "parle fab", which is exactly the misroute the
|
||||
# registry's own leading-brand rule refuses.
|
||||
from app.services.brand_registry import BRAND_ALIASES, resolve_parent_brand
|
||||
|
||||
parents = set(BRAND_ALIASES.values())
|
||||
words = text.split()
|
||||
for i in range(len(words) - 1):
|
||||
pair = f"{words[i]} {words[i + 1]}"
|
||||
if not all(w.isalpha() for w in pair.split()):
|
||||
continue
|
||||
if resolve_parent_brand(pair) in parents:
|
||||
return pair, i
|
||||
return None
|
||||
|
||||
|
||||
def _name_window(words: List[str], start: int) -> 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 "<job id>.<ext>"."""
|
||||
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()
|
||||
@@ -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"},
|
||||
|
||||
Reference in New Issue
Block a user