"""
Adaptador de origem GENERICO -- serve a maioria dos clientes sem escrever codigo.

    export CONTEXTIA_SOURCE_DRIVER=oracle
    export CONTEXTIA_SOURCE_HOST=10.0.0.5
    export CONTEXTIA_SOURCE_DATABASE=XEPDB1
    export CONTEXTIA_SOURCE_USER=leitura
    export CONTEXTIA_SOURCE_PASSWORD=...
    python3 connector.py --mapping=mappings/cliente.json --source=generic

Antes existia so um adaptador por cliente, que importava um `ssh_server.py` que
NAO esta neste repositorio -- o arquivo que a skill mandava copiar como modelo
nao rodava. Este aqui roda.

POR QUE PYTHON, E NAO PHP (que e a linguagem do Contextia)
---------------------------------------------------------
O PHP deste servidor nao fala Oracle: os drivers PDO sao mysql, odbc, pgsql e
sqlite. Habilitar oci8 exigiria instalar o Oracle Instant Client em toda maquina
onde a extracao rodasse. Em Python, `oracledb` em modo thin fala o protocolo
Oracle nativamente, sem client nenhum -- e o mesmo vale para MySQL, PostgreSQL e
SQL Server. Um cliente em Oracle APEX so e alcancavel sem
dor por este caminho.

O Contextia continua em PHP; ele so nao toca no banco do cliente.

TUNEL SSH
---------
Quando o banco nao esta exposto na rede -- o caso comum -- basta declarar o
salto:

    export CONTEXTIA_SSH_HOST=servidor.cliente.com.br
    export CONTEXTIA_SSH_USER=deploy
    export CONTEXTIA_SSH_KEY=~/.ssh/id_rsa_cliente

O tunel sobe sozinho, o banco e alcancado por 127.0.0.1 e tudo e fechado no
final.

DIAGNOSTICO
-----------
    python3 sources/generic.py

Conecta, lista alguns objetos e sai. E o primeiro comando a rodar num cliente
novo: separa "credencial errada" de "mapeamento errado" antes de a carga
comecar.
"""

import os
import re
import sys


# Nome de tabela/coluna que este adaptador aceita concatenar no SQL. O
# mapeamento e um arquivo de configuracao, nao uma fonte confiavel de SQL:
# quem escreve o mapeamento nao deve conseguir escrever a consulta.
IDENTIFICADOR = re.compile(r'^[A-Za-z_][A-Za-z0-9_$]*(\.[A-Za-z_][A-Za-z0-9_$]*)?$')

PORTA_PADRAO = {
    'oracle': 1521,
    'mysql': 3306,
    'mariadb': 3306,
    'postgres': 5432,
    'postgresql': 5432,
    'sqlserver': 1433,
    'mssql': 1433,
}

_conexao = None
_tunel = None
_driver = None


# Comando que trava a sessao em leitura no banco de origem. Este adaptador so
# monta SELECT -- a instrucao abaixo e a SEGUNDA barreira, para o caso de a
# credencial disponivel ter permissao de escrita (acontece: um ERP que so
# entrega o usuario dono do schema).
SOMENTE_LEITURA = {
    'oracle': 'SET TRANSACTION READ ONLY',
    'mysql': 'SET SESSION TRANSACTION READ ONLY',
    'mariadb': 'SET SESSION TRANSACTION READ ONLY',
    'postgres': 'SET SESSION CHARACTERISTICS AS TRANSACTION READ ONLY',
    'postgresql': 'SET SESSION CHARACTERISTICS AS TRANSACTION READ ONLY',
    # SQL Server nao tem equivalente de sessao; sobra a barreira do SELECT.
}


def _env(p_nome, p_padrao=None, p_obrigatorio=False):
    valor = os.environ.get(p_nome, p_padrao)
    if p_obrigatorio and not valor:
        raise SystemExit(
            'Falta a variavel de ambiente %s. Rode `python3 sources/generic.py`'
            ' para ver o diagnostico completo.' % p_nome
        )
    return valor


def _abrir_tunel(p_host, p_porta):
    """
    Sobe o tunel SSH quando configurado e devolve (host, porta) locais.

    Sem CONTEXTIA_SSH_HOST devolve o destino original -- conexao direta.
    """
    global _tunel

    ssh_host = os.environ.get('CONTEXTIA_SSH_HOST')
    if not ssh_host:
        return p_host, p_porta

    try:
        from sshtunnel import SSHTunnelForwarder  # type: ignore
    except ImportError:
        raise SystemExit('CONTEXTIA_SSH_HOST esta definido mas o pacote sshtunnel nao esta instalado.')

    chave = os.environ.get('CONTEXTIA_SSH_KEY')
    senha = os.environ.get('CONTEXTIA_SSH_PASSWORD')
    if not chave and not senha:
        raise SystemExit('Defina CONTEXTIA_SSH_KEY ou CONTEXTIA_SSH_PASSWORD para o tunel.')

    _tunel = SSHTunnelForwarder(
        (ssh_host, int(os.environ.get('CONTEXTIA_SSH_PORT', 22))),
        ssh_username=_env('CONTEXTIA_SSH_USER', p_obrigatorio=True),
        ssh_pkey=os.path.expanduser(chave) if chave else None,
        ssh_password=senha,
        remote_bind_address=(p_host, p_porta),
    )
    _tunel.start()

    return '127.0.0.1', _tunel.local_bind_port


def conectar():
    """Abre a conexao (uma so por execucao) e devolve o objeto do driver."""
    global _conexao, _driver

    if _conexao is not None:
        return _conexao

    driver = (_env('CONTEXTIA_SOURCE_DRIVER', p_obrigatorio=True) or '').lower()
    host = _env('CONTEXTIA_SOURCE_HOST', '127.0.0.1')
    porta = int(_env('CONTEXTIA_SOURCE_PORT', PORTA_PADRAO.get(driver, 0)) or 0)
    base = _env('CONTEXTIA_SOURCE_DATABASE', p_obrigatorio=True)
    usuario = _env('CONTEXTIA_SOURCE_USER', p_obrigatorio=True)
    senha = _env('CONTEXTIA_SOURCE_PASSWORD', '')

    host, porta = _abrir_tunel(host, porta)

    if driver == 'oracle':
        # Modo thin: fala o protocolo Oracle direto, sem Instant Client. E o que
        # torna um cliente em APEX alcancavel sem instalar nada.
        import oracledb  # type: ignore
        _conexao = oracledb.connect(
            user=usuario, password=senha,
            dsn='%s:%d/%s' % (host, porta, base),
        )
    elif driver in ('mysql', 'mariadb'):
        import pymysql  # type: ignore
        _conexao = pymysql.connect(
            host=host, port=porta, user=usuario, password=senha,
            database=base, cursorclass=pymysql.cursors.Cursor,
        )
    elif driver in ('postgres', 'postgresql'):
        import pg8000.dbapi  # type: ignore
        _conexao = pg8000.dbapi.connect(
            host=host, port=porta, user=usuario, password=senha, database=base,
        )
    elif driver in ('sqlserver', 'mssql'):
        import pytds  # type: ignore
        _conexao = pytds.connect(
            dsn=host, port=porta, user=usuario, password=senha, database=base,
        )
    else:
        raise SystemExit(
            'Driver %r nao suportado. Use oracle, mysql, postgres ou sqlserver.'
            ' (Firebird exige a biblioteca fbclient nativa e nao entra aqui.)' % driver
        )

    _driver = driver
    _somente_leitura(_conexao, driver)

    return _conexao


def _somente_leitura(p_conexao, p_driver):
    """
    Trava a sessao em leitura, quando o banco permite.

    Nao e paranoia: o conector roda com a credencial que o cliente entregou, e
    nem sempre ela e a que pedimos. Ja aconteceu de o unico usuario ser o DONO do
    schema. Um script de diagnostico mal colado nessa sessao viraria incidente
    na producao do cliente.

    No Oracle isto tem um segundo efeito, desejavel: a transacao passa a ver um
    retrato unico do banco, entao todos os datasets de uma carga saem do MESMO
    instante -- contas a pagar nao fica de um momento e o caixa de outro.
    """
    comando = SOMENTE_LEITURA.get(p_driver)
    if not comando:
        return

    cursor = p_conexao.cursor()
    try:
        cursor.execute(comando)
    except Exception as e:
        # Falhar aqui nao autoriza escrever: quem garante isso e o fato de este
        # modulo so montar SELECT. Mas o operador precisa saber que a rede de
        # protecao nao subiu.
        sys.stderr.write(
            'AVISO: nao foi possivel marcar a sessao como somente leitura'
            ' (%s: %s). A carga segue -- este adaptador nunca monta nada alem'
            ' de SELECT --, mas a barreira extra nao esta ativa.\n'
            % (type(e).__name__, e)
        )
    finally:
        cursor.close()


def _onde(p_sql, p_scope, p_since):
    """
    Acrescenta o WHERE a um SELECT e devolve (sql, parametros).

    Existe para que `fetch`, `contar` e `maximo` nao possam FILTRAR DIFERENTE.
    Ja seria ruim em qualquer lugar; aqui seria pior: a sonda contaria uma coisa
    e a leitura traria outra, e o ouvinte avancaria a marca de agua sobre linhas
    que nunca leu.

    O valor vai sempre como parametro do driver; a coluna e validada como
    identificador antes de ser concatenada.
    """
    condicoes = []

    if p_scope is not None:
        coluna, valor = p_scope
        if not IDENTIFICADOR.match(coluna):
            raise ValueError('Coluna de recorte invalida no mapeamento: %r' % coluna)
        condicoes.append((coluna, '=', 'escopo', valor))

    if p_since is not None:
        coluna, valor = p_since
        if not IDENTIFICADOR.match(coluna):
            raise ValueError('Coluna de mudanca invalida no mapeamento: %r' % coluna)
        # `>` e nao `>=`: com `>=` a leitura seguinte reenviaria sempre a ultima
        # linha. Reenviar e barato (volta como `iguais`), mas o log ficaria com
        # trafego eterno mesmo com a origem parada, e ninguem distinguiria isso
        # de movimento de verdade. A folga contra atraso de commit e do ouvinte,
        # que recua a marca antes de chamar aqui -- ver ouvinte.py.
        condicoes.append((coluna, '>', 'desde', valor))

    if not condicoes:
        return p_sql, ()

    conectar()  # garante que _driver esta resolvido

    partes = []
    nomeados = {}
    posicionais = []
    for coluna, operador, nome, valor in condicoes:
        if _driver == 'oracle':
            partes.append('%s %s :%s' % (coluna, operador, nome))
            nomeados[nome] = valor
        else:
            # pymysql, pg8000 e python-tds usam paramstyle `format`.
            partes.append('%s %s %%s' % (coluna, operador))
            posicionais.append(valor)

    return (p_sql + ' WHERE ' + ' AND '.join(partes),
            nomeados if _driver == 'oracle' else tuple(posicionais))


def fetch(p_object, p_columns, scope=None, since=None):
    """
    Le o objeto e devolve list[dict] com as chaves na grafia do mapeamento.

    `scope` e o recorte declarado no mapeamento: uma tupla (coluna, valor) que
    vira `WHERE <coluna> = <valor>`. Serve para a base multi-empresa, em que um
    schema so guarda varios clientes e cada um tem de virar um projeto separado
    no Contextia. O VALOR vai como parametro do driver, nunca concatenado; a
    COLUNA e validada como identificador, igual as demais.

    `since` e a tupla (coluna, valor) do OUVINTE: vira `WHERE <coluna> > <valor>`
    e le so o que mudou desde a ultima leitura. Combina com `scope` -- os dois
    juntos saem como `WHERE empresa = ? AND alterado_em > ?`. Sem `since` o SQL
    montado e byte a byte o de antes deste parametro existir.

    Contrato que o connector.py assume:

      * valor ausente vem como None, NUNCA a string "null". O None vira `null`
        no JSON e chega ao Contextia como NULL de verdade;
      * erro de leitura SOBE. Quem decide o que fazer e o conector, que pula o
        dataset inteiro sem fechar a carga -- devolver [] aqui faria o servidor
        remover tudo que existe na origem.
    """
    if not IDENTIFICADOR.match(p_object):
        raise ValueError('Nome de objeto invalido no mapeamento: %r' % p_object)

    for coluna in p_columns:
        if not IDENTIFICADOR.match(coluna):
            raise ValueError('Nome de coluna invalido no mapeamento: %r' % coluna)

    sql, parametros = _onde(
        'SELECT %s FROM %s' % (', '.join(p_columns), p_object), scope, since)

    cursor = conectar().cursor()
    try:
        cursor.execute(sql, parametros)

        # As chaves saem com a grafia DO MAPEAMENTO, nao a que o driver devolve.
        #
        # Forcar maiuscula aqui (o que este adaptador fazia) quebra um
        # mapeamento escrito em minusculas: os REGISTROS iam com "ID" e os
        # `fields` do envelope com "id", o dataset acumulava as duas grafias e a
        # view saia com coluna duplicada -- nome de coluna no MySQL nao
        # distingue maiuscula. O servidor respondia 500.
        #
        # A projecao foi construida a partir de p_columns, na ordem, entao casar
        # posicao por posicao e exato -- e vale para o Oracle, que devolve tudo
        # em maiuscula independentemente de como foi pedido.
        return [dict(zip(p_columns, linha)) for linha in cursor.fetchall()]
    finally:
        cursor.close()


def contar(p_object, scope=None, since=None):
    """
    `COUNT(*)` com os mesmos filtros do `fetch`.

    E a SONDA do ouvinte, e a razao de ela existir e economia na origem: na
    grande maioria dos ciclos nada mudou, e uma contagem por dataset e muito
    mais barata que trazer as linhas para descobrir que sao zero. Tambem e o
    que permite RECUSAR um ciclo grande demais antes de ele entrar na memoria
    -- uma alteracao em massa no ERP (recalculo de juros, fechamento de mes)
    devolveria centenas de milhares de linhas de uma vez.
    """
    if not IDENTIFICADOR.match(p_object):
        raise ValueError('Nome de objeto invalido no mapeamento: %r' % p_object)

    sql, parametros = _onde('SELECT COUNT(*) FROM %s' % p_object, scope, since)

    cursor = conectar().cursor()
    try:
        cursor.execute(sql, parametros)
        linha = cursor.fetchone()
        return int(linha[0]) if linha and linha[0] is not None else 0
    finally:
        cursor.close()


def renovar():
    """
    Descarta o retrato do banco e abre um novo, SEM derrubar a conexao nem o tunel.

    ISTO NAO E OPCIONAL NUM PROCESSO QUE VIVE. E o defeito que faria um ouvinte
    parecer perfeito e nunca mais ver dado novo:

      * no Oracle, `SET TRANSACTION READ ONLY` vale ate o commit/rollback. Sem
        renovar, a sessao fica presa ao retrato do instante em que subiu -- o
        ouvinte rodaria para sempre lendo o banco de ontem, sem erro nenhum;
      * no MySQL/InnoDB o efeito e o mesmo por outro caminho: sob REPEATABLE
        READ, o primeiro SELECT abre o retrato e todos os seguintes o reusam.

    Nao ha erro, nao ha aviso e o log fica cheio de "nada mudou" -- que e
    exatamente o que se espera de uma origem parada. O sintoma so aparece
    quando alguem compara o Contextia com o ERP, dias depois.

    Chamar isto entre ciclos e o que separa "ouvinte" de "leitor do passado".
    """
    global _conexao

    if _conexao is None:
        return

    try:
        _conexao.rollback()
    except Exception:
        # Conexao morta (rede, tunel, timeout do banco). Derruba para o proximo
        # ciclo reconectar -- insistir num handle morto daria o mesmo erro para
        # sempre.
        fechar()
        return

    _somente_leitura(_conexao, _driver)


def maximo(p_object, p_column, scope=None):
    """
    `MAX(<coluna>)` do objeto, com o recorte de empresa se houver.

    E como o ouvinte descobre onde ele COMECA. A alternativa seria comecar do
    relogio da maquina que roda a extracao, e ai um relogio adiantado em cinco
    minutos faria o ouvinte pular cinco minutos de movimento -- calado, e so na
    primeira execucao, que e justamente quando ninguem esta conferindo.

    Devolve None quando o objeto esta vazio; quem chama trata isso como "nada
    para acompanhar ainda", nunca como zero.
    """
    if not IDENTIFICADOR.match(p_object):
        raise ValueError('Nome de objeto invalido no mapeamento: %r' % p_object)
    if not IDENTIFICADOR.match(p_column):
        raise ValueError('Nome de coluna invalido no mapeamento: %r' % p_column)

    sql, parametros = _onde('SELECT MAX(%s) FROM %s' % (p_column, p_object), scope, None)

    cursor = conectar().cursor()
    try:
        cursor.execute(sql, parametros)
        linha = cursor.fetchone()
        return linha[0] if linha else None
    finally:
        cursor.close()


def fechar():
    global _conexao, _tunel

    if _conexao is not None:
        try:
            _conexao.close()
        except Exception:
            pass
        _conexao = None

    if _tunel is not None:
        try:
            _tunel.stop()
        except Exception:
            pass
        _tunel = None


def _diagnostico():
    """`python3 sources/generic.py` -- separa erro de acesso de erro de mapeamento."""
    driver = (os.environ.get('CONTEXTIA_SOURCE_DRIVER') or '').lower()

    print('Driver ...: %s' % (driver or '(nao definido)'))
    print('Host .....: %s:%s' % (
        os.environ.get('CONTEXTIA_SOURCE_HOST', '127.0.0.1'),
        os.environ.get('CONTEXTIA_SOURCE_PORT', PORTA_PADRAO.get(driver, '?')),
    ))
    print('Base .....: %s' % (os.environ.get('CONTEXTIA_SOURCE_DATABASE') or '(nao definida)'))
    print('Usuario ..: %s' % (os.environ.get('CONTEXTIA_SOURCE_USER') or '(nao definido)'))
    print('Tunel ....: %s' % (os.environ.get('CONTEXTIA_SSH_HOST') or 'nao'))
    print()

    inventario = {
        'oracle': "SELECT view_name FROM user_views ORDER BY view_name FETCH FIRST 15 ROWS ONLY",
        'mysql': 'SHOW FULL TABLES',
        'mariadb': 'SHOW FULL TABLES',
        'postgres': "SELECT table_name FROM information_schema.tables WHERE table_schema='public' ORDER BY table_name LIMIT 15",
        'postgresql': "SELECT table_name FROM information_schema.tables WHERE table_schema='public' ORDER BY table_name LIMIT 15",
        'sqlserver': 'SELECT TOP 15 name FROM sys.objects WHERE type IN (%s) ORDER BY name' % "'U','V'",
        'mssql': 'SELECT TOP 15 name FROM sys.objects WHERE type IN (%s) ORDER BY name' % "'U','V'",
    }

    try:
        cursor = conectar().cursor()
        cursor.execute(inventario.get(driver, 'SELECT 1'))
        linhas = cursor.fetchall()
        print('CONECTOU. Primeiros objetos visiveis para este usuario:')
        for linha in linhas[:15]:
            print('  %s' % (linha[0],))
        if not linhas:
            print('  (nenhum -- o usuario conecta mas nao enxerga objeto nenhum:')
            print('   quase sempre falta GRANT SELECT, ou o schema e outro)')
        cursor.close()
        return 0
    except Exception as e:
        print('NAO CONECTOU: %s: %s' % (type(e).__name__, e))
        print()
        print('Isto e problema de ACESSO (rede, credencial, tunel), nao de mapeamento.')
        return 1
    finally:
        fechar()


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