import json
import logging
from typing import Any

from sqlalchemy import text
from sqlalchemy.ext.asyncio import AsyncSession

from app.services.alerts import trigger_alerts
from app.services.coverage import update_coverage
from app.services.dedup import dedup_hash
from app.services.risk import update_company_risk

logger = logging.getLogger(__name__)


def _norm(value: Any) -> str | None:
    """Compara pelo conteúdo, não pela formatação.

    A fonte alterna `"1.ª Secção"` com `"1.ª Secção "` e datas chegam ora como
    `date` ora como string. Sem isto, espaço a mais é uma alteração do processo.
    """
    if value is None or value == "":
        return None
    return str(value).strip() or None


# Os campos cuja mudança é uma alteração do processo. O resto do `raw` muda a
# cada recolha (ordem de chaves, carimbos da fonte) e não diz nada a ninguém.
TRACKED_FIELDS = ("tribunal", "juizo", "species", "role_in_process", "date_filed")


async def upsert_process(session: AsyncSession, row: dict[str, Any]) -> tuple[str, bool]:
    """Upsert a process row. Returns (process_id, is_new).

    Grava um evento `discovered` na primeira vez, e depois **só quando algo muda**
    — tribunal, juízo, espécie, qualidade ou data de distribuição.

    Antes gravava um `updated` em cada re-observação, porque o `ON CONFLICT DO
    UPDATE SET last_seen_at = now()` faz sempre `xmax <> 0`. Um processo apanhado
    todos os dias juntava um evento por dia com o mesmo conteúdo: o 1846/26.5T8STS
    tinha 120 eventos, 118 deles idênticos, e a timeline deixou de ser o histórico
    do processo para ser o registo das corridas do scraper. Ter sido visto continua
    a ficar guardado — é o `last_seen_at`, que não precisa de uma linha por vez.

    If the dedup_hash is in the false_positive_reports blacklist (user flagged
    a prior ingestion of this same (company,process,tribunal,date) as a
    mismatch), silently skip insert and return ('', False)."""
    h = dedup_hash(
        str(row["company_id"]), row["process_number"], row["tribunal"], row["date_filed"]
    )
    blacklisted = (
        await session.execute(
            text("SELECT 1 FROM false_positive_reports WHERE dedup_hash = :h"),
            {"h": h},
        )
    ).first()
    if blacklisted:
        logger.debug("skipping blacklisted process %s/%s", row["process_number"], h[:8])
        return "", False
    # O estado anterior tem de ser lido **antes** do upsert: depois já não há
    # forma de saber se o que lá está agora é novo ou é o mesmo de ontem.
    before = (
        await session.execute(
            text(
                "SELECT tribunal, juizo, species, role_in_process, date_filed"
                "  FROM processes WHERE dedup_hash = :h"
            ),
            {"h": h},
        )
    ).mappings().first()

    result = await session.execute(
        text(
            """
            INSERT INTO processes (company_id, source, process_number, tribunal, juizo,
                                   species, role_in_process, date_filed, dedup_hash, raw, raw_html)
            VALUES (:company_id, :source, :process_number, :tribunal, :juizo,
                    :species, :role_in_process, :date_filed, :dedup_hash, CAST(:raw AS JSONB), :raw_html)
            ON CONFLICT (dedup_hash)
              DO UPDATE SET last_seen_at = now(),
                            juizo = EXCLUDED.juizo,
                            species = EXCLUDED.species,
                            role_in_process = EXCLUDED.role_in_process,
                            raw = EXCLUDED.raw
            RETURNING id, (xmax = 0) AS is_new
            """
        ),
        {
            "company_id": row["company_id"],
            "source": row["source"],
            "process_number": row["process_number"],
            "tribunal": row["tribunal"],
            "juizo": row.get("juizo"),
            "species": row.get("species"),
            "role_in_process": row.get("role_in_process"),
            "date_filed": row["date_filed"],
            "dedup_hash": h,
            "raw": row.get("raw_json", "{}"),
            "raw_html": row.get("raw_html"),
        },
    )
    record = result.first()
    process_id, is_new = record
    if is_new:
        await session.execute(
            text(
                """
                INSERT INTO process_events (process_id, event_type, payload)
                VALUES (:pid, 'discovered', CAST(:pl AS JSONB))
                """
            ),
            {"pid": process_id, "pl": row.get("raw_json", "{}")},
        )
    else:
        # A actualização só é evento se mudou alguma coisa, e o que se guarda é
        # o que mudou — não o `raw` inteiro outra vez. Os 16.736 eventos que
        # esta comparação evita ocupavam 107 MB a dizer "voltei a ver o mesmo".
        diff = {
            field: {"antes": _norm(before[field]), "agora": _norm(row.get(field))}
            for field in TRACKED_FIELDS
            if before is not None and _norm(before[field]) != _norm(row.get(field))
        }
        if diff:
            await session.execute(
                text(
                    """
                    INSERT INTO process_events (process_id, event_type, payload)
                    VALUES (:pid, 'updated', CAST(:pl AS JSONB))
                    """
                ),
                {"pid": process_id, "pl": json.dumps(diff, default=str)},
            )
    if is_new:
        # Insert process parties (for search + intelligence features)
        parties = row.get("parties") or row.get("_parties") or []
        if not parties:
            import json as _json
            try:
                raw_obj = _json.loads(row.get("raw_json") or "{}")
                parties = raw_obj.get("parties") or raw_obj.get("_parties") or []
            except Exception:
                parties = []
        for p in parties:
            name = (p.get("name") or "").strip()
            if not name:
                continue
            nif_val = (p.get("nif") or "").strip() or None
            if nif_val and (not nif_val.isdigit() or len(nif_val) != 9):
                nif_val = None
            await session.execute(
                text(
                    """
                    INSERT INTO process_parties (process_id, name, nif, role)
                    VALUES (:pid, :name, :nif, :role)
                    """
                ),
                {
                    "pid": process_id,
                    "name": name[:500],
                    "nif": nif_val,
                    "role": (p.get("role") or "")[:100] or None,
                },
            )

        # Load enriched row (needs company name + NIF for alert payload).
        meta = (
            await session.execute(
                text(
                    """
                    SELECT c.legal_name, c.nif, c.monitoring_type, p.first_seen_at
                    FROM processes p JOIN companies c ON c.id = p.company_id
                    WHERE p.id = :id
                    """
                ),
                {"id": process_id},
            )
        ).first()
        enriched = {
            **row,
            "company_name": meta[0] if meta else None,
            "company_nif": meta[1] if meta else None,
            "first_seen_at": meta[3] if meta else None,
        }
        # Alerts ONLY for internal companies (competitors/analysis are silent).
        if meta and meta[2] == "internal":
            try:
                await trigger_alerts(session, str(process_id), enriched)
            except Exception as e:
                logger.warning("trigger_alerts failed for process %s: %s", process_id, e)

        # Refresh coverage + risk for this company (cheap, small window)
        try:
            await update_coverage(session, str(row["company_id"]))
            await update_company_risk(session, str(row["company_id"]))
        except Exception as e:
            logger.warning("coverage/risk update failed: %s", e)
    return str(process_id), bool(is_new)
