456 files
This commit is contained in:
@@ -27,6 +27,8 @@ from ..helpers.integration_signatures import (
|
||||
discover_integration_setups,
|
||||
)
|
||||
from ..helpers.permissions import require_write
|
||||
from ..helpers.reset_wiring import apply_reset_action
|
||||
from ..helpers.task_origin import ORIGIN_KEY, integration_origin
|
||||
from . import ID_FIELD
|
||||
from .adopt_batch import AdoptBatch
|
||||
|
||||
@@ -123,77 +125,78 @@ async def ws_adopt_integration_setups(
|
||||
|
||||
batch = AdoptBatch(hass)
|
||||
|
||||
for sel in msg["selections"]:
|
||||
device_id = sel["device_id"]
|
||||
setup = setups.get(device_id)
|
||||
if setup is None:
|
||||
batch.errors.append({"device_id": device_id, "reason": "no suggestion for this device"})
|
||||
continue
|
||||
wanted = set(sel.get("task_names") or [t["task_name"] for t in setup["tasks"]])
|
||||
# Keyed by (task_name, direction): one integration can ship a task name
|
||||
# in two directions (LG ThinQ filter: hours vs percent), and the trigger
|
||||
# must be built from the signature matching the discovered direction.
|
||||
sig_by_key = {
|
||||
(s.task_name, s.direction): s for s in SIGNATURES[setup["integration"]].tasks
|
||||
}
|
||||
batch.begin()
|
||||
try:
|
||||
entry_id = sel.get("entry_id") or setup["suggested_entry_id"]
|
||||
if not entry_id:
|
||||
entry_id = await batch.create_object(
|
||||
name=sel.get("object_name") or setup["suggested_object_name"],
|
||||
ha_device_id=device_id,
|
||||
)
|
||||
|
||||
entry = batch.target_entry(entry_id)
|
||||
if entry is None:
|
||||
batch.errors.append({"device_id": device_id, "reason": "target object not found"})
|
||||
# Reload the touched objects even when a selection blew up mid-way —
|
||||
# their tasks are already stored.
|
||||
try:
|
||||
for sel in msg["selections"]:
|
||||
device_id = sel["device_id"]
|
||||
setup = setups.get(device_id)
|
||||
if setup is None:
|
||||
batch.errors.append({"device_id": device_id, "reason": "no suggestion for this device"})
|
||||
continue
|
||||
|
||||
# #105: adopting into a user-picked existing object that isn't
|
||||
# device-bound yet — bind it, so model/sibling gates work and
|
||||
# future discovery suggests this object instead of a new one.
|
||||
# (Objects reached via suggested_entry_id are bound by definition;
|
||||
# the guard makes this a no-op for them.)
|
||||
obj_data = entry.data.get(CONF_OBJECT, {})
|
||||
if not obj_data.get("ha_device_id"):
|
||||
new_data = dict(entry.data)
|
||||
new_data[CONF_OBJECT] = {**obj_data, "ha_device_id": device_id}
|
||||
hass.config_entries.async_update_entry(entry, data=new_data)
|
||||
refreshed = hass.config_entries.async_get_entry(entry_id)
|
||||
if refreshed is not None:
|
||||
entry = refreshed
|
||||
|
||||
# Second dedup layer (discovery already hides these): never create
|
||||
# a task whose name — in any language — already exists on the
|
||||
# target object.
|
||||
from ..helpers.integration_signatures import proposal_name_variants
|
||||
|
||||
existing_names = {
|
||||
str(t.get("name", "")).lower()
|
||||
for t in entry.data.get(CONF_TASKS, {}).values()
|
||||
wanted = set(sel.get("task_names") or [t["task_name"] for t in setup["tasks"]])
|
||||
# Keyed by (task_name, direction): one integration can ship a task name
|
||||
# in two directions (LG ThinQ filter: hours vs percent), and the trigger
|
||||
# must be built from the signature matching the discovered direction.
|
||||
sig_by_key = {
|
||||
(s.task_name, s.direction): s for s in SIGNATURES[setup["integration"]].tasks
|
||||
}
|
||||
batch.begin()
|
||||
try:
|
||||
entry_id = sel.get("entry_id") or setup["suggested_entry_id"]
|
||||
if not entry_id:
|
||||
entry_id = await batch.create_object(
|
||||
name=sel.get("object_name") or setup["suggested_object_name"],
|
||||
ha_device_id=device_id,
|
||||
)
|
||||
|
||||
baselines = sel.get("baselines") or {}
|
||||
for task in setup["tasks"]:
|
||||
if task["task_name"] not in wanted:
|
||||
entry = batch.target_entry(entry_id)
|
||||
if entry is None:
|
||||
batch.errors.append({"device_id": device_id, "reason": "target object not found"})
|
||||
continue
|
||||
# Per-entity duties (colour cartridges) carry the entity label
|
||||
# in task_name; the catalog key and the label travel separately.
|
||||
catalog_name = task.get("catalog_task_name") or task["task_name"]
|
||||
if existing_names & proposal_name_variants(catalog_name, task.get("entity_label")):
|
||||
continue
|
||||
sig = sig_by_key[(catalog_name, task["direction"])]
|
||||
trigger = build_setup_trigger(sig, hass, task["entity_ids"])
|
||||
# #102: "last service was at reading X" — usage_delta only.
|
||||
# Other directions either already have absolute semantics
|
||||
# (usage_above's explicit 0) or no baseline concept at all.
|
||||
baseline = baselines.get(task["task_name"])
|
||||
if baseline is not None and sig.direction == "usage_delta":
|
||||
trigger["trigger_baseline_value"] = float(baseline)
|
||||
await batch.persist_task(
|
||||
entry,
|
||||
{
|
||||
|
||||
# #105: adopting into a user-picked existing object that isn't
|
||||
# device-bound yet — bind it, so model/sibling gates work and
|
||||
# future discovery suggests this object instead of a new one.
|
||||
# (Objects reached via suggested_entry_id are bound by definition;
|
||||
# the guard makes this a no-op for them.)
|
||||
obj_data = entry.data.get(CONF_OBJECT, {})
|
||||
if not obj_data.get("ha_device_id"):
|
||||
new_data = dict(entry.data)
|
||||
new_data[CONF_OBJECT] = {**obj_data, "ha_device_id": device_id}
|
||||
hass.config_entries.async_update_entry(entry, data=new_data)
|
||||
refreshed = hass.config_entries.async_get_entry(entry_id)
|
||||
if refreshed is not None:
|
||||
entry = refreshed
|
||||
|
||||
# Second dedup layer (discovery already hides these): never create
|
||||
# a task whose name — in any language — already exists on the
|
||||
# target object.
|
||||
from ..helpers.integration_signatures import proposal_name_variants
|
||||
|
||||
existing_names = {
|
||||
str(t.get("name", "")).lower()
|
||||
for t in entry.data.get(CONF_TASKS, {}).values()
|
||||
}
|
||||
|
||||
baselines = sel.get("baselines") or {}
|
||||
for task in setup["tasks"]:
|
||||
if task["task_name"] not in wanted:
|
||||
continue
|
||||
# Per-entity duties (colour cartridges) carry the entity label
|
||||
# in task_name; the catalog key and the label travel separately.
|
||||
catalog_name = task.get("catalog_task_name") or task["task_name"]
|
||||
if existing_names & proposal_name_variants(catalog_name, task.get("entity_label")):
|
||||
continue
|
||||
sig = sig_by_key[(catalog_name, task["direction"])]
|
||||
trigger = build_setup_trigger(sig, hass, task["entity_ids"])
|
||||
# #102: "last service was at reading X" — usage_delta only.
|
||||
# Other directions either already have absolute semantics
|
||||
# (usage_above's explicit 0) or no baseline concept at all.
|
||||
baseline = baselines.get(task["task_name"])
|
||||
if baseline is not None and sig.direction == "usage_delta":
|
||||
trigger["trigger_baseline_value"] = float(baseline)
|
||||
task_data: dict[str, Any] = {
|
||||
"id": uuid4().hex,
|
||||
"object_id": entry.data.get(CONF_OBJECT, {}).get("id", ""),
|
||||
"name": task.get("task_name_localized")
|
||||
@@ -203,12 +206,61 @@ async def ws_adopt_integration_setups(
|
||||
"enabled": True,
|
||||
"schedule": {"kind": "manual"},
|
||||
"trigger_config": trigger,
|
||||
},
|
||||
)
|
||||
except (ValueError, KeyError) as err:
|
||||
# Removes an object created for this device together with the
|
||||
# tasks already persisted into it — and un-counts both (the
|
||||
# tasks were over-reported before the DRY audit 2026-09-26).
|
||||
await batch.fail({"device_id": device_id, "reason": str(err)})
|
||||
|
||||
# 2.95: the fingerprint — recognised as this duty after a rename.
|
||||
ORIGIN_KEY: integration_origin(setup["integration"], sig.task_name, sig.direction, device_id, task.get("entity_label")),
|
||||
}
|
||||
# 2.95: completing the task also resets the integration's
|
||||
# own counter (its reset button, enabled if it shipped off).
|
||||
if task.get("reset"):
|
||||
apply_reset_action(hass, task_data, task["reset"]["entity_id"], connection.user.id if connection.user else None)
|
||||
await batch.persist_task(entry, task_data)
|
||||
# One duty through two sensors (a softener's salt in % and in
|
||||
# days) is one task — the second proposal of the name is
|
||||
# skipped instead of creating a twin.
|
||||
existing_names.add(str(task_data["name"]).lower())
|
||||
except (ValueError, KeyError) as err:
|
||||
# Removes an object created for this device together with the
|
||||
# tasks already persisted into it — and un-counts both (the
|
||||
# tasks were over-reported before the DRY audit 2026-09-26).
|
||||
await batch.fail({"device_id": device_id, "reason": str(err)})
|
||||
finally:
|
||||
await batch.finish()
|
||||
connection.send_result(msg["id"], batch.result())
|
||||
|
||||
|
||||
@websocket_api.websocket_command({vol.Required("type"): f"{DOMAIN}/integration_setups/reset_offers"})
|
||||
@websocket_api.async_response
|
||||
async def ws_reset_offers(
|
||||
hass: HomeAssistant,
|
||||
connection: websocket_api.ActiveConnection,
|
||||
msg: dict[str, Any],
|
||||
) -> None:
|
||||
"""2.95: existing tasks whose completion could also reset the
|
||||
integration's own counter (helpers/reset_wiring.reset_offers)."""
|
||||
from ..helpers.reset_wiring import reset_offers
|
||||
|
||||
connection.send_result(msg["id"], {"offers": reset_offers(hass)})
|
||||
|
||||
|
||||
@websocket_api.websocket_command(
|
||||
{
|
||||
vol.Required("type"): f"{DOMAIN}/integration_setups/wire_resets",
|
||||
vol.Required("items"): vol.All(
|
||||
[{vol.Required("entry_id"): ID_FIELD, vol.Required("task_id"): ID_FIELD}],
|
||||
vol.Length(min=1, max=200),
|
||||
),
|
||||
}
|
||||
)
|
||||
@require_write
|
||||
@websocket_api.async_response
|
||||
async def ws_wire_resets(
|
||||
hass: HomeAssistant,
|
||||
connection: websocket_api.ActiveConnection,
|
||||
msg: dict[str, Any],
|
||||
) -> None:
|
||||
"""2.95: wire the chosen offers — the reset button becomes the task's
|
||||
completion action (enabled when the integration shipped it off)."""
|
||||
from ..helpers.reset_wiring import wire_resets
|
||||
|
||||
count = wire_resets(hass, msg["items"], connection.user.id if connection.user else None)
|
||||
connection.send_result(msg["id"], {"wired": count})
|
||||
|
||||
Reference in New Issue
Block a user