523 lines
21 KiB
Python
523 lines
21 KiB
Python
"""Task action WS handlers: complete / quick_complete / skip / reset."""
|
|
|
|
from __future__ import annotations
|
|
|
|
from datetime import datetime
|
|
from typing import Any
|
|
|
|
import voluptuous as vol
|
|
from homeassistant.components import websocket_api
|
|
from homeassistant.core import HomeAssistant
|
|
from homeassistant.exceptions import ServiceValidationError
|
|
from homeassistant.util import dt as dt_util
|
|
|
|
from ..const import (
|
|
MAX_CHECKLIST_ITEM_LENGTH,
|
|
MAX_CHECKLIST_ITEMS,
|
|
MAX_COST,
|
|
MAX_DATE_LENGTH,
|
|
MAX_DURATION_MINUTES,
|
|
MAX_TEXT_LENGTH,
|
|
MAX_TIMESTAMP_LENGTH,
|
|
)
|
|
from ..helpers.completion_photos import MAX_COMPLETION_PHOTOS, normalize_photo_doc_ids
|
|
from ..models.maintenance_task import MaintenanceTask
|
|
from . import (
|
|
ID_FIELD,
|
|
READING_VALUES_FIELD,
|
|
USED_PARTS_FIELD,
|
|
_load_object_task,
|
|
_parse_iso_date,
|
|
async_commit_store,
|
|
)
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Task Actions (Complete / Skip / Reset)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
def _completion_blocked(rd: Any, task_id: str) -> bool:
|
|
"""True iff the task's completion window forbids completing it right now.
|
|
|
|
Uses the coordinator's merged (live) task data so ``last_performed`` /
|
|
``next_due`` reflect the Store, not the stale config entry.
|
|
"""
|
|
coordinator = getattr(rd, "coordinator", None)
|
|
if coordinator is None:
|
|
return False
|
|
merged = coordinator._get_merged_tasks_data()
|
|
td = merged.get(task_id)
|
|
if not td:
|
|
return False
|
|
return not MaintenanceTask.from_dict(td).can_complete_now
|
|
|
|
|
|
def _refuse_too_early(connection: websocket_api.ActiveConnection, msg: dict[str, Any], rd: Any) -> bool:
|
|
"""Send ``too_early`` and return True when the completion window forbids
|
|
completing the task now — the ONE pre-check shared by the complete and
|
|
quick-complete commands (the coordinator choke point raises the same key
|
|
for callers that skip it, e.g. the HA service)."""
|
|
if not _completion_blocked(rd, msg["task_id"]):
|
|
return False
|
|
connection.send_error(msg["id"], "too_early", "Task can only be completed closer to its due date")
|
|
return True
|
|
|
|
|
|
@websocket_api.websocket_command(
|
|
{
|
|
vol.Required("type"): "maintenance_supporter/task/complete",
|
|
vol.Required("entry_id"): ID_FIELD,
|
|
vol.Required("task_id"): ID_FIELD,
|
|
vol.Optional("notes"): vol.Any(vol.All(str, vol.Length(max=MAX_TEXT_LENGTH)), None),
|
|
vol.Optional("cost"): vol.Any(vol.All(vol.Coerce(float), vol.Range(min=0, max=MAX_COST)), None),
|
|
vol.Optional("duration"): vol.Any(vol.All(vol.Coerce(int), vol.Range(min=0, max=MAX_DURATION_MINUTES)), None),
|
|
# Restrict checklist_state to {string-key (≤500): bool, ...} with
|
|
# a hard cap on entries. Without this, attackers (or bad clients)
|
|
# could inflate the per-task history with arbitrarily large dicts.
|
|
vol.Optional("checklist_state"): vol.Any(
|
|
vol.All(
|
|
{vol.All(str, vol.Length(max=MAX_CHECKLIST_ITEM_LENGTH)): bool},
|
|
vol.Length(max=MAX_CHECKLIST_ITEMS),
|
|
),
|
|
None,
|
|
),
|
|
vol.Optional("feedback"): vol.Any(vol.All(str, vol.Length(max=MAX_TEXT_LENGTH)), None),
|
|
# #133: when the maintenance was actually performed (ISO datetime,
|
|
# naive = local). Backfills a past completion; must not be in the
|
|
# future — the coordinator validates and splits latest-vs-backfill.
|
|
vol.Optional("completed_at"): vol.Any(vol.All(str, vol.Length(max=MAX_TIMESTAMP_LENGTH)), None),
|
|
# QR deep-link fallback: the panel asserts the scan on a tag-gated task
|
|
# whose quick-complete needs the full dialog (bug audit 2026-08-29).
|
|
vol.Optional("via_tag_scan"): bool,
|
|
# Optional completion photos (#161): doc_ids of already-uploaded
|
|
# images (via the document upload endpoint, tagged "photo"). The
|
|
# scalar form is what pre-2.75 clients send — merged into the list.
|
|
vol.Optional("photo_doc_ids"): vol.Any(
|
|
vol.All([ID_FIELD], vol.Length(max=MAX_COMPLETION_PHOTOS)),
|
|
None,
|
|
),
|
|
vol.Optional("photo_doc_id"): vol.Any(ID_FIELD, None),
|
|
# Meter readings (v2.20, #83): the recorded value for `reading` tasks.
|
|
# Wide numeric bounds — meters count high, temperatures go negative.
|
|
vol.Optional("reading_value"): vol.Any(vol.All(vol.Coerce(float), vol.Range(min=-1e12, max=1e12)), None),
|
|
# #161 phase 2: {slot_id: value} for a task with reading slots; None
|
|
# skips a meter this time (shared shape with the history edit).
|
|
vol.Optional("reading_values"): READING_VALUES_FIELD,
|
|
# Spare parts: on an auto-created "buy" task, how many units were
|
|
# actually bought (dialog override of the part's restock_quantity).
|
|
vol.Optional("restock_quantity"): vol.Any(vol.All(vol.Any(int, float), vol.Coerce(float), vol.Range(min=0.01, max=9999)), None),
|
|
# #99: the parts actually used on THIS completion. An explicit list
|
|
# (even an empty one) REPLACES the task's automatic consumes_parts
|
|
# deduction; omitting the key keeps the automatic behaviour.
|
|
vol.Optional("used_parts"): USED_PARTS_FIELD,
|
|
}
|
|
)
|
|
@websocket_api.async_response
|
|
async def ws_complete_task(
|
|
hass: HomeAssistant,
|
|
connection: websocket_api.ActiveConnection,
|
|
msg: dict[str, Any],
|
|
) -> None:
|
|
"""Mark a task as completed."""
|
|
ctx = _load_object_task(hass, connection, msg, merged=True, need_coordinator=True)
|
|
if ctx is None:
|
|
return
|
|
_entry, rd, slot_task = ctx
|
|
|
|
# #133: an optional backdated completion moment. Parsed here (the schema
|
|
# can only cheaply cap the length); range/future validation lives in the
|
|
# coordinator choke point shared with the HA service.
|
|
completed_at: datetime | None = None
|
|
if msg.get("completed_at"):
|
|
completed_at = dt_util.parse_datetime(msg["completed_at"])
|
|
if completed_at is None:
|
|
connection.send_error(
|
|
msg["id"],
|
|
"invalid_format",
|
|
f"completed_at is not a valid ISO datetime: {msg['completed_at']!r}",
|
|
)
|
|
return
|
|
|
|
# The earliest-completion window guards against doing the work too early
|
|
# — recording a completion that already happened on a PAST day is a
|
|
# history correction, not early work, so it bypasses the gate. A
|
|
# today-dated completed_at still honours it.
|
|
is_past_dated = completed_at is not None and completed_at.date() < dt_util.now().date()
|
|
if not is_past_dated and _refuse_too_early(connection, msg, rd):
|
|
return
|
|
|
|
# #99: validate the per-completion selection against the object's parts
|
|
# (unknown ids drop, quantities round) — absent stays absent so the
|
|
# automatic consumes_parts path is untouched.
|
|
used_parts = msg.get("used_parts")
|
|
if used_parts is not None:
|
|
from ..const import CONF_PARTS
|
|
from ..helpers.parts import sanitize_consumes_parts
|
|
from . import foreign_part_resolver
|
|
|
|
used_parts = sanitize_consumes_parts(
|
|
used_parts,
|
|
set(_entry.data.get(CONF_PARTS) or {}),
|
|
foreign_part_ids=foreign_part_resolver(hass),
|
|
)
|
|
|
|
# #161 phase 2: resolve the slot values into the history snapshot.
|
|
reading_values = None
|
|
if msg.get("reading_values"):
|
|
from ..helpers.reading_slots import resolve_reading_values
|
|
|
|
try:
|
|
reading_values = (
|
|
resolve_reading_values(
|
|
slot_task.get("readings") or [], msg["reading_values"], default_unit=slot_task.get("reading_unit")
|
|
)
|
|
or None
|
|
)
|
|
except ValueError as err:
|
|
connection.send_error(msg["id"], "invalid_input", str(err))
|
|
return
|
|
|
|
# Only this object's uploaded files count as completion photos — a
|
|
# foreign doc id satisfied a required photo and got linked (and later
|
|
# re-homed by task/move). This command is the one completion surface
|
|
# that carries photos (bug audit 2026-09-26).
|
|
from ..helpers.completion_requirements import own_photo_doc_ids
|
|
from . import object_id_for_entry
|
|
|
|
photo_doc_ids = own_photo_doc_ids(
|
|
hass, object_id_for_entry(_entry), normalize_photo_doc_ids(msg.get("photo_doc_ids"), msg.get("photo_doc_id"))
|
|
)
|
|
|
|
try:
|
|
await rd.coordinator.complete_maintenance(
|
|
source="panel",
|
|
task_id=msg["task_id"],
|
|
notes=msg.get("notes"),
|
|
cost=msg.get("cost"),
|
|
duration=msg.get("duration"),
|
|
checklist_state=msg.get("checklist_state"),
|
|
feedback=msg.get("feedback"),
|
|
photo_doc_ids=photo_doc_ids or None,
|
|
reading_value=msg.get("reading_value"),
|
|
reading_values=reading_values,
|
|
restock_quantity=msg.get("restock_quantity"),
|
|
used_parts=used_parts,
|
|
completed_at=completed_at,
|
|
# Who did it: taken from the authenticated connection, never from
|
|
# the payload — a client must not be able to credit someone else.
|
|
# This is also what feeds the `least_completed` rotation strategy
|
|
# and what satisfies a task requiring the "user" detail.
|
|
completed_by=connection.user.id if connection.user else None,
|
|
# The panel sets this only on the QR deep-link fallback (a scanned
|
|
# sticker whose task has no quick-complete defaults) - same trust
|
|
# level as task/quick_complete (bug audit 2026-08-29).
|
|
tag_verified=bool(msg.get("via_tag_scan")),
|
|
)
|
|
except ServiceValidationError as err:
|
|
# Validation refusals (missing required details, future completed_at).
|
|
# The dialog normally prevents these, so reaching here means an
|
|
# older/cached frontend or a scripted call — answer with the
|
|
# exception's own key rather than a traceback.
|
|
connection.send_error(msg["id"], err.translation_key or "completion_details_required", str(err))
|
|
return
|
|
connection.send_result(msg["id"], {"success": True})
|
|
|
|
|
|
# v1.3.0: One-tap completion using values pre-configured on the task.
|
|
# Used by the "quick_complete" QR scan path. Falls back with `no_defaults`
|
|
# error when the task has no quick_complete_defaults — frontend then
|
|
# routes the user to the normal complete dialog.
|
|
@websocket_api.websocket_command(
|
|
{
|
|
vol.Required("type"): "maintenance_supporter/task/quick_complete",
|
|
vol.Required("entry_id"): ID_FIELD,
|
|
vol.Required("task_id"): ID_FIELD,
|
|
}
|
|
)
|
|
@websocket_api.async_response
|
|
async def ws_quick_complete_task(
|
|
hass: HomeAssistant,
|
|
connection: websocket_api.ActiveConnection,
|
|
msg: dict[str, Any],
|
|
) -> None:
|
|
"""Complete a task using its pre-configured `quick_complete_defaults`."""
|
|
ctx = _load_object_task(hass, connection, msg, need_coordinator=True)
|
|
if ctx is None:
|
|
return
|
|
_entry, rd, task = ctx
|
|
|
|
if _refuse_too_early(connection, msg, rd):
|
|
return
|
|
|
|
defaults = task.get("quick_complete_defaults") or {}
|
|
if not isinstance(defaults, dict) or not defaults:
|
|
# Frontend fallback: open the normal complete dialog so the user
|
|
# is never stuck staring at a useless QR scan.
|
|
connection.send_error(
|
|
msg["id"],
|
|
"no_defaults",
|
|
"Task has no quick_complete_defaults; open complete dialog instead",
|
|
)
|
|
return
|
|
|
|
try:
|
|
await rd.coordinator.complete_maintenance(
|
|
source="qr",
|
|
task_id=msg["task_id"],
|
|
notes=defaults.get("notes"),
|
|
cost=defaults.get("cost"),
|
|
duration=defaults.get("duration"),
|
|
feedback=defaults.get("feedback"),
|
|
# The QR sticker hangs ON the thing — scanning it is presence.
|
|
tag_verified=True,
|
|
completed_by=connection.user.id if connection.user else None,
|
|
)
|
|
except ServiceValidationError as err:
|
|
# The task demands details the quick-complete defaults do not cover —
|
|
# same fallback as `no_defaults`: the caller opens the full dialog.
|
|
# Other refusals (too_early, task_inactive) keep their own key.
|
|
connection.send_error(msg["id"], err.translation_key or "completion_details_required", str(err))
|
|
return
|
|
connection.send_result(msg["id"], {"success": True, "via": "quick"})
|
|
|
|
|
|
@websocket_api.websocket_command(
|
|
{
|
|
vol.Required("type"): "maintenance_supporter/task/skip",
|
|
vol.Required("entry_id"): ID_FIELD,
|
|
vol.Required("task_id"): ID_FIELD,
|
|
vol.Optional("reason"): vol.Any(vol.All(str, vol.Length(max=MAX_TEXT_LENGTH)), None),
|
|
# Record the skipped cycle as MISSED (was due, never done) rather than a
|
|
# deliberate skip — clearer history + compliance views.
|
|
vol.Optional("as_missed", default=False): bool,
|
|
}
|
|
)
|
|
@websocket_api.async_response
|
|
async def ws_skip_task(
|
|
hass: HomeAssistant,
|
|
connection: websocket_api.ActiveConnection,
|
|
msg: dict[str, Any],
|
|
) -> None:
|
|
"""Skip the current maintenance cycle."""
|
|
ctx = _load_object_task(hass, connection, msg, need_coordinator=True)
|
|
if ctx is None:
|
|
return
|
|
_entry, rd, _task = ctx
|
|
|
|
try:
|
|
await rd.coordinator.skip_maintenance(
|
|
task_id=msg["task_id"],
|
|
reason=msg.get("reason"),
|
|
as_missed=msg.get("as_missed", False),
|
|
)
|
|
except ServiceValidationError as err:
|
|
# #150: the task carries a skip lock; an inactive task keeps its own
|
|
# key (task_inactive_skip) so the panel can say why.
|
|
connection.send_error(msg["id"], err.translation_key or "skip_disabled", str(err))
|
|
return
|
|
connection.send_result(msg["id"], {"success": True})
|
|
|
|
|
|
@websocket_api.websocket_command(
|
|
{
|
|
vol.Required("type"): "maintenance_supporter/task/reset",
|
|
vol.Required("entry_id"): ID_FIELD,
|
|
vol.Required("task_id"): ID_FIELD,
|
|
vol.Optional("date"): vol.Any(vol.All(str, vol.Length(max=MAX_DATE_LENGTH)), None),
|
|
}
|
|
)
|
|
@websocket_api.async_response
|
|
async def ws_reset_task(
|
|
hass: HomeAssistant,
|
|
connection: websocket_api.ActiveConnection,
|
|
msg: dict[str, Any],
|
|
) -> None:
|
|
"""Reset the last performed date."""
|
|
ctx = _load_object_task(hass, connection, msg, need_coordinator=True)
|
|
if ctx is None:
|
|
return
|
|
_entry, rd, _task = ctx
|
|
|
|
reset_date = None
|
|
if msg.get("date"):
|
|
# A reset marks a day the work WAS done — a future one (read tier:
|
|
# any household member) overflowed the schedule math and took the
|
|
# object down (bug audit 2026-09-27).
|
|
reset_date = _parse_iso_date(connection, msg["id"], msg["date"], field="date", not_future=True)
|
|
if reset_date is None:
|
|
return
|
|
|
|
try:
|
|
await rd.coordinator.reset_maintenance(
|
|
task_id=msg["task_id"],
|
|
date=reset_date,
|
|
)
|
|
except ServiceValidationError as err:
|
|
# An archived / disabled / paused task keeps its own key (task_inactive).
|
|
connection.send_error(msg["id"], err.translation_key or "task_inactive", str(err))
|
|
return
|
|
connection.send_result(msg["id"], {"success": True})
|
|
|
|
|
|
@websocket_api.websocket_command(
|
|
{
|
|
vol.Required("type"): "maintenance_supporter/task/set_phase",
|
|
vol.Required("entry_id"): ID_FIELD,
|
|
vol.Required("task_id"): ID_FIELD,
|
|
# Index into phase_sequence — which cycle step is due NEXT.
|
|
vol.Required("cursor"): vol.All(int, vol.Range(min=0, max=100)),
|
|
}
|
|
)
|
|
@websocket_api.async_response
|
|
async def ws_set_task_phase(
|
|
hass: HomeAssistant,
|
|
connection: websocket_api.ActiveConnection,
|
|
msg: dict[str, Any],
|
|
) -> None:
|
|
"""Set which phase of a cyclic task is due next (#139).
|
|
|
|
The explicit correction path ("I'm actually at step 3") — completions
|
|
rotate the cursor themselves and there is deliberately no phase picker
|
|
in the complete dialog.
|
|
"""
|
|
ctx = _load_object_task(hass, connection, msg, need_store=True, need_coordinator=True)
|
|
if ctx is None:
|
|
return
|
|
_entry, rd, task = ctx
|
|
sequence = task.get("phase_sequence") or []
|
|
if not (task.get("phases") and sequence):
|
|
connection.send_error(msg["id"], "no_phases", "Task has no phase cycle")
|
|
return
|
|
if msg["cursor"] >= len(sequence):
|
|
connection.send_error(
|
|
msg["id"], "invalid_cursor", f"cursor must be < {len(sequence)}"
|
|
)
|
|
return
|
|
rd.store.set_phase_cursor(msg["task_id"], msg["cursor"])
|
|
await async_commit_store(rd)
|
|
connection.send_result(msg["id"], {"success": True})
|
|
|
|
|
|
@websocket_api.websocket_command(
|
|
{
|
|
vol.Required("type"): "maintenance_supporter/task/postpone",
|
|
vol.Required("entry_id"): ID_FIELD,
|
|
vol.Required("task_id"): ID_FIELD,
|
|
vol.Required("until"): vol.All(str, vol.Length(max=MAX_DATE_LENGTH)),
|
|
}
|
|
)
|
|
@websocket_api.async_response
|
|
async def ws_postpone_task(
|
|
hass: HomeAssistant,
|
|
connection: websocket_api.ActiveConnection,
|
|
msg: dict[str, Any],
|
|
) -> None:
|
|
"""Postpone the current occurrence to a chosen date (per-occurrence defer)."""
|
|
ctx = _load_object_task(hass, connection, msg, need_coordinator=True)
|
|
if ctx is None:
|
|
return
|
|
_entry, rd, _task = ctx
|
|
|
|
until = _parse_iso_date(connection, msg["id"], msg["until"], field="until")
|
|
if until is None:
|
|
return
|
|
# Bug audit 2026-09-27: an override in year 9999 overflowed the calendar
|
|
# entity's "next day" — a postpone reaches at most one maximum interval out.
|
|
from datetime import timedelta
|
|
|
|
from ..const import MAX_INTERVAL_DAYS
|
|
|
|
if until > dt_util.now().date() + timedelta(days=MAX_INTERVAL_DAYS):
|
|
connection.send_error(msg["id"], "invalid_date", f"until must be within {MAX_INTERVAL_DAYS} days")
|
|
return
|
|
|
|
try:
|
|
await rd.coordinator.async_postpone_task(msg["task_id"], until)
|
|
except ServiceValidationError as err:
|
|
connection.send_error(msg["id"], err.translation_key or "task_inactive", str(err))
|
|
return
|
|
connection.send_result(msg["id"], {"success": True})
|
|
|
|
|
|
@websocket_api.websocket_command(
|
|
{
|
|
vol.Required("type"): "maintenance_supporter/task/snooze",
|
|
vol.Required("entry_id"): ID_FIELD,
|
|
vol.Required("task_id"): ID_FIELD,
|
|
}
|
|
)
|
|
@websocket_api.async_response
|
|
async def ws_snooze_task(
|
|
hass: HomeAssistant,
|
|
connection: websocket_api.ActiveConnection,
|
|
msg: dict[str, Any],
|
|
) -> None:
|
|
"""Snooze a task's notifications for the configured snooze duration.
|
|
|
|
Surfaces the existing notification-action snooze on the panel. Suppresses
|
|
due-soon/overdue/triggered reminders for ``snooze_duration_hours`` — it does
|
|
not change the task's schedule or state.
|
|
"""
|
|
from ..const import DOMAIN, NOTIFICATION_MANAGER_KEY
|
|
|
|
if _load_object_task(hass, connection, msg) is None:
|
|
return
|
|
|
|
nm = hass.data.get(DOMAIN, {}).get(NOTIFICATION_MANAGER_KEY)
|
|
if nm is None:
|
|
connection.send_error(msg["id"], "unavailable", "Notifications not configured")
|
|
return
|
|
nm.snooze_task(msg["entry_id"], msg["task_id"])
|
|
connection.send_result(msg["id"], {"success": True})
|
|
|
|
|
|
@websocket_api.websocket_command(
|
|
{
|
|
vol.Required("type"): "maintenance_supporter/task/checklist_progress",
|
|
vol.Required("entry_id"): ID_FIELD,
|
|
vol.Required("task_id"): ID_FIELD,
|
|
# Same shape/caps as task/complete's checklist_state.
|
|
vol.Required("checklist_state"): vol.All(
|
|
{vol.All(str, vol.Length(max=MAX_CHECKLIST_ITEM_LENGTH)): bool},
|
|
vol.Length(max=MAX_CHECKLIST_ITEMS),
|
|
),
|
|
}
|
|
)
|
|
@websocket_api.async_response
|
|
async def ws_checklist_progress(
|
|
hass: HomeAssistant,
|
|
connection: websocket_api.ActiveConnection,
|
|
msg: dict[str, Any],
|
|
) -> None:
|
|
"""#73: persist in-cycle checklist ticks WITHOUT completing the task.
|
|
|
|
The dict REPLACES the stored progress (the client sends the full current
|
|
state — idempotent, no per-item race). Keys are the item TEXTS, so ticks
|
|
survive reordering and a renamed step drops its tick. Unknown keys (items
|
|
no longer on the checklist) are dropped on write; completing or skipping
|
|
the task clears the progress for the next cycle.
|
|
|
|
Deliberately the same tier as task/complete (plain authenticated): ticking
|
|
a step IS doing the work — the same household member who may complete the
|
|
task must be able to record partial progress.
|
|
"""
|
|
from ..helpers.phases import effective_field
|
|
|
|
# Merged view: the phase cursor lives in the Store. A phased task (#139)
|
|
# shows the CURRENT phase's checklist; validating against the task-level
|
|
# list dropped every tick on it (bug audit 2026-09-26).
|
|
ctx = _load_object_task(hass, connection, msg, merged=True, need_coordinator=True)
|
|
if ctx is None:
|
|
return
|
|
_entry, rd, task = ctx
|
|
task_items = set(effective_field(task, "checklist") or [])
|
|
state = {item: bool(done) for item, done in msg["checklist_state"].items() if item in task_items}
|
|
# Progress lives ONLY in the Store (no legacy fallback) — degrade to a
|
|
# clean error instead of an AttributeError when it failed to load.
|
|
if rd.store is None:
|
|
connection.send_error(msg["id"], "storage_unavailable", "Task storage not loaded")
|
|
return
|
|
rd.store.set_checklist_progress(msg["task_id"], state)
|
|
await async_commit_store(rd)
|
|
connection.send_result(msg["id"], {"success": True, "checklist_state": state})
|