Skip to content

Commit 6c0d63f

Browse files
committed
Add analysis_types passthrough from Flotilla to SARA
Carry per-inspection analysis selections from Flotilla through ISAR to SARA without inspecting their semantics. The REST field 'analysisTypes' on each inspection definition is parsed into InspectionTask.analysis_types, copied onto InspectionMetadata, and re-emitted as 'required_analysis' on the MQTT inspection result payload that SARA consumes. The same list is also written to the metadata sidecar JSON uploaded alongside the blob. Backward compatible in both directions: missions without the field produce null on the wire, and SARA falls back to its existing default-analysis-by-file-extension behaviour.
1 parent 2c1b7d2 commit 6c0d63f

11 files changed

Lines changed: 133 additions & 16 deletions

File tree

src/isar/apis/models/start_mission_definition.py

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -40,6 +40,7 @@ class StartMissionInspectionDefinition(BaseModel):
4040
inspection_target: InputPosition
4141
inspection_description: str | None = None
4242
duration: float | None = None
43+
analysis_types: list[str] | None = Field(default=None)
4344

4445

4546
class StartMissionTaskDefinition(BaseModel):
@@ -123,6 +124,7 @@ def to_inspection_task(task_definition: StartMissionTaskDefinition) -> TASKS:
123124
inspection_description=task_definition.inspection.inspection_description,
124125
target=task_definition.inspection.inspection_target.to_alitra_position(),
125126
zoom=task_definition.zoom,
127+
analysis_types=inspection_definition.analysis_types,
126128
)
127129
elif inspection_definition.type == InspectionTypes.video:
128130
if inspection_definition.duration is None:
@@ -135,6 +137,7 @@ def to_inspection_task(task_definition: StartMissionTaskDefinition) -> TASKS:
135137
target=task_definition.inspection.inspection_target.to_alitra_position(),
136138
duration=inspection_definition.duration,
137139
zoom=task_definition.zoom,
140+
analysis_types=inspection_definition.analysis_types,
138141
)
139142
elif inspection_definition.type == InspectionTypes.thermal_image:
140143
return TakeThermalImage(
@@ -144,6 +147,7 @@ def to_inspection_task(task_definition: StartMissionTaskDefinition) -> TASKS:
144147
inspection_description=task_definition.inspection.inspection_description,
145148
target=task_definition.inspection.inspection_target.to_alitra_position(),
146149
zoom=task_definition.zoom,
150+
analysis_types=inspection_definition.analysis_types,
147151
)
148152
elif inspection_definition.type == InspectionTypes.thermal_video:
149153
if inspection_definition.duration is None:
@@ -156,6 +160,7 @@ def to_inspection_task(task_definition: StartMissionTaskDefinition) -> TASKS:
156160
target=task_definition.inspection.inspection_target.to_alitra_position(),
157161
duration=inspection_definition.duration,
158162
zoom=task_definition.zoom,
163+
analysis_types=inspection_definition.analysis_types,
159164
)
160165
elif inspection_definition.type == InspectionTypes.audio:
161166
if inspection_definition.duration is None:
@@ -167,13 +172,15 @@ def to_inspection_task(task_definition: StartMissionTaskDefinition) -> TASKS:
167172
inspection_description=task_definition.inspection.inspection_description,
168173
target=task_definition.inspection.inspection_target.to_alitra_position(),
169174
duration=inspection_definition.duration,
175+
analysis_types=inspection_definition.analysis_types,
170176
)
171177
elif inspection_definition.type == InspectionTypes.co2_measurement:
172178
return TakeCO2Measurement(
173179
id=task_definition.id if task_definition.id else str(uuid4()),
174180
robot_pose=task_definition.pose.to_alitra_pose(),
175181
tag_id=task_definition.tag,
176182
inspection_description=task_definition.inspection.inspection_description,
183+
analysis_types=inspection_definition.analysis_types,
177184
)
178185
else:
179186
raise ValueError(

src/isar/robot/robot_inspection_service.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -47,6 +47,7 @@ def robot_upload_inspection(
4747
)
4848

4949
inspection.metadata.tag_id = task.tag_id
50+
inspection.metadata.analysis_types = task.analysis_types
5051

5152
message: Tuple[Inspection, Mission] = (
5253
inspection,

src/isar/storage/uploader.py

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -259,9 +259,6 @@ def _publish_inspection_result(
259259
inspection: InspectionBlob,
260260
inspection_paths: StoragePaths[BlobStoragePath],
261261
) -> None:
262-
"""Publishes the reference of the inspection result to the MQTT Broker
263-
along with the analysis type
264-
"""
265262
if not self.mqtt_publisher:
266263
return
267264

@@ -275,6 +272,7 @@ def _publish_inspection_result(
275272
tag_id=inspection.metadata.tag_id,
276273
inspection_type=type(inspection).__name__,
277274
inspection_description=inspection.metadata.inspection_description,
275+
required_analysis=inspection.metadata.analysis_types,
278276
timestamp=inspection.metadata.start_time,
279277
)
280278
self.mqtt_publisher.publish(

src/isar/storage/utilities.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -43,6 +43,7 @@ def construct_metadata_file(
4343
"robot_name": settings.ROBOT_NAME,
4444
"inspection_description": inspection.metadata.inspection_description,
4545
"tag": inspection.metadata.tag_id,
46+
"analysis_types": inspection.metadata.analysis_types,
4647
"robot_pose": {
4748
"position": {
4849
"x": inspection.metadata.robot_pose.position.x,

src/robot_interface/models/inspection/inspection.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,7 @@ class InspectionMetadata(BaseModel):
1212
file_type: str
1313
tag_id: str | None = None
1414
inspection_description: str | None = None
15+
analysis_types: list[str] | None = None
1516

1617

1718
class ImageMetadata(InspectionMetadata):

src/robot_interface/models/mission/task.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -49,6 +49,7 @@ class InspectionTask(Task):
4949
robot_pose: Pose = Field()
5050
inspection_description: str | None = Field(default=None)
5151
zoom: ZoomDescription | None = Field(default=None)
52+
analysis_types: list[str] | None = Field(default=None)
5253

5354
@staticmethod
5455
def get_inspection_type() -> Type[Inspection]:

src/robot_interface/telemetry/payloads.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -114,6 +114,7 @@ class InspectionResultPayload(BaseModel):
114114
tag_id: str | None = None
115115
inspection_type: str | None = None
116116
inspection_description: str | None = None
117+
required_analysis: list[str] | None = None
117118
timestamp: datetime
118119

119120

tests/isar/apis/models/example_mission_definition.json

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@
2828
"z": 0,
2929
"frame_name": "robot"
3030
},
31+
"analysis_types": ["anonymize"],
3132
"id": "generated_inspection_id"
3233
},
3334
"tag": "MY-TAG-123",

tests/isar/apis/models/test_start_mission_definition.py

Lines changed: 40 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
import json
22
import os
33

4+
import pytest
45
from alitra import Frame, Orientation, Pose, Position
56

67
from isar.apis.models.models import InputOrientation, InputPose, InputPosition
@@ -13,7 +14,7 @@
1314
to_isar_mission,
1415
)
1516
from robot_interface.models.mission.mission import Mission
16-
from robot_interface.models.mission.task import TakeImage
17+
from robot_interface.models.mission.task import InspectionTask, TakeImage
1718

1819

1920
def test_to_isar_mission() -> None:
@@ -79,3 +80,41 @@ def test_mission_definition_from_json_to_isar_mission() -> None:
7980
)
8081
assert task.type == "take_image"
8182
assert task.target == Position(0.0, 0.0, 0.0, frame=Frame("robot"))
83+
assert task.analysis_types == ["anonymize"]
84+
85+
86+
def _build_mission_with_inspection_payload(inspection_payload: dict) -> Mission:
87+
payload = {
88+
"name": "mission",
89+
"tasks": [
90+
{
91+
"type": "inspection",
92+
"pose": {
93+
"position": {"x": 0, "y": 0, "z": 0},
94+
"orientation": {"x": 0, "y": 0, "z": 0, "w": 1},
95+
},
96+
"inspection": inspection_payload,
97+
}
98+
],
99+
}
100+
return to_isar_mission(StartMissionDefinition.model_validate(payload))
101+
102+
103+
@pytest.mark.parametrize(
104+
"inspection_type",
105+
[t.value for t in InspectionTypes],
106+
)
107+
def test_analysis_types_defaults_to_none_for_all_inspection_types(
108+
inspection_type: str,
109+
) -> None:
110+
payload: dict = {
111+
"type": inspection_type,
112+
"inspection_target": {"x": 0, "y": 0, "z": 0},
113+
}
114+
if inspection_type in {"Video", "ThermalVideo", "Audio"}:
115+
payload["duration"] = 1.0
116+
117+
mission = _build_mission_with_inspection_payload(payload)
118+
task = mission.tasks[0]
119+
assert isinstance(task, InspectionTask)
120+
assert task.analysis_types is None

tests/isar/storage/test_uploader.py

Lines changed: 72 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
1+
import json
12
import time
23
from datetime import datetime
3-
from typing import Tuple
4+
from typing import Callable, Tuple
45
from uuid import uuid4
56

67
from alitra import Frame, Orientation, Pose, Position
@@ -31,6 +32,17 @@
3132
)
3233

3334

35+
def _wait_until(
36+
predicate: Callable[[], bool], timeout: float = 5.0, interval: float = 0.01
37+
) -> None:
38+
deadline = time.monotonic() + timeout
39+
while time.monotonic() < deadline:
40+
if predicate():
41+
return
42+
time.sleep(interval)
43+
raise AssertionError(f"Timed out after {timeout}s waiting for predicate")
44+
45+
3446
def test_should_upload_from_queue(
3547
container: ApplicationContainer, uploader_thread: UploaderThreadMock
3648
) -> None:
@@ -65,12 +77,10 @@ def test_should_upload_from_queue(
6577
)
6678

6779
uploader: Uploader = container.uploader()
80+
storage_handler: StorageFake = uploader.storage_handlers[0] # type: ignore
6881

6982
uploader.upload_queue.put(message)
70-
time.sleep(0.01)
71-
72-
storage_handler: StorageFake = uploader.storage_handlers[0] # type: ignore
73-
assert inspection in storage_handler.stored_inspections
83+
_wait_until(lambda: inspection in storage_handler.stored_inspections)
7484

7585

7686
def test_should_retry_failed_upload_from_queue(
@@ -93,15 +103,15 @@ def test_should_retry_failed_upload_from_queue(
93103
# Need it to fail so that it retries
94104
storage_handler.will_fail = True
95105
uploader.upload_queue.put(message)
106+
# Fixed wait: asserting absence of upload requires a real elapsed window
96107
time.sleep(1)
97108

98109
# Should not upload, instead raise StorageException
99110
assert not storage_handler.blob_exists(inspection)
100111
storage_handler.will_fail = False
101-
time.sleep(3)
102112

103-
# After some time, it should have retried and now it should be successful
104-
assert storage_handler.blob_exists(inspection)
113+
# Retry succeeds once exponential backoff elapses
114+
_wait_until(lambda: storage_handler.blob_exists(inspection), timeout=5.0)
105115

106116

107117
def test_should_not_publish_when_blob_paths_are_empty(
@@ -127,8 +137,59 @@ def test_should_not_publish_when_blob_paths_are_empty(
127137
mission,
128138
)
129139
uploader.upload_queue.put(message)
130-
time.sleep(1)
131-
132-
assert inspection in storage_handler.stored
140+
_wait_until(lambda: inspection in storage_handler.stored)
133141

142+
# Brief quiet period to confirm no MQTT publish follows the store
143+
time.sleep(0.1)
134144
assert len(mqtt_fake.published) == 0
145+
146+
147+
def _put_inspection_with_analysis_types(
148+
container: ApplicationContainer,
149+
analysis_types: list[str] | None,
150+
) -> MqttPublisherFake:
151+
metadata = ImageMetadata(
152+
start_time=datetime.now(),
153+
robot_pose=Pose(
154+
Position(0, 0, 0, Frame("asset")),
155+
Orientation(x=0, y=0, z=0, w=1, frame=Frame("asset")),
156+
Frame("asset"),
157+
),
158+
target_position=Position(0, 0, 0, Frame("asset")),
159+
file_type="jpg",
160+
analysis_types=analysis_types,
161+
)
162+
inspection = InspectionBlob(metadata=metadata, id=str(uuid4()))
163+
mission = Mission(name="m")
164+
165+
uploader: Uploader = container.uploader()
166+
mqtt_fake = MqttPublisherFake()
167+
uploader.mqtt_publisher = mqtt_fake
168+
169+
message: Tuple[Inspection, Mission] = (inspection, mission)
170+
uploader.upload_queue.put(message)
171+
return mqtt_fake
172+
173+
174+
def test_publishes_required_analysis_when_present(
175+
container: ApplicationContainer, uploader_thread: UploaderThreadMock
176+
) -> None:
177+
uploader_thread.start()
178+
mqtt_fake = _put_inspection_with_analysis_types(
179+
container, ["anonymize", "thermal-reading"]
180+
)
181+
_wait_until(lambda: mqtt_fake.count() == 1)
182+
183+
payload = json.loads(mqtt_fake.last()["payload"])
184+
assert payload["required_analysis"] == ["anonymize", "thermal-reading"]
185+
186+
187+
def test_publishes_null_required_analysis_when_absent(
188+
container: ApplicationContainer, uploader_thread: UploaderThreadMock
189+
) -> None:
190+
uploader_thread.start()
191+
mqtt_fake = _put_inspection_with_analysis_types(container, None)
192+
_wait_until(lambda: mqtt_fake.count() == 1)
193+
194+
payload = json.loads(mqtt_fake.last()["payload"])
195+
assert payload["required_analysis"] is None

0 commit comments

Comments
 (0)