Merge pull request #5689 from natefoo/upload-memfix

Memory and performance fixes for upload tool
This commit is contained in:
Martin Cech
2018-03-19 12:36:57 -04:00
committed by GitHub
21 changed files with 401 additions and 335 deletions
+1
View File
@@ -74,6 +74,7 @@
<datatype extension="bowtie_base_index" type="galaxy.datatypes.ngsindex:BowtieBaseIndex" mimetype="text/html" display_in_upload="false"/>
<datatype extension="csfasta" type="galaxy.datatypes.sequence:csFasta" display_in_upload="true"/>
<datatype extension="data" type="galaxy.datatypes.data:Data" mimetype="application/octet-stream" max_optional_metadata_filesize="1048576" />
<datatype extension="binary" type="galaxy.datatypes.binary:Binary" mimetype="application/octet-stream" max_optional_metadata_filesize="1048576" />
<datatype extension="d3_hierarchy" type="galaxy.datatypes.text:Json" mimetype="application/json" subclass="true" display_in_upload="true"/>
<datatype extension="data_manager_json" type="galaxy.datatypes.text:Json" mimetype="application/json" subclass="true" display_in_upload="false"/>
<datatype extension="dbn" type="galaxy.datatypes.sequence:DotBracket" display_in_upload="true" description="Dot-Bracket format is a text-based format for storing both an RNA sequence and its corresponding 2D structure." description_url="https://wiki.galaxyproject.org/Learn/Datatypes#Dbn"/>
+2 -2
View File
@@ -188,7 +188,7 @@ class GenericAsn1Binary(Binary):
edam_data = "data_0849"
class BamNative(Binary):
class BamNative(CompressedArchive):
"""Class describing a BAM binary file that is not necessarily sorted"""
edam_format = "format_2572"
edam_data = "data_0863"
@@ -592,7 +592,7 @@ class CRAM(Binary):
return False
class BaseBcf(Binary):
class BaseBcf(CompressedArchive):
edam_format = "format_3020"
edam_data = "data_3498"
+11
View File
@@ -116,6 +116,15 @@ class Data(object):
self.composite_files = self.composite_files.copy()
self.display_applications = odict()
@property
def validate_mode(self):
"""Indicate that a sniffer should run, even if disabled.
Some sniffers (e.g. fastq.gz) work but are not enabled for certain reasons, but when running a sniffer to
"validate" a selected filetype, those sniffers should be enabled.
"""
return os.environ.get('GALAXY_SNIFFER_VALIDATE_MODE', '0') == '1'
def get_raw_data(self, dataset):
"""Returns the full data. To stream it open the file_name and read/write as needed"""
try:
@@ -752,6 +761,8 @@ class Text(Data):
file_ext = 'txt'
line_class = 'line'
is_binary = False
# Add metadata elements
MetadataElement(name="data_lines", default=0, desc="Number of data lines", readonly=True, optional=True, visible=False, no_value=0)
-1
View File
@@ -606,7 +606,6 @@ class RexpBase(Html):
MetadataElement(name="pheno_path", desc="Path to phenotype data for this experiment", default="rexpression.pheno", visible=True)
file_ext = 'rexpbase'
html_table = None
is_binary = True
composite_type = 'auto_primary_file'
allow_datatype_change = False
-1
View File
@@ -18,7 +18,6 @@ class BowtieIndex(Html):
MetadataElement(name="base_name", desc="base name for this index set", default='galaxy_generated_bowtie_index', set_in_upload=True, readonly=True)
MetadataElement(name="sequence_space", desc="sequence_space for this index set", default='unknown', set_in_upload=True, readonly=True)
is_binary = True
composite_type = 'auto_primary_file'
allow_datatype_change = False
+10 -7
View File
@@ -16,7 +16,10 @@ import bx.align.maf
from galaxy import util
from galaxy.datatypes import metadata
from galaxy.datatypes.binary import Binary
from galaxy.datatypes.binary import (
Binary,
CompressedArchive
)
from galaxy.datatypes.metadata import MetadataElement
from galaxy.datatypes.sniff import (
get_headers,
@@ -311,7 +314,7 @@ class Alignment(data.Text):
raise NotImplementedError("Can't split generic alignment files")
class FastaGz(Sequence, Binary):
class FastaGz(Sequence, CompressedArchive):
"""Class representing a generic compressed FASTA sequence"""
edam_format = "format_1929"
file_ext = "fasta.gz"
@@ -319,7 +322,7 @@ class FastaGz(Sequence, Binary):
def sniff(self, filename):
"""Determines whether the file is in gzip-compressed FASTA format"""
if not SNIFF_COMPRESSED_FASTAS:
if not SNIFF_COMPRESSED_FASTAS and not self.validate_mode:
return False
if not is_gzip(filename):
return False
@@ -738,7 +741,7 @@ class FastqCSSanger(Fastq):
file_ext = "fastqcssanger"
class FastqGz(BaseFastq, Binary):
class FastqGz(BaseFastq, CompressedArchive):
"""Class representing a generic compressed FASTQ sequence"""
edam_format = "format_1930"
file_ext = "fastq.gz"
@@ -746,7 +749,7 @@ class FastqGz(BaseFastq, Binary):
def sniff(self, filename):
"""Determines whether the file is in gzip-compressed FASTQ format"""
if not SNIFF_COMPRESSED_FASTQS:
if not SNIFF_COMPRESSED_FASTQS and not self.validate_mode:
return False
if not is_gzip(filename):
return False
@@ -776,7 +779,7 @@ class FastqCSSangerGz(FastqGz):
file_ext = "fastqcssanger.gz"
class FastqBz2(BaseFastq, Binary):
class FastqBz2(BaseFastq, CompressedArchive):
"""Class representing a generic compressed FASTQ sequence"""
edam_format = "format_1930"
file_ext = "fastq.bz2"
@@ -784,7 +787,7 @@ class FastqBz2(BaseFastq, Binary):
def sniff(self, filename):
"""Determine whether the file is in bzip2-compressed FASTQ format"""
if not SNIFF_COMPRESSED_FASTQS:
if not SNIFF_COMPRESSED_FASTQS and not self.validate_mode:
return False
if not is_bz2(filename):
return False
+122 -46
View File
@@ -14,15 +14,18 @@ import tempfile
import zipfile
from six import text_type
from six.moves import filter
from six.moves.urllib.request import urlopen
from galaxy import util
from galaxy.util import compression_utils
from galaxy.util.checkers import (
check_binary,
check_bz2,
check_gzip,
check_html,
is_bz2,
is_gzip
check_zip,
is_tar,
)
if sys.version_info < (3, 3):
@@ -278,7 +281,7 @@ def is_column_based(fname, sep='\t', skip=0):
return True
def guess_ext(fname, sniff_order):
def guess_ext(fname, sniff_order, is_binary=False):
"""
Returns an extension that can be used in the datatype factory to
generate a data for the 'fname' file
@@ -402,7 +405,8 @@ def guess_ext(fname, sniff_order):
successfully discovered.
"""
try:
if datatype.sniff(fname):
if ((is_binary and datatype.is_binary) or
(not is_binary)) and datatype.sniff(fname):
file_ext = datatype.file_ext
break
except Exception:
@@ -416,46 +420,78 @@ def guess_ext(fname, sniff_order):
if file_ext is not None:
return file_ext
# skip header check if data is already known to be binary
if is_binary:
return file_ext or 'binary'
try:
get_headers(fname, None)
except UnicodeDecodeError:
return 'data' # default binary data type file extension
return 'data' # default data type file extension
if is_column_based(fname, '\t', 1):
return 'tabular' # default tabular data type file extension
return 'txt' # default text data type file extension
def handle_compressed_file(filename, datatypes_registry, ext='auto'):
def zip_single_fileobj(path):
z = zipfile.ZipFile(path)
for name in z.namelist():
if not name.endswith('/'):
return z.open(name)
def handle_compressed_file(
filename,
datatypes_registry,
ext='auto',
tmp_prefix='sniff_uncompress_',
tmp_dir=None,
in_place=False,
check_content=True,
auto_decompress=True,
):
"""
Check uploaded files for compression, check compressed file contents, and uncompress if necessary.
Supports GZip, BZip2, and the first file in a Zip file.
For performance reasons, the temporary file used for uncompression is located in the same directory as the
input/output file. This behavior can be changed with the `tmp_dir` param.
``ext`` as returned will only be changed from the ``ext`` input param if the param was an autodetect type (``auto``)
and the file was sniffed as a keep-compressed datatype.
``is_valid`` as returned will only be set if the file is compressed and contains invalid contents (or the first file
in the case of a zip file), this is so lengthy decompression can be bypassed if there is invalid content in the
first 32KB. Otherwise the caller should be checking content.
"""
CHUNK_SIZE = 2 ** 20 # 1Mb
is_compressed = False
compressed_type = None
keep_compressed = False
is_valid = False
uncompressed = filename
tmp_dir = tmp_dir or os.path.dirname(filename)
for compressed_type, check_compressed_function in COMPRESSION_CHECK_FUNCTIONS:
is_compressed = check_compressed_function(filename)
is_compressed, is_valid = check_compressed_function(filename, check_content=check_content)
if is_compressed:
break # found compression type
if is_compressed:
if is_compressed and is_valid:
if ext in AUTO_DETECT_EXTENSIONS:
check_exts = COMPRESSION_DATATYPES[compressed_type]
elif ext in COMPRESSED_EXTENSIONS:
check_exts = [ext]
# attempt to sniff for a keep-compressed datatype (observing the sniff order)
sniff_datatypes = filter(lambda d: getattr(d, 'compressed', False), datatypes_registry.sniff_order)
for datatype in sniff_datatypes:
if datatype.sniff(filename):
ext = datatype.file_ext
keep_compressed = True
break
else:
check_exts = []
for compressed_ext in check_exts:
compressed_datatype = datatypes_registry.get_datatype_by_extension(compressed_ext)
if compressed_datatype.sniff(filename):
ext = compressed_ext
keep_compressed = True
is_valid = True
break
if not is_compressed:
is_valid = True
elif not keep_compressed:
is_valid = True
fd, uncompressed = tempfile.mkstemp()
datatype = datatypes_registry.get_datatype_by_extension(ext)
keep_compressed = getattr(datatype, 'compressed', False)
# don't waste time decompressing if we sniff invalid contents
if is_compressed and is_valid and auto_decompress and not keep_compressed:
fd, uncompressed = tempfile.mkstemp(prefix=tmp_prefix, dir=tmp_dir)
compressed_file = DECOMPRESSION_FUNCTIONS[compressed_type](filename)
# TODO: it'd be ideal to convert to posix newlines and space-to-tab here as well
while True:
try:
chunk = compressed_file.read(CHUNK_SIZE)
@@ -469,35 +505,75 @@ def handle_compressed_file(filename, datatypes_registry, ext='auto'):
os.write(fd, chunk)
os.close(fd)
compressed_file.close()
# Replace the compressed file with the uncompressed file
shutil.move(uncompressed, filename)
return is_valid, ext
if in_place:
# Replace the compressed file with the uncompressed file
shutil.move(uncompressed, filename)
uncompressed = filename
elif not is_compressed:
is_valid = True
return is_valid, ext, uncompressed, compressed_type
def handle_uploaded_dataset_file(filename, datatypes_registry, ext='auto'):
is_valid, ext = handle_compressed_file(filename, datatypes_registry, ext=ext)
def handle_uploaded_dataset_file(
filename,
datatypes_registry,
ext='auto',
tmp_prefix='sniff_upload_',
tmp_dir=None,
in_place=False,
check_content=True,
is_binary=None,
auto_decompress=True,
uploaded_file_ext=None,
convert_to_posix_lines=None,
convert_spaces_to_tabs=None,
):
is_valid, ext, converted_path, compressed_type = handle_compressed_file(
filename,
datatypes_registry,
ext=ext,
tmp_prefix=tmp_prefix,
tmp_dir=tmp_dir,
in_place=in_place,
check_content=check_content,
auto_decompress=auto_decompress,
)
try:
if not is_valid:
if is_tar(converted_path):
raise InappropriateDatasetContentError('TAR file uploads are not supported')
raise InappropriateDatasetContentError('The uploaded compressed file contains invalid content')
if not is_valid:
raise InappropriateDatasetContentError('The compressed uploaded file contains inappropriate content.')
# This needs to be checked again after decompression
is_binary = check_binary(converted_path)
if ext in AUTO_DETECT_EXTENSIONS:
ext = guess_ext(filename, sniff_order=datatypes_registry.sniff_order)
if not is_binary and convert_to_posix_lines:
# Convert universal line endings to Posix line endings, spaces to tabs (if desired)
if convert_spaces_to_tabs:
convert_fxn = convert_newlines_sep2tabs
else:
convert_fxn = convert_newlines
line_count, _converted_path = convert_fxn(converted_path, in_place=in_place, tmp_dir=tmp_dir, tmp_prefix=tmp_prefix)
if not in_place:
if converted_path and filename != converted_path:
os.unlink(converted_path)
converted_path = _converted_path
if check_binary(filename):
if not datatypes_registry.is_extension_unsniffable_binary(ext) and not datatypes_registry.get_datatype_by_extension(ext).sniff(filename):
raise InappropriateDatasetContentError('The binary uploaded file contains inappropriate content.')
elif check_html(filename):
raise InappropriateDatasetContentError('The uploaded file contains inappropriate HTML content.')
return ext
if ext in AUTO_DETECT_EXTENSIONS:
ext = guess_ext(converted_path, sniff_order=datatypes_registry.sniff_order, is_binary=is_binary)
if not is_binary and check_content and check_html(converted_path):
raise InappropriateDatasetContentError('The uploaded file contains invalid HTML content')
except Exception:
if filename != converted_path:
os.unlink(converted_path)
raise
return ext, converted_path, compressed_type
AUTO_DETECT_EXTENSIONS = ['auto'] # should 'data' also cause auto detect?
DECOMPRESSION_FUNCTIONS = dict(gzip=gzip.GzipFile, bz2=bz2.BZ2File)
COMPRESSION_CHECK_FUNCTIONS = [('gzip', is_gzip), ('bz2', is_bz2)]
COMPRESSION_DATATYPES = dict(gzip=['bam', 'fasta.gz', 'fastq.gz', 'fastqsanger.gz', 'fastqillumina.gz', 'fastqsolexa.gz', 'fastqcssanger.gz'], bz2=['fastq.bz2', 'fastqsanger.bz2', 'fastqillumina.bz2', 'fastqsolexa.bz2', 'fastqcssanger.bz2'])
COMPRESSED_EXTENSIONS = []
for exts in COMPRESSION_DATATYPES.values():
COMPRESSED_EXTENSIONS.extend(exts)
DECOMPRESSION_FUNCTIONS = dict(gz=gzip.GzipFile, bz2=bz2.BZ2File, zip=zip_single_fileobj)
COMPRESSION_CHECK_FUNCTIONS = [('gz', check_gzip), ('bz2', check_bz2), ('zip', check_zip)]
class InappropriateDatasetContentError(Exception):
+8 -2
View File
@@ -681,8 +681,14 @@ class BaseVcf(Tabular):
MetadataElement(name="sample_names", default=[], desc="Sample names", readonly=True, visible=False, optional=True, no_value=[])
def sniff(self, filename):
headers = get_headers(filename, '\n', count=1)
return headers[0][0].startswith("##fileformat=VCF")
# Because this sniffer is run on compressed files that might be BGZF (due to the VcfGz subclass), we should
# handle unicode decode errors. This should ultimately be done in get_headers(), but guess_ext() currently
# relies on get_headers() raising this exception.
try:
headers = get_headers(filename, '\n', count=1)
return headers[0][0].startswith("##fileformat=VCF")
except UnicodeDecodeError:
return False
def display_peek(self, dataset):
"""Returns formated html of peek"""
+25 -5
View File
@@ -3,6 +3,7 @@ Support for running a tool in Galaxy via an internal job management system
"""
import copy
import datetime
import errno
import logging
import os
import pwd
@@ -831,6 +832,18 @@ class JobWrapper(HasResourceParameters):
def get_version_string_path(self):
return os.path.abspath(os.path.join(self.app.config.new_file_path, "GALAXY_VERSION_STRING_%s" % self.job_id))
def __prepare_upload_paramfile(self, tool_evaluator):
"""Special case paramfile handling for the upload tool. Moves the paramfile to the working directory
"""
new = os.path.join(self.working_directory, 'upload_params.json')
try:
shutil.move(tool_evaluator.param_dict['paramfile'], new)
except (OSError, IOError) as exc:
# It won't exist at the old path if setup was interrupted and tried again later
if exc.errno != errno.ENOENT or not os.path.exists(new):
raise
tool_evaluator.param_dict['paramfile'] = new
def prepare(self, compute_environment=None):
"""
Prepare the job to run by creating the working directory and the
@@ -855,6 +868,10 @@ class JobWrapper(HasResourceParameters):
self.sa_session.flush()
# TODO: The upload tool actions that create the paramfile can probably be turned in to a configfile to remove this special casing
if job.tool_id == 'upload1':
self.__prepare_upload_paramfile(tool_evaluator)
self.command_line, self.extra_filenames, self.environment_variables = tool_evaluator.build()
# Ensure galaxy_lib_dir is set in case there are any later chdirs
self.galaxy_lib_dir
@@ -947,6 +964,11 @@ class JobWrapper(HasResourceParameters):
)
return tool_evaluator
def _fix_output_permissions(self):
for path in [dp.real_path for dp in self.get_mutable_output_fnames()]:
if os.path.exists(path):
util.umask_fix_perms(path, self.app.config.umask, 0o666, self.app.config.gid)
def fail(self, message, exception=False, stdout="", stderr="", exit_code=None):
"""
Indicate job failure by setting state and message on all output
@@ -1009,6 +1031,7 @@ class JobWrapper(HasResourceParameters):
# the partial files to the object store regardless of whether job.state == DELETED
self.__update_output(job, dataset, clean_only=True)
self._fix_output_permissions()
self._report_error()
# Perform email action even on failure.
for pja in [pjaa.post_job_action for pjaa in job.post_job_actions if pjaa.post_job_action.action_type == "EmailAction"]:
@@ -1273,7 +1296,7 @@ class JobWrapper(HasResourceParameters):
if retry_internally and not self.external_output_metadata.external_metadata_set_successfully(dataset, self.sa_session):
# 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)
extension = sniff.handle_uploaded_dataset_file(dataset.dataset.file_name, self.app.datatypes_registry)[0]
dataset.extension = extension
# call datatype.set_meta directly for the initial set_meta call during dataset creation
@@ -1418,10 +1441,7 @@ class JobWrapper(HasResourceParameters):
# user).
self.sa_session.flush()
# fix permissions
for path in [dp.real_path for dp in self.get_mutable_output_fnames()]:
if os.path.exists(path):
util.umask_fix_perms(path, self.app.config.umask, 0o666, self.app.config.gid)
self._fix_output_permissions()
# Finally set the job state. This should only happen *after* all
# dataset creation, and will allow us to eliminate force_history_refresh.
+3 -1
View File
@@ -433,7 +433,9 @@ class DiskObjectStore(ObjectStore):
if preserve_symlinks and os.path.islink(file_name):
force_symlink(os.readlink(file_name), self.get_filename(obj, **kwargs))
else:
shutil.copy(file_name, self.get_filename(obj, **kwargs))
path = self.get_filename(obj, **kwargs)
shutil.copy(file_name, path)
umask_fix_perms(path, self.config.umask, 0o666)
except IOError as ex:
log.critical('Error copying %s to %s: %s' % (file_name, self._get_filename(obj, **kwargs), ex))
raise ex
+31 -33
View File
@@ -6,7 +6,7 @@ import socket
import subprocess
import tempfile
from cgi import FieldStorage
from json import dumps
from json import dump, dumps
from six import StringIO
from sqlalchemy.orm import eagerload_all
@@ -308,10 +308,8 @@ def create_paramfile(trans, uploaded_datasets):
except Exception as e:
log.warning('Changing ownership of uploaded file %s failed: %s' % (path, str(e)))
# TODO: json_file should go in the working directory
json_file = tempfile.mkstemp()
json_file_path = json_file[1]
json_file = os.fdopen(json_file[0], 'w')
tool_params = []
json_file_path = None
for uploaded_dataset in uploaded_datasets:
data = uploaded_dataset.data
if uploaded_dataset.type == 'composite':
@@ -321,14 +319,14 @@ def create_paramfile(trans, uploaded_datasets):
setattr(data.metadata, meta_name, meta_value)
trans.sa_session.add(data)
trans.sa_session.flush()
json = dict(file_type=uploaded_dataset.file_type,
dataset_id=data.dataset.id,
dbkey=uploaded_dataset.dbkey,
type=uploaded_dataset.type,
metadata=uploaded_dataset.metadata,
primary_file=uploaded_dataset.primary_file,
composite_file_paths=uploaded_dataset.composite_files,
composite_files=dict((k, v.__dict__) for k, v in data.datatype.get_composite_files(data).items()))
params = dict(file_type=uploaded_dataset.file_type,
dataset_id=data.dataset.id,
dbkey=uploaded_dataset.dbkey,
type=uploaded_dataset.type,
metadata=uploaded_dataset.metadata,
primary_file=uploaded_dataset.primary_file,
composite_file_paths=uploaded_dataset.composite_files,
composite_files=dict((k, v.__dict__) for k, v in data.datatype.get_composite_files(data).items()))
else:
try:
is_binary = uploaded_dataset.datatype.is_binary
@@ -352,31 +350,31 @@ def create_paramfile(trans, uploaded_datasets):
user_ftp_dir = None
if user_ftp_dir and uploaded_dataset.path.startswith(user_ftp_dir):
uploaded_dataset.type = 'ftp_import'
json = dict(file_type=uploaded_dataset.file_type,
ext=uploaded_dataset.ext,
name=uploaded_dataset.name,
dataset_id=data.dataset.id,
dbkey=uploaded_dataset.dbkey,
type=uploaded_dataset.type,
is_binary=is_binary,
link_data_only=link_data_only,
uuid=uuid_str,
to_posix_lines=getattr(uploaded_dataset, "to_posix_lines", True),
auto_decompress=getattr(uploaded_dataset, "auto_decompress", True),
purge_source=purge_source,
space_to_tab=uploaded_dataset.space_to_tab,
run_as_real_user=trans.app.config.external_chown_script is not None,
check_content=trans.app.config.check_upload_content,
path=uploaded_dataset.path)
params = dict(file_type=uploaded_dataset.file_type,
ext=uploaded_dataset.ext,
name=uploaded_dataset.name,
dataset_id=data.dataset.id,
dbkey=uploaded_dataset.dbkey,
type=uploaded_dataset.type,
is_binary=is_binary,
link_data_only=link_data_only,
uuid=uuid_str,
to_posix_lines=getattr(uploaded_dataset, "to_posix_lines", True),
auto_decompress=getattr(uploaded_dataset, "auto_decompress", True),
purge_source=purge_source,
space_to_tab=uploaded_dataset.space_to_tab,
run_as_real_user=trans.app.config.external_chown_script is not None,
check_content=trans.app.config.check_upload_content,
path=uploaded_dataset.path)
# TODO: This will have to change when we start bundling inputs.
# Also, in_place above causes the file to be left behind since the
# user cannot remove it unless the parent directory is writable.
if link_data_only == 'copy_files' and trans.app.config.external_chown_script:
_chown(uploaded_dataset.path)
json_file.write(dumps(json) + '\n')
json_file.close()
if trans.app.config.external_chown_script:
_chown(json_file_path)
tool_params.append(params)
with tempfile.NamedTemporaryFile(prefix='upload_params_', delete=False) as fh:
json_file_path = fh.name
dump(tool_params, fh)
return json_file_path
+3 -3
View File
@@ -165,9 +165,9 @@ def looks_like_a_tool_xml(path):
if(checkers.check_binary(full_path) or
checkers.check_image(full_path) or
checkers.check_gzip(full_path)[0] or
checkers.check_bz2(full_path)[0] or
checkers.check_zip(full_path)):
checkers.is_gzip(full_path) or
checkers.is_bz2(full_path) or
checkers.is_zip(full_path)):
return False
with open(path, "r") as f:
+43 -10
View File
@@ -1,9 +1,11 @@
import gzip
import re
import sys
import tarfile
import zipfile
from six import StringIO
from six.moves import filter
from galaxy import util
from galaxy.util.image_util import image_type
@@ -52,19 +54,14 @@ def check_html(file_path, chunk=None):
def check_binary(name, file_path=True):
# Handles files if file_path is True or text if file_path is False
is_binary = False
if file_path:
temp = open(name, "U")
else:
temp = StringIO(name)
try:
for char in temp.read(100):
if util.is_binary(char):
is_binary = True
break
return util.is_binary(temp.read(1024))
finally:
temp.close()
return is_binary
def check_gzip(file_path, check_content=True):
@@ -124,10 +121,23 @@ def check_bz2(file_path, check_content=True):
return (True, True)
def check_zip(file_path):
if zipfile.is_zipfile(file_path):
return True
return False
def check_zip(file_path, check_content=True, files=1):
if not zipfile.is_zipfile(file_path):
return (False, False)
if not check_content:
return (True, True)
CHUNK_SIZE = 2 ** 15 # 32Kb
chunk = None
for filect, member in enumerate(iter_zip(file_path)):
handle, name = member
chunk = handle.read(CHUNK_SIZE)
if chunk and check_html(file_path, chunk):
return (True, False)
if filect >= files:
break
return (True, True)
def is_bz2(file_path):
@@ -140,6 +150,28 @@ def is_gzip(file_path):
return is_gzipped
def is_zip(file_path):
is_zipped, is_valid = check_zip(file_path, check_content=False)
return is_zipped
def is_single_file_zip(file_path):
for i, member in enumerate(iter_zip(file_path)):
if i > 1:
return False
return True
def is_tar(file_path):
return tarfile.is_tarfile(file_path)
def iter_zip(file_path):
with zipfile.ZipFile(file_path) as z:
for f in filter(lambda x: not x.endswith('/'), z.namelist()):
yield (z.open(f), f)
def check_image(file_path):
""" Simple wrapper around image_type to yield a True/False verdict """
if image_type(file_path):
@@ -156,4 +188,5 @@ __all__ = (
'check_zip',
'is_gzip',
'is_bz2',
'is_zip',
)
+1 -1
View File
@@ -44,7 +44,7 @@ def set_meta_with_tool_provided(dataset_instance, file_dict, set_meta_kwds, data
if extension == "_sniff_":
try:
from galaxy.datatypes import sniff
extension = sniff.handle_uploaded_dataset_file(dataset_instance.dataset.external_filename, datatypes_registry)
extension = sniff.handle_uploaded_dataset_file(dataset_instance.dataset.external_filename, datatypes_registry)[0]
# 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)
+1 -1
View File
@@ -301,7 +301,7 @@ def get_repository_file_contents(app, file_path, repository_id, is_admin=False):
return '<br/>gzip compressed file<br/>'
elif checkers.is_bz2(file_path):
return '<br/>bz2 compressed file<br/>'
elif checkers.check_zip(file_path):
elif checkers.is_zip(file_path):
return '<br/>zip compressed file<br/>'
elif checkers.check_binary(file_path):
return '<br/>Binary file<br/>'
+1 -1
View File
@@ -190,7 +190,7 @@ def is_data_index_sample_file(file_path):
return False
if checkers.is_gzip(file_path):
return False
if checkers.check_zip(file_path):
if checkers.is_zip(file_path):
return False
# Default to copying the file if none of the above are true.
return True
@@ -45,5 +45,7 @@
<datatype extension="vcf_bgzip" type="galaxy.datatypes.tabular:VcfGz" display_in_upload="true"/>
<datatype extension="html" type="galaxy.datatypes.text:Html" mimetype="text/html"/>
<datatype extension="data_manager_json" type="galaxy.datatypes.text:Json" mimetype="application/json" subclass="true" display_in_upload="false"/>
<datatype extension="data" type="galaxy.datatypes.data:Data" mimetype="application/octet-stream" max_optional_metadata_filesize="1048576" />
<datatype extension="binary" type="galaxy.datatypes.binary:Binary" mimetype="application/octet-stream" max_optional_metadata_filesize="1048576" />
</registration>
</datatypes>
@@ -246,7 +246,7 @@ class DefaultBinaryContentFiltersTestCase(BaseUploadContentConfigurationTestCase
self.history_id, 'file://%s/random-file' % TEST_DATA_DIRECTORY, file_type="auto", wait=True
)
dataset = self.dataset_populator.get_history_dataset_details(self.history_id, dataset=dataset)
assert dataset["file_ext"] == "data", dataset
assert dataset["file_ext"] == "binary", dataset
def test_gzipped_html_content_blocked_by_default(self):
dataset = self.dataset_populator.new_dataset(
@@ -287,7 +287,7 @@ class AutoDecompressTestCase(BaseUploadContentConfigurationTestCase):
self.history_id, 'file://%s/1.sam.gz' % TEST_DATA_DIRECTORY, file_type="auto", auto_decompress=False, wait=True
)
dataset = self.dataset_populator.get_history_dataset_details(self.history_id, dataset=dataset)
assert dataset["file_ext"] == "data", dataset
assert dataset["file_ext"] == "binary", dataset
def test_auto_decompress_on(self):
dataset = self.dataset_populator.new_dataset(
@@ -97,5 +97,7 @@
<datatype extension="vectorstrip" type="galaxy.datatypes.data:Text" subclass="True"/>
<datatype extension="wobble" type="galaxy.datatypes.data:Text" subclass="True"/>
<datatype extension="wordcount" type="galaxy.datatypes.data:Text" subclass="True"/>
<datatype extension="data" type="galaxy.datatypes.data:Data" mimetype="application/octet-stream" max_optional_metadata_filesize="1048576" />
<datatype extension="binary" type="galaxy.datatypes.binary:Binary" mimetype="application/octet-stream" max_optional_metadata_filesize="1048576" />
</registration>
</datatypes>
+130 -216
View File
@@ -6,53 +6,38 @@
from __future__ import print_function
import errno
import gzip
import os
import shutil
import sys
import tempfile
import zipfile
from json import dumps, loads
from json import dump, load, loads
from six.moves.urllib.request import urlopen
from galaxy import util
from galaxy.datatypes import sniff
from galaxy.datatypes.registry import Registry
from galaxy.datatypes.upload_util import (
handle_sniffable_binary_check,
handle_unsniffable_binary_check,
UploadProblemException,
)
from galaxy.datatypes.upload_util import UploadProblemException
from galaxy.util.checkers import (
check_binary,
check_bz2,
check_gzip,
check_html,
check_zip
is_single_file_zip,
is_zip,
)
if sys.version_info < (3, 3):
import bz2file as bz2
else:
import bz2
assert sys.version_info[:2] >= (2, 7)
def file_err(msg, dataset, json_file):
json_file.write(dumps(dict(type='dataset',
ext='data',
dataset_id=dataset.dataset_id,
stderr=msg,
failed=True)) + "\n")
def file_err(msg, dataset):
# never remove a server-side upload
if dataset.type in ('server_dir', 'path_paste'):
return
try:
os.remove(dataset.path)
except Exception:
pass
if dataset.type not in ('server_dir', 'path_paste'):
try:
os.remove(dataset.path)
except Exception:
pass
return dict(type='dataset',
ext='data',
dataset_id=dataset.dataset_id,
stderr=msg,
failed=True)
def safe_dict(d):
@@ -76,8 +61,9 @@ def parse_outputs(args):
return rval
def add_file(dataset, registry, json_file, output_path):
data_type = None
def add_file(dataset, registry, output_path):
ext = None
compression_type = None
line_count = None
converted_path = None
stdout = None
@@ -115,7 +101,7 @@ def add_file(dataset, registry, json_file, output_path):
# decompressing archive files before sniffing.
auto_decompress = dataset.get('auto_decompress', True)
try:
ext = dataset.file_type
dataset.file_type
except AttributeError:
raise UploadProblemException('Unable to process uploaded file, missing file_type parameter.')
@@ -132,184 +118,90 @@ def add_file(dataset, registry, json_file, output_path):
if not os.path.getsize(dataset.path) > 0:
raise UploadProblemException('The uploaded file is empty')
# Is dataset content supported sniffable binary?
# Does the first 1K contain a null?
is_binary = check_binary(dataset.path)
if is_binary:
data_type, ext = handle_sniffable_binary_check(data_type, ext, dataset.path, registry)
if not data_type:
root_datatype = registry.get_datatype_by_extension(dataset.file_type)
if getattr(root_datatype, 'compressed', False):
data_type = 'compressed archive'
ext = dataset.file_type
# Decompress if needed/desired and determine/validate filetype. If a keep-compressed datatype is explicitly selected
# or if autodetection is selected and the file sniffs as a keep-compressed datatype, it will not be decompressed.
if not link_data_only:
if is_zip(dataset.path) and not is_single_file_zip(dataset.path):
stdout = 'ZIP file contained more than one file, only the first file was added to Galaxy.'
try:
ext, converted_path, compression_type = sniff.handle_uploaded_dataset_file(
dataset.path,
registry,
ext=dataset.file_type,
tmp_prefix='data_id_%s_upload_' % dataset.dataset_id,
tmp_dir=output_adjacent_tmpdir(output_path),
in_place=in_place,
check_content=check_content,
is_binary=is_binary,
auto_decompress=auto_decompress,
uploaded_file_ext=os.path.splitext(dataset.name)[1].lower().lstrip('.'),
convert_to_posix_lines=dataset.to_posix_lines,
convert_spaces_to_tabs=dataset.space_to_tab,
)
except sniff.InappropriateDatasetContentError as exc:
raise UploadProblemException(str(exc))
elif dataset.file_type == 'auto':
# Link mode can't decompress anyway, so enable sniffing for keep-compressed datatypes even when auto_decompress
# is enabled
os.environ['GALAXY_SNIFFER_VALIDATE_MODE'] = '1'
ext = sniff.guess_ext(dataset.path, registry.sniff_order, is_binary=is_binary)
os.environ.pop('GALAXY_SNIFFER_VALIDATE_MODE')
# The converted path will be the same as the input path if no conversion was done (or in-place conversion is used)
converted_path = None if converted_path == dataset.path else converted_path
# Validate datasets where the filetype was explicitly set using the filetype's sniffer (if any)
if dataset.file_type != 'auto':
datatype = registry.get_datatype_by_extension(dataset.file_type)
# Enable sniffer "validate mode" (prevents certain sniffers from disabling themselves)
os.environ['GALAXY_SNIFFER_VALIDATE_MODE'] = '1'
if hasattr(datatype, 'sniff') and not datatype.sniff(dataset.path):
stdout = ("Warning: The file 'Type' was set to '{ext}' but the file does not appear to be of that"
" type".format(ext=dataset.file_type))
os.environ.pop('GALAXY_SNIFFER_VALIDATE_MODE')
# Handle unsniffable binaries
if is_binary and ext == 'binary':
upload_ext = os.path.splitext(dataset.name)[1].lower().lstrip('.')
if registry.is_extension_unsniffable_binary(upload_ext):
stdout = ("Warning: The file's datatype cannot be determined from its contents and was guessed based on"
" its extension, to avoid this warning, manually set the file 'Type' to '{ext}' when uploading"
" this type of file".format(ext=upload_ext))
ext = upload_ext
else:
# See if we have a gzipped file, which, if it passes our restrictions, we'll uncompress
is_gzipped, is_valid = check_gzip(dataset.path, check_content=check_content)
if is_gzipped and not is_valid:
raise UploadProblemException('The gzipped uploaded file contains inappropriate content')
elif is_gzipped and is_valid and auto_decompress:
if not link_data_only:
# We need to uncompress the temp_name file, but BAM files must remain compressed in the BGZF format
CHUNK_SIZE = 2 ** 20 # 1Mb
fd, uncompressed = tempfile.mkstemp(prefix='data_id_%s_upload_gunzip_' % dataset.dataset_id, dir=os.path.dirname(output_path), text=False)
gzipped_file = gzip.GzipFile(dataset.path, 'rb')
while 1:
try:
chunk = gzipped_file.read(CHUNK_SIZE)
except IOError:
os.close(fd)
os.remove(uncompressed)
raise UploadProblemException('Problem decompressing gzipped data')
if not chunk:
break
os.write(fd, chunk)
os.close(fd)
gzipped_file.close()
# Replace the gzipped file with the decompressed file if it's safe to do so
if not in_place:
dataset.path = uncompressed
else:
shutil.move(uncompressed, dataset.path)
os.chmod(dataset.path, 0o644)
dataset.name = dataset.name.rstrip('.gz')
data_type = 'gzip'
if not data_type:
# See if we have a bz2 file, much like gzip
is_bzipped, is_valid = check_bz2(dataset.path, check_content)
if is_bzipped and not is_valid:
raise UploadProblemException('The gzipped uploaded file contains inappropriate content')
elif is_bzipped and is_valid and auto_decompress:
if not link_data_only:
# We need to uncompress the temp_name file
CHUNK_SIZE = 2 ** 20 # 1Mb
fd, uncompressed = tempfile.mkstemp(prefix='data_id_%s_upload_bunzip2_' % dataset.dataset_id, dir=os.path.dirname(output_path), text=False)
bzipped_file = bz2.BZ2File(dataset.path, 'rb')
while 1:
try:
chunk = bzipped_file.read(CHUNK_SIZE)
except IOError:
os.close(fd)
os.remove(uncompressed)
raise UploadProblemException('Problem decompressing bz2 compressed data')
if not chunk:
break
os.write(fd, chunk)
os.close(fd)
bzipped_file.close()
# Replace the bzipped file with the decompressed file if it's safe to do so
if not in_place:
dataset.path = uncompressed
else:
shutil.move(uncompressed, dataset.path)
os.chmod(dataset.path, 0o644)
dataset.name = dataset.name.rstrip('.bz2')
data_type = 'bz2'
if not data_type:
# See if we have a zip archive
is_zipped = check_zip(dataset.path)
if is_zipped and auto_decompress:
if not link_data_only:
CHUNK_SIZE = 2 ** 20 # 1Mb
uncompressed = None
uncompressed_name = None
unzipped = False
z = zipfile.ZipFile(dataset.path)
for name in z.namelist():
if name.endswith('/'):
continue
if unzipped:
stdout = 'ZIP file contained more than one file, only the first file was added to Galaxy.'
break
fd, uncompressed = tempfile.mkstemp(prefix='data_id_%s_upload_zip_' % dataset.dataset_id, dir=os.path.dirname(output_path), text=False)
if sys.version_info[:2] >= (2, 6):
zipped_file = z.open(name)
while 1:
try:
chunk = zipped_file.read(CHUNK_SIZE)
except IOError:
os.close(fd)
os.remove(uncompressed)
raise UploadProblemException('Problem decompressing zipped data')
if not chunk:
break
os.write(fd, chunk)
os.close(fd)
zipped_file.close()
uncompressed_name = name
unzipped = True
else:
# python < 2.5 doesn't have a way to read members in chunks(!)
try:
with open(uncompressed, 'wb') as outfile:
outfile.write(z.read(name))
uncompressed_name = name
unzipped = True
except IOError:
os.close(fd)
os.remove(uncompressed)
raise UploadProblemException('Problem decompressing zipped data')
z.close()
# Replace the zipped file with the decompressed file if it's safe to do so
if uncompressed is not None:
if not in_place:
dataset.path = uncompressed
else:
shutil.move(uncompressed, dataset.path)
os.chmod(dataset.path, 0o644)
dataset.name = uncompressed_name
data_type = 'zip'
if not data_type:
data_type, ext = handle_unsniffable_binary_check(
data_type, ext, dataset.path, dataset.name, is_binary, dataset.file_type, check_content, registry
)
if not data_type:
# We must have a text file
if check_content and check_html(dataset.path):
raise UploadProblemException('The uploaded file contains inappropriate HTML content')
if data_type != 'binary':
if not link_data_only and data_type not in ('gzip', 'bz2', 'zip'):
# Convert universal line endings to Posix line endings if to_posix_lines is True
# and the data is not binary or gzip-, bz2- or zip-compressed.
if dataset.to_posix_lines:
tmpdir = output_adjacent_tmpdir(output_path)
tmp_prefix = 'data_id_%s_convert_' % dataset.dataset_id
if dataset.space_to_tab:
line_count, converted_path = sniff.convert_newlines_sep2tabs(dataset.path, in_place=in_place, tmp_dir=tmpdir, tmp_prefix=tmp_prefix)
else:
line_count, converted_path = sniff.convert_newlines(dataset.path, in_place=in_place, tmp_dir=tmpdir, tmp_prefix=tmp_prefix)
if dataset.file_type == 'auto':
ext = sniff.guess_ext(converted_path or dataset.path, registry.sniff_order)
else:
ext = dataset.file_type
data_type = ext
# Save job info for the framework
if ext == 'auto' and data_type == 'binary':
ext = 'data'
if ext == 'auto' and dataset.ext:
ext = dataset.ext
if ext == 'auto':
ext = 'data'
stdout = ("The uploaded binary file format cannot be determined automatically, please set the file 'Type'"
" manually")
datatype = registry.get_datatype_by_extension(ext)
# Strip compression extension from name
if compression_type and not getattr(datatype, 'compressed', False) and dataset.name.endswith('.' + compression_type):
dataset.name = dataset.name[:-len('.' + compression_type)]
# Move dataset
if link_data_only:
# Never alter a file that will not be copied to Galaxy's local file store.
if datatype.dataset_content_needs_grooming(dataset.path):
err_msg = 'The uploaded files need grooming, so change your <b>Copy data into Galaxy?</b> selection to be ' + \
'<b>Copy files into Galaxy</b> instead of <b>Link to files without copying into Galaxy</b> so grooming can be performed.'
raise UploadProblemException(err_msg)
if not link_data_only and converted_path:
# Move the dataset to its "real" path
try:
shutil.move(converted_path, output_path)
except OSError as e:
# We may not have permission to remove converted_path
if e.errno != errno.EACCES:
raise
elif not link_data_only:
if purge_source:
shutil.move(dataset.path, output_path)
if not link_data_only:
# Move the dataset to its "real" path. converted_path is a tempfile so we move it even if purge_source is False.
if purge_source or converted_path:
try:
shutil.move(converted_path or dataset.path, output_path)
except OSError as e:
# We may not have permission to remove the input
if e.errno != errno.EACCES:
raise
else:
shutil.copy(dataset.path, output_path)
# Write the job info
stdout = stdout or 'uploaded %s file' % data_type
stdout = stdout or 'uploaded %s file' % ext
info = dict(type='dataset',
dataset_id=dataset.dataset_id,
ext=ext,
@@ -318,13 +210,14 @@ def add_file(dataset, registry, json_file, output_path):
line_count=line_count)
if dataset.get('uuid', None) is not None:
info['uuid'] = dataset.get('uuid')
json_file.write(dumps(info) + "\n")
# FIXME: does this belong here? also not output-adjacent-tmpdir aware =/
if not link_data_only and datatype and datatype.dataset_content_needs_grooming(output_path):
# Groom the dataset content if necessary
datatype.groom_dataset_content(output_path)
return info
def add_composite_file(dataset, json_file, output_path, files_path):
def add_composite_file(dataset, output_path, files_path):
if dataset.composite_files:
os.mkdir(files_path)
for name, value in dataset.composite_files.items():
@@ -352,10 +245,33 @@ def add_composite_file(dataset, json_file, output_path, files_path):
# Move the dataset to its "real" path
shutil.move(dataset.primary_file, output_path)
# Write the job info
info = dict(type='dataset',
return dict(type='dataset',
dataset_id=dataset.dataset_id,
stdout='uploaded %s file' % dataset.file_type)
json_file.write(dumps(info) + "\n")
def __read_paramfile(path):
with open(path) as fh:
obj = load(fh)
# If there's a single dataset in an old-style paramfile it'll still parse, but it'll be a dict
assert type(obj) == list
return obj
def __read_old_paramfile(path):
datasets = []
with open(path) as fh:
for line in fh:
datasets.append(loads(line))
return datasets
def __write_job_metadata(metadata):
# TODO: make upload/set_metadata compatible with https://github.com/galaxyproject/galaxy/pull/4437
with open('galaxy.json', 'w') as fh:
for meta in metadata:
dump(meta, fh)
fh.write('\n')
def output_adjacent_tmpdir(output_path):
@@ -373,13 +289,17 @@ def __main__():
sys.exit(1)
output_paths = parse_outputs(sys.argv[4:])
json_file = open('galaxy.json', 'w')
registry = Registry()
registry.load_datatypes(root_dir=sys.argv[1], config=sys.argv[2])
for line in open(sys.argv[3], 'r'):
dataset = loads(line)
try:
datasets = __read_paramfile(sys.argv[3])
except (ValueError, AssertionError):
datasets = __read_old_paramfile(sys.argv[3])
metadata = []
for dataset in datasets:
dataset = util.bunch.Bunch(**safe_dict(dataset))
try:
output_path = output_paths[int(dataset.dataset_id)][0]
@@ -389,18 +309,12 @@ def __main__():
try:
if dataset.type == 'composite':
files_path = output_paths[int(dataset.dataset_id)][1]
add_composite_file(dataset, json_file, output_path, files_path)
metadata.append(add_composite_file(dataset, output_path, files_path))
else:
add_file(dataset, registry, json_file, output_path)
metadata.append(add_file(dataset, registry, output_path))
except UploadProblemException as e:
file_err(e.message, dataset, json_file)
# clean up paramfile
# TODO: this will not work when running as the actual user unless the
# parent directory is writable by the user.
try:
os.remove(sys.argv[3])
except Exception:
pass
metadata.append(file_err(e.message, dataset))
__write_job_metadata(metadata)
if __name__ == '__main__':
+3 -3
View File
@@ -1,12 +1,12 @@
<?xml version="1.0"?>
<tool name="Upload File" id="upload1" version="1.1.5" workflow_compatible="false">
<tool name="Upload File" id="upload1" version="1.1.6" workflow_compatible="false" profile="16.04">
<description>
from your computer
</description>
<action module="galaxy.tools.actions.upload" class="UploadToolAction"/>
<command interpreter="python">
upload.py $GALAXY_ROOT_DIR $GALAXY_DATATYPES_CONF_FILE $paramfile
<command>
python '$__tool_directory__/upload.py' $GALAXY_ROOT_DIR $GALAXY_DATATYPES_CONF_FILE $paramfile
#set $outnum = 0
#while $varExists('output%i' % $outnum):
#set $output = $getVar('output%i' % $outnum)