Electronics Catalog: API, MCP server, frontend and deployment

Verified catalogue of mobiles and laptops sold in India, collected from
real retail listings (FastAPI backend, React frontend, Postgres/pgvector).

- REST API under /api/elec (read-only catalogue; admin endpoints need login)
- MCP server (FastMCP) at /mcp/ with list_categories, search_products,
  get_product and price_history tools
- Real ratings and reviews read from product pages and search results
- Production Dockerfile (requirements-api.txt, no PyTorch) and
  .env.production.example; remote database only via an explicit
  ELEC_ALLOW_REMOTE_DB host/name allowlist
- docs/API.md: endpoint and MCP reference with live examples

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This commit is contained in:
sriram
2026-10-01 12:17:42 +05:30
commit c7e4d59188
115 changed files with 14329 additions and 0 deletions

View File

View File

@@ -0,0 +1,68 @@
"""Connections to the local electronics database.
settings.py has already refused to load unless DB_HOST is local and DB_NAME is
electronics_catalog, so nothing here can reach another database.
"""
from __future__ import annotations
import logging
from contextlib import contextmanager
from typing import Iterator
import psycopg
from psycopg.rows import dict_row
from app.infrastructure.settings import (
DB_CONNECT_TIMEOUT_SECONDS,
DB_HOST,
DB_NAME,
DB_PASSWORD,
DB_PORT,
DB_USER,
)
logger = logging.getLogger(__name__)
def connect(*, autocommit: bool = False) -> psycopg.Connection:
conn = psycopg.connect(
host=DB_HOST,
port=DB_PORT,
dbname=DB_NAME,
user=DB_USER,
password=DB_PASSWORD,
connect_timeout=DB_CONNECT_TIMEOUT_SECONDS,
autocommit=autocommit,
row_factory=dict_row,
)
try:
from pgvector.psycopg import register_vector
register_vector(conn)
except Exception: # noqa: BLE001 - the extension is created by migration 0001
pass
return conn
@contextmanager
def transaction() -> Iterator[psycopg.Connection]:
"""A connection whose work is committed on success, rolled back on error."""
conn = connect()
try:
yield conn
conn.commit()
except Exception:
conn.rollback()
raise
finally:
conn.close()
def check_connection() -> bool:
try:
with connect(autocommit=True) as conn:
conn.execute("SELECT 1")
return True
except Exception as exc: # noqa: BLE001 - a health probe reports, never raises
logger.debug("Database unreachable: %s", exc)
return False

View File

@@ -0,0 +1,63 @@
"""Tiny migration runner for the numbered SQL files in ./migrations.
Each file runs once, in its own transaction, and is recorded in
elec.schema_migrations with a checksum. Editing an applied file is refused
rather than silently ignored - add a new numbered file instead.
"""
from __future__ import annotations
import hashlib
import logging
from pathlib import Path
from typing import List
from app.electronics.db.connection import connect
logger = logging.getLogger(__name__)
MIGRATIONS_DIR = Path(__file__).resolve().parent / "migrations"
_BOOTSTRAP = """
CREATE SCHEMA IF NOT EXISTS elec;
CREATE TABLE IF NOT EXISTS elec.schema_migrations (
version TEXT PRIMARY KEY,
checksum TEXT NOT NULL,
applied_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
"""
def _files() -> List[Path]:
return sorted(MIGRATIONS_DIR.glob("[0-9][0-9][0-9][0-9]_*.sql"))
def run_migrations() -> List[str]:
"""Apply pending migrations. Returns the versions applied by this call."""
applied_now: List[str] = []
with connect() as conn:
conn.execute(_BOOTSTRAP)
conn.commit()
done = {
r["version"]: r["checksum"]
for r in conn.execute("SELECT version, checksum FROM elec.schema_migrations")
}
for path in _files():
sql = path.read_text(encoding="utf-8")
checksum = hashlib.sha256(sql.encode("utf-8")).hexdigest()
version = path.stem
if version in done:
if done[version] != checksum:
raise RuntimeError(
f"Migration {version} was edited after it was applied. "
f"Revert the edit and add a new numbered migration instead."
)
continue
logger.info("Applying migration %s", version)
with conn.transaction():
conn.execute(sql)
conn.execute(
"INSERT INTO elec.schema_migrations (version, checksum) VALUES (%s, %s)",
(version, checksum),
)
applied_now.append(version)
return applied_now

View File

@@ -0,0 +1,4 @@
-- Extensions and the dedicated schema. Everything this project owns lives in
-- schema `elec` of database `electronics_catalog`.
CREATE EXTENSION IF NOT EXISTS vector;
CREATE SCHEMA IF NOT EXISTS elec;

View File

@@ -0,0 +1,57 @@
-- Reference data: brands, categories, retail sites. Seeded from
-- app/electronics/reference/*.yaml by `elec seed-reference`.
CREATE TABLE elec.brand (
id SERIAL PRIMARY KEY,
name TEXT NOT NULL UNIQUE,
slug TEXT NOT NULL UNIQUE,
parent_brand_id INT REFERENCES elec.brand(id),
is_popular BOOLEAN NOT NULL DEFAULT TRUE,
official_domains TEXT[] NOT NULL DEFAULT '{}',
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
-- Every spelling that resolves to a brand. Sub-brands (Redmi, iQOO, Pixel)
-- resolve to their parent and are remembered as the product family.
CREATE TABLE elec.brand_alias (
alias TEXT PRIMARY KEY CHECK (alias = lower(alias)),
brand_id INT NOT NULL REFERENCES elec.brand(id) ON DELETE CASCADE,
is_sub_brand BOOLEAN NOT NULL DEFAULT FALSE
);
CREATE TABLE elec.category (
id SERIAL PRIMARY KEY,
slug TEXT NOT NULL UNIQUE,
name TEXT NOT NULL UNIQUE
);
CREATE TABLE elec.brand_category (
brand_id INT NOT NULL REFERENCES elec.brand(id) ON DELETE CASCADE,
category_id INT NOT NULL REFERENCES elec.category(id) ON DELETE CASCADE,
PRIMARY KEY (brand_id, category_id)
);
-- A retail platform or a brand's own site, with the outcome of its probe.
-- probe_outcome A = fetchable with structured product data (scraped)
-- B = fetchable, product data from page HTML/state (scraped)
-- C = not fetched: serp_only policy, robots.txt disallow,
-- block/CAPTCHA, or unreachable -> web search only
CREATE TABLE elec.site (
id SERIAL PRIMARY KEY,
domain TEXT NOT NULL UNIQUE,
name TEXT NOT NULL,
kind TEXT NOT NULL CHECK (kind IN ('marketplace','national_chain','tn_regional','brand_official')),
region TEXT NOT NULL CHECK (region IN ('national','TN')),
policy TEXT NOT NULL CHECK (policy IN ('probe','serp_only')),
brand_id INT REFERENCES elec.brand(id),
product_url TEXT,
pincode_param TEXT,
enabled BOOLEAN NOT NULL DEFAULT TRUE,
probe_outcome CHAR(1) CHECK (probe_outcome IN ('A','B','C')),
robots_allowed BOOLEAN,
probe_evidence JSONB NOT NULL DEFAULT '{}'::jsonb,
probed_at TIMESTAMPTZ,
breaker_until TIMESTAMPTZ,
breaker_reason TEXT,
CHECK (kind <> 'brand_official' OR brand_id IS NOT NULL)
);

View File

@@ -0,0 +1,160 @@
-- Runs, fetch audit trail, search cache, listings, prices, canonical products.
CREATE TABLE elec.crawl_run (
id BIGSERIAL PRIMARY KEY,
kind TEXT NOT NULL,
params JSONB NOT NULL DEFAULT '{}'::jsonb,
status TEXT NOT NULL DEFAULT 'running' CHECK (status IN ('running','done','failed')),
stats JSONB NOT NULL DEFAULT '{}'::jsonb,
error TEXT,
started_at TIMESTAMPTZ NOT NULL DEFAULT now(),
ended_at TIMESTAMPTZ
);
-- Every HTTP request made to a retail or brand site. Evidence that the
-- crawler obeyed robots.txt and its rate limits.
CREATE TABLE elec.fetch_log (
id BIGSERIAL PRIMARY KEY,
crawl_run_id BIGINT REFERENCES elec.crawl_run(id) ON DELETE SET NULL,
url TEXT NOT NULL,
host TEXT NOT NULL,
status INT,
bytes INT,
outcome TEXT NOT NULL,
robots_allowed BOOLEAN,
fetched_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
CREATE INDEX fetch_log_host_time ON elec.fetch_log (host, fetched_at DESC);
CREATE TABLE elec.search_cache (
provider TEXT NOT NULL,
kind TEXT NOT NULL CHECK (kind IN ('text','images')),
query TEXT NOT NULL,
results JSONB NOT NULL,
fetched_at TIMESTAMPTZ NOT NULL DEFAULT now(),
PRIMARY KEY (provider, kind, query)
);
-- One canonical product = one real-world variant (model + RAM + storage).
-- verification_status becomes 'verified' only when the product has a brand
-- official page, or listings on at least two different sites.
CREATE TABLE elec.product (
id BIGSERIAL PRIMARY KEY,
brand_id INT NOT NULL REFERENCES elec.brand(id),
category_id INT NOT NULL REFERENCES elec.category(id),
family TEXT,
model TEXT NOT NULL,
model_norm TEXT NOT NULL,
variant_key TEXT NOT NULL UNIQUE,
display_name TEXT NOT NULL,
ram_gb NUMERIC(6,1),
storage_gb NUMERIC(7,1),
processor TEXT,
mpn TEXT,
gtin TEXT,
canonical_specs JSONB NOT NULL DEFAULT '{}'::jsonb,
spec_sources JSONB NOT NULL DEFAULT '{}'::jsonb,
verification_status TEXT NOT NULL DEFAULT 'unverified'
CHECK (verification_status IN ('verified','unverified','rejected')),
evidence_count INT NOT NULL DEFAULT 0,
embedding vector(384),
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
CREATE INDEX product_brand_cat ON elec.product (brand_id, category_id);
CREATE INDEX product_specs_gin ON elec.product USING GIN (canonical_specs);
CREATE INDEX product_embedding_hnsw ON elec.product USING hnsw (embedding vector_cosine_ops);
-- The latest state of one product page on one site, or of one search result
-- that points at such a page. Nothing is stored without the URL it came from
-- and the text that the values were read from.
CREATE TABLE elec.source_listing (
id BIGSERIAL PRIMARY KEY,
site_id INT NOT NULL REFERENCES elec.site(id),
source_sku TEXT NOT NULL,
source_url TEXT NOT NULL CHECK (source_url ~ '^https?://'),
source_type TEXT NOT NULL CHECK (source_type IN ('scraped_page','search_snippet','brand_official')),
brand_id INT NOT NULL REFERENCES elec.brand(id),
category_id INT NOT NULL REFERENCES elec.category(id),
family TEXT,
title TEXT NOT NULL,
model TEXT,
model_number TEXT,
ram_gb NUMERIC(6,1),
storage_gb NUMERIC(7,1),
colour TEXT,
price NUMERIC(12,2) CHECK (price IS NULL OR price BETWEEN 500 AND 1000000),
mrp NUMERIC(12,2) CHECK (mrp IS NULL OR mrp BETWEEN 500 AND 1000000),
currency TEXT NOT NULL DEFAULT 'INR' CHECK (currency = 'INR'),
availability TEXT,
in_stock BOOLEAN,
pincode TEXT,
pincode_applied BOOLEAN NOT NULL DEFAULT FALSE,
rating NUMERIC(3,2) CHECK (rating IS NULL OR rating BETWEEN 0 AND 5),
review_count INT,
gtin TEXT,
image_urls TEXT[] NOT NULL DEFAULT '{}',
specs_raw JSONB NOT NULL DEFAULT '{}'::jsonb,
specs JSONB NOT NULL DEFAULT '{}'::jsonb,
evidence_text TEXT NOT NULL CHECK (length(evidence_text) > 0),
search_query TEXT,
confidence NUMERIC(3,2) NOT NULL CHECK (confidence BETWEEN 0 AND 1),
parser TEXT NOT NULL,
content_hash TEXT,
first_seen_at TIMESTAMPTZ NOT NULL DEFAULT now(),
last_seen_at TIMESTAMPTZ NOT NULL DEFAULT now(),
crawl_run_id BIGINT REFERENCES elec.crawl_run(id) ON DELETE SET NULL,
UNIQUE (site_id, source_sku),
CHECK (pincode_applied = FALSE OR pincode IS NOT NULL)
);
CREATE INDEX listing_brand_cat ON elec.source_listing (brand_id, category_id);
-- Append-only price observations. UPDATE is refused by a trigger.
CREATE TABLE elec.price_history (
id BIGSERIAL PRIMARY KEY,
listing_id BIGINT NOT NULL REFERENCES elec.source_listing(id) ON DELETE CASCADE,
price NUMERIC(12,2) CHECK (price IS NULL OR price BETWEEN 500 AND 1000000),
mrp NUMERIC(12,2),
availability TEXT,
in_stock BOOLEAN,
source_type TEXT NOT NULL,
pincode TEXT,
pincode_applied BOOLEAN NOT NULL DEFAULT FALSE,
evidence_text TEXT NOT NULL CHECK (length(evidence_text) > 0),
observed_at TIMESTAMPTZ NOT NULL DEFAULT now(),
crawl_run_id BIGINT REFERENCES elec.crawl_run(id) ON DELETE SET NULL
);
CREATE INDEX price_history_listing_time ON elec.price_history (listing_id, observed_at DESC);
CREATE FUNCTION elec.refuse_update() RETURNS trigger LANGUAGE plpgsql AS $$
BEGIN
RAISE EXCEPTION 'elec.price_history is append-only';
END $$;
CREATE TRIGGER price_history_append_only BEFORE UPDATE ON elec.price_history
FOR EACH ROW EXECUTE FUNCTION elec.refuse_update();
-- Which canonical product a listing belongs to, and how sure we are.
-- Only 'auto' and 'approved' links count as evidence or appear in views.
CREATE TABLE elec.product_listing_map (
listing_id BIGINT PRIMARY KEY REFERENCES elec.source_listing(id) ON DELETE CASCADE,
product_id BIGINT NOT NULL REFERENCES elec.product(id) ON DELETE CASCADE,
method TEXT NOT NULL CHECK (method IN ('gtin','mpn','variant_key','fuzzy','manual')),
confidence NUMERIC(3,2) NOT NULL CHECK (confidence BETWEEN 0 AND 1),
review_status TEXT NOT NULL CHECK (review_status IN ('auto','pending','approved','rejected')),
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
reviewed_at TIMESTAMPTZ
);
CREATE INDEX map_product ON elec.product_listing_map (product_id);
-- Images are URLs only (never downloaded), each tied to the listing it was
-- found on and checked live.
CREATE TABLE elec.product_image (
id BIGSERIAL PRIMARY KEY,
product_id BIGINT NOT NULL REFERENCES elec.product(id) ON DELETE CASCADE,
url TEXT NOT NULL CHECK (url ~ '^https?://'),
source_listing_id BIGINT NOT NULL REFERENCES elec.source_listing(id) ON DELETE CASCADE,
source_type TEXT NOT NULL,
rank INT NOT NULL DEFAULT 100,
validated_at TIMESTAMPTZ NOT NULL DEFAULT now(),
UNIQUE (product_id, url)
);

View File

@@ -0,0 +1,71 @@
-- Read-side views. Public views only ever show VERIFIED products and links
-- that are 'auto' or 'approved'.
CREATE VIEW elec.v_product_availability AS
SELECT p.id AS product_id,
b.name AS brand,
c.slug AS category,
p.display_name,
s.id AS site_id,
s.name AS site,
s.domain,
s.kind AS site_kind,
s.region AS site_region,
l.id AS listing_id,
l.source_url,
l.source_type,
l.title AS listing_title,
l.colour,
l.price,
l.mrp,
l.in_stock,
l.availability,
l.pincode,
l.pincode_applied,
l.confidence,
l.last_seen_at AS observed_at
FROM elec.product p
JOIN elec.brand b ON b.id = p.brand_id
JOIN elec.category c ON c.id = p.category_id
JOIN elec.product_listing_map m ON m.product_id = p.id AND m.review_status IN ('auto','approved')
JOIN elec.source_listing l ON l.id = m.listing_id
JOIN elec.site s ON s.id = l.site_id
WHERE p.verification_status = 'verified';
-- Cheapest known price per product. Scraped prices are preferred over search
-- snippet prices; a listing known to be out of stock is skipped.
CREATE VIEW elec.v_best_price AS
SELECT DISTINCT ON (product_id)
product_id, site, domain, source_url, source_type, price, mrp, in_stock, observed_at
FROM elec.v_product_availability
WHERE price IS NOT NULL AND in_stock IS DISTINCT FROM FALSE
ORDER BY product_id, (source_type = 'search_snippet'), price, observed_at DESC;
CREATE VIEW elec.v_brand_catalog AS
SELECT p.id AS product_id, b.name AS brand, b.slug AS brand_slug, c.slug AS category,
p.family, p.display_name, p.model, p.ram_gb, p.storage_gb, p.processor,
p.canonical_specs,
bp.price AS best_price,
bp.site AS best_price_site,
bp.source_type AS best_price_source_type,
(SELECT count(DISTINCT a.site_id) FROM elec.v_product_availability a
WHERE a.product_id = p.id) AS platform_count,
(SELECT coalesce(bool_or(a.site_region = 'TN'), FALSE) FROM elec.v_product_availability a
WHERE a.product_id = p.id) AS sold_by_tn_retailer,
(SELECT i.url FROM elec.product_image i WHERE i.product_id = p.id
ORDER BY i.rank, i.id LIMIT 1) AS image_url,
p.updated_at
FROM elec.product p
JOIN elec.brand b ON b.id = p.brand_id
JOIN elec.category c ON c.id = p.category_id
LEFT JOIN elec.v_best_price bp ON bp.product_id = p.id
WHERE p.verification_status = 'verified';
CREATE VIEW elec.v_brand_summary AS
SELECT brand, brand_slug, category,
count(*) AS product_count,
min(best_price) AS min_price,
max(best_price) AS max_price,
max(platform_count) AS max_platforms
FROM elec.v_brand_catalog
GROUP BY brand, brand_slug, category;

View File

@@ -0,0 +1,44 @@
-- Search results carry cached, sometimes seller-specific prices. A price that
-- disagrees sharply with the product-page price for the same product (or is
-- below what the category can cost) is kept with its evidence but flagged, and
-- is never used as the "best price". Set by repository.flag_price_outliers().
ALTER TABLE elec.source_listing ADD COLUMN price_outlier BOOLEAN NOT NULL DEFAULT FALSE;
CREATE OR REPLACE VIEW elec.v_product_availability AS
SELECT p.id AS product_id,
b.name AS brand,
c.slug AS category,
p.display_name,
s.id AS site_id,
s.name AS site,
s.domain,
s.kind AS site_kind,
s.region AS site_region,
l.id AS listing_id,
l.source_url,
l.source_type,
l.title AS listing_title,
l.colour,
l.price,
l.mrp,
l.in_stock,
l.availability,
l.pincode,
l.pincode_applied,
l.confidence,
l.last_seen_at AS observed_at,
l.price_outlier
FROM elec.product p
JOIN elec.brand b ON b.id = p.brand_id
JOIN elec.category c ON c.id = p.category_id
JOIN elec.product_listing_map m ON m.product_id = p.id AND m.review_status IN ('auto','approved')
JOIN elec.source_listing l ON l.id = m.listing_id
JOIN elec.site s ON s.id = l.site_id
WHERE p.verification_status = 'verified';
CREATE OR REPLACE VIEW elec.v_best_price AS
SELECT DISTINCT ON (product_id)
product_id, site, domain, source_url, source_type, price, mrp, in_stock, observed_at
FROM elec.v_product_availability
WHERE price IS NOT NULL AND NOT price_outlier AND in_stock IS DISTINCT FROM FALSE
ORDER BY product_id, (source_type = 'search_snippet'), price, observed_at DESC;

View File

@@ -0,0 +1,30 @@
-- Best price: a product whose every priced listing is out of stock still has a
-- price worth showing. In-stock (or unknown-stock) prices still win; an
-- out-of-stock price is used only when nothing else is priced. Same columns
-- as 0005, so v_brand_catalog keeps working unchanged.
CREATE OR REPLACE VIEW elec.v_best_price AS
SELECT DISTINCT ON (product_id)
product_id, site, domain, source_url, source_type, price, mrp, in_stock, observed_at
FROM elec.v_product_availability
WHERE price IS NOT NULL AND NOT price_outlier
ORDER BY product_id, (in_stock IS FALSE), (source_type = 'search_snippet'), price, observed_at DESC;
-- Individual customer reviews, exactly as a product page publishes them in its
-- schema.org JSON-LD. Nothing here is generated: every row is a review the
-- listing's own page stated. Sentiment is derived only from the reviewer's
-- own star rating (>=4 positive, >=3 neutral, <3 negative); NULL when the
-- review states no rating.
CREATE TABLE elec.listing_review (
id BIGSERIAL PRIMARY KEY,
listing_id BIGINT NOT NULL REFERENCES elec.source_listing(id) ON DELETE CASCADE,
author TEXT,
rating NUMERIC(2,1) CHECK (rating IS NULL OR rating BETWEEN 0 AND 5),
title TEXT,
body TEXT NOT NULL,
review_date TEXT,
sentiment TEXT CHECK (sentiment IS NULL OR sentiment IN ('positive','neutral','negative')),
content_hash TEXT NOT NULL,
fetched_at TIMESTAMPTZ NOT NULL DEFAULT now(),
UNIQUE (listing_id, content_hash)
);
CREATE INDEX listing_review_listing_idx ON elec.listing_review (listing_id);

View File

@@ -0,0 +1,588 @@
"""All SQL used by the pipeline. psycopg3, no ORM - the same style as the
original project, with each function owning one statement or one small unit
of work."""
from __future__ import annotations
import hashlib
import json
import re
from datetime import datetime, timezone
from decimal import Decimal
from typing import Any, Dict, List, Optional
from psycopg.types.json import Jsonb
from app.electronics.db.connection import connect, transaction
from app.electronics.models import Listing
from app.electronics.reference import Reference, slugify
def _json(value: Any) -> Jsonb:
return Jsonb(json.loads(json.dumps(value, default=str)))
# ---------------------------------------------------------------------------
# Reference data
# ---------------------------------------------------------------------------
def seed_reference(ref: Reference) -> Dict[str, int]:
"""Idempotent upsert of brands, aliases, categories and sites."""
counts = {"brands": 0, "aliases": 0, "categories": 0, "sites": 0}
with transaction() as conn:
for c in ref.categories.values():
conn.execute(
"INSERT INTO elec.category (slug, name) VALUES (%s, %s) "
"ON CONFLICT (slug) DO UPDATE SET name = EXCLUDED.name",
(c.slug, c.name),
)
counts["categories"] += 1
for b in ref.brands.values():
row = conn.execute(
"INSERT INTO elec.brand (name, slug, official_domains) VALUES (%s, %s, %s) "
"ON CONFLICT (slug) DO UPDATE SET name = EXCLUDED.name, official_domains = EXCLUDED.official_domains "
"RETURNING id",
(b.name, b.slug, list(b.official)),
).fetchone()
counts["brands"] += 1
for alias in b.aliases:
conn.execute(
"INSERT INTO elec.brand_alias (alias, brand_id, is_sub_brand) VALUES (%s, %s, FALSE) "
"ON CONFLICT (alias) DO UPDATE SET brand_id = EXCLUDED.brand_id, is_sub_brand = FALSE",
(alias, row["id"]),
)
counts["aliases"] += 1
for sub in b.sub_brands:
conn.execute(
"INSERT INTO elec.brand_alias (alias, brand_id, is_sub_brand) VALUES (%s, %s, TRUE) "
"ON CONFLICT (alias) DO UPDATE SET brand_id = EXCLUDED.brand_id, is_sub_brand = TRUE",
(sub, row["id"]),
)
counts["aliases"] += 1
for cat in b.categories:
conn.execute(
"INSERT INTO elec.brand_category (brand_id, category_id) "
"SELECT %s, id FROM elec.category WHERE slug = %s ON CONFLICT DO NOTHING",
(row["id"], cat),
)
for s in ref.sites.values():
conn.execute(
"""
INSERT INTO elec.site (domain, name, kind, region, policy, brand_id, product_url, pincode_param)
VALUES (%s, %s, %s, %s, %s, (SELECT id FROM elec.brand WHERE slug = %s), %s, %s)
ON CONFLICT (domain) DO UPDATE SET
name = EXCLUDED.name, kind = EXCLUDED.kind, region = EXCLUDED.region,
policy = EXCLUDED.policy, brand_id = EXCLUDED.brand_id,
product_url = EXCLUDED.product_url, pincode_param = EXCLUDED.pincode_param
""",
(s.domain, s.name, s.kind, s.region, s.policy, s.brand_slug, s.product_url, s.pincode_param),
)
counts["sites"] += 1
return counts
def id_maps() -> Dict[str, Dict[str, int]]:
with connect() as conn:
return {
"brand": {r["slug"]: r["id"] for r in conn.execute("SELECT id, slug FROM elec.brand")},
"category": {r["slug"]: r["id"] for r in conn.execute("SELECT id, slug FROM elec.category")},
"site": {r["domain"]: r["id"] for r in conn.execute("SELECT id, domain FROM elec.site")},
}
def sites() -> List[dict]:
with connect() as conn:
return list(conn.execute("SELECT * FROM elec.site ORDER BY kind, name"))
def set_probe_result(domain: str, outcome: str, robots_allowed: Optional[bool], evidence: dict) -> None:
with transaction() as conn:
conn.execute(
"UPDATE elec.site SET probe_outcome = %s, robots_allowed = %s, probe_evidence = %s, probed_at = now() "
"WHERE domain = %s",
(outcome, robots_allowed, _json(evidence), domain),
)
def trip_breaker(domain: str, reason: str, until_epoch: float) -> None:
with transaction() as conn:
conn.execute(
"UPDATE elec.site SET breaker_until = to_timestamp(%s), breaker_reason = %s, "
"probe_outcome = 'C' WHERE domain = %s OR %s LIKE '%%.' || domain",
(until_epoch, reason, domain, domain),
)
# ---------------------------------------------------------------------------
# Runs and fetch log
# ---------------------------------------------------------------------------
def start_run(kind: str, params: dict) -> int:
with transaction() as conn:
return conn.execute(
"INSERT INTO elec.crawl_run (kind, params) VALUES (%s, %s) RETURNING id", (kind, _json(params))
).fetchone()["id"]
def finish_run(run_id: int, status: str, stats: dict, error: Optional[str] = None) -> None:
with transaction() as conn:
conn.execute(
"UPDATE elec.crawl_run SET status = %s, stats = %s, error = %s, ended_at = now() WHERE id = %s",
(status, _json(stats), error, run_id),
)
def log_fetch(run_id: Optional[int], url: str, host: str, status: Optional[int], nbytes: int,
outcome: str, robots_allowed: Optional[bool]) -> None:
with transaction() as conn:
conn.execute(
"INSERT INTO elec.fetch_log (crawl_run_id, url, host, status, bytes, outcome, robots_allowed) "
"VALUES (%s, %s, %s, %s, %s, %s, %s)",
(run_id, url, host, status, nbytes, outcome, robots_allowed),
)
def recent_runs(limit: int = 20) -> List[dict]:
with connect() as conn:
return list(conn.execute("SELECT * FROM elec.crawl_run ORDER BY id DESC LIMIT %s", (limit,)))
# ---------------------------------------------------------------------------
# Search cache
# ---------------------------------------------------------------------------
def search_cache_get(provider: str, kind: str, query: str, ttl_hours: int) -> Optional[List[dict]]:
with connect() as conn:
row = conn.execute(
"SELECT results FROM elec.search_cache WHERE provider = %s AND kind = %s AND query = %s "
"AND fetched_at > now() - make_interval(hours => %s)",
(provider, kind, query, ttl_hours),
).fetchone()
return row["results"] if row else None
def search_cache_put(provider: str, kind: str, query: str, results: List[dict]) -> None:
with transaction() as conn:
conn.execute(
"INSERT INTO elec.search_cache (provider, kind, query, results) VALUES (%s, %s, %s, %s) "
"ON CONFLICT (provider, kind, query) DO UPDATE SET results = EXCLUDED.results, fetched_at = now()",
(provider, kind, query, _json(results)),
)
def google_queries_today() -> int:
with connect() as conn:
return conn.execute(
"SELECT count(*) AS n FROM elec.search_cache WHERE provider = 'google' AND fetched_at::date = current_date"
).fetchone()["n"]
# ---------------------------------------------------------------------------
# Listings and prices
# ---------------------------------------------------------------------------
def upsert_listing(listing: Listing, ids: Dict[str, Dict[str, int]], run_id: Optional[int]) -> int:
"""Write the latest state of a listing and append one price observation."""
listing.validate()
site_id = ids["site"][listing.site_domain]
brand_id = ids["brand"][listing.brand_slug]
category_id = ids["category"][listing.category]
with transaction() as conn:
existing = conn.execute(
"SELECT id, source_type, price FROM elec.source_listing WHERE site_id = %s AND source_sku = %s",
(site_id, listing.source_sku),
).fetchone()
# A scraped page is better evidence than a search snippet about the
# same page. Never let a later snippet overwrite scraped values.
if existing and existing["source_type"] in ("scraped_page", "brand_official") and listing.source_type == "search_snippet":
conn.execute("UPDATE elec.source_listing SET last_seen_at = now() WHERE id = %s", (existing["id"],))
return existing["id"]
params = dict(
site_id=site_id, source_sku=listing.source_sku, source_url=listing.source_url,
source_type=listing.source_type, brand_id=brand_id, category_id=category_id,
family=listing.family, title=listing.title[:500], model=listing.model,
model_number=listing.model_number, ram_gb=listing.ram_gb, storage_gb=listing.storage_gb,
colour=listing.colour, price=listing.price, mrp=listing.mrp, availability=listing.availability,
in_stock=listing.in_stock, pincode=listing.pincode, pincode_applied=listing.pincode_applied,
rating=listing.rating, review_count=listing.review_count, gtin=listing.gtin,
image_urls=listing.image_urls[:12], specs_raw=_json(listing.specs_raw), specs=_json(listing.specs),
evidence_text=listing.evidence_text[:4000], search_query=listing.search_query,
confidence=round(listing.confidence, 2), parser=listing.parser, content_hash=listing.content_hash,
crawl_run_id=run_id,
)
row = conn.execute(
"""
INSERT INTO elec.source_listing (
site_id, source_sku, source_url, source_type, brand_id, category_id, family, title, model,
model_number, ram_gb, storage_gb, colour, price, mrp, availability, in_stock, pincode,
pincode_applied, rating, review_count, gtin, image_urls, specs_raw, specs, evidence_text,
search_query, confidence, parser, content_hash, crawl_run_id)
VALUES (
%(site_id)s, %(source_sku)s, %(source_url)s, %(source_type)s, %(brand_id)s, %(category_id)s,
%(family)s, %(title)s, %(model)s, %(model_number)s, %(ram_gb)s, %(storage_gb)s, %(colour)s,
%(price)s, %(mrp)s, %(availability)s, %(in_stock)s, %(pincode)s, %(pincode_applied)s,
%(rating)s, %(review_count)s, %(gtin)s, %(image_urls)s, %(specs_raw)s, %(specs)s,
%(evidence_text)s, %(search_query)s, %(confidence)s, %(parser)s, %(content_hash)s,
%(crawl_run_id)s)
ON CONFLICT (site_id, source_sku) DO UPDATE SET
source_url = EXCLUDED.source_url, source_type = EXCLUDED.source_type,
brand_id = EXCLUDED.brand_id, category_id = EXCLUDED.category_id, family = EXCLUDED.family,
title = EXCLUDED.title, model = EXCLUDED.model, model_number = EXCLUDED.model_number,
ram_gb = EXCLUDED.ram_gb, storage_gb = EXCLUDED.storage_gb, colour = EXCLUDED.colour,
price = EXCLUDED.price, mrp = EXCLUDED.mrp, availability = EXCLUDED.availability,
in_stock = EXCLUDED.in_stock, pincode = EXCLUDED.pincode,
pincode_applied = EXCLUDED.pincode_applied, rating = EXCLUDED.rating,
review_count = EXCLUDED.review_count, gtin = EXCLUDED.gtin, image_urls = EXCLUDED.image_urls,
specs_raw = EXCLUDED.specs_raw, specs = EXCLUDED.specs, evidence_text = EXCLUDED.evidence_text,
search_query = EXCLUDED.search_query, confidence = EXCLUDED.confidence, parser = EXCLUDED.parser,
content_hash = EXCLUDED.content_hash, crawl_run_id = EXCLUDED.crawl_run_id, last_seen_at = now()
RETURNING id
""",
params,
).fetchone()
listing_id = row["id"]
if listing.price is not None or listing.in_stock is not None:
conn.execute(
"INSERT INTO elec.price_history (listing_id, price, mrp, availability, in_stock, source_type, "
"pincode, pincode_applied, evidence_text, crawl_run_id) VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s)",
(listing_id, listing.price, listing.mrp, listing.availability, listing.in_stock,
listing.source_type, listing.pincode, listing.pincode_applied,
listing.evidence_text[:2000], run_id),
)
return listing_id
# ---------------------------------------------------------------------------
# Ratings and reviews
# ---------------------------------------------------------------------------
def save_reviews(listing_id: int, reviews: List[Dict[str, Any]]) -> int:
"""Store the reviews a listing's page publishes. Idempotent per review
text; an empty list changes nothing (a later search-only sighting must
not erase what the page said). Returns the number of new rows."""
from app.electronics.reviews import sentiment_for
added = 0
with transaction() as conn:
for r in reviews:
body = (r.get("body") or "").strip()
if not body:
continue
digest = hashlib.sha1(f"{r.get('author') or ''}|{body}".encode("utf-8", "ignore")).hexdigest()
row = conn.execute(
"INSERT INTO elec.listing_review (listing_id, author, rating, title, body, review_date, sentiment, "
"content_hash) VALUES (%s,%s,%s,%s,%s,%s,%s,%s) "
"ON CONFLICT (listing_id, content_hash) DO NOTHING RETURNING id",
(listing_id, (r.get("author") or None) and str(r["author"])[:200], r.get("rating"),
(r.get("title") or None) and str(r["title"])[:300], body[:4000],
(r.get("review_date") or None) and str(r["review_date"])[:40],
sentiment_for(r.get("rating")), digest),
).fetchone()
added += 1 if row else 0
return added
def update_listing_rating(listing_id: int, rating: Optional[Decimal], review_count: Optional[int]) -> None:
"""Refresh only the rating fields of a listing (used by the review backfill)."""
with transaction() as conn:
conn.execute(
"UPDATE elec.source_listing SET rating = %s, review_count = %s WHERE id = %s",
(rating, review_count, listing_id),
)
def product_rating_and_reviews(conn, product_id: int) -> Dict[str, Any]:
"""Per-platform ratings and all stored reviews for a verified product's
approved listings, each with the page it was read from."""
sources = conn.execute(
"SELECT a.site, a.source_url, l.rating, l.review_count FROM elec.v_product_availability a "
"JOIN elec.source_listing l ON l.id = a.listing_id "
"WHERE a.product_id = %s AND l.rating > 0 ORDER BY l.review_count DESC NULLS LAST, a.site",
(product_id,),
).fetchall()
reviews = conn.execute(
"SELECT a.site, a.source_url, r.author, r.rating, r.title, r.body, r.review_date, r.sentiment "
"FROM elec.v_product_availability a JOIN elec.listing_review r ON r.listing_id = a.listing_id "
"WHERE a.product_id = %s",
(product_id,),
).fetchall()
return {"sources": [dict(s) for s in sources], "reviews": [dict(r) for r in reviews]}
def listings_for_review_backfill(category: Optional[str] = None) -> List[dict]:
"""Page-read listings of verified products, for re-reading ratings/reviews."""
sql = (
"SELECT a.listing_id, a.source_url, a.domain, a.site_kind, a.category, l.source_sku, l.title "
"FROM elec.v_product_availability a JOIN elec.source_listing l ON l.id = a.listing_id "
"WHERE a.source_type IN ('scraped_page','brand_official')"
)
params: tuple = ()
if category:
sql += " AND a.category = %s"
params = (category,)
with connect() as conn:
return list(conn.execute(sql + " ORDER BY a.listing_id", params))
# ---------------------------------------------------------------------------
# Products, matching, images
# ---------------------------------------------------------------------------
def product_candidates(brand_slug: str, category: str) -> List[dict]:
with connect() as conn:
return list(conn.execute(
"SELECT p.id, p.variant_key, p.model_norm, p.ram_gb, p.storage_gb, p.processor, p.mpn, p.gtin "
"FROM elec.product p JOIN elec.brand b ON b.id = p.brand_id JOIN elec.category c ON c.id = p.category_id "
"WHERE b.slug = %s AND c.slug = %s AND p.verification_status <> 'rejected'",
(brand_slug, category),
))
def _cpu_label(processor: Optional[str]) -> str:
""""ryzen 5 7530u" -> "Ryzen 5 7530U", "i5-1334u" -> "i5-1334U"."""
if not processor:
return ""
def fmt(t: str) -> str:
if re.fullmatch(r"i[3579]-\w+", t):
return "i" + t[1:].upper() # i5-1334U
if any(ch.isdigit() for ch in t):
return t.upper() # 7530U, M5
return t.title() # Ryzen, Core, Ultra
return " ".join(fmt(t) for t in processor.split())
def product_display_name(listing: Listing) -> str:
variant = [x for x in (
_cpu_label(listing.processor) if listing.category == "laptops" else "",
f"{_fmt_gb(listing.ram_gb)} RAM" if listing.ram_gb else "",
_fmt_gb(listing.storage_gb) if listing.storage_gb else "",
) if x]
if listing.category == "laptops" and not listing.processor and listing.model_number:
variant.insert(0, listing.model_number) # the part number is what tells it apart
return " ".join(x for x in [load_brand_name(listing.brand_slug), listing.model,
f"({', '.join(variant)})" if variant else ""] if x)
def create_product(listing: Listing, ids: Dict[str, Dict[str, int]]) -> int:
display = product_display_name(listing)
with transaction() as conn:
row = conn.execute(
"""
INSERT INTO elec.product (brand_id, category_id, family, model, model_norm, variant_key, display_name,
ram_gb, storage_gb, processor, mpn, gtin)
VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s)
ON CONFLICT (variant_key) DO UPDATE SET updated_at = now()
RETURNING id
""",
(ids["brand"][listing.brand_slug], ids["category"][listing.category], listing.family,
listing.model or listing.model_norm, listing.model_norm, listing.variant_key, display,
listing.ram_gb, listing.storage_gb, listing.processor, listing.model_number, listing.gtin),
).fetchone()
return row["id"]
def _fmt_gb(value: Optional[Decimal]) -> str:
if value is None:
return ""
if value >= 1024 and value % 1024 == 0:
return f"{int(value // 1024)}TB"
return f"{format(value.normalize(), 'f')}GB"
_BRAND_NAMES: Dict[str, str] = {}
def load_brand_name(slug: str) -> str:
if not _BRAND_NAMES:
from app.electronics.reference import load_reference
_BRAND_NAMES.update({s: b.name for s, b in load_reference().brands.items()})
return _BRAND_NAMES.get(slug, slug.title())
def map_listing(listing_id: int, product_id: int, method: str, confidence: float, review_status: str) -> None:
with transaction() as conn:
conn.execute(
"""
INSERT INTO elec.product_listing_map (listing_id, product_id, method, confidence, review_status)
VALUES (%s, %s, %s, %s, %s)
ON CONFLICT (listing_id) DO UPDATE SET
product_id = EXCLUDED.product_id, method = EXCLUDED.method, confidence = EXCLUDED.confidence,
review_status = CASE WHEN elec.product_listing_map.review_status IN ('approved','rejected')
AND elec.product_listing_map.product_id = EXCLUDED.product_id
THEN elec.product_listing_map.review_status
ELSE EXCLUDED.review_status END
""",
(listing_id, product_id, method, round(confidence, 2), review_status),
)
def merge_product_specs(product_id: int, specs: Dict[str, Any], sources: Dict[str, str], source_url: str) -> None:
"""Add spec keys the product does not have yet. Existing values win:
specs are only ever filled, never overwritten by a later source."""
if not specs:
return
with transaction() as conn:
row = conn.execute(
"SELECT canonical_specs, spec_sources FROM elec.product WHERE id = %s FOR UPDATE", (product_id,)
).fetchone()
current, cur_src = dict(row["canonical_specs"] or {}), dict(row["spec_sources"] or {})
changed = False
for key, value in specs.items():
if key not in current:
current[key] = value
cur_src[key] = {"url": source_url, "from": sources.get(key, "")}
changed = True
if changed:
conn.execute(
"UPDATE elec.product SET canonical_specs = %s, spec_sources = %s, updated_at = now() WHERE id = %s",
(_json(current), _json(cur_src), product_id),
)
def add_image(product_id: int, url: str, listing_id: int, source_type: str, rank: int) -> None:
with transaction() as conn:
conn.execute(
"INSERT INTO elec.product_image (product_id, url, source_listing_id, source_type, rank) "
"VALUES (%s, %s, %s, %s, %s) ON CONFLICT (product_id, url) DO UPDATE SET validated_at = now()",
(product_id, url, listing_id, source_type, rank),
)
def product_image_count(product_id: int) -> int:
with connect() as conn:
return conn.execute("SELECT count(*) AS n FROM elec.product_image WHERE product_id = %s",
(product_id,)).fetchone()["n"]
# The least a new device in the category can plausibly cost. Anything below is
# an accessory, an EMI or an offer amount that slipped through.
CATEGORY_MIN_PRICE = {"mobiles": 3000, "laptops": 15000}
OUTLIER_TOLERANCE = 0.35
def flag_price_outliers() -> int:
"""Flag prices that cannot be trusted as this product's price:
* below the category's floor (CATEGORY_MIN_PRICE);
* a search-result price more than OUTLIER_TOLERANCE away from the price
read off a product page for the same product;
* with no page price, a search-result price that far from the median of
at least three prices for the product.
Flagged prices stay stored with their evidence; they are just never used as
the best price. Returns the number flagged."""
floor_cases = " ".join(f"WHEN '{k}' THEN {v}" for k, v in CATEGORY_MIN_PRICE.items())
with transaction() as conn:
conn.execute("UPDATE elec.source_listing SET price_outlier = FALSE WHERE price_outlier")
cur = conn.execute(
f"""
WITH prices AS (
SELECT l.id, m.product_id, l.price, l.source_type, c.slug
FROM elec.source_listing l
JOIN elec.product_listing_map m ON m.listing_id = l.id AND m.review_status IN ('auto','approved')
JOIN elec.category c ON c.id = l.category_id
WHERE l.price IS NOT NULL
),
ref AS (
SELECT product_id,
percentile_cont(0.5) WITHIN GROUP (ORDER BY price)
FILTER (WHERE source_type <> 'search_snippet') AS page_median,
percentile_cont(0.5) WITHIN GROUP (ORDER BY price) AS all_median,
count(*) AS n
FROM prices GROUP BY product_id
)
UPDATE elec.source_listing l SET price_outlier = TRUE
FROM prices p JOIN ref r ON r.product_id = p.product_id
WHERE l.id = p.id AND (
p.price < CASE p.slug {floor_cases} ELSE 0 END
OR (p.source_type = 'search_snippet' AND r.page_median IS NOT NULL
AND abs(p.price - r.page_median) / r.page_median > %(tol)s)
OR (p.source_type = 'search_snippet' AND r.page_median IS NULL AND r.n >= 3
AND abs(p.price - r.all_median) / r.all_median > %(tol)s)
)
""",
{"tol": OUTLIER_TOLERANCE},
)
return cur.rowcount
def refresh_verification() -> Dict[str, int]:
flag_price_outliers()
"""A product is VERIFIED when auto/approved listings on at least two
different sites point at it, and at least one of them is a retailer
(so it is actually sold). Everything else stays unverified and hidden."""
with transaction() as conn:
conn.execute(
"""
WITH ev AS (
SELECT m.product_id,
count(DISTINCT l.site_id) AS sites,
count(DISTINCT l.site_id) FILTER (WHERE s.kind <> 'brand_official') AS retail_sites
FROM elec.product_listing_map m
JOIN elec.source_listing l ON l.id = m.listing_id
JOIN elec.site s ON s.id = l.site_id
WHERE m.review_status IN ('auto','approved')
GROUP BY m.product_id
)
UPDATE elec.product p SET
evidence_count = coalesce(ev.sites, 0),
verification_status = CASE
WHEN p.verification_status = 'rejected' THEN 'rejected'
WHEN coalesce(ev.sites, 0) >= 2 AND coalesce(ev.retail_sites, 0) >= 1 THEN 'verified'
ELSE 'unverified' END,
updated_at = now()
FROM elec.product p2 LEFT JOIN ev ON ev.product_id = p2.id
WHERE p.id = p2.id
"""
)
rows = conn.execute(
"SELECT verification_status AS s, count(*) AS n FROM elec.product GROUP BY 1"
).fetchall()
return {r["s"]: r["n"] for r in rows}
def products_without_embedding(limit: int = 500) -> List[dict]:
with connect() as conn:
return list(conn.execute(
"SELECT p.id, p.display_name, p.canonical_specs, b.name AS brand, c.name AS category "
"FROM elec.product p JOIN elec.brand b ON b.id = p.brand_id JOIN elec.category c ON c.id = p.category_id "
"WHERE p.embedding IS NULL AND p.verification_status = 'verified' LIMIT %s", (limit,)))
def set_embedding(product_id: int, vector: List[float]) -> None:
import numpy as np
with transaction() as conn:
conn.execute("UPDATE elec.product SET embedding = %s WHERE id = %s", (np.array(vector), product_id))
def review_queue(limit: int = 100) -> List[dict]:
with connect() as conn:
return list(conn.execute(
"""
SELECT m.listing_id, m.product_id, m.method, m.confidence, l.title AS listing_title, l.source_url,
s.name AS site, p.display_name AS product
FROM elec.product_listing_map m
JOIN elec.source_listing l ON l.id = m.listing_id
JOIN elec.site s ON s.id = l.site_id
JOIN elec.product p ON p.id = m.product_id
WHERE m.review_status = 'pending'
ORDER BY m.confidence DESC, m.listing_id LIMIT %s
""", (limit,)))
def set_review(listing_id: int, approve: bool) -> bool:
with transaction() as conn:
cur = conn.execute(
"UPDATE elec.product_listing_map SET review_status = %s, reviewed_at = now() "
"WHERE listing_id = %s AND review_status = 'pending'",
("approved" if approve else "rejected", listing_id),
)
return cur.rowcount > 0
def now_utc() -> datetime:
return datetime.now(timezone.utc)
def grounding_sample(n: int = 50) -> List[dict]:
with connect() as conn:
return list(conn.execute(
"SELECT id, source_url, source_type, price, evidence_text FROM elec.source_listing "
"WHERE price IS NOT NULL ORDER BY random() LIMIT %s", (n,)))
def slug(text: str) -> str:
return slugify(text)