mirror of
https://github.com/home-assistant/core.git
synced 2026-10-06 22:38:02 -04:00
Fix stale Nibe Heat Pump values at the end of polling (#184222)
This commit is contained in:
@@ -190,6 +190,7 @@ class CoilCoordinator(ContextCoordinator[dict[int, CoilData], int]):
|
||||
|
||||
async def _async_update_data_internal(self) -> dict[int, CoilData]:
|
||||
result: dict[int, CoilData] = {}
|
||||
read_error: ReadException | None = None
|
||||
|
||||
def _get_coils() -> Iterable[Coil]:
|
||||
for address in sorted(self.context_callbacks.keys()):
|
||||
@@ -208,14 +209,20 @@ class CoilCoordinator(ContextCoordinator[dict[int, CoilData], int]):
|
||||
try:
|
||||
async for data in self.connection.read_coils(_get_coils()):
|
||||
result[data.coil.address] = data
|
||||
self.seed.pop(data.coil.address, None)
|
||||
except ReadException as exception:
|
||||
if not result:
|
||||
raise UpdateFailed(f"Failed to update: {exception}") from exception
|
||||
self.logger.debug(
|
||||
"Some coils failed to update, and may be unsupported: %s", exception
|
||||
)
|
||||
read_error = exception
|
||||
|
||||
# Preserve broadcasts received while the polling batch was running.
|
||||
for address in self.context_callbacks:
|
||||
if seed := self.seed.pop(address, None):
|
||||
result[address] = seed
|
||||
|
||||
if read_error is not None:
|
||||
if not result:
|
||||
raise UpdateFailed(f"Failed to update: {read_error}") from read_error
|
||||
self.logger.debug(
|
||||
"Some coils failed to update, and may be unsupported: %s", read_error
|
||||
)
|
||||
return result
|
||||
|
||||
@override
|
||||
|
||||
@@ -97,6 +97,156 @@ async def test_pushed_update(
|
||||
assert hass.states.get(entity_id) == snapshot(name="4. final values")
|
||||
|
||||
|
||||
@pytest.mark.usefixtures("entity_registry_enabled_by_default")
|
||||
@pytest.mark.parametrize(
|
||||
("seeded_address", "read_address"),
|
||||
[
|
||||
pytest.param(40031, 40035, id="broadcast-after-seeded-value-captured"),
|
||||
pytest.param(40035, 40031, id="broadcast-before-read-response-consumed"),
|
||||
],
|
||||
)
|
||||
async def test_pushed_update_during_refresh(
|
||||
hass: HomeAssistant,
|
||||
coils: dict[int, float],
|
||||
mock_connection: MockConnection,
|
||||
seeded_address: int,
|
||||
read_address: int,
|
||||
) -> None:
|
||||
"""Test that a completed polling batch preserves newer pushed values."""
|
||||
entity_id = "number.heating_offset_climate_system_1_40031"
|
||||
coils[40031] = 10
|
||||
coils[40035] = 20
|
||||
|
||||
entry = await async_add_model(hass, Model.S320)
|
||||
coordinator = entry.runtime_data
|
||||
mock_connection.mock_coil_update(seeded_address, 20)
|
||||
|
||||
async def read_coil(coil: Coil, timeout: float = 0) -> CoilData:
|
||||
assert coil.address == read_address
|
||||
data = CoilData(coil, 20)
|
||||
# NibeGW publishes read replies before the polling iterator consumes them.
|
||||
mock_connection.heatpump.notify_coil_update(data)
|
||||
mock_connection.mock_coil_update(40031, 21)
|
||||
mock_connection.mock_coil_update(40031, 22)
|
||||
assert hass.states.get(entity_id).state == "22.0"
|
||||
return data
|
||||
|
||||
with patch.object(mock_connection, "read_coil", side_effect=read_coil):
|
||||
await coordinator.async_refresh()
|
||||
|
||||
assert hass.states.get(entity_id).state == "22.0"
|
||||
assert coordinator.data[40031].value == 22
|
||||
assert coordinator.data[40035].value == 20
|
||||
|
||||
# Without more broadcasts, the following refresh must read both coils again.
|
||||
coils[40031] = 30
|
||||
coils[40035] = 40
|
||||
await coordinator.async_refresh()
|
||||
|
||||
assert hass.states.get(entity_id).state == "30.0"
|
||||
assert coordinator.data[40035].value == 40
|
||||
|
||||
|
||||
@pytest.mark.usefixtures("entity_registry_enabled_by_default")
|
||||
async def test_pushed_update_during_partial_refresh(
|
||||
hass: HomeAssistant,
|
||||
coils: dict[int, float | None],
|
||||
mock_connection: MockConnection,
|
||||
) -> None:
|
||||
"""Test that a partial polling batch preserves broadcasts for failed coils."""
|
||||
entity_id = "number.heating_offset_climate_system_1_40031"
|
||||
coils[30002] = 10
|
||||
coils[40031] = 10
|
||||
coils[40035] = 20
|
||||
|
||||
entry = await async_add_model(hass, Model.S320)
|
||||
coordinator = entry.runtime_data
|
||||
assert 30002 not in coordinator.context_callbacks
|
||||
coils[40031] = None
|
||||
read_coil_original = mock_connection.read_coil
|
||||
|
||||
async def read_coil(coil: Coil, timeout: float = 0) -> CoilData:
|
||||
data = await read_coil_original(coil, timeout)
|
||||
assert coil.address == 40035
|
||||
# The earlier read failed, but its broadcast arrives during this read.
|
||||
mock_connection.mock_coil_update(40031, 22)
|
||||
mock_connection.mock_coil_update(30002, 30)
|
||||
mock_connection.heatpump.notify_coil_update(data)
|
||||
assert hass.states.get(entity_id).state == "22.0"
|
||||
return data
|
||||
|
||||
with patch.object(mock_connection, "read_coil", side_effect=read_coil):
|
||||
await coordinator.async_refresh()
|
||||
|
||||
assert hass.states.get(entity_id).state == "22.0"
|
||||
assert coordinator.data[40031].value == 22
|
||||
assert coordinator.data[40035].value == 20
|
||||
assert 30002 not in coordinator.data
|
||||
assert coordinator.last_update_success
|
||||
|
||||
# Without more broadcasts, the following refresh must read both coils again.
|
||||
coils[40031] = 30
|
||||
coils[40035] = 40
|
||||
await coordinator.async_refresh()
|
||||
|
||||
assert hass.states.get(entity_id).state == "30.0"
|
||||
assert coordinator.data[40035].value == 40
|
||||
|
||||
|
||||
@pytest.mark.usefixtures("entity_registry_enabled_by_default")
|
||||
@pytest.mark.parametrize(
|
||||
("broadcast_addresses", "expected_success", "expected_state"),
|
||||
[
|
||||
pytest.param((40031,), True, "22.0", id="active-broadcast"),
|
||||
pytest.param((30002,), False, "unavailable", id="inactive-broadcast"),
|
||||
pytest.param((), False, "unavailable", id="no-broadcast"),
|
||||
],
|
||||
)
|
||||
async def test_pushed_update_during_failed_refresh(
|
||||
hass: HomeAssistant,
|
||||
coils: dict[int, float | None],
|
||||
mock_connection: MockConnection,
|
||||
broadcast_addresses: tuple[int, ...],
|
||||
expected_success: bool,
|
||||
expected_state: str,
|
||||
) -> None:
|
||||
"""Test availability when all polling reads fail during a broadcast."""
|
||||
entity_id = "number.heating_offset_climate_system_1_40031"
|
||||
coils[30002] = 10
|
||||
coils[40031] = 10
|
||||
coils[40035] = 20
|
||||
|
||||
entry = await async_add_model(hass, Model.S320)
|
||||
coordinator = entry.runtime_data
|
||||
assert 30002 not in coordinator.context_callbacks
|
||||
coils[40031] = None
|
||||
coils[40035] = None
|
||||
read_coil_original = mock_connection.read_coil
|
||||
|
||||
async def read_coil(coil: Coil, timeout: float = 0) -> CoilData:
|
||||
try:
|
||||
return await read_coil_original(coil, timeout)
|
||||
finally:
|
||||
for address in broadcast_addresses:
|
||||
mock_connection.mock_coil_update(address, 22)
|
||||
|
||||
with patch.object(mock_connection, "read_coil", side_effect=read_coil):
|
||||
await coordinator.async_refresh()
|
||||
|
||||
assert coordinator.last_update_success is expected_success
|
||||
assert hass.states.get(entity_id).state == expected_state
|
||||
|
||||
# The following refresh must recover with fresh reads, without more broadcasts.
|
||||
coils[40031] = 30
|
||||
coils[40035] = 40
|
||||
await coordinator.async_refresh()
|
||||
|
||||
assert coordinator.last_update_success
|
||||
assert hass.states.get(entity_id).state == "30.0"
|
||||
assert coordinator.data[40035].value == 40
|
||||
assert 30002 not in coordinator.data
|
||||
|
||||
|
||||
@pytest.mark.usefixtures("entity_registry_enabled_by_default")
|
||||
async def test_shutdown(
|
||||
hass: HomeAssistant,
|
||||
|
||||
Reference in New Issue
Block a user