From 9690fe41e26d4b747a39956d12c66824f9b80425 Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 4 Sep 2026 06:41:14 -0600 Subject: [PATCH 1/2] test(mvp): gate collector runtime end to end --- deploy/mvp-multinode/README.md | 8 +- deploy/mvp-multinode/harness/acceptance.json | 8 ++ .../scripts/measure_collector_telemetry.py | 86 +++++++++++++++++++ deploy/mvp-multinode/scripts/mvp_evaluate.py | 63 ++++++++++++++ deploy/mvp-multinode/scripts/run_demo.sh | 8 ++ .../tests/test_measure_collector_telemetry.py | 34 ++++++++ .../scripts/tests/test_mvp_evaluate.py | 50 ++++++++++- 7 files changed, 255 insertions(+), 2 deletions(-) create mode 100644 deploy/mvp-multinode/scripts/measure_collector_telemetry.py create mode 100644 deploy/mvp-multinode/scripts/tests/test_measure_collector_telemetry.py diff --git a/deploy/mvp-multinode/README.md b/deploy/mvp-multinode/README.md index f76d60e3..e5d95021 100644 --- a/deploy/mvp-multinode/README.md +++ b/deploy/mvp-multinode/README.md @@ -52,7 +52,8 @@ mvp-multinode/ Each run preserves its manifest, copied acceptance/query inputs, replay JSONL, query responses, controller configuration and logs, freshness CSV, per-process -resource CSVs, host-NIC samples, and storage measurements. The evaluator emits +resource CSVs, host-NIC samples, storage measurements, and before/after +Collector self-telemetry counters. The evaluator emits `MVP_RESULTS.json` and `MVP_REPORT.md`; any missing, stale, empty, or incomparable required evidence makes the overall verdict `FAIL`. @@ -62,6 +63,11 @@ timestamps. Freshness is reported per query class and transmission mode. Resource evidence includes CPU time, steady-state and peak RSS, network traffic, storage, and every active ASAP/Thanos backend process. The checked-in cost weights convert only those measured quantities into the normalized comparison. +The Collector runtime gate additionally requires both agents to sustain at +least 90% of the declared base event rate, report no refused or failed-export +metric points, process at least 25k points per CPU-second, stay below four CPU +cores per agent, and keep per-agent peak RSS below 4 GiB. These thresholds are +declared in `harness/acceptance.json`; missing telemetry fails closed. The MVP excludes Serf comparisons, analytical/projected savings, fleet/figure sweeps, and archive fallback validation. diff --git a/deploy/mvp-multinode/harness/acceptance.json b/deploy/mvp-multinode/harness/acceptance.json index d0adb3c0..93fa8386 100644 --- a/deploy/mvp-multinode/harness/acceptance.json +++ b/deploy/mvp-multinode/harness/acceptance.json @@ -35,6 +35,14 @@ "require_asap_lower_p50": true, "require_asap_lower_p95": true }, + "collector_runtime": { + "minimum_offered_load_fraction": 0.90, + "minimum_points_per_cpu_second": 25000, + "maximum_peak_rss_mib_per_agent": 4096, + "maximum_cpu_cores_per_agent": 4.0, + "require_zero_refused_points": true, + "require_zero_export_failed_points": true + }, "cost": { "cpu_core_weight": 1.0, "rss_gib_weight": 0.1, diff --git a/deploy/mvp-multinode/scripts/measure_collector_telemetry.py b/deploy/mvp-multinode/scripts/measure_collector_telemetry.py new file mode 100644 index 00000000..a13f9a7f --- /dev/null +++ b/deploy/mvp-multinode/scripts/measure_collector_telemetry.py @@ -0,0 +1,86 @@ +#!/usr/bin/env python3 +"""Measure Collector ingest/export counters over an E2E load window.""" + +from __future__ import annotations + +import argparse +import json +import math +import re +import time +import urllib.request + + +SAMPLE_RE = re.compile(r"^([a-zA-Z_:][a-zA-Z0-9_:]*)(?:\{[^}]*\})?\s+([^\s]+)") + + +def parse_prometheus(text: str) -> dict[str, float]: + totals: dict[str, float] = {} + for line in text.splitlines(): + match = SAMPLE_RE.match(line) + if not match or line.startswith("#"): + continue + try: + value = float(match.group(2)) + except ValueError: + continue + if math.isfinite(value): + totals[match.group(1)] = totals.get(match.group(1), 0.0) + value + return totals + + +def scrape(url: str) -> dict[str, float]: + with urllib.request.urlopen(url, timeout=10) as response: + return parse_prometheus(response.read().decode("utf-8")) + + +def counter_total(values: dict[str, float], fragment: str) -> tuple[float, list[str]]: + names = sorted(name for name in values if name.endswith(fragment) or name.endswith(fragment + "_total")) + return sum(values[name] for name in names), names + + +def main() -> int: + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--agent", action="append", required=True, metavar="NAME=URL") + parser.add_argument("--duration", type=float, required=True) + parser.add_argument("--out", required=True) + args = parser.parse_args() + + endpoints = dict(item.split("=", 1) for item in args.agent) + started_at = time.time() + start = {name: scrape(url) for name, url in endpoints.items()} + time.sleep(max(args.duration, 0.0)) + end = {name: scrape(url) for name, url in endpoints.items()} + duration = max(time.time() - started_at, 1e-9) + + agents = {} + fragments = { + "accepted_metric_points": "receiver_accepted_metric_points", + "refused_metric_points": "receiver_refused_metric_points", + "sent_metric_points": "exporter_sent_metric_points", + "send_failed_metric_points": "exporter_send_failed_metric_points", + } + for agent in sorted(endpoints): + counters = {} + metric_names = {} + for key, fragment in fragments.items(): + before, before_names = counter_total(start[agent], fragment) + after, after_names = counter_total(end[agent], fragment) + counters[key] = after - before + metric_names[key] = sorted(set(before_names) | set(after_names)) + agents[agent] = {"endpoint": endpoints[agent], "counter_deltas": counters, + "metric_names": metric_names} + + result = { + "schema_version": 1, + "duration_s": duration, + "agents": agents, + } + with open(args.out, "w", encoding="utf-8") as handle: + json.dump(result, handle, indent=2, sort_keys=True) + handle.write("\n") + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/deploy/mvp-multinode/scripts/mvp_evaluate.py b/deploy/mvp-multinode/scripts/mvp_evaluate.py index 5f33076e..51a7133e 100644 --- a/deploy/mvp-multinode/scripts/mvp_evaluate.py +++ b/deploy/mvp-multinode/scripts/mvp_evaluate.py @@ -189,6 +189,8 @@ def freshness_summary(path: str, cfg: dict[str, Any]) -> dict[str, Any]: def resource_summary(arm_dir: str, weights: dict[str, Any]) -> dict[str, Any]: cpu = cpu_time_s = rss_mib = peak_rss_mib = network_bps = disk_mib = 0.0 collector_cpu = collector_rss = collector_network_bps = 0.0 + collector_peak_rss = collector_max_cpu = 0.0 + collector_containers: set[str] = set() stage_files = glob.glob(os.path.join(arm_dir, "stages-*.csv")) nic_files = glob.glob(os.path.join(arm_dir, "nic-*.csv")) storage_files = glob.glob(os.path.join(arm_dir, "storage-*.csv")) @@ -215,8 +217,11 @@ def resource_summary(arm_dir: str, weights: dict[str, Any]) -> dict[str, Any]: disk_mib += row_disk text = " ".join((row.get("stage") or "", row.get("container") or "")).lower() if "agent" in text or "collector" in text or "otel" in text: + collector_containers.add(row.get("container") or row.get("stage") or "unknown") collector_cpu += row_cpu collector_rss += row_rss + collector_peak_rss = max(collector_peak_rss, row_peak_rss) + collector_max_cpu = max(collector_max_cpu, row_cpu) if math.isfinite(row_net_out): collector_network_bps += row_net_out * 1024 for path in nic_files: @@ -249,11 +254,59 @@ def cost(c: float, r: float, n: float, d: float = 0.0) -> float: "peak_rss_mib": peak_rss_mib, "network_bytes_per_s": network_bps, "disk_mib": disk_mib, "normalized_cost": cost(cpu, rss_mib, network_bps, disk_mib), "collector_cpu_cores": collector_cpu, "collector_rss_mib": collector_rss, + "collector_peak_rss_mib_per_agent": collector_peak_rss, + "collector_max_cpu_cores_per_agent": collector_max_cpu, + "collector_containers": sorted(collector_containers), "collector_network_bytes_per_s": collector_network_bps, "collector_normalized_cost": cost(collector_cpu, collector_rss, collector_network_bps), } +def collector_runtime_summary(path: str, resources: dict[str, Any], workload: dict[str, Any], + cfg: dict[str, Any]) -> dict[str, Any]: + evidence = load_json(path) + agents = evidence.get("agents") or {} + duration = float(evidence.get("duration_s") or 0) + required_agents = int(workload.get("agents") or 0) + accepted = refused = send_failed = 0.0 + accepted_evidence = True + counters_nonnegative = True + for agent in agents.values(): + counters = agent.get("counter_deltas") or {} + names = agent.get("metric_names") or {} + accepted_evidence = accepted_evidence and bool(names.get("accepted_metric_points")) + values = [float(counters.get(key, math.nan)) for key in + ("accepted_metric_points", "refused_metric_points", "send_failed_metric_points")] + counters_nonnegative = counters_nonnegative and all(math.isfinite(value) and value >= 0 for value in values) + accepted += values[0] + refused += values[1] + send_failed += values[2] + accepted_per_s = accepted / duration if duration > 0 else 0.0 + offered_per_s = float(workload.get("series_cardinality") or 0) * float(workload.get("frequency_hz") or 0) + offered_fraction = accepted_per_s / offered_per_s if offered_per_s > 0 else 0.0 + cpu_cores = float(resources.get("collector_cpu_cores") or 0) + points_per_cpu_second = accepted_per_s / cpu_cores if cpu_cores > 0 else 0.0 + peak_rss = float(resources.get("collector_peak_rss_mib_per_agent") or 0) + max_cpu = float(resources.get("collector_max_cpu_cores_per_agent") or 0) + checks = { + "complete_evidence": len(agents) == required_agents and accepted_evidence and counters_nonnegative, + "sustained_offered_load": offered_fraction >= float(cfg["minimum_offered_load_fraction"]), + "cpu_efficiency": points_per_cpu_second >= float(cfg["minimum_points_per_cpu_second"]), + "peak_rss_bound": 0 < peak_rss <= float(cfg["maximum_peak_rss_mib_per_agent"]), + "cpu_bound": 0 < max_cpu <= float(cfg["maximum_cpu_cores_per_agent"]), + "zero_refused_points": not cfg.get("require_zero_refused_points", True) or refused == 0, + "zero_export_failed_points": not cfg.get("require_zero_export_failed_points", True) or send_failed == 0, + } + return { + "passed": all(checks.values()), "checks": checks, "duration_s": duration, + "agents": sorted(agents), "accepted_metric_points": accepted, + "refused_metric_points": refused, "send_failed_metric_points": send_failed, + "accepted_points_per_second": accepted_per_s, "declared_base_points_per_second": offered_per_s, + "offered_load_fraction": offered_fraction, "points_per_cpu_second": points_per_cpu_second, + "peak_rss_mib_per_agent": peak_rss, "max_cpu_cores_per_agent": max_cpu, + } + + def check_manifest(manifest: dict[str, Any], run_dir: str, baseline: str, asap: str, artifacts: list[str]) -> tuple[bool, list[str]]: errors = [] @@ -328,6 +381,7 @@ def evaluate(run_dir: str, config: dict[str, Any]) -> dict[str, Any]: baseline, asap = config["baseline_arm"], config["asap_arm"] failures: list[str] = [] required = ["run-manifest.json", f"{baseline}/replay.jsonl", f"{asap}/replay.jsonl", + f"{baseline}/collector-telemetry.json", f"{asap}/collector-telemetry.json", f"{baseline}/evaluation-start-ms.txt", f"{asap}/evaluation-start-ms.txt", f"{asap}/freshness-{asap}.csv", f"{asap}/controller-agents.json", f"{asap}/controller-config.yaml", f"{asap}/controller-config-agent-a.yaml", f"{asap}/controller-config-agent-b.yaml", @@ -418,12 +472,20 @@ def read_anchor(arm: str) -> int: e2e_cost_passed = (resource_complete and resources[asap]["normalized_cost"] < resources[baseline]["normalized_cost"]) if config["cost"]["require_end_to_end_total_lower"] else resource_complete if not collector_cost_passed: failures.append("collector cost gate failed") if not e2e_cost_passed: failures.append("end-to-end cost gate failed") + collector_runtime = { + arm: collector_runtime_summary(os.path.join(run_dir, arm, "collector-telemetry.json"), + resources[arm], manifest["workload"], config["collector_runtime"]) + for arm in (baseline, asap) + } + if not all(item["passed"] for item in collector_runtime.values()): + failures.append("collector runtime correctness/performance/resource gate failed") return { "schema_version": 1, "run_id": manifest.get("run_id"), "overall_verdict": "PASS" if not failures else "FAIL", "failures": failures, "manifest": {"passed": manifest_ok}, "run_manifest": manifest, "functional_correctness": correctness, "accuracy": accuracy, "freshness": freshness, "query_latency": {"passed": latency_passed, "arms": latency, "speedups": speedups}, "collector_cost": {"passed": collector_cost_passed, "resource_ratios": collector_ratios, "guardrail": guardrail}, + "collector_runtime": {"passed": all(item["passed"] for item in collector_runtime.values()), "arms": collector_runtime}, "end_to_end_cost": {"passed": e2e_cost_passed}, "resources": resources, } @@ -455,6 +517,7 @@ def latency_cell(arm: str) -> str: f"| Query accuracy | worst error | {worst_error:.6g} | ground truth | per-query SLA | {verdict('accuracy')} |", f"| Query freshness | worst p95 lag | {freshness_p95:.2f} ms | n/a | within SLA | {verdict('freshness')} |", f"| Query performance | per-query p50/p95 | {latency_cell('asap-gzip')} | {latency_cell('b1')} | ASAP lower | {verdict('query_latency')} |", + f"| Collector runtime | accepted points/s | {result.get('collector_runtime', {}).get('arms', {}).get('asap-gzip', {}).get('accepted_points_per_second', math.nan):.2f} | {result.get('collector_runtime', {}).get('arms', {}).get('b1', {}).get('accepted_points_per_second', math.nan):.2f} | no loss + throughput/CPU/RSS bounds | {verdict('collector_runtime')} |", f"| Collector cost | normalized cost | {asap.get('collector_normalized_cost', math.nan):.6g} | {baseline.get('collector_normalized_cost', math.nan):.6g} | ASAP lower | {verdict('collector_cost')} |", f"| End-to-end cost | normalized cost | {asap.get('normalized_cost', math.nan):.6g} | {baseline.get('normalized_cost', math.nan):.6g} | ASAP lower | {verdict('end_to_end_cost')} |", "", "Artifacts: [machine-readable results](MVP_RESULTS.json), [acceptance configuration](acceptance.json), [query workload](queries-e2e.json).", "", diff --git a/deploy/mvp-multinode/scripts/run_demo.sh b/deploy/mvp-multinode/scripts/run_demo.sh index c4db7b88..c5cdc22b 100755 --- a/deploy/mvp-multinode/scripts/run_demo.sh +++ b/deploy/mvp-multinode/scripts/run_demo.sh @@ -745,6 +745,13 @@ arm_measure() { fi log "[measure ${arm}] MetricsQL replay against ${query_endpoint} for ${SOAK_S}s" + python3 "${SCRIPT_DIR}/measure_collector_telemetry.py" \ + --duration "${SOAK_S}" \ + --agent "agent-a=http://${NODE0_IP}:8890/metrics" \ + --agent "agent-b=http://${NODE3_IP}:8890/metrics" \ + --out "${out}/collector-telemetry.json" \ + > "${out}/collector-telemetry.log" 2>&1 & + local COLLECTOR_TELEMETRY_PID=$! # The paired arms run sequentially, so their wall-clock timestamps cannot # be compared directly. Pin an explicit per-arm evaluation anchor and # preserve it as evidence. The reducer compares timestamps relative to @@ -765,6 +772,7 @@ arm_measure() { # Wait for soak to complete wait ${REPLAY_PID} 2>/dev/null || true + wait ${COLLECTOR_TELEMETRY_PID} 2>/dev/null || true if is_asap_arm "${arm}"; then # Controller-side acknowledgement is the proof that the supervisor diff --git a/deploy/mvp-multinode/scripts/tests/test_measure_collector_telemetry.py b/deploy/mvp-multinode/scripts/tests/test_measure_collector_telemetry.py new file mode 100644 index 00000000..48c1e1e9 --- /dev/null +++ b/deploy/mvp-multinode/scripts/tests/test_measure_collector_telemetry.py @@ -0,0 +1,34 @@ +from __future__ import annotations + +import importlib.util +import pathlib +import unittest + + +SCRIPT = pathlib.Path(__file__).parents[1] / "measure_collector_telemetry.py" +SPEC = importlib.util.spec_from_file_location("measure_collector_telemetry", SCRIPT) +MODULE = importlib.util.module_from_spec(SPEC) +assert SPEC.loader +SPEC.loader.exec_module(MODULE) + + +class CollectorTelemetryTest(unittest.TestCase): + def test_prometheus_parser_sums_label_series(self) -> None: + values = MODULE.parse_prometheus( + "# TYPE otelcol_receiver_accepted_metric_points counter\n" + 'otelcol_receiver_accepted_metric_points_total{receiver="otlp",transport="grpc"} 10\n' + 'otelcol_receiver_accepted_metric_points_total{receiver="otlp",transport="http"} 2\n' + "otelcol_receiver_accepted_metric_points_created 1700000000\n" + 'otelcol_receiver_refused_metric_points_total{receiver="otlp"} 0\n' + ) + total, names = MODULE.counter_total(values, "receiver_accepted_metric_points") + self.assertEqual(12, total) + self.assertEqual(["otelcol_receiver_accepted_metric_points_total"], names) + + def test_non_finite_samples_are_ignored(self) -> None: + values = MODULE.parse_prometheus("metric_a NaN\nmetric_b +Inf\nmetric_c 3\n") + self.assertEqual({"metric_c": 3}, values) + + +if __name__ == "__main__": + unittest.main() diff --git a/deploy/mvp-multinode/scripts/tests/test_mvp_evaluate.py b/deploy/mvp-multinode/scripts/tests/test_mvp_evaluate.py index 70a8f5e4..22326c67 100644 --- a/deploy/mvp-multinode/scripts/tests/test_mvp_evaluate.py +++ b/deploy/mvp-multinode/scripts/tests/test_mvp_evaluate.py @@ -28,6 +28,12 @@ def setUp(self) -> None: "freshness": {"p95_ms": 1000, "maximum_ms": 2000, "minimum_samples_per_tier": 2, "required_tiers": ["warm"]}, "query_latency": {"require_asap_lower_p50": True, "require_asap_lower_p95": True}, + "collector_runtime": {"minimum_offered_load_fraction": .9, + "minimum_points_per_cpu_second": 10, + "maximum_peak_rss_mib_per_agent": 1024, + "maximum_cpu_cores_per_agent": 4, + "require_zero_refused_points": True, + "require_zero_export_failed_points": True}, "cost": {"cpu_core_weight": 1.0, "rss_gib_weight": 0.1, "network_mib_per_s_weight": 0.01, "storage_gib_weight": 0.01, "baseline_storage_components": ["victoriametrics"], @@ -43,7 +49,8 @@ def write_fixture(self, *, plan=True, error=0.01) -> None: manifest = {"run_id": "run-1", "started_at": "2026-08-26T00:00:00Z", "collector_commit": "abc", "backend_commit": "def", "load_generator": {"commit": "abc"}, "exact_backend": {"name": "VictoriaMetrics"}, "images": {"collector": "sha256:1"}, "remote_image_digests": {"vm": "sha256:2"}, - "configuration_sha256": {"acceptance": "123"}, "workload": {"cardinality": 1}, + "configuration_sha256": {"acceptance": "123"}, + "workload": {"cardinality": 1, "frequency_hz": 10, "agents": 2, "series_cardinality": 2}, "time_alignment": {"mode": "result_timestamp"}, "arms": {"b1": {"seed": 42}, "asap-gzip": {"seed": 42}}} (self.run / "run-manifest.json").write_text(json.dumps(manifest)) @@ -58,6 +65,16 @@ def write_fixture(self, *, plan=True, error=0.01) -> None: "data_source": "sketch_store" if arm == "asap-gzip" else "victoriametrics", "result": [{"metric": {"zone": "a"}, "value": [anchor_ms / 1000, str(value)]}]} for i in range(3)] (arm_dir / "replay.jsonl").write_text("".join(json.dumps(row) + "\n" for row in records)) + telemetry = {"schema_version": 1, "duration_s": 10, "agents": {}} + for agent in ("agent-a", "agent-b"): + telemetry["agents"][agent] = { + "counter_deltas": {"accepted_metric_points": 100, "refused_metric_points": 0, + "sent_metric_points": 50, "send_failed_metric_points": 0}, + "metric_names": {"accepted_metric_points": ["otelcol_receiver_accepted_metric_points"], + "refused_metric_points": [], "sent_metric_points": [], + "send_failed_metric_points": []}, + } + (arm_dir / "collector-telemetry.json").write_text(json.dumps(telemetry)) with (arm_dir / "stages-node0.csv").open("w", newline="") as handle: writer = csv.writer(handle); writer.writerow(["baseline", "stage", "container", "cpu_cores", "cpu_time_s", "rss_mib", "peak_rss_mib"]) writer.writerow([arm, "agent", "asap-otel", 1 if arm == "b1" else .5, 60, 100, 110]) @@ -105,6 +122,37 @@ def test_missing_plan_evidence_fails(self) -> None: self.assertEqual("FAIL", result["overall_verdict"]) self.assertIn("functional correctness", " ".join(result["failures"])) + def test_collector_refused_points_fail(self) -> None: + self.write_fixture() + path = self.run / "asap-gzip" / "collector-telemetry.json" + value = json.loads(path.read_text()) + value["agents"]["agent-a"]["counter_deltas"]["refused_metric_points"] = 1 + path.write_text(json.dumps(value)) + result = MODULE.evaluate(str(self.run), self.config) + self.assertEqual("FAIL", result["overall_verdict"]) + self.assertFalse(result["collector_runtime"]["arms"]["asap-gzip"]["checks"]["zero_refused_points"]) + + def test_collector_low_throughput_fails(self) -> None: + self.write_fixture() + path = self.run / "asap-gzip" / "collector-telemetry.json" + value = json.loads(path.read_text()) + for agent in value["agents"].values(): + agent["counter_deltas"]["accepted_metric_points"] = 1 + path.write_text(json.dumps(value)) + result = MODULE.evaluate(str(self.run), self.config) + self.assertFalse(result["collector_runtime"]["arms"]["asap-gzip"]["checks"]["sustained_offered_load"]) + + def test_collector_peak_rss_bound_fails(self) -> None: + self.write_fixture() + path = self.run / "asap-gzip" / "stages-node0.csv" + with path.open() as handle: + rows = list(csv.reader(handle)) + rows[1][6] = "2048" + with path.open("w", newline="") as handle: + csv.writer(handle).writerows(rows) + result = MODULE.evaluate(str(self.run), self.config) + self.assertFalse(result["collector_runtime"]["arms"]["asap-gzip"]["checks"]["peak_rss_bound"]) + def test_archive_fallback_cannot_satisfy_warm_mvp_query(self) -> None: self.write_fixture() path = self.run / "asap-gzip" / "replay.jsonl" From 3ff90f705e00046f4d2b773ec9c63f8002c84ce5 Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 4 Sep 2026 07:26:57 -0600 Subject: [PATCH 2/2] chore: remove superseded manual test runners --- deploy/mvp-multinode/.gitignore | 2 + .../cmd/telemetrygen/run.sh | 2 - otel-app/abtest/.gitignore | 2 - otel-app/abtest/run_ab.sh | 76 ------------------- 4 files changed, 2 insertions(+), 80 deletions(-) delete mode 100644 opentelemetry-collector-contrib-patch/cmd/telemetrygen/run.sh delete mode 100644 otel-app/abtest/.gitignore delete mode 100755 otel-app/abtest/run_ab.sh diff --git a/deploy/mvp-multinode/.gitignore b/deploy/mvp-multinode/.gitignore index aa26dc01..6d66296b 100644 --- a/deploy/mvp-multinode/.gitignore +++ b/deploy/mvp-multinode/.gitignore @@ -4,3 +4,5 @@ results/ logs/ *.pid *.bak +__pycache__/ +*.py[cod] diff --git a/opentelemetry-collector-contrib-patch/cmd/telemetrygen/run.sh b/opentelemetry-collector-contrib-patch/cmd/telemetrygen/run.sh deleted file mode 100644 index c515f243..00000000 --- a/opentelemetry-collector-contrib-patch/cmd/telemetrygen/run.sh +++ /dev/null @@ -1,2 +0,0 @@ -go run . metrics --metric-type DDSketch --duration 100s --otlp-insecure - diff --git a/otel-app/abtest/.gitignore b/otel-app/abtest/.gitignore deleted file mode 100644 index a2c8e048..00000000 --- a/otel-app/abtest/.gitignore +++ /dev/null @@ -1,2 +0,0 @@ -results/ -run.log diff --git a/otel-app/abtest/run_ab.sh b/otel-app/abtest/run_ab.sh deleted file mode 100755 index 863c2f20..00000000 --- a/otel-app/abtest/run_ab.sh +++ /dev/null @@ -1,76 +0,0 @@ -#!/usr/bin/env bash -# A/B CPU comparison for otel-app across SDK aggregation modes. -# Identical workload; only -agg varies. Measures producer CPU-seconds -# (/usr/bin/time), max RSS, exported datapoints/bytes (from the discard -# sink), and a CPU pprof profile per mode. -set -u -cd "$(dirname "$0")/.." # otel-app/ -OUT="abtest/results" -mkdir -p "$OUT" -APP=./otel-app -SINK=./sink/sink - -FREQ=200 -CARD=500 -WINDOW=15s -DUR=60 # seconds -PROF_DELAY=12 # start profiling after warmup -PROF_SECS=40 # profile duration -PPROF_PORT=6060 -SINK_ADDR=:4317 - -: > "$OUT/summary.txt" -echo "params: freq_hz=$FREQ cardinality=$CARD sdk_window=$WINDOW duration=${DUR}s" | tee -a "$OUT/summary.txt" -echo "nproc=$(nproc) GOMAXPROCS=${GOMAXPROCS:-default}" | tee -a "$OUT/summary.txt" -echo "" | tee -a "$OUT/summary.txt" - -for AGG in default dd-full raw-buffer; do - echo "==================== AGG=$AGG ====================" | tee -a "$OUT/summary.txt" - - # fresh sink - $SINK "$SINK_ADDR" > "$OUT/sink-$AGG.log" 2>&1 & - SINK_PID=$! - # wait for sink to be listening - for i in $(seq 1 50); do - grep -q "sink listening" "$OUT/sink-$AGG.log" 2>/dev/null && break - kill -0 "$SINK_PID" 2>/dev/null || break - bash -c 'read -t 0.2 _ < <(:) || true' - done - - # producer under /usr/bin/time - /usr/bin/time -v -o "$OUT/time-$AGG.txt" \ - $APP -target=localhost:4317 -agg="$AGG" -freq-hz=$FREQ -cardinality=$CARD \ - -sdk-window=$WINDOW -duration=${DUR}s -pprof-addr=:$PPROF_PORT \ - > "$OUT/app-$AGG.log" 2>&1 & - APP_PID=$! - - # capture CPU profile mid-run - ( bash -c "read -t $PROF_DELAY _ < <(:) || true" - curl -s "http://localhost:$PPROF_PORT/debug/pprof/profile?seconds=$PROF_SECS" \ - -o "$OUT/cpu-$AGG.pprof" ) & - PROF_PID=$! - - wait "$APP_PID" - wait "$PROF_PID" 2>/dev/null - - # drain + stop sink to print stats - kill -TERM "$SINK_PID" 2>/dev/null - wait "$SINK_PID" 2>/dev/null - - # collect - USER_S=$(grep "User time" "$OUT/time-$AGG.txt" | grep -oE "[0-9.]+") - SYS_S=$(grep "System time" "$OUT/time-$AGG.txt" | grep -oE "[0-9.]+") - PCT=$(grep "Percent of CPU" "$OUT/time-$AGG.txt" | grep -oE "[0-9]+%") - RSS_KB=$(grep "Maximum resident" "$OUT/time-$AGG.txt" | grep -oE "[0-9]+") - STATS=$(grep "SINK_STATS" "$OUT/sink-$AGG.log" | tail -1) - CPU_TOTAL=$(awk "BEGIN{printf \"%.2f\", ${USER_S:-0}+${SYS_S:-0}}") - RSS_MB=$(awk "BEGIN{printf \"%.0f\", ${RSS_KB:-0}/1024}") - { - echo " CPU_total_s = $CPU_TOTAL (user=$USER_S sys=$SYS_S, %CPU=$PCT)" - echo " MaxRSS = ${RSS_MB} MB" - echo " sink = $STATS" - } | tee -a "$OUT/summary.txt" - echo "" | tee -a "$OUT/summary.txt" -done - -echo "DONE" | tee -a "$OUT/summary.txt"