BEHAVIOR-001: learn and predict actuator actions
This commit is contained in:
@@ -5,6 +5,7 @@ from dataclasses import dataclass
|
||||
from datetime import datetime
|
||||
import json
|
||||
import re
|
||||
from typing import Any
|
||||
from urllib.parse import quote
|
||||
|
||||
import requests
|
||||
@@ -19,6 +20,7 @@ from app.ha.exceptions import (
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
_ENTITY_ID_PATTERN = re.compile(r"^[a-z0-9_]+\.[a-z0-9_]+$")
|
||||
_SERVICE_PART_PATTERN = re.compile(r"^[a-z0-9_]+$")
|
||||
_MAX_HISTORY_SECONDS = 31 * 24 * 60 * 60
|
||||
|
||||
|
||||
@@ -84,6 +86,44 @@ class HaClient:
|
||||
)
|
||||
return payload
|
||||
|
||||
def get_logbook(
|
||||
self,
|
||||
entity_id: str,
|
||||
start_time: datetime,
|
||||
end_time: datetime,
|
||||
) -> list[object]:
|
||||
self._validate_period([entity_id], start_time, end_time)
|
||||
start = quote(start_time.isoformat(), safe=":+")
|
||||
payload = self._get_json(
|
||||
f"/api/logbook/{start}",
|
||||
params={
|
||||
"entity": entity_id,
|
||||
"end_time": end_time.isoformat(),
|
||||
},
|
||||
)
|
||||
if not isinstance(payload, list):
|
||||
raise HaUnexpectedPayloadError(
|
||||
"Logbook-Antwort von Home Assistant hat unerwartetes Format."
|
||||
)
|
||||
return payload
|
||||
|
||||
def call_service(
|
||||
self,
|
||||
domain: str,
|
||||
service: str,
|
||||
service_data: dict[str, object],
|
||||
) -> list[object]:
|
||||
if not _SERVICE_PART_PATTERN.fullmatch(domain):
|
||||
raise ValueError("Ungültige Service-Domain.")
|
||||
if not _SERVICE_PART_PATTERN.fullmatch(service):
|
||||
raise ValueError("Ungültiger Service-Name.")
|
||||
payload = self._post_json(f"/api/services/{domain}/{service}", service_data)
|
||||
if not isinstance(payload, list):
|
||||
raise HaUnexpectedPayloadError(
|
||||
"Service-Antwort von Home Assistant hat unerwartetes Format."
|
||||
)
|
||||
return payload
|
||||
|
||||
def list_entity_metadata(self, entity_ids: list[str]) -> dict[str, dict[str, str | None]]:
|
||||
if not entity_ids:
|
||||
return {}
|
||||
@@ -153,6 +193,36 @@ class HaClient:
|
||||
|
||||
return payload
|
||||
|
||||
def _post_json(self, path: str, payload: Any) -> object:
|
||||
try:
|
||||
response = self._session.post(
|
||||
f"{self._settings.url.rstrip('/')}{path}",
|
||||
json=payload,
|
||||
timeout=self._settings.timeout_seconds,
|
||||
)
|
||||
except requests.Timeout as exc:
|
||||
raise HaTimeoutError("Zeitüberschreitung beim Zugriff auf Home Assistant.") from exc
|
||||
except requests.RequestException as exc:
|
||||
raise HaHttpError(
|
||||
getattr(getattr(exc, "response", None), "status_code", 502),
|
||||
"Netzwerkfehler beim Zugriff auf Home Assistant.",
|
||||
) from exc
|
||||
if response.status_code in (401, 403):
|
||||
raise HaAuthError(
|
||||
response.status_code,
|
||||
"Authentifizierung bei Home Assistant fehlgeschlagen.",
|
||||
)
|
||||
try:
|
||||
response.raise_for_status()
|
||||
except requests.HTTPError as exc:
|
||||
raise HaHttpError(response.status_code, "Home Assistant meldet einen Fehler.") from exc
|
||||
try:
|
||||
return response.json()
|
||||
except ValueError as exc:
|
||||
raise HaUnexpectedPayloadError(
|
||||
"Antwort von Home Assistant ist kein gültiges JSON."
|
||||
) from exc
|
||||
|
||||
def _post_text(self, path: str, payload: dict[str, str]) -> str:
|
||||
try:
|
||||
response = self._session.post(
|
||||
@@ -179,6 +249,25 @@ class HaClient:
|
||||
raise HaHttpError(response.status_code, "Home Assistant meldet einen Fehler.") from exc
|
||||
return response.text
|
||||
|
||||
@staticmethod
|
||||
def _validate_period(
|
||||
entity_ids: list[str],
|
||||
start_time: datetime,
|
||||
end_time: datetime,
|
||||
) -> None:
|
||||
if not entity_ids:
|
||||
raise ValueError("Mindestens eine entity_id ist erforderlich.")
|
||||
if len(entity_ids) > 100:
|
||||
raise ValueError("Es können höchstens 100 Entities abgefragt werden.")
|
||||
if any(not _ENTITY_ID_PATTERN.fullmatch(entity_id) for entity_id in entity_ids):
|
||||
raise ValueError("entity_id enthält ein ungültiges Format.")
|
||||
if start_time.tzinfo is None or end_time.tzinfo is None:
|
||||
raise ValueError("start_time und end_time müssen eine Zeitzone enthalten.")
|
||||
if end_time <= start_time:
|
||||
raise ValueError("end_time muss nach start_time liegen.")
|
||||
if (end_time - start_time).total_seconds() > _MAX_HISTORY_SECONDS:
|
||||
raise ValueError("History-Abfragen sind auf 31 Tage begrenzt.")
|
||||
|
||||
|
||||
def _metadata_template(entity_ids: list[str]) -> str:
|
||||
ids = json.dumps(entity_ids, ensure_ascii=True)
|
||||
|
||||
@@ -18,6 +18,25 @@ class EntityHistorySeries(BaseModel):
|
||||
points: list[NumericHistoryPoint]
|
||||
|
||||
|
||||
class StateHistoryPoint(BaseModel):
|
||||
timestamp: datetime
|
||||
state: str
|
||||
|
||||
|
||||
class StateHistorySeries(BaseModel):
|
||||
entity_id: str
|
||||
points: list[StateHistoryPoint]
|
||||
|
||||
|
||||
class LogbookEntry(BaseModel):
|
||||
entity_id: str
|
||||
timestamp: datetime
|
||||
message: str = ""
|
||||
context_user_id: str | None = None
|
||||
context_domain: str | None = None
|
||||
context_service: str | None = None
|
||||
|
||||
|
||||
def normalize_history_payload(payload: object) -> list[EntityHistorySeries]:
|
||||
if not isinstance(payload, list):
|
||||
raise HaUnexpectedPayloadError("History-Payload muss eine Liste sein.")
|
||||
@@ -33,6 +52,68 @@ def normalize_history_payload(payload: object) -> list[EntityHistorySeries]:
|
||||
return sorted(normalized, key=lambda item: item.entity_id)
|
||||
|
||||
|
||||
def normalize_state_history_payload(payload: object) -> list[StateHistorySeries]:
|
||||
if not isinstance(payload, list):
|
||||
raise HaUnexpectedPayloadError("History-Payload muss eine Liste sein.")
|
||||
normalized: list[StateHistorySeries] = []
|
||||
for raw_series in payload:
|
||||
if not isinstance(raw_series, list):
|
||||
raise HaUnexpectedPayloadError("History-Serie muss eine Liste sein.")
|
||||
entity_id: str | None = None
|
||||
points: list[StateHistoryPoint] = []
|
||||
for raw_entry in raw_series:
|
||||
if not isinstance(raw_entry, dict):
|
||||
raise HaUnexpectedPayloadError("History-Eintrag muss ein Objekt sein.")
|
||||
raw_entity_id = raw_entry.get("entity_id")
|
||||
if raw_entity_id is not None:
|
||||
if not isinstance(raw_entity_id, str) or "." not in raw_entity_id:
|
||||
raise HaUnexpectedPayloadError(
|
||||
"History-Eintrag enthält ungültige entity_id."
|
||||
)
|
||||
if entity_id is not None and entity_id != raw_entity_id:
|
||||
raise HaUnexpectedPayloadError("History-Serie enthält mehrere Entities.")
|
||||
entity_id = raw_entity_id
|
||||
raw_state = raw_entry.get("state")
|
||||
if not isinstance(raw_state, str) or raw_state in {"unknown", "unavailable"}:
|
||||
continue
|
||||
if entity_id is None:
|
||||
raise HaUnexpectedPayloadError("History-Serie enthält keine entity_id.")
|
||||
timestamp = _parse_timestamp(
|
||||
raw_entry.get("last_changed") or raw_entry.get("last_updated")
|
||||
)
|
||||
if not points or points[-1].state != raw_state:
|
||||
points.append(StateHistoryPoint(timestamp=timestamp, state=raw_state))
|
||||
if entity_id is not None and points:
|
||||
points.sort(key=lambda point: point.timestamp)
|
||||
normalized.append(StateHistorySeries(entity_id=entity_id, points=points))
|
||||
return sorted(normalized, key=lambda item: item.entity_id)
|
||||
|
||||
|
||||
def normalize_logbook_payload(payload: object, entity_id: str) -> list[LogbookEntry]:
|
||||
if not isinstance(payload, list):
|
||||
raise HaUnexpectedPayloadError("Logbook-Payload muss eine Liste sein.")
|
||||
entries: list[LogbookEntry] = []
|
||||
for raw_entry in payload:
|
||||
if not isinstance(raw_entry, dict):
|
||||
raise HaUnexpectedPayloadError("Logbook-Eintrag muss ein Objekt sein.")
|
||||
raw_entity_id = raw_entry.get("entity_id")
|
||||
if raw_entity_id != entity_id:
|
||||
continue
|
||||
entries.append(
|
||||
LogbookEntry(
|
||||
entity_id=entity_id,
|
||||
timestamp=_parse_timestamp(raw_entry.get("when")),
|
||||
message=str(raw_entry.get("message") or ""),
|
||||
context_user_id=_optional_string(raw_entry.get("context_user_id")),
|
||||
context_domain=_optional_string(
|
||||
raw_entry.get("context_domain") or raw_entry.get("domain")
|
||||
),
|
||||
context_service=_optional_string(raw_entry.get("context_service")),
|
||||
)
|
||||
)
|
||||
return sorted(entries, key=lambda item: item.timestamp)
|
||||
|
||||
|
||||
def _normalize_series(raw_series: list[object]) -> EntityHistorySeries | None:
|
||||
entity_id: str | None = None
|
||||
points: list[NumericHistoryPoint] = []
|
||||
@@ -89,3 +170,9 @@ def _parse_timestamp(value: object) -> datetime:
|
||||
if parsed.tzinfo is None:
|
||||
raise HaUnexpectedPayloadError("History-Zeitstempel muss eine Zeitzone enthalten.")
|
||||
return parsed
|
||||
|
||||
|
||||
def _optional_string(value: object) -> str | None:
|
||||
if value is None or value == "":
|
||||
return None
|
||||
return str(value)
|
||||
|
||||
@@ -14,6 +14,7 @@ class HaState(BaseModel):
|
||||
class HaEntitySummary(BaseModel):
|
||||
entity_id: str
|
||||
domain: str
|
||||
state: str | None = None
|
||||
state_class: str | None = None
|
||||
device_class: str | None = None
|
||||
unit_of_measurement: str | None = None
|
||||
|
||||
@@ -9,7 +9,14 @@ from app.ha.exceptions import HaClientError
|
||||
|
||||
from app.ha.client import HaClient
|
||||
from app.ha.discovery import DiscoveredEntity, discover_entities
|
||||
from app.ha.history import EntityHistorySeries, normalize_history_payload
|
||||
from app.ha.history import (
|
||||
EntityHistorySeries,
|
||||
LogbookEntry,
|
||||
StateHistorySeries,
|
||||
normalize_history_payload,
|
||||
normalize_logbook_payload,
|
||||
normalize_state_history_payload,
|
||||
)
|
||||
from app.ha.models import HaEntitySummary
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
@@ -45,6 +52,7 @@ class HaReader:
|
||||
HaEntitySummary(
|
||||
entity_id=entity_id,
|
||||
domain=domain,
|
||||
state=_optional_str(item.get("state")),
|
||||
state_class=_optional_str(attributes.get("state_class")),
|
||||
device_class=_optional_str(attributes.get("device_class")),
|
||||
unit_of_measurement=_optional_str(attributes.get("unit_of_measurement")),
|
||||
@@ -77,6 +85,32 @@ class HaReader:
|
||||
payload = self._client.get_history(entity_ids, start_time, end_time)
|
||||
return normalize_history_payload(payload)
|
||||
|
||||
def read_state_history(
|
||||
self,
|
||||
entity_ids: list[str],
|
||||
start_time: datetime,
|
||||
end_time: datetime,
|
||||
) -> Sequence[StateHistorySeries]:
|
||||
payload = self._client.get_history(entity_ids, start_time, end_time)
|
||||
return normalize_state_history_payload(payload)
|
||||
|
||||
def read_logbook(
|
||||
self,
|
||||
entity_id: str,
|
||||
start_time: datetime,
|
||||
end_time: datetime,
|
||||
) -> Sequence[LogbookEntry]:
|
||||
payload = self._client.get_logbook(entity_id, start_time, end_time)
|
||||
return normalize_logbook_payload(payload, entity_id)
|
||||
|
||||
def call_service(
|
||||
self,
|
||||
domain: str,
|
||||
service: str,
|
||||
service_data: dict[str, object],
|
||||
) -> Sequence[object]:
|
||||
return self._client.call_service(domain, service, service_data)
|
||||
|
||||
|
||||
def _optional_str(value: object) -> str | None:
|
||||
if value is None or value == "":
|
||||
|
||||
Reference in New Issue
Block a user