This commit is contained in:
@@ -0,0 +1,78 @@
|
||||
import socket
|
||||
import struct
|
||||
|
||||
|
||||
class SimpleMQTT:
|
||||
def __init__(self, host, port=1883):
|
||||
self.host = host
|
||||
self.port = port
|
||||
self.sock = None
|
||||
|
||||
def connect(self, client_id, username=None, password=None):
|
||||
addr_info = socket.getaddrinfo(self.host, self.port)
|
||||
addr = addr_info[0][-1]
|
||||
self.sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
|
||||
self.sock.connect(addr)
|
||||
|
||||
# Flags: Clean session (bit 1) + optional credentials
|
||||
flags = 0x02
|
||||
if username:
|
||||
flags |= 0x80
|
||||
if password:
|
||||
flags |= 0x40
|
||||
|
||||
# Variable header: Protocol Name (MQTT) + Level (4) + Flags + KeepAlive (60s)
|
||||
var_header = bytearray(
|
||||
[0x00, 0x04, ord("M"), ord("Q"), ord("T"), ord("T"), 0x04, flags, 0x00, 0x3C]
|
||||
)
|
||||
|
||||
payload = bytearray()
|
||||
payload.extend(self._encode_str(client_id))
|
||||
if username:
|
||||
payload.extend(self._encode_str(username))
|
||||
if password:
|
||||
payload.extend(self._encode_str(password))
|
||||
|
||||
body = var_header + payload
|
||||
packet = bytearray([0x10]) + self._encode_len(len(body)) + body
|
||||
self.sock.write(packet)
|
||||
|
||||
# Read CONNACK (4 bytes: 0x20, 0x02, ack_flags, return_code)
|
||||
resp = self.sock.read(4)
|
||||
if not resp or resp[0] != 0x20 or resp[3] != 0x00:
|
||||
code = resp[3] if resp and len(resp) >= 4 else "nil"
|
||||
raise RuntimeError(f"MQTT connection failed: code {code}")
|
||||
|
||||
def publish(self, topic, payload, retain=False):
|
||||
cmd = 0x30 | (0x01 if retain else 0x00)
|
||||
if isinstance(payload, str):
|
||||
payload = payload.encode("utf-8")
|
||||
body = self._encode_str(topic) + payload
|
||||
packet = bytearray([cmd]) + self._encode_len(len(body)) + body
|
||||
self.sock.write(packet)
|
||||
|
||||
def disconnect(self):
|
||||
try:
|
||||
if self.sock:
|
||||
self.sock.write(bytearray([0xE0, 0x00]))
|
||||
self.sock.close()
|
||||
except Exception:
|
||||
pass
|
||||
finally:
|
||||
self.sock = None
|
||||
|
||||
def _encode_str(self, s):
|
||||
raw = s.encode("utf-8")
|
||||
return struct.pack("!H", len(raw)) + raw
|
||||
|
||||
def _encode_len(self, length):
|
||||
encoded = bytearray()
|
||||
while True:
|
||||
digit = length % 128
|
||||
length //= 128
|
||||
if length > 0:
|
||||
digit |= 0x80
|
||||
encoded.append(digit)
|
||||
if length == 0:
|
||||
break
|
||||
return encoded
|
||||
Reference in New Issue
Block a user