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.
This commit is contained in:
mvdbeek
2021-02-05 11:38:50 +01:00
parent 2c50b903b9
commit 7d0ec29fa0
6 changed files with 33 additions and 321 deletions
+7 -7
View File
@@ -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 <discover_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 <discover_datasets>, unpopulated
collections, etc) happens as part of the job.
:Default: ``directory``
:Type: str
+3 -3
View File
@@ -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 <discover_datasets>, unpopulated collections,
+16 -223
View File
@@ -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):
+6 -75
View File
@@ -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))
+1 -1
View File
@@ -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
-12
View File
@@ -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()