Files
gausse/tests/test_storage_database.py
jze9 7732577f31 Эволюция: в базу только улучшения лучших + счётчики по всем; выборка 10000
По требованию пользователя «не сохранять всё подряд, а только то, что
нужно для эволюции»: рядовые геномы поколения больше НЕ пишутся строками
в runs — писатель получает лёгкие StatsOnly-инкременты (счётчики,
гистограмма, причины отказов на дашборде остаются честными по ВСЕМ
оценкам), а полная строка пишется только когда лучший геном острова
улучшился (+ полировка). Вставка миллионов строк в HDD-базу с 6 индексами
была узким местом конвейера: GPU считал цикл ~23 мин, дозапись хвоста
очереди шла часами (~150 строк/с), реальная скорость была ~165 оценок/с
вместо ~1300/с. gpu_record_worker стал не нужен — удалён.

Выборка эволюции уменьшена до 10000 на поколение (8 островов x 1250,
было 8x4000=32000) — по просьбе пользователя, чтобы поколения сменялись
быстрее; поколений за цикл теперь 120 (было 60) — при освобождённом
писателе цикл углубляется вдвое дальше за то же время.
2026-07-12 22:29:17 +05:00

157 lines
5.3 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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 test_writer_stats_only_bumps_counters_without_rows(tmp_path):
"""StatsOnly-сообщения (эволюция шлёт их для рядовых геномов): агрегаты
дашборда растут, а строк в runs не появляется."""
from gausse.storage.database import StatsOnly
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("best-one")) # одна полная строка
queue.put(StatsOnly("evolve", True, 0.42, None))
queue.put(StatsOnly("evolve", False, None, "датчик: сигнал ниже порога"))
queue.put(StatsOnly("evolve", False, None, None))
queue.put(None)
writer.join(timeout=10)
assert not writer.is_alive()
conn = open_connection(db_path)
assert count_runs(conn) == 1 # строка только у полной записи
counters = dict(conn.execute("SELECT key, value FROM stats_counters"))
assert counters["total"] == 4 # но посчитаны все четыре
assert counters["feasible"] == 2
reasons = dict(conn.execute("SELECT reason, count FROM stats_reasons"))
assert reasons.get("датчик") == 1
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()