Files
IoT/loadtests/publisher.py
T

148 lines
5.7 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# 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