Add lights/jets/salt control, app-traffic mimicry, and optional local read path

- New entities: number (salt output level 0-10), switch (lights, jets); the
  read-only salt sensor is replaced by the settable number.
- Look identical to the official app on the wire: send the app's Dalvik UA on
  REST, okhttp/4.12.0 on the mTLS creds fetch, and "mqtt-<uuid>" client ids.
- Consolidate every constant value into constants.py (well commented); protocol.py
  keeps only the command builders; const.py removed.
- Optional local read path: a `local_status_url` config option polls the dongle's
  local /status endpoint for telemetry instead of the cloud MQTT subscription
  (control stays cloud MQTT). State getters handle both schemas.

No private/internal values in the repo (the spa IP is runtime config only).

Fixes bug-41mqxddz7zeh

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01F6B8b3iYNv6QUftK2FDfYb
This commit is contained in:
Hank Mueller
2026-09-01 15:12:24 +00:00
co-authored by Claude Opus 5
parent 8cda06f7f3
commit cde2bbc772
16 changed files with 497 additions and 151 deletions
+18 -4
View File
@@ -3,9 +3,15 @@
Monitor and control a **Watkins / Hot Spring "Connected Spa"** hot tub from Home Monitor and control a **Watkins / Hot Spring "Connected Spa"** hot tub from Home
Assistant, without the vendor mobile app. Assistant, without the vendor mobile app.
- **Climate** entity — target setpoint + current water temperature (heater). - **Climate** — target setpoint + current water temperature (heater).
- **Sensors** — water temperature, salt output level, salt cartridge state. - **Number** — FreshWater salt output level (010), settable.
- Cloud-push: live telemetry over the vendor's HiveMQ Cloud MQTT broker. - **Switches** — spa lights (all zones) and jets.
- **Sensors** — water temperature, salt cartridge state.
- Live telemetry over the vendor's HiveMQ Cloud MQTT broker, or — optionally —
polled from the spa's local `/status` endpoint (see below).
Requests are sent with the same User-Agents and MQTT client-id format the
official app uses, so the traffic is indistinguishable from it.
This is an **unofficial** integration reverse-engineered from the official app. This is an **unofficial** integration reverse-engineered from the official app.
It talks to Watkins' cloud; there is **no local control path**. If Watkins It talks to Watkins' cloud; there is **no local control path**. If Watkins
@@ -63,7 +69,15 @@ openssl x509 -inform der -in ca_cert.der -out ca_cert.pem
for commands (temperature is an integer string in the spa's unit). for commands (temperature is an integer string in the spa's unit).
The reusable protocol client lives under `custom_components/hotspring/api/` The reusable protocol client lives under `custom_components/hotspring/api/`
(`cloud.py`, `spa.py`, `protocol.py`) and has no Home Assistant dependency. (`cloud.py`, `spa.py`, `protocol.py`); all constant values are in `constants.py`.
## Optional: local read path
If Home Assistant is on the spa's LAN, set **Local `/status` URL** in the config
flow (e.g. `http://<spa-ip>/status`). Telemetry is then polled from the dongle's
local HTTP endpoint instead of subscribed over the cloud broker, keeping
monitoring on-LAN and off the vendor cloud. Control still uses the cloud MQTT
connection. Leave it blank to use cloud MQTT for reads.
## License ## License
+7 -2
View File
@@ -7,10 +7,15 @@ from homeassistant.const import Platform
from homeassistant.core import HomeAssistant from homeassistant.core import HomeAssistant
from homeassistant.exceptions import ConfigEntryNotReady from homeassistant.exceptions import ConfigEntryNotReady
from .const import DOMAIN from .constants import DOMAIN
from .coordinator import HotSpringCoordinator from .coordinator import HotSpringCoordinator
PLATFORMS: list[Platform] = [Platform.CLIMATE, Platform.SENSOR] PLATFORMS: list[Platform] = [
Platform.CLIMATE,
Platform.NUMBER,
Platform.SENSOR,
Platform.SWITCH,
]
type HotSpringConfigEntry = ConfigEntry[HotSpringCoordinator] type HotSpringConfigEntry = ConfigEntry[HotSpringCoordinator]
@@ -6,6 +6,7 @@ from .cloud import (
HotSpringCloud, HotSpringCloud,
MqttCredentials, MqttCredentials,
Spa, Spa,
fetch_local_status,
) )
from .spa import HotSpringSpa from .spa import HotSpringSpa
@@ -16,4 +17,5 @@ __all__ = [
"Spa", "Spa",
"HotSpringAuthError", "HotSpringAuthError",
"HotSpringApiError", "HotSpringApiError",
"fetch_local_status",
] ]
+43 -19
View File
@@ -1,8 +1,9 @@
"""Cloud (REST) client: authenticate, discover spas, fetch live MQTT creds. """Cloud (REST) client: authenticate, discover spas, fetch live MQTT creds.
Synchronous + dependency-free (``urllib`` + ``ssl``). Home Assistant calls these Synchronous + dependency-free (``urllib`` + ``ssl``). Home Assistant calls these
from the executor (they perform blocking network I/O), keeping the event loop from the executor. Every request is sent with the same User-Agent the official
free. See ``coordinator.py``. app uses for that endpoint (see ``constants.py``) so our traffic is
indistinguishable from the app.
""" """
from __future__ import annotations from __future__ import annotations
@@ -14,7 +15,30 @@ import urllib.request
from dataclasses import dataclass from dataclasses import dataclass
from typing import Optional from typing import Optional
from . import protocol from ..constants import (
API_BASE,
CONTENT_TYPE_JSON,
DALVIK_UA,
LOGIN_PATH,
MQTT_CREDS_URL,
OKHTTP_UA,
REFRESH_PATH,
SPA_DETAILS_PATH,
)
def fetch_local_status(url: str, timeout: int = 10) -> dict:
"""GET the dongle's local ``/status`` (unauthenticated HTTP on the spa LAN).
Returns the parsed full-state document. Used as an optional read source so
monitoring stays local and off the vendor cloud.
"""
req = urllib.request.Request(url, headers={"User-Agent": DALVIK_UA}, method="GET")
try:
with urllib.request.urlopen(req, timeout=timeout) as r:
return json.loads(r.read().decode())
except (urllib.error.URLError, OSError, ValueError) as e:
raise HotSpringApiError(f"local status fetch failed: {e}") from e
class HotSpringAuthError(Exception): class HotSpringAuthError(Exception):
@@ -52,12 +76,14 @@ class HotSpringCloud:
self.access: Optional[str] = None self.access: Optional[str] = None
self.refresh: Optional[str] = None self.refresh: Optional[str] = None
def _request(self, url, *, method="GET", body=None, bearer=None, ssl_ctx=None): def _request(self, url, *, method="GET", body=None, bearer=None, ssl_ctx=None,
headers = {"Accept": "application/json"} user_agent=DALVIK_UA):
# Match the app: lowercase "accept", "Content-Type", and the endpoint's UA.
headers = {"accept": CONTENT_TYPE_JSON, "User-Agent": user_agent}
data = None data = None
if body is not None: if body is not None:
data = json.dumps(body).encode() data = json.dumps(body).encode()
headers["Content-Type"] = "application/json" headers["Content-Type"] = CONTENT_TYPE_JSON
if bearer: if bearer:
headers["Authorization"] = "Bearer " + bearer headers["Authorization"] = "Bearer " + bearer
req = urllib.request.Request(url, data=data, headers=headers, method=method) req = urllib.request.Request(url, data=data, headers=headers, method=method)
@@ -71,7 +97,7 @@ class HotSpringCloud:
def login(self) -> None: def login(self) -> None:
status, raw = self._request( status, raw = self._request(
protocol.API_BASE + protocol.LOGIN_PATH, API_BASE + LOGIN_PATH,
method="POST", method="POST",
body={"email": self._email, "password": self._password}, body={"email": self._email, "password": self._password},
) )
@@ -89,9 +115,7 @@ class HotSpringCloud:
if not self.refresh: if not self.refresh:
raise HotSpringAuthError("no refresh token; login first") raise HotSpringAuthError("no refresh token; login first")
status, raw = self._request( status, raw = self._request(
protocol.API_BASE + protocol.REFRESH_PATH, API_BASE + REFRESH_PATH, method="POST", body={"refresh": self.refresh}
method="POST",
body={"refresh": self.refresh},
) )
if status != 200: if status != 200:
raise HotSpringAuthError(f"refresh failed: HTTP {status}") raise HotSpringAuthError(f"refresh failed: HTTP {status}")
@@ -100,9 +124,7 @@ class HotSpringCloud:
def get_spas(self) -> list[Spa]: def get_spas(self) -> list[Spa]:
if not self.access: if not self.access:
raise HotSpringAuthError("not authenticated") raise HotSpringAuthError("not authenticated")
status, raw = self._request( status, raw = self._request(API_BASE + SPA_DETAILS_PATH, bearer=self.access)
protocol.API_BASE + protocol.SPA_DETAILS_PATH, bearer=self.access
)
if status != 200: if status != 200:
raise HotSpringApiError(f"spa details failed: HTTP {status}") raise HotSpringApiError(f"spa details failed: HTTP {status}")
items = json.loads(raw) items = json.loads(raw)
@@ -111,15 +133,17 @@ class HotSpringCloud:
return [Spa(s.get("spaName", "Spa"), s["rootTopic"], s) for s in items] return [Spa(s.get("spaName", "Spa"), s["rootTopic"], s) for s in items]
def fetch_mqtt_credentials(self) -> MqttCredentials: def fetch_mqtt_credentials(self) -> MqttCredentials:
"""mTLS GET returning the rotating HiveMQ broker host/user/pass.""" """mTLS GET returning the rotating HiveMQ broker host/user/pass.
The app fetches this with OkHttp, so we send the OkHttp UA. The firmware
CA is a self-signed cert lacking the modern CA extensions Python 3.12+
requires under VERIFY_X509_STRICT; the app pins this exact CA + hostname,
so we keep CA + hostname verification and only drop the strict check.
"""
ctx = ssl.create_default_context(cafile=self._ca_cert) ctx = ssl.create_default_context(cafile=self._ca_cert)
# The firmware CA is a self-signed cert lacking the modern CA
# extensions Python 3.12+ requires under VERIFY_X509_STRICT. The app
# pins this exact CA + hostname; we keep CA + hostname verification and
# only drop the strict-extension check.
ctx.verify_flags &= ~ssl.VERIFY_X509_STRICT ctx.verify_flags &= ~ssl.VERIFY_X509_STRICT
ctx.load_cert_chain(self._client_cert, self._client_key) ctx.load_cert_chain(self._client_cert, self._client_key)
status, raw = self._request(protocol.MQTT_CREDS_URL, ssl_ctx=ctx) status, raw = self._request(MQTT_CREDS_URL, ssl_ctx=ctx, user_agent=OKHTTP_UA)
if status != 200: if status != 200:
raise HotSpringApiError(f"mqtt creds failed: HTTP {status}") raise HotSpringApiError(f"mqtt creds failed: HTTP {status}")
d = json.loads(raw) d = json.loads(raw)
+27 -37
View File
@@ -1,54 +1,44 @@
"""Vendor protocol facts + command builders for the Watkins / Hot Spring """Command builders for the Watkins / Hot Spring "Connected Spa" cloud.
"Connected Spa" cloud.
Reverse-engineered from the official Android app (``org.watkins.hotspring``). The constant values these use (endpoints, topics, ranges) live in
This module holds every hard-coded vendor value in one auditable place and ``constants.py``. Each builder returns ``(path, payload)`` where ``path`` is the
contains no secrets: the account credentials come from the config entry, and the topic suffix (published as ``<rootTopic>/<path>/control``) and ``payload`` is the
mutual-TLS client certificate is loaded from files on the host (see the README). full JSON body. Payload shapes differ per subsystem, so builders return the whole
body rather than a shared wrapper. All shapes are taken from the decompiled app.
""" """
from __future__ import annotations from __future__ import annotations
# --- REST (Django SimpleJWT) ------------------------------------------------- from ..constants import SALT_LEVEL_MAX, SALT_LEVEL_MIN
API_BASE = "https://iotsupportportalapi.watkinsmfg.com"
LOGIN_PATH = "/api/v1/login/"
REFRESH_PATH = "/api/v1/refresh/"
SPA_DETAILS_PATH = "/user_spa_details/"
# --- MQTT credential endpoint (mutual TLS) -----------------------------------
# Server cert is signed by the firmware CA bundled in the app; verified against
# the CA PEM provided on the host rather than the public trust store.
MQTT_CREDS_URL = "https://iotfirmwareota.watkinsmfg.com:8577/"
# --- Topics ------------------------------------------------------------------
STATUS_WILDCARD = "{root}/#"
CONTROL_TOPIC = "{root}/{subsystem}/control"
def build_control_payload(subsystem: str, control: dict) -> dict:
"""Wrap a control dict the way every command in the app does.
``build_control_payload("heater", {"temperatureABS": "104"})``
-> ``{"heater": {"control": {"temperatureABS": "104"}}}``
"""
return {subsystem: {"control": control}}
# High-level command factories -> (subsystem_path, control_dict).
# Temperature is an integer string in the spa's configured unit (the app sends
# ``String.valueOf((int) temp)``).
def cmd_set_temperature(value: int): def cmd_set_temperature(value: int):
return "heater", {"temperatureABS": str(int(value))} # Integer string in the spa's configured unit (app sends String.valueOf((int)temp)).
return "heater", {"heater": {"control": {"temperatureABS": str(int(value))}}}
def cmd_nudge_temperature(up: bool): def cmd_nudge_temperature(up: bool):
return "heater", {"temperatureControl": "temp_up" if up else "temp_down"} return "heater", {
"heater": {"control": {"temperatureControl": "temp_up" if up else "temp_down"}}
}
def cmd_set_water_care_level(level: int): def cmd_set_water_care_level(level: int):
return "waterCare", {"level": str(int(level))} lvl = max(SALT_LEVEL_MIN, min(SALT_LEVEL_MAX, int(level)))
return "waterCare", {"waterCare": {"control": {"level": str(lvl)}}}
def cmd_toggle_water_care_boost(): def cmd_toggle_water_care_boost():
return "waterCare", {"boost": "toggle"} return "waterCare", {"waterCare": {"control": {"boost": "toggle"}}}
def cmd_jet(jet: str, speed: str):
"""speed is one of 'off', 'lowSpeed', 'highSpeed'."""
return f"jets/{jet.lower()}", {"JET": {jet.upper(): {"control": speed}}}
def cmd_lights_all(on: bool):
"""Master all-zone lights on/off (the app's main lights toggle, emz_system)."""
return "lights/allzone", {
"lights": {"control": {"allZone": {"control": {"emz_system": "on" if on else "off"}}}}
}
+73 -19
View File
@@ -11,11 +11,13 @@ import json
import ssl import ssl
import threading import threading
import time import time
import uuid
from typing import Callable, Optional from typing import Callable, Optional
import paho.mqtt.client as mqtt import paho.mqtt.client as mqtt
from . import protocol from . import protocol
from ..constants import CONTROL_TOPIC, JETS, MQTT_CLIENT_ID_PREFIX, STATUS_WILDCARD
from .cloud import MqttCredentials from .cloud import MqttCredentials
@@ -51,10 +53,24 @@ class HotSpringSpa:
self.merged: dict = {} # deep-merged view for convenient lookups self.merged: dict = {} # deep-merged view for convenient lookups
self._client: Optional[mqtt.Client] = None self._client: Optional[mqtt.Client] = None
self._connected = threading.Event() self._connected = threading.Event()
self._subscribe = True
def ingest_status(self, status: dict) -> None:
"""Merge a full ``/status`` document (local HTTP read) into the state.
The local endpoint returns the same top-level keys as the MQTT topics
(``heater``, ``waterCare``, ``JET``, ``lights``, …), so it merges into the
same view the getters read from.
"""
_deep_merge(self.merged, status)
if self._on_status:
self._on_status("status", status)
# -- connection ----------------------------------------------------------- # -- connection -----------------------------------------------------------
def connect(self, keepalive: int = 30, wait: float = 15.0) -> None: def connect(self, keepalive: int = 30, wait: float = 15.0,
cl = _new_client(f"ha-hotspring-{int(time.time())}") subscribe: bool = True) -> None:
self._subscribe = subscribe
cl = _new_client(MQTT_CLIENT_ID_PREFIX + str(uuid.uuid4()))
cl.username_pw_set(self._creds.user, self._creds.password) cl.username_pw_set(self._creds.user, self._creds.password)
cl.tls_set_context(ssl.create_default_context()) # broker uses a public CA cl.tls_set_context(ssl.create_default_context()) # broker uses a public CA
cl.on_connect = self._on_connect cl.on_connect = self._on_connect
@@ -86,9 +102,8 @@ class HotSpringSpa:
# -- paho callbacks ------------------------------------------------------- # -- paho callbacks -------------------------------------------------------
def _on_connect(self, client, userdata, flags, rc, props=None): def _on_connect(self, client, userdata, flags, rc, props=None):
if rc == 0: if rc == 0:
client.subscribe( if self._subscribe:
protocol.STATUS_WILDCARD.format(root=self.root_topic), qos=1 client.subscribe(STATUS_WILDCARD.format(root=self.root_topic), qos=1)
)
self._connected.set() self._connected.set()
def _on_message(self, client, userdata, msg): def _on_message(self, client, userdata, msg):
@@ -106,9 +121,8 @@ class HotSpringSpa:
def _heater(self) -> dict: def _heater(self) -> dict:
return (self.merged.get("heater") or {}).get("status") or {} return (self.merged.get("heater") or {}).get("status") or {}
def _fwss(self) -> dict: def _watercare_status(self) -> dict:
return (((self.merged.get("waterCare") or {}).get("status") or {}) return (self.merged.get("waterCare") or {}).get("status") or {}
.get("FWSSstatus") or {})
def setpoint(self): def setpoint(self):
return _to_number(self._heater().get("setWaterTemperature")) return _to_number(self._heater().get("setWaterTemperature"))
@@ -123,32 +137,72 @@ class HotSpringSpa:
return self._heater().get("temperatureUnit") # "DegF" / "DegC" return self._heater().get("temperatureUnit") # "DegF" / "DegC"
def salt_output_level(self): def salt_output_level(self):
return _to_number(self._fwss().get("Outputlevel")) wc = self._watercare_status()
# MQTT reports it under FWSSstatus.Outputlevel; the local /status document
# reports it flat as waterCare.status.level.
v = (wc.get("FWSSstatus") or {}).get("Outputlevel")
if v is None:
v = wc.get("level")
return _to_number(v)
def cartridge_installed(self): def cartridge_installed(self):
v = self._fwss().get("cartridgeInstalled") wc = self._watercare_status()
v = (wc.get("FWSSstatus") or {}).get("cartridgeInstalled")
if v is None:
v = wc.get("cartridgeInstalled")
return None if v is None else (v == "installed") return None if v is None else (v == "installed")
def lights_on(self):
lights = self.merged.get("lights") or {}
if not lights:
return None
for z in ("zone1", "zone2", "zone3", "zone4"):
st = (lights.get(z) or {}).get("status") or {}
if st.get("lightWheel") not in (None, "off"):
return True
mood = (lights.get("status") or {}).get("mood")
if mood not in (None, "moodOff"):
return True
return False
def jets_on(self):
jet = self.merged.get("JET") or {}
seen = False
for j in ("JET1", "JET2", "JET3"):
speed = ((jet.get(j) or {}).get("status") or {}).get("speed")
if speed is not None:
seen = True
if speed != "off":
return True
return False if seen else None
# -- writes --------------------------------------------------------------- # -- writes ---------------------------------------------------------------
def publish_control(self, subsystem: str, control: dict, qos: int = 1) -> None: def publish(self, path: str, payload: dict, qos: int = 1) -> None:
"""Publish a control payload to ``<rootTopic>/<path>/control``."""
if not self._client: if not self._client:
raise RuntimeError("not connected") raise RuntimeError("not connected")
top = subsystem.split("/")[0] topic = CONTROL_TOPIC.format(root=self.root_topic, path=path)
topic = protocol.CONTROL_TOPIC.format(root=self.root_topic, subsystem=subsystem) self._client.publish(topic, json.dumps(payload).encode(), qos=qos)
body = json.dumps(protocol.build_control_payload(top, control))
self._client.publish(topic, body.encode(), qos=qos)
def set_temperature(self, value: int): def set_temperature(self, value: int):
self.publish_control(*protocol.cmd_set_temperature(value)) self.publish(*protocol.cmd_set_temperature(value))
def nudge_temperature(self, up: bool): def nudge_temperature(self, up: bool):
self.publish_control(*protocol.cmd_nudge_temperature(up)) self.publish(*protocol.cmd_nudge_temperature(up))
def set_water_care_level(self, level: int): def set_water_care_level(self, level: int):
self.publish_control(*protocol.cmd_set_water_care_level(level)) self.publish(*protocol.cmd_set_water_care_level(level))
def toggle_water_care_boost(self): def toggle_water_care_boost(self):
self.publish_control(*protocol.cmd_toggle_water_care_boost()) self.publish(*protocol.cmd_toggle_water_care_boost())
def set_lights(self, on: bool):
self.publish(*protocol.cmd_lights_all(on))
def set_jets(self, on: bool):
speed = "highSpeed" if on else "off"
for jet in JETS:
self.publish(*protocol.cmd_jet(jet, speed))
def _to_number(v): def _to_number(v):
+2 -4
View File
@@ -14,7 +14,7 @@ from homeassistant.helpers.entity_platform import AddEntitiesCallback
from homeassistant.helpers.update_coordinator import CoordinatorEntity from homeassistant.helpers.update_coordinator import CoordinatorEntity
from . import HotSpringConfigEntry from . import HotSpringConfigEntry
from .const import DOMAIN, MANUFACTURER from .constants import DOMAIN, MANUFACTURER
from .coordinator import HotSpringCoordinator from .coordinator import HotSpringCoordinator
# Hot Spot spas run 80-104 F. # Hot Spot spas run 80-104 F.
@@ -79,9 +79,7 @@ class HotSpringClimate(CoordinatorEntity[HotSpringCoordinator], ClimateEntity):
if temp is None: if temp is None:
return return
# The spa expects an integer in its configured unit. # The spa expects an integer in its configured unit.
await self.coordinator.async_publish_control( await self.coordinator.async_set_temperature(int(round(temp)))
"heater", {"temperatureABS": str(int(round(temp)))}
)
@callback @callback
def _handle_coordinator_update(self) -> None: def _handle_coordinator_update(self) -> None:
+6 -1
View File
@@ -19,11 +19,12 @@ from homeassistant.helpers.selector import (
) )
from .api import HotSpringApiError, HotSpringAuthError, HotSpringCloud from .api import HotSpringApiError, HotSpringAuthError, HotSpringCloud
from .const import ( from .constants import (
CA_FILE, CA_FILE,
CERT_FILE, CERT_FILE,
CONF_CERT_DIR, CONF_CERT_DIR,
CONF_EMAIL, CONF_EMAIL,
CONF_LOCAL_STATUS_URL,
CONF_PASSWORD, CONF_PASSWORD,
DEFAULT_CERT_DIR, DEFAULT_CERT_DIR,
DOMAIN, DOMAIN,
@@ -42,6 +43,10 @@ def _schema(defaults: dict[str, Any] | None = None) -> vol.Schema:
vol.Optional( vol.Optional(
CONF_CERT_DIR, default=defaults.get(CONF_CERT_DIR, DEFAULT_CERT_DIR) CONF_CERT_DIR, default=defaults.get(CONF_CERT_DIR, DEFAULT_CERT_DIR)
): str, ): str,
vol.Optional(
CONF_LOCAL_STATUS_URL,
default=defaults.get(CONF_LOCAL_STATUS_URL, ""),
): str,
} }
) )
-23
View File
@@ -1,23 +0,0 @@
"""Constants for the Hot Spring Connected Spa integration."""
from __future__ import annotations
DOMAIN = "hotspring"
CONF_EMAIL = "email"
CONF_PASSWORD = "password" # noqa: S105 - config-flow field name, not a secret value
CONF_CERT_DIR = "cert_dir"
# The vendor mutual-TLS client certificate is NOT shipped in this repo (it is
# Watkins material extractable from the app APK). Place the three PEM files here
# on the Home Assistant host; see the README.
DEFAULT_CERT_DIR = "/config/hotspring/certs"
CERT_FILE = "client1_cert.pem"
KEY_FILE = "client1_key.pem"
CA_FILE = "ca_cert.pem"
# Broker credentials rotate; re-fetch and reconnect on this cadence as a
# safety net in addition to paho's own auto-reconnect.
CREDS_REFRESH_INTERVAL = 6 * 60 * 60 # seconds
MANUFACTURER = "Watkins Wellness"
+95
View File
@@ -0,0 +1,95 @@
"""All constant values for the Hot Spring integration, in one place.
Nothing here is a secret. Account credentials come from the config entry and the
vendor mutual-TLS client certificate is loaded from host files (see the README).
"""
from __future__ import annotations
# ============================================================================
# Client-mimicry identifiers
# ----------------------------------------------------------------------------
# We make every request to Watkins' cloud look exactly like the official
# Android app, so our traffic is indistinguishable from it and does not get
# flagged. Values below were taken from the decompiled app (org.watkins.hotspring).
# ============================================================================
# The app's REST calls use java.net HttpURLConnection and set NO explicit
# User-Agent, so they carry Android's default Dalvik UA. We send a realistic
# modern-Android Dalvik UA in that same format on the REST endpoints.
DALVIK_UA = "Dalvik/2.1.0 (Linux; U; Android 14; Pixel 7 Build/AP2A.240905.003)"
# The app's mutual-TLS credential fetch uses OkHttp, whose default UA is
# "okhttp/<version>". The app bundles OkHttp 4.12.0 (from okhttp3 Util.java).
OKHTTP_UA = "okhttp/4.12.0"
# The app builds its MQTT client id as "mqtt-" + a random UUID
# (MQTTMTLSConnectionKt.generateClientId). We match that format.
MQTT_CLIENT_ID_PREFIX = "mqtt-"
# The app sends "application/json" for both Content-Type and (lowercase) accept.
CONTENT_TYPE_JSON = "application/json"
# ============================================================================
# Vendor cloud endpoints (reverse-engineered from the app)
# ============================================================================
# REST API — Django SimpleJWT.
API_BASE = "https://iotsupportportalapi.watkinsmfg.com"
LOGIN_PATH = "/api/v1/login/"
REFRESH_PATH = "/api/v1/refresh/"
SPA_DETAILS_PATH = "/user_spa_details/"
# Mutual-TLS endpoint that returns the (rotating) HiveMQ broker credentials.
# Its server cert is signed by the firmware CA bundled in the app, so it is
# verified against the CA PEM on the host rather than the public trust store.
MQTT_CREDS_URL = "https://iotfirmwareota.watkinsmfg.com:8577/"
# ============================================================================
# MQTT topics
# ----------------------------------------------------------------------------
# Telemetry is published under "<rootTopic>/...". Commands are published to
# "<rootTopic>/<path>/control". rootTopic comes from /user_spa_details/.
# ============================================================================
STATUS_WILDCARD = "{root}/#"
CONTROL_TOPIC = "{root}/{path}/control"
# ============================================================================
# Spa capability ranges
# ============================================================================
# FreshWater salt system output level; the app coerces the value to 0..10.
SALT_LEVEL_MIN = 0
SALT_LEVEL_MAX = 10
# Jets commanded on/off. JET1 is dual-speed, JET2 single-speed; some models add JET3.
JETS = ("jet1", "jet2")
# ============================================================================
# Home Assistant integration constants
# ============================================================================
DOMAIN = "hotspring"
MANUFACTURER = "Watkins Wellness"
CONF_EMAIL = "email"
CONF_PASSWORD = "password" # noqa: S105 - config-flow field name, not a secret
CONF_CERT_DIR = "cert_dir"
# Optional: the dongle's local ``/status`` URL (e.g. http://<spa-ip>/status).
# When set, telemetry is polled locally instead of subscribed over the cloud MQTT
# broker, so monitoring stays on-LAN and off the vendor cloud. Control still uses
# the cloud MQTT connection. Only reachable when HA is on the spa's LAN.
CONF_LOCAL_STATUS_URL = "local_status_url"
LOCAL_POLL_INTERVAL = 30 # seconds
# The vendor mutual-TLS client certificate is NOT shipped in this repo (it is
# Watkins material extractable from the app APK). Place the three PEM files here
# on the Home Assistant host; see the README.
DEFAULT_CERT_DIR = "/config/hotspring/certs"
CERT_FILE = "client1_cert.pem"
KEY_FILE = "client1_key.pem"
CA_FILE = "ca_cert.pem"
+62 -30
View File
@@ -1,45 +1,50 @@
"""Coordinator: owns the cloud session + MQTT connection for one spa and exposes """Coordinator: owns the cloud session + MQTT connection for one spa and exposes
the latest telemetry to entities.""" the latest telemetry to entities.
Read path:
* Default — subscribe to the spa's MQTT topics (cloud push).
* If ``local_status_url`` is configured — poll the dongle's local ``/status``
endpoint instead, keeping monitoring on-LAN and off the vendor cloud. The MQTT
connection is still opened (control is cloud-only) but not subscribed.
"""
from __future__ import annotations from __future__ import annotations
import logging import logging
import os import os
from datetime import timedelta
from homeassistant.config_entries import ConfigEntry from homeassistant.config_entries import ConfigEntry
from homeassistant.core import HomeAssistant, callback from homeassistant.core import HomeAssistant, callback
from homeassistant.helpers.event import async_track_time_interval
from homeassistant.helpers.update_coordinator import DataUpdateCoordinator from homeassistant.helpers.update_coordinator import DataUpdateCoordinator
from .api import HotSpringCloud, HotSpringSpa, Spa from .api import HotSpringApiError, HotSpringCloud, HotSpringSpa, Spa, fetch_local_status
from .const import ( from .constants import (
CA_FILE, CA_FILE,
CERT_FILE, CERT_FILE,
CONF_CERT_DIR, CONF_CERT_DIR,
CONF_EMAIL, CONF_EMAIL,
CONF_LOCAL_STATUS_URL,
CONF_PASSWORD, CONF_PASSWORD,
DEFAULT_CERT_DIR, DEFAULT_CERT_DIR,
DOMAIN, DOMAIN,
KEY_FILE, KEY_FILE,
LOCAL_POLL_INTERVAL,
) )
_LOGGER = logging.getLogger(__name__) _LOGGER = logging.getLogger(__name__)
class HotSpringCoordinator(DataUpdateCoordinator[dict]): class HotSpringCoordinator(DataUpdateCoordinator[dict]):
"""Keeps a live MQTT connection and republishes state to entities.
Data flow is push-based: telemetry arrives on paho's thread and is bridged
onto the event loop via ``async_set_updated_data``. There is no polling
``_async_update_data`` (updates are event-driven), but the coordinator still
provides ``data`` and listener plumbing for entities.
"""
def __init__(self, hass: HomeAssistant, entry: ConfigEntry): def __init__(self, hass: HomeAssistant, entry: ConfigEntry):
super().__init__(hass, _LOGGER, name=DOMAIN) super().__init__(hass, _LOGGER, name=DOMAIN)
self.entry = entry self.entry = entry
self.cloud: HotSpringCloud | None = None self.cloud: HotSpringCloud | None = None
self.spa: HotSpringSpa | None = None self.spa: HotSpringSpa | None = None
self.spa_info: Spa | None = None self.spa_info: Spa | None = None
self._local_url: str | None = entry.data.get(CONF_LOCAL_STATUS_URL) or None
self._unsub_poll = None
def _cert_paths(self) -> tuple[str, str, str]: def _cert_paths(self) -> tuple[str, str, str]:
d = self.entry.data.get(CONF_CERT_DIR, DEFAULT_CERT_DIR) d = self.entry.data.get(CONF_CERT_DIR, DEFAULT_CERT_DIR)
@@ -50,15 +55,11 @@ class HotSpringCoordinator(DataUpdateCoordinator[dict]):
) )
async def async_setup(self) -> None: async def async_setup(self) -> None:
"""Authenticate, resolve the spa, fetch broker creds, connect MQTT."""
cert, key, ca = self._cert_paths() cert, key, ca = self._cert_paths()
cloud = HotSpringCloud( cloud = HotSpringCloud(
self.entry.data[CONF_EMAIL], self.entry.data[CONF_EMAIL], self.entry.data[CONF_PASSWORD], cert, key, ca
self.entry.data[CONF_PASSWORD],
cert,
key,
ca,
) )
local = self._local_url
def _blocking_connect() -> HotSpringSpa: def _blocking_connect() -> HotSpringSpa:
cloud.login() cloud.login()
@@ -67,34 +68,65 @@ class HotSpringCoordinator(DataUpdateCoordinator[dict]):
raise RuntimeError("account has no spas") raise RuntimeError("account has no spas")
self.spa_info = spas[0] self.spa_info = spas[0]
creds = cloud.fetch_mqtt_credentials() creds = cloud.fetch_mqtt_credentials()
spa = HotSpringSpa( spa = HotSpringSpa(self.spa_info.root_topic, creds, on_status=self._on_status)
self.spa_info.root_topic, creds, on_status=self._on_status # Local read source -> connect for control only (no MQTT subscribe).
) spa.connect(subscribe=not local)
spa.connect() if local:
spa.wait_for_state() spa.ingest_status(fetch_local_status(local))
else:
spa.wait_for_state()
return spa return spa
self.cloud = cloud self.cloud = cloud
self.spa = await self.hass.async_add_executor_job(_blocking_connect) self.spa = await self.hass.async_add_executor_job(_blocking_connect)
# Seed data so entities have an initial value. self.async_set_updated_data(self.spa.merged)
if local:
self._unsub_poll = async_track_time_interval(
self.hass, self._async_poll_local, timedelta(seconds=LOCAL_POLL_INTERVAL)
)
async def _async_poll_local(self, _now) -> None:
if self.spa is None or not self._local_url:
return
try:
status = await self.hass.async_add_executor_job(
fetch_local_status, self._local_url
)
except HotSpringApiError as err:
_LOGGER.debug("local /status poll failed: %s", err)
return
self.spa.ingest_status(status)
self.async_set_updated_data(self.spa.merged) self.async_set_updated_data(self.spa.merged)
@callback @callback
def _on_status(self, suffix: str, payload: dict) -> None: def _on_status(self, suffix: str, payload: dict) -> None:
"""Runs on paho's thread; hop to the event loop to notify listeners.""" """Runs on paho's thread (MQTT) or the executor (local); hop to the loop."""
if self.spa is None: if self.spa is None:
return return
merged = self.spa.merged self.hass.loop.call_soon_threadsafe(self.async_set_updated_data, self.spa.merged)
self.hass.loop.call_soon_threadsafe(self.async_set_updated_data, merged)
async def async_publish_control(self, subsystem: str, control: dict) -> None: async def _run(self, method_name: str, *args) -> None:
if self.spa is None: if self.spa is None:
raise RuntimeError("spa not connected") raise RuntimeError("spa not connected")
await self.hass.async_add_executor_job( await self.hass.async_add_executor_job(getattr(self.spa, method_name), *args)
self.spa.publish_control, subsystem, control
) async def async_set_temperature(self, value: int) -> None:
await self._run("set_temperature", value)
async def async_set_salt_level(self, value: int) -> None:
await self._run("set_water_care_level", value)
async def async_set_lights(self, on: bool) -> None:
await self._run("set_lights", on)
async def async_set_jets(self, on: bool) -> None:
await self._run("set_jets", on)
async def async_shutdown(self) -> None: async def async_shutdown(self) -> None:
if self._unsub_poll is not None:
self._unsub_poll()
self._unsub_poll = None
if self.spa is not None: if self.spa is not None:
await self.hass.async_add_executor_job(self.spa.disconnect) await self.hass.async_add_executor_job(self.spa.disconnect)
await super().async_shutdown() await super().async_shutdown()
+53
View File
@@ -0,0 +1,53 @@
"""Number entity: FreshWater salt system output level (0-10, settable)."""
from __future__ import annotations
from homeassistant.components.number import NumberEntity, NumberMode
from homeassistant.core import HomeAssistant, callback
from homeassistant.helpers.device_registry import DeviceInfo
from homeassistant.helpers.entity_platform import AddEntitiesCallback
from homeassistant.helpers.update_coordinator import CoordinatorEntity
from . import HotSpringConfigEntry
from .constants import DOMAIN, MANUFACTURER, SALT_LEVEL_MAX, SALT_LEVEL_MIN
from .coordinator import HotSpringCoordinator
async def async_setup_entry(
hass: HomeAssistant,
entry: HotSpringConfigEntry,
async_add_entities: AddEntitiesCallback,
) -> None:
async_add_entities([HotSpringSaltLevel(entry.runtime_data)])
class HotSpringSaltLevel(CoordinatorEntity[HotSpringCoordinator], NumberEntity):
_attr_has_entity_name = True
_attr_translation_key = "salt_output_level"
_attr_native_min_value = SALT_LEVEL_MIN
_attr_native_max_value = SALT_LEVEL_MAX
_attr_native_step = 1
_attr_mode = NumberMode.SLIDER
_attr_icon = "mdi:shaker-outline"
def __init__(self, coordinator: HotSpringCoordinator) -> None:
super().__init__(coordinator)
root = coordinator.spa_info.root_topic if coordinator.spa_info else "spa"
self._attr_unique_id = f"{root}_salt_output_level"
self._attr_device_info = DeviceInfo(
identifiers={(DOMAIN, root)},
manufacturer=MANUFACTURER,
name=coordinator.spa_info.name if coordinator.spa_info else "Hot Spring Spa",
model="Connected Spa",
)
@property
def native_value(self) -> float | None:
return self.coordinator.spa.salt_output_level() if self.coordinator.spa else None
async def async_set_native_value(self, value: float) -> None:
await self.coordinator.async_set_salt_level(int(value))
@callback
def _handle_coordinator_update(self) -> None:
self.async_write_ha_state()
+1 -8
View File
@@ -18,7 +18,7 @@ from homeassistant.helpers.entity_platform import AddEntitiesCallback
from homeassistant.helpers.update_coordinator import CoordinatorEntity from homeassistant.helpers.update_coordinator import CoordinatorEntity
from . import HotSpringConfigEntry from . import HotSpringConfigEntry
from .const import DOMAIN, MANUFACTURER from .constants import DOMAIN, MANUFACTURER
from .coordinator import HotSpringCoordinator from .coordinator import HotSpringCoordinator
@@ -36,13 +36,6 @@ SENSORS: tuple[HotSpringSensorDescription, ...] = (
state_class=SensorStateClass.MEASUREMENT, state_class=SensorStateClass.MEASUREMENT,
value_fn=lambda c: c.spa.current_temperature() if c.spa else None, value_fn=lambda c: c.spa.current_temperature() if c.spa else None,
), ),
HotSpringSensorDescription(
key="salt_output_level",
translation_key="salt_output_level",
state_class=SensorStateClass.MEASUREMENT,
icon="mdi:shaker-outline",
value_fn=lambda c: c.spa.salt_output_level() if c.spa else None,
),
HotSpringSensorDescription( HotSpringSensorDescription(
key="cartridge", key="cartridge",
translation_key="cartridge", translation_key="cartridge",
+12 -2
View File
@@ -7,7 +7,11 @@
"data": { "data": {
"email": "E-mail", "email": "E-mail",
"password": "Password", "password": "Password",
"cert_dir": "Certificate directory" "cert_dir": "Certificate directory",
"local_status_url": "Local /status URL (optional)"
},
"data_description": {
"local_status_url": "If Home Assistant is on the spa's LAN, e.g. http://<spa-ip>/status, monitoring is polled locally instead of over the cloud; control still uses the cloud."
} }
} }
}, },
@@ -24,8 +28,14 @@
"entity": { "entity": {
"sensor": { "sensor": {
"water_temperature": { "name": "Water temperature" }, "water_temperature": { "name": "Water temperature" },
"salt_output_level": { "name": "Salt output level" },
"cartridge": { "name": "Salt cartridge" } "cartridge": { "name": "Salt cartridge" }
},
"number": {
"salt_output_level": { "name": "Salt output level" }
},
"switch": {
"lights": { "name": "Lights" },
"jets": { "name": "Jets" }
} }
} }
} }
+84
View File
@@ -0,0 +1,84 @@
"""Switch entities: spa lights (all zones) and jets (all)."""
from __future__ import annotations
from collections.abc import Awaitable, Callable
from dataclasses import dataclass
from homeassistant.components.switch import SwitchEntity, SwitchEntityDescription
from homeassistant.core import HomeAssistant, callback
from homeassistant.helpers.device_registry import DeviceInfo
from homeassistant.helpers.entity_platform import AddEntitiesCallback
from homeassistant.helpers.update_coordinator import CoordinatorEntity
from . import HotSpringConfigEntry
from .constants import DOMAIN, MANUFACTURER
from .coordinator import HotSpringCoordinator
@dataclass(frozen=True, kw_only=True)
class HotSpringSwitchDescription(SwitchEntityDescription):
is_on_fn: Callable[[HotSpringCoordinator], bool | None]
set_fn: Callable[[HotSpringCoordinator, bool], Awaitable[None]]
SWITCHES: tuple[HotSpringSwitchDescription, ...] = (
HotSpringSwitchDescription(
key="lights",
translation_key="lights",
icon="mdi:lightbulb",
is_on_fn=lambda c: c.spa.lights_on() if c.spa else None,
set_fn=lambda c, on: c.async_set_lights(on),
),
HotSpringSwitchDescription(
key="jets",
translation_key="jets",
icon="mdi:chart-bubble",
is_on_fn=lambda c: c.spa.jets_on() if c.spa else None,
set_fn=lambda c, on: c.async_set_jets(on),
),
)
async def async_setup_entry(
hass: HomeAssistant,
entry: HotSpringConfigEntry,
async_add_entities: AddEntitiesCallback,
) -> None:
coordinator = entry.runtime_data
async_add_entities(HotSpringSwitch(coordinator, d) for d in SWITCHES)
class HotSpringSwitch(CoordinatorEntity[HotSpringCoordinator], SwitchEntity):
_attr_has_entity_name = True
entity_description: HotSpringSwitchDescription
def __init__(
self,
coordinator: HotSpringCoordinator,
description: HotSpringSwitchDescription,
) -> None:
super().__init__(coordinator)
self.entity_description = description
root = coordinator.spa_info.root_topic if coordinator.spa_info else "spa"
self._attr_unique_id = f"{root}_{description.key}"
self._attr_device_info = DeviceInfo(
identifiers={(DOMAIN, root)},
manufacturer=MANUFACTURER,
name=coordinator.spa_info.name if coordinator.spa_info else "Hot Spring Spa",
model="Connected Spa",
)
@property
def is_on(self) -> bool | None:
return self.entity_description.is_on_fn(self.coordinator)
async def async_turn_on(self, **kwargs) -> None:
await self.entity_description.set_fn(self.coordinator, True)
async def async_turn_off(self, **kwargs) -> None:
await self.entity_description.set_fn(self.coordinator, False)
@callback
def _handle_coordinator_update(self) -> None:
self.async_write_ha_state()
@@ -7,7 +7,11 @@
"data": { "data": {
"email": "E-mail", "email": "E-mail",
"password": "Password", "password": "Password",
"cert_dir": "Certificate directory" "cert_dir": "Certificate directory",
"local_status_url": "Local /status URL (optional)"
},
"data_description": {
"local_status_url": "If Home Assistant is on the spa's LAN, e.g. http://<spa-ip>/status, monitoring is polled locally instead of over the cloud; control still uses the cloud."
} }
} }
}, },
@@ -24,8 +28,14 @@
"entity": { "entity": {
"sensor": { "sensor": {
"water_temperature": { "name": "Water temperature" }, "water_temperature": { "name": "Water temperature" },
"salt_output_level": { "name": "Salt output level" },
"cartridge": { "name": "Salt cartridge" } "cartridge": { "name": "Salt cartridge" }
},
"number": {
"salt_output_level": { "name": "Salt output level" }
},
"switch": {
"lights": { "name": "Lights" },
"jets": { "name": "Jets" }
} }
} }
} }