Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions roborock/devices/device.py
Original file line number Diff line number Diff line change
Expand Up @@ -202,6 +202,8 @@ async def connect(self) -> None:
await self.v1_properties.start()
elif self.b01_q10_properties is not None:
await self.b01_q10_properties.start()
if self.zeo is not None:
await self.zeo.start()
except RoborockException:
# Expected: start() can fail transiently. Unsubscribe before propagating
# so the retry by connect_loop() gets a clean channel.
Expand Down
85 changes: 83 additions & 2 deletions roborock/devices/traits/a01/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@
"""

import json
import logging
from collections.abc import Callable
from datetime import time
from typing import Any
Expand All @@ -40,6 +41,7 @@
ZeoDetergentType,
ZeoDryingMode,
ZeoError,
ZeoFeatureBits,
ZeoMode,
ZeoProgram,
ZeoRinse,
Expand All @@ -50,8 +52,18 @@
)
from roborock.devices.rpc.a01_channel import send_decoded_command
from roborock.devices.traits import Trait
from roborock.devices.traits.common import TraitUpdateListener
from roborock.devices.transport.mqtt_channel import MqttChannel
from roborock.roborock_message import RoborockDyadDataProtocol, RoborockZeoProtocol
from roborock.exceptions import RoborockException
from roborock.protocols.a01_protocol import decode_rpc_response
from roborock.roborock_message import (
RoborockDyadDataProtocol,
RoborockMessage,
RoborockMessageProtocol,
RoborockZeoProtocol,
)

_LOGGER = logging.getLogger(__name__)

__init__ = [
"DyadApi",
Expand Down Expand Up @@ -156,14 +168,80 @@ async def set_value(self, protocol: RoborockDyadDataProtocol, value: Any) -> dic
return await send_decoded_command(self._channel, params)


class ZeoApi(Trait):
class ZeoApi(Trait, TraitUpdateListener):
"""API for interacting with Zeo devices."""

name = "zeo"

def __init__(self, channel: MqttChannel) -> None:
"""Initialize the Zeo API."""
TraitUpdateListener.__init__(self, _LOGGER)
self._channel = channel
self._dps_cache: dict[int, Any] = {}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This appears unused in this PR. Can we explain how we expect this to be used here and what the semantics are? when it is ok to use vs when do we need to refresh, etc.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

See #897
ZeoCommandTrait._get_start_params() checks the cache for MODE/PROGRAM and only issues a device query on cache miss — avoiding redundant network round-trips when values were already received via push.

self._dps_unsub: Callable[[], None] | None = None
self._feature_bits: int = 0

async def start(self) -> None:
"""Subscribe to MQTT push and discover device features.

Subscribes to the DPS MQTT topic, then queries FEATURE_BITS
(DP 237) to wake the device and cache supported capabilities.
"""
await self._ensure_subscribed()
await self._discover_features()

def close(self) -> None:
"""Unsubscribe from MQTT push and release resources."""
if self._dps_unsub is not None:
self._dps_unsub()
self._dps_unsub = None

async def _ensure_subscribed(self) -> None:
"""Subscribe to MQTT DPS push (idempotent)."""
if self._dps_unsub is not None:
return
self._dps_unsub = await self._channel.subscribe(self._on_dps_message)

async def _discover_features(self) -> None:
"""Query FEATURE_BITS to wake the device and cache capabilities.

Sending an RPC query after subscribing triggers the device to
start pushing its full state — equivalent to how V1's
``discover_features()`` uses ``device_features.refresh()`` to
initiate the push cycle.

Only devices that support the FeatureBits DP will respond;
older or unsupported devices return nothing.
A failed query defaults to 0 — all feature-gated DPs are
disabled and the device operates in basic mode.
"""
try:
result = await self.query_values([RoborockZeoProtocol.FEATURE_BITS])
self._feature_bits = result.get(RoborockZeoProtocol.FEATURE_BITS, 0)
except RoborockException:
self._feature_bits = 0

def supports(self, feature: ZeoFeatureBits) -> bool:
"""Check whether the device supports a given feature bit."""
return bool(self._feature_bits & (1 << feature.value))

def _on_dps_message(self, message: RoborockMessage) -> None:
"""Handle unsolicited MQTT push (protocol 102 — RPC_RESPONSE).

Zeo devices broadcast status changes as ``{"dps": {...}}`` JSON
payloads. This callback decodes them and feeds the cache so
that ``query_values`` can skip the device round-trip when the
requested DPs are already up to date.
"""
if message.protocol != RoborockMessageProtocol.RPC_RESPONSE:
return
try:
decoded = decode_rpc_response(message)
except RoborockException:
_LOGGER.debug("Dropped malformed push message", exc_info=True)
return
self._dps_cache.update(decoded)
self._notify_update()

async def query_values(self, protocols: list[RoborockZeoProtocol]) -> dict[RoborockZeoProtocol, Any]:
"""Query the device for the values of the given protocols."""
Expand All @@ -172,6 +250,9 @@ async def query_values(self, protocols: list[RoborockZeoProtocol]) -> dict[Robor
{RoborockZeoProtocol.ID_QUERY: protocols},
value_encoder=json.dumps,
)
for protocol, value in response.items():
if value is not None:
self._dps_cache[int(protocol)] = value
return {protocol: convert_zeo_value(protocol, response.get(protocol)) for protocol in protocols}

async def set_value(self, protocol: RoborockZeoProtocol, value: Any) -> dict[RoborockZeoProtocol, Any]:
Expand Down
3 changes: 3 additions & 0 deletions roborock/roborock_message.py
Original file line number Diff line number Diff line change
Expand Up @@ -179,6 +179,9 @@ class RoborockZeoProtocol(RoborockEnum):
SILENT_MODE_ON = 240 # rw [independent] use set_silent_mode() for bundled set
SILENT_MODE_START_TIME = 241 # rw [independent] minute-of-day
SILENT_MODE_END_TIME = 242 # rw [independent] minute-of-day
UNKNOWN_243 = (
243 # unknown, not found in plugin bundle; present in MQTT push from some devices, increments with each push
)
DRY_CARE_MODE = 244 # rw [startWith]
SOFTENER_EXPANSION_TYPE = 245 # rw [independent]
SMILE_LIGHT_STATUS = 247 # rw [independent]
Expand Down
27 changes: 21 additions & 6 deletions tests/e2e/__snapshots__/test_device_manager.ambr
Original file line number Diff line number Diff line change
Expand Up @@ -16,17 +16,32 @@
00000000 30 6a 00 20 72 72 2f 6d 2f 69 2f 75 73 65 72 31 |0j. rr/m/i/user1|
00000010 32 33 2f 31 39 36 34 38 66 39 34 2f 7a 65 6f 5f |23/19648f94/zeo_|
00000020 64 75 69 64 00 41 30 31 00 00 23 82 00 00 23 83 |duid.A01..#...#.|
00000030 68 a6 a2 24 00 65 00 30 c5 de 2b f6 a9 ba 32 7e |h..$.e.0..+...2~|
00000040 6b 73 82 bb d8 67 d4 db 7d 80 60 67 80 96 b8 a1 |ks...g..}.`g....|
00000050 c6 bc 9e d2 da 07 fb d3 79 f5 6f 6d 04 9c 71 00 |........y.om..q.|
00000060 48 66 3d 7e 5d fe d3 df 18 e4 26 38 |Hf=~].....&8|
00000030 68 a6 a2 25 00 65 00 30 c5 de 2b f6 a9 ba 32 7e |h..%.e.0..+...2~|
00000040 6b 73 82 bb d8 67 d4 db 7d a3 e2 16 16 7e 83 5f |ks...g..}....~._|
00000050 bb 3d 2b 79 6b d0 52 9b 60 17 2d f4 06 3b a8 66 |.=+yk.R.`.-..;.f|
00000060 f7 20 c0 6a c2 94 1e 84 f8 91 41 67 |. .j......Ag|
[mqtt <]
00000000 30 5e 00 20 72 72 2f 6d 2f 6f 2f 75 73 65 72 31 |0^. rr/m/o/user1|
00000010 32 33 2f 31 39 36 34 38 66 39 34 2f 7a 65 6f 5f |23/19648f94/zeo_|
00000020 64 75 69 64 00 00 00 00 37 41 30 31 00 00 00 00 |duid....7A01....|
00000030 00 00 00 17 68 a6 a2 23 00 66 00 20 c6 d0 06 0c |....h..#.f. ....|
00000030 00 00 00 17 68 a6 a2 23 00 66 00 20 fe 9c 2a 27 |....h..#.f. ..*'|
00000040 da c4 6b 9f 0e cf 2c 56 ba 5b e6 99 1f a1 29 78 |..k...,V.[....)x|
00000050 37 37 42 66 bb d5 70 0e d4 c0 a7 d1 7f ef 8b 33 |77Bf..p........3|
[mqtt >]
00000000 30 6a 00 20 72 72 2f 6d 2f 69 2f 75 73 65 72 31 |0j. rr/m/i/user1|
00000010 32 33 2f 31 39 36 34 38 66 39 34 2f 7a 65 6f 5f |23/19648f94/zeo_|
00000020 64 75 69 64 00 41 30 31 00 00 23 84 00 00 23 85 |duid.A01..#...#.|
00000030 68 a6 a2 26 00 65 00 30 00 5e b7 20 93 7b 53 a7 |h..&.e.0.^. .{S.|
00000040 f9 47 d4 53 49 51 cd 59 c7 92 65 09 f3 7d 20 58 |.G.SIQ.Y..e..} X|
00000050 c6 9e 25 54 cd 4f ab c5 fc 48 0e 67 4a 97 28 59 |..%T.O...H.gJ.(Y|
00000060 14 a6 c6 77 39 ab 2a ef 77 d9 9d 63 |...w9.*.w..c|
[mqtt <]
00000000 30 5e 00 20 72 72 2f 6d 2f 6f 2f 75 73 65 72 31 |0^. rr/m/o/user1|
00000010 32 33 2f 31 39 36 34 38 66 39 34 2f 7a 65 6f 5f |23/19648f94/zeo_|
00000020 64 75 69 64 00 00 00 00 37 41 30 31 00 00 00 00 |duid....7A01....|
00000030 00 00 00 17 68 a6 a2 24 00 66 00 20 c6 d0 06 0c |....h..$.f. ....|
00000040 04 eb 86 8c 96 8c 51 45 4f 8e 96 93 9e 3d de 35 |......QEO....=.5|
00000050 bb a3 92 cf 68 49 69 ba 83 25 cc 5d 77 e8 62 8a |....hIi..%.]w.b.|
00000050 bb a3 92 cf 68 49 69 ba 83 25 cc 5d 43 da 0b 55 |....hIi..%.]C..U|
# ---
# name: test_l01_device
[mqtt >]
Expand Down
6 changes: 4 additions & 2 deletions tests/e2e/test_device_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -527,8 +527,10 @@ async def test_a01_device(
test_topic = TEST_TOPIC_FORMAT.format(duid="zeo_duid")
mqtt_responses: list[bytes] = [
*MQTT_DEFAULT_RESPONSES,
# ACK the Query state call sent below. id is deterministic based on deterministic_message_fixtures
mqtt_packet.gen_publish(test_topic, mid=2, payload=response_builder.build_a01_rpc({"203": 6})),
# ACK the FEATURE_BITS query sent by _discover_features()
mqtt_packet.gen_publish(test_topic, mid=2, payload=response_builder.build_a01_rpc({"237": 1})),
# ACK the Query state call sent below
mqtt_packet.gen_publish(test_topic, mid=3, payload=response_builder.build_a01_rpc({"203": 6})),
]
for response in mqtt_responses:
push_mqtt_response(response)
Expand Down
Loading