#!/usr/local/bin/python3.8
"""
Cliente de ingestao do Contextia -- montagem e envio do envelope de dataset.

Este modulo existe separado do conector por dois motivos:

  * O conector importa `models` / `ssh_server` (banco de origem do cliente).
    Quem quiser testar o formato do envelope, ou reaproveitar o envio em outro
    extrator, nao precisa ter o banco de origem a mao -- basta importar daqui.

  * A montagem do envelope e a unica parte que depende do contrato da API. Se o
    contrato mudar, muda so este arquivo.

Formato canonico (POST /api/v1/intake/json):

    {
      "dataset": "clientes",
      "description": "Cadastro de clientes do ERP",
      "key": "ID",
      "fields": {"ID": "Codigo interno do cliente", ...},
      "sync": "20260828T031500Z-9f3a1c",
      "complete": true,
      "records": [ {"ID": 1, "RAZAOSOCIAL": "..."} ]
    }

  * `dataset` (obrigatorio) identifica a "tabela" no Contextia. O servidor
    normaliza para minusculas/[a-z0-9_]; este modulo ja envia normalizado.
  * `key` (opcional) e o campo identificador. Com ele o servidor faz UPSERT:
    reenviar o mesmo ID ATUALIZA o registro em vez de duplicar. Sem ele, a
    deduplicacao e por hash do conteudo (dois registros diferentes do mesmo
    ID viram duas linhas).
  * `description`/`fields` (opcionais) viram definicao semantica automatica
    (source `client`) e o item de catalogo do dataset -- sem sobrescrever
    definicao ja existente feita por pessoa ou pelo agente.
  * `sync`/`complete` (opcionais) formam o CICLO DE CARGA -- ver abaixo.
  * `records` sao os registros CRUS. Nada e acrescentado a eles.

`description`/`key`/`fields` vao em TODOS os lotes do mesmo dataset. Isso e
proposital: o servidor faz upsert do metadado, entao repetir e idempotente e
qualquer lote isolado (um retry, por exemplo) chega auto-descrito.


CICLO DE CARGA (`sync` + `complete`)
------------------------------------
O upsert por `key` resolve inclusao e alteracao, mas nao EXCLUSAO: um registro
apagado na origem continuaria vivo no Contextia para sempre, e as agregacoes do
agente passariam a somar coisa que nao existe mais. O ciclo de carga fecha esse
buraco:

  * `sync` e o identificador da carga. O MESMO valor vai em TODOS os lotes
    daquele dataset naquela carga. Todo registro visto na carga (inserido,
    atualizado OU identico) fica marcado com ele no servidor.
  * `complete: true` vai APENAS no ULTIMO lote do dataset. Ao receber esse lote,
    o servidor marca como REMOVIDO (`removed_at`) todo registro do dataset que
    nao foi visto naquela carga -- e eles somem da view. Registro removido que
    volta numa carga futura e RESSUSCITADO.

  * SEM `sync`/`complete`, NADA e removido. Esse e o comportamento seguro, e e
    exatamente o que deve acontecer quando a carga falha no meio: metade dos
    registros da origem chegou, e se o `complete` fosse enviado assim mesmo o
    servidor apagaria a outra metade -- que existe na origem e so nao chegou.

REGRA DE OURO, implementada em `enviar_dataset()`: se QUALQUER lote daquele
dataset falhou, o `complete` NAO e enviado. A carga vira um upsert normal
(o que chegou entra/atualiza), nada e removido, e a proxima execucao completa
o servico. Falha de carga nunca vira perda de dado.

ESCOPO DO `sync`: este modulo gera UM identificador por EXECUCAO do processo
(`L_SYNC`, montado no import) e o usa como padrao para todos os datasets. E o
que faz sentido operacionalmente -- "a carga da noite de 28/08" e uma coisa so,
e o identificador aparece igual no log de todos os datasets. Nao ha risco em
compartilhar: a remocao e sempre por dataset (o servidor so remove registros do
dataset cujo lote trouxe `complete`), entao um `sync` compartilhado nunca faz um
dataset interferir no outro. Quem quiser um por dataset passa `p_sync=`
explicitamente; quem quiser DESLIGAR o ciclo (e nunca remover nada) passa
`p_sync=None`.

A resposta da API traz, no bloco `dataset`: `inserted`, `updated`,
`duplicates`, `sync`, `removed`, `resurrected`, `active_count` e
`removed_count`.
"""

import datetime
import decimal
import json
import os
import re
import uuid

import requests  # type: ignore


# ---------------------------------------------------------------------------
# Configuracao do destino (Contextia)
# ---------------------------------------------------------------------------
# A chave identifica tenant + projeto. Trocar a chave = trocar o destino dos
# dados; nao ha nenhum outro parametro de roteamento.
#
# Vem do ambiente porque connector.py e o MESMO arquivo para todo cliente: com
# a chave fixa aqui, atender um segundo cliente exigiria editar um arquivo
# compartilhado, e uma carga sairia no projeto errado no dia em que alguem
# esquecesse de trocar de volta. O valor abaixo e o default do primeiro
# cliente, para que o conector escrito a mao continue rodando sem mudanca.
#
#   export CONTEXTIA_API_KEY=ctx_...
#
# SEM PADRAO, de proposito. Um valor embutido aqui apontaria para a instalacao
# de OUTRO cliente: quem esquecesse de exportar mandaria a carga para o servidor
# errado, com a chave certa, e sem erro nenhum. Melhor falhar alto.
L_API_KEY = os.environ.get("CONTEXTIA_API_KEY", "")

# O Contextia e instalado por URL: cada instalacao tem a sua. UMA variavel
# aponta para a raiz, e todos os enderecos saem dela -- API, intake e MCP.
#
#   export CONTEXTIA_URL=https://contextia.cliente.com.br/ls/contextia/public
#
# As especificas continuam valendo e GANHAM da CONTEXTIA_URL, para o caso raro
# de intake e administracao ficarem em hosts diferentes.
L_URL = os.environ.get("CONTEXTIA_URL", "").rstrip("/")

L_API_BASE = os.environ.get("CONTEXTIA_API_BASE") or (L_URL + "/api/v1" if L_URL else "")
L_INTAKE_URL = os.environ.get("CONTEXTIA_INTAKE_URL") or (L_API_BASE + "/intake/json" if L_API_BASE else "")
L_MCP_URL = os.environ.get("CONTEXTIA_MCP_URL") or (L_URL + "/mcp" if L_URL else "")


def exigir_destino():
    """Aborta se nao souber PARA ONDE e COM QUAL chave enviar."""
    faltando = []
    if not L_URL and not L_INTAKE_URL:
        faltando.append('CONTEXTIA_URL (a instalacao DESTE cliente)')
    if not L_API_KEY:
        faltando.append('CONTEXTIA_API_KEY (a chave DESTE projeto)')
    if faltando:
        raise SystemExit(
            'Falta definir: %s.\n'
            'Nao ha padrao embutido: um valor default mandaria a carga para a\n'
            'instalacao de outro cliente, sem erro nenhum.' % ', '.join(faltando)
        )

# Registros por requisicao. O payload vai inteiro em memoria e no corpo do POST;
# um dataset grande estouraria os limites de upload/execucao do PHP.
L_CHUNK = 500

# Timeout por requisicao, em segundos.
L_TIMEOUT = 60


# ------------------------------------------  Identificador da carga
def gerar_sync():
    """
    Identificador de uma carga: '20260828T031500Z-9f3a1c'.

    Legivel (da para saber QUANDO a carga rodou olhando o log ou o banco) e
    unico o bastante (o sufixo aleatorio separa duas execucoes iniciadas no
    mesmo segundo, por exemplo dois cron concorrentes).
    """
    loc_agora = datetime.datetime.now(datetime.timezone.utc)
    return loc_agora.strftime("%Y%m%dT%H%M%SZ") + "-" + uuid.uuid4().hex[:6]


# Um `sync` por EXECUCAO do processo, compartilhado por todos os datasets --
# ver "ESCOPO DO `sync`" no docstring do modulo. E o padrao de `enviar_dataset`.
L_SYNC = gerar_sync()


# ------------------------------------------  Encoder para datas e decimais
def alchemyencoder(obj):
    """JSON encoder function for SQLAlchemy special classes."""
    if isinstance(obj, datetime.date):
        return obj.isoformat()
    elif isinstance(obj, decimal.Decimal):
        return float(obj)


# ------------------------------------------  Nome do dataset
def normalizar_dataset(p_nome):
    """
    'CONTAS_RECEBER' -> 'contas_receber'.

    O servidor normaliza do mesmo jeito; normalizar aqui deixa o log do conector
    igual ao nome que aparece no painel e nas views `ds_<projeto>_<dataset>`.
    """
    loc_nome = re.sub(r"[^a-z0-9_]+", "_", str(p_nome).strip().lower())
    return loc_nome.strip("_")


# ------------------------------------------  Montar o envelope de um lote
def montar_envelope(p_dataset, p_registros, p_key=None, p_description=None, p_fields=None,
                    p_sync=None, p_complete=False, p_group=None):
    """
    Monta o envelope de UM lote. Os opcionais so entram quando ha valor -- o
    contrato aceita envelope minimo (`dataset` + `records`).

    `p_sync` marca o lote como parte de uma carga; `p_complete=True` marca este
    lote como o ULTIMO da carga daquele dataset (o que autoriza o servidor a
    remover o que nao foi visto).

    `complete` NUNCA e emitido sem `sync`: fechar uma carga que nao foi
    identificada nao tem significado no servidor, e emitir o campo sozinho so
    criaria a ilusao de que o ciclo esta ativo.
    """
    loc_envelope = {"dataset": normalizar_dataset(p_dataset)}

    if p_description:
        loc_envelope["description"] = p_description

    # O grupo classifica o dataset por assunto no Contextia. A taxonomia vive no
    # TENANT e a associacao e por NOME de dataset, entao mandar em todo lote e
    # idempotente. Grupo invalido NAO derruba a carga: o servidor ingere o dado
    # e ignora a classificacao.
    if p_group:
        loc_envelope["group"] = p_group

    if p_key:
        loc_envelope["key"] = p_key

    if p_fields:
        loc_envelope["fields"] = p_fields

    if p_sync:
        loc_envelope["sync"] = str(p_sync)

        # `complete` so existe no ultimo lote -- nos demais o campo simplesmente
        # nao aparece (nao vai `complete: false`).
        if p_complete:
            loc_envelope["complete"] = True

    loc_envelope["records"] = list(p_registros)

    return loc_envelope


# ------------------------------------------  Enviar um envelope
def enviar_envelope(p_envelope, p_api_key=L_API_KEY, p_url=L_INTAKE_URL, p_timeout=L_TIMEOUT,
                    p_logger=None):
    """
    POST /api/v1/intake/json com um envelope.

    A chave no header identifica tenant + projeto -- nao ha parametro de
    roteamento. Devolve o corpo da resposta (dict) ou None em caso de falha.

    ATENCAO ao ler o corpo: `accepted`/`duplicates` contam ITENS DE INTAKE, e o
    envelope gera UM item de auditoria por lote (os registros nao entram na fila
    um a um). Entao `accepted: 1` significa "lote aceito", nao "1 registro".
    Os numeros reais do lote estao no bloco `dataset` do corpo (`inserted`,
    `updated`, `duplicates`, `removed`, `resurrected`, `active_count`, ...) --
    e `enviar_dataset` ja os agrega.
    """
    loc_headers = {
        "Authorization": "Bearer " + p_api_key,
        "Content-Type":  "application/json",
    }

    loc_corpo = json.dumps(
        p_envelope, default=alchemyencoder, ensure_ascii=False
    ).encode("utf-8")

    try:
        response = requests.post(
            p_url, headers=loc_headers, data=loc_corpo, timeout=p_timeout
        )
    except Exception as e:
        if p_logger:
            p_logger.error("  [FALHA] dataset=%s | %s", p_envelope.get("dataset"), str(e))
        return None

    if response.status_code == 201:
        return response.json()

    # 401 = chave invalida/revogada; 400 = payload fora do contrato.
    if p_logger:
        p_logger.warning(
            "  [ERRO] dataset=%s | HTTP %s | %s",
            p_envelope.get("dataset"), response.status_code, response.text[:500],
        )
    return None


# ------------------------------------------  Quebrar o dataset em lotes
def fatiar(p_registros, p_chunk):
    """
    Lista de lotes a enviar.

    DATASET VAZIO: devolve UM lote vazio, nao nenhum. Uma view que ficou sem
    linhas na origem precisa igualmente do envelope de fechamento (`sync` +
    `complete` + `records: []`); sem ele o servidor nunca saberia que a carga
    daquele dataset aconteceu e os registros antigos ficariam vivos para sempre.
    """
    if not p_registros:
        return [[]]

    return [p_registros[i:i + p_chunk] for i in range(0, len(p_registros), p_chunk)]


# ------------------------------------------  Enviar um dataset inteiro
def enviar_dataset(p_dataset, p_registros, p_key=None, p_description=None, p_fields=None,
                   p_chunk=L_CHUNK, p_api_key=L_API_KEY, p_url=L_INTAKE_URL,
                   p_timeout=L_TIMEOUT, p_logger=None, p_sync=L_SYNC, p_envio=None,
                   p_group=None):
    """
    Envia um dataset inteiro como UMA CARGA: quebra os registros em lotes de
    `p_chunk`, manda todos com o mesmo `sync` e marca `complete: true` APENAS no
    ultimo lote.

    `p_sync`  -- identificador da carga. Padrao: `L_SYNC` (um por execucao do
                 processo). Passe `None` para DESLIGAR o ciclo: os lotes vao sem
                 `sync`/`complete` e nada e removido no servidor.
    `p_envio` -- funcao de transporte (assinatura de `enviar_envelope`). Serve
                 para teste com transporte falso e para dry-run.

    SEGURANCA DO CICLO (a regra mais importante deste arquivo): se QUALQUER lote
    do dataset falhou, o `complete` NAO e enviado -- nem que a falha tenha sido
    no primeiro lote e todos os seguintes tenham passado. Motivo: o servidor
    remove tudo que nao viu na carga, e uma carga incompleta "nao viu" registros
    que EXISTEM na origem e apenas nao chegaram. Sem `complete` a carga vira um
    upsert comum: o que chegou entra, nada e apagado, e a proxima execucao
    (carga nova, `sync` novo) fecha o ciclo corretamente.

    Um lote que falha tambem nao interrompe o dataset: com `key`, o reenvio e um
    upsert, entao os lotes seguintes continuam sendo enviados.

    Devolve o resumo da carga -- contagem local (`registros`, `lotes`,
    `lotes_ok`, `lotes_erro`, `sync`, `complete_enviado`) MAIS os numeros reais
    respondidos pela API (`inserted`, `updated`, `duplicates`, `removed`,
    `resurrected` somados lote a lote; `active_count`, `removed_count` e `view`
    lidos do ULTIMO lote bem-sucedido, porque sao o estado do dataset depois da
    carga e nao a soma dos lotes).
    """
    loc_envio = p_envio or enviar_envelope

    loc_res = {
        "dataset": normalizar_dataset(p_dataset),
        "sync": p_sync,
        "registros": 0, "lotes": 0, "lotes_ok": 0, "lotes_erro": 0,
        "complete_enviado": False,
        "inserted": 0, "updated": 0, "duplicates": 0,
        "removed": 0, "resurrected": 0,
        "active_count": None, "removed_count": None, "view": None,
    }

    loc_lotes = fatiar(p_registros, p_chunk)
    loc_ultimo_indice = len(loc_lotes) - 1

    for loc_i, loc_lote in enumerate(loc_lotes):
        # O ultimo lote so fecha a carga se NENHUM lote anterior falhou.
        loc_complete = (
            p_sync is not None
            and loc_i == loc_ultimo_indice
            and loc_res["lotes_erro"] == 0
        )

        loc_envelope = montar_envelope(
            p_dataset, loc_lote,
            p_key=p_key, p_description=p_description, p_fields=p_fields,
            p_sync=p_sync, p_complete=loc_complete, p_group=p_group,
        )

        loc_res["lotes"] += 1
        loc_corpo = loc_envio(
            loc_envelope, p_api_key=p_api_key, p_url=p_url,
            p_timeout=p_timeout, p_logger=p_logger,
        )

        if loc_corpo is None:
            loc_res["lotes_erro"] += 1

            if loc_complete and p_logger:
                p_logger.warning(
                    "  [CICLO] dataset=%s | o lote de fechamento (complete) falhou;"
                    " nada foi removido nesta carga.", loc_res["dataset"],
                )
            continue

        loc_res["lotes_ok"] += 1
        loc_res["registros"] += len(loc_lote)

        loc_ds = loc_corpo.get("dataset") if isinstance(loc_corpo, dict) else None
        if not isinstance(loc_ds, dict):
            # Servidor antigo (sem bloco `dataset`): vale o que foi enviado.
            if loc_complete:
                loc_res["complete_enviado"] = True
            continue

        if loc_complete:
            # A resposta ECOA `complete`: se vier False, o servidor nao tratou
            # este lote como fechamento (sync recusado, por exemplo) e nada foi
            # varrido -- mentir no resumo esconderia justamente isso.
            loc_res["complete_enviado"] = bool(loc_ds.get("complete", True))

        for loc_campo in ("inserted", "updated", "duplicates", "removed", "resurrected"):
            loc_res[loc_campo] += int(loc_ds.get(loc_campo) or 0)

        # Estado do dataset APOS o lote -- vale o do ultimo lote, nao a soma.
        # `record_count` e o nome antigo de `active_count` (servidor sem ciclo).
        loc_res["view"] = loc_ds.get("view", loc_res["view"])
        loc_res["active_count"] = loc_ds.get("active_count", loc_ds.get("record_count"))
        loc_res["removed_count"] = loc_ds.get("removed_count")

    if p_sync is not None and not loc_res["complete_enviado"] and p_logger:
        p_logger.warning(
            "  [CICLO] dataset=%s | carga INCOMPLETA (%s lote(s) com erro):"
            " `complete` nao enviado, nenhum registro removido.",
            loc_res["dataset"], loc_res["lotes_erro"],
        )

    return loc_res
