Skip to content

Commit 62ba1ae

Browse files
committed
Replace mqtt publisher with mqtt queue
1 parent 7dce7b2 commit 62ba1ae

28 files changed

Lines changed: 273 additions & 360 deletions

src/isar/apis/api.py

Lines changed: 4 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,7 @@
1717
from isar.apis.schedule.scheduling_controller import SchedulingController
1818
from isar.apis.security.authentication import Authenticator
1919
from isar.config.settings import settings
20-
from robot_interface.telemetry.mqtt_client import MqttClientInterface
20+
from isar.models.mqtt_queue import MQTTQueue
2121
from robot_interface.telemetry.payloads import StartUpMessagePayload
2222

2323

@@ -27,15 +27,15 @@ def __init__(
2727
authenticator: Authenticator,
2828
scheduling_controller: SchedulingController,
2929
robot_controller: RobotController,
30-
mqtt_publisher: MqttClientInterface,
30+
mqtt_queue: MQTTQueue,
3131
port: int = settings.API_PORT,
3232
) -> None:
3333
self.authenticator: Authenticator = authenticator
3434
self.scheduling_controller: SchedulingController = scheduling_controller
3535
self.robot_controller: RobotController = robot_controller
3636
self.host: str = "0.0.0.0" # Locking uvicorn to use 0.0.0.0
3737
self.port: int = port
38-
self.mqtt_publisher: MqttClientInterface = mqtt_publisher
38+
self.mqtt_queue: MQTTQueue = mqtt_queue
3939

4040
self.logger: Logger = logging.getLogger("api")
4141

@@ -360,17 +360,14 @@ def _log_startup_message(self) -> None:
360360
)
361361

362362
def _publish_startup_message(self) -> None:
363-
if not self.mqtt_publisher:
364-
return
365-
366363
payload: StartUpMessagePayload = StartUpMessagePayload(
367364
isar_id=settings.ISAR_ID,
368365
timestamp=datetime.now(UTC),
369366
)
370367

371368
self.logger.info("Publishing startup message to MQTT broker")
372369

373-
self.mqtt_publisher.publish(
370+
self.mqtt_queue.publish(
374371
topic=settings.TOPIC_ISAR_STARTUP,
375372
payload=payload.model_dump_json(),
376373
qos=1,

src/isar/models/events.py

Lines changed: 6 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -8,13 +8,17 @@
88
MaintenanceResponse,
99
MissionStartResponse,
1010
)
11+
from isar.models.mqtt_queue import MQTTQueue
1112
from isar.state_machine.states_enum import States
13+
from robot_interface.models.exceptions.event_exceptions import (
14+
EventConflictError,
15+
EventTimeoutError,
16+
)
1217
from robot_interface.models.exceptions.robot_exceptions import ErrorMessage
1318
from robot_interface.models.inspection.inspection import Inspection
1419
from robot_interface.models.mission.mission import Mission
1520
from robot_interface.models.mission.status import RobotStatus
1621
from robot_interface.models.mission.task import InspectionTask
17-
from robot_interface.telemetry.mqtt_client import MQTTQueueMessage
1822

1923

2024
class EmptyMessage:
@@ -92,7 +96,7 @@ def __init__(self) -> None:
9296
"uploader", maxsize=10
9397
)
9498

95-
self.mqtt_queue: Queue[MQTTQueueMessage] = Queue[MQTTQueueMessage]()
99+
self.mqtt_queue: MQTTQueue = MQTTQueue(maxsize=10)
96100

97101
self.state: Event[States] = Event("state")
98102

@@ -192,11 +196,3 @@ def __init__(self) -> None:
192196
self.battery_above_recharge_threshold: Event[EmptyMessage] = Event(
193197
"battery_above_recharge_threshold"
194198
)
195-
196-
197-
class EventTimeoutError(Exception):
198-
pass
199-
200-
201-
class EventConflictError(Exception):
202-
pass

src/isar/models/mqtt_queue.py

Lines changed: 165 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,165 @@
1+
from datetime import UTC, datetime
2+
from queue import Empty, Queue
3+
4+
from isar.config.settings import settings
5+
from isar.models.status import IsarStatus
6+
from paho.mqtt.properties import Properties
7+
from robot_interface.models.exceptions.event_exceptions import EventTimeoutError
8+
from robot_interface.models.exceptions.robot_exceptions import ErrorMessage
9+
from robot_interface.models.mission.status import MissionStatus
10+
from robot_interface.models.mission.task import TASKS
11+
from robot_interface.telemetry.mqtt_client import (
12+
MQTTQueueMessage,
13+
props_expiry,
14+
)
15+
from robot_interface.telemetry.payloads import (
16+
InterventionNeededPayload,
17+
IsarStatusPayload,
18+
MissionAbortedPayload,
19+
MissionPayload,
20+
TaskPayload,
21+
)
22+
23+
24+
class MQTTQueue:
25+
queue: Queue[MQTTQueueMessage]
26+
27+
def __init__(self, maxsize: int = 10) -> None:
28+
self.queue = Queue(maxsize=maxsize)
29+
self.name = "MQTT queue"
30+
31+
def qsize(self) -> int:
32+
return self.queue.qsize()
33+
34+
def empty(self) -> bool:
35+
return self.queue.empty()
36+
37+
def publish(
38+
self,
39+
topic: str,
40+
payload: str,
41+
qos: int = 0,
42+
retain: bool = False,
43+
properties: Properties | None = None,
44+
) -> None:
45+
queue_message: MQTTQueueMessage = MQTTQueueMessage(
46+
topic=topic,
47+
payload=payload,
48+
qos=qos,
49+
retain=retain,
50+
properties=properties,
51+
)
52+
self.queue.put(queue_message)
53+
54+
def get(self, timeout: int = 0) -> MQTTQueueMessage:
55+
try:
56+
return self.queue.get(block=timeout is not None, timeout=timeout)
57+
except Empty:
58+
if timeout is not None:
59+
raise EventTimeoutError
60+
return None
61+
62+
def publish_task_status(self, task: TASKS, mission_id: str | None) -> None:
63+
"""Publishes the task status to the MQTT Broker"""
64+
65+
error_message: ErrorMessage | None = task.error_message
66+
67+
payload: TaskPayload = TaskPayload(
68+
isar_id=settings.ISAR_ID,
69+
robot_name=settings.ROBOT_NAME,
70+
mission_id=mission_id,
71+
task_id=task.id if task else None,
72+
status=task.status if task else None,
73+
task_type=task.type if task else None,
74+
error_reason=error_message.error_reason if error_message else None,
75+
error_description=(
76+
error_message.error_description if error_message else None
77+
),
78+
timestamp=datetime.now(UTC),
79+
)
80+
81+
self.publish(
82+
topic=settings.TOPIC_ISAR_TASK + f"/{task.id}",
83+
payload=payload.model_dump_json(),
84+
qos=1,
85+
retain=True,
86+
properties=props_expiry(settings.MQTT_MISSION_TASK_AND_STATUS_EXPIRY),
87+
)
88+
89+
def publish_mission_status(
90+
self,
91+
mission_id: str,
92+
mission_status: MissionStatus,
93+
error_message: ErrorMessage | None,
94+
) -> None:
95+
payload: MissionPayload = MissionPayload(
96+
isar_id=settings.ISAR_ID,
97+
robot_name=settings.ROBOT_NAME,
98+
mission_id=mission_id,
99+
status=mission_status,
100+
error_reason=error_message.error_reason if error_message else None,
101+
error_description=(
102+
error_message.error_description if error_message else None
103+
),
104+
timestamp=datetime.now(UTC),
105+
)
106+
107+
self.publish(
108+
topic=settings.TOPIC_ISAR_MISSION + f"/{mission_id}",
109+
payload=payload.model_dump_json(),
110+
qos=1,
111+
retain=True,
112+
properties=props_expiry(settings.MQTT_MISSION_TASK_AND_STATUS_EXPIRY),
113+
)
114+
115+
def publish_isar_status(self, status: IsarStatus) -> None:
116+
payload: IsarStatusPayload = IsarStatusPayload(
117+
isar_id=settings.ISAR_ID,
118+
robot_name=settings.ROBOT_NAME,
119+
status=status,
120+
timestamp=datetime.now(UTC),
121+
)
122+
123+
self.publish(
124+
topic=settings.TOPIC_ISAR_STATUS,
125+
payload=payload.model_dump_json(),
126+
qos=1,
127+
retain=True,
128+
properties=props_expiry(settings.MQTT_MISSION_TASK_AND_STATUS_EXPIRY),
129+
)
130+
131+
def publish_mission_aborted(
132+
self, current_mission_id: str | None, reason: str
133+
) -> None:
134+
payload: MissionAbortedPayload = MissionAbortedPayload(
135+
isar_id=settings.ISAR_ID,
136+
robot_name=settings.ROBOT_NAME,
137+
mission_id=current_mission_id,
138+
reason=reason,
139+
timestamp=datetime.now(UTC),
140+
)
141+
142+
self.publish(
143+
topic=settings.TOPIC_ISAR_MISSION_ABORTED,
144+
payload=payload.model_dump_json(),
145+
qos=1,
146+
retain=True,
147+
properties=props_expiry(settings.MQTT_MISSION_TASK_AND_STATUS_EXPIRY),
148+
)
149+
150+
def publish_intervention_needed(self, error_message: str) -> None:
151+
"""Publishes the intervention needed message to the MQTT Broker"""
152+
payload: InterventionNeededPayload = InterventionNeededPayload(
153+
isar_id=settings.ISAR_ID,
154+
robot_name=settings.ROBOT_NAME,
155+
reason=error_message,
156+
timestamp=datetime.now(UTC),
157+
)
158+
159+
self.publish(
160+
topic=settings.TOPIC_ISAR_INTERVENTION_NEEDED,
161+
payload=payload.model_dump_json(),
162+
qos=1,
163+
retain=True,
164+
properties=props_expiry(settings.MQTT_MISSION_TASK_AND_STATUS_EXPIRY),
165+
)

src/isar/modules.py

Lines changed: 6 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -10,13 +10,13 @@
1010
from isar.models.events import Events
1111
from isar.robot.robot_inspection_service import RobotInspectionService
1212
from isar.robot.robot_service import RobotService
13+
from isar.models.mqtt_queue import MQTTQueue
1314
from isar.services.utilities.robot_utilities import RobotUtilities
1415
from isar.services.utilities.scheduling_utilities import SchedulingUtilities
1516
from isar.state_machine.state_machine import StateMachine
1617
from isar.storage.blob_storage import BlobStorage
1718
from isar.storage.local_storage import LocalStorage
1819
from isar.storage.uploader import Uploader
19-
from robot_interface.telemetry.mqtt_client import MqttPublisher
2020

2121

2222
class ApplicationContainer(containers.DeclarativeContainer):
@@ -34,16 +34,13 @@ class ApplicationContainer(containers.DeclarativeContainer):
3434
robot_utilities = providers.Singleton(RobotUtilities, robot=robot_interface)
3535

3636
# Mqtt client
37-
mqtt_client = providers.Singleton(
38-
MqttPublisher,
39-
mqtt_queue=providers.Callable(events.provided.mqtt_queue),
40-
)
37+
mqtt_queue = providers.Object(events.provides().mqtt_queue)
4138

4239
# State machine
4340
state_machine = providers.Singleton(
4441
StateMachine,
4542
events=events,
46-
mqtt_publisher=mqtt_client,
43+
mqtt_queue=mqtt_queue,
4744
)
4845

4946
# API and controllers
@@ -64,7 +61,7 @@ class ApplicationContainer(containers.DeclarativeContainer):
6461
authenticator=authenticator,
6562
scheduling_controller=scheduling_controller,
6663
robot_controller=robot_controller,
67-
mqtt_publisher=mqtt_client,
64+
mqtt_queue=mqtt_queue,
6865
)
6966

7067
# Storage
@@ -82,14 +79,14 @@ class ApplicationContainer(containers.DeclarativeContainer):
8279
RobotService,
8380
events=events,
8481
robot=robot_interface,
85-
mqtt_publisher=mqtt_client,
82+
mqtt_queue=mqtt_queue,
8683
)
8784

8885
# Uploader
8986
uploader = providers.Singleton(
9087
Uploader,
9188
storage_handlers=storage_handlers,
92-
mqtt_publisher=mqtt_client,
89+
mqtt_queue=mqtt_queue,
9390
)
9491

9592
# Inspection data service

src/isar/robot/robot_monitor_mission.py

Lines changed: 8 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,7 @@
33
from collections.abc import Callable, Iterator
44

55
from isar.config.settings import settings
6-
from isar.services.utilities.mqtt_utilities import publish_task_status
6+
from isar.models.mqtt_queue import MQTTQueue
77
from robot_interface.models.exceptions.robot_exceptions import (
88
ErrorMessage,
99
RobotCommunicationException,
@@ -16,7 +16,6 @@
1616
from robot_interface.models.mission.status import MissionStatus, TaskStatus
1717
from robot_interface.models.mission.task import TASKS, InspectionTask
1818
from robot_interface.robot_interface import RobotInterface
19-
from robot_interface.telemetry.mqtt_client import MqttClientInterface
2019

2120

2221
def get_next_task(task_iterator: Iterator[TASKS]) -> TASKS | None:
@@ -148,7 +147,7 @@ async def get_and_report_task_status(
148147
current_task: TASKS,
149148
robot: RobotInterface,
150149
mission_id: str,
151-
mqtt_publisher: MqttClientInterface,
150+
mqtt_queue: MQTTQueue,
152151
) -> TaskStatus:
153152
logger = logging.getLogger("robot")
154153
new_task_status: TaskStatus | None = None
@@ -158,7 +157,7 @@ async def get_and_report_task_status(
158157
if current_task.status != new_task_status:
159158
current_task.status = new_task_status
160159
log_task_status(logger, current_task)
161-
publish_task_status(mqtt_publisher, current_task, mission_id)
160+
mqtt_queue.publish_task_status(current_task, mission_id)
162161
return new_task_status
163162

164163
except RobotTaskStatusException as e:
@@ -173,7 +172,7 @@ async def robot_monitor_mission(
173172
mission: Mission,
174173
robot: RobotInterface,
175174
request_inspection_upload: Callable[[InspectionTask], None],
176-
mqtt_publisher: MqttClientInterface,
175+
mqtt_queue: MQTTQueue,
177176
should_report_task_status: bool,
178177
) -> tuple[ErrorMessage | None, Mission, bool]:
179178
logger = logging.getLogger("robot")
@@ -227,19 +226,19 @@ async def robot_monitor_mission(
227226
if should_report_task_status and current_task:
228227
if mission_status == MissionStatus.Cancelled:
229228
current_task.status = TaskStatus.Cancelled
230-
publish_task_status(mqtt_publisher, current_task, mission.id)
229+
mqtt_queue.publish_task_status(current_task, mission.id)
231230
current_task = None
232231
continue
233232
if mission_status == MissionStatus.Failed:
234233
current_task.status = TaskStatus.Failed
235-
publish_task_status(mqtt_publisher, current_task, mission.id)
234+
mqtt_queue.publish_task_status(current_task, mission.id)
236235
current_task = None
237236
continue
238237
task_status = await get_and_report_task_status(
239238
current_task,
240239
robot,
241240
mission.id,
242-
mqtt_publisher,
241+
mqtt_queue,
243242
)
244243
if task_status is None:
245244
# Currently we only stop mission monitoring after failing to get mission status
@@ -253,7 +252,7 @@ async def robot_monitor_mission(
253252
except asyncio.CancelledError:
254253
if should_report_task_status and current_task is not None:
255254
current_task.status = TaskStatus.Cancelled
256-
publish_task_status(mqtt_publisher, current_task, mission.id)
255+
mqtt_queue.publish_task_status(current_task, mission.id)
257256
return None, mission, True
258257
finally:
259258
logger.info("Stopped monitoring mission")

0 commit comments

Comments
 (0)