mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
Datatype-based validation framework.
There was an attempt to do this over a decade ago that didn't take off: - xref925aec1bd0- xreff4f9e1d7b2
This commit is contained in:
@@ -22,7 +22,7 @@ from bx.seq.twobit import TWOBIT_MAGIC_NUMBER, TWOBIT_MAGIC_NUMBER_SWAP, TWOBIT_
|
||||
|
||||
from galaxy import util
|
||||
from galaxy.datatypes import metadata
|
||||
from galaxy.datatypes.data import get_file_peek
|
||||
from galaxy.datatypes.data import DatatypeValidation, get_file_peek
|
||||
from galaxy.datatypes.metadata import DictParameter, ListParameter, MetadataElement, MetadataParameter
|
||||
from galaxy.util import nice_size, sqlite
|
||||
from galaxy.util.checkers import is_bz2, is_gzip
|
||||
@@ -284,6 +284,10 @@ class BamNative(CompressedArchive):
|
||||
Binary.init_meta(self, dataset, copy_from=copy_from)
|
||||
|
||||
def sniff(self, filename):
|
||||
return BamNative.is_bam(filename)
|
||||
|
||||
@classmethod
|
||||
def is_bam(cls, filename):
|
||||
# BAM is compressed in the BGZF format, and must not be uncompressed in Galaxy.
|
||||
# The first 4 bytes of any bam file is 'BAM\1', and the file is binary.
|
||||
try:
|
||||
@@ -423,6 +427,13 @@ class BamNative(CompressedArchive):
|
||||
column_names=column_names,
|
||||
column_types=column_types)
|
||||
|
||||
def validate(self, dataset, **kwd):
|
||||
if not BamNative.is_bam(dataset.file_name):
|
||||
return DatatypeValidation.invalid("This dataset does not appear to a BAM file.")
|
||||
elif self.dataset_content_needs_grooming(dataset.file_name):
|
||||
return DatatypeValidation.invalid("This BAM file does not appear to have the correct sorting for declared datatype.")
|
||||
return DatatypeValidation.validated()
|
||||
|
||||
|
||||
@dataproviders.decorators.has_dataproviders
|
||||
class Bam(BamNative):
|
||||
|
||||
@@ -48,6 +48,36 @@ DOWNLOAD_FILENAME_PATTERN_DATASET = "Galaxy${hid}-[${name}].${ext}"
|
||||
DOWNLOAD_FILENAME_PATTERN_COLLECTION_ELEMENT = "Galaxy${hdca_hid}-[${hdca_name}__${element_identifier}].${ext}"
|
||||
|
||||
|
||||
class DatatypeValidation(object):
|
||||
|
||||
def __init__(self, state, message):
|
||||
self.state = state
|
||||
self.message = message
|
||||
|
||||
@staticmethod
|
||||
def validated():
|
||||
return DatatypeValidation("ok", "Dataset validated by datatype validator.")
|
||||
|
||||
@staticmethod
|
||||
def invalid(message):
|
||||
return DatatypeValidation("invalid", message)
|
||||
|
||||
@staticmethod
|
||||
def unvalidated():
|
||||
return DatatypeValidation("unknown", "Dataset validation unimplemented for this datatype.")
|
||||
|
||||
def __repr__(self):
|
||||
return "DatatypeValidation[state=%s,message=%s]" % (self.state, self.message)
|
||||
|
||||
|
||||
def validate(dataset_instance):
|
||||
try:
|
||||
datatype_validation = dataset_instance.datatype.validate(dataset_instance)
|
||||
except Exception as e:
|
||||
datatype_validation = DatatypeValidation.invalid("Problem running datatype validation method [%s]" % str(e))
|
||||
return datatype_validation
|
||||
|
||||
|
||||
class DataMeta(abc.ABCMeta):
|
||||
"""
|
||||
Metaclass for Data class. Sets up metadata spec.
|
||||
@@ -503,10 +533,6 @@ class Data(object):
|
||||
except Exception:
|
||||
return "info unavailable"
|
||||
|
||||
def validate(self, dataset):
|
||||
"""Unimplemented validate, return no exceptions"""
|
||||
return list()
|
||||
|
||||
def repair_methods(self, dataset):
|
||||
"""Unimplemented method, returns dict with method/option for repairing errors"""
|
||||
return None
|
||||
@@ -743,6 +769,9 @@ class Data(object):
|
||||
return self.dataproviders[data_format](self, dataset, **settings)
|
||||
raise dataproviders.exceptions.NoProviderAvailable(self, data_format)
|
||||
|
||||
def validate(self, dataset, **kwd):
|
||||
return DatatypeValidation.unvalidated()
|
||||
|
||||
@dataproviders.decorators.dataprovider_factory('base')
|
||||
def base_dataprovider(self, dataset, **settings):
|
||||
dataset_source = dataproviders.dataset.DatasetDataProvider(dataset)
|
||||
|
||||
@@ -20,7 +20,7 @@ from markupsafe import escape
|
||||
from six.moves.urllib.parse import quote_plus
|
||||
|
||||
from galaxy.datatypes import metadata
|
||||
from galaxy.datatypes.data import Text
|
||||
from galaxy.datatypes.data import DatatypeValidation, Text
|
||||
from galaxy.datatypes.metadata import MetadataElement
|
||||
from galaxy.datatypes.sniff import build_sniff_from_prefix
|
||||
from galaxy.datatypes.tabular import Tabular
|
||||
@@ -151,24 +151,17 @@ class GenomeGraphs(Tabular):
|
||||
out = "Can't create peek %s" % exc
|
||||
return out
|
||||
|
||||
def validate(self, dataset):
|
||||
def validate(self, dataset, **kwd):
|
||||
"""
|
||||
Validate a gg file - all numeric after header row
|
||||
"""
|
||||
errors = list()
|
||||
with open(dataset.file_name, "r") as infile:
|
||||
next(infile) # header
|
||||
for i, row in enumerate(infile):
|
||||
ll = row.strip().split('\t')[1:] # first is alpha feature identifier
|
||||
badvals = []
|
||||
for j, x in enumerate(ll):
|
||||
try:
|
||||
x = float(x)
|
||||
except Exception:
|
||||
badvals.append('col%d:%s' % (j + 1, x))
|
||||
if len(badvals) > 0:
|
||||
errors.append('row %d, %s' % (' '.join(badvals)))
|
||||
return errors
|
||||
x = float(x)
|
||||
return DatatypeValidation.validated()
|
||||
|
||||
def sniff_prefix(self, file_prefix):
|
||||
"""
|
||||
|
||||
@@ -11,6 +11,7 @@ from six.moves.urllib.parse import quote_plus
|
||||
|
||||
from galaxy import util
|
||||
from galaxy.datatypes import metadata
|
||||
from galaxy.datatypes.data import DatatypeValidation
|
||||
from galaxy.datatypes.metadata import MetadataElement
|
||||
from galaxy.datatypes.sniff import (
|
||||
build_sniff_from_prefix,
|
||||
@@ -269,9 +270,8 @@ class Interval(Tabular):
|
||||
ret_val.append((site_name, link))
|
||||
return ret_val
|
||||
|
||||
def validate(self, dataset):
|
||||
def validate(self, dataset, **kwd):
|
||||
"""Validate an interval file using the bx GenomicIntervalReader"""
|
||||
errors = list()
|
||||
c, s, e, t = dataset.metadata.chromCol, dataset.metadata.startCol, dataset.metadata.endCol, dataset.metadata.strandCol
|
||||
c, s, e, t = int(c) - 1, int(s) - 1, int(e) - 1, int(t) - 1
|
||||
with open(dataset.file_name, "r") as infile:
|
||||
@@ -286,9 +286,9 @@ class Interval(Tabular):
|
||||
try:
|
||||
next(reader)
|
||||
except ParseError as e:
|
||||
errors.append(e)
|
||||
return DatatypeValidation.invalid(str(e))
|
||||
except StopIteration:
|
||||
return errors
|
||||
return DatatypeValidation.valid()
|
||||
|
||||
def repair_methods(self, dataset):
|
||||
"""Return options for removing errors along with a description"""
|
||||
|
||||
@@ -19,6 +19,7 @@ from galaxy.datatypes import metadata
|
||||
from galaxy.datatypes.binary import (
|
||||
Binary
|
||||
)
|
||||
from galaxy.datatypes.data import DatatypeValidation
|
||||
from galaxy.datatypes.metadata import DictParameter, MetadataElement
|
||||
from galaxy.datatypes.sniff import (
|
||||
build_sniff_from_prefix,
|
||||
@@ -795,16 +796,36 @@ class BaseFastq(Sequence):
|
||||
@classmethod
|
||||
def check_first_block(cls, file_prefix):
|
||||
# check that first block looks like a fastq block
|
||||
headers = get_headers(file_prefix, sep='\n', count=4)
|
||||
if len(headers) == 4 and headers[0][0] and headers[0][0][0] == "@" and headers[2][0] and headers[2][0][0] == "+" and headers[1][0]:
|
||||
block = get_headers(file_prefix, sep='\n', count=4)
|
||||
return cls.check_block(block)
|
||||
|
||||
@classmethod
|
||||
def check_block(cls, block):
|
||||
if len(block) == 4 and block[0][0] and block[0][0][0] == "@" and block[2][0] and block[2][0][0] == "+" and block[1][0]:
|
||||
# Check the sequence line, make sure it contains only G/C/A/T/N
|
||||
match = cls.bases_regexp.match(headers[1][0])
|
||||
match = cls.bases_regexp.match(block[1][0])
|
||||
if match:
|
||||
start, end = match.span()
|
||||
if (end - start) == len(headers[1][0]):
|
||||
if (end - start) == len(block[1][0]):
|
||||
return True
|
||||
return False
|
||||
|
||||
def validate(self, dataset, **kwd):
|
||||
headers = iter_headers(dataset.file_name, sep='\n', count=-1)
|
||||
# check to see if the base qualities match
|
||||
if not self.quality_check(headers):
|
||||
return DatatypeValidation.invalid("Invalid quality score(s) found for this fastq datatype.")
|
||||
|
||||
headers = iter_headers(dataset.file_name, sep='\n', count=-1)
|
||||
while True:
|
||||
block = list(islice(headers, 4))
|
||||
if len(block) == 0:
|
||||
break
|
||||
if not self.check_block(block):
|
||||
return DatatypeValidation.invalid("Invalid FASTQ structure found.")
|
||||
|
||||
return DatatypeValidation.validated()
|
||||
|
||||
|
||||
class Fastq(BaseFastq):
|
||||
"""Class representing a generic FASTQ sequence"""
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
<tool id="__SET_METADATA__" name="Set External Metadata" version="1.0.1" tool_type="set_metadata">
|
||||
<tool id="__SET_METADATA__" name="Set External Metadata" version="1.0.2" tool_type="set_metadata">
|
||||
<type class="SetMetadataTool" module="galaxy.tools"/>
|
||||
<requirements>
|
||||
<requirement type="package" version="1.5">bcftools</requirement>
|
||||
|
||||
@@ -163,6 +163,11 @@ def iter_headers(fname_or_file_prefix, sep, count=60, comment_designator=None):
|
||||
break
|
||||
|
||||
|
||||
def validate_tabular(fname_or_file_prefix, validate_row, sep, comment_designator=None):
|
||||
for row in iter_headers(fname_or_file_prefix, sep, count=-1, comment_designator=comment_designator):
|
||||
validate_row(row)
|
||||
|
||||
|
||||
def get_headers(fname_or_file_prefix, sep, count=60, comment_designator=None):
|
||||
"""
|
||||
Returns a list with the first 'count' lines split by 'sep', ignoring lines
|
||||
|
||||
@@ -24,7 +24,8 @@ from galaxy.datatypes.metadata import MetadataElement
|
||||
from galaxy.datatypes.sniff import (
|
||||
build_sniff_from_prefix,
|
||||
get_headers,
|
||||
iter_headers
|
||||
iter_headers,
|
||||
validate_tabular,
|
||||
)
|
||||
from galaxy.util import compression_utils
|
||||
from . import dataproviders
|
||||
@@ -730,6 +731,19 @@ class BaseVcf(Tabular):
|
||||
if exit_code != 0:
|
||||
raise Exception("Error merging VCF files: %s" % stderr)
|
||||
|
||||
def validate(self, dataset, **kwd):
|
||||
# with tempfile.NamedTemporaryFile() as t:
|
||||
# try:
|
||||
# pysam.tabix_index(dataset.file_name, index=t.name, preset='vcf', force=False, keep_original=True)
|
||||
# except Exception as e:
|
||||
# data.DatatypeValidation.invalid("Failed to generate index [%s]" % e)
|
||||
# return data.DatatypeValidation.validated()
|
||||
def validate_row(row):
|
||||
if len(row) < 8:
|
||||
raise Exception("Not enough columns in row %s" % row.join("\t"))
|
||||
validate_tabular(dataset.file_name, sep='\t', validate_row=validate_row, comment_designator="#")
|
||||
return data.DatatypeValidation.validated()
|
||||
|
||||
# Dataproviders
|
||||
@dataproviders.decorators.dataprovider_factory('genomic-region',
|
||||
dataproviders.dataset.GenomicRegionDataProvider.settings)
|
||||
|
||||
@@ -998,6 +998,15 @@ class JobWrapper(HasResourceParameters):
|
||||
param_dict = self.tool.params_from_strings(param_dict, self.app)
|
||||
return param_dict
|
||||
|
||||
@property
|
||||
def validate_outputs(self):
|
||||
job = self.get_job()
|
||||
for p in job.parameters:
|
||||
if p.name == "__validate_outputs__":
|
||||
log.info("validate... %s" % p.value)
|
||||
return loads(p.value)
|
||||
return False
|
||||
|
||||
def get_version_string_path(self):
|
||||
return os.path.abspath(os.path.join(self.working_directory, COMMAND_VERSION_FILENAME))
|
||||
|
||||
@@ -1953,6 +1962,7 @@ class JobWrapper(HasResourceParameters):
|
||||
datatypes_config=datatypes_config,
|
||||
job_metadata=os.path.join(self.tool_working_directory, self.tool.provided_metadata_file),
|
||||
max_metadata_value_size=self.app.config.max_metadata_value_size,
|
||||
validate_outputs=self.validate_outputs,
|
||||
**kwds)
|
||||
if resolve_metadata_dependencies:
|
||||
metadata_tool = self.app.toolbox.get_tool("__SET_METADATA__")
|
||||
|
||||
@@ -68,6 +68,23 @@ class EmailAction(DefaultJobAction):
|
||||
return "Email the current user when this job is complete."
|
||||
|
||||
|
||||
class ValidateOutputsAction(DefaultJobAction):
|
||||
"""
|
||||
This action sends an email to the galaxy user responsible for a job.
|
||||
"""
|
||||
name = "ValidateOutputsAction"
|
||||
verbose_name = "Validate Tool Outputs"
|
||||
|
||||
@classmethod
|
||||
def execute(cls, app, sa_session, action, job, replacement_dict):
|
||||
# no-op: needs to inject metadata handling parameters ahead of time.
|
||||
pass
|
||||
|
||||
@classmethod
|
||||
def get_short_str(cls, pja):
|
||||
return "Validate tool outputs."
|
||||
|
||||
|
||||
class ChangeDatatypeAction(DefaultJobAction):
|
||||
name = "ChangeDatatypeAction"
|
||||
verbose_name = "Change Datatype"
|
||||
|
||||
@@ -370,7 +370,7 @@ class DatasetAssociationManager(base.ModelManager,
|
||||
else:
|
||||
raise exceptions.InsufficientPermissionsException('Changing datatype "%s" is not allowed.' % (data.extension))
|
||||
|
||||
def set_metadata(self, trans, dataset_assoc, overwrite=False):
|
||||
def set_metadata(self, trans, dataset_assoc, overwrite=False, validate=True):
|
||||
"""Trigger a job that detects and sets metadata on a given dataset association (ldda or hda)"""
|
||||
data = trans.sa_session.query(self.model_class).get(dataset_assoc.id)
|
||||
if not self.ok_to_edit_metadata(data.id):
|
||||
@@ -384,7 +384,7 @@ class DatasetAssociationManager(base.ModelManager,
|
||||
setattr(data.metadata, name, spec.unwrap(spec.get('default')))
|
||||
|
||||
self.app.datatypes_registry.set_external_metadata_tool.tool_action.execute(
|
||||
self.app.datatypes_registry.set_external_metadata_tool, trans, incoming={'input1': data},
|
||||
self.app.datatypes_registry.set_external_metadata_tool, trans, incoming={'input1': data, 'validate': validate},
|
||||
overwrite=overwrite)
|
||||
|
||||
def update_permissions(self, trans, dataset_assoc, **kwd):
|
||||
|
||||
@@ -295,6 +295,9 @@ class HDASerializer( # datasets._UnflattenedMetadataDatasetAssociationSerialize
|
||||
'display_types',
|
||||
'visualizations',
|
||||
|
||||
'validated_state',
|
||||
'validated_state_message',
|
||||
|
||||
# 'url',
|
||||
'download_url',
|
||||
|
||||
|
||||
@@ -107,6 +107,7 @@ class PortableDirectoryMetadataGenerator(MetadataCollectionStrategy):
|
||||
config_file=None, datatypes_config=None,
|
||||
job_metadata=None, compute_tmp_dir=None,
|
||||
include_command=True, max_metadata_value_size=0,
|
||||
validate_outputs=False,
|
||||
kwds=None):
|
||||
assert job_metadata, "setup_external_metadata must be supplied with job_metadata path"
|
||||
kwds = kwds or {}
|
||||
@@ -134,7 +135,8 @@ class PortableDirectoryMetadataGenerator(MetadataCollectionStrategy):
|
||||
_initialize_metadata_inputs(dataset, _metadata_path, tmp_dir, kwds)
|
||||
|
||||
outputs[name] = {
|
||||
"filename_override": _get_filename_override(output_fnames, dataset.file_name)
|
||||
"filename_override": _get_filename_override(output_fnames, dataset.file_name),
|
||||
"validate": validate_outputs,
|
||||
}
|
||||
|
||||
metadata_params_path = os.path.join(metadata_dir, "params.json")
|
||||
@@ -226,6 +228,7 @@ class JobExternalOutputMetadataWrapper(MetadataCollectionStrategy):
|
||||
config_file=None, datatypes_config=None,
|
||||
job_metadata=None, compute_tmp_dir=None,
|
||||
include_command=True, max_metadata_value_size=0,
|
||||
validate_outputs=False,
|
||||
kwds=None):
|
||||
kwds = kwds or {}
|
||||
tmp_dir = _init_tmp_dir(tmp_dir)
|
||||
|
||||
@@ -32,6 +32,18 @@ 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_validated_state(dataset_instance):
|
||||
from galaxy.datatypes.data import validate
|
||||
datatype_validation = validate(dataset_instance)
|
||||
|
||||
dataset_instance.validated_state = datatype_validation.state
|
||||
dataset_instance.validated_state_message = datatype_validation.message
|
||||
|
||||
# Set special metadata property that will reload this on server side.
|
||||
setattr(dataset_instance.metadata, "__validated_state__", datatype_validation.state)
|
||||
setattr(dataset_instance.metadata, "__validated_state_message__", datatype_validation.message)
|
||||
|
||||
|
||||
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.
|
||||
@@ -119,6 +131,8 @@ def set_metadata_portable():
|
||||
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)
|
||||
if output_dict.get("validate", False):
|
||||
set_validated_state(dataset)
|
||||
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
|
||||
|
||||
@@ -2428,10 +2428,15 @@ class DatasetInstance(object):
|
||||
states = Dataset.states
|
||||
conversion_messages = Dataset.conversion_messages
|
||||
permitted_actions = Dataset.permitted_actions
|
||||
validated_states = Bunch(
|
||||
UNKNOWN='unknown',
|
||||
INVALID='invalid',
|
||||
OK='ok',
|
||||
)
|
||||
|
||||
def __init__(self, id=None, hid=None, name=None, info=None, blurb=None, peek=None, tool_version=None, extension=None,
|
||||
dbkey=None, metadata=None, history=None, dataset=None, deleted=False, designation=None,
|
||||
parent_id=None, validation_errors=None, visible=True, create_dataset=False, sa_session=None,
|
||||
parent_id=None, validated_state='unknown', validated_state_message=None, visible=True, create_dataset=False, sa_session=None,
|
||||
extended_metadata=None, flush=True):
|
||||
self.name = name or "Unnamed dataset"
|
||||
self.id = id
|
||||
@@ -2449,6 +2454,8 @@ class DatasetInstance(object):
|
||||
self._metadata['dbkey'] = dbkey
|
||||
self.deleted = deleted
|
||||
self.visible = visible
|
||||
self.validated_state = validated_state
|
||||
self.validated_state_message = validated_state_message
|
||||
# Relationships
|
||||
if not dataset and create_dataset:
|
||||
# Had to pass the sqlalchemy session in order to create a new dataset
|
||||
@@ -2458,7 +2465,6 @@ class DatasetInstance(object):
|
||||
sa_session.flush()
|
||||
self.dataset = dataset
|
||||
self.parent_id = parent_id
|
||||
self.validation_errors = validation_errors
|
||||
|
||||
@property
|
||||
def peek(self):
|
||||
@@ -3179,6 +3185,8 @@ class HistoryDatasetAssociation(DatasetInstance, HasTags, Dictifiable, UsesAnnot
|
||||
update_time=hda.update_time.isoformat(),
|
||||
data_type=hda.datatype.__class__.__module__ + '.' + hda.datatype.__class__.__name__,
|
||||
genome_build=hda.dbkey,
|
||||
validated_state=hda.validated_state,
|
||||
validated_state_message=hda.validated_state_message,
|
||||
misc_info=hda.info.strip() if isinstance(hda.info, string_types) else hda.info,
|
||||
misc_blurb=hda.blurb)
|
||||
|
||||
|
||||
@@ -237,6 +237,8 @@ model.HistoryDatasetAssociation.table = Table(
|
||||
Column("version", Integer, default=1, nullable=True, index=True),
|
||||
Column("hid", Integer),
|
||||
Column("purged", Boolean, index=True, default=False),
|
||||
Column("validated_state", TrimmedString(64), default='unvalidated', nullable=False),
|
||||
Column("validated_state_message", TEXT),
|
||||
Column("hidden_beneath_collection_instance_id",
|
||||
ForeignKey("history_dataset_collection_association.id"), nullable=True))
|
||||
|
||||
@@ -516,6 +518,8 @@ model.LibraryDatasetDatasetAssociation.table = Table(
|
||||
Column("parent_id", Integer, ForeignKey("library_dataset_dataset_association.id"), nullable=True),
|
||||
Column("designation", TrimmedString(255)),
|
||||
Column("deleted", Boolean, index=True, default=False),
|
||||
Column("validated_state", TrimmedString(64), default='unvalidated', nullable=False),
|
||||
Column("validated_state_message", TEXT),
|
||||
Column("visible", Boolean),
|
||||
Column("extended_metadata_id", Integer, ForeignKey("extended_metadata.id"), index=True),
|
||||
Column("user_id", Integer, ForeignKey("galaxy_user.id"), index=True),
|
||||
|
||||
@@ -193,6 +193,10 @@ class MetadataCollection(object):
|
||||
dataset._metadata[name] = value
|
||||
if '__extension__' in JSONified_dict:
|
||||
dataset.extension = JSONified_dict['__extension__']
|
||||
if '__validated_state__' in JSONified_dict:
|
||||
dataset.validated_state = JSONified_dict['__validated_state__']
|
||||
if '__validated_state_message__' in JSONified_dict:
|
||||
dataset.validated_state_message = JSONified_dict['__validated_state_message__']
|
||||
|
||||
def to_JSON_dict(self, filename=None):
|
||||
# galaxy.model.customtypes.json_encoder.encode()
|
||||
@@ -203,6 +207,10 @@ class MetadataCollection(object):
|
||||
meta_dict[name] = spec.param.to_external_value(dataset_meta_dict[name])
|
||||
if '__extension__' in dataset_meta_dict:
|
||||
meta_dict['__extension__'] = dataset_meta_dict['__extension__']
|
||||
if '__validated_state__' in dataset_meta_dict:
|
||||
meta_dict['__validated_state__'] = dataset_meta_dict['__validated_state__']
|
||||
if '__validated_state_message__' in dataset_meta_dict:
|
||||
meta_dict['__validated_state_message__'] = dataset_meta_dict['__validated_state_message__']
|
||||
if filename is None:
|
||||
return json.dumps(meta_dict)
|
||||
json.dump(meta_dict, open(filename, 'wt+'))
|
||||
|
||||
@@ -0,0 +1,58 @@
|
||||
"""
|
||||
Rework dataset validation in database.
|
||||
"""
|
||||
from __future__ import print_function
|
||||
|
||||
import logging
|
||||
|
||||
from sqlalchemy import (
|
||||
Column,
|
||||
ForeignKey,
|
||||
Integer,
|
||||
MetaData,
|
||||
Table,
|
||||
TEXT,
|
||||
)
|
||||
|
||||
from galaxy.model.custom_types import TrimmedString
|
||||
from galaxy.model.migrate.versions.util import add_column, create_table, drop_column, drop_table
|
||||
|
||||
log = logging.getLogger(__name__)
|
||||
metadata = MetaData()
|
||||
|
||||
validation_error_table = Table("validation_error", metadata,
|
||||
Column("id", Integer, primary_key=True),
|
||||
Column("dataset_id", Integer, ForeignKey("history_dataset_association.id"), index=True),
|
||||
Column("message", TrimmedString(255)),
|
||||
Column("err_type", TrimmedString(64)),
|
||||
Column("attributes", TEXT))
|
||||
|
||||
|
||||
def upgrade(migrate_engine):
|
||||
print(__doc__)
|
||||
metadata.bind = migrate_engine
|
||||
metadata.reflect()
|
||||
|
||||
drop_table(validation_error_table)
|
||||
|
||||
history_dataset_association_table = Table("history_dataset_association", metadata, autoload=True)
|
||||
library_dataset_dataset_association_table = Table("library_dataset_dataset_association", metadata, autoload=True)
|
||||
for dataset_instance_table in [history_dataset_association_table, library_dataset_dataset_association_table]:
|
||||
validated_state_column = Column('validated_state', TrimmedString(64), default='unknown', server_default="unknown", nullable=False)
|
||||
add_column(validated_state_column, dataset_instance_table)
|
||||
|
||||
validated_state_message_column = Column('validated_state_message', TEXT)
|
||||
add_column(validated_state_message_column, dataset_instance_table)
|
||||
|
||||
|
||||
def downgrade(migrate_engine):
|
||||
metadata.bind = migrate_engine
|
||||
metadata.reflect()
|
||||
|
||||
create_table(validation_error_table)
|
||||
|
||||
history_dataset_association_table = Table("history_dataset_association", metadata, autoload=True)
|
||||
library_dataset_dataset_association_table = Table("library_dataset_dataset_association", metadata, autoload=True)
|
||||
for dataset_instance_table in [history_dataset_association_table, library_dataset_dataset_association_table]:
|
||||
drop_column('validated_state', dataset_instance_table)
|
||||
drop_column('validated_state_message', dataset_instance_table)
|
||||
@@ -5,6 +5,7 @@ from json import dumps
|
||||
|
||||
from galaxy.job_execution.datasets import DatasetPath
|
||||
from galaxy.metadata import get_metadata_compute_strategy
|
||||
from galaxy.util import asbool
|
||||
from . import ToolAction
|
||||
|
||||
log = logging.getLogger(__name__)
|
||||
@@ -34,7 +35,11 @@ class SetMetadataToolAction(ToolAction):
|
||||
"""
|
||||
Execute using application.
|
||||
"""
|
||||
|
||||
for name, value in incoming.items():
|
||||
# Why are we looping here and not just using a fixed input name? Needed?
|
||||
if not name.startswith("input"):
|
||||
continue
|
||||
if isinstance(value, app.model.HistoryDatasetAssociation):
|
||||
dataset = value
|
||||
dataset_name = name
|
||||
@@ -83,6 +88,7 @@ class SetMetadataToolAction(ToolAction):
|
||||
output_datatasets_dict = {
|
||||
dataset_name: dataset,
|
||||
}
|
||||
validate_outputs = asbool(incoming.get("validate", False))
|
||||
cmd_line = external_metadata_wrapper.setup_external_metadata(output_datatasets_dict,
|
||||
sa_session,
|
||||
exec_dir=None,
|
||||
@@ -94,6 +100,7 @@ class SetMetadataToolAction(ToolAction):
|
||||
datatypes_config=datatypes_config,
|
||||
job_metadata=os.path.join(job_working_dir, 'working', tool.provided_metadata_file),
|
||||
include_command=False,
|
||||
validate_outputs=validate_outputs,
|
||||
max_metadata_value_size=app.config.max_metadata_value_size,
|
||||
kwds={'overwrite': overwrite})
|
||||
incoming['__SET_EXTERNAL_METADATA_COMMAND_LINE__'] = cmd_line
|
||||
|
||||
@@ -29,7 +29,7 @@ class PartialJobExecution(Exception):
|
||||
MappingParameters = collections.namedtuple("MappingParameters", ["param_template", "param_combinations"])
|
||||
|
||||
|
||||
def execute(trans, tool, mapping_params, history, rerun_remap_job_id=None, collection_info=None, workflow_invocation_uuid=None, invocation_step=None, max_num_jobs=None, job_callback=None, completed_jobs=None, workflow_resource_parameters=None):
|
||||
def execute(trans, tool, mapping_params, history, rerun_remap_job_id=None, collection_info=None, workflow_invocation_uuid=None, invocation_step=None, max_num_jobs=None, job_callback=None, completed_jobs=None, workflow_resource_parameters=None, validate_outputs=False):
|
||||
"""
|
||||
Execute a tool and return object containing summary (output data, number of
|
||||
failures, etc...).
|
||||
@@ -61,6 +61,8 @@ def execute(trans, tool, mapping_params, history, rerun_remap_job_id=None, colle
|
||||
# Only workflow invocation code gets to set this, ignore user supplied
|
||||
# values or rerun parameters.
|
||||
del params['__workflow_resource_params__']
|
||||
if validate_outputs:
|
||||
params['__validate_outputs__'] = True
|
||||
job, result = tool.handle_single_execution(trans, rerun_remap_job_id, execution_slice, history, execution_cache, completed_job, collection_info)
|
||||
if job:
|
||||
message = EXECUTION_SUCCESS_MESSAGE % (tool.id, job.id, job_timer)
|
||||
|
||||
@@ -652,7 +652,7 @@ class HistoryContentsController(BaseAPIController, UsesLibraryMixin, UsesLibrary
|
||||
:type history_id: str
|
||||
:param history_id: encoded id string of the items's History
|
||||
:type id: str
|
||||
:param id: the encoded id of the history to update
|
||||
:param id: the encoded id of the history item to update
|
||||
:type payload: dict
|
||||
:param payload: a dictionary containing any or all the
|
||||
fields in :func:`galaxy.model.HistoryDatasetAssociation.to_dict`
|
||||
@@ -673,6 +673,29 @@ class HistoryContentsController(BaseAPIController, UsesLibraryMixin, UsesLibrary
|
||||
else:
|
||||
return self.__handle_unknown_contents_type(trans, contents_type)
|
||||
|
||||
@expose_api_anonymous
|
||||
def validate(self, trans, history_id, history_content_id, payload=None, **kwd):
|
||||
"""
|
||||
update( self, trans, history_id, id, payload, **kwd )
|
||||
* PUT /api/histories/{history_id}/contents/{id}/validate
|
||||
updates the values for the history content item with the given ``id``
|
||||
|
||||
:type history_id: str
|
||||
:param history_id: encoded id string of the items's History
|
||||
:type id: str
|
||||
:param id: the encoded id of the history item to validate
|
||||
|
||||
:rtype: dict
|
||||
:returns: TODO
|
||||
"""
|
||||
decoded_id = self.decode_id(history_content_id)
|
||||
history = self.history_manager.get_owned(self.decode_id(history_id), trans.user,
|
||||
current_history=trans.history)
|
||||
hda = self.hda_manager.get_owned_ids([decoded_id], history=history)[0]
|
||||
if hda:
|
||||
self.hda_manager.set_metadata(trans, hda, overwrite=True, validate=True)
|
||||
return {}
|
||||
|
||||
def __update_dataset(self, trans, history_id, id, payload, **kwd):
|
||||
# anon user: ensure that history ids match up and the history is the current,
|
||||
# check for uploading, and use only the subset of attribute keys manipulatable by anon users
|
||||
|
||||
@@ -236,6 +236,11 @@ def populate_api_routes(webapp, app):
|
||||
controller="history_contents",
|
||||
action="update_permissions",
|
||||
conditions=dict(method=["PUT"]))
|
||||
webapp.mapper.connect("history_contents_validate",
|
||||
"/api/histories/{history_id}/contents/{history_content_id}/validate",
|
||||
controller="history_contents",
|
||||
action="validate",
|
||||
conditions=dict(method=["PUT"]))
|
||||
webapp.mapper.connect("history_contents_extra_files",
|
||||
"/api/histories/{history_id}/contents/{history_content_id}/extra_files",
|
||||
controller="datasets",
|
||||
|
||||
@@ -458,7 +458,6 @@ class DatasetInterface(BaseUIController, UsesAnnotations, UsesItemRatings, UsesE
|
||||
return self.message_exception(trans, 'Changing datatype "%s" is not allowed.' % (data.extension))
|
||||
elif operation == 'autodetect':
|
||||
# The user clicked the Auto-detect button on the 'Edit Attributes' form
|
||||
# prevent modifying metadata when dataset is queued or running as input/output
|
||||
try:
|
||||
self.hda_manager.set_metadata(trans, data, overwrite=True)
|
||||
except MessageException as e:
|
||||
|
||||
@@ -1265,6 +1265,12 @@ class ToolModule(WorkflowModule):
|
||||
try:
|
||||
mapping_params = MappingParameters(tool_state.inputs, param_combinations)
|
||||
max_num_jobs = progress.maximum_jobs_to_schedule_or_none
|
||||
|
||||
validate_outputs = False
|
||||
for pja in step.post_job_actions:
|
||||
if pja.action_type == "ValidateOutputsAction":
|
||||
validate_outputs = True
|
||||
|
||||
execution_tracker = execute(
|
||||
trans=self.trans,
|
||||
tool=tool,
|
||||
@@ -1274,6 +1280,7 @@ class ToolModule(WorkflowModule):
|
||||
workflow_invocation_uuid=invocation.uuid.hex,
|
||||
invocation_step=invocation_step,
|
||||
max_num_jobs=max_num_jobs,
|
||||
validate_outputs=validate_outputs,
|
||||
job_callback=lambda job: self._handle_post_job_actions(step, job, invocation.replacement_dict),
|
||||
completed_jobs=completed_jobs,
|
||||
workflow_resource_parameters=resource_parameters
|
||||
|
||||
@@ -493,6 +493,28 @@ class ToolsUploadTestCase(api.ApiTestCase):
|
||||
history_id, new_dataset = self._upload('https://usegalaxy.org/api/version')
|
||||
self.dataset_populator.get_history_dataset_details(history_id, dataset_id=new_dataset["id"], assert_ok=True)
|
||||
|
||||
def test_upload_and_validate_invalid(self):
|
||||
path = TestDataResolver().get_filename("1.fastqsanger")
|
||||
with open(path, "rb") as fh:
|
||||
metadata = self._upload_and_get_details(fh, file_type="fastqcssanger")
|
||||
assert "validated_state" in metadata
|
||||
assert metadata['validated_state'] == 'unknown'
|
||||
history_id = metadata["history_id"]
|
||||
dataset_id = metadata["id"]
|
||||
terminal_validated_state = self.dataset_populator.validate_dataset_and_wait(history_id, dataset_id)
|
||||
assert terminal_validated_state == 'invalid', terminal_validated_state
|
||||
|
||||
def test_upload_and_validate_valid(self):
|
||||
path = TestDataResolver().get_filename("1.fastqsanger")
|
||||
with open(path, "rb") as fh:
|
||||
metadata = self._upload_and_get_details(fh, file_type="fastqsanger")
|
||||
assert "validated_state" in metadata
|
||||
assert metadata['validated_state'] == 'unknown'
|
||||
history_id = metadata["history_id"]
|
||||
dataset_id = metadata["id"]
|
||||
terminal_validated_state = self.dataset_populator.validate_dataset_and_wait(history_id, dataset_id)
|
||||
assert terminal_validated_state == 'ok', terminal_validated_state
|
||||
|
||||
def _velvet_upload(self, history_id, extra_inputs):
|
||||
payload = self.dataset_populator.upload_payload(
|
||||
history_id,
|
||||
|
||||
@@ -3148,6 +3148,69 @@ steps:
|
||||
print(hda3["deleted"])
|
||||
assert not hda4["deleted"]
|
||||
|
||||
@skip_without_tool("cat1")
|
||||
def test_validated_post_job_action_validated(self):
|
||||
with self.dataset_populator.test_history() as history_id:
|
||||
self._run_jobs("""
|
||||
class: GalaxyWorkflow
|
||||
inputs:
|
||||
input1: data
|
||||
outputs:
|
||||
wf_output_1:
|
||||
outputSource: first_cat/out_file1
|
||||
steps:
|
||||
first_cat:
|
||||
tool_id: cat1
|
||||
in:
|
||||
input1: input1
|
||||
post_job_actions:
|
||||
ValidateOutputsAction:
|
||||
action_type: ValidateOutputsAction
|
||||
""", test_data={"input1": {"type": "File", "file_type": "fastqsanger", "value": "1.fastqsanger"}}, history_id=history_id)
|
||||
hda2 = self.dataset_populator.get_history_dataset_details(history_id, hid=2)
|
||||
assert hda2["validated_state"] == "ok"
|
||||
|
||||
@skip_without_tool("cat1")
|
||||
def test_validated_post_job_action_unvalidated_default(self):
|
||||
with self.dataset_populator.test_history() as history_id:
|
||||
self._run_jobs("""
|
||||
class: GalaxyWorkflow
|
||||
inputs:
|
||||
input1: data
|
||||
outputs:
|
||||
wf_output_1:
|
||||
outputSource: first_cat/out_file1
|
||||
steps:
|
||||
first_cat:
|
||||
tool_id: cat1
|
||||
in:
|
||||
input1: input1
|
||||
""", test_data={"input1": {"type": "File", "file_type": "fastqsanger", "value": "1.fastqsanger"}}, history_id=history_id)
|
||||
hda2 = self.dataset_populator.get_history_dataset_details(history_id, hid=2)
|
||||
assert hda2["validated_state"] == "unknown"
|
||||
|
||||
@skip_without_tool("cat1")
|
||||
def test_validated_post_job_action_invalid(self):
|
||||
with self.dataset_populator.test_history() as history_id:
|
||||
self._run_jobs("""
|
||||
class: GalaxyWorkflow
|
||||
inputs:
|
||||
input1: data
|
||||
outputs:
|
||||
wf_output_1:
|
||||
outputSource: first_cat/out_file1
|
||||
steps:
|
||||
first_cat:
|
||||
tool_id: cat1
|
||||
in:
|
||||
input1: input1
|
||||
post_job_actions:
|
||||
ValidateOutputsAction:
|
||||
action_type: ValidateOutputsAction
|
||||
""", test_data={"input1": {"type": "File", "file_type": "fastqcssanger", "value": "1.fastqsanger"}}, history_id=history_id)
|
||||
hda2 = self.dataset_populator.get_history_dataset_details(history_id, hid=2)
|
||||
assert hda2["validated_state"] == "invalid"
|
||||
|
||||
@skip_without_tool("random_lines1")
|
||||
def test_run_replace_params_by_tool(self):
|
||||
workflow_request, history_id = self._setup_random_x2_workflow("test_for_replace_tool_params")
|
||||
|
||||
@@ -514,6 +514,28 @@ class BaseDatasetPopulator(object):
|
||||
assert update_response.status_code == 200, update_response.content
|
||||
return update_response.json()
|
||||
|
||||
def validate_dataset(self, history_id, dataset_id):
|
||||
url = "histories/%s/contents/%s/validate" % (history_id, dataset_id)
|
||||
update_response = self.galaxy_interactor._put(url, {})
|
||||
assert update_response.status_code == 200, update_response.content
|
||||
return update_response.json()
|
||||
|
||||
def validate_dataset_and_wait(self, history_id, dataset_id):
|
||||
self.validate_dataset(history_id, dataset_id)
|
||||
|
||||
def validated():
|
||||
metadata = self.get_history_dataset_details(history_id, dataset_id=dataset_id)
|
||||
validated_state = metadata['validated_state']
|
||||
if validated_state == 'unknown':
|
||||
return
|
||||
else:
|
||||
return validated_state
|
||||
|
||||
return wait_on(
|
||||
validated,
|
||||
"dataset validation"
|
||||
)
|
||||
|
||||
def export_url(self, history_id, data, check_download=True):
|
||||
url = "histories/%s/exports" % history_id
|
||||
put_response = self._put(url, data)
|
||||
|
||||
@@ -0,0 +1,62 @@
|
||||
import tempfile
|
||||
|
||||
from galaxy.datatypes.data import validate
|
||||
from galaxy.datatypes.registry import example_datatype_registry_for_sample
|
||||
from galaxy.datatypes.sniff import get_test_fname
|
||||
|
||||
datatypes_registry = example_datatype_registry_for_sample()
|
||||
|
||||
|
||||
def test_fastq_validation():
|
||||
_assert_valid("fastqsanger", "1.fastqsanger")
|
||||
_assert_invalid("fastqcssanger", "1.fastqsanger")
|
||||
|
||||
_assert_invalid("fastqsanger", "1.fastqcssanger")
|
||||
_assert_valid("fastqcssanger", "1.fastqcssanger")
|
||||
|
||||
|
||||
def test_bam_validation():
|
||||
_assert_valid("bam", "1.bam")
|
||||
_assert_invalid("bam", "1.qname_sorted.bam")
|
||||
_assert_invalid("bam", "3unsorted.bam")
|
||||
|
||||
_assert_valid("qname_sorted.bam", "1.qname_sorted.bam")
|
||||
|
||||
_assert_invalid("qname_sorted.bam", "3unsorted.bam")
|
||||
_assert_valid("unsorted.bam", "3unsorted.bam")
|
||||
|
||||
|
||||
def test_vcf_validation():
|
||||
_assert_valid("vcf", "1.vcf")
|
||||
_assert_invalid("vcf", _truncate("1.vcf", bytes=30))
|
||||
|
||||
|
||||
def _truncate(file_name, bytes=10):
|
||||
o = tempfile.NamedTemporaryFile(delete=False)
|
||||
with open(get_test_fname(file_name), "rb") as f:
|
||||
contents = f.read()
|
||||
o.write(contents[0:-bytes])
|
||||
return o.name
|
||||
|
||||
|
||||
def _assert_invalid(extension, file_name):
|
||||
validation = _run_validation(extension, file_name)
|
||||
assert validation.state == "invalid", validation
|
||||
|
||||
|
||||
def _assert_valid(extension, file_name):
|
||||
validation = _run_validation(extension, file_name)
|
||||
assert validation.state == "ok", validation
|
||||
|
||||
|
||||
def _run_validation(extension, file_name):
|
||||
datatype = datatypes_registry.datatypes_by_extension[extension]
|
||||
validation = validate(MockDataset(get_test_fname(file_name), datatype))
|
||||
return validation
|
||||
|
||||
|
||||
class MockDataset(object):
|
||||
|
||||
def __init__(self, file_name, datatype):
|
||||
self.file_name = file_name
|
||||
self.datatype = datatype
|
||||
Reference in New Issue
Block a user