diff --git a/homeassistant/components/recorder/table_managers/states.py b/homeassistant/components/recorder/table_managers/states.py index b1031a66b680..3a7ec739d417 100644 --- a/homeassistant/components/recorder/table_managers/states.py +++ b/homeassistant/components/recorder/table_managers/states.py @@ -81,14 +81,14 @@ class StatesManager: self._last_reported.clear() def reset(self) -> None: - """Reset after the database has been reset or changed. + """Reset pending and last-committed caches after the event session is closed. This call is not thread-safe and must be called from the recorder thread. """ self._last_committed_id.clear() self._pending.clear() - self._oldest_ts = None + # _oldest_ts is a database aggregate, not session state. def load_from_db(self, session: Session) -> None: """Update the cache. diff --git a/tests/components/recorder/table_managers/test_states.py b/tests/components/recorder/table_managers/test_states.py new file mode 100644 index 000000000000..f76e14c70335 --- /dev/null +++ b/tests/components/recorder/table_managers/test_states.py @@ -0,0 +1,25 @@ +"""The tests for the recorder states manager.""" + +from homeassistant.components.recorder.db_schema import States +from homeassistant.components.recorder.table_managers.states import StatesManager + + +def test_reset_preserves_oldest_ts() -> None: + """Test reset() keeps oldest_ts.""" + manager = StatesManager() + manager.add_pending("sensor.test", States(last_updated_ts=123.0)) + + manager.reset() + + assert manager.oldest_ts == 123.0 + assert manager.pop_pending("sensor.test") is None + + +def test_add_pending_after_reset_does_not_reseed_oldest_ts() -> None: + """Test add_pending after reset does not seed oldest_ts from the new state.""" + manager = StatesManager() + manager.add_pending("sensor.test", States(last_updated_ts=10.0)) + manager.reset() + manager.add_pending("sensor.test", States(last_updated_ts=99.0)) + + assert manager.oldest_ts == 10.0 diff --git a/tests/components/recorder/test_session_recovery.py b/tests/components/recorder/test_session_recovery.py new file mode 100644 index 000000000000..1364c6d12e9d --- /dev/null +++ b/tests/components/recorder/test_session_recovery.py @@ -0,0 +1,94 @@ +"""Tests for recorder session recovery.""" + +from datetime import timedelta +from typing import Any +from unittest.mock import patch + +from freezegun.api import FrozenDateTimeFactory +import pytest +from sqlalchemy.exc import SQLAlchemyError + +from homeassistant.components.recorder import get_instance, history +from homeassistant.components.recorder.db_schema import States +from homeassistant.components.recorder.util import session_scope +from homeassistant.core import HomeAssistant, State +from homeassistant.util import dt as dt_util + +from .common import async_wait_recording_done + +from tests.typing import RecorderInstanceContextManager + + +@pytest.fixture +async def mock_recorder_before_hass( + async_test_recorder: RecorderInstanceContextManager, +) -> None: + """Set up recorder.""" + + +@pytest.mark.usefixtures("recorder_mock") +async def test_oldest_ts_preserved_after_sqlalchemy_recovery( + hass: HomeAssistant, + caplog: pytest.LogCaptureFixture, + freezer: FrozenDateTimeFactory, +) -> None: + """Test oldest_ts is not reseeded from now after session recovery.""" + instance = get_instance(hass) + entity_id = "sensor.oldest_ts_recovery" + stable_entity_id = "sensor.oldest_ts_stable" + + hass.states.async_set(stable_entity_id, "on") + hass.states.async_set(entity_id, "before") + await async_wait_recording_done(hass) + + oldest_ts_before_recovery = instance.states_manager.oldest_ts + assert oldest_ts_before_recovery is not None + + freezer.tick(timedelta(seconds=1)) + start_time = dt_util.utcnow() + + event_session = instance.event_session + assert event_session is not None + + def _throw_if_state_in_session(*args: Any, **kwargs: Any) -> None: + for obj in event_session: + if isinstance(obj, States): + raise SQLAlchemyError( + "insert the state", "fake params", "forced to fail" + ) + + with ( + patch.object( + event_session, + "flush", + side_effect=_throw_if_state_in_session, + ), + ): + hass.states.async_set(entity_id, "fail") + await async_wait_recording_done(hass) + + assert "SQLAlchemyError error processing task" in caplog.text + + freezer.tick(timedelta(seconds=1)) + hass.states.async_set(entity_id, "after") + await async_wait_recording_done(hass) + + assert instance.states_manager.oldest_ts == oldest_ts_before_recovery + + with session_scope(hass=hass, read_only=True) as session: + recorded_states = {row.state for row in session.query(States)} + assert "fail" not in recorded_states + + hist = history.get_significant_states( + hass, + start_time, + dt_util.utcnow(), + entity_ids=[entity_id, stable_entity_id], + include_start_time_state=True, + ) + start_state = hist[entity_id][0] + assert isinstance(start_state, State) + assert start_state.state == "before" + stable_start_state = hist[stable_entity_id][0] + assert isinstance(stable_start_state, State) + assert stable_start_state.state == "on"