Move task creation into job creation

This commit is contained in:
mvdbeek
2022-05-26 12:48:43 +02:00
parent 38596d461c
commit b447b9659a
6 changed files with 162 additions and 117 deletions
+7 -7
View File
@@ -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))
+44 -20
View File
@@ -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,
+57 -74
View File
@@ -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.
+21 -8
View File
@@ -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
+26 -4
View File
@@ -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()
+7 -4
View File
@@ -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
)