"""Trabalhador remoto de recolha do IRN.

Corre fora do CT do Intel Grid — numa máquina com **outro IP** — e faz uma coisa
só: pede actos à fila, vai buscar o conteúdo ao portal do IRN e devolve o HTML.
Não conhece a base de dados, não parseia nada e não guarda estado.

É deliberado. O tecto de pedidos ao IRN é por endereço, portanto o que uma
segunda máquina traz de útil é o endereço dela; e deixar o parser num sítio só
evita que uma cópia desactualizada por aí escreva dados velhos sem ninguém dar
por isso.

Instalação e uso estão no `tools/README-worker.md`. Em resumo:

    IG_URL=http://127.0.0.1:8088 IG_WORKER_KEY=... python irn_worker.py

Um `--once` faz uma única volta e sai, que é como se confirma a instalação.
"""
import argparse
import asyncio
import logging
import os
import signal
import sys
import time

import httpx

# Os módulos do repositório que valem a pena reutilizar: o cliente do IRN (a
# sessão anónima, o cookie `nr2Users`, o `X-CSRFToken` e os `apiVersion`, que
# mudam quando o portal faz deploy) e o tecto de ritmo. Reescrever qualquer um
# deles aqui seria ter duas versões de algo que já custou a acertar.
_HERE = os.path.dirname(os.path.abspath(__file__))
# Instalado, o `backend` fica ao lado deste ficheiro; no repositório fica um
# nível acima, ao lado do `tools`. Funciona nos dois sítios para se poder
# experimentar aqui antes de mandar para lá.
for _candidate in (os.path.join(_HERE, "backend"), os.path.join(_HERE, "..", "backend")):
    if os.path.isdir(os.path.join(_candidate, "app")):
        sys.path.insert(0, os.path.abspath(_candidate))
        break

from app.scrapers.base import RateLimiter  # noqa: E402
from app.scrapers.mj_irn import IrnVersionChanged, MjIrnClient  # noqa: E402

logging.basicConfig(
    level=logging.INFO,
    format="%(asctime)s %(levelname)s %(message)s",
)
logger = logging.getLogger("irn_worker")

BASE_URL = os.environ.get("IG_URL", "http://127.0.0.1:8088").rstrip("/")
WORKER_KEY = os.environ.get("IG_WORKER_KEY", "")
WORKERS = int(os.environ.get("IG_WORKERS", "3"))
BATCH = int(os.environ.get("IG_BATCH", "8"))
# Fila vazia não é fim: o varrimento nacional acrescenta actos todos os dias.
IDLE_SLEEP = float(os.environ.get("IG_IDLE_SLEEP", "300"))
DELIVER_TRIES = 3


class Coordinator:
    """A fila, vista de fora: reclamar e entregar, por HTTP.

    Reclama em lotes porque a viagem até ao CT custa 23 ms e o tecto de ritmo
    obriga a esperar meio segundo por acto — pedir um de cada vez seria pagar a
    viagem sem necessidade. Entrega um a um, para que um acto que rebente não
    leve consigo os outros do lote.
    """

    def __init__(self, client: httpx.AsyncClient) -> None:
        self._http = client
        self._lock = asyncio.Lock()
        self._pending: list[dict] = []
        self._exhausted = False

    async def next_act(self) -> dict | None:
        async with self._lock:
            if self._pending:
                return self._pending.pop()
            if self._exhausted:
                return None
            r = await self._http.post("/api/v1/registry/claim", params={"n": BATCH})
            r.raise_for_status()
            acts = r.json().get("acts") or []
            if not acts:
                self._exhausted = True
                return None
            self._pending = acts
            return self._pending.pop()

    async def deliver(self, **result) -> None:
        """Entrega, insistindo.

        O HTML já custou um pedido ao IRN dentro do tecto de ritmo; deitá-lo fora
        porque o CT esteve em baixo dois segundos — um deploy, o túnel a
        restabelecer-se — é desperdiçar a única coisa que aqui é escassa. Três
        tentativas com espera a crescer chegam para atravessar um restart.
        """
        for attempt in range(DELIVER_TRIES):
            try:
                r = await self._http.post(
                    "/api/v1/registry/deliver", json={"results": [result]}
                )
                r.raise_for_status()
                return
            except httpx.HTTPError:
                if attempt == DELIVER_TRIES - 1:
                    raise
                await asyncio.sleep(2 ** attempt)

    async def release_pending(self) -> int:
        """Devolve à fila o que foi reclamado e não chegou a ser buscado.

        Sem isto, cada paragem — um deploy, um restart do serviço — deixava até
        duas dezenas de actos reclamados à espera dos 15 minutos da recuperação
        automática. Devolvidos assim, ficam disponíveis no instante seguinte e
        não gastam tentativa nenhuma.
        """
        async with self._lock:
            pending, self._pending = self._pending, []
        if not pending:
            return 0
        try:
            r = await self._http.post(
                "/api/v1/registry/deliver",
                json={
                    "results": [
                        {"act_id": a["act_id"], "released": True} for a in pending
                    ]
                },
            )
            r.raise_for_status()
        except httpx.HTTPError as e:
            # Não é grave: a recuperação de reclamações penduradas do CT trata
            # deles em 15 minutos. Só perdemos a pressa.
            logger.warning("não foi possível devolver %d actos (%s)", len(pending), e)
            return 0
        return len(pending)


async def run_once(shutdown: asyncio.Event | None = None) -> dict[str, int]:
    """Uma volta: esvazia o que a fila der e devolve as contagens.

    O `shutdown` vem do sinal de paragem. Quando dispara, os trabalhadores saem
    do ciclo e o que estiver reclamado por buscar é devolvido à fila antes de
    fechar — é a diferença entre um restart custar zero e custar 15 minutos.
    """
    totals = {"actos": 0, "erros": 0, "devolvidos": 0}
    limiter = RateLimiter()
    stop = shutdown or asyncio.Event()

    headers = {"X-Worker-Key": WORKER_KEY}
    async with httpx.AsyncClient(
        base_url=BASE_URL, headers=headers, timeout=60.0
    ) as http:
        coordinator = Coordinator(http)

        async def worker(n: int) -> None:
            async with MjIrnClient() as irn:
                while not stop.is_set():
                    try:
                        act = await coordinator.next_act()
                    except httpx.HTTPError as e:
                        # O túnel pode cair. Não é motivo para desistir da volta:
                        # os actos já reclamados voltam sozinhos à fila ao fim de
                        # 15 minutos e o próximo ciclo apanha-os.
                        logger.warning("worker %d: falha a reclamar (%s)", n, e)
                        await asyncio.sleep(10)
                        continue
                    if act is None:
                        return
                    act_id = act["act_id"]
                    try:
                        await limiter.acquire(3)
                        html = await irn.content_for(act["publication"])
                        await coordinator.deliver(act_id=act_id, html=html)
                        totals["actos"] += 1
                    except IrnVersionChanged as e:
                        # O portal mudou de versão: pára tudo em vez de martelar.
                        # O canário do CT trata de avisar.
                        logger.error("IRN mudou de versão: %s", e)
                        await coordinator.deliver(act_id=act_id, version_changed=True)
                        stop.set()
                        return
                    except Exception as e:  # noqa: BLE001
                        logger.warning("acto %s falhou: %s", act_id, e)
                        totals["erros"] += 1
                        try:
                            await coordinator.deliver(act_id=act_id, error=str(e)[:400])
                        except httpx.HTTPError:
                            # Sem forma de devolver: a recuperação de reclamações
                            # penduradas do CT trata dele em 15 minutos.
                            logger.warning("acto %s ficou por devolver", act_id)

        try:
            await asyncio.gather(*(worker(i) for i in range(WORKERS)))
        finally:
            # Corra bem ou mal, o que ficou reclamado sem ser buscado volta já
            # para a fila em vez de esperar pela recuperação automática.
            totals["devolvidos"] = await coordinator.release_pending()
    return totals


async def main() -> None:
    parser = argparse.ArgumentParser(description=__doc__)
    parser.add_argument("--once", action="store_true", help="uma volta e sai")
    args = parser.parse_args()

    if not WORKER_KEY:
        logger.error("falta o IG_WORKER_KEY")
        raise SystemExit(2)

    # O systemd manda SIGTERM ao parar. Sem isto o processo morria a meio e
    # deixava os actos reclamados pendurados 15 minutos — em cada restart.
    shutdown = asyncio.Event()
    loop = asyncio.get_running_loop()
    for sig in (signal.SIGTERM, signal.SIGINT):
        loop.add_signal_handler(sig, shutdown.set)

    logger.info("worker a arrancar: %s, %d corrotinas, lotes de %d",
                BASE_URL, WORKERS, BATCH)
    while not shutdown.is_set():
        started = time.monotonic()
        totals = await run_once(shutdown)
        elapsed = max(time.monotonic() - started, 0.001)
        logger.info(
            "volta: %d actos, %d erros, %d devolvidos, %.0f actos/min",
            totals["actos"], totals["erros"], totals["devolvidos"],
            totals["actos"] / elapsed * 60,
        )
        if args.once or shutdown.is_set():
            return
        # Espera acordável: um sinal durante a sesta não fica 5 minutos à espera.
        try:
            await asyncio.wait_for(shutdown.wait(), timeout=IDLE_SLEEP)
        except TimeoutError:
            pass
    logger.info("worker parado a pedido")


if __name__ == "__main__":
    asyncio.run(main())
