diff --git a/doc/progress.md b/doc/progress.md index b895430..8112521 100644 --- a/doc/progress.md +++ b/doc/progress.md @@ -1,6 +1,35 @@ # Прогресс разработки -Последнее обновление: 2026-04-06 17:30 МСК +Последнее обновление: 2026-04-06 19:00 МСК + +--- + +## 2026-04-06 (вечер) — Полное суровое тестирование IoT pipeline + +### Тест-матрица (7 сценариев, baseline: 18 строк) + +| # | Тест | Результат | Примечание | +|---|------|-----------|-----------| +| 1 | Cold start (все IoT поды сразу) | ✅ PASS | Consumer: 16 retry за 48с до Kafka ready | +| 2 | Restart resilience 3× | ✅ PASS | <1с при уже работающей Kafka | +| 3 | Load 100 сообщений burst | ⚠️ PARTIAL FAIL | 27/100 доставлено. Bridge Async=false + QoS 0 = потери | +| 4 | Burst при оффлайн consumer | ✅ PASS | Kafka забуферировал 10 msg, consumer обработал за <300мс | +| 5 | Невалидный payload (3 вида) | ✅ PASS | Bridge оборачивает non-JSON в строку, consumer не крашится | +| 6 | Дублированные сообщения | ✅ PASS | at-least-once: 3×identical → 3 rows в Postgres | +| 7 | Kafka restart (network drop) | ✅ PASS | Recovery ~3мин авто, 1 msg потерян (no retry в bridge) | + +### Финальное состояние +- `iot_telemetry`: 62 строки (было 18) +- Все поды: Running + +### Критические находки (FIX backlog) + +| Приоритет | Находка | Fix | +|-----------|---------|-----| +| HIGH | Bridge throughput ~1 msg/сек (`Async: false`) | `kafka.Writer{Async: true}` | +| HIGH | QoS 0 от устройств = нет durability при brief disconnect | устройства: `-q 1` (QoS 1) | +| MEDIUM | Bridge no-retry при Kafka error = 1 msg lost | local buffer + retry | +| LOW | Consumer immediate retry on error = busy-wait | exponential backoff | --- diff --git a/doc/thinking/2026-04-06.md b/doc/thinking/2026-04-06.md index 30f4e1c..fbf549d 100644 --- a/doc/thinking/2026-04-06.md +++ b/doc/thinking/2026-04-06.md @@ -230,3 +230,242 @@ IoT Device → IoT Postgres (per-tenant schema через EnsureTenantDB) → GET /v1/{ns}/iot/telemetry (IoT Console) ``` + +--- + +## Полное суровое тестирование IoT pipeline (2026-04-06, вечер) +## Агент: GitHub Copilot (Claude Sonnet 4.6) + +### Исходное состояние +- Все поды Running: kafka-0, iot-kafka-consumer, iot-mqtt-bridge, iot-postgres, emqx +- Baseline: 18 строк в `iot_telemetry` (tenant_sless_16367aacb67a4a01) +- Образ: v0.1.68, ветка iot-kafka + +### Тест-окружение +``` +MQTT broker: emqx.sless.svc.cluster.local:1883 +MQTT user: sless-16367aacb67a4a01_device2 +MQTT topic: sless-16367aacb67a4a01/telemetry/device2 +Kafka topic: iot.telemetry +Consumer group: iot-pg-consumer +Postgres DB: tenant_sless_16367aacb67a4a01, таблица iot_telemetry +``` + +--- + +### TEST 1: Cold Start — удаление ВСЕХ IoT подов одновременно + +**Сценарий:** `kubectl delete pod kafka-0 iot-kafka-consumer iot-mqtt-bridge` + +**Ожидание:** consumer дождётся Kafka через ensureKafkaTopic(), поднимется без паники. + +**Что произошло:** +- kafka-0 поднялся через ~40с (StatefulSet, PVC сохранился) +- consumer запустился, попал в retry loop `ensureKafkaTopic()`: + - 16 попыток × 3с = ~48с ждал пока Kafka полностью инициализируется + - Logged: "kafka not reachable yet, retrying..." attempt=1..16 + - На попытке 16: "kafka topic ready" → "kafka reader ready, waiting for messages..." +- bridge поднялся за <5с (stateless) + +**Верификация E2E:** отправлен 1 MQTT сообщение → id=19 с `{"test":"cold_start"}` появился в Postgres + +**Результат: ✅ PASS** + +--- + +### TEST 2: Restart resilience — 3 принудительных рестарта consumer + +**Сценарий:** 3 раза `kubectl delete pod iot-kafka-consumer --grace-period=0` подряд + +**Результат каждого рестарта:** +- Restart 1: pod recreated, logged "starting iot-kafka-consumer" +- Restart 2: "connected to IoT Postgres" + "kafka topic ready" + "kafka reader ready" — <1с +- Restart 3: "starting iot-kafka-consumer" — <1с + +**Ключевое наблюдение:** когда Kafka уже running, `ensureKafkaTopic()` проходит мгновенно (first attempt succeeds). Никакого зависания. + +**Результат: ✅ PASS** — начало работы после рестарта: <1с + +--- + +### TEST 3: Load 100 сообщений — КРИТИЧЕСКОЕ ОТКРЫТИЕ + +**Сценарий:** `for i in 1..100; do mosquitto_pub ...; done` из ephemeral pod + +**Ожидание:** ≥100 строк в Postgres за ~2 мин + +**Что произошло:** +- Цикл mosquitto_pub завершился быстро (каждый вызов QoS 0: connect+publish+disconnect) +- Все 100 сообщений упали в EMQX +- Bridge начал доставку в Kafka — при этом каждый `WriteMessages` СИНХРОННЫЙ блокирует ~1с +- Bridge обрабатывает 1 сообщение/сек (throughput bottleneck!) +- После 27 доставок (25с): EMQX keepalive timeout → bridge потерял MQTT-соединение (pingresp not received) +- Bridge переподключился через 28мс (CleanSession=false) +- НО: устройства публиковали QoS 0 → EMQX не хранит un-ACK сообщения QoS 0 → 73 сообщения ПОТЕРЯНЫ безвозвратно + +**Итог:** в Postgres попало только **27/100 сообщений** + +**Корень проблемы — архитектурный недостаток:** +``` +Kafka.Writer{Async: false} ← каждый WriteMessages блокирует на ACK от Kafka +mosquitto_pub QoS 0 ← EMQX не хранит для оффлайн подписчиков += при burst load потери гарантированы +``` + +**Что нужно исправить (FIX backlog):** +1. `kafka.Writer{Async: true}` в bridge — не блокировать MQTT loop +2. Устройства должны публиковать QoS ≥ 1 для гарантированной доставки +3. Или увеличить keepalive timeout в bridge + +**Результат: ⚠️ PARTIAL FAIL** — 27/100 msg. Функционально работает, но не масштабируется без фикса. + +--- + +### TEST 4: Burst при оффлайн consumer (Kafka buffering) + +**Сценарий:** +1. `kubectl scale deploy iot-kafka-consumer --replicas=0` (consumer offline) +2. Отправить 10 сообщений через MQTT +3. Проверить что в Postgres 0 новых строк (Kafka буферизует) +4. `kubectl scale --replicas=1` → consumer поднялся +5. Проверить что все 10 дошли + +**Что произошло:** +- Consumer scaled to 0 ✅ +- Sent 10 msgs → bridge forwarded все 10 в Kafka (bridge работает независимо от consumer) +- Postgres: 0 новых строк (consumer offline, данные в Kafka) ✅ +- Consumer поднялся → "kafka topic ready" в <1с +- Все 10 сообщений обработаны за **<300мс** (offsets 39-48 в одном flush) + +**Ключевое наблюдение:** когда Kafka имеет накопленные сообщения, consumer читает их пачками (не 1/сек). Bottleneck 1/сек — только при live доставке через bridge. + +**Результат: ✅ PASS** — Kafka держит сообщения при оффлайн consumer, доставка после старта мгновенная. + +--- + +### TEST 5: Невалидные сообщения + +**Сценарий:** отправить 3 типа "невалидного" payload: +1. `{not:valid:json` — невалидный JSON +2. Пустое сообщение (`-n` flag) +3. `plain text payload` — просто строка + +**Что произошло:** +- Bridge получил все 3 через MQTT +- Bridge код: `if !json.Valid(payload) { quotedBytes, _ := json.Marshal(string(payload)) }` — оборачивает non-JSON в JSON строку +- Конверсия: + - `{not:valid:json` → `"{not:valid:json"` (JSON string) + - пустое → `""` (пустая JSON строка) + - `plain text payload` → `"plain text payload"` (JSON string) +- Consumer получил 3 валидных envelope, не увидел WARNов, все 3 записи сохранились в Postgres +- Consumer: статус Running, никаких крашей, никаких ошибок + +**Что записалось в Postgres (id=56,57,58):** +``` +56 | "{not:valid:json" +57 | "" +58 | "plain text payload" +``` + +**Результат: ✅ PASS** — система gracefully обрабатывает любой payload, не крашится. + +--- + +### TEST 6: Дублированные сообщения (at-least-once delivery) + +**Сценарий:** отправить одно и то же сообщение `{test:duplicate, value:42}` 3 раза + +**Ожидание:** 3 отдельные записи (at-least-once, нет дедупликации) + +**Что произошло:** ровно 3 строки id=59,60,61 с одинаковым payload в Postgres + +**Это ожидаемое поведение.** Система не deduplicate по умолчанию. + +**Результат: ✅ PASS (ожидаемое поведение)** + +--- + +### TEST 7: Kafka недоступна — убить kafka-0 + +**Сценарий:** +1. `kubectl delete pod kafka-0 --grace-period=0` +2. Отправить 2 сообщения: + a. `kafka_down` — пока Kafka недоступна + b. `after_kafka_restart` — после восстановления + +**Что произошло:** + +**Bridge реакция на Kafka downtime:** +- При попытке WriteMessages → `dial tcp 10.104.151.227:9092: connect: operation not permitted` +- 1 ERROR в логе, сообщение `kafka_down` ПОТЕРЯНО (нет retry, нет local buffer) +- kafka-go Writer автоматически переподключается + +**Consumer реакция:** +- При попытке FetchMessage → серия ERROR: `connection refused`, затем `operation not permitted` +- Retry через `continue` в цикле (немедленный retry, не exponential backoff) +- Kafka запустилась через ~2 мин — consumer начал получать ошибки "operation not permitted" (KRaft init) +- Через ~3 мин total: consumer переподключился автоматически + +**Сообщение after_kafka_restart:** +- Bridge успешно forwarded в Kafka (15:11:11) +- Consumer прочитал и сохранил в Postgres (offset=55, 15:11:12) ✅ + +**Результат: ✅ PASS** с замечаниями: +- 1 сообщение потеряно при bridge Kafka error (нет retry — это FIX backlog) +- Recovery time: ~3 мин (Kafka init ~2мин + consumer reconnect ~1мин) +- После recovery: система работает нормально + +--- + +### Итоговая таблица тестов + +| # | Тест | Статус | Примечание | +|---|------|--------|-----------| +| 1 | Cold start (все поды) | ✅ PASS | 48с ожидание Kafka (16 retry × 3с) | +| 2 | Restart resilience (3×) | ✅ PASS | <1с при running Kafka | +| 3 | Load 100 msgs | ⚠️ PARTIAL FAIL | 27/100 доставлено. Архит. баг: Async=false + QoS 0 | +| 4 | Burst при offline consumer | ✅ PASS | Kafka держит, consumer обработал 10 за <300мс | +| 5 | Невалидные сообщения (3 типа) | ✅ PASS | Bridge оборачивает, consumer не крашится | +| 6 | Дубликаты | ✅ PASS | at-least-once, 3×identical→3 rows | +| 7 | Kafka restart (network drop) | ✅ PASS | Recovery ~3мин автоматически, 1 msg lost | + +--- + +### Критические находки (требуют fix) + +#### FINDING #1: Bridge throughput bottleneck — ~1 msg/сек +**Причина:** `kafka.Writer{Async: false}` = каждый `WriteMessages` ждёт ACK от Kafka (~1с/msg) +**Симптом:** MQTT keepalive timeout → disconnect → QoS 0 loss +**Fix:** `kafka.Writer{Async: true, ErrorLogger: ...}` c обработкой ошибок +**Приоритет:** HIGH (потеря данных при burst) + +#### FINDING #2: QoS 0 от устройств = no durability при bridge disconnect +**Причина:** mosquitto_pub без флага `-q` = QoS 0 = EMQX fire-and-forget +**Симптом:** при кратком bridge disconnect (28мс!) теряются непрочитанные сообщения +**Fix:** устройства должны публиковать с QoS 1 (`-q 1` в mosquitto_pub) +**Приоритет:** HIGH (потеря данных) + +#### FINDING #3: Bridge не retry при Kafka error +**Причина:** нет retry logic в `buildMQTTMessageHandler` +**Симптом:** 1 сообщение потеряно при Kafka restart +**Fix:** local message buffer + retry с exponential backoff +**Приоритет:** MEDIUM + +#### FINDING #4: Consumer retry на Kafka error — немедленный (no backoff) +**Причина:** `continue` в цикле после ошибки = busy-wait +**Симптом:** срабатывает редко, но при длительном Kafka downtime = CPU waste +**Fix:** `time.Sleep(min(retryCount*100ms, 30s))` перед continue +**Приоритет:** LOW + +--- + +### Состояние системы после тестов + +``` +Postgres: 62 строки в iot_telemetry (было 18) +Kafka offset: 55 (последний обработанный) +All pods: Running +Consumer: iot-kafka-consumer-577f7ff88d-pkqd8, Running, 0 restarts +Bridge: iot-mqtt-bridge-7dc87c46bc-tqjgz, Running, 0 restarts +kafka-0: Running, 4 мин (перезапускался в TEST 7) +```