diff --git a/doc/source/lib/galaxy.jobs.actions.rst b/doc/source/lib/galaxy.jobs.actions.rst index 849f193f77e..22358a11128 100644 --- a/doc/source/lib/galaxy.jobs.actions.rst +++ b/doc/source/lib/galaxy.jobs.actions.rst @@ -1,7 +1,7 @@ -galaxy.jobs.actions package +galaxy.job_execution.actions package =========================== -.. automodule:: galaxy.jobs.actions +.. automodule:: galaxy.job_execution.actions :members: :undoc-members: :show-inheritance: @@ -9,10 +9,10 @@ galaxy.jobs.actions package Submodules ---------- -galaxy.jobs.actions.post module +galaxy.job_execution.actions.post module ------------------------------- -.. automodule:: galaxy.jobs.actions.post +.. automodule:: galaxy.job_execution.actions.post :members: :undoc-members: :show-inheritance: diff --git a/doc/source/lib/galaxy.jobs.rst b/doc/source/lib/galaxy.jobs.rst index bb09c301003..aa64cbdbe49 100644 --- a/doc/source/lib/galaxy.jobs.rst +++ b/doc/source/lib/galaxy.jobs.rst @@ -12,7 +12,7 @@ Subpackages .. toctree:: :maxdepth: 4 - galaxy.jobs.actions + galaxy.job_execution.actions galaxy.jobs.rules galaxy.jobs.runners galaxy.jobs.splitters diff --git a/lib/galaxy/jobs/actions/__init__.py b/lib/galaxy/job_execution/actions/__init__.py similarity index 100% rename from lib/galaxy/jobs/actions/__init__.py rename to lib/galaxy/job_execution/actions/__init__.py diff --git a/lib/galaxy/jobs/actions/post.py b/lib/galaxy/job_execution/actions/post.py similarity index 99% rename from lib/galaxy/jobs/actions/post.py rename to lib/galaxy/job_execution/actions/post.py index 1f9deb2620a..eb34662d7f6 100644 --- a/lib/galaxy/jobs/actions/post.py +++ b/lib/galaxy/job_execution/actions/post.py @@ -27,6 +27,10 @@ class DefaultJobAction: def execute(cls, app, sa_session, action, job, replacement_dict=None, final_job_state=None): pass + @classmethod + def execute_on_mapped_over(cls, trans, sa_session, action, step_inputs, step_outputs, replacement_dict, final_job_state=None): + pass + @classmethod def get_short_str(cls, pja): if pja.action_arguments: diff --git a/lib/galaxy/job_execution/compute_environment.py b/lib/galaxy/job_execution/compute_environment.py new file mode 100644 index 00000000000..cd89a14c7db --- /dev/null +++ b/lib/galaxy/job_execution/compute_environment.py @@ -0,0 +1,156 @@ +import os +from abc import ( + ABCMeta, + abstractmethod, +) + +from galaxy.job_execution.setup import JobIO +from galaxy.model import Job + + +class ComputeEnvironment(metaclass=ABCMeta): + """ Definition of the job as it will be run on the (potentially) remote + compute server. + """ + + @abstractmethod + def output_names(self): + """ Output unqualified filenames defined by job. """ + + @abstractmethod + def input_path_rewrite(self, dataset): + """Input path for specified dataset.""" + + @abstractmethod + def output_path_rewrite(self, dataset): + """Output path for specified dataset.""" + + @abstractmethod + def input_extra_files_rewrite(self, dataset): + """Input extra files path rewrite for specified dataset.""" + + @abstractmethod + def output_extra_files_rewrite(self, dataset): + """Output extra files path rewrite for specified dataset.""" + + @abstractmethod + def input_metadata_rewrite(self, dataset, metadata_value): + """Input metadata path rewrite for specified dataset.""" + + @abstractmethod + def unstructured_path_rewrite(self, path): + """Rewrite loc file paths, etc..""" + + @abstractmethod + def working_directory(self): + """ Job working directory (potentially remote) """ + + @abstractmethod + def config_directory(self): + """ Directory containing config files (potentially remote) """ + + @abstractmethod + def env_config_directory(self): + """Working directory (possibly as environment variable evaluation).""" + + @abstractmethod + def sep(self): + """ os.path.sep for the platform this job will execute in. + """ + + @abstractmethod + def new_file_path(self): + """ Absolute path to dump new files for this job on compute server. """ + + @abstractmethod + def tool_directory(self): + """ Absolute path to tool files for this job on compute server. """ + + @abstractmethod + def version_path(self): + """ Location of the version file for the underlying tool. """ + + @abstractmethod + def home_directory(self): + """Home directory of target job - none if HOME should not be set.""" + + @abstractmethod + def tmp_directory(self): + """Temp directory of target job - none if HOME should not be set.""" + + @abstractmethod + def galaxy_url(self): + """URL to access Galaxy API from for this compute environment.""" + + +class SimpleComputeEnvironment: + + def config_directory(self): + return os.path.join(self.working_directory(), "configs") # type: ignore[attr-defined] + + def sep(self): + return os.path.sep + + +class SharedComputeEnvironment(SimpleComputeEnvironment, ComputeEnvironment): + """ Default ComputeEnvironment for job and task wrapper to pass + to ToolEvaluator - valid when Galaxy and compute share all the relevant + file systems. + """ + + def __init__(self, job_io: JobIO, job: Job): + self.job_io = job_io + self.job = job + + def output_names(self): + return self.job_io.get_output_basenames() + + def output_paths(self): + return self.job_io.get_output_fnames() + + def input_path_rewrite(self, dataset): + return self.job_io.get_input_path(dataset).false_path + + def output_path_rewrite(self, dataset): + dataset_path = self.job_io.get_output_path(dataset) + if hasattr(dataset_path, "false_path"): + return dataset_path.false_path + else: + return dataset_path + + def input_extra_files_rewrite(self, dataset): + return None + + def output_extra_files_rewrite(self, dataset): + return None + + def input_metadata_rewrite(self, dataset, metadata_value): + return None + + def unstructured_path_rewrite(self, path): + return None + + def working_directory(self): + return self.job_io.working_directory + + def env_config_directory(self): + """Working directory (possibly as environment variable evaluation).""" + return "$_GALAXY_JOB_DIR" + + def new_file_path(self): + return self.job_io.new_file_path + + def version_path(self): + return self.job_io.version_path + + def tool_directory(self): + return self.job_io.tool_directory + + def home_directory(self): + return self.job_io.home_directory + + def tmp_directory(self): + return self.job_io.tmp_directory + + def galaxy_url(self): + return self.job_io.galaxy_url diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index d90e8640403..c53f1deda27 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -12,10 +12,6 @@ import shutil import sys import time import traceback -from abc import ( - ABCMeta, - abstractmethod, -) from json import loads from typing import Any, Dict, List, TYPE_CHECKING @@ -33,6 +29,8 @@ from galaxy.exceptions import ( ObjectInvalid, ObjectNotFound, ) +from galaxy.job_execution.actions.post import ActionBox +from galaxy.job_execution.compute_environment import SharedComputeEnvironment from galaxy.job_execution.output_collect import ( collect_extra_files, collect_shrinked_content_from_path, @@ -45,14 +43,13 @@ from galaxy.job_execution.setup import ( # noqa: F401 TOOL_PROVIDED_JOB_METADATA_FILE, TOOL_PROVIDED_JOB_METADATA_KEYS, ) -from galaxy.jobs.actions.post import ActionBox from galaxy.jobs.mapper import ( JobMappingException, JobRunnerMapper, ) from galaxy.jobs.runners import BaseJobRunner, JobState from galaxy.metadata import get_metadata_compute_strategy -from galaxy.model import Job, store +from galaxy.model import store from galaxy.objectstore import ObjectStorePopulator from galaxy.structured_app import MinimalManagerApp from galaxy.tool_util.deps import requirements @@ -60,6 +57,7 @@ from galaxy.tool_util.output_checker import ( check_output, DETECTED_JOB_STATE, ) +from galaxy.tools.evaluation import ToolEvaluator from galaxy.util import ( parse_xml_string, RWXRWXRWX, @@ -1294,11 +1292,6 @@ class JobWrapper(HasResourceParameters): return job def _get_tool_evaluator(self, job): - # Hacky way to avoid circular import for now. - # Placing ToolEvaluator in either jobs or tools - # result in circular dependency. - from galaxy.tools.evaluation import ToolEvaluator - tool_evaluator = ToolEvaluator( app=self.app, job=job, @@ -2430,154 +2423,6 @@ class TaskWrapper(JobWrapper): return working_directory -class ComputeEnvironment(metaclass=ABCMeta): - """ Definition of the job as it will be run on the (potentially) remote - compute server. - """ - - @abstractmethod - def output_names(self): - """ Output unqualified filenames defined by job. """ - - @abstractmethod - def input_path_rewrite(self, dataset): - """Input path for specified dataset.""" - - @abstractmethod - def output_path_rewrite(self, dataset): - """Output path for specified dataset.""" - - @abstractmethod - def input_extra_files_rewrite(self, dataset): - """Input extra files path rewrite for specified dataset.""" - - @abstractmethod - def output_extra_files_rewrite(self, dataset): - """Output extra files path rewrite for specified dataset.""" - - @abstractmethod - def input_metadata_rewrite(self, dataset, metadata_value): - """Input metadata path rewrite for specified dataset.""" - - @abstractmethod - def unstructured_path_rewrite(self, path): - """Rewrite loc file paths, etc..""" - - @abstractmethod - def working_directory(self): - """ Job working directory (potentially remote) """ - - @abstractmethod - def config_directory(self): - """ Directory containing config files (potentially remote) """ - - @abstractmethod - def env_config_directory(self): - """Working directory (possibly as environment variable evaluation).""" - - @abstractmethod - def sep(self): - """ os.path.sep for the platform this job will execute in. - """ - - @abstractmethod - def new_file_path(self): - """ Absolute path to dump new files for this job on compute server. """ - - @abstractmethod - def tool_directory(self): - """ Absolute path to tool files for this job on compute server. """ - - @abstractmethod - def version_path(self): - """ Location of the version file for the underlying tool. """ - - @abstractmethod - def home_directory(self): - """Home directory of target job - none if HOME should not be set.""" - - @abstractmethod - def tmp_directory(self): - """Temp directory of target job - none if HOME should not be set.""" - - @abstractmethod - def galaxy_url(self): - """URL to access Galaxy API from for this compute environment.""" - - -class SimpleComputeEnvironment: - - def config_directory(self): - return os.path.join(self.working_directory(), "configs") - - def sep(self): - return os.path.sep - - -class SharedComputeEnvironment(SimpleComputeEnvironment, ComputeEnvironment): - """ Default ComputeEnvironment for job and task wrapper to pass - to ToolEvaluator - valid when Galaxy and compute share all the relevant - file systems. - """ - - def __init__(self, job_io: JobIO, job: Job): - self.job_io = job_io - self.job = job - - def output_names(self): - return self.job_io.get_output_basenames() - - def output_paths(self): - return self.job_io.get_output_fnames() - - def input_path_rewrite(self, dataset): - return self.job_io.get_input_path(dataset).false_path - - def output_path_rewrite(self, dataset): - dataset_path = self.job_io.get_output_path(dataset) - if hasattr(dataset_path, "false_path"): - return dataset_path.false_path - else: - return dataset_path - - def input_extra_files_rewrite(self, dataset): - return None - - def output_extra_files_rewrite(self, dataset): - return None - - def input_metadata_rewrite(self, dataset, metadata_value): - return None - - def unstructured_path_rewrite(self, path): - return None - - def working_directory(self): - return self.job_io.working_directory - - def env_config_directory(self): - """Working directory (possibly as environment variable evaluation).""" - return "$_GALAXY_JOB_DIR" - - def new_file_path(self): - return self.job_io.new_file_path - - def version_path(self): - return self.job_io.version_path - - def tool_directory(self): - return self.job_io.tool_directory - - def home_directory(self): - return self.job_io.home_directory - - def tmp_directory(self): - return self.job_io.tmp_directory - - def galaxy_url(self): - return self.job_io.galaxy_url - - class NoopQueue: """ Implements the JobQueue / JobStopQueue interface but does nothing diff --git a/lib/galaxy/jobs/runners/pulsar.py b/lib/galaxy/jobs/runners/pulsar.py index e6d6e3eca69..8906810bfaf 100644 --- a/lib/galaxy/jobs/runners/pulsar.py +++ b/lib/galaxy/jobs/runners/pulsar.py @@ -32,10 +32,8 @@ from pulsar.client import ( from pulsar.client.staging import DEFAULT_DYNAMIC_COLLECTION_PATTERN from galaxy import model -from galaxy.jobs import ( - ComputeEnvironment, - JobDestination -) +from galaxy.job_execution.compute_environment import ComputeEnvironment +from galaxy.jobs import JobDestination from galaxy.jobs.command_factory import build_command from galaxy.jobs.runners import ( AsynchronousJobRunner, diff --git a/lib/galaxy/managers/workflows.py b/lib/galaxy/managers/workflows.py index 82723d27006..ec3befb2072 100644 --- a/lib/galaxy/managers/workflows.py +++ b/lib/galaxy/managers/workflows.py @@ -26,7 +26,7 @@ from galaxy import ( model, util ) -from galaxy.jobs.actions.post import ActionBox +from galaxy.job_execution.actions.post import ActionBox from galaxy.model.item_attrs import UsesAnnotations from galaxy.structured_app import MinimalManagerApp from galaxy.tools.parameters import ( diff --git a/lib/galaxy/tools/__init__.py b/lib/galaxy/tools/__init__.py index c7b44e321fb..fa05a882a8b 100644 --- a/lib/galaxy/tools/__init__.py +++ b/lib/galaxy/tools/__init__.py @@ -9,7 +9,6 @@ import re import tarfile import tempfile import threading -from datetime import datetime from pathlib import Path from typing import ( Any, @@ -68,6 +67,7 @@ from galaxy.tools.actions.data_manager import DataManagerToolAction from galaxy.tools.actions.data_source import DataSourceToolAction from galaxy.tools.actions.model_operations import ModelOperationToolAction from galaxy.tools.cache import ToolDocumentCache +from galaxy.tools.evaluation import global_tool_errors from galaxy.tools.imp_exp import JobImportHistoryArchiveWrapper from galaxy.tools.parameters import ( check_param, @@ -241,25 +241,6 @@ WORKFLOW_SAFE_TOOL_VERSION_UPDATES = { } -class ToolErrorLog: - def __init__(self): - self.error_stack = [] - self.max_errors = 100 - - def add_error(self, file, phase, exception): - self.error_stack.insert(0, { - "file": file, - "time": str(datetime.now()), - "phase": phase, - "error": unicodify(exception) - }) - if len(self.error_stack) > self.max_errors: - self.error_stack.pop() - - -global_tool_errors = ToolErrorLog() - - class ToolNotFoundException(Exception): pass diff --git a/lib/galaxy/tools/actions/__init__.py b/lib/galaxy/tools/actions/__init__.py index b24f0eff6fa..55b71fb233a 100644 --- a/lib/galaxy/tools/actions/__init__.py +++ b/lib/galaxy/tools/actions/__init__.py @@ -8,7 +8,7 @@ from typing import Any, cast, Dict, List, Set, Union from galaxy import model from galaxy.exceptions import ItemAccessibilityException -from galaxy.jobs.actions.post import ActionBox +from galaxy.job_execution.actions.post import ActionBox from galaxy.model import LibraryDatasetDatasetAssociation, WorkflowRequestInputParameter from galaxy.model.dataset_collections.builder import CollectionBuilder from galaxy.model.none_like import NoneDataset diff --git a/lib/galaxy/tools/evaluation.py b/lib/galaxy/tools/evaluation.py index 002a7de4216..8d5405320ab 100644 --- a/lib/galaxy/tools/evaluation.py +++ b/lib/galaxy/tools/evaluation.py @@ -4,15 +4,15 @@ import os import shlex import string import tempfile -from typing import Any, Dict, List +from datetime import datetime +from typing import Any, Callable, Dict, List, Optional from galaxy import model from galaxy.files import ProvidesUserFileSourcesUserContext +from galaxy.job_execution.compute_environment import ComputeEnvironment from galaxy.job_execution.setup import ensure_configs_directory -from galaxy.jobs import ComputeEnvironment from galaxy.model.none_like import NoneDataset from galaxy.security.object_wrapper import wrap_with_safe_string -from galaxy.tools import global_tool_errors from galaxy.tools.parameters import ( visit_input_values, wrapped_json, @@ -92,6 +92,25 @@ What won't work: """ +class ToolErrorLog: + def __init__(self): + self.error_stack = [] + self.max_errors = 100 + + def add_error(self, file, phase, exception): + self.error_stack.insert(0, { + "file": file, + "time": str(datetime.now()), + "phase": phase, + "error": unicodify(exception) + }) + if len(self.error_stack) > self.max_errors: + self.error_stack.pop() + + +global_tool_errors = ToolErrorLog() + + class ToolEvaluator: """ An abstraction linking together a tool and a job runtime to evaluate tool inputs in an isolated, testable manner. diff --git a/lib/galaxy/tools/remote_tool_eval.py b/lib/galaxy/tools/remote_tool_eval.py index d0a2db7574a..b1536753d5b 100644 --- a/lib/galaxy/tools/remote_tool_eval.py +++ b/lib/galaxy/tools/remote_tool_eval.py @@ -13,8 +13,8 @@ from sqlalchemy.orm import scoped_session from galaxy import model from galaxy.datatypes.registry import Registry from galaxy.files import ConfiguredFileSources +from galaxy.job_execution.compute_environment import SharedComputeEnvironment from galaxy.job_execution.setup import JobIO -from galaxy.jobs import SharedComputeEnvironment from galaxy.metadata.set_metadata import ( get_metadata_params, get_object_store, diff --git a/lib/galaxy/tools/wrappers.py b/lib/galaxy/tools/wrappers.py index cb9c0297938..470b84678eb 100644 --- a/lib/galaxy/tools/wrappers.py +++ b/lib/galaxy/tools/wrappers.py @@ -39,7 +39,7 @@ if TYPE_CHECKING: from galaxy.tools import Tool from galaxy.tools.parameters.basic import SelectToolParameter, ToolParameter from galaxy.datatypes.registry import Registry - from galaxy.jobs import ComputeEnvironment + from galaxy.job_execution.compute_environment import ComputeEnvironment from galaxy.model.metadata import MetadataCollection log = logging.getLogger(__name__) diff --git a/lib/galaxy/webapps/galaxy/api/tools.py b/lib/galaxy/webapps/galaxy/api/tools.py index f26eebe6b3f..594a3a964c0 100644 --- a/lib/galaxy/webapps/galaxy/api/tools.py +++ b/lib/galaxy/webapps/galaxy/api/tools.py @@ -10,7 +10,7 @@ from galaxy.managers.collections_util import dictify_dataset_collection_instance from galaxy.managers.hdas import HDAManager from galaxy.managers.histories import HistoryManager from galaxy.model import PostJobAction -from galaxy.tools import global_tool_errors +from galaxy.tools.evaluation import global_tool_errors from galaxy.util.zipstream import ZipstreamWrapper from galaxy.web import ( expose_api, diff --git a/lib/galaxy/workflow/modules.py b/lib/galaxy/workflow/modules.py index 8e1185c6c2a..01558957579 100644 --- a/lib/galaxy/workflow/modules.py +++ b/lib/galaxy/workflow/modules.py @@ -16,7 +16,7 @@ from galaxy import ( web ) from galaxy.exceptions import ToolMissingException -from galaxy.jobs.actions.post import ActionBox +from galaxy.job_execution.actions.post import ActionBox from galaxy.model import PostJobAction, Workflow from galaxy.model.dataset_collections import matching from galaxy.tool_util.parser.output_objects import ToolExpressionOutput diff --git a/setup.cfg b/setup.cfg index 1f05ca9e99b..ff2ca524844 100644 --- a/setup.cfg +++ b/setup.cfg @@ -203,8 +203,6 @@ check_untyped_defs = False check_untyped_defs = False [mypy-galaxy.model.custom_types] check_untyped_defs = False -[mypy-galaxy.jobs.actions.post] -check_untyped_defs = False [mypy-galaxy.job_metrics.collectl.processes] check_untyped_defs = False [mypy-galaxy.datatypes.util.gff_util] diff --git a/test/unit/app/tools/test_evaluation.py b/test/unit/app/tools/test_evaluation.py index 3d4ffdb6998..501b7d81cb2 100644 --- a/test/unit/app/tools/test_evaluation.py +++ b/test/unit/app/tools/test_evaluation.py @@ -2,8 +2,8 @@ import os from unittest import TestCase from galaxy.app_unittest_utils.tools_support import UsesApp +from galaxy.job_execution.compute_environment import SimpleComputeEnvironment from galaxy.job_execution.datasets import DatasetPath -from galaxy.jobs import SimpleComputeEnvironment from galaxy.model import ( Dataset, History, diff --git a/test/unit/app/tools/test_wrappers.py b/test/unit/app/tools/test_wrappers.py index fb94a026280..257cbf3a1c3 100644 --- a/test/unit/app/tools/test_wrappers.py +++ b/test/unit/app/tools/test_wrappers.py @@ -6,8 +6,8 @@ from unittest.mock import Mock import pytest from galaxy.datatypes.metadata import MetadataSpecCollection +from galaxy.job_execution.compute_environment import ComputeEnvironment from galaxy.job_execution.datasets import DatasetPath -from galaxy.jobs import ComputeEnvironment from galaxy.model import DatasetInstance from galaxy.tools.parameters.basic import ( BooleanToolParameter,