Skip to content

Commit b49e9ec

Browse files
committed
Replace mqtt publisher with mqtt queue
1 parent 9eee345 commit b49e9ec

28 files changed

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

src/isar/modules.py

Lines changed: 5 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,6 @@
1616
from isar.storage.blob_storage import BlobStorage
1717
from isar.storage.local_storage import LocalStorage
1818
from isar.storage.uploader import Uploader
19-
from robot_interface.telemetry.mqtt_client import MqttPublisher
2019

2120

2221
class ApplicationContainer(containers.DeclarativeContainer):
@@ -34,16 +33,13 @@ class ApplicationContainer(containers.DeclarativeContainer):
3433
robot_utilities = providers.Singleton(RobotUtilities, robot=robot_interface)
3534

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

4238
# State machine
4339
state_machine = providers.Singleton(
4440
StateMachine,
4541
events=events,
46-
mqtt_publisher=mqtt_client,
42+
mqtt_queue=mqtt_queue,
4743
)
4844

4945
# API and controllers
@@ -64,7 +60,7 @@ class ApplicationContainer(containers.DeclarativeContainer):
6460
authenticator=authenticator,
6561
scheduling_controller=scheduling_controller,
6662
robot_controller=robot_controller,
67-
mqtt_publisher=mqtt_client,
63+
mqtt_queue=mqtt_queue,
6864
)
6965

7066
# Storage
@@ -82,14 +78,14 @@ class ApplicationContainer(containers.DeclarativeContainer):
8278
RobotService,
8379
events=events,
8480
robot=robot_interface,
85-
mqtt_publisher=mqtt_client,
81+
mqtt_queue=mqtt_queue,
8682
)
8783

8884
# Uploader
8985
uploader = providers.Singleton(
9086
Uploader,
9187
storage_handlers=storage_handlers,
92-
mqtt_publisher=mqtt_client,
88+
mqtt_queue=mqtt_queue,
9389
)
9490

9591
# 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)