"""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_TASKS, DOMAIN, MAX_CHECKLIST_ITEM_LENGTH, MAX_CHECKLIST_ITEMS, MAX_ENTITY_SLUG_LENGTH, MAX_ID_LENGTH, MAX_IMPORT_PAYLOAD_BYTES, MAX_JSON_IMPORT_PAYLOAD_BYTES, MAX_VACATION_EXEMPT_TASKS, ) from ..helpers.aggregate import get_store, object_name from ..helpers.dates import normalize_hhmm, parse_iso_date 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, _load_object_task _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 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: return s if parse_iso_date(s) is not None else 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_`` 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_``, 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 ..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) and parse_iso_date(val) is not None: 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()][:MAX_VACATION_EXEMPT_TASKS] 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_`` 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", "notify_enabled", # #185: notification icon override (shape-checked below). "notify_icon", "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", # D#183: mirror targets (shape-sanitized below). "mirror_todo_entities", "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 and parse_iso_date(lp) is None: 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) # D#183: mirror targets — todo.* ids only, deduped, capped; an # empty result drops the key (same rules as the WS write paths). if task_data.get("mirror_todo_entities") is not None: from ..helpers.sanitize import sanitize_mirror_todo_entities mirrors = sanitize_mirror_todo_entities(task_data["mirror_todo_entities"]) if mirrors: task_data["mirror_todo_entities"] = mirrors else: task_data.pop("mirror_todo_entities", None) # #185: notify_icon — same shape rule as the WS write paths; a # malformed or empty value drops the override (type default). if "notify_icon" in task_data: from ..helpers.notify_icons import normalize_icon icon = normalize_icon(task_data["notify_icon"]) if icon: task_data["notify_icon"] = icon else: task_data.pop("notify_icon", None) # schedule_time: canonical HH:MM. The options flow's TimeSelector # stores "HH:MM:SS" and the export writes it verbatim — that used # to be DROPPED here (strict HH:MM), so a backup lost the time. st = task_data.get("schedule_time") if st is not None: normalized = normalize_hhmm(st) if normalized is None: task_data.pop("schedule_time", None) else: task_data["schedule_time"] = normalized # entity_slug: the WS create/update paths reject anything but # [a-z0-9_]+ (it becomes part of the entity_id); import copied the # value verbatim. Normalise to that alphabet (HA's slugify would # turn all-junk into "unknown"), drop it when nothing valid # remains, and say so — a changed slug changes the entity ids # (bug audit 2026-09-12). raw_slug = task_data.get("entity_slug") if raw_slug is not None: slug = ( re.sub(r"[^a-z0-9_]+", "_", raw_slug.strip().lower()).strip("_")[:MAX_ENTITY_SLUG_LENGTH] if isinstance(raw_slug, str) else "" ) if not slug: task_data.pop("entity_slug", None) task_warnings.append(f"{task_name}: entity_slug dropped — not [a-z0-9_]+") elif slug != raw_slug: task_data["entity_slug"] = slug task_warnings.append(f"{task_name}: entity_slug normalised to {slug!r}") # 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] = {} # Outside the per-object try below on purpose (the docs must # exist before the entry is created) — so a crash here used to # abort the WHOLE import without a reply. The store skips # malformed records itself; this backstop turns anything it # still raises into a per-object warning (bug audit 2026-09-12). try: 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 ) except Exception: # one object's documents must not sink the import _LOGGER.exception("JSON import of %s: documents skipped", obj_name) task_warnings.append("documents: skipped — malformed document records") await _drop_imported_documents(doc_store, obj_id) 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) store_new = get_store(hass, result["result"].entry_id) 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.""" task_id = msg.get("task_id") task_name = None if task_id: ctx = _load_object_task(hass, connection, msg) if ctx is None: return entry, _rd, task = ctx task_name = task.get("name", "") else: entry = _load_object_entry(hass, connection, msg) if entry is None: return obj_data = entry.data.get(CONF_OBJECT, {}) 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": object_name(entry), "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 = object_name(entry) 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)})