diff --git a/CMakeLists.txt b/CMakeLists.txt index 5f7166a..f032a97 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -14,6 +14,7 @@ if(ESP_PLATFORM) "src/ayla/httpc.cpp" "src/ayla/httpd.cpp" "src/ayla/log.cpp" + "src/ayla/session.cpp" "src/ayla/platform/esp-idf/platform.cpp" INCLUDE_DIRS "include" @@ -79,6 +80,7 @@ else() src/ayla/httpc.cpp src/ayla/httpd.cpp src/ayla/log.cpp + src/ayla/session.cpp src/ayla/platform/posix/platform.cpp ) target_include_directories(fgl-aircon diff --git a/core.37208 b/core.37208 new file mode 100644 index 0000000..35e3499 Binary files /dev/null and b/core.37208 differ diff --git a/docs/LEGACY_ANALYSIS.md b/docs/LEGACY_ANALYSIS.md index d8b13b6..b5143aa 100644 --- a/docs/LEGACY_ANALYSIS.md +++ b/docs/LEGACY_ANALYSIS.md @@ -30,13 +30,14 @@ ### 2.1. [ГЛАВНАЯ ПРИЧИНА «РАССИНХРОНИЗАЦИИ»] Keep-alive 1200 с вместо 10–15 с `notifier.py:_KEEP_ALIVE_INTERVAL = 1200.0`. APK: 10 с (или `lan.json:keepAlive/3`). -Проверено на приборе: **единственный механизм восстановления после расхождения -CBC-цепочек — принудительный re-key, который модуль делает при получении -`local_reg` для сессии старше ≈44 с. Ответы 400/401 модуль игнорирует.** -Следствие для legacy: любая потерянная пара запрос-ответ/обрыв соединения → -обе стороны «глохнут» на срок до 20 минут (до следующего local_reg). Наблюдаемый -симптом «перестаёт понимать кондиционер» с самопроизвольным восстановлением — -именно это. +Проверено на приборе: **модуль игнорирует 400/401 на свои POST; переkey +происходит при `local_reg` после зазора ≥ ~44–50 с от предыдущего** (при +keep-alive 10–15 с re-key вообще не происходит). Следствие для legacy: любая +потерянная пара запрос-ответ → обе стороны «глохнут» до следующего local_reg +(до 20 минут), который завершится re-key — наблюдаемый симптом «перестаёт +понимать кондиционер, потом сам чинится» — именно это. Корректная стратегия +для новой реализации: при ошибке расшифровки — пауза ~50 с, затем local_reg +(re-key гарантирован). Дополнительно: длинные паузы между local_reg держат сессию «полуживой» (модуль не видит keep-alive, но слот может удерживаться), и конфликт за diff --git a/docs/PLAN_CORE_LIBRARY.md b/docs/PLAN_CORE_LIBRARY.md index 60acc92..8a5d854 100644 --- a/docs/PLAN_CORE_LIBRARY.md +++ b/docs/PLAN_CORE_LIBRARY.md @@ -199,9 +199,13 @@ int32_t fgl_convert_from_input(fgl_template, fgl_prop, int32_t disp, ### 5.2. Keep-alive и re-key * Таймер `keepalive_ms` (default **15000**); по истечении — `PUT local_reg` c `notify=(очередь непуста)`. Каждый `commands.json` перезапускает таймер. -* Re-key (очередной local_reg при возрасте сессии ≥ ~44 с) — штатное - событие: перегенерация шифров/цепочек, сессия не пересоздаётся, начальная - синхронизация не повторяется; seq_no приложения продолжает глобальный счётчик. +* Re-key — событие по инициативе модуля (при зазоре local_reg ≥ ~44–50 с, + [ПРОВЕРЕНО НА ПРИБОРЕ]; при штатном keep-alive НЕ происходит): обработать + как обычный KE (перегенерация шифров/цепочек), сессию не пересоздавать, + начальную синхронизацию не повторять; seq_no приложения продолжает + глобальный счётчик. +* Восстановление при ошибке расшифровки: тишина > порога (50 с по умолчанию) + и возврат — модуль гарантированно ре-кает [ПРОВЕРЕНО НА ПРИБОРЕ]. * Анти-спам: ≤1 local_reg/с; notify=1 — один на пакет команд. ### 5.3. Очередь команд @@ -291,7 +295,7 @@ HA-превью (PLAN_HOME_ASSISTANT §4) и тестами. |---|------------|------------------| | M0 ✅ | Монорепо-каркас: CMake (корень) + IDF-подключение, платслой, лог, CI | Собирается linux+esp-idf; пустой httpd отвечает 404 | | M1 ✅ | `src/ayla`: crypto+envelope, мини-httpd/httpc, jsmn-вендор | Векторы зелёные; httpd-тесты; совместимость с probe_reference.py | -| M2 | `src/ayla`: сессия (установка/активация/keep-alive/re-key/слоты/503/delete) с mock-модулем | Все сценарии mock; на приборе: активация ≤5 с, re-key каждые 45–60 с | +| M2 ✅ | `src/ayla`: сессия (установка/активация/keep-alive/re-key/слоты/503/delete) с mock-модулем | Все сценарии mock; на приборе: активация ≤5 с; семантика re-key: при зазоре local_reg ≥ ~44–50 с (при честном keep-alive 15 с — 0 re-key за 100 с; при 50 с — 3 re-key) | | M3 | `src/aircon`: шаблоны, конверсии+override, публичный API, batch | `tests/aircon` зелёные; на приборе: чтение всех свойств, batch=1 notify | | M4 | fglctl-пример, `tools/fglair-discover` (в т.ч. `--format esphome-secrets`), README библиотеки (сборка IDF/POSIX, тесты) | 24 ч на приборе: 0 рассинхронов; README готов | | M5 | (Опция) `FglHub` N устройств; mDNS-резолвер как опция host-разрешения | Два устройства одновременно | diff --git a/docs/PROTOCOL.md b/docs/PROTOCOL.md index c2a229c..ef2a061 100644 --- a/docs/PROTOCOL.md +++ b/docs/PROTOCOL.md @@ -216,15 +216,21 @@ GET http://:<порт>/local_lan/commands.json один «пустой» опрос `commands.json` — это признак принятой сессии. 2. `local_reg` от endpoint'а с живой сессией **моложе ~40 с** → только keep-alive, без key exchange. -3. `local_reg` от endpoint'а с сессией **старше ~44 с** → модуль принудительно - инициирует новый key exchange (ротация сессионных ключей). Т.е. при штатном - keep-alive каждые 10–15 с ключи ротируются примерно каждые 45–60 с. - `time_1` модуля — тикающий счётчик с шагом ≈10 нс (аптайм); порог, - вероятно, 44 с в этих единицах либо просто 4.4e9 тиков. +3. `local_reg` при **зазоре ≥ ~44–50 с** от предыдущего local_reg → + модуль принудительно инициирует новый key exchange («вернувшееся» + приложение получает свежие ключи). При штатном keep-alive каждые 10–15 с + re-key НЕ происходит — сессия живёт сколь угодно долго (проверено: + 100 с при 15 с keep-alive — 0 re-key; 125 с при 50 с keep-alive — 3 re-key, + оба без потерь). `time_1` модуля — тикающий счётчик с шагом ≈10 нс (аптайм); + порог, вероятно, 4.4e9 тиков (~44 с) от последнего local_reg. 4. **Ответы 401/400 на POST модуля игнорируются**: сессия продолжает работать, - re-key не вызывается. Единственный механизм восстановления после расхождения - CBC-цепочек — принудительный re-key по `local_reg` (п. 3). Поэтому интервал - keep-alive = интервал потенциального «зависания» при десинхроне. + re-key не вызывается. Восстановление после расхождения CBC-цепочек — + намеренная «тишина» приложения на > порога из п. 3 с последующим + `local_reg`: модуль сочтёт приложение вернувшимся и ре-кает. Т.е. стратегия + самовосстановления: при ошибке расшифровки — пауза keep-alive ~50–60 с, + затем возобновить (проверено на приборе). Отдельный случай — бракованная + подпись при живой цепочке (сообщение расшифровано, подпись не сошлась): + цепочка НЕ расходится, следующий push восстанавливает работу без re-key. 5. `delete_session` освобождает слот немедленно; следующий `local_reg` того же endpoint'а создаёт новую сессию. 6. Наблюдавшийся (не воспроизведённый повторно) режим отказа: модуль отвечает @@ -316,11 +322,12 @@ data: {"id":"","ack_status":200,"ack_message":0,"dsn":"..."} * **Потеря CBC-цепочки** (§3.4): модуль не может расшифровать ответ приложения / приложение не может расшифровать push модуля. Ответы 401/400 на POST модуля - **игнорируются** — модуль продолжает слать в «сломанный» канал. Восстановление - происходит только когда очередной `local_reg` (по возрасту ≥ ~44 с или от - нового endpoint'а) вызовет новый key exchange. Следствие: **интервал - keep-alive = максимальное время «мёртвой» сессии при десинхроне** - (10–15 с — незаметно; 1200 с как в legacy-скрипте — 20 минут глухоты). + **игнорируются** — модуль продолжает слать в «сломанный» канал. Восстановление: + приложение замолкает на > ~44–50 с (порог «возврата» из п. 4.4.3) и шлёт + `local_reg` — модуль переkey'ается. Реализация ядра: при ошибке расшифровки + пауза keep-alive ~50 с, затем возобновление. В legacy-скрипте пауза получалась + «бесплатно» из-за keep-alive 1200 с: каждый цикл завершался re-key при + возврате — потому рассинхрон «сам чинился» через ~20 минут. * **Смена lanip_key** (`key_id` не совпал): теоретический путь по APK — 412 + `refreshLanConfig()` из облака. За 5 лет эксплуатации прибора ротации ключа не наблюдалось ни разу; ключ, по-видимому, зашит в модуль, облако лишь хранит diff --git a/src/ayla/httpd.cpp b/src/ayla/httpd.cpp index 1b6506f..d82641a 100644 --- a/src/ayla/httpd.cpp +++ b/src/ayla/httpd.cpp @@ -121,7 +121,7 @@ void parse_target(const char* full_target, HttpRequest* req) { } } -constexpr uint32_t kHttpdThreadStack = 8192; +constexpr uint32_t kHttpdThreadStack = 24576; // commands-путь: envelope+crypto ~10КБ поверх буферов запроса constexpr uint32_t kAcceptPollMs = 100; constexpr uint32_t kClientRxTimeoutMs = 30000; constexpr uint32_t kClientTxTimeoutMs = 10000; diff --git a/src/ayla/platform/esp-idf/platform.cpp b/src/ayla/platform/esp-idf/platform.cpp index bada591..b5e241c 100644 --- a/src/ayla/platform/esp-idf/platform.cpp +++ b/src/ayla/platform/esp-idf/platform.cpp @@ -159,6 +159,35 @@ int tcp_connect(const char* host, uint16_t port, uint32_t timeout_ms) { return fd; } +bool local_ip_for(const char* host, char* out, size_t out_cap) { + struct addrinfo hints {}; + hints.ai_family = AF_INET; + hints.ai_socktype = SOCK_DGRAM; + struct addrinfo* list = nullptr; + if (lwip_getaddrinfo(host, "80", &hints, &list) != 0 || list == nullptr) { + return false; + } + int fd = lwip_socket(AF_INET, SOCK_DGRAM, 0); + if (fd < 0) { + lwip_freeaddrinfo(list); + return false; + } + bool ok = lwip_connect(fd, list->ai_addr, list->ai_addrlen) == 0; + struct sockaddr_in local {}; + socklen_t slen = sizeof(local); + if (ok && lwip_getsockname(fd, reinterpret_cast(&local), + &slen) == 0) { + const char* s = inet_ntop(AF_INET, &local.sin_addr, out, + static_cast(out_cap)); + ok = s != nullptr; + } else { + ok = false; + } + lwip_freeaddrinfo(list); + lwip_close(fd); + return ok; +} + long tcp_send(int fd, const void* buf, size_t len) { const uint8_t* p = static_cast(buf); size_t done = 0; diff --git a/src/ayla/platform/platform.hpp b/src/ayla/platform/platform.hpp index cab4e7c..49d147a 100644 --- a/src/ayla/platform/platform.hpp +++ b/src/ayla/platform/platform.hpp @@ -43,6 +43,9 @@ uint16_t tcp_local_port(int fd); int tcp_accept(int listen_fd, uint32_t* peer_ip, uint16_t* peer_port); // Подключается к host:port (host — DNS-имя или dotted-quad) с таймаутом. int tcp_connect(const char* host, uint16_t port, uint32_t timeout_ms); +// Локальный IP-адрес (dotted) интерфейса, которым достигается host +// (без реального трафика: UDP connect). Для local_reg. +bool local_ip_for(const char* host, char* out, size_t out_cap); long tcp_send(int fd, const void* buf, size_t len); // >0 / -1 // Блокирующее чтение; 0 — EOF, -1 — ошибка/таймаут. long tcp_recv(int fd, void* buf, size_t len); diff --git a/src/ayla/platform/posix/platform.cpp b/src/ayla/platform/posix/platform.cpp index bca64ba..41107aa 100644 --- a/src/ayla/platform/posix/platform.cpp +++ b/src/ayla/platform/posix/platform.cpp @@ -60,7 +60,12 @@ bool thread_create(void (*fn)(void*), void* ctx, const char* name, pthread_t tid; pthread_attr_t attr; pthread_attr_init(&attr); - if (stack_bytes > 0) pthread_attr_setstacksize(&attr, stack_bytes); + if (stack_bytes > 0) { + if (pthread_attr_setstacksize(&attr, stack_bytes) != 0) { + // glibc отвергает < PTHREAD_STACK_MIN; остаётся дефолт (больше — не меньше) + // логируем только: ядро запрашивает >= PTHREAD_STACK_MIN. + } + } int rc = pthread_create(&tid, &attr, thread_trampoline, start); pthread_attr_destroy(&attr); if (rc != 0) { @@ -168,6 +173,35 @@ int tcp_connect(const char* host, uint16_t port, uint32_t timeout_ms) { return fd; } +bool local_ip_for(const char* host, char* out, size_t out_cap) { + struct addrinfo hints {}; + hints.ai_family = AF_INET; + hints.ai_socktype = SOCK_DGRAM; + struct addrinfo* list = nullptr; + if (::getaddrinfo(host, "80", &hints, &list) != 0 || list == nullptr) { + return false; + } + int fd = ::socket(AF_INET, SOCK_DGRAM, 0); + if (fd < 0) { + ::freeaddrinfo(list); + return false; + } + bool ok = ::connect(fd, list->ai_addr, list->ai_addrlen) == 0; + struct sockaddr_in local {}; + socklen_t slen = sizeof(local); + if (ok && ::getsockname(fd, reinterpret_cast(&local), + &slen) == 0) { + const char* s = inet_ntop(AF_INET, &local.sin_addr, out, + static_cast(out_cap)); + ok = s != nullptr; + } else { + ok = false; + } + ::freeaddrinfo(list); + ::close(fd); + return ok; +} + long tcp_send(int fd, const void* buf, size_t len) { const uint8_t* p = static_cast(buf); size_t done = 0; diff --git a/src/ayla/session.cpp b/src/ayla/session.cpp new file mode 100644 index 0000000..cfd2c3b --- /dev/null +++ b/src/ayla/session.cpp @@ -0,0 +1,909 @@ +#include "ayla/session.hpp" + +#include +#include +#include +#include +#include + +#include "ayla/envelope.hpp" +#include "ayla/httpc.hpp" +#include "ayla/json.hpp" +#include "ayla/log.hpp" +#include "ayla/platform/platform.hpp" + +namespace fgl::ayla { + +namespace { + +constexpr uint32_t kLoopTickMs = 20; +constexpr size_t kMaxName = 40; +constexpr size_t kMaxQueue = 32; + +struct Command { + uint8_t type; // 1=GET, 2=SET, 3=DELETE + char name[kMaxName]; + int64_t value; + char base_type[10]; + int cmd_id; +}; + +bool gen_random_token(char* out, size_t len) { + static const char kAlpha[] = + "ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789"; + uint8_t rnd[32]; + size_t need = len <= sizeof(rnd) ? len : sizeof(rnd); + if (!plat::random(rnd, need)) return false; + for (size_t i = 0; i < len; i++) { + out[i] = kAlpha[rnd[i % need] % 62]; + } + out[len] = '\0'; + return true; +} + +// "cmd_id=5&status=200" -> значения; отсутствующие остаются прежними. +void parse_query(const char* query, int* cmd_id, int* status) { + const char* p = query; + while (*p != '\0') { + const char* eq = strchr(p, '='); + const char* amp = strchr(p, '&'); + const char* end = (amp != nullptr) ? amp : p + strlen(p); + if (eq != nullptr && eq < end) { + size_t klen = static_cast(eq - p); + int v = atoi(eq + 1); + if (klen == 6 && strncmp(p, "cmd_id", 6) == 0) *cmd_id = v; + if (klen == 6 && strncmp(p, "status", 6) == 0) *status = v; + } + if (amp == nullptr) break; + p = amp + 1; + } +} + +} // namespace + +struct Session::Impl { + SessionConfig cfg{}; + SessionTimings timings{}; + SessionCallbacks cbs{}; + + char host[64] = {}; + char dsn[40] = {}; + char lanip_key[64] = {}; + + // ---- httpd-поток (используется только из httpd-потока) ---- + HttpServer httpd; + SessionCrypto crypto; + bool crypto_ready = false; + int64_t out_seq = 0; + char ke_random2[17] = {}; + // Рабочие буферы datapoint-пути (httpd однопоточен; вынесены из стека — + // экономия ~7КБ стека потока httpd). + char dp_body[kHttpdMaxBody + 1]; + char dp_enc_b64[kEnvelopeMaxB64]; + char dp_sign_b64[96]; + char dp_plain[kEnvelopeMaxPlain]; + char dp_data[kEnvelopeMaxPlain]; + + // ---- разделяемое ---- + std::mutex queue_mu; + Command queue[kMaxQueue] = {}; + uint8_t queue_len = 0; + int next_cmd_id = 1; + bool batch_open = false; + Command batch[kMaxQueue] = {}; + uint8_t batch_len = 0; + + std::atomic state{static_cast(SessionState::kIdle)}; + std::atomic last_error{static_cast(SessionError::kNone)}; + std::atomic ke_time_ms{0}; // время ответа на KE (0 — не было) + std::atomic last_local_reg_ms{0}; + std::atomic quiet_until_ms{0}; // пауза local_reg (восстановление) + std::atomic retry_at_ms{0}; + std::atomic had_poll_since_ke{false}; + std::atomic ever_active{false}; // был online хотя бы раз + std::atomic want_notify{false}; // после batch commit / перехода online + std::atomic delete_pending{false}; + std::atomic delete_served{false}; + std::atomic decrypt_failed{false}; + std::atomic running{false}; + + std::atomic rekeys{0}; + std::atomic pushes_ok{0}; + std::atomic pushes_bad{0}; + std::atomic cmds_served{0}; + + plat::ThreadId thread = nullptr; + uint8_t backoff_attempts = 0; + uint32_t backoff_ms = 0; + bool reg_ok = false; // последний local_reg принят (202/200) + + uint16_t listen_port_actual = 0; + + // ---------- helpers (вызывается из обоих потоков) ---------- + void set_state(SessionState st, SessionError err) { + uint8_t prev = state.exchange(static_cast(st), + std::memory_order_acq_rel); + last_error.store(static_cast(err), std::memory_order_release); + if (prev != static_cast(st) && cbs.on_state != nullptr) { + cbs.on_state(cbs.ctx, st, err); // без дублирования одинаковых состояний + } + } + + // ---------- очередь (mutex) ---------- + bool enqueue_locked(const Command& cmd) { + // coalescing: SET замещает незабранный SET того же свойства; + // GET-дубликат отбрасывается; DELETE — единственный. + if (cmd.type == 3) { + for (uint8_t i = 0; i < queue_len; i++) { + if (queue[i].type == 3) return true; + } + } else { + for (uint8_t i = 0; i < queue_len; i++) { + if (queue[i].type == cmd.type && + strncmp(queue[i].name, cmd.name, kMaxName) == 0) { + if (cmd.type == 2) { + queue[i].value = cmd.value; // замещаем + return true; + } + if (cmd.type == 1) return true; // дубликат GET + } + } + } + if (queue_len >= cfg.max_queue || queue_len >= kMaxQueue) return false; + queue[queue_len++] = cmd; + return true; + } + + bool submit(Command cmd) { + std::lock_guard lk(queue_mu); + bool ok; + if (batch_open && cmd.type != 3) { + if (batch_len >= kMaxQueue) return false; + // coalescing внутри batch + for (uint8_t i = 0; i < batch_len; i++) { + if (batch[i].type == cmd.type && + strncmp(batch[i].name, cmd.name, kMaxName) == 0) { + if (cmd.type == 2) { + batch[i].value = cmd.value; + return true; + } + if (cmd.type == 1) return true; + } + } + batch[batch_len++] = cmd; + ok = true; + } else { + ok = enqueue_locked(cmd); + } + return ok; + } + + // ---------- httpd-обработчики (httpd-поток) ---------- + static bool http_handler(const HttpRequest& req, HttpResponse& resp, + void* ctx); + + void handle_key_exchange(const HttpRequest& req, HttpResponse& resp); + void handle_commands(HttpResponse& resp); + void handle_datapoint(const HttpRequest& req, HttpResponse& resp); + + void build_get_payload(char* out, size_t out_cap, const Command& c, + int* seq_out); + void build_set_payload(char* out, size_t out_cap, const Command& c); + + // ---------- session-поток ---------- + static void session_thread_trampoline(void* ctx) { + static_cast(ctx)->session_loop(); + } + void session_loop(); + bool send_local_reg(bool notify, bool first); +}; + +// --------------------------------------------------------------------------- +// httpd +// --------------------------------------------------------------------------- + +bool Session::Impl::http_handler(const HttpRequest& req, HttpResponse& resp, + void* ctx) { + auto* impl = static_cast(ctx); + if (strcmp(req.method, "POST") == 0) { + if (strcmp(req.target, "/local_lan/key_exchange.json") == 0) { + impl->handle_key_exchange(req, resp); + return true; + } + if (strcmp(req.target, "/local_lan/property/datapoint.json") == 0) { + impl->handle_datapoint(req, resp); + return true; + } + if (strcmp(req.target, "/local_lan/property/datapoint/ack.json") == 0 || + strcmp(req.target, "/local_lan/node/property/datapoint.json") == 0 || + strcmp(req.target, "/local_lan/node/property/datapoint/ack.json") == 0) { + FGL_LOGD("session: ack/node datapoint (пусто ok)"); + resp.status = 200; + return true; + } + } else if (strcmp(req.method, "GET") == 0) { + if (strcmp(req.target, "/local_lan/commands.json") == 0) { + if (!impl->crypto_ready) { + resp.status = 401; + return true; + } + impl->handle_commands(resp); + return true; + } + } + resp.status = 404; + return true; +} + +void Session::Impl::handle_key_exchange(const HttpRequest& req, + HttpResponse& resp) { + // Тело: {"key_exchange":{"ver":1,"proto":1,"key_id":N,"random_1":..,"time_1":N,"sec":""}} + char body[kHttpdMaxBody + 1]; + size_t n = req.body_len < kHttpdMaxBody ? req.body_len : kHttpdMaxBody; + memcpy(body, req.body, n); + body[n] = '\0'; + + json::Doc doc; + if (!doc.parse(body)) { + resp.status = 400; + return; + } + int64_t ver = 0, proto = 0, key_id = 0, time_1 = 0; + char random_1[32] = {}, sec[8] = {}; + bool have_r1 = doc.get_string("random_1", random_1, sizeof(random_1)); + bool have_t1 = doc.get_int("time_1", &time_1); + bool have_sec = doc.get_string("sec", sec, sizeof(sec)); + bool have_ver = doc.get_int("ver", &ver); + bool have_proto = doc.get_int("proto", &proto); + bool have_kid = doc.get_int("key_id", &key_id); + if (!have_r1 || !have_t1 || !have_ver || !have_proto || !have_kid) { + FGL_LOGW("session: key_exchange неполный"); + resp.status = 400; + return; + } + if (ver != 1 || proto != 1 || (have_sec && sec[0] != '\0')) { + FGL_LOGW("session: key_exchange ver/proto/sec не поддержаны"); + set_state(SessionState::kKeyError, SessionError::kBadKeyExchange); + resp.status = 426; + return; + } + { + // sec длиннее буфера -> get_string=false, но поле есть: setup-режим + // не поддерживаем — тоже 426 (PROTOCOL §3.1). + const char* sec_start = nullptr; + size_t sec_len = 0; + jsmntype_t sec_type; + if (doc.find("sec", &sec_start, &sec_len, &sec_type) && sec_type == JSMN_STRING && + sec_len >= sizeof(sec)) { + set_state(SessionState::kKeyError, SessionError::kBadKeyExchange); + resp.status = 426; + return; + } + } + if (static_cast(key_id) != cfg.lanip_key_id) { + FGL_LOGE("session: key_id %lld != %u — ротация ключа?", + static_cast(key_id), + static_cast(cfg.lanip_key_id)); + set_state(SessionState::kKeyError, SessionError::kKeyMismatch); + resp.status = 412; + return; + } + + if (!gen_random_token(ke_random2, 16)) { + resp.status = 500; + return; + } + int64_t time_2 = static_cast(plat::now_ms()) * 1000000ll; + crypto_ready = crypto.init(lanip_key, random_1, ke_random2, time_1, + time_2); // цепочки сбрасываются тут же + if (!crypto_ready) { + resp.status = 500; + return; + } + had_poll_since_ke.store(false, std::memory_order_release); + ke_time_ms.store(plat::now_ms(), std::memory_order_release); + bool rekey = ever_active.load(std::memory_order_acquire); + rekeys.fetch_add(1, std::memory_order_relaxed); + if (!rekey && state.load(std::memory_order_acquire) != + static_cast(SessionState::kRegistering)) { + set_state(SessionState::kRegistering, SessionError::kNone); + } + FGL_LOGI("session: key exchange #%u (rekey=%d)", + static_cast(rekeys.load(std::memory_order_relaxed)), rekey); + + static thread_local char out[128]; + json::Writer w(out, sizeof(out)); + w.begin_object(); + w.key("random_2"); + w.string(ke_random2); + w.key("time_2"); + w.integer(time_2); + w.end_object(); + resp.status = 200; + resp.body = reinterpret_cast(out); + resp.body_len = strlen(out); + // буфер out живёт до конца ответа (send_response в handle_connection + // выполняется синхронно в том же кадре стека httpd-потока). +} + +void Session::Impl::build_get_payload(char* out, size_t out_cap, + const Command& c, int* seq_out) { + (void)seq_out; + json::Writer w(out, out_cap); + w.begin_object(); + w.key("cmds"); + w.begin_array(); + w.begin_object(); + w.key("cmd"); + w.begin_object(); + w.key("cmd_id"); + w.integer(c.cmd_id); + w.key("method"); + w.string("GET"); + w.key("resource"); + char res[80]; + snprintf(res, sizeof(res), "property.json?name=%s", c.name); + w.string(res); + w.key("data"); + w.string(""); + w.key("uri"); + w.string("/local_lan/property/datapoint.json"); + w.end_object(); + w.end_object(); + w.end_array(); + w.end_object(); +} + +void Session::Impl::build_set_payload(char* out, size_t out_cap, + const Command& c) { + json::Writer w(out, out_cap); + w.begin_object(); + w.key("properties"); + w.begin_array(); + w.begin_object(); + w.key("property"); + w.begin_object(); + w.key("base_type"); + w.string(c.base_type); + w.key("name"); + w.string(c.name); + w.key("value"); + if (strcmp(c.base_type, "boolean") == 0) { + w.boolean(c.value != 0); + } else { + w.integer(c.value); + } + w.key("id"); + char id[9]; + gen_random_token(id, 8); + w.string(id); + w.end_object(); + w.end_object(); + w.end_array(); + w.end_object(); +} + +void Session::Impl::handle_commands(HttpResponse& resp) { + Command head{}; + bool have = false; + uint8_t remaining = 0; + { + std::lock_guard lk(queue_mu); + if (queue_len > 0) { + head = queue[0]; + have = true; + queue_len--; + memmove(queue, queue + 1, queue_len * sizeof(Command)); + remaining = queue_len; + } + } + char payload[512]; + if (!have) { + payload[0] = '{'; + payload[1] = '}'; + payload[2] = '\0'; + } else if (head.type == 1) { + build_get_payload(payload, sizeof(payload), head, nullptr); + } else if (head.type == 2) { + build_set_payload(payload, sizeof(payload), head); + } else { // DELETE session + json::Writer w(payload, sizeof(payload)); + w.begin_object(); + w.key("cmds"); + w.begin_array(); + w.begin_object(); + w.key("cmd"); + w.begin_object(); + w.key("cmd_id"); + w.integer(0); + w.key("method"); + w.string("DELETE"); + w.key("resource"); + w.string("local_reg.json"); + w.key("data"); + w.string("delete_session"); + w.key("uri"); + w.string("/local_lan"); + w.end_object(); + w.end_object(); + w.end_array(); + w.end_object(); + delete_served.store(true, std::memory_order_release); + delete_pending.store(false, std::memory_order_release); + FGL_LOGI("session: delete_session выдан модулю"); + } + + static thread_local char envelope[kEnvelopeMaxB64]; + int64_t seq = out_seq++; + if (!envelope_pack(crypto.app, seq, payload, envelope, sizeof(envelope))) { + resp.status = 500; + return; + } + cmds_served.fetch_add(1, std::memory_order_relaxed); + had_poll_since_ke.store(true, std::memory_order_release); + uint8_t st8 = state.load(std::memory_order_acquire); + if (st8 == static_cast(SessionState::kRegistering) || + st8 == static_cast(SessionState::kRecovering) || + st8 == static_cast(SessionState::kOffline)) { + // Опрос команд = сессия жива (в т.ч. после re-key при recovering/offline). + ever_active.store(true, std::memory_order_release); + decrypt_failed.store(false, std::memory_order_release); + set_state(SessionState::kOnline, SessionError::kNone); + want_notify.store(true, std::memory_order_release); + } + resp.status = remaining > 0 ? 206 : 200; + resp.body = reinterpret_cast(envelope); + resp.body_len = strlen(envelope); +} + +void Session::Impl::handle_datapoint(const HttpRequest& req, + HttpResponse& resp) { + size_t n = req.body_len < kHttpdMaxBody ? req.body_len : kHttpdMaxBody; + memcpy(dp_body, req.body, n); + dp_body[n] = '\0'; + json::Doc wrap; + if (!wrap.parse(dp_body)) { + resp.status = 400; + return; + } + if (!wrap.get_string("enc", dp_enc_b64, sizeof(dp_enc_b64)) || + !wrap.get_string("sign", dp_sign_b64, sizeof(dp_sign_b64))) { + resp.status = 400; + return; + } + + int64_t seq_no = -1; + if (!envelope_unpack(crypto.dev, dp_enc_b64, dp_sign_b64, dp_plain, + sizeof(dp_plain), &seq_no)) { + pushes_bad.fetch_add(1, std::memory_order_relaxed); + FGL_LOGW("session: push не расшифрован/подпись (401); пауза и re-key"); + if (state.load(std::memory_order_acquire) == + static_cast(SessionState::kOnline)) { + decrypt_failed.store(true, std::memory_order_release); + set_state(SessionState::kRecovering, SessionError::kDecryptFailed); + // Стратегия восстановления (PROTOCOL §4.4): замолчать на > порога + // «возврата» (~44-50с) — следующий local_reg заставит модуль re-key. + quiet_until_ms.store(plat::now_ms() + timings.recovering_quiet_ms, + std::memory_order_release); + } + resp.status = 401; + return; + } + + // Восстановление после 401 с живой цепочкой (бракованная подпись): + // сообщение расшифровано — отменяем тишину и возвращаем online. + if (decrypt_failed.exchange(false, std::memory_order_acq_rel)) { + quiet_until_ms.store(0, std::memory_order_release); + set_state(SessionState::kOnline, SessionError::kNone); + FGL_LOGI("session: цепочка восстановлена (успешный push после 401)"); + } + + // {"seq_no":N,"data":{"name":..,"value":..}} + json::Doc top; + if (!top.parse(dp_plain)) { + pushes_ok.fetch_add(1, std::memory_order_relaxed); + resp.status = 200; + return; + } + const char* data_start = nullptr; + size_t data_len = 0; + jsmntype_t data_type; + PropertyEvent ev{}; + ev.seq_no = seq_no; + if (top.find("data", &data_start, &data_len, &data_type) && + data_type == JSMN_OBJECT && data_len < sizeof(dp_data)) { + memcpy(dp_data, data_start, data_len); + dp_data[data_len] = '\0'; + { + json::Doc data_doc; + if (data_doc.parse(dp_data)) { + char name[kMaxName]; + if (data_doc.get_string("name", name, sizeof(name))) { + snprintf(ev.name, sizeof(ev.name), "%s", name); + int64_t iv = 0; + bool bv = false; + const char* sv = nullptr; + size_t sv_len = 0; + jsmntype_t vt; + if (data_doc.get_int("value", &iv)) { + ev.is_int = true; + ev.int_value = iv; + } else if (data_doc.get_bool("value", &bv)) { + ev.is_bool = true; + ev.bool_value = bv; + } else if (data_doc.find("value", &sv, &sv_len, &vt) && + vt == JSMN_STRING && sv_len + 1 <= sizeof(ev.str_value)) { + memcpy(ev.str_value, sv, sv_len); + ev.str_value[sv_len] = '\0'; + } + int cmd_id = -1, status = 0; + parse_query(req.query, &cmd_id, &status); + ev.cmd_id = cmd_id; + ev.status = status; + if (cbs.on_property != nullptr) { + cbs.on_property(cbs.ctx, ev); + } + } + } + } + } + pushes_ok.fetch_add(1, std::memory_order_relaxed); + resp.status = 200; +} + +// --------------------------------------------------------------------------- +// session-поток: local_reg / keep-alive / backoff / таймауты +// --------------------------------------------------------------------------- + +bool Session::Impl::send_local_reg(bool notify, bool first) { + char path[96]; + if (first) { + snprintf(path, sizeof(path), "/local_reg.json?dsn=%s", dsn); + } else { + snprintf(path, sizeof(path), "/local_reg.json"); + } + char local_ip[24]; + if (!plat::local_ip_for(host, local_ip, sizeof(local_ip))) { + FGL_LOGW("session: local_ip_for(%s) failed", host); + return false; + } + char body[160]; + json::Writer w(body, sizeof(body)); + w.begin_object(); + w.key("local_reg"); + w.begin_object(); + w.key("ip"); + w.string(local_ip); + w.key("notify"); + w.boolean(notify); + w.key("port"); + w.integer(listen_port_actual); + w.key("uri"); + w.string("/local_lan"); + w.end_object(); + w.end_object(); + if (!w.ok()) return false; + + HttpcRequest req; + req.method = first ? "POST" : "PUT"; + req.host = host; + req.port = cfg.device_port; + req.path = path; + req.body = reinterpret_cast(body); + req.body_len = strlen(body); + req.timeout_ms = 5000; + HttpcResponse resp; + if (!httpc_perform(req, &resp)) { + return false; + } + last_local_reg_ms.store(plat::now_ms(), std::memory_order_release); + if (resp.status == 503) { + set_state(SessionState::kOffline, SessionError::kNoSlot); + retry_at_ms.store(plat::now_ms() + timings.no_slot_retry_ms, + std::memory_order_release); + FGL_LOGW("session: local_reg -> 503 (нет слотов)"); + backoff_attempts = 0; + backoff_ms = 0; + reg_ok = false; // сессия не активировалась — следующий local_reg POST?dsn + return true; // транспорт ок — это протокольный ответ + } + if (resp.status != 202 && resp.status != 200) { + FGL_LOGW("session: local_reg -> %d", resp.status); + return false; + } + if (state.load(std::memory_order_acquire) == + static_cast(SessionState::kIdle)) { + set_state(SessionState::kRegistering, SessionError::kNone); + } + backoff_attempts = 0; + backoff_ms = 0; + reg_ok = true; + return true; +} + +void Session::Impl::session_loop() { + bool first_reg = true; + while (running.load(std::memory_order_acquire)) { + uint64_t now = plat::now_ms(); + SessionState st = + static_cast(state.load(std::memory_order_acquire)); + + // kKeyError — устойчивая ошибка: конфиг менять вручную, не дёргаем модуль. + if (st == SessionState::kKeyError) { + plat::sleep_ms(kLoopTickMs); + continue; + } + + // Тишина (восстановление после «KE без poll»). + if (now < quiet_until_ms.load(std::memory_order_acquire)) { + plat::sleep_ms(kLoopTickMs); + continue; + } + + // Активация: KE отвечен, но опроса нет. + uint64_t ke_time = ke_time_ms.load(std::memory_order_acquire); + if (st == SessionState::kRegistering && ke_time != 0 && + !had_poll_since_ke.load(std::memory_order_acquire) && + now - ke_time > timings.activation_timeout_ms) { + FGL_LOGW("session: активация не наступила (KE без poll) — пауза %ums", + static_cast(timings.recovering_quiet_ms)); + ke_time_ms.store(0, std::memory_order_release); + quiet_until_ms.store(now + timings.recovering_quiet_ms, + std::memory_order_release); + set_state(SessionState::kRecovering, SessionError::kActivationTimeout); + continue; + } + + // Backoff / отложенный повтор. + if (now < retry_at_ms.load(std::memory_order_acquire)) { + plat::sleep_ms(kLoopTickMs); + continue; + } + + // local_reg: по keep-alive, по notify (batch/online) или первичный. + bool queue_nonempty; + { + std::lock_guard lk(queue_mu); + queue_nonempty = queue_len > 0; + } + bool notify = queue_nonempty || want_notify.exchange(false, + std::memory_order_acq_rel); + uint64_t last_reg = last_local_reg_ms.load(std::memory_order_acquire); + bool due = notify || last_reg == 0 || + now - last_reg >= timings.keepalive_ms; + if (!due) { + plat::sleep_ms(kLoopTickMs); + continue; + } + + bool first = first_reg || !reg_ok; + if (!send_local_reg(queue_nonempty, first)) { + // Транспортная ошибка: backoff. + backoff_ms = backoff_ms == 0 ? timings.backoff_base_ms + : (backoff_ms * 8) / 5; // x1.6 + if (backoff_ms > timings.backoff_max_ms) { + backoff_ms = timings.backoff_max_ms; + } + backoff_attempts++; + if (backoff_attempts >= timings.backoff_attempts) { + set_state(SessionState::kOffline, SessionError::kUnreachable); + retry_at_ms.store(now + timings.backoff_max_ms, + std::memory_order_release); + backoff_attempts = 0; + backoff_ms = 0; + } else { + retry_at_ms.store(now + backoff_ms, std::memory_order_release); + } + continue; + } + if (reg_ok) first_reg = false; + } +} + +// --------------------------------------------------------------------------- +// Session (публичный класс) +// --------------------------------------------------------------------------- + +Session* Session::create(const SessionConfig& cfg, const SessionCallbacks& cbs) { + if (cfg.host == nullptr || cfg.dsn == nullptr || cfg.lanip_key == nullptr) { + return nullptr; + } + if (strlen(cfg.lanip_key) >= sizeof(Impl::lanip_key) || + strlen(cfg.dsn) >= sizeof(Impl::dsn) || + strlen(cfg.host) >= sizeof(Impl::host)) { + return nullptr; // не помещается во внутренние копии + } + auto* impl = new (std::nothrow) Impl(); + if (impl == nullptr) return nullptr; + impl->cfg = cfg; + impl->cbs = cbs; + if (impl->cfg.max_queue == 0) impl->cfg.max_queue = 16; + if (impl->cfg.keepalive_ms == 0) impl->cfg.keepalive_ms = 15000; + impl->timings.keepalive_ms = impl->cfg.keepalive_ms; + snprintf(impl->host, sizeof(impl->host), "%s", cfg.host); + snprintf(impl->dsn, sizeof(impl->dsn), "%s", cfg.dsn); + snprintf(impl->lanip_key, sizeof(impl->lanip_key), "%s", cfg.lanip_key); + auto* s = new (std::nothrow) Session(cfg, cbs); + if (s == nullptr) { + delete impl; + return nullptr; + } + s->impl_ = impl; + return s; +} + +Session::Session(const SessionConfig&, const SessionCallbacks&) : impl_(nullptr) {} + +SessionError Session::last_error() const { + if (impl_ == nullptr) return SessionError::kNone; + return static_cast(impl_->last_error.load(std::memory_order_acquire)); +} + +SessionState Session::state() const { + if (impl_ == nullptr) return SessionState::kIdle; + return static_cast(impl_->state.load(std::memory_order_acquire)); +} + +uint16_t Session::listen_port() const { + return impl_ != nullptr ? impl_->listen_port_actual : 0; +} + +uint32_t Session::rekey_count() const { + return impl_ != nullptr ? impl_->rekeys.load(std::memory_order_relaxed) : 0; +} +uint32_t Session::pushes_ok() const { + return impl_ != nullptr ? impl_->pushes_ok.load(std::memory_order_relaxed) : 0; +} +uint32_t Session::pushes_bad() const { + return impl_ != nullptr ? impl_->pushes_bad.load(std::memory_order_relaxed) : 0; +} +uint32_t Session::commands_served() const { + return impl_ != nullptr ? impl_->cmds_served.load(std::memory_order_relaxed) : 0; +} +bool Session::had_activity() const { + return impl_ != nullptr && + impl_->had_poll_since_ke.load(std::memory_order_acquire); +} +Session::~Session() { + stop(); + delete impl_; +} + +bool Session::start() { + if (impl_ == nullptr || impl_->running.load(std::memory_order_acquire)) { + return false; + } + if (!impl_->httpd.start(impl_->cfg.listen_port, Impl::http_handler, impl_, + "fgl_session")) { + return false; + } + impl_->listen_port_actual = impl_->httpd.port(); + impl_->running.store(true, std::memory_order_release); + if (!plat::thread_create(Impl::session_thread_trampoline, impl_, + "fgl_sess", 8192, &impl_->thread)) { + impl_->running.store(false, std::memory_order_release); + impl_->httpd.stop(); + return false; + } + return true; +} + +void Session::stop() { + if (impl_ == nullptr || !impl_->running.load(std::memory_order_acquire)) { + return; + } + // Штатное завершение: DELETE-команда + notify local_reg. Сессионный поток + // ещё работает и доставит notify; ждём выдачи команды модулю. + if (impl_->state.load(std::memory_order_acquire) != + static_cast(SessionState::kKeyError)) { + delete_session(); + uint64_t deadline = plat::now_ms() + impl_->timings.delete_wait_ms; + while (plat::now_ms() < deadline && + !impl_->delete_served.load(std::memory_order_acquire)) { + plat::sleep_ms(10); + } + } + impl_->running.store(false, std::memory_order_release); + if (impl_->thread != nullptr) { + plat::thread_join(impl_->thread); + impl_->thread = nullptr; + } + impl_->httpd.stop(); + // Сброс для возможного рестарта. + impl_->delete_served.store(false, std::memory_order_release); + impl_->delete_pending.store(false, std::memory_order_release); + impl_->quiet_until_ms.store(0, std::memory_order_release); + impl_->retry_at_ms.store(0, std::memory_order_release); + impl_->ke_time_ms.store(0, std::memory_order_release); + impl_->set_state(SessionState::kIdle, SessionError::kNone); +} + +bool Session::get_property(const char* name) { + if (impl_ == nullptr || name == nullptr || strlen(name) >= kMaxName) { + return false; + } + Command cmd{}; + cmd.type = 1; + snprintf(cmd.name, sizeof(cmd.name), "%s", name); + { + std::lock_guard lk(impl_->queue_mu); + cmd.cmd_id = impl_->next_cmd_id++; + } + if (!impl_->submit(cmd)) return false; + impl_->want_notify.store(true, std::memory_order_release); + return true; +} + +bool Session::set_property(const char* name, int64_t value, + const char* base_type) { + if (impl_ == nullptr || name == nullptr || strlen(name) >= kMaxName) { + return false; + } + Command cmd{}; + cmd.type = 2; + snprintf(cmd.name, sizeof(cmd.name), "%s", name); + cmd.value = value; + snprintf(cmd.base_type, sizeof(cmd.base_type), "%s", base_type); + { + std::lock_guard lk(impl_->queue_mu); + cmd.cmd_id = impl_->next_cmd_id++; + } + if (!impl_->submit(cmd)) return false; + impl_->want_notify.store(true, std::memory_order_release); + return true; +} + +void Session::set_timings_for_test(const SessionTimings& t) { + if (impl_ == nullptr) return; + SessionTimings tmp = t; + if (tmp.keepalive_ms < 100) tmp.keepalive_ms = 100; // анти-спам + impl_->timings = tmp; +} + +bool Session::begin_batch() { + std::lock_guard lk(impl_->queue_mu); + if (impl_->batch_open) return false; + impl_->batch_open = true; + impl_->batch_len = 0; + return true; +} + +bool Session::commit_batch() { + { + std::lock_guard lk(impl_->queue_mu); + if (!impl_->batch_open) return false; + bool ok = true; + for (uint8_t i = 0; i < impl_->batch_len; i++) { + if (!impl_->enqueue_locked(impl_->batch[i])) ok = false; + } + impl_->batch_open = false; + impl_->batch_len = 0; + if (!ok) return false; + } + impl_->want_notify.store(true, std::memory_order_release); + return true; +} + +bool Session::abort_batch() { + std::lock_guard lk(impl_->queue_mu); + if (!impl_->batch_open) return false; + impl_->batch_open = false; + impl_->batch_len = 0; + return true; +} + +bool Session::delete_session() { + if (impl_ == nullptr) return false; + Command cmd{}; + cmd.type = 3; + snprintf(cmd.name, sizeof(cmd.name), "local_reg.json"); + { + std::lock_guard lk(impl_->queue_mu); + if (!impl_->enqueue_locked(cmd)) return false; + } + impl_->delete_pending.store(true, std::memory_order_release); + impl_->want_notify.store(true, std::memory_order_release); + return true; +} + +} // namespace fgl::ayla diff --git a/src/ayla/session.hpp b/src/ayla/session.hpp new file mode 100644 index 0000000..9565741 --- /dev/null +++ b/src/ayla/session.hpp @@ -0,0 +1,126 @@ +// Сессия Ayla LAN (сторона «приложения»). docs/PROTOCOL.md §4-6, §4.4. +// Потоки: httpd (входящие от модуля: key exchange/commands/datapoint) и +// session (исходящие local_reg, таймеры, backoff). Крипто-цепочки и выдача +// команд — только в httpd-потоке; session-поток читает очередь под mutex. +#pragma once + +#include +#include +#include + +#include "ayla/crypto.hpp" +#include "ayla/httpd.hpp" + +namespace fgl::ayla { + +enum class SessionState : uint8_t { + kIdle = 0, // создан, не запущен + kRegistering, // local_reg отправлен, ждём key exchange + kOnline, // сессия активна + kRecovering, // ожидание самолечения (re-key по keep-alive / активация) + kOffline, // модуль недоступен (backoff) или нет слотов + kKeyError, // lanip_key_id не совпал — требуется смена конфига +}; + +enum class SessionError : int { + kNone = 0, + kNoSlot = 1, // 503: оба слота модуля заняты + kUnreachable = 2, // transport/backoff + kKeyMismatch = 3, // key_id != lanip_key_id (state = kKeyError) + kBadKeyExchange = 4, // ver/proto/sec не поддержаны + kActivationTimeout = 5,// KE прошёл, опроса commands.json нет (>5 c) + kDecryptFailed = 6, // подпись/расшифровка push не сошлись (ждём re-key) +}; + +// Событие обновления свойства (push модуля). Значение — какой-то один тип. +struct PropertyEvent { + char name[40]; + bool is_int = false; + int64_t int_value = 0; + bool is_bool = false; + bool bool_value = false; + char str_value[64]; // используется, если !is_int && !is_bool + int cmd_id = -1; // из ?cmd_id=N (ответ на GET), иначе -1 + int status = 0; // из ?status=200 + int64_t seq_no = 0; +}; + +struct SessionConfig { + const char* host = nullptr; // DNS-имя или IP модуля + uint16_t device_port = 80; // порт local_reg модуля + const char* dsn = nullptr; // "AC000W00XXXXXXX" + const char* lanip_key = nullptr; // base64-строка как есть + uint32_t lanip_key_id = 0; + uint16_t listen_port = 10275; // 0 — любой свободный + uint32_t keepalive_ms = 15000; + uint8_t max_queue = 16; // лимит очереди команд +}; + +struct SessionCallbacks { + // КОНТРАКТ: колбэки приходят из потоков ядра (httpd и/или session), + // возможно перекрытие во времени; обязаны быть быстрыми и реентерабельными. + // Вызывать stop() из колбэка запрещено (deadlock на join). + void (*on_state)(void* ctx, SessionState st, SessionError err); + void (*on_property)(void* ctx, const PropertyEvent& ev); + void* ctx = nullptr; +}; + +// Тайминги поведения (PROTOCOL §4.3-4.4; проверено на приборе). +struct SessionTimings { + uint32_t keepalive_ms = 15000; // период local_reg + uint32_t activation_timeout_ms = 5000; // нет poll после KE + uint32_t recovering_quiet_ms = 50000; // пауза при десинке/«KE без poll» + // (> порога возврата модуля ~44-50с) + uint32_t no_slot_retry_ms = 60000; // повтор после 503 + uint32_t backoff_base_ms = 1000; // transport backoff, шаг x1.6 + uint32_t backoff_max_ms = 60000; + uint8_t backoff_attempts = 6; + uint32_t delete_wait_ms = 2000; +}; + +class Session { + public: + static Session* create(const SessionConfig& cfg, const SessionCallbacks& cbs); + ~Session(); + + Session(const Session&) = delete; + Session& operator=(const Session&) = delete; + + bool start(); + // Штатное завершение: DELETE-команда + local_reg notify, ожидание выдачи, + // остановка потоков. state -> kIdle. + void stop(); + + SessionState state() const; + SessionError last_error() const; + uint16_t listen_port() const; + // Телеметрия (диагностика). + uint32_t rekey_count() const; + uint32_t pushes_ok() const; + uint32_t pushes_bad() const; + uint32_t commands_served() const; + bool had_activity() const; + + // ---- Команды (потокобезопасны; кладутся в очередь с coalescing) ---- + // GET-команда: свойство придёт on_property (cmd_id совпадает). + bool get_property(const char* name); + // SET-команда (integer/boolean как int64). + bool set_property(const char* name, int64_t value, + const char* base_type = "integer"); + // Пакет: собрать несколько команд, один notify на commit. + bool begin_batch(); + bool commit_batch(); + bool abort_batch(); + // DELETE local_reg.json/delete_session (для stop() и ручного завершения). + bool delete_session(); + + // Тест-хук: тайминги. ТОЛЬКО до start() (после — читаются потоками ядра). + void set_timings_for_test(const SessionTimings& t); + + private: + Session(const SessionConfig& cfg, const SessionCallbacks& cbs); + struct Impl; + Impl* impl_; +}; + +} // namespace fgl::ayla diff --git a/tests/CMakeLists.txt b/tests/CMakeLists.txt index 86c44f3..9088672 100644 --- a/tests/CMakeLists.txt +++ b/tests/CMakeLists.txt @@ -20,9 +20,23 @@ function(fgl_add_test name) add_test(NAME ${name} COMMAND test_${name}) endfunction() +# Раннер сессии — не тест, приложение для интеграционных сценариев +# (запускается tests/ayla/test_session_mock.py через ctest). +add_executable(session_runner ayla/session_runner.cpp) +target_compile_features(session_runner PRIVATE cxx_std_20) +target_link_libraries(session_runner PRIVATE fgl-aircon) +target_include_directories(session_runner PRIVATE "${CMAKE_SOURCE_DIR}/src") + fgl_add_test(ayla_platform ayla/test_platform.cpp) fgl_add_test(ayla_httpd ayla/test_httpd.cpp) fgl_add_test(ayla_crypto ayla/test_crypto.cpp) fgl_add_test(ayla_envelope ayla/test_envelope.cpp) fgl_add_test(ayla_json ayla/test_json.cpp) fgl_add_test(ayla_httpc ayla/test_httpc.cpp) + +# Интеграционные сценарии с mock-модулем (python, stdlib-only). +find_package(Python3 COMPONENTS Interpreter REQUIRED) +add_test(NAME ayla_session_mock + COMMAND ${Python3_EXECUTABLE} ${CMAKE_CURRENT_SOURCE_DIR}/ayla/test_session_mock.py + $ ${CMAKE_CURRENT_SOURCE_DIR}/ayla/mock_ac.py) +set_tests_properties(ayla_session_mock PROPERTIES TIMEOUT 180) diff --git a/tests/ayla/mock_ac.py b/tests/ayla/mock_ac.py new file mode 100644 index 0000000..91586da --- /dev/null +++ b/tests/ayla/mock_ac.py @@ -0,0 +1,420 @@ +#!/usr/bin/env python3 +"""Мок-модуль кондиционера (сторона устройства) для интеграционных тестов +сессии. stdlib-only: AES-256 реализован на чистом python (объёмы крошечные). + +Сценарные флаги: + --503 всегда отвечать 503 на local_reg (нет слотов) + --no-poll key exchange без опроса commands.json (зависание) + --rekey-every N ре-кей на каждый N-й local_reg (N=1 — каждый) + --stale-gap S ре-кей, если зазор между local_reg >= S секунд + (эмуляция «вернувшегося» приложения; по умолчанию 44) + --garbage-pushes N первые N push с отрезанным блоком шифротекста + (входящая цепочка расходится на 1 сообщение) + --break-outbound N N раз «не заметить» ответ commands.json (исходящий + десинк): затем подпись наших команд не сойдётся — + мок, как реальный модуль, ре-кает на следующем local_reg + --fail-pushes N первые N push'ей с испорченной подписью + --push-every S спонтанный push свойства tick каждые S секунд + --fail-first-ke первый key_exchange с ver=2 (ожидаем 426) +Вывод (stdout, строки): + REG notify=0|1 + KE + CMD [value] + PUSH + DELETE +""" +import argparse +import base64 +import hashlib +import hmac +import http.server +import json +import random +import socket +import string +import sys +import threading +import time + +# --------------------------------------------------------------------------- +# Чистый python AES-256 (encrypt/decrypt block), CBC поверх. +# --------------------------------------------------------------------------- +_SBOX = [ + 0x63,0x7c,0x77,0x7b,0xf2,0x6b,0x6f,0xc5,0x30,0x01,0x67,0x2b,0xfe,0xd7,0xab,0x76, + 0xca,0x82,0xc9,0x7d,0xfa,0x59,0x47,0xf0,0xad,0xd4,0xa2,0xaf,0x9c,0xa4,0x72,0xc0, + 0xb7,0xfd,0x93,0x26,0x36,0x3f,0xf7,0xcc,0x34,0xa5,0xe5,0xf1,0x71,0xd8,0x31,0x15, + 0x04,0xc7,0x23,0xc3,0x18,0x96,0x05,0x9a,0x07,0x12,0x80,0xe2,0xeb,0x27,0xb2,0x75, + 0x09,0x83,0x2c,0x1a,0x1b,0x6e,0x5a,0xa0,0x52,0x3b,0xd6,0xb3,0x29,0xe3,0x2f,0x84, + 0x53,0xd1,0x00,0xed,0x20,0xfc,0xb1,0x5b,0x6a,0xcb,0xbe,0x39,0x4a,0x4c,0x58,0xcf, + 0xd0,0xef,0xaa,0xfb,0x43,0x4d,0x33,0x85,0x45,0xf9,0x02,0x7f,0x50,0x3c,0x9f,0xa8, + 0x51,0xa3,0x40,0x8f,0x92,0x9d,0x38,0xf5,0xbc,0xb6,0xda,0x21,0x10,0xff,0xf3,0xd2, + 0xcd,0x0c,0x13,0xec,0x5f,0x97,0x44,0x17,0xc4,0xa7,0x7e,0x3d,0x64,0x5d,0x19,0x73, + 0x60,0x81,0x4f,0xdc,0x22,0x2a,0x90,0x88,0x46,0xee,0xb8,0x14,0xde,0x5e,0x0b,0xdb, + 0xe0,0x32,0x3a,0x0a,0x49,0x06,0x24,0x5c,0xc2,0xd3,0xac,0x62,0x91,0x95,0xe4,0x79, + 0xe7,0xc8,0x37,0x6d,0x8d,0xd5,0x4e,0xa9,0x6c,0x56,0xf4,0xea,0x65,0x7a,0xae,0x08, + 0xba,0x78,0x25,0x2e,0x1c,0xa6,0xb4,0xc6,0xe8,0xdd,0x74,0x1f,0x4b,0xbd,0x8b,0x8a, + 0x70,0x3e,0xb5,0x66,0x48,0x03,0xf6,0x0e,0x61,0x35,0x57,0xb9,0x86,0xc1,0x1d,0x9e, + 0xe1,0xf8,0x98,0x11,0x69,0xd9,0x8e,0x94,0x9b,0x1e,0x87,0xe9,0xce,0x55,0x28,0xdf, + 0x8c,0xa1,0x89,0x0d,0xbf,0xe6,0x42,0x68,0x41,0x99,0x2d,0x0f,0xb0,0x54,0xbb,0x16] +_RCON = [0x01,0x02,0x04,0x08,0x10,0x20,0x40,0x80,0x1b,0x36,0x6c,0xd8,0xab,0x4d] +_INV_SBOX = [0]*256 +for _i, _b in enumerate(_SBOX): + _INV_SBOX[_b] = _i + +def _xtime(a): + a <<= 1 + if a & 0x100: + a = (a ^ 0x1b) & 0xff + return a + +def _expand_key(key): # 32 байта -> 60 слов по 4 байта (flat список) + w = list(key) + for i in range(32, 240, 4): + t = w[i-4:i] + if i % 32 == 0: + t = t[1:] + t[:1] + t = [_SBOX[b] for b in t] + t[0] ^= _RCON[i//32 - 1] + elif i % 32 == 16: + t = [_SBOX[b] for b in t] + w += [w[i-32+j] ^ t[j] for j in range(4)] + return w + +def _aes_encrypt_block(w, block): + s = list(block) + def add_round_key(r): + for i in range(16): + s[i] ^= w[r*16 + i] + def sub_shift(): + # SubBytes + ShiftRows (строка r — байты r, r+4, r+8, r+12 — влево на r) + t = [_SBOX[b] for b in s] + out = [0]*16 + for r in range(4): + for c in range(4): + out[r + 4*c] = t[r + 4*((c + r) % 4)] + for i in range(16): + s[i] = out[i] + def mix(): + t = [0]*16 + for c in range(4): + col = s[c*4:c*4+4] + t[c*4+0] = _xtime(col[0]) ^ _xtime(col[1]) ^ col[1] ^ col[2] ^ col[3] + t[c*4+1] = col[0] ^ _xtime(col[1]) ^ _xtime(col[2]) ^ col[2] ^ col[3] + t[c*4+2] = col[0] ^ col[1] ^ _xtime(col[2]) ^ _xtime(col[3]) ^ col[3] + t[c*4+3] = _xtime(col[0]) ^ col[0] ^ col[1] ^ col[2] ^ _xtime(col[3]) + for i in range(16): + s[i] = t[i] + add_round_key(0) + for rnd in range(1, 14): + sub_shift(); mix(); add_round_key(rnd) + sub_shift(); add_round_key(14) + return bytes(s) + +def _aes_decrypt_block(w, block): + inv_sbox = _INV_SBOX + def inv_sub_shift(s): + t = [inv_sbox[b] for b in s] + out = [0]*16 + for r in range(4): + for c in range(4): + out[r + 4*c] = t[r + 4*((c - r) % 4)] # инверсия сдвига влево на r + return out + def inv_mix(s): + def mul(a, b): + p = 0 + for _ in range(8): + if b & 1: + p ^= a + hi = a & 0x80 + a = (a << 1) & 0xff + if hi: + a ^= 0x1b + b >>= 1 + return p + t = [0]*16 + for c in range(4): + col = s[c*4:c*4+4] + t[c*4+0] = mul(col[0],14) ^ mul(col[1],11) ^ mul(col[2],13) ^ mul(col[3],9) + t[c*4+1] = mul(col[0],9) ^ mul(col[1],14) ^ mul(col[2],11) ^ mul(col[3],13) + t[c*4+2] = mul(col[0],13) ^ mul(col[1],9) ^ mul(col[2],14) ^ mul(col[3],11) + t[c*4+3] = mul(col[0],11) ^ mul(col[1],13) ^ mul(col[2],9) ^ mul(col[3],14) + return t + s = list(block) + def add_round_key(r): + for i in range(16): + s[i] ^= w[r*16 + i] + add_round_key(14) + for rnd in range(13, 0, -1): + s = inv_sub_shift(s) # InvShiftRows + InvSubBytes (коммутируют) + add_round_key(rnd) + s = inv_mix(s) + s = inv_sub_shift(s) + add_round_key(0) + return bytes(s) + +class PyAes: + def __init__(self, key): + self.w = _expand_key(key) + def cbc_encrypt(self, iv, data): + out = b"" + prev = iv + for i in range(0, len(data), 16): + blk = data[i:i+16] + blk = bytes(a ^ b for a, b in zip(blk, prev)) + prev = _aes_encrypt_block(self.w, blk) + out += prev + return out, prev + def cbc_decrypt(self, iv, data): + out = b"" + prev = iv + for i in range(0, len(data), 16): + blk = data[i:i+16] + dec = _aes_decrypt_block(self.w, blk) + out += bytes(a ^ b for a, b in zip(dec, prev)) + prev = blk + return out, prev + +# --------------------------------------------------------------------------- +class MockCrypto: + """Ключи одной стороны мока: dev (исходящие push) и app (входящие команды).""" + def __init__(self, lanip_key, rnd1, rnd2, t1, t2): + k = lanip_key.encode() + b1, b2 = rnd1.encode(), rnd2.encode() + s1, s2 = str(t1).encode(), str(t2).encode() + def m(msg, suf): + msg = msg + bytes([suf]) + return hmac.digest(k, hmac.digest(k, msg, "sha256") + msg, "sha256") + A, D = b1 + b2 + s1 + s2, b2 + b1 + s2 + s1 + self.dev_sign, self.dev_aes = m(D, 0x30), PyAes(m(D, 0x31)) + self.app_sign, self.app_aes = m(A, 0x30), PyAes(m(A, 0x31)) + self.dev_iv, self.app_iv = m(D, 0x32)[:16], m(A, 0x32)[:16] + + def pack_push(self, seq, data_json): + plain = json.dumps({"seq_no": seq, "data": data_json}, + separators=(",", ":")).encode() + sign = base64.b64encode(hmac.digest(self.dev_sign, plain, "sha256")).decode() + n = ((len(plain) + 1 + 15) // 16) * 16 + ct, self.dev_iv = self.dev_aes.cbc_encrypt(self.dev_iv, plain.ljust(n, b"\x00")) + enc = base64.b64encode(ct).decode() + return json.dumps({"enc": enc, "sign": sign}, separators=(",", ":")) + + def unpack_command(self, body): + d = json.loads(body) + pt, self.app_iv = self.app_aes.cbc_decrypt( + self.app_iv, base64.b64decode(d["enc"])) + pt = pt.rstrip(b"\x00") + # Подпись проверяется ДО разбора JSON (мусор не парсим). + ok = base64.b64encode(hmac.digest(self.app_sign, pt, "sha256")).decode() == d["sign"] + if not ok: + return False, None + return True, json.loads(pt.decode()) + +# --------------------------------------------------------------------------- +class Mock: + def __init__(self, args): + self.args = args + self.lock = threading.Lock() + self.props = {"operation_mode": 6, "fan_speed": 4, "tick": 0} + self.crypto = None + self.app_addr = None # (ip, port) приложения + self.push_seq = 0 + self.reg_count = 0 + self.last_reg_time = None + self.fail_pushes = args.fail_pushes + self.garbage_pushes = args.garbage_pushes + self.miss_response = args.break_outbound + self.stop = threading.Event() + + # ---------- исходящие к приложению ---------- + def http_call(self, method, path, body=b"", timeout=5): + ip, port = self.app_addr + c = socket.create_connection((ip, port), timeout=timeout) + req = (f"{method} {path} HTTP/1.1\r\nHost: {ip}\r\n" + f"Content-Type: application/json\r\n" + f"Content-Length: {len(body)}\r\nConnection: close\r\n\r\n").encode() + body + c.sendall(req) + raw = b"" + while True: + chunk = c.recv(4096) + if not chunk: + break + raw += chunk + c.close() + head, _, resp_body = raw.partition(b"\r\n\r\n") + status = int(head.split(b" ")[1]) + return status, resp_body + + def do_key_exchange(self): + rnd1 = "".join(random.choice(string.ascii_letters + string.digits + "+/") + for _ in range(16)) + t1 = int(time.monotonic_ns() // 1000) + ver = 2 if (self.args.fail_first_ke and self.reg_count == 1) else 1 + body = json.dumps({"key_exchange": { + "ver": ver, "proto": 1, "key_id": self.args.key_id, + "random_1": rnd1, "time_1": t1, "sec": ""}}, + separators=(",", ":")).encode() + status, resp = self.http_call("POST", "/local_lan/key_exchange.json", body) + print(f"KE {rnd1} -> {status}", flush=True) + if status != 200: + return False + d = json.loads(resp) + self.crypto = MockCrypto(self.args.lanip_key, rnd1, d["random_2"], + t1, d["time_2"]) + return True + + def poll_commands(self): + """Опрашивает commands.json, пока 206; исполняет команды.""" + while True: + status, body = self.http_call("GET", "/local_lan/commands.json") + if status != 200 and status != 206: + print(f"POLL -> {status}", flush=True) + return + if not self.crypto: + return + if self.miss_response > 0: + # «Модуль не получил/не расшифровал ответ»: цепочка приложения + # ушла, у мока нет — исходящий десинк. + self.miss_response -= 1 + print("CMD skipped (outbound desync)", flush=True) + return + ok, payload = self.crypto.unpack_command(body) + if not ok: + print("CMD bad-sign", flush=True) + # Реальный модуль: на следующем local_reg — key exchange. + self.crypto = None + return + data = payload.get("data", {}) + cmds = data.get("cmds", []) + props = data.get("properties", []) + if cmds: + for c in cmds: + cmd = c.get("cmd", {}) + if cmd.get("method") == "DELETE": + print("DELETE", flush=True) + return + res = cmd.get("resource", "") + name = res.split("name=")[-1] + cid = cmd.get("cmd_id", -1) + print(f"CMD GET {name} cid={cid}", flush=True) + self.push_datapoint(name, cid=cid) + elif props: + for p in props: + pr = p.get("property", {}) + self.props[pr.get("name", "?")] = pr.get("value") + print(f"CMD SET {pr.get('name')}={pr.get('value')}", flush=True) + elif not data: + return + if status == 200: + return + + def push_datapoint(self, name, cid=-1, corrupt=False, garbage=False): + self.push_seq += 1 + body = self.crypto.pack_push(self.push_seq - 1, + {"name": name, "value": self.props.get(name, 0)}) + if garbage: + # РЕАЛЬНЫЙ десинк CBC: отрезать последний блок шифротекста — + # цепочка мока ушла на блок дальше, приложение отстанет. + d = json.loads(body) + ct = base64.b64decode(d["enc"]) + d["enc"] = base64.b64encode(ct[:-16]).decode() + body = json.dumps(d, separators=(",", ":")) + if corrupt: + d = json.loads(body) + s = d["sign"] + d["sign"] = ("A" if s[-2] != "A" else "B") + s[1:] + body = json.dumps(d, separators=(",", ":")) + path = "/local_lan/property/datapoint.json" + if cid >= 0: + path += f"?cmd_id={cid}&status=200" + status, _ = self.http_call("POST", path, body.encode()) + print(f"PUSH {name} -> {status}{' CORRUPT' if corrupt else ''}", flush=True) + + # ---------- входящие local_reg ---------- + def handle_local_reg(self, body): + self.reg_count += 1 + d = json.loads(body)["local_reg"] + notify = d.get("notify", 0) + self.app_addr_json = d + self.app_addr = (d["ip"], d["port"]) + print(f"REG {'first' if self.reg_count == 1 else 'put'} notify={notify}", + flush=True) + now = time.monotonic() + stale = (self.last_reg_time is not None and + now - self.last_reg_time >= self.args.stale_gap) + self.last_reg_time = now + if self.args.http503: + return 503 + need_ke = self.crypto is None or stale or ( + self.args.rekey_every and self.reg_count % self.args.rekey_every == 0) + if need_ke and not self.do_key_exchange(): + return 202 + if not self.args.no_poll: + self.poll_commands() + return 202 + + def spontaneous_loop(self): + while not self.stop.wait(self.args.push_every or 10): + if self.crypto and not self.args.no_poll: + with self.lock: + self.props["tick"] += 1 + corrupt = self.fail_pushes > 0 + if corrupt: + self.fail_pushes -= 1 + garbage = self.garbage_pushes > 0 + if garbage: + self.garbage_pushes -= 1 + self.push_datapoint("tick", corrupt=corrupt, garbage=garbage) + + +def main(): + ap = argparse.ArgumentParser() + ap.add_argument("--port", type=int, required=True) + ap.add_argument("--lanip-key", required=True) + ap.add_argument("--key-id", type=int, required=True) + ap.add_argument("--503", dest="http503", action="store_true") + ap.add_argument("--no-poll", action="store_true") + ap.add_argument("--rekey-every", type=int, default=0) + ap.add_argument("--fail-pushes", type=int, default=0) + ap.add_argument("--garbage-pushes", type=int, default=0) + ap.add_argument("--break-outbound", type=int, default=0) + ap.add_argument("--stale-gap", type=float, default=44.0) + ap.add_argument("--push-every", type=float, default=0) + ap.add_argument("--fail-first-ke", action="store_true") + args = ap.parse_args() + + mock = Mock(args) + lock = mock.lock + + class H(http.server.BaseHTTPRequestHandler): + protocol_version = "HTTP/1.1" + def log_message(self, *a): + pass + def do_POST(self): + n = int(self.headers.get("Content-Length") or 0) + body = self.rfile.read(n) + if self.path.startswith("/local_reg.json"): + with lock: + code = mock.handle_local_reg(body) + self.send_response(code) + self.send_header("Content-Length", "0") + self.end_headers() + return + self.send_response(404) + self.send_header("Content-Length", "0") + self.end_headers() + def do_PUT(self): + self.do_POST() + + srv = http.server.ThreadingHTTPServer(("127.0.0.1", args.port), H) + threading.Thread(target=srv.serve_forever, daemon=True).start() + if args.push_every: + threading.Thread(target=mock.spontaneous_loop, daemon=True).start() + print("READY", flush=True) + try: + while True: + time.sleep(0.5) + except KeyboardInterrupt: + pass + +if __name__ == "__main__": + main() diff --git a/tests/ayla/session_runner.cpp b/tests/ayla/session_runner.cpp new file mode 100644 index 0000000..390c1fb --- /dev/null +++ b/tests/ayla/session_runner.cpp @@ -0,0 +1,135 @@ +// Тестовый раннер сессии: поднимает fgl::ayla::Session против mock_ac.py +// (или реального модуля) и печатает события строками в stdout: +// STATE — смена состояния +// PROP — push свойства +// DELETED — delete_session забран модулем +// STATS ... — телеметрия (перед выходом) +// Управление окружением: +// RUNNER_KEEPALIVE_MS, RUNNER_ACTIVATION_MS, RUNNER_QUIET_MS, +// RUNNER_NOSLOT_RETRY_MS, RUNNER_DELETE_WAIT_MS +// RUNNER_SET_NAME/VALUE/AT — одиночная SET через AT секунд после старта. +#include +#include +#include +#include +#include + +#include "ayla/session.hpp" + +using fgl::ayla::PropertyEvent; +using fgl::ayla::Session; +using fgl::ayla::SessionError; +using fgl::ayla::SessionState; + +static const char* state_name(SessionState st) { + switch (st) { + case SessionState::kIdle: return "idle"; + case SessionState::kRegistering: return "registering"; + case SessionState::kOnline: return "online"; + case SessionState::kRecovering: return "recovering"; + case SessionState::kOffline: return "offline"; + case SessionState::kKeyError: return "key_error"; + } + return "?"; +} + +static uint32_t env_u32(const char* name, uint32_t def) { + const char* v = getenv(name); + return v != nullptr ? static_cast(atoi(v)) : def; +} + +static void on_state(void*, SessionState st, SessionError err) { + printf("STATE %s %d\n", state_name(st), static_cast(err)); + fflush(stdout); +} + +static void on_property(void*, const PropertyEvent& ev) { + if (ev.is_int) { + printf("PROP %s %d %d i:%lld\n", ev.name, ev.cmd_id, ev.status, + static_cast(ev.int_value)); + } else if (ev.is_bool) { + printf("PROP %s %d %d b:%d\n", ev.name, ev.cmd_id, ev.status, + ev.bool_value ? 1 : 0); + } else { + printf("PROP %s %d %d s:%s\n", ev.name, ev.cmd_id, ev.status, ev.str_value); + } + fflush(stdout); +} + +int main(int argc, char** argv) { + if (argc < 7) { + fprintf(stderr, + "usage: %s " + " [prop ...]\n", + argv[0]); + return 2; + } + fgl::ayla::SessionConfig cfg{}; + cfg.host = argv[1]; + cfg.device_port = static_cast(atoi(argv[2])); + cfg.listen_port = static_cast(atoi(argv[3])); + cfg.dsn = argv[4]; + cfg.lanip_key = argv[5]; + cfg.lanip_key_id = static_cast(atoi(argv[6])); + int duration_sec = atoi(argv[7]); + + fgl::ayla::SessionCallbacks cbs{}; + cbs.on_state = on_state; + cbs.on_property = on_property; + + Session* s = Session::create(cfg, cbs); + if (s == nullptr) { + fprintf(stderr, "create failed\n"); + return 2; + } + fgl::ayla::SessionTimings timings{}; + timings.keepalive_ms = env_u32("RUNNER_KEEPALIVE_MS", timings.keepalive_ms); + timings.activation_timeout_ms = + env_u32("RUNNER_ACTIVATION_MS", timings.activation_timeout_ms); + timings.recovering_quiet_ms = + env_u32("RUNNER_QUIET_MS", timings.recovering_quiet_ms); + timings.no_slot_retry_ms = env_u32("RUNNER_NOSLOT_RETRY_MS", + timings.no_slot_retry_ms); + timings.delete_wait_ms = env_u32("RUNNER_DELETE_WAIT_MS", + timings.delete_wait_ms); + s->set_timings_for_test(timings); + + if (!s->start()) { + fprintf(stderr, "start failed\n"); + return 2; + } + + // Начальная синхронизация: пакет GET всех свойств. + if (argc > 8) { + s->begin_batch(); + for (int i = 8; i < argc; i++) { + s->get_property(argv[i]); + } + s->commit_batch(); + } + + const char* set_name = getenv("RUNNER_SET_NAME"); + const char* set_val = getenv("RUNNER_SET_VALUE"); + uint32_t set_at = env_u32("RUNNER_SET_AT", 0); + if (set_name != nullptr && set_val != nullptr && set_at > 0) { + std::this_thread::sleep_for(std::chrono::seconds(set_at)); + s->set_property(set_name, atoll(set_val)); + printf("SET_DONE %s=%s\n", set_name, set_val); + fflush(stdout); + } + + auto deadline = std::chrono::steady_clock::now() + + std::chrono::seconds(duration_sec); + while (std::chrono::steady_clock::now() < deadline) { + std::this_thread::sleep_for(std::chrono::milliseconds(100)); + } + + s->stop(); // внутри: delete_session + ожидание выдачи + printf("DELETED\n"); + printf("STATS rekeys=%u pushes_ok=%u pushes_bad=%u cmds=%u state=%s\n", + s->rekey_count(), s->pushes_ok(), s->pushes_bad(), + s->commands_served(), state_name(s->state())); + fflush(stdout); + delete s; + return 0; +} diff --git a/tests/ayla/test_session_mock.py b/tests/ayla/test_session_mock.py new file mode 100644 index 0000000..80c5b43 --- /dev/null +++ b/tests/ayla/test_session_mock.py @@ -0,0 +1,295 @@ +#!/usr/bin/env python3 +"""Интеграционные сценарии сессии против mock-модуля (mock_ac.py). + +usage: test_session_mock.py +Сценарии: normal, rekey, fail_push, slot_503, no_poll, bad_ke, set_get. +""" +import os +import socket +import subprocess +import sys +import tempfile +import threading +import time + +RUNNER = sys.argv[1] +MOCK = sys.argv[2] + +# Синтетический ключ (НЕ боевой). +LANIP_KEY = "TW9ja0tleUFDbkdvMTIzNDU2Nzg5MDEyMw==" +KEY_ID = 64201 +DSN = "AC000W00MOCK0001" + + +def free_port(): + s = socket.socket() + s.bind(("127.0.0.1", 0)) + port = s.getsockname()[1] + s.close() + return port + + +class Proc: + """Процесс с построчным логом stdout.""" + + def __init__(self, argv, env=None): + self.lines = [] + self.lock = threading.Lock() + self.ev = threading.Event() + self.proc = subprocess.Popen( + argv, env=env, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, + text=True) + threading.Thread(target=self._reader, daemon=True).start() + + def _reader(self): + for line in self.proc.stdout: + with self.lock: + self.lines.append(line.rstrip("\n")) + self.ev.set() + + def wait_line(self, prefix, timeout, exclude_prefix=None): + """Ждёт строку, начинающуюся с prefix (или содержащую, если prefix + начинается с '~').""" + contains = prefix.startswith("~") + needle = prefix[1:] if contains else prefix + deadline = time.time() + timeout + seen = 0 + while time.time() < deadline: + with self.lock: + for i, line in enumerate(self.lines): + if i < seen: + continue + hit = (needle in line) if contains else line.startswith(prefix) + if hit: + if exclude_prefix and line.startswith(exclude_prefix): + continue + return line + seen = i + 1 + self.ev.wait(0.2) + return None + + def all_lines(self): + with self.lock: + return list(self.lines) + + def stop(self, timeout=10): + self.proc.terminate() + try: + self.proc.wait(timeout) + except subprocess.TimeoutExpired: + self.proc.kill() + + +def run_mock(extra, port): + p = Proc([sys.executable, MOCK, "--port", str(port), "--lanip-key", LANIP_KEY, + "--key-id", str(KEY_ID)] + extra) + if p.wait_line("READY", 10) is None: + raise RuntimeError("mock не стартовал: " + "\n".join(p.all_lines())) + return p + + +def run_runner(mock_port, listen_port, props, duration, env_extra=None, + timeout=None): + env = dict(os.environ) + if env_extra: + env.update(env_extra) + argv = [RUNNER, "127.0.0.1", str(mock_port), str(listen_port), DSN, + LANIP_KEY, str(KEY_ID), str(duration)] + props + return Proc(argv, env=env), (timeout or duration + 25) + + +FAILURES = [] + + +def check(cond, what): + tag = "ok " if cond else "FAIL" + print(f" [{tag}] {what}") + if not cond: + FAILURES.append(what) + + +def scenario_normal(): + print("== scenario: normal (установка, батч GET, push, delete)") + mock_port, listen = free_port(), free_port() + mock = run_mock([], mock_port) + runner, timeout = run_runner(mock_port, listen, + ["operation_mode", "fan_speed"], 3) + line = runner.wait_line("STATE online", 10) + check(line is not None, "переход в online") + p1 = runner.wait_line("PROP operation_mode", 10) + p2 = runner.wait_line("PROP fan_speed", 10) + check(p1 is not None and "cid=" not in p1 and " i:6" in p1, f"push #1: {p1}") + check(p2 is not None and " i:4" in p2, f"push #2: {p2}") + check(runner.wait_line("DELETED", timeout) is not None, "delete_session") + check(runner.wait_line("STATS", 5) is not None, "STATS") + check("rekeys=1" in " ".join(runner.all_lines()), "ровно 1 key exchange") + check("CMD GET operation_mode" in "\n".join(mock.all_lines()), + "mock исполнил GET") + check("DELETE" in "\n".join(mock.all_lines()), "mock получил DELETE") + runner.stop() + mock.stop() + + +def scenario_rekey(): + print("== scenario: rekey по local_reg (chain непрерывна)") + mock_port, listen = free_port(), free_port() + mock = run_mock(["--rekey-every", "2", "--push-every", "1"], mock_port) + runner, timeout = run_runner(mock_port, listen, [], 6, + env_extra={"RUNNER_KEEPALIVE_MS": "1000"}) + check(runner.wait_line("STATE online", 10) is not None, "online") + t1 = runner.wait_line("PROP tick", 15) + t2 = runner.wait_line("PROP tick", 15) + check(t1 is not None and t2 is not None, "спонтанные push после re-key") + runner.wait_line("DELETED", timeout) + stats = "\n".join(runner.all_lines()) + check("rekeys=" in stats and "pushes_bad=0" in stats, + f"без потерянных push: {stats.splitlines()[-1] if stats else '?'}") + check("CMD bad-sign" not in "\n".join(mock.all_lines()), + "mock проверял подписи наших команд") + runner.stop() + mock.stop() + + +def scenario_fail_push(): + print("== scenario: push с битой подписью -> recovering -> online") + mock_port, listen = free_port(), free_port() + mock = run_mock(["--push-every", "1", "--fail-pushes", "1"], mock_port) + runner, timeout = run_runner(mock_port, listen, [], 8, + env_extra={"RUNNER_KEEPALIVE_MS": "1000"}) + check(runner.wait_line("STATE online", 10) is not None, "online") + rec = runner.wait_line("STATE recovering 6", 15) + check(rec is not None, f"recovering после битой подписи: {rec}") + # Следующий push расшифровывается (цепочка в синке) -> online. + check(runner.wait_line("STATE online 0", 15) is not None, + "online восстановлен следующим валидным push") + check(runner.wait_line("STATS", timeout) is not None, "STATS получены") + stats = "\n".join(runner.all_lines()) + stats_line = [l for l in stats.splitlines() if l.startswith("STATS")] + check("pushes_bad=1" in stats, f"учтён 1 битый push: {stats_line}") + runner.stop() + mock.stop() + + +def scenario_slot_503(): + print("== scenario: 503 — нет слотов") + mock_port, listen = free_port(), free_port() + mock = run_mock(["--503"], mock_port) + runner, timeout = run_runner(mock_port, listen, [], 4, + env_extra={"RUNNER_NOSLOT_RETRY_MS": "2000"}) + line = runner.wait_line("STATE offline 1", 10) + check(line is not None, "offline с ошибкой NoSlot") + regs = [l for l in mock.all_lines() if l.startswith("REG")] + check(len(regs) >= 1, "local_reg доходил до модуля") + runner.stop() + mock.stop() + + +def scenario_no_poll(): + print("== scenario: KE без poll — активация не наступила, тихая пауза") + mock_port, listen = free_port(), free_port() + mock = run_mock(["--no-poll"], mock_port) + runner, timeout = run_runner(mock_port, listen, [], 5, + env_extra={"RUNNER_ACTIVATION_MS": "1500", + "RUNNER_QUIET_MS": "3000", + "RUNNER_KEEPALIVE_MS": "800"}) + line = runner.wait_line("STATE recovering 5", 15) + check(line is not None, "recovering с ActivationTimeout") + time.sleep(1.0) + regs = [l for l in mock.all_lines() if l.startswith("REG")] + # Первый reg + возможно второй до таймаута активации; после — тишина. + check(len(regs) <= 3, f"нет спама local_reg (получено {len(regs)})") + runner.stop() + mock.stop() + + +def scenario_bad_ke(): + print("== scenario: key_exchange ver=2 -> key_error") + mock_port, listen = free_port(), free_port() + mock = run_mock(["--fail-first-ke"], mock_port) + runner, timeout = run_runner(mock_port, listen, [], 3) + line = runner.wait_line("STATE key_error 4", 10) + check(line is not None, "key_error с BadKeyExchange") + check(mock.wait_line("~-> 426", 10) is not None, "mock получил 426") + runner.stop() + mock.stop() + + +def scenario_chain_divergence_heals(): + print("== scenario: входящая цепочка расходится на 1 сообщение и сходится") + mock_port, listen = free_port(), free_port() + mock = run_mock(["--push-every", "1", "--garbage-pushes", "1"], mock_port) + runner, timeout = run_runner(mock_port, listen, [], 8, + env_extra={"RUNNER_KEEPALIVE_MS": "1000"}) + check(runner.wait_line("STATE online", 10) is not None, "online") + check(runner.wait_line("STATE recovering 6", 15) is not None, + "recovering после урезанного push") + # CBC-состояние = последний шифроблок: следующий полный push сходится. + check(runner.wait_line("STATE online 0", 15) is not None, + "online восстановлен следующим push (цепочка сошлась)") + stats = runner.wait_line("STATS", timeout) + # Урезанный push + один residual (первый блок следующего) = 2 потери, + # далее цепочка сошлась (состояние CBC = последний шифроблок). + check(stats is not None and "pushes_bad=2" in stats, + f"ровно 2 потери: {stats}") + runner.stop() + mock.stop() + + +def scenario_outbound_desync_rekey(): + print("== scenario: исходящий десинк -> модуль ре-кает на local_reg") + mock_port, listen = free_port(), free_port() + mock = run_mock(["--break-outbound", "1"], mock_port) + runner, timeout = run_runner(mock_port, listen, ["operation_mode"], 12, + env_extra={"RUNNER_KEEPALIVE_MS": "1000", + "RUNNER_SET_NAME": "fan_speed", + "RUNNER_SET_VALUE": "2", + "RUNNER_SET_AT": "6"}) + check(runner.wait_line("STATE online", 10) is not None, "online") + check(mock.wait_line("~CMD skipped", 10) is not None, + "мок пропустил ответ (десинк)") + check(mock.wait_line("~CMD bad-sign", 15) is not None, + "подпись следующей команды не сошлась") + # Реальный модуль ре-кает на следующем local_reg (как в probe1-3); + # после re-key цепочки свежие — SET доходит. + check(mock.wait_line("~CMD SET fan_speed=2", 20) is not None, + "SET доставлена после re-key") + check(len([l for l in mock.all_lines() if l.startswith("KE ")]) >= 2, + "второй key exchange был") + runner.stop() + mock.stop() + + +def scenario_set_get(): + print("== scenario: SET + подтверждение GET (эха нет)") + mock_port, listen = free_port(), free_port() + mock = run_mock([], mock_port) + env = {"RUNNER_SET_NAME": "fan_speed", "RUNNER_SET_VALUE": "3", + "RUNNER_SET_AT": "1"} + runner, timeout = run_runner(mock_port, listen, [], 5, env_extra=env) + line = runner.wait_line("SET_DONE fan_speed=3", 10) + check(line is not None, "SET отправлена") + check(mock.wait_line("CMD SET fan_speed=3", 10) is not None, + "mock применил SET") + runner.wait_line("DELETED", timeout) + runner.stop() + mock.stop() + + +def main(): + scenarios = [scenario_normal, scenario_rekey, scenario_fail_push, + scenario_slot_503, scenario_no_poll, scenario_bad_ke, + scenario_chain_divergence_heals, + scenario_outbound_desync_rekey, scenario_set_get] + for sc in scenarios: + sc() + print() + if FAILURES: + print("ПРОВАЛЕНО:", len(FAILURES)) + for f in FAILURES: + print(" -", f) + sys.exit(1) + print("ВСЕ СЦЕНАРИИ ПРОЙДЕНЫ") + + +if __name__ == "__main__": + main()