"""Storage for every product ingestion source, keyed by (kind, external_ref).

Lives in the master database rather than a per-tenant schema for two reasons.
Shopify webhooks arrive keyed only by shop domain with no tenant, so per-tenant
storage would mean scanning every schema to route one. And the UNIQUE
constraint enforces one-shop-one-tenant globally, which per-schema cannot.
"""
import json
import logging

from cryptography.fernet import Fernet
from psycopg2.extras import Json, RealDictCursor

from app.core.config import settings
from app.services.infra.database import get_master_db_connection

logger = logging.getLogger(__name__)

KIND_SHOPIFY = "shopify"
KIND_HTTP_API = "http_api"
# The merchant's own website, read by the crawler. Products written before this
# became a registered source carry source_kind 'crawl' with a null source_ref;
# see reap_unregistered_crawl_rows for how those are retired.
KIND_CRAWL = "crawl"


def _fernet() -> Fernet:
    return Fernet(settings.SOURCE_CREDENTIALS_KEY.encode())


def encrypt_credentials(creds: dict) -> str:
    """Encrypt a credentials object.

    A dict rather than a bare string, so the same envelope holds a Shopify
    access token or an API key plus client secret without a schema change.
    sort_keys makes the plaintext stable, which matters for tests but not
    for security -- Fernet's IV already makes ciphertext non-deterministic.
    """
    return _fernet().encrypt(json.dumps(creds, sort_keys=True).encode()).decode()


def decrypt_credentials(blob: str) -> dict:
    """Decrypt a credentials blob. Raises InvalidToken if tampered or wrongly keyed."""
    return json.loads(_fernet().decrypt(blob.encode()).decode())


def ensure_table() -> None:
    conn = get_master_db_connection()
    try:
        with conn.cursor() as cur:
            cur.execute("""
                CREATE TABLE IF NOT EXISTS product_sources (
                    id                    SERIAL PRIMARY KEY,
                    tenant_id             TEXT NOT NULL,
                    kind                  TEXT NOT NULL,
                    external_ref          TEXT NOT NULL,
                    config                JSONB NOT NULL DEFAULT '{}',
                    credentials_encrypted TEXT NOT NULL,
                    status                TEXT NOT NULL DEFAULT 'active',
                    connected_at          TIMESTAMPTZ NOT NULL DEFAULT now(),
                    disconnected_at       TIMESTAMPTZ,
                    last_synced_at        TIMESTAMPTZ,
                    UNIQUE (kind, external_ref)
                )
            """)
            cur.execute(
                "CREATE INDEX IF NOT EXISTS product_sources_tenant_idx "
                "ON product_sources (tenant_id)"
            )
        conn.commit()
    except Exception:
        conn.rollback()
        logger.error("Could not ensure product_sources table", exc_info=True)
        raise
    finally:
        conn.close()


def upsert_source(tenant_id: str, kind: str, external_ref: str,
                   config: dict, credentials: dict) -> None:
    """Insert a source, or update it in place on (kind, external_ref) conflict.

    The ON CONFLICT target is what makes a reinstall update the existing row
    instead of creating a second row under a different tenant.
    """
    ensure_table()
    conn = get_master_db_connection()
    try:
        with conn.cursor() as cur:
            cur.execute("""
                INSERT INTO product_sources
                    (tenant_id, kind, external_ref, config,
                     credentials_encrypted, status, connected_at, disconnected_at)
                VALUES (%s, %s, %s, %s, %s, 'active', now(), NULL)
                ON CONFLICT (kind, external_ref) DO UPDATE SET
                    tenant_id             = EXCLUDED.tenant_id,
                    config                = EXCLUDED.config,
                    credentials_encrypted = EXCLUDED.credentials_encrypted,
                    status                = 'active',
                    connected_at          = now(),
                    disconnected_at       = NULL
            """, (tenant_id, kind, external_ref, Json(config),
                  encrypt_credentials(credentials)))
        conn.commit()
        logger.info("Stored %s source %s for tenant %s", kind, external_ref, tenant_id)
    except Exception:
        conn.rollback()
        logger.error("Could not store %s source for tenant %s", kind, tenant_id,
                     exc_info=True)
        raise
    finally:
        conn.close()


def get_sources(tenant_id: str) -> list[dict]:
    """Every active source for a tenant. Never includes credentials."""
    ensure_table()
    conn = get_master_db_connection()
    try:
        with conn.cursor(cursor_factory=RealDictCursor) as cur:
            cur.execute("""
                SELECT kind, external_ref, config, status,
                       connected_at, last_synced_at
                FROM product_sources
                WHERE tenant_id = %s AND status = 'active'
                ORDER BY connected_at DESC
            """, (tenant_id,))
            return [dict(r) for r in cur.fetchall()]
    finally:
        conn.close()


def get_credentials(kind: str, external_ref: str) -> dict | None:
    """Decrypted credentials for one active source, or None if not found.

    Callers must not log the return value.
    """
    ensure_table()
    conn = get_master_db_connection()
    try:
        with conn.cursor() as cur:
            cur.execute(
                "SELECT credentials_encrypted FROM product_sources "
                "WHERE kind = %s AND external_ref = %s AND status = 'active'",
                (kind, external_ref),
            )
            row = cur.fetchone()
            return decrypt_credentials(row[0]) if row else None
    finally:
        conn.close()


def update_source_config(kind: str, external_ref: str, config: dict) -> None:
    """Persist a config change (e.g. first-sync field inference) without
    touching credentials or status."""
    conn = get_master_db_connection()
    try:
        with conn.cursor() as cur:
            cur.execute(
                "UPDATE product_sources SET config = %s "
                "WHERE kind = %s AND external_ref = %s",
                (Json(config), kind, external_ref),
            )
        conn.commit()
    except Exception:
        conn.rollback()
        logger.error("Could not update config for %s:%s", kind, external_ref,
                     exc_info=True)
        raise
    finally:
        conn.close()


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 disconnect_all_sources(tenant_id: str) -> int:
    """Revoke every active source a tenant has connected. Pairs with wiping
    the whole catalogue: without this, a source stays connected after its
    products are deleted, and the next build silently re-crawls it and
    repopulates exactly what was just removed."""
    ensure_table()
    conn = get_master_db_connection()
    try:
        with conn.cursor() as cur:
            cur.execute(
                "UPDATE product_sources SET status = 'revoked', disconnected_at = now() "
                "WHERE tenant_id = %s AND status = 'active'",
                (tenant_id,),
            )
            count = cur.rowcount
        conn.commit()
        logger.info("Disconnected %d source(s) for tenant %s", count, tenant_id)
        return count
    except Exception:
        conn.rollback()
        logger.error("Could not disconnect sources for tenant %s", tenant_id, 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()
