"""Recolha do registo comercial: histórico por NIPC e expansão pelo grafo.

Duas vias, com objectivos diferentes e que se complementam:

* **Por NIPC** (aqui). O `publications_for_nif` do IRN é ilimitado no tempo,
  portanto dá o **histórico completo** de uma entidade — incluindo a constituição
  de 2008 que nenhum varrimento dos últimos anos encontraria. É a via das
  entidades que nos interessam: as monitorizadas e, em BFS, cada empresa
  descoberta como titular ou participada.
* **Nacional, dia a dia** (`registry_sweep.py`). A única forma de responder à
  pergunta inversa — "em que outras empresas é que esta pessoa é sócia" — porque
  o IRN **não** permite pesquisar pelo NIF de uma pessoa: com o NIF de um
  indivíduo devolve zero, medido.

As pessoas nunca se expandem por NIPC, por essa mesma razão. O BFS segue apenas
sócios que sejam entidades.
"""
import asyncio
import hashlib
import logging
import time
from datetime import date, datetime
from typing import Any

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

from app.config import settings
from app.scrapers.base import RateLimiter
from app.scrapers.mj_irn import IrnVersionChanged, MjIrnClient
from app.services.mj_company_enrichment import html_to_text
from app.services.registry_act_parser import parse_act, wants_content
from app.services.registry_ingest import apply_parse, enqueue_entity, upsert_act

logger = logging.getLogger(__name__)


class RateLimited(RuntimeError):
    """O IRN respondeu 429/403. Para o pool inteiro, não só quem apanhou."""


def _parse_date(value: str | None) -> date | None:
    if not value:
        return None
    try:
        return datetime.strptime(value[:10], "%Y-%m-%d").date()
    except ValueError:
        return None


async def _fetch_and_parse(
    session: AsyncSession,
    client: MjIrnClient,
    limiter: RateLimiter,
    *,
    act_id: str,
    nipc: str,
    pub: dict[str, Any],
) -> dict[str, int] | None:
    """Abre o conteúdo de um acto e grava nós e arestas. 3 pedidos HTTP."""
    await limiter.acquire(3)
    html = await client.content_for(pub)
    body = html_to_text(html) if html else None
    if not body:
        await session.execute(
            text(
                """
                UPDATE registry_acts
                   SET content_status = 'empty', content_attempts = content_attempts + 1,
                       content_fetched_at = now()
                 WHERE id = CAST(:aid AS UUID)
                """
            ),
            {"aid": act_id},
        )
        return None

    act_date = _parse_date(pub.get("Date"))
    await session.execute(
        text(
            """
            UPDATE registry_acts SET
                body = :body,
                body_sha256 = :sha,
                body_source = 'irn',
                body_truncated = FALSE,
                content_status = 'fetched',
                content_fetched_at = now(),
                content_attempts = content_attempts + 1,
                parse_error = NULL
            WHERE id = CAST(:aid AS UUID)
            """
        ),
        {"aid": act_id, "body": body, "sha": hashlib.sha256(body.encode()).hexdigest()},
    )
    parsed = parse_act(body, pub.get("Act_Fact"))
    return await apply_parse(
        session, act_id=act_id, nipc=nipc, parsed=parsed, act_date=act_date
    )


async def sync_nipc(
    session: AsyncSession,
    client: MjIrnClient,
    limiter: RateLimiter,
    nipc: str,
    *,
    priority: int = 0,
    fetch_content: bool = True,
    depth: int = 0,
) -> dict[str, int]:
    """Histórico completo de um NIPC, com conteúdo dos actos que trazem pessoas.

    Só se abre o conteúdo do que ainda não temos: a listagem é uma chamada, o
    conteúdo são três. Prestações de contas — mais de metade do volume de uma
    entidade — nunca se abrem.
    """
    await limiter.acquire()
    pubs = await client.publications_for_nif(nipc)
    return await ingest_publications(
        session, nipc, pubs,
        client=client, limiter=limiter,
        priority=priority, fetch_content=fetch_content,
    )


async def ingest_publications(
    session: AsyncSession,
    nipc: str,
    pubs: list[dict[str, Any]],
    *,
    client: MjIrnClient | None = None,
    limiter: RateLimiter | None = None,
    priority: int = 0,
    fetch_content: bool = True,
) -> dict[str, int]:
    """Grava a listagem de um NIPC. Não vai à rede se `fetch_content` for falso.

    Está separada do `sync_nipc` porque a listagem passou a poder chegar de duas
    origens: o CT, que a vai buscar ele próprio, e o trabalhador remoto, que a
    entrega já buscada. O que é escasso no IRN é o IP, não o processador — mas o
    que escreve na base de dados tem de ser um só, senão duas versões da mesma
    regra divergem em silêncio na primeira alteração. É a mesma razão pela qual o
    parser nunca saiu daqui.
    """
    stats = {"listados": len(pubs), "novos": 0, "conteudos": 0, "quotas": 0, "cargos": 0}
    if not pubs:
        # Um NIPC sem publicações não é um erro: marca-se como visto para não
        # voltar à fila todas as noites.
        await session.execute(
            text(
                "UPDATE registry_entities SET history_fetched_at = now(),"
                " updated_at = now() WHERE nif = :nipc"
            ),
            {"nipc": nipc},
        )
        return stats

    known_bodies = {
        r[0]
        for r in (
            await session.execute(
                text(
                    "SELECT irn_publication_id FROM registry_acts"
                    " WHERE nipc = :nipc AND body IS NOT NULL"
                ),
                {"nipc": nipc},
            )
        ).all()
        if r[0]
    }

    for pub in pubs:
        pub_id = str(pub.get("Id"))
        act_id, is_new = await upsert_act(
            session,
            nipc=nipc,
            irn_id=pub_id,
            act_type=pub.get("Act_Fact"),
            act_date=_parse_date(pub.get("Date")),
            entity_name=pub.get("Entity"),
            listing=pub,
            district=pub.get("District"),
            council=pub.get("Council"),
            priority=priority,
        )
        if not act_id:
            continue
        stats["novos"] += int(is_new)

        # Sem cliente não há rede: é o caso da entrega remota, em que o corpo
        # dos actos fica para a fila de conteúdo abrir depois.
        if not fetch_content or client is None or limiter is None:
            continue
        if pub_id in known_bodies:
            continue
        if not wants_content(pub.get("Act_Fact")):
            continue
        counts = await _fetch_and_parse(
            session, client, limiter, act_id=act_id, nipc=nipc, pub=pub
        )
        stats["conteudos"] += 1
        if counts:
            stats["quotas"] += counts["quotas"]
            stats["cargos"] += counts["officers"]

    await session.execute(
        text(
            """
            UPDATE registry_entities SET history_fetched_at = now(), updated_at = now()
             WHERE nif = :nipc
            """
        ),
        {"nipc": nipc},
    )
    return stats


async def seed_monitored(session: AsyncSession) -> int:
    """Põe na fila **todas** as empresas monitorizadas, os cinco tipos.

    O `run_mj_irn` corre só sobre internal/competitor/client, e é por isso que as
    7 de análise têm `registry_fetched_at` a NULL — zero informação de registo — e
    que as 17 `related` (descobertas como sócias) nunca são sincronizadas, o que
    trava o grafo no primeiro hop. Aqui entram as 205.
    """
    rows = (
        await session.execute(
            text(
                """
                SELECT nif, monitoring_type FROM companies
                 WHERE active = TRUE AND length(trim(nif)) = 9
                 ORDER BY CASE monitoring_type
                            WHEN 'internal' THEN 0 WHEN 'client' THEN 1
                            WHEN 'competitor' THEN 2 WHEN 'analysis' THEN 3 ELSE 4 END
                """
            )
        )
    ).all()
    added = 0
    for nif, mtype in rows:
        added += int(await enqueue_entity(session, nif=nif.strip(), reason=f"monitored:{mtype}", depth=0))
    return added


async def seed_all_companies(session: AsyncSession, depth: int = 9) -> int:
    """Põe na fila **todas** as entidades colectivas do país sem histórico próprio.

    A listagem diária do IRN não é um índice completo: numa amostra de 30
    empresas, **28 tinham actos que nunca vimos** — 379 no total, 83 deles de
    estrutura. A JOIN THE MOMENT tinha 4 actos no nosso índice e 30 na fonte.

    Entra a `depth = 9`, muito acima dos 0–3 da vizinhança, porque o
    `take_from_queue` ordena por `depth, requested_at`: as monitorizadas e os
    seus sócios continuam a ser servidos primeiro e o resto do país vem atrás
    sem lhes passar à frente. O `DO NOTHING` preserva quem já lá esteja mais
    perto.

    Uma só instrução, não 507 mil chamadas ao `enqueue_entity`.
    """
    result = await session.execute(
        text(
            """
            INSERT INTO registry_entity_fetch (nif, depth, reason)
            SELECT e.nif, :depth, 'nacional'
              FROM registry_entities e
             WHERE e.merged_into IS NULL
               AND e.nif IS NOT NULL
               AND e.kind IN ('company', 'public', 'foreign')
               AND e.history_fetched_at IS NULL
            ON CONFLICT (nif) DO NOTHING
            """
        ),
        {"depth": depth},
    )
    return result.rowcount or 0


async def enqueue_corporate_neighbours(session: AsyncSession, max_depth: int | None = None) -> int:
    """Alarga a fila às empresas ligadas ao que já temos.

    Só entidades: um sócio pessoa não se expande porque o IRN não pesquisa por
    NIF de pessoa. A profundidade de uma entrada nunca aumenta — um NIPC pedido
    como cliente não passa a hop 2 por aparecer também como sócio de um sócio.
    """
    max_depth = max_depth if max_depth is not None else settings.REGISTRY_MAX_DEPTH
    rows = (
        await session.execute(
            text(
                """
                WITH known AS (
                    SELECT f.nif, f.depth FROM registry_entity_fetch f
                )
                SELECT DISTINCT other.nif, min(k.depth) + 1 AS depth, min(k.nif) AS seed
                  FROM registry_edges e
                  JOIN registry_entities subj ON subj.id = e.subject_id
                  JOIN registry_entities other ON other.id = e.holder_id
                  JOIN known k ON k.nif = subj.nif
                 WHERE e.is_current
                   AND other.nif IS NOT NULL
                   AND other.kind IN ('company','public','foreign')
                   AND other.history_fetched_at IS NULL
                   AND k.depth < :maxd
                 GROUP BY other.nif
                UNION
                SELECT DISTINCT subj.nif, min(k.depth) + 1, min(k.nif)
                  FROM registry_edges e
                  JOIN registry_entities holder ON holder.id = e.holder_id
                  JOIN registry_entities subj ON subj.id = e.subject_id
                  JOIN known k ON k.nif = holder.nif
                 WHERE e.is_current
                   AND subj.nif IS NOT NULL
                   AND subj.kind IN ('company','public','foreign')
                   AND subj.history_fetched_at IS NULL
                   AND k.depth < :maxd
                 GROUP BY subj.nif
                """
            ),
            {"maxd": max_depth},
        )
    ).all()
    added = 0
    for nif, depth, seed in rows:
        added += int(
            await enqueue_entity(
                session, nif=nif.strip(), reason="socio-empresa", depth=int(depth), seed_nif=seed
            )
        )
    return added


# Nenhuma listagem por NIPC demora isto: quem estiver reclamado há mais tempo
# morreu. Mesmo valor que a fila de conteúdo usa, pela mesma razão.
STALE_CLAIM_MINUTES = 15


async def take_from_queue(session: AsyncSession, limit: int = 20) -> list[dict[str, Any]]:
    """Tira trabalho da fila, do mais perto do nosso universo para o mais longe.

    `FOR UPDATE SKIP LOCKED` para dois trabalhadores nunca pegarem no mesmo NIPC.

    **E recupera reclamações penduradas.** Enquanto só o CT consumia esta fila,
    em processo único, um restart levava tudo consigo e não havia nada preso.
    Com o trabalhador remoto pela frente — noutra máquina, noutra linha — um
    processo que morra deixaria entidades em `running` para sempre, invisíveis à
    fila. Com 507 mil na fila, ninguém daria por elas.
    """
    rows = (
        await session.execute(
            text(
                f"""
                WITH picked AS (
                    SELECT nif FROM registry_entity_fetch
                     WHERE attempts < 3
                       AND (status IN ('pending','error')
                            OR (status = 'running'
                                AND claimed_at
                                    < now() - interval '{STALE_CLAIM_MINUTES} minutes'))
                     ORDER BY depth, requested_at
                     FOR UPDATE SKIP LOCKED
                     LIMIT :lim
                )
                UPDATE registry_entity_fetch f
                   SET status = 'running', attempts = f.attempts + 1,
                       claimed_at = now()
                  FROM picked p
                 WHERE f.nif = p.nif
                RETURNING f.nif, f.depth, f.reason
                """
            ),
            {"lim": limit},
        )
    ).all()
    return [{"nif": r[0].strip(), "depth": r[1], "reason": r[2]} for r in rows]


async def finish_queue_item(
    session: AsyncSession, nif: str, *, status: str, stats: dict[str, int] | None = None,
    error: str | None = None,
) -> None:
    await session.execute(
        text(
            """
            UPDATE registry_entity_fetch SET
                status = :status,
                acts_seen = COALESCE(:seen, acts_seen),
                acts_new = COALESCE(:new, acts_new),
                last_error = :err,
                history_fetched_at = CASE WHEN :status = 'ok' THEN now() ELSE history_fetched_at END
             WHERE nif = :nif
            """
        ),
        {
            "nif": nif, "status": status, "err": error,
            "seen": (stats or {}).get("listados"),
            "new": (stats or {}).get("novos"),
        },
    )


async def process_queue(
    max_items: int = 50,
    max_seconds: float = 1800.0,
    workers: int | None = None,
    fetch_content: bool = True,
) -> dict[str, int]:
    """Consome a fila com um pool de trabalhadores e um tecto de ritmo comum.

    O orçamento é por tempo de parede **e** por contagem, o que chegar primeiro:
    a contagem sozinha mente, porque o custo de um NIPC varia entre uma chamada
    (uma entidade com dois actos) e trinta (uma com histórico longo).

    Com `fetch_content=False` só se lista, o que muda a ordem de grandeza: uma
    chamada por empresa em vez de uma mais três por cada acto que interesse.
    Medido, 102 empresas/min — o país inteiro em ~83 horas em vez de meses. Os
    actos ficam na fila de conteúdo e são abertos depois, pela escada de
    prioridades que já existe.
    """
    from app.db import AsyncSessionLocal

    workers = workers or settings.REGISTRY_WORKERS
    limiter = RateLimiter()
    deadline = time.monotonic() + max_seconds
    totals = {"entidades": 0, "listados": 0, "novos": 0, "conteudos": 0,
              "quotas": 0, "cargos": 0, "erros": 0}
    lock = asyncio.Lock()
    stop = asyncio.Event()

    async def worker(index: int) -> None:
        async with MjIrnClient() as client:
            while not stop.is_set() and time.monotonic() < deadline:
                async with lock:
                    if totals["entidades"] >= max_items:
                        return
                async with AsyncSessionLocal() as session:
                    batch = await take_from_queue(session, limit=1)
                    await session.commit()
                if not batch:
                    return
                item = batch[0]
                try:
                    async with AsyncSessionLocal() as session:
                        stats = await sync_nipc(
                            session, client, limiter, item["nif"],
                            priority=0 if item["depth"] == 0 else 10,
                            depth=item["depth"],
                            fetch_content=fetch_content,
                        )
                        await finish_queue_item(session, item["nif"], status="ok", stats=stats)
                        await session.commit()
                    async with lock:
                        totals["entidades"] += 1
                        for key in ("listados", "novos", "conteudos", "quotas", "cargos"):
                            totals[key] += stats.get(key, 0)
                except IrnVersionChanged:
                    # O IRN fez deploy e os apiVersion codificados ficaram
                    # obsoletos: nada nesta corrida vai funcionar. Devolve-se o
                    # item à fila e para-se o pool inteiro.
                    async with AsyncSessionLocal() as session:
                        await finish_queue_item(
                            session, item["nif"], status="pending", error="IrnVersionChanged"
                        )
                        await session.commit()
                    stop.set()
                    raise
                except Exception as e:  # noqa: BLE001
                    logger.warning("registo nif=%s falhou: %s", item["nif"], e)
                    async with AsyncSessionLocal() as session:
                        await finish_queue_item(
                            session, item["nif"], status="error", error=str(e)[:500]
                        )
                        await session.commit()
                    async with lock:
                        totals["erros"] += 1

    await asyncio.gather(*(worker(i) for i in range(workers)))
    totals["pedidos"] = limiter.acquired
    return totals
