Merge pull request #3619 from jmchilton/purge_jobs_rework

Improve interaction between job state and purged datasets.
This commit is contained in:
Martin Cech
2017-03-27 16:49:54 -04:00
committed by GitHub
9 changed files with 249 additions and 58 deletions
+51 -19
View File
@@ -1091,9 +1091,7 @@ class JobWrapper( object, HasResourceParameters ):
dataset.mark_unhidden()
if dataset.ext == 'auto':
dataset.extension = 'data'
# Update (non-library) job output datasets through the object store
if dataset not in job.output_library_datasets:
self.app.object_store.update_from_file(dataset.dataset, create=True)
self.__update_output(job, dataset)
# Pause any dependent jobs (and those jobs' outputs)
for dep_job_assoc in dataset.dependent_jobs:
self.pause( dep_job_assoc.job, "Execution of this dataset's job is paused because its input datasets are in an error state." )
@@ -1111,6 +1109,13 @@ class JobWrapper( object, HasResourceParameters ):
self.sa_session.add( job )
self.sa_session.flush()
else:
for dataset_assoc in job.output_datasets:
dataset = dataset_assoc.dataset
# Any reason for clean_only here? We should probably be more consistent and transfer
# the partial files to the object store regardless of whether job.state == DELETED
self.__update_output(job, dataset, clean_only=True)
self._report_error_to_sentry()
# 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"]:
@@ -1298,17 +1303,19 @@ class JobWrapper( object, HasResourceParameters ):
# lets not allow this to occur
# need to update all associated output hdas, i.e. history was shared with job running
for dataset in dataset_assoc.dataset.dataset.history_associations + dataset_assoc.dataset.dataset.library_associations:
trynum = 0
while trynum < self.app.config.retry_job_output_collection:
try:
# Attempt to short circuit NFS attribute caching
os.stat( dataset.dataset.file_name )
os.chown( dataset.dataset.file_name, os.getuid(), -1 )
trynum = self.app.config.retry_job_output_collection
except ( OSError, ObjectNotFound ) as e:
trynum += 1
log.warning( 'Error accessing %s, will retry: %s', dataset.dataset.file_name, e )
time.sleep( 2 )
purged = dataset.dataset.purged
if not purged:
trynum = 0
while trynum < self.app.config.retry_job_output_collection:
try:
# Attempt to short circuit NFS attribute caching
os.stat( dataset.dataset.file_name )
os.chown( dataset.dataset.file_name, os.getuid(), -1 )
trynum = self.app.config.retry_job_output_collection
except ( OSError, ObjectNotFound ) as e:
trynum += 1
log.warning( 'Error accessing %s, will retry: %s', dataset.dataset.file_name, e )
time.sleep( 2 )
if getattr( dataset, "hidden_beneath_collection_instance", None ):
dataset.visible = False
dataset.blurb = 'done'
@@ -1327,7 +1334,9 @@ class JobWrapper( object, HasResourceParameters ):
# Update (non-library) job output datasets through the object store
if dataset not in job.output_library_datasets:
self.app.object_store.update_from_file(dataset.dataset, create=True)
self._collect_extra_files(dataset.dataset, self.working_directory)
self.__update_output(job, dataset)
if not purged:
self._collect_extra_files(dataset.dataset, self.working_directory)
# Handle composite datatypes of auto_primary_file type
if dataset.datatype.composite_type == 'auto_primary_file' and not dataset.has_data():
try:
@@ -1341,7 +1350,7 @@ class JobWrapper( object, HasResourceParameters ):
if job.states.ERROR == final_job_state:
dataset.blurb = "error"
dataset.mark_unhidden()
elif dataset.has_data():
elif not purged and dataset.has_data():
# If the tool was expected to set the extension, attempt to retrieve it
if dataset.ext == 'auto':
dataset.extension = context.get( 'ext', 'data' )
@@ -1485,8 +1494,9 @@ class JobWrapper( object, HasResourceParameters ):
collected_bytes = 0
# Once datasets are collected, set the total dataset size (includes extra files)
for dataset_assoc in job.output_datasets:
dataset_assoc.dataset.dataset.set_total_size()
collected_bytes += dataset_assoc.dataset.dataset.get_total_size()
if not dataset_assoc.dataset.dataset.purged:
dataset_assoc.dataset.dataset.set_total_size()
collected_bytes += dataset_assoc.dataset.dataset.get_total_size()
if job.user:
job.user.adjust_total_disk_usage(collected_bytes)
@@ -1503,7 +1513,8 @@ class JobWrapper( object, HasResourceParameters ):
# fix permissions
for path in [ dp.real_path for dp in self.get_mutable_output_fnames() ]:
util.umask_fix_perms( path, self.app.config.umask, 0o666, self.app.config.gid )
if os.path.exists(path):
util.umask_fix_perms( path, self.app.config.umask, 0o666, self.app.config.gid )
# Finally set the job state. This should only happen *after* all
# dataset creation, and will allow us to eliminate force_history_refresh.
@@ -1800,6 +1811,27 @@ class JobWrapper( object, HasResourceParameters ):
else:
return 'anonymous@unknown'
def __update_output(self, job, dataset, clean_only=False):
"""Handle writing outputs to the object store.
This should be called regardless of whether the job was failed or not so
that writing of partial results happens and so that the object store is
cleaned up if the dataset has been purged.
"""
dataset = dataset.dataset
if dataset not in job.output_library_datasets:
purged = dataset.purged
if not purged and not clean_only:
self.app.object_store.update_from_file(dataset, create=True)
else:
# If the dataset is purged and Galaxy is configured to write directly
# to the object store from jobs - be sure that file is cleaned up. This
# is a bit of hack - our object store abstractions would be stronger
# and more consistent if tools weren't writing there directly.
target = dataset.file_name
if os.path.exists( target ):
os.remove( target )
def __link_file_check( self ):
""" outputs_to_working_directory breaks library uploads where data is
linked. This method is a hack that solves that problem, but is
+7 -14
View File
@@ -300,15 +300,14 @@ class DatasetAssociationManager( base.ModelManager,
# error here if disallowed - before jobs are stopped
# TODO: this check may belong in the controller
self.dataset_manager.error_unless_dataset_purge_allowed()
super( DatasetAssociationManager, self ).purge( dataset_assoc, flush=flush )
# We need to ignore a potential flush=False here and force the flush
# so that job cleanup associated with stop_creating_job will see
# the dataset as purged.
super( DatasetAssociationManager, self ).purge( dataset_assoc, flush=True )
# stop any jobs outputing the dataset_assoc
if dataset_assoc.creating_job_associations:
job = dataset_assoc.creating_job_associations[0].job
if not job.finished:
# signal to stop the creating job
job.mark_deleted( self.app.config.track_jobs_in_database )
self.app.job_manager.job_stop_queue.put( job.id )
self.stop_creating_job( dataset_assoc )
# more importantly, purge underlying dataset as well
if dataset_assoc.dataset.user_can_purge:
@@ -335,16 +334,10 @@ class DatasetAssociationManager( base.ModelManager,
"""
Stops an dataset_assoc's creating job if all the job's other outputs are deleted.
"""
# TODO: use in purge above
RUNNING_STATES = (
self.app.model.Job.states.QUEUED,
self.app.model.Job.states.RUNNING,
self.app.model.Job.states.NEW
)
if dataset_assoc.parent_id is None and len( dataset_assoc.creating_job_associations ) > 0:
# Mark associated job for deletion
job = dataset_assoc.creating_job_associations[0].job
if job.state in RUNNING_STATES:
if not job.finished:
# Are *all* of the job's other output datasets deleted?
if job.check_if_output_datasets_deleted():
job.mark_deleted( self.app.config.track_jobs_in_database )
+145 -10
View File
@@ -1,13 +1,26 @@
import datetime
import json
import os
import time
from operator import itemgetter
from base import api
from base.populators import TestsDatasets
from base.api_asserts import assert_status_code_is_ok
from base.populators import (
DatasetPopulator,
wait_on,
wait_on_state,
)
from requests import put
class JobsApiTestCase( api.ApiTestCase, TestsDatasets ):
class JobsApiTestCase( api.ApiTestCase ):
def setUp( self ):
super( JobsApiTestCase, self ).setUp()
self.dataset_populator = DatasetPopulator( self.galaxy_interactor )
def test_index( self ):
# Create HDA to ensure at least one job exists...
@@ -69,7 +82,7 @@ class JobsApiTestCase( api.ApiTestCase, TestsDatasets ):
jobs = self.__jobs_index( data={"history_id": history_id} )
assert len( jobs ) > 0
history_id = self._new_history()
history_id = self.dataset_populator.new_history()
jobs = self.__jobs_index( data={"history_id": history_id} )
assert len( jobs ) == 0
@@ -117,6 +130,128 @@ class JobsApiTestCase( api.ApiTestCase, TestsDatasets ):
show_jobs_response = self._get( "jobs/%s" % job_id, admin=True )
self._assert_has_keys( show_jobs_response.json(), "command_line", "external_id" )
def test_deleting_output_keep_running_until_all_deleted( self ):
history_id, job_state, outputs = self._setup_running_two_output_job( 60 )
# Delete one of the two outputs and make sure the job is still running.
self._raw_update_history_item( history_id, outputs[0]["id"], {"deleted": True} )
time.sleep( 1 )
state = job_state().json()["state"]
assert state == "running", state
# Delete the second output and make sure the job is cancelled.
self._raw_update_history_item( history_id, outputs[1]["id"], {"deleted": True} )
final_state = wait_on_state( job_state, assert_ok=False, timeout=15 )
assert final_state in ["deleted_new", "deleted"], final_state
def test_purging_output_keep_running_until_all_purged( self ):
history_id, job_state, outputs = self._setup_running_two_output_job( 60 )
# Pretty much right away after the job is running, these paths should be populated -
# if they are grab them and make sure they are deleted at the end of the job.
dataset_1 = self._get_history_item_as_admin( history_id, outputs[0]["id"] )
dataset_2 = self._get_history_item_as_admin( history_id, outputs[1]["id"] )
if "file_name" in dataset_1:
output_dataset_paths = [ dataset_1[ "file_name" ], dataset_2[ "file_name" ] ]
# This may or may not exist depending on if the test is local or not.
output_dataset_paths_exist = os.path.exists( output_dataset_paths[ 0 ] )
else:
output_dataset_paths = []
output_dataset_paths_exist = False
current_state = job_state().json()["state"]
assert current_state == "running", current_state
# Purge one of the two outputs and make sure the job is still running.
self._raw_update_history_item( history_id, outputs[0]["id"], {"purged": True} )
time.sleep( 1 )
current_state = job_state().json()["state"]
assert current_state == "running", current_state
# Purge the second output and make sure the job is cancelled.
self._raw_update_history_item( history_id, outputs[1]["id"], {"purged": True} )
final_state = wait_on_state( job_state, assert_ok=False, timeout=15 )
assert final_state in ["deleted_new", "deleted"], final_state
def paths_deleted():
if not os.path.exists( output_dataset_paths[ 0 ] ) and not os.path.exists( output_dataset_paths[ 1 ] ):
return True
if output_dataset_paths_exist:
wait_on(paths_deleted, "path deletion")
def test_purging_output_cleaned_after_ok_run( self ):
history_id, job_state, outputs = self._setup_running_two_output_job( 10 )
# Pretty much right away after the job is running, these paths should be populated -
# if they are grab them and make sure they are deleted at the end of the job.
dataset_1 = self._get_history_item_as_admin( history_id, outputs[0]["id"] )
dataset_2 = self._get_history_item_as_admin( history_id, outputs[1]["id"] )
if "file_name" in dataset_1:
output_dataset_paths = [ dataset_1[ "file_name" ], dataset_2[ "file_name" ] ]
# This may or may not exist depending on if the test is local or not.
output_dataset_paths_exist = os.path.exists( output_dataset_paths[ 0 ] )
else:
output_dataset_paths = []
output_dataset_paths_exist = False
if not output_dataset_paths_exist:
# Given this Galaxy configuration - there is nothing more to be tested here.
# Consider throwing a skip instead.
return
# Purge one of the two outputs and wait for the job to complete.
self._raw_update_history_item( history_id, outputs[0]["id"], {"purged": True} )
wait_on_state(job_state, assert_ok=True)
if output_dataset_paths_exist:
time.sleep( .5 )
# Make sure the non-purged dataset is on disk and the purged one is not.
assert os.path.exists( output_dataset_paths[ 1 ] )
assert not os.path.exists( output_dataset_paths[ 0 ] )
def _setup_running_two_output_job( self, sleep_time ):
history_id = self.dataset_populator.new_history()
payload = self.dataset_populator.run_tool_payload(
tool_id='create_2',
inputs=dict(
sleep_time=sleep_time,
),
history_id=history_id,
)
run_response = self._post( "tools", data=payload ).json()
outputs = run_response[ "outputs" ]
jobs = run_response[ "jobs" ]
assert len(outputs) == 2
assert len(jobs) == 1
def job_state():
jobs_response = self._get( "jobs/%s" % jobs[0]["id"] )
return jobs_response
# Give job some time to get up and running.
time.sleep( 2 )
running_state = wait_on_state( job_state, skip_states=["queued", "new"], assert_ok=False, timeout=15 )
assert running_state == "running", running_state
def job_state():
jobs_response = self._get( "jobs/%s" % jobs[0]["id"] )
return jobs_response
return history_id, job_state, outputs
def _raw_update_history_item( self, history_id, item_id, data ):
update_url = self._api_url( "histories/%s/contents/%s" % (history_id, item_id), use_key=True)
update_response = put(update_url, json=data)
assert_status_code_is_ok( update_response )
return update_response
def _get_history_item_as_admin( self, history_id, item_id ):
response = self._get( "histories/%s/contents/%s?view=detailed" % (history_id, item_id), admin=True )
assert_status_code_is_ok( response )
return response.json()
def test_search( self ):
history_id, dataset_id = self.__history_with_ok_dataset()
@@ -139,7 +274,7 @@ class JobsApiTestCase( api.ApiTestCase, TestsDatasets ):
self.assertEquals( len( empty_search_response.json() ), 0 )
self.__run_cat_tool( history_id, dataset_id )
self._wait_for_history( history_id, assert_ok=True )
self.dataset_populator.wait_for_history( history_id, assert_ok=True )
search_count = -1
# in case job and history aren't updated at exactly the same
@@ -160,7 +295,7 @@ class JobsApiTestCase( api.ApiTestCase, TestsDatasets ):
def __run_cat_tool( self, history_id, dataset_id ):
# Code duplication with test_jobs.py, eliminate
payload = self._run_tool_payload(
payload = self.dataset_populator.run_tool_payload(
tool_id='cat1',
inputs=dict(
input1=dict(
@@ -173,7 +308,7 @@ class JobsApiTestCase( api.ApiTestCase, TestsDatasets ):
self._post( "tools", data=payload )
def __run_randomlines_tool( self, lines, history_id, dataset_id ):
payload = self._run_tool_payload(
payload = self.dataset_populator.run_tool_payload(
tool_id="random_lines1",
inputs=dict(
num_lines=lines,
@@ -194,13 +329,13 @@ class JobsApiTestCase( api.ApiTestCase, TestsDatasets ):
return filter( lambda j: j[ "tool_id" ] == "upload1", jobs )
def __history_with_new_dataset( self ):
history_id = self._new_history()
dataset_id = self._new_dataset( history_id )[ "id" ]
history_id = self.dataset_populator.new_history()
dataset_id = self.dataset_populator.new_dataset( history_id )[ "id" ]
return history_id, dataset_id
def __history_with_ok_dataset( self ):
history_id = self._new_history()
dataset_id = self._new_dataset( history_id, wait=True )[ "id" ]
history_id = self.dataset_populator.new_history()
dataset_id = self.dataset_populator.new_dataset( history_id, wait=True )[ "id" ]
return history_id, dataset_id
def __jobs_index( self, **kwds ):
+18 -7
View File
@@ -1,18 +1,29 @@
""" Utility methods for making assertions about Galaxy API responses, etc...
"""
ASSERT_FAIL_ERROR_CODE = "Expected Galaxy error code %d, obtained %d"
ASSERT_FAIL_STATUS_CODE = "Request status code (%d) was not expected value %d. Body was %s"
ASSERT_FAIL_STATUS_CODE = "Request status code (%d) was not expected value %s. Body was %s"
def assert_status_code_is( response, expected_status_code ):
response_status_code = response.status_code
if expected_status_code != response_status_code:
try:
body = response.json()
except Exception:
body = "INVALID JSON RESPONSE <%s>" % response.content
assertion_message = ASSERT_FAIL_STATUS_CODE % ( response_status_code, expected_status_code, body )
raise AssertionError( assertion_message )
_report_status_code_error( response, expected_status_code )
def assert_status_code_is_ok( response ):
response_status_code = response.status_code
is_two_hundred_status_code = response_status_code >= 200 and response_status_code <= 300
if not is_two_hundred_status_code:
_report_status_code_error( response, "2XX" )
def _report_status_code_error( response, expected_status_code ):
try:
body = response.json()
except Exception:
body = "INVALID JSON RESPONSE <%s>" % response.content
assertion_message = ASSERT_FAIL_STATUS_CODE % ( response.status_code, expected_status_code, body )
raise AssertionError( assertion_message )
def assert_has_keys( response, *keys ):
+1
View File
@@ -177,6 +177,7 @@ def setup_galaxy_config(
cleanup_job='onsuccess',
data_manager_config_file=data_manager_config_file,
enable_beta_tool_formats=True,
expose_dataset_path=True,
file_path=file_path,
galaxy_data_manager_data_path=galaxy_data_manager_data_path,
id_secret='changethisinproductiontoo',
+4 -4
View File
@@ -469,17 +469,17 @@ class DatasetCollectionPopulator( BaseDatasetCollectionPopulator ):
return create_response
def wait_on_state( state_func, assert_ok=False, timeout=DEFAULT_TIMEOUT ):
def wait_on_state( state_func, skip_states=["running", "queued", "new", "ready"], assert_ok=False, timeout=DEFAULT_TIMEOUT ):
def get_state( ):
response = state_func()
assert response.status_code == 200, "Failed to fetch state update while waiting."
state = response.json()[ "state" ]
if state not in [ "running", "queued", "new", "ready" ]:
if state in skip_states:
return None
else:
if assert_ok:
assert state == "ok", "Final state - %s - not okay." % state
return state
else:
return None
return wait_on( get_state, desc="state", timeout=timeout)
+4 -4
View File
@@ -1,6 +1,6 @@
<tool id="create_10" name="Create 10">
<description>create 10</description>
<command>
<tool id="create_10" name="create_10">
<description>create 10 datasets for testing tools with many outputs</description>
<command><![CDATA[
echo "1" > 1;
echo "2" > 2;
echo "3" > 3;
@@ -11,7 +11,7 @@
echo "8" > 8;
echo "9" > 9;
echo "10" > 10;
</command>
]]></command>
<inputs>
<param name="input1" type="data" label="Concatenate Dataset"/>
<param name="input2" type="data" label="Concatenate Dataset"/>
+18
View File
@@ -0,0 +1,18 @@
<tool id="create_2" name="create_2">
<command><![CDATA[
echo "1" > '$out_file1';
echo "2" > '$out_file2';
sleep '$sleep_time';
]]></command>
<inputs>
<param name="sleep_time" type="integer" label="Sleep" help="Optionally simulates computation before creating collection" value="0" />
</inputs>
<outputs>
<data name="out_file1" format="txt" />
<data name="out_file2" format="txt" />
</outputs>
<tests>
</tests>
<help>
</help>
</tool>
@@ -54,6 +54,7 @@
<tool file="output_filter_exception_1.xml" />
<tool file="output_collection_filter.xml" />
<tool file="output_auto_format.xml" />
<tool file="create_2.xml" />
<tool file="create_10.xml" />
<tool file="disambiguate_repeats.xml" />
<tool file="min_repeat.xml" />