From 069ea848e26fddfc0d1bf443b76ca3fcc7bd8eff Mon Sep 17 00:00:00 2001 From: jze9 Date: Mon, 6 Jul 2026 20:21:21 +0500 Subject: [PATCH] Add SQLite results storage with a resilient multiprocess writer (Stage 5) - storage/schema.py: `runs` table recording every evaluated configuration, feasible or not, with an honest infeasible_reason instead of dropping failures - storage/database.py: WAL-mode SQLite, single-writer-over-a-queue pattern for safe concurrent writes from many sweep worker processes - Testing with real multiprocessing.Process (not mocks) caught a genuine bug: a duplicate run_id raised inside the writer loop and killed the writer process, silently halting queue drainage for the rest of a sweep. Fixed by catching the insert error per-record and logging to stderr instead of crashing; added a regression test for it. - Verified locally and inside the Docker image (43/43 tests both ways) Co-Authored-By: Claude Sonnet 5 --- PLAN.md | 2 +- src/gausse/storage/database.py | 116 ++++++++++++++++++++++++++++++ src/gausse/storage/schema.py | 50 +++++++++++++ tests/test_storage_database.py | 126 +++++++++++++++++++++++++++++++++ 4 files changed, 293 insertions(+), 1 deletion(-) create mode 100644 src/gausse/storage/database.py create mode 100644 src/gausse/storage/schema.py create mode 100644 tests/test_storage_database.py diff --git a/PLAN.md b/PLAN.md index 139a9c2..6a9b60d 100644 --- a/PLAN.md +++ b/PLAN.md @@ -47,7 +47,7 @@ L(x, I)-модель с coenergy-выводом силы — в разделе " - [x] **Этап 2 — Физическое ядро**: `physics/constants.py`, `inductance.py` (Уилер + ферромагнитный сердечник + размагничивание), `force.py`, `circuit.py` (ОДУ RLC). Юнит-тесты: аналитическое RLC-решение, согласованность dL/dx. - [x] **Этап 3 — Датчики и одна ступень**: `physics/sensors.py` (оба типа как события solve_ivp), `sim/stage.py` (полёт → триггер → разряд → энергобаланс). Тест на сохранение энергии + найден/исправлен баг насыщения (см. выше). - [x] **Этап 4 — Многоступенчатая цепочка**: `sim/coilgun.py` — сквозная координата, отбраковка нереализуемых конфигураций. Дымовой тест на реальной базе компонентов (`test_real_components_smoke.py`) подтверждает: полный путь реальные JSON → физика → цепочка работает (пример: 22.3 м/с, КПД 3.4% — честный неоптимизированный результат). -- [ ] **Этап 5 — Хранилище результатов**: `storage/schema.py` + `database.py` — SQLite (WAL), таблица `runs` со всеми прогонами (успех/провал + честная причина), однопроцессный писатель поверх многопроцессной очереди. +- [x] **Этап 5 — Хранилище результатов**: `storage/schema.py` + `database.py` — SQLite (WAL), таблица `runs` со всеми прогонами (успех/провал + честная причина), однопроцессный писатель поверх многопроцессной очереди. Тест с реальными `multiprocessing.Process` (не моками) поймал реальную проблему: коллизия `run_id` роняла writer и молча останавливала осушение очереди на весь sweep — писатель теперь переживает ошибку вставки одной записи (лог в stderr) и продолжает работу. - [ ] **Этап 6 — Поиск и оптимизация**: `optim/search_space.py` (геном переменной длины), `objective.py` (КПД), дешёвый квазистатический предфильтр, `optim/sweep.py` (Monte Carlo/LHS, миллионы прогонов), `optim/evolutionary.py` (ГА + coordinate-descent полировка) — всё пишется в общую таблицу `runs`. - [ ] **Этап 7 — Отчётность**: `report/plots.py`, `summary.py`, `bom.py`, раздел "ограничения модели", `cli.py` (`gausse sweep/evolve/simulate/report`). - [ ] `report/animate.py` — анимация одного прогона: положение снаряда в трубе по времени + визуализация поля/тока каждой катушки (свечение/интенсивность цвета ~ ток), сохранение в GIF/MP4 (matplotlib `FuncAnimation`). Нужна по запросу пользователя — "графика где будет показана симуляция пролёта цилиндра по трубе и электромагнитные поля в каждый момент времени". diff --git a/src/gausse/storage/database.py b/src/gausse/storage/database.py new file mode 100644 index 0000000..51657da --- /dev/null +++ b/src/gausse/storage/database.py @@ -0,0 +1,116 @@ +"""SQLite-хранилище прогонов: WAL-режим + однопроцессный писатель. + +Множество воркеров (Stage 6, `ProcessPoolExecutor`) не пишут в SQLite +напрямую — каждый кладёт `RunRecord` в `multiprocessing.Queue`, а +единственный процесс-писатель (`run_writer_process`) читает очередь и +пишет в базу последовательно. Это исключает конкурентную запись в один +файл SQLite, которая иначе могла бы повредить базу или потерять прогоны. +""" + +import sqlite3 +import sys +from dataclasses import asdict, fields +from pathlib import Path + +from gausse.storage.schema import CREATE_INDEXES_SQL, CREATE_RUNS_TABLE_SQL, RunRecord + +_COLUMNS = [f.name for f in fields(RunRecord)] + + +def open_connection(db_path: Path) -> sqlite3.Connection: + db_path = Path(db_path) + db_path.parent.mkdir(parents=True, exist_ok=True) + conn = sqlite3.connect(str(db_path)) + conn.execute("PRAGMA journal_mode=WAL") + conn.execute(CREATE_RUNS_TABLE_SQL) + for statement in CREATE_INDEXES_SQL: + conn.execute(statement) + conn.commit() + return conn + + +def _record_to_params(record: RunRecord) -> dict: + row = asdict(record) + row["feasible"] = int(row["feasible"]) + return row + + +def _row_to_record(row: sqlite3.Row) -> RunRecord: + data = dict(row) + data["feasible"] = bool(data["feasible"]) + return RunRecord(**data) + + +def insert_run(conn: sqlite3.Connection, record: RunRecord) -> None: + insert_runs(conn, [record]) + + +def insert_runs(conn: sqlite3.Connection, records: list[RunRecord]) -> None: + if not records: + return + placeholders = ", ".join(f":{name}" for name in _COLUMNS) + columns = ", ".join(_COLUMNS) + conn.executemany( + f"INSERT INTO runs ({columns}) VALUES ({placeholders})", + [_record_to_params(r) for r in records], + ) + conn.commit() + + +def count_runs(conn: sqlite3.Connection, feasible: bool | None = None) -> int: + if feasible is None: + cursor = conn.execute("SELECT COUNT(*) FROM runs") + else: + cursor = conn.execute("SELECT COUNT(*) FROM runs WHERE feasible = ?", (int(feasible),)) + return cursor.fetchone()[0] + + +def fetch_runs( + conn: sqlite3.Connection, + feasible: bool | None = None, + min_efficiency: float | None = None, + order_by_efficiency_desc: bool = False, + limit: int | None = None, +) -> list[RunRecord]: + conn.row_factory = sqlite3.Row + query = "SELECT * FROM runs WHERE 1=1" + params: list = [] + if feasible is not None: + query += " AND feasible = ?" + params.append(int(feasible)) + if min_efficiency is not None: + query += " AND efficiency >= ?" + params.append(min_efficiency) + if order_by_efficiency_desc: + query += " ORDER BY efficiency DESC" + if limit is not None: + query += " LIMIT ?" + params.append(limit) + cursor = conn.execute(query, params) + return [_row_to_record(row) for row in cursor.fetchall()] + + +def run_writer_process(queue, db_path: Path) -> None: + """Читает `RunRecord` из очереди и пишет в SQLite, пока не придёт None. + + Предназначен для запуска как отдельный `multiprocessing.Process` — + единственный writer на файл базы, пока воркеры sweep/evolve только + кладут результаты в очередь. + + Один "плохой" прогон (например, коллизия run_id) не должен уронить + писатель и остановить осушение очереди на много часов sweep — при + ошибке вставки конкретная запись пропускается с сообщением в stderr, + а не молча и не ценой падения всего процесса. + """ + conn = open_connection(db_path) + try: + while True: + record = queue.get() + if record is None: + break + try: + insert_run(conn, record) + except sqlite3.Error as exc: + print(f"[storage] не удалось записать run_id={record.run_id}: {exc}", file=sys.stderr) + finally: + conn.close() diff --git a/src/gausse/storage/schema.py b/src/gausse/storage/schema.py new file mode 100644 index 0000000..fb2524e --- /dev/null +++ b/src/gausse/storage/schema.py @@ -0,0 +1,50 @@ +"""Схема таблицы `runs` — ВСЕ прогоны (успешные и нет), без прикрас. + +Нереализуемые конфигурации хранятся наравне с успешными: `feasible=0` + +`infeasible_reason` с честной причиной. Ничего не отбрасывается перед +записью — отбор "хороших" результатов делается запросом к этой таблице +постфактум, а не фильтрацией на входе. +""" + +from dataclasses import dataclass + +CREATE_RUNS_TABLE_SQL = """ +CREATE TABLE IF NOT EXISTS runs ( + run_id TEXT PRIMARY KEY, + timestamp TEXT NOT NULL, + search_mode TEXT NOT NULL, + genome_json TEXT NOT NULL, + decoded_summary_json TEXT NOT NULL, + feasible INTEGER NOT NULL, + infeasible_reason TEXT, + failed_stage_index INTEGER, + efficiency REAL, + exit_velocity_mps REAL, + cost_rub REAL, + energy_breakdown_json TEXT, + model_version TEXT NOT NULL +) +""" + +CREATE_INDEXES_SQL = [ + "CREATE INDEX IF NOT EXISTS idx_runs_feasible ON runs(feasible)", + "CREATE INDEX IF NOT EXISTS idx_runs_efficiency ON runs(efficiency)", + "CREATE INDEX IF NOT EXISTS idx_runs_search_mode ON runs(search_mode)", +] + + +@dataclass(frozen=True) +class RunRecord: + run_id: str + timestamp: str + search_mode: str # "sweep" | "evolve" | "manual" + genome_json: str + decoded_summary_json: str + feasible: bool + model_version: str + infeasible_reason: str | None = None + failed_stage_index: int | None = None + efficiency: float | None = None + exit_velocity_mps: float | None = None + cost_rub: float | None = None + energy_breakdown_json: str | None = None diff --git a/tests/test_storage_database.py b/tests/test_storage_database.py new file mode 100644 index 0000000..e4d6fa6 --- /dev/null +++ b/tests/test_storage_database.py @@ -0,0 +1,126 @@ +import multiprocessing +from pathlib import Path + +from gausse.storage.database import ( + count_runs, + fetch_runs, + insert_run, + insert_runs, + open_connection, + run_writer_process, +) +from gausse.storage.schema import RunRecord + + +def _record(run_id: str, feasible: bool = True, efficiency: float | None = 0.1, reason=None) -> RunRecord: + return RunRecord( + run_id=run_id, + timestamp="2026-07-06T00:00:00", + search_mode="sweep", + genome_json="{}", + decoded_summary_json="{}", + feasible=feasible, + model_version="v1", + infeasible_reason=reason, + efficiency=efficiency, + exit_velocity_mps=20.0 if feasible else None, + cost_rub=500.0, + ) + + +def test_insert_and_fetch_roundtrip(tmp_path): + conn = open_connection(tmp_path / "runs.sqlite3") + insert_run(conn, _record("r1")) + rows = fetch_runs(conn) + assert len(rows) == 1 + assert rows[0].run_id == "r1" + assert rows[0].feasible is True + conn.close() + + +def test_infeasible_runs_are_stored_not_dropped(tmp_path): + conn = open_connection(tmp_path / "runs.sqlite3") + insert_run(conn, _record("r-fail", feasible=False, efficiency=None, reason="датчик не сработал")) + assert count_runs(conn) == 1 + assert count_runs(conn, feasible=False) == 1 + assert count_runs(conn, feasible=True) == 0 + row = fetch_runs(conn, feasible=False)[0] + assert row.infeasible_reason == "датчик не сработал" + conn.close() + + +def test_batch_insert_and_filters(tmp_path): + conn = open_connection(tmp_path / "runs.sqlite3") + records = [ + _record("r1", efficiency=0.05), + _record("r2", efficiency=0.30), + _record("r3", feasible=False, efficiency=None, reason="разряд не скоммутировался"), + ] + insert_runs(conn, records) + assert count_runs(conn) == 3 + best = fetch_runs(conn, feasible=True, min_efficiency=0.1) + assert [r.run_id for r in best] == ["r2"] + ranked = fetch_runs(conn, feasible=True, order_by_efficiency_desc=True) + assert [r.run_id for r in ranked] == ["r2", "r1"] + conn.close() + + +def test_reopening_database_preserves_previous_runs(tmp_path): + db_path = tmp_path / "runs.sqlite3" + conn = open_connection(db_path) + insert_run(conn, _record("r1")) + conn.close() + + conn2 = open_connection(db_path) + assert count_runs(conn2) == 1 + conn2.close() + + +def test_writer_survives_duplicate_run_id_and_keeps_draining_queue(tmp_path): + """Регрессия: одна коллизия run_id раньше роняла writer-процесс и молча + останавливала осушение очереди (см. историю разработки Этапа 5).""" + db_path = tmp_path / "runs.sqlite3" + ctx = multiprocessing.get_context("spawn") + queue = ctx.Queue() + + writer = ctx.Process(target=run_writer_process, args=(queue, db_path)) + writer.start() + + queue.put(_record("dup")) + queue.put(_record("dup")) # коллизия PRIMARY KEY + queue.put(_record("after-collision")) + queue.put(None) + writer.join(timeout=10) + + assert not writer.is_alive() + conn = open_connection(db_path) + assert count_runs(conn) == 2 # "dup" один раз + "after-collision" + assert {r.run_id for r in fetch_runs(conn)} == {"dup", "after-collision"} + conn.close() + + +def _writer_worker(queue, worker_id, n): + for i in range(n): + queue.put(_record(f"proc-{worker_id}-run-{i}")) + + +def test_multiprocess_writer_receives_all_records_from_multiple_workers(tmp_path): + db_path = tmp_path / "runs.sqlite3" + ctx = multiprocessing.get_context("spawn") + queue = ctx.Queue() + + writer = ctx.Process(target=run_writer_process, args=(queue, db_path)) + writer.start() + + workers = [ctx.Process(target=_writer_worker, args=(queue, worker_id, 20)) for worker_id in range(3)] + for w in workers: + w.start() + for w in workers: + w.join() + + queue.put(None) + writer.join() + + conn = open_connection(db_path) + assert count_runs(conn) == 60 + conn.close()