Files
wiggleverse-ecomm/backend/app/domains/products/service.py
T
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

228 lines
9.4 KiB
Python

"""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),
}