File size: 16,130 Bytes
094b8e7
 
ebc1354
 
 
 
 
 
 
 
 
 
 
 
 
05db690
 
 
 
 
 
094b8e7
05db690
094b8e7
05db690
094b8e7
7eea4b4
96a4177
094b8e7
 
 
 
 
 
96a4177
 
 
094b8e7
 
 
 
 
 
05db690
094b8e7
7eea4b4
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
094b8e7
 
05db690
094b8e7
05db690
094b8e7
05db690
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
094b8e7
 
 
 
 
05db690
 
094b8e7
05db690
 
7eea4b4
 
 
 
 
 
 
 
 
 
094b8e7
 
05db690
7eea4b4
 
 
05db690
 
7eea4b4
 
 
 
05db690
094b8e7
05db690
 
 
 
 
094b8e7
 
05db690
094b8e7
96a4177
 
 
05db690
96a4177
 
 
 
 
 
7eea4b4
 
96a4177
 
 
 
 
7eea4b4
96a4177
 
 
05db690
 
 
 
7eea4b4
 
 
 
 
 
 
 
 
 
 
 
 
05db690
 
 
 
 
 
 
96a4177
05db690
 
 
 
 
 
 
 
96a4177
05db690
 
 
 
 
96a4177
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
ebc1354
 
 
 
 
 
 
 
 
 
 
 
 
 
96a4177
 
05db690
7eea4b4
05db690
7eea4b4
e829bfc
05db690
 
 
 
e829bfc
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
05db690
 
ebc1354
06edf75
 
 
96a4177
 
 
 
 
7eea4b4
 
05db690
 
7eea4b4
094b8e7
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""Treino v4 — pipeline completo em DUAS FASES (doc 15):

  FASE A — rede aumentada (épocas 0..N-1, sem SOM):
           FASE Densa (épocas 0..2): máx conexões, MoE softmax-total, MTP com
           K ADAPTATIVO, PCGrad, ABMO, PRS V8, AgenteConfiança, MEMÓRIA INTERNA,
           GESTOR DE CRESCIMENTO + AUTO-ESCALA POR COMPUTO (doc 12); unidades
           NLP/NLG acopladas (nascem neutras — Teorema 15.1);
           FASE Foco (época 3): MoE top-k + máscara + QAT-W8A8 final
           (escalas POR GRUPO + correção de viés — doc 08 §§5-6);
           FASE DPO: pares (dourado, corrompido) com β adaptativo.
  FASE B — consolidação SOM com MÁXIMO de neurônios ativos (eq. 15.1):
           8 variantes SEM neurônios isolados (reseeding garantido),
           lr do tronco ×0.1, alvo A ≥ 0.90 (Teorema 15.2).

Uso (CPU em segmentos):
  python3 04_treinar.py --orcamento 1500        # segmento de 25 min
  python3 04_treinar.py --orcamento 1500        # retoma do último estado do HF
  ... quando 'pipeline completo' termina, roda DPO+SOM automaticamente.
"""
import argparse
import base64
import json
import os
import sys
import threading
import time

sys.path.insert(0, "/home/z/my-project/aurora/src")

import torch

from aurora.config import Config
from aurora.dados.checkpoints import GestorCheckpoints
from aurora.geracao import GeradorMultimodal
from aurora.qualidade import AgenteEngenheiro
from aurora.dados.tokenizador import TokenizadorAurora
from aurora.memoria.orquestrador import OrquestradorSOM
from aurora.nucleo.modelo import AURORAModel
from aurora.telemetria.hub import TelemetryHub
from aurora.treino.treinador import TreinadorExtensao

CORPUS = "/home/z/my-project/aurora/cache_dados/corpus_v2.jsonl"
TOKENIZADOR = "/home/z/my-project/aurora/cache_dados/tokenizador.json"
RESUMO = "/home/z/my-project/aurora/telemetria_out/resumo_treino_v7.json"
RAM_JSONL = "/home/z/my-project/aurora/telemetria_out/ram_treino_v7.jsonl"


class MonitorRAM:
    """v7 — observação do CONSUMO DE MEMÓRIA RAM durante o treino inteiro
    (pedido explícito do Agente Engenheiro). Amostra VmRSS do processo e
    MemAvailable do sistema a cada `intervalo_s` em thread daemon; grava JSONL
    (t, passo, rss_mb, disponivel_mb) e resume pico/média/mínimo no fim."""

    def __init__(self, caminho: str, passo_ref, intervalo_s: float = 2.0):
        self.caminho = caminho
        self.passo_ref = passo_ref
        self.intervalo_s = intervalo_s
        self._parar = threading.Event()
        self._thread = threading.Thread(target=self._laco, daemon=True)
        self.amostras: list[dict] = []

    @staticmethod
    def _rss_mb() -> float:
        try:
            with open("/proc/self/status", encoding="utf-8") as f:
                for linha in f:
                    if linha.startswith("VmRSS:"):
                        return float(linha.split()[1]) / 1024.0
        except Exception:
            pass
        return 0.0

    @staticmethod
    def _disponivel_mb() -> float:
        try:
            with open("/proc/meminfo", encoding="utf-8") as f:
                for linha in f:
                    if linha.startswith("MemAvailable:"):
                        return float(linha.split()[1]) / 1024.0
        except Exception:
            pass
        return 0.0

    def _laco(self):
        # modo "a": segmentos separados do pipeline acumulam no MESMO jsonl
        with open(self.caminho, "a", encoding="utf-8") as f:
            while not self._parar.is_set():
                amostra = {"t": round(time.time(), 2),
                           "rss_mb": round(self._rss_mb(), 1),
                           "disponivel_mb": round(self._disponivel_mb(), 1),
                           "passo": self.passo_ref()}
                f.write(json.dumps(amostra) + "\n")
                f.flush()
                self._parar.wait(self.intervalo_s)

    def iniciar(self):
        self._thread.start()

    def parar(self) -> dict:
        """Encerra a thread e resume o jsonl INTEIRO (todos os segmentos)."""
        self._parar.set()
        self._thread.join(timeout=5)
        amostras: list[dict] = []
        try:
            with open(self.caminho, encoding="utf-8") as f:
                for linha in f:
                    try:
                        amostras.append(json.loads(linha))
                    except Exception:
                        continue
        except FileNotFoundError:
            return {"amostras": 0}
        rss = [a["rss_mb"] for a in amostras]
        disp = [a["disponivel_mb"] for a in amostras]
        if not rss:
            return {"amostras": 0}
        return {"amostras": len(rss),
                "rss_min_mb": round(min(rss), 1),
                "rss_med_mb": round(sum(rss) / len(rss), 1),
                "rss_max_mb": round(max(rss), 1),
                "disponivel_min_mb": round(min(disp), 1),
                "duracao_s": round(amostras[-1]["t"] - amostras[0]["t"], 1)}


def carregar_registros(caminho: str) -> list[dict]:
    registros = []
    with open(caminho, encoding="utf-8") as f:
        for linha in f:
            try:
                r = json.loads(linha)
            except Exception:
                continue
            # multimodal volta a bytes
            if isinstance(r.get("imagem"), str):
                try:
                    r["imagem"] = base64.b64decode(r["imagem"])
                except Exception:
                    r["imagem"] = None
            if isinstance(r.get("audio"), str):
                try:
                    r["audio"] = base64.b64decode(r["audio"])
                except Exception:
                    r["audio"] = None
            registros.append(r)
    return registros


def main():
    ap = argparse.ArgumentParser()
    ap.add_argument("--orcamento", type=float, default=None,
                    help="segundos por segmento (CPU: retomada automática)")
    ap.add_argument("--max_passos_epoca", type=int, default=None)
    ap.add_argument("--sem_dpo", action="store_true")
    ap.add_argument("--sem_som", action="store_true")
    ap.add_argument("--sem_epocas", action="store_true",
                    help="v7: pula o laço de épocas (retomada já completa) e vai "
                         "direto a DPO→SOM→avaliação→resumo")
    ap.add_argument("--sem_difusao", action="store_true",
                    help="v7: pula o ciclo geracional de difusão (economiza RAM/tempo "
                         "do segmento final; evidência já validada em v5/v6)")
    ap.add_argument("--semente", type=int, default=None,
                    help="v7: semente torch p/ REPRODUTIBILIDADE da inicialização "
                         "(diagnóstico v7: init não semeada + lr 0.003 pode divergir; "
                         "padrão = cfg.dados.semente)")
    args = ap.parse_args()

    cfg = Config().de_arquivo("/home/z/my-project/aurora/configs/base.json")
    # v7 — reprodutibilidade: init de pesos, dropout e amostragens torch
    torch.manual_seed(args.semente if args.semente is not None
                      else cfg.dados.semente)
    os.makedirs("/home/z/my-project/aurora/telemetria_out", exist_ok=True)
    hub = TelemetryHub(cfg.telemetria.arquivo_jsonl)
    # v7 — Agente Engenheiro: observação do consumo de RAM durante TODO o treino
    ram = MonitorRAM(RAM_JSONL, passo_ref=lambda: 0)
    ram.iniciar()
    print("monitor RAM: amostrando", RAM_JSONL)
    tk = TokenizadorAurora(TOKENIZADOR)
    registros = carregar_registros(CORPUS)
    print(f"corpus: {len(registros)} registros")
    cont = {}
    for r in registros:
        cont[r["tarefa"]] = cont.get(r["tarefa"], 0) + 1
    print("tarefas:", cont)

    modelo = AURORAModel(cfg, usar_multimodal=True)
    print(f"modelo: {sum(p.numel() for p in modelo.parameters())/1e6:.2f}M parâmetros")

    checkpoints = GestorCheckpoints(dir_local=cfg.checkpoints.dir_local)
    orquestrador = OrquestradorSOM(cfg.modelo.d_modelo, cfg.som, hub=hub,
                    cfg_atencao=vars(cfg.som_atencao) if hasattr(cfg, "som_atencao") else None)
    treinar = TreinadorExtensao(cfg, modelo, tk, hub, registros,
                                orquestrador_som=orquestrador, checkpoints=checkpoints)
    # v5 — Agente Engenheiro: observação integral durante o treino (doc 16 §9)
    agente = AgenteEngenheiro(
        tolerancia_regressao=cfg.agente_engenheiro.tolerancia_regressao, hub=hub)
    treinar.agente_engenheiro = agente
    treinar.intervalo_agente = cfg.agente_engenheiro.intervalo
    # passo global do treino visível ao monitor de RAM
    ram.passo_ref = lambda: treinar.passo_global
    # v5 — gerador multimodal (difusão) para o ciclo geracional (doc 17)
    difusao = GeradorMultimodal(repo_id=cfg.difusao.repo_id,
                                altura=cfg.difusao.altura, largura=cfg.difusao.largura,
                                passos_inferencia=cfg.difusao.passos_inferencia,
                                guia_escala=cfg.difusao.guia_escala, hub=hub) \
        if cfg.difusao.ativo and not args.sem_difusao else None

    # item c: "salvar (e usar) estados" — retoma do checkpoint local mais avançado
    tag_retomada = treinar.retomar_ultimo_checkpoint()
    print(f"retomada do Hub: tag={tag_retomada} "
          f"(passo {treinar.passo_global}, época {treinar.epoca})")

    # ---------- fases 1+2: denso → foco (segmentos com retomada) ----------
    if not args.sem_epocas:
        while True:
            resultado = treinar.treinar("/home/z/my-project/aurora/checkpoints",
                                        max_passos_por_epoca=args.max_passos_epoca,
                                        orcamento_s=args.orcamento)
            print(f"segmento: {resultado['duracao_s']}s | passos={treinar.passo_global} "
                  f"| épocas={len(treinar.historico_epocas)} | interrompido={resultado['interrompido']}")
            if not resultado["interrompido"]:
                break
            print("→ estado salvo no Hub; rode o script novamente para o próximo segmento")
            if args.orcamento is not None:
                return                 # sai: próximo segmento em nova invocação
            return                     # sem orçamento: uma passada só por chamada

    # ---------- fase 3: DPO ----------
    resultado_dpo = {"pulado": True}
    if not args.sem_dpo:
        print("== fase DPO ==")
        resultado_dpo = treinar.treinar_dpo()
        print("DPO:", resultado_dpo)
        treinar.salvar_checkpoint("fase-dpo", {"dpo": resultado_dpo})

    # ---------- fase 4: SOM com máximo de neurônios ativos ----------
    resultado_som = {"pulado": True}
    if not args.sem_som:
        print("== fase SOM (máx neurônios ativos) ==")
        resultado_som = treinar.consolidar_som_max_ativos()
        for nome, st in resultado_som.get("variantes", {}).items():
            print(f"  {nome}: {st}")
        treinar.salvar_checkpoint("fase-som", {"som": {
            k: v for k, v in resultado_som.items() if k != "prototipos"}})

    aval_final = treinar.avaliar(max_lotes=4)
    print("avaliação final:", aval_final)

    # v5 — evidência de INFERÊNCIA para o Agente Engenheiro (métricas de geração)
    try:
        ids_prompt = tk.encode("Pergunta: o que e o Kohonen? Resposta:",
                               tarefa="instrucao", max_len=24)
        t0 = time.time()
        tokens = modelo.gerar(torch.tensor([ids_prompt]), max_novos=24)
        evid_inf = agente.observar_inferencia(modelo, torch.tensor([ids_prompt]),
                                              tokens, time.time() - t0)
        print("inferência (agente):", evid_inf)
        print("texto gerado:", tk.decode(tokens))
    except Exception as e:
        evid_inf = {"erro": str(e)}

    # v5 — ciclo de compreensão geracional (difusão; modo honesto)
    evid_difusao = None
    if difusao is not None:
        try:
            r = difusao.compreensao_geracional(
                "a floresta amazonica de manha", modelo,
                orquestrador=orquestrador, op="txt2img")
            evid_difusao = {"modo": r.get("modo"), "op": r.get("op"),
                            "latencia_s": r.get("latencia_s"),
                            "prompt_enriquecido": r.get("prompt_enriquecido"),
                            "gerou_pixels": r.get("modo") == "real",
                            "estatisticas": difusao.estatisticas()}
            print("difusão (modo honesto):", evid_difusao["modo"],
                  "—", evid_difusao["prompt_enriquecido"])
        except Exception as e:
            evid_difusao = {"erro": str(e)}

    # v4 — MÉTRICAS COMPLETAS (doc 15 §3): neurônios ativos por variante,
    # NLP/NLG, janela 1M, computo, crescimento
    metricas_completas = {
        "neuronios_ativos": (orquestrador.telemetria_ativos()
                             if orquestrador is not None else {}),
        "estruturais": modelo.metricas_estruturais(),
        "escala": modelo.info_escalacao,
        "regime": {"fase_final": treinar.regime.fase,
                   "alvo_ativos": treinar.regime.alvo_ativos},
    }
    print("== métricas completas (doc 15 §3) ==")
    print(json.dumps(metricas_completas, ensure_ascii=False, indent=1,
                     default=str)[:2200])

    rel_agente = agente.relatorio(orquestrador_som=orquestrador,
                                  avaliacao=aval_final)
    os.makedirs(os.path.dirname(RESUMO), exist_ok=True)
    ram_resumo = ram.parar()   # resume TODOS os segmentos (jsonl completo)
    with open(RESUMO, "w", encoding="utf-8") as f:
        json.dump({"versao": "v7-retreino-observado-ram",
                   "epocas": treinar.historico_epocas,
                   "passos": treinar.passo_global,
                   "punicoes": len(treinar.eventos_punicao),
                   "barreira_abmo_ativacoes": treinar.barreira_ativa_count,
                   "pcgrad_ultimo": treinar.ultimos_pcgrad,
                   "prs_v8": treinar.prs.estado(),
                   "confianca": {"h_ema": treinar.agente.ultimo_h,
                                 "ece": treinar.agente.ece(),
                                 "posterior_media": treinar.agente.posterior.media,
                                 "bernstein": treinar.agente.posterior.bernstein()},
                   "memoria": (treinar.memoria.estatisticas()
                               if treinar.memoria is not None else None),
                   "crescimento": (treinar.gestor.eventos
                                   if treinar.gestor is not None else []),
                   "microunidades": [
                       {"ramos": c.nome_ramos,
                        "ativos": [float(a) for a in c.ramo_ativo.tolist()],
                        "uso_ema": (c.uso_ema.tolist()
                                    if c.uso_ema is not None else None),
                        "passos_rec": c.ultimo_passos_rec}
                       for c in modelo.compositores],
                   "dpo": resultado_dpo, "som": resultado_som,
                   "avaliacao_final": aval_final,
                   "metricas_completas": metricas_completas,
                   "rpp": (treinar.rpp.estado() if treinar.rpp else None),
                   "roteador_ssom": (treinar.roteador.estatisticas()
                                     if treinar.roteador is not None else None),
                   "inferencia_agente": evid_inf,
                   "difusao": evid_difusao,
                   "agente_engenheiro": {"aprovadas": rel_agente["revisoes_aprovadas"],
                                          "bloqueadas": rel_agente["revisoes_bloqueadas"],
                                          "invariante_sem_orfaos": rel_agente.get("som_invariante_sem_orfaos")},
                   "checkpoints_eventos": checkpoints.eventos[-12:],
                   "ram": ram_resumo},
                  f, ensure_ascii=False, indent=2, default=str)
    print(f"resumo → {RESUMO}")
    print("ram:", json.dumps(ram_resumo))


if __name__ == "__main__":
    main()