From 8ef37d260506b64cb4e11d3c02b53ddfdbd6637b Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Wed, 18 Mar 2020 17:41:13 +0100 Subject: [PATCH 01/11] Set HDA state to job state on history import Prior to this you could export histories with failed datasets, but on import the dataset state would be set to OK (despite the job state being correctly set to error). --- lib/galaxy/model/store/__init__.py | 2 ++ lib/galaxy_test/api/test_histories.py | 26 +++++++++++++++++++++++--- 2 files changed, 25 insertions(+), 3 deletions(-) diff --git a/lib/galaxy/model/store/__init__.py b/lib/galaxy/model/store/__init__.py index 99121f2029c..830f4dbfeda 100644 --- a/lib/galaxy/model/store/__init__.py +++ b/lib/galaxy/model/store/__init__.py @@ -828,6 +828,7 @@ class DirectoryImportModelStore1901(BaseDirectoryImportModelStore): for output_key in job_attrs['output_datasets']: output_hda = _find_hda(output_key) if output_hda: + output_hda.state = imported_job.state imported_job.add_output_dataset(output_hda.name, output_hda) if 'input_mapping' in job_attrs: @@ -890,6 +891,7 @@ class DirectoryImportModelStoreLatest(BaseDirectoryImportModelStore): for output_key in output_keys: output_hda = _find_hda(output_key) if output_hda: + 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_test/api/test_histories.py b/lib/galaxy_test/api/test_histories.py index 1be49800e90..4868949fd2b 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 @@ -235,6 +236,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_deleted" + 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 +322,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 +345,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 From 5519fa79ab73e50e020120d720091aeb9b5c59e4 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Wed, 18 Mar 2020 19:06:18 +0100 Subject: [PATCH 02/11] Serialize HDA state during history export, fallback to job state --- lib/galaxy/model/__init__.py | 1 + lib/galaxy/model/store/__init__.py | 14 +++++++++++--- test/unit/test_model_store.py | 19 +++++++++++++++++-- 3 files changed, 29 insertions(+), 5 deletions(-) diff --git a/lib/galaxy/model/__init__.py b/lib/galaxy/model/__init__.py index 0f61eab86aa..621124e9ccd 100644 --- a/lib/galaxy/model/__init__.py +++ b/lib/galaxy/model/__init__.py @@ -3136,6 +3136,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 830f4dbfeda..ba8fa4d1f52 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. @@ -828,7 +832,9 @@ class DirectoryImportModelStore1901(BaseDirectoryImportModelStore): for output_key in job_attrs['output_datasets']: output_hda = _find_hda(output_key) if output_hda: - output_hda.state = imported_job.state + 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: @@ -891,7 +897,9 @@ class DirectoryImportModelStoreLatest(BaseDirectoryImportModelStore): for output_key in output_keys: output_hda = _find_hda(output_key) if output_hda: - output_hda.state = imported_job.state + 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/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) From 274ed99b3d882bd8e3f570825a33ca3214bd3362 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Wed, 18 Mar 2020 13:27:54 +0100 Subject: [PATCH 03/11] Make history import/export test more stable We need to wait for the second HDA to be finished, otherwise deleting the dataset wil cancel the job (if it is still running) and then the dataset will be in discarded state (either because that's what cancelling the job does, or if the export starts first it's the import logic that sets state discarded if the file was not exported). This means sometimes the state will be OK and sometimes it would be discarded. With this change the state will be OK. --- lib/galaxy_test/api/test_histories.py | 9 +++++---- 1 file changed, 5 insertions(+), 4 deletions(-) diff --git a/lib/galaxy_test/api/test_histories.py b/lib/galaxy_test/api/test_histories.py index 4868949fd2b..9ea6961cad8 100644 --- a/lib/galaxy_test/api/test_histories.py +++ b/lib/galaxy_test/api/test_histories.py @@ -182,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"] @@ -214,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"}) @@ -223,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, From fa08def5363989de592fab542fc107d0cea4fc2d Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Wed, 18 Mar 2020 19:52:42 +0100 Subject: [PATCH 04/11] Unit test fix --- test/unit/tools/test_history_imp_exp.py | 2 ++ 1 file changed, 2 insertions(+) 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) From 563969082eb16127166c292d15cb8e7497a98d07 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Wed, 18 Mar 2020 20:59:34 +0100 Subject: [PATCH 05/11] Make test history name unique --- lib/galaxy_test/api/test_histories.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/lib/galaxy_test/api/test_histories.py b/lib/galaxy_test/api/test_histories.py index 9ea6961cad8..7da3695a87a 100644 --- a/lib/galaxy_test/api/test_histories.py +++ b/lib/galaxy_test/api/test_histories.py @@ -239,7 +239,7 @@ class HistoriesApiTestCase(ApiTestCase): @skip_without_tool("job_properties") def test_import_export_failed_job(self): - history_name = "for_export_include_deleted" + 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) From 74efc01ccd9f283d8af028fd004fcfa1ecae6111 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Thu, 19 Mar 2020 12:16:37 +0100 Subject: [PATCH 06/11] Retry reading tool conf on IOError I think that might stabilize the test_tool_reload* unit tests. --- lib/galaxy/tools/toolbox/watcher.py | 19 ++++++++++++++----- 1 file changed, 14 insertions(+), 5 deletions(-) diff --git a/lib/galaxy/tools/toolbox/watcher.py b/lib/galaxy/tools/toolbox/watcher.py index 2a165cae02e..032ba55ec2d 100644 --- a/lib/galaxy/tools/toolbox/watcher.py +++ b/lib/galaxy/tools/toolbox/watcher.py @@ -71,6 +71,8 @@ class ToolConfWatcher(object): def __init__(self, reload_callback, tool_cache=None): self.paths = {} + self.drop_on_next_loop = set() + self.drop_now = set() self.cache = tool_cache self._active = False self._lock = threading.Lock() @@ -131,11 +133,16 @@ 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] - del paths[path] - except KeyError: - pass + if path in self.drop_now: + log.warning("'%s' could not be read, removing from watched files") + try: + del hashes[path] + del paths[path] + except KeyError: + pass + else: + log.debug("'%s could not be read") + self.drop_on_next_loop.add(path) if self.cache: self.cache.cleanup() do_reload = True @@ -145,6 +152,8 @@ class ToolConfWatcher(object): do_reload = True if do_reload: self.reload_callback() + self.drop_now = self.drop_on_next_loop + self.drop_on_next_loop = set() self.exit.wait(1) def monitor(self, path): From fbc70dae1904391fd91a53048ebf4a3d68d7e6c5 Mon Sep 17 00:00:00 2001 From: Marius van den Beek Date: Fri, 20 Mar 2020 09:17:45 +0100 Subject: [PATCH 07/11] fix logging Co-Authored-By: Nicola Soranzo --- lib/galaxy/tools/toolbox/watcher.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/lib/galaxy/tools/toolbox/watcher.py b/lib/galaxy/tools/toolbox/watcher.py index 032ba55ec2d..75b5902eff4 100644 --- a/lib/galaxy/tools/toolbox/watcher.py +++ b/lib/galaxy/tools/toolbox/watcher.py @@ -134,7 +134,7 @@ class ToolConfWatcher(object): # and reading the file from the filesystem. We do not want the watcher # thread to die in these cases. if path in self.drop_now: - log.warning("'%s' could not be read, removing from watched files") + log.warning("'%s' could not be read, removing from watched files", path) try: del hashes[path] del paths[path] From 6aea38443f7675d5d704033a1d876ac41cd48df1 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Fri, 20 Mar 2020 09:19:32 +0100 Subject: [PATCH 08/11] Make drop_now and drop_on_next_loop local variables --- lib/galaxy/tools/toolbox/watcher.py | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/lib/galaxy/tools/toolbox/watcher.py b/lib/galaxy/tools/toolbox/watcher.py index 75b5902eff4..f224bba505a 100644 --- a/lib/galaxy/tools/toolbox/watcher.py +++ b/lib/galaxy/tools/toolbox/watcher.py @@ -71,8 +71,6 @@ class ToolConfWatcher(object): def __init__(self, reload_callback, tool_cache=None): self.paths = {} - self.drop_on_next_loop = set() - self.drop_now = set() self.cache = tool_cache self._active = False self._lock = threading.Lock() @@ -103,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: @@ -133,7 +133,7 @@ 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. - if path in self.drop_now: + if path in drop_now: log.warning("'%s' could not be read, removing from watched files", path) try: del hashes[path] @@ -142,7 +142,7 @@ class ToolConfWatcher(object): pass else: log.debug("'%s could not be read") - self.drop_on_next_loop.add(path) + drop_on_next_loop.add(path) if self.cache: self.cache.cleanup() do_reload = True @@ -152,8 +152,8 @@ class ToolConfWatcher(object): do_reload = True if do_reload: self.reload_callback() - self.drop_now = self.drop_on_next_loop - self.drop_on_next_loop = set() + drop_now = drop_on_next_loop + drop_on_next_loop = set() self.exit.wait(1) def monitor(self, path): From 6208552e187fb67426dbd003be031de6e4ef1e32 Mon Sep 17 00:00:00 2001 From: Marius van den Beek Date: Fri, 20 Mar 2020 18:00:30 +0100 Subject: [PATCH 09/11] Fix removing tools from watched paths Co-Authored-By: Nicola Soranzo --- lib/galaxy/tools/toolbox/watcher.py | 6 ++---- 1 file changed, 2 insertions(+), 4 deletions(-) diff --git a/lib/galaxy/tools/toolbox/watcher.py b/lib/galaxy/tools/toolbox/watcher.py index f224bba505a..8f0a753fedb 100644 --- a/lib/galaxy/tools/toolbox/watcher.py +++ b/lib/galaxy/tools/toolbox/watcher.py @@ -135,11 +135,9 @@ class ToolConfWatcher(object): # thread to die in these cases. if path in drop_now: log.warning("'%s' could not be read, removing from watched files", path) - try: + del paths[path] + if path in hashes: del hashes[path] - del paths[path] - except KeyError: - pass else: log.debug("'%s could not be read") drop_on_next_loop.add(path) From bc3733204af57c82b118bb06c807b01e3c95c584 Mon Sep 17 00:00:00 2001 From: Nicola Soranzo Date: Fri, 20 Mar 2020 18:21:15 +0100 Subject: [PATCH 10/11] Fix another logging --- lib/galaxy/tools/toolbox/watcher.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/lib/galaxy/tools/toolbox/watcher.py b/lib/galaxy/tools/toolbox/watcher.py index 8f0a753fedb..7733ded05ab 100644 --- a/lib/galaxy/tools/toolbox/watcher.py +++ b/lib/galaxy/tools/toolbox/watcher.py @@ -139,7 +139,7 @@ class ToolConfWatcher(object): if path in hashes: del hashes[path] else: - log.debug("'%s could not be read") + log.debug("'%s could not be read", path) drop_on_next_loop.add(path) if self.cache: self.cache.cleanup() From 25332bdd2f3228cba76fd8c921d41810db96198e Mon Sep 17 00:00:00 2001 From: Nicola Soranzo Date: Fri, 20 Mar 2020 17:24:32 +0000 Subject: [PATCH 11/11] Ensure miniumum virtualenv version for py27-unit build --- .circleci/config.yml | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/.circleci/config.yml b/.circleci/config.yml index a87c6269b74..3afbbc0bb98 100644 --- a/.circleci/config.yml +++ b/.circleci/config.yml @@ -95,7 +95,8 @@ jobs: <<: *set_workdir steps: - *restore_repo_cache - - *install_tox + # Ensure minimum virtualenv version due to https://github.com/pypa/virtualenv/issues/1670 + - run: sudo pip install tox 'virtualenv>=20.0.8' - run: tox -e py27-unit py35_docstring: docker: