diff --git a/custom_components/spa_pool/binary_sensor.py b/custom_components/spa_pool/binary_sensor.py index 586c4bc..0a3d37a 100644 --- a/custom_components/spa_pool/binary_sensor.py +++ b/custom_components/spa_pool/binary_sensor.py @@ -55,7 +55,7 @@ def _state_available( ) -> bool: """Return whether a current decoded state is available.""" - return client.available and state is not None + return client.state_available and state is not None BINARY_SENSOR_DESCRIPTIONS: Final[ @@ -66,14 +66,20 @@ def _state_available( translation_key="status_stream", device_class=BinarySensorDeviceClass.CONNECTIVITY, entity_category=EntityCategory.DIAGNOSTIC, - value_fn=lambda client, state: client.available, + value_fn=lambda client, state: client.transport_available, attributes_fn=lambda client, state: { "tcp_connected": client.connected, + "state_available": client.state_available, "last_valid_message": ( client.last_message_at.isoformat() if client.last_message_at is not None else None ), + "last_status_message": ( + client.last_status_at.isoformat() + if client.last_status_at is not None + else None + ), }, availability_fn=_always_available, ), @@ -252,7 +258,7 @@ def available(self) -> bool: availability_fn = self.entity_description.availability_fn if availability_fn is None: - return self._client.available + return self._client.transport_available return availability_fn(self._client, self._client.state) diff --git a/custom_components/spa_pool/button.py b/custom_components/spa_pool/button.py index f9aea5b..25d8bbf 100644 --- a/custom_components/spa_pool/button.py +++ b/custom_components/spa_pool/button.py @@ -225,7 +225,13 @@ def available(self) -> bool: if not self.entity_description.requires_stream: return True - return self._client.available + if self.entity_description.action in ( + SpaPoolButtonAction.SYNC_CLOCK, + SpaPoolButtonAction.CLEAR_REMINDER, + ): + return self._client.state_available + + return self._client.transport_available @override async def async_press(self) -> None: @@ -305,7 +311,7 @@ def _require_state(self) -> SpaState: """Return current state or raise a user-visible communication error.""" state = self._client.state - if state is None or not self._client.available: + if state is None or not self._client.state_available: raise HomeAssistantError("Spa status stream is unavailable") return state @@ -344,7 +350,7 @@ def __init__( def available(self) -> bool: """Return whether the bridge currently accepts commands.""" - return self._client.available + return self._client.transport_available @override async def async_press(self) -> None: diff --git a/custom_components/spa_pool/client.py b/custom_components/spa_pool/client.py index ba60769..342d7e7 100644 --- a/custom_components/spa_pool/client.py +++ b/custom_components/spa_pool/client.py @@ -24,6 +24,7 @@ _CONNECT_TIMEOUT = 10.0 _FIRST_FRAME_TIMEOUT = 20.0 _STREAM_STALE_TIMEOUT = 15.0 +_STATUS_STALE_TIMEOUT = 10.0 _WRITE_TIMEOUT = 5.0 _COMMAND_CONNECTION_TIMEOUT = 10.0 _CLOSE_TIMEOUT = 5.0 @@ -87,7 +88,8 @@ def __init__( self._device_configuration_condition = asyncio.Condition() self._first_frame_event = asyncio.Event() - self._available_event = asyncio.Event() + self._transport_available_event = asyncio.Event() + self._state_available_event = asyncio.Event() self._listeners: set[Listener] = set() self._message_listeners: set[Listener] = set() @@ -109,12 +111,15 @@ def __init__( self._ready_revision = 0 self._connected = False - self._available = False + self._transport_available = False + self._state_available = False + self._status_stale_handle: asyncio.TimerHandle | None = None self._current_connection_received_frame = False self._last_frame: bytes | None = None self._last_message_type: str | None = None self._last_message_at: datetime | None = None + self._last_status_at: datetime | None = None self._last_error: str | None = None self._connection_attempts = 0 @@ -182,9 +187,25 @@ def connected(self) -> bool: @property def available(self) -> bool: + """Return whether a recent regular status frame is available. + + Kept as the state-availability interface used by existing entities. + Transport-only diagnostics should use ``transport_available``. + """ + + return self._state_available + + @property + def transport_available(self) -> bool: """Return whether the current connection has supplied a valid frame.""" - return self._available + return self._transport_available + + @property + def state_available(self) -> bool: + """Return whether spa state is backed by a recent status frame.""" + + return self._state_available @property def last_message_at(self) -> datetime | None: @@ -192,6 +213,12 @@ def last_message_at(self) -> datetime | None: return self._last_message_at + @property + def last_status_at(self) -> datetime | None: + """Return the UTC time of the most recent regular status frame.""" + + return self._last_status_at + @property def last_frame(self) -> bytes | None: """Return the most recent valid raw frame.""" @@ -269,7 +296,8 @@ async def async_stop(self) -> None: await task await self._async_close_connection() - self._set_available(False) + self._set_transport_available(False) + self._set_state_available(False) async def async_reconnect(self) -> None: """Manually rebuild the stream without restarting Home Assistant.""" @@ -292,7 +320,7 @@ async def async_send_frame(self, frame: bytes) -> None: """ async with self._command_lock: - await self._async_wait_until_available() + await self._async_wait_until_transport_available() await self._async_write_frame(frame) async def async_send_and_wait( @@ -305,7 +333,7 @@ async def async_send_and_wait( """Queue a command and confirm it from a newer status message.""" async with self._command_lock: - await self._async_wait_until_available() + await self._async_wait_until_state_available() revision = self._state_revision await self._async_write_frame(frame) @@ -329,7 +357,7 @@ async def async_send_and_wait_for_fault( """Queue a fault-log request and await a newer fault entry.""" async with self._command_lock: - await self._async_wait_until_available() + await self._async_wait_until_transport_available() revision = self._fault_revision await self._async_write_frame(frame) @@ -352,7 +380,7 @@ async def async_send_and_wait_for_device_configuration( """Queue the panel request and await a newer ``0A BF 2E`` frame.""" async with self._command_lock: - await self._async_wait_until_available() + await self._async_wait_until_transport_available() revision = self._device_configuration_revision await self._async_write_frame(frame) @@ -444,12 +472,23 @@ async def async_wait_for_fault( await self._fault_condition.wait() - async def _async_wait_until_available(self) -> None: + async def _async_wait_until_transport_available(self) -> None: + """Wait briefly for a checksum-valid stream.""" + + try: + async with asyncio.timeout(_COMMAND_CONNECTION_TIMEOUT): + await self._transport_available_event.wait() + except TimeoutError as err: + raise SpaPoolNotConnectedError( + f"Spa bridge {self.host}:{self.port} is unavailable" + ) from err + + async def _async_wait_until_state_available(self) -> None: """Wait briefly for a usable status stream.""" try: async with asyncio.timeout(_COMMAND_CONNECTION_TIMEOUT): - await self._available_event.wait() + await self._state_available_event.wait() except TimeoutError as err: raise SpaPoolNotConnectedError( f"Spa bridge {self.host}:{self.port} is unavailable" @@ -508,7 +547,8 @@ async def _async_run(self) -> None: ) finally: await self._async_close_connection() - self._set_available(False) + self._set_transport_available(False) + self._set_state_available(False) self._protocol.reset_stream() delay_before_retry = reconnect_delay @@ -605,10 +645,9 @@ async def _async_process_update(self, update: Any) -> None: self._last_message_at = datetime.now(UTC) self._last_error = None - first_available_frame = not self._available - if first_available_frame: - self._available = True - self._available_event.set() + first_transport_frame = not self._transport_available + if first_transport_frame: + self._set_transport_available(True) self._first_frame_event.set() @@ -647,6 +686,11 @@ async def _async_process_update(self, update: Any) -> None: self._valid_status_count += 1 self._state_condition.notify_all() + if regular_status_received: + self._last_status_at = self._last_message_at + self._set_state_available(True) + self._arm_status_stale_timer() + fault = update.fault if fault is not None: async with self._fault_condition: @@ -657,7 +701,7 @@ async def _async_process_update(self, update: Any) -> None: self._notify_fault_listeners(fault) if ( - first_available_frame + first_transport_frame or (state_received and self._state != previous_state) ): self._notify_listeners() @@ -678,22 +722,61 @@ async def _async_close_connection(self) -> None: async with asyncio.timeout(_CLOSE_TIMEOUT): await writer.wait_closed() - def _set_available(self, available: bool) -> None: - """Set protocol availability and notify entities when it changes.""" + def _set_transport_available(self, available: bool) -> None: + """Set transport availability and notify entities when it changes.""" + + if self._transport_available == available: + if not available: + self._transport_available_event.clear() + return + + self._transport_available = available + if available: + self._transport_available_event.set() + else: + self._transport_available_event.clear() + self._notify_listeners() + + def _set_state_available(self, available: bool) -> None: + """Set status freshness and notify entities when it changes.""" - if self._available == available: + if self._state_available == available: if not available: - self._available_event.clear() + self._state_available_event.clear() + self._cancel_status_stale_timer() return - self._available = available + self._state_available = available if available: - self._available_event.set() + self._state_available_event.set() else: - self._available_event.clear() + self._state_available_event.clear() + self._cancel_status_stale_timer() self._notify_listeners() + def _arm_status_stale_timer(self) -> None: + """Expire state independently if regular status traffic stops.""" + + self._cancel_status_stale_timer() + self._status_stale_handle = asyncio.get_running_loop().call_later( + _STATUS_STALE_TIMEOUT, + self._handle_status_stale, + ) + + def _cancel_status_stale_timer(self) -> None: + """Cancel the current state-freshness deadline.""" + + if self._status_stale_handle is not None: + self._status_stale_handle.cancel() + self._status_stale_handle = None + + def _handle_status_stale(self) -> None: + """Mark retained state stale without dropping a healthy transport.""" + + self._status_stale_handle = None + self._set_state_available(False) + def _notify_listeners(self) -> None: """Call listeners without allowing one entity to break the stream.""" @@ -726,7 +809,9 @@ def diagnostics(self) -> dict[str, Any]: return { "connected": self._connected, - "available": self._available, + "available": self._state_available, + "transport_available": self._transport_available, + "state_available": self._state_available, "state_revision": self._state_revision, "fault_revision": self._fault_revision, "device_configuration": { @@ -764,6 +849,11 @@ def diagnostics(self) -> dict[str, Any]: else None ), "last_message_type": self._last_message_type, + "last_status_at": ( + self._last_status_at.isoformat() + if self._last_status_at is not None + else None + ), "last_error": self._last_error, "connection_attempts": self._connection_attempts, "reconnect_count": self._reconnect_count, diff --git a/custom_components/spa_pool/climate.py b/custom_components/spa_pool/climate.py index 6e2d656..6d95eb5 100644 --- a/custom_components/spa_pool/climate.py +++ b/custom_components/spa_pool/climate.py @@ -96,7 +96,7 @@ def __init__(self, entry: SpaPoolConfigEntry) -> None: def available(self) -> bool: """Return whether a current status stream is available.""" - return self._client.available and self._client.state is not None + return self._client.state_available and self._client.state is not None @property @override @@ -221,7 +221,7 @@ async def async_set_temperature(self, **kwargs: Any) -> None: return state = self._client.state - if state is None or not self._client.available: + if state is None or not self._client.state_available: raise HomeAssistantError("Spa status stream is unavailable") target = float(temperature) diff --git a/custom_components/spa_pool/event.py b/custom_components/spa_pool/event.py index f6fbed3..053cd4d 100644 --- a/custom_components/spa_pool/event.py +++ b/custom_components/spa_pool/event.py @@ -63,7 +63,7 @@ def __init__(self, entry: SpaPoolConfigEntry) -> None: def available(self) -> bool: """Return whether the spa status stream is available.""" - return self._client.available + return self._client.transport_available @override async def async_added_to_hass(self) -> None: diff --git a/custom_components/spa_pool/sensor.py b/custom_components/spa_pool/sensor.py index 1831065..08a2ccf 100644 --- a/custom_components/spa_pool/sensor.py +++ b/custom_components/spa_pool/sensor.py @@ -61,7 +61,7 @@ def _state_available( ) -> bool: """Return whether a current decoded state is available.""" - return client.available and state is not None + return client.state_available and state is not None def _message_available( @@ -413,7 +413,7 @@ def available(self) -> bool: if availability_fn is not None: return availability_fn(self._client, state) - return self._client.available + return self._client.state_available @property @override @@ -523,7 +523,10 @@ def _decoded_value(self) -> SpaIntEnum | None: def available(self) -> bool: """Return whether this field is present in the current status frame.""" - return self._client.available and self._decoded_value() is not None + return ( + self._client.state_available + and self._decoded_value() is not None + ) @property @override diff --git a/tests/test_client_availability.py b/tests/test_client_availability.py new file mode 100644 index 0000000..eb87de2 --- /dev/null +++ b/tests/test_client_availability.py @@ -0,0 +1,110 @@ +"""Tests for independent transport and status availability.""" + +from __future__ import annotations + +import asyncio +from pathlib import Path +import sys +from types import ModuleType, SimpleNamespace +import unittest + + +# Import the Home Assistant-independent client without executing the +# integration package's Home Assistant entrypoint. +ROOT = Path(__file__).parents[1] +CUSTOM_COMPONENTS = ROOT / "custom_components" +SPA_POOL = CUSTOM_COMPONENTS / "spa_pool" + +custom_components = ModuleType("custom_components") +custom_components.__path__ = [str(CUSTOM_COMPONENTS)] +spa_pool = ModuleType("custom_components.spa_pool") +spa_pool.__path__ = [str(SPA_POOL)] +sys.modules.setdefault("custom_components", custom_components) +sys.modules.setdefault("custom_components.spa_pool", spa_pool) + +from custom_components.spa_pool import client as client_module # noqa: E402 +from custom_components.spa_pool.client import SpaPoolClient # noqa: E402 + + +def _update(message_type: str, *, state: object | None = None) -> object: + """Build the subset of a decoded update used by the client.""" + + return SimpleNamespace( + message_type=message_type, + raw_frame=b"\x7e\x00\x7e", + payload=b"", + state=state, + fault=None, + ) + + +class SpaPoolAvailabilityTests(unittest.IsolatedAsyncioTestCase): + """Exercise transport and state availability as separate signals.""" + + async def asyncTearDown(self) -> None: + """Allow any cancelled timer callbacks to be discarded.""" + + await asyncio.sleep(0) + + async def test_non_status_frame_only_makes_transport_available(self) -> None: + """Bus-management traffic must not validate retained spa state.""" + + client = SpaPoolClient("spa.local", 8899) + + await client._async_process_update(_update("ready_to_send")) + + self.assertTrue(client.transport_available) + self.assertFalse(client.state_available) + self.assertFalse(client.available) + self.assertIsNone(client.last_status_at) + + async def test_status_freshness_expires_without_dropping_transport( + self, + ) -> None: + """Status state becomes unavailable while valid transport remains.""" + + client = SpaPoolClient("spa.local", 8899) + retained_state = object() + original_timeout = client_module._STATUS_STALE_TIMEOUT + client_module._STATUS_STALE_TIMEOUT = 0.01 + self.addCleanup( + setattr, + client_module, + "_STATUS_STALE_TIMEOUT", + original_timeout, + ) + + await client._async_process_update( + _update("status_update", state=retained_state) + ) + + self.assertTrue(client.transport_available) + self.assertTrue(client.state_available) + self.assertIs(client.state, retained_state) + self.assertIsNotNone(client.last_status_at) + + await asyncio.sleep(0.02) + + self.assertTrue(client.transport_available) + self.assertFalse(client.state_available) + self.assertIs(client.state, retained_state) + + async def test_disconnect_clears_both_availability_signals(self) -> None: + """A disconnect invalidates transport and status immediately.""" + + client = SpaPoolClient("spa.local", 8899) + await client._async_process_update( + _update("status_update", state=object()) + ) + + client._set_transport_available(False) + client._set_state_available(False) + + self.assertFalse(client.transport_available) + self.assertFalse(client.state_available) + self.assertFalse(client._transport_available_event.is_set()) + self.assertFalse(client._state_available_event.is_set()) + + +if __name__ == "__main__": + unittest.main()