From 5e64a8edf64809aa9ff85fafaa38b581dc6970f7 Mon Sep 17 00:00:00 2001 From: davidabuch Date: Sun, 20 Sep 2026 15:50:08 -0700 Subject: [PATCH 1/4] Allow newer same-epoch cleanup verification evidence --- poolos/thermal_automatic_execution.py | 18 ++++++++++++++++-- 1 file changed, 16 insertions(+), 2 deletions(-) diff --git a/poolos/thermal_automatic_execution.py b/poolos/thermal_automatic_execution.py index 09024e4..26c7c2f 100644 --- a/poolos/thermal_automatic_execution.py +++ b/poolos/thermal_automatic_execution.py @@ -393,6 +393,15 @@ def reserve_circulation_candidate( def last_epoch_identity(self) -> str | None: return self._last_epoch_identity + def verification_reobservation_required(self) -> bool: + """Return whether newer same-epoch evidence may advance bounded cleanup.""" + + return ( + self.termination_attempt is not None + or self.cleanup_provenance is not None + or self.cleanup_attempt is not None + ) + def probe_execution_evidence(self) -> PoolTemperatureProbeExecutionEvidence | None: """Return positive in-memory probe provenance for the evaluator.""" @@ -619,8 +628,13 @@ async def process_epoch( command_delivery_performed=False, ) if frame.epoch_identity == self._last_epoch_identity: - assert self.assessment is not None - return self.assessment + if ( + self._last_epoch_at is None + or frame.observed_at <= self._last_epoch_at + or not self.verification_reobservation_required() + ): + assert self.assessment is not None + return self.assessment if self._last_epoch_at is not None and frame.observed_at < self._last_epoch_at: assert self.assessment is not None return self.assessment From 1de22b1e8cb715c57f983fc699801e6a4019b71c Mon Sep 17 00:00:00 2001 From: davidabuch Date: Sun, 20 Sep 2026 15:50:22 -0700 Subject: [PATCH 2/4] Resubmit newer same-epoch cleanup evidence --- .../poolos/thermal_automatic_runtime.py | 31 ++++++++++++++++--- 1 file changed, 26 insertions(+), 5 deletions(-) diff --git a/custom_components/poolos/thermal_automatic_runtime.py b/custom_components/poolos/thermal_automatic_runtime.py index 461823b..ab20206 100644 --- a/custom_components/poolos/thermal_automatic_runtime.py +++ b/custom_components/poolos/thermal_automatic_runtime.py @@ -566,11 +566,28 @@ def observe( self.spa_automatic_control.state.suppressed ), ) - if ( + same_epoch_reobservation = bool( self._latest_frame is not None and self._latest_frame.epoch_identity == frame.epoch_identity - ): + ) + if same_epoch_reobservation: + if ( + frame.observed_at <= self._latest_frame.observed_at + or not self.driver.verification_reobservation_required() + ): + return + # Newer verification evidence may advance an accepted + # termination/cleanup consequence without creating a new semantic + # authority epoch. State/policy identity is unchanged; only the + # authoritative chronology has advanced. + self._latest_frame = frame + if not self.driver.requested_enabled: + self.driver.note_disabled_epoch(frame) + self.coordinator.async_update_listeners() + return + self._schedule_if_idle() return + self._latest_frame = frame self.circulation_ownership.begin_epoch(frame.epoch_identity) self.authority.begin_automatic_thermal_epoch(frame.epoch_identity) @@ -850,9 +867,13 @@ def _task_done(self, task: asyncio.Task[object]) -> None: if self._unloaded or not self.driver.requested_enabled: return latest = self._latest_frame - if ( - latest is not None - and latest.epoch_identity != self.driver.last_epoch_identity + if latest is not None and ( + latest.epoch_identity != self.driver.last_epoch_identity + or ( + self.driver.verification_reobservation_required() + and self.driver.assessment is not None + and latest.observed_at > self.driver.assessment.evaluated_at + ) ): self._schedule_if_idle() From bc88e155c8525033cace0262531e0a309d6e7337 Mon Sep 17 00:00:00 2001 From: davidabuch Date: Sun, 20 Sep 2026 15:51:06 -0700 Subject: [PATCH 3/4] Regress same-epoch cleanup evidence resubmission --- ...ome_assistant_thermal_automatic_runtime.py | 86 +++++++++++++++++++ 1 file changed, 86 insertions(+) diff --git a/tests/test_home_assistant_thermal_automatic_runtime.py b/tests/test_home_assistant_thermal_automatic_runtime.py index 9952103..2d49ccf 100644 --- a/tests/test_home_assistant_thermal_automatic_runtime.py +++ b/tests/test_home_assistant_thermal_automatic_runtime.py @@ -69,6 +69,11 @@ class FakeDriver: pump_session_purpose: object | None = None probe_evidence: object | None = None cleanup_provenance: object | None = None + termination_attempt: object | None = None + cleanup_attempt: object | None = None + verification_required: bool = False + assessment: object | None = None + processed_at: list[datetime] = field(default_factory=list) def set_enabled(self, enabled: bool, **_: object) -> None: self.requested_enabled = enabled @@ -84,9 +89,14 @@ def restrictive_authority_changed(self, **_: object) -> None: async def process_epoch(self, frame: object, **_: object) -> None: self.processed.append(frame.epoch_identity) + self.processed_at.append(frame.observed_at) self.started.set() await self.release.wait() self.last_epoch_identity = frame.epoch_identity + self.assessment = SimpleNamespace(evaluated_at=frame.observed_at) + + def verification_reobservation_required(self) -> bool: + return self.verification_required def fail_closed(self, *, reason: str, **_: object) -> None: self.failed.append(reason) @@ -248,6 +258,82 @@ async def scenario() -> None: asyncio.run(scenario()) +def test_newer_same_epoch_cleanup_evidence_resubmits_without_new_authority_epoch() -> None: + async def scenario() -> None: + module = _load_module() + runtime, hass, authority, _, driver = _runtime(module) + runtime.set_enabled(True) + driver.release.set() + + runtime.observe( + _snapshot(NOW), + None, + _orchestration(NOW, "epoch-1"), + ) + await hass.tasks[0] + + assert driver.processed_at == [NOW] + assert authority.epochs == ["epoch-1"] + + runtime.observe( + _snapshot(NOW + timedelta(seconds=1)), + None, + _orchestration(NOW + timedelta(seconds=1), "epoch-1"), + ) + await asyncio.sleep(0) + assert driver.processed_at == [NOW] + assert authority.epochs == ["epoch-1"] + + driver.verification_required = True + runtime.observe( + _snapshot(NOW + timedelta(seconds=2)), + None, + _orchestration(NOW + timedelta(seconds=2), "epoch-1"), + ) + await hass.tasks[-1] + + assert driver.processed_at == [NOW, NOW + timedelta(seconds=2)] + assert authority.epochs == ["epoch-1"] + + asyncio.run(scenario()) + + +def test_same_epoch_cleanup_evidence_arriving_inflight_is_replayed_after_task() -> None: + async def scenario() -> None: + module = _load_module() + runtime, hass, authority, _, driver = _runtime(module) + runtime.set_enabled(True) + driver.verification_required = True + + runtime.observe( + _snapshot(NOW), + None, + _orchestration(NOW, "epoch-1"), + ) + first = hass.tasks[0] + await driver.started.wait() + + runtime.observe( + _snapshot(NOW + timedelta(seconds=2)), + None, + _orchestration(NOW + timedelta(seconds=2), "epoch-1"), + ) + + assert len(hass.tasks) == 1 + assert authority.epochs == ["epoch-1"] + + driver.release.set() + await first + await asyncio.sleep(0) + + assert len(hass.tasks) == 2 + await hass.tasks[1] + assert driver.processed_at == [NOW, NOW + timedelta(seconds=2)] + assert authority.epochs == ["epoch-1"] + + asyncio.run(scenario()) + + def test_owned_priming_and_probe_sessions_actively_reobserve_unchanged_native_evidence() -> None: async def scenario(*, priming: bool) -> None: module = _load_module() From d554f87486691186c2695742c63301c5383f92ca Mon Sep 17 00:00:00 2001 From: davidabuch Date: Sun, 20 Sep 2026 15:51:16 -0700 Subject: [PATCH 4/4] Regress newer same-epoch termination verification --- tests/test_thermal_automatic_execution.py | 27 +++++++++++++++++++++++ 1 file changed, 27 insertions(+) diff --git a/tests/test_thermal_automatic_execution.py b/tests/test_thermal_automatic_execution.py index a916a87..fda4b81 100644 --- a/tests/test_thermal_automatic_execution.py +++ b/tests/test_thermal_automatic_execution.py @@ -5263,6 +5263,33 @@ def test_accepted_termination_delivery_needs_a_later_authoritative_epoch() -> No assert len(factory.delivery.calls) == 4 +def test_newer_same_epoch_source_off_evidence_advances_termination_verification() -> None: + orchestrator, driver, factory, ending, _ = ( + _driver_awaiting_source_off_verification() + ) + confirmed = _frame( + orchestrator, + NOW + timedelta(seconds=66), + pool_active=True, + pump_rpm=3000, + configured_rpm=3000, + pool_heater="00000", + mode=ThermalRequestedMode.OFF, + ) + same_epoch_confirmed = replace( + confirmed, + epoch_identity=ending.epoch_identity, + ) + + result = asyncio.run( + driver.process_epoch(same_epoch_confirmed, delivery_factory=factory) + ) + + assert result.state is ThermalAutomaticDriverState.CONVERGED + assert driver.termination_attempt is None + assert len(factory.delivery.calls) == 4 + + def test_native_gas_does_not_verify_accepted_source_off_delivery() -> None: orchestrator, driver, factory, _, _ = _driver_awaiting_source_off_verification() still_gas = _frame(