mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
Resolve circular dependency for ToolEvaluator
by placing ComputeEnvironment and ActionBox under job_execution.
This commit is contained in:
@@ -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:
|
||||
|
||||
@@ -12,7 +12,7 @@ Subpackages
|
||||
.. toctree::
|
||||
:maxdepth: 4
|
||||
|
||||
galaxy.jobs.actions
|
||||
galaxy.job_execution.actions
|
||||
galaxy.jobs.rules
|
||||
galaxy.jobs.runners
|
||||
galaxy.jobs.splitters
|
||||
|
||||
@@ -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:
|
||||
@@ -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
|
||||
+4
-159
@@ -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
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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 (
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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__)
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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]
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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,
|
||||
|
||||
Reference in New Issue
Block a user