По требованию пользователя «не сохранять всё подряд, а только то, что нужно для эволюции»: рядовые геномы поколения больше НЕ пишутся строками в runs — писатель получает лёгкие StatsOnly-инкременты (счётчики, гистограмма, причины отказов на дашборде остаются честными по ВСЕМ оценкам), а полная строка пишется только когда лучший геном острова улучшился (+ полировка). Вставка миллионов строк в HDD-базу с 6 индексами была узким местом конвейера: GPU считал цикл ~23 мин, дозапись хвоста очереди шла часами (~150 строк/с), реальная скорость была ~165 оценок/с вместо ~1300/с. gpu_record_worker стал не нужен — удалён. Выборка эволюции уменьшена до 10000 на поколение (8 островов x 1250, было 8x4000=32000) — по просьбе пользователя, чтобы поколения сменялись быстрее; поколений за цикл теперь 120 (было 60) — при освобождённом писателе цикл углубляется вдвое дальше за то же время.
157 lines
5.3 KiB
Python
157 lines
5.3 KiB
Python
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()
|