mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
Merge pull request #1052 from jmchilton/unicode_stream_3
Improved Encoding Handling for Jobs
This commit is contained in:
@@ -37,9 +37,6 @@ from .datasets import DatasetPath
|
||||
|
||||
log = logging.getLogger( __name__ )
|
||||
|
||||
DATABASE_MAX_STRING_SIZE = util.DATABASE_MAX_STRING_SIZE
|
||||
DATABASE_MAX_STRING_SIZE_PRETTY = util.DATABASE_MAX_STRING_SIZE_PRETTY
|
||||
|
||||
# This file, if created in the job's working directory, will be used for
|
||||
# setting advanced metadata properties on the job and its associated outputs.
|
||||
# This interface is currently experimental, is only used by the upload tool,
|
||||
@@ -869,7 +866,7 @@ class JobWrapper( object ):
|
||||
self.dependency_shell_commands = self.tool.build_dependency_shell_commands()
|
||||
# We need command_line persisted to the db in order for Galaxy to re-queue the job
|
||||
# if the server was stopped and restarted before the job finished
|
||||
job.command_line = self.command_line
|
||||
job.command_line = unicodify(self.command_line)
|
||||
self.sa_session.add( job )
|
||||
self.sa_session.flush()
|
||||
# Return list of all extra files
|
||||
@@ -990,18 +987,11 @@ class JobWrapper( object ):
|
||||
self.sa_session.add( dataset )
|
||||
self.sa_session.flush()
|
||||
job.set_final_state( job.states.ERROR )
|
||||
job.command_line = self.command_line
|
||||
job.command_line = unicodify(self.command_line)
|
||||
job.info = message
|
||||
# TODO: Put setting the stdout, stderr, and exit code in one place
|
||||
# (not duplicated with the finish method).
|
||||
if ( len( stdout ) > DATABASE_MAX_STRING_SIZE ):
|
||||
stdout = util.shrink_string_by_size( stdout, DATABASE_MAX_STRING_SIZE, join_by="\n..\n", left_larger=True, beginning_on_size_error=True )
|
||||
log.info( "stdout for job %d is greater than %s, only a portion will be logged to database" % ( job.id, DATABASE_MAX_STRING_SIZE_PRETTY ) )
|
||||
job.stdout = stdout
|
||||
if ( len( stderr ) > DATABASE_MAX_STRING_SIZE ):
|
||||
stderr = util.shrink_string_by_size( stderr, DATABASE_MAX_STRING_SIZE, join_by="\n..\n", left_larger=True, beginning_on_size_error=True )
|
||||
log.info( "stderr for job %d is greater than %s, only a portion will be logged to database" % ( job.id, DATABASE_MAX_STRING_SIZE_PRETTY ) )
|
||||
job.stderr = stderr
|
||||
job.set_streams( stdout, stderr )
|
||||
# Let the exit code be Null if one is not provided:
|
||||
if ( exit_code is not None ):
|
||||
job.exit_code = exit_code
|
||||
@@ -1093,8 +1083,6 @@ class JobWrapper( object ):
|
||||
the contents of the output files.
|
||||
"""
|
||||
finish_timer = util.ExecutionTimer()
|
||||
stdout = unicodify( stdout )
|
||||
stderr = unicodify( stderr )
|
||||
|
||||
# default post job setup
|
||||
self.sa_session.expunge_all()
|
||||
@@ -1272,13 +1260,10 @@ class JobWrapper( object ):
|
||||
# Flush all the dataset and job changes above. Dataset state changes
|
||||
# will now be seen by the user.
|
||||
self.sa_session.flush()
|
||||
# Save stdout and stderr
|
||||
if len( job.stdout ) > DATABASE_MAX_STRING_SIZE:
|
||||
log.info( "stdout for job %d is greater than %s, only a portion will be logged to database" % ( job.id, DATABASE_MAX_STRING_SIZE_PRETTY ) )
|
||||
job.stdout = util.shrink_string_by_size( job.stdout, DATABASE_MAX_STRING_SIZE, join_by="\n..\n", left_larger=True, beginning_on_size_error=True )
|
||||
if len( job.stderr ) > DATABASE_MAX_STRING_SIZE:
|
||||
log.info( "stderr for job %d is greater than %s, only a portion will be logged to database" % ( job.id, DATABASE_MAX_STRING_SIZE_PRETTY ) )
|
||||
job.stderr = util.shrink_string_by_size( job.stderr, DATABASE_MAX_STRING_SIZE, join_by="\n..\n", left_larger=True, beginning_on_size_error=True )
|
||||
|
||||
# Shrink streams and ensure unicode.
|
||||
job.set_streams( job.stdout, job.stderr )
|
||||
|
||||
# The exit code will be null if there is no exit code to be set.
|
||||
# This is so that we don't assign an exit code, such as 0, that
|
||||
# is either incorrect or has the wrong semantics.
|
||||
@@ -1324,7 +1309,7 @@ class JobWrapper( object ):
|
||||
self.tool.call_hook( 'exec_after_process', self.queue.app, inp_data=inp_data,
|
||||
out_data=out_data, param_dict=param_dict,
|
||||
tool=self.tool, stdout=job.stdout, stderr=job.stderr )
|
||||
job.command_line = self.command_line
|
||||
job.command_line = unicodify(self.command_line)
|
||||
|
||||
collected_bytes = 0
|
||||
# Once datasets are collected, set the total dataset size (includes extra files)
|
||||
@@ -1810,8 +1795,6 @@ class TaskWrapper(JobWrapper):
|
||||
the output datasets based on stderr and stdout from the command, and
|
||||
the contents of the output files.
|
||||
"""
|
||||
stdout = unicodify( stdout )
|
||||
stderr = unicodify( stderr )
|
||||
|
||||
# This may have ended too soon
|
||||
log.debug( 'task %s for job %d ended; exit code: %d'
|
||||
@@ -1839,13 +1822,8 @@ class TaskWrapper(JobWrapper):
|
||||
task.state = task.states.ERROR
|
||||
|
||||
# Save stdout and stderr
|
||||
if len( stdout ) > DATABASE_MAX_STRING_SIZE:
|
||||
log.error( "stdout for task %d is greater than %s, only a portion will be logged to database" % ( task.id, DATABASE_MAX_STRING_SIZE_PRETTY ) )
|
||||
task.stdout = util.shrink_string_by_size( stdout, DATABASE_MAX_STRING_SIZE, join_by="\n..\n", left_larger=True, beginning_on_size_error=True )
|
||||
if len( stderr ) > DATABASE_MAX_STRING_SIZE:
|
||||
log.error( "stderr for task %d is greater than %s, only a portion will be logged to database" % ( task.id, DATABASE_MAX_STRING_SIZE_PRETTY ) )
|
||||
task.set_streams( stdout, stderr )
|
||||
self._collect_metrics( task )
|
||||
task.stderr = util.shrink_string_by_size( stderr, DATABASE_MAX_STRING_SIZE, join_by="\n..\n", left_larger=True, beginning_on_size_error=True )
|
||||
task.exit_code = tool_exit_code
|
||||
task.command_line = self.command_line
|
||||
self.sa_session.flush()
|
||||
|
||||
@@ -3,6 +3,8 @@ from os import chmod
|
||||
from os.path import join
|
||||
from os.path import abspath
|
||||
|
||||
from galaxy import util
|
||||
|
||||
from logging import getLogger
|
||||
log = getLogger( __name__ )
|
||||
|
||||
@@ -158,7 +160,8 @@ class CommandsBuilder(object):
|
||||
def __init__(self, initial_command):
|
||||
# Remove trailing semi-colon so we can start hacking up this command.
|
||||
# TODO: Refactor to compose a list and join with ';', would be more clean.
|
||||
commands = initial_command.rstrip("; ")
|
||||
initial_command = util.unicodify(initial_command)
|
||||
commands = initial_command.rstrip(u"; ")
|
||||
self.commands = commands
|
||||
|
||||
# Coping work dir outputs or setting metadata will mask return code of
|
||||
@@ -167,17 +170,19 @@ class CommandsBuilder(object):
|
||||
self.return_code_captured = False
|
||||
|
||||
def prepend_command(self, command):
|
||||
self.commands = "%s; %s" % (command, self.commands)
|
||||
self.commands = u"%s; %s" % (command,
|
||||
self.commands)
|
||||
return self
|
||||
|
||||
def prepend_commands(self, commands):
|
||||
return self.prepend_command("; ".join(commands))
|
||||
return self.prepend_command(u"; ".join(commands))
|
||||
|
||||
def append_command(self, command):
|
||||
self.commands = "%s; %s" % (self.commands, command)
|
||||
self.commands = u"%s; %s" % (self.commands,
|
||||
command)
|
||||
|
||||
def append_commands(self, commands):
|
||||
self.append_command("; ".join(commands))
|
||||
self.append_command(u"; ".join(commands))
|
||||
|
||||
def capture_return_code(self):
|
||||
if not self.return_code_captured:
|
||||
|
||||
@@ -133,9 +133,8 @@ def check_output( tool, stdout, stderr, tool_exit_code, job ):
|
||||
success = True
|
||||
|
||||
# Store the modified stdout and stderr in the job:
|
||||
if None is not job:
|
||||
job.stdout = stdout
|
||||
job.stderr = stderr
|
||||
if job is not None:
|
||||
job.set_streams( stdout, stderr )
|
||||
|
||||
return success
|
||||
|
||||
|
||||
@@ -305,6 +305,8 @@ class BaseJobRunner( object ):
|
||||
|
||||
def write_executable_script( self, path, contents, mode=0o755 ):
|
||||
with open( path, 'w' ) as f:
|
||||
if isinstance(contents, unicode):
|
||||
contents = contents.encode("UTF-8")
|
||||
f.write( contents )
|
||||
os.chmod( path, mode )
|
||||
self._handle_script_integrity( path )
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
from string import Template
|
||||
from pkg_resources import resource_string
|
||||
from galaxy.util import unicodify
|
||||
|
||||
DEFAULT_JOB_FILE_TEMPLATE = Template(
|
||||
resource_string(__name__, 'DEFAULT_JOB_FILE_TEMPLATE.sh').decode('UTF-8')
|
||||
@@ -61,6 +62,8 @@ def job_script(template=DEFAULT_JOB_FILE_TEMPLATE, **kwds):
|
||||
template_params.update(**kwds)
|
||||
env_setup_commands_str = "\n".join(template_params["env_setup_commands"])
|
||||
template_params["env_setup_commands"] = env_setup_commands_str
|
||||
for key, value in template_params.items():
|
||||
template_params[key] = unicodify(value)
|
||||
if not isinstance(template, Template):
|
||||
template = Template(template)
|
||||
return template.safe_substitute(template_params)
|
||||
|
||||
@@ -99,7 +99,7 @@ class HasName:
|
||||
return name
|
||||
|
||||
|
||||
class HasJobMetrics:
|
||||
class JobLike:
|
||||
|
||||
def _init_metrics( self ):
|
||||
self.text_metrics = []
|
||||
@@ -130,6 +130,18 @@ class HasJobMetrics:
|
||||
# TODO: Make iterable, concatenate with chain
|
||||
return self.text_metrics + self.numeric_metrics
|
||||
|
||||
def set_streams( self, stdout, stderr ):
|
||||
stdout = galaxy.util.unicodify( stdout )
|
||||
stderr = galaxy.util.unicodify( stderr )
|
||||
if ( len( stdout ) > galaxy.util.DATABASE_MAX_STRING_SIZE ):
|
||||
stdout = galaxy.util.shrink_string_by_size( stdout, galaxy.util.DATABASE_MAX_STRING_SIZE, join_by="\n..\n", left_larger=True, beginning_on_size_error=True )
|
||||
log.info( "stdout for %s %d is greater than %s, only a portion will be logged to database", type(self), self.id, galaxy.util.DATABASE_MAX_STRING_SIZE_PRETTY )
|
||||
self.stdout = stdout
|
||||
if ( len( stderr ) > galaxy.util.DATABASE_MAX_STRING_SIZE ):
|
||||
stderr = galaxy.util.shrink_string_by_size( stderr, galaxy.util.DATABASE_MAX_STRING_SIZE, join_by="\n..\n", left_larger=True, beginning_on_size_error=True )
|
||||
log.info( "stderr for %s %d is greater than %s, only a portion will be logged to database", type(self), self.id, galaxy.util.DATABASE_MAX_STRING_SIZE_PRETTY )
|
||||
self.stderr = stderr
|
||||
|
||||
|
||||
class User( object, Dictifiable ):
|
||||
use_pbkdf2 = True
|
||||
@@ -307,7 +319,7 @@ class TaskMetricNumeric( BaseJobMetric ):
|
||||
pass
|
||||
|
||||
|
||||
class Job( object, HasJobMetrics, Dictifiable ):
|
||||
class Job( object, JobLike, Dictifiable ):
|
||||
dict_collection_visible_keys = [ 'id', 'state', 'exit_code', 'update_time', 'create_time' ]
|
||||
dict_element_visible_keys = [ 'id', 'state', 'exit_code', 'update_time', 'create_time' ]
|
||||
|
||||
@@ -356,6 +368,7 @@ class Job( object, HasJobMetrics, Dictifiable ):
|
||||
self.destination_id = None
|
||||
self.destination_params = None
|
||||
self.post_job_actions = []
|
||||
self.state_history = []
|
||||
self.imported = False
|
||||
self.handler = None
|
||||
self.exit_code = None
|
||||
@@ -657,7 +670,7 @@ class Job( object, HasJobMetrics, Dictifiable ):
|
||||
self.workflow_invocation_step.update()
|
||||
|
||||
|
||||
class Task( object, HasJobMetrics ):
|
||||
class Task( object, JobLike ):
|
||||
"""
|
||||
A task represents a single component of a job.
|
||||
"""
|
||||
|
||||
@@ -273,9 +273,10 @@ class TextToolParameter( ToolParameter ):
|
||||
def to_string( self, value, app ):
|
||||
"""Convert a value to a string representation suitable for persisting"""
|
||||
if value is None:
|
||||
return ''
|
||||
rval = ''
|
||||
else:
|
||||
return str( value )
|
||||
rval = util.smart_str( value )
|
||||
return rval
|
||||
|
||||
def to_html_value( self, value, app ):
|
||||
if value is None:
|
||||
|
||||
@@ -18,6 +18,7 @@
|
||||
<tool file="dbkey_filter_input.xml" />
|
||||
<tool file="dbkey_filter_multi_input.xml" />
|
||||
<tool file="composite_output.xml" />
|
||||
<tool file="unicode_stream.xml" />
|
||||
<tool file="metadata.xml" />
|
||||
<tool file="metadata_bam.xml" />
|
||||
<tool file="metadata_bcf.xml" />
|
||||
|
||||
@@ -0,0 +1,41 @@
|
||||
<tool id="unicode_stream" name="unicode_stream" version="0.1.0">
|
||||
<description>
|
||||
</description>
|
||||
<configfiles>
|
||||
<configfile name="cf">ვეპხის ტყაოსანი შოთა რუსთაველი
|
||||
</configfile>
|
||||
</configfiles>
|
||||
<command detect_errors="exit_code">
|
||||
echo '$input1' > $out_file1;
|
||||
cat $cf;
|
||||
>&2 cat $cf;
|
||||
sh -c "exit $exit"
|
||||
</command>
|
||||
<inputs>
|
||||
<param name="input1" type="text" label="Input">
|
||||
<sanitizer sanitize="False">
|
||||
</sanitizer>
|
||||
</param>
|
||||
<param name="exit" type="integer" value="0" label="Exit Code" />
|
||||
</inputs>
|
||||
<outputs>
|
||||
<data name="out_file1" format="txt" />
|
||||
</outputs>
|
||||
<tests>
|
||||
<test expect_exit_code="0" expect_failure="false">
|
||||
<param name="input1" value="This is a line of text."/>
|
||||
<param name="exit" value="0" />
|
||||
<output name="out_file1" file="simple_line.txt" />
|
||||
</test>
|
||||
<test expect_exit_code="1" expect_failure="true">
|
||||
<param name="input1" value="This is a line of text."/>
|
||||
<param name="exit" value="1" />
|
||||
</test>
|
||||
<test expect_exit_code="0" expect_failure="false">
|
||||
<param name="input1" value="ვვვვვ"/>
|
||||
<param name="exit" value="0" />
|
||||
</test>
|
||||
</tests>
|
||||
<help>
|
||||
</help>
|
||||
</tool>
|
||||
@@ -3,6 +3,7 @@ from galaxy.util.bunch import Bunch
|
||||
from galaxy.jobs.output_checker import check_output
|
||||
from galaxy.jobs.error_level import StdioErrorLevel
|
||||
from galaxy.tools.parser.interface import ToolStdioRegex
|
||||
from galaxy.model import Job
|
||||
|
||||
|
||||
class OutputCheckerTestCase( TestCase ):
|
||||
@@ -12,11 +13,8 @@ class OutputCheckerTestCase( TestCase ):
|
||||
stdio_regexes=[],
|
||||
stdio_exit_codes=[],
|
||||
)
|
||||
self.job = Bunch(
|
||||
stdout=None,
|
||||
stderr=None,
|
||||
get_id_tag=lambda: "test_id",
|
||||
)
|
||||
self.job = Job()
|
||||
self.job.id = "test_id"
|
||||
self.stdout = ''
|
||||
self.stderr = ''
|
||||
self.tool_exit_code = None
|
||||
|
||||
Reference in New Issue
Block a user