Spaces:
Running
Running
Isaac Quarenta
fix(contexto): dedup de contexto, isolamento por conversa e placeholder de busca
04bb7b4 Download modules/database_pg.py from akra35567/AKIRA-SOFTEDGE: direct link, hf CLI and curl.
- Browser
- Download file 91.2 kB
-
https://huggingface.co/spaces/akra35567/AKIRA-SOFTEDGE/resolve/main/modules/database_pg.py
- Command line
-
hf download hf://spaces/akra35567/AKIRA-SOFTEDGE/modules/database_pg.py
-
curl -L -o database_pg.py https://huggingface.co/spaces/akra35567/AKIRA-SOFTEDGE/resolve/main/modules/database_pg.py
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 | |
| # ================================================================ | |
| 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 "") | |