diff --git a/orchestrator/activities/formatters.py b/orchestrator/activities/formatters.py index 9bb17ff..bf6f0cf 100644 --- a/orchestrator/activities/formatters.py +++ b/orchestrator/activities/formatters.py @@ -336,12 +336,15 @@ class Formatters(BaseActivity): attachment_content=json.dumps(attachment, indent=4, sort_keys=True) ) - def parse_report_schedule(self, input_data: dict[str, Any]) -> tuple[list[str], list[str]]: + def parse_report_schedule(self, input_data: dict[str, Any]) -> tuple[list[str], dict[str, Any]]: success_keys = [f"{value['namespace']}/{value['schedule_name']}" for value in input_data if value['success']] - error_keys = [f"{value['namespace']}/{value['schedule_name']}: {value['message']}" - for value in input_data if not value['success']] + error_keys = {f"{value['namespace']}/{value['schedule_name']}": { + 'message': value['message'], + 'attachment': value.get('attachment', None) + } + for value in input_data if not value['success']} return success_keys, error_keys @@ -384,15 +387,24 @@ class Formatters(BaseActivity): self.send_success_report( metadata=metadata, message=f"Created schedules: \n {', '.join(success_keys)}", - notification_id="REPORT_ORCHESTRATION_CREATED_SCHEDULES" + notification_id="REPORT_ORCHESTRATION_CREATED_SCHEDULES", + attachment=created_schedules ) if len(error_keys) > 0: + attachment = [] + for value in error_keys.values(): + if value['attachment'] is not None: + attachment.append( + f"{value['message']}\n{value['attachment']}") + else: + attachment.append(value['message']) + self.send_error_report( metadata=metadata, message=f"Failed to create schedules: \n {', '.join(error_keys)}", notification_id="REPORT_ORCHESTRATION_CREATED_SCHEDULES_ERROR", - attachment=created_schedules + attachment="\n ========== \n".join(attachment) ) # Send report for updated schedules @@ -404,15 +416,24 @@ class Formatters(BaseActivity): self.send_success_report( metadata=metadata, message=f"Updated schedules: \n {', '.join(success_keys)}", - notification_id="REPORT_ORCHESTRATION_UPDATED_SCHEDULES" + notification_id="REPORT_ORCHESTRATION_UPDATED_SCHEDULES", + attachment=updated_schedules ) if len(error_keys) > 0: + attachment = [] + for value in error_keys.values(): + if value['attachment'] is not None: + attachment.append( + f"{value['message']}\n{value['attachment']}") + else: + attachment.append(value['message']) + self.send_error_report( metadata=metadata, message=f"Failed to update schedules: \n {', '.join(error_keys)}", notification_id="REPORT_ORCHESTRATION_UPDATED_SCHEDULES_ERROR", - attachment=updated_schedules + attachment="\n ========== \n".join(attachment) ) if len(deleted_schedules) > 0: @@ -423,15 +444,24 @@ class Formatters(BaseActivity): self.send_success_report( metadata=metadata, message=f"Deleted schedules: \n {', '.join(success_keys)}", - notification_id="REPORT_ORCHESTRATION_DELETED_SCHEDULES" + notification_id="REPORT_ORCHESTRATION_DELETED_SCHEDULES", + attachment=deleted_schedules ) if len(error_keys) > 0: + attachment = [] + for value in error_keys.values(): + if value['attachment'] is not None: + attachment.append( + f"{value['message']}\n{value['attachment']}") + else: + attachment.append(value['message']) + self.send_error_report( metadata=metadata, message=f"Failed to delete schedules: \n {', '.join(error_keys)}", notification_id="REPORT_ORCHESTRATION_DELETED_SCHEDULES_ERROR", - attachment=deleted_schedules + attachment="\n ========== \n".join(attachment) ) @activity.defn(name="report_slot_orchestration") diff --git a/tests/orchestrator/activities/test_formatters.py b/tests/orchestrator/activities/test_formatters.py index f3c1397..c480997 100644 --- a/tests/orchestrator/activities/test_formatters.py +++ b/tests/orchestrator/activities/test_formatters.py @@ -527,14 +527,18 @@ def test_parse_report_schedule(formatters): "namespace": "test_namespace", "schedule_name": "test_schedule_name_to_create_error", "success": False, - "message": "test_error" + "message": "test_error", + "attachment": "test_attachment" } ] result = formatters.parse_report_schedule(input_data) assert result == ( ["test_namespace/test_schedule_name_to_create"], - ["test_namespace/test_schedule_name_to_create_error: test_error"] + {"test_namespace/test_schedule_name_to_create_error": { + 'message': 'test_error', + 'attachment': 'test_attachment' + }} ) @@ -558,7 +562,14 @@ async def test_report_schedule_orchestration(formatters): "namespace": "test_namespace", "schedule_name": "test_schedule_name_to_create_error", "success": False, - "message": "test_error" + "message": "test_error1", + "attachment": "test_attachment1" + }, + { + "namespace": "test_namespace", + "schedule_name": "test_schedule_name_to_create_error2", + "success": False, + "message": "test_error2" } ], "updated_schedules": [ @@ -571,7 +582,21 @@ async def test_report_schedule_orchestration(formatters): "namespace": "test_namespace", "schedule_name": "test_schedule_name_to_update_error", "success": False, - "message": "test_error" + "message": "test_error2", + "attachment": "test_attachment2" + }, + { + "namespace": "test_namespace", + "schedule_name": "test_schedule_name_to_update_error2", + "success": False, + "message": "test_error3", + "attachment": "test_attachment3" + }, + { + "namespace": "test_namespace", + "schedule_name": "test_schedule_name_to_update_error3", + "success": False, + "message": "test_error4" } ], "deleted_schedules": [ @@ -584,7 +609,14 @@ async def test_report_schedule_orchestration(formatters): "namespace": "test_namespace", "schedule_name": "test_schedule_name_to_delete_error", "success": False, - "message": "test_error" + "message": "test_error4", + "attachment": "test_attachment4" + }, + { + "namespace": "test_namespace", + "schedule_name": "test_schedule_name_to_delete_error2", + "success": False, + "message": "test_error5" } ] } @@ -600,39 +632,42 @@ async def test_report_schedule_orchestration(formatters): call( metadata=metadata['metadata'], message="Created schedules: \n test_namespace/test_schedule_name_to_create", - notification_id="REPORT_ORCHESTRATION_CREATED_SCHEDULES" - ), - call( - metadata=metadata['metadata'], - message="Updated schedules: \n test_namespace/test_schedule_name_to_update", - notification_id="REPORT_ORCHESTRATION_UPDATED_SCHEDULES" - ), - call( - metadata=metadata['metadata'], - message="Deleted schedules: \n test_namespace/test_schedule_name_to_delete", - notification_id="REPORT_ORCHESTRATION_DELETED_SCHEDULES" - ) - ]) - formatters.send_error_report.assert_has_calls([ - call( - metadata=metadata['metadata'], - message="Failed to create schedules: \n test_namespace/test_schedule_name_to_create_error: test_error", notification_id="REPORT_ORCHESTRATION_CREATED_SCHEDULES", attachment=input_data['created_schedules'] ), call( metadata=metadata['metadata'], - message="Failed to update schedules: \n test_namespace/test_schedule_name_to_update_error: test_error", + message="Updated schedules: \n test_namespace/test_schedule_name_to_update", notification_id="REPORT_ORCHESTRATION_UPDATED_SCHEDULES", attachment=input_data['updated_schedules'] ), call( metadata=metadata['metadata'], - message="Failed to delete schedules: \n test_namespace/test_schedule_name_to_delete_error: test_error", + message="Deleted schedules: \n test_namespace/test_schedule_name_to_delete", notification_id="REPORT_ORCHESTRATION_DELETED_SCHEDULES", attachment=input_data['deleted_schedules'] ) ]) + formatters.send_error_report.assert_has_calls([ + call( + metadata=metadata['metadata'], + message="Failed to create schedules: \n test_namespace/test_schedule_name_to_create_error, test_namespace/test_schedule_name_to_create_error2", + notification_id="REPORT_ORCHESTRATION_CREATED_SCHEDULES_ERROR", + attachment="test_error1\ntest_attachment1\n ========== \ntest_error2" + ), + call( + metadata=metadata['metadata'], + message="Failed to update schedules: \n test_namespace/test_schedule_name_to_update_error, test_namespace/test_schedule_name_to_update_error2, test_namespace/test_schedule_name_to_update_error3", + notification_id="REPORT_ORCHESTRATION_UPDATED_SCHEDULES_ERROR", + attachment="test_error2\ntest_attachment2\n ========== \ntest_error3\ntest_attachment3\n ========== \ntest_error4" + ), + call( + metadata=metadata['metadata'], + message="Failed to delete schedules: \n test_namespace/test_schedule_name_to_delete_error, test_namespace/test_schedule_name_to_delete_error2", + notification_id="REPORT_ORCHESTRATION_DELETED_SCHEDULES_ERROR", + attachment="test_error4\ntest_attachment4\n ========== \ntest_error5" + ) + ]) @mark.asyncio diff --git a/tests/orchestrator/activities/test_mongo_db.py b/tests/orchestrator/activities/test_mongo_db.py index e3fa686..742d42f 100644 --- a/tests/orchestrator/activities/test_mongo_db.py +++ b/tests/orchestrator/activities/test_mongo_db.py @@ -280,7 +280,7 @@ async def test_update_pipelines_timestamps_success(datetime_mock, mongo_db): {"schedule_name": "test2", "namespace": "test2"} ]}, {"$set": { - "updated_at": datetime_mock.now.return_value.strftime.return_value}} + "updated_at": datetime_mock.now.return_value}} ) @@ -325,9 +325,9 @@ async def test_create_pipelines_timestamps_success(datetime_mock, mongo_db): mongo_db.database["pipelines"].insert_many.assert_called_once_with( [ {"schedule_name": "test1", "namespace": "test1", - "updated_at": datetime_mock.now.return_value.strftime.return_value}, + "updated_at": datetime_mock.now.return_value}, {"schedule_name": "test2", "namespace": "test2", - "updated_at": datetime_mock.now.return_value.strftime.return_value} + "updated_at": datetime_mock.now.return_value} ] ) diff --git a/tests/orchestrator/activities/test_temporal_manager.py b/tests/orchestrator/activities/test_temporal_manager.py index 058aa70..dbb1258 100644 --- a/tests/orchestrator/activities/test_temporal_manager.py +++ b/tests/orchestrator/activities/test_temporal_manager.py @@ -213,21 +213,24 @@ async def test_create_schedule( input_data['schedules']['scouter']['test-schedule'], id="test-schedule", task_queue="test-workflow-queue", - execution_timeout=ANY + execution_timeout=ANY, + typed_search_attributes=mock_typed_search_attributes.return_value, ), call( "test-workflow", input_data['schedules']['scouter']['test-schedule-invalid-frequency'], id="test-schedule-invalid-frequency", task_queue="test-workflow-queue", - execution_timeout=ANY + execution_timeout=ANY, + typed_search_attributes=mock_typed_search_attributes.return_value, ), call( "test-workflow", input_data['schedules']['laborious']['test-schedule-laborious'], id="test-schedule-laborious", task_queue="test-workflow-queue", - execution_timeout=ANY + execution_timeout=ANY, + typed_search_attributes=mock_typed_search_attributes.return_value, ) ]) @@ -546,7 +549,8 @@ async def test_delete_schedules(temporal_manager): "schedule_name": "test-schedule_no_handler", "namespace": "scouter", "success": False, - "message": "Schedule test-schedule_no_handler not found" + "message": "Schedule test-schedule_no_handler not found", + "attachment": ANY }, { "schedule_name": "test-schedule-laborious",