Convert mqtt platforms to async (#6145)
* Convert mqtt platforms to async * fix lint * add more platforms * convert mqtt_eventstream * fix lint / add mqtt_room * fix lint * fix test part 1 * fix test part 2 * fix out of memory bug * address comments
This commit is contained in:
parent
65ed85c6eb
commit
b0d3bbed79
15 changed files with 509 additions and 291 deletions
|
@ -4,12 +4,14 @@ Support for MQTT room presence detection.
|
|||
For more details about this platform, please refer to the documentation at
|
||||
https://home-assistant.io/components/sensor.mqtt_room/
|
||||
"""
|
||||
import asyncio
|
||||
import logging
|
||||
import json
|
||||
from datetime import timedelta
|
||||
|
||||
import voluptuous as vol
|
||||
|
||||
from homeassistant.core import callback
|
||||
import homeassistant.components.mqtt as mqtt
|
||||
from homeassistant.components.sensor import PLATFORM_SCHEMA
|
||||
from homeassistant.const import (
|
||||
|
@ -54,11 +56,10 @@ MQTT_PAYLOAD = vol.Schema(vol.All(json.loads, vol.Schema({
|
|||
}, extra=vol.ALLOW_EXTRA)))
|
||||
|
||||
|
||||
# pylint: disable=unused-argument
|
||||
def setup_platform(hass, config, add_devices, discovery_info=None):
|
||||
@asyncio.coroutine
|
||||
def async_setup_platform(hass, config, async_add_devices, discovery_info=None):
|
||||
"""Setup MQTT Sensor."""
|
||||
add_devices([MQTTRoomSensor(
|
||||
hass,
|
||||
yield from async_add_devices([MQTTRoomSensor(
|
||||
config.get(CONF_NAME),
|
||||
config.get(CONF_STATE_TOPIC),
|
||||
config.get(CONF_DEVICE_ID),
|
||||
|
@ -70,11 +71,9 @@ def setup_platform(hass, config, add_devices, discovery_info=None):
|
|||
class MQTTRoomSensor(Entity):
|
||||
"""Representation of a room sensor that is updated via MQTT."""
|
||||
|
||||
def __init__(self, hass, name, state_topic, device_id, timeout,
|
||||
consider_home):
|
||||
def __init__(self, name, state_topic, device_id, timeout, consider_home):
|
||||
"""Initialize the sensor."""
|
||||
self._state = STATE_AWAY
|
||||
self._hass = hass
|
||||
self._name = name
|
||||
self._state_topic = '{}{}'.format(state_topic, '/+')
|
||||
self._device_id = slugify(device_id).upper()
|
||||
|
@ -85,13 +84,19 @@ class MQTTRoomSensor(Entity):
|
|||
self._distance = None
|
||||
self._updated = None
|
||||
|
||||
def async_added_to_hass(self):
|
||||
"""Subscribe mqtt events.
|
||||
|
||||
This method must be run in the event loop and returns a coroutine.
|
||||
"""
|
||||
@callback
|
||||
def update_state(device_id, room, distance):
|
||||
"""Update the sensor state."""
|
||||
self._state = room
|
||||
self._distance = distance
|
||||
self._updated = dt.utcnow()
|
||||
|
||||
self.update_ha_state()
|
||||
self.hass.async_add_job(self.async_update_ha_state())
|
||||
|
||||
def message_received(topic, payload, qos):
|
||||
"""A new MQTT message has been received."""
|
||||
|
@ -117,7 +122,8 @@ class MQTTRoomSensor(Entity):
|
|||
or timediff.seconds >= self._timeout:
|
||||
update_state(**device)
|
||||
|
||||
mqtt.subscribe(hass, self._state_topic, message_received, 1)
|
||||
return mqtt.async_subscribe(
|
||||
self.hass, self._state_topic, message_received, 1)
|
||||
|
||||
@property
|
||||
def name(self):
|
||||
|
|
Loading…
Add table
Add a link
Reference in a new issue