mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
merge and pack scripts
This commit is contained in:
@@ -400,7 +400,6 @@ class JobWrapper( object ):
|
||||
# the state or whether the tool used exit codes and regular
|
||||
# expressions to do so. So we use
|
||||
# job.state == job.states.ERROR to replace this same test.
|
||||
#elif not self.external_output_metadata.external_metadata_set_successfully( dataset, self.sa_session ) and not context['stderr']:
|
||||
elif not self.external_output_metadata.external_metadata_set_successfully( dataset, self.sa_session ) and job.states.ERROR != job.state:
|
||||
dataset._state = model.Dataset.states.FAILED_METADATA
|
||||
else:
|
||||
|
||||
@@ -183,8 +183,9 @@ class LocalJobRunner( BaseJobRunner ):
|
||||
|
||||
def stop_job( self, job ):
|
||||
#if our local job has JobExternalOutputMetadata associated, then our primary job has to have already finished
|
||||
if job.get_external_output_metadata():
|
||||
pid = job.get_external_output_metadata()[0].job_runner_external_pid #every JobExternalOutputMetadata has a pid set, we just need to take from one of them
|
||||
job_ext_output_metadata = job.get_external_output_metadata()
|
||||
if job_ext_output_metadata:
|
||||
pid = job_ext_output_metadata[0].job_runner_external_pid #every JobExternalOutputMetadata has a pid set, we just need to take from one of them
|
||||
else:
|
||||
pid = job.get_job_runner_external_id()
|
||||
if pid in [ None, '' ]:
|
||||
|
||||
@@ -18,7 +18,7 @@ from galaxy.web.form_builder import *
|
||||
from galaxy.model.item_attrs import UsesAnnotations, APIItem
|
||||
from sqlalchemy.orm import object_session
|
||||
from sqlalchemy.sql.expression import func
|
||||
import sys, os.path, os, errno, codecs, operator, socket, pexpect, logging, time, shutil
|
||||
import os.path, os, errno, codecs, operator, socket, pexpect, logging, time, shutil
|
||||
|
||||
if sys.version_info[:2] < ( 2, 5 ):
|
||||
from sets import Set as set
|
||||
@@ -138,6 +138,12 @@ class Job( object ):
|
||||
|
||||
# TODO: Add accessors for members defined in SQL Alchemy for the Job table and
|
||||
# for the mapper defined to the Job table.
|
||||
def get_external_output_metadata( self ):
|
||||
"""
|
||||
The external_output_metadata is currently a reference from Job to
|
||||
JobExternalOutputMetadata. It exists for a job but not a task.
|
||||
"""
|
||||
return self.external_output_metadata
|
||||
def get_session_id( self ):
|
||||
return self.session_id
|
||||
def get_user_id( self ):
|
||||
@@ -370,6 +376,13 @@ class Task( object ):
|
||||
# (e.g., for a session) or never use the member (e.g., external output
|
||||
# metdata). These can be filled in as needed.
|
||||
def get_external_output_metadata( self ):
|
||||
"""
|
||||
The external_output_metadata is currently a backref to
|
||||
JobExternalOutputMetadata. It exists for a job but not a task,
|
||||
and when a task is cancelled its corresponding parent Job will
|
||||
be cancelled. So None is returned now, but that could be changed
|
||||
to self.get_job().get_external_output_metadata().
|
||||
"""
|
||||
return None
|
||||
def get_job_runner_name( self ):
|
||||
"""
|
||||
|
||||
@@ -19,6 +19,14 @@ def create_env_var_dict( elem, tool_dependency_install_dir=None, tool_shed_repos
|
||||
else:
|
||||
env_var_text = elem.text.replace( '$INSTALL_DIR', tool_shed_repository_install_dir )
|
||||
return dict( name=env_var_name, action=env_var_action, value=env_var_text )
|
||||
if elem.text:
|
||||
# Allow for environment variables that contain neither REPOSITORY_INSTALL_DIR nor INSTALL_DIR since there may be command line
|
||||
# parameters that are tuned for a Galaxy instance. Allowing them to be set in one location rather than being hard coded into
|
||||
# each tool config is the best approach. For example:
|
||||
# <environment_variable name="GATK2_SITE_OPTIONS" action="set_to">
|
||||
# "--num_threads 4 --num_cpu_threads_per_data_thread 3 --phone_home STANDARD"
|
||||
# </environment_variable>
|
||||
return dict( name=env_var_name, action=env_var_action, value=elem.text)
|
||||
return None
|
||||
def create_or_update_env_shell_file( install_dir, env_var_dict ):
|
||||
env_var_name = env_var_dict[ 'name' ]
|
||||
|
||||
@@ -1731,7 +1731,7 @@ class Tool:
|
||||
callback( "", input, value[input.name] )
|
||||
else:
|
||||
input.visit_inputs( "", value[input.name], callback )
|
||||
def handle_input( self, trans, incoming, history=None ):
|
||||
def handle_input( self, trans, incoming, history=None, old_errors=None ):
|
||||
"""
|
||||
Process incoming parameters for this tool from the dict `incoming`,
|
||||
update the tool state (or create if none existed), and either return
|
||||
@@ -1766,7 +1766,7 @@ class Tool:
|
||||
else:
|
||||
# Update state for all inputs on the current page taking new
|
||||
# values from `incoming`.
|
||||
errors = self.update_state( trans, self.inputs_by_page[state.page], state.inputs, incoming )
|
||||
errors = self.update_state( trans, self.inputs_by_page[state.page], state.inputs, incoming, old_errors=old_errors or {} )
|
||||
# If the tool provides a `validate_input` hook, call it.
|
||||
validate_input = self.get_hook( 'validate_input' )
|
||||
if validate_input:
|
||||
@@ -1895,7 +1895,10 @@ class Tool:
|
||||
any_group_errors = True
|
||||
# Only need to find one that can't be removed due to size, since only
|
||||
# one removal is processed at # a time anyway
|
||||
break
|
||||
break
|
||||
elif group_old_errors and group_old_errors[i]:
|
||||
group_errors[i] = group_old_errors[i]
|
||||
any_group_errors = True
|
||||
# Update state
|
||||
max_index = -1
|
||||
for i, rep_state in enumerate( group_state ):
|
||||
@@ -1978,6 +1981,8 @@ class Tool:
|
||||
update_only=update_only,
|
||||
old_errors=group_old_errors,
|
||||
item_callback=item_callback )
|
||||
if input.test_param.name in group_old_errors and not test_param_error:
|
||||
test_param_error = group_old_errors[ input.test_param.name ]
|
||||
if test_param_error:
|
||||
group_errors[ input.test_param.name ] = test_param_error
|
||||
if group_errors:
|
||||
|
||||
@@ -94,3 +94,26 @@ def params_from_strings( params, param_values, app, ignore_errors=False ):
|
||||
value = params[key].value_from_basic( value, app, ignore_errors )
|
||||
rval[ key ] = value
|
||||
return rval
|
||||
|
||||
def params_to_incoming( incoming, inputs, input_values, app, name_prefix="" ):
|
||||
"""
|
||||
Given a tool's parameter definition (`inputs`) and a specific set of
|
||||
parameter `input_values` objects, populate `incoming` with the html values.
|
||||
|
||||
Useful for e.g. the rerun function.
|
||||
"""
|
||||
for input in inputs.itervalues():
|
||||
if isinstance( input, Repeat ) or isinstance( input, UploadDataset ):
|
||||
for i, d in enumerate( input_values[ input.name ] ):
|
||||
index = d['__index__']
|
||||
new_name_prefix = name_prefix + "%s_%d|" % ( input.name, index )
|
||||
params_to_incoming( incoming, input.inputs, d, app, new_name_prefix )
|
||||
elif isinstance( input, Conditional ):
|
||||
values = input_values[ input.name ]
|
||||
current = values["__current_case__"]
|
||||
new_name_prefix = name_prefix + input.name + "|"
|
||||
incoming[ new_name_prefix + input.test_param.name ] = values[ input.test_param.name ]
|
||||
params_to_incoming( incoming, input.cases[current].inputs, values, app, new_name_prefix )
|
||||
else:
|
||||
incoming[ name_prefix + input.name ] = input.to_string( input_values.get( input.name ), app )
|
||||
|
||||
|
||||
@@ -1539,25 +1539,38 @@ class DataToolParameter( ToolParameter ):
|
||||
if trans.workflow_building_mode:
|
||||
return None
|
||||
if not value:
|
||||
raise ValueError( "History does not include a dataset of the required format / build" )
|
||||
raise ValueError( "History does not include a dataset of the required format / build" )
|
||||
if value in [None, "None"]:
|
||||
return None
|
||||
if isinstance( value, list ):
|
||||
return [ trans.sa_session.query( trans.app.model.HistoryDatasetAssociation ).get( v ) for v in value ]
|
||||
rval = [ trans.sa_session.query( trans.app.model.HistoryDatasetAssociation ).get( v ) for v in value ]
|
||||
elif isinstance( value, trans.app.model.HistoryDatasetAssociation ):
|
||||
return value
|
||||
rval = value
|
||||
else:
|
||||
return trans.sa_session.query( trans.app.model.HistoryDatasetAssociation ).get( value )
|
||||
rval = trans.sa_session.query( trans.app.model.HistoryDatasetAssociation ).get( value )
|
||||
if isinstance( rval, list ):
|
||||
values = rval
|
||||
else:
|
||||
values = [ rval ]
|
||||
for v in values:
|
||||
if v:
|
||||
if v.deleted:
|
||||
raise ValueError( "The previously selected dataset has been previously deleted" )
|
||||
if v.dataset.state in [galaxy.model.Dataset.states.ERROR, galaxy.model.Dataset.states.DISCARDED ]:
|
||||
raise ValueError( "The previously selected dataset has entered an unusable state" )
|
||||
return rval
|
||||
|
||||
def to_string( self, value, app ):
|
||||
if value is None or isinstance( value, str ):
|
||||
if value is None or isinstance( value, basestring ):
|
||||
return value
|
||||
elif isinstance( value, int ):
|
||||
return str( value )
|
||||
elif isinstance( value, DummyDataset ):
|
||||
return None
|
||||
elif isinstance( value, list) and len(value) > 0 and isinstance( value[0], DummyDataset):
|
||||
return None
|
||||
elif isinstance( value, list ):
|
||||
return ",".join( [ val if isinstance( val, str ) else str(val.id) for val in value] )
|
||||
return ",".join( [ val if isinstance( val, basestring ) else str(val.id) for val in value] )
|
||||
return value.id
|
||||
|
||||
def to_python( self, value, app ):
|
||||
|
||||
@@ -6,6 +6,7 @@ from galaxy.web.base.controller import *
|
||||
from galaxy.util.bunch import Bunch
|
||||
from galaxy.tools import DefaultToolState
|
||||
from galaxy.tools.parameters.basic import UnvalidatedValue
|
||||
from galaxy.tools.parameters import params_to_incoming
|
||||
from galaxy.tools.actions import upload_common
|
||||
|
||||
import logging
|
||||
@@ -192,25 +193,29 @@ class ToolRunner( BaseUIController ):
|
||||
if isinstance(value,list):
|
||||
values = []
|
||||
for val in value:
|
||||
if val not in history.datasets and val in hda_source_dict:
|
||||
if val in history.datasets:
|
||||
values.append( val )
|
||||
elif val in hda_source_dict:
|
||||
values.append( hda_source_dict[ val ])
|
||||
return values
|
||||
if value not in history.datasets and value in hda_source_dict:
|
||||
return hda_source_dict[ value ]
|
||||
visit_input_values( tool.inputs, params_objects, rerun_callback )
|
||||
# Create a fake tool_state for the tool, with the parameters values
|
||||
# Create a fake tool_state for the tool, with the parameters values
|
||||
state = tool.new_state( trans )
|
||||
state.inputs = params_objects
|
||||
tool_state_string = util.object_to_string(state.encode(tool, trans.app))
|
||||
# Setup context for template
|
||||
vars = dict( tool_state=state, errors = upgrade_messages )
|
||||
#create an incoming object from the original job's dataset-modified param objects
|
||||
incoming = {}
|
||||
params_to_incoming( incoming, tool.inputs, params_objects, trans.app )
|
||||
incoming[ "tool_state" ] = util.object_to_string( state.encode( tool, trans.app ) )
|
||||
template, vars = tool.handle_input( trans, incoming, old_errors=upgrade_messages ) #update new state with old parameters
|
||||
# Is the "add frame" stuff neccesary here?
|
||||
add_frame = AddFrameData()
|
||||
add_frame.debug = trans.debug
|
||||
if from_noframe is not None:
|
||||
add_frame.wiki_url = trans.app.config.wiki_url
|
||||
add_frame.from_noframe = True
|
||||
return trans.fill_template( "tool_form.mako",
|
||||
return trans.fill_template( template,
|
||||
history=history,
|
||||
toolbox=self.get_toolbox(),
|
||||
tool_version_select_field=tool_version_select_field,
|
||||
|
||||
Reference in New Issue
Block a user