# Integration Conformance and Ingestion Completeness Implementation Plan


**Goal:** Make the Shopify connection conform to the platform's existing three-table integration pattern, capture the product fields the recommendation matching specification needs, and make the HTTP product-API source connectable.

**Architecture:** A new `platform_integrations.py` service owns every read and write against the platform's Sequelize-managed tables (`platform_configs`, `integrations`, `integration_connect_tracker`), so cross-service schema coupling lives in exactly one file. The existing Shopify callback calls into it. Ingestion changes are additive columns plus five optional field-map keys, with currency normalisation so prices from two sources are comparable at all. A new `sources.py` router adds the HTTP connect flow, guarded by an SSRF check that resolves hostnames before any socket opens.

**A recurring rule across this plan:** a product that cannot be made usable is **stored, flagged in `missing_fields`, and excluded from serving** — never silently rejected and never silently wrong. It applies to a missing URL and to a price that cannot be converted. Both failures are then visible in the sync report instead of surfacing months later as "the recommendations are bad".

**Tech Stack:** FastAPI, psycopg2, httpx, Azure OpenAI (via the existing client), pytest.

**Spec:** `docs/specs/2026-08-10-ingestion-completeness-design.md`

**Depends on:** sub-projects A and B, both built and live-verified in this worktree.

## Global Constraints

- **Never commit, never push, never stage.** The user commits their own work. Every task ends with `git status --short` and a report. 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 is not evident from the code. A reviewer who knows Python must find zero noise comments.
- **`integrations`, `platform_configs` and `integration_connect_tracker` are owned by the Node application** and managed by Sequelize migrations. This plan **reads** all three and **writes** `integrations` and `integration_connect_tracker` rows. It must never ALTER, CREATE or DROP any of them, and never assume a column this plan did not verify exists.
- **Tenant id mapping:** the strategist uses `org_<uuid>`; the platform uses `<uuid>`. Verified identical apart from the prefix, and the uuid matches `Tenants.id`. Implement the strip once, in one helper.
- **`integrations.accessToken` receives Fernet ciphertext, never plaintext.** Matching the other providers' apparent plaintext storage is not worth reducing protection on a live credential. Set `metadata.token_encoding = "fernet"` so the next reader knows.
- **`integrations.user` is NOT NULL with a foreign key to `Users(id)`.** Derive it from `Tenants.userId`. Two of 28 tenants have a null `userId` — those must produce a clear 409, never a constraint crash or a silent skip.
- **Money is integer minor units. `in_stock` never derives from inventory quantity. Empty string is never stored — normalise to `None`.** (Carried from B.)
- **Never log or return a credential, token, API key, authorization code, or full response body.**
- Match conventions: `psycopg2` with `sql.SQL`/`sql.Identifier` for interpolated identifiers, module-level `logger = logging.getLogger(__name__)`, snake_case, double quotes, 4-space indent, type hints on public signatures, imports stdlib → third-party → `app.*`.
- Tests: `.venv/bin/python -m pytest`. Baselines: `tests/unit/` **418 passing**, `tests/integration/` **29 passing**. `.venv` is a symlink to the main checkout's virtualenv. Use `PYTHONPATH=.` only when running a script directly.
- Do **not** modify `app/services/products.py` beyond what B already changed. Recommendation behaviour is owned by separate work.

---

## File Structure

| File | Responsibility |
|---|---|
| `app/services/platform_integrations.py` (new) | Every read/write against `platform_configs`, `integrations`, `integration_connect_tracker`. The only file that knows those schemas. |
| `app/services/ssrf.py` (new) | Resolve-and-reject guard for user-supplied URLs. Pure, no I/O beyond DNS. |
| `app/services/fieldmap_proposer.py` (new) | Sample-structure summarisation, the single LLM call, validation of inferred paths, and url_template inference. Called by the sync pipeline, never by a route. |
| `app/api/sources.py` (new) | `APIRouter(prefix="/sources")` — create, list, delete. |
| `app/services/catalog_shopify.py` (modify) | Complementary-products metafield → `tenant_relations`. |
| `app/services/catalog_http.py` (modify) | Five new optional field-map keys; SSRF guard on fetch. |
| `app/services/catalog_normalise.py` (modify) | Carry the new fields; drop description from `content_hash`. |
| `app/services/database.py` (modify) | Four new columns in `_PRODUCT_COLUMNS` and `bootstrap_tenant`. |
| `app/api/shopify.py` (modify) | Callback writes `integrations` + tracker rows. |
| `app/core/config.py` (modify) | Client credentials read from `platform_configs`, `.env` as fallback. |
| `app/main.py` (modify) | Mount the sources router. |

---

### Task 1: Schema additions for the matching specification

**Files:**
- Modify: `app/services/database.py` — `_PRODUCT_COLUMNS` and `bootstrap_tenant`'s CREATE TABLE
- Modify: `tests/integration/test_products_table.py` — `EXPECTED_COLUMNS`
- Modify: `tests/integration/test_products_migration.py` — `NEW_COLUMNS`

**Interfaces:**
- Consumes: `migrate_products_table(tenant_id)` from B.
- Produces: four columns — `tenant_relations JSONB`, `rating REAL`, `review_count INT`, `featured_rank INT` — and `product_url` becomes nullable.

**Why `product_url` must become nullable:** a source that exposes no product URLs at all
(the reference HTTP API is one) would otherwise have every product rejected at validation,
and the merchant would connect and see an empty catalog with no explanation. Storing the
product with a null URL keeps it searchable, countable and usable as an anchor; only its
use as a clickable suggestion is affected. See the spec's §10.

**This task breaks two existing tests unless you update them.** `test_products_table.py` asserts an exact column set; `test_products_migration.py` asserts `NEW_COLUMNS <= columns`. Both need the four names added.

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

Add to `tests/integration/test_products_migration.py`, and add the same four names to that file's existing `NEW_COLUMNS` set:

```python
MATCHING_COLUMNS = {"tenant_relations", "rating", "review_count", "featured_rank",
                    "missing_fields", "price_reference_cents", "fx_rate_used"}


def test_matching_spec_columns_are_added(old_table):
    assert not (MATCHING_COLUMNS & _columns(old_table))
    migrate_products_table(old_table)
    assert MATCHING_COLUMNS <= _columns(old_table)


def test_product_url_becomes_nullable(old_table):
    # A source with no product URLs would otherwise have every row rejected, and
    # the merchant would see an empty catalog with no explanation.
    migrate_products_table(old_table)
    _exec(f"""INSERT INTO "{old_table}".strategist_products
              (product_key, name, product_url) VALUES ('h:1','No Link',NULL)""")
    rows = _query(
        f"SELECT product_url FROM \"{old_table}\".strategist_products "
        "WHERE product_key = 'h:1'")
    assert rows[0][0] is None


def test_tenant_relations_defaults_to_empty_object(old_table):
    # Never null: downstream code reads .get("related") without a guard, and a
    # null would force every caller to branch.
    migrate_products_table(old_table)
    _exec(f"""INSERT INTO "{old_table}".strategist_products
              (product_key, name, product_url) VALUES ('c:x','X','https://x/1')""")
    rows = _query(
        f"SELECT tenant_relations FROM \"{old_table}\".strategist_products "
        "WHERE product_key = 'c:x'")
    assert rows[0][0] == {}
```

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

Run: `.venv/bin/python -m pytest tests/integration/test_products_migration.py -v`
Expected: FAIL — the four columns do not exist.

- [ ] **Step 3: Add the columns**

In `app/services/database.py`, append to `_PRODUCT_COLUMNS`:

```python
    ("tenant_relations", "JSONB NOT NULL DEFAULT '{}'"),
    ("rating", "REAL"),
    ("review_count", "INT"),
    ("featured_rank", "INT"),
    ("missing_fields", "TEXT[] NOT NULL DEFAULT '{}'"),
    ("price_reference_cents", "INT"),
    ("fx_rate_used", "REAL"),
```

Add the same four to `bootstrap_tenant`'s `strategist_products` CREATE TABLE, matching the
style of the columns B added, and change its `product_url` from `TEXT NOT NULL` to `TEXT`.
Note the existing brace-escaping convention in that CREATE TABLE: literal `{}` is written
`{{}}` because the statement goes through `.format()`.

Then, in `migrate_products_table`, after the ADD COLUMN loop and before the index
statements, relax the constraint on existing tenants:

```python
            # A source that exposes no product URLs would otherwise have every
            # row rejected at validation, leaving the merchant with an empty
            # catalog and no explanation.
            cur.execute(sql.SQL(
                "ALTER TABLE {}.strategist_products "
                "ALTER COLUMN product_url DROP NOT NULL"
            ).format(sql.Identifier(tenant_id)))
```

`DROP NOT NULL` is idempotent in Postgres — running it on an already-nullable column
succeeds and changes nothing — so the migration stays safe to re-run.

- [ ] **Step 4: Update the two existing tests**

Add `"tenant_relations"`, `"rating"`, `"review_count"`, `"featured_rank"` to
`EXPECTED_COLUMNS` in `tests/integration/test_products_table.py`, with a comment recording
that they exist for the recommendation matching specification. Do not weaken the `==`
assertion to a subset check — its exactness is what catches an accidental column.

- [ ] **Step 5: Run the migration against the live tenant, then the suite**

```bash
PYTHONPATH=. .venv/bin/python -c "
from app.services.database import migrate_products_table
print(migrate_products_table('org_8c32bf3e-6a18-4739-9b1c-94c0cf11125f'))"
.venv/bin/python -m pytest tests/integration/ -q
```

Expected: the migration reports the four new columns, then all integration tests pass.
Run the migration a second time and confirm it reports zero added — it must stay idempotent.

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

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

Do not commit and do not stage.

---

### Task 2: SSRF guard

The highest-risk code in this plan. It runs before any user-supplied URL is fetched.

**Files:**
- Create: `app/services/ssrf.py`
- Test: `tests/unit/test_ssrf.py`

**Interfaces:**
- Consumes: nothing.
- Produces:
  - `class UnsafeUrlError(Exception)`
  - `assert_safe_url(url: str) -> str` — returns the URL unchanged, or raises `UnsafeUrlError` naming the rule that failed.

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

Create `tests/unit/test_ssrf.py`:

```python
import pytest

from app.services.ssrf import UnsafeUrlError, assert_safe_url


@pytest.mark.parametrize("url", [
    "https://dummyjson.com/products",
    "http://example.com:8080/api",
    "https://api.example.co.uk/v1/products",
])
def test_public_urls_allowed(url):
    assert assert_safe_url(url) == url


@pytest.mark.parametrize("url", [
    "http://127.0.0.1/admin",
    "http://127.0.0.1:8001/api/shopify/status",
    "http://localhost/",
    "http://[::1]/",
    "http://169.254.169.254/latest/meta-data/",   # cloud instance metadata
    "http://10.0.0.1/",
    "http://192.168.1.1/",
    "http://172.16.0.1/",
    "http://0.0.0.0/",
])
def test_internal_addresses_rejected(url):
    with pytest.raises(UnsafeUrlError):
        assert_safe_url(url)


@pytest.mark.parametrize("url", [
    "file:///etc/passwd",
    "gopher://example.com/",
    "ftp://example.com/",
    "not-a-url",
    "",
    None,
])
def test_non_http_schemes_rejected(url):
    with pytest.raises(UnsafeUrlError):
        assert_safe_url(url)


def test_hostname_resolving_to_loopback_is_rejected(monkeypatch):
    # The load-bearing case. Validating the hostname STRING would let this
    # through, because "sneaky.example.com" looks entirely public.
    import app.services.ssrf as mod
    monkeypatch.setattr(mod, "_resolve", lambda host: ["127.0.0.1"])
    with pytest.raises(UnsafeUrlError):
        assert_safe_url("https://sneaky.example.com/products")


def test_unresolvable_host_rejected(monkeypatch):
    import app.services.ssrf as mod

    def boom(host):
        raise OSError("nodename nor servname provided")

    monkeypatch.setattr(mod, "_resolve", boom)
    with pytest.raises(UnsafeUrlError):
        assert_safe_url("https://does-not-exist.example/")


def test_error_names_the_rule():
    with pytest.raises(UnsafeUrlError) as ex:
        assert_safe_url("file:///etc/passwd")
    assert "scheme" in str(ex.value).lower()
```

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

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

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

Create `app/services/ssrf.py`:

```python
"""Rejects user-supplied URLs that point inside the network.

Any endpoint that fetches a URL a user chose is an SSRF sink: pointed at cloud
instance metadata it leaks credentials, pointed at an RFC1918 address it reaches
services that were never meant to be public.
"""
import ipaddress
import logging
import socket
from urllib.parse import urlparse

logger = logging.getLogger(__name__)

ALLOWED_SCHEMES = ("http", "https")


class UnsafeUrlError(Exception):
    """The URL points somewhere this service must not fetch from."""


def _resolve(host: str) -> list:
    return [info[4][0] for info in socket.getaddrinfo(host, None)]


def assert_safe_url(url) -> str:
    if not url or not isinstance(url, str):
        raise UnsafeUrlError("url must be a non-empty string")

    parsed = urlparse(url)
    if parsed.scheme not in ALLOWED_SCHEMES:
        raise UnsafeUrlError(f"scheme {parsed.scheme!r} is not allowed")
    if not parsed.hostname:
        raise UnsafeUrlError("url has no host")

    try:
        addresses = _resolve(parsed.hostname)
    except OSError as ex:
        raise UnsafeUrlError(f"host {parsed.hostname!r} does not resolve") from ex

    for raw in addresses:
        # Check the RESOLVED address, never the hostname string: a public-looking
        # name can resolve to 127.0.0.1 and would otherwise pass every check.
        ip = ipaddress.ip_address(raw)
        if (ip.is_private or ip.is_loopback or ip.is_link_local
                or ip.is_reserved or ip.is_multicast or ip.is_unspecified):
            raise UnsafeUrlError(
                f"host {parsed.hostname!r} resolves to non-public address {raw}")

    return url
```

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

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

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

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

Do not commit and do not stage.

---

### Task 3: Platform integration service

The only file that touches the Node application's tables.

**Files:**
- Create: `app/services/platform_integrations.py`
- Test: `tests/unit/test_platform_integrations.py`
- Test: `tests/integration/test_platform_integrations_db.py`

**Interfaces:**
- Consumes: `get_master_db_connection` from `app.services.database`.
- Produces:
  - `strip_tenant_prefix(tenant_id: str) -> str` — `org_<uuid>` → `<uuid>`; raises `ValueError` otherwise
  - `get_platform_config(platform: str) -> dict | None` — `{client_id, client_secret, redirect_uri, scope}`
  - `resolve_tenant_user(tenant_uuid: str) -> str | None`
  - `class TenantHasNoUserError(Exception)`
  - `upsert_integration(tenant_id, provider, auth_type, access_token_ciphertext, metadata: dict) -> str` returns the integration id; raises `TenantHasNoUserError`
  - `set_integration_status(tenant_id, provider, status) -> None`
  - `get_integration(tenant_id, provider) -> dict | None` — never returns the token
  - `track_connect(tenant_id, platform, delivery_status, outcome, error_message=None, extra: dict | None = None) -> None`

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

Create `tests/unit/test_platform_integrations.py`:

```python
import pytest

from app.services.platform_integrations import strip_tenant_prefix


def test_strips_the_org_prefix():
    assert strip_tenant_prefix("org_8c32bf3e-6a18-4739-9b1c-94c0cf11125f") == \
        "8c32bf3e-6a18-4739-9b1c-94c0cf11125f"


@pytest.mark.parametrize("bad", [
    "8c32bf3e-6a18-4739-9b1c-94c0cf11125f",   # no prefix
    "org_not-a-uuid",
    "org_",
    "",
    None,
    "tenant_8c32bf3e-6a18-4739-9b1c-94c0cf11125f",
])
def test_rejects_anything_that_is_not_a_prefixed_uuid(bad):
    # The value goes into a uuid column with a foreign key; a malformed one must
    # fail here rather than as a database error halfway through a callback.
    with pytest.raises(ValueError):
        strip_tenant_prefix(bad)
```

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

Run: `.venv/bin/python -m pytest tests/unit/test_platform_integrations.py -v`
Expected: FAIL with `ModuleNotFoundError`

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

Create `app/services/platform_integrations.py`:

```python
"""Reads and writes the platform's integration tables.

Those tables belong to the Node application and are managed by its Sequelize
migrations. Every access is funnelled through this module so the cross-service
coupling sits in one reviewable place, and so a schema change over there breaks
one file rather than several.
"""
import json
import logging
import uuid as uuid_lib

from psycopg2.extras import Json, RealDictCursor

from app.services.database import get_master_db_connection

logger = logging.getLogger(__name__)

TENANT_PREFIX = "org_"
JOB_NAME = "integration_oauth_connect"


class TenantHasNoUserError(Exception):
    """integrations.user is NOT NULL, and this tenant has no owning user."""


def strip_tenant_prefix(tenant_id) -> str:
    if not tenant_id or not isinstance(tenant_id, str):
        raise ValueError("tenant_id must be a non-empty string")
    if not tenant_id.startswith(TENANT_PREFIX):
        raise ValueError(f"tenant_id {tenant_id!r} is missing the org_ prefix")
    raw = tenant_id[len(TENANT_PREFIX):]
    # Parsed rather than pattern-matched: the value lands in a uuid column with a
    # foreign key, so a malformed one must fail here, not mid-callback.
    return str(uuid_lib.UUID(raw))


def get_platform_config(platform: str):
    conn = get_master_db_connection()
    try:
        with conn.cursor(cursor_factory=RealDictCursor) as cur:
            cur.execute(
                'SELECT "clientId", "clientSecret", "redirectUri", scope '
                'FROM platform_configs WHERE platform = %s AND "isActive" = true',
                (platform,))
            row = cur.fetchone()
            if not row:
                return None
            return {"client_id": row["clientId"],
                    "client_secret": row["clientSecret"],
                    "redirect_uri": row["redirectUri"],
                    "scope": row["scope"]}
    finally:
        conn.close()


def resolve_tenant_user(tenant_uuid: str):
    conn = get_master_db_connection()
    try:
        with conn.cursor() as cur:
            cur.execute('SELECT "userId" FROM "Tenants" WHERE id = %s', (tenant_uuid,))
            row = cur.fetchone()
            return row[0] if row and row[0] else None
    finally:
        conn.close()


def upsert_integration(tenant_id: str, provider: str, auth_type: str,
                       access_token_ciphertext: str, metadata: dict) -> str:
    tenant_uuid = strip_tenant_prefix(tenant_id)
    user_id = resolve_tenant_user(tenant_uuid)
    if not user_id:
        raise TenantHasNoUserError(
            f"Tenant {tenant_uuid} has no userId; integrations.user is NOT NULL")

    payload = dict(metadata or {})
    # Records that accessToken holds ciphertext, not a raw token, so a future
    # reader does not send it to Shopify verbatim.
    payload["token_encoding"] = "fernet"

    conn = get_master_db_connection()
    try:
        with conn.cursor() as cur:
            cur.execute(
                'SELECT id FROM integrations WHERE "tenantId" = %s AND provider = %s',
                (tenant_uuid, provider))
            existing = cur.fetchone()

            if existing:
                cur.execute(
                    'UPDATE integrations SET "accessToken" = %s, "authType" = %s, '
                    'metadata = %s, status = %s, "updatedAt" = now() WHERE id = %s',
                    (access_token_ciphertext, auth_type, Json(payload),
                     "ACTIVE", existing[0]))
                integration_id = str(existing[0])
            else:
                integration_id = str(uuid_lib.uuid4())
                cur.execute(
                    'INSERT INTO integrations (id, "user", provider, "authType", '
                    '"accessToken", metadata, status, "tenantId", "createdAt", '
                    '"updatedAt") VALUES (%s,%s,%s,%s,%s,%s,%s,%s, now(), now())',
                    (integration_id, user_id, provider, auth_type,
                     access_token_ciphertext, Json(payload), "ACTIVE", tenant_uuid))
        conn.commit()
        logger.info("Recorded %s integration for tenant %s", provider, tenant_uuid)
        return integration_id
    except Exception:
        conn.rollback()
        logger.error("Could not record %s integration for tenant %s",
                     provider, tenant_uuid, exc_info=True)
        raise
    finally:
        conn.close()


def set_integration_status(tenant_id: str, provider: str, status: str) -> None:
    tenant_uuid = strip_tenant_prefix(tenant_id)
    conn = get_master_db_connection()
    try:
        with conn.cursor() as cur:
            cur.execute(
                'UPDATE integrations SET status = %s, "updatedAt" = now() '
                'WHERE "tenantId" = %s AND provider = %s',
                (status, tenant_uuid, provider))
        conn.commit()
    except Exception:
        conn.rollback()
        logger.error("Could not set %s status for tenant %s", provider, tenant_uuid,
                     exc_info=True)
        raise
    finally:
        conn.close()


def get_integration(tenant_id: str, provider: str):
    """Never returns accessToken -- callers that need the token read it from
    product_sources, where it is stored with the rest of the sync state."""
    tenant_uuid = strip_tenant_prefix(tenant_id)
    conn = get_master_db_connection()
    try:
        with conn.cursor(cursor_factory=RealDictCursor) as cur:
            cur.execute(
                'SELECT id, provider, "authType", status, metadata, "createdAt", '
                '"updatedAt" FROM integrations WHERE "tenantId" = %s AND provider = %s',
                (tenant_uuid, provider))
            row = cur.fetchone()
            return dict(row) if row else None
    finally:
        conn.close()


def track_connect(tenant_id: str, platform: str, delivery_status: str,
                  outcome: str, error_message=None, extra=None) -> None:
    """Best-effort audit. A tracker failure must never fail a connection that
    otherwise succeeded, so this logs and returns rather than raising."""
    try:
        tenant_uuid = strip_tenant_prefix(tenant_id)
    except ValueError:
        logger.warning("Cannot track connect for malformed tenant %r", tenant_id)
        return

    metadata = {"jobName": JOB_NAME, "providerRaw": platform,
                "connectStep": delivery_status}
    metadata.update(extra or {})

    conn = get_master_db_connection()
    try:
        with conn.cursor() as cur:
            cur.execute(
                'INSERT INTO integration_connect_tracker ("tenantId", platform, '
                'outcome, "deliveryStatus", "errorMessage", metadata, "createdAt", '
                '"updatedAt") VALUES (%s,%s,%s,%s,%s,%s, now(), now())',
                (tenant_uuid, platform, outcome, delivery_status,
                 error_message, Json(metadata)))
        conn.commit()
    except Exception:
        conn.rollback()
        logger.error("Could not write connect tracker row for %s", platform,
                     exc_info=True)
    finally:
        conn.close()
```

- [ ] **Step 4: Run the unit test**

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

- [ ] **Step 5: Write the database test**

Create `tests/integration/test_platform_integrations_db.py`:

```python
"""Read-only checks against the platform's real integration tables.

These tables belong to the Node application. The tests below read them and write
only rows scoped to the known test tenant, which is cleaned up afterwards.
"""
import pytest

from app.services.database import get_master_db_connection
from app.services.platform_integrations import (
    TenantHasNoUserError, get_platform_config, resolve_tenant_user,
    strip_tenant_prefix, upsert_integration,
)

TENANT = "org_8c32bf3e-6a18-4739-9b1c-94c0cf11125f"
TENANT_UUID = "8c32bf3e-6a18-4739-9b1c-94c0cf11125f"


@pytest.fixture
def clean_integration():
    yield
    conn = get_master_db_connection()
    try:
        with conn.cursor() as cur:
            cur.execute('DELETE FROM integrations WHERE "tenantId" = %s '
                        "AND provider = 'shopify'", (TENANT_UUID,))
        conn.commit()
    finally:
        conn.close()


def test_the_test_tenant_has_an_owning_user():
    assert resolve_tenant_user(TENANT_UUID) is not None


def test_a_tenant_without_a_user_is_reported_not_crashed(clean_integration):
    # Two of 28 tenants have a null userId. integrations.user is NOT NULL, so
    # this must surface as a clear error rather than a constraint violation.
    conn = get_master_db_connection()
    try:
        with conn.cursor() as cur:
            cur.execute('SELECT id FROM "Tenants" WHERE "userId" IS NULL LIMIT 1')
            row = cur.fetchone()
    finally:
        conn.close()
    if not row:
        pytest.skip("no userless tenant in this database")

    with pytest.raises(TenantHasNoUserError):
        upsert_integration(f"org_{row[0]}", "shopify", "oauth2", "gAAAA", {})


def test_upsert_creates_then_updates_one_row(clean_integration):
    first = upsert_integration(TENANT, "shopify", "oauth2", "gAAAAciphertext1",
                               {"shop_domain": "a.myshopify.com"})
    second = upsert_integration(TENANT, "shopify", "oauth2", "gAAAAciphertext2",
                                {"shop_domain": "a.myshopify.com"})
    assert first == second

    conn = get_master_db_connection()
    try:
        with conn.cursor() as cur:
            cur.execute('SELECT count(*), max("accessToken"), max(metadata->>%s) '
                        'FROM integrations WHERE "tenantId" = %s AND provider = %s',
                        ("token_encoding", TENANT_UUID, "shopify"))
            count, token, encoding = cur.fetchone()
    finally:
        conn.close()

    assert count == 1
    assert token == "gAAAAciphertext2"
    assert encoding == "fernet"


def test_shopify_platform_config_is_absent_or_complete():
    # Task 4 adds this row. Before that it is absent, and the code must cope.
    cfg = get_platform_config("shopify")
    if cfg is not None:
        assert cfg["client_id"]
        assert cfg["client_secret"]


def test_strip_prefix_matches_a_real_tenant_row():
    assert strip_tenant_prefix(TENANT) == TENANT_UUID
```

- [ ] **Step 6: Run the database test**

Run: `.venv/bin/python -m pytest tests/integration/test_platform_integrations_db.py -v`
Expected: PASS, 5 tests (one may skip)

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

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

Do not commit and do not stage.

---

### Task 4: Read client credentials from `platform_configs`

**Files:**
- Modify: `app/services/shopify_oauth.py`
- Test: `tests/unit/test_shopify_oauth.py`

**Interfaces:**
- Consumes: `platform_integrations.get_platform_config`.
- Produces: `shopify_credentials() -> tuple[str, str]` returning `(client_id, client_secret)`.

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

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

```python
from app.services.shopify_oauth import shopify_credentials


def test_credentials_prefer_platform_configs(monkeypatch):
    import app.services.shopify_oauth as mod
    monkeypatch.setattr(mod, "get_platform_config", lambda p: {
        "client_id": "from-platform", "client_secret": "secret-from-platform",
        "redirect_uri": None, "scope": None})
    assert shopify_credentials() == ("from-platform", "secret-from-platform")


def test_credentials_fall_back_to_env_when_no_row(monkeypatch):
    # A missing platform_configs row must not break a working connector.
    import app.services.shopify_oauth as mod
    monkeypatch.setattr(mod, "get_platform_config", lambda p: None)
    client_id, secret = shopify_credentials()
    assert client_id == settings.SHOPIFY_CLIENT_ID
    assert secret == settings.SHOPIFY_CLIENT_SECRET


def test_credentials_fall_back_when_the_lookup_raises(monkeypatch):
    # The platform database being unreachable must not take the connector down.
    import app.services.shopify_oauth as mod

    def boom(platform):
        raise RuntimeError("master db unreachable")

    monkeypatch.setattr(mod, "get_platform_config", boom)
    assert shopify_credentials() == (settings.SHOPIFY_CLIENT_ID,
                                     settings.SHOPIFY_CLIENT_SECRET)
```

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

Run: `.venv/bin/python -m pytest tests/unit/test_shopify_oauth.py -k credentials -v`
Expected: FAIL with `ImportError: cannot import name 'shopify_credentials'`

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

Add to the top import block of `app/services/shopify_oauth.py`:

```python
from app.services.platform_integrations import get_platform_config
```

And add:

```python
PLATFORM = "shopify"


def shopify_credentials() -> tuple:
    """Client credentials from platform_configs, falling back to the environment.

    platform_configs is where the other thirteen providers keep theirs, and
    rotating a secret in one place every service reads beats rotating a copy per
    service. The fallback exists so a missing row cannot break a working
    connector.
    """
    try:
        cfg = get_platform_config(PLATFORM)
    except Exception:
        logger.warning("platform_configs lookup failed; using environment "
                       "credentials", exc_info=True)
        cfg = None

    if cfg and cfg.get("client_id") and cfg.get("client_secret"):
        return cfg["client_id"], cfg["client_secret"]
    return settings.SHOPIFY_CLIENT_ID, settings.SHOPIFY_CLIENT_SECRET
```

Then replace every direct use of `settings.SHOPIFY_CLIENT_ID` and
`settings.SHOPIFY_CLIENT_SECRET` inside `build_authorize_url`, `verify_hmac` and
`exchange_code` with a call to `shopify_credentials()`. Leave the settings themselves in
place — they are the fallback.

- [ ] **Step 4: Run the tests**

Run: `.venv/bin/python -m pytest tests/unit/test_shopify_oauth.py -v`
Expected: PASS. The existing tests still pass because with no `platform_configs` row the
function returns the settings values they already assert against.

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

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

Do not commit and do not stage. Do **not** insert a `shopify` row into `platform_configs`
from code — that is an operator action, and the row carries a real client secret.

---

### Task 5: Callback writes `integrations` and tracker rows

**Files:**
- Modify: `app/api/shopify.py`
- Test: `tests/unit/test_shopify_routes.py`

**Interfaces:**
- Consumes: `platform_integrations.upsert_integration`, `track_connect`, `TenantHasNoUserError`.
- Produces: no new public functions; the callback gains conformance writes.

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

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

```python
def test_install_tracks_authorize_url_ready(client, fake_redis, monkeypatch):
    tracked = []
    monkeypatch.setattr(shopify_api, "track_connect",
                        lambda *a, **kw: tracked.append((a, kw)))

    client.get("/api/shopify/install",
               params={"shop": "galaxiq-braexaal.myshopify.com",
                       "tenant_id": "org_test"},
               follow_redirects=False)

    assert tracked, "install must record that an authorize URL was issued"
    assert "oauth_authorize_url_ready" in str(tracked[0])


def test_callback_records_the_integration(client, fake_redis, happy_path, monkeypatch):
    recorded = {}
    monkeypatch.setattr(shopify_api, "upsert_integration",
                        lambda **kw: recorded.update(kw) or "int-1")
    monkeypatch.setattr(shopify_api, "track_connect", lambda *a, **kw: None)

    state = _seed_state(fake_redis)
    resp = _callback(client, shop="galaxiq-braexaal.myshopify.com",
                     code="c", state=state, hmac="x")

    assert resp.status_code == 307
    assert recorded["provider"] == "shopify"
    assert recorded["auth_type"] == "oauth2"
    assert recorded["metadata"]["shop_domain"] == "galaxiq-braexaal.myshopify.com"
    # Ciphertext, never a raw shpat_ token.
    assert not recorded["access_token_ciphertext"].startswith("shpat_")


def test_callback_reports_a_tenant_without_a_user(client, fake_redis, happy_path,
                                                  monkeypatch):
    from app.services.platform_integrations import TenantHasNoUserError

    def boom(**kw):
        raise TenantHasNoUserError("no userId")

    monkeypatch.setattr(shopify_api, "upsert_integration", boom)
    monkeypatch.setattr(shopify_api, "track_connect", lambda *a, **kw: None)

    state = _seed_state(fake_redis)
    resp = _callback(client, shop="galaxiq-braexaal.myshopify.com",
                     code="c", state=state, hmac="x")

    assert resp.status_code == 409
    assert "no_owning_user" in resp.headers["location"]


def test_a_tracker_failure_does_not_fail_the_connection(client, fake_redis,
                                                        happy_path, monkeypatch):
    def boom(*a, **kw):
        raise RuntimeError("tracker table unreachable")

    monkeypatch.setattr(shopify_api, "track_connect", boom)
    monkeypatch.setattr(shopify_api, "upsert_integration", lambda **kw: "int-1")

    state = _seed_state(fake_redis)
    resp = _callback(client, shop="galaxiq-braexaal.myshopify.com",
                     code="c", state=state, hmac="x")
    assert resp.status_code == 307
```

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

Run: `.venv/bin/python -m pytest tests/unit/test_shopify_routes.py -k "track or integration or owning" -v`
Expected: FAIL — `shopify_api` has no `track_connect` or `upsert_integration`.

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

Add to the top import block of `app/api/shopify.py`:

```python
from app.services.platform_integrations import (
    TenantHasNoUserError, track_connect, upsert_integration,
)
from app.services.product_sources import encrypt_credentials
```

In `install`, after the Redis write and before returning the redirect:

```python
    track_connect(tenant_id, "shopify", "oauth_authorize_url_ready", "success")
```

In `callback`, after `upsert_source` succeeds, record the connection on the platform's
tables. Wrap the tracker call so a bookkeeping failure cannot fail a completed connection:

```python
    try:
        upsert_integration(
            tenant_id=stored["tenant_id"],
            provider="shopify",
            auth_type="oauth2",
            access_token_ciphertext=encrypt_credentials(
                {"access_token": token_data["access_token"]}),
            metadata={
                "shop_domain": domain,
                "shop_name": info["name"],
                "scopes": token_data["scope"],
                "api_version": SHOPIFY_API_VERSION,
                "currency_code": info["currency_code"],
                "primary_domain": info["primary_domain"],
            },
        )
    except TenantHasNoUserError:
        # integrations.user is NOT NULL and this tenant has no owner, so the
        # platform cannot record the connection. Say so rather than leaving a
        # half-connected state the merchant cannot see.
        logger.error("Tenant %s has no owning user; cannot record integration",
                     stored["tenant_id"])
        _safe_track(stored["tenant_id"], "not_connected", "failed",
                    "tenant has no owning user")
        return _error_redirect("no_owning_user", 409)

    _safe_track(stored["tenant_id"], "connected", "success")
```

And add the helper beside `_error_redirect`:

```python
def _safe_track(tenant_id: str, delivery_status: str, outcome: str,
                error_message=None) -> None:
    """The audit trail must never be the reason a connection fails."""
    try:
        track_connect(tenant_id, "shopify", delivery_status, outcome, error_message)
    except Exception:
        logger.warning("Connect tracking failed for %s", tenant_id, exc_info=True)
```

Also call `_safe_track(..., "not_connected", "failed", <reason>)` before each existing
`_error_redirect` in the callback, passing the same reason string already used there. Never
pass a token, code, or response body as the error message.

- [ ] **Step 4: Run the tests**

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

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

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

Do not commit and do not stage.

---

### Task 6: Capture relations and prominence

**Files:**
- Modify: `app/services/catalog_shopify.py`
- Modify: `app/services/catalog_http.py`
- Modify: `app/services/catalog_normalise.py`
- Test: `tests/unit/test_catalog_shopify.py`, `tests/unit/test_catalog_http.py`, `tests/unit/test_catalog_normalise.py`

**Interfaces:**
- Consumes: existing adapters and `normalise`.
- Produces: `SourceProduct` gains `tenant_relations: dict`, `rating`, `review_count`, `featured_rank`; the normalised row gains the same four keys.

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

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

```python
def test_complementary_metafield_becomes_tenant_relations():
    node = dict(load("no-category"))
    node["complementary"] = {"value": '["gid://shopify/Product/1","gid://shopify/Product/2"]'}
    sp = to_source_product(node, SHOP)
    assert sp["tenant_relations"]["related"] == [
        "gid://shopify/Product/1", "gid://shopify/Product/2"]


def test_absent_complementary_metafield_yields_an_empty_object():
    # Never None: downstream reads .get("related") without a guard.
    sp = to_source_product(load("no-category"), SHOP)
    assert sp["tenant_relations"] == {}


def test_malformed_complementary_value_is_ignored():
    # A merchant or app can write anything into a metafield; a sync must not die.
    node = dict(load("no-category"))
    node["complementary"] = {"value": "not json at all"}
    assert to_source_product(node, SHOP)["tenant_relations"] == {}


def test_query_requests_the_complementary_metafield():
    assert "shopify--discovery--product_recommendation" in PRODUCTS_QUERY
    assert "complementary_products" in PRODUCTS_QUERY


def test_shopify_has_no_rating_or_review_count():
    # The Admin API exposes neither; the prominence term drops out, which the
    # matching specification permits as its option 4.
    sp = to_source_product(load("no-category"), SHOP)
    assert sp["rating"] is None
    assert sp["review_count"] is None
```

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

```python
def test_relations_and_prominence_map_from_config():
    cfg = dict(CONFIG)
    cfg["fields"] = dict(CONFIG["fields"], related_ids="$.related",
                         cross_sells="$.crossSells", rating="$.rating",
                         review_count="$.reviewCount")
    record = dict(PAGE["products"][0], related=[7, 8], crossSells=[9],
                  rating=4.5, reviewCount=120)

    sp = to_source_product(record, cfg, "dummyjson.com")

    assert sp["tenant_relations"] == {"related": ["7", "8"], "cross_sells": ["9"]}
    assert sp["rating"] == 4.5
    assert sp["review_count"] == 120


def test_unmapped_relations_yield_an_empty_object():
    sp = to_source_product(PAGE["products"][0], CONFIG, "dummyjson.com")
    assert sp["tenant_relations"] == {}
    assert sp["rating"] is None
```

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

```python
def test_description_no_longer_affects_content_hash():
    # The matching specification's embedding input excludes description, so a
    # description-only edit must not trigger a re-embed.
    a = normalise(sp(description="One description"), CONTEXT)
    b = normalise(sp(description="A completely different description"), CONTEXT)
    assert content_hash(a) == content_hash(b)
    assert record_hash(a) != record_hash(b)


def test_a_product_without_a_url_is_kept_and_flagged():
    # A URL is genuinely needed, but rejecting the product hides the problem:
    # the merchant sees an empty catalog and no reason for it.
    out = normalise(sp(handle=None, product_url=None), CONTEXT)
    assert out is not None
    assert out["product_url"] is None
    assert out["missing_fields"] == ["product_url"]


def test_a_product_with_a_url_has_no_missing_fields():
    assert normalise(sp(), CONTEXT)["missing_fields"] == []


def test_missing_fields_is_never_none():
    # Downstream reads it as a list without a guard.
    assert isinstance(normalise(sp(), CONTEXT)["missing_fields"], list)


def test_the_hard_rejections_still_reject():
    # These cannot be flagged and carried -- there is nothing usable left.
    for bad in ({"external_id": None}, {"title": None}, {"title": "ab"},
                {"variants": []}):
        assert normalise(sp(**bad), CONTEXT) is None


def test_relations_and_prominence_are_carried_through():
    out = normalise(sp(tenant_relations={"related": ["1"]}, rating=4.5,
                       review_count=120, featured_rank=3), CONTEXT)
    assert out["tenant_relations"] == {"related": ["1"]}
    assert out["rating"] == 4.5
    assert out["review_count"] == 120
    assert out["featured_rank"] == 3


def test_missing_relations_normalise_to_an_empty_object():
    assert normalise(sp(), CONTEXT)["tenant_relations"] == {}
```

Add `tenant_relations`, `rating`, `review_count` and `featured_rank` to the `sp()` helper's
base dict in `test_catalog_normalise.py`, defaulting to `{}`, `None`, `None`, `None`.

- [ ] **Step 2: Run the tests to verify they fail**

Run: `.venv/bin/python -m pytest tests/unit/test_catalog_shopify.py tests/unit/test_catalog_http.py tests/unit/test_catalog_normalise.py -v`
Expected: FAIL — the new keys do not exist.

- [ ] **Step 3: Shopify adapter**

In `app/services/catalog_shopify.py`, add to `PRODUCTS_QUERY` inside the product node,
beside `metafields`:

```graphql
      complementary: metafield(
        namespace: "shopify--discovery--product_recommendation"
        key: "complementary_products"
      ) { value }
```

And in `to_source_product`, add:

```python
        "tenant_relations": _complementary(node),
        "rating": None,
        "review_count": None,
        "featured_rank": None,
```

with the helper:

```python
def _complementary(node: dict) -> dict:
    """Shopify has no native related-products field; the Search & Discovery app
    stores hand-linked products in this metafield. A merchant or app can write
    anything into it, so a malformed value is ignored rather than fatal."""
    raw = (node.get("complementary") or {}).get("value")
    if not raw:
        return {}
    try:
        ids = json.loads(raw)
    except (TypeError, ValueError):
        logger.warning("Ignoring malformed complementary_products metafield")
        return {}
    if not isinstance(ids, list) or not ids:
        return {}
    return {"related": [str(i) for i in ids]}
```

Add `import json` to the top import block.

- [ ] **Step 4: HTTP adapter**

In `app/services/catalog_http.py`, inside `to_source_product`, build the relations from the
mapped values and add the four keys:

```python
    relations = {}
    for target, key in (("related_ids", "related"),
                        ("cross_sells", "cross_sells"),
                        ("upsells", "upsells")):
        values = mapped.get(target)
        if isinstance(values, list) and values:
            relations[key] = [str(v) for v in values]
```

then in the returned dict:

```python
        "tenant_relations": relations,
        "rating": mapped.get("rating"),
        "review_count": mapped.get("review_count"),
        "featured_rank": mapped.get("featured_rank"),
```

- [ ] **Step 5: Normaliser**

In `app/services/catalog_normalise.py`, **remove the early return that rejects a product
with no resolvable URL**, and instead record it. Replace that rejection with:

```python
    url = resolve_product_url(source_product,
                              context.get("primary_domain"),
                              context.get("url_template"))

    # A URL is genuinely needed, but rejecting the product hides the problem:
    # the merchant would see an empty catalog with no explanation. Keep it,
    # flag it, and let it be excluded from serving instead.
    missing_fields = [] if url else ["product_url"]
```

and add to the returned dict:

```python
        "missing_fields": missing_fields,
        "tenant_relations": source_product.get("tenant_relations") or {},
        "rating": source_product.get("rating"),
        "review_count": source_product.get("review_count"),
        "featured_rank": source_product.get("featured_rank"),
```

Leave every other rejection exactly as it is — no `external_id`, no title or a title under
three characters, no price or a negative price, no variants. Those cannot be flagged and
carried because nothing usable remains.

And remove the description entry from `content_hash`'s `parts` list, leaving a comment:

```python
        # Description is deliberately absent: the matching specification's
        # embedding input excludes it, so hashing it would re-embed a catalog
        # after an edit that cannot change any embedding.
```

- [ ] **Step 6: Run the tests**

Run: `.venv/bin/python -m pytest tests/unit/ -q`
Expected: PASS. Note the count rises by the new tests.

- [ ] **Step 7: Persist the new fields and report incompleteness**

In `app/services/catalog_sync.py`'s `persist`, add `tenant_relations`, `rating`,
`review_count`, `featured_rank`, `missing_fields`, `price_reference_cents` and
`fx_rate_used` to both the INSERT column list and the `ON CONFLICT DO UPDATE SET` clause,
passing `Json(p["tenant_relations"])`, the scalars, and `p["missing_fields"]` in the values
tuple. Follow the existing ordering exactly.

Then in `sync_source`, count incomplete products alongside rejected ones and include the
figure in the returned report:

```python
    incomplete = sum(1 for p in normalised if p["missing_fields"])
```

Add `"incomplete": incomplete` to the returned dict, and log a warning when it is non-zero,
naming the most common missing field — a source producing unlinkable products must be
obvious in the sync response rather than discovered later by a merchant wondering why
nothing is recommended.

Add a unit test to `tests/unit/test_catalog_sync.py` asserting the report carries
`incomplete` and that a product with a populated `missing_fields` is counted in it.

- [ ] **Step 8: Run the integration suite**

Run: `.venv/bin/python -m pytest tests/integration/ -q`
Expected: PASS.

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

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

Do not commit and do not stage.

---

### Task 6b: Currency normalisation

Without this, every price comparison across two sources is silently wrong.

**Files:**
- Create: `app/services/fx.py`
- Modify: `app/services/catalog_normalise.py`
- Test: `tests/unit/test_fx.py`
- Test: `tests/integration/test_fx_table.py`

**Interfaces:**
- Consumes: `get_master_db_connection`.
- Produces:
  - `ensure_fx_table() -> None`
  - `get_rate(base: str, quote: str) -> float | None` — returns `1.0` when base equals quote
  - `to_reference_cents(price_cents, source_currency, reference_currency) -> (int | None, float | None)` returning `(converted, rate_used)`

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

Create `tests/unit/test_fx.py`:

```python
import pytest

from app.services.fx import to_reference_cents


def test_same_currency_needs_no_rate(monkeypatch):
    import app.services.fx as mod
    monkeypatch.setattr(mod, "get_rate", lambda b, q: None)  # never consulted
    assert to_reference_cents(88595, "INR", "INR") == (88595, 1.0)


def test_converts_with_the_stored_rate(monkeypatch):
    import app.services.fx as mod
    monkeypatch.setattr(mod, "get_rate", lambda b, q: 0.012)
    converted, rate = to_reference_cents(88595, "INR", "USD")
    assert converted == 1063          # round(88595 * 0.012)
    assert rate == 0.012


def test_missing_rate_returns_none_rather_than_guessing(monkeypatch):
    # A guessed rate produces confident, wrong recommendations. A null makes the
    # product incomplete, which is visible.
    import app.services.fx as mod
    monkeypatch.setattr(mod, "get_rate", lambda b, q: None)
    assert to_reference_cents(88595, "INR", "USD") == (None, None)


@pytest.mark.parametrize("price,src,ref", [
    (None, "INR", "USD"),
    (1000, None, "USD"),
    (1000, "INR", None),
])
def test_missing_inputs_return_none(price, src, ref):
    assert to_reference_cents(price, src, ref) == (None, None)


def test_conversion_stays_an_integer(monkeypatch):
    import app.services.fx as mod
    monkeypatch.setattr(mod, "get_rate", lambda b, q: 0.0123456)
    converted, _ = to_reference_cents(99999, "INR", "USD")
    assert isinstance(converted, int)
```

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

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

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

Create `app/services/fx.py`:

```python
"""Converts prices into a tenant's reference currency so they can be compared.

Every price rule in the matching specification divides one price by another. Across
currencies that is meaningless and silent: an INR price against a USD one yields a
ratio dozens of times out, and the output looks like a ranking opinion rather than a
data fault. Conversion happens once, at ingestion.
"""
import logging

from app.services.database import get_master_db_connection

logger = logging.getLogger(__name__)


def ensure_fx_table() -> None:
    conn = get_master_db_connection()
    try:
        with conn.cursor() as cur:
            cur.execute("""
                CREATE TABLE IF NOT EXISTS fx_rates (
                    base_currency  TEXT NOT NULL,
                    quote_currency TEXT NOT NULL,
                    rate           REAL NOT NULL,
                    fetched_at     TIMESTAMPTZ NOT NULL DEFAULT now(),
                    PRIMARY KEY (base_currency, quote_currency)
                )
            """)
        conn.commit()
    finally:
        conn.close()


def get_rate(base: str, quote: str):
    if not base or not quote:
        return None
    if base == quote:
        return 1.0

    ensure_fx_table()
    conn = get_master_db_connection()
    try:
        with conn.cursor() as cur:
            cur.execute(
                "SELECT rate FROM fx_rates WHERE base_currency = %s "
                "AND quote_currency = %s", (base, quote))
            row = cur.fetchone()
            return float(row[0]) if row else None
    finally:
        conn.close()


def to_reference_cents(price_cents, source_currency, reference_currency):
    """Returns (converted_cents, rate_used), or (None, None) when no rate exists.

    Deliberately does not fall back to a default rate: a guessed conversion is
    indistinguishable from a real one in the output, whereas a null makes the
    product incomplete and therefore visible.
    """
    if price_cents is None or not source_currency or not reference_currency:
        return None, None

    if source_currency == reference_currency:
        return int(price_cents), 1.0

    rate = get_rate(source_currency, reference_currency)
    if rate is None:
        logger.warning("No fx rate for %s -> %s; product will be incomplete",
                       source_currency, reference_currency)
        return None, None

    return int(round(price_cents * rate)), rate
```

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

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

- [ ] **Step 5: Wire it into the normaliser**

In `app/services/catalog_normalise.py`, after the price range is computed, convert and
extend `missing_fields`:

```python
    reference_currency = context.get("reference_currency")
    price_reference_cents, fx_rate_used = to_reference_cents(
        lo, context.get("currency"), reference_currency)

    if price_reference_cents is None:
        # Comparing an unconverted price is worse than not recommending the
        # product: the failure is invisible and looks like a ranking decision.
        missing_fields.append("price_reference")
```

and add both to the returned dict. Import `to_reference_cents` from `app.services.fx`.

- [ ] **Step 6: Set the reference currency per tenant**

In `app/services/catalog_sync.py`'s `sync_source`, resolve the tenant's reference currency
once and pass it in `context`:

```python
    reference_currency = get_setting(tenant_id, "reference_currency")
    if not reference_currency:
        # The first connected source defines it. Recording it explicitly stops it
        # drifting when a second source with a different currency is added.
        reference_currency = context_currency
        update_setting(tenant_id, "reference_currency", reference_currency)
```

using the existing `get_setting` / `update_setting` from `app.services.database`, where
`context_currency` is the currency already resolved for this source. Add
`"reference_currency": reference_currency` to the `context` dict.

- [ ] **Step 7: Write the database test**

Create `tests/integration/test_fx_table.py`:

```python
from app.services.database import get_master_db_connection
from app.services.fx import ensure_fx_table, get_rate


def test_table_is_created_and_idempotent():
    ensure_fx_table()
    ensure_fx_table()
    conn = get_master_db_connection()
    try:
        with conn.cursor() as cur:
            cur.execute("SELECT to_regclass('public.fx_rates')")
            assert cur.fetchone()[0] is not None
    finally:
        conn.close()


def test_same_currency_is_one_without_a_row():
    assert get_rate("INR", "INR") == 1.0


def test_absent_pair_returns_none():
    assert get_rate("XXX", "YYY") is None


def test_a_seeded_rate_is_read_back():
    ensure_fx_table()
    conn = get_master_db_connection()
    try:
        with conn.cursor() as cur:
            cur.execute(
                "INSERT INTO fx_rates (base_currency, quote_currency, rate) "
                "VALUES ('TST','REF',2.5) ON CONFLICT (base_currency, quote_currency) "
                "DO UPDATE SET rate = 2.5")
        conn.commit()
    finally:
        conn.close()

    assert get_rate("TST", "REF") == 2.5

    conn = get_master_db_connection()
    try:
        with conn.cursor() as cur:
            cur.execute("DELETE FROM fx_rates WHERE base_currency = 'TST'")
        conn.commit()
    finally:
        conn.close()
```

- [ ] **Step 8: Run both suites**

Run: `.venv/bin/python -m pytest tests/unit/ -q` then `.venv/bin/python -m pytest tests/integration/ -q`
Expected: PASS. Note the Shopify products already in the live tenant are `INR`, which will
become the reference currency, so they convert at `1.0` and gain no `missing_fields`.

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

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

Do not commit and do not stage.

---

### Task 7: Field-map proposer

**Files:**
- Create: `app/services/fieldmap_proposer.py`
- Test: `tests/unit/test_fieldmap_proposer.py`

**Interfaces:**
- Consumes: `catalog_fieldmap.resolve_path`, `FieldMapError`; `app.core.llm_client.client`; `settings.LLM_MODEL`.
- Produces:
  - `summarise_record(record: dict, max_value_len: int = 80) -> dict`
  - `validate_proposal(proposal: dict, sample: dict) -> tuple[dict, list]` returning `(kept, dropped)`
  - `propose_field_map(sample: dict) -> dict` returning `{"fields": {...}, "dropped": [...]}`
  - `infer_url_template(base_url: str, fields: dict) -> str | None`
  - `TARGET_FIELDS: tuple[str, ...]`
  - `NOT_MODELLED: tuple[str, ...]` — `("url_template", "currency")`, inferred separately rather than asked of the model

**This module is called by the sync pipeline, not by a route.** The merchant never sees its
output. It runs on a source's first sync and its result is cached in
`product_sources.config.fields`, so a steady-state sync makes zero LLM calls.

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

Create `tests/unit/test_fieldmap_proposer.py`:

```python
import json
import pathlib

from app.services.fieldmap_proposer import (
    NOT_MODELLED, TARGET_FIELDS, infer_url_template, propose_field_map, summarise_record,
    validate_proposal,
)

SAMPLE = json.loads(
    pathlib.Path("tests/fixtures/http/dummyjson-page.json").read_text()
)["products"][0]


def test_summary_keeps_shape_and_drops_bulk():
    long_text = "x" * 5000
    summary = summarise_record({"id": 1, "title": "T", "description": long_text,
                                "tags": ["a", "b"], "dims": {"w": 1}})
    assert summary["id"]["type"] == "int"
    assert summary["tags"]["type"] == "list"
    assert summary["dims"]["type"] == "dict"
    # A single record can be tens of kilobytes of prose that teaches the model
    # nothing about shape.
    assert len(summary["description"]["example"]) <= 80


def test_validate_drops_paths_that_do_not_resolve():
    kept, dropped = validate_proposal(
        {"title": "$.title", "brand": "$.nope", "price": "$.price"}, SAMPLE)
    assert kept == {"title": "$.title", "price": "$.price"}
    assert dropped == ["brand"]


def test_validate_drops_malformed_paths_without_raising():
    kept, dropped = validate_proposal({"title": "title", "price": "$.price"}, SAMPLE)
    assert kept == {"price": "$.price"}
    assert "title" in dropped


def test_validate_drops_unknown_target_names():
    kept, dropped = validate_proposal(
        {"title": "$.title", "not_a_target": "$.id"}, SAMPLE)
    assert "not_a_target" not in kept


def test_url_template_and_currency_are_never_asked_of_the_model():
    # Neither appears in a payload that does not contain them, so asking the
    # model invites a hallucinated path. They are inferred separately.
    for name in NOT_MODELLED:
        assert name not in TARGET_FIELDS


@pytest.mark.parametrize("base_url,expected", [
    ("https://dummyjson.com/products",
     "https://dummyjson.com/products/{external_id}"),
    ("https://dummyjson.com/products?limit=30",
     "https://dummyjson.com/products/{external_id}"),
    ("https://shop.example.com/api/v2/items/",
     "https://shop.example.com/api/v2/items/{external_id}"),
])
def test_url_template_inferred_from_the_base_url(base_url, expected):
    # A REST collection endpoint almost always addresses one member by id.
    assert infer_url_template(base_url, {"external_id": "$.id"}) == expected


def test_no_url_template_without_an_id_mapping():
    # Without an id there is nothing to substitute, so guessing would produce a
    # template that renders a broken link for every product.
    assert infer_url_template("https://dummyjson.com/products", {}) is None


def test_propose_returns_only_fields(monkeypatch):
    import app.services.fieldmap_proposer as mod
    monkeypatch.setattr(mod, "_call_llm",
                        lambda summary: {"title": "$.title", "price": "$.price"})

    out = propose_field_map(SAMPLE)
    assert out["fields"] == {"title": "$.title", "price": "$.price"}


def test_propose_survives_an_llm_failure(monkeypatch):
    # A model outage must leave the source connected and retryable, not broken.
    import app.services.fieldmap_proposer as mod

    def boom(summary):
        raise RuntimeError("model unavailable")

    monkeypatch.setattr(mod, "_call_llm", boom)
    assert propose_field_map(SAMPLE)["fields"] == {}


def test_propose_survives_non_json_from_the_model(monkeypatch):
    import app.services.fieldmap_proposer as mod
    monkeypatch.setattr(mod, "_call_llm", lambda summary: "not a dict")
    assert propose_field_map(SAMPLE)["fields"] == {}
```

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

Run: `.venv/bin/python -m pytest tests/unit/test_fieldmap_proposer.py -v`
Expected: FAIL with `ModuleNotFoundError`

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

Create `app/services/fieldmap_proposer.py`:

```python
"""Proposes a field map for an unknown product API from one sample record.

The proposal is never trusted: every path is executed against the sample before
it is returned, so a merchant is never shown a mapping that would fail at sync
time. Nothing is stored here -- the caller approves a map before it is saved.
"""
import json
import logging

from app.core.config import settings
from app.core.llm_client import client
from app.services.catalog_fieldmap import FieldMapError, resolve_path

logger = logging.getLogger(__name__)

TARGET_FIELDS = (
    "external_id", "title", "description", "brand", "raw_category",
    "image_url", "price", "discount_percentage", "availability", "sku",
    "tags", "related_ids", "cross_sells", "upsells", "rating", "review_count",
)

# Not asked of the model: neither appears in a payload that lacks them, so a
# request invites a hallucinated path. url_template is derived from the base URL
# below; currency falls back to the tenant's existing one.
NOT_MODELLED = ("url_template", "currency")

_PROMPT = (
    "You map an unknown e-commerce product API onto a fixed schema.\n"
    "You are given a summary of ONE product record: each key, its type, and a "
    "truncated example value.\n"
    "Return JSON only: an object whose keys are target field names and whose "
    "values are paths.\n"
    "Valid target names: {targets}\n"
    "Path syntax is exactly $.key, $.key.subkey and $.key[0]. Nothing else.\n"
    "Omit a target you cannot map. Never invent a key that is not in the summary."
)


def summarise_record(record: dict, max_value_len: int = 80) -> dict:
    summary = {}
    for key, value in (record or {}).items():
        example = value
        if isinstance(value, str):
            example = value[:max_value_len]
        elif isinstance(value, list):
            example = value[:2]
        elif isinstance(value, dict):
            example = list(value)[:5]
        summary[key] = {"type": type(value).__name__, "example": example}
    return summary


def validate_proposal(proposal: dict, sample: dict) -> tuple:
    kept, dropped = {}, []
    for target, path in (proposal or {}).items():
        if target not in TARGET_FIELDS:
            dropped.append(target)
            continue
        try:
            if resolve_path(sample, path) is None:
                dropped.append(target)
                continue
        except FieldMapError:
            dropped.append(target)
            continue
        kept[target] = path
    return kept, sorted(dropped)


def _call_llm(summary: dict) -> dict:
    response = client.chat.completions.create(
        model=settings.LLM_MODEL,
        response_format={"type": "json_object"},
        messages=[
            {"role": "system",
             "content": _PROMPT.format(targets=", ".join(TARGET_FIELDS))},
            {"role": "user", "content": json.dumps(summary, default=str)},
        ],
    )
    return json.loads(response.choices[0].message.content)


def propose_field_map(sample: dict) -> dict:
    try:
        raw = _call_llm(summarise_record(sample))
    except Exception:
        # A model outage must leave the merchant able to map by hand.
        logger.error("Field-map proposal failed", exc_info=True)
        raw = {}

    if not isinstance(raw, dict):
        logger.warning("Field-map proposal was not an object; ignoring")
        raw = {}

    kept, dropped = validate_proposal(raw, sample)
    return {"fields": kept, "dropped": dropped}


def infer_url_template(base_url: str, fields: dict):
    """Derive a per-product URL from a collection endpoint.

    A REST collection almost always addresses one member by id, so
    ".../products" plus an id mapping yields ".../products/{external_id}".
    Without an id mapping there is nothing to substitute, and guessing would
    give every product a broken link -- so this returns None and those products
    are stored with a null URL instead.
    """
    if not base_url or "external_id" not in (fields or {}):
        return None
    root = base_url.split("?")[0].rstrip("/")
    return f"{root}/{{external_id}}"
```

- [ ] **Step 4: Run the tests**

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

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

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

Do not commit and do not stage.

---

### Task 8: Sources router

**Scope note — an earlier draft had a fourth endpoint, `POST /sources/http/preview`, which
proposed a field map for a merchant to approve before connecting. It was removed at the
user's direction.** A schema mapping is a machine problem, not a merchant decision, and
requiring approval put a human gate in front of every connection. Human approval exists in
exactly one place in this system — low-confidence pairings — and that is a later
sub-project. **Do not build a preview endpoint, and do not require a field map,
`url_template` or `currency` at connect time.** A merchant supplies a base URL and
credentials; everything else is inferred by the pipeline on first sync.

**Files:**
- Create: `app/api/sources.py`
- Modify: `app/main.py`
- Modify: `app/services/catalog_http.py` — SSRF guard on fetch, plus first-sync inference
- Modify: `app/services/platform_integrations.py` — `set_integration_status` used on delete
- Test: `tests/integration/test_sources_route.py`

**Interfaces:**
- Consumes: `ssrf.assert_safe_url`, `product_sources.upsert_source`/`get_sources`/`mark_source_status`, `platform_integrations.get_integration`/`set_integration_status`.
- Produces: `router: APIRouter(prefix="/sources")` with `POST /http`, `GET ""`, `DELETE /{kind}/{external_ref}`.

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

Create `tests/integration/test_sources_route.py`:

```python
import pytest
from fastapi.testclient import TestClient

from app.main import app

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


def test_connect_needs_only_a_base_url(monkeypatch):
    # No field map, no url_template, no currency. Everything else is inferred by
    # the pipeline on first sync.
    from app.api import sources
    saved = {}
    monkeypatch.setattr(sources, "upsert_source",
                        lambda *a, **kw: saved.update({"args": a}))

    resp = client.post("/sources/http", json={
        "tenant_id": TENANT, "base_url": "https://dummyjson.com/products"})

    assert resp.status_code == 200
    assert resp.json()["external_ref"] == "dummyjson.com"
    assert saved, "the source must be stored"


def test_connect_rejects_an_internal_base_url():
    resp = client.post("/sources/http", json={
        "tenant_id": TENANT, "base_url": "http://127.0.0.1:8001/products"})
    assert resp.status_code == 400


def test_connect_rejects_a_non_http_scheme():
    resp = client.post("/sources/http", json={
        "tenant_id": TENANT, "base_url": "file:///etc/passwd"})
    assert resp.status_code == 400


def test_connect_rejects_cloud_metadata():
    resp = client.post("/sources/http", json={
        "tenant_id": TENANT, "base_url": "http://169.254.169.254/latest/meta-data/"})
    assert resp.status_code == 400


def test_optional_overrides_are_accepted(monkeypatch):
    # A merchant whose API defeats inference can still supply values by hand;
    # none of them is required.
    from app.api import sources
    saved = {}
    monkeypatch.setattr(sources, "upsert_source",
                        lambda t, k, r, config, creds: saved.update(config))

    client.post("/sources/http", json={
        "tenant_id": TENANT, "base_url": "https://dummyjson.com/products",
        "currency": "USD",
        "url_template": "https://dummyjson.com/products/{external_id}"})

    assert saved["currency"] == "USD"
    assert saved["url_template"].endswith("{external_id}")


def test_list_sources_includes_the_shopify_connection():
    resp = client.get("/sources", params={"tenant_id": TENANT})
    assert resp.status_code == 200
    kinds = {s["kind"] for s in resp.json()["sources"]}
    assert "shopify" in kinds


def test_list_never_returns_credentials():
    body = client.get("/sources", params={"tenant_id": TENANT}).text
    for leak in ("shpat_", "credentials", "accessToken", "api_key"):
        assert leak not in body


def test_delete_unknown_source_404s():
    resp = client.delete("/sources/http_api/not-connected.example")
    assert resp.status_code in (404, 422)


def test_tenant_id_is_required():
    assert client.get("/sources").status_code == 422
```

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

Run: `.venv/bin/python -m pytest tests/integration/test_sources_route.py -v`
Expected: FAIL — every route 404s.

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

Create `app/api/sources.py`:

```python
"""Connect and manage product ingestion sources."""
import logging

import httpx
from fastapi import APIRouter, Body, HTTPException, Query

from app.services.platform_integrations import get_integration, set_integration_status
from app.services.product_sources import (
    KIND_HTTP_API, KIND_SHOPIFY, encrypt_credentials, get_sources, upsert_source,
)
from app.services.ssrf import UnsafeUrlError, assert_safe_url

logger = logging.getLogger(__name__)

router = APIRouter(prefix="/sources", tags=["sources"])

@router.post("/http")
async def create_http_source(payload: dict = Body(...)):
    """Connect a custom product API.

    Only a base URL and credentials are required. The field map, url_template
    and currency are inferred by the pipeline on first sync -- a merchant should
    not have to describe their own API's shape, and an API cannot be made to
    return something it does not have.
    """
    tenant_id = payload.get("tenant_id")
    base_url = payload.get("base_url")
    if not tenant_id or not base_url:
        raise HTTPException(status_code=422,
                            detail="tenant_id and base_url are required")

    try:
        assert_safe_url(base_url)
    except UnsafeUrlError as ex:
        raise HTTPException(status_code=400, detail=str(ex))

    external_ref = httpx.URL(base_url).host
    config = {
        "base_url": base_url,
        "records_path": payload.get("records_path") or "$.products",
        "pagination": payload.get("pagination") or {
            "style": "offset", "limit_param": "limit",
            "offset_param": "skip", "page_size": 30, "total_path": "$.total"},
        # Empty until first sync infers it. Overrides are accepted for an API
        # that defeats inference, but none is required.
        "fields": payload.get("fields") or {},
        "url_template": payload.get("url_template"),
        "currency": payload.get("currency"),
        "auth": payload.get("auth") or {"style": "none"},
    }
    credentials = {"api_key": payload["api_key"]} if payload.get("api_key") else {}

    upsert_source(tenant_id, KIND_HTTP_API, external_ref, config, credentials)
    return {"kind": KIND_HTTP_API, "external_ref": external_ref,
            "status": "active", "mapping": "pending_first_sync"}


@router.get("")
async def list_sources(tenant_id: str = Query(...)):
    """Every source for a tenant. Never includes credentials."""
    sources = [
        {"kind": s["kind"], "external_ref": s["external_ref"],
         "config": s["config"], "status": s["status"],
         "connected_at": s["connected_at"], "last_synced_at": s["last_synced_at"]}
        for s in get_sources(tenant_id)
    ]

    integration = get_integration(tenant_id, KIND_SHOPIFY)
    return {"sources": sources,
            "platform_integration": {
                "shopify": integration["status"] if integration else None}}


@router.delete("/{kind}/{external_ref}")
async def delete_source(kind: str, external_ref: str, tenant_id: str = Query(...)):
    matches = [s for s in get_sources(tenant_id)
               if s["kind"] == kind and s["external_ref"] == external_ref]
    if not matches:
        raise HTTPException(status_code=404, detail="No such connected source")

    from app.services.product_sources import mark_source_status
    mark_source_status(kind, external_ref, "revoked")
    return {"kind": kind, "external_ref": external_ref, "status": "revoked"}
```

- [ ] **Step 4: Mount the router**

In `app/main.py`, add beside the other router imports:

```python
from app.api.sources import router as sources_router
```

and beside the other `include_router` calls:

```python
app.include_router(sources_router)
```

- [ ] **Step 5: Guard the sync fetch**

In `app/services/catalog_http.py`'s `fetch_products`, immediately after reading `base`:

```python
    # A stored base URL is as user-supplied as a previewed one, and it is
    # fetched on every scheduled sync.
    assert_safe_url(base)
```

Add `from app.services.ssrf import assert_safe_url` to the top import block, and let
`UnsafeUrlError` propagate — `sync_source` already reports a failed source.

- [ ] **Step 6: Run the tests**

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

- [ ] **Step 7: Run both suites**

Run: `.venv/bin/python -m pytest tests/unit/ -q` then `.venv/bin/python -m pytest tests/integration/ -q`
Expected: PASS, no regressions.

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

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

Do not commit and do not stage.

---

### Task 9: Live verification

**Files:** none — verification only.

**Prerequisites:** uvicorn running from this worktree on port 8001; Postgres and Redis reachable; the Shopify source already connected from sub-project A.

- [ ] **Step 1: Start the server**

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

Expected: `200`.

- [ ] **Step 2: Confirm the SSRF guard rejects internal targets over HTTP**

```bash
for u in "http://169.254.169.254/latest/meta-data/" "http://127.0.0.1:8001/health" "file:///etc/passwd"; do
  printf "%-45s -> " "$u"
  curl -s -o /dev/null -w "%{http_code}\n" -X POST http://localhost:8001/sources/http \
    -H "Content-Type: application/json" \
    -d "{\"tenant_id\":\"org_8c32bf3e-6a18-4739-9b1c-94c0cf11125f\",\"base_url\":\"$u\"}"
done
```

Expected: `400` for all three, and **no entry for them in the server log's outbound
requests** — the guard must reject before any socket opens.

- [ ] **Step 3: Connect a real API — authenticated, with nothing but a URL and a key**

This exercises the **authenticated** path against the reference API's full catalogue of
**194 products**, which at the default page size is **7 pages** — real pagination, not a
single-page happy path.

First obtain a token. The reference API issues a JWT; ask for a long expiry so the run does
not fail halfway through for an unrelated reason:

```bash
TOKEN=$(curl -s -X POST https://dummyjson.com/auth/login \
  -H "Content-Type: application/json" \
  -d '{"username":"emilys","password":"emilyspass","expiresInMins":600}' \
  | python3 -c "import sys,json; print(json.load(sys.stdin)['accessToken'])")
```

Confirm the endpoint genuinely requires it — this proves the auth path is being tested
rather than silently bypassed:

```bash
curl -s "https://dummyjson.com/auth/products?limit=1"
```

Expected: `{"message":"Access Token is required"}`.

Now connect, supplying only the URL and the key:

```bash
curl -s -X POST http://localhost:8001/sources/http \
  -H "Content-Type: application/json" \
  -d "{\"tenant_id\":\"org_8c32bf3e-6a18-4739-9b1c-94c0cf11125f\",
       \"base_url\":\"https://dummyjson.com/auth/products\",
       \"api_key\":\"$TOKEN\",
       \"auth\":{\"style\":\"bearer\"}}" | python3 -m json.tool
```

Expected: `{"kind": "http_api", "external_ref": "dummyjson.com", "status": "active",
"mapping": "pending_first_sync"}`. No field map, no `url_template`, no `currency` supplied
— all three are inferred on first sync.

- [ ] **Step 4: Confirm nothing was inferred yet**

```bash
PYTHONPATH=. .venv/bin/python -c "
from app.services.product_sources import get_sources
s = [x for x in get_sources('org_8c32bf3e-6a18-4739-9b1c-94c0cf11125f')
     if x['kind'] == 'http_api'][0]
print('fields at connect time:', s['config']['fields'])
print('url_template:', s['config']['url_template'])"
```

Expected: an empty `fields` object and a null `url_template` — inference happens on sync,
not at connect.

- [ ] **Step 5: Sync it**

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

Expected: two results — the Shopify source with 17 products, and `dummyjson.com` with
**`fetched: 194`** and `normalised: 194`, `rejected: 0`.

**194 is the whole catalogue and the point of the exercise.** At the default page size of 30
that is seven pages, so a wrong offset, a premature loop exit, or a mishandled `total` shows
up as a count lower than 194 rather than passing unnoticed on a single page. A count of 30
means pagination stopped after page one. A hang means the `MAX_PAGES` bound is not working.

The report also carries `incomplete`, which should be **0** if `infer_url_template` derived
`https://dummyjson.com/auth/products/{external_id}` from the base URL. A non-zero `rejected`
means a mapping is wrong; a non-zero `incomplete` means URL inference failed. Read the logs
before changing code in either case.

Also confirm the **authenticated** fetch is what ran, not an accidental anonymous one: the
server log must show no `401`, and removing the token from the stored source must make the
next sync fail. That distinction matters because the unauthenticated endpoint returns the
same shape, so a broken auth header could otherwise look like success.

Then confirm the inference was cached rather than repeated:

```bash
PYTHONPATH=. .venv/bin/python -c "
from app.services.product_sources import get_sources
s = [x for x in get_sources('org_8c32bf3e-6a18-4739-9b1c-94c0cf11125f')
     if x['kind'] == 'http_api'][0]
print('inferred fields:', s['config']['fields'])
print('inferred url_template:', s['config']['url_template'])"
```

Expected: a populated field map and a derived template. Re-running the sync must make **no
further LLM call** — the stored map is what executes.

- [ ] **Step 6: Confirm both sources coexist and neither deleted the other**

```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 ORDER BY 1')
print('by source:', cur.fetchall())
cur.execute('SELECT name, product_url, currency, rating FROM \"'+T+'\".strategist_products WHERE source_kind=%s LIMIT 3', ('http_api',))
for r in cur.fetchall(): print(' ', r)
c.close()"
```

Expected: three source kinds — `crawl` 7, `shopify` 17, `http_api` **194**. HTTP products
carry resolvable URLs and a rating.

This is also the strongest test the scoped-delete guarantee has had: **194 rows arriving
from a third producer must not disturb the other 24.** Earlier runs proved it with 17.

Then confirm currency normalisation did the right thing across two currencies:

```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 currency, count(*), count(price_reference_cents) FROM \"'+T+'\".strategist_products WHERE price_cents IS NOT NULL GROUP BY 1')
print('currency, rows, with_reference:', cur.fetchall())
cur.execute('SELECT count(*) FROM \"'+T+'\".strategist_products WHERE %s = ANY(missing_fields)', ('price_reference',))
print('incomplete for want of an fx rate:', cur.fetchone()[0])
c.close()"
```

The reference currency is `INR`, set by the first-connected source. So the 17 Shopify rows
convert at `1.0`, and the 194 USD rows need a `USD -> INR` row in `fx_rates`. **Without one
they are correctly flagged incomplete rather than silently compared as bare numbers** —
which is the entire point of the currency work. Seed a rate and re-sync to watch the count
fall to zero:

```bash
PYTHONPATH=. .venv/bin/python -c "
from app.services.fx import ensure_fx_table
from app.services.database import get_master_db_connection
ensure_fx_table()
c=get_master_db_connection(); cur=c.cursor()
cur.execute(\"INSERT INTO fx_rates (base_currency, quote_currency, rate) VALUES ('USD','INR',83.0) ON CONFLICT (base_currency, quote_currency) DO UPDATE SET rate=83.0\")
c.commit(); c.close(); print('seeded USD -> INR')"
```

- [ ] **Step 7: Confirm the platform integration row**

```bash
PYTHONPATH=. .venv/bin/python -c "
from app.services.platform_integrations import get_integration
print(get_integration('org_8c32bf3e-6a18-4739-9b1c-94c0cf11125f', 'shopify'))"
```

Expected: `None` until a Shopify reconnect happens (Task 5 writes on callback, not
retroactively). Re-run the OAuth flow from the chat UI, then repeat — expect a row with
`status ACTIVE` and `metadata.token_encoding == "fernet"`.

- [ ] **Step 8: Report**

Summarise: changed files, test counts, the proposal the LLM returned, the sync numbers for
both sources, the row counts by source kind, and whether the integration row appeared.
Remind the user to commit, and that a `shopify` row in `platform_configs` is still an
operator action.

---

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

- **C2:** A4 attribute extraction — the matching specification's stated prerequisite.
- **C3:** `product_neighbors` and the offline matching job.
- **C4:** complement inference plus the merchandiser override surface.
- **C5:** online ranking B1–B5, replacing `decide()`'s internals.
- **C6:** D1 measurement counters.
- **Operator actions:** insert the `shopify` row into `platform_configs` with a rotated
  secret; publish the Online Store channel so synced Shopify URLs resolve; close
  `TODO(auth)`, which would remove the `Tenants.userId` derivation entirely.
