From 7d0ec29fa06ae6f04c0ae6fc193a390bc11cddcd Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Thu, 4 Feb 2021 12:48:02 +0100 Subject: [PATCH] Drop legacy metadata mode And always serialize input dataset metadata using the model exports. This means we don't have to pickle anymore, which breaks with weakrefs used by sqlalchemy-json's Mutable class. --- doc/source/admin/galaxy_options.rst | 14 +- lib/galaxy/config/sample/galaxy.yml.sample | 6 +- lib/galaxy/metadata/__init__.py | 239 ++------------------ lib/galaxy/metadata/set_metadata.py | 81 +------ lib/galaxy/webapps/galaxy/config_schema.yml | 2 +- test/unit/tools/test_metadata.py | 12 - 6 files changed, 33 insertions(+), 321 deletions(-) diff --git a/doc/source/admin/galaxy_options.rst b/doc/source/admin/galaxy_options.rst index 6469d27017f..86479656e0d 100644 --- a/doc/source/admin/galaxy_options.rst +++ b/doc/source/admin/galaxy_options.rst @@ -3851,13 +3851,13 @@ ~~~~~~~~~~~~~~~~~~~~~ :Description: - Determines how metadata will be set. Valid values are `directory`, - `extended` and `legacy`. In extended mode jobs will decide if a - tool run failed, the object stores configuration is serialized and - made available to the job and is used for writing output datasets - to the object store as part of the job and dynamic output - discovery (e.g. discovered datasets , - unpopulated collections, etc) happens as part of the job. + Determines how metadata will be set. Valid values are `directory` + and `extended`. In extended mode jobs will decide if a tool run + failed, the object stores configuration is serialized and made + available to the job and is used for writing output datasets to + the object store as part of the job and dynamic output discovery + (e.g. discovered datasets , unpopulated + collections, etc) happens as part of the job. :Default: ``directory`` :Type: str diff --git a/lib/galaxy/config/sample/galaxy.yml.sample b/lib/galaxy/config/sample/galaxy.yml.sample index 7fccb510097..9bf13d71da0 100644 --- a/lib/galaxy/config/sample/galaxy.yml.sample +++ b/lib/galaxy/config/sample/galaxy.yml.sample @@ -1902,9 +1902,9 @@ galaxy: # database. #enable_job_recovery: true - # Determines how metadata will be set. Valid values are `directory`, - # `extended` and `legacy`. In extended mode jobs will decide if a tool - # run failed, the object stores configuration is serialized and made + # Determines how metadata will be set. Valid values are `directory` + # and `extended`. In extended mode jobs will decide if a tool run + # failed, the object stores configuration is serialized and made # available to the job and is used for writing output datasets to the # object store as part of the job and dynamic output discovery (e.g. # discovered datasets , unpopulated collections, diff --git a/lib/galaxy/metadata/__init__.py b/lib/galaxy/metadata/__init__.py index bff6ac4baf1..f072015bd55 100644 --- a/lib/galaxy/metadata/__init__.py +++ b/lib/galaxy/metadata/__init__.py @@ -3,17 +3,14 @@ import abc import json import os -import pickle import shutil -import tempfile from logging import getLogger -from os.path import abspath import galaxy.model from galaxy.model import store from galaxy.model.metadata import FileParameter, MetadataTempFile from galaxy.model.store import DirectoryModelExportStore -from galaxy.util import in_directory, safe_makedirs +from galaxy.util import safe_makedirs log = getLogger(__name__) @@ -23,7 +20,7 @@ SET_METADATA_SCRIPT = 'from galaxy_ext.metadata.set_metadata import set_metadata def get_metadata_compute_strategy(config, job_id, metadata_strategy_override=None, tool_id=None): metadata_strategy = metadata_strategy_override or config.metadata_strategy if metadata_strategy == "legacy": - return JobExternalOutputMetadataWrapper(job_id) + raise Exception('legacy metadata_strategy has been removed') elif metadata_strategy == "extended" and tool_id != "__SET_METADATA__": return ExtendedDirectoryMetadataGenerator(job_id) else: @@ -130,17 +127,15 @@ class PortableDirectoryMetadataGenerator(MetadataCollectionStrategy): outputs = {} output_collections = {} - real_metadata_object = self.write_object_store_conf for name, dataset in datasets_dict.items(): assert name is not None assert name not in outputs - key = name def _metadata_path(what): return os.path.join(metadata_dir, f"metadata_{what}_{key}") - _initialize_metadata_inputs(dataset, _metadata_path, tmp_dir, kwds, real_metadata_object=real_metadata_object) + _initialize_metadata_inputs(dataset, _metadata_path, tmp_dir, kwds, real_metadata_object=self.write_object_store_conf) outputs[name] = { "filename_override": _get_filename_override(output_fnames, dataset.file_name), @@ -159,19 +154,19 @@ class PortableDirectoryMetadataGenerator(MetadataCollectionStrategy): "outputs": outputs, } + # export model objects and object store configuration for extended metadata also. + export_directory = os.path.join(metadata_dir, "outputs_new") + with DirectoryModelExportStore(export_directory, for_edit=True, serialize_dataset_objects=True) as export_store: + for dataset in datasets_dict.values(): + export_store.add_dataset(dataset) + + for name, dataset_collection in out_collections.items(): + export_store.add_dataset_collection(dataset_collection) + output_collections[name] = { + 'id': dataset_collection.id, + } + if self.write_object_store_conf: - # export model objects and object store configuration for extended metadata also. - export_directory = os.path.join(metadata_dir, "outputs_new") - with DirectoryModelExportStore(export_directory, for_edit=True, serialize_dataset_objects=True) as export_store: - for dataset in datasets_dict.values(): - export_store.add_dataset(dataset) - - for name, dataset_collection in out_collections.items(): - export_store.add_dataset_collection(dataset_collection) - output_collections[name] = { - 'id': dataset_collection.id, - } - with open(os.path.join(metadata_dir, "object_store_conf.json"), "w") as f: json.dump(object_store_conf, f) @@ -241,195 +236,12 @@ class ExtendedDirectoryMetadataGenerator(PortableDirectoryMetadataGenerator): return dataset -class JobExternalOutputMetadataWrapper(MetadataCollectionStrategy): - """ - Class with methods allowing set_meta() to be called externally to the - Galaxy head. - This class allows access to external metadata filenames for all outputs - associated with a job. - We will use JSON as the medium of exchange of information, except for the - DatasetInstance object which will use pickle (in the future this could be - JSONified as well) - """ - portable = False - - def __init__(self, job_id): - self.job_id = job_id - - def _get_output_filenames_by_dataset(self, dataset, sa_session): - if isinstance(dataset, galaxy.model.HistoryDatasetAssociation): - return sa_session.query(galaxy.model.JobExternalOutputMetadata) \ - .filter_by(job_id=self.job_id, - history_dataset_association_id=dataset.id, - is_valid=True) \ - .first() # there should only be one or None - elif isinstance(dataset, galaxy.model.LibraryDatasetDatasetAssociation): - return sa_session.query(galaxy.model.JobExternalOutputMetadata) \ - .filter_by(job_id=self.job_id, - library_dataset_dataset_association_id=dataset.id, - is_valid=True) \ - .first() # there should only be one or None - return None - - def _get_dataset_metadata_key(self, dataset): - # Set meta can be called on library items and history items, - # need to make different keys for them, since ids can overlap - return "%s_%d" % (dataset.__class__.__name__, dataset.id) - - def invalidate_external_metadata(self, datasets, sa_session): - for dataset in datasets: - jeom = self._get_output_filenames_by_dataset(dataset, sa_session) - # shouldn't be more than one valid, but you never know - while jeom: - jeom.is_valid = False - sa_session.add(jeom) - sa_session.flush() - jeom = self._get_output_filenames_by_dataset(dataset, sa_session) - - def setup_external_metadata(self, datasets_dict, out_collections, sa_session, exec_dir=None, - tmp_dir=None, dataset_files_path=None, - output_fnames=None, config_root=None, use_bin=False, - config_file=None, datatypes_config=None, - job_metadata=None, provided_metadata_style=None, compute_tmp_dir=None, - include_command=True, max_metadata_value_size=0, - validate_outputs=False, - object_store_conf=None, tool=None, job=None, - kwds=None): - kwds = kwds or {} - if not job: - job = sa_session.query(galaxy.model.Job).get(self.job_id) - tmp_dir = _init_tmp_dir(tmp_dir) - _assert_datatypes_config(datatypes_config) - - # path is calculated for Galaxy, may be different on compute - rewrite - # for the compute server. - def metadata_path_on_compute(path): - compute_path = path - if compute_tmp_dir and tmp_dir and in_directory(path, tmp_dir): - path_relative = os.path.relpath(path, tmp_dir) - compute_path = os.path.join(compute_tmp_dir, path_relative) - return compute_path - - # fill in metadata_files_dict and return the command with args required to set metadata - def __metadata_files_list_to_cmd_line(metadata_files): - line = '"{},{},{},{},{},{}"'.format( - metadata_path_on_compute(metadata_files.filename_in), - metadata_path_on_compute(metadata_files.filename_kwds), - metadata_path_on_compute(metadata_files.filename_out), - metadata_path_on_compute(metadata_files.filename_results_code), - _get_filename_override(output_fnames, metadata_files.dataset.file_name), - metadata_path_on_compute(metadata_files.filename_override_metadata), - ) - return line - - datasets = list(datasets_dict.values()) - if exec_dir is None: - exec_dir = os.path.abspath(os.getcwd()) - if dataset_files_path is None: - dataset_files_path = galaxy.model.Dataset.file_path - if config_root is None: - config_root = os.path.abspath(os.getcwd()) - metadata_files_list = [] - for dataset in datasets: - key = self._get_dataset_metadata_key(dataset) - # future note: - # wonkiness in job execution causes build command line to be called more than once - # when setting metadata externally, via 'auto-detect' button in edit attributes, etc., - # we don't want to overwrite (losing the ability to cleanup) our existing dataset keys and files, - # so we will only populate the dictionary once - metadata_files = self._get_output_filenames_by_dataset(dataset, sa_session) - if not metadata_files: - metadata_files = galaxy.model.JobExternalOutputMetadata(job=job, dataset=dataset) - # we are using tempfile to create unique filenames, tempfile always returns an absolute path - # we will use pathnames relative to the galaxy root, to accommodate instances where the galaxy root - # is located differently, i.e. on a cluster node with a different filesystem structure - - def _metadata_path(what): - return abspath(tempfile.NamedTemporaryFile(dir=tmp_dir, prefix=f"metadata_{what}_{key}_").name) - - filename_in, filename_out, filename_results_code, filename_kwds, filename_override_metadata = _initialize_metadata_inputs(dataset, _metadata_path, tmp_dir, kwds) - - # file to store existing dataset - metadata_files.filename_in = filename_in - - # file to store metadata results of set_meta() - metadata_files.filename_out = filename_out - - # file to store a 'return code' indicating the results of the set_meta() call - # results code is like (True/False - if setting metadata was successful/failed , exception or string of reason of success/failure ) - metadata_files.filename_results_code = filename_results_code - - # file to store kwds passed to set_meta() - metadata_files.filename_kwds = filename_kwds - - # existing metadata file parameters need to be overridden with cluster-writable file locations - metadata_files.filename_override_metadata = filename_override_metadata - - # add to session and flush - sa_session.add(metadata_files) - sa_session.flush() - metadata_files_list.append(metadata_files) - args = '"{}" "{}" {} {}'.format(metadata_path_on_compute(datatypes_config), - 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 - with tempfile.NamedTemporaryFile(mode='w', suffix='.py', dir=tmp_dir, prefix="set_metadata_", delete=False) as temp: - temp.write(SET_METADATA_SCRIPT) - return 'python "{}" {}'.format(metadata_path_on_compute(temp.name), args) - else: - # return args to galaxy_ext.metadata.set_metadata required to build - return args - - def external_metadata_set_successfully(self, dataset, name, sa_session, working_directory): - metadata_files = self._get_output_filenames_by_dataset(dataset, sa_session) - if not metadata_files: - return False # this file doesn't exist - return self._metadata_results_from_file(dataset, metadata_files.filename_results_code) - - def cleanup_external_metadata(self, sa_session): - log.debug('Cleaning up external metadata files') - for metadata_files in sa_session.query(galaxy.model.Job).get(self.job_id).external_output_metadata: - # we need to confirm that any MetadataTempFile files were removed, if not we need to remove them - # can occur if the job was stopped before completion, but a MetadataTempFile is used in the set_meta - MetadataTempFile.cleanup_from_JSON_dict_filename(metadata_files.filename_out) - dataset_key = self._get_dataset_metadata_key(metadata_files.dataset) - for key, fname in [('filename_in', metadata_files.filename_in), - ('filename_out', metadata_files.filename_out), - ('filename_results_code', metadata_files.filename_results_code), - ('filename_kwds', metadata_files.filename_kwds), - ('filename_override_metadata', metadata_files.filename_override_metadata)]: - try: - os.remove(fname) - except Exception as e: - log.debug(f'Failed to cleanup external metadata file ({key}) for {dataset_key}: {e}') - - def set_job_runner_external_pid(self, pid, sa_session): - for metadata_files in sa_session.query(galaxy.model.Job).get(self.job_id).external_output_metadata: - metadata_files.job_runner_external_pid = pid - sa_session.add(metadata_files) - sa_session.flush() - - def load_metadata(self, dataset, name, sa_session, working_directory, remote_metadata_directory=None): - # load metadata from file - # we need to no longer allow metadata to be edited while the job is still running, - # since if it is edited, the metadata changed on the running output will no longer match - # the metadata that was stored to disk for use via the external process, - # and the changes made by the user will be lost, without warning or notice - output_filename = self._get_output_filenames_by_dataset(dataset, sa_session).filename_out - self._load_metadata_from_path(dataset, output_filename, working_directory, remote_metadata_directory) - - def _initialize_metadata_inputs(dataset, path_for_part, tmp_dir, kwds, real_metadata_object=True): - filename_in = path_for_part("in") filename_out = path_for_part("out") filename_results_code = path_for_part("results") filename_kwds = path_for_part("kwds") filename_override_metadata = path_for_part("override") - _dump_dataset_instance_to(dataset, filename_in) open(filename_out, 'wt+') # create the file on disk, so it cannot be reused by tempfile (unlikely, but possible) # create the file on disk, so it cannot be reused by tempfile (unlikely, but possible) json.dump((False, 'External set_meta() not called'), open(filename_results_code, 'wt+')) @@ -446,26 +258,7 @@ def _initialize_metadata_inputs(dataset, path_for_part, tmp_dir, kwds, real_meta json.dump(override_metadata, open(filename_override_metadata, 'wt+')) - return filename_in, filename_out, filename_results_code, filename_kwds, filename_override_metadata - - -def _assert_datatypes_config(datatypes_config): - if datatypes_config is None: - raise Exception('In setup_external_metadata, the received datatypes_config is None.') - - -def _dump_dataset_instance_to(dataset_instance, file_path): - # FIXME: HACK - # sqlalchemy introduced 'expire_on_commit' flag for sessionmaker at version 0.5x - # This may be causing the dataset attribute of the dataset_association object to no-longer be loaded into memory when needed for pickling. - # For now, we'll simply 'touch' dataset_association.dataset to force it back into memory. - dataset_instance.dataset # force dataset_association.dataset to be loaded before pickling - # A better fix could be setting 'expire_on_commit=False' on the session, or modifying where commits occur, or ? - - # Touch also deferred column - dataset_instance._metadata - - pickle.dump(dataset_instance, open(file_path, 'wb+')) + return filename_out, filename_results_code, filename_kwds, filename_override_metadata def _get_filename_override(output_fnames, file_name): diff --git a/lib/galaxy/metadata/set_metadata.py b/lib/galaxy/metadata/set_metadata.py index 9fb3769e395..e47adfaafc0 100644 --- a/lib/galaxy/metadata/set_metadata.py +++ b/lib/galaxy/metadata/set_metadata.py @@ -13,7 +13,6 @@ constructed automatically). import json import logging import os -import pickle import sys import traceback @@ -23,10 +22,7 @@ import galaxy.model.mapping # need to load this before we unpickle, in order to from galaxy.model import store 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, - unicodify, -) +from galaxy.util import stringify_dictionary_keys logging.basicConfig() log = logging.getLogger(__name__) @@ -78,10 +74,7 @@ def set_meta_with_tool_provided(dataset_instance, file_dict, set_meta_kwds, data def set_metadata(): - if len(sys.argv) == 1: - set_metadata_portable() - else: - set_metadata_legacy() + set_metadata_portable() def set_metadata_portable(): @@ -172,17 +165,13 @@ def set_metadata_portable(): job_context = ExpressionContext(dict(stdout=tool_stdout, stderr=tool_stderr)) # Load outputs. - import_model_store = store.imported_store_for_metadata('metadata/outputs_new', object_store=object_store) export_store = store.DirectoryModelExportStore('metadata/outputs_populated', serialize_dataset_objects=True, for_edit=True, strip_metadata_files=False) + import_model_store = store.imported_store_for_metadata('metadata/outputs_new', object_store=object_store) for output_name, output_dict in outputs.items(): - if extended_metadata_collection: - dataset_instance_id = output_dict["id"] - dataset = import_model_store.sa_session.query(galaxy.model.HistoryDatasetAssociation).find(dataset_instance_id) - assert dataset is not None - else: - filename_in = os.path.join("metadata/metadata_in_%s" % output_name) - dataset = pickle.load(open(filename_in, 'rb')) # load DatasetInstance + dataset_instance_id = output_dict["id"] + dataset = import_model_store.sa_session.query(galaxy.model.HistoryDatasetAssociation).find(dataset_instance_id) + assert dataset is not None filename_kwds = os.path.join("metadata/metadata_kwds_%s" % output_name) filename_out = os.path.join("metadata/metadata_out_%s" % output_name) @@ -311,64 +300,6 @@ def set_metadata_portable(): 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, None) - - 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 = pickle.load(open(filename_in, 'rb')) # load DatasetInstance - dataset.dataset.external_filename = dataset_filename_override - 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, "working", extra_files_dir_name)) - dataset.dataset.external_extra_files_path = files_path - file_dict = tool_provided_metadata.get_dataset_meta(None, dataset.dataset.id, dataset.dataset.uuid) - 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, unicodify(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)) diff --git a/lib/galaxy/webapps/galaxy/config_schema.yml b/lib/galaxy/webapps/galaxy/config_schema.yml index 5471ac209ae..ed024e7d944 100644 --- a/lib/galaxy/webapps/galaxy/config_schema.yml +++ b/lib/galaxy/webapps/galaxy/config_schema.yml @@ -2818,7 +2818,7 @@ mapping: required: false default: directory desc: | - Determines how metadata will be set. Valid values are `directory`, `extended` and `legacy`. + Determines how metadata will be set. Valid values are `directory` and `extended`. In extended mode jobs will decide if a tool run failed, the object stores configuration is serialized and made available to the job and is used for writing output datasets to the object store as part of the job and dynamic diff --git a/test/unit/tools/test_metadata.py b/test/unit/tools/test_metadata.py index bbb2dd7db1b..8f802a11071 100644 --- a/test/unit/tools/test_metadata.py +++ b/test/unit/tools/test_metadata.py @@ -33,10 +33,6 @@ class MetadataTestCase(unittest.TestCase, tools_support.UsesApp, tools_support.U super().tearDown() self.metadata_compute_strategy = None - def test_simple_output_legacy(self): - self.app.config.metadata_strategy = "legacy" - self._test_simple_output() - def test_simple_output_directory(self): self.app.config.metadata_strategy = "directory" self._test_simple_output() @@ -66,10 +62,6 @@ class MetadataTestCase(unittest.TestCase, tools_support.UsesApp, tools_support.U assert output_dataset.metadata.data_lines == 2 assert output_dataset.metadata.sequences == 1 - def test_primary_dataset_output_extension_legacy(self): - self.app.config.metadata_strategy = "legacy" - self._test_primary_dataset_output_extension() - def test_primary_dataset_output_extension_directory(self): self.app.config.metadata_strategy = "directory" self._test_primary_dataset_output_extension() @@ -99,10 +91,6 @@ class MetadataTestCase(unittest.TestCase, tools_support.UsesApp, tools_support.U assert output_dataset.metadata.data_lines == 2 assert output_dataset.metadata.sequences == 1 - def test_primary_dataset_output_metadata_override_legacy(self): - self.app.config.metadata_strategy = "legacy" - self._test_primary_dataset_output_metadata_override() - def test_primary_dataset_output_metadata_override_directory(self): self.app.config.metadata_strategy = "directory" self._test_primary_dataset_output_metadata_override()