feat: complete local integration of AI mind codebase (brain, memory, etc.) in text2video, supporting both Direct Codebase execution and Remote HTTP fallback
3d7a63c | # core/memory.py | |
| import sqlite3 | |
| from datetime import datetime | |
| import threading | |
| import os | |
| import json | |
| import secrets | |
| import hashlib | |
| import time | |
| import re | |
| from werkzeug.security import generate_password_hash, check_password_hash | |
| BASE_DIR = os.path.dirname(os.path.abspath(__file__)) | |
| class ConversationMemory: | |
| def __init__(self, db_path=None): | |
| self.db_path = db_path or "/data/invicta_data/memory.db" | |
| os.makedirs(os.path.dirname(self.db_path), exist_ok=True) | |
| self.lock = threading.Lock() | |
| self._init_tables() | |
| def _get_connection(self): | |
| """Create a new SQLite connection with WAL mode and busy timeout.""" | |
| conn = sqlite3.connect( | |
| self.db_path, | |
| check_same_thread=False, | |
| timeout=30.0, | |
| isolation_level=None | |
| ) | |
| conn.execute("PRAGMA foreign_keys = ON") | |
| conn.execute("PRAGMA journal_mode = WAL") | |
| conn.execute("PRAGMA synchronous = NORMAL") | |
| conn.execute("PRAGMA temp_store = MEMORY") | |
| conn.execute("PRAGMA mmap_size = 30000000000") | |
| conn.execute("PRAGMA busy_timeout = 30000") | |
| conn.row_factory = sqlite3.Row | |
| return conn | |
| def _get_cursor(self): | |
| """Legacy compatibility wrapper.""" | |
| conn = self._get_connection() | |
| return conn, conn.cursor() | |
| def _execute_with_retry(self, sql, params=(), max_retries=5, sleep_base=0.05): | |
| """Execute SQL with automatic retry on database lock.""" | |
| for attempt in range(max_retries): | |
| conn = None | |
| try: | |
| conn = self._get_connection() | |
| cur = conn.cursor() | |
| cur.execute("BEGIN IMMEDIATE") | |
| cur.execute(sql, params) | |
| conn.commit() | |
| return cur | |
| except sqlite3.OperationalError as e: | |
| if conn: | |
| try: | |
| conn.rollback() | |
| except Exception: | |
| pass | |
| if "database is locked" in str(e).lower() and attempt < max_retries - 1: | |
| time.sleep(sleep_base * (attempt + 1)) | |
| continue | |
| raise | |
| finally: | |
| if conn: | |
| try: | |
| conn.close() | |
| except Exception: | |
| pass | |
| return None | |
| def _init_tables(self): | |
| with self.lock: | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute("BEGIN IMMEDIATE") | |
| cur.execute(""" | |
| CREATE TABLE IF NOT EXISTS users ( | |
| id INTEGER PRIMARY KEY AUTOINCREMENT, | |
| username TEXT UNIQUE NOT NULL, | |
| password_hash TEXT NOT NULL, | |
| google_id TEXT, | |
| email TEXT, | |
| created_at TEXT DEFAULT (datetime('now')) | |
| ); | |
| """) | |
| cur.execute(""" | |
| CREATE TABLE IF NOT EXISTS auth_tokens ( | |
| token TEXT PRIMARY KEY, | |
| user_id INTEGER NOT NULL REFERENCES users(id) ON DELETE CASCADE, | |
| created_at TEXT DEFAULT (datetime('now')), | |
| last_used TEXT DEFAULT (datetime('now')) | |
| ); | |
| """) | |
| cur.execute(""" | |
| CREATE TABLE IF NOT EXISTS chat_sessions ( | |
| id INTEGER PRIMARY KEY AUTOINCREMENT, | |
| user_id INTEGER NOT NULL REFERENCES users(id) ON DELETE CASCADE, | |
| title TEXT, | |
| created_at TEXT DEFAULT (datetime('now')) | |
| ); | |
| """) | |
| cur.execute(""" | |
| CREATE TABLE IF NOT EXISTS conversations ( | |
| id INTEGER PRIMARY KEY AUTOINCREMENT, | |
| user_id INTEGER REFERENCES users(id) ON DELETE CASCADE, | |
| session_id INTEGER REFERENCES chat_sessions(id) ON DELETE CASCADE, | |
| role TEXT NOT NULL, | |
| content TEXT NOT NULL, | |
| timestamp TEXT DEFAULT (datetime('now')) | |
| ); | |
| """) | |
| cur.execute(""" | |
| CREATE TABLE IF NOT EXISTS user_profile ( | |
| user_id INTEGER NOT NULL DEFAULT 0 REFERENCES users(id) ON DELETE CASCADE, | |
| key TEXT NOT NULL, | |
| value TEXT, | |
| category TEXT DEFAULT 'general', | |
| updated_at TEXT DEFAULT (datetime('now')), | |
| PRIMARY KEY (user_id, key) | |
| ); | |
| """) | |
| cur.execute(""" | |
| CREATE TABLE IF NOT EXISTS user_learning ( | |
| id INTEGER PRIMARY KEY AUTOINCREMENT, | |
| user_id INTEGER NOT NULL REFERENCES users(id) ON DELETE CASCADE, | |
| category TEXT NOT NULL, | |
| key TEXT NOT NULL, | |
| value TEXT NOT NULL, | |
| confidence REAL DEFAULT 0.5, | |
| times_observed INTEGER DEFAULT 1, | |
| first_seen TEXT DEFAULT (datetime('now')), | |
| updated_at TEXT DEFAULT (datetime('now')), | |
| UNIQUE(user_id, category, key) | |
| ); | |
| """) | |
| cur.execute(""" | |
| CREATE TABLE IF NOT EXISTS user_insights ( | |
| user_id INTEGER PRIMARY KEY REFERENCES users(id) ON DELETE CASCADE, | |
| preferred_formality TEXT DEFAULT 'casual', | |
| preferred_response_length TEXT DEFAULT 'medium', | |
| preferred_tone TEXT DEFAULT 'friendly', | |
| topics_of_interest TEXT DEFAULT '[]', | |
| communication_patterns TEXT DEFAULT '{}', | |
| total_interactions INTEGER DEFAULT 0, | |
| last_analyzed TEXT DEFAULT (datetime('now')), | |
| updated_at TEXT DEFAULT (datetime('now')) | |
| ); | |
| """) | |
| cur.execute(""" | |
| CREATE TABLE IF NOT EXISTS api_keys ( | |
| id INTEGER PRIMARY KEY AUTOINCREMENT, | |
| user_id INTEGER NOT NULL REFERENCES users(id) ON DELETE CASCADE, | |
| key_hash TEXT NOT NULL, | |
| name TEXT DEFAULT 'API Key', | |
| prefix TEXT DEFAULT '', | |
| created_at TEXT DEFAULT (datetime('now')), | |
| last_used TEXT, | |
| revoked INTEGER DEFAULT 0, | |
| rate_limit INTEGER DEFAULT 60 | |
| ); | |
| """) | |
| cur.execute("CREATE INDEX IF NOT EXISTS idx_conversations_user ON conversations(user_id);") | |
| cur.execute("CREATE INDEX IF NOT EXISTS idx_conversations_session ON conversations(session_id);") | |
| cur.execute("CREATE INDEX IF NOT EXISTS idx_conversations_user_session ON conversations(user_id, session_id);") | |
| cur.execute("CREATE INDEX IF NOT EXISTS idx_conversations_timestamp ON conversations(timestamp);") | |
| cur.execute("CREATE INDEX IF NOT EXISTS idx_learning_user ON user_learning(user_id);") | |
| cur.execute("CREATE INDEX IF NOT EXISTS idx_learning_category ON user_learning(user_id, category);") | |
| cur.execute("CREATE INDEX IF NOT EXISTS idx_learning_unique ON user_learning(user_id, category, key);") | |
| cur.execute("CREATE INDEX IF NOT EXISTS idx_auth_tokens_user ON auth_tokens(user_id);") | |
| cur.execute("CREATE INDEX IF NOT EXISTS idx_api_keys_user ON api_keys(user_id);") | |
| cur.execute("CREATE INDEX IF NOT EXISTS idx_api_keys_hash ON api_keys(key_hash);") | |
| conn.commit() | |
| except Exception as e: | |
| conn.rollback() | |
| if "already exists" not in str(e).lower() and "duplicate" not in str(e).lower(): | |
| raise | |
| finally: | |
| conn.close() | |
| # ββ Auth Tokens ββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| def create_token(self, user_id): | |
| token = secrets.token_urlsafe(32) | |
| with self.lock: | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute("BEGIN IMMEDIATE") | |
| cur.execute( | |
| "INSERT INTO auth_tokens (token, user_id, created_at, last_used) VALUES (?, ?, ?, ?)", | |
| (token, user_id, datetime.now().isoformat(sep=' ', timespec='seconds'), datetime.now().isoformat(sep=' ', timespec='seconds')) | |
| ) | |
| conn.commit() | |
| finally: | |
| conn.close() | |
| return token | |
| def validate_token(self, token): | |
| if not token or len(token) < 10: | |
| return None | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute("SELECT user_id FROM auth_tokens WHERE token = ?", (token,)) | |
| row = cur.fetchone() | |
| return row["user_id"] if row else None | |
| finally: | |
| conn.close() | |
| def touch_token(self, token): | |
| """Update token last_used with retry logic for concurrent access.""" | |
| now = datetime.now().isoformat(sep=' ', timespec='seconds') | |
| for attempt in range(5): | |
| conn = None | |
| try: | |
| conn = self._get_connection() | |
| cur = conn.cursor() | |
| cur.execute("BEGIN IMMEDIATE") | |
| cur.execute( | |
| "UPDATE auth_tokens SET last_used = ? WHERE token = ?", | |
| (now, token) | |
| ) | |
| conn.commit() | |
| return | |
| except sqlite3.OperationalError as e: | |
| if conn: | |
| try: | |
| conn.rollback() | |
| except Exception: | |
| pass | |
| if "database is locked" in str(e).lower() and attempt < 4: | |
| time.sleep(0.05 * (attempt + 1)) | |
| continue | |
| raise | |
| finally: | |
| if conn: | |
| try: | |
| conn.close() | |
| except Exception: | |
| pass | |
| def delete_token(self, token): | |
| with self.lock: | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute("BEGIN IMMEDIATE") | |
| cur.execute("DELETE FROM auth_tokens WHERE token = ?", (token,)) | |
| conn.commit() | |
| finally: | |
| conn.close() | |
| def revoke_user_tokens(self, user_id): | |
| with self.lock: | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute("BEGIN IMMEDIATE") | |
| cur.execute("DELETE FROM auth_tokens WHERE user_id = ?", (user_id,)) | |
| conn.commit() | |
| finally: | |
| conn.close() | |
| def cleanup_old_tokens(self, max_age_days=30): | |
| with self.lock: | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute("BEGIN IMMEDIATE") | |
| cur.execute( | |
| "DELETE FROM auth_tokens WHERE last_used < datetime('now', '-{} days')".format(max_age_days) | |
| ) | |
| conn.commit() | |
| finally: | |
| conn.close() | |
| # ββ User Auth βββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| def create_user(self, username, password): | |
| """Create a user with validation. Returns True on success, False on failure.""" | |
| # Defense-in-depth validations (also enforced in app.py) | |
| if not username or not password: | |
| return False | |
| if len(username) < 3 or len(username) > 30: | |
| return False | |
| if not re.match(r'^[a-zA-Z0-9_\-]+$', username): | |
| return False | |
| if username.lower().startswith("__guest_"): | |
| return False | |
| if len(password) < 8: | |
| return False | |
| if not re.search(r'[A-Z]', password) or not re.search(r'[a-z]', password) or not re.search(r'[0-9]', password): | |
| return False | |
| with self.lock: | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute("BEGIN IMMEDIATE") | |
| pw_hash = generate_password_hash(password) | |
| # Store username lowercase for case-normalization | |
| cur.execute( | |
| "INSERT INTO users (username, password_hash) VALUES (?, ?)", | |
| (username.strip().lower(), pw_hash) | |
| ) | |
| conn.commit() | |
| return True | |
| except sqlite3.IntegrityError: | |
| conn.rollback() | |
| return False | |
| finally: | |
| conn.close() | |
| def verify_user(self, username, password): | |
| if not username or not password: | |
| return None | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| # Case-normalized lookup | |
| cur.execute( | |
| "SELECT id, password_hash FROM users WHERE username = ?", (username.strip().lower(),) | |
| ) | |
| row = cur.fetchone() | |
| if row and check_password_hash(row["password_hash"], password): | |
| return row["id"] | |
| return None | |
| finally: | |
| conn.close() | |
| def get_username(self, user_id): | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute("SELECT username FROM users WHERE id = ?", (user_id,)) | |
| row = cur.fetchone() | |
| return row["username"] if row else None | |
| finally: | |
| conn.close() | |
| def is_guest_user(self, user_id): | |
| username = self.get_username(user_id) | |
| return username.startswith("__guest_") if username else True | |
| def create_guest_user(self, username=None, password=None): | |
| if not username: | |
| username = f"__guest_{secrets.token_hex(8)}" | |
| if not password: | |
| password = secrets.token_hex(32) | |
| with self.lock: | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute("BEGIN IMMEDIATE") | |
| pw_hash = generate_password_hash(password) | |
| cur.execute( | |
| "INSERT INTO users (username, password_hash) VALUES (?, ?)", | |
| (username, pw_hash) | |
| ) | |
| user_id = cur.lastrowid | |
| conn.commit() | |
| return user_id | |
| except sqlite3.IntegrityError: | |
| conn.rollback() | |
| return None | |
| finally: | |
| conn.close() | |
| def find_or_create_google_user(self, google_id, name, email): | |
| with self.lock: | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute("BEGIN IMMEDIATE") | |
| cur.execute("SELECT id FROM users WHERE google_id = ?", (google_id,)) | |
| row = cur.fetchone() | |
| if row: | |
| conn.commit() | |
| return row["id"] | |
| cur.execute("SELECT id FROM users WHERE email = ?", (email,)) | |
| row = cur.fetchone() | |
| if row: | |
| cur.execute("UPDATE users SET google_id = ? WHERE id = ?", (google_id, row["id"])) | |
| conn.commit() | |
| return row["id"] | |
| base_username = name.replace(" ", "_").lower() | |
| base_username = ''.join(c for c in base_username if c.isalnum() or c == '_') | |
| if not base_username: | |
| base_username = "user" | |
| username = base_username | |
| counter = 1 | |
| while True: | |
| cur.execute("SELECT id FROM users WHERE username = ?", (username,)) | |
| if not cur.fetchone(): | |
| break | |
| username = f"{base_username}_{counter}" | |
| counter += 1 | |
| random_pw = secrets.token_hex(32) | |
| pw_hash = generate_password_hash(random_pw) | |
| cur.execute( | |
| "INSERT INTO users (username, password_hash, google_id, email) VALUES (?, ?, ?, ?)", | |
| (username, pw_hash, google_id, email) | |
| ) | |
| user_id = cur.lastrowid | |
| conn.commit() | |
| return user_id | |
| except sqlite3.IntegrityError: | |
| conn.rollback() | |
| return None | |
| finally: | |
| conn.close() | |
| def transfer_guest_data(self, guest_user_id, real_user_id): | |
| """ | |
| Transfer guest user data to real user account. | |
| Handles UNIQUE constraint conflicts by merging learning data. | |
| """ | |
| with self.lock: | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute("BEGIN IMMEDIATE") | |
| # 1. Transfer conversations | |
| cur.execute( | |
| "UPDATE conversations SET user_id = ? WHERE user_id = ?", | |
| (real_user_id, guest_user_id) | |
| ) | |
| # 2. Transfer chat sessions | |
| cur.execute( | |
| "UPDATE chat_sessions SET user_id = ? WHERE user_id = ?", | |
| (real_user_id, guest_user_id) | |
| ) | |
| # 3. Transfer user_profile (ON CONFLICT handles duplicates) | |
| cur.execute( | |
| "SELECT user_id, key, value, category, updated_at FROM user_profile WHERE user_id = ?", | |
| (guest_user_id,) | |
| ) | |
| for row in cur.fetchall(): | |
| cur.execute( | |
| """INSERT INTO user_profile (user_id, key, value, category, updated_at) | |
| VALUES (?, ?, ?, ?, ?) | |
| ON CONFLICT (user_id, key) DO UPDATE SET | |
| value = excluded.value, | |
| category = excluded.category, | |
| updated_at = excluded.updated_at""", | |
| (real_user_id, row["key"], row["value"], row["category"], row["updated_at"]) | |
| ) | |
| cur.execute("DELETE FROM user_profile WHERE user_id = ?", (guest_user_id,)) | |
| # 4. Transfer user_learning with conflict resolution | |
| cur.execute( | |
| "SELECT category, key, value, confidence, times_observed, first_seen, updated_at " | |
| "FROM user_learning WHERE user_id = ?", | |
| (guest_user_id,) | |
| ) | |
| guest_rows = cur.fetchall() | |
| for row in guest_rows: | |
| cur.execute( | |
| "SELECT id, confidence, times_observed, value FROM user_learning " | |
| "WHERE user_id = ? AND category = ? AND key = ?", | |
| (real_user_id, row["category"], row["key"]) | |
| ) | |
| existing = cur.fetchone() | |
| if existing: | |
| new_confidence = max(float(existing["confidence"]), float(row["confidence"])) | |
| new_times = int(existing["times_observed"]) + int(row["times_observed"]) | |
| new_value = row["value"] if float(row["confidence"]) >= float(existing["confidence"]) else existing["value"] | |
| cur.execute( | |
| "UPDATE user_learning SET value = ?, confidence = ?, times_observed = ?, updated_at = ? " | |
| "WHERE id = ?", | |
| (new_value, new_confidence, new_times, | |
| datetime.now().isoformat(sep=' ', timespec='seconds'), existing["id"]) | |
| ) | |
| else: | |
| cur.execute( | |
| "INSERT INTO user_learning (user_id, category, key, value, confidence, " | |
| "times_observed, first_seen, updated_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?)", | |
| (real_user_id, row["category"], row["key"], row["value"], | |
| row["confidence"], row["times_observed"], row["first_seen"], | |
| datetime.now().isoformat(sep=' ', timespec='seconds')) | |
| ) | |
| cur.execute("DELETE FROM user_learning WHERE user_id = ?", (guest_user_id,)) | |
| # 5. Transfer user_insights (ON CONFLICT handles duplicates) | |
| cur.execute( | |
| "SELECT * FROM user_insights WHERE user_id = ?", | |
| (guest_user_id,) | |
| ) | |
| guest_insights = cur.fetchone() | |
| if guest_insights: | |
| insights_dict = dict(guest_insights) | |
| insights_dict.pop("user_id", None) | |
| cols = [k for k in insights_dict.keys()] | |
| vals = [insights_dict[k] for k in cols] | |
| placeholders = ",".join(["?"] * len(vals)) | |
| updates = ",".join([f"{c}=excluded.{c}" for c in cols]) | |
| cur.execute( | |
| f"""INSERT INTO user_insights (user_id, {','.join(cols)}) | |
| VALUES (? , {placeholders}) | |
| ON CONFLICT (user_id) DO UPDATE SET {updates}""", | |
| [real_user_id] + vals | |
| ) | |
| cur.execute("DELETE FROM user_insights WHERE user_id = ?", (guest_user_id,)) | |
| # 6. Transfer API keys | |
| cur.execute( | |
| "UPDATE api_keys SET user_id = ? WHERE user_id = ?", | |
| (real_user_id, guest_user_id) | |
| ) | |
| # 7. Delete guest auth tokens | |
| cur.execute("DELETE FROM auth_tokens WHERE user_id = ?", (guest_user_id,)) | |
| # 8. Delete guest user | |
| cur.execute("DELETE FROM users WHERE id = ?", (guest_user_id,)) | |
| conn.commit() | |
| print(f"β Transferred guest {guest_user_id} β user {real_user_id}") | |
| except Exception as e: | |
| conn.rollback() | |
| print(f"β Transfer error: {e}") | |
| raise | |
| finally: | |
| conn.close() | |
| def cleanup_guest_users(self, max_age_days=14): | |
| with self.lock: | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute("BEGIN IMMEDIATE") | |
| cur.execute( | |
| "SELECT id FROM users WHERE username LIKE '__guest_%' AND created_at < datetime('now', '-{} days')".format(max_age_days) | |
| ) | |
| guest_ids = [row["id"] for row in cur.fetchall()] | |
| if not guest_ids: | |
| conn.commit() | |
| return | |
| placeholders = ','.join('?' * len(guest_ids)) | |
| cur.execute(f"DELETE FROM conversations WHERE user_id IN ({placeholders})", guest_ids) | |
| cur.execute(f"DELETE FROM chat_sessions WHERE user_id IN ({placeholders})", guest_ids) | |
| cur.execute(f"DELETE FROM user_learning WHERE user_id IN ({placeholders})", guest_ids) | |
| cur.execute(f"DELETE FROM user_insights WHERE user_id IN ({placeholders})", guest_ids) | |
| cur.execute(f"DELETE FROM auth_tokens WHERE user_id IN ({placeholders})", guest_ids) | |
| cur.execute(f"DELETE FROM api_keys WHERE user_id IN ({placeholders})", guest_ids) | |
| cur.execute(f"DELETE FROM user_profile WHERE user_id IN ({placeholders})", guest_ids) | |
| cur.execute(f"DELETE FROM users WHERE id IN ({placeholders})", guest_ids) | |
| conn.commit() | |
| print(f"π§Ή Cleaned up {len(guest_ids)} old guest users") | |
| finally: | |
| conn.close() | |
| # ββ Sessions βββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| def create_session(self, user_id, title=None): | |
| with self.lock: | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute("BEGIN IMMEDIATE") | |
| cur.execute( | |
| "INSERT INTO chat_sessions (user_id, title) VALUES (?, ?)", | |
| (user_id, title or "New Chat") | |
| ) | |
| session_id = cur.lastrowid | |
| conn.commit() | |
| return session_id | |
| finally: | |
| conn.close() | |
| def rename_session(self, session_id, title): | |
| with self.lock: | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute("BEGIN IMMEDIATE") | |
| cur.execute("UPDATE chat_sessions SET title = ? WHERE id = ?", (title[:60], session_id)) | |
| conn.commit() | |
| finally: | |
| conn.close() | |
| def rename_session_if_default(self, session_id, title, user_id): | |
| with self.lock: | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute("BEGIN IMMEDIATE") | |
| cur.execute( | |
| "SELECT title FROM chat_sessions WHERE id = ? AND user_id = ?", | |
| (session_id, user_id) | |
| ) | |
| row = cur.fetchone() | |
| if row and (row["title"] == "New Chat" or not row["title"]): | |
| cur.execute("UPDATE chat_sessions SET title = ? WHERE id = ?", (title[:60], session_id)) | |
| conn.commit() | |
| else: | |
| conn.commit() | |
| finally: | |
| conn.close() | |
| def delete_session(self, session_id, user_id): | |
| with self.lock: | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute("BEGIN IMMEDIATE") | |
| cur.execute("DELETE FROM conversations WHERE session_id = ? AND user_id = ?", (session_id, user_id)) | |
| cur.execute("DELETE FROM chat_sessions WHERE id = ? AND user_id = ?", (session_id, user_id)) | |
| conn.commit() | |
| finally: | |
| conn.close() | |
| def get_sessions(self, user_id): | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute( | |
| """SELECT cs.id, cs.title, cs.created_at, | |
| COUNT(c.id) as msg_count | |
| FROM chat_sessions cs | |
| LEFT JOIN conversations c ON c.session_id = cs.id | |
| WHERE cs.user_id = ? | |
| GROUP BY cs.id | |
| ORDER BY cs.created_at DESC | |
| LIMIT 50""", | |
| (user_id,) | |
| ) | |
| rows = cur.fetchall() | |
| return [dict(r) for r in rows] | |
| finally: | |
| conn.close() | |
| # ββ Chat Messages ββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| def save(self, role, content, user_id=None, session_id=None): | |
| with self.lock: | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute("BEGIN IMMEDIATE") | |
| cur.execute( | |
| "INSERT INTO conversations (user_id, session_id, role, content, timestamp) VALUES (?, ?, ?, ?, ?)", | |
| (user_id, session_id, role, content, datetime.now().isoformat(sep=' ', timespec='seconds')) | |
| ) | |
| conn.commit() | |
| finally: | |
| conn.close() | |
| def get_recent(self, user_id, limit=20, session_id=None): | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| if session_id: | |
| cur.execute( | |
| "SELECT role, content, timestamp FROM conversations " | |
| "WHERE user_id = ? AND session_id = ? ORDER BY id DESC LIMIT ?", | |
| (user_id, session_id, limit) | |
| ) | |
| else: | |
| cur.execute( | |
| "SELECT role, content, timestamp FROM conversations " | |
| "WHERE user_id = ? ORDER BY id DESC LIMIT ?", | |
| (user_id, limit) | |
| ) | |
| rows = cur.fetchall() | |
| return list(reversed([(r["role"], r["content"], r["timestamp"]) for r in rows])) | |
| finally: | |
| conn.close() | |
| def get_session_messages(self, session_id, user_id): | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute( | |
| "SELECT role, content, timestamp FROM conversations " | |
| "WHERE session_id = ? AND user_id = ? ORDER BY id ASC", | |
| (session_id, user_id) | |
| ) | |
| rows = cur.fetchall() | |
| return [dict(r) for r in rows] | |
| finally: | |
| conn.close() | |
| def get_chat_history_list(self, user_id): | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute( | |
| "SELECT id, content FROM conversations " | |
| "WHERE user_id = ? AND role = 'user' ORDER BY id DESC LIMIT 20", | |
| (user_id,) | |
| ) | |
| rows = cur.fetchall() | |
| return [(r["id"], r["content"]) for r in rows] | |
| finally: | |
| conn.close() | |
| def search_history(self, user_id, query, limit=10): | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute( | |
| "SELECT role, content, timestamp FROM conversations " | |
| "WHERE user_id = ? AND content LIKE ? ORDER BY id DESC LIMIT ?", | |
| (user_id, f"%{query}%", limit) | |
| ) | |
| rows = cur.fetchall() | |
| return [dict(r) for r in rows] | |
| finally: | |
| conn.close() | |
| def get_user_message_count(self, user_id): | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute( | |
| "SELECT COUNT(*) as c FROM conversations WHERE user_id = ? AND role = 'user'", | |
| (user_id,) | |
| ) | |
| row = cur.fetchone() | |
| return row["c"] if row else 0 | |
| finally: | |
| conn.close() | |
| # ββ User Profile (user-specific) βββββββββββββββββββββββββββββββββββββ | |
| def update_user_profile(self, user_id, key, value, category="general"): | |
| with self.lock: | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute("BEGIN IMMEDIATE") | |
| cur.execute(""" | |
| INSERT INTO user_profile (user_id, key, value, category, updated_at) | |
| VALUES (?, ?, ?, ?, ?) | |
| ON CONFLICT (user_id, key) DO UPDATE SET | |
| value = excluded.value, category = excluded.category, updated_at = excluded.updated_at | |
| """, (user_id, key, str(value), category, datetime.now().isoformat(sep=' ', timespec='seconds'))) | |
| conn.commit() | |
| finally: | |
| conn.close() | |
| def get_user_profile(self, user_id): | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute( | |
| "SELECT key, value FROM user_profile WHERE user_id = ? ORDER BY key", | |
| (user_id,) | |
| ) | |
| rows = cur.fetchall() | |
| return {r["key"]: r["value"] for r in rows} | |
| finally: | |
| conn.close() | |
| def delete_profile_key(self, user_id, key): | |
| with self.lock: | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute("BEGIN IMMEDIATE") | |
| cur.execute("DELETE FROM user_profile WHERE user_id = ? AND key = ?", (user_id, key)) | |
| conn.commit() | |
| finally: | |
| conn.close() | |
| def clear_profile(self, user_id): | |
| with self.lock: | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute("BEGIN IMMEDIATE") | |
| cur.execute("DELETE FROM user_profile WHERE user_id = ?", (user_id,)) | |
| conn.commit() | |
| finally: | |
| conn.close() | |
| # ββ Self-Learning ββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| def store_learning(self, user_id, category, key, value, confidence=0.5): | |
| with self.lock: | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute("BEGIN IMMEDIATE") | |
| cur.execute( | |
| "SELECT id, confidence, times_observed FROM user_learning " | |
| "WHERE user_id = ? AND category = ? AND key = ?", | |
| (user_id, category, key) | |
| ) | |
| existing = cur.fetchone() | |
| if existing: | |
| new_confidence = min(existing["confidence"] + 0.1, 2.0) | |
| new_times = existing["times_observed"] + 1 | |
| cur.execute( | |
| "UPDATE user_learning SET confidence = ?, times_observed = ?, value = ?, updated_at = ? WHERE id = ?", | |
| (new_confidence, new_times, value, datetime.now().isoformat(sep=' ', timespec='seconds'), existing["id"]) | |
| ) | |
| else: | |
| cur.execute( | |
| "INSERT INTO user_learning (user_id, category, key, value, confidence, first_seen, updated_at) VALUES (?, ?, ?, ?, ?, ?, ?)", | |
| (user_id, category, key, value, confidence, datetime.now().isoformat(sep=' ', timespec='seconds'), datetime.now().isoformat(sep=' ', timespec='seconds')) | |
| ) | |
| conn.commit() | |
| finally: | |
| conn.close() | |
| def get_user_learning(self, user_id, category=None): | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| if category: | |
| cur.execute( | |
| "SELECT category, key, value, confidence, times_observed FROM user_learning WHERE user_id = ? AND category = ? ORDER BY confidence DESC", | |
| (user_id, category) | |
| ) | |
| else: | |
| cur.execute( | |
| "SELECT category, key, value, confidence, times_observed FROM user_learning WHERE user_id = ? ORDER BY confidence DESC", | |
| (user_id,) | |
| ) | |
| rows = cur.fetchall() | |
| return [dict(r) for r in rows] | |
| finally: | |
| conn.close() | |
| def get_learning_by_category(self, user_id, category): | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute( | |
| "SELECT key, value FROM user_learning WHERE user_id = ? AND category = ? ORDER BY confidence DESC", | |
| (user_id, category) | |
| ) | |
| rows = cur.fetchall() | |
| return {r["key"]: r["value"] for r in rows} | |
| finally: | |
| conn.close() | |
| def get_top_learning_topics(self, user_id, limit=10): | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute( | |
| "SELECT key, value, confidence, times_observed FROM user_learning " | |
| "WHERE user_id = ? AND category = 'interest' ORDER BY times_observed DESC, confidence DESC LIMIT ?", | |
| (user_id, limit) | |
| ) | |
| rows = cur.fetchall() | |
| return [dict(r) for r in rows] | |
| finally: | |
| conn.close() | |
| def prune_old_learning(self, user_id, max_entries=500): | |
| with self.lock: | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute("BEGIN IMMEDIATE") | |
| cur.execute("SELECT COUNT(*) as c FROM user_learning WHERE user_id = ?", (user_id,)) | |
| count = cur.fetchone()["c"] | |
| if count > max_entries: | |
| cur.execute( | |
| "DELETE FROM user_learning WHERE id IN (SELECT id FROM user_learning WHERE user_id = ? ORDER BY confidence ASC, times_observed ASC LIMIT ?)", | |
| (user_id, count - max_entries) | |
| ) | |
| conn.commit() | |
| else: | |
| conn.commit() | |
| finally: | |
| conn.close() | |
| # ββ User Insights ββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| def update_user_insights(self, user_id, insights_dict): | |
| with self.lock: | |
| allowed_keys = ["preferred_formality", "preferred_response_length", "preferred_tone", "topics_of_interest", "communication_patterns", "total_interactions"] | |
| insert_cols = [k for k in insights_dict if k in allowed_keys] | |
| if not insert_cols: | |
| return | |
| insert_vals = [json.dumps(insights_dict[k]) if isinstance(insights_dict[k], (dict, list)) else str(insights_dict[k]) for k in insert_cols] | |
| insert_cols.append("updated_at") | |
| insert_vals.append(datetime.now().isoformat(sep=' ', timespec='seconds')) | |
| updates = [f"{column} = excluded.{column}" for column in insert_cols] | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute("BEGIN IMMEDIATE") | |
| sql = f"""INSERT INTO user_insights (user_id, {', '.join(insert_cols)}) | |
| VALUES (?, {', '.join(['?'] * len(insert_vals))}) | |
| ON CONFLICT (user_id) DO UPDATE SET {', '.join(updates)}""" | |
| cur.execute(sql, [user_id] + insert_vals) | |
| conn.commit() | |
| finally: | |
| conn.close() | |
| def get_user_insights(self, user_id): | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute("SELECT * FROM user_insights WHERE user_id = ?", (user_id,)) | |
| row = cur.fetchone() | |
| if row: | |
| result = dict(row) | |
| for field in ["topics_of_interest", "communication_patterns"]: | |
| if result.get(field): | |
| try: | |
| result[field] = json.loads(result[field]) | |
| except (json.JSONDecodeError, TypeError): | |
| result[field] = [] if field == "topics_of_interest" else {} | |
| return result | |
| return None | |
| finally: | |
| conn.close() | |
| def increment_interaction_count(self, user_id): | |
| with self.lock: | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute("BEGIN IMMEDIATE") | |
| cur.execute( | |
| """INSERT INTO user_insights (user_id, total_interactions, updated_at) | |
| VALUES (?, 1, ?) | |
| ON CONFLICT (user_id) DO UPDATE SET | |
| total_interactions = total_interactions + 1, updated_at = excluded.updated_at""", | |
| (user_id, datetime.now().isoformat(sep=' ', timespec='seconds')) | |
| ) | |
| conn.commit() | |
| finally: | |
| conn.close() | |
| # ββ API Keys βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| def create_api_key(self, user_id, name="API Key"): | |
| raw_key = "invicta_sk_" + secrets.token_urlsafe(32) | |
| key_hash = hashlib.sha256(raw_key.encode()).hexdigest() | |
| prefix = raw_key[:14] + "..." | |
| with self.lock: | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute("BEGIN IMMEDIATE") | |
| cur.execute( | |
| "INSERT INTO api_keys (user_id, key_hash, name, prefix, created_at) VALUES (?, ?, ?, ?, ?)", | |
| (user_id, key_hash, name, prefix, datetime.now().isoformat(sep=' ', timespec='seconds')) | |
| ) | |
| key_id = cur.lastrowid | |
| conn.commit() | |
| return key_id, raw_key | |
| finally: | |
| conn.close() | |
| def validate_api_key(self, raw_key): | |
| if not raw_key or not raw_key.startswith("invicta_sk_"): | |
| return None | |
| key_hash = hashlib.sha256(raw_key.encode()).hexdigest() | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute( | |
| "SELECT id, user_id, revoked FROM api_keys WHERE key_hash = ?", | |
| (key_hash,) | |
| ) | |
| row = cur.fetchone() | |
| if row and not row["revoked"]: | |
| cur.execute( | |
| "UPDATE api_keys SET last_used = ? WHERE id = ?", | |
| (datetime.now().isoformat(sep=' ', timespec='seconds'), row["id"]) | |
| ) | |
| conn.commit() | |
| return row["user_id"] | |
| return None | |
| finally: | |
| conn.close() | |
| def get_api_keys(self, user_id): | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute( | |
| "SELECT id, name, prefix, created_at, last_used, revoked, rate_limit FROM api_keys WHERE user_id = ? ORDER BY created_at DESC", | |
| (user_id,) | |
| ) | |
| rows = cur.fetchall() | |
| return [dict(r) for r in rows] | |
| finally: | |
| conn.close() | |
| def revoke_api_key(self, key_id, user_id): | |
| with self.lock: | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute("BEGIN IMMEDIATE") | |
| cur.execute( | |
| "UPDATE api_keys SET revoked = 1 WHERE id = ? AND user_id = ?", | |
| (key_id, user_id) | |
| ) | |
| conn.commit() | |
| return cur.rowcount > 0 | |
| finally: | |
| conn.close() | |
| def delete_api_key(self, key_id, user_id): | |
| with self.lock: | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute("BEGIN IMMEDIATE") | |
| cur.execute("DELETE FROM api_keys WHERE id = ? AND user_id = ?", (key_id, user_id)) | |
| conn.commit() | |
| finally: | |
| conn.close() | |
| # ββ Stats ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| def get_stats(self, user_id): | |
| conn = self._get_connection() | |
| try: | |
| cur = conn.cursor() | |
| cur.execute("SELECT COUNT(*) as c FROM conversations WHERE user_id = ?", (user_id,)) | |
| total = cur.fetchone()["c"] | |
| cur.execute("SELECT COUNT(*) as c FROM chat_sessions WHERE user_id = ?", (user_id,)) | |
| sessions = cur.fetchone()["c"] | |
| cur.execute("SELECT COUNT(*) as c FROM user_learning WHERE user_id = ?", (user_id,)) | |
| learning = cur.fetchone()["c"] | |
| return {"total_messages": total, "total_sessions": sessions, "learning_observations": learning} | |
| finally: | |
| conn.close() | |