Skip to content

Commit 4e369c4

Browse files
committed
Simplify uploader and make tests deterministic
1 parent 77f6fa7 commit 4e369c4

12 files changed

Lines changed: 219 additions & 395 deletions

File tree

src/isar/models/events.py

Lines changed: 4 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -16,8 +16,6 @@
1616
from robot_interface.models.mission.task import InspectionTask
1717
from robot_interface.telemetry.mqtt_client import MQTTQueueType
1818

19-
InspectionQueueTuple = tuple[Inspection, Mission]
20-
2119

2220
class EmptyMessage:
2321
def __str__(self) -> str:
@@ -28,8 +26,8 @@ def __str__(self) -> str:
2826

2927

3028
class Event[T](Queue[T]):
31-
def __init__(self, name: str) -> None:
32-
super().__init__(maxsize=1)
29+
def __init__(self, name: str, maxsize: int = 1) -> None:
30+
super().__init__(maxsize=maxsize)
3331
self.name = name
3432

3533
def trigger_event(self, data: T, timeout: int | None = None) -> None:
@@ -90,8 +88,8 @@ def __init__(self) -> None:
9088
self.state_machine_events: StateMachineEvents = StateMachineEvents()
9189
self.robot_service_events: RobotServiceEvents = RobotServiceEvents()
9290

93-
self.upload_queue: Queue[InspectionQueueTuple] = Queue[InspectionQueueTuple](
94-
maxsize=10
91+
self.upload_event: Event[tuple[Inspection, Mission]] = Event(
92+
"uploader", maxsize=10
9593
)
9694

9795
self.mqtt_queue: Queue[MQTTQueueType] = Queue[MQTTQueueType]()

src/isar/modules.py

Lines changed: 5 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -85,22 +85,18 @@ class ApplicationContainer(containers.DeclarativeContainer):
8585
mqtt_publisher=mqtt_client,
8686
)
8787

88-
# Inspection data service
89-
inspection_service = providers.Singleton(
90-
RobotInspectionService,
91-
events=events,
92-
robot=robot_interface,
93-
mqtt_publisher=mqtt_client,
94-
)
95-
9688
# Uploader
9789
uploader = providers.Singleton(
9890
Uploader,
99-
events=events,
10091
storage_handlers=storage_handlers,
10192
mqtt_publisher=mqtt_client,
10293
)
10394

95+
# Inspection data service
96+
inspection_service = providers.Singleton(
97+
RobotInspectionService, events=events, robot=robot_interface, uploader=uploader
98+
)
99+
104100

105101
def get_injector() -> ApplicationContainer:
106102
container = ApplicationContainer()

src/isar/robot/robot_inspection_service.py

Lines changed: 41 additions & 36 deletions
Original file line numberDiff line numberDiff line change
@@ -1,18 +1,12 @@
11
import logging
22
from collections.abc import Callable
3-
from queue import Queue
43
from threading import Event as ThreadEvent
54
from threading import Thread
65

76
from isar.config.settings import settings
8-
from isar.models.events import (
9-
EventConflictError,
10-
Events,
11-
EventTimeoutError,
12-
RobotServiceEvents,
13-
StateMachineEvents,
14-
)
7+
from isar.models.events import Event, EventConflictError, Events, EventTimeoutError
158
from isar.robot.function_thread import FunctionThread
9+
from isar.storage.uploader import Uploader
1610
from robot_interface.models.exceptions.robot_exceptions import (
1711
RobotException,
1812
RobotRetrieveInspectionException,
@@ -21,18 +15,17 @@
2115
from robot_interface.models.mission.mission import Mission
2216
from robot_interface.models.mission.task import InspectionTask
2317
from robot_interface.robot_interface import RobotInterface
24-
from robot_interface.telemetry.mqtt_client import MqttClientInterface
2518

2619

27-
def robot_upload_inspection(
28-
robot: RobotInterface,
20+
def fetch_and_upload_inspection(
21+
get_inspection_function: Callable[[InspectionTask], Inspection],
2922
logger: logging.Logger,
23+
upload_function: Callable[[Inspection, Mission], None],
3024
task: InspectionTask,
3125
mission: Mission,
32-
upload_queue: Queue,
3326
) -> None:
3427
try:
35-
inspection: Inspection = robot.get_inspection(task=task)
28+
inspection: Inspection = get_inspection_function(task)
3629
if task.id != inspection.id:
3730
logger.warning(
3831
f"The id of task ({task.id}) "
@@ -45,33 +38,30 @@ def robot_upload_inspection(
4538
return
4639

4740
if not inspection:
48-
logger.warning(
49-
f"No inspection result data retrieved for task {str(task.id)[:8]}"
50-
)
41+
logger.error(f"No inspection result data retrieved for task {str(task.id)[:8]}")
42+
return
5143

5244
inspection.metadata.tag_id = task.tag_id
5345
inspection.metadata.analysis_types = task.analysis_types
5446

55-
message: tuple[Inspection, Mission] = (
56-
inspection,
57-
mission,
58-
)
59-
upload_queue.put(message)
60-
logger.info(f"Inspection result: {str(inspection.id)[:8]} queued for upload")
47+
upload_function(inspection, mission)
6148

6249

6350
class RobotInspectionService:
6451
def __init__(
6552
self,
6653
events: Events,
6754
robot: RobotInterface,
68-
mqtt_publisher: MqttClientInterface,
55+
uploader: Uploader,
6956
) -> None:
7057
self.logger = logging.getLogger("uploader")
71-
self.state_machine_events: StateMachineEvents = events.state_machine_events
72-
self.robot_service_events: RobotServiceEvents = events.robot_service_events
73-
self.mqtt_publisher: MqttClientInterface = mqtt_publisher
74-
self.upload_queue: Queue = events.upload_queue
58+
self.upload_task_event: Event[tuple[InspectionTask, Mission]] = (
59+
events.robot_service_events.request_inspection_upload
60+
)
61+
self.upload_inspection_event: Event[tuple[Inspection, Mission]] = (
62+
events.upload_event
63+
)
64+
self.uploader: Uploader = uploader
7565
self.robot: RobotInterface = robot
7666
self.upload_inspection_threads: list[FunctionThread] = []
7767
self.signal_exit: ThreadEvent = ThreadEvent()
@@ -107,7 +97,7 @@ def _restart_inspection_thread_if_stopped(self) -> None:
10797

10898
def register_and_monitor_inspection_callback(
10999
self,
110-
callback_function: Callable,
100+
callback_function: Callable[[Inspection, Mission], None],
111101
) -> None:
112102
self.inspection_callback_function = callback_function
113103

@@ -121,18 +111,33 @@ def register_and_monitor_inspection_callback(
121111
def run(self) -> None:
122112
try:
123113
while not self.signal_exit.wait(0):
124-
upload_request: tuple[(InspectionTask, Mission)] | None = (
125-
self.robot_service_events.request_inspection_upload.consume_event()
114+
115+
upload_task_request: tuple[(InspectionTask, Mission)] | None = (
116+
self.upload_task_event.consume_event()
126117
)
127-
if upload_request is not None:
118+
119+
if upload_task_request is not None:
128120
self.upload_inspection_threads.append(
129121
FunctionThread(
130-
robot_upload_inspection,
131-
self.robot,
122+
fetch_and_upload_inspection,
123+
self.robot.get_inspection,
132124
self.logger,
133-
upload_request[0],
134-
upload_request[1],
135-
self.upload_queue,
125+
self.uploader.upload_inspection,
126+
upload_task_request[0],
127+
upload_task_request[1],
128+
)
129+
)
130+
131+
upload_inspection_request: tuple[Inspection, Mission] | None = (
132+
self.upload_inspection_event.consume_event()
133+
)
134+
135+
if upload_inspection_request is not None:
136+
self.upload_inspection_threads.append(
137+
FunctionThread(
138+
self.uploader.upload_inspection,
139+
upload_inspection_request[0],
140+
upload_inspection_request[1],
136141
)
137142
)
138143

src/isar/script.py

Lines changed: 2 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,7 @@
1010
from isar.config.log import setup_loggers
1111
from isar.config.open_telemetry import instrument_fastapi, setup_open_telemetry
1212
from isar.config.settings import robot_settings, settings
13-
from isar.models.events import Events, InspectionQueueTuple
13+
from isar.models.events import Events
1414
from isar.modules import ApplicationContainer, get_injector
1515
from isar.robot.robot_inspection_service import RobotInspectionService
1616
from isar.robot.robot_service import RobotService
@@ -22,7 +22,6 @@
2222
RobotInfoPublisher,
2323
)
2424
from isar.state_machine.state_machine import StateMachine
25-
from isar.storage.uploader import Uploader
2625
from robot_interface.models.inspection.inspection import Inspection
2726
from robot_interface.models.mission.mission import Mission
2827
from robot_interface.robot_interface import RobotInterface
@@ -83,7 +82,6 @@ def start() -> None:
8382
print_startup_info()
8483

8584
state_machine: StateMachine = injector.state_machine()
86-
uploader: Uploader = injector.uploader()
8785
robot_interface: RobotInterface = injector.robot_interface()
8886
events: Events = injector.events()
8987
robot: RobotService = injector.robot()
@@ -110,12 +108,6 @@ def start() -> None:
110108
inspection_service_thread.start()
111109
threads.append(inspection_service_thread)
112110

113-
uploader_thread: Thread = Thread(
114-
target=uploader.run, name="ISAR Uploader", daemon=True
115-
)
116-
uploader_thread.start()
117-
threads.append(uploader_thread)
118-
119111
robot_service_thread: Thread = Thread(
120112
target=robot.run, name="Robot service", daemon=True
121113
)
@@ -125,11 +117,7 @@ def start() -> None:
125117
if settings.UPLOAD_INSPECTIONS_ASYNC:
126118

127119
def inspections_callback(inspection: Inspection, mission: Mission) -> None:
128-
message: InspectionQueueTuple = (
129-
inspection,
130-
mission,
131-
)
132-
state_machine.events.upload_queue.put(message)
120+
state_machine.events.upload_event.trigger_event((inspection, mission))
133121

134122
inspection_service.register_and_monitor_inspection_callback(
135123
inspections_callback

src/isar/services/utilities/robot_utilities.py

Lines changed: 0 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,3 @@
1-
import logging
2-
31
from robot_interface.models.robots.media import MediaConfig
42
from robot_interface.robot_interface import RobotInterface
53

@@ -14,7 +12,6 @@ def __init__(
1412
robot: RobotInterface,
1513
):
1614
self.robot: RobotInterface = robot
17-
self.logger = logging.getLogger("api")
1815

1916
def generate_media_config(self) -> MediaConfig | None:
2017
return self.robot.generate_media_config()

0 commit comments

Comments
 (0)