SWE-Race › Tasks › python-roborock-python-roborock-632 ← prevnext →

python-roborock-python-roborock-632

Python-roborock/python-roborockcleansinglemerged 2025-12-05Apache-2.0fix: 1 file, +56 −73 fail-to-pass · 6 pass-to-pass
Results
Modelsolved / attemptsmedian stepsmedian costattempts
GPT-5.6 Luna6/613$0.0121✓ 2✓ 3✓ 4✓ 5✓ 6✓
DeepSeek V4 Flash2/245$0.0731✓ 2✓
GLM-5.3 Flash2/211$0.0031✓ 2✓
The prompt the agent sees

MQTT topic subscriptions are currently removed as soon as their last callback is unsubscribed, which can disrupt later RPCs that reuse the same topic connection. When a topic becomes inactive, keep the underlying subscription available during a configurable idle period, then unsubscribe it only if it remains unused until that period expires.

If a topic has multiple callbacks, removing one callback must not begin idle cleanup while other callbacks remain active. Calling `subscribe(topic, callback)` again before the idle period expires must prevent cleanup and reuse the existing subscription. Once no callbacks remain and the topic stays idle through the configured period, the MQTT client must call `unsubscribe(topic)` exactly once. An `aiomqtt.MqttError` during this unsubscribe is logged rather than raised. `close()` must also clean up pending topic idle handling along with the connection tasks.

`RoborockMqttSession.__init__` must accept `topic_idle_timeout` as a `datetime.timedelta`, defaulting to 60 seconds, and use it as the topic idle period. Tests construct `RoborockMqttSession(params, topic_idle_timeout=datetime.timedelta(milliseconds=50))`. `subscribe(topic, callback)` remains awaitable and returns an unsubscribe callable. Calling that callable removes only its callback immediately; cleanup begins only when it was the last callback.

The existing session error behavior must remain intact. `create_mqtt_session` must create a connected session, and `publish("topic-1", message=b"payload")` must raise `MqttSessionException` with the message matching `"Error publishing message"` when the mocked MQTT client's publish operation raises `aiomqtt.MqttError`. If the mocked client's subscribe operation raises `aiomqtt.MqttError`, `subscribe("topic-1", subscriber.append)` must raise `MqttSessionException` with the message matching `"Error subscribing to topic"` and must not invoke the callback. The session exposes its connected state through `connected`, and callers can use `start()` and `close()`.

Tests replace `aiomqtt.Client` at `roborock.mqtt.roborock_session.aiomqtt.Client` with an asynchronous mock client and assert the client's `unsubscribe` call count and arguments. After a topic remains idle for the configured duration, it must be called exactly once with that topic; if the topic is resubscribed before the duration expires, it must not be called.

Hidden tests · 3 fail-to-pass, 6 pass-to-passrun after the agent submits, in a clean verifier
test_idle_timeout_multiple_callbackstest_idle_timeout_resubscribetest_idle_timeout_unsubscribe
Test patch · 198 lines
diff --git a/tests/mqtt/test_roborock_session.py b/tests/mqtt/test_roborock_session.py
index bb3f1bdb..f3b10139 100644
--- a/tests/mqtt/test_roborock_session.py
+++ b/tests/mqtt/test_roborock_session.py
@@ -11,7 +11,7 @@
 import paho.mqtt.client as mqtt
 import pytest
 
-from roborock.mqtt.roborock_session import create_mqtt_session
+from roborock.mqtt.roborock_session import RoborockMqttSession, create_mqtt_session
 from roborock.mqtt.session import MqttParams, MqttSessionException
 from tests import mqtt_packet
 from tests.conftest import FakeSocketHandler
@@ -80,6 +80,23 @@ def fast_backoff_fixture() -> Generator[None, None, None]:
         yield
 
 
+@pytest.fixture
+def mock_mqtt_client() -> Generator[AsyncMock, None, None]:
+    """Fixture to create a mock MQTT client with patched aiomqtt.Client."""
+    mock_client = AsyncMock()
+    mock_client.messages = FakeAsyncIterator()
+
+    mock_aenter = AsyncMock()
+    mock_aenter.return_value = mock_client
+
+    mock_shim = Mock()
+    mock_shim.return_value.__aenter__ = mock_aenter
+    mock_shim.return_value.__aexit__ = AsyncMock()
+
+    with patch("roborock.mqtt.roborock_session.aiomqtt.Client", mock_shim):
+        yield mock_client
+
+
 @pytest.fixture
 def push_response(response_queue: Queue, fake_socket_handler: FakeSocketHandler) -> Callable[[bytes], None]:
     """Fixtures to push messages."""
@@ -195,52 +212,34 @@ async def __anext__(self) -> None:
             await asyncio.sleep(1)
 
 
-async def test_publish_failure() -> None:
+async def test_publish_failure(mock_mqtt_client: AsyncMock) -> None:
     """Test an MQTT error is received when publishing a message."""
 
-    mock_client = AsyncMock()
-    mock_client.messages = FakeAsyncIterator()
-
-    mock_aenter = AsyncMock()
-    mock_aenter.return_value = mock_client
-
-    with patch("roborock.mqtt.roborock_session.aiomqtt.Client.__aenter__", mock_aenter):
-        session = await create_mqtt_session(FAKE_PARAMS)
-        assert session.connected
+    session = await create_mqtt_session(FAKE_PARAMS)
+    assert session.connected
 
-        mock_client.publish.side_effect = aiomqtt.MqttError
+    mock_mqtt_client.publish.side_effect = aiomqtt.MqttError
 
-        with pytest.raises(MqttSessionException, match="Error publishing message"):
-            await session.publish("topic-1", message=b"payload")
+    with pytest.raises(MqttSessionException, match="Error publishing message"):
+        await session.publish("topic-1", message=b"payload")
 
-        await session.close()
+    await session.close()
 
 
-async def test_subscribe_failure() -> None:
+async def test_subscribe_failure(mock_mqtt_client: AsyncMock) -> None:
     """Test an MQTT error while subscribing."""
 
-    mock_client = AsyncMock()
-    mock_client.messages = FakeAsyncIterator()
-
-    mock_aenter = AsyncMock()
-    mock_aenter.return_value = mock_client
-
-    mock_shim = Mock()
-    mock_shim.return_value.__aenter__ = mock_aenter
-    mock_shim.return_value.__aexit__ = AsyncMock()
-
-    with patch("roborock.mqtt.roborock_session.aiomqtt.Client", mock_shim):
-        session = await create_mqtt_session(FAKE_PARAMS)
-        assert session.connected
+    session = await create_mqtt_session(FAKE_PARAMS)
+    assert session.connected
 
-        mock_client.subscribe.side_effect = aiomqtt.MqttError
+    mock_mqtt_client.subscribe.side_effect = aiomqtt.MqttError
 
-        subscriber1 = Subscriber()
-        with pytest.raises(MqttSessionException, match="Error subscribing to topic"):
-            await session.subscribe("topic-1", subscriber1.append)
+    subscriber1 = Subscriber()
+    with pytest.raises(MqttSessionException, match="Error subscribing to topic"):
+        await session.subscribe("topic-1", subscriber1.append)
 
-        assert not subscriber1.messages
-        await session.close()
+    assert not subscriber1.messages
+    await session.close()
 
 
 async def test_restart(push_response: Callable[[bytes], None]) -> None:
@@ -279,3 +278,91 @@ async def test_restart(push_response: Callable[[bytes], None]) -> None:
     assert subscriber.messages == [b"12345", b"67890"]
 
     await session.close()
+
+
+async def test_idle_timeout_resubscribe(mock_mqtt_client: AsyncMock) -> None:
+    """Test that resubscribing before idle timeout cancels the unsubscribe."""
+
+    # Create session with idle timeout
+    session = RoborockMqttSession(FAKE_PARAMS, topic_idle_timeout=datetime.timedelta(seconds=5))
+    await session.start()
+    assert session.connected
+
+    topic = "test/topic"
+    subscriber1 = Subscriber()
+    unsub1 = await session.subscribe(topic, subscriber1.append)
+
+    # Unsubscribe to start idle timer
+    unsub1()
+
+    # Resubscribe before idle timeout expires (should cancel timer)
+    subscriber2 = Subscriber()
+    await session.subscribe(topic, subscriber2.append)
+
+    # Give a brief moment for any async operations to complete
+    await asyncio.sleep(0.01)
+
+    # unsubscribe should NOT have been called because we resubscribed
+    mock_mqtt_client.unsubscribe.assert_not_called()
+
+    await session.close()
+
+
+async def test_idle_timeout_unsubscribe(mock_mqtt_client: AsyncMock) -> None:
+    """Test that unsubscribe happens after idle timeout expires."""
+
+    # Create session with very short idle timeout for fast test
+    session = RoborockMqttSession(FAKE_PARAMS, topic_idle_timeout=datetime.timedelta(milliseconds=50))
+    await session.start()
+    assert session.connected
+
+    topic = "test/topic"
+    subscriber = Subscriber()
+    unsub = await session.subscribe(topic, subscriber.append)
+
+    # Unsubscribe to start idle timer
+    unsub()
+
+    # Wait for idle timeout plus a small buffer
+    await asyncio.sleep(0.1)
+
+    # unsubscribe should have been called after idle timeout
+    mock_mqtt_client.unsubscribe.assert_called_once_with(topic)
+
+    await session.close()
+
+
+async def test_idle_timeout_multiple_callbacks(mock_mqtt_client: AsyncMock) -> None:
+    """Test that unsubscribe is delayed when multiple subscribers exist."""
+
+    # Create session with very short idle timeout for fast test
+    session = RoborockMqttSession(FAKE_PARAMS, topic_idle_timeout=datetime.timedelta(milliseconds=50))
+    await session.start()
+    assert session.connected
+
+    topic = "test/topic"
+    subscriber1 = Subscriber()
+    subscriber2 = Subscriber()
+
+    unsub1 = await session.subscribe(topic, subscriber1.append)
+    unsub2 = await session.subscribe(topic, subscriber2.append)
+
+    # Unsubscribe first callback (should NOT start timer, subscriber2 still active)
+    unsub1()
+
+    # Brief wait to ensure no timer fires
+    await asyncio.sleep(0.1)
+
+    # unsubscribe should NOT have been called because subscriber2 is still active
+    mock_mqtt_client.unsubscribe.assert_not_called()
+
+    # Unsubscribe second callback (NOW timer should start)
+    unsub2()
+
+    # Wait for idle timeout plus a small buffer
+    await asyncio.sleep(0.1)
+
+    # Now unsubscribe should have been called
+    mock_mqtt_client.unsubscribe.assert_called_once_with(topic)
+
+    await session.close()
Reference fix · 1 file, +56 −7the upstream merge, used only for grading calibration

The agent could not see this: the repository holds one commit and the sandbox has no network. Leak audit.

roborock/mqtt/roborock_session.py

diff --git a/roborock/mqtt/roborock_session.py b/roborock/mqtt/roborock_session.py
index 53a337e7..b35e3910 100644
--- a/roborock/mqtt/roborock_session.py
+++ b/roborock/mqtt/roborock_session.py
@@ -24,7 +24,8 @@
 _LOGGER = logging.getLogger(__name__)
 _MQTT_LOGGER = logging.getLogger(f"{__name__}.aiomqtt")
 
-KEEPALIVE = 60
+CLIENT_KEEPALIVE = datetime.timedelta(seconds=120)
+TOPIC_KEEPALIVE = datetime.timedelta(seconds=60)
 
 # Exponential backoff parameters
 MIN_BACKOFF_INTERVAL = datetime.timedelta(seconds=10)
@@ -47,7 +48,11 @@ class RoborockMqttSession(MqttSession):
     re-established.
     """
 
-    def __init__(self, params: MqttParams):
+    def __init__(
+        self,
+        params: MqttParams,
+        topic_idle_timeout: datetime.timedelta = TOPIC_KEEPALIVE,
+    ):
         self._params = params
         self._reconnect_task: asyncio.Task[None] | None = None
         self._healthy = False
@@ -57,6 +62,8 @@ def __init__(self, params: MqttParams):
         self._client_lock = asyncio.Lock()
         self._listeners: CallbackMap[str, bytes] = CallbackMap(_LOGGER)
         self._connection_task: asyncio.Task[None] | None = None
+        self._topic_idle_timeout = topic_idle_timeout
+        self._idle_timers: dict[str, asyncio.Task[None]] = {}
 
     @property
     def connected(self) -> bool:
@@ -86,11 +93,15 @@ async def start(self) -> None:
     async def close(self) -> None:
         """Cancels the MQTT loop and shutdown the client library."""
         self._stop = True
-        tasks = [task for task in [self._connection_task, self._reconnect_task] if task]
+        tasks = [task for task in [self._connection_task, self._reconnect_task, *self._idle_timers.values()] if task]
+        self._connection_task = None
+        self._reconnect_task = None
+        self._idle_timers.clear()
+
         for task in tasks:
             task.cancel()
         try:
-            await asyncio.gather(*tasks)
+            await asyncio.gather(*tasks, return_exceptions=True)
         except asyncio.CancelledError:
             pass
 
@@ -183,7 +194,7 @@ async def _mqtt_client(self, params: MqttParams) -> aiomqtt.Client:
                 port=params.port,
                 username=params.username,
                 password=params.password,
-                keepalive=KEEPALIVE,
+                keepalive=int(CLIENT_KEEPALIVE.total_seconds()),
                 protocol=aiomqtt.ProtocolVersion.V5,
                 tls_params=TLSParameters() if params.tls else None,
                 timeout=params.timeout,
@@ -210,9 +221,17 @@ async def subscribe(self, topic: str, callback: Callable[[bytes], None]) -> Call
         The callback will be called with the message payload as a bytes object. The callback
         should not block since it runs in the async loop. It should not raise any exceptions.
 
-        The returned callable unsubscribes from the topic when called.
+        The returned callable unsubscribes from the topic when called, but will delay actual
+        unsubscription for the idle timeout period. If a new subscription comes in during the
+        timeout, the timer is cancelled and the subscription is reused.
         """
         _LOGGER.debug("Subscribing to topic %s", topic)
+
+        # If there is an idle timer for this topic, cancel it (reuse subscription)
+        if idle_timer := self._idle_timers.pop(topic, None):
+            idle_timer.cancel()
+            _LOGGER.debug("Cancelled idle timer for topic %s (reused subscription)", topic)
+
         unsub = self._listeners.add_callback(topic, callback)
 
         async with self._client_lock:
@@ -221,11 +240,41 @@ async def subscribe(self, topic: str, callback: Callable[[bytes], None]) -> Call
                 try:
                     await self._client.subscribe(topic)
                 except MqttError as err:
+                    # Clean up the callback if subscription fails
+                    unsub()
                     raise MqttSessionException(f"Error subscribing to topic: {err}") from err
             else:
                 _LOGGER.debug("Client not connected, will establish subscription later")
 
-        return unsub
+        def schedule_unsubscribe():
+            async def idle_unsubscribe():
+                try:
+                    await asyncio.sleep(self._topic_idle_timeout.total_seconds())
+                    # Only unsubscribe if there are no callbacks left for this topic
+                    if not self._listeners.get_callbacks(topic):
+                        async with self._client_lock:
+                            if self._client:
+                                _LOGGER.debug("Idle timeout expired, unsubscribing from topic %s", topic)
+                                try:
+                                    await self._client.unsubscribe(topic)
+                                except MqttError as err:
+                                    _LOGGER.warning("Error unsubscribing from topic %s: %s", topic, err)
+                    # Clean up timer from dict
+                    self._idle_timers.pop(topic, None)
+                except asyncio.CancelledError:
+                    _LOGGER.debug("Idle unsubscribe for topic %s cancelled", topic)
+
+            # Start the idle timer task
+            task = asyncio.create_task(idle_unsubscribe())
+            self._idle_timers[topic] = task
+
+        def delayed_unsub():
+            unsub()  # Remove the callback from CallbackMap
+            # If no more callbacks for this topic, start idle timer
+            if not self._listeners.get_callbacks(topic):
+                schedule_unsubscribe()
+
+        return delayed_unsub
 
     async def publish(self, topic: str, message: bytes) -> None:
         """Publish a message on the topic."""

What each model did

Submitted diff first, then every command. Outputs show 12 lines; expand for the rest.
GLM-5.3 Flash · attempt 2 · solved · 11 steps · $0.003
Submitted diff · 86 lines
diff --git a/roborock/mqtt/roborock_session.py b/roborock/mqtt/roborock_session.py
index 53a337e..ce7bd56 100644
--- a/roborock/mqtt/roborock_session.py
+++ b/roborock/mqtt/roborock_session.py
@@ -47,8 +47,11 @@ class RoborockMqttSession(MqttSession):
     re-established.
     """
 
-    def __init__(self, params: MqttParams):
+    def __init__(
+        self, params: MqttParams, topic_idle_timeout: datetime.timedelta = datetime.timedelta(seconds=60)
+    ):
         self._params = params
+        self._topic_idle_timeout = topic_idle_timeout
         self._reconnect_task: asyncio.Task[None] | None = None
         self._healthy = False
         self._stop = False
@@ -57,6 +60,7 @@ class RoborockMqttSession(MqttSession):
         self._client_lock = asyncio.Lock()
         self._listeners: CallbackMap[str, bytes] = CallbackMap(_LOGGER)
         self._connection_task: asyncio.Task[None] | None = None
+        self._topic_idle_tasks: dict[str, asyncio.Task[None]] = {}
 
     @property
     def connected(self) -> bool:
@@ -87,6 +91,9 @@ class RoborockMqttSession(MqttSession):
         """Cancels the MQTT loop and shutdown the client library."""
         self._stop = True
         tasks = [task for task in [self._connection_task, self._reconnect_task] if task]
+        for idle_task in self._topic_idle_tasks.values():
+            idle_task.cancel()
+        self._topic_idle_tasks.clear()
         for task in tasks:
             task.cancel()
         try:
@@ -213,6 +220,12 @@ class RoborockMqttSession(MqttSession):
         The returned callable unsubscribes from the topic when called.
         """
         _LOGGER.debug("Subscribing to topic %s", topic)
+        # If there is a pending idle cleanup for this topic, cancel it and reuse
+        # the existing subscription.
+        if idle_task := self._topic_idle_tasks.pop(topic, None):
+            idle_task.cancel()
+            _LOGGER.debug("Cancelled idle cleanup for topic %s", topic)
+
         unsub = self._listeners.add_callback(topic, callback)
 
         async with self._client_lock:
@@ -221,11 +234,36 @@ class RoborockMqttSession(MqttSession):
                 try:
                     await self._client.subscribe(topic)
                 except MqttError as err:
+                    unsub()
                     raise MqttSessionException(f"Error subscribing to topic: {err}") from err
             else:
                 _LOGGER.debug("Client not connected, will establish subscription later")
 
-        return unsub
+        def unsubscribe() -> None:
+            """Remove the callback, starting idle cleanup if it was the last one."""
+            unsub()
+            if not self._listeners.get_callbacks(topic):
+                self._start_topic_idle_cleanup(topic)
+
+        return unsubscribe
+
+    def _start_topic_idle_cleanup(self, topic: str) -> None:
+        """Schedule unsubscribing from a topic after it has been idle for a while."""
+
+        async def _idle_unsubscribe() -> None:
+            """Unsubscribe from the topic if it is still idle."""
+            await asyncio.sleep(self._topic_idle_timeout.total_seconds())
+            self._topic_idle_tasks.pop(topic, None)
+            _LOGGER.debug("Topic %s idle for %s, unsubscribing", topic, self._topic_idle_timeout)
+            async with self._client_lock:
+                if self._client is None:
+                    return
+                try:
+                    await self._client.unsubscribe(topic)
+                except MqttError as err:
+                    _LOGGER.warning("Error unsubscribing from topic %s: %s", topic, err)
+
+        self._topic_idle_tasks[topic] = asyncio.create_task(_idle_unsubscribe())
 
     async def publish(self, topic: str, message: bytes) -> None:
         """Publish a message on the topic."""
step 1
$ find /repo -name "*.py" | grep -i mqtt; ls /repo
output · 24 lines
/repo/roborock/devices/mqtt_channel.py
/repo/roborock/mqtt/health_manager.py
/repo/roborock/mqtt/roborock_session.py
/repo/roborock/mqtt/session.py
/repo/roborock/mqtt/__init__.py
/repo/roborock/version_1_apis/roborock_mqtt_client_v1.py
/repo/roborock/version_a01_apis/roborock_mqtt_client_a01.py
/repo/tests/devices/test_mqtt_channel.py
/repo/tests/mqtt/test_health_manager.py
/repo/tests/mqtt/test_roborock_session.py
/repo/tests/mqtt_packet.py
CHANGELOG.md
12 more lines
LICENSE
README.md
SUPPORTED_FEATURES.md
commitlint.config.mjs
device_info.yaml
docs
examples
mypy.ini
pyproject.toml
roborock
tests
uv.lock
step 2
$ cat /repo/roborock/mqtt/roborock_session.py; sed -n 1,80p /repo/roborock/mqtt/session.py
output · 45 lines
"""An MQTT session for sending and receiving messages.

See create_mqtt_session for a factory function to create an MQTT session.

This is a thin wrapper around the async MQTT client that handles dispatching messages
from a topic to a callback function, since the async MQTT client does not
support this out of the box. It also handles the authentication process and
receiving messages from the vacuum cleaner.
"""

import asyncio
import datetime
33 more lines
import logging
from collections.abc import Callable
from contextlib import asynccontextmanager

import aiomqtt
from aiomqtt import MqttError, TLSParameters

from roborock.callbacks import CallbackMap

from .session import MqttParams, MqttSession, MqttSessionException

_LOGGER = logging.getLogger(__name__)
_MQTT_LOGGER = logging.getLogger(f"{__name__}.aiomqtt")

KEEPALIVE = 60

# Exponential backoff parameters
MIN_BACKOFF_INTERVAL = datetime.timed
…[8629 chars omitted]…
   This will raise an exception if the message could not be sent.
        """

    @abstractmethod
    async def restart(self) -> None:
        """Force the session to disconnect and reconnect."""

    @abstractmethod
    async def close(self) -> None:
        """Cancels the mqtt loop"""


class MqttSessionException(RoborockException):
    """ "Raised when there is an error communicating with MQTT."""
step 3
$ sed -n 100,220p /repo/roborock/mqtt/roborock_session.py; echo ====; cat /repo/roborock/callbacks.py 2>/dev/null | head -60; echo ====; sed -n 1,120p /repo/tests/mqtt/test_roborock_session.py
output · 35 lines
"""Force the session to disconnect and reconnect.

        The active connection task will be cancelled and restarted in the background, retried by
        the reconnect loop. This is a no-op if there is no active connection.
        """
        _LOGGER.info("Forcing MQTT session restart")
        if self._connection_task:
            self._connection_task.cancel()
        else:
            _LOGGER.debug("No message loop task to cancel")

    async def _run_reconnect_loop(self, start_future: asyncio.Future[None] | None) -> None:
23 more lines
        """Run the MQTT loop."""
        _LOGGER.info("Starting MQTT session")
        while True:
            try:
                self._connection_task = asyncio.create_task(self._run_connection(start_future))
                await self._connection_task
            except asyncio.CancelledError:
                _LOGGER.debug("MQTT connection task can
…[8629 chars omitted]…
 a message to the subscriber."""
        self.messages.append(message)
        self.event.set()

    async def wait(self) -> None:
        """Wait for a message to be received."""
        await self.event.wait()
        self.event.clear()


async def test_session(push_response: Callable[[bytes], None]) -> None:
    """Test the MQTT session."""

    push_response(mqtt_packet.gen_connack(rc=0, flags=2))
step 4
$ sed -n 200,260p /repo/roborock/mqtt/roborock_session.py; sed -n 120,400p /repo/tests/mqtt/test_roborock_session.py
output · 33 lines
await client.subscribe(topic)

                yield client
        finally:
            async with self._client_lock:
                self._client = None

    async def subscribe(self, topic: str, callback: Callable[[bytes], None]) -> Callable[[], None]:
        """Subscribe to messages on the specified topic and invoke the callback for new messages.

        The callback will be called with the message payload as a bytes object. The callback
        should not block since it runs in the async loop. It should not raise any exceptions.
21 more lines
        The returned callable unsubscribes from the topic when called.
        """
        _LOGGER.debug("Subscribing to topic %s", topic)
        unsub = self._listeners.add_callback(topic, callback)

        async with self._client_lock:
            if self._client:
                _LOGGER.debug("Establishing subscription to t
…[6669 chars omitted]…
   push_response(mqtt_packet.gen_connack(rc=0, flags=2))

    # And a suback for the resubscription. Since we created a new client,
    # the message ID resets to 1.
    push_response(mqtt_packet.gen_suback(mid=1))

    push_response(mqtt_packet.gen_publish("topic-1", mid=4, payload=b"67890"))
    await subscriber.wait()
    assert subscriber.messages == [b"12345", b"67890"]

    await session.close()
step 5
$ sed -n 60,120p /repo/roborock/callbacks.py
output · 41 lines
if not cb_list:
                    del self._callbacks[key]

        return remove_callback

    def get_callbacks(self, key: K) -> list[Callable[[V], None]]:
        """Get all callbacks for a specific key."""
        return self._callbacks.get(key, [])

    def __call__(self, key: K, value: V) -> None:
        """Invoke all callbacks for a specific key."""
        for callback in self.get_callbacks(key):
29 more lines
            safe_callback(callback, self._logger)(value)


class CallbackList(Generic[V]):
    """A list of callbacks that can be invoked.

    This combines a list of callbacks into a single callable. Callers can add
    additional callbacks to the list at any time.
    """

    def __init__(self, logger: logging.Logger | None = None) -> None:
        self._callbacks: list[Callable[[V], None]] = []
        self._logger = logger or _LOGGER

    def add_callback(self
…[805 chars omitted]…
to a list of values. The callback is then invoked
    for each value in the list.

    Any failures during decoding or invoking the callbacks will be logged.
    """
    if logger is None:
        logger = _LOGGER

    safe_cb = safe_callback(callback, logger)

    def wrapper(data: K) -> None:
        if not (messages := decoder(data)):
            logger.warning("Failed to decode message: %s", data)
step 6
$ cd /repo && python - <<'PY'
output · 2 lines
56:        def remove_callback() -> None:
63:        return remove_callback
step 7
$ cd /repo && python - <<'PY'
step 8
$ cd /repo && cat > /tmp/repro.py <<'EOF'
output · 4 lines
unsubscribe calls: 1 call('topic-1')
after resub: 1
after close: 1
OK
step 9
$ cd /repo && python -m pytest tests/mqtt/test_roborock_session.py -q 2>&1 | tail -5
output · 4 lines
self._warn_or_fail_if_strict(f"Unknown config option: {key}\n")

-- Docs: https://docs.pytest.org/en/stable/how-to/capture-warnings.html
6 passed, 1 warning in 1.74s
step 10
$ cd /repo && python -m pytest tests/mqtt -q 2>&1 | tail -3 && git diff --stat
output · 4 lines
-- Docs: https://docs.pytest.org/en/stable/how-to/capture-warnings.html
9 passed, 1 warning in 1.77s
 roborock/mqtt/roborock_session.py | 42 +++++++++++++++++++++++++++++++++++++--
 1 file changed, 40 insertions(+), 2 deletions(-)
step 11
$ done: {"notes": "Added topic idle timeout handling to RoborockMqttSession: `topic_idle_timeout` timedelta param (default 60s), idle cleanup tasks per topic that unsubscribe via the MQTT client after the idle period (MqttError logged), cancelled on resubscribe and in close(). Unsubscribe callable removes only its callback; cleanup starts only when last callback removed. Subscribe failure removes the callback before raising. All existing mqtt tests pass."}
output · 1 lines
submitted.