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