diff --git a/swarms/regression.py b/swarms/regression.py index 0ae807492..3664904db 100755 --- a/swarms/regression.py +++ b/swarms/regression.py @@ -1,4 +1,6 @@ #!/usr/bin/env python3 +import os +import random import sys from testflows.core import * @@ -18,13 +20,23 @@ ) from helpers.feature_support import validate_feature_support, setting_supported + +def argparser(parser): + """Custom argparser that adds swarms-specific options.""" + argparser_minio(parser) + + parser.add_argument( + "--seed", + action="store", + type=int, + default=None, + help="seed for sampling the join and union combinations, default: $GITHUB_RUN_ID in CI, otherwise random (printed in the log)", + ) + from swarms.requirements.requirements import * xfails = { - "/swarms/feature/swarm joins/join clause/join 455 of 816480*": [ - (Fail, "https://github.com/Altinity/ClickHouse/issues/1244"), - ], "/swarms/feature/task rescheduling/rescheduling with bucket granularity": [ ( Fail, @@ -64,7 +76,7 @@ @Name("swarms") @FFails(ffails) @XFails(xfails) -@ArgumentParser(argparser_minio) +@ArgumentParser(argparser) @Specifications(SRS_044_Swarm_Cluster_Query_Execution) @CaptureClusterArgs @CaptureMinioArgs @@ -75,6 +87,7 @@ def regression( stress=None, with_analyzer=False, minio_args=None, + seed=None, ): """Run tests for Swarm clusters.""" nodes = { @@ -92,6 +105,14 @@ def regression( if stress is not None: self.context.stress = stress + # One seed per CI workflow run: every job of the run, and a rerun of + # any of them, samples the same join and union combinations. + if seed is None: + seed = os.environ.get("GITHUB_RUN_ID") + if seed is None: + seed = random.SystemRandom().randrange(2**32) + self.context.seed = int(seed) + minio_root_user = minio_args["minio_root_user"].value minio_root_password = minio_args["minio_root_password"].value diff --git a/swarms/tests/joins.py b/swarms/tests/joins.py index f9c9e1016..fc9b9c179 100644 --- a/swarms/tests/joins.py +++ b/swarms/tests/joins.py @@ -27,9 +27,12 @@ def __init__( cluster_name=None, minio_root_user=None, minio_root_password=None, + label=None, ): self.location = location self.table_type = table_type + # Stable name for the scenario name: the table names carry getuid(). + self.label = label self.cluster_name = cluster_name self.minio_root_user = minio_root_user or JoinTable.minio_root_user self.minio_root_password = minio_root_password or JoinTable.minio_root_password @@ -538,7 +541,7 @@ def join_clause(self, minio_root_user, minio_root_password, node=None): ) with Given("create iceberg tables in different locations"): - for location in locations: + for i, location in enumerate(locations, 1): _, table_name, namespace = ( swarm_steps.iceberg_table_with_all_basic_data_types( minio_root_user=minio_root_user, @@ -551,6 +554,7 @@ def join_clause(self, minio_root_user, minio_root_password, node=None): database_name=database_name, namespace=namespace, table_name=table_name, + label=f"iceberg_table_data{i}", ) ) @@ -561,37 +565,46 @@ def join_clause(self, minio_root_user, minio_root_password, node=None): url, minio_root_user=minio_root_user, minio_root_password=minio_root_password, + label=f"s3_data{i}", ) - for url in urls + for i, url in enumerate(urls, 1) ] iceberg_table_functions = [ JoinTable.create_iceberg_table_function( url, minio_root_user=minio_root_user, minio_root_password=minio_root_password, + label=f"iceberg_data{i}", ) - for url in urls + for i, url in enumerate(urls, 1) ] iceberg_s3_cluster_table_functions = [ - JoinTable.create_icebergS3Cluster_table_function(url, "replicated_cluster") - for url in urls + JoinTable.create_icebergS3Cluster_table_function( + url, "replicated_cluster", label=f"icebergS3Cluster_data{i}" + ) + for i, url in enumerate(urls, 1) ] s3_cluster_table_functions = [ - JoinTable.create_s3Cluster_table_function(url, "replicated_cluster") - for url in urls + JoinTable.create_s3Cluster_table_function( + url, "replicated_cluster", label=f"s3Cluster_data{i}" + ) + for i, url in enumerate(urls, 1) ] with Given( "create merge tree tables from iceberg tables with same schema and data" ): merge_tree_tables = [] - for iceberg_table in iceberg_tables: + for i, iceberg_table in enumerate(iceberg_tables, 1): merge_tree_table_name = f"merge_tree_table_{getuid()}" create_table_as_select( as_select_from=iceberg_table, table_name=merge_tree_table_name ) merge_tree_tables.append( - JoinTable.create_merge_tree_table(table_name=merge_tree_table_name) + JoinTable.create_merge_tree_table( + table_name=merge_tree_table_name, + label=f"merge_tree_table_data{i}", + ) ) left_tables = ( @@ -652,19 +665,6 @@ def join_clause(self, minio_root_user, minio_root_password, node=None): "FULL OUTER JOIN", ] - length = len( - list( - product( - left_tables, - right_tables, - modes, - join_conditions, - object_storage_clusters, - join_clauses, - ) - ) - ) - all_possible_combinations = list( product( left_tables, @@ -676,21 +676,38 @@ def join_clause(self, minio_root_user, minio_root_password, node=None): ) ) + # A scenario is named after its combination: "join N of " is the + # position of the combination in the full product, so "join N" is the + # same combination in every run. Each run executes a different sample of + # combinations, drawn from the run's seed (see regression.py), so the + # numbers of a run are not sequential. To reproduce a failure, use the + # seed printed below, otherwise "join N" is probably not in the sample: + # --seed --only "/swarms/feature/swarm joins/join clause/join N of*" + total = len(all_possible_combinations) + combination_ids = range(total) + if not self.context.stress: - all_possible_combinations = random.sample( - all_possible_combinations, min(1000, len(all_possible_combinations)) + note(f"join combinations seed: {self.context.seed}") + combination_ids = sorted( + random.Random(self.context.seed).sample(combination_ids, min(1000, total)) ) with Pool() as pool: - for num, ( - left_table, - right_table, - mode, - join_condition, - object_storage_cluster, - join_clause, - ) in enumerate(all_possible_combinations): - name = f"join {num} of {length}: {left_table} with {right_table} in {mode} mode on {object_storage_cluster} cluster with {join_clause} clause" + for combination_id in combination_ids: + ( + left_table, + right_table, + mode, + join_condition, + object_storage_cluster, + join_clause, + ) = all_possible_combinations[combination_id] + condition = join_condition.replace("t1.", "").replace("t2.", "") + name = ( + f"join {combination_id} of {total}: {left_table.label} with {right_table.label}" + f" in {mode} mode on {object_storage_cluster} cluster" + f" with {join_clause} on {condition}" + ) Scenario( name=name, test=check_join, diff --git a/swarms/tests/union.py b/swarms/tests/union.py index ec9321326..5b974297a 100644 --- a/swarms/tests/union.py +++ b/swarms/tests/union.py @@ -27,9 +27,12 @@ def __init__( cluster_name=None, minio_root_user=None, minio_root_password=None, + label=None, ): self.location = location self.table_type = table_type + # Stable name for the scenario name: the table names carry getuid(). + self.label = label self.cluster_name = cluster_name self.minio_root_user = minio_root_user or JoinTable.minio_root_user self.minio_root_password = minio_root_password or JoinTable.minio_root_password @@ -154,18 +157,27 @@ def check_union( left_table, right_table, node=None, - object_storage_cluster_1="replicated_cluster", - object_storage_cluster_2="replicated_cluster", + all_or_distinct="UNION ALL", + object_storage_cluster_1=None, + object_storage_cluster_2=None, order_by="tuple(*)", ): """Check union operation.""" if node is None: node = self.context.node - with Given("choose to run query with UNION ALL or UNION DISTINCT"): - all_or_distinct = random.choice(["UNION ALL", "UNION DISTINCT"]) + object_storage_cluster_setting_1 = ( + f"SETTINGS object_storage_cluster = '{object_storage_cluster_1}'" + if object_storage_cluster_1 + else "" + ) + object_storage_cluster_setting_2 = ( + f"SETTINGS object_storage_cluster = '{object_storage_cluster_2}'" + if object_storage_cluster_2 + else "" + ) - with And("create merge tree tables as left and right tables"): + with Given("create merge tree tables as left and right tables"): left_merge_tree_table = create_table_as_select( as_select_from=left_table, ) @@ -184,12 +196,6 @@ def check_union( expected_result = node.query(query) with Then("check UNION on iceberg tables and cluster functions"): - object_storage_cluster_setting_1 = random.choice( - ["", f"SETTINGS object_storage_cluster = '{object_storage_cluster_1}'"] - ) - object_storage_cluster_setting_2 = random.choice( - ["", f"SETTINGS object_storage_cluster = '{object_storage_cluster_2}'"] - ) query = _union_query( self, left_table, @@ -229,7 +235,7 @@ def union_clause(self, minio_root_user, minio_root_password, node=None): ) with Given("create iceberg tables in different locations"): - for location in locations: + for i, location in enumerate(locations, 1): _, table_name, namespace = ( swarm_steps.iceberg_table_with_all_basic_data_types( minio_root_user=minio_root_user, @@ -242,6 +248,7 @@ def union_clause(self, minio_root_user, minio_root_password, node=None): database_name=database_name, namespace=namespace, table_name=table_name, + label=f"iceberg_table_data{i}", ) ) @@ -252,37 +259,46 @@ def union_clause(self, minio_root_user, minio_root_password, node=None): url, minio_root_user=minio_root_user, minio_root_password=minio_root_password, + label=f"s3_data{i}", ) - for url in urls + for i, url in enumerate(urls, 1) ] iceberg_table_functions = [ JoinTable.create_iceberg_table_function( url, minio_root_user=minio_root_user, minio_root_password=minio_root_password, + label=f"iceberg_data{i}", ) - for url in urls + for i, url in enumerate(urls, 1) ] iceberg_s3_cluster_table_functions = [ - JoinTable.create_icebergS3Cluster_table_function(url, "replicated_cluster") - for url in urls + JoinTable.create_icebergS3Cluster_table_function( + url, "replicated_cluster", label=f"icebergS3Cluster_data{i}" + ) + for i, url in enumerate(urls, 1) ] s3_cluster_table_functions = [ - JoinTable.create_s3Cluster_table_function(url, "replicated_cluster") - for url in urls + JoinTable.create_s3Cluster_table_function( + url, "replicated_cluster", label=f"s3Cluster_data{i}" + ) + for i, url in enumerate(urls, 1) ] with Given( "create merge tree tables from iceberg tables with same schema and data" ): merge_tree_tables = [] - for iceberg_table in iceberg_tables: + for i, iceberg_table in enumerate(iceberg_tables, 1): merge_tree_table_name = f"merge_tree_table_{getuid()}" create_table_as_select( as_select_from=iceberg_table, table_name=merge_tree_table_name ) merge_tree_tables.append( - JoinTable.create_merge_tree_table(table_name=merge_tree_table_name) + JoinTable.create_merge_tree_table( + table_name=merge_tree_table_name, + label=f"merge_tree_table_data{i}", + ) ) with And("define all possible left and right tables combinations for union"): @@ -305,6 +321,7 @@ def union_clause(self, minio_root_user, minio_root_password, node=None): with And("define object storage clusters options"): object_storage_clusters = [ + None, "replicated_cluster_three_nodes", "replicated_cluster_two_nodes", "replicated_cluster_two_nodes_version_2", @@ -317,24 +334,41 @@ def union_clause(self, minio_root_user, minio_root_password, node=None): product( left_tables, right_tables, + ["UNION ALL", "UNION DISTINCT"], object_storage_clusters, object_storage_clusters, ) ) - if not self.context.stress: - all_possible_combinations = random.sample(all_possible_combinations, 500) + # A scenario is named after its combination: "union N of " is the + # position of the combination in the full product, so "union N" is the + # same combination in every run. Each run executes a different sample of + # combinations, drawn from the run's seed (see regression.py), so the + # numbers of a run are not sequential. To reproduce a failure, use the + # seed printed below, otherwise "union N" is probably not in the sample: + # --seed --only "/swarms/feature/swarm union/union clause/union N of*" + total = len(all_possible_combinations) + combination_ids = range(total) - length = len(all_possible_combinations) + if not self.context.stress: + note(f"union combinations seed: {self.context.seed}") + combination_ids = sorted( + random.Random(self.context.seed).sample(combination_ids, min(500, total)) + ) with Pool(10) as pool: - for num, ( - left_table, - right_table, - object_storage_cluster_1, - object_storage_cluster_2, - ) in enumerate(all_possible_combinations): - name = f"union {num} of {length}: {left_table} with {right_table} in {object_storage_cluster_1} cluster and {object_storage_cluster_2} cluster" + for combination_id in combination_ids: + ( + left_table, + right_table, + all_or_distinct, + object_storage_cluster_1, + object_storage_cluster_2, + ) = all_possible_combinations[combination_id] + name = ( + f"union {combination_id} of {total}: {left_table.label} {all_or_distinct} {right_table.label}" + f" in {object_storage_cluster_1} cluster and {object_storage_cluster_2} cluster" + ) Scenario( name=name, test=check_union, @@ -343,6 +377,7 @@ def union_clause(self, minio_root_user, minio_root_password, node=None): )( left_table=left_table, right_table=right_table, + all_or_distinct=all_or_distinct, object_storage_cluster_1=object_storage_cluster_1, object_storage_cluster_2=object_storage_cluster_2, )