Compare commits
3
Commits
4c8bb39354
...
3fedf9317c
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
3fedf9317c | ||
|
|
090cedca67 | ||
|
|
a4a11bae54 |
+1
-1
@@ -28,7 +28,7 @@ ENV API_PORT=9090 \
|
|||||||
SQS_QUEUE_NAME=iot-telemetry \
|
SQS_QUEUE_NAME=iot-telemetry \
|
||||||
SQS_REGION=us-east-1 \
|
SQS_REGION=us-east-1 \
|
||||||
SQS_LONG_POLL_SECONDS=20 \
|
SQS_LONG_POLL_SECONDS=20 \
|
||||||
SQS_VISIBILITY_TIMEOUT=30 \
|
SQS_VISIBILITY_TIMEOUT=120 \
|
||||||
MQTT_CLIENT_ID=iot-bridge \
|
MQTT_CLIENT_ID=iot-bridge \
|
||||||
MQTT_PORT=8083 \
|
MQTT_PORT=8083 \
|
||||||
MQTT_WS_PATH=/mqtt \
|
MQTT_WS_PATH=/mqtt \
|
||||||
|
|||||||
@@ -1041,3 +1041,140 @@ OSS-образ emqx/emqx:5.5.1 СОДЕРЖИТ postgres authn/authz
|
|||||||
| 4. Процедура деплоя | doc/deployment-nubes-production.md. |
|
| 4. Процедура деплоя | doc/deployment-nubes-production.md. |
|
||||||
|
|
||||||
Версии: iot-service v0.1.2, iot-emqx v0.2.4 (Docker Hub, latest обновлены).
|
Версии: iot-service v0.1.2, iot-emqx v0.2.4 (Docker Hub, latest обновлены).
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## 30. Ресурсы кластера + новая директива (15:55 GMT+03)
|
||||||
|
|
||||||
|
### 30.1 Замер ресурсов кластера (kubectl, фактические цифры)
|
||||||
|
Ноды (4): control-plane-xb699, workers-6f74n, workers-bhbvs, workers-v8zq4.
|
||||||
|
Факт (kubectl top nodes): control-plane 33% CPU / 53% RAM; workers 14–19% CPU /
|
||||||
|
24–40% RAM. Свободных ресурсов много.
|
||||||
|
Аллокация requests/limits: workers ~3.5–4.0 CPU requests, limits 11.5–14.3 CPU
|
||||||
|
(оверкоммит), память 16–31% / 45–77%.
|
||||||
|
Наши контейнеры:
|
||||||
|
- EMQX: requests 10m/12Mi, limits 2 CPU/1Gi; факт 24m/202Mi.
|
||||||
|
- iot-service: requests 10m/12Mi, limits 200m/512Mi; факт 1m/4Mi.
|
||||||
|
- ResourceQuota в наших namespace НЕТ.
|
||||||
|
Вывод: повышать ресурсы не требуется; при желании «на вырост» через deck —
|
||||||
|
EMQX 4 CPU/2Gi, iot-service 1 CPU/1Gi.
|
||||||
|
|
||||||
|
### 30.2 Новая директива пользователя
|
||||||
|
- «всё документируй» — принято, этот файл ведётся.
|
||||||
|
- «надо сочинить тесты — МНОГО РАЗНЫХ нагрузочных тестов» — создаётся набор
|
||||||
|
нагрузочных тестов (отдельный бинарь/скрипты, см. секцию 31).
|
||||||
|
- «попросим соннет сначала код ревью провести и про тесты пусть расскажет» —
|
||||||
|
Sonnet-ревью кода и рекомендации по тестам запрошены (см. 30.3).
|
||||||
|
|
||||||
|
### 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 перед большими нагрузками).
|
||||||
|
- Прогон полных сценариев — по команде.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## 32. Фиксы ревью внедрены (v0.1.6) + нагрузочные прогоны (17:00 GMT+03)
|
||||||
|
|
||||||
|
### 32.1 Внедрённые фиксы (код)
|
||||||
|
| # | Что | Где |
|
||||||
|
|---|---|---|
|
||||||
|
| 1 | Асинхронный SQS-диспатчер: канал 10000 + 4 worker'а, ретраи 1/2/4с, дроп при переполнении | bridge/sender.go (новый), handler.go, bridge.go |
|
||||||
|
| 2 | VisibilityTimeout 120с (config дефолт + Dockerfile + env на деплойменте) | config.go, Dockerfile |
|
||||||
|
| 3 | Backoff (1с→60с, x2) при ошибках обработки PG | consumer.go |
|
||||||
|
| 4 | Shutdown: bridge+consumer → flush → HTTP последним | main.go |
|
||||||
|
| 5 | Батч-INSERT (100 строк / 200мс flusher, final flush в Close) | iotpg |
|
||||||
|
| 6 | Disconnect(2000) вместо 250мс при остановке | bridge.go |
|
||||||
|
| 7 | (whitelist tenant-БД) — ОТЛОЖЕНО, отдельная фича | — |
|
||||||
|
| 8 | Пулы: devices 20 (env IOT_PG_MAX_CONNS), admin 10 | store/open.go, iotpg |
|
||||||
|
| 9 | JWT: опциональная HS256-проверка (env JWT_HMAC_SECRET) | middleware/auth.go |
|
||||||
|
| 10 | Close rows вместо defer в цикле | iotpg AdminStats |
|
||||||
|
| 11 | Пустое SQS-body: лог + удаление из очереди | consumer.go |
|
||||||
|
| 12 | Валидация топика (ns/deviceID непустые, "telemetry") | bridge/handler.go |
|
||||||
|
| 13 | CREATE USER: проверка роли + QuoteIdentifier/QuoteLiteral (не Sprintf) | iotpg |
|
||||||
|
| 14 | Валидация MQTT_BROKER_URL (ws:// или wss://) | config.go |
|
||||||
|
| 15 | generateMQTTPassword: 3 попытки crypto/rand | handler/devices.go |
|
||||||
|
| 16 | Баг DO-блока: параметры $n в DO недопустимы (найден прогоном) | iotpg |
|
||||||
|
| 17 | Пагинация: List limit/offset; telemetry limit+offset | devices.go, telemetry.go, iotpg |
|
||||||
|
| 18 | Admin-пул 10 | iotpg |
|
||||||
|
| 19 | Лог битого JSON — уже был | — |
|
||||||
|
| 20 | Subscribe-таймаут/ошибка → принудительный реконнект | bridge.go |
|
||||||
|
| + | Оversize >250KB → явный дроп с логом (лимит SQS 256KB) | sender.go |
|
||||||
|
| + | Гонка getTenantDB → per-ns мьютекс | iotpg |
|
||||||
|
|
||||||
|
### 32.2 Баги, всплывшие при выкатке
|
||||||
|
- DO-блок с $n-параметрами («got 2 parameters but statement requires 0») —
|
||||||
|
старые сообщения крутились с backoff; исправлено (32.1 №16).
|
||||||
|
- Зеркало платформы кэширует теги: повторный push v0.1.3 не подтянулся
|
||||||
|
(rollout restart) → правило: КАЖДОЕ изменение = НОВЫЙ тег (v0.1.4).
|
||||||
|
- SQS-лимит 256KB: 300KB-сообщения дропались SendMessage 400
|
||||||
|
(InvalidParameterValue) → явный пре-чек 250KB + документировано.
|
||||||
|
- Шлюз платформы рвёт ответы API >~15КБ (эмпирика: limit=50 OK 15KB/0.27с,
|
||||||
|
limit=200 завис на 15.6KB) → добавлен offset в API, верификатор ходит
|
||||||
|
страницами по 50; large-payload верификация только по логам consumer.
|
||||||
|
|
||||||
|
### 32.3 Результаты нагрузочных прогонов (прод, GMT+03)
|
||||||
|
| Сценарий | Итог |
|
||||||
|
|---|---|
|
||||||
|
| baseline 20 уст-в × 2 msg/s × 60с | sent 2390, delivered 2390, lost 0 (среди sent), dup 0; p50 466мс / p95 785мс / p99 1486мс |
|
||||||
|
| burst x5, 20с | sent_high 1958, delivered 1948, dup 0, recovery 112.3с; латентность пика p50 17.4с |
|
||||||
|
| large-payload 200KB | доставка по логам consumer: 6/6 сохранено, oversize-дропов 0 |
|
||||||
|
| reconnect-storm 30с-циклы, 120с | 30 форс-реконнектов, 1191/1191, dup 0 |
|
||||||
|
| multitenant 5 тенантов | first-msg: p50 1376мс / p95 3085мс (создание БД) |
|
||||||
|
| acl-violation | 11/11 чужих публикаций разорвали сессию, accepted 0 |
|
||||||
|
| auth-neg | rc=4 / rc=5 (отказы) |
|
||||||
|
| api-crud 5×10 | POST/GET/DELETE p50 ~240-255мс, 0 ошибок |
|
||||||
|
| telemetry-query 5×20 | 17.1 rps, p50 241мс, 0 ошибок |
|
||||||
|
|
||||||
|
### 32.4 Версии и артефакты
|
||||||
|
- iot-service v0.1.6 (digest 5868e9481cfd) + latest — задеплоен.
|
||||||
|
- loadtests: API-ретраи (3×), парсинг 204-пустого тела, offset-пагинация,
|
||||||
|
burst считает по фактическому sent, large-payload без API-верификации.
|
||||||
|
- Repo-память: /memories/repo/iot.md (факты платформы, правила деплоя).
|
||||||
|
- ОТЛОЖЕНО: tenant-whitelist (#7), DLQ в shared-sqs, soak 24ч.
|
||||||
|
|||||||
@@ -1,7 +1,7 @@
|
|||||||
# Makefile — монолит iot-service (образ naeel/iot-service).
|
# Makefile — монолит iot-service (образ naeel/iot-service).
|
||||||
# Старые k8s-цели — в legacy/Makefile.old.
|
# Старые k8s-цели — в legacy/Makefile.old.
|
||||||
|
|
||||||
VERSION ?= v0.1.2
|
VERSION ?= v0.1.6
|
||||||
IMAGE ?= naeel/iot-service
|
IMAGE ?= naeel/iot-service
|
||||||
LDFLAGS = -X main.version=$(VERSION)
|
LDFLAGS = -X main.version=$(VERSION)
|
||||||
|
|
||||||
|
|||||||
@@ -89,7 +89,7 @@ func main() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// HTTP API — в основном потоке.
|
// HTTP API — в основном потоке.
|
||||||
router := api.NewRouter(h, log, cfg.AuthTestMode, version)
|
router := api.NewRouter(h, log, cfg.AuthTestMode, cfg.JwtHMACSecret, version)
|
||||||
server := &http.Server{
|
server := &http.Server{
|
||||||
Addr: ":" + cfg.APIPort,
|
Addr: ":" + cfg.APIPort,
|
||||||
Handler: router,
|
Handler: router,
|
||||||
@@ -114,9 +114,13 @@ func main() {
|
|||||||
}
|
}
|
||||||
}()
|
}()
|
||||||
|
|
||||||
// Сервер: при завершении ctx — гасим.
|
// Сервер гаснет ПОСЛЕДНИМ: сначала останавливаются bridge+consumer
|
||||||
|
// (ctx → disconnect MQTT → flush SQS), потом HTTP. Иначе платформа
|
||||||
|
// может убить контейнер по liveness, пока bridge/consumer ещё работают
|
||||||
|
// (фикс MEDIUM из ревью 2026-08-16).
|
||||||
go func() {
|
go func() {
|
||||||
<-ctx.Done()
|
<-ctx.Done()
|
||||||
|
wg.Wait()
|
||||||
shutdownCtx, shutdownCancel := context.WithTimeout(context.Background(), 10*time.Second)
|
shutdownCtx, shutdownCancel := context.WithTimeout(context.Background(), 10*time.Second)
|
||||||
defer shutdownCancel()
|
defer shutdownCancel()
|
||||||
_ = server.Shutdown(shutdownCtx)
|
_ = server.Shutdown(shutdownCtx)
|
||||||
|
|||||||
@@ -36,7 +36,7 @@ func (h *Handler) ListIoTTelemetry(w http.ResponseWriter, r *http.Request) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
rows, err := h.IoTPG.QueryTelemetry(r.Context(), ns, deviceID, limit)
|
rows, err := h.IoTPG.QueryTelemetry(r.Context(), ns, deviceID, limit, 0)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
h.Log.Error("query IoT telemetry", "namespace", ns, "device", deviceID, "err", err)
|
h.Log.Error("query IoT telemetry", "namespace", ns, "device", deviceID, "err", err)
|
||||||
writeJSON(w, http.StatusInternalServerError, errResp("failed to query telemetry"))
|
writeJSON(w, http.StatusInternalServerError, errResp("failed to query telemetry"))
|
||||||
|
|||||||
@@ -15,6 +15,7 @@ import (
|
|||||||
"encoding/json"
|
"encoding/json"
|
||||||
"errors"
|
"errors"
|
||||||
"net/http"
|
"net/http"
|
||||||
|
"strconv"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"gitea.services.ngcloud.ru/Nail/IoT/internal/service/store"
|
"gitea.services.ngcloud.ru/Nail/IoT/internal/service/store"
|
||||||
@@ -72,12 +73,19 @@ func deviceToResponse(d *store.Device, password string) iotDeviceResponse {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// generateMQTTPassword — 32 случайных байта в hex (как старый контроллер).
|
// generateMQTTPassword — 32 случайных байта в hex (как старый контроллер).
|
||||||
|
// crypto/rand ошибку даёт только при сбое системного PRNG — 3 попытки
|
||||||
|
// (ревью 2026-08-16: не возвращать пустой пароль).
|
||||||
func generateMQTTPassword() (string, error) {
|
func generateMQTTPassword() (string, error) {
|
||||||
|
var lastErr error
|
||||||
|
for attempt := 0; attempt < 3; attempt++ {
|
||||||
buf := make([]byte, 32)
|
buf := make([]byte, 32)
|
||||||
if _, err := rand.Read(buf); err != nil {
|
if _, err := rand.Read(buf); err == nil {
|
||||||
return "", err
|
|
||||||
}
|
|
||||||
return hex.EncodeToString(buf), nil
|
return hex.EncodeToString(buf), nil
|
||||||
|
} else {
|
||||||
|
lastErr = err
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return "", lastErr
|
||||||
}
|
}
|
||||||
|
|
||||||
// CreateIoTDevice — POST /v1/namespaces/{ns}/iot/devices.
|
// CreateIoTDevice — POST /v1/namespaces/{ns}/iot/devices.
|
||||||
@@ -142,13 +150,17 @@ func (h *Handler) CreateIoTDevice(w http.ResponseWriter, r *http.Request) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// ListIoTDevices — GET /v1/namespaces/{ns}/iot/devices (без паролей).
|
// ListIoTDevices — GET /v1/namespaces/{ns}/iot/devices (без паролей).
|
||||||
|
// Пагинация: ?limit=N&offset=M (0 = без ограничения).
|
||||||
func (h *Handler) ListIoTDevices(w http.ResponseWriter, r *http.Request) {
|
func (h *Handler) ListIoTDevices(w http.ResponseWriter, r *http.Request) {
|
||||||
ns := pathVar(r, "namespace")
|
ns := pathVar(r, "namespace")
|
||||||
if ns == "" {
|
if ns == "" {
|
||||||
ns = "default"
|
ns = "default"
|
||||||
}
|
}
|
||||||
|
|
||||||
devices, err := h.Devices.List(r.Context(), ns)
|
limit, _ := strconv.Atoi(r.URL.Query().Get("limit"))
|
||||||
|
offset, _ := strconv.Atoi(r.URL.Query().Get("offset"))
|
||||||
|
|
||||||
|
devices, err := h.Devices.List(r.Context(), ns, limit, offset)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
h.Log.Error("list devices", "namespace", ns, "err", err)
|
h.Log.Error("list devices", "namespace", ns, "err", err)
|
||||||
writeJSON(w, http.StatusInternalServerError, errResp("failed to list devices"))
|
writeJSON(w, http.StatusInternalServerError, errResp("failed to list devices"))
|
||||||
|
|||||||
@@ -25,8 +25,14 @@ func (h *Handler) ListIoTTelemetry(w http.ResponseWriter, r *http.Request) {
|
|||||||
limit = n
|
limit = n
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
offset := 0
|
||||||
|
if os := r.URL.Query().Get("offset"); os != "" {
|
||||||
|
if n, err := strconv.Atoi(os); err == nil && n >= 0 {
|
||||||
|
offset = n
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
rows, err := h.IoTPG.QueryTelemetry(r.Context(), ns, deviceID, limit)
|
rows, err := h.IoTPG.QueryTelemetry(r.Context(), ns, deviceID, limit, offset)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
h.Log.Error("query IoT telemetry", "namespace", ns, "device", deviceID, "err", err)
|
h.Log.Error("query IoT telemetry", "namespace", ns, "device", deviceID, "err", err)
|
||||||
writeJSON(w, http.StatusInternalServerError, errResp("failed to query telemetry"))
|
writeJSON(w, http.StatusInternalServerError, errResp("failed to query telemetry"))
|
||||||
|
|||||||
@@ -2,6 +2,8 @@
|
|||||||
package middleware
|
package middleware
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"crypto/hmac"
|
||||||
|
"crypto/sha256"
|
||||||
"encoding/base64"
|
"encoding/base64"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
"net/http"
|
"net/http"
|
||||||
@@ -12,11 +14,10 @@ import (
|
|||||||
//
|
//
|
||||||
// authTestMode (env AUTH_TEST_MODE, дефолт false):
|
// authTestMode (env AUTH_TEST_MODE, дефолт false):
|
||||||
// - true — принимается любая строка без пробелов (для локальных тестов);
|
// - true — принимается любая строка без пробелов (для локальных тестов);
|
||||||
// - false — структурная проверка JWT (sub + exp), как в старом коде.
|
// - false — проверка JWT: sub + exp; если задан hmacSecret — обязательна
|
||||||
//
|
// HS256-подпись (HMAC-SHA256, env JWT_HMAC_SECRET). Без секрета —
|
||||||
// В новой архитектуре подпись JWT не проверяется (как и раньше): внешний
|
// структурная проверка (периметр обеспечивает платформа).
|
||||||
// периметр обеспечивает платформа, полная валидация — на стороне шлюза.
|
func Auth(authTestMode bool, hmacSecret string, log *slog.Logger, next http.Handler) http.Handler {
|
||||||
func Auth(authTestMode bool, log *slog.Logger, next http.Handler) http.Handler {
|
|
||||||
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||||
header := r.Header.Get("Authorization")
|
header := r.Header.Get("Authorization")
|
||||||
if header == "" {
|
if header == "" {
|
||||||
@@ -38,7 +39,7 @@ func Auth(authTestMode bool, log *slog.Logger, next http.Handler) http.Handler {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
if err := validateJWT(token); err != nil {
|
if err := validateJWT(token, hmacSecret); err != nil {
|
||||||
log.Warn("auth: invalid token", "remote", r.RemoteAddr, "path", r.URL.Path, "reason", err.Error())
|
log.Warn("auth: invalid token", "remote", r.RemoteAddr, "path", r.URL.Path, "reason", err.Error())
|
||||||
http.Error(w, `{"error":"invalid token"}`, http.StatusForbidden)
|
http.Error(w, `{"error":"invalid token"}`, http.StatusForbidden)
|
||||||
return
|
return
|
||||||
@@ -60,13 +61,37 @@ type jwtError struct{ msg string }
|
|||||||
|
|
||||||
func (e *jwtError) Error() string { return e.msg }
|
func (e *jwtError) Error() string { return e.msg }
|
||||||
|
|
||||||
// validateJWT проверяет структуру JWT: три части, корректный payload, sub и exp.
|
// validateJWT проверяет структуру JWT (sub + exp) и, если задан secret,
|
||||||
// Подпись НЕ проверяется — см. комментарий пакета.
|
// HS256-подпись. Без secret подпись не проверяется — см. комментарий Auth.
|
||||||
func validateJWT(token string) error {
|
func validateJWT(token, secret string) error {
|
||||||
jwtParts := strings.Split(token, ".")
|
jwtParts := strings.Split(token, ".")
|
||||||
if len(jwtParts) != 3 {
|
if len(jwtParts) != 3 {
|
||||||
return &jwtError{"not a JWT: expected 3 parts"}
|
return &jwtError{"not a JWT: expected 3 parts"}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if secret != "" {
|
||||||
|
var header struct {
|
||||||
|
Alg string `json:"alg"`
|
||||||
|
}
|
||||||
|
headerBytes, err := base64.RawURLEncoding.DecodeString(jwtParts[0])
|
||||||
|
if err != nil {
|
||||||
|
headerBytes, err = base64.StdEncoding.DecodeString(jwtParts[0])
|
||||||
|
}
|
||||||
|
if err != nil || jsonUnmarshal(headerBytes, &header) != nil || header.Alg != "HS256" {
|
||||||
|
return &jwtError{"JWT must be HS256 when JWT_HMAC_SECRET is set"}
|
||||||
|
}
|
||||||
|
mac := hmac.New(sha256.New, []byte(secret))
|
||||||
|
mac.Write([]byte(jwtParts[0] + "." + jwtParts[1]))
|
||||||
|
expected := mac.Sum(nil)
|
||||||
|
sig, err := base64.RawURLEncoding.DecodeString(jwtParts[2])
|
||||||
|
if err != nil {
|
||||||
|
return &jwtError{"cannot decode JWT signature"}
|
||||||
|
}
|
||||||
|
if !hmac.Equal(expected, sig) {
|
||||||
|
return &jwtError{"JWT signature mismatch"}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
payload := jwtParts[1]
|
payload := jwtParts[1]
|
||||||
switch len(payload) % 4 {
|
switch len(payload) % 4 {
|
||||||
case 2:
|
case 2:
|
||||||
|
|||||||
@@ -26,7 +26,7 @@ func corsMiddleware(next http.Handler) http.Handler {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// NewRouter собирает все маршруты сервиса.
|
// NewRouter собирает все маршруты сервиса.
|
||||||
func NewRouter(h *handler.Handler, log *slog.Logger, authTestMode bool, version string) http.Handler {
|
func NewRouter(h *handler.Handler, log *slog.Logger, authTestMode bool, jwtHMACSecret string, version string) http.Handler {
|
||||||
r := mux.NewRouter()
|
r := mux.NewRouter()
|
||||||
|
|
||||||
// Health — для платформенных проверок контейнера.
|
// Health — для платформенных проверок контейнера.
|
||||||
@@ -54,7 +54,7 @@ func NewRouter(h *handler.Handler, log *slog.Logger, authTestMode bool, version
|
|||||||
v1.HandleFunc("/namespaces/{namespace}/iot/telemetry", h.ListIoTTelemetry).Methods(http.MethodGet)
|
v1.HandleFunc("/namespaces/{namespace}/iot/telemetry", h.ListIoTTelemetry).Methods(http.MethodGet)
|
||||||
|
|
||||||
v1.Use(func(next http.Handler) http.Handler {
|
v1.Use(func(next http.Handler) http.Handler {
|
||||||
return middleware.Auth(authTestMode, log, next)
|
return middleware.Auth(authTestMode, jwtHMACSecret, log, next)
|
||||||
})
|
})
|
||||||
|
|
||||||
return corsMiddleware(middleware.Logging(log, r))
|
return corsMiddleware(middleware.Logging(log, r))
|
||||||
|
|||||||
@@ -34,16 +34,25 @@ func Run(ctx context.Context, cfg *config.Config, sqsClient *sqs.Client, log *sl
|
|||||||
// OnConnectHandler, т.е. повторяется при каждом (ре)подключении.
|
// OnConnectHandler, т.е. повторяется при каждом (ре)подключении.
|
||||||
// Иначе после потери сессии EMQX бридж оставался бы без подписки
|
// Иначе после потери сессии EMQX бридж оставался бы без подписки
|
||||||
// (прецедент 2026-08-16: реконнект без resubscribe → телеметрия терялась).
|
// (прецедент 2026-08-16: реконнект без resubscribe → телеметрия терялась).
|
||||||
handler := newMessageHandler(ctx, sqsClient, queueURL, log)
|
//
|
||||||
|
// Отправка в SQS — через диспатчер (канал + worker-пул): колбэк MQTT
|
||||||
|
// не блокируется сетью (фикс CRITICAL из ревью 2026-08-16).
|
||||||
|
dispatcher := newSQSDispatcher(sqsClient, queueURL, log)
|
||||||
|
dispatcher.start(ctx)
|
||||||
|
|
||||||
|
handler := newMessageHandler(dispatcher, log)
|
||||||
|
|
||||||
client, err := connectMQTT(ctx, cfg, log, handler)
|
client, err := connectMQTT(ctx, cfg, log, handler)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
dispatcher.closeAndWait()
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
defer client.Disconnect(250)
|
|
||||||
|
|
||||||
<-ctx.Done()
|
<-ctx.Done()
|
||||||
log.Info("bridge: shutting down")
|
log.Info("bridge: shutting down")
|
||||||
|
// Порядок: сначала отключить MQTT (источник), затем слить остаток в SQS.
|
||||||
|
client.Disconnect(2000)
|
||||||
|
dispatcher.closeAndWait()
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -80,10 +89,12 @@ func connectMQTT(ctx context.Context, cfg *config.Config, log *slog.Logger, hand
|
|||||||
token := c.Subscribe(telemetryTopicFilter, 1, handler)
|
token := c.Subscribe(telemetryTopicFilter, 1, handler)
|
||||||
if !token.WaitTimeout(10 * time.Second) {
|
if !token.WaitTimeout(10 * time.Second) {
|
||||||
log.Error("bridge: resubscribe timeout", "filter", telemetryTopicFilter)
|
log.Error("bridge: resubscribe timeout", "filter", telemetryTopicFilter)
|
||||||
|
c.Disconnect(100) // принудительный реконнект (подписка обязательна)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
if token.Error() != nil {
|
if token.Error() != nil {
|
||||||
log.Error("bridge: resubscribe failed", "filter", telemetryTopicFilter, "err", token.Error())
|
log.Error("bridge: resubscribe failed", "filter", telemetryTopicFilter, "err", token.Error())
|
||||||
|
c.Disconnect(100) // принудительный реконнект (подписка обязательна)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
log.Info("bridge: subscribed", "filter", telemetryTopicFilter)
|
log.Info("bridge: subscribed", "filter", telemetryTopicFilter)
|
||||||
|
|||||||
@@ -1,15 +1,15 @@
|
|||||||
// handler.go — обработчик MQTT-сообщений: envelope → SQS SendMessage.
|
// handler.go — обработчик MQTT-сообщений: envelope → канал диспатчера SQS.
|
||||||
|
//
|
||||||
|
// Колбэк НЕ делает сетевых вызовов: только парсинг топика и неблокирующий
|
||||||
|
// enqueue. Отправка в SQS — worker-пул sqsDispatcher (см. sender.go).
|
||||||
package bridge
|
package bridge
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
"strings"
|
"strings"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/aws/aws-sdk-go-v2/aws"
|
|
||||||
"github.com/aws/aws-sdk-go-v2/service/sqs"
|
|
||||||
mqtt "github.com/eclipse/paho.mqtt.golang"
|
mqtt "github.com/eclipse/paho.mqtt.golang"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -22,15 +22,16 @@ type TelemetryEnvelope struct {
|
|||||||
ReceivedAt string `json:"received_at"`
|
ReceivedAt string `json:"received_at"`
|
||||||
}
|
}
|
||||||
|
|
||||||
// newMessageHandler возвращает обработчик MQTT-сообщений.
|
// newMessageHandler возвращает обработчик MQTT-сообщений:
|
||||||
func newMessageHandler(ctx context.Context, sqsClient *sqs.Client, queueURL string, log *slog.Logger) mqtt.MessageHandler {
|
// валидация топика "{namespace}/telemetry/{deviceId}" → enqueue в диспатчер.
|
||||||
|
func newMessageHandler(dispatcher *sqsDispatcher, log *slog.Logger) mqtt.MessageHandler {
|
||||||
return func(_ mqtt.Client, msg mqtt.Message) {
|
return func(_ mqtt.Client, msg mqtt.Message) {
|
||||||
topic := msg.Topic()
|
topic := msg.Topic()
|
||||||
payload := msg.Payload()
|
payload := msg.Payload()
|
||||||
|
|
||||||
// Топик: "{namespace}/telemetry/{deviceId}".
|
// Топик: "{namespace}/telemetry/{deviceId}". Пустые части — мусор.
|
||||||
parts := strings.SplitN(topic, "/", 3)
|
parts := strings.SplitN(topic, "/", 3)
|
||||||
if len(parts) != 3 {
|
if len(parts) != 3 || parts[0] == "" || parts[2] == "" || parts[1] != "telemetry" {
|
||||||
log.Warn("bridge: unexpected topic format, skipping", "topic", topic)
|
log.Warn("bridge: unexpected topic format, skipping", "topic", topic)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
@@ -44,28 +45,12 @@ func newMessageHandler(ctx context.Context, sqsClient *sqs.Client, queueURL stri
|
|||||||
rawPayload = json.RawMessage(quoted)
|
rawPayload = json.RawMessage(quoted)
|
||||||
}
|
}
|
||||||
|
|
||||||
envelope := TelemetryEnvelope{
|
dispatcher.enqueue(TelemetryEnvelope{
|
||||||
Namespace: ns,
|
Namespace: ns,
|
||||||
DeviceID: deviceID,
|
DeviceID: deviceID,
|
||||||
Topic: topic,
|
Topic: topic,
|
||||||
Payload: rawPayload,
|
Payload: rawPayload,
|
||||||
ReceivedAt: time.Now().UTC().Format(time.RFC3339),
|
ReceivedAt: time.Now().UTC().Format(time.RFC3339),
|
||||||
}
|
|
||||||
body, err := json.Marshal(envelope)
|
|
||||||
if err != nil {
|
|
||||||
log.Error("bridge: marshal envelope", "topic", topic, "err", err)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
_, err = sqsClient.SendMessage(ctx, &sqs.SendMessageInput{
|
|
||||||
QueueUrl: aws.String(queueURL),
|
|
||||||
MessageBody: aws.String(string(body)),
|
|
||||||
})
|
})
|
||||||
if err != nil {
|
|
||||||
log.Error("bridge: SQS SendMessage failed", "topic", topic, "err", err)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
log.Info("bridge: forwarded telemetry to SQS",
|
|
||||||
"mqtt_topic", topic, "namespace", ns, "device", deviceID)
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,145 @@
|
|||||||
|
// sender.go — асинхронная отправка envelope в SQS из бриджа.
|
||||||
|
//
|
||||||
|
// Мотивация (Sonnet-ревью 2026-08-16, находка CRITICAL): SendMessage был
|
||||||
|
// синхронным внутри MQTT-колбэка paho. При недоступности SQS (таймаут ~30с)
|
||||||
|
// блокировался приём ВСЕХ MQTT-сообщений → потери телеметрии.
|
||||||
|
//
|
||||||
|
// Теперь: колбэк кладёт envelope в буферизованный канал НЕблокирующе
|
||||||
|
// (переполнение — дроп со счётчиком), worker-пул отправляет в SQS с
|
||||||
|
// ограниченными ретраями и backoff.
|
||||||
|
package bridge
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"encoding/json"
|
||||||
|
"log/slog"
|
||||||
|
"sync"
|
||||||
|
"sync/atomic"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/aws/aws-sdk-go-v2/aws"
|
||||||
|
"github.com/aws/aws-sdk-go-v2/service/sqs"
|
||||||
|
)
|
||||||
|
|
||||||
|
// Константы диспатчера.
|
||||||
|
const (
|
||||||
|
// dispatcherQueueSize — буфер envelope перед SQS (~100 msg/s * 100с буфера).
|
||||||
|
dispatcherQueueSize = 10000
|
||||||
|
// dispatcherWorkers — горутин, отправляющих в SQS.
|
||||||
|
dispatcherWorkers = 4
|
||||||
|
// sendMaxAttempts — попыток SendMessage на одно сообщение.
|
||||||
|
sendMaxAttempts = 3
|
||||||
|
|
||||||
|
// maxSQSBodyBytes — потолок размера envelope для SQS.
|
||||||
|
// SQS (и shared-sqs) режут сообщения >256KB — запас на оверхед.
|
||||||
|
maxSQSBodyBytes = 250 * 1024
|
||||||
|
)
|
||||||
|
|
||||||
|
// sqsDispatcher — канал + worker-пул для отправки в SQS.
|
||||||
|
type sqsDispatcher struct {
|
||||||
|
ch chan TelemetryEnvelope
|
||||||
|
client *sqs.Client
|
||||||
|
queueURL string
|
||||||
|
log *slog.Logger
|
||||||
|
wg sync.WaitGroup
|
||||||
|
dropped atomic.Int64
|
||||||
|
}
|
||||||
|
|
||||||
|
// newSQSDispatcher создаёт диспатчер. start() запускает worker-пул.
|
||||||
|
func newSQSDispatcher(client *sqs.Client, queueURL string, log *slog.Logger) *sqsDispatcher {
|
||||||
|
return &sqsDispatcher{
|
||||||
|
ch: make(chan TelemetryEnvelope, dispatcherQueueSize),
|
||||||
|
client: client,
|
||||||
|
queueURL: queueURL,
|
||||||
|
log: log,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// start запускает worker-пул.
|
||||||
|
func (d *sqsDispatcher) start(ctx context.Context) {
|
||||||
|
for i := 0; i < dispatcherWorkers; i++ {
|
||||||
|
d.wg.Add(1)
|
||||||
|
go func(wid int) {
|
||||||
|
defer d.wg.Done()
|
||||||
|
d.worker(ctx, wid)
|
||||||
|
}(i)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// enqueue кладёт envelope в канал НЕблокирующе. false = буфер полон (дроп).
|
||||||
|
func (d *sqsDispatcher) enqueue(e TelemetryEnvelope) bool {
|
||||||
|
select {
|
||||||
|
case d.ch <- e:
|
||||||
|
return true
|
||||||
|
default:
|
||||||
|
n := d.dropped.Add(1)
|
||||||
|
if n%100 == 1 {
|
||||||
|
d.log.Error("bridge: SQS dispatch queue full, dropping",
|
||||||
|
"namespace", e.Namespace, "device", e.DeviceID,
|
||||||
|
"dropped_total", n)
|
||||||
|
}
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// closeAndWait закрывает канал и ждёт завершения worker'ов (flush).
|
||||||
|
func (d *sqsDispatcher) closeAndWait() {
|
||||||
|
close(d.ch)
|
||||||
|
d.wg.Wait()
|
||||||
|
d.log.Info("bridge: SQS dispatcher stopped", "dropped_total", d.dropped.Load())
|
||||||
|
}
|
||||||
|
|
||||||
|
// worker — цикл отправки из канала с ограниченными ретраями.
|
||||||
|
func (d *sqsDispatcher) worker(ctx context.Context, wid int) {
|
||||||
|
for e := range d.ch {
|
||||||
|
body, err := json.Marshal(e)
|
||||||
|
if err != nil {
|
||||||
|
d.log.Error("bridge: marshal envelope", "worker", wid,
|
||||||
|
"namespace", e.Namespace, "err", err)
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if len(body) > maxSQSBodyBytes {
|
||||||
|
// Лимит SQS 256KB — такие сообщения не пройдут в принципе
|
||||||
|
// (проверено 2026-08-16: InvalidParameterValue message size
|
||||||
|
// exceeds the limit). Дроп с явным логом.
|
||||||
|
d.log.Error("bridge: payload too large for SQS, dropping",
|
||||||
|
"worker", wid, "namespace", e.Namespace,
|
||||||
|
"device", e.DeviceID, "bytes", len(body),
|
||||||
|
"limit_bytes", maxSQSBodyBytes)
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if err := d.sendWithRetry(ctx, string(body)); err != nil {
|
||||||
|
d.log.Error("bridge: SQS send failed after retries, dropping",
|
||||||
|
"worker", wid, "namespace", e.Namespace,
|
||||||
|
"device", e.DeviceID, "err", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// sendWithRetry — SendMessage с backoff 1с/2с/4с.
|
||||||
|
func (d *sqsDispatcher) sendWithRetry(ctx context.Context, body string) error {
|
||||||
|
var lastErr error
|
||||||
|
backoff := time.Second
|
||||||
|
for attempt := 1; attempt <= sendMaxAttempts; attempt++ {
|
||||||
|
_, err := d.client.SendMessage(ctx, &sqs.SendMessageInput{
|
||||||
|
QueueUrl: aws.String(d.queueURL),
|
||||||
|
MessageBody: aws.String(body),
|
||||||
|
})
|
||||||
|
if err == nil {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
lastErr = err
|
||||||
|
if ctx.Err() != nil {
|
||||||
|
return lastErr
|
||||||
|
}
|
||||||
|
if attempt < sendMaxAttempts {
|
||||||
|
select {
|
||||||
|
case <-ctx.Done():
|
||||||
|
return lastErr
|
||||||
|
case <-time.After(backoff):
|
||||||
|
}
|
||||||
|
backoff *= 2
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return lastErr
|
||||||
|
}
|
||||||
@@ -43,6 +43,9 @@ type Config struct {
|
|||||||
|
|
||||||
// Безопасность
|
// Безопасность
|
||||||
AuthTestMode bool
|
AuthTestMode bool
|
||||||
|
// JwtHMACSecret — если задан, JWT проверяется по HS256 (HMAC-SHA256).
|
||||||
|
// Пусто — структурная проверка (как раньше, периметр = платформа).
|
||||||
|
JwtHMACSecret string
|
||||||
|
|
||||||
// Логирование
|
// Логирование
|
||||||
LogLevel slog.Level
|
LogLevel slog.Level
|
||||||
@@ -108,12 +111,13 @@ func Load() (*Config, error) {
|
|||||||
SQSQueueName: getEnv("SQS_QUEUE_NAME", DefaultSQSQueueName),
|
SQSQueueName: getEnv("SQS_QUEUE_NAME", DefaultSQSQueueName),
|
||||||
SQSRegion: getEnv("SQS_REGION", DefaultSQSRegion),
|
SQSRegion: getEnv("SQS_REGION", DefaultSQSRegion),
|
||||||
SQSLongPollSeconds: getEnvInt("SQS_LONG_POLL_SECONDS", 20),
|
SQSLongPollSeconds: getEnvInt("SQS_LONG_POLL_SECONDS", 20),
|
||||||
SQSVisibilityTimeout: getEnvInt("SQS_VISIBILITY_TIMEOUT", 30),
|
SQSVisibilityTimeout: getEnvInt("SQS_VISIBILITY_TIMEOUT", 120),
|
||||||
MQTTUsername: os.Getenv("MQTT_USERNAME"),
|
MQTTUsername: os.Getenv("MQTT_USERNAME"),
|
||||||
MQTTPassword: os.Getenv("MQTT_PASSWORD"),
|
MQTTPassword: os.Getenv("MQTT_PASSWORD"),
|
||||||
MQTTClientID: getEnv("MQTT_CLIENT_ID", DefaultMQTTClientID),
|
MQTTClientID: getEnv("MQTT_CLIENT_ID", DefaultMQTTClientID),
|
||||||
AdminStatsToken: os.Getenv("ADMIN_STATS_TOKEN"),
|
AdminStatsToken: os.Getenv("ADMIN_STATS_TOKEN"),
|
||||||
AuthTestMode: getEnvBool("AUTH_TEST_MODE", false),
|
AuthTestMode: getEnvBool("AUTH_TEST_MODE", false),
|
||||||
|
JwtHMACSecret: os.Getenv("JWT_HMAC_SECRET"),
|
||||||
LogLevel: parseLogLevel(os.Getenv("LOG_LEVEL"), slog.LevelInfo),
|
LogLevel: parseLogLevel(os.Getenv("LOG_LEVEL"), slog.LevelInfo),
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -124,6 +128,9 @@ func Load() (*Config, error) {
|
|||||||
port := getEnv("MQTT_PORT", DefaultMQTTPort)
|
port := getEnv("MQTT_PORT", DefaultMQTTPort)
|
||||||
path := getEnv("MQTT_WS_PATH", DefaultMQTTWSPath)
|
path := getEnv("MQTT_WS_PATH", DefaultMQTTWSPath)
|
||||||
cfg.MQTTBrokerURL = "ws://" + host + ":" + port + path
|
cfg.MQTTBrokerURL = "ws://" + host + ":" + port + path
|
||||||
|
} else if !strings.HasPrefix(cfg.MQTTBrokerURL, "ws://") &&
|
||||||
|
!strings.HasPrefix(cfg.MQTTBrokerURL, "wss://") {
|
||||||
|
return nil, fmt.Errorf("MQTT_BROKER_URL must start with ws:// or wss://, got %q", cfg.MQTTBrokerURL)
|
||||||
}
|
}
|
||||||
|
|
||||||
var missing []string
|
var missing []string
|
||||||
|
|||||||
@@ -20,6 +20,9 @@ import (
|
|||||||
// errorBackoff — пауза между попытками при ошибках SQS.
|
// errorBackoff — пауза между попытками при ошибках SQS.
|
||||||
const errorBackoff = 5 * time.Second
|
const errorBackoff = 5 * time.Second
|
||||||
|
|
||||||
|
// maxProcessBackoff — потолок backoff после ошибок обработки (PG).
|
||||||
|
const maxProcessBackoff = 60 * time.Second
|
||||||
|
|
||||||
// Run — бесконечный long-poll цикл потребителя.
|
// Run — бесконечный long-poll цикл потребителя.
|
||||||
func Run(ctx context.Context, cfg *config.Config, sqsClient *sqs.Client, store *iotpg.IoTPostgresStore, log *slog.Logger) error {
|
func Run(ctx context.Context, cfg *config.Config, sqsClient *sqs.Client, store *iotpg.IoTPostgresStore, log *slog.Logger) error {
|
||||||
queueURL, err := sqsclient.ResolveQueueURL(ctx, sqsClient, cfg.SQSQueueName)
|
queueURL, err := sqsclient.ResolveQueueURL(ctx, sqsClient, cfg.SQSQueueName)
|
||||||
@@ -50,15 +53,31 @@ func Run(ctx context.Context, cfg *config.Config, sqsClient *sqs.Client, store *
|
|||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
|
||||||
|
processBackoff := time.Second
|
||||||
for _, msg := range resp.Messages {
|
for _, msg := range resp.Messages {
|
||||||
if msg.Body == nil {
|
if msg.Body == nil || *msg.Body == "" {
|
||||||
|
log.Warn("consumer: empty message body, deleting",
|
||||||
|
"message_id", aws.ToString(msg.MessageId))
|
||||||
|
if err := deleteMessage(ctx, sqsClient, queueURL, msg.ReceiptHandle, log); err != nil {
|
||||||
|
log.Error("consumer: DeleteMessage failed", "err", err, "message_id", aws.ToString(msg.MessageId))
|
||||||
|
}
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
if err := processTelemetry(ctx, *msg.Body, store, log); err != nil {
|
if err := processTelemetry(ctx, *msg.Body, store, log); err != nil {
|
||||||
log.Error("consumer: process telemetry", "err", err, "message_id", aws.ToString(msg.MessageId))
|
log.Error("consumer: process telemetry", "err", err, "message_id", aws.ToString(msg.MessageId))
|
||||||
// Сообщение вернётся в очередь после visibility timeout.
|
// Сообщение вернётся в очередь после visibility timeout.
|
||||||
|
// Backoff с jitter — не долбить упавший PG в цикле
|
||||||
|
// (фикс HIGH из ревью 2026-08-16: лавина ретраев).
|
||||||
|
if !sleepCtx(ctx, processBackoff) {
|
||||||
|
break
|
||||||
|
}
|
||||||
|
processBackoff *= 2
|
||||||
|
if processBackoff > maxProcessBackoff {
|
||||||
|
processBackoff = maxProcessBackoff
|
||||||
|
}
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
processBackoff = time.Second
|
||||||
if err := deleteMessage(ctx, sqsClient, queueURL, msg.ReceiptHandle, log); err != nil {
|
if err := deleteMessage(ctx, sqsClient, queueURL, msg.ReceiptHandle, log); err != nil {
|
||||||
log.Error("consumer: DeleteMessage failed", "err", err, "message_id", aws.ToString(msg.MessageId))
|
log.Error("consumer: DeleteMessage failed", "err", err, "message_id", aws.ToString(msg.MessageId))
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -29,12 +29,22 @@ VALUES ($1, $2, $3, $4, $5, $6)`,
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// List возвращает устройства namespace (без сортировки по паролям — пароли на месте, но API их не отдаёт).
|
// List возвращает устройства namespace (без паролей — API их не отдаёт).
|
||||||
func (s *DeviceStore) List(ctx context.Context, namespace string) ([]Device, error) {
|
// limit <= 0 — без ограничения; offset — смещение (пагинация API).
|
||||||
rows, err := s.db.QueryContext(ctx, `
|
func (s *DeviceStore) List(ctx context.Context, namespace string, limit, offset int) ([]Device, error) {
|
||||||
SELECT namespace, name, device_id, enabled, mqtt_password, metadata, phase,
|
query := `SELECT namespace, name, device_id, enabled, mqtt_password, metadata, phase,
|
||||||
last_connected, created_at
|
last_connected, created_at
|
||||||
FROM iot_devices WHERE namespace = $1 ORDER BY name`, namespace)
|
FROM iot_devices WHERE namespace = $1 ORDER BY name`
|
||||||
|
args := []any{namespace}
|
||||||
|
if limit > 0 {
|
||||||
|
query += fmt.Sprintf(" LIMIT $%d", len(args)+1)
|
||||||
|
args = append(args, limit)
|
||||||
|
}
|
||||||
|
if offset > 0 {
|
||||||
|
query += fmt.Sprintf(" OFFSET $%d", len(args)+1)
|
||||||
|
args = append(args, offset)
|
||||||
|
}
|
||||||
|
rows, err := s.db.QueryContext(ctx, query, args...)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("store: list devices: %w", err)
|
return nil, fmt.Errorf("store: list devices: %w", err)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -4,6 +4,8 @@ import (
|
|||||||
"context"
|
"context"
|
||||||
"database/sql"
|
"database/sql"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"os"
|
||||||
|
"strconv"
|
||||||
"time"
|
"time"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -39,7 +41,15 @@ func Open(ctx context.Context, dsn string) (*DeviceStore, error) {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("store: open DB: %w", err)
|
return nil, fmt.Errorf("store: open DB: %w", err)
|
||||||
}
|
}
|
||||||
db.SetMaxOpenConns(5)
|
// Пул устройств: дефолт 20 (ревью 2026-08-16: 5 мало для API+MQTT auth),
|
||||||
|
// переопределяется IOT_PG_MAX_CONNS.
|
||||||
|
maxConns := 20
|
||||||
|
if v := os.Getenv("IOT_PG_MAX_CONNS"); v != "" {
|
||||||
|
if n, convErr := strconv.Atoi(v); convErr == nil && n > 0 {
|
||||||
|
maxConns = n
|
||||||
|
}
|
||||||
|
}
|
||||||
|
db.SetMaxOpenConns(maxConns)
|
||||||
db.SetMaxIdleConns(2)
|
db.SetMaxIdleConns(2)
|
||||||
db.SetConnMaxLifetime(5 * time.Minute)
|
db.SetConnMaxLifetime(5 * time.Minute)
|
||||||
|
|
||||||
|
|||||||
@@ -37,9 +37,29 @@ type IoTPostgresStore struct {
|
|||||||
adminDB *sql.DB
|
adminDB *sql.DB
|
||||||
adminDSN string
|
adminDSN string
|
||||||
tenants sync.Map
|
tenants sync.Map
|
||||||
|
tenantMu sync.Map // namespace → *sync.Mutex (дедупликация открытия)
|
||||||
|
ensured sync.Map // namespace → struct{} (EnsureTenantDB уже выполнен)
|
||||||
log *slog.Logger
|
log *slog.Logger
|
||||||
|
|
||||||
|
// батчинг вставок (фикс MEDIUM ревью 2026-08-16: 1000 msg/s = 1000 INSERT)
|
||||||
|
mu sync.Mutex
|
||||||
|
batches map[string][]telemetryInsert
|
||||||
|
stopFlush chan struct{}
|
||||||
|
flushWG sync.WaitGroup
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// telemetryInsert — строка для батч-вставки.
|
||||||
|
type telemetryInsert struct {
|
||||||
|
deviceID string
|
||||||
|
payload []byte
|
||||||
|
}
|
||||||
|
|
||||||
|
// telemetryBatchSize — размер батча перед синхронным flush.
|
||||||
|
const telemetryBatchSize = 100
|
||||||
|
|
||||||
|
// telemetryFlushInterval — период фонового flush неполных батчей.
|
||||||
|
const telemetryFlushInterval = 200 * time.Millisecond
|
||||||
|
|
||||||
// TelemetryRow — одна запись телеметрии из таблицы iot_telemetry.
|
// TelemetryRow — одна запись телеметрии из таблицы iot_telemetry.
|
||||||
type TelemetryRow struct {
|
type TelemetryRow struct {
|
||||||
ID int64 `json:"id"`
|
ID int64 `json:"id"`
|
||||||
@@ -60,19 +80,90 @@ func New(adminDSN string, log *slog.Logger) (*IoTPostgresStore, error) {
|
|||||||
db.Close()
|
db.Close()
|
||||||
return nil, fmt.Errorf("iotpg: ping admin DB: %w", err)
|
return nil, fmt.Errorf("iotpg: ping admin DB: %w", err)
|
||||||
}
|
}
|
||||||
db.SetMaxOpenConns(5)
|
db.SetMaxOpenConns(10)
|
||||||
db.SetMaxIdleConns(2)
|
db.SetMaxIdleConns(2)
|
||||||
db.SetConnMaxLifetime(5 * time.Minute)
|
db.SetConnMaxLifetime(5 * time.Minute)
|
||||||
|
|
||||||
store := &IoTPostgresStore{adminDB: db, adminDSN: adminDSN, log: log}
|
store := &IoTPostgresStore{
|
||||||
|
adminDB: db,
|
||||||
|
adminDSN: adminDSN,
|
||||||
|
batches: make(map[string][]telemetryInsert),
|
||||||
|
stopFlush: make(chan struct{}),
|
||||||
|
log: log,
|
||||||
|
}
|
||||||
if err := store.initManagementSchema(ctx); err != nil {
|
if err := store.initManagementSchema(ctx); err != nil {
|
||||||
db.Close()
|
db.Close()
|
||||||
return nil, fmt.Errorf("iotpg: init management schema: %w", err)
|
return nil, fmt.Errorf("iotpg: init management schema: %w", err)
|
||||||
}
|
}
|
||||||
|
store.startFlusher()
|
||||||
log.Info("iotpg: connected to IoT Postgres management DB")
|
log.Info("iotpg: connected to IoT Postgres management DB")
|
||||||
return store, nil
|
return store, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// startFlusher — фоновая горутина периодического flush неполных батчей.
|
||||||
|
func (s *IoTPostgresStore) startFlusher() {
|
||||||
|
s.flushWG.Add(1)
|
||||||
|
go func() {
|
||||||
|
defer s.flushWG.Done()
|
||||||
|
t := time.NewTicker(telemetryFlushInterval)
|
||||||
|
defer t.Stop()
|
||||||
|
for {
|
||||||
|
select {
|
||||||
|
case <-s.stopFlush:
|
||||||
|
return
|
||||||
|
case <-t.C:
|
||||||
|
s.flushAll(false)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
}
|
||||||
|
|
||||||
|
// flushAll — flush всех накопленных батчей. reportErr=false: только лог.
|
||||||
|
func (s *IoTPostgresStore) flushAll(reportErr bool) error {
|
||||||
|
s.mu.Lock()
|
||||||
|
if len(s.batches) == 0 {
|
||||||
|
s.mu.Unlock()
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
batches := s.batches
|
||||||
|
s.batches = make(map[string][]telemetryInsert)
|
||||||
|
s.mu.Unlock()
|
||||||
|
|
||||||
|
var firstErr error
|
||||||
|
for ns, rows := range batches {
|
||||||
|
if err := s.flushTenant(ns, rows); err != nil {
|
||||||
|
s.log.Error("iotpg: flush batch", "namespace", ns, "rows", len(rows), "err", err)
|
||||||
|
if firstErr == nil {
|
||||||
|
firstErr = err
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return firstErr
|
||||||
|
}
|
||||||
|
|
||||||
|
// flushTenant — INSERT батча строк одного тенанта (multi-VALUES).
|
||||||
|
func (s *IoTPostgresStore) flushTenant(ns string, rows []telemetryInsert) error {
|
||||||
|
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
|
||||||
|
defer cancel()
|
||||||
|
tenantDB, err := s.getTenantDB(ctx, ns)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
// INSERT INTO iot_telemetry (device_id, payload) VALUES ($1,$2),($3,$4),...
|
||||||
|
var sb strings.Builder
|
||||||
|
sb.WriteString("INSERT INTO iot_telemetry (device_id, payload) VALUES ")
|
||||||
|
args := make([]any, 0, len(rows)*2)
|
||||||
|
for i, r := range rows {
|
||||||
|
if i > 0 {
|
||||||
|
sb.WriteString(",")
|
||||||
|
}
|
||||||
|
sb.WriteString(fmt.Sprintf("($%d,$%d)", i*2+1, i*2+2))
|
||||||
|
args = append(args, r.deviceID, r.payload)
|
||||||
|
}
|
||||||
|
_, err = tenantDB.ExecContext(ctx, sb.String(), args...)
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
// NewFromEnv создаёт store из env var IOT_PG_DSN.
|
// NewFromEnv создаёт store из env var IOT_PG_DSN.
|
||||||
// Возвращает (nil, nil) если переменная не задана — IoT Postgres опционален.
|
// Возвращает (nil, nil) если переменная не задана — IoT Postgres опционален.
|
||||||
func NewFromEnv(log *slog.Logger) (*IoTPostgresStore, error) {
|
func NewFromEnv(log *slog.Logger) (*IoTPostgresStore, error) {
|
||||||
@@ -97,9 +188,12 @@ created_at TIMESTAMPTZ DEFAULT now()
|
|||||||
}
|
}
|
||||||
|
|
||||||
// EnsureTenantDB создаёт DATABASE, USER и таблицу iot_telemetry для namespace.
|
// EnsureTenantDB создаёт DATABASE, USER и таблицу iot_telemetry для namespace.
|
||||||
// Идемпотентен — повторный вызов безопасен.
|
// Идемпотентен — повторный вызов безопасен (кэш ensured пропускает проверки).
|
||||||
// Вызывается mqtt-bridge при первом сообщении от нового tenant.
|
|
||||||
func (s *IoTPostgresStore) EnsureTenantDB(ctx context.Context, namespace string) error {
|
func (s *IoTPostgresStore) EnsureTenantDB(ctx context.Context, namespace string) error {
|
||||||
|
if _, ok := s.ensured.Load(namespace); ok {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
dbName := tenantDBName(namespace)
|
dbName := tenantDBName(namespace)
|
||||||
userName := dbName
|
userName := dbName
|
||||||
|
|
||||||
@@ -114,28 +208,39 @@ func (s *IoTPostgresStore) EnsureTenantDB(ctx context.Context, namespace string)
|
|||||||
if !exists {
|
if !exists {
|
||||||
password := uuid.New().String()
|
password := uuid.New().String()
|
||||||
|
|
||||||
// CREATE USER через DO block — pg не поддерживает CREATE USER IF NOT EXISTS
|
// Роль создаём только если её нет (DO-блоки НЕ принимают параметры —
|
||||||
_, err = s.adminDB.ExecContext(ctx, fmt.Sprintf(
|
// прецедент 2026-08-16: "got 2 parameters but the statement requires 0").
|
||||||
`DO $$ BEGIN
|
var roleExists bool
|
||||||
IF NOT EXISTS (SELECT FROM pg_roles WHERE rolname = '%s') THEN
|
err = s.adminDB.QueryRowContext(ctx,
|
||||||
CREATE USER %s WITH PASSWORD '%s';
|
`SELECT EXISTS(SELECT 1 FROM pg_roles WHERE rolname = $1)`, userName,
|
||||||
END IF;
|
).Scan(&roleExists)
|
||||||
END $$`, userName, userName, password,
|
|
||||||
))
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
return fmt.Errorf("iotpg: check role %s: %w", userName, err)
|
||||||
|
}
|
||||||
|
if !roleExists {
|
||||||
|
_, err = s.adminDB.ExecContext(ctx,
|
||||||
|
`CREATE USER `+pq.QuoteIdentifier(userName)+` WITH PASSWORD `+pq.QuoteLiteral(password),
|
||||||
|
)
|
||||||
|
if err != nil {
|
||||||
|
var pqErr *pq.Error
|
||||||
|
if errors.As(err, &pqErr) && pqErr.Code == "42710" { // duplicate_object
|
||||||
|
// гонка: роль создал параллельный вызов — ок
|
||||||
|
} else {
|
||||||
return fmt.Errorf("iotpg: create user %s: %w", userName, err)
|
return fmt.Errorf("iotpg: create user %s: %w", userName, err)
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// PG15+: GRANT role TO current_user перед CREATE DATABASE ... OWNER
|
// PG15+: GRANT role TO current_user перед CREATE DATABASE ... OWNER
|
||||||
if _, err = s.adminDB.ExecContext(ctx,
|
if _, err = s.adminDB.ExecContext(ctx,
|
||||||
fmt.Sprintf(`GRANT %s TO CURRENT_USER`, userName),
|
`GRANT `+pq.QuoteIdentifier(userName)+` TO CURRENT_USER`,
|
||||||
); err != nil {
|
); err != nil {
|
||||||
return fmt.Errorf("iotpg: grant role %s: %w", userName, err)
|
return fmt.Errorf("iotpg: grant role %s: %w", userName, err)
|
||||||
}
|
}
|
||||||
|
|
||||||
// CREATE DATABASE нельзя в транзакции
|
// CREATE DATABASE нельзя в транзакции
|
||||||
if _, err = s.adminDB.ExecContext(ctx,
|
if _, err = s.adminDB.ExecContext(ctx,
|
||||||
fmt.Sprintf(`CREATE DATABASE %s OWNER %s`, dbName, userName),
|
`CREATE DATABASE `+pq.QuoteIdentifier(dbName)+` OWNER `+pq.QuoteIdentifier(userName),
|
||||||
); err != nil {
|
); err != nil {
|
||||||
return fmt.Errorf("iotpg: create database %s: %w", dbName, err)
|
return fmt.Errorf("iotpg: create database %s: %w", dbName, err)
|
||||||
}
|
}
|
||||||
@@ -165,32 +270,50 @@ payload JSONB NOT NULL
|
|||||||
CREATE INDEX IF NOT EXISTS idx_iot_telemetry_device_ts
|
CREATE INDEX IF NOT EXISTS idx_iot_telemetry_device_ts
|
||||||
ON iot_telemetry (device_id, ts DESC);
|
ON iot_telemetry (device_id, ts DESC);
|
||||||
`)
|
`)
|
||||||
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
s.ensured.Store(namespace, struct{}{})
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
// InsertTelemetry записывает строку телеметрии в tenant DB.
|
// InsertTelemetry ставит строку в батч-буфер тенанта.
|
||||||
|
// При накоплении telemetryBatchSize строк батч пишется синхронно (ошибка
|
||||||
|
// возвращается вызывающему); неполные батчи дописывает фоновый flusher.
|
||||||
func (s *IoTPostgresStore) InsertTelemetry(ctx context.Context, namespace, deviceID string, payload json.RawMessage) error {
|
func (s *IoTPostgresStore) InsertTelemetry(ctx context.Context, namespace, deviceID string, payload json.RawMessage) error {
|
||||||
tenantDB, err := s.getTenantDB(ctx, namespace)
|
s.mu.Lock()
|
||||||
if err != nil {
|
s.batches[namespace] = append(s.batches[namespace], telemetryInsert{
|
||||||
return fmt.Errorf("iotpg: get tenant DB for insert: %w", err)
|
deviceID: deviceID,
|
||||||
|
payload: append([]byte(nil), payload...),
|
||||||
|
})
|
||||||
|
full := len(s.batches[namespace]) >= telemetryBatchSize
|
||||||
|
var rows []telemetryInsert
|
||||||
|
if full {
|
||||||
|
rows = s.batches[namespace]
|
||||||
|
delete(s.batches, namespace)
|
||||||
}
|
}
|
||||||
_, err = tenantDB.ExecContext(ctx,
|
s.mu.Unlock()
|
||||||
`INSERT INTO iot_telemetry (device_id, payload) VALUES ($1, $2)`,
|
|
||||||
deviceID, []byte(payload),
|
if full {
|
||||||
)
|
return s.flushTenant(namespace, rows)
|
||||||
return err
|
}
|
||||||
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// QueryTelemetry читает телеметрию из tenant DB (ts DESC).
|
// QueryTelemetry читает телеметрию из tenant DB (ts DESC).
|
||||||
// deviceID — фильтр (пустая строка = все устройства). limit — max записей (50..1000).
|
// deviceID — фильтр (пустая строка = все устройства). limit — max записей
|
||||||
// Если tenant DB не существует (данных ещё нет) — возвращает пустой срез без ошибки.
|
// (50..1000), offset — смещение для пагинации.
|
||||||
func (s *IoTPostgresStore) QueryTelemetry(ctx context.Context, namespace, deviceID string, limit int) ([]TelemetryRow, error) {
|
// Если tenant DB не существует (данных ещё нет) — пустой срез без ошибки.
|
||||||
|
func (s *IoTPostgresStore) QueryTelemetry(ctx context.Context, namespace, deviceID string, limit, offset int) ([]TelemetryRow, error) {
|
||||||
if limit <= 0 {
|
if limit <= 0 {
|
||||||
limit = 50
|
limit = 50
|
||||||
}
|
}
|
||||||
if limit > 1000 {
|
if limit > 1000 {
|
||||||
limit = 1000
|
limit = 1000
|
||||||
}
|
}
|
||||||
|
if offset < 0 {
|
||||||
|
offset = 0
|
||||||
|
}
|
||||||
tenantDB, err := s.getTenantDB(ctx, namespace)
|
tenantDB, err := s.getTenantDB(ctx, namespace)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
// Если DB не существует — тенант ещё не отправлял данные, это нормально
|
// Если DB не существует — тенант ещё не отправлял данные, это нормально
|
||||||
@@ -204,14 +327,14 @@ func (s *IoTPostgresStore) QueryTelemetry(ctx context.Context, namespace, device
|
|||||||
if deviceID != "" {
|
if deviceID != "" {
|
||||||
rows, err = tenantDB.QueryContext(ctx,
|
rows, err = tenantDB.QueryContext(ctx,
|
||||||
`SELECT id, device_id, ts, payload FROM iot_telemetry
|
`SELECT id, device_id, ts, payload FROM iot_telemetry
|
||||||
WHERE device_id = $1 ORDER BY ts DESC LIMIT $2`,
|
WHERE device_id = $1 ORDER BY ts DESC LIMIT $2 OFFSET $3`,
|
||||||
deviceID, limit,
|
deviceID, limit, offset,
|
||||||
)
|
)
|
||||||
} else {
|
} else {
|
||||||
rows, err = tenantDB.QueryContext(ctx,
|
rows, err = tenantDB.QueryContext(ctx,
|
||||||
`SELECT id, device_id, ts, payload FROM iot_telemetry
|
`SELECT id, device_id, ts, payload FROM iot_telemetry
|
||||||
ORDER BY ts DESC LIMIT $1`,
|
ORDER BY ts DESC LIMIT $1 OFFSET $2`,
|
||||||
limit,
|
limit, offset,
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -308,7 +431,6 @@ FROM iot_telemetry`).Scan(&stats.Total, &stats.Last1h, &stats.Last24h)
|
|||||||
latestRows, err := tenantDB.QueryContext(ctx,
|
latestRows, err := tenantDB.QueryContext(ctx,
|
||||||
`SELECT id, device_id, ts, payload FROM iot_telemetry ORDER BY ts DESC LIMIT 5`)
|
`SELECT id, device_id, ts, payload FROM iot_telemetry ORDER BY ts DESC LIMIT 5`)
|
||||||
if err == nil {
|
if err == nil {
|
||||||
defer latestRows.Close()
|
|
||||||
for latestRows.Next() {
|
for latestRows.Next() {
|
||||||
var r TelemetryRow
|
var r TelemetryRow
|
||||||
var rawPayload []byte
|
var rawPayload []byte
|
||||||
@@ -317,6 +439,7 @@ FROM iot_telemetry`).Scan(&stats.Total, &stats.Last1h, &stats.Last24h)
|
|||||||
stats.Latest = append(stats.Latest, r)
|
stats.Latest = append(stats.Latest, r)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
latestRows.Close() // закрываем сразу, не defer в цикле
|
||||||
}
|
}
|
||||||
|
|
||||||
result.Tenants = append(result.Tenants, stats)
|
result.Tenants = append(result.Tenants, stats)
|
||||||
@@ -336,8 +459,11 @@ func isDBNotExistErr(err error) bool {
|
|||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
|
|
||||||
// Close закрывает все подключения (admin + tenant кэш).
|
// Close останавливает flusher, дописывает остаток и закрывает все подключения.
|
||||||
func (s *IoTPostgresStore) Close() error {
|
func (s *IoTPostgresStore) Close() error {
|
||||||
|
close(s.stopFlush)
|
||||||
|
s.flushWG.Wait()
|
||||||
|
_ = s.flushAll(true)
|
||||||
s.tenants.Range(func(_, value any) bool {
|
s.tenants.Range(func(_, value any) bool {
|
||||||
if db, ok := value.(*sql.DB); ok {
|
if db, ok := value.(*sql.DB); ok {
|
||||||
db.Close()
|
db.Close()
|
||||||
@@ -348,7 +474,16 @@ func (s *IoTPostgresStore) Close() error {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// getTenantDB возвращает *sql.DB для tenant DB из кэша или открывает новый.
|
// getTenantDB возвращает *sql.DB для tenant DB из кэша или открывает новый.
|
||||||
|
// Открытие дедуплицируется per-namespace мьютексом (гонка Load→LoadOrStore
|
||||||
|
// из ревью 2026-08-16).
|
||||||
func (s *IoTPostgresStore) getTenantDB(ctx context.Context, namespace string) (*sql.DB, error) {
|
func (s *IoTPostgresStore) getTenantDB(ctx context.Context, namespace string) (*sql.DB, error) {
|
||||||
|
if cached, ok := s.tenants.Load(namespace); ok {
|
||||||
|
return cached.(*sql.DB), nil
|
||||||
|
}
|
||||||
|
m, _ := s.tenantMu.LoadOrStore(namespace, &sync.Mutex{})
|
||||||
|
mu := m.(*sync.Mutex)
|
||||||
|
mu.Lock()
|
||||||
|
defer mu.Unlock()
|
||||||
if cached, ok := s.tenants.Load(namespace); ok {
|
if cached, ok := s.tenants.Load(namespace); ok {
|
||||||
return cached.(*sql.DB), nil
|
return cached.(*sql.DB), nil
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,99 @@
|
|||||||
|
# Нагрузочные тесты IoT
|
||||||
|
|
||||||
|
Контур: устройство → wss (EMQX) → бридж → SQS `iot-telemetry` → consumer → PG → API.
|
||||||
|
|
||||||
|
## Запуск
|
||||||
|
|
||||||
|
```bash
|
||||||
|
cd loadtests
|
||||||
|
python3 -m loadtests.run <scenario> [опции]
|
||||||
|
# или из корня репо:
|
||||||
|
python3 -m loadtests.run <scenario> [опции]
|
||||||
|
```
|
||||||
|
|
||||||
|
Требования: 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_<run-id>` и удаляет их после прогона при `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
|
||||||
|
```
|
||||||
|
|
||||||
|
## Ограничения
|
||||||
|
|
||||||
|
- **SQS-лимит 256KB** на сообщение (shared-sqs отклоняет: `InvalidParameterValue
|
||||||
|
message size exceeds the limit`) → payload устройств ≤ ~250KB с учётом
|
||||||
|
envelope. Бридж дропает oversize с логом `payload too large for SQS`.
|
||||||
|
EMQX пропускает до 1MB, но узким местом конвейера является SQS.
|
||||||
|
- API отдаёт максимум **1000 строк** телеметрии на запрос; шлюз платформы
|
||||||
|
рвёт ответы >~15КБ (баг MSS/MTU, тикет Nubes) → верификатор ходит
|
||||||
|
страницами по 50 (`offset`). Большие страницы НЕ использовать.
|
||||||
|
- Задержка = `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 <pg-ns> 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.
|
||||||
@@ -0,0 +1 @@
|
|||||||
|
# loadtests — набор нагрузочных тестов контура IoT.
|
||||||
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
@@ -0,0 +1,139 @@
|
|||||||
|
# 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
|
||||||
|
# Внешний путь платформы даёт таймауты ~31-33с (~5.5% запросов) —
|
||||||
|
# ретраи 3 с паузами (см. doc/sqs-integration.md, тикет Nubes).
|
||||||
|
last = None
|
||||||
|
for attempt in range(3):
|
||||||
|
try:
|
||||||
|
req = urllib.request.Request(
|
||||||
|
url, data=data, method=method,
|
||||||
|
headers={"Authorization": "Bearer " + self.token,
|
||||||
|
"Content-Type": "application/json"},
|
||||||
|
)
|
||||||
|
with urllib.request.urlopen(req, timeout=15, context=self.ctx) as r:
|
||||||
|
raw = r.read().decode()
|
||||||
|
return json.loads(raw) if raw.strip() else {}
|
||||||
|
except Exception as e:
|
||||||
|
last = e
|
||||||
|
time.sleep(1 + attempt)
|
||||||
|
raise last
|
||||||
|
|
||||||
|
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, offset=0):
|
||||||
|
q = f"?limit={limit}&offset={offset}"
|
||||||
|
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)
|
||||||
@@ -0,0 +1,147 @@
|
|||||||
|
# publisher.py — MQTT-паблишер: один поток на устройство.
|
||||||
|
#
|
||||||
|
# Каждый DevicePublisher:
|
||||||
|
# - подключается по wss (subprotocol mqtt), проверяет CONNACK через on_connect;
|
||||||
|
# - публикует в <ns>/telemetry/<device_id> с заданной частотой;
|
||||||
|
# - нумерует сообщения (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
|
||||||
@@ -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()
|
||||||
@@ -0,0 +1,434 @@
|
|||||||
|
# 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)
|
||||||
|
# Фазы различаются суффиксом run_id — верифицируем фазу "-h" отдельно.
|
||||||
|
high_sent = high.get("sent", 0)
|
||||||
|
expected_per_device = max(1, high_sent // max(1, len(devices)))
|
||||||
|
rep_high = Verifier(api, ns, run_id + "-h").run_report(
|
||||||
|
devices, expected_per_device)
|
||||||
|
# recovery: сколько времени после конца всплеска очередь отдала всё
|
||||||
|
_, recovery_s = _wait_delivery(
|
||||||
|
api, ns, run_id + "-h", devices, high_sent, timeout=300)
|
||||||
|
return {
|
||||||
|
"scenario": "burst",
|
||||||
|
"ns": ns,
|
||||||
|
"quiet_phase": low,
|
||||||
|
"spike_phase": high,
|
||||||
|
"tail_phase": tail,
|
||||||
|
"sent_high_total": high_sent,
|
||||||
|
"recovery_sec": round(recovery_s, 1),
|
||||||
|
"high": rep_high,
|
||||||
|
}
|
||||||
|
finally:
|
||||||
|
_cleanup(api, ns, devices)
|
||||||
|
|
||||||
|
|
||||||
|
# --------------------------------------------------------------------------
|
||||||
|
# 3. LARGE_PAYLOAD — крупные сообщения
|
||||||
|
# Верификация через API невозможна: строка с большим payload в ответе >15 КБ,
|
||||||
|
# а шлюз платформы рвёт такие ответы (баг MSS/MTU, тикет Nubes). Доставку
|
||||||
|
# проверять по логам монолита: 'consumer: telemetry saved'.
|
||||||
|
# --------------------------------------------------------------------------
|
||||||
|
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:
|
||||||
|
counters = _run_publishers(
|
||||||
|
run_id, ns, devices, rate=args.rate, duration=args.duration,
|
||||||
|
qos=args.qos, payload_size=args.payload_size)
|
||||||
|
return {
|
||||||
|
"scenario": "large_payload",
|
||||||
|
"ns": ns,
|
||||||
|
"payload_size": args.payload_size,
|
||||||
|
"publisher": counters,
|
||||||
|
"note": "доставку смотреть в логах монолита (telemetry saved); "
|
||||||
|
"API-верификация невозможна из-за бага шлюза (>15KB ответ)",
|
||||||
|
}
|
||||||
|
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)
|
||||||
@@ -0,0 +1,102 @@
|
|||||||
|
# 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 ждали.
|
||||||
|
Пагинация по 50 строк: шлюз платформы рвёт ответы >~15 КБ
|
||||||
|
(баг MSS/MTU, тикет Nubes) — большие страницы НЕ использовать."""
|
||||||
|
rows = []
|
||||||
|
offset = 0
|
||||||
|
pages = 0
|
||||||
|
while pages < 40:
|
||||||
|
page = self.api.telemetry(self.ns, device_id=device_id,
|
||||||
|
limit=50, offset=offset)
|
||||||
|
items = page.get("items", [])
|
||||||
|
rows.extend(items)
|
||||||
|
pages += 1
|
||||||
|
offset += 50
|
||||||
|
if len(items) < 50 or len(rows) >= max(expected + 100, 200):
|
||||||
|
break
|
||||||
|
seqs = {}
|
||||||
|
for r in rows:
|
||||||
|
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,
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user