AKIRA-SOFTEDGE / modules /database_pg.py
Isaac Quarenta
fix(contexto): dedup de contexto, isolamento por conversa e placeholder de busca
04bb7b4
Raw History Blame Contribute Delete
91.2 kB
"""
================================================================================
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
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"""
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()
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
return conn
except Exception as 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:
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):
# ✅ CONNECTION POOL: Usar pool quando disponível
if DatabasePG._pool:
try:
return DatabasePG._pool.get_connection()
except Exception as e:
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
return conn
except psycopg2.OperationalError as e:
if attempt < self.max_retries - 1:
time.sleep(self.retry_delay * (2 ** attempt))
continue
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."""
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"):
return cur.fetchall()
return None
except Exception as 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 "")