From 090cedca67df962647b8314f5515fc0f6c0c876c Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E2=80=9CNaeel=E2=80=9D?= Date: Sun, 16 Aug 2026 16:59:29 +0400 Subject: [PATCH] test(loadtests): 10 load scenarios + verifier; docs: Sonnet review findings --- HISTORY/2026-08-16-session-log.md | 49 ++ loadtests/README.md | 95 ++++ loadtests/__init__.py | 1 + .../__pycache__/__init__.cpython-312.pyc | Bin 0 -> 118 bytes loadtests/__pycache__/common.cpython-312.pyc | Bin 0 -> 7981 bytes .../__pycache__/publisher.cpython-312.pyc | Bin 0 -> 7400 bytes loadtests/__pycache__/run.cpython-312.pyc | Bin 0 -> 4293 bytes .../__pycache__/scenarios.cpython-312.pyc | Bin 0 -> 21240 bytes .../__pycache__/verifier.cpython-312.pyc | Bin 0 -> 4278 bytes loadtests/common.py | 129 ++++++ loadtests/publisher.py | 147 ++++++ loadtests/run.py | 82 ++++ loadtests/scenarios.py | 430 ++++++++++++++++++ loadtests/verifier.py | 91 ++++ 14 files changed, 1024 insertions(+) create mode 100644 loadtests/README.md create mode 100644 loadtests/__init__.py create mode 100644 loadtests/__pycache__/__init__.cpython-312.pyc create mode 100644 loadtests/__pycache__/common.cpython-312.pyc create mode 100644 loadtests/__pycache__/publisher.cpython-312.pyc create mode 100644 loadtests/__pycache__/run.cpython-312.pyc create mode 100644 loadtests/__pycache__/scenarios.cpython-312.pyc create mode 100644 loadtests/__pycache__/verifier.cpython-312.pyc create mode 100644 loadtests/common.py create mode 100644 loadtests/publisher.py create mode 100644 loadtests/run.py create mode 100644 loadtests/scenarios.py create mode 100644 loadtests/verifier.py 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 0000000000000000000000000000000000000000..6e32802a1ceaf7a7469f835e3558aafb24e290cb GIT binary patch literal 118 zcmX@j%ge<81U8!+vx0&2V-N=&d}aZPOlPQM&}8&m$xy@u2KczG$)vkyYsEQGYi$RQ!%#4hTMa)1J077~gMF0Q* literal 0 HcmV?d00001 diff --git a/loadtests/__pycache__/common.cpython-312.pyc b/loadtests/__pycache__/common.cpython-312.pyc new file mode 100644 index 0000000000000000000000000000000000000000..0b7a6da562206bab29d6512af726d0ddb38731d7 GIT binary patch literal 7981 zcmcIJYj9IndgtmsbR|pj8wdd-8*mUYHV&_az_P}ecLc%S64V5Z(7nbIvSiM^2CNkw zIm;|oNoq6MDW)OCZFf2_Fk5H4nWYIKwEfeW&UBh~S# zUddh)hIV>D-+6uCx#xWEbNr8Px068m^k>7-J5_}I8+Obh6dTk3$q{myC`936Bv7ci z0LS{afQ{9BfM>N35Lhh+L{>`yJJhzA-0uiDIHKC3b>{M=I=ZmbHd_j>2sdqoUcmJN zbfO|zbnX|Rvn#Si=Xn7-hvGEp6z9oRgaj&WM0Kd{o4_?Jfhx5s+Ei$l!?T=b%B zYl|fG)^H-Ohaz#6p8Ss1s;H-1oA|rV8$FR$AQAHSg8~}OvO)E@c z^`Za>iI!tAR07Iu3MeWThsr9&X3g5mhKtRc#rmA>+rOdXj#!B@Gc$XIIz@E1IigpZ zt*FOrm1^^Q=AKeq5+)V?YH63%C2X@bm+Jp;PCjK*1pH=-XI~Z(b`ziE*A_$T2Th^a ziKZ=0@kCs0YH8Ak2B5R=NKezbX_Ub>pUvR4!Q%$6o*kIsT9ZPn7CGT_7}D{Orfyts zRD+&rpVo_l3muC_#5^=dlw9KZl zH#4t~FUVCa%L>b;VJ{?%z$~T;Pw#+{%S0vUfx-K70bUh)ImHInFDgz&fLc;Suw}a< z0hARxK!+j&bSf2!1LiA+c%RGQ!JY@V!D`t*&)&^`p8YgCnH%*&^W)s5TsrqI4tQaF zGW&Dy;jZo;FW`H#cNp&7>}T2AfHK)S%L38*m<2v$n@n+0%LZYfSW1?cB0+0jx2(|A zTbRYQSowBKDx;803dNQpQJ2{()j`Xib=e$^Lg<@YolR;O&}D9#Bfet5ZV2{T&0 zP7yoBo6aq238=J?Dd0Hh{5i2qaPzp;sO3~k^eG`FrZ|OzUC!fPnM<>wJG~{`Ebi@` z^qjBE<8V`8ckqRp(s#|)E$L?Q^c%=Idl@I@$%az)`SKx?1JRevHA}TQXL2x%IT`#j zqdf<%dO2w)8g~XVj-!Nwe?#Off0Ufze0G1*zC#UBm3os^p@D%|BplKsiFg}2Tu?lB z$DrPqpphgSO}aW!f8+X^o}mHN#~XqY(nE&SuL6={*!$FwqEhVW00j z_C7b$VOkIKbO8c9k0C(-h((SYcB-Bl1W#?)4;LDAB~II6Ix(QeDWq4#5YeA&h5&Y@ zP-LJ_rY%s>cAOOkLV6zLJ4UTH4wHOqz6f8&x!cW^3Fnzv4^-1LO!>NkVOuUt7zgi7ZK|63;ks$-dB_ub2L z?&TAmlQp+{zIf|z-paa{=iR#=6VB284~-9M7hc|SY0G%%TJM$K5Bshqt|ameZTZ@@ z-|#U14I({7Fr%Tf`@040oz?u`SBssC?X%1kjf$D;Mrbf|wLC8j6pI33B*lFQnZ-vC z@l%k3LxxdQT!{9f-w+NR?CzPVY&+fFh7L;`2!&OxZN`P^U>jtB5$K|~DtNQPKro`* z;AjVY1EMqv%CN;Xx(gZ995S(kW9}<#`5&OQ*zp>4p3jby?L+dBbsqs%2!4a# z-PJ>P0{M)){Pb;f4~*U53o5x6+U%$re%(?!=LuDZh`>(h!4sDH-vd0&3!rqo`iZv5 zrpda=ckZpf*PQk4%~$Tr%KHk3GSdPeAQ=0Yj9a1c{6l@jn+}2t-KGF%a+7bZ2dx*EQ-W~bM&aB+|{~{G>XV-zQp02_<-Z0NW z{kWl7B5*5eOl5Joq-V2+P(kMiV$no!^6=!&+wObHz2D0&Kaj6{BP+jAIEMh}X0f#E z5RAahnM?7YKQ@SiHY%qCFf!C9TM{A_jFEYhBrrZ^iegHbqk9r~?Cr5ge?;GwtZ24u zY+FfyNe|OV9Zm?2nGO|XDVAi(3=bSBE@E(}G|as~IZdDFuqcym6s^@ob=JNI0KH@V zm^yYO>smU(eN|C6<#La`K0Y*lGF!bm=UP3|VcMptyvOc!+*_CR?#@^4$;x{e&u(yj zP%sSPkNyK>11Tw_Gw3;jLBkou%x5r$ZC5aOYA_Tl&e(%NB@qq=X(JHQH3-nR(0vFF zB6t%)4}!M<+#u{@%)1n5knk1wX&4&Lldpu`50^m3EsATVJr1$`=>nVB!9A5ZaX4W{`nlgqg1EQfOV~PLptg!>htmazQl6eeYEIx;L{@l~q6^TI6bGc}oSC8z*?msK z-kBJL&`Jll!=l-1*}ITl--R3;^6~5^*mV!-_oqPhtfX@T;V>PwwhySqGT#jt(?Zh^7Sh%EX;T~qp1FzD(6bXO zo1fwqAQGrh(Ov)s4}RWP2n_~*QXMk5_-wBPH`{Qs!ekV5jsTg48UYsPt{y#-Ir0QT zTKm-EhHK4Nny;9b#1_mUV1RkGS!`*Uf*cYa99SM1p4gPHeHGQC1$ong$zw~AH%r#> ze=Tmc&+;CYZVWt!H>Fz*=!gKwbISkbc2dxjf(5}~ky-x+3ST`-(Aoc)gz*Y+a zZ-EpEX2}G-IrHe~RGyXS>6V?&CyiUy^Cxqbn6b<$r6*9hE%H8L@H8O6b1sA)_z$l_ z_!9I5{C*DbHXHez1*bc0Fz^|~s&{keA&R9jGQl$v#;rR@`3ZYmx($;Uxd08;xa0lT z^pCxTNaxLd0YrDRKSNB6ZLRPW?Fg$>$E%vwJp#4^PF2$4cp`3w#7B7Zk;o zghWWUN^#G|PeS5mu3**0{hec1IJ{m%de8fR!#)Fmq8xpqA!!M!tBS$HOQ^xe z;C79n*bphCAqucq0xt}Aaf%Vh;537a%(pI$dl%1Jk+^Jr~ou4VExd zgf1*w*OFyveGh-n^k$Z)SCW=q%K5&mIOe`Yr+#4G1}O7VN>^uLbd*R&EgP7(7qL z!$Uz$4}nSi5ul@nXqdI6$-ld+#ui?D`$5gZtIml{*LUU@wNJL^Yc@Q8;H#E(KVEs2 zAK(3J@?#`Em+$=7jaOHLZ z+=Wyu0CRBJKozHp5X|<>h7?bn0nX7em`lS?LnO#cD?ISz18)o$e=y#Yt8C0yG-ZV* zwqtnakj=j=MAC0^#MeNpfr!0r8McI0@J(SH!o>_v;gPEqoEW}wQGAIpT*X25XEut5 zHtcHrbc812g)q{HL$)K`hNJYiF}wLUF+O-kGi2-e*nXrtxCb6+;1LXL+raLjgNJ)4 zX0C-^N6#L^)3tMVS8&JSu8#e|?ygSuWX*ySbB9b{m=-Wz%;D3FytQyG`!1b;MbSL9 ze*y5A=QwVWbaiI+Xltf*eCc&Q?`)am zbIx^X`>!R>2>oFUmU6mBTQV)bbS@chO53rwDzoYj+s49H=a!d7lsfPuKEX(^HqoRXwe;bJBDj$1P_9{~rJow1ofw literal 0 HcmV?d00001 diff --git a/loadtests/__pycache__/publisher.cpython-312.pyc b/loadtests/__pycache__/publisher.cpython-312.pyc new file mode 100644 index 0000000000000000000000000000000000000000..6505d76857950fa4d2319a63fd8efe28e8096fb9 GIT binary patch literal 7400 zcmbsuTWlOx_0B$L-)ryIk9h6aiS1;Q*lE&+6ew+jn@}~46Jk?T7dmWrCdsb%l{@3u zUacWWk-Dy0Wer8P)uOni2uF?7zz2k{0zZ5pkT>3Fr&E7e0;$Shs9RL!i*sgnX4dYu z#Nt}MbI-Z&^S)>Q=JUA;ln;LQ)8twUA^$``FZNnv<1#en2~T(?MaHx*GsdW}Wz3>} ztz%a8YfIbX7KWr9V|Iq{7T%h2rr9wz?HY5X-DB>wXUvoKj(O9*F&{%L;H}3!(%R2#b70DDc<{lwYj}mxVxRi(sXC_`wCBzE?jr;3% z$IbD7<0f3`7tB17jLunqfxHgn=b0Op_wd3{7*b#ggtKT6Aq5LS`UObA3Xsj{vjL>D zj}+_xIgCCBKu!a40>l~+3y{lzTmY%!ftZ3DAdk`KnPNmYL!@>c^Iok1qkKB9;04II zbHN8tK=0!TKtUancI&NVuii4hb!&%o>;umT>oW?LHwkYoSl=YSL~1>x;~#h>P!O$5 zlc;TYG|DIrJ|?8InY`zqB&3A2AkpcA(N@JSPUJu|S&mC&5)#KL4my$H61-x`h>D9> zg#tl#LKJ8wmKGE?7Zb(FEaep&jY$CW6Yx!BGm7^JGQRE(N!m1f7WDP}XQzR`#YSC7# z+OR@Fp*yhZzzWq}adBKc1sCTy5#^$LQ@MP<1+Od#q9h)y-9n&)xoL`H^KDynYfvXP zZZ3WR)#sPVrz~;#XPsrXZOOL6wpG|Zxg}Dxz2=dd!|-v+?YrRPm)pAG=VuUGEqM{JwXyQ2%xfwf5AwbP5IS0e-J z+_iA${OP&VtKmLv{T(G~;hRg-OXDjK9R`eht*z%;U@@@T25Y|V+kh*OxS_WIAe&W- zaUUA=3W#FUp209~o$Xe2^r3x0WQ1 z*iv|A6Hidiwv?pdG0culJh2Qb9@BP+QJbdg`(Zfk#_9l6z!;=H9l#2yL?6cL5L6;w zM5VSyY(Em)Trl6&h~8%ZQCtGK1SH=?>&^R zRUg+cBLc0&LhrSK#ewDU`<>T8500*}q0%w%cC8)N);;CcJw*mQ*N(30j)C%y0RX^E z!9-M%W#;DD6&B_^CI^B=R#TUn=)&O$!H`i^mv^86Zv)0u7*z%8$KrS>x(MV!(`PJ_ zVYnKp3y8K6lxnlFE-;cuhjiluBfUD#{1)Kl@s^=LAaJ%}g_%+v2G}6f5A+5X6{iTy z;9`;*z$n&SjE_1MD~OXCvnaNt2oac_pGfCKO_cCd>foqVGi4_caSAFC$pHpTV>HXv zBDMx`O0DMvPI?9^Fj@YlH-;9Ng-5F$2g)4>R>K2j|3J~Y76@I<&*tTp9rNRJ<8r8F zt+n&#Q*wLfwXVf30RQcAbhtN&Bj~wD94?Q>rF!#@G*-Fv7iet9A?Vb_M@OR&LPG6P zi)8SrVA0=e>MVm(`bdy~S`+#kyC%bwm~&v&nkgP3*O}p{Q}gBacuI)D;}3m4xMPYp zrRbv+cPg9BaiSEXl46^TB_y3w1K8`75@{hjA#LSWZJh?(s`Z?K+2RmXnuAdtaYvc$ zsIZ+f>%Z!s^}jLsu5-Dy!X5%xQ@;w^Ap;9%W@oB_o^qgPY0qk4-wM0$E4+^v?w~Tt z@ISC+t?1;NfhQ>_kiqCip2!#r@0}Tl6B{WR&x}6}R118W&&mS}GYoH;vGdj$#{=#M z@5x50n}oHt7X6GRi}=xltzNL^!TL8=`Ef8gB@6w1!67~kBaC;ZxvS~##6vKmftO4g zzR{e;*Tt)Kr)=OYdBfl^8mY=@yxR3WHz_BR#3$O8qMu4t-mjLG9*n$nYOnd8QG~k0| zL?Q4(Vv&Zt5r3FefSsR*etRB_4rmCBRM)n{F7`a!UDcYMTZjOlX^>`IU>tiIj3CKf z4X}_+X_VWz;nL0KjJx1Y_J47vv{%O^_vtP3tG?TMwjl{_lLp`>J?i}}na*vj;5Ao3 z%K$SmxZhIfZvl@n{+2?2y^WoG=1yNq>jrxiT#z|4=&DNW!&{C4VDgbK;&{Klqxu$p zP#Y1@XE7vmWE=Mjn)6^*qcdW_xutBgTBhviv4xE^HsLze^^>-DOq0oy$i#G z|I3z4lifu6ntt|PLahvW0p9qNWXk#inPj3(!)ny_u{Px1lP4(Z0DT;)QKbozT5*Br zi~w9buZ>%OGRmmAC(Kl3B#ORL@$0EDfzm7$l|U^e#?|n5wNW~j7cvRRfdzFq7LN-# zN#LUqicySk?#!+9UmD{lWJOets6&u zAFELu>N-A!hlhtx)Km8aKZ$)7>0!u}L(Y(cbVveo>X%|EEt>+76;QwrDLTodv0{a^ zI-E?)NGc2!J+KT!zIxFjU@N3N^P#QJ96;i!#P`9DLnPJo(0pVr^7F2uL-q$C;Po|K zeQEZkRbOY(A~%Hq+rRHOiMJCsd)`Tx9$5%qYhP^t*^{?IMOF^C7hP*wB)r-lE&A7+ zTjry4(Q5Oqa`Ub`&Am5z-ygg_xVrBfx6YRLJ^80c|N5;ze`_^*=ED$V?mFgY=4O_j zue2VJ!)^1=%{^BQ_m;!GOD|T!ht>e|?A)^p7vDG~ha+;NyX3m(Ai;1ou={pkH}($C z4X?uBAD6>NZiTB~KUMzvsoUXGaxgsao%2?M-Q{5SO3$N};1N04QVsT$gFThtZnceF zmEbN6;3JEXH@iv>Inq;dt@ZT%rt9sl)t*Bo|A)c0f5A-#=LT1scNK@^Ky&H&R|{*f zxO>k1)*}m-maU70<&o8nLzUoR7%Y36OBc%CuHp$f@X&(w)q>m}S(sXcrPtxrke*8~ z4Xw89Eq({~pZCxC7o=EawaXIejh-rksw{)L6O8$P-HEgy1!@@uvW$BC1Q8PTlPpRhW?qxu!t>NJlbOeAX za~DGu_)!Rc6~L4~{3_t)xOA4ENMYN{aTh0IDXmA%aI47&HQ}Jz8`U6C)tF4cLW-}J zJC6&g6vu(N>jSKaf!R~40#ejUwZeVHgHU~bg{(huOn&0qa!>zf&KCQ=jl)iRWTV+) z@BV~!*_%FTVeOucu+6^zKmNza26M3s*+I;=3~Jd6 z)h?;oJxH)-(_k9apC+`|89o`+CmKbcMKQmkI7ZbUAkE0J*GC= zmyWb^&iv>5ujfDC%>J#RAW%ZMg+`A17e2U`&zU?q3X?B6!Vx~rov}y$3{QMS zNc%Fv8G+}Bc!4`3Dt_Ysu#@AcRL|a$Njm8zgNcfU(#z9PQ<6H`Kp1z1c#FBr;(VpN*-~u<{8|PEO zr;%VUNTgC|87ug>ao@P$df>Trw8EtBnOZ_Bw1%&G&8j@EkcLv@3}y}2x?Y+{^E17y zozVvJ+%vU=RDnn3Ti@Fr_NDN zoI7gBYo`aldSBM=q+LY*KS!~j}`bbdPO zGerVxrbLwyJxLXkSSE2W9e^Ju{bUD7zeLjZt~$O6uPxI(@s$J*Ij)TtHRnJAnfhzM z2R!&-4S2~DPl$wp1VW`cw>)sX*N&>*qt#G@mqt%iO^n9OPj{G6nT+7TZd8hukJ(Ljn;DV8Naf?o`N~ItSYWqs#%A%r3@hop=|8JeRnt#^ zJvLi8&+f86gRS%Ic25Lw<;o{)A@a{pKaapdnccRgu6$Pn92mCm2<``NM4?`mEt7Hv zTEW?t_M;SEzJ-P~Y>VPqsY?1INEVw0G~@$D)O%mEBgg z9L?xQQ1h;pbQzMGgFr9ZVMCPj5vwvc!SEHx=}H9BSzz;sg4)c(ugq>$k$J4X1EN^5 zAG5m26reNBI?N!c!RKLC)89mo-+}G|y8)49;~7QG#}kTEggGmgZ>@Sefz3Kl%liWR zrgGW#U70Mv39`+qRixWO>+*?`^9ps=s(2_U()Q3FzsBpYee^jLA#G9Av(L`+Fe@~~ z6KQ!gnM+%T%xVvI7dqe@*iP6tXpB#xM$jCyk;+Aoe+$9g!W~}lY%E&<0btq9f)v-a zVL7V|={v!^Rkm3xMSGiY*RY~cTx^bkUkncR54;_DeNQ)N=96-Q>O_}q&DG5ari{qI zsb3#oqbaSr3{xde)GB_IVlP4cZQ@W1O`4KcGKxmW* zbt;~MGAmR3+$=Nz1)(K5oy;V)Dgha-hnp#?x%h{AH`<@W3fd+rpIKdc3l4*ImY@;X z9jm+GW);Lf59LCgEURfCFWeB^yewkrI=E4S<_GSOkOnwRin|M$>(9FpQZTj zHhI+)En~iOSu+6IIjX5aSc%vdMY=RaBvN_|rWl<19wcAn9!t#+q^%Wc>(vkLOI^kO zhvB*h;kHV+?P}OK&BAT>!ac=-pTvOwB){0qh3Y38FE>sVSm(J*jrXLXV*evCe5vcg zTR%#5*2?c#M;cc$%PS8>|HZ-cgHuOJenTq`-WPja>>rr?E@DqFW1SVYv6Jn|;ZlBi z$uVG1Ts#w!Twb|kDNIwzEWg`z$eba1-EXz_xkdZ!iR zG)j70%z^2aYa_EAwG(y5DxACXIe*P_fk@B2af2r_yVe=zV`sC!Iz=>^^ GvHt>AA$?T< literal 0 HcmV?d00001 diff --git a/loadtests/__pycache__/scenarios.cpython-312.pyc b/loadtests/__pycache__/scenarios.cpython-312.pyc new file mode 100644 index 0000000000000000000000000000000000000000..2c87eabb225037485b10214b3936226d552747e3 GIT binary patch literal 21240 zcmd6P3vd+Ixn}ot_e{@Q(u|}L5_%v38bG{^ZEOQJLSV3s{J@*oBwR(P2N+o|caJPX zMvn5RVg}3N8evF}al?osLQ)90K%ea`7~{{Q>`^Pl-IZnuL&xc-B$_txLbasQKU%wWkz zuD*lF1gCHc-^cAWuKZq}-COor6pJGC3H{<;v0vIN^;`E^`)zw|{qkP9-@ez*bI2F_ z9R1F{P8OH?T>b97Zp5wM;`Vx!O2zgqey>-NaaX3;LCcjY#euX6#fdwg;sW(6ZqQ1_ z16rkcL93NA(1217TBB5eE>e7;wTd6K?l2#${+K$@`7sX~wCcjvM|SD9Z4YhV+WD^@ zL1b`vAQDn_d*`06M;_b0W7iY9yRSPE8t6IPrA4|Unr%Ug1qkS2t5<-2|Y5r zzb~vE3aLnvzYtQx2f`t>hpIP&fC&9x4S-B=5i^~81%np_eO1`Q1udOA-xcIl8?sTA zRUIDa3M(2_)Oe$+Z?Id5gtUmZM(YU;bgSV(ZS~M$UFaIn)G`$D;$QQDjB&G`>e!bi z`N?IIfuwUuQe1Ks$y|>_G#g}GS&M#!$OLzg8%uXegk$S1u^8#v=6-y0LaTT&x zk*ZMt1yp9N#5yZfJxM7oBxkUG%aAIjbib!*zlpsEB# zU8H2yj~aD3S4u7@AgSF02k~H!8p0&#GHTWx2SX7vZtWf#3JoYhyDoGOg>?%)hllEu zdT1u;?rdX*)X;(Ov%1Z=)pS0rR#P!v(d@W3nn9)SsC&DLUZ|-{kxCy|i-U}D8Bf{R z{j=_h@!OyO@)gb@Ib(dLqG7xtQJxTz-jrZ!KlnJJ%HH`APK z-I8|Q8{;pzyeB&2os;V)BWYI)9+tf)JaJFPAIQ`!n!M*!eI`&dxh54@VO}rWrGVof zImzyvbK-S#Rh-K=)~S}GH9dv1!HNVnuX+(Foaq{Wn2TC4pRfc*1v|IkED{u+zOXqr z@O2l7++nVni-<^Op@mZgo_mrzzI=pxR(O&d;jwnua0UgH<{FK9Unu{*?AxF12@OTU zg9AZPtw#a1fyfdfEg%|IO<%wmeYvZrFVsCSJfsGZ-b@8LK*qS=IV&=G+sl9JDw`8i+`4=;c=Ao4BViD$*F9ka{c!xgN{D#YmANx$%#3SXGiSez+PdM>1B< z0Z9>~VwjH#uk+vLPg;&Qh*;O7Qq+oXQh|DGQCq2XYQ7oY!H-o|Sh~cDomX1yrMW#X zm}T~$2R!~qr4cSD4fml~aw0jMd_VQg)Yp^$0kI#YzR{R`ACb3{f18|1#gH(TdVW=7 z^4}+am^w+B?;`UZmW<5ge?*}Z$bKjJ6Xd*;8cR;6;@NsqUr!ySXHrL#GpKxKHCy#6 zF=UnIn{H$46--=GX%ed>0Cfp0UybNu?_hX9m$iZJA??s$Bq*~bo}Waz<*5;MB|Sy6 zNf(EP_hW?*X~bn|T2U|;HgR?28xliP*P}=q{mCuLA?w_-ARg3>9SzPs( z%F3@voV#vvN1}a-pIUj=GsDlUc+Z`wTRho&YWHPv)6SlbR1v|=16_?!Ql#ZE9>DIXQ1)0*zFF7kNI%`wT+O%_VQe6C} zD_+#5(aJipWa~PyhX%+Dl*(58PcU1L+6!kbn~xS#5*)MVXKJr#KAyjt^J3Kc*~&$& zNAQiJ(o>Kk>a&sDY>e3t@CyGeIVyB;W}T+27`4UirRHQ*P%Li=CeML9k`*RrqT*9u zrU??2KTn$uSoUxQF(3Etj-kyNai$J$JhVB zT_5AHsw!GCRn>C>=krgjIk_g$lz8yW>glF*#fIrGeOR&eG8Om}lRMvBvXob;rD&Zm zH?~6|HH$FR+mZJ%jaJaAI&sm(NEl$o-aj}n7#SQ0_vnH)+^;)`vl!h1-|7p6hJua) zVMm~ZNj%+?dkOo1x{azE9FFLg$XZ>%i**6js2lNsMspX;%Lw-wJj_!(8tM*|xtj|A z90Z;JNG==qjqjgamX@2ww$9q!C+g$%X?x9B`>exrVs(6V+OcTzv6Q1BvHaJL;H;YMIO2}+b;sS;+?>-}{I;|sIOpM9?h~8hnq%=rJz$8+<7!J$zm^7pyY+iu1yJ4#Y4tAED35 zLXm=};ry2HG1 z1g$e$WiB+}48or*Q`F_ELT$Wk)} z7;Cv1&tDa*lGB9>zaFHJd;}}rd=JGHyLqQLP8Ud#s69epl?B#%(QG}>?IHBMU~UdW zJ3Eg^`@6MJn@Tf1DCwdKoky3IVO)Wob>ZnjO?M4-A0}P2OA8+f>0UL|GdM5+dC(OC zk3IY`O-@~=YWl(hAzjo$0};mEKK=u0{h0AN2~aAl!!it|+fYh@POaN>$VivzNa{@&f-N)-9u^HL5jJhf%sIC zZj4Yvv^&65puTziFYP()f34C<^O)o*_!Fc-`)52zErS1y>wf;z9T8u z&bociKM}Lc`j%u|731xbwuJIV_)K`F<^8tz+A?MS3HhX)h+JI0F|~YSrmpc~-P%;$ z+G{q!zdj~klR00_i>;H7ootKQOn$olxM$W=HLgvzC#-4D@{6AJDbMW?o z0Kb-<#CbD;C7B9;rn){;(=;1cJbCD~zEgcucTBHK*WdnQVP@;h$glkGJ(u3tnX2zh zKJsWf@YrRi)$h2*S)o~axH6w9JeGebFS&kPW`cy`H#A~v=e^rDZLi^eUgO=~Diom! zns+p9uaa@%|8EAZK;;2LBTaB#0Q4pGp9NqY;Rs?`wpxNsN2wO5Zu7%i;+t834z^fW zpLC1*xL&e^%nKw8qSjl!NiKSm^%ifsb<6D$GLNvJ*owqh$&n~2vK*hQj!)@*$=aO} zf}$)Gi;GcL)YV&S=H;&t8BUXVkGdmJ-i=^9tr%YiWeUVlPB>8u6!|T?`DMQhMmp*l z^(qd`517fK_JyB|SwnI1KCB&>@Hr-Kh?jm7CQm9J{_2?qy1S`}pP}K_w4?|>-5^l7anyWiltGfD! z`y%QWDJv&57}nK=zR-aBBxR6L4VI~2qBsNC%&wFZ7rMkc12|6Iz;lBm>H(Amf=5?$ zR}P3HLwD!ysNqbYKnRA0y4WpoKS2Ck2Q^v-CNMXz44^<= z;d3p&Y(3u!6z_eAj|m@n>nER?5>Jgz?VnzrX=#09*O^_@;=A6ry{VRaGp%jY8{Xac z_Qq?HrEGIdBvV<#O~q%Va~poS@%%Hc&gc^|Y4QKrMl^{My;CrJA>#t4uZTNH^S{4(yD%3k4YDm9J`I%gHT?KNxi4}5)V?3@IamF+Urx2{%J{3%%9@(EJI2TEo|CznMv`&0 z^_kj+%%b{C^V%QYHa+r#O{wO4XX~00E2g^Bb!#*Jy36JEGRL(l&R(9AVx=V7cI%F2 z?ibD89qWrF8sQiw(K@l#|0N|_gr+|WrKOf)F}Yfi8j_Pdif~K*T~dyfpUXv4VPa=IG!M`jweZk|+L3!e&4wS2$rY+I`S&h(;$m9|posc&F7&QxJw<|H^9|Z(nuETX74}T+_a~qmuhsrMIJH0T#3AUzo*UVLI@6 z=7l15@>jQNe0~dtVnMr0KxWt_O0kX-tf$EOumGG~07~Yq55zihY$rd22J0x@LNJsZ zv}6n=$DT^s4i?NUWGI5#leHYda@O{}KSUZ6z=)qYrNi#dw!n`coVEC-17=_sxg zhbf`+t2ySPM$np@_h-}*Br|fB@~?mp%cywiwwEFkhff}U@vCWHbF!i(>0bINETc$U z_x@9PMR{Vw#pXLx&36)?07Chkd?FJF5=yCQ0z5P7!5%JY+O&u3`V5+O`^iB-21oN-&6$7|5R2Jac2_fD1(lozu+6NNnIZrZw7N63wv-^p&*7&(R{3CH-bz zc*|;Prr3wMLoD3XYMfUTx6nd!hzZ}+dCRr~JmlJ*M>p(SnFtG7w~St{!QU^JHpU;Xk2*^kXU!7%79YP&cea?rg?imj+nU~J8?*_*#v z)Il4Vk(FjkVI9bh3N&-;EytvyU`*O_buepyIVQ%YJz>JXTx~g)56xfbkF=>Y$8|me zmA(lv1}4=m#6Bhk!0@dpI}nD2pv&~5f}z(*%5@lSRow#-Y-^DWz&;ea)%1$2lS~(F zf*Faro9$I)x5PBv5$eNUJoX?_m*1$pUprXvhG2Dp$Y2KJuTgUdzN^m>A*F#Ck#%e6 zS-eCsRQ)_mB_5+oXk)*zOK+Hy@1w_@MD8F$3rRghgy5%P5AYV?R%QyIH|x&szP?=Z zJw~t0dzH*(hWSH8-!*lZ?zhonj@-~`Y)3A;SE*_=x~ATry&CHOC_PK$r})=6AiGZy zVTPSKQ?vMD&DvDW+N4+op)vMgrm}jhll+S|Ol&&2DeZ2&!bwt7jGuLvjqe!WIJs`} z(aCM&4QY21tjy5!H8>|&lWW^EYt~MePYY9DB#r%!n0RS%Q_M52rsPE#*>_Q1gvjKE z59P)=C#sus1*Dc+;Aq2;oe34VGEvUFK6wlz};?;XLt zd7RJqD<>X4`EX)u;;A#8(>v3CO7r<)MM?QuQ@#xufBi-O(v*K`qC44ocWT)^ANn_4 zw&Bq^33bD;U9&7xT{{sy8O#r`l#))=jTSw`{)bwv{>N zD!7{E1-<3%EoNu^<9jBTOe!C`nt}f9E!QVKpQ+ppjWFs5-fdO)`?#Mk+iJt*m%ioq z-(4)%31CB_b5}2c3;iYK`sdx^gg#9R7J+SNX5}yVHpy<`)tt6ngaZV~W*%@l`=lZf zS0~5CY|Ivw@U%6Gtv1Rb7e&ZHTF5Mza&pvOh!x1TYLv&-kR3KdUMQmC0Km3mi_ZC$ z>&&6CNHVUZiD1{vI`HCCwwD;z4!_9Rl9>sfFEU<#mw(xwrZGKyJ zUT%(im%#$i-P6~mK8jZyv9-aG6>b0q_Xf-0jyTsu;6vAsO|ipk=J`r7tHI-7+j(Zm z)cKWoPIrQvs@<^24#6?Oo88arIxwh)!UqR*H{RLxOn4BY06U7qe}gBoS>S~O;gGT! zep_Ep{%BRBncUq&u6xK>+t|8&_oGiV!maBa=KV$G;MVnZczu2Q?nXR*xT_nT7#!HF zgthFhWbkGM1MHssSAxz00l|a;lM;j*@(y6?D!L`vBq3-mrB|3F0?Sd;$*A&L6YEKsB{H=wlPeQnLasBl zt|Pg2dor*CM%0R$7lYUr%rahix;E9iIaAk?*oz&+>DHOMQ|otRnwG#U{OCDbruFu7 z5o{~EVY}sFH?=#+_w2Ym>uh#hz7$x_bd<)6HNjL(FkQ29&cavk<}cf=u$-3N!g4CL zJGa$yKdtv}Yb)ktq&qN9MuPsRa8x`h9km{{5l>UNKCYi9$HRWhF)Q;tXI|>ZgnsFm z*l#^1_1liwcrGTy#F!Mb#%$y``7L>rw}AgB65Fb$x)E54tk@NYiMN2;ETb01MXOlx zoF=DkZztRR&f3^?+ZXqA?Rs_R4C=G9hYIFKy4Jjl!Q{ApJWd)4ZA!#8+YdY z!ZC1i%nCm{@mqHIq>53YD?2D0)E?BAQ_zfe5RJq44t{Uha_p@%@#lp+E6g+EQl2Y{ zwW$3i)UkX`aVs;60x%l^{=>}>9(g92^P?!Gg2f>+Na=eLM;u0_=On7lh>h?$c53g` z9V4Or+8|bXLRn;d5HXW{#$TYCC*%dnzL2l~31EGU?-!<^1Q2$UQpO7yez|>%#%Gbej0*hkX#>Hd1n6tO~RlJ1B3|$C4 z+jExh47zmdHr7|&gBh%0?+!B)F9VaiDA%Oq8nE1*7v#3z%D(Zf$36cj^7h(G)xgrnU2_(zeMJWT zXuIOO#t%-8Oj$0tR?_`L@rRCg&f*;p#2*-c`uJ{2!|RUkhAVHGFXp%omte!Qe8p_d z;>qq;stpE;4<)_u0PIyU$DijU-oE07MquZB#a6k!j{Bg_zD*WBSQ^+S2tO4h#E&@p zB5jxyZFJ({h>MPwkV?|FcI~<&_7Qb(;9%Rj^|$3p(dX2i1B2Rd&mnfGLYGGl!G#7} zu9)#2_5ijDjMEjmYk1%(9F-Z!mNWOM5Mnm+1kV~MHnPzkGXgX>^F93$t=0`~o}&CCS%4^Uoxk4W)UVHFuOg_JqW zl=~oYJ;wdL*qjube=lxHi(5Vs%O&3&2XeWUt6h?4o31}sn_jg0 zw}FO)dunv1Hyzk@SrWWfh+}LaHfO4~z)j{de7qasr7blYXDPBg8P#(fS;UDsTzv?D z>?lVf?5KRye$?^T0q4gYH~SwRvmV6Rpo1`(ofeMCd9zzdYdy6Y0e;?YBbEog$7~R# zT+9}eWA>QifK?G7QZ0Fxk$6czM3ckg#DkDM@c{k3&?mB#`wQT(aDa#5j$kYT&N|fL ztb+(yC}z((NHB-wxW7?~XkKElvkrx&tMjGISaE4ZiaSd80;vJ>nfx`9GuFjr#Y(Hx z>}^pFQ&vjmFm=StiI>~~Fk7hDPs_Q}76MLNI0*a?h1=l%EIvYLM7V$Z{am?&ZJ;I z5WsTAI0IlEuqXt#ka9r45kMLDG-r%^neGq3pMT%htkH|d?8zyHW`yP_TVQ?_a*b@+ zxDdBs)agW?B@V^;)trh5P%3n)AD|q21MGG*6>;pPo|2l0&}=u%QNjpLAkskR1?CNT zQ(t%gex-ZMkrkLkIH@+Q-GqW$<^xzn00FC!iwy5Kzp0TOeipa~tf*nV9&5yN*Q{O;gbg>Q>_$9F} z=6c8c{Y@O8{v1Uf0_hZ%8$7n0{==pKjigZ8!9Ka9LO$Vfr*2e>KCWmf1RK@W~&y2SDs+6t@fx z>6QZn%vxD6d)P>mxP^IDRR>jIROtM=h?5gfv8lz-*jJd4WpgWw#;nE!IuclDYHg$r zyo-wejDxP_RpXI~ubuo_x_tQ+&TVfdn;qO#$3s~Q<>pw&rE-9y2a{zFWy-5Q3~Wr6 z-IXb?`LK3Nvh3bDn^^0)!a3c|C#~a_@%C#DuBz_EC&n#v7OA}XxBjY^I$#;aQH{zw zGVm{IKV_p62d_W#-DjtloEg0kxcw6WS=TR@qlU{8UjK<7_j6SK&mUF8X3Cd0XZ(vM z9ys~Hr23ss9IYyEzM-u^+s{|n+x_DC+QxR-cK#k2@egFFz1;bM+eUH!qV{#R57zRa zXDzxF2$`Id*=WkJ3(kfBUffV`B*&>Hcq- z{yUKd5M~si=~0Bj(v(ePEzJTC$p885R*{Uaa&rGvd&;o}p4CqOi7oLhX=m+N$E@4` z{F7rlJ`p`q^JNa?+7hm&9-bX=T}}s9!jSC6c@zIN0hs%uw;|)WCtu?*#K=K*j& z5|yX+roFAR6%7d?vGa^OU9mb-5uh_g>568!n7OQToL#d1>53a|U^?Ijd_zZ*SR}@C zQ1XDlpd>)#QRh+DQTJa7O1euxNqZhX#_e%;l+RmVA)-q9Lz{fuwCavpu&0`{uSR*) zf?pYcEmkVre=InVLn`&WP`vke%$eV#E#=B#LWyIpyf1Ms%>gjW#T)=Ju9zF37BB_i zl&~BN=o2y0dBFr00JUNPs1*y^LBeWYfLn$LP@CJi%z|2JRLpw|o9%+tTgn361-#14 zlE(gEQE8I@^Fndy!moDoyNyFPilgw@$|!8)ad+v_Kx)gp*Wnha%Wh81@2ckaOyjvd z(k74L5MeL|T9<6G>!|gg7(zG<5ixlczWo=jw`<=(Nno+rO6-ijcJHZs z)0Hc)a5i`AI6qriYl4wh_#jX0KDj%wWUAxhs!gd?o5puv@Ndqdll4F+uQ^URUcckJ z8>jZ4*?b|ej*!Xv>zD1Q=o1-H0wwEBP=eYFP%`ox|8hbltv9qFn(@7cts5=p#l>4! zN#~nnim#Hk-flm?&PMT#i`pxt5Bxl+nnG#y5fI%wuaC?&SO`W`ChOTJ_#gCuQw=@M zun76++p;bJ1_Y@o1wo91USFgalKNnr?7EL)zaeswNHHi;Gn7UjrC|Fe*8szD{}AQo zP>Qjy<0;tJVXvBpSGLhJM?iF2*89NFCrY6gIt=g?%n23>%t(Jvbr9J{-~O-(U_PgV z8GKULO)KyAz1^2yvn}oEpdFs2Qu9Z0^`tdnN!&l(G_(I)_lNS%g7aIoOmw& zT-w_Z6KBf$+s?MW0#MzTlD1tpdWh`5IBsL_s8#lVHakWz3P&YUTLK7i0^n|C#6?@ zR8cwJeR`R(0Y0@MRRcJ)HfFnIc&%mk$`{qYw&v6txY=KN~)x8~E>!jgJR;H5mmgio5iq$irF$Rt^@HTmi$klYVt_zChw`BK(DbE0Q|7^8}S z7yU3s)*PnGVXdp1e#|0hQSU^ZDt!)(plSApYSgR!semk61@SsD_wmTv2O9RY_3rk;`ypJCJrQC8t7mT0K+li&^n&Aecc@;)U|I3!XM} z3p6tE+{x!IxR${=W#XQb_e?&M_AQHbWGZ|Ut4^-MzA2xsSP|QXv+gH$$9K;*EO~v$ zcXv*yXLe6No@%};-Eem%P?sU6yVT;g%VhJr#))n(tbOi^ms)7rPJSWL4l7GV>jhWv znve5Vm~Hsb(>CYN_Tpv;^QlLGGUgAMByfDEEq{E+I<C8AdE&Po9qkbn)`Aph~R zcqq77{Rf&g3@8)pP=8DW;Duv9lV;#%(mD+&Oe=#O>oR`)O2)n}`z=Gm-<3GC0mf<8 z8#sIi#X5<6g@{UIgvcn77l_OdA(I`0y=(@nuhDHe5ibZf6#561;l9uo^$KoDo@xCc zz==G6)y46yzvIgOj&uH=d+;Ln;P1F)zvJq!Sb4trs_f={R~l`6`4zu|ufMv;$v@0r ySyIC{U)jiNbiJNf0j^~zp;1#0;U-^4ew=KMdhFb4qu literal 0 HcmV?d00001 diff --git a/loadtests/__pycache__/verifier.cpython-312.pyc b/loadtests/__pycache__/verifier.cpython-312.pyc new file mode 100644 index 0000000000000000000000000000000000000000..99e923cdffa8bdef54f63316be2a3381b2463a33 GIT binary patch literal 4278 zcmZ`+Yj6|S6~0%ydaY%zWXrY;30<3qkqN@!1Sc+K3V|dh3DChI#e}Mgc1`5y<=qw7 z?kY8pN#(&A?7+lgLL1S^Ov#W=Jp@_@9K!$p*^Sa+v%`;0+D^ki*G{G>zk2RUD;Z+W z%$|Gix%ZxPU*~*h|5{tiAfSJKa$@lFI)wf~8Wz#z%Hme2Od}rgIF1siA7RAXcq)z) zwtgGb>C>p6vZ3b@Z$FKA#}!*?UOzqJ47%=-Mi3h=J|YWpEFs*(B|&W*1j*<~Sdt^M z6eUuNh+zu8#Zh3>$c^Mm?_gQEHI$SJx>D{Rv@eyH>Q~Al9*;ur0kj&)4r`R`vT6gi zHAti|I)p}S$I&o#2pz{kDnoYOzNK@^){UJ{Z`{1KZ}XEqTe^C-YzcL3?b`a}(;GUw zdOADXxv$@TmlM*d=o?%@;?}>?_iQk9?;F^`-k@D18!>7J#8e_Cr3O-BA|e}ZauBfZ zMA~2uM&cP^w+>cY_cSgZE^P-cuSqjUk za6%H{<`4^?^b=rXXr89*{=01R9d@n8t}VP?WP^9uE{*LfvQN#@8oNz(ZTtSNr|F_w zt!Z08*zW&|ZJ4K^yf1k`!^ySl(Pg%W`V{w2QD|4lx`xG7Aeu%3>c>3lxAB;~0Nw_i z<|*Lz5jto$TnEdmjan~>C?XeVUT(-|tT^`^&i^gA_+(1oI8(H7M}#oWX;;M4E|2d$nZ<>1m8WjTyodKi`z z^Y~-imGby36+>BCp$45xl$EPIbre|+whJ9OgHYD4*q0qiTV64YC2&*F=};WYjzB8;%rDbU$LY!Ec?KSvYZSttb)H$bOkR&ke9dVD`)IV<)zXS5yO340 z>}HvxC|kojvmTg9p3_5QjfZ>t1wyaE!_GRhwMs3|w7^KSU6*B*IuG$5A4~e)GT`&KFcmfij!yFTn{dS zCypq-tY_kvGy=DU(^K3EqtwW)mR+UBnx%M7)XA$Wat~i?t%$F?=sJLSXQ>RbnQriR z^^$s3y{?U^H`EzU{an4pskhW?>Sw^`v{RgRN;{#Qgzj7Fr=+G{)lP*tVK^;BWr63u z9y`xzCt&U^^`a04; zfdN=0EbTMx3`~c`-ULI}v@>u3GoXB;x-Io8_trCD^pTcAB8d<|o8jh#W3i|J!O*bB zV~Ln-*kiJgkPK%!GD6V9jEDvWr?OWh5#Df0LQ)P#WRbj+AAoNc2p10mdl9gngtsGb z#aUd1VFg3A#(DGcBcWHNxlhO?Z}7~X-HD9K?f!AVHLdJ~)g_l3BS5M*)0 z#3_S13UvZj1|tbFFAPL7aoKQ0(rF>diyY~o5|LrJxJ<%u93z>6Bof~@5#L~9QY;B+ zL^3KER4fU%J^+_3izEv&>|!dD_*)iuto~!Jm#3(2oZX` zXZMZo%kRjGMNiuuPf+s&bzl8-<5c6Md~Rhv@_}$cDD)Org;ejR?6+zxAWW zK%PF=n{(e~>nAr)zENn=nmUT?h8%sjF#u2Pa^Ig0e{}fDkuL+i#m3$oqt|(Jl@DA9P&kc>jq7#Sd6 zxexni=xfs4?wR=C);@dJ-<+rNuN3+&yjJvgOk&;JFwIOc`OW8R^#*?)pWCW)ty8`8 zEU2HaM=MtA0Zwme1LOLNHF{Ia-R3og1I6Z`-u+Xx`=t-2<~JVoY>LaqVXlpWf1rEOOn(6(3F4!|;o zcU0Q0i9L`(zTO1pxUAARwTaHV(4ufL44gDig?px`{J4LS|qu zx&4FAEe65vJ@N}At^pQw8P1ZAN@jehh62M4kxmMW5aZ0KAa=lTgGmcwxHJKNZaAW5 zCt)P=isX06j6T7dQrOuI0;9nSqv7gYqfSoAk$BkZGT2hP8fTf3pcHM(otd;`P@u$c zM&fZ2Uk%KQ&wvdHjvp-vK&9ZrkROB+R*W`OB2QIFRwt>z#?aqep7>ar>As?9EkDlH z=}$hTx}M6h!15i1(UQZ$FM(6{1s0IQy_y6Y{~d3e=4~r%E*!qJUF~|V=zSi<9`@{> z@ja766EEwYmHDB<$_sC4o{hPEx~D1Mnt$#BM3vRKJ*C!57dAnRA+3hVp?qNKSDL47 z{Q3D>Fz~Gx1y<&}&nbFfo!+uu_pkbD=l+Y0`O(XJwIB7&?)b80yL#{tI5F|ySBuRX z^%ajlt)J p1#K5era^p7A_*+~hhqGlAK_L1Lahr9JNEzAy9;Ax0TFJ_|3BhGTP6Si literal 0 HcmV?d00001 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, + }