From e0ad253a9e7c830284fb7bd4c2cf3da8ef230555 Mon Sep 17 00:00:00 2001 From: CarlosFelipeOR Date: Wed, 23 Sep 2026 19:53:22 -0300 Subject: [PATCH 1/2] swarms: deterministic join/union scenario names, seeded sampling Scenario names embedded the tables' SQL, which carries getuid() table names, so every run produced different names. Names are now just "join N of 1000" / "union N of 500"; the parameters are in the scenario's arguments in the log. The combination sample is drawn from a per-run seed (random, or --seed) with its own random.Random, so each run covers different combinations and any run can be reproduced. The union per-scenario choices (UNION ALL/DISTINCT, object_storage_cluster settings) are drawn in the feature instead of inside the parallel scenarios. Remove the "join 455 of 816480" xfail for #1244: it had already stopped matching the intended combination. Co-Authored-By: Claude Opus 5.5 (1M context) --- swarms/regression.py | 21 ++++++++++++--- swarms/tests/joins.py | 29 ++++++++++---------- swarms/tests/union.py | 63 ++++++++++++++++++++++++++++++++----------- 3 files changed, 78 insertions(+), 35 deletions(-) diff --git a/swarms/regression.py b/swarms/regression.py index 0ae807492..479537fa6 100755 --- a/swarms/regression.py +++ b/swarms/regression.py @@ -18,13 +18,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: 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 +74,7 @@ @Name("swarms") @FFails(ffails) @XFails(xfails) -@ArgumentParser(argparser_minio) +@ArgumentParser(argparser) @Specifications(SRS_044_Swarm_Cluster_Query_Execution) @CaptureClusterArgs @CaptureMinioArgs @@ -75,6 +85,7 @@ def regression( stress=None, with_analyzer=False, minio_args=None, + seed=None, ): """Run tests for Swarm clusters.""" nodes = { @@ -92,6 +103,8 @@ def regression( if stress is not None: self.context.stress = stress + self.context.seed = 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..ecf8423e1 100644 --- a/swarms/tests/joins.py +++ b/swarms/tests/joins.py @@ -652,19 +652,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, @@ -677,10 +664,22 @@ def join_clause(self, minio_root_user, minio_root_password, node=None): ) if not self.context.stress: - all_possible_combinations = random.sample( + # Each run samples a different set of combinations, so "join N" names + # a different combination in every run (and in a CI rerun). Its + # parameters are in the scenario's arguments in the log; to reproduce + # a failure, rerun with the seed printed below: + # --seed --only "/swarms/feature/swarm joins/join clause/join N of*" + seed = self.context.seed + if seed is None: + seed = random.SystemRandom().randrange(2**32) + note(f"join combinations seed: {seed}") + + all_possible_combinations = random.Random(seed).sample( all_possible_combinations, min(1000, len(all_possible_combinations)) ) + length = len(all_possible_combinations) + with Pool() as pool: for num, ( left_table, @@ -690,7 +689,7 @@ def join_clause(self, minio_root_user, minio_root_password, node=None): 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" + name = f"join {num} of {length}" Scenario( name=name, test=check_join, diff --git a/swarms/tests/union.py b/swarms/tests/union.py index ec9321326..bcff5b04f 100644 --- a/swarms/tests/union.py +++ b/swarms/tests/union.py @@ -154,18 +154,16 @@ 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_setting_1="", + object_storage_cluster_setting_2="", 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"]) - - 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 +182,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, @@ -322,19 +314,57 @@ def union_clause(self, minio_root_user, minio_root_password, node=None): ) ) + # Each run samples a different set of combinations, so "union N" names + # a different combination in every run (and in a CI rerun). Its + # parameters are in the scenario's arguments in the log; to reproduce + # a failure, rerun with the seed printed below: + # --seed --only "/swarms/feature/swarm union/union clause/union N of*" + seed = self.context.seed + if seed is None: + seed = random.SystemRandom().randrange(2**32) + note(f"union combinations seed: {seed}") + rng = random.Random(seed) + if not self.context.stress: - all_possible_combinations = random.sample(all_possible_combinations, 500) + all_possible_combinations = rng.sample(all_possible_combinations, 500) length = len(all_possible_combinations) + # Draw the per-scenario choices here, in order, for every combination: + # drawing them inside the parallel scenarios would not be reproducible. + all_possible_combinations = [ + ( + left_table, + right_table, + object_storage_cluster_1, + object_storage_cluster_2, + rng.choice(["UNION ALL", "UNION DISTINCT"]), + rng.choice( + ["", f"SETTINGS object_storage_cluster = '{object_storage_cluster_1}'"] + ), + rng.choice( + ["", f"SETTINGS object_storage_cluster = '{object_storage_cluster_2}'"] + ), + ) + for ( + left_table, + right_table, + object_storage_cluster_1, + object_storage_cluster_2, + ) in all_possible_combinations + ] + with Pool(10) as pool: for num, ( left_table, right_table, object_storage_cluster_1, object_storage_cluster_2, + all_or_distinct, + object_storage_cluster_setting_1, + object_storage_cluster_setting_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" + name = f"union {num} of {length}" Scenario( name=name, test=check_union, @@ -343,8 +373,9 @@ def union_clause(self, minio_root_user, minio_root_password, node=None): )( left_table=left_table, right_table=right_table, - object_storage_cluster_1=object_storage_cluster_1, - object_storage_cluster_2=object_storage_cluster_2, + all_or_distinct=all_or_distinct, + object_storage_cluster_setting_1=object_storage_cluster_setting_1, + object_storage_cluster_setting_2=object_storage_cluster_setting_2, ) join() From 834c93a72835f286fa1c6d8e95ee440dfb8e282f Mon Sep 17 00:00:00 2001 From: CarlosFelipeOR Date: Tue, 29 Sep 2026 16:53:01 -0300 Subject: [PATCH 2/2] swarms: name join/union scenarios by combination, seed from GITHUB_RUN_ID "join N of " is now the combination's position in the full product plus its parameters, so the same name is always the same combination. The seed defaults to GITHUB_RUN_ID, so all jobs of a CI run and their reruns sample the same combinations. Union ALL/DISTINCT and object_storage_cluster (None = unset) are now part of the combination. Co-Authored-By: Claude Opus 5.5 (1M context) --- swarms/regression.py | 12 ++++- swarms/tests/joins.py | 82 ++++++++++++++++++----------- swarms/tests/union.py | 120 ++++++++++++++++++++++-------------------- 3 files changed, 122 insertions(+), 92 deletions(-) diff --git a/swarms/regression.py b/swarms/regression.py index 479537fa6..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 * @@ -28,7 +30,7 @@ def argparser(parser): action="store", type=int, default=None, - help="seed for sampling the join and union combinations, default: random (printed in the log)", + 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 * @@ -103,7 +105,13 @@ def regression( if stress is not None: self.context.stress = stress - self.context.seed = seed + # 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 ecf8423e1..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 = ( @@ -663,33 +676,38 @@ def join_clause(self, minio_root_user, minio_root_password, node=None): ) ) - if not self.context.stress: - # Each run samples a different set of combinations, so "join N" names - # a different combination in every run (and in a CI rerun). Its - # parameters are in the scenario's arguments in the log; to reproduce - # a failure, rerun with the seed printed below: - # --seed --only "/swarms/feature/swarm joins/join clause/join N of*" - seed = self.context.seed - if seed is None: - seed = random.SystemRandom().randrange(2**32) - note(f"join combinations seed: {seed}") + # 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) - all_possible_combinations = random.Random(seed).sample( - all_possible_combinations, min(1000, len(all_possible_combinations)) + if not self.context.stress: + note(f"join combinations seed: {self.context.seed}") + combination_ids = sorted( + random.Random(self.context.seed).sample(combination_ids, min(1000, total)) ) - length = len(all_possible_combinations) - 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}" + 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 bcff5b04f..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 @@ -155,14 +158,25 @@ def check_union( right_table, node=None, all_or_distinct="UNION ALL", - object_storage_cluster_setting_1="", - object_storage_cluster_setting_2="", + object_storage_cluster_1=None, + object_storage_cluster_2=None, order_by="tuple(*)", ): """Check union operation.""" if node is None: node = self.context.node + 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 Given("create merge tree tables as left and right tables"): left_merge_tree_table = create_table_as_select( as_select_from=left_table, @@ -221,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, @@ -234,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}", ) ) @@ -244,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"): @@ -297,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", @@ -309,62 +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, ) ) - # Each run samples a different set of combinations, so "union N" names - # a different combination in every run (and in a CI rerun). Its - # parameters are in the scenario's arguments in the log; to reproduce - # a failure, rerun with the seed printed below: + # 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*" - seed = self.context.seed - if seed is None: - seed = random.SystemRandom().randrange(2**32) - note(f"union combinations seed: {seed}") - rng = random.Random(seed) + total = len(all_possible_combinations) + combination_ids = range(total) if not self.context.stress: - all_possible_combinations = rng.sample(all_possible_combinations, 500) - - length = len(all_possible_combinations) - - # Draw the per-scenario choices here, in order, for every combination: - # drawing them inside the parallel scenarios would not be reproducible. - all_possible_combinations = [ - ( - left_table, - right_table, - object_storage_cluster_1, - object_storage_cluster_2, - rng.choice(["UNION ALL", "UNION DISTINCT"]), - rng.choice( - ["", f"SETTINGS object_storage_cluster = '{object_storage_cluster_1}'"] - ), - rng.choice( - ["", f"SETTINGS object_storage_cluster = '{object_storage_cluster_2}'"] - ), + note(f"union combinations seed: {self.context.seed}") + combination_ids = sorted( + random.Random(self.context.seed).sample(combination_ids, min(500, total)) ) - for ( - left_table, - right_table, - object_storage_cluster_1, - object_storage_cluster_2, - ) in all_possible_combinations - ] with Pool(10) as pool: - for num, ( - left_table, - right_table, - object_storage_cluster_1, - object_storage_cluster_2, - all_or_distinct, - object_storage_cluster_setting_1, - object_storage_cluster_setting_2, - ) in enumerate(all_possible_combinations): - name = f"union {num} of {length}" + 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, @@ -374,8 +378,8 @@ 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_setting_1=object_storage_cluster_setting_1, - object_storage_cluster_setting_2=object_storage_cluster_setting_2, + object_storage_cluster_1=object_storage_cluster_1, + object_storage_cluster_2=object_storage_cluster_2, ) join()