# Catalog Sync and Normalisation Implementation Plan


**Goal:** Pull a merchant's catalog from the Shopify Admin API and from an authenticated HTTP product API, normalise both deterministically into one shape, and land them in `strategist_products` so the chatbot recommends real products.

**Architecture:** Two source adapters emit the same `SourceProduct` dict, so normalisation is written once. Normalisation is split into four pure modules (text, taxonomy, attributes, commercial) with no I/O, each testable against committed fixtures captured from a real store. An orchestrator fetches, normalises, validates, persists scoped by source, and reports. No LLM anywhere in this plan.

**Tech Stack:** FastAPI, psycopg2, httpx, pytest.

**Spec:** `docs/specs/2026-08-10-catalog-sync-design.md`

**Depends on:** sub-project A (`app/services/product_sources.py`, `app/services/shopify_oauth.py`), already built and live-verified.

## Global Constraints

- **Never commit, never push, never stage.** The user commits all work themselves. Every task ends by running `git status --short` and reporting the changed files. Do not run `git commit`, `git push`, or `git add`.
- **No assistant attribution** in any file, message, or commit.
- **Comments explain WHY, not WHAT.** A short comment only where the reason for a decision is not evident from the code. Do not narrate mechanics, label obvious blocks, or restate a function name in its docstring. A reviewer who knows Python must find zero noise comments.
- **Match existing conventions:** module-level `logger = logging.getLogger(__name__)`, snake_case, double quotes, 4-space indent, type hints on public signatures, import ordering stdlib → third-party → `app.*`, `psycopg2` with `sql.SQL`/`sql.Identifier` for any interpolated identifier.
- **Every write and every delete against `strategist_products` in this plan is scoped by `source_kind` AND `source_ref`.** The existing unscoped delete in `products.save_products` would erase all crawled products after a Shopify sync. Never widen a delete beyond one source.
- **Money is integer minor units.** Never floats, never a single price where a range exists.
- **`in_stock` derives from availability flags, never from inventory quantity.** Overselling merchants carry negative quantities on purchasable products.
- **Empty string is never stored.** Normalise it to `None` so hashes and code paths do not fork.
- **Sort every array before hashing.** Unsorted tags or attributes change the hash on every sync and re-embed the whole catalog nightly.
- Tests: `.venv/bin/python -m pytest`. `tests/unit/` baseline is **224 passing**. `.venv` is a symlink to the main checkout's virtualenv. Use `PYTHONPATH=.` only when running a script directly; pytest does not need it.
- Fixtures already captured from the live store are in `tests/fixtures/shopify/`. Use them; do not invent product JSON.

---

## File Structure

| File | Responsibility |
|---|---|
| `app/services/database.py` (modify) | Add the new columns to `bootstrap_tenant`'s CREATE TABLE, plus a `migrate_products_table(tenant_id)` for existing tenants. |
| `app/services/catalog_fieldmap.py` (new) | Execute a declarative field map against a JSON record. Pure. |
| `app/services/catalog_text.py` (new) | Text normalisation, title cleaning, brand resolution. Pure. |
| `app/services/catalog_taxonomy.py` (new) | The category resolution chain. Pure. |
| `app/services/catalog_attributes.py` (new) | Key/value canonicalisation, tag parsing, source precedence. Pure. |
| `app/services/catalog_commercial.py` (new) | Price range, sale, stock, product URL, image cleaning. Pure. |
| `app/services/catalog_normalise.py` (new) | Orchestrates the passes; computes `quality_score`, `content_hash`, `record_hash`; validates. Pure. |
| `app/services/catalog_shopify.py` (new) | Cursor-paginated GraphQL fetch → `SourceProduct` list. |
| `app/services/catalog_http.py` (new) | Offset-paginated authenticated fetch → `SourceProduct` list. |
| `app/services/catalog_sync.py` (new) | Fetch → normalise → validate → persist scoped → report. |
| `app/api/endpoints.py` (modify) | `POST /catalog/sync`. |
`app/services/products.py` is **not** modified by this plan. Recommendation behaviour —
which products are eligible, how they rank — is owned by separate work. B stores `status`,
`in_stock` and prices accurately and stops there.

---

### Task 1: Schema migration

The riskiest task in the plan: it alters a table that already holds live crawled rows. It runs first because everything else writes to the new columns.

**Files:**
- Modify: `app/services/database.py`
- Modify: `tests/integration/test_products_table.py:18-27` (its `EXPECTED_COLUMNS` is an exact set and this task adds 19 columns)
- Create: `tests/integration/conftest.py` addition — a `temp_tenant` fixture
- Test: `tests/integration/test_products_migration.py`

**Interfaces:**
- Consumes: `get_db_connection`, `bootstrap_tenant` (both exist).
- Produces:
  - `migrate_products_table(tenant_id: str) -> dict` returning `{"added": [column names]}`
  - `temp_tenant` pytest fixture yielding a throwaway schema name, dropped on teardown.

**This task breaks an existing test unless you fix it.** `tests/integration/test_products_table.py` asserts `set(columns) == EXPECTED_COLUMNS` exactly. Adding columns without updating that set turns a passing suite red. Step 6 handles it.

- [ ] **Step 1: Add the throwaway-schema fixture**

Append to `tests/integration/conftest.py`:

```python
import pytest

from app.services.database import get_db_connection


@pytest.fixture
def temp_tenant():
    """A throwaway schema so destructive tests never touch a real tenant.

    The scoped-delete behaviour this exercises is the one bug in this feature
    that would silently destroy a tenant's crawled catalog, so it is tested
    against real Postgres rather than by inspecting source text.
    """
    schema = "test_catalog_tmp"
    conn = get_db_connection()
    try:
        with conn.cursor() as cur:
            cur.execute(f'DROP SCHEMA IF EXISTS "{schema}" CASCADE')
            cur.execute(f'CREATE SCHEMA "{schema}"')
        conn.commit()
    finally:
        conn.close()

    yield schema

    conn = get_db_connection()
    try:
        with conn.cursor() as cur:
            cur.execute(f'DROP SCHEMA IF EXISTS "{schema}" CASCADE')
        conn.commit()
    finally:
        conn.close()
```

- [ ] **Step 2: Write the failing test**

Create `tests/integration/test_products_migration.py`:

```python
"""Real-database tests for the products table migration.

Asserted against Postgres rather than source text: the migration runs against
live tenant data, so "the words appear in the function" is not evidence.
"""
import pytest

from app.services.database import (
    bootstrap_tenant, get_db_connection, migrate_products_table,
)

NEW_COLUMNS = {
    "source_kind", "source_ref", "external_id", "brand",
    "taxonomy_path", "taxonomy_source", "raw_category",
    "price_cents", "price_max_cents", "compare_at_cents", "currency",
    "on_sale", "in_stock", "status", "attributes",
    "quality_score", "content_hash", "record_hash", "synced_at",
}

OLD_SHAPE = """
    CREATE TABLE "{s}".strategist_products (
        id           SERIAL PRIMARY KEY,
        product_key  TEXT NOT NULL UNIQUE,
        name         TEXT NOT NULL,
        description  TEXT,
        image_url    TEXT,
        product_url  TEXT NOT NULL,
        category     TEXT,
        ctas         JSONB NOT NULL DEFAULT '[]',
        options      JSONB NOT NULL DEFAULT '[]',
        raw          JSONB NOT NULL DEFAULT '{{}}',
        related_keys TEXT[] NOT NULL DEFAULT '{{}}',
        extracted_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP
    )
"""


def _exec(sql_text):
    conn = get_db_connection()
    try:
        with conn.cursor() as cur:
            cur.execute(sql_text)
        conn.commit()
    finally:
        conn.close()


def _query(sql_text):
    conn = get_db_connection()
    try:
        with conn.cursor() as cur:
            cur.execute(sql_text)
            return cur.fetchall()
    finally:
        conn.close()


def _columns(schema):
    return {r[0] for r in _query(
        "SELECT column_name FROM information_schema.columns "
        f"WHERE table_schema = '{schema}' AND table_name = 'strategist_products'")}


@pytest.fixture
def old_table(temp_tenant):
    _exec(OLD_SHAPE.format(s=temp_tenant))
    _exec(f"""INSERT INTO "{temp_tenant}".strategist_products
              (product_key, name, product_url) VALUES
              ('crawl:a', 'Crawled A', 'https://x.example/a'),
              ('crawl:b', 'Crawled B', 'https://x.example/b')""")
    return temp_tenant


def test_migration_adds_every_new_column(old_table):
    assert not (NEW_COLUMNS & _columns(old_table))
    migrate_products_table(old_table)
    assert NEW_COLUMNS <= _columns(old_table)


def test_migration_loses_no_rows(old_table):
    before = _query(f'SELECT count(*) FROM "{old_table}".strategist_products')[0][0]
    migrate_products_table(old_table)
    after = _query(f'SELECT count(*) FROM "{old_table}".strategist_products')[0][0]
    assert after == before == 2


def test_pre_existing_rows_are_labelled_crawl(old_table):
    # Any other default would let a Shopify sync delete them as stale.
    migrate_products_table(old_table)
    kinds = _query(
        f'SELECT DISTINCT source_kind FROM "{old_table}".strategist_products')
    assert kinds == [("crawl",)]


def test_migration_is_idempotent(old_table):
    migrate_products_table(old_table)
    first = _columns(old_table)
    migrate_products_table(old_table)
    assert _columns(old_table) == first
    rows = _query(f'SELECT count(*) FROM "{old_table}".strategist_products')[0][0]
    assert rows == 2


def test_bootstrap_creates_the_full_shape_for_a_new_tenant(temp_tenant):
    bootstrap_tenant(temp_tenant)
    assert NEW_COLUMNS <= _columns(temp_tenant)


def test_the_source_index_exists(old_table):
    migrate_products_table(old_table)
    idx = {r[0] for r in _query(
        f"SELECT indexname FROM pg_indexes WHERE schemaname = '{old_table}'")}
    assert "strategist_products_source_idx" in idx
```

- [ ] **Step 3: Run test to verify it fails**

Run: `.venv/bin/python -m pytest tests/integration/test_products_migration.py -v`
Expected: FAIL with `ImportError: cannot import name 'migrate_products_table'`

- [ ] **Step 4: Add the columns to bootstrap_tenant**

In `app/services/database.py`, inside `bootstrap_tenant`, replace the `strategist_products` CREATE TABLE body with:

```python
            cur.execute(sql.SQL("""
                CREATE TABLE IF NOT EXISTS {}.strategist_products (
                    id           SERIAL PRIMARY KEY,
                    product_key  TEXT NOT NULL UNIQUE,
                    name         TEXT NOT NULL,
                    description  TEXT,
                    image_url    TEXT,
                    product_url  TEXT NOT NULL,
                    category     TEXT,
                    ctas         JSONB NOT NULL DEFAULT '[]',
                    options      JSONB NOT NULL DEFAULT '[]',
                    raw          JSONB NOT NULL DEFAULT '{{}}',
                    related_keys TEXT[] NOT NULL DEFAULT '{{}}',
                    extracted_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP,
                    source_kind      TEXT NOT NULL DEFAULT 'crawl',
                    source_ref       TEXT,
                    external_id      TEXT,
                    brand            TEXT,
                    taxonomy_path    TEXT[] NOT NULL DEFAULT '{{}}',
                    taxonomy_source  TEXT,
                    raw_category     TEXT,
                    price_cents      INT,
                    price_max_cents  INT,
                    compare_at_cents INT,
                    currency         TEXT,
                    on_sale          BOOLEAN NOT NULL DEFAULT false,
                    in_stock         BOOLEAN NOT NULL DEFAULT true,
                    status           TEXT,
                    attributes       JSONB NOT NULL DEFAULT '[]',
                    quality_score    REAL,
                    content_hash     TEXT,
                    record_hash      TEXT,
                    synced_at        TIMESTAMP WITH TIME ZONE
                )
            """).format(sql.Identifier(tenant_id)))
```

Then add these indexes beside the existing two:

```python
            cur.execute(sql.SQL(
                "CREATE INDEX IF NOT EXISTS strategist_products_source_idx "
                "ON {}.strategist_products (source_kind, source_ref)"
            ).format(sql.Identifier(tenant_id)))
            cur.execute(sql.SQL(
                "CREATE INDEX IF NOT EXISTS strategist_products_status_idx "
                "ON {}.strategist_products (status, in_stock)"
            ).format(sql.Identifier(tenant_id)))
```

- [ ] **Step 5: Add the migration function**

Append to `app/services/database.py`:

```python
_PRODUCT_COLUMNS = [
    ("source_kind", "TEXT NOT NULL DEFAULT 'crawl'"),
    ("source_ref", "TEXT"),
    ("external_id", "TEXT"),
    ("brand", "TEXT"),
    ("taxonomy_path", "TEXT[] NOT NULL DEFAULT '{}'"),
    ("taxonomy_source", "TEXT"),
    ("raw_category", "TEXT"),
    ("price_cents", "INT"),
    ("price_max_cents", "INT"),
    ("compare_at_cents", "INT"),
    ("currency", "TEXT"),
    ("on_sale", "BOOLEAN NOT NULL DEFAULT false"),
    ("in_stock", "BOOLEAN NOT NULL DEFAULT true"),
    ("status", "TEXT"),
    ("attributes", "JSONB NOT NULL DEFAULT '[]'"),
    ("quality_score", "REAL"),
    ("content_hash", "TEXT"),
    ("record_hash", "TEXT"),
    ("synced_at", "TIMESTAMP WITH TIME ZONE"),
]


def migrate_products_table(tenant_id: str) -> dict:
    """Bring an existing tenant's products table up to the catalog-sync shape.

    bootstrap_tenant uses CREATE TABLE IF NOT EXISTS, which will not alter a
    table that already holds crawled rows. source_kind defaults to 'crawl' so
    every pre-existing row is labelled correctly -- any other default would let
    a later sync delete them as stale.
    """
    added = []
    conn = get_db_connection()
    try:
        with conn.cursor() as cur:
            for name, ddl in _PRODUCT_COLUMNS:
                # ddl is appended as its own sql.SQL fragment rather than folded
                # into the format string: taxonomy_path's default is a literal
                # '{}', which .format() would otherwise read as a placeholder.
                stmt = sql.SQL(
                    "ALTER TABLE {}.strategist_products ADD COLUMN IF NOT EXISTS {} "
                ).format(sql.Identifier(tenant_id), sql.Identifier(name)) + sql.SQL(ddl)
                cur.execute(stmt)
                added.append(name)
            cur.execute(sql.SQL(
                "CREATE INDEX IF NOT EXISTS strategist_products_source_idx "
                "ON {}.strategist_products (source_kind, source_ref)"
            ).format(sql.Identifier(tenant_id)))
            cur.execute(sql.SQL(
                "CREATE INDEX IF NOT EXISTS strategist_products_status_idx "
                "ON {}.strategist_products (status, in_stock)"
            ).format(sql.Identifier(tenant_id)))
        conn.commit()
        logger.info(f"Migrated products table for {tenant_id}: {len(added)} columns")
        return {"added": added}
    except Exception:
        conn.rollback()
        logger.error(f"Could not migrate products table for {tenant_id}", exc_info=True)
        raise
    finally:
        conn.close()
```

- [ ] **Step 6: Update the existing exact-columns test**

`tests/integration/test_products_table.py` asserts `set(columns) == EXPECTED_COLUMNS`.
Add all 19 new names to that set, with a comment recording why they exist:

```python
EXPECTED_COLUMNS = {
    "id", "product_key", "name", "description",
    "image_url", "product_url", "category", "ctas", "options",
    # The page's structured data, kept verbatim. Nothing reads it yet -- it
    # exists so that adding a field later is a query rather than a re-crawl of
    # every tenant's site.
    "raw",
    "related_keys", "extracted_at",
    # Catalog sync: the table now serves the crawler, Shopify and HTTP APIs, so
    # every row records which producer wrote it and every delete is scoped to one.
    "source_kind", "source_ref", "external_id", "brand",
    "taxonomy_path", "taxonomy_source", "raw_category",
    "price_cents", "price_max_cents", "compare_at_cents", "currency",
    "on_sale", "in_stock", "status", "attributes",
    "quality_score", "content_hash", "record_hash", "synced_at",
}
```

- [ ] **Step 7: Run the migration against the live test tenant and prove no row was lost**

Do this BEFORE the full integration suite. `bootstrap_tenant` now creates indexes on the
new columns, and the live tenant's existing table does not have them yet — running the
suite first fails with `UndefinedColumn` until this migration has run once.

```bash
PYTHONPATH=. .venv/bin/python -c "
from app.services.database import get_db_connection, migrate_products_table
T='org_8c32bf3e-6a18-4739-9b1c-94c0cf11125f'
c=get_db_connection(); cur=c.cursor()
cur.execute('SELECT count(*) FROM \"%s\".strategist_products' % T); before=cur.fetchone()[0]
c.close()
print('rows before:', before)
print(migrate_products_table(T))
c=get_db_connection(); cur=c.cursor()
cur.execute('SELECT count(*), count(*) FILTER (WHERE source_kind=%s) FROM \"'+T+'\".strategist_products', ('crawl',))
after, crawl = cur.fetchone(); c.close()
print('rows after:', after, 'labelled crawl:', crawl)
assert after == before, 'ROWS LOST'
assert crawl == after, 'pre-existing rows not labelled crawl'
print('migration safe')
"
```

Expected: `migration safe`, with `rows after` equal to `rows before`.

Run it a second time and confirm it succeeds again — the migration must be idempotent.

- [ ] **Step 9: Report changed files**

```bash
git status --short
```

Do not commit and do not stage.

---

### Task 2: Declarative field map executor

Pure function, no I/O. This is the piece C's LLM will later author configs for.

**Files:**
- Create: `app/services/catalog_fieldmap.py`
- Test: `tests/unit/test_catalog_fieldmap.py`

**Interfaces:**
- Consumes: nothing.
- Produces:
  - `resolve_path(record: dict, path: str)` → the value at a `$.a.b[0]` path, or `None`
  - `apply_field_map(record: dict, fields: dict) -> dict` → `{target_name: value}`
  - `class FieldMapError(Exception)`

- [ ] **Step 1: Write the failing test**

Create `tests/unit/test_catalog_fieldmap.py`:

```python
import pytest

from app.services.catalog_fieldmap import (
    FieldMapError, apply_field_map, resolve_path,
)

RECORD = {
    "id": 1,
    "title": "Essence Mascara Lash Princess",
    "price": 9.99,
    "dimensions": {"width": 15.14, "height": 13.08},
    "images": ["https://cdn.example.com/1.webp", "https://cdn.example.com/2.webp"],
    "tags": ["beauty", "mascara"],
    "meta": {"barcode": "5784719087687"},
    "empty": None,
}


@pytest.mark.parametrize("path,expected", [
    ("$.id", 1),
    ("$.title", "Essence Mascara Lash Princess"),
    ("$.price", 9.99),
    ("$.dimensions.width", 15.14),
    ("$.images[0]", "https://cdn.example.com/1.webp"),
    ("$.images[1]", "https://cdn.example.com/2.webp"),
    ("$.meta.barcode", "5784719087687"),
    ("$.tags", ["beauty", "mascara"]),
    ("$.empty", None),
])
def test_resolve_path(path, expected):
    assert resolve_path(RECORD, path) == expected


@pytest.mark.parametrize("path", [
    "$.missing",
    "$.dimensions.depth",
    "$.images[9]",
    "$.title.nested",
    "$.tags[0].deeper",
])
def test_missing_paths_return_none_rather_than_raising(path):
    # A merchant API with a sparse record must degrade, not abort the sync.
    assert resolve_path(RECORD, path) is None


@pytest.mark.parametrize("path", ["title", "", "$", "$..title", None])
def test_malformed_paths_raise(path):
    with pytest.raises(FieldMapError):
        resolve_path(RECORD, path)


def test_apply_field_map():
    out = apply_field_map(RECORD, {
        "external_id": "$.id",
        "title": "$.title",
        "price": "$.price",
        "image_url": "$.images[0]",
        "tags": "$.tags",
        "brand": "$.missing",
    })
    assert out == {
        "external_id": 1,
        "title": "Essence Mascara Lash Princess",
        "price": 9.99,
        "image_url": "https://cdn.example.com/1.webp",
        "tags": ["beauty", "mascara"],
        "brand": None,
    }


def test_apply_field_map_rejects_a_bad_path_loudly():
    # A typo in a stored config must surface at sync time, not silently null a field.
    with pytest.raises(FieldMapError):
        apply_field_map(RECORD, {"title": "title"})
```

- [ ] **Step 2: Run test to verify it fails**

Run: `.venv/bin/python -m pytest tests/unit/test_catalog_fieldmap.py -v`
Expected: FAIL with `ModuleNotFoundError: No module named 'app.services.catalog_fieldmap'`

- [ ] **Step 3: Write the implementation**

Create `app/services/catalog_fieldmap.py`:

```python
"""Executes a declarative field map against a source record.

The syntax is deliberately tiny: $.a.b and $.a[0]. Anything richer would make
the maps harder for a human to review, and C's LLM authors these configs.
"""
import logging
import re

logger = logging.getLogger(__name__)

_SEGMENT = re.compile(r"\A([A-Za-z_][A-Za-z0-9_]*)((?:\[\d+\])*)\Z")
_INDEX = re.compile(r"\[(\d+)\]")


class FieldMapError(Exception):
    """The field map itself is malformed, as opposed to the data being sparse."""


def resolve_path(record: dict, path):
    if not isinstance(path, str) or not path.startswith("$.") or path == "$.":
        raise FieldMapError(f"malformed path: {path!r}")

    current = record
    for raw in path[2:].split("."):
        m = _SEGMENT.match(raw)
        if not m:
            raise FieldMapError(f"malformed path segment {raw!r} in {path!r}")

        key, indexes = m.group(1), _INDEX.findall(m.group(2))
        if not isinstance(current, dict):
            return None
        current = current.get(key)

        for i in indexes:
            if not isinstance(current, list):
                return None
            idx = int(i)
            current = current[idx] if idx < len(current) else None

    return current


def apply_field_map(record: dict, fields: dict) -> dict:
    return {target: resolve_path(record, path) for target, path in fields.items()}
```

- [ ] **Step 4: Run test to verify it passes**

Run: `.venv/bin/python -m pytest tests/unit/test_catalog_fieldmap.py -v`
Expected: PASS, 20 tests

- [ ] **Step 5: Report changed files**

```bash
git status --short
```

Do not commit and do not stage.

---

### Task 3: Text, title and brand normalisation

**Files:**
- Create: `app/services/catalog_text.py`
- Test: `tests/unit/test_catalog_text.py`

**Interfaces:**
- Consumes: nothing.
- Produces:
  - `normalise_text(s) -> str | None`
  - `clean_title(raw, brand=None) -> str | None`
  - `normalise_brand(vendor, shop_name, known_brands=()) -> str | None`

- [ ] **Step 1: Write the failing test**

Create `tests/unit/test_catalog_text.py`:

```python
import pytest

from app.services.catalog_text import clean_title, normalise_brand, normalise_text


@pytest.mark.parametrize("raw,expected", [
    ("  Trail Runner  ", "Trail Runner"),
    ("Men’s Jacket", "Men’s Jacket"),
    ("ﬁne wool", "fine wool"),          # NFKC folds the fi ligature
    ("２００", "200"),           # full-width digits
    ("a​b", "ab"),                      # zero-width space
    ("a b", "a b"),                     # non-breaking space
    ("a   b", "a b"),
    ("", None),
    ("   ", None),
    (None, None),
])
def test_normalise_text(raw, expected):
    assert normalise_text(raw) == expected


def test_empty_string_becomes_none_not_empty():
    # "" and None must not produce two hashes or two code paths.
    assert normalise_text("") is None


@pytest.mark.parametrize("raw,brand,expected", [
    ("Salomon - Trail Runner", "Salomon", "Trail Runner"),
    ("Salomon | Trail Runner", "Salomon", "Trail Runner"),
    ("Salomon: Trail Runner", "Salomon", "Trail Runner"),
    ("Trail Runner", "Salomon", "Trail Runner"),
    ("Trail Runner Pro (SKU-88421-B)", None, "Trail Runner Pro"),
    ("Trail Runner Pro [item: 42]", None, "Trail Runner Pro"),
    ("THE COMPLETE SNOWBOARD", None, "The Complete Snowboard"),
    ("XL", None, "XL"),
])
def test_clean_title(raw, brand, expected):
    assert clean_title(raw, brand) == expected


def test_brand_strip_never_empties_the_title():
    # A title that is only the brand needs the brand more than the separation.
    assert clean_title("Salomon", "Salomon") == "Salomon"
    assert clean_title("Salomon - S", "Salomon") == "Salomon - S"


def test_shop_name_vendor_is_rejected():
    # 8 of 17 products in the live store carry the shop name as vendor, which is
    # Shopify's default when the merchant leaves it blank.
    assert normalise_brand("Galaxiq", "Galaxiq") is None
    assert normalise_brand("galaxiq", "Galaxiq") is None
    assert normalise_brand("Galaxiq Store", "Galaxiq") is None


@pytest.mark.parametrize("v", ["n/a", "N/A", "none", "unknown", "default",
                               "generic", "no brand", "-", "", None])
def test_placeholder_vendors_rejected(v):
    assert normalise_brand(v, "Galaxiq") is None


def test_real_vendor_kept():
    assert normalise_brand("Snowboard Vendor", "Galaxiq") == "Snowboard Vendor"


def test_casing_canonicalised_against_known_brands():
    assert normalise_brand("salomon", "Galaxiq", ["Salomon"]) == "Salomon"
    assert normalise_brand("SALOMON", "Galaxiq", ["Salomon"]) == "Salomon"
```

- [ ] **Step 2: Run test to verify it fails**

Run: `.venv/bin/python -m pytest tests/unit/test_catalog_text.py -v`
Expected: FAIL with `ModuleNotFoundError: No module named 'app.services.catalog_text'`

- [ ] **Step 3: Write the implementation**

Create `app/services/catalog_text.py`:

```python
"""Text, title and brand normalisation. Pure functions, no I/O."""
import logging
import re
import unicodedata

logger = logging.getLogger(__name__)

_ZERO_WIDTH = re.compile(r"[​-‍﻿]")
_ODD_SPACES = re.compile(r"[   ]")
_WHITESPACE = re.compile(r"\s+")
_SKU_SUFFIX = re.compile(
    r"\s*[\(\[]\s*(sku|item|ref|art)[-.: ]*[\w-]+\s*[\)\]]\s*\Z", re.I)
_PLACEHOLDER_BRAND = re.compile(
    r"\A(n/?a|none|unknown|default|generic|no brand|-)\Z", re.I)



def normalise_text(s):
    """NFKC first: merchants paste from Word, Excel and PDFs, so ligatures and
    full-width characters arrive routinely and would never match otherwise."""
    if s is None or not isinstance(s, str):
        return None
    out = unicodedata.normalize("NFKC", s)
    out = _ZERO_WIDTH.sub("", out)
    out = _ODD_SPACES.sub(" ", out)
    out = _WHITESPACE.sub(" ", out).strip()
    return out or None


def _title_case(s: str) -> str:
    return " ".join(w.capitalize() for w in s.split(" "))


def clean_title(raw, brand=None):
    t = normalise_text(raw)
    if not t:
        return None

    if brand:
        stripped = re.sub(
            rf"\A{re.escape(brand)}\s*[-–—|:]\s*", "", t, flags=re.I)
        # Guard: a title that is only the brand needs it more than the separation.
        if len(stripped.strip()) > 1:
            t = stripped

    t = _SKU_SUFFIX.sub("", t)

    if len(t) > 8 and t == t.upper() and re.search(r"[A-Z]{4,}", t):
        t = _title_case(t)

    return t.strip() or None


def _is_shop_name(vendor: str, shop_name: str) -> bool:
    """Exact-or-prefix rather than fuzzy similarity: no threshold separates
    "Galaxiq Store" (0.700) from the real brand "Galaxy" (0.769), so any score
    low enough to catch the first destroys the second."""
    v, s = vendor.lower(), shop_name.lower()
    if v == s:
        return True
    return v.startswith(f"{s} ") or s.startswith(f"{v} ")


def normalise_brand(vendor, shop_name, known_brands=()):
    """Shopify defaults vendor to the shop name when a merchant leaves it blank,
    so without the shop-name check a whole catalog shares one brand and any
    brand-match signal becomes noise."""
    b = normalise_text(vendor)
    if not b:
        return None

    if _PLACEHOLDER_BRAND.match(b):
        return None

    if shop_name and _is_shop_name(b, shop_name):
        return None

    for known in known_brands:
        if known.lower() == b.lower():
            return known

    return b
```

- [ ] **Step 4: Run test to verify it passes**

Run: `.venv/bin/python -m pytest tests/unit/test_catalog_text.py -v`
Expected: PASS, 31 tests

- [ ] **Step 5: Report changed files**

```bash
git status --short
```

Do not commit and do not stage.

---

### Task 4: Taxonomy resolution

**Files:**
- Create: `app/services/catalog_taxonomy.py`
- Test: `tests/unit/test_catalog_taxonomy.py`

**Interfaces:**
- Consumes: `catalog_text.normalise_text`.
- Produces: `resolve_taxonomy(source_product: dict) -> dict` returning
  `{"path": list[str], "source": str, "raw": str | None}` where `source` is one of
  `native | product_type | collection | tag | none`.

- [ ] **Step 1: Write the failing test**

Create `tests/unit/test_catalog_taxonomy.py`:

```python
from app.services.catalog_taxonomy import MAX_DEPTH, resolve_taxonomy


def sp(**kw):
    base = {"raw_category": None, "product_type": None,
            "collections": [], "tags": []}
    base.update(kw)
    return base


def test_native_category_wins():
    out = resolve_taxonomy(sp(
        raw_category="Apparel & Accessories > Shoes > Athletic",
        product_type="Shoes", collections=["Sale"], tags=["shoes"]))
    assert out["source"] == "native"
    assert out["path"] == ["Apparel & Accessories", "Shoes", "Athletic"]


def test_falls_back_to_product_type():
    out = resolve_taxonomy(sp(product_type="snowboard"))
    assert out["source"] == "product_type"
    assert out["path"] == ["snowboard"]
    assert out["raw"] == "snowboard"


def test_empty_product_type_is_not_a_match():
    # One live product has productType "" -- it must fall through, not resolve
    # to an empty path labelled product_type.
    out = resolve_taxonomy(sp(product_type="   "))
    assert out["source"] == "none"
    assert out["path"] == []


def test_falls_back_to_collections():
    out = resolve_taxonomy(sp(collections=["Snowboards"]))
    assert out["source"] == "collection"
    assert out["path"] == ["Snowboards"]


def test_merchandising_collections_are_skipped():
    # The live store has "Automated Collection" on 8 products and "Home page" on
    # one. Resolving to either makes unrelated products category siblings.
    out = resolve_taxonomy(sp(collections=["Home page", "Automated Collection"]))
    assert out["source"] == "none"


def test_merchandising_skipped_but_a_real_collection_still_wins():
    out = resolve_taxonomy(sp(collections=["Sale", "Snowboards"]))
    assert out["source"] == "collection"
    assert out["path"] == ["Snowboards"]


def test_falls_back_to_tags():
    out = resolve_taxonomy(sp(tags=["Snowboard"]))
    assert out["source"] == "tag"
    assert out["path"] == ["Snowboard"]


def test_unresolved():
    out = resolve_taxonomy(sp())
    assert out == {"path": [], "source": "none", "raw": None}


def test_depth_is_truncated():
    deep = " > ".join(f"L{i}" for i in range(8))
    out = resolve_taxonomy(sp(raw_category=deep))
    assert len(out["path"]) == MAX_DEPTH


def test_raw_category_is_preserved_for_debugging():
    out = resolve_taxonomy(sp(raw_category="Gift Cards"))
    assert out["raw"] == "Gift Cards"
```

- [ ] **Step 2: Run test to verify it fails**

Run: `.venv/bin/python -m pytest tests/unit/test_catalog_taxonomy.py -v`
Expected: FAIL with `ModuleNotFoundError: No module named 'app.services.catalog_taxonomy'`

- [ ] **Step 3: Write the implementation**

Create `app/services/catalog_taxonomy.py`:

```python
"""Resolves a product's category from whichever source actually carries one."""
import logging
import re

from app.services.catalog_text import normalise_text

logger = logging.getLogger(__name__)

MAX_DEPTH = 4

# Nearly every store has a collection holding products from unrelated categories.
# Resolving to one makes every discounted product a category sibling of every
# other, which is exactly the wrong grouping. "Automated Collection" covers 8 of
# 17 products in the live test store and carries no meaning at all.
MERCH_COLLECTIONS = re.compile(
    r"\A(sale|clearance|new|new in|featured|best ?sellers?|trending|all|"
    r"home ?page|staff picks|automated collection)\Z", re.I)


def _truncate(path: list) -> list:
    # A 7-level path fragments the catalog into single-product leaves, which
    # makes category affinity useless because nothing shares a leaf.
    return path[:MAX_DEPTH]


def resolve_taxonomy(source_product: dict) -> dict:
    raw = normalise_text(source_product.get("raw_category"))
    if raw:
        path = [normalise_text(p) for p in raw.split(" > ")]
        return {"path": _truncate([p for p in path if p]),
                "source": "native", "raw": raw}

    ptype = normalise_text(source_product.get("product_type"))
    if ptype:
        return {"path": [ptype], "source": "product_type", "raw": ptype}

    for title in source_product.get("collections") or []:
        c = normalise_text(title)
        if c and not MERCH_COLLECTIONS.match(c):
            return {"path": [c], "source": "collection", "raw": c}

    for tag in source_product.get("tags") or []:
        t = normalise_text(tag)
        if t:
            return {"path": [t], "source": "tag", "raw": t}

    return {"path": [], "source": "none", "raw": None}
```

- [ ] **Step 4: Run test to verify it passes**

Run: `.venv/bin/python -m pytest tests/unit/test_catalog_taxonomy.py -v`
Expected: PASS, 10 tests

- [ ] **Step 5: Report changed files**

```bash
git status --short
```

Do not commit and do not stage.

---

### Task 5: Attribute canonicalisation

**Files:**
- Create: `app/services/catalog_attributes.py`
- Test: `tests/unit/test_catalog_attributes.py`

**Interfaces:**
- Consumes: `catalog_text.normalise_text`.
- Produces:
  - `normalise_key(k) -> str | None`
  - `normalise_value(key, raw) -> tuple[str | None, str | None]` returning `(value, unit)`
  - `build_attributes(source_product: dict) -> list[dict]` with each dict
    `{key, raw_value, value, unit, source, confidence}`

- [ ] **Step 1: Write the failing test**

Create `tests/unit/test_catalog_attributes.py`:

```python
import pytest

from app.services.catalog_attributes import (
    build_attributes, normalise_key, normalise_value,
)


@pytest.mark.parametrize("raw,expected", [
    ("Colour", "color"), ("color", "color"), ("Shade", "color"),
    ("Size", "size"), ("size (uk)", "size"),
    ("Material", "material"), ("Fabric", "material"), ("Composition", "material"),
    ("Gender", "gender"), ("Department", "gender"),
    ("Flavour", "flavor"), ("Scent", "scent"),
    ("Some Custom Thing", "some_custom_thing"),
])
def test_normalise_key(raw, expected):
    assert normalise_key(raw) == expected


def test_shopify_title_placeholder_is_dropped():
    # "Title"/"Default Title" is Shopify's placeholder on single-variant
    # products. It appears on all 17 live products. Kept, every product gains a
    # meaningless attribute that inflates attribute count and quality_score.
    assert normalise_key("Title") is None


@pytest.mark.parametrize("raw,expected", [
    ("Navy", "blue"), ("Midnight Navy", "blue"), ("Charcoal", "grey"),
    ("Ivory", "white"), ("Burgundy", "red"), ("Camel", "brown"),
    ("Black / White", "black"), ("Red, Blue", "red"),
])
def test_colour_maps_to_base_palette(raw, expected):
    value, unit = normalise_value("color", raw)
    assert value == expected


@pytest.mark.parametrize("raw,value,unit", [
    ("XS", "xs", "alpha"), ("Extra Small", "xs", "alpha"),
    ("M", "m", "alpha"), ("Large", "l", "alpha"), ("2XL", "xxl", "alpha"),
    ("UK 9", "9", "uk"), ("EU 43", "43", "eu"), ("10", "10", "numeric"),
])
def test_size_parses_system(raw, value, unit):
    assert normalise_value("size", raw) == (value, unit)


def test_measured_sizes_convert_to_a_base_unit():
    # 500ml and 0.5L must agree, or two identical products score as unrelated.
    assert normalise_value("size", "500ml") == normalise_value("size", "0.5 L")


@pytest.mark.parametrize("raw,expected", [
    ("Mens", "male"), ("Women", "female"), ("Ladies", "female"),
    ("Unisex", "unisex"), ("Kids", "kids"),
])
def test_gender_closed_set(raw, expected):
    value, _ = normalise_value("gender", raw)
    assert value == expected


def test_key_value_tags_parsed():
    attrs = build_attributes({
        "tags": ["material:leather", "color = Navy"],
        "variants": [], "attributes_raw": [],
    })
    by_key = {a["key"]: a for a in attrs}
    assert by_key["material"]["value"] == "leather"
    assert by_key["material"]["confidence"] == 0.9
    assert by_key["color"]["value"] == "blue"
    assert by_key["color"]["raw_value"] == "Navy"


def test_variant_options_become_attributes_but_title_does_not():
    attrs = build_attributes({
        "tags": [],
        "variants": [{"options": {"Title": "Default Title", "Color": "Navy"}}],
        "attributes_raw": [],
    })
    keys = {a["key"] for a in attrs}
    assert "color" in keys
    assert "title" not in keys


def test_metafield_beats_a_bare_tag_on_conflict():
    attrs = build_attributes({
        "tags": ["suede"],
        "variants": [],
        "attributes_raw": [{"key": "material", "value": "leather",
                            "source": "metafield"}],
    })
    material = [a for a in attrs if a["key"] == "material"]
    assert len(material) == 1
    assert material[0]["value"] == "leather"
    assert material[0]["source"] == "metafield"


def test_raw_value_is_always_preserved():
    attrs = build_attributes({
        "tags": [], "attributes_raw": [],
        "variants": [{"options": {"Color": "Midnight Navy"}}],
    })
    colour = next(a for a in attrs if a["key"] == "color")
    assert colour["raw_value"] == "Midnight Navy"
    assert colour["value"] == "blue"
```

- [ ] **Step 2: Run test to verify it fails**

Run: `.venv/bin/python -m pytest tests/unit/test_catalog_attributes.py -v`
Expected: FAIL with `ModuleNotFoundError: No module named 'app.services.catalog_attributes'`

- [ ] **Step 3: Write the implementation**

Create `app/services/catalog_attributes.py`:

```python
"""Canonicalises attribute keys and values from options, tags and metafields."""
import logging
import re

from app.services.catalog_text import normalise_text

logger = logging.getLogger(__name__)

# Shopify names the sole option of a single-variant product "Title" with the
# value "Default Title". It is a placeholder, not data.
DROPPED_KEYS = {"title"}

KEY_ALIASES = {
    "colour": "color", "color": "color", "shade": "color", "colourway": "color",
    "size": "size", "sizes": "size", "size (uk)": "size", "fit": "size",
    "material": "material", "fabric": "material", "composition": "material",
    "gender": "gender", "sex": "gender", "for": "gender", "department": "gender",
    "style": "style", "design": "style", "pattern": "pattern",
    "capacity": "capacity", "volume": "capacity",
    "weight": "weight", "length": "length", "width": "width",
    "flavour": "flavor", "flavor": "flavor",
    "scent": "scent", "fragrance": "scent",
}

COLOR_BASE = {
    "navy": "blue", "midnight blue": "blue", "cobalt": "blue", "teal": "blue",
    "charcoal": "grey", "slate": "grey", "gunmetal": "grey", "graphite": "grey",
    "ivory": "white", "cream": "white", "ecru": "white", "off white": "white",
    "burgundy": "red", "maroon": "red", "wine": "red", "crimson": "red",
    "tan": "brown", "camel": "brown", "khaki": "brown", "beige": "brown",
}

_ALPHA_SIZE = {
    "xs": "xs", "extrasmall": "xs", "s": "s", "small": "s", "m": "m",
    "medium": "m", "l": "l", "large": "l", "xl": "xl", "extralarge": "xl",
    "xxl": "xxl", "2xl": "xxl",
}

GENDER = {
    "mens": "male", "men": "male", "man": "male", "male": "male", "boys": "male",
    "womens": "female", "women": "female", "woman": "female", "female": "female",
    "ladies": "female", "girls": "female",
    "unisex": "unisex", "all": "unisex", "adult": "unisex", "kids": "kids",
}

_TO_BASE = {"ml": 1.0, "l": 1000.0, "g": 1.0, "kg": 1000.0,
            "mm": 1.0, "cm": 10.0, "m": 1000.0}
_BASE_UNIT = {"ml": "ml", "l": "ml", "g": "g", "kg": "g",
              "mm": "mm", "cm": "mm", "m": "mm"}

_KV_TAG = re.compile(r"\A([a-z][a-z0-9 _-]*?)\s*[:=]\s*(.+)\Z", re.I)
_REGIONAL = re.compile(r"\A(uk|us|eu|jp)?[- ]?(\d+(?:\.\d+)?)\Z")
_MEASURED = re.compile(r"\A(\d+(?:\.\d+)?)\s*(ml|l|g|kg|oz|lb|cm|mm|m|in)\Z")

# Structured data wins. A metafield is an explicit merchant statement; a bare
# tag is a guess from a keyword.
SOURCE_CONFIDENCE = {"metafield": 1.0, "option": 1.0, "tag_kv": 0.9,
                     "tag": 0.7, "inferred": 0.5}


def normalise_key(k):
    c = normalise_text(k)
    if not c:
        return None
    c = c.lower()
    if c in DROPPED_KEYS:
        return None
    if c in KEY_ALIASES:
        return KEY_ALIASES[c]
    return re.sub(r"[^a-z0-9]+", "_", c).strip("_") or None


def _normalise_colour(raw):
    v = normalise_text(raw)
    if not v:
        return None
    v = v.lower()
    first = re.split(r"[/,|+]| and ", v)[0].strip()
    if first in COLOR_BASE:
        return COLOR_BASE[first]
    for k, base in COLOR_BASE.items():
        if k in first:
            return base
    return first or None


def _normalise_size(raw):
    v = normalise_text(raw)
    if not v:
        return None, None
    flat = v.lower().replace(" ", "")

    if flat in _ALPHA_SIZE:
        return _ALPHA_SIZE[flat], "alpha"

    m = _MEASURED.match(flat)
    if m and m.group(2) in _TO_BASE:
        amount = float(m.group(1)) * _TO_BASE[m.group(2)]
        return f"{amount:g}", _BASE_UNIT[m.group(2)]

    m = _REGIONAL.match(flat)
    if m:
        return m.group(2), m.group(1) or "numeric"

    return flat, None


def normalise_value(key, raw):
    if key == "color":
        return _normalise_colour(raw), None
    if key == "size":
        return _normalise_size(raw)
    if key == "gender":
        v = normalise_text(raw)
        return (GENDER.get(v.lower().replace("'", "")) if v else None), None
    v = normalise_text(raw)
    return (v.lower() if v else None), None


def _attr(key, raw_value, source):
    value, unit = normalise_value(key, raw_value)
    if value is None:
        return None
    return {"key": key, "raw_value": normalise_text(raw_value), "value": value,
            "unit": unit, "source": source,
            "confidence": SOURCE_CONFIDENCE[source]}


def build_attributes(source_product: dict) -> list:
    """Collect attributes from every source, keeping the most trusted per key."""
    found = []

    for a in source_product.get("attributes_raw") or []:
        key = normalise_key(a.get("key"))
        if key:
            attr = _attr(key, a.get("value"), a.get("source", "metafield"))
            if attr:
                found.append(attr)

    for variant in source_product.get("variants") or []:
        for name, value in (variant.get("options") or {}).items():
            key = normalise_key(name)
            if key:
                attr = _attr(key, value, "option")
                if attr:
                    found.append(attr)

    for tag in source_product.get("tags") or []:
        t = normalise_text(tag)
        if not t:
            continue
        m = _KV_TAG.match(t)
        if m:
            key = normalise_key(m.group(1))
            if key:
                attr = _attr(key, m.group(2), "tag_kv")
                if attr:
                    found.append(attr)

    best = {}
    for a in found:
        seen = best.get((a["key"], a["value"]))
        if seen is None or a["confidence"] > seen["confidence"]:
            best[(a["key"], a["value"])] = a

    by_key = {}
    for a in best.values():
        seen = by_key.get(a["key"])
        if seen is None or a["confidence"] > seen["confidence"]:
            by_key[a["key"]] = a

    return sorted(by_key.values(), key=lambda a: (a["key"], a["value"]))
```

- [ ] **Step 4: Run test to verify it passes**

Run: `.venv/bin/python -m pytest tests/unit/test_catalog_attributes.py -v`
Expected: PASS, 35 tests

- [ ] **Step 5: Report changed files**

```bash
git status --short
```

Do not commit and do not stage.

---

### Task 6: Commercial fields, stock, URL and media

**Files:**
- Create: `app/services/catalog_commercial.py`
- Test: `tests/unit/test_catalog_commercial.py`

**Interfaces:**
- Consumes: nothing.
- Produces:
  - `to_cents(price) -> int | None`
  - `price_range(variants) -> tuple[int | None, int | None, int | None, bool]` returning `(min_cents, max_cents, compare_at_cents, on_sale)`
  - `is_in_stock(source_product: dict) -> bool`
  - `resolve_product_url(source_product: dict, primary_domain=None, url_template=None) -> str | None`
  - `clean_image_url(url) -> str | None`

- [ ] **Step 1: Write the failing test**

Create `tests/unit/test_catalog_commercial.py`:

```python
import pytest

from app.services.catalog_commercial import (
    clean_image_url, is_in_stock, price_range, resolve_product_url, to_cents,
)


@pytest.mark.parametrize("raw,expected", [
    ("9.99", 999), (9.99, 999), ("885.95", 88595), ("0", 0),
    ("2629.95", 262995), (None, None), ("", None), ("abc", None),
])
def test_to_cents(raw, expected):
    assert to_cents(raw) == expected


def test_to_cents_never_returns_a_float():
    assert isinstance(to_cents("9.99"), int)


def test_price_range_uses_min_and_max_across_variants():
    # The live Gift Card spans 10.00 to 100.00; a single price is wrong for it.
    lo, hi, cmp_at, on_sale = price_range([
        {"price": "10.00", "compare_at_price": None},
        {"price": "50.00", "compare_at_price": None},
        {"price": "100.00", "compare_at_price": None},
    ])
    assert (lo, hi) == (1000, 10000)
    assert cmp_at is None
    assert on_sale is False


def test_single_variant_has_equal_min_and_max():
    lo, hi, _, _ = price_range([{"price": "949.95", "compare_at_price": None}])
    assert lo == hi == 94995


def test_on_sale_when_compare_at_exceeds_price():
    lo, hi, cmp_at, on_sale = price_range([
        {"price": "80.00", "compare_at_price": "100.00"}])
    assert on_sale is True
    assert cmp_at == 10000


def test_compare_at_below_price_is_not_a_sale():
    _, _, _, on_sale = price_range([
        {"price": "100.00", "compare_at_price": "80.00"}])
    assert on_sale is False


def test_no_variants_yields_none():
    assert price_range([]) == (None, None, None, False)


def test_in_stock_from_availability_not_quantity():
    # Overselling merchants carry negative quantities on purchasable products.
    assert is_in_stock({"tracks_inventory": True, "variants": [
        {"available": True, "inventory_quantity": -5}]}) is True
    assert is_in_stock({"tracks_inventory": True, "variants": [
        {"available": False, "inventory_quantity": 100}]}) is False


def test_untracked_inventory_is_in_stock():
    assert is_in_stock({"tracks_inventory": False, "variants": [
        {"available": False, "inventory_quantity": 0}]}) is True


def test_any_available_variant_makes_the_product_in_stock():
    assert is_in_stock({"tracks_inventory": True, "variants": [
        {"available": False}, {"available": True}]}) is True


def test_url_prefers_the_explicit_one():
    assert resolve_product_url(
        {"product_url": "https://shop.example.com/products/x", "handle": "y"},
        primary_domain="https://shop.example.com") == \
        "https://shop.example.com/products/x"


def test_url_falls_back_to_handle():
    # All 17 live products have onlineStoreUrl null, so this is the only path.
    assert resolve_product_url(
        {"product_url": None, "handle": "the-minimal-snowboard"},
        primary_domain="https://galaxiq-braexaal.myshopify.com") == \
        "https://galaxiq-braexaal.myshopify.com/products/the-minimal-snowboard"


def test_url_template_used_for_http_sources():
    assert resolve_product_url(
        {"product_url": None, "handle": None, "external_id": 42},
        url_template="https://shop.example.com/product/{external_id}") == \
        "https://shop.example.com/product/42"


def test_url_is_none_when_nothing_resolves():
    assert resolve_product_url({"product_url": None, "handle": None}) is None


@pytest.mark.parametrize("raw,expected", [
    ("https://cdn.shopify.com/x/snowboard.jpg?v=1786334565",
     "https://cdn.shopify.com/x/snowboard.jpg"),
    ("https://cdn.shopify.com/x/snowboard_800x.jpg",
     "https://cdn.shopify.com/x/snowboard.jpg"),
    ("https://cdn.shopify.com/x/snowboard_800x600.jpg?v=1",
     "https://cdn.shopify.com/x/snowboard.jpg"),
    (None, None),
])
def test_clean_image_url(raw, expected):
    # Size suffixes and cache-busting params change between syncs and would
    # churn record_hash for no reason.
    assert clean_image_url(raw) == expected
```

- [ ] **Step 2: Run test to verify it fails**

Run: `.venv/bin/python -m pytest tests/unit/test_catalog_commercial.py -v`
Expected: FAIL with `ModuleNotFoundError: No module named 'app.services.catalog_commercial'`

- [ ] **Step 3: Write the implementation**

Create `app/services/catalog_commercial.py`:

```python
"""Price, stock, URL and media normalisation."""
import logging
import re

logger = logging.getLogger(__name__)

_CDN_SIZE = re.compile(r"_(\d+x\d*|\d*x\d+)(?=\.\w+\Z)")


def to_cents(price):
    """Integer minor units. Floats accumulate error across a catalog and make
    equality comparisons unreliable."""
    if price is None or price == "":
        return None
    try:
        return int(round(float(price) * 100))
    except (TypeError, ValueError):
        return None


def price_range(variants):
    prices = [c for c in (to_cents(v.get("price")) for v in variants or [])
              if c is not None]
    if not prices:
        return None, None, None, False

    lo, hi = min(prices), max(prices)

    compare_at, on_sale = None, False
    for v in variants:
        cmp_c, price_c = to_cents(v.get("compare_at_price")), to_cents(v.get("price"))
        if cmp_c is not None and price_c is not None and cmp_c > price_c:
            on_sale = True
            compare_at = cmp_c if compare_at is None else max(compare_at, cmp_c)

    return lo, hi, compare_at, on_sale


def is_in_stock(source_product: dict) -> bool:
    """Availability flags only. Merchants with overselling enabled carry
    negative quantities on perfectly purchasable products."""
    if not source_product.get("tracks_inventory"):
        return True
    return any(v.get("available") for v in source_product.get("variants") or [])


def resolve_product_url(source_product: dict, primary_domain=None, url_template=None):
    explicit = source_product.get("product_url")
    if explicit:
        return explicit

    handle = source_product.get("handle")
    if handle and primary_domain:
        return f"{primary_domain.rstrip('/')}/products/{handle}"

    if url_template:
        try:
            return url_template.format(**source_product)
        except (KeyError, IndexError):
            logger.warning("url_template does not match the record's fields")
            return None

    return None


def clean_image_url(url):
    if not url:
        return None
    return _CDN_SIZE.sub("", url.split("?")[0])
```

- [ ] **Step 4: Run test to verify it passes**

Run: `.venv/bin/python -m pytest tests/unit/test_catalog_commercial.py -v`
Expected: PASS, 27 tests

- [ ] **Step 5: Report changed files**

```bash
git status --short
```

Do not commit and do not stage.

---

### Task 7: Normalisation orchestrator, hashes, quality score, validation

**Files:**
- Create: `app/services/catalog_normalise.py`
- Test: `tests/unit/test_catalog_normalise.py`

**Interfaces:**
- Consumes: `catalog_text`, `catalog_taxonomy`, `catalog_attributes`, `catalog_commercial`.
- Produces:
  - `normalise(source_product: dict, context: dict) -> dict | None` — returns `None` when the product fails validation. `context` carries `{"shop_name", "primary_domain", "url_template", "currency", "known_brands"}`.
  - `quality_score(product: dict) -> float`
  - `content_hash(product: dict) -> str`
  - `record_hash(product: dict) -> str`
  - `REJECT_REASONS: tuple[str, ...]`

- [ ] **Step 1: Write the failing test**

Create `tests/unit/test_catalog_normalise.py`:

```python
import copy

import pytest

from app.services.catalog_normalise import (
    content_hash, normalise, quality_score, record_hash,
)

CONTEXT = {
    "shop_name": "Galaxiq",
    "primary_domain": "https://galaxiq-braexaal.myshopify.com",
    "url_template": None,
    "currency": "INR",
    "known_brands": [],
}


def sp(**kw):
    base = {
        "external_id": "gid://shopify/Product/1",
        "title": "The Minimal Snowboard",
        "description": None,
        "brand": "Galaxiq",
        "handle": "the-minimal-snowboard",
        "product_url": None,
        "image_url": None,
        "raw_category": None,
        "product_type": "snowboard",
        "tags": ["Sport", "Winter"],
        "collections": [],
        "status": "ACTIVE",
        "tracks_inventory": True,
        "variants": [{"price": "885.95", "compare_at_price": None,
                      "available": True, "options": {"Title": "Default Title"}}],
        "attributes_raw": [],
        "source_kind": "shopify",
        "source_ref": "galaxiq-braexaal.myshopify.com",
    }
    base.update(kw)
    return base


def test_happy_path():
    out = normalise(sp(), CONTEXT)
    assert out["product_key"] == "shopify:gid://shopify/Product/1"
    assert out["name"] == "The Minimal Snowboard"
    assert out["price_cents"] == 88595
    assert out["price_max_cents"] == 88595
    assert out["currency"] == "INR"
    assert out["in_stock"] is True
    assert out["status"] == "ACTIVE"
    assert out["taxonomy_path"] == ["snowboard"]
    assert out["taxonomy_source"] == "product_type"
    assert out["product_url"] == \
        "https://galaxiq-braexaal.myshopify.com/products/the-minimal-snowboard"


def test_shop_name_brand_is_nulled():
    assert normalise(sp(), CONTEXT)["brand"] is None


def test_title_placeholder_option_produces_no_attribute():
    assert normalise(sp(), CONTEXT)["attributes"] == []


@pytest.mark.parametrize("bad", [
    {"external_id": None},
    {"title": None},
    {"title": "ab"},
    {"variants": []},
    {"handle": None, "product_url": None},
])
def test_validation_rejects(bad):
    assert normalise(sp(**bad), CONTEXT) is None


def test_negative_price_rejected():
    assert normalise(sp(variants=[{"price": "-1.00", "available": True}]),
                     CONTEXT) is None


def test_missing_brand_and_category_degrade_rather_than_reject():
    out = normalise(sp(brand=None, product_type=None, tags=[]), CONTEXT)
    assert out is not None
    assert out["brand"] is None
    assert out["taxonomy_source"] == "none"


def test_normalising_twice_is_identical():
    a = normalise(sp(), CONTEXT)
    b = normalise(sp(), CONTEXT)
    assert a == b


def test_content_hash_is_stable_across_tag_order():
    # Unsorted hashing changes the hash every sync and re-embeds the whole
    # catalog nightly. This is the most expensive normalisation bug available.
    a = normalise(sp(tags=["Sport", "Winter"]), CONTEXT)
    b = normalise(sp(tags=["Winter", "Sport"]), CONTEXT)
    assert content_hash(a) == content_hash(b)


def test_content_hash_changes_when_matching_content_changes():
    a = normalise(sp(), CONTEXT)
    b = normalise(sp(title="A Different Snowboard"), CONTEXT)
    assert content_hash(a) != content_hash(b)


def test_price_change_does_not_change_content_hash():
    # A price change must update the row without triggering a re-embed.
    a = normalise(sp(), CONTEXT)
    b = normalise(sp(variants=[{"price": "999.00", "compare_at_price": None,
                                "available": True, "options": {}}]), CONTEXT)
    assert content_hash(a) == content_hash(b)
    assert record_hash(a) != record_hash(b)


def test_no_field_is_ever_an_empty_string():
    out = normalise(sp(description="   ", brand="  "), CONTEXT)
    assert all(v != "" for v in out.values())


def test_price_min_never_exceeds_max():
    out = normalise(sp(variants=[
        {"price": "100.00", "available": True, "options": {}},
        {"price": "10.00", "available": True, "options": {}},
    ]), CONTEXT)
    assert out["price_cents"] <= out["price_max_cents"]


def test_quality_score_range():
    out = normalise(sp(), CONTEXT)
    assert 0.0 <= quality_score(out) <= 1.0


def test_quality_score_rewards_completeness():
    sparse = normalise(sp(brand=None, product_type=None, tags=[],
                          image_url=None, description=None), CONTEXT)
    rich = normalise(sp(brand="Snowboard Vendor", description="A good board.",
                        image_url="https://cdn.example.com/x.jpg",
                        raw_category="Sporting Goods > Snowboards",
                        tags=["color:Navy", "size:M"]), CONTEXT)
    assert quality_score(rich) > quality_score(sparse)
```

- [ ] **Step 2: Run test to verify it fails**

Run: `.venv/bin/python -m pytest tests/unit/test_catalog_normalise.py -v`
Expected: FAIL with `ModuleNotFoundError: No module named 'app.services.catalog_normalise'`

- [ ] **Step 3: Write the implementation**

Create `app/services/catalog_normalise.py`:

```python
"""Turns a SourceProduct into the row shape strategist_products stores.

Deterministic and pure: no network, no database, no LLM. Every rule here is
reproducible from a fixture, which is what makes dictionary changes safe to make.
"""
import hashlib
import json
import logging

from app.services.catalog_attributes import build_attributes
from app.services.catalog_commercial import (
    clean_image_url, is_in_stock, price_range, resolve_product_url,
)
from app.services.catalog_taxonomy import resolve_taxonomy
from app.services.catalog_text import clean_title, normalise_brand, normalise_text

logger = logging.getLogger(__name__)

MIN_TITLE_LENGTH = 3

REJECT_REASONS = (
    "no_external_id", "no_title", "title_too_short",
    "no_price", "negative_price", "no_product_url", "no_variants",
)


def normalise(source_product: dict, context: dict):
    external_id = source_product.get("external_id")
    if external_id in (None, ""):
        return None

    variants = source_product.get("variants") or []
    if not variants:
        return None

    brand = normalise_brand(source_product.get("brand"),
                            context.get("shop_name"),
                            context.get("known_brands") or ())

    name = clean_title(source_product.get("title"), brand)
    if not name or len(name) < MIN_TITLE_LENGTH:
        return None

    lo, hi, compare_at, on_sale = price_range(variants)
    if lo is None or lo < 0:
        return None

    url = resolve_product_url(source_product,
                              context.get("primary_domain"),
                              context.get("url_template"))
    if not url:
        return None

    taxonomy = resolve_taxonomy(source_product)
    attributes = build_attributes(source_product)
    tags = sorted({t for t in (normalise_text(x) for x in
                               source_product.get("tags") or []) if t})

    return {
        "product_key": f"{source_product['source_kind']}:{external_id}",
        "source_kind": source_product["source_kind"],
        "source_ref": source_product["source_ref"],
        "external_id": str(external_id),
        "name": name,
        "description": normalise_text(source_product.get("description")),
        "brand": brand,
        "image_url": clean_image_url(source_product.get("image_url")),
        "product_url": url,
        "taxonomy_path": taxonomy["path"],
        "taxonomy_source": taxonomy["source"],
        "raw_category": taxonomy["raw"],
        "category": taxonomy["path"][-1] if taxonomy["path"] else None,
        "tags": tags,
        "attributes": attributes,
        "price_cents": lo,
        "price_max_cents": hi,
        "compare_at_cents": compare_at,
        "currency": context.get("currency"),
        "on_sale": on_sale,
        "in_stock": is_in_stock(source_product),
        "status": normalise_text(source_product.get("status")),
    }


def quality_score(product: dict) -> float:
    s = 0.0
    name = product.get("name") or ""
    s += 0.25 if len(name) >= 10 else 0.0
    s += 0.25 if len(product.get("taxonomy_path") or []) >= 2 else 0.0
    s += 0.15 if product.get("brand") else 0.0
    s += min(len(product.get("attributes") or []), 4) / 4 * 0.20
    s += 0.10 if product.get("image_url") else 0.0
    s += 0.05 if product.get("description") else 0.0
    return round(s, 2)


def content_hash(product: dict) -> str:
    """Covers only fields that affect matching, so a price change does not
    trigger a needless re-embed. Arrays are sorted: unsorted hashing changes on
    every sync and re-embeds the entire catalog."""
    parts = [
        product.get("name") or "",
        product.get("brand") or "",
        ">".join(product.get("taxonomy_path") or []),
        ",".join(sorted(product.get("tags") or [])),
        ",".join(sorted(f"{a['key']}={a['value']}"
                        for a in product.get("attributes") or [])),
        product.get("description") or "",
    ]
    return hashlib.sha256("|".join(parts).encode()).hexdigest()


def record_hash(product: dict) -> str:
    return hashlib.sha256(
        json.dumps(product, sort_keys=True, default=str).encode()).hexdigest()
```

- [ ] **Step 4: Run test to verify it passes**

Run: `.venv/bin/python -m pytest tests/unit/test_catalog_normalise.py -v`
Expected: PASS, 20 tests

- [ ] **Step 5: Report changed files**

```bash
git status --short
```

Do not commit and do not stage.

---

### Task 8: Shopify adapter and golden fixtures

**Files:**
- Create: `app/services/catalog_shopify.py`
- Test: `tests/unit/test_catalog_shopify.py`
- Read (already captured, do not regenerate): `tests/fixtures/shopify/*.json`

**Interfaces:**
- Consumes: `shopify_oauth.SHOPIFY_API_VERSION`, `product_sources.get_credentials`.
- Produces:
  - `PRODUCTS_QUERY: str`
  - `to_source_product(node: dict, shop_domain: str) -> dict`
  - `async fetch_products(shop: str, token: str, *, client=None) -> list[dict]`

- [ ] **Step 1: Confirm the fixtures exist**

```bash
ls tests/fixtures/shopify/
```

Expected: `all-products.json` plus `no-category.json`, `vendor-is-shop-name.json`, `title-placeholder-option.json`, `multi-variant-price-spread.json`, `out-of-stock.json`, `archived.json`, `draft.json`, `no-image.json`, `empty-product-type.json`, `native-category.json`, `has-metafields.json`, `no-online-store-url.json`.

These were captured from the live store. Do not regenerate or edit them.

- [ ] **Step 2: Write the failing test**

Create `tests/unit/test_catalog_shopify.py`:

```python
import json
import pathlib

import httpx
import pytest

from app.services.catalog_shopify import (
    PRODUCTS_QUERY, fetch_products, to_source_product,
)
from app.services.catalog_normalise import normalise

FIXTURES = pathlib.Path("tests/fixtures/shopify")
SHOP = "galaxiq-braexaal.myshopify.com"

CONTEXT = {
    "shop_name": "Galaxiq",
    "primary_domain": "https://galaxiq-braexaal.myshopify.com",
    "url_template": None,
    "currency": "INR",
    "known_brands": [],
}


def load(name):
    return json.loads((FIXTURES / f"{name}.json").read_text())


def test_query_requests_every_field_normalisation_needs():
    for field in ["category", "productType", "vendor", "tags", "handle",
                  "status", "featuredImage", "collections", "tracksInventory",
                  "variants", "compareAtPrice", "availableForSale",
                  "selectedOptions", "metafields"]:
        assert field in PRODUCTS_QUERY


def test_maps_a_real_node():
    node = load("no-category")
    sp = to_source_product(node, SHOP)
    assert sp["source_kind"] == "shopify"
    assert sp["source_ref"] == SHOP
    assert sp["external_id"] == node["id"]
    assert sp["title"] == node["title"]
    assert sp["handle"] == node["handle"]
    assert sp["raw_category"] is None
    assert sp["variants"][0]["price"] == node["variants"]["nodes"][0]["price"]


def test_native_category_is_flattened():
    sp = to_source_product(load("native-category"), SHOP)
    assert sp["raw_category"] == "Gift Cards"


def test_collections_flattened_to_titles():
    sp = to_source_product(load("has-metafields"), SHOP)
    assert all(isinstance(c, str) for c in sp["collections"])


def test_variant_options_flattened_to_a_dict():
    sp = to_source_product(load("title-placeholder-option"), SHOP)
    assert sp["variants"][0]["options"]["Title"] == "Default Title"


def test_missing_image_maps_to_none():
    sp = to_source_product(load("no-image"), SHOP)
    assert sp["image_url"] is None


def test_every_fixture_normalises_without_raising():
    for node in load("all-products"):
        normalise(to_source_product(node, SHOP), CONTEXT)


def test_the_whole_live_catalog_survives_validation():
    nodes = load("all-products")
    out = [normalise(to_source_product(n, SHOP), CONTEXT) for n in nodes]
    kept = [p for p in out if p is not None]
    # Every live product has a handle, so the URL fallback resolves for all.
    assert len(kept) == len(nodes), "some real products were rejected"


def test_shop_name_vendor_nulled_across_the_real_catalog():
    nodes = load("all-products")
    products = [normalise(to_source_product(n, SHOP), CONTEXT) for n in nodes]
    galaxiq_vendors = [n for n in nodes if n["vendor"] == "Galaxiq"]
    assert len(galaxiq_vendors) == 8, "fixture drifted from the profiled store"
    assert all(p["brand"] is None for p, n in zip(products, nodes)
               if n["vendor"] == "Galaxiq")


def test_no_product_gains_a_title_attribute():
    for node in load("all-products"):
        p = normalise(to_source_product(node, SHOP), CONTEXT)
        assert all(a["key"] != "title" for a in p["attributes"])


def test_archived_and_draft_are_preserved_not_dropped():
    # Status filtering belongs at recommendation time, not ingest: an archived
    # product may still be the anchor a visitor is looking at.
    for name, expected in [("archived", "ARCHIVED"), ("draft", "DRAFT")]:
        p = normalise(to_source_product(load(name), SHOP), CONTEXT)
        assert p is not None
        assert p["status"] == expected


def test_out_of_stock_detected():
    p = normalise(to_source_product(load("out-of-stock"), SHOP), CONTEXT)
    assert p["in_stock"] is False


def test_multi_variant_spread_produces_a_range():
    p = normalise(to_source_product(load("multi-variant-price-spread"), SHOP),
                  CONTEXT)
    assert p["price_cents"] < p["price_max_cents"]


def test_url_built_from_handle_since_no_product_has_online_store_url():
    for node in load("all-products"):
        assert node["onlineStoreUrl"] is None
        p = normalise(to_source_product(node, SHOP), CONTEXT)
        assert p["product_url"].startswith(
            "https://galaxiq-braexaal.myshopify.com/products/")


@pytest.mark.asyncio
async def test_fetch_products_follows_pagination():
    pages = [
        {"data": {"products": {
            "nodes": [{"id": "gid://1"}],
            "pageInfo": {"hasNextPage": True, "endCursor": "CUR"}}}},
        {"data": {"products": {
            "nodes": [{"id": "gid://2"}],
            "pageInfo": {"hasNextPage": False, "endCursor": None}}}},
    ]
    seen = []

    def handler(request):
        body = json.loads(request.content)
        seen.append(body["variables"]["cursor"])
        return httpx.Response(200, json=pages[len(seen) - 1])

    async with httpx.AsyncClient(transport=httpx.MockTransport(handler)) as c:
        out = await fetch_products(SHOP, "shpat_x", client=c)

    assert [n["id"] for n in out] == ["gid://1", "gid://2"]
    assert seen == [None, "CUR"]


@pytest.mark.asyncio
async def test_fetch_products_raises_on_graphql_errors():
    def handler(request):
        return httpx.Response(200, json={"errors": [{"message": "Access denied"}]})

    from app.services.shopify_oauth import ShopQueryError
    async with httpx.AsyncClient(transport=httpx.MockTransport(handler)) as c:
        with pytest.raises(ShopQueryError):
            await fetch_products(SHOP, "shpat_x", client=c)
```

- [ ] **Step 3: Run test to verify it fails**

Run: `.venv/bin/python -m pytest tests/unit/test_catalog_shopify.py -v`
Expected: FAIL with `ModuleNotFoundError: No module named 'app.services.catalog_shopify'`

- [ ] **Step 4: Write the implementation**

Create `app/services/catalog_shopify.py`:

```python
"""Fetches a Shopify catalog and maps it to SourceProduct.

Shopify's schema is fixed, so this mapping is code rather than a declarative
field map -- only heterogeneous HTTP sources need configuration.
"""
import logging

import httpx

from app.services.shopify_oauth import SHOPIFY_API_VERSION, ShopQueryError

logger = logging.getLogger(__name__)

PAGE_SIZE = 250

PRODUCTS_QUERY = """
query($cursor: String) {
  products(first: %d, after: $cursor) {
    nodes {
      id title handle status productType vendor tags onlineStoreUrl description
      category { fullName }
      featuredImage { url }
      collections(first: 20) { nodes { title } }
      tracksInventory
      variants(first: 100) {
        nodes {
          sku price compareAtPrice availableForSale inventoryQuantity
          selectedOptions { name value }
        }
      }
      metafields(first: 20) { nodes { namespace key value } }
    }
    pageInfo { hasNextPage endCursor }
  }
}
""" % PAGE_SIZE


def to_source_product(node: dict, shop_domain: str) -> dict:
    return {
        "external_id": node.get("id"),
        "title": node.get("title"),
        "description": node.get("description"),
        "brand": node.get("vendor"),
        "handle": node.get("handle"),
        "product_url": node.get("onlineStoreUrl"),
        "image_url": (node.get("featuredImage") or {}).get("url"),
        "raw_category": (node.get("category") or {}).get("fullName"),
        "product_type": node.get("productType"),
        "tags": node.get("tags") or [],
        "collections": [c["title"] for c in
                        (node.get("collections") or {}).get("nodes", [])],
        "status": node.get("status"),
        "tracks_inventory": node.get("tracksInventory"),
        "variants": [
            {
                "sku": v.get("sku"),
                "price": v.get("price"),
                "compare_at_price": v.get("compareAtPrice"),
                "available": v.get("availableForSale"),
                "inventory_quantity": v.get("inventoryQuantity"),
                "options": {o["name"]: o["value"]
                            for o in v.get("selectedOptions") or []},
            }
            for v in (node.get("variants") or {}).get("nodes", [])
        ],
        "attributes_raw": [
            {"key": m["key"], "value": m["value"], "source": "metafield"}
            for m in (node.get("metafields") or {}).get("nodes", [])
        ],
        "source_kind": "shopify",
        "source_ref": shop_domain,
    }


async def fetch_products(shop: str, token: str, *, client=None) -> list:
    """Every page or nothing: a partial fetch must not be mistaken for a
    complete catalog, or the sync would delete products it simply never saw."""
    owned = client is None
    c = client or httpx.AsyncClient(timeout=60.0)
    nodes, cursor = [], None
    try:
        while True:
            resp = await c.post(
                f"https://{shop}/admin/api/{SHOPIFY_API_VERSION}/graphql.json",
                headers={"X-Shopify-Access-Token": token,
                         "Content-Type": "application/json"},
                json={"query": PRODUCTS_QUERY, "variables": {"cursor": cursor}},
            )
            if resp.status_code != 200:
                logger.error("Catalog fetch for %s failed: HTTP %s",
                             shop, resp.status_code)
                raise ShopQueryError(f"HTTP {resp.status_code}")

            body = resp.json()
            if body.get("errors"):
                logger.error("Catalog fetch for %s returned GraphQL errors", shop)
                raise ShopQueryError(str(body["errors"]))

            page = body["data"]["products"]
            nodes.extend(page["nodes"])
            if not page["pageInfo"]["hasNextPage"]:
                return nodes
            cursor = page["pageInfo"]["endCursor"]
    finally:
        if owned:
            await c.aclose()
```

- [ ] **Step 5: Run test to verify it passes**

Run: `.venv/bin/python -m pytest tests/unit/test_catalog_shopify.py -v`
Expected: PASS, 16 tests

- [ ] **Step 6: Report changed files**

```bash
git status --short
```

Do not commit and do not stage.

---

### Task 9: HTTP API adapter

**Files:**
- Create: `app/services/catalog_http.py`
- Create: `tests/fixtures/http/dummyjson-page.json`
- Test: `tests/unit/test_catalog_http.py`

**Interfaces:**
- Consumes: `catalog_fieldmap.apply_field_map`, `catalog_fieldmap.FieldMapError`.
- Produces:
  - `to_source_product(record: dict, config: dict, source_ref: str) -> dict`
  - `async fetch_products(config: dict, credentials: dict, source_ref: str, *, client=None) -> list[dict]`
  - `class HttpSourceError(Exception)`
  - `auth_headers(credentials: dict, config: dict) -> dict`

- [ ] **Step 1: Create the fixture**

Create `tests/fixtures/http/dummyjson-page.json` with a two-record page in the reference API's shape:

```json
{
  "products": [
    {
      "id": 1,
      "title": "Essence Mascara Lash Princess",
      "description": "A popular mascara known for its volumizing effects.",
      "category": "beauty",
      "price": 9.99,
      "discountPercentage": 10.48,
      "rating": 2.56,
      "stock": 99,
      "tags": ["beauty", "mascara"],
      "brand": "Essence",
      "sku": "BEA-ESS-ESS-001",
      "availabilityStatus": "In Stock",
      "images": ["https://cdn.dummyjson.com/1.webp"],
      "thumbnail": "https://cdn.dummyjson.com/thumbnail.webp"
    },
    {
      "id": 2,
      "title": "Eyeshadow Palette with Mirror",
      "description": "A versatile range of eyeshadow shades.",
      "category": "beauty",
      "price": 19.99,
      "discountPercentage": 18.19,
      "stock": 0,
      "tags": ["beauty", "eyeshadow"],
      "brand": "Glamour Beauty",
      "availabilityStatus": "Out of Stock",
      "images": ["https://cdn.dummyjson.com/2.webp"],
      "thumbnail": "https://cdn.dummyjson.com/thumb2.webp"
    }
  ],
  "total": 2,
  "skip": 0,
  "limit": 30
}
```

- [ ] **Step 2: Write the failing test**

Create `tests/unit/test_catalog_http.py`:

```python
import json
import pathlib

import httpx
import pytest

from app.services.catalog_http import (
    HttpSourceError, auth_headers, fetch_products, to_source_product,
)
from app.services.catalog_normalise import normalise

PAGE = json.loads(
    pathlib.Path("tests/fixtures/http/dummyjson-page.json").read_text())

CONFIG = {
    "base_url": "https://dummyjson.com/products",
    "records_path": "$.products",
    "pagination": {"style": "offset", "limit_param": "limit",
                   "offset_param": "skip", "page_size": 30,
                   "total_path": "$.total"},
    "fields": {
        "external_id": "$.id",
        "title": "$.title",
        "description": "$.description",
        "brand": "$.brand",
        "raw_category": "$.category",
        "image_url": "$.thumbnail",
        "price": "$.price",
        "discount_percentage": "$.discountPercentage",
        "availability": "$.availabilityStatus",
        "tags": "$.tags",
    },
    "url_template": "https://shop.example.com/product/{external_id}",
    "currency": "USD",
    "auth": {"style": "bearer"},
}

CONTEXT = {"shop_name": None, "primary_domain": None,
           "url_template": CONFIG["url_template"], "currency": "USD",
           "known_brands": []}


def test_maps_a_record():
    sp = to_source_product(PAGE["products"][0], CONFIG, "dummyjson.com")
    assert sp["external_id"] == 1
    assert sp["title"] == "Essence Mascara Lash Princess"
    assert sp["brand"] == "Essence"
    assert sp["raw_category"] == "beauty"
    assert sp["source_kind"] == "http_api"
    assert sp["source_ref"] == "dummyjson.com"


def test_price_becomes_a_single_synthetic_variant():
    # The reference API has no variants; one record is one product.
    sp = to_source_product(PAGE["products"][0], CONFIG, "dummyjson.com")
    assert len(sp["variants"]) == 1
    assert sp["variants"][0]["price"] == 9.99


def test_discount_percentage_derives_compare_at():
    sp = to_source_product(PAGE["products"][0], CONFIG, "dummyjson.com")
    p = normalise(sp, CONTEXT)
    assert p["on_sale"] is True
    # 9.99 / (1 - 0.1048) == 11.16
    assert p["compare_at_cents"] == 1116


def test_availability_string_maps_to_stock():
    in_stock = to_source_product(PAGE["products"][0], CONFIG, "dummyjson.com")
    out = to_source_product(PAGE["products"][1], CONFIG, "dummyjson.com")
    assert normalise(in_stock, CONTEXT)["in_stock"] is True
    assert normalise(out, CONTEXT)["in_stock"] is False


def test_url_comes_from_the_template():
    # The reference API returns no product URL at all; without a template every
    # record is rejected for having no resolvable URL.
    p = normalise(to_source_product(PAGE["products"][0], CONFIG,
                                    "dummyjson.com"), CONTEXT)
    assert p["product_url"] == "https://shop.example.com/product/1"


def test_missing_url_template_rejects_everything():
    ctx = dict(CONTEXT, url_template=None)
    p = normalise(to_source_product(PAGE["products"][0], CONFIG,
                                    "dummyjson.com"), ctx)
    assert p is None


def test_currency_is_never_defaulted():
    with pytest.raises(HttpSourceError):
        to_source_product(PAGE["products"][0],
                          dict(CONFIG, currency=None), "dummyjson.com")


def test_auth_headers_bearer():
    assert auth_headers({"api_key": "k"}, {"auth": {"style": "bearer"}}) == \
        {"Authorization": "Bearer k"}


def test_auth_headers_api_key_header():
    assert auth_headers(
        {"api_key": "k"},
        {"auth": {"style": "header", "header": "X-API-Key"}}) == {"X-API-Key": "k"}


def test_auth_headers_none():
    assert auth_headers({}, {"auth": {"style": "none"}}) == {}


@pytest.mark.asyncio
async def test_fetch_follows_offset_pagination():
    seen = []

    def handler(request):
        seen.append(str(request.url))
        skip = int(dict(request.url.params).get("skip", 0))
        if skip == 0:
            return httpx.Response(200, json={"products": [{"id": 1}], "total": 2})
        return httpx.Response(200, json={"products": [{"id": 2}], "total": 2})

    cfg = dict(CONFIG)
    cfg["pagination"] = dict(CONFIG["pagination"], page_size=1)
    async with httpx.AsyncClient(transport=httpx.MockTransport(handler)) as c:
        out = await fetch_products(cfg, {"api_key": "k"}, "dummyjson.com", client=c)

    assert [r["id"] for r in out] == [1, 2]
    assert len(seen) == 2


@pytest.mark.asyncio
async def test_fetch_raises_on_non_200():
    def handler(request):
        return httpx.Response(401, text="unauthorized")

    async with httpx.AsyncClient(transport=httpx.MockTransport(handler)) as c:
        with pytest.raises(HttpSourceError):
            await fetch_products(CONFIG, {"api_key": "k"}, "x", client=c)
```

- [ ] **Step 3: Run test to verify it fails**

Run: `.venv/bin/python -m pytest tests/unit/test_catalog_http.py -v`
Expected: FAIL with `ModuleNotFoundError: No module named 'app.services.catalog_http'`

- [ ] **Step 4: Write the implementation**

Create `app/services/catalog_http.py`:

```python
"""Fetches a catalog from an authenticated HTTP product API.

Unlike Shopify, the shape is unknown, so the mapping comes from the source's
stored config. C's LLM authors that config; this module only executes it.
"""
import logging

import httpx

from app.services.catalog_fieldmap import apply_field_map, resolve_path

logger = logging.getLogger(__name__)

_OUT_OF_STOCK = {"out of stock", "outofstock", "unavailable", "sold out"}


class HttpSourceError(Exception):
    """The source is misconfigured or unreachable."""


def auth_headers(credentials: dict, config: dict) -> dict:
    auth = config.get("auth") or {}
    style = auth.get("style", "none")
    if style == "none":
        return {}
    key = credentials.get("api_key") or credentials.get("access_token")
    if not key:
        raise HttpSourceError("auth style requires a key, none stored")
    if style == "bearer":
        return {"Authorization": f"Bearer {key}"}
    if style == "header":
        header = auth.get("header")
        if not header:
            raise HttpSourceError("header auth style requires a header name")
        return {header: key}
    raise HttpSourceError(f"unknown auth style: {style}")


def to_source_product(record: dict, config: dict, source_ref: str) -> dict:
    currency = config.get("currency")
    if not currency:
        # No API in this class reliably carries a currency, and guessing it
        # silently mis-prices an entire catalog.
        raise HttpSourceError("source config must declare a currency")

    mapped = apply_field_map(record, config["fields"])

    price = mapped.get("price")
    discount = mapped.get("discount_percentage")
    compare_at = None
    if price is not None and discount:
        try:
            compare_at = round(float(price) / (1 - float(discount) / 100), 2)
        except (TypeError, ValueError, ZeroDivisionError):
            compare_at = None

    availability = (mapped.get("availability") or "").strip().lower()
    available = availability not in _OUT_OF_STOCK

    return {
        "external_id": mapped.get("external_id"),
        "title": mapped.get("title"),
        "description": mapped.get("description"),
        "brand": mapped.get("brand"),
        "handle": None,
        "product_url": mapped.get("product_url"),
        "image_url": mapped.get("image_url"),
        "raw_category": mapped.get("raw_category"),
        "product_type": None,
        "tags": mapped.get("tags") or [],
        "collections": [],
        "status": "ACTIVE",
        "tracks_inventory": True,
        "variants": [{
            "sku": mapped.get("sku"),
            "price": price,
            "compare_at_price": compare_at,
            "available": available,
            "inventory_quantity": None,
            "options": {},
        }],
        "attributes_raw": [],
        "source_kind": "http_api",
        "source_ref": source_ref,
    }


async def fetch_products(config: dict, credentials: dict, source_ref: str,
                         *, client=None) -> list:
    pg = config.get("pagination") or {}
    if pg.get("style") != "offset":
        raise HttpSourceError(f"unsupported pagination style: {pg.get('style')}")

    base = config.get("base_url")
    if not base:
        raise HttpSourceError("source config must declare a base_url")

    size = int(pg.get("page_size", 30))
    headers = auth_headers(credentials, config)

    owned = client is None
    c = client or httpx.AsyncClient(timeout=60.0)
    records, offset = [], 0
    try:
        while True:
            resp = await c.get(base, headers=headers, params={
                pg.get("limit_param", "limit"): size,
                pg.get("offset_param", "skip"): offset,
            })
            if resp.status_code != 200:
                logger.error("Catalog fetch for %s failed: HTTP %s",
                             source_ref, resp.status_code)
                raise HttpSourceError(f"HTTP {resp.status_code}")

            body = resp.json()
            page = resolve_path(body, config["records_path"]) or []
            records.extend(page)

            total = resolve_path(body, pg["total_path"]) if pg.get("total_path") else None
            offset += size
            if not page or (total is not None and offset >= int(total)):
                return records
    finally:
        if owned:
            await c.aclose()
```

Note for the implementer: `resolve_product_url` formats `url_template` with `**source_product`, so `{external_id}` in a template resolves against the `external_id` key returned above.

- [ ] **Step 5: Run test to verify it passes**

Run: `.venv/bin/python -m pytest tests/unit/test_catalog_http.py -v`
Expected: PASS, 12 tests

- [ ] **Step 6: Report changed files**

```bash
git status --short
```

Do not commit and do not stage.

---

### Task 10: Sync orchestration and scoped persistence

The task that carries the §4 hazard. Every write and delete is scoped to one source.

**Files:**
- Create: `app/services/catalog_sync.py`
- Test: `tests/unit/test_catalog_sync.py` (orchestration, no DB)
- Test: `tests/integration/test_catalog_persistence.py` (the scoped delete, real DB)

**Interfaces:**
- Consumes: `product_sources.get_sources`, `product_sources.get_credentials`, `catalog_shopify`, `catalog_http`, `catalog_normalise`, `database.migrate_products_table`, `database.get_db_connection`.
- Produces:
  - `async sync_source(tenant_id: str, source: dict) -> dict` returning
    `{"fetched", "normalised", "rejected", "saved", "deleted", "rejection_rate", "quality_median"}`
  - `persist(tenant_id: str, products: list, source_kind: str, source_ref: str, run_started) -> dict`
  - `delete_stale(tenant_id, source_kind, source_ref, run_started) -> int`

- [ ] **Step 1: Write the failing test**

Create `tests/unit/test_catalog_sync.py`:

```python
import inspect

import pytest

from app.services import catalog_sync


def test_delete_never_runs_without_a_source_ref():
    # A blank ref would match every row of this kind across every source.
    # Checked before any connection is opened, so this needs no database.
    with pytest.raises(ValueError):
        catalog_sync.delete_stale("t", "shopify", None, "2026-01-01")
    with pytest.raises(ValueError):
        catalog_sync.delete_stale("t", "shopify", "", "2026-01-01")


@pytest.mark.asyncio
async def test_partial_fetch_deletes_nothing(monkeypatch):
    # A failure mid-fetch must not be mistaken for "these products are gone".
    async def boom(*a, **kw):
        raise RuntimeError("network died on page 2")

    deleted = []
    monkeypatch.setattr(catalog_sync.catalog_shopify, "fetch_products", boom)
    monkeypatch.setattr(catalog_sync, "delete_stale",
                        lambda *a, **kw: deleted.append(a))
    monkeypatch.setattr(catalog_sync, "get_credentials",
                        lambda *a: {"access_token": "shpat_x"})

    with pytest.raises(RuntimeError):
        await catalog_sync.sync_source("t", {
            "kind": "shopify", "external_ref": "s.myshopify.com",
            "config": {"shop_name": "S", "primary_domain": "https://s.com",
                       "currency": "INR"}})

    assert deleted == []


@pytest.mark.asyncio
async def test_rejected_products_are_counted_not_fatal(monkeypatch):
    async def two_nodes(*a, **kw):
        return [
            {"id": "gid://1", "title": "A Real Snowboard", "handle": "a",
             "status": "ACTIVE", "productType": "snowboard", "vendor": "V",
             "tags": [], "onlineStoreUrl": None, "description": None,
             "category": None, "featuredImage": None,
             "collections": {"nodes": []}, "tracksInventory": True,
             "variants": {"nodes": [{"price": "10.00", "compareAtPrice": None,
                                     "availableForSale": True, "sku": None,
                                     "inventoryQuantity": 1,
                                     "selectedOptions": []}]},
             "metafields": {"nodes": []}},
            {"id": None, "title": None, "handle": None, "status": "ACTIVE",
             "productType": None, "vendor": None, "tags": [],
             "onlineStoreUrl": None, "description": None, "category": None,
             "featuredImage": None, "collections": {"nodes": []},
             "tracksInventory": True, "variants": {"nodes": []},
             "metafields": {"nodes": []}},
        ]

    monkeypatch.setattr(catalog_sync.catalog_shopify, "fetch_products", two_nodes)
    monkeypatch.setattr(catalog_sync, "get_credentials",
                        lambda *a: {"access_token": "shpat_x"})
    monkeypatch.setattr(catalog_sync, "persist",
                        lambda *a, **kw: {"saved": 1})
    monkeypatch.setattr(catalog_sync, "delete_stale", lambda *a, **kw: 0)
    monkeypatch.setattr(catalog_sync, "migrate_products_table", lambda t: None)

    out = await catalog_sync.sync_source("t", {
        "kind": "shopify", "external_ref": "s.myshopify.com",
        "config": {"shop_name": "S", "primary_domain": "https://s.com",
                   "currency": "INR"}})

    assert out["fetched"] == 2
    assert out["normalised"] == 1
    assert out["rejected"] == 1
    assert out["rejection_rate"] == 0.5
```

- [ ] **Step 2: Run test to verify it fails**

Run: `.venv/bin/python -m pytest tests/unit/test_catalog_sync.py -v`
Expected: FAIL with `ModuleNotFoundError: No module named 'app.services.catalog_sync'`

- [ ] **Step 3: Write the implementation**

Create `app/services/catalog_sync.py`:

```python
"""Fetch, normalise, persist and report one source's catalog.

Every write and delete is scoped by (source_kind, source_ref). The crawler, a
Shopify store and an HTTP API share one table, so an unscoped delete here would
erase another producer's products.
"""
import logging
import statistics
from datetime import datetime, timezone

from psycopg2 import sql
from psycopg2.extras import Json

from app.services import catalog_http, catalog_shopify
from app.services.catalog_normalise import (
    content_hash, normalise, quality_score, record_hash,
)
from app.services.database import get_db_connection, migrate_products_table
from app.services.product_sources import get_credentials

logger = logging.getLogger(__name__)

REJECTION_ALERT_RATIO = 0.05


def delete_stale(tenant_id: str, source_kind: str, source_ref: str,
                 run_started) -> int:
    if not source_ref:
        # A blank ref would match every row of this kind across all sources.
        raise ValueError("delete_stale requires a source_ref")

    conn = get_db_connection()
    try:
        with conn.cursor() as cur:
            cur.execute(sql.SQL(
                "DELETE FROM {}.strategist_products "
                "WHERE source_kind = %s AND source_ref = %s "
                "AND (synced_at IS NULL OR synced_at < %s)"
            ).format(sql.Identifier(tenant_id)),
                (source_kind, source_ref, run_started))
            removed = cur.rowcount
        conn.commit()
        return removed
    except Exception:
        conn.rollback()
        logger.error("Could not delete stale products for %s", tenant_id,
                     exc_info=True)
        raise
    finally:
        conn.close()


def persist(tenant_id: str, products: list, source_kind: str, source_ref: str,
            run_started) -> dict:
    saved = 0
    conn = get_db_connection()
    try:
        with conn.cursor() as cur:
            for p in products:
                cur.execute(sql.SQL("""
                    INSERT INTO {}.strategist_products
                        (product_key, name, description, image_url, product_url,
                         category, source_kind, source_ref, external_id, brand,
                         taxonomy_path, taxonomy_source, raw_category,
                         price_cents, price_max_cents, compare_at_cents, currency,
                         on_sale, in_stock, status, attributes, quality_score,
                         content_hash, record_hash, synced_at, extracted_at)
                    VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,
                            %s,%s,%s,%s,%s,%s,%s,%s, CURRENT_TIMESTAMP)
                    ON CONFLICT (product_key) DO UPDATE SET
                        name = EXCLUDED.name,
                        description = EXCLUDED.description,
                        image_url = EXCLUDED.image_url,
                        product_url = EXCLUDED.product_url,
                        category = EXCLUDED.category,
                        source_kind = EXCLUDED.source_kind,
                        source_ref = EXCLUDED.source_ref,
                        external_id = EXCLUDED.external_id,
                        brand = EXCLUDED.brand,
                        taxonomy_path = EXCLUDED.taxonomy_path,
                        taxonomy_source = EXCLUDED.taxonomy_source,
                        raw_category = EXCLUDED.raw_category,
                        price_cents = EXCLUDED.price_cents,
                        price_max_cents = EXCLUDED.price_max_cents,
                        compare_at_cents = EXCLUDED.compare_at_cents,
                        currency = EXCLUDED.currency,
                        on_sale = EXCLUDED.on_sale,
                        in_stock = EXCLUDED.in_stock,
                        status = EXCLUDED.status,
                        attributes = EXCLUDED.attributes,
                        quality_score = EXCLUDED.quality_score,
                        content_hash = EXCLUDED.content_hash,
                        record_hash = EXCLUDED.record_hash,
                        synced_at = EXCLUDED.synced_at,
                        extracted_at = CURRENT_TIMESTAMP
                """).format(sql.Identifier(tenant_id)), (
                    p["product_key"], p["name"], p["description"], p["image_url"],
                    p["product_url"], p["category"], source_kind, source_ref,
                    p["external_id"], p["brand"], p["taxonomy_path"],
                    p["taxonomy_source"], p["raw_category"], p["price_cents"],
                    p["price_max_cents"], p["compare_at_cents"], p["currency"],
                    p["on_sale"], p["in_stock"], p["status"],
                    Json(p["attributes"]), quality_score(p),
                    content_hash(p), record_hash(p), run_started,
                ))
                saved += 1
        conn.commit()
        return {"saved": saved}
    except Exception:
        conn.rollback()
        logger.error("Could not persist products for %s", tenant_id, exc_info=True)
        raise
    finally:
        conn.close()


async def sync_source(tenant_id: str, source: dict) -> dict:
    run_started = datetime.now(timezone.utc)
    kind, ref, config = source["kind"], source["external_ref"], source["config"]

    migrate_products_table(tenant_id)
    credentials = get_credentials(kind, ref)
    if not credentials:
        raise ValueError(f"no credentials stored for {kind}:{ref}")

    if kind == "shopify":
        nodes = await catalog_shopify.fetch_products(
            ref, credentials["access_token"])
        source_products = [catalog_shopify.to_source_product(n, ref) for n in nodes]
        context = {
            "shop_name": config.get("shop_name"),
            "primary_domain": config.get("primary_domain"),
            "url_template": None,
            "currency": config.get("currency_code") or config.get("currency"),
            "known_brands": [],
        }
    elif kind == "http_api":
        records = await catalog_http.fetch_products(config, credentials, ref)
        source_products = [catalog_http.to_source_product(r, config, ref)
                           for r in records]
        context = {
            "shop_name": config.get("shop_name"),
            "primary_domain": None,
            "url_template": config.get("url_template"),
            "currency": config.get("currency"),
            "known_brands": [],
        }
    else:
        raise ValueError(f"unsupported source kind: {kind}")

    normalised = [n for n in
                  (normalise(sp, context) for sp in source_products)
                  if n is not None]

    fetched, kept = len(source_products), len(normalised)
    rejected = fetched - kept
    rate = round(rejected / fetched, 4) if fetched else 0.0

    persist(tenant_id, normalised, kind, ref, run_started)
    # Only reached when every page fetched and every write succeeded, so a
    # partial run can never be mistaken for "these products are gone".
    removed = delete_stale(tenant_id, kind, ref, run_started)

    scores = [quality_score(p) for p in normalised]
    if rate > REJECTION_ALERT_RATIO:
        logger.warning("Rejection rate %.1f%% for %s:%s suggests a normaliser "
                       "bug rather than bad merchant data", rate * 100, kind, ref)

    return {
        "fetched": fetched,
        "normalised": kept,
        "rejected": rejected,
        "rejection_rate": rate,
        "saved": kept,
        "deleted": removed,
        "quality_median": round(statistics.median(scores), 2) if scores else None,
    }
```

- [ ] **Step 4: Run test to verify it passes**

Run: `.venv/bin/python -m pytest tests/unit/test_catalog_sync.py -v`
Expected: PASS, 3 tests

- [ ] **Step 4b: Write the real-database test for the scoped delete**

This is the single most dangerous behaviour in the feature. It gets a test against real
Postgres, not a source-text grep: a refactor that keeps the words but widens the `WHERE`
clause would destroy a tenant's crawled catalog and pass a grep.

Create `tests/integration/test_catalog_persistence.py`:

```python
"""Real-database tests for scoped persistence.

The crawler, a Shopify store and an HTTP API all write to one table. An unscoped
delete here erases another producer's products, so it is proven against Postgres.
"""
from datetime import datetime, timedelta, timezone

import pytest

from app.services.catalog_sync import delete_stale, persist
from app.services.database import bootstrap_tenant, get_db_connection

NOW = datetime(2026, 8, 10, 12, 0, 0, tzinfo=timezone.utc)
EARLIER = NOW - timedelta(hours=1)


def product(key, name="A Product"):
    return {
        "product_key": key, "name": name, "description": None, "brand": None,
        "image_url": None, "product_url": f"https://x.example/{key}",
        "category": None, "taxonomy_path": [], "taxonomy_source": "none",
        "raw_category": None, "tags": [], "attributes": [],
        "price_cents": 1000, "price_max_cents": 1000, "compare_at_cents": None,
        "currency": "INR", "on_sale": False, "in_stock": True,
        "status": "ACTIVE", "external_id": key.split(":", 1)[1],
        "source_kind": key.split(":", 1)[0], "source_ref": "ref",
    }


def rows(schema, where=""):
    conn = get_db_connection()
    try:
        with conn.cursor() as cur:
            cur.execute(
                f'SELECT product_key FROM "{schema}".strategist_products {where} '
                "ORDER BY product_key")
            return [r[0] for r in cur.fetchall()]
    finally:
        conn.close()


@pytest.fixture
def populated(temp_tenant):
    bootstrap_tenant(temp_tenant)
    conn = get_db_connection()
    try:
        with conn.cursor() as cur:
            cur.execute(
                f'INSERT INTO "{temp_tenant}".strategist_products '
                "(product_key, name, product_url, source_kind, source_ref) VALUES "
                "('crawl:a', 'Crawled A', 'https://x.example/a', 'crawl', 'site'),"
                "('crawl:b', 'Crawled B', 'https://x.example/b', 'crawl', 'site')")
        conn.commit()
    finally:
        conn.close()
    return temp_tenant


def test_a_shopify_sync_does_not_delete_crawled_products(populated):
    # THE hazard. products.save_products' unscoped delete would remove both
    # crawled rows here.
    persist(populated, [product("shopify:1")], "shopify", "shop.myshopify.com", NOW)
    delete_stale(populated, "shopify", "shop.myshopify.com", NOW)

    assert rows(populated) == ["crawl:a", "crawl:b", "shopify:1"]


def test_stale_rows_of_the_same_source_are_deleted(populated):
    persist(populated, [product("shopify:1"), product("shopify:2")],
            "shopify", "shop.myshopify.com", EARLIER)
    persist(populated, [product("shopify:1")],
            "shopify", "shop.myshopify.com", NOW)

    removed = delete_stale(populated, "shopify", "shop.myshopify.com", NOW)

    assert removed == 1
    assert rows(populated) == ["crawl:a", "crawl:b", "shopify:1"]


def test_a_second_source_of_the_same_kind_is_untouched(populated):
    persist(populated, [product("shopify:1")], "shopify", "shop-a.myshopify.com", NOW)
    persist(populated, [product("shopify:2")], "shopify", "shop-b.myshopify.com", EARLIER)

    delete_stale(populated, "shopify", "shop-a.myshopify.com", NOW)

    assert "shopify:2" in rows(populated)


def test_upsert_does_not_duplicate(populated):
    persist(populated, [product("shopify:1", "First")], "shopify", "s", EARLIER)
    persist(populated, [product("shopify:1", "Second")], "shopify", "s", NOW)

    assert rows(populated, "WHERE product_key = 'shopify:1'") == ["shopify:1"]

    conn = get_db_connection()
    try:
        with conn.cursor() as cur:
            cur.execute(f'SELECT name FROM "{populated}".strategist_products '
                        "WHERE product_key = 'shopify:1'")
            assert cur.fetchone()[0] == "Second"
    finally:
        conn.close()


def test_persist_writes_the_commercial_fields(populated):
    persist(populated, [product("shopify:1")], "shopify", "s", NOW)
    conn = get_db_connection()
    try:
        with conn.cursor() as cur:
            cur.execute(
                f'SELECT price_cents, currency, in_stock, status, content_hash '
                f'FROM "{populated}".strategist_products '
                "WHERE product_key = 'shopify:1'")
            price, currency, in_stock, status, chash = cur.fetchone()
    finally:
        conn.close()

    assert price == 1000
    assert currency == "INR"
    assert in_stock is True
    assert status == "ACTIVE"
    assert chash
```

- [ ] **Step 4c: Run the real-database test**

Run: `.venv/bin/python -m pytest tests/integration/test_catalog_persistence.py -v`
Expected: PASS, 5 tests. If `test_a_shopify_sync_does_not_delete_crawled_products` fails,
stop — the delete is not scoped and would destroy live data.

- [ ] **Step 5: Write the failing test for source status tracking**

Spec §11 step 5 requires `last_synced_at` to be updated, and §13 requires a source whose
token was rejected to be marked `status='error'`. Neither exists yet.

Append to `tests/unit/test_catalog_sync.py`:

```python
def test_source_status_helpers_exist():
    from app.services import product_sources
    assert hasattr(product_sources, "mark_source_status")
    assert hasattr(product_sources, "touch_last_synced")


@pytest.mark.asyncio
async def test_rejected_token_marks_the_source_as_error(monkeypatch):
    from app.services.shopify_oauth import ShopQueryError

    marked = []

    async def unauthorised(*a, **kw):
        raise ShopQueryError("HTTP 401")

    monkeypatch.setattr(catalog_sync.catalog_shopify, "fetch_products", unauthorised)
    monkeypatch.setattr(catalog_sync, "get_credentials",
                        lambda *a: {"access_token": "shpat_dead"})
    monkeypatch.setattr(catalog_sync, "migrate_products_table", lambda t: None)
    monkeypatch.setattr(catalog_sync, "mark_source_status",
                        lambda *a: marked.append(a))

    with pytest.raises(ShopQueryError):
        await catalog_sync.sync_source("t", {
            "kind": "shopify", "external_ref": "s.myshopify.com",
            "config": {"shop_name": "S", "primary_domain": "https://s.com",
                       "currency": "INR"}})

    assert marked, "a rejected token must mark the source as error"
    assert marked[0][2] == "error"


@pytest.mark.asyncio
async def test_successful_sync_touches_last_synced_at(monkeypatch):
    touched = []

    async def none_at_all(*a, **kw):
        return []

    monkeypatch.setattr(catalog_sync.catalog_shopify, "fetch_products", none_at_all)
    monkeypatch.setattr(catalog_sync, "get_credentials",
                        lambda *a: {"access_token": "shpat_x"})
    monkeypatch.setattr(catalog_sync, "migrate_products_table", lambda t: None)
    monkeypatch.setattr(catalog_sync, "persist", lambda *a, **kw: {"saved": 0})
    monkeypatch.setattr(catalog_sync, "delete_stale", lambda *a, **kw: 0)
    monkeypatch.setattr(catalog_sync, "touch_last_synced",
                        lambda *a: touched.append(a))

    await catalog_sync.sync_source("t", {
        "kind": "shopify", "external_ref": "s.myshopify.com",
        "config": {"shop_name": "S", "primary_domain": "https://s.com",
                   "currency": "INR"}})

    assert touched, "a completed sync must record last_synced_at"
```

- [ ] **Step 6: Run test to verify it fails**

Run: `.venv/bin/python -m pytest tests/unit/test_catalog_sync.py -k "status or last_synced" -v`
Expected: FAIL with `AssertionError` on the `hasattr` check.

- [ ] **Step 7: Add the two helpers to product_sources.py**

Append to `app/services/product_sources.py`:

```python
def mark_source_status(kind: str, external_ref: str, status: str) -> None:
    """A source whose credentials stopped working must stop being retried
    silently; the dashboard needs to show that it needs reconnecting."""
    conn = get_master_db_connection()
    try:
        with conn.cursor() as cur:
            cur.execute(
                "UPDATE product_sources SET status = %s "
                "WHERE kind = %s AND external_ref = %s",
                (status, kind, external_ref),
            )
        conn.commit()
    except Exception:
        conn.rollback()
        logger.error("Could not set status for %s:%s", kind, external_ref,
                     exc_info=True)
        raise
    finally:
        conn.close()


def touch_last_synced(kind: str, external_ref: str, when) -> None:
    conn = get_master_db_connection()
    try:
        with conn.cursor() as cur:
            cur.execute(
                "UPDATE product_sources SET last_synced_at = %s "
                "WHERE kind = %s AND external_ref = %s",
                (when, kind, external_ref),
            )
        conn.commit()
    except Exception:
        conn.rollback()
        logger.error("Could not set last_synced_at for %s:%s", kind, external_ref,
                     exc_info=True)
        raise
    finally:
        conn.close()
```

- [ ] **Step 8: Wire them into sync_source**

In `app/services/catalog_sync.py`, extend the import:

```python
from app.services.product_sources import (
    get_credentials, mark_source_status, touch_last_synced,
)
```

Wrap the fetch so a rejected token marks the source, and record the timestamp on success.
Replace the fetch-and-map block in `sync_source` with:

```python
    try:
        if kind == "shopify":
            nodes = await catalog_shopify.fetch_products(
                ref, credentials["access_token"])
            source_products = [catalog_shopify.to_source_product(n, ref)
                               for n in nodes]
            context = {
                "shop_name": config.get("shop_name"),
                "primary_domain": config.get("primary_domain"),
                "url_template": None,
                "currency": config.get("currency_code") or config.get("currency"),
                "known_brands": [],
            }
        elif kind == "http_api":
            records = await catalog_http.fetch_products(config, credentials, ref)
            source_products = [catalog_http.to_source_product(r, config, ref)
                               for r in records]
            context = {
                "shop_name": config.get("shop_name"),
                "primary_domain": None,
                "url_template": config.get("url_template"),
                "currency": config.get("currency"),
                "known_brands": [],
            }
        else:
            raise ValueError(f"unsupported source kind: {kind}")
    except (ShopQueryError, HttpSourceError):
        # The credential no longer works. Surface it rather than retrying
        # silently every night and filling the logs with 401s.
        mark_source_status(kind, ref, "error")
        raise
```

Add the two exception imports at the top:

```python
from app.services.catalog_http import HttpSourceError
from app.services.shopify_oauth import ShopQueryError
```

Then, immediately before the `return` in `sync_source`:

```python
    touch_last_synced(kind, ref, run_started)
```

- [ ] **Step 9: Run test to verify it passes**

Run: `.venv/bin/python -m pytest tests/unit/test_catalog_sync.py -v`
Expected: PASS, 8 tests

- [ ] **Step 10: Report changed files**

```bash
git status --short
```

Do not commit and do not stage.

---

### Task 11: Sync API route

**Scope note:** an earlier draft of this task also filtered DRAFT / ARCHIVED / out-of-stock
products out of recommendation queries. That was removed at the user's direction —
recommendation behaviour is covered by separate work. **Do not modify
`app/services/products.py`, and do not change any query or ranking behaviour.** B's job
ends at storing `status` and `in_stock` accurately; whoever owns recommendations decides
how to use them.

**Files:**
- Modify: `app/api/endpoints.py`
- Test: `tests/integration/test_catalog_route.py`

**Interfaces:**
- Consumes: `catalog_sync.sync_source`, `product_sources.get_sources`.
- Produces: `POST /catalog/sync`.

- [ ] **Step 1: Write the failing test**

Create `tests/integration/test_catalog_route.py`:

```python
"""The sync route. Storage correctness is covered elsewhere; this checks the
route's contract and that a caller cannot trigger a sync for a source that is
not connected.
"""
import pytest
from fastapi.testclient import TestClient

from app.main import app

client = TestClient(app)
TENANT = "org_8c32bf3e-6a18-4739-9b1c-94c0cf11125f"


def test_unknown_tenant_returns_404():
    resp = client.post("/catalog/sync", data={"tenant_id": "org_does_not_exist"})
    assert resp.status_code == 404


def test_unknown_source_ref_returns_404():
    resp = client.post("/catalog/sync", data={
        "tenant_id": TENANT, "source_ref": "not-a-real-shop.myshopify.com"})
    assert resp.status_code == 404


def test_tenant_id_is_required():
    assert client.post("/catalog/sync", data={}).status_code == 422


def test_route_is_registered():
    paths = {r.path for r in app.routes}
    assert "/catalog/sync" in paths


def test_a_failing_source_reports_an_error_without_failing_the_request(monkeypatch):
    # One broken source must not hide the results of the others.
    from app.api import endpoints

    async def boom(tenant_id, source):
        raise RuntimeError("token rejected")

    monkeypatch.setattr(endpoints, "sync_source", boom)
    monkeypatch.setattr(endpoints, "get_sources", lambda t: [
        {"kind": "shopify", "external_ref": "s.myshopify.com", "config": {}}])

    resp = client.post("/catalog/sync", data={"tenant_id": TENANT})
    assert resp.status_code == 200
    assert "error" in resp.json()["results"][0]
```

- [ ] **Step 2: Run test to verify it fails**

Run: `.venv/bin/python -m pytest tests/integration/test_catalog_route.py -v`
Expected: FAIL — the route does not exist, so `/catalog/sync` 404s on every call
including the ones expecting 422.

- [ ] **Step 3: Add the sync route**

In `app/api/endpoints.py`, add to the imports:

```python
from app.services.catalog_sync import sync_source
from app.services.product_sources import get_sources
```

And add the route:

```python
@router.post("/catalog/sync")
async def catalog_sync_endpoint(
    tenant_id: str = Form(...),
    source_ref: Optional[str] = Form(None),
):
    sources = get_sources(tenant_id)
    if source_ref:
        sources = [s for s in sources if s["external_ref"] == source_ref]
    if not sources:
        raise HTTPException(status_code=404, detail="No connected source found")

    results = []
    for source in sources:
        try:
            results.append({"source": source["external_ref"],
                            **await sync_source(tenant_id, source)})
        except Exception as ex:
            logger.error("Catalog sync failed for %s", source["external_ref"],
                         exc_info=True)
            results.append({"source": source["external_ref"], "error": str(ex)})

    return {"results": results}
```

- [ ] **Step 4: Run test to verify it passes**

Run: `.venv/bin/python -m pytest tests/integration/test_catalog_route.py -v`
Expected: PASS, 5 tests

- [ ] **Step 5: Run the whole suite**

Run: `.venv/bin/python -m pytest tests/unit/ -q`
Expected: PASS. Baseline was 224; this plan adds roughly 150 tests.

- [ ] **Step 7: Report changed files**

```bash
git status --short
```

Do not commit and do not stage.

---

### Task 12: Live sync against the real store

**Files:** none — verification only.

**Prerequisites:** uvicorn running from this worktree on port 8001, Postgres and Redis reachable, and the Shopify source connected (it is — A stored it).

- [ ] **Step 1: Start the server from this worktree**

```bash
.venv/bin/uvicorn app.main:app --port 8001 &
sleep 5
curl -s -o /dev/null -w "%{http_code}\n" http://localhost:8001/health
```

Expected: `200`.

- [ ] **Step 2: Run the sync**

```bash
curl -s -X POST http://localhost:8001/catalog/sync \
  -F "tenant_id=org_8c32bf3e-6a18-4739-9b1c-94c0cf11125f" | python3 -m json.tool
```

Expected: `fetched: 17`, `normalised: 17`, `rejected: 0`, `rejection_rate: 0.0`,
`saved: 17`, and a `quality_median` value. A non-zero `rejected` means a normalisation
rule is wrong — read the logs before changing anything.

- [ ] **Step 3: Confirm the rows and that crawled products survived**

```bash
PYTHONPATH=. .venv/bin/python -c "
from app.services.database import get_db_connection
T='org_8c32bf3e-6a18-4739-9b1c-94c0cf11125f'
c=get_db_connection(); cur=c.cursor()
cur.execute('SELECT source_kind, count(*) FROM \"'+T+'\".strategist_products GROUP BY 1')
print('by source:', cur.fetchall())
cur.execute('SELECT count(*) FROM \"'+T+'\".strategist_products WHERE brand IS NULL AND source_kind=%s', ('shopify',))
print('shopify rows with brand nulled:', cur.fetchone()[0])
cur.execute('SELECT name, price_cents, price_max_cents, currency, in_stock, status FROM \"'+T+'\".strategist_products WHERE source_kind=%s ORDER BY price_cents LIMIT 5', ('shopify',))
for r in cur.fetchall(): print(' ', r)
c.close()
"
```

Expected: both `crawl` and `shopify` present in the source breakdown — the §4 hazard
would show up here as crawled products having vanished. Eight Shopify rows with `brand`
nulled. Prices as integers, currency `INR`.

- [ ] **Step 4: Confirm status and stock landed accurately**

B stores these; it does not act on them. Whoever owns recommendations decides how.

```bash
PYTHONPATH=. .venv/bin/python -c "
from app.services.database import get_db_connection
T='org_8c32bf3e-6a18-4739-9b1c-94c0cf11125f'
c=get_db_connection(); cur=c.cursor()
cur.execute('SELECT status, count(*) FROM \"'+T+'\".strategist_products WHERE source_kind=%s GROUP BY 1', ('shopify',))
print('status:', cur.fetchall())
cur.execute('SELECT in_stock, count(*) FROM \"'+T+'\".strategist_products WHERE source_kind=%s GROUP BY 1', ('shopify',))
print('in_stock:', cur.fetchall())
c.close()"
```

Expected: 15 ACTIVE, 1 DRAFT, 1 ARCHIVED; 15 in stock, 2 out. These match the live store
profile exactly, so any other numbers mean a mapping bug.

- [ ] **Step 5: Confirm a second sync is idempotent**

Re-run Step 2, then:

```bash
PYTHONPATH=. .venv/bin/python -c "
from app.services.database import get_db_connection
T='org_8c32bf3e-6a18-4739-9b1c-94c0cf11125f'
c=get_db_connection(); cur=c.cursor()
cur.execute('SELECT count(*) FROM \"'+T+'\".strategist_products WHERE source_kind=%s', ('shopify',))
print('shopify rows after second sync:', cur.fetchone()[0])
c.close()"
```

Expected: still 17. A larger number means the upsert key is wrong.

- [ ] **Step 6: Confirm the rows are queryable**

```bash
PYTHONPATH=. .venv/bin/python -c "
from app.services.products import list_products
rows = list_products('org_8c32bf3e-6a18-4739-9b1c-94c0cf11125f', limit=5)
for r in rows: print(' ', r['name'], '|', r['product_url'])
"
```

Expected: real product titles from the store. Do not assert anything about which products
appear or in what order — recommendation behaviour is out of scope for this plan.

- [ ] **Step 7: Report**

Summarise: changed files, test counts, the sync report numbers, the source breakdown, and
confirmation that crawled products survived. Remind the user to commit, and that the
synced product URLs will 404 until the Online Store channel is published on the dev store.

---

## Follow-on work (not in this plan)

- **C:** LLM field-map authoring, enrichment for low-scoring products, canonical taxonomy
  tree, quality gating and tenant health reporting.
- **Before production:** close `TODO(auth)` from A, register webhooks so prices stop being
  sync-cycle stale, and revisit whether cards should then show prices.
