#!/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/khtst/src") import torch from khtst.config import Config from khtst.dados.checkpoints import GestorCheckpoints from khtst.geracao import GeradorMultimodal from khtst.qualidade import AgenteEngenheiro from khtst.dados.tokenizador import TokenizadorKHTST from khtst.memoria.orquestrador import OrquestradorSOM from khtst.nucleo.modelo import KHTSTModel from khtst.telemetria.hub import TelemetryHub from khtst.treino.treinador import TreinadorExtensao CORPUS = "/home/z/my-project/khtst/cache_dados/corpus_v2.jsonl" TOKENIZADOR = "/home/z/my-project/khtst/cache_dados/tokenizador.json" RESUMO = "/home/z/my-project/khtst/telemetria_out/resumo_treino_v10.json" RAM_JSONL = "/home/z/my-project/khtst/telemetria_out/ram_treino_v10.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 TEMPO_ACUM = "/home/z/my-project/khtst/telemetria_out/tempo_treino_acumulado.json" def tempo_acumulado_h() -> float: import json as _json try: with open(TEMPO_ACUM) as f: return float(_json.load(f).get("segundos", 0.0)) / 3600.0 except Exception: return 0.0 def registrar_tempo(segundos: float): import json as _json total = tempo_acumulado_h() * 3600.0 + max(0.0, segundos) os.makedirs(os.path.dirname(TEMPO_ACUM), exist_ok=True) with open(TEMPO_ACUM, "w") as f: _json.dump({"segundos": total, "horas": total / 3600.0}, f) def definir_tempo(segundos: float): """v10.1 — FONTE ÚNICA de verdade: sobrescreve o contador com a duração wall do MonitorRAM (jsonl completo). Corrige a dupla contagem v9.0 (soma dos segmentos + duração total = contagem inflada).""" import json as _json os.makedirs(os.path.dirname(TEMPO_ACUM), exist_ok=True) with open(TEMPO_ACUM, "w") as f: _json.dump({"segundos": max(0.0, segundos), "horas": max(0.0, segundos) / 3600.0, "nota": "wall time do MonitorRAM (fonte única)"}, f) 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("--teto_horas", type=float, default=9.0, help="v11 (REQUISITO): TETO RÍGIDO de tempo ACUMULADO " "de treino (h) = orçamento 8 h + tolerância 1 h. " "Ao atingir, o treino para (requisito: parar se " "ultrapassar 9 horas).") ap.add_argument("--continuar_ate_horas", type=float, default=8.0, help="v11 (REQUISITO): orçamento ALVO de treino (h) — " "épocas extras 'sobrando tempo' até esgotá-lo " "(nunca além do teto rígido 8h+1h).") 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)") ap.add_argument("--estados_hf", action="store_true", help="v11 (REQUISITO): envia os estados safetensors ao Hub " "ANTES de 900 MB locais (limiar da config), apaga os " "locais confirmados e permite retomada do Hub — o treino " "continua. Requer HF_TOKEN no ambiente.") args = ap.parse_args() # v8 — TETO RÍGIDO (requisito): parar se o treino acumulado passar de 5h acum = tempo_acumulado_h() if acum >= args.teto_horas: print(f"[TETO] treino acumulado {acum:.2f} h ≥ teto {args.teto_horas:.1f} h — " "PARANDO as tarefas de treino (requisito atendido).") return print(f"teto de treino: {acum:.2f}/{args.teto_horas:.1f} h acumuladas") cfg = Config().de_arquivo("/home/z/my-project/khtst/configs/base.json") # v11 — gestor de estados no Hub (limiar 900 MB + retomada do Hub) if args.estados_hf: cfg.checkpoints.estados_hf["ativo"] = True # 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/khtst/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 = TokenizadorKHTST(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 = KHTSTModel(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 # v9 — ESTADO INTEGRAL: roteador S-SOM + orquestrador SOM nos snapshots # (Teorema 20.1; bug v8: estado fora do state_dict perdido no reload) treinar.preparar_estado_integral() # 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/khtst/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']}") registrar_tempo(resultado.get("duracao_s", 0.0)) print(f"→ tempo acumulado de treino: {tempo_acumulado_h():.2f} h " f"(teto {args.teto_horas:.1f} h)") 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 # ---------- v10: 'sobrando tempo: continuar treinamento' (doc 21) ------ # Épocas EXTRAS enquanto houver orçamento (alvo continuar_ate_horas), # sempre sob o TETO RÍGIDO. Cada época é avaliada; o gate de estado # (doc 20) decide o que é promovido — regressão nunca é publicada. # Persistência: o nº de épocas é gravado no configs/base.json para que a # retomada por segmentos NÃO perca as épocas extras entre invocações. def _persistir_epocas(n: int): cam = "/home/z/my-project/khtst/configs/base.json" try: j = json.load(open(cam, encoding="utf-8")) j["treino"]["epocas"] = n json.dump(j, open(cam, "w", encoding="utf-8"), ensure_ascii=False, indent=2) except Exception as e: print(f"AVISO: não persistiu épocas={n}: {e}") if args.continuar_ate_horas is not None and not args.sem_epocas: while tempo_acumulado_h() < min(args.continuar_ate_horas, args.teto_horas): cfg.treino.epocas += 1 _persistir_epocas(cfg.treino.epocas) # ponteiro p/ a PRIMEIRA época NÃO concluída (evita duplicar a # última época já completa — retomada parcial preservada) treinar.epoca = max(treinar.epoca, len(treinar.historico_epocas)) print(f"== v10: sobrou orçamento ({tempo_acumulado_h():.2f} h < " f"{args.continuar_ate_horas:.2f} h) → ÉPOCA EXTRA " f"{treinar.epoca} (de {cfg.treino.epocas}) ==") modelo.definir_fase_moe("foco") res_extra = treinar.treinar("/home/z/my-project/khtst/checkpoints", orcamento_s=args.orcamento) registrar_tempo(res_extra.get("duracao_s", 0.0)) print(f"época extra: {res_extra['duracao_s']}s | acumulado " f"{tempo_acumulado_h():.2f} h | interrompido={res_extra['interrompido']}") if res_extra["interrompido"]: print("→ orçamento do segmento esgotado no meio da época extra; " "estado salvo — próxima invocação retoma a MESMA época") return if tempo_acumulado_h() >= args.teto_horas: print("→ teto rígido atingido — sem mais épocas extras") break # ---------- fase 3: DPO (com GATE de rollback — Teorema 20.3) ---------- resultado_dpo = {"pulado": True} evento_gate_dpo = None if not args.sem_dpo: print("== fase DPO ==") pre_dpo = treinar.avaliar(max_lotes=3) treinar.salvar_checkpoint("pre-dpo", {"gate": {"fase": "pre-dpo", **pre_dpo}}) resultado_dpo = treinar.treinar_dpo() print("DPO:", resultado_dpo) treinar.salvar_checkpoint("fase-dpo", {"dpo": resultado_dpo}) evento_gate_dpo = treinar.gate_pos_fase("dpo", "pre-dpo", "fase-dpo", pre_dpo) print("GATE DPO:", evento_gate_dpo) # ---------- fase 4: SOM com máximo de neurônios ativos (idem gate) ------ resultado_som = {"pulado": True} evento_gate_som = None if not args.sem_som: print("== fase SOM (máx neurônios ativos) ==") pre_som = treinar.avaliar(max_lotes=3) treinar.salvar_checkpoint("pre-som", {"gate": {"fase": "pre-som", **pre_som}}) 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"}}) evento_gate_som = treinar.gate_pos_fase("som", "pre-som", "fase-som", pre_som) print("GATE SOM:", evento_gate_som) # v9 — avaliação final no objeto vivo (pós-rollback, se houve) + honesta aval_final = treinar.avaliar(max_lotes=4) print("avaliação final (vivo):", aval_final) aval_recarregado = treinar._avaliar_estado_do_disco("fase-som", max_lotes=4) \ if not args.sem_som else treinar._avaliar_estado_do_disco("fase-dpo", max_lotes=4) print("avaliação final (RECARREGADA do disco — honesta):", aval_recarregado) # 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) # v12 — CICLO PDCA com FERRAMENTAS REAIS (órfãos conectados: raciocinio/ # ciclo.py + ferramentas.py): o relatório final passa por um ciclo # Planejar→Distribuir→Executar→Avaliar com a calculadora no hub # (ferramenta executada de fato, não módulo morto) evid_pdca = None try: from khtst.raciocinio.ferramentas import HubFerramentas from khtst.raciocinio.ciclo import CicloPDCA, Traco hub_fer = HubFerramentas() hub_fer.registrar_padrao(busca_memoria=None) pdca = CicloPDCA(modelo, tk, hub_fer, cfg, hub=hub) traco = Traco() plano = pdca.planejar("lm", "revisao das métricas finais", traco) _, conf = pdca.executar("métricas finais", plano.get("acao", "resumo"), traco) calc = hub_fer.executar("calculadora", f"{treinar.passo_global} * 1.0") evid_pdca = {"planejado": plano.get("acao"), "confianca": round(conf, 4), "ferramenta_calculadora": calc[:60], "estatisticas_hub": {"hits": hub_fer.stats.hits, "misses": hub_fer.stats.misses}} print("PDCA v12:", evid_pdca["planejado"], "conf=", evid_pdca["confianca"]) except Exception as e: evid_pdca = {"erro": f"{type(e).__name__}: {str(e)[:120]}"} 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: # v9.1 — contador final = duração wall do monitor (sem dupla contagem) definir_tempo(ram_resumo["duracao_s"] if isinstance(ram_resumo, dict) else 0.0) json.dump({"versao": "v12-khtst-tronco-mamba3-muuonclip", "teto_horas": args.teto_horas, "tempo_acumulado_h": tempo_acumulado_h(), "memoria_v8": treinar.gestor_memoria.estatisticas(), "autoajuste_v9": getattr(treinar, "ultimo_autoajuste", None), "ciar": {k: (round(v, 6) if isinstance(v, float) else v) for k, v in treinar.ciar.items()}, "gate_fase": {"dpo": evento_gate_dpo, "som": evento_gate_som, "eventos": treinar.eventos_gate}, "avaliacao_recarregada_honesta": aval_recarregado, "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, "pdca_ferramentas": evid_pdca, "muonclip": {"tau_qk": treinar.tau_qk, "eventos_qk_clip": treinar.qk_clip_eventos[-20:]}, "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()