mirror of
https://github.com/home-assistant/core.git
synced 2025-11-14 21:40:16 +00:00
Refactor RFLink component (#17402)
* Start refactor of RFLink component * alias _id not added correctly Aliases for sensor not added correctly And some debug traces. * Update rflink.py * Cleanup, fix review comments * Call event handlers directly when processing initial event * Use hass.async_create_task when adding discovered device * Review comments * Review comments
This commit is contained in:
committed by
Martin Hjelmare
parent
0904ff45fe
commit
2ceb4d2d1e
@@ -6,7 +6,6 @@ https://home-assistant.io/components/rflink/
|
||||
"""
|
||||
import asyncio
|
||||
from collections import defaultdict
|
||||
import functools as ft
|
||||
import logging
|
||||
import async_timeout
|
||||
|
||||
@@ -14,7 +13,7 @@ import voluptuous as vol
|
||||
|
||||
from homeassistant.const import (
|
||||
ATTR_ENTITY_ID, CONF_COMMAND, CONF_HOST, CONF_PORT,
|
||||
EVENT_HOMEASSISTANT_STOP, STATE_UNKNOWN)
|
||||
EVENT_HOMEASSISTANT_STOP)
|
||||
from homeassistant.core import CoreState, callback
|
||||
from homeassistant.exceptions import HomeAssistantError
|
||||
import homeassistant.helpers.config_validation as cv
|
||||
@@ -68,6 +67,9 @@ DOMAIN = 'rflink'
|
||||
SERVICE_SEND_COMMAND = 'send_command'
|
||||
|
||||
SIGNAL_AVAILABILITY = 'rflink_device_available'
|
||||
SIGNAL_HANDLE_EVENT = 'rflink_handle_event_{}'
|
||||
|
||||
TMP_ENTITY = 'tmp.{}'
|
||||
|
||||
DEVICE_DEFAULTS_SCHEMA = vol.Schema({
|
||||
vol.Optional(CONF_FIRE_EVENT, default=False): cv.boolean,
|
||||
@@ -153,28 +155,38 @@ async def async_setup(hass, config):
|
||||
return
|
||||
|
||||
# Lookup entities who registered this device id as device id or alias
|
||||
event_id = event.get('id', None)
|
||||
event_id = event.get(EVENT_KEY_ID, None)
|
||||
|
||||
is_group_event = (event_type == EVENT_KEY_COMMAND and
|
||||
event[EVENT_KEY_COMMAND] in RFLINK_GROUP_COMMANDS)
|
||||
if is_group_event:
|
||||
entities = hass.data[DATA_ENTITY_GROUP_LOOKUP][event_type].get(
|
||||
entity_ids = hass.data[DATA_ENTITY_GROUP_LOOKUP][event_type].get(
|
||||
event_id, [])
|
||||
else:
|
||||
entities = hass.data[DATA_ENTITY_LOOKUP][event_type][event_id]
|
||||
entity_ids = hass.data[DATA_ENTITY_LOOKUP][event_type][event_id]
|
||||
|
||||
if entities:
|
||||
_LOGGER.debug('entity_ids: %s', entity_ids)
|
||||
if entity_ids:
|
||||
# Propagate event to every entity matching the device id
|
||||
for entity in entities:
|
||||
_LOGGER.debug('passing event to %s', entities)
|
||||
entity.handle_event(event)
|
||||
else:
|
||||
_LOGGER.debug('device_id not known, adding new device')
|
||||
|
||||
for entity in entity_ids:
|
||||
_LOGGER.debug('passing event to %s', entity)
|
||||
async_dispatcher_send(hass,
|
||||
SIGNAL_HANDLE_EVENT.format(entity),
|
||||
event)
|
||||
elif not is_group_event:
|
||||
# If device is not yet known, register with platform (if loaded)
|
||||
if event_type in hass.data[DATA_DEVICE_REGISTER]:
|
||||
hass.async_run_job(
|
||||
hass.data[DATA_DEVICE_REGISTER][event_type], event)
|
||||
_LOGGER.debug('device_id not known, adding new device')
|
||||
# Add bogus event_id first to avoid race if we get another
|
||||
# event before the device is created
|
||||
# Any additional events recevied before the device has been
|
||||
# created will thus be ignored.
|
||||
hass.data[DATA_ENTITY_LOOKUP][event_type][
|
||||
event_id].append(TMP_ENTITY.format(event_id))
|
||||
hass.async_create_task(
|
||||
hass.data[DATA_DEVICE_REGISTER][event_type](event))
|
||||
else:
|
||||
_LOGGER.debug('device_id not known and automatic add disabled')
|
||||
|
||||
# When connecting to tcp host instead of serial port (optional)
|
||||
host = config[DOMAIN].get(CONF_HOST)
|
||||
@@ -192,7 +204,7 @@ async def async_setup(hass, config):
|
||||
# If HA is not stopping, initiate new connection
|
||||
if hass.state != CoreState.stopping:
|
||||
_LOGGER.warning('disconnected from Rflink, reconnecting')
|
||||
hass.async_add_job(connect)
|
||||
hass.async_create_task(connect())
|
||||
|
||||
async def connect():
|
||||
"""Set up connection and hook it into HA for reconnect/shutdown."""
|
||||
@@ -242,7 +254,7 @@ async def async_setup(hass, config):
|
||||
|
||||
_LOGGER.info('Connected to Rflink')
|
||||
|
||||
hass.async_add_job(connect)
|
||||
hass.async_create_task(connect())
|
||||
return True
|
||||
|
||||
|
||||
@@ -253,26 +265,31 @@ class RflinkDevice(Entity):
|
||||
"""
|
||||
|
||||
platform = None
|
||||
_state = STATE_UNKNOWN
|
||||
_state = None
|
||||
_available = True
|
||||
|
||||
def __init__(self, device_id, hass, name=None, aliases=None, group=True,
|
||||
group_aliases=None, nogroup_aliases=None, fire_event=False,
|
||||
def __init__(self, device_id, initial_event=None, name=None, aliases=None,
|
||||
group=True, group_aliases=None, nogroup_aliases=None,
|
||||
fire_event=False,
|
||||
signal_repetitions=DEFAULT_SIGNAL_REPETITIONS):
|
||||
"""Initialize the device."""
|
||||
self.hass = hass
|
||||
|
||||
# Rflink specific attributes for every component type
|
||||
self._initial_event = initial_event
|
||||
self._device_id = device_id
|
||||
if name:
|
||||
self._name = name
|
||||
else:
|
||||
self._name = device_id
|
||||
|
||||
self._aliases = aliases
|
||||
self._group = group
|
||||
self._group_aliases = group_aliases
|
||||
self._nogroup_aliases = nogroup_aliases
|
||||
self._should_fire_event = fire_event
|
||||
self._signal_repetitions = signal_repetitions
|
||||
|
||||
def handle_event(self, event):
|
||||
@callback
|
||||
def handle_event_callback(self, event):
|
||||
"""Handle incoming event for device type."""
|
||||
# Call platform specific event handler
|
||||
self._handle_event(event)
|
||||
@@ -283,7 +300,7 @@ class RflinkDevice(Entity):
|
||||
# Put command onto bus for user to subscribe to
|
||||
if self._should_fire_event and identify_event_type(
|
||||
event) == EVENT_KEY_COMMAND:
|
||||
self.hass.bus.fire(EVENT_BUTTON_PRESSED, {
|
||||
self.hass.bus.async_fire(EVENT_BUTTON_PRESSED, {
|
||||
ATTR_ENTITY_ID: self.entity_id,
|
||||
ATTR_STATE: event[EVENT_KEY_COMMAND],
|
||||
})
|
||||
@@ -314,7 +331,7 @@ class RflinkDevice(Entity):
|
||||
@property
|
||||
def assumed_state(self):
|
||||
"""Assume device state until first device event sets state."""
|
||||
return self._state is STATE_UNKNOWN
|
||||
return self._state is None
|
||||
|
||||
@property
|
||||
def available(self):
|
||||
@@ -322,15 +339,52 @@ class RflinkDevice(Entity):
|
||||
return self._available
|
||||
|
||||
@callback
|
||||
def set_availability(self, availability):
|
||||
def _availability_callback(self, availability):
|
||||
"""Update availability state."""
|
||||
self._available = availability
|
||||
self.async_schedule_update_ha_state()
|
||||
|
||||
async def async_added_to_hass(self):
|
||||
"""Register update callback."""
|
||||
# Remove temporary bogus entity_id if added
|
||||
tmp_entity = TMP_ENTITY.format(self._device_id)
|
||||
if tmp_entity in self.hass.data[DATA_ENTITY_LOOKUP][
|
||||
EVENT_KEY_SENSOR][self._device_id]:
|
||||
self.hass.data[DATA_ENTITY_LOOKUP][
|
||||
EVENT_KEY_SENSOR][self._device_id].remove(tmp_entity)
|
||||
|
||||
# Register id and aliases
|
||||
self.hass.data[DATA_ENTITY_LOOKUP][
|
||||
EVENT_KEY_COMMAND][self._device_id].append(self.entity_id)
|
||||
if self._group:
|
||||
self.hass.data[DATA_ENTITY_GROUP_LOOKUP][
|
||||
EVENT_KEY_COMMAND][self._device_id].append(self.entity_id)
|
||||
# aliases respond to both normal and group commands (allon/alloff)
|
||||
if self._aliases:
|
||||
for _id in self._aliases:
|
||||
self.hass.data[DATA_ENTITY_LOOKUP][
|
||||
EVENT_KEY_COMMAND][_id].append(self.entity_id)
|
||||
self.hass.data[DATA_ENTITY_GROUP_LOOKUP][
|
||||
EVENT_KEY_COMMAND][_id].append(self.entity_id)
|
||||
# group_aliases only respond to group commands (allon/alloff)
|
||||
if self._group_aliases:
|
||||
for _id in self._group_aliases:
|
||||
self.hass.data[DATA_ENTITY_GROUP_LOOKUP][
|
||||
EVENT_KEY_COMMAND][_id].append(self.entity_id)
|
||||
# nogroup_aliases only respond to normal commands
|
||||
if self._nogroup_aliases:
|
||||
for _id in self._nogroup_aliases:
|
||||
self.hass.data[DATA_ENTITY_LOOKUP][
|
||||
EVENT_KEY_COMMAND][_id].append(self.entity_id)
|
||||
async_dispatcher_connect(self.hass, SIGNAL_AVAILABILITY,
|
||||
self.set_availability)
|
||||
self._availability_callback)
|
||||
async_dispatcher_connect(self.hass,
|
||||
SIGNAL_HANDLE_EVENT.format(self.entity_id),
|
||||
self.handle_event_callback)
|
||||
|
||||
# Process the initial event now that the entity is created
|
||||
if self._initial_event:
|
||||
self.handle_event_callback(self._initial_event)
|
||||
|
||||
|
||||
class RflinkCommand(RflinkDevice):
|
||||
@@ -388,7 +442,7 @@ class RflinkCommand(RflinkDevice):
|
||||
cmd = 'on'
|
||||
# if the state is unknown or false, it gets set as true
|
||||
# if the state is true, it gets set as false
|
||||
self._state = self._state in [STATE_UNKNOWN, False]
|
||||
self._state = self._state in [None, False]
|
||||
|
||||
# Cover options for RFlink
|
||||
elif command == 'close_cover':
|
||||
@@ -439,8 +493,8 @@ class RflinkCommand(RflinkDevice):
|
||||
# Rflink protocol/transport handles asynchronous writing of buffer
|
||||
# to serial/tcp device. Does not wait for command send
|
||||
# confirmation.
|
||||
self.hass.async_add_job(ft.partial(
|
||||
self._protocol.send_command, self._device_id, cmd))
|
||||
self.hass.async_create_task(self._protocol.send_command(
|
||||
self._device_id, cmd))
|
||||
|
||||
if repetitions > 1:
|
||||
self._repetition_task = self.hass.async_create_task(
|
||||
|
||||
Reference in New Issue
Block a user