148 lines
5.7 KiB
Python
148 lines
5.7 KiB
Python
# publisher.py — MQTT-паблишер: один поток на устройство.
|
||
#
|
||
# Каждый DevicePublisher:
|
||
# - подключается по wss (subprotocol mqtt), проверяет CONNACK через on_connect;
|
||
# - публикует в <ns>/telemetry/<device_id> с заданной частотой;
|
||
# - нумерует сообщения (seq) и кладёт sent_at (epoch ms) — для верификатора;
|
||
# - умеет периодически рвать соединение (reconnect_every) — эмуляция
|
||
# edge-разрывов платформы (~150с) и reconnect-storm;
|
||
# - ведёт метрики: sent, conn_ok, conn_fail, disconnects, connack_rc{...}.
|
||
import json
|
||
import ssl
|
||
import threading
|
||
import time
|
||
|
||
import paho.mqtt.client as mqtt
|
||
|
||
from .common import Counter, WS_HOST, WS_PATH, WS_PORT, now_ms, ws_tls_opts
|
||
|
||
|
||
class DevicePublisher(threading.Thread):
|
||
def __init__(self, run_id, ns, device_id, username, password,
|
||
rate=1.0, duration=60.0, qos=0, payload_size=128,
|
||
reconnect_every=0, counter=None, extra_topic=None,
|
||
publish_foreign=False):
|
||
super().__init__(daemon=True)
|
||
self.run_id = run_id
|
||
self.ns = ns
|
||
self.device_id = device_id
|
||
self.username = username
|
||
self.password = password
|
||
self.rate = rate # сообщений/сек
|
||
self.duration = duration # секунд
|
||
self.qos = qos
|
||
self.payload_size = payload_size
|
||
self.reconnect_every = reconnect_every # 0 = не рвать
|
||
self.counter = counter or Counter()
|
||
self.extra_topic = extra_topic # дополнительный (чужой) топик
|
||
self.publish_foreign = publish_foreign # публиковать в чужой топик
|
||
self.connack_rc = None
|
||
self.disconnected = threading.Event()
|
||
self.stop = threading.Event()
|
||
self.topic = f"{ns}/telemetry/{device_id}"
|
||
self._connack = threading.Event()
|
||
|
||
# --- MQTT ---
|
||
def _on_connect(self, client, userdata, flags, rc, props=None):
|
||
self.connack_rc = rc
|
||
self.counter.inc(f"connack_rc{rc}")
|
||
if rc == 0:
|
||
self.counter.inc("conn_ok")
|
||
else:
|
||
self.counter.inc("conn_fail")
|
||
self._connack.set()
|
||
|
||
def _on_disconnect(self, client, userdata, rc, props=None):
|
||
self.counter.inc("disconnects")
|
||
self.disconnected.set()
|
||
|
||
def _client(self):
|
||
c = mqtt.Client(client_id=f"lt-{self.run_id}-{self.device_id}",
|
||
transport="websockets")
|
||
c.ws_set_options(path=WS_PATH)
|
||
c.username_pw_set(self.username, self.password)
|
||
ws_tls_opts(c)
|
||
c.on_connect = self._on_connect
|
||
c.on_disconnect = self._on_disconnect
|
||
c.reconnect_delay_set(min_delay=2, max_delay=10)
|
||
return c
|
||
|
||
def _payload(self, seq):
|
||
pad = "x" * max(0, self.payload_size)
|
||
return json.dumps({
|
||
"run_id": self.run_id,
|
||
"seq": seq,
|
||
"sent_at": now_ms(),
|
||
"device_id": self.device_id,
|
||
"pad": pad,
|
||
})
|
||
|
||
def _connect(self, client, timeout=20):
|
||
self._connack.clear()
|
||
client.connect(WS_HOST, WS_PORT, timeout)
|
||
client.loop_start()
|
||
return self._connack.wait(timeout)
|
||
|
||
# --- поток ---
|
||
def run(self):
|
||
client = self._client()
|
||
seq = 0
|
||
t0 = time.monotonic()
|
||
next_t = t0
|
||
interval = 1.0 / self.rate if self.rate > 0 else 1.0
|
||
connected = False
|
||
last_reconnect = time.monotonic()
|
||
|
||
while not self.stop.is_set() and (time.monotonic() - t0) < self.duration:
|
||
if not connected:
|
||
if not self._connect(client):
|
||
self.counter.inc("conn_fail")
|
||
time.sleep(2)
|
||
continue
|
||
connected = True
|
||
|
||
# периодический разрыв — эмуляция edge-разрывов платформы
|
||
if self.reconnect_every and \
|
||
time.monotonic() - last_reconnect >= self.reconnect_every:
|
||
self.counter.inc("forced_reconnects")
|
||
client.disconnect()
|
||
client.loop_stop()
|
||
client = self._client()
|
||
connected = False
|
||
last_reconnect = time.monotonic()
|
||
continue
|
||
|
||
now = time.monotonic()
|
||
if now < next_t:
|
||
time.sleep(min(next_t - now, 0.5))
|
||
continue
|
||
next_t += interval
|
||
if next_t < now: # отстали — не навёрстываем очередью
|
||
next_t = now + interval
|
||
|
||
payload = self._payload(seq)
|
||
r = client.publish(self.topic, payload, qos=self.qos)
|
||
if r.rc == mqtt.MQTT_ERR_SUCCESS:
|
||
self.counter.inc("sent")
|
||
seq += 1
|
||
else:
|
||
self.counter.inc("publish_errors")
|
||
# потеряли соединение — переподключимся
|
||
connected = False
|
||
|
||
if self.publish_foreign and self.extra_topic:
|
||
# ACL-тест: публикация в чужой топик должна разорвать сессию
|
||
self.counter.inc("foreign_publish")
|
||
client.publish(self.extra_topic, payload, qos=0)
|
||
if self.disconnected.wait(3):
|
||
self.counter.inc("foreign_denied")
|
||
else:
|
||
self.counter.inc("foreign_accepted")
|
||
self.disconnected.clear()
|
||
|
||
try:
|
||
client.disconnect()
|
||
client.loop_stop()
|
||
except Exception:
|
||
pass
|