"""Arquivo nacional da Distribuição do CITIUS: ingestão e cruzamento.

A varredura sempre descarregou o país inteiro — a consulta é feita sem filtro de
parte — mas guardava só o que dava match e deitava o resto fora em memória. Como
a fonte só serve ~180 dias, isso era informação a perder-se todos os dias.

Aqui a ordem inverte-se: **primeiro guarda-se, depois cruza-se**. É a mesma forma
do `cire_announcement_ingest` + `insolvency_matching`, e a razão é a mesma: um
cruzamento que vive à parte pode voltar a correr. Quando entra um cliente novo,
ou quando a regra melhora, reaplica-se ao arquivo — em vez de valer só daí para
a frente, que é o que acontece hoje com uma correspondência falhada.

**Dois critérios, com pesos diferentes.** O NIF é o bom: o `insolvency_matching`
já escolheu só esse porque um falso positivo custa caro — alguém vai agir
convencido de que um cliente está em tribunal. Mas a Distribuição não publica
NIF, só nomes. O que o arquivo permite é ir buscá-lo ao grafo do registo
comercial, que tem 1,27 milhões de entidades com nome e NIF. O que não resolver
fica com o critério antigo, o distintivo, marcado como tal no `match_kind` para
se poder rever depois só o que é frágil.
"""
import json
import logging
from datetime import date
from typing import Any

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

from app.config import settings
from app.scrapers.distribuicao import build_distintivos_map, match_row_to_companies
from app.services.dedup import filing_hash
from app.services.ingest import upsert_process
from app.services.registry_act_parser import normalize_name, parse_money

logger = logging.getLogger(__name__)

# Quantos processos por lote no cruzamento. O trabalho é de strings e de dois
# UPDATEs; o que interessa é não segurar uma transacção durante um backfill de
# 128 mil linhas.
LOTE = 500


async def upsert_filing(session: AsyncSession, row: dict[str, Any]) -> tuple[str, bool]:
    """Grava uma linha da Distribuição e as suas partes. Devolve (id, é_nova).

    O `raw` **não leva as partes** — elas têm tabela própria. É a lição que o
    `processes.raw` deu ao duplicá-las e custar 7,7 kB por linha; a esta escala,
    ~230 mil processos por ano a 6 partes cada, seriam ~80 MB/ano por nada.
    """
    data_dist = row.get("data_distribuicao_d")
    h = filing_hash(row.get("process_number"), row.get("tribunal"), data_dist)
    valor, _moeda = parse_money(row.get("valor"))
    raw = {k: v for k, v in row.items() if k not in ("parties", "data_distribuicao_d", "data_entrada_d")}

    result = await session.execute(
        text(
            """
            INSERT INTO court_filings (
                process_number, tribunal, unorganica, especie, valor_raw, valor,
                data_distribuicao, data_entrada, observacoes, raw, dedup_hash
            ) VALUES (
                :process_number, :tribunal, :unorganica, :especie, :valor_raw, :valor,
                :data_distribuicao, :data_entrada, :observacoes, CAST(:raw AS JSONB), :hash
            )
            ON CONFLICT (dedup_hash) DO UPDATE SET
                last_seen_at = now(),
                -- Só se enriquece: uma re-observação não apaga o que já se sabia.
                unorganica = COALESCE(EXCLUDED.unorganica, court_filings.unorganica),
                especie = COALESCE(EXCLUDED.especie, court_filings.especie),
                valor = COALESCE(EXCLUDED.valor, court_filings.valor),
                observacoes = COALESCE(EXCLUDED.observacoes, court_filings.observacoes)
            RETURNING id::text, (xmax = 0) AS is_new
            """
        ),
        {
            "process_number": (row.get("process_number") or "").strip(),
            "tribunal": (row.get("tribunal") or "").strip(),
            "unorganica": row.get("unorganica"),
            "especie": row.get("especie"),
            "valor_raw": row.get("valor"),
            "valor": valor,
            "data_distribuicao": data_dist,
            "data_entrada": row.get("data_entrada_d"),
            "observacoes": row.get("observacoes"),
            "raw": json.dumps(raw, ensure_ascii=False, default=str),
            "hash": h,
        },
    )
    filing_id, is_new = result.first()

    if is_new:
        for p in row.get("parties") or []:
            name = (p.get("name") or "").strip()
            if not name:
                continue
            await session.execute(
                text(
                    """
                    INSERT INTO court_filing_parties (filing_id, name, name_norm, role)
                    VALUES (:fid, :name, :norm, :role)
                    """
                ),
                {
                    "fid": filing_id,
                    "name": name[:500],
                    # `normalize_name` e não a normalização do scraper: é esta a
                    # que produziu o `registry_entities.name_norm`, e as duas não
                    # dão o mesmo resultado (uma tira a pontuação, a outra não).
                    "norm": normalize_name(name)[:500],
                    "role": (p.get("role") or "").strip()[:100] or None,
                },
            )
    return str(filing_id), bool(is_new)


async def resolve_party_nifs(session: AsyncSession, filing_ids: list[str]) -> int:
    """Dá NIF às partes que o grafo do registo consegue identificar sem dúvida.

    **Só entidades colectivas.** Uma pessoa singular com nome raro podia ser
    única no registo e ainda assim não ser a pessoa que está em tribunal — o
    nome não chega para identificar gente, e é a mesma razão pela qual o
    `is_company_party` do scraper exige forma societária.

    **E só quando não há ambiguidade.** Se dois NIPCs partilham o nome
    normalizado, não há como escolher e fica por resolver: preferimos falhar um
    match a inventar um.
    """
    result = await session.execute(
        text(
            """
            WITH alvo AS (
                SELECT DISTINCT p.name_norm
                  FROM court_filing_parties p
                 WHERE p.nif IS NULL
                   AND p.name_norm <> ''
                   AND p.filing_id = ANY(CAST(:ids AS UUID[]))
            ),
            resolvidas AS (
                SELECT a.name_norm, min(e.nif) AS nif
                  FROM alvo a
                  JOIN registry_entities e ON e.name_norm = a.name_norm
                 WHERE e.nif IS NOT NULL
                   AND e.merged_into IS NULL
                   AND e.kind IN ('company', 'public', 'foreign')
                 GROUP BY a.name_norm
                HAVING count(*) = 1
            )
            UPDATE court_filing_parties p
               SET nif = r.nif, nif_source = 'grafo'
              FROM resolvidas r
             WHERE p.name_norm = r.name_norm
               AND p.nif IS NULL
               AND p.filing_id = ANY(CAST(:ids AS UUID[]))
            """
        ),
        {"ids": filing_ids},
    )
    return (result.rowcount or 0) + await _resolver_truncadas(session, filing_ids)


# A Distribuição corta o nome das partes aos 80 caracteres. Uma firma com nome
# comprido chega partida a meio de uma palavra — a sucursal da IGPS PROTEK vem
# como "…SOCIETÉ A RESPONSABIL" — e a igualdade exacta contra o registo nunca
# casa. São 2.206 partes no arquivo, e entre elas ficava a própria empresa do
# grupo: os processos dela não eram atribuídos a ninguém.
_TRUNCADO = 79


async def _resolver_truncadas(session: AsyncSession, filing_ids: list[str]) -> int:
    """Para os nomes cortados pela fonte, o prefixo vale — se for de um só.

    Mesma exigência do resto: uma entidade no registo, ou não se resolve. Um
    prefixo que apanhe duas empresas do mesmo grupo é ambíguo por definição, e
    entre inventar e falhar continua a preferir-se falhar.
    """
    result = await session.execute(
        text(
            """
            WITH truncadas AS (
                SELECT DISTINCT p.name_norm
                  FROM court_filing_parties p
                 WHERE p.nif IS NULL
                   AND length(p.name) >= :corte
                   AND p.name_norm <> ''
                   AND p.filing_id = ANY(CAST(:ids AS UUID[]))
            ),
            resolvidas AS (
                SELECT t.name_norm, min(e.nif) AS nif
                  FROM truncadas t
                  JOIN registry_entities e ON e.name_norm LIKE t.name_norm || '%'
                 WHERE e.nif IS NOT NULL
                   AND e.merged_into IS NULL
                   AND e.kind IN ('company', 'public', 'foreign')
                 GROUP BY t.name_norm
                HAVING count(*) = 1
            )
            UPDATE court_filing_parties p
               SET nif = r.nif, nif_source = 'grafo-prefixo'
              FROM resolvidas r
             WHERE p.name_norm = r.name_norm
               AND p.nif IS NULL
               AND p.filing_id = ANY(CAST(:ids AS UUID[]))
            """
        ),
        {"ids": filing_ids, "corte": _TRUNCADO},
    )
    return result.rowcount or 0


async def _hits_por_nif(session: AsyncSession, filing_ids: list[str]) -> list[str]:
    """Acertos pelo NIF resolvido. Devolve os ids dos acertos novos."""
    rows = (
        await session.execute(
            text(
                """
                INSERT INTO court_filing_hits (
                    filing_id, party_id, company_id, watchlist_id,
                    matched_nif, matched_name, matched_role, match_kind, score
                )
                SELECT p.filing_id, p.id, c.id, w.id, p.nif,
                       COALESCE(c.legal_name, w.name, p.name), p.role, 'nif', 1.0
                  FROM court_filing_parties p
                  LEFT JOIN companies c ON c.nif = p.nif
                  LEFT JOIN insolvency_watchlist w ON w.nif = p.nif AND w.active
                 WHERE p.filing_id = ANY(CAST(:ids AS UUID[]))
                   AND p.nif IS NOT NULL
                   AND (c.id IS NOT NULL OR w.id IS NOT NULL)
                ON CONFLICT DO NOTHING
                RETURNING id::text
                """
            ),
            {"ids": filing_ids},
        )
    ).all()
    return [r[0] for r in rows]


async def _hits_por_nome(
    session: AsyncSession, filing_ids: list[str], distintivos: dict[str, dict[str, Any]]
) -> list[str]:
    """Acertos pelo nome, para o que o NIF não resolveu — com duas guardas.

    O `match_row_to_companies` compara o **prefixo** de duas palavras. Sozinho,
    isso atribuiu 316 processos e, das 119 vezes em que houve NIF para conferir,
    **errou 119** — nem uma acertou. "Centro de Jardinagem Vila Jasmim" ia parar
    ao Centro de Cultura e Desporto de Ronfe, e os processos da Prosegur
    Logística, da Prosegur Alarmes e até da Prosegur *Seguros* eram todos
    carimbados à Prosegur - Companhia de Segurança.

    Por isso o prefixo passa a ser só o primeiro crivo, seguido de dois:

    1. **O NIF manda, quando existe.** Se a parte foi resolvida para um NIF e ele
       não é o da empresa, a correspondência está provada errada. Não é uma
       suspeita, é uma contradição.
    2. **O nome inteiro tem de parecer-se**, não só o princípio — o mesmo
       `token_set_ratio` a `FUZZY_MATCH_THRESHOLD` que o caminho por empresa já
       usava. É o que separa "Prosegur - Companhia de Segurança, Lda." de
       "Prosegur - Logística e Tratamento de Valores".
    """
    partes = (
        await session.execute(
            text(
                """
                SELECT p.filing_id::text, p.id::text, p.name, p.role, p.nif
                  FROM court_filing_parties p
                 WHERE p.filing_id = ANY(CAST(:ids AS UUID[]))
                 ORDER BY p.filing_id
                """
            ),
            {"ids": filing_ids},
        )
    ).all()

    por_processo: dict[str, list[dict[str, Any]]] = {}
    for filing_id, party_id, name, role, nif in partes:
        por_processo.setdefault(filing_id, []).append(
            {"id": party_id, "name": name, "role": role, "nif": nif}
        )

    novos: list[str] = []
    for filing_id, parties in por_processo.items():
        for company, party in match_row_to_companies({"parties": parties}, distintivos):
            nif_parte = (party.get("nif") or "").strip()
            nif_empresa = (company.get("nif") or "").strip()
            if nif_parte and nif_parte != nif_empresa:
                continue
            score = (
                fuzz.token_set_ratio(
                    (party.get("name") or "").lower(),
                    (company.get("legal_name") or "").lower(),
                )
                / 100.0
            )
            if score < settings.FUZZY_MATCH_THRESHOLD:
                continue
            row = (
                await session.execute(
                    text(
                        """
                        INSERT INTO court_filing_hits (
                            filing_id, party_id, company_id,
                            matched_name, matched_role, match_kind, score
                        ) VALUES (
                            CAST(:fid AS UUID), CAST(:pid AS UUID), CAST(:cid AS UUID),
                            :name, :role, 'nome', :score
                        )
                        ON CONFLICT DO NOTHING
                        RETURNING id::text
                        """
                    ),
                    {
                        "fid": filing_id,
                        "pid": party["id"],
                        "cid": company["id"],
                        "name": company.get("legal_name") or party.get("name"),
                        "role": party.get("role"),
                        "score": round(score, 3),
                    },
                )
            ).first()
            if row:
                novos.append(row[0])
    return novos


async def _processos_dos_hits(session: AsyncSession, hit_ids: list[str]) -> int:
    """Escreve na `processes` o que os acertos novos trouxeram.

    A `processes` continua a ser exactamente o que era — a vista das empresas
    monitorizadas — e continua a ser escrita pelo `upsert_process`, com a sua
    lista negra de falsos positivos e os seus eventos. Só mudou de onde vem a
    linha: do arquivo, em vez de directamente do scraper.
    """
    if not hit_ids:
        return 0
    linhas = (
        await session.execute(
            text(
                """
                SELECT h.company_id::text, f.id::text, f.process_number, f.tribunal,
                       f.unorganica, f.especie, h.matched_role, f.data_distribuicao,
                       f.data_entrada, f.raw
                  FROM court_filing_hits h
                  JOIN court_filings f ON f.id = h.filing_id
                 WHERE h.id = ANY(CAST(:ids AS UUID[]))
                   AND h.company_id IS NOT NULL
                """
            ),
            {"ids": hit_ids},
        )
    ).all()

    escritos = 0
    for (company_id, filing_id, numero, tribunal, unorganica, especie,
         papel, data_dist, data_ent, raw) in linhas:
        partes = (
            await session.execute(
                text(
                    "SELECT name, role, nif FROM court_filing_parties"
                    " WHERE filing_id = CAST(:fid AS UUID)"
                ),
                {"fid": filing_id},
            )
        ).all()
        entry = dict(raw or {})
        entry["parties"] = [{"name": n, "role": r, "nif": nif} for n, r, nif in partes]
        _pid, is_new = await upsert_process(
            session,
            {
                "company_id": company_id,
                "source": "distribuicao",
                "process_number": numero,
                "tribunal": tribunal,
                "juizo": unorganica,
                "species": especie,
                "role_in_process": papel,
                "date_filed": data_dist or data_ent,
                "parties": entry["parties"],
                "raw_json": json.dumps(entry, ensure_ascii=False, default=str),
                "raw_html": None,
            },
        )
        escritos += 1 if is_new else 0
    return escritos


async def carregar_distintivos(session: AsyncSession) -> dict[str, dict[str, Any]]:
    """O mapa de distintivos das empresas activas.

    Fica à parte para se construir **uma vez por corrida**, não uma vez por lote:
    o `build_distintivos_map` avisa a cada nome ambíguo que encontra ("freguesia
    de", "associação desportiva") e, chamado por tribunal e por dia, enterrava o
    log em milhares de avisos iguais.
    """
    empresas = (
        await session.execute(
            text("SELECT id::text, legal_name, nif FROM companies WHERE active")
        )
    ).mappings().all()
    return build_distintivos_map([dict(c) for c in empresas])


async def match_filings(
    session: AsyncSession,
    filing_ids: list[str],
    distintivos: dict[str, dict[str, Any]] | None = None,
) -> dict[str, int]:
    """Cruza um conjunto de processos do arquivo. Idempotente."""
    if not filing_ids:
        return {"nifs": 0, "hits_nif": 0, "hits_nome": 0, "processos": 0}

    if distintivos is None:
        distintivos = await carregar_distintivos(session)

    nifs = await resolve_party_nifs(session, filing_ids)
    por_nif = await _hits_por_nif(session, filing_ids)
    por_nome = await _hits_por_nome(session, filing_ids, distintivos)
    processos = await _processos_dos_hits(session, por_nif + por_nome)
    return {
        "nifs": nifs,
        "hits_nif": len(por_nif),
        "hits_nome": len(por_nome),
        "processos": processos,
    }


async def rematch_all(
    session: AsyncSession,
    *,
    since: date | None = None,
    limit: int | None = None,
) -> dict[str, int]:
    """Reaplica o cruzamento ao arquivo inteiro, por lotes.

    É a razão de existir do arquivo. Corre-se quando entra uma empresa nova — e
    aí o alcance é tudo o que estiver guardado, não os 180 dias que o CITIUS
    ainda serve — ou quando a regra de correspondência muda e o que falhou em
    silêncio merece uma segunda leitura.
    """
    total = {"processos_vistos": 0, "nifs": 0, "hits_nif": 0, "hits_nome": 0, "processos": 0}
    distintivos = await carregar_distintivos(session)
    offset = 0
    while True:
        ids = [
            r[0]
            for r in (
                await session.execute(
                    text(
                        """
                        SELECT id::text FROM court_filings
                         WHERE (CAST(:desde AS DATE) IS NULL
                                OR data_distribuicao >= CAST(:desde AS DATE))
                         ORDER BY data_distribuicao DESC NULLS LAST, id
                         LIMIT :lote OFFSET :off
                        """
                    ),
                    {"desde": since, "lote": LOTE, "off": offset},
                )
            ).all()
        ]
        if not ids:
            break
        parcial = await match_filings(session, ids, distintivos)
        for k, v in parcial.items():
            total[k] += v
        total["processos_vistos"] += len(ids)
        await session.commit()
        offset += LOTE
        if offset % (LOTE * 10) == 0:
            logger.info("rematch: %d processos do arquivo revistos", total["processos_vistos"])
        if limit and total["processos_vistos"] >= limit:
            break
    logger.info("rematch terminado: %s", total)
    return total
