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 @@ +