BiGRU_T_version / scripts /train_v6_5_v2.py
PowerMachine's picture
V6.7: upload scripts/train_v6_5_v2.py (som_auto_adjust_runner + train script V6.7)
acb0bc3 verified
Raw History Blame Contribute Delete
183 kB
"""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"<error: {e}>"
try:
pred_proba, prob = kls.predict_proba(q)
except Exception as e:
pred_proba, prob = f"<error: {e}>", -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"<error: {e}>"
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"<predict_error: {e}>"
try:
reasoning_text = kls.reason_sync(q)
except Exception as e:
reasoning_text = f"<reasoning_error: {e}>"
latency_ms = (time.time() - t_start) * 1000.0
think_match = re.search(r"<think>(.*?)</think>", reasoning_text, re.DOTALL)
plan_match = re.search(r"<plan>(.*?)</plan>", reasoning_text, re.DOTALL)
answer_match = re.search(r"<answer>(.*?)</answer>", reasoning_text, re.DOTALL)
decompose_match = re.search(r"<decompose>(.*?)</decompose>", 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 ["<think>", "<plan>", "<decompose>", "<answer>"] 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_<alphanumeric>{32,}) e substitui por '<HF_TOKEN_REMOVED>'.
"""
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("<HF_TOKEN_REMOVED>", 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 <think>/<plan>/<answer>
# 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())