Merge branch 'release_20.09' into dev

This commit is contained in:
mvdbeek
2020-10-17 10:43:51 +02:00
36 changed files with 235 additions and 135 deletions
+5 -4
View File
@@ -96,10 +96,11 @@ class ItemGrabber:
HANDLER_ASSIGNMENT_METHODS.DB_TRANSACTION_ISOLATION,
HANDLER_ASSIGNMENT_METHODS.DB_SKIP_LOCKED,
}
try:
return [m for m in handler_assignment_methods if m in grabbable_methods][0]
except IndexError:
return
if handler_assignment_methods:
try:
return [m for m in handler_assignment_methods if m in grabbable_methods][0]
except IndexError:
return
def grab_unhandled_items(self):
"""
+38 -49
View File
@@ -6,7 +6,6 @@ import logging
import math
import os
import re
from time import sleep
import yaml
@@ -20,11 +19,12 @@ from galaxy.jobs.runners.util.pykube_util import (
DEFAULT_JOB_API_VERSION,
ensure_pykube,
find_job_object_by_name,
find_pod_object_by_name,
galaxy_instance_id,
Job,
job_object_dict,
Pod,
produce_unique_k8s_job_name,
produce_k8s_job_prefix,
pull_policy,
pykube_client_from_dict,
stop_job,
@@ -88,7 +88,10 @@ class KubernetesJobRunner(AsynchronousJobRunner):
self.setup_volumes()
def setup_volumes(self):
volume_claims = dict(volume.split(":") for volume in self.runner_params['k8s_persistent_volume_claims'].split(','))
if self.runner_params.get('k8s_persistent_volume_claims'):
volume_claims = dict(volume.split(":") for volume in self.runner_params['k8s_persistent_volume_claims'].split(','))
else:
volume_claims = {}
mountable_volumes = [{'name': claim_name, 'persistentVolumeClaim': {'claimName': claim_name}} for claim_name in volume_claims]
self.runner_params['k8s_mountable_volumes'] = mountable_volumes
volume_mounts = [{'name': claim_name, 'mountPath': mount_path} for claim_name, mount_path in volume_claims.items()]
@@ -120,45 +123,21 @@ class KubernetesJobRunner(AsynchronousJobRunner):
return
# Construction of the Kubernetes Job object follows: http://kubernetes.io/docs/user-guide/persistent-volumes/
k8s_job_name = self.__produce_unique_k8s_job_name(job_wrapper.get_id_tag())
k8s_job_prefix = self.__produce_k8s_job_prefix()
k8s_job_obj = job_object_dict(
self.runner_params,
k8s_job_name,
k8s_job_prefix,
self.__get_k8s_job_spec(ajs)
)
# Checks if job exists and is trusted, or if it needs re-creation.
job = Job(self._pykube_api, k8s_job_obj)
job_exists = job.exists()
if job_exists and not self._galaxy_instance_id:
# if galaxy instance id is not set, then we don't trust matching jobs and we simply delete and
# re-create the job
log.debug("Matching job exists, but Job is not trusted, so it will be deleted and a new one created.")
job.delete()
elapsed_seconds = 0
while job.exists():
sleep(3)
elapsed_seconds += 3
if elapsed_seconds > self.runner_params['k8s_timeout_seconds_job_deletion']:
log.debug(
"Timed out before k8s could delete existing untrusted job %s, not queuing associated Galaxy job."
% k8s_job_name)
return
log.debug("Waiting for job to be deleted " + k8s_job_name)
Job(self._pykube_api, k8s_job_obj).create()
elif job_exists and self._galaxy_instance_id:
# The job exists and we trust the identifier.
log.debug("Matching job exists, but Job is trusted, so we simply use the existing one for " + k8s_job_name)
# We simply leave the k8s job to be handled later on by check_watched_item().
else:
# Creates the Kubernetes Job if it doesn't exist.
job.create()
job.create()
job_id = job.metadata['name']
# define job attributes in the AsyncronousJobState for follow-up
ajs.job_id = k8s_job_name
ajs.job_id = job_id
# store runner information for tracking if Galaxy restarts
job_wrapper.set_external_id(k8s_job_name)
job_wrapper.set_external_id(job_id)
self.monitor_queue.put(ajs)
def __get_pull_policy(self):
@@ -206,10 +185,9 @@ class KubernetesJobRunner(AsynchronousJobRunner):
"""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.
def __produce_k8s_job_prefix(self):
instance_id = self._galaxy_instance_id or ''
return produce_unique_k8s_job_name(app_prefix='gxy', instance_id=instance_id, job_id=galaxy_internal_job_id)
return produce_k8s_job_prefix(app_prefix='gxy', instance_id=instance_id)
def __get_k8s_job_spec(self, ajs):
"""Creates the k8s Job spec. For a Job spec, the only requirement is to have a .spec.template.
@@ -227,7 +205,21 @@ class KubernetesJobRunner(AsynchronousJobRunner):
(see pod selector) and an appropriate restart policy."""
k8s_spec_template = {
"metadata": {
"labels": {"app": self.__produce_unique_k8s_job_name(ajs.job_wrapper.get_id_tag())[:-5]}
"labels": {
"app.kubernetes.io/name": ajs.job_wrapper.tool.old_id,
"app.kubernetes.io/instance": self.__produce_k8s_job_prefix(),
"app.kubernetes.io/version": ajs.job_wrapper.tool.version,
"app.kubernetes.io/component": "tool",
"app.kubernetes.io/part-of": "galaxy",
"app.kubernetes.io/managed-by": "galaxy",
"app.galaxyproject.org/job_id": ajs.job_wrapper.get_id_tag(),
"app.galaxyproject.org/instance": self._galaxy_instance_id or "",
"app.galaxyproject.org/handler": self.app.config.server_name,
"app.galaxyproject.org/destination": ajs.job_wrapper.job_destination.id,
},
"annotations": {
"app.galaxyproject.org/tool_id": ajs.job_wrapper.tool.id,
}
},
"spec": {
"volumes": self.runner_params['k8s_mountable_volumes'],
@@ -401,8 +393,8 @@ class KubernetesJobRunner(AsynchronousJobRunner):
def check_watched_item(self, job_state):
"""Checks the state of a job already submitted on k8s. Job state is an AsynchronousJobState"""
jobs = Job.objects(self._pykube_api).filter(selector="app=" + job_state.job_id,
namespace=self.runner_params['k8s_namespace'])
jobs = find_job_object_by_name(self._pykube_api, job_state.job_id, self.runner_params['k8s_namespace'])
if len(jobs.response['items']) == 1:
job = Job(self._pykube_api, jobs.response['items'][0])
job_destination = job_state.job_wrapper.job_destination
@@ -506,8 +498,7 @@ class KubernetesJobRunner(AsynchronousJobRunner):
marks the job for resubmission (resubmit logic is part of destinations).
"""
pods = Pod.objects(self._pykube_api).filter(selector="app=%s" % job_state.job_id,
namespace=self.runner_params['k8s_namespace'])
pods = find_pod_object_by_name(self._pykube_api, job_state.job_id, self.runner_params['k8s_namespace'])
if not pods.response['items']:
return False
@@ -522,17 +513,16 @@ class KubernetesJobRunner(AsynchronousJobRunner):
"""Attempts to delete a dispatched job to the k8s cluster"""
job = job_wrapper.get_job()
try:
name = job.job_runner_external_id
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)
job_to_delete = find_job_object_by_name(self._pykube_api, job.get_job_runner_external_id(), self.runner_params['k8s_namespace'])
if job_to_delete and len(job_to_delete.response['items']) > 0:
k8s_job = Job(self._pykube_api, job_to_delete.response['items'][0])
self.__cleanup_k8s_job(k8s_job)
# TODO assert whether job parallelism == 0
# assert not job_to_delete.exists(), "Could not delete job,"+job.job_runner_external_id+" it still exists"
log.debug(f"({job.id}/{job.job_runner_external_id}) Terminated at user's request")
except Exception as e:
log.exception("({}/{}) User killed running job, but error encountered during termination: {}".format(
job.id, job.job_runner_external_id, e))
job.id, job.get_job_runner_external_id(), e))
def recover(self, job, job_wrapper):
"""Recovers jobs stuck in the queued/running state when Galaxy started"""
@@ -561,8 +551,7 @@ class KubernetesJobRunner(AsynchronousJobRunner):
def finish_job(self, job_state):
super().finish_job(job_state)
jobs = Job.objects(self._pykube_api).filter(selector="app=" + job_state.job_id,
namespace=self.runner_params['k8s_namespace'])
jobs = find_job_object_by_name(self._pykube_api, job_state.job_id, self.runner_params['k8s_namespace'])
if len(jobs.response['items']) != 1:
log.warning("More than one job matches selector. Possible configuration error"
" in job id '%s'", job_state.job_id)
+11 -34
View File
@@ -2,7 +2,6 @@
import logging
import os
import re
import uuid
try:
from pykube.config import KubeConfig
@@ -46,18 +45,9 @@ def pykube_client_from_dict(params):
return pykube_client
def produce_unique_k8s_job_name(app_prefix=None, instance_id=None, job_id=None):
if job_id is None:
job_id = str(uuid.uuid4())
job_name = ""
if app_prefix:
job_name += "%s-" % app_prefix
if instance_id and len(instance_id) > 0:
job_name += "%s-" % instance_id
return "{}{}-{}".format(job_name, job_id, uuid.uuid4())[-63:]
def produce_k8s_job_prefix(app_prefix=None, instance_id=None):
job_name_elems = [app_prefix or "", instance_id or ""]
return '-'.join(elem for elem in job_name_elems if elem)
def pull_policy(params):
@@ -69,23 +59,13 @@ def pull_policy(params):
def find_job_object_by_name(pykube_api, job_name, namespace=None):
return _find_object_by_name(Job, pykube_api, job_name, namespace=namespace)
if not job_name:
raise ValueError("job name must not be empty")
return Job.objects(pykube_api).filter(field_selector={"metadata.name": 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 find_pod_object_by_name(pykube_api, job_name, namespace=None):
return Pod.objects(pykube_api).filter(selector="job-name=" + job_name, namespace=namespace)
def stop_job(job, cleanup="always"):
@@ -106,16 +86,13 @@ def stop_job(job, cleanup="always"):
job.api.raise_for_status(r)
def job_object_dict(params, job_name, spec):
def job_object_dict(params, job_prefix, spec):
k8s_job_obj = {
"apiVersion": params.get('k8s_job_api_version', DEFAULT_JOB_API_VERSION),
"kind": "Job",
"metadata": {
# metadata.name is the name of the pod resource created, and must be unique
# http://kubernetes.io/docs/user-guide/configuring-containers/
"name": job_name,
"generateName": job_prefix + "-",
"namespace": params.get('k8s_namespace', DEFAULT_NAMESPACE),
"labels": {"app": job_name}
},
"spec": spec,
}
@@ -152,7 +129,7 @@ __all__ = (
"Job",
"job_object_dict",
"Pod",
"produce_unique_k8s_job_name",
"produce_k8s_job_prefix",
"pull_policy",
"pykube_client_from_dict",
"stop_job",
+1 -1
View File
@@ -548,7 +548,7 @@ class WeblessApplicationStack(ApplicationStack):
for m in remove_methods:
try:
job_config.handler_assignment_methods.remove(m)
log.debug("%s: Removed '%s' from handler assignment methods due to use of mules", conf_class_name, m)
log.debug("%s: Removed '%s' from handler assignment methods due to use of --attach-to-pool", conf_class_name, m)
except ValueError:
pass
if add_method not in job_config.handler_assignment_methods:
+12 -13
View File
@@ -267,19 +267,18 @@ class WorkflowRequestMonitor(Monitors):
self.workflow_scheduling_manager = workflow_scheduling_manager
self._init_monitor_thread(name="WorkflowRequestMonitor.monitor_thread", target=self.__monitor, config=app.config)
self.invocation_grabber = None
if self.workflow_scheduling_manager.handler_assignment_methods_configured:
self_handler_tags = set(self.app.job_config.self_handler_tags)
self_handler_tags.add(self.workflow_scheduling_manager.default_handler_id)
handler_assignment_method = ItemGrabber.get_grabbable_handler_assignment_method(self.workflow_scheduling_manager.handler_assignment_methods)
if handler_assignment_method:
self.invocation_grabber = ItemGrabber(
app=app,
grab_type='WorkflowInvocation',
handler_assignment_method=handler_assignment_method,
max_grab=self.workflow_scheduling_manager.handler_max_grab,
self_handler_tags=self_handler_tags,
handler_tags=self_handler_tags,
)
self_handler_tags = set(self.app.job_config.self_handler_tags)
self_handler_tags.add(self.workflow_scheduling_manager.default_handler_id)
handler_assignment_method = ItemGrabber.get_grabbable_handler_assignment_method(self.workflow_scheduling_manager.handler_assignment_methods)
if handler_assignment_method:
self.invocation_grabber = ItemGrabber(
app=app,
grab_type='WorkflowInvocation',
handler_assignment_method=handler_assignment_method,
max_grab=self.workflow_scheduling_manager.handler_max_grab,
self_handler_tags=self_handler_tags,
handler_tags=self_handler_tags,
)
def __monitor(self):
to_monitor = self.workflow_scheduling_manager.active_workflow_schedulers
+7 -1
View File
@@ -6,11 +6,17 @@ History
.. to_doc
---------------------
20.5.1.dev0
20.9.1.dev0
---------------------
---------------------
20.9.0 (2020-10-15)
---------------------
* First release from the 20.09 branch of Galaxy.
---------------------
20.5.0 (2020-07-04)
---------------------
+3 -1
View File
@@ -1,4 +1,6 @@
__version__ = '20.5.1.dev0'
# -*- coding: utf-8 -*-
__version__ = '20.9.1.dev0'
PROJECT_NAME = "galaxy-app"
PROJECT_OWNER = PROJECT_USERAME = "galaxyproject"
+7 -1
View File
@@ -6,11 +6,17 @@ History
.. to_doc
---------------------
20.5.1.dev0
20.9.1.dev0
---------------------
---------------------
20.9.0 (2020-10-15)
---------------------
* First release from the 20.09 branch of Galaxy.
---------------------
20.5.0 (2020-07-04)
---------------------
+3 -1
View File
@@ -1,4 +1,6 @@
__version__ = '20.5.1.dev0'
# -*- coding: utf-8 -*-
__version__ = '20.9.1.dev0'
PROJECT_NAME = "galaxy-auth"
PROJECT_OWNER = PROJECT_USERAME = "galaxyproject"
+7 -1
View File
@@ -6,11 +6,17 @@ History
.. to_doc
---------------------
20.5.1.dev0
20.9.1.dev0
---------------------
---------------------
20.9.0 (2020-10-15)
---------------------
* First release from the 20.09 branch of Galaxy.
---------------------
20.5.0 (2020-07-04)
+3 -1
View File
@@ -1,4 +1,6 @@
__version__ = '20.5.1.dev0'
# -*- coding: utf-8 -*-
__version__ = '20.9.1.dev0'
PROJECT_NAME = "galaxy-data"
PROJECT_OWNER = PROJECT_USERAME = "galaxyproject"
+7 -1
View File
@@ -6,11 +6,17 @@ History
.. to_doc
---------------------
20.5.1.dev0
20.9.1.dev0
---------------------
---------------------
20.9.0 (2020-10-15)
---------------------
* First release from the 20.09 branch of Galaxy.
---------------------
20.5.0 (2020-07-04)
---------------------
@@ -1,4 +1,6 @@
__version__ = '20.5.1.dev0'
# -*- coding: utf-8 -*-
__version__ = '20.9.1.dev0'
PROJECT_NAME = "galaxy-job-execution"
PROJECT_OWNER = PROJECT_USERAME = "galaxyproject"
+7 -1
View File
@@ -6,11 +6,17 @@ History
.. to_doc
---------------------
20.5.1.dev0
20.9.1.dev0
---------------------
---------------------
20.9.0 (2020-10-15)
---------------------
* First release from the 20.09 branch of Galaxy.
---------------------
20.5.0 (2020-07-04)
---------------------
@@ -1,4 +1,6 @@
__version__ = '20.5.1.dev0'
# -*- coding: utf-8 -*-
__version__ = '20.9.1.dev0'
PROJECT_NAME = "galaxy-job-metrics"
PROJECT_OWNER = PROJECT_USERAME = "galaxyproject"
+7 -1
View File
@@ -6,11 +6,17 @@ History
.. to_doc
---------------------
20.5.1.dev0
20.9.1.dev0
---------------------
---------------------
20.9.0 (2020-10-15)
---------------------
* First release from the 20.09 branch of Galaxy.
---------------------
20.5.0 (2020-07-04)
---------------------
@@ -1,4 +1,6 @@
__version__ = '20.5.1.dev0'
# -*- coding: utf-8 -*-
__version__ = '20.9.1.dev0'
PROJECT_NAME = "galaxy-objectstore"
PROJECT_OWNER = PROJECT_USERAME = "galaxyproject"
+7 -1
View File
@@ -6,11 +6,17 @@ History
.. to_doc
---------------------
20.5.1.dev0
20.9.1.dev0
---------------------
---------------------
20.9.0 (2020-10-15)
---------------------
* First release from the 20.09 branch of Galaxy.
---------------------
20.5.0 (2020-07-04)
---------------------
@@ -1,4 +1,6 @@
__version__ = '20.5.1.dev0'
# -*- coding: utf-8 -*-
__version__ = '20.9.1.dev0'
PROJECT_NAME = "galaxy-selenium"
PROJECT_OWNER = PROJECT_USERAME = "galaxyproject"
+9 -2
View File
@@ -6,7 +6,14 @@ History
.. to_doc
---------------------
20.1.0.dev0
20.9.1.dev0
---------------------
* Initial import from dev branch of Galaxy during 20.01 development cycle.
---------------------
20.9.0 (2020-10-15)
---------------------
* Initial import from dev branch of Galaxy during 20.09 branch of Galaxy.
@@ -1,4 +1,6 @@
__version__ = '20.1.0.dev0'
# -*- coding: utf-8 -*-
__version__ = '20.9.1.dev0'
PROJECT_NAME = "galaxy-test-api"
PROJECT_OWNER = PROJECT_USERAME = "galaxyproject"
+1 -2
View File
@@ -1,2 +1 @@
galaxy-test-driver
galaxy-test-base
+7 -1
View File
@@ -6,11 +6,17 @@ History
.. to_doc
---------------------
20.5.1.dev0
20.9.1.dev0
---------------------
---------------------
20.9.0 (2020-10-15)
---------------------
* Initial release from the 20.09 branch of Galaxy.
---------------------
20.5.0 (2020-07-04)
---------------------
@@ -1,4 +1,6 @@
__version__ = '20.5.1.dev0'
# -*- coding: utf-8 -*-
__version__ = '20.9.1.dev0'
PROJECT_NAME = "galaxy-test-base"
PROJECT_OWNER = PROJECT_USERAME = "galaxyproject"
+7 -1
View File
@@ -6,11 +6,17 @@ History
.. to_doc
---------------------
20.5.1.dev0
20.9.1.dev0
---------------------
---------------------
20.9.0 (2020-10-15)
---------------------
* First release from the 20.09 branch of Galaxy.
---------------------
20.5.0 (2020-07-04)
---------------------
@@ -1,4 +1,6 @@
__version__ = '20.5.1.dev0'
# -*- coding: utf-8 -*-
__version__ = '20.9.1.dev0'
PROJECT_NAME = "galaxy-test-driver"
PROJECT_OWNER = PROJECT_USERAME = "galaxyproject"
+8 -2
View File
@@ -6,7 +6,13 @@ History
.. to_doc
---------------------
20.1.0.dev0
20.9.1.dev0
---------------------
* Initial import from dev branch of Galaxy during 20.01 development cycle.
---------------------
20.9.0 (2020-10-15)
---------------------
* First release from the 20.09 branch of Galaxy.
@@ -1,4 +1,6 @@
__version__ = '20.1.0.dev0'
# -*- coding: utf-8 -*-
__version__ = '20.9.1.dev0'
PROJECT_NAME = "galaxy-test-selenium"
PROJECT_OWNER = PROJECT_USERAME = "galaxyproject"
+2 -3
View File
@@ -1,4 +1,3 @@
galaxy-data
nose
pytest
galaxy-test-base
+7
View File
@@ -5,10 +5,17 @@ History
.. to_doc
---------------------
20.9.1.dev0
---------------------
---------------------
20.9.0.dev2
---------------------
* First release from the 20.09 branch of Galaxy.
---------------------
20.5.0 (2020-07-04)
@@ -1,4 +1,6 @@
__version__ = '20.9.0.dev3'
# -*- coding: utf-8 -*-
__version__ = '20.9.1.dev0'
PROJECT_NAME = "galaxy-tool-util"
PROJECT_OWNER = PROJECT_USERAME = "galaxyproject"
+7 -1
View File
@@ -6,11 +6,17 @@ History
.. to_doc
---------------------
20.5.1.dev0
20.9.1.dev0
---------------------
---------------------
20.9.0 (2020-10-15)
---------------------
* First release from the 20.09 branch of Galaxy.
---------------------
20.5.0 (2020-07-03)
---------------------
+3 -1
View File
@@ -1,4 +1,6 @@
__version__ = '20.5.1.dev0'
# -*- coding: utf-8 -*-
__version__ = '20.9.1.dev0'
PROJECT_NAME = "galaxy-util"
PROJECT_OWNER = PROJECT_USERAME = "galaxyproject"
+7 -1
View File
@@ -6,11 +6,17 @@ History
.. to_doc
---------------------
20.5.1.dev0
20.9.1.dev0
---------------------
---------------------
20.9.0 (2020-10-15)
---------------------
* First release from the 20.09 branch of Galaxy.
---------------------
20.5.0 (2020-07-04)
---------------------
@@ -1,4 +1,6 @@
__version__ = '20.5.1.dev0'
# -*- coding: utf-8 -*-
__version__ = '20.9.1.dev0'
PROJECT_NAME = "galaxy-web-framework"
PROJECT_OWNER = PROJECT_USERAME = "galaxyproject"
@@ -42,6 +42,19 @@ WORKFLOW_HANDLER_JOB_CONFIG_TEMPLATE = string.Template("""
</job_conf>
""")
POOL_JOB_CONFIG_TEMPLATE = string.Template("""<job_conf>
<plugins>
<plugin id="local" type="runner" load="galaxy.jobs.runners.local:LocalJobRunner" workers="2"/>
</plugins>
<handlers $assign_with>
</handlers>
<destinations default="local">
<destination id="local" runner="local">
</destination>
</destinations>
</job_conf>
""")
WORKFLOW_SCHEDULERS_CONFIG_TEMPLATE = string.Template("""
<workflow_schedulers default="core">
<core id="core" />
@@ -220,6 +233,15 @@ class JobHandlerAsWorkflowHandlerWithDbSkipLocked(BaseWorkflowHandlerConfigurati
assert self.is_app_workflow_scheduler
class JobHandlerAsWorkflowHandlerWithDbSkipLockedAttachToPool(JobHandlerAsWorkflowHandlerWithDbSkipLocked):
@classmethod
def handle_galaxy_config_kwds(cls, config):
config["job_config_file"] = config_file(POOL_JOB_CONFIG_TEMPLATE, assign_with=cls.assign_with)
config["server_name"] = "handler0"
config["attach_to_pools"] = ["job-handlers"]
class DefaultWorkflowHandlerIfJobHandlerOffTestCase(BaseWorkflowHandlerConfigurationTestCase):
@classmethod