Merge pull request #236 from natefoo/clean_job_workdir

[15.05] Fix job resubmission
This commit is contained in:
John Chilton
2015-05-07 12:54:28 -05:00
6 changed files with 188 additions and 73 deletions
+16 -2
View File
@@ -688,11 +688,15 @@ class JobExternalOutputMetadataWrapper( object ):
def get_output_filenames_by_dataset( self, dataset, sa_session ):
if isinstance( dataset, galaxy.model.HistoryDatasetAssociation ):
return sa_session.query( galaxy.model.JobExternalOutputMetadata ) \
.filter_by( job_id=self.job_id, history_dataset_association_id=dataset.id ) \
.filter_by( job_id=self.job_id,
history_dataset_association_id=dataset.id,
is_valid=True ) \
.first() # there should only be one or None
elif isinstance( dataset, galaxy.model.LibraryDatasetDatasetAssociation ):
return sa_session.query( galaxy.model.JobExternalOutputMetadata ) \
.filter_by( job_id=self.job_id, library_dataset_dataset_association_id=dataset.id ) \
.filter_by( job_id=self.job_id,
library_dataset_dataset_association_id=dataset.id,
is_valid=True ) \
.first() # there should only be one or None
return None
@@ -701,6 +705,16 @@ class JobExternalOutputMetadataWrapper( object ):
# need to make different keys for them, since ids can overlap
return "%s_%d" % ( dataset.__class__.__name__, dataset.id )
def invalidate_external_metadata( self, datasets, sa_session ):
for dataset in datasets:
jeom = self.get_output_filenames_by_dataset( dataset, sa_session )
# shouldn't be more than one valid, but you never know
while jeom:
jeom.is_valid = False
sa_session.add( jeom )
sa_session.flush()
jeom = self.get_output_filenames_by_dataset( dataset, sa_session )
def setup_external_metadata( self, datasets, sa_session, exec_dir=None,
tmp_dir=None, dataset_files_path=None,
output_fnames=None, config_root=None,
+44 -7
View File
@@ -750,12 +750,7 @@ class JobWrapper( object ):
# directory to be set before prepare is run, or else premature deletion
# and job recovery fail.
# Create the working dir if necessary
try:
self.app.object_store.create(job, base_dir='job_work', dir_only=True, extra_dir=str(self.job_id))
self.working_directory = self.app.object_store.get_filename(job, base_dir='job_work', dir_only=True, extra_dir=str(self.job_id))
log.debug('(%s) Working directory for job is: %s' % (self.job_id, self.working_directory))
except ObjectInvalid:
raise Exception('Unable to create job working directory, job failure')
self.create_working_directory()
self.dataset_path_rewriter = self._job_dataset_path_rewriter( self.working_directory )
self.output_paths = None
self.output_hdas_and_paths = None
@@ -879,6 +874,41 @@ class JobWrapper( object ):
self.write_version_cmd = None
return self.extra_filenames
def create_working_directory( self ):
job = self.get_job()
try:
self.app.object_store.create(
job, base_dir='job_work', dir_only=True, obj_dir=True )
self.working_directory = self.app.object_store.get_filename(
job, base_dir='job_work', dir_only=True, obj_dir=True )
log.debug( '(%s) Working directory for job is: %s',
self.job_id, self.working_directory )
except ObjectInvalid:
raise Exception( '(%s) Unable to create job working directory',
job.id )
def clear_working_directory( self ):
job = self.get_job()
if not os.path.exists( self.working_directory ):
log.warning( '(%s): Working directory clear requested but %s does '
'not exist',
self.job_id,
self.working_directory )
return
self.app.object_store.create(
job, base_dir='job_work', dir_only=True, obj_dir=True,
extra_dir='_cleared_contents', extra_dir_at_root=True )
base = self.app.object_store.get_filename(
job, base_dir='job_work', dir_only=True, obj_dir=True,
extra_dir='_cleared_contents', extra_dir_at_root=True )
date_str = datetime.datetime.now().strftime( '%Y%m%d-%H%M%S' )
arc_dir = os.path.join( base, date_str )
shutil.move( self.working_directory, arc_dir )
self.create_working_directory()
log.debug( '(%s) Previous working directory moved to %s',
self.job_id, arc_dir )
def default_compute_environment( self, job=None ):
if not job:
job = self.get_job()
@@ -1334,7 +1364,7 @@ class JobWrapper( object ):
galaxy.tools.imp_exp.JobExportHistoryArchiveWrapper( self.job_id ).cleanup_after_job( self.sa_session )
galaxy.tools.imp_exp.JobImportHistoryArchiveWrapper( self.app, self.job_id ).cleanup_after_job()
if delete_files:
self.app.object_store.delete(self.get_job(), base_dir='job_work', entire_dir=True, dir_only=True, extra_dir=str(self.job_id))
self.app.object_store.delete(self.get_job(), base_dir='job_work', entire_dir=True, dir_only=True, obj_dir=True)
except:
log.exception( "Unable to cleanup job %d" % self.job_id )
@@ -1525,6 +1555,13 @@ class JobWrapper( object ):
return ExpressionContext( meta, job_context )
return job_context
def invalidate_external_metadata( self ):
job = self.get_job()
self.external_output_metadata.invalidate_external_metadata( [ output_dataset_assoc.dataset for
output_dataset_assoc in
job.output_datasets + job.output_library_datasets ],
self.sa_session )
def setup_external_metadata( self, exec_dir=None, tmp_dir=None,
dataset_files_path=None, config_root=None,
config_file=None, datatypes_config=None,
@@ -14,53 +14,57 @@ MESSAGES = dict(
def failure(app, job_runner, job_state):
if (getattr(job_state, 'runner_state', None)
and job_state.runner_state in
(job_state.runner_states.WALLTIME_REACHED,
job_state.runner_states.MEMORY_LIMIT_REACHED)):
# Intercept jobs that hit the walltime and have a walltime or
# nonspecific resubmit destination configured
for resubmit in job_state.job_destination.get('resubmit'):
if (resubmit.get('condition', None) and resubmit['condition'] !=
job_state.runner_state):
# There is a resubmit defined for the destination but
# its condition is not for walltime_reached
continue
log.info("(%s/%s) Job will be resubmitted to '%s' because %s at "
"the '%s' destination",
job_state.job_wrapper.job_id,
job_state.job_id,
resubmit['destination'],
MESSAGES[job_state.runner_state],
job_state.job_wrapper.job_destination.id )
# fetch JobDestination for the id or tag
new_destination = app.job_config.get_destination(
resubmit['destination'])
# Resolve dynamic if necessary
new_destination = (job_state.job_wrapper.job_runner_mapper
.cache_job_destination(new_destination))
# Reset job state
job = job_state.job_wrapper.get_job()
if resubmit.get('handler', None):
log.debug('(%s/%s) Job reassigned to handler %s',
job_state.job_wrapper.job_id, job_state.job_id,
resubmit['handler'])
job.set_handler(resubmit['handler'])
job_runner.sa_session.add( job )
# Is this safe to do here?
job_runner.sa_session.flush()
# Cache the destination to prevent rerunning dynamic after
# resubmit
job_state.job_wrapper.job_runner_mapper \
.cached_job_destination = new_destination
job_state.job_wrapper.set_job_destination(new_destination)
# Clear external ID (state change below flushes the change)
job.job_runner_external_id = None
# Allow the UI to query for resubmitted state
if job.params is None:
job.params = {}
job_state.runner_state_handled = True
info = "This job was resubmitted to the queue because %s on its " \
"compute resource." % MESSAGES[job_state.runner_state]
job_runner.mark_as_resubmitted(job_state, info=info)
return
runner_state = getattr(job_state, 'runner_state', None)
if (not runner_state
or runner_state not in (job_state.runner_states.WALLTIME_REACHED,
job_state.runner_states.MEMORY_LIMIT_REACHED)):
# not set or not a handleable runner state
return
# Intercept jobs that hit the walltime and have a walltime or
# nonspecific resubmit destination configured
for resubmit in job_state.job_destination.get('resubmit'):
condition = resubmit.get('condition', None)
if condition and condition != runner_state:
# There is a resubmit defined for the destination but
# its condition is not for the encountered state
continue
log.info("(%s/%s) Job will be resubmitted to '%s' because %s at "
"the '%s' destination",
job_state.job_wrapper.job_id,
job_state.job_id,
resubmit['destination'],
MESSAGES[job_state.runner_state],
job_state.job_wrapper.job_destination.id )
# fetch JobDestination for the id or tag
new_destination = app.job_config.get_destination(
resubmit['destination'])
# Resolve dynamic if necessary
new_destination = (job_state.job_wrapper.job_runner_mapper
.cache_job_destination(new_destination))
# Reset job state
job_state.job_wrapper.clear_working_directory()
job_state.job_wrapper.invalidate_external_metadata()
job = job_state.job_wrapper.get_job()
if resubmit.get('handler', None):
log.debug('(%s/%s) Job reassigned to handler %s',
job_state.job_wrapper.job_id, job_state.job_id,
resubmit['handler'])
job.set_handler(resubmit['handler'])
job_runner.sa_session.add( job )
# Is this safe to do here?
job_runner.sa_session.flush()
# Cache the destination to prevent rerunning dynamic after
# resubmit
job_state.job_wrapper.job_runner_mapper \
.cached_job_destination = new_destination
job_state.job_wrapper.set_job_destination(new_destination)
# Clear external ID (state change below flushes the change)
job.job_runner_external_id = None
# Allow the UI to query for resubmitted state
if job.params is None:
job.params = {}
job_state.runner_state_handled = True
info = "This job was resubmitted to the queue because %s on its " \
"compute resource." % MESSAGES[job_state.runner_state]
job_runner.mark_as_resubmitted(job_state, info=info)
return
+1
View File
@@ -485,6 +485,7 @@ model.JobExternalOutputMetadata.table = Table( "job_external_output_metadata", m
Column( "job_id", Integer, ForeignKey( "job.id" ), index=True ),
Column( "history_dataset_association_id", Integer, ForeignKey( "history_dataset_association.id" ), index=True, nullable=True ),
Column( "library_dataset_dataset_association_id", Integer, ForeignKey( "library_dataset_dataset_association.id" ), index=True, nullable=True ),
Column( "is_valid", Boolean, default=True ),
Column( "filename_in", String( 255 ) ),
Column( "filename_out", String( 255 ) ),
Column( "filename_results_code", String( 255 ) ),
@@ -0,0 +1,50 @@
"""
Migration script to allow invalidation of job external output metadata temp files
"""
from sqlalchemy import *
from sqlalchemy.orm import *
from migrate import *
from migrate.changeset import *
from galaxy.model.custom_types import *
import datetime
now = datetime.datetime.utcnow
import logging
log = logging.getLogger( __name__ )
metadata = MetaData()
def upgrade(migrate_engine):
metadata.bind = migrate_engine
print __doc__
metadata.reflect()
isvalid_column = Column( "is_valid", Boolean, default=True )
__add_column( isvalid_column, "job_external_output_metadata", metadata )
def downgrade(migrate_engine):
metadata.bind = migrate_engine
metadata.reflect()
__drop_column( isvalid_column, "job_external_output_metadata", metadata )
def __add_column(column, table_name, metadata, **kwds):
try:
table = Table( table_name, metadata, autoload=True )
column.create( table, **kwds )
except Exception as e:
print str(e)
log.exception( "Adding column %s failed." % column)
def __drop_column( column_name, table_name, metadata ):
try:
table = Table( table_name, metadata, autoload=True )
getattr( table.c, column_name ).drop()
except Exception as e:
print str(e)
log.exception( "Dropping column %s failed." % column_name )
+23 -14
View File
@@ -74,15 +74,19 @@ class ObjectStore(object):
:type alt_name: string
:param alt_name: Use this name as the alternative name for the created
dataset rather than the default.
:type obj_dir: bool
:param obj_dir: Append a subdirectory named with the object's ID (e.g.
000/obj.id)
"""
raise NotImplementedError()
def file_ready(self, obj, base_dir=None, dir_only=False, extra_dir=None, extra_dir_at_root=False, alt_name=None):
def file_ready(self, obj, base_dir=None, dir_only=False, extra_dir=None, extra_dir_at_root=False, alt_name=None, obj_dir=False):
""" A helper method that checks if a file corresponding to a dataset
is ready and available to be used. Return True if so, False otherwise."""
return True
def create(self, obj, base_dir=None, dir_only=False, extra_dir=None, extra_dir_at_root=False, alt_name=None):
def create(self, obj, base_dir=None, dir_only=False, extra_dir=None, extra_dir_at_root=False, alt_name=None, obj_dir=False):
"""
Mark the object identified by `obj` as existing in the store, but with
no content. This method will create a proper directory structure for
@@ -91,7 +95,7 @@ class ObjectStore(object):
"""
raise NotImplementedError()
def empty(self, obj, base_dir=None, extra_dir=None, extra_dir_at_root=False, alt_name=None):
def empty(self, obj, base_dir=None, extra_dir=None, extra_dir_at_root=False, alt_name=None, obj_dir=False):
"""
Test if the object identified by `obj` has content.
If the object does not exist raises `ObjectNotFound`.
@@ -99,7 +103,7 @@ class ObjectStore(object):
"""
raise NotImplementedError()
def size(self, obj, extra_dir=None, extra_dir_at_root=False, alt_name=None):
def size(self, obj, extra_dir=None, extra_dir_at_root=False, alt_name=None, obj_dir=False):
"""
Return size of the object identified by `obj`.
If the object does not exist, return 0.
@@ -107,7 +111,7 @@ class ObjectStore(object):
"""
raise NotImplementedError()
def delete(self, obj, entire_dir=False, base_dir=None, extra_dir=None, extra_dir_at_root=False, alt_name=None):
def delete(self, obj, entire_dir=False, base_dir=None, extra_dir=None, extra_dir_at_root=False, alt_name=None, obj_dir=False):
"""
Deletes the object identified by `obj`.
See `exists` method for the description of other fields.
@@ -115,11 +119,12 @@ class ObjectStore(object):
:type entire_dir: bool
:param entire_dir: If True, delete the entire directory pointed to by
extra_dir. For safety reasons, this option applies
only for and in conjunction with the extra_dir option.
only for and in conjunction with the extra_dir or
obj_dir options.
"""
raise NotImplementedError()
def get_data(self, obj, start=0, count=-1, base_dir=None, extra_dir=None, extra_dir_at_root=False, alt_name=None):
def get_data(self, obj, start=0, count=-1, base_dir=None, extra_dir=None, extra_dir_at_root=False, alt_name=None, obj_dir=False):
"""
Fetch `count` bytes of data starting at offset `start` from the
object identified uniquely by `obj`.
@@ -134,7 +139,7 @@ class ObjectStore(object):
"""
raise NotImplementedError()
def get_filename(self, obj, base_dir=None, dir_only=False, extra_dir=None, extra_dir_at_root=False, alt_name=None):
def get_filename(self, obj, base_dir=None, dir_only=False, extra_dir=None, extra_dir_at_root=False, alt_name=None, obj_dir=False):
"""
Get the expected filename (including the absolute path) which can be used
to access the contents of the object uniquely identified by `obj`.
@@ -142,7 +147,7 @@ class ObjectStore(object):
"""
raise NotImplementedError()
def update_from_file(self, obj, base_dir=None, extra_dir=None, extra_dir_at_root=False, alt_name=None, file_name=None, create=False):
def update_from_file(self, obj, base_dir=None, extra_dir=None, extra_dir_at_root=False, alt_name=None, obj_dir=False, file_name=None, create=False):
"""
Inform the store that the file associated with the object has been
updated. If `file_name` is provided, update from that file instead
@@ -159,7 +164,7 @@ class ObjectStore(object):
"""
raise NotImplementedError()
def get_object_url(self, obj, extra_dir=None, extra_dir_at_root=False, alt_name=None):
def get_object_url(self, obj, extra_dir=None, extra_dir_at_root=False, alt_name=None, obj_dir=False):
"""
If the store supports direct URL access, return a URL. Otherwise return
None.
@@ -207,18 +212,18 @@ class DiskObjectStore(ObjectStore):
if extra_dirs is not None:
self.extra_dirs.update( extra_dirs )
def _get_filename(self, obj, base_dir=None, dir_only=False, extra_dir=None, extra_dir_at_root=False, alt_name=None):
def _get_filename(self, obj, base_dir=None, dir_only=False, extra_dir=None, extra_dir_at_root=False, alt_name=None, obj_dir=False):
"""Class method that returns the absolute path for the file corresponding
to the `obj`.id regardless of whether the file exists.
"""
path = self._construct_path(obj, base_dir=base_dir, dir_only=dir_only, extra_dir=extra_dir, extra_dir_at_root=extra_dir_at_root, alt_name=alt_name, old_style=True)
path = self._construct_path(obj, base_dir=base_dir, dir_only=dir_only, extra_dir=extra_dir, extra_dir_at_root=extra_dir_at_root, alt_name=alt_name, obj_dir=False, old_style=True)
# For backward compatibility, check the old style root path first; otherwise,
# construct hashed path
if not os.path.exists(path):
return self._construct_path(obj, base_dir=base_dir, dir_only=dir_only, extra_dir=extra_dir, extra_dir_at_root=extra_dir_at_root, alt_name=alt_name)
# TODO: rename to _disk_path or something like that to avoid conflicts with children that'll use the local_extra_dirs decorator, e.g. S3
def _construct_path(self, obj, old_style=False, base_dir=None, dir_only=False, extra_dir=None, extra_dir_at_root=False, alt_name=None, **kwargs):
def _construct_path(self, obj, old_style=False, base_dir=None, dir_only=False, extra_dir=None, extra_dir_at_root=False, alt_name=None, obj_dir=False, **kwargs):
""" Construct the expected absolute path for accessing the object
identified by `obj`.id.
@@ -256,6 +261,9 @@ class DiskObjectStore(ObjectStore):
else:
# Construct hashed path
rel_path = os.path.join(*directory_hash_id(obj.id))
# Create a subdirectory for the object ID
if obj_dir:
rel_path = os.path.join(rel_path, str(obj.id))
# Optionally append extra_dir
if extra_dir is not None:
if extra_dir_at_root:
@@ -304,8 +312,9 @@ class DiskObjectStore(ObjectStore):
def delete(self, obj, entire_dir=False, **kwargs):
path = self.get_filename(obj, **kwargs)
extra_dir = kwargs.get('extra_dir', None)
obj_dir = kwargs.get('obj_dir', False)
try:
if entire_dir and extra_dir:
if entire_dir and (extra_dir or obj_dir):
shutil.rmtree(path)
return True
if self.exists(obj, **kwargs):