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