Files
Home-Assistant/custom_components/hilo/__init__.py
T
2026-06-11 11:50:50 -04:00

1189 lines
47 KiB
Python

"""Support for Hilo automation systems."""
from __future__ import annotations
import asyncio
from collections import OrderedDict
from datetime import datetime, timedelta
import traceback
from typing import List, Optional
from aiohttp import CookieJar, client_exceptions
from homeassistant.components.select import (
ATTR_OPTION,
DOMAIN as SELECT_DOMAIN,
SERVICE_SELECT_OPTION,
)
from homeassistant.components.sensor import SensorDeviceClass
from homeassistant.config_entries import ConfigEntry
from homeassistant.const import (
ATTR_DEVICE_CLASS,
ATTR_UNIT_OF_MEASUREMENT,
CONF_SCAN_INTERVAL,
CONF_USERNAME,
EVENT_HOMEASSISTANT_STOP,
Platform,
)
from homeassistant.core import Context, HomeAssistant, callback
from homeassistant.exceptions import ConfigEntryAuthFailed, ConfigEntryNotReady
from homeassistant.helpers import (
config_entry_oauth2_flow,
device_registry as dr,
entity_registry as er,
)
from homeassistant.helpers.aiohttp_client import async_create_clientsession
from homeassistant.helpers.dispatcher import async_dispatcher_send
from homeassistant.helpers.update_coordinator import DataUpdateCoordinator
from pyhilo import API
from pyhilo.device import HiloDevice
from pyhilo.devices import Devices
from pyhilo.event import Event
from pyhilo.exceptions import HiloError
from pyhilo.graphql import GraphQlHelper
from pyhilo.signalr import SignalREvent, signalr_event_from_payload
from pyhilo.util import from_utc_timestamp, time_diff
from pysignalr.exceptions import ServerError as SignalRServerError
from .config_flow import STEP_OPTION_SCHEMA, HiloFlowHandler
from .const import (
CONF_APPRECIATION_PHASE,
CONF_CHALLENGE_LOCK,
CONF_GENERATE_ENERGY_METERS,
CONF_HQ_PLAN_NAME,
CONF_LOG_TRACES,
CONF_PRE_COLD_PHASE,
CONF_TARIFF,
CONF_TRACK_UNKNOWN_SOURCES,
CONF_UNTARIFICATED_DEVICES,
DEFAULT_APPRECIATION_PHASE,
DEFAULT_CHALLENGE_LOCK,
DEFAULT_GENERATE_ENERGY_METERS,
DEFAULT_HQ_PLAN_NAME,
DEFAULT_LOG_TRACES,
DEFAULT_PRE_COLD_PHASE,
DEFAULT_SCAN_INTERVAL,
DEFAULT_TRACK_UNKNOWN_SOURCES,
DEFAULT_UNTARIFICATED_DEVICES,
DOMAIN,
HILO_ENERGY_TOTAL,
LOG,
MIN_SCAN_INTERVAL,
)
from .oauth2 import AuthCodeWithPKCEImplementation
DISPATCHER_TOPIC_SIGNALR_EVENT = "pyhilo_signalr_event"
SIGNAL_UPDATE_ENTITY = "pyhilo_device_update_{}"
COORDINATOR_AWARE_PLATFORMS = [Platform.SENSOR]
PLATFORMS = COORDINATOR_AWARE_PLATFORMS + [
Platform.CLIMATE,
Platform.LIGHT,
Platform.SWITCH,
]
@callback
def _async_standardize_config_entry(hass: HomeAssistant, entry: ConfigEntry) -> None:
"""Bring a config entry up to current standards."""
entry_updates = {}
config_keys = STEP_OPTION_SCHEMA.schema.keys()
if not entry.unique_id:
# If the config entry doesn't already have a unique ID, set one:
entry_updates["unique_id"] = entry.data[CONF_USERNAME]
if any(x in entry.data for x in config_keys):
# If an option was provided as part of configuration.yaml, pop it out of
# the config entry's data and move it to options and the same with other
# possible options.
data = {**entry.data}
entry_updates["data"] = data
options = {}
for conf_item in config_keys:
if attr := data.pop(conf_item, None):
options[conf_item] = attr
if len(options):
entry_updates["options"] = {**entry.options, **options}
if entry_updates:
hass.config_entries.async_update_entry(entry, **entry_updates)
@callback
def _async_register_custom_device(
hass: HomeAssistant, entry: ConfigEntry, device: HiloDevice
) -> None:
"""Register a custom device.
This is used to register the Hilo gateway and the unknown source tracker.
"""
LOG.debug("Generating custom device %s", device)
device_registry = dr.async_get(hass)
device_registry.async_get_or_create(
config_entry_id=entry.entry_id,
identifiers={(DOMAIN, device.identifier)},
manufacturer=device.manufacturer,
model=device.model,
name=device.name,
)
async def async_setup_entry( # noqa: C901
hass: HomeAssistant, entry: ConfigEntry
) -> bool:
"""Set up Hilo as config entry."""
HiloFlowHandler.async_register_implementation(
hass, AuthCodeWithPKCEImplementation(hass)
)
implementation = (
await config_entry_oauth2_flow.async_get_config_entry_implementation(
hass, entry
)
)
current_options = {**entry.options}
try:
api = await API.async_create(
session=async_create_clientsession(
hass, cookie_jar=CookieJar(quote_cookie=False)
),
oauth_session=config_entry_oauth2_flow.OAuth2Session(
hass, entry, implementation
),
log_traces=current_options.get(CONF_LOG_TRACES, DEFAULT_LOG_TRACES),
)
except (TimeoutError, client_exceptions.ClientConnectorError):
LOG.debug("Timeout")
raise ConfigEntryNotReady
except Exception as err:
raise ConfigEntryAuthFailed(err) from err
_async_standardize_config_entry(hass, entry)
scan_interval = current_options.get(CONF_SCAN_INTERVAL, DEFAULT_SCAN_INTERVAL)
scan_interval = (
scan_interval if scan_interval >= MIN_SCAN_INTERVAL else MIN_SCAN_INTERVAL
)
hilo = Hilo(hass, entry, api)
try:
await hilo.async_init(scan_interval)
except HiloError as err:
raise ConfigEntryNotReady from err
hass.data.setdefault(DOMAIN, {})
hass.data[DOMAIN][entry.entry_id] = hilo
hass.async_create_task(
hass.config_entries.async_forward_entry_setups(entry, PLATFORMS)
)
# Note (ic-dev21): This is a bit of a hack to rename some entities that were created with non-standard names in early versions
# HA has changed the way they name entities linked to a device by default and this breaks the naming scheme of the gateway entities.
# If we let the HA default behaviour, all gateway entities (such as challenge sensor) will be prefixed with "hilo_gateway_",
# This code will check for the old entity IDs and rename them to the new format if they exist.
entity_registry = er.async_get(hass)
gateway_entity_renames = {
"sensor.hilo_gateway_defi_hilo": "sensor.defi_hilo",
"sensor.hilo_gateway_notifications_hilo": "sensor.notifications_hilo",
"sensor.hilo_gateway_recompenses_hilo": "sensor.recompenses_hilo",
"sensor.hilo_gateway_outdoor_weather_hilo": "sensor.outdoor_weather_hilo",
"sensor.hilo_gateway_hilo_rate_current": "sensor.hilo_rate_current",
"sensor.hilo_gateway_hilo_rate_low": "sensor.hilo_rate_low",
"sensor.hilo_gateway_hilo_rate_medium": "sensor.hilo_rate_medium",
"sensor.hilo_gateway_hilo_rate_access": "sensor.hilo_rate_access",
"sensor.hilo_gateway_hilo_rate_low_threshold": "sensor.hilo_rate_low_threshold",
"sensor.hilo_gateway_hilo_rate_reward_rate": "sensor.hilo_rate_reward_rate",
"sensor.hilo_gateway_hilo_cost_total": "sensor.hilo_cost_total",
"sensor.hilo_gateway": "sensor.hilo_gateway",
}
for old_id, new_id in gateway_entity_renames.items():
if old_id != new_id and entity_registry.async_get(old_id):
entity_registry.async_update_entity(old_id, new_entity_id=new_id)
LOG.info("Migrated entity ID %s -> %s", old_id, new_id)
# Note (ic-dev21): this new renaming by HA also breaks sensor.hilo_energy_total, which is used in check_tarif
energy_entity_renames = {
"sensor.meter00_hilo_energy_total": "sensor.hilo_energy_total"
}
for old_id, new_id in energy_entity_renames.items():
if old_id != new_id and entity_registry.async_get(old_id):
entity_registry.async_update_entity(old_id, new_entity_id=new_id)
LOG.info("Migrated entity ID %s -> %s", old_id, new_id)
async def handle_debug_event(event: Event):
"""Handle an event."""
LOG.debug("HILO_DEBUG: Event received: %s", event)
log_traces = current_options.get(CONF_LOG_TRACES)
LOG.debug("HILO_DEBUG: log_traces is %s", log_traces)
signalr_event = signalr_event_from_payload(event.data)
LOG.debug("HILO_DEBUG: SignalR event parsed: %s", signalr_event)
await hilo.on_signalr_event(signalr_event)
log_traces = current_options.get(CONF_LOG_TRACES)
if log_traces:
LOG.debug("HILO_DEBUG: log_traces is %s", log_traces)
hass.bus.async_listen("hilo_debug", handle_debug_event)
async def async_reload_entry(_: HomeAssistant, updated_entry: ConfigEntry) -> None:
"""Handle an options update.
This method will get called in two scenarios:
1. When HiloOptionsFlowHandler is initiated
2. When a new refresh token is saved to the config entry data
We only want #1 to trigger an actual reload.
"""
updated_options = {**updated_entry.options}
if updated_options == current_options:
return
await hass.config_entries.async_reload(entry.entry_id)
entry.async_on_unload(entry.add_update_listener(async_reload_entry))
return True
async def async_unload_entry(hass: HomeAssistant, entry: ConfigEntry) -> bool:
"""Unload a Hilo config entry."""
LOG.debug("Unloading Hilo Integration")
hilo = hass.data[DOMAIN][entry.entry_id]
hilo.should_signalr_reconnect = False
for task in list(hilo._signalr_reconnect_tasks):
if not task.done():
task.cancel()
try:
await task
except asyncio.CancelledError:
pass
try:
await hilo._api.signalr_devices.disconnect()
except Exception as err:
LOG.error("Error disconnecting device SignalR hub: %s", err)
try:
await hilo._api.signalr_challenges.disconnect()
except Exception as err:
LOG.error("Error disconnecting challenge SignalR hub: %s", err)
unload_ok = await hass.config_entries.async_unload_platforms(entry, PLATFORMS)
if unload_ok:
try:
if hasattr(hilo, "_api") and hilo._api and hasattr(hilo._api, "session"):
if hilo._api.session and not hilo._api.session.closed:
await hilo._api.session.close()
LOG.debug("Session closed")
except Exception as err:
LOG.error("Error closing session: %s", err)
LOG.debug("Hilo Integration unloaded")
hass.data[DOMAIN].pop(entry.entry_id)
return unload_ok
async def async_migrate_entry(hass, config_entry: ConfigEntry):
"""Migrate old entry."""
LOG.debug("Migrating from version %s", config_entry.version)
if config_entry.version > 1:
# This means the user has downgraded from a future version
return False
if config_entry.version == 1:
config_entry.version = 2
hass.config_entries.async_update_entry(
config_entry, unique_id="hilo", data={"auth_implementation": "hilo"}
)
LOG.debug("Migration to version %s successful", config_entry.version)
return True
class Hilo:
"""Define a Hilo data object."""
def __init__(self, hass: HomeAssistant, entry: ConfigEntry, api: API) -> None:
"""Initialize."""
self._api = api
self._hass = hass
self.find_meter(self._hass)
self.entry = entry
self.devices: Devices = Devices(api)
self.graphql_helper: GraphQlHelper = GraphQlHelper(api, self.devices)
self.challenge_id = 0
self._should_signalr_reconnect = True
self._signalr_reconnect_tasks: list[asyncio.Task | None] = [None, None]
self._update_task: list[asyncio.Task | None] = [None, None]
self.subscriptions: List[Optional[asyncio.Task]] = [None]
self.hq_plan_name = entry.options.get(CONF_HQ_PLAN_NAME, DEFAULT_HQ_PLAN_NAME)
self.appreciation = entry.options.get(
CONF_APPRECIATION_PHASE, DEFAULT_APPRECIATION_PHASE
)
self.pre_cold = entry.options.get(CONF_PRE_COLD_PHASE, DEFAULT_PRE_COLD_PHASE)
self.challenge_lock = entry.options.get(
CONF_CHALLENGE_LOCK, DEFAULT_CHALLENGE_LOCK
)
self.track_unknown_sources = entry.options.get(
CONF_TRACK_UNKNOWN_SOURCES, DEFAULT_TRACK_UNKNOWN_SOURCES
)
self.untarificated_devices = entry.options.get(
CONF_UNTARIFICATED_DEVICES, DEFAULT_UNTARIFICATED_DEVICES
)
self.generate_energy_meters = entry.options.get(
CONF_GENERATE_ENERGY_METERS, DEFAULT_GENERATE_ENERGY_METERS
)
# This will get filled in by async_init:
self.coordinator: DataUpdateCoordinator | None = None
self.unknown_tracker_device: HiloDevice | None = None
self._events: dict = {}
if self.track_unknown_sources:
self._api._get_device_callbacks = [self._get_unknown_source_tracker]
self._signalr_listeners = []
async def _on_devices_connected(self) -> None:
"""Trigger device subscriptions after the device hub connects."""
await self.subscribe_to_location()
async def _on_challenges_connected(self) -> None:
"""Trigger challenge subscriptions after the challenge hub connects."""
await self.subscribe_to_challenge()
await self.subscribe_to_challengelist()
def validate_heartbeat(self, event: SignalREvent) -> None:
"""Validate heartbeat messages from SignalR."""
heartbeat_time = from_utc_timestamp(event.arguments[0]) # type: ignore
if self._api.log_traces:
LOG.debug("Heartbeat: %s", time_diff(heartbeat_time, event.timestamp))
def register_signalr_listener(self, listener):
"""Register a listener for SignalR events."""
LOG.debug("Registering SignalR listener: %s", listener.__class__.__name__)
self._signalr_listeners.append(listener)
async def _handle_signalr_message(self, event):
"""Process SignalR messages and notify listeners."""
# TODO: ic-dev21: This needs to be cleaned up and optimized
LOG.debug("Received SignalR message type: %s", event)
target = event.target
LOG.debug("handle_signalr_message_target %s", target)
msg_data = event
LOG.debug("handle_signalr_message_ msg_data %s", msg_data)
if target in [
"ChallengeListInitialValuesReceived",
"EventListInitialValuesReceived",
]:
msg_type = "challenge_list_initial"
elif target in ["ChallengeAdded", "EventAdded"]:
msg_type = "challenge_added"
elif target in [
"ChallengeDetailsUpdated",
"ChallengeConsumptionUpdatedValuesReceived",
"EventCHConsumptionUpdatedValuesReceived",
"ChallengeDetailsUpdatedValuesReceived",
"EventCHDetailsUpdatedValuesReceived",
"EventFlexDetailsUpdatedValuesReceived",
"ChallengeDetailsInitialValuesReceived",
"EventCHDetailsInitialValuesReceived",
"EventFlexDetailsInitialValuesReceived",
"ChallengeListUpdatedValuesReceived",
"EventListUpdatedValuesReceived",
]:
msg_type = "challenge_details_update"
elif target == "EventFlexConsumptionUpdatedValuesReceived":
LOG.debug("%s message received", target)
LOG.debug("%s data: %s", target, msg_data)
return
# ic-dev21 Notify listeners
for listener in self._signalr_listeners:
handler_name = f"handle_{msg_type}"
if hasattr(listener, handler_name):
handler = getattr(listener, handler_name)
try:
# ic-dev21 Extract the arguments from the SignalREvent object
if isinstance(msg_data, SignalREvent):
arguments = msg_data.arguments
if arguments: # ic-dev21 check if there are arguments
await handler(arguments[0])
else:
LOG.warning(
f"SHOULD NOT HAPPEN: Received empty arguments for {msg_type}"
)
else:
LOG.warning(f"SHOULD NOT HAPPEN: Not SignalREvent: {msg_data}")
await handler(msg_data)
except Exception as e:
LOG.error("Error in SignalR handler %s: %s", handler_name, e)
LOG.error(traceback.format_exc())
async def _handle_challenge_events(self, event: SignalREvent) -> None:
"""Handle all challenge-related SignalR events."""
if event.target == "ChallengeDetailsInitialValuesReceived":
challenge = event.arguments[0]
LOG.debug(
"ChallengeDetailsInitialValuesReceived, challenge = %s", challenge
)
self.challenge_id = challenge.get("id")
elif event.target == "ChallengeDetailsUpdatedValuesReceived":
LOG.debug("ChallengeDetailsUpdatedValuesReceived")
elif event.target == "ChallengeListUpdatedValuesReceived":
LOG.debug("ChallengeListUpdatedValuesReceived")
self.challenge_phase = event.arguments[0][0]["currentPhase"]
elif event.target == "ChallengeAdded":
LOG.debug("ChallengeAdded")
challenge = event.arguments[0]
self.challenge_id = challenge.get("id")
await self.subscribe_to_challenge(self.challenge_id)
elif event.target == "ChallengeListInitialValuesReceived":
LOG.debug("ChallengeListInitialValuesReceived")
challenges = event.arguments[0]
for challenge in challenges:
challenge_id = challenge.get("id")
self.challenge_phase = challenge.get("currentPhase")
self.challenge_id = challenge.get("id")
await self.subscribe_to_challenge(challenge_id)
elif event.target == "EventCHDetailsUpdatedValuesReceived":
LOG.debug("EventCHDetailsUpdatedValuesReceived")
data = event.arguments[0]
if "report" in data:
report = data["report"]
event_id = data.get("id")
LOG.debug("Report for event %s: %s", event_id, report)
async def _handle_device_events(self, event: SignalREvent) -> None:
"""Handle all device-related SignalR events."""
if event.target == "DevicesValuesReceived":
new_devices = any(
self.devices.find_device(item["deviceId"]) is None
for item in event.arguments[0]
)
if new_devices:
LOG.warning(
"Device list appears to be desynchronized, "
"waiting for next DeviceListInitialValuesReceived to refresh..."
)
# Device list will refresh on next SignalR reconnect/subscribe
updated_devices = self.devices.parse_values_received(event.arguments[0])
# NOTE(dvd): If we don't do this, we need to wait until the coordinator
# runs (scan_interval) to have updated data in the dashboard.
for device in updated_devices:
async_dispatcher_send(
self._hass, SIGNAL_UPDATE_ENTITY.format(device.id)
)
elif event.target == "DeviceListInitialValuesReceived":
await self.devices.update_devicelist_from_signalr(event.arguments[0])
elif event.target == "DeviceListUpdatedValuesReceived":
# This message only contains display information, such as the Device's name (as set in the app), it's groupid, icon, etc.
# Updating the device name causes issues in the integration, it detects it as a new device and creates a new entity.
# Ignore this call, for now... (update_devicelist_from_signalr does work, but causes the issue above)
# await self.devices.update_devicelist_from_signalr(event.arguments[0])
LOG.debug(
"Received 'DeviceListUpdatedValuesReceived' message, not implemented yet."
)
elif event.target == "DevicesListChanged":
LOG.debug("Received 'DevicesListChanged' message, not implemented yet.")
elif event.target == "DeviceAdded":
devices = [event.arguments[0]]
await self.devices.add_device_from_signalr(devices)
elif event.target == "DeviceDeleted":
LOG.debug("Received 'DeviceDeleted' message, not implemented yet.")
elif event.target == "GatewayValuesReceived":
gateway = self.devices.find_device(1)
if gateway:
gateway.id = event.arguments[0][0]["deviceId"]
LOG.debug("Updated Gateway's deviceId from default 1 to %s", gateway.id)
updated_devices = self.devices.parse_values_received(event.arguments[0])
for device in updated_devices:
async_dispatcher_send(
self._hass, SIGNAL_UPDATE_ENTITY.format(device.id)
)
@callback
async def on_signalr_event(self, event: SignalREvent) -> None:
"""Define a callback for receiving a SignalR event."""
async_dispatcher_send(self._hass, DISPATCHER_TOPIC_SIGNALR_EVENT, event)
if event.target == "Heartbeat":
self.validate_heartbeat(event)
elif "Challenge" in event.target or "Event" in event.target:
LOG.debug("HILO_DEBUG: Handling challenge/event SignalR event: %s", event)
await self._handle_challenge_events(event)
await self._handle_signalr_message(event)
elif "Device" in event.target or event.target == "GatewayValuesReceived":
await self._handle_device_events(event)
else:
LOG.warning("Unhandled SignalR event: %s", event)
@callback
async def subscribe_to_location(self) -> None:
"""Send the json payload to receive updates from the location."""
LOG.debug("Subscribing to location %s", self.devices.location_id)
await self._api.signalr_devices.invoke(
"SubscribeToLocation", [self.devices.location_id]
)
@callback
async def subscribe_to_challenge(self, event_id: int = 0) -> None:
"""Send the json payload to receive updates from the challenge."""
LOG.debug("Subscribing to challenge : %s or %s", event_id, self.challenge_id)
event_id = event_id or self.challenge_id
LOG.debug("API URN is %s", self._api.urn)
# Get plan name to connect to the correct challenge hub list
tarif_config = self.hq_plan_name
LOG.debug("Event list needed is %s", tarif_config)
# TODO: This is a fallback but will eventually need to be removed, I expect it to create
# websocket disconnects once the split is complete.
LOG.warning(
"Starting legacy connection to ChallengeHub. Your tarif is %s, and will also attempt connection. This can be safely ignored. This will be deprecated",
tarif_config,
)
await self._api.signalr_challenges.invoke(
"SubscribeToChallenge",
[{"locationId": self.devices.location_id, "eventId": event_id}],
)
# Subscribe to the correct challenge hub
if tarif_config == "rate d":
await self._api.signalr_challenges.invoke(
"SubscribeToEventCH",
[{"locationHiloId": self._api.urn, "eventId": event_id}],
)
elif tarif_config == "flex d":
await self._api.signalr_challenges.invoke(
"SubscribeToEventFlex",
[{"locationHiloId": self._api.urn, "eventId": event_id}],
)
else:
LOG.warning("Unknown plan name %s, falling back to default", tarif_config)
await self._api.signalr_challenges.invoke(
"SubscribeToChallenge",
[{"locationId": self.devices.location_id, "eventId": event_id}],
)
@callback
async def subscribe_to_challengelist(self) -> None:
"""Send the json payload to receive updates from the challenge list."""
# TODO : Rename challenge functions to Event, fallback on challenge for now
LOG.debug(
"Subscribing to challenge list at location %s", self.devices.location_id
)
LOG.debug("API URN is %s", self._api.urn)
await self._api.signalr_challenges.invoke(
"SubscribeToChallengeList",
[{"locationId": self.devices.location_id}],
)
LOG.debug("Subscribing to event list at location %s", self.devices.location_id)
await self._api.signalr_challenges.invoke(
"SubscribeToEventList",
[{"locationHiloId": self._api.urn}],
)
@callback
async def request_challenge_consumption_update(self, event_id: int = 0) -> None:
"""Send the json payload to receive energy consumption updates from the challenge."""
event_id = event_id or self.challenge_id
# TODO: Remove fallback once split is complete
LOG.debug(
"Requesting challenge %s consumption update at location %s",
event_id,
self.devices.location_id,
)
await self._api.signalr_challenges.invoke(
"RequestChallengeConsumptionUpdate",
[{"locationId": self.devices.location_id, "eventId": event_id}],
)
# Get plan name to request the correct consumption update
tarif_config = self.hq_plan_name
LOG.debug("API URN is %s", self._api.urn)
if tarif_config == "rate d":
LOG.debug(
"Requesting event CH consumption update at location %s",
self.devices.location_id,
)
await self._api.signalr_challenges.invoke(
"RequestEventCHConsumptionUpdate",
[{"locationHiloId": self._api.urn, "eventId": event_id}],
)
elif tarif_config == "flex d":
LOG.debug(
"Requesting event Flex consumption update at location %s",
self.devices.location_id,
)
await self._api.signalr_challenges.invoke(
"RequestEventFlexConsumptionUpdate",
[{"locationHiloId": self._api.urn, "eventId": event_id}],
)
else:
LOG.debug(
"Requesting challenge %s consumption update at location %s",
event_id,
self.devices.location_id,
)
await self._api.signalr_challenges.invoke(
"RequestChallengeConsumptionUpdate",
[{"locationId": self.devices.location_id, "eventId": event_id}],
)
@callback
def _get_unknown_source_tracker(self) -> HiloDevice:
return {
"name": "Unknown Source Tracker",
"Disconnected": False,
"type": "Tracker",
"category": "Tracker",
"supportedAttributes": "Power",
"settableAttributes": "",
"id": 69420,
"hilo_id": "AB-A2025",
"identifier": "hass-hilo-unknown_source_tracker",
"provider": 0,
"model_number": "Hass-hilo-2025.5",
"sw_version": "0.0.1",
}
async def get_event_details(self, event_id: int):
"""Get events from Hilo only when necessary, otherwise, we hit the cache.
When preheat is started and our last update is before
the preheat_start, we refresh. This should update the
allowed_kWh, etc. values.
"""
if event_data := self._events.get(event_id):
event = Event(**event_data)
if event.invalid:
LOG.debug(
"Invalidating cache for event %s during %s phase (event.current_phase_times=%s event.last_update=%s)",
event_id,
event.state,
event.current_phase_times,
event.last_update,
)
del self._events[event_id]
"""
Note ic-dev21: temp fix until we an make it prettier.
During appreciation, pre-heat and reduction we delete
the event attributes and reload them with the next if,
the rest of time time we're reading it from cache
"""
if event.state in ["appreciation", "pre_heat", "reduction"]:
LOG.debug(
"Invalidating cache for event %s during appreciation, pre_heat or reduction phase (event.last_update=%s)",
event_id,
event.last_update,
)
del self._events[event_id]
if event_id not in self._events:
self._events[event_id] = await self._api.get_gd_events(
self.devices.location_id, event_id=event_id
)
return self._events[event_id]
async def async_init(self, scan_interval: int) -> None:
"""Initialize the Hilo "manager" class.
Flow:
1. Get location IDs (REST - kept)
2. Register websocket callbacks and connect
3. Websocket subscribes to location, triggering DeviceListInitialValuesReceived
4. Wait for device cache to be populated from websocket
5. Build device list (websocket devices + gateway from REST)
6. Initialize GraphQL, register custom devices, start coordinator
"""
# Step 1: Get location IDs (still REST)
await self.devices.async_init()
# Step 2: Register SignalR callbacks and start connections
self._api.signalr_devices.add_connect_callback(self._on_devices_connected)
self._api.signalr_challenges.add_connect_callback(self._on_challenges_connected)
self._api.signalr_devices.add_event_callback(self.on_signalr_event)
self._api.signalr_challenges.add_event_callback(self.on_signalr_event)
self._signalr_reconnect_tasks[0] = asyncio.create_task(
self.start_signalr_loop(self._api.signalr_devices, 0)
)
self._signalr_reconnect_tasks[1] = asyncio.create_task(
self.start_signalr_loop(self._api.signalr_challenges, 1)
)
# Step 3: Wait for DeviceListInitialValuesReceived from SignalR
await self._api.wait_for_device_cache(timeout=30.0)
# Step 4: Build device list (websocket devices + gateway REST + callbacks)
await self.devices.update()
# Step 5: Initialize GraphQL
await self.graphql_helper.async_init()
self.subscriptions[0] = asyncio.create_task(
self.graphql_helper.subscribe_to_device_updated(
self.devices.location_hilo_id,
self.handle_subscription_result,
)
)
# Step 6: Register custom devices in HA
_async_register_custom_device(
self._hass, self.entry, self.devices.find_device(1)
)
if self.track_unknown_sources:
if not self.unknown_tracker_device:
self.unknown_tracker_device = self.devices.generate_device(
self._get_unknown_source_tracker()
)
self.unknown_tracker_device.net_consumption = True
_async_register_custom_device(
self._hass, self.entry, self.unknown_tracker_device
)
async def signalr_disconnect_listener(_: Event) -> None:
"""Define an event handler to disconnect from the SignalR hubs."""
await self._api.signalr_devices.disconnect()
await self._api.signalr_challenges.disconnect()
self.entry.async_on_unload(
self._hass.bus.async_listen_once(
EVENT_HOMEASSISTANT_STOP, signalr_disconnect_listener
)
)
self.coordinator = DataUpdateCoordinator(
self._hass,
LOG,
name="hilo",
update_interval=timedelta(seconds=scan_interval),
update_method=self.async_update,
)
async def start_signalr_loop(self, hub, id) -> None:
"""Start a SignalR reconnection loop that retries forever until HA stops."""
backoff = 5 # seconds; doubles on each error, reset to 5 after a clean run
while self.should_signalr_reconnect:
try:
LOG.info("SignalRHub[%s]: connecting", id)
await hub.run()
# hub.run() returned without raising — server closed the connection.
# That's a normal disconnect; reset backoff and reconnect quickly.
LOG.warning(
"SignalRHub[%s]: connection closed by server; reconnecting in %ss",
id,
backoff,
)
backoff = 5
except asyncio.CancelledError:
LOG.debug("SignalRHub[%s]: loop cancelled — stopping", id)
return
except SignalRServerError as err:
LOG.warning(
"SignalRHub[%s]: server-initiated close; reconnecting in %ss — %s",
id,
backoff,
err,
)
except Exception as err: # pylint: disable=broad-except
LOG.warning(
"SignalRHub[%s]: connection error; reconnecting in %ss — %s",
id,
backoff,
err,
)
if not self.should_signalr_reconnect:
return
try:
await asyncio.sleep(backoff)
except asyncio.CancelledError:
LOG.debug("SignalRHub[%s]: sleep cancelled — stopping", id)
return
backoff = min(backoff * 2, 300)
async def cancel_task(self, task) -> None:
"""Cancel a task."""
LOG.debug("Cancelling task %s", task)
if task:
task.cancel()
try:
await task
except asyncio.CancelledError:
LOG.debug("Task %s successfully canceled", task)
task = None
return task
@property
def should_signalr_reconnect(self) -> bool:
"""Determine if a SignalR hub should reconnect when the connection is lost.
Currently only used to disable SignalR in the unit tests.
"""
return self._should_signalr_reconnect
@should_signalr_reconnect.setter
def should_signalr_reconnect(self, value: bool) -> None:
"""Set if SignalR hub should reconnect on disconnection."""
self._should_signalr_reconnect = value
async def async_update(self) -> None:
"""Update tarif periodically."""
if self.generate_energy_meters or self.track_unknown_sources:
self.check_tarif()
if self.track_unknown_sources:
self.handle_unknown_power()
def find_meter(self, hass):
"""Find the smart meter entity in Home Assistant."""
entity_registry_dict = {}
registry = hass.data.get("entity_registry")
if registry is None:
return entity_registry_dict
# ic-dev21: Get names of all entities
for entity_id, entity_entry in registry.entities.items():
entity_registry_dict[entity_id] = {
"name": entity_entry.entity_id,
}
sorted_entity_registry_dict = OrderedDict(sorted(entity_registry_dict.items()))
LOG.debug("Entities Ordered dict is %s", sorted_entity_registry_dict)
# Initialize empty list to put meter name into
filtered_names = []
# ic-dev21: Let's grab the meter from our dict
for entity_id, entity_data in sorted_entity_registry_dict.items():
if all(
substring in entity_data["name"] for substring in ["meter", "_power"]
):
filtered_names.append(entity_data["name"])
LOG.debug("Hilo Smart meter name is: %s", filtered_names)
# Format output to use in check_tarif
return ", ".join(filtered_names) if filtered_names else ""
def set_state(self, entity, state, new_attrs={}, keep_state=False, force=False):
"""Set the state of an entity."""
params = f"{entity=} {state=} {new_attrs=} {keep_state=}"
current = self._hass.states.get(entity)
if not current:
if not force:
LOG.warning(
"Unable to set state because there's no current: %s", params
)
return
attrs = {}
else:
attrs = dict(current.as_dict()["attributes"])
attrs["last_update"] = datetime.now()
attrs["hilo_update"] = True
attrs = {**attrs, **new_attrs}
if keep_state and current:
state = current.state
if "Cost" in attrs:
attrs["Cost"] = state
LOG.debug("Setting state %s current=%s attrs=%s", params, current, attrs)
self._hass.states.async_set(entity, state, attrs, force_update=force)
@property
def high_times(self):
"""Check if the current time is within high tariff periods."""
challenge_sensor = self._hass.states.get("sensor.defi_hilo")
if challenge_sensor is None:
LOG.warning("high_times check tarif challenge sensor not found")
return False
LOG.debug(
"high_times check tarif challenge sensor is %s", challenge_sensor.state
)
return challenge_sensor.state == "reduction"
def check_season(self):
"""Determine if we are using a winter or summer rate."""
current_month = datetime.now().month
LOG.debug("check_season current month is %s", current_month)
return current_month in [12, 1, 2, 3]
def check_tarif(self):
"""Determine which tarif to select depending on season and user-selected rate."""
if self.generate_energy_meters:
season = self.check_season()
LOG.debug("check_tarif current season state is %s", season)
tarif = "low"
base_sensor = f"sensor.{HILO_ENERGY_TOTAL}_low"
energy_used = self._hass.states.get(base_sensor)
if not energy_used:
LOG.warning("check_tarif: Unable to find state for %s", base_sensor)
return tarif
user_selected_plan_name = self.hq_plan_name
if user_selected_plan_name == "flex d":
if season:
plan_name = "flex d"
else:
plan_name = "rate d"
else:
plan_name = user_selected_plan_name
tarif_config = CONF_TARIFF.get(plan_name)
for tarif_name, rate in tarif_config.items():
if rate > 0 and tarif_name in ["low", "medium", "high"]:
if hasattr(self, "cost_sensors") and tarif_name in self.cost_sensors:
sensor = self.cost_sensors[tarif_name]
if sensor._cost != rate:
sensor._cost = rate
sensor.async_write_ha_state()
LOG.debug(
"check_tarif Updated %s sensor from %s to %s",
tarif_name,
sensor._cost,
rate,
)
current_cost = self._hass.states.get("sensor.hilo_rate_current")
if not current_cost:
LOG.warning(
"check_tarif: Unable to find state for sensor.hilo_rate_current"
)
return
try:
if float(energy_used.state) >= tarif_config.get("low_threshold"):
tarif = "medium"
except ValueError:
LOG.warning(
"Unable to restore a valid state of %s: %s",
base_sensor,
energy_used.state,
)
if tarif_config.get("high", 0) > 0 and self.high_times:
tarif = "high"
target_cost = self._hass.states.get(f"sensor.hilo_rate_{tarif}")
if not target_cost:
LOG.warning("check_tarif: sensor.hilo_rate_%s not available yet", tarif)
return
if target_cost.state != current_cost.state:
LOG.debug(
"check_tarif: Updating current cost, was %s now %s",
current_cost.state,
target_cost.state,
)
self.set_state("sensor.hilo_rate_current", target_cost.state)
if "current" in self.cost_sensors:
self.cost_sensors["current"]._cost = target_cost.state
LOG.debug(
"check_tarif: Current plan: %s Target Tarif: %s Energy used: %s Peak: %s",
plan_name,
tarif,
energy_used.state,
self.high_times,
)
# ic-dev21 : make sure the select for all meters still work by moving this here
for state in self._hass.states.async_all():
entity = state.entity_id
self.set_tarif(entity, state.state, tarif)
def handle_unknown_power(self):
"""Take care of the unknown source meter."""
known_power = 0
smart_meter = self.find_meter(self._hass)
LOG.debug("Smart meter used currently is: %s", smart_meter)
unknown_source_tracker = "sensor.unknown_source_tracker_power"
for state in self._hass.states.async_all():
entity = state.entity_id
if entity.endswith("hilo_rate_current"):
continue
if entity.endswith("_power") and entity not in [
unknown_source_tracker,
smart_meter,
]:
try:
known_power += int(float(state.state))
except ValueError:
pass
if not entity.endswith("_hilo_energy") or entity.endswith("_cost"):
continue
self.fix_utility_sensor(entity, state)
if self.track_unknown_sources:
total_power = self._hass.states.get(smart_meter)
try:
if known_power <= int(total_power.state):
unknown_power = int(total_power.state) - known_power
else:
unknown_power = 0
except ValueError:
unknown_power = known_power
LOG.warning(
f"value of total_power ({total_power} not initialized correctly)"
)
self.devices.parse_values_received(
[
{
"deviceId": 69420,
"locationId": self.devices.location_id,
"timeStampUTC": datetime.utcnow().isoformat(),
"attribute": "Power",
"value": unknown_power,
"valueType": "Watt",
}
]
)
LOG.debug(
"Currently in use: Total: %s Known sources: %s Unknown sources: %s",
total_power.state,
known_power,
unknown_power,
)
@callback
def fix_utility_sensor(self, entity, state):
"""Not sure why this doesn't get created with a proper device_class."""
current_state = state.as_dict()
attrs = current_state.get("attributes", {})
if entity.startswith("select.") or entity.find("hilo_rate") > 0:
return
if not attrs.get("source"):
LOG.debug("No source entity defined on %s: %s", entity, current_state)
return
parent_unit_state = self._hass.states.get(attrs.get("source"))
parent_unit = (
"kWh"
if parent_unit_state is None
else parent_unit_state.attributes.get("unit_of_measurement")
)
if not parent_unit:
LOG.warning("Unable to find state for parent unit: %s", current_state)
return
new_attrs = {
ATTR_UNIT_OF_MEASUREMENT: parent_unit, # note ic-dev21: now uses parent_unit directly
ATTR_DEVICE_CLASS: SensorDeviceClass.ENERGY,
}
if not all(a in attrs.keys() for a in new_attrs.keys()):
LOG.warning(
f"Fixing utility sensor: {entity} {current_state} new_attrs: {new_attrs}"
)
self.set_state(entity, None, new_attrs=new_attrs, keep_state=True)
@callback
def set_tarif(self, entity, current, new):
"""Set the tarif on the select entity if needed."""
if self.untarificated_devices and entity != f"select.{HILO_ENERGY_TOTAL}":
return
if entity.startswith("select.hilo_energy") and current != new:
LOG.debug(
"check_tarif: Changing tarif of %s from %s to %s", entity, current, new
)
context = Context()
data = {ATTR_OPTION: new, "entity_id": entity}
self._hass.async_create_task(
self._hass.services.async_call(
SELECT_DOMAIN, SERVICE_SELECT_OPTION, data, context=context
)
)
if (
entity.startswith("select.")
and entity.endswith("_hilo_energy")
and current != new
):
LOG.debug(
"check_tarif: Changing tarif of %s from %s to %s", entity, current, new
)
context = Context()
data = {ATTR_OPTION: new, "entity_id": entity}
self._hass.async_create_task(
self._hass.services.async_call(
SELECT_DOMAIN, SERVICE_SELECT_OPTION, data, context=context
)
)
@callback
def async_migrate_unique_id(
self, old_unique_id: str, new_unique_id: str | None, platform: str
) -> None:
"""Migrate legacy unique IDs to new format."""
assert new_unique_id is not None
LOG.debug(
"Checking if unique ID %s on %s needs to be migrated",
old_unique_id,
platform,
)
entity_registry = er.async_get(self._hass)
# async_get_entity_id wants the "HILO" domain
# in the platform field and the actual platform in the domain
# field for historical reasons since everything used to be
# PLATFORM.INTEGRATION instead of INTEGRATION.PLATFORM
if (
entity_id := entity_registry.async_get_entity_id(
platform, DOMAIN, old_unique_id
)
) is None:
LOG.debug("Unique ID %s does not need to be migrated", old_unique_id)
return
if new_entity_id := entity_registry.async_get_entity_id(
platform, DOMAIN, new_unique_id
):
LOG.debug(
(
"Unique ID %s is already in use by %s (system may have been"
" downgraded)"
),
new_unique_id,
new_entity_id,
)
return
LOG.debug(
"Migrating unique ID for entity %s (%s -> %s)",
entity_id,
old_unique_id,
new_unique_id,
)
entity_registry.async_update_entity(entity_id, new_unique_id=new_unique_id)
@callback
def handle_subscription_result(self, hilo_id: str) -> None:
"""Handle subscription result by notifying entities."""
async_dispatcher_send(self._hass, SIGNAL_UPDATE_ENTITY.format(hilo_id))