Skip to content

Commit 1c4109a

Browse files
committed
Replace mqtt tuple with dataclass
1 parent a453d6a commit 1c4109a

9 files changed

Lines changed: 55 additions & 48 deletions

File tree

src/isar/models/events.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -14,7 +14,7 @@
1414
from robot_interface.models.mission.mission import Mission
1515
from robot_interface.models.mission.status import RobotStatus
1616
from robot_interface.models.mission.task import InspectionTask
17-
from robot_interface.telemetry.mqtt_client import MQTTQueueType
17+
from robot_interface.telemetry.mqtt_client import MQTTQueueMessage
1818

1919

2020
class EmptyMessage:
@@ -92,7 +92,7 @@ def __init__(self) -> None:
9292
"uploader", maxsize=10
9393
)
9494

95-
self.mqtt_queue: Queue[MQTTQueueType] = Queue[MQTTQueueType]()
95+
self.mqtt_queue: Queue[MQTTQueueMessage] = Queue[MQTTQueueMessage]()
9696

9797
self.state: Event[States] = Event("state")
9898

src/isar/services/service_connections/mqtt/mqtt_client.py

Lines changed: 9 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -14,7 +14,7 @@
1414
from paho.mqtt.reasoncodes import ReasonCode
1515

1616
from isar.config.settings import settings
17-
from robot_interface.telemetry.mqtt_client import MqttClientInterface
17+
from robot_interface.telemetry.mqtt_client import MqttClientInterface, MQTTQueueMessage
1818

1919

2020
def props_expiry(seconds: int) -> Properties:
@@ -44,10 +44,10 @@ def _on_giveup(data: Details) -> None:
4444

4545

4646
class MqttClient(MqttClientInterface):
47-
def __init__(self, mqtt_queue: Queue) -> None:
47+
def __init__(self, mqtt_queue: Queue[MQTTQueueMessage]) -> None:
4848
self.logger = logging.getLogger("mqtt_client")
4949
self.logger.setLevel("INFO")
50-
self.mqtt_queue: Queue = mqtt_queue
50+
self.mqtt_queue: Queue[MQTTQueueMessage] = mqtt_queue
5151

5252
username: str = settings.MQTT_USERNAME
5353
password: str = ""
@@ -93,23 +93,16 @@ def run(self) -> None:
9393
time.sleep(0) # avoid CPU spin
9494
continue
9595
try:
96-
item: tuple[str, str, int, bool, Properties | None] = (
97-
self.mqtt_queue.get(timeout=1)
98-
)
99-
if len(item) == 4:
100-
topic, payload, qos, retain = item
101-
properties = None
102-
else:
103-
topic, payload, qos, retain, properties = item
96+
item: MQTTQueueMessage = self.mqtt_queue.get(timeout=1)
10497
except Empty:
10598
continue
10699

107100
self.publish(
108-
topic=topic,
109-
payload=payload,
110-
qos=qos,
111-
retain=retain,
112-
properties=properties,
101+
topic=item.topic,
102+
payload=item.payload,
103+
qos=item.qos,
104+
retain=item.retain,
105+
properties=item.properties,
113106
)
114107

115108
def on_connect(

src/isar/services/service_connections/mqtt/robot_heartbeat_publisher.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -4,12 +4,12 @@
44

55
from isar.config.settings import settings
66
from isar.services.service_connections.mqtt.mqtt_client import props_expiry
7-
from robot_interface.telemetry.mqtt_client import MqttPublisher
7+
from robot_interface.telemetry.mqtt_client import MqttPublisher, MQTTQueueMessage
88
from robot_interface.telemetry.payloads import RobotHeartbeatPayload
99

1010

1111
class RobotHeartbeatPublisher:
12-
def __init__(self, mqtt_queue: Queue):
12+
def __init__(self, mqtt_queue: Queue[MQTTQueueMessage]):
1313
self.mqtt_publisher: MqttPublisher = MqttPublisher(mqtt_queue=mqtt_queue)
1414

1515
def run(self) -> None:

src/isar/services/service_connections/mqtt/robot_info_publisher.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -3,12 +3,12 @@
33
from queue import Queue
44

55
from isar.config.settings import robot_settings, settings
6-
from robot_interface.telemetry.mqtt_client import MqttPublisher
6+
from robot_interface.telemetry.mqtt_client import MqttPublisher, MQTTQueueMessage
77
from robot_interface.telemetry.payloads import RobotInfoPayload
88

99

1010
class RobotInfoPublisher:
11-
def __init__(self, mqtt_queue: Queue):
11+
def __init__(self, mqtt_queue: Queue[MQTTQueueMessage]):
1212
self.mqtt_publisher: MqttPublisher = MqttPublisher(mqtt_queue=mqtt_queue)
1313

1414
def run(self) -> None:

src/isar/services/utilities/mqtt_utilities.py

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,7 @@
1010
from robot_interface.telemetry.mqtt_client import (
1111
MqttClientInterface,
1212
MqttPublisher,
13-
MQTTQueueType,
13+
MQTTQueueMessage,
1414
)
1515
from robot_interface.telemetry.payloads import (
1616
InterventionNeededPayload,
@@ -50,7 +50,7 @@ def publish_task_status(
5050

5151

5252
def publish_mission_status(
53-
mqtt_queue: Queue[MQTTQueueType],
53+
mqtt_queue: Queue[MQTTQueueMessage],
5454
mission_id: str,
5555
mission_status: MissionStatus,
5656
error_message: ErrorMessage | None,
@@ -97,7 +97,7 @@ def publish_isar_status(
9797

9898

9999
def publish_mission_aborted(
100-
mqtt_queue: Queue[MQTTQueueType], current_mission_id: str | None, reason: str
100+
mqtt_queue: Queue[MQTTQueueMessage], current_mission_id: str | None, reason: str
101101
) -> None:
102102
mqtt_publisher: MqttPublisher = MqttPublisher(mqtt_queue=mqtt_queue)
103103

@@ -119,7 +119,7 @@ def publish_mission_aborted(
119119

120120

121121
def publish_intervention_needed(
122-
mqtt_queue: Queue[MQTTQueueType], error_message: str
122+
mqtt_queue: Queue[MQTTQueueMessage], error_message: str
123123
) -> None:
124124
"""Publishes the intervention needed message to the MQTT Broker"""
125125
mqtt_publisher: MqttPublisher = MqttPublisher(mqtt_queue=mqtt_queue)

src/robot_interface/robot_interface.py

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -8,7 +8,10 @@
88
from robot_interface.models.mission.status import MissionStatus, RobotStatus, TaskStatus
99
from robot_interface.models.mission.task import InspectionTask
1010
from robot_interface.models.robots.media import MediaConfig
11-
from robot_interface.telemetry.mqtt_client import MqttTelemetryPublisher
11+
from robot_interface.telemetry.mqtt_client import (
12+
MQTTQueueMessage,
13+
MqttTelemetryPublisher,
14+
)
1215

1316

1417
class RobotInterface(metaclass=ABCMeta):
@@ -216,7 +219,7 @@ def generate_media_config(self) -> MediaConfig | None:
216219

217220
@abstractmethod
218221
def get_telemetry_publishers(
219-
self, queue: Queue, isar_id: str, robot_name: str
222+
self, queue: Queue[MQTTQueueMessage], isar_id: str, robot_name: str
220223
) -> list[MqttTelemetryPublisher]:
221224
"""
222225
Set up telemetry publisher threads to publish regular updates for pose, battery

src/robot_interface/telemetry/mqtt_client.py

Lines changed: 25 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@
33
import time
44
from abc import ABCMeta, abstractmethod
55
from collections.abc import Callable
6+
from dataclasses import dataclass
67
from datetime import UTC, datetime
78
from logging import Logger
89
from queue import Queue
@@ -19,7 +20,14 @@
1920
)
2021
from robot_interface.telemetry.payloads import CloudHealthPayload
2122

22-
MQTTQueueType = tuple[str, str, int, bool, Properties | None]
23+
24+
@dataclass
25+
class MQTTQueueMessage:
26+
topic: str
27+
payload: str
28+
qos: int
29+
retain: bool
30+
properties: Properties | None
2331

2432

2533
def props_expiry(seconds: int) -> Properties:
@@ -56,8 +64,8 @@ def publish(
5664

5765

5866
class MqttPublisher(MqttClientInterface):
59-
def __init__(self, mqtt_queue: Queue[MQTTQueueType]) -> None:
60-
self.mqtt_queue: Queue[MQTTQueueType] = mqtt_queue
67+
def __init__(self, mqtt_queue: Queue[MQTTQueueMessage]) -> None:
68+
self.mqtt_queue: Queue[MQTTQueueMessage] = mqtt_queue
6169

6270
def publish(
6371
self,
@@ -67,12 +75,12 @@ def publish(
6775
retain: bool = False,
6876
properties: Properties | None = None,
6977
) -> None:
70-
queue_message: tuple[str, str, int, bool, Properties | None] = (
71-
topic,
72-
payload,
73-
qos,
74-
retain,
75-
properties,
78+
queue_message: MQTTQueueMessage = MQTTQueueMessage(
79+
topic=topic,
80+
payload=payload,
81+
qos=qos,
82+
retain=retain,
83+
properties=properties,
7684
)
7785
self.mqtt_queue.put(queue_message)
7886

@@ -81,13 +89,13 @@ class MqttTelemetryPublisher(Thread):
8189
def __init__(
8290
self,
8391
name: str,
84-
mqtt_queue: Queue[MQTTQueueType],
92+
mqtt_queue: Queue[MQTTQueueMessage],
8593
telemetry_method: Callable,
8694
topic: str,
8795
interval: float,
8896
should_expire: bool,
8997
) -> None:
90-
self.mqtt_queue: Queue[MQTTQueueType] = mqtt_queue
98+
self.mqtt_queue: Queue[MQTTQueueMessage] = mqtt_queue
9199
self.telemetry_method: Callable = telemetry_method
92100
self.topic: str = topic
93101
self.interval: float = interval
@@ -130,11 +138,11 @@ def run(self) -> None:
130138
if self.should_expire:
131139
properties = props_expiry(settings.MQTT_TELEMETRY_EXPIRY)
132140

133-
queue_message: MQTTQueueType = (
134-
topic,
135-
payload,
136-
0,
137-
False,
138-
properties,
141+
queue_message: MQTTQueueMessage = MQTTQueueMessage(
142+
topic=topic,
143+
payload=payload,
144+
qos=0,
145+
retain=False,
146+
properties=properties,
139147
)
140148
self.mqtt_queue.put(queue_message)

tests/isar/state_machine/states/test_going_to_lockdown_state.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -60,7 +60,7 @@ def test_stopping_lockdown_transitions_to_going_to_lockdown(events: Events) -> N
6060
assert not events.mqtt_queue.empty()
6161
mqtt_message = events.mqtt_queue.get(block=False)
6262
assert mqtt_message is not None
63-
mqtt_payload_topic = mqtt_message[0]
63+
mqtt_payload_topic = mqtt_message.topic
6464
assert mqtt_payload_topic is settings.TOPIC_ISAR_MISSION_ABORTED
6565

6666

tests/test_mocks/robot_interface.py

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,10 @@
1717
from robot_interface.models.mission.task import InspectionTask
1818
from robot_interface.models.robots.media import MediaConfig, MediaConnectionType
1919
from robot_interface.robot_interface import RobotInterface
20-
from robot_interface.telemetry.mqtt_client import MqttTelemetryPublisher
20+
from robot_interface.telemetry.mqtt_client import (
21+
MQTTQueueMessage,
22+
MqttTelemetryPublisher,
23+
)
2124
from tests.test_mocks.inspection import stub_image_metadata
2225

2326
_ROBOT_FRAME = Frame(name="robot")
@@ -85,7 +88,7 @@ def register_inspection_callback(
8588
return
8689

8790
def get_telemetry_publishers(
88-
self, queue: Queue, isar_id: str, robot_name: str
91+
self, queue: Queue[MQTTQueueMessage], isar_id: str, robot_name: str
8992
) -> list[MqttTelemetryPublisher]:
9093
return []
9194

0 commit comments

Comments
 (0)