/
a1ex_ma
/
ds2
Обзор
Документация
Войти
/
a1ex_ma
/
ds2
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
Аналитика
Безопасность
main
src/adapters/sqlite_adapter.py
512 строк
19 KB
A1eksMa
Alpha 0.2.0a1: auto-detect act on load, loader docs, release archives
06 июл 2026, 22:28
06 июл 2026, 22:28
029e5f2
Код
Авторство
О чём код?
from __future__ import annotations import sqlite3 from typing import Dict, List, Optional, Tuple, Union from src.domain.entities import ( ActId, CntId, IdId, LbId, SrcId, ValId, Lb, Src, Transaction, TransactionInput, ) from src.domain.enums import Act from src.domain.errors import StorageError from src.domain.result import Err, Ok from src.persistence.mappers import lb_mapper, src_mapper, transaction_mapper from src.persistence.records.lb_record import LbRecord from src.persistence.records.src_record import SrcRecord from src.persistence.records.transaction_record import TransactionRecord _SCHEMA = """ PRAGMA foreign_keys = ON; CREATE TABLE IF NOT EXISTS lbs ( lb_id INTEGER PRIMARY KEY AUTOINCREMENT, name TEXT UNIQUE NOT NULL, description TEXT, p REAL NOT NULL DEFAULT 0.5 ); CREATE TABLE IF NOT EXISTS acts ( act_id INTEGER PRIMARY KEY AUTOINCREMENT, name TEXT UNIQUE NOT NULL ); CREATE TABLE IF NOT EXISTS srcs ( src_id INTEGER PRIMARY KEY AUTOINCREMENT, name TEXT UNIQUE NOT NULL, description TEXT, p REAL NOT NULL DEFAULT 0.5, key_label INTEGER NOT NULL, FOREIGN KEY (key_label) REFERENCES lbs(lb_id) ); CREATE TABLE IF NOT EXISTS ids ( id_id INTEGER PRIMARY KEY AUTOINCREMENT, value TEXT UNIQUE NOT NULL ); CREATE TABLE IF NOT EXISTS vals ( val_id INTEGER PRIMARY KEY AUTOINCREMENT, value TEXT UNIQUE NOT NULL ); CREATE TABLE IF NOT EXISTS cnts ( cnt_id INTEGER PRIMARY KEY AUTOINCREMENT, created_at REAL NOT NULL DEFAULT (julianday('now')) ); CREATE TABLE IF NOT EXISTS transactions ( cnt INTEGER PRIMARY KEY, act INTEGER NOT NULL, dt REAL NOT NULL, src INTEGER NOT NULL, lb INTEGER NOT NULL, id INTEGER NOT NULL, val INTEGER, p REAL NOT NULL DEFAULT 1.0, FOREIGN KEY (cnt) REFERENCES cnts(cnt_id) ON DELETE RESTRICT, FOREIGN KEY (act) REFERENCES acts(act_id) ON DELETE RESTRICT, FOREIGN KEY (src) REFERENCES srcs(src_id) ON DELETE RESTRICT, FOREIGN KEY (lb) REFERENCES lbs(lb_id) ON DELETE RESTRICT, FOREIGN KEY (id) REFERENCES ids(id_id) ON DELETE RESTRICT, FOREIGN KEY (val) REFERENCES vals(val_id) ON DELETE RESTRICT ); CREATE TABLE IF NOT EXISTS transactions_archive ( cnt INTEGER PRIMARY KEY, act INTEGER NOT NULL, dt REAL NOT NULL, src INTEGER NOT NULL, lb INTEGER NOT NULL, id INTEGER NOT NULL, val INTEGER, p REAL NOT NULL DEFAULT 1.0, FOREIGN KEY (cnt) REFERENCES cnts(cnt_id) ON DELETE RESTRICT, FOREIGN KEY (act) REFERENCES acts(act_id) ON DELETE RESTRICT, FOREIGN KEY (src) REFERENCES srcs(src_id) ON DELETE RESTRICT, FOREIGN KEY (lb) REFERENCES lbs(lb_id) ON DELETE RESTRICT, FOREIGN KEY (id) REFERENCES ids(id_id) ON DELETE RESTRICT, FOREIGN KEY (val) REFERENCES vals(val_id) ON DELETE RESTRICT ); CREATE INDEX IF NOT EXISTS idx_txn_dt ON transactions(dt); CREATE INDEX IF NOT EXISTS idx_txn_src_lb_id ON transactions(src, lb, id); CREATE INDEX IF NOT EXISTS idx_arch_dt ON transactions_archive(dt); CREATE INDEX IF NOT EXISTS idx_arch_src_lb_id ON transactions_archive(src, lb, id); CREATE VIEW IF NOT EXISTS transactions_full AS SELECT * FROM transactions UNION ALL SELECT * FROM transactions_archive; """ _SQLITE_MAX_VARS = 900 # stay comfortably under SQLite's bound-parameter limit def _chunks(values: List[int], size: int) -> List[List[int]]: return [values[i:i + size] for i in range(0, len(values), size)] or [[]] class SQLiteAdapter: def __init__(self, db_path: str) -> None: self._conn = sqlite3.connect(db_path, check_same_thread=False) self._conn.executescript(_SCHEMA) self._conn.commit() self._conn.isolation_level = None # autocommit; transactions managed explicitly self._conn.execute("PRAGMA journal_mode=WAL") self._conn.execute("PRAGMA synchronous=NORMAL") self._lb_cache: Dict[str, int] = {} self._act_cache: Dict[str, int] = {} self._src_cache: Dict[str, int] = {} self._id_cache: Dict[str, int] = {} self._val_cache: Dict[str, int] = {} self._load_small_pools() def _load_small_pools(self) -> None: for row in self._conn.execute("SELECT name, lb_id FROM lbs"): self._lb_cache[row[0]] = row[1] for row in self._conn.execute("SELECT name, act_id FROM acts"): self._act_cache[row[0]] = row[1] for row in self._conn.execute("SELECT name, src_id FROM srcs"): self._src_cache[row[0]] = row[1] for row in self._conn.execute("SELECT value, id_id FROM ids"): self._id_cache[row[0]] = row[1] # --- String pool operations --- def begin(self) -> None: self._conn.execute("BEGIN") def commit(self) -> None: self._conn.commit() def rollback(self) -> None: self._conn.rollback() def lb_intern(self, name: str) -> Union[Ok[LbId], Err[StorageError]]: try: if name in self._lb_cache: return Ok(LbId(self._lb_cache[name])) self._conn.execute( "INSERT OR IGNORE INTO lbs (name) VALUES (?)", (name,) ) row = self._conn.execute( "SELECT lb_id FROM lbs WHERE name = ?", (name,) ).fetchone() lb_id = row[0] self._lb_cache[name] = lb_id return Ok(LbId(lb_id)) except sqlite3.Error as exc: return Err(StorageError(str(exc))) def id_intern(self, value: str) -> Union[Ok[IdId], Err[StorageError]]: try: if value in self._id_cache: return Ok(IdId(self._id_cache[value])) self._conn.execute( "INSERT OR IGNORE INTO ids (value) VALUES (?)", (value,) ) row = self._conn.execute( "SELECT id_id FROM ids WHERE value = ?", (value,) ).fetchone() id_id = row[0] self._id_cache[value] = id_id return Ok(IdId(id_id)) except sqlite3.Error as exc: return Err(StorageError(str(exc))) def id_get(self, id_id: IdId) -> Union[Ok[str], Err[StorageError]]: try: row = self._conn.execute( "SELECT value FROM ids WHERE id_id = ?", (int(id_id),) ).fetchone() if row is None: return Err(StorageError(f"id {id_id} not found")) return Ok(row[0]) except sqlite3.Error as exc: return Err(StorageError(str(exc))) def val_intern(self, value: str) -> Union[Ok[ValId], Err[StorageError]]: try: if value in self._val_cache: return Ok(ValId(self._val_cache[value])) self._conn.execute( "INSERT OR IGNORE INTO vals (value) VALUES (?)", (value,) ) row = self._conn.execute( "SELECT val_id FROM vals WHERE value = ?", (value,) ).fetchone() val_id = row[0] self._val_cache[value] = val_id return Ok(ValId(val_id)) except sqlite3.Error as exc: return Err(StorageError(str(exc))) def val_get(self, val_id: ValId) -> Union[Ok[str], Err[StorageError]]: try: row = self._conn.execute( "SELECT value FROM vals WHERE val_id = ?", (int(val_id),) ).fetchone() if row is None: return Err(StorageError(f"val {val_id} not found")) return Ok(row[0]) except sqlite3.Error as exc: return Err(StorageError(str(exc))) def act_intern(self, act: Act) -> Union[Ok[ActId], Err[StorageError]]: try: name = act.value if name in self._act_cache: return Ok(ActId(self._act_cache[name])) self._conn.execute( "INSERT OR IGNORE INTO acts (name) VALUES (?)", (name,) ) row = self._conn.execute( "SELECT act_id FROM acts WHERE name = ?", (name,) ).fetchone() act_id = row[0] self._act_cache[name] = act_id return Ok(ActId(act_id)) except sqlite3.Error as exc: return Err(StorageError(str(exc))) # --- Source metadata --- def src_get_or_create( self, name: str, key_label_id: LbId ) -> Union[Ok[Src], Err[StorageError]]: try: if name in self._src_cache: return self.src_get(SrcId(self._src_cache[name])) self._conn.execute( "INSERT OR IGNORE INTO srcs (name, key_label) VALUES (?, ?)", (name, int(key_label_id)), ) row = self._conn.execute( "SELECT src_id, name, description, p, key_label FROM srcs WHERE name = ?", (name,), ).fetchone() record = SrcRecord.from_row(row) self._src_cache[name] = record.src_id result = src_mapper.record_to_domain(record) if isinstance(result, Err): return Err(StorageError(result.error.message)) return result except sqlite3.Error as exc: return Err(StorageError(str(exc))) def src_get(self, src_id: SrcId) -> Union[Ok[Src], Err[StorageError]]: try: row = self._conn.execute( "SELECT src_id, name, description, p, key_label FROM srcs WHERE src_id = ?", (int(src_id),), ).fetchone() if row is None: return Err(StorageError(f"src {src_id} not found")) result = src_mapper.record_to_domain(SrcRecord.from_row(row)) if isinstance(result, Err): return Err(StorageError(result.error.message)) return result except sqlite3.Error as exc: return Err(StorageError(str(exc))) def src_update(self, src: Src) -> Union[Ok[None], Err[StorageError]]: try: self._conn.execute( "UPDATE srcs SET description = ?, p = ?, key_label = ? WHERE src_id = ?", (src.description, src.p, int(src.key_label), int(src.src_id)), ) self._conn.commit() # standalone call — commit immediately return Ok(None) except sqlite3.Error as exc: return Err(StorageError(str(exc))) def src_list(self) -> Union[Ok[List[Src]], Err[StorageError]]: try: rows = self._conn.execute( "SELECT src_id, name, description, p, key_label FROM srcs" ).fetchall() result: List[Src] = [] for row in rows: mapped = src_mapper.record_to_domain(SrcRecord.from_row(row)) if isinstance(mapped, Err): return Err(StorageError(mapped.error.message)) result.append(mapped.value) return Ok(result) except sqlite3.Error as exc: return Err(StorageError(str(exc))) # --- Label metadata --- def lb_get(self, lb_id: LbId) -> Union[Ok[Lb], Err[StorageError]]: try: row = self._conn.execute( "SELECT lb_id, name, description, p FROM lbs WHERE lb_id = ?", (int(lb_id),), ).fetchone() if row is None: return Err(StorageError(f"lb {lb_id} not found")) result = lb_mapper.record_to_domain(LbRecord.from_row(row)) if isinstance(result, Err): return Err(StorageError(result.error.message)) return result except sqlite3.Error as exc: return Err(StorageError(str(exc))) def lb_update(self, lb: Lb) -> Union[Ok[None], Err[StorageError]]: try: self._conn.execute( "UPDATE lbs SET description = ?, p = ? WHERE lb_id = ?", (lb.description, lb.p, int(lb.lb_id)), ) self._conn.commit() # standalone call — commit immediately return Ok(None) except sqlite3.Error as exc: return Err(StorageError(str(exc))) def lb_list(self) -> Union[Ok[List[Lb]], Err[StorageError]]: try: rows = self._conn.execute( "SELECT lb_id, name, description, p FROM lbs" ).fetchall() result: List[Lb] = [] for row in rows: mapped = lb_mapper.record_to_domain(LbRecord.from_row(row)) if isinstance(mapped, Err): return Err(StorageError(mapped.error.message)) result.append(mapped.value) return Ok(result) except sqlite3.Error as exc: return Err(StorageError(str(exc))) # --- Transactions --- def txn_insert( self, txn: TransactionInput ) -> Union[Ok[Transaction], Err[StorageError]]: try: cur = self._conn.execute("INSERT INTO cnts (created_at) VALUES (julianday('now'))") cnt_id = cur.lastrowid self._conn.execute( "INSERT INTO transactions (cnt, act, dt, src, lb, id, val, p) " "VALUES (?, ?, ?, ?, ?, ?, ?, ?)", ( cnt_id, int(txn.act), txn.dt, int(txn.src), int(txn.lb), int(txn.id), int(txn.val) or None, # 0 → NULL in SQLite (DELETE semantics) txn.p, ), ) result = transaction_mapper.record_to_domain( TransactionRecord( cnt=cnt_id, act=int(txn.act), dt=txn.dt, src=int(txn.src), lb=int(txn.lb), id=int(txn.id), p=txn.p, val=int(txn.val) or None, # 0 → NULL in SQLite (DELETE semantics) ) ) if isinstance(result, Err): return Err(StorageError(result.error.message)) return result except sqlite3.Error as exc: return Err(StorageError(str(exc))) def txn_query( self, src_id: Optional[SrcId] = None, lb_id: Optional[LbId] = None, id_id: Optional[IdId] = None, until_dt: Optional[float] = None, from_cnt: Optional[CntId] = None, include_archive: bool = False, ) -> Union[Ok[List[Transaction]], Err[StorageError]]: try: table = "transactions_full" if include_archive else "transactions" clauses = [] params = [] if src_id is not None: clauses.append("src = ?") params.append(int(src_id)) if lb_id is not None: clauses.append("lb = ?") params.append(int(lb_id)) if id_id is not None: clauses.append("id = ?") params.append(int(id_id)) if until_dt is not None: clauses.append("dt <= ?") params.append(until_dt) if from_cnt is not None: clauses.append("cnt > ?") params.append(int(from_cnt)) where = ("WHERE " + " AND ".join(clauses)) if clauses else "" rows = self._conn.execute( f"SELECT cnt, act, dt, src, lb, id, val, p FROM {table} {where} ORDER BY cnt", params, ).fetchall() result: List[Transaction] = [] for row in rows: mapped = transaction_mapper.record_to_domain(TransactionRecord.from_row(row)) if isinstance(mapped, Err): return Err(StorageError(mapped.error.message)) result.append(mapped.value) return Ok(result) except sqlite3.Error as exc: return Err(StorageError(str(exc))) def txn_archive(self, until_dt: float) -> Union[Ok[int], Err[StorageError]]: try: self._conn.execute("BEGIN") cur = self._conn.execute( "INSERT INTO transactions_archive " "SELECT cnt, act, dt, src, lb, id, val, p " "FROM transactions WHERE dt <= ?", (until_dt,), ) count = cur.rowcount self._conn.execute("DELETE FROM transactions WHERE dt <= ?", (until_dt,)) self._conn.commit() return Ok(count) except sqlite3.Error as exc: self._conn.rollback() return Err(StorageError(str(exc))) def txn_delete( self, src_id: SrcId, lb_id: Optional[LbId] = None, ) -> Union[Ok[int], Err[StorageError]]: try: self._conn.execute("BEGIN") if lb_id is not None: cur = self._conn.execute( "DELETE FROM transactions WHERE src = ? AND lb = ?", (int(src_id), int(lb_id)), ) self._conn.execute( "DELETE FROM transactions_archive WHERE src = ? AND lb = ?", (int(src_id), int(lb_id)), ) else: cur = self._conn.execute( "DELETE FROM transactions WHERE src = ?", (int(src_id),) ) self._conn.execute( "DELETE FROM transactions_archive WHERE src = ?", (int(src_id),) ) self._conn.commit() return Ok(cur.rowcount) except sqlite3.Error as exc: self._conn.rollback() return Err(StorageError(str(exc))) # --- Bulk history lookup --- def txn_last_values( self, src_id: SrcId, lb_ids: List[LbId], id_ids: List[IdId], ) -> Union[Ok[Dict[Tuple[LbId, IdId], ValId]], Err[StorageError]]: lb_list = sorted({int(x) for x in lb_ids}) id_list = sorted({int(x) for x in id_ids}) if not lb_list or not id_list: return Ok({}) try: result: Dict[Tuple[LbId, IdId], ValId] = {} for lb_chunk in _chunks(lb_list, _SQLITE_MAX_VARS): if not lb_chunk: continue for id_chunk in _chunks(id_list, _SQLITE_MAX_VARS): if not id_chunk: continue lb_ph = ",".join("?" * len(lb_chunk)) id_ph = ",".join("?" * len(id_chunk)) rows = self._conn.execute( f""" SELECT t.lb, t.id, t.val FROM transactions_full t JOIN ( SELECT lb, id, MAX(cnt) AS mx FROM transactions_full WHERE src = ? AND lb IN ({lb_ph}) AND id IN ({id_ph}) GROUP BY lb, id ) last ON t.lb = last.lb AND t.id = last.id AND t.cnt = last.mx WHERE t.src = ? """, [int(src_id), *lb_chunk, *id_chunk, int(src_id)], ).fetchall() for lb, id_, val in rows: result[(LbId(lb), IdId(id_))] = ValId(val or 0) return Ok(result) except sqlite3.Error as exc: return Err(StorageError(str(exc)))