"""Aproveitar o corpus que já está guardado, sem um único pedido ao IRN.

O `dre_publications` tem 2252 publicações do IRN com corpo, e 858 delas contêm
blocos `NIF/NIPC:`. Nunca foram exploradas porque o parser antigo exigia essa
linha colada ao fim de linha e a fonte traz um espaço à esquerda — devolvia 15
pessoas no total. Com o parser novo, as mesmas publicações dão ~2180 arestas.

Correr isto antes de qualquer recolha é deliberado: valida a extracção contra
dados reais de produção em minutos, e se estiver errada descobre-se agora e não
depois de 780 mil actos.

O `dedup_hash` usa a mesma receita do varrimento nacional, portanto quando o mesmo
acto voltar por essa via aterra nesta linha em vez de criar outra.
"""
import logging
from typing import Any

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

from app.services.registry_act_parser import (
    REGISTRY_PARSE_VERSION,
    clean_entity_name,
    normalize_name,
    parse_act,
)
from app.services.registry_graph import refresh_degrees
from app.services.registry_ingest import apply_parse, upsert_act

logger = logging.getLogger(__name__)

# O `mj_irn_sync` grava `summary = (body or "")[:4000]`. Um corpo que chegue a
# esse limite pode estar cortado a meio de uma cap table, e uma cap table cortada
# é pior do que nenhuma: parece completa. Estes ficam marcados para refetch em
# vez de parseados.
TRUNCATION_LIMIT = 3990


async def backfill_from_publications(
    session: AsyncSession, limit: int | None = None, batch: int = 100
) -> dict[str, int]:
    rows = (
        await session.execute(
            text(
                """
                SELECT p.id::text AS pub_id,
                       p.raw_json->>'Id' AS irn_id,
                       c.nif,
                       p.type,
                       p.date,
                       p.summary,
                       length(p.summary) AS len
                  FROM dre_publications p
                  JOIN companies c ON c.id = p.company_id
                 WHERE p.source = 'mj' AND p.summary IS NOT NULL
                 ORDER BY p.date DESC NULLS LAST
                 """
                + (" LIMIT :lim" if limit else "")
            ),
            {"lim": limit} if limit else {},
        )
    ).mappings().all()

    stats = {
        "publicacoes": len(rows), "actos": 0, "novos": 0, "truncados": 0,
        "quotas": 0, "cargos": 0, "conjuges": 0, "sem_nif": 0,
    }
    touched: set[str] = set()

    for i, row in enumerate(rows, start=1):
        nipc = (row["nif"] or "").strip()
        if len(nipc) != 9:
            continue
        body = row["summary"]
        truncated = (row["len"] or 0) >= TRUNCATION_LIMIT

        act_id, is_new = await upsert_act(
            session,
            nipc=nipc,
            irn_id=row["irn_id"],
            act_type=row["type"],
            act_date=row["date"],
            body=None if truncated else body,
            body_source="dre_publications",
            body_truncated=truncated,
            # o que já monitorizamos é o que interessa primeiro na fila de recolha
            priority=0,
        )
        if not act_id:
            continue
        stats["actos"] += 1
        stats["novos"] += int(is_new)

        if truncated:
            # Sem corpo fiável: fica na fila para ser buscado inteiro ao IRN.
            stats["truncados"] += 1
            await session.execute(
                text(
                    """
                    UPDATE registry_acts
                       SET content_status = 'pending', wants_content = TRUE, priority = 0,
                           parse_error = 'corpo truncado aos 4000 chars pelo mj_irn_sync'
                     WHERE id = CAST(:aid AS UUID)
                    """
                ),
                {"aid": act_id},
            )
            continue

        parsed = parse_act(body, row["type"])
        counts = await apply_parse(
            session, act_id=act_id, nipc=nipc, parsed=parsed, act_date=row["date"]
        )
        stats["quotas"] += counts["quotas"]
        stats["cargos"] += counts["officers"]
        stats["conjuges"] += counts["spouses"]
        if parsed.get("quotas") and not parsed.get("captable_linkable"):
            stats["sem_nif"] += 1
        touched.add(nipc)

        if i % batch == 0:
            await session.commit()
            logger.info("backfill registo: %d/%d publicações", i, len(rows))

    await session.commit()
    await refresh_degrees(session)
    await session.commit()
    stats["entidades"] = len(touched)
    logger.info("backfill registo terminado: %s", stats)
    return stats


async def reparse_stored_acts(
    session: AsyncSession, limit: int | None = None, batch: int = 100
) -> dict[str, int]:
    """Reprocessa os actos cujo corpo já está guardado, sem tocar na rede.

    É para isto que o `registry_acts.body` existe. Quando a extracção muda de
    forma que valha a pena — o `REGISTRY_PARSE_VERSION` sobe — o corpus inteiro
    volta a passar pelo parser e as arestas são reescritas por `source_act_id`,
    que é a operação que o `apply_parse` já garante ser idempotente.

    Corre por ordem de `parse_version` para os actos nunca parseados virem
    primeiro, e commita a cada lote: são milhares de actos e uma transacção única
    seguraria as escritas do resto do sistema durante minutos.
    """
    rows = (
        await session.execute(
            text(
                """
                SELECT id::text, nipc, act_type, act_date, body
                  FROM registry_acts
                 WHERE body IS NOT NULL
                   AND COALESCE(parse_version, 0) < :pv
                 ORDER BY parse_version NULLS FIRST, act_date DESC NULLS LAST
                """
                + (" LIMIT :lim" if limit else "")
            ),
            {"pv": REGISTRY_PARSE_VERSION, **({"lim": limit} if limit else {})},
        )
    ).mappings().all()

    stats = {"candidatos": len(rows), "reparseados": 0, "quotas": 0, "cargos": 0}
    touched: set[str] = set()
    for i, row in enumerate(rows, start=1):
        parsed = parse_act(row["body"], row["act_type"])
        counts = await apply_parse(
            session,
            act_id=row["id"],
            nipc=row["nipc"],
            parsed=parsed,
            act_date=row["act_date"],
        )
        stats["reparseados"] += 1
        stats["quotas"] += counts["quotas"]
        stats["cargos"] += counts["officers"]
        touched.add(row["nipc"])
        if i % batch == 0:
            await session.commit()
            logger.info("reparse: %d/%d actos", i, len(rows))

    await session.commit()
    await refresh_degrees(session)
    await session.commit()
    stats["entidades"] = len(touched)
    logger.info("reparse terminado: %s", stats)
    return stats


async def clean_entity_names(session: AsyncSession) -> dict[str, int]:
    """Limpa nomes já guardados que trazem colado o que não é o nome.

    Tem de ser uma passagem própria e não um reparse: o `upsert_entity` só
    substitui o nome **quando o novo é mais longo** — regra que existe por bom
    motivo, porque a fonte alterna `Crescentimob` com `CRESCENTIMOB -
    IMOBILIÁRIA, LDA` e a forma completa é a útil. Consequência: um nome
    encurtado nunca ganharia ao que lá está, e reparsear os 60 mil actos não
    corrigia uma única linha.

    Depois disto o parser já entrega nomes limpos, portanto o problema não
    regressa. Todas as entidades afectadas têm NIF, logo a `entity_key` vem do
    NIF e nenhuma muda de identidade — ninguém fica órfão das suas arestas.
    """
    rows = (
        await session.execute(
            text(
                """
                SELECT id::text, name FROM registry_entities
                 WHERE name ~* ' represent[ao]d[ao] por '
                    OR name ~ '^\\s*"'
                    OR name ~* ',\\s*n\\.?º\\s*[0-9]+\\s*$'
                """
            )
        )
    ).all()

    stats = {"candidatos": len(rows), "limpos": 0}
    for entity_id, name in rows:
        novo = clean_entity_name(name)
        if not novo or novo == name:
            continue
        await session.execute(
            text(
                "UPDATE registry_entities SET name = :n, name_norm = :nn,"
                " updated_at = now() WHERE id = CAST(:id AS UUID)"
            ),
            {"n": novo, "nn": normalize_name(novo), "id": entity_id},
        )
        stats["limpos"] += 1
    await session.commit()
    logger.info("limpeza de nomes: %s", stats)
    return stats


async def summary(session: AsyncSession) -> dict[str, Any]:
    row = (
        await session.execute(
            text(
                """
                SELECT
                  (SELECT count(*) FROM registry_acts) AS actos,
                  (SELECT count(*) FROM registry_acts WHERE body IS NOT NULL) AS com_corpo,
                  (SELECT count(*) FROM registry_entities) AS entidades,
                  (SELECT count(*) FROM registry_entities WHERE kind = 'person') AS pessoas,
                  (SELECT count(*) FROM registry_entities WHERE kind = 'company') AS empresas,
                  (SELECT count(*) FROM registry_edges) AS arestas,
                  (SELECT count(*) FROM registry_edges WHERE is_current) AS arestas_actuais,
                  (SELECT count(*) FROM registry_edges
                     WHERE edge_type = 'holds_quota' AND is_current) AS quotas_actuais,
                  (SELECT count(*) FROM registry_edges WHERE edge_type = 'spouse_of') AS conjuges,
                  (SELECT count(*) FROM registry_entities
                     WHERE captable_as_of IS NOT NULL) AS com_captable,
                  (SELECT count(*) FROM registry_entities WHERE captable_complete) AS captable_completa
                """
            )
        )
    ).mappings().first()
    return dict(row) if row else {}
