"""Fila de conteúdo do registo comercial, servida a trabalhadores remotos.

Existe por uma razão só: **o que é escasso é o IP, não o processador.** O ritmo de
pedidos ao portal do IRN é limitado de propósito para não parecermos um ataque, e
esse tecto é por endereço. Um segundo trabalhador na mesma máquina não adianta
nada; um que saia por outra linha duplica a recolha sem que nenhum dos dois se
porte pior. Faltavam 805 mil actos por abrir e o CT sozinho levava sete dias e
meio.

**O trabalhador remoto só busca HTML.** Quem parseia e escreve arestas continua a
ser este processo, e isso não é detalhe de arrumação: no dia em que o
`REGISTRY_PARSE_VERSION` sobe — subiu hoje, para 2 — uma máquina lá fora com a
versão antiga escreveria dados desactualizados sem nada acusar. Assim o parser
está num sítio só, e o trabalhador não conhece a base de dados nem precisa de
credenciais dela.

A fila já estava preparada: `take_content` reclama com `FOR UPDATE SKIP LOCKED` e
devolve sozinha os actos reclamados há mais de 15 minutos. Um trabalhador que
morra a meio não perde nada — só atrasa aqueles actos um quarto de hora.
"""
import logging
import secrets
from typing import Annotated, Any
from uuid import UUID

from fastapi import APIRouter, Depends, Header, HTTPException, Query, status
from pydantic import BaseModel, Field
from sqlalchemy import text
from sqlalchemy.ext.asyncio import AsyncSession

from app.db import get_session
from app.services import app_settings
from app.services.registry_sweep import (
    ingest_content,
    publication_payload,
    release_content,
    take_content,
)
from app.services.registry_sync import (
    finish_queue_item,
    ingest_publications,
    take_from_queue,
)

router = APIRouter(prefix="/registry", tags=["public-api"])
logger = logging.getLogger(__name__)

# A chave vive no `app_settings` e não na tabela `api_keys` de propósito: as
# chaves de lá são de leitura (o dossiê do Sabichão) e reclamar trabalho é outra
# coisa. Assim nenhuma chave existente ganha poderes novos, e revoga-se esta
# apagando uma linha.
SETTING_KEY = "registry_worker_key"


async def require_worker_key(
    session: Annotated[AsyncSession, Depends(get_session)],
    x_worker_key: Annotated[str | None, Header(alias="X-Worker-Key")] = None,
) -> None:
    expected = await app_settings.get(session, SETTING_KEY)
    if not expected:
        raise HTTPException(
            status.HTTP_503_SERVICE_UNAVAILABLE,
            detail="recolha remota desligada: não há chave configurada",
        )
    # Comparação em tempo constante — a chave é o único guarda desta porta.
    if not x_worker_key or not secrets.compare_digest(x_worker_key, str(expected)):
        raise HTTPException(status.HTTP_401_UNAUTHORIZED, detail="chave inválida")


class Delivery(BaseModel):
    act_id: str
    html: str | None = None
    # Erro do lado do trabalhador. Conta para o tecto de três tentativas: sem
    # isso, um acto que rebente sempre voltava à fila para sempre.
    error: str | None = None
    # Uma mudança de `apiVersion` do IRN não é culpa do acto — devolve-se à fila
    # sem gastar tentativa, como faz o ciclo local.
    version_changed: bool = False
    # Devolvido sem ter chegado a ser tentado: o trabalhador foi mandado parar e
    # entrega o que tinha reclamado. Sem isto, esses actos ficavam à espera dos
    # 15 minutos da recuperação de reclamações penduradas em **cada** restart.
    released: bool = False


class DeliverIn(BaseModel):
    results: list[Delivery] = Field(default_factory=list, max_length=50)


@router.post("/claim", dependencies=[Depends(require_worker_key)])
async def claim(
    session: Annotated[AsyncSession, Depends(get_session)],
    n: int = Query(8, ge=1, le=20),
) -> dict[str, Any]:
    """Reclama actos da fila e devolve o que é preciso para os ir buscar.

    O trabalhador recebe o `id` (para depois entregar) e a publicação tal como o
    `content_for` a consome. Nada mais: nem NIPC para cruzar, nem o corpo de
    outros actos, nem contexto da base de dados.
    """
    batch = await take_content(session, limit=n)
    await session.commit()
    return {
        "acts": [
            {"act_id": act["id"], "publication": publication_payload(act)}
            for act in batch
        ]
    }


@router.post("/deliver", dependencies=[Depends(require_worker_key)])
async def deliver(
    payload: DeliverIn,
    session: Annotated[AsyncSession, Depends(get_session)],
) -> dict[str, int]:
    """Recebe o HTML buscado e fecha cada acto.

    Cada resultado é tratado à parte e com commit próprio: um acto cujo parse
    rebente não pode arrastar consigo os outros nove que vieram bem. É
    idempotente — reentregar o mesmo corpo reescreve as mesmas arestas, que é o
    que o `apply_parse` já garante ao apagar por `source_act_id`.
    """
    totals = {"guardados": 0, "vazios": 0, "erros": 0, "devolvidos": 0, "desconhecidos": 0}
    for item in payload.results:
        act = await _act_for(session, item.act_id)
        if not act:
            totals["desconhecidos"] += 1
            continue
        try:
            if item.released:
                # Não chegou a ser tentado: volta à fila intacto.
                await release_content(session, item.act_id, error=None, count_attempt=False)
                totals["devolvidos"] += 1
            elif item.version_changed or item.error:
                await release_content(
                    session,
                    item.act_id,
                    error=item.error,
                    count_attempt=not item.version_changed,
                )
                totals["erros"] += 1
            else:
                counts = await ingest_content(session, act, item.html)
                totals["vazios"] += counts["vazio"]
                totals["guardados"] += 1 - counts["vazio"]
            await session.commit()
        except Exception as e:  # noqa: BLE001 — um acto mau não derruba a entrega
            await session.rollback()
            logger.warning("entrega remota do acto %s falhou: %s", item.act_id, e)
            await release_content(session, item.act_id, error=str(e), count_attempt=True)
            await session.commit()
            totals["erros"] += 1
    return totals


class NipcDelivery(BaseModel):
    nif: str = Field(max_length=9)
    # A listagem tal como o IRN a devolve. O trabalhador não lhe toca: quem
    # decide o que é um acto de estrutura, o que vale a pena abrir e como se
    # datam os factos continua a ser este processo.
    publications: list[dict[str, Any]] = Field(default_factory=list)
    error: str | None = None
    version_changed: bool = False
    released: bool = False


class DeliverNipcIn(BaseModel):
    results: list[NipcDelivery] = Field(default_factory=list, max_length=20)


@router.post("/claim-nipc", dependencies=[Depends(require_worker_key)])
async def claim_nipc(
    session: Annotated[AsyncSession, Depends(get_session)],
    n: int = Query(4, ge=1, le=20),
) -> dict[str, Any]:
    """Reclama entidades da fila do histórico por NIPC.

    Lotes pequenos por defeito: cada NIPC são uma a trinta chamadas ao IRN,
    conforme o histórico, e um lote grande fica reclamado muito tempo. A
    recuperação de reclamações penduradas existe, mas é melhor não a exercitar.
    """
    batch = await take_from_queue(session, limit=n)
    await session.commit()
    return {"entities": [{"nif": item["nif"], "depth": item["depth"]} for item in batch]}


@router.post("/deliver-nipc", dependencies=[Depends(require_worker_key)])
async def deliver_nipc(
    payload: DeliverNipcIn,
    session: Annotated[AsyncSession, Depends(get_session)],
) -> dict[str, int]:
    """Recebe a listagem de cada NIPC e fecha-o na fila.

    **Sem rede deste lado**: o `ingest_publications` grava os actos como
    `pending` e a fila de conteúdo abre-os depois, pela escada de prioridades que
    já põe a cap table mais recente à frente do resto do histórico.

    Commit por entidade, como na entrega de conteúdo: um NIPC cuja escrita
    rebente não pode arrastar os outros do lote.
    """
    totais = {"listados": 0, "actos_novos": 0, "erros": 0, "devolvidos": 0}
    for item in payload.results:
        nif = (item.nif or "").strip()
        if len(nif) != 9:
            totais["erros"] += 1
            continue
        try:
            if item.released:
                await finish_queue_item(session, nif, status="pending", error=None)
                totais["devolvidos"] += 1
            elif item.version_changed:
                # O IRN fez deploy: não é culpa deste NIPC e não gasta tentativa.
                await finish_queue_item(
                    session, nif, status="pending", error="IrnVersionChanged"
                )
                totais["devolvidos"] += 1
            elif item.error:
                await finish_queue_item(session, nif, status="error", error=item.error)
                totais["erros"] += 1
            else:
                stats = await ingest_publications(
                    session, nif, item.publications, fetch_content=False
                )
                await finish_queue_item(session, nif, status="ok", stats=stats)
                totais["listados"] += stats["listados"]
                totais["actos_novos"] += stats["novos"]
            await session.commit()
        except Exception as e:  # noqa: BLE001 — um NIPC mau não derruba a entrega
            await session.rollback()
            logger.warning("entrega remota do NIPC %s falhou: %s", nif, e)
            await finish_queue_item(session, nif, status="error", error=str(e)[:500])
            await session.commit()
            totais["erros"] += 1
    return totais


async def _act_for(session: AsyncSession, act_id: str) -> dict[str, Any] | None:
    # O `act_id` vem de fora: um valor que não seja UUID rebentaria no cast e
    # levaria com ele o resto do lote.
    try:
        UUID(act_id)
    except (ValueError, AttributeError, TypeError):
        return None
    row = (
        await session.execute(
            text(
                "SELECT id::text, nipc, irn_publication_id, act_type, act_date,"
                "       entity_name, listing_json"
                "  FROM registry_acts WHERE id = CAST(:aid AS UUID)"
            ),
            {"aid": act_id},
        )
    ).mappings().first()
    return dict(row) if row else None
