"""Search-first collection of real product listings. For one category and a set of allow-listed brands: 1. DISCOVER web search `site: ` 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 ` 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