Working version
This commit is contained in:
+57
@@ -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
|
||||||
@@ -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}
|
||||||
+143
@@ -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."
|
||||||
|
)
|
||||||
@@ -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()
|
||||||
@@ -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
|
||||||
@@ -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
|
||||||
@@ -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("<H", crc))
|
||||||
|
|
||||||
|
return bytes(frame)
|
||||||
|
|
||||||
|
@staticmethod
|
||||||
|
def decode_frame(data: bytes) -> 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", data[-2:])[0]
|
||||||
|
calculated_crc = ModbusRTUDecoder.calculate_crc(data[:-2])
|
||||||
|
|
||||||
|
if received_crc != calculated_crc:
|
||||||
|
return None
|
||||||
|
|
||||||
|
decoded = ModbusRTUFrame(
|
||||||
|
slave_id=slave_id,
|
||||||
|
function_code=function_code,
|
||||||
|
data=frame_data,
|
||||||
|
is_request=False,
|
||||||
|
is_response=False,
|
||||||
|
)
|
||||||
|
|
||||||
|
# Detect request (function 1, Read Coils - with address and count)
|
||||||
|
if function_code == 1 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 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")
|
||||||
@@ -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"
|
||||||
|
}
|
||||||
@@ -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}
|
||||||
@@ -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}"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user