Skip to content
Open
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
31 changes: 26 additions & 5 deletions custom_components/poolos/thermal_automatic_runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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()

Expand Down
18 changes: 16 additions & 2 deletions poolos/thermal_automatic_execution.py
Original file line number Diff line number Diff line change
Expand Up @@ -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."""

Expand Down Expand Up @@ -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
Expand Down
86 changes: 86 additions & 0 deletions tests/test_home_assistant_thermal_automatic_runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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)
Expand Down Expand Up @@ -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()
Expand Down
27 changes: 27 additions & 0 deletions tests/test_thermal_automatic_execution.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down
Loading