Branch was auto-updated.

This commit is contained in:
Bhavin Patel
2021-11-22 01:30:02 -08:00
committed by GitHub
2 changed files with 18 additions and 20 deletions
@@ -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:
+3 -8
View File
@@ -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