From 02d04aa8dc25521d201124f7d4b116db20466589 Mon Sep 17 00:00:00 2001 From: Milind Srivastava Date: Tue, 18 Aug 2026 14:22:12 -0400 Subject: [PATCH 1/3] fix(tools): make Args the single source of truth for remote_write_ip MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit cfg.streaming.remote_write.ip has been a dead null since PR #469 replaced the OmegaConf resolver that used to populate it with a plain per-script `args.remote_write_ip = provider.get_node_ip(...)` — a value only arroyo.py ever read. generate_prometheus_config still reads cfg.streaming.remote_write.ip directly, so it silently baked http://None:/... into prometheus.yml for every non-arroyo streaming engine, causing remote_write to fail and precompute-engine queries to hang forever waiting on data that never arrives (#546). Args.__init__ now builds the provider itself and sets both self.remote_write_ip and cfg.streaming.remote_write.ip from that one place, so every script gets it automatically instead of each one re-deriving (and potentially forgetting) it. Scripts that used to call create_provider(cfg) right alongside config.Args(cfg) now just read args.provider. Co-Authored-By: Claude Sonnet 5 --- .../experiment_only_ingest_path.py | 7 ++----- asap-tools/experiments/experiment_run_e2e.py | 7 ++----- ...experiment_run_exporters_and_prometheus.py | 6 ++---- .../experiment_run_grafana_demo.py | 7 ++----- .../experiment_teardown_everything.py | 5 +---- .../experiments/experiment_utils/config.py | 20 +++++++++++++++---- 6 files changed, 25 insertions(+), 27 deletions(-) diff --git a/asap-tools/experiments/experiment_only_ingest_path.py b/asap-tools/experiments/experiment_only_ingest_path.py index 942e72f..e03bc07 100644 --- a/asap-tools/experiments/experiment_only_ingest_path.py +++ b/asap-tools/experiments/experiment_only_ingest_path.py @@ -8,7 +8,6 @@ import constants import experiment_utils from experiment_utils import sync, config -from experiment_utils.providers.factory import create_provider from experiment_utils.services import ( KafkaService, ExporterServiceFactory, @@ -74,11 +73,9 @@ def main(cfg: DictConfig): is_v2 = experiment_mode == constants.SKETCHDB_EXPERIMENT_NAME # Convert config to args-like object for backward compatibility + # (also constructs the infrastructure provider, exposed as args.provider) args = config.Args(cfg) - - # Create infrastructure provider - provider = create_provider(cfg) - args.remote_write_ip = provider.get_node_ip(args.node_offset) + provider = args.provider local_experiment_root_dir = os.path.join( constants.LOCAL_EXPERIMENT_DIR, args.experiment_name diff --git a/asap-tools/experiments/experiment_run_e2e.py b/asap-tools/experiments/experiment_run_e2e.py index 6edbf07..41cdd45 100644 --- a/asap-tools/experiments/experiment_run_e2e.py +++ b/asap-tools/experiments/experiment_run_e2e.py @@ -8,7 +8,6 @@ import constants import experiment_utils from experiment_utils import sync, config -from experiment_utils.providers.factory import create_provider from experiment_utils.services import ( KafkaService, FlinkService, @@ -55,12 +54,10 @@ def main(cfg: DictConfig): # Validate experiment configuration config.validate_experiment_config(cfg.experiment_params) - # Create infrastructure provider - provider = create_provider(cfg) - # Convert config to args-like object for backward compatibility + # (also constructs the infrastructure provider, exposed as args.provider) args = config.Args(cfg) - args.remote_write_ip = provider.get_node_ip(args.node_offset) + provider = args.provider if provider.is_remote(): local_experiment_root_dir = os.path.join( diff --git a/asap-tools/experiments/experiment_run_exporters_and_prometheus.py b/asap-tools/experiments/experiment_run_exporters_and_prometheus.py index ba11bab..3be0b61 100644 --- a/asap-tools/experiments/experiment_run_exporters_and_prometheus.py +++ b/asap-tools/experiments/experiment_run_exporters_and_prometheus.py @@ -7,7 +7,6 @@ import constants import experiment_utils from experiment_utils import sync, config -from experiment_utils.providers.factory import create_provider from experiment_utils.services import ( ExporterServiceFactory, SystemExportersService, @@ -30,10 +29,9 @@ def main(cfg: DictConfig): # Validate experiment configuration config.validate_experiment_config(cfg.experiment_params) # Convert config to args-like object for backward compatibility + # (also constructs the infrastructure provider, exposed as args.provider) args = config.Args(cfg) - - # Create infrastructure provider - provider = create_provider(cfg) + provider = args.provider local_experiment_root_dir = os.path.join( constants.LOCAL_EXPERIMENT_DIR, args.experiment_name diff --git a/asap-tools/experiments/experiment_run_grafana_demo.py b/asap-tools/experiments/experiment_run_grafana_demo.py index b7383ba..a585411 100644 --- a/asap-tools/experiments/experiment_run_grafana_demo.py +++ b/asap-tools/experiments/experiment_run_grafana_demo.py @@ -8,7 +8,6 @@ import constants import experiment_utils from experiment_utils import sync, config -from experiment_utils.providers.factory import create_provider from experiment_utils.services import ( KafkaService, QueryEngineRustService, @@ -51,11 +50,9 @@ def main(cfg: DictConfig): # Validate experiment configuration config.validate_experiment_config(cfg.experiment_params) # Convert config to args-like object for backward compatibility + # (also constructs the infrastructure provider, exposed as args.provider) args = config.Args(cfg) - - # Create infrastructure provider - provider = create_provider(cfg) - args.remote_write_ip = provider.get_node_ip(args.node_offset) + provider = args.provider args.forward_unsupported_queries = True print("Forcing forward_unsupported_queries to True for Grafana demo") diff --git a/asap-tools/experiments/experiment_teardown_everything.py b/asap-tools/experiments/experiment_teardown_everything.py index 534fbd3..c14276f 100644 --- a/asap-tools/experiments/experiment_teardown_everything.py +++ b/asap-tools/experiments/experiment_teardown_everything.py @@ -14,7 +14,6 @@ import constants from experiment_utils import config -from experiment_utils.providers.factory import create_provider from experiment_utils.services import ( KafkaService, FlinkService, @@ -55,9 +54,7 @@ def main(cfg: DictConfig): # Validate configuration (minimal validation for provider setup) config.validate_config(cfg) args = config.Args(cfg) - - # Create infrastructure provider - provider = create_provider(cfg) + provider = args.provider num_nodes_in_experiment = args.num_nodes diff --git a/asap-tools/experiments/experiment_utils/config.py b/asap-tools/experiments/experiment_utils/config.py index 64bcba4..966dbf9 100644 --- a/asap-tools/experiments/experiment_utils/config.py +++ b/asap-tools/experiments/experiment_utils/config.py @@ -11,6 +11,7 @@ from omegaconf import DictConfig, ListConfig, OmegaConf import constants +from experiment_utils.providers.factory import create_provider def validate_basic_config( @@ -532,6 +533,19 @@ def __init__(self, cfg: DictConfig): self.cloudlab_username = cfg.providers.cloudlab.username self.hostname_suffix = cfg.providers.cloudlab.hostname_suffix + # Single source of truth for the infrastructure provider: build it here + # (from cfg alone) so every node-dependent value derived from it — + # remote_write_ip included — is computed in exactly one place instead + # of being re-derived (and potentially forgotten) in each script. + self.provider = create_provider(cfg) + + # Remote write IP (127.0.0.1 for local mode, CloudLab's 10.10.1.x + # scheme otherwise). Written back onto cfg.streaming.remote_write.ip + # too, since generate_prometheus_config reads it from cfg directly + # rather than through this Args object. + self.remote_write_ip = self.provider.get_node_ip(self.node_offset) + cfg.streaming.remote_write.ip = self.remote_write_ip + # Logging and debugging self.log_level = cfg.logging.level @@ -565,10 +579,8 @@ def __init__(self, cfg: DictConfig): self.do_local_flink = cfg.streaming.do_local_flink self.forward_unsupported_queries = cfg.streaming.forward_unsupported_queries self.use_kafka_ingest = cfg.streaming.use_kafka_ingest - # Remote write configuration - self.remote_write_ip = ( - None # set by caller: provider.get_node_ip(self.node_offset) - ) + # Remote write configuration (self.remote_write_ip set above, near + # provider construction) self.remote_write_base_port = cfg.streaming.remote_write.base_port self.remote_write_path = cfg.streaming.remote_write.path From c52307a5798bf90d6ee2d81a0a9f2b6798cfed75 Mon Sep 17 00:00:00 2001 From: Milind Srivastava Date: Tue, 18 Aug 2026 15:01:10 -0400 Subject: [PATCH 2/3] fix(tools): fix cmdline_args.txt dump after Args gained a provider field MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit vars(args) broke once Args started holding a live provider object (added in the previous commit) — CloudLabProvider/LocalProvider aren't JSON-serializable, so every script crashed writing cmdline_args.txt. Args.to_dict() now swaps provider for its repr() (both provider classes already define one) instead of dropping it, keeping the debug dump both working and informative. Co-Authored-By: Claude Sonnet 5 --- asap-tools/experiments/experiment_only_ingest_path.py | 2 +- asap-tools/experiments/experiment_run_e2e.py | 2 +- .../experiment_run_exporters_and_prometheus.py | 2 +- asap-tools/experiments/experiment_run_grafana_demo.py | 2 +- asap-tools/experiments/experiment_utils/config.py | 9 +++++++++ 5 files changed, 13 insertions(+), 4 deletions(-) diff --git a/asap-tools/experiments/experiment_only_ingest_path.py b/asap-tools/experiments/experiment_only_ingest_path.py index e03bc07..a7c16dc 100644 --- a/asap-tools/experiments/experiment_only_ingest_path.py +++ b/asap-tools/experiments/experiment_only_ingest_path.py @@ -88,7 +88,7 @@ def main(cfg: DictConfig): # Also dump args to a file for backward compatibility with open(os.path.join(local_experiment_root_dir, "cmdline_args.txt"), "w") as f: - json.dump(vars(args), f) + json.dump(args.to_dict(), f) experiment_root_output_dir = ( f"{constants.CLOUDLAB_HOME_DIR}/experiment_outputs/{args.experiment_name}" diff --git a/asap-tools/experiments/experiment_run_e2e.py b/asap-tools/experiments/experiment_run_e2e.py index 41cdd45..ce9ccba 100644 --- a/asap-tools/experiments/experiment_run_e2e.py +++ b/asap-tools/experiments/experiment_run_e2e.py @@ -82,7 +82,7 @@ def main(cfg: DictConfig): # Also dump args to a file for backward compatibility with open(os.path.join(local_experiment_root_dir, "cmdline_args.txt"), "w") as f: - json.dump(vars(args), f) + json.dump(args.to_dict(), f) global CONTROLLER_REMOTE_OUTPUT_DIR, CONTROLLER_LOCAL_OUTPUT_DIR CONTROLLER_LOCAL_OUTPUT_DIR = os.path.join( diff --git a/asap-tools/experiments/experiment_run_exporters_and_prometheus.py b/asap-tools/experiments/experiment_run_exporters_and_prometheus.py index 3be0b61..414d7da 100644 --- a/asap-tools/experiments/experiment_run_exporters_and_prometheus.py +++ b/asap-tools/experiments/experiment_run_exporters_and_prometheus.py @@ -44,7 +44,7 @@ def main(cfg: DictConfig): # Also dump args to a file for backward compatibility with open(os.path.join(local_experiment_root_dir, "cmdline_args.txt"), "w") as f: - json.dump(vars(args), f) + json.dump(args.to_dict(), f) experiment_root_output_dir = ( f"{constants.CLOUDLAB_HOME_DIR}/experiment_outputs/{args.experiment_name}" diff --git a/asap-tools/experiments/experiment_run_grafana_demo.py b/asap-tools/experiments/experiment_run_grafana_demo.py index a585411..d7e98b4 100644 --- a/asap-tools/experiments/experiment_run_grafana_demo.py +++ b/asap-tools/experiments/experiment_run_grafana_demo.py @@ -68,7 +68,7 @@ def main(cfg: DictConfig): # Also dump args to a file for backward compatibility with open(os.path.join(local_experiment_root_dir, "cmdline_args.txt"), "w") as f: - json.dump(vars(args), f) + json.dump(args.to_dict(), f) experiment_root_output_dir = ( f"{constants.CLOUDLAB_HOME_DIR}/experiment_outputs/{args.experiment_name}" diff --git a/asap-tools/experiments/experiment_utils/config.py b/asap-tools/experiments/experiment_utils/config.py index 966dbf9..891f159 100644 --- a/asap-tools/experiments/experiment_utils/config.py +++ b/asap-tools/experiments/experiment_utils/config.py @@ -636,6 +636,15 @@ def get_coordinator_node(self) -> int: """Get the coordinator node index (first node in the range).""" return self.node_offset + def to_dict(self) -> Dict[str, Any]: + """JSON-serializable snapshot of this Args instance. + + `provider` holds a live InfrastructureProvider object, so it's + swapped for its repr() (e.g. "CloudLabProvider(username='...', ...)") + rather than dropped. + """ + return {k: (repr(v) if k == "provider" else v) for k, v in vars(self).items()} + def validate_config(cfg: DictConfig, script_name: str = "experiment_run_e2e"): """ From d36e48e142f8e360ac2cad972a181a3a314d9b9e Mon Sep 17 00:00:00 2001 From: Milind Srivastava Date: Wed, 19 Aug 2026 09:08:33 -0400 Subject: [PATCH 3/3] Updated scrape_interval to 1s, from 10s --- asap-tools/experiments/config/config.yaml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/asap-tools/experiments/config/config.yaml b/asap-tools/experiments/config/config.yaml index 7b69745..0c6d5ef 100644 --- a/asap-tools/experiments/config/config.yaml +++ b/asap-tools/experiments/config/config.yaml @@ -70,7 +70,7 @@ streaming: # Prometheus configuration prometheus: # Global Prometheus settings - scrape_interval: "10s" # How frequently to scrape targets + scrape_interval: "1s" # How frequently to scrape targets evaluation_interval: "10s" # How frequently to evaluate rules # query_log_file: "/scratch/sketch_db_for_prometheus/prometheus/queries.log" # Disabled to avoid permission issues in Docker # Recording rules settings