Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 4 additions & 4 deletions endpoints/base
Original file line number Diff line number Diff line change
Expand Up @@ -758,7 +758,7 @@ function process_bench_roadblocks() {

echo "Initializing data structures:"
while read -u 9 line; do
iter_samp=`echo ${line} | awk '{print $1}'`
iter_samp="${line}"
iter_id=`echo ${iter_samp} | awk -F- '{print $1}'`
samp_id=`echo ${iter_samp} | awk -F- '{print $2}'`
iter_array_idx=${total_tests}
Expand All @@ -772,7 +772,7 @@ function process_bench_roadblocks() {
sample_data_attempt_fail[${iter_array_idx}]=0

(( total_tests += 1 ))
done 9< "${engine_bench_cmds_dir}/client/1/start"
done 9< <(xzcat "${engine_bench_cmds_dir}/client/1/start.json.xz" | jq -r '.[] | .test')

echo "Total tests: ${total_tests}"

Expand All @@ -788,7 +788,7 @@ function process_bench_roadblocks() {

(( current_test += 1 ))

iter_samp=`echo ${line} | awk '{print $1}'`
iter_samp="${line}"
iter_id=`echo ${iter_samp} | awk -F- '{print $1}'`
samp_id=`echo ${iter_samp} | awk -F- '{print $2}'`
let iter_array_idx=${current_test}-1
Expand Down Expand Up @@ -978,7 +978,7 @@ function process_bench_roadblocks() {
${max_sample_failures} \
${sample_result}
done
done 9< "$engine_bench_cmds_dir/client/1/start"
done 9< <(xzcat "$engine_bench_cmds_dir/client/1/start.json.xz" | jq -r '.[] | .test')
}

function process_roadblocks() {
Expand Down
42 changes: 22 additions & 20 deletions endpoints/endpoints.py
Original file line number Diff line number Diff line change
Expand Up @@ -1229,26 +1229,28 @@ def process_bench_roadblocks(callbacks = None, roadblock_id = None, endpoint_lab
iteration_sample_data = []

logger.info("Initializing data structures")
with open(engine_commands_dir + "/client/1/start") as bench_cmds_fp:
for line in bench_cmds_fp:
split = line.split(" ")
iteration_sample = split[0]
split = iteration_sample.split("-")
iteration_id = int(split[0])
sample_id = int(split[1])

logger.info("iteration_sample=%s iteration_id=%s sample_id=%s" % (iteration_sample, iteration_id, sample_id))

obj = {
"iteration-sample": iteration_sample,
"iteration-id": iteration_id,
"sample-id": sample_id,
"failures": 0,
"complete": False,
"attempt-num": 0,
"attempt-fail": 0
}
iteration_sample_data.append(obj)
bench_cmds, err = load_json_file(engine_commands_dir + "/client/1/start.json.xz", uselzma = True)
if bench_cmds is None:
logger.error("Failed to load bench commands from %s/client/1/start.json.xz: %s" % (engine_commands_dir, err))
return 1
for entry in bench_cmds:
iteration_sample = entry["test"]
split = iteration_sample.split("-")
iteration_id = int(split[0])
sample_id = int(split[1])

logger.info("iteration_sample=%s iteration_id=%s sample_id=%s" % (iteration_sample, iteration_id, sample_id))

obj = {
"iteration-sample": iteration_sample,
"iteration-id": iteration_id,
"sample-id": sample_id,
"failures": 0,
"complete": False,
"attempt-num": 0,
"attempt-fail": 0
}
iteration_sample_data.append(obj)

logger.info("Total tests: %d" % (len(iteration_sample_data)))

Expand Down
104 changes: 54 additions & 50 deletions engine/engine_lib.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
import logging
import os
import re
import shlex
import shutil
import sys
import tempfile
Expand Down Expand Up @@ -533,16 +534,16 @@ def get_data(self):
"bench-cmds",
self.cs_type,
self.cs_id,
"start",
"start.json.xz",
),
"bench-start-cmds",
"bench-start-cmds.json.xz",
)
else:
self.scp_from_controller(
os.path.join(
self.engine_config_dir, "bench-cmds", "client", "1", "start"
self.engine_config_dir, "bench-cmds", "client", "1", "start.json.xz"
),
"bench-start-cmds",
"bench-start-cmds.json.xz",
)

if self.cs_type == "client":
Expand All @@ -552,9 +553,9 @@ def get_data(self):
"bench-cmds",
self.cs_type,
self.cs_id,
"infra",
"infra.json.xz",
),
"bench-infra-cmds",
"bench-infra-cmds.json.xz",
)
if self.cs_id == "1":
self.scp_from_controller(
Expand All @@ -563,9 +564,9 @@ def get_data(self):
"bench-cmds",
self.cs_type,
self.cs_id,
"runtime",
"runtime.json.xz",
),
"bench-runtime-cmds",
"bench-runtime-cmds.json.xz",
)
elif self.cs_type == "server":
self.scp_from_controller(
Expand All @@ -574,9 +575,9 @@ def get_data(self):
"bench-cmds",
self.cs_type,
self.cs_id,
"stop",
"stop.json.xz",
),
"bench-stop-cmds",
"bench-stop-cmds.json.xz",
)

if self.cs_type in ("client", "server"):
Expand Down Expand Up @@ -625,7 +626,7 @@ def _parse_tool_commands(self, cmds_file):
raise EngineError(
"Failed to load tool commands from %s: %s" % (cmds_file, err)
)
return [(t["name"], t["command"]) for t in data.get("tools", [])]
return [(t["name"], t["argv"]) for t in data.get("tools", [])]

def start_tools(self, one_tool=None):
logger.info("Starting tools")
Expand All @@ -645,7 +646,7 @@ def start_tools(self, one_tool=None):

tools = self._parse_tool_commands(self.tool_start_cmds)
total = 0
for tool_name, tool_command in tools:
for tool_name, tool_argv in tools:
if one_tool and one_tool != tool_name:
logger.info(
"Skipping tool '%s' (engine runs '%s' only)",
Expand All @@ -658,14 +659,14 @@ def start_tools(self, one_tool=None):
tool_dir = os.path.join("tool-data", tool_name)
os.makedirs(tool_dir, exist_ok=True)
logger.info("Starting tool '%s'", tool_name)
run_command("cd %s && %s" % (tool_dir, tool_command))
run_command("cd %s && %s" % (tool_dir, shlex.join(tool_argv)))

if total == 0:
logger.info("No tools configured for this engine")

# ---- Benchmark execution ---------------------------------------------

def run_bench_cmd(self, matching_type, cmd_type, cmd, force=False):
def run_bench_cmd(self, matching_type, cmd_type, argv, force=False):
if self.cs_type != matching_type:
return 0
if not force and (self.abort or self.quit):
Expand All @@ -676,18 +677,21 @@ def run_bench_cmd(self, matching_type, cmd_type, cmd, force=False):
self.quit,
)
return 0
if not cmd:
if not argv:
logger.info("No %s command to run", cmd_type)
return 0
logger.info("Running %s command", cmd_type)
result = run_command(cmd)
result = run_command(shlex.join(argv))
return result.return_code

def _load_bench_cmds(self, filename):
if not os.path.exists(filename):
return []
with open(filename) as fp:
return [line.strip() for line in fp if line.strip()]
data, err = load_json_file(filename, uselzma=True)
if data is None:
logger.error("Failed to load bench commands from %s: %s", filename, err)
return []
return data

def _roadblock_and_evaluate(self, rb_name, timeout, iter_idx,
sample_data, msgs_file=None, do_abort=False):
Expand Down Expand Up @@ -731,19 +735,19 @@ def process_bench_roadblocks(self):
rc = self.do_roadblock("setup-bench-begin", self.default_timeout)
self.roadblock_exit_on_error(rc)

bench_start_cmds = self._load_bench_cmds("bench-start-cmds")
bench_infra_cmds = self._load_bench_cmds("bench-infra-cmds")
bench_runtime_cmds = self._load_bench_cmds("bench-runtime-cmds")
bench_stop_cmds = self._load_bench_cmds("bench-stop-cmds")
bench_start_cmds = self._load_bench_cmds("bench-start-cmds.json.xz")
bench_infra_cmds = self._load_bench_cmds("bench-infra-cmds.json.xz")
bench_runtime_cmds = self._load_bench_cmds("bench-runtime-cmds.json.xz")
bench_stop_cmds = self._load_bench_cmds("bench-stop-cmds.json.xz")

if not bench_start_cmds:
self.abort_error("bench-start-cmds not found", "setup-bench-end")
return

total_tests = len(bench_start_cmds)
sample_data = []
for i, line in enumerate(bench_start_cmds):
iter_samp = line.split()[0]
for i, entry in enumerate(bench_start_cmds):
iter_samp = entry["test"]
parts = iter_samp.split("-")
iter_id = int(parts[0])
samp_id = int(parts[1])
Expand Down Expand Up @@ -772,7 +776,7 @@ def process_bench_roadblocks(self):
sd = sample_data[i]
iter_id = sd["iteration-id"]
samp_id = sd["sample-id"]
iter_samp = bench_start_cmds[i].split()[0]
iter_samp = bench_start_cmds[i]["test"]

iter_samp_dir = os.path.join(
self.cs_dir,
Expand All @@ -784,17 +788,17 @@ def process_bench_roadblocks(self):
sd["rx-msgs-dir"] = cs_rx_msgs_dir

if self.cs_type == "client":
start_cmd = bench_start_cmds[i].split(None, 1)[1] if len(bench_start_cmds[i].split()) > 1 else ""
runtime_cmd = bench_runtime_cmds[i].split(None, 1)[1] if i < len(bench_runtime_cmds) and len(bench_runtime_cmds[i].split()) > 1 else ""
infra_cmd = bench_infra_cmds[i].split(None, 1)[1] if i < len(bench_infra_cmds) and len(bench_infra_cmds[i].split()) > 1 else ""
stop_cmd = ""
start_argv = bench_start_cmds[i].get("argv", [])
runtime_argv = bench_runtime_cmds[i].get("argv", []) if i < len(bench_runtime_cmds) else []
infra_argv = bench_infra_cmds[i].get("argv", []) if i < len(bench_infra_cmds) else []
stop_argv = []
elif self.cs_type == "server":
start_cmd = bench_start_cmds[i].split(None, 1)[1] if len(bench_start_cmds[i].split()) > 1 else ""
stop_cmd = bench_stop_cmds[i].split(None, 1)[1] if i < len(bench_stop_cmds) and len(bench_stop_cmds[i].split()) > 1 else ""
runtime_cmd = ""
infra_cmd = ""
start_argv = bench_start_cmds[i].get("argv", [])
stop_argv = bench_stop_cmds[i].get("argv", []) if i < len(bench_stop_cmds) else []
runtime_argv = []
infra_argv = []
else:
start_cmd = runtime_cmd = infra_cmd = stop_cmd = ""
start_argv = runtime_argv = infra_argv = stop_argv = []

self.abort = False

Expand Down Expand Up @@ -844,7 +848,7 @@ def process_bench_roadblocks(self):
)
self._roadblock_and_evaluate(rb_name, timeout, i, sample_data, msgs_file)

abort_rc = self.run_bench_cmd("client", "infra", infra_cmd)
abort_rc = self.run_bench_cmd("client", "infra", infra_argv)
do_abort_arg = abort_rc != 0
if do_abort_arg:
self.abort = True
Expand All @@ -865,7 +869,7 @@ def process_bench_roadblocks(self):
)
self._roadblock_and_evaluate(rb_name, timeout, i, sample_data, msgs_file)

abort_rc = self.run_bench_cmd("server", "server", start_cmd)
abort_rc = self.run_bench_cmd("server", "server", start_argv)
do_abort_arg = abort_rc != 0
if do_abort_arg:
self.abort = True
Expand Down Expand Up @@ -900,8 +904,8 @@ def process_bench_roadblocks(self):
and not self.quit
and self.cs_type == "client"
and self.cs_id == "1"
and runtime_cmd):
result = run_command(runtime_cmd)
and runtime_argv):
result = run_command(shlex.join(runtime_argv))
runtime_output = result.stdout.strip()

if result.return_code == 0 and runtime_output:
Expand Down Expand Up @@ -959,11 +963,11 @@ def process_bench_roadblocks(self):
msgs_file = prepare_user_msgs_file(
cs_tx_msgs_dir, iter_samp_dir, rb_name, default_recipients
)
wait_for_cmd = (
"python3 /usr/local/bin/engine_lib.py"
" run_bench_cmd '%s' 'client' 'client' '%s' '%s' '0' '%s'"
% (self.cs_type, self.abort, self.quit, start_cmd)
)
wait_for_cmd = [
"python3", "/usr/local/bin/engine_lib.py", "run_bench_cmd",
self.cs_type, "client", "client",
str(self.abort), str(self.quit), "0",
] + start_argv
rc = self.do_roadblock(
rb_name, timeout, messages=msgs_file,
wait_for=wait_for_cmd,
Expand All @@ -978,7 +982,7 @@ def process_bench_roadblocks(self):
if result["is_abort"]:
self.abort = True
else:
abort_rc = self.run_bench_cmd("client", "client", start_cmd)
abort_rc = self.run_bench_cmd("client", "client", start_argv)
do_abort_arg = abort_rc != 0
if do_abort_arg:
self.abort = True
Expand Down Expand Up @@ -1030,7 +1034,7 @@ def process_bench_roadblocks(self):
self._roadblock_and_evaluate(rb_name, timeout, i, sample_data, msgs_file)

self.run_bench_cmd(
"server", "server", stop_cmd, force=force_server_stop
"server", "server", stop_argv, force=force_server_stop
)
abort_rc = 0
do_abort_arg = self.abort
Expand Down Expand Up @@ -1169,7 +1173,7 @@ def cli_stop_tools(working_dir, tool_cmds_file, disabled, one_tool=""):

for tool in data.get("tools", []):
tool_name = tool["name"]
tool_command = tool["command"]
tool_argv = tool["argv"]

if one_tool and one_tool != tool_name:
log.info("Skipping tool '%s' (engine runs '%s' only)", tool_name, one_tool)
Expand All @@ -1180,7 +1184,7 @@ def cli_stop_tools(working_dir, tool_cmds_file, disabled, one_tool=""):
continue

log.info("Stopping tool '%s'", tool_name)
run_command("cd %s && %s" % (tool_dir, tool_command))
run_command("cd %s && %s" % (tool_dir, shlex.join(tool_argv)))

if os.path.isfile(env_file):
shutil.copy2(env_file, tool_dir)
Expand Down Expand Up @@ -1230,7 +1234,7 @@ def cli_send_data(ssh_id_file, src_dir, dest_host, dest_path):


def cli_run_bench_cmd(cs_type, matching_type, cmd_type, abort_str, quit_str,
force_str, cmd):
force_str, *cmd_argv):
"""Standalone entry point for roadblock wait-for: run bench command."""
logging.basicConfig(level=logging.INFO, format="%(message)s")
if cs_type != matching_type:
Expand All @@ -1240,9 +1244,9 @@ def cli_run_bench_cmd(cs_type, matching_type, cmd_type, abort_str, quit_str,
force = force_str not in ("0", "False", "false", "")
if not force and (abort or quit_flag):
return
if not cmd:
if not cmd_argv:
return
result = run_command(cmd)
result = run_command(shlex.join(cmd_argv))
sys.exit(result.return_code)


Expand Down
Loading
Loading