/
dsboikov
/
aiBoardRoom
Обзор
Документация
Войти
/
dsboikov
/
aiBoardRoom
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
develop
src/storage/sqlite.py
601 строка
20 KB
Денис
Interface updated
03 авг 2026, 18:52
03 авг 2026, 18:52
57955c1
Код
Авторство
О чём код?
"""SQLite metadata repository for rooms, roles, messages, usage.""" from __future__ import annotations import json import sqlite3 from contextlib import contextmanager from dataclasses import dataclass from datetime import date, datetime, timezone from pathlib import Path from typing import Any, Generator, Iterator from config import ( DEFAULT_DAILY_LIMITS, DEFAULT_PROVIDER_MODEL, MODE_JUDGE, SQLITE_PATH, ) from storage.migrator import migrate from utils.logger import get_logger logger = get_logger("sqlite") def _utcnow() -> str: return datetime.now(timezone.utc).replace(microsecond=0).isoformat() @dataclass class Room: id: int name: str mode: str created_at: str updated_at: str starred: bool = False @dataclass class Role: id: int room_id: int name: str provider: str model: str system_prompt: str is_analyst: bool = False is_judge: bool = False sort_order: int = 0 @dataclass class Message: id: int room_id: int sender: str content: str session_id: int | None = None role_id: int | None = None model: str | None = None is_draft: bool = False is_consensus: bool = False meta_json: str | None = None created_at: str = "" class Database: """Thin SQLite access layer.""" def __init__(self, path: Path | str | None = None) -> None: self.path = Path(path or SQLITE_PATH) migrate(self.path) @contextmanager def connect(self) -> Generator[sqlite3.Connection, None, None]: conn = sqlite3.connect(str(self.path)) conn.row_factory = sqlite3.Row conn.execute("PRAGMA foreign_keys = ON") try: yield conn conn.commit() except Exception: conn.rollback() raise finally: conn.close() # --- Provider settings --- def ensure_provider_defaults(self) -> None: with self.connect() as conn: for provider, model in DEFAULT_PROVIDER_MODEL.items(): conn.execute( """ INSERT OR IGNORE INTO provider_settings (provider, model, daily_limit, enabled) VALUES (?, ?, ?, 1) """, (provider, model, DEFAULT_DAILY_LIMITS.get(provider, 0)), ) def get_provider_settings(self) -> list[dict[str, Any]]: with self.connect() as conn: rows = conn.execute( "SELECT provider, model, daily_limit, enabled FROM provider_settings ORDER BY provider" ).fetchall() return [dict(r) for r in rows] def upsert_provider_setting( self, provider: str, model: str, daily_limit: int, enabled: bool = True, ) -> None: with self.connect() as conn: conn.execute( """ INSERT INTO provider_settings (provider, model, daily_limit, enabled) VALUES (?, ?, ?, ?) ON CONFLICT(provider) DO UPDATE SET model=excluded.model, daily_limit=excluded.daily_limit, enabled=excluded.enabled """, (provider, model, daily_limit, int(enabled)), ) # --- Rooms --- def create_room(self, name: str, mode: str = MODE_JUDGE) -> Room: now = _utcnow() with self.connect() as conn: cur = conn.execute( "INSERT INTO rooms (name, mode, created_at, updated_at) VALUES (?, ?, ?, ?)", (name, mode, now, now), ) room_id = int(cur.lastrowid) return self.get_room(room_id) # type: ignore[return-value] def get_room(self, room_id: int) -> Room | None: with self.connect() as conn: row = conn.execute("SELECT * FROM rooms WHERE id = ?", (room_id,)).fetchone() return self._row_to_room(row) if row else None def list_rooms(self) -> list[Room]: with self.connect() as conn: rows = conn.execute( "SELECT * FROM rooms ORDER BY starred DESC, updated_at DESC" ).fetchall() return [self._row_to_room(r) for r in rows] def update_room( self, room_id: int, *, name: str | None = None, mode: str | None = None, starred: bool | None = None, ) -> None: fields: list[str] = ["updated_at = ?"] values: list[Any] = [_utcnow()] if name is not None: fields.append("name = ?") values.append(name) if mode is not None: fields.append("mode = ?") values.append(mode) if starred is not None: fields.append("starred = ?") values.append(int(starred)) values.append(room_id) with self.connect() as conn: conn.execute(f"UPDATE rooms SET {', '.join(fields)} WHERE id = ?", values) def delete_room(self, room_id: int) -> None: with self.connect() as conn: conn.execute("DELETE FROM rooms WHERE id = ?", (room_id,)) @staticmethod def _row_to_room(row: sqlite3.Row) -> Room: return Room( id=row["id"], name=row["name"], mode=row["mode"], created_at=row["created_at"], updated_at=row["updated_at"], starred=bool(row["starred"]), ) # --- Roles --- def create_role( self, room_id: int, name: str, provider: str, model: str, system_prompt: str = "", *, is_analyst: bool = False, is_judge: bool = False, sort_order: int = 0, ) -> Role: with self.connect() as conn: cur = conn.execute( """ INSERT INTO roles (room_id, name, provider, model, system_prompt, is_analyst, is_judge, sort_order) VALUES (?, ?, ?, ?, ?, ?, ?, ?) """, ( room_id, name, provider, model, system_prompt, int(is_analyst), int(is_judge), sort_order, ), ) role_id = int(cur.lastrowid) conn.execute( """ INSERT INTO prompt_versions (role_id, system_prompt, version) VALUES (?, ?, 1) """, (role_id, system_prompt), ) conn.execute( "UPDATE rooms SET updated_at = ? WHERE id = ?", (_utcnow(), room_id), ) return self.get_role(role_id) # type: ignore[return-value] def get_role(self, role_id: int) -> Role | None: with self.connect() as conn: row = conn.execute("SELECT * FROM roles WHERE id = ?", (role_id,)).fetchone() return self._row_to_role(row) if row else None def list_roles(self, room_id: int) -> list[Role]: with self.connect() as conn: rows = conn.execute( "SELECT * FROM roles WHERE room_id = ? ORDER BY sort_order, id", (room_id,), ).fetchall() return [self._row_to_role(r) for r in rows] def update_role( self, role_id: int, *, name: str | None = None, provider: str | None = None, model: str | None = None, system_prompt: str | None = None, is_analyst: bool | None = None, is_judge: bool | None = None, sort_order: int | None = None, ) -> None: role = self.get_role(role_id) if not role: return fields: list[str] = [] values: list[Any] = [] if name is not None: fields.append("name = ?") values.append(name) if provider is not None: fields.append("provider = ?") values.append(provider) if model is not None: fields.append("model = ?") values.append(model) if is_analyst is not None: fields.append("is_analyst = ?") values.append(int(is_analyst)) if is_judge is not None: fields.append("is_judge = ?") values.append(int(is_judge)) if sort_order is not None: fields.append("sort_order = ?") values.append(sort_order) if system_prompt is not None and system_prompt != role.system_prompt: fields.append("system_prompt = ?") values.append(system_prompt) with self.connect() as conn: ver = conn.execute( "SELECT COALESCE(MAX(version), 0) + 1 FROM prompt_versions WHERE role_id = ?", (role_id,), ).fetchone()[0] conn.execute( "INSERT INTO prompt_versions (role_id, system_prompt, version) VALUES (?, ?, ?)", (role_id, system_prompt, ver), ) if not fields: return values.append(role_id) with self.connect() as conn: conn.execute(f"UPDATE roles SET {', '.join(fields)} WHERE id = ?", values) def delete_role(self, role_id: int) -> None: with self.connect() as conn: conn.execute("DELETE FROM roles WHERE id = ?", (role_id,)) def list_prompt_versions(self, role_id: int) -> list[dict[str, Any]]: with self.connect() as conn: rows = conn.execute( """ SELECT id, version, system_prompt, created_at FROM prompt_versions WHERE role_id = ? ORDER BY version DESC """, (role_id,), ).fetchall() return [dict(r) for r in rows] @staticmethod def _row_to_role(row: sqlite3.Row) -> Role: return Role( id=row["id"], room_id=row["room_id"], name=row["name"], provider=row["provider"], model=row["model"], system_prompt=row["system_prompt"], is_analyst=bool(row["is_analyst"]), is_judge=bool(row["is_judge"]), sort_order=row["sort_order"], ) # --- Sessions / messages --- def create_session(self, room_id: int, mode: str) -> int: with self.connect() as conn: cur = conn.execute( "INSERT INTO sessions (room_id, mode, status) VALUES (?, ?, 'running')", (room_id, mode), ) return int(cur.lastrowid) def finish_session(self, session_id: int, status: str = "done") -> None: with self.connect() as conn: conn.execute( "UPDATE sessions SET status = ?, finished_at = ? WHERE id = ?", (status, _utcnow(), session_id), ) def add_message( self, room_id: int, sender: str, content: str, *, session_id: int | None = None, role_id: int | None = None, model: str | None = None, is_draft: bool = False, is_consensus: bool = False, meta: dict[str, Any] | None = None, ) -> Message: meta_json = json.dumps(meta, ensure_ascii=False) if meta else None now = _utcnow() with self.connect() as conn: cur = conn.execute( """ INSERT INTO messages (room_id, session_id, role_id, sender, content, model, is_draft, is_consensus, meta_json, created_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) """, ( room_id, session_id, role_id, sender, content, model, int(is_draft), int(is_consensus), meta_json, now, ), ) msg_id = int(cur.lastrowid) conn.execute( "UPDATE rooms SET updated_at = ? WHERE id = ?", (now, room_id), ) return Message( id=msg_id, room_id=room_id, session_id=session_id, role_id=role_id, sender=sender, content=content, model=model, is_draft=is_draft, is_consensus=is_consensus, meta_json=meta_json, created_at=now, ) def list_messages(self, room_id: int, *, include_drafts: bool = True) -> list[Message]: query = "SELECT * FROM messages WHERE room_id = ?" if not include_drafts: query += " AND is_draft = 0" query += " ORDER BY created_at, id" with self.connect() as conn: rows = conn.execute(query, (room_id,)).fetchall() return [self._row_to_message(r) for r in rows] def clear_room_messages(self, room_id: int) -> None: with self.connect() as conn: conn.execute("DELETE FROM messages WHERE room_id = ?", (room_id,)) conn.execute("DELETE FROM sessions WHERE room_id = ?", (room_id,)) @staticmethod def _row_to_message(row: sqlite3.Row) -> Message: return Message( id=row["id"], room_id=row["room_id"], session_id=row["session_id"], role_id=row["role_id"], sender=row["sender"], content=row["content"], model=row["model"], is_draft=bool(row["is_draft"]), is_consensus=bool(row["is_consensus"]), meta_json=row["meta_json"], created_at=row["created_at"], ) # --- Usage / limits --- def add_usage( self, provider: str, model: str, input_tokens: int, output_tokens: int, cost_usd: float, day: str | None = None, ) -> None: day = day or date.today().isoformat() with self.connect() as conn: conn.execute( """ INSERT INTO usage_daily (day, provider, model, input_tokens, output_tokens, cost_usd) VALUES (?, ?, ?, ?, ?, ?) ON CONFLICT(day, provider, model) DO UPDATE SET input_tokens = input_tokens + excluded.input_tokens, output_tokens = output_tokens + excluded.output_tokens, cost_usd = cost_usd + excluded.cost_usd """, (day, provider, model, input_tokens, output_tokens, cost_usd), ) def get_daily_usage(self, provider: str, day: str | None = None) -> int: day = day or date.today().isoformat() with self.connect() as conn: row = conn.execute( """ SELECT COALESCE(SUM(input_tokens + output_tokens), 0) AS total FROM usage_daily WHERE day = ? AND provider = ? """, (day, provider), ).fetchone() return int(row["total"]) def get_usage_report(self, day: str | None = None) -> list[dict[str, Any]]: day = day or date.today().isoformat() with self.connect() as conn: rows = conn.execute( """ SELECT provider, model, input_tokens, output_tokens, cost_usd FROM usage_daily WHERE day = ? ORDER BY provider, model """, (day,), ).fetchall() return [dict(r) for r in rows] def get_usage_history(self, days: int = 14) -> list[dict[str, Any]]: """Per-day totals for the last *days* calendar days (oldest first).""" with self.connect() as conn: rows = conn.execute( """ SELECT day, SUM(input_tokens) AS input_tokens, SUM(output_tokens) AS output_tokens, SUM(cost_usd) AS cost_usd FROM usage_daily WHERE day >= date('now', ?) GROUP BY day ORDER BY day ASC """, (f"-{max(1, days) - 1} days",), ).fetchall() return [dict(r) for r in rows] def get_usage_by_provider(self, day: str | None = None) -> list[dict[str, Any]]: """Aggregate today's usage per provider (all models).""" day = day or date.today().isoformat() with self.connect() as conn: rows = conn.execute( """ SELECT provider, SUM(input_tokens) AS input_tokens, SUM(output_tokens) AS output_tokens, SUM(cost_usd) AS cost_usd FROM usage_daily WHERE day = ? GROUP BY provider ORDER BY (SUM(input_tokens) + SUM(output_tokens)) DESC """, (day,), ).fetchall() return [dict(r) for r in rows] def get_daily_limit(self, provider: str) -> int: with self.connect() as conn: row = conn.execute( "SELECT daily_limit FROM provider_settings WHERE provider = ?", (provider,), ).fetchone() if row: return int(row["daily_limit"]) return DEFAULT_DAILY_LIMITS.get(provider, 0) def export_settings_payload(self) -> dict[str, Any]: """Export rooms, roles, prompts, provider settings (no API keys).""" rooms_out: list[dict[str, Any]] = [] for room in self.list_rooms(): roles = [ { "name": r.name, "provider": r.provider, "model": r.model, "system_prompt": r.system_prompt, "is_analyst": r.is_analyst, "is_judge": r.is_judge, "sort_order": r.sort_order, "prompt_versions": self.list_prompt_versions(r.id), } for r in self.list_roles(room.id) ] rooms_out.append( { "name": room.name, "mode": room.mode, "starred": room.starred, "roles": roles, } ) return { "version": 1, "exported_at": _utcnow(), "provider_settings": self.get_provider_settings(), "rooms": rooms_out, } def import_settings_payload(self, payload: dict[str, Any], *, merge: bool = True) -> None: """Import rooms/roles/settings from an export payload.""" for ps in payload.get("provider_settings", []): self.upsert_provider_setting( ps["provider"], ps["model"], int(ps.get("daily_limit", 0)), bool(ps.get("enabled", True)), ) if not merge: for room in self.list_rooms(): self.delete_room(room.id) for room_data in payload.get("rooms", []): room = self.create_room(room_data["name"], room_data.get("mode", MODE_JUDGE)) if room_data.get("starred"): self.update_room(room.id, starred=True) for role_data in room_data.get("roles", []): self.create_role( room.id, role_data["name"], role_data["provider"], role_data["model"], role_data.get("system_prompt", ""), is_analyst=bool(role_data.get("is_analyst")), is_judge=bool(role_data.get("is_judge")), sort_order=int(role_data.get("sort_order", 0)), ) _db: Database | None = None def get_db(path: Path | str | None = None) -> Database: """Return process-wide Database singleton (or a new one for a custom path).""" global _db if path is not None: return Database(path) if _db is None: _db = Database() _db.ensure_provider_defaults() return _db