diff --git a/config/dev.ini b/config/dev.ini new file mode 100644 index 00000000000..1bf93b5a6ba --- /dev/null +++ b/config/dev.ini @@ -0,0 +1,34 @@ +[circus] +debug = True + +[env] +GALAXY_CONFIG_FILE=config/galaxy.yml +PYTHONPATH=lib + +[watcher:web] +cmd = gunicorn 'galaxy.webapps.galaxy.fast_factory:factory()' --pythonpath lib -k uvicorn.workers.UvicornWorker -b fd://$(circus.sockets.web) +send_hup = true +numprocesses = 1 +use_sockets = True +stop_signal = TERM +stop_children = True +copy_env = True + +[watcher:client] +working_dir = client +cmd = yarn watch +numprocesses = 1 +singleton = True +copy_env = True +stop_signal = TERM +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 858fc8a21de..ddc26a4658c 100644 --- a/doc/source/admin/galaxy_options.rst +++ b/doc/source/admin/galaxy_options.rst @@ -4360,6 +4360,18 @@ :Type: str +~~~~~~~~~~~~~~~~~~~~~~~ +``enable_celery_tasks`` +~~~~~~~~~~~~~~~~~~~~~~~ + +:Description: + Offload long-running tasks to a Celery task queue. Activate this + only if you have setup a Celery worker for Galaxy. For details, + see https://docs.galaxyproject.org/en/master/admin/production.html +:Default: ``false`` +:Type: bool + + ~~~~~~~~~~~~~~ ``use_pbkdf2`` ~~~~~~~~~~~~~~ diff --git a/lib/galaxy/app.py b/lib/galaxy/app.py index eb69840d8a5..1549a7a9e82 100644 --- a/lib/galaxy/app.py +++ b/lib/galaxy/app.py @@ -2,7 +2,7 @@ import logging import signal import sys import time -from typing import Any +from typing import Any, Callable, List, Tuple from sqlalchemy.orm.scoping import ( scoped_session, @@ -70,19 +70,41 @@ from galaxy.web_stack import application_stack_instance, ApplicationStack from galaxy.webhooks import WebhooksRegistry from galaxy.workflow.trs_proxy import TrsProxy from .di import Container -from .structured_app import BasicApp, StructuredApp +from .structured_app import BasicApp, MinimalManagerApp, StructuredApp log = logging.getLogger(__name__) app = None -class UniverseApplication(StructuredApp, config.ConfiguresGalaxyMixin, Container): - """Encapsulates the state of a Universe application""" +class HaltableContainer(Container): + haltables: List[Tuple[str, Callable]] - def __init__(self, **kwargs) -> None: + def __init__(self) -> None: super().__init__() + self.haltables = [] + + def shutdown(self): + exception = None + for what, haltable in self.haltables: + try: + haltable() + except Exception as e: + log.exception(f"Failed to shutdown {what} cleanly") + exception = exception or e + if exception is not None: + raise exception + + +class MinimalGalaxyApplication(BasicApp, config.ConfiguresGalaxyMixin, HaltableContainer): + """Encapsulates the state of a minimal Galaxy application""" + + def __init__(self, fsmon=False, configure_logging=True, **kwargs) -> None: + super().__init__() + self.haltables = [ + ("object store", self._shutdown_object_store), + ("database connection", self._shutdown_model), + ] self._register_singleton(BasicApp, self) - self._register_singleton(StructuredApp, self) if not log.handlers: # Paste didn't handle it, so we need a temporary basic log # configured. The handler added here gets dumped and replaced with @@ -91,38 +113,47 @@ class UniverseApplication(StructuredApp, config.ConfiguresGalaxyMixin, Container log.debug("python path is: %s", ", ".join(sys.path)) self.name = 'galaxy' self.is_webapp = False - startup_timer = ExecutionTimer() self.new_installation = False # Read config file and check for errors self.config: Any = self._register_singleton(config.Configuration, config.Configuration(**kwargs)) self.config.check() - config.configure_logging(self.config) - self.execution_timer_factory = self._register_singleton(ExecutionTimerFactory, ExecutionTimerFactory(self.config)) - self.configure_fluent_log() - # A lot of postfork initialization depends on the server name, ensure it is set immediately after forking before other postfork functions - self.application_stack = self._register_singleton(ApplicationStack, application_stack_instance(app=self)) - self.application_stack.register_postfork_function(self.application_stack.set_postfork_server_name, self) - self.config.reload_sanitize_allowlist(explicit='sanitize_allowlist_file' in kwargs) - self.amqp_internal_connection_obj = galaxy.queues.connection_from_config(self.config) - # queue_worker *can* be initialized with a queue, but here we don't - # want to and we'll allow postfork to bind and start it. - self.queue_worker = self._register_singleton(GalaxyQueueWorker, GalaxyQueueWorker(self)) - - self._configure_tool_shed_registry() + if configure_logging: + config.configure_logging(self.config) self._configure_object_store(fsmon=True) - # Setup the database engine and ORM config_file = kwargs.get('global_conf', {}).get('__file__', None) if config_file: log.debug('Using "galaxy.ini" config file: %s', config_file) check_migrate_tools = self.config.check_migrate_tools self._configure_models(check_migrate_databases=self.config.check_migrate_databases, check_migrate_tools=check_migrate_tools, config_file=config_file) - # Security helper self._configure_security() self._register_singleton(IdEncodingHelper, self.security) self._register_singleton(SharedModelMapping, self.model) self._register_singleton(GalaxyModelMapping, self.model) self._register_singleton(scoped_session, self.model.context) + + def configure_fluent_log(self): + if self.config.fluent_log: + from galaxy.util.custom_logging.fluent_log import FluentTraceLogger + self.trace_logger = FluentTraceLogger('galaxy', self.config.fluent_host, self.config.fluent_port) + else: + self.trace_logger = None + + def _shutdown_object_store(self): + self.object_store.shutdown() + + def _shutdown_model(self): + self.model.engine.dispose() + + +class GalaxyManagerApplication(MinimalManagerApp, MinimalGalaxyApplication): + """Extends the MinimalGalaxyApplication with most managers that are not tied to a web or job handling context.""" + def __init__(self, **kwargs): + super().__init__(**kwargs) + self._register_singleton(MinimalManagerApp, self) + self.execution_timer_factory = self._register_singleton(ExecutionTimerFactory, ExecutionTimerFactory(self.config)) + self.configure_fluent_log() + # Tag handler self.tag_handler = self._register_singleton(GalaxyTagHandler) self.user_manager = self._register_singleton(UserManager) @@ -133,17 +164,66 @@ class UniverseApplication(StructuredApp, config.ConfiguresGalaxyMixin, Container self.dataset_collections_service = self._register_singleton(DatasetCollectionManager) self.workflow_manager = self._register_singleton(WorkflowsManager) self.workflow_contents_manager = self._register_singleton(WorkflowContentsManager) - self.dependency_resolvers_view = self._register_singleton(DependencyResolversView, DependencyResolversView(self)) - self.test_data_resolver = self._register_singleton(TestDataResolver, TestDataResolver(file_dirs=self.config.tool_test_data_directories)) self.library_folder_manager = self._register_singleton(FolderManager) self.library_manager = self._register_singleton(LibraryManager) self.role_manager = self._register_singleton(RoleManager) - self.dynamic_tool_manager = self._register_singleton(DynamicToolManager) - self.api_keys_manager = self._register_singleton(ApiKeyManager) # ConfiguredFileSources self.file_sources = self._register_singleton(ConfiguredFileSources, ConfiguredFileSources.from_app_config(self.config)) + # We need the datatype registry for running certain tasks that modify HDAs, and to build the registry we need + # to setup the installed repositories ... this is not ideal + self._configure_tool_config_files() + self.installed_repository_manager = self._register_singleton(InstalledRepositoryManager, InstalledRepositoryManager(self)) + self._configure_datatypes_registry(self.installed_repository_manager) + self._register_singleton(Registry, self.datatypes_registry) + galaxy.model.set_datatypes_registry(self.datatypes_registry) + + self.sentry_client = None + if self.config.sentry_dsn: + + def postfork_sentry_client(): + import raven + self.sentry_client = raven.Client(self.config.sentry_dsn, transport=raven.transport.HTTPTransport) + + self.application_stack.register_postfork_function(postfork_sentry_client) + + +class UniverseApplication(StructuredApp, GalaxyManagerApplication): + """Encapsulates the state of a Universe application""" + + def __init__(self, **kwargs) -> None: + startup_timer = ExecutionTimer() + super().__init__(fsmon=True, **kwargs) + self.haltables = [ + ("queue worker", self._shutdown_queue_worker), + ("file watcher", self._shutdown_watcher), + ("database heartbeat", self._shutdown_database_heartbeat), + ("workflow scheduler", self._shutdown_scheduling_manager), + ("object store", self._shutdown_object_store), + ("job manager", self._shutdown_job_manager), + ("application heartbeat", self._shutdown_heartbeat), + ("repository manager", self._shutdown_repo_manager), + ("database connection", self._shutdown_model), + ("application stack", self._shutdown_application_stack), + ] + self._register_singleton(StructuredApp, self) + # A lot of postfork initialization depends on the server name, ensure it is set immediately after forking before other postfork functions + self.application_stack = self._register_singleton(ApplicationStack, application_stack_instance(app=self)) + self.application_stack.register_postfork_function(self.application_stack.set_postfork_server_name, self) + self.config.reload_sanitize_allowlist(explicit='sanitize_allowlist_file' in kwargs) + self.amqp_internal_connection_obj = galaxy.queues.connection_from_config(self.config) + # queue_worker *can* be initialized with a queue, but here we don't + # want to and we'll allow postfork to bind and start it. + self.queue_worker = self._register_singleton(GalaxyQueueWorker, GalaxyQueueWorker(self)) + + self._configure_tool_shed_registry() + + self.dependency_resolvers_view = self._register_singleton(DependencyResolversView, DependencyResolversView(self)) + self.test_data_resolver = self._register_singleton(TestDataResolver, TestDataResolver(file_dirs=self.config.tool_test_data_directories)) + self.dynamic_tool_manager = self._register_singleton(DynamicToolManager) + self.api_keys_manager = self._register_singleton(ApiKeyManager) + # Tool Data Tables self._configure_tool_data_tables(from_shed_config=False) # Load dbkey / genome build manager @@ -169,14 +249,7 @@ class UniverseApplication(StructuredApp, config.ConfiguresGalaxyMixin, Container self.tool_shed_repository_cache = self._register_singleton(ToolShedRepositoryCache) # Watch various config files for immediate reload self.watchers = self._register_singleton(ConfigWatchers) - self._configure_tool_config_files() - self.installed_repository_manager = self._register_singleton(InstalledRepositoryManager, InstalledRepositoryManager(self)) - self._configure_datatypes_registry(self.installed_repository_manager) - self._register_singleton(Registry, self.datatypes_registry) - galaxy.model.set_datatypes_registry(self.datatypes_registry) - self._configure_toolbox() - # Load Data Manager self.data_managers = self._register_singleton(DataManagers) # Load the update repository manager. @@ -231,16 +304,6 @@ class UniverseApplication(StructuredApp, config.ConfiguresGalaxyMixin, Container self.authnz_manager = managers.AuthnzManager(self, self.config.oidc_config_file, self.config.oidc_backends_config_file) - - self.sentry_client = None - if self.config.sentry_dsn: - - def postfork_sentry_client(): - import raven - self.sentry_client = raven.Client(self.config.sentry_dsn, transport=raven.transport.HTTPTransport) - - self.application_stack.register_postfork_function(postfork_sentry_client) - # Start the job manager from galaxy.jobs import manager self.job_manager = self._register_singleton(manager.JobManager) @@ -284,8 +347,6 @@ class UniverseApplication(StructuredApp, config.ConfiguresGalaxyMixin, Container # Delay toolbox index until after startup self.application_stack.register_postfork_function(lambda: send_local_control_task(self, 'rebuild_toolbox_search_index')) - self.model.engine.dispose() - # Inject url_for for components to more easily optionally depend # on url_for. self.url_for = url_for @@ -293,74 +354,30 @@ class UniverseApplication(StructuredApp, config.ConfiguresGalaxyMixin, Container self.server_starttime = int(time.time()) # used for cachebusting log.info("Galaxy app startup finished %s" % startup_timer) - def shutdown(self): - log.debug('Shutting down') - exception = None - try: - self.queue_worker.shutdown() - except Exception as e: - exception = exception or e - log.exception("Failed to shutdown control worker cleanly") - try: - self.watchers.shutdown() - except Exception as e: - exception = exception or e - log.exception("Failed to shutdown configuration watchers cleanly") - try: - self.database_heartbeat.shutdown() - except Exception as e: - exception = exception or e - log.exception("Failed to shutdown database heartbeat cleanly") - try: - self.workflow_scheduling_manager.shutdown() - except Exception as e: - exception = exception or e - log.exception("Failed to shutdown workflow scheduling manager cleanly") - try: - self.job_manager.shutdown() - except Exception as e: - exception = exception or e - log.exception("Failed to shutdown job manager cleanly") - try: - self.object_store.shutdown() - except Exception as e: - exception = exception or e - log.exception("Failed to shutdown object store cleanly") - try: - if self.heartbeat: - self.heartbeat.shutdown() - except Exception as e: - exception = exception or e - log.exception("Failed to shutdown heartbeat cleanly") - try: - self.update_repository_manager.shutdown() - except Exception as e: - exception = exception or e - log.exception("Failed to shutdown update repository manager cleanly") + def _shutdown_queue_worker(self): + self.queue_worker.shutdown() - try: - self.model.engine.dispose() - except Exception as e: - exception = exception or e - log.exception("Failed to shutdown SA database engine cleanly") + def _shutdown_watcher(self): + self.watchers.shutdown() - try: - self.application_stack.shutdown() - except Exception as e: - exception = exception or e - log.exception("Failed to shutdown application stack interface cleanly") + def _shutdown_database_heartbeat(self): + self.database_heartbeat.shutdown() - if exception: - raise exception - else: - log.debug('Finished shutting down') + def _shutdown_scheduling_manager(self): + self.workflow_scheduling_manager.shutdown() - def configure_fluent_log(self): - if self.config.fluent_log: - from galaxy.util.custom_logging.fluent_log import FluentTraceLogger - self.trace_logger = FluentTraceLogger('galaxy', self.config.fluent_host, self.config.fluent_port) - else: - self.trace_logger = None + def _shutdown_job_manager(self): + self.job_manager.shutdown() + + def _shutdown_heartbeat(self): + if self.heartbeat: + self.heartbeat.shutdown() + + def _shutdown_repo_manager(self): + self.update_repository_manager.shutdown() + + def _shutdown_application_stack(self): + self.application_stack.shutdown() @property def is_job_handler(self) -> bool: diff --git a/lib/galaxy/celery/__init__.py b/lib/galaxy/celery/__init__.py new file mode 100644 index 00000000000..59d4bc410c7 --- /dev/null +++ b/lib/galaxy/celery/__init__.py @@ -0,0 +1,55 @@ +import os +from functools import lru_cache + +from celery import Celery + +from galaxy.config import Configuration +from galaxy.util.custom_logging import get_logger +from galaxy.util.properties import load_app_properties + +log = get_logger(__name__) + + +@lru_cache(maxsize=1) +def get_galaxy_app(): + import galaxy.app + if galaxy.app.app: + return galaxy.app.app + kwargs = get_app_properties() + if kwargs: + kwargs['check_migrate_tools'] = False + kwargs['check_migrate_databases'] = False + galaxy_app = galaxy.app.GalaxyManagerApplication(configure_logging=False, **kwargs) + return galaxy_app + + +@lru_cache(maxsize=1) +def get_app_properties(): + config_file = os.environ.get("GALAXY_CONFIG_FILE") + if config_file: + return load_app_properties( + config_file=os.path.abspath(config_file), + config_section='galaxy', + ) + + +@lru_cache(maxsize=1) +def get_config(): + kwargs = get_app_properties() + if kwargs: + kwargs['override_tempdir'] = False + return Configuration(**kwargs) + + +def get_broker(): + config = get_config() + if config: + return config.amqp_internal_connection + + +broker = get_broker() +celery_app = Celery('galaxy', broker=broker, include=['galaxy.celery.tasks']) + + +if __name__ == '__main__': + celery_app.start() diff --git a/lib/galaxy/celery/tasks.py b/lib/galaxy/celery/tasks.py new file mode 100644 index 00000000000..5f0f02e0a22 --- /dev/null +++ b/lib/galaxy/celery/tasks.py @@ -0,0 +1,51 @@ +from lagom import magic_bind_to_container +from sqlalchemy.orm.scoping import ( + scoped_session, +) + +from galaxy.celery import celery_app +from galaxy.managers.hdas import HDAManager +from galaxy.managers.lddas import LDDAManager +from galaxy.model import User +from galaxy.util.custom_logging import get_logger +from . import get_galaxy_app + +log = get_logger(__name__) + + +def galaxy_task(func): + app = get_galaxy_app() + if app: + return magic_bind_to_container(app)(func) + return func + + +@celery_app.task(ignore_result=True) +@galaxy_task +def recalculate_user_disk_usage(session: scoped_session, user_id=None): + if user_id: + user = session.query(User).get(user_id) + if user: + user.calculate_and_set_disk_usage() + log.info(f"New user disk usage is {user.disk_usage}") + else: + log.error("Recalculate user disk usage task failed, user %s not found" % user_id) + else: + log.error("Recalculate user disk usage task received without user_id.") + + +@celery_app.task(ignore_result=True) +@galaxy_task +def purge_hda(hda_manager: HDAManager, hda_id): + hda = hda_manager.by_id(hda_id) + hda_manager._purge(hda) + + +@celery_app.task +@galaxy_task +def set_metadata(hda_manager: HDAManager, ldda_manager: LDDAManager, dataset_id, model_class='HistoryDatasetAssociation'): + if model_class == 'HistoryDatasetAssociation': + dataset = hda_manager.by_id(dataset_id) + elif model_class == 'LibraryDatasetDatasetAssociation': + dataset = ldda_manager.by_id(dataset_id) + dataset.datatype.set_meta(dataset) diff --git a/lib/galaxy/config/sample/galaxy.yml.sample b/lib/galaxy/config/sample/galaxy.yml.sample index e630d68e858..a65ad1df8fb 100644 --- a/lib/galaxy/config/sample/galaxy.yml.sample +++ b/lib/galaxy/config/sample/galaxy.yml.sample @@ -2159,6 +2159,11 @@ galaxy: # commented out line below). #amqp_internal_connection: sqlalchemy+sqlite:///./database/control.sqlite?isolation_level=IMMEDIATE + # Offload long-running tasks to a Celery task queue. Activate this + # only if you have setup a Celery worker for Galaxy. For details, see + # https://docs.galaxyproject.org/en/master/admin/production.html + #enable_celery_tasks: false + # Allow disabling pbkdf2 hashing of passwords for legacy situations. # This should normally be left enabled unless there is a specific # reason to disable it. diff --git a/lib/galaxy/dependencies/dev-requirements.txt b/lib/galaxy/dependencies/dev-requirements.txt index 82e26ec3016..55047974fff 100644 --- a/lib/galaxy/dependencies/dev-requirements.txt +++ b/lib/galaxy/dependencies/dev-requirements.txt @@ -17,6 +17,7 @@ bagit==1.8.1; python_version >= "3.6" and python_full_version < "3.0.0" and pyth bcrypt==3.2.0; python_version >= "3.6" bdbag==1.6.1; (python_version >= "2.7" and python_full_version < "3.0.0") or (python_full_version >= "3.5.0" and python_version < "4") beaker==1.11.0 +billiard==3.6.3.0; python_version >= "3.6" bioblend==0.15.0; python_version >= "3.6" bleach==3.3.0; (python_version >= "2.7" and python_full_version < "3.0.0") or (python_full_version >= "3.5.0") boltons==20.2.1 @@ -27,11 +28,16 @@ bx-python==0.8.11; python_version >= "3.6" cachecontrol==0.11.7; python_version >= "3.6" and python_version < "4" cached-property==1.5.2; python_version < "3.8" and python_version >= "3.6" cachetools==4.2.1; python_version >= "3.5" and python_version < "4.0" and (python_version >= "2.7" and python_full_version < "3.0.0" or python_full_version >= "3.6.0") +celery==5.0.5; python_version >= "3.6" certifi==2020.12.5; python_version >= "3.6" and python_full_version < "3.0.0" or python_full_version >= "3.5.0" and python_version < "4" and python_version >= "3.6" -cffi==1.14.5; python_version >= "3.6" and python_full_version < "3.0.0" or python_full_version >= "3.4.0" and python_version >= "3.6" +cffi==1.14.5; implementation_name == "pypy" and python_version >= "3.6" and (python_version >= "2.7" and python_full_version < "3.0.0" or python_full_version >= "3.4.0") chardet==4.0.0; python_version >= "2.7" and python_full_version < "3.0.0" or python_full_version >= "3.5.0" and python_version < "4" cheetah3==3.2.6.post1; (python_version >= "2.7" and python_full_version < "3.0.0") or (python_full_version >= "3.4.0") -click==7.1.2; python_version >= "2.7" and python_full_version < "3.0.0" or python_full_version >= "3.5.0" +circus==0.17.1 +click-didyoumean==0.0.3; python_version >= "3.6" +click-plugins==1.1.1; python_version >= "3.6" +click-repl==0.1.6; python_version >= "3.6" +click==7.1.2; python_version >= "3.6" and python_full_version < "3.0.0" or python_full_version >= "3.5.0" and python_version >= "3.6" cliff==3.7.0; python_version >= "3.6" cloudauthz==0.6.0 cloudbridge==2.1.0 @@ -132,14 +138,15 @@ pbr==5.5.1; python_version >= "3.6" pluggy==0.13.1; python_version >= "3.6" and python_full_version < "3.0.0" or python_full_version >= "3.4.0" and python_version >= "3.6" port-for==0.4; python_version >= "3.6" prettytable==0.7.2; python_version >= "3.6" +prompt-toolkit==3.0.3; python_version >= "3.6" protobuf==3.15.6; python_version >= "3.6" and python_full_version < "3.0.0" or python_full_version >= "3.6.0" and python_version >= "3.6" prov==1.5.1; python_version >= "3.6" and python_version < "4" psutil==5.8.0; (python_version >= "2.6" and python_full_version < "3.0.0") or (python_full_version >= "3.4.0") pulsar-galaxy-lib==0.14.2 -py==1.10.0; python_version >= "3.6" and python_full_version < "3.0.0" or python_full_version >= "3.4.0" and python_version >= "3.6" +py==1.10.0; python_version >= "3.6" and python_full_version < "3.0.0" and implementation_name == "pypy" or python_full_version >= "3.4.0" and python_version >= "3.6" and implementation_name == "pypy" pyasn1-modules==0.2.8; python_version >= "2.7" and python_full_version < "3.0.0" or python_full_version >= "3.6.0" pyasn1==0.4.8; python_version >= "3.5" and python_version < "4" -pycparser==2.20; python_version >= "3.6" and python_full_version < "3.0.0" or python_full_version >= "3.4.0" and python_version >= "3.6" +pycparser==2.20; python_version >= "3.6" and python_full_version < "3.0.0" and implementation_name == "pypy" or implementation_name == "pypy" and python_version >= "3.6" and python_full_version >= "3.4.0" pycryptodome==3.10.1; (python_version >= "2.7" and python_full_version < "3.0.0") or (python_full_version >= "3.5.0") pydantic==1.7.3; python_version >= "3.6" and python_version < "4.0" pydot==1.4.2; python_version >= "3.6" and python_full_version < "3.0.0" and python_version < "4" or python_version >= "3.6" and python_version < "4" and python_full_version >= "3.4.0" @@ -178,6 +185,7 @@ python3-openid==3.2.0; python_version >= "3.0" pytz==2021.1; python_version >= "3.6" and python_full_version < "3.0.0" or python_full_version >= "3.6.0" and python_version < "4" and python_version >= "3.6" pyuwsgi==2.0.19.1.post0 pyyaml==5.4.1; (python_version >= "2.7" and python_full_version < "3.0.0") or (python_full_version >= "3.6.0") +pyzmq==22.0.3; python_version >= "3.6" rdflib-jsonld==0.5.0; python_version >= "3.6" and python_version < "4" rdflib==4.2.2; python_version >= "3.6" and python_version < "4" recommonmark==0.7.1 @@ -226,6 +234,7 @@ tempita==0.5.2 tenacity==7.0.0 testfixtures==6.17.1 toml==0.10.2; python_version >= "3.6" and python_full_version < "3.0.0" or python_full_version >= "3.3.0" and python_version >= "3.6" +tornado==6.1; python_version >= "3.5" tqdm==4.59.0; python_version >= "2.7" and python_full_version < "3.0.0" or python_full_version >= "3.4.0" twill==3.0 typing-extensions==3.7.4.3; python_version >= "3.6" and python_version < "3.8" diff --git a/lib/galaxy/dependencies/pinned-requirements.txt b/lib/galaxy/dependencies/pinned-requirements.txt index c4b41ed44ed..f52b8054d07 100644 --- a/lib/galaxy/dependencies/pinned-requirements.txt +++ b/lib/galaxy/dependencies/pinned-requirements.txt @@ -15,6 +15,7 @@ bagit==1.8.1; python_version >= "3.6" and python_full_version < "3.0.0" and pyth bcrypt==3.2.0; python_version >= "3.6" bdbag==1.6.1; (python_version >= "2.7" and python_full_version < "3.0.0") or (python_full_version >= "3.5.0" and python_version < "4") beaker==1.11.0 +billiard==3.6.3.0; python_version >= "3.6" bioblend==0.15.0; python_version >= "3.6" bleach==3.3.0; (python_version >= "2.7" and python_full_version < "3.0.0") or (python_full_version >= "3.5.0") boltons==20.2.1 @@ -25,11 +26,16 @@ bx-python==0.8.11; python_version >= "3.6" cachecontrol==0.11.7; python_version >= "3.6" and python_version < "4" cached-property==1.5.2; python_version < "3.8" and python_version >= "3.6" cachetools==4.2.1; python_version >= "3.5" and python_version < "4.0" and (python_version >= "2.7" and python_full_version < "3.0.0" or python_full_version >= "3.6.0") +celery==5.0.5; python_version >= "3.6" certifi==2020.12.5; python_version >= "2.7" and python_full_version < "3.0.0" or python_full_version >= "3.5.0" and python_version < "4" -cffi==1.14.5; python_version >= "3.6" and python_full_version < "3.0.0" or python_full_version >= "3.4.0" and python_version >= "3.6" +cffi==1.14.5; implementation_name == "pypy" and python_version >= "3.6" and (python_version >= "2.7" and python_full_version < "3.0.0" or python_full_version >= "3.4.0") chardet==4.0.0; python_version >= "2.7" and python_full_version < "3.0.0" or python_full_version >= "3.5.0" and python_version < "4" cheetah3==3.2.6.post1; (python_version >= "2.7" and python_full_version < "3.0.0") or (python_full_version >= "3.4.0") -click==7.1.2; python_version >= "2.7" and python_full_version < "3.0.0" or python_full_version >= "3.5.0" +circus==0.17.1 +click-didyoumean==0.0.3; python_version >= "3.6" +click-plugins==1.1.1; python_version >= "3.6" +click-repl==0.1.6; python_version >= "3.6" +click==7.1.2; python_version >= "3.6" and python_full_version < "3.0.0" or python_full_version >= "3.5.0" and python_version >= "3.6" cliff==3.7.0; python_version >= "3.6" cloudauthz==0.6.0 cloudbridge==2.1.0 @@ -115,13 +121,15 @@ paste==3.5.0 pastedeploy==2.1.1 pbr==5.5.1; python_version >= "3.6" prettytable==0.7.2; python_version >= "3.6" +prompt-toolkit==3.0.3; python_version >= "3.6" protobuf==3.15.6; python_version >= "3.6" and python_full_version < "3.0.0" or python_full_version >= "3.6.0" and python_version >= "3.6" prov==1.5.1; python_version >= "3.6" and python_version < "4" psutil==5.8.0; (python_version >= "2.6" and python_full_version < "3.0.0") or (python_full_version >= "3.4.0") pulsar-galaxy-lib==0.14.2 +py==1.10.0; python_version >= "3.6" and python_full_version < "3.0.0" and implementation_name == "pypy" or implementation_name == "pypy" and python_version >= "3.6" and python_full_version >= "3.4.0" pyasn1-modules==0.2.8; python_version >= "2.7" and python_full_version < "3.0.0" or python_full_version >= "3.6.0" pyasn1==0.4.8; python_version >= "3.5" and python_version < "4" -pycparser==2.20; python_version >= "3.6" and python_full_version < "3.0.0" or python_full_version >= "3.4.0" and python_version >= "3.6" +pycparser==2.20; python_version >= "3.6" and python_full_version < "3.0.0" and implementation_name == "pypy" or implementation_name == "pypy" and python_version >= "3.6" and python_full_version >= "3.4.0" pycryptodome==3.10.1; (python_version >= "2.7" and python_full_version < "3.0.0") or (python_full_version >= "3.5.0") pydantic==1.7.3; python_version >= "3.6" and python_version < "4.0" pydot==1.4.2; python_version >= "3.6" and python_full_version < "3.0.0" and python_version < "4" or python_version >= "3.6" and python_version < "4" and python_full_version >= "3.4.0" @@ -147,6 +155,7 @@ python3-openid==3.2.0; python_version >= "3.0" pytz==2021.1; python_version >= "3.6" and python_full_version < "3.0.0" or python_full_version >= "3.6.0" and python_version < "4" and python_version >= "3.6" pyuwsgi==2.0.19.1.post0 pyyaml==5.4.1; (python_version >= "2.7" and python_full_version < "3.0.0") or (python_full_version >= "3.6.0") +pyzmq==22.0.3; python_version >= "3.6" rdflib-jsonld==0.5.0; python_version >= "3.6" and python_version < "4" rdflib==4.2.2; python_version >= "3.6" and python_version < "4" refgenconf==0.9.3 @@ -178,6 +187,7 @@ stevedore==3.3.0; python_version >= "3.6" svgwrite==1.4.1; python_version >= "3.6" tempita==0.5.2 tenacity==7.0.0 +tornado==6.1; python_version >= "3.5" tqdm==4.59.0; python_version >= "2.7" and python_full_version < "3.0.0" or python_full_version >= "3.4.0" typing-extensions==3.7.4.3; python_version >= "3.6" and python_version < "3.8" tzlocal==2.1; python_version >= "2.7" and python_full_version < "3.0.0" or python_full_version >= "3.5.0" and python_version < "4" diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index c8a6d25ed4b..2ef3f7a3f63 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -57,7 +57,7 @@ from galaxy.jobs.runners import BaseJobRunner, JobState from galaxy.metadata import get_metadata_compute_strategy from galaxy.model import store from galaxy.objectstore import ObjectStorePopulator -from galaxy.structured_app import StructuredApp +from galaxy.structured_app import MinimalManagerApp from galaxy.tool_util.deps import requirements from galaxy.tool_util.output_checker import ( check_output, @@ -301,7 +301,7 @@ class JobConfiguration(ConfiguresHandlers): """ - def __init__(self, app: StructuredApp): + def __init__(self, app: MinimalManagerApp): """Parse the job configuration XML. """ self.app = app diff --git a/lib/galaxy/jobs/manager.py b/lib/galaxy/jobs/manager.py index a9740288e75..fab77182b5d 100644 --- a/lib/galaxy/jobs/manager.py +++ b/lib/galaxy/jobs/manager.py @@ -9,7 +9,7 @@ from sqlalchemy.sql.expression import null from galaxy.exceptions import HandlerAssignmentError, ToolExecutionError from galaxy.jobs import handler, NoopQueue from galaxy.model import Job -from galaxy.structured_app import StructuredApp +from galaxy.structured_app import MinimalManagerApp from galaxy.web_stack.message import JobHandlerMessage log = logging.getLogger(__name__) @@ -21,7 +21,7 @@ class JobManager: """ job_handler: handler.JobHandlerI - def __init__(self, app: StructuredApp): + def __init__(self, app: MinimalManagerApp): self.app = app self.job_lock = False if self.app.is_job_handler: diff --git a/lib/galaxy/managers/base.py b/lib/galaxy/managers/base.py index de18f4c279c..bac9e9198ab 100644 --- a/lib/galaxy/managers/base.py +++ b/lib/galaxy/managers/base.py @@ -37,7 +37,7 @@ from sqlalchemy.orm.scoping import scoped_session from galaxy import exceptions from galaxy import model from galaxy.model import tool_shed_install -from galaxy.structured_app import BasicApp, StructuredApp +from galaxy.structured_app import BasicApp, MinimalManagerApp from galaxy.util import namedtuple log = logging.getLogger(__name__) @@ -488,7 +488,7 @@ class HasAModelManager: # examples where this doesn't really work are ConfigurationSerializer (no manager) # and contents (2 managers) - def __init__(self, app: StructuredApp, manager=None, **kwargs): + def __init__(self, app: MinimalManagerApp, manager=None, **kwargs): self._manager = manager @property @@ -542,7 +542,7 @@ class ModelSerializer(HasAModelManager): default_view: Optional[str] views: Dict[str, List[str]] - def __init__(self, app: StructuredApp, **kwargs): + def __init__(self, app: MinimalManagerApp, **kwargs): """ Set up serializer map, any additional serializable keys, and views here. """ @@ -713,7 +713,7 @@ class ModelDeserializer(HasAModelManager): """ # TODO:?? a larger question is: which should be first? Deserialize then validate - or - validate then deserialize? - def __init__(self, app: StructuredApp, validator=None, **kwargs): + def __init__(self, app: MinimalManagerApp, validator=None, **kwargs): """ Set up deserializers and validator. """ @@ -904,7 +904,7 @@ class ModelFilterParser(HasAModelManager): orm_filter_parsers: Dict[str, Dict] fn_filter_parsers: Dict[str, Dict] - def __init__(self, app: StructuredApp, **kwargs): + def __init__(self, app: MinimalManagerApp, **kwargs): """ Set up serializer map, any additional serializable keys, and views here. """ diff --git a/lib/galaxy/managers/configuration.py b/lib/galaxy/managers/configuration.py index 13ddaa3b7a9..cbf2bd1be09 100644 --- a/lib/galaxy/managers/configuration.py +++ b/lib/galaxy/managers/configuration.py @@ -15,7 +15,7 @@ from typing import ( List, ) -from galaxy.app import StructuredApp +from galaxy.app import MinimalManagerApp from galaxy.managers import base from galaxy.managers.context import ProvidesUserContext from galaxy.schema.fields import EncodedDatabaseIdField @@ -30,7 +30,7 @@ VERSION_JSON_FILE = 'version.json' class ConfigurationManager: """Interface/service object for interacting with configuration and related data.""" - def __init__(self, app: StructuredApp): + def __init__(self, app: MinimalManagerApp): self._app = app def get_configuration( diff --git a/lib/galaxy/managers/context.py b/lib/galaxy/managers/context.py index 4b9069a6f85..cbdf15cbfb5 100644 --- a/lib/galaxy/managers/context.py +++ b/lib/galaxy/managers/context.py @@ -50,7 +50,7 @@ from galaxy.model import ( ) from galaxy.model.base import ModelMapping from galaxy.security.idencoding import IdEncodingHelper -from galaxy.structured_app import StructuredApp +from galaxy.structured_app import MinimalManagerApp from galaxy.util import bunch @@ -62,7 +62,7 @@ class ProvidesAppContext: """ @abc.abstractproperty - def app(self) -> StructuredApp: + def app(self) -> MinimalManagerApp: """Provide access to the Galaxy ``app`` object. """ diff --git a/lib/galaxy/managers/datasets.py b/lib/galaxy/managers/datasets.py index 0f7092e3a5e..483e13004fc 100644 --- a/lib/galaxy/managers/datasets.py +++ b/lib/galaxy/managers/datasets.py @@ -18,7 +18,7 @@ from galaxy.managers import ( secured, users ) -from galaxy.structured_app import StructuredApp +from galaxy.structured_app import MinimalManagerApp from galaxy.util.checkers import check_binary log = logging.getLogger(__name__) @@ -33,7 +33,7 @@ class DatasetManager(base.ModelManager, secured.AccessibleManagerMixin, deletabl # TODO:?? get + error_if_uploading is common pattern, should upload check be worked into access/owed? - def __init__(self, app: StructuredApp): + def __init__(self, app: MinimalManagerApp): super().__init__(app) self.permissions = DatasetRBACPermissions(app) # needed for admin test @@ -143,7 +143,7 @@ class DatasetRBACPermissions: class DatasetSerializer(base.ModelSerializer, deletable.PurgableSerializerMixin): model_manager_class = DatasetManager - def __init__(self, app: StructuredApp, user_manager: users.UserManager): + def __init__(self, app: MinimalManagerApp, user_manager: users.UserManager): super().__init__(app) self.dataset_manager = self.manager # needed for admin test @@ -302,8 +302,10 @@ class DatasetAssociationManager(base.ModelManager, if not job.finished: # Are *all* of the job's other output datasets deleted? if job.check_if_output_datasets_deleted(): - job.mark_deleted(self.app.config.track_jobs_in_database) - self.app.job_manager.stop(job) + track_jobs_in_database = self.app.config.track_jobs_in_database + job.mark_deleted(track_jobs_in_database) + if not track_jobs_in_database: + self.app.job_manager.stop(job) return True return False diff --git a/lib/galaxy/managers/group_roles.py b/lib/galaxy/managers/group_roles.py index 55ad0475729..a6e488272d0 100644 --- a/lib/galaxy/managers/group_roles.py +++ b/lib/galaxy/managers/group_roles.py @@ -7,7 +7,7 @@ from typing import ( ) from galaxy import model -from galaxy.app import StructuredApp +from galaxy.app import MinimalManagerApp from galaxy.exceptions import ( ObjectNotFound, ) @@ -22,7 +22,7 @@ log = logging.getLogger(__name__) class GroupRolesManager: """Interface/service object shared by controllers for interacting with group roles.""" - def __init__(self, app: StructuredApp) -> None: + def __init__(self, app: MinimalManagerApp) -> None: self._app = app def index(self, trans: ProvidesAppContext, group_id: EncodedDatabaseIdField) -> List[Dict[str, Any]]: diff --git a/lib/galaxy/managers/group_users.py b/lib/galaxy/managers/group_users.py index 9dd194311ef..1bbef7ec4a2 100644 --- a/lib/galaxy/managers/group_users.py +++ b/lib/galaxy/managers/group_users.py @@ -7,7 +7,7 @@ from typing import ( ) from galaxy import model -from galaxy.app import StructuredApp +from galaxy.app import MinimalManagerApp from galaxy.exceptions import ( ObjectNotFound, ) @@ -22,7 +22,7 @@ log = logging.getLogger(__name__) class GroupUsersManager: """Interface/service object shared by controllers for interacting with group users.""" - def __init__(self, app: StructuredApp) -> None: + def __init__(self, app: MinimalManagerApp) -> None: self._app = app def index(self, trans: ProvidesAppContext, group_id: EncodedDatabaseIdField) -> List[Dict[str, Any]]: diff --git a/lib/galaxy/managers/groups.py b/lib/galaxy/managers/groups.py index 8162bf3ebfe..a35abb4838c 100644 --- a/lib/galaxy/managers/groups.py +++ b/lib/galaxy/managers/groups.py @@ -7,7 +7,7 @@ from typing import ( from sqlalchemy import false from galaxy import model -from galaxy.app import StructuredApp +from galaxy.app import MinimalManagerApp from galaxy.exceptions import ( Conflict, ObjectAttributeMissingException, @@ -22,7 +22,7 @@ from galaxy.web import url_for class GroupsManager: """Interface/service object shared by controllers for interacting with groups.""" - def __init__(self, app: StructuredApp) -> None: + def __init__(self, app: MinimalManagerApp) -> None: self._app = app def index(self, trans: ProvidesAppContext): diff --git a/lib/galaxy/managers/hdas.py b/lib/galaxy/managers/hdas.py index c4d4093a452..dc6c2b2aeba 100644 --- a/lib/galaxy/managers/hdas.py +++ b/lib/galaxy/managers/hdas.py @@ -20,7 +20,7 @@ from galaxy.managers import ( taggable, users, ) -from galaxy.structured_app import StructuredApp +from galaxy.structured_app import MinimalManagerApp log = logging.getLogger(__name__) @@ -41,7 +41,7 @@ class HDAManager(datasets.DatasetAssociationManager, # TODO: move what makes sense into DatasetManager # TODO: which of these are common with LDDAs and can be pushed down into DatasetAssociationManager? - def __init__(self, app: StructuredApp, user_manager: users.UserManager): + def __init__(self, app: MinimalManagerApp, user_manager: users.UserManager): """ Set up and initialize other managers needed by hdas. """ @@ -135,6 +135,13 @@ class HDAManager(datasets.DatasetAssociationManager, # .... deletion and purging def purge(self, hda, flush=True): + if self.app.config.enable_celery_tasks: + from galaxy.celery.tasks import purge_hda + purge_hda.delay(hda_id=hda.id) + else: + self._purge(hda, flush=flush) + + def _purge(self, hda, flush=True): """ Purge this HDA and the dataset underlying it. """ @@ -146,7 +153,6 @@ class HDAManager(datasets.DatasetAssociationManager, # decrease the user's space used if quota_amount_reduction: user.adjust_total_disk_usage(-quota_amount_reduction) - return hda # .... states def error_if_uploading(self, hda): @@ -248,7 +254,7 @@ class HDASerializer( # datasets._UnflattenedMetadataDatasetAssociationSerialize annotatable.AnnotatableSerializerMixin): model_manager_class = HDAManager - def __init__(self, app: StructuredApp): + def __init__(self, app: MinimalManagerApp): super().__init__(app) self.hda_manager = self.manager @@ -507,7 +513,7 @@ class HDADeserializer(datasets.DatasetAssociationDeserializer, """ model_manager_class = HDAManager - def __init__(self, app: StructuredApp): + def __init__(self, app: MinimalManagerApp): super().__init__(app) self.hda_manager = self.manager diff --git a/lib/galaxy/managers/hdcas.py b/lib/galaxy/managers/hdcas.py index ac25761b6d0..d810fcc9757 100644 --- a/lib/galaxy/managers/hdcas.py +++ b/lib/galaxy/managers/hdcas.py @@ -16,7 +16,7 @@ from galaxy.managers import ( taggable ) from galaxy.managers.collections_util import get_hda_and_element_identifiers -from galaxy.structured_app import StructuredApp +from galaxy.structured_app import MinimalManagerApp from galaxy.util.zipstream import ZipstreamWrapper @@ -87,7 +87,7 @@ class DCESerializer(base.ModelSerializer): Serializer for DatasetCollectionElements. """ - def __init__(self, app: StructuredApp): + def __init__(self, app: MinimalManagerApp): super().__init__(app) self.hda_serializer = hdas.HDASerializer(app) self.dc_serializer = DCSerializer(app, dce_serializer=self) @@ -121,7 +121,7 @@ class DCSerializer(base.ModelSerializer): Serializer for DatasetCollections. """ - def __init__(self, app: StructuredApp, dce_serializer=None): + def __init__(self, app: MinimalManagerApp, dce_serializer=None): super().__init__(app) self.dce_serializer = dce_serializer or DCESerializer(app) @@ -160,7 +160,7 @@ class DCASerializer(base.ModelSerializer): Base (abstract) Serializer class for HDCAs and LDCAs. """ - def __init__(self, app: StructuredApp, dce_serializer=None): + def __init__(self, app: MinimalManagerApp, dce_serializer=None): super().__init__(app) self.dce_serializer = dce_serializer or DCESerializer(app) @@ -215,7 +215,7 @@ class HDCASerializer( Serializer for HistoryDatasetCollectionAssociations. """ - def __init__(self, app: StructuredApp): + def __init__(self, app: MinimalManagerApp): super().__init__(app) self.hdca_manager = HDCAManager(app) diff --git a/lib/galaxy/managers/histories.py b/lib/galaxy/managers/histories.py index d1426d112b1..b3975772f94 100644 --- a/lib/galaxy/managers/histories.py +++ b/lib/galaxy/managers/histories.py @@ -23,7 +23,7 @@ from galaxy.managers import ( sharable ) from galaxy.schema.fields import EncodedDatabaseIdField -from galaxy.structured_app import StructuredApp +from galaxy.structured_app import MinimalManagerApp log = logging.getLogger(__name__) @@ -40,7 +40,7 @@ class HistoryManager(sharable.SharableModelManager, deletable.PurgableManagerMix # TODO: incorporate imp/exp (or alias to) - def __init__(self, app: StructuredApp, hda_manager: hdas.HDAManager, contents_manager: history_contents.HistoryContentsManager, contents_filters: history_contents.HistoryContentsFilters): + def __init__(self, app: MinimalManagerApp, hda_manager: hdas.HDAManager, contents_manager: history_contents.HistoryContentsManager, contents_filters: history_contents.HistoryContentsFilters): super().__init__(app) self.hda_manager = hda_manager self.contents_manager = contents_manager @@ -174,7 +174,7 @@ class HistoryManager(sharable.SharableModelManager, deletable.PurgableManagerMix class HistoryExportView: - def __init__(self, app: StructuredApp): + def __init__(self, app: MinimalManagerApp): self.app = app def get_exports(self, trans, history_id): @@ -226,7 +226,7 @@ class HistorySerializer(sharable.SharableModelSerializer, deletable.PurgableSeri model_manager_class = HistoryManager SINGLE_CHAR_ABBR = 'h' - def __init__(self, app: StructuredApp, hda_manager: hdas.HDAManager, hda_serializer: hdas.HDASerializer, history_contents_serializer: history_contents.HistoryContentsSerializer): + def __init__(self, app: MinimalManagerApp, hda_manager: hdas.HDAManager, hda_serializer: hdas.HDASerializer, history_contents_serializer: history_contents.HistoryContentsSerializer): super().__init__(app) self.history_manager = self.manager @@ -445,7 +445,7 @@ class HistoryDeserializer(sharable.SharableModelDeserializer, deletable.Purgable """ model_manager_class = HistoryManager - def __init__(self, app: StructuredApp): + def __init__(self, app: MinimalManagerApp): super().__init__(app) self.history_manager = self.manager @@ -482,7 +482,7 @@ class HistoriesService: and pydantic models to declare its parameters and return types. """ - def __init__(self, app: StructuredApp, manager: HistoryManager, serializer: HistorySerializer): + def __init__(self, app: MinimalManagerApp, manager: HistoryManager, serializer: HistorySerializer): self.app = app self.manager = manager self.serializer = serializer diff --git a/lib/galaxy/managers/history_contents.py b/lib/galaxy/managers/history_contents.py index 9765a7c0a95..d8d4db39973 100644 --- a/lib/galaxy/managers/history_contents.py +++ b/lib/galaxy/managers/history_contents.py @@ -32,7 +32,7 @@ from galaxy.managers import ( taggable, tools ) -from galaxy.structured_app import StructuredApp +from galaxy.structured_app import MinimalManagerApp log = logging.getLogger(__name__) @@ -73,7 +73,7 @@ class HistoryContentsManager(containers.ContainerManagerMixin): ) default_order_by = 'hid' - def __init__(self, app: StructuredApp): + def __init__(self, app: MinimalManagerApp): self.app = app self.contained_manager = app[self.contained_class_manager_class] self.subcontainer_manager = app[self.subcontainer_class_manager_class] @@ -427,7 +427,7 @@ class HistoryContentsSerializer(base.ModelSerializer, deletable.PurgableSerializ """ model_manager_class = HistoryContentsManager - def __init__(self, app: StructuredApp, **kwargs): + def __init__(self, app: MinimalManagerApp, **kwargs): super().__init__(app, **kwargs) self.default_view = 'summary' diff --git a/lib/galaxy/managers/jobs.py b/lib/galaxy/managers/jobs.py index c96de062caf..db08253505c 100644 --- a/lib/galaxy/managers/jobs.py +++ b/lib/galaxy/managers/jobs.py @@ -9,6 +9,7 @@ from pydantic import ( ) from sqlalchemy import and_, false, func, or_ from sqlalchemy.orm import aliased +from sqlalchemy.orm.scoping import scoped_session from sqlalchemy.sql import select from galaxy import model @@ -21,6 +22,7 @@ from galaxy.managers.collections import DatasetCollectionManager from galaxy.managers.datasets import DatasetManager from galaxy.managers.hdas import HDAManager from galaxy.managers.lddas import LDDAManager +from galaxy.security.idencoding import IdEncodingHelper from galaxy.structured_app import StructuredApp from galaxy.util import ( defaultdict, @@ -96,17 +98,17 @@ class JobSearch: """Search for jobs using tool inputs or other jobs""" def __init__( self, - app: StructuredApp, + sa_session: scoped_session, hda_manager: HDAManager, dataset_collection_manager: DatasetCollectionManager, - ldda_manager: LDDAManager + ldda_manager: LDDAManager, + id_encoding_helper: IdEncodingHelper, ): - self.app = app - self.sa_session = app.model.context + self.sa_session = sa_session self.hda_manager = hda_manager self.dataset_collection_manager = dataset_collection_manager self.ldda_manager = ldda_manager - self.decode_id = self.app.security.decode_id + self.decode_id = id_encoding_helper.decode_id def by_tool_input(self, trans, tool_id, tool_version, param=None, param_dump=None, job_state='ok'): """Search for jobs producing same results using the 'inputs' part of a tool POST.""" diff --git a/lib/galaxy/managers/lddas.py b/lib/galaxy/managers/lddas.py index 53d2a0a9b28..5f088b32078 100644 --- a/lib/galaxy/managers/lddas.py +++ b/lib/galaxy/managers/lddas.py @@ -3,7 +3,7 @@ import logging from galaxy import model, util from galaxy.managers import base as manager_base from galaxy.managers.datasets import DatasetAssociationManager -from galaxy.structured_app import StructuredApp +from galaxy.structured_app import MinimalManagerApp log = logging.getLogger(__name__) @@ -14,7 +14,7 @@ class LDDAManager(DatasetAssociationManager): """ model_class = model.LibraryDatasetDatasetAssociation - def __init__(self, app: StructuredApp): + def __init__(self, app: MinimalManagerApp): """ Set up and initialize other managers needed by lddas. """ diff --git a/lib/galaxy/managers/metrics.py b/lib/galaxy/managers/metrics.py index 66df5c3e3cd..6ff09336451 100644 --- a/lib/galaxy/managers/metrics.py +++ b/lib/galaxy/managers/metrics.py @@ -13,7 +13,7 @@ from pydantic import ( Field, ) -from galaxy.app import StructuredApp +from galaxy.app import MinimalManagerApp log = logging.getLogger(__name__) @@ -72,7 +72,7 @@ TimeSeriesTupleGenerator = Generator[TimeSeriesTuple, None, None] class MetricsManager: """Interface/service object shared by controllers for interacting with metrics.""" - def __init__(self, app: StructuredApp) -> None: + def __init__(self, app: MinimalManagerApp) -> None: self._app = app #: set to true to send additional debugging info to the log self.debugging = True diff --git a/lib/galaxy/managers/pages.py b/lib/galaxy/managers/pages.py index 6d7582e8a75..18ef7921366 100644 --- a/lib/galaxy/managers/pages.py +++ b/lib/galaxy/managers/pages.py @@ -32,7 +32,7 @@ from galaxy.managers.markdown_util import ( ) from galaxy.model.item_attrs import UsesAnnotations from galaxy.schema.fields import EncodedDatabaseIdField -from galaxy.structured_app import StructuredApp +from galaxy.structured_app import MinimalManagerApp from galaxy.util import unicodify from galaxy.util.sanitize_html import sanitize_html @@ -197,7 +197,7 @@ class PagesService: and pydantic models to declare its parameters and return types. """ - def __init__(self, app: StructuredApp): + def __init__(self, app: MinimalManagerApp): self.manager = PageManager(app) self.serializer = PageSerializer(app) self.shareable_service = sharable.ShareableService(self.manager, self.serializer) @@ -304,7 +304,7 @@ class PageManager(sharable.SharableModelManager, UsesAnnotations): annotation_assoc = model.PageAnnotationAssociation rating_assoc = model.PageRatingAssociation - def __init__(self, app: StructuredApp): + def __init__(self, app: MinimalManagerApp): """ """ super().__init__(app) @@ -427,7 +427,7 @@ class PageSerializer(sharable.SharableModelSerializer): model_manager_class = PageManager SINGLE_CHAR_ABBR = 'p' - def __init__(self, app: StructuredApp): + def __init__(self, app: MinimalManagerApp): super().__init__(app) self.page_manager = PageManager(app) @@ -448,7 +448,7 @@ class PageDeserializer(sharable.SharableModelDeserializer): """ model_manager_class = PageManager - def __init__(self, app: StructuredApp): + def __init__(self, app: MinimalManagerApp): super().__init__(app) self.page_manager = self.manager diff --git a/lib/galaxy/managers/remote_files.py b/lib/galaxy/managers/remote_files.py index a7dc50a550e..9d40209fc83 100644 --- a/lib/galaxy/managers/remote_files.py +++ b/lib/galaxy/managers/remote_files.py @@ -12,7 +12,7 @@ from typing import ( from pydantic.tools import parse_obj_as from galaxy import exceptions -from galaxy.app import StructuredApp +from galaxy.app import MinimalManagerApp from galaxy.files import ( ConfiguredFileSources, ProvidesUserFileSourcesUserContext, @@ -37,7 +37,7 @@ class RemoteFilesManager: Interface/service object for interacting with remote files. """ - def __init__(self, app: StructuredApp): + def __init__(self, app: MinimalManagerApp): self._app = app def index( diff --git a/lib/galaxy/managers/sharable.py b/lib/galaxy/managers/sharable.py index 1ec13776546..5782895eae8 100644 --- a/lib/galaxy/managers/sharable.py +++ b/lib/galaxy/managers/sharable.py @@ -34,7 +34,7 @@ from galaxy.managers import ( ) from galaxy.model import UserShareAssociation from galaxy.schema.fields import EncodedDatabaseIdField -from galaxy.structured_app import StructuredApp +from galaxy.structured_app import MinimalManagerApp from galaxy.util import ready_name_for_url log = logging.getLogger(__name__) @@ -51,7 +51,7 @@ class SharableModelManager(base.ModelManager, secured.OwnableManagerMixin, secur #: the single character abbreviation used in username_and_slug: e.g. 'h' for histories: u/user/h/slug SINGLE_CHAR_ABBR: Optional[str] = None - def __init__(self, app: StructuredApp): + def __init__(self, app: MinimalManagerApp): super().__init__(app) # user manager is needed to check access/ownership/admin self.user_manager = users.UserManager(app) diff --git a/lib/galaxy/managers/users.py b/lib/galaxy/managers/users.py index e8b0493bc31..cb9f4c67b73 100644 --- a/lib/galaxy/managers/users.py +++ b/lib/galaxy/managers/users.py @@ -30,7 +30,7 @@ from galaxy.security.validate_user_input import ( validate_password, validate_publicname ) -from galaxy.structured_app import BasicApp, StructuredApp +from galaxy.structured_app import BasicApp, MinimalManagerApp from galaxy.util.hash_util import new_secure_hash from galaxy.web import url_for @@ -611,7 +611,7 @@ class UserManager(base.ModelManager, deletable.PurgableManagerMixin): class UserSerializer(base.ModelSerializer, deletable.PurgableSerializerMixin): model_manager_class = UserManager - def __init__(self, app: StructuredApp): + def __init__(self, app: MinimalManagerApp): """ Convert a User and associated data to a dictionary representation. """ diff --git a/lib/galaxy/managers/visualizations.py b/lib/galaxy/managers/visualizations.py index 5c3c55be16d..84f0d1fa1f9 100644 --- a/lib/galaxy/managers/visualizations.py +++ b/lib/galaxy/managers/visualizations.py @@ -10,7 +10,7 @@ from typing import Optional from galaxy import model from galaxy.managers import sharable from galaxy.schema.fields import EncodedDatabaseIdField -from galaxy.structured_app import StructuredApp +from galaxy.structured_app import MinimalManagerApp log = logging.getLogger(__name__) @@ -43,7 +43,7 @@ class VisualizationSerializer(sharable.SharableModelSerializer): model_manager_class = VisualizationManager SINGLE_CHAR_ABBR = 'v' - def __init__(self, app: StructuredApp): + def __init__(self, app: MinimalManagerApp): super().__init__(app) self.visualization_manager = self.manager @@ -82,7 +82,7 @@ class VisualizationsService: and pydantic models to declare its parameters and return types. """ - def __init__(self, app: StructuredApp, manager: VisualizationManager, serializer: VisualizationSerializer): + def __init__(self, app: MinimalManagerApp, manager: VisualizationManager, serializer: VisualizationSerializer): self.app = app self.manager = manager self.serializer = serializer diff --git a/lib/galaxy/managers/workflows.py b/lib/galaxy/managers/workflows.py index 0877596279d..5ee98861d8e 100644 --- a/lib/galaxy/managers/workflows.py +++ b/lib/galaxy/managers/workflows.py @@ -26,7 +26,7 @@ from galaxy import ( ) from galaxy.jobs.actions.post import ActionBox from galaxy.model.item_attrs import UsesAnnotations -from galaxy.structured_app import StructuredApp +from galaxy.structured_app import MinimalManagerApp from galaxy.tools.parameters import ( params_to_incoming, visit_input_values @@ -69,7 +69,7 @@ class WorkflowsManager: the galaxy.workflow module. """ - def __init__(self, app: StructuredApp): + def __init__(self, app: MinimalManagerApp): self.app = app def get_stored_workflow(self, trans, workflow_id, by_stored_id=True): @@ -306,7 +306,7 @@ CreatedWorkflow = namedtuple("CreatedWorkflow", ["stored_workflow", "workflow", class WorkflowContentsManager(UsesAnnotations): - def __init__(self, app: StructuredApp): + def __init__(self, app: MinimalManagerApp): self.app = app self._resource_mapper_function = get_resource_mapper_function(app) diff --git a/lib/galaxy/model/orm/engine_factory.py b/lib/galaxy/model/orm/engine_factory.py index 8d23f4288f7..5748525b421 100644 --- a/lib/galaxy/model/orm/engine_factory.py +++ b/lib/galaxy/model/orm/engine_factory.py @@ -3,6 +3,7 @@ import logging import os import threading import time +from multiprocessing.util import register_after_fork from sqlalchemy import ( create_engine, @@ -87,4 +88,5 @@ def build_engine(url, engine_options, database_query_profiling_proxy=False, trac # Create the database engine engine = create_engine(url, **engine_options) + register_after_fork(engine, lambda e: e.dispose()) return engine diff --git a/lib/galaxy/structured_app.py b/lib/galaxy/structured_app.py index 7728929a0ee..1cb49fb8d39 100644 --- a/lib/galaxy/structured_app.py +++ b/lib/galaxy/structured_app.py @@ -41,7 +41,43 @@ class BasicApp(Container): datatypes_registry: Registry -class StructuredApp(BasicApp): +class MinimalApp(BasicApp): + is_webapp: bool # is_webapp will be set to true when building WSGI app + new_installation: bool + tag_handler: GalaxyTagHandler + model: GalaxyModelMapping + install_model: ModelMapping + security_agent: GalaxyRBACAgent + host_security_agent: HostAgent + + +class MinimalManagerApp(MinimalApp): + is_webapp: bool # is_webapp will be set to true when building WSGI app + new_installation: bool + tag_handler: GalaxyTagHandler + file_sources: ConfiguredFileSources + genome_builds: GenomeBuilds + model: GalaxyModelMapping + install_model: ModelMapping + security_agent: GalaxyRBACAgent + host_security_agent: HostAgent + dataset_collections_service: Any # 'galaxy.managers.collections.DatasetCollectionManager' + history_manager: Any # 'galaxy.managers.histories.HistoryManager' + hda_manager: Any # 'galaxy.managers.hdas.HDAManager' + workflow_manager: Any # 'galaxy.managers.workflows.WorkflowsManager' + workflow_contents_manager: Any # 'galaxy.managers.workflows.WorkflowContentsManager' + library_folder_manager: Any # 'galaxy.managers.folders.FolderManager' + library_manager: Any # 'galaxy.managers.libraries.LibraryManager' + role_manager: Any # 'galaxy.managers.roles.RoleManager' + installed_repository_manager: Any # 'galaxy.tool_shed.galaxy_install.installed_repository_manager.InstalledRepositoryManager' + user_manager: Any + + @property + def is_job_handler(self) -> bool: + pass + + +class StructuredApp(MinimalManagerApp): """Interface defining typed description of the Galaxy UniverseApplication. Ideally nothing that depends on StructuredApp should require @@ -92,7 +128,3 @@ class StructuredApp(BasicApp): job_manager: Any # galaxy.jobs.manager.JobManager user_manager: Any api_keys_manager: Any - - @property - def is_job_handler(self) -> bool: - pass diff --git a/lib/galaxy/tools/cache.py b/lib/galaxy/tools/cache.py index 94e0b6d2996..351e654ba67 100644 --- a/lib/galaxy/tools/cache.py +++ b/lib/galaxy/tools/cache.py @@ -16,7 +16,7 @@ from sqlalchemy.orm import ( from sqlitedict import SqliteDict from galaxy.model.tool_shed_install import ToolShedRepository -from galaxy.structured_app import StructuredApp +from galaxy.structured_app import MinimalManagerApp from galaxy.tools.toolbox.base import ToolConfRepository from galaxy.util import unicodify from galaxy.util.hash_util import md5_hash_file @@ -286,7 +286,7 @@ class ToolShedRepositoryCache: repositories: List[ToolShedRepository] repos_by_tuple: Dict[Tuple[str, str, str], List[ToolConfRepository]] - def __init__(self, app: StructuredApp): + def __init__(self, app: MinimalManagerApp): self.app = app # Contains ToolConfRepository objects created from shed_tool_conf.xml entries self.local_repositories = [] diff --git a/lib/galaxy/tools/data_manager/manager.py b/lib/galaxy/tools/data_manager/manager.py index 4805f1320ab..dd823b2e925 100644 --- a/lib/galaxy/tools/data_manager/manager.py +++ b/lib/galaxy/tools/data_manager/manager.py @@ -5,7 +5,7 @@ import os from typing import Dict from galaxy import util -from galaxy.structured_app import StructuredApp +from galaxy.structured_app import MinimalManagerApp from galaxy.tools.data import TabularToolDataTable from galaxy.util.template import fill_template @@ -20,7 +20,7 @@ class DataManagers: data_managers: Dict[str, 'DataManager'] managed_data_tables: Dict[str, 'DataManager'] - def __init__(self, app: StructuredApp, xml_filename=None): + def __init__(self, app: MinimalManagerApp, xml_filename=None): self.app = app self.data_managers = {} self.managed_data_tables = {} diff --git a/lib/galaxy/visualization/genomes.py b/lib/galaxy/visualization/genomes.py index 1737b9e1802..3869a5178ec 100644 --- a/lib/galaxy/visualization/genomes.py +++ b/lib/galaxy/visualization/genomes.py @@ -11,7 +11,7 @@ from galaxy.exceptions import ( ObjectNotFound, ReferenceDataError, ) -from galaxy.structured_app import StructuredApp +from galaxy.structured_app import MinimalManagerApp from galaxy.util.bunch import Bunch log = logging.getLogger(__name__) @@ -197,7 +197,7 @@ class Genomes: Provides information about available genome data and methods for manipulating that data. """ - def __init__(self, app: StructuredApp): + def __init__(self, app: MinimalManagerApp): self.app = app # Create list of genomes from app.genome_builds self.genomes: Dict[str, Genome] = {} diff --git a/lib/galaxy/webapps/galaxy/config_schema.yml b/lib/galaxy/webapps/galaxy/config_schema.yml index a6fd5931f03..378e1b9bb4b 100644 --- a/lib/galaxy/webapps/galaxy/config_schema.yml +++ b/lib/galaxy/webapps/galaxy/config_schema.yml @@ -3183,6 +3183,15 @@ mapping: will automatically create and use a separate sqlite database located in your /database folder (indicated in the commented out line below). + enable_celery_tasks: + type: bool + default: false + required: false + desc: | + Offload long-running tasks to a Celery task queue. + Activate this only if you have setup a Celery worker for Galaxy. + For details, see https://docs.galaxyproject.org/en/master/admin/production.html + use_pbkdf2: type: bool default: true diff --git a/lib/galaxy/webapps/galaxy/controllers/user.py b/lib/galaxy/webapps/galaxy/controllers/user.py index e1749d5d790..ffba90b02ab 100644 --- a/lib/galaxy/webapps/galaxy/controllers/user.py +++ b/lib/galaxy/webapps/galaxy/controllers/user.py @@ -222,11 +222,15 @@ class User(BaseUIController, UsesFormDefinitionsMixin, CreatesApiKeysMixin): if message: return self.message_exception(trans, message) if trans.user: - # Queue a quota recalculation (async) task -- this takes a - # while sometimes, so we don't want to block on logout. - send_local_control_task(trans.app, - "recalculate_user_disk_usage", - kwargs={"user_id": trans.security.encode_id(trans.user.id)}) + if trans.app.config.enable_celery_tasks: + # Queue a quota recalculation (async) task -- this takes a + # while sometimes, so we don't want to block on logout. + from galaxy.celery.tasks import recalculate_user_disk_usage + recalculate_user_disk_usage.delay(user_id=trans.user.id) + else: + send_local_control_task(trans.app, + "recalculate_user_disk_usage", + kwargs={"user_id": trans.security.encode_id(trans.user.id)}) # Since logging an event requires a session, we'll log prior to ending the session trans.log_event("User logged out") trans.handle_user_logout(logout_all=logout_all) diff --git a/packages/app/galaxy/celery b/packages/app/galaxy/celery new file mode 120000 index 00000000000..6c7157a6e38 --- /dev/null +++ b/packages/app/galaxy/celery @@ -0,0 +1 @@ +../../../lib/galaxy/celery \ No newline at end of file diff --git a/packages/app/requirements.txt b/packages/app/requirements.txt index 4be8956dc10..c7ff5fbe1b6 100644 --- a/packages/app/requirements.txt +++ b/packages/app/requirements.txt @@ -5,6 +5,7 @@ galaxy-tool-util galaxy-web-framework galaxy-web-stack +celery kombu Beaker pykwalify diff --git a/pyproject.toml b/pyproject.toml index 96bfbba2693..6cb68917621 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -25,10 +25,12 @@ bleach = "*" boltons = "*" boto = "*" bx-python = "*" +celery = "*" Cheetah3 = "*" cloudauthz = "==0.6.0" cloudbridge = "*" contextvars = {version = "*", python = "~3.6"} +circus = "*" cwltool = "==3.0.20201109103151" dictobj = "*" docutils = "*" diff --git a/test/unit/managers/test_HDAManager.py b/test/unit/managers/test_HDAManager.py index 8f44aaeb47b..bdf491c1de1 100644 --- a/test/unit/managers/test_HDAManager.py +++ b/test/unit/managers/test_HDAManager.py @@ -175,7 +175,7 @@ class HDAManagerTestCase(HDATestCase): self.log("should purge an hda if config does allow") self.assertFalse(item1.purged) - self.assertEqual(self.hda_manager.purge(item1), item1) + self.hda_manager.purge(item1) self.assertTrue(item1.deleted) self.assertTrue(item1.purged) diff --git a/test/unit/unittest_utils/galaxy_mock.py b/test/unit/unittest_utils/galaxy_mock.py index 0d660173437..02e44efe61a 100644 --- a/test/unit/unittest_utils/galaxy_mock.py +++ b/test/unit/unittest_utils/galaxy_mock.py @@ -21,7 +21,7 @@ from galaxy.model import mapping, tags from galaxy.model.base import SharedModelMapping from galaxy.model.mapping import GalaxyModelMapping from galaxy.security import idencoding -from galaxy.structured_app import BasicApp, StructuredApp +from galaxy.structured_app import BasicApp, MinimalManagerApp, StructuredApp from galaxy.tool_util.deps.containers import NullContainerFinder from galaxy.util import StructuredExecutionTimer from galaxy.util.bunch import Bunch @@ -67,6 +67,7 @@ class MockApp(di.Container): def __init__(self, config=None, **kwargs): super().__init__() self[BasicApp] = self + self[MinimalManagerApp] = self self[StructuredApp] = self self.config = config or MockAppConfig(**kwargs) self.security = self.config.security @@ -139,6 +140,7 @@ class MockAppConfig(Bunch): self.security = idencoding.IdEncodingHelper(id_secret='6e46ed6483a833c100e68cc3f1d0dd76') self.database_connection = kwargs.get('database_connection', "sqlite:///:memory:") self.use_remote_user = kwargs.get('use_remote_user', False) + self.enable_celery_tasks = False self.data_dir = os.path.join(root, 'database') self.file_path = os.path.join(self.data_dir, 'files') self.jobs_directory = os.path.join(self.data_dir, 'jobs_directory')