From 47a0d2c7c5a93ee2b7548fe068aa98bfd640140d Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Sat, 26 Sep 2020 17:53:30 +0200 Subject: [PATCH 1/8] Implement collection.dataset_instances with query Instead of recursively loading elements fro child collections. We use this in a bunch of different places, but this makes the check_inputs_ready check about twice as fast. Before: *** PROFILER RESULTS *** check_inputs_ready (/Users/mvandenb/src/galaxy/lib/galaxy/tools/actions/model_operations.py:17) function called 1 times 305330 function calls (299676 primitive calls) in 0.509 seconds Ordered by: cumulative time, internal time, call count List reduced from 717 to 40 due to restriction <40> ncalls tottime percall cumtime percall filename:lineno(function) 1 0.000 0.000 0.510 0.510 model_operations.py:17(check_inputs_ready) 1 0.000 0.000 0.507 0.507 __init__.py:246(_collect_inputs) 5314/4312 0.003 0.000 0.494 0.000 attributes.py:279(__get__) 1202/601 0.002 0.000 0.491 0.001 attributes.py:699(get) 601 0.003 0.000 0.482 0.001 strategies.py:665(_load_for_state) 501 0.001 0.000 0.473 0.001 :1() 501 0.006 0.000 0.472 0.001 strategies.py:772(_emit_lazyload) 2 0.000 0.000 0.427 0.213 __init__.py:1513(visit_inputs) 2 0.000 0.000 0.427 0.213 __init__.py:21(visit_input_values) 4 0.000 0.000 0.427 0.107 __init__.py:117(callback_helper) 1 0.000 0.000 0.427 0.427 __init__.py:65(_collect_input_datasets) 2 0.000 0.000 0.427 0.213 __init__.py:81(visitor) 202/2 0.001 0.000 0.336 0.168 __init__.py:4116(dataset_instances) 501 0.005 0.000 0.243 0.000 baked.py:421(__iter__) 501 0.002 0.000 0.214 0.000 query.py:3501(_execute_and_instances) 200 0.005 0.000 0.204 0.001 baked.py:557(_load_on_pk_identity) 301 0.003 0.000 0.193 0.001 baked.py:539(all) 502 0.001 0.000 0.186 0.000 base.py:952(execute) 502 0.001 0.000 0.185 0.000 elements.py:296(_execute_on_connection) 502 0.003 0.000 0.184 0.000 base.py:1088(_execute_clauseelement) 502 0.004 0.000 0.175 0.000 base.py:1195(_execute_context) 502 0.000 0.000 0.152 0.000 default.py:589(do_execute) 200 0.000 0.000 0.152 0.001 __init__.py:4565(dataset_instance) 502 0.147 0.000 0.152 0.000 {method 'execute' of 'psycopg2.extensions.cursor' objects} 200 0.000 0.000 0.151 0.001 __init__.py:4554(element_object) 1001 0.005 0.000 0.144 0.000 loading.py:35(instances) 100 0.001 0.000 0.082 0.001 __init__.py:83(process_dataset) 101 0.000 0.000 0.080 0.001 __init__.py:155(auto_propagated_tags) 100 0.000 0.000 0.076 0.001 security.py:473(can_access_dataset) 100 0.000 0.000 0.076 0.001 security.py:1031(dataset_is_public) 501 0.001 0.000 0.060 0.000 loading.py:59() 501 0.002 0.000 0.060 0.000 query.py:4345(row_processor) 502 0.001 0.000 0.059 0.000 result.py:1268(fetchall) 601/501 0.017 0.000 0.057 0.000 loading.py:354(_instance_processor) 502 0.001 0.000 0.054 0.000 result.py:926(_soft_close) 502 0.001 0.000 0.052 0.000 base.py:899(close) 502 0.001 0.000 0.051 0.000 base.py:1031(close) 502 0.001 0.000 0.050 0.000 base.py:858(_checkin) 502 0.002 0.000 0.050 0.000 base.py:671(_finalize_fairy) 501 0.002 0.000 0.042 0.000 baked.py:180(_add_lazyload_options) after: *** PROFILER RESULTS *** check_inputs_ready (/Users/mvandenb/src/galaxy/lib/galaxy/tools/actions/model_operations.py:17) function called 1 times 186085 function calls (183671 primitive calls) in 0.251 seconds Ordered by: cumulative time, internal time, call count List reduced from 738 to 40 due to restriction <40> ncalls tottime percall cumtime percall filename:lineno(function) 1 0.000 0.000 0.251 0.251 model_operations.py:17(check_inputs_ready) 2920/2718 0.002 0.000 0.200 0.000 attributes.py:279(__get__) 802/401 0.001 0.000 0.198 0.000 attributes.py:699(get) 401 0.002 0.000 0.192 0.000 strategies.py:665(_load_for_state) 301 0.001 0.000 0.187 0.001 :1() 301 0.003 0.000 0.187 0.001 strategies.py:772(_emit_lazyload) 1 0.000 0.000 0.161 0.161 __init__.py:246(_collect_inputs) 201 0.002 0.000 0.111 0.001 baked.py:539(all) 301 0.002 0.000 0.107 0.000 baked.py:421(__iter__) 303 0.001 0.000 0.102 0.000 query.py:3501(_execute_and_instances) 2 0.000 0.000 0.099 0.049 __init__.py:1513(visit_inputs) 2 0.000 0.000 0.099 0.049 __init__.py:21(visit_input_values) 4 0.000 0.000 0.099 0.025 __init__.py:117(callback_helper) 1 0.000 0.000 0.099 0.099 __init__.py:65(_collect_input_datasets) 2 0.000 0.000 0.099 0.049 __init__.py:81(visitor) 1 0.000 0.000 0.091 0.091 __init__.py:2760(check_inputs_ready) 304 0.000 0.000 0.088 0.000 base.py:952(execute) 304 0.000 0.000 0.088 0.000 elements.py:296(_execute_on_connection) 304 0.002 0.000 0.087 0.000 base.py:1088(_execute_clauseelement) 304 0.002 0.000 0.078 0.000 base.py:1195(_execute_context) 101/1 0.000 0.000 0.072 0.072 __init__.py:4049(populated) 76 0.000 0.000 0.069 0.001 {built-in method builtins.all} 101 0.000 0.000 0.069 0.001 __init__.py:4053() 100 0.001 0.000 0.069 0.001 __init__.py:83(process_dataset) 304 0.000 0.000 0.065 0.000 default.py:589(do_execute) 304 0.063 0.000 0.064 0.000 {method 'execute' of 'psycopg2.extensions.cursor' objects} 100 0.000 0.000 0.063 0.001 security.py:473(can_access_dataset) 100 0.000 0.000 0.063 0.001 security.py:1031(dataset_is_public) 803 0.003 0.000 0.062 0.000 loading.py:35(instances) 101 0.000 0.000 0.062 0.001 __init__.py:155(auto_propagated_tags) 100 0.002 0.000 0.053 0.001 baked.py:557(_load_on_pk_identity) 2 0.000 0.000 0.039 0.020 __init__.py:4116(dataset_instances) 2 0.000 0.000 0.032 0.016 query.py:3303(all) 304 0.001 0.000 0.029 0.000 result.py:1268(fetchall) 304 0.001 0.000 0.026 0.000 result.py:926(_soft_close) 304 0.001 0.000 0.025 0.000 base.py:899(close) 304 0.000 0.000 0.025 0.000 base.py:1031(close) 304 0.000 0.000 0.024 0.000 base.py:858(_checkin) 304 0.001 0.000 0.024 0.000 base.py:671(_finalize_fairy) 2 0.000 0.000 0.022 0.011 query.py:3476(__iter__) --- lib/galaxy/model/__init__.py | 38 ++++++++++++++++++++++++++++-------- 1 file changed, 30 insertions(+), 8 deletions(-) diff --git a/lib/galaxy/model/__init__.py b/lib/galaxy/model/__init__.py index 6cf72fd1006..be537bed3df 100644 --- a/lib/galaxy/model/__init__.py +++ b/lib/galaxy/model/__init__.py @@ -4109,14 +4109,36 @@ class DatasetCollection(Dictifiable, UsesAnnotations, RepresentById): @property def dataset_instances(self): - instances = [] - for element in self.elements: - if element.is_collection: - instances.extend(element.child_collection.dataset_instances) - else: - instance = element.dataset_instance - instances.append(instance) - return instances + db_session = object_session(self) + if db_session and self.id: + dc = alias(DatasetCollection.table) + de = alias(DatasetCollectionElement.table) + hda = alias(HistoryDatasetAssociation.table) + + depth_collection_type = self.collection_type + select_from = dc.outerjoin(de, de.c.dataset_collection_id == dc.c.id) + + while ":" in depth_collection_type: + child_collection = alias(DatasetCollection.table) + child_collection_element = alias(DatasetCollectionElement.table) + select_from = select_from.outerjoin(child_collection, child_collection.c.id == de.c.child_collection_id) + select_from = select_from.outerjoin(child_collection_element, child_collection_element.c.dataset_collection_id == child_collection.c.id) + + de = child_collection_element + depth_collection_type = depth_collection_type.split(":", 1)[1] + select_from = select_from.outerjoin(hda, hda.c.id == de.c.hda_id) + select_stmt = select([hda]).select_from(select_from).where(dc.c.id == self.id).distinct() + return db_session.query(HistoryDatasetAssociation).select_entity_from(select_stmt).all() + else: + # Sessionless context + instances = [] + for element in self.elements: + if element.is_collection: + instances.extend(element.child_collection.dataset_instances) + else: + instance = element.dataset_instance + instances.append(instance) + return instances @property def dataset_elements(self): From 0b5275913f3b4d74d6c06ad19949af6982620a65 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Sat, 26 Sep 2020 18:00:28 +0200 Subject: [PATCH 2/8] Speed up collection input dataset ready check We can use the optimized dataset_states_and_extensions_summary query to not load up all collection dataset instances here. --- lib/galaxy/tools/__init__.py | 15 ++++++++------- 1 file changed, 8 insertions(+), 7 deletions(-) diff --git a/lib/galaxy/tools/__init__.py b/lib/galaxy/tools/__init__.py index 2ad7345053a..9b22928641d 100644 --- a/lib/galaxy/tools/__init__.py +++ b/lib/galaxy/tools/__init__.py @@ -2757,24 +2757,25 @@ class DatabaseOperationTool(Tool): return not self.require_dataset_ok def check_inputs_ready(self, input_datasets, input_dataset_collections): - def check_dataset_instance(input_dataset): - if input_dataset.is_pending: + def check_dataset_state(state): + if state in model.Dataset.non_ready_states: raise ToolInputsNotReadyException("An input dataset is pending.") if self.require_dataset_ok: - if input_dataset.state != input_dataset.dataset.states.OK: + if state != model.Dataset.states.OK: raise ValueError("Tool requires inputs to be in valid state.") for input_dataset in input_datasets.values(): - check_dataset_instance(input_dataset) + check_dataset_state(input_dataset.state) for input_dataset_collection_pairs in input_dataset_collections.values(): for input_dataset_collection, _ in input_dataset_collection_pairs: - if not input_dataset_collection.collection.populated: + if not input_dataset_collection.collection.populated_optimized: raise ToolInputsNotReadyException("An input collection is not populated.") - for dataset_instance in input_dataset_collection.dataset_instances: - check_dataset_instance(dataset_instance) + states, _ = input_dataset_collection.collection.dataset_states_and_extensions_summary + for state in states: + check_dataset_state(state) def _add_datasets_to_history(self, history, elements, datasets_visible=False): datasets = [] From d4fd4d597dac7efe5dfbbdb0aafd8e1bdc11e725 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Sat, 26 Sep 2020 18:03:03 +0200 Subject: [PATCH 3/8] Speed up DefaultToolAction._collect_input_datasets We can again use collection.dataset_states_and_extensions_summary to avoid having to load all dataset inputs. With this commit: *** PROFILER RESULTS *** check_inputs_ready (/Users/mvandenb/src/galaxy/lib/galaxy/tools/actions/model_operations.py:17) function called 1 times 17507 function calls (17346 primitive calls) in 0.019 seconds Ordered by: cumulative time, internal time, call count List reduced from 374 to 40 due to restriction <40> ncalls tottime percall cumtime percall filename:lineno(function) 1 0.000 0.000 0.019 0.019 model_operations.py:17(check_inputs_ready) 1 0.000 0.000 0.015 0.015 __init__.py:251(_collect_inputs) 2 0.000 0.000 0.015 0.008 __init__.py:1513(visit_inputs) 2 0.000 0.000 0.015 0.008 __init__.py:21(visit_input_values) 1 0.000 0.000 0.015 0.015 __init__.py:65(_collect_input_datasets) 4 0.000 0.000 0.015 0.004 __init__.py:117(callback_helper) 2 0.000 0.000 0.015 0.008 __init__.py:81(visitor) 107 0.000 0.000 0.010 0.000 langhelpers.py:880(__get__) 17 0.000 0.000 0.010 0.001 selectable.py:634(columns) 17 0.000 0.000 0.009 0.001 selectable.py:1395(_populate_column_collection) 170 0.001 0.000 0.009 0.000 schema.py:1658(_make_proxy) 1 0.000 0.000 0.009 0.009 __init__.py:4056(dataset_action_tuples) 3 0.000 0.000 0.007 0.002 session.py:1155(execute) 3 0.000 0.000 0.007 0.002 base.py:952(execute) 3 0.000 0.000 0.007 0.002 elements.py:296(_execute_on_connection) 3 0.000 0.000 0.007 0.002 base.py:1088(_execute_clauseelement) 2 0.000 0.000 0.007 0.003 __init__.py:3977(dataset_states_and_extensions_summary) 170 0.002 0.000 0.005 0.000 schema.py:1089(__init__) 3 0.000 0.000 0.005 0.002 base.py:1195(_execute_context) 3 0.000 0.000 0.005 0.002 default.py:589(do_execute) 3 0.005 0.002 0.005 0.002 {method 'execute' of 'psycopg2.extensions.cursor' objects} 1 0.000 0.000 0.004 0.004 __init__.py:2760(check_inputs_ready) 1 0.000 0.000 0.004 0.004 __init__.py:4011(populated_optimized) 170 0.000 0.000 0.002 0.000 schema.py:102(_init_items) 40 0.000 0.000 0.002 0.000 base.py:461(_set_parent_with_dispatch) 40 0.000 0.000 0.002 0.000 schema.py:2153(_set_parent) 40 0.000 0.000 0.002 0.000 schema.py:1609(_on_table_attach) 40 0.000 0.000 0.002 0.000 api.py:34(listen) 3 0.000 0.000 0.002 0.001 :1() 3 0.000 0.000 0.002 0.001 elements.py:412(compile) 3 0.000 0.000 0.002 0.001 elements.py:478(_compiler) 3 0.000 0.000 0.002 0.001 compiler.py:527(__init__) 3 0.000 0.000 0.002 0.001 compiler.py:274(__init__) 3 0.000 0.000 0.002 0.001 compiler.py:349(process) 108/3 0.000 0.000 0.002 0.001 visitors.py:86(_compiler_dispatch) 3 0.000 0.000 0.002 0.001 compiler.py:2045(visit_select) 40 0.000 0.000 0.001 0.000 registry.py:193(listen) 3 0.000 0.000 0.001 0.000 :1(select) 4 0.000 0.000 0.001 0.000 deprecations.py:126(warned) 3 0.001 0.000 0.001 0.000 selectable.py:2826(__init__) Which is about 25 times faster as compared to the starting situation prior to 7b4ed38453e30e476f950c599ef9ab8a7e66942d. On the whole for list:list with 100 inner elements and the apply_rules tool this is about 25% faster (1683 ms vs 1254 ms now). --- lib/galaxy/tools/__init__.py | 2 +- lib/galaxy/tools/actions/__init__.py | 23 +++++++++++++++-------- 2 files changed, 16 insertions(+), 9 deletions(-) diff --git a/lib/galaxy/tools/__init__.py b/lib/galaxy/tools/__init__.py index 9b22928641d..97c5b3dfaee 100644 --- a/lib/galaxy/tools/__init__.py +++ b/lib/galaxy/tools/__init__.py @@ -2773,7 +2773,7 @@ class DatabaseOperationTool(Tool): if not input_dataset_collection.collection.populated_optimized: raise ToolInputsNotReadyException("An input collection is not populated.") - states, _ = input_dataset_collection.collection.dataset_states_and_extensions_summary + states, _ = input_dataset_collection.collection.dataset_states_and_extensions_summary for state in states: check_dataset_state(state) diff --git a/lib/galaxy/tools/actions/__init__.py b/lib/galaxy/tools/actions/__init__.py index d655dc7ac54..068397d9bcf 100644 --- a/lib/galaxy/tools/actions/__init__.py +++ b/lib/galaxy/tools/actions/__init__.py @@ -177,16 +177,23 @@ class DefaultToolAction: for action, role_id in action_tuples: record_permission(action, role_id) - replace_collection = False + _, extensions = collection.dataset_states_and_extensions_summary + conversion_required = False + for ext in extensions: + if ext: + datatype = trans.app.datatypes_registry.get_datatype_by_extension(ext) + if not datatype.matches_any(input.formats): + conversion_required = True + break processed_dataset_dict = {} for i, v in enumerate(collection.dataset_instances): - processed_dataset = process_dataset(v) - if processed_dataset is not v: - replace_collection = True - processed_dataset_dict[v] = processed_dataset - input_datasets[prefix + input.name + str(i + 1)] = processed_dataset - - if replace_collection: + processed_dataset = None + if conversion_required: + processed_dataset = process_dataset(v) + if processed_dataset is not v: + processed_dataset_dict[v] = processed_dataset + input_datasets[prefix + input.name + str(i + 1)] = processed_dataset or v + if conversion_required: collection_type_description = trans.app.dataset_collections_service.collection_type_descriptions.for_collection_type(collection.collection_type) collection_builder = CollectionBuilder(collection_type_description) collection_builder.replace_elements_in_collection( From 38682c8f60126d3b27228cc3040d126a1c90ceb1 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Sun, 27 Sep 2020 12:08:16 +0200 Subject: [PATCH 4/8] Flush once for apply_rules tool Before commit: ``` *** PROFILER RESULTS *** _produce_outputs (/Users/mvandenb/src/galaxy/lib/galaxy/tools/actions/model_operations.py:74) function called 1 times 897129 function calls (883135 primitive calls) in 0.817 seconds Ordered by: cumulative time, internal time, call count List reduced from 1500 to 40 due to restriction <40> ncalls tottime percall cumtime percall filename:lineno(function) 1 0.000 0.000 0.820 0.820 model_operations.py:74(_produce_outputs) 1 0.000 0.000 0.820 0.820 __init__.py:3158(produce_outputs) 4 0.001 0.000 0.488 0.122 session.py:2489(flush) 3 0.001 0.000 0.485 0.162 session.py:2542(_flush) 1 0.000 0.000 0.472 0.472 __init__.py:829(create_collection) 1 0.000 0.000 0.464 0.464 collections.py:97(create) 3 0.000 0.000 0.437 0.146 unitofwork.py:402(execute) 1 0.000 0.000 0.428 0.428 collections.py:138(_create_instance_for_collection) 1010 0.001 0.000 0.417 0.000 base.py:952(execute) 1010 0.001 0.000 0.415 0.000 elements.py:296(_execute_on_connection) 1010 0.005 0.000 0.414 0.000 base.py:1088(_execute_clauseelement) 9085/8883 0.005 0.000 0.356 0.000 attributes.py:699(get) 5 0.000 0.000 0.318 0.064 scoping.py:162(do) 1 0.000 0.000 0.313 0.313 collections.py:347(__persist) 1 0.000 0.000 0.309 0.309 __init__.py:2781(_add_datasets_to_history) 1 0.000 0.000 0.309 0.309 __init__.py:1814(add_datasets) 1010 0.007 0.000 0.297 0.000 base.py:1195(_execute_context) 26 0.000 0.000 0.286 0.011 unitofwork.py:520(execute_aggregate) 303 0.001 0.000 0.280 0.001 state.py:640(_load_expired) 203 0.001 0.000 0.278 0.001 loading.py:938(load_scalar_attributes) 203 0.000 0.000 0.272 0.001 loading.py:190(load_on_ident) 203 0.001 0.000 0.271 0.001 loading.py:211(load_on_pk_identity) 203 0.000 0.000 0.269 0.001 query.py:3417(one) 203 0.003 0.000 0.269 0.001 query.py:3381(one_or_none) 15 0.000 0.000 0.240 0.016 persistence.py:184(save_obj) 203 0.000 0.000 0.222 0.001 query.py:3476(__iter__) 6909/6908 0.003 0.000 0.221 0.000 attributes.py:279(__get__) 1608 0.002 0.000 0.219 0.000 strategies.py:665(_load_for_state) 304 0.001 0.000 0.215 0.001 query.py:3501(_execute_and_instances) 15 0.006 0.000 0.213 0.014 persistence.py:1039(_emit_insert_statements) 1 0.000 0.000 0.212 0.212 __init__.py:1822() 100 0.000 0.000 0.212 0.002 __init__.py:2750(get_total_size) 1009 0.001 0.000 0.194 0.000 default.py:589(do_execute) 1009 0.188 0.000 0.193 0.000 {method 'execute' of 'psycopg2.extensions.cursor' objects} 13 0.000 0.000 0.190 0.015 unitofwork.py:585(execute) 2 0.000 0.000 0.183 0.092 mapping.py:2835(db_next_hid) 5/3 0.000 0.000 0.175 0.058 session.py:899(begin) 5/3 0.000 0.000 0.175 0.058 session.py:221(__init__) 5/3 0.000 0.000 0.175 0.058 session.py:338(_take_snapshot) 14 0.002 0.000 0.154 0.011 dependency.py:792(process_saves) ``` after commit: ``` *** PROFILER RESULTS *** _produce_outputs (/Users/mvandenb/src/galaxy/lib/galaxy/tools/actions/model_operations.py:73) function called 1 times 158126 function calls (156268 primitive calls) in 0.093 seconds Ordered by: cumulative time, internal time, call count List reduced from 615 to 40 due to restriction <40> ncalls tottime percall cumtime percall filename:lineno(function) 1 0.000 0.000 0.094 0.094 model_operations.py:73(_produce_outputs) 1 0.000 0.000 0.094 0.094 __init__.py:3158(produce_outputs) 1 0.000 0.000 0.053 0.053 __init__.py:829(create_collection) 1 0.000 0.000 0.046 0.046 collections.py:97(create) 5816 0.004 0.000 0.040 0.000 attributes.py:271(__set__) 1 0.000 0.000 0.040 0.040 collections.py:487(apply_rules) 1 0.001 0.001 0.038 0.038 collections.py:500(_build_elements_from_rule_data) 100 0.000 0.000 0.036 0.000 __init__.py:3163(copy_dataset) 100 0.001 0.000 0.036 0.000 __init__.py:3149(copy) 201/1 0.001 0.000 0.035 0.035 collections.py:178(create_dataset_collection) 201 0.000 0.000 0.032 0.000 builder.py:7(build_collection) 101/1 0.001 0.000 0.028 0.028 collections.py:377(__recursively_create_collections_for_elements) 201 0.001 0.000 0.028 0.000 builder.py:18(set_collection_elements) 602 0.001 0.000 0.027 0.000 state.py:423(_initialize_instance) 101 0.000 0.000 0.018 0.000 session.py:1988(add) 101 0.000 0.000 0.018 0.000 session.py:2019(_save_or_update_state) 1403 0.002 0.000 0.015 0.000 attributes.py:976(set) 602 0.006 0.000 0.015 0.000 mapper.py:3035(cascade_iterator) 100 0.000 0.000 0.014 0.000 __init__.py:3099(__init__) 4211 0.005 0.000 0.013 0.000 attributes.py:849(set) 1403 0.002 0.000 0.012 0.000 attributes.py:1031(fire_replace_event) 100 0.001 0.000 0.012 0.000 __init__.py:2600(__init__) 1 0.000 0.000 0.011 0.011 collections.py:138(_create_instance_for_collection) 200 0.000 0.000 0.011 0.000 __init__.py:2706(set_metadata) 2 0.000 0.000 0.011 0.005 scoping.py:162(do) 1 0.000 0.000 0.011 0.011 collections.py:347(__persist) 501 0.000 0.000 0.010 0.000 list.py:13(generate_elements) 300 0.000 0.000 0.010 0.000 :1(__init__) 202 0.001 0.000 0.008 0.000 attributes.py:1268(set) 300 0.001 0.000 0.008 0.000 __init__.py:4528(__init__) 1 0.000 0.000 0.007 0.007 __init__.py:775(get_output_name) 1 0.000 0.000 0.007 0.007 template.py:40(fill_template) 1 0.000 0.000 0.007 0.007 Template.py:353(compile) 701 0.001 0.000 0.006 0.000 attributes.py:1418(emit_backref_from_scalar_set_event) 3617 0.003 0.000 0.006 0.000 relationships.py:1925(cascade_iterator) 202 0.001 0.000 0.006 0.000 collections.py:767(bulk_replace) 300 0.000 0.000 0.006 0.000 attributes.py:1237(append) 200 0.000 0.000 0.005 0.000 metadata.py:140(make_dict_copy) 400 0.000 0.000 0.004 0.000 attributes.py:871(fire_replace_event) 6216 0.004 0.000 0.004 0.000 state.py:716(_modified_event) ``` --- lib/galaxy/tools/__init__.py | 37 ++++++++++---------- lib/galaxy/tools/actions/model_operations.py | 10 +++--- 2 files changed, 24 insertions(+), 23 deletions(-) diff --git a/lib/galaxy/tools/__init__.py b/lib/galaxy/tools/__init__.py index 97c5b3dfaee..33bc0df6993 100644 --- a/lib/galaxy/tools/__init__.py +++ b/lib/galaxy/tools/__init__.py @@ -27,7 +27,6 @@ from galaxy import ( from galaxy.job_execution import output_collect from galaxy.managers.jobs import JobSearch from galaxy.metadata import get_metadata_compute_strategy -from galaxy.model.tags import GalaxyTagHandler from galaxy.tool_shed.util.repository_util import get_installed_repository from galaxy.tool_shed.util.shed_util_common import set_image_paths from galaxy.tool_util.deps import ( @@ -2787,7 +2786,7 @@ class DatabaseOperationTool(Tool): if datasets: history.add_datasets(self.sa_session, datasets, set_hid=True) - def produce_outputs(self, trans, out_data, output_collections, incoming, history): + def produce_outputs(self, trans, out_data, output_collections, incoming, history, **kwds): return self._outputs_dict() def _outputs_dict(self): @@ -2797,7 +2796,7 @@ class DatabaseOperationTool(Tool): class UnzipCollectionTool(DatabaseOperationTool): tool_type = 'unzip_collection' - def produce_outputs(self, trans, out_data, output_collections, incoming, history, tags=None): + def produce_outputs(self, trans, out_data, output_collections, incoming, history, tags=None, **kwds): has_collection = incoming["input"] if hasattr(has_collection, "element_type"): # It is a DCE @@ -2835,7 +2834,7 @@ class ZipCollectionTool(DatabaseOperationTool): class BuildListCollectionTool(DatabaseOperationTool): tool_type = 'build_list' - def produce_outputs(self, trans, out_data, output_collections, incoming, history, tags=None): + def produce_outputs(self, trans, out_data, output_collections, incoming, history, tags=None, **kwds): new_elements = OrderedDict() for i, incoming_repeat in enumerate(incoming["datasets"]): @@ -2851,7 +2850,7 @@ class BuildListCollectionTool(DatabaseOperationTool): class ExtractDatasetCollectionTool(DatabaseOperationTool): tool_type = 'extract_dataset' - def produce_outputs(self, trans, out_data, output_collections, incoming, history, tags=None): + def produce_outputs(self, trans, out_data, output_collections, incoming, history, tags=None, **kwds): has_collection = incoming["input"] if hasattr(has_collection, "element_type"): # It is a DCE @@ -3154,36 +3153,38 @@ class RelabelFromFileTool(DatabaseOperationTool): class ApplyRulesTool(DatabaseOperationTool): tool_type = 'apply_rules' - def produce_outputs(self, trans, out_data, output_collections, incoming, history, **kwds): + def produce_outputs(self, trans, out_data, output_collections, incoming, history, tag_handler, **kwds): hdca = incoming["input"] rule_set = RuleSet(incoming["rules"]) - copied_datasets = [] + datasets_to_persist = [] def copy_dataset(dataset, tags): copied_dataset = dataset.copy(flush=False) - copied_datasets.append(copied_dataset) if tags is not None: - self.app.tag_handler.set_tags_from_list(trans.get_user(), copied_dataset, tags) + tag_handler.set_tags_from_list(trans.get_user(), copied_dataset, tags, flush=False) + copied_dataset.history_id = history.id + datasets_to_persist.append(copied_dataset) return copied_dataset new_elements = self.app.dataset_collections_service.apply_rules( hdca, rule_set, copy_dataset ) - self._add_datasets_to_history(history, copied_datasets) - output_collections.create_collection( - next(iter(self.outputs.values())), "output", collection_type=rule_set.collection_type, elements=new_elements + hdca = output_collections.create_collection( + next(iter(self.outputs.values())), "output", collection_type=rule_set.collection_type, elements=new_elements, flush=False, set_hid=False, ) + if hdca: + datasets_to_persist.append(hdca) + return datasets_to_persist class TagFromFileTool(DatabaseOperationTool): tool_type = 'tag_from_file' - def produce_outputs(self, trans, out_data, output_collections, incoming, history, **kwds): + def produce_outputs(self, trans, out_data, output_collections, incoming, history, tag_handler, **kwds): hdca = incoming["input"] how = incoming['how'] new_tags_dataset_assoc = incoming["tags"] new_elements = OrderedDict() - tags_manager = GalaxyTagHandler(trans.app.model.context) new_datasets = [] def add_copied_value_to_new_elements(new_tags_dict, dce): @@ -3196,13 +3197,13 @@ class TagFromFileTool(DatabaseOperationTool): if new_tags: if how in ('add', 'remove') and dce.element_object.tags: # We need get the original tags and update them with the new tags - old_tags = {tag for tag in tags_manager.get_tags_str(dce.element_object.tags).split(',') if tag} + old_tags = {tag for tag in tag_handler.get_tags_str(dce.element_object.tags).split(',') if tag} if how == 'add': old_tags.update(set(new_tags)) elif how == 'remove': old_tags = old_tags - set(new_tags) new_tags = old_tags - tags_manager.add_tags_from_list(user=history.user, item=copied_value, new_tags_list=new_tags) + tag_handler.add_tags_from_list(user=history.user, item=copied_value, new_tags_list=new_tags, flush=False) else: # We have a collection, and we copy the elements so that we don't manipulate the original tags copied_value = dce.element_object.copy(element_destination=history) @@ -3212,14 +3213,14 @@ class TagFromFileTool(DatabaseOperationTool): new_element.element_object.visible = False new_tags = new_tags_dict.get(new_element.element_identifier) if how in ('add', 'remove'): - old_tags = {tag for tag in tags_manager.get_tags_str(old_element.element_object.tags).split(',') if tag} + old_tags = {tag for tag in tag_handler.get_tags_str(old_element.element_object.tags).split(',') if tag} if new_tags: if how == 'add': old_tags.update(set(new_tags)) elif how == 'remove': old_tags = old_tags - set(new_tags) new_tags = old_tags - tags_manager.add_tags_from_list(user=history.user, item=new_element.element_object, new_tags_list=new_tags) + tag_handler.add_tags_from_list(user=history.user, item=new_element.element_object, new_tags_list=new_tags, flush=False) new_elements[dce.element_identifier] = copied_value new_tags_path = new_tags_dataset_assoc.file_name diff --git a/lib/galaxy/tools/actions/model_operations.py b/lib/galaxy/tools/actions/model_operations.py index 4f362ff6bbb..81f545240ca 100644 --- a/lib/galaxy/tools/actions/model_operations.py +++ b/lib/galaxy/tools/actions/model_operations.py @@ -57,21 +57,21 @@ class ModelOperationToolAction(DefaultToolAction): # Create job. # job, galaxy_session = self._new_job_for_session(trans, tool, history) - self._produce_outputs(trans, tool, out_data, output_collections, incoming=incoming, history=history, tags=preserved_tags) + datasets_to_persist = self._produce_outputs(trans, tool, out_data, output_collections, incoming=incoming, history=history, tags=preserved_tags) self._record_inputs(trans, tool, job, incoming, inp_data, inp_dataset_collections) self._record_outputs(job, out_data, output_collections) job.state = job.states.OK trans.sa_session.add(job) - trans.sa_session.flush() # ensure job.id are available # Queue the job for execution # trans.app.job_manager.job_queue.put( job.id, tool.id ) # trans.log_event( "Added database job action to the job queue, id: %s" % str(job.id), tool_id=job.tool_id ) log.info("Calling produce_outputs, tool is %s" % tool) - return job, out_data + return job, out_data, datasets_to_persist, history def _produce_outputs(self, trans, tool, out_data, output_collections, incoming, history, tags): - tool.produce_outputs(trans, out_data, output_collections, incoming, history=history, tags=tags) + tag_handler = trans.app.tag_handler.create_tag_handler_session() + datasets_to_persist = tool.produce_outputs(trans, out_data, output_collections, incoming, history=history, tags=tags, tag_handler=tag_handler) mapped_over_elements = output_collections.dataset_collection_elements if mapped_over_elements: for name, value in out_data.items(): @@ -80,4 +80,4 @@ class ModelOperationToolAction(DefaultToolAction): mapped_over_elements[name].hda = value trans.sa_session.add_all(out_data.values()) - trans.sa_session.flush() + return datasets_to_persist From b3ace06a60b29b90d6e8ec0527ae608c066bcbf7 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Sun, 27 Sep 2020 17:29:57 +0200 Subject: [PATCH 5/8] Drop history session workaround from threading work This was added in 41d3dc62f7e967ea3b409be410310563d0737bc4 because history may have been created in another thread. But we don't do that anymore, so we can drop this. --- lib/galaxy/tools/actions/__init__.py | 3 --- 1 file changed, 3 deletions(-) diff --git a/lib/galaxy/tools/actions/__init__.py b/lib/galaxy/tools/actions/__init__.py index 068397d9bcf..6d517914c78 100644 --- a/lib/galaxy/tools/actions/__init__.py +++ b/lib/galaxy/tools/actions/__init__.py @@ -252,12 +252,9 @@ class DefaultToolAction: def _collect_inputs(self, tool, trans, incoming, history, current_user_roles, collection_info): """ Collect history as well as input datasets and collections. """ - app = trans.app # Set history. if not history: history = tool.get_default_history_by_trans(trans, create=True) - if history not in trans.sa_session: - history = trans.sa_session.query(app.model.History).get(history.id) # Track input dataset collections - but replace with simply lists so collect # input datasets can process these normally. From 12d80f71595ccaad284f7d6ef264fe78662ba93c Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Sun, 27 Sep 2020 17:52:45 +0200 Subject: [PATCH 6/8] Generalize tracking objects to add to history This will use the optimized route of minimizing flushes during execution handling for all database operation tools. This also seems signifcantly simpler than passing around the datasets_to_persist list. --- lib/galaxy/model/__init__.py | 23 ++++++++++++++++ lib/galaxy/tools/__init__.py | 25 ++++++----------- lib/galaxy/tools/actions/__init__.py | 29 ++++++++------------ lib/galaxy/tools/actions/model_operations.py | 7 ++--- lib/galaxy/tools/execute.py | 11 +++----- test/unit/tools/test_actions.py | 2 +- 6 files changed, 51 insertions(+), 46 deletions(-) diff --git a/lib/galaxy/model/__init__.py b/lib/galaxy/model/__init__.py index be537bed3df..db92f3f64e8 100644 --- a/lib/galaxy/model/__init__.py +++ b/lib/galaxy/model/__init__.py @@ -40,6 +40,7 @@ from sqlalchemy.orm import ( joinedload, object_session, Query, + reconstructor, ) from sqlalchemy.schema import UniqueConstraint @@ -1757,11 +1758,33 @@ class History(HasTags, Dictifiable, UsesAnnotations, HasName, RepresentById): self.datasets = [] self.galaxy_sessions = [] self.tags = [] + # Objects to eventually add to history + self._pending_additions = [] + + @reconstructor + def init_on_load(self): + # Restores properties that are not tracked in the database + self._pending_additions = [] + + def stage_addition(self, items): + history_id = self.id + for item in listify(items): + if history_id: + item.history_id = history_id + else: + item.history = self + self._pending_additions.append(item) @property def empty(self): return self.hid_counter == 1 + def add_pending_datasets(self, set_output_hid=True): + # These are assumed to be either copies of existing datasets or new, empty datasets, + # so we don't need to set the quota. + self.add_datasets(object_session(self), self._pending_additions, set_hid=set_output_hid, quota=False, flush=False) + self._pending_additions = [] + def _next_hid(self, n=1): # this is overriden in mapping.py db_next_hid() method if len(self.datasets) == 0: diff --git a/lib/galaxy/tools/__init__.py b/lib/galaxy/tools/__init__.py index 33bc0df6993..f6517f45f80 100644 --- a/lib/galaxy/tools/__init__.py +++ b/lib/galaxy/tools/__init__.py @@ -1646,9 +1646,8 @@ class Tool(Dictifiable): ) job = rval[0] out_data = rval[1] - if len(rval) == 4: - execution_slice.datasets_to_persist = rval[2] - execution_slice.history = rval[3] + if len(rval) > 2: + execution_slice.history = rval[2] except (webob.exc.HTTPFound, exceptions.MessageException) as e: # if it's a webob redirect exception, pass it up the stack raise e @@ -2777,14 +2776,10 @@ class DatabaseOperationTool(Tool): check_dataset_state(state) def _add_datasets_to_history(self, history, elements, datasets_visible=False): - datasets = [] for element_object in elements: if getattr(element_object, "history_content_type", None) == "dataset": element_object.visible = datasets_visible - datasets.append(element_object) - - if datasets: - history.add_datasets(self.sa_session, datasets, set_hid=True) + history.stage_addition(element_object) def produce_outputs(self, trans, out_data, output_collections, incoming, history, **kwds): return self._outputs_dict() @@ -3156,25 +3151,23 @@ class ApplyRulesTool(DatabaseOperationTool): def produce_outputs(self, trans, out_data, output_collections, incoming, history, tag_handler, **kwds): hdca = incoming["input"] rule_set = RuleSet(incoming["rules"]) - datasets_to_persist = [] + copied_datasets = [] def copy_dataset(dataset, tags): copied_dataset = dataset.copy(flush=False) if tags is not None: tag_handler.set_tags_from_list(trans.get_user(), copied_dataset, tags, flush=False) copied_dataset.history_id = history.id - datasets_to_persist.append(copied_dataset) + copied_datasets.append(copied_dataset) return copied_dataset new_elements = self.app.dataset_collections_service.apply_rules( hdca, rule_set, copy_dataset ) - hdca = output_collections.create_collection( - next(iter(self.outputs.values())), "output", collection_type=rule_set.collection_type, elements=new_elements, flush=False, set_hid=False, + self._add_datasets_to_history(history, copied_datasets) + output_collections.create_collection( + next(iter(self.outputs.values())), "output", collection_type=rule_set.collection_type, elements=new_elements, ) - if hdca: - datasets_to_persist.append(hdca) - return datasets_to_persist class TagFromFileTool(DatabaseOperationTool): @@ -3270,10 +3263,10 @@ class FilterFromFileTool(DatabaseOperationTool): discarded_elements[element_identifier] = copied_value self._add_datasets_to_history(history, filtered_elements.values()) - self._add_datasets_to_history(history, discarded_elements.values()) output_collections.create_collection( self.outputs["output_filtered"], "output_filtered", elements=filtered_elements ) + self._add_datasets_to_history(history, discarded_elements.values()) output_collections.create_collection( self.outputs["output_discarded"], "output_discarded", elements=discarded_elements ) diff --git a/lib/galaxy/tools/actions/__init__.py b/lib/galaxy/tools/actions/__init__.py index 6d517914c78..138c503268d 100644 --- a/lib/galaxy/tools/actions/__init__.py +++ b/lib/galaxy/tools/actions/__init__.py @@ -459,7 +459,6 @@ class DefaultToolAction: # Flush all datasets at once. return data - datasets_to_persist = [] for name, output in tool.outputs.items(): if not filter_output(output, incoming): handle_output_timer = ExecutionTimer() @@ -496,7 +495,7 @@ class DefaultToolAction: effective_output_name = output_part_def.effective_output_name element = handle_output(effective_output_name, output_part_def.output_def, hidden=True) - datasets_to_persist.append(element) + history.stage_addition(element) # TODO: this shouldn't exist in the top-level of the history at all # but for now we are still working around that by hiding the contents # there. @@ -513,17 +512,13 @@ class DefaultToolAction: element_kwds = dict(elements=collections_manager.ELEMENTS_UNINITIALIZED) else: element_kwds = dict(element_identifiers=element_identifiers) - hdca = output_collections.create_collection( + output_collections.create_collection( output=output, name=name, - set_hid=True if flush_job else False, - flush=flush_job, completed_job=completed_job, **element_kwds ) - if hdca: - datasets_to_persist.append(hdca) - log.info(f"Handled collection output named {name} for tool {tool.id} {handle_output_timer}") + log.info("Handled collection output named {} for tool {} {}".format(name, tool.id, handle_output_timer)) else: handle_output(name, output) log.info(f"Handled output named {name} for tool {tool.id} {handle_output_timer}") @@ -535,7 +530,7 @@ class DefaultToolAction: # Add all the top-level (non-child) datasets to the history unless otherwise specified for name, data in out_data.items(): if name not in child_dataset_names and name not in incoming: # don't add children; or already existing datasets, i.e. async created - datasets_to_persist.append(data) + history.stage_addition(data) # Add all the children to their parents for parent_name, child_name in parent_to_child_pairs: @@ -560,14 +555,13 @@ class DefaultToolAction: if app.config.track_jobs_in_database and rerun_remap_job_id is not None: # We need a flush here and get hids in order to rewrite jobs parameter, # but remapping jobs should only affect single jobs anyway, so this is not too costly. + history.add_pending_datasets(set_output_hid=set_output_hid) trans.sa_session.flush() - history.add_datasets(trans.sa_session, datasets_to_persist, set_hid=set_output_hid, quota=False, flush=False) self._remap_job_on_rerun(trans=trans, galaxy_session=galaxy_session, rerun_remap_job_id=rerun_remap_job_id, current_job=job, out_data=out_data) - datasets_to_persist = [] log.info(f"Setup for job {job.log_str()} complete, ready to be enqueued {job_setup_timer}") # Some tools are not really executable, but jobs are still created for them ( for record keeping ). @@ -593,8 +587,7 @@ class DefaultToolAction: else: if flush_job: # Set HID and add to history. - # This is brand new and certainly empty so don't worry about quota. - history.add_datasets(trans.sa_session, datasets_to_persist, set_hid=set_output_hid, quota=False, flush=False) + history.add_pending_datasets(set_output_hid=set_output_hid) job_flush_timer = ExecutionTimer() trans.sa_session.flush() log.info(f"Flushed transaction for job {job.log_str()} {job_flush_timer}") @@ -602,7 +595,7 @@ class DefaultToolAction: # Dispatch to a job handler. enqueue() is responsible for flushing the job app.job_manager.enqueue(job, tool=tool) trans.log_event("Added job to the job queue, id: %s" % str(job.id), tool_id=job.tool_id) - return job, out_data, datasets_to_persist, history + return job, out_data, history def _remap_job_on_rerun(self, trans, galaxy_session, rerun_remap_job_id, current_job, out_data): """ @@ -825,7 +818,7 @@ class OutputCollections: self.out_collection_instances = {} self.tags = tags - def create_collection(self, output, name, collection_type=None, set_hid=True, flush=True, completed_job=None, **element_kwds): + def create_collection(self, output, name, collection_type=None, completed_job=None, **element_kwds): input_collections = self.input_collections collections_manager = self.trans.app.dataset_collections_service collection_type = collection_type or output.structure.collection_type @@ -904,16 +897,16 @@ class OutputCollections: collection_type=collection_type, trusted_identifiers=True, tags=self.tags, - set_hid=set_hid, - flush=flush, + set_hid=False, + flush=False, completed_job=completed_job, output_name=name, **element_kwds ) # name here is name of the output element - not name # of the hdca. + self.history.stage_addition(hdca) self.out_collection_instances[name] = hdca - return hdca def on_text_for_names(input_names): diff --git a/lib/galaxy/tools/actions/model_operations.py b/lib/galaxy/tools/actions/model_operations.py index 81f545240ca..f45dd4c6baa 100644 --- a/lib/galaxy/tools/actions/model_operations.py +++ b/lib/galaxy/tools/actions/model_operations.py @@ -57,7 +57,7 @@ class ModelOperationToolAction(DefaultToolAction): # Create job. # job, galaxy_session = self._new_job_for_session(trans, tool, history) - datasets_to_persist = self._produce_outputs(trans, tool, out_data, output_collections, incoming=incoming, history=history, tags=preserved_tags) + self._produce_outputs(trans, tool, out_data, output_collections, incoming=incoming, history=history, tags=preserved_tags) self._record_inputs(trans, tool, job, incoming, inp_data, inp_dataset_collections) self._record_outputs(job, out_data, output_collections) job.state = job.states.OK @@ -67,11 +67,11 @@ class ModelOperationToolAction(DefaultToolAction): # trans.app.job_manager.job_queue.put( job.id, tool.id ) # trans.log_event( "Added database job action to the job queue, id: %s" % str(job.id), tool_id=job.tool_id ) log.info("Calling produce_outputs, tool is %s" % tool) - return job, out_data, datasets_to_persist, history + return job, out_data, history def _produce_outputs(self, trans, tool, out_data, output_collections, incoming, history, tags): tag_handler = trans.app.tag_handler.create_tag_handler_session() - datasets_to_persist = tool.produce_outputs(trans, out_data, output_collections, incoming, history=history, tags=tags, tag_handler=tag_handler) + tool.produce_outputs(trans, out_data, output_collections, incoming, history=history, tags=tags, tag_handler=tag_handler) mapped_over_elements = output_collections.dataset_collection_elements if mapped_over_elements: for name, value in out_data.items(): @@ -80,4 +80,3 @@ class ModelOperationToolAction(DefaultToolAction): mapped_over_elements[name].hda = value trans.sa_session.add_all(out_data.values()) - return datasets_to_persist diff --git a/lib/galaxy/tools/execute.py b/lib/galaxy/tools/execute.py index 86008f91544..c8f0a1cff3c 100644 --- a/lib/galaxy/tools/execute.py +++ b/lib/galaxy/tools/execute.py @@ -90,7 +90,7 @@ def execute(trans, tool, mapping_params, history, rerun_remap_job_id=None, colle jobs_executed = 0 has_remaining_jobs = False - datasets_to_persist = [] + execution_slice = None for i, execution_slice in enumerate(execution_tracker.new_execution_slices()): if max_num_jobs and jobs_executed >= max_num_jobs: @@ -100,12 +100,10 @@ def execute(trans, tool, mapping_params, history, rerun_remap_job_id=None, colle execute_single_job(execution_slice, completed_jobs[i]) history = execution_slice.history or history jobs_executed += 1 - if execution_slice.datasets_to_persist: - datasets_to_persist.extend(execution_slice.datasets_to_persist) - if datasets_to_persist: - history.add_datasets(trans.sa_session, datasets_to_persist, set_hid=True, quota=False, flush=False) - # a side effect of history.add_datasets is a commit within db_next_hid (even with flush=False). + if execution_slice: + # a side effect of adding datasets to a history is a commit within db_next_hid (even with flush=False). + history.add_pending_datasets() else: # Make sure collections, implicit jobs etc are flushed even if there are no precreated output datasets trans.sa_session.flush() @@ -130,7 +128,6 @@ class ExecutionSlice: self.job_index = job_index self.param_combination = param_combination self.dataset_collection_elements = dataset_collection_elements - self.datasets_to_persist = None self.history = None diff --git a/test/unit/tools/test_actions.py b/test/unit/tools/test_actions.py index d8d6e69bbf8..3f7347ee286 100644 --- a/test/unit/tools/test_actions.py +++ b/test/unit/tools/test_actions.py @@ -134,7 +134,7 @@ class DefaultToolActionTestCase(unittest.TestCase, tools_support.UsesApp, tools_ if incoming is None: incoming = dict(param1="moo") self._init_tool(contents) - job, out_data, _, _ = self.action.execute( + job, out_data, _, = self.action.execute( tool=self.tool, trans=self.trans, history=self.history, From 0334b2f39b7194c0da9041f818adc7cf96dd3b77 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Sun, 25 Oct 2020 16:19:30 +0100 Subject: [PATCH 7/8] Fix derived permission setting logic --- lib/galaxy/tools/actions/__init__.py | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/lib/galaxy/tools/actions/__init__.py b/lib/galaxy/tools/actions/__init__.py index 138c503268d..373fe8e472c 100644 --- a/lib/galaxy/tools/actions/__init__.py +++ b/lib/galaxy/tools/actions/__init__.py @@ -331,8 +331,7 @@ class DefaultToolAction: if not completed_job: # Determine output dataset permission/roles list - existing_datasets = [inp for inp in inp_data.values() if inp] - if existing_datasets: + if all_permissions: output_permissions = app.security_agent.guess_derived_permissions(all_permissions) else: # No valid inputs, we will use history defaults From a8191d4b8a4bcf1b7d0cd303a42838b608026b2c Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Sun, 25 Oct 2020 16:48:59 +0100 Subject: [PATCH 8/8] Move one more flush out of dicovering outputs --- lib/galaxy/model/__init__.py | 5 +++++ lib/galaxy/model/store/discover.py | 1 - lib/galaxy/objectstore/__init__.py | 6 +++++- test/unit/jobs/test_job_context.py | 3 ++- 4 files changed, 12 insertions(+), 3 deletions(-) diff --git a/lib/galaxy/model/__init__.py b/lib/galaxy/model/__init__.py index db92f3f64e8..bac0d7dd509 100644 --- a/lib/galaxy/model/__init__.py +++ b/lib/galaxy/model/__init__.py @@ -2307,6 +2307,11 @@ class StorableObject: else: self.uuid = UUID(str(uuid)) + def flush(self): + sa_session = object_session(self) + if sa_session: + sa_session.flush() + class Dataset(StorableObject, RepresentById): states = Bunch(NEW='new', diff --git a/lib/galaxy/model/store/discover.py b/lib/galaxy/model/store/discover.py index 263aa8d794a..bfee249db29 100644 --- a/lib/galaxy/model/store/discover.py +++ b/lib/galaxy/model/store/discover.py @@ -283,7 +283,6 @@ class ModelPersistenceContext(metaclass=abc.ABCMeta): association_name = f'__new_primary_file_{name}|{element_identifier_str}__' self.add_output_dataset_association(association_name, dataset) - self.flush() self.update_object_store_with_datasets(datasets=element_datasets['datasets'], paths=element_datasets['paths'], extra_files=element_datasets['extra_files']) add_datasets_timer = ExecutionTimer() self.add_datasets_to_history(element_datasets['datasets']) diff --git a/lib/galaxy/objectstore/__init__.py b/lib/galaxy/objectstore/__init__.py index 2166b2a8120..c9cfaecd6fc 100644 --- a/lib/galaxy/objectstore/__init__.py +++ b/lib/galaxy/objectstore/__init__.py @@ -256,7 +256,11 @@ class BaseObjectStore(ObjectStore): def _get_object_id(self, obj): if hasattr(obj, self.store_by): - return getattr(obj, self.store_by) + obj_id = getattr(obj, self.store_by) + if obj_id is None: + obj.flush() + return obj.id + return obj_id else: # job's don't have uuids, so always use ID in this case when creating # job working directories. diff --git a/test/unit/jobs/test_job_context.py b/test/unit/jobs/test_job_context.py index faf0ed992f8..6bbacc62331 100644 --- a/test/unit/jobs/test_job_context.py +++ b/test/unit/jobs/test_job_context.py @@ -81,6 +81,7 @@ def test_job_context_discover_outputs_flushes_once(mocker): final_job_state=job_context.final_job_state, ) collection_builder.populate() - assert spy.call_count == 1 + assert spy.call_count == 0 + sa_session.flush() assert len(collection.dataset_instances) == 10 assert collection.dataset_instances[0].dataset.file_size == 1