diff --git a/lib/galaxy/workflow/scheduling_manager.py b/lib/galaxy/workflow/scheduling_manager.py index 0e24b2f27e3..222b878c22f 100644 --- a/lib/galaxy/workflow/scheduling_manager.py +++ b/lib/galaxy/workflow/scheduling_manager.py @@ -73,6 +73,7 @@ class WorkflowSchedulingManager( object ): workflow_invocation.state = model.WorkflowInvocation.states.NEW scheduler = request_params.get( "scheduler", None ) or self.default_scheduler_id handler = self._get_handler() + log.info("Queueing workflow invocation for handler [%s]" % handler) workflow_invocation.scheduler = scheduler workflow_invocation.handler = handler diff --git a/scripts/summarize_timings.py b/scripts/summarize_timings.py index e4e9a6f422a..1740563358a 100644 --- a/scripts/summarize_timings.py +++ b/scripts/summarize_timings.py @@ -1,9 +1,7 @@ +"""Script to parse timings out of a Galaxy log and summarize.""" from __future__ import print_function -try: - from argparse import ArgumentParser -except ImportError: - ArgumentParser = None +from argparse import ArgumentParser import re import numpy @@ -15,19 +13,19 @@ TIMING_LINE_PATTERN = re.compile("\((\d+.\d+) ms\)") def main(argv=None): - if ArgumentParser is None: - raise Exception("Test requires Python 2.7") + """Entry point for script.""" arg_parser = ArgumentParser(description=DESCRIPTION) arg_parser.add_argument("--file", default="paster.log") arg_parser.add_argument("--print_lines", default=False, action="store_true") - arg_parser.add_argument("--pattern") + arg_parser.add_argument("--pattern", default=None) args = arg_parser.parse_args(argv) print_lines = args.print_lines - filter_pattern = re.compile(args.pattern) + pattern_str = args.pattern + filter_pattern = re.compile(pattern_str) if pattern_str is not None else None times = [] for line in open(args.file, "r"): - if not filter_pattern.search(line): + if filter_pattern and not filter_pattern.search(line): continue match = TIMING_LINE_PATTERN.search(line) diff --git a/test/functional/tools/for_workflows/create_input_collection.xml b/test/functional/tools/for_workflows/create_input_collection.xml new file mode 100644 index 00000000000..04c7093a6b8 --- /dev/null +++ b/test/functional/tools/for_workflows/create_input_collection.xml @@ -0,0 +1,42 @@ + + This tool is used to create a collection of text files. + + mkdir outputs; cd outputs; python $script + + + +import os + +for i in range($collection_size): + template = "File number %s\n" + contents = template % i + with open(str(i), "w") as f: + f.write(contents) + + + + + + + + + + + + + + + + + + + + + + + + + + + + diff --git a/test/functional/tools/for_workflows/split.xml b/test/functional/tools/for_workflows/split.xml new file mode 100644 index 00000000000..def9c17a8d0 --- /dev/null +++ b/test/functional/tools/for_workflows/split.xml @@ -0,0 +1,33 @@ + + + bash $script + + + + mkdir outputs; + cd outputs; + i=1; + while read -r line || [[ -n "\$line" ]]; do + printf "\$line\n" > \$i ; + i=\$[\$i +1]; + done < "$input1"; + + + + + + + + + + + + + + + + + + + + diff --git a/test/functional/tools/samples_tool_conf.xml b/test/functional/tools/samples_tool_conf.xml index 5ef9bfa6e12..bc9a64b7c6a 100644 --- a/test/functional/tools/samples_tool_conf.xml +++ b/test/functional/tools/samples_tool_conf.xml @@ -96,6 +96,8 @@ + + diff --git a/test/manual/launch_and_run.sh b/test/manual/launch_and_run.sh new file mode 100755 index 00000000000..de389301573 --- /dev/null +++ b/test/manual/launch_and_run.sh @@ -0,0 +1,78 @@ +#!/bin/bash +#set -e + +# Open and few the contents of a docker-galaxy-stable container. +# docker run --rm -i -t bgruening/galaxy-stable /bin/bash + +pwd_dir=$(pwd) +GALAXY_ROOT=`dirname $0`/../.. +cd $GALAXY_ROOT +GALAXY_ROOT=$(pwd) +SCRIPT_DIR="$GALAXY_ROOT/test/manual" + +manual_test_script=$1 +shift +manual_test_script_args="$@" + +GALAXY_VIRTUAL_ENV="${GALAXY_VIRTUAL_ENV:-.venv}" + +# Docker options defined to reflect run_tests.sh names and behavior. +DOCKER_DEFAULT_IMAGE='bgruening/galaxy-stable' + +DOCKER_EXTRA_ARGS=${DOCKER_ARGS:-""} +DOCKER_RUN_EXTRA_ARGS=${DOCKER_RUN_EXTRA_ARGS:-""} +DOCKER_IMAGE=${DOCKER_IMAGE:-${DOCKER_DEFAULT_IMAGE}} +# Root for Galaxy in the docker container +DOCKER_GALAXY_ROOT=${DOCKER_GALAXY_ROOT:-/galaxy-central} + +# Location of this script's directory when mounted into the container. +DOCKER_SCRIPT_DIR=/etc/galaxy/manual + +GALAXY_PORT=${GALAXY_PORT:-"any_free"} +if [ "$GALAXY_PORT" == "any_free" ]; +then + GALAXY_PORT=`python -c 'import socket; s=socket.socket(); s.bind(("", 0)); print(s.getsockname()[1]); s.close()'` +fi + +GALAXY_URL=${GALAXY_URL:-http://localhost:${GALAXY_PORT}} +GALAXY_MASTER_API_KEY=${GALAXY_MASTER_API_KEY:-HSNiugRFvgT574F43jZ7N9F3} + +LOGS_DIR=`cd "$LOGS_DIR"; pwd` +WORK_DIR=`mktemp --tmpdir=$LOGS_DIR -d -t gxperfXXXX` +echo "WORK_DIR is ${WORK_DIR}" +NAME=`basename $WORK_DIR` + +GALAXY_HANDLER_NUMPROCS=${GALAXY_HANDLER_NUMPROCS:-1} + +DOCKER_ENVIRONMENT="\ +-e NONUSE=nodejs,proftp,reports \ +-e GALAXY_HANDLER_NUMPROCS=$GALAXY_HANDLER_NUMPROCS \ +-e GALAXY_CONFIG_OVERRIDE_TOOL_CONFIG_FILE=$DOCKER_GALAXY_ROOT/test/functional/tools/samples_tool_conf.xml \ +-e GALAXY_CONFIG_ENABLE_BETA_WORKFLOW_MODULES=true \ +-e GALAXY_CONFIG_OVERRIDE_ENABLE_BETA_TOOL_FORMATS=true \ +" + +if [ $manual_test_script == "workflows_scaling" ]; +then + DOCKER_ENVIRONMENT="$DOCKER_ENVIRONMENT -e GALAXY_CONFIG_JOB_CONFIG_FILE=$DOCKER_SCRIPT_DIR/workflow_job_conf.xml " +fi + +# Mount logs, local galaxy changes, and local galaxy config. +DOCKER_VOLUMES="\ +-v $WORK_DIR:/galaxy_logs \ +-v $GALAXY_ROOT/lib:$DOCKER_GALAXY_ROOT/lib \ +-v $GALAXY_ROOT/test:/galaxy-central/test \ +-v $SCRIPT_DIR:$DOCKER_SCRIPT_DIR \ +" +DOCKER_RUN_ARGS="$DOCKER_RUN_EXTRA_ARGS -d -p ${GALAXY_PORT}:80 -i -t $DOCKER_VOLUMES $DOCKER_ENVIRONMENT" + +docker_image_id=`docker $DOCKER_EXTRA_ARGS run $DOCKER_RUN_ARGS ${DOCKER_IMAGE}` + +echo "Docker container with id $docker_image_id launched. Inspect with 'docker exec -i -t $docker_image_id /bin/bash'." + +# Wait for Galaxy to be available +for i in {1..40}; do curl --silent --fail ${GALAXY_URL}/api/version && break || sleep 5; done + +${GALAXY_VIRTUAL_ENV}/bin/python test/manual/$manual_test_script.py --api_key ${GALAXY_MASTER_API_KEY} --host ${GALAXY_URL} $manual_test_script_args +docker exec -i -t $docker_image_id /bin/bash -c "cp /home/galaxy/*log /galaxy_logs" +docker kill $docker_image_id diff --git a/test/manual/workflow_job_conf.xml b/test/manual/workflow_job_conf.xml new file mode 100644 index 00000000000..0e5ebef11fb --- /dev/null +++ b/test/manual/workflow_job_conf.xml @@ -0,0 +1,26 @@ + + + + + /usr/lib/slurm-drmaa/lib/libdrmaa.so + + + + + + + + + + + + + + + + + + + + + diff --git a/test/manual/workflows_scaling.py b/test/manual/workflows_scaling.py index 13c55eef9b7..f9d47396a02 100644 --- a/test/manual/workflows_scaling.py +++ b/test/manual/workflows_scaling.py @@ -1,3 +1,10 @@ +#!/usr/bin/env python +"""A small script to drive workflow performance testing. + +% ./test/manual/launch_and_run.sh workflows_scaling --collection_size 500 --workflow_depth 4 +$ .venv/bin/python scripts/summarize_timings.py --file /tmp//handler1.log --pattern 'Workflow step' +$ .venv/bin/python scripts/summarize_timings.py --file /tmp//handler1.log --pattern 'Created step' +""" import functools import json import os @@ -9,10 +16,7 @@ from uuid import uuid4 galaxy_root = os.path.abspath(os.path.join(os.path.dirname(__file__), os.path.pardir, os.path.pardir)) sys.path[1:1] = [ os.path.join( galaxy_root, "lib" ), os.path.join( galaxy_root, "test" ) ] -try: - from argparse import ArgumentParser -except ImportError: - ArgumentParser = None +from argparse import ArgumentParser import requests from bioblend import galaxy @@ -24,20 +28,30 @@ DESCRIPTION = "Script to exercise the workflow engine." def main(argv=None): - if ArgumentParser is None: - raise Exception("Test requires Python 2.7") + """Entry point for workflow driving.""" arg_parser = ArgumentParser(description=DESCRIPTION) arg_parser.add_argument("--api_key", default="testmasterapikey") arg_parser.add_argument("--host", default="http://localhost:8080/") arg_parser.add_argument("--collection_size", type=int, default=20) + + arg_parser.add_argument("--schedule_only_test", default=False, action="store_true") arg_parser.add_argument("--workflow_depth", type=int, default=10) - arg_parser.add_argument("--two_outputs", default=False, action="store_true") arg_parser.add_argument("--workflow_count", type=int, default=1) + group = arg_parser.add_mutually_exclusive_group() + group.add_argument("--two_outputs", default=False, action="store_true") + group.add_argument("--wave_simple", default=False, action="store_true") + args = arg_parser.parse_args(argv) + uuid = str(uuid4()) workflow_struct = _workflow_struct(args, uuid) + + has_input = any([s.get("type", "tool") == "input_collection" for s in workflow_struct]) + if not has_input: + uuid = None + gi = _gi(args) workflow = yaml_to_workflow.python_to_workflow(workflow_struct) @@ -61,13 +75,16 @@ def _run(args, gi, workflow_id, uuid): dataset_collection_populator = GiDatasetCollectionPopulator(gi) history_id = dataset_populator.new_history() - contents = [] - for i in range(args.collection_size): - contents.append("random dataset number #%d" % i) - hdca = dataset_collection_populator.create_list_in_history( history_id, contents=contents ).json() - label_map = { - uuid: {"src": "hdca", "id": hdca["id"]}, - } + if uuid is not None: + contents = [] + for i in range(args.collection_size): + contents.append("random dataset number #%d" % i) + hdca = dataset_collection_populator.create_list_in_history( history_id, contents=contents ).json() + label_map = { + uuid: {"src": "hdca", "id": hdca["id"]}, + } + else: + label_map = {} workflow_request = dict( history="hist_id=%s" % history_id, @@ -77,10 +94,23 @@ def _run(args, gi, workflow_id, uuid): invoke_response = dataset_populator._post( url, data=workflow_request ).json() invocation_id = invoke_response["id"] workflow_populator = GiWorkflowPopulator(gi) - workflow_populator.wait_for_workflow( workflow_id, invocation_id, history_id, timeout=LONG_TIMEOUT ) + if args.schedule_only_test: + workflow_populator.wait_for_invocation( + workflow_id, + invocation_id, + timeout=LONG_TIMEOUT, + ) + else: + workflow_populator.wait_for_workflow( + workflow_id, + invocation_id, + history_id, + timeout=LONG_TIMEOUT, + ) class GiPostGetMixin: + """Mixin for adapting Galaxy API testing helpers to bioblend.""" def _get(self, route): return self._gi.make_get_request(self.__url(route)) @@ -95,14 +125,18 @@ class GiPostGetMixin: class GiDatasetPopulator(helpers.BaseDatasetPopulator, GiPostGetMixin): + """Utility class for dealing with datasets and histories.""" def __init__(self, gi): + """Construct a dataset populator from a bioblend GalaxyInstance.""" self._gi = gi class GiDatasetCollectionPopulator(helpers.BaseDatasetCollectionPopulator, GiPostGetMixin): + """Utility class for dealing with dataset collections.""" def __init__(self, gi): + """Construct a dataset collection populator from a bioblend GalaxyInstance.""" self._gi = gi self.dataset_populator = GiDatasetPopulator(gi) @@ -112,8 +146,10 @@ class GiDatasetCollectionPopulator(helpers.BaseDatasetCollectionPopulator, GiPos class GiWorkflowPopulator(helpers.BaseWorkflowPopulator, GiPostGetMixin): + """Utility class for dealing with workflows.""" def __init__(self, gi): + """Construct a workflow populator from a bioblend GalaxyInstance.""" self._gi = gi self.dataset_populator = GiDatasetPopulator(gi) @@ -121,21 +157,23 @@ class GiWorkflowPopulator(helpers.BaseWorkflowPopulator, GiPostGetMixin): def _workflow_struct(args, input_uuid): if args.two_outputs: return _workflow_struct_two_outputs(args, input_uuid) + elif args.wave_simple: + return _workflow_struct_wave(args, input_uuid) else: return _workflow_struct_simple(args, input_uuid) def _workflow_struct_simple(args, input_uuid): workflow_struct = [ - {"type": "input_collection", "uuid": input_uuid}, - {"tool_id": "cat1", "state": {"input1": _link(0)}} + {"tool_id": "create_input_collection", "state": {"collection_size": args.collection_size}}, + {"tool_id": "cat", "state": {"input1": _link(0, "output")}} ] workflow_depth = args.workflow_depth for i in range(workflow_depth): link = str(i + 1) + "#out_file1" workflow_struct.append( - {"tool_id": "cat1", "state": {"input1": _link(link)}} + {"tool_id": "cat", "state": {"input1": _link(link)}} ) return workflow_struct @@ -143,7 +181,7 @@ def _workflow_struct_simple(args, input_uuid): def _workflow_struct_two_outputs(args, input_uuid): workflow_struct = [ {"type": "input_collection", "uuid": input_uuid}, - {"tool_id": "cat1", "state": {"input1": _link(0), "input2": _link(0)}} + {"tool_id": "cat", "state": {"input1": _link(0), "input2": _link(0)}} ] workflow_depth = args.workflow_depth @@ -151,12 +189,30 @@ def _workflow_struct_two_outputs(args, input_uuid): link1 = str(i + 1) + "#out_file1" link2 = str(i + 1) + "#out_file2" workflow_struct.append( - {"tool_id": "cat1", "state": {"input1": _link(link1), "input2": _link(link2)}} + {"tool_id": "cat", "state": {"input1": _link(link1), "input2": _link(link2)}} ) return workflow_struct -def _link(link): +def _workflow_struct_wave(args, input_uuid): + workflow_struct = [ + {"tool_id": "create_input_collection", "state": {"collection_size": args.collection_size}}, + {"tool_id": "cat_list", "state": {"input1": _link(0, "output")}} + ] + + workflow_depth = args.workflow_depth + for i in range(workflow_depth): + step = i + 2 + if step % 2 == 1: + workflow_struct += [{"tool_id": "cat_list", "state": {"input1": _link(step - 1, "output")}}] + else: + workflow_struct += [{"tool_id": "split", "state": {"input1": _link(step - 1, "out_file1") }}] + return workflow_struct + + +def _link(link, output_name=None): + if output_name is not None: + link = str(link) + "#" + output_name return {"$link": link}