Compare commits

..

12 Commits

Author SHA1 Message Date
ben.stull a8538d3ecc fix(products): cleared Variant Position resolves to file order, not NULL
A blank Variant Position cell previewed position->null and aborted the
confirm transaction (variant.position is NOT NULL). Resolve the clear to
the variant's file order at diff time so preview and apply stay in
lockstep; carry file_order on VariantPlan instead of an identity-keyed
map; release the read snapshot before PreviewStale/NothingToApply.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-06-11 16:14:14 -07:00
ben.stull 0c7865e9e1 feat(products): confirm/apply in one transaction, runs, summary — INV-10/11/14 + TEL-2/6
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-06-11 16:02:27 -07:00
ben.stull fcbf1393f5 feat(products): import_validate → draft, preview records, discard + TEL-1
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-06-11 15:51:03 -07:00
ben.stull 138126ab17 feat(products): repo — catalog snapshot, draft/run SQL, upsert primitives
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-06-11 15:42:42 -07:00
ben.stull 667a462e0b fix(products): honest position compare for matched variants in diff
A Variant Position column previously produced a spurious update with a
false before:None (CatalogVariant keeps position as an attribute, not in
fields{}). Compare against the attribute, like images already do.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-06-11 15:37:59 -07:00
ben.stull 9bc6e4dbd2 feat(products): diff engine — add/update/unchanged/error + INV-11 fingerprint
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-06-11 15:28:03 -07:00
ben.stull 19ee695c20 feat(products): row validation — §6.5.1 rules + INV-15 sanitization
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-06-11 15:17:01 -07:00
ben.stull f64c3fddf9 feat(products): CSV codec — file-level parse + INV-18 caps (PUC-5a)
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-06-11 14:58:53 -07:00
ben.stull c39bbd4728 feat(products): domain skeleton — errors, column registry, canonical row model
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-06-11 14:55:37 -07:00
ben.stull 0ee948b34d feat(products): SD-0002 §6.3 schema — migration 0002 + nh3/multipart deps (SLICE-5)
Also updates test_migrate_from_empty_applies_0001 → test_migrate_from_empty_applies_all
to assert both migrations apply on a fresh DB (mechanical consequence of adding 0002).

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-06-11 14:52:18 -07:00
ben.stull 1267d4f29d docs: SLICE-5 implementation plan (SD-0002 §7.2 import spine)
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-06-11 14:48:47 -07:00
ben.stull 6b405a0f2d Merge pull request 'fix(signin): key the resend cooldown to the address it was set for (#20)' (#22) from fix-20-signin-cooldown-keying into main
ci / check (push) Has been cancelled
2026-06-11 08:22:08 +00:00
18 changed files with 4725 additions and 2 deletions
+36
View File
@@ -0,0 +1,36 @@
"""products domain — catalog + bulk CSV import/export (SD-0002 §6.2).
Owns the canonical row model, codec, validation, diff engine, and import
drafts/runs. Storefront-scoped throughout (INV-14); upsert is the only mutation
(INV-10). Imported via this package surface only.
"""
from __future__ import annotations
from .errors import (
DraftExpired,
DraftNotFound,
FileRejected,
NothingToApply,
PreviewStale,
ProductsError,
RunNotFound,
)
from .models import MAX_DATA_ROWS, MAX_FILE_BYTES
from .service import (
confirm_draft,
discard_draft,
get_draft,
get_draft_records,
get_run,
import_validate,
list_runs,
summary,
)
__all__ = [
"ProductsError", "FileRejected", "DraftNotFound", "DraftExpired",
"PreviewStale", "NothingToApply", "RunNotFound",
"MAX_DATA_ROWS", "MAX_FILE_BYTES",
"import_validate", "get_draft", "get_draft_records", "discard_draft",
"confirm_draft", "list_runs", "get_run", "summary",
]
+65
View File
@@ -0,0 +1,65 @@
"""CSV codec — bytes → ParsedFile (SD-0002 §6.5.1). File-level gates only (PUC-5a);
row semantics live in validate.py. Dialect detection is the INV-17 seam (SLICE-8
adds Shopify)."""
from __future__ import annotations
import csv
import io
from .errors import FileRejected
from .models import KNOWN_COLUMNS, MAX_DATA_ROWS, MAX_FILE_BYTES, ParsedFile, Row
_REQUIRED_HEADER_COLUMNS = ("Handle", "Title")
def detect_dialect(header: list[str]) -> str:
"""The INV-17 seam: SLICE-8 recognizes Shopify's exact header set here."""
return "canonical"
def parse_csv(data: bytes) -> ParsedFile:
if len(data) > MAX_FILE_BYTES:
raise FileRejected("file_too_large", "This file is larger than 10 MB.")
try:
text = data.decode("utf-8-sig")
except UnicodeDecodeError:
raise FileRejected("not_csv", "This file isn't readable as CSV.") from None
reader = csv.reader(io.StringIO(text))
try:
try:
raw_header = next(reader)
except StopIteration:
raise FileRejected("not_csv", "This file isn't readable as CSV.") from None
header = [h.strip() for h in raw_header]
for col in _REQUIRED_HEADER_COLUMNS:
if col not in header:
raise FileRejected(
"missing_required_column",
f"This file is missing the required column '{col}'.",
)
# First occurrence of a duplicated column wins.
col_index: dict[str, int] = {}
for i, name in enumerate(header):
if name and name not in col_index:
col_index[name] = i
known_present = [c for c in col_index if c in KNOWN_COLUMNS]
unknown = [c for c in col_index if c not in KNOWN_COLUMNS]
rows: list[Row] = []
for raw in reader:
if not any(cell.strip() for cell in raw):
continue
if len(rows) >= MAX_DATA_ROWS:
raise FileRejected(
"too_many_rows",
f"This file has more than {MAX_DATA_ROWS:,} rows — split it and import in parts.",
)
cells = {
c: (raw[col_index[c]].strip() if col_index[c] < len(raw) else "")
for c in known_present
}
rows.append(Row(line_number=reader.line_num, cells=cells))
except csv.Error:
raise FileRejected("not_csv", "This file isn't readable as CSV.") from None
return ParsedFile(
dialect=detect_dialect(header), header=header, unknown_columns=unknown, rows=rows
)
+322
View File
@@ -0,0 +1,322 @@
"""Diff engine — catalog × canonical products → apply plan + preview records (SD-0002 §6.5.2).
Classifies each canonical product against the storefront's current catalog as
add / update / unchanged / error. One walk produces two views of the same
computation: a typed apply *plan* carrying resolved native values (Decimal etc.)
for the confirm transaction, and JSON-ready preview records derived from that
walk (stored as draft JSONB, served verbatim to the SPA) — so what confirm
applies is exactly what preview showed (INV-11). A summary and a deterministic
fingerprint over the records detect catalog drift between preview and confirm.
Deliberately DB-free: the catalog snapshot dataclasses are defined here and the
repo layer builds them. Only fields present in the file's canonical fields{}
participate in a comparison — an absent column is untouched, never a change
(§6.5.1); a present-but-empty cell resolves to the field's CLEAR_DEFAULTS entry
— except a variant's position, whose default is the variant's 1-based file
order within its product ("defaults to file order"), resolved here at diff time.
Catalog variants/images absent from the file are likewise untouched (INV-10);
file variants match catalog variants by their option-value combination (INV-13).
"""
from __future__ import annotations
import hashlib
import json
from dataclasses import dataclass, field
from decimal import Decimal
from .models import CLEAR_DEFAULTS, CanonicalProduct, CanonicalVariant
@dataclass
class CatalogVariant:
id: int
options: tuple[str | None, str | None, str | None]
position: int
# sku, barcode, price (Decimal|None), cost, weight, weight_unit, volume,
# volume_unit, tax_id_1, tax_id_2, inventory_tracker, inventory_qty,
# variant_image (linked image source_url or None)
fields: dict[str, object]
@dataclass
class CatalogImage:
id: int
source_url: str
position: int
alt_text: str | None
@dataclass
class CatalogProduct:
id: int
handle: str
title: str
option_names: tuple[str | None, str | None, str | None]
# title, description_html, vendor, product_type, google_product_category,
# tags (list[str]), status, published (bool)
fields: dict[str, object]
variants: list[CatalogVariant]
images: list[CatalogImage]
# ---------------------------------------------------------------------------
# Apply plan — the typed twin of the preview records. confirm_draft executes
# these; the values are resolved natives (Decimal, bool, list), never the
# json-safe strings the records carry for display.
# ---------------------------------------------------------------------------
@dataclass
class VariantPlan:
kind: str # "add" | "update"
canonical: CanonicalVariant
catalog_id: int | None # None for add
# 1-based index within the product's file variants — the position default.
file_order: int = 0
# field -> resolved after-value (update only); adds resolve from canonical.fields
changes: dict[str, object] = field(default_factory=dict)
@dataclass
class ImagePlan:
kind: str # "add" | "update"
source_url: str
position: int
alt_text: str | None
image_id: int | None # None for add
changes: dict[str, object] = field(default_factory=dict) # subset of {"position","alt_text"}
@dataclass
class ProductPlan:
kind: str # "add" | "update" | "unchanged" | "error"
canonical: CanonicalProduct
catalog: CatalogProduct | None
# field -> resolved after-value (update only; may include title/option*_name)
product_changes: dict[str, object] = field(default_factory=dict)
variant_plans: list[VariantPlan] = field(default_factory=list)
image_plans: list[ImagePlan] = field(default_factory=list)
@dataclass(frozen=True)
class DiffResult:
records: list[dict]
summary: dict
fingerprint: str
plan: list[ProductPlan]
# Option-name presence markers in fields{} (see validate.py); the values are
# compared via the option_names attribute, not through the generic field loop.
_OPTION_NAME_FIELDS = ("option1_name", "option2_name", "option3_name")
_SUMMARY_KEY = {"add": "adds", "update": "updates", "unchanged": "unchanged", "error": "errors"}
def compute_diff(catalog: dict[str, CatalogProduct], products: list[CanonicalProduct]) -> DiffResult:
records: list[dict] = []
plan: list[ProductPlan] = []
summary = {"adds": 0, "updates": 0, "unchanged": 0, "errors": 0}
for product in products:
current = catalog.get(product.handle)
if product.errors:
product_plan = ProductPlan(kind="error", canonical=product, catalog=current)
detail: dict = {"errors": [e.as_json() for e in product.errors]}
elif current is None:
product_plan = _add_plan(product)
detail = _add_detail(product_plan)
else:
product_plan, detail = _update_plan(product, current)
kind = product_plan.kind
plan.append(product_plan)
records.append(
{
"handle": product.handle,
"title": product.title or (current.title if current else ""),
"kind": kind,
"variant_count": len(current.variants) if kind == "unchanged" else len(product.variants),
"detail": detail,
}
)
summary[_SUMMARY_KEY[kind]] += 1
fingerprint = hashlib.sha256(
json.dumps(records, sort_keys=True, separators=(",", ":")).encode()
).hexdigest()
return DiffResult(records=records, summary=summary, fingerprint=fingerprint, plan=plan)
def _add_plan(product: CanonicalProduct) -> ProductPlan:
return ProductPlan(
kind="add",
canonical=product,
catalog=None,
variant_plans=[
VariantPlan(kind="add", canonical=v, catalog_id=None, file_order=order)
for order, v in enumerate(product.variants, start=1)
],
image_plans=[
ImagePlan(kind="add", source_url=i.source_url, position=i.position,
alt_text=i.alt_text, image_id=None)
for i in product.images
],
)
def _add_detail(plan: ProductPlan) -> dict:
# On add, absent fields fall back to their defaults for display where one
# exists — the detail shows what will actually be set.
set_fields = resolved_product_fields(plan.canonical)
for field_name, default in CLEAR_DEFAULTS.items():
set_fields.setdefault(field_name, default)
return {
"set": {f: _json_safe(v) for f, v in set_fields.items()},
"option_names": list(plan.canonical.option_names),
"variants": [
{"options": list(vp.canonical.options),
"set": _resolved_variant_fields(vp.canonical, vp.file_order)}
for vp in plan.variant_plans
],
"images": [
{"src": ip.source_url, "position": ip.position, "alt_text": ip.alt_text}
for ip in plan.image_plans
],
}
def _update_plan(product: CanonicalProduct, current: CatalogProduct) -> tuple[ProductPlan, dict]:
"""One walk, two outputs: the resolved-value plan entries and the json-safe
record detail entries are appended side by side, so they can never diverge."""
product_changes: dict[str, object] = {}
changes: list[dict] = []
# Title is always file-present (required header column); "" means the block
# already carries an error and never reaches here.
if product.title and product.title != current.title:
product_changes["title"] = product.title
changes.append(_change("title", current.title, product.title))
for field_name, resolved in resolved_product_fields(product).items():
before = current.fields.get(field_name)
if resolved != before:
product_changes[field_name] = resolved
changes.append(_change(field_name, before, resolved))
for slot in (1, 2, 3):
if f"option{slot}_name" not in product.fields:
continue
before, after = current.option_names[slot - 1], product.option_names[slot - 1]
if after != before:
product_changes[f"option{slot}_name"] = after
changes.append(_change(f"option{slot}_name", before, after))
variant_plans: list[VariantPlan] = []
variant_entries: list[dict] = []
by_options = {v.options: v for v in current.variants}
for file_order, variant in enumerate(product.variants, start=1):
match = by_options.get(variant.options)
if match is None:
variant_plans.append(
VariantPlan(kind="add", canonical=variant, catalog_id=None, file_order=file_order)
)
variant_entries.append(
{"options": list(variant.options), "kind": "add",
"set": _resolved_variant_fields(variant, file_order)}
)
continue
variant_changes: dict[str, object] = {}
variant_change_entries: list[dict] = []
for f, resolved in resolved_variant_fields(variant, file_order).items():
# position lives on the catalog variant as an attribute, not in
# fields{} — compare it explicitly (as images do for theirs).
before = match.position if f == "position" else match.fields.get(f)
if resolved != before:
variant_changes[f] = resolved
variant_change_entries.append(_change(f, before, resolved))
if variant_changes:
variant_plans.append(
VariantPlan(kind="update", canonical=variant, catalog_id=match.id,
file_order=file_order, changes=variant_changes)
)
variant_entries.append(
{"options": list(variant.options), "kind": "update", "changes": variant_change_entries}
)
image_plans: list[ImagePlan] = []
image_entries: list[dict] = []
by_src = {i.source_url: i for i in current.images}
for image in product.images:
match = by_src.get(image.source_url)
if match is None:
image_plans.append(
ImagePlan(kind="add", source_url=image.source_url, position=image.position,
alt_text=image.alt_text, image_id=None)
)
image_entries.append(
{"src": image.source_url, "kind": "add", "position": image.position, "alt_text": image.alt_text}
)
continue
image_changes: dict[str, object] = {}
image_change_entries: list[dict] = []
for f, before, after in (
("position", match.position, image.position),
("alt_text", match.alt_text, image.alt_text),
):
if after != before:
image_changes[f] = after
image_change_entries.append(_change(f, before, after))
if image_changes:
image_plans.append(
ImagePlan(kind="update", source_url=image.source_url, position=image.position,
alt_text=image.alt_text, image_id=match.id, changes=image_changes)
)
image_entries.append(
{"src": image.source_url, "kind": "update", "changes": image_change_entries}
)
detail: dict = {}
if changes:
detail["changes"] = changes
if variant_entries:
detail["variants"] = variant_entries
if image_entries:
detail["images"] = image_entries
kind = "update" if detail else "unchanged"
return (
ProductPlan(kind=kind, canonical=product, catalog=current, product_changes=product_changes,
variant_plans=variant_plans, image_plans=image_plans),
detail,
)
def resolved_fields(fields: dict[str, object]) -> dict[str, object]:
"""File-present fields with None (an explicit clear) resolved to the default."""
return {f: (v if v is not None else CLEAR_DEFAULTS.get(f)) for f, v in fields.items()}
def resolved_product_fields(product: CanonicalProduct) -> dict[str, object]:
"""Resolved product-level fields, minus the option-name presence markers."""
return {
f: v for f, v in resolved_fields(product.fields).items() if f not in _OPTION_NAME_FIELDS
}
def resolved_variant_fields(variant: CanonicalVariant, file_order: int) -> dict[str, object]:
"""File-present variant fields with clears resolved. A cleared position has
no CLEAR_DEFAULTS entry — it resets to the variant's 1-based file order
within its product ("defaults to file order"), never to NULL."""
resolved = resolved_fields(variant.fields)
if "position" in resolved and resolved["position"] is None:
resolved["position"] = file_order
return resolved
def _resolved_variant_fields(variant: CanonicalVariant, file_order: int) -> dict[str, object]:
return {f: _json_safe(v) for f, v in resolved_variant_fields(variant, file_order).items()}
def _change(field_name: str, before: object, after: object) -> dict:
return {"field": field_name, "before": _json_safe(before), "after": _json_safe(after)}
def _json_safe(value: object) -> object:
if isinstance(value, Decimal):
return str(value)
if isinstance(value, (tuple, list)):
return [_json_safe(v) for v in value]
return value
+37
View File
@@ -0,0 +1,37 @@
"""products domain errors (SD-0002 §6.4 error envelope codes)."""
from __future__ import annotations
class ProductsError(Exception):
"""Base for products-domain errors."""
class FileRejected(ProductsError):
"""PUC-5a: the whole file is unusable; no draft is created. `code` is the §6.4
error code (not_csv | missing_required_column | unknown_dialect | too_many_rows |
file_too_large)."""
def __init__(self, code: str, message: str):
super().__init__(message)
self.code = code
self.message = message
class DraftNotFound(ProductsError):
"""No such draft for this storefront (or already discarded)."""
class DraftExpired(ProductsError):
"""The draft's validity window passed (§6.3 ~1 h)."""
class PreviewStale(ProductsError):
"""INV-11: the catalog changed since validation — the previewed diff no longer holds."""
class NothingToApply(ProductsError):
"""PUC-10: no adds and no updates — confirming would be a no-op."""
class RunNotFound(ProductsError):
"""No such import run for this storefront."""
+121
View File
@@ -0,0 +1,121 @@
"""Canonical row model + column registry — the one model every dialect maps to (INV-17)."""
from __future__ import annotations
from dataclasses import dataclass, field
# §6.5.1 canonical columns, by level. Header detection, unknown-column warnings, and
# validation all read from this registry.
PRODUCT_COLUMNS: dict[str, str] = {
# column -> product field name
"Title": "title",
"Description": "description_html",
"Vendor": "vendor",
"Type": "product_type",
"Google Product Category": "google_product_category",
"Tags": "tags",
"Status": "status",
"Published": "published",
"Option1 Name": "option1_name",
"Option2 Name": "option2_name",
"Option3 Name": "option3_name",
}
VARIANT_COLUMNS: dict[str, str] = {
"Variant SKU": "sku",
"Variant Barcode": "barcode",
"Variant Price": "price",
"Variant Cost": "cost",
"Variant Weight": "weight",
"Variant Weight Unit": "weight_unit",
"Variant Volume": "volume",
"Variant Volume Unit": "volume_unit",
"Variant Tax ID 1": "tax_id_1",
"Variant Tax ID 2": "tax_id_2",
"Variant Inventory Tracker": "inventory_tracker",
"Variant Inventory Qty": "inventory_qty",
"Variant Position": "position",
"Variant Image": "variant_image",
}
OPTION_VALUE_COLUMNS = ("Option1 Value", "Option2 Value", "Option3 Value")
IMAGE_COLUMNS = ("Image Src", "Image Position", "Image Alt Text")
COMPONENT_COLUMNS = tuple(
f"Component {i} {kind}" for i in range(1, 11) for kind in ("SKU", "Quantity")
)
KNOWN_COLUMNS = (
{"Handle"}
| set(PRODUCT_COLUMNS)
| set(VARIANT_COLUMNS)
| set(OPTION_VALUE_COLUMNS)
| set(IMAGE_COLUMNS)
| set(COMPONENT_COLUMNS)
)
# Clearing a field (present-but-empty cell, §6.5.1) resets it to its default.
CLEAR_DEFAULTS: dict[str, object] = {
"status": "active",
"published": True,
"product_type": "standalone",
"tags": [],
}
MAX_DATA_ROWS = 5_000 # INV-18
MAX_FILE_BYTES = 10 * 1024 * 1024 # INV-18
@dataclass(frozen=True)
class Row:
"""One CSV data row: 1-based file line number + the cells of known columns
present in the header (column name -> raw string, possibly empty)."""
line_number: int
cells: dict[str, str]
@dataclass(frozen=True)
class ParsedFile:
dialect: str
header: list[str]
unknown_columns: list[str]
rows: list[Row]
@dataclass(frozen=True)
class RowError:
line_number: int
column: str | None
message: str
def as_json(self) -> dict:
return {"line": self.line_number, "column": self.column, "message": self.message}
@dataclass
class CanonicalVariant:
line_number: int
options: tuple[str | None, str | None, str | None]
# field name -> normalized value; present only for columns in the file.
# value None == clear (reset to default/NULL).
fields: dict[str, object] = field(default_factory=dict)
@dataclass
class CanonicalImage:
line_number: int
source_url: str
position: int
alt_text: str | None
@dataclass
class CanonicalProduct:
first_line: int
handle: str
title: str # "" when missing (the block then carries an error)
option_names: tuple[str | None, str | None, str | None] = (None, None, None)
fields: dict[str, object] = field(default_factory=dict) # product-level, same semantics
variants: list[CanonicalVariant] = field(default_factory=list)
images: list[CanonicalImage] = field(default_factory=list)
errors: list[RowError] = field(default_factory=list)
@property
def valid(self) -> bool:
return not self.errors
+479
View File
@@ -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])
+227
View File
@@ -0,0 +1,227 @@
"""products service — the import/export use-case orchestration (SD-0002 §6.5).
Coordinates codec → validate → diff → repo; owns transaction boundaries (repo
never commits). Preview is read-only against catalog tables (INV-11): validation
writes exactly one row — the import_draft. TEL events per §9.1.
"""
from __future__ import annotations
import time
from datetime import datetime, timezone
import psycopg
from app.platform import telemetry
from . import codec, diff, repo, validate
from .errors import DraftExpired, DraftNotFound, NothingToApply, PreviewStale, RunNotFound
def import_validate(conn: psycopg.Connection, storefront_id: int, account_id: int,
file_name: str, data: bytes) -> dict:
"""Upload → validate → diff → persist draft (PUC-2/3; INV-11). Raises FileRejected."""
started = time.monotonic()
# Commit the sweep before parsing: a FileRejected mid-parse must not roll
# back expired-draft cleanup along with it.
repo.sweep_expired_drafts(conn)
conn.commit()
parsed = codec.parse_csv(data)
products = validate.build_products(parsed)
catalog = repo.load_catalog(conn, storefront_id)
diff_result = diff.compute_diff(catalog, products)
draft = repo.insert_draft(
conn, storefront_id, account_id, file_name, parsed.dialect, data,
diff_result.summary, diff_result.records, diff_result.fingerprint,
parsed.unknown_columns,
)
conn.commit()
telemetry.emit(
"import_draft_created",
storefront_id=storefront_id,
dialect=parsed.dialect,
row_count=len(parsed.rows),
adds=diff_result.summary["adds"],
updates=diff_result.summary["updates"],
unchanged=diff_result.summary["unchanged"],
errors=diff_result.summary["errors"],
unknown_columns_count=len(parsed.unknown_columns),
duration_ms=int((time.monotonic() - started) * 1000),
)
return draft
def _live_draft_row(conn: psycopg.Connection, storefront_id: int, draft_id: int) -> dict:
"""The draft row if it exists and hasn't expired; expiry deletes lazily (§6.3)."""
row = repo.get_draft_row(conn, storefront_id, draft_id)
if row is None:
raise DraftNotFound()
if row["expires_at"] < datetime.now(timezone.utc):
repo.delete_draft(conn, storefront_id, draft_id)
conn.commit()
raise DraftExpired()
return row
def get_draft(conn: psycopg.Connection, storefront_id: int, draft_id: int) -> dict:
"""The §6.4 draft payload — never file_bytes or the full records list."""
row = _live_draft_row(conn, storefront_id, draft_id)
return {
"id": row["id"],
"file_name": row["file_name"],
"dialect": row["dialect"],
"summary": row["summary"],
"unknown_columns": row["unknown_columns"],
"expires_at": row["expires_at"].isoformat(),
}
def get_draft_records(conn: psycopg.Connection, storefront_id: int, draft_id: int,
kind: str | None = None, limit: int = 100, offset: int = 0) -> list[dict]:
"""The draft's preview records, paged, optionally filtered by kind (PUC-3)."""
_live_draft_row(conn, storefront_id, draft_id)
return repo.draft_records(conn, storefront_id, draft_id, kind, limit, offset)
def discard_draft(conn: psycopg.Connection, storefront_id: int, draft_id: int) -> None:
"""Delete the draft, no trace kept; idempotent — an absent draft is fine (PUC-3a)."""
repo.delete_draft(conn, storefront_id, draft_id)
conn.commit()
def confirm_draft(conn: psycopg.Connection, storefront_id: int, account_id: int,
draft_id: int) -> int:
"""Apply the previewed diff in one transaction (PUC-4; INV-10/11).
Everything is re-derived from the draft's stored file bytes against the live
catalog; a fingerprint mismatch means the catalog drifted since preview
(PreviewStale — the draft is kept so the merchant can re-validate). The apply
executes the typed plan compute_diff built alongside the preview records, so
what lands is exactly what the preview showed. rows_errored counts the
import_run_error rows recorded (one per RowError), which is what the run
detail's error table shows; the preview's errors tile counts error *products*.
"""
started = time.monotonic()
row = _live_draft_row(conn, storefront_id, draft_id)
parsed = codec.parse_csv(row["file_bytes"])
products = validate.build_products(parsed)
catalog = repo.load_catalog(conn, storefront_id)
diff_result = diff.compute_diff(catalog, products)
if diff_result.fingerprint != row["fingerprint"]:
# Release the read snapshot; nothing written.
conn.rollback()
raise PreviewStale()
summary_counts = diff_result.summary
if summary_counts["adds"] + summary_counts["updates"] == 0:
# Release the read snapshot; nothing written.
conn.rollback()
raise NothingToApply()
error_rows = [
error.as_json()
for plan in diff_result.plan if plan.kind == "error"
for error in plan.canonical.errors
]
try:
run_id = repo.insert_run(
conn, storefront_id, account_id, row["file_name"], row["dialect"],
added=summary_counts["adds"], updated=summary_counts["updates"],
errored=len(error_rows), status="complete",
)
for plan in diff_result.plan:
_apply_product_plan(conn, storefront_id, plan, run_id)
repo.insert_run_errors(conn, run_id, error_rows)
repo.delete_draft(conn, storefront_id, draft_id)
conn.commit()
except Exception as exc:
conn.rollback()
telemetry.emit(
"import_apply_failed",
draft_id=draft_id,
storefront_id=storefront_id,
error_class=type(exc).__name__,
)
raise
telemetry.emit(
"import_run_completed",
run_id=run_id,
storefront_id=storefront_id,
added=summary_counts["adds"],
updated=summary_counts["updates"],
errored=len(error_rows),
duration_ms=int((time.monotonic() - started) * 1000),
)
return run_id
def _apply_product_plan(conn: psycopg.Connection, storefront_id: int,
plan: diff.ProductPlan, run_id: int) -> None:
"""Execute one product's plan inside the confirm transaction (no commits here)."""
if plan.kind == "add":
# Title is a canonical attribute, not a fields{} entry — non-error
# products always carry one (validate guarantees it).
product_fields = {"title": plan.canonical.title}
product_fields.update(diff.resolved_product_fields(plan.canonical))
product_id = repo.insert_product(
conn, storefront_id, plan.canonical.handle, product_fields, plan.canonical.option_names
)
image_ids: dict[str, int] = {}
elif plan.kind == "update":
product_id = plan.catalog.id
repo.update_product(conn, product_id, plan.product_changes)
image_ids = {image.source_url: image.id for image in plan.catalog.images}
else:
return
# Images first, so variants' variant_image URLs resolve to ids: validate puts
# every variant_image URL into canonical.images, so each URL is in either the
# catalog map (existing image) or the adds below.
for image_plan in plan.image_plans:
if image_plan.kind == "add":
image_ids[image_plan.source_url] = repo.get_or_create_image(
conn, product_id, image_plan.source_url, image_plan.position,
image_plan.alt_text, run_id,
)
else:
repo.update_image(conn, image_plan.image_id, image_plan.changes)
for variant_plan in plan.variant_plans:
if variant_plan.kind == "add":
fields = diff.resolved_variant_fields(variant_plan.canonical, variant_plan.file_order)
# diff time resolved any cleared position to file order; the
# file_order fallback covers an absent position column.
position = fields.get("position") or variant_plan.file_order
url = fields.get("variant_image")
image_id = image_ids[url] if url else None
repo.insert_variant(
conn, product_id, position, variant_plan.canonical.options, fields, image_id
)
elif "variant_image" in variant_plan.changes:
url = variant_plan.changes["variant_image"]
repo.update_variant(
conn, variant_plan.catalog_id, variant_plan.changes,
image_id=image_ids[url] if url else None,
)
else:
repo.update_variant(conn, variant_plan.catalog_id, variant_plan.changes)
def list_runs(conn: psycopg.Connection, storefront_id: int,
limit: int = 50, offset: int = 0) -> list[dict]:
"""The storefront's import history, newest first (PUC-8)."""
return repo.list_runs(conn, storefront_id, limit, offset)
def get_run(conn: psycopg.Connection, storefront_id: int, run_id: int) -> dict:
"""One run's §6.4 detail payload, errors included."""
run = repo.get_run(conn, storefront_id, run_id)
if run is None:
raise RunNotFound()
return run
def summary(conn: psycopg.Connection, storefront_id: int) -> dict:
"""The products dashboard counts (§6.4)."""
return {
"product_count": repo.product_count(conn, storefront_id),
"image_problem_count": repo.image_problem_count(conn, storefront_id),
"latest_run_id": repo.latest_run_id(conn, storefront_id),
}
+299
View File
@@ -0,0 +1,299 @@
"""Row validation — ParsedFile rows → canonical products + row errors (SD-0002 §6.5.1).
The codec (codec.py) handles file-level gates; this module is the row-semantics
half of the PUC-5 import spine. It groups consecutive rows sharing a Handle into
product blocks (Shopify's grammar), normalizes product/variant/image fields, and
records every rule violation as a merchant-language RowError. Errors never raise:
an error poisons its whole product block (the product previews as kind="error"
and is excluded from apply) while parsing continues so the merchant gets a
complete accounting in one pass (BUC-1a). Description HTML is sanitized with nh3
on the way in (INV-15).
"""
from __future__ import annotations
import re
from decimal import Decimal, InvalidOperation
import nh3
from .models import (
COMPONENT_COLUMNS,
OPTION_VALUE_COLUMNS,
PRODUCT_COLUMNS,
VARIANT_COLUMNS,
CanonicalImage,
CanonicalProduct,
CanonicalVariant,
ParsedFile,
Row,
RowError,
)
_HANDLE_RE = re.compile(r"^[a-z0-9-]+$")
_STATUSES = {"draft", "active", "archived"}
# Title is handled as an attribute, not via fields{}. Option names live in BOTH
# .option_names (the values) and fields{} (the file-presence signal the diff
# needs: absent column == untouched, never a clear).
_ATTRIBUTE_COLUMNS = {"Title"}
# Sentinel: the cell failed normalization (the error is already recorded).
_INVALID = object()
def build_products(parsed: ParsedFile) -> list[CanonicalProduct]:
"""Group rows into product blocks and validate every §6.5.1 rule."""
products: list[CanonicalProduct] = []
closed_handles: set[str] = set()
block: list[Row] = []
def flush() -> None:
nonlocal block
if block:
closed_handles.add(block[0].cells["Handle"])
products.append(_build_block(block))
block = []
for row in parsed.rows:
handle = row.cells.get("Handle", "")
if block and handle == block[0].cells["Handle"]:
block.append(row)
continue
flush()
if not handle:
products.append(
_error_block(row, "(missing)", RowError(row.line_number, "Handle", "a row needs a Handle"))
)
elif not _HANDLE_RE.match(handle):
products.append(
_error_block(
row,
handle,
RowError(
row.line_number,
"Handle",
f"'{handle}' isn't a valid handle — lowercase letters, numbers, and dashes only",
),
)
)
elif handle in closed_handles:
products.append(
_error_block(
row,
handle,
RowError(
row.line_number,
"Handle",
f"rows for '{handle}' must be consecutive — it already appeared earlier in the file",
),
)
)
else:
block = [row]
flush()
return products
def all_errors(products: list[CanonicalProduct]) -> list[RowError]:
return [error for product in products for error in product.errors]
def _error_block(row: Row, handle: str, error: RowError) -> CanonicalProduct:
return CanonicalProduct(first_line=row.line_number, handle=handle, title="", errors=[error])
def _build_block(rows: list[Row]) -> CanonicalProduct:
first = rows[0]
handle = first.cells["Handle"]
errors: list[RowError] = []
title = first.cells.get("Title", "")
if not title:
errors.append(RowError(first.line_number, "Title", f"'{handle}' is missing its Title"))
option_names = tuple(first.cells.get(f"Option{n} Name") or None for n in (1, 2, 3))
has_options = any(option_names)
# Product-level fields come from the first row only; empty cell == clear (None).
fields: dict[str, object] = {}
for column, field_name in PRODUCT_COLUMNS.items():
if column in _ATTRIBUTE_COLUMNS or column not in first.cells:
continue
cell = first.cells[column]
if not cell:
fields[field_name] = None
continue
value = _product_value(first.line_number, column, field_name, cell, errors)
if value is not _INVALID:
fields[field_name] = value
variants: list[CanonicalVariant] = []
images: list[CanonicalImage] = []
seen_combos: set[tuple[str | None, str | None, str | None]] = set()
for index, row in enumerate(rows):
cells = row.cells
line = row.line_number
for column in COMPONENT_COLUMNS:
if cells.get(column):
errors.append(
RowError(line, column, "kits arrive in a coming release — leave the Component columns empty")
)
has_image = bool(cells.get("Image Src"))
if has_image:
_collect_image(row, images, errors)
# A row carries a variant iff any option value / Variant-* cell is filled;
# the first row of a no-option product always carries the single variant.
carries_variant = (
any(cells.get(c) for c in OPTION_VALUE_COLUMNS)
or any(cells.get(c) for c in VARIANT_COLUMNS)
or (index == 0 and not has_options)
)
if carries_variant:
options = tuple(cells.get(c) or None for c in OPTION_VALUE_COLUMNS)
for n in (1, 2, 3):
value, name = options[n - 1], option_names[n - 1]
if value and not name:
errors.append(
RowError(
line,
f"Option{n} Value",
f"Option{n} Value given but the product has no Option{n} Name",
)
)
elif name and not value:
errors.append(
RowError(
line,
f"Option{n} Value",
f"this variant is missing its Option{n} Value ('{name}')",
)
)
if not has_options:
if variants:
errors.append(RowError(line, None, "a product without options can have only one variant"))
elif options in seen_combos:
errors.append(
RowError(line, None, f"duplicate variant — '{handle}' already has a variant with these options")
)
seen_combos.add(options)
variant_fields: dict[str, object] = {}
for column, field_name in VARIANT_COLUMNS.items():
if column not in cells:
continue
cell = cells[column]
if not cell:
variant_fields[field_name] = None
continue
value = _variant_value(line, column, field_name, cell, errors)
if value is _INVALID:
continue
variant_fields[field_name] = value
if field_name == "variant_image" and cell not in {i.source_url for i in images}:
images.append(
CanonicalImage(line_number=line, source_url=cell, position=len(images) + 1, alt_text=None)
)
variants.append(CanonicalVariant(line_number=line, options=options, fields=variant_fields))
elif index > 0 and not has_image:
errors.append(RowError(line, None, "this row has no variant or image data"))
return CanonicalProduct(
first_line=first.line_number,
handle=handle,
title=title,
option_names=option_names,
fields=fields,
variants=variants,
images=images,
errors=errors,
)
def _collect_image(row: Row, images: list[CanonicalImage], errors: list[RowError]) -> None:
cells = row.cells
source_url = cells["Image Src"]
position_cell = cells.get("Image Position", "")
position: int | None = None
if position_cell:
try:
position = int(position_cell)
if position < 1:
raise ValueError
except ValueError:
position = None
errors.append(RowError(row.line_number, "Image Position", f"'{position_cell}' is not a position"))
# Dedupe by source URL within the block — first occurrence wins.
if source_url in {i.source_url for i in images}:
return
images.append(
CanonicalImage(
line_number=row.line_number,
source_url=source_url,
position=position if position is not None else len(images) + 1,
alt_text=cells.get("Image Alt Text") or None,
)
)
def _product_value(line: int, column: str, field_name: str, cell: str, errors: list[RowError]) -> object:
if field_name == "tags":
return [tag.strip() for tag in cell.split(",") if tag.strip()]
if field_name == "status":
status = cell.lower()
if status not in _STATUSES:
errors.append(RowError(line, column, f"'{cell}' is not a status — use draft, active, or archived"))
return _INVALID
return status
if field_name == "published":
flag = cell.upper()
if flag not in ("TRUE", "FALSE"):
errors.append(RowError(line, column, f"'{cell}' is not TRUE or FALSE"))
return _INVALID
return flag == "TRUE"
if field_name == "description_html":
return nh3.clean(cell)
if field_name == "product_type":
if cell != "standalone":
errors.append(RowError(line, column, "kits arrive in a coming release — Type must be 'standalone'"))
return _INVALID
return cell
return cell
def _variant_value(line: int, column: str, field_name: str, cell: str, errors: list[RowError]) -> object:
if field_name in ("price", "cost"):
return _decimal_or_error(line, column, cell, f"'{cell}' is not a price", errors)
if field_name in ("weight", "volume"):
return _decimal_or_error(line, column, cell, f"'{cell}' is not a number", errors)
if field_name == "inventory_qty":
try:
quantity = int(cell)
if quantity < 0:
raise ValueError
except ValueError:
errors.append(RowError(line, column, f"'{cell}' is not a whole number"))
return _INVALID
return quantity
if field_name == "position":
try:
position = int(cell)
if position < 1:
raise ValueError
except ValueError:
errors.append(RowError(line, column, f"'{cell}' is not a position"))
return _INVALID
return position
return cell
def _decimal_or_error(line: int, column: str, cell: str, message: str, errors: list[RowError]) -> object:
try:
value = Decimal(cell)
if not value.is_finite() or value < 0:
raise InvalidOperation
except InvalidOperation:
errors.append(RowError(line, column, message))
return _INVALID
return value
+13
View File
@@ -0,0 +1,13 @@
"""Structured log-event telemetry (SD-0002 §9.1). One JSON object per event on the
`ecomm.telemetry` logger — counts and durations only; never file names, URLs,
catalog content, or secret bytes (§6.3-handbook)."""
from __future__ import annotations
import json
import logging
_logger = logging.getLogger("ecomm.telemetry")
def emit(event: str, **fields: object) -> None:
_logger.info(json.dumps({"event": event, **fields}, sort_keys=True, default=str))
+130
View File
@@ -0,0 +1,130 @@
-- 0002_products.sql — SD-0002 §6.3 data model (SLICE-5). Forward-only (INV-7):
-- never edit once merged; add a new numbered migration.
-- product — one catalog product per (storefront, handle) (INV-13/14). Option *names*
-- live here; product_type's kit values are schema-open but a service-layer rule
-- rejects non-'standalone' until #15 (same pattern as INV-4).
CREATE TABLE product (
id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
storefront_id BIGINT NOT NULL REFERENCES storefront (id),
handle TEXT NOT NULL,
title TEXT NOT NULL,
description_html TEXT,
vendor TEXT,
product_type TEXT NOT NULL DEFAULT 'standalone'
CHECK (product_type IN ('standalone', 'kit_virtual', 'kit_assembled')),
google_product_category TEXT,
tags TEXT[] NOT NULL DEFAULT '{}',
status TEXT NOT NULL DEFAULT 'active' CHECK (status IN ('draft', 'active', 'archived')),
published BOOLEAN NOT NULL DEFAULT TRUE,
option1_name TEXT,
option2_name TEXT,
option3_name TEXT,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
CREATE UNIQUE INDEX product_handle_key ON product (storefront_id, handle); -- INV-13
-- variant — one purchasable form, identified by its option-value combo (INV-13).
-- SKU is indexed data, never identity. image_id FK is added after product_image.
CREATE TABLE variant (
id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
product_id BIGINT NOT NULL REFERENCES product (id),
position INTEGER NOT NULL,
option1_value TEXT,
option2_value TEXT,
option3_value TEXT,
sku TEXT,
barcode TEXT,
price NUMERIC,
cost NUMERIC,
weight NUMERIC,
weight_unit TEXT,
volume NUMERIC,
volume_unit TEXT,
tax_id_1 TEXT,
tax_id_2 TEXT,
inventory_tracker TEXT,
inventory_qty INTEGER,
image_id BIGINT,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
-- Postgres 16 (compose + Cloud SQL pin): NULLS NOT DISTINCT makes the all-NULL
-- no-option combo unique too (INV-13).
CREATE UNIQUE INDEX variant_option_combo_key
ON variant (product_id, option1_value, option2_value, option3_value)
NULLS NOT DISTINCT;
CREATE INDEX variant_sku_idx ON variant (sku);
-- product_image — identity within a product is source_url (§6.3); bytes live in
-- object storage from SLICE-7 (keys nullable until fetched). status starts 'pending';
-- SLICE-5 stubs the fetch phase so rows simply stay pending.
CREATE TABLE product_image (
id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
product_id BIGINT NOT NULL REFERENCES product (id),
position INTEGER NOT NULL,
source_url TEXT NOT NULL,
alt_text TEXT,
status TEXT NOT NULL DEFAULT 'pending'
CHECK (status IN ('pending', 'fetched', 'rejected_low_res', 'rejected_not_image', 'failed')),
failure_reason TEXT,
key_original TEXT,
key_thumb TEXT,
key_card TEXT,
key_detail TEXT,
import_run_id BIGINT,
fetched_at TIMESTAMPTZ,
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
CREATE UNIQUE INDEX product_image_src_key ON product_image (product_id, source_url);
ALTER TABLE variant
ADD CONSTRAINT variant_image_fk FOREIGN KEY (image_id) REFERENCES product_image (id);
-- import_draft — the preview's server side (INV-11). file_bytes is the SLICE-5
-- interim home for the upload (objectstore key from SLICE-7). Deleted outright on
-- cancel/expiry — drafts never appear in history (PUC-3a).
CREATE TABLE import_draft (
id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
storefront_id BIGINT NOT NULL REFERENCES storefront (id),
account_id BIGINT NOT NULL REFERENCES account (id),
file_name TEXT NOT NULL,
dialect TEXT NOT NULL,
file_bytes BYTEA NOT NULL,
summary JSONB NOT NULL,
records JSONB NOT NULL,
fingerprint TEXT NOT NULL,
unknown_columns TEXT[] NOT NULL DEFAULT '{}',
expires_at TIMESTAMPTZ NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
-- import_run — the durable record of one confirmed import (PUC-8); created only at
-- confirm (PUC-4).
CREATE TABLE import_run (
id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
storefront_id BIGINT NOT NULL REFERENCES storefront (id),
account_id BIGINT NOT NULL REFERENCES account (id),
file_name TEXT NOT NULL,
dialect TEXT NOT NULL,
products_added INTEGER NOT NULL,
products_updated INTEGER NOT NULL,
rows_errored INTEGER NOT NULL,
status TEXT NOT NULL
CHECK (status IN ('applying', 'fetching_images', 'complete', 'complete_with_problems')),
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
completed_at TIMESTAMPTZ
);
CREATE INDEX import_run_history_idx ON import_run (storefront_id, created_at DESC);
-- import_run_error — one row per rejected CSV row (PUC-5), merchant-language message.
CREATE TABLE import_run_error (
id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
run_id BIGINT NOT NULL REFERENCES import_run (id),
line_number INTEGER NOT NULL,
column_name TEXT,
message TEXT NOT NULL
);
ALTER TABLE product_image
ADD CONSTRAINT product_image_run_fk FOREIGN KEY (import_run_id) REFERENCES import_run (id);
+2
View File
@@ -5,3 +5,5 @@ psycopg[binary]>=3.1
psycopg-pool>=3.2
pytest>=8.0
import-linter>=2.0
nh3>=0.2
python-multipart>=0.0.9
+24 -2
View File
@@ -1,4 +1,5 @@
import psycopg
import pytest
from app.platform import db
@@ -12,10 +13,10 @@ def _table_names(conn) -> set[str]:
return {r[0] for r in rows}
def test_migrate_from_empty_applies_0001(fresh_db_url):
def test_migrate_from_empty_applies_all(fresh_db_url):
with psycopg.connect(fresh_db_url) as conn:
applied = db.migrate(conn)
assert applied == ["0001_init.sql"]
assert applied == ["0001_init.sql", "0002_products.sql"]
with psycopg.connect(fresh_db_url) as conn:
assert _TABLES.issubset(_table_names(conn))
@@ -41,3 +42,24 @@ def test_membership_has_no_unique_account_constraint(fresh_db_url):
).fetchall()
defs = " ".join(r[0] for r in rows).lower()
assert "unique" not in defs.replace("primary key", "") or "(account_id)" not in defs
def test_0002_products_tables_exist(fresh_db_url):
with psycopg.connect(fresh_db_url) as conn:
db.migrate(conn)
for table in ("product", "variant", "product_image", "import_draft", "import_run", "import_run_error"):
assert conn.execute("SELECT to_regclass(%s)", (f"public.{table}",)).fetchone()[0] == table
def test_0002_variant_option_combo_unique_treats_nulls_as_equal(fresh_db_url):
# INV-13: the no-option product's single variant has NULL option values; a second
# all-NULL combo must collide (NULLS NOT DISTINCT).
with psycopg.connect(fresh_db_url) as conn:
db.migrate(conn)
sf = conn.execute("INSERT INTO storefront (name) VALUES ('s') RETURNING id").fetchone()[0]
pid = conn.execute(
"INSERT INTO product (storefront_id, handle, title) VALUES (%s,'h','T') RETURNING id", (sf,)
).fetchone()[0]
conn.execute("INSERT INTO variant (product_id, position) VALUES (%s, 1)", (pid,))
with pytest.raises(psycopg.errors.UniqueViolation):
conn.execute("INSERT INTO variant (product_id, position) VALUES (%s, 2)", (pid,))
+69
View File
@@ -0,0 +1,69 @@
"""§6.5.1 file-level codec — parse, caps, required columns (PUC-5a fixtures)."""
import pytest
from app.domains.products import FileRejected
from app.domains.products.codec import parse_csv
def _csv(*lines: str) -> bytes:
return ("\n".join(lines) + "\n").encode()
GOOD = _csv(
"Handle,Title,Variant Price,Bogus Column",
"moon-mug,Moon Mug,18.00,x",
"star-tee,Star Tee,24.00,y",
)
def test_parses_header_rows_and_unknown_columns():
parsed = parse_csv(GOOD)
assert parsed.dialect == "canonical"
assert parsed.unknown_columns == ["Bogus Column"]
assert [r.line_number for r in parsed.rows] == [2, 3]
assert parsed.rows[0].cells["Handle"] == "moon-mug"
assert parsed.rows[0].cells["Variant Price"] == "18.00"
assert "Bogus Column" not in parsed.rows[0].cells
def test_bom_tolerated():
parsed = parse_csv(b"\xef\xbb\xbf" + GOOD)
assert parsed.rows[0].cells["Handle"] == "moon-mug"
def test_quoted_cells_rfc4180():
parsed = parse_csv(_csv("Handle,Title,Tags", 'mug,"The ""Best"" Mug","a, b"'))
assert parsed.rows[0].cells["Title"] == 'The "Best" Mug'
assert parsed.rows[0].cells["Tags"] == "a, b"
def test_empty_rows_skipped_short_rows_padded():
parsed = parse_csv(_csv("Handle,Title,Vendor", "mug,Mug", "", ",,", "tee,Tee,Acme"))
assert [r.cells["Handle"] for r in parsed.rows] == ["mug", "tee"]
assert parsed.rows[0].cells["Vendor"] == ""
@pytest.mark.parametrize(
"data,code",
[
(b"\xff\xfe\x00garbage\x00", "not_csv"),
(b"", "not_csv"),
(_csv("Title,Vendor", "Mug,Acme"), "missing_required_column"),
(_csv("Handle,Vendor", "mug,Acme"), "missing_required_column"),
(
_csv("Handle,Title", *(f"h{i},T{i}" for i in range(5001))),
"too_many_rows",
),
(b"Handle,Title\n" + b"x" * (10 * 1024 * 1024), "file_too_large"),
],
)
def test_file_level_rejections(data, code):
with pytest.raises(FileRejected) as exc:
parse_csv(data)
assert exc.value.code == code
def test_missing_column_message_names_the_column():
with pytest.raises(FileRejected) as exc:
parse_csv(_csv("Handle,Vendor", "mug,Acme"))
assert "'Title'" in exc.value.message
+146
View File
@@ -0,0 +1,146 @@
"""Diff engine — classification, blank-vs-absent, option matching, fingerprint (§6.8)."""
from decimal import Decimal
from app.domains.products.codec import parse_csv
from app.domains.products.diff import (
CatalogImage, CatalogProduct, CatalogVariant, compute_diff,
)
from app.domains.products.validate import build_products
def _canon(*lines: str):
return build_products(parse_csv(("\n".join(lines) + "\n").encode()))
def _catalog_mug(**overrides):
fields = {
"title": "Moon Mug", "description_html": None, "vendor": "Acme",
"product_type": "standalone", "google_product_category": None,
"tags": ["kitchen"], "status": "active", "published": True,
} | overrides
return {
"moon-mug": CatalogProduct(
id=1, handle="moon-mug", title=fields["title"],
option_names=(None, None, None), fields=fields,
variants=[CatalogVariant(id=10, options=(None, None, None), position=1,
fields={"sku": "SKU-1", "barcode": None, "price": Decimal("18.00"),
"cost": None, "weight": None, "weight_unit": None,
"volume": None, "volume_unit": None, "tax_id_1": None,
"tax_id_2": None, "inventory_tracker": None,
"inventory_qty": 40, "variant_image": None})],
images=[CatalogImage(id=100, source_url="https://x/a.jpg", position=1, alt_text=None)],
)
}
HEADER = "Handle,Title,Vendor,Tags,Status,Variant SKU,Variant Price,Variant Inventory Qty,Image Src"
MUG_ROW = 'moon-mug,Moon Mug,Acme,kitchen,active,SKU-1,18.00,40,https://x/a.jpg'
def test_new_handle_classifies_add():
diff = compute_diff({}, _canon(HEADER, MUG_ROW))
[rec] = diff.records
assert rec["kind"] == "add" and rec["handle"] == "moon-mug" and rec["variant_count"] == 1
assert diff.summary == {"adds": 1, "updates": 0, "unchanged": 0, "errors": 0}
assert rec["detail"]["set"]["status"] == "active"
def test_identical_file_classifies_unchanged():
diff = compute_diff(_catalog_mug(), _canon(HEADER, MUG_ROW))
assert diff.records[0]["kind"] == "unchanged"
assert diff.summary["unchanged"] == 1
def test_changed_price_classifies_update_with_before_after():
row = MUG_ROW.replace("18.00", "21.00")
diff = compute_diff(_catalog_mug(), _canon(HEADER, row))
[rec] = diff.records
assert rec["kind"] == "update"
[vchange] = rec["detail"]["variants"]
assert {"field": "price", "before": "18.00", "after": "21.00"} in vchange["changes"]
def test_absent_column_untouched_empty_cell_clears():
# Vendor column absent: vendor stays Acme. Status present-but-empty: clears to default 'active' (already active -> no change).
diff = compute_diff(
_catalog_mug(),
_canon("Handle,Title,Status,Variant SKU,Variant Price,Variant Inventory Qty,Image Src",
"moon-mug,Moon Mug,,SKU-1,18.00,40,https://x/a.jpg"),
)
assert diff.records[0]["kind"] == "unchanged"
def test_empty_cell_clear_shows_in_diff():
# Vendor present-but-empty clears Acme -> None: an explicit, previewable change (§6.5.1).
diff = compute_diff(
_catalog_mug(),
_canon("Handle,Title,Vendor,Variant SKU,Variant Price,Variant Inventory Qty,Image Src",
"moon-mug,Moon Mug,,SKU-1,18.00,40,https://x/a.jpg"),
)
[rec] = diff.records
assert rec["kind"] == "update"
assert {"field": "vendor", "before": "Acme", "after": None} in rec["detail"]["changes"]
def test_new_option_combo_is_variant_add_existing_untouched():
catalog = _catalog_mug()
diff = compute_diff(
catalog,
_canon("Handle,Title,Option1 Name,Option1 Value,Variant Price",
"moon-mug,Moon Mug,Size,Large,25.00"),
)
[rec] = diff.records
assert rec["kind"] == "update"
kinds = [v.get("kind") for v in rec["detail"]["variants"]]
assert "add" in kinds
POSITION_HEADER = HEADER + ",Variant Position"
def test_matching_variant_position_column_classifies_unchanged():
# position is an attribute on CatalogVariant (not in fields{}); the compare
# must read it from there, not invent a before:None.
diff = compute_diff(_catalog_mug(), _canon(POSITION_HEADER, MUG_ROW + ",1"))
assert diff.records[0]["kind"] == "unchanged"
def test_changed_variant_position_reports_honest_before():
diff = compute_diff(_catalog_mug(), _canon(POSITION_HEADER, MUG_ROW + ",2"))
[rec] = diff.records
assert rec["kind"] == "update"
[ventry] = rec["detail"]["variants"]
assert {"field": "position", "before": 1, "after": 2} in ventry["changes"]
def test_blank_position_cell_resolves_to_file_order_unchanged():
# A present-but-empty Variant Position cell resets to file order (the spec's
# "defaults to file order"), never to NULL — here file order matches the
# catalog position, so nothing changes.
diff = compute_diff(_catalog_mug(), _canon(POSITION_HEADER, MUG_ROW + ","))
assert diff.records[0]["kind"] == "unchanged"
def test_blank_position_cell_updates_to_file_order():
catalog = _catalog_mug()
catalog["moon-mug"].variants[0].position = 2
diff = compute_diff(catalog, _canon(POSITION_HEADER, MUG_ROW + ","))
[rec] = diff.records
assert rec["kind"] == "update"
[ventry] = rec["detail"]["variants"]
assert {"field": "position", "before": 2, "after": 1} in ventry["changes"]
def test_error_product_classifies_error():
diff = compute_diff({}, _canon("Handle,Title,Variant Price", "mug,Mug,nope"))
[rec] = diff.records
assert rec["kind"] == "error"
assert rec["detail"]["errors"][0]["column"] == "Variant Price"
def test_fingerprint_stable_and_drift_sensitive():
d1 = compute_diff(_catalog_mug(), _canon(HEADER, MUG_ROW))
d2 = compute_diff(_catalog_mug(), _canon(HEADER, MUG_ROW))
d3 = compute_diff(_catalog_mug(title="Renamed"), _canon(HEADER, MUG_ROW))
assert d1.fingerprint == d2.fingerprint
assert d1.fingerprint != d3.fingerprint
+79
View File
@@ -0,0 +1,79 @@
"""SD-0002 invariants: INV-10 (never deletes), INV-14 (two-storefront zero bleed),
apply transactionality (§6.8), TEL-6."""
import json
import logging
import psycopg
import pytest
from app.domains import products
from app.domains.products import repo, service
from app.platform import db
CSV_A = b"Handle,Title,Variant Price\nmug,Mug,10.00\ntee,Tee,20.00\n"
CSV_PARTIAL = b"Handle,Title,Variant Price\nmug,Mug,12.00\n"
@pytest.fixture()
def migrated_conn(fresh_db_url):
with psycopg.connect(fresh_db_url) as conn:
db.migrate(conn)
yield conn
def _merchant(conn, email="m@example.com", shop="Shop"):
acct = conn.execute("INSERT INTO account (email) VALUES (%s) RETURNING id", (email,)).fetchone()[0]
sf = conn.execute("INSERT INTO storefront (name) VALUES (%s) RETURNING id", (shop,)).fetchone()[0]
conn.execute("INSERT INTO storefront_membership (account_id, storefront_id) VALUES (%s,%s)", (acct, sf))
conn.commit()
return acct, sf
def _import(conn, acct, sf, data):
d = products.import_validate(conn, sf, acct, "f.csv", data)
return products.confirm_draft(conn, sf, acct, d["id"])
def test_inv10_partial_file_never_deletes(migrated_conn):
acct, sf = _merchant(migrated_conn)
_import(migrated_conn, acct, sf, CSV_A)
before = migrated_conn.execute("SELECT count(*) FROM product").fetchone()[0]
_import(migrated_conn, acct, sf, CSV_PARTIAL)
after = migrated_conn.execute("SELECT count(*) FROM product").fetchone()[0]
assert after >= before == 2
def test_inv14_two_storefronts_zero_bleed(migrated_conn):
acct1, sf1 = _merchant(migrated_conn)
acct2, sf2 = _merchant(migrated_conn, "n@example.com", "Other")
_import(migrated_conn, acct1, sf1, CSV_A)
assert products.summary(migrated_conn, sf2)["product_count"] == 0
assert products.list_runs(migrated_conn, sf2) == []
_import(migrated_conn, acct2, sf2, CSV_A)
assert products.summary(migrated_conn, sf2)["product_count"] == 2
run1 = products.list_runs(migrated_conn, sf1)[0]
with pytest.raises(products.RunNotFound):
products.get_run(migrated_conn, sf2, run1["id"])
def test_apply_failure_rolls_back_whole_transaction_tel6(migrated_conn, monkeypatch, caplog):
acct, sf = _merchant(migrated_conn)
d = products.import_validate(migrated_conn, sf, acct, "f.csv", CSV_A)
def boom(*a, **k):
raise RuntimeError("mid-apply crash")
monkeypatch.setattr(service.repo, "insert_run_errors", boom)
lg = logging.getLogger("ecomm")
prior = lg.propagate
lg.propagate = True
try:
with caplog.at_level(logging.INFO, logger="ecomm.telemetry"):
with pytest.raises(RuntimeError):
products.confirm_draft(migrated_conn, sf, acct, d["id"])
finally:
lg.propagate = prior
assert migrated_conn.execute("SELECT count(*) FROM product").fetchone()[0] == 0
assert migrated_conn.execute("SELECT count(*) FROM import_run").fetchone()[0] == 0
assert migrated_conn.execute("SELECT count(*) FROM import_draft").fetchone()[0] == 1
events = [json.loads(r.message) for r in caplog.records if r.name == "ecomm.telemetry"]
assert any(e["event"] == "import_apply_failed" and e["error_class"] == "RuntimeError" for e in events)
+222
View File
@@ -0,0 +1,222 @@
"""products service — drafts: validate/preview/discard (PUC-2/3/3a/5a; INV-11)."""
import json
import logging
import psycopg
import pytest
from app.domains import products
from app.platform import db
GOOD_CSV = b"Handle,Title,Vendor,Variant Price\nmoon-mug,Moon Mug,Acme,18.00\nstar-tee,Star Tee,Acme,24.00\n"
@pytest.fixture()
def migrated_conn(fresh_db_url):
with psycopg.connect(fresh_db_url) as conn:
db.migrate(conn)
yield conn
@pytest.fixture()
def merchant(migrated_conn):
acct = migrated_conn.execute(
"INSERT INTO account (email) VALUES ('m@example.com') RETURNING id").fetchone()[0]
sf = migrated_conn.execute(
"INSERT INTO storefront (name) VALUES ('Shop') RETURNING id").fetchone()[0]
migrated_conn.execute(
"INSERT INTO storefront_membership (account_id, storefront_id) VALUES (%s,%s)", (acct, sf))
migrated_conn.commit()
return {"account_id": acct, "storefront_id": sf}
def test_import_validate_creates_draft_with_summary(migrated_conn, merchant):
draft = products.import_validate(
migrated_conn, merchant["storefront_id"], merchant["account_id"], "cat.csv", GOOD_CSV)
assert draft["dialect"] == "canonical"
assert draft["summary"] == {"adds": 2, "updates": 0, "unchanged": 0, "errors": 0}
assert draft["expires_at"]
def test_validate_writes_nothing_to_catalog_inv11(migrated_conn, merchant):
products.import_validate(
migrated_conn, merchant["storefront_id"], merchant["account_id"], "cat.csv", GOOD_CSV)
assert migrated_conn.execute("SELECT count(*) FROM product").fetchone()[0] == 0
assert migrated_conn.execute("SELECT count(*) FROM variant").fetchone()[0] == 0
def test_file_rejection_leaves_no_draft(migrated_conn, merchant):
with pytest.raises(products.FileRejected):
products.import_validate(
migrated_conn, merchant["storefront_id"], merchant["account_id"], "bad.csv",
b"Vendor,Price\nAcme,1\n")
assert migrated_conn.execute("SELECT count(*) FROM import_draft").fetchone()[0] == 0
def test_records_paging_and_kind_filter(migrated_conn, merchant):
draft = products.import_validate(
migrated_conn, merchant["storefront_id"], merchant["account_id"], "cat.csv", GOOD_CSV)
recs = products.get_draft_records(migrated_conn, merchant["storefront_id"], draft["id"])
assert [r["handle"] for r in recs] == ["moon-mug", "star-tee"]
adds = products.get_draft_records(
migrated_conn, merchant["storefront_id"], draft["id"], kind="add", limit=1)
assert len(adds) == 1 and adds[0]["kind"] == "add"
def test_discard_deletes_no_trace_puc3a(migrated_conn, merchant):
draft = products.import_validate(
migrated_conn, merchant["storefront_id"], merchant["account_id"], "cat.csv", GOOD_CSV)
products.discard_draft(migrated_conn, merchant["storefront_id"], draft["id"])
assert migrated_conn.execute("SELECT count(*) FROM import_draft").fetchone()[0] == 0
products.discard_draft(migrated_conn, merchant["storefront_id"], draft["id"]) # idempotent
def test_draft_scoped_to_storefront_inv14(migrated_conn, merchant):
draft = products.import_validate(
migrated_conn, merchant["storefront_id"], merchant["account_id"], "cat.csv", GOOD_CSV)
other_sf = migrated_conn.execute(
"INSERT INTO storefront (name) VALUES ('Other') RETURNING id").fetchone()[0]
migrated_conn.commit()
with pytest.raises(products.DraftNotFound):
products.get_draft(migrated_conn, other_sf, draft["id"])
def test_expired_draft_raises_and_lazily_deletes(migrated_conn, merchant):
draft = products.import_validate(
migrated_conn, merchant["storefront_id"], merchant["account_id"], "cat.csv", GOOD_CSV)
migrated_conn.execute(
"UPDATE import_draft SET expires_at = now() - interval '1 minute' WHERE id = %s",
(draft["id"],))
migrated_conn.commit()
with pytest.raises(products.DraftExpired):
products.get_draft(migrated_conn, merchant["storefront_id"], draft["id"])
assert migrated_conn.execute("SELECT count(*) FROM import_draft").fetchone()[0] == 0
@pytest.fixture()
def telemetry_propagation():
"""create_app() sets propagate=False on the parent "ecomm" logger
(main._ensure_app_logging), which hides ecomm.telemetry records from caplog's
root-logger handler whenever an API test ran first. Restore propagation here."""
lg = logging.getLogger("ecomm")
prior = lg.propagate
lg.propagate = True
yield
lg.propagate = prior
def test_tel1_emitted(migrated_conn, merchant, caplog, telemetry_propagation):
with caplog.at_level(logging.INFO, logger="ecomm.telemetry"):
products.import_validate(
migrated_conn, merchant["storefront_id"], merchant["account_id"], "cat.csv", GOOD_CSV)
events = [json.loads(r.message) for r in caplog.records if r.name == "ecomm.telemetry"]
assert any(
e["event"] == "import_draft_created" and e["adds"] == 2 and e["row_count"] == 2
and "duration_ms" in e and e["unknown_columns_count"] == 0
for e in events
)
UPDATE_CSV = b"Handle,Title,Vendor,Variant Price\nmoon-mug,Moon Mug,Acme,21.00\nstar-tee,Star Tee,Acme,24.00\n"
MIXED_CSV = b"Handle,Title,Variant Price\ngood-mug,Mug,10.00\nbad-tee,Tee,not-a-price\n"
def _validate(conn, m, data=GOOD_CSV):
return products.import_validate(conn, m["storefront_id"], m["account_id"], "cat.csv", data)
def test_confirm_applies_adds_and_records_run(migrated_conn, merchant):
draft = _validate(migrated_conn, merchant)
run_id = products.confirm_draft(
migrated_conn, merchant["storefront_id"], merchant["account_id"], draft["id"])
assert migrated_conn.execute("SELECT count(*) FROM product").fetchone()[0] == 2
run = products.get_run(migrated_conn, merchant["storefront_id"], run_id)
assert run["products_added"] == 2 and run["status"] == "complete"
assert run["by"] == "m@example.com"
assert run["image_progress"] == {"done": 0, "total": 0} and run["image_outcomes"] == []
assert migrated_conn.execute("SELECT count(*) FROM import_draft").fetchone()[0] == 0
def test_confirm_update_changes_only_diffed_fields(migrated_conn, merchant):
d1 = _validate(migrated_conn, merchant)
products.confirm_draft(migrated_conn, merchant["storefront_id"], merchant["account_id"], d1["id"])
d2 = _validate(migrated_conn, merchant, UPDATE_CSV)
assert d2["summary"] == {"adds": 0, "updates": 1, "unchanged": 1, "errors": 0}
products.confirm_draft(migrated_conn, merchant["storefront_id"], merchant["account_id"], d2["id"])
price = migrated_conn.execute(
"SELECT v.price FROM variant v JOIN product p ON p.id = v.product_id WHERE p.handle='moon-mug'"
).fetchone()[0]
assert str(price) == "21.00"
assert migrated_conn.execute("SELECT count(*) FROM product").fetchone()[0] == 2 # no dupes (BUC-3)
def test_confirm_blank_position_cell_round_trips(migrated_conn, merchant):
# A present-but-empty Variant Position cell resolves to file order at diff
# time — never SET position = NULL (which would abort the confirm on the
# NOT NULL constraint).
d1 = _validate(migrated_conn, merchant)
products.confirm_draft(migrated_conn, merchant["storefront_id"], merchant["account_id"], d1["id"])
blank_position_csv = (
b"Handle,Title,Vendor,Variant Price,Variant Position\n"
b"moon-mug,Moon Mug,Acme,21.00,\n"
b"star-tee,Star Tee,Acme,24.00,\n"
)
d2 = _validate(migrated_conn, merchant, blank_position_csv)
products.confirm_draft(migrated_conn, merchant["storefront_id"], merchant["account_id"], d2["id"])
price = migrated_conn.execute(
"SELECT v.price FROM variant v JOIN product p ON p.id = v.product_id WHERE p.handle='moon-mug'"
).fetchone()[0]
assert str(price) == "21.00"
def test_confirm_mixed_applies_valid_records_errors(migrated_conn, merchant):
draft = _validate(migrated_conn, merchant, MIXED_CSV)
run_id = products.confirm_draft(
migrated_conn, merchant["storefront_id"], merchant["account_id"], draft["id"])
assert migrated_conn.execute("SELECT count(*) FROM product").fetchone()[0] == 1
run = products.get_run(migrated_conn, merchant["storefront_id"], run_id)
assert run["rows_errored"] == 1
assert run["errors"][0]["column"] == "Variant Price"
def test_confirm_stale_fingerprint_409_inv11(migrated_conn, merchant):
draft = _validate(migrated_conn, merchant)
other = _validate(migrated_conn, merchant)
products.confirm_draft(migrated_conn, merchant["storefront_id"], merchant["account_id"], other["id"])
with pytest.raises(products.PreviewStale):
products.confirm_draft(migrated_conn, merchant["storefront_id"], merchant["account_id"], draft["id"])
def test_confirm_nothing_to_apply_puc10(migrated_conn, merchant):
d1 = _validate(migrated_conn, merchant)
products.confirm_draft(migrated_conn, merchant["storefront_id"], merchant["account_id"], d1["id"])
d2 = _validate(migrated_conn, merchant)
assert d2["summary"]["unchanged"] == 2
with pytest.raises(products.NothingToApply):
products.confirm_draft(migrated_conn, merchant["storefront_id"], merchant["account_id"], d2["id"])
def test_runs_history_newest_first(migrated_conn, merchant):
d1 = _validate(migrated_conn, merchant)
products.confirm_draft(migrated_conn, merchant["storefront_id"], merchant["account_id"], d1["id"])
d2 = _validate(migrated_conn, merchant, UPDATE_CSV)
r2 = products.confirm_draft(migrated_conn, merchant["storefront_id"], merchant["account_id"], d2["id"])
runs = products.list_runs(migrated_conn, merchant["storefront_id"])
assert [r["id"] for r in runs][0] == r2
def test_summary_counts(migrated_conn, merchant):
assert products.summary(migrated_conn, merchant["storefront_id"]) == {
"product_count": 0, "image_problem_count": 0, "latest_run_id": None}
d = _validate(migrated_conn, merchant)
rid = products.confirm_draft(migrated_conn, merchant["storefront_id"], merchant["account_id"], d["id"])
s = products.summary(migrated_conn, merchant["storefront_id"])
assert s == {"product_count": 2, "image_problem_count": 0, "latest_run_id": rid}
def test_tel2_emitted_on_confirm(migrated_conn, merchant, caplog, telemetry_propagation):
draft = _validate(migrated_conn, merchant)
with caplog.at_level(logging.INFO, logger="ecomm.telemetry"):
products.confirm_draft(migrated_conn, merchant["storefront_id"], merchant["account_id"], draft["id"])
events = [json.loads(r.message) for r in caplog.records if r.name == "ecomm.telemetry"]
assert any(e["event"] == "import_run_completed" and e["added"] == 2 for e in events)
+164
View File
@@ -0,0 +1,164 @@
"""§6.5.1 row validation — every row-error rule has a fixture (SD-0002 §6.8)."""
from decimal import Decimal
import pytest
from app.domains.products.codec import parse_csv
from app.domains.products.validate import build_products
def _products(*lines: str):
return build_products(parse_csv(("\n".join(lines) + "\n").encode()))
def _errors(*lines: str):
return [e for p in _products(*lines) for e in p.errors]
def test_simple_product_parses_clean():
[p] = _products(
"Handle,Title,Vendor,Tags,Status,Published,Variant Price,Variant SKU",
"moon-mug,Moon Mug,Acme,\"kitchen, mugs\",active,TRUE,18.00,SKU-1",
)
assert p.valid and p.handle == "moon-mug" and p.title == "Moon Mug"
assert p.fields["tags"] == ["kitchen", "mugs"]
assert p.fields["status"] == "active" and p.fields["published"] is True
[v] = p.variants
assert v.options == (None, None, None)
assert v.fields["price"] == Decimal("18.00") and v.fields["sku"] == "SKU-1"
def test_option_product_groups_consecutive_rows():
[p] = _products(
"Handle,Title,Option1 Name,Option1 Value,Variant Price",
"tee,Tee,Size,S,24.00",
"tee,,,M,24.00",
"tee,,,L,26.00",
)
assert p.valid and p.option_names == ("Size", None, None)
assert [v.options[0] for v in p.variants] == ["S", "M", "L"]
def test_image_only_rows_and_dedupe():
[p] = _products(
"Handle,Title,Image Src,Image Position,Image Alt Text",
"mug,Mug,https://x/a.jpg,1,front",
"mug,,https://x/b.jpg,2,back",
"mug,,https://x/a.jpg,3,dupe",
)
assert p.valid
assert [(i.source_url, i.position) for i in p.images] == [
("https://x/a.jpg", 1), ("https://x/b.jpg", 2),
]
def test_variant_image_joins_product_images():
[p] = _products(
"Handle,Title,Variant Image",
"mug,Mug,https://x/v.jpg",
)
assert p.variants[0].fields["variant_image"] == "https://x/v.jpg"
assert [i.source_url for i in p.images] == ["https://x/v.jpg"]
def test_description_sanitized_inv15():
[p] = _products(
"Handle,Title,Description",
'mug,Mug,"<p onclick=\'x()\'>hi</p><script>evil()</script>"',
)
html = p.fields["description_html"]
assert "<p>" in html and "script" not in html and "onclick" not in html
def test_blank_cell_clears_absent_column_missing():
[p] = _products("Handle,Title,Vendor", "mug,Mug,")
assert p.fields["vendor"] is None # present-but-empty == clear
assert "status" not in p.fields # absent column == untouched
@pytest.mark.parametrize(
"header,row,column,fragment",
[
("Handle,Title", ",NoHandle", "Handle", "needs a Handle"),
("Handle,Title", "Bad_Handle!,T", "Handle", "isn't a valid handle"),
("Handle,Title", "mug,", "Title", "missing its Title"),
("Handle,Title,Type", "mug,Mug,kit_virtual", "Type", "kits arrive"),
("Handle,Title,Component 1 SKU", "mug,Mug,ABC", "Component 1 SKU", "kits arrive"),
("Handle,Title,Status", "mug,Mug,live", "Status", "is not a status"),
("Handle,Title,Published", "mug,Mug,YES", "Published", "is not TRUE or FALSE"),
("Handle,Title,Variant Price", 'mug,Mug,"12,50"', "Variant Price", "is not a price"),
("Handle,Title,Variant Cost", "mug,Mug,-3", "Variant Cost", "is not a price"),
("Handle,Title,Variant Weight", "mug,Mug,heavy", "Variant Weight", "is not a number"),
("Handle,Title,Variant Inventory Qty", "mug,Mug,3.5", "Variant Inventory Qty", "is not a whole number"),
("Handle,Title,Variant Position", "mug,Mug,0", "Variant Position", "is not a position"),
("Handle,Title,Option1 Value", "mug,Mug,Red", "Option1 Value", "no Option1 Name"),
],
)
def test_row_error_rules(header, row, column, fragment):
errors = _errors(header, row)
assert any(e.column == column and fragment in e.message for e in errors), errors
def test_missing_option_value_for_named_option():
errors = _errors(
"Handle,Title,Option1 Name,Option1 Value,Variant SKU",
"tee,Tee,Size,S,A",
"tee,,,,B",
)
assert any("missing its Option1 Value" in e.message and e.line_number == 3 for e in errors)
def test_duplicate_option_combo_is_error():
errors = _errors(
"Handle,Title,Option1 Name,Option1 Value",
"tee,Tee,Size,S",
"tee,,,S",
)
assert any("duplicate variant" in e.message for e in errors)
def test_second_variant_on_no_option_product_is_error():
errors = _errors(
"Handle,Title,Variant SKU",
"mug,Mug,A",
"mug,,B",
)
assert any("without options can have only one variant" in e.message for e in errors)
def test_non_consecutive_handle_is_error():
errors = _errors(
"Handle,Title",
"mug,Mug",
"tee,Tee",
"mug,",
)
assert any("must be consecutive" in e.message and e.line_number == 4 for e in errors)
def test_no_data_row_is_error():
errors = _errors(
"Handle,Title,Variant SKU,Image Src",
"mug,Mug,A,",
"mug,,,",
)
assert any("no variant or image data" in e.message for e in errors)
def test_errors_poison_their_product_only():
products = _products(
"Handle,Title,Variant Price",
"good-mug,Mug,10.00",
"bad-tee,Tee,not-a-price",
)
by_handle = {p.handle: p for p in products}
assert by_handle["good-mug"].valid
assert not by_handle["bad-tee"].valid
def test_all_errors_collected_not_first_only():
[p] = _products(
"Handle,Title,Status,Variant Price",
"mug,,bogus,abc",
)
assert len(p.errors) == 3 # missing Title + bad status + bad price
File diff suppressed because it is too large Load Diff