""" ================================================================================ AKIRA V21 ULTIMATE - POSTGRESQL DATABASE MODULE ================================================================================ Drop-in replacement para database.py (SQLite) usando PostgreSQL. Permite 2+ workers sem problemas de concorrência. Env vars necessárias: DATABASE_URL=postgresql://user:pass@localhost:5432/akira ou PGHOST, PGPORT, PGDATABASE, PGUSER, PGPASSWORD ================================================================================ """ import os import time import json import hashlib import random import re import threading from typing import Optional, List, Dict, Any, Tuple, Union from datetime import datetime from loguru import logger try: import psycopg2 import psycopg2.extras import psycopg2.errors HAS_PG = True except ImportError: HAS_PG = False logger.warning("psycopg2 não instalado — PostgreSQL indisponível") def _safe_json_load(value): """Parse JSON safely — returns dict/list as-is if already deserialized by psycopg2.""" if value is None: return None if isinstance(value, (dict, list)): return value if isinstance(value, str): try: return json.loads(value) except (json.JSONDecodeError, TypeError): return value return value # ===================================================================== # CIRCUIT BREAKER PG — Space sem Postgres válido não deve tentar ligar # a cada mensagem (spam de logs + ~1.5s de latência por retry). Após N # falhas seguidas o PG desativa-se na sessão e os callers caem no SQLite. # ===================================================================== _PG_FAIL_LIMIT = int(os.getenv("AKIRA_PG_FAIL_LIMIT", "3") or 3) _pg_state = {"fails": 0, "disabled": False} def _pg_disabled() -> bool: return _pg_state["disabled"] def _pg_failure(err) -> None: if _pg_state["disabled"]: return _pg_state["fails"] += 1 if _pg_state["fails"] >= _PG_FAIL_LIMIT: _pg_state["disabled"] = True logger.warning( f"⚠️ PostgreSQL indisponível ({_pg_state['fails']} falhas seguidas) — " f"PG desativado nesta sessão; a usar SQLite. Último erro: {str(err)[:200]}" ) def _pg_success() -> None: _pg_state["fails"] = 0 def _pg_unavailable_error(): if HAS_PG: return psycopg2.OperationalError("PG desativado apos falhas consecutivas (circuit breaker)") return Exception("PG desativado apos falhas consecutivas (circuit breaker)") class ConnectionPool: """ Pool de conexões simples para PostgreSQL. Reduz overhead de criar/conexões frequentemente. """ def __init__(self, conn_params: dict, max_connections: int = 10): self._conn_params = conn_params self._max_connections = max_connections self._pool: list = [] self._lock = threading.Lock() self._created = 0 def get_connection(self): """Obtém conexão do pool ou cria nova""" if _pg_disabled(): raise _pg_unavailable_error() with self._lock: # Tentar reutilizar conexão existente while self._pool: conn = self._pool.pop() try: # Verificar se conexão ainda está válida if conn.closed: continue # Testar com ping simples cur = conn.cursor() cur.execute("SELECT 1") cur.close() _pg_success() return conn except Exception: try: conn.close() except: pass continue # Criar nova conexão se pool vazio if self._created < self._max_connections: try: params = self._conn_params if 'dsn' in params: conn = psycopg2.connect(params['dsn'], cursor_factory=psycopg2.extras.RealDictCursor) else: conn = psycopg2.connect(**params, cursor_factory=psycopg2.extras.RealDictCursor) conn.autocommit = False self._created += 1 _pg_success() return conn except Exception as e: _pg_failure(e) logger.error(f"Erro ao criar conexão: {e}") raise # Pool cheio — aguardar curto e retry (max 5s) logger.warning(f"⚠️ Pool de conexões cheio ({self._created}/{self._max_connections}), aguardando...") self._lock.release() for _wait in range(50): time.sleep(0.1) with self._lock: if len(self._pool) > 0: conn = self._pool.pop() if not conn.closed: self._lock.release() return conn else: self._created = max(0, self._created - 1) elif self._created < self._max_connections: self._lock.release() try: params = self._conn_params if 'dsn' in params: conn = psycopg2.connect(params['dsn'], cursor_factory=psycopg2.extras.RealDictCursor) else: conn = psycopg2.connect(**params, cursor_factory=psycopg2.extras.RealDictCursor) conn.autocommit = False self._created += 1 return conn except Exception as e: _pg_failure(e) logger.error(f"Erro ao criar conexão: {e}") raise else: self._lock.release() raise Exception("Pool de conexões esgotado após 5s de espera") def return_connection(self, conn): """Retorna conexão ao pool""" if conn is None or conn.closed: with self._lock: self._created = max(0, self._created - 1) return with self._lock: if len(self._pool) < self._max_connections: try: # Resetar estado da conexão conn.rollback() self._pool.append(conn) self._created = max(0, self._created - 1) except Exception: try: conn.close() except: pass self._created = max(0, self._created - 1) else: try: conn.close() except: pass self._created = max(0, self._created - 1) def close_all(self): """Fecha todas as conexões""" with self._lock: for conn in self._pool: try: conn.close() except: pass self._pool.clear() self._created = 0 class DatabasePG: """ PostgreSQL drop-in replacement para a classe Database (SQLite). Mantém a mesma interface pública. """ _instances: Dict[str, 'DatabasePG'] = {} _initialized: Dict[str, bool] = {} _url_cache: Optional[str] = None # Cache do DATABASE_URL CODIGOS_VERIFICACAO: Dict[str, str] = {} _pool: Optional['ConnectionPool'] = None # Pool de conexões def __new__(cls, db_path: str = ""): if cls._url_cache is None: cls._url_cache = os.environ.get('DATABASE_URL', 'pg_default') key = cls._url_cache if key not in cls._instances: cls._instances[key] = super(DatabasePG, cls).__new__(cls) cls._initialized[key] = False return cls._instances[key] def __init__(self, db_path: str = ""): if self._url_cache is None: self.__class__._url_cache = os.environ.get('DATABASE_URL', 'pg_default') key = self._url_cache if self._initialized.get(key, False): return # Já inicializado — retorno rápido self.max_retries = 5 self.retry_delay = 0.1 self._conn_params = self._get_conn_params() # ✅ CONNECTION POOL: Inicializa pool de conexões if DatabasePG._pool is None: DatabasePG._pool = ConnectionPool(self._conn_params, max_connections=10) logger.info("✅ Connection pool inicializado (max 10 conexões)") if hasattr(self, '_init_db'): self._init_db() if hasattr(self, '_init_context_isolation_tables'): self._init_context_isolation_tables() if hasattr(self, '_ensure_all_columns_and_indexes'): self._ensure_all_columns_and_indexes() DatabasePG._initialized[key] = True logger.success("Database PostgreSQL inicializado") # ================================================================ # CONEXÃO # ================================================================ def _get_conn_params(self) -> dict: url = os.environ.get('DATABASE_URL', '') if url: return {'dsn': url} return { 'host': os.environ.get('PGHOST', 'localhost'), 'port': os.environ.get('PGPORT', '5432'), 'dbname': os.environ.get('PGDATABASE', 'akira'), 'user': os.environ.get('PGUSER', 'akira'), 'password': os.environ.get('PGPASSWORD', 'akira'), } def _get_connection(self): # ✅ CIRCUIT BREAKER: PG já provou estar indisponível → falha rápida if _pg_disabled(): raise _pg_unavailable_error() # ✅ CONNECTION POOL: Usar pool quando disponível if DatabasePG._pool: try: return DatabasePG._pool.get_connection() except Exception as e: if _pg_disabled(): raise logger.warning(f"⚠️ Pool falhou, criando conexão direta: {e}") # Fallback: criar conexão direta for attempt in range(self.max_retries): try: params = self._conn_params if 'dsn' in params: conn = psycopg2.connect(params['dsn'], cursor_factory=psycopg2.extras.RealDictCursor) else: conn = psycopg2.connect(**params, cursor_factory=psycopg2.extras.RealDictCursor) conn.autocommit = False _pg_success() return conn except psycopg2.OperationalError as e: if attempt < self.max_retries - 1: time.sleep(self.retry_delay * (2 ** attempt)) continue _pg_failure(e) logger.error(f"Erro ao conectar ao PostgreSQL: {e}") raise raise psycopg2.OperationalError("Falha ao conectar ao PostgreSQL após retries") def return_connection(self, conn): """Retorna conexão ao pool. Delega ao ConnectionPool.""" if DatabasePG._pool: try: DatabasePG._pool.return_connection(conn) return except Exception: pass # Fallback: fecha se pool não disponível try: if conn and not conn.closed: conn.close() except Exception: pass def _return_conn(self, conn): """Alias seguro para return_connection - usado internamente.""" self.return_connection(conn) # ================================================================ # DISTRIBUTED LOCK — pg_advisory_lock entre workers # ================================================================ @staticmethod def _hash_lock_id(key: str) -> int: """Gera um bigint de 64 bits a partir de uma string para advisory lock.""" h = hashlib.sha256(key.encode()).hexdigest()[:16] return int(h, 16) & 0x7FFFFFFFFFFFFFFF def acquire_advisory_lock(self, lock_key: str) -> Optional[Any]: """ Tenta adquirir um lock distribuído via pg_try_advisory_lock. Retorna a conexão (que mantém o lock) ou None se não conseguiu. A conexão DEVE ser fechada com release_advisory_lock() para liberar o lock. """ lock_id = self._hash_lock_id(lock_key) try: conn = self._get_connection() conn.autocommit = True cur = conn.cursor() cur.execute("SELECT pg_try_advisory_lock(%s)", (lock_id,)) acquired = cur.fetchone()[0] if acquired: return conn _pool_return(conn) return None except Exception as e: logger.debug(f"[LOCK] Erro ao adquirir lock '{lock_key[:40]}': {e}") if conn: try: _pool_return(conn) except Exception: pass return None def release_advisory_lock(self, lock_key: str, conn: Optional[Any]) -> bool: """Libera o lock e fecha a conexão.""" if conn is None: return False lock_id = self._hash_lock_id(lock_key) try: cur = conn.cursor() cur.execute("SELECT pg_advisory_unlock(%s)", (lock_id,)) _pool_return(conn) return True except Exception as e: logger.debug(f"[LOCK] Erro ao liberar lock: {e}") try: _pool_return(conn) except Exception: pass return False def wait_for_advisory_lock(self, lock_key: str, timeout: float = 15.0, poll_interval: float = 0.2) -> Optional[Any]: """ Polling com backoff até adquirir o lock ou estourar timeout. Retorna conexão com lock adquirido, ou None se timeout. """ start = time.time() attempt = 0 while time.time() - start < timeout: conn = self.acquire_advisory_lock(lock_key) if conn is not None: elapsed = time.time() - start if elapsed > 0.5: logger.info(f"[LOCK] Lock '{lock_key[:40]}' adquirido após {elapsed:.1f}s ({attempt} tentativas)") return conn attempt += 1 sleep_time = min(poll_interval * (1 + attempt * 0.15), 1.5) time.sleep(sleep_time) logger.warning(f"[LOCK] Timeout ({timeout}s) ao aguardar lock '{lock_key[:40]}'") return None def _execute_with_retry(self, query: str, params: Optional[tuple] = None, commit: bool = False): """Executa query com retry. Converte syntax SQLite→PG automaticamente.""" pg_query = self._convert_query(query) for attempt in range(self.max_retries): conn = None try: conn = self._get_connection() cur = conn.cursor() cur.execute(pg_query, params or ()) if pg_query.strip().upper().startswith("SELECT"): result = cur.fetchall() if commit: conn.commit() return result if commit: conn.commit() return None except psycopg2.errors.UniqueViolation: if conn: conn.rollback() return None except psycopg2.OperationalError as e: if "locked" in str(e) and attempt < self.max_retries - 1: time.sleep(self.retry_delay * (2 ** attempt)) if conn: conn.rollback() continue if conn: conn.rollback() logger.error(f"Erro SQL PG: {e}") raise except Exception as e: if conn: conn.rollback() import traceback as _tb logger.error(f"Erro SQL PG: {type(e).__name__}: {e}\nquery={pg_query[:120]}\nparams={params}\ntraceback={_tb.format_exc()}") raise finally: if conn: try: _pool_return(conn) except Exception: try: conn.close() except: pass raise Exception("Query falhou após retries") def _execute_simple(self, query: str, params: Optional[tuple] = None): """Query simples com cursor tuple — evita bug psycopg2 com RealDictCursor.""" if _pg_disabled(): return [] pg_query = self._convert_query(query) conn = None try: params_dict = self._conn_params if 'dsn' in params_dict: conn = psycopg2.connect(params_dict['dsn']) else: conn = psycopg2.connect(**params_dict) cur = conn.cursor() cur.execute(pg_query, params or ()) if pg_query.strip().upper().startswith("SELECT"): _pg_success() return cur.fetchall() _pg_success() return None except Exception as e: _pg_failure(e) if conn: try: conn.rollback() except: pass return [] finally: if conn: try: conn.close() except: pass # ================================================================ # PUBLIC API - Cursor Management (para operações de embedding) # ================================================================ def get_connection_context(self): """ Retorna uma conexão em contexto manager para operações que precisam de cursor. Uso: with db.get_connection_context() as conn: cur = conn.cursor() cur.execute("INSERT INTO ... VALUES ...") conn.commit() """ class ConnectionContextManager: def __init__(self, db_instance): self.db = db_instance self.conn = None def __enter__(self): self.conn = self.db._get_connection() return self.conn def __exit__(self, exc_type, exc_val, exc_tb): if self.conn: try: if exc_type: self.conn.rollback() else: self.conn.commit() finally: _pool_return(self.conn) return ConnectionContextManager(self) # ================================================================ # CONVERSÃO SQLite → PostgreSQL # ================================================================ def _convert_query(self, query: str) -> str: q = query.strip() # 1. INSERT OR IGNORE → INSERT ... ON CONFLICT DO NOTHING q = re.sub( r'INSERT\s+OR\s+IGNORE\s+INTO', 'INSERT INTO', q, flags=re.IGNORECASE ) # 2. INSERT OR REPLACE → INSERT ... ON CONFLICT (pk) DO UPDATE SET ALL # Padrão: INSERT OR REPLACE INTO tabela (col1, col2, ...) VALUES (...) or_match = re.match( r'INSERT\s+OR\s+REPLACE\s+INTO\s+(\w+)\s*\(([^)]+)\)\s*VALUES\s*\(([^)]+)\)', q, flags=re.IGNORECASE ) if or_match: table = or_match.group(1) cols_str = or_match.group(2) vals_str = or_match.group(3) cols = [c.strip() for c in cols_str.split(',')] # PKs compostas conhecidas composite_pks = { 'lstm_contexto': ['context_id', 'numero_usuario'], 'lstm_message_links': ['context_id', 'message_id', 'numero_usuario'], } pk_cols = composite_pks.get(table, [cols[0]]) # Gera SET para todas as colunas exceto as PKs set_parts = [] for col in cols: if col not in pk_cols: set_parts.append(f"{col} = EXCLUDED.{col}") set_clause = ', '.join(set_parts) conflict_cols = ', '.join(pk_cols) q = f"INSERT INTO {table} ({cols_str}) VALUES ({vals_str}) ON CONFLICT ({conflict_cols}) DO UPDATE SET {set_clause}" # 3. Placeholders ? → %s q = q.replace('?', '%s') # 4. AUTOINCREMENT → SERIAL q = re.sub( r'(?:INTEGER\s+)?PRIMARY\s+KEY\s+AUTOINCREMENT', 'SERIAL PRIMARY KEY', q, flags=re.IGNORECASE ) # 5. strftime('%s', 'now') → EXTRACT(EPOCH FROM NOW()) q = re.sub( r"strftime\('%s',\s*'now'\)", 'EXTRACT(EPOCH FROM NOW())', q, flags=re.IGNORECASE ) # 6. CURRENT_TIMESTAMP → NOW() q = re.sub(r'CURRENT_TIMESTAMP', 'NOW()', q, flags=re.IGNORECASE) # 7. datetime('now') / datetime('now', '-24 hours') → NOW() / NOW() - INTERVAL '24 hours' q = re.sub( r"datetime\('now'(?:\s*,\s*'([^']+)')?\)", lambda m: f"NOW() - INTERVAL '{m.group(1)}'" if m.group(1) else "NOW()", q, flags=re.IGNORECASE ) # 8. sqlite_master → information_schema.tables q = re.sub( r"sqlite_master", "information_schema.tables", q, flags=re.IGNORECASE ) q = re.sub( r"WHERE\s+type\s*=\s*'table'\s*AND\s+name\s*=\s*'(\w+)'", r"WHERE table_name = '\1' AND table_schema = 'public'", q, flags=re.IGNORECASE ) # 9. CREATE TABLE (sem IF NOT EXISTS) → adiciona if q.upper().startswith("CREATE TABLE") and "IF NOT EXISTS" not in q.upper(): q = q.replace("CREATE TABLE", "CREATE TABLE IF NOT EXISTS", 1) # 10. CREATE INDEX → IF NOT EXISTS q = re.sub( r'CREATE\s+INDEX\s+(?!IF)', 'CREATE INDEX IF NOT EXISTS ', q, flags=re.IGNORECASE ) return q # ================================================================ # SCHEMA # ================================================================ def _init_db(self): """Inicializa schema com cada DDL isolado - rollback explícito em erro.""" conn = self._get_connection() conn.autocommit = True cur = conn.cursor() # Helper para executar DDL isoladamente def exec_ddl(sql: str, desc: str = "") -> bool: try: cur.execute(sql) return True except Exception as e: # Rollback explícito para limpar estado de transação abortada try: conn.rollback() except Exception: pass logger.debug(f"DDL {desc} falhou (ignorado): {e}") return False # Tabelas principais exec_ddl(""" CREATE TABLE IF NOT EXISTS mensagens ( id SERIAL PRIMARY KEY, usuario TEXT, mensagem TEXT, resposta TEXT, numero TEXT, is_reply BOOLEAN DEFAULT FALSE, mensagem_original TEXT, humor TEXT DEFAULT 'neutro', modo_resposta TEXT DEFAULT 'normal', nivel_transicao INTEGER DEFAULT 1, usuario_privilegiado BOOLEAN DEFAULT FALSE, modelo_usado TEXT DEFAULT 'desconhecido', conversation_id TEXT DEFAULT '', nome_usuario TEXT DEFAULT '', message_id TEXT UNIQUE, created_at TIMESTAMP DEFAULT NOW() ) """, "mensagens") exec_ddl(""" CREATE TABLE IF NOT EXISTS usuarios_privilegiados ( id SERIAL PRIMARY KEY, numero TEXT UNIQUE, nome TEXT, apelido TEXT, modo_fala TEXT, codigo_verificacao TEXT, ativo BOOLEAN DEFAULT TRUE, privilegio_temporario_ativo BOOLEAN DEFAULT FALSE, expira_em DOUBLE PRECISION, created_at TIMESTAMP DEFAULT NOW() ) """, "usuarios_privilegiados") exec_ddl(""" CREATE TABLE IF NOT EXISTS embeddings ( id SERIAL PRIMARY KEY, numero_usuario TEXT, source_type TEXT, texto TEXT, embedding BYTEA ) """, "embeddings") exec_ddl(""" CREATE TABLE IF NOT EXISTS aprendizados ( id SERIAL PRIMARY KEY, numero_usuario TEXT, chave TEXT, valor TEXT, created_at TIMESTAMP DEFAULT NOW(), UNIQUE(numero_usuario, chave) ) """, "aprendizados") exec_ddl(""" CREATE TABLE IF NOT EXISTS girias_aprendidas ( id SERIAL PRIMARY KEY, numero_usuario TEXT, giria TEXT, significado TEXT, contexto TEXT, frequencia INTEGER DEFAULT 1, created_at TIMESTAMP DEFAULT NOW(), updated_at TIMESTAMP DEFAULT NOW(), UNIQUE(numero_usuario, giria) ) """, "girias_aprendidas") exec_ddl(""" CREATE TABLE IF NOT EXISTS vocabulario_autonomo ( id SERIAL PRIMARY KEY, termo TEXT NOT NULL, significado_inferido TEXT, categoria VARCHAR(50) DEFAULT 'desconhecida', contexto_original TEXT, frequencia INTEGER DEFAULT 1, ultimo_uso TIMESTAMP DEFAULT NOW(), confianca FLOAT DEFAULT 0.0, usuario_origem TEXT DEFAULT '', grupo_origem TEXT DEFAULT '', exemplos_uso JSONB DEFAULT '[]'::jsonb, sentimento_associado VARCHAR(50) DEFAULT 'neutro', quando_usar TEXT DEFAULT '', dimensao VARCHAR(50) DEFAULT 'vocabulario', padrao_json JSONB DEFAULT '{}'::jsonb, revisado BOOLEAN DEFAULT FALSE, created_at TIMESTAMP DEFAULT NOW(), updated_at TIMESTAMP DEFAULT NOW() ) """, "vocabulario_autonomo") # Migrações ALTER TABLE (idempotentes) exec_ddl("ALTER TABLE vocabulario_autonomo ADD COLUMN IF NOT EXISTS dimensao VARCHAR(50) DEFAULT 'vocabulario'", "ALTER vocab dimensao") exec_ddl("ALTER TABLE vocabulario_autonomo ADD COLUMN IF NOT EXISTS padrao_json JSONB DEFAULT '{}'::jsonb", "ALTER vocab padrao_json") # Índices (DROP IF EXISTS pode falhar em PG antigo, isolar) exec_ddl("DROP INDEX IF EXISTS idx_vocab_termo", "DROP idx_vocab_termo") exec_ddl("CREATE INDEX IF NOT EXISTS idx_vocab_termo_dim ON vocabulario_autonomo(termo, dimensao)", "CREATE idx_vocab_termo_dim") exec_ddl("CREATE INDEX IF NOT EXISTS idx_vocab_usuario ON vocabulario_autonomo(usuario_origem)", "CREATE idx_vocab_usuario") exec_ddl("CREATE INDEX IF NOT EXISTS idx_vocab_dimensao ON vocabulario_autonomo(dimensao)", "CREATE idx_vocab_dimensao") # UNIQUE indexes - podem falhar se houver duplicatas, isolar cada um exec_ddl("CREATE UNIQUE INDEX IF NOT EXISTS idx_aprendizados_user_chave ON aprendizados(numero_usuario, chave)", "UNIQUE idx_aprendizados") exec_ddl("CREATE UNIQUE INDEX IF NOT EXISTS idx_girias_user_giria ON girias_aprendidas(numero_usuario, giria)", "UNIQUE idx_girias") exec_ddl(""" CREATE TABLE IF NOT EXISTS tom_usuario ( id SERIAL PRIMARY KEY, numero_usuario TEXT, tom_detectado TEXT, intensidade REAL DEFAULT 0.5, contexto TEXT, humor TEXT DEFAULT 'neutro', created_at TIMESTAMP DEFAULT NOW() ) """, "tom_usuario") exec_ddl(""" CREATE TABLE IF NOT EXISTS contexto ( user_key TEXT PRIMARY KEY, historico TEXT, emocao_atual TEXT, humor_atual TEXT DEFAULT 'neutro', modo_resposta TEXT DEFAULT 'normal', nivel_transicao INTEGER DEFAULT 1, usuario_privilegiado BOOLEAN DEFAULT FALSE, termos TEXT, girias TEXT, tom TEXT, updated_at TIMESTAMP DEFAULT NOW() ) """, "contexto") exec_ddl(""" CREATE TABLE IF NOT EXISTS pronomes_por_tom ( tom TEXT PRIMARY KEY, pronomes TEXT ) """, "pronomes_por_tom") exec_ddl(""" CREATE TABLE IF NOT EXISTS persona_usuario ( numero_usuario TEXT PRIMARY KEY, personalidade TEXT, nome TEXT DEFAULT '', vicios_linguagem TEXT, gostos TEXT, desgostos TEXT, emocional TEXT, created_at TIMESTAMP DEFAULT NOW(), updated_at TIMESTAMP DEFAULT NOW() ) """, "persona_usuario") exec_ddl(""" CREATE TABLE IF NOT EXISTS lstm_contexto ( context_id VARCHAR(255) NOT NULL, numero_usuario VARCHAR(50) NOT NULL, topic_principal VARCHAR(255), subtopicas JSONB, conversation_path JSONB, interaction_pattern VARCHAR(50), emotional_state VARCHAR(50), unanswered_questions JSONB, assumed_knowledge JSONB, last_key_message TEXT, context_switches INTEGER DEFAULT 0, contradictions JSONB, created_at TIMESTAMP DEFAULT NOW(), last_updated TIMESTAMP DEFAULT NOW(), metadata JSONB, PRIMARY KEY (context_id, numero_usuario) ) """, "lstm_contexto") exec_ddl(""" CREATE TABLE IF NOT EXISTS lstm_message_links ( id SERIAL PRIMARY KEY, context_id VARCHAR(255) NOT NULL, message_id VARCHAR(255) NOT NULL, numero_usuario VARCHAR(50) NOT NULL, speaker_name VARCHAR(255), parent_message_id VARCHAR(255), topic_changed BOOLEAN DEFAULT FALSE, context_switch_type VARCHAR(50), relevance_score DOUBLE PRECISION DEFAULT 0.0, created_at TIMESTAMP DEFAULT NOW(), UNIQUE(context_id, message_id, numero_usuario), FOREIGN KEY (context_id, numero_usuario) REFERENCES lstm_contexto(context_id, numero_usuario) ON DELETE CASCADE ) """, "lstm_message_links") exec_ddl("CREATE INDEX IF NOT EXISTS idx_lstm_usuario ON lstm_contexto(numero_usuario)", "idx_lstm_usuario") exec_ddl("CREATE INDEX IF NOT EXISTS idx_lstm_created ON lstm_contexto(created_at)", "idx_lstm_created") exec_ddl("CREATE INDEX IF NOT EXISTS idx_lstm_msg_context ON lstm_message_links(context_id)", "idx_lstm_msg_context") exec_ddl("CREATE INDEX IF NOT EXISTS idx_lstm_msg_message ON lstm_message_links(message_id)", "idx_lstm_msg_message") # ✅ FIX 2026-07-28: tabela de gírias angolanas para injeção dinâmica no prompt # Lista hardcoded eliminada — todas as gírias vivem aqui para crescer/adicionar sem mexer em código. exec_ddl(""" CREATE TABLE IF NOT EXISTS girias ( id SERIAL PRIMARY KEY, giria TEXT NOT NULL UNIQUE, intensidade INT NOT NULL CHECK (intensidade BETWEEN 1 AND 5), tom_adequado TEXT[] NOT NULL DEFAULT '{}', exemplo_uso TEXT, ativo BOOLEAN NOT NULL DEFAULT TRUE, created_at TIMESTAMP DEFAULT NOW() ) """, "girias_table") exec_ddl("CREATE INDEX IF NOT EXISTS idx_girias_tom ON girias USING GIN(tom_adequado) WHERE ativo", "idx_girias_tom") exec_ddl(""" INSERT INTO girias (giria, intensidade, tom_adequado, exemplo_uso) VALUES ('cassules', 4, ARRAY['casual', 'aggressive'], 'Cassules, isso é trabalho sério.'), ('puto', 5, ARRAY['aggressive'], 'Puto, foda-se esse bug.'), ('mambo', 3, ARRAY['casual', 'playful'], 'Que mambo é esse?'), ('mano', 1, ARRAY['formal', 'casual'], 'Mano, vê se tu consegue.'), ('mana', 1, ARRAY['formal', 'casual'], 'Mana, vê se tu consegue.'), ('parceiro', 1, ARRAY['formal', 'casual'], 'Parceiro, isto é formal.'), ('parceira', 1, ARRAY['formal', 'casual'], 'Parceira, deixa comigo.'), ('fera', 2, ARRAY['casual', 'playful'], 'Fera, isto aqui resolve.'), ('cria', 2, ARRAY['casual', 'playful'], 'Cria, escuta só isto.'), ('kenga', 4, ARRAY['aggressive', 'playful'], 'Kenga, bora despachar.'), ('bana', 2, ARRAY['casual'], 'Bana, sabes que sim.'), ('kala', 3, ARRAY['casual', 'playful'], 'Kala, não vou esquecer.'), ('kwasa', 4, ARRAY['aggressive'], 'Kwasa, isto é hard.'), ('tio', 1, ARRAY['formal', 'casual'], 'Tio, presta atenção.'), ('tia', 1, ARRAY['formal', 'casual'], 'Tia, deixa eu explicar.'), ('velho', 2, ARRAY['casual'], 'Velho, tu sabes.'), ('patrão', 2, ARRAY['formal'], 'Patrão, bora despachar.') ON CONFLICT (giria) DO UPDATE SET exemplo_uso = EXCLUDED.exemplo_uso, tom_adequado = EXCLUDED.tom_adequado, intensidade = EXCLUDED.intensidade """, "INSERT girias_seed") exec_ddl(""" INSERT INTO pronomes_por_tom (tom, pronomes) VALUES ('neutro', 'tu/você'), ('formal', 'o senhor/a senhora'), ('informal', 'tu/você'), ('tecnico_formal', 'senhor') ON CONFLICT (tom) DO NOTHING """, "INSERT pronomes_por_tom") exec_ddl(""" INSERT INTO usuarios_privilegiados (numero, nome, apelido, modo_fala) VALUES ('244937035662', 'Isaac Quarenta', 'Isaac', 'tecnico_formal'), ('244978787009', 'Isaac Quarenta 2', 'Isaac', 'tecnico_formal') ON CONFLICT (numero) DO NOTHING """, "INSERT usuarios_privilegiados") exec_ddl(""" CREATE TABLE IF NOT EXISTS user_emotional_profiles ( id SERIAL PRIMARY KEY, user_id TEXT UNIQUE NOT NULL, numero_usuario TEXT, profile_data TEXT NOT NULL, created_at TIMESTAMP DEFAULT NOW(), updated_at TIMESTAMP DEFAULT NOW() ) """, "user_emotional_profiles") exec_ddl(""" CREATE TABLE IF NOT EXISTS continuous_learning ( id SERIAL PRIMARY KEY, ts DOUBLE PRECISION NOT NULL, usuario TEXT, numero TEXT, nome_usuario TEXT, tipo_conversa TEXT DEFAULT 'pv', mensagem TEXT, resposta_do_bot BOOLEAN DEFAULT FALSE, resposta_gerada TEXT, is_reply BOOLEAN DEFAULT FALSE, reply_to_bot BOOLEAN DEFAULT FALSE, contexto_grupo TEXT, modelo_usado TEXT DEFAULT 'desconhecido', message_id TEXT UNIQUE, qualidade DOUBLE PRECISION DEFAULT 0.0, tipo_conteudo TEXT, created_at TIMESTAMP DEFAULT NOW() ) """, "continuous_learning") exec_ddl("CREATE INDEX IF NOT EXISTS idx_cl_user ON continuous_learning(usuario)", "idx_cl_user") exec_ddl("CREATE INDEX IF NOT EXISTS idx_cl_ts ON continuous_learning(ts)", "idx_cl_ts") exec_ddl("CREATE INDEX IF NOT EXISTS idx_cl_quality ON continuous_learning(qualidade)", "idx_cl_quality") exec_ddl("CREATE INDEX IF NOT EXISTS idx_cl_msgid ON continuous_learning(message_id)", "idx_cl_msgid") exec_ddl(""" CREATE TABLE IF NOT EXISTS finetuning_examples ( id SERIAL PRIMARY KEY, user_id TEXT NOT NULL, conversation_id TEXT NOT NULL, input_message TEXT NOT NULL, expected_response TEXT NOT NULL, actual_response TEXT, quality_score INT DEFAULT 50, tone_level VARCHAR(50), hostility_score INT DEFAULT 0, embedding_vector BYTEA, embedding_input BYTEA, embedding_output BYTEA, similarity_score FLOAT DEFAULT 0.0, emotion_label VARCHAR(50), created_at TIMESTAMP DEFAULT NOW(), updated_at TIMESTAMP DEFAULT NOW(), indexed BOOLEAN DEFAULT FALSE ) """, "finetuning_examples") exec_ddl("CREATE INDEX IF NOT EXISTS idx_ft_user ON finetuning_examples(user_id)", "idx_ft_user") exec_ddl("CREATE INDEX IF NOT EXISTS idx_ft_quality ON finetuning_examples(quality_score)", "idx_ft_quality") exec_ddl("CREATE INDEX IF NOT EXISTS idx_ft_tone ON finetuning_examples(tone_level)", "idx_ft_tone") exec_ddl("CREATE INDEX IF NOT EXISTS idx_ft_emotion ON finetuning_examples(emotion_label)", "idx_ft_emotion") exec_ddl(""" CREATE TABLE IF NOT EXISTS training_metrics ( id SERIAL PRIMARY KEY, training_session_id TEXT UNIQUE NOT NULL, examples_used INT, avg_quality FLOAT, model_accuracy FLOAT, embedding_loss FLOAT, emotion_accuracy FLOAT, loss FLOAT, created_at TIMESTAMP DEFAULT NOW(), weights_checkpoint BYTEA, embedding_weights BYTEA, status VARCHAR(50) ) """, "training_metrics") exec_ddl(""" CREATE TABLE IF NOT EXISTS training_cycles ( id SERIAL PRIMARY KEY, cycle_number INT, cycle_type VARCHAR(50), started_at TIMESTAMP DEFAULT NOW(), completed_at TIMESTAMP, examples_processed INT, improvement_pct FLOAT, embedding_improvement FLOAT, emotion_improvement FLOAT, status VARCHAR(50) ) """, "training_cycles") exec_ddl("CREATE INDEX IF NOT EXISTS idx_emb_user ON embeddings(numero_usuario)", "idx_emb_user") exec_ddl("CREATE INDEX IF NOT EXISTS idx_aprend_user ON aprendizados(numero_usuario)", "idx_aprend_user") exec_ddl("CREATE INDEX IF NOT EXISTS idx_aprend_chave ON aprendizados(chave)", "idx_aprend_chave") exec_ddl(""" CREATE TABLE IF NOT EXISTS akira_timers ( id SERIAL PRIMARY KEY, fire_at TIMESTAMP NOT NULL, message TEXT NOT NULL, chat_context TEXT DEFAULT '', fired BOOLEAN DEFAULT FALSE, created_at TIMESTAMP DEFAULT NOW() ) """, "akira_timers") exec_ddl("CREATE INDEX IF NOT EXISTS idx_timers_fired ON akira_timers(fired, fire_at)", "idx_timers_fired") _pool_return(conn) logger.info("Tabelas PostgreSQL criadas/garantidas") def _init_context_isolation_tables(self): try: conn = self._get_connection() conn.autocommit = True cur = conn.cursor() def exec_ddl(sql: str, desc: str = "") -> bool: try: cur.execute(sql) return True except Exception as e: try: conn.rollback() except Exception: pass logger.debug(f"DDL {desc} falhou (ignorado): {e}") return False exec_ddl(""" CREATE TABLE IF NOT EXISTS contextos_isolados ( context_id TEXT PRIMARY KEY, numero_usuario TEXT NOT NULL, grupo_id TEXT, tipo_conversa TEXT DEFAULT 'pv', estado_emocional TEXT DEFAULT 'neutral', nivel_intimidade INTEGER DEFAULT 1, short_memory TEXT DEFAULT '[]', metadata TEXT DEFAULT '{}', created_at DOUBLE PRECISION DEFAULT EXTRACT(EPOCH FROM NOW()), last_interaction DOUBLE PRECISION DEFAULT EXTRACT(EPOCH FROM NOW()) ) """, "contextos_isolados") exec_ddl("CREATE INDEX IF NOT EXISTS idx_contextos_user ON contextos_isolados(numero_usuario)", "idx_contextos_user") exec_ddl("CREATE INDEX IF NOT EXISTS idx_contextos_tipo ON contextos_isolados(tipo_conversa)", "idx_contextos_tipo") except Exception as e: logger.warning(f"Erro ao criar contextos_isolados: {e}") finally: if conn: _pool_return(conn) def _ensure_all_columns_and_indexes(self): """Garante que todas as colunas e índices existam (migrations PostgreSQL).""" conn = self._get_connection() conn.autocommit = True cur = conn.cursor() def exec_ddl(sql: str, desc: str = "") -> bool: try: cur.execute(sql) return True except Exception as e: try: conn.rollback() except Exception: pass logger.debug(f"DDL {desc} falhou (ignorado): {e}") return False # Lista de colunas a adicionar: (tabela, coluna, tipo) migrations = [ ('mensagens', 'humor', "TEXT DEFAULT 'neutro'"), ('mensagens', 'modo_resposta', "TEXT DEFAULT 'normal'"), ('mensagens', 'nivel_transicao', "INTEGER DEFAULT 1"), ('mensagens', 'usuario_privilegiado', "BOOLEAN DEFAULT FALSE"), ('mensagens', 'modelo_usado', "TEXT DEFAULT 'desconhecido'"), ('mensagens', 'conversation_id', "TEXT DEFAULT ''"), ('mensagens', 'nome_usuario', "TEXT DEFAULT ''"), ('tom_usuario', 'humor', "TEXT DEFAULT 'neutro'"), ('contexto', 'humor_atual', "TEXT DEFAULT 'neutro'"), ('contexto', 'modo_resposta', "TEXT DEFAULT 'normal'"), ('contexto', 'nivel_transicao', "INTEGER DEFAULT 1"), ('contexto', 'usuario_privilegiado', "BOOLEAN DEFAULT FALSE"), ('usuarios_privilegiados', 'privilegio_temporario_ativo', "BOOLEAN DEFAULT FALSE"), ('usuarios_privilegiados', 'expira_em', "DOUBLE PRECISION"), ('persona_usuario', 'nome', "TEXT DEFAULT ''"), ] for table, col, col_type in migrations: exec_ddl(f"ALTER TABLE {table} ADD COLUMN IF NOT EXISTS {col} {col_type}", f"ALTER {table}.{col}") # Índices adicionais indexes = [ ("CREATE INDEX IF NOT EXISTS idx_emb_user ON embeddings(numero_usuario)", "idx_emb_user"), ("CREATE INDEX IF NOT EXISTS idx_aprend_user ON aprendizados(numero_usuario)", "idx_aprend_user"), ("CREATE INDEX IF NOT EXISTS idx_aprend_chave ON aprendizados(chave)", "idx_aprend_chave"), ] for sql, desc in indexes: exec_ddl(sql, desc) logger.debug("Migrations PG concluídas") _pool_return(conn) # ================================================================ # MÉTODOS PÚBLICOS (mesma interface do Database SQLite) # ================================================================ def adicionar_usuario_privilegiado(self, numero, nome, apelido, modo_fala="tecnico_formal"): try: codigo = str(random.randint(100000, 999999)) self._execute_with_retry( """INSERT INTO usuarios_privilegiados (numero, nome, apelido, modo_fala, codigo_verificacao) VALUES (%s, %s, %s, %s, %s) ON CONFLICT (numero) DO UPDATE SET nome=EXCLUDED.nome, apelido=EXCLUDED.apelido, modo_fala=EXCLUDED.modo_fala, codigo_verificacao=EXCLUDED.codigo_verificacao""", (numero, nome, apelido, modo_fala, codigo), commit=True ) return True, codigo except Exception as e: logger.error(f"Erro ao adicionar privilegiado: {e}") return False, str(e) def eh_privilegiado(self, numero): try: rows = self._execute_with_retry( "SELECT ativo FROM usuarios_privilegiados WHERE numero = %s AND ativo = TRUE", (numero,) ) return rows is not None and len(rows) > 0 except: return False def verificar_privilegios_usuario(self, numero): try: rows = self._execute_with_retry( "SELECT ativo, privilegio_temporario_ativo, expira_em FROM usuarios_privilegiados WHERE numero = %s", (numero,) ) if rows: r = rows[0] ativo = r['ativo'] if isinstance(r, dict) else r[0] temp_ativo = r['privilegio_temporario_ativo'] if isinstance(r, dict) else r[1] expira = r['expira_em'] if isinstance(r, dict) else r[2] return {'ativo': ativo, 'temporario_ativo': temp_ativo, 'expira_em': expira} return {'ativo': False, 'temporario_ativo': False, 'expira_em': None} except: return {'ativo': False, 'temporario_ativo': False, 'expira_em': None} def verificar_codigo(self, numero, codigo): try: rows = self._execute_with_retry( "SELECT codigo_verificacao FROM usuarios_privilegiados WHERE numero = %s", (numero,) ) if rows: r = rows[0] cod = r['codigo_verificacao'] if isinstance(r, dict) else r[0] return str(cod) == str(codigo) return False except: return False def obter_modo_fala_privilegiado(self, numero): try: rows = self._execute_with_retry( "SELECT modo_fala FROM usuarios_privilegiados WHERE numero = %s", (numero,) ) if rows: r = rows[0] return r['modo_fala'] if isinstance(r, dict) else r[0] return None except: return None def salvar_mensagem(self, usuario, mensagem, resposta, numero=None, is_reply=False, mensagem_original=None, humor="neutro", modo_resposta="normal", nivel_transicao=1, usuario_privilegiado=False, modelo_usado="desconhecido", **kwargs): try: cols = ['usuario', 'mensagem', 'resposta', 'humor', 'modo_resposta', 'nivel_transicao', 'usuario_privilegiado', 'is_reply', 'modelo_usado'] vals = [usuario, mensagem, resposta, humor, modo_resposta, nivel_transicao, usuario_privilegiado, is_reply, modelo_usado] message_id = kwargs.get('message_id') if message_id: cols.append('message_id') vals.append(message_id) conversation_id = kwargs.get('conversation_id') if conversation_id: cols.append('conversation_id') vals.append(conversation_id) if numero: cols.append('numero') vals.append(numero) if mensagem_original: cols.append('mensagem_original') vals.append(mensagem_original) nome_usuario = kwargs.get('nome_usuario') or usuario if nome_usuario: cols.append('nome_usuario') vals.append(nome_usuario) placeholders = ', '.join(['%s'] * len(cols)) # ✅ FIX #3-CAMADA: ON CONFLICT DO UPDATE ao invés de DO NOTHING # Motivo: ON CONFLICT DO NOTHING falha SILENCIOSAMENTE em duplicatas # Resultado: A tentativa é registrada sem erro, causando corridas de dedup # Solução: ON CONFLICT (message_id) DO UPDATE SET + logging explícito if message_id: # Build UPDATE clause for all columns except message_id (PK) update_cols = [col for col in cols if col != 'message_id'] set_clause = ', '.join([f"{col} = EXCLUDED.{col}" for col in update_cols]) # RETURNING xmax: xmax=0 significa INSERT real, xmax>0 significa UPDATE (upsert) query = f"INSERT INTO mensagens ({', '.join(cols)}) VALUES ({placeholders}) ON CONFLICT (message_id) DO UPDATE SET {set_clause} RETURNING xmax" else: # Fallback: se não houver message_id, não faz update (compatibilidade) query = f"INSERT INTO mensagens ({', '.join(cols)}) VALUES ({placeholders}) ON CONFLICT DO NOTHING" try: conn = self._get_connection() cur = conn.cursor() cur.execute(self._convert_query(query), tuple(vals)) if message_id: row = cur.fetchone() conn.commit() _pool_return(conn) # xmax=0 → INSERT real; xmax>0 → UPDATE (registro já existia) if row is not None: xmax = row[0] if not isinstance(row, dict) else row.get('xmax', 0) if int(xmax) == 0: logger.info(f"✅ [DB INSERT OK] message_id={message_id} | usuario={usuario} | modelo={modelo_usado}") else: logger.debug(f"🔄 [DB UPSERT] message_id={message_id} já existia — actualizado (sem duplicado)") else: logger.debug(f"🔄 [DB UPSERT SKIP] message_id={message_id} sem retorno (ON CONFLICT DO NOTHING activado)") else: conn.commit() _pool_return(conn) return True except Exception as db_err: # ❌ Log de falha com contexto completo logger.error(f"❌ [DB INSERT FAIL] Erro ao salvar mensagem: {db_err} | message_id={message_id} | usuario={usuario}") try: conn.rollback() _pool_return(conn) except Exception: pass return False except Exception as e: logger.warning(f"Erro salvar_mensagem (outer): {e}") return False def recuperar_mensagens(self, usuario, limite=5): try: result = self._execute_with_retry( """SELECT mensagem, resposta FROM mensagens WHERE usuario=%s OR numero=%s ORDER BY id DESC LIMIT %s""", (usuario, usuario, limite) ) if not result: return [] return [(row['mensagem'], row['resposta']) for row in result] except: return [] def recuperar_mensagens_por_contexto(self, context_id, limite=50): try: rows = self._execute_with_retry( "SELECT usuario, mensagem, resposta, created_at FROM mensagens WHERE conversation_id = %s ORDER BY id DESC LIMIT %s", (context_id, limite) ) if not rows: return [] return [dict(row) for row in rows] except: return [] def recuperar_todas_mensagens(self, usuario, limite: int = 500, offset: int = 0): """Retorna mensagens paginadas para treino batch.""" try: rows = self._execute_with_retry( """SELECT mensagem, resposta, modelo_usado, humor, created_at FROM mensagens WHERE usuario=%s OR numero=%s ORDER BY id DESC LIMIT %s OFFSET %s""", (usuario, usuario, limite, offset) ) if not rows: return [] return [dict(row) for row in rows] except: return [] def contar_mensagens_usuario(self, usuario) -> int: """Total de mensagens de um utilizador para planeamento de treino.""" try: rows = self._execute_with_retry( "SELECT COUNT(*) as total FROM mensagens WHERE usuario=%s OR numero=%s", (usuario, usuario) ) if rows: r = rows[0] return r['total'] if isinstance(r, dict) else int(r[0]) return 0 except: return 0 def recuperar_humor(self, numero_usuario): try: rows = self._execute_with_retry( "SELECT humor FROM contexto WHERE user_key = %s", (numero_usuario,) ) if rows: r = rows[0] return r['humor'] if isinstance(r, dict) else r[0] return "neutro" except: return "neutro" def salvar_contexto(self, user_key, historico, emocao_atual="neutro", humor_atual="neutro", modo_resposta="normal", nivel_transicao=1, usuario_privilegiado=False, termos=None, girias=None, tom=None): try: self._execute_with_retry( """INSERT INTO contexto (user_key, historico, emocao_atual, humor_atual, modo_resposta, nivel_transicao, usuario_privilegiado, termos, girias, tom, updated_at) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, NOW()) ON CONFLICT (user_key) DO UPDATE SET historico=EXCLUDED.historico, emocao_atual=EXCLUDED.emocao_atual, humor_atual=EXCLUDED.humor_atual, modo_resposta=EXCLUDED.modo_resposta, nivel_transicao=EXCLUDED.nivel_transicao, usuario_privilegiado=EXCLUDED.usuario_privilegiado, termos=EXCLUDED.termos, girias=EXCLUDED.girias, tom=EXCLUDED.tom, updated_at=NOW()""", (user_key, historico, emocao_atual, humor_atual, modo_resposta, nivel_transicao, usuario_privilegiado, termos, girias, tom), commit=True ) return True except: return False def recuperar_contexto(self, user_key): try: rows = self._execute_with_retry("SELECT * FROM contexto WHERE user_key = %s", (user_key,)) if rows: return dict(rows[0]) return None except: return None def registrar_tom_usuario(self, numero_usuario, tom_detectado, intensidade=0.5, contexto="", humor="neutro"): try: self._execute_with_retry( """INSERT INTO tom_usuario (numero_usuario, tom_detectado, intensidade, contexto, humor) VALUES (%s, %s, %s, %s, %s)""", (numero_usuario, tom_detectado, intensidade, contexto, humor), commit=True ) return True except: return False def obter_tom_predominante(self, numero_usuario): try: rows = self._execute_with_retry( """SELECT tom_detectado, COUNT(*) as cnt FROM tom_usuario WHERE numero_usuario=%s GROUP BY tom_detectado ORDER BY cnt DESC LIMIT 1""", (numero_usuario,) ) if rows: r = rows[0] return r['tom_detectado'] if isinstance(r, dict) else r[0] return None except: return None def salvar_aprendizado_detalhado(self, numero_usuario, chave, valor): try: self._execute_with_retry( """INSERT INTO aprendizados (numero_usuario, chave, valor) VALUES (%s, %s, %s) ON CONFLICT ON CONSTRAINT aprendizados_numero_usuario_chave_key DO UPDATE SET valor = EXCLUDED.valor""", (numero_usuario, chave, valor), commit=True ) return True except Exception: # Fallback: se não existe unique constraint, tenta sem upsert try: self._execute_with_retry( """INSERT INTO aprendizados (numero_usuario, chave, valor) VALUES (%s, %s, %s)""", (numero_usuario, chave, valor), commit=True ) return True except Exception: return False def recuperar_aprendizado_detalhado(self, numero_usuario, chave=None): try: if chave: rows = self._execute_with_retry( "SELECT valor FROM aprendizados WHERE numero_usuario=%s AND chave=%s", (numero_usuario, chave) ) if rows: r = rows[0] return r['valor'] if isinstance(r, dict) else r[0] return None else: rows = self._execute_with_retry( "SELECT chave, valor FROM aprendizados WHERE numero_usuario=%s ORDER BY id DESC LIMIT 20", (numero_usuario,) ) if not rows: return {} return {r['chave']: r['valor'] for r in rows} except: return {} if not chave else None def salvar_giria_aprendida(self, numero_usuario, giria, significado, contexto=""): try: self._execute_with_retry( """INSERT INTO girias_aprendidas (numero_usuario, giria, significado, contexto) VALUES (%s, %s, %s, %s) ON CONFLICT DO NOTHING""", (numero_usuario, giria, significado, contexto), commit=True ) return True except: return False def recuperar_girias_usuario(self, numero_usuario): try: rows = self._execute_with_retry( "SELECT giria, significado, frequencia FROM girias_aprendidas WHERE numero_usuario=%s ORDER BY id DESC LIMIT 10", (numero_usuario,) ) if not rows: return [] return [{'giria': r['giria'], 'significado': r['significado'], 'frequencia': r.get('frequencia', 1)} for r in rows] except: return [] def salvar_termo_autonomo(self, termo, significado, categoria, contexto_original, usuario_origem='', grupo_origem='', confianca=0.0, sentimento='neutro', quando_usar='', dimensao='vocabulario', padrao_json=None): try: import json as _json if padrao_json is None: padrao_json = {} existing = self._execute_with_retry( "SELECT id, frequencia, exemplos_uso FROM vocabulario_autonomo WHERE termo=%s AND dimensao=%s", (termo, dimensao) ) if existing: row = existing[0] exemplos = row.get('exemplos_uso', []) if isinstance(exemplos, str): exemplos = _json.loads(exemplos) novo_exemplo = { 'contexto': contexto_original[:200], 'usuario': usuario_origem, 'grupo': grupo_origem } if not any(e.get('contexto') == novo_exemplo['contexto'] for e in exemplos): exemplos.append(novo_exemplo) self._execute_with_retry( """UPDATE vocabulario_autonomo SET frequencia = frequencia + 1, ultimo_uso = NOW(), exemplos_uso = %s::jsonb, padrao_json = %s::jsonb, updated_at = NOW() WHERE id = %s""", (_json.dumps(exemplos), _json.dumps(padrao_json), row['id']), commit=True ) else: exemplos = _json.dumps([{ 'contexto': contexto_original[:200], 'usuario': usuario_origem, 'grupo': grupo_origem }]) self._execute_with_retry( """INSERT INTO vocabulario_autonomo (termo, significado_inferido, categoria, contexto_original, usuario_origem, grupo_origem, confianca, sentimento_associado, quando_usar, exemplos_uso, dimensao, padrao_json) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s::jsonb, %s, %s::jsonb)""", (termo, significado, categoria, contexto_original, usuario_origem, grupo_origem, confianca, sentimento, quando_usar, exemplos, dimensao, _json.dumps(padrao_json)), commit=True ) return True except Exception as e: import traceback traceback.print_exc() return False def recuperar_vocabulario_autonomo(self, limite=50, confianca_min=0.3, categoria=None, dimensao=None): try: clausulas = ["confianca >= %s"] params = [confianca_min] if categoria: clausulas.append("categoria = %s") params.append(categoria) if dimensao: clausulas.append("dimensao = %s") params.append(dimensao) where = " AND ".join(clausulas) params.append(limite) rows = self._execute_with_retry( f"""SELECT termo, significado_inferido, categoria, frequencia, confianca, sentimento_associado, quando_usar, dimensao FROM vocabulario_autonomo WHERE {where} ORDER BY frequencia DESC, confianca DESC LIMIT %s""", tuple(params) ) if not rows: return [] return [{ 'termo': r['termo'], 'significado': r['significado_inferido'], 'categoria': r['categoria'], 'frequencia': r['frequencia'], 'confianca': float(r['confianca']), 'sentimento': r['sentimento_associado'], 'quando_usar': r['quando_usar'], 'dimensao': r['dimensao'] } for r in rows] except: return [] def buscar_termo_autonomo(self, termo, dimensao='vocabulario'): try: rows = self._execute_with_retry( """SELECT * FROM vocabulario_autonomo WHERE termo = %s AND dimensao = %s""", (termo, dimensao) ) if rows: return dict(rows[0]) return None except: return None def salvar_embedding(self, numero_usuario, source_type, texto, embedding): try: import numpy as np if isinstance(embedding, np.ndarray): embedding = embedding.tobytes() self._execute_with_retry( """INSERT INTO embeddings (numero_usuario, source_type, texto, embedding) VALUES (%s, %s, %s, %s)""", (numero_usuario, source_type, texto, embedding), commit=True ) return True except: return False def recuperar_embeddings(self, numero_usuario, limite: int = 100, offset: int = 0): try: rows = self._execute_with_retry( "SELECT source_type, texto, embedding FROM embeddings WHERE numero_usuario=%s ORDER BY id DESC LIMIT %s OFFSET %s", (numero_usuario, limite, offset) ) if not rows: return [] import numpy as np results = [] for r in rows: emb = r['embedding'] if isinstance(emb, memoryview): emb = bytes(emb) elif isinstance(emb, (bytes, bytearray)): pass else: emb = bytes(emb) results.append({ 'source_type': r['source_type'], 'texto': r['texto'], 'embedding': np.frombuffer(emb, dtype=np.float32) }) return results except: return [] def atualizar_persona(self, numero_usuario, campos): if not campos: return False try: self._execute_with_retry( "INSERT INTO persona_usuario (numero_usuario) VALUES (%s) ON CONFLICT DO NOTHING", (numero_usuario,), commit=True ) ALLOWED_PERSONA_COLS = { 'personalidade', 'nome', 'vicios_linguagem', 'gostos', 'desgostos', 'emocional' } parts = [] vals = [] for k, v in campos.items(): if k not in ALLOWED_PERSONA_COLS: continue parts.append(f"{k} = %s") vals.append(v) if not parts: return False parts.append("updated_at = NOW()") vals.append(numero_usuario) self._execute_with_retry( f"UPDATE persona_usuario SET {', '.join(parts)} WHERE numero_usuario = %s", tuple(vals), commit=True ) return True except: return False def recuperar_persona(self, numero_usuario): try: rows = self._execute_with_retry("SELECT * FROM persona_usuario WHERE numero_usuario=%s", (numero_usuario,)) if rows: return dict(rows[0]) return {} except: return {} def salvar_contexto_isolado(self, context_data): try: self._execute_with_retry( """INSERT INTO contextos_isolados (context_id, numero_usuario, grupo_id, tipo_conversa, estado_emocional, nivel_intimidade, short_memory, metadata, created_at, last_interaction) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s) ON CONFLICT (context_id) DO UPDATE SET estado_emocional=EXCLUDED.estado_emocional, nivel_intimidade=EXCLUDED.nivel_intimidade, short_memory=EXCLUDED.short_memory, metadata=EXCLUDED.metadata, last_interaction=EXCLUDED.last_interaction""", (context_data.get('context_id'), context_data.get('numero_usuario'), context_data.get('grupo_id'), context_data.get('tipo_conversa', 'pv'), context_data.get('estado_emocional', 'neutral'), context_data.get('nivel_intimidade', 1), json.dumps(context_data.get('short_memory', [])), json.dumps(context_data.get('metadata', {})), context_data.get('created_at', time.time()), context_data.get('last_interaction', time.time())), commit=True ) return True except: return False def recuperar_contexto_isolado(self, context_id): try: rows = self._execute_with_retry("SELECT * FROM contextos_isolados WHERE context_id = %s", (context_id,)) if rows: row = dict(rows[0]) row['short_memory'] = _safe_json_load(row.get('short_memory')) or [] row['metadata'] = _safe_json_load(row.get('metadata')) or {} return row return None except: return None def deletar_contexto_isolado(self, context_id): try: self._execute_with_retry("DELETE FROM contextos_isolados WHERE context_id = %s", (context_id,), commit=True) return True except: return False def listar_contextos_usuario(self, numero_usuario): try: rows = self._execute_with_retry( "SELECT * FROM contextos_isolados WHERE numero_usuario = %s ORDER BY last_interaction DESC", (numero_usuario,) ) if not rows: return [] results = [] for r in rows: d = dict(r) d['short_memory'] = _safe_json_load(d.get('short_memory')) or [] d['metadata'] = _safe_json_load(d.get('metadata')) or {} results.append(d) return results except: return [] def recuperar_historico(self, usuario="", numero="", conversation_id="", limite=10): try: if conversation_id: rows = self._execute_with_retry( """SELECT usuario, mensagem, resposta, created_at FROM mensagens WHERE conversation_id = %s ORDER BY id DESC LIMIT %s""", (conversation_id, limite) ) else: rows = self._execute_with_retry( """SELECT usuario, mensagem, resposta, created_at FROM mensagens WHERE usuario = %s OR numero = %s ORDER BY id DESC LIMIT %s""", (usuario, numero, limite) ) # FIX 2026-10-02: ORDER BY id DESC devolve o MAIS RECENTE primeiro, mas # todos os consumidores (api.py::_build_context_history via # context_history[-N:] e contexto.obter_historico) assumem ordem # CRONOLÓGICA. Sem este reversed() o LLM via a pergunta nova no início # e a antiga no fim => repetia respostas e respondia a mensagens que já # passaram (a versão SQLite em database.py já reverte). return [dict(r) for r in reversed(rows)] if rows else [] except: return [] def recuperar_resposta_por_id(self, message_id): try: rows = self._execute_with_retry( "SELECT usuario, mensagem, resposta, modelo_usado, created_at FROM mensagens WHERE message_id = %s LIMIT 1", (message_id,) ) if rows: return dict(rows[0]) return None except: return None def registrar_mensagem_conversation_id(self, usuario, mensagem, resposta, conversation_id, **kwargs): return self.salvar_mensagem(usuario, mensagem, resposta, conversation_id=conversation_id, **kwargs) def get_recent_conversation_turns(self, usuario="", numero="", conversation_id="", limite=10): """Retorna turnos de conversa formatados para contexto LLM (user/assistant pairs).""" try: if conversation_id: rows = self._execute_with_retry( """SELECT usuario, mensagem, resposta, created_at FROM mensagens WHERE conversation_id = %s ORDER BY id DESC LIMIT %s""", (conversation_id, limite) ) else: rows = self._execute_with_retry( """SELECT usuario, mensagem, resposta, created_at FROM mensagens WHERE usuario = %s OR numero = %s ORDER BY id DESC LIMIT %s""", (usuario, numero, limite) ) if not rows: return [] turns = [] for r in reversed(rows): d = dict(r) if not isinstance(r, dict) else r msg = d.get('mensagem', '') resp = d.get('resposta', '') if msg: turns.append({"role": "user", "content": str(msg)}) if resp: turns.append({"role": "assistant", "content": str(resp)}) return turns except: return [] def limpar_contexto_usuario(self, usuario="", numero=""): try: key = usuario or numero self._execute_with_retry("DELETE FROM contexto WHERE user_key = %s", (key,), commit=True) return True except: return False def fazer_checkpoint_hf_sync(self): """Backup PostgreSQL via pg_dump com lock para evitar dupla execução entre workers.""" import subprocess import fcntl from pathlib import Path try: cloud_sync_dir = Path("/akira/data/cloud_sync") cloud_sync_dir.mkdir(parents=True, exist_ok=True) dump_path = cloud_sync_dir / "akira_dump.sql" lock_path = cloud_sync_dir / ".backup.lock" # Lock file para garantir que apenas 1 worker executa pg_dump lock_fd = open(lock_path, 'w') try: fcntl.flock(lock_fd, fcntl.LOCK_EX | fcntl.LOCK_NB) except BlockingIOError: logger.info("⏳ Backup já em execução por outro worker — pulando") lock_fd.close() return True try: params = self._conn_params if 'dsn' in params: cmd = f'pg_dump "{params["dsn"]}" > "{dump_path}"' else: cmd = (f'PGHOST={params["host"]} PGPORT={params["port"]} ' f'PGDATABASE={params["dbname"]} PGUSER={params["user"]} ' f'PGPASSWORD={params["password"]} pg_dump > "{dump_path}"') result = subprocess.run(cmd, shell=True, capture_output=True, text=True, timeout=120) if result.returncode == 0: logger.info(f"Checkpoint HF Sync concluído: {dump_path}") return True elif "server version mismatch" in result.stderr or "pg_dump version" in result.stderr: logger.warning(f"pg_dump version mismatch ignorado (HF Space infra): {result.stderr[:120]}") return True else: logger.error(f"pg_dump falhou: {result.stderr}") return False finally: fcntl.flock(lock_fd, fcntl.LOCK_UN) lock_fd.close() except Exception as e: logger.error(f"Erro no checkpoint: {e}") return False # ================================================================ # CONTINUOUS LEARNING — CRUD completo # ================================================================ def salvar_continuous_learning(self, usuario, mensagem, resposta_gerada=None, numero=None, nome_usuario=None, tipo_conversa='pv', resposta_do_bot=False, is_reply=False, reply_to_bot=False, contexto_grupo=None, modelo_usado='desconhecido', message_id=None, qualidade=0.0, tipo_conteudo='texto'): """Insere ou atualiza um registro de aprendizado contínuo.""" try: ts = time.time() self._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 UPDATE SET qualidade = EXCLUDED.qualidade, resposta_gerada = COALESCE(EXCLUDED.resposta_gerada, continuous_learning.resposta_gerada)""", (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), commit=True ) return True except Exception as e: logger.debug(f"[CL] Erro ao salvar continuous_learning: {e}") return False def recuperar_continuous_learning(self, usuario=None, limite=50, min_qualidade=0.0): """Recupera registros de aprendizado contínuo, opcionalmente filtrado por usuário.""" try: if usuario: rows = self._execute_with_retry( """SELECT ts, usuario, mensagem, resposta_do_bot, resposta_gerada, tipo_conversa, modelo_usado, qualidade, message_id, created_at FROM continuous_learning WHERE usuario = %s AND qualidade >= %s ORDER BY ts DESC LIMIT %s""", (usuario, min_qualidade, limite) ) else: rows = self._execute_with_retry( """SELECT ts, usuario, mensagem, resposta_do_bot, resposta_gerada, tipo_conversa, modelo_usado, qualidade, message_id, created_at FROM continuous_learning WHERE qualidade >= %s ORDER BY ts DESC LIMIT %s""", (min_qualidade, limite) ) return [dict(r) for r in rows] if rows else [] except Exception as e: logger.debug(f"[CL] Erro ao recuperar continuous_learning: {e}") return [] def contar_continuous_learning(self, usuario=None): """Conta registros de aprendizado contínuo.""" try: if usuario: rows = self._execute_with_retry( "SELECT COUNT(*) as cnt FROM continuous_learning WHERE usuario = %s", (usuario,) ) else: rows = self._execute_with_retry( "SELECT COUNT(*) as cnt FROM continuous_learning" ) if rows: r = rows[0] return r['cnt'] if isinstance(r, dict) else r[0] return 0 except: return 0 def cleanup_continuous_learning(self, days=90): """Remove registros de aprendizado contínuo mais antigos que N dias.""" try: conn = self._get_connection() cur = conn.cursor() cur.execute( "DELETE FROM continuous_learning WHERE created_at < NOW() - INTERVAL '%s days'", (days,) ) deleted = cur.rowcount conn.commit() _pool_return(conn) if deleted > 0: logger.info(f"[CL] Limpeza: {deleted} registros removidos (>{days} dias)") return deleted except: return 0 # ================================================================ # FINETUNING_EXAMPLES — CRUD completo # ================================================================ def salvar_exemplo_treino(self, user_id: str, conversation_id: str, input_message: str, expected_response: str, tone_level: str = "very_serious", hostility_score: int = 0, emotion_label: str = "neutro", quality_score: int = 50, embedding_input: Optional[bytes] = None, embedding_output: Optional[bytes] = None, similarity_score: float = 0.0, actual_response: Optional[str] = None) -> int: """ Insere um exemplo de treinamento com embeddings opcionais. Retorna o ID do exemplo inserido, ou -1 em caso de erro. Se embedding_input/output forem passados (bytes), são armazenados. Caso contrário, o caller deve chamar atualizar_embeddings_exemplo() depois. """ try: with self.get_connection_context() as conn: cur = conn.cursor() cur.execute(""" INSERT INTO finetuning_examples (user_id, conversation_id, input_message, expected_response, actual_response, quality_score, tone_level, hostility_score, embedding_input, embedding_output, similarity_score, emotion_label) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s) RETURNING id """, (user_id, conversation_id, input_message, expected_response, actual_response, quality_score, tone_level, hostility_score, embedding_input, embedding_output, similarity_score, emotion_label)) row = cur.fetchone() example_id = row['id'] if row else -1 if example_id > 0: logger.debug( f"[FT] Exemplo #{example_id} salvo | emotion={emotion_label} " f"| similarity={similarity_score:.3f}" ) return example_id except Exception as e: logger.error(f"[FT] Erro ao salvar exemplo: {e}") return -1 def atualizar_embeddings_exemplo(self, example_id: int, embedding_input: bytes, embedding_output: bytes) -> bool: """Atualiza os embeddings de um exemplo já existente.""" try: self._execute_with_retry( """UPDATE finetuning_examples SET embedding_input = %s, embedding_output = %s, indexed = TRUE, updated_at = NOW() WHERE id = %s""", (embedding_input, embedding_output, example_id), commit=True ) return True except Exception as e: logger.error(f"[FT] Erro ao atualizar embeddings exemplo #{example_id}: {e}") return False def get_exemplos_treino(self, min_quality: int = 0, user_id: Optional[str] = None, emotion_label: Optional[str] = None, tone_level: Optional[str] = None, limit: int = 100, offset: int = 0) -> List[Dict]: """ Recupera exemplos de treinamento com filtros. Args: min_quality: Qualidade mínima (0-100) user_id: Filtrar por utilizador específico emotion_label: Filtrar por emoção (e.g. 'alegria', 'raiva') tone_level: Filtrar por tom (e.g. 'formal', 'informal') limit: Máximo de resultados offset: Paginação Returns: Lista de dicionários com os exemplos """ try: conditions = ["quality_score >= %s"] params: List[Any] = [min_quality] if user_id: conditions.append("user_id = %s") params.append(user_id) if emotion_label: conditions.append("emotion_label = %s") params.append(emotion_label) if tone_level: conditions.append("tone_level = %s") params.append(tone_level) params.extend([limit, offset]) query = f""" SELECT id, user_id, conversation_id, input_message, expected_response, actual_response, quality_score, tone_level, hostility_score, embedding_input, embedding_output, similarity_score, emotion_label, created_at, updated_at, indexed FROM finetuning_examples WHERE {' AND '.join(conditions)} ORDER BY quality_score DESC, similarity_score DESC, created_at DESC LIMIT %s OFFSET %s """ rows = self._execute_with_retry(query, tuple(params)) if not rows: return [] results = [] for r in rows: d = dict(r) # Deserialize embeddings se presentes for emb_field in ('embedding_input', 'embedding_output'): if d.get(emb_field): emb = d[emb_field] if isinstance(emb, memoryview): emb = bytes(emb) d[emb_field] = emb results.append(d) return results except Exception as e: logger.error(f"[FT] Erro ao recuperar exemplos: {e}") return [] def contar_exemplos_treino(self, min_quality: int = 0, user_id: Optional[str] = None) -> int: """Conta exemplos de treinamento com filtros.""" try: params: List[Any] = [min_quality] conditions = ["quality_score >= %s"] if user_id: conditions.append("user_id = %s") params.append(user_id) rows = self._execute_with_retry( f"SELECT COUNT(*) as total FROM finetuning_examples WHERE {' AND '.join(conditions)}", tuple(params) ) if rows: r = rows[0] return r['total'] if isinstance(r, dict) else int(r[0]) return 0 except Exception as e: logger.error(f"[FT] Erro ao contar exemplos: {e}") return 0 def atualizar_quality_exemplo(self, example_id: int, quality_score: int) -> bool: """Atualiza score de qualidade de um exemplo.""" try: self._execute_with_retry( """UPDATE finetuning_examples SET quality_score = %s, updated_at = NOW() WHERE id = %s""", (max(0, min(100, quality_score)), example_id), commit=True ) return True except Exception as e: logger.error(f"[FT] Erro ao atualizar quality exemplo #{example_id}: {e}") return False def marcar_exemplo_indexado(self, example_id: int) -> bool: """Marca exemplo como tendo embeddings indexados.""" try: self._execute_with_retry( """UPDATE finetuning_examples SET indexed = TRUE, updated_at = NOW() WHERE id = %s""", (example_id,), commit=True ) return True except Exception: return False def get_estatisticas_finetuning(self) -> Dict: """Estatísticas agregadas da tabela de fine-tuning.""" try: conn = self._get_connection() cur = conn.cursor() stats: Dict[str, Any] = {} cur.execute("SELECT COUNT(*) as total FROM finetuning_examples") r = cur.fetchone() stats['total'] = r['total'] if isinstance(r, dict) else r[0] cur.execute(""" SELECT AVG(quality_score) as avg_q, AVG(similarity_score) as avg_s, COUNT(CASE WHEN indexed THEN 1 END) as indexed_count FROM finetuning_examples """) r = cur.fetchone() if r: d = dict(r) if not isinstance(r, dict) else r stats['avg_quality'] = round(float(d.get('avg_q') or 0), 2) stats['avg_similarity'] = round(float(d.get('avg_s') or 0), 3) stats['indexed_count'] = int(d.get('indexed_count') or 0) cur.execute(""" SELECT emotion_label, COUNT(*) as cnt, AVG(quality_score) as avg_q FROM finetuning_examples GROUP BY emotion_label ORDER BY cnt DESC """) stats['por_emocao'] = { row['emotion_label'] if isinstance(row, dict) else row[0]: { 'count': int(row['cnt'] if isinstance(row, dict) else row[1]), 'avg_quality': round(float(row['avg_q'] if isinstance(row, dict) else row[2] or 0), 2) } for row in cur.fetchall() } cur.execute(""" SELECT tone_level, COUNT(*) as cnt, AVG(quality_score) as avg_q FROM finetuning_examples GROUP BY tone_level ORDER BY cnt DESC """) stats['por_tom'] = { row['tone_level'] if isinstance(row, dict) else row[0]: { 'count': int(row['cnt'] if isinstance(row, dict) else row[1]), 'avg_quality': round(float(row['avg_q'] if isinstance(row, dict) else row[2] or 0), 2) } for row in cur.fetchall() } _pool_return(conn) return stats except Exception as e: logger.error(f"[FT] Erro ao gerar estatísticas: {e}") return {} # ================================================================ # EMBEDDINGS — Limpeza # ================================================================ def cleanup_embeddings(self, max_per_user=200): """Mantém apenas os N embeddings mais recentes por usuário (evita crescimento infinito).""" try: conn = self._get_connection() cur = conn.cursor() cur.execute(""" DELETE FROM embeddings WHERE id NOT IN ( SELECT id FROM ( SELECT id, ROW_NUMBER() OVER (PARTITION BY numero_usuario ORDER BY id DESC) as rn FROM embeddings ) sub WHERE rn <= %s ) """, (max_per_user,)) deleted = cur.rowcount conn.commit() _pool_return(conn) if deleted > 0: logger.info(f"[EMB] Limpeza: {deleted} embeddings antigos removidos (max {max_per_user}/user)") return deleted except: return 0 def contar_embeddings(self, numero_usuario=None): """Conta embeddings, opcionalmente filtrado por usuário.""" try: if numero_usuario: rows = self._execute_with_retry( "SELECT COUNT(*) as cnt FROM embeddings WHERE numero_usuario = %s", (numero_usuario,) ) else: rows = self._execute_with_retry( "SELECT COUNT(*) as cnt FROM embeddings" ) if rows: r = rows[0] return r['cnt'] if isinstance(r, dict) else r[0] return 0 except: return 0 def _pool_return(conn): """Retorna conexão ao pool de forma segura. Função de nível de módulo.""" if DatabasePG._pool: try: DatabasePG._pool.return_connection(conn) return except Exception: pass try: if conn and not conn.closed: conn.close() except Exception: pass def get_database(db_path=None): return DatabasePG(db_path or "")