core(M2): машина состояний сессии Ayla LAN + mock-модуль + интеграционные сценарии

- session.{hpp,cpp}: state machine (idle/registering/online/recovering/
  offline/key_error); httpd-обработчики key_exchange (200/426/412, re-key
  прозрачно), commands (одна команда, 206/200, envelope, глобальный seq_no),
  datapoint (unpack -> PropertyEvent / 401+тишина 50с для re-key-восстановления);
  сессионный поток: local_reg POST?dsn/PUT (local_ip_for), keep-alive, backoff
  x1.6->60с, 503->offline/NoSlot, activation-timeout->recovering, delete_session
  с ожиданием выдачи; очередь с coalescing + batch; телеметрия; колбэки из
  двух потоков с задокументированным контрактом; буферы datapoint-пути в Impl.
- platform: local_ip_for (UDP-connect) posix+esp-idf; стек httpd 24576
  (переполнение 16КБ поймано gdb на Release).
- mock_ac.py: мок-модуль, stdlib-only чистый python AES-256 (свёрстан с
  pycryptodome); сценарии: 503, no-poll, rekey-every, stale-gap (эмуляция
  'вернувшегося' приложения), fail-pushes (битая подпись), garbage-pushes
  (обрыв блока), break-outbound (исходящий десинк -> модуль ре-кает на
  local_reg, как probe1-3), push-every, fail-first-ke.
- session_runner + test_session_mock.py: 9 сценариев через ctest, включая
  самосинхронизацию CBC и восстановление после исходящего десинка.
- Прибор AP-WC1E: активация <=1с; re-key семантика ИСПРАВЛЕНА по живым
  тестам: re-key при зазоре local_reg >= ~44-50с (не по возрасту сессии!);
  при честном keep-alive 15с сессия стабильна без re-key; PROTOCOL/LEGACY/
  PLAN обновлены; восстановление = тишина >порога + возврат.
- CI: 7/7 x3 (gcc-Rel, gcc-ASan/UBSan, clang); ESP-IDF esp32 build complete.
Ревью под-агентом: 2 круга (стек httpd, залипание состояний, dangling cfg,
физика десинка) — APPROVED.
This commit is contained in:
2026-09-27 10:53:28 +03:00
parent 345fe19ca7
commit e74f3dc67a
15 changed files with 2009 additions and 26 deletions

420
tests/ayla/mock_ac.py Normal file
View File

@@ -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 <first|put> notify=0|1
KE <random1>
CMD <method> <resource|name> [value]
PUSH <name> <value>
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()

View File

@@ -0,0 +1,135 @@
// Тестовый раннер сессии: поднимает fgl::ayla::Session против mock_ac.py
// (или реального модуля) и печатает события строками в stdout:
// STATE <state> <err> — смена состояния
// PROP <name> <cmd_id> <status> <kind:value> — 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 <chrono>
#include <cstdio>
#include <cstdlib>
#include <cstring>
#include <thread>
#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<uint32_t>(atoi(v)) : def;
}
static void on_state(void*, SessionState st, SessionError err) {
printf("STATE %s %d\n", state_name(st), static_cast<int>(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<long long>(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 <host> <device_port> <listen_port> <dsn> <lanip_key> "
"<key_id> <duration_sec> [prop ...]\n",
argv[0]);
return 2;
}
fgl::ayla::SessionConfig cfg{};
cfg.host = argv[1];
cfg.device_port = static_cast<uint16_t>(atoi(argv[2]));
cfg.listen_port = static_cast<uint16_t>(atoi(argv[3]));
cfg.dsn = argv[4];
cfg.lanip_key = argv[5];
cfg.lanip_key_id = static_cast<uint32_t>(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;
}

View File

@@ -0,0 +1,295 @@
#!/usr/bin/env python3
"""Интеграционные сценарии сессии против mock-модуля (mock_ac.py).
usage: test_session_mock.py <session_runner> <mock_ac.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()