mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-19 10:51:34 +08:00
Seperate tool and job script standard I/O.
This should allow Galaxy to reason about what the tool reports vs. what dependency resolution and metadata reports, allow us to enable logging during metadata collection, and allow us to reason about the tool's I/O from inside the job script during cleanup/metadata collection at the end.
This commit is contained in:
+16
-14
@@ -987,7 +987,7 @@ class JobWrapper(HasResourceParameters):
|
||||
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):
|
||||
def fail(self, message, exception=False, tool_stdout="", tool_stderr="", exit_code=None, job_stdout=None, job_stderr=None):
|
||||
"""
|
||||
Indicate job failure by setting state and message on all output
|
||||
datasets.
|
||||
@@ -1035,7 +1035,7 @@ class JobWrapper(HasResourceParameters):
|
||||
job.info = message
|
||||
# TODO: Put setting the stdout, stderr, and exit code in one place
|
||||
# (not duplicated with the finish method).
|
||||
job.set_streams(stdout, stderr)
|
||||
job.set_streams(tool_stdout, tool_stderr, job_stdout=job_stdout, job_stderr=job_stderr)
|
||||
# Let the exit code be Null if one is not provided:
|
||||
if (exit_code is not None):
|
||||
job.exit_code = exit_code
|
||||
@@ -1213,9 +1213,11 @@ class JobWrapper(HasResourceParameters):
|
||||
|
||||
def finish(
|
||||
self,
|
||||
stdout,
|
||||
stderr,
|
||||
tool_stdout,
|
||||
tool_stderr,
|
||||
tool_exit_code=None,
|
||||
job_stdout=None,
|
||||
job_stderr=None,
|
||||
check_output_detected_state=None,
|
||||
remote_metadata_directory=None,
|
||||
):
|
||||
@@ -1230,12 +1232,15 @@ class JobWrapper(HasResourceParameters):
|
||||
self.sa_session.expunge_all()
|
||||
job = self.get_job()
|
||||
|
||||
def fail():
|
||||
return self.fail(job.info, tool_stdout=tool_stdout, tool_stderr=tool_stderr, exit_code=tool_exit_code, job_stdout=job_stdout, job_stderr=job_stderr)
|
||||
|
||||
# TODO: After failing here, consider returning from the function.
|
||||
try:
|
||||
self.reclaim_ownership()
|
||||
except Exception:
|
||||
log.exception('(%s) Failed to change ownership of %s, failing' % (job.id, self.working_directory))
|
||||
return self.fail(job.info, stdout=stdout, stderr=stderr, exit_code=tool_exit_code)
|
||||
return fail()
|
||||
|
||||
# if the job was deleted, don't finish it
|
||||
if job.state == job.states.DELETED or job.state == job.states.ERROR:
|
||||
@@ -1244,7 +1249,7 @@ class JobWrapper(HasResourceParameters):
|
||||
# was deleted by an administrator (based on old comments), but it
|
||||
# could also mean that a job was broken up into tasks and one of
|
||||
# the tasks failed. So include the stderr, stdout, and exit code:
|
||||
return self.fail(job.info, stderr=stderr, stdout=stdout, exit_code=tool_exit_code)
|
||||
return fail()
|
||||
|
||||
# We collect the stderr from tools that write their stderr to galaxy.json
|
||||
tool_provided_metadata = self.get_tool_provided_job_metadata()
|
||||
@@ -1256,7 +1261,7 @@ class JobWrapper(HasResourceParameters):
|
||||
# We set final_job_state to use for dataset management, but *don't* set
|
||||
# job.state until after dataset discovery to prevent history issues
|
||||
if check_output_detected_state is None:
|
||||
check_output_detected_state = self.check_tool_output(stdout, stderr, tool_exit_code, job)
|
||||
check_output_detected_state = self.check_tool_output(tool_stdout, tool_stderr, tool_exit_code=tool_exit_code, job=job, job_stdout=job_stdout, job_stderr=job_stderr)
|
||||
|
||||
if check_output_detected_state == DETECTED_JOB_STATE.OK and not tool_provided_metadata.has_failed_outputs():
|
||||
final_job_state = job.states.OK
|
||||
@@ -1406,9 +1411,6 @@ class JobWrapper(HasResourceParameters):
|
||||
# will now be seen by the user.
|
||||
self.sa_session.flush()
|
||||
|
||||
# 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.
|
||||
@@ -1489,16 +1491,16 @@ class JobWrapper(HasResourceParameters):
|
||||
delete_files = cleanup_job == 'always' or (job.state == job.states.OK and cleanup_job == 'onsuccess')
|
||||
self.cleanup(delete_files=delete_files)
|
||||
|
||||
def check_tool_output(self, stdout, stderr, tool_exit_code, job):
|
||||
def check_tool_output(self, tool_stdout, tool_stderr, tool_exit_code, job, job_stdout=None, job_stderr=None):
|
||||
job_id_tag = "<unknown job id>"
|
||||
if job is not None:
|
||||
job_id_tag = job.get_id_tag()
|
||||
|
||||
state, stdout_unicodified, stderr_unicodified, job_messages = check_output(self.tool.stdio_regexes, self.tool.stdio_exit_codes, stdout, stderr, tool_exit_code, job_id_tag)
|
||||
state, tool_stdout, tool_stderr, job_messages = check_output(self.tool.stdio_regexes, self.tool.stdio_exit_codes, tool_stdout, tool_stderr, tool_exit_code, job_id_tag)
|
||||
|
||||
# Store the modified stdout and stderr in the job:
|
||||
if job is not None:
|
||||
job.set_streams(stdout_unicodified, stderr_unicodified, job_messages=job_messages)
|
||||
job.set_streams(tool_stdout, tool_stderr, job_messages=job_messages, job_stdout=job_stdout, job_stderr=job_stderr)
|
||||
|
||||
return state
|
||||
|
||||
@@ -2036,7 +2038,7 @@ class TaskWrapper(JobWrapper):
|
||||
# Check what the tool returned. If the stdout or stderr matched
|
||||
# regular expressions that indicate errors, then set an error.
|
||||
# The same goes if the tool's exit code was in a given range.
|
||||
if self.check_tool_output(stdout, stderr, tool_exit_code, task) == DETECTED_JOB_STATE.OK:
|
||||
if self.check_tool_output(stdout, stderr, tool_exit_code=tool_exit_code, job=task) == DETECTED_JOB_STATE.OK:
|
||||
task.state = task.states.OK
|
||||
else:
|
||||
task.state = task.states.ERROR
|
||||
|
||||
@@ -143,6 +143,7 @@ def __externalize_commands(job_wrapper, shell, commands_builder, remote_command_
|
||||
commands = local_container_script
|
||||
if 'working_directory' in remote_command_params:
|
||||
commands = "%s %s" % (shell, join(remote_command_params['working_directory'], script_name))
|
||||
commands += " > ../tool_stdout 2> ../tool_stderr"
|
||||
log.info("Built script [%s] for tool command [%s]" % (local_container_script, tool_commands))
|
||||
return commands
|
||||
|
||||
|
||||
@@ -361,8 +361,10 @@ class JobHandlerQueue(Monitors):
|
||||
job.text_metrics = copied_from_job.text_metrics
|
||||
job.dependencies = copied_from_job.dependencies
|
||||
job.state = copied_from_job.state
|
||||
job.stderr = copied_from_job.stderr
|
||||
job.stdout = copied_from_job.stdout
|
||||
job.job_stderr = copied_from_job.job_stderr
|
||||
job.job_stdout = copied_from_job.job_stdout
|
||||
job.tool_stderr = copied_from_job.tool_stderr
|
||||
job.tool_stdout = copied_from_job.tool_stdout
|
||||
job.command_line = copied_from_job.command_line
|
||||
job.traceback = copied_from_job.traceback
|
||||
job.tool_version = copied_from_job.tool_version
|
||||
|
||||
@@ -458,11 +458,33 @@ class BaseJobRunner(object):
|
||||
def _job_io_for_db(self, stream):
|
||||
return shrink_stream_by_size(stream, DATABASE_MAX_STRING_SIZE, join_by="\n..\n", left_larger=True, beginning_on_size_error=True)
|
||||
|
||||
def _finish_or_resubmit_job(self, job_state, stdout, stderr, job_id=None, external_job_id=None):
|
||||
def _finish_or_resubmit_job(self, job_state, job_stdout, job_stderr, job_id=None, external_job_id=None):
|
||||
job_wrapper = job_state.job_wrapper
|
||||
try:
|
||||
job = job_state.job_wrapper.get_job()
|
||||
exit_code = job_state.read_exit_code()
|
||||
check_output_detected_state = job_state.job_wrapper.check_tool_output(stdout, stderr, exit_code, job)
|
||||
|
||||
tool_stdout_path = os.path.join(job_wrapper.working_directory, "tool_stdout")
|
||||
tool_stderr_path = os.path.join(job_wrapper.working_directory, "tool_stderr")
|
||||
# TODO: These might not exist for running jobs at the upgrade to 19.XX, remove that
|
||||
# assumption in 20.XX.
|
||||
if os.path.exists(tool_stdout_path):
|
||||
with open(tool_stdout_path, "rb") as stdout_file:
|
||||
tool_stdout = self._job_io_for_db(stdout_file)
|
||||
else:
|
||||
# Legacy job, were getting a merged output - assume it is mostly tool output.
|
||||
tool_stdout = job_stdout
|
||||
job_stdout = None
|
||||
|
||||
if os.path.exists(tool_stderr_path):
|
||||
with open(tool_stderr_path, "rb") as stdout_file:
|
||||
tool_stderr = self._job_io_for_db(stdout_file)
|
||||
else:
|
||||
# Legacy job, were getting a merged output - assume it is mostly tool output.
|
||||
tool_stderr = job_stderr
|
||||
job_stderr = None
|
||||
|
||||
check_output_detected_state = job_wrapper.check_tool_output(tool_stdout, tool_stderr, tool_exit_code=exit_code, job=job, job_stdout=job_stdout, job_stderr=job_stderr)
|
||||
job_not_ok = check_output_detected_state != DETECTED_JOB_STATE.OK
|
||||
|
||||
# clean up the job files
|
||||
@@ -484,10 +506,10 @@ class BaseJobRunner(object):
|
||||
if job_state.runner_state_handled:
|
||||
return
|
||||
|
||||
job_state.job_wrapper.finish(stdout, stderr, exit_code, check_output_detected_state=check_output_detected_state)
|
||||
job_wrapper.finish(tool_stdout, tool_stderr, exit_code, check_output_detected_state=check_output_detected_state, job_stdout=job_stdout, job_stderr=job_stderr)
|
||||
except Exception:
|
||||
log.exception("(%s/%s) Job wrapper finish method failed" % (job_id or '', external_job_id or ''))
|
||||
job_state.job_wrapper.fail("Unable to finish job", exception=True)
|
||||
job_wrapper.fail("Unable to finish job", exception=True)
|
||||
|
||||
|
||||
class JobState(object):
|
||||
|
||||
@@ -534,7 +534,7 @@ class PulsarJobRunner(AsynchronousJobRunner):
|
||||
self._handle_runner_state('failure', job_state)
|
||||
if not job_state.runner_state_handled:
|
||||
job_state.job_wrapper.fail(getattr(job_state, "fail_message", message),
|
||||
stdout=stdout, stderr=stderr, exception=exception)
|
||||
tool_stdout=stdout, tool_stderr=stderr, exception=exception)
|
||||
|
||||
def check_pid(self, pid):
|
||||
try:
|
||||
|
||||
@@ -230,17 +230,26 @@ class JobLike(object):
|
||||
# TODO: Make iterable, concatenate with chain
|
||||
return self.text_metrics + self.numeric_metrics
|
||||
|
||||
def set_streams(self, stdout, stderr, job_messages=None):
|
||||
stdout = galaxy.util.unicodify(stdout) or u''
|
||||
stderr = galaxy.util.unicodify(stderr) or u''
|
||||
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
|
||||
def set_streams(self, tool_stdout, tool_stderr, job_stdout=None, job_stderr=None, job_messages=None):
|
||||
def shrink_and_unicodify(what, stream):
|
||||
stream = galaxy.util.unicodify(stream) or u''
|
||||
if (len(stream) > galaxy.util.DATABASE_MAX_STRING_SIZE):
|
||||
stream = galaxy.util.shrink_string_by_size(tool_stdout, galaxy.util.DATABASE_MAX_STRING_SIZE, join_by="\n..\n", left_larger=True, beginning_on_size_error=True)
|
||||
log.info("%s for %s %d is greater than %s, only a portion will be logged to database", what, type(self), self.id, galaxy.util.DATABASE_MAX_STRING_SIZE_PRETTY)
|
||||
return stream
|
||||
|
||||
self.tool_stdout = shrink_and_unicodify('tool_stdout', tool_stdout)
|
||||
self.tool_stderr = shrink_and_unicodify('tool_stderr', tool_stderr)
|
||||
if job_stdout is not None:
|
||||
self.job_stdout = shrink_and_unicodify('job_stdout', job_stdout)
|
||||
else:
|
||||
self.job_stdout = None
|
||||
|
||||
if job_stderr is not None:
|
||||
self.job_stderr = shrink_and_unicodify('job_stderr', job_stderr)
|
||||
else:
|
||||
self.job_stderr = None
|
||||
|
||||
if job_messages is not None:
|
||||
self.job_messages = job_messages
|
||||
|
||||
@@ -254,6 +263,27 @@ class JobLike(object):
|
||||
|
||||
return "%s[%s,tool_id=%s]" % (self.__class__.__name__, extra, self.tool_id)
|
||||
|
||||
def get_stdout(self):
|
||||
stdout = self.tool_stdout
|
||||
if self.job_stdout:
|
||||
stdout += "\n" + self.job_stdout
|
||||
return stdout
|
||||
|
||||
def set_stdout(self, stdout):
|
||||
raise NotImplementedError("Attempt to set stdout, must set tool_stdout or job_stdout")
|
||||
|
||||
def get_stderr(self):
|
||||
stderr = self.tool_stderr
|
||||
if self.job_stderr:
|
||||
stderr += "\n" + self.job_stderr
|
||||
return stderr
|
||||
|
||||
def set_stderr(self, stderr):
|
||||
raise NotImplementedError("Attempt to set stdout, must set tool_stderr or job_stderr")
|
||||
|
||||
stdout = property(get_stdout, set_stdout)
|
||||
stderr = property(get_stderr, set_stderr)
|
||||
|
||||
|
||||
class User(Dictifiable, RepresentById):
|
||||
use_pbkdf2 = True
|
||||
@@ -983,8 +1013,6 @@ class Task(JobLike, RepresentById):
|
||||
self.task_runner_name = None
|
||||
self.task_runner_external_id = None
|
||||
self.job = job
|
||||
self.stdout = ""
|
||||
self.stderr = ""
|
||||
self.exit_code = None
|
||||
self.prepare_input_files_cmd = prepare_files_cmd
|
||||
self._init_metrics()
|
||||
@@ -1030,12 +1058,6 @@ class Task(JobLike, RepresentById):
|
||||
def get_job(self):
|
||||
return self.job
|
||||
|
||||
def get_stdout(self):
|
||||
return self.stdout
|
||||
|
||||
def get_stderr(self):
|
||||
return self.stderr
|
||||
|
||||
def get_prepare_input_files_cmd(self):
|
||||
return self.prepare_input_files_cmd
|
||||
|
||||
@@ -1112,12 +1134,6 @@ class Task(JobLike, RepresentById):
|
||||
def set_job(self, job):
|
||||
self.job = job
|
||||
|
||||
def set_stdout(self, stdout):
|
||||
self.stdout = stdout
|
||||
|
||||
def set_stderr(self, stderr):
|
||||
self.stderr = stderr
|
||||
|
||||
def set_prepare_input_files_cmd(self, prepare_input_files_cmd):
|
||||
self.prepare_input_files_cmd = prepare_input_files_cmd
|
||||
|
||||
|
||||
@@ -530,8 +530,10 @@ model.Job.table = Table(
|
||||
Column("job_messages", JSONType, nullable=True),
|
||||
Column("param_filename", String(1024)),
|
||||
Column("runner_name", String(255)),
|
||||
Column("stdout", TEXT),
|
||||
Column("stderr", TEXT),
|
||||
Column("job_stdout", TEXT),
|
||||
Column("job_stderr", TEXT),
|
||||
Column("tool_stdout", TEXT),
|
||||
Column("tool_stderr", TEXT),
|
||||
Column("exit_code", Integer, nullable=True),
|
||||
Column("traceback", TEXT),
|
||||
Column("session_id", Integer, ForeignKey("galaxy_session.id"), index=True, nullable=True),
|
||||
@@ -722,8 +724,10 @@ model.Task.table = Table(
|
||||
Column("command_line", TEXT),
|
||||
Column("param_filename", String(1024)),
|
||||
Column("runner_name", String(255)),
|
||||
Column("stdout", TEXT),
|
||||
Column("stderr", TEXT),
|
||||
Column("job_stdout", TEXT), # job_stdout makes sense here because it is short for job script standard out.
|
||||
Column("job_stderr", TEXT),
|
||||
Column("tool_stdout", TEXT),
|
||||
Column("tool_stderr", TEXT),
|
||||
Column("exit_code", Integer, nullable=True),
|
||||
Column("job_messages", JSONType, nullable=True),
|
||||
Column("info", TrimmedString(255)),
|
||||
|
||||
@@ -5,13 +5,17 @@ from __future__ import print_function
|
||||
|
||||
import logging
|
||||
|
||||
from sqlalchemy import Column, MetaData, Table
|
||||
from sqlalchemy import Column, MetaData, Table, TEXT
|
||||
|
||||
from galaxy.model.custom_types import JSONType
|
||||
|
||||
log = logging.getLogger(__name__)
|
||||
job_messages_column = Column("job_messages", JSONType, nullable=True)
|
||||
task_job_messages_column = Column("job_messages", JSONType, nullable=True)
|
||||
job_job_stdout_column = Column("job_stdout", TEXT, nullable=True)
|
||||
job_job_stderr_column = Column("job_stderr", TEXT, nullable=True)
|
||||
task_job_stdout_column = Column("job_stdout", TEXT, nullable=True)
|
||||
task_job_stderr_column = Column("job_stderr", TEXT, nullable=True)
|
||||
|
||||
|
||||
def upgrade(migrate_engine):
|
||||
@@ -24,15 +28,24 @@ def upgrade(migrate_engine):
|
||||
jobs_table = Table("job", metadata, autoload=True)
|
||||
job_messages_column.create(jobs_table)
|
||||
assert job_messages_column is jobs_table.c.job_messages
|
||||
except Exception:
|
||||
log.exception("Adding column 'job_messages' to job table failed.")
|
||||
|
||||
try:
|
||||
job_job_stdout_column.create(jobs_table)
|
||||
job_job_stderr_column.create(jobs_table)
|
||||
|
||||
jobs_table.c.stdout.alter(name="tool_stdout")
|
||||
jobs_table.c.stderr.alter(name="tool_stderr")
|
||||
|
||||
tasks_table = Table("task", metadata, autoload=True)
|
||||
task_job_messages_column.create(tasks_table)
|
||||
assert task_job_messages_column is tasks_table.c.job_messages
|
||||
|
||||
task_job_stdout_column.create(tasks_table)
|
||||
task_job_stderr_column.create(tasks_table)
|
||||
|
||||
tasks_table.c.stdout.alter(name="tool_stdout")
|
||||
tasks_table.c.stderr.alter(name="tool_stderr")
|
||||
except Exception:
|
||||
log.exception("Adding column 'job_messages' to task table failed.")
|
||||
log.exception("Failed to alter jobs and/or tasks tables.")
|
||||
|
||||
|
||||
def downgrade(migrate_engine):
|
||||
@@ -42,13 +55,17 @@ def downgrade(migrate_engine):
|
||||
|
||||
try:
|
||||
jobs_table = Table("job", metadata, autoload=True)
|
||||
job_messages = jobs_table.c.job_messages
|
||||
job_messages.drop()
|
||||
except Exception:
|
||||
log.exception("Dropping 'job_messages' column from job table failed.")
|
||||
try:
|
||||
tasks_table = Table("task", metadata, autoload=True)
|
||||
job_messages = tasks_table.c.job_messages
|
||||
job_messages.drop()
|
||||
|
||||
for colname in ["job_messages", "job_stdout", "job_stderr"]:
|
||||
job_col = getattr(jobs_table.c, colname)
|
||||
job_col.drop()
|
||||
task_col = getattr(tasks_table.c, colname)
|
||||
task_col.drop()
|
||||
|
||||
for table in [jobs_table, tasks_table]:
|
||||
table.c.tool_stdout.alter(name="stdout")
|
||||
table.c.tool_stderr.alter(name="stderr")
|
||||
|
||||
except Exception:
|
||||
log.exception("Dropping 'job_messages' column from task table failed.")
|
||||
log.exception("Failed to alter jobs and/or tasks tables.")
|
||||
|
||||
@@ -188,8 +188,8 @@ class JobImportHistoryArchiveWrapper:
|
||||
imported_job.info = job_attrs.get('info', None)
|
||||
imported_job.exit_code = job_attrs.get('exit_code', None)
|
||||
imported_job.traceback = job_attrs.get('traceback', None)
|
||||
imported_job.stdout = job_attrs.get('stdout', None)
|
||||
imported_job.stderr = job_attrs.get('stderr', None)
|
||||
imported_job.tool_stdout = job_attrs.get('stdout', None)
|
||||
imported_job.tool_stderr = job_attrs.get('stderr', None)
|
||||
imported_job.command_line = job_attrs.get('command_line', None)
|
||||
try:
|
||||
imported_job.create_time = datetime.datetime.strptime(job_attrs["create_time"], "%Y-%m-%dT%H:%M:%S.%f")
|
||||
@@ -242,7 +242,7 @@ class JobImportHistoryArchiveWrapper:
|
||||
if os.path.exists(archive_dir):
|
||||
shutil.rmtree(archive_dir)
|
||||
except Exception as e:
|
||||
jiha.job.stderr += "Error cleaning up history import job: %s" % e
|
||||
jiha.job.tool_stderr += "Error cleaning up history import job: %s" % e
|
||||
self.sa_session.flush()
|
||||
raise
|
||||
|
||||
|
||||
@@ -883,7 +883,7 @@ def _verify_outputs(testdef, history, jobs, tool_id, data_list, data_collection_
|
||||
"stdout": "Standard output of the job",
|
||||
"stderr": "Standard error of the job",
|
||||
}
|
||||
# TODO: Only hack the stdio like this for older profkle, for newer tool profiles
|
||||
# TODO: Only hack the stdio like this for older profile, for newer tool profiles
|
||||
# add some syntax for asserting job messages maybe - or just drop this because exit
|
||||
# code and regex on stdio can be tested directly - so this is really testing Galaxy
|
||||
# core handling more than the tool.
|
||||
|
||||
@@ -134,7 +134,17 @@ class JobController(BaseAPIController, UsesLibraryMixinItems):
|
||||
job_dict = self.encode_all_ids(trans, job.to_dict('element', system_details=is_admin), True)
|
||||
full_output = util.asbool(kwd.get('full', 'false'))
|
||||
if full_output:
|
||||
job_dict.update(dict(stderr=job.stderr, stdout=job.stdout, job_messages=job.job_messages))
|
||||
|
||||
job_dict.update(dict(
|
||||
tool_stdout=job.tool_stdout,
|
||||
tool_stderr=job.tool_stderr,
|
||||
job_stdout=job.job_stdout,
|
||||
job_stderr=job.job_stderr,
|
||||
stderr=job.stderr,
|
||||
stdout=job.stdout,
|
||||
job_messages=job.job_messages
|
||||
))
|
||||
|
||||
if is_admin:
|
||||
if job.user:
|
||||
job_dict['user_email'] = job.user.email
|
||||
|
||||
@@ -1623,7 +1623,7 @@ class AdminGalaxy(controller.JSAppLauncher, AdminActions, UsesQuotaMixin, QuotaP
|
||||
% (stop_msg, self.app.config.get("support_url", "https://galaxyproject.org/support/"))
|
||||
if trans.app.config.track_jobs_in_database:
|
||||
job = trans.sa_session.query(trans.app.model.Job).get(job_id)
|
||||
job.stderr = error_msg
|
||||
job.job_stderr = error_msg
|
||||
job.set_state(trans.app.model.Job.states.DELETED_NEW)
|
||||
trans.sa_session.add(job)
|
||||
else:
|
||||
|
||||
@@ -3,6 +3,8 @@
|
||||
<container type="docker">busybox:ubuntu-14.04</container>
|
||||
</requirements>
|
||||
<command detect_errors="exit_code"><![CDATA[
|
||||
echo 'Writing environment properties to output files.';
|
||||
(>&2 echo 'Example tool stderr output.');
|
||||
echo `id -u` > '$user_id';
|
||||
echo `id -g` > '$group_id';
|
||||
echo `pwd` > '$pwd';
|
||||
|
||||
@@ -0,0 +1,18 @@
|
||||
<?xml version="1.0"?>
|
||||
<job_conf>
|
||||
<plugins>
|
||||
<plugin id="local" type="runner" load="galaxy.jobs.runners.local:LocalJobRunner" workers="2"/>
|
||||
</plugins>
|
||||
|
||||
<handlers>
|
||||
<handler id="main"/>
|
||||
</handlers>
|
||||
|
||||
<destinations>
|
||||
<destination id="local_dest" runner="local">
|
||||
<env exec="echo 'moo std cow'" />
|
||||
<env exec="(>&2 echo 'moo err cow')" />
|
||||
</destination>
|
||||
</destinations>
|
||||
|
||||
</job_conf>
|
||||
@@ -12,6 +12,7 @@ from base.populators import (
|
||||
|
||||
SCRIPT_DIRECTORY = os.path.abspath(os.path.dirname(__file__))
|
||||
SIMPLE_JOB_CONFIG_FILE = os.path.join(SCRIPT_DIRECTORY, "simple_job_conf.xml")
|
||||
IO_INJECTION_JOB_CONFIG_FILE = os.path.join(SCRIPT_DIRECTORY, "io_injection_job_conf.xml")
|
||||
SETS_TMP_DIR_TO_TRUE_JOB_CONFIG = os.path.join(SCRIPT_DIRECTORY, "sets_tmp_dir_to_true_job_conf.xml")
|
||||
SETS_TMP_DIR_AS_EXPRESSION_JOB_CONFIG = os.path.join(SCRIPT_DIRECTORY, "sets_tmp_dir_expression_job_conf.xml")
|
||||
|
||||
@@ -31,6 +32,7 @@ class RunsEnvironmentJobs(object):
|
||||
with self.dataset_populator.test_history() as history_id:
|
||||
self.dataset_populator.run_tool(tool_id, {}, history_id)
|
||||
self.dataset_populator.wait_for_history(history_id, assert_ok=True)
|
||||
self._check_completed_history(history_id)
|
||||
return self._environment_properties(history_id)
|
||||
|
||||
def _environment_properties(self, history_id):
|
||||
@@ -42,6 +44,9 @@ class RunsEnvironmentJobs(object):
|
||||
some_env = self.dataset_populator.get_history_dataset_content(history_id, hid=6).strip()
|
||||
return JobEnviromentProperties(user_id, group_id, pwd, home, tmp, some_env)
|
||||
|
||||
def _check_completed_history(self, history_id):
|
||||
"""Extension point that lets subclasses investigate the completed job."""
|
||||
|
||||
|
||||
class BaseJobEnvironmentIntegrationTestCase(integration_util.IntegrationTestCase, RunsEnvironmentJobs):
|
||||
|
||||
@@ -164,3 +169,28 @@ class SharedHomeJobEnvironmentIntegrationTestCase(BaseJobEnvironmentIntegrationT
|
||||
# shared_home_dir used for newer tools if forced in tool XML
|
||||
job_env = self._run_and_get_environment_properties("job_environment_explicit_shared_home")
|
||||
assert job_env.home == self.shared_home_directory, job_env.home
|
||||
|
||||
|
||||
class JobIOEnvironmentIntegrationTestCase(BaseJobEnvironmentIntegrationTestCase):
|
||||
|
||||
@classmethod
|
||||
def handle_galaxy_config_kwds(cls, config):
|
||||
cls.jobs_directory = tempfile.mkdtemp()
|
||||
config["jobs_directory"] = cls.jobs_directory
|
||||
config["job_config_file"] = IO_INJECTION_JOB_CONFIG_FILE
|
||||
|
||||
@skip_without_tool("job_environment_default")
|
||||
def test_io_separation(self):
|
||||
self._run_and_get_environment_properties()
|
||||
|
||||
def _check_completed_history(self, history_id):
|
||||
jobs = self.dataset_populator.history_jobs(history_id)
|
||||
assert len(jobs) == 1
|
||||
job = jobs[0]
|
||||
job_details = self.dataset_populator.get_job_details(job["id"], full=True)
|
||||
job_details = job_details.json()
|
||||
print(job_details)
|
||||
assert 'moo std cow' in job_details['job_stdout']
|
||||
assert 'moo err cow' in job_details['job_stderr']
|
||||
assert 'Writing environment properties to output files.' in job_details['tool_stdout']
|
||||
assert 'Example tool stderr output.' in job_details['tool_stderr']
|
||||
|
||||
@@ -45,7 +45,7 @@ class TestCommandFactory(TestCase):
|
||||
self.include_work_dir_outputs = False
|
||||
dep_commands = [". /opt/galaxy/tools/bowtie/default/env.sh"]
|
||||
self.job_wrapper.dependency_shell_commands = dep_commands
|
||||
self.__assert_command_is(_surround_command("%s/tool_script.sh; return_code=$?" % self.job_wrapper.working_directory))
|
||||
self.__assert_command_is(_surround_command("%s/tool_script.sh > ../tool_stdout 2> ../tool_stderr; return_code=$?" % self.job_wrapper.working_directory))
|
||||
self.__assert_tool_script_is("#!/bin/sh\n%s; %s" % (dep_commands[0], MOCK_COMMAND_LINE))
|
||||
|
||||
def test_remote_dependency_resolution(self):
|
||||
|
||||
@@ -21,7 +21,7 @@ JOBS_ATTRS = '''[{"info": null, "tool_id": "upload1", "update_time": "2016-02-08
|
||||
def _run_jihaw_cleanup(history_archive, msg):
|
||||
app = MockApp()
|
||||
job = model.Job()
|
||||
job.stderr = ''
|
||||
job.tool_stderr = ''
|
||||
jiha = model.JobImportHistoryArchive(job=job, archive_dir=history_archive.arc_directory)
|
||||
app.model.context.current.add_all([job, jiha])
|
||||
app.model.context.flush()
|
||||
|
||||
Reference in New Issue
Block a user