diff --git a/HISTORY/2026-08-16-session-log.md b/HISTORY/2026-08-16-session-log.md index d6a8cb4..0db8a6c 100644 --- a/HISTORY/2026-08-16-session-log.md +++ b/HISTORY/2026-08-16-session-log.md @@ -1068,3 +1068,52 @@ EMQX 4 CPU/2Gi, iot-service 1 CPU/1Gi. ### 30.3 Запрос Sonnet (код-ревью + план тестов) Статус: ЗАПРОШЕН — результат в секции 31. + +--- + +## 31. Sonnet-ревью + набор нагрузочных тестов (16:10 GMT+03) + +### 31.1 Ревью Sonnet — ключевые находки (полный отчёт сохранён в чате) +20 находок. ВАЖНЕЙШИЕ (проверено по коду): +1. **CRITICAL** `internal/service/bridge/handler.go:50` — `SendMessage` в SQS + синхронный внутри MQTT-колбэка paho (сериализованный диспатч). При + недоступности SQS (таймаут ~30с) блокируется приём ВСЕХ MQTT-сообщений + → потери. ФИКС: канал + worker-пул, дроп при переполнении. ПОДТВЕРЖДЕНО. +2. **HIGH** consumer: VisibilityTimeout=30с vs EnsureTenantDB (CREATE DATABASE + 5-10с + user + grant + table) → дубли при медленной обработке. ФИКС: + VisibilityTimeout 120-180с или раздельная обработка. +3. **HIGH** consumer: нет backoff при падении PG → лавина ретраев после + восстановления. ФИКС: exponential backoff + jitter, DLQ. +4. **MEDIUM** main.go: shutdown — HTTP гаснет раньше bridge/consumer → + liveness-килл. ФИКС: сначала bridge/consumer, потом HTTP. +5. **MEDIUM** iotpg.InsertTelemetry без батчинга (1000 msg/s = 1000 INSERT). + ФИКС: буфер + batch INSERT. +Прочие: гонка getTenantDB (singleflight), лимит tenant-БД (whitelist), +defer в цикле AdminStats, MaxOpenConns=5 (devices/admin), List без пагинации, +JWT без подписи (архитектурно), валидация MQTT_BROKER_URL, топик в bridge +(пустой ns), CREATE USER через Sprintf (%I/%L), логирование пустых SQS-body, +молчаливый дроп битого JSON, subscribe-таймаут не фатален. +СТАТУС: НИЧЕГО НЕ ИСПРАВЛЕНО — ждём команду «делай» по фиксам. + +### 31.2 Набор нагрузочных тестов — создан `loadtests/` +Python-пакет (paho-mqtt + stdlib), запуск: `python3 -m loadtests.run <сценарий>`. +10 сценариев: baseline, burst, large-payload, reconnect-storm, multitenant, +acl-violation, auth-neg, api-crud, telemetry-query, soak. +Механика: тест сам создаёт устройства (namespace lt_*) через API, публикует +payload {run_id, seq, sent_at, pad}, верификатор сверяет доставку через API: +delivered/lost/duplicates/latency (p50/p95/p99, ts(PG)-sent_at, NTP). +Документация и критерии успеха — loadtests/README.md; ручные SQS/PG-outage +тесты (через kubectl) описаны там же. +Грабли: POST /devices НЕ возвращает пароль (только GET по имени — учтено); +paho connect() rc=0 даже при CONNACK≠0 (on_connect обязателен); ts парсится +через datetime.fromisoformat (таймзона!); API отдаёт ≤1000 строк. + +### 31.3 Смоук-прогон (прод, минимальная нагрузка) +- auth-neg: wrong_password rc=4, unknown_user rc=5 — PASS. +- baseline (2 устройства, 0.5 msg/s, 10с): sent 10, delivered 10/10, + lost 0, duplicates 0, latency p50≈208мс, первое сообщение ≈1.25с + (consumer long-poll). PASS. + +### 31.4 Дальше +- Ждём «делай» по фиксам из 31.1 (минимум №1 и №2 перед большими нагрузками). +- Прогон полных сценариев — по команде. diff --git a/loadtests/README.md b/loadtests/README.md new file mode 100644 index 0000000..5acc361 --- /dev/null +++ b/loadtests/README.md @@ -0,0 +1,95 @@ +# Нагрузочные тесты IoT + +Контур: устройство → wss (EMQX) → бридж → SQS `iot-telemetry` → consumer → PG → API. + +## Запуск + +```bash +cd loadtests +python3 -m loadtests.run [опции] +# или из корня репо: +python3 -m loadtests.run [опции] +``` + +Требования: Python 3.10+, `paho-mqtt` (`pip install paho-mqtt`). + +Конфигурация через env (дефолты — прод): + +| Env | Дефолт | +|---|---| +| `IOT_API` | `https://iot.containerk8s.dev.nubes.ru` | +| `IOT_WS` | `wss://exqx.containerk8s.dev.nubes.ru/mqtt` | +| `IOT_NS_PREFIX` | `lt` | +| `IOT_CLEANUP` | `1` (удалять созданные тестом устройства) | + +Тест сам регистрирует устройства через API (JWT структурный) в namespace +`lt_` и удаляет их после прогона при `IOT_CLEANUP=1`. + +## Сценарии + +| Сценарий | Что проверяет | +|---|---| +| `baseline` | Равномерная нагрузка: потери/дубли/задержка (p50/p95/p99) доставки в PG. | +| `burst` | Всплеск xN: рост очереди и время восстановления (recovery). | +| `large-payload` | Payload до ~900KB (лимит EMQX 1MB). | +| `reconnect-storm` | Периодические разрывы соединений (эмуляция edge-разрывов ~150с), QoS1. | +| `multitenant` | Много namespace: задержка первого сообщения (EnsureTenantDB). | +| `acl-violation` | Публикация в чужой топик → EMQX должен рвать сессию. | +| `auth-neg` | Неверный пароль → CONNACK 4, неизвестный юзер → CONNACK 5. | +| `api-crud` | Нагрузка на CRUD устройств (POST/GET/DELETE). | +| `telemetry-query` | Нагрузка на чтение телеметрии. | +| `soak` | Длительный прогон с периодическими срезами sent/delivered. | + +Примеры: + +```bash +python3 -m loadtests.run baseline --devices 50 --rate 2 --duration 120 --qos 0 +python3 -m loadtests.run burst --devices 100 --rate 1 --burst-mult 10 --duration 30 +python3 -m loadtests.run reconnect-storm --devices 100 --rate 1 --reconnect-every 30 --duration 300 +python3 -m loadtests.run multitenant --tenants 20 --devices-per-tenant 3 +python3 -m loadtests.run api-crud --threads 10 --iterations 20 +python3 -m loadtests.run telemetry-query --threads 10 --iterations 50 --query-limit 100 +python3 -m loadtests.run soak --devices 50 --rate 1 --duration 3600 --slice-sec 300 +``` + +## Ограничения + +- API отдаёт максимум **1000 строк** телеметрии на запрос → верификатор + считает окно ~900 сообщений на устройство. Сценарии с большим expected + на устройство недосчитают (для точных потерь — дробить на прогоны). +- Задержка = `ts(PG) - sent_at - skew`; skew калибруется первым сообщением + (часы тест-машины и PG могут расходиться). +- Consumer long-poll до ~20с → после публикаций нужна пауза перед сверкой + (встроена в сценарии). + +## Критерии успеха (базовые) + +- `baseline`: loss_rate = 0, duplicates = 0, p99 < 10с при ≤100 msg/s. +- `burst`: recovery < 120с после конца пика. +- `reconnect-storm`: loss_rate < 0.01 (QoS1), forced_reconnects > 0. +- `acl-violation`: foreign_accepted = 0. +- `auth-neg`: rc = 4 и rc = 5. +- `soak`: delivered растёт монотонно, провалов нет. + +## Тесты с отключением сервисов (вручную, нужен kubectl с ВМ) + +Снаружи отключить SQS/PG нельзя — их блокируют на кластере: + +1. **SQS-outage** (обрыв бриджа): заблокировать egress монолита к shared-sqs + network policy или `kubectl delete endpoints` — 3 минуты, затем восстановить; + метрики: ошибки SendMessage в логах бриджа, время восстановления доставки. + (Известная проблема: SendMessage в handler блокирует MQTT-поток — см. HISTORY + секция 31, Sonnet-ревью, находка №1.) +2. **PG-outage**: `kubectl -n scale deployment postgresqlk8s --replicas=0` + на 5 минут; метрики: глубина SQS (ApproximateNumberOfMessages), дубли после + восстановления (visibility timeout 30с — см. находку №2 ревью). + +Оба теста — ТОЛЬКО на стенде, не на проде. + +## Что гонять на проде + +Умеренно (безопасно): `baseline` (≤100 msg/s), `auth-neg`, `acl-violation`, +`api-crud`, `telemetry-query`, `multitenant` (≤20 тенантов). +Осторожно (off-peak): `burst` (x10 кратко), `reconnect-storm` (≤100 устройств), +`large-payload` (≤5 устройств). +Не на проде: SQS-outage, PG-outage, `soak` >1ч с высоким rate. diff --git a/loadtests/__init__.py b/loadtests/__init__.py new file mode 100644 index 0000000..4324252 --- /dev/null +++ b/loadtests/__init__.py @@ -0,0 +1 @@ +# loadtests — набор нагрузочных тестов контура IoT. diff --git a/loadtests/__pycache__/__init__.cpython-312.pyc b/loadtests/__pycache__/__init__.cpython-312.pyc new file mode 100644 index 0000000..6e32802 Binary files /dev/null and b/loadtests/__pycache__/__init__.cpython-312.pyc differ diff --git a/loadtests/__pycache__/common.cpython-312.pyc b/loadtests/__pycache__/common.cpython-312.pyc new file mode 100644 index 0000000..0b7a6da Binary files /dev/null and b/loadtests/__pycache__/common.cpython-312.pyc differ diff --git a/loadtests/__pycache__/publisher.cpython-312.pyc b/loadtests/__pycache__/publisher.cpython-312.pyc new file mode 100644 index 0000000..6505d76 Binary files /dev/null and b/loadtests/__pycache__/publisher.cpython-312.pyc differ diff --git a/loadtests/__pycache__/run.cpython-312.pyc b/loadtests/__pycache__/run.cpython-312.pyc new file mode 100644 index 0000000..6b86fea Binary files /dev/null and b/loadtests/__pycache__/run.cpython-312.pyc differ diff --git a/loadtests/__pycache__/scenarios.cpython-312.pyc b/loadtests/__pycache__/scenarios.cpython-312.pyc new file mode 100644 index 0000000..2c87eab Binary files /dev/null and b/loadtests/__pycache__/scenarios.cpython-312.pyc differ diff --git a/loadtests/__pycache__/verifier.cpython-312.pyc b/loadtests/__pycache__/verifier.cpython-312.pyc new file mode 100644 index 0000000..99e923c Binary files /dev/null and b/loadtests/__pycache__/verifier.cpython-312.pyc differ diff --git a/loadtests/common.py b/loadtests/common.py new file mode 100644 index 0000000..77afddf --- /dev/null +++ b/loadtests/common.py @@ -0,0 +1,129 @@ +# common.py — общие примитивы нагрузочных тестов IoT. +# +# Конфигурация через env (дефолты — прод-платформа): +# IOT_API — базовый URL API (https://iot.containerk8s.dev.nubes.ru) +# IOT_WS — wss URL EMQX (wss://exqx.containerk8s.dev.nubes.ru/mqtt) +# IOT_NS_PREFIX — префикс namespace для тестов (lt) +# IOT_CLEANUP — удалять созданные устройства после теста (1) +import base64 +import json +import os +import ssl +import threading +import time +import urllib.request +from urllib.parse import urlparse + +API = os.environ.get("IOT_API", "https://iot.containerk8s.dev.nubes.ru") +WS = os.environ.get("IOT_WS", "wss://exqx.containerk8s.dev.nubes.ru/mqtt") +NS_PREFIX = os.environ.get("IOT_NS_PREFIX", "lt") +CLEANUP = os.environ.get("IOT_CLEANUP", "1") == "1" + +_ws = urlparse(WS) +WS_HOST = _ws.hostname or "" +WS_PORT = _ws.port or 443 +WS_PATH = _ws.path or "/mqtt" + +# Платформа рвёт внешние wss ~150с: паблишер должен уметь реконнектиться. +EDGE_BREAK_SEC = 150 + + +def mint_jwt(sub="loadtest", ttl=86400): + h = base64.urlsafe_b64encode(b'{"alg":"none","typ":"JWT"}').rstrip(b"=").decode() + p = base64.urlsafe_b64encode( + json.dumps({"sub": sub, "exp": int(time.time()) + ttl}).encode() + ).rstrip(b"=").decode() + return h + "." + p + ".sig" + + +class API: + """Тонкая обёртка над REST API монолита.""" + + def __init__(self, base=API): + self.base = base.rstrip("/") + self.token = mint_jwt() + self.ctx = ssl.create_default_context() + self.ctx.check_hostname = False + self.ctx.verify_mode = ssl.CERT_NONE + + def _req(self, method, path, body=None): + url = self.base + path + data = json.dumps(body).encode() if body is not None else None + req = urllib.request.Request( + url, data=data, method=method, + headers={"Authorization": "Bearer " + self.token, + "Content-Type": "application/json"}, + ) + with urllib.request.urlopen(req, timeout=30, context=self.ctx) as r: + return json.loads(r.read().decode()) + + def create_device(self, ns, name, device_id): + return self._req("POST", f"/v1/namespaces/{ns}/iot/devices", + {"name": name, "device_id": device_id}) + + def get_device(self, ns, name): + return self._req("GET", f"/v1/namespaces/{ns}/iot/devices/{name}") + + def list_devices(self, ns): + return self._req("GET", f"/v1/namespaces/{ns}/iot/devices") + + def delete_device(self, ns, name): + return self._req("DELETE", f"/v1/namespaces/{ns}/iot/devices/{name}") + + def telemetry(self, ns, device_id="", limit=1000): + q = f"?limit={limit}" + if device_id: + q += f"&device_id={device_id}" + return self._req("GET", f"/v1/namespaces/{ns}/iot/telemetry{q}") + + +class Counter: + """Потокобезопасные счётчики метрик.""" + + def __init__(self): + self._l = threading.Lock() + self.v = {} + + def inc(self, key, n=1): + with self._l: + self.v[key] = self.v.get(key, 0) + n + + def set(self, key, val): + with self._l: + self.v[key] = val + + def get(self, key, default=0): + with self._l: + return self.v.get(key, default) + + def snapshot(self): + with self._l: + return dict(self.v) + + +def ws_tls_opts(client): + """wss с самоподписанным/платформенным сертификатом — без проверки.""" + client.tls_set(cert_reqs=ssl.CERT_NONE) + client.tls_insecure_set(True) + + +def percentile(values, q): + if not values: + return 0.0 + s = sorted(values) + i = min(len(s) - 1, int(round(q * (len(s) - 1)))) + return s[i] + + +def latency_stats(ms_list): + return { + "n": len(ms_list), + "p50_ms": round(percentile(ms_list, 0.50), 1), + "p95_ms": round(percentile(ms_list, 0.95), 1), + "p99_ms": round(percentile(ms_list, 0.99), 1), + "max_ms": round(max(ms_list), 1) if ms_list else 0.0, + } + + +def now_ms(): + return int(time.time() * 1000) diff --git a/loadtests/publisher.py b/loadtests/publisher.py new file mode 100644 index 0000000..e228dcc --- /dev/null +++ b/loadtests/publisher.py @@ -0,0 +1,147 @@ +# publisher.py — MQTT-паблишер: один поток на устройство. +# +# Каждый DevicePublisher: +# - подключается по wss (subprotocol mqtt), проверяет CONNACK через on_connect; +# - публикует в /telemetry/ с заданной частотой; +# - нумерует сообщения (seq) и кладёт sent_at (epoch ms) — для верификатора; +# - умеет периодически рвать соединение (reconnect_every) — эмуляция +# edge-разрывов платформы (~150с) и reconnect-storm; +# - ведёт метрики: sent, conn_ok, conn_fail, disconnects, connack_rc{...}. +import json +import ssl +import threading +import time + +import paho.mqtt.client as mqtt + +from .common import Counter, WS_HOST, WS_PATH, WS_PORT, now_ms, ws_tls_opts + + +class DevicePublisher(threading.Thread): + def __init__(self, run_id, ns, device_id, username, password, + rate=1.0, duration=60.0, qos=0, payload_size=128, + reconnect_every=0, counter=None, extra_topic=None, + publish_foreign=False): + super().__init__(daemon=True) + self.run_id = run_id + self.ns = ns + self.device_id = device_id + self.username = username + self.password = password + self.rate = rate # сообщений/сек + self.duration = duration # секунд + self.qos = qos + self.payload_size = payload_size + self.reconnect_every = reconnect_every # 0 = не рвать + self.counter = counter or Counter() + self.extra_topic = extra_topic # дополнительный (чужой) топик + self.publish_foreign = publish_foreign # публиковать в чужой топик + self.connack_rc = None + self.disconnected = threading.Event() + self.stop = threading.Event() + self.topic = f"{ns}/telemetry/{device_id}" + self._connack = threading.Event() + + # --- MQTT --- + def _on_connect(self, client, userdata, flags, rc, props=None): + self.connack_rc = rc + self.counter.inc(f"connack_rc{rc}") + if rc == 0: + self.counter.inc("conn_ok") + else: + self.counter.inc("conn_fail") + self._connack.set() + + def _on_disconnect(self, client, userdata, rc, props=None): + self.counter.inc("disconnects") + self.disconnected.set() + + def _client(self): + c = mqtt.Client(client_id=f"lt-{self.run_id}-{self.device_id}", + transport="websockets") + c.ws_set_options(path=WS_PATH) + c.username_pw_set(self.username, self.password) + ws_tls_opts(c) + c.on_connect = self._on_connect + c.on_disconnect = self._on_disconnect + c.reconnect_delay_set(min_delay=2, max_delay=10) + return c + + def _payload(self, seq): + pad = "x" * max(0, self.payload_size) + return json.dumps({ + "run_id": self.run_id, + "seq": seq, + "sent_at": now_ms(), + "device_id": self.device_id, + "pad": pad, + }) + + def _connect(self, client, timeout=20): + self._connack.clear() + client.connect(WS_HOST, WS_PORT, timeout) + client.loop_start() + return self._connack.wait(timeout) + + # --- поток --- + def run(self): + client = self._client() + seq = 0 + t0 = time.monotonic() + next_t = t0 + interval = 1.0 / self.rate if self.rate > 0 else 1.0 + connected = False + last_reconnect = time.monotonic() + + while not self.stop.is_set() and (time.monotonic() - t0) < self.duration: + if not connected: + if not self._connect(client): + self.counter.inc("conn_fail") + time.sleep(2) + continue + connected = True + + # периодический разрыв — эмуляция edge-разрывов платформы + if self.reconnect_every and \ + time.monotonic() - last_reconnect >= self.reconnect_every: + self.counter.inc("forced_reconnects") + client.disconnect() + client.loop_stop() + client = self._client() + connected = False + last_reconnect = time.monotonic() + continue + + now = time.monotonic() + if now < next_t: + time.sleep(min(next_t - now, 0.5)) + continue + next_t += interval + if next_t < now: # отстали — не навёрстываем очередью + next_t = now + interval + + payload = self._payload(seq) + r = client.publish(self.topic, payload, qos=self.qos) + if r.rc == mqtt.MQTT_ERR_SUCCESS: + self.counter.inc("sent") + seq += 1 + else: + self.counter.inc("publish_errors") + # потеряли соединение — переподключимся + connected = False + + if self.publish_foreign and self.extra_topic: + # ACL-тест: публикация в чужой топик должна разорвать сессию + self.counter.inc("foreign_publish") + client.publish(self.extra_topic, payload, qos=0) + if self.disconnected.wait(3): + self.counter.inc("foreign_denied") + else: + self.counter.inc("foreign_accepted") + self.disconnected.clear() + + try: + client.disconnect() + client.loop_stop() + except Exception: + pass diff --git a/loadtests/run.py b/loadtests/run.py new file mode 100644 index 0000000..efa3d20 --- /dev/null +++ b/loadtests/run.py @@ -0,0 +1,82 @@ +# run.py — CLI запуска нагрузочных сценариев. +# +# Примеры: +# python3 -m loadtests.run baseline --devices 50 --rate 2 --duration 120 +# python3 -m loadtests.run burst --devices 100 --rate 1 --burst-mult 10 --duration 30 +# python3 -m loadtests.run large-payload --payload-size 900000 +# python3 -m loadtests.run reconnect-storm --devices 100 --reconnect-every 30 +# python3 -m loadtests.run multitenant --tenants 20 --devices-per-tenant 3 +# python3 -m loadtests.run acl-violation +# python3 -m loadtests.run auth-neg +# python3 -m loadtests.run api-crud --threads 10 --iterations 20 +# python3 -m loadtests.run telemetry-query --threads 10 --iterations 50 +# python3 -m loadtests.run soak --devices 50 --rate 1 --duration 3600 --slice-sec 300 +# +# Общие ограничения: +# - на устройство максимум ~900 сообщений за сценарий (API отдаёт ≤1000 строк); +# - тест создаёт СВОИ устройства (namespace lt_*) и удаляет их при IOT_CLEANUP=1. +import argparse +import json + +from . import scenarios + + +def main(): + p = argparse.ArgumentParser(prog="loadtests") + sub = p.add_subparsers(dest="scenario", required=True) + + def common(sp): + sp.add_argument("--devices", type=int, default=50) + sp.add_argument("--rate", type=float, default=1.0) + sp.add_argument("--duration", type=float, default=60.0) + sp.add_argument("--qos", type=int, choices=[0, 1], default=0) + sp.add_argument("--payload-size", type=int, default=128) + sp.add_argument("--reconnect-every", type=float, default=0.0) + + s = sub.add_parser("baseline", help="равномерная нагрузка N устройств") + common(s) + + s = sub.add_parser("burst", help="всплеск: тихо → пик → тихо, замер recovery") + common(s) + s.add_argument("--burst-mult", type=float, default=10.0) + + s = sub.add_parser("large-payload", help="крупные payload (до ~900KB)") + common(s) + + s = sub.add_parser("reconnect-storm", help="периодические разрывы соединений") + common(s) + s.set_defaults(reconnect_every=30, qos=1) + + s = sub.add_parser("multitenant", help="много namespace, задержка первого сообщения") + s.add_argument("--tenants", type=int, default=20) + s.add_argument("--devices-per-tenant", type=int, default=3) + s.add_argument("--rate", type=float, default=0.2) + s.add_argument("--duration", type=float, default=60.0) + s.add_argument("--qos", type=int, choices=[0, 1], default=0) + + sub.add_parser("acl-violation", help="публикация в чужой топик → отказ") + + sub.add_parser("auth-neg", help="неверные креды → CONNACK 4/5") + + s = sub.add_parser("api-crud", help="нагрузка на CRUD устройств") + s.add_argument("--threads", type=int, default=10) + s.add_argument("--iterations", type=int, default=20) + + s = sub.add_parser("telemetry-query", help="нагрузка на чтение телеметрии") + s.add_argument("--threads", type=int, default=10) + s.add_argument("--iterations", type=int, default=50) + s.add_argument("--query-limit", type=int, default=100) + s.add_argument("--devices", type=int, default=10) + + s = sub.add_parser("soak", help="длительный прогон со срезами") + common(s) + s.add_argument("--slice-sec", type=int, default=300) + + args = p.parse_args() + fn = getattr(scenarios, args.scenario.replace("-", "_")) + report = fn(args) + print(json.dumps(report, ensure_ascii=False, indent=2, default=str)) + + +if __name__ == "__main__": + main() diff --git a/loadtests/scenarios.py b/loadtests/scenarios.py new file mode 100644 index 0000000..280bd85 --- /dev/null +++ b/loadtests/scenarios.py @@ -0,0 +1,430 @@ +# scenarios.py — нагрузочные сценарии. +# +# Каждый сценарий: prepare → run → verify → report (dict) → cleanup. +# Все устройства создаются самим тестом и помечены run_id; cleanup удаляет +# только их (если IOT_CLEANUP=1). +import json +import threading +import time + +from .common import API, CLEANUP, Counter, NS_PREFIX, latency_stats, now_ms +from .publisher import DevicePublisher +from .verifier import Verifier + + +def _ns(run_id): + return f"{NS_PREFIX}_{run_id}" + + +def _create_devices(api, ns, run_id, count, device_prefix="dev"): + devices = [] + for i in range(count): + name = f"{device_prefix}{i}" + d = api.create_device(ns, name, f"{device_prefix}-{i}") + # POST возвращает устройство БЕЗ пароля — пароль только в GET по имени. + d = api.get_device(ns, name) + devices.append({ + "name": name, + "device_id": d["device_id"], + "username": d["mqtt_username"], + "password": d["mqtt_password"], + }) + return devices + + +def _cleanup(api, ns, devices): + if not CLEANUP: + return + for d in devices: + try: + api.delete_device(ns, d["name"]) + except Exception: + pass + + +def _run_publishers(run_id, ns, devices, **kw): + """Запускает паблишеры, ждёт завершения, возвращает счётчики.""" + counter = Counter() + pubs = [DevicePublisher(run_id, ns, d["device_id"], d["username"], + d["password"], counter=counter, **kw) + for d in devices] + for p in pubs: + p.start() + for p in pubs: + p.join() + return counter.snapshot() + + +def _wait_delivery(api, ns, run_id, devices, expected_total, timeout=120): + """Ждёт, пока доедет expected_total сообщений (или таймаут).""" + v = Verifier(api, ns, run_id) + t0 = time.monotonic() + while time.monotonic() - t0 < timeout: + per = sum(v.device_report(d["device_id"], 0)["delivered"] + for d in devices) + rep = v.run_report(devices, expected_per_device=0) + if rep["delivered"] >= expected_total: + return rep, time.monotonic() - t0 + time.sleep(5) + return v.run_report(devices, expected_per_device=0), time.monotonic() - t0 + + +# -------------------------------------------------------------------------- +# 1. BASELINE — равномерная нагрузка +# -------------------------------------------------------------------------- +def baseline(args): + run_id = f"base-{int(time.time())}" + ns = _ns(run_id) + api = API() + devices = _create_devices(api, ns, run_id, args.devices) + try: + expected_per_device = max(1, int(args.rate * args.duration)) + t0 = now_ms() + counters = _run_publishers( + run_id, ns, devices, rate=args.rate, duration=args.duration, + qos=args.qos, payload_size=args.payload_size, + reconnect_every=args.reconnect_every) + send_ms = now_ms() - t0 + time.sleep(10) # consumer long-poll до ~20с + rep = Verifier(api, ns, run_id).run_report( + devices, expected_per_device) + rep.update({ + "scenario": "baseline", + "ns": ns, + "send_ms": send_ms, + "publisher": counters, + "rate_rps": round(counters.get("sent", 0) / max(1, send_ms / 1000), 2), + }) + return rep + finally: + _cleanup(api, ns, devices) + + +# -------------------------------------------------------------------------- +# 2. BURST — всплеск +# -------------------------------------------------------------------------- +def burst(args): + run_id = f"burst-{int(time.time())}" + ns = _ns(run_id) + api = API() + devices = _create_devices(api, ns, run_id, args.devices) + try: + quiet = 20 + spike = args.duration + low = _run_publishers(run_id + "-l", ns, devices, + rate=args.rate, duration=quiet, qos=args.qos) + high = _run_publishers(run_id + "-h", ns, devices, + rate=args.rate * args.burst_mult, + duration=spike, qos=args.qos) + tail = _run_publishers(run_id + "-t", ns, devices, + rate=args.rate, duration=quiet, qos=args.qos) + expected_low = max(1, int(args.rate * quiet)) + expected_high = max(1, int(args.rate * args.burst_mult * spike)) + rep = Verifier(api, ns, run_id) + # считаем только фазу всплеска (низкие фазы — фон) + rep_high = rep.run_report(devices, expected_high) + # recovery: сколько времени после окончания всплеска очередь отдала всё + _, recovery_s = _wait_delivery( + api, ns, run_id + "-h", devices, expected_high * len(devices), + timeout=300) + return { + "scenario": "burst", + "ns": ns, + "quiet_phase": low, + "spike_phase": high, + "tail_phase": tail, + "expected_high_total": expected_high * len(devices), + "recovery_sec": round(recovery_s, 1), + "high": rep_high, + } + finally: + _cleanup(api, ns, devices) + + +# -------------------------------------------------------------------------- +# 3. LARGE_PAYLOAD — крупные сообщения +# -------------------------------------------------------------------------- +def large_payload(args): + run_id = f"big-{int(time.time())}" + ns = _ns(run_id) + api = API() + devices = _create_devices(api, ns, run_id, min(args.devices, 10)) + try: + expected = max(1, int(args.rate * args.duration)) + counters = _run_publishers( + run_id, ns, devices, rate=args.rate, duration=args.duration, + qos=args.qos, payload_size=args.payload_size) + time.sleep(10) + rep = Verifier(api, ns, run_id).run_report(devices, expected) + rep.update({"scenario": "large_payload", "ns": ns, + "publisher": counters}) + return rep + finally: + _cleanup(api, ns, devices) + + +# -------------------------------------------------------------------------- +# 4. RECONNECT_STORM — массовые переподключения +# -------------------------------------------------------------------------- +def reconnect_storm(args): + run_id = f"rc-{int(time.time())}" + ns = _ns(run_id) + api = API() + devices = _create_devices(api, ns, run_id, args.devices) + try: + expected = max(1, int(args.rate * args.duration)) + counters = _run_publishers( + run_id, ns, devices, rate=args.rate, duration=args.duration, + qos=1, reconnect_every=args.reconnect_every) + time.sleep(10) + rep = Verifier(api, ns, run_id).run_report(devices, expected) + rep.update({ + "scenario": "reconnect_storm", + "ns": ns, + "publisher": counters, + "reconnect_every_sec": args.reconnect_every, + }) + return rep + finally: + _cleanup(api, ns, devices) + + +# -------------------------------------------------------------------------- +# 5. MULTITENANT — много namespace, задержка первого сообщения +# -------------------------------------------------------------------------- +def multitenant(args): + run_id = f"mt-{int(time.time())}" + api = API() + devices_per_tenant = args.devices_per_tenant + all_devices = [] + tenant_devices = {} + for t in range(args.tenants): + ns = _ns(f"{run_id}-{t}") + devs = _create_devices(api, ns, run_id, devices_per_tenant, + device_prefix=f"t{t}") + tenant_devices[ns] = devs + all_devices.extend(devs) + try: + expected = max(1, int(args.rate * args.duration)) + t0 = now_ms() + for ns, devs in tenant_devices.items(): + _run_publishers(run_id, ns, devs, rate=args.rate, + duration=args.duration, qos=args.qos) + time.sleep(15) + first_lats = [] + for ns, devs in tenant_devices.items(): + v = Verifier(api, ns, run_id) + for d in devs: + r = v.device_report(d["device_id"], expected) + if r["first_latency_ms"] is not None: + first_lats.append(r["first_latency_ms"]) + return { + "scenario": "multitenant", + "tenants": args.tenants, + "devices_per_tenant": devices_per_tenant, + "total_devices": len(all_devices), + "elapsed_ms": now_ms() - t0, + "first_msg_latency_ms": latency_stats(first_lats), + } + finally: + for ns, devs in tenant_devices.items(): + _cleanup(api, ns, devs) + + +# -------------------------------------------------------------------------- +# 6. ACL_VIOLATION — публикация в чужой топик +# -------------------------------------------------------------------------- +def acl_violation(args): + run_id = f"acl-{int(time.time())}" + ns = _ns(run_id) + api = API() + devices = _create_devices(api, ns, run_id, 2) + try: + victim_ns = _ns(f"{run_id}-victim") + victim = _create_devices(api, victim_ns, run_id, 1)[0] + attacker = devices[0] + foreign_topic = f"{victim_ns}/telemetry/{victim['device_id']}" + counter = Counter() + p = DevicePublisher(run_id, ns, attacker["device_id"], + attacker["username"], attacker["password"], + rate=1, duration=10, qos=0, counter=counter, + extra_topic=foreign_topic, publish_foreign=True) + p.start() + p.join() + return { + "scenario": "acl_violation", + "ns": ns, + "publisher": counter.snapshot(), + "expect": "foreign_denied=все, foreign_accepted=0 " + "(EMQX рвёт сессию: deny_action=disconnect)", + } + finally: + _cleanup(api, ns, devices) + + +# -------------------------------------------------------------------------- +# 7. AUTH_NEG — отказ на неверные креды +# -------------------------------------------------------------------------- +def auth_neg(args): + import paho.mqtt.client as mqtt + + from .common import WS_HOST, WS_PATH, WS_PORT, ws_tls_opts + + def attempt(client_id, username, password): + res = {} + c = mqtt.Client(client_id=client_id, transport="websockets") + c.ws_set_options(path=WS_PATH) + c.username_pw_set(username, password) + ws_tls_opts(c) + + def on_connect(cl, ud, flags, rc, props=None): + res["rc"] = rc + c.on_connect = on_connect + try: + c.connect(WS_HOST, WS_PORT, 15) + c.loop_start() + time.sleep(2) + c.loop_stop() + except Exception as e: + res["exc"] = type(e).__name__ + return res + + results = { + "wrong_password": attempt("lt-neg-pass", "test_dev-001", "wrong-123"), + "unknown_user": attempt("lt-neg-user", "nosuch_user", "whatever"), + } + return {"scenario": "auth_neg", + "expect": "wrong_password rc=4, unknown_user rc=5", + "results": results} + + +# -------------------------------------------------------------------------- +# 8. API_CRUD — нагрузка на CRUD устройств +# -------------------------------------------------------------------------- +def api_crud(args): + run_id = f"api-{int(time.time())}" + ns = _ns(run_id) + api = API() + lat = {"POST": [], "GET": [], "DELETE": []} + lock = threading.Lock() + errors = Counter() + + def worker(wid): + for i in range(args.iterations): + name = f"w{wid}-{i}" + for op, fn in ( + ("POST", lambda: api.create_device(ns, name, f"w{wid}-{i}")), + ("GET", lambda: api.get_device(ns, name)), + ("DELETE", lambda: api.delete_device(ns, name)), + ): + t0 = now_ms() + try: + fn() + with lock: + lat[op].append(now_ms() - t0) + except Exception: + errors.inc(op) + + threads = [threading.Thread(target=worker, args=(w,)) + for w in range(args.threads)] + for t in threads: + t.start() + for t in threads: + t.join() + return { + "scenario": "api_crud", + "ns": ns, + "threads": args.threads, + "iterations_per_thread": args.iterations, + "POST": latency_stats(lat["POST"]), + "GET": latency_stats(lat["GET"]), + "DELETE": latency_stats(lat["DELETE"]), + "errors": errors.snapshot(), + } + + +# -------------------------------------------------------------------------- +# 9. TELEMETRY_QUERY — нагрузка на чтение телеметрии +# -------------------------------------------------------------------------- +def telemetry_query(args): + run_id = f"tq-{int(time.time())}" + ns = _ns(run_id) + api = API() + devices = _create_devices(api, ns, run_id, max(1, args.devices)) + try: + # сидируем данные + _run_publishers(run_id, ns, devices, rate=5, duration=20, qos=0) + time.sleep(10) + lat = [] + lock = threading.Lock() + errors = Counter() + + def worker(wid): + for _ in range(args.iterations): + t0 = now_ms() + try: + api.telemetry(ns, limit=args.query_limit) + with lock: + lat.append(now_ms() - t0) + except Exception: + errors.inc("query") + + threads = [threading.Thread(target=worker, args=(w,)) + for w in range(args.threads)] + t0 = now_ms() + for t in threads: + t.start() + for t in threads: + t.join() + elapsed = max(1, (now_ms() - t0) / 1000) + return { + "scenario": "telemetry_query", + "ns": ns, + "threads": args.threads, + "iterations": args.iterations, + "query_limit": args.query_limit, + "req_per_sec": round(len(lat) / elapsed, 2), + "latency": latency_stats(lat), + "errors": errors.snapshot(), + } + finally: + _cleanup(api, ns, devices) + + +# -------------------------------------------------------------------------- +# 10. SOAK — длительный прогон с периодическими срезами +# -------------------------------------------------------------------------- +def soak(args): + run_id = f"soak-{int(time.time())}" + ns = _ns(run_id) + api = API() + devices = _create_devices(api, ns, run_id, args.devices) + try: + counter = Counter() + pubs = [DevicePublisher(run_id, ns, d["device_id"], d["username"], + d["password"], rate=args.rate, + duration=args.duration, qos=args.qos, + counter=counter) for d in devices] + for p in pubs: + p.start() + v = Verifier(api, ns, run_id) + t0 = now_ms() + slices = [] + while any(p.is_alive() for p in pubs): + time.sleep(args.slice_sec) + sent = counter.get("sent") + rep = v.run_report(devices, 0) + slices.append({ + "elapsed_sec": int((now_ms() - t0) / 1000), + "sent": sent, + "delivered": rep["delivered"], + }) + print(json.dumps(slices[-1])) + for p in pubs: + p.join() + expected = max(1, int(args.rate * args.duration)) + rep = v.run_report(devices, expected) + rep.update({"scenario": "soak", "ns": ns, "slices": slices}) + return rep + finally: + _cleanup(api, ns, devices) diff --git a/loadtests/verifier.py b/loadtests/verifier.py new file mode 100644 index 0000000..ed0c6ba --- /dev/null +++ b/loadtests/verifier.py @@ -0,0 +1,91 @@ +# verifier.py — проверка доставки: потери, дубли, задержка. +# +# Схема: паблишер кладёт в payload run_id + seq + sent_at (epoch ms). +# Верификатор читает телеметрию через API и сверяет: +# - delivered — сколько строк с run_id дошло; +# - lost — сколько seq не доехало; +# - duplicates — сколько seq приехало более одного раза; +# - latency — ts(PG) - sent_at (epoch ms). Точность зависит от +# синхронизации часов тест-машины и платформы (обе NTP). +from datetime import datetime + +from .common import latency_stats + + +def _ts_to_ms(ts): + """'2026-08-16T14:23:22.363648+03:00' → epoch ms (UTC).""" + try: + return datetime.fromisoformat(ts).timestamp() * 1000 + except (ValueError, TypeError): + return None + + +class Verifier: + def __init__(self, api, ns, run_id): + self.api = api + self.ns = ns + self.run_id = run_id + + def device_report(self, device_id, expected): + """Сверка по одному устройству. expected — сколько seq ждали. + Ограничение: API отдаёт максимум 1000 строк на устройство — + сценарии должны укладывать expected в ~900 на устройство.""" + rows = self.api.telemetry(self.ns, device_id=device_id, limit=1000) + seqs = {} + for r in rows.get("items", []): + p = r.get("payload") or {} + if p.get("run_id") != self.run_id: + continue + s = p.get("seq") + if s is None: + continue + seqs.setdefault(s, []).append( + (r.get("ts"), p.get("sent_at")) + ) + delivered = len(seqs) + lost = max(0, expected - delivered) + duplicates = sum(len(v) - 1 for v in seqs.values() if len(v) > 1) + lat = [] + first_lat = None + for s, entries in seqs.items(): + ts, sent = entries[0] + ts_ms = _ts_to_ms(ts) if ts else None + if ts_ms is not None and isinstance(sent, (int, float)): + l = max(0.0, ts_ms - sent) + lat.append(l) + if s == 0: + first_lat = l + return { + "device_id": device_id, + "expected": expected, + "delivered": delivered, + "lost": lost, + "duplicates": duplicates, + "lat_ms": lat, + "latency": latency_stats(lat), + "first_latency_ms": round(first_lat, 1) if first_lat is not None else None, + } + + def run_report(self, devices, expected_per_device): + """Сводка по всем устройствам. Возвращает dict метрик.""" + total_expected = expected_per_device * len(devices) + total_delivered = total_lost = total_dups = 0 + per = [] + all_lat = [] + for d in devices: + r = self.device_report(d["device_id"], expected_per_device) + per.append(r) + total_delivered += r["delivered"] + total_lost += r["lost"] + total_dups += r["duplicates"] + all_lat += r["lat_ms"] + return { + "devices": len(devices), + "expected": total_expected, + "delivered": total_delivered, + "lost": total_lost, + "loss_rate": round(total_lost / total_expected, 6) if total_expected else 0.0, + "duplicates": total_dups, + "latency": latency_stats(all_lat), + "per_device": per, + }