mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
Merge branch 'release_21.05' into dev
This commit is contained in:
@@ -111,20 +111,13 @@ export default {
|
||||
return document.getElementById("galaxy_main");
|
||||
},
|
||||
},
|
||||
updated() {
|
||||
if (this.$refs.dropdown && this.galaxyIframe) {
|
||||
this.galaxyIframe.addEventListener("load", this.iframeListener);
|
||||
}
|
||||
mounted() {
|
||||
window.addEventListener("blur", 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();
|
||||
|
||||
@@ -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
|
||||
+1
-6
@@ -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
|
||||
|
||||
@@ -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``
|
||||
~~~~~~~~~~~~~
|
||||
@@ -1602,6 +1613,29 @@
|
||||
:Type: str
|
||||
|
||||
|
||||
~~~~~~~~~~~~~~~~~~~~~~~~~~~
|
||||
``interactivetools_prefix``
|
||||
~~~~~~~~~~~~~~~~~~~~~~~~~~~
|
||||
|
||||
:Description:
|
||||
Prefix to use in the formation of the subdomain or path for
|
||||
interactive tools
|
||||
:Default: ``interactivetool``
|
||||
: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 +2743,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``
|
||||
~~~~~~~~~~~~~~~~~~~~~~
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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__':
|
||||
|
||||
@@ -9,6 +9,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 +74,12 @@ 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."""
|
||||
timer = ExecutionTimer()
|
||||
model.HistoryAudit.prune(sa_session)
|
||||
log.debug(f"Successfully pruned history_audit table {timer}")
|
||||
|
||||
@@ -847,8 +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_proxy_host = kwargs.get("interactivetools_proxy_host", None)
|
||||
|
||||
self.containers_conf = parse_containers_config(self.containers_config_file)
|
||||
|
||||
|
||||
@@ -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'
|
||||
@@ -878,6 +882,15 @@ galaxy:
|
||||
# <data_dir>.
|
||||
#interactivetools_map: interactivetools_map.sqlite
|
||||
|
||||
# Prefix to use in the formation of the subdomain or path for
|
||||
# interactive tools
|
||||
#interactivetools_prefix: interactivetool
|
||||
|
||||
# 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 +1378,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
|
||||
|
||||
@@ -314,6 +314,15 @@
|
||||
a requirement of the other.
|
||||
|
||||
-->
|
||||
<!-- <param id="k8s_run_as_user_id">101</param> -->
|
||||
<!-- ID of user that should be used by the tool command in the job containers -->
|
||||
<!-- <param id="k8s_run_as_group_id">101</param> -->
|
||||
<!-- ID of group that should be used by the tool command in the job containers -->
|
||||
<!-- <param id="k8s_interactivetools_use_ssl">True</param> -->
|
||||
<!-- Whether to enabled HTTPS access in the k8s ingress spec generated for each interactive tools -->
|
||||
<!-- <param id="k8s_interactivetools_ingress_annotations">{'cert-manager.io/cluster-issuer': 'letsencrypt-prod'}</param> -->
|
||||
<!-- Annotations to add to the metadata section in the k8s ingress generated for each interactive tools -->
|
||||
|
||||
</plugin>
|
||||
<plugin id="godocker" type="runner" load="galaxy.jobs.runners.godocker:GodockerJobRunner">
|
||||
<!-- Go-Docker is a batch computing/cluster management tool using Docker
|
||||
@@ -1021,6 +1030,11 @@
|
||||
<!-- REQUIRED: To play nicely with the existing galaxy setup for containers. This could be set though
|
||||
internally by the runner. -->
|
||||
<param id="docker_enabled">true</param>
|
||||
<param id="container_monitor">true</param>
|
||||
<!-- Flag for appending the container monitor command to jobs. This should be set to false
|
||||
when running Interactive Tools with the kubernetes runner, but is needed for running with the
|
||||
local docker runner.
|
||||
-->
|
||||
</destination>
|
||||
<destination id="kubernetes" runner="k8s">
|
||||
<!-- <param id="requests_cpu">500m</param> -->
|
||||
|
||||
@@ -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):
|
||||
|
||||
@@ -76,13 +76,15 @@ 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),
|
||||
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()
|
||||
@@ -347,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": {
|
||||
@@ -372,6 +375,13 @@ 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):
|
||||
@@ -661,6 +671,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 +696,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 +755,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 +778,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 = f"{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 +821,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)
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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 = f"{rval.rstrip('/')}/{entry_point.entry_url.lstrip('/')}"
|
||||
if rval[0] != "/":
|
||||
|
||||
@@ -35,6 +35,7 @@ from sqlalchemy import (
|
||||
select,
|
||||
text,
|
||||
true,
|
||||
tuple_,
|
||||
type_coerce,
|
||||
types)
|
||||
from sqlalchemy.exc import OperationalError
|
||||
@@ -1845,6 +1846,27 @@ 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
|
||||
|
||||
@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):
|
||||
|
||||
dict_collection_visible_keys = ['id', 'name', 'published', 'deleted']
|
||||
@@ -1860,6 +1882,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 = []
|
||||
|
||||
+22
-11
@@ -22,6 +22,7 @@ from sqlalchemy import (
|
||||
MetaData,
|
||||
not_,
|
||||
Numeric,
|
||||
PrimaryKeyConstraint,
|
||||
select,
|
||||
String, Table,
|
||||
TEXT,
|
||||
@@ -46,10 +47,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 +203,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 +217,13 @@ model.History.table = Table(
|
||||
Index('ix_history_slug', 'slug', mysql_length=200),
|
||||
)
|
||||
|
||||
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(
|
||||
"history_user_share_association", metadata,
|
||||
Column("id", Integer, primary_key=True),
|
||||
@@ -228,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),
|
||||
@@ -505,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'),
|
||||
@@ -602,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)),
|
||||
@@ -920,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,
|
||||
@@ -983,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),
|
||||
@@ -1113,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),
|
||||
@@ -1889,7 +1897,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 +1916,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)),
|
||||
|
||||
+1
-7
@@ -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
|
||||
@@ -0,0 +1,172 @@
|
||||
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(engine) 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(engine):
|
||||
"""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 statement_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 DISTINCT {id_field}, CURRENT_TIMESTAMP AT TIME ZONE 'UTC'
|
||||
FROM new_table
|
||||
WHERE {id_field} IS NOT NULL
|
||||
ON CONFLICT DO NOTHING;
|
||||
RETURN NULL;
|
||||
END;
|
||||
$BODY$
|
||||
"""
|
||||
|
||||
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
|
||||
# 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}();
|
||||
"""
|
||||
|
||||
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))
|
||||
|
||||
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
|
||||
|
||||
|
||||
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():
|
||||
# delete old stuff first
|
||||
sql = _sqlite_remove()
|
||||
|
||||
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)
|
||||
SELECT NEW.{id_field}, strftime('%%Y-%%m-%%d %%H:%%M:%%f', 'now')
|
||||
WHERE NEW.{id_field} IS NOT NULL;
|
||||
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
|
||||
|
||||
|
||||
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}"
|
||||
@@ -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()
|
||||
|
||||
@@ -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,
|
||||
)
|
||||
|
||||
@@ -0,0 +1,83 @@
|
||||
"""
|
||||
Add history audit table and associated triggers
|
||||
"""
|
||||
|
||||
import datetime
|
||||
import logging
|
||||
|
||||
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
|
||||
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()
|
||||
|
||||
AuditTable = 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')
|
||||
)
|
||||
|
||||
|
||||
def upgrade(migrate_engine):
|
||||
print(__doc__)
|
||||
metadata.bind = migrate_engine
|
||||
metadata.reflect()
|
||||
|
||||
# create table + index
|
||||
AuditTable.drop(migrate_engine, checkfirst=True)
|
||||
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
|
||||
) 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)
|
||||
@@ -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)
|
||||
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
@@ -283,7 +283,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)
|
||||
|
||||
@@ -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
|
||||
@@ -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')
|
||||
|
||||
@@ -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/'
|
||||
|
||||
@@ -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)
|
||||
@@ -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
|
||||
|
||||
|
||||
|
||||
@@ -1,11 +1,13 @@
|
||||
"""
|
||||
API operations on the contents of a history.
|
||||
"""
|
||||
import datetime
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import re
|
||||
from datetime import datetime
|
||||
|
||||
import dateutil.parser
|
||||
|
||||
from galaxy import (
|
||||
exceptions,
|
||||
@@ -1082,8 +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_str = self.history_contents_filters.parse_date(since)
|
||||
since_date = datetime.fromisoformat(since_str)
|
||||
# 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
|
||||
|
||||
@@ -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
|
||||
@@ -1155,6 +1163,22 @@ mapping:
|
||||
desc: |
|
||||
Map for interactivetool proxy.
|
||||
|
||||
interactivetools_prefix:
|
||||
type: str
|
||||
default: "interactivetool"
|
||||
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
|
||||
|
||||
@@ -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):
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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))
|
||||
|
||||
@@ -472,6 +472,46 @@ 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")
|
||||
h1 = model.History(name="HistoryAuditHistory", user=u)
|
||||
h2 = model.History(name="HistoryAuditHistory", user=u)
|
||||
|
||||
def get_audit_table_entries(history):
|
||||
return self.session().query(model.HistoryAudit.table).filter(
|
||||
model.HistoryAudit.table.c.history_id == history.id).all()
|
||||
|
||||
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 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())
|
||||
|
||||
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
|
||||
lf = model.LibraryFolder(name="RootFolder")
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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
|
||||
Reference in New Issue
Block a user