mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
Improve interaction between job state and purged datasets.
There are two basic behavior changes here and bunch of testing to verify the new behaviors. 1) Previously purging a dataset would cause a creating job to be cancelled but all outputs would need to be "deleted" in order to cause the creating job to be cancelled. After this change the behavior is the same for both deletion and purging - the job will run until all outputs are deleted or purged. 2) Previously Galaxy jobs wouldn't attempt to cleanup purged datasets that were created by a running tool - this cleanup now occurs. These are tackled as a pair because I assume the difference in the behaviors between deletion and purging was a hack around tools populating datasets that had been purged while running. Cancelling the job wasn't a sure fix - but it would reduce the likelihood of such problems. I think the new approach is a bit more robust and explicit.
This commit is contained in:
+51
-19
@@ -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
|
||||
|
||||
@@ -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:
|
||||
|
||||
+132
-1
@@ -1,10 +1,19 @@
|
||||
import datetime
|
||||
import json
|
||||
import os
|
||||
import time
|
||||
|
||||
from operator import itemgetter
|
||||
|
||||
from base import api
|
||||
from base.populators import DatasetPopulator
|
||||
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 ):
|
||||
@@ -121,6 +130,128 @@ class JobsApiTestCase( api.ApiTestCase ):
|
||||
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()
|
||||
|
||||
|
||||
@@ -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 ):
|
||||
|
||||
@@ -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',
|
||||
|
||||
@@ -465,17 +465,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)
|
||||
|
||||
|
||||
|
||||
@@ -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"/>
|
||||
|
||||
@@ -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" />
|
||||
|
||||
Reference in New Issue
Block a user