From bc976e33b3e0b1d4acb48517e114ef4fe924a23d Mon Sep 17 00:00:00 2001 From: Constantin Pascal Date: Sun, 28 Dec 2025 10:58:51 +0200 Subject: [PATCH] Working version --- __init__.py | 57 ++++ binary_sensor.py | 118 ++++++++ climate.py | 143 ++++++++++ config_flow.py | 76 +++++ const.py | 73 +++++ coordinator.py | 48 ++++ hub.py | 727 +++++++++++++++++++++++++++++++++++++++++++++++ manifest.json | 11 + sensor.py | 253 +++++++++++++++++ strings.json | 36 +++ 10 files changed, 1542 insertions(+) create mode 100644 __init__.py create mode 100644 binary_sensor.py create mode 100644 climate.py create mode 100644 config_flow.py create mode 100644 const.py create mode 100644 coordinator.py create mode 100644 hub.py create mode 100644 manifest.json create mode 100644 sensor.py create mode 100644 strings.json diff --git a/__init__.py b/__init__.py new file mode 100644 index 0000000..30a13cc --- /dev/null +++ b/__init__.py @@ -0,0 +1,57 @@ +"""The Modbus RTU Monitor integration.""" + +from __future__ import annotations + +from homeassistant.config_entries import ConfigEntry +from homeassistant.const import CONF_HOST, CONF_PORT, Platform +from homeassistant.core import HomeAssistant + +from .coordinator import ModbusRTUMonitorCoordinator +from .hub import ModbusRTUMonitorHub + +type ModbusRTUMonitorConfigEntry = ConfigEntry[ModbusRTUMonitorCoordinator] + +PLATFORMS: list[Platform] = [Platform.BINARY_SENSOR, Platform.CLIMATE, Platform.SENSOR] + + +async def async_setup_entry( + hass: HomeAssistant, entry: ModbusRTUMonitorConfigEntry +) -> bool: + """Set up Modbus RTU Monitor from a config entry.""" + # Create hub + hub = ModbusRTUMonitorHub( + hass, + entry, + entry.data[CONF_HOST], + entry.data[CONF_PORT], + ) + + # Create coordinator + coordinator = ModbusRTUMonitorCoordinator(hass, entry, hub) + hub.coordinator = coordinator # Link back to coordinator + + # Setup hub (connect and start monitoring) + if not await hub.async_setup(): + return False + + # Store in runtime_data + entry.runtime_data = coordinator + + # Forward to platforms + await hass.config_entries.async_forward_entry_setups(entry, PLATFORMS) + + return True + + +async def async_unload_entry( + hass: HomeAssistant, entry: ModbusRTUMonitorConfigEntry +) -> bool: + """Unload a config entry.""" + # Unload platforms + unload_ok = await hass.config_entries.async_unload_platforms(entry, PLATFORMS) + + if unload_ok: + # Close hub + await entry.runtime_data.hub.async_close() + + return unload_ok diff --git a/binary_sensor.py b/binary_sensor.py new file mode 100644 index 0000000..174bf8a --- /dev/null +++ b/binary_sensor.py @@ -0,0 +1,118 @@ +"""Binary sensor platform for Modbus RTU Monitor.""" + +from __future__ import annotations + +from homeassistant.components.binary_sensor import BinarySensorEntity +from homeassistant.core import HomeAssistant, callback +from homeassistant.helpers.device_registry import DeviceInfo +from homeassistant.helpers.entity_platform import AddEntitiesCallback +from homeassistant.helpers.update_coordinator import CoordinatorEntity + +from . import ModbusRTUMonitorConfigEntry +from .const import COIL_COUNT, DOMAIN, MANUFACTURER, MODEL +from .coordinator import ModbusRTUMonitorCoordinator + + +async def async_setup_entry( + hass: HomeAssistant, + entry: ModbusRTUMonitorConfigEntry, + async_add_entities: AddEntitiesCallback, +) -> None: + """Set up binary sensor entities from a config entry.""" + coordinator = entry.runtime_data + + # Track which slaves have binary sensor entities + added_slaves: set[int] = set() + + @callback + def _async_add_binary_sensor_entities() -> None: + """Add binary sensor entities for newly discovered slaves with coil data.""" + if coordinator.data is None: + return + + current_slaves = { + slave_id + for slave_id, slave_data in coordinator.data.items() + if slave_data.coils is not None + } + new_slaves = current_slaves - added_slaves + + if new_slaves: + entities = [] + for slave_id in new_slaves: + # Create 40 binary sensors (coil1-40) for each slave + for coil_idx in range(COIL_COUNT): + entities.append( + ModbusRTUMonitorCoilSensor( + coordinator, slave_id, coil_idx + 1 + ) + ) + async_add_entities(entities) + added_slaves.update(new_slaves) + + # Subscribe to coordinator updates + entry.async_on_unload( + coordinator.async_add_listener(_async_add_binary_sensor_entities) + ) + + # Add entities for already discovered slaves with coil data + _async_add_binary_sensor_entities() + + +class ModbusRTUMonitorCoilSensor( + CoordinatorEntity[ModbusRTUMonitorCoordinator], BinarySensorEntity +): + """Binary sensor entity for Modbus RTU coil state.""" + + _attr_has_entity_name = True + + def __init__( + self, + coordinator: ModbusRTUMonitorCoordinator, + slave_id: int, + coil_number: int, + ) -> None: + """Initialize binary sensor entity.""" + super().__init__(coordinator) + + self._slave_id = slave_id + self._coil_number = coil_number + self._coil_idx = coil_number - 1 # Convert to 0-based index + + self._attr_unique_id = ( + f"{coordinator.config_entry.entry_id}_{slave_id}_coil{coil_number}" + ) + self._attr_name = f"Coil{coil_number}" + + # Device info (groups with climate and humidity sensor) + self._attr_device_info = DeviceInfo( + identifiers={(DOMAIN, f"{coordinator.config_entry.entry_id}_{slave_id}")}, + name=f"Modbus Slave {slave_id}", + manufacturer=MANUFACTURER, + model=MODEL, + ) + + @property + def available(self) -> bool: + """Return if entity is available.""" + if self._slave_id not in self.coordinator.data: + return False + slave_data = self.coordinator.data[self._slave_id] + if not slave_data.available or slave_data.coils is None: + return False + return self._coil_idx < len(slave_data.coils) + + @property + def is_on(self) -> bool | None: + """Return true if the coil is on.""" + if self._slave_id not in self.coordinator.data: + return None + slave_data = self.coordinator.data[self._slave_id] + if slave_data.coils is None or self._coil_idx >= len(slave_data.coils): + return None + return slave_data.coils[self._coil_idx] + + @property + def extra_state_attributes(self) -> dict[str, int]: + """Return additional state attributes.""" + return {"coil_number": self._coil_number} \ No newline at end of file diff --git a/climate.py b/climate.py new file mode 100644 index 0000000..220ef5a --- /dev/null +++ b/climate.py @@ -0,0 +1,143 @@ +"""Climate platform for Modbus RTU Monitor.""" + +from __future__ import annotations + +from typing import Any + +from homeassistant.components.climate import ( + ClimateEntity, + ClimateEntityFeature, + HVACMode, +) +from homeassistant.const import ATTR_TEMPERATURE, UnitOfTemperature +from homeassistant.core import HomeAssistant, callback +from homeassistant.exceptions import HomeAssistantError +from homeassistant.helpers.device_registry import DeviceInfo +from homeassistant.helpers.entity_platform import AddEntitiesCallback +from homeassistant.helpers.update_coordinator import CoordinatorEntity + +from . import ModbusRTUMonitorConfigEntry +from .const import DOMAIN, MANUFACTURER, MODEL +from .coordinator import ModbusRTUMonitorCoordinator + + +async def async_setup_entry( + hass: HomeAssistant, + entry: ModbusRTUMonitorConfigEntry, + async_add_entities: AddEntitiesCallback, +) -> None: + """Set up climate entities from a config entry.""" + coordinator = entry.runtime_data + + # Track which slaves have climate entities + added_slaves: set[int] = set() + + @callback + def _async_add_climate_entities() -> None: + """Add climate entities for newly discovered slaves.""" + if coordinator.data is None: + return + + current_slaves = set(coordinator.data.keys()) + new_slaves = current_slaves - added_slaves + + if new_slaves: + async_add_entities( + ModbusRTUMonitorClimate(coordinator, slave_id) + for slave_id in new_slaves + ) + added_slaves.update(new_slaves) + + # Subscribe to coordinator updates + entry.async_on_unload( + coordinator.async_add_listener(_async_add_climate_entities) + ) + + # Add entities for already discovered slaves + _async_add_climate_entities() + + +class ModbusRTUMonitorClimate( + CoordinatorEntity[ModbusRTUMonitorCoordinator], ClimateEntity +): + """Climate entity for Modbus RTU temperature monitor.""" + + _attr_has_entity_name = True + _attr_name = None # Device name will be used + _attr_temperature_unit = UnitOfTemperature.CELSIUS + _attr_hvac_modes = [HVACMode.HEAT] # Always on, no off mode + _attr_supported_features = ClimateEntityFeature.TARGET_TEMPERATURE + + def __init__( + self, + coordinator: ModbusRTUMonitorCoordinator, + slave_id: int, + ) -> None: + """Initialize climate entity.""" + super().__init__(coordinator) + + self._slave_id = slave_id + self._attr_unique_id = ( + f"{coordinator.config_entry.entry_id}_{slave_id}_climate" + ) + + # Device info (groups climate + humidity sensor) + self._attr_device_info = DeviceInfo( + identifiers={(DOMAIN, f"{coordinator.config_entry.entry_id}_{slave_id}")}, + name=f"Modbus Slave {slave_id}", + manufacturer=MANUFACTURER, + model=MODEL, + ) + + # Target temperature storage (not available from passive monitoring) + self._attr_target_temperature = 20.0 # Default + + @property + def available(self) -> bool: + """Return if entity is available.""" + if self._slave_id not in self.coordinator.data: + return False + return self.coordinator.data[self._slave_id].available + + @property + def current_temperature(self) -> float | None: + """Return current temperature from passive monitoring.""" + if self._slave_id not in self.coordinator.data: + return None + return self.coordinator.data[self._slave_id].temperature + + @property + def hvac_mode(self) -> HVACMode: + """Return current HVAC mode.""" + # Always HEAT mode (no on/off control as per requirements) + return HVACMode.HEAT + + async def async_set_temperature(self, **kwargs: Any) -> None: + """Set new target temperature and write to device.""" + if (temperature := kwargs.get(ATTR_TEMPERATURE)) is None: + return + + # Update local target temperature + self._attr_target_temperature = temperature + + # Write to device via hub + try: + await self.coordinator.hub.async_write_setpoint( + self._slave_id, + temperature, + ) + except Exception as ex: + raise HomeAssistantError( + f"Failed to set temperature for slave {self._slave_id}" + ) from ex + + # Update state + self.async_write_ha_state() + + async def async_set_hvac_mode(self, hvac_mode: HVACMode) -> None: + """Set HVAC mode.""" + # No on/off control - only HEAT mode supported + if hvac_mode != HVACMode.HEAT: + raise HomeAssistantError( + f"HVAC mode {hvac_mode} not supported. Only HEAT mode is available." + ) diff --git a/config_flow.py b/config_flow.py new file mode 100644 index 0000000..b102f63 --- /dev/null +++ b/config_flow.py @@ -0,0 +1,76 @@ +"""Config flow for Modbus RTU Monitor integration.""" + +from __future__ import annotations + +import asyncio +import logging +from typing import Any + +import voluptuous as vol + +from homeassistant.config_entries import ConfigFlow, ConfigFlowResult +from homeassistant.const import CONF_HOST, CONF_PORT +import homeassistant.helpers.config_validation as cv + +from .const import DEFAULT_PORT, DOMAIN + +_LOGGER = logging.getLogger(__name__) + + +class ModbusRTUMonitorConfigFlow(ConfigFlow, domain=DOMAIN): + """Handle a config flow for Modbus RTU Monitor.""" + + VERSION = 1 + MINOR_VERSION = 1 + + async def async_step_user( + self, user_input: dict[str, Any] | None = None + ) -> ConfigFlowResult: + """Handle the initial step.""" + errors: dict[str, str] = {} + + if user_input is not None: + # Validate connection by attempting to connect + try: + await self._test_connection( + user_input[CONF_HOST], + user_input[CONF_PORT], + ) + except asyncio.TimeoutError: + errors["base"] = "cannot_connect" + except Exception: + _LOGGER.exception("Unexpected exception during connection test") + errors["base"] = "unknown" + else: + # Prevent duplicate entries for same host:port + self._async_abort_entries_match( + { + CONF_HOST: user_input[CONF_HOST], + CONF_PORT: user_input[CONF_PORT], + } + ) + + title = f"Modbus RTU Monitor ({user_input[CONF_HOST]}:{user_input[CONF_PORT]})" + return self.async_create_entry(title=title, data=user_input) + + data_schema = vol.Schema( + { + vol.Required(CONF_HOST): str, + vol.Required(CONF_PORT, default=DEFAULT_PORT): cv.port, + } + ) + + return self.async_show_form( + step_id="user", + data_schema=data_schema, + errors=errors, + ) + + async def _test_connection(self, host: str, port: int) -> None: + """Test TCP connection to Modbus gateway.""" + reader, writer = await asyncio.wait_for( + asyncio.open_connection(host, port), + timeout=5.0, + ) + writer.close() + await writer.wait_closed() diff --git a/const.py b/const.py new file mode 100644 index 0000000..072b183 --- /dev/null +++ b/const.py @@ -0,0 +1,73 @@ +"""Constants for the Modbus RTU Monitor integration.""" + +from dataclasses import dataclass +from datetime import datetime, timedelta +from typing import TYPE_CHECKING + +if TYPE_CHECKING: + pass + +DOMAIN = "modbus_rtu_monitor" + +# Default values +DEFAULT_PORT = 502 + +# Modbus register configuration +TEMP_HUMIDITY_REGISTER = 0x83 # Register 131 (decimal) +SETPOINT_REGISTER = 144 # Register 144 (decimal, 0x90 hex) +REGISTER_COUNT = 4 + +# Coil configuration +COIL_START_ADDRESS = 0x01 # Starting coil address (1 in decimal) +COIL_COUNT = 40 # Number of coils to monitor + +# Register monitoring configuration +REGISTER_START_ADDRESS = 0xA5 # Register 165 (decimal) +REGISTER_MONITOR_COUNT = 20 # Number of registers to monitor (165-184) + +# Additional register monitoring +REGISTER_START_ADDRESS_2 = 0xD2 # Register 210 (decimal) +REGISTER_MONITOR_COUNT_2 = 8 # Number of registers to monitor (210-217) + +# Third register range +REGISTER_START_ADDRESS_3 = 0x8C # Register 140 (decimal) +REGISTER_MONITOR_COUNT_3 = 23 # Number of registers to monitor (140-162) + +# Availability timeout (if no data for 5 minutes, mark unavailable) +AVAILABILITY_TIMEOUT = timedelta(minutes=5) + +# Update throttling (prevent overwhelming HA with high-frequency Modbus traffic) +COORDINATOR_UPDATE_INTERVAL = timedelta(seconds=5) # Minimum time between coordinator updates +FRAME_PROCESSING_DELAY = 0.1 # Seconds to sleep between frame processing batches + +# Device info +MANUFACTURER = "Modbus RTU" +MODEL = "Temperature/Humidity Monitor" + + +@dataclass +class SlaveData: + """Data for a discovered Modbus slave.""" + + slave_id: int + temperature: float | None + humidity: float | None + last_seen: datetime + available: bool = True + coils: list[bool] | None = None # 40 coil states (True=ON, False=OFF) + registers: dict[int, int] | None = None # Register address -> int16 value + + +@dataclass +class ModbusRTUFrame: + """Decoded Modbus RTU frame.""" + + slave_id: int + function_code: int + data: bytes + is_request: bool + is_response: bool + start_address: int | None = None + register_count: int | None = None + values: list[int] | None = None + coil_values: list[bool] | None = None # For function code 01 (Read Coils) responses diff --git a/coordinator.py b/coordinator.py new file mode 100644 index 0000000..2d6bbe9 --- /dev/null +++ b/coordinator.py @@ -0,0 +1,48 @@ +"""DataUpdateCoordinator for Modbus RTU Monitor.""" + +from __future__ import annotations + +import logging +from typing import TYPE_CHECKING + +from homeassistant.helpers.update_coordinator import DataUpdateCoordinator + +from .const import DOMAIN, SlaveData + +if TYPE_CHECKING: + from homeassistant.config_entries import ConfigEntry + from homeassistant.core import HomeAssistant + + from .hub import ModbusRTUMonitorHub + +_LOGGER = logging.getLogger(__name__) + + +class ModbusRTUMonitorCoordinator(DataUpdateCoordinator[dict[int, SlaveData]]): + """Coordinator for Modbus RTU Monitor state updates.""" + + def __init__( + self, + hass: HomeAssistant, + config_entry: ConfigEntry, + hub: ModbusRTUMonitorHub, + ) -> None: + """Initialize coordinator.""" + self.hub = hub + + # No update_interval - updates come from passive monitoring + super().__init__( + hass, + _LOGGER, + config_entry=config_entry, + name=DOMAIN, + update_interval=None, # Push-based updates + ) + + async def _async_update_data(self) -> dict[int, SlaveData]: + """Not used - data comes from passive monitoring. + + This method won't be called due to update_interval=None. + Updates happen via hub calling async_set_updated_data(). + """ + return self.hub.discovered_slaves diff --git a/hub.py b/hub.py new file mode 100644 index 0000000..6f1d21a --- /dev/null +++ b/hub.py @@ -0,0 +1,727 @@ +"""Modbus RTU Monitor Hub for passive monitoring and active writes.""" + +from __future__ import annotations + +import asyncio +import logging +import struct +from typing import TYPE_CHECKING + +from homeassistant.core import HomeAssistant +from homeassistant.exceptions import ConfigEntryNotReady, HomeAssistantError +from homeassistant.helpers.event import async_call_later +from homeassistant.util.dt import utcnow + +from .const import ( + AVAILABILITY_TIMEOUT, + COIL_COUNT, + COIL_START_ADDRESS, + COORDINATOR_UPDATE_INTERVAL, + DOMAIN, + FRAME_PROCESSING_DELAY, + REGISTER_COUNT, + REGISTER_MONITOR_COUNT, + REGISTER_MONITOR_COUNT_2, + REGISTER_MONITOR_COUNT_3, + REGISTER_START_ADDRESS, + REGISTER_START_ADDRESS_2, + REGISTER_START_ADDRESS_3, + SETPOINT_REGISTER, + TEMP_HUMIDITY_REGISTER, + ModbusRTUFrame, + SlaveData, +) + +if TYPE_CHECKING: + from homeassistant.config_entries import ConfigEntry + + from .coordinator import ModbusRTUMonitorCoordinator + +_LOGGER = logging.getLogger(__name__) + + +class ModbusRTUDecoder: + """Decode Modbus RTU frames.""" + + @staticmethod + def calculate_crc(data: bytes) -> int: + """Calculate Modbus RTU CRC.""" + crc = 0xFFFF + for byte in data: + crc ^= byte + for _ in range(8): + if crc & 0x0001: + crc = (crc >> 1) ^ 0xA001 + else: + crc >>= 1 + return crc + + @staticmethod + def build_write_register_frame(slave_id: int, register: int, value: int) -> bytes: + """Build a Modbus RTU Write Single Register frame (function code 0x06). + + Args: + slave_id: Modbus slave ID (1-247) + register: Register address (0-65535) + value: Register value to write (0-65535) + + Returns: + Complete RTU frame with CRC + """ + # Function code 0x06 = Write Single Register + frame = bytearray() + frame.append(slave_id) + frame.append(0x06) # Function code + frame.extend(struct.pack(">H", register)) # Register address (big-endian) + frame.extend(struct.pack(">H", value)) # Register value (big-endian) + + # Calculate and append CRC (little-endian) + crc = ModbusRTUDecoder.calculate_crc(bytes(frame)) + frame.extend(struct.pack(" ModbusRTUFrame | None: + """Decode a Modbus RTU frame.""" + if len(data) < 4: + return None + + slave_id = data[0] + function_code = data[1] + frame_data = data[2:-2] + received_crc = struct.unpack("H", frame_data[0:2])[0] + count = struct.unpack(">H", frame_data[2:4])[0] + decoded.is_request = True + decoded.start_address = start_addr + decoded.register_count = count + + # Detect response (function 1, Read Coils - with byte count and coil values) + elif function_code == 1 and len(frame_data) > 1: + byte_count = frame_data[0] + if len(frame_data) >= byte_count + 1: + # Decode coil bitmap + coil_values = [] + for byte_idx in range(byte_count): + byte_val = frame_data[byte_idx + 1] + # Each byte contains 8 coil states (LSB first) + for bit_idx in range(8): + coil_values.append(bool(byte_val & (1 << bit_idx))) + decoded.is_response = True + decoded.coil_values = coil_values + + # Detect request (function 3, with address and count) + elif function_code == 3 and len(frame_data) == 4: + start_addr = struct.unpack(">H", frame_data[0:2])[0] + count = struct.unpack(">H", frame_data[2:4])[0] + decoded.is_request = True + decoded.start_address = start_addr + decoded.register_count = count + + # Detect response (function 3, with byte count and values) + elif function_code == 3 and len(frame_data) > 1: + byte_count = frame_data[0] + if len(frame_data) >= byte_count + 1: + values = [] + for i in range(0, byte_count, 2): + if i + 1 < byte_count: + val = struct.unpack(">H", frame_data[i + 1 : i + 3])[0] + values.append(val) + decoded.is_response = True + decoded.values = values + + return decoded + + +class ModbusRTUMonitorHub: + """Hub for passive Modbus RTU monitoring and active writes.""" + + def __init__( + self, hass: HomeAssistant, config_entry: ConfigEntry, host: str, port: int + ) -> None: + """Initialize the hub.""" + self.hass = hass + self.config_entry = config_entry + self.host = host + self.port = port + + # Asyncio TCP connection + self._reader: asyncio.StreamReader | None = None + self._writer: asyncio.StreamWriter | None = None + self._monitor_task: asyncio.Task | None = None + + # Coordination lock for write operations + self._write_lock = asyncio.Lock() + + # Slave discovery and state + self.discovered_slaves: dict[int, SlaveData] = {} + self._pending_requests: dict[int | str, dict] = {} + + # Frame decoder + self._decoder = ModbusRTUDecoder() + + # Coordinator reference (set after initialization) + self.coordinator: ModbusRTUMonitorCoordinator | None = None + + # Availability check task + self._availability_task: asyncio.Task | None = None + + # Throttling for coordinator updates + self._last_coordinator_update = utcnow() - COORDINATOR_UPDATE_INTERVAL + + def _should_update_coordinator(self) -> bool: + """Check if enough time has passed to update the coordinator.""" + return (utcnow() - self._last_coordinator_update) >= COORDINATOR_UPDATE_INTERVAL + + def _update_coordinator_throttled(self) -> None: + """Update coordinator with throttling to prevent overwhelming HA.""" + if self.coordinator and self._should_update_coordinator(): + self.coordinator.async_set_updated_data(self.discovered_slaves) + self._last_coordinator_update = utcnow() + + async def async_setup(self) -> bool: + """Set up the hub and start monitoring.""" + try: + await self._connect() + self._monitor_task = asyncio.create_task(self._monitor_loop()) + self._availability_task = asyncio.create_task( + self._availability_check_loop() + ) + return True + except Exception as ex: + raise ConfigEntryNotReady( + f"Failed to connect to {self.host}:{self.port}" + ) from ex + + async def _connect(self) -> None: + """Establish TCP connection.""" + self._reader, self._writer = await asyncio.wait_for( + asyncio.open_connection(self.host, self.port), + timeout=10.0, + ) + _LOGGER.info( + "Connected to Modbus RTU gateway at %s:%s", self.host, self.port + ) + + async def _reconnect(self) -> None: + """Reconnect to Modbus gateway after connection loss.""" + # Close existing connection + if self._writer: + self._writer.close() + await self._writer.wait_closed() + + # Wait before reconnecting + await asyncio.sleep(5) + + # Attempt reconnection + try: + await self._connect() + _LOGGER.info("Reconnected to Modbus gateway") + except Exception as ex: + _LOGGER.error("Reconnection failed: %s", ex) + # Will retry on next loop iteration + + async def _monitor_loop(self) -> None: + """Main monitoring loop for passive frame reception.""" + buffer = bytearray() + + while True: + try: + # Read with timeout to allow graceful shutdown + data = await asyncio.wait_for(self._reader.read(1024), timeout=1.0) + + if not data: + # Connection closed + _LOGGER.warning("Connection closed, attempting reconnect") + await self._reconnect() + continue + + buffer.extend(data) + + # Process frames from buffer (limit iterations to prevent blocking) + max_iterations = 100 + iterations = 0 + while len(buffer) >= 4 and iterations < max_iterations: + frame_found = False + + # Try different frame lengths + for frame_len in range(4, min(len(buffer) + 1, 256)): + potential_frame = bytes(buffer[:frame_len]) + decoded = self._decoder.decode_frame(potential_frame) + + if decoded: + await self._handle_frame(decoded) + buffer = buffer[frame_len:] + frame_found = True + break + + if not frame_found: + # Invalid data, discard first byte + buffer = buffer[1:] + + iterations += 1 + + # If we hit the iteration limit, yield to event loop + if iterations >= max_iterations and len(buffer) >= 4: + await asyncio.sleep(0) # Yield control to event loop + + # Rate limiting: sleep after processing to avoid overwhelming HA + if iterations > 0: + await asyncio.sleep(FRAME_PROCESSING_DELAY) + + except asyncio.TimeoutError: + # Normal timeout, continue + continue + except asyncio.CancelledError: + break + except Exception as ex: + _LOGGER.error("Error in monitor loop: %s", ex) + await asyncio.sleep(5) + + async def _handle_frame(self, frame: ModbusRTUFrame) -> None: + """Handle a decoded Modbus RTU frame.""" + slave_id = frame.slave_id + + # Detect Read Coils request (function 01) + if ( + frame.is_request + and frame.function_code == 1 + and frame.start_address == COIL_START_ADDRESS + and frame.register_count == COIL_COUNT + ): + self._pending_requests[f"{slave_id}_coils"] = { + "timestamp": utcnow(), + "slave_id": slave_id, + "type": "coils", + } + + # Handle Read Coils response + elif ( + frame.is_response + and frame.function_code == 1 + and f"{slave_id}_coils" in self._pending_requests + and frame.coil_values + ): + # Ensure slave exists + if slave_id not in self.discovered_slaves: + _LOGGER.info("Discovered new slave via coils: %s", slave_id) + self.discovered_slaves[slave_id] = SlaveData( + slave_id=slave_id, + temperature=None, + humidity=None, + last_seen=utcnow(), + available=False, + ) + + # Update coil data (limit to COIL_COUNT coils) + slave_data = self.discovered_slaves[slave_id] + slave_data.coils = frame.coil_values[:COIL_COUNT] + slave_data.last_seen = utcnow() + slave_data.available = True + + # Clear pending request + del self._pending_requests[f"{slave_id}_coils"] + + _LOGGER.debug( + "Slave %s: Updated %d coils", + slave_id, + len(slave_data.coils), + ) + + # Notify coordinator (throttled) + self._update_coordinator_throttled() + + # Detect Read Holding Registers request (function 03) for monitored registers + elif ( + frame.is_request + and frame.function_code == 3 + and frame.start_address == REGISTER_START_ADDRESS + and frame.register_count == REGISTER_MONITOR_COUNT + ): + self._pending_requests[f"{slave_id}_registers"] = { + "timestamp": utcnow(), + "slave_id": slave_id, + "type": "registers", + "start_address": frame.start_address, + } + + # Handle Read Holding Registers response for monitored registers + elif ( + frame.is_response + and frame.function_code == 3 + and f"{slave_id}_registers" in self._pending_requests + and frame.values + ): + # Ensure slave exists + if slave_id not in self.discovered_slaves: + _LOGGER.info("Discovered new slave via registers: %s", slave_id) + self.discovered_slaves[slave_id] = SlaveData( + slave_id=slave_id, + temperature=None, + humidity=None, + last_seen=utcnow(), + available=False, + ) + + # Get start address from pending request + start_addr = self._pending_requests[f"{slave_id}_registers"]["start_address"] + + # Update register data - convert unsigned to signed int16 + slave_data = self.discovered_slaves[slave_id] + if slave_data.registers is None: + slave_data.registers = {} + + for idx, value in enumerate(frame.values): + register_addr = start_addr + idx + # Convert unsigned 16-bit to signed int16 + if value > 32767: + signed_value = value - 65536 + else: + signed_value = value + slave_data.registers[register_addr] = signed_value + + slave_data.last_seen = utcnow() + slave_data.available = True + + # Clear pending request + del self._pending_requests[f"{slave_id}_registers"] + + _LOGGER.debug( + "Slave %s: Updated %d registers starting at %d", + slave_id, + len(frame.values), + start_addr, + ) + + # Notify coordinator (throttled) + self._update_coordinator_throttled() + + # Detect Read Holding Registers request (function 03) for second monitored register range + elif ( + frame.is_request + and frame.function_code == 3 + and frame.start_address == REGISTER_START_ADDRESS_2 + and frame.register_count == REGISTER_MONITOR_COUNT_2 + ): + self._pending_requests[f"{slave_id}_registers2"] = { + "timestamp": utcnow(), + "slave_id": slave_id, + "type": "registers2", + "start_address": frame.start_address, + } + + # Handle Read Holding Registers response for second monitored register range + elif ( + frame.is_response + and frame.function_code == 3 + and f"{slave_id}_registers2" in self._pending_requests + and frame.values + ): + # Ensure slave exists + if slave_id not in self.discovered_slaves: + _LOGGER.info("Discovered new slave via registers2: %s", slave_id) + self.discovered_slaves[slave_id] = SlaveData( + slave_id=slave_id, + temperature=None, + humidity=None, + last_seen=utcnow(), + available=False, + ) + + # Get start address from pending request + start_addr = self._pending_requests[f"{slave_id}_registers2"]["start_address"] + + # Update register data - convert unsigned to signed int16 + slave_data = self.discovered_slaves[slave_id] + if slave_data.registers is None: + slave_data.registers = {} + + for idx, value in enumerate(frame.values): + register_addr = start_addr + idx + # Convert unsigned 16-bit to signed int16 + if value > 32767: + signed_value = value - 65536 + else: + signed_value = value + slave_data.registers[register_addr] = signed_value + + slave_data.last_seen = utcnow() + slave_data.available = True + + # Clear pending request + del self._pending_requests[f"{slave_id}_registers2"] + + _LOGGER.debug( + "Slave %s: Updated %d registers (range 2) starting at %d", + slave_id, + len(frame.values), + start_addr, + ) + + # Notify coordinator (throttled) + self._update_coordinator_throttled() + + # Detect Read Holding Registers request (function 03) for third monitored register range + elif ( + frame.is_request + and frame.function_code == 3 + and frame.start_address == REGISTER_START_ADDRESS_3 + and frame.register_count == REGISTER_MONITOR_COUNT_3 + ): + self._pending_requests[f"{slave_id}_registers3"] = { + "timestamp": utcnow(), + "slave_id": slave_id, + "type": "registers3", + "start_address": frame.start_address, + } + + # Handle Read Holding Registers response for third monitored register range + elif ( + frame.is_response + and frame.function_code == 3 + and f"{slave_id}_registers3" in self._pending_requests + and frame.values + ): + # Ensure slave exists + if slave_id not in self.discovered_slaves: + _LOGGER.info("Discovered new slave via registers3: %s", slave_id) + self.discovered_slaves[slave_id] = SlaveData( + slave_id=slave_id, + temperature=None, + humidity=None, + last_seen=utcnow(), + available=False, + ) + + # Get start address from pending request + start_addr = self._pending_requests[f"{slave_id}_registers3"]["start_address"] + + # Update register data - convert unsigned to signed int16 + slave_data = self.discovered_slaves[slave_id] + if slave_data.registers is None: + slave_data.registers = {} + + for idx, value in enumerate(frame.values): + register_addr = start_addr + idx + # Convert unsigned 16-bit to signed int16 + if value > 32767: + signed_value = value - 65536 + else: + signed_value = value + slave_data.registers[register_addr] = signed_value + + slave_data.last_seen = utcnow() + slave_data.available = True + + # Clear pending request + del self._pending_requests[f"{slave_id}_registers3"] + + _LOGGER.debug( + "Slave %s: Updated %d registers (range 3) starting at %d", + slave_id, + len(frame.values), + start_addr, + ) + + # Notify coordinator (throttled) + self._update_coordinator_throttled() + + # Detect discovery: request to register 0x83 with 4 registers + elif ( + frame.is_request + and frame.start_address == TEMP_HUMIDITY_REGISTER + and frame.register_count == REGISTER_COUNT + ): + self._pending_requests[slave_id] = { + "timestamp": utcnow(), + "slave_id": slave_id, + } + + # Auto-discover new slave + if slave_id not in self.discovered_slaves: + _LOGGER.info("Discovered new slave: %s", slave_id) + self.discovered_slaves[slave_id] = SlaveData( + slave_id=slave_id, + temperature=None, + humidity=None, + last_seen=utcnow(), + available=False, + ) + # Trigger entity creation (throttled) + self._update_coordinator_throttled() + + # Handle response with temperature/humidity data + elif ( + frame.is_response + and slave_id in self._pending_requests + and frame.values + and len(frame.values) >= 4 + ): + # Extract temperature (index 2) and humidity (index 3) + temp_raw = frame.values[2] + humidity_raw = frame.values[3] + + # Scale values (divide by 10) + temperature = temp_raw / 10.0 + humidity = humidity_raw / 10.0 + + # Update slave data + slave_data = self.discovered_slaves[slave_id] + slave_data.temperature = temperature + slave_data.humidity = humidity + slave_data.last_seen = utcnow() + slave_data.available = True + + # Clear pending request + del self._pending_requests[slave_id] + + _LOGGER.debug( + "Slave %s: Temperature=%.1f°C, Humidity=%.1f%%", + slave_id, + temperature, + humidity, + ) + + # Notify coordinator (throttled) + self._update_coordinator_throttled() + + async def _availability_check_loop(self) -> None: + """Periodically check slave availability.""" + while True: + try: + await asyncio.sleep(60) # Check every minute + self._check_slave_availability() + except asyncio.CancelledError: + break + except Exception as ex: + _LOGGER.error("Error in availability check loop: %s", ex) + + def _check_slave_availability(self) -> None: + """Mark slaves unavailable if no data received recently.""" + now = utcnow() + updated = False + + for slave_data in self.discovered_slaves.values(): + if (now - slave_data.last_seen) > AVAILABILITY_TIMEOUT: + if slave_data.available: + slave_data.available = False + updated = True + _LOGGER.warning("Slave %s marked unavailable", slave_data.slave_id) + + if updated: + self._update_coordinator_throttled() + + async def async_write_setpoint( + self, slave_id: int, temperature: float + ) -> None: + """Write temperature setpoint to slave using RTU protocol over existing connection.""" + async with self._write_lock: + # Ensure writer is available + if not self._writer: + raise HomeAssistantError( + "No active connection to Modbus gateway - cannot write setpoint" + ) + + # Scale temperature value (multiply by 10) + value = int(temperature * 10) + + # Ensure value fits in uint16 + if not 0 <= value <= 65535: + raise HomeAssistantError( + f"Setpoint value {value} out of valid range (0-65535)" + ) + + _LOGGER.info( + "Writing setpoint via RTU - slave_id=%s, temp=%.1f°C, scaled_value=%d, register=%s", + slave_id, + temperature, + value, + SETPOINT_REGISTER, + ) + + try: + # Build Modbus RTU Write Single Register frame + frame = self._decoder.build_write_register_frame( + slave_id=slave_id, + register=SETPOINT_REGISTER, + value=value, + ) + + _LOGGER.debug( + "Sending RTU frame: %s (length=%d)", + frame.hex(" "), + len(frame), + ) + + # Send frame via existing TCP connection + self._writer.write(frame) + await self._writer.drain() + + _LOGGER.info( + "Successfully sent setpoint %.1f°C (value=%d) to slave %s register %s", + temperature, + value, + slave_id, + SETPOINT_REGISTER, + ) + + # Note: We don't wait for response here as the passive monitoring + # loop will receive and process it. The Write Single Register + # response echoes the request if successful, or returns an error code. + + except OSError as ex: + _LOGGER.exception( + "Connection error writing setpoint to slave %s: %s", + slave_id, + ex, + ) + raise HomeAssistantError( + f"Connection error writing setpoint to slave {slave_id}" + ) from ex + except Exception as ex: + _LOGGER.exception( + "Unexpected error writing setpoint to slave %s: %s", + slave_id, + ex, + ) + raise HomeAssistantError( + f"Unexpected error writing setpoint to slave {slave_id}" + ) from ex + + async def async_close(self) -> None: + """Close hub and cleanup resources.""" + if self._monitor_task: + self._monitor_task.cancel() + try: + await self._monitor_task + except asyncio.CancelledError: + pass + + if self._availability_task: + self._availability_task.cancel() + try: + await self._availability_task + except asyncio.CancelledError: + pass + + if self._writer: + self._writer.close() + await self._writer.wait_closed() + + _LOGGER.info("Modbus RTU Monitor hub closed") diff --git a/manifest.json b/manifest.json new file mode 100644 index 0000000..698a06a --- /dev/null +++ b/manifest.json @@ -0,0 +1,11 @@ +{ + "domain": "modbus_rtu_monitor", + "name": "Modbus RTU Monitor", + "codeowners": ["@constantin"], + "config_flow": true, + "documentation": "https://github.com/constantin/modbus_rtu_monitor", + "integration_type": "hub", + "iot_class": "local_push", + "requirements": ["pymodbus==3.11.2"], + "version": "1.0.0" +} diff --git a/sensor.py b/sensor.py new file mode 100644 index 0000000..aa456f4 --- /dev/null +++ b/sensor.py @@ -0,0 +1,253 @@ +"""Sensor platform for Modbus RTU Monitor.""" + +from __future__ import annotations + +from homeassistant.components.sensor import ( + SensorDeviceClass, + SensorEntity, + SensorStateClass, +) +from homeassistant.const import PERCENTAGE +from homeassistant.core import HomeAssistant, callback +from homeassistant.helpers.device_registry import DeviceInfo +from homeassistant.helpers.entity_platform import AddEntitiesCallback +from homeassistant.helpers.update_coordinator import CoordinatorEntity + +from . import ModbusRTUMonitorConfigEntry +from .const import ( + DOMAIN, + MANUFACTURER, + MODEL, + REGISTER_MONITOR_COUNT, + REGISTER_MONITOR_COUNT_2, + REGISTER_MONITOR_COUNT_3, + REGISTER_START_ADDRESS, + REGISTER_START_ADDRESS_2, + REGISTER_START_ADDRESS_3, +) +from .coordinator import ModbusRTUMonitorCoordinator + + +async def async_setup_entry( + hass: HomeAssistant, + entry: ModbusRTUMonitorConfigEntry, + async_add_entities: AddEntitiesCallback, +) -> None: + """Set up sensor entities from a config entry.""" + coordinator = entry.runtime_data + + # Track which slaves have sensor entities + added_slaves: set[int] = set() + added_slaves_registers: set[int] = set() + added_slaves_registers2: set[int] = set() + added_slaves_registers3: set[int] = set() + + @callback + def _async_add_sensor_entities() -> None: + """Add sensor entities for newly discovered slaves.""" + if coordinator.data is None: + return + + entities = [] + + # Add humidity sensors for new slaves + current_slaves = set(coordinator.data.keys()) + new_slaves = current_slaves - added_slaves + + if new_slaves: + for slave_id in new_slaves: + entities.append( + ModbusRTUMonitorHumiditySensor(coordinator, slave_id) + ) + added_slaves.update(new_slaves) + + # Add register sensors for slaves with register data (range 1: 165-184) + current_slaves_registers = { + slave_id + for slave_id, slave_data in coordinator.data.items() + if slave_data.registers is not None + and any( + REGISTER_START_ADDRESS <= addr < REGISTER_START_ADDRESS + REGISTER_MONITOR_COUNT + for addr in slave_data.registers + ) + } + new_slaves_registers = current_slaves_registers - added_slaves_registers + + if new_slaves_registers: + for slave_id in new_slaves_registers: + # Create sensors for each monitored register (range 1) + for reg_offset in range(REGISTER_MONITOR_COUNT): + register_addr = REGISTER_START_ADDRESS + reg_offset + entities.append( + ModbusRTUMonitorRegisterSensor( + coordinator, slave_id, register_addr + ) + ) + added_slaves_registers.update(new_slaves_registers) + + # Add register sensors for slaves with register data (range 2: 210-217) + current_slaves_registers2 = { + slave_id + for slave_id, slave_data in coordinator.data.items() + if slave_data.registers is not None + and any( + REGISTER_START_ADDRESS_2 <= addr < REGISTER_START_ADDRESS_2 + REGISTER_MONITOR_COUNT_2 + for addr in slave_data.registers + ) + } + new_slaves_registers2 = current_slaves_registers2 - added_slaves_registers2 + + if new_slaves_registers2: + for slave_id in new_slaves_registers2: + # Create sensors for each monitored register (range 2) + for reg_offset in range(REGISTER_MONITOR_COUNT_2): + register_addr = REGISTER_START_ADDRESS_2 + reg_offset + entities.append( + ModbusRTUMonitorRegisterSensor( + coordinator, slave_id, register_addr + ) + ) + added_slaves_registers2.update(new_slaves_registers2) + + # Add register sensors for slaves with register data (range 3: 140-147) + current_slaves_registers3 = { + slave_id + for slave_id, slave_data in coordinator.data.items() + if slave_data.registers is not None + and any( + REGISTER_START_ADDRESS_3 <= addr < REGISTER_START_ADDRESS_3 + REGISTER_MONITOR_COUNT_3 + for addr in slave_data.registers + ) + } + new_slaves_registers3 = current_slaves_registers3 - added_slaves_registers3 + + if new_slaves_registers3: + for slave_id in new_slaves_registers3: + # Create sensors for each monitored register (range 3) + for reg_offset in range(REGISTER_MONITOR_COUNT_3): + register_addr = REGISTER_START_ADDRESS_3 + reg_offset + entities.append( + ModbusRTUMonitorRegisterSensor( + coordinator, slave_id, register_addr + ) + ) + added_slaves_registers3.update(new_slaves_registers3) + + if entities: + async_add_entities(entities) + + # Subscribe to coordinator updates + entry.async_on_unload( + coordinator.async_add_listener(_async_add_sensor_entities) + ) + + # Add entities for already discovered slaves + _async_add_sensor_entities() + + +class ModbusRTUMonitorHumiditySensor( + CoordinatorEntity[ModbusRTUMonitorCoordinator], + SensorEntity, +): + """Humidity sensor entity for Modbus RTU monitor.""" + + _attr_has_entity_name = True + _attr_translation_key = "humidity" + _attr_device_class = SensorDeviceClass.HUMIDITY + _attr_native_unit_of_measurement = PERCENTAGE + _attr_state_class = SensorStateClass.MEASUREMENT + + def __init__( + self, + coordinator: ModbusRTUMonitorCoordinator, + slave_id: int, + ) -> None: + """Initialize humidity sensor.""" + super().__init__(coordinator) + + self._slave_id = slave_id + self._attr_unique_id = ( + f"{coordinator.config_entry.entry_id}_{slave_id}_humidity" + ) + + # Same device as climate entity + self._attr_device_info = DeviceInfo( + identifiers={(DOMAIN, f"{coordinator.config_entry.entry_id}_{slave_id}")}, + name=f"Modbus Slave {slave_id}", + manufacturer=MANUFACTURER, + model=MODEL, + ) + + @property + def available(self) -> bool: + """Return if entity is available.""" + if self._slave_id not in self.coordinator.data: + return False + return self.coordinator.data[self._slave_id].available + + @property + def native_value(self) -> float | None: + """Return humidity value from passive monitoring.""" + if self._slave_id not in self.coordinator.data: + return None + return self.coordinator.data[self._slave_id].humidity + + +class ModbusRTUMonitorRegisterSensor( + CoordinatorEntity[ModbusRTUMonitorCoordinator], + SensorEntity, +): + """Register sensor entity for Modbus RTU monitor.""" + + _attr_has_entity_name = True + _attr_state_class = SensorStateClass.MEASUREMENT + + def __init__( + self, + coordinator: ModbusRTUMonitorCoordinator, + slave_id: int, + register_addr: int, + ) -> None: + """Initialize register sensor.""" + super().__init__(coordinator) + + self._slave_id = slave_id + self._register_addr = register_addr + + self._attr_unique_id = ( + f"{coordinator.config_entry.entry_id}_{slave_id}_register{register_addr}" + ) + self._attr_name = f"Register{register_addr}" + + # Same device as climate and humidity entities + self._attr_device_info = DeviceInfo( + identifiers={(DOMAIN, f"{coordinator.config_entry.entry_id}_{slave_id}")}, + name=f"Modbus Slave {slave_id}", + manufacturer=MANUFACTURER, + model=MODEL, + ) + + @property + def available(self) -> bool: + """Return if entity is available.""" + if self._slave_id not in self.coordinator.data: + return False + slave_data = self.coordinator.data[self._slave_id] + if not slave_data.available or slave_data.registers is None: + return False + return self._register_addr in slave_data.registers + + @property + def native_value(self) -> int | None: + """Return register value (signed int16).""" + if self._slave_id not in self.coordinator.data: + return None + slave_data = self.coordinator.data[self._slave_id] + if slave_data.registers is None: + return None + return slave_data.registers.get(self._register_addr) + + @property + def extra_state_attributes(self) -> dict[str, int]: + """Return additional state attributes.""" + return {"register_address": self._register_addr} diff --git a/strings.json b/strings.json new file mode 100644 index 0000000..e5dd609 --- /dev/null +++ b/strings.json @@ -0,0 +1,36 @@ +{ + "config": { + "step": { + "user": { + "title": "Set up Modbus RTU monitor", + "description": "Enter the host and port of your Modbus RTU over TCP gateway. Slaves will be discovered automatically.", + "data": { + "host": "Host", + "port": "Port" + } + } + }, + "error": { + "cannot_connect": "Failed to connect to the Modbus gateway", + "unknown": "An unexpected error occurred" + }, + "abort": { + "already_configured": "This Modbus gateway is already configured" + } + }, + "entity": { + "sensor": { + "humidity": { + "name": "Humidity" + }, + "register": { + "name": "Register {register_number}" + } + }, + "binary_sensor": { + "coil": { + "name": "Coil {coil_number}" + } + } + } +}