From 7cb2f56737f27952d724d31114469b300e1193ad Mon Sep 17 00:00:00 2001 From: Tavis Date: Tue, 2 Jun 2026 11:37:32 -0700 Subject: [PATCH 1/4] add serial2mqtt script --- scripts/aircube_serial2mqtt.py | 256 +++++++++++++++++++++++++++++++++ 1 file changed, 256 insertions(+) create mode 100644 scripts/aircube_serial2mqtt.py diff --git a/scripts/aircube_serial2mqtt.py b/scripts/aircube_serial2mqtt.py new file mode 100644 index 0000000..8cfacad --- /dev/null +++ b/scripts/aircube_serial2mqtt.py @@ -0,0 +1,256 @@ +#!/usr/bin/env python3 +""" +AirCube → MQTT Bridge +Reads JSON sensor data from AirCube over USB serial and publishes to MQTT +with Home Assistant auto-discovery support. + +Usage: + pip install pyserial paho-mqtt + python aircube_mqtt_bridge.py + +Config via environment variables or edit the DEFAULTS below. +""" + +import argparse +import collections +import json +import logging +import os +import sys +import time +import serial +import paho.mqtt.client as mqtt + +# --------------------------------------------------------------------------- +# Configuration — override with environment variables +# --------------------------------------------------------------------------- +SERIAL_PORT = os.getenv("AIRCUBE_PORT", "/dev/cu.usbmodem101") +SERIAL_BAUD = int(os.getenv("AIRCUBE_BAUD", "115200")) +MQTT_HOST = os.getenv("MQTT_HOST", "localhost") +MQTT_PORT = int(os.getenv("MQTT_PORT", "1883")) +MQTT_USER = os.getenv("MQTT_USER", "") +MQTT_PASS = os.getenv("MQTT_PASS", "") +DEVICE_NAME = os.getenv("AIRCUBE_NAME", "AirCube") +DEVICE_ID = os.getenv("AIRCUBE_ID", "aircube_1") # unique per device +DISCOVERY_PREFIX = os.getenv("HA_DISCOVERY_PREFIX", "homeassistant") + +STATE_TOPIC = f"aircube/{DEVICE_ID}/state" +AVAIL_TOPIC = f"aircube/{DEVICE_ID}/availability" + +# --------------------------------------------------------------------------- +# Sensor definitions: (ha_key, unit, device_class, state_class, icon) +# ha_key must match the key we publish into the state JSON payload +# --------------------------------------------------------------------------- +SENSORS = [ + ("temperature_c", "°C", "temperature", "measurement", None), + ("humidity", "%", "humidity", "measurement", None), + ("eco2", "ppm", "carbon_dioxide", "measurement", None), + ("etvoc", "ppb", None, "measurement", "mdi:chemical-weapon"), + ("aqi", None, "aqi", "measurement", None), +] + +logging.basicConfig( + level=logging.INFO, + format="%(asctime)s [%(levelname)s] %(message)s", + datefmt="%Y-%m-%d %H:%M:%S", +) +log = logging.getLogger(__name__) + + +# --------------------------------------------------------------------------- +# MQTT helpers +# --------------------------------------------------------------------------- +def publish_discovery(client: mqtt.Client) -> None: + """Publish HA MQTT discovery config for each sensor (retained).""" + device_info = { + "identifiers": [DEVICE_ID], + "name": DEVICE_NAME, + "model": "AirCube", + "manufacturer": "StuckAtPrototype", + } + + for ha_key, unit, device_class, state_class, icon in SENSORS: + friendly = ha_key.replace("_", " ").replace("c", "").title().strip() + unique_id = f"{DEVICE_ID}_{ha_key}" + config_topic = f"{DISCOVERY_PREFIX}/sensor/{DEVICE_ID}/{ha_key}/config" + + payload: dict = { + "name": friendly, + "unique_id": unique_id, + "state_topic": STATE_TOPIC, + "value_template": f"{{{{ value_json.{ha_key} }}}}", + "availability_topic": AVAIL_TOPIC, + "device": device_info, + "state_class": state_class, + } + if unit: + payload["unit_of_measurement"] = unit + if device_class: + payload["device_class"] = device_class + if icon: + payload["icon"] = icon + + client.publish(config_topic, json.dumps(payload), retain=True) + log.info("Discovery published: %s", config_topic) + + +def on_connect(client: mqtt.Client, userdata, flags, rc) -> None: + if rc == 0: + log.info("MQTT connected to %s:%s", MQTT_HOST, MQTT_PORT) + publish_discovery(client) + client.publish(AVAIL_TOPIC, "online", retain=True) + else: + log.error("MQTT connection failed, rc=%s", rc) + + +def on_disconnect(client: mqtt.Client, userdata, rc) -> None: + log.warning("MQTT disconnected (rc=%s), will retry", rc) + + +def build_mqtt_client() -> mqtt.Client: + client = mqtt.Client(client_id=DEVICE_ID, clean_session=True) + client.will_set(AVAIL_TOPIC, "offline", retain=True) + client.on_connect = on_connect + client.on_disconnect = on_disconnect + if MQTT_USER: + client.username_pw_set(MQTT_USER, MQTT_PASS) + return client + + +# --------------------------------------------------------------------------- +# Serial helpers +# --------------------------------------------------------------------------- +def open_serial(port: str, baud: int, retries: int = 10) -> serial.Serial: + for attempt in range(1, retries + 1): + try: + ser = serial.Serial(port, baud, timeout=5) + log.info("Serial opened: %s @ %d baud", port, baud) + return ser + except serial.SerialException as exc: + log.warning("Serial open failed (attempt %d/%d): %s", attempt, retries, exc) + time.sleep(3) + log.error("Could not open serial port %s after %d attempts", port, retries) + sys.exit(1) + + +def parse_aircube(raw: str) -> dict | None: + """ + Parse one line of AirCube JSON output into a flat dict for MQTT. + Returns None if the line isn't valid AirCube data. + + Expected input: + {"ens210": {"temperature_c": 23.45, "humidity": 52.30}, + "ens16x": {"etvoc": 42, "eco2": 415, "aqi": 3}, + "timestamp": 12345} + """ + try: + data = json.loads(raw.strip()) + except json.JSONDecodeError: + return None + + # Must have both sensor blocks + if "ens210" not in data or "ens16x" not in data: + return None + + return { + "temperature_c": data["ens210"].get("temperature_c"), + "humidity": data["ens210"].get("humidity"), + "eco2": data["ens16x"].get("eco2"), + "etvoc": data["ens16x"].get("etvoc"), + "aqi": data["ens16x"].get("aqi"), + } + + +# --------------------------------------------------------------------------- +# Main loop +# --------------------------------------------------------------------------- +def main() -> None: + parser = argparse.ArgumentParser(description="AirCube to MQTT Bridge") + parser.add_argument( + "--no-mqtt", action="store_true", help="Don't connect to MQTT, just print data" + ) + args = parser.parse_args() + + log.info("Starting AirCube MQTT bridge (device_id=%s)", DEVICE_ID) + log.info("Serial: %s", SERIAL_PORT) + + client = None + if not args.no_mqtt: + log.info("MQTT: %s:%s State topic: %s", MQTT_HOST, MQTT_PORT, STATE_TOPIC) + client = build_mqtt_client() + # Connect MQTT (non-blocking loop in background thread) + try: + client.connect(MQTT_HOST, MQTT_PORT, keepalive=60) + except Exception as exc: + log.error("Could not connect to MQTT broker: %s", exc) + sys.exit(1) + client.loop_start() + else: + log.info("MQTT disabled: running in print-only mode") + + ser = open_serial(SERIAL_PORT, SERIAL_BAUD) + + consecutive_errors = 0 + buffer = collections.defaultdict(list) + last_publish_time = time.time() + + while True: + try: + raw = ser.readline().decode("utf-8", errors="replace") + if not raw.strip(): + continue + + payload = parse_aircube(raw) + if payload is None: + log.debug("Skipping non-sensor line: %s", raw.strip()) + continue + + # Add values to buffer + for key, val in payload.items(): + if val is not None: + buffer[key].append(val) + + current_time = time.time() + if current_time - last_publish_time >= 10: + if buffer: + # Calculate averages + avg_payload = {} + for key, values in buffer.items(): + avg_payload[key] = round(sum(values) / len(values), 2) + + if client: + client.publish(STATE_TOPIC, json.dumps(avg_payload)) + log.info("Published averaged data: %s", avg_payload) + else: + print( + f"WOULD PUBLISH averaged data to {STATE_TOPIC}: {json.dumps(avg_payload)}" + ) + + buffer.clear() + + last_publish_time = current_time + + consecutive_errors = 0 + + except serial.SerialException as exc: + consecutive_errors += 1 + log.error("Serial error: %s (attempt %d)", exc, consecutive_errors) + if client: + client.publish(AVAIL_TOPIC, "offline", retain=True) + ser.close() + time.sleep(5) + ser = open_serial(SERIAL_PORT, SERIAL_BAUD) + if client: + client.publish(AVAIL_TOPIC, "online", retain=True) + + except KeyboardInterrupt: + log.info("Shutting down") + if client: + client.publish(AVAIL_TOPIC, "offline", retain=True) + client.loop_stop() + ser.close() + sys.exit(0) + + +if __name__ == "__main__": + main() From fa19ce20d1d433183b21a379967df4d3062821ef Mon Sep 17 00:00:00 2001 From: Tavis Date: Tue, 2 Jun 2026 11:51:02 -0700 Subject: [PATCH 2/4] make publish interval configurable --- scripts/aircube_serial2mqtt.py | 26 +++++++++++++++++++++++++- 1 file changed, 25 insertions(+), 1 deletion(-) diff --git a/scripts/aircube_serial2mqtt.py b/scripts/aircube_serial2mqtt.py index 8cfacad..38926b3 100644 --- a/scripts/aircube_serial2mqtt.py +++ b/scripts/aircube_serial2mqtt.py @@ -21,6 +21,13 @@ import serial import paho.mqtt.client as mqtt +try: + from dotenv import load_dotenv + + load_dotenv() +except ImportError: + pass + # --------------------------------------------------------------------------- # Configuration — override with environment variables # --------------------------------------------------------------------------- @@ -33,6 +40,7 @@ DEVICE_NAME = os.getenv("AIRCUBE_NAME", "AirCube") DEVICE_ID = os.getenv("AIRCUBE_ID", "aircube_1") # unique per device DISCOVERY_PREFIX = os.getenv("HA_DISCOVERY_PREFIX", "homeassistant") +PUBLISH_INTERVAL = int(os.getenv("AIRCUBE_INTERVAL", "0")) STATE_TOPIC = f"aircube/{DEVICE_ID}/state" AVAIL_TOPIC = f"aircube/{DEVICE_ID}/availability" @@ -169,10 +177,17 @@ def main() -> None: parser.add_argument( "--no-mqtt", action="store_true", help="Don't connect to MQTT, just print data" ) + parser.add_argument( + "--interval", + type=int, + default=PUBLISH_INTERVAL, + help=f"Publish interval in seconds. Set to 0 to publish every packet immediately. (default: {PUBLISH_INTERVAL})", + ) args = parser.parse_args() log.info("Starting AirCube MQTT bridge (device_id=%s)", DEVICE_ID) log.info("Serial: %s", SERIAL_PORT) + log.info("Interval: %s seconds", args.interval) client = None if not args.no_mqtt: @@ -205,13 +220,22 @@ def main() -> None: log.debug("Skipping non-sensor line: %s", raw.strip()) continue + if args.interval <= 0: + # Publish immediately + if client: + client.publish(STATE_TOPIC, json.dumps(payload)) + log.info("Published data: %s", payload) + else: + print(f"WOULD PUBLISH data to {STATE_TOPIC}: {json.dumps(payload)}") + continue + # Add values to buffer for key, val in payload.items(): if val is not None: buffer[key].append(val) current_time = time.time() - if current_time - last_publish_time >= 10: + if current_time - last_publish_time >= args.interval: if buffer: # Calculate averages avg_payload = {} From 73da2fc86e4d56a00ecfc1713dbee7791d039590 Mon Sep 17 00:00:00 2001 From: Tavis Date: Mon, 27 Jul 2026 15:27:49 -0700 Subject: [PATCH 3/4] add brightness control passthrough for HA --- scripts/.env.serial2mqtt | 14 ++ scripts/aircube_serial2mqtt.py | 332 ++++++++++++++++++++++++++++----- 2 files changed, 303 insertions(+), 43 deletions(-) create mode 100644 scripts/.env.serial2mqtt diff --git a/scripts/.env.serial2mqtt b/scripts/.env.serial2mqtt new file mode 100644 index 0000000..29d566b --- /dev/null +++ b/scripts/.env.serial2mqtt @@ -0,0 +1,14 @@ +# AirCube MQTT bridge — copy to the Pi and run from scripts/ (or set path in load_dotenv) + +AIRCUBE_PORT=/dev/cu.usbmodem21101 + +MQTT_HOST=10.10.1.10 +MQTT_PORT=1883 +# MQTT_USER= +# MQTT_PASS= + +AIRCUBE_NAME=AirCube +AIRCUBE_ID=aircube_1 + +HA_DISCOVERY_PREFIX=homeassistant +AIRCUBE_INTERVAL=5 diff --git a/scripts/aircube_serial2mqtt.py b/scripts/aircube_serial2mqtt.py index 38926b3..1ce2ad7 100644 --- a/scripts/aircube_serial2mqtt.py +++ b/scripts/aircube_serial2mqtt.py @@ -1,35 +1,88 @@ #!/usr/bin/env python3 """ -AirCube → MQTT Bridge +AirCube -> MQTT Bridge Reads JSON sensor data from AirCube over USB serial and publishes to MQTT -with Home Assistant auto-discovery support. +with Home Assistant auto-discovery support. Accepts brightness commands on +MQTT and forwards them to the device as set_intensity serial commands. + +Requires Python 3.2+. Usage: - pip install pyserial paho-mqtt - python aircube_mqtt_bridge.py + pip install pyserial paho-mqtt python-dotenv + python aircube_serial2mqtt.py Config via environment variables or edit the DEFAULTS below. """ +from __future__ import print_function + import argparse import collections import json import logging import os +import queue import sys +import threading import time + import serial import paho.mqtt.client as mqtt +try: + JSON_DECODE_ERROR = json.JSONDecodeError +except AttributeError: + JSON_DECODE_ERROR = ValueError + +SCRIPT_DIR = os.path.dirname(os.path.abspath(__file__)) +ENV_FILE = os.path.join(SCRIPT_DIR, ".env") +DOTENV_AVAILABLE = False +ENV_LOADED = False + try: from dotenv import load_dotenv - load_dotenv() + DOTENV_AVAILABLE = True + if os.path.isfile(ENV_FILE): + load_dotenv(ENV_FILE) + ENV_LOADED = True except ImportError: pass + +def log_env_status(logger): + """Warn at startup if a .env file exists but could not be loaded.""" + if ENV_LOADED: + logger.info("Loaded config from %s", ENV_FILE) + return + + if not os.path.isfile(ENV_FILE): + return + + if DOTENV_AVAILABLE: + logger.warning( + "Found %s but dotenv did not load it; using script defaults", ENV_FILE + ) + else: + logger.warning( + "Found %s but python-dotenv is not installed - .env ignored, " + "using script defaults. Install with: pip install python-dotenv or edit script directly" , + ENV_FILE, + ) + + +def serial_port_open(ser): + """Return True when a pyserial handle is open (2.x and 3.x compatible).""" + if ser is None: + return False + is_open = getattr(ser, "is_open", None) + if is_open is not None: + return is_open + return ser.isOpen() + + # --------------------------------------------------------------------------- -# Configuration — override with environment variables +# Configuration - override with environment variables # --------------------------------------------------------------------------- SERIAL_PORT = os.getenv("AIRCUBE_PORT", "/dev/cu.usbmodem101") SERIAL_BAUD = int(os.getenv("AIRCUBE_BAUD", "115200")) @@ -42,15 +95,17 @@ DISCOVERY_PREFIX = os.getenv("HA_DISCOVERY_PREFIX", "homeassistant") PUBLISH_INTERVAL = int(os.getenv("AIRCUBE_INTERVAL", "0")) -STATE_TOPIC = f"aircube/{DEVICE_ID}/state" -AVAIL_TOPIC = f"aircube/{DEVICE_ID}/availability" +STATE_TOPIC = "aircube/{0}/state".format(DEVICE_ID) +AVAIL_TOPIC = "aircube/{0}/availability".format(DEVICE_ID) +BRIGHTNESS_CMD_TOPIC = "aircube/{0}/brightness/set".format(DEVICE_ID) +BRIGHTNESS_STATE_TOPIC = "aircube/{0}/brightness/state".format(DEVICE_ID) # --------------------------------------------------------------------------- # Sensor definitions: (ha_key, unit, device_class, state_class, icon) # ha_key must match the key we publish into the state JSON payload # --------------------------------------------------------------------------- SENSORS = [ - ("temperature_c", "°C", "temperature", "measurement", None), + ("temperature_c", "C", "temperature", "measurement", None), ("humidity", "%", "humidity", "measurement", None), ("eco2", "ppm", "carbon_dioxide", "measurement", None), ("etvoc", "ppb", None, "measurement", "mdi:chemical-weapon"), @@ -65,10 +120,22 @@ log = logging.getLogger(__name__) +class BridgeState(object): + """Shared state between the serial reader, main loop, and MQTT callbacks.""" + + def __init__(self): + self.ser = None + self.write_lock = threading.Lock() + self.intensity_ack = threading.Event() + self.intensity_value = None + self.sensor_queue = queue.Queue() + self.stop = threading.Event() + + # --------------------------------------------------------------------------- # MQTT helpers # --------------------------------------------------------------------------- -def publish_discovery(client: mqtt.Client) -> None: +def publish_discovery(client): """Publish HA MQTT discovery config for each sensor (retained).""" device_info = { "identifiers": [DEVICE_ID], @@ -79,14 +146,16 @@ def publish_discovery(client: mqtt.Client) -> None: for ha_key, unit, device_class, state_class, icon in SENSORS: friendly = ha_key.replace("_", " ").replace("c", "").title().strip() - unique_id = f"{DEVICE_ID}_{ha_key}" - config_topic = f"{DISCOVERY_PREFIX}/sensor/{DEVICE_ID}/{ha_key}/config" + unique_id = "{0}_{1}".format(DEVICE_ID, ha_key) + config_topic = "{0}/sensor/{1}/{2}/config".format( + DISCOVERY_PREFIX, DEVICE_ID, ha_key + ) - payload: dict = { + payload = { "name": friendly, "unique_id": unique_id, "state_topic": STATE_TOPIC, - "value_template": f"{{{{ value_json.{ha_key} }}}}", + "value_template": "{{{{ value_json.{0} }}}}".format(ha_key), "availability_topic": AVAIL_TOPIC, "device": device_info, "state_class": state_class, @@ -101,25 +170,46 @@ def publish_discovery(client: mqtt.Client) -> None: client.publish(config_topic, json.dumps(payload), retain=True) log.info("Discovery published: %s", config_topic) + brightness_config_topic = "{0}/number/{1}/brightness/config".format( + DISCOVERY_PREFIX, DEVICE_ID + ) + brightness_payload = { + "name": "Brightness", + "unique_id": "{0}_brightness".format(DEVICE_ID), + "command_topic": BRIGHTNESS_CMD_TOPIC, + "state_topic": BRIGHTNESS_STATE_TOPIC, + "availability_topic": AVAIL_TOPIC, + "device": device_info, + "min": 0, + "max": 1, + "step": 0.01, + "mode": "slider", + } + client.publish(brightness_config_topic, json.dumps(brightness_payload), retain=True) + log.info("Discovery published: %s", brightness_config_topic) + -def on_connect(client: mqtt.Client, userdata, flags, rc) -> None: +def on_connect(client, userdata, flags, rc): if rc == 0: log.info("MQTT connected to %s:%s", MQTT_HOST, MQTT_PORT) publish_discovery(client) + client.subscribe(BRIGHTNESS_CMD_TOPIC) + log.info("Subscribed to brightness command topic: %s", BRIGHTNESS_CMD_TOPIC) client.publish(AVAIL_TOPIC, "online", retain=True) else: log.error("MQTT connection failed, rc=%s", rc) -def on_disconnect(client: mqtt.Client, userdata, rc) -> None: +def on_disconnect(client, userdata, rc): log.warning("MQTT disconnected (rc=%s), will retry", rc) -def build_mqtt_client() -> mqtt.Client: - client = mqtt.Client(client_id=DEVICE_ID, clean_session=True) +def build_mqtt_client(bridge): + client = mqtt.Client(client_id=DEVICE_ID, clean_session=True, userdata=bridge) client.will_set(AVAIL_TOPIC, "offline", retain=True) client.on_connect = on_connect client.on_disconnect = on_disconnect + client.on_message = on_message if MQTT_USER: client.username_pw_set(MQTT_USER, MQTT_PASS) return client @@ -128,35 +218,176 @@ def build_mqtt_client() -> mqtt.Client: # --------------------------------------------------------------------------- # Serial helpers # --------------------------------------------------------------------------- -def open_serial(port: str, baud: int, retries: int = 10) -> serial.Serial: +def open_serial(port, baud, retries=10): for attempt in range(1, retries + 1): try: ser = serial.Serial(port, baud, timeout=5) log.info("Serial opened: %s @ %d baud", port, baud) return ser except serial.SerialException as exc: - log.warning("Serial open failed (attempt %d/%d): %s", attempt, retries, exc) + log.warning( + "Serial open failed (attempt %d/%d): %s", attempt, retries, exc + ) time.sleep(3) log.error("Could not open serial port %s after %d attempts", port, retries) sys.exit(1) -def parse_aircube(raw: str) -> dict | None: +def parse_brightness_payload(raw): + """ + Parse an MQTT brightness command payload. + + Accepts: + - plain number: "0.3" (0.0-1.0) or "30" (0-100 percent) + - JSON object: {"value": 0.3} or {"brightness": 30} + """ + text = raw.strip() + if not text: + return None + + try: + data = json.loads(text) + except JSON_DECODE_ERROR: + data = text + + if isinstance(data, dict): + value = data.get("value", data.get("brightness")) + if value is None: + return None + else: + value = data + + try: + intensity = float(value) + except (TypeError, ValueError): + return None + + if intensity > 1.0: + intensity /= 100.0 + + return max(0.0, min(1.0, intensity)) + + +def parse_set_intensity_response(raw): + """Return applied intensity if raw is a set_intensity ack, else None.""" + try: + parsed = json.loads(raw.strip()) + except JSON_DECODE_ERROR: + return None + + if parsed.get("status") == "ok" and parsed.get("cmd") == "set_intensity": + return float(parsed.get("value", 0)) + return None + + +def dispatch_serial_line(bridge, line): + """Route one serial line to brightness ack handling or the sensor queue.""" + stripped = line.strip() + if not stripped: + return + + applied = parse_set_intensity_response(stripped) + if applied is not None: + bridge.intensity_value = applied + bridge.intensity_ack.set() + return + + payload = parse_aircube(stripped) + if payload is not None: + bridge.sensor_queue.put(payload) + return + + log.debug("Skipping non-sensor line: %s", stripped) + + +def serial_reader_loop(bridge): + """Single reader for all serial input so command acks are never dropped.""" + log.info("Serial reader thread started") + while not bridge.stop.is_set(): + ser = bridge.ser + if not serial_port_open(ser): + time.sleep(0.1) + continue + + try: + raw = ser.readline().decode("utf-8", "replace") + if raw: + dispatch_serial_line(bridge, raw) + except serial.SerialException as exc: + log.error("Serial reader error: %s", exc) + bridge.sensor_queue.put(None) + + +def send_set_intensity(bridge, intensity): + """Send set_intensity to the device and wait for ack from the reader thread.""" + if bridge is None: + log.error("Internal bridge state missing; cannot set brightness") + return False + + if not serial_port_open(bridge.ser): + log.error("Serial port not open; cannot set brightness") + return False + + # Firmware parser expects compact JSON (no spaces). + command = json.dumps( + {"cmd": "set_intensity", "value": intensity}, + separators=(",", ":"), + ) + "\n" + bridge.intensity_ack.clear() + bridge.intensity_value = None + + with bridge.write_lock: + try: + bridge.ser.write(command.encode("utf-8")) + bridge.ser.flush() + log.info("Sent set_intensity command: %.2f", intensity) + except serial.SerialException as exc: + log.error("Serial write failed while setting brightness: %s", exc) + return False + + if bridge.intensity_ack.wait(5.0): + log.info("Brightness set to %.2f via serial", bridge.intensity_value) + return True + + log.warning("Timed out waiting for set_intensity response (command was sent)") + return False + + +def on_message(client, userdata, msg): + if msg.topic != BRIGHTNESS_CMD_TOPIC: + return + + intensity = parse_brightness_payload(msg.payload.decode("utf-8", "replace")) + if intensity is None: + log.warning("Invalid brightness payload on %s: %r", msg.topic, msg.payload) + return + + bridge = userdata + if not isinstance(bridge, BridgeState): + log.error( + "MQTT userdata is not bridge state (%r); cannot set brightness", userdata + ) + return + + if send_set_intensity(bridge, intensity): + applied = ( + bridge.intensity_value + if bridge.intensity_value is not None + else intensity + ) + client.publish(BRIGHTNESS_STATE_TOPIC, "{0:.2f}".format(applied), retain=True) + + +def parse_aircube(raw): """ Parse one line of AirCube JSON output into a flat dict for MQTT. Returns None if the line isn't valid AirCube data. - - Expected input: - {"ens210": {"temperature_c": 23.45, "humidity": 52.30}, - "ens16x": {"etvoc": 42, "eco2": 415, "aqi": 3}, - "timestamp": 12345} """ try: data = json.loads(raw.strip()) - except json.JSONDecodeError: + except JSON_DECODE_ERROR: return None - # Must have both sensor blocks if "ens210" not in data or "ens16x" not in data: return None @@ -172,7 +403,7 @@ def parse_aircube(raw: str) -> dict | None: # --------------------------------------------------------------------------- # Main loop # --------------------------------------------------------------------------- -def main() -> None: +def main(): parser = argparse.ArgumentParser(description="AirCube to MQTT Bridge") parser.add_argument( "--no-mqtt", action="store_true", help="Don't connect to MQTT, just print data" @@ -181,19 +412,24 @@ def main() -> None: "--interval", type=int, default=PUBLISH_INTERVAL, - help=f"Publish interval in seconds. Set to 0 to publish every packet immediately. (default: {PUBLISH_INTERVAL})", + help=( + "Publish interval in seconds. Set to 0 to publish every packet " + "immediately. (default: {0})".format(PUBLISH_INTERVAL) + ), ) args = parser.parse_args() log.info("Starting AirCube MQTT bridge (device_id=%s)", DEVICE_ID) + log_env_status(log) log.info("Serial: %s", SERIAL_PORT) log.info("Interval: %s seconds", args.interval) + bridge = BridgeState() client = None if not args.no_mqtt: log.info("MQTT: %s:%s State topic: %s", MQTT_HOST, MQTT_PORT, STATE_TOPIC) - client = build_mqtt_client() - # Connect MQTT (non-blocking loop in background thread) + log.info("Brightness command topic: %s", BRIGHTNESS_CMD_TOPIC) + client = build_mqtt_client(bridge) try: client.connect(MQTT_HOST, MQTT_PORT, keepalive=60) except Exception as exc: @@ -204,6 +440,11 @@ def main() -> None: log.info("MQTT disabled: running in print-only mode") ser = open_serial(SERIAL_PORT, SERIAL_BAUD) + bridge.ser = ser + + reader = threading.Thread(target=serial_reader_loop, args=(bridge,)) + reader.daemon = True + reader.start() consecutive_errors = 0 buffer = collections.defaultdict(list) @@ -211,25 +452,26 @@ def main() -> None: while True: try: - raw = ser.readline().decode("utf-8", errors="replace") - if not raw.strip(): + try: + payload = bridge.sensor_queue.get(True, 1) + except queue.Empty: continue - payload = parse_aircube(raw) if payload is None: - log.debug("Skipping non-sensor line: %s", raw.strip()) - continue + raise serial.SerialException("serial reader reported disconnect") if args.interval <= 0: - # Publish immediately if client: client.publish(STATE_TOPIC, json.dumps(payload)) log.info("Published data: %s", payload) else: - print(f"WOULD PUBLISH data to {STATE_TOPIC}: {json.dumps(payload)}") + print( + "WOULD PUBLISH data to {0}: {1}".format( + STATE_TOPIC, json.dumps(payload) + ) + ) continue - # Add values to buffer for key, val in payload.items(): if val is not None: buffer[key].append(val) @@ -237,17 +479,18 @@ def main() -> None: current_time = time.time() if current_time - last_publish_time >= args.interval: if buffer: - # Calculate averages avg_payload = {} for key, values in buffer.items(): - avg_payload[key] = round(sum(values) / len(values), 2) + avg_payload[key] = round(sum(values) / float(len(values)), 2) if client: client.publish(STATE_TOPIC, json.dumps(avg_payload)) log.info("Published averaged data: %s", avg_payload) else: print( - f"WOULD PUBLISH averaged data to {STATE_TOPIC}: {json.dumps(avg_payload)}" + "WOULD PUBLISH averaged data to {0}: {1}".format( + STATE_TOPIC, json.dumps(avg_payload) + ) ) buffer.clear() @@ -262,13 +505,16 @@ def main() -> None: if client: client.publish(AVAIL_TOPIC, "offline", retain=True) ser.close() + bridge.ser = None time.sleep(5) ser = open_serial(SERIAL_PORT, SERIAL_BAUD) + bridge.ser = ser if client: client.publish(AVAIL_TOPIC, "online", retain=True) except KeyboardInterrupt: log.info("Shutting down") + bridge.stop.set() if client: client.publish(AVAIL_TOPIC, "offline", retain=True) client.loop_stop() From 8eb0bd0f8e14d9109d8c04bf9bd352412ab6354a Mon Sep 17 00:00:00 2001 From: Tavis Date: Mon, 27 Jul 2026 18:20:00 -0700 Subject: [PATCH 4/4] thrashing to get serial more consistent on old pi --- scripts/.env.serial2mqtt | 5 +- scripts/aircube_serial2mqtt.py | 336 ++++++++++++++++++++++++--------- 2 files changed, 246 insertions(+), 95 deletions(-) diff --git a/scripts/.env.serial2mqtt b/scripts/.env.serial2mqtt index 29d566b..9a42244 100644 --- a/scripts/.env.serial2mqtt +++ b/scripts/.env.serial2mqtt @@ -1,5 +1,8 @@ -# AirCube MQTT bridge — copy to the Pi and run from scripts/ (or set path in load_dotenv) +# AirCube MQTT bridge — copy to the Pi as ~/.env or scripts/.env +# Pi / Linux: +#AIRCUBE_PORT=/dev/ttyACM0 +# macOS (uncomment and comment out ttyACM0 above): AIRCUBE_PORT=/dev/cu.usbmodem21101 MQTT_HOST=10.10.1.10 diff --git a/scripts/aircube_serial2mqtt.py b/scripts/aircube_serial2mqtt.py index 1ce2ad7..7997294 100644 --- a/scripts/aircube_serial2mqtt.py +++ b/scripts/aircube_serial2mqtt.py @@ -26,7 +26,19 @@ import threading import time -import serial +try: + import serial +except ImportError: + serial = None + +if serial is None or not hasattr(serial, "Serial"): + sys.stderr.write( + "ERROR: pyserial is not installed (or the wrong 'serial' package is).\n" + " pip uninstall serial\n" + " pip install pyserial\n" + ) + sys.exit(1) + import paho.mqtt.client as mqtt try: @@ -35,19 +47,55 @@ JSON_DECODE_ERROR = ValueError SCRIPT_DIR = os.path.dirname(os.path.abspath(__file__)) -ENV_FILE = os.path.join(SCRIPT_DIR, ".env") +ENV_CANDIDATES = [ + os.path.join(SCRIPT_DIR, ".env"), + os.path.join(SCRIPT_DIR, ".env.serial2mqtt"), + os.path.expanduser("~/.env"), +] +ENV_FILE = "" DOTENV_AVAILABLE = False ENV_LOADED = False -try: - from dotenv import load_dotenv - DOTENV_AVAILABLE = True - if os.path.isfile(ENV_FILE): - load_dotenv(ENV_FILE) - ENV_LOADED = True -except ImportError: - pass +def load_env_file(path): + """Minimal .env loader (no python-dotenv required).""" + with open(path, "r") as handle: + for line in handle: + line = line.strip() + if not line or line.startswith("#") or "=" not in line: + continue + key, _, value = line.partition("=") + key = key.strip() + value = value.strip().strip('"').strip("'") + if key: + os.environ[key] = value + + +def init_env(): + global ENV_FILE, ENV_LOADED, DOTENV_AVAILABLE + + try: + from dotenv import load_dotenv + + for path in ENV_CANDIDATES: + if os.path.isfile(path): + load_dotenv(path) + ENV_FILE = path + ENV_LOADED = True + DOTENV_AVAILABLE = True + return + except ImportError: + pass + + for path in ENV_CANDIDATES: + if os.path.isfile(path): + load_env_file(path) + ENV_FILE = path + ENV_LOADED = True + return + + +init_env() def log_env_status(logger): @@ -56,19 +104,12 @@ def log_env_status(logger): logger.info("Loaded config from %s", ENV_FILE) return - if not os.path.isfile(ENV_FILE): - return - - if DOTENV_AVAILABLE: - logger.warning( - "Found %s but dotenv did not load it; using script defaults", ENV_FILE - ) - else: - logger.warning( - "Found %s but python-dotenv is not installed - .env ignored, " - "using script defaults. Install with: pip install python-dotenv or edit script directly" , - ENV_FILE, - ) + for path in ENV_CANDIDATES: + if os.path.isfile(path): + logger.warning( + "Found %s but could not load it; using script defaults", path + ) + return def serial_port_open(ser): @@ -84,7 +125,16 @@ def serial_port_open(ser): # --------------------------------------------------------------------------- # Configuration - override with environment variables # --------------------------------------------------------------------------- -SERIAL_PORT = os.getenv("AIRCUBE_PORT", "/dev/cu.usbmodem101") +if sys.platform == "darwin": + _DEFAULT_SERIAL_PORT = "/dev/cu.usbmodem101" +else: + _DEFAULT_SERIAL_PORT = "/dev/ttyACM0" + +SERIAL_PORT = os.getenv("AIRCUBE_PORT", _DEFAULT_SERIAL_PORT) +if sys.platform != "darwin" and ( + "usbmodem" in SERIAL_PORT or SERIAL_PORT.startswith("/dev/cu.") +): + SERIAL_PORT = _DEFAULT_SERIAL_PORT SERIAL_BAUD = int(os.getenv("AIRCUBE_BAUD", "115200")) MQTT_HOST = os.getenv("MQTT_HOST", "localhost") MQTT_PORT = int(os.getenv("MQTT_PORT", "1883")) @@ -94,6 +144,9 @@ def serial_port_open(ser): DEVICE_ID = os.getenv("AIRCUBE_ID", "aircube_1") # unique per device DISCOVERY_PREFIX = os.getenv("HA_DISCOVERY_PREFIX", "homeassistant") PUBLISH_INTERVAL = int(os.getenv("AIRCUBE_INTERVAL", "0")) +SENSOR_INTERVAL = float(os.getenv("AIRCUBE_SENSOR_INTERVAL", "1.0")) +COMMAND_GAP = float(os.getenv("AIRCUBE_COMMAND_GAP", "0.15")) +COMMAND_MIN_INTERVAL = float(os.getenv("AIRCUBE_COMMAND_INTERVAL", "1.0")) STATE_TOPIC = "aircube/{0}/state".format(DEVICE_ID) AVAIL_TOPIC = "aircube/{0}/availability".format(DEVICE_ID) @@ -125,11 +178,19 @@ class BridgeState(object): def __init__(self): self.ser = None - self.write_lock = threading.Lock() - self.intensity_ack = threading.Event() + self.mqtt_client = None self.intensity_value = None + self.expected_intensity = None self.sensor_queue = queue.Queue() self.stop = threading.Event() + self.reader_stop = threading.Event() + self.reader_thread = None + self.brightness_lock = threading.Lock() + self.pending_brightness = None + self.last_brightness_at = 0.0 + self.last_sensor_at = 0.0 + self.brightness_due_at = 0.0 + self.read_errors = 0 # --------------------------------------------------------------------------- @@ -204,8 +265,18 @@ def on_disconnect(client, userdata, rc): log.warning("MQTT disconnected (rc=%s), will retry", rc) +def mqtt_text(value): + if isinstance(value, bytes): + return value.decode("utf-8", "replace") + return value + + def build_mqtt_client(bridge): - client = mqtt.Client(client_id=DEVICE_ID, clean_session=True, userdata=bridge) + try: + client = mqtt.Client(client_id=DEVICE_ID, clean_session=True, userdata=bridge) + except TypeError: + client = mqtt.Client(client_id=DEVICE_ID, clean_session=True) + client._userdata = bridge client.will_set(AVAIL_TOPIC, "offline", retain=True) client.on_connect = on_connect client.on_disconnect = on_disconnect @@ -221,7 +292,7 @@ def build_mqtt_client(bridge): def open_serial(port, baud, retries=10): for attempt in range(1, retries + 1): try: - ser = serial.Serial(port, baud, timeout=5) + ser = serial.Serial(port, baud, timeout=1) log.info("Serial opened: %s @ %d baud", port, baud) return ser except serial.SerialException as exc: @@ -233,6 +304,37 @@ def open_serial(port, baud, retries=10): sys.exit(1) +def make_set_intensity_command(intensity): + return '{{"cmd":"set_intensity","value":{0:.2f}}}\n'.format(float(intensity)) + + +def extract_json_from_line(line): + start = line.find("{") + end = line.rfind("}") + if start >= 0 and end > start: + return line[start : end + 1] + return line.strip() + + +def start_reader(bridge): + if bridge.reader_thread is not None and bridge.reader_thread.is_alive(): + return + bridge.reader_stop.clear() + bridge.reader_thread = threading.Thread( + target=serial_reader_loop, args=(bridge,) + ) + bridge.reader_thread.daemon = True + bridge.reader_thread.start() + + +def stop_reader(bridge, timeout=3.0): + bridge.reader_stop.set() + thread = bridge.reader_thread + if thread is not None and thread.is_alive(): + thread.join(timeout) + bridge.reader_thread = None + + def parse_brightness_payload(raw): """ Parse an MQTT brightness command payload. @@ -269,113 +371,145 @@ def parse_brightness_payload(raw): def parse_set_intensity_response(raw): - """Return applied intensity if raw is a set_intensity ack, else None.""" try: - parsed = json.loads(raw.strip()) + parsed = json.loads(extract_json_from_line(raw)) except JSON_DECODE_ERROR: return None - if parsed.get("status") == "ok" and parsed.get("cmd") == "set_intensity": return float(parsed.get("value", 0)) return None def dispatch_serial_line(bridge, line): - """Route one serial line to brightness ack handling or the sensor queue.""" stripped = line.strip() if not stripped: - return + return False applied = parse_set_intensity_response(stripped) if applied is not None: - bridge.intensity_value = applied - bridge.intensity_ack.set() - return + if bridge.expected_intensity is not None: + if abs(applied - bridge.expected_intensity) <= 0.02: + bridge.intensity_value = applied + bridge.expected_intensity = None + if bridge.mqtt_client is not None: + bridge.mqtt_client.publish( + BRIGHTNESS_STATE_TOPIC, + "{0:.2f}".format(applied), + retain=True, + ) + log.info("Brightness set to %.2f", applied) + return False payload = parse_aircube(stripped) if payload is not None: + bridge.last_sensor_at = time.time() + bridge.brightness_due_at = bridge.last_sensor_at + COMMAND_GAP bridge.sensor_queue.put(payload) + return True + + return False + + +def maybe_send_brightness(bridge): + """Send queued brightness once we are past the post-sensor delay.""" + if bridge.pending_brightness is None: + return + if bridge.last_sensor_at <= 0: + return + now = time.time() + if now < bridge.brightness_due_at: + return + if now - bridge.last_brightness_at < COMMAND_MIN_INTERVAL: + return + + intensity = None + with bridge.brightness_lock: + if bridge.pending_brightness is None: + return + intensity = bridge.pending_brightness + bridge.pending_brightness = None + + if not serial_port_open(bridge.ser): + with bridge.brightness_lock: + bridge.pending_brightness = intensity return - log.debug("Skipping non-sensor line: %s", stripped) + command = make_set_intensity_command(intensity) + bridge.expected_intensity = float(intensity) + try: + bridge.ser.write(command.encode("utf-8")) + bridge.ser.flush() + bridge.last_brightness_at = time.time() + log.info( + "Sent set_intensity: %.2f (%.0fms after sensor)", + intensity, + (bridge.last_brightness_at - bridge.last_sensor_at) * 1000, + ) + if bridge.mqtt_client is not None: + bridge.mqtt_client.publish( + BRIGHTNESS_STATE_TOPIC, + "{0:.2f}".format(intensity), + retain=True, + ) + except serial.SerialException as exc: + log.error("Serial write failed: %s", exc) + bridge.expected_intensity = None + with bridge.brightness_lock: + bridge.pending_brightness = intensity + bridge.sensor_queue.put(None) def serial_reader_loop(bridge): - """Single reader for all serial input so command acks are never dropped.""" + """Single owner of serial I/O — old Pi USB wedges on concurrent read+write.""" log.info("Serial reader thread started") - while not bridge.stop.is_set(): - ser = bridge.ser - if not serial_port_open(ser): + while not bridge.reader_stop.is_set() and not bridge.stop.is_set(): + if not serial_port_open(bridge.ser): time.sleep(0.1) continue + maybe_send_brightness(bridge) + try: - raw = ser.readline().decode("utf-8", "replace") + raw = bridge.ser.readline().decode("utf-8", "replace") if raw: + bridge.read_errors = 0 dispatch_serial_line(bridge, raw) except serial.SerialException as exc: - log.error("Serial reader error: %s", exc) - bridge.sensor_queue.put(None) - - -def send_set_intensity(bridge, intensity): - """Send set_intensity to the device and wait for ack from the reader thread.""" - if bridge is None: - log.error("Internal bridge state missing; cannot set brightness") - return False - - if not serial_port_open(bridge.ser): - log.error("Serial port not open; cannot set brightness") - return False - - # Firmware parser expects compact JSON (no spaces). - command = json.dumps( - {"cmd": "set_intensity", "value": intensity}, - separators=(",", ":"), - ) + "\n" - bridge.intensity_ack.clear() - bridge.intensity_value = None - - with bridge.write_lock: - try: - bridge.ser.write(command.encode("utf-8")) - bridge.ser.flush() - log.info("Sent set_intensity command: %.2f", intensity) - except serial.SerialException as exc: - log.error("Serial write failed while setting brightness: %s", exc) - return False + bridge.read_errors += 1 + log.warning( + "Serial read glitch (%d/5): %s", bridge.read_errors, exc + ) + if bridge.read_errors >= 5: + log.error("Too many serial read errors, reconnecting") + bridge.sensor_queue.put(None) + break + time.sleep(0.5) + log.info("Serial reader thread stopped") - if bridge.intensity_ack.wait(5.0): - log.info("Brightness set to %.2f via serial", bridge.intensity_value) - return True - log.warning("Timed out waiting for set_intensity response (command was sent)") - return False +def queue_brightness(bridge, intensity): + with bridge.brightness_lock: + bridge.pending_brightness = intensity def on_message(client, userdata, msg): - if msg.topic != BRIGHTNESS_CMD_TOPIC: + topic = mqtt_text(msg.topic) + if topic != BRIGHTNESS_CMD_TOPIC: return - intensity = parse_brightness_payload(msg.payload.decode("utf-8", "replace")) + intensity = parse_brightness_payload(mqtt_text(msg.payload)) if intensity is None: - log.warning("Invalid brightness payload on %s: %r", msg.topic, msg.payload) + log.warning("Invalid brightness payload on %s: %r", topic, msg.payload) return bridge = userdata + if bridge is None: + bridge = getattr(client, "_userdata", None) if not isinstance(bridge, BridgeState): - log.error( - "MQTT userdata is not bridge state (%r); cannot set brightness", userdata - ) + log.error("MQTT userdata is not bridge state (%r)", userdata) return - if send_set_intensity(bridge, intensity): - applied = ( - bridge.intensity_value - if bridge.intensity_value is not None - else intensity - ) - client.publish(BRIGHTNESS_STATE_TOPIC, "{0:.2f}".format(applied), retain=True) + queue_brightness(bridge, intensity) def parse_aircube(raw): @@ -423,8 +557,14 @@ def main(): log_env_status(log) log.info("Serial: %s", SERIAL_PORT) log.info("Interval: %s seconds", args.interval) + log.info( + "Brightness: %.0fms after each sensor (sensor period %.1fs)", + COMMAND_GAP * 1000, + SENSOR_INTERVAL, + ) bridge = BridgeState() + bridge.mqtt_client = None client = None if not args.no_mqtt: log.info("MQTT: %s:%s State topic: %s", MQTT_HOST, MQTT_PORT, STATE_TOPIC) @@ -436,15 +576,14 @@ def main(): log.error("Could not connect to MQTT broker: %s", exc) sys.exit(1) client.loop_start() + bridge.mqtt_client = client else: log.info("MQTT disabled: running in print-only mode") ser = open_serial(SERIAL_PORT, SERIAL_BAUD) bridge.ser = ser - reader = threading.Thread(target=serial_reader_loop, args=(bridge,)) - reader.daemon = True - reader.start() + start_reader(bridge) consecutive_errors = 0 buffer = collections.defaultdict(list) @@ -504,17 +643,26 @@ def main(): log.error("Serial error: %s (attempt %d)", exc, consecutive_errors) if client: client.publish(AVAIL_TOPIC, "offline", retain=True) - ser.close() + stop_reader(bridge) + try: + ser.close() + except serial.SerialException: + pass bridge.ser = None + bridge.last_sensor_at = 0.0 + bridge.brightness_due_at = 0.0 + bridge.read_errors = 0 time.sleep(5) ser = open_serial(SERIAL_PORT, SERIAL_BAUD) bridge.ser = ser + start_reader(bridge) if client: client.publish(AVAIL_TOPIC, "online", retain=True) except KeyboardInterrupt: log.info("Shutting down") bridge.stop.set() + stop_reader(bridge) if client: client.publish(AVAIL_TOPIC, "offline", retain=True) client.loop_stop()