Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
25 commits
Select commit Hold shift + click to select a range
b8bdf3a
fix: request Q10 maps without starting cleaning
hCoureau Aug 29, 2026
95d268e
feat: parse Q10 archived map packets
hCoureau Aug 31, 2026
eaac141
feat: expose Q10 map archives
hCoureau Aug 31, 2026
2161763
feat: render Q10 map obstacle markers
hCoureau Aug 29, 2026
ded170d
chore: merge main and resolve Q10 CLI timeout conflict
hCoureau Sep 8, 2026
1a4ffde
chore: merge main for Q10 archive review fixes
hCoureau Sep 11, 2026
cbd6abc
refactor: give Q10 clean-record maps ownership of historical paths
hCoureau Sep 11, 2026
3f8ec38
fix: integrate Q10 archive review changes and upstream timeout
hCoureau Sep 11, 2026
9a68c0d
fix: integrate Q10 archive review changes into obstacle maps
hCoureau Sep 11, 2026
60248c4
fix: assert historical paths are exclusive to clean-record maps
hCoureau Sep 11, 2026
d6979b4
fix: narrow clean-record packet types in obstacle tests
hCoureau Sep 11, 2026
4eb0b9b
refactor: compose Q10 clean-record details and address parser review
hCoureau Sep 17, 2026
49dfd31
refactor: consume composed Q10 clean-record details
hCoureau Sep 17, 2026
dd92094
refactor: compose archived obstacle maps with historical paths
hCoureau Sep 17, 2026
565c7d0
fix: preserve immutable Q10 points in model conformance checks
hCoureau Sep 17, 2026
ac6d83b
fix: include standalone Q10 point conformance correction
hCoureau Sep 17, 2026
d4dbd86
fix: include standalone Q10 point conformance correction
hCoureau Sep 17, 2026
db54bd3
fix: include standalone Q10 point conformance correction
hCoureau Sep 17, 2026
3d3a337
Merge remote-tracking branch 'upstream/main' into fix/pr936-review
hCoureau Sep 27, 2026
9d5d9fe
Merge branch 'fix/pr936-review' into fix/pr937-review
hCoureau Sep 27, 2026
569e6b3
fix: bound Q10 archive selections and preserve response correlation
hCoureau Sep 27, 2026
f7dfee3
Merge branch 'fix/pr937-review' into fix/pr938-conflicts
hCoureau Sep 27, 2026
9b5e77d
Merge remote-tracking branch 'upstream/main' into fix/pr937-review
hCoureau Sep 27, 2026
d6b7c59
test: split Q10 archive decoder cases
hCoureau Sep 27, 2026
f157a8c
Merge branch 'fix/pr937-review' into fix/pr938-conflicts
hCoureau Sep 27, 2026
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
13 changes: 9 additions & 4 deletions roborock/cli.py
Original file line number Diff line number Diff line change
Expand Up @@ -605,6 +605,7 @@ async def maps(ctx, device_id: str):
async def _await_q10_map_push(
properties: Q10PropertiesApi,
predicate: Callable[[], bool],
revision: Callable[[], int],
*,
timeout: float = _Q10_MAP_PUSH_TIMEOUT,
allow_cached_on_timeout: bool = False,
Expand All @@ -617,9 +618,10 @@ async def _await_q10_map_push(
"""
loop = asyncio.get_running_loop()
updated: asyncio.Future[None] = loop.create_future()
initial_revision = revision()

def on_update() -> None:
if predicate() and not updated.done():
if revision() > initial_revision and predicate() and not updated.done():
updated.set_result(None)

unsub = properties.map.add_update_listener(on_update)
Expand Down Expand Up @@ -649,6 +651,7 @@ async def map_image(ctx, device_id: str, output_file: str):
await _await_q10_map_push(
properties,
lambda: properties.map.image_content is not None,
lambda: properties.map.map_revision,
allow_cached_on_timeout=True,
)
image_content = properties.map.image_content
Expand Down Expand Up @@ -697,8 +700,8 @@ async def map_data(ctx, device_id: str, include_path: bool):
async def q10_position(ctx, device_id: str, include_path: bool):
"""Get the current Q10 robot position and live cleaning path.

The Q10 only streams its position/path while it is actively cleaning, so this
will report that no live trace is available for an idle/docked robot.
The Q10 normally streams position/path while it is actively cleaning, so an
idle device may report that no fresh live trace is available.
"""
context: RoborockContext = ctx.obj
device_manager = await context.get_device_manager()
Expand All @@ -710,9 +713,10 @@ async def q10_position(ctx, device_id: str, include_path: bool):
got_trace = await _await_q10_map_push(
properties,
lambda: bool(properties.map.path),
lambda: properties.map.trace_revision,
)
if not got_trace:
click.echo("No live trace available (the robot only reports position while cleaning).")
click.echo("No fresh live trace available.")
return
map_trait = properties.map
position = map_trait.robot_position
Expand Down Expand Up @@ -889,6 +893,7 @@ async def rooms(ctx, device_id: str):
await _await_q10_map_push(
properties,
lambda: properties.map.image_content is not None,
lambda: properties.map.map_revision,
allow_cached_on_timeout=True,
)
click.echo(dump_json({room.id: room.name for room in properties.map.rooms}))
Expand Down
4 changes: 3 additions & 1 deletion roborock/data/b01_q10/b01_q10_containers.py
Original file line number Diff line number Diff line change
Expand Up @@ -144,7 +144,9 @@ class Q10MapInfo(RoborockBase):
"""A saved map reported by ``dpMultiMap``.

Q10 firmware represents the map identifier as a string on the wire. The
value is sent back unchanged in a subsequent ``{"op": "get"}`` request.
value is sent back unchanged in a subsequent ``{"op": "select"}`` detail
request. On Q10 firmware, ``select`` previews a saved map without applying
it as the active map.
"""

id: str
Expand Down
13 changes: 12 additions & 1 deletion roborock/devices/device_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
from roborock.devices.device import DeviceReadyCallback, RoborockDevice
from roborock.diagnostics import Diagnostics, redact_device_data
from roborock.exceptions import RoborockException
from roborock.map.b01_q10_map_parser import B01Q10MapParserConfig
from roborock.map.map_parser import MapParserConfig
from roborock.mqtt.roborock_session import create_lazy_mqtt_session
from roborock.mqtt.session import MqttSession, SessionUnauthorizedHook
Expand Down Expand Up @@ -262,7 +263,17 @@ def device_creator(home_data: HomeData, device: HomeDataDevice, product: HomeDat
if "ss" in model_part:
b01_q10_channel = create_b01_q10_channel(mqtt_channel)
channel = b01_q10_channel
trait = b01.q10.create(channel)
trait = b01.q10.create(
channel,
map_parser_config=(
B01Q10MapParserConfig(
map_scale=map_parser_config.map_scale,
drawables=map_parser_config.drawables,
)
if map_parser_config
else None
),
)
elif "sc" in model_part:
# Q7 devices start with 'sc' in their model naming.
b01_q7_channel = create_b01_q7_channel(device, product, mqtt_channel)
Expand Down
48 changes: 40 additions & 8 deletions roborock/devices/traits/b01/q10/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,13 @@
from roborock.data.containers import RoborockBase
from roborock.devices.rpc.b01_q10_channel import B01Q10Channel
from roborock.devices.traits import Trait
from roborock.map.b01_q10_map_parser import Q10MapPacket, Q10TracePacket
from roborock.map.b01_q10_map_parser import (
B01Q10MapParserConfig,
Q10CleanRecordDetail,
Q10MapPacket,
Q10MapPacketKind,
Q10TracePacket,
)
from roborock.protocols.b01_q10_protocol import Q10DpsUpdate, Q10Message

from .button_light import ButtonLightTrait
Expand Down Expand Up @@ -92,7 +98,12 @@ class Q10PropertiesApi(Trait):
clean_history: CleanHistoryTrait
"""Trait for fetching the device clean-record history (``dpCleanRecord``)."""

def __init__(self, channel: B01Q10Channel) -> None:
def __init__(
self,
channel: B01Q10Channel,
*,
map_parser_config: B01Q10MapParserConfig,
) -> None:
"""Initialize the B01Props API."""
self._channel = channel
self.command = CommandTrait(channel)
Expand All @@ -106,10 +117,17 @@ def __init__(self, channel: B01Q10Channel) -> None:
self.network_info = NetworkInfoTrait()
self.consumable = ConsumableTrait()
self._map_dps = MapDpsTrait()
self.maps = MapsTrait(self.command)
self.map = MapContentTrait(self._map_dps, self.maps, self.command)
self.maps = MapsTrait(self.command, map_parser_config=map_parser_config)
self.map = MapContentTrait(
self._map_dps,
self.command,
map_parser_config=map_parser_config,
)
self.clean_history = CleanHistoryTrait(
self.command,
map_parser_config=map_parser_config,
)
self.vacuum = VacuumTrait(self.command, self.status, self.map)
self.clean_history = CleanHistoryTrait(self.command)
# Read-model traits updated from the device's DPS push stream.
self._updatable_traits = [
self.status,
Expand All @@ -131,6 +149,8 @@ async def start(self) -> None:

async def close(self) -> None:
"""Close any resources held by the trait."""
self.maps.close()
self.clean_history.close()
await self.vacuum.close()
if self._subscribe_task is not None:
self._subscribe_task.cancel()
Expand Down Expand Up @@ -158,7 +178,12 @@ def _handle_message(self, message: Q10Message) -> None:
Map-list DPS responses and other DPS updates feed the read-model traits.
"""
if isinstance(message, Q10MapPacket):
self.map.update_from_map_packet(message)
if message.kind is Q10MapPacketKind.CURRENT:
self.map.update_from_map_packet(message)
elif message.kind is Q10MapPacketKind.SAVED_MAP_DETAIL:
self.maps.update_from_map_packet(message)
elif isinstance(message, Q10CleanRecordDetail):
self.clean_history.update_from_detail(message)
elif isinstance(message, Q10TracePacket):
self.map.update_from_trace_packet(message)
elif isinstance(message, Q10DpsUpdate):
Expand All @@ -179,6 +204,13 @@ def as_dict(self) -> dict[str, Any]:
return result


def create(channel: B01Q10Channel) -> Q10PropertiesApi:
def create(
channel: B01Q10Channel,
*,
map_parser_config: B01Q10MapParserConfig | None = None,
) -> Q10PropertiesApi:
"""Create traits for B01 devices."""
return Q10PropertiesApi(channel)
return Q10PropertiesApi(
channel,
map_parser_config=map_parser_config or B01Q10MapParserConfig(),
)
143 changes: 139 additions & 4 deletions roborock/devices/traits/b01/q10/clean_history.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,8 +10,9 @@
a ``dpCleanRecord`` envelope into a :class:`CleanRecordPush`, and the trait applies it.
"""

import asyncio
import logging
from dataclasses import dataclass, field
from dataclasses import dataclass, field, replace
from typing import Any

from roborock.data.b01_q10.b01_q10_code_mappings import (
Expand All @@ -22,6 +23,17 @@
YXStartMethod,
)
from roborock.data.b01_q10.b01_q10_containers import Q10CleanRecord
from roborock.exceptions import RoborockException, RoborockTimeout
from roborock.map.b01_q10_map_parser import (
B01Q10MapParserConfig,
Q10CleanRecordDetail,
Q10HistoricalTracePacket,
Q10MapPacket,
Q10MapPacketKind,
Q10Obstacle,
Q10Point,
)
from roborock.map.b01_q10_render import Q10MapOverlays, render_q10_map

from .command import CommandTrait
from .common import UpdatableTrait
Expand All @@ -33,6 +45,7 @@
]

_LOGGER = logging.getLogger(__name__)
_DETAIL_TIMEOUT = 30.0

_RECORD_FIELD_COUNT = 12

Expand Down Expand Up @@ -115,12 +128,32 @@ class CleanHistoryTrait(UpdatableTrait):
or a single ``op:"notify"`` record) rather than a flat data-point-to-field map.
"""

def __init__(self, command: CommandTrait) -> None:
_command: CommandTrait

def __init__(
self,
command: CommandTrait,
*,
map_parser_config: B01Q10MapParserConfig,
) -> None:
"""Initialize the clean history trait."""
UpdatableTrait.__init__(self, command, _LOGGER)
self._command = command
self._converter = CleanRecordConverter()
self._map_parser_config = map_parser_config
self.records: list[Q10CleanRecord] = []
"""Decoded clean records, most recent first."""
self.detail: Q10CleanRecordDetail | None = None
"""Most recently pushed ``03 01`` clean-record map detail."""
self.detail_record: Q10CleanRecord | None = None
"""Record associated with :attr:`detail_packet`, when requested here."""
self.detail_image_content: bytes | None = None
"""Rendered clean-record detail image, if decoding succeeded."""
self._pending_detail_record: Q10CleanRecord | None = None
self._detail_response: asyncio.Future[None] | None = None
self._detail_task: asyncio.Task[None] | None = None
self._detail_requested = False
self._discard_detail_response = False

@property
def last_record(self) -> Q10CleanRecord | None:
Expand All @@ -134,13 +167,84 @@ async def refresh(self) -> None:
asynchronously on the device stream and populate :attr:`records` once
:meth:`update_from_dps` processes the ``dpCleanRecord`` push.
"""
if self._command is None:
raise ValueError("Trait is read-only; no command channel was provided")
await self._command.send(
B01_Q10_DP.COMMON,
params={str(B01_Q10_DP.CLEAN_RECORD.code): {"op": "list"}},
)

async def refresh_detail(self, record: Q10CleanRecord) -> None:
"""Request the saved map and path for one clean record.

The complete 12-field raw record is the firmware's detail identifier;
the shorter human-facing record ID is not accepted. Only one request
may be outstanding because ``03 01`` responses carry no correlation ID.
Wait up to 30 seconds for its response. Timeout, cancellation, or send
failure blocks further selections until a late response is discarded.
Reconnecting does not reset this barrier. Other clients selecting
records concurrently cannot be distinguished by this protocol.
"""
if not record.raw or not record.map_len:
raise RoborockException("The Q10 clean record has no saved map detail")
if self._pending_detail_record is not None:
raise RoborockException("A Q10 clean-record detail request is already pending")
if self._discard_detail_response:
raise RoborockException("A previous Q10 clean-record selection is unresolved; wait for its late response")
response: asyncio.Future[None] = asyncio.get_running_loop().create_future()
self._detail_response = response
self._detail_task = asyncio.current_task()
self._detail_requested = True
self._pending_detail_record = replace(record)
try:
async with asyncio.timeout(_DETAIL_TIMEOUT):
await self._command.send(
B01_Q10_DP.COMMON,
{str(B01_Q10_DP.CLEAN_RECORD.code): {"op": "select", "id": record.raw}},
)
await response
except TimeoutError as ex:
raise RoborockTimeout("Q10 archive detail request timed out") from ex
finally:
if self._detail_response is response:
if self._pending_detail_record is not None:
self._discard_detail_response = True
self._pending_detail_record = None
self._detail_response = None
self._detail_task = None
if not response.done():
response.cancel()

def close(self) -> None:
"""Cancel an active archive selection during device teardown."""
if self._pending_detail_record is not None:
self._discard_detail_response = True
if self._detail_task is not None:
self._detail_task.cancel()
if self._detail_response is not None:
self._detail_response.cancel()
self._pending_detail_record = None
self._detail_response = None
self._detail_task = None

@property
def detail_packet(self) -> Q10MapPacket | None:
"""The map from the most recently received clean-record detail."""
return self.detail.map if self.detail else None

@property
def detail_trace(self) -> Q10HistoricalTracePacket | None:
"""Historical path embedded in the selected clean-record detail."""
return self.detail.trace if self.detail else None

@property
def detail_path(self) -> list[Q10Point]:
"""Historical path points for the selected clean record."""
return list(self.detail_trace.points) if self.detail_trace else []

@property
def detail_obstacles(self) -> list[Q10Obstacle]:
"""Obstacle markers embedded in the selected clean-record map."""
return list(self.detail_packet.obstacles) if self.detail_packet else []

def update_from_dps(self, decoded_dps: dict[B01_Q10_DP, Any]) -> None:
"""Apply a ``dpCleanRecord`` push (a full list reply or a single notify)."""
envelope = decoded_dps.get(B01_Q10_DP.CLEAN_RECORD)
Expand All @@ -151,6 +255,37 @@ def update_from_dps(self, decoded_dps: dict[B01_Q10_DP, Any]) -> None:
return
self._apply(push)

def update_from_detail(self, detail: Q10CleanRecordDetail) -> None:
"""Store and render a pushed clean-record detail map."""
if detail.map.kind is not Q10MapPacketKind.CLEAN_RECORD_DETAIL:
raise ValueError(f"Expected a Q10 clean-record detail packet, got {detail.map.kind.value}")
if self._discard_detail_response:
self._discard_detail_response = False
return
response = self._detail_response
if self._detail_requested and (response is None or response.done()):
# Cancellation can finish the future before refresh_detail's
# finally block runs. Drain this response in that interval too.
if response is not None and response.cancelled():
self._pending_detail_record = None
return
self.detail_record = self._pending_detail_record
self._pending_detail_record = None
self.detail = detail
try:
self.detail_image_content = render_q10_map(
detail.map,
detail.trace,
Q10MapOverlays(),
config=self._map_parser_config,
)
except RoborockException:
_LOGGER.debug("Failed to render Q10 clean-record detail", exc_info=True)
self.detail_image_content = None
if response is not None and not response.done():
response.set_result(None)
self._notify_update()

def _apply(self, push: CleanRecordPush) -> None:
"""Merge or replace the records from ``push``, then sort newest-first and notify."""
if push.replace:
Expand Down
Loading
Loading