Brand Discovery-LLM Updates

This commit is contained in:
sriram
2026-09-29 16:42:49 +05:30
parent c0489d89d6
commit c0601b65fe
16 changed files with 1971 additions and 5 deletions

View File

@@ -294,6 +294,18 @@ BRAND_SYNC_INTERVAL_SECONDS=300
# Food Facts runs first and is never subject to it.
#BRAND_DISCOVERY_DEADLINE_SECONDS=300
# Brand Discovery's optional "Web & retail listings" source, for brands Open
# Food Facts does not carry. Reads search-engine results that point at retailer
# product pages (Blinkit, Zepto, Instamart, BigBasket, JioMart, Amazon,
# Flipkart...) - never the pages themselves, never a language model. Uses the
# Google Custom Search API when USE_GOOGLE_CSE=true, else DuckDuckGo.
#WEB_DISCOVERY_ENABLED=true
#WEB_DISCOVERY_MAX_QUERIES=80
#WEB_DISCOVERY_PAUSE_SECONDS=2.0
#WEB_DISCOVERY_RESULTS_PER_QUERY=20
#WEB_DISCOVERY_CACHE_DAYS=7
#WEB_DISCOVERY_DIR=/app/data/cache/web_discovery
USE_S3=true
S3_ACCESS_KEY=your-do-spaces-key

View File

@@ -33,6 +33,7 @@ from app.infrastructure.settings import (
BATCH_MAX_TOTAL_ROWS,
BRAND_DISCOVERY_DEADLINE_SECONDS,
BRAND_DISCOVERY_MAX_PRODUCTS,
WEB_DISCOVERY_ENABLED,
)
from app.services import active_brands, brand_discovery
@@ -72,6 +73,11 @@ class DiscoveryPreviewRequest(BaseModel):
deadline_seconds: float = Field(default=90.0, ge=0.0,
le=BRAND_DISCOVERY_DEADLINE_SECONDS)
refresh_corpus: bool = False
# Web & retail listings. Reads a FINISHED web-discovery job (start one at
# POST /web-jobs first); off by default, and off means the preview is
# exactly what it was before this source existed.
use_web: bool = False
web_job_id: Optional[str] = None
class DiscoveredProductIn(BaseModel):
@@ -93,6 +99,19 @@ class DiscoveredProductIn(BaseModel):
fssai_license: Optional[str] = None
barcode: Optional[str] = None
image_url: Optional[str] = None
# Web-found rows only: the listings that justified the product. Not
# written by the pipeline; recorded afterwards into field_sources by
# web_discovery.provenance, which also fills a blank price_range from them.
listings: List[Dict[str, Any]] = Field(default_factory=list)
price_range: Optional[str] = None
retailer_count: int = 0
class WebJobRequest(BaseModel):
brand: str
# Re-search even when a finished job for this brand is less than a day old.
# Answers already cached are still reused, so this is cheap.
refresh: bool = False
class DiscoveryIngestRequest(BaseModel):
@@ -131,6 +150,8 @@ async def preview_brand_discovery(payload: DiscoveryPreviewRequest) -> Dict[str,
use_llm=payload.use_llm,
require_evidence=payload.require_evidence,
refresh_corpus=payload.refresh_corpus,
use_web=payload.use_web and WEB_DISCOVERY_ENABLED,
web_job_id=payload.web_job_id,
)
except ValueError as exc:
raise HTTPException(status_code=400, detail=str(exc)) from exc
@@ -220,9 +241,51 @@ async def ingest_brand_discovery(payload: DiscoveryIngestRequest) -> batch_commo
"press Resume on it once the current batch finishes."
),
)
# Web-found rows: record their retailer listings once the batch has stored
# them. Keyed by CSV position, which is how stage 11 reports source rows.
web_entries = {
index: {"listings": item.listings, "price_range": item.price_range or "",
"retailer_count": item.retailer_count}
for index, item in enumerate(payload.products) if item.listings
}
if web_entries:
from app.services.web_discovery import provenance
provenance.watch(manifest.batch_id, brand, web_entries)
return batch_common.to_out(manifest)
# ---------------------------------------------------------------------------
# Web & retail listings - a background search the panel polls
# ---------------------------------------------------------------------------
def _require_web() -> None:
if not WEB_DISCOVERY_ENABLED:
raise HTTPException(status_code=404, detail="Web discovery is switched off (WEB_DISCOVERY_ENABLED).")
@router.post("/web-jobs", status_code=status.HTTP_202_ACCEPTED,
dependencies=[Depends(require_admin)])
def start_web_job(payload: WebJobRequest) -> Dict[str, Any]:
"""Start searching retailer listings for a brand, or return the job that is
already running or finished within the last day. Poll GET /web-jobs/{id}."""
_require_web()
from app.services.web_discovery import jobs as web_jobs
try:
job = web_jobs.start_job(payload.brand, refresh=payload.refresh)
except ValueError as exc:
raise HTTPException(status_code=400, detail=str(exc)) from exc
return job.summary()
@router.get("/web-jobs/{job_id}", dependencies=[Depends(require_admin)])
def get_web_job(job_id: str) -> Dict[str, Any]:
_require_web()
from app.services.web_discovery import jobs as web_jobs
job = web_jobs.get_job(job_id)
if job is None:
raise HTTPException(status_code=404, detail="No web discovery job with that id.")
return job.summary()
# ---------------------------------------------------------------------------
# The ACTIVE_BRANDS message, in one place
# ---------------------------------------------------------------------------

View File

@@ -347,6 +347,25 @@ BRAND_DISCOVERY_DEADLINE_SECONDS = float(
os.getenv("BRAND_DISCOVERY_DEADLINE_SECONDS", "300")
)
# Web & retail listing discovery (app/services/web_discovery/) - the optional
# Brand Discovery source for brands Open Food Facts does not carry. It reads
# ONLY search-engine results (title, URL, snippet) that point at Indian
# retailers' product pages; it never fetches a retailer page and never asks a
# language model. A product is kept only when a real listing names the brand
# and states a pack size. Off in a request unless the admin ticks it.
WEB_DISCOVERY_ENABLED = _bool("WEB_DISCOVERY_ENABLED", "true")
# Searches per brand job. A parent brand fans out to its sub-brands (Reckitt ->
# Dettol, Harpic, Lizol, ...) times the retailers, so this is the cost ceiling.
WEB_DISCOVERY_MAX_QUERIES = int(os.getenv("WEB_DISCOVERY_MAX_QUERIES", "80"))
# Seconds between two LIVE searches (cached ones are free). DuckDuckGo throttles
# a burst; a throttled search is "could not ask", never "nothing found".
WEB_DISCOVERY_PAUSE_SECONDS = float(os.getenv("WEB_DISCOVERY_PAUSE_SECONDS", "2.0"))
WEB_DISCOVERY_RESULTS_PER_QUERY = int(os.getenv("WEB_DISCOVERY_RESULTS_PER_QUERY", "20"))
# How long a search answer is reused. An EMPTY answer is kept one day only:
# search reach is unstable, and a flaky empty page must not become a week-long fact.
WEB_DISCOVERY_CACHE_DAYS = float(os.getenv("WEB_DISCOVERY_CACHE_DAYS", "7"))
WEB_DISCOVERY_DIR = _dir("WEB_DISCOVERY_DIR", DATA_DIR / "cache" / "web_discovery")
# ---------------------------------------------------------------------------
# S3 / DigitalOcean Spaces (product image storage) - optional
# ---------------------------------------------------------------------------

View File

@@ -227,8 +227,23 @@ class DiscoveredProduct:
confidence: float = 0.0
matches_existing: Optional[str] = None
notes: List[str] = field(default_factory=list)
# Web & retail listings (app/services/web_discovery). Empty unless that
# source found this product, and then shown in the preview and handed back
# on ingest for the provenance post-pass - never written to the CSV.
listings: List[Dict[str, Any]] = field(default_factory=list)
price_range: str = ""
retailer_count: int = 0
def as_preview(self) -> Dict[str, Any]:
body = self._preview_body()
# Only for rows the web source found, so a preview without it is
# exactly what it was before that source existed.
if self.listings:
body.update({"listings": list(self.listings), "price_range": self.price_range,
"retailer_count": self.retailer_count})
return body
def _preview_body(self) -> Dict[str, Any]:
return {
"brand": self.brand,
"product_name": self.product_name,
@@ -918,6 +933,8 @@ def discover_brand_products(
use_store: bool = True,
require_evidence: bool = True,
refresh_corpus: bool = False,
use_web: bool = False,
web_job_id: Optional[str] = None,
) -> DiscoveryResult:
"""Discover a brand's products. Reads the catalog; writes nothing.
@@ -987,7 +1004,15 @@ def discover_brand_products(
"brand_registry.BRAND_STORE_DOMAINS and run "
"scripts/backfill_brand_stores.py to populate a real catalogue."
)
if not off_candidates and not store_candidates:
# Web & retail listings: what a finished web-discovery job found. Read,
# never run, here - the job takes minutes and the panel starts it first.
web_candidates: List[Dict[str, Any]] = []
if use_web:
web_candidates = _from_web(brand, web_job_id, warnings)
# With the web source on, this warning is only true when the model is on
# too; without it, it is left exactly as it has always been.
if (not off_candidates and not store_candidates and not web_candidates
and (use_llm or not use_web)):
warnings.append(
"Every row below rests on the language model alone - review them "
"individually."
@@ -1018,12 +1043,17 @@ def discover_brand_products(
# STORE FIRST, then Open Food Facts, then the model. The merge fills a
# blank from whatever comes later, so order IS precedence: the brand's own
# shop outranks a crowd-sourced database, which outranks a guess.
for candidate in store_candidates + off_candidates + llm_candidates:
for candidate in store_candidates + off_candidates + web_candidates + llm_candidates:
title = candidate["title"]
key = _normalise_title(brand, title)
target = None
if key:
for index, existing in enumerate(keys):
# Two web rows were already told apart by web_discovery's
# stricter rule (a variant word on one side only - "Citrus"
# vs "Floral" scores 0.857 here, over the 0.85 floor).
if candidate.get("source") == "web" and merged[index].get("sources") == ["web"]:
continue
if not _same_product(key, existing, title, merged[index]["title"]):
continue
held = merged[index].get("size")
@@ -1050,6 +1080,16 @@ def discover_brand_products(
entry["sizes"] = candidate["sizes"]
if not entry.get("providers") and candidate.get("providers"):
entry["providers"] = candidate["providers"]
if candidate.get("listings"):
# A web listing joining an OFF / store row: keep the evidence and
# the retailers, the other source keeps the name and the barcode.
entry["listings"] = (entry.get("listings") or []) + list(candidate["listings"])
entry["retailer_count"] = max(int(entry.get("retailer_count") or 0),
int(candidate.get("retailer_count") or 0))
entry["providers"] = list(dict.fromkeys(
list(candidate.get("providers") or []) + list(entry.get("providers") or [])))
if not entry.get("price_range") and candidate.get("price_range"):
entry["price_range"] = candidate["price_range"]
# EVERY PACK OF ONE PRODUCT MUST BE NAMED THE SAME WAY.
#
@@ -1154,6 +1194,8 @@ def discover_brand_products(
)
if product is None:
continue
if "web" in product.sources:
_apply_web(product, candidate)
if (require_evidence and "off" not in product.sources
and "store" not in product.sources and not product.evidence):
dropped += 1
@@ -1173,6 +1215,8 @@ def discover_brand_products(
"with_barcode": sum(1 for p in products if p.barcode),
"dropped_without_evidence": dropped,
}
if use_web:
counts["from_web"] = sum(1 for p in products if "web" in p.sources)
return DiscoveryResult(
brand=brand,
@@ -1186,6 +1230,50 @@ def discover_brand_products(
)
# ---------------------------------------------------------------------------
# Web & retail listings - the optional fourth source
# ---------------------------------------------------------------------------
# A product a real Indian retailer is listing now, found from search results
# only (app/services/web_discovery). Two retailers or more is strong evidence;
# one is real but thin, so it is shown and left unticked.
WEB_CONFIDENCE_MULTI = 0.9
WEB_CONFIDENCE_SINGLE = 0.75
def _from_web(brand: str, job_id: Optional[str], warnings: List[str]) -> List[Dict[str, Any]]:
from app.services.web_discovery import jobs as web_jobs
job = web_jobs.candidates_for(brand, job_id)
if job is None:
warnings.append(
"Web & retail listings have not been searched for this brand yet - "
"tick the box and run Discover again to start the search."
)
return []
kept = len(job.candidates)
note = (f"Web & retail listings: {kept} product(s) from {job.listings_kept} retailer "
f"listing(s), {job.queries_done} of {job.queries_total} searches via {job.backend}.")
if job.status != web_jobs.DONE:
note += f" Incomplete: {job.detail}"
warnings.append(note)
return [dict(c) for c in job.candidates]
def _apply_web(product: DiscoveredProduct, candidate: Dict[str, Any]) -> None:
product.listings = list(candidate.get("listings") or [])
product.price_range = candidate.get("price_range") or ""
product.retailer_count = int(candidate.get("retailer_count") or 0)
# A live retailer listing outranks our own earlier output ("catalog") and a
# hardcoded list ("registry"); only Open Food Facts stays ahead of it.
if product.evidence in (None, "catalog", "registry"):
product.evidence = "retail"
if product.sources == ["web"]:
product.confidence = (WEB_CONFIDENCE_MULTI if product.retailer_count >= 2
else WEB_CONFIDENCE_SINGLE)
retailers = ", ".join(product.providers[:product.retailer_count] or product.providers)
product.notes.append(f"listed by {retailers}")
# ---------------------------------------------------------------------------
# Serialisation - the bridge itself
# ---------------------------------------------------------------------------

View File

@@ -0,0 +1,37 @@
"""Find a brand's real products from web search results that point at Indian retailers.
WHY THIS EXISTS
---------------
Brand Discovery's sources are Open Food Facts, the brand's own storefront and a
language model. For a non-food brand such as Reckitt Benckiser the first two are
empty - OFF is a FOOD database and no storefront is registered - which left the
model inventing products no shop has ever sold. This package is a fourth
source that only reports what a real retailer is listing.
WHAT IT READS, AND WHAT IT NEVER DOES
-------------------------------------
It reads the search ENGINE's answer - result title, URL and snippet - for
queries like `"Dettol site:blinkit.com"`. It never fetches a Blinkit, Zepto,
Instamart, Amazon or Flipkart page, never calls their private APIs, and never
asks a language model anything. A product is kept only when
1. the result URL is a PRODUCT page on an allowlisted retailer (not a search,
category or brand page) - listings.is_product_url;
2. the listing title names the brand or one of its sub-brands, as a word; and
3. the title states an explicit pack size. No size, no product: the size is
what separates a pack that exists from an invented one (see
retail_presence's "Anil Wheat Vermicelli 12g" measurement).
MODULES
-------
search.py DuckDuckGo, or the Google Custom Search API when configured
listings.py retailer + product-URL detection, title cleaning, size, price
brands.py a parent brand -> the sub-brands it trades under
cache.py SQLite cache of search answers (a throttle is never cached)
discover.py queries -> listings -> one candidate per (product, pack size)
jobs.py one background worker; the admin panel polls its progress
provenance.py after an ingest batch: listing URLs into field_sources
Nothing here is imported by the 11-stage pipeline, and nothing here runs
unless an admin ticks "Web & retail listings" in the Brand Discovery tab.
"""

View File

@@ -0,0 +1,118 @@
"""A brand name -> the names its products are actually sold under.
Nobody shops for "Reckitt Benckiser". Its products are listed as Dettol,
Harpic, Lizol, Mortein, Durex..., so a search for the parent finds almost
nothing. `brand_registry.BRAND_ALIASES` already knows the family, but spells
it with an abbreviation - "rb dettol" -> "reckitt benckiser" - and
`get_known_sub_brands` only accepts aliases that START with the parent name,
so for Reckitt it returns []. This reads the same registry without changing it.
WHICH WORDS COUNT AS "THE BRAND"
--------------------------------
A listing is kept only if it names the brand (listings.names_brand). For an
abbreviation alias the sub-brand itself is the brand word: "rb dettol" ->
"dettol". For an alias that starts with the parent ("nestle cerelac",
"amul butter") the parent always counts, and the rest counts too ONLY when it
is a product-line name rather than a kind of food: "Cerelac" is sold without
"Nestle" in the title, but "butter" is not a brand, and accepting it would
let another company's butter in under Amul. A rest is a line name when it is
not a category keyword, is at least 5 characters, and belongs to exactly one
parent in the registry.
"""
from __future__ import annotations
import re
from collections import defaultdict
from dataclasses import dataclass
from typing import Dict, List, Set, Tuple
from app.services.brand_registry import BRAND_ALIASES, resolve_parent_brand
@dataclass(frozen=True)
class Target:
"""One name to search for, and the words a listing must lead with to count."""
query: str # what goes before `site:` - "Dettol", "Amul Butter"
brand_terms: Tuple[str, ...] # "dettol"; or ("nestle", "cerelac")
def _words(text: str) -> List[str]:
return re.findall(r"[a-z0-9]+", (text or "").lower())
def _is_abbreviation(token: str, parent: str) -> bool:
""""rb" for Reckitt Benckiser, "hul" for Hindustan Unilever, "jnj" for
Johnson & Johnson. Short, not itself a word of the parent, and starting
with the parent's first letter - which is what keeps "kit kat" (Nestle)
from being read as the abbreviation "kit" plus a sub-brand "kat"."""
parent_words = _words(parent)
return (bool(parent_words) and 1 < len(token) <= 4 and token not in parent_words
and token[0] == parent_words[0][0])
def _sub_brand_part(alias: str, parent: str) -> str:
words = alias.split()
if len(words) > 1 and _is_abbreviation(words[0], parent):
return " ".join(words[1:])
return alias
def _rest_owners() -> Dict[str, Set[str]]:
"""For every "<parent> <rest>" alias, which parents use that rest."""
owners: Dict[str, Set[str]] = defaultdict(set)
for alias, owner in BRAND_ALIASES.items():
parent = owner.lower().strip()
if alias.startswith(parent + " "):
owners[alias[len(parent) + 1:].strip()].add(parent)
return owners
def _is_line_name(rest: str, owners: Dict[str, Set[str]]) -> bool:
from app.services.category_registry import detect_category_from_text
if len(rest.replace(" ", "")) < 5 or len(owners.get(rest, ())) != 1:
return False
return detect_category_from_text(rest, exact_only=True) is None
def _family(parent: str) -> List[Target]:
parent_key = parent.lower().strip()
owners = _rest_owners()
out: List[Target] = []
seen: Set[str] = set()
for alias, owner in sorted(BRAND_ALIASES.items()):
if owner.lower().strip() != parent_key or alias == parent_key:
continue
if alias.startswith(parent_key + " "):
rest = alias[len(parent_key) + 1:].strip()
terms = (parent_key, rest) if _is_line_name(rest, owners) else (parent_key,)
target = Target(query=alias.title(), brand_terms=terms)
else:
rest = _sub_brand_part(alias, parent)
target = Target(query=rest.title(), brand_terms=(rest,))
if target.query.lower() not in seen:
seen.add(target.query.lower())
out.append(target)
return out
def targets_for(brand: str) -> List[Target]:
"""Sub-brands first, the brand itself last. Deterministic order.
A sub-brand typed on its own ("Dettol") searches only that sub-brand."""
name = (brand or "").strip()
if not name:
return []
parent = resolve_parent_brand(name) or name
parent_key = parent.lower().strip()
typed_key = " ".join(_words(name))
family = _family(parent)
if typed_key != " ".join(_words(parent_key)):
own = [t for t in family
if " ".join(_words(t.query)) == typed_key
or typed_key in {" ".join(_words(term)) for term in t.brand_terms[1:] or t.brand_terms}]
return own or [Target(query=name, brand_terms=(name.lower(),))]
display = parent.title() if parent.islower() else parent
return family + [Target(query=display, brand_terms=(parent_key,))]

View File

@@ -0,0 +1,82 @@
"""Search answers, kept so a re-run costs nothing and a throttled run can finish.
Only ANSWERS are stored. A search that could not be asked (search.search
returned None) is never written: caching it would turn "we could not ask" into
"there is nothing", which is the one confusion this whole source must avoid.
An empty answer is kept for a day only - the same query minutes apart has been
measured returning 7 shop results and then 0 (see retail_presence).
"""
from __future__ import annotations
import json
import logging
import sqlite3
import threading
import time
from contextlib import closing
from pathlib import Path
from typing import List, Optional
from app.infrastructure import settings
from app.services.web_discovery.search import Hit
logger = logging.getLogger(__name__)
_lock = threading.Lock()
_EMPTY_TTL_SECONDS = 24 * 3600
def _db_path() -> Path:
return Path(settings.WEB_DISCOVERY_DIR) / "search_cache.db"
def _connect() -> sqlite3.Connection:
path = _db_path()
path.parent.mkdir(parents=True, exist_ok=True)
conn = sqlite3.connect(str(path), timeout=10)
conn.execute(
"CREATE TABLE IF NOT EXISTS search_cache ("
" query TEXT NOT NULL, backend TEXT NOT NULL, hits TEXT NOT NULL,"
" created_at REAL NOT NULL, PRIMARY KEY (query, backend))"
)
return conn
def _key(query: str) -> str:
return " ".join((query or "").lower().split())
def get(query: str, backend: str) -> Optional[List[Hit]]:
"""The cached answer, or None for a miss or an expired entry. Never raises."""
try:
with _lock, closing(_connect()) as conn:
row = conn.execute(
"SELECT hits, created_at FROM search_cache WHERE query = ? AND backend = ?",
(_key(query), backend),
).fetchone()
except Exception as exc: # noqa: BLE001 - a cache failure is a cache miss
logger.debug("web discovery cache read failed: %s", exc)
return None
if not row:
return None
hits = [Hit.from_dict(h) for h in json.loads(row[0])]
ttl = settings.WEB_DISCOVERY_CACHE_DAYS * 86400 if hits else _EMPTY_TTL_SECONDS
if time.time() - row[1] > ttl:
return None
return hits
def put(query: str, backend: str, hits: List[Hit]) -> None:
"""Store an ANSWER. Callers must never pass the None of a failed search."""
if hits is None: # defensive: see the module docstring
return
try:
with _lock, closing(_connect()) as conn:
conn.execute(
"INSERT OR REPLACE INTO search_cache (query, backend, hits, created_at) "
"VALUES (?, ?, ?, ?)",
(_key(query), backend, json.dumps([h.as_dict() for h in hits]), time.time()),
)
conn.commit()
except Exception as exc: # noqa: BLE001
logger.debug("web discovery cache write failed: %s", exc)

View File

@@ -0,0 +1,225 @@
"""Queries -> listings -> one candidate per (product, pack size).
For each sub-brand (brands.targets_for) and each retailer
(listings.SEARCH_ORDER) one query, `"<sub-brand> site:<retailer>"`, capped at
WEB_DISCOVERY_MAX_QUERIES. The order is retailer-major, so when the budget
runs out it is the last retailers that go unasked for every sub-brand, not
the last sub-brands that go unasked everywhere.
Each hit either becomes a listings.Listing or is refused with a reason; the
reasons are counted and reported, so a filter never throws products away
silently. Listings of the same pack on several retailers are folded into one
candidate, and the number of distinct retailers is its corroboration.
A search that could not be asked is counted as such. Three in a row end the
run as PARTIAL - a throttled provider only gets worse if pushed - and what was
learned so far is kept; the cached answers make the next run pick up there.
"""
from __future__ import annotations
import logging
import re
import time
from collections import Counter
from dataclasses import dataclass, field
from typing import Callable, Dict, List, Optional, Sequence
from app.infrastructure import settings
from app.services.web_discovery import cache, listings, search
from app.services.web_discovery.brands import Target, targets_for
logger = logging.getLogger(__name__)
DONE = "done"
PARTIAL = "partial"
_MAX_CONSECUTIVE_FAILURES = 3
_MAX_LISTINGS_PER_CANDIDATE = 10
# Two names are one product when their words (brand and packaging words set
# aside) overlap this much, Jaccard. Deliberately strict: "Dettol Soap" and
# "Dettol Original Soap" stay apart (0.5), because folding a generic title into
# a variant would name two products as one. Two rows the admin can see are
# better than one silently wrong row.
_SAME_NAME_JACCARD = 0.6
_PACKAGING_WORDS = frozenset({"bar", "pouch", "bottle", "jar", "box", "tin", "carton", "pack"})
@dataclass
class Query:
text: str
target: Target
retailer: str
@dataclass
class RunResult:
status: str = DONE
candidates: List[Dict] = field(default_factory=list)
queries_total: int = 0
queries_done: int = 0
queries_failed: int = 0
queries_cached: int = 0
listings_kept: int = 0
rejected: Dict[str, int] = field(default_factory=dict)
backend: str = ""
detail: str = ""
def plan_queries(brand: str, max_queries: Optional[int] = None) -> List[Query]:
budget = max_queries if max_queries is not None else settings.WEB_DISCOVERY_MAX_QUERIES
out: List[Query] = []
for retailer in listings.SEARCH_ORDER:
site = listings.RETAILERS[retailer][1]
for target in targets_for(brand):
out.append(Query(text=f"{target.query} site:{site}", target=target, retailer=retailer))
return out[:max(0, budget)]
def _identity(listing: listings.Listing) -> List[str]:
brand_words = set()
for term in (listing.brand_term,):
brand_words.update(listings._words(term))
return [t for t in listing.tokens if t not in brand_words and t not in _PACKAGING_WORDS]
def _variant_words(listing: listings.Listing) -> set:
"""Words the retailer put in brackets - how Blinkit and others mark a
variant: "Lizol Floor Cleaner (Citrus)", "(Floral)", "(Original)"."""
return {w for group in re.findall(r"\(([^()]*)\)", listing.name) for w in listings._words(group)}
def _same_product(a: listings.Listing, b: listings.Listing) -> bool:
"""Same brand word, same pack, near-identical name, and no bracketed
variant word on one side only. Measured: "(Citrus) 2l" and "(Floral) 2l"
share 4 of 6 words - enough for the overlap test alone, and two products."""
if a.brand_term != b.brand_term or a.size != b.size:
return False
ident_a, ident_b = _identity(a), _identity(b)
if _jaccard(ident_a, ident_b) < _SAME_NAME_JACCARD:
return False
differing = set(ident_a) ^ set(ident_b)
return not (differing & (_variant_words(a) | _variant_words(b)))
def _jaccard(a: Sequence[str], b: Sequence[str]) -> float:
sa, sb = set(a), set(b)
if not sa and not sb:
return 1.0
return len(sa & sb) / len(sa | sb)
def _price_range(prices: Sequence[float]) -> str:
if not prices:
return ""
low, high = min(prices), max(prices)
fmt = lambda v: f"₹{int(v)}" if float(v).is_integer() else f"₹{v:.2f}" # noqa: E731
return fmt(low) if low == high else f"{fmt(low)} - {fmt(high)}"
def cluster(found: Sequence[listings.Listing]) -> List[Dict]:
"""Fold listings of the same pack into candidates, best-corroborated first."""
groups: List[List[listings.Listing]] = []
seen_urls = set()
for listing in found:
if listing.url in seen_urls:
continue
seen_urls.add(listing.url)
for members in groups:
# Against the group's FIRST listing, so a chain of near-misses
# cannot drift a group from one product to another.
if _same_product(members[0], listing):
members.append(listing)
break
else:
groups.append([listing])
candidates: List[Dict] = []
for members in groups:
brand_term, size = members[0].brand_term, members[0].size
names = Counter(m.name for m in members)
# The spelling most retailers agree on; then the shortest; then A-Z.
name = sorted(names, key=lambda n: (-names[n], len(n), n))[0]
retailers = sorted({m.retailer for m in members}, key=listings.SEARCH_ORDER.index)
prices = [m.price for m in members if m.price]
# Brackets come off, their words stay: Brand Discovery's name
# normaliser (and off_bulk.strip_sizes) drops bracketed text, which
# turned "Lizol Floor Cleaner (Citrus)" and "(Floral)" into one name.
plain = re.sub(r"\s+", " ", re.sub(r"[()]", " ", name)).strip()
candidates.append({
"title": f"{plain} {size}",
"size": size,
"sizes": [size],
"providers": [listings.display_name(r) for r in retailers],
"retailer_count": len(retailers),
"price_range": _price_range(prices),
"listings": [m.as_dict() for m in members[:_MAX_LISTINGS_PER_CANDIDATE]],
"brand_term": brand_term,
"source": "web",
})
candidates.sort(key=lambda c: (-c["retailer_count"], c["title"].lower()))
return candidates
def run(brand: str, *, progress: Optional[Callable[[RunResult], None]] = None,
should_stop: Optional[Callable[[], bool]] = None,
max_queries: Optional[int] = None, sleep: Callable[[float], None] = time.sleep) -> RunResult:
"""Search, parse, fold. Never raises for a search failure."""
queries = plan_queries(brand, max_queries)
backend = search.backend_name()
result = RunResult(queries_total=len(queries), backend=backend)
rejected: Counter = Counter()
kept: List[listings.Listing] = []
consecutive_failures = 0
for query in queries:
if should_stop and should_stop():
result.status = PARTIAL
result.detail = "stopped"
break
hits = cache.get(query.text, backend)
if hits is not None:
result.queries_cached += 1
else:
hits = search.search(query.text)
if hits is None:
result.queries_failed += 1
consecutive_failures += 1
if consecutive_failures >= _MAX_CONSECUTIVE_FAILURES:
result.status = PARTIAL
result.detail = (f"the search provider stopped answering after "
f"{result.queries_done} of {len(queries)} searches; "
f"run it again later - answers so far are cached")
result.queries_done += 1
break
result.queries_done += 1
if progress:
progress(result)
sleep(settings.WEB_DISCOVERY_PAUSE_SECONDS)
continue
cache.put(query.text, backend, hits)
sleep(settings.WEB_DISCOVERY_PAUSE_SECONDS)
consecutive_failures = 0
for hit in hits:
listing, reason = listings.parse_hit(hit.title, hit.url, hit.snippet,
query.target.brand_terms)
if listing is None:
rejected[reason] += 1
else:
kept.append(listing)
result.queries_done += 1
if progress:
progress(result)
result.listings_kept = len(kept)
result.rejected = dict(rejected)
result.candidates = cluster(kept)
if result.status == DONE and result.queries_failed:
# Some searches went unasked, so absence of a product proves nothing.
result.status = PARTIAL
result.detail = (f"{result.queries_failed} search(es) could not be asked; run it again "
f"to fill them in - answered searches are cached")
logger.info("web discovery %s: %d candidates from %d listings (%d/%d searches, %d cached, %d failed)",
brand, len(result.candidates), len(kept), result.queries_done, len(queries),
result.queries_cached, result.queries_failed)
return result

View File

@@ -0,0 +1,225 @@
"""Web discovery as a background job the Brand Discovery panel polls.
A brand fans out to ~80 searches paced 2 s apart, which is minutes - far past
what a preview request may block for. So "Discover products" with the web box
ticked starts (or reuses) a job, the panel polls its progress, and the preview
then reads the finished job's candidates.
One worker thread: the search provider throttles bursts, and two brands
searched in parallel would only get both throttled. Job records are JSON under
WEB_DISCOVERY_DIR/jobs so a restart does not lose a finished job; a job that
was running when the process died is marked `interrupted`.
Reuse: a finished job for the same brand younger than REUSE_SECONDS is
returned as is unless the caller asks to refresh. Its searches are cached for
days anyway, so even a refresh is cheap for what was already answered.
"""
from __future__ import annotations
import json
import logging
import queue
import threading
import time
import uuid
from dataclasses import asdict, dataclass, field
from pathlib import Path
from typing import Dict, List, Optional
from app.infrastructure import settings
from app.services.web_discovery import discover
logger = logging.getLogger(__name__)
QUEUED = "queued"
RUNNING = "running"
DONE = discover.DONE
PARTIAL = discover.PARTIAL
FAILED = "failed"
INTERRUPTED = "interrupted"
FINISHED = {DONE, PARTIAL, FAILED, INTERRUPTED}
REUSE_SECONDS = 24 * 3600
@dataclass
class WebJob:
job_id: str
brand: str
status: str = QUEUED
created_at: float = field(default_factory=time.time)
updated_at: float = field(default_factory=time.time)
backend: str = ""
queries_total: int = 0
queries_done: int = 0
queries_failed: int = 0
queries_cached: int = 0
listings_kept: int = 0
rejected: Dict[str, int] = field(default_factory=dict)
candidates: List[Dict] = field(default_factory=list)
detail: str = ""
def summary(self) -> Dict:
"""Everything but the candidates - what the progress poll needs."""
body = asdict(self)
body.pop("candidates")
body["candidate_count"] = len(self.candidates)
return body
_jobs: Dict[str, WebJob] = {}
_lock = threading.Lock()
_queue: "queue.Queue[str]" = queue.Queue()
_worker: Optional[threading.Thread] = None
def _brand_key(brand: str) -> str:
return " ".join((brand or "").lower().split())
def _jobs_dir() -> Path:
path = Path(settings.WEB_DISCOVERY_DIR) / "jobs"
path.mkdir(parents=True, exist_ok=True)
return path
def _save(job: WebJob) -> None:
job.updated_at = time.time()
try:
(_jobs_dir() / f"{job.job_id}.json").write_text(json.dumps(asdict(job)), encoding="utf-8")
except Exception as exc: # noqa: BLE001 - the in-memory job is still served
logger.warning("web discovery: could not save job %s: %s", job.job_id, exc)
def _load(job_id: str) -> Optional[WebJob]:
if not job_id or not all(c in "0123456789abcdef" for c in job_id):
return None
path = _jobs_dir() / f"{job_id}.json"
if not path.exists():
return None
try:
job = WebJob(**json.loads(path.read_text(encoding="utf-8")))
except Exception: # noqa: BLE001 - an unreadable record is no record
return None
if job.status in (QUEUED, RUNNING):
# Loaded from disk, so no worker of THIS process is running it.
job.status = INTERRUPTED
job.detail = "the server restarted while this was running; start it again"
return job
def get_job(job_id: str) -> Optional[WebJob]:
with _lock:
job = _jobs.get(job_id)
if job is None:
job = _load(job_id)
if job is not None:
with _lock:
_jobs.setdefault(job.job_id, job)
return job
def latest_for(brand: str) -> Optional[WebJob]:
"""The newest FINISHED job with results for `brand`, in memory or on disk."""
key = _brand_key(brand)
best: Optional[WebJob] = None
with _lock:
known = list(_jobs.values())
try:
for path in _jobs_dir().glob("*.json"):
if path.stem not in {j.job_id for j in known}:
loaded = _load(path.stem)
if loaded:
known.append(loaded)
except Exception: # noqa: BLE001
pass
for job in known:
if _brand_key(job.brand) != key or job.status not in (DONE, PARTIAL):
continue
if best is None or job.updated_at > best.updated_at:
best = job
return best
def start_job(brand: str, *, refresh: bool = False) -> WebJob:
"""Queue a job for `brand`, or return the one already running or recent."""
brand = (brand or "").strip()
if not brand:
raise ValueError("a brand name is required")
key = _brand_key(brand)
with _lock:
for job in _jobs.values():
if _brand_key(job.brand) == key and job.status in (QUEUED, RUNNING):
return job
if not refresh:
recent = latest_for(brand)
if recent and time.time() - recent.updated_at < REUSE_SECONDS:
return recent
job = WebJob(job_id=uuid.uuid4().hex, brand=brand)
with _lock:
_jobs[job.job_id] = job
_save(job)
_ensure_worker()
_queue.put(job.job_id)
return job
def candidates_for(brand: str, job_id: Optional[str] = None) -> Optional[WebJob]:
"""The job whose candidates a preview should use: the named one if it is
for this brand and finished, else the brand's latest finished job."""
if job_id:
job = get_job(job_id)
if job and _brand_key(job.brand) == _brand_key(brand) and job.status in (DONE, PARTIAL):
return job
return latest_for(brand)
def _ensure_worker() -> None:
global _worker
with _lock:
if _worker is not None and _worker.is_alive():
return
_worker = threading.Thread(target=_work, name="web-discovery", daemon=True)
_worker.start()
def _work() -> None:
while True:
job_id = _queue.get()
job = get_job(job_id)
if job is None:
continue
run_job(job)
def run_job(job: WebJob) -> WebJob:
"""Run one job to the end on the calling thread. The worker's body; also
what tests call directly."""
job.status = RUNNING
_save(job)
def progress(partial: discover.RunResult) -> None:
job.backend = partial.backend
job.queries_total = partial.queries_total
job.queries_done = partial.queries_done
job.queries_failed = partial.queries_failed
job.queries_cached = partial.queries_cached
job.updated_at = time.time()
try:
result = discover.run(job.brand, progress=progress)
except Exception as exc: # noqa: BLE001 - a failed job must say so, not vanish
logger.exception("web discovery failed for %r", job.brand)
job.status = FAILED
job.detail = str(exc)
_save(job)
return job
progress(result)
job.status = result.status
job.detail = result.detail
job.listings_kept = result.listings_kept
job.rejected = result.rejected
job.candidates = result.candidates
_save(job)
return job

View File

@@ -0,0 +1,292 @@
"""A search hit -> a retail listing we can trust, or a reason it was refused.
Every rule here is a refusal. A listing is accepted only when it is a PRODUCT
page on a known retailer, names the brand as a word, states exactly one pack
size, and is not a multipack or a combo. Nothing is fetched: the title, URL and
snippet are what the search engine returned.
"""
from __future__ import annotations
import re
from dataclasses import dataclass, field
from typing import Dict, List, Optional, Sequence, Tuple
from urllib.parse import urlparse
# ---------------------------------------------------------------------------
# Retailers
# ---------------------------------------------------------------------------
# domain -> (display name as the `providers` column spells it, `site:` scope
# for the query, the path shape of ONE product's page). A search, category,
# brand or offer page is not a listing of a product and must not count as one;
# these patterns are what tell them apart. Taken from the URL grammar each site
# uses today (sku_service._MARKETPLACE_PATTERNS covers several of the same).
RETAILERS: Dict[str, Tuple[str, str, re.Pattern]] = {
"blinkit.com": ("Blinkit", "blinkit.com", re.compile(r"/prn/[^/]+/prid/\d+", re.I)),
"zeptonow.com": ("Zepto", "zeptonow.com", re.compile(r"/pn/[^/]+/pvid/[0-9a-f-]{8,}", re.I)),
"swiggy.com": ("Swiggy Instamart", "swiggy.com/instamart",
re.compile(r"/instamart/item/[A-Za-z0-9]+", re.I)),
"bigbasket.com": ("BigBasket", "bigbasket.com", re.compile(r"/pd/\d+", re.I)),
"jiomart.com": ("JioMart", "jiomart.com", re.compile(r"/p/[a-z-]+/[^/]+/\d+", re.I)),
"amazon.in": ("Amazon", "amazon.in", re.compile(r"/(?:dp|gp/product)/[A-Z0-9]{10}", re.I)),
"flipkart.com": ("Flipkart", "flipkart.com", re.compile(r"/p/itm[a-z0-9]+", re.I)),
}
# The retailers queried with `site:`, in the order they are asked. Quick
# commerce first: it lists what is in Indian shops this week.
SEARCH_ORDER: Tuple[str, ...] = (
"blinkit.com", "zeptonow.com", "swiggy.com", "bigbasket.com",
"jiomart.com", "amazon.in", "flipkart.com",
)
def retailer_for(url: str) -> Optional[str]:
"""The allowlisted retailer domain `url` is on, or None."""
host = urlparse(url or "").netloc.lower().split(":")[0]
for domain in RETAILERS:
if host == domain or host.endswith("." + domain):
return domain
return None
def is_product_url(retailer: str, url: str) -> bool:
spec = RETAILERS.get(retailer)
if not spec:
return False
parsed = urlparse(url or "")
return bool(spec[2].search(parsed.path + ("?" + parsed.query if parsed.query else "")))
def display_name(retailer: str) -> str:
return RETAILERS.get(retailer, (retailer,))[0]
# ---------------------------------------------------------------------------
# Titles
# ---------------------------------------------------------------------------
# What surrounds the product name in a result title. Each retailer wraps it
# differently:
# blinkit "Dettol Original Soap (125 g) - Buy Online at Best Price | Blinkit"
# zepto "Buy Dettol Original Soap 125 g Online at Best Price | Zepto"
# bigbasket "Buy Dettol Original Soap 125 g Online at Best Price of Rs 58 - bigbasket"
# amazon "Dettol Original Soap Bar, 125g : Amazon.in: Beauty"
# instamart "Buy Dettol Original Soap 125 g Online | Swiggy Instamart"
# blinkit "Harpic Toilet Cleaner - (Original) - 500 ml Price - Buy Online at ₹120 in India"
# so the name can run over several segments: they are joined from the first
# one that leads with the brand until the first piece of shop furniture.
_SEGMENT_SPLIT = re.compile(r"\s+[|:]\s+|\s+[-–—]\s+|\s*\|\s*")
_LEADING_BUY = re.compile(r"^\s*(?:buy|shop|order)\s+", re.I)
# Where the product name ends and the shop's wording begins. Everything from
# the first of these on is furniture: "... Online at Best Price", "... Price -
# Buy", "...: Amazon.in: Beauty", "... | Blinkit".
_FURNITURE = re.compile(
r"\s*\b(?:online|price|blinkit|zepto|zeptonow|swiggy|instamart|bigbasket|jiomart|"
r"amazon(?:\.in)?|flipkart(?:\.com)?|grofers|free\s+delivery|delivered\s+in|"
r"in\s+\d+\s+minutes?)\b.*$",
re.I,
)
# A listing of several packs, or of a bundle, is a different SKU from the pack
# itself; taken as the pack it would put a carton's size or a combo's name on
# a catalogue row.
_MULTIPACK = re.compile(
r"\b(?:pack\s+of\s+(?:[2-9]|\d{2,})|set\s+of\s+\d+|combo|bundle|"
r"(?:[2-9]|\d{2,})\s*x\s*\d|\d+\s*x\s*(?:[2-9]|\d{2,})\b|multipack|value\s+pack\s+of)\b"
r"|(?:\d\s*(?:g|gm|gms|kg|ml|l|ltr)\s*[x×*]\s*\d)",
re.I,
)
# A long title is several listings run together by the provider - see
# retail_presence.MAX_LISTING_TITLE_CHARS for the measured case.
MAX_TITLE_CHARS = 160
# Pack sizes. Mass and volume in the units the catalogue stores; counts become
# "pcs" because that is the count unit brand_discovery's title parser reads.
_MASS_VOLUME = re.compile(
r"(\d+(?:\.\d+)?)\s*(kilograms?|kgs?|grams?|gms?|gm|g|millilitres?|milliliters?|mls?|ml|"
r"litres?|liters?|ltrs?|ltr|l)\b", re.I)
_COUNT = re.compile(
r"(\d+)\s*(pcs|pc|pieces?|count|ct|tablets?|capsules?|lozenges?|sachets?|rolls?|"
r"sticks?|condoms?|wipes|napkins|pads|units?|n)\b", re.I)
_UNIT_CANON = {
"kilogram": "kg", "kilograms": "kg", "kg": "kg", "kgs": "kg",
"gram": "g", "grams": "g", "gm": "g", "gms": "g", "g": "g",
"millilitre": "ml", "millilitres": "ml", "milliliter": "ml", "milliliters": "ml",
"ml": "ml", "mls": "ml",
"litre": "l", "litres": "l", "liter": "l", "liters": "l", "ltr": "l", "ltrs": "l", "l": "l",
}
_PRICE = re.compile(r"(?:₹|rs\.?|inr)\s*(\d{1,6}(?:\.\d{1,2})?)", re.I)
def _number(raw: str) -> str:
value = float(raw)
return str(int(value)) if value.is_integer() else str(value)
def extract_size(text: str) -> Tuple[Optional[str], int]:
"""(the pack size in canonical form, how many DIFFERENT sizes were named).
"Dettol Soap (125 g)" -> ("125g", 1); "Harpic 500ml + 200ml" -> (.., 2),
which the caller refuses: one listing, one pack.
"""
found: List[str] = []
for num, unit in _MASS_VOLUME.findall(text or ""):
canon = _UNIT_CANON.get(unit.lower())
if canon and float(num) > 0:
size = f"{_number(num)}{canon}"
if size not in found:
found.append(size)
if not found:
for num, _unit in _COUNT.findall(text or ""):
if int(num) > 0:
size = f"{int(num)}pcs"
if size not in found:
found.append(size)
return (found[0] if found else None), len(found)
def extract_price(text: str) -> Optional[float]:
"""The first rupee amount in a snippet, if it is a plausible shelf price."""
for raw in _PRICE.findall(text or ""):
value = float(raw)
if 1 <= value <= 100_000:
return value
return None
def _strip_sizes(text: str) -> str:
"""The name without its pack size - keeping a variant that shares a
bracket with it: "(Lavender - 500 ml)" -> "(Lavender)", "(125 g)" -> ""."""
text = re.sub(r"\bpack\s+of\s+1\b", " ", text, flags=re.I)
text = _MASS_VOLUME.sub(" ", text)
text = _COUNT.sub(" ", text)
text = re.sub(r"\(\s*[,;\-]*\s*", "(", text)
text = re.sub(r"\s*[,;\-]*\s*\)", ")", text)
text = re.sub(r"\(\s*\)", " ", text)
return text
def _tidy(text: str) -> str:
text = re.sub(r"\s+", " ", text)
return text.strip(" ,;:.-/&+–—")
def _words(text: str) -> List[str]:
return re.findall(r"[a-z0-9]+", (text or "").lower())
# Words a retail title may put before the brand ("New Dettol ...").
_LEAD_IN = {"new", "the", "original"}
def names_brand(title: str, brand_terms: Sequence[str]) -> bool:
"""True when `title` LEADS with one of `brand_terms`, as whole words.
Retail titles put the brand first ("Dettol Original Soap", "Maggi 2-Minute
Noodles"). Requiring it there - optionally after "New"/"The" - is what
keeps "Haldiram's Soan Papdi Mithai" from counting as an "Amul Mithai"
listing, or a Savlon listing that mentions Dettol from counting as Dettol.
"""
words = _words(title)
starts = [0] + ([1] if words and words[0] in _LEAD_IN else [])
for term in brand_terms:
term_words = _words(term)
if not term_words:
continue
for start in starts:
if words[start:start + len(term_words)] == term_words:
return True
return False
def _protect_brackets(title: str) -> str:
"""A separator INSIDE brackets is part of the variant, not a segment break:
"(Lavender - 500 ml)" must not split into "(Lavender" and "500 ml)"."""
return re.sub(r"\(([^()]*)\)",
lambda m: "(" + re.sub(r"\s+[-–—|:]\s+", " ", m.group(1)) + ")",
title or "")
def clean_title(title: str, brand_terms: Sequence[str]) -> str:
"""The product name inside a result title: shop furniture removed.
Starts at the first segment that leads with the brand and keeps following
segments - a variant ("(Original)") or the size ("500 ml") - until the
shop's own wording begins. "Buy X Online | Blinkit", "X : Amazon.in:
Beauty" and "X - (Original) - 500 ml Price - Buy Online" all reduce to X
with its variant and size.
"""
collected: List[str] = []
for segment in _SEGMENT_SPLIT.split(_protect_brackets(title)):
segment = _LEADING_BUY.sub("", segment)
cut = _FURNITURE.search(segment)
content = _tidy(segment[:cut.start()] if cut else segment)
if not collected and not names_brand(content, brand_terms):
continue # before the name: "Buy", a site banner
if content:
collected.append(content)
if cut:
break
return _tidy(" ".join(collected))
# ---------------------------------------------------------------------------
# One hit -> one listing
# ---------------------------------------------------------------------------
@dataclass
class Listing:
retailer: str # domain, e.g. "blinkit.com"
url: str
raw_title: str
name: str # cleaned product name, size removed
size: str # canonical, e.g. "125g", "10pcs"
brand_term: str # which brand/sub-brand it named
price: Optional[float] = None
tokens: List[str] = field(default_factory=list)
def as_dict(self) -> dict:
return {"retailer": display_name(self.retailer), "domain": self.retailer,
"url": self.url, "title": self.raw_title, "price": self.price}
# Why a hit was refused. Counted per job so the admin can see what was thrown
# away and why, rather than trusting a silent filter.
NOT_A_RETAILER = "not_a_retailer"
NOT_A_PRODUCT_PAGE = "not_a_product_page"
TITLE_TOO_LONG = "title_too_long"
WRONG_BRAND = "wrong_brand"
NO_PACK_SIZE = "no_pack_size"
SEVERAL_SIZES = "several_sizes"
MULTIPACK = "multipack_or_combo"
NO_NAME = "no_product_name"
def parse_hit(title: str, url: str, snippet: str,
brand_terms: Sequence[str]) -> Tuple[Optional[Listing], Optional[str]]:
"""(listing, None) when the hit is a trustworthy product listing, else (None, reason)."""
retailer = retailer_for(url)
if retailer is None:
return None, NOT_A_RETAILER
if not is_product_url(retailer, url):
return None, NOT_A_PRODUCT_PAGE
if len(title or "") > MAX_TITLE_CHARS:
return None, TITLE_TOO_LONG
cleaned = clean_title(title, brand_terms)
if not cleaned:
return None, WRONG_BRAND
if _MULTIPACK.search(title or ""):
return None, MULTIPACK
size, distinct = extract_size(cleaned)
if size is None:
# Some retailers put the size only in the title's tail ("... - 125 g | Zepto").
size, distinct = extract_size(title)
if size is None:
return None, NO_PACK_SIZE
if distinct > 1:
return None, SEVERAL_SIZES
name = _tidy(_strip_sizes(cleaned))
if not name or not re.search(r"[a-z]", name, re.I):
return None, NO_NAME
term = next((t for t in brand_terms if names_brand(cleaned, [t])), brand_terms[0])
return Listing(
retailer=retailer, url=url, raw_title=title.strip(), name=name, size=size,
brand_term=term, price=extract_price(snippet) or extract_price(title),
tokens=_words(name),
), None

View File

@@ -0,0 +1,117 @@
"""After a Brand Discovery ingest: record WHERE each web-found product was seen.
The ingest CSV has no column for it, and the 11-stage pipeline is not changed
to carry one, so this runs AFTER the batch, the same way
capture_discovery.mark_captured_row does. For each row the batch stored from a
web-found product it
* merges `field_sources.web_listings` = the retailer listings (retailer, URL,
listing title, price) that justified the product, with when they were seen;
* fills `price_range` from the listing prices when the pipeline left it blank;
* sets `validation_status = 'needs_review'` when only ONE retailer listed it
(never lifting a `rejected` verdict). Two or more retailers keep whatever the
validator decided.
The join is exact: stage 11 reports every stored row's `image_id` with the
1-based sheet row it came from (header = row 1), and the ingest wrote the web
products in a known order, so CSV row i is `source_row` i + 2.
"""
from __future__ import annotations
import logging
import threading
import time
from typing import Any, Dict, List, Optional
logger = logging.getLogger(__name__)
_POLL_SECONDS = 5.0
_GIVE_UP_SECONDS = 6 * 3600
def web_updates(manifest_files: List[Any], entries: Dict[int, Dict[str, Any]]) -> List[Dict[str, Any]]:
"""(image_id, table-brand, entry) for every stored row that came from a
web-found product. `entries` is keyed by the product's 0-based CSV index."""
out: List[Dict[str, Any]] = []
for file in manifest_files:
result = getattr(file, "result", None) or {}
for product in result.get("products") or []:
row = product.get("source_row")
image_id = product.get("image_id")
if not image_id or not isinstance(row, int):
continue
entry = entries.get(row - 2)
if entry:
out.append({"image_id": image_id, "brand": product.get("brand"), "entry": entry})
return out
def apply(brand: str, updates: List[Dict[str, Any]]) -> int:
"""Write the provenance. Returns rows updated. Never raises."""
if not updates:
return 0
from psycopg.types.json import Json
from app.services.vector_store import _connect, _table_name
conn = _connect()
if conn is None:
logger.warning("web discovery provenance: no database connection")
return 0
seen_at = time.strftime("%Y-%m-%dT%H:%M:%S%z")
written = 0
try:
with conn, conn.cursor() as cur:
for update in updates:
entry = update["entry"]
table = _table_name(update.get("brand") or brand)
single = int(entry.get("retailer_count") or 0) < 2
provenance = {"web_listings": {
"seen_at": seen_at,
"retailer_count": int(entry.get("retailer_count") or 0),
"listings": list(entry.get("listings") or [])[:10],
}}
cur.execute(
f"UPDATE {table} SET "
f"field_sources = COALESCE(field_sources, '{{}}'::jsonb) || %s, "
f"price_range = CASE WHEN COALESCE(price_range, '') = '' AND %s <> '' "
f"THEN %s ELSE price_range END, "
f"validation_status = CASE WHEN %s AND COALESCE(validation_status, '') <> 'rejected' "
f"THEN 'needs_review' ELSE validation_status END "
f"WHERE image_id = %s",
(Json(provenance), entry.get("price_range") or "", entry.get("price_range") or "",
single, update["image_id"]),
)
written += cur.rowcount
except Exception as exc: # noqa: BLE001 - provenance must not take the batch down
logger.exception("web discovery provenance failed for %s: %s", brand, exc)
finally:
conn.close()
logger.info("web discovery provenance: %d row(s) of %s marked with their listings", written, brand)
return written
def watch_batch(batch_id: str, brand: str, entries: Dict[int, Dict[str, Any]], *,
read_manifest=None, sleep=time.sleep, give_up_seconds: float = _GIVE_UP_SECONDS) -> Optional[int]:
"""Block until the batch ends, then apply. Returns rows written, or None
when the batch never finished in time. Run it on a thread - see `watch`."""
if read_manifest is None:
from app.core.batch_ingest import read_manifest
from app.core.batch_ingest import TERMINAL_BATCH_STATES
deadline = time.monotonic() + give_up_seconds
while time.monotonic() < deadline:
manifest = read_manifest(batch_id)
if manifest is not None and manifest.status in TERMINAL_BATCH_STATES:
return apply(brand, web_updates(manifest.files, entries))
sleep(_POLL_SECONDS)
logger.warning("web discovery provenance: batch %s did not finish in time", batch_id)
return None
def watch(batch_id: str, brand: str, entries: Dict[int, Dict[str, Any]]) -> None:
"""Fire-and-forget `watch_batch` on a daemon thread."""
if not entries:
return
threading.Thread(target=watch_batch, args=(batch_id, brand, entries),
name=f"web-provenance-{batch_id[:8]}", daemon=True).start()

View File

@@ -0,0 +1,105 @@
"""One web search, from whichever backend is available.
search(query) -> list[Hit] the engine answered (possibly with nothing)
-> None the engine could not be asked (throttled,
unreachable, not installed)
The two are different answers and every caller depends on it: None is "we do
not know" and must never be recorded or reported as "no products", which is
the same rule retail_presence._search keeps.
Backends, in order:
* Google Custom Search JSON API, when USE_GOOGLE_CSE is on and both keys are
set. The official API - we never scrape google.com result pages. 100 free
queries a day; a quota error falls through to DuckDuckGo.
* DuckDuckGo through the `ddgs` package the project already uses for images
and retail presence. Free, no key, throttles bursts.
"""
from __future__ import annotations
import logging
from dataclasses import dataclass
from typing import List, Optional
from app.infrastructure import settings
logger = logging.getLogger(__name__)
_GOOGLE_URL = "https://www.googleapis.com/customsearch/v1"
_GOOGLE_MAX_NUM = 10 # the API's per-request ceiling
_TIMEOUT_SECONDS = 15
@dataclass(frozen=True)
class Hit:
title: str
url: str
snippet: str = ""
def as_dict(self) -> dict:
return {"title": self.title, "url": self.url, "snippet": self.snippet}
@classmethod
def from_dict(cls, data: dict) -> "Hit":
return cls(title=str(data.get("title") or ""), url=str(data.get("url") or ""),
snippet=str(data.get("snippet") or ""))
def google_configured() -> bool:
return bool(settings.USE_GOOGLE_CSE and settings.GOOGLE_API_KEY and settings.GOOGLE_CSE_ID)
def backend_name() -> str:
return "google_cse" if google_configured() else "duckduckgo"
def _google(query: str, max_results: int) -> Optional[List[Hit]]:
import requests
params = {
"key": settings.GOOGLE_API_KEY, "cx": settings.GOOGLE_CSE_ID, "q": query,
"num": min(_GOOGLE_MAX_NUM, max(1, max_results)), "gl": "in",
}
try:
response = requests.get(_GOOGLE_URL, params=params, timeout=_TIMEOUT_SECONDS)
except Exception as exc: # noqa: BLE001 - unreachable is "could not ask"
logger.debug("Google CSE unreachable for %r: %s", query, exc)
return None
if response.status_code != 200:
# 429 / 403 dailyLimitExceeded: a quota, not an answer.
logger.info("Google CSE answered %s for %r; falling back", response.status_code, query)
return None
try:
items = response.json().get("items") or []
except ValueError:
return None
return [Hit(title=str(i.get("title") or ""), url=str(i.get("link") or ""),
snippet=str(i.get("snippet") or "")) for i in items if i.get("link")]
def _duckduckgo(query: str, max_results: int) -> Optional[List[Hit]]:
try:
from ddgs import DDGS
except ImportError:
logger.debug("web discovery: 'ddgs' is not installed")
return None
try:
with DDGS(timeout=_TIMEOUT_SECONDS) as ddgs:
rows = list(ddgs.text(query, region="in-en", safesearch="off",
max_results=max_results) or [])
except Exception as exc: # noqa: BLE001 - a throttle is not a verdict
logger.debug("DuckDuckGo failed for %r: %s", query, exc)
return None
return [Hit(title=str(r.get("title") or ""), url=str(r.get("href") or r.get("url") or ""),
snippet=str(r.get("body") or "")) for r in rows if (r.get("href") or r.get("url"))]
def search(query: str, max_results: Optional[int] = None) -> Optional[List[Hit]]:
"""The hits for `query`, or None when no backend could be asked."""
limit = max_results or settings.WEB_DISCOVERY_RESULTS_PER_QUERY
if google_configured():
hits = _google(query, limit)
if hits is not None:
return hits
return _duckduckgo(query, limit)

View File

@@ -1,7 +1,10 @@
{
"brand": "reckitt benckiser",
"schema": 2,
"endpoint": "https://world.openfoodfacts.org/api/v2/search",
"brand": "Reckitt Benckiser",
"brand_tag": "reckitt-benckiser",
"country": "india",
"fetched_at": 1788851995.867008,
"fetched_at_human": "2026-09-08 12:49:55",
"fetched_at": 1790675563.0787868,
"fetched_at_human": "2026-09-29 15:22:43",
"hits": []
}

BIN
data/cache/web_discovery/search_cache.db vendored Normal file

Binary file not shown.

453
tests/test_web_discovery.py Normal file
View File

@@ -0,0 +1,453 @@
"""Web & retail listing discovery (app/services/web_discovery) - no network.
Every search is a recorded fixture. What is pinned here is the anti-invention
contract: a product exists only if a real retailer's PRODUCT page names the
brand and a pack size; a throttled search is "could not ask", never "nothing
there"; and with the source switched off Brand Discovery is exactly what it was.
"""
from __future__ import annotations
from typing import Dict, List, Optional
import pytest
from app.infrastructure import settings
from app.services import brand_discovery as bd
from app.services.web_discovery import brands, cache, discover, jobs, listings, provenance, search
from app.services.web_discovery.search import Hit
DETTOL = ("dettol",)
# ---------------------------------------------------------------------------
# listings.parse_hit - one hit -> one trustworthy listing, or a reason
# ---------------------------------------------------------------------------
@pytest.mark.parametrize("title,url,snippet,name,size,price", [
("Dettol Original Soap (125 g) - Buy Online at Best Price | Blinkit",
"https://blinkit.com/prn/dettol-original-soap/prid/12345", "", "Dettol Original Soap", "125g", None),
("Buy Dettol Original Soap 125 g Online at Best Price | Zepto",
"https://www.zeptonow.com/pn/dettol-original-soap/pvid/1a2b3c4d-1111", "", "Dettol Original Soap", "125g", None),
("Buy Dettol Original Soap 125 g Online at Best Price of Rs 58 - bigbasket",
"https://www.bigbasket.com/pd/40001234/dettol-soap/", "MRP Rs 58", "Dettol Original Soap", "125g", 58.0),
("Dettol Original Soap Bar, 125g : Amazon.in: Beauty",
"https://www.amazon.in/Dettol-Original/dp/B00ABCDEFG", "₹ 55", "Dettol Original Soap Bar", "125g", 55.0),
("Buy Dettol Antiseptic Liquid 1 L Online | Swiggy Instamart",
"https://www.swiggy.com/instamart/item/ABC123XYZ", "", "Dettol Antiseptic Liquid", "1l", None),
("Dettol Antiseptic Liquid 550 ml - JioMart",
"https://www.jiomart.com/p/groceries/dettol-antiseptic/490001234", "", "Dettol Antiseptic Liquid", "550ml", None),
])
def test_each_retailers_title_shape_reduces_to_the_product(title, url, snippet, name, size, price):
listing, reason = listings.parse_hit(title, url, snippet, DETTOL)
assert reason is None
assert (listing.name, listing.size, listing.price) == (name, size, price)
@pytest.mark.parametrize("title,url,reason", [
("Dettol Soaps - Buy Dettol Soaps Online | Blinkit",
"https://blinkit.com/cn/dettol/cid/1", listings.NOT_A_PRODUCT_PAGE),
("Dettol Original Soap 125g", "https://www.example-blog.com/dettol-review", listings.NOT_A_RETAILER),
("Dettol Original Soap 125g", "https://www.amazon.in/s?k=dettol", listings.NOT_A_PRODUCT_PAGE),
("Dettol Liquid Handwash Refill | Blinkit", "https://blinkit.com/prn/x/prid/7", listings.NO_PACK_SIZE),
("Dettol Original Soap 125g (Pack of 4) : Amazon.in", "https://www.amazon.in/x/dp/B00ABCDEFH",
listings.MULTIPACK),
("Dettol Soap 125g x 3 | Zepto", "https://www.zeptonow.com/pn/x/pvid/1a2b3c4d-9999", listings.MULTIPACK),
("Dettol Handwash 200ml + Refill 175ml | Blinkit", "https://blinkit.com/prn/x/prid/8",
listings.SEVERAL_SIZES),
])
def test_what_is_not_a_listing_of_one_pack_is_refused(title, url, reason):
assert listings.parse_hit(title, url, "", DETTOL) == (None, reason)
def test_another_brand_is_refused_even_when_it_mentions_ours():
listing, reason = listings.parse_hit("Savlon Soap 125g better than Dettol | Blinkit",
"https://blinkit.com/prn/savlon/prid/9", "", DETTOL)
assert listing is None and reason == listings.WRONG_BRAND
def test_the_brand_must_lead_the_title():
assert listings.names_brand("Dettol Original Soap", DETTOL)
assert listings.names_brand("New Dettol Original Soap", DETTOL)
assert not listings.names_brand("Haldiram Mithai with Dettol", DETTOL)
def test_counts_become_pcs_the_unit_discovery_reads():
listing, _ = listings.parse_hit("Durex Extra Thin Condoms 10 Count : Amazon.in",
"https://www.amazon.in/x/dp/B00ABCDEFJ", "", ("durex",))
assert listing.size == "10pcs" and listing.name == "Durex Extra Thin Condoms"
def test_a_run_together_title_is_refused():
title = "Dettol Original Soap 125g " + "x" * 200
assert listings.parse_hit(title, "https://blinkit.com/prn/x/prid/1", "", DETTOL)[1] == listings.TITLE_TOO_LONG
# ---------------------------------------------------------------------------
# brands.targets_for
# ---------------------------------------------------------------------------
def test_reckitt_fans_out_to_the_names_its_products_are_sold_under():
queries = [t.query for t in brands.targets_for("Reckitt Benckiser")]
for sub in ("Dettol", "Harpic", "Lizol", "Mortein", "Durex"):
assert sub in queries
assert queries[-1] == "Reckitt Benckiser"
def test_a_typed_sub_brand_searches_only_itself():
assert [(t.query, t.brand_terms) for t in brands.targets_for("Dettol")] == [("Dettol", ("dettol",))]
def test_a_generic_product_word_is_never_a_brand_word():
butter = next(t for t in brands.targets_for("Amul") if t.query == "Amul Butter")
assert butter.brand_terms == ("amul",)
def test_a_line_name_sold_without_the_parent_counts_on_its_own():
cerelac = next(t for t in brands.targets_for("Nestle") if t.query == "Nestle Cerelac")
assert "cerelac" in cerelac.brand_terms
def test_kit_kat_is_not_an_abbreviation_plus_a_sub_brand():
assert not brands._is_abbreviation("kit", "nestle")
assert brands._is_abbreviation("rb", "reckitt benckiser")
# ---------------------------------------------------------------------------
# discover.run - the loop, with a fake search engine
# ---------------------------------------------------------------------------
@pytest.fixture(autouse=True)
def _isolated_dir(tmp_path, monkeypatch):
monkeypatch.setattr(settings, "WEB_DISCOVERY_DIR", tmp_path / "web_discovery")
monkeypatch.setattr(settings, "WEB_DISCOVERY_PAUSE_SECONDS", 0.0)
monkeypatch.setattr(search, "google_configured", lambda: False)
jobs._jobs.clear()
yield
jobs._jobs.clear()
class FakeEngine:
def __init__(self, answers: Dict[str, Optional[List[Hit]]]):
self.answers = answers
self.asked: List[str] = []
def __call__(self, query, max_results=None):
self.asked.append(query)
for key, hits in self.answers.items():
if key in query:
return hits
return []
BLINKIT_SOAP = Hit("Dettol Original Soap (125 g) - Buy Online at Best Price | Blinkit",
"https://blinkit.com/prn/dettol-original-soap/prid/1", "₹ 50")
ZEPTO_SOAP = Hit("Buy Dettol Original Soap 125 g Online at Best Price | Zepto",
"https://www.zeptonow.com/pn/dettol-original-soap/pvid/1a2b3c4d-0001", "₹ 55")
ZEPTO_LIQUID = Hit("Buy Dettol Antiseptic Liquid 250 ml Online | Zepto",
"https://www.zeptonow.com/pn/dettol-liquid/pvid/1a2b3c4d-0002", "")
def test_the_same_pack_on_two_retailers_is_one_candidate_with_both(monkeypatch):
engine = FakeEngine({"site:blinkit.com": [BLINKIT_SOAP], "site:zeptonow.com": [ZEPTO_SOAP, ZEPTO_LIQUID]})
monkeypatch.setattr(search, "search", engine)
result = discover.run("Dettol", sleep=lambda s: None)
assert result.status == discover.DONE
by_title = {c["title"]: c for c in result.candidates}
soap = by_title["Dettol Original Soap 125g"]
assert soap["providers"] == ["Blinkit", "Zepto"] and soap["retailer_count"] == 2
assert soap["price_range"] == "₹50 - ₹55"
assert by_title["Dettol Antiseptic Liquid 250ml"]["retailer_count"] == 1
assert result.candidates[0]["title"] == "Dettol Original Soap 125g" # best corroborated first
def test_a_generic_title_does_not_swallow_a_variant():
a, _ = listings.parse_hit("Dettol Soap 125g | Blinkit", "https://blinkit.com/prn/a/prid/1", "", DETTOL)
b, _ = listings.parse_hit("Dettol Cool Soap 125g | Zepto",
"https://www.zeptonow.com/pn/b/pvid/1a2b3c4d-0003", "", DETTOL)
assert len(discover.cluster([a, b])) == 2
def test_three_unanswered_searches_end_the_run_as_partial_and_are_not_cached(monkeypatch):
engine = FakeEngine({"site:": None})
monkeypatch.setattr(search, "search", engine)
result = discover.run("Dettol", sleep=lambda s: None)
assert result.status == discover.PARTIAL and result.queries_failed == 3
assert len(engine.asked) == 3 and result.candidates == []
assert cache.get(engine.asked[0], search.backend_name()) is None # "could not ask" is never stored
def test_a_single_unanswered_search_still_makes_the_run_partial(monkeypatch):
answers = {"site:blinkit.com": None, "site:zeptonow.com": [ZEPTO_SOAP]}
monkeypatch.setattr(search, "search", FakeEngine(answers))
result = discover.run("Dettol", sleep=lambda s: None)
assert result.status == discover.PARTIAL and result.queries_failed == 1
assert [c["title"] for c in result.candidates] == ["Dettol Original Soap 125g"]
def test_answers_are_cached_so_a_rerun_asks_nothing(monkeypatch):
engine = FakeEngine({"site:blinkit.com": [BLINKIT_SOAP]})
monkeypatch.setattr(search, "search", engine)
discover.run("Dettol", sleep=lambda s: None)
asked = len(engine.asked)
again = discover.run("Dettol", sleep=lambda s: None)
assert len(engine.asked) == asked and again.queries_cached == again.queries_total
def test_the_budget_cuts_the_last_retailers_not_the_last_sub_brands():
planned = discover.plan_queries("Reckitt Benckiser", max_queries=11)
assert {q.retailer for q in planned} == {"blinkit.com"}
assert len({q.target.query for q in planned}) == 11
def test_refused_hits_are_counted_by_reason(monkeypatch):
category = Hit("Dettol - Buy Online | Blinkit", "https://blinkit.com/cn/dettol/cid/1", "")
monkeypatch.setattr(search, "search", FakeEngine({"site:blinkit.com": [category, BLINKIT_SOAP]}))
result = discover.run("Dettol", sleep=lambda s: None)
assert result.rejected == {listings.NOT_A_PRODUCT_PAGE: 1}
assert result.listings_kept == 1
# ---------------------------------------------------------------------------
# jobs
# ---------------------------------------------------------------------------
def test_a_finished_job_is_reused_and_found_from_disk(monkeypatch):
monkeypatch.setattr(search, "search", FakeEngine({"site:blinkit.com": [BLINKIT_SOAP]}))
job = jobs.WebJob(job_id="a" * 32, brand="Dettol")
jobs._jobs[job.job_id] = job
jobs.run_job(job)
jobs._jobs.clear() # as after a restart
found = jobs.candidates_for("dettol")
assert found.job_id == job.job_id and found.candidates[0]["title"] == "Dettol Original Soap 125g"
assert jobs.start_job("Dettol").job_id == job.job_id
def test_a_job_that_was_running_when_the_process_died_is_interrupted():
job = jobs.WebJob(job_id="b" * 32, brand="Dettol", status=jobs.RUNNING)
jobs._save(job)
assert jobs._load(job.job_id).status == jobs.INTERRUPTED
def test_a_job_id_that_is_not_hex_is_never_a_path():
assert jobs.get_job("../../etc/passwd") is None
# ---------------------------------------------------------------------------
# brand_discovery with the web source
# ---------------------------------------------------------------------------
@pytest.fixture
def no_other_sources(monkeypatch):
monkeypatch.setattr(bd, "_from_open_facts", lambda brand, refresh=False: [])
monkeypatch.setattr(bd, "_from_llm", lambda brand, deadline, budget: [])
monkeypatch.setattr(bd, "_from_brand_store", lambda brand: [])
monkeypatch.setattr(bd, "get_products_by_brand", lambda brand, **kw: [])
def _finished_job(candidates, status=jobs.DONE) -> jobs.WebJob:
job = jobs.WebJob(job_id="c" * 32, brand="Reckitt Benckiser", status=status, backend="duckduckgo",
queries_total=77, queries_done=77, listings_kept=3, candidates=candidates)
jobs._jobs[job.job_id] = job
return job
def _candidate(title, size, providers):
return {"title": title, "size": size, "sizes": [size], "providers": providers,
"retailer_count": len(providers), "price_range": "₹58",
"listings": [{"retailer": p, "url": f"https://{p.lower()}.example/{i}", "title": title}
for i, p in enumerate(providers)],
"brand_term": "dettol", "source": "web"}
def test_web_rows_arrive_with_their_listings_and_honest_confidence(no_other_sources):
job = _finished_job([_candidate("Dettol Original Soap 125g", "125g", ["Blinkit", "Zepto"]),
_candidate("Harpic Power Plus 500ml", "500ml", ["Amazon"])])
result = bd.discover_brand_products("Reckitt Benckiser", use_llm=False, use_web=True,
web_job_id=job.job_id)
rows = {p.product_name: p for p in result.products}
soap = rows["Dettol Original Soap"]
assert soap.size_variants == ["125g"] and soap.evidence == "retail"
assert soap.confidence == bd.WEB_CONFIDENCE_MULTI and soap.as_preview()["selected"] is True
harpic = rows["Harpic Power Plus"]
assert harpic.confidence == bd.WEB_CONFIDENCE_SINGLE and harpic.as_preview()["retailer_count"] == 1
assert soap.as_preview()["listings"][0]["retailer"] == "Blinkit"
assert result.counts["from_web"] == 2
assert not any("language model alone" in w for w in result.warnings)
assert any(w.startswith("Web & retail listings: 2 product(s)") for w in result.warnings)
def test_web_on_but_never_searched_says_so(no_other_sources):
result = bd.discover_brand_products("Reckitt Benckiser", use_llm=False, use_web=True)
assert result.products == []
assert any("have not been searched" in w for w in result.warnings)
def test_a_partial_job_says_it_is_incomplete(no_other_sources):
job = _finished_job([_candidate("Dettol Original Soap 125g", "125g", ["Blinkit"])], status=jobs.PARTIAL)
job.detail = "the search provider stopped answering"
result = bd.discover_brand_products("Reckitt Benckiser", use_llm=False, use_web=True)
assert any("Incomplete: the search provider stopped answering" in w for w in result.warnings)
def test_with_the_web_source_off_nothing_changes(no_other_sources):
_finished_job([_candidate("Dettol Original Soap 125g", "125g", ["Blinkit", "Zepto"])])
result = bd.discover_brand_products("Reckitt Benckiser", use_llm=False)
assert result.products == [] and "from_web" not in result.counts
assert "Every row below rests on the language model alone - review them individually." in result.warnings
def test_a_web_listing_joins_an_open_food_facts_row_instead_of_duplicating_it(no_other_sources, monkeypatch):
monkeypatch.setattr(bd, "_from_open_facts", lambda brand, refresh=False: [
{"title": "Dettol Original Soap", "barcode": "8901396393009", "size": "125g", "source": "off"}])
_finished_job([_candidate("Dettol Original Soap 125g", "125g", ["Blinkit", "Zepto"])])
result = bd.discover_brand_products("Reckitt Benckiser", use_llm=False, use_web=True)
[row] = [p for p in result.products if "Soap" in p.product_name]
assert row.sources == ["off", "web"] and row.barcode == "8901396393009"
assert row.providers[:2] == ["Blinkit", "Zepto"] and row.retailer_count == 2
assert row.evidence == "openfacts"
# ---------------------------------------------------------------------------
# provenance
# ---------------------------------------------------------------------------
class _File:
def __init__(self, products):
self.result = {"products": products}
class _Manifest:
def __init__(self, status, files):
self.status, self.files = status, files
def test_stored_rows_are_joined_back_by_csv_position():
entries = {0: {"retailer_count": 2}, 2: {"retailer_count": 1}}
files = [_File([
{"image_id": "rb_dettol_soap_125g", "brand": "Reckitt Benckiser", "source_row": 2},
{"image_id": "rb_other", "brand": "Reckitt Benckiser", "source_row": 3},
{"image_id": "rb_harpic_500ml", "brand": "Reckitt Benckiser", "source_row": 4},
])]
updates = provenance.web_updates(files, entries)
assert [(u["image_id"], u["entry"]["retailer_count"]) for u in updates] == [
("rb_dettol_soap_125g", 2), ("rb_harpic_500ml", 1)]
def test_the_watcher_waits_for_the_batch_then_applies(monkeypatch):
states = iter([_Manifest("running", []), _Manifest("done", [_File([
{"image_id": "x", "brand": "Reckitt Benckiser", "source_row": 2}])])])
applied = []
monkeypatch.setattr(provenance, "apply", lambda brand, updates: applied.append(updates) or len(updates))
written = provenance.watch_batch("batch1", "Reckitt Benckiser", {0: {"retailer_count": 1}},
read_manifest=lambda _id: next(states), sleep=lambda s: None)
assert written == 1 and applied[0][0]["image_id"] == "x"
class _Cursor:
def __init__(self):
self.calls = []
self.rowcount = 1
def __enter__(self):
return self
def __exit__(self, *exc):
return False
def execute(self, sql, params):
self.calls.append((sql, params))
class _Conn(_Cursor):
def __init__(self):
super().__init__()
self.cur = _Cursor()
def cursor(self):
return self.cur
def close(self):
pass
def test_apply_touches_only_provenance_price_and_review_status(monkeypatch):
conn = _Conn()
monkeypatch.setattr("app.services.vector_store._connect", lambda: conn)
provenance.apply("Reckitt Benckiser", [
{"image_id": "one_shop", "brand": "Reckitt Benckiser",
"entry": {"retailer_count": 1, "price_range": "₹58", "listings": [{"url": "u"}]}},
])
sql, params = conn.cur.calls[0]
assert sql.startswith("UPDATE brand_reckitt_benckiser SET field_sources")
assert "price_range = CASE WHEN COALESCE(price_range, '') = ''" in sql
assert "'rejected'" in sql and params[-2] is True and params[-1] == "one_shop"
for column in ("product_name", "barcode", "image_url", "category"):
assert f"{column} =" not in sql
# ---------------------------------------------------------------------------
# Found on the first live run (Reckitt Benckiser, Blinkit titles)
# ---------------------------------------------------------------------------
@pytest.mark.parametrize("title,name,size", [
("Dettol Original Hand Wash Refill 675 ml Price - Buy Online at \u20b992 in...",
"Dettol Original Hand Wash Refill", "675ml"),
("Harpic Disinfectant Liquid Toilet Cleaner - (Original) - 500 ml Price - Buy Online at Best",
"Harpic Disinfectant Liquid Toilet Cleaner (Original)", "500ml"),
("Lizol Disinfectant Surface & Floor Cleaner (Lavender - 500 ml) Price - Buy Online at Best",
"Lizol Disinfectant Surface & Floor Cleaner (Lavender)", "500ml"),
("Buy Lizol Disinfectant Surface & Floor Cleaner (Citrus, 625 ml) Online",
"Lizol Disinfectant Surface & Floor Cleaner (Citrus)", "625ml"),
])
def test_blinkit_names_keep_their_variant_and_lose_the_price_wording(title, name, size):
term = (title.split()[1] if title.startswith("Buy") else title.split()[0]).lower()
listing, reason = listings.parse_hit(title, "https://blinkit.com/prn/x/prid/1", "", (term,))
assert reason is None and (listing.name, listing.size) == (name, size)
def test_two_scents_of_one_product_stay_two_products():
citrus, _ = listings.parse_hit("Lizol Disinfectant Surface & Floor Cleaner (Citrus) 2 l Price - Buy",
"https://blinkit.com/prn/a/prid/1", "", ("lizol",))
floral, _ = listings.parse_hit("Lizol Disinfectant Surface & Floor Cleaner (Floral) - 2 l Price - Buy",
"https://blinkit.com/prn/b/prid/2", "", ("lizol",))
assert len(discover.cluster([citrus, floral])) == 2
def test_discovery_keeps_two_web_variants_its_own_merge_would_fold(no_other_sources):
"""Its title similarity scores "... Cleaner Citrus" vs "... Cleaner Floral"
at 0.857, over its 0.85 floor; web rows were already told apart."""
_finished_job([_candidate("Lizol Floor Cleaner Citrus 2l", "2l", ["Blinkit"]),
_candidate("Lizol Floor Cleaner Floral 2l", "2l", ["Blinkit"])])
result = bd.discover_brand_products("Reckitt Benckiser", use_llm=False, use_web=True)
assert sorted(p.product_name for p in result.products) == [
"Lizol Floor Cleaner Citrus", "Lizol Floor Cleaner Floral"]
def test_a_retail_listing_outranks_our_own_earlier_output_as_evidence(no_other_sources, monkeypatch):
monkeypatch.setattr(bd, "get_products_by_brand", lambda brand, **kw: [
{"product_name": "Dettol Antiseptic Liquid", "title": "Dettol Antiseptic Liquid"}])
_finished_job([_candidate("Dettol Antiseptic Liquid 250ml", "250ml", ["Blinkit"])])
[row] = bd.discover_brand_products("Reckitt Benckiser", use_llm=False, use_web=True).products
assert row.evidence == "retail" and row.matches_existing == "Dettol Antiseptic Liquid"
def test_bracketed_variant_words_survive_into_the_candidate_title():
a, _ = listings.parse_hit("Lizol Floor Cleaner (Citrus) 2 l Price - Buy",
"https://blinkit.com/prn/a/prid/1", "", ("lizol",))
assert discover.cluster([a])[0]["title"] == "Lizol Floor Cleaner Citrus 2l"

View File

@@ -0,0 +1,127 @@
"""HTTP tests for the web & retail listings additions to /api/admin/brand-discovery.
The batch fixtures mirror test_brand_discovery_api.py (read the docstring
there before removing `batch_root` or `no_background_worker`: without them a
test writes batch manifests into the repository's data directory).
"""
from __future__ import annotations
import pytest
from app.api.routers import brand_discovery as router_module
from app.core import batch_ingest
from app.infrastructure import settings
from app.services import brand_discovery as bd
from app.services.web_discovery import jobs, provenance
PREVIEW = "/api/admin/brand-discovery/preview"
INGEST = "/api/admin/brand-discovery/ingest"
WEB_JOBS = "/api/admin/brand-discovery/web-jobs"
@pytest.fixture(autouse=True)
def _isolate(tmp_path, monkeypatch):
from app.services import active_brands, sku_service
from app.core import batch_worker
monkeypatch.setattr(sku_service, "_data_dir", tmp_path / "sku_sequences")
monkeypatch.setattr(batch_ingest, "BATCH_UPLOAD_DIR", tmp_path / "batch_uploads")
monkeypatch.setattr(batch_worker, "submit", lambda batch_id: None)
monkeypatch.setattr(active_brands, "filtering_enabled", lambda: False)
monkeypatch.setattr(active_brands, "is_active_brand", lambda brand: True)
monkeypatch.setattr(settings, "WEB_DISCOVERY_DIR", tmp_path / "web_discovery")
monkeypatch.setattr(router_module, "WEB_DISCOVERY_ENABLED", True)
from app.api.batch_job_store import batch_job_store
batch_job_store._batches.clear()
batch_job_store._cancelled.clear()
jobs._jobs.clear()
yield
batch_job_store._batches.clear()
batch_job_store._cancelled.clear()
jobs._jobs.clear()
@pytest.fixture
def no_worker(monkeypatch):
"""start_job must not spawn a real search thread in a test."""
queued = []
monkeypatch.setattr(jobs, "_ensure_worker", lambda: None)
monkeypatch.setattr(jobs._queue, "put", queued.append)
return queued
def test_the_web_job_routes_are_admin_only(client, user_headers):
assert client.post(WEB_JOBS, json={"brand": "Dettol"}).status_code in (401, 403)
assert client.post(WEB_JOBS, json={"brand": "Dettol"}, headers=user_headers).status_code == 403
assert client.get(f"{WEB_JOBS}/{'a' * 32}").status_code in (401, 403)
def test_starting_a_job_queues_it_and_polling_returns_progress(client, admin_headers, no_worker):
started = client.post(WEB_JOBS, json={"brand": "Reckitt Benckiser"}, headers=admin_headers)
assert started.status_code == 202, started.text
body = started.json()
assert body["status"] == jobs.QUEUED and body["brand"] == "Reckitt Benckiser"
assert "candidates" not in body and body["candidate_count"] == 0
assert no_worker == [body["job_id"]]
polled = client.get(f"{WEB_JOBS}/{body['job_id']}", headers=admin_headers)
assert polled.status_code == 200 and polled.json()["job_id"] == body["job_id"]
def test_a_second_start_while_running_returns_the_same_job(client, admin_headers, no_worker):
first = client.post(WEB_JOBS, json={"brand": "Dettol"}, headers=admin_headers).json()
second = client.post(WEB_JOBS, json={"brand": "dettol"}, headers=admin_headers).json()
assert first["job_id"] == second["job_id"] and len(no_worker) == 1
def test_an_unknown_job_is_404(client, admin_headers):
assert client.get(f"{WEB_JOBS}/{'f' * 32}", headers=admin_headers).status_code == 404
def test_switched_off_the_routes_say_so(client, admin_headers, monkeypatch):
monkeypatch.setattr(router_module, "WEB_DISCOVERY_ENABLED", False)
assert client.post(WEB_JOBS, json={"brand": "Dettol"}, headers=admin_headers).status_code == 404
def test_preview_passes_the_web_flag_through_and_defaults_it_off(client, admin_headers, monkeypatch):
seen = []
def fake_discover(brand, **kwargs):
seen.append(kwargs)
return bd.DiscoveryResult(brand=brand, parent_brand=brand.lower(), table="brand_x",
brand_active=True, filtering_enabled=False)
monkeypatch.setattr(router_module.brand_discovery, "discover_brand_products", fake_discover)
client.post(PREVIEW, json={"brand": "Dettol"}, headers=admin_headers)
client.post(PREVIEW, json={"brand": "Dettol", "use_web": True, "web_job_id": "c" * 32},
headers=admin_headers)
assert (seen[0]["use_web"], seen[0]["web_job_id"]) == (False, None)
assert (seen[1]["use_web"], seen[1]["web_job_id"]) == (True, "c" * 32)
def test_ingest_hands_web_rows_to_the_provenance_watcher_by_csv_position(client, admin_headers, monkeypatch):
watched = []
monkeypatch.setattr(provenance, "watch", lambda batch_id, brand, entries: watched.append(
(batch_id, brand, entries)))
listing = {"retailer": "Blinkit", "url": "https://blinkit.com/prn/x/prid/1", "title": "Dettol Soap 125g"}
resp = client.post(INGEST, json={"brand": "Reckitt Benckiser", "products": [
{"product_name": "Dettol Original Soap", "size_variants": ["125g"]},
{"product_name": "Harpic Power Plus", "size_variants": ["500ml"], "listings": [listing],
"price_range": "₹99", "retailer_count": 1},
]}, headers=admin_headers)
assert resp.status_code == 202, resp.text
[(batch_id, brand, entries)] = watched
assert batch_id == resp.json()["batch_id"] and brand == "Reckitt Benckiser"
assert list(entries) == [1] and entries[1]["listings"] == [listing]
assert entries[1]["price_range"] == "₹99"
def test_ingest_without_web_rows_starts_no_watcher(client, admin_headers, monkeypatch):
watched = []
monkeypatch.setattr(provenance, "watch", lambda *a: watched.append(a))
resp = client.post(INGEST, json={"brand": "Britannia", "products": [
{"product_name": "Marie Gold", "size_variants": ["250g"]}]}, headers=admin_headers)
assert resp.status_code == 202 and watched == []