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
2 changes: 2 additions & 0 deletions deploy/mvp-multinode/.gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -4,3 +4,5 @@ results/
logs/
*.pid
*.bak
__pycache__/
*.py[cod]
8 changes: 7 additions & 1 deletion deploy/mvp-multinode/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`.

Expand All @@ -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.
8 changes: 8 additions & 0 deletions deploy/mvp-multinode/harness/acceptance.json
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
86 changes: 86 additions & 0 deletions deploy/mvp-multinode/scripts/measure_collector_telemetry.py
Original file line number Diff line number Diff line change
@@ -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())
63 changes: 63 additions & 0 deletions deploy/mvp-multinode/scripts/mvp_evaluate.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"))
Expand All @@ -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:
Expand Down Expand Up @@ -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 = []
Expand Down Expand Up @@ -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",
Expand Down Expand Up @@ -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,
}

Expand Down Expand Up @@ -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).", "",
Expand Down
8 changes: 8 additions & 0 deletions deploy/mvp-multinode/scripts/run_demo.sh
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand Down
Original file line number Diff line number Diff line change
@@ -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()
Loading
Loading