601 lines
31 KiB
Python
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
|