From 138126ab17bc1ce44162c588b8bb22b66f23ef74 Mon Sep 17 00:00:00 2001 From: Ben Stull Date: Thu, 11 Jun 2026 15:42:42 -0700 Subject: [PATCH] =?UTF-8?q?feat(products):=20repo=20=E2=80=94=20catalog=20?= =?UTF-8?q?snapshot,=20draft/run=20SQL,=20upsert=20primitives?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Fable 5 --- backend/app/domains/products/repo.py | 479 +++++++++++++++++++++++++++ 1 file changed, 479 insertions(+) create mode 100644 backend/app/domains/products/repo.py diff --git a/backend/app/domains/products/repo.py b/backend/app/domains/products/repo.py new file mode 100644 index 0000000..bcdd1fb --- /dev/null +++ b/backend/app/domains/products/repo.py @@ -0,0 +1,479 @@ +"""products repo — the SQL layer for the import spine (SD-0002 §6.3 data model). + +Owns SQL only: the catalog snapshot the diff engine reads, import draft/run +CRUD, and the apply primitives the confirm transaction calls. Business rules +live in service.py and diff.py — nothing here validates, diffs, commits, or +rolls back (the confirm flow runs the apply primitives inside its own +transaction). Every catalog/draft/run query is storefront-scoped (INV-14). + +Dict payload conventions: functions feeding §6.4 API payloads (insert_draft, +list_runs, get_run) return datetimes as `.isoformat()` strings; get_draft_row +returns raw datetimes for the service's expiry check. TEXT[] columns bind/load +as Python lists and NUMERIC loads as Decimal natively under psycopg 3. +""" +from __future__ import annotations + +import psycopg +from psycopg import sql +from psycopg.types.json import Jsonb + +from .diff import CatalogImage, CatalogProduct, CatalogVariant + +# --------------------------------------------------------------------------- +# Catalog snapshot (diff input) + dashboard counts +# --------------------------------------------------------------------------- + + +def load_catalog(conn: psycopg.Connection, storefront_id: int) -> dict[str, CatalogProduct]: + """The storefront's full catalog, keyed by handle, in diff.py's snapshot shape.""" + catalog: dict[str, CatalogProduct] = {} + by_id: dict[int, CatalogProduct] = {} + for row in conn.execute( + "SELECT id, handle, title, description_html, vendor, product_type," + " google_product_category, tags, status, published," + " option1_name, option2_name, option3_name" + " FROM product WHERE storefront_id = %s", + (storefront_id,), + ): + product = CatalogProduct( + id=row[0], + handle=row[1], + title=row[2], + option_names=(row[10], row[11], row[12]), + fields={ + "title": row[2], + "description_html": row[3], + "vendor": row[4], + "product_type": row[5], + "google_product_category": row[6], + "tags": row[7], + "status": row[8], + "published": row[9], + }, + variants=[], + images=[], + ) + catalog[product.handle] = product + by_id[product.id] = product + for row in conn.execute( + "SELECT v.product_id, v.id, v.position," + " v.option1_value, v.option2_value, v.option3_value," + " v.sku, v.barcode, v.price, v.cost, v.weight, v.weight_unit," + " v.volume, v.volume_unit, v.tax_id_1, v.tax_id_2," + " v.inventory_tracker, v.inventory_qty, i.source_url" + " FROM variant v" + " JOIN product p ON p.id = v.product_id" + " LEFT JOIN product_image i ON i.id = v.image_id" + " WHERE p.storefront_id = %s" + " ORDER BY v.product_id, v.position, v.id", + (storefront_id,), + ): + by_id[row[0]].variants.append( + CatalogVariant( + id=row[1], + options=(row[3], row[4], row[5]), + position=row[2], + fields={ + "sku": row[6], + "barcode": row[7], + "price": row[8], + "cost": row[9], + "weight": row[10], + "weight_unit": row[11], + "volume": row[12], + "volume_unit": row[13], + "tax_id_1": row[14], + "tax_id_2": row[15], + "inventory_tracker": row[16], + "inventory_qty": row[17], + "variant_image": row[18], + }, + ) + ) + for row in conn.execute( + "SELECT i.product_id, i.id, i.source_url, i.position, i.alt_text" + " FROM product_image i" + " JOIN product p ON p.id = i.product_id" + " WHERE p.storefront_id = %s" + " ORDER BY i.product_id, i.position, i.id", + (storefront_id,), + ): + by_id[row[0]].images.append( + CatalogImage(id=row[1], source_url=row[2], position=row[3], alt_text=row[4]) + ) + return catalog + + +def product_count(conn: psycopg.Connection, storefront_id: int) -> int: + return conn.execute( + "SELECT count(*) FROM product WHERE storefront_id = %s", (storefront_id,) + ).fetchone()[0] + + +def image_problem_count(conn: psycopg.Connection, storefront_id: int) -> int: + return conn.execute( + "SELECT count(*) FROM product_image i" + " JOIN product p ON p.id = i.product_id" + " WHERE p.storefront_id = %s" + " AND i.status IN ('rejected_low_res', 'rejected_not_image', 'failed')", + (storefront_id,), + ).fetchone()[0] + + +def latest_run_id(conn: psycopg.Connection, storefront_id: int) -> int | None: + row = conn.execute( + "SELECT id FROM import_run WHERE storefront_id = %s" + " ORDER BY created_at DESC, id DESC LIMIT 1", + (storefront_id,), + ).fetchone() + return row[0] if row else None + + +# --------------------------------------------------------------------------- +# Import drafts (preview server side, INV-11) +# --------------------------------------------------------------------------- + + +def insert_draft( + conn: psycopg.Connection, + storefront_id: int, + account_id: int, + file_name: str, + dialect: str, + file_bytes: bytes, + summary: dict, + records: list, + fingerprint: str, + unknown_columns: list[str], +) -> dict: + """Create a draft (expires in 1 hour); returns the §6.4 draft payload.""" + row = conn.execute( + "INSERT INTO import_draft" + " (storefront_id, account_id, file_name, dialect, file_bytes," + " summary, records, fingerprint, unknown_columns, expires_at)" + " VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, now() + interval '1 hour')" + " RETURNING id, expires_at", + ( + storefront_id, + account_id, + file_name, + dialect, + file_bytes, + Jsonb(summary), + Jsonb(records), + fingerprint, + unknown_columns, + ), + ).fetchone() + return { + "id": row[0], + "file_name": file_name, + "dialect": dialect, + "summary": summary, + "unknown_columns": unknown_columns, + "expires_at": row[1].isoformat(), + } + + +def get_draft_row(conn: psycopg.Connection, storefront_id: int, draft_id: int) -> dict | None: + row = conn.execute( + "SELECT id, storefront_id, account_id, file_name, dialect, file_bytes," + " summary, records, fingerprint, unknown_columns, expires_at, created_at" + " FROM import_draft WHERE id = %s AND storefront_id = %s", + (draft_id, storefront_id), + ).fetchone() + if row is None: + return None + columns = ( + "id", + "storefront_id", + "account_id", + "file_name", + "dialect", + "file_bytes", + "summary", + "records", + "fingerprint", + "unknown_columns", + "expires_at", + "created_at", + ) + record = dict(zip(columns, row)) + # BYTEA loads as memoryview; the service expects bytes. + record["file_bytes"] = bytes(record["file_bytes"]) + return record + + +def draft_records( + conn: psycopg.Connection, + storefront_id: int, + draft_id: int, + kind: str | None, + limit: int, + offset: int, +) -> list[dict]: + """The draft's preview records, order-preserving, optionally filtered by kind.""" + rows = conn.execute( + "SELECT rec FROM import_draft d," + " jsonb_array_elements(d.records) WITH ORDINALITY AS r(rec, ord)" + " WHERE d.id = %(draft_id)s AND d.storefront_id = %(storefront_id)s" + " AND (%(kind)s::text IS NULL OR rec->>'kind' = %(kind)s)" + " ORDER BY ord LIMIT %(limit)s OFFSET %(offset)s", + { + "draft_id": draft_id, + "storefront_id": storefront_id, + "kind": kind, + "limit": limit, + "offset": offset, + }, + ).fetchall() + return [row[0] for row in rows] + + +def delete_draft(conn: psycopg.Connection, storefront_id: int, draft_id: int) -> None: + conn.execute( + "DELETE FROM import_draft WHERE id = %s AND storefront_id = %s", + (draft_id, storefront_id), + ) + + +def sweep_expired_drafts(conn: psycopg.Connection) -> None: + conn.execute("DELETE FROM import_draft WHERE expires_at < now()") + + +# --------------------------------------------------------------------------- +# Import runs (history, PUC-8) +# --------------------------------------------------------------------------- + +_TERMINAL_RUN_STATUSES = ("complete", "complete_with_problems") + + +def insert_run( + conn: psycopg.Connection, + storefront_id: int, + account_id: int, + file_name: str, + dialect: str, + added: int, + updated: int, + errored: int, + status: str, +) -> int: + return conn.execute( + "INSERT INTO import_run" + " (storefront_id, account_id, file_name, dialect," + " products_added, products_updated, rows_errored, status, completed_at)" + " VALUES (%s, %s, %s, %s, %s, %s, %s, %s, CASE WHEN %s THEN now() END)" + " RETURNING id", + ( + storefront_id, + account_id, + file_name, + dialect, + added, + updated, + errored, + status, + status in _TERMINAL_RUN_STATUSES, + ), + ).fetchone()[0] + + +def insert_run_errors(conn: psycopg.Connection, run_id: int, errors: list[dict]) -> None: + """Record per-row errors (RowError.as_json shape: line/column/message).""" + if not errors: + return + with conn.cursor() as cur: + cur.executemany( + "INSERT INTO import_run_error (run_id, line_number, column_name, message)" + " VALUES (%s, %s, %s, %s)", + [(run_id, e["line"], e["column"], e["message"]) for e in errors], + ) + + +_RUN_SELECT = ( + "SELECT r.id, r.file_name, r.dialect, r.created_at, r.completed_at, r.status," + " a.email, r.products_added, r.products_updated, r.rows_errored" + " FROM import_run r JOIN account a ON a.id = r.account_id" +) + + +def _run_dict(row: tuple) -> dict: + return { + "id": row[0], + "file_name": row[1], + "dialect": row[2], + "created_at": row[3].isoformat(), + "completed_at": row[4].isoformat() if row[4] is not None else None, + "status": row[5], + "by": row[6], + "products_added": row[7], + "products_updated": row[8], + "rows_errored": row[9], + } + + +def list_runs( + conn: psycopg.Connection, storefront_id: int, limit: int, offset: int +) -> list[dict]: + rows = conn.execute( + _RUN_SELECT + + " WHERE r.storefront_id = %s ORDER BY r.created_at DESC, r.id DESC" + " LIMIT %s OFFSET %s", + (storefront_id, limit, offset), + ).fetchall() + return [_run_dict(row) for row in rows] + + +def get_run(conn: psycopg.Connection, storefront_id: int, run_id: int) -> dict | None: + row = conn.execute( + _RUN_SELECT + " WHERE r.storefront_id = %s AND r.id = %s", + (storefront_id, run_id), + ).fetchone() + if row is None: + return None + run = _run_dict(row) + run["errors"] = [ + {"line": line, "column": column, "message": message} + for line, column, message in conn.execute( + "SELECT line_number, column_name, message FROM import_run_error" + " WHERE run_id = %s ORDER BY line_number, id", + (run_id,), + ) + ] + # SLICE-7 fills these; the §6.4 payload shape is stable from SLICE-5 on. + run["image_progress"] = {"done": 0, "total": 0} + run["image_outcomes"] = [] + return run + + +# --------------------------------------------------------------------------- +# Apply primitives — called inside the confirm transaction (Task 8); no commits. +# --------------------------------------------------------------------------- + + +def insert_product( + conn: psycopg.Connection, + storefront_id: int, + handle: str, + resolved_fields: dict, + option_names: tuple[str | None, str | None, str | None], +) -> int: + """INSERT with only the file-present fields; absent ones take column defaults.""" + columns = ["storefront_id", "handle", "option1_name", "option2_name", "option3_name"] + values: list[object] = [storefront_id, handle, *option_names] + for field_name, value in resolved_fields.items(): + columns.append(field_name) + values.append(value) + query = sql.SQL("INSERT INTO product ({}) VALUES ({}) RETURNING id").format( + sql.SQL(", ").join(sql.Identifier(c) for c in columns), + sql.SQL(", ").join(sql.Placeholder() for _ in columns), + ) + return conn.execute(query, values).fetchone()[0] + + +def update_product(conn: psycopg.Connection, product_id: int, changed_fields: dict) -> None: + if not changed_fields: + return + assignments = [ + sql.SQL("{} = {}").format(sql.Identifier(f), sql.Placeholder()) + for f in changed_fields + ] + query = sql.SQL("UPDATE product SET {}, updated_at = now() WHERE id = {}").format( + sql.SQL(", ").join(assignments), sql.Placeholder() + ) + conn.execute(query, [*changed_fields.values(), product_id]) + + +def insert_variant( + conn: psycopg.Connection, + product_id: int, + position: int, + options: tuple[str | None, str | None, str | None], + resolved_fields: dict, + image_id: int | None, +) -> int: + # variant_image is not a column — the caller translates it to image_id; position + # is the explicit param. Filter both defensively. + fields = { + k: v for k, v in resolved_fields.items() if k != "variant_image" and k != "position" + } + columns = [ + "product_id", + "position", + "option1_value", + "option2_value", + "option3_value", + "image_id", + ] + values: list[object] = [product_id, position, *options, image_id] + for field_name, value in fields.items(): + columns.append(field_name) + values.append(value) + query = sql.SQL("INSERT INTO variant ({}) VALUES ({}) RETURNING id").format( + sql.SQL(", ").join(sql.Identifier(c) for c in columns), + sql.SQL(", ").join(sql.Placeholder() for _ in columns), + ) + return conn.execute(query, values).fetchone()[0] + + +def update_variant( + conn: psycopg.Connection, + variant_id: int, + changed_fields: dict, + image_id: int | None | type(...) = ..., +) -> None: + """Dynamic UPDATE; image_id's Ellipsis default means "don't touch image_id".""" + fields = { + k: v for k, v in changed_fields.items() if k != "variant_image" + } + assignments = [ + sql.SQL("{} = {}").format(sql.Identifier(f), sql.Placeholder()) for f in fields + ] + values: list[object] = list(fields.values()) + if image_id is not ...: + assignments.append(sql.SQL("image_id = {}").format(sql.Placeholder())) + values.append(image_id) + if not assignments: + return + query = sql.SQL("UPDATE variant SET {}, updated_at = now() WHERE id = {}").format( + sql.SQL(", ").join(assignments), sql.Placeholder() + ) + conn.execute(query, [*values, variant_id]) + + +def get_or_create_image( + conn: psycopg.Connection, + product_id: int, + source_url: str, + position: int, + alt_text: str | None, + run_id: int, +) -> int: + """Image identity within a product is source_url (§6.3); existing rows are + returned untouched — diff emits explicit image update entries for position/alt.""" + row = conn.execute( + "SELECT id FROM product_image WHERE product_id = %s AND source_url = %s", + (product_id, source_url), + ).fetchone() + if row is not None: + return row[0] + return conn.execute( + "INSERT INTO product_image (product_id, source_url, position, alt_text, import_run_id)" + " VALUES (%s, %s, %s, %s, %s) RETURNING id", + (product_id, source_url, position, alt_text, run_id), + ).fetchone()[0] + + +def update_image(conn: psycopg.Connection, image_id: int, changed_fields: dict) -> None: + """Subset of {position, alt_text}.""" + if not changed_fields: + return + assignments = [ + sql.SQL("{} = {}").format(sql.Identifier(f), sql.Placeholder()) + for f in changed_fields + ] + query = sql.SQL("UPDATE product_image SET {} WHERE id = {}").format( + sql.SQL(", ").join(assignments), sql.Placeholder() + ) + conn.execute(query, [*changed_fields.values(), image_id])