#!/usr/bin/env python3
"""
OUVINTE -- o dado chega ao Contextia em segundos, nao na madrugada seguinte.

    python3 ouvinte.py --mapping=mappings/<cliente>.json --source=generic \\
        --chaves=/opt/contextia/chaves.tsv --estado=/opt/contextia/estado

Um processo que fica de pe e faz, sozinho, TUDO o que o cron da recarga fazia --
mais o que ele nao conseguia fazer: acompanhar a origem de perto.

    ciclo de DADOS       a cada 60s     o que mudou desde a ultima leitura
    ciclo de PROJETOS    a cada 300s    empresa nova ganha projeto, chave e carga
    VARREDURA            1x por dia     a carga completa -- e a UNICA que remove

Ele nao reimplementa nada disso: o ciclo de projetos chama o `provisionar.py` e a
varredura chama o `carregar-todas.sh`, os mesmos que o cron chamava, com os
mesmos argumentos. O que este arquivo acrescenta e o ciclo de dados, o relogio,
a trava e o estado. Toda a parte perigosa -- criar chave, entregar token
assinado, fechar carga -- continua morando onde ja estava testada.

O cron continua valendo, e para muito cliente ele e o certo. Este arquivo e para
quando a pergunta *"o titulo que eu baixei agora aparece?"* precisa ser "sim".


O QUE "AO VIVO" QUER DIZER AQUI, EXATAMENTE
===========================================
Nao e CDC de log de transacao. E uma SONDA CURTA numa coluna de mudanca:

    SELECT COUNT(*) FROM <view> WHERE <coluna_de_mudanca> > <marca>

E a escolha e deliberada, nao uma limitacao aceita por comodidade. As
alternativas, e por que nenhuma serve a um SDK que roda em qualquer cliente:

    binlog / LogMiner / logical decoding
        A latencia ideal, e cada banco e um bicho diferente. Exige privilegio
        de replicacao, configuracao no servidor do cliente e, no Oracle,
        frequentemente produto licenciado. Nao ha "um codigo" que atenda
        Oracle, MySQL, Postgres e SQL Server por esse caminho.

    trigger + tabela de outbox na origem
        Funciona bem e tem latencia otima -- e exige DDL e ESCRITA no banco do
        cliente. A Etapa 0 deste SDK pede usuario SOMENTE LEITURA, e o
        adaptador trava a sessao em leitura ao conectar. Um recurso que obriga a
        desfazer isso nao pode ser o padrao.

    varredura completa curta
        Simples e correta, e cara: 27 datasets vezes 11 empresas a cada minuto e
        a producao do cliente pagando a conta.

A sonda e o unico dos quatro que roda com `SELECT` e credencial de leitura, em
qualquer um dos quatro bancos, sem pedir nada ao DBA do cliente. Latencia: o
intervalo do ciclo, tipicamente **de 30 a 60 segundos**.

Quando o cliente PUDER oferecer trigger+outbox ou CDC nativo, o caminho e
apontar `source_object` para a view/tabela que ele expuser: o ouvinte nao muda,
porque para ele continua sendo uma coluna de mudanca. E isso vale a
recomendacao inteira -- a regra fica na casa do cliente, e nos continuamos so
lendo.


O QUE ELE NAO VE: EXCLUSAO
==========================
Diga isto ao cliente antes de ligar, porque e a unica surpresa possivel.

Linha APAGADA na origem nao tem coluna de mudanca -- ela simplesmente deixa de
existir, e nenhum `WHERE alterado > x` a encontra. E, do outro lado, o intake do
Contextia nao tem como receber "remova este registro": remocao acontece so pelo
CICLO DE CARGA (`sync` + `complete`), que marca como removido tudo que nao
apareceu na carga.

Entao:

    INCLUSAO e ALTERACAO      em segundos, pelo ciclo de dados
    EXCLUSAO                  na VARREDURA, 1x por dia

Na pratica isso incomoda menos do que parece, e vale conferir no cliente: ERP
raramente apaga titulo -- ele marca `SITUACAO = 'C'`, que e uma ALTERACAO e o
ouvinte pega em segundos. Onde houver exclusao fisica de verdade, a saida e
aumentar a frequencia da varredura naquele dataset (a carga completa e
idempotente e segura de repetir), ou pedir ao cliente uma coluna de cancelamento
em vez do DELETE.

O que consertaria isso de vez e um campo `remove: [<chaves>]` no envelope do
intake, que HOJE NAO EXISTE. Enquanto nao existir, nao prometa exclusao ao vivo.


A MARCA DE AGUA, e os tres jeitos de perder dado com ela
========================================================
Cada dataset guarda a maior mudanca que ja foi lida. Tres cuidados, todos
implementados aqui, e cada um deles ja e um jeito conhecido de perder linha em
silencio:

  1. A marca avanca para o MAIOR VALOR LIDO, nunca para o relogio da maquina.
     Um relogio adiantado em cinco minutos pularia cinco minutos de movimento.

  2. A leitura RECUA a marca em `--folga` segundos (padrao 180). Uma transacao
     que comecou antes e comitou depois grava um timestamp ANTERIOR a marca que
     ja avancou -- sem a folga, aquela linha nunca seria vista. Reler e barato:
     `key` faz upsert, e a linha repetida volta como `iguais`.

  3. A marca so avanca se TODOS os envios do ciclo derem certo. Falha de rede
     mantem a marca velha, e o ciclo seguinte rele. E a mesma regra do
     `complete` do conector: falha de carga nunca vira perda de dado.

Coluna de mudanca com granularidade de DIA (um `DATE` sem hora) funciona, mas
recua um dia inteiro -- o ouvinte rele o movimento de hoje a cada ciclo. Custa
trafego e nao perde nada. Coluna numerica (sequence, `ROWVERSION`) nao recua:
para ela a folga nao tem unidade, e quem cobre transacao em voo e a varredura.


ONDE FICA O QUE
===============
    mapeamento      `watch_column` -- e uma coluna da ORIGEM, entao e mapeamento
    linha de comando  intervalos, folga, teto -- e ajuste de IMPLANTACAO
    estado          `<estado>/ouvinte-estado.json`  marcas de agua e escopos vistos
    saude           `<estado>/ouvinte-status.txt`   o arquivo que um alerta le
    trava           `<estado>/ouvinte.lock`         um ouvinte por implantacao

O status e um arquivo e nao um log de propósito: alerta que precisa varrer log
para saber se o servico esta vivo nao e alerta. Se a ultima linha dele estiver
velha, o ouvinte parou -- e isso da para checar com um `find -mmin`.
"""

import argparse
import datetime
import importlib
import io
import json
import os
import subprocess
import sys
import time


def _sem_logzero():
    """
    Substitutos de `logzero` e `tqdm` quando eles nao estao instalados.

    Existe por causa do `--autoteste`: ele foi feito para rodar na maquina de
    quem esta escrevendo o mapeamento, ANTES de instalar qualquer coisa e antes
    de tocar na origem do cliente. Um `ModuleNotFoundError` ali transformaria a
    unica prova que da para rodar de graca numa prova que exige preparar
    ambiente -- e prova que da trabalho nao e rodada.

    So entra em acao quando o pacote de verdade falta. Com ele instalado (o
    normal, via requirements.txt) nada aqui e usado, e o log continua sendo o
    do logzero, com cor e arquivo rotativo.
    """
    import logging
    import types

    logging.basicConfig(level=logging.INFO,
                        format='[%(levelname).1s %(asctime)s] %(message)s')

    loc_logzero = types.ModuleType('logzero')
    loc_logzero.logger = logging.getLogger('contextia')

    def _logfile(p_path, **p_kwargs):
        from logging.handlers import RotatingFileHandler
        loc_h = RotatingFileHandler(p_path, maxBytes=p_kwargs.get('maxBytes', 0),
                                    backupCount=p_kwargs.get('backupCount', 0))
        loc_h.setFormatter(logging.Formatter('[%(levelname).1s %(asctime)s] %(message)s'))
        loc_logzero.logger.addHandler(loc_h)

    loc_logzero.logfile = _logfile
    sys.modules.setdefault('logzero', loc_logzero)

    loc_tqdm = types.ModuleType('tqdm')
    loc_tqdm.tqdm = lambda p_iteravel=None, **p_kwargs: (
        p_iteravel if p_iteravel is not None else [])
    sys.modules.setdefault('tqdm', loc_tqdm)


try:
    from logzero import logger  # type: ignore
except ImportError:
    _sem_logzero()
    from logzero import logger  # type: ignore

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

import connector                                     # noqa: E402
import provisionar                                    # noqa: E402
import regra_projeto                                  # noqa: E402
from contextia_client import (                        # noqa: E402
    L_INTAKE_URL,
    L_URL,
    enviar_dataset,
    gerar_sync,
    normalizar_dataset,
)


AQUI = os.path.dirname(os.path.abspath(__file__))

# Teto de linhas que um ciclo aceita enviar. Acima disso nao e movimento, e
# operacao em massa: recalculo de juros, fechamento de mes, correcao de cadastro
# em lote. Mandar isso pela porta incremental gastaria memoria e tempo, e a
# operacao em massa costuma vir junto de exclusao -- que o incremental nao ve.
# Entao o ouvinte NAO envia: ele entrega o dataset a varredura, que e a
# ferramenta certa para o volume e a unica que remove.
MAXIMO_POR_CICLO = 20000

# Espera entre falhas consecutivas, em segundos. Dobra, e para de dobrar aqui:
# 15 minutos e o suficiente para nao martelar uma origem fora do ar, e pouco
# o bastante para o servico voltar sozinho quando ela voltar.
TETO_ESPERA = 900

# ---------------------------------------------------------------------------
# AS ESTRATEGIAS DE DETECCAO -- a decisao que vem antes de tudo
# ---------------------------------------------------------------------------
#
# "Como eu descubro o que mudou" nao tem UMA resposta: tem cinco, e qual serve
# depende do que o banco do cliente oferece e do que o DBA dele aceita liberar.
# A tabela abaixo e o catalogo, e ela vive no codigo -- e nao so na skill --
# porque e ela que o ouvinte imprime quando recusa uma estrategia que nao
# implementa. Uma recusa que so diz "nao suportado" manda a pessoa procurar; uma
# que diz o que seria preciso resolve a conversa com o DBA na mesma hora.
#
# IMPLEMENTADA HOJE: so a `coluna`.
#
# As outras quatro estao aqui DECLARADAS e NAO implementadas, de proposito. Nao
# vou escrever a abstracao de quatro estrategias imaginadas: interface desenhada
# para requisito que ninguem viu ainda nao serve a nenhum dos reais. Quando
# houver a SEGUNDA de verdade -- e sabendo se e binlog ou Change Tracking --,
# ela mostra onde a costura fica. Ate la, `coluna` roda e as outras recusam
# dizendo o que falta.
ESTRATEGIAS = {
    'coluna': {
        'implementada': True,
        'o_que_e': 'sonda numa coluna de alteracao da propria view (`watch_column`)',
        'exige': 'uma coluna que a origem atualize a cada alteracao. Nada mais:'
                 ' SELECT com a credencial de leitura que ja usamos',
        've_delete': False,
        'nota': 'exclusao so na varredura diaria',
    },
    'binlog': {
        'implementada': False,
        'o_que_e': 'MySQL/MariaDB: ler o binlog pelo protocolo de replicacao',
        'exige': 'log_bin=ON, binlog_format=ROW e GRANT REPLICATION SLAVE.'
                 ' NENHUM DDL e NENHUMA escrita -- e virar uma replica de leitura.'
                 ' Em Python puro (python-mysql-replication), sem client nativo',
        've_delete': True,
        'nota': 'o melhor caso. Cuidado com a expiracao do binlog: fora do ar por'
                ' mais tempo que a retencao, o ouvinte NAO sabe o que perdeu e'
                ' precisa cair na varredura',
    },
    'change_tracking': {
        'implementada': False,
        'o_que_e': 'SQL Server: CHANGETABLE(CHANGES <tabela>, @versao)',
        'exige': 'ALTER DATABASE ... CHANGE_TRACKING = ON uma vez (DDL do DBA, sem'
                 ' trigger) + permissao VIEW CHANGE TRACKING',
        've_delete': True,
        'nota': 'rastreia TABELA, nao view -- e o mapeamento aponta para views.'
                ' Le-se a chave que mudou e relê a view por ela',
    },
    'slot_logico': {
        'implementada': False,
        'o_que_e': 'PostgreSQL: pg_logical_slot_get_changes(), por SQL comum',
        'exige': 'wal_level=logical (EXIGE REINICIO do banco do cliente) +'
                 ' privilegio REPLICATION',
        've_delete': True,
        'nota': 'PERIGO OPERACIONAL: slot nao consumido RETEM WAL. Ouvinte morto e'
                ' sem ninguem olhando enche o disco do cliente. Slot abandonado e'
                ' incidente, nao inconveniente',
    },
    'flashback': {
        'implementada': False,
        'o_que_e': 'Oracle: VERSIONS BETWEEN SCN :a AND :b',
        'exige': 'privilegio FLASHBACK no objeto + undo_retention MAIOR que o'
                 ' intervalo do ciclo',
        've_delete': True,
        'nota': 'nao confunda com ORA_ROWSCN, que custa so um SELECT e parece a'
                ' solucao magica: a granularidade e de BLOCO e a Oracle NAO garante'
                ' o valor para deteccao de mudanca (delayed block cleanout deixa o'
                ' SCN atras do commit). Falso positivo aqui e inofensivo -- o upsert'
                ' deduplica --, mas falso negativo PERDE LINHA em silencio',
    },
    'reconciliacao': {
        'implementada': False,
        'o_que_e': 'diferenca do conjunto de CHAVES: SELECT <pk> FROM <view>',
        'exige': 'nada. Uma coluna, index-only, barato mesmo com 100 mil linhas,'
                 ' em qualquer banco, sem DDL e sem privilegio nenhum',
        've_delete': True,
        'nota': 'a saida universal para quem nao tem coluna nem privilegio: pega'
                ' INCLUSAO e EXCLUSAO ao vivo, e nao pega ALTERACAO. DEPENDE de um'
                ' campo `remove:` no envelope do intake, que HOJE NAO EXISTE --'
                ' hoje ela veria a exclusao e nao teria como contar ao servidor',
    },
}


def exigir_estrategia(p_nome):
    """
    Resolve a estrategia pedida, ou recusa dizendo o que seria preciso.

    A recusa e o produto aqui. Quem pede `binlog` num cliente MySQL nao esta
    errado -- e a melhor escolha para aquele banco --, e a resposta util nao e
    "nao suportado": e a linha que ele leva para o DBA.
    """
    loc = ESTRATEGIAS.get(p_nome)

    if loc is None:
        raise SystemExit(
            'Estrategia %r nao existe. As que existem: %s.'
            % (p_nome, ', '.join(sorted(ESTRATEGIAS)))
        )

    if loc['implementada']:
        return loc

    loc_linhas = [
        'Estrategia %r ainda NAO esta implementada no ouvinte.' % p_nome,
        '',
        '  o que e ..: %s' % loc['o_que_e'],
        '  exigiria .: %s' % loc['exige'],
        '  ve DELETE : %s' % ('sim' if loc['ve_delete'] else 'nao'),
        '  atencao ..: %s' % loc['nota'],
        '',
        'Implementada hoje: `coluna` -- %s.' % ESTRATEGIAS['coluna']['exige'],
        '',
        'Se este cliente nao tem coluna de alteracao, isso NAO e detalhe de',
        'configuracao: e a decisao da Etapa 10, e ela precisa ser tomada com o DBA',
        'antes de prometer dado ao vivo. Ver a tabela de estrategias na skill.',
    ]
    raise SystemExit('\n'.join(loc_linhas))


_parar = False


def _pedir_parada(p_sinal, p_quadro):     # noqa: ARG001
    """
    SIGTERM/SIGINT nao derrubam o ciclo no meio: levantam a bandeira.

    Um `kill` no meio de um envio nao corromperia nada (o intake e idempotente),
    mas deixaria a marca de agua sem avancar e o log sem dizer por que o processo
    sumiu. Terminar o ciclo custa segundos e o desfecho fica explicado.
    """
    global _parar
    _parar = True
    logger.warning('parada pedida; terminando o ciclo atual antes de sair')


# ---------------------------------------------------------------------------
# Estado
# ---------------------------------------------------------------------------

class Estado(object):
    """
    Marcas de agua e escopos ja conhecidos, num JSON gravado de forma atomica.

    Atomica porque o ouvinte pode ser morto a qualquer instante: um arquivo de
    estado truncado faria a proxima subida achar que nunca leu nada, e o
    ouvinte comecaria do `MAX()` da origem -- pulando tudo o que aconteceu
    enquanto ele estava fora.
    """

    def __init__(self, p_path):
        self.path = p_path
        self.dados = {'marca_agua': {}, 'escopos_vistos': [], 'ultima_varredura': None,
                       'versao': 1}
        if os.path.exists(p_path):
            try:
                with io.open(p_path, encoding='utf-8') as f:
                    loc = json.load(f)
                if isinstance(loc, dict):
                    self.dados.update(loc)
            except Exception as e:            # noqa: BLE001
                # Nao apaga e nao continua as cegas: um estado ilegivel e
                # decisao de pessoa. Continuar zerado reprocessaria ou pularia
                # dado sem ninguem saber qual dos dois.
                raise SystemExit(
                    'Estado ilegivel em %s (%s: %s).\n'
                    'Conserte ou mova o arquivo. Comecar do zero faria o ouvinte'
                    ' pular tudo que aconteceu desde a ultima leitura.'
                    % (p_path, type(e).__name__, e)
                )

    def marca(self, p_dataset):
        return (self.dados.get('marca_agua') or {}).get(p_dataset)

    def definir_marca(self, p_dataset, p_valor):
        self.dados.setdefault('marca_agua', {})[p_dataset] = _texto_marca(p_valor)

    def escopos(self):
        return set(self.dados.get('escopos_vistos') or [])

    def registrar_escopos(self, p_escopos):
        loc = self.escopos() | set(p_escopos)
        self.dados['escopos_vistos'] = sorted(loc)

    def varredura_feita_em(self):
        return self.dados.get('ultima_varredura')

    def registrar_varredura(self, p_dia):
        self.dados['ultima_varredura'] = p_dia

    def gravar(self):
        loc_tmp = self.path + '.tmp'
        with io.open(loc_tmp, 'w', encoding='utf-8') as f:
            json.dump(self.dados, f, ensure_ascii=False, indent=2, sort_keys=True)
            f.flush()
            os.fsync(f.fileno())
        os.replace(loc_tmp, self.path)


def _texto_marca(p_valor):
    """A marca vai para o JSON como texto ISO ou numero -- nunca como objeto."""
    if isinstance(p_valor, (datetime.datetime, datetime.date)):
        return p_valor.isoformat()
    if isinstance(p_valor, (int, float)):
        return p_valor
    return None if p_valor is None else str(p_valor)


# ---------------------------------------------------------------------------
# A marca de agua: recuar e avancar
# ---------------------------------------------------------------------------

def recuar(p_marca, p_folga):
    """
    A marca menos a folga -- o valor que vai no `WHERE coluna > ?`.

    A folga cobre a transacao que comecou antes da leitura anterior e comitou
    depois dela: o timestamp gravado e ANTERIOR a marca que ja avancou, e sem
    recuar aquela linha nunca seria lida. E o modo classico de um incremental
    perder exatamente as linhas de maior movimento.
    """
    loc = _como_valor(p_marca)

    if isinstance(loc, datetime.datetime):
        return loc - datetime.timedelta(seconds=p_folga)

    if isinstance(loc, datetime.date):
        # Granularidade de dia: recuar segundos nao muda nada, e `> hoje`
        # perderia toda alteracao de hoje. Recua um dia inteiro -- rele o dia
        # corrente a cada ciclo, o que custa trafego e nao perde linha.
        return loc - datetime.timedelta(days=1)

    # Numero (sequence, ROWVERSION) ou texto opaco: nao ha unidade para a
    # folga. Quem cobre transacao em voo aqui e a varredura diaria.
    return loc


def avancar(p_linhas, p_coluna, p_marca_atual):
    """
    O maior valor da coluna de mudanca entre as linhas lidas.

    Do MAIOR VALOR LIDO, e nao do relogio da maquina: os dois relogios sao
    diferentes, e usar o nosso faria o ouvinte pular a diferenca -- calado.
    Nenhuma linha lida, nenhum avanco: a marca velha continua valendo.
    """
    loc_maior = None
    for linha in p_linhas:
        loc_valor = linha.get(p_coluna)
        if loc_valor is None:
            continue
        if loc_maior is None or _comparavel(loc_valor) > _comparavel(loc_maior):
            loc_maior = loc_valor

    if loc_maior is None:
        return p_marca_atual

    # Nunca RETROCEDE. Uma linha com data futura (digitacao errada, relogio de
    # outro servidor) elevaria a marca; a linha seguinte, normal, nao pode
    # baixa-la de volta -- isso reprocessaria em loop.
    if p_marca_atual is not None:
        try:
            if _comparavel(_como_valor(p_marca_atual)) > _comparavel(loc_maior):
                return p_marca_atual
        except TypeError:
            pass

    return loc_maior


def _como_valor(p_marca):
    """Texto do estado de volta a datetime/date/numero, quando der."""
    if isinstance(p_marca, (datetime.datetime, datetime.date, int, float)):
        return p_marca
    if not isinstance(p_marca, str):
        return p_marca

    for formato in ('%Y-%m-%dT%H:%M:%S.%f', '%Y-%m-%dT%H:%M:%S',
                    '%Y-%m-%d %H:%M:%S.%f', '%Y-%m-%d %H:%M:%S'):
        try:
            return datetime.datetime.strptime(p_marca, formato)
        except ValueError:
            pass
    try:
        return datetime.datetime.strptime(p_marca, '%Y-%m-%d').date()
    except ValueError:
        return p_marca


def _comparavel(p_valor):
    """
    Chave de ordenacao que nao explode ao comparar `date` com `datetime`.

    Python recusa `date > datetime` com TypeError, e uma coluna que devolve os
    dois tipos existe -- driver diferente, view com `CAST`. Um TypeError aqui
    derrubaria o ciclo inteiro por causa de uma linha.
    """
    loc = _como_valor(p_valor)
    if isinstance(loc, datetime.datetime):
        return loc
    if isinstance(loc, datetime.date):
        return datetime.datetime(loc.year, loc.month, loc.day)
    return loc


# ---------------------------------------------------------------------------
# Que datasets dao para acompanhar
# ---------------------------------------------------------------------------

def datasets_ao_vivo(p_mapa):
    """
    Devolve (acompanhaveis, de_fora) -- e a segunda lista importa tanto quanto a
    primeira.

    Dataset sem coluna de mudanca NAO e um erro: tabela de dominio nao muda, e
    view que o ERP nao datou nao tem como ser acompanhada. O que seria erro e
    isso passar calado -- o operador ligaria o ouvinte achando que 27 datasets
    estao ao vivo quando so 12 estao, e descobriria a diferenca quando alguem
    reclamasse de dado velho.

    `watch_column` vem do dataset, ou do `defaults`. `null` no dataset DESLIGA
    o acompanhamento -- e o que se usa quando a coluna existe mas nao presta
    (sempre vazia, ou preenchida por carga em lote em vez de por alteracao).

    A COLUNA FORA DA PROJECAO tem dois desfechos, e a diferenca importa:

        declarada NO DATASET      ABORTA. Alguem escreveu aquele nome ali de
                                  proposito e errou; seguir em frente deixaria
                                  o dataset silenciosamente fora do ao-vivo.
        herdada do `defaults`     sai do ao-vivo, com o motivo no log. Num ERP
                                  real nem toda view tem a mesma coluna de
                                  alteracao, e derrubar o ouvinte inteiro
                                  porque uma das 27 nao a tem obrigaria a
                                  escrever `watch_column: null` uma por uma --
                                  fazendo do erro a primeira experiencia com o
                                  recurso.

    Nos dois casos o dataset continua entrando na varredura diaria: ficar fora
    do ao-vivo nao e ficar fora da integracao.
    """
    loc_padrao = (p_mapa.get('defaults') or {}).get('watch_column')

    loc_vivos = []
    loc_fora = []
    for row in p_mapa['datasets']:
        loc_declarada = 'watch_column' in row
        loc_coluna = row['watch_column'] if loc_declarada else loc_padrao

        if not loc_coluna:
            loc_fora.append((row['dataset'],
                             'watch_column: null' if loc_declarada else 'sem watch_column'))
            continue

        if loc_coluna not in row['columns']:
            if loc_declarada:
                raise SystemExit(
                    'Dataset %r acompanha por %r, mas essa coluna nao esta em'
                    ' `columns`. Sem ela na projecao a marca de agua nunca avanca,'
                    ' e o ouvinte releria o mesmo intervalo indefinidamente.'
                    % (row['dataset'], loc_coluna)
                )
            loc_fora.append((row['dataset'],
                             'nao tem %s (herdada do defaults)' % loc_coluna))
            continue

        loc_copia = dict(row)
        loc_copia['_watch'] = loc_coluna
        loc_vivos.append(loc_copia)

    return loc_vivos, loc_fora


# ---------------------------------------------------------------------------
# Agrupar o que mudou por empresa
# ---------------------------------------------------------------------------

def agrupar_por_escopo(p_linhas, p_mapa, p_row, p_chaves):
    """
    Devolve ({escopo: [registros]}, [escopos desconhecidos]).

    UMA leitura por dataset serve TODAS as empresas, e e por isso que o ciclo
    e barato: a coluna de mudanca nao sabe de empresa, entao filtrar por ela
    empresa a empresa multiplicaria as consultas na producao do cliente por
    onze sem trazer nada de novo. Le-se uma vez e reparte-se aqui.

    Dataset sem recorte (`scope_column: null` -- tabela de dominio) vai para
    TODAS as empresas, que e exatamente o que a carga completa faz: cada projeto
    tem a sua copia do dominio.

    Escopo que aparece no dado e nao esta no arquivo de chaves e uma EMPRESA
    NOVA que ainda nao foi provisionada. Nao se inventa chave para ela: ela e
    devolvida na segunda lista, e o ciclo de projetos -- que sabe criar projeto,
    chave e webhook -- resolve.
    """
    loc_coluna = connector.coluna_de_recorte(p_mapa, p_row)

    if not loc_coluna:
        return {escopo: list(p_linhas) for escopo in p_chaves}, []

    loc_grupos = {}
    loc_desconhecidos = set()
    for linha in p_linhas:
        loc_escopo = linha.get(loc_coluna)
        loc_escopo = '' if loc_escopo is None else str(loc_escopo).strip()

        if loc_escopo not in p_chaves:
            loc_desconhecidos.add(loc_escopo or '(vazio)')
            continue

        loc_grupos.setdefault(loc_escopo, []).append(linha)

    return loc_grupos, sorted(loc_desconhecidos)


# ---------------------------------------------------------------------------
# Ciclo 1 -- dados
# ---------------------------------------------------------------------------

def ciclo_dados(p_ctx):
    """O que mudou desde a ultima leitura, empresa por empresa. Nada e removido."""
    loc_chaves = provisionar.ler_chaves(p_ctx['chaves'])
    if not loc_chaves:
        logger.warning('ciclo de dados: nenhuma chave em %s -- nada a atualizar',
                       p_ctx['chaves'])
        return {'datasets': 0, 'registros': 0, 'erros': 0}

    loc_resumo = {'datasets': 0, 'registros': 0, 'erros': 0, 'grandes': []}
    loc_origem = p_ctx['origem']
    loc_mapa = p_ctx['mapa']

    for row in p_ctx['vivos']:
        if _parar:
            break

        loc_nome = normalizar_dataset(row['dataset'])
        loc_watch = row['_watch']
        loc_marca = p_ctx['estado'].marca(loc_nome)

        # PRIMEIRA VEZ: o ouvinte NAO despeja a tabela inteira. Quem faz carga
        # completa e o conector, e ele ja rodou (ou vai rodar na varredura).
        # Aqui so se fixa o ponto de partida -- de outro modo, ligar o ouvinte
        # reenviaria a base toda pela porta incremental, sem `complete`, e o
        # operador acharia que o servico enlouqueceu.
        if loc_marca is None:
            loc_inicial = loc_origem.maximo(row['source_object'], loc_watch)
            p_ctx['estado'].definir_marca(loc_nome, loc_inicial)
            logger.info('%-28s marca inicial = %s (nada enviado; a carga completa'
                        ' e da varredura)', loc_nome, _texto_marca(loc_inicial))
            continue

        loc_desde = recuar(loc_marca, p_ctx['folga'])

        loc_quantos = loc_origem.contar(row['source_object'],
                                        since=(loc_watch, loc_desde))
        if loc_quantos == 0:
            continue

        if loc_quantos > p_ctx['maximo']:
            # Operacao em massa. A varredura e a ferramenta certa: ela fecha a
            # carga, remove o que sumiu e nao passa nada disso pela memoria de
            # um processo que precisa ficar de pe.
            logger.warning('%-28s %s linhas mudaram (teto %s): entregue a VARREDURA,'
                           ' nao ao incremental', loc_nome, loc_quantos, p_ctx['maximo'])
            loc_resumo['grandes'].append('%s (%s linhas)' % (loc_nome, loc_quantos))
            continue

        loc_linhas = loc_origem.fetch(row['source_object'], row['columns'],
                                      since=(loc_watch, loc_desde))
        loc_grupos, loc_novos = agrupar_por_escopo(loc_linhas, loc_mapa, row, loc_chaves)

        if loc_novos:
            logger.warning('%-28s dado de escopo sem chave: %s -- o ciclo de projetos'
                           ' provisiona', loc_nome, ', '.join(loc_novos[:6]))

        loc_desc, loc_fields = connector.semantica(loc_mapa, row)
        loc_key = row.get('key') or (loc_mapa.get('defaults') or {}).get('key')
        if loc_key not in row['columns']:
            loc_key = None

        if loc_key is None:
            # Sem `key` o servidor deduplica por HASH do conteudo: uma linha
            # ALTERADA nao atualiza -- entra como registro novo, e o antigo fica.
            # Num incremental isso duplica a cada alteracao, e so a varredura
            # limpa. Aviso e nao recusa: dataset sem chave estavel existe, e o
            # incremental ainda vale para as INCLUSOES.
            logger.warning('%-28s sem `key`: alteracao vira registro NOVO (dedup por'
                           ' hash). Ate a varredura, o antigo fica.', loc_nome)

        loc_ok = True
        loc_enviados = 0
        for loc_escopo, loc_registros in sorted(loc_grupos.items()):
            if p_ctx['seco']:
                logger.info('%-28s [SECO] escopo=%-8s %s registro(s)',
                            loc_nome, loc_escopo, len(loc_registros))
                loc_enviados += len(loc_registros)
                continue

            loc_res = enviar_dataset(
                loc_nome, loc_registros,
                p_key=loc_key,
                p_description=loc_desc,
                p_fields=loc_fields,
                p_group=row.get('group'),
                p_chunk=row.get('chunk') or (loc_mapa.get('defaults') or {}).get('chunk') or 500,
                p_api_key=loc_chaves[loc_escopo],
                p_logger=logger,
                # SEM CICLO DE CARGA, e isto e o coracao da seguranca deste
                # arquivo. `sync=None` faz o servidor tratar tudo como upsert e
                # NAO remover nada. Um incremental que mandasse `complete`
                # apagaria todo registro que nao mudou nos ultimos 60 segundos
                # -- ou seja, o dataset inteiro.
                p_sync=None,
            )

            if loc_res['lotes_erro']:
                loc_ok = False
                loc_resumo['erros'] += loc_res['lotes_erro']

            loc_enviados += loc_res['registros']
            logger.info('%-28s escopo=%-8s enviados=%-5s novos=%-4s atualizados=%-4s'
                        ' iguais=%-4s', loc_nome, loc_escopo, loc_res['registros'],
                        loc_res['inserted'], loc_res['updated'], loc_res['duplicates'])

        # A MARCA SO AVANCA SE TUDO DEU CERTO. Mesma regra do `complete`: falha
        # de envio mantem a marca velha e o proximo ciclo rele. Avancar sobre um
        # envio que falhou perderia aquelas linhas para sempre -- e o incremental
        # nao tem varredura propria para consertar antes da noite.
        if loc_ok and not p_ctx['seco']:
            # A marca avanca sobre as linhas LIDAS, nao sobre as ENVIADAS -- e a
            # diferenca sao as linhas de empresa ainda nao provisionada. Elas
            # ficam para tras de proposito: se a marca esperasse por elas, um
            # dataset cujo unico movimento fosse de uma empresa desconhecida
            # releria o mesmo intervalo para sempre.
            #
            # Nada se perde porque o ciclo de projetos faz CARGA COMPLETA da
            # empresa nova assim que a provisiona -- e essa carga traz o
            # historico inteiro dela, nao so o que mudou. As duas coisas
            # dependem uma da outra: mexer numa, reveja a outra.
            p_ctx['estado'].definir_marca(
                loc_nome, avancar(loc_linhas, loc_watch, loc_marca))
        elif not loc_ok:
            logger.warning('%-28s envio com erro: marca de agua NAO avancou'
                           ' (o proximo ciclo rele)', loc_nome)

        loc_resumo['datasets'] += 1
        loc_resumo['registros'] += loc_enviados

    if not p_ctx['seco']:
        p_ctx['estado'].gravar()

    return loc_resumo


# ---------------------------------------------------------------------------
# Ciclo 2 -- projetos (empresa nova)
# ---------------------------------------------------------------------------

def ciclo_projetos(p_ctx):
    """
    Empresa nova no ERP ganha projeto, chave, grupo, webhook e carga completa.

    NAO reimplementa o provisionamento: chama o `provisionar.py`, que e onde a
    chave e gravada antes de qualquer coisa e o token e entregue assinado. Um
    segundo caminho para criar chave seria um segundo caminho para perder uma.
    """
    loc_regra = p_ctx['regra']
    loc_conhecidos = set(provisionar.ler_chaves(p_ctx['chaves']))
    loc_excluidos = set((p_ctx['mapa'].get('source') or {}).get('excluded_scopes') or {})

    loc_linhas = p_ctx['origem'].fetch(loc_regra['source_object'], loc_regra['columns'])

    loc_novos = []
    loc_recusados = []
    for linha in loc_linhas:
        loc_escopo = regra_projeto.escopo_de(linha, loc_regra)

        if loc_escopo in loc_conhecidos or loc_escopo in loc_excluidos:
            continue

        loc_ok, loc_motivo = regra_projeto.aceita(linha, loc_regra)
        if not loc_ok:
            loc_recusados.append('%s (%s)' % (loc_escopo or '?', loc_motivo))
            continue

        loc_novos.append((loc_escopo, regra_projeto.nome_de(linha, loc_regra)))

    if loc_recusados:
        logger.info('ciclo de projetos: %s escopo(s) recusados pela regra: %s',
                    len(loc_recusados), '; '.join(loc_recusados[:5]))

    if not loc_novos:
        return {'criados': 0, 'carregados': 0, 'erros': 0}

    logger.warning('*** %s EMPRESA(S) NOVA(S): %s', len(loc_novos),
                   ', '.join('%s=%s' % (e, n[:24]) for e, n in loc_novos))

    loc_resumo = {'criados': 0, 'carregados': 0, 'erros': 0}

    for loc_escopo, loc_nome in loc_novos:
        if _parar:
            break

        # Uma chamada por empresa, e nao uma para todas: com `group_column` os
        # grupos saem do cadastro e diferem de empresa para empresa.
        loc_argv = [p_ctx['python'], 'provisionar.py',
                    '--mapping=' + p_ctx['mapping'], '--source=' + p_ctx['source'],
                    '--arquivo=' + p_ctx['chaves'],
                    '--somente-escopo=' + loc_escopo]
        for grupo in p_ctx['grupos']:
            loc_argv.append('--grupo=' + grupo)
        if p_ctx['seco']:
            loc_argv.append('--dry-run')

        if _rodar(loc_argv, 'provisionar %s' % loc_escopo, p_ctx) != 0:
            loc_resumo['erros'] += 1
            continue
        loc_resumo['criados'] += 1

        if p_ctx['seco']:
            continue

        # CARGA COMPLETA para a empresa nova, agora. Ela nao tem nada no
        # Contextia, e o ciclo incremental so traz o que MUDAR daqui para
        # frente -- deixa-la para a varredura significaria um assistente vazio
        # ate a madrugada, no dia em que o cliente foi cadastrado.
        loc_chaves = provisionar.ler_chaves(p_ctx['chaves'])
        if loc_escopo not in loc_chaves:
            logger.error('escopo %s provisionado e sem chave no arquivo -- carga'
                         ' inicial nao feita', loc_escopo)
            loc_resumo['erros'] += 1
            continue

        loc_env = dict(os.environ, CONTEXTIA_API_KEY=loc_chaves[loc_escopo])
        loc_argv = [p_ctx['python'], 'connector.py',
                    '--mapping=' + p_ctx['mapping'], '--source=' + p_ctx['source'],
                    '--scope=' + loc_escopo]
        if _rodar(loc_argv, 'carga inicial %s' % loc_escopo, p_ctx, loc_env) != 0:
            loc_resumo['erros'] += 1
            continue
        loc_resumo['carregados'] += 1

    p_ctx['estado'].registrar_escopos(e for e, _ in loc_novos)
    p_ctx['estado'].gravar()

    return loc_resumo


# ---------------------------------------------------------------------------
# Ciclo 3 -- varredura (a carga completa, e a unica que remove)
# ---------------------------------------------------------------------------

def ciclo_varredura(p_ctx):
    """
    A recarga completa, identica a que o cron fazia -- as duas fases, na ordem.

    Fase 1 provisiona (rede de seguranca: se o ciclo de projetos vinha
    falhando, e aqui que aparece); fase 2 carrega tudo com `sync`+`complete`,
    e e o unico momento em que registro apagado na origem sai do Contextia.
    """
    logger.warning('===== varredura: carga completa de todas as empresas =====')

    loc_erros = 0

    loc_argv = [p_ctx['python'], 'provisionar.py',
                '--mapping=' + p_ctx['mapping'], '--source=' + p_ctx['source'],
                '--arquivo=' + p_ctx['chaves']]
    for grupo in p_ctx['grupos']:
        loc_argv.append('--grupo=' + grupo)
    if p_ctx['seco']:
        loc_argv.append('--dry-run')
    if _rodar(loc_argv, 'varredura/provisionar', p_ctx) != 0:
        loc_erros += 1

    loc_env = dict(os.environ, CONTEXTIA_MAPPING=p_ctx['mapping'],
                   CONTEXTIA_PYTHON=p_ctx['python'], CONTEXTIA_COLETA=AQUI)
    loc_argv = ['bash', os.path.join(AQUI, 'carregar-todas.sh'), p_ctx['chaves']]
    if p_ctx['seco']:
        loc_argv.append('--dry-run')
    if _rodar(loc_argv, 'varredura/carregar-todas', p_ctx, loc_env) != 0:
        loc_erros += 1

    # A varredura mexeu na origem por fora deste processo e pode ter demorado
    # horas. O retrato da nossa sessao esta velho.
    p_ctx['origem'].renovar()

    if not p_ctx['seco']:
        p_ctx['estado'].registrar_varredura(datetime.date.today().isoformat())
        p_ctx['estado'].gravar()

    return {'erros': loc_erros}


def hora_de_varrer(p_ctx, p_agora=None):
    """
    Chegou a hora da varredura de hoje?

    Duas decisoes aqui, e as duas foram aprendidas errando:

    **Compara com o DIA ja varrido, nao com um intervalo.** "As 03:30" tem de
    acontecer uma vez por dia mesmo que o ouvinte tenha sido reiniciado tres
    vezes, e nao acontecer duas vezes porque ele subiu 03:29.

    **E so dentro de uma JANELA depois da hora.** Sem ela, um ouvinte que sobe
    as 14h com o estado zerado (implantacao nova, disco trocado, arquivo movido)
    ve "ainda nao varri hoje" e dispara na hora -- carga completa de todas as
    empresas em plena tarde, que e exatamente o que a janela combinada na Etapa
    0 existe para evitar. Com a janela de 3h, ele nao varre hoje, varre amanha
    as 03:30, e quem quiser a carga inicial agora pede com `--varrer-agora`.

    A janela tambem serve de recuperacao: se o ouvinte estava fora as 03:30 e
    subiu as 05:00, ele ainda varre. Perder a varredura por ter ficado fora dez
    minutos seria pior que varrer com atraso.
    """
    if not p_ctx['varredura']:
        return False

    loc_agora = p_agora or datetime.datetime.now()

    if p_ctx['estado'].varredura_feita_em() == loc_agora.date().isoformat():
        return False

    loc_h, loc_m = p_ctx['varredura']
    loc_marcada = loc_agora.replace(hour=loc_h, minute=loc_m, second=0, microsecond=0)
    loc_janela = datetime.timedelta(minutes=p_ctx.get('janela', 180))

    return loc_marcada <= loc_agora <= loc_marcada + loc_janela


# ---------------------------------------------------------------------------
# Rodar as ferramentas de sempre
# ---------------------------------------------------------------------------

def _rodar(p_argv, p_rotulo, p_ctx, p_env=None):
    """
    Executa uma das ferramentas do SDK e devolve o codigo de saida.

    Subprocesso, e nao `import` + chamada de funcao, por um motivo que vale
    dizer: `provisionar.py` e `connector.py` chamam `SystemExit` em erro de
    configuracao. Importados dentro de um processo que precisa FICAR DE PE, um
    `SystemExit` derrubaria o ouvinte junto. Isolado, ele volta como codigo de
    saida e o ouvinte segue.

    A saida vai INTEIRA para o log: sao os numeros da carga (`novos`,
    `removidos`, `campos descritos`) que alguem vai querer conferir depois, e
    guardar so o codigo de saida jogaria tudo isso fora.
    """
    logger.info('$ %s', ' '.join(p_argv))
    try:
        loc = subprocess.run(
            p_argv, cwd=AQUI, env=p_env or os.environ,
            stdout=subprocess.PIPE, stderr=subprocess.STDOUT,
            timeout=p_ctx['timeout'],
        )
    except subprocess.TimeoutExpired:
        logger.error('%s: estourou o tempo de %ss e foi interrompido',
                     p_rotulo, p_ctx['timeout'])
        return 124
    except Exception as e:                    # noqa: BLE001
        logger.error('%s: nao foi possivel executar (%s: %s)', p_rotulo,
                     type(e).__name__, e)
        return 125

    for linha in (loc.stdout or b'').decode('utf-8', 'replace').splitlines():
        logger.info('  | %s', linha)

    if loc.returncode != 0:
        logger.error('%s: saiu com %s', p_rotulo, loc.returncode)
    return loc.returncode


# ---------------------------------------------------------------------------
# Saude
# ---------------------------------------------------------------------------

def escrever_status(p_ctx, p_linhas):
    """
    O arquivo que um alerta le -- e a razao de ele ser um ARQUIVO.

    Alerta que precisa varrer log para saber se o servico esta vivo nao e
    alerta. Aqui basta olhar a idade do arquivo: se ele nao foi tocado no ultimo
    ciclo, o ouvinte parou.

        find <estado>/ouvinte-status.txt -mmin +10   # nada = ele parou
    """
    loc_path = os.path.join(p_ctx['estado_dir'], 'ouvinte-status.txt')
    loc_tmp = loc_path + '.tmp'
    try:
        with io.open(loc_tmp, 'w', encoding='utf-8') as f:
            f.write('ouvinte ...........: %s\n' % ('SECO (nao envia)' if p_ctx['seco'] else 'ativo'))
            f.write('ultimo ciclo ......: %s\n' % datetime.datetime.now().strftime('%F %T %Z').strip())
            f.write('cliente ...........: %s\n' % (p_ctx['mapa'].get('client') or '-'))
            f.write('destino ...........: %s\n' % (L_URL or L_INTAKE_URL))
            f.write('datasets ao vivo ..: %s de %s\n' % (
                len(p_ctx['vivos']), len(p_ctx['mapa']['datasets'])))
            f.write('falhas seguidas ...: %s\n' % p_ctx['falhas'])
            f.write('ultima varredura ..: %s\n' % (p_ctx['estado'].varredura_feita_em() or 'nenhuma'))
            for linha in p_linhas:
                f.write('%s\n' % linha)
            f.flush()
            os.fsync(f.fileno())
        os.replace(loc_tmp, loc_path)
    except Exception as e:                    # noqa: BLE001
        # Nao derruba o ouvinte por causa do arquivo de status: o servico dele e
        # levar dado, e disco cheio nao e motivo para parar de levar.
        logger.warning('nao foi possivel gravar o status (%s: %s)', type(e).__name__, e)


def travar(p_estado_dir):
    """
    Um ouvinte por implantacao.

    Dois ouvintes leriam o MESMO estado e se sobrescreveriam: cada um avancaria
    a marca de agua sem ver o que o outro leu, e o intervalo lido por um seria
    considerado lido pelo outro. Perda de dado, sem erro nenhum.
    """
    loc_path = os.path.join(p_estado_dir, 'ouvinte.lock')
    try:
        import fcntl
    except ImportError:
        logger.warning('sem fcntl nesta plataforma: NAO ha trava. Garanta por fora'
                       ' que so um ouvinte roda por implantacao.')
        return None

    loc_f = io.open(loc_path, 'w')
    try:
        fcntl.flock(loc_f.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB)
    except (IOError, OSError):
        raise SystemExit(
            'Ja existe um ouvinte rodando nesta implantacao (%s).\n'
            'Dois ouvintes disputariam a mesma marca de agua, e o intervalo lido'
            ' por um seria dado como lido pelo outro.' % loc_path
        )
    loc_f.write('%s\n' % os.getpid())
    loc_f.flush()
    return loc_f


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

def _hora(p_texto):
    if not p_texto or p_texto.lower() in ('nao', 'no', 'off', 'nenhuma'):
        return None
    try:
        loc_h, loc_m = p_texto.split(':')
        loc_h, loc_m = int(loc_h), int(loc_m)
        if not (0 <= loc_h <= 23 and 0 <= loc_m <= 59):
            raise ValueError
        return loc_h, loc_m
    except Exception:
        raise SystemExit('--varredura espera HH:MM (ou "nao"), recebi %r' % p_texto)


def main(p_argv=None):
    # O autoteste nao exige mapeamento, chaves nem estado -- ele nao toca em
    # nenhum dos tres. Atendido ANTES do argparse porque `required=True` o
    # obrigaria a inventar valores para argumentos que ele nao usa, e a primeira
    # coisa que alguem faz com um autoteste e roda-lo sem ler nada.
    if '--autoteste' in (p_argv if p_argv is not None else sys.argv[1:]):
        return _autoteste()

    p = argparse.ArgumentParser(
        description='Ouvinte: leva ao Contextia o que mudou na origem, em segundos.')
    p.add_argument('--mapping', required=True)
    p.add_argument('--source', required=True, help='modulo em sources/ (sem .py)')
    p.add_argument('--chaves', required=True,
                   help='o arquivo de chaves do provisionador (uma por empresa)')
    p.add_argument('--estado', required=True,
                   help='diretorio das marcas de agua, do status e da trava')
    p.add_argument('--intervalo-dados', type=int, default=60, metavar='S',
                   help='segundos entre leituras do que mudou (padrao 60)')
    p.add_argument('--intervalo-projetos', type=int, default=300, metavar='S',
                   help='segundos entre buscas por empresa nova (padrao 300)')
    p.add_argument('--varredura', default='03:30', metavar='HH:MM',
                   help='hora da carga completa diaria, a UNICA que remove'
                        ' (padrao 03:30; "nao" desliga)')
    p.add_argument('--janela-varredura', type=int, default=180, metavar='MIN',
                   help='por quantos minutos depois da hora a varredura ainda'
                        ' pode disparar (padrao 180). Fora da janela ela espera o'
                        ' dia seguinte, para nao cair no horario comercial')
    p.add_argument('--estrategia', default=None, metavar='NOME',
                   help='como detectar o que mudou: %s. Padrao: o'
                        ' `listener.strategy` do mapeamento, ou `coluna`'
                        % ', '.join(sorted(ESTRATEGIAS)))
    p.add_argument('--folga', type=int, default=180, metavar='S',
                   help='quanto a marca de agua recua a cada leitura, para nao'
                        ' perder transacao que comitou atrasada (padrao 180)')
    p.add_argument('--maximo-por-ciclo', type=int, default=MAXIMO_POR_CICLO, metavar='N',
                   help='acima disso o dataset e entregue a varredura (padrao %d)'
                        % MAXIMO_POR_CICLO)
    p.add_argument('--grupo', action='append', default=[], metavar='SLUG',
                   help='grupo de empresas para o projeto de empresa nova;'
                        ' repita para mais de um')
    p.add_argument('--timeout-ferramenta', type=int, default=7200, metavar='S',
                   help='tempo maximo de uma carga completa (padrao 7200)')
    p.add_argument('--log', default=None, help='arquivo de log (alem da saida padrao)')
    p.add_argument('--varrer-agora', action='store_true',
                   help='roda a varredura completa nesta subida, sem esperar a hora')
    p.add_argument('--once', action='store_true',
                   help='um ciclo de cada e sai. E o modo de QA e de teste do cron')
    p.add_argument('--dry-run', action='store_true',
                   help='le tudo e nao envia, nao cria e nao grava marca de agua')
    p.add_argument('--autoteste', action='store_true',
                   help='prova a logica pura sem banco e sem rede')
    args = p.parse_args(p_argv)

    if args.autoteste:
        return _autoteste()

    if args.log:
        import logzero                        # type: ignore
        logzero.logfile(args.log, maxBytes=20 * 1024 * 1024, backupCount=10)

    import signal
    for loc_sinal in (signal.SIGTERM, signal.SIGINT):
        signal.signal(loc_sinal, _pedir_parada)

    # A ESTRATEGIA e a PRIMEIRA coisa resolvida, antes do destino, do arquivo de
    # chaves e de qualquer conexao. Dois motivos:
    #
    #   * conceitualmente ela vem primeiro -- e a decisao da Etapa 10, tomada com
    #     o DBA do cliente, e as outras sao consequencia dela;
    #   * praticamente, `--estrategia=binlog` e a pergunta "o que isso exigiria?".
    #     Quem quer a resposta nao deveria ter de configurar URL e provisionar
    #     chaves primeiro para consegui-la.
    #
    # A linha de comando vence o mapeamento: quem digitou esta olhando a tela
    # agora, e o mapeamento foi escrito semanas atras.
    loc_mapa = connector.carregar_mapeamento(args.mapping)
    loc_estrategia = (args.estrategia
                      or (loc_mapa.get('listener') or {}).get('strategy')
                      or 'coluna')
    exigir_estrategia(loc_estrategia)

    if not L_URL and not L_INTAKE_URL:
        raise SystemExit(
            'Falta CONTEXTIA_URL (a instalacao DESTE cliente).\n'
            'Nao ha padrao embutido: um valor default mandaria o dado para a'
            ' instalacao de outro cliente, sem erro nenhum.'
        )
    if not os.path.exists(args.chaves):
        raise SystemExit(
            'Arquivo de chaves nao existe: %s\n'
            'O ouvinte precisa de uma chave por empresa. Rode o provisionar.py'
            ' antes -- ou aponte --chaves para o arquivo que ele gerou.' % args.chaves
        )

    os.makedirs(args.estado, exist_ok=True)
    loc_trava = travar(args.estado)

    loc_vivos, loc_fora = datasets_ao_vivo(loc_mapa)
    loc_regra = regra_projeto.resolver(loc_mapa)
    loc_origem = importlib.import_module(args.source)

    for loc_nome in ('fetch', 'contar', 'maximo', 'renovar', 'fechar'):
        if not hasattr(loc_origem, loc_nome):
            raise SystemExit(
                'sources/%s.py nao expoe `%s`. O ouvinte precisa de sonda'
                ' (`contar`), leitura incremental (`fetch(since=)`), ponto de'
                ' partida (`maximo`) e renovacao de retrato (`renovar`).'
                ' Use sources/generic.py.' % (args.source, loc_nome))

    loc_ctx = {
        'mapa': loc_mapa, 'vivos': loc_vivos, 'regra': loc_regra, 'origem': loc_origem,
        'estado': Estado(os.path.join(args.estado, 'ouvinte-estado.json')),
        'estado_dir': args.estado,
        'mapping': args.mapping, 'source': args.source, 'chaves': args.chaves,
        'folga': args.folga, 'maximo': args.maximo_por_ciclo,
        'grupos': args.grupo, 'seco': args.dry_run,
        'varredura': _hora(args.varredura),
        'janela': args.janela_varredura,
        'timeout': args.timeout_ferramenta,
        'python': sys.executable or 'python3',
        'falhas': 0,
    }

    logger.info('===== ouvinte =====')
    logger.info('cliente ..........: %s', loc_mapa.get('client') or '(sem nome)')
    logger.info('destino ..........: %s', L_URL or L_INTAKE_URL)
    logger.info('chaves ...........: %s (%s empresa(s))',
                args.chaves, len(provisionar.ler_chaves(args.chaves)))
    logger.info('estrategia .......: %s -- %s',
                loc_estrategia, ESTRATEGIAS[loc_estrategia]['o_que_e'])
    if not ESTRATEGIAS[loc_estrategia]['ve_delete']:
        logger.warning('esta estrategia NAO ve exclusao: linha apagada na origem so'
                       ' sai do Contextia na varredura. O cliente precisa saber disso'
                       ' antes de ligar.')
    logger.info('dados ............: a cada %ss, folga de %ss, teto de %s linhas',
                args.intervalo_dados, args.folga, args.maximo_por_ciclo)
    logger.info('projetos .........: a cada %ss', args.intervalo_projetos)
    logger.info('varredura ........: %s',
                '%02d:%02d (+%smin de janela; a UNICA que remove)'
                % (loc_ctx['varredura'][0], loc_ctx['varredura'][1], args.janela_varredura)
                if loc_ctx['varredura'] else 'DESLIGADA -- exclusao na origem NUNCA'
                                             ' sai do Contextia')
    if loc_ctx['varredura'] and not loc_ctx['estado'].varredura_feita_em():
        logger.warning('nenhuma varredura registrada ainda. O ouvinte NAO faz carga'
                       ' completa fora da janela -- ela cairia no horario comercial.'
                       ' Para a carga inicial, suba com --varrer-agora.')
    if args.dry_run:
        logger.warning('MODO SECO: nada e enviado, nada e criado, a marca de agua'
                       ' nao avanca')

    logger.info('datasets ao vivo .: %s de %s', len(loc_vivos), len(loc_mapa['datasets']))
    for loc_nome, loc_motivo in loc_fora:
        logger.info('   fora do ao-vivo: %-28s %s (so muda na varredura)',
                    loc_nome, loc_motivo)
    logger.info('regra de projeto:\n%s', regra_projeto.descrever(loc_regra))

    if not loc_vivos:
        logger.warning('NENHUM dataset tem coluna de mudanca: o ouvinte so fara'
                       ' varredura, o que o cron ja fazia. Declare `watch_column`'
                       ' no mapeamento para o ao-vivo valer a pena.')

    loc_proximo_dados = 0.0
    loc_proximo_projetos = 0.0
    loc_saida = 0

    try:
        while True:
            loc_agora = time.monotonic()
            loc_linhas_status = []

            try:
                # Retrato novo a cada rodada. Sem isto o Oracle e o MySQL
                # devolveriam para sempre o banco do instante em que o ouvinte
                # subiu -- ver `renovar()` em sources/generic.py.
                loc_origem.renovar()

                if loc_agora >= loc_proximo_projetos or args.once:
                    loc_res = ciclo_projetos(loc_ctx)
                    loc_proximo_projetos = time.monotonic() + args.intervalo_projetos
                    loc_linhas_status.append(
                        'projetos ..........: criados=%s carregados=%s erros=%s'
                        % (loc_res['criados'], loc_res['carregados'], loc_res['erros']))

                if loc_agora >= loc_proximo_dados or args.once:
                    loc_res = ciclo_dados(loc_ctx)
                    loc_proximo_dados = time.monotonic() + args.intervalo_dados
                    loc_linhas_status.append(
                        'dados .............: datasets=%s registros=%s erros=%s'
                        % (loc_res['datasets'], loc_res['registros'], loc_res['erros']))
                    if loc_res.get('grandes'):
                        loc_linhas_status.append(
                            'entregues a varredura: %s' % ', '.join(loc_res['grandes']))

                if hora_de_varrer(loc_ctx) or args.varrer_agora:
                    loc_res = ciclo_varredura(loc_ctx)
                    loc_linhas_status.append('varredura .........: erros=%s'
                                             % loc_res['erros'])
                    # Uma vez, nesta subida. Sem isto o laco varreria a cada
                    # ciclo -- carga completa de todas as empresas a cada 60s.
                    args.varrer_agora = False

                loc_ctx['falhas'] = 0

            except Exception as e:            # noqa: BLE001
                # O ouvinte NAO morre por uma falha de ciclo. Origem fora do ar,
                # tunel caido, Contextia em manutencao: tudo isso passa, e um
                # processo que morre no primeiro soluco vira um servico que
                # ninguem confia -- e que so volta quando alguem percebe.
                loc_ctx['falhas'] += 1
                logger.exception('ciclo falhou (%sa vez seguida): %s: %s',
                                 loc_ctx['falhas'], type(e).__name__, e)
                loc_linhas_status.append('ULTIMA FALHA ......: %s: %s'
                                         % (type(e).__name__, str(e)[:200]))
                # Conexao possivelmente morta: derruba para o proximo ciclo
                # reconectar em vez de insistir num handle quebrado.
                try:
                    loc_origem.fechar()
                except Exception:             # noqa: BLE001
                    pass
                loc_saida = 1

            escrever_status(loc_ctx, loc_linhas_status)

            if args.once or _parar:
                break

            loc_espera = args.intervalo_dados
            if loc_ctx['falhas']:
                loc_espera = min(args.intervalo_dados * (2 ** loc_ctx['falhas']),
                                 TETO_ESPERA)
                logger.warning('esperando %ss antes de tentar de novo', loc_espera)

            # Dorme em fatias para que SIGTERM nao espere o intervalo inteiro.
            loc_fim = time.monotonic() + loc_espera
            while time.monotonic() < loc_fim and not _parar:
                time.sleep(min(1.0, loc_fim - time.monotonic()))

            if _parar:
                break
    finally:
        try:
            loc_origem.fechar()
        except Exception:                     # noqa: BLE001
            pass
        if loc_trava is not None:
            loc_trava.close()
        logger.info('ouvinte encerrado')

    return loc_saida if args.once else 0


# ---------------------------------------------------------------------------
# Autoteste -- sem banco, sem rede, sem mapeamento de cliente
# ---------------------------------------------------------------------------

def _autoteste():
    """
    `python3 ouvinte.py --autoteste` -- prova a logica que nao da para ver a olho.

    Sao justamente as partes onde um erro perde dado em silencio: a marca de
    agua, o repartir por empresa e a montagem do WHERE. Roda em qualquer maquina,
    antes de tocar na origem do cliente.
    """
    loc_falhas = []

    def igual(p_rotulo, p_obtido, p_esperado):
        if p_obtido != p_esperado:
            loc_falhas.append('%s: obtido %r, esperado %r' % (p_rotulo, p_obtido, p_esperado))
            print('  FALHOU  %s' % p_rotulo)
        else:
            print('  ok      %s' % p_rotulo)

    print('marca de agua -- recuar')
    loc_dt = datetime.datetime(2026, 9, 2, 10, 0, 0)
    igual('datetime recua os segundos da folga',
          recuar(loc_dt, 180), datetime.datetime(2026, 9, 2, 9, 57, 0))
    igual('texto ISO do estado volta a datetime e recua',
          recuar('2026-09-02T10:00:00', 60), datetime.datetime(2026, 9, 2, 9, 59, 0))
    igual('coluna de DIA recua um dia inteiro',
          recuar(datetime.date(2026, 9, 2), 180), datetime.date(2026, 9, 1))
    igual('numero (sequence) nao recua', recuar(918273, 180), 918273)

    print('\nmarca de agua -- avancar')
    loc_linhas = [
        {'ID': 1, 'ALT': datetime.datetime(2026, 9, 2, 10, 0, 0)},
        {'ID': 2, 'ALT': datetime.datetime(2026, 9, 2, 10, 5, 0)},
        {'ID': 3, 'ALT': None},
    ]
    igual('avanca para o MAIOR valor lido',
          avancar(loc_linhas, 'ALT', loc_dt), datetime.datetime(2026, 9, 2, 10, 5, 0))
    igual('nenhuma linha: a marca velha continua',
          avancar([], 'ALT', loc_dt), loc_dt)
    igual('so linha sem valor: a marca velha continua',
          avancar([{'ID': 9, 'ALT': None}], 'ALT', loc_dt), loc_dt)
    igual('nao retrocede quando o lido e menor',
          avancar([{'ALT': datetime.datetime(2026, 9, 1, 8, 0, 0)}], 'ALT', loc_dt), loc_dt)
    igual('date e datetime na mesma coluna nao explodem',
          avancar([{'ALT': datetime.date(2026, 9, 3)},
                   {'ALT': datetime.datetime(2026, 9, 2, 10, 5)}], 'ALT', loc_dt),
          datetime.date(2026, 9, 3))

    print('\nrepartir por empresa')
    loc_mapa = {'defaults': {'scope_column': 'EMPRESA_ID', 'key': 'ID'}, 'datasets': []}
    loc_row = {'dataset': 'contas', 'source_object': 'V', 'columns': ['ID', 'EMPRESA_ID']}
    loc_chaves = {'297': 'ctx_a', '301': 'ctx_b'}
    loc_grupos, loc_novos = agrupar_por_escopo(
        [{'ID': 1, 'EMPRESA_ID': 297}, {'ID': 2, 'EMPRESA_ID': '301'},
         {'ID': 3, 'EMPRESA_ID': 999}, {'ID': 4, 'EMPRESA_ID': None}],
        loc_mapa, loc_row, loc_chaves)
    igual('numero 297 casa com a chave "297"', sorted(loc_grupos), ['297', '301'])
    igual('escopo sem chave nao inventa destino', loc_novos, ['(vazio)', '999'])

    loc_dominio = {'dataset': 'status', 'source_object': 'V', 'scope_column': None,
                   'columns': ['ID']}
    loc_grupos, _ = agrupar_por_escopo([{'ID': 1}], loc_mapa, loc_dominio, loc_chaves)
    igual('tabela de dominio vai para TODAS as empresas',
          sorted(loc_grupos), ['297', '301'])
    igual('e com o mesmo registro em cada uma',
          [len(v) for _, v in sorted(loc_grupos.items())], [1, 1])

    print('\nque datasets ficam ao vivo')
    loc_mapa2 = {'defaults': {'watch_column': 'ALT'}, 'datasets': [
        {'dataset': 'a', 'source_object': 'V', 'columns': ['ID', 'ALT']},
        {'dataset': 'b', 'source_object': 'V', 'columns': ['ID'], 'watch_column': None},
        {'dataset': 'c', 'source_object': 'V', 'columns': ['ID', 'M'], 'watch_column': 'M'},
    ]}
    loc_mapa2['datasets'].append(
        {'dataset': 'd', 'source_object': 'V', 'columns': ['ID', 'OUTRA']})
    loc_vivos, loc_fora = datasets_ao_vivo(loc_mapa2)
    igual('herda do defaults e aceita o especifico',
          [d['dataset'] for d in loc_vivos], ['a', 'c'])
    igual('null no dataset desliga; view sem a coluna herdada sai sem derrubar',
          [n for n, _ in loc_fora], ['b', 'd'])
    igual('e o motivo diz que a coluna foi herdada',
          'herdada' in dict(loc_fora)['d'], True)

    try:
        datasets_ao_vivo({'defaults': {}, 'datasets': [
            {'dataset': 'x', 'source_object': 'V', 'columns': ['ID'], 'watch_column': 'ALT'}]})
        loc_falhas.append('watch_column DECLARADA fora de columns: NAO abortou')
        print('  FALHOU  watch_column DECLARADA fora de columns ABORTA')
    except SystemExit:
        print('  ok      watch_column DECLARADA fora de columns ABORTA')

    print('\nWHERE montado pelo adaptador')
    import generic                            # noqa: E402
    loc_original = generic.conectar
    try:
        generic.conectar = lambda: None
        for loc_driver, loc_esperados in (
            ('mysql', [
                ('so escopo (identico ao de antes do ouvinte)',
                 (None, ('EMPRESA_ID', 297)), 'SELECT 1 WHERE EMPRESA_ID = %s', (297,)),
                ('so mudanca', (('ALT', 'x'), None), 'SELECT 1 WHERE ALT > %s', ('x',)),
                ('os dois, escopo primeiro', (('ALT', 'x'), ('E', 1)),
                 'SELECT 1 WHERE E = %s AND ALT > %s', (1, 'x')),
            ]),
            ('oracle', [
                ('so escopo, parametro nomeado',
                 (None, ('EMPRESA_ID', 297)), 'SELECT 1 WHERE EMPRESA_ID = :escopo',
                 {'escopo': 297}),
                ('os dois, nomeados', (('ALT', 'x'), ('E', 1)),
                 'SELECT 1 WHERE E = :escopo AND ALT > :desde',
                 {'escopo': 1, 'desde': 'x'}),
            ]),
        ):
            generic._driver = loc_driver
            for loc_rotulo, (loc_since, loc_scope), loc_sql, loc_par in loc_esperados:
                loc_obtido = generic._onde('SELECT 1', loc_scope, loc_since)
                igual('%s: %s' % (loc_driver, loc_rotulo), loc_obtido, (loc_sql, loc_par))

        igual('sem filtro nenhum: SQL intacto',
              generic._onde('SELECT 1', None, None), ('SELECT 1', ()))
    finally:
        generic.conectar = loc_original
        generic._driver = None

    print('\nhora da varredura')
    class _E(object):
        def __init__(self, p_dia):
            self.dia = p_dia
        def varredura_feita_em(self):
            return self.dia
    igual('antes da hora, nao varre',
          hora_de_varrer({'varredura': (3, 30), 'estado': _E(None)},
                         datetime.datetime(2026, 9, 2, 3, 29)), False)
    igual('na hora, varre',
          hora_de_varrer({'varredura': (3, 30), 'estado': _E(None)},
                         datetime.datetime(2026, 9, 2, 3, 30)), True)
    igual('ja varreu hoje: nao repete',
          hora_de_varrer({'varredura': (3, 30), 'estado': _E('2026-09-02')},
                         datetime.datetime(2026, 9, 2, 23, 0)), False)
    igual('FORA da janela nao varre (subida as 18h com estado zerado)',
          hora_de_varrer({'varredura': (3, 30), 'estado': _E(None)},
                         datetime.datetime(2026, 9, 2, 18, 0)), False)
    igual('dentro da janela recupera atraso (subiu 05:00)',
          hora_de_varrer({'varredura': (3, 30), 'estado': _E(None)},
                         datetime.datetime(2026, 9, 2, 5, 0)), True)
    igual('no limite da janela ainda varre',
          hora_de_varrer({'varredura': (3, 30), 'estado': _E(None), 'janela': 180},
                         datetime.datetime(2026, 9, 2, 6, 30)), True)
    igual('um minuto depois do limite, nao',
          hora_de_varrer({'varredura': (3, 30), 'estado': _E(None), 'janela': 180},
                         datetime.datetime(2026, 9, 2, 6, 31)), False)
    igual('varreu ontem, e esta na janela: varre de novo',
          hora_de_varrer({'varredura': (3, 30), 'estado': _E('2026-09-01')},
                         datetime.datetime(2026, 9, 2, 4, 0)), True)
    igual('desligada nunca varre',
          hora_de_varrer({'varredura': None, 'estado': _E(None)},
                         datetime.datetime(2026, 9, 2, 23, 0)), False)

    print('\nestrategias')
    igual('a implementada de hoje resolve',
          exigir_estrategia('coluna')['implementada'], True)
    for loc_nome in ('binlog', 'change_tracking', 'slot_logico', 'flashback',
                     'reconciliacao'):
        try:
            exigir_estrategia(loc_nome)
            loc_falhas.append('%s: NAO recusou' % loc_nome)
            print('  FALHOU  %s recusa dizendo o que exigiria' % loc_nome)
        except SystemExit as e:
            # A recusa tem de ENSINAR: sem o "exigiria", ela manda a pessoa
            # procurar em vez de resolver.
            if 'exigiria' in str(e) and ESTRATEGIAS[loc_nome]['exige'][:20] in str(e):
                print('  ok      %s recusa dizendo o que exigiria' % loc_nome)
            else:
                loc_falhas.append('%s: recusou sem dizer o que exigiria' % loc_nome)
                print('  FALHOU  %s recusa dizendo o que exigiria' % loc_nome)
    try:
        exigir_estrategia('inventada')
        loc_falhas.append('nome inventado: NAO recusou')
        print('  FALHOU  nome inventado recusa listando as que existem')
    except SystemExit as e:
        igual('nome inventado lista as que existem', 'coluna' in str(e), True)

    igual('so `coluna` esta implementada',
          sorted(n for n, v in ESTRATEGIAS.items() if v['implementada']), ['coluna'])
    igual('e todas as outras dizem se veem DELETE',
          all('ve_delete' in v for v in ESTRATEGIAS.values()), True)

    print('\nhora invalida na linha de comando')
    for loc_ruim in ('25:00', '3h30', '03:70', 'abc'):
        try:
            _hora(loc_ruim)
            loc_falhas.append('--varredura=%s: NAO abortou' % loc_ruim)
            print('  FALHOU  --varredura=%s aborta' % loc_ruim)
        except SystemExit:
            print('  ok      --varredura=%s aborta' % loc_ruim)
    igual('"nao" desliga', _hora('nao'), None)

    print()
    if loc_falhas:
        print('%d FALHA(S):' % len(loc_falhas))
        for f in loc_falhas:
            print('  - %s' % f)
        return 1
    print('tudo passou.')
    return 0


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