diff --git a/homeassistant/components/nibe_heatpump/coordinator.py b/homeassistant/components/nibe_heatpump/coordinator.py index 3758963c49a1..fe1f2173c6b7 100644 --- a/homeassistant/components/nibe_heatpump/coordinator.py +++ b/homeassistant/components/nibe_heatpump/coordinator.py @@ -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 diff --git a/tests/components/nibe_heatpump/test_coordinator.py b/tests/components/nibe_heatpump/test_coordinator.py index ae9b7f386d4a..18f05575751c 100644 --- a/tests/components/nibe_heatpump/test_coordinator.py +++ b/tests/components/nibe_heatpump/test_coordinator.py @@ -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,