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 <noreply@anthropic.com>
This commit is contained in:
2
PLAN.md
2
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] **Этап 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] **Этап 3 — Датчики и одна ступень**: `physics/sensors.py` (оба типа как события solve_ivp), `sim/stage.py` (полёт → триггер → разряд → энергобаланс). Тест на сохранение энергии + найден/исправлен баг насыщения (см. выше).
|
||||||
- [x] **Этап 4 — Многоступенчатая цепочка**: `sim/coilgun.py` — сквозная координата, отбраковка нереализуемых конфигураций. Дымовой тест на реальной базе компонентов (`test_real_components_smoke.py`) подтверждает: полный путь реальные JSON → физика → цепочка работает (пример: 22.3 м/с, КПД 3.4% — честный неоптимизированный результат).
|
- [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`.
|
- [ ] **Этап 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`).
|
- [ ] **Этап 7 — Отчётность**: `report/plots.py`, `summary.py`, `bom.py`, раздел "ограничения модели", `cli.py` (`gausse sweep/evolve/simulate/report`).
|
||||||
- [ ] `report/animate.py` — анимация одного прогона: положение снаряда в трубе по времени + визуализация поля/тока каждой катушки (свечение/интенсивность цвета ~ ток), сохранение в GIF/MP4 (matplotlib `FuncAnimation`). Нужна по запросу пользователя — "графика где будет показана симуляция пролёта цилиндра по трубе и электромагнитные поля в каждый момент времени".
|
- [ ] `report/animate.py` — анимация одного прогона: положение снаряда в трубе по времени + визуализация поля/тока каждой катушки (свечение/интенсивность цвета ~ ток), сохранение в GIF/MP4 (matplotlib `FuncAnimation`). Нужна по запросу пользователя — "графика где будет показана симуляция пролёта цилиндра по трубе и электромагнитные поля в каждый момент времени".
|
||||||
|
|||||||
116
src/gausse/storage/database.py
Normal file
116
src/gausse/storage/database.py
Normal file
@@ -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()
|
||||||
50
src/gausse/storage/schema.py
Normal file
50
src/gausse/storage/schema.py
Normal file
@@ -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
|
||||||
126
tests/test_storage_database.py
Normal file
126
tests/test_storage_database.py
Normal file
@@ -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()
|
||||||
Reference in New Issue
Block a user