From 33fd1b1a8a46203d73a1f82b1460fc9841638254 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Wed, 5 Jun 2019 08:54:31 -0400 Subject: [PATCH 01/15] Refactor output_checker.py into galaxy-tool-util. Needed for remote job completion. --- lib/galaxy/jobs/__init__.py | 2 +- lib/galaxy/jobs/runners/__init__.py | 2 +- lib/galaxy/{jobs => tool_util}/output_checker.py | 0 packages/tool_util/tests/test_output_checker.py | 1 + .../test_job_output_checker.py => tools/test_output_checker.py} | 2 +- 5 files changed, 4 insertions(+), 3 deletions(-) rename lib/galaxy/{jobs => tool_util}/output_checker.py (100%) create mode 120000 packages/tool_util/tests/test_output_checker.py rename test/unit/{jobs/test_job_output_checker.py => tools/test_output_checker.py} (97%) diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index d6c003eea1d..92c264ed649 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -33,6 +33,7 @@ from galaxy.jobs.runners import BaseJobRunner, JobState from galaxy.metadata import get_metadata_compute_strategy from galaxy.objectstore import ObjectStorePopulator from galaxy.tool_util.deps import requirements +from galaxy.tool_util.output_checker import check_output, DETECTED_JOB_STATE from galaxy.util import safe_makedirs, unicodify from galaxy.util.bunch import Bunch from galaxy.util.expressions import ExpressionContext @@ -40,7 +41,6 @@ from galaxy.util.xml_macros import load from galaxy.web.stack.handlers import ConfiguresHandlers from .datasets import (DatasetPath, NullDatasetPathRewriter, OutputsToWorkingDirectoryPathRewriter, TaskPathRewriter) -from .output_checker import check_output, DETECTED_JOB_STATE log = logging.getLogger(__name__) diff --git a/lib/galaxy/jobs/runners/__init__.py b/lib/galaxy/jobs/runners/__init__.py index ca559898c64..d7036a48ac7 100644 --- a/lib/galaxy/jobs/runners/__init__.py +++ b/lib/galaxy/jobs/runners/__init__.py @@ -19,7 +19,6 @@ from six.moves.queue import ( import galaxy.jobs from galaxy import model from galaxy.jobs.command_factory import build_command -from galaxy.jobs.output_checker import DETECTED_JOB_STATE from galaxy.jobs.runners.util.env import env_to_statement from galaxy.jobs.runners.util.job_script import ( job_script, @@ -29,6 +28,7 @@ from galaxy.tool_util.deps.dependencies import ( JobInfo, ToolInfo ) +from galaxy.tool_util.output_checker import DETECTED_JOB_STATE from galaxy.util import ( DATABASE_MAX_STRING_SIZE, ExecutionTimer, diff --git a/lib/galaxy/jobs/output_checker.py b/lib/galaxy/tool_util/output_checker.py similarity index 100% rename from lib/galaxy/jobs/output_checker.py rename to lib/galaxy/tool_util/output_checker.py diff --git a/packages/tool_util/tests/test_output_checker.py b/packages/tool_util/tests/test_output_checker.py new file mode 120000 index 00000000000..2a54b77ca94 --- /dev/null +++ b/packages/tool_util/tests/test_output_checker.py @@ -0,0 +1 @@ +../../../test/unit/tools/test_output_checker.py \ No newline at end of file diff --git a/test/unit/jobs/test_job_output_checker.py b/test/unit/tools/test_output_checker.py similarity index 97% rename from test/unit/jobs/test_job_output_checker.py rename to test/unit/tools/test_output_checker.py index 112222931f5..b2158f5a327 100644 --- a/test/unit/jobs/test_job_output_checker.py +++ b/test/unit/tools/test_output_checker.py @@ -1,6 +1,6 @@ from unittest import TestCase -from galaxy.jobs.output_checker import check_output, DETECTED_JOB_STATE +from galaxy.tool_util.output_checker import check_output, DETECTED_JOB_STATE from galaxy.tool_util.parser.error_level import StdioErrorLevel from galaxy.tool_util.parser.interface import ToolStdioRegex from galaxy.util.bunch import Bunch From de30bde5e4b52441272b9326c26cd59b10e50001 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Wed, 5 Jun 2019 10:16:52 -0400 Subject: [PATCH 02/15] Refactor galaxy.tools.provided_metadata -> galaxy.tool_util. Needed for remote job output metadata collection. --- lib/galaxy/{tools => tool_util}/provided_metadata.py | 0 lib/galaxy/tools/__init__.py | 2 +- lib/galaxy_ext/metadata/_provided_metadata.py | 2 +- test/unit/tools/test_collect_primary_datasets.py | 2 +- 4 files changed, 3 insertions(+), 3 deletions(-) rename lib/galaxy/{tools => tool_util}/provided_metadata.py (100%) diff --git a/lib/galaxy/tools/provided_metadata.py b/lib/galaxy/tool_util/provided_metadata.py similarity index 100% rename from lib/galaxy/tools/provided_metadata.py rename to lib/galaxy/tool_util/provided_metadata.py diff --git a/lib/galaxy/tools/__init__.py b/lib/galaxy/tools/__init__.py index db5175818a0..89141d202d0 100755 --- a/lib/galaxy/tools/__init__.py +++ b/lib/galaxy/tools/__init__.py @@ -44,6 +44,7 @@ from galaxy.tool_util.parser import ( ToolOutputCollectionPart ) from galaxy.tool_util.parser.xml import XmlPageSource +from galaxy.tool_util.provided_metadata import parse_tool_provided_metadata from galaxy.tools import expressions from galaxy.tools.actions import DefaultToolAction from galaxy.tools.actions.data_manager import DataManagerToolAction @@ -101,7 +102,6 @@ from .execute import ( execute as execute_job, MappingParameters, ) -from .provided_metadata import parse_tool_provided_metadata log = logging.getLogger(__name__) diff --git a/lib/galaxy_ext/metadata/_provided_metadata.py b/lib/galaxy_ext/metadata/_provided_metadata.py index 6e81f9dc9f3..b44719c00dc 120000 --- a/lib/galaxy_ext/metadata/_provided_metadata.py +++ b/lib/galaxy_ext/metadata/_provided_metadata.py @@ -1 +1 @@ -../../galaxy/tools/provided_metadata.py \ No newline at end of file +../../galaxy/tool_util/provided_metadata.py \ No newline at end of file diff --git a/test/unit/tools/test_collect_primary_datasets.py b/test/unit/tools/test_collect_primary_datasets.py index 2b63db0a3d1..d33d5e4bb84 100644 --- a/test/unit/tools/test_collect_primary_datasets.py +++ b/test/unit/tools/test_collect_primary_datasets.py @@ -7,7 +7,7 @@ from galaxy import ( util ) from galaxy.tool_util.parser import output_collection_def -from galaxy.tools.provided_metadata import LegacyToolProvidedMetadata, NullToolProvidedMetadata +from galaxy.tool_util.provided_metadata import LegacyToolProvidedMetadata, NullToolProvidedMetadata from .. import tools_support DEFAULT_TOOL_OUTPUT = "out1" From 1405446ec14acd9508ab36679e4a4f88d9b39310 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Wed, 5 Jun 2019 08:49:07 -0400 Subject: [PATCH 03/15] Refactor output_collect into galaxy-tool-util. Needed for dynamic discovery of outputs on remote servers. --- lib/galaxy/{tools/parameters => tool_util}/output_collect.py | 0 lib/galaxy/tools/__init__.py | 2 +- 2 files changed, 1 insertion(+), 1 deletion(-) rename lib/galaxy/{tools/parameters => tool_util}/output_collect.py (100%) diff --git a/lib/galaxy/tools/parameters/output_collect.py b/lib/galaxy/tool_util/output_collect.py similarity index 100% rename from lib/galaxy/tools/parameters/output_collect.py rename to lib/galaxy/tool_util/output_collect.py diff --git a/lib/galaxy/tools/__init__.py b/lib/galaxy/tools/__init__.py index 89141d202d0..8227ef08ba0 100755 --- a/lib/galaxy/tools/__init__.py +++ b/lib/galaxy/tools/__init__.py @@ -29,6 +29,7 @@ from galaxy.managers.jobs import JobSearch from galaxy.metadata import get_metadata_compute_strategy from galaxy.model.tags import GalaxyTagHandler from galaxy.queue_worker import send_control_task +from galaxy.tool_util import output_collect from galaxy.tool_util.deps import ( CachedDependencyManager, ) @@ -58,7 +59,6 @@ from galaxy.tools.parameters import ( populate_state, visit_input_values ) -from galaxy.tools.parameters import output_collect from galaxy.tools.parameters.basic import ( BaseURLToolParameter, DataCollectionToolParameter, From f9baa0b4839c70e14d5e66f988214c3b05abebdd Mon Sep 17 00:00:00 2001 From: John Chilton Date: Wed, 5 Jun 2019 17:47:59 -0400 Subject: [PATCH 04/15] Stub out job_exection package for remote job finish code. --- lib/galaxy/job_execution/__init__.py | 0 .../output_collect.py | 0 lib/galaxy/tools/__init__.py | 4 +- packages/job_execution/HISTORY.rst | 12 +++ packages/job_execution/LICENSE | 1 + packages/job_execution/MANIFEST.in | 2 + packages/job_execution/Makefile | 1 + packages/job_execution/README.rst | 14 +++ packages/job_execution/dev-requirements.txt | 1 + packages/job_execution/galaxy/__init__.py | 1 + packages/job_execution/galaxy/job_execution | 1 + .../galaxy/project_galaxy_job_execution.py | 13 +++ packages/job_execution/requirements.txt | 2 + packages/job_execution/scripts | 1 + packages/job_execution/setup.cfg | 1 + packages/job_execution/setup.py | 100 ++++++++++++++++++ packages/test.sh | 3 +- 17 files changed, 155 insertions(+), 2 deletions(-) create mode 100644 lib/galaxy/job_execution/__init__.py rename lib/galaxy/{tool_util => job_execution}/output_collect.py (100%) create mode 100644 packages/job_execution/HISTORY.rst create mode 120000 packages/job_execution/LICENSE create mode 100644 packages/job_execution/MANIFEST.in create mode 120000 packages/job_execution/Makefile create mode 100644 packages/job_execution/README.rst create mode 120000 packages/job_execution/dev-requirements.txt create mode 100644 packages/job_execution/galaxy/__init__.py create mode 120000 packages/job_execution/galaxy/job_execution create mode 100644 packages/job_execution/galaxy/project_galaxy_job_execution.py create mode 100644 packages/job_execution/requirements.txt create mode 120000 packages/job_execution/scripts create mode 120000 packages/job_execution/setup.cfg create mode 100644 packages/job_execution/setup.py diff --git a/lib/galaxy/job_execution/__init__.py b/lib/galaxy/job_execution/__init__.py new file mode 100644 index 00000000000..e69de29bb2d diff --git a/lib/galaxy/tool_util/output_collect.py b/lib/galaxy/job_execution/output_collect.py similarity index 100% rename from lib/galaxy/tool_util/output_collect.py rename to lib/galaxy/job_execution/output_collect.py diff --git a/lib/galaxy/tools/__init__.py b/lib/galaxy/tools/__init__.py index 8227ef08ba0..7aa13edf6d9 100755 --- a/lib/galaxy/tools/__init__.py +++ b/lib/galaxy/tools/__init__.py @@ -25,11 +25,11 @@ from galaxy import ( exceptions, model ) +from galaxy.job_execution import output_collect from galaxy.managers.jobs import JobSearch from galaxy.metadata import get_metadata_compute_strategy from galaxy.model.tags import GalaxyTagHandler from galaxy.queue_worker import send_control_task -from galaxy.tool_util import output_collect from galaxy.tool_util.deps import ( CachedDependencyManager, ) @@ -1737,6 +1737,8 @@ class Tool(Dictifiable): Find any additional datasets generated by a tool and attach (for cases where number of outputs is not known in advance). """ + # given the job_execution import is the only one, probably makes sense to refactor this out + # into job_wrapper. tool = self permission_provider = output_collect.PermissionProvider(inp_data, tool.app.security_agent, job) metadata_source_provider = output_collect.MetadataSourceProvider(inp_data) diff --git a/packages/job_execution/HISTORY.rst b/packages/job_execution/HISTORY.rst new file mode 100644 index 00000000000..ec2f0e746e9 --- /dev/null +++ b/packages/job_execution/HISTORY.rst @@ -0,0 +1,12 @@ +.. :changelog: + +History +------- + +.. to_doc + +--------------------- +19.9.0.dev0 +--------------------- + +* Initial import from dev branch of Galaxy during 19.09 development cycle. diff --git a/packages/job_execution/LICENSE b/packages/job_execution/LICENSE new file mode 120000 index 00000000000..1ef648f64b3 --- /dev/null +++ b/packages/job_execution/LICENSE @@ -0,0 +1 @@ +../../LICENSE.txt \ No newline at end of file diff --git a/packages/job_execution/MANIFEST.in b/packages/job_execution/MANIFEST.in new file mode 100644 index 00000000000..2bf148a5764 --- /dev/null +++ b/packages/job_execution/MANIFEST.in @@ -0,0 +1,2 @@ +include *.rst LICENSE + diff --git a/packages/job_execution/Makefile b/packages/job_execution/Makefile new file mode 120000 index 00000000000..37af8bae5ba --- /dev/null +++ b/packages/job_execution/Makefile @@ -0,0 +1 @@ +../package.Makefile \ No newline at end of file diff --git a/packages/job_execution/README.rst b/packages/job_execution/README.rst new file mode 100644 index 00000000000..897422a8aa3 --- /dev/null +++ b/packages/job_execution/README.rst @@ -0,0 +1,14 @@ + +.. image:: https://badge.fury.io/py/galaxy-util.svg + :target: https://pypi.python.org/pypi/galaxy-util/ + + +Overview +-------- + +The Galaxy_ job execution runtime module. + +* Free software: Academic Free License version 3.0 +* Code: https://github.com/galaxyproject/galaxy + +.. _Galaxy: http://galaxyproject.org/ diff --git a/packages/job_execution/dev-requirements.txt b/packages/job_execution/dev-requirements.txt new file mode 120000 index 00000000000..467b90d7a23 --- /dev/null +++ b/packages/job_execution/dev-requirements.txt @@ -0,0 +1 @@ +../package-dev-requirements.txt \ No newline at end of file diff --git a/packages/job_execution/galaxy/__init__.py b/packages/job_execution/galaxy/__init__.py new file mode 100644 index 00000000000..69e3be50dac --- /dev/null +++ b/packages/job_execution/galaxy/__init__.py @@ -0,0 +1 @@ +__path__ = __import__('pkgutil').extend_path(__path__, __name__) diff --git a/packages/job_execution/galaxy/job_execution b/packages/job_execution/galaxy/job_execution new file mode 120000 index 00000000000..4ef157cced2 --- /dev/null +++ b/packages/job_execution/galaxy/job_execution @@ -0,0 +1 @@ +../../../lib/galaxy/job_execution \ No newline at end of file diff --git a/packages/job_execution/galaxy/project_galaxy_job_execution.py b/packages/job_execution/galaxy/project_galaxy_job_execution.py new file mode 100644 index 00000000000..a30f915133d --- /dev/null +++ b/packages/job_execution/galaxy/project_galaxy_job_execution.py @@ -0,0 +1,13 @@ +# -*- coding: utf-8 -*- + +__version__ = '19.9.0.dev0' + +PROJECT_NAME = "galaxy-job-execution" +PROJECT_OWNER = PROJECT_USERAME = "galaxyproject" +PROJECT_URL = "https://github.com/galaxyproject/galaxy" +PROJECT_AUTHOR = 'Galaxy Project and Community' +PROJECT_DESCRIPTION = 'Galaxy Job Execution Runtime Utilities' +PROJECT_EMAIL = 'jmchilton@gmail.com' +RAW_CONTENT_URL = "https://raw.github.com/%s/%s/master/" % ( + PROJECT_USERAME, PROJECT_NAME +) diff --git a/packages/job_execution/requirements.txt b/packages/job_execution/requirements.txt new file mode 100644 index 00000000000..77de615ea9c --- /dev/null +++ b/packages/job_execution/requirements.txt @@ -0,0 +1,2 @@ +galaxy-data + diff --git a/packages/job_execution/scripts b/packages/job_execution/scripts new file mode 120000 index 00000000000..9aec9dc5a06 --- /dev/null +++ b/packages/job_execution/scripts @@ -0,0 +1 @@ +../build_scripts \ No newline at end of file diff --git a/packages/job_execution/setup.cfg b/packages/job_execution/setup.cfg new file mode 120000 index 00000000000..eb7cf09393f --- /dev/null +++ b/packages/job_execution/setup.cfg @@ -0,0 +1 @@ +../../setup.cfg \ No newline at end of file diff --git a/packages/job_execution/setup.py b/packages/job_execution/setup.py new file mode 100644 index 00000000000..0ea5361f7d5 --- /dev/null +++ b/packages/job_execution/setup.py @@ -0,0 +1,100 @@ +#!/usr/bin/env python +# -*- coding: utf-8 -*- + +import ast +import os +import re +try: + from setuptools import setup +except ImportError: + from distutils.core import setup + +SOURCE_DIR = "galaxy" + +_version_re = re.compile(r'__version__\s+=\s+(.*)') + +with open('%s/project_galaxy_job_execution.py' % SOURCE_DIR, 'rb') as f: + init_contents = f.read().decode('utf-8') + + def get_var(var_name): + pattern = re.compile(r'%s\s+=\s+(.*)' % var_name) + match = pattern.search(init_contents).group(1) + return str(ast.literal_eval(match)) + + version = get_var("__version__") + PROJECT_NAME = get_var("PROJECT_NAME") + PROJECT_URL = get_var("PROJECT_URL") + PROJECT_AUTHOR = get_var("PROJECT_AUTHOR") + PROJECT_EMAIL = get_var("PROJECT_EMAIL") + PROJECT_DESCRIPTION = get_var("PROJECT_DESCRIPTION") + +TEST_DIR = 'tests' +PACKAGES = [ + 'galaxy', + 'galaxy.job_execution', +] +ENTRY_POINTS = ''' + [console_scripts] +''' +PACKAGE_DATA = { + # Be sure to update MANIFEST.in for source dist. + 'galaxy': [ + ], +} +PACKAGE_DIR = { + SOURCE_DIR: SOURCE_DIR, +} + +readme = open('README.rst').read() +history = open('HISTORY.rst').read().replace('.. :changelog:', '') + +if os.path.exists("requirements.txt"): + requirements = open("requirements.txt").read().split("\n") +else: + # In tox, it will cover them anyway. + requirements = [] + + +test_requirements = [ + # TODO: put package test requirements here +] + + +setup( + name=PROJECT_NAME, + version=version, + description=PROJECT_DESCRIPTION, + long_description=readme + '\n\n' + history, + long_description_content_type='text/x-rst', + author=PROJECT_AUTHOR, + author_email=PROJECT_EMAIL, + url=PROJECT_URL, + packages=PACKAGES, + entry_points=ENTRY_POINTS, + package_data=PACKAGE_DATA, + package_dir=PACKAGE_DIR, + include_package_data=True, + install_requires=requirements, + extras_require={}, + license="AFL", + zip_safe=False, + keywords='galaxy', + classifiers=[ + 'Development Status :: 5 - Production/Stable', + 'Intended Audience :: Developers', + 'Environment :: Console', + 'License :: OSI Approved :: Academic Free License (AFL)', + 'Operating System :: POSIX', + 'Topic :: Software Development', + 'Topic :: Software Development :: Code Generators', + 'Topic :: Software Development :: Testing', + 'Natural Language :: English', + "Programming Language :: Python :: 2", + 'Programming Language :: Python :: 2.7', + 'Programming Language :: Python :: 3.5', + 'Programming Language :: Python :: 3.6', + 'Programming Language :: Python :: 3.7', + ], + test_suite=TEST_DIR, + tests_require=test_requirements +) diff --git a/packages/test.sh b/packages/test.sh index 5deb29bf919..a597d69ee1a 100755 --- a/packages/test.sh +++ b/packages/test.sh @@ -20,10 +20,11 @@ PACKAGE_DIRS=( containers tool_util data + job_execution ) # containers has no tests, tool_util not yet working 100%, # data has many problems quota, tool shed install database, etc.. -RUN_TESTS=(1 1 1 0 0 0 0) +RUN_TESTS=(1 1 1 0 0 0 0 0) for ((i=0; i<${#PACKAGE_DIRS[@]}; i++)); do package_dir=${PACKAGE_DIRS[$i]} From fea4356acf9c09098c6c4ce7493ab166bb8727fb Mon Sep 17 00:00:00 2001 From: John Chilton Date: Mon, 17 Dec 2018 14:00:19 -0500 Subject: [PATCH 05/15] Refactor toward alternative extended-metadata collection. --- lib/galaxy/jobs/__init__.py | 183 ++++++++++++++++++------------------ 1 file changed, 94 insertions(+), 89 deletions(-) diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index 92c264ed649..ec2902c655d 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -1390,6 +1390,97 @@ class JobWrapper(HasResourceParameters): job.object_store_id = object_store_populator.object_store_id self._setup_working_directory(job=job) + def _finish_dataset(self, output_name, dataset, job, context, final_job_state, remote_metadata_directory): + implicit_collection_jobs = job.implicit_collection_jobs_association + purged = dataset.dataset.purged + if not purged and dataset.dataset.external_filename is None: + trynum = 0 + while trynum < self.app.config.retry_job_output_collection: + try: + # Attempt to short circuit NFS attribute caching + os.stat(dataset.dataset.file_name) + os.chown(dataset.dataset.file_name, os.getuid(), -1) + trynum = self.app.config.retry_job_output_collection + except (OSError, ObjectNotFound) as e: + trynum += 1 + log.warning('Error accessing dataset with ID %i, will retry: %s', dataset.dataset.id, e) + time.sleep(2) + if getattr(dataset, "hidden_beneath_collection_instance", None): + dataset.visible = False + dataset.blurb = 'done' + dataset.peek = 'no peek' + dataset.info = (dataset.info or '') + if context['stdout'].strip(): + # Ensure white space between entries + dataset.info = dataset.info.rstrip() + "\n" + context['stdout'].strip() + if context['stderr'].strip(): + # Ensure white space between entries + dataset.info = dataset.info.rstrip() + "\n" + context['stderr'].strip() + dataset.tool_version = self.version_string + dataset.set_size() + if 'uuid' in context: + dataset.dataset.uuid = context['uuid'] + self.__update_output(job, dataset) + if not purged: + self._collect_extra_files(dataset.dataset, self.working_directory) + # Handle composite datatypes of auto_primary_file type + if dataset.datatype.composite_type == 'auto_primary_file' and not dataset.has_data(): + try: + with NamedTemporaryFile(mode='w') as temp_fh: + temp_fh.write(dataset.datatype.generate_primary_file(dataset)) + temp_fh.flush() + self.object_store.update_from_file(dataset.dataset, file_name=temp_fh.name, create=True) + dataset.set_size() + except Exception as e: + log.warning('Unable to generate primary composite file automatically for %s: %s', dataset.dataset.id, e) + if job.states.ERROR == final_job_state: + dataset.blurb = "error" + if not implicit_collection_jobs: + # Only unhide dataset outputs that are not part of a implicit collection + dataset.mark_unhidden() + elif not purged: + # If the tool was expected to set the extension, attempt to retrieve it + if dataset.ext == 'auto': + dataset.extension = context.get('ext', 'data') + dataset.init_meta(copy_from=dataset) + # if a dataset was copied, it won't appear in our dictionary: + # either use the metadata from originating output dataset, or call set_meta on the copies + # it would be quicker to just copy the metadata from the originating output dataset, + # but somewhat trickier (need to recurse up the copied_from tree), for now we'll call set_meta() + retry_internally = util.asbool(self.get_destination_configuration("retry_metadata_internally", True)) + metadata_set_successfully = self.external_output_metadata.external_metadata_set_successfully(dataset, output_name, self.sa_session, working_directory=self.working_directory) + if retry_internally and not metadata_set_successfully: + # If Galaxy was expected to sniff type and didn't - do so. + if dataset.ext == "_sniff_": + extension = sniff.handle_uploaded_dataset_file(dataset.dataset.file_name, self.app.datatypes_registry) + dataset.extension = extension + + # call datatype.set_meta directly for the initial set_meta call during dataset creation + dataset.datatype.set_meta(dataset, overwrite=False) + elif (job.states.ERROR != final_job_state and not metadata_set_successfully): + dataset._state = model.Dataset.states.FAILED_METADATA + else: + self.external_output_metadata.load_metadata(dataset, output_name, self.sa_session, working_directory=self.working_directory, remote_metadata_directory=remote_metadata_directory) + line_count = context.get('line_count', None) + try: + # Certain datatype's set_peek methods contain a line_count argument + dataset.set_peek(line_count=line_count) + except TypeError: + # ... and others don't + dataset.set_peek() + else: + # Handle purged datasets. + dataset.blurb = "empty" + if dataset.ext == 'auto': + dataset.extension = context.get('ext', 'txt') + + for context_key in TOOL_PROVIDED_JOB_METADATA_KEYS: + if context_key in context: + context_value = context[context_key] + setattr(dataset, context_key, context_value) + + self.sa_session.add(dataset) + def finish( self, tool_stdout, @@ -1476,101 +1567,15 @@ class JobWrapper(HasResourceParameters): return self.fail("Job %s's output dataset(s) could not be read" % job.id) job_context = ExpressionContext(dict(stdout=job.stdout, stderr=job.stderr)) - implicit_collection_jobs = job.implicit_collection_jobs_association for dataset_assoc in job.output_datasets + job.output_library_datasets: context = self.get_dataset_finish_context(job_context, dataset_assoc) # should this also be checking library associations? - can a library item be added from a history before the job has ended? - # lets not allow this to occur # need to update all associated output hdas, i.e. history was shared with job running for dataset in dataset_assoc.dataset.dataset.history_associations + dataset_assoc.dataset.dataset.library_associations: - purged = dataset.dataset.purged - if not purged and dataset.dataset.external_filename is None: - trynum = 0 - while trynum < self.app.config.retry_job_output_collection: - try: - # Attempt to short circuit NFS attribute caching - os.stat(dataset.dataset.file_name) - os.chown(dataset.dataset.file_name, os.getuid(), -1) - trynum = self.app.config.retry_job_output_collection - except (OSError, ObjectNotFound) as e: - trynum += 1 - log.warning('Error accessing dataset with ID %i, will retry: %s', dataset.dataset.id, e) - time.sleep(2) - if getattr(dataset, "hidden_beneath_collection_instance", None): - dataset.visible = False - dataset.blurb = 'done' - dataset.peek = 'no peek' - dataset.info = (dataset.info or '') - if context['stdout'].strip(): - # Ensure white space between entries - dataset.info = dataset.info.rstrip() + "\n" + context['stdout'].strip() - if context['stderr'].strip(): - # Ensure white space between entries - dataset.info = dataset.info.rstrip() + "\n" + context['stderr'].strip() - dataset.tool_version = self.version_string - dataset.set_size() - if 'uuid' in context: - dataset.dataset.uuid = context['uuid'] - self.__update_output(job, dataset) - if not purged: - self._collect_extra_files(dataset.dataset, self.working_directory) - # Handle composite datatypes of auto_primary_file type - if dataset.datatype.composite_type == 'auto_primary_file' and not dataset.has_data(): - try: - with NamedTemporaryFile(mode='w') as temp_fh: - temp_fh.write(dataset.datatype.generate_primary_file(dataset)) - temp_fh.flush() - self.object_store.update_from_file(dataset.dataset, file_name=temp_fh.name, create=True) - dataset.set_size() - except Exception as e: - log.warning('Unable to generate primary composite file automatically for %s: %s', dataset.dataset.id, e) - if job.states.ERROR == final_job_state: - dataset.blurb = "error" - if not implicit_collection_jobs: - # Only unhide dataset outputs that are not part of a implicit collection - dataset.mark_unhidden() - elif not purged: - # If the tool was expected to set the extension, attempt to retrieve it - if dataset.ext == 'auto': - dataset.extension = context.get('ext', 'data') - dataset.init_meta(copy_from=dataset) - # if a dataset was copied, it won't appear in our dictionary: - # either use the metadata from originating output dataset, or call set_meta on the copies - # it would be quicker to just copy the metadata from the originating output dataset, - # but somewhat trickier (need to recurse up the copied_from tree), for now we'll call set_meta() - retry_internally = util.asbool(self.get_destination_configuration("retry_metadata_internally", True)) - metadata_set_successfully = self.external_output_metadata.external_metadata_set_successfully(dataset, dataset_assoc.name, self.sa_session, working_directory=self.working_directory) - if retry_internally and not metadata_set_successfully: - # If Galaxy was expected to sniff type and didn't - do so. - if dataset.ext == "_sniff_": - extension = sniff.handle_uploaded_dataset_file(dataset.dataset.file_name, self.app.datatypes_registry) - dataset.extension = extension - - # call datatype.set_meta directly for the initial set_meta call during dataset creation - dataset.datatype.set_meta(dataset, overwrite=False) - elif (job.states.ERROR != final_job_state and not metadata_set_successfully): - dataset._state = model.Dataset.states.FAILED_METADATA - else: - self.external_output_metadata.load_metadata(dataset, dataset_assoc.name, self.sa_session, working_directory=self.working_directory, remote_metadata_directory=remote_metadata_directory) - line_count = context.get('line_count', None) - try: - # Certain datatype's set_peek methods contain a line_count argument - dataset.set_peek(line_count=line_count) - except TypeError: - # ... and others don't - dataset.set_peek() - else: - # Handle purged datasets. - dataset.blurb = "empty" - if dataset.ext == 'auto': - dataset.extension = context.get('ext', 'txt') - - for context_key in TOOL_PROVIDED_JOB_METADATA_KEYS: - if context_key in context: - context_value = context[context_key] - setattr(dataset, context_key, context_value) - - self.sa_session.add(dataset) + self._finish_dataset( + dataset_assoc.name, dataset, job, context, final_job_state, remote_metadata_directory + ) if job.states.ERROR == final_job_state: log.debug("(%s) setting dataset %s state to ERROR", job.id, dataset_assoc.dataset.dataset.id) # TODO: This is where the state is being set to error. Change it! From 65c3f6476af3471699a70e9fc48678853688a632 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Wed, 3 Apr 2019 09:37:06 -0400 Subject: [PATCH 06/15] Refactor job wrapper/runner utilities for reuse in extended set metadata. --- lib/galaxy/job_execution/output_collect.py | 61 ++++++++++++++++++++++ lib/galaxy/jobs/__init__.py | 39 +------------- lib/galaxy/jobs/runners/__init__.py | 24 ++------- lib/galaxy/jobs/runners/local.py | 5 +- 4 files changed, 69 insertions(+), 60 deletions(-) diff --git a/lib/galaxy/job_execution/output_collect.py b/lib/galaxy/job_execution/output_collect.py index 21a7aaf935b..a5e21fb68c5 100644 --- a/lib/galaxy/job_execution/output_collect.py +++ b/lib/galaxy/job_execution/output_collect.py @@ -4,6 +4,7 @@ import logging import operator import os import re +from tempfile import NamedTemporaryFile import galaxy.model from galaxy.model.dataset_collections import builder @@ -463,3 +464,63 @@ def _compose(f, g): DEFAULT_DATASET_COLLECTOR = DatasetCollector(DEFAULT_DATASET_COLLECTOR_DESCRIPTION) DEFAULT_TOOL_PROVIDED_DATASET_COLLECTOR = ToolMetadataDatasetCollector(ToolProvidedMetadataDatasetCollection()) + + +def read_exit_code_from(exit_code_file, id_tag): + """Read exit code reported for a Galaxy job.""" + try: + # This should be an 8-bit exit code, but read ahead anyway: + exit_code_str = open(exit_code_file, "r").read(32) + except Exception: + # By default, the exit code is 0, which typically indicates success. + exit_code_str = "0" + + try: + # Decode the exit code. If it's bogus, then just use 0. + exit_code = int(exit_code_str) + except ValueError: + galaxy_id_tag = id_tag + log.warning("(%s) Exit code '%s' invalid. Using 0." % (galaxy_id_tag, exit_code_str)) + exit_code = 0 + + return exit_code + + +def default_exit_code_file(files_dir, id_tag): + return os.path.join(files_dir, 'galaxy_%s.ec' % id_tag) + + +def collect_extra_files(object_store, dataset, job_working_directory): + store_by = getattr(object_store, "store_by", "id") + file_name = "dataset_%s_files" % getattr(dataset.dataset, store_by) + temp_file_path = os.path.join(job_working_directory, file_name) + extra_dir = None + try: + # This skips creation of directories - object store + # automatically creates them. However, empty directories will + # not be created in the object store at all, which might be a + # problem. + for root, dirs, files in os.walk(temp_file_path): + extra_dir = root.replace(job_working_directory, '', 1).lstrip(os.path.sep) + for f in files: + object_store.update_from_file( + dataset.dataset, + extra_dir=extra_dir, + alt_name=f, + file_name=os.path.join(root, f), + create=True, + preserve_symlinks=True + ) + except Exception as e: + log.debug("Error in collect_associated_files: %s" % (e)) + + # Handle composite datatypes of auto_primary_file type + if dataset.datatype.composite_type == 'auto_primary_file' and not dataset.has_data(): + try: + with NamedTemporaryFile(mode='w') as temp_fh: + temp_fh.write(dataset.datatype.generate_primary_file(dataset)) + temp_fh.flush() + object_store.update_from_file(dataset.dataset, file_name=temp_fh.name, create=True) + dataset.set_size() + except Exception as e: + log.warning('Unable to generate primary composite file automatically for %s: %s', dataset.dataset.id, e) diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index ec2902c655d..96f557ef29e 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -16,7 +16,6 @@ import time import traceback from abc import ABCMeta, abstractmethod from json import loads -from tempfile import NamedTemporaryFile from xml.etree import ElementTree import six @@ -27,6 +26,7 @@ import galaxy from galaxy import model, util from galaxy.datatypes import sniff from galaxy.exceptions import ObjectInvalid, ObjectNotFound +from galaxy.job_execution.output_collect import collect_extra_files from galaxy.jobs.actions.post import ActionBox from galaxy.jobs.mapper import JobMappingException, JobRunnerMapper from galaxy.jobs.runners import BaseJobRunner, JobState @@ -1422,17 +1422,7 @@ class JobWrapper(HasResourceParameters): dataset.dataset.uuid = context['uuid'] self.__update_output(job, dataset) if not purged: - self._collect_extra_files(dataset.dataset, self.working_directory) - # Handle composite datatypes of auto_primary_file type - if dataset.datatype.composite_type == 'auto_primary_file' and not dataset.has_data(): - try: - with NamedTemporaryFile(mode='w') as temp_fh: - temp_fh.write(dataset.datatype.generate_primary_file(dataset)) - temp_fh.flush() - self.object_store.update_from_file(dataset.dataset, file_name=temp_fh.name, create=True) - dataset.set_size() - except Exception as e: - log.warning('Unable to generate primary composite file automatically for %s: %s', dataset.dataset.id, e) + collect_extra_files(self.object_store, dataset, self.working_directory) if job.states.ERROR == final_job_state: dataset.blurb = "error" if not implicit_collection_jobs: @@ -1706,31 +1696,6 @@ class JobWrapper(HasResourceParameters): except Exception: log.exception("Unable to cleanup job %d", self.job_id) - def _collect_extra_files(self, dataset, job_working_directory): - object_store = self.app.object_store - store_by = getattr(object_store, "store_by", "id") - file_name = "dataset_%s_files" % getattr(dataset, store_by) - temp_file_path = os.path.join(job_working_directory, file_name) - extra_dir = None - try: - # This skips creation of directories - object store - # automatically creates them. However, empty directories will - # not be created in the object store at all, which might be a - # problem. - for root, dirs, files in os.walk(temp_file_path): - extra_dir = root.replace(job_working_directory, '', 1).lstrip(os.path.sep) - for f in files: - self.object_store.update_from_file( - dataset, - extra_dir=extra_dir, - alt_name=f, - file_name=os.path.join(root, f), - create=True, - preserve_symlinks=True - ) - except Exception as e: - log.debug("Error in collect_associated_files: %s" % (e)) - def _collect_metrics(self, has_metrics, job_metrics_directory=None): job = has_metrics.get_job() job_metrics_directory = job_metrics_directory or self.working_directory diff --git a/lib/galaxy/jobs/runners/__init__.py b/lib/galaxy/jobs/runners/__init__.py index d7036a48ac7..095f6e3527b 100644 --- a/lib/galaxy/jobs/runners/__init__.py +++ b/lib/galaxy/jobs/runners/__init__.py @@ -18,6 +18,7 @@ from six.moves.queue import ( import galaxy.jobs from galaxy import model +from galaxy.job_execution.output_collect import default_exit_code_file, read_exit_code_from from galaxy.jobs.command_factory import build_command from galaxy.jobs.runners.util.env import env_to_statement from galaxy.jobs.runners.util.job_script import ( @@ -544,7 +545,7 @@ class JobState(object): self.job_file = JobState.default_job_file(files_dir, id_tag) self.output_file = os.path.join(files_dir, 'galaxy_%s.o' % id_tag) self.error_file = os.path.join(files_dir, 'galaxy_%s.e' % id_tag) - self.exit_code_file = os.path.join(files_dir, 'galaxy_%s.ec' % id_tag) + self.exit_code_file = default_exit_code_file(files_dir, id_tag) job_name = 'g%s' % id_tag if self.job_wrapper.tool.old_id: job_name += '_%s' % self.job_wrapper.tool.old_id @@ -556,27 +557,8 @@ class JobState(object): def default_job_file(files_dir, id_tag): return os.path.join(files_dir, 'galaxy_%s.sh' % id_tag) - @staticmethod - def default_exit_code_file(files_dir, id_tag): - return os.path.join(files_dir, 'galaxy_%s.ec' % id_tag) - def read_exit_code(self): - try: - # This should be an 8-bit exit code, but read ahead anyway: - exit_code_str = open(self.exit_code_file, "r").read(32) - except Exception: - # By default, the exit code is 0, which typically indicates success. - exit_code_str = "0" - - try: - # Decode the exit code. If it's bogus, then just use 0. - exit_code = int(exit_code_str) - except ValueError: - galaxy_id_tag = self.job_wrapper.get_id_tag() - log.warning("(%s) Exit code '%s' invalid. Using 0." % (galaxy_id_tag, exit_code_str)) - exit_code = 0 - - return exit_code + return read_exit_code_from(self.exit_code_file, self.job_wrapper.get_id_tag()) def cleanup(self): for file in [getattr(self, a) for a in self.cleanup_file_attributes if hasattr(self, a)]: diff --git a/lib/galaxy/jobs/runners/local.py b/lib/galaxy/jobs/runners/local.py index 799187b127e..0f64d64129e 100644 --- a/lib/galaxy/jobs/runners/local.py +++ b/lib/galaxy/jobs/runners/local.py @@ -11,6 +11,7 @@ import threading from time import sleep from galaxy import model +from galaxy.job_execution.output_collect import default_exit_code_file from galaxy.util import ( asbool, ) @@ -65,7 +66,7 @@ class LocalJobRunner(BaseJobRunner): job_id = job_wrapper.get_id_tag() job_file = JobState.default_job_file(job_wrapper.working_directory, job_id) - exit_code_path = JobState.default_exit_code_file(job_wrapper.working_directory, job_id) + exit_code_path = default_exit_code_file(job_wrapper.working_directory, job_id) job_script_props = { 'slots_statement': slots_statement, 'command': command_line, @@ -137,7 +138,7 @@ class LocalJobRunner(BaseJobRunner): job_destination = job_wrapper.job_destination job_state = JobState(job_wrapper, job_destination) - job_state.exit_code_file = JobState.default_exit_code_file(job_wrapper.working_directory, job_id) + job_state.exit_code_file = default_exit_code_file(job_wrapper.working_directory, job_id) job_state.stop_job = False self._finish_or_resubmit_job(job_state, stdout, stderr, job_id=job_id) From 83ba4517c939214e1fb2910f2651e0697ff9a2c8 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Thu, 6 Jun 2019 09:08:38 -0400 Subject: [PATCH 07/15] Refactor galaxy.jobs.datasets into galaxy.job_execution.datasets. Needed for metadata evaluation stuff (at least in testing) and removes some circular dependencies in the form of stuff in galaxy.tools depending on stuff in galaxy.jobs. --- lib/galaxy/{jobs => job_execution}/datasets.py | 0 lib/galaxy/jobs/__init__.py | 8 ++++++-- lib/galaxy/tools/actions/metadata.py | 2 +- lib/galaxy/tools/evaluation.py | 2 +- packages/job_execution/tests/test_datasets.py | 1 + test/unit/jobs/test_datasets.py | 2 +- test/unit/tools/test_evaluation.py | 2 +- test/unit/tools/test_metadata.py | 2 +- test/unit/tools/test_wrappers.py | 2 +- 9 files changed, 13 insertions(+), 8 deletions(-) rename lib/galaxy/{jobs => job_execution}/datasets.py (100%) create mode 120000 packages/job_execution/tests/test_datasets.py diff --git a/lib/galaxy/jobs/datasets.py b/lib/galaxy/job_execution/datasets.py similarity index 100% rename from lib/galaxy/jobs/datasets.py rename to lib/galaxy/job_execution/datasets.py diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index 96f557ef29e..e69f8cb30d2 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -26,6 +26,12 @@ import galaxy from galaxy import model, util from galaxy.datatypes import sniff from galaxy.exceptions import ObjectInvalid, ObjectNotFound +from galaxy.job_execution.datasets import ( + DatasetPath, + NullDatasetPathRewriter, + OutputsToWorkingDirectoryPathRewriter, + TaskPathRewriter +) from galaxy.job_execution.output_collect import collect_extra_files from galaxy.jobs.actions.post import ActionBox from galaxy.jobs.mapper import JobMappingException, JobRunnerMapper @@ -39,8 +45,6 @@ from galaxy.util.bunch import Bunch from galaxy.util.expressions import ExpressionContext from galaxy.util.xml_macros import load from galaxy.web.stack.handlers import ConfiguresHandlers -from .datasets import (DatasetPath, NullDatasetPathRewriter, - OutputsToWorkingDirectoryPathRewriter, TaskPathRewriter) log = logging.getLogger(__name__) diff --git a/lib/galaxy/tools/actions/metadata.py b/lib/galaxy/tools/actions/metadata.py index f82206a8249..da5ec68af9f 100644 --- a/lib/galaxy/tools/actions/metadata.py +++ b/lib/galaxy/tools/actions/metadata.py @@ -2,7 +2,7 @@ import logging import os from json import dumps -from galaxy.jobs.datasets import DatasetPath +from galaxy.job_execution.datasets import DatasetPath from galaxy.metadata import get_metadata_compute_strategy from galaxy.util.odict import odict from . import ToolAction diff --git a/lib/galaxy/tools/evaluation.py b/lib/galaxy/tools/evaluation.py index 01d25922a9b..9d95987e8cf 100644 --- a/lib/galaxy/tools/evaluation.py +++ b/lib/galaxy/tools/evaluation.py @@ -7,7 +7,7 @@ import tempfile from six import string_types from galaxy import model -from galaxy.jobs.datasets import dataset_path_rewrites +from galaxy.job_execution.datasets import dataset_path_rewrites from galaxy.model.none_like import NoneDataset from galaxy.tools import global_tool_errors from galaxy.tools.parameters import ( diff --git a/packages/job_execution/tests/test_datasets.py b/packages/job_execution/tests/test_datasets.py new file mode 120000 index 00000000000..8b6edf2669d --- /dev/null +++ b/packages/job_execution/tests/test_datasets.py @@ -0,0 +1 @@ +../../../test/unit/jobs/test_datasets.py \ No newline at end of file diff --git a/test/unit/jobs/test_datasets.py b/test/unit/jobs/test_datasets.py index a12dd7ce603..1a920d83bfa 100644 --- a/test/unit/jobs/test_datasets.py +++ b/test/unit/jobs/test_datasets.py @@ -1,4 +1,4 @@ -from galaxy.jobs.datasets import DatasetPath +from galaxy.job_execution.datasets import DatasetPath def test_dataset_path(): diff --git a/test/unit/tools/test_evaluation.py b/test/unit/tools/test_evaluation.py index 93345de7215..37243fbfdbf 100644 --- a/test/unit/tools/test_evaluation.py +++ b/test/unit/tools/test_evaluation.py @@ -2,8 +2,8 @@ import os from unittest import TestCase from xml.etree.ElementTree import XML +from galaxy.job_execution.datasets import DatasetPath from galaxy.jobs import SimpleComputeEnvironment -from galaxy.jobs.datasets import DatasetPath from galaxy.model import ( Dataset, History, diff --git a/test/unit/tools/test_metadata.py b/test/unit/tools/test_metadata.py index 2c7e2e99fd0..e707877fced 100644 --- a/test/unit/tools/test_metadata.py +++ b/test/unit/tools/test_metadata.py @@ -3,7 +3,7 @@ import subprocess import unittest from galaxy import model -from galaxy.jobs.datasets import DatasetPath +from galaxy.job_execution.datasets import DatasetPath from galaxy.metadata import get_metadata_compute_strategy from galaxy.objectstore import ObjectStorePopulator from .. import tools_support diff --git a/test/unit/tools/test_wrappers.py b/test/unit/tools/test_wrappers.py index 4f58df39bef..57248754e06 100644 --- a/test/unit/tools/test_wrappers.py +++ b/test/unit/tools/test_wrappers.py @@ -3,7 +3,7 @@ import tempfile from xml.etree.ElementTree import XML from galaxy.datatypes.metadata import MetadataSpecCollection -from galaxy.jobs.datasets import DatasetPath +from galaxy.job_execution.datasets import DatasetPath from galaxy.tools.parameters.basic import ( DrillDownSelectToolParameter, IntegerToolParameter, From fd7bc2d0da0b0de38db642be11ee8aa826febfa2 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Thu, 6 Jun 2019 09:17:09 -0400 Subject: [PATCH 08/15] Place galaxy.metadata into galaxy-job-execution. --- packages/job_execution/galaxy/metadata | 1 + packages/job_execution/setup.py | 1 + packages/job_execution/tests/test_metadata.py | 1 + 3 files changed, 3 insertions(+) create mode 120000 packages/job_execution/galaxy/metadata create mode 120000 packages/job_execution/tests/test_metadata.py diff --git a/packages/job_execution/galaxy/metadata b/packages/job_execution/galaxy/metadata new file mode 120000 index 00000000000..668ce86f349 --- /dev/null +++ b/packages/job_execution/galaxy/metadata @@ -0,0 +1 @@ +../../../lib/galaxy/metadata \ No newline at end of file diff --git a/packages/job_execution/setup.py b/packages/job_execution/setup.py index 0ea5361f7d5..8b988e66385 100644 --- a/packages/job_execution/setup.py +++ b/packages/job_execution/setup.py @@ -32,6 +32,7 @@ TEST_DIR = 'tests' PACKAGES = [ 'galaxy', 'galaxy.job_execution', + 'galaxy.metadata', ] ENTRY_POINTS = ''' [console_scripts] diff --git a/packages/job_execution/tests/test_metadata.py b/packages/job_execution/tests/test_metadata.py new file mode 120000 index 00000000000..9aa7e9cb83b --- /dev/null +++ b/packages/job_execution/tests/test_metadata.py @@ -0,0 +1 @@ +../../../test/unit/tools/test_metadata.py \ No newline at end of file From b4a1c21456fb820e5012501d8ea65db906a672c2 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Thu, 6 Jun 2019 09:20:47 -0400 Subject: [PATCH 09/15] Enable testing for galaxy-job-execution. --- packages/test.sh | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/test.sh b/packages/test.sh index a597d69ee1a..b9dfa0b8ddb 100755 --- a/packages/test.sh +++ b/packages/test.sh @@ -24,7 +24,7 @@ PACKAGE_DIRS=( ) # containers has no tests, tool_util not yet working 100%, # data has many problems quota, tool shed install database, etc.. -RUN_TESTS=(1 1 1 0 0 0 0 0) +RUN_TESTS=(1 1 1 0 0 0 0 1) for ((i=0; i<${#PACKAGE_DIRS[@]}; i++)); do package_dir=${PACKAGE_DIRS[$i]} From 7b6d06bd4ee6b1526f2bffd26e03279fce2e0f3f Mon Sep 17 00:00:00 2001 From: John Chilton Date: Thu, 6 Jun 2019 09:36:28 -0400 Subject: [PATCH 10/15] Refactor get_metadata_compute_strategy to not depend on ``app``. This allows me to claim galaxy-job-execution doesn't depend on ``app``. --- lib/galaxy/jobs/__init__.py | 2 +- lib/galaxy/metadata/__init__.py | 4 ++-- lib/galaxy/tools/__init__.py | 2 +- lib/galaxy/tools/actions/metadata.py | 2 +- test/unit/tools/test_metadata.py | 2 +- 5 files changed, 6 insertions(+), 6 deletions(-) diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index e69f8cb30d2..8cfafec9291 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -875,7 +875,7 @@ class JobWrapper(HasResourceParameters): self.output_hdas_and_paths = None self.tool_provided_job_metadata = None # Wrapper holding the info required to restore and clean up from files used for setting metadata externally - self.external_output_metadata = get_metadata_compute_strategy(self.app, job.id) + self.external_output_metadata = get_metadata_compute_strategy(self.app.config, job.id) self.job_runner_mapper = JobRunnerMapper(self, queue.dispatcher.url_to_destination, self.app.job_config) self.params = None if job.params: diff --git a/lib/galaxy/metadata/__init__.py b/lib/galaxy/metadata/__init__.py index b095abe60cb..bd16701e6b0 100644 --- a/lib/galaxy/metadata/__init__.py +++ b/lib/galaxy/metadata/__init__.py @@ -20,8 +20,8 @@ log = getLogger(__name__) SET_METADATA_SCRIPT = 'from galaxy_ext.metadata.set_metadata import set_metadata; set_metadata()' -def get_metadata_compute_strategy(app, job_id): - metadata_strategy = app.config.metadata_strategy +def get_metadata_compute_strategy(config, job_id): + metadata_strategy = config.metadata_strategy if metadata_strategy == "legacy": return JobExternalOutputMetadataWrapper(job_id) else: diff --git a/lib/galaxy/tools/__init__.py b/lib/galaxy/tools/__init__.py index 7aa13edf6d9..920dec40203 100755 --- a/lib/galaxy/tools/__init__.py +++ b/lib/galaxy/tools/__init__.py @@ -2394,7 +2394,7 @@ class SetMetadataTool(Tool): job, base_dir='job_work', dir_only=True, obj_dir=True ) for name, dataset in inp_data.items(): - external_metadata = get_metadata_compute_strategy(app, job.id) + external_metadata = get_metadata_compute_strategy(app.config, job.id) sa_session = app.model.context if external_metadata.external_metadata_set_successfully(dataset, name, sa_session, working_directory=working_directory): external_metadata.load_metadata(dataset, name, sa_session, working_directory=working_directory) diff --git a/lib/galaxy/tools/actions/metadata.py b/lib/galaxy/tools/actions/metadata.py index da5ec68af9f..c15439cb706 100644 --- a/lib/galaxy/tools/actions/metadata.py +++ b/lib/galaxy/tools/actions/metadata.py @@ -78,7 +78,7 @@ class SetMetadataToolAction(ToolAction): job_working_dir = app.object_store.get_filename(job, base_dir='job_work', dir_only=True, extra_dir=str(job.id)) datatypes_config = os.path.join(job_working_dir, 'registry.xml') app.datatypes_registry.to_xml_file(path=datatypes_config) - external_metadata_wrapper = get_metadata_compute_strategy(app, job.id) + external_metadata_wrapper = get_metadata_compute_strategy(app.config, job.id) output_datatasets_dict = { dataset_name: dataset, } diff --git a/test/unit/tools/test_metadata.py b/test/unit/tools/test_metadata.py index e707877fced..e86f83360d0 100644 --- a/test/unit/tools/test_metadata.py +++ b/test/unit/tools/test_metadata.py @@ -143,7 +143,7 @@ class MetadataTestCase(unittest.TestCase, tools_support.UsesApp, tools_support.U f.write(contents) def metadata_command(self, output_datasets): - metadata_compute_strategy = get_metadata_compute_strategy(self.app, self.job.id) + metadata_compute_strategy = get_metadata_compute_strategy(self.app.config, self.job.id) self.metadata_compute_strategy = metadata_compute_strategy exec_dir = None From 16b7d6a031821ca8180fb7357414853b44e6085d Mon Sep 17 00:00:00 2001 From: John Chilton Date: Thu, 6 Jun 2019 13:57:41 -0400 Subject: [PATCH 11/15] Move set metadata. --- lib/galaxy/metadata/set_metadata.py | 219 ++++++++++++++++++ lib/galaxy_ext/metadata/_provided_metadata.py | 1 - lib/galaxy_ext/metadata/set_metadata.py | 210 +---------------- packages/job_execution/setup.py | 1 + 4 files changed, 223 insertions(+), 208 deletions(-) create mode 100644 lib/galaxy/metadata/set_metadata.py delete mode 120000 lib/galaxy_ext/metadata/_provided_metadata.py diff --git a/lib/galaxy/metadata/set_metadata.py b/lib/galaxy/metadata/set_metadata.py new file mode 100644 index 00000000000..09fcd1f5403 --- /dev/null +++ b/lib/galaxy/metadata/set_metadata.py @@ -0,0 +1,219 @@ +""" +Execute an external process to set_meta() on a provided list of pickled datasets. + +This was formerly scripts/set_metadata.py and expects these arguments: + + %prog datatypes_conf.xml job_metadata_file metadata_in,metadata_kwds,metadata_out,metadata_results_code,output_filename_override,metadata_override... max_metadata_value_size + +Galaxy should be importable on sys.path and output_filename_override should be +set to the path of the dataset on which metadata is being set +(output_filename_override could previously be left empty and the path would be +constructed automatically). +""" +import json +import logging +import os +import sys + +from six.moves import cPickle +from sqlalchemy.orm import clear_mappers + +import galaxy.model.mapping # need to load this before we unpickle, in order to setup properties assigned by the mappers +from galaxy.model.custom_types import total_size +from galaxy.tool_util.provided_metadata import parse_tool_provided_metadata +from galaxy.util import stringify_dictionary_keys + +logging.basicConfig() +log = logging.getLogger(__name__) + +galaxy.model.Job() # this looks REAL stupid, but it is REQUIRED in order for SA to insert parameters into the classes defined by the mappers --> it appears that instantiating ANY mapper'ed class would suffice here + + +def set_meta_with_tool_provided(dataset_instance, file_dict, set_meta_kwds, datatypes_registry, max_metadata_value_size): + # This method is somewhat odd, in that we set the metadata attributes from tool, + # then call set_meta, then set metadata attributes from tool again. + # This is intentional due to interplay of overwrite kwd, the fact that some metadata + # parameters may rely on the values of others, and that we are accepting the + # values provided by the tool as Truth. + extension = dataset_instance.extension + if extension == "_sniff_": + try: + from galaxy.datatypes import sniff + extension = sniff.handle_uploaded_dataset_file(dataset_instance.dataset.external_filename, datatypes_registry) + # We need to both set the extension so it is available to set_meta + # and record it in the metadata so it can be reloaded on the server + # side and the model updated (see MetadataCollection.{from,to}_JSON_dict) + dataset_instance.extension = extension + # Set special metadata property that will reload this on server side. + setattr(dataset_instance.metadata, "__extension__", extension) + except Exception: + log.exception("Problem sniffing datatype.") + + for metadata_name, metadata_value in file_dict.get('metadata', {}).items(): + setattr(dataset_instance.metadata, metadata_name, metadata_value) + dataset_instance.datatype.set_meta(dataset_instance, **set_meta_kwds) + for metadata_name, metadata_value in file_dict.get('metadata', {}).items(): + setattr(dataset_instance.metadata, metadata_name, metadata_value) + + if max_metadata_value_size: + for k, v in list(dataset_instance.metadata.items()): + if total_size(v) > max_metadata_value_size: + log.info("Key %s too large for metadata, discarding" % k) + dataset_instance.metadata.remove_key(k) + + +def set_metadata(): + if len(sys.argv) == 1: + set_metadata_portable() + else: + set_metadata_legacy() + + +def set_metadata_portable(): + import galaxy.model + galaxy.model.metadata.MetadataTempFile.tmp_dir = tool_job_working_directory = os.path.abspath(os.getcwd()) + + metadata_params_path = os.path.join("metadata", "params.json") + try: + with open(metadata_params_path, "r") as f: + metadata_params = json.load(f) + except IOError: + raise Exception("Failed to find metadata/params.json from cwd [%s]" % tool_job_working_directory) + datatypes_config = metadata_params["datatypes_config"] + job_metadata = metadata_params["job_metadata"] + max_metadata_value_size = metadata_params.get("max_metadata_value_size") or 0 + outputs = metadata_params["outputs"] + + datatypes_registry = validate_and_load_datatypes_config(datatypes_config) + tool_provided_metadata = load_job_metadata(job_metadata) + + def set_meta(new_dataset_instance, file_dict): + set_meta_with_tool_provided(new_dataset_instance, file_dict, set_meta_kwds, datatypes_registry, max_metadata_value_size) + + for output_name, output_dict in outputs.items(): + filename_in = os.path.join("metadata/metadata_in_%s" % output_name) + filename_kwds = os.path.join("metadata/metadata_kwds_%s" % output_name) + filename_out = os.path.join("metadata/metadata_out_%s" % output_name) + filename_results_code = os.path.join("metadata/metadata_results_%s" % output_name) + override_metadata = os.path.join("metadata/metadata_override_%s" % output_name) + dataset_filename_override = output_dict["filename_override"] + + # Same block as below... + set_meta_kwds = stringify_dictionary_keys(json.load(open(filename_kwds))) # load kwds; need to ensure our keywords are not unicode + try: + dataset = cPickle.load(open(filename_in, 'rb')) # load DatasetInstance + dataset.dataset.external_filename = dataset_filename_override + store_by = metadata_params.get("object_store_store_by", "id") + extra_files_dir_name = "dataset_%s_files" % getattr(dataset.dataset, store_by) + files_path = os.path.abspath(os.path.join(tool_job_working_directory, extra_files_dir_name)) + dataset.dataset.external_extra_files_path = files_path + file_dict = tool_provided_metadata.get_dataset_meta(output_name, dataset.dataset.id) + if 'ext' in file_dict: + dataset.extension = file_dict['ext'] + # Metadata FileParameter types may not be writable on a cluster node, and are therefore temporarily substituted with MetadataTempFiles + override_metadata = json.load(open(override_metadata)) + for metadata_name, metadata_file_override in override_metadata: + if galaxy.datatypes.metadata.MetadataTempFile.is_JSONified_value(metadata_file_override): + metadata_file_override = galaxy.datatypes.metadata.MetadataTempFile.from_JSON(metadata_file_override) + setattr(dataset.metadata, metadata_name, metadata_file_override) + set_meta(dataset, file_dict) + dataset.metadata.to_JSON_dict(filename_out) # write out results of set_meta + json.dump((True, 'Metadata has been set successfully'), open(filename_results_code, 'wt+')) # setting metadata has succeeded + except Exception as e: + json.dump((False, str(e)), open(filename_results_code, 'wt+')) # setting metadata has failed somehow + + write_job_metadata(tool_job_working_directory, job_metadata, set_meta, tool_provided_metadata) + + +def set_metadata_legacy(): + import galaxy.model + galaxy.model.metadata.MetadataTempFile.tmp_dir = tool_job_working_directory = os.path.abspath(os.getcwd()) + + # This is ugly, but to transition from existing jobs without this parameter + # to ones with, smoothly, it has to be the last optional parameter and we + # have to sniff it. + try: + max_metadata_value_size = int(sys.argv[-1]) + sys.argv = sys.argv[:-1] + except ValueError: + max_metadata_value_size = 0 + # max_metadata_value_size is unspecified and should be 0 + + # Set up datatypes registry + datatypes_config = sys.argv.pop(1) + datatypes_registry = validate_and_load_datatypes_config(datatypes_config) + + job_metadata = sys.argv.pop(1) + tool_provided_metadata = load_job_metadata(job_metadata) + + def set_meta(new_dataset_instance, file_dict): + set_meta_with_tool_provided(new_dataset_instance, file_dict, set_meta_kwds, datatypes_registry, max_metadata_value_size) + + for filenames in sys.argv[1:]: + fields = filenames.split(',') + filename_in = fields.pop(0) + filename_kwds = fields.pop(0) + filename_out = fields.pop(0) + filename_results_code = fields.pop(0) + dataset_filename_override = fields.pop(0) + override_metadata = fields.pop(0) + set_meta_kwds = stringify_dictionary_keys(json.load(open(filename_kwds))) # load kwds; need to ensure our keywords are not unicode + try: + dataset = cPickle.load(open(filename_in, 'rb')) # load DatasetInstance + dataset.dataset.external_filename = dataset_filename_override + files_path = os.path.abspath(os.path.join(tool_job_working_directory, "dataset_%s_files" % (dataset.dataset.id))) + dataset.dataset.external_extra_files_path = files_path + file_dict = tool_provided_metadata.get_dataset_meta(None, dataset.dataset.id) + if 'ext' in file_dict: + dataset.extension = file_dict['ext'] + # Metadata FileParameter types may not be writable on a cluster node, and are therefore temporarily substituted with MetadataTempFiles + override_metadata = json.load(open(override_metadata)) + for metadata_name, metadata_file_override in override_metadata: + if galaxy.datatypes.metadata.MetadataTempFile.is_JSONified_value(metadata_file_override): + metadata_file_override = galaxy.datatypes.metadata.MetadataTempFile.from_JSON(metadata_file_override) + setattr(dataset.metadata, metadata_name, metadata_file_override) + set_meta(dataset, file_dict) + dataset.metadata.to_JSON_dict(filename_out) # write out results of set_meta + json.dump((True, 'Metadata has been set successfully'), open(filename_results_code, 'wt+')) # setting metadata has succeeded + except Exception as e: + json.dump((False, str(e)), open(filename_results_code, 'wt+')) # setting metadata has failed somehow + + write_job_metadata(tool_job_working_directory, job_metadata, set_meta, tool_provided_metadata) + + +def validate_and_load_datatypes_config(datatypes_config): + galaxy_root = os.path.abspath(os.path.join(os.path.dirname(__file__), os.pardir, os.pardir, os.pardir)) + + if not os.path.exists(datatypes_config): + # Hack for Pulsar on usegalaxy.org, drop ASAP. + datatypes_config = "configs/registry.xml" + + if not os.path.exists(datatypes_config): + print("Metadata setting failed because registry.xml [%s] could not be found. You may retry setting metadata." % datatypes_config) + sys.exit(1) + import galaxy.datatypes.registry + datatypes_registry = galaxy.datatypes.registry.Registry() + datatypes_registry.load_datatypes(root_dir=galaxy_root, config=datatypes_config) + galaxy.model.set_datatypes_registry(datatypes_registry) + return datatypes_registry + + +def load_job_metadata(job_metadata): + return parse_tool_provided_metadata(job_metadata) + + +def write_job_metadata(tool_job_working_directory, job_metadata, set_meta, tool_provided_metadata): + for i, file_dict in enumerate(tool_provided_metadata.get_new_datasets_for_metadata_collection(), start=1): + filename = file_dict["filename"] + new_dataset_filename = os.path.join(tool_job_working_directory, "working", filename) + new_dataset = galaxy.model.Dataset(id=-i, external_filename=new_dataset_filename) + extra_files = file_dict.get('extra_files', None) + if extra_files is not None: + new_dataset._extra_files_path = os.path.join(tool_job_working_directory, "working", extra_files) + new_dataset.state = new_dataset.states.OK + new_dataset_instance = galaxy.model.HistoryDatasetAssociation(id=-i, dataset=new_dataset, extension=file_dict.get('ext', 'data')) + set_meta(new_dataset_instance, file_dict) + file_dict['metadata'] = json.loads(new_dataset_instance.metadata.to_JSON_dict()) # storing metadata in external form, need to turn back into dict, then later jsonify + + tool_provided_metadata.rewrite() + clear_mappers() diff --git a/lib/galaxy_ext/metadata/_provided_metadata.py b/lib/galaxy_ext/metadata/_provided_metadata.py deleted file mode 120000 index b44719c00dc..00000000000 --- a/lib/galaxy_ext/metadata/_provided_metadata.py +++ /dev/null @@ -1 +0,0 @@ -../../galaxy/tool_util/provided_metadata.py \ No newline at end of file diff --git a/lib/galaxy_ext/metadata/set_metadata.py b/lib/galaxy_ext/metadata/set_metadata.py index 9121bd8a4e6..380e3de3336 100644 --- a/lib/galaxy_ext/metadata/set_metadata.py +++ b/lib/galaxy_ext/metadata/set_metadata.py @@ -10,217 +10,13 @@ set to the path of the dataset on which metadata is being set (output_filename_override could previously be left empty and the path would be constructed automatically). """ -import json -import logging + import os import sys # insert *this* galaxy before all others on sys.path sys.path.insert(1, os.path.abspath(os.path.join(os.path.dirname(__file__), os.pardir, os.pardir))) -from six.moves import cPickle -from sqlalchemy.orm import clear_mappers +from galaxy.metadata.set_metadata import set_metadata -import galaxy.model.mapping # need to load this before we unpickle, in order to setup properties assigned by the mappers -from galaxy.model.custom_types import total_size -from galaxy.util import stringify_dictionary_keys -from ._provided_metadata import parse_tool_provided_metadata - -# ensure supported version -assert sys.version_info[:2] >= (2, 7), 'Python version must be at least 2.7, this is: %s' % sys.version - -logging.basicConfig() -log = logging.getLogger(__name__) - -galaxy.model.Job() # this looks REAL stupid, but it is REQUIRED in order for SA to insert parameters into the classes defined by the mappers --> it appears that instantiating ANY mapper'ed class would suffice here - - -def set_meta_with_tool_provided(dataset_instance, file_dict, set_meta_kwds, datatypes_registry, max_metadata_value_size): - # This method is somewhat odd, in that we set the metadata attributes from tool, - # then call set_meta, then set metadata attributes from tool again. - # This is intentional due to interplay of overwrite kwd, the fact that some metadata - # parameters may rely on the values of others, and that we are accepting the - # values provided by the tool as Truth. - extension = dataset_instance.extension - if extension == "_sniff_": - try: - from galaxy.datatypes import sniff - extension = sniff.handle_uploaded_dataset_file(dataset_instance.dataset.external_filename, datatypes_registry) - # We need to both set the extension so it is available to set_meta - # and record it in the metadata so it can be reloaded on the server - # side and the model updated (see MetadataCollection.{from,to}_JSON_dict) - dataset_instance.extension = extension - # Set special metadata property that will reload this on server side. - setattr(dataset_instance.metadata, "__extension__", extension) - except Exception: - log.exception("Problem sniffing datatype.") - - for metadata_name, metadata_value in file_dict.get('metadata', {}).items(): - setattr(dataset_instance.metadata, metadata_name, metadata_value) - dataset_instance.datatype.set_meta(dataset_instance, **set_meta_kwds) - for metadata_name, metadata_value in file_dict.get('metadata', {}).items(): - setattr(dataset_instance.metadata, metadata_name, metadata_value) - - if max_metadata_value_size: - for k, v in list(dataset_instance.metadata.items()): - if total_size(v) > max_metadata_value_size: - log.info("Key %s too large for metadata, discarding" % k) - dataset_instance.metadata.remove_key(k) - - -def set_metadata(): - if len(sys.argv) == 1: - set_metadata_portable() - else: - set_metadata_legacy() - - -def set_metadata_portable(): - import galaxy.model - tool_job_working_directory = os.path.abspath(os.getcwd()) - galaxy.model.metadata.MetadataTempFile.tmp_dir = os.path.join(tool_job_working_directory, "metadata") - - metadata_params_path = os.path.join("metadata", "params.json") - try: - with open(metadata_params_path, "r") as f: - metadata_params = json.load(f) - except IOError: - raise Exception("Failed to find metadata/params.json from cwd [%s]" % tool_job_working_directory) - datatypes_config = metadata_params["datatypes_config"] - job_metadata = metadata_params["job_metadata"] - max_metadata_value_size = metadata_params.get("max_metadata_value_size") or 0 - outputs = metadata_params["outputs"] - - datatypes_registry = validate_and_load_datatypes_config(datatypes_config) - tool_provided_metadata = load_job_metadata(job_metadata) - - def set_meta(new_dataset_instance, file_dict): - set_meta_with_tool_provided(new_dataset_instance, file_dict, set_meta_kwds, datatypes_registry, max_metadata_value_size) - - for output_name, output_dict in outputs.items(): - filename_in = os.path.join("metadata/metadata_in_%s" % output_name) - filename_kwds = os.path.join("metadata/metadata_kwds_%s" % output_name) - filename_out = os.path.join("metadata/metadata_out_%s" % output_name) - filename_results_code = os.path.join("metadata/metadata_results_%s" % output_name) - override_metadata = os.path.join("metadata/metadata_override_%s" % output_name) - dataset_filename_override = output_dict["filename_override"] - - # Same block as below... - set_meta_kwds = stringify_dictionary_keys(json.load(open(filename_kwds))) # load kwds; need to ensure our keywords are not unicode - try: - dataset = cPickle.load(open(filename_in, 'rb')) # load DatasetInstance - dataset.dataset.external_filename = dataset_filename_override - store_by = metadata_params.get("object_store_store_by", "id") - extra_files_dir_name = "dataset_%s_files" % getattr(dataset.dataset, store_by) - files_path = os.path.abspath(os.path.join(tool_job_working_directory, extra_files_dir_name)) - dataset.dataset.external_extra_files_path = files_path - file_dict = tool_provided_metadata.get_dataset_meta(output_name, dataset.dataset.id) - if 'ext' in file_dict: - dataset.extension = file_dict['ext'] - # Metadata FileParameter types may not be writable on a cluster node, and are therefore temporarily substituted with MetadataTempFiles - override_metadata = json.load(open(override_metadata)) - for metadata_name, metadata_file_override in override_metadata: - if galaxy.datatypes.metadata.MetadataTempFile.is_JSONified_value(metadata_file_override): - metadata_file_override = galaxy.datatypes.metadata.MetadataTempFile.from_JSON(metadata_file_override) - setattr(dataset.metadata, metadata_name, metadata_file_override) - set_meta(dataset, file_dict) - dataset.metadata.to_JSON_dict(filename_out) # write out results of set_meta - json.dump((True, 'Metadata has been set successfully'), open(filename_results_code, 'wt+')) # setting metadata has succeeded - except Exception as e: - json.dump((False, str(e)), open(filename_results_code, 'wt+')) # setting metadata has failed somehow - - write_job_metadata(tool_job_working_directory, job_metadata, set_meta, tool_provided_metadata) - - -def set_metadata_legacy(): - import galaxy.model - galaxy.model.metadata.MetadataTempFile.tmp_dir = tool_job_working_directory = os.path.abspath(os.getcwd()) - - # This is ugly, but to transition from existing jobs without this parameter - # to ones with, smoothly, it has to be the last optional parameter and we - # have to sniff it. - try: - max_metadata_value_size = int(sys.argv[-1]) - sys.argv = sys.argv[:-1] - except ValueError: - max_metadata_value_size = 0 - # max_metadata_value_size is unspecified and should be 0 - - # Set up datatypes registry - datatypes_config = sys.argv.pop(1) - datatypes_registry = validate_and_load_datatypes_config(datatypes_config) - - job_metadata = sys.argv.pop(1) - tool_provided_metadata = load_job_metadata(job_metadata) - - def set_meta(new_dataset_instance, file_dict): - set_meta_with_tool_provided(new_dataset_instance, file_dict, set_meta_kwds, datatypes_registry, max_metadata_value_size) - - for filenames in sys.argv[1:]: - fields = filenames.split(',') - filename_in = fields.pop(0) - filename_kwds = fields.pop(0) - filename_out = fields.pop(0) - filename_results_code = fields.pop(0) - dataset_filename_override = fields.pop(0) - override_metadata = fields.pop(0) - set_meta_kwds = stringify_dictionary_keys(json.load(open(filename_kwds))) # load kwds; need to ensure our keywords are not unicode - try: - dataset = cPickle.load(open(filename_in, 'rb')) # load DatasetInstance - dataset.dataset.external_filename = dataset_filename_override - files_path = os.path.abspath(os.path.join(tool_job_working_directory, "dataset_%s_files" % (dataset.dataset.id))) - dataset.dataset.external_extra_files_path = files_path - file_dict = tool_provided_metadata.get_dataset_meta(None, dataset.dataset.id) - if 'ext' in file_dict: - dataset.extension = file_dict['ext'] - # Metadata FileParameter types may not be writable on a cluster node, and are therefore temporarily substituted with MetadataTempFiles - override_metadata = json.load(open(override_metadata)) - for metadata_name, metadata_file_override in override_metadata: - if galaxy.datatypes.metadata.MetadataTempFile.is_JSONified_value(metadata_file_override): - metadata_file_override = galaxy.datatypes.metadata.MetadataTempFile.from_JSON(metadata_file_override) - setattr(dataset.metadata, metadata_name, metadata_file_override) - set_meta(dataset, file_dict) - dataset.metadata.to_JSON_dict(filename_out) # write out results of set_meta - json.dump((True, 'Metadata has been set successfully'), open(filename_results_code, 'wt+')) # setting metadata has succeeded - except Exception as e: - json.dump((False, str(e)), open(filename_results_code, 'wt+')) # setting metadata has failed somehow - - write_job_metadata(tool_job_working_directory, job_metadata, set_meta, tool_provided_metadata) - - -def validate_and_load_datatypes_config(datatypes_config): - galaxy_root = os.path.abspath(os.path.join(os.path.dirname(__file__), os.pardir, os.pardir, os.pardir)) - - if not os.path.exists(datatypes_config): - # Hack for Pulsar on usegalaxy.org, drop ASAP. - datatypes_config = "configs/registry.xml" - - if not os.path.exists(datatypes_config): - print("Metadata setting failed because registry.xml [%s] could not be found. You may retry setting metadata." % datatypes_config) - sys.exit(1) - import galaxy.datatypes.registry - datatypes_registry = galaxy.datatypes.registry.Registry() - datatypes_registry.load_datatypes(root_dir=galaxy_root, config=datatypes_config) - galaxy.model.set_datatypes_registry(datatypes_registry) - return datatypes_registry - - -def load_job_metadata(job_metadata): - return parse_tool_provided_metadata(job_metadata) - - -def write_job_metadata(tool_job_working_directory, job_metadata, set_meta, tool_provided_metadata): - for i, file_dict in enumerate(tool_provided_metadata.get_new_datasets_for_metadata_collection(), start=1): - filename = file_dict["filename"] - new_dataset_filename = os.path.join(tool_job_working_directory, "working", filename) - new_dataset = galaxy.model.Dataset(id=-i, external_filename=new_dataset_filename) - extra_files = file_dict.get('extra_files', None) - if extra_files is not None: - new_dataset._extra_files_path = os.path.join(tool_job_working_directory, "working", extra_files) - new_dataset.state = new_dataset.states.OK - new_dataset_instance = galaxy.model.HistoryDatasetAssociation(id=-i, dataset=new_dataset, extension=file_dict.get('ext', 'data')) - set_meta(new_dataset_instance, file_dict) - file_dict['metadata'] = json.loads(new_dataset_instance.metadata.to_JSON_dict()) # storing metadata in external form, need to turn back into dict, then later jsonify - - tool_provided_metadata.rewrite() - clear_mappers() +__all__ = ('set_metadata', ) diff --git a/packages/job_execution/setup.py b/packages/job_execution/setup.py index 8b988e66385..e0d5cb5c879 100644 --- a/packages/job_execution/setup.py +++ b/packages/job_execution/setup.py @@ -36,6 +36,7 @@ PACKAGES = [ ] ENTRY_POINTS = ''' [console_scripts] + galaxy-set-metadata=galaxy.metadata.set_metadata:set_metadata ''' PACKAGE_DATA = { # Be sure to update MANIFEST.in for source dist. From d4e1a488bfbf1b463f4f6e761462cc4f1ef2a9bd Mon Sep 17 00:00:00 2001 From: John Chilton Date: Thu, 6 Jun 2019 20:41:02 -0400 Subject: [PATCH 12/15] Allow using galaxy-set-metadata to set metadata. This option, combined with new "portable" metadata generation, greatly simplifies remote metadata collection with Pulsar. --- lib/galaxy/jobs/__init__.py | 4 ++++ lib/galaxy/jobs/command_factory.py | 1 + lib/galaxy/jobs/runners/pulsar.py | 2 +- lib/galaxy/metadata/__init__.py | 16 ++++++++++------ test/unit/jobs/test_command_factory.py | 1 + test/unit/jobs/test_runner_local.py | 1 + 6 files changed, 18 insertions(+), 7 deletions(-) diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index 8cfafec9291..c826f2de2e4 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -921,6 +921,10 @@ class JobWrapper(HasResourceParameters): def requires_containerization(self): return util.asbool(self.get_destination_configuration("require_container", "False")) + @property + def use_metadata_binary(self): + return util.asbool(self.get_destination_configuration('use_metadata_binary', "False")) + def can_split(self): # Should the job handler split this job up? return self.app.config.use_tasked_jobs and self.tool.parallelism diff --git a/lib/galaxy/jobs/command_factory.py b/lib/galaxy/jobs/command_factory.py index db5cff813ce..adccf42b427 100644 --- a/lib/galaxy/jobs/command_factory.py +++ b/lib/galaxy/jobs/command_factory.py @@ -204,6 +204,7 @@ def __handle_metadata(commands_builder, job_wrapper, runner, remote_command_para datatypes_config=datatypes_config, compute_tmp_dir=compute_tmp_dir, resolve_metadata_dependencies=resolve_metadata_dependencies, + use_bin=job_wrapper.use_metadata_binary, kwds={'overwrite': False} ) or '' metadata_command = metadata_command.strip() diff --git a/lib/galaxy/jobs/runners/pulsar.py b/lib/galaxy/jobs/runners/pulsar.py index eda05554b6f..16829353477 100644 --- a/lib/galaxy/jobs/runners/pulsar.py +++ b/lib/galaxy/jobs/runners/pulsar.py @@ -699,7 +699,7 @@ class PulsarJobRunner(AsynchronousJobRunner): def __build_metadata_configuration(self, client, job_wrapper, remote_metadata, remote_job_config): metadata_kwds = {} - if remote_metadata: + if remote_metadata and not job_wrapper.use_metadata_binary: remote_system_properties = remote_job_config.get("system_properties", {}) remote_galaxy_home = remote_system_properties.get("galaxy_home", None) if not remote_galaxy_home: diff --git a/lib/galaxy/metadata/__init__.py b/lib/galaxy/metadata/__init__.py index bd16701e6b0..7b78d39072c 100644 --- a/lib/galaxy/metadata/__init__.py +++ b/lib/galaxy/metadata/__init__.py @@ -46,7 +46,7 @@ class MetadataCollectionStrategy(object): @abc.abstractmethod def setup_external_metadata(self, datasets_dict, sa_session, exec_dir=None, tmp_dir=None, dataset_files_path=None, - output_fnames=None, config_root=None, + output_fnames=None, config_root=None, use_bin=False, config_file=None, datatypes_config=None, job_metadata=None, compute_tmp_dir=None, include_command=True, max_metadata_value_size=0, @@ -99,7 +99,7 @@ class PortableDirectoryMetadataGenerator(MetadataCollectionStrategy): def setup_external_metadata(self, datasets_dict, sa_session, exec_dir=None, tmp_dir=None, dataset_files_path=None, - output_fnames=None, config_root=None, + output_fnames=None, config_root=None, use_bin=False, config_file=None, datatypes_config=None, job_metadata=None, compute_tmp_dir=None, include_command=True, max_metadata_value_size=0, @@ -147,9 +147,12 @@ class PortableDirectoryMetadataGenerator(MetadataCollectionStrategy): if include_command: # return command required to build script_path = os.path.join(metadata_dir, "set.py") - with open(script_path, "w") as f: - f.write(SET_METADATA_SCRIPT) - return 'python "metadata/set.py"' + if use_bin: + return "galaxy-set-metadata" + else: + with open(script_path, "w") as f: + f.write(SET_METADATA_SCRIPT) + return 'python "metadata/set.py"' else: # return args to galaxy_ext.metadata.set_metadata required to build return '' @@ -214,7 +217,7 @@ class JobExternalOutputMetadataWrapper(MetadataCollectionStrategy): def setup_external_metadata(self, datasets_dict, sa_session, exec_dir=None, tmp_dir=None, dataset_files_path=None, - output_fnames=None, config_root=None, + output_fnames=None, config_root=None, use_bin=False, config_file=None, datatypes_config=None, job_metadata=None, compute_tmp_dir=None, include_command=True, max_metadata_value_size=0, @@ -296,6 +299,7 @@ class JobExternalOutputMetadataWrapper(MetadataCollectionStrategy): job_metadata, " ".join(map(__metadata_files_list_to_cmd_line, metadata_files_list)), max_metadata_value_size) + assert not use_bin if include_command: # return command required to build fd, fp = tempfile.mkstemp(suffix='.py', dir=tmp_dir, prefix="set_metadata_") diff --git a/test/unit/jobs/test_command_factory.py b/test/unit/jobs/test_command_factory.py index 28f16c58f7b..f997190badf 100644 --- a/test/unit/jobs/test_command_factory.py +++ b/test/unit/jobs/test_command_factory.py @@ -178,6 +178,7 @@ class MockJobWrapper(object): ) ) self.shell = "/bin/sh" + self.use_metadata_binary = False def get_command_line(self): return self.command_line diff --git a/test/unit/jobs/test_runner_local.py b/test/unit/jobs/test_runner_local.py index 610da86a66f..7675b96a577 100644 --- a/test/unit/jobs/test_runner_local.py +++ b/test/unit/jobs/test_runner_local.py @@ -144,6 +144,7 @@ class MockJobWrapper(object): self.shell = "/bin/bash" self.cleanup_job = "never" self.tmp_dir_creation_statement = "" + self.use_metadata_binary = False # Cruft for setting metadata externally, axe at some point. self.external_output_metadata = bunch.Bunch( From 729c6108fe96bd67aaffbd0b15a604b075681018 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Mon, 17 Jun 2019 15:44:26 -0400 Subject: [PATCH 13/15] Update packages/job_execution/README.rst Co-Authored-By: Marius van den Beek --- packages/job_execution/README.rst | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/job_execution/README.rst b/packages/job_execution/README.rst index 897422a8aa3..27e1508c806 100644 --- a/packages/job_execution/README.rst +++ b/packages/job_execution/README.rst @@ -1,5 +1,5 @@ -.. image:: https://badge.fury.io/py/galaxy-util.svg +.. image:: https://badge.fury.io/py/galaxy-job-execution.svg :target: https://pypi.python.org/pypi/galaxy-util/ From 055ecac4c1fdf1271811d7ff90b11d8663968598 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Mon, 17 Jun 2019 15:44:33 -0400 Subject: [PATCH 14/15] Update packages/job_execution/requirements.txt Co-Authored-By: Marius van den Beek --- packages/job_execution/requirements.txt | 1 - 1 file changed, 1 deletion(-) diff --git a/packages/job_execution/requirements.txt b/packages/job_execution/requirements.txt index 77de615ea9c..7dcf8f4def7 100644 --- a/packages/job_execution/requirements.txt +++ b/packages/job_execution/requirements.txt @@ -1,2 +1 @@ galaxy-data - From 5dde802db0cd408bf9009e8a6bb2be522e103aa7 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Mon, 17 Jun 2019 15:44:46 -0400 Subject: [PATCH 15/15] Update packages/job_execution/README.rst Co-Authored-By: Marius van den Beek --- packages/job_execution/README.rst | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/job_execution/README.rst b/packages/job_execution/README.rst index 27e1508c806..3a7574eab70 100644 --- a/packages/job_execution/README.rst +++ b/packages/job_execution/README.rst @@ -1,6 +1,6 @@ .. image:: https://badge.fury.io/py/galaxy-job-execution.svg - :target: https://pypi.python.org/pypi/galaxy-util/ + :target: https://pypi.python.org/pypi/galaxy-job-execution/ Overview