← back to Reid Witlin Onboarding

verify_onboarding.py

105 lines

#!/usr/bin/env python3
"""Read-only E2E verifier for the Reid Witlin onboarding batch."""
import collections, csv, json, os, subprocess
import create2

HERE = os.path.dirname(os.path.abspath(__file__))
rows = list(csv.DictReader(open(os.path.join(HERE, "targets_ready_v2.csv"))))
VERIFY = """
query($query: String!) {
  products(first: 250, query: $query) {
    nodes {
      id handle title status vendor tags
      variants(first: 5) { nodes { sku price } }
      media(first: 2) { nodes { mediaContentType status } }
    }
    pageInfo { hasNextPage endCursor }
  }
}"""

def scalar(sql):
    r = subprocess.run(["psql", "host=/tmp dbname=dw_unified", "-tA", "-c", sql],
                       capture_output=True, text=True)
    if r.returncode:
        raise RuntimeError(r.stderr.strip())
    return r.stdout.strip()

def fetch_products():
    found, after = [], None
    query = "vendor:'Architectural Fabrics'"
    while True:
        q = VERIFY.replace("products(first: 250, query: $query)",
                           "products(first: 250, after: %s, query: $query)" %
                           ("null" if after is None else json.dumps(after)))
        data = create2.gql(q, {"query": query})
        if data.get("errors"):
            raise RuntimeError(data["errors"])
        conn = (data.get("data") or {}).get("products") or {}
        found.extend(conn.get("nodes") or [])
        page = conn.get("pageInfo") or {}
        if not page.get("hasNextPage"):
            return found
        after = page.get("endCursor")

products = fetch_products()
sku_counts = collections.Counter()
sku_products = collections.defaultdict(list)
handle_counts = collections.Counter(p.get("handle") for p in products if p.get("handle"))
for p in products:
    for v in (p.get("variants") or {}).get("nodes") or []:
        if v.get("sku"):
            sku_counts[v["sku"]] += 1
            sku_products[v["sku"]].append({"id": p.get("id"), "handle": p.get("handle"), "status": p.get("status")})
duplicate_handles = sorted(k for k, n in handle_counts.items() if n > 1)
duplicate_skus = sorted(k for k, n in sku_counts.items() if n > 1)
archive_plan = []
for variant_sku in duplicate_skus:
    base_sku = variant_sku.removesuffix("-Sample").replace("'", "''")
    canonical_pid = scalar(
        f"SELECT COALESCE(MAX(shopify_product_id),'') FROM rwltd_catalog WHERE dw_sku='{base_sku}';")
    canonical_gid = canonical_pid if canonical_pid.startswith("gid://") else (
        f"gid://shopify/Product/{canonical_pid}" if canonical_pid else "")
    candidates = sku_products[variant_sku]
    extras = [p for p in candidates if p.get("id") != canonical_gid]
    archive_plan.append({"sku": variant_sku, "retain": canonical_gid, "archive": extras})
expected = {r["sku"] + "-Sample": r for r in rows}
matched = []
for p in products:
    for v in (p.get("variants") or {}).get("nodes") or []:
        if v.get("sku") in expected:
            matched.append((p, v, expected[v["sku"]]))

failures = []
for p, v, row in matched:
    checks = {
        "status": p.get("status") == "DRAFT",
        "vendor": p.get("vendor") == "Architectural Fabrics",
        "price": v.get("price") == "4.25",
        "sku": v.get("sku") == row["sku"] + "-Sample",
        "tags": {"quotes", "Commercial"}.issubset(set(p.get("tags") or [])),
        "media": bool((p.get("media") or {}).get("nodes")),
    }
    if not all(checks.values()):
        failures.append({"sku": row["sku"], "checks": checks})

catalog_count = int(scalar("""
  SELECT COUNT(DISTINCT mfr_sku) FROM rwltd_catalog
  WHERE created_at::date >= '2026-09-01' AND COALESCE(shopify_product_id,'')<>'' AND COALESCE(dw_sku,'')<>'';
""") or 0)
registry_count = int(scalar("""
  SELECT COUNT(*) FROM dw_sku_registry
  WHERE vendor_prefix='DWKR' AND status='draft' AND COALESCE(shopify_product_id,'')<>'';
""") or 0)
out = {
    "expected_batch": len(expected), "shopify_matched": len(matched),
    "failures": failures[:20], "catalog_linked_since_rescrape": catalog_count,
    "registry_draft_linked": registry_count,
    "duplicate_handles": duplicate_handles,
    "duplicate_skus": duplicate_skus,
    "duplicate_sku_products": {k: sku_products[k] for k in duplicate_skus},
    "archive_plan": archive_plan,
}
print(json.dumps(out, indent=2))
if failures or duplicate_handles or duplicate_skus:
    raise SystemExit(1)