From 67432c63c25544cbc4b75625ec0f94ab268bc20f Mon Sep 17 00:00:00 2001 From: almahmoud Date: Wed, 5 May 2021 17:31:21 -0400 Subject: [PATCH 01/40] Add shorten url --- lib/galaxy/config/__init__.py | 1 + lib/galaxy/managers/interactivetool.py | 10 +++++-- .../interactive/interactivetool_terminal.xml | 30 +++++++++++++++++++ 3 files changed, 39 insertions(+), 2 deletions(-) create mode 100644 tools/interactive/interactivetool_terminal.xml diff --git a/lib/galaxy/config/__init__.py b/lib/galaxy/config/__init__.py index 1134db7e5d6..9d8b012e7f3 100644 --- a/lib/galaxy/config/__init__.py +++ b/lib/galaxy/config/__init__.py @@ -848,6 +848,7 @@ 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.interactivetools_prefix = kwargs.get("interactivetools_prefix", "interactivetool") + self.interactivetools_shorten_url = kwargs.get("interactivetools_shorten_url", False) self.interactivetools_proxy_host = kwargs.get("interactivetools_proxy_host", None) self.containers_conf = parse_containers_config(self.containers_config_file) diff --git a/lib/galaxy/managers/interactivetool.py b/lib/galaxy/managers/interactivetool.py index b42392256e1..207f802cb35 100644 --- a/lib/galaxy/managers/interactivetool.py +++ b/lib/galaxy/managers/interactivetool.py @@ -270,7 +270,10 @@ class InteractiveToolManager: entry_point_encoded_id = trans.security.encode_id(entry_point.id) entry_point_class = entry_point.__class__.__name__.lower() entry_point_prefix = self.app.config.interactivetools_prefix - return f'{entry_point_encoded_id}-{entry_point.token}.{entry_point_class}.{entry_point_prefix}' + entry_point_token = entry_point.token + if self.app.config.interactivetools_shorten_url: + return f'{entry_point_encoded_id}-{entry_point_token[:10]}.{entry_point_prefix}' + return f'{entry_point_encoded_id}-{entry_point_token}.{entry_point_class}.{entry_point_prefix}' def get_entry_point_path(self, trans, entry_point): entry_point_encoded_id = trans.security.encode_id(entry_point.id) @@ -278,7 +281,10 @@ class InteractiveToolManager: entry_point_prefix = self.app.config.interactivetools_prefix rval = "/" if not entry_point.requires_domain: - rval = self.app.url_for(f'/{entry_point_prefix}/access/{entry_point_class}/{entry_point_encoded_id}/{entry_point.token}/') + if self.app.config.interactivetools_shorten_url: + rval = self.app.url_for(f'/{entry_point_prefix}/{entry_point_encoded_id}/{entry_point.token[:10]}/') + else: + rval = self.app.url_for(f'/{entry_point_prefix}/access/{entry_point_class}/{entry_point_encoded_id}/{entry_point.token}/') if entry_point.entry_url: rval = '{}/{}'.format(rval.rstrip('/'), entry_point.entry_url.lstrip('/')) if rval[0] != "/": diff --git a/tools/interactive/interactivetool_terminal.xml b/tools/interactive/interactivetool_terminal.xml new file mode 100644 index 00000000000..eafdc96835f --- /dev/null +++ b/tools/interactive/interactivetool_terminal.xml @@ -0,0 +1,30 @@ + + + cloudve/ttyd:latest + + + + 7681 + terminal + + + + ${__app__.security.encode_id($terminal.history_id)} + ${__app__.config.galaxy_infrastructure_url} + 80 + $__galaxy_url__ + true + true + + + + + + + + + + + From f7a4ac94e223db942040bfbaf787a9ce85382e3d Mon Sep 17 00:00:00 2001 From: almahmoud Date: Wed, 5 May 2021 17:31:52 -0400 Subject: [PATCH 02/40] Accept user/group IDs as ints --- lib/galaxy/jobs/runners/kubernetes.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index dffcfeda8cb..0e9095acae0 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -76,9 +76,9 @@ class KubernetesJobRunner(AsynchronousJobRunner): k8s_job_metadata=dict(map=str, default=None), k8s_supplemental_group_id=dict(map=str), k8s_pull_policy=dict(map=str, default="Default"), - k8s_run_as_user_id=dict(map=str, valid=lambda s: s == "$uid" or s.isdigit()), - k8s_run_as_group_id=dict(map=str, valid=lambda s: s == "$gid" or s.isdigit()), - k8s_fs_group_id=dict(map=int), + k8s_run_as_user_id=dict(map=str, valid=lambda s: s == "$uid" or isinstance(s, int) or s.isdigit()), + k8s_run_as_group_id=dict(map=str, valid=lambda s: s == "$gid" or isinstance(s, int) or s.isdigit()), + k8s_fs_group_id=dict(map=str, valid=lambda s: s == "$gid" or isinstance(s, int) or s.isdigit()), k8s_cleanup_job=dict(map=str, valid=lambda s: s in {"onsuccess", "always", "never"}, default="always"), k8s_pod_retries=dict(map=int, valid=lambda x: int(x) >= 0, default=3), k8s_walltime_limit=dict(map=int, valid=lambda x: int(x) >= 0, default=172800), From e1eef06243bf7d042c4a7ddd1c5c4d4d8096cf51 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Thu, 6 May 2021 14:34:40 +0200 Subject: [PATCH 03/40] Make guest_ports more reliable --- lib/galaxy/jobs/__init__.py | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index ff47ad0ad0f..1aae91bd45e 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -1174,10 +1174,13 @@ class JobWrapper(HasResourceParameters): @property def guest_ports(self): if hasattr(self, "interactivetools"): + # This works when the job is being prepared guest_ports = [ep.get('port') for ep in self.interactivetools] return guest_ports else: - return [] + # This works when handling a running job + job = self._load_job() + return [ep.tool_port for ep in job.interactivetool_entry_points] @property def working_directory(self): From 6254961cfd629387986f154ebc84498776225d7c Mon Sep 17 00:00:00 2001 From: almahmoud Date: Thu, 6 May 2021 11:36:20 -0400 Subject: [PATCH 04/40] Consolidate IT cleanup --- lib/galaxy/jobs/runners/kubernetes.py | 59 +++++++------------ .../interactive/interactivetool_terminal.xml | 30 ---------- 2 files changed, 20 insertions(+), 69 deletions(-) delete mode 100644 tools/interactive/interactivetool_terminal.xml diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 0e9095acae0..a5e51ade8f6 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -661,6 +661,8 @@ class KubernetesJobRunner(AsynchronousJobRunner): job_state.running = False self.mark_as_failed(job_state) try: + if job_state.job_wrapper.guest_ports: + self.__cleanup_k8s_interactivetools(job_state.job_wrapper, job) self.__cleanup_k8s_job(job) except Exception: log.exception("Could not clean up k8s batch job. Ignoring...") @@ -684,19 +686,7 @@ class KubernetesJobRunner(AsynchronousJobRunner): self.mark_as_failed(job_state) try: if job_state.job_wrapper.guest_ports: - k8s_job_prefix = self.__produce_k8s_job_prefix() - k8s_job_name = self.__get_k8s_job_name(k8s_job_prefix, job_state.job_wrapper) - job_failed = (job.obj['status']['failed'] > 0 - if 'failed' in job.obj['status'] else False) - log.debug(f'Deleting service/ingress for job with ID {job_state.job_wrapper.get_id_tag()}') - ingress_to_delete = find_ingress_object_by_name(self._pykube_api, k8s_job_name, self.runner_params['k8s_namespace']) - if ingress_to_delete and len(ingress_to_delete.response['items']) > 0: - k8s_ingress = Ingress(self._pykube_api, ingress_to_delete.response['items'][0]) - self.__cleanup_k8s_ingress(k8s_ingress, job_failed) - service_to_delete = find_service_object_by_name(self._pykube_api, k8s_job_name, self.runner_params['k8s_namespace']) - if service_to_delete and len(service_to_delete.response['items']) > 0: - k8s_service = Service(self._pykube_api, service_to_delete.response['items'][0]) - self.__cleanup_k8s_service(k8s_service, job_failed) + self.__cleanup_k8s_interactivetools(job_state.job_wrapper, job) self.__cleanup_k8s_job(job) except Exception: log.exception("Could not clean up k8s batch job. Ignoring...") @@ -755,6 +745,21 @@ class KubernetesJobRunner(AsynchronousJobRunner): pod = Pod(self._pykube_api, pods.response['items'][0]) return is_pod_unschedulable(self._pykube_api, pod, self.runner_params['k8s_namespace']) + def __cleanup_k8s_interactivetools(self, job_wrapper, k8s_job): + k8s_job_prefix = self.__produce_k8s_job_prefix() + k8s_job_name = "{}-{}".format(k8s_job_prefix, self.__force_label_conformity(job_wrapper.get_id_tag())) + log.debug(f'Deleting service/ingress for job with ID {job_wrapper.get_id_tag()}') + job_failed = (k8s_job.obj['status']['failed'] > 0 + if 'failed' in k8s_job.obj['status'] else False) + ingress_to_delete = find_ingress_object_by_name(self._pykube_api, k8s_job_name, self.runner_params['k8s_namespace']) + if ingress_to_delete and len(ingress_to_delete.response['items']) > 0: + k8s_ingress = Ingress(self._pykube_api, ingress_to_delete.response['items'][0]) + self.__cleanup_k8s_ingress(k8s_ingress, job_failed) + service_to_delete = find_service_object_by_name(self._pykube_api, k8s_job_name, self.runner_params['k8s_namespace']) + if service_to_delete and len(service_to_delete.response['items']) > 0: + k8s_service = Service(self._pykube_api, service_to_delete.response['items'][0]) + self.__cleanup_k8s_service(k8s_service, job_failed) + def stop_job(self, job_wrapper): """Attempts to delete a dispatched job to the k8s cluster""" job = job_wrapper.get_job() @@ -763,19 +768,7 @@ class KubernetesJobRunner(AsynchronousJobRunner): if job_to_delete and len(job_to_delete.response['items']) > 0: k8s_job = Job(self._pykube_api, job_to_delete.response['items'][0]) if job_wrapper.guest_ports: - k8s_job_prefix = self.__produce_k8s_job_prefix() - k8s_job_name = "{}-{}".format(k8s_job_prefix, self.__force_label_conformity(job_wrapper.get_id_tag())) - log.debug(f'Deleting service/ingress for job with ID {job_wrapper.get_id_tag()}') - job_failed = (k8s_job.obj['status']['failed'] > 0 - if 'failed' in k8s_job.obj['status'] else False) - ingress_to_delete = find_ingress_object_by_name(self._pykube_api, k8s_job_name, self.runner_params['k8s_namespace']) - if ingress_to_delete and len(ingress_to_delete.response['items']) > 0: - k8s_ingress = Ingress(self._pykube_api, ingress_to_delete.response['items'][0]) - self.__cleanup_k8s_ingress(k8s_ingress, job_failed) - service_to_delete = find_service_object_by_name(self._pykube_api, k8s_job_name, self.runner_params['k8s_namespace']) - if service_to_delete and len(service_to_delete.response['items']) > 0: - k8s_service = Service(self._pykube_api, service_to_delete.response['items'][0]) - self.__cleanup_k8s_service(k8s_service, job_failed) + self.__cleanup_k8s_interactivetools(job_wrapper, k8s_job) 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" @@ -818,17 +811,5 @@ class KubernetesJobRunner(AsynchronousJobRunner): " in job id '%s'", job_state.job_id) job = Job(self._pykube_api, jobs.response['items'][0]) if job_state.job_wrapper.guest_ports: - k8s_job_prefix = self.__produce_k8s_job_prefix() - k8s_job_name = self.__get_k8s_job_name(k8s_job_prefix, job_state.job_wrapper) - job_failed = (job.obj['status']['failed'] > 0 - if 'failed' in job.obj['status'] else False) - log.debug(f'Deleting service/ingress for job with ID {job_state.job_wrapper.get_id_tag()}') - ingress_to_delete = find_ingress_object_by_name(self._pykube_api, k8s_job_name, self.runner_params['k8s_namespace']) - if ingress_to_delete and len(ingress_to_delete.response['items']) > 0: - k8s_ingress = Ingress(self._pykube_api, ingress_to_delete.response['items'][0]) - self.__cleanup_k8s_ingress(k8s_ingress, job_failed) - service_to_delete = find_service_object_by_name(self._pykube_api, k8s_job_name, self.runner_params['k8s_namespace']) - if service_to_delete and len(service_to_delete.response['items']) > 0: - k8s_service = Service(self._pykube_api, service_to_delete.response['items'][0]) - self.__cleanup_k8s_service(k8s_service, job_failed) + self.__cleanup_k8s_interactivetools(job_state.job_wrapper, job) self.__cleanup_k8s_job(job) diff --git a/tools/interactive/interactivetool_terminal.xml b/tools/interactive/interactivetool_terminal.xml deleted file mode 100644 index eafdc96835f..00000000000 --- a/tools/interactive/interactivetool_terminal.xml +++ /dev/null @@ -1,30 +0,0 @@ - - - cloudve/ttyd:latest - - - - 7681 - terminal - - - - ${__app__.security.encode_id($terminal.history_id)} - ${__app__.config.galaxy_infrastructure_url} - 80 - $__galaxy_url__ - true - true - - - - - - - - - - - From 84e563cfceda924cd86930168512952f7d22aca0 Mon Sep 17 00:00:00 2001 From: almahmoud Date: Thu, 6 May 2021 12:31:43 -0400 Subject: [PATCH 05/40] SSL support for ITs --- lib/galaxy/jobs/runners/kubernetes.py | 12 +++++++++++- 1 file changed, 11 insertions(+), 1 deletion(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index a5e51ade8f6..0cbc2356d61 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -82,7 +82,9 @@ class KubernetesJobRunner(AsynchronousJobRunner): k8s_cleanup_job=dict(map=str, valid=lambda s: s in {"onsuccess", "always", "never"}, default="always"), k8s_pod_retries=dict(map=int, valid=lambda x: int(x) >= 0, default=3), k8s_walltime_limit=dict(map=int, valid=lambda x: int(x) >= 0, default=172800), - k8s_unschedulable_walltime_limit=dict(map=int, valid=lambda x: int(x) >= 0, default=1800),) + k8s_unschedulable_walltime_limit=dict(map=int, valid=lambda x: int(x) >= 0, default=1800), + k8s_interactivetools_use_ssl=dict(map=bool, default=False), + k8s_interactivetools_ingress_annotations=dict(map=str),) if 'runner_param_specs' not in kwargs: kwargs['runner_param_specs'] = dict() @@ -372,6 +374,14 @@ class KubernetesJobRunner(AsynchronousJobRunner): }]}} for ep in entry_points] } } + if self.runner_params.get("k8s_interactivetools_use_ssl"): + domains = list(set([e["domain"] for e in entry_points])) + k8s_spec_template["spec"]["tls"] = [{"hosts": [domain], + "secretName": re.sub("[^a-z0-9\-]", "-", domain) + } for domain in domains] + if self.runner_params.get("k8s_interactivetools_ingress_annotations"): + new_ann = yaml.safe_load(self.runner_params.get("k8s_interactivetools_ingress_annotations")) + k8s_spec_template["metadata"]["annotations"].update(new_ann) return k8s_spec_template def __get_k8s_security_context(self): From c50171dad7216b6a12f97e90ee61f560e20a657e Mon Sep 17 00:00:00 2001 From: almahmoud Date: Thu, 6 May 2021 13:19:06 -0400 Subject: [PATCH 06/40] Adding example conf files --- config/galaxy.yml.k8s_interactivetools | 20 +++++++++++ config/job_conf.yml.k8s_interactivetools | 45 ++++++++++++++++++++++++ 2 files changed, 65 insertions(+) create mode 100644 config/galaxy.yml.k8s_interactivetools create mode 100644 config/job_conf.yml.k8s_interactivetools diff --git a/config/galaxy.yml.k8s_interactivetools b/config/galaxy.yml.k8s_interactivetools new file mode 100644 index 00000000000..46d03989062 --- /dev/null +++ b/config/galaxy.yml.k8s_interactivetools @@ -0,0 +1,20 @@ +galaxy: + conda_auto_init: false + enable_data_manager_user_view: true + enable_tool_document_cache: true + integrated_tool_panel_config: /galaxy/server/config/mutable/integrated_tool_panel.xml + interactivetools_enable: true + interactivetools_map: database/interactivetools_map.sqlite + interactivetools_prefix: "its" + interactivetools_proxy_host: my-galaxy.usegvl.org + interactivetools_shorten_url: true + nginx_x_accel_redirect_base: '/galaxy/_x_accel_redirect' + outputs_to_working_directory: true + sanitize_allowlist_file: /galaxy/server/config/mutable/sanitize_allowlist.txt + shed_data_manager_config_file: /galaxy/server/config/mutable/shed_data_manager_conf.xml + shed_tool_config_file: /galaxy/server/config/mutable/editable_shed_tool_conf.xml + shed_tool_data_table_config: /galaxy/server/config/mutable/shed_tool_data_table_conf.xml + tool_data_path: '/galaxy/server/database/tool-data' + tool_dependency_dir: '/galaxy/server/database/deps' + tool_path: '/galaxy/server/database/tools' + \ No newline at end of file diff --git a/config/job_conf.yml.k8s_interactivetools b/config/job_conf.yml.k8s_interactivetools new file mode 100644 index 00000000000..7b7790f78aa --- /dev/null +++ b/config/job_conf.yml.k8s_interactivetools @@ -0,0 +1,45 @@ +execution: + default: dynamic_k8s_dispatcher + environments: + dynamic_k8s_dispatcher: + docker_default_container_id: 'galaxy/galaxy-min:latest' + container_monitor: false + docker_enabled: true + function: k8s_container_mapper + container_monitor: false + runner: dynamic + type: python + local: + runner: local +handling: + assign: + - db-skip-locked +limits: +- type: registered_user_concurrent_jobs + value: 5 +- type: anonymous_user_concurrent_jobs + value: 2 +runners: + k8s: + k8s_cleanup_job: onsuccess + k8s_fs_group_id: "101" + k8s_galaxy_instance_id: 'my-galaxy-instance' + k8s_namespace: 'galaxy-namespace' + k8s_persistent_volume_claims: |- + my-galaxy-instance-galaxy-pvc:/galaxy/server/database,my-galaxy-instance-cvmfs-gxy-data-pvc:/cvmfs/data.galaxyproject.org + ,my-galaxy-instance-galaxy-pvc/cvmfsclone:/cvmfs/cloud.galaxyproject.org + k8s_pod_priority_class: 'my-galaxy-instance-job-priority' + k8s_pull_policy: IfNotPresent + k8s_run_as_group_id: "101" + k8s_supplemental_group_id: "101" + k8s_use_service_account: true + k8s_run_as_user_id: "0" + k8s_run_as_group_id: "0" + k8s_interactivetools_use_ssl: true + k8s_interactivetools_ingress_annotations: + cert-manager.io/cluster-issuer: letsencrypt-prod + kubernetes.io/tls-acme: "true" + load: galaxy.jobs.runners.kubernetes:KubernetesJobRunner + local: + load: galaxy.jobs.runners.local:LocalJobRunner + workers: 4 From 1dae08a6865871df10c939a14f7e5ee6316c1454 Mon Sep 17 00:00:00 2001 From: Alexandru Mahmoud Date: Thu, 6 May 2021 13:29:50 -0400 Subject: [PATCH 07/40] Remove CVMFS-dependent portion --- config/job_conf.yml.k8s_interactivetools | 3 --- 1 file changed, 3 deletions(-) diff --git a/config/job_conf.yml.k8s_interactivetools b/config/job_conf.yml.k8s_interactivetools index 7b7790f78aa..4f7a510656c 100644 --- a/config/job_conf.yml.k8s_interactivetools +++ b/config/job_conf.yml.k8s_interactivetools @@ -25,9 +25,6 @@ runners: k8s_fs_group_id: "101" k8s_galaxy_instance_id: 'my-galaxy-instance' k8s_namespace: 'galaxy-namespace' - k8s_persistent_volume_claims: |- - my-galaxy-instance-galaxy-pvc:/galaxy/server/database,my-galaxy-instance-cvmfs-gxy-data-pvc:/cvmfs/data.galaxyproject.org - ,my-galaxy-instance-galaxy-pvc/cvmfsclone:/cvmfs/cloud.galaxyproject.org k8s_pod_priority_class: 'my-galaxy-instance-job-priority' k8s_pull_policy: IfNotPresent k8s_run_as_group_id: "101" From 18a9dd42b90033cbef70e1f81a190f529f02c5e2 Mon Sep 17 00:00:00 2001 From: almahmoud Date: Thu, 6 May 2021 13:35:11 -0400 Subject: [PATCH 08/40] Linting --- lib/galaxy/jobs/runners/kubernetes.py | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 0cbc2356d61..b3452a92117 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -377,8 +377,7 @@ class KubernetesJobRunner(AsynchronousJobRunner): if self.runner_params.get("k8s_interactivetools_use_ssl"): domains = list(set([e["domain"] for e in entry_points])) k8s_spec_template["spec"]["tls"] = [{"hosts": [domain], - "secretName": re.sub("[^a-z0-9\-]", "-", domain) - } for domain in domains] + "secretName": re.sub("[^a-z0-9-]", "-", domain)} for domain in domains] if self.runner_params.get("k8s_interactivetools_ingress_annotations"): new_ann = yaml.safe_load(self.runner_params.get("k8s_interactivetools_ingress_annotations")) k8s_spec_template["metadata"]["annotations"].update(new_ann) From 70215cdede1c999f40831abc63a671e9d8bddc07 Mon Sep 17 00:00:00 2001 From: Alexandru Mahmoud Date: Thu, 6 May 2021 16:16:38 -0400 Subject: [PATCH 09/40] Remove duplicate --- config/job_conf.yml.k8s_interactivetools | 1 - 1 file changed, 1 deletion(-) diff --git a/config/job_conf.yml.k8s_interactivetools b/config/job_conf.yml.k8s_interactivetools index 4f7a510656c..4f33951cac8 100644 --- a/config/job_conf.yml.k8s_interactivetools +++ b/config/job_conf.yml.k8s_interactivetools @@ -27,7 +27,6 @@ runners: k8s_namespace: 'galaxy-namespace' k8s_pod_priority_class: 'my-galaxy-instance-job-priority' k8s_pull_policy: IfNotPresent - k8s_run_as_group_id: "101" k8s_supplemental_group_id: "101" k8s_use_service_account: true k8s_run_as_user_id: "0" From 71656c0372faa86ffb1cfb90bbd1d298f390a8f4 Mon Sep 17 00:00:00 2001 From: Mason Houtz Date: Sat, 1 May 2021 10:04:12 -0700 Subject: [PATCH 10/40] Implemented new trigger scheme that updates rows in an audit table to avoid deadlocking writes to history --- lib/galaxy/model/__init__.py | 6 + lib/galaxy/model/mapping.py | 22 ++- lib/galaxy/model/migrate/triggers/__init__.py | 0 .../triggers/history_update_time_field.py} | 8 +- .../migrate/triggers/update_audit_table.py | 143 ++++++++++++++++++ .../versions/0165_add_content_update_time.py | 2 +- .../0174_readd_update_time_triggers.py | 2 +- .../migrate/versions/0175_history_audit.py | 86 +++++++++++ lib/galaxy/model/migrate/versions/util.py | 8 + 9 files changed, 263 insertions(+), 14 deletions(-) create mode 100644 lib/galaxy/model/migrate/triggers/__init__.py rename lib/galaxy/model/{triggers.py => migrate/triggers/history_update_time_field.py} (97%) create mode 100644 lib/galaxy/model/migrate/triggers/update_audit_table.py create mode 100644 lib/galaxy/model/migrate/versions/0175_history_audit.py diff --git a/lib/galaxy/model/__init__.py b/lib/galaxy/model/__init__.py index 873edab4fb4..015214682c6 100644 --- a/lib/galaxy/model/__init__.py +++ b/lib/galaxy/model/__init__.py @@ -1845,6 +1845,12 @@ def is_hda(d): return isinstance(d, HistoryDatasetAssociation) +class HistoryAudit(RepresentById): + def __init__(self, history, update_time): + self.history = history + self.update_time = update_time + + class History(HasTags, Dictifiable, UsesAnnotations, HasName, RepresentById): dict_collection_visible_keys = ['id', 'name', 'published', 'deleted'] diff --git a/lib/galaxy/model/mapping.py b/lib/galaxy/model/mapping.py index e8fc136bf89..42fb6a25145 100644 --- a/lib/galaxy/model/mapping.py +++ b/lib/galaxy/model/mapping.py @@ -46,10 +46,10 @@ from galaxy.model.custom_types import ( TrimmedString, UUIDType, ) +from galaxy.model.migrate.triggers.update_audit_table import install as install_timestamp_triggers from galaxy.model.orm.engine_factory import build_engine from galaxy.model.orm.now import now from galaxy.model.security import GalaxyRBACAgent -from galaxy.model.triggers import install_timestamp_triggers from galaxy.model.view import HistoryDatasetCollectionJobStateSummary from galaxy.model.view.utils import install_views @@ -202,7 +202,7 @@ model.History.table = Table( "history", metadata, Column("id", Integer, primary_key=True), Column("create_time", DateTime, default=now), - Column("update_time", DateTime, index=True, default=now, onupdate=now), + Column("update_time", DateTime, key="_update_time", index=True, default=now, onupdate=now), Column("user_id", Integer, ForeignKey("galaxy_user.id"), index=True), Column("name", TrimmedString(255)), Column("hid_counter", Integer, default=1), @@ -216,6 +216,15 @@ model.History.table = Table( Index('ix_history_slug', 'slug', mysql_length=200), ) +model.HistoryAudit.table = Table( + "history_audit", metadata, + Column("id", Integer, primary_key=True), + Column("history_id", Integer, ForeignKey("history.id"), nullable=False), + Column("update_time", DateTime, default=now, nullable=False), +) + +Index('ix_history_audit_history_id_update_time_desc', model.HistoryAudit.table.c.history_id.desc(), model.HistoryAudit.table.c.update_time.desc()) + model.HistoryUserShareAssociation.table = Table( "history_user_share_association", metadata, Column("id", Integer, primary_key=True), @@ -1889,7 +1898,10 @@ mapper(model.History, model.History.table, properties=dict( users_shared_with_count=column_property( select([func.count(model.HistoryUserShareAssociation.table.c.id)]).where(model.History.table.c.id == model.HistoryUserShareAssociation.table.c.history_id), deferred=True - ) + ), + update_time=column_property( + select([func.max(model.HistoryAudit.table.c.update_time)]).where(model.HistoryAudit.table.c.history_id == model.History.table.c.id), + ), )) # Set up proxy so that @@ -1905,13 +1917,13 @@ mapper(model.HistoryUserShareAssociation, model.HistoryUserShareAssociation.tabl mapper(model.User, model.User.table, properties=dict( histories=relation(model.History, backref="user", - order_by=desc(model.History.table.c.update_time)), + order_by=desc(model.History.update_time)), active_histories=relation(model.History, primaryjoin=( (model.History.table.c.user_id == model.User.table.c.id) & (not_(model.History.table.c.deleted)) ), - order_by=desc(model.History.table.c.update_time)), + order_by=desc(model.History.update_time)), galaxy_sessions=relation(model.GalaxySession, order_by=desc(model.GalaxySession.table.c.update_time)), diff --git a/lib/galaxy/model/migrate/triggers/__init__.py b/lib/galaxy/model/migrate/triggers/__init__.py new file mode 100644 index 00000000000..e69de29bb2d diff --git a/lib/galaxy/model/triggers.py b/lib/galaxy/model/migrate/triggers/history_update_time_field.py similarity index 97% rename from lib/galaxy/model/triggers.py rename to lib/galaxy/model/migrate/triggers/history_update_time_field.py index 8464c9822d6..3cbfa83eb9f 100644 --- a/lib/galaxy/model/triggers.py +++ b/lib/galaxy/model/migrate/triggers/history_update_time_field.py @@ -2,7 +2,7 @@ Database trigger installation and removal """ -from sqlalchemy import DDL +from galaxy.model.migrate.versions.util import execute_statements def install_timestamp_triggers(engine): @@ -21,12 +21,6 @@ def drop_timestamp_triggers(engine): execute_statements(engine, statements) -def execute_statements(engine, statements): - for sql in statements: - cmd = DDL(sql) - cmd.execute(bind=engine) - - def get_timestamp_install_sql(variant): """ Generate a list of SQL statements for installation of timestamp triggers diff --git a/lib/galaxy/model/migrate/triggers/update_audit_table.py b/lib/galaxy/model/migrate/triggers/update_audit_table.py new file mode 100644 index 00000000000..41fc1b8104b --- /dev/null +++ b/lib/galaxy/model/migrate/triggers/update_audit_table.py @@ -0,0 +1,143 @@ +from galaxy.model.migrate.versions.util import execute_statements + + +# function name prefix +fn_prefix = "fn_audit_history_by" + +# map between source table and associated incoming id field +trigger_config = { + 'history_dataset_association': "history_id", + 'history_dataset_collection_association': "history_id", + 'history': "id", +} + + +def install(engine): + """Install history audit table triggers""" + sql = _postgres_install() if 'postgres' in engine.name else _sqlite_install() + execute_statements(engine, sql) + + +def remove(engine): + """Uninstall history audit table triggers""" + sql = _postgres_remove() if 'postgres' in engine.name else _sqlite_remove() + execute_statements(engine, sql) + + +# Postgres trigger installation + + +def _postgres_remove(): + """postgres trigger removal sql""" + + sql = [] + sql.append(f"DROP FUNCTION IF EXISTS {fn_prefix}_history_id() CASCADE;") + sql.append(f"DROP FUNCTION IF EXISTS {fn_prefix}_id() CASCADE;") + + return sql + + +def _postgres_install(): + """postgres trigger installation sql""" + + sql = [] + + # postgres trigger function template + # need to make separate functions purely because the incoming history_id field name will be + # different for different source tables. There may be a fancier way to dynamically choose + # between incoming fields, but having 2 triggers fns seems straightforward + + def trigger_fn(id_field): + fn = f"{fn_prefix}_{id_field}" + + return f""" + CREATE OR REPLACE FUNCTION {fn}() + RETURNS TRIGGER + LANGUAGE 'plpgsql' + AS $BODY$ + BEGIN + INSERT INTO history_audit (history_id, update_time) + SELECT {id_field}, CURRENT_TIMESTAMP AT TIME ZONE 'UTC' + FROM new_table; + RETURN NULL; + END; + $BODY$ + """ + + def trigger_def(source_table, id_field, operation, when="AFTER"): + fn = f"{fn_prefix}_{id_field}" + + # Postgres supports many triggers per operation/table so the label can + # be indicative of what's happening + label = f"history_audit_by_{id_field}" + trigger_name = get_trigger_name(label, operation, when, statement=True) + + return f""" + CREATE TRIGGER {trigger_name} + {when} {operation} ON {source_table} + REFERENCING NEW TABLE AS new_table + FOR EACH STATEMENT EXECUTE FUNCTION {fn}(); + """ + + # trigger functions, each reads a different incoming id + for id_field in ["history_id", "id"]: + sql.append(trigger_fn(id_field)) + + # add triggers for each configured table (history, hda, hdca) + # picking the appropriate function via the config + for source_table, id_field in trigger_config.items(): + for operation in ["UPDATE", "INSERT"]: + sql.append(trigger_def(source_table, id_field, operation)) + + return sql + + +# Other DBs + + +def _sqlite_remove(): + sql = [] + + for source_table in trigger_config: + for operation in ["UPDATE", "INSERT"]: + trigger_name = get_trigger_name(source_table, operation, "AFTER") + sql.append(f"DROP TRIGGER IF EXISTS {trigger_name};") + + return sql + + +def _sqlite_install(): + sql = [] + + def trigger_def(source_table, id_field, operation, when="AFTER"): + + # only one trigger per operation/table in simple databases, so + # trigger name is less descriptive + trigger_name = get_trigger_name(source_table, operation, when) + + return f""" + CREATE TRIGGER {trigger_name} + {when} {operation} + ON {source_table} + FOR EACH ROW + BEGIN + INSERT INTO history_audit (history_id, update_time) + VALUES (NEW.{id_field}, CURRENT_TIMESTAMP); + END; + """ + + for source_table, id_field in trigger_config.items(): + for operation in ["UPDATE", "INSERT"]: + sql.append(trigger_def(source_table, id_field, operation)) + + return sql + + +# Utils + + +def get_trigger_name(label, operation, when, statement=False): + op_initial = operation.lower()[0] + when_initial = when.lower()[0] + rs = "s" if statement else "r" + return f"trigger_{label}_{when_initial}{op_initial}{rs}" diff --git a/lib/galaxy/model/migrate/versions/0165_add_content_update_time.py b/lib/galaxy/model/migrate/versions/0165_add_content_update_time.py index dc7b0f9647c..607dee1150d 100644 --- a/lib/galaxy/model/migrate/versions/0165_add_content_update_time.py +++ b/lib/galaxy/model/migrate/versions/0165_add_content_update_time.py @@ -7,9 +7,9 @@ import logging from sqlalchemy import Column, DateTime, MetaData, Table +from galaxy.model.migrate.triggers.history_update_time_field import drop_timestamp_triggers, install_timestamp_triggers from galaxy.model.migrate.versions.util import add_column, drop_column from galaxy.model.orm.now import now -from galaxy.model.triggers import drop_timestamp_triggers, install_timestamp_triggers log = logging.getLogger(__name__) metadata = MetaData() diff --git a/lib/galaxy/model/migrate/versions/0174_readd_update_time_triggers.py b/lib/galaxy/model/migrate/versions/0174_readd_update_time_triggers.py index 031ce925ef3..bd1ad2931c8 100644 --- a/lib/galaxy/model/migrate/versions/0174_readd_update_time_triggers.py +++ b/lib/galaxy/model/migrate/versions/0174_readd_update_time_triggers.py @@ -6,7 +6,7 @@ import logging from sqlalchemy import MetaData -from galaxy.model.triggers import ( +from galaxy.model.migrate.triggers.history_update_time_field import ( drop_timestamp_triggers, install_timestamp_triggers, ) diff --git a/lib/galaxy/model/migrate/versions/0175_history_audit.py b/lib/galaxy/model/migrate/versions/0175_history_audit.py new file mode 100644 index 00000000000..2457e6a35cf --- /dev/null +++ b/lib/galaxy/model/migrate/versions/0175_history_audit.py @@ -0,0 +1,86 @@ +""" +Add history audit table and associated triggers +""" + +import datetime +import logging + +from sqlalchemy import Column, DateTime, ForeignKey, Index, Integer, MetaData, Table + +from galaxy.model.migrate.triggers import ( + history_update_time_field as old_triggers, # rollback to old ones + update_audit_table as new_triggers, # install me +) +from galaxy.model.migrate.versions.util import ( + create_table, + drop_table +) + +log = logging.getLogger(__name__) +now = datetime.datetime.utcnow +metadata = MetaData() + +# NOTE: A normal incrementing PK is not typically used in an audit table, +# just foreign keys, but the ORM will probably choke without it +AuditTable = Table( + "history_audit", + metadata, + Column("id", Integer, primary_key=True), + Column("history_id", Integer, ForeignKey("history.id"), nullable=False), + Column("update_time", DateTime, default=now, nullable=False), +) + +Index('ix_history_audit_history_id_update_time_desc', AuditTable.c.history_id.desc(), AuditTable.c.update_time.desc()) + + +def upgrade(migrate_engine): + print(__doc__) + metadata.bind = migrate_engine + metadata.reflect() + + # create table + index + create_table(AuditTable) + + # populate with update_time from every history + copy_update_times = """ + INSERT INTO history_audit (history_id, update_time) + SELECT id, update_time FROM history + """ + migrate_engine.execute(copy_update_times) + + # drop existing timestamp triggers + old_triggers.drop_timestamp_triggers(migrate_engine) + + # install new timestamp triggers + new_triggers.install(migrate_engine) + + +def downgrade(migrate_engine): + print(__doc__) + metadata.bind = migrate_engine + metadata.reflect() + + # drop existing timestamp triggers + new_triggers.remove(migrate_engine) + + try: + # update history.update_time with vals from audit table + put_em_back = """ + UPDATE history h + SET update_time = a.max_update_time + FROM ( + SELECT history_id, max(update_time) as max_update_time + FROM history_audit + GROUP BY history_id, update_time + ) a + WHERE h.id = a.history_id + """ + migrate_engine.execute(put_em_back) + except Exception: + print("Unable to put update_times back") + + # drop audit table + drop_table(AuditTable) + + # install old timestamp triggers + old_triggers.install_timestamp_triggers(migrate_engine) diff --git a/lib/galaxy/model/migrate/versions/util.py b/lib/galaxy/model/migrate/versions/util.py index a9265f9a5e1..e961e07b4c2 100644 --- a/lib/galaxy/model/migrate/versions/util.py +++ b/lib/galaxy/model/migrate/versions/util.py @@ -3,6 +3,7 @@ import logging from sqlalchemy import ( BLOB, + DDL, Index, Table, Text @@ -193,3 +194,10 @@ def drop_index(index, table, column_name=None, metadata=None): index.drop() except Exception: log.exception("Dropping index '%s' from table '%s' failed", index, table) + + +def execute_statements(engine, raw_sql): + statements = raw_sql if isinstance(raw_sql, list) else [raw_sql] + for sql in statements: + cmd = DDL(sql) + cmd.execute(bind=engine) From 7bdf0f8dd0872994ff2b88d0f8251fa4c23965ca Mon Sep 17 00:00:00 2001 From: Mason Houtz Date: Sat, 1 May 2021 10:57:14 -0700 Subject: [PATCH 11/40] amended triggers to filter out updates with a NULL history_id --- lib/galaxy/model/migrate/triggers/update_audit_table.py | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/lib/galaxy/model/migrate/triggers/update_audit_table.py b/lib/galaxy/model/migrate/triggers/update_audit_table.py index 41fc1b8104b..be75154862e 100644 --- a/lib/galaxy/model/migrate/triggers/update_audit_table.py +++ b/lib/galaxy/model/migrate/triggers/update_audit_table.py @@ -58,7 +58,8 @@ def _postgres_install(): BEGIN INSERT INTO history_audit (history_id, update_time) SELECT {id_field}, CURRENT_TIMESTAMP AT TIME ZONE 'UTC' - FROM new_table; + FROM new_table + WHERE {id_field} IS NOT NULL; RETURN NULL; END; $BODY$ @@ -122,7 +123,8 @@ def _sqlite_install(): FOR EACH ROW BEGIN INSERT INTO history_audit (history_id, update_time) - VALUES (NEW.{id_field}, CURRENT_TIMESTAMP); + SELECT NEW.{id_field}, CURRENT_TIMESTAMP + WHERE NEW.{id_field} IS NOT NULL; END; """ From 7c70eb53b9117cb50faf0975e1ddf901b795fb44 Mon Sep 17 00:00:00 2001 From: Mason Houtz Date: Sat, 1 May 2021 11:09:35 -0700 Subject: [PATCH 12/40] deleted sqlite triggers first since sqlite does not support create or replace trigger syntax --- lib/galaxy/model/migrate/triggers/update_audit_table.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/lib/galaxy/model/migrate/triggers/update_audit_table.py b/lib/galaxy/model/migrate/triggers/update_audit_table.py index be75154862e..69ffb8b80b9 100644 --- a/lib/galaxy/model/migrate/triggers/update_audit_table.py +++ b/lib/galaxy/model/migrate/triggers/update_audit_table.py @@ -108,7 +108,8 @@ def _sqlite_remove(): def _sqlite_install(): - sql = [] + # delete old stuff first + sql = _sqlite_remove() def trigger_def(source_table, id_field, operation, when="AFTER"): From aa60fbd0b7bc70dfbe84e7748b9e2ad901c984a3 Mon Sep 17 00:00:00 2001 From: Mason Houtz Date: Sat, 1 May 2021 11:22:03 -0700 Subject: [PATCH 13/40] removed references to old update_time field --- lib/galaxy/managers/histories.py | 2 +- lib/galaxy/webapps/galaxy/controllers/history.py | 4 ++-- lib/galaxy/webapps/reports/controllers/system.py | 4 ++-- 3 files changed, 5 insertions(+), 5 deletions(-) diff --git a/lib/galaxy/managers/histories.py b/lib/galaxy/managers/histories.py index 534d6f123dc..026f2155579 100644 --- a/lib/galaxy/managers/histories.py +++ b/lib/galaxy/managers/histories.py @@ -85,7 +85,7 @@ class HistoryManager(sharable.SharableModelManager, deletable.PurgableManagerMix """ if self.user_manager.is_anonymous(user): return None if (not current_history or current_history.deleted) else current_history - desc_update_time = desc(self.model_class.table.c.update_time) + desc_update_time = desc(self.model_class.update_time) filters = self._munge_filters(filters, self.model_class.user_id == user.id) # TODO: normalize this return value return self.query(filters=filters, order_by=desc_update_time, limit=1, **kwargs).first() diff --git a/lib/galaxy/webapps/galaxy/controllers/history.py b/lib/galaxy/webapps/galaxy/controllers/history.py index 96dba4ba74f..77a2b183631 100644 --- a/lib/galaxy/webapps/galaxy/controllers/history.py +++ b/lib/galaxy/webapps/galaxy/controllers/history.py @@ -86,9 +86,9 @@ class HistoryListGrid(grids.Grid): def sort(self, trans, query, ascending, column_name=None): if ascending: - query = query.order_by(self.model_class.table.c.purged.asc(), self.model_class.table.c.update_time.desc()) + query = query.order_by(self.model_class.table.c.purged.asc(), self.model_class.update_time.desc()) else: - query = query.order_by(self.model_class.table.c.purged.desc(), self.model_class.table.c.update_time.desc()) + query = query.order_by(self.model_class.table.c.purged.desc(), self.model_class.update_time.desc()) return query def build_initial_query(self, trans, **kwargs): diff --git a/lib/galaxy/webapps/reports/controllers/system.py b/lib/galaxy/webapps/reports/controllers/system.py index b31a1123d3b..1d56d955075 100644 --- a/lib/galaxy/webapps/reports/controllers/system.py +++ b/lib/galaxy/webapps/reports/controllers/system.py @@ -60,7 +60,7 @@ class System(BaseUIController): for history in trans.sa_session.query(model.History) \ .filter(and_(model.History.table.c.user_id == null(), model.History.table.c.deleted == true(), - model.History.table.c.update_time < cutoff_time)): + model.History.update_time < cutoff_time)): for dataset in history.datasets: if not dataset.deleted: dataset_count += 1 @@ -86,7 +86,7 @@ class System(BaseUIController): histories = trans.sa_session.query(model.History) \ .filter(and_(model.History.table.c.deleted == true(), model.History.table.c.purged == false(), - model.History.table.c.update_time < cutoff_time)) \ + model.History.update_time < cutoff_time)) \ .options(eagerload('datasets')) for history in histories: From a2e923fbb36163986213773e69459a12cae25ea3 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Sat, 1 May 2021 22:23:04 +0200 Subject: [PATCH 14/40] Add update_time to History class for mypy --- lib/galaxy/model/__init__.py | 1 + 1 file changed, 1 insertion(+) diff --git a/lib/galaxy/model/__init__.py b/lib/galaxy/model/__init__.py index 015214682c6..1b6f8f233b8 100644 --- a/lib/galaxy/model/__init__.py +++ b/lib/galaxy/model/__init__.py @@ -1866,6 +1866,7 @@ class History(HasTags, Dictifiable, UsesAnnotations, HasName, RepresentById): self.importing = False self.genome_build = None self.published = False + self.update_time = None # Relationships self.user = user self.datasets = [] From fc42afd5f1cb19ab1c982f48362123f2649b4f4f Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Sat, 1 May 2021 22:29:06 +0200 Subject: [PATCH 15/40] Add microseconds to sqlite timestamp --- lib/galaxy/model/migrate/triggers/update_audit_table.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/lib/galaxy/model/migrate/triggers/update_audit_table.py b/lib/galaxy/model/migrate/triggers/update_audit_table.py index 69ffb8b80b9..a3ca8a811b9 100644 --- a/lib/galaxy/model/migrate/triggers/update_audit_table.py +++ b/lib/galaxy/model/migrate/triggers/update_audit_table.py @@ -124,7 +124,7 @@ def _sqlite_install(): FOR EACH ROW BEGIN INSERT INTO history_audit (history_id, update_time) - SELECT NEW.{id_field}, CURRENT_TIMESTAMP + SELECT NEW.{id_field}, strftime('%%Y-%%m-%%d %%H:%%M:%%f', 'now') WHERE NEW.{id_field} IS NOT NULL; END; """ From f9fb1085bc91c093ab4608951dbb8fb9f74800ec Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Sat, 1 May 2021 23:04:36 +0200 Subject: [PATCH 16/40] Allow grid sorting on class attributes if no mapped column --- lib/galaxy/web/framework/helpers/grids.py | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/lib/galaxy/web/framework/helpers/grids.py b/lib/galaxy/web/framework/helpers/grids.py index 4568987d3ac..f741e931bf4 100644 --- a/lib/galaxy/web/framework/helpers/grids.py +++ b/lib/galaxy/web/framework/helpers/grids.py @@ -76,10 +76,13 @@ class GridColumn: """Sort query using this column.""" if column_name is None: column_name = self.key + column = self.model_class.table.c.get(column_name) + if column is None: + column = getattr(self.model_class, column_name) if ascending: - query = query.order_by(self.model_class.table.c.get(column_name).asc()) + query = query.order_by(column.asc()) else: - query = query.order_by(self.model_class.table.c.get(column_name).desc()) + query = query.order_by(column.desc()) return query From e253e49a2f1bf93d0df838fbd94bcc29bf2d1779 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Mon, 3 May 2021 13:46:23 +0200 Subject: [PATCH 17/40] Drop comments --- lib/galaxy/model/migrate/triggers/update_audit_table.py | 6 ------ 1 file changed, 6 deletions(-) diff --git a/lib/galaxy/model/migrate/triggers/update_audit_table.py b/lib/galaxy/model/migrate/triggers/update_audit_table.py index a3ca8a811b9..ef718ae72a5 100644 --- a/lib/galaxy/model/migrate/triggers/update_audit_table.py +++ b/lib/galaxy/model/migrate/triggers/update_audit_table.py @@ -93,9 +93,6 @@ def _postgres_install(): return sql -# Other DBs - - def _sqlite_remove(): sql = [] @@ -136,9 +133,6 @@ def _sqlite_install(): return sql -# Utils - - def get_trigger_name(label, operation, when, statement=False): op_initial = operation.lower()[0] when_initial = when.lower()[0] From 04889e4b4027117cc51379e71648f1d43a8b1356 Mon Sep 17 00:00:00 2001 From: Mason Houtz Date: Mon, 3 May 2021 12:26:34 -0700 Subject: [PATCH 18/40] Added distinct for statement triggers to account for large multi-row collection or dataset updates --- lib/galaxy/model/migrate/triggers/update_audit_table.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/lib/galaxy/model/migrate/triggers/update_audit_table.py b/lib/galaxy/model/migrate/triggers/update_audit_table.py index ef718ae72a5..a3d90cbba01 100644 --- a/lib/galaxy/model/migrate/triggers/update_audit_table.py +++ b/lib/galaxy/model/migrate/triggers/update_audit_table.py @@ -57,7 +57,7 @@ def _postgres_install(): AS $BODY$ BEGIN INSERT INTO history_audit (history_id, update_time) - SELECT {id_field}, CURRENT_TIMESTAMP AT TIME ZONE 'UTC' + SELECT DISTINCT {id_field}, CURRENT_TIMESTAMP AT TIME ZONE 'UTC' FROM new_table WHERE {id_field} IS NOT NULL; RETURN NULL; From 5206b06d6cd37a0c4d0b544a347b0f36fd759402 Mon Sep 17 00:00:00 2001 From: Mason Houtz Date: Tue, 4 May 2021 08:30:35 -0700 Subject: [PATCH 19/40] added compound primary key on audit table and conflict ignore clause on triggers --- lib/galaxy/model/mapping.py | 7 ++----- .../model/migrate/triggers/update_audit_table.py | 10 ++++++---- .../model/migrate/versions/0175_history_audit.py | 13 ++++--------- 3 files changed, 12 insertions(+), 18 deletions(-) diff --git a/lib/galaxy/model/mapping.py b/lib/galaxy/model/mapping.py index 42fb6a25145..54df51a382b 100644 --- a/lib/galaxy/model/mapping.py +++ b/lib/galaxy/model/mapping.py @@ -218,13 +218,10 @@ model.History.table = Table( model.HistoryAudit.table = Table( "history_audit", metadata, - Column("id", Integer, primary_key=True), - Column("history_id", Integer, ForeignKey("history.id"), nullable=False), - Column("update_time", DateTime, default=now, nullable=False), + Column("history_id", Integer, ForeignKey("history.id"), primary_key=True, nullable=False), + Column("update_time", DateTime, default=now, primary_key=True, nullable=False), ) -Index('ix_history_audit_history_id_update_time_desc', model.HistoryAudit.table.c.history_id.desc(), model.HistoryAudit.table.c.update_time.desc()) - model.HistoryUserShareAssociation.table = Table( "history_user_share_association", metadata, Column("id", Integer, primary_key=True), diff --git a/lib/galaxy/model/migrate/triggers/update_audit_table.py b/lib/galaxy/model/migrate/triggers/update_audit_table.py index a3d90cbba01..607a12ea70e 100644 --- a/lib/galaxy/model/migrate/triggers/update_audit_table.py +++ b/lib/galaxy/model/migrate/triggers/update_audit_table.py @@ -14,7 +14,7 @@ trigger_config = { def install(engine): """Install history audit table triggers""" - sql = _postgres_install() if 'postgres' in engine.name else _sqlite_install() + sql = _postgres_install(engine) if 'postgres' in engine.name else _sqlite_install() execute_statements(engine, sql) @@ -37,7 +37,7 @@ def _postgres_remove(): return sql -def _postgres_install(): +def _postgres_install(engine): """postgres trigger installation sql""" sql = [] @@ -59,7 +59,8 @@ def _postgres_install(): INSERT INTO history_audit (history_id, update_time) SELECT DISTINCT {id_field}, CURRENT_TIMESTAMP AT TIME ZONE 'UTC' FROM new_table - WHERE {id_field} IS NOT NULL; + WHERE {id_field} IS NOT NULL + ON CONFLICT DO NOTHING; RETURN NULL; END; $BODY$ @@ -122,7 +123,8 @@ def _sqlite_install(): BEGIN INSERT INTO history_audit (history_id, update_time) SELECT NEW.{id_field}, strftime('%%Y-%%m-%%d %%H:%%M:%%f', 'now') - WHERE NEW.{id_field} IS NOT NULL; + WHERE NEW.{id_field} IS NOT NULL + ON CONFLICT IGNORE; END; """ diff --git a/lib/galaxy/model/migrate/versions/0175_history_audit.py b/lib/galaxy/model/migrate/versions/0175_history_audit.py index 2457e6a35cf..9b502304ad9 100644 --- a/lib/galaxy/model/migrate/versions/0175_history_audit.py +++ b/lib/galaxy/model/migrate/versions/0175_history_audit.py @@ -5,7 +5,7 @@ Add history audit table and associated triggers import datetime import logging -from sqlalchemy import Column, DateTime, ForeignKey, Index, Integer, MetaData, Table +from sqlalchemy import Column, DateTime, ForeignKey, Integer, MetaData, Table from galaxy.model.migrate.triggers import ( history_update_time_field as old_triggers, # rollback to old ones @@ -20,18 +20,13 @@ log = logging.getLogger(__name__) now = datetime.datetime.utcnow metadata = MetaData() -# NOTE: A normal incrementing PK is not typically used in an audit table, -# just foreign keys, but the ORM will probably choke without it AuditTable = Table( "history_audit", metadata, - Column("id", Integer, primary_key=True), - Column("history_id", Integer, ForeignKey("history.id"), nullable=False), - Column("update_time", DateTime, default=now, nullable=False), + Column("history_id", Integer, ForeignKey("history.id"), primary_key=True, nullable=False), + Column("update_time", DateTime, default=now, primary_key=True, nullable=False), ) -Index('ix_history_audit_history_id_update_time_desc', AuditTable.c.history_id.desc(), AuditTable.c.update_time.desc()) - def upgrade(migrate_engine): print(__doc__) @@ -71,7 +66,7 @@ def downgrade(migrate_engine): FROM ( SELECT history_id, max(update_time) as max_update_time FROM history_audit - GROUP BY history_id, update_time + GROUP BY history_id ) a WHERE h.id = a.history_id """ From deadd268c569e3a6e46ce43463197e5c1490a647 Mon Sep 17 00:00:00 2001 From: Mason Houtz Date: Tue, 4 May 2021 09:12:15 -0700 Subject: [PATCH 20/40] added sqlite-specific conflict clause to history audit table --- lib/galaxy/model/mapping.py | 2 ++ lib/galaxy/model/migrate/triggers/update_audit_table.py | 3 +-- lib/galaxy/model/migrate/versions/0175_history_audit.py | 3 ++- 3 files changed, 5 insertions(+), 3 deletions(-) diff --git a/lib/galaxy/model/mapping.py b/lib/galaxy/model/mapping.py index 54df51a382b..bbffbd36d5d 100644 --- a/lib/galaxy/model/mapping.py +++ b/lib/galaxy/model/mapping.py @@ -22,6 +22,7 @@ from sqlalchemy import ( MetaData, not_, Numeric, + PrimaryKeyConstraint, select, String, Table, TEXT, @@ -220,6 +221,7 @@ model.HistoryAudit.table = Table( "history_audit", metadata, Column("history_id", Integer, ForeignKey("history.id"), primary_key=True, nullable=False), Column("update_time", DateTime, default=now, primary_key=True, nullable=False), + PrimaryKeyConstraint(sqlite_on_conflict='IGNORE') ) model.HistoryUserShareAssociation.table = Table( diff --git a/lib/galaxy/model/migrate/triggers/update_audit_table.py b/lib/galaxy/model/migrate/triggers/update_audit_table.py index 607a12ea70e..f91dac11952 100644 --- a/lib/galaxy/model/migrate/triggers/update_audit_table.py +++ b/lib/galaxy/model/migrate/triggers/update_audit_table.py @@ -123,8 +123,7 @@ def _sqlite_install(): BEGIN INSERT INTO history_audit (history_id, update_time) SELECT NEW.{id_field}, strftime('%%Y-%%m-%%d %%H:%%M:%%f', 'now') - WHERE NEW.{id_field} IS NOT NULL - ON CONFLICT IGNORE; + WHERE NEW.{id_field} IS NOT NULL; END; """ diff --git a/lib/galaxy/model/migrate/versions/0175_history_audit.py b/lib/galaxy/model/migrate/versions/0175_history_audit.py index 9b502304ad9..e322270bb97 100644 --- a/lib/galaxy/model/migrate/versions/0175_history_audit.py +++ b/lib/galaxy/model/migrate/versions/0175_history_audit.py @@ -5,7 +5,7 @@ Add history audit table and associated triggers import datetime import logging -from sqlalchemy import Column, DateTime, ForeignKey, Integer, MetaData, Table +from sqlalchemy import Column, DateTime, ForeignKey, Integer, MetaData, PrimaryKeyConstraint, Table from galaxy.model.migrate.triggers import ( history_update_time_field as old_triggers, # rollback to old ones @@ -25,6 +25,7 @@ AuditTable = Table( metadata, Column("history_id", Integer, ForeignKey("history.id"), primary_key=True, nullable=False), Column("update_time", DateTime, default=now, primary_key=True, nullable=False), + PrimaryKeyConstraint(sqlite_on_conflict='IGNORE') ) From 014a5b524162155a45a97843fd86c34f6cff4f40 Mon Sep 17 00:00:00 2001 From: Mason Houtz Date: Tue, 4 May 2021 10:43:55 -0700 Subject: [PATCH 21/40] added row triggers for older postgres versions --- .../migrate/triggers/update_audit_table.py | 41 ++++++++++++++++--- .../migrate/versions/0175_history_audit.py | 1 + 2 files changed, 37 insertions(+), 5 deletions(-) diff --git a/lib/galaxy/model/migrate/triggers/update_audit_table.py b/lib/galaxy/model/migrate/triggers/update_audit_table.py index f91dac11952..d3340f8b727 100644 --- a/lib/galaxy/model/migrate/triggers/update_audit_table.py +++ b/lib/galaxy/model/migrate/triggers/update_audit_table.py @@ -47,7 +47,7 @@ def _postgres_install(engine): # different for different source tables. There may be a fancier way to dynamically choose # between incoming fields, but having 2 triggers fns seems straightforward - def trigger_fn(id_field): + def statement_trigger_fn(id_field): fn = f"{fn_prefix}_{id_field}" return f""" @@ -66,7 +66,24 @@ def _postgres_install(engine): $BODY$ """ - def trigger_def(source_table, id_field, operation, when="AFTER"): + def row_trigger_fn(id_field): + fn = f"{fn_prefix}_{id_field}" + + return f""" + CREATE OR REPLACE FUNCTION {fn}() + RETURNS TRIGGER + LANGUAGE 'plpgsql' + AS $BODY$ + BEGIN + INSERT INTO history_audit (history_id, update_time) + VALUES (NEW.{id_field}, CURRENT_TIMESTAMP AT TIME ZONE 'UTC') + ON CONFLICT DO NOTHING; + RETURN NULL; + END; + $BODY$ + """ + + def statement_trigger_def(source_table, id_field, operation, when="AFTER"): fn = f"{fn_prefix}_{id_field}" # Postgres supports many triggers per operation/table so the label can @@ -81,12 +98,26 @@ def _postgres_install(engine): FOR EACH STATEMENT EXECUTE FUNCTION {fn}(); """ - # trigger functions, each reads a different incoming id + def row_trigger_def(source_table, id_field, operation, when="AFTER"): + fn = f"{fn_prefix}_{id_field}" + + label = f"history_audit_by_{id_field}" + trigger_name = get_trigger_name(label, operation, when, statement=True) + + return f""" + CREATE TRIGGER {trigger_name} + {when} {operation} ON {source_table} + FOR EACH ROW EXECUTE FUNCTION {fn}(); + """ + + # pick row or statement triggers depending on postgres version + version = engine.dialect.server_version_info[0] + trigger_fn = statement_trigger_fn if version > 10 else row_trigger_fn + trigger_def = statement_trigger_def if version > 10 else row_trigger_def + for id_field in ["history_id", "id"]: sql.append(trigger_fn(id_field)) - # add triggers for each configured table (history, hda, hdca) - # picking the appropriate function via the config for source_table, id_field in trigger_config.items(): for operation in ["UPDATE", "INSERT"]: sql.append(trigger_def(source_table, id_field, operation)) diff --git a/lib/galaxy/model/migrate/versions/0175_history_audit.py b/lib/galaxy/model/migrate/versions/0175_history_audit.py index e322270bb97..521bb762396 100644 --- a/lib/galaxy/model/migrate/versions/0175_history_audit.py +++ b/lib/galaxy/model/migrate/versions/0175_history_audit.py @@ -35,6 +35,7 @@ def upgrade(migrate_engine): metadata.reflect() # create table + index + AuditTable.drop(migrate_engine, checkfirst=True) create_table(AuditTable) # populate with update_time from every history From 3d597171428ac7e664460c1f9bdf3d9b1749c7f2 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Fri, 7 May 2021 20:06:01 +0200 Subject: [PATCH 22/40] Add celery setup that regularly cleans audit_table --- config/celery.ini | 14 +++++++++++ config/dev.ini | 7 +----- doc/source/admin/galaxy_options.rst | 22 +++++++++++++++++ lib/galaxy/celery/__init__.py | 17 ++++++++++++++ lib/galaxy/celery/tasks.py | 26 +++++++++++++++++++++ lib/galaxy/config/sample/galaxy.yml.sample | 8 +++++++ lib/galaxy/webapps/galaxy/config_schema.yml | 8 +++++++ 7 files changed, 96 insertions(+), 6 deletions(-) create mode 100644 config/celery.ini diff --git a/config/celery.ini b/config/celery.ini new file mode 100644 index 00000000000..53a358c5843 --- /dev/null +++ b/config/celery.ini @@ -0,0 +1,14 @@ +[env] +PYTHONPATH = lib + +[watcher:celery] +cmd = celery +args = --app galaxy.celery worker --concurrency 2 -l debug +copy_env = True +numprocesses = 1 + +[watcher:celery-beat] +cmd = celery +args = --app galaxy.celery beat -l debug +copy_env = True +numprocesses = 1 diff --git a/config/dev.ini b/config/dev.ini index d4e5494fe7b..56bd44ee53f 100644 --- a/config/dev.ini +++ b/config/dev.ini @@ -1,5 +1,6 @@ [circus] debug = True +include = celery.ini [env] PYTHONPATH=lib @@ -25,9 +26,3 @@ stop_children = True [socket:web] host = 0.0.0.0 port = 8080 - -[watcher:celery] -cmd = celery -args = --app galaxy.celery worker -l debug -copy_env = True -numprocesses = 1 diff --git a/doc/source/admin/galaxy_options.rst b/doc/source/admin/galaxy_options.rst index 8d1c05cb446..2de2b71318f 100644 --- a/doc/source/admin/galaxy_options.rst +++ b/doc/source/admin/galaxy_options.rst @@ -246,6 +246,17 @@ :Type: float +~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~ +``history_audit_table_prune_interval`` +~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~ + +:Description: + Time (in seconds) between attempts to remove old rows from the + history_audit database table. Set to 0 to disable pruning. +:Default: ``3600`` +:Type: int + + ~~~~~~~~~~~~~ ``file_path`` ~~~~~~~~~~~~~ @@ -2709,6 +2720,17 @@ :Type: bool +~~~~~~~~~~~~~~~~~~~~~ +``statsd_mock_calls`` +~~~~~~~~~~~~~~~~~~~~~ + +:Description: + Mock out statsd client calls - only used by testing infrastructure + really. Do not set this in production environments. +:Default: ``false`` +:Type: bool + + ~~~~~~~~~~~~~~~~~~~~~~ ``library_import_dir`` ~~~~~~~~~~~~~~~~~~~~~~ diff --git a/lib/galaxy/celery/__init__.py b/lib/galaxy/celery/__init__.py index d8f916e3738..4d7c573ecdf 100644 --- a/lib/galaxy/celery/__init__.py +++ b/lib/galaxy/celery/__init__.py @@ -52,8 +52,25 @@ def get_broker(): return config.amqp_internal_connection +def get_history_audit_table_prune_interval(): + config = get_config() + if config: + return config.history_audit_table_prune_interval + else: + return 3600 + + broker = get_broker() celery_app = Celery('galaxy', broker=broker, include=['galaxy.celery.tasks']) +prune_interval = get_history_audit_table_prune_interval() +if prune_interval > 0: + celery_app.conf.beat_schedule = { + 'prune-history-audit-table': { + 'task': 'galaxy.celery.tasks.prune_history_audit_table', + 'schedule': prune_interval, + }, + } +celery_app.conf.timezone = 'UTC' if __name__ == '__main__': diff --git a/lib/galaxy/celery/tasks.py b/lib/galaxy/celery/tasks.py index 878c080f64b..2275420b603 100644 --- a/lib/galaxy/celery/tasks.py +++ b/lib/galaxy/celery/tasks.py @@ -1,4 +1,9 @@ from lagom import magic_bind_to_container +from sqlalchemy import ( + and_, + func, + tuple_ +) from sqlalchemy.orm.scoping import ( scoped_session, ) @@ -9,6 +14,7 @@ from galaxy.celery import celery_app from galaxy.jobs.manager import JobManager from galaxy.managers.hdas import HDAManager from galaxy.managers.lddas import LDDAManager +from galaxy.util import ExecutionTimer from galaxy.util.custom_logging import get_logger from . import get_galaxy_app @@ -73,3 +79,23 @@ def export_history( job.state = model.Job.states.NEW sa_session.flush() job_manager.enqueue(job) + + +@celery_app.task +@galaxy_task +def prune_history_audit_table(sa_session: scoped_session): + """Prune ever growing history_audit table.""" + history_audit_table = model.HistoryAudit.table + latest_subq = sa_session.query( + history_audit_table.c.history_id, + func.max(history_audit_table.c.update_time).label('max_update_time')).group_by(history_audit_table.c.history_id).subquery() + not_latest_query = sa_session.query( + history_audit_table.c.history_id, history_audit_table.c.update_time + ).select_from(latest_subq).join( + history_audit_table, and_( + history_audit_table.c.update_time < latest_subq.columns.max_update_time, + history_audit_table.c.history_id == latest_subq.columns.history_id)) + d = history_audit_table.delete() + timer = ExecutionTimer() + sa_session.execute(d.where(tuple_(history_audit_table.c.history_id, history_audit_table.c.update_time).in_(not_latest_query))) + log.debug(f"Successfully pruned history_audit table {timer}") diff --git a/lib/galaxy/config/sample/galaxy.yml.sample b/lib/galaxy/config/sample/galaxy.yml.sample index 3f90f6c42c3..7265362f579 100644 --- a/lib/galaxy/config/sample/galaxy.yml.sample +++ b/lib/galaxy/config/sample/galaxy.yml.sample @@ -222,6 +222,10 @@ galaxy: # seconds). #database_wait_sleep: 1.0 + # Time (in seconds) between attempts to remove old rows from the + # history_audit database table. Set to 0 to disable pruning. + #history_audit_table_prune_interval: 3600 + # Where dataset files are stored. It must be accessible at the same # path on any cluster nodes that will run Galaxy jobs, unless using # Pulsar. The default value has been changed from 'files' to 'objects' @@ -1365,6 +1369,10 @@ galaxy: # prefix with a tag path set to the page url #statsd_influxdb: false + # Mock out statsd client calls - only used by testing infrastructure + # really. Do not set this in production environments. + #statsd_mock_calls: false + # Add an option to the library upload form which allows administrators # to upload a directory of files. #library_import_dir: null diff --git a/lib/galaxy/webapps/galaxy/config_schema.yml b/lib/galaxy/webapps/galaxy/config_schema.yml index 87d9be76f7c..0efb5b9846b 100644 --- a/lib/galaxy/webapps/galaxy/config_schema.yml +++ b/lib/galaxy/webapps/galaxy/config_schema.yml @@ -197,6 +197,14 @@ mapping: desc: | Time to sleep between attempts if database_wait is enabled (in seconds). + history_audit_table_prune_interval: + type: int + default: 3600 + required: false + desc: | + Time (in seconds) between attempts to remove old rows from the history_audit database table. + Set to 0 to disable pruning. + file_path: type: str default: objects From 4265dc4b62921dc7487b8657f444af851b951dce Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Mon, 10 May 2021 18:05:37 +0200 Subject: [PATCH 23/40] Move history audit pruning to HistroyAudit class --- lib/galaxy/celery/tasks.py | 18 +----------------- lib/galaxy/model/__init__.py | 16 ++++++++++++++++ 2 files changed, 17 insertions(+), 17 deletions(-) diff --git a/lib/galaxy/celery/tasks.py b/lib/galaxy/celery/tasks.py index 2275420b603..a5238e8b7c1 100644 --- a/lib/galaxy/celery/tasks.py +++ b/lib/galaxy/celery/tasks.py @@ -1,9 +1,4 @@ from lagom import magic_bind_to_container -from sqlalchemy import ( - and_, - func, - tuple_ -) from sqlalchemy.orm.scoping import ( scoped_session, ) @@ -85,17 +80,6 @@ def export_history( @galaxy_task def prune_history_audit_table(sa_session: scoped_session): """Prune ever growing history_audit table.""" - history_audit_table = model.HistoryAudit.table - latest_subq = sa_session.query( - history_audit_table.c.history_id, - func.max(history_audit_table.c.update_time).label('max_update_time')).group_by(history_audit_table.c.history_id).subquery() - not_latest_query = sa_session.query( - history_audit_table.c.history_id, history_audit_table.c.update_time - ).select_from(latest_subq).join( - history_audit_table, and_( - history_audit_table.c.update_time < latest_subq.columns.max_update_time, - history_audit_table.c.history_id == latest_subq.columns.history_id)) - d = history_audit_table.delete() timer = ExecutionTimer() - sa_session.execute(d.where(tuple_(history_audit_table.c.history_id, history_audit_table.c.update_time).in_(not_latest_query))) + model.HistoryAudit.prune(sa_session) log.debug(f"Successfully pruned history_audit table {timer}") diff --git a/lib/galaxy/model/__init__.py b/lib/galaxy/model/__init__.py index 1b6f8f233b8..af5c156ecf2 100644 --- a/lib/galaxy/model/__init__.py +++ b/lib/galaxy/model/__init__.py @@ -35,6 +35,7 @@ from sqlalchemy import ( select, text, true, + tuple_, type_coerce, types) from sqlalchemy.exc import OperationalError @@ -1850,6 +1851,21 @@ class HistoryAudit(RepresentById): self.history = history self.update_time = update_time + @classmethod + def prune(cls, sa_session): + history_audit_table = cls.table + latest_subq = sa_session.query( + history_audit_table.c.history_id, + func.max(history_audit_table.c.update_time).label('max_update_time')).group_by(history_audit_table.c.history_id).subquery() + not_latest_query = sa_session.query( + history_audit_table.c.history_id, history_audit_table.c.update_time + ).select_from(latest_subq).join( + history_audit_table, and_( + history_audit_table.c.update_time < latest_subq.columns.max_update_time, + history_audit_table.c.history_id == latest_subq.columns.history_id)) + d = history_audit_table.delete() + sa_session.execute(d.where(tuple_(history_audit_table.c.history_id, history_audit_table.c.update_time).in_(not_latest_query))) + class History(HasTags, Dictifiable, UsesAnnotations, HasName, RepresentById): From e0f0fa338e45365058d73745f300a396bd5e3c2f Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Mon, 10 May 2021 20:45:42 +0200 Subject: [PATCH 24/40] Add periodic history audit prune task that runs when not using celery. --- lib/galaxy/app.py | 11 ++++++++ lib/galaxy/util/task.py | 51 +++++++++++++++++++++++++++++++++++++ test/unit/util/test_task.py | 28 ++++++++++++++++++++ 3 files changed, 90 insertions(+) create mode 100644 lib/galaxy/util/task.py create mode 100644 test/unit/util/test_task.py diff --git a/lib/galaxy/app.py b/lib/galaxy/app.py index fa6b7e6b62a..dffa996a080 100644 --- a/lib/galaxy/app.py +++ b/lib/galaxy/app.py @@ -61,6 +61,7 @@ from galaxy.util import ( heartbeat, StructuredExecutionTimer, ) +from galaxy.util.task import IntervalTask from galaxy.visualization.data_providers.registry import DataProviderRegistry from galaxy.visualization.genomes import Genomes from galaxy.visualization.plugins.registry import VisualizationsRegistry @@ -304,6 +305,16 @@ class UniverseApplication(StructuredApp, GalaxyManagerApplication): self.authnz_manager = managers.AuthnzManager(self, self.config.oidc_config_file, self.config.oidc_backends_config_file) + + if not self.config.enable_celery_tasks and self.config.history_audit_table_prune_interval > 0: + self.prune_history_audit_task = IntervalTask( + func=lambda: galaxy.model.HistoryAudit.prune(self.model.session), + name="HistoryAuditTablePruneTask", + interval=self.config.history_audit_table_prune_interval, + immediate_start=False, + time_execution=True) + self.application_stack.register_postfork_function(self.prune_history_audit_task.start) + self.haltables.append(("HistoryAuditTablePruneTask", self.prune_history_audit_task.shutdown)) # Start the job manager self.application_stack.register_postfork_function(self.job_manager.start) self.proxy_manager = ProxyManager(self.config) diff --git a/lib/galaxy/util/task.py b/lib/galaxy/util/task.py new file mode 100644 index 00000000000..9da0e347d08 --- /dev/null +++ b/lib/galaxy/util/task.py @@ -0,0 +1,51 @@ +import logging +from threading import ( + Event, + Thread, +) + +from galaxy.util import ExecutionTimer + +log = logging.getLogger(__name__) + + +class IntervalTask: + + def __init__(self, func, name="Periodic task", interval=3600, immediate_start=False, time_execution=False): + """ + Run an arbitrary function `func` every `interval` seconds. + + Set `immediate_start` to True to run `func` when task is started. + """ + self.func = func + self.name = name + self.interval = interval + self.time_execution = time_execution + self.immediate_start = immediate_start + self.event = Event() + self.thread = Thread(target=self.run, name=self.name, daemon=True) + self.running = False + + def start(self): + self.running = True + self.thread.start() + + def _exec(self): + if self.time_execution: + timer = ExecutionTimer() + self.func() + if self.time_execution: + log.debug(f"Executed periodic task {self.name} {timer}") + + def run(self): + if self.immediate_start: + self._exec() + while not self.event.isSet(): + self.event.wait(self.interval) + if self.running: + self._exec() + + def shutdown(self): + self.running = False + self.event.set() + self.thread.join(5) diff --git a/test/unit/util/test_task.py b/test/unit/util/test_task.py new file mode 100644 index 00000000000..a1b7b317634 --- /dev/null +++ b/test/unit/util/test_task.py @@ -0,0 +1,28 @@ +import time + +from galaxy.util.task import IntervalTask + + +def test_interval_task_immediate_start(): + results = [] + task = IntervalTask(lambda: results.append(1), name="test_task", interval=0.2, immediate_start=True) + task.start() + task.shutdown() + assert len(results) == 1 + + +def test_interval_task_delayed_start(): + results = [] + task = IntervalTask(lambda: results.append(1), name="test_task", interval=0.2, immediate_start=False) + task.start() + task.shutdown() + assert len(results) == 0 + + +def test_interval_task_delayed_start_run_once(): + results = [] + task = IntervalTask(lambda: results.append(1), name="test_task", interval=0.2, immediate_start=False) + task.start() + time.sleep(0.25) + task.shutdown() + assert len(results) == 1 From 9e9343ccd2634183b201e849cf3070dfbfe6b3b4 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Mon, 10 May 2021 21:53:32 +0200 Subject: [PATCH 25/40] Fix cleanup_datasets.py script --- scripts/cleanup_datasets/cleanup_datasets.py | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/scripts/cleanup_datasets/cleanup_datasets.py b/scripts/cleanup_datasets/cleanup_datasets.py index 9934b5c297b..6f3311d0c1b 100755 --- a/scripts/cleanup_datasets/cleanup_datasets.py +++ b/scripts/cleanup_datasets/cleanup_datasets.py @@ -140,12 +140,12 @@ def delete_userless_histories(app, cutoff_time, info_only=False, force_retry=Fal if force_retry: histories = app.sa_session.query(app.model.History) \ .filter(and_(app.model.History.table.c.user_id == null(), - app.model.History.table.c.update_time < cutoff_time)) + app.model.History.update_time < cutoff_time)) else: histories = app.sa_session.query(app.model.History) \ .filter(and_(app.model.History.table.c.user_id == null(), app.model.History.table.c.deleted == false(), - app.model.History.table.c.update_time < cutoff_time)) + app.model.History.update_time < cutoff_time)) for history in histories: if not info_only: log.info("Deleting history id %d", history.id) @@ -170,13 +170,13 @@ def purge_histories(app, cutoff_time, remove_from_disk, info_only=False, force_r if force_retry: histories = app.sa_session.query(app.model.History) \ .filter(and_(app.model.History.table.c.deleted == true(), - app.model.History.table.c.update_time < cutoff_time)) \ + app.model.History.update_time < cutoff_time)) \ .options(eagerload('datasets')) else: histories = app.sa_session.query(app.model.History) \ .filter(and_(app.model.History.table.c.deleted == true(), app.model.History.table.c.purged == false(), - app.model.History.table.c.update_time < cutoff_time)) \ + app.model.History.update_time < cutoff_time)) \ .options(eagerload('datasets')) for history in histories: log.info("### Processing history id %d (%s)", history.id, unicodify(history.name)) From b62a0dc70338909cf227f92024cfce49f1bbf516 Mon Sep 17 00:00:00 2001 From: Oleg Zharkov Date: Tue, 11 May 2021 20:20:55 +0200 Subject: [PATCH 26/40] fix dowpdown close --- client/src/components/Masthead/MastheadItem.vue | 14 ++++---------- 1 file changed, 4 insertions(+), 10 deletions(-) diff --git a/client/src/components/Masthead/MastheadItem.vue b/client/src/components/Masthead/MastheadItem.vue index f98a3eb2ccc..ec7877ae9ca 100644 --- a/client/src/components/Masthead/MastheadItem.vue +++ b/client/src/components/Masthead/MastheadItem.vue @@ -22,6 +22,7 @@ this.hideDropdown()); }, destroyed() { - if (this.$refs.dropdown && this.galaxyIframe) { - this.galaxyIframe.removeEventListener("load", this.iframeListener); - } + window.removeEventListener("blur", () => this.hideDropdown()); }, methods: { - iframeListener() { - return this.galaxyIframe.contentDocument.addEventListener("click", this.hideDropdown); - }, hideDropdown() { if (this.$refs.dropdown) this.$refs.dropdown.hide(); }, From 0a088624a6c035b9a9fcc26cb4af95ab26e5eb5a Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Tue, 11 May 2021 21:40:51 +0200 Subject: [PATCH 27/40] Use dateutil for python 3.6 compatibility datetime.fromisoformat was added in python 3.6 --- lib/galaxy/webapps/galaxy/api/history_contents.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/lib/galaxy/webapps/galaxy/api/history_contents.py b/lib/galaxy/webapps/galaxy/api/history_contents.py index 5bff0f937d4..b73c08457cf 100644 --- a/lib/galaxy/webapps/galaxy/api/history_contents.py +++ b/lib/galaxy/webapps/galaxy/api/history_contents.py @@ -5,7 +5,8 @@ import json import logging import os import re -from datetime import datetime + +import dateutil.parser from galaxy import ( exceptions, @@ -1082,8 +1083,7 @@ class HistoryContentsController(BaseGalaxyAPIController, UsesLibraryMixinItems, # if it hasn't then we can short-circuit the poll request since = kwd.get('update_time-gt', None) if since: - since_str = self.history_contents_filters.parse_date(since) - since_date = datetime.fromisoformat(since_str) + since_date = dateutil.parser.isoparse(since) if history.update_time <= since_date: trans.response.status = 204 return From cef6bd4fb265b526da25213e2c6ef310c799c140 Mon Sep 17 00:00:00 2001 From: Oleg Zharkov Date: Tue, 11 May 2021 20:24:17 +0200 Subject: [PATCH 28/40] simlify listener --- client/src/components/Masthead/MastheadItem.vue | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/client/src/components/Masthead/MastheadItem.vue b/client/src/components/Masthead/MastheadItem.vue index ec7877ae9ca..8221f9a8300 100644 --- a/client/src/components/Masthead/MastheadItem.vue +++ b/client/src/components/Masthead/MastheadItem.vue @@ -22,7 +22,6 @@ this.hideDropdown()); + window.addEventListener("blur", this.hideDropdown); }, destroyed() { - window.removeEventListener("blur", () => this.hideDropdown()); + window.removeEventListener("blur", this.hideDropdown); }, methods: { hideDropdown() { From d0971cc1b5e7bd1fccf9d61f6128717b1f258f48 Mon Sep 17 00:00:00 2001 From: almahmoud Date: Tue, 11 May 2021 16:43:46 -0400 Subject: [PATCH 29/40] Move configs to schema --- doc/source/admin/galaxy_options.rst | 34 +++++++++++++++++++++ lib/galaxy/config/__init__.py | 3 -- lib/galaxy/config/sample/galaxy.yml.sample | 13 ++++++++ lib/galaxy/webapps/galaxy/config_schema.yml | 16 ++++++++++ 4 files changed, 63 insertions(+), 3 deletions(-) diff --git a/doc/source/admin/galaxy_options.rst b/doc/source/admin/galaxy_options.rst index 8d1c05cb446..f14f01aed1a 100644 --- a/doc/source/admin/galaxy_options.rst +++ b/doc/source/admin/galaxy_options.rst @@ -1602,6 +1602,29 @@ :Type: str +~~~~~~~~~~~~~~~~~~~~~~~~~~~ +``interactivetools_prefix`` +~~~~~~~~~~~~~~~~~~~~~~~~~~~ + +:Description: + Prefix to use in the formation of the subdomain or path for + interactive tools +:Default: ``interactivetools`` +:Type: str + + +~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~ +``interactivetools_shorten_url`` +~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~ + +:Description: + Shorten the uuid portion of the subdomain or path for interactive + tools. Especially useful for avoiding the need for wildcard + certificates by keeping subdomain under 63 chars +:Default: ``false`` +:Type: bool + + ~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~ ``retry_interactivetool_metadata_internally`` ~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~ @@ -2709,6 +2732,17 @@ :Type: bool +~~~~~~~~~~~~~~~~~~~~~ +``statsd_mock_calls`` +~~~~~~~~~~~~~~~~~~~~~ + +:Description: + Mock out statsd client calls - only used by testing infrastructure + really. Do not set this in production environments. +:Default: ``false`` +:Type: bool + + ~~~~~~~~~~~~~~~~~~~~~~ ``library_import_dir`` ~~~~~~~~~~~~~~~~~~~~~~ diff --git a/lib/galaxy/config/__init__.py b/lib/galaxy/config/__init__.py index 9d8b012e7f3..89c4b2fafbd 100644 --- a/lib/galaxy/config/__init__.py +++ b/lib/galaxy/config/__init__.py @@ -847,9 +847,6 @@ 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.interactivetools_prefix = kwargs.get("interactivetools_prefix", "interactivetool") - self.interactivetools_shorten_url = kwargs.get("interactivetools_shorten_url", False) - self.interactivetools_proxy_host = kwargs.get("interactivetools_proxy_host", None) self.containers_conf = parse_containers_config(self.containers_config_file) diff --git a/lib/galaxy/config/sample/galaxy.yml.sample b/lib/galaxy/config/sample/galaxy.yml.sample index 3f90f6c42c3..f9026c7042e 100644 --- a/lib/galaxy/config/sample/galaxy.yml.sample +++ b/lib/galaxy/config/sample/galaxy.yml.sample @@ -878,6 +878,15 @@ galaxy: # . #interactivetools_map: interactivetools_map.sqlite + # Prefix to use in the formation of the subdomain or path for + # interactive tools + #interactivetools_prefix: interactivetools + + # Shorten the uuid portion of the subdomain or path for interactive + # tools. Especially useful for avoiding the need for wildcard + # certificates by keeping subdomain under 63 chars + #interactivetools_shorten_url: false + # Galaxy Interactive Tools (GxITs) can be stopped from within the # Galaxy interface, killing the GxIT job without completing its # metadata setting post-job steps. In such a case it may be desirable @@ -1365,6 +1374,10 @@ galaxy: # prefix with a tag path set to the page url #statsd_influxdb: false + # Mock out statsd client calls - only used by testing infrastructure + # really. Do not set this in production environments. + #statsd_mock_calls: false + # Add an option to the library upload form which allows administrators # to upload a directory of files. #library_import_dir: null diff --git a/lib/galaxy/webapps/galaxy/config_schema.yml b/lib/galaxy/webapps/galaxy/config_schema.yml index 87d9be76f7c..83ea2df6742 100644 --- a/lib/galaxy/webapps/galaxy/config_schema.yml +++ b/lib/galaxy/webapps/galaxy/config_schema.yml @@ -1155,6 +1155,22 @@ mapping: desc: | Map for interactivetool proxy. + interactivetools_prefix: + type: str + default: "interactivetools" + required: false + desc: | + Prefix to use in the formation of the subdomain or path for interactive tools + + interactivetools_shorten_url: + type: bool + default: false + required: false + desc: | + Shorten the uuid portion of the subdomain or path for interactive tools. + Especially useful for avoiding the need for wildcard certificates by keeping + subdomain under 63 chars + retry_interactivetool_metadata_internally: type: bool default: true From be2967b24d33dc33576b1e90c56335664473a141 Mon Sep 17 00:00:00 2001 From: almahmoud Date: Tue, 11 May 2021 20:19:48 -0400 Subject: [PATCH 30/40] Fix for subdomain ITs --- lib/galaxy/jobs/runners/kubernetes.py | 1 + 1 file changed, 1 insertion(+) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index b3452a92117..6942acfc9a7 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -349,6 +349,7 @@ class KubernetesJobRunner(AsynchronousJobRunner): if entry_point.requires_domain: entry_point_subdomain = self.app.interactivetool_manager.get_entry_point_subdomain(self.app, entry_point) entry_point_domain = f'{entry_point_subdomain}.{entry_point_domain}' + entry_point_path = '/' entry_points.append({"tool_port": entry_point.tool_port, "domain": entry_point_domain, "entry_path": entry_point_path}) k8s_spec_template = { "metadata": { From e80e1a80d6e737ddd69899c6ef39dd2effecd7f0 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Wed, 12 May 2021 11:51:29 +0200 Subject: [PATCH 31/40] Convert timezone aware datetime to naive datetime sqlalchemy DateTime columns are not timezone aware, so parse `since` into timezone-aware datetime and then convert to naive datetime object representing UTC, assuming history.update_time represents UTC time. --- lib/galaxy/webapps/galaxy/api/history_contents.py | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/lib/galaxy/webapps/galaxy/api/history_contents.py b/lib/galaxy/webapps/galaxy/api/history_contents.py index b73c08457cf..a892c81bf6e 100644 --- a/lib/galaxy/webapps/galaxy/api/history_contents.py +++ b/lib/galaxy/webapps/galaxy/api/history_contents.py @@ -1,6 +1,7 @@ """ API operations on the contents of a history. """ +import datetime import json import logging import os @@ -1083,7 +1084,10 @@ class HistoryContentsController(BaseGalaxyAPIController, UsesLibraryMixinItems, # if it hasn't then we can short-circuit the poll request since = kwd.get('update_time-gt', None) if since: - since_date = dateutil.parser.isoparse(since) + # sqlalchemy DateTime columns are not timezone aware, so parse `since` into timezone-aware + # datetime and then convert to naive datetime object representing UTC, + # assuming history.update_time represents UTC time. + since_date = dateutil.parser.isoparse(since).astimezone(datetime.timezone.utc).replace(tzinfo=None) if history.update_time <= since_date: trans.response.status = 204 return From f466edda9acc188519533690c3b69a4bdbff1d2a Mon Sep 17 00:00:00 2001 From: Simon Bray <32272674+simonbray@users.noreply.github.com> Date: Wed, 12 May 2021 13:06:56 +0200 Subject: [PATCH 32/40] force galaxy_id to string in galactic_job_json --- lib/galaxy/tool_util/cwl/util.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/lib/galaxy/tool_util/cwl/util.py b/lib/galaxy/tool_util/cwl/util.py index f8fcc038e4c..411df08965a 100644 --- a/lib/galaxy/tool_util/cwl/util.py +++ b/lib/galaxy/tool_util/cwl/util.py @@ -172,7 +172,7 @@ def galactic_job_json( def replacement_file(value): if value.get('galaxy_id'): - return {"src": "hda", "id": value['galaxy_id']} + return {"src": "hda", "id": str(value['galaxy_id'])} file_path = value.get("location", None) or value.get("path", None) # format to match output definitions in tool, where did filetype come from? filetype = value.get("filetype", None) or value.get("format", None) @@ -281,7 +281,7 @@ def galactic_job_json( def replacement_collection(value): if value.get('galaxy_id'): - return {"src": "hdca", "id": value['galaxy_id']} + return {"src": "hdca", "id": str(value['galaxy_id'])} assert "collection_type" in value collection_type = value["collection_type"] elements = to_elements(value, collection_type) From 56148b972cd588eac00cf8ae93a6fc497af6d83a Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Wed, 12 May 2021 14:25:11 +0200 Subject: [PATCH 33/40] Unit test for history audit and pruning --- test/unit/data/test_galaxy_mapping.py | 19 +++++++++++++++++++ 1 file changed, 19 insertions(+) diff --git a/test/unit/data/test_galaxy_mapping.py b/test/unit/data/test_galaxy_mapping.py index 94064615a2d..86e0011e912 100644 --- a/test/unit/data/test_galaxy_mapping.py +++ b/test/unit/data/test_galaxy_mapping.py @@ -472,6 +472,25 @@ class MappingTests(BaseModelTestCase): assert contents_iter_names(ids=[d1.id, d3.id]) == ["1", "3"] + def test_history_audit(self): + model = self.model + u = model.User(email="contents@foo.bar.baz", password="password") + # gs = model.GalaxySession() + h1 = model.History(name="HistoryAuditHistory", user=u) + + def count_audit_table_entries(): + return self.session().query(model.HistoryAudit.table.c.history_id).filter(model.HistoryAudit.table.c.history_id == h1.id).count() + + self.persist(u, h1, expunge=False) + assert count_audit_table_entries() == 1 + + self.new_hda(h1, name="1") + self.session().flush() + # db_next_hid modifies history, plus trigger on HDA means 2 additional audit rows + assert count_audit_table_entries() == 3 + model.HistoryAudit.prune(self.session()) + assert count_audit_table_entries() == 1 + def _non_empty_flush(self): model = self.model lf = model.LibraryFolder(name="RootFolder") From 93af8ec7808f6161dcffd9db82bd7432227858e4 Mon Sep 17 00:00:00 2001 From: Sergey Golitsynskiy Date: Wed, 12 May 2021 12:14:27 -0400 Subject: [PATCH 34/40] Modify test to ensure correct rows are pruned --- test/unit/data/test_galaxy_mapping.py | 37 +++++++++++++++++++++------ 1 file changed, 29 insertions(+), 8 deletions(-) diff --git a/test/unit/data/test_galaxy_mapping.py b/test/unit/data/test_galaxy_mapping.py index 86e0011e912..d68a984d45b 100644 --- a/test/unit/data/test_galaxy_mapping.py +++ b/test/unit/data/test_galaxy_mapping.py @@ -475,21 +475,42 @@ class MappingTests(BaseModelTestCase): def test_history_audit(self): model = self.model u = model.User(email="contents@foo.bar.baz", password="password") - # gs = model.GalaxySession() h1 = model.History(name="HistoryAuditHistory", user=u) + h2 = model.History(name="HistoryAuditHistory", user=u) - def count_audit_table_entries(): - return self.session().query(model.HistoryAudit.table.c.history_id).filter(model.HistoryAudit.table.c.history_id == h1.id).count() + def get_audit_table_entries(history): + return self.session().query(model.HistoryAudit.table).filter( + model.HistoryAudit.table.c.history_id == history.id).all() - self.persist(u, h1, expunge=False) - assert count_audit_table_entries() == 1 + def get_latest_entry(entries): + # key ensures result is correct if new columns are added + return max(entries, key=lambda x: x.update_time) + + self.persist(u, h1, h2, expunge=False) + assert len(get_audit_table_entries(h1)) == 1 + assert len(get_audit_table_entries(h2)) == 1 self.new_hda(h1, name="1") + self.new_hda(h2, name="2") self.session().flush() - # db_next_hid modifies history, plus trigger on HDA means 2 additional audit rows - assert count_audit_table_entries() == 3 + # db_next_hid modifies history, plus trigger on HDA means 2 additional audit rows per history + + h1_audits = get_audit_table_entries(h1) + h2_audits = get_audit_table_entries(h2) + assert len(h1_audits) == 3 + assert len(h2_audits) == 3 + + h1_latest = get_latest_entry(h1_audits) + h2_latest = get_latest_entry(h2_audits) + model.HistoryAudit.prune(self.session()) - assert count_audit_table_entries() == 1 + + h1_audits = get_audit_table_entries(h1) + h2_audits = get_audit_table_entries(h2) + assert len(h1_audits) == 1 + assert len(h2_audits) == 1 + assert h1_audits[0] == h1_latest + assert h2_audits[0] == h2_latest def _non_empty_flush(self): model = self.model From e1f517cf4b74ad8b3882b965caecda127c4b1773 Mon Sep 17 00:00:00 2001 From: almahmoud Date: Wed, 12 May 2021 13:34:39 -0400 Subject: [PATCH 35/40] Remove separate conf files --- config/galaxy.yml.k8s_interactivetools | 20 --------- config/job_conf.yml.k8s_interactivetools | 41 ------------------- .../sample/job_conf.xml.sample_advanced | 14 +++++++ 3 files changed, 14 insertions(+), 61 deletions(-) delete mode 100644 config/galaxy.yml.k8s_interactivetools delete mode 100644 config/job_conf.yml.k8s_interactivetools diff --git a/config/galaxy.yml.k8s_interactivetools b/config/galaxy.yml.k8s_interactivetools deleted file mode 100644 index 46d03989062..00000000000 --- a/config/galaxy.yml.k8s_interactivetools +++ /dev/null @@ -1,20 +0,0 @@ -galaxy: - conda_auto_init: false - enable_data_manager_user_view: true - enable_tool_document_cache: true - integrated_tool_panel_config: /galaxy/server/config/mutable/integrated_tool_panel.xml - interactivetools_enable: true - interactivetools_map: database/interactivetools_map.sqlite - interactivetools_prefix: "its" - interactivetools_proxy_host: my-galaxy.usegvl.org - interactivetools_shorten_url: true - nginx_x_accel_redirect_base: '/galaxy/_x_accel_redirect' - outputs_to_working_directory: true - sanitize_allowlist_file: /galaxy/server/config/mutable/sanitize_allowlist.txt - shed_data_manager_config_file: /galaxy/server/config/mutable/shed_data_manager_conf.xml - shed_tool_config_file: /galaxy/server/config/mutable/editable_shed_tool_conf.xml - shed_tool_data_table_config: /galaxy/server/config/mutable/shed_tool_data_table_conf.xml - tool_data_path: '/galaxy/server/database/tool-data' - tool_dependency_dir: '/galaxy/server/database/deps' - tool_path: '/galaxy/server/database/tools' - \ No newline at end of file diff --git a/config/job_conf.yml.k8s_interactivetools b/config/job_conf.yml.k8s_interactivetools deleted file mode 100644 index 4f33951cac8..00000000000 --- a/config/job_conf.yml.k8s_interactivetools +++ /dev/null @@ -1,41 +0,0 @@ -execution: - default: dynamic_k8s_dispatcher - environments: - dynamic_k8s_dispatcher: - docker_default_container_id: 'galaxy/galaxy-min:latest' - container_monitor: false - docker_enabled: true - function: k8s_container_mapper - container_monitor: false - runner: dynamic - type: python - local: - runner: local -handling: - assign: - - db-skip-locked -limits: -- type: registered_user_concurrent_jobs - value: 5 -- type: anonymous_user_concurrent_jobs - value: 2 -runners: - k8s: - k8s_cleanup_job: onsuccess - k8s_fs_group_id: "101" - k8s_galaxy_instance_id: 'my-galaxy-instance' - k8s_namespace: 'galaxy-namespace' - k8s_pod_priority_class: 'my-galaxy-instance-job-priority' - k8s_pull_policy: IfNotPresent - k8s_supplemental_group_id: "101" - k8s_use_service_account: true - k8s_run_as_user_id: "0" - k8s_run_as_group_id: "0" - k8s_interactivetools_use_ssl: true - k8s_interactivetools_ingress_annotations: - cert-manager.io/cluster-issuer: letsencrypt-prod - kubernetes.io/tls-acme: "true" - load: galaxy.jobs.runners.kubernetes:KubernetesJobRunner - local: - load: galaxy.jobs.runners.local:LocalJobRunner - workers: 4 diff --git a/lib/galaxy/config/sample/job_conf.xml.sample_advanced b/lib/galaxy/config/sample/job_conf.xml.sample_advanced index e50748de254..0a344004129 100644 --- a/lib/galaxy/config/sample/job_conf.xml.sample_advanced +++ b/lib/galaxy/config/sample/job_conf.xml.sample_advanced @@ -314,6 +314,15 @@ a requirement of the other. --> + + + + + + + + + true + true + From d6992b2bc57539d3faeba2378e9ec3a58446b589 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Wed, 12 May 2021 18:30:34 +0200 Subject: [PATCH 36/40] Update default base images and make DEFAULT_CHANNELS configurable via env var --- lib/galaxy/tool_util/deps/mulled/invfile.lua | 2 +- lib/galaxy/tool_util/deps/mulled/mulled_build.py | 9 ++++++--- 2 files changed, 7 insertions(+), 4 deletions(-) diff --git a/lib/galaxy/tool_util/deps/mulled/invfile.lua b/lib/galaxy/tool_util/deps/mulled/invfile.lua index a883abdd297..26a586c04a7 100644 --- a/lib/galaxy/tool_util/deps/mulled/invfile.lua +++ b/lib/galaxy/tool_util/deps/mulled/invfile.lua @@ -60,7 +60,7 @@ end local destination_base_image = VAR.DEST_BASE_IMAGE if destination_base_image == '' then - destination_base_image = 'bgruening/busybox-bash:0.1' + destination_base_image = 'quay.io/bioconda/base-glibc-busybox-bash:latest' end local verbose = VAR.VERBOSE diff --git a/lib/galaxy/tool_util/deps/mulled/mulled_build.py b/lib/galaxy/tool_util/deps/mulled/mulled_build.py index e5e1f21e40f..41345a0b676 100644 --- a/lib/galaxy/tool_util/deps/mulled/mulled_build.py +++ b/lib/galaxy/tool_util/deps/mulled/mulled_build.py @@ -46,9 +46,12 @@ from ..conda_compat import MetaData log = logging.getLogger(__name__) DIRNAME = os.path.dirname(__file__) -DEFAULT_BASE_IMAGE = "bgruening/busybox-bash:0.1" -DEFAULT_EXTENDED_BASE_IMAGE = "bioconda/extended-base-image:latest" -DEFAULT_CHANNELS = ["conda-forge", "bioconda"] +DEFAULT_BASE_IMAGE = os.environ.get("DEFAULT_BASE_IMAGE", "quay.io/bioconda/base-glibc-busybox-bash:latest") +DEFAULT_EXTENDED_BASE_IMAGE = os.environ.get("DEFAULT_EXTENDED_BASE_IMAGE", "quay.io/bioconda/base-glibc-debian-bash:latest") +if 'DEFAULT_MULLED_CONDA_CHANNELS' in os.environ: + DEFAULT_CHANNELS = os.environ['DEFAULT_MULLED_CONDA_CHANNELS'].split(',') +else: + DEFAULT_CHANNELS = ["conda-forge", "bioconda"] DEFAULT_REPOSITORY_TEMPLATE = "quay.io/${namespace}/${image}" DEFAULT_BINDS = ["build/dist:/usr/local/"] DEFAULT_WORKING_DIR = '/source/' From ab8bfe4b100d0d145c55b2311a5df16f65558fb8 Mon Sep 17 00:00:00 2001 From: almahmoud Date: Thu, 13 May 2021 12:03:57 -0400 Subject: [PATCH 37/40] Change default prefix --- doc/source/admin/galaxy_options.rst | 2 +- lib/galaxy/config/sample/galaxy.yml.sample | 2 +- lib/galaxy/webapps/galaxy/config_schema.yml | 2 +- 3 files changed, 3 insertions(+), 3 deletions(-) diff --git a/doc/source/admin/galaxy_options.rst b/doc/source/admin/galaxy_options.rst index f14f01aed1a..b6c85e4077f 100644 --- a/doc/source/admin/galaxy_options.rst +++ b/doc/source/admin/galaxy_options.rst @@ -1609,7 +1609,7 @@ :Description: Prefix to use in the formation of the subdomain or path for interactive tools -:Default: ``interactivetools`` +:Default: ``interactivetool`` :Type: str diff --git a/lib/galaxy/config/sample/galaxy.yml.sample b/lib/galaxy/config/sample/galaxy.yml.sample index f9026c7042e..ee9b2506af6 100644 --- a/lib/galaxy/config/sample/galaxy.yml.sample +++ b/lib/galaxy/config/sample/galaxy.yml.sample @@ -880,7 +880,7 @@ galaxy: # Prefix to use in the formation of the subdomain or path for # interactive tools - #interactivetools_prefix: interactivetools + #interactivetools_prefix: interactivetool # Shorten the uuid portion of the subdomain or path for interactive # tools. Especially useful for avoiding the need for wildcard diff --git a/lib/galaxy/webapps/galaxy/config_schema.yml b/lib/galaxy/webapps/galaxy/config_schema.yml index 83ea2df6742..3dd769c1756 100644 --- a/lib/galaxy/webapps/galaxy/config_schema.yml +++ b/lib/galaxy/webapps/galaxy/config_schema.yml @@ -1157,7 +1157,7 @@ mapping: interactivetools_prefix: type: str - default: "interactivetools" + default: "interactivetool" required: false desc: | Prefix to use in the formation of the subdomain or path for interactive tools From fd03c446c0347517fb483fc4c65700ca9b342411 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Fri, 14 May 2021 18:22:15 +0200 Subject: [PATCH 38/40] Fix mulled singularity building This used to fail with ``` [May 13 12:53:06] DEBU Container [8a432293fb2f step-adb59c6b2d] started, waiting for completion [May 13 12:53:06] SERR mkdir: missing operand [May 13 12:53:06] SERR Try 'mkdir --help' for more information. [May 13 12:53:06] ERRO Task processing failed: Unexpected exit code [1] of container [8a432293fb2f step-adb59c6b2d], container preserved ``` Involucro interpretes comma-separated arguments in `.run` as individual commands. You can see the full error in https://github.com/BioContainers/multi-package-containers/pull/1713/checks?check_run_id=2575472382 Broke in https://github.com/mvdbeek/galaxy/commit/71a70ea1f12cca48c33ee5569bd74035f96376f7 --- lib/galaxy/tool_util/deps/mulled/invfile.lua | 8 ++++---- test/unit/tool_util/mulled/test_mulled_build.py | 9 +++++++++ 2 files changed, 13 insertions(+), 4 deletions(-) diff --git a/lib/galaxy/tool_util/deps/mulled/invfile.lua b/lib/galaxy/tool_util/deps/mulled/invfile.lua index 26a586c04a7..db1abe49df6 100644 --- a/lib/galaxy/tool_util/deps/mulled/invfile.lua +++ b/lib/galaxy/tool_util/deps/mulled/invfile.lua @@ -101,11 +101,11 @@ inv.task('build') if VAR.SINGULARITY ~= '' then inv.task('singularity') .using(singularity_image) - .withHostConfig({binds = {"build:/data", "'" .. singularity_image_dir .. "':/import"}, privileged = true}) + .withHostConfig({binds = {"build:/data", singularity_image_dir .. ":/import"}, privileged = true}) .withConfig({entrypoint = {'/bin/sh', '-c'}}) - .run('mkdir', '-p', '/usr/local/var/singularity/mnt/container') - .run('singularity', 'build', '/import/' .. VAR.SINGULARITY_IMAGE_NAME, '/import/Singularity.def') - .run('chown', VAR.USER_ID, '/import/' .. VAR.SINGULARITY_IMAGE_NAME) + .run('mkdir -p /usr/local/var/singularity/mnt/container && ' + .. 'singularity build /import/' .. VAR.SINGULARITY_IMAGE_NAME .. ' /import/Singularity.def && ' + .. 'chown ' .. VAR.USER_ID .. ' /import/' .. VAR.SINGULARITY_IMAGE_NAME) end inv.task('cleanup') diff --git a/test/unit/tool_util/mulled/test_mulled_build.py b/test/unit/tool_util/mulled/test_mulled_build.py index 2aead594a72..7ef5cedfcf0 100644 --- a/test/unit/tool_util/mulled/test_mulled_build.py +++ b/test/unit/tool_util/mulled/test_mulled_build.py @@ -5,6 +5,7 @@ from galaxy.tool_util.deps.mulled.mulled_build import ( build_target, DEFAULT_BASE_IMAGE, DEFAULT_EXTENDED_BASE_IMAGE, + mull_targets, ) from ..util import external_dependency_management @@ -18,3 +19,11 @@ from ..util import external_dependency_management def test_base_image_for_targets(target, version, base_image): target = build_target(target, version=version) assert base_image_for_targets([target]) == base_image + + +@external_dependency_management +def test_mulled_build_files_cli(tmpdir): + singularity_image_dir = tmpdir.mkdir('singularity_image_dir') + target = build_target('zlib') + mull_targets([target], command='build-and-test', singularity=True, singularity_image_dir=singularity_image_dir) + assert singularity_image_dir.join('zlib').exists() From 8184ab089eba2d71a98b8e6af7a2a3bb073d9976 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Fri, 14 May 2021 21:32:22 +0200 Subject: [PATCH 39/40] Test with space in path --- test/unit/tool_util/mulled/test_mulled_build.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/test/unit/tool_util/mulled/test_mulled_build.py b/test/unit/tool_util/mulled/test_mulled_build.py index 7ef5cedfcf0..568a713d9f7 100644 --- a/test/unit/tool_util/mulled/test_mulled_build.py +++ b/test/unit/tool_util/mulled/test_mulled_build.py @@ -23,7 +23,7 @@ def test_base_image_for_targets(target, version, base_image): @external_dependency_management def test_mulled_build_files_cli(tmpdir): - singularity_image_dir = tmpdir.mkdir('singularity_image_dir') + singularity_image_dir = tmpdir.mkdir('singularity image dir') target = build_target('zlib') mull_targets([target], command='build-and-test', singularity=True, singularity_image_dir=singularity_image_dir) assert singularity_image_dir.join('zlib').exists() From 8107975f3583c9ec79906314e516ec5f5f49d381 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Mon, 17 May 2021 19:56:57 +0200 Subject: [PATCH 40/40] Add update_time indexes on columns used in ORDER BY clauses Without an index here the database typically needs to do a seq scan.. Noticed this with an extremly fast query to api/jobs with a user that has few datasets, and a timeout for a user that has many datasets. --- lib/galaxy/model/mapping.py | 12 ++-- .../0176_add_indexes_on_update_time.py | 65 +++++++++++++++++++ 2 files changed, 71 insertions(+), 6 deletions(-) create mode 100644 lib/galaxy/model/migrate/versions/0176_add_indexes_on_update_time.py diff --git a/lib/galaxy/model/mapping.py b/lib/galaxy/model/mapping.py index bbffbd36d5d..8917c0be34e 100644 --- a/lib/galaxy/model/mapping.py +++ b/lib/galaxy/model/mapping.py @@ -236,7 +236,7 @@ model.HistoryDatasetAssociation.table = Table( Column("history_id", Integer, ForeignKey("history.id"), index=True), Column("dataset_id", Integer, ForeignKey("dataset.id"), index=True), Column("create_time", DateTime, default=now), - Column("update_time", DateTime, default=now, onupdate=now), + Column("update_time", DateTime, default=now, onupdate=now, index=True), Column("state", TrimmedString(64), index=True, key="_state"), Column("copied_from_history_dataset_association_id", Integer, ForeignKey("history_dataset_association.id"), nullable=True), @@ -513,7 +513,7 @@ model.LibraryDatasetDatasetAssociation.table = Table( Column("library_dataset_id", Integer, ForeignKey("library_dataset.id"), index=True), Column("dataset_id", Integer, ForeignKey("dataset.id"), index=True), Column("create_time", DateTime, default=now), - Column("update_time", DateTime, default=now, onupdate=now), + Column("update_time", DateTime, default=now, onupdate=now, index=True), Column("state", TrimmedString(64), index=True, key="_state"), Column("copied_from_history_dataset_association_id", Integer, ForeignKey("history_dataset_association.id", use_alter=True, name='history_dataset_association_dataset_id_fkey'), @@ -610,7 +610,7 @@ model.Job.table = Table( "job", metadata, Column("id", Integer, primary_key=True), Column("create_time", DateTime, default=now), - Column("update_time", DateTime, default=now, onupdate=now), + Column("update_time", DateTime, default=now, onupdate=now, index=True), Column("history_id", Integer, ForeignKey("history.id"), index=True), Column("library_folder_id", Integer, ForeignKey("library_folder.id"), index=True), Column("tool_id", String(255)), @@ -928,7 +928,7 @@ model.HistoryDatasetCollectionAssociation.table = Table( Column("job_id", ForeignKey("job.id"), index=True, nullable=True), Column("implicit_collection_jobs_id", ForeignKey("implicit_collection_jobs.id"), index=True, nullable=True), Column("create_time", DateTime, default=now), - Column("update_time", DateTime, default=now, onupdate=now)) + Column("update_time", DateTime, default=now, onupdate=now, index=True)) model.LibraryDatasetCollectionAssociation.table = Table( "library_dataset_collection_association", metadata, @@ -991,7 +991,7 @@ model.StoredWorkflow.table = Table( "stored_workflow", metadata, Column("id", Integer, primary_key=True), Column("create_time", DateTime, default=now), - Column("update_time", DateTime, default=now, onupdate=now), + Column("update_time", DateTime, default=now, onupdate=now, index=True), Column("user_id", Integer, ForeignKey("galaxy_user.id"), index=True, nullable=False), Column("latest_workflow_id", Integer, ForeignKey("workflow.id", use_alter=True, name='stored_workflow_latest_workflow_id_fk'), index=True), @@ -1121,7 +1121,7 @@ model.WorkflowInvocation.table = Table( "workflow_invocation", metadata, Column("id", Integer, primary_key=True), Column("create_time", DateTime, default=now), - Column("update_time", DateTime, default=now, onupdate=now), + Column("update_time", DateTime, default=now, onupdate=now, index=True), Column("workflow_id", Integer, ForeignKey("workflow.id"), index=True, nullable=False), Column("state", TrimmedString(64), index=True), Column("scheduler", TrimmedString(255), index=True), diff --git a/lib/galaxy/model/migrate/versions/0176_add_indexes_on_update_time.py b/lib/galaxy/model/migrate/versions/0176_add_indexes_on_update_time.py new file mode 100644 index 00000000000..87b77802ed4 --- /dev/null +++ b/lib/galaxy/model/migrate/versions/0176_add_indexes_on_update_time.py @@ -0,0 +1,65 @@ +""" +Migration script to add indexes on update_time columns that are frequently used in ORDER BY clauses. +""" + +import logging + +from sqlalchemy import MetaData + +from galaxy.model.migrate.versions.util import ( + add_index, + drop_index +) + +log = logging.getLogger(__name__) +metadata = MetaData() + +indexes = [ + [ + "ix_history_dataset_association_update_time", + "history_dataset_association", + "update_time" + ], + [ + "ix_library_dataset_dataset_association_update_time", + "library_dataset_dataset_association", + "update_time" + ], + [ + "ix_job_update_time", + "job", + "update_time" + ], + [ + "ix_history_dataset_collection_association_update_time", + "history_dataset_collection_association", + "update_time" + ], + [ + "ix_workflow_invocation_update_time", + "workflow_invocation", + "update_time" + ], + [ + "ix_stored_workflow_update_time", + "stored_workflow", + "update_time" + ], +] + + +def upgrade(migrate_engine): + print(__doc__) + metadata.bind = migrate_engine + metadata.reflect() + + for ix, table, col in indexes: + add_index(ix, table, col, metadata) + + +def downgrade(migrate_engine): + metadata.bind = migrate_engine + metadata.reflect() + + for ix, table, col in indexes: + drop_index(ix, table, col, metadata)