Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ services:
"--dataset", "{{ dataset }}",
"--num-labels", "{{ num_labels }}",
"--num-values-per-label", "{{ num_values_per_label | string }}",
"--metric-type", "{{ metric_type }}"
"--metric-type", "{{ metric_type }}"{% if seed is defined %},
"--seed", "{{ seed }}"{% endif %}
]
restart: unless-stopped
Original file line number Diff line number Diff line change
Expand Up @@ -73,7 +73,7 @@ const CONST_1M: u64 = 1_000_000;
const CONST_2M: u64 = 2_000_000;
const CONST_3M: u64 = 3_000_000;

const RNG_SEED: u64 = 0; // seed for rng used by all distributions
const DEFAULT_RNG_SEED: u64 = 0; // default seed for rng used by all distributions

const ZIPF_ALPHA: f64 = 1.01; // zipf parameter

Expand Down Expand Up @@ -217,6 +217,7 @@ impl FakeCollector {
label_names: Option<String>,
label_value_prefixes: Option<String>,
add_pattern_label: bool,
seed: u64,
) -> Self {
let num_values_per_label = get_num_vals_per_label(num_values_per_label, num_labels);
let prefixes: Option<Vec<String>> = label_value_prefixes
Expand Down Expand Up @@ -308,7 +309,7 @@ impl FakeCollector {
metric_name,
label_names,
add_pattern_label,
rng: Mutex::new(SmallRng::seed_from_u64(RNG_SEED)),
rng: Mutex::new(SmallRng::seed_from_u64(seed)),
zipf_dist,
normal_dist,
uniform_dist,
Expand Down Expand Up @@ -626,6 +627,9 @@ struct Args {

#[arg(long, default_value = "false", help = "Add 'pattern' label to metrics with dataset name")]
add_pattern_label: bool,

#[arg(long, default_value_t = DEFAULT_RNG_SEED, help = "Seed for random value generation")]
seed: u64,
}

#[tokio::main]
Expand All @@ -642,6 +646,7 @@ async fn main() -> Result<(), BoxedErr> {
args.label_names,
args.label_value_prefixes,
args.add_pattern_label,
args.seed,
));

// Register collector and start serving
Expand Down
6 changes: 6 additions & 0 deletions asap-tools/experiments/CONFIG_PARAMETERS_REFERENCE.md
Original file line number Diff line number Diff line change
Expand Up @@ -482,6 +482,11 @@ These parameters come from the `experiment_type` config group and are prefixed w
- **Example**: `5`
- **SCHEMA ISSUE**: Values vary widely (1-10) across configs without clear pattern

#### `experiment_params.exporters.exporter_list.fake_exporter.seed` (int, optional)
- **Description**: Base seed for Rust fake exporter data generation. Each exporter receives a distinct seed derived from this base seed, worker ordinal, and port ordinal.
- **Default**: `0`
- **Example**: `42`

#### `experiment_params.exporters.exporter_list.fake_exporter.start_port` (int, optional)
- **Description**: Starting port number for fake exporters
- **Default**: `50000`
Expand Down Expand Up @@ -727,6 +732,7 @@ python experiment_run_e2e.py \
python experiment_run_e2e.py \
experiment_type=cloud_demo \
experiment_params.exporters.exporter_list.fake_exporter.num_ports_per_server=5 \
experiment_params.exporters.exporter_list.fake_exporter.seed=42 \
experiment_params.exporters.exporter_list.fake_exporter.metric_type=gauge \
[required params...]

Expand Down
207 changes: 149 additions & 58 deletions asap-tools/experiments/experiment_utils/services/fake_exporters.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,11 +4,77 @@

import os
from abc import abstractmethod
from concurrent.futures import ThreadPoolExecutor
from typing import Tuple, List, Dict, Any

from .base import BaseService
from experiment_utils.providers.base import InfrastructureProvider

DEFAULT_FAKE_EXPORTER_SEED = 0
MAX_FAKE_EXPORTER_SEED = (1 << 64) - 1


def get_fake_exporter_seed(
config: Dict[str, Any], worker_ordinal: int, port_ordinal: int
) -> int:
"""Derive a unique, reproducible seed for one exporter instance."""
base_seed = config.get("seed", DEFAULT_FAKE_EXPORTER_SEED)
if (
not isinstance(base_seed, int)
or isinstance(base_seed, bool)
or not 0 <= base_seed <= MAX_FAKE_EXPORTER_SEED
):
raise ValueError(
"fake_exporter.seed must be an integer in the range "
f"0..{MAX_FAKE_EXPORTER_SEED}, got {base_seed!r}"
)

num_ports_per_server = config["num_ports_per_server"]
exporter_ordinal = worker_ordinal * num_ports_per_server + port_ordinal
derived_seed = base_seed + exporter_ordinal
if derived_seed > MAX_FAKE_EXPORTER_SEED:
raise ValueError(
"derived fake exporter seed exceeds the Rust u64 range: "
f"base seed {base_seed} with exporter ordinal {exporter_ordinal} "
f"produces {derived_seed}, maximum {MAX_FAKE_EXPORTER_SEED}"
)

return derived_seed


def get_fake_exporter_target_nodes(args) -> List[int]:
"""Return exporter nodes, preserving the single-node local fallback."""
worker_nodes = args.get_node_range(include_coordinator=False)
return worker_nodes or [args.get_coordinator_node()]


def execute_fake_exporter_commands_in_parallel(
provider: InfrastructureProvider,
node_commands: List[Tuple[int, str]],
**kwargs,
) -> None:
"""Run node-specific exporter commands concurrently.

The provider method applies one command to every node in ``node_idxs``.
Fake exporters need different commands per node because their seeds are
node-specific, so invoke it once per node while preserving concurrency.
"""
if not node_commands:
return

with ThreadPoolExecutor(max_workers=len(node_commands)) as executor:
futures = [
executor.submit(
provider.execute_command_parallel,
node_idxs=[node_idx],
cmd=cmd,
**kwargs,
)
for node_idx, cmd in node_commands
]
for future in futures:
future.result()


class BaseExporterService(BaseService):
"""Base class for exporter services."""
Expand Down Expand Up @@ -362,17 +428,27 @@ def _start_bare_metal(
num_ports = config["num_ports_per_server"]
dataset = config["dataset"]

cmds = []
for port in range(num_ports):
cmd = "./target/release/fake_exporter --port {} --valuescale {} --dataset {} --num-labels {} --num-values-per-label {} --metric-type {}".format(
port + config["start_port"],
config["synthetic_data_value_scale"],
dataset,
config["num_labels"],
config["num_values_per_label"],
config["metric_type"],
)
cmds.append(cmd)
commands_by_port: List[List[Tuple[int, str]]] = []
target_nodes = get_fake_exporter_target_nodes(self.args)
for port_ordinal in range(num_ports):
port_commands = []
for worker_ordinal, node_idx in enumerate(target_nodes):
seed = get_fake_exporter_seed(config, worker_ordinal, port_ordinal)
cmd = "./target/release/fake_exporter --port {} --valuescale {} --dataset {} --num-labels {} --num-values-per-label {} --metric-type {} --seed {}".format(
port_ordinal + config["start_port"],
config["synthetic_data_value_scale"],
dataset,
config["num_labels"],
config["num_values_per_label"],
config["metric_type"],
seed,
)
port_commands.append((node_idx, cmd))
commands_by_port.append(port_commands)

cmds = [
command for port_commands in commands_by_port for command in port_commands
]

cmd_dir = f"{self.provider.get_home_dir()}/code/asap-tools/data-sources/prometheus-exporters/fake_exporter/fake_exporter_rust/fake_exporter"

Expand All @@ -383,13 +459,13 @@ def _start_bare_metal(
with open(
os.path.join(local_experiment_dir, "fake_exporter_config", "cmds.sh"), "w"
) as f:
f.write("\n".join(cmds))
f.write("\n".join(cmd for _, cmd in cmds))

# Run commands in parallel across nodes
for cmd in cmds:
self.provider.execute_command_parallel(
node_idxs=self.args.get_node_range(include_coordinator=False),
cmd=cmd,
# Preserve cross-node fan-out while assigning each node its own seed.
for port_commands in commands_by_port:
execute_fake_exporter_commands_in_parallel(
self.provider,
port_commands,
cmd_dir=cmd_dir,
nohup=False,
popen=True,
Expand All @@ -410,30 +486,34 @@ def _start_containerized(
dataset = config["dataset"]

# Build docker run commands for each port
docker_run_cmds: List[str] = []
docker_run_cmds: List[Tuple[int, str]] = []
container_names: List[str] = []
target_nodes = get_fake_exporter_target_nodes(self.args)

for worker_ordinal, node_idx in enumerate(target_nodes):
for port_ordinal in range(num_ports):
actual_port = port_ordinal + config["start_port"]
seed = get_fake_exporter_seed(config, worker_ordinal, port_ordinal)
container_name = f"{BaseExporterService.FAKE_EXPORTER_BASE_CONTAINER_NAME}-{node_idx}-{actual_port}-rust"

# Build docker run command
docker_cmd = (
f"docker run -d "
f"--name {container_name} "
f"-p {actual_port}:{actual_port} "
f"--restart unless-stopped "
f"sketchdb-fake-exporter-rust:latest "
f"--port {actual_port} "
f"--valuescale {config['synthetic_data_value_scale']} "
f"--dataset {dataset} "
f"--num-labels {config['num_labels']} "
f"--num-values-per-label {config['num_values_per_label']} "
f"--metric-type {config['metric_type']} "
f"--seed {seed}"
)

for port in range(num_ports):
actual_port = port + config["start_port"]
container_name = f"{BaseExporterService.FAKE_EXPORTER_BASE_CONTAINER_NAME}-{actual_port}-rust"

# Build docker run command
docker_cmd = (
f"docker run -d "
f"--name {container_name} "
f"-p {actual_port}:{actual_port} "
f"--restart unless-stopped "
f"sketchdb-fake-exporter-rust:latest "
f"--port {actual_port} "
f"--valuescale {config['synthetic_data_value_scale']} "
f"--dataset {dataset} "
f"--num-labels {config['num_labels']} "
f"--num-values-per-label {config['num_values_per_label']} "
f"--metric-type {config['metric_type']}"
)

container_names.append(container_name)
docker_run_cmds.append(docker_cmd)
container_names.append(container_name)
docker_run_cmds.append((node_idx, docker_cmd))

self.container_names = container_names

Expand All @@ -447,24 +527,33 @@ def _start_containerized(
),
"w",
) as f:
f.write("\n".join(docker_run_cmds))
f.write("\n".join(cmd for _, cmd in docker_run_cmds))

# Start containers in batches to avoid overwhelming Docker daemon
# Start containers in batches to avoid overwhelming Docker daemon.
BATCH_SIZE = 5
for i in range(0, len(docker_run_cmds), BATCH_SIZE):
batch = docker_run_cmds[i : i + BATCH_SIZE]
# Combine docker run commands in batch into single SSH command
batch_cmd = "; ".join(batch)

self.provider.execute_command_parallel(
node_idxs=self.args.get_node_range(include_coordinator=False),
cmd=batch_cmd,
cmd_dir="",
nohup=False,
popen=True,
redirect=True,
wait=True, # Wait for batch to complete
)
node_batch_commands: List[Tuple[int, str]] = []
for node_idx in target_nodes:
node_commands = [
cmd
for command_node_idx, cmd in docker_run_cmds
if command_node_idx == node_idx
]
for i in range(0, len(node_commands), BATCH_SIZE):
batch = node_commands[i : i + BATCH_SIZE]
# Combine docker run commands in batch for one node.
batch_cmd = "; ".join(batch)
node_batch_commands.append((node_idx, batch_cmd))

# Start each worker's batches concurrently.
execute_fake_exporter_commands_in_parallel(
self.provider,
node_batch_commands,
cmd_dir="",
nohup=False,
popen=True,
redirect=True,
wait=True, # Wait for batch to complete
)

return

Expand All @@ -488,8 +577,9 @@ def _stop_bare_metal(self, **kwargs) -> None:
**kwargs: Additional configuration (currently unused)
"""
cmd = "pkill -f fake_exporter"
target_nodes = get_fake_exporter_target_nodes(self.args)
self.provider.execute_command_parallel(
node_idxs=self.args.get_node_range(include_coordinator=False),
node_idxs=target_nodes,
cmd=cmd,
cmd_dir=None,
nohup=False,
Expand All @@ -499,6 +589,7 @@ def _stop_bare_metal(self, **kwargs) -> None:

def _stop_containerized(self, **kwargs) -> None:
"""Stop fake exporters using containerized deployment."""
target_nodes = get_fake_exporter_target_nodes(self.args)
try:
if self.container_names is not None and len(self.container_names) > 0:
# Stop and remove containers by name
Expand All @@ -510,7 +601,7 @@ def _stop_containerized(self, **kwargs) -> None:
cmd = f"docker stop {container_list} 2>/dev/null || true; docker rm {container_list} 2>/dev/null || true"

self.provider.execute_command_parallel(
node_idxs=self.args.get_node_range(include_coordinator=False),
node_idxs=target_nodes,
cmd=cmd,
cmd_dir=None,
nohup=False,
Expand All @@ -521,7 +612,7 @@ def _stop_containerized(self, **kwargs) -> None:
# Fallback: stop all containers matching the base name pattern
cmd = f"docker ps -a --filter name={BaseExporterService.FAKE_EXPORTER_BASE_CONTAINER_NAME} --format '{{{{.Names}}}}' | xargs -r docker stop; docker ps -a --filter name={BaseExporterService.FAKE_EXPORTER_BASE_CONTAINER_NAME} --format '{{{{.Names}}}}' | xargs -r docker rm"
self.provider.execute_command_parallel(
node_idxs=self.args.get_node_range(include_coordinator=False),
node_idxs=target_nodes,
cmd=cmd,
cmd_dir=None,
nohup=False,
Expand Down
7 changes: 7 additions & 0 deletions asap-tools/experiments/generate_fake_exporter_compose.py
Original file line number Diff line number Diff line change
Expand Up @@ -88,6 +88,11 @@ def main():
required=True,
help="Type of metric to generate (e.g., gauge, counter)",
)
parser.add_argument(
"--seed",
type=int,
help="Base seed for random value generation",
)

parser.add_argument(
"--template-path",
Expand Down Expand Up @@ -132,6 +137,8 @@ def main():
# Only include container_name if provided, so Jinja2 default filter can work
if args.container_name:
template_vars["container_name"] = args.container_name
if args.seed is not None:
template_vars["seed"] = args.seed

# Set up file paths
script_dir = Path(__file__).parent
Expand Down
1 change: 1 addition & 0 deletions asap-tools/experiments/generate_workload.py
Original file line number Diff line number Diff line change
Expand Up @@ -366,6 +366,7 @@ def get_base_config() -> Dict:
"fake_exporter": {
"num_ports_per_server": 1,
"start_port": 50000,
"seed": 0,
"dataset": "zipf",
"synthetic_data_value_scale": 10000,
"num_labels": 3,
Expand Down
Loading