diff --git a/bin/ssa-end-to-end-testing/modules/test_ssa_detections.py b/bin/ssa-end-to-end-testing/modules/test_ssa_detections.py index c9042cc1c5..1c7d8992e3 100644 --- a/bin/ssa-end-to-end-testing/modules/test_ssa_detections.py +++ b/bin/ssa-end-to-end-testing/modules/test_ssa_detections.py @@ -2,6 +2,7 @@ import logging import os import time import sys +import uuid from http import HTTPStatus from modules.streams_service_api_helper import DSPApi @@ -46,7 +47,8 @@ class SSADetectionTesting: test_results = [] for i in range(0, len(test_spls)): self.max_execution_time = MAX_EXECUTION_TIME_LIMIT - test_result = self.ssa_detection_test(read_spl(file_path_spl, test_spls[i]), file_path_data, test_names[i]) + test_id = str(uuid.uuid4()) + test_result = self.ssa_detection_test(read_spl(file_path_spl, test_spls[i]), file_path_data, test_names[i], test_id) test_results.append(test_result.copy()) passed = True @@ -65,9 +67,9 @@ class SSADetectionTesting: LOGGER.info('Test SSA Detection: ' + test_obj["detection_obj"]["name"]) self.max_execution_time = MAX_EXECUTION_TIME_LIMIT file_path_attack_data = test_obj["attack_data_file_path"] - + test_id = str(uuid.uuid4()) test_results = self.ssa_detection_test(test_obj["detection_obj"]["search"], file_path_attack_data, - "SSA Smoke Test " + test_obj["test_obj"]["name"], + "SSA Smoke Test " + test_obj["test_obj"]["name"], test_id, test_obj['test_obj']['tests'][0]['pass_condition']) return test_results @@ -89,7 +91,7 @@ class SSADetectionTesting: self.cleanup_old_pipelines() self.test_results["result"] = True self.test_results["msg"] = "" - self.results_index = self.api.create_temp_index("mc") + #self.results_index = self.api.create_temp_index("mc") self.created_pipelines = [] self.activated_pipelines = [] @@ -111,13 +113,13 @@ class SSADetectionTesting: else: LOGGER.warning("Found and deleted an old pipeline: %s", pipeline['name']) - def ssa_detection_test_main(self, spl, source, test_name, pass_condition): + def ssa_detection_test_main(self, spl, source, test_name, pass_condition, test_id): self.execution_passed = True self.wait_time(SLEEP_TIME_CREATE_INDEX) check_ssa_spl = check_source_sink(spl) - spl = manipulate_spl(self.api.env, spl, self.results_index) + spl = manipulate_spl(self.api.env, spl, test_id) assert spl is not None, "fail to manipulate spl file" upl = self.api.compile_spl(spl) @@ -160,8 +162,9 @@ class SSADetectionTesting: while not (search_results or max_execution_time_reached): max_execution_time_reached = self.wait_time(WAIT_CYCLE) - query = f"from indexes('{self.results_index['name']}') | search source!=\"Search Catalog\" " - sid = self.api.submit_search_job(self.results_index['module'], query) + query = f"from indexes('detection_testing') | search test_id=\"{test_id}\" " + LOGGER.info(f"Executing search query: {query}") + sid = self.api.submit_search_job('mc', query) assert sid is not None, f"Failed to create a Search Job" job_finished = False @@ -198,10 +201,10 @@ class SSADetectionTesting: """ deactivate_pipeline = lambda p: self.api.deactivate_pipeline(p)[0].status_code == HTTPStatus.OK delete_pipeline = lambda p: self.api.delete_pipeline(p).status_code == HTTPStatus.NO_CONTENT - delete_index = lambda p: self.api.delete_temp_index(p["id"]) == HTTPStatus.NO_CONTENT + #delete_index = lambda p: self.api.delete_temp_index(p["id"]) == HTTPStatus.NO_CONTENT self.activated_pipelines = [p for p in self.activated_pipelines if not deactivate_pipeline(p)] self.created_pipelines = [p for p in self.created_pipelines if not delete_pipeline(p)] - if len(self.activated_pipelines) > 0 or len(self.created_pipelines) > 0 or not delete_index(self.results_index): + if len(self.activated_pipelines) > 0 or len(self.created_pipelines) > 0: LOGGER.warning("Not all SCS resources freed up") LOGGER.info(f"Created Pipelines: {','.join(self.created_pipelines)}") LOGGER.info(f"Active Pipelines: {','.join(self.activated_pipelines)}") @@ -209,10 +212,10 @@ class SSADetectionTesting: else: LOGGER.info("Testing successfully cleaned up") - def ssa_detection_test(self, spl, source, test_name, pass_condition='@count_gt(0)'): + def ssa_detection_test(self, spl, source, test_name, test_id, pass_condition='@count_gt(0)'): self.ssa_detection_test_init() try: - test_result = self.ssa_detection_test_main(spl, source, test_name, pass_condition) + test_result = self.ssa_detection_test_main(spl, source, test_name, pass_condition, test_id) self.ssa_detection_test_teardown() return test_result except AssertionError as e: diff --git a/bin/ssa-end-to-end-testing/modules/utils.py b/bin/ssa-end-to-end-testing/modules/utils.py index 77dcb226b5..14111281e8 100644 --- a/bin/ssa-end-to-end-testing/modules/utils.py +++ b/bin/ssa-end-to-end-testing/modules/utils.py @@ -78,19 +78,14 @@ def check_source_sink(spl): return match_sink -def manipulate_spl(env, spl, results_index): +def manipulate_spl(env, spl, test_id): # Obtain the SSA source pulsar_source_connection_id, pulsar_source_topic = return_macros(env) source = READ_SSA_ENRICHED_EVENTS_EXPANDED\ .replace("__PULSAR_SOURCE_CONNECTION_ID__", pulsar_source_connection_id)\ .replace("__PULSAR_SOURCE_TOPIC__", pulsar_source_topic) # Obtain the test sink - if results_index is not None: - module = results_index["module"] - index = results_index["name"] - sink = f"index(\"{module}\", \"{index}\")" - else: - sink = "write_null()" + sink = f" eval test_id=\"{test_id}\" | into index(\"mc\", \"detection_testing\")" # Replace spl template with its `source` and `sink` spl = replace_ssa_macros(source, sink, spl) LOGGER.info(f"spl: {spl}") @@ -105,7 +100,7 @@ def read_spl(file_path, file_name): def replace_ssa_macros(source, sink, spl): spl = re.sub(r'read_ssa_enriched_events\(\s*\)', source, spl, flags=re.IGNORECASE) - spl = re.sub(r'write_ssa_detected_events\(\s*\)', sink, spl, flags=re.IGNORECASE) + spl = re.sub(r'into write_ssa_detected_events()\(\s*\)', sink, spl, flags=re.IGNORECASE) return spl