Files
HomeAssistantVS/custom_components/maintenance_supporter/websocket/io.py
T

1286 lines
54 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""WebSocket handlers for export, import, CSV, QR, and templates."""
from __future__ import annotations
import json as json_mod
import logging
import re
from functools import lru_cache
from typing import Any
from uuid import uuid4
import voluptuous as vol
from homeassistant.components import websocket_api
from homeassistant.core import HomeAssistant
from ..const import (
BATTERY_FLEET_DUE_WITHOUT_SENSOR,
BATTERY_FLEET_EXCLUDED,
BATTERY_FLEET_INCLUDED,
BATTERY_FLEET_OBJECT_FLAG,
BATTERY_FLEET_REMOVED_PARTS,
BATTERY_FLEET_TASK_FLAG,
BATTERY_FLEET_TRACK_SELF_CHARGING,
CONF_OBJECT,
CONF_OBJECT_MANUFACTURER,
CONF_OBJECT_MODEL,
CONF_OBJECT_NAME,
CONF_TASKS,
DOMAIN,
MAX_CHECKLIST_ITEM_LENGTH,
MAX_CHECKLIST_ITEMS,
MAX_ID_LENGTH,
MAX_IMPORT_PAYLOAD_BYTES,
MAX_JSON_IMPORT_PAYLOAD_BYTES,
)
from ..helpers.global_options import get_default_warning_days
from ..helpers.phases import clamp_phase_cursor, sanitize_phase_defs, sanitize_phase_sequence
from ..helpers.qr_generator import (
_ACTION_ICON_MAP,
build_qr_url,
generate_qr_svg,
generate_qr_svg_data_uri,
)
from ..websocket.tasks import _check_nfc_tag_duplicate, _validate_trigger_config
from . import _get_object_entries, _load_object_entry
_LOGGER = logging.getLogger(__name__)
def _ref_or_none(value: Any) -> int | None:
"""#170: a reference number / counter from a backup — positive int or nothing."""
return value if isinstance(value, int) and not isinstance(value, bool) and value > 0 else None
def _iso_marker(value: Any) -> str | None:
"""Keep ``value`` only if it parses as an ISO date/datetime, else drop it.
``paused_at`` is a *marker* whose mere presence means "paused"; a garbage
value imported from a hand-edited/foreign backup would otherwise freeze the
object as paused forever (and a malformed ``paused_until`` means auto-resume
never fires). Validate on import so only a real timestamp restores the state.
"""
from datetime import date, datetime
if not isinstance(value, str) or not value.strip():
return None
s = value.strip()
try:
datetime.fromisoformat(s.replace("Z", "+00:00"))
return s
except ValueError:
try:
date.fromisoformat(s)
return s
except ValueError:
return None
def _sanitize_history(history: Any) -> list[dict[str, Any]]:
"""Scrub imported history entries: drop a non-finite/negative ``cost``.
Every live write path range-guards cost, but import copied history verbatim
and ``json.loads``/``yaml.safe_load`` both accept ``NaN``/``Infinity``. Such
a value would poison budget aggregation (a `+inf` fake "budget exceeded"
alert, or `nan` silently disabling all alerts). The completion still counts;
only the bad cost is removed.
"""
import math
if not isinstance(history, list):
return []
out: list[dict[str, Any]] = []
for entry in history:
if not isinstance(entry, dict):
continue
clean = dict(entry)
cost = clean.get("cost")
if isinstance(cost, bool) or not isinstance(cost, (int, float)) or not math.isfinite(cost) or cost < 0:
clean.pop("cost", None)
# Readings (#83 / #161 phase 2): same NaN/Infinity hole — a poisoned
# value would break every delta after it. Malformed slot snapshots
# are dropped item-wise, the completion itself is kept.
rv = clean.get("reading_value")
if rv is not None and (isinstance(rv, bool) or not isinstance(rv, (int, float)) or not math.isfinite(rv)):
clean.pop("reading_value", None)
if "reading_values" in clean:
from ..helpers.reading_slots import history_reading_values
snapshot = history_reading_values(clean)
if snapshot:
clean["reading_values"] = snapshot
else:
clean.pop("reading_values", None)
out.append(clean)
return out
def _remap_document_refs(
import_tasks: dict[str, dict[str, Any]],
import_parts: dict[str, dict[str, Any]],
doc_id_map: dict[str, str],
) -> None:
"""Re-point document references at the freshly minted doc ids.
History entries carry completion photos (``photo_doc_ids``, or the
pre-2.75 ``photo_doc_id`` scalar — folded into the list here) and spare
parts carry a ``doc_id``. Ids the export did not carry (a hand-written
file, a doc that vanished before the export) stay verbatim: a dangling
reference renders as a missing picture, which is what it is.
"""
from ..helpers.completion_photos import history_photo_ids
for task_data in import_tasks.values():
for hist_entry in task_data.get("history") or []:
if not isinstance(hist_entry, dict):
continue
photos = history_photo_ids(hist_entry)
if not photos:
continue
hist_entry.pop("photo_doc_id", None)
hist_entry["photo_doc_ids"] = [doc_id_map.get(p, p) for p in photos]
for part in import_parts.values():
old = part.get("doc_id")
if isinstance(old, str) and old in doc_id_map:
part["doc_id"] = doc_id_map[old]
async def _drop_imported_documents(doc_store: Any, object_id: str) -> None:
"""Undo a pre-flow document import when the object never came to be.
Documents are recreated BEFORE the entry flow (their fresh ids must be
known to remap history photos and part doc_ids); if the flow then
fails they would linger as orphans nobody can reach. ``object_id`` is
freshly minted per import, so every doc under it is ours to drop.
"""
if doc_store is None:
return
for orphan in list(doc_store.for_object(object_id)):
await doc_store.async_remove(orphan["id"])
def _import_fleet_identity(
hass: HomeAssistant,
obj_data: dict[str, Any],
import_obj: dict[str, Any],
obj_name: str,
) -> bool:
"""Restore the battery-fleet markers onto ``import_obj``; True if it is the fleet.
The exported flag is honoured only while this instance has NO fleet yet
(``find_fleet_entry`` returns the FIRST flagged entry, so a second flagged
object would silently shadow or be shadowed by the existing one). The
exclude/include lists are re-validated like the live WS writes: entity
ids only, deduped + sorted, capped at ``FLEET_LIST_CAP``.
"""
if obj_data.get(BATTERY_FLEET_OBJECT_FLAG) is not True:
return False
from homeassistant.core import valid_entity_id
from ..helpers.battery_fleet_setup import FLEET_LIST_CAP, find_fleet_entry
if find_fleet_entry(hass) is not None:
_LOGGER.warning(
"JSON import: %r is flagged as the Battery Fleet but this instance already has one — importing it as a plain object",
obj_name,
)
return False
import_obj[BATTERY_FLEET_OBJECT_FLAG] = True
for key in (BATTERY_FLEET_EXCLUDED, BATTERY_FLEET_INCLUDED):
raw_list = obj_data.get(key)
if not isinstance(raw_list, list):
continue
cleaned = sorted({e.strip() for e in raw_list if isinstance(e, str) and valid_entity_id(e.strip())})
if cleaned:
import_obj[key] = cleaned[:FLEET_LIST_CAP]
if obj_data.get(BATTERY_FLEET_TRACK_SELF_CHARGING) is True:
import_obj[BATTERY_FLEET_TRACK_SELF_CHARGING] = True
if obj_data.get(BATTERY_FLEET_DUE_WITHOUT_SENSOR) is False:
import_obj[BATTERY_FLEET_DUE_WITHOUT_SENSOR] = False
# Deleted type-parts stay deleted after a restore too — same id rule as
# _keep_fleet_part_id, so nothing but ``batt_<type>`` ids get through.
raw_removed = obj_data.get(BATTERY_FLEET_REMOVED_PARTS)
if isinstance(raw_removed, list):
removed = sorted({p.strip() for p in raw_removed if isinstance(p, str) and _keep_fleet_part_id(True, p.strip())})
if removed:
import_obj[BATTERY_FLEET_REMOVED_PARTS] = removed[:FLEET_LIST_CAP]
return True
def _keep_fleet_part_id(is_fleet: bool, old_id: str) -> bool:
"""Whether an imported part keeps its id instead of getting a fresh uuid.
Only the fleet's deterministic type-part ids (``batt_<type>``, minted by
battery_fleet_setup._type_part) — every other part id is re-minted so an
import can never collide with or impersonate an existing part.
"""
return is_fleet and old_id.startswith("batt_") and len(old_id) <= MAX_ID_LENGTH
@websocket_api.websocket_command({vol.Required("type"): f"{DOMAIN}/version"})
@websocket_api.async_response
async def ws_version(hass: HomeAssistant, connection: websocket_api.ActiveConnection, msg: dict[str, Any]) -> None:
"""The installed integration version (manifest).
Roadmap guard 2 — stale-bundle handshake: the panel compares this against
the version esbuild stamped into its bundle and offers a reload when a
cached old frontend is talking to a newer backend (HA's service worker
updates stale-while-revalidate, so this happens routinely after updates).
"""
from homeassistant.loader import async_get_integration
integration = await async_get_integration(hass, DOMAIN)
connection.send_result(msg["id"], {"version": integration.version})
@websocket_api.websocket_command(
{
vol.Required("type"): f"{DOMAIN}/templates",
# v2.21.1: the caller's UI language — template/task names arrive
# localized. Falls back to the server language.
vol.Optional("language"): vol.All(str, vol.Length(max=10)),
}
)
@websocket_api.async_response
async def ws_get_templates(
hass: HomeAssistant,
connection: websocket_api.ActiveConnection,
msg: dict[str, Any],
) -> None:
"""Return all maintenance templates.
Every template is returned with a ``disabled`` flag (v2.21 gallery
curation): the pickers hide disabled ones client-side, while the Settings
section needs the full list to render the toggles.
"""
from ..helpers.i18n import normalize_language, normalize_language_code
from ..templates import (
TEMPLATE_CATEGORIES,
TEMPLATES,
get_disabled_template_ids,
localize_template_text,
)
disabled = get_disabled_template_ids(hass)
lang = normalize_language_code(msg.get("language")) if msg.get("language") else normalize_language(hass)
result = {
"categories": {cat_id: {k: v for k, v in cat.items()} for cat_id, cat in TEMPLATE_CATEGORIES.items()},
"templates": [
{
"id": t.id,
"name": localize_template_text(t.name, lang),
"category": t.category,
"disabled": t.id in disabled,
"tasks": [
{
"name": localize_template_text(tt.name, lang),
"type": tt.type,
"schedule_type": tt.schedule_type,
"interval_days": tt.interval_days,
"warning_days": tt.warning_days,
}
for tt in t.tasks
],
}
for t in TEMPLATES
],
}
connection.send_result(msg["id"], result)
@websocket_api.websocket_command(
{
vol.Required("type"): f"{DOMAIN}/export",
vol.Optional("format", default="json"): vol.In(["json", "yaml"]),
vol.Optional("include_history", default=True): bool,
# Selective export: restrict to these object entry_ids (omit = all).
vol.Optional("entry_ids"): [vol.All(str, vol.Length(max=MAX_ID_LENGTH))],
}
)
@websocket_api.require_admin
@websocket_api.async_response
async def ws_export_data(
hass: HomeAssistant,
connection: websocket_api.ActiveConnection,
msg: dict[str, Any],
) -> None:
"""Export all (or a selection of) maintenance data as JSON or YAML."""
from ..export import build_export_data, serialize_export
fmt = msg.get("format", "json")
include_history = msg.get("include_history", True)
entry_ids = set(msg["entry_ids"]) if msg.get("entry_ids") else None
# Phase 1: gather data on the event loop (accesses HA APIs)
data = build_export_data(hass, include_history=include_history, entry_ids=entry_ids)
# Phase 2: serialize in executor (CPU-bound, no HA API calls)
result = await hass.async_add_executor_job(serialize_export, data, fmt)
connection.send_result(msg["id"], {"format": fmt, "data": result})
@websocket_api.websocket_command(
{
vol.Required("type"): f"{DOMAIN}/csv/export",
vol.Optional("entry_ids"): [vol.All(str, vol.Length(max=MAX_ID_LENGTH))],
}
)
@websocket_api.require_admin
@websocket_api.async_response
async def ws_export_csv(
hass: HomeAssistant,
connection: websocket_api.ActiveConnection,
msg: dict[str, Any],
) -> None:
"""Export all (or a selection of) maintenance data as CSV."""
from ..helpers.csv_handler import export_objects_csv
entry_ids = set(msg["entry_ids"]) if msg.get("entry_ids") else None
csv_data = export_objects_csv(hass, entry_ids=entry_ids)
connection.send_result(msg["id"], {"csv": csv_data})
@websocket_api.websocket_command(
{
vol.Required("type"): f"{DOMAIN}/objects/csv",
vol.Optional("entry_ids"): [vol.All(str, vol.Length(max=MAX_ID_LENGTH))],
}
)
@websocket_api.async_response
async def ws_export_objects_csv(
hass: HomeAssistant,
connection: websocket_api.ActiveConnection,
msg: dict[str, Any],
) -> None:
"""Export one row per maintenance object as CSV (#67), all or a selection.
Not admin-gated: it exposes only the asset fields the panel already sends
to every user via ``maintenance_supporter/objects`` (no cost/history).
"""
from ..helpers.csv_handler import export_object_records_csv
entry_ids = set(msg["entry_ids"]) if msg.get("entry_ids") else None
csv_data = export_object_records_csv(hass, entry_ids=entry_ids)
connection.send_result(msg["id"], {"csv": csv_data})
@websocket_api.websocket_command(
{
vol.Required("type"): f"{DOMAIN}/csv/import",
vol.Required("csv_content"): str,
}
)
@websocket_api.require_admin
@websocket_api.async_response
async def ws_import_csv(
hass: HomeAssistant,
connection: websocket_api.ActiveConnection,
msg: dict[str, Any],
) -> None:
"""Import maintenance objects from CSV content."""
from ..helpers.csv_handler import import_objects_csv
csv_content = msg["csv_content"]
# Guard against oversized payloads (max 1MB / 1000 objects)
if len(csv_content) > MAX_IMPORT_PAYLOAD_BYTES:
connection.send_error(msg["id"], "too_large", "CSV content exceeds 1MB limit")
return
objects = import_objects_csv(csv_content, hass=hass)
if len(objects) > 1000:
connection.send_error(msg["id"], "too_many", "CSV contains more than 1000 objects")
return
if not objects:
connection.send_error(msg["id"], "empty_csv", "No valid objects found in CSV")
return
created = []
errors: list[dict[str, str]] = []
for idx, obj_data in enumerate(objects):
# Check for NFC tag duplicates in CSV-imported tasks
nfc_warnings: list[str] = []
for t_data in obj_data.get("tasks", {}).values():
nfc_val = t_data.get("nfc_tag_id")
if nfc_val:
nfc_warn = _check_nfc_tag_duplicate(hass, nfc_val)
if nfc_warn:
nfc_warnings.append(nfc_warn)
try:
result = await hass.config_entries.flow.async_init(
DOMAIN,
context={"source": "websocket"},
data={
CONF_OBJECT: obj_data["object"],
CONF_TASKS: obj_data["tasks"],
},
)
except Exception:
obj_name = obj_data.get("object", {}).get("name", f"row {idx + 1}")
_LOGGER.exception("CSV import failed for %s", obj_name)
errors.append({"name": obj_name, "reason": "unexpected error"})
continue
if result["type"] == "create_entry":
entry_info: dict[str, Any] = {
"entry_id": result["result"].entry_id,
"name": obj_data["object"].get("name", ""),
"task_count": len(obj_data["tasks"]),
}
if nfc_warnings:
entry_info["warnings"] = nfc_warnings
created.append(entry_info)
else:
obj_name = obj_data.get("object", {}).get("name", f"row {idx + 1}")
errors.append({"name": obj_name, "reason": result.get("reason", "unknown")})
resp: dict[str, Any] = {
"imported": created,
"total": len(objects),
"created": len(created),
}
if errors:
resp["errors"] = errors
connection.send_result(msg["id"], resp)
def _parse_structured(raw: str) -> Any:
"""Parse JSON *or* YAML export content into a Python object.
Both formats are accepted so every structured export (JSON and YAML)
round-trips back through the importer. Raises ValueError if the content
parses to neither a mapping nor a list.
"""
try:
return json_mod.loads(raw)
except (json_mod.JSONDecodeError, ValueError):
pass
import yaml # type: ignore[import-untyped]
try:
loaded = yaml.safe_load(raw)
except yaml.YAMLError as err:
raise ValueError("not valid JSON or YAML") from err
# safe_load returns a bare string/scalar for non-structured text (e.g. a
# CSV blob) — require an object/array so those route elsewhere cleanly.
if not isinstance(loaded, (dict, list)):
raise ValueError("not valid JSON or YAML")
return loaded
@websocket_api.websocket_command(
{
vol.Required("type"): f"{DOMAIN}/settings/export",
}
)
@websocket_api.require_admin
@websocket_api.async_response
async def ws_export_settings(
hass: HomeAssistant,
connection: websocket_api.ActiveConnection,
msg: dict[str, Any],
) -> None:
"""Export the global entry's settings as JSON.
The objects export deliberately excludes the global scope (groups, saved
views, vacation, notification/budget settings, feature toggles) — this is
its second half. Import goes through the regular json/import command,
which recognizes the ``global_settings`` section.
"""
from ..export import build_settings_export
connection.send_result(
msg["id"],
{"format": "json", "data": json_mod.dumps(build_settings_export(hass), indent=2)},
)
def _apply_settings_import(hass: HomeAssistant, raw: dict[str, Any]) -> list[str]:
"""Apply an imported ``global_settings`` payload; returns the applied keys.
Scalar settings run through the SAME validation as the ``global/update``
WS command (``sanitize_settings_input``); an invalid notify_service is
dropped rather than failing the import. The structured sections reuse
their own sanitizers: saved views via ``sanitize_view``, groups shape-
checked here, vacation dates validated like ``vacation/update``. Group
task_refs and vacation exempt ids may point at objects of the SOURCE
instance — they are kept verbatim (same-instance restores keep them
valid; elsewhere they degrade gracefully like every stale reference).
"""
from datetime import date as date_cls
from ..const import (
CONF_GROUPS,
CONF_NOTIFY_SERVICE,
CONF_SAVED_FILTER_VIEWS,
CONF_VACATION_BUFFER_DAYS,
CONF_VACATION_ENABLED,
CONF_VACATION_END,
CONF_VACATION_EXEMPT_TASK_IDS,
CONF_VACATION_START,
MAX_GROUP_TASK_REFS,
MAX_NAME_LENGTH,
)
from ..export import _NON_PORTABLE_SETTINGS
from ..helpers.global_options import get_global_entry
from ..helpers.saved_views import MAX_SAVED_VIEWS, sanitize_view
from ..helpers.settings_registry import ALLOWED_SETTING_KEYS
from .dashboard import sanitize_settings_input
entry = get_global_entry(hass)
if entry is None or not isinstance(raw, dict):
return []
scalars = {k: v for k, v in raw.items() if k in ALLOWED_SETTING_KEYS and k not in _NON_PORTABLE_SETTINGS}
filtered, notify_error = sanitize_settings_input(scalars)
if notify_error:
filtered.pop(CONF_NOTIFY_SERVICE, None)
groups_in = raw.get(CONF_GROUPS)
if isinstance(groups_in, dict):
groups: dict[str, dict[str, Any]] = {}
for gid, g in groups_in.items():
if not isinstance(g, dict) or not str(g.get("name") or "").strip():
continue
refs = [
{"entry_id": str(r["entry_id"]), "task_id": str(r["task_id"])}
for r in (g.get("task_refs") or [])
if isinstance(r, dict) and r.get("entry_id") and r.get("task_id")
][:MAX_GROUP_TASK_REFS]
groups[str(gid)] = {
"name": str(g["name"]).strip()[:MAX_NAME_LENGTH],
"description": str(g.get("description") or "")[:MAX_NAME_LENGTH],
"task_refs": refs,
}
if groups:
filtered[CONF_GROUPS] = groups
views_in = raw.get(CONF_SAVED_FILTER_VIEWS)
if isinstance(views_in, list):
views = []
for v in views_in[:MAX_SAVED_VIEWS]:
clean = sanitize_view(v, view_id=str(v.get("id")) if isinstance(v, dict) and v.get("id") else None)
if clean is not None:
views.append(clean)
if views:
filtered[CONF_SAVED_FILTER_VIEWS] = views
if isinstance(raw.get(CONF_VACATION_ENABLED), bool):
filtered[CONF_VACATION_ENABLED] = raw[CONF_VACATION_ENABLED]
for key in (CONF_VACATION_START, CONF_VACATION_END):
val = raw.get(key)
if isinstance(val, str):
try:
date_cls.fromisoformat(val)
except ValueError:
continue
filtered[key] = val
if isinstance(raw.get(CONF_VACATION_BUFFER_DAYS), int) and not isinstance(raw.get(CONF_VACATION_BUFFER_DAYS), bool):
filtered[CONF_VACATION_BUFFER_DAYS] = raw[CONF_VACATION_BUFFER_DAYS]
exempt = raw.get(CONF_VACATION_EXEMPT_TASK_IDS)
if isinstance(exempt, list):
cleaned = [t.strip() for t in exempt if isinstance(t, str) and t.strip()][:2000]
filtered[CONF_VACATION_EXEMPT_TASK_IDS] = cleaned
if not filtered:
return []
merged = dict(entry.options or entry.data)
merged.update(filtered)
hass.config_entries.async_update_entry(entry, options=merged)
_LOGGER.info("Settings import applied %d key(s)", len(filtered))
return sorted(filtered)
@websocket_api.websocket_command(
{
vol.Required("type"): f"{DOMAIN}/json/import",
vol.Required("json_content"): str,
}
)
@websocket_api.require_admin
@websocket_api.async_response
async def ws_import_json(
hass: HomeAssistant,
connection: websocket_api.ActiveConnection,
msg: dict[str, Any],
) -> None:
"""Import maintenance objects from JSON or YAML content (from /export)."""
raw = msg["json_content"]
if len(raw) > MAX_JSON_IMPORT_PAYLOAD_BYTES:
connection.send_error(msg["id"], "too_large", "Content exceeds 10MB limit")
return
try:
data = _parse_structured(raw)
except ValueError:
connection.send_error(msg["id"], "invalid_format", "Content is not valid JSON or YAML")
return
has_settings = isinstance(data, dict) and isinstance(data.get("global_settings"), dict)
if not isinstance(data, dict) or ("objects" not in data and not has_settings):
connection.send_error(msg["id"], "invalid_format", "JSON must contain an 'objects' array")
return
# A settings export (see export.build_settings_export) may travel alone or
# alongside an objects payload — apply it first either way.
settings_applied: list[str] = []
if has_settings:
settings_applied = _apply_settings_import(hass, data["global_settings"])
objects = data.get("objects", [])
if not isinstance(objects, list):
connection.send_error(msg["id"], "invalid_format", "'objects' must be an array")
return
if len(objects) > 1000:
connection.send_error(msg["id"], "too_many", "JSON contains more than 1000 objects")
return
if not objects and not settings_applied:
connection.send_error(msg["id"], "empty", "No objects found in JSON")
return
created = []
errors: list[dict[str, str]] = []
for idx, obj_entry in enumerate(objects):
# Guard against malformed-but-schema-valid input (the schema only checks
# json_content is a str): a non-dict entry / non-dict object would raise
# AttributeError and escape the per-object try/except below.
if not isinstance(obj_entry, dict):
errors.append({"name": f"object {idx + 1}", "reason": "not an object"})
continue
obj_data = obj_entry.get("object", {})
if not isinstance(obj_data, dict):
errors.append({"name": f"object {idx + 1}", "reason": "invalid object data"})
continue
obj_name = (obj_data.get("name") or "").strip()
if not obj_name:
errors.append({"name": f"object {idx + 1}", "reason": "missing name"})
continue
obj_id = uuid4().hex
import_obj: dict[str, Any] = {
"id": obj_id,
"name": obj_name,
"manufacturer": obj_data.get("manufacturer"),
"model": obj_data.get("model"),
"serial_number": obj_data.get("serial_number"),
"area_id": obj_data.get("area_id"),
"installation_date": obj_data.get("installation_date"),
"warranty_expiry": obj_data.get("warranty_expiry"),
# Imported counterparts of the export fields above; length-capped by
# cap_object_fields and the frontend only renders http(s) doc URLs.
"documentation_url": obj_data.get("documentation_url"),
"notes": obj_data.get("notes"),
# 2.19: device link / parent hierarchy — same-instance restores
# keep them valid; stale ids degrade gracefully at read time.
"ha_device_id": obj_data.get("ha_device_id"),
"parent_entry_id": obj_data.get("parent_entry_id"),
# 2.20: seasonal pause round-trips (a paused pool restored in
# winter stays paused); replace-flow lineage ids are the same
# instance-specific story as parent_entry_id above.
"paused_at": _iso_marker(obj_data.get("paused_at")),
"paused_until": _iso_marker(obj_data.get("paused_until")),
"predecessor_entry_id": obj_data.get("predecessor_entry_id"),
"replaced_by_entry_id": obj_data.get("replaced_by_entry_id"),
# Object-level archive marker — same presence-means-archived
# semantics as paused_at, so it gets the same ISO validation. Its
# tasks carry their own archived_* pair (mirrored below).
"archived_at": _iso_marker(obj_data.get("archived_at")),
# #170: keep the numbers a backup carries (collisions are
# renumbered by the setup pass); bool/negative junk is dropped.
"ref_no": _ref_or_none(obj_data.get("ref_no")),
"next_task_ref": _ref_or_none(obj_data.get("next_task_ref")),
"task_ids": [],
}
# Battery fleet identity (object flag + exclude/include lists + the
# self-charging opt-in). The fleet is ONE object by invariant
# (find_fleet_entry returns the first flagged entry), so the flag is
# only restored when this instance has no fleet yet — otherwise the
# payload imports as a plain object (fresh part ids, no task flag).
is_fleet = _import_fleet_identity(hass, obj_data, import_obj, obj_name)
# Spare parts: regenerate ids (like tasks) and remember the mapping so
# task-side links (consumes_parts / part_ref) can be rewritten below.
# Stock is dynamic Store state — collected here, written after setup.
# Fleet type-parts keep their deterministic ``batt_<type>`` ids: the
# fleet reconcile / mark-replaced paths key on them, so a re-minted
# uuid would orphan the whole battery-type ↔ part mapping.
from uuid import uuid4 as _uuid4
part_id_map: dict[str, str] = {}
import_parts: dict[str, dict[str, Any]] = {}
part_stocks: dict[str, float] = {}
parts_list = obj_entry.get("parts", [])
if isinstance(parts_list, list):
for part_entry in parts_list:
if not isinstance(part_entry, dict) or not (part_entry.get("name") or "").strip():
continue
old_id = str(part_entry.get("id") or "")
new_id = old_id if _keep_fleet_part_id(is_fleet, old_id) else _uuid4().hex
pdata = {k: v for k, v in part_entry.items() if k != "stock"}
pdata["id"] = new_id
# Drop a non-http(s) product_url — the WS write path validates it
# via _clean_url, but import copied it verbatim, so a crafted
# backup could persist a javascript: link (the panel now also
# guards the href, but keep bad data out of storage).
_purl = pdata.get("product_url")
if isinstance(_purl, str) and _purl.strip().lower().startswith(("http://", "https://")):
pdata["product_url"] = _purl.strip() # store trimmed so the render guard matches
else:
pdata.pop("product_url", None)
import_parts[new_id] = pdata
if old_id:
part_id_map[old_id] = new_id
stock = part_entry.get("stock")
if isinstance(stock, (int, float)) and not isinstance(stock, bool) and stock >= 0:
part_stocks[new_id] = stock
import_tasks: dict[str, dict[str, Any]] = {}
# old task id → new id, so document task-links (task_ids) can be
# remapped onto the freshly generated tasks (mirrors part_id_map).
task_id_map: dict[str, str] = {}
fleet_task_seen = False
# Per-task import losses (an invalid trigger is dropped, not fatal) —
# reported next to the NFC warnings instead of vanishing silently.
task_warnings: list[str] = []
tasks_list = obj_entry.get("tasks", [])
if not isinstance(tasks_list, list):
tasks_list = []
for task_entry in tasks_list:
if not isinstance(task_entry, dict):
continue
task_name = (task_entry.get("name") or "").strip()
if not task_name:
continue
task_id = uuid4().hex
old_task_id = str(task_entry.get("id") or "")
if old_task_id:
task_id_map[old_task_id] = task_id
task_data: dict[str, Any] = {
"id": task_id,
"object_id": obj_id,
"name": task_name,
"type": task_entry.get("type", "custom"),
"enabled": task_entry.get("enabled", True),
"schedule_type": task_entry.get("schedule_type", "time_based"),
"warning_days": task_entry.get("warning_days", get_default_warning_days(hass)),
"history": _sanitize_history(task_entry.get("history", [])),
}
for key in (
# Provenance + lifecycle — mirror the export builder so an
# archived task stays archived and created_at (the next_due
# fallback anchor) survives the round trip.
"created_at",
"archived_at",
"archived_reason",
"interval_days",
"interval_unit",
"due_date",
"interval_anchor",
"last_planned_due",
# per-occurrence postpone (round-trips like last_planned_due)
"due_override",
# nested recurrence (calendar kinds) — config-flow normalize
# treats it as authoritative when present.
"schedule",
"last_performed",
"notes",
"documentation_url",
"custom_icon",
"nfc_tag_id",
"require_tag_scan",
"allow_skip",
"responsible_user_id",
"entity_slug",
"trigger_config",
"adaptive_config",
"checklist",
"schedule_time",
# v2.17+ / #83 fields — mirror the export builder so a JSON
# backup round-trips them (validated/clamped just below).
"priority",
"labels",
"earliest_completion_days",
"ref_no",
"on_complete_action",
"quick_complete_defaults",
"assignee_pool",
"required_completion_fields",
"rotation_strategy",
"reading_unit",
"readings",
# spare parts (ids remapped below)
"consumes_parts",
"part_ref",
):
val = task_entry.get(key)
if val is not None:
task_data[key] = val
# The fleet's single aggregate task keeps its marker (detail view
# renders the battery section; the fleet reconcile repairs its
# trigger). Only ONE task may carry it, and only on the fleet.
if is_fleet and task_entry.get(BATTERY_FLEET_TASK_FLAG) is True and not fleet_task_seen:
task_data[BATTERY_FLEET_TASK_FLAG] = True
fleet_task_seen = True
# In-cycle checklist ticks: keyed by item TEXT so they survive the
# id regeneration; keys are filtered against the imported checklist
# exactly like the live checklist_progress WS write. Rides
# entry.data until the fresh entry's first setup migrates it into
# the Store (split-only field — storage._SPLIT_ONLY_TASK_FIELDS).
raw_progress = task_entry.get("checklist_progress")
if isinstance(raw_progress, dict):
items = set(task_data.get("checklist") or [])
progress = {k: bool(v) for k, v in raw_progress.items() if isinstance(k, str) and k in items}
if progress:
task_data["checklist_progress"] = progress
# #130: history entries carry used_parts, and since they are
# editable (stock reconciled by delta), the part ids must follow
# the regenerated ones. Own-part ids remap via part_id_map; links
# into another object's pool (entry_id set) are kept verbatim —
# if that entry doesn't exist in this instance they degrade to
# the safe recorded-only path, name preserved.
for hist_entry in task_data.get("history") or []:
used = hist_entry.get("used_parts")
if not isinstance(used, list):
continue
for link in used:
if (
isinstance(link, dict)
and not link.get("entry_id")
and link.get("part_id") in part_id_map
):
link["part_id"] = part_id_map[link["part_id"]]
# Remap part links to the regenerated part ids; drop dangling ones.
links = task_data.get("consumes_parts")
if isinstance(links, list):
remapped = []
for link in links:
if not isinstance(link, dict):
continue
foreign = str(link.get("entry_id") or "").strip()
if foreign:
# A link to another object's pool (#111). Import mints
# new entry ids, so the reference only means anything
# if that object is present in THIS instance — keep it
# then, drop it otherwise rather than restore a link
# that points nowhere.
if hass.config_entries.async_get_entry(foreign) is not None:
remapped.append(dict(link))
elif link.get("part_id") in part_id_map:
remapped.append(
{"part_id": part_id_map[link["part_id"]], "quantity": link.get("quantity", 1)}
)
if remapped:
task_data["consumes_parts"] = remapped
else:
task_data.pop("consumes_parts", None)
elif links is not None:
task_data.pop("consumes_parts", None)
ref = task_data.get("part_ref")
if isinstance(ref, dict) and ref.get("part_id") in part_id_map:
task_data["part_ref"] = {"part_id": part_id_map[ref["part_id"]]}
elif ref is not None:
task_data.pop("part_ref", None)
# Task phases (#139): sanitize like the live WS write, remap each
# phase's part links to the regenerated ids (same rules as the
# task-level links above), and clamp the cursor to the imported
# sequence. The cursor rides entry.data until the fresh entry's
# first setup migrates it into the Store (dynamic field), so a
# restore resumes mid-cycle.
raw_defs = task_entry.get("phases")
raw_seq = task_entry.get("phase_sequence")
if isinstance(raw_defs, dict) and isinstance(raw_seq, list):
defs = sanitize_phase_defs(raw_defs)
for pdef in defs.values():
plinks = pdef.get("consumes_parts")
if not isinstance(plinks, list):
continue
kept = []
for link in plinks:
if not isinstance(link, dict):
continue
foreign = str(link.get("entry_id") or "").strip()
if foreign:
if hass.config_entries.async_get_entry(foreign) is not None:
kept.append(dict(link))
elif link.get("part_id") in part_id_map:
kept.append(
{"part_id": part_id_map[link["part_id"]], "quantity": link.get("quantity", 1)}
)
if kept:
pdef["consumes_parts"] = kept
else:
pdef.pop("consumes_parts", None)
seq = sanitize_phase_sequence(raw_seq, defs)
if defs and seq:
task_data["phases"] = defs
task_data["phase_sequence"] = seq
task_data["phase_cursor"] = clamp_phase_cursor(task_entry.get("phase_cursor"), len(seq))
# Sanitize critical fields from import data
iv = task_data.get("interval_days")
if iv is not None and (not isinstance(iv, int) or iv < 1):
task_data.pop("interval_days", None)
lp = task_data.get("last_performed")
if lp is not None:
try:
from datetime import date
date.fromisoformat(lp)
except (ValueError, TypeError):
task_data.pop("last_performed", None)
wd = task_data.get("warning_days")
if not isinstance(wd, int) or wd < 0 or wd > 365:
task_data["warning_days"] = get_default_warning_days(hass)
# A rotation task must carry its effective assignee (imports from
# pre-seeding exports may lack one) — same rule as create/update.
from ..helpers.sanitize import seed_rotation_assignee
seed_rotation_assignee(task_data)
# Sanitize checklist: only keep string items within length budget,
# cap total items. Drops malformed entries silently rather than
# rejecting the whole import — same forgiving model as the other
# fields above.
cl = task_data.get("checklist")
if cl is not None:
if not isinstance(cl, list):
task_data.pop("checklist", None)
else:
cleaned = [item.strip() for item in cl if isinstance(item, str) and len(item) <= MAX_CHECKLIST_ITEM_LENGTH]
cleaned = [c for c in cleaned if c]
task_data["checklist"] = cleaned[:MAX_CHECKLIST_ITEMS]
# #161 phase 2: reading slots — same shape rules as the WS write.
if task_data.get("readings") is not None:
from ..helpers.reading_slots import sanitize_reading_slots
slots = sanitize_reading_slots(task_data["readings"])
if slots:
task_data["readings"] = slots
else:
task_data.pop("readings", None)
# schedule_time: strict HH:MM, otherwise drop
st = task_data.get("schedule_time")
if st is not None:
if not isinstance(st, str) or not re.fullmatch(r"^([01]\d|2[0-3]):[0-5]\d$", st):
task_data.pop("schedule_time", None)
# Validate an imported trigger_config the same way the WS create/update
# path does — strip unknown keys, normalize entity_ids, and drop it
# entirely if invalid — so import isn't a hole around trigger validation.
tc = task_data.get("trigger_config")
if isinstance(tc, dict):
# The export carries the live per-entity trigger state
# (accumulated runtime hours, counter baseline, change count)
# merged in as ``_trigger_state``. The validator strips it as
# an unknown key, so a restore silently started every
# sensor trigger from zero (bug review 2026-09-04). Keep it
# aside and re-attach it: the fresh entry's first setup
# migrates it into the Store like any other dynamic field.
trigger_state = tc.pop("_trigger_state", None)
tc_errors, _warnings = _validate_trigger_config(hass, tc)
if tc_errors:
task_data.pop("trigger_config", None)
task_warnings.append(f"{task_name}: trigger dropped — {tc_errors[0]}")
elif isinstance(trigger_state, dict) and trigger_state:
tc["_trigger_state"] = trigger_state
elif tc is not None:
task_data.pop("trigger_config", None)
task_warnings.append(f"{task_name}: trigger dropped — not a mapping")
import_tasks[task_id] = task_data
import_obj["task_ids"].append(task_id)
# Check for NFC tag duplicates across imported tasks
nfc_warnings: list[str] = []
for t_data in import_tasks.values():
nfc_val = t_data.get("nfc_tag_id")
if nfc_val:
nfc_warn = _check_nfc_tag_duplicate(hass, nfc_val)
if nfc_warn:
nfc_warnings.append(nfc_warn)
# (roadmap P6) recreate document metadata + web-links for the object
# (blobs travel via the /config backup; a JSON-only import leaves
# file docs dangling, which the storage-hygiene repair issue catches).
# Done BEFORE the entry is created: the docs get fresh ids, and the
# history entries (completion photos, #161) and spare parts (doc_id)
# that point at them by id must be re-pointed before they are
# persisted — the export carries the old ids for exactly this.
doc_store = None
import_docs = obj_entry.get("documents")
if isinstance(import_docs, list) and import_docs:
from .. import DOCUMENT_STORE_KEY
doc_store = hass.data.get(DOMAIN, {}).get(DOCUMENT_STORE_KEY)
if doc_store is not None:
doc_id_map: dict[str, str] = {}
await doc_store.async_import_documents(
obj_id, import_docs, task_id_map=task_id_map, part_id_map=part_id_map, id_map=doc_id_map
)
if doc_id_map:
_remap_document_refs(import_tasks, import_parts, doc_id_map)
try:
result = await hass.config_entries.flow.async_init(
DOMAIN,
context={"source": "websocket"},
data={
CONF_OBJECT: import_obj,
CONF_TASKS: import_tasks,
"parts": import_parts,
},
)
except Exception:
_LOGGER.exception("JSON import failed for %s", obj_name)
errors.append({"name": obj_name, "reason": "unexpected error"})
await _drop_imported_documents(doc_store, obj_id)
continue
if result["type"] == "create_entry":
entry_info: dict[str, Any] = {
"entry_id": result["result"].entry_id,
"name": obj_name,
"task_count": len(import_tasks),
}
if nfc_warnings or task_warnings:
entry_info["warnings"] = nfc_warnings + task_warnings
for warning in task_warnings:
_LOGGER.warning("JSON import of %s: %s", obj_name, warning)
created.append(entry_info)
# Restore tracked part stocks into the new entry's Store.
if part_stocks:
new_entry = hass.config_entries.async_get_entry(result["result"].entry_id)
rd_new = getattr(new_entry, "runtime_data", None) if new_entry else None
store_new = getattr(rd_new, "store", None) if rd_new else None
if store_new is not None:
for pid, stock_val in part_stocks.items():
store_new.set_part_stock(pid, stock_val)
await store_new.async_save()
# Restored stocks can sit below min_stock — reconcile buy
# tasks like every other stock mutation does, or the
# shopping list stays silent until the next unrelated
# stock change (bug audit 2026-08-22).
from ..parts_runtime import schedule_buy_task_reconcile
if new_entry is not None:
schedule_buy_task_reconcile(hass, new_entry)
else:
errors.append({"name": obj_name, "reason": result.get("reason", "unknown")})
await _drop_imported_documents(doc_store, obj_id)
resp: dict[str, Any] = {
"imported": created,
"total": len(objects),
"created": len(created),
}
if settings_applied:
resp["settings_applied"] = settings_applied
if errors:
resp["errors"] = errors
connection.send_result(msg["id"], resp)
@websocket_api.websocket_command(
{
vol.Required("type"): "maintenance_supporter/qr/generate",
vol.Required("entry_id"): vol.All(str, vol.Length(max=MAX_ID_LENGTH)),
vol.Optional("task_id"): vol.All(str, vol.Length(max=MAX_ID_LENGTH)),
vol.Optional("action", default="view"): vol.In(["view", "complete", "quick_complete"]),
vol.Optional("url_mode", default="server"): vol.In(["server", "local", "companion"]),
vol.Optional("base_url"): vol.All(vol.Url(), vol.Length(max=512)),
}
)
@websocket_api.async_response
async def ws_generate_qr(
hass: HomeAssistant,
connection: websocket_api.ActiveConnection,
msg: dict[str, Any],
) -> None:
"""Generate a QR code for a maintenance object or task."""
entry = _load_object_entry(hass, connection, msg)
if entry is None:
return
obj_data = entry.data.get(CONF_OBJECT, {})
task_id = msg.get("task_id")
task_name = None
if task_id:
tasks_data = entry.data.get(CONF_TASKS, {})
if task_id not in tasks_data:
connection.send_error(msg["id"], "not_found", "Task not found")
return
task_name = tasks_data[task_id].get("name", "")
action = msg.get("action", "view")
url_mode = msg.get("url_mode", "server")
base_url = msg.get("base_url")
try:
url = build_qr_url(
hass,
entry.entry_id,
task_id=task_id,
action=action,
base_url_override=base_url,
url_mode=url_mode,
)
except ValueError as err:
connection.send_error(msg["id"], "no_url", str(err))
return
from functools import partial
icon = _ACTION_ICON_MAP.get(action)
gen_fn = partial(generate_qr_svg_data_uri, url, border=2, icon=icon)
svg_data_uri = await hass.async_add_executor_job(gen_fn)
connection.send_result(
msg["id"],
{
"svg_data_uri": svg_data_uri,
"url": url,
"label": {
"object_name": obj_data.get(CONF_OBJECT_NAME, ""),
"manufacturer": obj_data.get(CONF_OBJECT_MANUFACTURER, ""),
"model": obj_data.get(CONF_OBJECT_MODEL, ""),
"task_name": task_name,
},
},
)
# Batch QR generation — used by the "Print QR codes" panel section.
#
# Typical household: 20-30 tasks × 2 actions = 40-60 QRs. Benchmarked at
# ~40 ms each with icon embed (HIGH ECC) → 2.5 s for 60, 7 s for 200.
# The raw SVG is ~32 KB each, so 200 × 32 KB = ~6 MB over the websocket;
# we cap at 200 to keep the payload bounded and the print layout sane
# (generous 6 QRs/A4 page = 34 pages).
_MAX_BATCH_QRS = 200
# LRU cache keyed on (url, icon). Two users printing the same task twice
# in a session hit this cache; so does re-running the batch after
# narrowing the filter. Bounded size so long-running HA instances with
# thousands of task-action combos can't grow the cache forever.
@lru_cache(maxsize=512)
def _cached_qr_svg(url: str, icon: str | None) -> str:
return generate_qr_svg(url, border=2, icon=icon)
@websocket_api.websocket_command(
{
vol.Required("type"): "maintenance_supporter/qr/batch_generate",
vol.Optional("entry_ids"): vol.All(
[vol.All(str, vol.Length(max=MAX_ID_LENGTH))],
vol.Length(max=1000),
),
vol.Optional("task_ids"): vol.All(
[vol.All(str, vol.Length(max=MAX_ID_LENGTH))],
vol.Length(max=2000),
),
vol.Required("actions"): vol.All(
[vol.In(["view", "complete", "skip", "quick_complete"])],
vol.Length(min=1, max=4),
),
vol.Optional("url_mode", default="server"): vol.In(["server", "local", "companion"]),
vol.Optional("base_url"): vol.All(vol.Url(), vol.Length(max=512)),
}
)
@websocket_api.async_response
async def ws_batch_generate_qr(
hass: HomeAssistant,
connection: websocket_api.ActiveConnection,
msg: dict[str, Any],
) -> None:
"""Generate multiple QR codes in one call for the print-all-QRs page.
Resolves (entry × task × action) combinations and returns SVG strings
ready to inline into a printable grid. Empty ``entry_ids`` / ``task_ids``
filters mean "all" at that level.
"""
# Resolve target entries (always exclude the global config entry).
all_entries = _get_object_entries(hass)
entry_filter = msg.get("entry_ids")
if entry_filter:
wanted = set(entry_filter)
entries = [e for e in all_entries if e.entry_id in wanted]
else:
entries = all_entries
# Build the flat (entry_id, object_name, task_id, task_name) target list,
# honouring the optional task_ids filter.
task_filter = set(msg["task_ids"]) if msg.get("task_ids") else None
targets: list[tuple[str, str, str, str]] = []
for entry in entries:
obj_name = entry.data.get(CONF_OBJECT, {}).get(CONF_OBJECT_NAME, "")
tasks_data = entry.data.get(CONF_TASKS, {})
for task_id, task_data in tasks_data.items():
if task_filter is not None and task_id not in task_filter:
continue
targets.append((entry.entry_id, obj_name, task_id, task_data.get("name", "")))
actions: list[str] = msg["actions"]
total = len(targets) * len(actions)
if total == 0:
connection.send_result(msg["id"], {"qrs": [], "total": 0})
return
if total > _MAX_BATCH_QRS:
connection.send_error(
msg["id"],
"too_many",
f"Batch would produce {total} QR codes; the per-request cap is "
f"{_MAX_BATCH_QRS}. Narrow the object/task/action filter.",
)
return
url_mode = msg.get("url_mode", "server")
base_url = msg.get("base_url")
# Generate URL first (fast), then offload the SVG encoding to the executor
# since it's CPU-bound (~30-40 ms/QR). Each SVG passes through the LRU
# cache so re-runs after a filter change are near-instant.
results: list[dict[str, Any]] = []
for entry_id, obj_name, task_id, task_name in targets:
for action in actions:
try:
url = build_qr_url(
hass,
entry_id,
task_id=task_id,
action=action,
base_url_override=base_url,
url_mode=url_mode,
)
except ValueError:
# No HA URL configured — skip this row rather than fail the
# whole batch. "server" mode is the only path that raises;
# "companion" and "local" always resolve.
continue
icon = _ACTION_ICON_MAP.get(action) # None for "skip" (no icon)
svg = await hass.async_add_executor_job(_cached_qr_svg, url, icon)
results.append(
{
"entry_id": entry_id,
"task_id": task_id,
"object_name": obj_name,
"task_name": task_name,
"action": action,
"svg": svg,
}
)
connection.send_result(msg["id"], {"qrs": results, "total": len(results)})