Merge branch 'release_20.01' into dev

This commit is contained in:
Nicola Soranzo
2020-01-09 22:27:51 +00:00
27 changed files with 310 additions and 193 deletions
+48
View File
@@ -0,0 +1,48 @@
name: Integration
on: [push, pull_request]
env:
GALAXY_TEST_DBURI: 'postgres://postgres:postgres@localhost:5432/galaxy?client_encoding=utf8'
jobs:
test:
name: Test
runs-on: ubuntu-18.04
strategy:
matrix:
python-version: [3.7]
subset: ['upload_datatype', 'extended_metadata', 'kubernetes', 'not (upload_datatype or extended_metadata or kubernetes)']
services:
postgres:
image: postgres:11
env:
POSTGRES_USER: postgres
POSTGRES_PASSWORD: postgres
POSTGRES_DB: postgres
ports:
- 5432:5432
steps:
- name: Setup Minikube
if: matrix.subset == 'kubernetes'
id: minikube
uses: CodingNagger/minikube-setup-action@v1.0.2
- name: Launch Minikube
if: matrix.subset == 'kubernetes'
run: eval ${{ steps.minikube.outputs.launcher }}
- name: Check pods
if: matrix.subset == 'kubernetes'
run: |
kubectl get pods
- uses: actions/checkout@v1
with:
fetch-depth: 1
- uses: actions/setup-python@v1
with:
python-version: ${{ matrix.python-version }}
- name: Cache venv dir
uses: actions/cache@v1
id: cache-venv-galaxy
with:
path: .venv
key: cache-venv-galaxy-${{ matrix.python-version }}-${{ hashFiles('requirements.txt') }}
- name: Run tests
run: './run_tests.sh -integration test/integration -- -k "${{ matrix.subset }}"'
+6
View File
@@ -101,6 +101,12 @@ UWSGI_OPTIONS = OrderedDict([
'default': '/favicon.ico=static/favicon.ico',
'type': 'str',
}),
('static-safe', {
'key': 'static-safe',
'desc': """Allow serving images out of `client`. Most modern Galaxy interfaces bundle all of this, but some older pages still serve these via symlink, requiring this rule.""",
'default': 'client/galaxy/images',
'type': 'str',
}),
('master', {
'desc': """Enable the master process manager. Disabled by default for maximum compatibility with CTRL+C, but should be enabled for use with --daemon and/or production deployments.""",
'default': False,
@@ -50,6 +50,11 @@ uwsgi:
# Mapping to serve the favicon.
static-map: /favicon.ico=static/favicon.ico
# Allow serving images out of `client`. Most modern Galaxy interfaces
# bundle all of this, but some older pages still serve these via
# symlink, requiring this rule.
static-safe: client/galaxy/images
# Enable the master process manager. Disabled by default for maximum
# compatibility with CTRL+C, but should be enabled for use with
# --daemon and/or production deployments.
@@ -32,6 +32,11 @@ uwsgi:
# Mapping to serve the favicon.
static-map: /favicon.ico=static/favicon.ico
# Allow serving images out of `client`. Most modern Galaxy interfaces
# bundle all of this, but some older pages still serve these via
# symlink, requiring this rule.
static-safe: client/galaxy/images
# Enable the master process manager. Disabled by default for maximum
# compatibility with CTRL+C, but should be enabled for use with
# --daemon and/or production deployments.
@@ -32,6 +32,11 @@ uwsgi:
# Mapping to serve the favicon.
static-map: /favicon.ico=static/favicon.ico
# Allow serving images out of `client`. Most modern Galaxy interfaces
# bundle all of this, but some older pages still serve these via
# symlink, requiring this rule.
static-safe: client/galaxy/images
# Enable the master process manager. Disabled by default for maximum
# compatibility with CTRL+C, but should be enabled for use with
# --daemon and/or production deployments.
+3
View File
@@ -37,6 +37,7 @@ DEFAULT_DB_CONN = 'sqlite:///./database/universe.sqlite?isolation_level=IMMEDIAT
SAMPLES_PATH = os.path.abspath(os.path.join(os.path.dirname(__file__), 'sample'))
GALAXY_CONFIG_TEMPLATE_FILE = os.path.join(SAMPLES_PATH, 'galaxy.yml.sample')
STATIC_PATH = os.path.abspath(os.path.join(os.path.dirname(__file__), os.pardir, 'web', 'framework', 'static'))
CLIENT_PATH = os.path.abspath(os.path.join(os.path.dirname(__file__), os.pardir, os.pardir, os.pardir, 'client'))
MSG_CONFIG_SUMMARY = """
For help on configuring Galaxy, consult the documentation at: \n {}
@@ -55,6 +56,7 @@ GALAXY_CONFIG_SUBSTITUTIONS = {
' static-map: /static/style=static/style/blue': ' static-map: /static=${static_path}/style/blue',
' static-map: /static=static': ' static-map: /static=${static_path}',
' static-map: /favicon.ico=static/favicon.ico': ' static-map: /static=${static_path}/favicon.ico',
' static-safe: client/galaxy/images': ' ${client_path}/galaxy/images',
' virtualenv: .venv': ' #venv: .venv # not used when running installed',
' pythonpath: lib': ' #pythonpath: lib # not used when running installed',
' #config_dir: false': ' config_dir: ${config_dir}',
@@ -147,6 +149,7 @@ def _handle_galaxy_yml(args, config_dir, data_dir):
uwsgi_transport=uwsgi_transport,
config_dir=config_dir,
data_dir=data_dir,
client_dir=CLIENT_PATH,
static_path=STATIC_PATH,
database_connection=args.db_conn,
)
+6 -3
View File
@@ -11,7 +11,10 @@ import pkg_resources
import yaml
from galaxy.containers import parse_containers_config
from galaxy.util import asbool
from galaxy.util import (
asbool,
which,
)
from galaxy.util.properties import (
find_config_file,
load_app_properties
@@ -145,7 +148,7 @@ class ConditionalDependencies(object):
return "galaxy.jobs.runners.pbs:PBSJobRunner" in self.job_runners
def check_pykube(self):
return "galaxy.jobs.runners.kubernetes:KubernetesJobRunner" in self.job_runners
return "galaxy.jobs.runners.kubernetes:KubernetesJobRunner" in self.job_runners or which('kubectl')
def check_chronos_python(self):
return "galaxy.jobs.runners.chronos:ChronosJobRunner" in self.job_runners
@@ -194,7 +197,7 @@ class ConditionalDependencies(object):
def optional(config_file=None):
if not config_file:
config_file = find_config_file(['galaxy', 'universe_wsgi'])
config_file = find_config_file(['galaxy', 'universe_wsgi'], include_samples=True)
if not config_file:
print("galaxy.dependencies.optional: no config file found", file=sys.stderr)
return []
+1 -1
View File
@@ -2111,7 +2111,7 @@ class JobWrapper(HasResourceParameters):
"connection_configuration": container.connection_configuration,
}, f)
return "(python '%s'/lib/galaxy_ext/container_monitor/monitor.py &) " % exec_dir
return "(python '%s'/lib/galaxy_ext/container_monitor/monitor.py &); sleep 1 " % exec_dir
@property
def user(self):
+6 -4
View File
@@ -161,10 +161,6 @@ class ShellJobRunner(AsynchronousJobRunner):
if ajs.job_wrapper.get_state() == model.Job.states.DELETED:
continue
external_metadata = not asbool(ajs.job_wrapper.job_destination.params.get("embed_metadata_in_job", DEFAULT_EMBED_METADATA_IN_JOB))
if external_metadata:
self._handle_metadata_externally(ajs.job_wrapper, resolve_requirements=True)
log.debug("(%s/%s) job not found in batch state check" % (id_tag, external_job_id))
shell_params, job_params = self.parse_destination_params(ajs.job_destination.params)
shell, job_interface = self.get_cli_plugins(shell_params, job_params)
@@ -187,6 +183,9 @@ class ShellJobRunner(AsynchronousJobRunner):
ajs.running = True
ajs.old_state = state
if state == model.Job.states.OK:
external_metadata = not asbool(ajs.job_wrapper.job_destination.params.get("embed_metadata_in_job", DEFAULT_EMBED_METADATA_IN_JOB))
if external_metadata:
self.work_queue.put((self.handle_metadata_externally, ajs))
log.debug('(%s/%s) job execution finished, running job wrapper finish method' % (id_tag, external_job_id))
self.work_queue.put((self.finish_job, ajs))
else:
@@ -194,6 +193,9 @@ class ShellJobRunner(AsynchronousJobRunner):
# Replace the watch list with the updated version
self.watched = new_watched
def handle_metadata_externally(self, ajs):
self._handle_metadata_externally(ajs.job_wrapper, resolve_requirements=True)
def __handle_out_of_memory(self, ajs, external_job_id):
shell_params, job_params = self.parse_destination_params(ajs.job_destination.params)
shell, job_interface = self.get_cli_plugins(shell_params, job_params)
@@ -6,7 +6,8 @@ from pulsar.managers.util.retry import RetryActionExecutor
from galaxy.util import (
smart_str,
unicodify
string_as_bool,
unicodify,
)
from galaxy.util.bunch import Bunch
from .local import LocalShell
@@ -40,11 +41,11 @@ class RemoteShell(LocalShell):
class SecureShell(RemoteShell):
SSH_NEW_KEY_STRING = 'Are you sure you want to continue connecting'
def __init__(self, rsh='ssh', rcp='scp', private_key=None, port=None, strict_host_key_checking=True, **kwargs):
strict_host_key_checking = "yes" if strict_host_key_checking else "no"
options = ["-o", "StrictHostKeyChecking=%s" % strict_host_key_checking]
options = []
if not string_as_bool(strict_host_key_checking):
options.extend(["-o", "StrictHostKeyChecking=no", "-o", "UserKnownHostsFile=/dev/null"])
options.extend(["-o", "ConnectTimeout=60"])
if private_key:
options.extend(['-i', private_key])
@@ -55,13 +56,14 @@ class SecureShell(RemoteShell):
class ParamikoShell(object):
def __init__(self, username, hostname, password=None, private_key=None, port=22, timeout=60, **kwargs):
def __init__(self, username, hostname, password=None, private_key=None, port=22, timeout=60, strict_host_key_checking=True, **kwargs):
self.username = username
self.hostname = hostname
self.password = password
self.private_key = private_key
self.port = int(port) if port else port
self.timeout = int(timeout) if timeout else timeout
self.strict_host_key_checking = string_as_bool(strict_host_key_checking)
self.ssh = None
self.retry_action_executor = RetryActionExecutor(max_retries=100, interval_max=300)
self.connect()
@@ -69,7 +71,7 @@ class ParamikoShell(object):
def connect(self):
log.info("Attempting establishment of new paramiko SSH channel")
self.ssh = paramiko.SSHClient()
self.ssh.set_missing_host_key_policy(paramiko.AutoAddPolicy())
self.ssh.set_missing_host_key_policy(paramiko.RejectPolicy() if self.strict_host_key_checking else paramiko.WarningPolicy())
self.ssh.connect(hostname=self.hostname,
port=self.port,
username=self.username,
+1 -1
View File
@@ -270,7 +270,7 @@ class InteractiveToolManager(object):
entry_point = trans.sa_session.query(model.InteractiveToolEntryPoint).get(entry_point_id)
if self.app.interactivetool_manager.can_access_entry_point(trans, entry_point):
if entry_point.active:
return self.target_if_active(entry_point)
return self.target_if_active(trans, entry_point)
elif entry_point.deleted:
raise exceptions.MessageException('InteractiveTool has ended. You will have to start a new one.')
else:
@@ -180,7 +180,7 @@ def mull_targets(
if DEST_BASE_IMAGE:
dest_base_image = DEST_BASE_IMAGE
else:
dest_base_image = DEFAULT_EXTENDED_BASE_IMAGE if not any_target_requires_extended_base(targets) else DEST_BASE_IMAGE
dest_base_image = DEFAULT_EXTENDED_BASE_IMAGE if any_target_requires_extended_base(targets) else DEST_BASE_IMAGE
targets = list(targets)
if involucro_context is None:
+1 -1
View File
@@ -271,7 +271,7 @@ def xml_to_string(elem, pretty=False):
except TypeError as e:
# we assume this is a comment
if hasattr(elem, 'text'):
return "<!-- %s -->\n" % elem.text
return u"<!-- %s -->\n" % elem.text
else:
raise e
if xml_str and pretty:
+4
View File
@@ -563,6 +563,10 @@ class GalaxyWebTransaction(base.DefaultWebTransaction,
return
except IndexError:
pass
authnz_controller_base = url_for(controller='authnz', action='index')
if self.request.path.startswith(authnz_controller_base):
# All authnz requests pass through
return
# redirect to root if the path is not in the list above
if self.request.path not in allowed_paths:
login_url = url_for(controller='root', action='login', redirect=self.request.path)
+1 -1
View File
@@ -22,7 +22,7 @@ def parse_ports(container_name, connection_configuration):
preexec_fn=os.setpgrp)
if exit_code == 0:
stdout_file.seek(0)
ports_raw = stdout_file.read()
ports_raw = stdout_file.read().decode('utf-8')
return ports_raw
+12 -8
View File
@@ -4,6 +4,7 @@ import fcntl
import logging
import os
import random
import re
import shutil
import signal
import socket
@@ -62,13 +63,15 @@ INSTALLED_TOOL_PANEL_CONFIGS = [
]
REALTIME_PROXY_TEMPLATE = string.Template(r"""
uwsgi:
realtime_map: $tempdir/realtime_map.sqlite
python-raw: scripts/realtime/key_type_token_mapping.py
route-host: ^([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.(realtime\.$test_host:$test_port)$ goto:realtime
http-raw-body: true
interactivetools_map: $tempdir/interactivetools_map.sqlite
python-raw: scripts/interactivetools/key_type_token_mapping.py
route-host: ^([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)-([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.(interactivetool\.$test_host:$test_port)$ goto:interactivetool
route-run: goto:endendend
route-label: realtime
route-host: ^([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.(realtime\.$test_host:$test_port)$ rpcvar:TARGET_HOST rtt_key_type_token_mapper_cached $2 $1 $3 $4 $0 5
route-label: interactivetool
route-host: ^([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)-([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.(interactivetool\.$test_host:$test_port)$ rpcvar:TARGET_HOST rtt_key_type_token_mapper_cached $1 $3 $2 $4 $0 5
route-if-not: empty:${TARGET_HOST} httpdumb:${TARGET_HOST}
route: .* break:404 Not Found
route-label: endendend
""")
@@ -729,8 +732,9 @@ def launch_uwsgi(kwargs, tempdir, prefix=DEFAULT_CONFIG_PREFIX, config_object=No
enable_realtime_mapping = getattr(config_object, "enable_realtime_mapping", False)
if enable_realtime_mapping:
config["galaxy"]["realtime_prefix"] = "realtime"
config["galaxy"]["realtime_map"] = os.path.join(tempdir, "realtime_map.sqlite")
config["galaxy"]["interactivetools_prefix"] = "interactivetool"
config["galaxy"]["interactivetools_map"] = os.path.join(tempdir, "interactivetools_map.sqlite")
config['galaxy']['interactivetools_enable'] = True
yaml_config_path = os.path.join(tempdir, "galaxy.yml")
with open(yaml_config_path, "w") as f:
@@ -743,7 +747,7 @@ def launch_uwsgi(kwargs, tempdir, prefix=DEFAULT_CONFIG_PREFIX, config_object=No
old_contents = f.read()
with open(yaml_config_path, "w") as f:
test_port = str(port) if port else r"[0-9]+"
test_host = host or "localhost"
test_host = re.escape(host) if host else "localhost"
uwsgi_section = REALTIME_PROXY_TEMPLATE.safe_substitute(test_host=test_host, test_port=test_port, tempdir=tempdir)
f.write(uwsgi_section)
f.write(old_contents)
+3 -3
View File
@@ -31,7 +31,7 @@ def skip_if_jenkins(cls):
def skip_unless_executable(executable):
if which(executable):
return _identity
return skip("PATH doesn't contain executable %s" % executable)
return pytest.mark.skip("PATH doesn't contain executable %s" % executable)
def skip_unless_docker():
@@ -47,10 +47,10 @@ def k8s_config_path():
def skip_unless_fixed_port():
if os.environ.get("GALAXY_TEST_PORT"):
if os.environ.get("GALAXY_TEST_PORT_RANDOM") != "1":
return _identity
return skip("GALAXY_TEST_PORT must be set for this test.")
return pytest.mark.skip("GALAXY_TEST_PORT must be set for this test.")
class IntegrationInstance(UsesApiTestCaseMixin):
+1
View File
@@ -602,6 +602,7 @@ if [ -z "$skip_common_startup" ]; then
export GALAXY_CONFIG_OVERRIDE_DATABASE_CONNECTION
fi
./scripts/common_startup.sh $skip_venv $no_create_venv $no_replace_pip $replace_pip $skip_client_build --dev-wheels || exit 1
unset GALAXY_CONFIG_OVERRIDE_DATABASE_CONNECTION
fi
. ./scripts/common_startup_functions.sh
+115 -119
View File
@@ -146,150 +146,146 @@ def main(argv):
blacklisted_tools = config['sanitization']['tools']
annotate('export_jobs_start', 'Exporting Jobs')
handle_job = io.open(REPORT_BASE + '.jobs.tsv', 'w', encoding='utf-8')
handle_job.write('\t'.join(('id', 'tool_id', 'tool_version', 'state', 'create_time')) + '\n')
for offset_start in range(last_job_sent, end_job_id, args.batch_size):
logging.debug("Processing %s:%s", offset_start, min(end_job_id, offset_start + args.batch_size))
for job in sa_session.query(model.Job.id, model.Job.user_id, model.Job.tool_id, model.Job.tool_version, model.Job.state, model.Job.create_time) \
.filter(model.Job.id > offset_start) \
.filter(model.Job.id <= min(end_job_id, offset_start + args.batch_size)) \
.all():
# If the tool is blacklisted, exclude everywhere
if job[2] in blacklisted_tools:
continue
with io.open(REPORT_BASE + '.jobs.tsv', 'w', encoding='utf-8') as handle_job:
handle_job.write(u'\t'.join(('id', 'tool_id', 'tool_version', 'state', 'create_time')) + '\n')
for offset_start in range(last_job_sent, end_job_id, args.batch_size):
logging.debug("Processing %s:%s", offset_start, min(end_job_id, offset_start + args.batch_size))
for job in sa_session.query(model.Job.id, model.Job.user_id, model.Job.tool_id, model.Job.tool_version, model.Job.state, model.Job.create_time) \
.filter(model.Job.id > offset_start) \
.filter(model.Job.id <= min(end_job_id, offset_start + args.batch_size)) \
.all():
# If the tool is blacklisted, exclude everywhere
if job[2] in blacklisted_tools:
continue
try:
line = [
str(job[0]), # id
job[2], # tool_id
job[3], # tool_version
job[4], # state
str(job[5]) # create_time
]
cline = unicodify('\t'.join(line) + '\n')
handle_job.write(cline)
except Exception:
logging.warning("Unable to write out a 'handle_job' row. Ignoring the row.", exc_info=True)
continue
# meta counts
job_state_data[job[4]] += 1
active_users[job[1]] += 1
job_tool_map[job[0]] = job[2]
handle_job.close()
try:
line = [
str(job[0]), # id
job[2], # tool_id
job[3], # tool_version
job[4], # state
str(job[5]) # create_time
]
cline = unicodify('\t'.join(line) + '\n')
handle_job.write(cline)
except Exception:
logging.warning("Unable to write out a 'handle_job' row. Ignoring the row.", exc_info=True)
continue
# meta counts
job_state_data[job[4]] += 1
active_users[job[1]] += 1
job_tool_map[job[0]] = job[2]
annotate('export_jobs_end')
annotate('export_datasets_start', 'Exporting Datasets')
handle_datasets = io.open(REPORT_BASE + '.datasets.tsv', 'w', encoding='utf-8')
handle_datasets.write('\t'.join(('job_id', 'dataset_id', 'extension', 'file_size', 'param_name', 'type')) + '\n')
for offset_start in range(last_job_sent, end_job_id, args.batch_size):
logging.debug("Processing %s:%s", offset_start, min(end_job_id, offset_start + args.batch_size))
with io.open(REPORT_BASE + '.datasets.tsv', 'w', encoding='utf-8') as handle_datasets:
handle_datasets.write(u'\t'.join(('job_id', 'dataset_id', 'extension', 'file_size', 'param_name', 'type')) + '\n')
for offset_start in range(last_job_sent, end_job_id, args.batch_size):
logging.debug("Processing %s:%s", offset_start, min(end_job_id, offset_start + args.batch_size))
# four queries: JobToInputDatasetAssociation, JobToOutputDatasetAssociation, HistoryDatasetAssociation, Dataset
# four queries: JobToInputDatasetAssociation, JobToOutputDatasetAssociation, HistoryDatasetAssociation, Dataset
job_to_input_hda_ids = sa_session.query(model.JobToInputDatasetAssociation.job_id, model.JobToInputDatasetAssociation.dataset_id,
model.JobToInputDatasetAssociation.name) \
.filter(model.JobToInputDatasetAssociation.job_id > offset_start) \
.filter(model.JobToInputDatasetAssociation.job_id <= min(end_job_id, offset_start + args.batch_size)) \
.all()
job_to_input_hda_ids = sa_session.query(model.JobToInputDatasetAssociation.job_id, model.JobToInputDatasetAssociation.dataset_id,
model.JobToInputDatasetAssociation.name) \
.filter(model.JobToInputDatasetAssociation.job_id > offset_start) \
.filter(model.JobToInputDatasetAssociation.job_id <= min(end_job_id, offset_start + args.batch_size)) \
.all()
job_to_output_hda_ids = sa_session.query(model.JobToOutputDatasetAssociation.job_id, model.JobToOutputDatasetAssociation.dataset_id,
model.JobToOutputDatasetAssociation.name) \
.filter(model.JobToOutputDatasetAssociation.job_id > offset_start) \
.filter(model.JobToOutputDatasetAssociation.job_id <= min(end_job_id, offset_start + args.batch_size)) \
.all()
job_to_output_hda_ids = sa_session.query(model.JobToOutputDatasetAssociation.job_id, model.JobToOutputDatasetAssociation.dataset_id,
model.JobToOutputDatasetAssociation.name) \
.filter(model.JobToOutputDatasetAssociation.job_id > offset_start) \
.filter(model.JobToOutputDatasetAssociation.job_id <= min(end_job_id, offset_start + args.batch_size)) \
.all()
# add type and concat
job_to_hda_ids = [[list(i), "input"] for i in job_to_input_hda_ids] + [[list(i), "output"] for i in job_to_output_hda_ids]
# add type and concat
job_to_hda_ids = [[list(i), "input"] for i in job_to_input_hda_ids] + [[list(i), "output"] for i in job_to_output_hda_ids]
# put all of the hda_ids into a list
hda_ids = [i[0][1] for i in job_to_hda_ids]
# put all of the hda_ids into a list
hda_ids = [i[0][1] for i in job_to_hda_ids]
hdas = sa_session.query(model.HistoryDatasetAssociation.id, model.HistoryDatasetAssociation.dataset_id,
model.HistoryDatasetAssociation.extension) \
.filter(model.HistoryDatasetAssociation.id.in_(hda_ids)) \
.all()
hdas = sa_session.query(model.HistoryDatasetAssociation.id, model.HistoryDatasetAssociation.dataset_id,
model.HistoryDatasetAssociation.extension) \
.filter(model.HistoryDatasetAssociation.id.in_(hda_ids)) \
.all()
# put all the dataset ids into a list
dataset_ids = [i[1] for i in hdas]
# put all the dataset ids into a list
dataset_ids = [i[1] for i in hdas]
# get the sizes of the datasets
datasets = sa_session.query(model.Dataset.id, model.Dataset.total_size) \
.filter(model.Dataset.id.in_(dataset_ids)) \
.all()
# get the sizes of the datasets
datasets = sa_session.query(model.Dataset.id, model.Dataset.total_size) \
.filter(model.Dataset.id.in_(dataset_ids)) \
.all()
# datasets to dictionay for easy search
hdas = {i[0]: i[1:] for i in hdas}
datasets = {i[0]: i[1:] for i in datasets}
# datasets to dictionay for easy search
hdas = {i[0]: i[1:] for i in hdas}
datasets = {i[0]: i[1:] for i in datasets}
for job_to_hda in job_to_hda_ids:
for job_to_hda in job_to_hda_ids:
job = job_to_hda[0] # job_id, hda_id, name
filetype = job_to_hda[1] # input|output
job = job_to_hda[0] # job_id, hda_id, name
filetype = job_to_hda[1] # input|output
# No associated job
if job[0] not in job_tool_map:
continue
# No associated job
if job[0] not in job_tool_map:
continue
# If the tool is blacklisted, exclude everywhere
if job_tool_map[job[0]] in blacklisted_tools:
continue
# If the tool is blacklisted, exclude everywhere
if job_tool_map[job[0]] in blacklisted_tools:
continue
hda_id = job[1]
if hda_id is None:
continue
hda_id = job[1]
if hda_id is None:
continue
dataset_id = hdas[hda_id][0]
if dataset_id is None:
continue
dataset_id = hdas[hda_id][0]
if dataset_id is None:
continue
try:
line = [
str(job[0]), # Job ID
str(hda_id), # HDA ID
str(hdas[hda_id][1]), # Extension
round_to_2sd(datasets[dataset_id][0]), # File size
job[2], # Parameter name
str(filetype) # input/output
]
cline = unicodify('\t'.join(line) + '\n')
handle_datasets.write(cline)
except Exception:
logging.warning("Unable to write out a 'handle_datasets' row. Ignoring the row.", exc_info=True)
continue
handle_datasets.close()
try:
line = [
str(job[0]), # Job ID
str(hda_id), # HDA ID
str(hdas[hda_id][1]), # Extension
round_to_2sd(datasets[dataset_id][0]), # File size
job[2], # Parameter name
str(filetype) # input/output
]
cline = unicodify('\t'.join(line) + '\n')
handle_datasets.write(cline)
except Exception:
logging.warning("Unable to write out a 'handle_datasets' row. Ignoring the row.", exc_info=True)
continue
annotate('export_datasets_end')
annotate('export_metric_num_start', 'Exporting Metrics (Numeric)')
handle_metric_num = io.open(REPORT_BASE + '.metric_num.tsv', 'w', encoding='utf-8')
handle_metric_num.write('\t'.join(('job_id', 'plugin', 'name', 'value')) + '\n')
for offset_start in range(last_job_sent, end_job_id, args.batch_size):
logging.debug("Processing %s:%s", offset_start, min(end_job_id, offset_start + args.batch_size))
for metric in sa_session.query(model.JobMetricNumeric.job_id, model.JobMetricNumeric.plugin, model.JobMetricNumeric.metric_name, model.JobMetricNumeric.metric_value) \
.filter(model.JobMetricNumeric.job_id > offset_start) \
.filter(model.JobMetricNumeric.job_id <= min(end_job_id, offset_start + args.batch_size)) \
.all():
# No associated job
if metric[0] not in job_tool_map:
continue
# If the tool is blacklisted, exclude everywhere
if job_tool_map[metric[0]] in blacklisted_tools:
continue
with io.open(REPORT_BASE + '.metric_num.tsv', 'w', encoding='utf-8') as handle_metric_num:
handle_metric_num.write(u'\t'.join(('job_id', 'plugin', 'name', 'value')) + '\n')
for offset_start in range(last_job_sent, end_job_id, args.batch_size):
logging.debug("Processing %s:%s", offset_start, min(end_job_id, offset_start + args.batch_size))
for metric in sa_session.query(model.JobMetricNumeric.job_id, model.JobMetricNumeric.plugin, model.JobMetricNumeric.metric_name, model.JobMetricNumeric.metric_value) \
.filter(model.JobMetricNumeric.job_id > offset_start) \
.filter(model.JobMetricNumeric.job_id <= min(end_job_id, offset_start + args.batch_size)) \
.all():
# No associated job
if metric[0] not in job_tool_map:
continue
# If the tool is blacklisted, exclude everywhere
if job_tool_map[metric[0]] in blacklisted_tools:
continue
try:
line = [
str(metric[0]), # job id
metric[1], # plugin
metric[2], # name
str(metric[3]) # value
]
try:
line = [
str(metric[0]), # job id
metric[1], # plugin
metric[2], # name
str(metric[3]) # value
]
cline = unicodify('\t'.join(line) + '\n')
handle_metric_num.write(cline)
except Exception:
logging.warning("Unable to write out a 'handle_metric_num' row. Ignoring the row.", exc_info=True)
continue
handle_metric_num.close()
cline = unicodify('\t'.join(line) + '\n')
handle_metric_num.write(cline)
except Exception:
logging.warning("Unable to write out a 'handle_metric_num' row. Ignoring the row.", exc_info=True)
continue
annotate('export_metric_num_end')
# Now on to outputs.
@@ -6,7 +6,7 @@ from time import time
import uwsgi
realtime_db_file = uwsgi.opt["interactivetools_map"]
realtime_db_file = uwsgi.opt["interactivetools_map"].decode('utf-8')
db_conn = sqlite3.connect(realtime_db_file)
DATABASE_TABLE_NAME = 'gxitproxy'
@@ -51,6 +51,14 @@ class CacheList():
key_type_token_mapped_cache = CacheList()
def args_as_unicode(func):
def wrap_args(*args):
args = (arg.decode('utf-8') if isinstance(arg, bytes) else arg for arg in args)
return func(*args)
return wrap_args
@args_as_unicode
def key_type_token_mapper_cached(key, key_type, token, route_extra, url, ttl):
global key_type_token_mapped_cache
cache_key = (key, key_type, token)
@@ -58,11 +66,12 @@ def key_type_token_mapper_cached(key, key_type, token, route_extra, url, ttl):
if entry is None:
entry = key_type_token_mapper(key, key_type, token, route_extra, url)
if entry is not None:
# Should we cache empt/not authorized entries, perhaps for shorter time?
# Should we cache empty/not authorized entries, perhaps for shorter time?
key_type_token_mapped_cache.add_entry(cache_key, entry, ttl=float(ttl))
return entry
@args_as_unicode
def key_type_token_mapper(key, key_type, token, route_extra, url):
global db_conn
# print 'key %s key_type %s token %s route_extra %s url %s\n' % (key, key_type, token, route_extra, url)
+9 -5
View File
@@ -3,6 +3,7 @@ import collections
import os
import string
import subprocess
import sys
import tempfile
import unittest
@@ -11,11 +12,9 @@ from galaxy_test.base.ssh_util import generate_ssh_keys
from galaxy_test.driver import integration_util
from .test_job_environments import BaseJobEnvironmentIntegrationTestCase
RemoteConnection = collections.namedtuple('remote_connection', ['hostname', 'username', 'password', 'port', 'private_key', 'public_key'])
RemoteConnection = collections.namedtuple('remote_connection', ['hostname', 'username', 'port', 'private_key', 'public_key'])
@integration_util.skip_unless_docker()
def start_ssh_docker(container_name, jobs_directory, port=10022, image='agaveapi/slurm'):
ssh_keys = generate_ssh_keys()
START_SLURM_DOCKER = ['docker',
@@ -28,6 +27,7 @@ def start_ssh_docker(container_name, jobs_directory, port=10022, image='agaveapi
'--name',
container_name,
'--rm',
'--privileged', # for torque
'-v',
"{jobs_directory}:{jobs_directory}".format(jobs_directory=jobs_directory),
"-v",
@@ -36,7 +36,9 @@ def start_ssh_docker(container_name, jobs_directory, port=10022, image='agaveapi
'nofile=2048:2048',
image]
subprocess.check_call(START_SLURM_DOCKER)
return RemoteConnection('localhost', 'testuser', 'testuser', port, ssh_keys.private_key_file, ssh_keys.public_key_file)
if sys.platform != 'darwin':
subprocess.check_call(['docker', 'exec', container_name, 'usermod', '-u', str(os.getuid()), 'testuser'])
return RemoteConnection('localhost', 'testuser', port, ssh_keys.private_key_file, ssh_keys.public_key_file)
def stop_ssh_docker(container_name, remote_connection):
@@ -58,6 +60,7 @@ def cli_job_config(remote_connection, shell_plugin='ParamikoShell', job_plugin='
<param id="shell_private_key">$private_key</param>
<param id="shell_hostname">$hostname</param>
<param id="shell_port">$port</param>
<param id="shell_strict_host_key_checking">False</param>
<param id="embed_metadata_in_job">False</param>
<env id="SOME_ENV_VAR">42</env>
</destination>
@@ -72,6 +75,7 @@ def cli_job_config(remote_connection, shell_plugin='ParamikoShell', job_plugin='
return job_conf.name
@integration_util.skip_unless_docker()
class BaseCliIntegrationTestCase(BaseJobEnvironmentIntegrationTestCase):
@classmethod
@@ -91,7 +95,7 @@ class BaseCliIntegrationTestCase(BaseJobEnvironmentIntegrationTestCase):
super(BaseCliIntegrationTestCase, cls).tearDownClass()
@classmethod
def handle_galaxy_config_kwds(cls, config, ):
def handle_galaxy_config_kwds(cls, config):
config["jobs_directory"] = cls.jobs_directory
config["file_path"] = cls.jobs_directory
config["job_config_file"] = cli_job_config(remote_connection=cls.remote_connection,
+16 -8
View File
@@ -1,5 +1,6 @@
"""Integration tests for running tools in Docker containers."""
import json
import os
import unittest
@@ -37,6 +38,11 @@ class ContainerizedIntegrationTestCase(integration_util.IntegrationTestCase):
skip_if_container_type_unavailable(cls)
super(ContainerizedIntegrationTestCase, cls).setUpClass()
@classmethod
def handle_galaxy_config_kwds(cls, config):
config["job_config_file"] = DOCKERIZED_JOB_CONFIG_FILE
disable_dependency_resolution(config)
def disable_dependency_resolution(config):
# Disable tool dependency resolution.
@@ -56,7 +62,6 @@ class DockerizedJobsIntegrationTestCase(integration_util.IntegrationTestCase, Ru
job_config_file = DOCKERIZED_JOB_CONFIG_FILE
build_mulled_resolver = 'build_mulled'
container_type = 'docker'
default_container_home_dir = '/'
@classmethod
def handle_galaxy_config_kwds(cls, config):
@@ -100,19 +105,24 @@ class DockerizedJobsIntegrationTestCase(integration_util.IntegrationTestCase, Ru
assert job_env.pwd.endswith("/working")
# Should we change env_pass_through to just always include TMP and HOME for docker?
# I'm not sure, if yes this would change.
assert job_env.home == self.default_container_home_dir, job_env.home
assert not job_env.home.endswith('/home')
def test_build_mulled(self):
if not which('docker'):
raise unittest.SkipTest("Docker not found on PATH, required for building images via involucro")
resolver_type = self.build_mulled_resolver
tool_id = 'mulled_example_multi_1'
endpoint = "tools/%s/dependencies" % tool_id
data = {'id': tool_id, 'resolver_type': resolver_type}
tool_ids = ['mulled_example_multi_1']
endpoint = "dependency_resolvers/toolbox/install"
data = {'tool_ids': json.dumps(tool_ids), 'resolver_type': resolver_type, 'container_type': self.container_type, 'include_containers': True}
create_response = self._post(endpoint, data=data, admin=True)
self._assert_status_code_is(create_response, 200)
create_response = self._get("dependency_resolvers/toolbox", data={'tool_ids': tool_ids, 'container_type': self.container_type, 'include_containers': True, 'index_by': 'tools'}, admin=True)
response = create_response.json()
assert any([True for d in response if d['dependency_type'] == self.container_type])
assert len(response) == 1
status = response[0]['status']
assert status[0]['model_class'] == 'ContainerDependency'
assert status[0]['dependency_type'] == self.container_type
assert status[0]['container_description']['identifier'].startswith('quay.io/local/mulled-v2-')
class MappingContainerResolverTestCase(integration_util.IntegrationTestCase):
@@ -164,5 +174,3 @@ class SingularityJobsIntegrationTestCase(DockerizedJobsIntegrationTestCase):
job_config_file = SINGULARITY_JOB_CONFIG_FILE
build_mulled_resolver = 'build_mulled_singularity'
container_type = 'singularity'
# singularity passes $HOME by default
default_container_home_dir = os.environ.get('HOME', '/')
@@ -47,8 +47,7 @@ class InteractiveToolsIntegrationTestCase(ContainerizedIntegrationTestCase):
assert isinstance(jobs, list)
assert len(jobs) == 1
job0 = jobs[0]
entry_points = self.wait_on_entry_points_active(job0["id"])
assert len(entry_points) == 2
entry_points = self.wait_on_entry_points_active(job0["id"], expected_num=2)
entry_point0 = entry_points[0]
entry_point1 = entry_points[1]
target0 = self.entry_point_target(entry_point0["id"])
@@ -59,19 +58,18 @@ class InteractiveToolsIntegrationTestCase(ContainerizedIntegrationTestCase):
content1 = self.wait_on_proxied_content(target1)
assert content1 == "moo cow\n", content1
assert False
def wait_on_proxied_content(self, target):
def get_hosted_content():
try:
scheme, rest = target.split("://", 1)
prefix, host_and_port = rest.split(".realtime.")
print(rest)
prefix, host_and_port = rest.split(".interactivetool.")
faked_host = rest
if "/" in rest:
faked_host = rest.split("/", 1)[0]
response = requests.get("%s://%s" % (scheme, host_and_port), timeout=1, headers={"Host": faked_host})
return response.content
url = "%s://%s" % (scheme, host_and_port)
response = requests.get(url, timeout=1, headers={"Host": faked_host})
return response.text
except Exception as e:
print(e)
return None
+11 -5
View File
@@ -1,5 +1,6 @@
"""Integration tests for the Kubernetes runner."""
# Tested on docker for mac 18.06.1-ce-mac73 using the default kubernetes setup
# Tested on docker for mac 18.06.1-ce-mac73 using the default kubernetes setup,
# also works on minikube
import collections
import json
import os
@@ -91,37 +92,40 @@ def job_config(jobs_directory):
<param id="k8s_persistent_volume_claims">jobs-directory-claim:$jobs_directory,tool-directory-claim:$tool_directory</param>
<param id="k8s_config_path">$k8s_config_path</param>
<param id="k8s_galaxy_instance_id">gx-short-id</param>
<param id="k8s_run_as_user_id">$$uid</param>
</plugin>
<plugin id="k8s_walltime_short" type="runner" load="galaxy.jobs.runners.kubernetes:KubernetesJobRunner">
<param id="k8s_persistent_volume_claims">jobs-directory-claim:$jobs_directory,tool-directory-claim:$tool_directory</param>
<param id="k8s_config_path">$k8s_config_path</param>
<param id="k8s_galaxy_instance_id">gx-short-id</param>
<param id="k8s_walltime_limit">10</param>
<param id="k8s_run_as_user_id">$$uid</param>
</plugin>
<plugin id="k8s_no_cleanup" type="runner" load="galaxy.jobs.runners.kubernetes:KubernetesJobRunner">
<param id="k8s_persistent_volume_claims">jobs-directory-claim:$jobs_directory,tool-directory-claim:$tool_directory</param>
<param id="k8s_config_path">$k8s_config_path</param>
<param id="k8s_galaxy_instance_id">gx-short-id</param>
<param id="k8s_cleanup_job">never</param>
<param id="k8s_run_as_user_id">$$uid</param>
</plugin>
</plugins>
<destinations default="k8s_destination">
<destination id="k8s_destination" runner="k8s">
<param id="limits_cpu">1.9</param>
<param id="limits_cpu">1.1</param>
<param id="limits_memory">10M</param>
<param id="docker_enabled">true</param>
<param id="docker_default_container_id">busybox:ubuntu-14.04</param>
<env id="SOME_ENV_VAR">42</env>
</destination>
<destination id="k8s_destination_walltime_short" runner="k8s_walltime_short">
<param id="limits_cpu">1.9</param>
<param id="limits_cpu">1.1</param>
<param id="limits_memory">10M</param>
<param id="docker_enabled">true</param>
<param id="docker_default_container_id">busybox:ubuntu-14.04</param>
<env id="SOME_ENV_VAR">42</env>
</destination>
<destination id="k8s_destination_no_cleanup" runner="k8s_no_cleanup">
<param id="limits_cpu">1.9</param>
<param id="limits_cpu">1.1</param>
<param id="limits_memory">10M</param>
<param id="docker_enabled">true</param>
<param id="docker_default_container_id">busybox:ubuntu-14.04</param>
@@ -182,7 +186,9 @@ class BaseKubernetesIntegrationTestCase(BaseJobEnvironmentIntegrationTestCase, M
super(BaseKubernetesIntegrationTestCase, cls).tearDownClass()
@classmethod
def handle_galaxy_config_kwds(cls, config, ):
def handle_galaxy_config_kwds(cls, config):
# TODO: implement metadata setting as separate job, as service or side-car
config['retry_metadata_internally'] = True
config["jobs_directory"] = cls.jobs_directory
config["file_path"] = cls.jobs_directory
config["job_config_file"] = cls.job_config.path
+11 -4
View File
@@ -116,7 +116,6 @@ class ScriptsIntegrationTestCase(integration_util.IntegrationTestCase):
output = self._scripts_check_output(script, ["-c", config_file])
assert "Complete" in output
@integration_util.skip_if_jenkins
def test_grt_export(self):
script = "grt/export.py"
self._scripts_check_argparse_help(script)
@@ -134,7 +133,7 @@ class ScriptsIntegrationTestCase(integration_util.IntegrationTestCase):
json_file = os.path.join(self.config_dir, json_files[0])
with open(json_file, "r") as f:
export = json.load(f)
assert export["version"] == 2
assert export["version"] == 3
def test_admin_cleanup_datasets(self):
self._scripts_check_argparse_help("cleanup_datasets/admin_cleanup_datasets.py")
@@ -179,13 +178,21 @@ class ScriptsIntegrationTestCase(integration_util.IntegrationTestCase):
clean_env = {
"PATH": os.environ.get("PATH", None),
} # Don't let testing environment variables interfere with config.
return unicodify(subprocess.check_output(cmd, cwd=cwd, env=clean_env))
try:
return unicodify(subprocess.check_output(cmd, cwd=cwd, env=clean_env))
except Exception as e:
if isinstance(e, subprocess.CalledProcessError):
raise Exception("%s\nOutput was:\n%s" % (unicodify(e), e.output))
raise
def write_config_file(self):
config_dir = self.config_dir
path = os.path.join(config_dir, "galaxy.yml")
self._test_driver.temp_directories.extend([config_dir])
config = self._raw_config
# Update config dict with database_connection, which might be set through env variables
config['database_connection'] = self._app.config.database_connection
with open(path, "w") as f:
yaml.dump({"galaxy": self._raw_config}, f)
yaml.dump({"galaxy": config}, f)
return path
+1
View File
@@ -23,6 +23,7 @@ class TestCliInterface(unittest.TestCase):
cls.username = 'testuser'
cls.shell_params = {'username': cls.username,
'private_key': cls.ssh_keys.private_key_file,
'strict_host_key_checking': False,
'hostname': 'localhost'}
cls.cli_interface = CliInterface()
+14 -14
View File
@@ -4,47 +4,47 @@ from contextlib import contextmanager
from os import path
from shutil import rmtree
import pytest
from galaxy.tools.toolbox import watcher
from galaxy.util import bunch
@pytest.mark.skipif(not watcher.can_watch, reason="watchdog not available")
def test_watcher():
if not watcher.can_watch:
from nose.plugins.skip import SkipTest
raise SkipTest()
with __test_directory() as t:
tool_path = path.join(t, "test.xml")
toolbox = Toolbox()
open(tool_path, "w").write("a")
with open(tool_path, "w") as f:
f.write("a")
tool_watcher = watcher.get_tool_watcher(toolbox, bunch.Bunch(
watch_tools=True
))
tool_watcher.start()
time.sleep(1)
tool_watcher.watch_file(tool_path, "cool_tool")
time.sleep(2)
assert not toolbox.was_reloaded("cool_tool")
open(tool_path, "w").write("b")
with open(tool_path, "w") as f:
f.write("b")
wait_for_reload(lambda: toolbox.was_reloaded("cool_tool"))
tool_watcher.shutdown()
assert tool_watcher.observer is None
@pytest.mark.skipif(not watcher.can_watch, reason="watchdog not available")
def test_tool_conf_watcher():
if not watcher.can_watch:
from nose.plugins.skip import SkipTest
raise SkipTest()
callback = CallbackRecorder()
conf_watcher = watcher.get_tool_conf_watcher(callback.call)
conf_watcher.start()
with __test_directory() as t:
tool_conf_path = path.join(t, "test_conf.xml")
open(tool_conf_path, "w").write("a")
with open(tool_conf_path, "w") as f:
f.write("a")
conf_watcher.watch_file(tool_conf_path)
time.sleep(1)
open(tool_conf_path, "w").write("b")
time.sleep(2)
with open(tool_conf_path, "w") as f:
f.write("b")
wait_for_reload(lambda: callback.called)
conf_watcher.shutdown()
assert conf_watcher.thread is None