fix: import result follows the commit; dropped route runs reach the store (#495)

Apply re-validated the preview token after both halves were durable, so a
TTL that lapsed during the write, or an eviction by a newer preview of the
same user, answered "preview expired" for a plan that was already replaced
and withheld both update events. The token is now simply spent after the
commit; validity is decided once, on entry under write_lock.

Route runs dropped by drop_unknown_routes never reached the store unless a
marker happened to be orphaned in the same pass; an idle robot kept them in
memory only and a restart brought them back. The drop now happens under
_refresh_lock, leaves with the orphan transaction or with an immediate
write of its own, and rolls back like #335 when the store refuses.

Tests: HA harness for both apply races and a live config/set route drop,
recorder stubs for the four durability cases; five mutants caught by the
standard runner.

Issue: #495
User-Visible: yes
This commit is contained in:
Codex
2026-09-09 08:34:46 +03:00
parent adbd349920
commit 2946e28352
8 changed files with 427 additions and 50 deletions
+105 -48
View File
@@ -317,15 +317,53 @@ class TrailRecorder:
A tombstone deliberately stays in config so discovery cannot resurrect
a deleted device. For live tracking and trail ownership it is absent:
this is the same boundary used by ``async_refresh`` above.
Returns the number of erased markers. Runs dropped because their route
vanished (#162) are not counted, but they are just as durable (#495):
they leave with the orphan transaction, or with a write of their own.
"""
live_marker_ids = {
str(marker.get("id"))
for marker in config.get("markers") or []
if marker.get("id") is not None and marker.get("removed") is not True
}
# #162: a route that vanished (deleted, or re-targeted to another
# space, which is a new identity) takes its own runs with it, while the
# marker and its other routes stay untouched.
fire = False
try:
async with self._refresh_lock:
# #162: a route that vanished (deleted, or re-targeted to
# another space, which is a new identity) takes its own runs
# with it, while the marker and its other routes stay untouched.
dropped = self._drop_unknown_routes(config)
orphan_ids = set(self.book.data) - live_marker_ids
if not orphan_ids and not dropped:
return 0
try:
if orphan_ids:
removed = await self._delete_many_locked(orphan_ids)
else:
removed = 0
await self._save_now_locked()
except Exception:
# The store is the durable authority (#335): a drop that
# never reached it must not survive only in memory, or a
# restart brings the runs back while the card believed
# them gone. Put them back so the next commit retries.
self._restore_dropped(dropped)
raise
fire = True
return removed
except Exception: # noqa: BLE001 — config commit already succeeded
_LOGGER.exception(
"House Plan: reconciling vacuum trails after a config commit failed",
)
return 0
finally:
if fire:
self.hass.bus.async_fire("houseplan_trail_updated", {})
def _drop_unknown_routes(self, config: dict[str, Any]) -> dict[str, dict[str, Any]]:
"""Forget runs of routes the config no longer has; return what was dropped."""
dropped: dict[str, dict[str, Any]] = {}
for marker in config.get("markers") or []:
marker_id = str(marker.get("id"))
if marker.get("removed") is True or marker_id not in self.book.data:
@@ -333,56 +371,75 @@ class TrailRecorder:
vacuum = marker.get("vacuum") or {}
routes = effective_routes(
marker_id, vacuum, str(marker.get("space") or ""), vacuum.get("source"))
self.book.drop_unknown_routes(marker_id, {str(r.get("id")) for r in routes})
orphan_ids = set(self.book.data) - live_marker_ids
if not orphan_ids:
return 0
before = {
slot: run for slot, run in (self.book.data.get(marker_id) or {}).items()
if slot in ("current", "previous")
}
if self.book.drop_unknown_routes(marker_id, {str(r.get("id")) for r in routes}):
after = self.book.data.get(marker_id) or {}
dropped[marker_id] = {
slot: run for slot, run in before.items() if slot not in after
}
return dropped
def _restore_dropped(self, dropped: dict[str, dict[str, Any]]) -> None:
for marker_id, runs in dropped.items():
self.book.data.setdefault(marker_id, {}).update(runs)
async def _save_now_locked(self) -> None:
"""Write the book immediately, superseding a pending debounced write.
A failed write puts the debounce back: the points it was going to
persist are still only in memory.
"""
had_pending_save = self._unsub_save is not None
if self._unsub_save:
self._unsub_save()
self._unsub_save = None
try:
return await self._async_delete_many(orphan_ids)
except Exception: # noqa: BLE001 — config commit already succeeded
_LOGGER.exception(
"House Plan: removing orphan vacuum trails failed: markers=%s",
sorted(orphan_ids),
)
return 0
await self.store.async_save(self.book.data)
except Exception:
if had_pending_save:
self._schedule_save()
raise
async def _async_delete_many(self, markers: set[str]) -> int:
"""Delete one or more books with one subscription/store transaction."""
async with self._refresh_lock:
# The trail book owns deletion. When it has no such marker, this
# is a no-op and must not silently damage the live tracking graph.
removed = {
marker: self.book.data.pop(marker)
for marker in markers
if marker in self.book.data
}
if not removed:
return 0
previous_pairs = {src: list(pairs) for src, pairs in self.pairs.items()}
had_pending_save = self._unsub_save is not None
try:
for src in list(self.pairs):
kept = [pair for pair in self.pairs[src] if pair[0] not in removed]
if kept:
self.pairs[src] = kept
else:
del self.pairs[src]
self._resubscribe()
if self._unsub_save:
self._unsub_save()
self._unsub_save = None
await self.store.async_save(self.book.data)
except Exception:
# The store is the durable authority. Restore the in-memory
# owner graph so the next successful config sync can retry
# instead of leaving an orphan on disk forever (#335).
self.book.data.update(removed)
self.pairs = previous_pairs
self._resubscribe()
if had_pending_save:
self._schedule_save()
raise
self.hass.bus.async_fire("houseplan_trail_updated", {})
removed = await self._delete_many_locked(markers)
if removed:
self.hass.bus.async_fire("houseplan_trail_updated", {})
return removed
async def _delete_many_locked(self, markers: set[str]) -> int:
"""The transaction of ``_async_delete_many``; the caller holds the lock."""
# The trail book owns deletion. When it has no such marker, this
# is a no-op and must not silently damage the live tracking graph.
removed = {
marker: self.book.data.pop(marker)
for marker in markers
if marker in self.book.data
}
if not removed:
return 0
previous_pairs = {src: list(pairs) for src, pairs in self.pairs.items()}
try:
for src in list(self.pairs):
kept = [pair for pair in self.pairs[src] if pair[0] not in removed]
if kept:
self.pairs[src] = kept
else:
del self.pairs[src]
self._resubscribe()
await self._save_now_locked()
except Exception:
# The store is the durable authority. Restore the in-memory
# owner graph so the next successful config sync can retry
# instead of leaving an orphan on disk forever (#335).
self.book.data.update(removed)
self.pairs = previous_pairs
self._resubscribe()
raise
return len(removed)
def _resubscribe(self) -> None:
+6 -2
View File
@@ -630,8 +630,12 @@ async def ws_import_apply(hass: HomeAssistant, connection, msg: dict[str, Any])
"final_metadata": original_metadata,
}
await _commit_pair(rt, pending, rollback)
# A token becomes single-use only after both durable halves land.
get_candidate(rt, msg["token"], _connection_user_id(connection), consume=True)
# Both halves are durable: the token is spent no matter what
# happened to the preview registry meanwhile. Re-validating it
# here (TTL lapsed during the write, eviction by a newer preview
# of the same user) would report a failure for a plan that is
# already replaced — the outcome follows the commit (#495).
rt.import_previews.pop(msg["token"], None)
except ImportFailure as err:
_send_import_error(connection, msg["id"], err)
return
+8
View File
@@ -2,6 +2,14 @@
## Unreleased
- Applying a backup import no longer reports "preview expired" after the plan
was already replaced: once both halves of the import are written, the result
and the update events follow the commit even if the preview timed out or was
displaced by a newer preview while the write was in flight. Deleting or
re-targeting a robot map route now removes its recorded runs from disk as
well as from memory, so a Home Assistant restart cannot bring them back
([#495](https://github.com/Matysh/houseplan-card/issues/495)).
## v1.73.0-beta.7 — 2026-09-09
- Summary-panel settings stay responsive with large Home Assistant installs:
+8
View File
@@ -8,6 +8,14 @@
## Не выпущено
- Применение импорта из резервной копии больше не отвечает «preview expired»
после того, как план уже заменён: когда обе половины импорта записаны, ответ
и события обновления следуют за записью, даже если срок предпросмотра истёк
или его вытеснил более новый предпросмотр во время записи. Удаление или
перенацеливание маршрута карты робота теперь убирает записанные проезды и с
диска, а не только из памяти — перезапуск Home Assistant их не вернёт
([#495](https://github.com/Matysh/houseplan-card/issues/495)).
## v1.73.0-beta.7 — 2026-09-09
- Настройки сводной панели теперь остаются отзывчивыми даже при большом числе
+69
View File
@@ -7949,6 +7949,75 @@ const MUTANT_DEFINITIONS = [
replace: ' checked = _normalize(msg["config"]) # mutant: Optimize bypasses summary guard',
}],
},
{
id: 'import-apply-rechecks-preview-after-commit',
guard: 'node scripts/backend-test-guard.mjs '
+ 'issue_495_apply_result '
+ 'tests_backend/test_ha_import_export.py',
because: 'once both halves are durable, a lapsed TTL or an eviction by a newer preview must '
+ 'not turn the accepted import into an error (#495 AC1/AC2)',
patches: [{
file: 'custom_components/houseplan/websocket_api.py',
find: ' rt.import_previews.pop(msg["token"], None)\n',
replace: ' get_candidate(rt, msg["token"], _connection_user_id(connection), consume=True)\n',
}],
},
{
id: 'trail-purge-drops-routes-in-memory-only',
guard: 'node scripts/backend-test-guard.mjs '
+ 'issue_495_dropped_route_runs_reach_the_store '
+ 'tests_backend/test_trail_recorder.py',
because: 'a route deleted from a live marker without any orphan must still write the book, or a '
+ 'restart brings its runs back (#495 AC3)',
patches: [{
file: 'custom_components/houseplan/trails.py',
find: ' removed = 0\n await self._save_now_locked()\n',
replace: ' removed = 0 # mutant: memory only\n',
}],
},
{
id: 'trail-purge-keeps-dropped-runs-after-failed-write',
guard: 'node scripts/backend-test-guard.mjs '
+ 'issue_495_dropped_route_runs_roll_back '
+ 'tests_backend/test_trail_recorder.py',
because: 'when the store refuses the write the dropped runs must return to memory so the next '
+ 'commit retries instead of diverging from the disk (#495 AC4)',
patches: [{
file: 'custom_components/houseplan/trails.py',
find: ' self._restore_dropped(dropped)\n raise\n',
replace: ' raise # mutant: memory keeps the drop the disk never saw\n',
}],
},
{
id: 'trail-purge-forgets-routes-when-orphans-exist',
guard: 'node scripts/backend-test-guard.mjs '
+ 'issue_495_dropped_route_runs_share_the_orphan_transaction '
+ 'tests_backend/test_trail_recorder.py',
because: 'orphan deletion and route drops are one transaction: the orphan path must not write a '
+ 'book that still holds the dropped runs (#495 AC5)',
patches: [{
file: 'custom_components/houseplan/trails.py',
find: ' if orphan_ids:\n removed = await self._delete_many_locked(orphan_ids)\n',
replace: ' if orphan_ids:\n self._restore_dropped(dropped) # mutant\n'
+ ' removed = await self._delete_many_locked(orphan_ids)\n',
}],
},
{
id: 'trail-purge-failed-orphan-write-loses-dropped-runs',
guard: 'node scripts/backend-test-guard.mjs '
+ 'issue_495_failed_orphan_transaction_restores_dropped_route_runs_too '
+ 'tests_backend/test_trail_recorder.py',
because: 'the #335 rollback of the orphan transaction must cover the route drops made in the same '
+ 'pass, not only the orphan books (#495 AC4)',
patches: [{
file: 'custom_components/houseplan/trails.py',
find: ' except Exception:\n'
+ ' # The store is the durable authority (#335): a drop that\n',
replace: ' except Exception:\n'
+ ' dropped = {} # mutant: nothing to restore\n'
+ ' # The store is the durable authority (#335): a drop that\n',
}],
},
];
const mutationCardSource = readFileSync(join(repoRoot, 'src/houseplan-card.ts'), 'utf8');
+70
View File
@@ -2522,6 +2522,76 @@ async def test_success_events_are_emitted_only_after_both_target_writes(
assert all(config_done and layout_done for _event, config_done, layout_done in observed)
async def _apply_with_registry_hook(
hass: HomeAssistant, tmp_path: Path, monkeypatch, hook,
) -> tuple[Any, dict[str, Any], _Connection, list[str]]:
"""Apply a full import; ``hook(rt, token)`` runs after both halves are durable."""
await _setup(hass)
rt, response, _ = await _candidate(hass, tmp_path)
real_commit = wsapi._commit_pair
async def commit_then_hook(runtime: Any, pending: dict[str, Any], rollback: dict[str, Any]) -> None:
await real_commit(runtime, pending, rollback)
hook(runtime, response["token"])
fired: list[str] = []
def fire(_bus: Any, event_type: str, _event_data: dict[str, Any] | None = None, **_kwargs: Any) -> None:
if event_type in {"houseplan_config_updated", "houseplan_layout_updated"}:
fired.append(event_type)
monkeypatch.setattr(wsapi, "_commit_pair", commit_then_hook)
monkeypatch.setattr(type(hass.bus), "async_fire", fire)
connection = await _apply(hass, response)
return rt, response, connection, fired
async def test_issue_495_apply_result_follows_the_commit_when_the_preview_expires_meanwhile(
hass: HomeAssistant, tmp_path: Path, monkeypatch,
) -> None:
"""#495 AC1: a TTL that lapses during the write cannot turn a durable import into an error."""
def expire(rt: Any, token: str) -> None:
rt.import_previews[token]["expires"] = 0
rt, response, connection, fired = await _apply_with_registry_hook(hass, tmp_path, monkeypatch, expire)
assert connection.error is None, connection.error
assert connection.result and connection.result["ok"]
assert connection.result["config_rev"] == connection.result["layout_rev"] == 2
assert fired == ["houseplan_config_updated", "houseplan_layout_updated"]
assert response["token"] not in rt.import_previews
assert (await rt.config_store.async_load())["rev"] == 2
assert (await rt.config_store.async_load())["config"]["spaces"][0]["title"] == "Imported"
async def test_issue_495_apply_result_survives_eviction_by_a_newer_preview(
hass: HomeAssistant, tmp_path: Path, monkeypatch,
) -> None:
"""#495 AC2: previews made by the same user during the write may evict the token; the import stands."""
monkeypatch.setattr(import_export_api, "MAX_IMPORT_PREVIEWS_PER_USER", 1)
newer: list[str] = []
def evict(rt: Any, token: str) -> None:
raw = json.dumps(create_export(
rt, {"config": _config(), "rev": 1}, {"layout": {}, "rev": 1},
kind="full", space_id=None, card_version="1.61.0",
config_root=Path(hass.config.path("")),
)[0]).encode()
newer.append(create_preview(
rt, raw, owner_id="review-owner", duplicate_policy="skip",
current_config_data={"config": _config(), "rev": 1},
current_layout_data={"layout": {}, "rev": 1}, config_root=Path(hass.config.path("")),
)["token"])
assert token not in rt.import_previews, "the hook must evict the applying token"
rt, response, connection, fired = await _apply_with_registry_hook(hass, tmp_path, monkeypatch, evict)
assert connection.error is None, connection.error
assert connection.result and connection.result["ok"]
assert fired == ["houseplan_config_updated", "houseplan_layout_updated"]
assert response["token"] not in rt.import_previews
assert newer and all(token in rt.import_previews for token in newer)
assert (await rt.store.async_load())["rev"] == 2
async def test_pair_rolls_back_before_reporting_a_persistent_target_failure(
hass: HomeAssistant, tmp_path: Path, monkeypatch,
) -> None:
+49
View File
@@ -474,6 +474,55 @@ async def test_config_set_purges_tombstoned_and_absent_trails_durably(
assert "late_orphan" in (await recorder.store.async_load() or {})
async def test_issue_495_config_set_dropping_a_route_purges_its_runs_durably(
hass: HomeAssistant, hass_ws_client: WebSocketGenerator
) -> None:
"""#495 AC6: a route deleted from a live marker takes its runs off the disk, not just out of memory."""
await _setup(hass)
client = await hass_ws_client(hass)
def robot(*route_ids: str) -> dict:
return {
"id": "robot", "binding": "entity:vacuum.robot", "space": "f1",
"vacuum": {"source": "camera.map", "map_routes": [
{"id": route_id, "source": "camera.map", "map_id": route_id,
"space": "f1", "calibration": [1, 0, 0, 0, 1, 0]}
for route_id in route_ids
]},
}
initial = {"spaces": [_space("f1", "r1")], "markers": [robot("vr_old", "vr_keep")], "settings": {}}
await client.send_json_auto_id({
"type": "houseplan/config/set", "config": initial, "expected_rev": 0,
})
first = await client.receive_json()
assert first["success"], first
await hass.async_block_till_done()
recorder = hass.data[DOMAIN]["trail_recorder"]
recorder.book.data = {"robot": {
"current": {"route_id": "vr_old", "map_id": "vr_old", "points": [[1, 2]]},
"previous": {"route_id": "vr_keep", "map_id": "vr_keep", "points": [[3, 4]]},
}}
await recorder.store.async_save(copy.deepcopy(recorder.book.data))
candidate = {**copy.deepcopy(initial), "markers": [robot("vr_keep")]}
await client.send_json_auto_id({
"type": "houseplan/config/set", "config": candidate,
"expected_rev": first["result"]["rev"],
})
dropped = await client.receive_json()
assert dropped["success"], dropped
await client.send_json_auto_id({"type": "houseplan/trail/get"})
trails = (await client.receive_json())["result"]["trails"]
assert trails["robot"].get("current") is None
assert trails["robot"]["previous"]["route_id"] == "vr_keep"
durable = await recorder.store.async_load() or {}
assert durable["robot"].get("current") is None, "the dropped run must not survive a restart"
assert durable["robot"]["previous"]["route_id"] == "vr_keep"
async def test_config_rev_conflict(hass: HomeAssistant, hass_ws_client: WebSocketGenerator) -> None:
await _setup(hass)
client = await hass_ws_client(hass)
+112
View File
@@ -294,6 +294,118 @@ def test_failed_orphan_store_write_rolls_back_and_can_be_retried():
assert hass.bus.fired == [("houseplan_trail_updated", {})]
def _route_book():
"""One live marker: a run filed under a route that is about to vanish, one to keep."""
return {
"live": {
"current": {"route_id": "vr_old", "map_id": "1", "points": [[1, 2]]},
"previous": {"route_id": "vr_keep", "map_id": "1", "points": [[3, 4]]},
},
}
def _live_marker_without_vr_old():
return {
"id": "live", "space": "ground",
"vacuum": {"map_routes": [
{"id": "vr_keep", "source": "camera.map", "map_id": "1", "space": "ground"},
]},
}
class _RecordingStore:
def __init__(self):
self.saved = []
async def async_save(self, data):
self.saved.append(json.loads(json.dumps(data)))
def test_issue_495_dropped_route_runs_reach_the_store_without_orphans():
"""#495 AC3: deleting a route persists the drop even when no marker is orphaned."""
rec, hass, _states = _rec()
rec.book.data = _route_book()
rec.pairs = {"camera.map": [("live", "vacuum.x50")]}
rec.store = _RecordingStore()
pending_cancelled = []
rec._unsub_save = lambda: pending_cancelled.append(True)
removed = _run_isolated(rec.async_purge_orphans({"markers": [_live_marker_without_vr_old()]}))
assert removed == 0, "the return value still counts markers, not runs"
assert len(rec.store.saved) == 1
assert rec.store.saved[0] == {
"live": {"previous": {"route_id": "vr_keep", "map_id": "1", "points": [[3, 4]]}},
}
assert pending_cancelled == [True] and rec._unsub_save is None
assert hass.bus.fired == [("houseplan_trail_updated", {})]
# Restart model: a fresh book from what the store holds has no vr_old run.
restarted = trails.TrailBook(json.loads(json.dumps(rec.store.saved[-1])))
assert restarted.data["live"].get("current") is None
assert restarted.data["live"]["previous"]["route_id"] == "vr_keep"
# No change → no write, no event.
_run_isolated(rec.async_purge_orphans({"markers": [_live_marker_without_vr_old()]}))
assert len(rec.store.saved) == 1
assert hass.bus.fired == [("houseplan_trail_updated", {})]
def test_issue_495_dropped_route_runs_roll_back_when_the_store_write_fails():
"""#495 AC4: memory follows the store — a failed write keeps the run for a retry."""
rec, hass, _states = _rec()
rec.book.data = _route_book()
rec.pairs = {"camera.map": [("live", "vacuum.x50")]}
class FailingStore:
async def async_save(self, _data):
raise OSError("disk full")
rec.store = FailingStore()
assert _run_isolated(rec.async_purge_orphans({"markers": [_live_marker_without_vr_old()]})) == 0
assert rec.book.data == _route_book()
assert hass.bus.fired == []
rec.store = _RecordingStore()
assert _run_isolated(rec.async_purge_orphans({"markers": [_live_marker_without_vr_old()]})) == 0
assert rec.store.saved[-1]["live"] == {
"previous": {"route_id": "vr_keep", "map_id": "1", "points": [[3, 4]]},
}
assert hass.bus.fired == [("houseplan_trail_updated", {})]
def test_issue_495_dropped_route_runs_share_the_orphan_transaction():
"""#495 AC5: orphans and dropped routes leave in one write with one event."""
rec, hass, _states = _rec()
rec.book.data = {**_route_book(), "orphan": {"current": {"points": [[5, 6]]}}}
rec.pairs = {"camera.map": [("live", "vacuum.x50"), ("orphan", "vacuum.x50")]}
rec.store = _RecordingStore()
removed = _run_isolated(rec.async_purge_orphans({"markers": [_live_marker_without_vr_old()]}))
assert removed == 1
assert len(rec.store.saved) == 1
assert rec.store.saved[0] == {
"live": {"previous": {"route_id": "vr_keep", "map_id": "1", "points": [[3, 4]]}},
}
assert rec.pairs == {"camera.map": [("live", "vacuum.x50")]}
assert hass.bus.fired == [("houseplan_trail_updated", {})]
def test_issue_495_failed_orphan_transaction_restores_dropped_route_runs_too():
"""#495 AC4 in the orphan path: rollback covers both halves of the write."""
rec, hass, _states = _rec()
rec.book.data = {**_route_book(), "orphan": {"current": {"points": [[5, 6]]}}}
rec.pairs = {"camera.map": [("live", "vacuum.x50"), ("orphan", "vacuum.x50")]}
class FailingStore:
async def async_save(self, _data):
raise OSError("disk full")
rec.store = FailingStore()
assert _run_isolated(rec.async_purge_orphans({"markers": [_live_marker_without_vr_old()]})) == 0
assert rec.book.data == {**_route_book(), "orphan": {"current": {"points": [[5, 6]]}}}
assert hass.bus.fired == []
def test_object_style_position_is_read():
# Tasshack in-memory attributes hold a Point OBJECT, not a dict
class Point: