mirror of
https://github.com/andvikt/mega_hacs.git
synced 2025-12-11 00:54:28 +05:00
add http support
This commit is contained in:
@@ -4,18 +4,22 @@ import logging
|
||||
from functools import partial
|
||||
|
||||
import voluptuous as vol
|
||||
|
||||
from homeassistant.const import (
|
||||
CONF_SCAN_INTERVAL, CONF_ID, CONF_NAME, CONF_DOMAIN,
|
||||
CONF_UNIT_OF_MEASUREMENT,
|
||||
CONF_UNIT_OF_MEASUREMENT, CONF_HOST
|
||||
)
|
||||
from homeassistant.core import HomeAssistant, ServiceCall
|
||||
from homeassistant.helpers.service import bind_hass
|
||||
from homeassistant.helpers.template import Template
|
||||
from homeassistant.helpers import config_validation as cv
|
||||
from homeassistant.components import mqtt
|
||||
from homeassistant.config_entries import ConfigEntry
|
||||
from .const import DOMAIN, CONF_INVERT, CONF_RELOAD, PLATFORMS, CONF_PORTS, CONF_CUSTOM, CONF_SKIP, CONF_PORT_TO_SCAN
|
||||
from .const import DOMAIN, CONF_INVERT, CONF_RELOAD, PLATFORMS, CONF_PORTS, CONF_CUSTOM, CONF_SKIP, CONF_PORT_TO_SCAN, \
|
||||
CONF_MQTT_INPUTS, CONF_HTTP, CONF_RESPONSE_TEMPLATE, CONF_ACTION, CONF_GET_VALUE
|
||||
from .hub import MegaD
|
||||
from .config_flow import ConfigFlow
|
||||
|
||||
from .http import MegaView
|
||||
|
||||
_LOGGER = logging.getLogger(__name__)
|
||||
|
||||
@@ -34,6 +38,12 @@ CONFIG_SCHEMA = vol.Schema(
|
||||
vol.Any(str, {
|
||||
vol.Required(str): str
|
||||
}),
|
||||
vol.Optional(
|
||||
CONF_RESPONSE_TEMPLATE,
|
||||
description='шаблон ответа когда на этот порт приходит'
|
||||
'сообщение из меги '): cv.template,
|
||||
vol.Optional(CONF_ACTION): cv.script_action,
|
||||
vol.Optional(CONF_GET_VALUE, default=True): bool,
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -51,7 +61,8 @@ _subs = {}
|
||||
async def async_setup(hass: HomeAssistant, config: dict):
|
||||
"""YAML-конфигурация содержит только кастомизации портов"""
|
||||
hass.data[DOMAIN] = {CONF_CUSTOM: config.get(DOMAIN, {})}
|
||||
|
||||
hass.data[DOMAIN][CONF_HTTP] = view = MegaView(cfg=config.get(DOMAIN, {}))
|
||||
hass.http.register_view(view)
|
||||
hass.services.async_register(
|
||||
DOMAIN, 'save', partial(_save_service, hass), schema=vol.Schema({
|
||||
vol.Optional('mega_id'): str
|
||||
@@ -70,6 +81,7 @@ async def async_setup(hass: HomeAssistant, config: dict):
|
||||
vol.Optional('mega_id'): str,
|
||||
})
|
||||
)
|
||||
|
||||
return True
|
||||
|
||||
|
||||
@@ -78,9 +90,17 @@ async def get_hub(hass, entry):
|
||||
data = dict(entry.data)
|
||||
data.update(entry.options or {})
|
||||
data.update(id=id)
|
||||
_mqtt = hass.data.get(mqtt.DOMAIN)
|
||||
if _mqtt is None:
|
||||
raise Exception('mqtt not configured, please configure mqtt first')
|
||||
use_mqtt = data.get(CONF_MQTT_INPUTS, True)
|
||||
|
||||
_mqtt = hass.data.get(mqtt.DOMAIN) if use_mqtt else None
|
||||
if _mqtt is None and use_mqtt:
|
||||
for x in range(5):
|
||||
await asyncio.sleep(5)
|
||||
_mqtt = hass.data.get(mqtt.DOMAIN)
|
||||
if _mqtt is not None:
|
||||
break
|
||||
if _mqtt is None:
|
||||
raise Exception('mqtt not configured, please configure mqtt first')
|
||||
hub = MegaD(hass, **data, mqtt=_mqtt, lg=_LOGGER, loop=asyncio.get_event_loop())
|
||||
hub.mqtt_id = await hub.get_mqtt_id()
|
||||
return hub
|
||||
@@ -89,7 +109,8 @@ async def get_hub(hass, entry):
|
||||
async def _add_mega(hass: HomeAssistant, entry: ConfigEntry):
|
||||
id = entry.data.get('id', entry.entry_id)
|
||||
hub = await get_hub(hass, entry)
|
||||
hass.data[DOMAIN][id] = hub
|
||||
hass.data[DOMAIN][id] = hass.data[DOMAIN]['__def'] = hub
|
||||
hass.data[DOMAIN][entry.data.get(CONF_HOST)] = hub
|
||||
if not await hub.authenticate():
|
||||
raise Exception("not authentificated")
|
||||
mid = await hub.get_mqtt_id()
|
||||
|
||||
@@ -73,20 +73,13 @@ class MegaBinarySensor(BinarySensorEntity, MegaPushEntity):
|
||||
|
||||
@property
|
||||
def is_on(self) -> bool:
|
||||
if self._is_on is not None:
|
||||
return self._is_on
|
||||
return self._state == 'ON'
|
||||
val = self.mega.values.get(self.port, {}).get("value") \
|
||||
or self.mega.values.get(self.port, {}).get('m')
|
||||
if val is None and self._state is not None:
|
||||
return self._state == 'ON'
|
||||
elif val is not None:
|
||||
return val == 'ON' or val == 1
|
||||
|
||||
def _update(self, payload: dict):
|
||||
data = {CONF_ENTITY_ID: self.entity_id}
|
||||
payload = payload.copy()
|
||||
payload.pop(CONF_PORT)
|
||||
data.update(payload)
|
||||
if not self.is_first_update:
|
||||
self.hass.bus.async_fire(
|
||||
EVENT_BINARY_SENSOR,
|
||||
data,
|
||||
)
|
||||
val = payload.get("value")
|
||||
self._is_on = val == 'ON'
|
||||
self._attrs = data
|
||||
self.mega.values[self.port] = payload
|
||||
|
||||
|
||||
@@ -9,7 +9,8 @@ from homeassistant.components import mqtt
|
||||
from homeassistant.config_entries import ConfigEntry
|
||||
from homeassistant.const import CONF_HOST, CONF_ID, CONF_PASSWORD, CONF_SCAN_INTERVAL
|
||||
from homeassistant.core import callback, HomeAssistant
|
||||
from .const import DOMAIN, CONF_PORT_TO_SCAN, CONF_RELOAD, PLATFORMS # pylint:disable=unused-import
|
||||
from .const import DOMAIN, CONF_PORT_TO_SCAN, CONF_RELOAD, PLATFORMS, CONF_MQTT_INPUTS, \
|
||||
CONF_NPORTS, CONF_UPDATE_ALL # pylint:disable=unused-import
|
||||
from .hub import MegaD
|
||||
from . import exceptions
|
||||
|
||||
@@ -22,6 +23,9 @@ STEP_USER_DATA_SCHEMA = vol.Schema(
|
||||
vol.Required(CONF_PASSWORD, default="sec"): str,
|
||||
vol.Optional(CONF_SCAN_INTERVAL, default=0): int,
|
||||
vol.Optional(CONF_PORT_TO_SCAN, default=0): int,
|
||||
vol.Optional(CONF_MQTT_INPUTS, default=True): bool,
|
||||
vol.Optional(CONF_NPORTS, default=37): int,
|
||||
vol.Optional(CONF_UPDATE_ALL, default=True): bool,
|
||||
},
|
||||
)
|
||||
|
||||
@@ -52,7 +56,7 @@ async def validate_input(hass: core.HomeAssistant, data):
|
||||
class ConfigFlow(config_entries.ConfigFlow, domain=DOMAIN):
|
||||
"""Handle a config flow for mega."""
|
||||
|
||||
VERSION = 3
|
||||
VERSION = 4
|
||||
CONNECTION_CLASS = config_entries.CONN_CLASS_ASSUMED
|
||||
|
||||
async def async_step_user(self, user_input=None):
|
||||
@@ -67,7 +71,7 @@ class ConfigFlow(config_entries.ConfigFlow, domain=DOMAIN):
|
||||
try:
|
||||
hub = await validate_input(self.hass, user_input)
|
||||
await hub.start()
|
||||
config = await hub.get_config()
|
||||
config = await hub.get_config(nports=user_input.get(CONF_NPORTS, 37))
|
||||
await hub.stop()
|
||||
hub.lg.debug(f'config loaded: %s', config)
|
||||
config.update(user_input)
|
||||
@@ -110,7 +114,7 @@ class OptionsFlowHandler(config_entries.OptionsFlow):
|
||||
hub = await get_hub(self.hass, self.config_entry.data)
|
||||
if reload:
|
||||
await hub.start()
|
||||
new = await hub.get_config()
|
||||
new = await hub.get_config(nports=user_input.get(CONF_NPORTS, 37))
|
||||
await hub.stop()
|
||||
|
||||
_LOGGER.debug(f'new config: %s', new)
|
||||
@@ -128,7 +132,10 @@ class OptionsFlowHandler(config_entries.OptionsFlow):
|
||||
data_schema=vol.Schema({
|
||||
vol.Optional(CONF_SCAN_INTERVAL, default=e.get(CONF_SCAN_INTERVAL, 0)): int,
|
||||
vol.Optional(CONF_PORT_TO_SCAN, default=e.get(CONF_PORT_TO_SCAN, 0)): int,
|
||||
vol.Optional(CONF_MQTT_INPUTS, default=e.get(CONF_MQTT_INPUTS, True)): bool,
|
||||
vol.Optional(CONF_NPORTS, default=e.get(CONF_NPORTS, 37)): int,
|
||||
vol.Optional(CONF_RELOAD, default=False): bool,
|
||||
# vol.Optional(CONF_UPDATE_ALL, default=e.get(CONF_UPDATE_ALL, True)): bool,
|
||||
# vol.Optional(CONF_INVERT, default=''): str,
|
||||
}),
|
||||
)
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
"""Constants for the mega integration."""
|
||||
import re
|
||||
|
||||
DOMAIN = "mega"
|
||||
CONF_MEGA_ID = "mega_id"
|
||||
@@ -14,11 +15,19 @@ CONF_RELOAD = 'reload'
|
||||
CONF_INVERT = 'invert'
|
||||
CONF_PORTS = 'ports'
|
||||
CONF_CUSTOM = '__custom'
|
||||
CONF_HTTP = '__http'
|
||||
CONF_SKIP = 'skip'
|
||||
CONF_MQTT_INPUTS = 'mqtt_inputs'
|
||||
CONF_NPORTS = 'nports'
|
||||
CONF_RESPONSE_TEMPLATE = 'response_template'
|
||||
CONF_ACTION = 'action'
|
||||
CONF_UPDATE_ALL = 'update_all'
|
||||
CONF_GET_VALUE = 'get_value'
|
||||
PLATFORMS = [
|
||||
"light",
|
||||
"switch",
|
||||
"binary_sensor",
|
||||
"sensor",
|
||||
]
|
||||
EVENT_BINARY_SENSOR = f'{DOMAIN}.sensor'
|
||||
EVENT_BINARY_SENSOR = f'{DOMAIN}.sensor'
|
||||
PATT_SPLIT = re.compile('[;/]')
|
||||
@@ -29,6 +29,7 @@ class BaseMegaEntity(CoordinatorEntity, RestoreEntity):
|
||||
self.port = port
|
||||
self.config_entry = config_entry
|
||||
self.mega = mega
|
||||
mega.entities.append(self)
|
||||
self._mega_id = mega.id
|
||||
self._lg = None
|
||||
self._unique_id = unique_id or f"mega_{mega.id}_{port}" + \
|
||||
@@ -90,6 +91,10 @@ class BaseMegaEntity(CoordinatorEntity, RestoreEntity):
|
||||
await super().async_added_to_hass()
|
||||
self._state = await self.async_get_last_state()
|
||||
|
||||
async def get_state(self):
|
||||
if self.mega.mqtt is None:
|
||||
self.async_write_ha_state()
|
||||
|
||||
|
||||
class MegaPushEntity(BaseMegaEntity):
|
||||
|
||||
@@ -114,7 +119,8 @@ class MegaPushEntity(BaseMegaEntity):
|
||||
|
||||
async def async_added_to_hass(self) -> None:
|
||||
await super().async_added_to_hass()
|
||||
asyncio.create_task(self.mega.get_port(self.port))
|
||||
if self.mega.mqtt is not None:
|
||||
asyncio.create_task(self.mega.get_port(self.port))
|
||||
|
||||
|
||||
class MegaOutPort(MegaPushEntity):
|
||||
@@ -137,16 +143,26 @@ class MegaOutPort(MegaPushEntity):
|
||||
|
||||
@property
|
||||
def brightness(self):
|
||||
if self._brightness is not None:
|
||||
return self._brightness
|
||||
if self._state:
|
||||
val = self.mega.values.get(self.port, {}).get("value")
|
||||
if val is None and self._state is not None:
|
||||
return self._state.attributes.get("brightness")
|
||||
elif val is not None:
|
||||
try:
|
||||
val = int(val)
|
||||
return val
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
@property
|
||||
def is_on(self) -> bool:
|
||||
if self._is_on is not None:
|
||||
return self._is_on
|
||||
return self._state == 'ON'
|
||||
val = self.mega.values.get(self.port, {}).get("value")
|
||||
if val is None and self._state is not None:
|
||||
return self._state == 'ON'
|
||||
elif val is not None:
|
||||
if not self.invert:
|
||||
return val == 'ON' or str(val) == '1' or (safe_int(val) is not None and safe_int(val) > 0)
|
||||
else:
|
||||
return val == 'OFF' or str(val) == '0' or (safe_int(val) is not None and safe_int(val) == 0)
|
||||
|
||||
async def async_turn_on(self, brightness=None, **kwargs) -> None:
|
||||
brightness = brightness or self.brightness or 255
|
||||
@@ -157,31 +173,23 @@ class MegaOutPort(MegaPushEntity):
|
||||
cmd = brightness
|
||||
else:
|
||||
cmd = 1 if not self.invert else 0
|
||||
if await self.mega.send_command(self.port, f"{self.port}:{cmd}"):
|
||||
self._is_on = True
|
||||
self._brightness = brightness
|
||||
await self.async_update_ha_state()
|
||||
await self.mega.send_command(self.port, f"{self.port}:{cmd}")
|
||||
self.mega.values[self.port] = {'value': cmd}
|
||||
await self.get_state()
|
||||
|
||||
|
||||
async def async_turn_off(self, **kwargs) -> None:
|
||||
|
||||
cmd = "0" if not self.invert else "1"
|
||||
|
||||
if await self.mega.send_command(self.port, f"{self.port}:{cmd}"):
|
||||
self._is_on = False
|
||||
await self.async_update_ha_state()
|
||||
await self.mega.send_command(self.port, f"{self.port}:{cmd}")
|
||||
self.mega.values[self.port] = {'value': cmd}
|
||||
await self.get_state()
|
||||
|
||||
def _update(self, payload: dict):
|
||||
val = payload.get("value")
|
||||
try:
|
||||
val = int(val)
|
||||
except Exception:
|
||||
pass
|
||||
if isinstance(val, int):
|
||||
self._is_on = val
|
||||
if val > 0:
|
||||
self._brightness = val
|
||||
else:
|
||||
if not self.invert:
|
||||
self._is_on = val == 'ON'
|
||||
else:
|
||||
self._is_on = val == 'OFF'
|
||||
def safe_int(v):
|
||||
if v in ['ON', 'OFF']:
|
||||
return None
|
||||
try:
|
||||
return int(v)
|
||||
except ValueError:
|
||||
return None
|
||||
95
custom_components/mega/http.py
Normal file
95
custom_components/mega/http.py
Normal file
@@ -0,0 +1,95 @@
|
||||
import asyncio
|
||||
import logging
|
||||
|
||||
import typing
|
||||
from collections import defaultdict
|
||||
|
||||
from aiohttp.web_request import Request
|
||||
from aiohttp.web_response import Response
|
||||
|
||||
from homeassistant.helpers.template import Template
|
||||
from .const import EVENT_BINARY_SENSOR, CONF_HTTP, DOMAIN, CONF_CUSTOM, CONF_RESPONSE_TEMPLATE
|
||||
from homeassistant.components.http import HomeAssistantView
|
||||
from homeassistant.core import callback, HomeAssistant
|
||||
from . import hub
|
||||
|
||||
_LOGGER = logging.getLogger(__name__).getChild('http')
|
||||
|
||||
|
||||
class MegaView(HomeAssistantView):
|
||||
"""Handle Yandex Smart Home unauthorized requests."""
|
||||
|
||||
url = '/mega'
|
||||
name = 'mega'
|
||||
requires_auth = False
|
||||
|
||||
def __init__(self, cfg: dict):
|
||||
self._try = 0
|
||||
self.allowed_hosts = {'::1'}
|
||||
self.callbacks: typing.DefaultDict[int, typing.List[typing.Callable[[dict], typing.Coroutine]]] \
|
||||
= defaultdict(list)
|
||||
self.templates: typing.Dict[str, typing.Dict[str, Template]] = {
|
||||
mid: {
|
||||
pt: cfg[mid][pt][CONF_RESPONSE_TEMPLATE]
|
||||
for pt in cfg[mid]
|
||||
if CONF_RESPONSE_TEMPLATE in cfg[mid][pt]
|
||||
} for mid in cfg
|
||||
}
|
||||
_LOGGER.debug('templates: %s', self.templates)
|
||||
|
||||
async def get(self, request: Request) -> Response:
|
||||
auth = False
|
||||
for x in self.allowed_hosts:
|
||||
if request.remote.startswith(x):
|
||||
auth = True
|
||||
break
|
||||
if not auth:
|
||||
_LOGGER.warning(f'unauthorised attempt to connect from {request.remote}')
|
||||
return Response(status=401)
|
||||
|
||||
hass: HomeAssistant = request.app['hass']
|
||||
hub: 'hub.MegaD' = hass.data.get(DOMAIN).get(request.remote) # TODO: проверить какой remote
|
||||
if hub is None and request.remote == '::1':
|
||||
hub = hass.data.get(DOMAIN).get('__def')
|
||||
if hub is None:
|
||||
return Response(status=400)
|
||||
data = dict(request.query)
|
||||
hass.bus.async_fire(
|
||||
EVENT_BINARY_SENSOR,
|
||||
data,
|
||||
)
|
||||
_LOGGER.debug(f"Request: %s from '%s'", data, request.remote)
|
||||
make_ints(data)
|
||||
port = data.get('pt')
|
||||
data = data.copy()
|
||||
ret = 'd'
|
||||
if port is not None:
|
||||
for cb in self.callbacks[port]:
|
||||
cb(data)
|
||||
template: Template = self.templates.get(hub.id, {}).get(port)
|
||||
if hub.update_all:
|
||||
asyncio.create_task(self.later_update(hub))
|
||||
if template is not None:
|
||||
template.hass = hass
|
||||
ret = template.async_render(data)
|
||||
_LOGGER.debug('response %s', ret)
|
||||
ret = Response(body=ret or 'd', content_type='text/plain', headers={})
|
||||
ret.headers.clear()
|
||||
return ret
|
||||
|
||||
async def later_update(self, hub):
|
||||
_LOGGER.debug('force update')
|
||||
await asyncio.sleep(1)
|
||||
await hub.updater.async_refresh()
|
||||
|
||||
|
||||
def make_ints(d: dict):
|
||||
for x in d:
|
||||
try:
|
||||
d[x] = float(d[x])
|
||||
except ValueError:
|
||||
pass
|
||||
if 'm' not in d:
|
||||
d['m'] = 0
|
||||
if 'click' not in d:
|
||||
d['click'] = 0
|
||||
@@ -14,8 +14,9 @@ from homeassistant.const import DEVICE_CLASS_TEMPERATURE, DEVICE_CLASS_HUMIDITY
|
||||
from homeassistant.core import HomeAssistant
|
||||
from homeassistant.helpers.entity import Entity
|
||||
from homeassistant.helpers.update_coordinator import DataUpdateCoordinator
|
||||
from .const import TEMP, HUM
|
||||
from .exceptions import CannotConnect
|
||||
from .const import TEMP, HUM, PATT_SPLIT, DOMAIN, CONF_HTTP
|
||||
from .exceptions import CannotConnect, MqttNotConfigured
|
||||
from .http import MegaView
|
||||
|
||||
TEMP_PATT = re.compile(r'temp:([01234567890\.]+)')
|
||||
HUM_PATT = re.compile(r'hum:([01234567890\.]+)')
|
||||
@@ -32,6 +33,10 @@ CLASSES = {
|
||||
HUM: DEVICE_CLASS_HUMIDITY
|
||||
}
|
||||
|
||||
class NoPort(Exception):
|
||||
pass
|
||||
|
||||
|
||||
class MegaD:
|
||||
"""MegaD Hub"""
|
||||
|
||||
@@ -44,13 +49,24 @@ class MegaD:
|
||||
mqtt: mqtt.MQTT,
|
||||
lg: logging.Logger,
|
||||
id: str,
|
||||
mqtt_inputs: bool = True,
|
||||
mqtt_id: str = None,
|
||||
scan_interval=60,
|
||||
port_to_scan=0,
|
||||
nports=38,
|
||||
inverted: typing.List[int] = None,
|
||||
update_all=True,
|
||||
**kwargs,
|
||||
):
|
||||
"""Initialize."""
|
||||
if mqtt_inputs is None or mqtt_inputs == 'None' or mqtt_inputs is False:
|
||||
self.http = hass.data[DOMAIN][CONF_HTTP]
|
||||
self.http.allowed_hosts |= {host}
|
||||
else:
|
||||
self.http = None
|
||||
self.update_all = update_all if update_all is not None else True
|
||||
self.nports = nports
|
||||
self.mqtt_inputs = mqtt_inputs
|
||||
self.loop: asyncio.AbstractEventLoop = None
|
||||
self.hass = hass
|
||||
self.host = host
|
||||
@@ -58,6 +74,7 @@ class MegaD:
|
||||
self.mqtt = mqtt
|
||||
self.id = id
|
||||
self.lck = asyncio.Lock()
|
||||
self._http_lck = asyncio.Lock()
|
||||
self._notif_lck = asyncio.Lock()
|
||||
self.cnd = asyncio.Condition()
|
||||
self.online = True
|
||||
@@ -89,14 +106,16 @@ class MegaD:
|
||||
|
||||
async def start(self):
|
||||
self.loop = asyncio.get_event_loop()
|
||||
self.subs = await self.mqtt.async_subscribe(
|
||||
topic=f"{self.mqtt_id}/+",
|
||||
msg_callback=self._process_msg,
|
||||
qos=0,
|
||||
)
|
||||
if self.mqtt is not None:
|
||||
self.subs = await self.mqtt.async_subscribe(
|
||||
topic=f"{self.mqtt_id}/+",
|
||||
msg_callback=self._process_msg,
|
||||
qos=0,
|
||||
)
|
||||
|
||||
async def stop(self):
|
||||
self.subs()
|
||||
if self.subs is not None:
|
||||
self.subs()
|
||||
for x in self._callbacks.values():
|
||||
x.clear()
|
||||
|
||||
@@ -136,6 +155,9 @@ class MegaD:
|
||||
offline
|
||||
"""
|
||||
self.lg.debug('poll')
|
||||
if self.mqtt is None:
|
||||
await self.get_all_ports()
|
||||
return
|
||||
if len(self.sensors) > 0:
|
||||
await self.get_sensors()
|
||||
else:
|
||||
@@ -154,25 +176,51 @@ class MegaD:
|
||||
return _id or 'megad/' + self.host.split('.')[-1]
|
||||
|
||||
async def send_command(self, port=None, cmd=None):
|
||||
if port:
|
||||
url = f"http://{self.host}/{self.sec}/?pt={port}&cmd={cmd}"
|
||||
else:
|
||||
url = f"http://{self.host}/{self.sec}/?cmd={cmd}"
|
||||
self.lg.debug('run command: %s', url)
|
||||
async with self.lck:
|
||||
return await self.request(pt=port, cmd=cmd)
|
||||
|
||||
async def request(self, **kwargs):
|
||||
cmd = '&'.join([f'{k}={v}' for k, v in kwargs.items() if v is not None])
|
||||
url = f"http://{self.host}/{self.sec}/?{cmd}"
|
||||
self.lg.debug('request: %s', url)
|
||||
async with self._http_lck:
|
||||
async with aiohttp.request("get", url=url) as req:
|
||||
if req.status != 200:
|
||||
self.lg.warning('%s returned %s (%s)', url, req.status, await req.text())
|
||||
return False
|
||||
return None
|
||||
else:
|
||||
return True
|
||||
return await req.text()
|
||||
|
||||
async def save(self):
|
||||
await self.send_command(cmd='s')
|
||||
|
||||
def parse_response(self, ret):
|
||||
if ret is None:
|
||||
raise NoPort()
|
||||
if ':' in ret:
|
||||
ret = PATT_SPLIT.split(ret)
|
||||
ret = dict([
|
||||
x.split(':') for x in ret if x.count(':') == 1
|
||||
])
|
||||
elif 'ON' in ret:
|
||||
ret = {'value': 'ON'}
|
||||
elif 'OFF' in ret:
|
||||
ret = {'value': 'OFF'}
|
||||
else:
|
||||
ret = {'value': ret}
|
||||
return ret
|
||||
|
||||
async def get_port(self, port):
|
||||
"""Запрос состояния порта. Блокируется пока не придет какое-нибудь сообщение от меги или таймаут"""
|
||||
"""
|
||||
Запрос состояния порта. Состояние всегда возвращается в виде объекта, всегда сохраняется в центральное
|
||||
хранилище values
|
||||
"""
|
||||
self.lg.debug(f'get port %s', port)
|
||||
if self.mqtt is None:
|
||||
ret = await self.request(pt=port, cmd='get')
|
||||
ret = self.parse_response(ret)
|
||||
self.values[port] = ret
|
||||
return ret
|
||||
|
||||
async with self._notif_lck:
|
||||
async with self.notifiers[port]:
|
||||
cnd = self.notifiers[port]
|
||||
@@ -184,13 +232,19 @@ class MegaD:
|
||||
)
|
||||
try:
|
||||
await asyncio.wait_for(cnd.wait(), timeout=10)
|
||||
return self.values[port]
|
||||
return self.values.get(port)
|
||||
except asyncio.TimeoutError:
|
||||
self.lg.error(f'timeout when getting port {port}')
|
||||
|
||||
async def get_all_ports(self):
|
||||
for x in range(37):
|
||||
await self.get_port(x)
|
||||
if not self.mqtt_inputs:
|
||||
ret = await self.request(cmd='all')
|
||||
for port, x in enumerate(ret.split(';')):
|
||||
ret = self.parse_response(x)
|
||||
self.values[port] = ret
|
||||
else:
|
||||
for x in range(self.nports + 1):
|
||||
await self.get_port(x)
|
||||
|
||||
async def reboot(self, save=True):
|
||||
await self.save()
|
||||
@@ -198,7 +252,6 @@ class MegaD:
|
||||
async def _notify(self, port, value):
|
||||
async with self.notifiers[port]:
|
||||
cnd = self.notifiers[port]
|
||||
self.values[port] = value
|
||||
cnd.notify_all()
|
||||
|
||||
def _process_msg(self, msg):
|
||||
@@ -222,6 +275,7 @@ class MegaD:
|
||||
value = None
|
||||
try:
|
||||
value = json.loads(msg.payload)
|
||||
self.values[port] = value
|
||||
for cb in self._callbacks[port]:
|
||||
cb(value)
|
||||
except Exception as exc:
|
||||
@@ -235,7 +289,10 @@ class MegaD:
|
||||
self.lg.debug(
|
||||
f'subscribe %s %s', port, callback
|
||||
)
|
||||
self._callbacks[port].append(callback)
|
||||
if self.mqtt_inputs:
|
||||
self._callbacks[port].append(callback)
|
||||
else:
|
||||
self.http.callbacks[port].append(callback)
|
||||
|
||||
async def authenticate(self) -> bool:
|
||||
"""Test if we can authenticate with the host."""
|
||||
@@ -263,6 +320,8 @@ class MegaD:
|
||||
)
|
||||
async with aiohttp.request('get', url) as req:
|
||||
html = await req.text()
|
||||
if req.status != 200:
|
||||
return
|
||||
tree = BeautifulSoup(html, features="lxml")
|
||||
pty = tree.find('select', attrs={'name': 'pty'})
|
||||
if pty is None:
|
||||
@@ -286,15 +345,16 @@ class MegaD:
|
||||
self._scanned[port] = (pty, m)
|
||||
return pty, m
|
||||
|
||||
async def scan_ports(self,):
|
||||
for x in range(38):
|
||||
async def scan_ports(self, nports=37):
|
||||
for x in range(nports+1):
|
||||
ret = await self.scan_port(x)
|
||||
if ret:
|
||||
yield [x, *ret]
|
||||
self.nports = nports+1
|
||||
|
||||
async def get_config(self):
|
||||
async def get_config(self, nports=37):
|
||||
ret = defaultdict(lambda: defaultdict(list))
|
||||
async for port, pty, m in self.scan_ports():
|
||||
async for port, pty, m in self.scan_ports(nports):
|
||||
if pty == "0":
|
||||
ret['binary_sensor'][port].append({})
|
||||
elif pty == "1" and m in ['0', '1']:
|
||||
|
||||
@@ -10,9 +10,7 @@
|
||||
"ssdp": [],
|
||||
"zeroconf": [],
|
||||
"homekit": {},
|
||||
"dependencies": [
|
||||
"mqtt"
|
||||
],
|
||||
"dependencies": [],
|
||||
"codeowners": [
|
||||
"@andvikt"
|
||||
],
|
||||
|
||||
@@ -111,7 +111,6 @@ class Mega1WSensor(MegaPushEntity):
|
||||
:param patt: pattern to extract value, must have at least one group that will contain parsed value
|
||||
"""
|
||||
super().__init__(*args, **kwargs)
|
||||
self.mega.sensors.append(self)
|
||||
self._value = None
|
||||
self.key = key
|
||||
self._device_class = device_class
|
||||
@@ -152,9 +151,6 @@ class Mega1WSensor(MegaPushEntity):
|
||||
ret = self._state.state
|
||||
return ret
|
||||
|
||||
def _update(self, payload: dict):
|
||||
self.mega.values[self.port] = payload
|
||||
|
||||
@property
|
||||
def name(self):
|
||||
n = super().name
|
||||
|
||||
@@ -11,7 +11,10 @@
|
||||
"mqtt_id": "[%key:common::config_flow::data::mqtt_id%]",
|
||||
"scan_interval": "[%key:common::config_flow::data::mqtt_id%]",
|
||||
"port_to_scan": "[%key:common::config_flow::data::port_to_scan%]",
|
||||
"invert": "[%key:common::config_flow::data::invert%]"
|
||||
"invert": "[%key:common::config_flow::data::invert%]",
|
||||
"mqtt_inputs": "[%key:common::config_flow::data::mqtt_inputs%]",
|
||||
"nports": "[%key:common::config_flow::data::nports%]",
|
||||
"update_all": "[%key:common::config_flow::data::update_all%]"
|
||||
}
|
||||
}
|
||||
},
|
||||
@@ -32,7 +35,9 @@
|
||||
"scan_interval": "[%key:common::config_flow::data::scan_interval%]",
|
||||
"port_to_scan": "[%key:common::config_flow::data::port_to_scan%]",
|
||||
"reload": "[%key:common::config_flow::data::reload%]",
|
||||
"invert": "[%key:common::config_flow::data::invert%]"
|
||||
"invert": "[%key:common::config_flow::data::invert%]",
|
||||
"mqtt_inputs": "[%key:common::config_flow::data::mqtt_inputs%]",
|
||||
"nports": "[%key:common::config_flow::data::nports%]"
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -7,7 +7,8 @@
|
||||
"cannot_connect": "Failed to connect",
|
||||
"invalid_auth": "Invalid authentication",
|
||||
"unknown": "Unexpected error",
|
||||
"duplicate_id": "Duplicate ID"
|
||||
"duplicate_id": "Duplicate ID",
|
||||
"mqtt_inputs": "Use MQTT"
|
||||
},
|
||||
"step": {
|
||||
"user": {
|
||||
@@ -18,7 +19,9 @@
|
||||
"id": "ID",
|
||||
"mqtt_id": "MQTT id",
|
||||
"scan_interval": "Scan interval (sec), 0 - don't update",
|
||||
"port_to_scan": "Port to poll aliveness (needed only if no sensors used)"
|
||||
"port_to_scan": "Port to poll aliveness (needed only if no sensors used)",
|
||||
"nports": "Number of ports",
|
||||
"update_all": "Update all outs when input"
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -29,7 +32,9 @@
|
||||
"data": {
|
||||
"scan_interval": "Scan interval (sec), 0 - don't update",
|
||||
"port_to_scan": "Port to poll aliveness (needed only if no sensors used)",
|
||||
"reload": "Reload objects"
|
||||
"reload": "Reload objects",
|
||||
"mqtt_inputs": "Use MQTT",
|
||||
"update_all": "Update all outs when input"
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -18,7 +18,10 @@
|
||||
"id": "ID",
|
||||
"mqtt_id": "MQTT id",
|
||||
"scan_interval": "Периодичность обновлений (сек.), 0 - не обновлять",
|
||||
"port_to_scan": "Порт, который сканируется когда нет датчиков"
|
||||
"port_to_scan": "Порт, который сканируется когда нет датчиков",
|
||||
"mqtt_inputs": "Использовать MQTT",
|
||||
"nports": "Кол-во портов",
|
||||
"update_all": "Обновить все выходы когда срабатывает вход"
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -30,7 +33,10 @@
|
||||
"scan_interval": "Периодичность обновлений (сек.), 0 - не обновлять",
|
||||
"port_to_scan": "Порт, который сканируется когда нет датчиков",
|
||||
"reload": "Обновить объекты",
|
||||
"invert": "Список портов (через ,) с инвертированной логикой"
|
||||
"invert": "Список портов (через ,) с инвертированной логикой",
|
||||
"mqtt_inputs": "Использовать MQTT",
|
||||
"nports": "Кол-во портов",
|
||||
"update_all": "Обновить все выходы когда срабатывает вход"
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user