Files
2026-10-05 11:21:39 +05:30

601 lines
31 KiB
Python

"""Search-first collection of real product listings.
For one category and a set of allow-listed brands:
1. DISCOVER web search `site:<platform> <brand> <category term>` on every
registered platform (marketplaces, national chains, Tamil Nadu
chains, the brand's own site). Only URLs that are single product
pages on a registered platform are kept.
2. EXPAND for each model found, search `<brand> <model> price` to find the
same model on other platforms.
3. COLLECT per URL, by the platform's probe grade:
A/B and breaker closed -> fetch the page politely and read
JSON-LD / meta / spec tables
C (or fetch refused) -> use the search result itself: its
title, snippet price and stock text
4. MATCH link the listing to one canonical variant (match.matcher)
5. ENRICH specs (deterministic, LLM gap-fill grounded in page text) and
images (only from the product's own listings, validated live)
6. VERIFY products with listings on ≥2 sites (≥1 a retailer) become
verified and visible.
Nothing in this module invents a product, price or image: every value is read
from a page or a search result, and stored with that URL and text.
"""
from __future__ import annotations
import hashlib
import json
import logging
import re
from dataclasses import dataclass, field
from decimal import Decimal
from typing import Callable, Dict, List, Optional, Tuple
from urllib.parse import urlparse
from rapidfuzz import fuzz
from app.electronics.db import repository as repo
from app.electronics.extract.embedded_ratings import IGNORED_DOMAINS, embedded_rating
from app.electronics.extract.html_fallback import extract_page, spec_tables, visible_text
from app.electronics.extract.jsonld import extract_products
from app.electronics.extract.serp_parser import clean_result_title, read_price, read_rating, read_stock
from app.electronics.match.matcher import decide
from app.electronics.models import Listing
from app.electronics.net.breaker import CircuitBreaker
from app.electronics.net.polite_client import PoliteClient
from app.electronics.normalise.brand_alias import looks_like_device_title
from app.electronics.normalise.llm_fill import fill_missing
from app.electronics.normalise.spec_normaliser import normalise_specs
from app.electronics.normalise.title_parser import ParsedTitle, fill_from_context, parse_title, variant_key
from app.electronics.probe.site_probe import probe_site
from app.electronics.reference import BrandRef, SiteRef, load_reference, site_for_url
from app.electronics.search.engine import SearchEngine
from app.electronics.search.providers import SearchHit
from app.infrastructure.settings import ELEC_PROBE_TTL_DAYS, MIN_IMAGE_BYTES
from bs4 import BeautifulSoup
logger = logging.getLogger(__name__)
# Titles that belong to another category even when the brand matches.
_OFF_CATEGORY = {
"mobiles": re.compile(r"\b(?:tab|tablet|pad|watch|buds|earbuds|laptop|book|monitor|tv|television|band)\b", re.I),
"laptops": re.compile(r"\b(?:tablet|tab|monitor|mouse|keyboard|phone|smartphone|printer|desktop|all[- ]in[- ]one)\b", re.I),
}
_LISTING_PAGE = re.compile(r"/(?:search|s|c|category|categories|brand|brands|compare|offers?|deals?)(?:/|\?|$)", re.I)
@dataclass
class RunOptions:
category: str
brands: List[str] # brand slugs
max_products_per_brand: int = 15
expand_per_brand: int = 8 # models to look up on other platforms
search_budget: int = 200
use_llm: bool = True
fetch_pages: bool = True
find_images: bool = True
reprobe: bool = False
@dataclass
class RunStats:
counts: Dict[str, int] = field(default_factory=dict)
def inc(self, key: str, n: int = 1) -> None:
self.counts[key] = self.counts.get(key, 0) + n
def source_sku(site: SiteRef, url: str) -> str:
rx = site.product_url_re
if rx is not None:
m = rx.search(url)
if m and m.groups() and m.group(1):
return m.group(1)
p = urlparse(url)
return (p.netloc.lower().removeprefix("www.") + p.path.rstrip("/").lower())[:300]
_INDIA_PATH = re.compile(r"^/(?:in|in-en|en-in|en_in|in_en)(?:/|$)", re.I)
def is_product_url(site: SiteRef, url: str) -> bool:
p = urlparse(url)
if _LISTING_PAGE.search(p.path):
return False
if site.kind == "brand_official":
# Only the brand's Indian storefront: www.samsung.com/in/..., not
# us.samsung.com or news.samsung.com. Domains that are Indian already
# (oneplus.in, motorola.co.in) qualify as they are.
host = (p.hostname or "").lower()
if host not in (site.domain, "www." + site.domain, "in." + site.domain):
return False
if not site.domain.endswith((".in", ".co.in")) and not _INDIA_PATH.search(p.path) and not host.startswith("in."):
return False
rx = site.product_url_re
if rx is not None:
return bool(rx.search(url))
return p.path.count("/") >= 2 # brand sites: at least /section/product
def embed_verified_products() -> int:
"""MiniLM vectors for verified products that do not have one yet."""
rows = repo.products_without_embedding()
if not rows:
return 0
from app.services.embeddings_service import embed_texts
texts = [
f"{r['brand']} {r['display_name']} {r['category']} "
+ " ".join(f"{k} {v}" for k, v in (r["canonical_specs"] or {}).items())
for r in rows
]
for r, vec in zip(rows, embed_texts(texts)):
repo.set_embedding(r["id"], vec)
return len(rows)
class Collector:
def __init__(self, options: RunOptions, *, progress: Optional[Callable[[str], None]] = None) -> None:
self.opt = options
self.ref = load_reference()
self.stats = RunStats()
self.progress = progress or (lambda msg: logger.info(msg))
self.run_id: Optional[int] = None
self.ids = repo.id_maps()
self.site_rows = {r["domain"]: r for r in repo.sites()}
self.breaker = CircuitBreaker(on_trip=self._on_trip)
for r in self.site_rows.values():
if r.get("breaker_until"):
self.breaker.preload(r["domain"], r["breaker_until"].timestamp(), r.get("breaker_reason") or "")
self.client = PoliteClient(breaker=self.breaker, on_fetch=self._on_fetch)
self.engine = SearchEngine(budget=options.search_budget)
self._touched_products: Dict[int, List[Tuple[int, Listing]]] = {}
# -- callbacks -----------------------------------------------------------
def _on_trip(self, host: str, reason: str, until: float) -> None:
self.stats.inc("breaker_trips")
self.progress(f"Circuit breaker opened for {host}: {reason}. Falling back to web search for it.")
repo.trip_breaker(host, reason, until)
def _on_fetch(self, result, host: str) -> None:
self.stats.inc(f"fetch_{result.outcome}")
repo.log_fetch(self.run_id, result.url, host, result.status, result.bytes, result.outcome, result.robots_allowed)
# -- grading -------------------------------------------------------------
def _grade(self, site: SiteRef) -> str:
if site.policy == "serp_only":
return "C"
host = site.domain
if self.breaker.is_open(host) or self.breaker.is_open("www." + host):
return "C"
return (self.site_rows.get(site.domain) or {}).get("probe_outcome") or "C"
def ensure_probes(self, sites: List[SiteRef]) -> None:
from datetime import datetime, timedelta, timezone
stale_before = datetime.now(timezone.utc) - timedelta(days=ELEC_PROBE_TTL_DAYS)
for site in sites:
row = self.site_rows.get(site.domain) or {}
if not self.opt.reprobe and row.get("probed_at") and row["probed_at"] > stale_before:
continue
self.progress(f"Probing {site.name} ({site.domain})")
result = probe_site(site, self.client, self.engine)
repo.set_probe_result(site.domain, result["outcome"], result["robots_allowed"], result["evidence"])
row.update(probe_outcome=result["outcome"])
self.site_rows[site.domain] = row
self.stats.inc(f"probe_{result['outcome']}")
self.progress(f" -> grade {result['outcome']}: {result['evidence'].get('reason')}")
# -- discovery -----------------------------------------------------------
def _accept_hit(self, hit: SearchHit, brand: BrandRef) -> Optional[Tuple[SiteRef, ParsedTitle]]:
site = site_for_url(hit.url)
if site is None or not self._site_allowed_for(site, brand):
return None
if not is_product_url(site, hit.url):
return None
title = clean_result_title(hit.title)
if not looks_like_device_title(title) or _OFF_CATEGORY[self.opt.category].search(title):
return None
parsed = parse_title(title, self.opt.category, expected_brand=brand.slug)
if parsed.brand is None or not parsed.model_norm:
return None
fill_from_context(parsed, self.opt.category, snippet=hit.snippet or "")
return site, parsed
def _site_allowed_for(self, site: SiteRef, brand: BrandRef) -> bool:
if site.kind == "brand_official":
return site.brand_slug == brand.slug
return True
def _platforms_for(self, brand: BrandRef) -> List[SiteRef]:
return [s for s in self.ref.sites.values()
if s.kind != "brand_official" or s.brand_slug == brand.slug]
def _search_names(self, brand: BrandRef) -> List[str]:
"""The brand, plus the sub-brands phones are actually sold under
("Redmi", "POCO", "iQOO") - a search for "Xiaomi" alone misses most
Redmi listings."""
names = [brand.name]
if self.opt.category == "mobiles":
names += [s.upper() if len(s) <= 4 else s.title() for s in brand.sub_brands
if s not in ("mi", "iphone", "pixel", "narzo")][:2]
return names
def _collect_hits(self, hits: Optional[List[SearchHit]], brand: BrandRef, query: str,
found: Dict, models: Dict[str, ParsedTitle]) -> int:
new = 0
for hit in hits or []:
accepted = self._accept_hit(hit, brand)
if not accepted:
continue
site_ref, parsed = accepted
key = (site_ref.domain, source_sku(site_ref, hit.url))
if key not in found:
found[key] = (hit, site_ref, parsed, query)
new += 1
models.setdefault(parsed.model_norm, parsed)
return new
def discover(self, brand: BrandRef) -> Dict[Tuple[str, str], Tuple[SearchHit, SiteRef, ParsedTitle, str]]:
category = self.ref.categories[self.opt.category]
found: Dict[Tuple[str, str], Tuple[SearchHit, SiteRef, ParsedTitle, str]] = {}
models: Dict[str, ParsedTitle] = {}
per_site_target = max(4, self.opt.max_products_per_brand)
for site in self._platforms_for(brand):
got = 0
for name in self._search_names(brand):
for terms in category.query_terms or category.search_terms:
query = f"site:{site.domain} {name} {terms}"
hits = self.engine.text(query, max_results=20)
if hits is None:
self.stats.inc("search_unavailable")
continue
got += self._collect_hits(hits, brand, query, found, models)
if got >= per_site_target:
break
if got >= per_site_target:
break
self.progress(f"{brand.name}: {len(found)} listing URLs, {len(models)} models from platform searches")
# Cross-platform: look each variant up by name, to find the same product
# on platforms the site: searches missed. Phones are grouped by model
# line; laptops by full configuration (line + CPU + RAM + storage),
# because one laptop line is sold in dozens of configurations and only
# the exact one confirms a product.
for p in self._expansion_targets(found):
query = f"{brand.name} {p.model or p.model_norm} {self._variant_terms(p)} price"
self._collect_hits(self.engine.text(re.sub(r"\s+", " ", query), max_results=20), brand, query, found, models)
self.progress(f"{brand.name}: {len(found)} listing URLs after cross-platform search")
return found
def _expansion_targets(self, found: Dict) -> List[ParsedTitle]:
"""Which variants to look up on other platforms, most useful first:
variants seen on the most sites, then ones whose page we can read with
a price (grade A/B platforms), since one more site verifies those."""
groups: Dict[str, Dict] = {}
for (domain, _), (_, site, parsed, _) in found.items():
if self.opt.category == "laptops":
key = variant_key(parsed, "laptops")
if key is None:
continue
else:
key = parsed.model_norm
g = groups.setdefault(key, {"parsed": parsed, "sites": set(), "readable": False})
g["sites"].add(domain)
g["readable"] = g["readable"] or self._grade(site) in ("A", "B")
ranked = sorted(groups.values(), key=lambda g: (-len(g["sites"]), not g["readable"]))
return [g["parsed"] for g in ranked[: self.opt.expand_per_brand]]
@staticmethod
def _variant_terms(p: ParsedTitle) -> str:
parts = []
if p.processor:
# "ryzen 5 7530u" -> "Ryzen 5 7530U", "i5-1334u" -> "i5-1334U"
parts.append(" ".join(t.upper() if any(ch.isdigit() for ch in t) and len(t) > 2 else t.title()
for t in p.processor.split()))
if p.ram_gb:
parts.append(f"{format(p.ram_gb.normalize(), 'f')}GB RAM")
if p.storage_gb:
parts.append(f"{format(p.storage_gb.normalize(), 'f')}GB")
return " ".join(parts)
# -- listing construction --------------------------------------------------
def _base_listing(self, site: SiteRef, url: str, parsed: ParsedTitle, title: str, source_type: str,
evidence: str, confidence: float, parser: str, query: str) -> Listing:
l = Listing(
site_domain=site.domain, source_sku=source_sku(site, url), source_url=url, source_type=source_type,
brand_slug=parsed.brand.brand_slug, category=self.opt.category, title=title,
evidence_text=evidence, confidence=confidence, parser=parser, family=parsed.brand.family,
model=parsed.model, model_number=parsed.mpn, ram_gb=parsed.ram_gb, storage_gb=parsed.storage_gb,
colour=parsed.colour, search_query=query,
)
l.model_norm = parsed.model_norm
l.processor = parsed.processor
l.variant_key = variant_key(parsed, self.opt.category)
return l
def listing_from_search(self, hit: SearchHit, site: SiteRef, parsed: ParsedTitle, query: str) -> Listing:
title = clean_result_title(hit.title)
evidence = f"{hit.title} — {hit.snippet}".strip(" —")
reading = read_price(f"{hit.title} {hit.snippet}")
price, mrp = reading.price, reading.mrp
# A snippet that names a different RAM/storage than the title is about
# another variant; its price cannot be trusted for this one.
snippet_variant = parse_title(hit.snippet or "", self.opt.category)
for a, b in ((snippet_variant.storage_gb, parsed.storage_gb), (snippet_variant.ram_gb, parsed.ram_gb)):
if a is not None and b is not None and a != b:
price = mrp = None
self.stats.inc("snippet_price_variant_conflict")
# Truncated titles ("... - (16 GB ...") lose the variant; the snippet
# of the same result usually states it.
fill_from_context(parsed, self.opt.category, snippet=hit.snippet or "")
parser, confidence = f"serp:{hit.provider}", (0.55 if price is not None else 0.45)
in_stock = read_stock(hit.snippet or "")
availability = None if in_stock is None else ("InStock" if in_stock else "OutOfStock")
# Structured offer the search engine read from the page itself (Google
# pagemap). Better than snippet text, and still no request to the site.
if hit.offer and price is None:
try:
offered = Decimal(str(hit.offer["price"]).replace(",", ""))
except Exception: # noqa: BLE001
offered = None
if offered is not None and Decimal(500) <= offered <= Decimal(1000000):
price, mrp = offered, None
evidence = f"{evidence} || search-engine offer data: {json.dumps(hit.offer['raw'], default=str)[:600]}"
parser, confidence = f"serp:{hit.provider}:pagemap", 0.65
av = (hit.offer.get("availability") or "").lower().replace(" ", "")
if "instock" in av:
in_stock, availability = True, "InStock"
elif "outofstock" in av or "soldout" in av:
in_stock, availability = False, "OutOfStock"
# Rating: the engine's structured data first (read from the page
# itself), else an explicit "x out of 5" in this result's own text.
rating, review_count = None, None
if hit.rating and hit.rating.get("rating") is not None:
rating = Decimal(str(hit.rating["rating"]))
review_count = hit.rating.get("review_count")
evidence = f"{evidence} || search-engine rating data: {json.dumps(hit.rating.get('raw'), default=str)[:300]}"
else:
stated = read_rating(f"{hit.title} {hit.snippet}")
if stated.rating is not None:
rating, review_count = stated.rating, stated.review_count
listing = self._base_listing(site, hit.url, parsed, title, "search_snippet", evidence,
confidence, parser, query)
listing.price, listing.mrp = price, mrp
listing.in_stock, listing.availability = in_stock, availability
listing.rating, listing.review_count = rating, review_count
return listing
def listing_from_page(self, hit: SearchHit, site: SiteRef, parsed_hit: ParsedTitle, query: str,
html: str, final_url: str) -> Optional[Listing]:
products = extract_products(html)
page = extract_page(html)
if site.kind == "brand_official":
# A brand's own store names products without the brand ("Galaxy A56 5G
# (8 GB Memory)"); on its own site the brand is not in doubt.
brand_name = self.ref.brands[parsed_hit.brand.brand_slug].name
for p in products:
if not p["name"].lower().startswith(brand_name.lower()):
p["name"] = f"{brand_name} {p['name']}"
name = None
product = None
for p in products:
pp = parse_title(p["name"], self.opt.category, expected_brand=parsed_hit.brand.brand_slug)
if pp.brand and pp.model_norm and fuzz.token_set_ratio(pp.model_norm, parsed_hit.model_norm) >= 85:
product, name = p, p["name"]
break
if name is None:
name = page.get("name")
if not name:
return None
parsed = parse_title(name, self.opt.category, expected_brand=parsed_hit.brand.brand_slug)
if parsed.brand is None or not parsed.model_norm:
return None
if fuzz.token_set_ratio(parsed.model_norm, parsed_hit.model_norm) < 85:
# The URL did not lead to the product the search result named.
self.stats.inc("page_title_mismatch")
return None
# Fill variant fields the page name leaves out from the search title
# of the same URL (both are statements by the same site).
for attr in ("ram_gb", "storage_gb", "colour", "mpn", "processor"):
if getattr(parsed, attr) is None and getattr(parsed_hit, attr) is not None:
setattr(parsed, attr, getattr(parsed_hit, attr))
grade = self._grade(site)
source_type = "brand_official" if site.kind == "brand_official" else "scraped_page"
if product is not None:
price = product.get("price")
currency = product.get("currency")
evidence = product["evidence"]
parser = "jsonld"
else:
price, currency, evidence, parser = page.get("price"), page.get("currency"), page.get("evidence") or "", "html_meta"
if currency not in (None, "INR"):
price = None
if currency is None and site.kind == "brand_official":
price = None # a brand's global site may not be quoting rupees
if price is not None and not (Decimal(500) <= price <= Decimal(1000000)):
price = None
evidence = evidence or f"{name} ({final_url})"
listing = self._base_listing(
site, final_url if final_url.startswith("http") else hit.url, parsed, name, source_type,
evidence, 0.9 if (price is not None and parser == "jsonld") else 0.75, f"{parser}:grade{grade}", query,
)
# The site's own SKU when the page states it; otherwise its canonical URL.
listing.source_sku = ((product or {}).get("sku") or source_sku(site, final_url or hit.url))[:300]
listing.price = price
listing.in_stock = (product or page).get("in_stock")
listing.availability = (product or page).get("availability")
listing.image_urls = list(dict.fromkeys((product or {}).get("images", []) + page.get("images", [])))[:8]
if product:
listing.gtin = product.get("gtin")
listing.model_number = listing.model_number or product.get("mpn")
listing.rating = product.get("rating")
listing.review_count = product.get("review_count")
listing.reviews = list(product.get("reviews") or [])
listing.colour = listing.colour or product.get("color")
# Retailers that embed their rating outside JSON-LD (page's own product
# only): the rating when JSON-LD has none, and the star breakdown.
embedded = embedded_rating(site.domain, html, final_url or hit.url)
if embedded:
if listing.rating is None:
listing.rating, listing.review_count = embedded["rating"], embedded["review_count"]
listing.rating_breakdown = embedded["breakdown"]
if site.domain in IGNORED_DOMAINS:
# Known placeholder rating markup: nothing rating-shaped is kept.
listing.rating = listing.review_count = listing.rating_breakdown = None
listing.reviews = []
raw_specs = dict((product or {}).get("properties") or {})
raw_specs.update({k: v for k, v in spec_tables(BeautifulSoup(html, "lxml")).items() if k not in raw_specs})
listing.specs_raw = dict(list(raw_specs.items())[:150])
listing.specs, listing.spec_sources = normalise_specs(self.opt.category, raw_specs)
if self.opt.use_llm:
wanted = [k for k in self.ref.spec_keys.get(self.opt.category, {}) if k not in listing.specs]
if wanted:
text = "\n".join(f"{k}: {v}" for k, v in raw_specs.items()) or visible_text(html, 3500)
extra, extra_src = fill_missing(self.opt.category, text, wanted)
listing.specs.update(extra)
listing.spec_sources.update(extra_src)
self.stats.inc("llm_specs_kept", len(extra))
if self.opt.category == "laptops":
# "13th Gen Intel Core i7/ 16GB RAM" in a title names no CPU model;
# the page's own spec table usually does.
spec_texts = tuple(str(v) for k, v in raw_specs.items() if "processor" in k.lower() or "cpu" in k.lower())
spec_texts += (str(listing.specs.get("processor") or ""),)
before = parsed.processor
fill_from_context(parsed, self.opt.category, spec_texts=spec_texts)
if parsed.processor != before:
listing.processor = parsed.processor
listing.variant_key = variant_key(parsed, self.opt.category)
listing.content_hash = hashlib.sha1(html.encode("utf-8", "ignore")).hexdigest()
return listing
# -- persistence -----------------------------------------------------------
def store(self, listing: Listing) -> Optional[int]:
try:
listing_id = repo.upsert_listing(listing, self.ids, self.run_id)
except ValueError as exc:
self.stats.inc("rejected_listing")
logger.info("Listing rejected (%s): %s", exc, listing.source_url)
return None
self.stats.inc(f"listing_{listing.source_type}")
if listing.reviews:
self.stats.inc("reviews_stored", repo.save_reviews(listing_id, listing.reviews))
if listing.price is not None:
self.stats.inc("listing_with_price")
decision = decide(listing, repo.product_candidates(listing.brand_slug, listing.category))
if decision is None:
self.stats.inc("listing_unmatched_no_variant")
return listing_id
product_id = decision.product_id or repo.create_product(listing, self.ids)
if decision.product_id is None:
self.stats.inc("product_created")
repo.map_listing(listing_id, product_id, decision.method, decision.confidence, decision.review_status)
if decision.review_status == "pending":
self.stats.inc("match_pending_review")
else:
repo.merge_product_specs(product_id, listing.specs, listing.spec_sources, listing.source_url)
self._touched_products.setdefault(product_id, []).append((listing_id, listing))
return listing_id
# -- images ----------------------------------------------------------------
def attach_images(self) -> None:
rank_for = {"brand_official": 10, "scraped_page": 20, "search_snippet": 50}
for product_id, entries in self._touched_products.items():
if repo.product_image_count(product_id) >= 3:
continue
added = 0
for listing_id, listing in sorted(entries, key=lambda e: rank_for[e[1].source_type]):
for url in listing.image_urls:
if added >= 3:
break
if self.client.check_image(url, MIN_IMAGE_BYTES):
repo.add_image(product_id, url, listing_id, listing.source_type, rank_for[listing.source_type])
added += 1
self.stats.inc("images_from_pages")
if added or not self.opt.find_images:
continue
added = self._images_from_search(product_id, entries)
self.stats.inc("images_from_search", added)
def _images_from_search(self, product_id: int, entries: List[Tuple[int, Listing]]) -> int:
"""Image search results are used only when the page an image sits on is
one of THIS product's own listings (same site, same product), and the
image result's title names the model."""
_, listing = entries[0]
brand = self.ref.brands[listing.brand_slug]
listing_by_site = {l.site_domain: lid for lid, l in entries}
hits = self.engine.images(f"{brand.name} {listing.model or listing.model_norm}", max_results=15) or []
added = 0
for hit in hits:
site = site_for_url(hit.url)
if site is None or site.domain not in listing_by_site:
continue
parsed = parse_title(clean_result_title(hit.title), listing.category, expected_brand=brand.slug)
if not parsed.model_norm or fuzz.token_set_ratio(parsed.model_norm, listing.model_norm or "") < 90:
continue
if self.client.check_image(hit.image_url, MIN_IMAGE_BYTES):
repo.add_image(product_id, hit.image_url, listing_by_site[site.domain], "search_image", 60)
added += 1
if added >= 2:
break
return added
# -- embeddings ------------------------------------------------------------
def embed(self) -> int:
return embed_verified_products()
# -- the run ---------------------------------------------------------------
def run(self, *, embed: bool = True) -> Dict[str, int]:
self.run_id = repo.start_run("collect", {
"category": self.opt.category, "brands": self.opt.brands,
"max_products_per_brand": self.opt.max_products_per_brand,
"search_budget": self.opt.search_budget,
})
status, error = "done", None
try:
brands = [self.ref.brands[b] for b in self.opt.brands]
probe_targets = {s.domain: s for b in brands for s in self._platforms_for(b) if s.policy == "probe"}
if self.opt.fetch_pages:
self.ensure_probes(list(probe_targets.values()))
for brand in brands:
found = self.discover(brand)
# Keep the most common models first, up to the per-brand limit.
by_model: Dict[str, int] = {}
for (_, _), (_, _, parsed, _) in found.items():
by_model[parsed.model_norm] = by_model.get(parsed.model_norm, 0) + 1
keep = set(sorted(by_model, key=lambda m: -by_model[m])[: self.opt.max_products_per_brand])
for (domain, _sku), (hit, site, parsed, query) in found.items():
if parsed.model_norm not in keep:
continue
listing = None
if self.opt.fetch_pages and self._grade(site) in ("A", "B"):
res = self.client.get(hit.url)
if res.ok:
listing = self.listing_from_page(hit, site, parsed, query, res.text, res.final_url)
if listing is None:
self.stats.inc("page_unusable_fell_back_to_search")
if listing is None:
listing = self.listing_from_search(hit, site, parsed, query)
self.store(listing)
self.progress(f"{brand.name}: stored listings; {self.stats.counts}")
self.attach_images()
verification = repo.refresh_verification()
self.stats.counts.update({f"products_{k}": v for k, v in verification.items()})
if embed:
try:
self.stats.inc("embedded", self.embed())
except Exception as exc: # noqa: BLE001 - embeddings are optional
logger.warning("Embedding step skipped: %s", exc)
self.stats.counts.update({f"search_{k}": v for k, v in self.engine.stats.items()})
except Exception as exc:
status, error = "failed", repr(exc)
logger.exception("Collection run failed")
raise
finally:
repo.finish_run(self.run_id, status, self.stats.counts, error)
self.client.close()
return self.stats.counts