# type: ignore """ Aprendizado contínuo simples para AKIRA V21 - Registra todas as mensagens (PV/Grupo), replies e respostas geradas - Persiste em JSONL em data/continuous_learning.jsonl - Fornece contexto global resumido para alimentar o LLM quando solicitado - Sugere melhor API baseada em heurísticas leves """ import os import json import time import threading from pathlib import Path from typing import Optional, Dict, Any, List, Set # Imports robustos com fallback try: from . import config from .database import Database except ImportError: try: import modules.config as config from modules.database import Database except ImportError: config = None Database = None DATA_DIR: Path = getattr(config, 'DATA_DIR', Path('./data')) DATA_DIR.mkdir(parents=True, exist_ok=True) JSONL_PATH: Path = DATA_DIR / 'continuous_learning.jsonl' LOCK = threading.Lock() class AprendizadoContinuo: def __init__(self, jsonl_path: Path, db: Optional[Any] = None): self.path = jsonl_path self.path.parent.mkdir(parents=True, exist_ok=True) self.db = db # índice leve em memória (opcional) self._buffer: List[Dict[str, Any]] = [] self._buffer_limit = 2000 self._idempotency_cache: Set[str] = set() # ✅ Cache de IDs para evitar duplicados no mesmo worker def _append_jsonl(self, row: Dict[str, Any]) -> None: """Salva no PostgreSQL (ou JSONL como fallback se DB não disponível).""" if self.db: try: self.db._execute_with_retry(""" INSERT INTO continuous_learning (ts, usuario, numero, nome_usuario, tipo_conversa, mensagem, resposta_do_bot, resposta_gerada, is_reply, reply_to_bot, contexto_grupo, modelo_usado, message_id, qualidade, tipo_conteudo) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s) ON CONFLICT (message_id) DO NOTHING """, ( row.get('ts'), row.get('usuario'), row.get('numero'), row.get('nome_usuario'), row.get('tipo_conversa'), row.get('mensagem'), row.get('resposta_do_bot'), row.get('resposta_gerada'), row.get('is_reply'), row.get('reply_to_bot'), row.get('contexto_grupo'), row.get('modelo_usado'), row.get('message_id'), row.get('qualidade', 0.0), row.get('tipo_conteudo', 'desconhecido') ), commit=True) return except Exception: pass # Fallback: JSONL se DB não disponível with LOCK: with self.path.open('a', encoding='utf-8') as f: f.write(json.dumps(row, ensure_ascii=False) + '\n') self._buffer.append(row) if len(self._buffer) > self._buffer_limit: self._buffer = self._buffer[-self._buffer_limit:] def _now_ts(self) -> float: return time.time() def processar_mensagem( self, mensagem: str, usuario: str, numero: str, nome_usuario: Optional[str] = None, tipo_conversa: str = 'pv', # 'pv' ou 'grupo' resposta_do_bot: bool = False, resposta_gerada: Optional[str] = None, is_reply: bool = False, reply_to_bot: bool = False, # ✅ Restaurado contexto_grupo: Optional[str] = None, modelo_usado: Optional[str] = None, message_id: Optional[str] = None, # ✅ Adicionado para idempotência ) -> Dict[str, Any]: """Registra evento para aprendizado contínuo com filtragem de qualidade.""" # 0. Verificação de idempotência if message_id and message_id in self._idempotency_cache: return {'status': 'ignored', 'motivo': 'duplicado_idempotencia'} if message_id: self._idempotency_cache.add(message_id) # Limpa cache se ficar muito grande if len(self._idempotency_cache) > 5000: self._idempotency_cache.clear() mensagem_norm = (mensagem or '').strip() if not mensagem_norm: return {'status': 'ignored', 'motivo': 'mensagem_vazia'} # ============================================================ # FILTRO DE QUALIDADE — decide se deve ser aprendido ou descartado # ============================================================ # 1. Descarta mensagens muito curtas (spam/ruído) palavras = mensagem_norm.split() if len(palavras) < 2 and not resposta_do_bot: return {'status': 'discarded', 'motivo': 'muito_curta', 'analise': {'comprimento': len(palavras)}} # 2. Descarta spam de links puros if mensagem_norm.startswith('http://') or mensagem_norm.startswith('https://'): return {'status': 'discarded', 'motivo': 'link_puro'} # 3. Descarta caracteres repetidos (ex: "kkkkkkkk", "aaaaa") if len(set(mensagem_norm.lower())) < 4 and len(mensagem_norm) > 5: return {'status': 'discarded', 'motivo': 'repetitivo'} # 4. Detecta tipo de conteúdo para priorizar treino qualidade = self._avaliar_qualidade(mensagem_norm, resposta_do_bot) row = { 'ts': self._now_ts(), 'usuario': usuario, 'numero': numero, 'nome_usuario': nome_usuario or usuario, 'tipo_conversa': tipo_conversa, 'mensagem': mensagem_norm[:4000], 'resposta_do_bot': bool(resposta_do_bot), 'resposta_gerada': (resposta_gerada or '')[:4000] if resposta_do_bot else None, 'is_reply': bool(is_reply), 'reply_to_bot': bool(reply_to_bot), 'contexto_grupo': contexto_grupo or '', 'modelo_usado': modelo_usado or 'desconhecido', 'message_id': message_id or '', # ✅ Salva ID para rastreio 'qualidade': qualidade, # score 0.0-1.0 } self._append_jsonl(row) # ============================================================ # INTEGRAÇÃO COM DATABASE (Sincronização SQLite) # ============================================================ if self.db: try: # Salva metadados de qualidade para o RAG/Contexto self.db.salvar_aprendizado_detalhado( usuario=usuario, chave=f"qualidade_msg_{int(row['ts'])}", valor=json.dumps({ "msg_prefix": mensagem_norm[:50], "qualidade": qualidade, "tipo": row['tipo_conteudo'] if 'tipo_conteudo' in locals() else self._classificar_conteudo(mensagem_norm, resposta_do_bot) }, ensure_ascii=False) ) except Exception as e: if hasattr(config, 'DEBUG_MODE') and config.DEBUG_MODE: print(f"Erro ao persistir qualidade no DB: {e}") analise = { 'comprimento': len(palavras), 'tem_link': ('http://' in mensagem_norm) or ('https://' in mensagem_norm), 'tem_interrogacao': '?' in mensagem_norm, 'qualidade': qualidade, 'tipo_conteudo': self._classificar_conteudo(mensagem_norm, resposta_do_bot), } aprendizado = {'armazenado_em': str(self.path)} return {'ok': True, 'analise': analise, 'aprendizado': aprendizado} def _avaliar_qualidade(self, mensagem: str, resposta_do_bot: bool) -> float: """Avalia qualidade de uma mensagem para aprendizado (0.0-1.0).""" score = 0.3 # baseline palavras = mensagem.split() n_palavras = len(palavras) # Comprimento: mensagens médias são mais úteis if 5 <= n_palavras <= 50: score += 0.2 elif n_palavras > 50: score += 0.1 # Perguntas são valiosas (curiosidade do usuário) if '?' in mensagem: score += 0.2 # Pares Q&A do bot são ouro para treino if resposta_do_bot: score += 0.3 # Replies ao bot indicam engajamento if n_palavras > 3: score += 0.1 # Hashtags/comandos são menos úteis para treino de linguagem if mensagem.startswith('!') or mensagem.startswith('.'): score -= 0.2 return min(max(score, 0.0), 1.0) def _classificar_conteudo(self, mensagem: str, resposta_do_bot: bool) -> str: """Classifica tipo de conteúdo para saber O QUE treinar.""" if resposta_do_bot: return 'resposta_bot' if '?' in mensagem: return 'pergunta' if mensagem.startswith('!') or mensagem.startswith('.'): return 'comando' palavras_lower = mensagem.lower() if any(w in palavras_lower for w in ['porque', 'por que', 'como', 'quando', 'onde', 'quem', 'qual']): return 'pergunta_indireta' if 'http://' in palavras_lower or 'https://' in palavras_lower: return 'com_link' if len(mensagem.split()) > 20: return 'texto_longo' return 'conversa_comum' def obter_contexto_para_llm(self, topico: Optional[str] = None, limite: int = 10) -> List[str]: """Retorna últimas N mensagens do PostgreSQL (ou JSONL como fallback).""" registros: List[Dict[str, Any]] = [] # Tenta ler do PostgreSQL primeiro if self.db: try: rows = self.db._execute_with_retry( """SELECT usuario, nome_usuario, mensagem, resposta_do_bot, resposta_gerada, tipo_conversa, modelo_usado, qualidade FROM continuous_learning WHERE mensagem IS NOT NULL AND LENGTH(mensagem) > 5 ORDER BY ts DESC LIMIT %s""", (min(limite * 50, 500),) ) if rows: for r in rows: registros.append({ 'usuario': r['usuario'] if isinstance(r, dict) else r[0], 'nome_usuario': r['nome_usuario'] if isinstance(r, dict) else r[1], 'mensagem': r['mensagem'] if isinstance(r, dict) else r[2], 'resposta_do_bot': r['resposta_do_bot'] if isinstance(r, dict) else r[3], 'resposta_gerada': r['resposta_gerada'] if isinstance(r, dict) else r[4], 'tipo_conversa': r['tipo_conversa'] if isinstance(r, dict) else r[5], 'modelo_usado': r['modelo_usado'] if isinstance(r, dict) else r[6], 'qualidade': r['qualidade'] if isinstance(r, dict) else r[7], }) except Exception: pass # Fallback: JSONL if not registros: try: if self.path.exists(): linhas: List[str] = [] with self.path.open('r', encoding='utf-8') as f: for line in f: linhas.append(line) linhas = linhas[-2000:] for line in linhas[-500:]: try: registros.append(json.loads(line)) except Exception: continue except Exception: pass # filtra if topico: t = topico.lower().strip() registros = [r for r in registros if t in (r.get('mensagem', '').lower())] # monta blocos curtos para contexto blocos: List[str] = [] for r in registros[-limite:]: autor = r.get('nome_usuario') or r.get('usuario') msg = r.get('mensagem', '') tipo = r.get('tipo_conversa', 'pv') blocos.append(f"[{tipo}] {autor}: {msg}") return blocos def get_best_api_for_context( self, complexidade: float = 0.5, emocao: str = 'neutral', intencao: str = 'afirmacao', tipo_conversa: str = 'pv', ) -> str: """ HEURÍSTICA DELEGADA AO LOCAL_LLM / MOE ROUTER. Mantido para compatibilidade, mas agora apenas sugere o padrão. """ return 'moe_router' _singleton: Optional[AprendizadoContinuo] = None def get_aprendizado_continuo(db: Optional[Any] = None) -> AprendizadoContinuo: global _singleton if _singleton is None: _singleton = AprendizadoContinuo(JSONL_PATH, db=db) elif db and _singleton.db is None: _singleton.db = db return _singleton # ============================================================ # COMPATIBILIDADE — aliases para imports legados # ============================================================ def processar_conversa_global( mensagem: str, usuario: str, numero: str, nome_usuario: Optional[str] = None, tipo_conversa: str = 'pv', resposta_do_bot: bool = False, resposta_gerada: Optional[str] = None, is_reply: bool = False, reply_to_bot: bool = False, contexto_grupo: Optional[str] = None, modelo_usado: Optional[str] = None, ) -> Dict[str, Any]: """Wrapper legado — delega para o singleton.""" ac = get_aprendizado_continuo() return ac.processar_mensagem( mensagem=mensagem, usuario=usuario, numero=numero, nome_usuario=nome_usuario, tipo_conversa=tipo_conversa, resposta_do_bot=resposta_do_bot, resposta_gerada=resposta_gerada, is_reply=is_reply, reply_to_bot=reply_to_bot, contexto_grupo=contexto_grupo, modelo_usado=modelo_usado, message_id=None ) # Aliases de classe para compatibilidade ConversaGlobal = AprendizadoContinuo APIContextScore = type('APIContextScore', (), { 'score': 0.5, 'api': 'gemini', '__init__': lambda self, **kw: self.__dict__.update(kw), })