Download scripts/train_v6_5_v2.py from PowerMachine/BiGRU_T_version: direct link, hf CLI and curl.
- Browser
- Download file 183 kB
-
https://huggingface.co/PowerMachine/BiGRU_T_version/resolve/main/scripts/train_v6_5_v2.py
- Command line
-
hf download hf://PowerMachine/BiGRU_T_version/scripts/train_v6_5_v2.py
-
curl -L -o train_v6_5_v2.py https://huggingface.co/PowerMachine/BiGRU_T_version/resolve/main/scripts/train_v6_5_v2.py
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()) | |