#!/usr/bin/env python3
"""
Conector generico Contextia -- o MESMO arquivo para todo cliente.

    python3 connector.py --mapping=mappings/<cliente>.json --source=generic

O que era um script por cliente virou tres pecas com donos diferentes:

    connector.py            este arquivo. Nunca muda de cliente para cliente.
    mappings/<cliente>.json O QUE extrair e o que cada coisa SIGNIFICA.
                            Declarativo, sem codigo. E o artefato que a skill
                            de onboarding gera lendo o schema e a documentacao,
                            e o mesmo que o agente instalado vai consumir
                            quando existir -- por isso ele nao pode ter logica
                            dentro.
    sources/<cliente>.py    COMO alcancar a origem: tunel SSH, driver, DSN.
                            Isso e ambiente, muda de cliente para cliente e nao
                            tem como ser declarativo. ~20 linhas.

O adaptador de origem precisa expor uma funcao so:

    fetch(object_name, columns) -> list[dict]

Erro de leitura ele deixa SUBIR (o conector trata); lista vazia significa que a
origem leu e nao tem nada, que e coisa completamente diferente -- ver abaixo.

CICLO DE CARGA (`sync` + `complete`)
------------------------------------
O upsert por `key` cobre inclusao e alteracao, mas nao EXCLUSAO: um titulo
cancelado no ERP continuaria vivo no Contextia e o agente somaria dinheiro que
nao existe mais. Por isso cada execucao e uma CARGA identificada: todos os lotes
levam o mesmo `sync`, e so o ULTIMO lote de cada dataset leva `complete: true`
-- e so entao o servidor marca como removido o que nao apareceu.

A distincao que protege o dado do cliente:

    erro de leitura  -> o dataset e PULADO INTEIRO, nem o fechamento vai.
                        Nao sabemos o que a origem tem; remover seria inventar.
    zero linhas      -> o fechamento VAI, com `records: []`. A origem esvaziou
                        de verdade e o dataset deve esvaziar junto.

E a mesma regra dentro do lote: se qualquer lote falhou, `enviar_dataset` nao
manda o `complete`, a carga vira upsert comum e a proxima execucao completa o
servico. Sem isso, uma queda de rede viraria perda de dado.
"""

import argparse
import inspect
import importlib
import json
import os
import sys

from logzero import logger  # type: ignore
from tqdm import tqdm  # type: ignore

from contextia_client import (
    L_API_KEY,
    L_INTAKE_URL,
    L_SYNC,
    enviar_dataset,
    normalizar_dataset,
)


# ---------------------------------------------------------------------------
# Mapeamento
# ---------------------------------------------------------------------------

def carregar_mapeamento(p_path):
    """Le e valida o mapeamento. Erro aqui aborta ANTES de tocar a origem."""
    with open(p_path, encoding='utf-8') as f:
        loc_map = json.load(f)

    if loc_map.get('mapping_version') != 1:
        raise SystemExit(
            'Mapeamento em versao %r; este conector le a versao 1.'
            % loc_map.get('mapping_version')
        )

    loc_datasets = loc_map.get('datasets') or []
    if not loc_datasets:
        raise SystemExit('Mapeamento sem nenhum dataset em `datasets`.')

    loc_vistos = set()
    for row in loc_datasets:
        for campo in ('dataset', 'source_object', 'columns'):
            if not row.get(campo):
                raise SystemExit(
                    'Dataset %r sem `%s` no mapeamento.'
                    % (row.get('dataset') or row.get('source_object'), campo)
                )

        # Nome repetido faria a segunda carga sobrescrever a primeira no mesmo
        # `sync` -- e a primeira sairia como "removida" no fechamento.
        loc_nome = normalizar_dataset(row['dataset'])
        if loc_nome in loc_vistos:
            raise SystemExit('Dataset %r aparece duas vezes no mapeamento.' % loc_nome)
        loc_vistos.add(loc_nome)

    return loc_map


def coluna_de_recorte(p_map, p_row):
    """
    Qual coluna recorta ESTE dataset, segundo o mapeamento.

    O mapeamento declara o recorte, nunca o valor: `scope_column` diz QUAL
    coluna separa um cliente do outro, e o `--scope` da linha de comando diz
    QUAL cliente esta sendo carregado agora. Assim um unico mapeamento serve as
    N empresas de uma base compartilhada, e a chave de API -- que e o que
    escolhe tenant e projeto -- continua sendo o unico roteamento.

    `scope_column` ausente herda o `defaults`; `scope_column: null` no dataset
    desliga o recorte, e e o que se usa em tabela de dominio (status, tipos)
    que vale para todas as empresas.
    """
    if 'scope_column' in p_row:
        return p_row['scope_column']
    return (p_map.get('defaults') or {}).get('scope_column')


def exigir_suporte_a_escopo(p_fetch, p_escopo):
    """
    Aborta se um recorte foi pedido e o adaptador nao sabe recebe-lo.

    Nao e preciosismo: um adaptador que ignora o recorte em silencio carrega as
    33 empresas da base dentro do projeto de UMA, e o agente daquele cliente
    passa a responder sobre o financeiro dos outros 32. E o pior defeito
    possivel aqui -- nao da erro, so vaza. Por isso e SystemExit, e nao aviso.
    """
    if p_escopo is None:
        return

    try:
        inspect.signature(p_fetch).bind('x', ['y'], scope=('c', 'v'))
    except TypeError:
        raise SystemExit(
            'Foi pedido --scope=%s, mas o adaptador de origem nao aceita'
            ' `scope`. Carregar assim colocaria TODAS as empresas da base no'
            ' projeto desta chave. Use sources/generic.py, ou acrescente o'
            ' parametro `scope=None` ao fetch do adaptador.' % p_escopo
        )


def ler_origem(p_fetch, p_map, p_row, p_escopo):
    """Le um dataset ja aplicando o recorte que o mapeamento declara para ele."""
    loc_coluna = coluna_de_recorte(p_map, p_row)

    if p_escopo is None or not loc_coluna:
        return p_fetch(p_row['source_object'], p_row['columns'])

    if loc_coluna not in p_row['columns']:
        raise SystemExit(
            'Dataset %r recorta por %r, mas essa coluna nao esta em `columns`.'
            ' Sem ela na projecao o recorte nao pode ser conferido depois.'
            % (p_row['dataset'], loc_coluna)
        )

    return p_fetch(p_row['source_object'], p_row['columns'],
                   scope=(loc_coluna, p_escopo))


def semantica(p_map, p_row):
    """
    Descricao e campos de um dataset: os campos padrao do mapeamento mesclados
    com os especificos, e os especificos ganham.

    So descreve campo que a view REALMENTE traz -- descrever coluna que nao
    existe polui o catalogo e engana o agente.
    """
    loc_defaults = (p_map.get('defaults') or {}).get('fields') or {}
    loc_todos = dict(loc_defaults)
    loc_todos.update(p_row.get('fields') or {})

    loc_fields = {}
    for col in p_row['columns']:
        if col in loc_todos:
            loc_fields[col] = loc_todos[col]

    return p_row.get('description'), loc_fields


# ---------------------------------------------------------------------------
# Execucao de um dataset
# ---------------------------------------------------------------------------

def processar(p_map, p_row, p_fetch, p_escopo=None):
    """Devolve (registros_enviados, lotes_com_erro, removidos)."""
    loc_objeto = p_row['source_object']
    loc_dataset = normalizar_dataset(p_row['dataset'])

    logger.info('Consultando %s ...', loc_objeto)

    try:
        loc_registros = ler_origem(p_fetch, p_map, p_row, p_escopo)
    except Exception as e:  # noqa: BLE001 -- qualquer falha de leitura vale
        logger.error('  [ERRO] %s: %s', loc_objeto, e)
        logger.error('  Carga de %s ABORTADA (nada enviado, nada removido).', loc_dataset)
        return 0, 1, 0

    loc_description, loc_fields = semantica(p_map, p_row)

    if not loc_registros:
        logger.warning(
            '  %s respondeu ZERO linhas. Enviando so o fechamento:'
            ' o dataset %s sera esvaziado no Contextia.', loc_objeto, loc_dataset,
        )

    # `key` ausente faz o servidor deduplicar por hash do conteudo. So mandamos
    # a chave se a coluna estiver mesmo na lista -- prometer uma chave que nao
    # vem no registro quebraria o upsert.
    loc_key = p_row.get('key') or (p_map.get('defaults') or {}).get('key')
    if loc_key not in p_row['columns']:
        loc_key = None

    loc_res = enviar_dataset(
        loc_dataset, loc_registros,
        p_key=loc_key,
        p_description=loc_description,
        p_fields=loc_fields,
        p_group=p_row.get('group'),
        p_chunk=p_row.get('chunk') or (p_map.get('defaults') or {}).get('chunk') or 500,
        p_logger=logger,
        p_sync=L_SYNC,
    )

    logger.info(
        '  [OK] dataset=%s | lidos=%s | enviados=%s | lotes=%s (erro=%s)'
        ' | novos=%s atualizados=%s iguais=%s removidos=%s ressuscitados=%s'
        ' | ativos=%s | carga fechada=%s | campos descritos=%s',
        loc_dataset, len(loc_registros), loc_res['registros'],
        loc_res['lotes'], loc_res['lotes_erro'],
        loc_res['inserted'], loc_res['updated'], loc_res['duplicates'],
        loc_res['removed'], loc_res['resurrected'],
        loc_res['active_count'], 'sim' if loc_res['complete_enviado'] else 'NAO',
        len(loc_fields),
    )

    return loc_res['registros'], loc_res['lotes_erro'], loc_res['removed']


# ---------------------------------------------------------------------------
# Principal
# ---------------------------------------------------------------------------

def main():
    p = argparse.ArgumentParser(description='Conector generico Contextia.')
    p.add_argument('--mapping', required=True, help='caminho do mappings/<cliente>.json')
    p.add_argument('--source', required=True, help='nome do modulo em sources/ (sem .py)')
    p.add_argument('--only', default=None,
                   help='carrega so estes datasets (nomes separados por virgula)')
    p.add_argument('--scope', default=None,
                   help='valor do recorte declarado em `scope_column` (ex.: o'
                        ' EMPRESA_ID, numa base com varias empresas)')
    p.add_argument('--dry-run', action='store_true',
                   help='le a origem e mostra o que iria, sem enviar nada')
    args = p.parse_args()

    sys.path.insert(0, os.path.join(os.path.dirname(os.path.abspath(__file__)), 'sources'))

    loc_map = carregar_mapeamento(args.mapping)

    try:
        loc_source = importlib.import_module(args.source)
    except ImportError as e:
        raise SystemExit('Nao consegui carregar sources/%s.py: %s' % (args.source, e))

    if not hasattr(loc_source, 'fetch'):
        raise SystemExit('sources/%s.py precisa expor fetch(object_name, columns).' % args.source)

    exigir_suporte_a_escopo(loc_source.fetch, args.scope)

    # Empresa marcada como teste no mapeamento nao carrega. A lista e aplicada
    # aqui, e nao apenas documentada: dado de demonstracao entra na mesma camada
    # que o real e o agente cita os dois com a mesma confianca.
    loc_excluidos = (loc_map.get('source') or {}).get('excluded_scopes') or {}
    if args.scope and str(args.scope) in loc_excluidos:
        raise SystemExit(
            'O escopo %s esta marcado como excluido no mapeamento: %s\n'
            'Se isso mudou, tire-o de `source.excluded_scopes`.'
            % (args.scope, loc_excluidos[str(args.scope)])
        )

    loc_datasets = loc_map['datasets']
    if args.only:
        loc_pedidos = {normalizar_dataset(x) for x in args.only.split(',') if x.strip()}
        loc_datasets = [d for d in loc_datasets
                        if normalizar_dataset(d['dataset']) in loc_pedidos]
        loc_faltando = loc_pedidos - {normalizar_dataset(d['dataset']) for d in loc_datasets}
        if loc_faltando:
            raise SystemExit('Nao existem no mapeamento: %s' % ', '.join(sorted(loc_faltando)))

    logger.info('Cliente ..: %s', loc_map.get('client') or '(sem nome)')
    logger.info('Mapeamento: %s (%s datasets)', args.mapping, len(loc_datasets))
    logger.info('Origem ...: sources/%s.py', args.source)
    if args.scope:
        logger.info('Recorte ..: %s = %s -- so estes registros entram nesta chave',
                    (loc_map.get('defaults') or {}).get('scope_column') or '?',
                    args.scope)

    if args.dry_run:
        for row in loc_datasets:
            loc_desc, loc_fields = semantica(loc_map, row)
            try:
                loc_n = len(ler_origem(loc_source.fetch, loc_map, row, args.scope))
            except Exception as e:  # noqa: BLE001
                print('%-30s ERRO na leitura: %s' % (row['dataset'], e))
                continue
            print('%-30s %7s linhas  grupo=%-14s campos descritos=%s/%s  recorte=%s'
                  % (row['dataset'], loc_n, row.get('group') or '-',
                     len(loc_fields), len(row['columns']),
                     coluna_de_recorte(loc_map, row) or 'global'))
        print('\nNADA foi enviado (--dry-run).')
        return 0

    from contextia_client import exigir_destino
    exigir_destino()          # antes de qualquer byte sair

    logger.info('Enviando para %s', L_INTAKE_URL)
    logger.info('Chave: %s... (identifica tenant e projeto no Contextia)', L_API_KEY[:12])
    logger.info('Carga (sync): %s -- o mesmo em todos os lotes desta execucao', L_SYNC)

    loc_total_registros = 0
    loc_total_erros = 0
    loc_total_removidos = 0

    for row in tqdm(loc_datasets):
        registros, erros, removidos = processar(loc_map, row, loc_source.fetch, args.scope)
        loc_total_registros += registros
        loc_total_erros += erros
        loc_total_removidos += removidos

    print('|----------------------------------------|')
    print('Processo finalizado!')
    print('  Carga (sync) ........: %s' % L_SYNC)
    print('  Datasets ............: %s' % len(loc_datasets))
    print('  Registros enviados ..: %s' % loc_total_registros)
    print('  Registros removidos .: %s' % loc_total_removidos)
    print('  Lotes com erro ......: %s' % loc_total_erros)
    print('|----------------------------------------|')

    if loc_total_removidos:
        print('ATENCAO: %s registro(s) foram marcados como REMOVIDOS nesta carga --'
              % loc_total_removidos)
        print('estavam no Contextia e nao vieram mais da origem. Numero alto sem')
        print('motivo de negocio quase sempre e problema na ORIGEM (view quebrada,')
        print('filtro errado, banco em manutencao). Eles somem das views, mas nao sao')
        print('apagados: voltam ressuscitados assim que reaparecerem numa carga.')

    if loc_total_erros:
        print('ATENCAO: %s lote(s) falharam. Os datasets afetados NAO tiveram a carga'
              % loc_total_erros)
        print('fechada (`complete` nao enviado), entao nada foi removido neles.')
        print('Rode o conector de novo para completar.')

    return 1 if loc_total_erros else 0


if __name__ == '__main__':
    sys.exit(main())
