diff --git a/README.md b/README.md index b747f8a..acb2061 100644 --- a/README.md +++ b/README.md @@ -87,7 +87,7 @@ flyguard/ ядро: стадии обработки, память, сч pipeline.py сборка synth.py вставка предметов трассировкой лучей tools/ обучение, оценка, разбор -tests/ 27 тестов, запускаются без данных и без ROS +tests/ 29 тестов, запускаются без данных и без ROS docs/ методика и результаты artifacts/ обученные модели ``` @@ -119,6 +119,17 @@ python tools/make_benchmark.py --memory artifacts/mushroom_body.npz \ python tools/plot_benchmark.py # кривые и график ``` +Тяжёлые шаги сами раскладываются по бэгам на процессы — записей пять, физических +ядер шесть, и это вся доступная зернистость: конвейер держит состояние между +кадрами, поэтому разрезать одну запись нельзя. Замерено: полигон 134 → 36 с, +сбор выборки 96 → 26 с на облегчённой конфигурации, то есть 3.7–3.8×, и файл на +выходе совпадает с последовательным **побайтово**. Отключается `--jobs 1`. + +Для замера задержки кадра `--jobs 1` обязателен: под пятью процессами время +кадра растёт с 32 до 56 мс. Это свойство замера, а не конвейера, поэтому +`evaluate.py` в параллельном режиме печатает задержку как `nan` — чтобы такое +число нельзя было случайно привести в отчёте. + --- ## Где мы сейчас diff --git a/tests/test_pipeline.py b/tests/test_pipeline.py index d6f1080..6690f83 100644 --- a/tests/test_pipeline.py +++ b/tests/test_pipeline.py @@ -404,3 +404,49 @@ def test_real_bag_projects_without_angular_error(): r = img.r_near[img.valid] assert r.min() > 0 and r.max() < 250 assert np.all(img.r_far[img.valid] >= img.r_near[img.valid] - 1e-3) + + +# ------------------------------------------------------------ раскладка по бэгам + +_TOOLS = ROOT / "tools" # выгрузка: tools лежит рядом с тестами +if not _TOOLS.exists(): + _TOOLS = ROOT.parents[2] / "tools" # основной проект: ros2_ws/src/flyguard +if str(_TOOLS) not in sys.path: + sys.path.insert(0, str(_TOOLS)) + +import _parallel as _P # noqa: E402 + + +def _twice(x): + """Задача для проверки. Верхнего уровня: иначе её не передать в процесс.""" + return x * 2 + + +def test_parallel_keeps_task_order_when_results_arrive_out_of_order(): + """Считается по готовности, складывается по номеру задачи. + + Ломается это незаметно и опасно: цифры остаются правдоподобными, просто + приписываются не тому бэгу. Поэтому проверяется не «столько же строк», а + что результат каждой задачи лёг на своё место. + """ + tasks = list(range(7)) + want = [(t, t * 2) for t in tasks] + + seq = [None] * len(tasks) + for i, t, r, _ in _P.run(_twice, tasks, jobs=1): + seq[i] = (t, r) + + par = [None] * len(tasks) + for i, t, r, _ in _P.run(_twice, tasks, jobs=4): + par[i] = (t, r) + + assert seq == want + assert par == want + + +def test_parallel_never_starts_more_processes_than_there_are_bags(): + """Бэгов пять, и шестой процесс занять нечем.""" + assert _P.resolve(0, 1) == 1 # одна задача — без пула вовсе + assert _P.resolve(8, 3) == 3 # просили больше, чем есть работы + assert _P.resolve(1, 5) == 1 # явная последовательная отладка + assert 1 <= _P.resolve(0, 5) <= 5 diff --git a/tools/_parallel.py b/tools/_parallel.py new file mode 100644 index 0000000..5778094 --- /dev/null +++ b/tools/_parallel.py @@ -0,0 +1,91 @@ +"""Раскладка задач по бэгам на процессы. + +Конвейер держит состояние между кадрами — пройденный путь, треки, накопитель +складчатого тела, — поэтому разрезать один бэг нельзя. Зато бэги независимы +друг от друга, и это ровно та зернистость, которая нужна: пять записей на +шесть физических ядер. + +Результат совпадает с последовательным прогоном **точно**, а не «примерно»: +генератор случайных чисел создаётся внутри задачи от того же зерна, общей +изменяемой памяти между бэгами нет. Порядок родитель восстанавливает сам — +печатает по готовности, а сохраняет в исходном. + +Ускорение упирается не в ядра, а в число бэгов: их пять, и больше пяти +процессов тут просто нечем занять. +""" +from __future__ import annotations + +import os +import time +from concurrent.futures import ProcessPoolExecutor, as_completed + + +def limit_threads() -> None: + """Один поток BLAS на процесс. + + Вызывать ДО создания пула: переменные наследуются дочерним процессом при + запуске, а число потоков BLAS выбирает один раз при импорте numpy и потом + не меняет. Без этого каждый из пяти воркеров разворачивается на все ядра + и они дерутся за те же шесть. + """ + for v in ("OMP_NUM_THREADS", "OPENBLAS_NUM_THREADS", "MKL_NUM_THREADS", + "NUMEXPR_NUM_THREADS", "VECLIB_MAXIMUM_THREADS"): + os.environ.setdefault(v, "1") + + +def resolve(jobs: int, n_tasks: int) -> int: + """Сколько процессов поднимать: 0 — по числу задач, но не больше ядер.""" + if n_tasks <= 1: + return 1 + if jobs > 0: + return max(1, min(jobs, n_tasks)) + cpu = os.cpu_count() or 2 + return max(1, min(n_tasks, cpu // 2)) # логических ядер вдвое больше физических + + +class _Timed: + """Часы держит сам воркер. + + В родителе видно только момент готовности, а он у всех задач, стартовавших + разом, почти один и тот же — по такому замеру не понять, какой бэг тяжёлый. + """ + + def __init__(self, fn): + self.fn = fn + + def __call__(self, task): + t0 = time.time() + return self.fn(task), time.time() - t0 + + +def run(work, tasks, jobs: int): + """Выполнить `work(task)` по всем задачам, отдавая `(i, task, res, с)`. + + Отдаёт по мере готовности, поэтому `i` — исходный номер задачи, и по нему + вызывающий раскладывает результаты обратно в порядок бэгов. + + При одном процессе всё считается прямо здесь, без пула: остаётся чем + отлаживать, и трассировка ошибки не проходит через межпроцессную передачу. + """ + tasks = list(tasks) + n = resolve(jobs, len(tasks)) + if n <= 1: + for i, t in enumerate(tasks): + t0 = time.time() + yield i, t, work(t), time.time() - t0 + return + + limit_threads() + with ProcessPoolExecutor(max_workers=n) as ex: + fut = {ex.submit(_Timed(work), t): (i, t) for i, t in enumerate(tasks)} + for f in as_completed(fut): + i, t = fut[f] + res, secs = f.result() + yield i, t, res, secs + + +def add_argument(ap) -> None: + """Один и тот же флаг во всех инструментах, чтобы не помнить разные.""" + ap.add_argument("--jobs", type=int, default=0, + help="сколько бэгов считать разом; 0 — по числу бэгов, " + "но не больше физических ядер; 1 — в один процесс") diff --git a/tools/evaluate.py b/tools/evaluate.py index 38b54bc..81f2b0f 100644 --- a/tools/evaluate.py +++ b/tools/evaluate.py @@ -1,246 +1,281 @@ -"""Оценка обобщаемости: leave-one-bag-out. - -Память тоннеля обучается на всех данных, **кроме** проверяемого бэга, и только -после этого конвейер прогоняется по нему. Иначе цифры лгут: подавлять -конструкции, которые сам же и запомнил, умеет кто угодно, а на приватном тесте -будет новый участок тоннеля. - -Отчёт: ложные тревоги на километр пути и доля кадров с тревогой для пустых -бэгов; для `doubleT_obstacle` — ещё и доля кадров, в которых найден настоящий -объект на ~55 м. - - python tools/evaluate.py --device cuda --out artifacts/generalisation.json -""" -from __future__ import annotations - -import argparse -import json -import time - -import numpy as np - -import _bootstrap as B # noqa: F401 -from flyguard.bag import Bag, find_bags -from flyguard.mushroom_body import MushroomBody, MushroomBodyConfig -from flyguard.pipeline import FlyGuard, Params - -OBSTACLE_BAG = "doubleT_obstacle" -TRUE_D = (50.0, 62.0) - - -def train_excluding(per_bag: dict[str, np.ndarray], extra: np.ndarray | None, - exclude: str, target: float, device: str) -> MushroomBody: - parts = [v for k, v in per_bag.items() if k not in (exclude, OBSTACLE_BAG)] - if extra is not None: - parts.append(extra) - X = np.concatenate(parts).astype(np.float32) - mb = MushroomBody(MushroomBodyConfig()) - mb.fit_normalizer(X) - mb.learn(X, rate=mb.auto_rate(X.shape[0], target), device=device) - return mb - - -def run_bag(bag_path, memory, limit: int, params: Params | None = None, - readout=None) -> dict: - fg = FlyGuard(params or Params(), memory=memory, readout=readout) - bag = Bag(bag_path) - n = alarms = obj_hits = fp_objects = 0 - path_m = 0.0 - fp_dists, times = [], [] - # один и тот же лоток, попавший в треки, виден сотню кадров подряд; для - # эксплуатации важно не это, а сколько РАЗНЫХ ложных объектов возникло — - # именно столько раз поезд затормозил бы напрасно - fp_tracks: set[int] = set() - for _, pc in bag.frames(stop=limit): - res = fg.process(pc) - if res is None: - continue - n += 1 - path_m += res.ego.ds if res.ego else 0.0 - times.append(res.total_ms) - d = res.decision - is_obstacle_bag = bag.path.name == OBSTACLE_BAG - mine = [o for o in d.objects if TRUE_D[0] < o.distance < TRUE_D[1]] \ - if is_obstacle_bag else [] - others = [o for o in d.objects if o not in mine] - if mine: - obj_hits += 1 - if others: - alarms += 1 - fp_objects += len(others) - fp_tracks.update(o.track_id for o in others) - fp_dists.extend(o.distance for o in others) - km = max(path_m / 1000.0, 1e-6) - return dict(bag=bag.path.name, frames=n, path_m=path_m, - alarm_frames=alarms, alarm_rate=alarms / max(n, 1), - fp_objects=fp_objects, fp_tracks=len(fp_tracks), - fp_per_km=(len(fp_tracks) / km) if path_m > 5 else float("nan"), - fp_median_d=float(np.median(fp_dists)) if fp_dists else float("nan"), - obj_rate=obj_hits / max(n, 1) if bag.path.name == OBSTACLE_BAG else None, - ms_p50=float(np.median(times)) if times else 0.0, - ms_p95=float(np.percentile(times, 95)) if times else 0.0) - - -def main() -> None: - ap = argparse.ArgumentParser(description=__doc__) - ap.add_argument("--root", default=str(B.DATA / "for_hackathon")) - ap.add_argument("--cache", default=str(B.CACHE / "tune_candidates.npz")) - ap.add_argument("--extra-cache", default=str(B.CACHE / "new_data_candidates.npz")) - ap.add_argument("--limit", type=int, default=250) - ap.add_argument("--target", type=float, default=0.4) - ap.add_argument("--device", default="cpu") - ap.add_argument("--split-adv", type=float, default=None, - help="порог разделения фигуры и фона; 0 — выключить") - ap.add_argument("--split-gap", type=float, default=None, - help="порог разреза по контрасту ламины, м; 0 — выключить") - ap.add_argument("--split-near", type=float, default=None, - help="ближе этой дальности не резать, м") - ap.add_argument("--split-top", type=int, default=None, - help="сколько фигур выносить из одной компоненты; 0 — все") - ap.add_argument("--no-acc", action="store_true", help="выключить накопитель") - ap.add_argument("--no-hab", action="store_true", help="выключить привыкание") - ap.add_argument("--mbon", default="", help="путь к обученному считыванию MBON") - ap.add_argument("--mbon-dir", default="", - help="каталог с моделями по складкам (mbon_<бэг>.npz): " - "для каждого бэга берётся модель, его не видевшая") - ap.add_argument("--mbon-blend", type=float, default=None, - help="1 — только модель, 0 — только ручная формула") - ap.add_argument("--mbon-power", type=float, default=None, - help="резкость вероятности модели") - ap.add_argument("--warn", type=float, default=None, - help="порог улики для тревоги (по умолчанию 0.5)") - ap.add_argument("--clear", type=float, default=None, - help="нижний порог гистерезиса (по умолчанию 0.3)") - ap.add_argument("--min-hits", type=int, default=None, - help="наблюдений, без которых трек не считается") - ap.add_argument("--warn-far", type=float, default=None, - help="порог тревоги на дальнем краю; ниже обычного — послабление далёким трекам") - ap.add_argument("--warn-far-from", type=float, default=None, - help="с какой дальности порог начинает падать, м") - ap.add_argument("--novelty-floor", type=float, default=None, - help="ниже этой новизны трек не считается") - ap.add_argument("--near-long", type=float, default=None, - help="разброс дальности, выше которого компоненту режут и вблизи; 0 — не резать вблизи") - ap.add_argument("--min-rays", type=int, default=None, - help="сколько лучей минимум образуют кандидата") - ap.add_argument("--min-rays-far", type=int, default=None, - help="порог по лучам за `--min-rays-far-from`; 0 — не различать") - ap.add_argument("--min-rays-far-from", type=float, default=None) - ap.add_argument("--leak-far", type=float, default=None, - help="утечка улики на дальнем краю; ниже обычной — далёкий трек прощает промахи") - ap.add_argument("--leak-far-from", type=float, default=None, - help="с какой дальности утечка начинает падать, м") - ap.add_argument("--nov-fade-from", type=float, default=None, - help="с какой дальности гасить вклад знакомости; 0 — не гасить") - ap.add_argument("--nov-fade-to", type=float, default=None, - help="к какой дальности вклад знакомости обнуляется") - ap.add_argument("--mbon-prior-from", type=float, default=None, - help="с какой дальности поправлять оценку модели на распространённость предметов; 0 — не поправлять") - ap.add_argument("--no-memory", action="store_true", - help="совсем без памяти тоннеля — так выглядит первый проезд по новой линии") - ap.add_argument("--out", default=str(B.ARTIFACTS / "generalisation.json")) - args = ap.parse_args() - - d = np.load(args.cache, allow_pickle=True) - per_bag = {str(k): d[f"X_{k}"].astype(np.float32) for k in d["names"]} - extra = None - from pathlib import Path - if Path(args.extra_cache).exists(): - extra = np.load(args.extra_cache)["X"].astype(np.float32) - print(f"дополнительно в обучение: {extra.shape[0]} кандидатов из new_data") - - over = {} - if args.split_adv is not None: - over["split_adv"] = args.split_adv - if args.split_gap is not None: - over["split_gap"] = args.split_gap - if args.split_near is not None: - over["split_near"] = args.split_near - if args.split_top is not None: - over["split_top"] = args.split_top - if args.no_acc: - over["enable_accumulator"] = False - if args.no_hab: - over["enable_habituation"] = False - if args.mbon_blend is not None: - over["mbon_blend"] = args.mbon_blend - if args.mbon_power is not None: - over["mbon_power"] = args.mbon_power - if args.warn is not None: - over["warn_evidence"] = args.warn - over["clear_evidence"] = args.warn * 0.6 if args.clear is None else args.clear - elif args.clear is not None: - over["clear_evidence"] = args.clear - if args.min_hits is not None: - over["min_hits"] = args.min_hits - if args.warn_far is not None: - over["warn_far"] = args.warn_far - if args.warn_far_from is not None: - over["warn_far_from"] = args.warn_far_from - if args.novelty_floor is not None: - over["novelty_floor"] = args.novelty_floor - if args.near_long is not None: - over["near_long"] = args.near_long - if args.min_rays is not None: - over["min_rays"] = args.min_rays - if args.min_rays_far is not None: - over["min_rays_far"] = args.min_rays_far - if args.min_rays_far_from is not None: - over["min_rays_far_from"] = args.min_rays_far_from - if args.leak_far is not None: - over["leak_far"] = args.leak_far - if args.leak_far_from is not None: - over["leak_far_from"] = args.leak_far_from - if args.nov_fade_from is not None: - over["nov_fade_from"] = args.nov_fade_from - if args.nov_fade_to is not None: - over["nov_fade_to"] = args.nov_fade_to - if args.mbon_prior_from is not None: - over["mbon_prior_from"] = args.mbon_prior_from - params = Params(**over) - readout = None - folds = {} - if args.mbon_dir: - from pathlib import Path as _P - from flyguard.mbon_readout import MbonReadout - for f in _P(args.mbon_dir).glob("mbon_*.npz"): - folds[f.stem[len("mbon_"):]] = MbonReadout.load(f) - print(f"считывание MBON по складкам: {args.mbon_dir} " - f"({len(folds)} моделей)") - elif args.mbon: - from flyguard.mbon_readout import MbonReadout - readout = MbonReadout.load(args.mbon) - print(f"считывание MBON: {args.mbon}") - rows = [] - for p in find_bags(args.root): - t0 = time.time() - mem = (None if args.no_memory else - train_excluding(per_bag, extra, p.name, args.target, args.device)) - rd = folds.get(p.name, readout) if folds else readout - if folds and p.name not in folds and p.name != OBSTACLE_BAG: - print(f" внимание: для {p.name} нет своей складки") - r = run_bag(p, mem, args.limit, params, readout=rd) - r["train_size"] = int(mem.n_seen) if mem is not None else 0 - rows.append(r) - obj = f"объект {r['obj_rate']:6.1%} | " if r["obj_rate"] is not None else "" - print(f"{r['bag']:40s} кадров {r['frames']:4d} путь {r['path_m']:6.0f} м | " - f"{obj}тревог {r['alarm_rate']:6.1%} | ложных треков {r['fp_tracks']:3d} " - f"({r['fp_per_km']:6.1f} на км) | {r['ms_p50']:5.1f}/{r['ms_p95']:5.1f} мс | " - f"{time.time()-t0:5.0f} с", flush=True) - - B.ARTIFACTS.mkdir(parents=True, exist_ok=True) - with open(args.out, "w", encoding="utf-8") as f: - json.dump(rows, f, ensure_ascii=False, indent=1) - empty = [r for r in rows if r["bag"] != OBSTACLE_BAG] - print("\nсводка по пустым бэгам:") - print(f" доля кадров с ложной тревогой: {np.mean([r['alarm_rate'] for r in empty]):.2%}") - fp_km = [r["fp_per_km"] for r in empty if np.isfinite(r["fp_per_km"])] - if fp_km: - print(f" разных ложных треков на километр: {np.mean(fp_km):.1f} " - f"(медиана {np.median(fp_km):.1f})") - print(f" суммарный путь: {sum(r['path_m'] for r in empty):.0f} м") - print("сохранено:", args.out) - - -if __name__ == "__main__": - main() +"""Оценка обобщаемости: leave-one-bag-out. + +Память тоннеля обучается на всех данных, **кроме** проверяемого бэга, и только +после этого конвейер прогоняется по нему. Иначе цифры лгут: подавлять +конструкции, которые сам же и запомнил, умеет кто угодно, а на приватном тесте +будет новый участок тоннеля. + +Отчёт: ложные тревоги на километр пути и доля кадров с тревогой для пустых +бэгов; для `doubleT_obstacle` — ещё и доля кадров, в которых найден настоящий +объект на ~55 м. + + python tools/evaluate.py --device cuda --out artifacts/generalisation.json +""" +from __future__ import annotations + +import argparse +import json +import tempfile + +import numpy as np + +import _bootstrap as B # noqa: F401 +import _parallel as P +from flyguard.bag import Bag, find_bags +from flyguard.mushroom_body import MushroomBody, MushroomBodyConfig +from flyguard.pipeline import FlyGuard, Params + +OBSTACLE_BAG = "doubleT_obstacle" +TRUE_D = (50.0, 62.0) + + +def _work(task): + """Одна задача — один бэг. Память и считывание грузятся по пути уже здесь.""" + path, params, limit, mem_path, rd_path = task + memory = MushroomBody.load(mem_path) if mem_path else None + readout = None + if rd_path: + from flyguard.mbon_readout import MbonReadout + readout = MbonReadout.load(rd_path) + return run_bag(path, memory, limit, params, readout=readout) + + +def train_excluding(per_bag: dict[str, np.ndarray], extra: np.ndarray | None, + exclude: str, target: float, device: str) -> MushroomBody: + parts = [v for k, v in per_bag.items() if k not in (exclude, OBSTACLE_BAG)] + if extra is not None: + parts.append(extra) + X = np.concatenate(parts).astype(np.float32) + mb = MushroomBody(MushroomBodyConfig()) + mb.fit_normalizer(X) + mb.learn(X, rate=mb.auto_rate(X.shape[0], target), device=device) + return mb + + +def run_bag(bag_path, memory, limit: int, params: Params | None = None, + readout=None) -> dict: + fg = FlyGuard(params or Params(), memory=memory, readout=readout) + bag = Bag(bag_path) + n = alarms = obj_hits = fp_objects = 0 + path_m = 0.0 + fp_dists, times = [], [] + # один и тот же лоток, попавший в треки, виден сотню кадров подряд; для + # эксплуатации важно не это, а сколько РАЗНЫХ ложных объектов возникло — + # именно столько раз поезд затормозил бы напрасно + fp_tracks: set[int] = set() + for _, pc in bag.frames(stop=limit): + res = fg.process(pc) + if res is None: + continue + n += 1 + path_m += res.ego.ds if res.ego else 0.0 + times.append(res.total_ms) + d = res.decision + is_obstacle_bag = bag.path.name == OBSTACLE_BAG + mine = [o for o in d.objects if TRUE_D[0] < o.distance < TRUE_D[1]] \ + if is_obstacle_bag else [] + others = [o for o in d.objects if o not in mine] + if mine: + obj_hits += 1 + if others: + alarms += 1 + fp_objects += len(others) + fp_tracks.update(o.track_id for o in others) + fp_dists.extend(o.distance for o in others) + km = max(path_m / 1000.0, 1e-6) + return dict(bag=bag.path.name, frames=n, path_m=path_m, + alarm_frames=alarms, alarm_rate=alarms / max(n, 1), + fp_objects=fp_objects, fp_tracks=len(fp_tracks), + fp_per_km=(len(fp_tracks) / km) if path_m > 5 else float("nan"), + fp_median_d=float(np.median(fp_dists)) if fp_dists else float("nan"), + obj_rate=obj_hits / max(n, 1) if bag.path.name == OBSTACLE_BAG else None, + ms_p50=float(np.median(times)) if times else 0.0, + ms_p95=float(np.percentile(times, 95)) if times else 0.0) + + +def main() -> None: + ap = argparse.ArgumentParser(description=__doc__) + ap.add_argument("--root", default=str(B.DATA / "for_hackathon")) + ap.add_argument("--cache", default=str(B.CACHE / "tune_candidates.npz")) + ap.add_argument("--extra-cache", default=str(B.CACHE / "new_data_candidates.npz")) + ap.add_argument("--limit", type=int, default=250) + ap.add_argument("--target", type=float, default=0.4) + ap.add_argument("--device", default="cpu") + ap.add_argument("--split-adv", type=float, default=None, + help="порог разделения фигуры и фона; 0 — выключить") + ap.add_argument("--split-gap", type=float, default=None, + help="порог разреза по контрасту ламины, м; 0 — выключить") + ap.add_argument("--split-near", type=float, default=None, + help="ближе этой дальности не резать, м") + ap.add_argument("--split-top", type=int, default=None, + help="сколько фигур выносить из одной компоненты; 0 — все") + ap.add_argument("--no-acc", action="store_true", help="выключить накопитель") + ap.add_argument("--no-hab", action="store_true", help="выключить привыкание") + ap.add_argument("--mbon", default="", help="путь к обученному считыванию MBON") + ap.add_argument("--mbon-dir", default="", + help="каталог с моделями по складкам (mbon_<бэг>.npz): " + "для каждого бэга берётся модель, его не видевшая") + ap.add_argument("--mbon-blend", type=float, default=None, + help="1 — только модель, 0 — только ручная формула") + ap.add_argument("--mbon-power", type=float, default=None, + help="резкость вероятности модели") + ap.add_argument("--warn", type=float, default=None, + help="порог улики для тревоги (по умолчанию 0.5)") + ap.add_argument("--clear", type=float, default=None, + help="нижний порог гистерезиса (по умолчанию 0.3)") + ap.add_argument("--min-hits", type=int, default=None, + help="наблюдений, без которых трек не считается") + ap.add_argument("--warn-far", type=float, default=None, + help="порог тревоги на дальнем краю; ниже обычного — послабление далёким трекам") + ap.add_argument("--warn-far-from", type=float, default=None, + help="с какой дальности порог начинает падать, м") + ap.add_argument("--novelty-floor", type=float, default=None, + help="ниже этой новизны трек не считается") + ap.add_argument("--near-long", type=float, default=None, + help="разброс дальности, выше которого компоненту режут и вблизи; 0 — не резать вблизи") + ap.add_argument("--min-rays", type=int, default=None, + help="сколько лучей минимум образуют кандидата") + ap.add_argument("--min-rays-far", type=int, default=None, + help="порог по лучам за `--min-rays-far-from`; 0 — не различать") + ap.add_argument("--min-rays-far-from", type=float, default=None) + ap.add_argument("--leak-far", type=float, default=None, + help="утечка улики на дальнем краю; ниже обычной — далёкий трек прощает промахи") + ap.add_argument("--leak-far-from", type=float, default=None, + help="с какой дальности утечка начинает падать, м") + ap.add_argument("--nov-fade-from", type=float, default=None, + help="с какой дальности гасить вклад знакомости; 0 — не гасить") + ap.add_argument("--nov-fade-to", type=float, default=None, + help="к какой дальности вклад знакомости обнуляется") + ap.add_argument("--mbon-prior-from", type=float, default=None, + help="с какой дальности поправлять оценку модели на распространённость предметов; 0 — не поправлять") + ap.add_argument("--no-memory", action="store_true", + help="совсем без памяти тоннеля — так выглядит первый проезд по новой линии") + ap.add_argument("--out", default=str(B.ARTIFACTS / "generalisation.json")) + P.add_argument(ap) + args = ap.parse_args() + + d = np.load(args.cache, allow_pickle=True) + per_bag = {str(k): d[f"X_{k}"].astype(np.float32) for k in d["names"]} + extra = None + from pathlib import Path + if Path(args.extra_cache).exists(): + extra = np.load(args.extra_cache)["X"].astype(np.float32) + print(f"дополнительно в обучение: {extra.shape[0]} кандидатов из new_data") + + over = {} + if args.split_adv is not None: + over["split_adv"] = args.split_adv + if args.split_gap is not None: + over["split_gap"] = args.split_gap + if args.split_near is not None: + over["split_near"] = args.split_near + if args.split_top is not None: + over["split_top"] = args.split_top + if args.no_acc: + over["enable_accumulator"] = False + if args.no_hab: + over["enable_habituation"] = False + if args.mbon_blend is not None: + over["mbon_blend"] = args.mbon_blend + if args.mbon_power is not None: + over["mbon_power"] = args.mbon_power + if args.warn is not None: + over["warn_evidence"] = args.warn + over["clear_evidence"] = args.warn * 0.6 if args.clear is None else args.clear + elif args.clear is not None: + over["clear_evidence"] = args.clear + if args.min_hits is not None: + over["min_hits"] = args.min_hits + if args.warn_far is not None: + over["warn_far"] = args.warn_far + if args.warn_far_from is not None: + over["warn_far_from"] = args.warn_far_from + if args.novelty_floor is not None: + over["novelty_floor"] = args.novelty_floor + if args.near_long is not None: + over["near_long"] = args.near_long + if args.min_rays is not None: + over["min_rays"] = args.min_rays + if args.min_rays_far is not None: + over["min_rays_far"] = args.min_rays_far + if args.min_rays_far_from is not None: + over["min_rays_far_from"] = args.min_rays_far_from + if args.leak_far is not None: + over["leak_far"] = args.leak_far + if args.leak_far_from is not None: + over["leak_far_from"] = args.leak_far_from + if args.nov_fade_from is not None: + over["nov_fade_from"] = args.nov_fade_from + if args.nov_fade_to is not None: + over["nov_fade_to"] = args.nov_fade_to + if args.mbon_prior_from is not None: + over["mbon_prior_from"] = args.mbon_prior_from + params = Params(**over) + fold_paths: dict[str, str] = {} + if args.mbon_dir: + from pathlib import Path as _P + for f in _P(args.mbon_dir).glob("mbon_*.npz"): + fold_paths[f.stem[len("mbon_"):]] = str(f) + print(f"считывание MBON по складкам: {args.mbon_dir} " + f"({len(fold_paths)} моделей)") + elif args.mbon: + print(f"считывание MBON: {args.mbon}") + + # Память обучается ЗДЕСЬ, а не в воркере: обучение идёт на видеокарте, и + # делить одну карту на пять процессов незачем. Воркеру достаётся готовый + # файл — долгая часть это проход по записи, а не обучение. + tmp = None if args.no_memory else tempfile.TemporaryDirectory(prefix="fg_mem_") + tasks, sizes = [], [] + for p in find_bags(args.root): + mem_path, n_seen = "", 0 + if tmp is not None: + mem = train_excluding(per_bag, extra, p.name, args.target, args.device) + mem_path = str(Path(tmp.name) / f"mem_{p.name}.npz") + mem.save(mem_path) + n_seen = int(mem.n_seen) + del mem + rd = fold_paths.get(p.name, args.mbon) if fold_paths else args.mbon + if fold_paths and p.name not in fold_paths and p.name != OBSTACLE_BAG: + print(f" внимание: для {p.name} нет своей складки") + tasks.append((p, params, args.limit, mem_path, rd)) + sizes.append(n_seen) + + # Задержка кадра под пятью процессами вырастает с 32 до 56 мс — это + # свойство замера, а не конвейера, и такое число нельзя показывать как + # запас по бюджету. Поэтому в параллельном режиме оно не печатается + # вовсе: перепутать nan с честным замером невозможно, а предупреждение + # в шапке пролистывается. + par = P.resolve(args.jobs, len(tasks)) > 1 + if par: + print("параллельно: задержка кадра не измеряется, для неё нужен --jobs 1") + slots: list = [None] * len(tasks) + for i, task, r, secs in P.run(_work, tasks, args.jobs): + r["train_size"] = sizes[i] + if par: + r["ms_p50"] = r["ms_p95"] = float("nan") + slots[i] = r + obj = f"объект {r['obj_rate']:6.1%} | " if r["obj_rate"] is not None else "" + print(f"{r['bag']:40s} кадров {r['frames']:4d} путь {r['path_m']:6.0f} м | " + f"{obj}тревог {r['alarm_rate']:6.1%} | ложных треков {r['fp_tracks']:3d} " + f"({r['fp_per_km']:6.1f} на км) | {r['ms_p50']:5.1f}/{r['ms_p95']:5.1f} мс | " + f"{secs:5.0f} с", flush=True) + rows = [r for r in slots if r is not None] + if tmp is not None: + tmp.cleanup() + + B.ARTIFACTS.mkdir(parents=True, exist_ok=True) + with open(args.out, "w", encoding="utf-8") as f: + json.dump(rows, f, ensure_ascii=False, indent=1) + empty = [r for r in rows if r["bag"] != OBSTACLE_BAG] + print("\nсводка по пустым бэгам:") + print(f" доля кадров с ложной тревогой: {np.mean([r['alarm_rate'] for r in empty]):.2%}") + fp_km = [r["fp_per_km"] for r in empty if np.isfinite(r["fp_per_km"])] + if fp_km: + print(f" разных ложных треков на километр: {np.mean(fp_km):.1f} " + f"(медиана {np.median(fp_km):.1f})") + print(f" суммарный путь: {sum(r['path_m'] for r in empty):.0f} м") + print("сохранено:", args.out) + + +if __name__ == "__main__": + main() diff --git a/tools/make_benchmark.py b/tools/make_benchmark.py index dac01ea..280c5d7 100644 --- a/tools/make_benchmark.py +++ b/tools/make_benchmark.py @@ -1,236 +1,253 @@ -"""Размеченный полигон: сценарии сближения с синтетическим препятствием. - -Для каждого бэга и каждого типа предмета строится сценарий: предмет ставится -в фиксированную точку тоннеля далеко впереди, поезд к нему подъезжает, и на -каждом кадре известна истинная дистанция. Отсюда получаются именно те цифры, -которые просит ТЗ: с какой дальности предмет уверенно виден, сколько ложных -тревог и как это зависит от размера. - -Первый проход считает собственное движение по чистым данным (это и есть -разметка по дистанции), второй — гоняет конвейер по кадрам со вставленным -предметом. Все сценарии одного бэга обрабатываются в одном проходе по файлу: -чтение данных дороже самой обработки. - - python tools/make_benchmark.py --out artifacts/benchmark.npz --memory artifacts/mushroom_body.npz -""" -from __future__ import annotations - -import argparse -import json -import time - -import numpy as np - -import _bootstrap as B # noqa: F401 -from flyguard.bag import Bag, find_bags -from flyguard.mushroom_body import MushroomBody -from flyguard.pipeline import FlyGuard, Params -from flyguard.synth import IntensityEnv, Placement, catalogue, inject - -HOLDOUT = "doubleT_obstacle" # там уже есть настоящий объект - - -def ego_track(bag: Bag, params: Params, limit: int | None): - """Первый проход: пройденный путь на каждом кадре (разметка по дистанции).""" - fg = FlyGuard(params, memory=None) - s, stamps = [], [] - total = 0.0 - for _, pc in bag.frames(stop=limit): - res = fg.process(pc) - if res is None: - s.append(None); stamps.append(pc.stamp); continue - total += res.ego.ds if res.ego else 0.0 - s.append(total); stamps.append(pc.stamp) - return s, stamps - - -def run_bag(bag_path, params: Params, memory, limit: int, d_start: float, - laterals: tuple[float, ...], seed: int, readout=None) -> list[dict]: - bag = Bag(bag_path) - s_track, _ = ego_track(bag, params, limit) - have = [x for x in s_track if x is not None] - if len(have) < 20: - return [] - travel = have[-1] - have[0] - - cat = catalogue() - scen = [(name, lat) for name in cat for lat in laterals] - pipes = [FlyGuard(params, memory=memory, readout=readout) for _ in scen] - rng = np.random.default_rng(seed) - records = [[] for _ in scen] - - # решётка и поза нужны для вставки — берутся из отдельного «чистого» конвейера - guide = FlyGuard(params, memory=None) - - for k, (_, pc) in enumerate(bag.frames(stop=limit)): - gres = guide.process(pc) - if gres is None or s_track[k] is None: - continue - s_now = s_track[k] - have[0] - env = IntensityEnv(pc) # один раз на кадр, общий для сценариев - for i, (name, lat) in enumerate(scen): - d_true = d_start - s_now - if d_true < 6.0: - continue - # Предмет лежит НА ПУТИ, а путь в кривой уходит вбок: на 150 м при - # радиусе 1300 м это 8.6 м. Если ставить его в поперечных координатах - # сенсора, он окажется в стене, а не в габарите. - u_obj = float(gres.corridor.centre(np.array([d_true], np.float32))[0]) + lat - pc2, lab = inject(pc, guide.layout_full, gres.plane, cat[name], - Placement(d=d_true, u=u_obj), rng=rng, env=env) - res = pipes[i].process(pc2) - if res is None: - continue - tol = max(3.0, 0.12 * d_true) - hit = any(abs(o.distance - d_true) < tol for o in res.decision.objects) - fp = sum(1 for o in res.decision.objects if abs(o.distance - d_true) >= tol) - # Воронка потерь. Лучи в предмет попали — а дальше он может - # пропасть на любой из трёх ступеней, и лечатся они по-разному: - # нет кандидата — вопрос к кластеризации и разделению фигуры и - # фона, нет трека — к сопоставлению по кадрам, нет решения — - # к порогу. Без этого разбиения улучшать нечего, кроме удачи. - cand = any(abs(c.d - d_true) < tol for c in res.candidates) - cx = pipes[i].cx - trk = any(abs(tr.distance(cx.s_world) - d_true) < tol - for tr in cx.tracks) - records[i].append((d_true, int(hit), fp, lab["hit_rays"], - int(cand), int(trk))) - - out = [] - for (name, lat), rec in zip(scen, records): - if not rec: - continue - a = np.array(rec, np.float32) - out.append(dict(bag=bag.path.name, obj=name, lateral=lat, - d=a[:, 0].tolist(), hit=a[:, 1].tolist(), - fp=a[:, 2].tolist(), rays=a[:, 3].tolist(), - cand=a[:, 4].tolist(), trk=a[:, 5].tolist(), - travel=float(travel))) - return out - - -def main() -> None: - ap = argparse.ArgumentParser(description=__doc__) - ap.add_argument("--root", default=str(B.DATA / "for_hackathon")) - ap.add_argument("--memory") - ap.add_argument("--mbon", default="", - help="обученное считывание MBON одной моделью; она видела " - "эти бэги — годится только для отладки") - ap.add_argument("--mbon-dir", default="", - help="каталог с моделями по складкам (mbon_<бэг>.npz): для " - "каждого бэга берётся модель, его НЕ видевшая. Иначе " - "дальность завышена: считывание обучалось ровно на " - "таких же вставках в этот же тоннель") - ap.add_argument("--out", default=str(B.ARTIFACTS / "benchmark.json")) - ap.add_argument("--limit", type=int, default=250) - ap.add_argument("--d-start", type=float, default=200.0) - ap.add_argument("--laterals", default="0.0,0.9") - ap.add_argument("--seed", type=int, default=12345) - ap.add_argument("--mbon-blend", type=float, default=None, - help="1 — только модель, 0 — только ручная формула") - ap.add_argument("--warn", type=float, default=None, - help="порог улики для тревоги (по умолчанию 0.5)") - ap.add_argument("--clear", type=float, default=None) - ap.add_argument("--min-hits", type=int, default=None) - ap.add_argument("--warn-far", type=float, default=None, - help="порог тревоги на дальнем краю; ниже обычного — послабление далёким трекам") - ap.add_argument("--warn-far-from", type=float, default=None, - help="с какой дальности порог начинает падать, м") - ap.add_argument("--novelty-floor", type=float, default=None, - help="ниже этой новизны трек не считается") - ap.add_argument("--near-long", type=float, default=None, - help="разброс дальности, выше которого компоненту режут и вблизи; 0 — не резать вблизи") - ap.add_argument("--min-rays", type=int, default=None, - help="сколько лучей минимум образуют кандидата") - ap.add_argument("--min-rays-far", type=int, default=None, - help="порог по лучам за `--min-rays-far-from`; 0 — не различать") - ap.add_argument("--min-rays-far-from", type=float, default=None) - ap.add_argument("--leak-far", type=float, default=None, - help="утечка улики на дальнем краю; ниже обычной — далёкий трек прощает промахи") - ap.add_argument("--leak-far-from", type=float, default=None, - help="с какой дальности утечка начинает падать, м") - ap.add_argument("--nov-fade-from", type=float, default=None, - help="с какой дальности гасить вклад знакомости; 0 — не гасить") - ap.add_argument("--nov-fade-to", type=float, default=None, - help="к какой дальности вклад знакомости обнуляется") - ap.add_argument("--mbon-prior-from", type=float, default=None, - help="с какой дальности поправлять оценку модели на распространённость предметов; 0 — не поправлять") - args = ap.parse_args() - - B.ARTIFACTS.mkdir(parents=True, exist_ok=True) - memory = MushroomBody.load(args.memory) if args.memory else None - readout = None - folds: dict = {} - if args.mbon_dir: - from pathlib import Path as _P - from flyguard.mbon_readout import MbonReadout - for f in _P(args.mbon_dir).glob("mbon_*.npz"): - folds[f.stem[len("mbon_"):]] = MbonReadout.load(f) - print(f"считывание MBON по складкам: {args.mbon_dir} " - f"({len(folds)} моделей)") - elif args.mbon: - from flyguard.mbon_readout import MbonReadout - readout = MbonReadout.load(args.mbon) - print(f"считывание MBON: {args.mbon}") - over = {} - if args.mbon_blend is not None: - over["mbon_blend"] = args.mbon_blend - if args.warn is not None: - over["warn_evidence"] = args.warn - over["clear_evidence"] = args.warn * 0.6 if args.clear is None else args.clear - elif args.clear is not None: - over["clear_evidence"] = args.clear - if args.min_hits is not None: - over["min_hits"] = args.min_hits - if args.warn_far is not None: - over["warn_far"] = args.warn_far - if args.warn_far_from is not None: - over["warn_far_from"] = args.warn_far_from - if args.novelty_floor is not None: - over["novelty_floor"] = args.novelty_floor - if args.near_long is not None: - over["near_long"] = args.near_long - if args.min_rays is not None: - over["min_rays"] = args.min_rays - if args.min_rays_far is not None: - over["min_rays_far"] = args.min_rays_far - if args.min_rays_far_from is not None: - over["min_rays_far_from"] = args.min_rays_far_from - if args.leak_far is not None: - over["leak_far"] = args.leak_far - if args.leak_far_from is not None: - over["leak_far_from"] = args.leak_far_from - if args.nov_fade_from is not None: - over["nov_fade_from"] = args.nov_fade_from - if args.nov_fade_to is not None: - over["nov_fade_to"] = args.nov_fade_to - if args.mbon_prior_from is not None: - over["mbon_prior_from"] = args.mbon_prior_from - params = Params(**over) - laterals = tuple(float(x) for x in args.laterals.split(",")) - - all_rec = [] - for p in find_bags(args.root): - if p.name == HOLDOUT: - continue - t0 = time.time() - rd = readout - if folds: - rd = folds.get(p.name) - if rd is None: - print(f" внимание: для {p.name} нет своей складки — пропуск") - continue - rec = run_bag(p, params, memory, args.limit, args.d_start, laterals, - args.seed, readout=rd) - all_rec.extend(rec) - n = sum(len(r["d"]) for r in rec) - print(f" {p.name:42s} сценариев {len(rec):3d}, наблюдений {n:6d}, " - f"{time.time()-t0:6.1f} с", flush=True) - - with open(args.out, "w", encoding="utf-8") as f: - json.dump(all_rec, f, ensure_ascii=False) - print(f"сохранено: {args.out} ({len(all_rec)} сценариев)") - - -if __name__ == "__main__": - main() +"""Размеченный полигон: сценарии сближения с синтетическим препятствием. + +Для каждого бэга и каждого типа предмета строится сценарий: предмет ставится +в фиксированную точку тоннеля далеко впереди, поезд к нему подъезжает, и на +каждом кадре известна истинная дистанция. Отсюда получаются именно те цифры, +которые просит ТЗ: с какой дальности предмет уверенно виден, сколько ложных +тревог и как это зависит от размера. + +Первый проход считает собственное движение по чистым данным (это и есть +разметка по дистанции), второй — гоняет конвейер по кадрам со вставленным +предметом. Все сценарии одного бэга обрабатываются в одном проходе по файлу: +чтение данных дороже самой обработки. + + python tools/make_benchmark.py --out artifacts/benchmark.npz --memory artifacts/mushroom_body.npz +""" +from __future__ import annotations + +import argparse +import json + +import numpy as np + +import _bootstrap as B # noqa: F401 +import _parallel as P +from flyguard.bag import Bag, find_bags +from flyguard.mushroom_body import MushroomBody +from flyguard.pipeline import FlyGuard, Params +from flyguard.synth import IntensityEnv, Placement, catalogue, inject + +HOLDOUT = "doubleT_obstacle" # там уже есть настоящий объект + + +def _work(task): + """Одна задача — один бэг. + + Обученное грузится путями и уже внутри процесса: передавать модели через + межпроцессную границу незачем, а свою складку каждый воркер берёт сам. + """ + path, params, limit, d_start, laterals, seed, mem_path, rd_path = task + memory = MushroomBody.load(mem_path) if mem_path else None + readout = None + if rd_path: + from flyguard.mbon_readout import MbonReadout + readout = MbonReadout.load(rd_path) + return run_bag(path, params, memory, limit, d_start, laterals, seed, + readout=readout) + + +def ego_track(bag: Bag, params: Params, limit: int | None): + """Первый проход: пройденный путь на каждом кадре (разметка по дистанции).""" + fg = FlyGuard(params, memory=None) + s, stamps = [], [] + total = 0.0 + for _, pc in bag.frames(stop=limit): + res = fg.process(pc) + if res is None: + s.append(None); stamps.append(pc.stamp); continue + total += res.ego.ds if res.ego else 0.0 + s.append(total); stamps.append(pc.stamp) + return s, stamps + + +def run_bag(bag_path, params: Params, memory, limit: int, d_start: float, + laterals: tuple[float, ...], seed: int, readout=None) -> list[dict]: + bag = Bag(bag_path) + s_track, _ = ego_track(bag, params, limit) + have = [x for x in s_track if x is not None] + if len(have) < 20: + return [] + travel = have[-1] - have[0] + + cat = catalogue() + scen = [(name, lat) for name in cat for lat in laterals] + pipes = [FlyGuard(params, memory=memory, readout=readout) for _ in scen] + rng = np.random.default_rng(seed) + records = [[] for _ in scen] + + # решётка и поза нужны для вставки — берутся из отдельного «чистого» конвейера + guide = FlyGuard(params, memory=None) + + for k, (_, pc) in enumerate(bag.frames(stop=limit)): + gres = guide.process(pc) + if gres is None or s_track[k] is None: + continue + s_now = s_track[k] - have[0] + env = IntensityEnv(pc) # один раз на кадр, общий для сценариев + for i, (name, lat) in enumerate(scen): + d_true = d_start - s_now + if d_true < 6.0: + continue + # Предмет лежит НА ПУТИ, а путь в кривой уходит вбок: на 150 м при + # радиусе 1300 м это 8.6 м. Если ставить его в поперечных координатах + # сенсора, он окажется в стене, а не в габарите. + u_obj = float(gres.corridor.centre(np.array([d_true], np.float32))[0]) + lat + pc2, lab = inject(pc, guide.layout_full, gres.plane, cat[name], + Placement(d=d_true, u=u_obj), rng=rng, env=env) + res = pipes[i].process(pc2) + if res is None: + continue + tol = max(3.0, 0.12 * d_true) + hit = any(abs(o.distance - d_true) < tol for o in res.decision.objects) + fp = sum(1 for o in res.decision.objects if abs(o.distance - d_true) >= tol) + # Воронка потерь. Лучи в предмет попали — а дальше он может + # пропасть на любой из трёх ступеней, и лечатся они по-разному: + # нет кандидата — вопрос к кластеризации и разделению фигуры и + # фона, нет трека — к сопоставлению по кадрам, нет решения — + # к порогу. Без этого разбиения улучшать нечего, кроме удачи. + cand = any(abs(c.d - d_true) < tol for c in res.candidates) + cx = pipes[i].cx + trk = any(abs(tr.distance(cx.s_world) - d_true) < tol + for tr in cx.tracks) + records[i].append((d_true, int(hit), fp, lab["hit_rays"], + int(cand), int(trk))) + + out = [] + for (name, lat), rec in zip(scen, records): + if not rec: + continue + a = np.array(rec, np.float32) + out.append(dict(bag=bag.path.name, obj=name, lateral=lat, + d=a[:, 0].tolist(), hit=a[:, 1].tolist(), + fp=a[:, 2].tolist(), rays=a[:, 3].tolist(), + cand=a[:, 4].tolist(), trk=a[:, 5].tolist(), + travel=float(travel))) + return out + + +def main() -> None: + ap = argparse.ArgumentParser(description=__doc__) + ap.add_argument("--root", default=str(B.DATA / "for_hackathon")) + ap.add_argument("--memory") + ap.add_argument("--mbon", default="", + help="обученное считывание MBON одной моделью; она видела " + "эти бэги — годится только для отладки") + ap.add_argument("--mbon-dir", default="", + help="каталог с моделями по складкам (mbon_<бэг>.npz): для " + "каждого бэга берётся модель, его НЕ видевшая. Иначе " + "дальность завышена: считывание обучалось ровно на " + "таких же вставках в этот же тоннель") + ap.add_argument("--out", default=str(B.ARTIFACTS / "benchmark.json")) + ap.add_argument("--limit", type=int, default=250) + ap.add_argument("--d-start", type=float, default=200.0) + ap.add_argument("--laterals", default="0.0,0.9") + ap.add_argument("--seed", type=int, default=12345) + ap.add_argument("--mbon-blend", type=float, default=None, + help="1 — только модель, 0 — только ручная формула") + ap.add_argument("--warn", type=float, default=None, + help="порог улики для тревоги (по умолчанию 0.5)") + ap.add_argument("--clear", type=float, default=None) + ap.add_argument("--min-hits", type=int, default=None) + ap.add_argument("--warn-far", type=float, default=None, + help="порог тревоги на дальнем краю; ниже обычного — послабление далёким трекам") + ap.add_argument("--warn-far-from", type=float, default=None, + help="с какой дальности порог начинает падать, м") + ap.add_argument("--novelty-floor", type=float, default=None, + help="ниже этой новизны трек не считается") + ap.add_argument("--near-long", type=float, default=None, + help="разброс дальности, выше которого компоненту режут и вблизи; 0 — не резать вблизи") + ap.add_argument("--min-rays", type=int, default=None, + help="сколько лучей минимум образуют кандидата") + ap.add_argument("--min-rays-far", type=int, default=None, + help="порог по лучам за `--min-rays-far-from`; 0 — не различать") + ap.add_argument("--min-rays-far-from", type=float, default=None) + ap.add_argument("--leak-far", type=float, default=None, + help="утечка улики на дальнем краю; ниже обычной — далёкий трек прощает промахи") + ap.add_argument("--leak-far-from", type=float, default=None, + help="с какой дальности утечка начинает падать, м") + ap.add_argument("--nov-fade-from", type=float, default=None, + help="с какой дальности гасить вклад знакомости; 0 — не гасить") + ap.add_argument("--nov-fade-to", type=float, default=None, + help="к какой дальности вклад знакомости обнуляется") + ap.add_argument("--mbon-prior-from", type=float, default=None, + help="с какой дальности поправлять оценку модели на распространённость предметов; 0 — не поправлять") + P.add_argument(ap) + args = ap.parse_args() + + B.ARTIFACTS.mkdir(parents=True, exist_ok=True) + fold_paths: dict[str, str] = {} + if args.mbon_dir: + from pathlib import Path as _P + for f in _P(args.mbon_dir).glob("mbon_*.npz"): + fold_paths[f.stem[len("mbon_"):]] = str(f) + print(f"считывание MBON по складкам: {args.mbon_dir} " + f"({len(fold_paths)} моделей)") + elif args.mbon: + print(f"считывание MBON: {args.mbon}") + over = {} + if args.mbon_blend is not None: + over["mbon_blend"] = args.mbon_blend + if args.warn is not None: + over["warn_evidence"] = args.warn + over["clear_evidence"] = args.warn * 0.6 if args.clear is None else args.clear + elif args.clear is not None: + over["clear_evidence"] = args.clear + if args.min_hits is not None: + over["min_hits"] = args.min_hits + if args.warn_far is not None: + over["warn_far"] = args.warn_far + if args.warn_far_from is not None: + over["warn_far_from"] = args.warn_far_from + if args.novelty_floor is not None: + over["novelty_floor"] = args.novelty_floor + if args.near_long is not None: + over["near_long"] = args.near_long + if args.min_rays is not None: + over["min_rays"] = args.min_rays + if args.min_rays_far is not None: + over["min_rays_far"] = args.min_rays_far + if args.min_rays_far_from is not None: + over["min_rays_far_from"] = args.min_rays_far_from + if args.leak_far is not None: + over["leak_far"] = args.leak_far + if args.leak_far_from is not None: + over["leak_far_from"] = args.leak_far_from + if args.nov_fade_from is not None: + over["nov_fade_from"] = args.nov_fade_from + if args.nov_fade_to is not None: + over["nov_fade_to"] = args.nov_fade_to + if args.mbon_prior_from is not None: + over["mbon_prior_from"] = args.mbon_prior_from + params = Params(**over) + laterals = tuple(float(x) for x in args.laterals.split(",")) + + tasks = [] + for p in find_bags(args.root): + if p.name == HOLDOUT: + continue + rd = args.mbon + if fold_paths: + rd = fold_paths.get(p.name, "") + if not rd: + print(f" внимание: для {p.name} нет своей складки — пропуск") + continue + tasks.append((p, params, args.limit, args.d_start, laterals, args.seed, + args.memory or "", rd)) + + # Печатается по готовности, собирается по номеру задачи: порядок сценариев + # в файле не должен зависеть от того, какой бэг досчитался первым. + slots: list = [None] * len(tasks) + for i, task, rec, secs in P.run(_work, tasks, args.jobs): + slots[i] = rec + n = sum(len(r["d"]) for r in rec) + print(f" {task[0].name:42s} сценариев {len(rec):3d}, наблюдений {n:6d}, " + f"{secs:6.1f} с", flush=True) + all_rec = [r for rec in slots if rec for r in rec] + + with open(args.out, "w", encoding="utf-8") as f: + json.dump(all_rec, f, ensure_ascii=False) + print(f"сохранено: {args.out} ({len(all_rec)} сценариев)") + + +if __name__ == "__main__": + main() diff --git a/tools/make_training_set.py b/tools/make_training_set.py index 8bdcc41..d855323 100644 --- a/tools/make_training_set.py +++ b/tools/make_training_set.py @@ -1,177 +1,186 @@ -"""Размеченная выборка кандидатов: предмет против тоннельной обстановки. - -Разметки в датасете нет, и до сих пор это определяло архитектуру: грибовидное -тело учится **без меток**, запоминая частоту обстановки. Но у нас есть -физически обоснованный генератор предметов (`flyguard.synth`), сверенный с -единственным реальным объектом: настоящий 0.67 × 1.35 м на 55 м даёт 47–69 -лучей, синтетический человек 0.44 × 1.71 м на 60 м — 50. Значит, метки можно -изготовить, и изготовить достоверно. - -Каждый кандидат помечается по **пересечению лучей**, а не «по дальности -примерно»: `inject` возвращает индексы лучей, в которые предмет действительно -записан, и кандидат считается предметом, если его ядро состоит из этих лучей. -Так структура тоннеля, случайно оказавшаяся на той же дальности, в -положительные не попадает. - -Сценарии намеренно ставят предмет в РАЗНЫЕ точки тоннеля (`--d-starts`): замер -показал, что одна и та же дальность в разных местах перегона ведёт себя -совершенно по-разному — где-то предмет виден целиком, где-то за поворотом -(EXPERIMENTS п. 9.6). Обучаться на одной точке постановки значит выучить эту -точку. - -`doubleT_obstacle` исключён целиком: там настоящий объект, и он остаётся -независимой проверкой того, что модель выучила предмет, а не «синтетику». - - python tools/make_training_set.py --out data/cache/training_set.npz -""" -from __future__ import annotations - -import argparse -import time - -import numpy as np - -import _bootstrap as B # noqa: F401 -from flyguard.bag import Bag, find_bags -from flyguard.mushroom_body import FEATURES, describe -from flyguard.pipeline import FlyGuard, Params -from flyguard.synth import IntensityEnv, Placement, catalogue, inject - -HOLDOUT = "doubleT_obstacle" # там реальный объект — только для проверки -MIN_OVERLAP = 0.5 # доля лучей ядра, пришедших от предмета - - -def ego_track(bag: Bag, params: Params, limit: int): - """Первый проход: пройденный путь на каждом кадре и общая решётка.""" - fg = FlyGuard(params, memory=None) - s, total = [], 0.0 - for _, pc in bag.frames(stop=limit): - res = fg.process(pc) - if res is None: - s.append(None) - continue - total += res.ego.ds if res.ego else 0.0 - s.append(total) - return s, fg - - -def collect_bag(path, params: Params, limit: int, d_starts, laterals, seed: int): - bag = Bag(path) - s_track, ego_fg = ego_track(bag, params, limit) - have = [x for x in s_track if x is not None] - if len(have) < 30: - return None - s0 = have[0] - layout = ego_fg.layout_full - col0 = ego_fg.cols.start - - cat = catalogue() - scen = [(n, lat, d0) for n in cat for lat in laterals for d0 in d_starts] - guide = FlyGuard(params, memory=None, layout=layout) - pipes = [FlyGuard(params, memory=None, layout=layout) for _ in scen] - rng = np.random.default_rng(seed) - - X, y, dd, obj, lat_out = [], [], [], [], [] - for k, (_, pc) in enumerate(bag.frames(stop=limit)): - g = guide.process(pc) - if g is None or s_track[k] is None: - continue - s_now = s_track[k] - s0 - env = IntensityEnv(pc) # один раз на кадр, общий для сценариев - for i, (name, lat, d0) in enumerate(scen): - d_true = d0 - s_now - if d_true < 6.0: - continue - u = float(g.corridor.centre(np.array([d_true], np.float32))[0]) + lat - pc2, lab = inject(pc, layout, g.plane, cat[name], - Placement(d=d_true, u=u), rng=rng, env=env) - res = pipes[i].process(pc2) - if res is None or not res.candidates: - continue - rr = lab.get("rays") - if rr is None or lab["hit_rays"] == 0: - truth = None - else: - # лучи предмета в координатах полной решётки → плоский индекс; - # отсортированный массив, а не множество: проверка идёт - # сотни тысяч раз, и `in` по множеству тут заметно дороже - truth = np.sort((rr[0].astype(np.int64) << 20) - | rr[1].astype(np.int64)) - for c in res.candidates: - ii, jj = c.extra.get("rays", (None, None)) - if ii is None: - continue - lbl = 0 - if truth is not None and truth.size: - key = ((ii.astype(np.int64) << 20) - | (jj.astype(np.int64) + col0)) - pos = np.searchsorted(truth, key) - np.clip(pos, 0, truth.size - 1, out=pos) - frac = float((truth[pos] == key).mean()) - lbl = int(frac >= MIN_OVERLAP) - v = describe(c) - acc = float(c.extra.get("acc_support", 0.0)) if c.extra else 0.0 - X.append(np.append(v, acc).astype(np.float32)) - y.append(lbl) - dd.append(c.d) - obj.append(name if lbl else "") - lat_out.append(lat) - return (np.asarray(X, np.float32), np.asarray(y, np.int8), - np.asarray(dd, np.float32), np.asarray(obj), - np.asarray(lat_out, np.float32)) - - -def main() -> None: - ap = argparse.ArgumentParser(description=__doc__) - ap.add_argument("--root", default=str(B.DATA / "for_hackathon")) - ap.add_argument("--out", default=str(B.CACHE / "training_set.npz")) - ap.add_argument("--limit", type=int, default=250) - ap.add_argument("--d-starts", default="80,140,200") - # Поперечные положения намеренно кроют ВЕСЬ габарит, а не только ось. - # Иначе выборка вырождается: предметы у оси, обстановка у стен, и модель - # выучивает «всё, что у оси — предмет» вместо признаков предмета. Проверено - # на пробном прогоне с одной постановкой: AUC 1.000 по всем бэгам — цифра - # красивая и бессмысленная, а в тоннеле у оси полно штатных конструкций. - ap.add_argument("--laterals", default="0,-0.6,0.6,-1.2,1.2") - ap.add_argument("--seed", type=int, default=20260921) - args = ap.parse_args() - - d_starts = tuple(float(x) for x in args.d_starts.split(",")) - laterals = tuple(float(x) for x in args.laterals.split(",")) - params = Params() - B.CACHE.mkdir(parents=True, exist_ok=True) - - parts = {} - for p in find_bags(args.root): - if p.name == HOLDOUT: - continue - t0 = time.time() - got = collect_bag(p, params, args.limit, d_starts, laterals, args.seed) - if got is None: - print(f" {p.name:42s} пропущен") - continue - parts[p.name] = got - Xb, yb = got[0], got[1] - print(f" {p.name:42s} {Xb.shape[0]:7d} кандидатов, " - f"предметов {int(yb.sum()):6d} ({yb.mean():5.1%}), " - f"{time.time() - t0:6.0f} с", flush=True) - - if not parts: - raise SystemExit("ничего не собрано") - out = {} - for name, (Xb, yb, db, ob, lb) in parts.items(): - out[f"X_{name}"] = Xb - out[f"y_{name}"] = yb - out[f"d_{name}"] = db - out[f"obj_{name}"] = ob - out[f"lat_{name}"] = lb - np.savez_compressed(args.out, names=np.array(list(parts)), - features=np.array(FEATURES + ("acc_support",)), **out) - tot = sum(v[0].shape[0] for v in parts.values()) - pos = sum(int(v[1].sum()) for v in parts.values()) - print(f"\nвсего {tot} кандидатов, предметов {pos} ({pos / tot:.1%})") - print("сохранено:", args.out) - - -if __name__ == "__main__": - main() +"""Размеченная выборка кандидатов: предмет против тоннельной обстановки. + +Разметки в датасете нет, и до сих пор это определяло архитектуру: грибовидное +тело учится **без меток**, запоминая частоту обстановки. Но у нас есть +физически обоснованный генератор предметов (`flyguard.synth`), сверенный с +единственным реальным объектом: настоящий 0.67 × 1.35 м на 55 м даёт 47–69 +лучей, синтетический человек 0.44 × 1.71 м на 60 м — 50. Значит, метки можно +изготовить, и изготовить достоверно. + +Каждый кандидат помечается по **пересечению лучей**, а не «по дальности +примерно»: `inject` возвращает индексы лучей, в которые предмет действительно +записан, и кандидат считается предметом, если его ядро состоит из этих лучей. +Так структура тоннеля, случайно оказавшаяся на той же дальности, в +положительные не попадает. + +Сценарии намеренно ставят предмет в РАЗНЫЕ точки тоннеля (`--d-starts`): замер +показал, что одна и та же дальность в разных местах перегона ведёт себя +совершенно по-разному — где-то предмет виден целиком, где-то за поворотом +(EXPERIMENTS п. 9.6). Обучаться на одной точке постановки значит выучить эту +точку. + +`doubleT_obstacle` исключён целиком: там настоящий объект, и он остаётся +независимой проверкой того, что модель выучила предмет, а не «синтетику». + + python tools/make_training_set.py --out data/cache/training_set.npz +""" +from __future__ import annotations + +import argparse + +import numpy as np + +import _bootstrap as B # noqa: F401 +import _parallel as P +from flyguard.bag import Bag, find_bags +from flyguard.mushroom_body import FEATURES, describe +from flyguard.pipeline import FlyGuard, Params +from flyguard.synth import IntensityEnv, Placement, catalogue, inject + +HOLDOUT = "doubleT_obstacle" # там реальный объект — только для проверки +MIN_OVERLAP = 0.5 # доля лучей ядра, пришедших от предмета + + +def _work(task): + """Одна задача — один бэг. Верхнего уровня: иначе не передать в процесс.""" + path, params, limit, d_starts, laterals, seed = task + return collect_bag(path, params, limit, d_starts, laterals, seed) + + +def ego_track(bag: Bag, params: Params, limit: int): + """Первый проход: пройденный путь на каждом кадре и общая решётка.""" + fg = FlyGuard(params, memory=None) + s, total = [], 0.0 + for _, pc in bag.frames(stop=limit): + res = fg.process(pc) + if res is None: + s.append(None) + continue + total += res.ego.ds if res.ego else 0.0 + s.append(total) + return s, fg + + +def collect_bag(path, params: Params, limit: int, d_starts, laterals, seed: int): + bag = Bag(path) + s_track, ego_fg = ego_track(bag, params, limit) + have = [x for x in s_track if x is not None] + if len(have) < 30: + return None + s0 = have[0] + layout = ego_fg.layout_full + col0 = ego_fg.cols.start + + cat = catalogue() + scen = [(n, lat, d0) for n in cat for lat in laterals for d0 in d_starts] + guide = FlyGuard(params, memory=None, layout=layout) + pipes = [FlyGuard(params, memory=None, layout=layout) for _ in scen] + rng = np.random.default_rng(seed) + + X, y, dd, obj, lat_out = [], [], [], [], [] + for k, (_, pc) in enumerate(bag.frames(stop=limit)): + g = guide.process(pc) + if g is None or s_track[k] is None: + continue + s_now = s_track[k] - s0 + env = IntensityEnv(pc) # один раз на кадр, общий для сценариев + for i, (name, lat, d0) in enumerate(scen): + d_true = d0 - s_now + if d_true < 6.0: + continue + u = float(g.corridor.centre(np.array([d_true], np.float32))[0]) + lat + pc2, lab = inject(pc, layout, g.plane, cat[name], + Placement(d=d_true, u=u), rng=rng, env=env) + res = pipes[i].process(pc2) + if res is None or not res.candidates: + continue + rr = lab.get("rays") + if rr is None or lab["hit_rays"] == 0: + truth = None + else: + # лучи предмета в координатах полной решётки → плоский индекс; + # отсортированный массив, а не множество: проверка идёт + # сотни тысяч раз, и `in` по множеству тут заметно дороже + truth = np.sort((rr[0].astype(np.int64) << 20) + | rr[1].astype(np.int64)) + for c in res.candidates: + ii, jj = c.extra.get("rays", (None, None)) + if ii is None: + continue + lbl = 0 + if truth is not None and truth.size: + key = ((ii.astype(np.int64) << 20) + | (jj.astype(np.int64) + col0)) + pos = np.searchsorted(truth, key) + np.clip(pos, 0, truth.size - 1, out=pos) + frac = float((truth[pos] == key).mean()) + lbl = int(frac >= MIN_OVERLAP) + v = describe(c) + acc = float(c.extra.get("acc_support", 0.0)) if c.extra else 0.0 + X.append(np.append(v, acc).astype(np.float32)) + y.append(lbl) + dd.append(c.d) + obj.append(name if lbl else "") + lat_out.append(lat) + return (np.asarray(X, np.float32), np.asarray(y, np.int8), + np.asarray(dd, np.float32), np.asarray(obj), + np.asarray(lat_out, np.float32)) + + +def main() -> None: + ap = argparse.ArgumentParser(description=__doc__) + ap.add_argument("--root", default=str(B.DATA / "for_hackathon")) + ap.add_argument("--out", default=str(B.CACHE / "training_set.npz")) + ap.add_argument("--limit", type=int, default=250) + ap.add_argument("--d-starts", default="80,140,200") + # Поперечные положения намеренно кроют ВЕСЬ габарит, а не только ось. + # Иначе выборка вырождается: предметы у оси, обстановка у стен, и модель + # выучивает «всё, что у оси — предмет» вместо признаков предмета. Проверено + # на пробном прогоне с одной постановкой: AUC 1.000 по всем бэгам — цифра + # красивая и бессмысленная, а в тоннеле у оси полно штатных конструкций. + ap.add_argument("--laterals", default="0,-0.6,0.6,-1.2,1.2") + ap.add_argument("--seed", type=int, default=20260921) + P.add_argument(ap) + args = ap.parse_args() + + d_starts = tuple(float(x) for x in args.d_starts.split(",")) + laterals = tuple(float(x) for x in args.laterals.split(",")) + params = Params() + B.CACHE.mkdir(parents=True, exist_ok=True) + + bags = [p for p in find_bags(args.root) if p.name != HOLDOUT] + tasks = [(p, params, args.limit, d_starts, laterals, args.seed) for p in bags] + # Печатается по готовности, складывается по номеру задачи: порядок бэгов + # в файле не должен зависеть от того, какой из них досчитался первым. + slots: list = [None] * len(tasks) + for i, task, got, secs in P.run(_work, tasks, args.jobs): + name = task[0].name + if got is None: + print(f" {name:42s} пропущен") + continue + slots[i] = (name, got) + Xb, yb = got[0], got[1] + print(f" {name:42s} {Xb.shape[0]:7d} кандидатов, " + f"предметов {int(yb.sum()):6d} ({yb.mean():5.1%}), " + f"{secs:6.0f} с", flush=True) + parts = {name: got for name, got in (s for s in slots if s is not None)} + + if not parts: + raise SystemExit("ничего не собрано") + out = {} + for name, (Xb, yb, db, ob, lb) in parts.items(): + out[f"X_{name}"] = Xb + out[f"y_{name}"] = yb + out[f"d_{name}"] = db + out[f"obj_{name}"] = ob + out[f"lat_{name}"] = lb + np.savez_compressed(args.out, names=np.array(list(parts)), + features=np.array(FEATURES + ("acc_support",)), **out) + tot = sum(v[0].shape[0] for v in parts.values()) + pos = sum(int(v[1].sum()) for v in parts.values()) + print(f"\nвсего {tot} кандидатов, предметов {pos} ({pos / tot:.1%})") + print("сохранено:", args.out) + + +if __name__ == "__main__": + main()