"""monitoring.py — Monitoramento completo e evolução do CNN-BiGRU. Rastreia métricas de treinamento, inferência, uso de recursos, evolução dos componentes (EWC, Medusa, CyclicReasoning, etc.) e exporta relatórios. ============================================================================== FUNCIONALIDADES ============================================================================== 1. Métricas de treinamento: - Loss total, loss_main, loss_medusa, loss_ewc, loss_hallucination - Perplexidade (PPL) - Learning rate (G, V) - Norma do gradiente - Hypothesis activations - Synergy history 2. Métricas de inferência: - Tokens gerados - Medusa accept rate - Tempo por token (ms) - Throughput (tokens/s) - Memória usada (MB) 3. Métricas de evolução: - Loss/PPL ao longo do tempo (por época) - Convergência do CyclicReasoning (n_cycles, deltas) - Codebook usage do VQ-VAE-2 - Quantização stats (memória economizada, etc.) - EWC penalty evolution 4. Métricas de sistema: - CPU/Memory usage - Device info - Runtime info (Xeon AVX512/AMX) 5. Exportação: - JSON report - CSV time series - Markdown summary ============================================================================== USO ============================================================================== from cnn_bigru.utils.monitoring import Monitor monitor = Monitor(output_dir="/path/to/logs") monitor.start_training() for epoch in range(N): for batch in dataloader: ... monitor.log_batch({ "loss": loss.item(), "ppl": ppl, "lr": lr, "grad_norm": grad_norm, }) monitor.log_epoch({"epoch": epoch, "avg_loss": ...}) monitor.end_training() monitor.export_report() Autor: CNN-BiGRU Project """ from __future__ import annotations import csv import json import logging import os import time from collections import defaultdict, deque from dataclasses import dataclass, field, asdict from pathlib import Path from typing import Any, Dict, List, Optional, Union import torch logger = logging.getLogger(__name__) # ============================================================================ # Data Classes # ============================================================================ @dataclass class BatchMetrics: """Métricas de um batch.""" step: int epoch: int batch_in_epoch: int timestamp: float loss: float = 0.0 loss_main: float = 0.0 loss_medusa: float = 0.0 loss_ewc: float = 0.0 loss_hallucination: float = 0.0 loss_total: float = 0.0 ppl: float = 0.0 lr_g: float = 0.0 lr_v: float = 0.0 grad_norm_g: float = 0.0 grad_norm_v: float = 0.0 v_mean: float = 0.0 hypothesis_activations: int = 0 elapsed_ms: float = 0.0 extra: Dict[str, Any] = field(default_factory=dict) @dataclass class EpochMetrics: """Métricas agregadas de uma época.""" epoch: int avg_loss: float = 0.0 avg_ppl: float = 0.0 n_batches: int = 0 n_hypothesis_activations: int = 0 elapsed_s: float = 0.0 best_loss: float = float("inf") worst_loss: float = 0.0 synergy_history: List[Dict] = field(default_factory=list) extra: Dict[str, Any] = field(default_factory=dict) @dataclass class EvolutionMetrics: """Métricas de evolução dos componentes ao longo do tempo.""" cyclic_reasoning_stats: List[Dict] = field(default_factory=list) vqvae2_codebook_usage: List[Dict] = field(default_factory=list) quantization_stats: List[Dict] = field(default_factory=list) ewc_penalty_evolution: List[Dict] = field(default_factory=list) medusa_accept_rate: List[Dict] = field(default_factory=list) @dataclass class SystemMetrics: """Métricas de sistema.""" timestamp: float cpu_percent: float = 0.0 memory_mb: float = 0.0 gpu_memory_mb: float = 0.0 n_active_threads: int = 0 # ============================================================================ # Monitor # ============================================================================ class Monitor: """Monitor completo de treinamento, inferência e evolução. Args: output_dir: diretório para salvar logs e relatórios max_history: número máximo de métricas em memória (FIFO) log_every: intervalo de log (em batches) track_system: se True, rastreia CPU/memória (requer psutil) """ def __init__( self, output_dir: Optional[Union[str, Path]] = None, max_history: int = 10000, log_every: int = 1, track_system: bool = True, ): self.output_dir = Path(output_dir) if output_dir else None if self.output_dir: self.output_dir.mkdir(parents=True, exist_ok=True) self.max_history = max_history self.log_every = log_every self.track_system = track_system and self._psutil_available() # Histórico self.batch_history: deque = deque(maxlen=max_history) self.epoch_history: List[EpochMetrics] = [] self.evolution = EvolutionMetrics() self.system_history: deque = deque(maxlen=max_history) # Estado self._training_active = False self._inference_active = False self._start_time: Optional[float] = None self._epoch_start_time: Optional[float] = None self._current_epoch = 0 self._current_batch = 0 self._global_step = 0 # Agregados para época atual self._epoch_losses: List[float] = [] self._epoch_ppls: List[float] = [] self._epoch_hypothesis_activations = 0 self._epoch_synergy_history: List[Dict] = [] # Inferência self.inference_stats: Dict[str, Any] = { "total_tokens": 0, "total_time_ms": 0.0, "medusa_accepted": 0, "n_calls": 0, "fallback_to_greedy": 0, } # Component tracking self.component_status: Dict[str, bool] = { "ewc": False, "medusa": False, "cyclic_reasoning": False, "vqvae2": False, "quantization": False, "multimodal_attention": False, "context_window": False, } def _psutil_available(self) -> bool: try: import psutil # noqa: F401 return True except ImportError: return False # ---------------------------------------------------------------------- # Training lifecycle # ---------------------------------------------------------------------- def start_training(self) -> None: """Inicia uma sessão de monitoramento de treino.""" self._training_active = True self._start_time = time.time() self._current_epoch = 0 self._current_batch = 0 self._global_step = 0 logger.info("Monitor: sessão de treino iniciada") def start_epoch(self, epoch: int) -> None: """Inicia uma nova época.""" self._current_epoch = epoch self._epoch_start_time = time.time() self._epoch_losses = [] self._epoch_ppls = [] self._epoch_hypothesis_activations = 0 self._epoch_synergy_history = [] def log_batch(self, metrics: Dict[str, Any]) -> None: """Loga métricas de um batch. Args: metrics: dict com chaves como loss, ppl, lr, grad_norm, etc. """ if not self._training_active: return self._global_step += 1 self._current_batch += 1 ts = time.time() batch_metrics = BatchMetrics( step=self._global_step, epoch=self._current_epoch, batch_in_epoch=self._current_batch, timestamp=ts, loss=metrics.get("loss", 0.0), loss_main=metrics.get("loss_main", 0.0), loss_medusa=metrics.get("loss_medusa", 0.0), loss_ewc=metrics.get("loss_ewc", 0.0), loss_hallucination=metrics.get("loss_hallucination", 0.0), loss_total=metrics.get("loss_total", metrics.get("loss", 0.0)), ppl=metrics.get("ppl", 0.0), lr_g=metrics.get("lr_g", metrics.get("lr", 0.0)), lr_v=metrics.get("lr_v", 0.0), grad_norm_g=metrics.get("grad_norm_g", metrics.get("grad_norm", 0.0)), grad_norm_v=metrics.get("grad_norm_v", 0.0), v_mean=metrics.get("v_mean", 0.0), hypothesis_activations=metrics.get("hypothesis_activations", 0), elapsed_ms=metrics.get("elapsed_ms", 0.0), extra={k: v for k, v in metrics.items() if k not in {"loss", "loss_main", "loss_medusa", "loss_ewc", "loss_hallucination", "loss_total", "ppl", "lr_g", "lr_v", "grad_norm_g", "grad_norm_v", "v_mean", "hypothesis_activations", "elapsed_ms", "lr", "grad_norm"}}, ) self.batch_history.append(batch_metrics) # Acumular para época if batch_metrics.loss > 0: self._epoch_losses.append(batch_metrics.loss) if batch_metrics.ppl > 0: self._epoch_ppls.append(batch_metrics.ppl) self._epoch_hypothesis_activations += batch_metrics.hypothesis_activations if "synergy_history" in metrics: self._epoch_synergy_history.extend(metrics["synergy_history"]) # Log if self._global_step % self.log_every == 0: logger.info( "Batch %d (epoch %d) | loss=%.4f | ppl=%.2f | lr_g=%.2e | hyp=%d", self._global_step, self._current_epoch, batch_metrics.loss, batch_metrics.ppl, batch_metrics.lr_g, batch_metrics.hypothesis_activations, ) # Track system metrics occasionally if self.track_system and self._global_step % 50 == 0: self._track_system() def end_epoch(self, extra: Optional[Dict[str, Any]] = None) -> EpochMetrics: """Finaliza a época atual e retorna métricas agregadas.""" if self._epoch_start_time is None: logger.warning("end_epoch chamado sem start_epoch") return EpochMetrics(epoch=self._current_epoch) elapsed = time.time() - self._epoch_start_time # Agregar if self._epoch_losses: avg_loss = sum(self._epoch_losses) / len(self._epoch_losses) best_loss = min(self._epoch_losses) worst_loss = max(self._epoch_losses) else: avg_loss = best_loss = worst_loss = 0.0 if self._epoch_ppls: avg_ppl = sum(self._epoch_ppls) / len(self._epoch_ppls) else: avg_ppl = 0.0 epoch_metrics = EpochMetrics( epoch=self._current_epoch, avg_loss=avg_loss, avg_ppl=avg_ppl, n_batches=len(self._epoch_losses), n_hypothesis_activations=self._epoch_hypothesis_activations, elapsed_s=elapsed, best_loss=best_loss, worst_loss=worst_loss, synergy_history=list(self._epoch_synergy_history), extra=extra or {}, ) self.epoch_history.append(epoch_metrics) logger.info( "Época %d concluída | avg_loss=%.4f | avg_ppl=%.2f | hyp_act=%d | batches=%d | %.1fs", self._current_epoch, avg_loss, avg_ppl, self._epoch_hypothesis_activations, len(self._epoch_losses), elapsed, ) return epoch_metrics def end_training(self) -> Dict[str, Any]: """Finaliza a sessão de treino e retorna resumo.""" if self._start_time is None: return {} total_time = time.time() - self._start_time self._training_active = False summary = { "total_time_s": total_time, "total_epochs": len(self.epoch_history), "total_batches": self._global_step, "final_loss": self.epoch_history[-1].avg_loss if self.epoch_history else 0.0, "final_ppl": self.epoch_history[-1].avg_ppl if self.epoch_history else 0.0, "best_loss": min((e.avg_loss for e in self.epoch_history), default=0.0), "best_ppl": min((e.avg_ppl for e in self.epoch_history), default=0.0), "total_hypothesis_activations": sum(e.n_hypothesis_activations for e in self.epoch_history), } logger.info("Monitor: sessão de treino finalizada — %s", summary) return summary # ---------------------------------------------------------------------- # Inference tracking # ---------------------------------------------------------------------- def log_inference( self, n_tokens: int, elapsed_ms: float, medusa_accepted: int = 0, fallback_to_greedy: int = 0, ) -> None: """Loga uma chamada de inferência.""" self.inference_stats["total_tokens"] += n_tokens self.inference_stats["total_time_ms"] += elapsed_ms self.inference_stats["medusa_accepted"] += medusa_accepted self.inference_stats["fallback_to_greedy"] += fallback_to_greedy self.inference_stats["n_calls"] += 1 def get_inference_throughput(self) -> Dict[str, float]: """Retorna throughput de inferência.""" s = self.inference_stats if s["total_time_ms"] == 0: return {"tokens_per_s": 0.0, "ms_per_token": 0.0, "medusa_accept_rate": 0.0} total_s = s["total_time_ms"] / 1000.0 tps = s["total_tokens"] / total_s if total_s > 0 else 0.0 mspt = s["total_time_ms"] / s["total_tokens"] if s["total_tokens"] > 0 else 0.0 accept_rate = s["medusa_accepted"] / max(1, s["total_tokens"]) return { "tokens_per_s": tps, "ms_per_token": mspt, "medusa_accept_rate": accept_rate, "n_calls": s["n_calls"], "total_tokens": s["total_tokens"], } # ---------------------------------------------------------------------- # Evolution tracking (componentes específicos) # ---------------------------------------------------------------------- def log_cyclic_reasoning(self, stats: Dict[str, Any]) -> None: """Loga estatísticas do CyclicReasoning.""" stats["step"] = self._global_step stats["timestamp"] = time.time() self.evolution.cyclic_reasoning_stats.append(stats) def log_vqvae2_usage(self, stats: Dict[str, Any]) -> None: """Loga uso do codebook VQ-VAE-2.""" stats["step"] = self._global_step stats["timestamp"] = time.time() self.evolution.vqvae2_codebook_usage.append(stats) def log_quantization(self, stats: Dict[str, Any]) -> None: """Loga estatísticas de quantização W8A8.""" stats["step"] = self._global_step stats["timestamp"] = time.time() self.evolution.quantization_stats.append(stats) def log_ewc_penalty(self, penalty: float, num_tasks: int) -> None: """Loga evolução do penalty EWC.""" self.evolution.ewc_penalty_evolution.append({ "step": self._global_step, "timestamp": time.time(), "penalty": penalty, "num_tasks": num_tasks, }) def log_medusa_accept(self, accepted: int, total: int) -> None: """Loga taxa de aceitação Medusa.""" rate = accepted / max(1, total) self.evolution.medusa_accept_rate.append({ "step": self._global_step, "timestamp": time.time(), "accepted": accepted, "total": total, "rate": rate, }) def register_component(self, name: str, active: bool = True) -> None: """Registra que um componente está ativo.""" if name in self.component_status: self.component_status[name] = active # ---------------------------------------------------------------------- # System tracking # ---------------------------------------------------------------------- def _track_system(self) -> None: """Rastreia métricas de sistema.""" if not self.track_system: return try: import psutil cpu = psutil.cpu_percent(interval=None) mem = psutil.virtual_memory() sm = SystemMetrics( timestamp=time.time(), cpu_percent=cpu, memory_mb=mem.used / (1024 * 1024), n_active_threads=psutil.cpu_count() or 0, ) # GPU memory (se disponível) if torch.cuda.is_available(): try: gpu_mem = torch.cuda.memory_allocated() / (1024 * 1024) sm.gpu_memory_mb = gpu_mem except Exception: pass self.system_history.append(sm) except Exception as e: logger.debug("System tracking falhou: %s", e) # ---------------------------------------------------------------------- # Exportação # ---------------------------------------------------------------------- def export_report(self, output_path: Optional[Union[str, Path]] = None) -> Path: """Exporta relatório completo em JSON. Args: output_path: caminho do arquivo (default: output_dir/monitor_report.json) Returns: Path do arquivo salvo """ if output_path is None: if self.output_dir is None: raise ValueError("output_dir ou output_path deve ser fornecido") output_path = self.output_dir / "monitor_report.json" output_path = Path(output_path) output_path.parent.mkdir(parents=True, exist_ok=True) # Construir relatório report = { "metadata": { "generated_at": time.time(), "training_active": self._training_active, "global_step": self._global_step, "current_epoch": self._current_epoch, }, "components": dict(self.component_status), "training_summary": self._get_training_summary(), "epoch_history": [asdict(e) for e in self.epoch_history], "inference_stats": { **self.inference_stats, "throughput": self.get_inference_throughput(), }, "evolution": { "cyclic_reasoning": self.evolution.cyclic_reasoning_stats[-50:], "vqvae2_codebook_usage": self.evolution.vqvae2_codebook_usage[-50:], "quantization_stats": self.evolution.quantization_stats[-10:], "ewc_penalty_evolution": self.evolution.ewc_penalty_evolution[-50:], "medusa_accept_rate": self.evolution.medusa_accept_rate[-50:], }, "system_metrics": [asdict(s) for s in list(self.system_history)[-50:]], } with open(output_path, "w", encoding="utf-8") as f: json.dump(report, f, indent=2, default=str, ensure_ascii=False) logger.info("Monitor report salvo em: %s", output_path) return output_path def export_csv(self, output_path: Optional[Union[str, Path]] = None) -> Path: """Exporta métricas de batch em CSV.""" if output_path is None: if self.output_dir is None: raise ValueError("output_dir ou output_path deve ser fornecido") output_path = self.output_dir / "batch_metrics.csv" output_path = Path(output_path) output_path.parent.mkdir(parents=True, exist_ok=True) if not self.batch_history: logger.warning("Sem batch history para exportar") return output_path # Pegar chaves do primeiro item + extras first = self.batch_history[0] fieldnames = ["step", "epoch", "batch_in_epoch", "timestamp", "loss", "loss_main", "loss_medusa", "loss_ewc", "loss_hallucination", "loss_total", "ppl", "lr_g", "lr_v", "grad_norm_g", "grad_norm_v", "v_mean", "hypothesis_activations", "elapsed_ms"] with open(output_path, "w", newline="", encoding="utf-8") as f: writer = csv.DictWriter(f, fieldnames=fieldnames, extrasaction="ignore") writer.writeheader() for m in self.batch_history: writer.writerow(asdict(m)) logger.info("CSV metrics salvo em: %s (%d rows)", output_path, len(self.batch_history)) return output_path def export_markdown_summary(self, output_path: Optional[Union[str, Path]] = None) -> Path: """Exporta resumo em Markdown.""" if output_path is None: if self.output_dir is None: raise ValueError("output_dir ou output_path deve ser fornecido") output_path = self.output_dir / "monitor_summary.md" output_path = Path(output_path) output_path.parent.mkdir(parents=True, exist_ok=True) summary = self._get_training_summary() throughput = self.get_inference_throughput() lines = [ "# Monitor Report — CNN-BiGRU", "", f"**Generated at:** {time.strftime('%Y-%m-%d %H:%M:%S')}", "", "## Components Status", "", ] for comp, active in self.component_status.items(): icon = "[x]" if active else "[ ]" lines.append(f"- {icon} {comp}") lines.extend([ "", "## Training Summary", "", f"- Total time: {summary.get('total_time_s', 0):.1f}s", f"- Total epochs: {summary.get('total_epochs', 0)}", f"- Total batches: {summary.get('total_batches', 0)}", f"- Final loss: {summary.get('final_loss', 0):.4f}", f"- Final PPL: {summary.get('final_ppl', 0):.2f}", f"- Best loss: {summary.get('best_loss', 0):.4f}", f"- Best PPL: {summary.get('best_ppl', 0):.2f}", f"- Total hypothesis activations: {summary.get('total_hypothesis_activations', 0)}", "", "## Inference Throughput", "", f"- Tokens/s: {throughput.get('tokens_per_s', 0):.1f}", f"- ms/token: {throughput.get('ms_per_token', 0):.2f}", f"- Medusa accept rate: {throughput.get('medusa_accept_rate', 0):.2%}", f"- Total tokens generated: {throughput.get('total_tokens', 0)}", f"- Total inference calls: {throughput.get('n_calls', 0)}", "", "## Evolution", "", f"- CyclicReasoning events: {len(self.evolution.cyclic_reasoning_stats)}", f"- VQ-VAE-2 usage events: {len(self.evolution.vqvae2_codebook_usage)}", f"- Quantization events: {len(self.evolution.quantization_stats)}", f"- EWC penalty events: {len(self.evolution.ewc_penalty_evolution)}", f"- Medusa accept events: {len(self.evolution.medusa_accept_rate)}", "", "## Epoch History", "", "| Epoch | Avg Loss | Avg PPL | Batches | Hyp Act | Time(s) |", "|-------|----------|---------|---------|---------|---------|", ]) for e in self.epoch_history: lines.append( f"| {e.epoch} | {e.avg_loss:.4f} | {e.avg_ppl:.2f} | " f"{e.n_batches} | {e.n_hypothesis_activations} | {e.elapsed_s:.1f} |" ) with open(output_path, "w", encoding="utf-8") as f: f.write("\n".join(lines)) logger.info("Markdown summary salvo em: %s", output_path) return output_path def _get_training_summary(self) -> Dict[str, Any]: """Constrói sumário de treino.""" if not self.epoch_history: return {"total_time_s": 0, "total_epochs": 0, "total_batches": self._global_step} return { "total_time_s": time.time() - (self._start_time or time.time()), "total_epochs": len(self.epoch_history), "total_batches": self._global_step, "final_loss": self.epoch_history[-1].avg_loss, "final_ppl": self.epoch_history[-1].avg_ppl, "best_loss": min((e.avg_loss for e in self.epoch_history), default=0.0), "best_ppl": min((e.avg_ppl for e in self.epoch_history), default=0.0), "total_hypothesis_activations": sum( e.n_hypothesis_activations for e in self.epoch_history ), } # ---------------------------------------------------------------------- # Reset # ---------------------------------------------------------------------- def reset(self) -> None: """Limpa todo o histórico.""" self.batch_history.clear() self.epoch_history.clear() self.evolution = EvolutionMetrics() self.system_history.clear() self._global_step = 0 self._current_epoch = 0 self._current_batch = 0 self._start_time = None self._epoch_start_time = None self.inference_stats = {k: 0 if isinstance(v, int) else 0.0 for k, v in self.inference_stats.items()} # ============================================================================ # Singleton instance (opcional) # ============================================================================ _global_monitor: Optional[Monitor] = None def get_monitor(output_dir: Optional[Union[str, Path]] = None) -> Monitor: """Retorna a instância global do monitor (singleton).""" global _global_monitor if _global_monitor is None: _global_monitor = Monitor(output_dir=output_dir) return _global_monitor # ============================================================================ # Self-test # ============================================================================ def _self_test(): """Teste rápido do monitor.""" import tempfile with tempfile.TemporaryDirectory() as tmpdir: mon = Monitor(output_dir=tmpdir, log_every=1, track_system=True) mon.start_training() mon.start_epoch(0) for i in range(5): mon.log_batch({ "loss": 5.0 - i * 0.5, "ppl": 100.0 - i * 10, "lr_g": 1e-3, "grad_norm": 0.5 + i * 0.1, "hypothesis_activations": i, "elapsed_ms": 10 + i, }) mon.end_epoch() mon.end_training() mon.log_inference(n_tokens=50, elapsed_ms=500, medusa_accepted=10) mon.log_cyclic_reasoning({"n_cycles": 3, "converged": True}) mon.log_vqvae2_usage({"top_usage": 0.5, "bottom_usage": 0.7}) mon.log_quantization({"reduction_pct": 60.0}) mon.log_ewc_penalty(penalty=0.001, num_tasks=1) mon.log_medusa_accept(accepted=10, total=50) report = mon.export_report() csv_path = mon.export_csv() md_path = mon.export_markdown_summary() print(f"Report: {report}") print(f"CSV: {csv_path}") print(f"MD: {md_path}") print(f"Throughput: {mon.get_inference_throughput()}") if __name__ == "__main__": _self_test() __all__ = [ "BatchMetrics", "EpochMetrics", "EvolutionMetrics", "SystemMetrics", "Monitor", "get_monitor", ]