"""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 collections.abc import Iterator from datetime import datetime, timezone import psycopg from app.platform import config, telemetry from . import codec, diff, repo, serialize, validate from .errors import ( DraftExpired, DraftNotFound, EmptyCatalog, 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="applying", ) for plan in diff_result.plan: _apply_product_plan(conn, storefront_id, plan, run_id) repo.insert_run_errors(conn, run_id, error_rows) pending = repo.run_image_counts(conn, run_id)["pending"] repo.set_run_status(conn, run_id, "fetching_images" if pending else "complete") 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 export_catalog( conn: psycopg.Connection, storefront_id: int, status_filter: str ) -> Iterator[str]: """Stream the storefront's catalog as canonical CSV (PUC-9; INV-12 codec). Read-only: builds the snapshot, then streams the serializer over it. TEL-3 is emitted once the stream is exhausted, with the product count and elapsed time. Raises EmptyCatalog before yielding anything if the (filtered) catalog is empty, so the BFF can answer 409 cleanly with no partial body. """ started = time.monotonic() snapshot = repo.export_catalog(conn, storefront_id, status_filter) if not snapshot: raise EmptyCatalog() def _stream() -> Iterator[str]: yield from serialize.catalog_to_csv(snapshot, base_url=config.public_base_url()) telemetry.emit( "catalog_exported", storefront_id=storefront_id, status_filter=status_filter, product_count=len(snapshot), duration_ms=int((time.monotonic() - started) * 1000), ) return _stream() 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), }