diff --git a/lib/galaxy/app.py b/lib/galaxy/app.py index 21ee7f8144f..92bb20dc7e0 100644 --- a/lib/galaxy/app.py +++ b/lib/galaxy/app.py @@ -541,6 +541,13 @@ class GalaxyManagerApplication(MinimalManagerApp, MinimalGalaxyApplication): ) self.vault = self._register_singleton(Vault, VaultFactory.from_app(self)) + # Load security policy. + self.security_agent = self.model.security_agent + self.host_security_agent = galaxy.model.security.HostAgent( + model=self.security_agent.model, permitted_actions=self.security_agent.permitted_actions + ) + # Load quota management. + self.quota_agent = self._register_singleton(QuotaAgent, get_quota_agent(self.config, self.model)) # 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 @@ -649,13 +656,6 @@ class UniverseApplication(StructuredApp, GalaxyManagerApplication): self[ToursRegistry] = tour_registry # type: ignore[misc] # Webhooks registry self.webhooks_registry = self._register_singleton(WebhooksRegistry, WebhooksRegistry(self.config.webhooks_dir)) - # Load security policy. - self.security_agent = self.model.security_agent - self.host_security_agent = galaxy.model.security.HostAgent( - model=self.security_agent.model, permitted_actions=self.security_agent.permitted_actions - ) - # Load quota management. - self.quota_agent = self._register_singleton(QuotaAgent, get_quota_agent(self.config, self.model)) # Heartbeat for thread profiling self.heartbeat = None self.auth_manager = self._register_singleton(auth.AuthManager, auth.AuthManager(self.config)) diff --git a/lib/galaxy/celery/tasks.py b/lib/galaxy/celery/tasks.py index 1a397bf2425..1fdc362e165 100644 --- a/lib/galaxy/celery/tasks.py +++ b/lib/galaxy/celery/tasks.py @@ -1,16 +1,17 @@ -import os +import json +from functools import lru_cache from pathlib import Path from galaxy import model from galaxy.celery import galaxy_task from galaxy.config import GalaxyAppConfiguration from galaxy.datatypes.registry import Registry as DatatypesRegistry +from galaxy.jobs import MinimalJobWrapper from galaxy.managers.collections import DatasetCollectionManager from galaxy.managers.hdas import HDAManager from galaxy.managers.lddas import LDDAManager from galaxy.managers.markdown_util import generate_branded_pdf from galaxy.managers.model_stores import ModelStoreManager -from galaxy.metadata import PortableDirectoryMetadataGenerator from galaxy.metadata.set_metadata import set_metadata_portable from galaxy.model.scoped_session import galaxy_scoped_session from galaxy.objectstore import BaseObjectStore @@ -28,6 +29,7 @@ from galaxy.schema.tasks import ( WriteInvocationTo, ) from galaxy.structured_app import MinimalManagerApp +from galaxy.tools import create_tool_from_representation from galaxy.tools.data_fetch import do_fetch from galaxy.util.custom_logging import get_logger from galaxy.web.short_term_storage import ShortTermStorageMonitor @@ -35,6 +37,13 @@ from galaxy.web.short_term_storage import ShortTermStorageMonitor log = get_logger(__name__) +@lru_cache() +def cached_create_tool_from_representation(app, raw_tool_source): + return create_tool_from_representation( + app=app, raw_tool_source=raw_tool_source, tool_dir="", tool_source_class="XmlToolSource" + ) + + @galaxy_task(ignore_result=True, action="recalculate a user's disk usage") def recalculate_user_disk_usage(session: galaxy_scoped_session, user_id=None): if user_id: @@ -70,7 +79,6 @@ def set_job_metadata( datatypes_registry: DatatypesRegistry, object_store: BaseObjectStore, ): - log.error("set_job_metatdata, %s", tool_job_working_directory) set_metadata_portable( tool_job_working_directory, datatypes_registry=datatypes_registry, @@ -79,16 +87,6 @@ def set_job_metadata( ) -@galaxy_task(action="Create metadata setup") -def create_metadata_setup(tool_job_working_directory, job_id, sa_session: galaxy_scoped_session): - log.error("workign dir in create_metadata_setup is %s", tool_job_working_directory) - job_metadata = os.path.join(tool_job_working_directory, "galaxy.json") - PortableDirectoryMetadataGenerator(job_id=job_id).setup_external_metadata( - {}, {}, sa_session=sa_session, tmp_dir=tool_job_working_directory, job_metadata=job_metadata - ) - return tool_job_working_directory - - @galaxy_task(action="set dataset association metadata") def set_metadata( hda_manager: HDAManager, ldda_manager: LDDAManager, dataset_id, model_class="HistoryDatasetAssociation" @@ -101,19 +99,45 @@ def set_metadata( @galaxy_task -def setup_fetch_data(job_id: int, app: MinimalManagerApp, sa_session: galaxy_scoped_session): - sa_session.query(model.Job).get(job_id) - # TODO: split up JobWrapper some more so we can share the setup_metdata logic and finish logic. - # JobWrapper(job, queue) +def setup_fetch_data(job_id: int, raw_tool_source: str, app: MinimalManagerApp, sa_session: galaxy_scoped_session): + tool = cached_create_tool_from_representation(app=app, raw_tool_source=raw_tool_source) + job = sa_session.query(model.Job).get(job_id) + # TODO: assert state + mini_job_wrapper = MinimalJobWrapper(job=job, app=app, tool=tool) + mini_job_wrapper.change_state(model.Job.states.QUEUED, flush=False, job=job) + # Set object store after job destination so can leverage parameters... + mini_job_wrapper._set_object_store_ids(job) + request_json = Path(mini_job_wrapper.working_directory) / "request.json" + request_json_value = next(iter(p.value for p in job.parameters if p.name == "request_json")) + request_json.write_text(json.loads(request_json_value)) + mini_job_wrapper.setup_external_metadata( + output_fnames=mini_job_wrapper.job_io.get_output_fnames(), + set_extension=True, + tmp_dir=mini_job_wrapper.working_directory, + # We don't want to overwrite metadata that was copied over in init_meta(), as per established behavior + kwds={"overwrite": False}, + ) + mini_job_wrapper.prepare() + # Technically this should be changed in fetch_data + mini_job_wrapper.change_state(model.Job.states.RUNNING, flush=True, job=job) + return mini_job_wrapper.working_directory, str(request_json), mini_job_wrapper.job_io.file_sources_dict + + +@galaxy_task +def finish_job(job_id: int, raw_tool_source: str, app: MinimalManagerApp, sa_session: galaxy_scoped_session): + tool = cached_create_tool_from_representation(app=app, raw_tool_source=raw_tool_source) + job = sa_session.query(model.Job).get(job_id) + # TODO: assert state ? + mini_job_wrapper = MinimalJobWrapper(job=job, app=app, tool=tool) + mini_job_wrapper.finish("", "") @galaxy_task(action="Run fetch_data") def fetch_data( - tool_job_working_directory, - request_path, - file_sources_dict, + setup_return, datatypes_registry: DatatypesRegistry, ): + tool_job_working_directory, request_path, file_sources_dict = setup_return working_directory = Path(tool_job_working_directory) / "working" do_fetch( request_path=request_path, diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index 7a7a812280d..250e266643a 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -957,7 +957,7 @@ class HasResourceParameters: return resource_params -class JobWrapper(HasResourceParameters): +class MinimalJobWrapper(HasResourceParameters): """ Wraps a 'model.Job' with convenience methods for running processes and state management. @@ -965,13 +965,12 @@ class JobWrapper(HasResourceParameters): is_task = False - def __init__(self, job, queue: "JobHandlerQueue", use_persisted_destination=False, app=None): + def __init__(self, job: model.Job, app: MinimalManagerApp, use_persisted_destination: bool = False, tool=None): self.job_id = job.id self.session_id = job.session_id self.user_id = job.user_id - self.tool = queue.app.toolbox.get_tool(job.tool_id, job.tool_version, exact=True) - self.queue = queue - self.app: MinimalManagerApp = queue.app + self.app: MinimalManagerApp = app + self.tool = tool self.sa_session = self.app.model.context self.extra_filenames: List[str] = [] self.environment_variables: List[Dict[str, str]] = [] @@ -991,12 +990,10 @@ class JobWrapper(HasResourceParameters): # resolved self._job_io = None self.tool_provided_job_metadata = None - self.job_runner_mapper = JobRunnerMapper(self, queue.dispatcher.url_to_destination, self.app.job_config) self.params = None if job.params: self.params = loads(job.params) - if use_persisted_destination: - self.job_runner_mapper.cached_job_destination = JobDestination(from_job=job) + # Wrapper holding the info required to restore and clean up from files used for setting metadata externally self.__external_output_metadata = None self.__has_tasks = bool(job.tasks) @@ -1156,18 +1153,9 @@ class JobWrapper(HasResourceParameters): get_job_runner = get_job_runner_url @property - def job_destination(self): - """Return the JobDestination that this job will use to run. This will - either be a configured destination, a randomly selected destination if - the configured destination was a tag, or a dynamically generated - destination from the dynamic runner. - - Calling this method for the first time causes the dynamic runner to do - its calculation, if any. - - :returns: ``JobDestination`` - """ - return self.job_runner_mapper.get_job_destination(self.params) + def job_destination(self) -> JobDestination: + """Subclassses can return a configured job destination.""" + return JobDestination() @property def galaxy_url(self): @@ -1248,8 +1236,9 @@ class JobWrapper(HasResourceParameters): self.environment_variables, ) = tool_evaluator.build() job.command_line = self.command_line - self.interactivetools = tool_evaluator.populate_interactivetools() - self.app.interactivetool_manager.create_interactivetool(job, self.tool, self.interactivetools) + if hasattr(self.app, "interactivetool_manager"): + self.interactivetools = tool_evaluator.populate_interactivetools() + self.app.interactivetool_manager.create_interactivetool(job, self.tool, self.interactivetools) # Ensure galaxy_lib_dir is set in case there are any later chdirs self.galaxy_lib_dir @@ -1548,24 +1537,6 @@ class JobWrapper(HasResourceParameters): log.warning("set_runner() is deprecated, use set_job_destination()") self.set_job_destination(self.job_destination, external_id) - def set_job_destination(self, job_destination, external_id=None, flush=True, job=None): - """ - Persist job destination params in the database for recovery. - - self.job_destination is not used because a runner may choose to rewrite - parts of the destination (e.g. the params). - """ - if job is None: - job = self.get_job() - log.debug(f"({job.id}) Persisting job destination (destination id: {job_destination.id})") - job.destination_id = job_destination.id - job.destination_params = job_destination.params - job.job_runner_name = job_destination.runner - job.job_runner_external_id = external_id - self.sa_session.add(job) - if flush: - self.sa_session.flush() - def set_external_id(self, external_id, job=None, flush=True): if job is None: job = self.get_job() @@ -1598,43 +1569,13 @@ class JobWrapper(HasResourceParameters): self.set_job_destination(self.job_destination, None, flush=False, job=job) # Set object store after job destination so can leverage parameters... self._set_object_store_ids(job) - if self.app.config.enable_celery_tasks and job.tool_id == "__DATA_FETCH__": - # Move into task - import pathlib - - request_json = pathlib.Path(self.working_directory) / "request.json" - request_json_value = next(iter(p.value for p in job.parameters if p.name == "request_json")) - request_json.write_text(json.loads(request_json_value)) - from galaxy.celery.tasks import ( - fetch_data, - set_job_metadata, - ) - - self.change_state(model.Job.states.RUNNING, flush=False, job=job) - self.sa_session.flush() - self.setup_external_metadata( - output_fnames=self.job_io.get_output_fnames(), - set_extension=True, - tmp_dir=self.working_directory, - # We don't want to overwrite metadata that was copied over in init_meta(), as per established behavior - kwds={"overwrite": False}, - ) - self.prepare() - try: - ( - fetch_data.s(self.working_directory, str(request_json), self.job_io.file_sources_dict) - | set_job_metadata.s(extended_metadata_collection="extended" in self.metadata_strategy) - )().get() - except Exception as e: - # TODO ... error handler on task / chain level ? - self.fail(message=str(e)) - # Reset the Null Metadata. That would normally happen because the finish context rebuilds another JobWrapper instance. - self.tool_provided_job_metadata = None - self.finish(tool_stdout="", tool_stderr="") - return False self.sa_session.flush() return True + def set_job_destination(self, job_destination, external_id=None, flush=True, job=None): + """Subclasses should implement this to persist a destination, if necessary.""" + pass + def _set_object_store_ids(self, job): if job.object_store_id: # We aren't setting this during job creation anymore, but some existing @@ -2454,6 +2395,48 @@ class JobWrapper(HasResourceParameters): self.sa_session.flush() +class JobWrapper(MinimalJobWrapper): + def __init__(self, job, queue: "JobHandlerQueue", use_persisted_destination=False, app=None): + super().__init__(job, app=queue.app, use_persisted_destination=use_persisted_destination) + self.queue = queue + self.tool = self.app.toolbox.get_tool(job.tool_id, job.tool_version, exact=True) + self.job_runner_mapper = JobRunnerMapper(self, queue.dispatcher.url_to_destination, self.app.job_config) + if use_persisted_destination: + self.job_runner_mapper.cached_job_destination = JobDestination(from_job=job) + + @property + def job_destination(self): + """Return the JobDestination that this job will use to run. This will + either be a configured destination, a randomly selected destination if + the configured destination was a tag, or a dynamically generated + destination from the dynamic runner. + + Calling this method for the first time causes the dynamic runner to do + its calculation, if any. + + :returns: ``JobDestination`` + """ + return self.job_runner_mapper.get_job_destination(self.params) + + def set_job_destination(self, job_destination, external_id=None, flush=True, job=None): + """ + Persist job destination params in the database for recovery. + + self.job_destination is not used because a runner may choose to rewrite + parts of the destination (e.g. the params). + """ + if job is None: + job = self.get_job() + log.debug(f"({job.id}) Persisting job destination (destination id: {job_destination.id})") + job.destination_id = job_destination.id + job.destination_params = job_destination.params + job.job_runner_name = job_destination.runner + job.job_runner_external_id = external_id + self.sa_session.add(job) + if flush: + self.sa_session.flush() + + class TaskWrapper(JobWrapper): """ Extension of JobWrapper intended for running tasks. diff --git a/lib/galaxy/tools/__init__.py b/lib/galaxy/tools/__init__.py index 36e2de0885d..ecf0319216d 100644 --- a/lib/galaxy/tools/__init__.py +++ b/lib/galaxy/tools/__init__.py @@ -288,6 +288,13 @@ def create_tool_from_source(app, tool_source, config_file=None, **kwds): return tool +def create_tool_from_representation( + app, raw_tool_source: str, tool_dir: str, tool_source_class="XmlToolSource" +) -> "Tool": + tool_source = get_tool_source(tool_source_class=tool_source_class, raw_tool_source=raw_tool_source) + return create_tool_from_source(app, tool_source=tool_source, tool_dir=tool_dir) + + class ToolBox(BaseGalaxyToolBox): """A derivative of AbstractToolBox with knowledge about Tool internals - how to construct them, action types, dependency management, etc.... @@ -1046,14 +1053,20 @@ class Tool(Dictifiable): self.required_files = required_files self.citations = self._parse_citations(tool_source) - ontology_data = expand_ontology_data( - tool_source, - self.all_ids, - self.app.biotools_metadata_source, - ) - self.xrefs = ontology_data.xrefs - self.edam_operations = ontology_data.edam_operations - self.edam_topics = ontology_data.edam_topics + biotools_metadata_source = getattr(self.app, "biotools_metadata_source", None) + if biotools_metadata_source: + ontology_data = expand_ontology_data( + tool_source, + self.all_ids, + self.app.biotools_metadata_source, + ) + self.xrefs = ontology_data.xrefs + self.edam_operations = ontology_data.edam_operations + self.edam_topics = ontology_data.edam_topics + else: + self.xrefs = [] + self.edam_operations = None + self.edam_topics = None self.__parse_trackster_conf(tool_source) # Record macro paths so we can reload a tool if any of its macro has changes diff --git a/lib/galaxy/tools/execute.py b/lib/galaxy/tools/execute.py index ae2305f6dce..5326f2e5c35 100644 --- a/lib/galaxy/tools/execute.py +++ b/lib/galaxy/tools/execute.py @@ -5,6 +5,7 @@ collections from matched collections. """ import collections import logging +import typing from abc import abstractmethod from typing import ( Dict, @@ -26,6 +27,9 @@ from galaxy.tools.actions import ( ) from galaxy.tools.parameters.basic import is_runtime_value +if typing.TYPE_CHECKING: + from galaxy.tools import Tool + log = logging.getLogger(__name__) SINGLE_EXECUTION_SUCCESS_MESSAGE = "Tool ${tool_id} created job ${job_id}" @@ -42,9 +46,9 @@ MappingParameters = collections.namedtuple("MappingParameters", ["param_template def execute( trans, - tool, + tool: "Tool", mapping_params, - history, + history: model.History, rerun_remap_job_id=None, collection_info=None, workflow_invocation_uuid=None, @@ -121,10 +125,11 @@ def execute( execution_tracker.record_error(result) tool_action = tool.tool_action - if hasattr(tool_action, "check_inputs_ready"): + check_inputs_ready = getattr(tool_action, "check_inputs_ready", None) + if check_inputs_ready: for params in execution_tracker.param_combinations: # This will throw an exception if the tool is not ready. - tool_action.check_inputs_ready( + check_inputs_ready( tool, trans, params, @@ -163,6 +168,23 @@ def execute( tool_id = tool.id for job2 in execution_tracker.successful_jobs: # Put the job in the queue if tracking in memory + if tool_id == "__DATA_FETCH__" and tool.app.config.enable_celery_tasks: + job_id = job2.id + from galaxy.celery.tasks import ( + fetch_data, + finish_job, + set_job_metadata, + setup_fetch_data, + ) + + raw_tool_source = tool.tool_source.to_string() + ( + setup_fetch_data.s(job_id, raw_tool_source=raw_tool_source) + | fetch_data.s() + | set_job_metadata.s(extended_metadata_collection="extended" in tool.app.config.metadata_strategy) + | finish_job.si(job_id=job_id, raw_tool_source=raw_tool_source) + )() + continue tool.app.job_manager.enqueue(job2, tool=tool, flush=False) trans.log_event(f"Added job to the job queue, id: {str(job2.id)}", tool_id=tool_id) trans.sa_session.flush() diff --git a/lib/galaxy/tools/remote_tool_eval.py b/lib/galaxy/tools/remote_tool_eval.py index a9d6513954b..0616100d8ef 100644 --- a/lib/galaxy/tools/remote_tool_eval.py +++ b/lib/galaxy/tools/remote_tool_eval.py @@ -21,9 +21,8 @@ from galaxy.model import store from galaxy.model.store import SessionlessContext from galaxy.objectstore import ObjectStore from galaxy.structured_app import MinimalToolApp -from galaxy.tool_util.parser.factory import get_tool_source from galaxy.tools import ( - create_tool_from_source, + create_tool_from_representation, evaluation, ) from galaxy.tools.data import ToolDataTableManager @@ -98,8 +97,12 @@ def main(TMPDIR, WORKING_DIRECTORY, IMPORT_STORE_DIRECTORY): file_sources=job_io.file_sources, ) # TODO: could try to serialize just a minimal tool variant instead of the whole thing ? - tool_source = get_tool_source(tool_source_class=job_io.tool_source_class, raw_tool_source=job_io.tool_source) - tool = create_tool_from_source(app, tool_source=tool_source, tool_dir=job_io.tool_dir) + tool = create_tool_from_representation( + app=app, + raw_tool_source=job_io.tool_source, + tool_dir=job_io.tool_dir, + tool_source_class=job_io.tool_source_class, + ) tool_evaluator = evaluation.RemoteToolEvaluator( app=app, tool=tool, job=job_io.job, local_working_directory=WORKING_DIRECTORY )