← back to Hollywood Import
Fix nightly clobber: patch consolidate-vendor-catalogs.py to preserve list_price/hw_price/retail_price/cost/image_alt + Hollywood name-scrub on ON CONFLICT upsert (deployed to Kamatera, validated by dry-run + merge sim)
724cdd91ef99c86aedccb61f84473ddfa9367f0a · 2026-06-18 07:08:59 -0700 · SteveStudio2
Files touched
A consolidate-vendor-catalogs.pyA thib-price.mjs
Diff
commit 724cdd91ef99c86aedccb61f84473ddfa9367f0a
Author: SteveStudio2 <stevestudio2@SteveStacStudio.lan>
Date: Thu Jun 18 07:08:59 2026 -0700
Fix nightly clobber: patch consolidate-vendor-catalogs.py to preserve list_price/hw_price/retail_price/cost/image_alt + Hollywood name-scrub on ON CONFLICT upsert (deployed to Kamatera, validated by dry-run + merge sim)
---
consolidate-vendor-catalogs.py | 920 +++++++++++++++++++++++++++++++++++++++++
thib-price.mjs | 25 ++
2 files changed, 945 insertions(+)
diff --git a/consolidate-vendor-catalogs.py b/consolidate-vendor-catalogs.py
new file mode 100644
index 0000000..0e88a1d
--- /dev/null
+++ b/consolidate-vendor-catalogs.py
@@ -0,0 +1,920 @@
+#!/usr/bin/env python3
+"""
+consolidate-vendor-catalogs.py
+-------------------------------
+Consolidates all individual vendor catalog tables into the unified `vendor_catalog`
+table using upsert semantics. COALESCE is used throughout so that a NULL value
+from the source row never overwrites a non-NULL value that already exists in the
+target.
+
+Usage:
+ python consolidate-vendor-catalogs.py --all
+ python consolidate-vendor-catalogs.py --vendor cowtan_tout
+ python consolidate-vendor-catalogs.py --all --dry-run
+ python consolidate-vendor-catalogs.py --all --shopify-first
+"""
+
+from __future__ import annotations
+
+import argparse
+import logging
+import sys
+import time
+from dataclasses import dataclass, field
+from datetime import datetime
+from pathlib import Path
+from typing import Any
+
+import psycopg2
+import psycopg2.extras
+
+# ---------------------------------------------------------------------------
+# Configuration
+# ---------------------------------------------------------------------------
+
+DSN = "postgresql://dw_admin@127.0.0.1:5432/dw_unified"
+
+LOG_DIR = Path("/root/DW-Agents/logs")
+LOG_FILE = LOG_DIR / "consolidation.log"
+
+# Tables to skip entirely — competitive intel, not vendor products.
+SKIP_TABLES: frozenset[str] = frozenset(
+ {
+ "connie_competitor_catalog",
+ "connie_our_catalog",
+ "connie_products",
+ "vendor_catalog", # the target table itself
+ "catalog_elli_popp", # no matching pg_stat row yet
+ }
+)
+
+# vendor_codes for which we set sync_status='never_push' after consolidation.
+NEVER_PUSH_VENDORS: frozenset[str] = frozenset({"cowtan_tout", "colefax_fowler"})
+
+# Tables with non-standard primary-key / name columns that need special mapping.
+# Keys are table names; values are dicts describing the override mapping.
+SPECIAL_TABLES: dict[str, dict[str, str]] = {
+ "arte_catalog": {
+ # arte has both arte_sku and mfr_sku; prefer mfr_sku, fall back to arte_sku.
+ "mfr_sku": "COALESCE(mfr_sku, arte_sku)",
+ },
+ "bespoke_catalog": {
+ "mfr_sku": "COALESCE(mfr_sku, sku)",
+ "pattern_name": "COALESCE(pattern_name, title)",
+ },
+ "jeffrey_stevens_catalog": {
+ "mfr_sku": "COALESCE(mfr_sku, sku)",
+ "pattern_name": "COALESCE(pattern_name, name)",
+ },
+ "malibu_catalog": {
+ "mfr_sku": "COALESCE(mfr_sku, sku)",
+ "pattern_name": "COALESCE(pattern_name, name)",
+ },
+ "relativity_textiles_catalog": {
+ "mfr_sku": "COALESCE(mfr_sku, sku)",
+ "pattern_name": "name", # source uses 'name' column, not 'pattern_name'
+ },
+ "glitter_catalog": {
+ "specs": "specs::jsonb", # source is text, target is jsonb
+ },
+ "osborne_catalog": {
+ # source 'specs' is text holding a JSON object; target is jsonb.
+ # NULLIF guards empty-string text so it doesn't fail the ::jsonb cast.
+ "specs": "NULLIF(specs, '')::jsonb",
+ },
+ # momentum_colorways uses a completely different shape.
+ "momentum_colorways": {
+ "mfr_sku": "COALESCE(momentum_sku, alt_sku)",
+ "collection": "collection_name",
+ "composition": "content",
+ # vendor_code is dynamic per row (pl_brand field), handled separately.
+ },
+}
+
+# Target column varchar limits — expressions for these columns get LEFT() truncated.
+# Columns not listed here are TEXT (unlimited) and need no truncation.
+TARGET_COL_LIMITS: dict[str, int] = {
+ "color_hex": 7, "dominant_color_hex": 7, "product_category": 20,
+ "repeat_classification": 20, "color_primary": 100, "color_secondary": 100,
+ "fire_rating_eu": 100, "fire_rating_us": 100, "length": 100, "match_type": 100,
+ "mfr_sku": 100, "pattern_repeat": 100, "product_type": 100, "sync_status": 100,
+ "vendor_code": 100, "width": 100, "adhesive": 200, "collection": 255,
+ "color_name": 255, "coverage": 255, "design": 255, "finish": 255,
+ "pattern_name": 255, "us_distributor": 255,
+}
+
+# Columns in vendor catalog tables that map to non-identically-named target columns.
+# Any table that has the *source* column name but NOT the *target* column name will
+# use the alias listed here. A None target means "no mapping — skip this column".
+GENERIC_COLUMN_ALIASES: dict[str, str | None] = {
+ "roll_length": "length", # source col -> target col
+ "source_url": "product_url",
+ "spec_sheet": "spec_sheet_url",
+ "description": None, # some tables have 'description' not mapped anywhere
+}
+
+# vendor_catalog target columns we care about (order matches INSERT statement).
+TARGET_COLS: list[str] = [
+ "vendor_code",
+ "mfr_sku",
+ "pattern_name",
+ "color_name",
+ "collection",
+ "product_type",
+ "composition",
+ "width",
+ "length",
+ "pattern_repeat",
+ "fire_rating",
+ "fire_rating_us",
+ "fire_rating_eu",
+ "image_url",
+ "product_url",
+ "specs",
+ "dw_sku",
+ "ai_colors",
+ "ai_background_color",
+ "ai_styles",
+ "ai_patterns",
+ "ai_tags",
+ "ai_description",
+ "color_hex",
+ "width_inches",
+ "dominant_color_hex",
+ "on_shopify",
+ "all_images",
+ "finish",
+ "application",
+ "match_type",
+ "coverage",
+ "features",
+ "design",
+ "rooms",
+ "color_primary",
+ "color_secondary",
+ "first_seen_at",
+ "last_scraped_at",
+ "abrasion",
+ "sustainability",
+ "adhesive",
+ "maintenance",
+ "spec_sheet_url",
+ "designer",
+ "us_distributor",
+ "showroom_locations",
+ "about_vendor",
+]
+
+# For the ON CONFLICT DO UPDATE clause we skip the upsert key cols and identity cols.
+UPDATE_SKIP_COLS: frozenset[str] = frozenset({"vendor_code", "mfr_sku", "first_seen_at"})
+
+# ---------------------------------------------------------------------------
+# Logging setup
+# ---------------------------------------------------------------------------
+
+LOG_DIR.mkdir(parents=True, exist_ok=True)
+
+logging.basicConfig(
+ level=logging.INFO,
+ format="%(asctime)s [%(levelname)s] %(message)s",
+ handlers=[
+ logging.FileHandler(LOG_FILE),
+ logging.StreamHandler(sys.stdout),
+ ],
+)
+log = logging.getLogger(__name__)
+
+
+# ---------------------------------------------------------------------------
+# Data structures
+# ---------------------------------------------------------------------------
+
+@dataclass
+class VendorResult:
+ vendor_code: str
+ table_name: str
+ inserted: int = 0
+ updated: int = 0
+ skipped: int = 0
+ duration_s: float = 0.0
+ error: str | None = None
+
+
+@dataclass
+class ConsolidationSummary:
+ results: list[VendorResult] = field(default_factory=list)
+ shopify_marked: int = 0
+ never_push_marked: int = 0
+ total_duration_s: float = 0.0
+
+
+# ---------------------------------------------------------------------------
+# Database helpers
+# ---------------------------------------------------------------------------
+
+def get_connection() -> psycopg2.extensions.connection:
+ return psycopg2.connect(DSN)
+
+
+def get_table_columns(cur: psycopg2.extensions.cursor, table_name: str) -> set[str]:
+ """Return the set of column names for a given table."""
+ cur.execute(
+ """
+ SELECT column_name
+ FROM information_schema.columns
+ WHERE table_schema = 'public'
+ AND table_name = %s
+ """,
+ (table_name,),
+ )
+ return {row[0] for row in cur.fetchall()}
+
+
+def get_live_row_count(cur: psycopg2.extensions.cursor, table_name: str) -> int:
+ """Return approximate live row count from pg_stat_user_tables."""
+ cur.execute(
+ "SELECT n_live_tup FROM pg_stat_user_tables WHERE relname = %s",
+ (table_name,),
+ )
+ row = cur.fetchone()
+ return row[0] if row else 0
+
+
+def get_vendor_registry(
+ cur: psycopg2.extensions.cursor,
+) -> dict[str, dict[str, Any]]:
+ """Return a dict keyed by catalog_table -> vendor info."""
+ cur.execute(
+ """
+ SELECT vendor_code, catalog_table, vendor_name
+ FROM vendor_registry
+ WHERE catalog_table IS NOT NULL
+ """
+ )
+ mapping: dict[str, dict[str, Any]] = {}
+ for vendor_code, catalog_table, vendor_name in cur.fetchall():
+ if catalog_table and catalog_table not in ("vendor_catalog",):
+ mapping[catalog_table] = {
+ "vendor_code": vendor_code,
+ "vendor_name": vendor_name,
+ }
+ return mapping
+
+
+def get_all_catalog_tables(
+ cur: psycopg2.extensions.cursor,
+) -> list[tuple[str, int]]:
+ """Return (table_name, n_live_tup) for every catalog-like table with rows."""
+ cur.execute(
+ """
+ SELECT relname, n_live_tup
+ FROM pg_stat_user_tables
+ WHERE (relname LIKE '%_catalog' OR relname = 'momentum_colorways')
+ AND n_live_tup > 0
+ ORDER BY n_live_tup DESC
+ """
+ )
+ return cur.fetchall()
+
+
+def get_shopify_mfr_skus(cur: psycopg2.extensions.cursor) -> set[str]:
+ """Return mfr_skus of ACTIVE Shopify products."""
+ cur.execute(
+ "SELECT DISTINCT mfr_sku FROM shopify_products WHERE status = 'ACTIVE' AND mfr_sku IS NOT NULL"
+ )
+ return {row[0] for row in cur.fetchall()}
+
+
+# ---------------------------------------------------------------------------
+# SQL builder
+# ---------------------------------------------------------------------------
+
+def resolve_vendor_code_expr(
+ table_name: str,
+ vendor_code: str,
+ source_cols: set[str],
+) -> str:
+ """
+ Return the SQL expression that produces vendor_code for each source row.
+
+ For momentum_colorways the vendor_code is derived from the pl_brand column.
+ For all other tables it is the literal string from vendor_registry.
+ """
+ if table_name == "momentum_colorways" and "pl_brand" in source_cols:
+ # Map pl_brand values to vendor_code. Only 'Hollywood Wallcoverings'
+ # is known at the time of writing; fall back to 'momentum' for all others.
+ return (
+ "CASE pl_brand "
+ "WHEN 'Hollywood Wallcoverings' THEN 'hollywood' "
+ "ELSE 'momentum' END"
+ )
+ return f"'{vendor_code}'"
+
+
+def build_source_expr(
+ col: str,
+ table_name: str,
+ source_cols: set[str],
+ special_mapping: dict[str, str],
+) -> str | None:
+ """
+ Return the SQL expression for a single target column from this source table,
+ or None if the column cannot be mapped.
+
+ Resolution order:
+ 1. Special per-table override (SPECIAL_TABLES).
+ 2. Direct name match in source_cols.
+ 3. Generic alias (GENERIC_COLUMN_ALIASES).
+ 4. Derived expression (pattern_repeat from repeat_v/repeat_h).
+ 5. None — column not available in this table.
+ """
+ # 1. Special override
+ if col in special_mapping:
+ return special_mapping[col]
+
+ # 2. Direct match
+ if col in source_cols:
+ return col
+
+ # 3. Generic aliases (exact rename: source col -> target col)
+ for src_col, tgt_col in GENERIC_COLUMN_ALIASES.items():
+ if tgt_col == col and src_col in source_cols and col not in source_cols:
+ return src_col
+
+ # 4. Derived: pattern_repeat from repeat_v / repeat_h
+ if col == "pattern_repeat":
+ has_v = "repeat_v" in source_cols
+ has_h = "repeat_h" in source_cols
+ if has_v and has_h:
+ # NULL-safe: only concatenate when at least one repeat is not null
+ return (
+ "CASE WHEN repeat_v IS NOT NULL OR repeat_h IS NOT NULL "
+ "THEN COALESCE(repeat_v::text,'') || ' x ' || COALESCE(repeat_h::text,'') "
+ "END"
+ )
+ if has_v:
+ return "repeat_v::text"
+ if has_h:
+ return "repeat_h::text"
+
+ # 5. composition: prefer 'composition' col, fall back to 'material'
+ if col == "composition" and "composition" not in source_cols:
+ if "material" in source_cols:
+ return "material"
+
+ # 6. last_scraped_at from crawled_at, scraped_at, last_scraped, updated_at
+ if col == "last_scraped_at":
+ for candidate in ("last_scraped_at", "crawled_at", "scraped_at", "last_scraped", "updated_at"):
+ if candidate in source_cols:
+ return candidate
+
+ # 7. first_seen_at from created_at, crawled_at
+ if col == "first_seen_at":
+ for candidate in ("first_seen_at", "created_at", "crawled_at"):
+ if candidate in source_cols:
+ return candidate
+
+ # 8. spec_sheet_url from spec_sheet (also covered by generic alias, explicit here for clarity)
+ if col == "spec_sheet_url" and "spec_sheet_url" not in source_cols:
+ if "spec_sheet" in source_cols:
+ return "spec_sheet"
+
+ return None
+
+
+def build_upsert_sql(
+ table_name: str,
+ vendor_code: str,
+ source_cols: set[str],
+ dry_run: bool = False,
+) -> tuple[str, list[str]] | None:
+ """
+ Build the INSERT ... ON CONFLICT DO UPDATE SQL for one vendor table.
+
+ Returns (sql_string, list_of_mapped_target_cols) or None if the table
+ cannot be mapped (missing mfr_sku).
+ """
+ special_mapping = SPECIAL_TABLES.get(table_name, {})
+
+ # Verify mfr_sku can be resolved — it is the primary business key.
+ mfr_expr = build_source_expr("mfr_sku", table_name, source_cols, special_mapping)
+ if mfr_expr is None:
+ log.warning("Table %s has no resolvable mfr_sku — skipping.", table_name)
+ return None
+
+ vendor_expr = resolve_vendor_code_expr(table_name, vendor_code, source_cols)
+
+ select_parts: list[str] = []
+ mapped_cols: list[str] = []
+
+ for col in TARGET_COLS:
+ if col == "vendor_code":
+ select_parts.append(f"{vendor_expr} AS vendor_code")
+ mapped_cols.append("vendor_code")
+ continue
+ if col == "mfr_sku":
+ select_parts.append(f"({mfr_expr}) AS mfr_sku")
+ mapped_cols.append("mfr_sku")
+ continue
+
+ expr = build_source_expr(col, table_name, source_cols, special_mapping)
+ if expr is not None:
+ # Truncate to target varchar limit if one exists
+ limit = TARGET_COL_LIMITS.get(col)
+ if limit:
+ expr = f"LEFT(({expr})::text, {limit})"
+ select_parts.append(f"{expr} AS {col}")
+ mapped_cols.append(col)
+
+ if not mapped_cols:
+ log.warning("Table %s produced no mappable columns.", table_name)
+ return None
+
+ col_list = ", ".join(mapped_cols)
+ select_list = ", ".join(select_parts)
+
+ # Filter out rows where the resolved mfr_sku is NULL or empty string.
+ where_clause = f"WHERE ({mfr_expr}) IS NOT NULL AND ({mfr_expr}) <> ''"
+
+ # ON CONFLICT DO UPDATE: use COALESCE(NULLIF(...), existing) so we never
+ # overwrite non-null with null OR empty string (NULLIF closes the '' gap).
+ update_parts: list[str] = []
+ # Columns where NULLIF(value, '') doesn't apply (jsonb, boolean, numeric, etc.)
+ no_nullif_cols = frozenset({
+ "ai_colors", "ai_styles", "ai_patterns", "ai_tags", "specs",
+ "on_shopify", "width_inches", "all_images",
+ })
+ for col in mapped_cols:
+ if col in UPDATE_SKIP_COLS:
+ continue
+ if col == "specs":
+ # PRESERVE manually-managed keys (pricing + alt) the scrape source does
+ # not carry, while still refreshing everything the scrape provides:
+ # scraped_specs || {protected keys from the existing row}
+ update_parts.append(
+ " specs = COALESCE(EXCLUDED.specs, '{}'::jsonb) "
+ "|| jsonb_strip_nulls(jsonb_build_object("
+ "'list_price', vendor_catalog.specs->'list_price', "
+ "'hw_price', vendor_catalog.specs->'hw_price', "
+ "'retail_price', vendor_catalog.specs->'retail_price', "
+ "'cost', vendor_catalog.specs->'cost', "
+ "'image_alt', vendor_catalog.specs->'image_alt'))"
+ )
+ continue
+ if col == "pattern_name":
+ # Keep the private-label name scrub: strip any "| Hollywood Wallcoverings"
+ # vendor suffix the scrape re-supplies (no-op for non-hollywood rows).
+ update_parts.append(
+ " pattern_name = COALESCE(NULLIF(regexp_replace("
+ "EXCLUDED.pattern_name, '\\s*\\|\\s*Hollywood Wallcoverings.*$', '', 'g'), ''), "
+ "vendor_catalog.pattern_name)"
+ )
+ continue
+ if col in no_nullif_cols:
+ update_parts.append(
+ f" {col} = COALESCE(EXCLUDED.{col}, vendor_catalog.{col})"
+ )
+ else:
+ update_parts.append(
+ f" {col} = COALESCE(NULLIF(EXCLUDED.{col}::text, ''), vendor_catalog.{col})"
+ )
+ # Always bump last_scraped_at to the max of the two values.
+ if "last_scraped_at" in mapped_cols:
+ # Remove the generic COALESCE update and replace with GREATEST.
+ update_parts = [
+ p for p in update_parts if not p.strip().startswith("last_scraped_at")
+ ]
+ update_parts.append(
+ " last_scraped_at = GREATEST(EXCLUDED.last_scraped_at, vendor_catalog.last_scraped_at)"
+ )
+ # Always bump updated_at — consolidation touches every row nightly.
+ # NOTE: This means updated_at reflects "last consolidated" not "last changed".
+ # Use last_scraped_at or column-specific checks for true change detection.
+ update_parts.append(" updated_at = NOW()")
+
+ update_clause = ",\n".join(update_parts)
+
+ sql = f"""
+INSERT INTO vendor_catalog ({col_list})
+SELECT {select_list}
+FROM {table_name}
+{where_clause}
+ON CONFLICT (vendor_code, mfr_sku) DO UPDATE SET
+{update_clause}
+""".strip()
+
+ return sql, mapped_cols
+
+
+# ---------------------------------------------------------------------------
+# Per-vendor consolidation
+# ---------------------------------------------------------------------------
+
+def consolidate_vendor(
+ conn: psycopg2.extensions.connection,
+ table_name: str,
+ vendor_code: str,
+ dry_run: bool,
+) -> VendorResult:
+ result = VendorResult(vendor_code=vendor_code, table_name=table_name)
+ t0 = time.monotonic()
+
+ try:
+ with conn.cursor() as cur:
+ source_cols = get_table_columns(cur, table_name)
+
+ if not source_cols:
+ result.error = f"Table {table_name} not found or has no columns"
+ log.warning("[%s] %s", vendor_code, result.error)
+ return result
+
+ built = build_upsert_sql(table_name, vendor_code, source_cols, dry_run)
+ if built is None:
+ result.error = "Could not build upsert SQL (no mfr_sku mapping)"
+ return result
+
+ sql, mapped_cols = built
+ log.debug("[%s] Mapped columns: %s", vendor_code, mapped_cols)
+
+ if dry_run:
+ log.info("[%s] DRY RUN — would execute:\n%s\n", vendor_code, sql)
+ # Estimate eligible rows using the same mfr_sku expression used in
+ # the real query so special tables (arte, momentum, etc.) are counted
+ # correctly.
+ special_mapping = SPECIAL_TABLES.get(table_name, {})
+ mfr_expr = build_source_expr("mfr_sku", table_name, source_cols, special_mapping)
+ count_sql = (
+ f"SELECT COUNT(*) FROM {table_name} "
+ f"WHERE ({mfr_expr}) IS NOT NULL AND ({mfr_expr}) <> ''"
+ )
+ cur.execute(count_sql)
+ row = cur.fetchone()
+ result.inserted = row[0] if row else 0
+ else:
+ cur.execute(sql)
+ rows_affected = cur.rowcount
+ conn.commit()
+
+ # PostgreSQL does not differentiate INSERT vs UPDATE in rowcount for
+ # INSERT ... ON CONFLICT DO UPDATE. We report total rows affected.
+ result.inserted = rows_affected
+ result.updated = 0 # Cannot distinguish without xmax trick
+
+ except Exception as exc: # noqa: BLE001
+ conn.rollback()
+ result.error = str(exc)
+ log.error("[%s] Failed: %s", vendor_code, exc, exc_info=True)
+ finally:
+ result.duration_s = time.monotonic() - t0
+
+ return result
+
+
+# ---------------------------------------------------------------------------
+# Post-consolidation steps
+# ---------------------------------------------------------------------------
+
+def mark_shopify_products(
+ conn: psycopg2.extensions.connection,
+ dry_run: bool,
+) -> int:
+ """Set on_shopify=true AND shopify_product_id for rows matched by dw_sku.
+
+ FIX (2026-04-12): Three bugs in the old version caused 6,903 false positives:
+ 1. Set on_shopify=true WITHOUT setting shopify_product_id (ghost flags)
+ 2. Joined on mfr_sku alone — 5,316 ambiguous SKUs caused cross-vendor contamination
+ 3. Didn't set updated_at, making drift invisible
+
+ Now: join on dw_sku only (unique, safe). Products without dw_sku in
+ shopify_products rely on Quinn's per-product confirmation at push time.
+ """
+ sql = """
+ UPDATE vendor_catalog vc
+ SET on_shopify = true,
+ shopify_product_id = REPLACE(sp.shopify_id, 'gid://shopify/Product/', '')::bigint,
+ updated_at = NOW()
+ FROM shopify_products sp
+ WHERE sp.status = 'ACTIVE'
+ AND sp.dw_sku IS NOT NULL
+ AND sp.dw_sku <> ''
+ AND vc.dw_sku = sp.dw_sku
+ AND (vc.on_shopify IS NOT TRUE OR vc.shopify_product_id IS NULL)
+ """
+ if dry_run:
+ check_sql = """
+ SELECT COUNT(*) FROM vendor_catalog vc
+ JOIN shopify_products sp ON sp.dw_sku = vc.dw_sku
+ WHERE sp.status = 'ACTIVE'
+ AND sp.dw_sku IS NOT NULL AND sp.dw_sku <> ''
+ AND (vc.on_shopify IS NOT TRUE OR vc.shopify_product_id IS NULL)
+ """
+ with conn.cursor() as cur:
+ cur.execute(check_sql)
+ row = cur.fetchone()
+ count = row[0] if row else 0
+ log.info("DRY RUN — would mark ~%d rows as on_shopify=true (dw_sku join)", count)
+ return count
+
+ with conn.cursor() as cur:
+ cur.execute(sql)
+ count = cur.rowcount
+ conn.commit()
+ log.info("Marked %d vendor_catalog rows as on_shopify=true (dw_sku join)", count)
+ return count
+
+
+def mark_never_push(
+ conn: psycopg2.extensions.connection,
+ dry_run: bool,
+) -> int:
+ """Set sync_status='never_push' for cowtan_tout and colefax_fowler."""
+ placeholders = ", ".join(f"'{v}'" for v in sorted(NEVER_PUSH_VENDORS))
+ sql = f"""
+ UPDATE vendor_catalog
+ SET sync_status = 'never_push'
+ WHERE vendor_code IN ({placeholders})
+ AND sync_status IS DISTINCT FROM 'never_push'
+ """
+ if dry_run:
+ check_sql = f"""
+ SELECT COUNT(*) FROM vendor_catalog
+ WHERE vendor_code IN ({placeholders})
+ AND sync_status IS DISTINCT FROM 'never_push'
+ """
+ with conn.cursor() as cur:
+ cur.execute(check_sql)
+ row = cur.fetchone()
+ count = row[0] if row else 0
+ log.info("DRY RUN — would mark ~%d rows as never_push", count)
+ return count
+
+ with conn.cursor() as cur:
+ cur.execute(sql)
+ count = cur.rowcount
+ conn.commit()
+ log.info("Marked %d vendor_catalog rows as never_push", count)
+ return count
+
+
+# ---------------------------------------------------------------------------
+# Orchestration
+# ---------------------------------------------------------------------------
+
+def build_vendor_queue(
+ cur: psycopg2.extensions.cursor,
+ registry: dict[str, dict[str, Any]],
+ all_tables: list[tuple[str, int]],
+ vendor_filter: str | None,
+ shopify_first: bool,
+) -> list[tuple[str, str, int]]:
+ """
+ Return an ordered list of (table_name, vendor_code, row_count) tuples
+ representing the work queue.
+ """
+ # Build a set of all known table names (with rows).
+ table_row_counts = {t: n for t, n in all_tables}
+
+ # Collect candidates.
+ candidates: list[tuple[str, str, int]] = []
+
+ # From vendor_registry: each registry entry uniquely owns its catalog_table.
+ # If multiple vendor_codes map to the same table (e.g. DWDP, GRD -> vendor_catalog),
+ # we skip the target table itself (already in SKIP_TABLES).
+ seen_tables: set[str] = set()
+ for table_name, info in registry.items():
+ if table_name in SKIP_TABLES:
+ continue
+ if table_name not in table_row_counts:
+ continue # table exists in registry but has no rows
+ vc = info["vendor_code"]
+ normalized_vc = vc.lower()
+ if vendor_filter and normalized_vc != vendor_filter:
+ continue
+ # Prefer the first registry entry if the same table appears multiple times.
+ if table_name in seen_tables:
+ log.debug("Table %s already registered; skipping duplicate vendor_code %s", table_name, vc)
+ continue
+ candidates.append((table_name, vc, table_row_counts[table_name]))
+ seen_tables.add(table_name)
+
+ # Tables discovered in pg_stat but not in vendor_registry.
+ for table_name, row_count in all_tables:
+ if table_name in seen_tables or table_name in SKIP_TABLES:
+ continue
+ # Infer vendor_code by stripping _catalog suffix.
+ inferred_vc = table_name.removesuffix("_catalog")
+ if vendor_filter and inferred_vc.lower() != vendor_filter:
+ continue
+ log.info(
+ "Discovered table %s not in vendor_registry; inferring vendor_code=%s",
+ table_name,
+ inferred_vc,
+ )
+ candidates.append((table_name, inferred_vc, row_count))
+ seen_tables.add(table_name)
+
+ if shopify_first:
+ # Get mfr_skus of ACTIVE Shopify products to prioritise tables that have them.
+ shopify_skus = get_shopify_mfr_skus(cur)
+
+ def sort_key(item: tuple[str, str, int]) -> tuple[int, int]:
+ table_name, _, row_count = item
+ # We cannot easily check overlap without a query per table; use row_count
+ # as a proxy for shopify coverage likelihood.
+ return (0 if _has_shopify_overlap(cur, table_name, shopify_skus) else 1, -row_count)
+
+ candidates.sort(key=sort_key)
+ else:
+ # Largest tables first.
+ candidates.sort(key=lambda x: -x[2])
+
+ return candidates
+
+
+def _has_shopify_overlap(
+ cur: psycopg2.extensions.cursor,
+ table_name: str,
+ shopify_skus: set[str],
+) -> bool:
+ """Quick heuristic: sample 50 mfr_skus and check against Shopify set."""
+ try:
+ cols = get_table_columns(cur, table_name)
+ sku_col = "mfr_sku" if "mfr_sku" in cols else ("sku" if "sku" in cols else None)
+ if sku_col is None:
+ return False
+ cur.execute(f"SELECT {sku_col} FROM {table_name} WHERE {sku_col} IS NOT NULL LIMIT 50")
+ sample = {row[0] for row in cur.fetchall()}
+ return bool(sample & shopify_skus)
+ except Exception: # noqa: BLE001
+ return False
+
+
+def run(args: argparse.Namespace) -> ConsolidationSummary:
+ summary = ConsolidationSummary()
+ t_start = time.monotonic()
+
+ log.info("=" * 70)
+ log.info(
+ "Consolidation started at %s | dry_run=%s vendor=%s shopify_first=%s",
+ datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
+ args.dry_run,
+ args.vendor or "ALL",
+ args.shopify_first,
+ )
+ log.info("=" * 70)
+
+ try:
+ conn = get_connection()
+ except Exception as exc:
+ log.critical("Cannot connect to database: %s", exc)
+ sys.exit(1)
+
+ try:
+ with conn.cursor() as cur:
+ registry = get_vendor_registry(cur)
+ all_tables = get_all_catalog_tables(cur)
+ queue = build_vendor_queue(
+ cur, registry, all_tables, args.vendor, args.shopify_first
+ )
+ except Exception as exc:
+ log.critical("Failed to build vendor queue: %s", exc)
+ conn.close()
+ sys.exit(1)
+
+ if not queue:
+ log.warning("No vendors matched the given filter — nothing to do.")
+ conn.close()
+ return summary
+
+ log.info("Vendors to process: %d", len(queue))
+
+ for table_name, vendor_code, row_count in queue:
+ log.info(
+ "-> [%s] table=%s rows~%d",
+ vendor_code,
+ table_name,
+ row_count,
+ )
+ result = consolidate_vendor(conn, table_name, vendor_code, args.dry_run)
+ summary.results.append(result)
+
+ status = "ERROR" if result.error else "OK"
+ log.info(
+ " [%s] %s | affected=%d | %.2fs%s",
+ vendor_code,
+ status,
+ result.inserted,
+ result.duration_s,
+ f" | error={result.error}" if result.error else "",
+ )
+
+ # Post-consolidation steps (skip if single-vendor run).
+ if not args.vendor:
+ log.info("Post-consolidation: marking Shopify products ...")
+ summary.shopify_marked = mark_shopify_products(conn, args.dry_run)
+
+ log.info("Post-consolidation: marking never_push vendors ...")
+ summary.never_push_marked = mark_never_push(conn, args.dry_run)
+
+ conn.close()
+ summary.total_duration_s = time.monotonic() - t_start
+ return summary
+
+
+# ---------------------------------------------------------------------------
+# Reporting
+# ---------------------------------------------------------------------------
+
+def print_summary(summary: ConsolidationSummary, dry_run: bool) -> None:
+ mode = "DRY RUN" if dry_run else "LIVE"
+ log.info("")
+ log.info("=" * 70)
+ log.info("CONSOLIDATION SUMMARY [%s]", mode)
+ log.info("=" * 70)
+
+ successes = [r for r in summary.results if not r.error]
+ failures = [r for r in summary.results if r.error]
+
+ total_rows = sum(r.inserted for r in successes)
+ log.info("Vendors processed : %d", len(summary.results))
+ log.info("Succeeded : %d", len(successes))
+ log.info("Failed : %d", len(failures))
+ log.info("Total rows upserted: %d", total_rows)
+ log.info("Shopify marked : %d", summary.shopify_marked)
+ log.info("Never-push marked : %d", summary.never_push_marked)
+ log.info("Total duration : %.1fs", summary.total_duration_s)
+
+ if failures:
+ log.info("")
+ log.info("FAILURES:")
+ for r in failures:
+ log.info(" %-25s [%s] %s", r.vendor_code, r.table_name, r.error)
+
+ if successes:
+ log.info("")
+ log.info("TOP VENDORS BY ROWS AFFECTED:")
+ for r in sorted(successes, key=lambda x: -x.inserted)[:20]:
+ log.info(
+ " %-25s [%s] %d rows in %.2fs",
+ r.vendor_code,
+ r.table_name,
+ r.inserted,
+ r.duration_s,
+ )
+
+ log.info("=" * 70)
+
+
+# ---------------------------------------------------------------------------
+# CLI
+# ---------------------------------------------------------------------------
+
+def parse_args() -> argparse.Namespace:
+ parser = argparse.ArgumentParser(
+ description="Consolidate vendor catalog tables into vendor_catalog.",
+ formatter_class=argparse.RawDescriptionHelpFormatter,
+ epilog=__doc__,
+ )
+ group = parser.add_mutually_exclusive_group(required=True)
+ group.add_argument(
+ "--all",
+ action="store_true",
+ help="Process all vendor catalog tables.",
+ )
+ group.add_argument(
+ "--vendor",
+ metavar="CODE",
+ help="Process a single vendor by vendor_code.",
+ )
+ parser.add_argument(
+ "--dry-run",
+ action="store_true",
+ help="Build and log SQL but do not execute writes.",
+ )
+ parser.add_argument(
+ "--shopify-first",
+ action="store_true",
+ help="Process tables with Shopify ACTIVE products before others.",
+ )
+ return parser.parse_args()
+
+
+def main() -> None:
+ args = parse_args()
+
+ # Normalise: lowercase vendor filter, None when --all was used.
+ if args.vendor:
+ args.vendor = args.vendor.lower().strip()
+ else:
+ args.vendor = None
+
+ summary = run(args)
+ print_summary(summary, args.dry_run)
+
+ # Exit non-zero if any vendor failed.
+ n_failed = sum(1 for r in summary.results if r.error)
+ sys.exit(1 if n_failed else 0)
+
+
+if __name__ == "__main__":
+ main()
diff --git a/thib-price.mjs b/thib-price.mjs
new file mode 100644
index 0000000..2139e2b
--- /dev/null
+++ b/thib-price.mjs
@@ -0,0 +1,25 @@
+import { chromium } from 'playwright';
+const b=await chromium.launch({channel:'chrome',headless:true});
+const p=await (await b.newContext({viewport:{width:1400,height:1000},userAgent:'Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15) AppleWebKit/537.36 Chrome/124 Safari/537.36'})).newPage();
+const tries=['https://www.thibautdesign.com/search?q=T75099','https://www.thibautdesign.com/search?q=crocodile'];
+for(const u of tries){
+ await p.goto(u,{waitUntil:'domcontentloaded',timeout:45000}).catch(()=>{});
+ await p.waitForTimeout(5000);
+ const r=await p.evaluate(()=>{
+ const txt=document.body.innerText;
+ const prods=[...document.querySelectorAll('a[href*="/products/"],a[href*="/product/"]')].map(a=>a.getAttribute('href')).filter((v,i,a)=>a.indexOf(v)===i).slice(0,6);
+ const prices=[...txt.matchAll(/\$\s?\d[\d,]*\.\d{2}/g)].map(m=>m[0]).slice(0,8);
+ const croc=/crocodile|crowley/i.test(txt);
+ return {url:location.href,prods,prices,croc};
+ });
+ console.log(JSON.stringify(r));
+ if(r.prods.length){ // open first product, read price
+ const u2=r.prods[0].startsWith('http')?r.prods[0]:'https://www.thibautdesign.com'+r.prods[0];
+ await p.goto(u2,{waitUntil:'domcontentloaded',timeout:45000}).catch(()=>{});
+ await p.waitForTimeout(5000);
+ const d=await p.evaluate(()=>{const t=document.body.innerText;const m=[...t.matchAll(/\$\s?\d[\d,]*\.\d{2}\s*(\/\s*(roll|yard|yd))?/gi)].map(x=>x[0]).slice(0,6);return {url:location.href,title:document.title,prices:m,sku:(t.match(/T7509[0-9]|T7510[0-9]/)||[])[0]};});
+ console.log("PRODUCT:",JSON.stringify(d));
+ break;
+ }
+}
+await b.close();
← e06a701 Innovations cost map (224 colorways, list×0.90 per-yd); 182
·
back to Hollywood Import
·
Phase-1 Hollywood create pipeline: sample-safe inline-2-vari ae20a4f →