"""Varrimento nacional dos actos de registo, dia a dia, e a fila de conteúdo.

É esta via que responde à pergunta que motivou tudo — *"o quim é sócio na
concorrente abc e por sua vez tem quotas na empresa hhi e na xyz"*. A via por
NIPC não chega para isso: o IRN só pesquisa por NIPC de entidade, e com o NIF de
um indivíduo devolve zero. Saber onde mais é que uma pessoa é sócia **só se
responde com um índice nacional construído localmente**, e a resposta vale
exactamente o que valer a janela recolhida — daí a cobertura ir junta com qualquer
ecrã que a use.

Duas metades, deliberadamente separadas:

* **Listagem** (`sweep_days`) — barata: 31 páginas e ~4 s por dia útil. Cinco anos
  são ~1,5 h e ficam com NIPC, nome, tipo de acto e data de todos os ~780 mil
  actos. É o que permite decidir prioridades **sem** abrir conteúdo.
* **Conteúdo** (`process_content`) — três chamadas por acto. Só se abre o que
  traz pessoas ou quotas; prestações de contas, dissoluções e mudanças de sede
  nunca. Corre em contínuo com tecto de ritmo, e o que interessa aterra primeiro:
  o nosso universo, depois a vizinhança do grafo, depois a cap table mais recente
  de cada entidade, e só no fim o resto do histórico.
"""
import asyncio
import hashlib
import logging
import time
from datetime import date, datetime, timedelta
from typing import Any

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

from app.config import settings
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
from app.services.registry_ingest import apply_parse, upsert_act
from app.scrapers.base import RateLimiter
from app.services.registry_sync import _parse_date

logger = logging.getLogger(__name__)

# Um dia útil ronda os 618 actos = 31 páginas. O cap é folga, não expectativa; se
# alguma vez for atingido o dia fecha como `partial` e não como `ok`.
PAGE_CAP = 120


async def enqueue_days(session: AsyncSession, date_from: date, date_to: date) -> int:
    """Marca os dias a varrer. Idempotente.

    Um registo por dia, em vez do `MAX(params->>'date_to')` que o resto do
    scheduler usa: esse padrão **não consegue representar um buraco**, e uma
    corrida manual sobre uma janela posterior deixaria os dias anteriores em falta
    órfãos para sempre. A 618 actos por dia, isso não é um detalhe.
    """
    result = await session.execute(
        text(
            """
            INSERT INTO registry_sweep_days (day)
            SELECT d::date FROM generate_series(
                CAST(:d1 AS DATE), CAST(:d2 AS DATE), interval '1 day') AS d
            ON CONFLICT (day) DO NOTHING
            """
        ),
        {"d1": date_from, "d2": date_to},
    )
    return result.rowcount or 0


async def take_days(session: AsyncSession, limit: int = 5) -> list[date]:
    """Dias por varrer, do mais recente para o mais antigo.

    Do mais recente primeiro de propósito: a estrutura societária **actual** é o
    que responde à pergunta inversa, e vem dos actos recentes. O histórico mais
    fundo continua a encher por trás.
    """
    rows = (
        await session.execute(
            text(
                """
                WITH picked AS (
                    SELECT day FROM registry_sweep_days
                     WHERE status IN ('pending','error','partial') AND attempts < 4
                     ORDER BY day DESC
                     FOR UPDATE SKIP LOCKED
                     LIMIT :lim
                )
                UPDATE registry_sweep_days s
                   SET status = 'running', attempts = s.attempts + 1, started_at = now()
                  FROM picked p
                 WHERE s.day = p.day
                RETURNING s.day
                """
            ),
            {"lim": limit},
        )
    ).all()
    return [r[0] for r in rows]


async def sweep_day(
    session: AsyncSession, client: MjIrnClient, limiter: RateLimiter, day: date
) -> dict[str, Any]:
    """Lista um dia nacional e grava os actos sem conteúdo."""
    await limiter.acquire()
    rows = await client.publications_for_day(day, max_pages=PAGE_CAP)
    stats = {"listados": len(rows), "novos": 0, "paginas": (len(rows) + 19) // 20}

    for pub in rows:
        act_id, is_new = await upsert_act(
            session,
            nipc=(pub.get("NIF_NIPC") or "").strip(),
            irn_id=str(pub.get("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=60,
        )
        if act_id:
            stats["novos"] += int(is_new)

    truncated = stats["paginas"] >= PAGE_CAP
    await session.execute(
        text(
            """
            UPDATE registry_sweep_days SET
                status = :status, listed_count = :n, pages = :p, new_count = :new,
                finished_at = now(), last_error = :err
             WHERE day = CAST(:day AS DATE)
            """
        ),
        {
            "day": day, "n": stats["listados"], "p": stats["paginas"],
            "new": stats["novos"],
            "status": "partial" if truncated else "ok",
            "err": f"cap de {PAGE_CAP} páginas atingido" if truncated else None,
        },
    )
    return stats


async def refresh_priorities(session: AsyncSession) -> int:
    """Escada de prioridades da fila de conteúdo.

    Recalcula-se porque a vizinhança do grafo cresce: um NIPC que ontem era
    anónimo pode hoje ser sócio de um cliente. Sem isto, a fila ficaria congelada
    na fotografia do dia em que cada acto foi listado.

    0  monitorizadas · 10 vizinhas no grafo · 20 watchlist e entidades com
    anúncios CIRE · 30 cap table mais recente de cada entidade · 50 restantes
    actos com pessoas · 90 nunca.
    """
    total = 0
    steps = (
        (0, "SELECT nif FROM companies WHERE active"),
        (
            10,
            """
            SELECT DISTINCT e.nif FROM registry_entities e
             WHERE e.nif IS NOT NULL AND EXISTS (
                SELECT 1 FROM registry_edges r
                  JOIN registry_entities o
                    ON o.id = CASE WHEN r.holder_id = e.id THEN r.subject_id ELSE r.holder_id END
                 WHERE (r.holder_id = e.id OR r.subject_id = e.id) AND r.is_current
                   AND o.nif IN (SELECT nif FROM companies WHERE active)
             )
            """,
        ),
        (
            20,
            """
            SELECT nif FROM insolvency_watchlist WHERE active
             UNION SELECT DISTINCT nif FROM cire_announcement_parties WHERE nif IS NOT NULL
            """,
        ),
    )
    for priority, query in steps:
        result = await session.execute(
            text(
                f"""
                UPDATE registry_acts a SET priority = :p
                 WHERE a.content_status = 'pending' AND a.wants_content
                   AND a.priority > :p
                   AND a.nipc IN ({query})
                """
            ),
            {"p": priority},
        )
        total += result.rowcount or 0

    # O acto de cap table mais recente de cada entidade: é o que dá a estrutura
    # actual, e é por isso que sobe acima do resto do histórico. Uma entidade fica
    # útil ao primeiro acto aberto em vez de precisar dos trinta.
    result = await session.execute(
        text(
            """
            WITH latest AS (
                SELECT DISTINCT ON (nipc) id
                  FROM registry_acts
                 WHERE content_status = 'pending' AND wants_content
                   AND act_class IN ('constituicao','alteracao_contrato','aumento_capital',
                                     'reducao_capital','transmissao_quotas','fusao')
                 ORDER BY nipc, act_date DESC NULLS LAST
            )
            UPDATE registry_acts a SET priority = 30
              FROM latest l
             WHERE a.id = l.id AND a.priority > 30
            """
        )
    )
    total += result.rowcount or 0
    return total


async def refresh_acts_count(session: AsyncSession) -> int:
    """Põe o `registry_entities.acts_count` de acordo com os actos listados.

    Era escrito só no `apply_parse`, ou seja quando um acto era **parseado**. A
    listagem nacional acrescenta actos sem os parsear — é o que ela faz —, e o
    contador ficava a contar outra coisa: a Trablisa dizia 59 actos e tinha 97.

    Importa porque é este número que decide a frase "sem actos do registo
    comercial recolhidos" na ficha de uma entidade.
    """
    result = await session.execute(
        text(
            """
            WITH reais AS (
                SELECT nipc, count(*) AS n FROM registry_acts GROUP BY nipc
            )
            UPDATE registry_entities e SET acts_count = r.n, updated_at = now()
              FROM reais r
             WHERE e.nif = r.nipc AND e.acts_count IS DISTINCT FROM r.n
            """
        )
    )
    return result.rowcount or 0


STALE_CLAIM_MINUTES = 15


async def take_content(session: AsyncSession, limit: int = 1) -> list[dict[str, Any]]:
    """Reclama actos da fila, e **recupera reclamações penduradas**.

    O `running` marca o acto como tomado, para dois trabalhadores nunca pegarem no
    mesmo. Mas um trabalhador que morra — restart do contentor, deploy, OOM —
    deixa o acto tomado e invisível para a fila. Por isso a condição inclui
    também os que estão `running` há mais de 15 minutos: nenhum acto demora isso,
    portanto quem lá estiver está pendurado.

    Sem esta recuperação, cada restart perdia definitivamente os actos em curso.
    Nos dias do varrimento isso resolveu-se com um `UPDATE` à mão; a 193 mil
    actos não é uma opção.
    """
    rows = (
        await session.execute(
            text(
                f"""
                WITH picked AS (
                    SELECT id FROM registry_acts
                     WHERE wants_content
                       AND content_attempts < 3
                       AND (
                         content_status = 'pending'
                         OR (content_status = 'running'
                             AND content_claimed_at
                                 < now() - interval '{STALE_CLAIM_MINUTES} minutes')
                       )
                     ORDER BY priority, act_date DESC NULLS LAST
                     FOR UPDATE SKIP LOCKED
                     LIMIT :lim
                )
                UPDATE registry_acts a
                   SET content_status = 'running', content_claimed_at = now()
                  FROM picked p
                 WHERE a.id = p.id
                RETURNING a.id::text, a.nipc, a.irn_publication_id, a.act_type,
                          a.act_date, a.entity_name, a.listing_json
                """
            ),
            {"lim": limit},
        )
    ).mappings().all()
    return [dict(r) for r in rows]


def publication_payload(act: dict[str, Any]) -> dict[str, Any]:
    """O que o `MjIrnClient.content_for` precisa de saber sobre uma publicação.

    A listagem guardada é a fonte preferida; quando falta, reconstrói-se a partir
    das colunas do acto. Vive aqui em vez de dentro do ciclo porque o worker
    remoto recebe exactamente isto e mais nada — é a fronteira entre "o que se vai
    buscar" e "o que se faz com o que veio".
    """
    return act.get("listing_json") or {
        "Id": act["irn_publication_id"],
        "Date": act["act_date"].isoformat() if act["act_date"] else None,
        "NIF_NIPC": act["nipc"],
        "Entity": act["entity_name"],
    }


async def ingest_content(
    session: AsyncSession, act: dict[str, Any], html: str | None
) -> dict[str, int]:
    """Guarda e parseia o conteúdo de um acto. **Não** faz commit.

    É o que acontece depois de o HTML chegar, seja ele buscado aqui ou por um
    trabalhador noutro sítio. Existe como função própria para haver um só caminho:
    o parser e a escrita de arestas ficam do lado da base de dados, e uma subida
    de `REGISTRY_PARSE_VERSION` não tem de ser replicada em máquina nenhuma.
    """
    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 {"vazio": 1, "quotas": 0, "officers": 0}

    await session.execute(
        text(
            """
            UPDATE registry_acts SET
                body = :body, body_sha256 = :sha, body_source = 'irn',
                content_status = 'fetched', content_fetched_at = now(),
                content_attempts = content_attempts + 1
             WHERE id = CAST(:aid AS UUID)
            """
        ),
        {
            "aid": act["id"], "body": body,
            "sha": hashlib.sha256(body.encode()).hexdigest(),
        },
    )
    parsed = parse_act(body, act["act_type"])
    counts = await apply_parse(
        session, act_id=act["id"], nipc=act["nipc"],
        parsed=parsed, act_date=act["act_date"],
    )
    return {"vazio": 0, "quotas": counts["quotas"], "officers": counts["officers"]}


async def release_content(
    session: AsyncSession, act_id: str, *, error: str | None, count_attempt: bool
) -> None:
    """Devolve um acto à fila. **Não** faz commit.

    `count_attempt` distingue as duas razões para largar um acto: uma falha ao
    abri-lo conta para o tecto de três tentativas, uma mudança de `apiVersion` do
    IRN não — nesse caso a culpa não é do acto e ele há-de abrir na corrida
    seguinte.
    """
    await session.execute(
        text(
            f"""
            UPDATE registry_acts SET content_status = 'pending'
                {", content_attempts = content_attempts + 1" if count_attempt else ""}
                {", parse_error = :err" if error else ""}
             WHERE id = CAST(:aid AS UUID)
            """
        ),
        {"aid": act_id, **({"err": error[:500]} if error else {})},
    )


async def process_content(
    max_items: int = 2000,
    max_seconds: float = 3600.0,
    workers: int | None = None,
) -> dict[str, int]:
    """Abre e parseia conteúdo da fila, com um pool e tecto de ritmo comum.

    Orçamento por tempo de parede **e** contagem, o que chegar primeiro: 2000
    actos são 9 minutos sem pausa nenhuma e mais de meia hora ao ritmo com que
    isto corre de facto.
    """
    from app.db import AsyncSessionLocal

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

    async def worker() -> None:
        async with MjIrnClient() as client:
            while not stop.is_set() and time.monotonic() < deadline:
                async with lock:
                    if totals["actos"] >= max_items:
                        return
                async with AsyncSessionLocal() as session:
                    batch = await take_content(session, limit=1)
                    await session.commit()
                if not batch:
                    return
                act = batch[0]
                try:
                    await limiter.acquire(3)
                    html = await client.content_for(publication_payload(act))
                    async with AsyncSessionLocal() as session:
                        counts = await ingest_content(session, act, html)
                        await session.commit()
                    async with lock:
                        totals["actos"] += 1
                        totals["vazios"] += counts["vazio"]
                        totals["quotas"] += counts["quotas"]
                        totals["cargos"] += counts["officers"]
                except IrnVersionChanged:
                    async with AsyncSessionLocal() as session:
                        await release_content(
                            session, act["id"], error=None, count_attempt=False
                        )
                        await session.commit()
                    stop.set()
                    raise
                except Exception as e:  # noqa: BLE001
                    logger.warning("conteúdo do acto %s falhou: %s", act["id"], e)
                    async with AsyncSessionLocal() as session:
                        await release_content(
                            session, act["id"], error=str(e), count_attempt=True
                        )
                        await session.commit()
                    async with lock:
                        totals["erros"] += 1

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


async def sweep_days(
    max_days: int = 5, max_seconds: float = 1800.0, workers: int | None = None
) -> dict[str, int]:
    from app.db import AsyncSessionLocal

    workers = workers or max(1, min(settings.REGISTRY_WORKERS, 2))
    limiter = RateLimiter()
    deadline = time.monotonic() + max_seconds
    totals = {"dias": 0, "listados": 0, "novos": 0, "erros": 0}
    lock = asyncio.Lock()

    async def worker() -> None:
        async with MjIrnClient() as client:
            while time.monotonic() < deadline:
                async with lock:
                    if totals["dias"] >= max_days:
                        return
                async with AsyncSessionLocal() as session:
                    days = await take_days(session, limit=1)
                    await session.commit()
                if not days:
                    return
                day = days[0]
                try:
                    async with AsyncSessionLocal() as session:
                        stats = await sweep_day(session, client, limiter, day)
                        await session.commit()
                    async with lock:
                        totals["dias"] += 1
                        totals["listados"] += stats["listados"]
                        totals["novos"] += stats["novos"]
                except IrnVersionChanged:
                    async with AsyncSessionLocal() as session:
                        await session.execute(
                            text(
                                "UPDATE registry_sweep_days SET status = 'pending',"
                                " last_error = 'IrnVersionChanged'"
                                " WHERE day = CAST(:d AS DATE)"
                            ),
                            {"d": day},
                        )
                        await session.commit()
                    raise
                except Exception as e:  # noqa: BLE001
                    logger.warning("varrimento de %s falhou: %s", day, e)
                    async with AsyncSessionLocal() as session:
                        await session.execute(
                            text(
                                "UPDATE registry_sweep_days SET status = 'error',"
                                " last_error = :err WHERE day = CAST(:d AS DATE)"
                            ),
                            {"d": day, "err": str(e)[:500]},
                        )
                        await session.commit()
                    async with lock:
                        totals["erros"] += 1

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


async def sweep_progress(session: AsyncSession) -> dict[str, Any]:
    row = (
        await session.execute(
            text(
                """
                SELECT
                  count(*) FILTER (WHERE status = 'ok') AS dias_ok,
                  count(*) FILTER (WHERE status IN ('pending','error','partial')) AS dias_por_fazer,
                  min(day) FILTER (WHERE status = 'ok') AS de,
                  max(day) FILTER (WHERE status = 'ok') AS ate,
                  (SELECT count(*) FROM registry_acts) AS actos,
                  -- `running` conta como por abrir: um acto reclamado não está
                  -- lido, e se o trabalhador morreu volta à fila em 15 minutos.
                  (SELECT count(*) FROM registry_acts
                    WHERE content_status IN ('pending', 'running')
                      AND wants_content) AS conteudo_por_abrir,
                  (SELECT count(*) FROM registry_acts WHERE body IS NOT NULL) AS actos_lidos
                FROM registry_sweep_days
                """
            )
        )
    ).mappings().first()
    return dict(row) if row else {}
