From ca60cfdd16515ccb47261f44b3bcdfd29bf630e5 Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 2 Oct 2026 03:05:37 +0000 Subject: [PATCH 1/4] fix(asap-tools): run cluster_data_exporter in local mode and find Alibaba data files The service asserted exactly one worker node, which local mode (0 nodes) never satisfies, and validated Alibaba data against MsResource_*.csv.gz, a name the exporter never reads, so valid MSMetrics/NodeMetrics data was rejected. Co-Authored-By: Claude Opus 5.5 --- .../services/cluster_data_exporter.py | 38 ++++++------- .../test_cluster_data_exporter_service.py | 56 +++++++++++++++++++ 2 files changed, 73 insertions(+), 21 deletions(-) create mode 100644 asap-tools/experiments/tests/test_cluster_data_exporter_service.py diff --git a/asap-tools/experiments/experiment_utils/services/cluster_data_exporter.py b/asap-tools/experiments/experiment_utils/services/cluster_data_exporter.py index e855c675..4eb93679 100644 --- a/asap-tools/experiments/experiment_utils/services/cluster_data_exporter.py +++ b/asap-tools/experiments/experiment_utils/services/cluster_data_exporter.py @@ -71,10 +71,11 @@ def start( # Get number of nodes from provider (assuming it has this info) num_nodes = kwargs.get("num_nodes", 1) - # Assert that we have exactly 2 nodes - assert num_nodes == 1, ( - f"cluster_data_exporter requires exactly 1 node (num_nodes==1), " - f"got {num_nodes}" + # One worker node next to the coordinator; local mode reports 0 nodes + # and runs everything on this machine. + assert num_nodes in (0, 1), ( + f"cluster_data_exporter requires one worker node (num_nodes==1) " + f"or local mode (num_nodes==0), got {num_nodes}" ) # Extract configuration @@ -400,23 +401,18 @@ def _validate_alibaba_data( Raises: ValueError: If required files are missing """ - # Determine expected file pattern based on data type and year - if data_type == "node": - if data_year == 2021 or data_year == 2022: - pattern = "Node_*.csv.gz" - else: - raise ValueError( - f"Invalid data_year for Alibaba node data: {data_year}" - ) - elif data_type == "msresource": - if data_year == 2021 or data_year == 2022: - pattern = "MsResource_*.csv.gz" - else: - raise ValueError( - f"Invalid data_year for Alibaba msresource data: {data_year}" - ) - else: - raise ValueError(f"Invalid data_type for Alibaba: {data_type}") + # File names the exporter reads (alibaba_metrics/{node,ms_resource}.rs). + patterns = { + ("node", 2021): "Node_*.csv.gz", + ("node", 2022): "NodeMetrics_*.csv.gz", + ("msresource", 2021): "MSResource_*.csv.gz", + ("msresource", 2022): "MSMetrics_*.csv.gz", + } + pattern = patterns.get((data_type, int(data_year))) + if pattern is None: + raise ValueError( + f"Invalid Alibaba data_type/data_year: {data_type}/{data_year}" + ) # Check for data files on remote node target_node = self.node_offset + 1 diff --git a/asap-tools/experiments/tests/test_cluster_data_exporter_service.py b/asap-tools/experiments/tests/test_cluster_data_exporter_service.py new file mode 100644 index 00000000..8dca383e --- /dev/null +++ b/asap-tools/experiments/tests/test_cluster_data_exporter_service.py @@ -0,0 +1,56 @@ +"""Tests for ClusterDataExporterService data-file validation and node checks.""" + +import subprocess +import unittest + +from experiment_utils.services.cluster_data_exporter import ClusterDataExporterService + + +class RecordingProvider: + """Answers every command with one matching file and records the commands.""" + + def __init__(self): + self.commands = [] + + def execute_command(self, node_idx, cmd, cmd_dir, nohup, popen): + self.commands.append(cmd) + return subprocess.CompletedProcess([], 0, "1\n", "") + + +class ValidateAlibabaDataTest(unittest.TestCase): + def _counted_pattern(self, data_type, data_year): + provider = RecordingProvider() + service = ClusterDataExporterService(provider, 0, "/traces") + service._validate_alibaba_data(data_type, data_year) + return [c for c in provider.commands if "wc -l" in c][0] + + def test_patterns_match_exporter_file_names(self): + # The check used to look for MsResource_*.csv.gz for both years, a + # name the exporter never reads, so valid data was rejected. + cases = { + ("node", 2021): "Node_*.csv.gz", + ("node", 2022): "NodeMetrics_*.csv.gz", + ("msresource", 2021): "MSResource_*.csv.gz", + ("msresource", 2022): "MSMetrics_*.csv.gz", + } + for (data_type, data_year), pattern in cases.items(): + with self.subTest(data_type=data_type, data_year=data_year): + self.assertIn( + f"/traces/{pattern}", self._counted_pattern(data_type, data_year) + ) + + def test_unknown_year_raises(self): + service = ClusterDataExporterService(RecordingProvider(), 0, "/traces") + with self.assertRaises(ValueError): + service._validate_alibaba_data("msresource", 2020) + + +class NodeCountTest(unittest.TestCase): + def test_more_than_one_worker_node_is_rejected(self): + service = ClusterDataExporterService(RecordingProvider(), 0, "/traces") + with self.assertRaises(AssertionError): + service.start({"provider": "google"}, "/out", "/local", num_nodes=2) + + +if __name__ == "__main__": + unittest.main() From 24db2ba0dcaefc4e1c1aa10e3d0ac50026434133 Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 2 Oct 2026 03:57:18 +0000 Subject: [PATCH 2/4] fix(asap-tools): run e2e experiments with the local provider The exporter health poll aborted on the first failed curl because the local provider raises on non-zero exits, and the monitor keyword list kept its SSH-escaped quotes locally, so no process matched and the query client never ran. Co-Authored-By: Claude Opus 5.5 --- .../services/cluster_data_exporter.py | 3 + .../services/remote_monitor_service.py | 14 +++-- .../tests/test_remote_monitor_keywords.py | 61 +++++++++++++++++++ 3 files changed, 74 insertions(+), 4 deletions(-) create mode 100644 asap-tools/experiments/tests/test_remote_monitor_keywords.py diff --git a/asap-tools/experiments/experiment_utils/services/cluster_data_exporter.py b/asap-tools/experiments/experiment_utils/services/cluster_data_exporter.py index 4eb93679..4eb73596 100644 --- a/asap-tools/experiments/experiment_utils/services/cluster_data_exporter.py +++ b/asap-tools/experiments/experiment_utils/services/cluster_data_exporter.py @@ -493,12 +493,15 @@ def _wait_for_health( while time.time() - start_time < timeout: # Run curl from the remote node to check health check_cmd = f"curl -s -o /dev/null -w '%{{http_code}}' {url}" + # curl exits non-zero while the exporter is still starting; the + # local provider raises on that, so keep polling instead. result = self.provider.execute_command( node_idx=node_idx, cmd=check_cmd, cmd_dir="", nohup=False, popen=False, + ignore_errors=True, ) try: diff --git a/asap-tools/experiments/experiment_utils/services/remote_monitor_service.py b/asap-tools/experiments/experiment_utils/services/remote_monitor_service.py index 1e10ceb8..80912449 100644 --- a/asap-tools/experiments/experiment_utils/services/remote_monitor_service.py +++ b/asap-tools/experiments/experiment_utils/services/remote_monitor_service.py @@ -100,13 +100,19 @@ def start( else: keywords.append(constants.QUERY_ENGINE_RS_PROCESS_KEYWORD) + # Over SSH the command is parsed by two shells, so the quotes that keep + # the keyword list one argument must be escaped once more. + keywords_arg = '"{}"'.format(",".join(keywords)) + if self.provider.is_remote(): + keywords_arg = r"\"{}\"".format(",".join(keywords)) + if use_timed_mode: # Build command for timed mode (skip_querying) cmd = ( "python3 -u remote_monitor.py " "--execution_mode timed " "--experiment_mode {} " - r"--keywords \"{}\" " + "--keywords {} " "--config_file {} " "--experiment_output_dir {} " "--monitor_output_file {} " @@ -116,7 +122,7 @@ def start( "--monitor_interval_seconds {} " ).format( experiment_mode, - ",".join(keywords), + keywords_arg, os.path.join( os.path.dirname(experiment_output_dir), "controller_client_configs", @@ -162,7 +168,7 @@ def start( "python3 -u remote_monitor.py " "--execution_mode prometheus_client " "--experiment_mode {} " - r"--keywords \"{}\" " + "--keywords {} " "--config_file {} " "--experiment_output_dir {} " "--monitor_output_file {} " @@ -172,7 +178,7 @@ def start( "--monitor_interval_seconds {} " ).format( experiment_mode, - ",".join(keywords), + keywords_arg, os.path.join( os.path.dirname(experiment_output_dir), "controller_client_configs", diff --git a/asap-tools/experiments/tests/test_remote_monitor_keywords.py b/asap-tools/experiments/tests/test_remote_monitor_keywords.py new file mode 100644 index 00000000..48bb1634 --- /dev/null +++ b/asap-tools/experiments/tests/test_remote_monitor_keywords.py @@ -0,0 +1,61 @@ +"""Tests that remote_monitor.py receives its keyword list as one clean argument.""" + +import shlex +import unittest + +from experiment_utils.services.remote_monitor_service import RemoteMonitorService + + +class RecordingProvider: + def __init__(self, remote): + self.remote = remote + self.commands = [] + + def is_remote(self): + return self.remote + + def get_home_dir(self): + return "/home" + + def execute_command(self, **kwargs): + self.commands.append(kwargs["cmd"]) + + +def keywords_seen_by_monitor(remote): + provider = RecordingProvider(remote) + RemoteMonitorService(provider, 0).start( + controller_client_config="/out/controller_client_configs/baseline.yaml", + experiment_output_dir="/out/baseline", + experiment_mode="baseline", + profile_query_engine=False, + profile_prometheus_time=None, + manual_mode=False, + streaming_engine="precompute", + query_engine_service=None, + controller_remote_output_dir="/out/controller_output", + use_container_prometheus_client=True, + prometheus_client_parallel=False, + backend_protocol="prometheus", + pre_query_wait_seconds=0, + monitor_interval_seconds=1.0, + timed_duration=10, + ) + args = shlex.split(provider.commands[0]) + if remote: + # The SSH provider wraps the command in one more shell. + args = shlex.split(" ".join(args)) + return args[args.index("--keywords") + 1] + + +class KeywordQuotingTest(unittest.TestCase): + def test_local_provider_gets_unescaped_keywords(self): + # Escaped quotes meant for SSH used to reach remote_monitor.py + # verbatim in local mode, so no process matched and no queries ran. + self.assertEqual(keywords_seen_by_monitor(remote=False), "prometheus.yml") + + def test_remote_provider_keeps_escaped_quotes_for_ssh(self): + self.assertEqual(keywords_seen_by_monitor(remote=True), "prometheus.yml") + + +if __name__ == "__main__": + unittest.main() From 4b480a4311725e6dd3aa451f27179b81eaf9e53d Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 2 Oct 2026 03:57:18 +0000 Subject: [PATCH 3/4] fix(asap-tools): keep series that miss some repetitions in fidelity results get_all_timeseries raised KeyError when a key was absent from one repetition, e.g. a series that appears mid-run, so calculate_fidelity crashed on cluster-trace runs. Co-Authored-By: Claude Opus 5.5 --- .../promql_utilities/query_results/classes.py | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/asap-common/dependencies/py/promql_utilities/promql_utilities/query_results/classes.py b/asap-common/dependencies/py/promql_utilities/promql_utilities/query_results/classes.py index 7f43d742..f8a4bfaf 100644 --- a/asap-common/dependencies/py/promql_utilities/promql_utilities/query_results/classes.py +++ b/asap-common/dependencies/py/promql_utilities/promql_utilities/query_results/classes.py @@ -65,8 +65,10 @@ def get_all_timeseries(self) -> Dict[frozenset, TimeSeries]: for k in keys: for repetition_idx, result in enumerate(self.query_results): + # A key absent from one repetition (e.g. a series that appears + # mid-run) stays None for that repetition. if result.result: - intermediate_ret[k][repetition_idx] = result.result[k] + intermediate_ret[k][repetition_idx] = result.result.get(k) ret[k] = TimeSeries(k, intermediate_ret[k]) From a10ff471804cfddcf1d955cb9f8d588d3cb5b17b Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 2 Oct 2026 04:53:54 +0000 Subject: [PATCH 4/4] fix(asap-tools): pass ms-resource to the Alibaba cluster_data_exporter The exporter's --data-type accepts node or ms-resource; the config value msresource was passed through, so the container exited at startup. Co-Authored-By: Claude Opus 5.5 --- .../services/cluster_data_exporter.py | 7 ++++++- .../tests/test_cluster_data_exporter_service.py | 14 ++++++++++++++ 2 files changed, 20 insertions(+), 1 deletion(-) diff --git a/asap-tools/experiments/experiment_utils/services/cluster_data_exporter.py b/asap-tools/experiments/experiment_utils/services/cluster_data_exporter.py index 4eb73596..ada32405 100644 --- a/asap-tools/experiments/experiment_utils/services/cluster_data_exporter.py +++ b/asap-tools/experiments/experiment_utils/services/cluster_data_exporter.py @@ -13,6 +13,10 @@ from experiment_utils.providers.base import InfrastructureProvider +# Experiment-config data_type -> the exporter's --data-type value. +ALIBABA_DATA_TYPE_CLI_VALUES = {"node": "node", "msresource": "ms-resource"} + + class ClusterDataExporterService(BaseService): """ Service for managing cluster_data_exporter via Docker. @@ -259,7 +263,8 @@ def _build_docker_command( elif provider == "alibaba": if "data_type" in config: - cmd_parts.append(f"--data-type={config['data_type']}") + data_type = ALIBABA_DATA_TYPE_CLI_VALUES[config["data_type"]] + cmd_parts.append(f"--data-type={data_type}") if "data_year" in config: cmd_parts.append(f"--data-year={config['data_year']}") diff --git a/asap-tools/experiments/tests/test_cluster_data_exporter_service.py b/asap-tools/experiments/tests/test_cluster_data_exporter_service.py index 8dca383e..9aa322a2 100644 --- a/asap-tools/experiments/tests/test_cluster_data_exporter_service.py +++ b/asap-tools/experiments/tests/test_cluster_data_exporter_service.py @@ -45,6 +45,20 @@ def test_unknown_year_raises(self): service._validate_alibaba_data("msresource", 2020) +class DockerCommandTest(unittest.TestCase): + def test_msresource_uses_exporter_cli_value(self): + # The exporter's clap enum spells it ms-resource; passing the config + # value through made the container exit before serving metrics. + service = ClusterDataExporterService(RecordingProvider(), 0, "/traces") + service.container_name = "cde" + cmd = service._build_docker_command( + {"provider": "alibaba", "data_type": "msresource", "data_year": 2022}, + port=40000, + output_dir="/out", + ) + self.assertIn("--data-type=ms-resource", cmd) + + class NodeCountTest(unittest.TestCase): def test_more_than_one_worker_node_is_rejected(self): service = ClusterDataExporterService(RecordingProvider(), 0, "/traces")