Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion docs/ARCHITECTURE_RU.md
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@
| Типы и состояния транспорта | `hadleys/enums.py`; данные — `hadleys/models.py` |
| Интернет, люди, финансы, аварии, безопасность | Одноимённые файлы в `hadleys/domains/` |
| Деньги домохозяйств, счета, кредиты, банкротство | `hadleys/domains/households.py`; правила — `docs/FINANCE_RU.md` |
| Профилирование и замеры до/после | `scripts/profile_run.py`, `scripts/bench_ticks.py`; результаты — `docs/PERF_RU.md` |
| Профилирование и замеры до/после | `scripts/profile_run.py`, `scripts/bench_ticks.py`, `scripts/bench_mqtt.py`; результаты — `docs/PERF_RU.md` |
| Баланс воды и стоков, счета объёмов | `hadleys/domains/hydraulics.py` (`water_balance`); описание — `docs/WATER_RU.md`; сценарий — `scripts/water_scenario.py` |
| Финансовые сценарии A/B/C и эталонный расчёт | `hadleys/scenarios.py`, `hadleys/finance_reference.py`, `scripts/finance_acceptance.py` |
| Данные аттракторов | `hadleys/domains/attractors.py` |
Expand Down
76 changes: 75 additions & 1 deletion docs/PERF_RU.md
Original file line number Diff line number Diff line change
Expand Up @@ -83,4 +83,78 @@ python3 scripts/bench_ticks.py data/late # сравнени
2. **Структура решателя давления.** Оставшиеся ~1–1,5 мс на тик — это накладные расходы numpy на 22 уровнях дерева. Можно заменить проходы по уровням одним умножением на разреженную матрицу путей. Это меняет порядок сложения (биты в последних разрядах), поэтому нужно отдельно решить, как обновлять эталонную фикстуру. Воспроизведение: `cProfile` в `profile_run.py` (`solve`, `aggregate`, `headloss`).
3. **Скорость больше ~190 минут в секунду недостижима.** `sim_loop` отбрасывает тики сверх 200 за шаг, а интерфейс показывает выставленную скорость. Предложение: показывать фактическую скорость или ограничить ползунок.
4. **Затор на кольцевой дороге** (уже вынесен отдельной задачей). Он косвенно создаёт поздний мир с сотнями прорванных труб: без него решатель давления реже работает с отрезанными домами.
5. **`world.pkl` 5,1 МБ.** Журнал операций домохозяйств хранит 300 × 256 строк, включая 92 пустых дома. Можно хранить только занятые дома или уменьшить длину журнала. Сохранение раз в 5 минут, так что это не срочно.
5. **`world.pkl` 5,1 МБ.** Журнал операций домохозяйств хранит 300 × 256 строк, включая 92 пустых дома. Можно хранить только занятые дома или уменьшить длину журнала. Сохранение раз в 5 минут, так что это не срочно (замер ниже: 9–11 мс раз в 5 минут).

# Хранение и шина MQTT

Все замеры: образ из `Dockerfile` (`python:3.12-slim`, numpy 2.3.5, paho-mqtt 2.1.0), брокер `eclipse-mosquitto:2`, контроллеры домов — `houses_runtime.py` в том же образе. Мир 30-го дня собирается так же, как выше: `scripts/finance_acceptance.py --days 30 --out data/late`.

## Как воспроизвести

Хранение (`world.pkl` и `history.db`) — отдельный раздел в `profile_run.py`, свежий мир и сохранённый:

```bash
python3 scripts/profile_run.py --days 0.5 # свежий мир
python3 scripts/profile_run.py --days 0.5 --data data/late # мир из каталога; скрипт его не перезаписывает
```

Раздел `persistence` печатается двумя строками и попадает в `--out`: время `save_world` и `load_world` (по 5 раз), `record_hour` (60 раз), запрос `history(720)` за 30 суток, размер `world.pkl` и десять самых больших полей мира в нём. Всё пишется во временный каталог, настоящие данные не трогаются.

Шина MQTT — `scripts/bench_mqtt.py` против живого брокера с отвечающими контроллерами:

```bash
docker compose -f scripts/compose.bench.yaml up --build --attach bench --abort-on-container-exit --exit-code-from bench
docker compose -f scripts/compose.bench.yaml down
```

`scripts/compose.bench.yaml` поднимает Mosquitto, контроллеры домов и сам замер; код монтируется из текущей копии репозитория, так что для сравнения «до/после» достаточно переключить файл (`git stash push -- hadleys/integrations/mqtt.py`, замер, `git stash pop`), пересобирать образ не нужно. Число тиков — переменная `TICKS` (по умолчанию 600). Без Docker: локальный `mosquitto`, `python3 houses_runtime.py` и `python3 scripts/bench_mqtt.py localhost:1883`.

Замер в CI не запускается: время на общих раннерах слишком плавает для порогов. Поведение шины (сколько домов уходит на тик) проверяет unit-тест, он в CI есть.

`bench_mqtt.py` берёт свежий мир, подключает `MqttBridge` и 600 тиков подряд меряет отдельно `world_tick`, `bridge.apply()` (приём команд) и `bridge.publish()` (отправка датчиков), с паузой 20 мс на тик, чтобы контроллеры успевали отвечать. Печатает, сколько домов ушло на каждом тике, и времена (среднее, p95, максимум). До и после сравнивались поочерёдно, по два прогона: версия до — `hadleys/integrations/mqtt.py` из коммита `612f587`.

## Хранение: узкого места нет

| | свежий мир | мир 30-го дня | как часто |
|---|---:|---:|---|
| `save_world`, медиана / максимум | 8,8 / 10,2 мс | 11,1 / 12,0 мс | раз в 5 минут, под блокировкой мира |
| `load_world` | 51 мс | 54 мс | при старте |
| `record_hour`, медиана / максимум | 2,7 / 4,1 мс | 2,8 / 5,4 мс | раз в 60 тиков |
| `history(720)` | 0,4 мс | 0,3 мс | по запросу `/history` |
| `world.pkl` | 5,10 МБ | 5,33 МБ | |

В пересчёте на тик это около 0,05 мс (`record_hour`) и тысячные доли миллисекунды (`save_world`) при тике ~5 мс. В SQLite пишется одна строка в игровой час, пакетировать нечего.

Состав `world.pkl` (одинаков в обоих мирах): `hh_led` (журнал операций домохозяйств) 3,2 МБ — 61–63 %, `utilities` 0,65 МБ, `h_hist` и `h_hist_draw` по 0,35 МБ, `hh_hist` 0,29 МБ, остальное меньше 60 КБ на поле.

Что стоит учитывать не ради скорости, а ради надёжности: мир живёт только в памяти, при падении процесса теряется до 5 минут; `save_world` не вызывает `fsync` перед `os.replace`.

## MQTT: синхронная пачка раз в 10 тиков

**Что было.** `publish()` отправлял дом, если он изменился или прошло ≥ 10 тиков с его прошлой отправки. После подключения уходят все 300 домов сразу, у всех одинаковый `last_pub_t`, и дальше они обновляются тоже все сразу. В свежем мире дома почти не меняются, поэтому 9 тиков из 10 не отправляли ничего, а каждый десятый — все 300 сообщений, ~15 мс под блокировкой мира. Страдал и сам тик: с отвечающими контроллерами его p95 был 12,6–14,8 мс, в прогоне без контроллеров — 6,1 мс. Вероятная причина — пачка из 300 ответов, которую сетевой поток paho разбирает под GIL одновременно с тиком (отдельно не профилировалось).

**Что сделано** (`hadleys/integrations/mqtt.py`, `publish`). Периодическое обновление идёт в собственной фазе дома: `id % 10 == t % 10`. Каждый дом по-прежнему уходит ровно раз в 10 тиков, но по 30 за тик. Порог 10 не меняется: дом без команды дольше 15 тиков переходит на свой термостат (`houses.py`, `ext`). После (пере)подключения, пока `last_pub_t < 0`, по-прежнему уходят все дома за один тик — это разовая и намеренная отправка. Изменившиеся дома отправляются сразу, как и раньше.

Одно отличие: раньше отправка из-за изменения сбрасывала отсчёт 10 тиков, теперь дом всё равно уйдёт в свою фазу, даже если его только что отправили. Сверху это не больше `N / 10` = 30 лишних сообщений за тик. На практике разницы нет: с подменённым клиентом paho по 600 тиков на свежем мире, на свежем мире со снежной бурей (`inject(w, "storm")`) и на мире 30-го дня обе версии отправляют в среднем ровно 30,0 сообщения за тик — дома почти не пересекают порог изменения (0,2 °C или смена флагов). Максимум за тик: 300 до, 30 после.

Без `--mqtt` мост не создаётся, так что на симуляцию, эталонную фикстуру и `bench_ticks.py` правка не влияет.

| | до | после |
|---|---|---|
| домов на тик | 0 на 540 тиках, 300 на 60 | 30 на всех 600 |
| `publish()`, среднее | 2,1 / 1,8 мс | 1,9 / 1,9 мс |
| `publish()`, p95 | 15,5 / 15,2 мс | 2,2 / 2,2 мс |
| `publish()`, максимум | 35 / 24 мс | 8,6 / 8,9 мс |
| тик, p95 | 14,8 / 12,6 мс | 6,1 / 6,1 мс |
| тик, максимум | 39 / 18 мс | 7,7 / 8,4 мс |

Среднее не изменилось: сообщений столько же, исчезли только пики.

**Брокер не узкое место:** Mosquitto во время замера — 0,03 % CPU и 3 МБ памяти при ~30 сообщениях в каждую сторону за тик. Время уходит в нашем процессе: `json.dumps` и `client.publish` paho, около 60 мкс на сообщение, то есть ~1,9 мс на тик. Сравнивать брокеры имеет смысл только после того, как нагрузка вырастет на порядки.

**Как проверено.** `tests/test_mqtt.py::test_periodic_refresh_is_spread_over_ticks` без брокера (клиент paho подменён `MagicMock`): после первой полной отправки 30 тиков подряд, ничего в домах не меняется; на каждом тике уходит ровно `N // 10` домов, каждый дом — ровно через 10 тиков после своей прошлой отправки, за 30 тиков отправлены все. На старом коде тест падает (на этих тиках уходит 0 домов). Полный набор тестов, `npm run check`, сборка образа и сутки в образе проходят; браузерный smoke (`npm run test:browser`) не прогонялся — правка его не затрагивает, но в CI он пройдёт отдельно.

## Что осталось

1. **Сериализация сообщений.** ~60 мкс на сообщение в `publish()`: 30 отдельных `json.dumps` и `client.publish` за тик. Можно отправлять одно сообщение на сектор или на тик, но это меняет протокол с контроллерами (`hh/house/{id}/sensors`), поэтому нужно отдельное решение. Воспроизведение: `bench_mqtt.py`, строка `publish_ms`.
2. **Надёжность хранения.** Потеря до 5 минут при падении и отсутствие `fsync` — см. выше. Скорость здесь ни при чём.
5 changes: 4 additions & 1 deletion hadleys/integrations/mqtt.py
Original file line number Diff line number Diff line change
Expand Up @@ -135,7 +135,10 @@ def publish(self):
changed = (
(np.abs(w.h_t_in - self.last_t_in) >= 0.2)
| (flags != self.last_flags).any(axis=1)
| (w.t - self.last_pub_t >= 10)
| (self.last_pub_t < 0) # never sent since (re)connect: everything at once
# periodic refresh in each house's own phase: every house still every 10 ticks, but 30 per tick
# instead of all 300 on the same tick (a synchronized burst held the world lock for ~15 ms)
| (np.arange(w.N) % 10 == w.t % 10)
)
for i in np.flatnonzero(changed):
payload = {
Expand Down
76 changes: 76 additions & 0 deletions scripts/bench_mqtt.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,76 @@
"""Tick and MQTT bridge cost against a live broker with the house controllers answering.

python3 scripts/bench_mqtt.py host:port [ticks]

Needs a broker and houses_runtime.py connected to it (see docs/PERF_RU.md). Fresh world, the bridge's apply()
and publish() are timed separately from world_tick, with a 20 ms pause per tick so the controllers answer as in
the real loop. Prints how many houses go out per tick and the timings; one JSON line at the end.
"""

import collections
import json
import os
import statistics
import sys
import time

sys.path.insert(0, os.path.join(os.path.dirname(os.path.abspath(__file__)), ".."))

from hadleys.integrations.mqtt import MqttBridge # noqa: E402
from hadleys.simulation import world_tick # noqa: E402
from hadleys.world import World # noqa: E402


def stats(ms):
s = sorted(ms)
return {"mean": round(statistics.fmean(s), 2), "p95": round(s[int(0.95 * (len(s) - 1))], 2), "max": round(s[-1], 2)}


def main():
url = sys.argv[1]
ticks = int(sys.argv[2]) if len(sys.argv) > 2 else 600
w = World()
w.bridge = None # the bridge calls are timed here, outside world_tick
for _ in range(30):
world_tick(w)
b = MqttBridge(w, url)
for _ in range(50):
if b.connected:
break
time.sleep(0.1)
if not b.connected:
raise SystemExit(f"broker {url} not reachable")
b.publish() # the full send after a connect
time.sleep(1)
tick, apply, publish, per_tick = [], [], [], []
sent0, received0 = b.sent, b.received
for _ in range(ticks):
t0 = time.perf_counter()
world_tick(w)
tick.append((time.perf_counter() - t0) * 1000)
t0 = time.perf_counter()
b.apply()
apply.append((time.perf_counter() - t0) * 1000)
s0 = b.sent
t0 = time.perf_counter()
b.publish()
publish.append((time.perf_counter() - t0) * 1000)
per_tick.append(b.sent - s0)
time.sleep(0.02)
result = {
"ticks": ticks,
"houses_per_tick": dict(collections.Counter(per_tick).most_common(6)),
"tick_ms": stats(tick),
"apply_ms": stats(apply),
"publish_ms": stats(publish),
"sent": b.sent - sent0,
"received": b.received - received0,
}
print("houses sent per tick -> number of ticks:", result["houses_per_tick"])
for k in ("tick_ms", "apply_ms", "publish_ms"):
print(f"{k}: {result[k]}")
print(json.dumps(result))


if __name__ == "__main__":
main()
21 changes: 21 additions & 0 deletions scripts/compose.bench.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,21 @@
# MQTT bench in one command: broker + house controllers + scripts/bench_mqtt.py, code mounted from this checkout.
# docker compose -f scripts/compose.bench.yaml up --build --attach bench --abort-on-container-exit --exit-code-from bench
# docker compose -f scripts/compose.bench.yaml down
services:
mosquitto:
image: eclipse-mosquitto:2
command: ["sh", "-c", "printf 'listener 1883\\nallow_anonymous true\\n' > /tmp/m.conf && exec mosquitto -c /tmp/m.conf"]

houses:
build: ..
volumes: ["..:/app"]
working_dir: /app
command: ["python3", "houses_runtime.py", "--mqtt", "mosquitto:1883"]
depends_on: [mosquitto]

bench:
build: ..
volumes: ["..:/app"]
working_dir: /app
command: ["python3", "scripts/bench_mqtt.py", "mosquitto:1883", "${TICKS:-600}"]
depends_on: [houses]
Loading
Loading