From 009e2cbf5e9592fd8af782b2cde36d8d1687c003 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Mon, 13 Feb 2017 11:16:36 -0500 Subject: [PATCH] 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. --- lib/galaxy/jobs/__init__.py | 70 ++++++++--- lib/galaxy/managers/datasets.py | 13 +- test/api/test_jobs.py | 133 +++++++++++++++++++- test/base/api_asserts.py | 25 ++-- test/base/driver_util.py | 1 + test/base/populators.py | 8 +- test/functional/tools/create_10.xml | 8 +- test/functional/tools/create_2.xml | 18 +++ test/functional/tools/samples_tool_conf.xml | 1 + 9 files changed, 235 insertions(+), 42 deletions(-) create mode 100644 test/functional/tools/create_2.xml diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index 56298a3ed6f..dd457bf1c22 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -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 diff --git a/lib/galaxy/managers/datasets.py b/lib/galaxy/managers/datasets.py index ff2828f8971..9104f1fcd92 100644 --- a/lib/galaxy/managers/datasets.py +++ b/lib/galaxy/managers/datasets.py @@ -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: diff --git a/test/api/test_jobs.py b/test/api/test_jobs.py index 89beeb7c1f4..c4bc99f30f6 100644 --- a/test/api/test_jobs.py +++ b/test/api/test_jobs.py @@ -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() diff --git a/test/base/api_asserts.py b/test/base/api_asserts.py index 19878851663..ea316954a1c 100644 --- a/test/base/api_asserts.py +++ b/test/base/api_asserts.py @@ -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 ): diff --git a/test/base/driver_util.py b/test/base/driver_util.py index 4df2aba42fb..f65b0dd9e4b 100644 --- a/test/base/driver_util.py +++ b/test/base/driver_util.py @@ -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', diff --git a/test/base/populators.py b/test/base/populators.py index 183d58c9903..f1440f04eb3 100644 --- a/test/base/populators.py +++ b/test/base/populators.py @@ -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) diff --git a/test/functional/tools/create_10.xml b/test/functional/tools/create_10.xml index af872252d68..e459b13ec03 100644 --- a/test/functional/tools/create_10.xml +++ b/test/functional/tools/create_10.xml @@ -1,6 +1,6 @@ - - create 10 - + + create 10 datasets for testing tools with many outputs + 1; echo "2" > 2; echo "3" > 3; @@ -11,7 +11,7 @@ echo "8" > 8; echo "9" > 9; echo "10" > 10; - + ]]> diff --git a/test/functional/tools/create_2.xml b/test/functional/tools/create_2.xml new file mode 100644 index 00000000000..f9827846f64 --- /dev/null +++ b/test/functional/tools/create_2.xml @@ -0,0 +1,18 @@ + + '$out_file1'; + echo "2" > '$out_file2'; + sleep '$sleep_time'; + ]]> + + + + + + + + + + + + diff --git a/test/functional/tools/samples_tool_conf.xml b/test/functional/tools/samples_tool_conf.xml index faae608f365..3af0d3644f5 100644 --- a/test/functional/tools/samples_tool_conf.xml +++ b/test/functional/tools/samples_tool_conf.xml @@ -54,6 +54,7 @@ +