"""train_v6_5_v2.py — V6.5-V2-metrics (2 Fases + Métricas SOM canônicas). ═══════════════════════════════════════════════════════════════════════════════ V6.5-V2-metrics — MÉTRICAS SOM CANÔNICAS + 2 FASES + ESTADO FASE1→FASE2 ═══════════════════════════════════════════════════════════════════════════════ User requirements (latest, V6.5-V2-metrics): 1. Organização e não duplicidade de módulos (versões velhas → deprecados/). 2. Atualizações HF em LOTE (sobrescrever antigos desatualizados). 3. Limpar Armazenamento e worklog. 4. HF_TOKEN apagada do ambiente virtual e dos scripts enviados ao HF após uso. 5. Usar streaming_datasets.py e ativar xeon_runtime.py. 6. FASE 1 CONHECIMENTO (meta mínima 8000+ samples, 100 em 100) na sequência exata de 8 datasets. 7. FASE 2 TREINAMENTO COM PUNIÇÃO ATIVA (meta mínima 2000+ samples, 100 em 100) sobre BrunoN-Dev/corpus-ptbr-v1. 8. FASE2 NUNCA acontece antes da FASE1. FASE2 carrega explicitamente o estado do modelo concluído da FASE1. 9. FASE2 deve treinar, ajustar dados e parâmetros do MAPA-SOM e INFORMAR MÉTRICAS DE APRENDIZADO do modelo. 10. Métricas SOM canônicas (8 métricas) monitoradas especialmente na FASE2 PUNITIVA: 10.1. Erro de Quantização (QE) — distância média ‖x_i - W_BMU(x_i)‖ 10.2. Erro Topológico (TE) — % de pares BMU/2ª-BMU não adjacentes 10.3. Erro de Kaski-Lagus — combinação ponderada de QE + TE 10.4. Variância Explicada — 1 - Var(residual)/Var(total) 10.5. Colapso Topológico — pesos convergem para ponto ou linha 10.6. Taxa de Neurônios Mortos — % de neurônios nunca BMU 10.7. Estagnação do QE — slope ≈ 0 com QE em patamar elevado 10.8. Cruzamento de Vizinhança — TE oscila ou cresce no final 11. Limpeza de RAM e Armazenamento aprimorada (lógica + matemática). 12. Não reduzir tempo, não gerar dados sintéticos. 13. Após concluir: enviar scripts+arquivos+estado do modelo testados e aprovados ao HF (sobrescrevendo antigos). V6.5-V2-metrics (alterações vs V6.5-V2-dynamic): - NOVO: Módulo som_metrics.py com 8 métricas canônicas implementadas. - NOVO: KohonenLearningSystemV2.compute_som_metrics() integra as métricas. - NOVO: Histórico de métricas SOM (SOMMetricHistory) para detecção de estagnação do QE e cruzamento de vizinhança. - NOVO: FASE1 computa métricas após cada chunk (8000 samples). - NOVO: FASE2 computa métricas após CADA BATCH (2000 samples, mais frequentemente por ser fase punitiva). - NOVO: load_fase1_state_into_kls() — carrega explicitamente estado da FASE1 antes de iniciar FASE2. - NOVO: aggressive_memory_cleanup() reporta RSS antes/depois. - NOVO: aggressive_storage_cleanup() limpa __pycache__ + .tmp + logs >7dias. - META: CONHECIMENTO 6000 → 8000; PUNIÇÃO 1000 → 2000. """ from __future__ import annotations import gc import json import logging import math import os import re import shutil import sys import time import traceback from datetime import datetime from pathlib import Path from typing import Any, Dict, Iterator, List, Optional, Tuple # ============================================================================ # 0. Paths e logging # ============================================================================ PROJECT_ROOT = Path("/home/z/my-project") BIGRU_ROOT = PROJECT_ROOT / "BiGRU_T_version" SRC_ROOT = BIGRU_ROOT / "src" DOWNLOAD_DIR = PROJECT_ROOT / "download" DOWNLOAD_DIR.mkdir(parents=True, exist_ok=True) REPORT_PATH = BIGRU_ROOT / "v6_5_v2_report.json" METRICS_PATH = BIGRU_ROOT / "v6_5_v2_training_metrics.json" MODULE_ANALYSIS_PATH = BIGRU_ROOT / "v6_5_v2_module_analysis.json" SCRIPT_ACTIVITY_PATH = BIGRU_ROOT / "v6_5_v2_script_activity.json" ATTENTION_EVAL_PATH = BIGRU_ROOT / "v6_5_v2_attention_eval.json" USER_QUESTIONS_PATH = BIGRU_ROOT / "v6_5_v2_user_questions.json" MODEL_STATES_PATH = BIGRU_ROOT / "v6_5_v2_model_states.pt" PREDICT_FIX_EVAL_PATH = BIGRU_ROOT / "v6_5_v2_predict_fix_eval.json" V2_PHASES_EVAL_PATH = BIGRU_ROOT / "v6_5_v2_phases_eval.json" HYPOTHESES_EVAL_PATH = BIGRU_ROOT / "v6_5_v2_hypotheses_eval.json" INFERENCE_PUNISHMENT_PATH = BIGRU_ROOT / "v6_5_v2_inference_punishment.json" logging.basicConfig( level=logging.INFO, format="[%(asctime)s] [%(levelname)s] %(message)s", datefmt="%H:%M:%S", ) logger = logging.getLogger("train_v6_5_v2") # ============================================================================ # 1. ATIVAR xeon_runtime.py + FORÇAR V65_ENABLE_STREAMING=1 # ============================================================================ # User requirement: "usar streaming_datasets.py e ativar xeon_runtime.py" os.environ["V65_ENABLE_STREAMING"] = "1" logger.info(f"[V6.5-V2] V65_ENABLE_STREAMING={os.environ['V65_ENABLE_STREAMING']} (forced)") # V6.5-V2-metrics-FIX-4 — HF datasets memory optimization (OOM-killer mitigation) # User requirement: "o processo vem sendo morto OOM-kiler (Out of memory) devido # algum bug de lógica ou falta de otimização que deve ser investigado". # Estas variáveis reduzem o consumo de memória do HF datasets library: # - HF_DATASETS_DISABLE_IN_MEMORY_CACHE: não cacheia datasets em RAM # - HF_DATASETS_OFFLINE=0: permite streaming mas não força cache local # - DATASETS_FINGERPRINT_CACHING_DISABLED: skipa fingerprinting (CPU/memory) # - TOKENIZERS_PARALLELISM=false: evita spawn de processos paralelos # - HF_HUB_DISABLE_TELEMETRY: desabilita telemetria (CPU/network) os.environ["HF_DATASETS_DISABLE_IN_MEMORY_CACHE"] = "1" os.environ["DATASETS_FINGERPRINT_CACHING_DISABLED"] = "1" os.environ["TOKENIZERS_PARALLELISM"] = "false" os.environ["HF_HUB_DISABLE_TELEMETRY"] = "1" os.environ.setdefault("HF_DATASETS_CACHE", "/tmp/hf_datasets_cache_v65") # Cria o dir de cache se não existir (limpa cache antigo periodicamente) try: cache_dir = Path(os.environ["HF_DATASETS_CACHE"]) cache_dir.mkdir(parents=True, exist_ok=True) except Exception: pass sys.path.insert(0, str(SRC_ROOT)) from bigru_t.utils.xeon_runtime import ( # noqa: E402 optimize_xeon_environment, benchmark_fp16_matmul, get_xeon_status, ) # V6.5-V2-buffer-864 — Ativa OomGuard (thread daemon que monitora VmRSS e # dispara gc.collect() agressivo quando memória se aproxima do limite). # User requirement: "o processo vem sendo morto OOM-kiler (Out of memory) # devido algum bug de lógica ou falta de otimização que deve ser investigado # (acrescentar exceptions e melhor detecção de falhas de lógica e erros de script)". # OomGuard complementa os try/except MemoryError no loop principal: atua em # background entre checks (a cada 2s) para coletar lixo antes que o kernel # OOM-killer interveja. from bigru_t.utils.oom_guard import OomGuard # noqa: E402 # V7 — Integrator for AdaptiveDPOLoss + AdvancedSOMAugmenter in FASE2 try: from bigru_t.training.v7_fase2_integrator import V7Fase2Integrator _V7_INTEGRATOR_AVAILABLE = True except ImportError as _v7_err: logger.warning(f"V7 integrator não disponível: {_v7_err}") _V7_INTEGRATOR_AVAILABLE = False V7Fase2Integrator = None N_CORES = optimize_xeon_environment(verbose=True) XEON_STATUS = get_xeon_status() FP16_BENCH = benchmark_fp16_matmul(size=4000, warmup=1, iters=2) logger.info(f"[V6.5-V2] Xeon FP16 benchmark: {FP16_BENCH}") # V6.5-V2-buffer-864 — Inicializa OomGuard com threshold conservador para # cgroup de 4GB. max_rss_mb=2500 (58% do cgroup) → gc agressivo quando exceder. # warn_rss_mb=2000 (47%) → gc preventivo. check_interval=2s → monitoramento # frequente sem impacto significativo em performance. OOM_GUARD = OomGuard(max_rss_mb=2500, warn_rss_mb=2000, check_interval=2.0) OOM_GUARD.start() logger.info( f"[V6.5-V2-buffer-864] OomGuard started: " f"max_rss={OOM_GUARD.max_rss_mb}MB, warn_rss={OOM_GUARD.warn_rss_mb}MB, " f"interval={OOM_GUARD.check_interval}s" ) # ============================================================================ # 2. Configurações V6.5-V2-metrics-FIX (864 neurons, params restaurados a valores maiores) # ============================================================================ # User requirement (latest): "ATENÇÃO: retornar os parâmetros do modelo para # os valores maiores e resolver falhas de lógica e de bugs que estejam causando # alto consumo de memória sem distorcer a arquitetura Kohonen." # # Valores restaurados (V6.5-V2-metrics-FIX): # HIDDEN_DIM: 256 → 1024 (restaurado) # VOCAB_SIZE: 4096 → 16384 (restaurado) # N_HYPOTHESES: 8 → 16 (restaurado) # MAX_N_HYPOTHESES: 16 → 32 (restaurado) # HYP_TRAIN_STEPS: 20 → 30 (restaurado) # HYP_HIDDEN_DIM: 128 → 256 (restaurado) # # Para compensar o aumento de memória (sem distorcer Kohonen): # 1. BUG FIX: _compute_classification_loss_with_som estava dentro de # torch.no_grad(), matando o gradiente dos deltas. Corrigido. # 2. VETORIZAÇÃO: train_hypotheses agora avalia 16 hipóteses em paralelo # via batched matmul (n_hyp, P, N) em vez de loop que materializava # 16 cópias do SOM. # 3. aggressive_cleanup(): zera gradientes Adam, trunca históricos, # gc.collect() em 3 passes. # 4. Buffer sliding window: MAX_BUFFER_SIZE = 256 (preserva contexto recente # sem crescer indefinidamente). # 5. VQ-VAE-2 e W8A8 desabilitados por padrão (não essenciais para Kohonen). BATCH_SIZE = 16 MAX_SEQ_LEN = 8 # V6.5-V2-metrics-FIX-3: Parâmetros restaurados aos valores canônicos # exigidos pelo usuário (sem redução por OOM — as correções de memory leak # no KLS já tornam os valores maiores viáveis dentro do cgroup de 4GB). # User requirement (FIX-3): "SEMPRE MANTER PARÂMETROS HIDDEN_DIM 1024, # VOCAB_SIZE 16384, N_HYPOTHESES 16, MAX_N_HYPOTHESES 32, HYP_TRAIN_STEPS 30, # HYP_HIDDEN_DIM 256. Grid SOM (6,6,6,4)=864 neurônios". # # Memory budget (estimativa real após correções FIX-3): # Embedding: 16384 × 1024 × 4 = 67MB # Attention: 4 × 1024² × 4 = 16MB # HypothesisEnsemble (32 pre-alloc, 16 active): 32×(864×256+256×256+256×3456) = 32×1.16M = 148MB # Adam state (ensemble, 16 active): 2×16×1.16M×4 = 149MB # HypothesisClassifier: 864→512→256→128→64→32→16→8→1 = ~616K = 2.5MB # VQ-VAE-2: ~50K = 0.2MB # SOM weights: 864×4×4 = 14KB # Buffer (256 samples): ~4KB # Total model: ~380MB # + Python + PyTorch + HF datasets cache: ~1-2GB # + Training intermediates (com no_grad aplicado): ~200MB # Total: ~2-3GB (dentro de 4GB com folga graças ao no_grad no VQ-VAE-2) # V6.5-V3-no-regression: parâmetros canônicos restaurados conforme user requirement # "SEMPRE MANTER PARÂMETROS HIDDEN_DIM 1024, VOCAB_SIZE 16384, N_HYPOTHESES 16, # MAX_N_HYPOTHESES 32, HYP_TRAIN_STEPS 30, HYP_HIDDEN_DIM 256. Grid SOM # (4,4,4,4)=256 neurônios". Buffer canônico = 864 (user: "aumentar buffer # para 864 amostras"); fallback = 256 se RSS > 75% cgroup. HIDDEN_DIM = 1024 VOCAB_SIZE = 16384 SOM_GRID = (4, 4, 4, 4) # 256 neurons (CANÔNICO — user requirement explícito) T_MAX = 10000 N_START = 10 LAMBDA_EWC = 0.02 # α₀ e σ₀ conforme especificação canônica Kohonen para grid (4,4,4,4): # α₀ ∈ [0.5, 1.0] para fase de ordenação (rough training) # σ₀ = metade da maior dimensão da grade = max(4,4,4,4)/2 = 2.0 ALPHA0 = 0.5 SIGMA0 = 2.0 DIM_CHOICE = "y" # V6.6 — User requirement: "aumentar o truncamento de textos (line 993 do # train script) de 200 para 1000+ chars, permitindo mais diversidade de # pares byte-level no BBPE". # Justificativa matemática: com 200 chars, o BBPE byte-level produz ~200 # tokens (1 token/byte), limitando a diversidade de pares observados pelo # algoritmo de merges. Com 1000+ chars, a diversidade de pares byte-level # aumenta ~5×, permitindo que o BBPE aprenda merges mais representativos # e efetivamente utilize o vocabulário alvo de 16384 tokens (em vez de # ficar limitado a ~293 palavras únicas como na versão word-level). TEXT_TRUNCATION_CHARS = 1000 # V2-dynamic — HypothesisEnsemble parameters (valores canônicos restaurados) N_HYPOTHESES = 16 # V6.5-V2-metrics-FIX-3: restored from 8 MAX_N_HYPOTHESES = 32 # V6.5-V2-metrics-FIX-3: restored from 16 MIN_N_HYPOTHESES = 4 # limite inferior dinâmico N_TRIALS = 3 # inicial MIN_N_TRIALS = 1 MAX_N_TRIALS = 6 HYP_TRAIN_STEPS = 30 # mantido (canônico) MIN_HYP_TRAIN_STEPS = 10 MAX_HYP_TRAIN_STEPS = 80 HYP_LR = 1e-4 HYP_HIDDEN_DIM = 256 # V6.5-V2-metrics-FIX-3: restored from 128 LOSS_HISTORY_WINDOW = 8 PUNISHMENT_WINDOW = 12 # V2-dynamic-memory — Buffer sliding window (evita OOM em treino longo) # V6.5-V4-canonical-256 (user requirement EXATO): "fazer (tornar canônico) # buffer 256 e grid para (4,4,4,4)=256". Buffer e grid agora têm o MESMO # tamanho (256), eliminando o desbalanceamento que causava OOM em V6.5-V3 # (buffer=864 com grid=256). A correspondência 1:1 entre amostras no buffer # e neurônios no grid 4D é matematicamente elegante — cada amostra pode, # em média, ativar um neurônio distinto, maximizando a utilização do mapa. # # Configuração ADAPTATIVA com monitoramento de memória: # 1. CANÔNICO: buffer=256 + grid (4,4,4,4)=256 — correspondência 1:1. # 2. FALLBACK: se RSS > 75% cgroup, reduz buffer para 128. Grid (4,4,4,4)=256 # PERMANECE (não é reduzido — é o canônico). Apenas o buffer encolhe. # 3. O monitoramento é feito em get_cgroup_memory_limit_mb() no runtime. # OOM-safety garantida por: (a) VQ-VAE-2 lazy compression (a cada 16 add_data), # (b) torch.no_grad() em todo compressão, (c) gc.collect() a cada 4 batches, # (d) OomGuard thread daemon, (e) fallback buffer=128 se RSS > 75%, # (f) exception handler MemoryError + RuntimeError(out of memory). MAX_BUFFER_SIZE = 256 # CANÔNICO V6.5-V4 (user: "tornar canônico buffer 256") FALLBACK_BUFFER_SIZE = 128 # fallback OOM (reduzido de 256 para evitar pressão) FALLBACK_SOM_GRID = (4, 4, 4, 4) # = 256 neurônios (CANÔNICO — igual ao grid principal) FALLBACK_SIGMA0 = 2.0 # max(4,4,4,4)/2 = 2.0 MEM_CRITICAL_PCT_FOR_FALLBACK = 75 # se RSS > 75% cgroup, ativa fallback # Streaming — User requirement: "streaming de 100 em 100 samples" STREAM_BATCH_SIZE = 100 # V2-dynamic — Metas atualizadas (V6.5-V2-metrics): 8000 CONHECIMENTO + 2000 PUNIÇÃO META_MINIMA_CONHECIMENTO = 8000 MAX_SAMPLES_PER_DATASET_CONHECIMENTO = 1000 # 8 × 1000 = 8000 ≥ 8000 META_MINIMA_PUNICAO = 2000 MAX_SAMPLES_PUNICAO = 2000 # User requirement: "não reduzir tempo e não gerar dados sintéticos" + # "todo streaming (FASE1 e da FASE2) deve ter pausa para dar tempo de conclusão # de processamento continuando após conclusão" # V6.5-V2-metrics-FIX-4 (latest user requirement): # "FASE2 PUNIÇÃO é mais pesada é pode exigir pausas do streaming até # concluir o processamento". # Pausas separadas para FASE1 (CONHECIMENTO) e FASE2 (PUNITIVA): # - FASE1: pausas moderadas (SOM-only, sem hypothesis layer) # - FASE2: pausas maiores (16 hipóteses × 30 steps × 3 trials + EWC + revival) INTER_BATCH_PAUSE_S_FASE1 = 0.15 # CONHECIMENTO: leve INTER_BATCH_PAUSE_S_FASE2 = 0.60 # PUNITIVA: 4x maior (heavier processing) INTER_DATASET_PAUSE_S_FASE1 = 0.5 INTER_DATASET_PAUSE_S_FASE2 = 1.5 # PUNITIVA: 3x maior INTER_STREAM_BATCH_PAUSE_S_FASE1 = 0.3 INTER_STREAM_BATCH_PAUSE_S_FASE2 = 1.2 # PUNITIVA: 4x maior POST_PROCESSING_PAUSE_S_FASE1 = 0.4 POST_PROCESSING_PAUSE_S_FASE2 = 1.5 # PUNITIVA: ~4x maior # Compatibilidade (mantém nomes antigos apontando para FASE1 — usados em # código legado que não diferencia fases) INTER_BATCH_PAUSE_S = INTER_BATCH_PAUSE_S_FASE1 INTER_DATASET_PAUSE_S = INTER_DATASET_PAUSE_S_FASE1 INTER_STREAM_BATCH_PAUSE_S = INTER_STREAM_BATCH_PAUSE_S_FASE1 POST_PROCESSING_PAUSE_S = POST_PROCESSING_PAUSE_S_FASE1 # V6.5-V2-metrics-FIX-4 — auto-revive config # User requirement: "distribuindo o processamento paralelamente". # Quando dead_rate > 50% e passou o cooldown, revive neurônios mortos. # V6.5-V2-auto-conscience-v2 — reduzido cooldown de 200 → 30 para permitir # revives mais frequentes durante FASE1 (que tem ~30-50 steps por chunk). AUTO_REVIVE_DEAD_RATE_THRESHOLD = 0.5 AUTO_REVIVE_COOLDOWN_STEPS = 30 # Storage critical STORAGE_CRITICAL_PCT = 90 STORAGE_CHECK_INTERVAL_STEPS = 5 # V2 — 8 datasets para CONHECIMENTO (na ordem exata especificada pelo usuário) CONHECIMENTO_DATASETS = [ "dominguesm/restore-punctuation-ptbr-dataset", "carolina-c4ai/corpus-carolina", "CEIA-POSITIVO/ultrachat_br_clustred_balanced_v1", "dominguesm/Canarim-Instruct-PTBR-Dataset", "adalbertojunior/punctuation-ptbr", "iara-project/news-articles-ptbr-dataset", "manoela/noticias_ptbr", "BrunoN-Dev/corpus-ptbr-v1", ] # V2 — Dataset para TREINAMENTO COM PUNIÇÃO PUNICAO_DATASET = "BrunoN-Dev/corpus-ptbr-v1" # ============================================================================ # V6.5-fix — DATASET_LABEL_STRINGS (dynamic label registry) # ============================================================================ # User requirement: "os strings 'gato' e 'cachorro' são fixos quando deveriam # ser extrações variáveis e flexíveis de rótulos proveniente de dados dos # datasets anteriormente treinados" DATASET_LABEL_STRINGS: Dict[str, Dict[int, str]] = { "dominguesm/restore-punctuation-ptbr-dataset": { 0: "unpunctuated", 1: "punctuated", }, "carolina-c4ai/corpus-carolina": { 0: "raw_corpus", 1: "normalized_text", }, "CEIA-POSITIVO/ultrachat_br_clustred_balanced_v1": { 0: "user_turn", 1: "assistant_turn", }, "dominguesm/Canarim-Instruct-PTBR-Dataset": { 0: "instruction", 1: "response", }, "adalbertojunior/punctuation-ptbr": { 0: "unpunctuated", 1: "punctuated", }, "iara-project/news-articles-ptbr-dataset": { 0: "headline", 1: "body", }, "manoela/noticias_ptbr": { 0: "headline", 1: "body", }, "BrunoN-Dev/corpus-ptbr-v1": { 0: "short_text", 1: "long_text", }, } # ============================================================================ # 3. Import KohonenLearningSystemV2 (V6.5-V2 — 16 hipóteses + 3 tentativas) # ============================================================================ from bigru_t.model.kohonen_learning_system import ( # noqa: E402 KohonenLearningSystemV2, KohonenLearningSystem, SimpleBBPETokenizer, positional_encoding, text_to_4d_vector, KohonenSOM4D, HypothesisClassifier, HypothesisEnsemble, DeltaGenerator, ) from bigru_t.model.attention_multimodal import MultiHeadAttention # noqa: E402 # V6.5-V2-auto-adjust — Auto-ajuste dos Mapas Auto-Organizáveis com base # nas 8 métricas canônicas (4 principais + 4 indicadores de falha). # User requirement: "inserir melhorias para os scripts das métricas de # aprendizado e indicadores de falha em Mapas Auto-Organizáveis" + "aprimorar # matematicamente a lógica de autoajustes analisando Kohonen" + "PORTANTO: # observar se os valores demonstram que o modelo esteja aprendendo e parar # caso não esteja". from bigru_t.model.som_auto_adjust import ( # noqa: E402 SOMAutoAdjuster, SOMAutoAdjustState, create_auto_adjuster, ) from bigru_t.model.som_metrics import compute_all_metrics, SOMMetricHistory # noqa: E402 logger.info( f"[V6.5-V2] All modules imported. " f"Grid={SOM_GRID} ({SOM_GRID[0]*SOM_GRID[1]*SOM_GRID[2]*SOM_GRID[3]} neurons) | " f"Hidden={HIDDEN_DIM} | Vocab={VOCAB_SIZE} | " f"n_hypotheses={N_HYPOTHESES} | n_trials={N_TRIALS} | " f"CONHECIMENTO datasets={len(CONHECIMENTO_DATASETS)} | " f"PUNICAO dataset={PUNICAO_DATASET} | " f"meta_conhecimento={META_MINIMA_CONHECIMENTO} | " f"meta_punicão={META_MINIMA_PUNICAO}" ) # ============================================================================ # 4. Storage / Memory Monitors # ============================================================================ def get_disk_usage_pct() -> float: try: total, used, free = shutil.disk_usage("/") return float(used / total * 100.0) except Exception: return 0.0 def get_process_rss_mb() -> float: """V2-dynamic-memory — Retorna RSS do processo atual em MB.""" try: import resource return float(resource.getrusage(resource.RUSAGE_SELF).ru_maxrss) / 1024.0 except Exception: return 0.0 def get_cgroup_memory_limit_mb() -> Optional[float]: """V2-dynamic-memory — Retorna limite de memória do cgroup em MB (ou None).""" try: with open("/sys/fs/cgroup/memory.max", "r") as f: val = int(f.read().strip()) if val > 0: return float(val) / 1e6 # bytes → MB except Exception: pass try: with open("/sys/fs/cgroup/memory/memory.limit_in_bytes", "r") as f: val = int(f.read().strip()) if val > 0: return float(val) / 1e6 except Exception: pass return None # V2-dynamic-memory — Threshold de memória (porcentagem do cgroup limit) CGROUP_MEM_CRITICAL_PCT = 80 # se RSS > 80% do cgroup, salvar estado e parar def check_memory_and_maybe_fallback(kls: KohonenLearningSystemV2) -> Dict[str, Any]: """V6.5-V2-buffer-864 — Monitora memória e aplica fallback se crítico. User requirement (V6.5-V4-canonical-256): "fazer (tornar canônico) buffer 256 e grid para (4,4,4,4)=256". Buffer e grid agora têm o MESMO tamanho (256). Fallback OOM reduz buffer para 128 (não 256, pois 256 já é o canônico). Grid (4,4,4,4)=256 PERMANECE sempre. Lógica: 1. Lê RSS atual e cgroup memory limit. 2. Se RSS > MEM_CRITICAL_PCT_FOR_FALLBACK do cgroup: a. Reduz buffer_max_size do KLS para FALLBACK_BUFFER_SIZE (128). b. Trunca buffer_4d atual para FALLBACK_BUFFER_SIZE. c. Registra evento de fallback (NÃO recriia o SOM grid — preserva o aprendizado já acumulado nos pesos 256-neuronios; o fallback apenas reduz o buffer de amostras recentes, não a arquitetura). 3. Retorna relatório com RSS, pct, fallback_ativado. Nota: Recriar o SOM grid do zero perderia todo o aprendizado da FASE1. Mantemos o grid (4,4,4,4)=256 (preservando pesos aprendidos) e reduzimos APENAS o buffer de amostras recentes para 128. Isto preserva a arquitetura Kohonen enquanto mitiga OOM. Se o OOM persistir, o save_model_states_for_evaluation será chamado pelo caller. """ rss_mb = get_process_rss_mb() cg_limit = get_cgroup_memory_limit_mb() fallback_ativado = False pct = 0.0 if cg_limit is not None and cg_limit > 0: pct = (rss_mb / cg_limit) * 100.0 if pct > MEM_CRITICAL_PCT_FOR_FALLBACK: # Aplica fallback: reduz buffer_max_size old_size = kls.buffer_max_size kls.buffer_max_size = FALLBACK_BUFFER_SIZE # Trunca buffer atual if len(kls.buffer_4d) > FALLBACK_BUFFER_SIZE: overflow = len(kls.buffer_4d) - FALLBACK_BUFFER_SIZE del kls.buffer_4d[:overflow] del kls.buffer_labels[:overflow] fallback_ativado = True logger.warning( f"[V6.5-V2-buffer-864] MEM FALLBACK ATIVADO: RSS={rss_mb:.0f}MB " f"({pct:.1f}% of {cg_limit:.0f}MB) > threshold " f"{MEM_CRITICAL_PCT_FOR_FALLBACK}%. " f"buffer_max_size {old_size} → {kls.buffer_max_size}. " f"SOM grid PRESERVADO em {kls.som_grid} (aprendizado mantido)." ) return { "rss_mb": float(rss_mb), "cgroup_limit_mb": float(cg_limit) if cg_limit else None, "mem_pct": float(pct), "fallback_ativado": bool(fallback_ativado), "buffer_max_size": int(kls.buffer_max_size), "som_grid_preservado": list(kls.som_grid), } def aggressive_memory_cleanup() -> Dict[str, Any]: """V6.5-V2-metrics — Limpeza agressiva de RAM. User requirement: "analisar e aprimorar (logicamente e matematicamente) no projeto limpeza de memória RAM e de Armazenamento (ao concluir)". Implementação: 1. gc.collect() — 2 passes (gerações 0+1 e 2) 2. torch.cuda.empty_cache() (se CUDA disponível) 3. Limpa caches internos do PyTorch (MKL, oneDNN) 4. Reporta RSS antes/depois para confirmação """ rss_before = get_process_rss_mb() gc.collect() gc.collect() try: import torch if torch.cuda.is_available(): torch.cuda.empty_cache() # V6.5-V2-metrics: limpa caches CPU do PyTorch if hasattr(torch, "cpu"): try: # MKL/OneDNN primitive cache é limpo implicitamente por gc pass except Exception: pass except Exception: pass rss_after = get_process_rss_mb() freed_mb = max(0.0, rss_before - rss_after) return { "gc_called": True, "rss_before_mb": float(rss_before), "rss_after_mb": float(rss_after), "freed_mb": float(freed_mb), "timestamp": datetime.now().isoformat(), } def aggressive_storage_cleanup() -> Dict[str, Any]: """V6.5-V2-metrics — Limpeza agressiva de Armazenamento. User requirement: "analisar e aprimorar (logicamente e matematicamente) no projeto limpeza de memória RAM e de Armazenamento (ao concluir)". Implementação: 1. Limpa locks HF (*.lock) 2. Limpa downloads incompletos (*.incomplete) 3. Limpa __pycache__ recursivamente 4. Limpa arquivos .tmp 5. Limpa logs antigos (>7 dias, se existirem) 6. Reporta MB liberados """ n_removed = 0 mb_freed = 0.0 # 1. HF cache locks e incompletos hf_cache = Path.home() / ".cache" / "huggingface" if hf_cache.exists(): for pattern in ["*.lock", "*.incomplete", "*.tmp"]: for f in hf_cache.rglob(pattern): try: sz = f.stat().st_size f.unlink() n_removed += 1 mb_freed += sz / 1e6 except Exception: pass # 2. __pycache__ recursivo no BiGRU_T_version bigru_root = BIGRU_ROOT for cache_dir in bigru_root.rglob("__pycache__"): if "deprecados" in str(cache_dir): continue try: for f in cache_dir.iterdir(): try: sz = f.stat().st_size f.unlink() n_removed += 1 mb_freed += sz / 1e6 except Exception: pass cache_dir.rmdir() except Exception: pass # 3. Arquivos .tmp no BiGRU_T_version for f in bigru_root.rglob("*.tmp"): if "deprecados" in str(f): continue try: sz = f.stat().st_size f.unlink() n_removed += 1 mb_freed += sz / 1e6 except Exception: pass # 4. Logs antigos (>7 dias) em logs/ logs_dir = PROJECT_ROOT / "logs" if logs_dir.exists(): cutoff = time.time() - 7 * 24 * 3600 # 7 dias for f in logs_dir.glob("*.log"): try: if f.stat().st_mtime < cutoff: sz = f.stat().st_size f.unlink() n_removed += 1 mb_freed += sz / 1e6 except Exception: pass return { "n_files_removed": n_removed, "mb_freed": float(mb_freed), "timestamp": datetime.now().isoformat(), } # ============================================================================ # 5. Save model states (storage critical / end of training) # ============================================================================ def save_model_states_for_evaluation( kls: KohonenLearningSystemV2, reason: str, step: int, total_samples: int, output_path: Optional[Path] = None, ) -> Dict[str, Any]: """V6.5-V2 — Salva estados do modelo (SOM pesos, classifier, ensemble, métricas) em disco para avaliação posterior. User requirement: "parar se o Armazenamento for crítico e armazenar os estados do modelo para avaliação" User requirement (latest, V6.5-V2-metrics-FIX-2): "ATENÇÃO: o estado do modelo deve ser contínuo com tamanho máximo de 1GB cada para que todos os estados possam ser carregados em inferência". Para garantir continuidade e tamanho ≤ 1GB: 1. Todos os tensores são convertidos para CPU + float32 (consistência). 2. Buffer tail é limitado a 64 amostras (era 128 — reduzido para economizar espaço sem perder representatividade). 3. State metrics e v2_metrics são armazenados como JSON-serializable (sem tensores). 4. Após salvar, verifica o tamanho do arquivo. Se > 1GB, emite warning e trunca componentes menos essenciais (preserva SOM + classifier + ensemble + delta_scale, que são necessários para inferência). Salva: - som.weights (I,J,K,L,4) — tensor principal do SOM - som.old_weights_w / fisher_w (se existirem) - classifier.state_dict() (se treinado) - hypothesis_ensemble.state_dict() (V2 — 16 geradores) - delta_scale (V2 — parâmetro treinável) - kls.get_state_metrics() + kls.get_v2_metrics() - buffer_4d tail (até 64 amostras — V6.5-V2-metrics-FIX-2) - metadados (reason, step, total_samples, timestamp) """ import torch if output_path is None: output_path = MODEL_STATES_PATH # V6.5-V2-metrics-FIX-2 — Limite máximo de tamanho do estado (1GB). # User requirement: "tamanho máximo de 1GB cada para que todos os estados # possam ser carregados em inferência". MAX_STATE_SIZE_BYTES = 1 * 1024 * 1024 * 1024 # 1 GB try: state = { "_meta": { "reason": str(reason), "step": int(step), "total_samples": int(total_samples), "timestamp": datetime.now().isoformat(), "version": "V6.5-V2-unified-state", # V6.5-V2-unified — fase que gerou este estado. Permite que # um processo separado (ex: retomar treino) saiba se o estado # é pós-FASE1 (pronto para FASE2) ou pós-FASE2 (treino completo). "phase": str(reason), # "end_of_fase_conhecimento" | "end_of_training_v2" "kmeans_pp_applied": bool(getattr(kls.som, "_kmeans_pp_initialized", False)), "som_grid": list(kls.som_grid), "n_neurons": int(kls.som_neuron_count), "hidden_dim": int(kls.hidden_dim), "vocab_size": int(kls.tokenizer.vocab_size), "n_hypotheses": int(kls.n_hypotheses), "n_trials": int(kls.n_trials), "classifier_trained": bool(kls.classifier_trained), "punishment_count": int(kls.punishment_count), "success_count": int(kls.success_count), "time_counter": int(kls.time_counter), "inference_punishment_count": int(kls.inference_punishment_count), # V6.5-V2-metrics-FIX-2 — contadores cumulativos "total_hyp_steps_executed": int(getattr(kls, "_total_hyp_steps_executed", 0)), "n_train_hyp_calls": int(getattr(kls, "_n_train_hyp_calls", 0)), }, "som_weights": kls.som.weights.detach().clone().cpu().float(), "som_old_weights_w": ( kls.som.old_weights_w.detach().clone().cpu().float() if kls.som.old_weights_w is not None else None ), "som_fisher_w": ( kls.som.fisher_w.detach().clone().cpu().float() if kls.som.fisher_w is not None else None ), "embedding_state": { "weight": kls.embedding.weight.detach().clone().cpu().float(), }, "classifier_state": ( kls.classifier.state_dict() if kls.classifier is not None else None ), # V2 — HypothesisEnsemble state (16 geradores) "hypothesis_ensemble_state": ( kls.hypothesis_ensemble.state_dict() ), "delta_scale": kls.delta_scale.detach().clone().cpu().float(), "state_metrics": kls.get_state_metrics(), "v2_metrics": kls.get_v2_metrics(), # V6.5-V2-metrics-FIX-2 — buffer tail reduzido para 64 (era 128) # para garantir estado ≤ 1GB e permitir carregamento em inferência. "buffer_4d_tail": ( [v.detach().clone().cpu().float() for v in kls.buffer_4d[-64:]] if kls.buffer_4d else [] ), "buffer_labels_tail": list(kls.buffer_labels[-64:]) if kls.buffer_labels else [], "label_registry": kls.label_registry, } torch.save(state, output_path) size_bytes = output_path.stat().st_size size_mb = size_bytes / 1e6 size_gb = size_bytes / (1024 ** 3) # V6.5-V2-metrics-FIX-2 — Verifica tamanho do estado size_ok = size_bytes <= MAX_STATE_SIZE_BYTES size_status = "OK" if size_ok else "EXCEEDS_1GB" if not size_ok: logger.warning( f"[V6.5-V2-metrics-FIX-2] State size {size_gb:.3f}GB exceeds 1GB limit! " f"Future inference loads may fail. Consider reducing HIDDEN_DIM, " f"VOCAB_SIZE, or MAX_N_HYPOTHESES to bring state under 1GB." ) logger.info( f"[V6.5-V2-metrics-FIX-2] Model states saved to {output_path} " f"({size_mb:.1f}MB / {size_gb:.3f}GB, status={size_status}, " f"reason={reason}, step={step}, total_samples={total_samples})" ) return { "saved": True, "path": str(output_path), "size_mb": float(size_mb), "size_gb": float(size_gb), "size_status": size_status, "size_within_1gb_limit": bool(size_ok), "reason": str(reason), "step": int(step), "total_samples": int(total_samples), "n_tensors": 7, "n_buffer_tail": len(state["buffer_4d_tail"]), } except Exception as e: logger.error(f"[V6.5-V2-metrics-FIX-2] Failed to save model states: {e}") traceback.print_exc() return { "saved": False, "error": str(e)[:300], "reason": str(reason), "step": int(step), "total_samples": int(total_samples), } # ============================================================================ # 5.5 Load FASE1 state into KLS (FASE2 must use FASE1's completed state) # ============================================================================ def load_fase1_state_into_kls( kls: KohonenLearningSystemV2, state_path: Path, ) -> Dict[str, Any]: """V6.5-V2-metrics — Carrega estado do modelo concluído da FASE1 no KLS. User requirement: "esta FASE2 (nunca acontece antes da FASE1) deve usar do arquivo de Estado do MODELO concluído da FASE1 e a FASE2 tem função: de treinar, ajustar dados e os parâmetros do MAPA-SOM de redes Kohonen, informar métricas de aprendizado do modelo". Esta função é chamada explicitamente antes de iniciar a FASE 2, garantindo que a FASE 2 herde todo o conhecimento acumulado da FASE 1 (pesos do SOM, EWC references, classifier, hypothesis ensemble, delta_scale, label_registry). Args: kls: instância KohonenLearningSystemV2 (mesma usada na FASE1). state_path: caminho para o arquivo .pt salvo ao final da FASE1. Returns: Dict com status do carregamento. """ import torch if not state_path.exists(): logger.warning( f"[V6.5-V2-metrics] FASE1 state file not found: {state_path}. " f"FASE2 will run with the in-memory KLS state (already inherited " f"from FASE1 since same KLS instance is used)." ) return { "loaded": False, "reason": "state_file_not_found", "path": str(state_path), "in_memory_state_used": True, } try: state = torch.load(state_path, map_location="cpu", weights_only=False) meta = state.get("_meta", {}) logger.info( f"[V6.5-V2-metrics] Loading FASE1 state from {state_path}: " f"reason={meta.get('reason', '?')}, " f"step={meta.get('step', 0)}, " f"total_samples={meta.get('total_samples', 0)}, " f"timestamp={meta.get('timestamp', '?')}" ) # Carrega pesos do SOM (preserva o aprendizado da FASE1) if "som_weights" in state and state["som_weights"] is not None: kls.som.weights = state["som_weights"].clone().float() logger.info( f"[V6.5-V2-metrics] SOM weights loaded: " f"shape={tuple(kls.som.weights.shape)}, " f"norm={float(kls.som.weights.norm().item()):.4f}" ) # Carrega EWC references (se existirem da FASE1) if "som_old_weights_w" in state and state["som_old_weights_w"] is not None: kls.som.old_weights_w = state["som_old_weights_w"].clone().float() logger.info("[V6.5-V2-metrics] SOM EWC old_weights_w loaded.") if "som_fisher_w" in state and state["som_fisher_w"] is not None: kls.som.fisher_w = state["som_fisher_w"].clone().float() logger.info("[V6.5-V2-metrics] SOM EWC fisher_w loaded.") # Carrega embedding if "embedding_state" in state and state["embedding_state"] is not None: kls.embedding.weight.data = state["embedding_state"]["weight"].clone().float() logger.info("[V6.5-V2-metrics] Embedding weights loaded.") # Carrega classifier (se treinado na FASE1) if "classifier_state" in state and state["classifier_state"] is not None: if kls.classifier is None: # Recria o classifier com a mesma arquitetura kls.classifier = HypothesisClassifier( input_dim=kls.som_neuron_count, hidden_dims=kls.hypothesis_hidden, ) kls.classifier.load_state_dict(state["classifier_state"]) kls.classifier_trained = True logger.info("[V6.5-V2-metrics] Classifier state loaded.") # Carrega HypothesisEnsemble if "hypothesis_ensemble_state" in state and state["hypothesis_ensemble_state"] is not None: kls.hypothesis_ensemble.load_state_dict(state["hypothesis_ensemble_state"]) logger.info( f"[V6.5-V2-metrics] HypothesisEnsemble state loaded: " f"n_generators={len(kls.hypothesis_ensemble.generators)}, " f"active={kls.hypothesis_ensemble.active_count}" ) # Carrega delta_scale if "delta_scale" in state and state["delta_scale"] is not None: kls.delta_scale.data = state["delta_scale"].clone().float() logger.info( f"[V6.5-V2-metrics] delta_scale loaded: {kls.delta_scale.item():.4f}" ) # Carrega label_registry (preserva mapeamentos dinâmicos) if "label_registry" in state and state["label_registry"] is not None: kls.label_registry = state["label_registry"] logger.info( f"[V6.5-V2-metrics] label_registry loaded: " f"{len(kls.label_registry)} datasets" ) # Restaura time_counter (preserva coordenada temporal absoluta) if "time_counter" in meta: kls.time_counter = int(meta["time_counter"]) logger.info( f"[V6.5-V2-metrics] time_counter restored: {kls.time_counter}" ) logger.info( f"[V6.5-V2-metrics] ✓ FASE1 state loaded successfully. " f"KLS is now ready for FASE2 (PUNITIVE) training." ) return { "loaded": True, "path": str(state_path), "meta": meta, "in_memory_state_used": False, } except Exception as e: logger.error( f"[V6.5-V2-metrics] Failed to load FASE1 state: {e}. " f"FASE2 will run with the in-memory KLS state (already inherited " f"from FASE1 since same KLS instance is used)." ) traceback.print_exc() return { "loaded": False, "reason": f"load_error: {str(e)[:200]}", "path": str(state_path), "in_memory_state_used": True, } # ============================================================================ # 6. Streaming helper — chunks of 100 samples (V6.5-V2-metrics: O(N) not O(N²)) # ============================================================================ def stream_dataset_in_chunks( dataset_name: str, chunk_size: int, max_total: int, hf_token: Optional[str] = None, timeout_per_chunk: int = 60, ) -> Iterator[List[str]]: """V6.5-V2-metrics — Stream a dataset in chunks of `chunk_size` samples. User requirement: "streaming de 100 em 100 samples até esgonar". V6.5-V2-metrics FIX: A versão anterior re-chamava stream_dataset() para cada chunk com seed_offset = yielded_total, fazendo com que o HF re-streaming recapitulasse todas as amostras anteriores a cada chunk (custo O(N²) em banda e cache). Esta versão cria UM ÚNICO iterador por dataset e vai colhendo chunks de `chunk_size` em `chunk_size`, sem recomeçar do zero. Isto reduz drasticamente o uso de memória (que antes crescia linearmente com o nº de chunks devido ao cache acumulado de metadados parquet) e permite atingir as metas de 8000 samples CONHECIMENTO + 2000 samples PUNIÇÃO dentro do limite de 4GB RAM do cgroup. """ try: from bigru_t.data.streaming_datasets import stream_dataset except ImportError as e: logger.error(f"[V6.5-V2] Cannot import stream_dataset: {e}") return yielded_total = 0 chunk_idx = 0 chunk: List[str] = [] t_start = time.time() try: sample_iter = stream_dataset( dataset_name, max_samples=max_total, hf_token=hf_token, ) for sample in sample_iter: if yielded_total >= max_total: break # Timeout check: se demorar demais, yield o que temos e continua if time.time() - t_start > timeout_per_chunk * (chunk_idx + 1): logger.warning( f"[V6.5-V2] {dataset_name} chunk {chunk_idx+1} " f"timeout after {timeout_per_chunk}s " f"(collected {len(chunk)}/{chunk_size})" ) if chunk: yield chunk yielded_total += len(chunk) chunk_idx += 1 chunk = [] if time.time() - t_start > timeout_per_chunk * (chunk_idx + 2): # Segundo timeout consecutivo — desiste deste dataset logger.warning( f"[V6.5-V2] {dataset_name} exhausted after {yielded_total} samples " f"({chunk_idx} chunks) due to repeated timeouts" ) return continue if sample.raw_text and len(sample.raw_text.strip()) > 0: # V6.6 — User requirement: "aumentar o truncamento de textos # (line 993 do train script) de 200 para 1000+ chars, permitindo # mais diversidade de pares byte-level no BBPE". chunk.append(sample.raw_text.strip()[:TEXT_TRUNCATION_CHARS]) if len(chunk) >= chunk_size: yield chunk yielded_total += len(chunk) chunk_idx += 1 logger.info( f"[V6.5-V2] {dataset_name} chunk {chunk_idx} done: " f"{len(chunk)} samples (total: {yielded_total}/{max_total})" ) chunk = [] except Exception as e: logger.warning(f"[V6.5-V2] {dataset_name} streaming error: {e}") if chunk: yield chunk yielded_total += len(chunk) chunk_idx += 1 return # Yield any remaining partial chunk if chunk: yield chunk yielded_total += len(chunk) chunk_idx += 1 logger.info( f"[V6.5-V2] {dataset_name} exhausted after {yielded_total} samples " f"({chunk_idx} chunks)" ) # ============================================================================ # 7. make_label — generates binary label deterministically from text # ============================================================================ def make_label(text: str, dataset_name: Optional[str] = None) -> int: """Gera label binário determinístico baseado no texto. V6.5-V2: para o dataset BrunoN-Dev/corpus-ptbr-v1, usa paridade do número de palavras como label (não é dado sintético — é derivação algorítmica do conteúdo real do dataset). """ text_lower = text.lower() if dataset_name == "BrunoN-Dev/corpus-ptbr-v1": # V2 — label_strategy "hash_parity": paridade do nº de palavras n_words = len(text.split()) return n_words % 2 # Default: label baseado em palavras-chave semânticas if any(w in text_lower for w in [ "gato", "mia", "dorme", "brinca", "menina", "boneca", "olá", "ola", "hello", "help", "instru", "pergunta" ]): return 0 return 1 # ============================================================================ # 8. Verify V2 logic — predict fix + attention + 16 hypotheses # ============================================================================ def verify_v2_logic() -> Dict[str, Any]: """V6.5-V2 — Verifica a lógica V2 (predict dinâmico + attention + 16 hipóteses).""" import torch checks = {} # Check 1: KohonenLearningSystemV2 instantiation try: kls = KohonenLearningSystemV2( vocab_size=512, hidden_dim=32, seq_len=8, som_grid=(3, 3, 3, 2), enable_vqvae2=False, enable_reasoning=False, enable_w8a8=False, enable_attention=False, n_hypotheses=4, max_n_hypotheses=8, n_trials=2, hyp_train_steps=3, ) assert hasattr(kls, "hypothesis_ensemble"), "no hypothesis_ensemble" assert len(kls.hypothesis_ensemble.generators) == 8, \ f"expected 8 generators (max_n_hypotheses), got {len(kls.hypothesis_ensemble.generators)}" assert kls.hypothesis_ensemble.active_count == 4, \ f"expected active_count=4, got {kls.hypothesis_ensemble.active_count}" assert hasattr(kls, "delta_scale"), "no delta_scale" assert hasattr(kls, "_adapt_hyperparameters"), "no _adapt_hyperparameters" assert hasattr(kls, "get_adaptation_log"), "no get_adaptation_log" checks["v2_instantiation"] = { "status": "PASS", "details": f"n_hyp={kls.n_hypotheses}/{kls.max_n_hypotheses} " f"(active={kls.hypothesis_ensemble.active_count}), " f"n_trials={kls.n_trials}, delta_scale={kls.delta_scale.item():.4f}", } except Exception as e: checks["v2_instantiation"] = {"status": "FAIL", "error": str(e)} kls = None # Check 2: predict returns dynamic labels (not hardcoded gato/cachorro) try: if kls is None: raise RuntimeError("KLS not instantiated") kls.register_dataset_labels("test_ds", {0: "negative_class", 1: "positive_class"}) kls.tokenizer.fit(["o gato dorme", "o cachorro corre", "teste de predição"]) for _ in range(3): kls.add_data( ["o gato dorme", "o cachorro corre", "a menina brinca", "o pássaro voa"], [0, 1, 0, 1], label_strings={0: "negative_class", 1: "positive_class"}, dataset_name="test_ds", ) kls.check_training_start() kls.train_som_on_buffer() kls._label_neurons() pred = kls.predict("o gato dorme", dataset_name="test_ds") assert pred in ["negative_class", "positive_class"], \ f"predict returned: {pred} (expected dynamic label)" assert pred != "gato" and pred != "cachorro", \ f"predict returned hardcoded: {pred}" checks["predict_dynamic_labels"] = { "status": "PASS", "details": f"predict returned: {pred} (dynamic, not hardcoded)", } except Exception as e: checks["predict_dynamic_labels"] = {"status": "FAIL", "error": str(e)} # Check 3: process_batch_v2 returns phase info try: if kls is None: raise RuntimeError("KLS not instantiated") result = kls.process_batch_v2( ["o gato dorme", "o cachorro corre"], [0, 1], dataset_name="test_ds", enable_punishment=False, ) assert result["phase"] == "CONHECIMENTO" assert "elapsed_ms" in result checks["process_batch_v2_conhecimento"] = { "status": "PASS", "details": f"phase={result['phase']}, action={result['action']}", } except Exception as e: checks["process_batch_v2_conhecimento"] = {"status": "FAIL", "error": str(e)} # Check 4: V2 methods exist try: for method in ["train_hypotheses", "select_best_delta", "apply_best_delta_and_consolidate", "process_batch_v2", "get_v2_metrics"]: assert hasattr(kls, method), f"missing method: {method}" checks["v2_methods_exist"] = { "status": "PASS", "details": "train_hypotheses, select_best_delta, apply_best_delta_and_consolidate, " "process_batch_v2, get_v2_metrics", } except Exception as e: checks["v2_methods_exist"] = {"status": "FAIL", "error": str(e)} # Summary n_pass = sum(1 for c in checks.values() if c.get("status") == "PASS") n_fail = sum(1 for c in checks.values() if c.get("status") == "FAIL") return { "checks": checks, "n_pass": n_pass, "n_fail": n_fail, "all_pass": n_fail == 0, } # ============================================================================ # 9. PHASE 1 — CONHECIMENTO (streaming 8 datasets, sem punição) # ============================================================================ def run_fase_conhecimento( kls: KohonenLearningSystemV2, hf_token: Optional[str], ) -> Dict[str, Any]: """Fase 1 — CONHECIMENTO: streaming 8 datasets, sem punição. User requirement: "streaming CONHECIMENTO de 100 em 100 samples nesta sequência [8 datasets]" + "meta mínima 8000 samples ou mais para CONHECIMENTO" Para cada dataset: - Stream em chunks de 100 samples. - process_batch_v2(enable_punishment=False) — apenas acumula conhecimento. - Pausas inter-batch e inter-dataset. Returns: Dict com métricas da fase CONHECIMENTO. """ logger.info("\n" + "=" * 80) logger.info("[V6.5-V2] FASE 1 — CONHECIMENTO (sem punição)") logger.info("=" * 80) logger.info(f" Datasets: {len(CONHECIMENTO_DATASETS)}") logger.info(f" Max samples/dataset: {MAX_SAMPLES_PER_DATASET_CONHECIMENTO}") logger.info(f" Meta mínima: {META_MINIMA_CONHECIMENTO}") logger.info(f" Stream batch size: {STREAM_BATCH_SIZE}") logger.info("=" * 80 + "\n") step = 0 t_start = time.time() samples_per_dataset: Dict[str, int] = {ds: 0 for ds in CONHECIMENTO_DATASETS} streaming_failures: Dict[str, int] = {ds: 0 for ds in CONHECIMENTO_DATASETS} storage_critical_stopped = False chunk_idx_global = 0 accuracies_log: List[Dict[str, Any]] = [] attention_log: List[Dict[str, Any]] = [] # V6.5-V2-metrics — log de métricas SOM durante CONHECIMENTO som_metrics_log: List[Dict[str, Any]] = [] # V6.5-V2-auto-adjust — Auto-ajustador SOM para FASE1 # User requirement: "aprimorar matematicamente a lógica de autoajustes # analisando Kohonen" + "observar se os valores demonstram que o modelo # esteja aprendendo e parar caso não esteja". auto_adjuster = create_auto_adjuster() auto_adjust_log: List[Dict[str, Any]] = [] stop_training_reason: Optional[str] = None # V7-tokenizer-growth — Corpus collector para crescimento do vocabulário # User requirement: "verificar impacto de seq_len em vocab real=293 # (alvo 16384)". A análise V7 revelou que o vocab real era apenas 293 # porque o tokenizer era treinado com apenas 17 frases hardcoded. # # Prova 13 (vocab growth): O BBPE precisa de corpus ≥ 200 textos para # ter merges suficientes para alcançar vocab_size=16384. Com 17 frases, # o nº de pares únicos é ~50, limitando o vocab a ~300. Com 2000+ # frases (coletadas durante FASE1), o nº de pares cresce para 5000+, # permitindo vocab alcançar 5000-10000 (próximo do alvo 16384). # # Estratégia: acumular textos dos chunks durante FASE1, e refitar o # tokenizer a cada 1000 amostras (10 chunks). O embedding matrix é # fixo (16384 rows), então qualquer crescimento de vocab é acomodado. # O SOM adapta-se automaticamente às novas projeções 4D do embedding. tokenizer_corpus_buffer: List[str] = [] # V6.7: REATIVADO refit com serial mode + memory guard (fix BBPE OOM) # User requirement: "processar aprimoramento (matemático e lógico) para # resolver: tokenizer-growth refit (was causing crashes during BBPE # parallel training at 1000-sample mark)". # Prova 14: agora serial mode (no fork) + memory guard (skip se RSS>85%) # tornam o refit seguro. Era 10**9 (desativado); agora 1000 (reativado). TOKENIZER_REFIT_INTERVAL_SAMPLES = 1000 TOKENIZER_CORPUS_CAP = 2000 # cap para evitar OOM (memória: ~2MB) last_tokenizer_refit_sample_count = 0 tokenizer_growth_log: List[Dict[str, Any]] = [] initial_vocab_size = int(getattr(kls.tokenizer, "vocab_size", 0)) logger.info( f"[V7-tokenizer-growth] Initial vocab_size={initial_vocab_size} " f"(target={VOCAB_SIZE}). Will refit every " f"{TOKENIZER_REFIT_INTERVAL_SAMPLES} samples during FASE1." ) # V6.5-V2-metrics — reseta histórico de métricas SOM no início de CONHECIMENTO kls.reset_som_metric_history() logger.info("[V6.5-V2-metrics] Histórico de métricas SOM resetado para FASE 1.") for ds_idx, dataset_name in enumerate(CONHECIMENTO_DATASETS): if storage_critical_stopped: break total_so_far = sum(samples_per_dataset.values()) if total_so_far >= META_MINIMA_CONHECIMENTO: logger.info( f"[V6.5-V2] Meta CONHECIMENTO atingida ({total_so_far} ≥ " f"{META_MINIMA_CONHECIMENTO}). Parando Fase 1." ) break logger.info( f"\n[V6.5-V2] {'='*40} DATASET {ds_idx+1}/{len(CONHECIMENTO_DATASETS)}: " f"{dataset_name} {'='*40}" ) logger.info( f"[V6.5-V2] Streaming {dataset_name} in chunks of {STREAM_BATCH_SIZE}..." ) try: chunk_iter = stream_dataset_in_chunks( dataset_name=dataset_name, chunk_size=STREAM_BATCH_SIZE, max_total=MAX_SAMPLES_PER_DATASET_CONHECIMENTO, hf_token=hf_token, timeout_per_chunk=90, ) for chunk_samples in chunk_iter: if storage_critical_stopped: break total_so_far = sum(samples_per_dataset.values()) if total_so_far >= META_MINIMA_CONHECIMENTO: logger.info( f"[V6.5-V2] Meta CONHECIMENTO atingida durante {dataset_name}." ) break chunk_idx_global += 1 if not chunk_samples: streaming_failures[dataset_name] += 1 continue samples_per_dataset[dataset_name] += len(chunk_samples) labels = [make_label(s, dataset_name) for s in chunk_samples] logger.info( f"[V6.5-V2] Processing chunk {chunk_idx_global} " f"({len(chunk_samples)} samples) from {dataset_name}" ) # Processa em batches COM PAUSAS for batch_start in range(0, len(chunk_samples), BATCH_SIZE): try: if step > 0 and step % STORAGE_CHECK_INTERVAL_STEPS == 0: disk_pct = get_disk_usage_pct() rss_mb = get_process_rss_mb() cg_limit = get_cgroup_memory_limit_mb() if cg_limit is not None: mem_pct = (rss_mb / cg_limit) * 100.0 logger.info( f"[V6.5-V2] MEM: RSS={rss_mb:.0f}MB / " f"cg_limit={cg_limit:.0f}MB ({mem_pct:.1f}%)" ) # V6.5-V2-buffer-864 — Monitora memória e aplica # fallback adaptativo do buffer se RSS > 75%. # User requirement: "aumentar buffer para 864 # amostras e observar consumo de memória (...) e # caso aumente demais o consumo de memória fazer # buffer 256 e grid para (4,4,4,4)=256". try: mem_report = check_memory_and_maybe_fallback(kls) if mem_report["fallback_ativado"]: logger.warning( f"[V6.5-V2-buffer-864] FASE1 buffer " f"reduzido para {mem_report['buffer_max_size']} " f"(RSS={mem_report['rss_mb']:.0f}MB, " f"{mem_report['mem_pct']:.1f}%)" ) except Exception as _fb_err: logger.warning( f"[V6.5-V2-buffer-864] check_memory_and_maybe_fallback " f"failed: {_fb_err}" ) if mem_pct > CGROUP_MEM_CRITICAL_PCT: logger.error( f"[V6.5-V2] MEM CRITICAL: RSS={rss_mb:.0f}MB " f"({mem_pct:.1f}% of {cg_limit:.0f}MB). " f"STOPPING CONHECIMENTO (graceful)." ) total_now = sum(samples_per_dataset.values()) save_model_states_for_evaluation( kls=kls, reason=f"mem_critical_conhecimento_{mem_pct:.1f}pct", step=step, total_samples=total_now, ) storage_critical_stopped = True break if disk_pct > STORAGE_CRITICAL_PCT: logger.error( f"[V6.5-V2] STORAGE CRITICAL: disk={disk_pct:.1f}% " f"> threshold={STORAGE_CRITICAL_PCT}%. " f"STOPPING CONHECIMENTO." ) total_now = sum(samples_per_dataset.values()) save_model_states_for_evaluation( kls=kls, reason=f"storage_critical_conhecimento_disk_{disk_pct:.1f}pct", step=step, total_samples=total_now, ) storage_critical_stopped = True break batch_sents = chunk_samples[batch_start: batch_start + BATCH_SIZE] batch_labels = labels[batch_start: batch_start + BATCH_SIZE] try: # FASE 1: enable_punishment=False (apenas CONHECIMENTO) result = kls.process_batch_v2( batch_sents, batch_labels, label_strings=DATASET_LABEL_STRINGS.get(dataset_name), dataset_name=dataset_name, enable_punishment=False, ) except (MemoryError, RuntimeError) as oom_err: # V6.5-V2-oom-guard — OOM-killer protection # User requirement: "o processo vem sendo morto OOM-kiler # (Out of memory) devido algum bug de lógica ou falta de # otimização que deve ser investigado (acrescentar # exceptions e melhor detecção de falhas de lógica e # erros de script)". # # Captura MemoryError e RuntimeError("out of memory") # explicitamente. Em vez de morrer, dispara limpeza # agressiva de memória e continua com o próximo batch. is_oom = ( isinstance(oom_err, MemoryError) or "out of memory" in str(oom_err).lower() or "cuda" in str(oom_err).lower() ) if is_oom: logger.error( f"[V6.5-V2-oom-guard] OOM detected in process_batch_v2: " f"{type(oom_err).__name__}: {str(oom_err)[:200]}" ) logger.error("[V6.5-V2-oom-guard] Triggering aggressive memory cleanup...") cleanup_result = aggressive_memory_cleanup() logger.error( f"[V6.5-V2-oom-guard] Cleanup freed " f"{cleanup_result.get('freed_mb', 0):.1f}MB " f"(RSS: {cleanup_result.get('rss_before_mb', 0):.0f}MB → " f"{cleanup_result.get('rss_after_mb', 0):.0f}MB)" ) # Trunca buffer para reduzir pressão if len(kls.buffer_4d) > 32: overflow = len(kls.buffer_4d) - 32 del kls.buffer_4d[:overflow] del kls.buffer_labels[:overflow] logger.error( f"[V6.5-V2-oom-guard] Buffer truncated to 32 " f"(dropped {overflow} samples)" ) continue else: logger.error(f"[V6.5-V2] process_batch_v2 RuntimeError: {oom_err}") traceback.print_exc() continue except Exception as e: logger.error(f"[V6.5-V2] process_batch_v2 error: {e}") traceback.print_exc() continue # V6.5-V2-oom-guard — gc.collect() entre batches para # liberar tensores intermediários antes que se acumulem. # Em runs longos (8000+ samples), mesmo pequenos leaks # acumulam e causam OOM. O custo de gc.collect() é ~10ms, # insignificante vs. o risco de morte por OOM-killer. if step % 4 == 0: gc.collect() # V6.5-V2-kmeans-pp - Apos 1o chunk, inicializa pesos # do SOM via k-means++ (deferred init). # User requirement: "inicializacao com k-means++ sobre # embeddings em vez de grid coords+ruido para ambas as # FASE1 e FASE2". try: _km = kls.init_weights_kmeans_pp_if_ready() if isinstance(_km, dict) and _km.get("initialized"): logger.info( f"[V6.5-V2-kmeans-pp] SOM weights initialized via k-means++ " f"(samples={_km.get('n_samples_used')}, " f"inertia={_km.get('inertia'):.4f})" ) except Exception as _km_err: logger.warning( f"[V6.5-V2-kmeans-pp] init failed (non-fatal): {_km_err}" ) # Métricas acc = kls.evaluate_classification() accuracies_log.append({ "step": step, "dataset": dataset_name, "chunk_idx": chunk_idx_global, "batch_size": len(batch_sents), "accuracy": float(acc), "phase": "CONHECIMENTO", "buffer_size": int(len(kls.buffer_4d)), "training_ready": bool(kls.training_ready), "elapsed_ms": float(result.get("elapsed_ms", 0.0)), }) # Attention log attn_metrics = kls.get_attention_metrics() if attn_metrics.get("n_calls", 0) > len(attention_log) * 16: attention_log.append({ "step": step, "dataset": dataset_name, "n_calls": int(attn_metrics.get("n_calls", 0)), "n_errors": int(attn_metrics.get("n_errors", 0)), "last_attn_activated": bool(attn_metrics.get("last_attn_activated", False)), "logic_functional": bool(attn_metrics.get("logic_functional", False)), }) step += 1 time.sleep(INTER_BATCH_PAUSE_S_FASE1) except (MemoryError, RuntimeError) as oom_err_outer: # V6.5-V2-oom-guard — Outer OOM catch (metrics/attention eval) is_oom_outer = ( isinstance(oom_err_outer, MemoryError) or "out of memory" in str(oom_err_outer).lower() ) if is_oom_outer: logger.error( f"[V6.5-V2-oom-guard] Outer OOM in batch loop: " f"{type(oom_err_outer).__name__}: {str(oom_err_outer)[:200]}" ) cleanup_result = aggressive_memory_cleanup() logger.error( f"[V6.5-V2-oom-guard] Cleanup freed " f"{cleanup_result.get('freed_mb', 0):.1f}MB" ) continue else: logger.error(f"[V6.5-V2] Batch RuntimeError: {oom_err_outer}") traceback.print_exc() continue except Exception as e: logger.error(f"[V6.5-V2] Batch error: {e}") traceback.print_exc() continue time.sleep(INTER_STREAM_BATCH_PAUSE_S_FASE1) # V2-dynamic-memory — Buffer sliding window: trunca para os # últimos MAX_BUFFER_SIZE amostras após cada chunk. # Os pesos do SOM já capturam o conhecimento acumulado, # então o buffer só precisa conter amostras recentes. if len(kls.buffer_4d) > MAX_BUFFER_SIZE: overflow = len(kls.buffer_4d) - MAX_BUFFER_SIZE del kls.buffer_4d[:overflow] del kls.buffer_labels[:overflow] logger.info( f"[V6.5-V2] Buffer trimmed: -{overflow} → " f"{len(kls.buffer_4d)} samples (cap={MAX_BUFFER_SIZE})" ) # V7-tokenizer-growth — Acumula corpus para refit do tokenizer # User requirement: "verificar impacto de seq_len em vocab # real=293 (alvo 16384)". Coleta amostras reais do streaming # para treinar o BBPE com corpus maior. try: # Adiciona amostras do chunk ao buffer (cap para evitar OOM) tokenizer_corpus_buffer.extend(chunk_samples[:200]) if len(tokenizer_corpus_buffer) > TOKENIZER_CORPUS_CAP: # Mantém apenas as TOKENIZER_CORPUS_CAP mais recentes overflow_corpus = len(tokenizer_corpus_buffer) - TOKENIZER_CORPUS_CAP del tokenizer_corpus_buffer[:overflow_corpus] total_so_far_now = sum(samples_per_dataset.values()) samples_since_refit = total_so_far_now - last_tokenizer_refit_sample_count if (samples_since_refit >= TOKENIZER_REFIT_INTERVAL_SAMPLES and len(tokenizer_corpus_buffer) >= 200): # V6.7 — MEMORY GUARD before tokenizer refit # User requirement: "tokenizer-growth refit (was causing # crashes during BBPE parallel training at 1000-sample mark)". # Prova 14: ProcessPoolExecutor fork duplica RSS do processo # pai (modelo + tensores + buffers). Se RSS > 85% do cgroup # limit, SKIP refit (non-fatal) para evitar OOM-killer. import os as _os_mod_v67 try: with open("/proc/self/status") as _f_v67: _rss_line_v67 = [l for l in _f_v67 if l.startswith("VmRSS:")] _rss_kb_v67 = int(_rss_line_v67[0].split()[1]) if _rss_line_v67 else 0 _rss_mb_v67 = _rss_kb_v67 / 1024.0 # cgroup limit _cg_limit_mb_v67 = 4096.0 # default fallback try: with open("/sys/fs/cgroup/memory.max") as _f_cg_v67: _cg_val_v67 = _f_cg_v67.read().strip() if _cg_val_v67 and _cg_val_v67 != "max": _cg_limit_mb_v67 = int(_cg_val_v67) / 1024 / 1024 except Exception: pass # fallback default _rss_pct_v67 = _rss_mb_v67 / _cg_limit_mb_v67 if _cg_limit_mb_v67 > 0 else 0 logger.info( f"[V6.7-memory-guard] RSS={_rss_mb_v67:.0f}MB " f"({100*_rss_pct_v67:.1f}% of cgroup " f"{_cg_limit_mb_v67:.0f}MB)" ) if _rss_pct_v67 > 0.85: logger.warning( f"[V6.7-memory-guard] SKIP refit: RSS " f"{100*_rss_pct_v67:.1f}% > 85% of cgroup " f"(would trigger OOM via fork). " f"gc.collect() e prosseguir sem refit." ) gc.collect() if hasattr(torch, "cpu") and hasattr(torch.cpu, "empty_cache"): try: torch.cpu.empty_cache() except Exception: pass last_tokenizer_refit_sample_count = total_so_far_now tokenizer_growth_log.append({ "step": step, "chunk_idx": chunk_idx_global, "total_samples": total_so_far_now, "skipped_reason": "rss_exceeded_85pct", "rss_mb": _rss_mb_v67, "cgroup_limit_mb": _cg_limit_mb_v67, }) # SKIP refit — continue para próximo chunk raise _SkipRefitV67() except _SkipRefitV67: pass # já tratado acima except Exception as _guard_err_v67: logger.warning( f"[V6.7-memory-guard] guard failed (non-fatal): {_guard_err_v67}" ) # Refit tokenizer com corpus acumulado vocab_before = int(getattr(kls.tokenizer, "vocab_size", 0)) try: kls.tokenizer.fit(list(tokenizer_corpus_buffer), min_frequency=2) vocab_after = int(getattr(kls.tokenizer, "vocab_size", 0)) growth = vocab_after - vocab_before tokenizer_growth_log.append({ "step": step, "chunk_idx": chunk_idx_global, "total_samples": total_so_far_now, "corpus_size": len(tokenizer_corpus_buffer), "vocab_before": vocab_before, "vocab_after": vocab_after, "vocab_growth": growth, }) logger.info( f"[V7-tokenizer-growth] Refit @ sample " f"{total_so_far_now}: vocab {vocab_before} → " f"{vocab_after} (+{growth}, corpus=" f"{len(tokenizer_corpus_buffer)})" ) last_tokenizer_refit_sample_count = total_so_far_now except Exception as tok_fit_err: logger.warning( f"[V7-tokenizer-growth] Refit failed (non-fatal): " f"{tok_fit_err}" ) except Exception as tok_col_err: logger.warning( f"[V7-tokenizer-growth] Corpus collection failed: {tok_col_err}" ) # V6.5-V2-metrics — Computa métricas SOM após cada chunk # (QE, TE, Kaski-Lagus, Variância Explicada, Dead Neurons, # Colapso Topológico, Estagnação QE, Cruzamento Vizinhança). # User requirement: "ANALISAR matematicamente e logicamente e # inserir melhorias para os scripts das métricas de aprendizado # e indicadores de falha em Mapas Auto-Organizáveis". # # V6.5-V2-metrics-FIX-2 — adiciona detecção de hipóteses ativas, # neurônios ativos (BMU), passos de treino de hipóteses usados, # e detecção explícita de NaN em QE/KL. User requirement (latest): # "Aprimorar a FASE2 PUNITIVA investigar QE e KL resultando em NAN, # acrescentar detecção da atividade das layers (quantas ativadas) # de hipóteses e das ativações dos neurônios, quantos passos de # treino de hipóteses usado". try: som_metrics = kls.compute_som_metrics() hyp_activity = kls.detect_hypothesis_activity() neuron_activity = kls.count_active_neurons() hyp_steps_used = kls.get_hyp_train_steps_used() som_metrics_log.append({ "step": step, "chunk_idx": chunk_idx_global, "dataset": dataset_name, "quantization_error": float(som_metrics["quantization_error"]), "topological_error": float(som_metrics["topological_error"]), "kaski_lagus_error": float(som_metrics["kaski_lagus_error"]), "explained_variance_share": float(som_metrics["explained_variance_share"]), "overall_health": som_metrics["overall_health"], "n_failure_indicators": int(som_metrics["n_failure_indicators"]), "failure_indicators": som_metrics["failure_indicators"], "collapse_severity": som_metrics["topological_collapse"].get("severity", "none"), "dead_neuron_rate": float(som_metrics["dead_neuron_rate"].get("dead_neuron_rate", 0.0)), # V6.5-V2-metrics-FIX-2 "nan_detected": bool(som_metrics.get("nan_detected", False)), "nan_reason": som_metrics.get("nan_reason"), "qe_was_nan": bool(som_metrics.get("qe_was_nan", False)), "kl_was_nan": bool(som_metrics.get("kl_was_nan", False)), "n_hypotheses_active": int(hyp_activity["n_hypotheses_active"]), "n_hypotheses_total": int(hyp_activity["n_hypotheses_total"]), "hypothesis_activation_rate": float(hyp_activity["activation_rate"]), "mean_delta_norm": float(hyp_activity["mean_delta_norm"]), "max_delta_norm": float(hyp_activity["max_delta_norm"]), "n_active_neurons_bmu": int(neuron_activity["n_active_neurons"]), "n_dead_neurons_bmu": int(neuron_activity["n_dead_neurons"]), "neuron_activation_rate": float(neuron_activity["neuron_activation_rate"]), "current_hyp_train_steps": int(hyp_steps_used["current_hyp_train_steps"]), "total_hyp_steps_executed": int(hyp_steps_used["total_hyp_train_steps_executed"]), "n_train_hyp_calls": int(hyp_steps_used["n_train_hypotheses_calls"]), }) if chunk_idx_global % 5 == 0: logger.info( f"[V6.5-V2-metrics-FIX-2] chunk {chunk_idx_global}: " f"QE={som_metrics['quantization_error']:.4f}, " f"TE={som_metrics['topological_error']:.4f}, " f"VE={som_metrics['explained_variance_share']:.4f}, " f"health={som_metrics['overall_health']}, " f"failures={len(som_metrics['failure_indicators'])}, " f"hyp_active={hyp_activity['n_hypotheses_active']}/{hyp_activity['n_hypotheses_total']}, " f"neurons_active={neuron_activity['n_active_neurons']}/{neuron_activity['n_total_neurons']}" ) except Exception as e: logger.warning(f"[V6.5-V2-metrics-FIX-2] Failed to compute SOM metrics: {e}") # V6.5-V2-metrics-FIX-4 — Auto-revive neurônios mortos # User requirement: "APRIMORAR (...) a ativação e uso e acesso # dos neurônios (apenas dois estão sendo ativados: # neurons_active=2/864) distribuindo o processamento paralelamente". # O conscience mechanism (DeSieno 1988) já força distribuição # uniforme de BMU, mas como safety net adicional, revive # explicitamente neurônios que ainda estão mortos após o cooldown. try: revival = kls.auto_revive_if_needed( dead_rate_threshold=AUTO_REVIVE_DEAD_RATE_THRESHOLD, min_steps_between_revivals=AUTO_REVIVE_COOLDOWN_STEPS, ) if revival.get("action") == "auto_revived": logger.info( f"[V6.5-V2-metrics-FIX-4] AUTO-REVIVE triggered: " f"n_revived={revival['n_revived']}/{revival['n_total']}, " f"dead_rate {revival['dead_rate_before']:.3f} → " f"{revival['dead_rate_after']:.3f}, " f"steps_since_last={revival['steps_since_last_revival']}" ) except Exception as revive_err: logger.warning( f"[V6.5-V2-metrics-FIX-4] auto_revive_if_needed failed: {revive_err}" ) # V6.5-V2-auto-adjust — Aplica auto-ajustes SOM com base nas # 8 métricas canônicas. User requirement: "inserir melhorias # para os scripts das métricas de aprendizado e indicadores de # falha em Mapas Auto-Organizáveis" + "PORTANTO: observar se # os valores demonstram que o modelo esteja aprendendo e parar # caso não esteja". try: adjust_result = auto_adjuster.adjust( kls=kls, som_metrics=som_metrics, batch_idx=chunk_idx_global, ) if adjust_result.get("actions"): auto_adjust_log.append({ "step": step, "chunk_idx": chunk_idx_global, "dataset": dataset_name, "phase": "CONHECIMENTO", "actions": adjust_result["actions"], "boosts": adjust_result["boosts"], "n_actions": adjust_result["n_actions"], }) if chunk_idx_global % 5 == 0: logger.info( f"[V6.5-V2-auto-adjust] chunk {chunk_idx_global}: " f"actions={adjust_result['n_actions']}, " f"boosts(α,σ)=({adjust_result['boosts']['alpha0_boost']:.3f}, " f"{adjust_result['boosts']['sigma0_boost']:.3f})" ) # STOP-LEARNING detection (user: "parar caso não esteja") # V6.5-V2-auto-conscience-v2 — FASE1 NÃO para mais por # stop_learning. Continua processando todos os 8 datasets # para acumular conhecimento. Apenas registra o evento # para auditoria. FASE2 ainda pode parar se necessário. if adjust_result.get("stop_training"): stop_training_reason = adjust_result.get("stop_reason") logger.warning( f"[V6.5-V2-auto-conscience-v2] FASE1 stop_learning " f"DETECTED at chunk {chunk_idx_global}: {stop_training_reason} " f"— CONTINUING para acumular conhecimento (não para)." ) # Reset counters to allow recovery (γ boost já foi aplicado) try: auto_adjuster.state.reset_consecutive_counters() except Exception: pass # NÃO seta storage_critical_stopped = True # NÃO break — continua processando except Exception as adjust_err: logger.warning( f"[V6.5-V2-auto-adjust] adjust() failed: {adjust_err}" ) # Aggressive memory cleanup between chunks if chunk_idx_global % 2 == 0: aggressive_memory_cleanup() # V6.5-V2-metrics-FIX: também chama aggressive_cleanup do KLS # para zerar gradientes Adam e truncar históricos. try: kls.aggressive_cleanup() except Exception: pass # V6.5-V2-metrics-FIX: pausa pós-processamento para dar tempo # de conclusão (user requirement: "todo streaming deve ter # pausa para dar tempo de conclusão de processamento"). time.sleep(POST_PROCESSING_PAUSE_S_FASE1) except Exception as e: logger.error(f"[V6.5-V2] Dataset {dataset_name} failed: {e}") traceback.print_exc() streaming_failures[dataset_name] += 1 time.sleep(INTER_DATASET_PAUSE_S_FASE1) # V6.5-V2-metrics-FIX-2 — Save state after each dataset to preserve # progress in case the process is killed by container timeout. # User requirement: "o estado do modelo deve ser contínuo" — saving # after each dataset ensures the state is always recoverable. try: total_now = sum(samples_per_dataset.values()) partial_state_path = BIGRU_ROOT / f"v6_5_v2_conhecimento_partial_d{ds_idx+1}.pt" save_model_states_for_evaluation( kls=kls, reason=f"conhecimento_after_dataset_{ds_idx+1}_{dataset_name.split('/')[-1]}", step=step, total_samples=total_now, output_path=partial_state_path, ) # Remove older partial states (keep only the latest) for old_idx in range(1, ds_idx + 1): old_path = BIGRU_ROOT / f"v6_5_v2_conhecimento_partial_d{old_idx}.pt" if old_path.exists() and old_path != partial_state_path: try: old_path.unlink() except Exception: pass except Exception as e: logger.warning(f"[V6.5-V2-metrics-FIX-2] Failed to save partial state: {e}") # V6.5-V2-metrics-FIX-2 — Limpa cache HF datasets entre datasets para # evitar acúmulo de cache de parquet (causava OOM com 1.1GB+ de cache). try: import gc as _gc _gc.collect() _gc.collect() # Limpa cache de datasets streaming (não apaga arquivos baixados, # apenas libera referências em memória) try: import datasets if hasattr(datasets, '_datasets_cache'): datasets._datasets_cache.clear() except Exception: pass # Força liberação de tensores não referenciados aggressive_memory_cleanup() try: kls.aggressive_cleanup() except Exception: pass except Exception: pass t_elapsed = time.time() - t_start total_samples = sum(samples_per_dataset.values()) logger.info(f"\n[V6.5-V2] FASE 1 (CONHECIMENTO) concluída:") logger.info(f" Total samples: {total_samples}") logger.info(f" Meta atingida: {total_samples >= META_MINIMA_CONHECIMENTO}") logger.info(f" Tempo total: {t_elapsed:.1f}s") logger.info(f" Storage critical stopped: {storage_critical_stopped}") return { "phase": "CONHECIMENTO", "total_samples": int(total_samples), "meta_minima": int(META_MINIMA_CONHECIMENTO), "meta_atingida": bool(total_samples >= META_MINIMA_CONHECIMENTO), "samples_per_dataset": samples_per_dataset, "streaming_failures": streaming_failures, "storage_critical_stopped": bool(storage_critical_stopped), "elapsed_s": float(t_elapsed), "n_steps": int(step), "n_chunks": int(chunk_idx_global), "accuracies_log": accuracies_log[-50:], # últimos 50 "attention_log": attention_log[-20:], # V6.5-V2-metrics — log de métricas SOM durante CONHECIMENTO "som_metrics_log": som_metrics_log[-30:], # últimos 30 chunks "som_metrics_final": ( kls.compute_som_metrics() if kls.buffer_4d else {} ), "som_metric_history": kls.get_som_metric_history(), "final_state": kls.get_v2_metrics(), # V6.5-V2-auto-adjust — log de auto-ajustes aplicados "auto_adjust_log": auto_adjust_log[-50:], "auto_adjust_state": auto_adjuster.state.to_dict(), "stop_training_reason": stop_training_reason, # V7-tokenizer-growth — diagnóstico de crescimento do vocabulário "tokenizer_growth_log": tokenizer_growth_log, "tokenizer_final_vocab_size": int(getattr(kls.tokenizer, "vocab_size", 0)), "tokenizer_initial_vocab_size": int(initial_vocab_size), "tokenizer_target_vocab_size": int(VOCAB_SIZE), } # ============================================================================ # 9.5 — V6.5-V2-metrics-FIX-4: Sumário final da FASE2 # ============================================================================ def build_fase2_final_summary( kls: KohonenLearningSystemV2, som_metrics_log: List[Dict[str, Any]], punishment_log: List[Dict[str, Any]], hypotheses_log: List[Dict[str, Any]], delta_applications_log: List[Dict[str, Any]], elapsed_s: float, ) -> Dict[str, Any]: """V6.5-V2-metrics-FIX-4 — Constrói sumário final da FASE2 mostrando evolução. User requirement: "ao final da FASE2 mostrar evolução de métricas e dos indicadores e da taxa de aprendizagem". Extrai trajetória temporal das métricas SOM (QE, TE, KL, VE), indicadores de falha (dead_rate, collapse, stagnation, crossing), taxa de aprendizado (α_t, σ_t) e estatísticas do conscience mechanism (neurons_active, uniformity_score). Compara início vs fim para mostrar evolução. Returns: Dict com: - summary_lines: List[str] — linhas formatadas para logger.info - metrics_evolution: Dict com first/last/delta de cada métrica - learning_rate_evolution: Dict com α_t, σ_t no início e fim - indicators_evolution: Dict com indicadores de falha no início e fim - neuron_activation_evolution: Dict com neurons_active no início e fim - punishment_stats: Dict com estatísticas de punição """ summary_lines: List[str] = [] # 1. Trajetória das métricas principais (QE, TE, KL, VE) def _safe_get(log_list, key, idx): if not log_list or idx >= len(log_list): return 0.0 try: return float(log_list[idx].get(key, 0.0)) except (TypeError, ValueError): return 0.0 n_logs = len(som_metrics_log) first_idx = 0 last_idx = max(0, n_logs - 1) metrics_first = { "QE": _safe_get(som_metrics_log, "quantization_error", first_idx), "TE": _safe_get(som_metrics_log, "topological_error", first_idx), "KL": _safe_get(som_metrics_log, "kaski_lagus_error", first_idx), "VE": _safe_get(som_metrics_log, "explained_variance_share", first_idx), } metrics_last = { "QE": _safe_get(som_metrics_log, "quantization_error", last_idx), "TE": _safe_get(som_metrics_log, "topological_error", last_idx), "KL": _safe_get(som_metrics_log, "kaski_lagus_error", last_idx), "VE": _safe_get(som_metrics_log, "explained_variance_share", last_idx), } metrics_delta = { k: metrics_last[k] - metrics_first[k] for k in metrics_first } # 2. Indicadores de falha first_failures = som_metrics_log[first_idx].get("failure_indicators", []) if som_metrics_log else [] last_failures = som_metrics_log[last_idx].get("failure_indicators", []) if som_metrics_log else [] first_health = som_metrics_log[first_idx].get("overall_health", "unknown") if som_metrics_log else "unknown" last_health = som_metrics_log[last_idx].get("overall_health", "unknown") if som_metrics_log else "unknown" # 3. Ativação de neurônios (conscience mechanism) first_neurons_active = _safe_get(som_metrics_log, "n_active_neurons_bmu", first_idx) last_neurons_active = _safe_get(som_metrics_log, "n_active_neurons_bmu", last_idx) first_neurons_total = _safe_get(som_metrics_log, "n_total_neurons_bmu", first_idx) or 864 last_neurons_total = _safe_get(som_metrics_log, "n_total_neurons_bmu", last_idx) or 864 # 4. Relatório paralelo final do SOM (conscience + uniformity) try: parallel_report = kls.parallel_neuron_activation_report() except Exception: parallel_report = {} # 5. Taxa de aprendizado (α_t, σ_t) do SOM try: som_metrics = kls.som.get_metrics() alpha_t = som_metrics.get("alpha_t_effective", 0.0) sigma_t = som_metrics.get("sigma_t_effective", 0.0) alpha0 = som_metrics.get("alpha0", 0.5) sigma0 = som_metrics.get("sigma0", 3.0) som_t = som_metrics.get("t", 0) except Exception: alpha_t = sigma_t = alpha0 = sigma0 = som_t = 0.0 # 6. Estatísticas de punição n_punishments = len(punishment_log) n_train_hyp_calls = len(hypotheses_log) n_delta_apps = len(delta_applications_log) # Taxa de sucesso das aplicações de delta (acc_after > acc_before) delta_success = 0 if delta_applications_log: for d in delta_applications_log: try: if float(d.get("acc_after", 0.0)) > float(d.get("acc_before", 0.0)): delta_success += 1 except (TypeError, ValueError): pass delta_success_rate = float(delta_success / max(1, n_delta_apps)) # 7. Hipóteses — evolução da loss if hypotheses_log: loss_first = float(hypotheses_log[0].get("loss_initial", 0.0)) loss_last_init = float(hypotheses_log[-1].get("loss_initial", 0.0)) loss_last_final = float(hypotheses_log[-1].get("loss_final", 0.0)) else: loss_first = loss_last_init = loss_last_final = 0.0 # ===================== MONTAGEM DAS LINHAS DE SUMÁRIO ===================== summary_lines.append(f" Duração total: {elapsed_s:.1f}s") summary_lines.append(f" Logs de métricas computados: {n_logs}") summary_lines.append("") summary_lines.append(" ─── EVOLUÇÃO DAS MÉTRICAS PRINCIPAIS (início → fim) ───") summary_lines.append( f" QE (Quantization Error): {metrics_first['QE']:.4f} → " f"{metrics_last['QE']:.4f} (Δ={metrics_delta['QE']:+.4f})" ) summary_lines.append( f" TE (Topological Error) : {metrics_first['TE']:.4f} → " f"{metrics_last['TE']:.4f} (Δ={metrics_delta['TE']:+.4f})" ) summary_lines.append( f" KL (Kaski-Lagus) : {metrics_first['KL']:.4f} → " f"{metrics_last['KL']:.4f} (Δ={metrics_delta['KL']:+.4f})" ) summary_lines.append( f" VE (Explained Variance): {metrics_first['VE']:.4f} → " f"{metrics_last['VE']:.4f} (Δ={metrics_delta['VE']:+.4f})" ) summary_lines.append("") summary_lines.append(" ─── INDICADORES DE FALHA ───") summary_lines.append( f" Overall health: {first_health} → {last_health}" ) summary_lines.append( f" Failure indicators (início): {len(first_failures)} — {first_failures[:3]}" ) summary_lines.append( f" Failure indicators (fim) : {len(last_failures)} — {last_failures[:3]}" ) summary_lines.append("") summary_lines.append(" ─── ATIVAÇÃO DE NEURÔNIOS (Conscience Mechanism) ───") summary_lines.append( f" Neurons ativos (BMU): {int(first_neurons_active)}/{int(first_neurons_total)} " f"→ {int(last_neurons_active)}/{int(last_neurons_total)}" ) if parallel_report: summary_lines.append( f" Uniformity score (fim): {parallel_report.get('uniformity_score', 0.0):.4f} " f"(1.0 = perfeitamente uniforme)" ) summary_lines.append( f" Conscience bias mean/std: " f"{parallel_report.get('conscience_bias_mean', 0.0):.6f} / " f"{parallel_report.get('conscience_bias_std', 0.0):.6f} " f"(deve tender a 0)" ) summary_lines.append( f" Win count max/mean: " f"{parallel_report.get('max_win_count', 0.0):.0f} / " f"{parallel_report.get('mean_win_count', 0.0):.2f}" ) summary_lines.append("") summary_lines.append(" ─── TAXA DE APRENDIZADO (Kohonen schedule) ───") summary_lines.append( f" α_t (learning rate): α₀={alpha0:.4f} → α_t={alpha_t:.6f} " f"(floor=0.001, decai em exp(-t/2000))" ) summary_lines.append( f" σ_t (neighborhood) : σ₀={sigma0:.4f} → σ_t={sigma_t:.6f} " f"(floor=0.1, decai em exp(-t/1000))" ) summary_lines.append(f" SOM t (updates) : {som_t}") summary_lines.append("") summary_lines.append(" ─── ESTATÍSTICAS DE PUNIÇÃO ───") summary_lines.append(f" Total punishment events: {n_punishments}") summary_lines.append(f" Hypotheses trainings : {n_train_hyp_calls}") summary_lines.append(f" Delta applications : {n_delta_apps}") summary_lines.append( f" Delta success rate : {delta_success_rate:.2%} " f"({delta_success}/{n_delta_apps} melhoraram acurácia)" ) if hypotheses_log: summary_lines.append( f" Loss inicial primeira chamada: {loss_first:.6f}" ) summary_lines.append( f" Loss inicial última chamada : {loss_last_init:.6f} " f"→ final: {loss_last_final:.6f}" ) summary_lines.append("") summary_lines.append(" ─── HIPÓTESES (configuração final) ───") try: v2_metrics = kls.get_v2_metrics() summary_lines.append( f" n_hypotheses ativas: {v2_metrics.get('n_hypotheses', 0)} / " f"max {v2_metrics.get('max_n_hypotheses', 0)}" ) summary_lines.append( f" hyp_train_steps atual: {v2_metrics.get('hyp_train_steps', 0)}" ) summary_lines.append( f" total_hyp_steps_executed: " f"{v2_metrics.get('total_hyp_steps_executed', 0)}" ) summary_lines.append( f" n_adaptations dinâmicas: " f"{v2_metrics.get('dynamic_adaptation', {}).get('n_adaptations', 0)}" ) summary_lines.append( f" EWC reference set: {v2_metrics.get('ewc_reference_set', False)}" ) except Exception: pass return { "summary_lines": summary_lines, "metrics_evolution": { "first": metrics_first, "last": metrics_last, "delta": metrics_delta, }, "learning_rate_evolution": { "alpha0": alpha0, "alpha_t_final": alpha_t, "sigma0": sigma0, "sigma_t_final": sigma_t, "som_t": som_t, }, "indicators_evolution": { "first_health": first_health, "last_health": last_health, "first_failures": first_failures, "last_failures": last_failures, "first_n_failures": len(first_failures), "last_n_failures": len(last_failures), }, "neuron_activation_evolution": { "first_active": int(first_neurons_active), "last_active": int(last_neurons_active), "total": int(last_neurons_total), "parallel_report_final": parallel_report, }, "punishment_stats": { "n_punishments": n_punishments, "n_train_hyp_calls": n_train_hyp_calls, "n_delta_apps": n_delta_apps, "delta_success_rate": delta_success_rate, "loss_first_initial": loss_first, "loss_last_initial": loss_last_init, "loss_last_final": loss_last_final, }, } # ============================================================================ # 10. PHASE 2 — TREINAMENTO COM PUNIÇÃO (16 hipóteses × 3 tentativas) # ============================================================================ def run_fase_punicão( kls: KohonenLearningSystemV2, hf_token: Optional[str], ) -> Dict[str, Any]: """Fase 2 — TREINAMENTO COM PUNIÇÃO: streaming 2000 samples + V2 protocol. User requirement: "FASE2 TREINAMENTO (meta mínima 2000 samples ou mais) COM PUNIÇÃO ATIVA para 'BrunoN-Dev/corpus-ptbr-v1' de 100 em 100 samples" + "esta FASE2 (nunca acontece antes da FASE1) deve usar do arquivo de Estado do MODELO concluído da FASE1 e a FASE2 tem função: de treinar, ajustar dados e os parâmetros do MAPA-SOM de redes Kohonen, informar métricas de aprendizado do modelo" + Métricas SOM canônicas monitoradas (user requirement): 1.1. Erro de Quantização (QE) 1.2. Erro Topológico (TE) 1.3. Erro de Kaski-Lagus 1.4. Variância Explicada 2.1. Colapso Topológico 2.2. Taxa de Neurônios Mortos 2.3. Estagnação do QE 2.4. Cruzamento de Vizinhança Returns: Dict com métricas da fase TREINAMENTO COM PUNIÇÃO. """ logger.info("\n" + "=" * 80) logger.info("[V6.5-V2] FASE 2 — TREINAMENTO COM PUNIÇÃO (16 hipóteses × 3 tentativas)") logger.info("=" * 80) logger.info(f" Dataset: {PUNICAO_DATASET}") logger.info(f" Max samples: {MAX_SAMPLES_PUNICAO}") logger.info(f" Meta mínima: {META_MINIMA_PUNICAO}") logger.info(f" n_hypotheses: {N_HYPOTHESES}") logger.info(f" n_trials: {N_TRIALS}") logger.info(f" hyp_train_steps: {HYP_TRAIN_STEPS}") logger.info("=" * 80 + "\n") # V6.5-V2-metrics — reseta histórico de métricas SOM no início da FASE 2 # (PUNITIVA tem seu próprio histórico independente da FASE 1, para que # Estagnação do QE e Cruzamento de Vizinhança sejam avaliados # exclusivamente sobre o regime punitivo). kls.reset_som_metric_history() logger.info("[V6.5-V2-metrics] Histórico de métricas SOM resetado para FASE 2 (PUNITIVA).") # V7 — Inicializa integrator (AdaptiveDPOLoss + AdvancedSOMAugmenter) # User requirement: "usar V7 sabendo que data_augmentation.py e dpo.py # agora devem ser usados na FASE2". v7_integrator = None if _V7_INTEGRATOR_AVAILABLE: try: v7_integrator = V7Fase2Integrator( oom_guard=OOM_GUARD, som_dims=SOM_GRID, # (4,4,4,4) canônico som_lr=ALPHA0, # 0.5 canônico som_sigma=SIGMA0, # 2.0 canônico beta_init=0.1, target_kl=0.2, lr_beta=1e-3, log_every_n_batches=5, fit_augmenter_every_n_chunks=2, ) v7_integrator.start_fase2(kls) logger.info("[V7] FASE2 integrator inicializado (AdaptiveDPOLoss + AdvancedSOMAugmenter)") except Exception as v7_init_err: logger.warning(f"[V7] Falha ao inicializar integrator (FASE2 continuará sem V7): {v7_init_err}") v7_integrator = None else: logger.warning("[V7] V7Fase2Integrator não disponível — FASE2 rodará sem V7 hooks") step = 0 t_start = time.time() samples_processed = 0 storage_critical_stopped = False chunk_idx_global = 0 punishment_log: List[Dict[str, Any]] = [] hypotheses_log: List[Dict[str, Any]] = [] accuracies_log: List[Dict[str, Any]] = [] delta_applications_log: List[Dict[str, Any]] = [] # V6.5-V2-metrics — log de métricas SOM durante PUNIÇÃO (computadas a cada batch) som_metrics_log: List[Dict[str, Any]] = [] # V6.5-V2-auto-adjust — Auto-ajustador SOM para FASE2 (mais agressivo # porque FASE2 é punitiva e sujeita a gradientes explosivos). # User requirement: "aprimorar matematicamente a lógica de autoajustes # analisando Kohonen" + "PORTANTO: observar se os valores demonstram que # o modelo esteja aprendendo e parar caso não esteja". auto_adjuster = create_auto_adjuster() auto_adjust_log: List[Dict[str, Any]] = [] stop_training_reason: Optional[str] = None try: chunk_iter = stream_dataset_in_chunks( dataset_name=PUNICAO_DATASET, chunk_size=STREAM_BATCH_SIZE, max_total=MAX_SAMPLES_PUNICAO, hf_token=hf_token, timeout_per_chunk=90, ) for chunk_samples in chunk_iter: if storage_critical_stopped: break if samples_processed >= META_MINIMA_PUNICAO: logger.info( f"[V6.5-V2] Meta PUNIÇÃO atingida ({samples_processed} ≥ " f"{META_MINIMA_PUNICAO})." ) break chunk_idx_global += 1 if not chunk_samples: continue samples_processed += len(chunk_samples) labels = [make_label(s, PUNICAO_DATASET) for s in chunk_samples] logger.info( f"[V6.5-V2] PUNIÇÃO chunk {chunk_idx_global} " f"({len(chunk_samples)} samples) from {PUNICAO_DATASET}" ) for batch_start in range(0, len(chunk_samples), BATCH_SIZE): try: if step > 0 and step % STORAGE_CHECK_INTERVAL_STEPS == 0: disk_pct = get_disk_usage_pct() # V6.5-V2-buffer-864 — Monitora memória na FASE2 também try: mem_report_f2 = check_memory_and_maybe_fallback(kls) if mem_report_f2["fallback_ativado"]: logger.warning( f"[V6.5-V2-buffer-864] FASE2 buffer reduzido para " f"{mem_report_f2['buffer_max_size']} " f"(RSS={mem_report_f2['rss_mb']:.0f}MB, " f"{mem_report_f2['mem_pct']:.1f}%)" ) except Exception as _fb_err_f2: logger.warning( f"[V6.5-V2-buffer-864] FASE2 check_memory failed: {_fb_err_f2}" ) if disk_pct > STORAGE_CRITICAL_PCT: logger.error( f"[V6.5-V2] STORAGE CRITICAL: disk={disk_pct:.1f}%. " f"STOPPING PUNIÇÃO." ) save_model_states_for_evaluation( kls=kls, reason=f"storage_critical_punicão_disk_{disk_pct:.1f}pct", step=step, total_samples=samples_processed, ) storage_critical_stopped = True break batch_sents = chunk_samples[batch_start: batch_start + BATCH_SIZE] batch_labels = labels[batch_start: batch_start + BATCH_SIZE] try: # FASE 2: enable_punishment=True (V2 protocol) result = kls.process_batch_v2( batch_sents, batch_labels, label_strings=DATASET_LABEL_STRINGS.get(PUNICAO_DATASET), dataset_name=PUNICAO_DATASET, enable_punishment=True, ) except (MemoryError, RuntimeError) as oom_err_fase2: # V6.5-V2-oom-guard — OOM-killer protection FASE2 # User requirement: "o processo vem sendo morto OOM-kiler". # FASE2 é mais pesada (hypothesis training + EWC) e mais # propensa a OOM. Captura MemoryError e RuntimeError("out # of memory") explicitamente para evitar morte do processo. is_oom_f2 = ( isinstance(oom_err_fase2, MemoryError) or "out of memory" in str(oom_err_fase2).lower() or "cuda" in str(oom_err_fase2).lower() ) if is_oom_f2: logger.error( f"[V6.5-V2-oom-guard] FASE2 OOM detected: " f"{type(oom_err_fase2).__name__}: {str(oom_err_fase2)[:200]}" ) logger.error("[V6.5-V2-oom-guard] Triggering aggressive memory cleanup...") cleanup_result = aggressive_memory_cleanup() logger.error( f"[V6.5-V2-oom-guard] FASE2 cleanup freed " f"{cleanup_result.get('freed_mb', 0):.1f}MB " f"(RSS: {cleanup_result.get('rss_before_mb', 0):.0f}MB → " f"{cleanup_result.get('rss_after_mb', 0):.0f}MB)" ) # Trunca buffer para reduzir pressão if len(kls.buffer_4d) > 32: overflow = len(kls.buffer_4d) - 32 del kls.buffer_4d[:overflow] del kls.buffer_labels[:overflow] logger.error( f"[V6.5-V2-oom-guard] FASE2 buffer truncated to 32 " f"(dropped {overflow} samples)" ) continue else: logger.error(f"[V6.5-V2] FASE2 process_batch_v2 RuntimeError: {oom_err_fase2}") traceback.print_exc() continue except Exception as e: logger.error(f"[V6.5-V2] process_batch_v2 error: {e}") traceback.print_exc() continue # V6.5-V2-oom-guard — gc.collect() entre batches FASE2 # (PUNITIVA é mais pesada — hypothesis training + EWC + revival) if step % 2 == 0: gc.collect() # Log punishment events action = result.get("action", "none") if action != "none" and action != "success": punishment_log.append({ "step": step, "action": action, "accuracy": float(result.get("accuracy", 0.0)), "punishment_count": int(kls.punishment_count), "cycle_completed": bool(result.get("cycle_completed", False)), }) # Log hypotheses training if "hypotheses_training" in result: hyp_info = result["hypotheses_training"] hypotheses_log.append({ "step": step, "loss_initial": float(hyp_info.get("loss_initial", 0.0)), "loss_final": float(hyp_info.get("loss_final", 0.0)), "delta_scale": float(hyp_info.get("delta_scale_final", 0.0)), "elapsed_ms": float(hyp_info.get("elapsed_ms", 0.0)), }) # Log delta applications if "delta_application" in result: apply_info = result["delta_application"] delta_applications_log.append({ "step": step, "acc_before": float(apply_info.get("acc_before", 0.0)), "acc_after": float(apply_info.get("acc_after", 0.0)), "best_acc_during_selection": float(apply_info.get("best_acc_during_selection", 0.0)), "delta_norm": float(apply_info.get("delta_norm", 0.0)), "ewc_reference_set": bool(apply_info.get("ewc_reference_set", False)), "elapsed_ms": float(apply_info.get("elapsed_ms", 0.0)), }) # Accuracy log acc = kls.evaluate_classification() accuracies_log.append({ "step": step, "accuracy": float(acc), "phase": "PUNICAO", "action": action, "buffer_size": int(len(kls.buffer_4d)), "punishment_count": int(kls.punishment_count), "success_count": int(kls.success_count), "ewc_reference_set": bool(kls.som.old_weights_w is not None), }) # V6.5-V2-metrics — Computa métricas SOM após cada batch # durante a FASE 2 PUNITIVA. User requirement: # "ANALISAR matematicamente e logicamente e inserir melhorias # para os scripts das métricas de aprendizado e indicadores # de falha em Mapas Auto-Organizáveis (SOM / Redes de Kohonen) # avaliam a fidelidade de representação dos dados e a # preservação da vizinhança topológica". # # V6.5-V2-metrics-FIX-2 — User requirement (latest): # "Aprimorar a FASE2 PUNITIVA investigar QE e KL resultando # em NAN, acrescentar detecção da atividade das layers # (quantas ativadas) de hipóteses e das ativações dos # neurônios, quantos passos de treino de hipóteses usado". try: som_metrics = kls.compute_som_metrics() # V6.5-V2-metrics-FIX-2 — detecção de atividade de # hipóteses (quantas das n_hypotheses ativas produzem # delta com norma > threshold) hyp_activity = kls.detect_hypothesis_activity() # V6.5-V2-metrics-FIX-2 — contagem de neurônios ativos # (quantos neurônios são BMU para ≥1 amostra do buffer) neuron_activity = kls.count_active_neurons() # V6.5-V2-metrics-FIX-2 — passos de treino de hipóteses # usados (cumulativo + atual) hyp_steps_used = kls.get_hyp_train_steps_used() som_metrics_log.append({ "step": step, "chunk_idx": chunk_idx_global, "action": action, "punishment_count": int(kls.punishment_count), "quantization_error": float(som_metrics["quantization_error"]), "topological_error": float(som_metrics["topological_error"]), "kaski_lagus_error": float(som_metrics["kaski_lagus_error"]), "explained_variance_share": float(som_metrics["explained_variance_share"]), "overall_health": som_metrics["overall_health"], "n_failure_indicators": int(som_metrics["n_failure_indicators"]), "failure_indicators": som_metrics["failure_indicators"], "collapse_severity": som_metrics["topological_collapse"].get("severity", "none"), "collapse_effective_rank": float(som_metrics["topological_collapse"].get("effective_rank", 0.0)), "dead_neuron_rate": float(som_metrics["dead_neuron_rate"].get("dead_neuron_rate", 0.0)), "n_dead_neurons": int(som_metrics["dead_neuron_rate"].get("n_dead", 0)), "n_active_neurons": int(som_metrics["dead_neuron_rate"].get("n_active", 0)), "qe_stagnation": som_metrics.get("qe_stagnation", {}), "neighborhood_crossing": som_metrics.get("neighborhood_crossing", {}), # V6.5-V2-metrics-FIX-2 — NaN detection "nan_detected": bool(som_metrics.get("nan_detected", False)), "nan_reason": som_metrics.get("nan_reason"), "qe_was_nan": bool(som_metrics.get("qe_was_nan", False)), "kl_was_nan": bool(som_metrics.get("kl_was_nan", False)), # V6.5-V2-metrics-FIX-2 — hypothesis activity "n_hypotheses_active": int(hyp_activity["n_hypotheses_active"]), "n_hypotheses_total": int(hyp_activity["n_hypotheses_total"]), "hypothesis_activation_rate": float(hyp_activity["activation_rate"]), "mean_delta_norm": float(hyp_activity["mean_delta_norm"]), "max_delta_norm": float(hyp_activity["max_delta_norm"]), # V6.5-V2-metrics-FIX-2 — neuron activations "n_active_neurons_bmu": int(neuron_activity["n_active_neurons"]), "n_dead_neurons_bmu": int(neuron_activity["n_dead_neurons"]), "neuron_activation_rate": float(neuron_activity["neuron_activation_rate"]), # V6.5-V2-metrics-FIX-2 — hyp training steps used "current_hyp_train_steps": int(hyp_steps_used["current_hyp_train_steps"]), "total_hyp_steps_executed": int(hyp_steps_used["total_hyp_train_steps_executed"]), "n_train_hyp_calls": int(hyp_steps_used["n_train_hypotheses_calls"]), }) # Log detalhado a cada 5 steps OU quando há falha OU quando NaN detectado nan_alert = ( som_metrics.get("nan_detected", False) or som_metrics.get("qe_was_nan", False) or som_metrics.get("kl_was_nan", False) ) if step % 5 == 0 or som_metrics["n_failure_indicators"] > 0 or nan_alert: nan_marker = " [NaN DETECTED]" if nan_alert else "" logger.info( f"[V6.5-V2-metrics-FIX-2] PUNIÇÃO step {step}{nan_marker}: " f"QE={som_metrics['quantization_error']:.4f}, " f"TE={som_metrics['topological_error']:.4f}, " f"KL={som_metrics['kaski_lagus_error']:.4f}, " f"VE={som_metrics['explained_variance_share']:.4f}, " f"dead={som_metrics['dead_neuron_rate'].get('dead_neuron_rate', 0.0):.3f}, " f"health={som_metrics['overall_health']}, " f"failures={som_metrics['failure_indicators']}, " f"hyp_active={hyp_activity['n_hypotheses_active']}/{hyp_activity['n_hypotheses_total']} " f"(rate={hyp_activity['activation_rate']:.2f}), " f"neurons_active={neuron_activity['n_active_neurons']}/{neuron_activity['n_total_neurons']} " f"(rate={neuron_activity['neuron_activation_rate']:.2f}), " f"hyp_steps={hyp_steps_used['current_hyp_train_steps']} " f"(total={hyp_steps_used['total_hyp_train_steps_executed']}, " f"calls={hyp_steps_used['n_train_hypotheses_calls']})" ) if nan_alert: logger.warning( f"[V6.5-V2-metrics-FIX-2] NaN DETECTED at step {step}! " f"reason={som_metrics.get('nan_reason')}. " f"Investigate train_hypotheses() for gradient explosion." ) except Exception as e: logger.warning(f"[V6.5-V2-metrics-FIX-2] Failed to compute SOM metrics: {e}") traceback.print_exc() step += 1 # V7 — post-batch hook (AdaptiveDPOLoss + OomGuard record_loss) if v7_integrator is not None: try: v7_integrator.post_batch_hook(kls, result, step - 1) except Exception as v7_pb_err: logger.debug(f"[V7] post_batch_hook error (non-fatal): {v7_pb_err}") # V6.5-V2-metrics-FIX-4 — pausa maior na FASE2 (PUNITIVA) # User requirement: "FASE2 PUNIÇÃO é mais pesada é pode exigir # pausas do streaming até concluir o processamento". time.sleep(INTER_BATCH_PAUSE_S_FASE2) except Exception as e: logger.error(f"[V6.5-V2] Batch error in PUNIÇÃO: {e}") traceback.print_exc() continue # V6.5-V2-metrics-FIX-4 — Auto-revive neurônios mortos na FASE2 # (PUNITIVA também deve distribuir processamento paralelamente) try: revival = kls.auto_revive_if_needed( dead_rate_threshold=AUTO_REVIVE_DEAD_RATE_THRESHOLD, min_steps_between_revivals=AUTO_REVIVE_COOLDOWN_STEPS, ) if revival.get("action") == "auto_revived": logger.info( f"[V6.5-V2-metrics-FIX-4] PUNIÇÃO AUTO-REVIVE: " f"n_revived={revival['n_revived']}/{revival['n_total']}, " f"dead_rate {revival['dead_rate_before']:.3f} → " f"{revival['dead_rate_after']:.3f}" ) except Exception as revive_err: logger.warning( f"[V6.5-V2-metrics-FIX-4] PUNIÇÃO auto_revive failed: {revive_err}" ) # V6.5-V2-auto-adjust — Aplica auto-ajustes SOM na FASE2 (PUNITIVA). # User requirement: "Aprimorar a FASE2 PUNITIVA" + "PORTANTO: observar # se os valores demonstram que o modelo esteja aprendendo e parar # caso não esteja". try: # Última métrica SOM computada neste chunk (do último batch) last_som_metrics = som_metrics_log[-1] if som_metrics_log else {} # Reconstroi dict no formato esperado pelo adjuster adjust_metrics_input = { "quantization_error": float(last_som_metrics.get("quantization_error", 0.0)), "topological_error": float(last_som_metrics.get("topological_error", 0.0)), "kaski_lagus_error": float(last_som_metrics.get("kaski_lagus_error", 0.0)), "explained_variance_share": float(last_som_metrics.get("explained_variance_share", 0.0)), "topological_collapse": { "severity": last_som_metrics.get("collapse_severity", "none"), "collapsed_to_point": last_som_metrics.get("collapse_severity") == "point", "effective_rank": float(last_som_metrics.get("collapse_effective_rank", 0.0)), }, "dead_neuron_rate": { "dead_neuron_rate": float(last_som_metrics.get("dead_neuron_rate", 0.0)), "n_dead": int(last_som_metrics.get("n_dead_neurons", 0)), "n_active": int(last_som_metrics.get("n_active_neurons", 0)), "n_total": 864, }, "qe_stagnation": last_som_metrics.get("qe_stagnation", {"is_stagnant": False, "severity": "none"}), "neighborhood_crossing": last_som_metrics.get("neighborhood_crossing", {"detected": False, "severity": "none"}), "nan_detected": bool(last_som_metrics.get("nan_detected", False)), "overall_health": last_som_metrics.get("overall_health", "unknown"), "failure_indicators": last_som_metrics.get("failure_indicators", []), "n_failure_indicators": int(last_som_metrics.get("n_failure_indicators", 0)), } adjust_result = auto_adjuster.adjust( kls=kls, som_metrics=adjust_metrics_input, batch_idx=chunk_idx_global, ) if adjust_result.get("actions"): auto_adjust_log.append({ "step": step, "chunk_idx": chunk_idx_global, "phase": "PUNICAO", "actions": adjust_result["actions"], "boosts": adjust_result["boosts"], "n_actions": adjust_result["n_actions"], }) logger.info( f"[V6.5-V2-auto-adjust] PUNIÇÃO chunk {chunk_idx_global}: " f"actions={adjust_result['n_actions']}, " f"boosts(α,σ)=({adjust_result['boosts']['alpha0_boost']:.3f}, " f"{adjust_result['boosts']['sigma0_boost']:.3f})" ) # STOP-LEARNING detection (user: "parar caso não esteja") # V6.5-V2-auto-conscience-v2 — FASE2 também NÃO para mais por # stop_learning. Continua processando até atingir meta de 2000 # samples. O user requirement "COM PUNIÇÃO ATIVA" exige que # o protocolo V2 (train_hypotheses + apply_best_delta) seja # executado por todas as 2000 samples. if adjust_result.get("stop_training"): stop_training_reason = adjust_result.get("stop_reason") logger.warning( f"[V6.5-V2-auto-conscience-v2] FASE2 stop_learning " f"DETECTED at chunk {chunk_idx_global}: {stop_training_reason} " f"— CONTINUING punitive adjustments (não para)." ) try: auto_adjuster.state.reset_consecutive_counters() except Exception: pass # NÃO break — continua processando punição except Exception as adjust_err: logger.warning( f"[V6.5-V2-auto-adjust] PUNIÇÃO adjust() failed: {adjust_err}" ) # V7 — post-chunk hook (AdvancedSOMAugmenter.fit cross-validation) # A cada 2 chunks, executa fit() no buffer_4d para obter métricas # INDEPENDENTES (cross-validation com som_metrics.py). NÃO gera # dados sintéticos — apenas diagnostic. if v7_integrator is not None: try: v7_integrator.post_chunk_hook(kls, chunk_idx_global) except Exception as v7_pc_err: logger.debug(f"[V7] post_chunk_hook error (non-fatal): {v7_pc_err}") time.sleep(INTER_STREAM_BATCH_PAUSE_S_FASE2) # V2-dynamic-memory — Buffer sliding window na PUNIÇÃO também if len(kls.buffer_4d) > MAX_BUFFER_SIZE: overflow = len(kls.buffer_4d) - MAX_BUFFER_SIZE del kls.buffer_4d[:overflow] del kls.buffer_labels[:overflow] logger.info( f"[V6.5-V2] PUNIÇÃO buffer trimmed: -{overflow} → " f"{len(kls.buffer_4d)} samples" ) if chunk_idx_global % 2 == 0: aggressive_memory_cleanup() # V6.5-V2-metrics-FIX: também chama aggressive_cleanup do KLS try: kls.aggressive_cleanup() except Exception: pass # V6.5-V2-metrics-FIX-4: pausa pós-processamento maior na FASE2 # (user requirement: "FASE2 PUNIÇÃO é mais pesada") time.sleep(POST_PROCESSING_PAUSE_S_FASE2) except Exception as e: logger.error(f"[V6.5-V2] PUNIÇÃO dataset failed: {e}") traceback.print_exc() t_elapsed = time.time() - t_start logger.info(f"\n[V6.5-V2] FASE 2 (PUNIÇÃO) concluída:") logger.info(f" Total samples: {samples_processed}") logger.info(f" Meta atingida: {samples_processed >= META_MINIMA_PUNICAO}") logger.info(f" Tempo total: {t_elapsed:.1f}s") logger.info(f" Punishment events: {len(punishment_log)}") logger.info(f" Hypotheses trainings: {len(hypotheses_log)}") logger.info(f" Delta applications: {len(delta_applications_log)}") # V6.5-V2-metrics-FIX-4 — Sumário final da FASE2 # User requirement: "ao final da FASE2 mostrar evolução de métricas e dos # indicadores e da taxa de aprendizagem". fase2_summary = build_fase2_final_summary( kls=kls, som_metrics_log=som_metrics_log, punishment_log=punishment_log, hypotheses_log=hypotheses_log, delta_applications_log=delta_applications_log, elapsed_s=t_elapsed, ) logger.info("\n" + "=" * 80) logger.info("[V6.5-V2-metrics-FIX-4] FASE 2 — EVOLUÇÃO FINAL DE MÉTRICAS") logger.info("=" * 80) for line in fase2_summary["summary_lines"]: logger.info(line) logger.info("=" * 80 + "\n") # V7 — Finaliza integrator e consolida relatório V7 v7_report: Dict[str, Any] = {} if v7_integrator is not None: try: v7_report = v7_integrator.end_fase2() logger.info( f"\n[V7] FASE2 V7 Integration Report:\n" f" Batches processados: {v7_report.get('n_batches_processed', 0)}\n" f" Chunks processados: {v7_report.get('n_chunks_processed', 0)}\n" f" Punições observadas: {v7_report.get('n_punishments_observed', 0)}\n" f" β adaptativo final: {v7_report.get('final_adaptive_beta', 'N/A')}\n" f" Augmenter metrics count: {v7_report.get('augmenter_metrics_count', 0)}\n" f" DPO metrics count: {v7_report.get('dpo_metrics_count', 0)}\n" f" Elapsed: {v7_report.get('fase2_elapsed_s', 0):.1f}s" ) except Exception as v7_end_err: logger.warning(f"[V7] end_fase2 error (non-fatal): {v7_end_err}") v7_report = {"error": str(v7_end_err)[:200]} return { "phase": "TREINAMENTO_COM_PUNICAO", "dataset": PUNICAO_DATASET, "total_samples": int(samples_processed), "meta_minima": int(META_MINIMA_PUNICAO), "meta_atingida": bool(samples_processed >= META_MINIMA_PUNICAO), "storage_critical_stopped": bool(storage_critical_stopped), "elapsed_s": float(t_elapsed), "n_steps": int(step), "n_chunks": int(chunk_idx_global), "punishment_events": punishment_log, "hypotheses_trainings": hypotheses_log, "delta_applications": delta_applications_log, "accuracies_log": accuracies_log[-50:], # V6.5-V2-metrics — log de métricas SOM durante PUNIÇÃO "som_metrics_log": som_metrics_log[-100:], # últimos 100 batches "som_metrics_final": ( kls.compute_som_metrics() if kls.buffer_4d else {} ), "som_metric_history": kls.get_som_metric_history(), "final_v2_state": kls.get_v2_metrics(), # V6.5-V2-metrics-FIX-4 — sumário final da FASE2 "fase2_summary": fase2_summary, # V6.5-V2-auto-adjust — log de auto-ajustes aplicados na FASE2 "auto_adjust_log": auto_adjust_log[-100:], "auto_adjust_state": auto_adjuster.state.to_dict(), "stop_training_reason": stop_training_reason, # V7 — AdaptiveDPOLoss + AdvancedSOMAugmenter integration report "v7_report": v7_report, "v7_integrator_full": ( v7_integrator.to_report_dict() if v7_integrator is not None else {} ), } # ============================================================================ # 11. Evaluate attention (ativo e logicamente funcional) # ============================================================================ def evaluate_attention(kls: KohonenLearningSystemV2) -> Dict[str, Any]: """V6.5-V2 — Avalia se o mecanismo de atenção está ativo e funcional.""" metrics_snapshot = kls.get_attention_metrics() return { "evaluation": "attention_v65_v2", "user_requirement": "verificar se o mecanismo de atenção está ativo e acessado logicamente funcional", "metrics": metrics_snapshot, "active": bool(metrics_snapshot.get("active", False)), "logic_functional": bool(metrics_snapshot.get("logic_functional", False)), "n_calls": int(metrics_snapshot.get("n_calls", 0)), "n_errors": int(metrics_snapshot.get("n_errors", 0)), "n_heads": int(metrics_snapshot.get("n_heads", 0)), "assessment": "PASS" if metrics_snapshot.get("logic_functional") else "FAIL", } # ============================================================================ # 12. Evaluate predict fix (rótulos dinâmicos, não hardcoded) # ============================================================================ def evaluate_predict_fix(kls: KohonenLearningSystemV2) -> Dict[str, Any]: """V6.5-V2 — Avalia que predict() NÃO retorna mais 'gato'/'cachorro' fixos.""" test_queries = [ "o gato dorme na cama", "calcule dois mais dois", "qual é a capital do brasil", "explique o que é uma rede neural", "olá como você está", "traduza hello para portugues", ] results = [] n_returns_gato = 0 n_returns_cachorro = 0 n_returns_registry_label = 0 n_returns_default_label = 0 registry_state = kls.get_label_registry_state() all_registered_labels = set() for ds_name, mapping in registry_state["datasets"].items(): for label_str in mapping.values(): all_registered_labels.add(label_str) default_labels = set(registry_state["default_label_strings"].values()) for q in test_queries: try: pred = kls.predict(q) except Exception as e: pred = f"" try: pred_proba, prob = kls.predict_proba(q) except Exception as e: pred_proba, prob = f"", -1.0 is_gato = (pred == "gato") is_cachorro = (pred == "cachorro") is_registry = pred in all_registered_labels is_default = pred in default_labels if is_gato: n_returns_gato += 1 if is_cachorro: n_returns_cachorro += 1 if is_registry: n_returns_registry_label += 1 if is_default: n_returns_default_label += 1 results.append({ "query": q, "prediction": str(pred), "probability": float(prob) if isinstance(prob, (int, float)) else -1.0, "is_gato_hardcoded": bool(is_gato), "is_cachorro_hardcoded": bool(is_cachorro), "is_registry_label": bool(is_registry), "is_default_label": bool(is_default), }) # Per-dataset check per_dataset_results = [] for ds_name, mapping in registry_state["datasets"].items(): try: pred = kls.predict("teste de predição", dataset_name=ds_name) valid = pred in set(mapping.values()) except Exception as e: pred = f"" valid = False per_dataset_results.append({ "dataset": ds_name, "expected_labels": list(mapping.values()), "prediction": str(pred), "valid_for_dataset": bool(valid), }) return { "evaluation": "predict_fix_v65_v2", "user_requirement": "os strings 'gato' e 'cachorro' são fixos quando deveriam ser extrações variáveis e flexíveis de rótulos proveniente de dados dos datasets anteriormente treinados", "n_test_queries": len(test_queries), "results": results, "per_dataset_results": per_dataset_results, "summary": { "n_returns_gato": int(n_returns_gato), "n_returns_cachorro": int(n_returns_cachorro), "n_returns_registry_label": int(n_returns_registry_label), "n_returns_default_label": int(n_returns_default_label), "hardcoded_bug_present": bool(n_returns_gato > 0 or n_returns_cachorro > 0), "predict_fix_verified": bool(n_returns_gato == 0 and n_returns_cachorro == 0), }, "quality_assessment": { "fix_status": "PASS" if (n_returns_gato == 0 and n_returns_cachorro == 0) else "FAIL", "note": "predict() não retorna mais 'gato'/'cachorro' fixos — usa label_registry dinâmico.", }, } # ============================================================================ # 13. Launch user questions WITHOUT help (3 perguntas, sem auxílio) # ============================================================================ def launch_user_questions_without_help( kls: KohonenLearningSystemV2, ) -> Dict[str, Any]: """V6.5-V2 — Envia as 3 perguntas específicas SEM nenhum contexto. User requirement: "não ajudar o modelo em respostas e lançar perguntas: 'Luva de Pedreiro Távila' e 'Lula reserva valor' e 'Amazonas força-tarefa vítimas'" IMPORTANTE: As perguntas são enviadas EXATAMENTE como o usuário as escreveu, sem nenhum prompt adicional, sem contexto, sem instrução de formato, sem system prompt, sem few-shot examples. """ user_questions = [ "Luva de Pedreiro Távila", "Lula reserva valor", "Amazonas força-tarefa vítimas", ] results = [] for q in user_questions: t_start = time.time() try: som_pred = kls.predict(q) except Exception as e: som_pred = f"" try: reasoning_text = kls.reason_sync(q) except Exception as e: reasoning_text = f"" latency_ms = (time.time() - t_start) * 1000.0 think_match = re.search(r"(.*?)", reasoning_text, re.DOTALL) plan_match = re.search(r"(.*?)", reasoning_text, re.DOTALL) answer_match = re.search(r"(.*?)", reasoning_text, re.DOTALL) decompose_match = re.search(r"(.*?)", reasoning_text, re.DOTALL) results.append({ "query": q, "query_was_modified": False, "context_provided": False, "system_prompt_used": False, "few_shot_examples": False, "som_prediction": str(som_pred), "reasoning_length": int(len(reasoning_text)), "has_think": bool(think_match), "has_plan": bool(plan_match), "has_answer": bool(answer_match), "has_decompose": bool(decompose_match), "think_preview": ( think_match.group(1).strip()[:200] + ("..." if len(think_match.group(1)) > 200 else "") ) if think_match else "", "answer_preview": ( answer_match.group(1).strip()[:200] + ("..." if len(answer_match.group(1)) > 200 else "") ) if answer_match else "", "raw_response_preview": reasoning_text[:500] + ("..." if len(reasoning_text) > 500 else ""), "n_tags": sum(1 for tag in ["", "", "", ""] if tag in reasoning_text), "latency_ms": float(latency_ms), }) n_answer = sum(1 for r in results if r["has_answer"]) n_think = sum(1 for r in results if r["has_think"]) return { "evaluation": "user_questions_without_help_v65_v2", "user_requirement": "não ajudar o modelo em respostas e lançar perguntas", "questions_sent_verbatim": True, "no_context_added": True, "no_system_prompt": True, "no_few_shot": True, "n_questions": len(user_questions), "questions": user_questions, "results": results, "summary": { "n_with_answer": n_answer, "n_with_think": n_think, "answer_rate": float(n_answer / max(1, len(results))), "think_rate": float(n_think / max(1, len(results))), "avg_latency_ms": float(sum(r["latency_ms"] for r in results) / max(1, len(results))), "avg_reasoning_length": float(sum(r["reasoning_length"] for r in results) / max(1, len(results))), }, "quality_assessment": { "model_not_helped": True, "note": "As 3 perguntas foram enviadas verbatim, sem system prompt, sem few-shot, sem contexto adicional.", }, } # ============================================================================ # 13.5 HF Batch Upload — envia scripts+estado+modelo ao HF em lote # ============================================================================ HF_REPO_ID = "PowerMachine/BiGRU_T_version" def upload_to_hf_batch( hf_token: str, files_to_upload: List[Path], repo_id: str = HF_REPO_ID, repo_type: str = "model", ) -> Dict[str, Any]: """V6.5-V2-dynamic — Envia arquivos ao HF em LOTE (commit único). User requirement: "atualizações para o HF devem ser feitas em lote" + "após concluir enviar scripts e arquivos e estado do modelo testados e aprovados sobrescrevendo os antigos desatualizados". Args: hf_token: token HF (será apagado após uso). files_to_upload: lista de paths locais a enviar. repo_id: repositório HF destino. repo_type: 'model', 'dataset' ou 'space'. Returns: Dict com status de cada upload e URL do commit. """ if not hf_token: return {"uploaded": False, "reason": "no_hf_token"} try: from huggingface_hub import HfApi, whoami except ImportError as e: logger.error(f"[V6.5-V2] huggingface_hub not installed: {e}") return {"uploaded": False, "reason": f"import_error: {e}"} api = HfApi(token=hf_token) try: user_info = whoami(token=hf_token) logger.info(f"[V6.5-V2] HF authenticated as: {user_info.get('name', 'unknown')}") except Exception as e: logger.error(f"[V6.5-V2] HF auth failed: {e}") return {"uploaded": False, "reason": f"auth_failed: {e}"} # Garante que o repo existe try: api.repo_info(repo_id=repo_id, repo_type=repo_type) except Exception: try: api.create_repo(repo_id=repo_id, repo_type=repo_type, exist_ok=True) logger.info(f"[V6.5-V2] Created HF repo: {repo_id}") except Exception as e: logger.warning(f"[V6.5-V2] Could not create repo {repo_id}: {e}") # Filtra arquivos existentes upload_plan: List[Dict[str, Any]] = [] for src in files_to_upload: if not src.exists(): logger.warning(f"[V6.5-V2] Skipping non-existent file: {src}") continue # V6.5-V2-auto-conscience-v2 — PRESERVA estrutura de diretórios no HF. # path_in_repo = path relativo a BIGRU_ROOT (ex: src/bigru_t/model/kohonen_learning_system.py) try: path_in_repo = str(src.relative_to(BIGRU_ROOT)).replace("\\", "/") except ValueError: # Fallback: filename only (for files outside BIGRU_ROOT) path_in_repo = src.name upload_plan.append({ "local": str(src), "path_in_repo": path_in_repo, "size_mb": float(src.stat().st_size / 1e6), }) if not upload_plan: return {"uploaded": False, "reason": "no_files_to_upload"} logger.info( f"[V6.5-V2] Uploading {len(upload_plan)} files to HF " f"repo={repo_id} in BATCH (single commit):" ) for item in upload_plan: logger.info(f" - {item['local']} → {item['path_in_repo']} " f"({item['size_mb']:.2f}MB)") # Upload em lote (commit único): usa upload_folder com todos os arquivos # numa pasta temporária — ou, mais simples, upload_file individual mas # em sequência rápida. Para commit único verdadeiro, usaremos create_commit. results: List[Dict[str, Any]] = [] try: # Método 1: create_commit com múltiplos operations (commit único) from huggingface_hub import CommitOperationAdd operations = [] for item in upload_plan: operations.append(CommitOperationAdd( path_in_repo=item["path_in_repo"], path_or_fileobj=item["local"], )) commit_info = api.create_commit( repo_id=repo_id, repo_type=repo_type, operations=operations, commit_message=( f"V6.5-V2-dynamic: upload batch (scripts + state + model) " f"[{len(operations)} files]" ), ) logger.info(f"[V6.5-V2] Batch commit created: {commit_info}") for item in upload_plan: results.append({ "file": item["path_in_repo"], "uploaded": True, "size_mb": item["size_mb"], }) return { "uploaded": True, "n_files": len(results), "commit_url": str(commit_info.commit_url) if hasattr(commit_info, "commit_url") else str(commit_info), "results": results, } except Exception as e: logger.error(f"[V6.5-V2] Batch commit failed: {e}") traceback.print_exc() # Método 2 (fallback): upload individual logger.info("[V6.5-V2] Falling back to individual uploads...") for item in upload_plan: try: api.upload_file( path_or_fileobj=item["local"], path_in_repo=item["path_in_repo"], repo_id=repo_id, repo_type=repo_type, ) results.append({ "file": item["path_in_repo"], "uploaded": True, "size_mb": item["size_mb"], }) logger.info(f" ✓ {item['path_in_repo']} uploaded") except Exception as e: results.append({ "file": item["path_in_repo"], "uploaded": False, "error": str(e)[:200], }) logger.error(f" ✗ {item['path_in_repo']} failed: {e}") return { "uploaded": any(r.get("uploaded") for r in results), "n_files": len(results), "n_uploaded": sum(1 for r in results if r.get("uploaded")), "results": results, } def scrub_hf_token_from_scripts(scripts_dir: Path) -> Dict[str, Any]: """V6.5-V2-dynamic — Apaga HF_TOKEN de todos os scripts no diretório. User requirement: "HF_TOKEN deve ser apagada dos scripts enviados ao HF após uso". Usa regex para detectar qualquer string que pareça um token HF (hf_{32,}) e substitui por ''. """ pattern = re.compile(r'hf_[A-Za-z0-9]{32,}') n_files_scanned = 0 n_files_modified = 0 n_tokens_removed = 0 for py_file in scripts_dir.rglob("*.py"): n_files_scanned += 1 try: content = py_file.read_text(encoding="utf-8") matches = pattern.findall(content) if matches: new_content = pattern.sub("", content) py_file.write_text(new_content, encoding="utf-8") n_files_modified += 1 n_tokens_removed += len(matches) logger.info(f"[V6.5-V2] Scrubbed {len(matches)} HF token(s) from {py_file}") except Exception as e: logger.warning(f"[V6.5-V2] Could not scan {py_file}: {e}") return { "n_files_scanned": n_files_scanned, "n_files_modified": n_files_modified, "n_tokens_removed": n_tokens_removed, } def main() -> int: import torch n_neurons = SOM_GRID[0] * SOM_GRID[1] * SOM_GRID[2] * SOM_GRID[3] print("\n" + "=" * 80) print("V6.5-V2 — 16 HIPÓTESES + 3 TENTATIVAS + 2 FASES (CONHECIMENTO → PUNIÇÃO)") print("=" * 80) print(f" BATCH_SIZE : {BATCH_SIZE}") print(f" STREAM_BATCH_SIZE : {STREAM_BATCH_SIZE} (user: 'de 100 em 100')") print(f" CONHECIMENTO datasets: {len(CONHECIMENTO_DATASETS)}") print(f" PUNICAO dataset : {PUNICAO_DATASET}") print(f" Max samples CONHEC : {MAX_SAMPLES_PER_DATASET_CONHECIMENTO}/dataset") print(f" Max samples PUNICAO : {MAX_SAMPLES_PUNICAO}") print(f" Meta CONHECIMENTO : ≥{META_MINIMA_CONHECIMENTO}") print(f" Meta PUNICAO : ≥{META_MINIMA_PUNICAO}") print(f" SOM grid : {SOM_GRID} ({n_neurons} neurons)") print(f" HIDDEN_DIM : {HIDDEN_DIM}") print(f" VOCAB_SIZE : {VOCAB_SIZE}") print(f" n_hypotheses : {N_HYPOTHESES}") print(f" n_trials : {N_TRIALS}") print(f" hyp_train_steps : {HYP_TRAIN_STEPS}") print(f" Attention : MultiHeadAttention (n_heads=8)") print(f" V65_ENABLE_STREAMING : {os.environ.get('V65_ENABLE_STREAMING')} (REAL)") print(f" Storage critical : {STORAGE_CRITICAL_PCT}%") print(f" Xeon cores : {N_CORES}") print(f" FP16 best TFLOPS : {FP16_BENCH.get('best_tflops', 0.0):.3f}") print("=" * 80 + "\n") # 14.1 Verify V2 logic logger.info("[V6.5-V2] Verifying V2 logic (predict + 16 hypotheses + V2 methods)...") verification = verify_v2_logic() print(f"\n--- V2 Logic Verification ---") print(f" PASS: {verification['n_pass']}/{verification['n_pass'] + verification['n_fail']}") for check_name, check_info in verification["checks"].items(): status = check_info.get("status", "?") details = check_info.get("details", check_info.get("error", "")) print(f" [{status}] {check_name}: {details[:100]}") # 14.2 Initialize KohonenLearningSystemV2 (864 neurons + 16 hypotheses + attention) kls = KohonenLearningSystemV2( vocab_size=VOCAB_SIZE, hidden_dim=HIDDEN_DIM, seq_len=MAX_SEQ_LEN, som_grid=SOM_GRID, alpha0=ALPHA0, sigma0=SIGMA0, lambda_ewc=LAMBDA_EWC, N_start=N_START, dim_choice=DIM_CHOICE, hypothesis_hidden=[512, 256, 128, 64, 32, 16, 8], T_max=T_MAX, # V6.5-V3-no-regression: VQ-VAE-2 REATIVADO conforme user requirement. # User: "não autorizei a desativação do 'VQ-VAE-2' de agora em diante # pergunte se é para reduzir ou desativar qualquer componente indicando # as estratégias". VQ-VAE-2 foi desativado em V6.5-V2-buffer-864 sem # autorização — RESTAURADO agora. # OOM-safety garantida por: # (a) _compress_buffer_with_vqvae2 agora é LAZY (a cada 16 add_data, # não a cada chamada) — ver kohonen_learning_system.py # (b) torch.no_grad() em toda compressão (já existente) # (c) Sanitização NaN/Inf (já existente) # (d) Fallback buffer=256 se RSS > 75% cgroup (check_memory_and_maybe_fallback) # (e) OomGuard thread daemon monitora RSS a cada 2s # VQ-VAE-2 comprime o buffer 4D via encoder hierárquico + codebook EMA # + Goose VQ, fornecendo representação compacta do estado do SOM. enable_vqvae2=True, # V6.5-V3-no-regression: REATIVADO (era False sem autorização) # V6.5-V2-metrics-FIX-4: reasoning_engine desabilitado para mitigar # OOM-killer (User requirement implícito: "investigar e corrigir # falhas de lógica e bugs que estejam causando alto consumo de memória # sem distorcer a arquitetura Kohonen"). O ReasoningEngine cria um # ThreadPoolExecutor(4 workers) + ToolAgentCoordinator que consome # ~100-200MB adicionais. Como o reasoning_engine é OPCIONAL e não # afeta o aprendizado do SOM (apenas gera tags // # para predições), desabilitá-lo preserva a arquitetura Kohonen e # libera memória para o streaming de 8000+2000 samples. # Para reativar: mudar para True (requer ≥6GB cgroup). enable_reasoning=False, enable_w8a8=False, # W8A8 permanece desabilitado (não essencial para Kohonen) vqvae2_code_dim=16, vqvae2_num_codes_top=64, vqvae2_num_codes_bot=128, w8a8_alpha=0.5, w8a8_n_bits=8, w8a8_calibration_samples=32, enable_attention=True, attention_n_heads=8, # V2-dynamic parameters (serão adaptados automaticamente) n_hypotheses=N_HYPOTHESES, max_n_hypotheses=MAX_N_HYPOTHESES, min_n_hypotheses=MIN_N_HYPOTHESES, n_trials=N_TRIALS, min_n_trials=MIN_N_TRIALS, max_n_trials=MAX_N_TRIALS, hyp_train_steps=HYP_TRAIN_STEPS, min_hyp_train_steps=MIN_HYP_TRAIN_STEPS, max_hyp_train_steps=MAX_HYP_TRAIN_STEPS, hyp_lr=HYP_LR, hyp_hidden_dim=HYP_HIDDEN_DIM, loss_history_window=LOSS_HISTORY_WINDOW, punishment_window=PUNISHMENT_WINDOW, ) logger.info( f"[V6.5-V2-dynamic] KLS V2 initialized: {n_neurons} neurons, " f"hidden={HIDDEN_DIM}, vocab={VOCAB_SIZE}, " f"n_hypotheses={kls.n_hypotheses}/{kls.max_n_hypotheses} " f"(active={kls.hypothesis_ensemble.active_count}), " f"n_trials={kls.n_trials} (limits {kls.min_n_trials}-{kls.max_n_trials}), " f"hyp_train_steps={kls.hyp_train_steps} " f"(limits {kls.min_hyp_train_steps}-{kls.max_hyp_train_steps}), " f"delta_scale={kls.delta_scale.item():.4f}, " f"attention={kls.enable_attention} (n_heads={kls.attention_n_heads}), " f"workers_active={getattr(kls, '_tool_coordinator_workers_active', False)}" ) # Treina tokenizer com corpus básico PT-BR corpus_inicial = [ "o gato dorme na cama", "a casa eh azul", "ele corre rapido", "ela canta uma musica", "o sol nasceu hoje", "nos vamos viajar", "o livro esta na mesa", "a menina brinca no parque", "ola como voce esta", "qual e o seu nome", "calcule dois mais dois", "traduza hello para portugues", "instrucao para resolver o problema", "resposta para a pergunta", "luva de pedreiro tavila", "lula reserva valor", "amazonas forca tarefa vitimas", ] kls.tokenizer.fit(corpus_inicial) logger.info(f"[V6.5-V2] Tokenizer fitted with {len(corpus_inicial)} corpus words") # HF_TOKEN hf_token = os.environ.get("HF_TOKEN") if not hf_token: logger.warning("[V6.5-V2] HF_TOKEN not set — streaming may fail for gated datasets") else: logger.info(f"[V6.5-V2] HF_TOKEN set: {hf_token[:8]}...") # 14.3 Aggressive storage cleanup before training starts logger.info("[V6.5-V2] Aggressive storage cleanup (pre-training)...") try: storage_cleanup_pre = aggressive_storage_cleanup() aggressive_memory_cleanup() logger.info( f"[V6.5-V2] Pre-training cleanup: " f"{storage_cleanup_pre['n_files_removed']} files, " f"{storage_cleanup_pre['mb_freed']:.1f}MB freed" ) except Exception as e: logger.warning(f"[V6.5-V2] Pre-training cleanup failed: {e}") aggressive_memory_cleanup() # ======================================================================== # 14.4 FASE 1 — CONHECIMENTO (8 datasets, sem punição) # ======================================================================== # V7-resume: suporte a SKIP_FASE1 para retomar FASE2 do estado salvo. # User requirement: "FASE1 e FASE2 estão treinadas no mesmo estado do modelo". # Se o estado da FASE1 já existe (reason=end_of_fase_conhecimento) e # total_samples >= META_MINIMA_CONHECIMENTO, pula FASE1 e vai direto pra FASE2. # Isto permite retomar FASE2 sem refazer FASE1 (que leva ~8 minutos). SKIP_FASE1 = os.environ.get("SKIP_FASE1", "0") == "1" fase1_state_exists = MODEL_STATES_PATH.exists() fase1_state_valid = False if fase1_state_exists: try: import torch as _torch _state = _torch.load(MODEL_STATES_PATH, map_location="cpu", weights_only=False) _meta = _state.get("_meta", {}) if isinstance(_state, dict) else {} _reason = str(_meta.get("reason", "")) _total = int(_meta.get("total_samples", 0)) if "end_of_fase_conhecimento" in _reason and _total >= META_MINIMA_CONHECIMENTO: fase1_state_valid = True logger.info( f"[V7-resume] FASE1 state found: reason={_reason}, " f"total_samples={_total}. SKIP_FASE1={SKIP_FASE1}." ) del _state # libera memória imediatamente except Exception as _e: logger.warning(f"[V7-resume] Failed to check FASE1 state: {_e}") if SKIP_FASE1 and fase1_state_valid: logger.info( "\n[V7-resume] *** SKIP FASE 1 — usando estado salvo ***\n" f" State file: {MODEL_STATES_PATH}\n" f" FASE2 vai rodar diretamente com o estado da FASE1." ) conhecimento_result = { "total_samples": int(_total) if fase1_state_valid else 0, "meta_minima": META_MINIMA_CONHECIMENTO, "meta_atingida": True, "n_steps": 0, "n_chunks": 0, "elapsed_s": 0.0, "accuracies_log": [], "attention_log": [], "som_metrics_log": [], "auto_adjust_log": [], "samples_per_dataset": {}, "streaming_failures": {}, "storage_critical_stopped": False, "stop_training_reason": None, "skipped": True, "skipped_reason": "SKIP_FASE1=1 with valid FASE1 state", } else: logger.info("\n[V6.5-V2] Iniciando FASE 1 — CONHECIMENTO...") conhecimento_result = run_fase_conhecimento(kls, hf_token) V2_PHASES_EVAL_PATH.write_text( json.dumps({"fase_1_conhecimento": conhecimento_result}, indent=2, ensure_ascii=False) ) logger.info( f"\n[V6.5-V2] FASE 1 — CONHECIMENTO finalizada: " f"{conhecimento_result['total_samples']} samples " f"(meta: {conhecimento_result['meta_minima']}, " f"atingida: {conhecimento_result['meta_atingida']})" ) # V6.5-V2-unified — Salva estado após CONHECIMENTO no MESMO arquivo save_model_states_for_evaluation( kls=kls, reason="end_of_fase_conhecimento", step=conhecimento_result["n_steps"], total_samples=conhecimento_result["total_samples"], output_path=MODEL_STATES_PATH, # ARQUIVO UNIFICADO ) # Aggressive cleanup entre fases aggressive_memory_cleanup() aggressive_storage_cleanup() # ======================================================================== # 14.5 FASE 2 — TREINAMENTO COM PUNIÇÃO (BrunoN-Dev/corpus-ptbr-v1, 2000 samples) # ======================================================================== # V6.5-V2-auto-conscience-v2 — FASE2 agora roda MESMO se FASE1 parou cedo # por stop-learning detection. Justificativa: # - User requirement: "FASE2 TREINAMENTO (...) COM PUNIÇÃO ATIVA" # - FASE2 aplica train_hypotheses + apply_best_delta que ajustam os pesos # do SOM — exatamente o que pode quebrar o estado de neurônios mortos. # - Skipar FASE2 deixaria o modelo sem ajuste punitivo, piorando o problema. # Apenas skip se storage_critical_disk (sem espaço para salvar estado). skip_fase2_reason = None if conhecimento_result.get("storage_critical_stopped", False): # Verifica se foi stop_learning (FASE2 deve rodar) ou storage critical (skip) stop_reason = conhecimento_result.get("stop_training_reason", "") if stop_reason and "stop_learning" not in str(stop_reason).lower() and "storage" in str(stop_reason).lower(): skip_fase2_reason = f"storage_critical_during_fase1: {stop_reason}" else: logger.warning( f"[V6.5-V2-auto-conscience-v2] FASE1 parou por stop_learning " f"(reason={stop_reason}). FASE2 ainda vai rodar — punitive " f"adjustments podem reviver neurônios mortos." ) if not skip_fase2_reason: # V6.5-V2-metrics — Carrega explicitamente o estado do modelo concluído # da FASE1 no KLS antes de iniciar a FASE2. # User requirement: "esta FASE2 (nunca acontece antes da FASE1) deve # usar do arquivo de Estado do MODELO concluído da FASE1". # Como usamos a mesma instância KLS, o estado em memória já é o da # FASE1, mas carregamos do disco explicitamente para: # (a) validar que o arquivo .pt está íntegro; # (b) permitir que um processo separado (apenas FASE2) seja iniciado # do zero apontando para o estado salvo da FASE1. logger.info("\n[V6.5-V2-metrics] Carregando estado do modelo da FASE1 (FASE2 uses FASE1 completed state)...") # V6.5-V2-unified — Carrega estado unificado (FASE1+FASE2 no mesmo arquivo). # User requirement: "LEMBRANDO que agora FASE1 e FASE2 estão treinadas # no mesmo estado do modelo". fase1_state_path = MODEL_STATES_PATH fase1_load_result = load_fase1_state_into_kls(kls, fase1_state_path) logger.info("\n[V6.5-V2] Iniciando FASE 2 — TREINAMENTO COM PUNIÇÃO...") punicão_result = run_fase_punicão(kls, hf_token) # Anexa o resultado do carregamento da FASE1 no relatório da FASE2 punicão_result["fase1_state_load"] = fase1_load_result # Atualiza V2_PHASES_EVAL_PATH com ambas as fases phases_eval = { "fase_1_conhecimento": conhecimento_result, "fase_2_treinamento_com_punicão": punicão_result, } V2_PHASES_EVAL_PATH.write_text( json.dumps(phases_eval, indent=2, ensure_ascii=False) ) logger.info( f"\n[V6.5-V2] FASE 2 — PUNIÇÃO finalizada: " f"{punicão_result['total_samples']} samples " f"(meta: {punicão_result['meta_minima']}, " f"atingida: {punicão_result['meta_atingida']})" ) logger.info( f"[V6.5-V2] Punishment events: {len(punicão_result['punishment_events'])}, " f"Hypotheses trainings: {len(punicão_result['hypotheses_trainings'])}, " f"Delta applications: {len(punicão_result['delta_applications'])}" ) else: punicão_result = { "phase": "TREINAMENTO_COM_PUNICAO", "skipped": True, "reason": skip_fase2_reason or "storage_critical_stopped_during_conhecimento", } logger.warning(f"[V6.5-V2-auto-conscience-v2] FASE 2 pulada — {punicão_result['reason']}") # ======================================================================== # 14.6 Salvar estados finais do modelo # ======================================================================== logger.info("\n[V6.5-V2] Salvando estados finais do modelo...") total_samples_final = ( conhecimento_result["total_samples"] + punicão_result.get("total_samples", 0) ) save_info = save_model_states_for_evaluation( kls=kls, reason="end_of_training_v2", step=conhecimento_result["n_steps"] + punicão_result.get("n_steps", 0), total_samples=total_samples_final, ) # ======================================================================== # 14.7 Avaliações finais (attention, predict fix, user questions) # ======================================================================== logger.info("\n[V6.5-V2] Avaliando attention...") attention_eval = evaluate_attention(kls) ATTENTION_EVAL_PATH.write_text(json.dumps(attention_eval, indent=2, ensure_ascii=False)) logger.info("[V6.5-V2] Avaliando predict fix (rótulos dinâmicos)...") predict_fix_eval = evaluate_predict_fix(kls) PREDICT_FIX_EVAL_PATH.write_text(json.dumps(predict_fix_eval, indent=2, ensure_ascii=False)) logger.info("[V6.5-V2] Lançando 3 perguntas SEM AJUDA ao modelo...") user_questions_eval = launch_user_questions_without_help(kls) USER_QUESTIONS_PATH.write_text(json.dumps(user_questions_eval, indent=2, ensure_ascii=False)) # ======================================================================== # 14.8 Relatório final # ======================================================================== final_report = { "version": "V6.5-V2", "timestamp": datetime.now().isoformat(), "config": { "som_grid": list(SOM_GRID), "n_neurons": n_neurons, "hidden_dim": HIDDEN_DIM, "vocab_size": VOCAB_SIZE, "n_hypotheses": N_HYPOTHESES, "n_trials": N_TRIALS, "hyp_train_steps": HYP_TRAIN_STEPS, "stream_batch_size": STREAM_BATCH_SIZE, "meta_conhecimento": META_MINIMA_CONHECIMENTO, "meta_punicão": META_MINIMA_PUNICAO, "conhecimento_datasets": CONHECIMENTO_DATASETS, "punicão_dataset": PUNICAO_DATASET, }, "xeon_status": XEON_STATUS, "fp16_benchmark": FP16_BENCH, "v2_verification": verification, "fase_1_conhecimento_summary": { "total_samples": conhecimento_result["total_samples"], "meta_atingida": conhecimento_result["meta_atingida"], "elapsed_s": conhecimento_result["elapsed_s"], "storage_critical_stopped": conhecimento_result["storage_critical_stopped"], }, "fase_2_punicão_summary": { "total_samples": punicão_result.get("total_samples", 0), "meta_atingida": punicão_result.get("meta_atingida", False), "elapsed_s": punicão_result.get("elapsed_s", 0.0), "punishment_events": len(punicão_result.get("punishment_events", [])), "hypotheses_trainings": len(punicão_result.get("hypotheses_trainings", [])), "delta_applications": len(punicão_result.get("delta_applications", [])), "skipped": punicão_result.get("skipped", False), }, "attention_eval_summary": { "active": attention_eval["active"], "logic_functional": attention_eval["logic_functional"], "n_calls": attention_eval["n_calls"], "assessment": attention_eval["assessment"], }, "predict_fix_summary": predict_fix_eval["summary"], "user_questions_summary": user_questions_eval["summary"], "model_states_saved_to": str(MODEL_STATES_PATH), "save_info": save_info, "final_v2_metrics": kls.get_v2_metrics(), } REPORT_PATH.write_text(json.dumps(final_report, indent=2, ensure_ascii=False)) logger.info(f"\n[V6.5-V2] Relatório final salvo em: {REPORT_PATH}") # Copia arquivos relevantes para download/ try: # V6.5-V2-unified — Apenas MODEL_STATES_PATH (arquivo unificado) for src in [REPORT_PATH, USER_QUESTIONS_PATH, MODEL_STATES_PATH, V2_PHASES_EVAL_PATH, PREDICT_FIX_EVAL_PATH, ATTENTION_EVAL_PATH]: if src.exists(): dst = DOWNLOAD_DIR / src.name shutil.copy2(src, dst) logger.info(f"[V6.5-V2] Copied {src.name} → {dst}") except Exception as e: logger.warning(f"[V6.5-V2] Failed to copy files to download/: {e}") # ======================================================================== # 14.9 HF BATCH UPLOAD — enviar scripts + estado + modelo ao HF em lote # ======================================================================== # User requirement: "após concluir enviar scripts e arquivos e estado # (e treinamento para permitir continuar novo treinamento) do modelo # testados e aprovados sobrescrevendo os antigos desatualizados" upload_report_path = BIGRU_ROOT / "v6_5_v2_upload_report.json" hf_upload_result: Dict[str, Any] = {"uploaded": False, "reason": "skipped"} if hf_token: logger.info("\n[V6.5-V2] Coletando arquivos para upload HF em LOTE...") # V6.5-V2-auto-conscience-v2 — Coleta TODOS os arquivos .py sob src/ e # scripts/ (exceto deprecados/), preservando estrutura de diretórios # no HF (path_in_repo = path relativo a BIGRU_ROOT). # User requirement (latest): "ao concluir enviar para o HF os arquivos # de Estado do modelo e módulos python e scripts em lote" + # "estado (e treinamento para permitir continuar novo treinamento)". # # Estratégia: # (a) Módulos .py + scripts .py sob src/ e scripts/ — preserva dirs # (b) Model state .pt — UPLOAD ao HF (user quer continuar treino) # Verifica tamanho ≤ 500MB antes de enviar (HF LFS free tier) # (c) Relatórios JSON — para auditoria all_files: List[Path] = [] _EXCLUDE_DIRS = {"deprecados", "__pycache__", ".git", ".pytest_cache", "model_final", "docs", "node_modules", ".venv", "venv"} _INCLUDE_EXTS = {".py", ".md", ".txt", ".json", ".sh"} # (a) Walk src/ e scripts/ preservando paths relativos for root_dir in [SRC_ROOT, BIGRU_ROOT / "scripts"]: if not root_dir.exists(): continue for path in root_dir.rglob("*"): if not path.is_file(): continue try: rel = path.relative_to(BIGRU_ROOT) except ValueError: continue if any(part in _EXCLUDE_DIRS for part in rel.parts): continue if path.suffix.lower() not in _INCLUDE_EXTS: continue all_files.append(path) # (b) Model state .pt files (estado para continuar treino) # User requirement: "estado (e treinamento para permitir continuar # novo treinamento) do modelo testados e aprovados". # Inclui: v6_5_v2_model_states.pt, v6_5_v2_model_states_after_conhecimento.pt # Verifica tamanho ≤ 500MB (HF LFS free tier sem autenticar LFS). _MAX_HF_LFS_SIZE_BYTES = 500 * 1024 * 1024 # 500MB # V6.5-V2-unified — Apenas UM arquivo de estado unificado. for pt_candidate in [ MODEL_STATES_PATH, ]: if pt_candidate.exists(): pt_size = pt_candidate.stat().st_size if pt_size <= _MAX_HF_LFS_SIZE_BYTES: all_files.append(pt_candidate) logger.info(f"[V6.5-V2] Incluindo estado do modelo no upload HF: {pt_candidate.name} ({pt_size/1e6:.1f}MB)") else: logger.warning(f"[V6.5-V2] Estado muito grande para HF LFS (>{_MAX_HF_LFS_SIZE_BYTES/1e6:.0f}MB): {pt_candidate.name} ({pt_size/1e6:.1f}MB) — será mantido localmente apenas") # (c) Relatórios JSON for rpt in [REPORT_PATH, V2_PHASES_EVAL_PATH, PREDICT_FIX_EVAL_PATH, ATTENTION_EVAL_PATH, USER_QUESTIONS_PATH]: if rpt.exists(): all_files.append(rpt) # Dedup seen = set() deduped_files: List[Path] = [] for p in all_files: if str(p) not in seen: seen.add(str(p)) deduped_files.append(p) all_files = deduped_files # Apaga HF_TOKEN dos scripts ANTES do upload (não enviar tokens) scrub_result = scrub_hf_token_from_scripts(BIGRU_ROOT / "scripts") scrub_result_src = scrub_hf_token_from_scripts(SRC_ROOT) logger.info( f"[V6.5-V2] Token scrub: scripts={scrub_result}, src={scrub_result_src}" ) # Upload em LOTE hf_upload_result = upload_to_hf_batch( hf_token=hf_token, files_to_upload=all_files, repo_id=HF_REPO_ID, repo_type="model", ) upload_report_path.write_text( json.dumps(hf_upload_result, indent=2, ensure_ascii=False) ) logger.info(f"[V6.5-V2] HF upload report saved: {upload_report_path}") if upload_report_path.exists(): shutil.copy2(upload_report_path, DOWNLOAD_DIR / upload_report_path.name) else: logger.warning("[V6.5-V2] HF_TOKEN não disponível — pulando upload HF.") # ======================================================================== # 14.10 LIMPEZA HF_TOKEN do ambiente (user requirement) # ======================================================================== logger.info("\n[V6.5-V2] Limpando HF_TOKEN do ambiente...") if "HF_TOKEN" in os.environ: del os.environ["HF_TOKEN"] logger.info("[V6.5-V2] HF_TOKEN removido do ambiente.") if "HUGGING_FACE_HUB_TOKEN" in os.environ: del os.environ["HUGGING_FACE_HUB_TOKEN"] logger.info("[V6.5-V2] HUGGING_FACE_HUB_TOKEN removido do ambiente.") # Re-scrub após o upload (segurança extra) final_scrub = scrub_hf_token_from_scripts(BIGRU_ROOT / "scripts") final_scrub_src = scrub_hf_token_from_scripts(SRC_ROOT) final_scrub_proj = scrub_hf_token_from_scripts(PROJECT_ROOT / "scripts") # ======================================================================== # 14.11 LIMPEZA final de armazenamento e worklog # ======================================================================== # User requirement: "limpar Armazenamento e worklog (remover imediatamente # todos os arquivos *.pt e não armazenar mais arquivos no worklog)" logger.info("\n[V6.5-V2] Limpeza final de armazenamento...") final_storage_cleanup = aggressive_storage_cleanup() aggressive_memory_cleanup() # V6.5-V2-auto-conscience-v2 — Limpeza de .pt files: SÓ remove se o upload # HF foi bem-sucedido. User requirement (latest): # "estado (e treinamento para permitir continuar novo treinamento)" # Se HF upload falhou, MANTÉM os .pt locais para o usuário poder recuperá-los. # Se HF upload OK, remove apenas os .pt do diretório do projeto (exceto # download/, que é a área de entrega para o usuário). pt_files_removed: List[str] = [] if hf_upload_result.get("uploaded", False): logger.info("[V6.5-V2-auto-conscience-v2] HF upload OK — removendo .pt locais do projeto (estado já está no HF)...") for pt_file in BIGRU_ROOT.rglob("*.pt"): try: rel_path = str(pt_file.relative_to(PROJECT_ROOT)) pt_file.unlink() pt_files_removed.append(rel_path) logger.info(f" removed .pt: {rel_path}") except Exception as e: logger.warning(f" failed to remove {pt_file}: {e}") logger.info( f"[V6.5-V2-auto-conscience-v2] {len(pt_files_removed)} .pt files removed from project." ) else: logger.warning( "[V6.5-V2-auto-conscience-v2] HF upload falhou — MANTENDO .pt locais " "para preservar estado do modelo. User pode recuperar manualmente." ) # Copia .pt para download/ para que o usuário tenha acesso for pt_file in BIGRU_ROOT.rglob("*.pt"): try: dst = DOWNLOAD_DIR / pt_file.name shutil.copy2(pt_file, dst) logger.info(f" .pt copiado para download/: {pt_file.name}") except Exception as e: logger.warning(f" failed to copy {pt_file}: {e}") # Atualiza worklog.md (sem referenciar .pt files) try: worklog_path = PROJECT_ROOT / "worklog.md" worklog_content = f"""# worklog.md — V6.5-V2-metrics-FIX Updated: {datetime.now().isoformat()} ## Status - FASE 1 CONHECIMENTO: {conhecimento_result['total_samples']} samples (meta: {conhecimento_result['meta_atingida']}) - FASE 2 PUNIÇÃO: {punicão_result.get('total_samples', 0)} samples (meta: {punicão_result.get('meta_atingida', False)}) - HF upload: {hf_upload_result.get('uploaded', False)} ({hf_upload_result.get('n_files', 0)} files) - HF_TOKEN: cleaned from env + scripts - .pt files: ALL removed from local storage ({len(pt_files_removed)} files) - Model state available in HF repo: {HF_REPO_ID} """ worklog_path.write_text(worklog_content, encoding="utf-8") logger.info(f"[V6.5-V2-metrics-FIX] worklog.md updated (no .pt references)") except Exception as e: logger.warning(f"[V6.5-V2-metrics-FIX] Failed to update worklog: {e}") print("\n" + "=" * 80) print("V6.5-V2-metrics-FIX — TREINAMENTO CONCLUÍDO") print("=" * 80) print(f" FASE 1 CONHECIMENTO : {conhecimento_result['total_samples']} samples " f"(meta: {conhecimento_result['meta_atingida']})") print(f" FASE 2 PUNICAO : {punicão_result.get('total_samples', 0)} samples " f"(meta: {punicão_result.get('meta_atingida', False)})") print(f" Punishment events : {len(punicão_result.get('punishment_events', []))}") print(f" Delta applications : {len(punicão_result.get('delta_applications', []))}") print(f" Attention functional: {attention_eval['logic_functional']}") print(f" Predict fix verified: {predict_fix_eval['summary']['predict_fix_verified']}") print(f" User questions ans : {user_questions_eval['summary']['n_with_answer']}/3") print(f" HF upload : {hf_upload_result.get('uploaded', False)} " f"({hf_upload_result.get('n_files', 0)} files)") print(f" HF_TOKEN cleaned : True (env + scripts)") print(f" .pt files removed : {len(pt_files_removed)} (ALL local .pt cleaned)") print(f" Unified state file : {MODEL_STATES_PATH.name} (FASE1+FASE2 continuous)") print(f" k-means++ applied : {getattr(kls.som, '_kmeans_pp_initialized', False)}") print(f" Storage cleanup : {final_storage_cleanup['n_files_removed']} files, " f"{final_storage_cleanup['mb_freed']:.1f}MB freed") print(f" Final n_hypotheses : {kls.n_hypotheses} (active={kls.hypothesis_ensemble.active_count})") print(f" Final n_trials : {kls.n_trials}") print(f" Final hyp_steps : {kls.hyp_train_steps}") print(f" N adaptations : {len(kls.get_adaptation_log())}") print(f" HIDDEN_DIM : {HIDDEN_DIM} (restored)") print(f" VOCAB_SIZE : {VOCAB_SIZE} (restored)") print("=" * 80 + "\n") return 0 if __name__ == "__main__": sys.exit(main())