Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 9 additions & 3 deletions custom_components/spa_pool/binary_sensor.py
Original file line number Diff line number Diff line change
Expand Up @@ -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[
Expand All @@ -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,
),
Expand Down Expand Up @@ -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)

Expand Down
12 changes: 9 additions & 3 deletions custom_components/spa_pool/button.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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:
Expand Down
138 changes: 114 additions & 24 deletions custom_components/spa_pool/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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()
Expand All @@ -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
Expand Down Expand Up @@ -182,16 +187,38 @@ 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:
"""Return the UTC time of the most recent valid protocol frame."""

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."""
Expand Down Expand Up @@ -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."""
Expand All @@ -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(
Expand All @@ -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)

Expand All @@ -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)

Expand All @@ -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)

Expand Down Expand Up @@ -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"
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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()

Expand Down Expand Up @@ -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:
Expand All @@ -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()
Expand All @@ -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."""

Expand Down Expand Up @@ -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": {
Expand Down Expand Up @@ -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,
Expand Down
4 changes: 2 additions & 2 deletions custom_components/spa_pool/climate.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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)
Expand Down
2 changes: 1 addition & 1 deletion custom_components/spa_pool/event.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
Loading