diff --git a/display_applications/igv/vcf.xml b/display_applications/igv/vcf.xml
index f2c1e07966a..3c66ffd677a 100644
--- a/display_applications/igv/vcf.xml
+++ b/display_applications/igv/vcf.xml
@@ -12,8 +12,8 @@
${$site_id.startswith( 'local_' ) or $dataset.dbkey in $site_dbkeys}
${redirect_url}
-
-
+
+
#if ($dataset.dbkey in $site_dbkeys)
$site_organisms[ $site_dbkeys.index( $bgzip_file.dbkey ) ]
@@ -94,8 +94,8 @@
${ $dataset.dbkey == $value }
http://www.broadinstitute.org/igv/projects/current/igv.php?sessionURL=${bgzip_file.qp}&genome=${qp($bgzip_file.dbkey)}&merge=true&name=${qp( ( $bgzip_file.name or $DATASET_HASH ).replace( ',', ';' ) )}
-
-
+
+
diff --git a/lib/galaxy/config/__init__.py b/lib/galaxy/config/__init__.py
index bc28f5e0e01..b3a88ff6776 100644
--- a/lib/galaxy/config/__init__.py
+++ b/lib/galaxy/config/__init__.py
@@ -671,8 +671,8 @@ class GalaxyAppConfiguration(BaseAppConfiguration, CommonConfigurationMixin):
# InteractiveTools propagator mapping file
self.interactivetools_map = self._in_root_dir(kwargs.get("interactivetools_map", self._in_data_dir("interactivetools_map.sqlite")))
- self.interactivetool_prefix = kwargs.get("interactivetools_prefix", "interactivetool")
- self.interactivetool_proxy_host = kwargs.get("interactivetool_proxy_host", None)
+ self.interactivetools_prefix = kwargs.get("interactivetools_prefix", "interactivetool")
+ self.interactivetools_proxy_host = kwargs.get("interactivetool_proxy_host", None)
self.containers_conf = parse_containers_config(self.containers_config_file)
diff --git a/lib/galaxy/dependencies/pipfiles/default/pinned-requirements.txt b/lib/galaxy/dependencies/pipfiles/default/pinned-requirements.txt
index 948a30e797e..049894aa2ab 100644
--- a/lib/galaxy/dependencies/pipfiles/default/pinned-requirements.txt
+++ b/lib/galaxy/dependencies/pipfiles/default/pinned-requirements.txt
@@ -122,7 +122,7 @@ pbr==5.4.5
prettytable==0.7.2
prov==1.5.1
psutil==5.7.0
-pulsar-galaxy-lib==0.14.0.dev1
+pulsar-galaxy-lib==0.14.0.dev3
pyasn1-modules==0.2.8
pyasn1==0.4.8
pycparser==2.20
diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py
index 2451197b3aa..7d4907acc68 100644
--- a/lib/galaxy/jobs/__init__.py
+++ b/lib/galaxy/jobs/__init__.py
@@ -1141,6 +1141,14 @@ class JobWrapper(HasResourceParameters):
raise Exception('(%s) Unable to create job working directory',
job.id)
+ @property
+ def guest_ports(self):
+ if hasattr(self, "interactivetools"):
+ guest_ports = [ep.get('port') for ep in self.interactivetools]
+ return guest_ports
+ else:
+ return []
+
@property
def working_directory(self):
if self.__working_directory is None:
diff --git a/lib/galaxy/jobs/runners/__init__.py b/lib/galaxy/jobs/runners/__init__.py
index f5023efc1dd..26a75c19065 100644
--- a/lib/galaxy/jobs/runners/__init__.py
+++ b/lib/galaxy/jobs/runners/__init__.py
@@ -425,7 +425,7 @@ class BaseJobRunner(object):
compute_tmp_directory = job_wrapper.tmp_directory()
tool = job_wrapper.tool
- guest_ports = [ep.get('port') for ep in getattr(job_wrapper, 'interactivetools', [])]
+ guest_ports = job_wrapper.guest_ports
tool_info = ToolInfo(
tool.containers,
tool.requirements,
diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py
index ad68c60db77..aaac45ebf11 100644
--- a/lib/galaxy/jobs/runners/kubernetes.py
+++ b/lib/galaxy/jobs/runners/kubernetes.py
@@ -19,12 +19,15 @@ from galaxy.jobs.runners import (
from galaxy.jobs.runners.util.pykube_util import (
DEFAULT_JOB_API_VERSION,
ensure_pykube,
+ find_job_object_by_name,
+ galaxy_instance_id,
Job,
job_object_dict,
Pod,
produce_unique_k8s_job_name,
pull_policy,
pykube_client_from_dict,
+ stop_job,
)
from galaxy.util.bytesize import ByteSize
@@ -201,25 +204,8 @@ class KubernetesJobRunner(AsynchronousJobRunner):
return None
def __get_galaxy_instance_id(self):
- """
- Gets the id of the Galaxy instance. This will be added to Jobs and Pods names, so it needs to be DNS friendly,
- this means: `The Internet standards (Requests for Comments) for protocols mandate that component hostname labels
- may contain only the ASCII letters 'a' through 'z' (in a case-insensitive manner), the digits '0' through '9',
- and the minus sign ('-').`
-
- It looks for the value set on self.runner_params['k8s_galaxy_instance_id'], which might or not be set. The
- idea behind this is to allow the Galaxy instance to trust (or not) existing k8s Jobs and Pods that match the
- setup of a Job that is being recovered or restarted after a downtime/reboot.
- :return:
- :rtype:
- """
- if "k8s_galaxy_instance_id" in self.runner_params:
- if re.match(r"(?!-)[a-z\d-]{1,20}(?20 characters) or it includes non DNS acceptable characters, ignoring it.')
- return None
+ """Parse the ID of the Galaxy instance from runner params."""
+ return galaxy_instance_id(self.runner_params)
def __produce_unique_k8s_job_name(self, galaxy_internal_job_id):
# wrapper.get_id_tag() instead of job_id for compatibility with TaskWrappers.
@@ -509,19 +495,7 @@ class KubernetesJobRunner(AsynchronousJobRunner):
def __cleanup_k8s_job(self, job):
k8s_cleanup_job = self.runner_params['k8s_cleanup_job']
- job_failed = (job.obj['status']['failed'] > 0
- if 'failed' in job.obj['status'] else False)
- # Scale down the job just in case even if cleanup is never
- job.scale(replicas=0)
- if (k8s_cleanup_job == "always" or
- (k8s_cleanup_job == "onsuccess" and not job_failed)):
- delete_options = {
- "apiVersion": "v1",
- "kind": "DeleteOptions",
- "propagationPolicy": "Background"
- }
- r = job.api.delete(json=delete_options, **job.api_kwargs())
- job.api.raise_for_status(r)
+ stop_job(job, k8s_cleanup_job)
def __job_failed_due_to_walltime_limit(self, job):
conditions = job.obj['status'].get('conditions') or []
@@ -550,11 +524,10 @@ class KubernetesJobRunner(AsynchronousJobRunner):
"""Attempts to delete a dispatched job to the k8s cluster"""
job = job_wrapper.get_job()
try:
- jobs = Job.objects(self._pykube_api).filter(
- selector="app=" + self.__produce_unique_k8s_job_name(job.get_id_tag()),
- namespace=self.runner_params['k8s_namespace'])
- if len(jobs.response['items']) > 0:
- job_to_delete = Job(self._pykube_api, jobs.response['items'][0])
+ name = self.__produce_unique_k8s_job_name(job.get_id_tag())
+ namespace = self.runner_params['k8s_namespace']
+ job_to_delete = find_job_object_by_name(self._pykube_api, name, namespace)
+ if job_to_delete:
self.__cleanup_k8s_job(job_to_delete)
# TODO assert whether job parallelism == 0
# assert not job_to_delete.exists(), "Could not delete job,"+job.job_runner_external_id+" it still exists"
diff --git a/lib/galaxy/jobs/runners/pulsar.py b/lib/galaxy/jobs/runners/pulsar.py
index ebd72250e58..68befa97b86 100644
--- a/lib/galaxy/jobs/runners/pulsar.py
+++ b/lib/galaxy/jobs/runners/pulsar.py
@@ -189,6 +189,7 @@ class PulsarJobRunner(AsynchronousJobRunner):
runner_name = "PulsarJobRunner"
default_build_pulsar_app = False
use_mq = False
+ poll = True
def __init__(self, app, nworkers, **kwds):
"""Start the job runner."""
@@ -207,12 +208,13 @@ class PulsarJobRunner(AsynchronousJobRunner):
if self.use_mq:
# This is a message queue driven runner, don't monitor
# just setup required callback.
- self._init_noop_monitor()
-
self.client_manager.ensure_has_status_update_callback(self.__async_update)
self.client_manager.ensure_has_ack_consumers()
- else:
+
+ if self.poll:
self._init_monitor_thread()
+ else:
+ self._init_noop_monitor()
def __init_client_manager(self):
pulsar_conf = self.runner_params.get('pulsar_app_config', None)
@@ -262,6 +264,23 @@ class PulsarJobRunner(AsynchronousJobRunner):
return JobDestination(runner="pulsar", params=url_to_destination_params(url))
def check_watched_item(self, job_state):
+ if self.use_mq:
+ # Might still need to check pod IPs.
+ job_wrapper = job_state.job_wrapper
+ guest_ports = job_wrapper.guest_ports
+ if len(guest_ports) > 0:
+ client = self.get_client_from_state(job_state)
+ job_ip = client.job_ip()
+ if job_ip:
+ ports_dict = {}
+ for guest_port in guest_ports:
+ ports_dict[str(guest_port)] = dict(host=job_ip, port=guest_port, protocol="http")
+ self.app.interactivetool_manager.configure_entry_points(job_wrapper.get_job(), ports_dict)
+ return job_state
+ else:
+ return self.check_watched_item_state(job_state)
+
+ def check_watched_item_state(self, job_state):
try:
client = self.get_client_from_state(job_state)
status = client.get_status()
@@ -366,6 +385,7 @@ class PulsarJobRunner(AsynchronousJobRunner):
remote_pulsar_app_config=remote_pulsar_app_config,
job_directory_files=job_directory_files,
container=None if not remote_container else remote_container.container_id,
+ guest_ports=job_wrapper.guest_ports,
)
job_id = pulsar_submit_job(client, client_job_description, remote_job_config)
log.info("Pulsar job submitted with job_id %s" % job_id)
@@ -853,6 +873,7 @@ class PulsarLegacyJobRunner(PulsarJobRunner):
class PulsarMQJobRunner(PulsarJobRunner):
"""Flavor of Pulsar job runner with sensible defaults for message queue communication."""
use_mq = True
+ poll = False
destination_defaults = dict(
default_file_action="remote_transfer",
@@ -868,7 +889,7 @@ KUBERNETES_DESTINATION_DEFAULTS = {
"default_file_action": "remote_transfer",
"rewrite_parameters": "true",
"jobs_directory": "/pulsar_staging",
- "pulsar_container_image": "galaxy/pulsar-pod-staging:0.13.0",
+ "pulsar_container_image": "galaxy/pulsar-pod-staging:0.14.0",
"remote_container_handling": True,
"k8s_enabled": True,
"url": PARAMETER_SPECIFICATION_IGNORED,
@@ -878,6 +899,7 @@ KUBERNETES_DESTINATION_DEFAULTS = {
class PulsarKubernetesJobRunner(PulsarMQJobRunner):
destination_defaults = KUBERNETES_DESTINATION_DEFAULTS
+ poll = True # Poll so we can check API for pod IP for ITs.
def _populate_parameter_defaults(self, job_destination):
super(PulsarKubernetesJobRunner, self)._populate_parameter_defaults(job_destination)
diff --git a/lib/galaxy/jobs/runners/util/pykube_util.py b/lib/galaxy/jobs/runners/util/pykube_util.py
index 0203b9f5046..f2e06f78c59 100644
--- a/lib/galaxy/jobs/runners/util/pykube_util.py
+++ b/lib/galaxy/jobs/runners/util/pykube_util.py
@@ -1,5 +1,7 @@
"""Interface layer for pykube library shared between Galaxy and Pulsar."""
+import logging
import os
+import re
import uuid
try:
@@ -17,8 +19,13 @@ except ImportError as exc:
'this feature, please install it or correct the '
'following error:\nImportError %s' % str(exc))
+log = logging.getLogger(__name__)
+
DEFAULT_JOB_API_VERSION = "batch/v1"
DEFAULT_NAMESPACE = "default"
+INSTANCE_ID_INVALID_MESSAGE = ("Galaxy instance [%s] is either too long "
+ "(>20 characters) or it includes non DNS "
+ "acceptable characters, ignoring it.")
def ensure_pykube():
@@ -61,6 +68,44 @@ def pull_policy(params):
return None
+def find_job_object_by_name(pykube_api, job_name, namespace=None):
+ return _find_object_by_name(Job, pykube_api, job_name, namespace=namespace)
+
+
+def find_pod_object_by_name(pykube_api, pod_name, namespace=None):
+ return _find_object_by_name(Pod, pykube_api, pod_name, namespace=namespace)
+
+
+def _find_object_by_name(clazz, pykube_api, object_name, namespace=None):
+ filter_kwd = dict(selector="app=%s" % object_name)
+ if namespace is not None:
+ filter_kwd["namespace"] = namespace
+
+ objs = clazz.objects(pykube_api).filter(**filter_kwd)
+ obj = None
+ if len(objs.response['items']) > 0:
+ obj = clazz(pykube_api, objs.response['items'][0])
+ return obj
+
+
+def stop_job(job, cleanup="always"):
+ job_failed = (job.obj['status']['failed'] > 0
+ if 'failed' in job.obj['status'] else False)
+ # Scale down the job just in case even if cleanup is never
+ job.scale(replicas=0)
+ api_delete = cleanup == "always"
+ if not api_delete and cleanup == "onsuccess" and not job_failed:
+ api_delete = True
+ if api_delete:
+ delete_options = {
+ "apiVersion": "v1",
+ "kind": "DeleteOptions",
+ "propagationPolicy": "Background"
+ }
+ r = job.api.delete(json=delete_options, **job.api_kwargs())
+ job.api.raise_for_status(r)
+
+
def job_object_dict(params, job_name, spec):
k8s_job_obj = {
"apiVersion": params.get('k8s_job_api_version', DEFAULT_JOB_API_VERSION),
@@ -77,13 +122,38 @@ def job_object_dict(params, job_name, spec):
return k8s_job_obj
+def galaxy_instance_id(params):
+ """Parse and validate the id of the Galaxy instance from supplied dict.
+
+ This will be added to Jobs and Pods names, so it needs to be DNS friendly,
+ this means: `The Internet standards (Requests for Comments) for protocols mandate that component hostname labels
+ may contain only the ASCII letters 'a' through 'z' (in a case-insensitive manner), the digits '0' through '9',
+ and the minus sign ('-').`
+
+ It looks for the value set on params['k8s_galaxy_instance_id'], which might or not be set. The
+ idea behind this is to allow the Galaxy instance to trust (or not) existing k8s Jobs and Pods that match the
+ setup of a Job that is being recovered or restarted after a downtime/reboot.
+ """
+ if "k8s_galaxy_instance_id" in params:
+ raw_value = params['k8s_galaxy_instance_id']
+ if re.match(r"(?!-)[a-z\d-]{1,20}(?