diff --git a/lib/galaxy/model/__init__.py b/lib/galaxy/model/__init__.py index 720dd3a8314..f8cafc59af6 100644 --- a/lib/galaxy/model/__init__.py +++ b/lib/galaxy/model/__init__.py @@ -3155,6 +3155,7 @@ class HistoryDatasetAssociation(DatasetInstance, HasTags, Dictifiable, UsesAnnot return rval rval = super(HistoryDatasetAssociation, self).serialize(id_encoder, serialization_options) + rval['state'] = self.state rval["hid"] = self.hid rval["annotation"] = unicodify(getattr(self, 'annotation', '')) rval["tags"] = self.make_tag_string_list() diff --git a/lib/galaxy/model/store/__init__.py b/lib/galaxy/model/store/__init__.py index a92cb1c1ce1..d88c22330d8 100644 --- a/lib/galaxy/model/store/__init__.py +++ b/lib/galaxy/model/store/__init__.py @@ -82,6 +82,7 @@ class ModelImportStore(object): self.sessionless = True self.user = user self.import_options = import_options or ImportOptions() + self.dataset_state_serialized = True @abc.abstractmethod def defines_new_history(self): @@ -184,6 +185,9 @@ class ModelImportStore(object): for dataset_attrs in datasets_attrs: + if 'state' not in dataset_attrs: + self.dataset_state_serialized = False + def handle_dataset_object_edit(dataset_instance): if "dataset" in dataset_attrs: assert self.import_options.allow_dataset_object_edit @@ -321,7 +325,7 @@ class ModelImportStore(object): dataset_instance.dataset.deleted = True dataset_instance.dataset.purged = True else: - dataset_instance.state = dataset_instance.states.OK + dataset_instance.state = dataset_attrs.get('state', dataset_instance.states.OK) self.object_store.update_from_file(dataset_instance.dataset, file_name=temp_dataset_file_name, create=True) # Import additional files if present. Histories exported previously might not have this attribute set. @@ -827,6 +831,9 @@ class DirectoryImportModelStore1901(BaseDirectoryImportModelStore): for output_key in job_attrs['output_datasets']: output_hda = _find_hda(output_key) if output_hda: + if not self.dataset_state_serialized: + # dataset state has not been serialized, get state from job + output_hda.state = imported_job.state imported_job.add_output_dataset(output_hda.name, output_hda) if 'input_mapping' in job_attrs: @@ -889,6 +896,9 @@ class DirectoryImportModelStoreLatest(BaseDirectoryImportModelStore): for output_key in output_keys: output_hda = _find_hda(output_key) if output_hda: + if not self.dataset_state_serialized: + # dataset state has not been serialized, get state from job + output_hda.state = imported_job.state imported_job.add_output_dataset(output_name, output_hda) if 'output_dataset_collection_mapping' in job_attrs: diff --git a/lib/galaxy/tools/toolbox/watcher.py b/lib/galaxy/tools/toolbox/watcher.py index 2a165cae02e..7733ded05ab 100644 --- a/lib/galaxy/tools/toolbox/watcher.py +++ b/lib/galaxy/tools/toolbox/watcher.py @@ -101,6 +101,8 @@ class ToolConfWatcher(object): hashes = {} while self._active and not self.exit.isSet(): do_reload = False + drop_on_next_loop = set() + drop_now = set() with self._lock: paths = list(self.paths.keys()) for path in paths: @@ -131,11 +133,14 @@ class ToolConfWatcher(object): # in rare cases `path` may be deleted between `os.path.exists` calls # and reading the file from the filesystem. We do not want the watcher # thread to die in these cases. - try: - del hashes[path] + if path in drop_now: + log.warning("'%s' could not be read, removing from watched files", path) del paths[path] - except KeyError: - pass + if path in hashes: + del hashes[path] + else: + log.debug("'%s could not be read", path) + drop_on_next_loop.add(path) if self.cache: self.cache.cleanup() do_reload = True @@ -145,6 +150,8 @@ class ToolConfWatcher(object): do_reload = True if do_reload: self.reload_callback() + drop_now = drop_on_next_loop + drop_on_next_loop = set() self.exit.wait(1) def monitor(self, path): diff --git a/lib/galaxy_test/api/test_histories.py b/lib/galaxy_test/api/test_histories.py index 1be49800e90..7da3695a87a 100644 --- a/lib/galaxy_test/api/test_histories.py +++ b/lib/galaxy_test/api/test_histories.py @@ -9,6 +9,7 @@ from requests import ( from galaxy_test.base.populators import ( DatasetCollectionPopulator, DatasetPopulator, + skip_without_tool, ) from ._framework import ApiTestCase @@ -181,7 +182,7 @@ class HistoriesApiTestCase(ApiTestCase): history_name = "for_export_default" history_id = self.dataset_populator.new_history(name=history_name) self.dataset_populator.new_dataset(history_id, content="1 2 3") - deleted_hda = self.dataset_populator.new_dataset(history_id, content="1 2 3") + deleted_hda = self.dataset_populator.new_dataset(history_id, content="1 2 3", wait=True) self.dataset_populator.delete_dataset(history_id, deleted_hda["id"]) deleted_details = self.dataset_populator.get_history_dataset_details(history_id, id=deleted_hda["id"]) assert deleted_details["deleted"] @@ -213,7 +214,7 @@ class HistoriesApiTestCase(ApiTestCase): history_name = "for_export_include_deleted" history_id = self.dataset_populator.new_history(name=history_name) self.dataset_populator.new_dataset(history_id, content="1 2 3") - deleted_hda = self.dataset_populator.new_dataset(history_id, content="1 2 3") + deleted_hda = self.dataset_populator.new_dataset(history_id, content="1 2 3", wait=True) self.dataset_populator.delete_dataset(history_id, deleted_hda["id"]) imported_history_id = self._reimport_history(history_id, history_name, wait_on_history_length=2, export_kwds={"include_deleted": "True"}) @@ -222,12 +223,13 @@ class HistoriesApiTestCase(ApiTestCase): def upload_job_check(job): assert job["tool_id"] == "upload1" - def check_ok(hda): + def check_deleted_not_purged(hda): assert hda["state"] == "ok", hda assert hda["deleted"] is True, hda + assert hda["purged"] is False, hda self._check_imported_dataset(history_id=imported_history_id, hid=1, job_checker=upload_job_check) - self._check_imported_dataset(history_id=imported_history_id, hid=2, hda_checker=check_ok, job_checker=upload_job_check) + self._check_imported_dataset(history_id=imported_history_id, hid=2, hda_checker=check_deleted_not_purged, job_checker=upload_job_check) imported_content = self.dataset_populator.get_history_dataset_content( history_id=imported_history_id, @@ -235,6 +237,24 @@ class HistoriesApiTestCase(ApiTestCase): ) assert imported_content == "1 2 3\n" + @skip_without_tool("job_properties") + def test_import_export_failed_job(self): + history_name = "for_export_include_failed_job" + history_id = self.dataset_populator.new_history(name=history_name) + self.dataset_populator.run_tool('job_properties', inputs={'failbool': True}, history_id=history_id, assert_ok=False) + self.dataset_populator.wait_for_history(history_id, assert_ok=False) + + imported_history_id = self._reimport_history(history_id, history_name, assert_ok=False, wait_on_history_length=4, export_kwds={"include_deleted": "True"}) + self._assert_history_length(imported_history_id, 4) + + def check_failed(hda_or_job): + print(hda_or_job) + assert hda_or_job["state"] == "error", hda_or_job + + self.dataset_populator._summarize_history(imported_history_id) + + self._check_imported_dataset(history_id=imported_history_id, hid=1, assert_ok=False, hda_checker=check_failed, job_checker=check_failed) + def test_import_metadata_regeneration(self): history_name = "for_import_metadata_regeneration" history_id = self.dataset_populator.new_history(name=history_name) @@ -303,9 +323,9 @@ class HistoriesApiTestCase(ApiTestCase): self._check_imported_collection(imported_history_id, hid=1, collection_type="list:paired", elements_checker=check_elements) - def _reimport_history(self, history_id, history_name, wait_on_history_length=None, export_kwds={}): + def _reimport_history(self, history_id, history_name, wait_on_history_length=None, assert_ok=True, export_kwds={}): # Ensure the history is ready to go... - self.dataset_populator.wait_for_history(history_id, assert_ok=True) + self.dataset_populator.wait_for_history(history_id, assert_ok=assert_ok) return self.dataset_populator.reimport_history( history_id, history_name, wait_on_history_length=wait_on_history_length, export_kwds=export_kwds, url=self.url, api_key=self.galaxy_interactor.api_key @@ -326,10 +346,11 @@ class HistoriesApiTestCase(ApiTestCase): contents = contents_response.json() assert len(contents) == n, contents - def _check_imported_dataset(self, history_id, hid, has_job=True, hda_checker=None, job_checker=None): + def _check_imported_dataset(self, history_id, hid, assert_ok=True, has_job=True, hda_checker=None, job_checker=None): imported_dataset_metadata = self.dataset_populator.get_history_dataset_details( history_id=history_id, hid=hid, + assert_ok=assert_ok, ) assert imported_dataset_metadata["history_content_type"] == "dataset" assert imported_dataset_metadata["history_id"] == history_id diff --git a/test/unit/test_model_store.py b/test/unit/test_model_store.py index 4f26f5b2574..58fd0d5d2ce 100644 --- a/test/unit/test_model_store.py +++ b/test/unit/test_model_store.py @@ -21,6 +21,17 @@ def test_import_export_history(): _assert_simple_cat_job_imported(imported_history) +def test_import_export_history_failed_job(): + """Test a simple job import/export, make sure state is maintained correctly.""" + app = _mock_app() + + u, h, d1, d2, j = _setup_simple_cat_job(app, state='error') + + imported_history = _import_export_history(app, h, export_files="copy") + + _assert_simple_cat_job_imported(imported_history, state='error') + + def test_import_export_bag_archive(): """Test a simple job import/export using a BagIt archive.""" dest_parent = mkdtemp() @@ -402,13 +413,15 @@ def _setup_simple_export(export_kwds): return app, h, temp_directory, import_history -def _assert_simple_cat_job_imported(imported_history): +def _assert_simple_cat_job_imported(imported_history, state='ok'): assert imported_history.name == "imported from archive: Test History" datasets = imported_history.datasets assert len(datasets) == 2 + assert datasets[0].state == datasets[1].state == state imported_job = datasets[1].creating_job assert imported_job + assert imported_job.state == state assert imported_job.output_datasets assert imported_job.output_datasets[0].dataset == datasets[1] @@ -421,17 +434,19 @@ def _assert_simple_cat_job_imported(imported_history): assert f.read().startswith("chr1\t147962192\t147962580\tNM_005997_cds_0_0_chr1_147962193_r\t0\t-") -def _setup_simple_cat_job(app): +def _setup_simple_cat_job(app, state='ok'): sa_session = app.model.context u = model.User(email="collection@example.com", password="password") h = model.History(name="Test History", user=u) d1, d2 = _create_datasets(sa_session, h, 2) + d1.state = d2.state = state j = model.Job() j.user = u j.tool_id = "cat1" + j.state = state j.add_input_dataset("input1", d1) j.add_output_dataset("out_file1", d2) diff --git a/test/unit/tools/test_history_imp_exp.py b/test/unit/tools/test_history_imp_exp.py index 7f9f63610ee..63eb9d71dc1 100644 --- a/test/unit/tools/test_history_imp_exp.py +++ b/test/unit/tools/test_history_imp_exp.py @@ -115,10 +115,12 @@ def test_export_dataset(): app, sa_session, h = _setup_history_for_export("Datasets History") d1, d2 = _create_datasets(sa_session, h, 2) + d1.state = d2.state = 'ok' j = model.Job() j.user = h.user j.tool_id = "cat1" + j.state = 'ok' j.add_input_dataset("input1", d1) j.add_output_dataset("out_file1", d2)