Files
greptimedb/tests/perf/query_regression_runner.py
T
discord9 efc23d79bd feat(perf): add mixed read/write saturated case for scheduler validation
Add optional [scenario.write_measure.mix] to write_throughput: runs the
remote-write ingestion in a background thread while query-loop threads
(query_parallelism x query_interval_ms) hammer the same frontend, so
query and write tasks contend on the datanode runtime under a dual
backlog. Adds query gates (failure rate, p99 regression) and a
best-effort scheduler poll-share diagnostic scraped from the datanode
/metrics. Includes write_read_mixed_scheduler built-in case.

Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>
2026-08-06 14:24:52 +08:00

2322 lines
114 KiB
Python

#!/usr/bin/env python3
# Copyright 2023 Greptime Team
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
"""Base-vs-candidate query performance regression runner."""
from __future__ import annotations
import argparse
import fcntl
import hashlib
import json
import math
import os
import re
import shutil
import socket
import statistics
import subprocess
import sys
import threading
import time
import tomllib
import urllib.error
import urllib.parse
import urllib.request
from dataclasses import dataclass, field
from pathlib import Path
from typing import Any
FICLONE = 0x40049409
OTLP_TRACE_METRICS = {
"greptime_frontend_otlp_traces_rows",
"greptime_frontend_otlp_traces_failure_count",
"greptime_servers_http_otlp_traces_elapsed_sum",
"greptime_servers_http_otlp_traces_elapsed_count",
}
PROMETHEUS_SAMPLE_RE = re.compile(r"^([A-Za-z_:][A-Za-z0-9_:]*)(?:\{.*\})?\s+([^\s]+)(?:\s+.*)?$")
# Workload-scheduler cumulative poll counters on the datanode /metrics
# endpoint, e.g. greptime_workload_scheduler_polls{workload="write"} 1234.
SCHEDULER_POLL_RE = re.compile(r'^greptime_workload_scheduler_polls\{[^}]*workload="(query|write)"[^}]*\}\s+(\d+)(?:\s+.*)?$')
@dataclass(frozen=True)
class RunTarget:
name: str
binary: Path
work_dir: Path
data_dir: Path
fixture_dir: Path
report_path: Path
http_port: int
grpc_port: int
mysql_port: int
postgres_port: int
metasrv_rpc_port: int
metasrv_http_port: int
datanode_rpc_port: int
datanode_http_port: int
datanode_data_dir: Path
scheduler_env: dict[str, str] = field(default_factory=dict)
def parse_args() -> argparse.Namespace:
p = argparse.ArgumentParser(description="Run a query perf case against base and candidate binaries.")
p.add_argument("--case", required=True, type=Path)
p.add_argument("--base-bin", required=True, type=Path)
p.add_argument("--candidate-bin", required=True, type=Path)
p.add_argument("--fixture-generator", type=Path)
p.add_argument("--otelgen-bin", type=Path)
p.add_argument("--remote-write-generator", type=Path, help="deprecated: use --fixture-generator query_perf_fixture")
p.add_argument("--storage-inspector", type=Path, help="deprecated: use --fixture-generator query_perf_fixture")
p.add_argument("--work-dir", required=True, type=Path)
p.add_argument("--output", type=Path, help="write the final JSON report to this file instead of stdout")
p.add_argument("--fixture-cache-dir", type=Path, help="persistent directory for generated fixtures, keyed by case content")
p.add_argument("--reuse-fixture", action="store_true")
p.add_argument("--allow-large-fixture", action="store_true")
p.add_argument("--dry-run", action="store_true")
p.add_argument("--fixture-only", action="store_true", help="old smoke mode: generate/materialize fixture only")
p.add_argument("--reuse-work-dir", action="store_true", help="allow non-empty base/candidate work dirs")
p.add_argument("--http-timeout", type=float, default=120.0, help="HTTP SQL timeout seconds")
return p.parse_args()
def load_case(path: Path) -> dict[str, Any]:
with path.open("rb") as f:
return tomllib.load(f)
def load_normalized_case(case_path: Path, fixture_generator: Path | None) -> dict[str, Any]:
if fixture_generator is None:
raise ValueError("--fixture-generator query_perf_fixture is required for Rust-owned case planning")
result = run_command([str(fixture_generator), "plan", "--case", str(case_path)])
if result["returncode"] != 0:
raise RuntimeError(f"query_perf_fixture plan failed: {result['stderr'][:2000]}")
plan = json.loads(result["stdout"])
return {"case": load_case(case_path).get("case", {}), "scenario": plan["scenario"], "schema_version": plan.get("schema_version")}
def require_binary(path: Path, name: str, *, dry_run: bool) -> None:
if dry_run:
return
if not path.exists():
raise FileNotFoundError(f"{name} binary does not exist: {path}")
if not os.access(path, os.X_OK):
raise PermissionError(f"{name} binary is not executable: {path}")
def allocate_ports(n: int) -> list[int]:
socks = []
try:
for _ in range(n):
s = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
s.bind(("127.0.0.1", 0))
socks.append(s)
return [s.getsockname()[1] for s in socks]
finally:
for s in socks:
s.close()
def make_target(name: str, binary: Path, root: Path, ports: list[int], fixture_dir: Path | None = None, scheduler_env: dict[str, str] | None = None) -> RunTarget:
work_dir = root / name
return RunTarget(
name,
binary,
work_dir,
work_dir / "cluster_data",
fixture_dir or root / "fixture",
work_dir / "report.json",
ports[4],
ports[5],
ports[6],
ports[7],
ports[0],
ports[1],
ports[2],
ports[3],
work_dir / "datanode-0" / "data",
scheduler_env or {},
)
def scheduler_env(enable: bool, scheduler: dict[str, Any] | None) -> dict[str, str]:
"""Derive the datanode workload-scheduler environment variables for a target.
The datanode is spawned with CLI args only (no --config-file), so the
scheduler must be injected via environment. Config env vars use the
`PREFIX__KEY__SUBKEY` form (see src/common/config/src/config.rs
load_layered_options), and the datanode env prefix defaults to
`GREPTIMEDB_DATANODE`; the scheduler lives at
`runtime.experimental_workload_scheduler`.
Returns an empty dict when the case has no `[scenario.scheduler]` section:
both base and candidate then run with the datanode default (scheduler
disabled). When the section is present the base target only pins
`ENABLE=false` while the candidate pins `ENABLE=true` and forwards the
scheduler weights/max-polls so both targets are configured identically
except for the enable flag.
"""
if scheduler is None:
return {}
prefix = "GREPTIMEDB_DATANODE__RUNTIME__EXPERIMENTAL_WORKLOAD_SCHEDULER"
env = {f"{prefix}__ENABLE": "true" if enable else "false"}
if enable:
for key in ("max_concurrent_polls", "query_weight", "write_weight"):
if key in scheduler:
env[f"{prefix}__{key.upper()}"] = str(scheduler[key])
return env
def scheduler_report_entry(target_name: str, scheduler: dict[str, Any] | None) -> dict[str, Any]:
"""Scheduler configuration applied to a target, for the target report entry.
The base target always runs with the scheduler disabled; the candidate
runs with it enabled iff the case has a `[scenario.scheduler]` section.
Defaults mirror the datanode runtime defaults (max_concurrent_polls 0 =
four times global_rt_size, query_weight 2, write_weight 8).
"""
enabled = target_name == "candidate" and scheduler is not None
return {
"enabled": enabled,
"max_concurrent_polls": int(scheduler["max_concurrent_polls"]) if scheduler else 0,
"query_weight": int(scheduler["query_weight"]) if scheduler else 2,
"write_weight": int(scheduler["write_weight"]) if scheduler else 8,
}
def run_command(cmd: list[str], cwd: Path | None = None) -> dict[str, Any]:
started = time.monotonic()
proc = subprocess.run(cmd, cwd=cwd, text=True, stdout=subprocess.PIPE, stderr=subprocess.PIPE, check=False)
return {"cmd": cmd, "cwd": str(cwd) if cwd else None, "returncode": proc.returncode, "elapsed_seconds": time.monotonic() - started, "stdout": proc.stdout, "stderr": proc.stderr}
def sql_ident(name: str) -> str:
return '"' + name.replace('"', '""') + '"'
def sql_string(value: str) -> str:
return "'" + value.replace("'", "''") + "'"
def column_sql(col: dict[str, Any]) -> str:
return f"{sql_ident(col['name'])} {col['type']}"
def scenario(case: dict[str, Any]) -> dict[str, Any]:
value = case.get("scenario")
if not isinstance(value, dict):
raise ValueError("case requires [scenario] with a supported kind")
kind = value.get("kind")
supported = ("direct_readable_sst", "prom_remote_write_then_query", "otlp_trace_load", "write_throughput")
if kind not in supported:
raise ValueError(f"unsupported scenario kind {kind!r}; supported: {', '.join(supported)}")
return value
def case_tables(case: dict[str, Any]) -> list[dict[str, Any]]:
value = scenario(case)
if value.get("kind") == "otlp_trace_load":
load = value["load"]
return [{"database": load["database"], "name": load["table"], "engine": "trace", "validate_show_create_engine": False}]
if value.get("kind") in ("prom_remote_write_then_query", "write_throughput"):
remote = value["remote_write"]
metric = remote["metric"]
database = remote["database"]
return [{"database": database, "name": metric, "engine": "metric", "validate_show_create_engine": False}]
tables = value.get("tables") or []
if not tables or value.get("layout", {}).get("regions") != 1:
raise ValueError("runner supports one or more tables and exactly one region per table")
pairs = [(table.get("database"), table.get("name")) for table in tables]
if len(set(pairs)) != len(pairs):
raise ValueError("duplicate (database, name) table entries are not supported")
names = [table.get("name") for table in tables]
if len(set(names)) != len(names):
raise ValueError("duplicate table names are not supported because fixture generator --table selects by name")
return list(tables)
def fixture_subdir(table: dict[str, Any], index: int) -> str:
raw = f"{index:02d}_{table['database']}_{table['name']}"
safe = re.sub(r"[^A-Za-z0-9_.-]+", "_", raw).strip("._-")
return safe or f"table_{index:02d}"
def safe_path_component(raw: str) -> str:
safe = re.sub(r"[^A-Za-z0-9_.-]+", "_", raw).strip("._-")
return safe or "query_perf_case"
def fixture_root(work_root: Path, case_path: Path, case: dict[str, Any], fixture_cache_dir: Path | None) -> Path:
if fixture_cache_dir is None:
return work_root / "fixture"
case_name = case.get("case", {}).get("name") or case_path.parent.name or case_path.stem
fixture_config = dict(scenario(case))
fixture_config.pop("queries", None)
digest = hashlib.sha256(json.dumps(fixture_config, sort_keys=True).encode()).hexdigest()[:16]
return fixture_cache_dir.resolve() / f"{safe_path_component(case_name)}-{digest}"
def table_fixture_dir(root: Path, tables: list[dict[str, Any]], table: dict[str, Any], index: int) -> Path:
return root if len(tables) == 1 else root / fixture_subdir(table, index)
def create_table_sql(table: dict[str, Any]) -> str:
cols = ",\n ".join(column_sql(c) for c in table["columns"])
pk = ", ".join(sql_ident(c) for c in table.get("primary_key", []))
opts = []
if "append_mode" in table:
opts.append(("append_mode", str(table["append_mode"]).lower()))
if table.get("sst_format"):
opts.append(("sst_format", table["sst_format"]))
with_sql = ""
if opts:
with_sql = "\nWITH (" + ", ".join(f"'{k}'='{v}'" for k, v in opts) + ")"
return f"CREATE TABLE {sql_ident(table['name'])} (\n {cols},\n TIME INDEX ({sql_ident(table['time_index'])}),\n PRIMARY KEY ({pk})\n) ENGINE=mito{with_sql};"
def planned_queries(case: dict[str, Any]) -> list[dict[str, Any]]:
queries = scenario(case).get("queries", [])
if isinstance(queries, dict):
return [{"name": name, **value} for name, value in queries.items()]
return list(queries)
def http_post_sql(port: int, sql: str, db: str, timeout: float) -> dict[str, Any]:
data = urllib.parse.urlencode({"sql": sql, "db": db, "format": "json"}).encode()
req = urllib.request.Request(f"http://127.0.0.1:{port}/v1/sql", data=data, method="POST")
started = time.monotonic()
elapsed_ms = (time.monotonic() - started) * 1000.0
try:
with urllib.request.urlopen(req, timeout=timeout) as resp:
raw = resp.read().decode()
status = resp.status
elapsed_ms = (time.monotonic() - started) * 1000.0
try:
body: Any = json.loads(raw)
except json.JSONDecodeError:
body = {"raw": raw}
ok = status < 400 and not response_has_error(body)
return {"ok": ok, "status": status, "latency_ms": elapsed_ms, "response": body, "sql": sql}
except urllib.error.HTTPError as e:
elapsed_ms = (time.monotonic() - started) * 1000.0
raw = e.read().decode(errors="replace")
try:
body = json.loads(raw)
except json.JSONDecodeError:
body = {"raw": raw}
return {
"ok": False,
"status": e.code,
"latency_ms": elapsed_ms,
"response": body,
"error": repr(e),
"sql": sql,
}
except Exception as e: # noqa: BLE001 - report HTTP/query failures as JSON
elapsed_ms = (time.monotonic() - started) * 1000.0
return {"ok": False, "status": None, "latency_ms": elapsed_ms, "error": repr(e), "sql": sql}
def response_has_error(body: Any) -> bool:
"""Detect top-level GreptimeDB HTTP error envelopes without inspecting rows.
Query outputs may legitimately contain columns named `code` or `error`, so
this intentionally avoids recursive checks through result rows.
"""
if isinstance(body, dict):
if body.get("error") or body.get("err_msg") or body.get("error_msg"):
return True
if "error_code" in body and str(body.get("error_code", "")).lower() not in ("", "0", "success"):
return True
if "code" in body and "output" not in body and str(body.get("code", "")).lower() not in ("", "0", "success"):
return True
return False
def wait_health(port: int, timeout_s: float = 60.0) -> None:
deadline = time.monotonic() + timeout_s
last = None
while time.monotonic() < deadline:
try:
with urllib.request.urlopen(f"http://127.0.0.1:{port}/health", timeout=2) as r:
if r.status < 500:
return
except Exception as e: # noqa: BLE001 - diagnostic loop
last = e
time.sleep(0.5)
raise TimeoutError(f"health check timed out on port {port}: {last}")
class DistributedCluster:
"""Local metasrv + one datanode + frontend cluster for query-mode runs."""
def __init__(self, target: RunTarget):
self.target = target
self.procs: dict[str, subprocess.Popen[bytes]] = {}
def component_report(self) -> dict[str, Any]:
return {
"metasrv": {
"grpc": f"127.0.0.1:{self.target.metasrv_rpc_port}",
"http": f"127.0.0.1:{self.target.metasrv_http_port}",
"logs": str(self.target.work_dir / "logs" / "metasrv"),
},
"datanode_0": {
"node_id": 0,
"grpc": f"127.0.0.1:{self.target.datanode_rpc_port}",
"http": f"127.0.0.1:{self.target.datanode_http_port}",
"data_home": str(self.target.datanode_data_dir),
"logs": str(self.target.work_dir / "logs" / "datanode-0"),
},
"frontend": {
"http": f"127.0.0.1:{self.target.http_port}",
"grpc": f"127.0.0.1:{self.target.grpc_port}",
"mysql": f"127.0.0.1:{self.target.mysql_port}",
"postgres": f"127.0.0.1:{self.target.postgres_port}",
"logs": str(self.target.work_dir / "logs" / "frontend"),
},
}
def _spawn(self, name: str, args: list[str], env: dict[str, str] | None = None) -> None:
logs = self.target.work_dir / "logs" / name
logs.mkdir(parents=True, exist_ok=True)
with (logs / "stdout.log").open("ab") as out, (logs / "stderr.log").open("ab") as err:
self.procs[name] = subprocess.Popen(args, stdout=out, stderr=err, env=env)
def _ensure_metasrv_alive(self) -> None:
proc = self.procs.get("metasrv")
if proc is None or proc.poll() is not None:
raise RuntimeError("metasrv exited; memory-store metadata is no longer valid")
def start_metasrv(self) -> None:
log_dir = self.target.work_dir / "logs" / "metasrv"
self._spawn(
"metasrv",
[
str(self.target.binary),
"metasrv",
"start",
"--grpc-bind-addr",
f"127.0.0.1:{self.target.metasrv_rpc_port}",
"--grpc-server-addr",
f"127.0.0.1:{self.target.metasrv_rpc_port}",
"--http-addr",
f"127.0.0.1:{self.target.metasrv_http_port}",
"--backend",
"memory-store",
"--enable-region-failover",
"false",
"--log-dir",
str(log_dir),
],
)
wait_health(self.target.metasrv_http_port)
def start_datanode(self) -> None:
self._ensure_metasrv_alive()
log_dir = self.target.work_dir / "logs" / "datanode-0"
self.target.datanode_data_dir.mkdir(parents=True, exist_ok=True)
self._spawn(
"datanode",
[
str(self.target.binary),
"datanode",
"start",
"--grpc-bind-addr",
f"127.0.0.1:{self.target.datanode_rpc_port}",
"--grpc-server-addr",
f"127.0.0.1:{self.target.datanode_rpc_port}",
"--http-addr",
f"127.0.0.1:{self.target.datanode_http_port}",
"--data-home",
str(self.target.datanode_data_dir),
"--log-dir",
str(log_dir),
"--node-id",
"0",
"--metasrv-addrs",
f"127.0.0.1:{self.target.metasrv_rpc_port}",
],
env=None
if not self.target.scheduler_env
else {**os.environ, **self.target.scheduler_env},
)
wait_health(self.target.datanode_http_port)
def start_frontend(self, config_file: Path | None = None) -> None:
self._ensure_metasrv_alive()
log_dir = self.target.work_dir / "logs" / "frontend"
cmd = [
str(self.target.binary),
"frontend",
"start",
]
if config_file is not None:
cmd += ["--config-file", str(config_file)]
cmd += [
"--metasrv-addrs",
f"127.0.0.1:{self.target.metasrv_rpc_port}",
"--http-addr",
f"127.0.0.1:{self.target.http_port}",
"--grpc-bind-addr",
f"127.0.0.1:{self.target.grpc_port}",
"--grpc-server-addr",
f"127.0.0.1:{self.target.grpc_port}",
"--mysql-addr",
f"127.0.0.1:{self.target.mysql_port}",
"--postgres-addr",
f"127.0.0.1:{self.target.postgres_port}",
"--log-dir",
str(log_dir),
]
self._spawn("frontend", cmd)
wait_health(self.target.http_port)
def start_all(self, frontend_config: Path | None = None) -> None:
self.start_metasrv()
self.start_datanode()
self.start_frontend(frontend_config)
def stop_component(self, name: str) -> None:
proc = self.procs.pop(name, None)
if proc is None:
return
proc.terminate()
try:
proc.wait(timeout=20)
except subprocess.TimeoutExpired:
proc.kill()
proc.wait(timeout=20)
def stop_all(self) -> None:
for name in ("frontend", "datanode", "metasrv"):
self.stop_component(name)
def restart_frontend(self) -> None:
self.stop_component("frontend")
self.start_frontend()
def extract_rows(body: Any) -> list[Any]:
"""Extract rows from GreptimeDB HTTP JSON in a tolerant way."""
rows: list[Any] = []
if isinstance(body, dict):
for key in ("data", "rows", "records", "output"):
if key in body:
value = body[key]
if key in ("data", "rows") and isinstance(value, list):
rows.extend(value)
else:
rows.extend(extract_rows(value))
elif isinstance(body, list):
if body and all(not isinstance(v, (dict, list)) for v in body):
rows.append(body)
else:
for item in body:
rows.extend(extract_rows(item))
return rows
def row_value(row: Any, idx: int, name: str) -> Any:
if isinstance(row, dict):
for key in (name, name.upper(), name.lower()):
if key in row:
return row[key]
return None
if isinstance(row, list):
return row[idx]
return row
def discover_region_via_frontend(target: RunTarget, table: dict[str, Any], http_timeout: float) -> dict[str, Any]:
schema = table["database"]
table_name = table["name"].replace("'", "''")
schema_name = schema.replace("'", "''")
table_sql = f"SELECT table_id FROM information_schema.tables WHERE table_schema = '{schema_name}' AND table_name = '{table_name}'"
table_result = http_post_sql(target.http_port, table_sql, schema, http_timeout)
if not table_result["ok"]:
raise RuntimeError(f"table_id discovery failed: {table_result}")
table_rows = extract_rows(table_result.get("response"))
if len(table_rows) != 1:
raise RuntimeError(f"expected one information_schema.tables row, got {len(table_rows)}: {table_result}")
table_id = int(row_value(table_rows[0], 0, "table_id"))
region_sql = f"SELECT region_id, peer_id, peer_addr, is_leader, status FROM information_schema.region_peers WHERE table_schema = '{schema_name}' AND table_name = '{table_name}'"
region_result = http_post_sql(target.http_port, region_sql, schema, http_timeout)
if not region_result["ok"]:
raise RuntimeError(f"region_peers discovery failed: {region_result}")
region_rows = extract_rows(region_result.get("response"))
if len(region_rows) != 1:
raise RuntimeError(f"expected one information_schema.region_peers row, got {len(region_rows)}: {region_result}")
row = region_rows[0]
region_id = int(row_value(row, 0, "region_id"))
peer_id = int(row_value(row, 1, "peer_id"))
peer_addr = row_value(row, 2, "peer_addr")
is_leader = row_value(row, 3, "is_leader")
status = row_value(row, 4, "status")
if str(is_leader).lower() not in ("yes", "true"):
raise RuntimeError(f"expected leader region peer, got is_leader={is_leader}: {region_result}")
if str(status).upper() != "ALIVE":
raise RuntimeError(f"expected ALIVE region peer, got status={status}: {region_result}")
if peer_id != 0:
raise RuntimeError(f"expected region leader on datanode peer_id=0, got {peer_id}")
computed_table_id = region_id >> 32
region_seq = region_id & 0xFFFFFFFF
if computed_table_id != table_id:
raise RuntimeError(f"table_id mismatch: tables={table_id}, region_id-derived={computed_table_id}")
table_dir = f"data/greptime/{schema}/{table_id}/"
region_dir = f"data/greptime/{schema}/{table_id}/{table_id}_{region_seq:010}"
return {
"catalog": "greptime",
"schema": schema,
"table_id": table_id,
"region_id": region_id,
"region_seq": region_seq,
"table_dir": table_dir,
"region_dir": region_dir,
"peer_id": peer_id,
"peer_addr": peer_addr,
"is_leader": is_leader,
"status": status,
"discovery_queries": {"table": table_result, "region": region_result},
}
def assert_fixture_summary(summary: dict[str, Any], *, region_id: int | None, table_dir: str | None, region_dir: str | None = None, table_name: str | None = None, database: str | None = None) -> None:
if table_name is not None and summary.get("table") != table_name:
raise RuntimeError(f"fixture table mismatch: {summary.get('table')} != {table_name}")
if database is not None and summary.get("database") != database:
raise RuntimeError(f"fixture database mismatch: {summary.get('database')} != {database}")
if region_id is not None and int(summary["region_id"]) != region_id:
raise RuntimeError(f"fixture region_id mismatch: {summary['region_id']} != {region_id}")
if table_dir is not None and summary["table_dir"] != table_dir:
raise RuntimeError(f"fixture table_dir mismatch: {summary['table_dir']} != {table_dir}")
if region_dir is not None and summary["region_dir"].strip("/") != region_dir.strip("/"):
raise RuntimeError(f"fixture region_dir mismatch: {summary['region_dir']} != {region_dir}")
def generate_fixture(generator: Path | None, case_path: Path, fixture_dir: Path, *, dry_run: bool, reuse_fixture: bool, allow_large_fixture: bool, table_name: str | None = None, database: str | None = None, region_id: int | None = None, table_dir: str | None = None, region_dir: str | None = None) -> dict[str, Any]:
if generator is None:
return {"status": "skipped", "reason": "no fixture generator provided"}
summary_path = fixture_dir / "summary.json"
reuse_miss: str | None = None
if reuse_fixture and summary_path.exists():
summary = json.loads(summary_path.read_text())
try:
assert_fixture_summary(summary, region_id=region_id, table_dir=table_dir, region_dir=region_dir, table_name=table_name, database=database)
return {"status": "reused", "fixture_dir": str(fixture_dir), "summary": summary}
except RuntimeError as e:
reuse_miss = str(e)
shutil.rmtree(fixture_dir)
if fixture_dir.exists() and not reuse_fixture:
shutil.rmtree(fixture_dir)
fixture_dir.mkdir(parents=True, exist_ok=True)
cmd = [str(generator), "direct-sst", "--case", str(case_path), "--out-dir", str(fixture_dir)]
if table_name is not None:
cmd += ["--table", table_name]
if region_id is not None:
cmd += ["--region-id", str(region_id)]
if table_dir is not None:
cmd += ["--table-dir", table_dir]
if allow_large_fixture:
cmd.append("--allow-large")
if dry_run:
return {"status": "dry-run", "cmd": cmd}
result = run_command(cmd)
result["status"] = "ok" if result["returncode"] == 0 else "failed"
if result["returncode"] != 0:
raise RuntimeError(f"fixture generator failed: {result['stderr'][:2000]}")
summary = json.loads(summary_path.read_text())
assert_fixture_summary(summary, region_id=region_id, table_dir=table_dir, region_dir=region_dir, table_name=table_name, database=database)
result["summary"] = summary
if reuse_miss is not None:
result["reuse_miss"] = reuse_miss
return result
def copy_stats() -> dict[str, int]:
return {"files": 0, "dirs": 0, "reflinked": 0, "hardlinked": 0, "copied": 0}
def merge_copy_stats(left: dict[str, int], right: dict[str, int]) -> None:
for key, value in right.items():
left[key] = left.get(key, 0) + value
def reflink_file(src: Path, dst: Path) -> bool:
try:
dst.parent.mkdir(parents=True, exist_ok=True)
sfd = os.open(src, os.O_RDONLY)
try:
dfd = os.open(dst, os.O_WRONLY | os.O_CREAT | os.O_TRUNC, src.stat().st_mode & 0o777)
try:
fcntl.ioctl(dfd, FICLONE, sfd)
shutil.copystat(src, dst, follow_symlinks=True)
return True
finally:
os.close(dfd)
finally:
os.close(sfd)
except OSError:
try:
dst.unlink()
except OSError:
pass
return False
def materialize_file(src: Path, dst: Path, *, allow_hardlink: bool) -> str:
if reflink_file(src, dst):
return "reflinked"
if allow_hardlink:
try:
dst.parent.mkdir(parents=True, exist_ok=True)
try:
dst.unlink()
except FileNotFoundError:
pass
os.link(src, dst)
return "hardlinked"
except OSError:
pass
dst.parent.mkdir(parents=True, exist_ok=True)
shutil.copy2(src, dst)
return "copied"
def copy_tree_contents(src: Path, dst: Path, *, allow_hardlinks: bool = True) -> dict[str, int]:
if not src.exists():
raise FileNotFoundError(f"fixture path does not exist: {src}")
stats = copy_stats()
dst.mkdir(parents=True, exist_ok=True)
for child in src.iterdir():
target = dst / child.name
if child.is_dir():
stats["dirs"] += 1
merge_copy_stats(stats, copy_tree_contents(child, target, allow_hardlinks=allow_hardlinks))
else:
mode = materialize_file(child, target, allow_hardlink=allow_hardlinks)
stats["files"] += 1
stats[mode] += 1
return stats
def materialize_fixture(target: RunTarget, *, dry_run: bool, preserve_state: bool, expected_region_dir: str | None = None, fixture_dir: Path | None = None, reset_data: bool = True) -> dict[str, Any]:
fixture_dir = fixture_dir or target.fixture_dir
summary_path = fixture_dir / "summary.json"
materialize_root = target.datanode_data_dir if preserve_state else target.data_dir
if dry_run:
return {"status": "dry-run", "summary_path": str(summary_path), "data_dir": str(materialize_root)}
summary = json.loads(summary_path.read_text())
object_store_dir = fixture_dir / "object-store"
manifest_dir = fixture_dir / "manifest"
region_dir = summary["region_dir"].strip("/")
target_region_dir = materialize_root / region_dir
data_root = materialize_root.resolve()
resolved_region_dir = target_region_dir.resolve(strict=False)
if data_root != resolved_region_dir and data_root not in resolved_region_dir.parents:
raise RuntimeError(f"unsafe fixture region_dir escapes data_dir: {region_dir}")
if expected_region_dir is not None and region_dir != expected_region_dir.strip("/"):
raise RuntimeError(f"fixture region_dir {region_dir} does not equal discovered {expected_region_dir}")
target_manifest_dir = target_region_dir / "manifest"
if not preserve_state and reset_data and materialize_root.exists():
shutil.rmtree(materialize_root)
materialize_root.mkdir(parents=True, exist_ok=True)
if preserve_state and target_region_dir.exists():
shutil.rmtree(target_region_dir)
object_stats: dict[str, int]
if preserve_state:
fixture_region_dir = object_store_dir / region_dir
object_stats = copy_tree_contents(fixture_region_dir, target_region_dir, allow_hardlinks=True)
else:
object_stats = copy_tree_contents(object_store_dir, materialize_root, allow_hardlinks=True)
if target_manifest_dir.exists():
shutil.rmtree(target_manifest_dir)
manifest_stats = copy_tree_contents(manifest_dir, target_manifest_dir, allow_hardlinks=False)
return {"status": "ok", "summary_path": str(summary_path), "data_dir": str(materialize_root), "region_dir": region_dir, "manifest_dir": str(target_manifest_dir), "preserve_state": preserve_state, "copy_stats": {"object_store": object_stats, "manifest": manifest_stats}}
def response_text(body: Any) -> str:
return json.dumps(body, sort_keys=True) if not isinstance(body, str) else body
def validate_show_create(result: dict[str, Any], table: dict[str, Any]) -> list[str]:
text = response_text(result.get("response", {})).lower()
errors = []
if table["name"].lower() not in text:
errors.append("SHOW CREATE output does not contain table name")
if table.get("validate_show_create_engine", True) and ("engine" not in text or "mito" not in text):
errors.append("SHOW CREATE output does not mention ENGINE=mito")
if "append_mode" in table and "append_mode" not in text:
errors.append("SHOW CREATE output does not mention append_mode")
if table.get("sst_format") and "sst_format" not in text:
errors.append("SHOW CREATE output does not mention sst_format")
return errors
def run_queries(target: RunTarget, case: dict[str, Any], tables: list[dict[str, Any]], http_timeout: float) -> dict[str, Any]:
table = tables[0]
db = table["database"]
queries = planned_queries(case)
if not queries:
queries = [{"name": "count_all", "kind": "sql", "query": f"SELECT count(*) FROM {sql_ident(table['name'])}", "warmup": 0, "iterations": 1}]
validations = []
validation_errors = []
for table in tables:
for sql in [f"SHOW CREATE TABLE {sql_ident(table['name'])}"]:
result = http_post_sql(target.http_port, sql, table["database"], http_timeout)
if not result["ok"]:
validation_errors.append({"sql": sql, "error": result.get("error"), "response": result.get("response")})
else:
for error in validate_show_create(result, table):
validation_errors.append({"sql": sql, "error": error, "response": result.get("response")})
validations.append(result)
if queries:
result = http_post_sql(target.http_port, queries[0]["query"], db, http_timeout)
if not result["ok"]:
validation_errors.append({"sql": queries[0]["query"], "error": result.get("error"), "response": result.get("response")})
validations.append(result)
measurements = []
for q in queries:
for _ in range(int(q.get("warmup", 0))):
warmup = http_post_sql(target.http_port, q["query"], db, http_timeout)
if not warmup["ok"]:
validation_errors.append({"sql": q["query"], "phase": "warmup", "error": warmup.get("error"), "response": warmup.get("response")})
samples = []
for _ in range(int(q.get("iterations", 1))):
result = http_post_sql(target.http_port, q["query"], db, http_timeout)
result["execution_time_ms"] = extract_execution_time(result.get("response"))
samples.append(result)
good_lats = [s["latency_ms"] for s in samples if s["ok"]]
measurements.append({"name": q.get("name"), "kind": q.get("kind"), "iterations": len(samples), "samples": samples, "latency_ms_median": statistics.median(good_lats) if good_lats else None, "latency_ms_p95": percentile(good_lats, 95) if good_lats else None, "status": "ok" if len(good_lats) == len(samples) else "failed"})
return {"validation": validations, "validation_errors": validation_errors, "measurements": measurements, "status": "failed" if validation_errors or any(m["status"] == "failed" for m in measurements) else "ok"}
def extract_execution_time(body: Any) -> Any:
if isinstance(body, dict):
for key in ("execution_time_ms", "execution_time", "elapsed"):
if key in body:
return body[key]
for value in body.values():
found = extract_execution_time(value)
if found is not None:
return found
if isinstance(body, list):
for value in body:
found = extract_execution_time(value)
if found is not None:
return found
return None
def percentile(values: list[float], pct: float) -> float:
if not values:
return 0.0
ordered = sorted(values)
idx = min(len(ordered) - 1, max(0, round((pct / 100.0) * (len(ordered) - 1))))
return ordered[idx]
def enforce_thresholds(case: dict[str, Any], base: dict[str, Any], candidate: dict[str, Any]) -> list[dict[str, Any]]:
results = []
base_by_name = {m["name"]: m for m in base.get("measurements", [])}
for cm in candidate.get("measurements", []):
bm = base_by_name.get(cm["name"])
qcfg = next((q for q in planned_queries(case) if q.get("name") == cm["name"]), {})
th = qcfg.get("thresholds") or {}
max_reg = th.get("max_candidate_latency_regression_pct")
if max_reg is not None and not bm:
results.append({"query": cm["name"], "threshold": "max_candidate_latency_regression_pct", "status": "failed", "reason": "missing base measurement"})
elif max_reg is not None and bm and bm.get("latency_ms_median") in (None, 0):
results.append({"query": cm["name"], "threshold": "max_candidate_latency_regression_pct", "status": "failed", "reason": "base median latency is missing or zero", "base_latency_ms_median": bm.get("latency_ms_median")})
elif max_reg is not None and bm and cm.get("latency_ms_median") is None:
results.append({"query": cm["name"], "threshold": "max_candidate_latency_regression_pct", "status": "failed", "reason": "missing candidate measurement"})
elif bm and max_reg is not None:
ratio = (cm["latency_ms_median"] - bm["latency_ms_median"]) / bm["latency_ms_median"] * 100.0
results.append({"query": cm["name"], "threshold": "max_candidate_latency_regression_pct", "status": "passed" if ratio <= max_reg else "failed", "actual_pct": ratio, "limit_pct": max_reg})
for key in th:
if key != "max_candidate_latency_regression_pct":
results.append({"query": cm["name"], "threshold": key, "status": "failed", "reason": "unsupported threshold"})
return results
def storage_config(remote: dict[str, Any]) -> dict[str, Any] | None:
value = remote["storage"]
if not isinstance(value, dict):
return None
if value["inspect"] is False:
return None
return value
def run_storage_inspection(helper: Path | None, target: RunTarget, storage: dict[str, Any], *, dry_run: bool) -> dict[str, Any]:
if helper is None:
if not dry_run:
raise ValueError("--fixture-generator query_perf_fixture is required when storage inspection is enabled")
helper = Path("query_perf_fixture")
root = target.datanode_data_dir
if storage.get("root_suffix"):
root = root / str(storage["root_suffix"])
cmd = [
str(helper),
"inspect-footer",
"--root", str(root),
"--column", str(storage["column"]),
]
if storage["include_metadata_files"]:
cmd.append("--include-metadata-files")
if dry_run:
return {"status": "dry-run", "cmd": cmd, "root": str(root)}
result = run_command(cmd)
result["status"] = "ok" if result["returncode"] == 0 else "failed"
if result["returncode"] != 0:
raise RuntimeError(f"storage inspector failed for {target.name}: {result['stderr'][:2000]}")
try:
result["summary"] = json.loads(result["stdout"])
except json.JSONDecodeError:
result["summary_parse_error"] = result["stdout"]
result["status"] = "failed"
return result
def storage_summary(inspection: dict[str, Any]) -> dict[str, Any]:
summary = inspection.get("summary")
if isinstance(summary, dict):
nested = summary.get("summary")
if isinstance(nested, dict):
return nested
return {}
def enforce_storage_thresholds(storage: dict[str, Any] | None, base: dict[str, Any] | None, candidate: dict[str, Any] | None) -> list[dict[str, Any]]:
if not storage:
return []
results: list[dict[str, Any]] = []
target_summaries = [("base", storage_summary(base or {})), ("candidate", storage_summary(candidate or {}))]
cand = target_summaries[1][1]
base_summary = storage_summary(base or {})
def add_abs(name: str, field: str) -> None:
limit = storage.get(name)
if limit is None:
return
for target_name, summary in target_summaries:
actual = summary.get(field)
ok = actual is not None and float(actual) <= float(limit)
results.append({"target": target_name, "threshold": name, "status": "passed" if ok else "failed", "actual": actual, "limit": limit})
def add_min(name: str, field: str) -> None:
limit = storage[name]
if limit is None:
return
for target_name, summary in target_summaries:
actual = summary.get(field)
ok = actual is not None and int(actual) >= int(limit)
results.append({"target": target_name, "threshold": name, "status": "passed" if ok else "failed", "actual": actual, "limit": limit})
def add_cmp(name: str, field: str) -> None:
limit = storage.get(name)
if limit is None:
return
base_value = base_summary.get(field)
candidate_value = cand.get(field)
if base_value in (None, 0) or candidate_value is None:
results.append({"threshold": name, "status": "failed", "reason": "missing or zero base/candidate value", "base": base_value, "candidate": candidate_value})
return
actual = (float(candidate_value) - float(base_value)) / float(base_value) * 100.0
results.append({"threshold": name, "status": "passed" if actual <= float(limit) else "failed", "actual_pct": actual, "limit_pct": limit, "base": base_value, "candidate": candidate_value})
add_min("min_files", "file_count")
add_min("min_files_with_column", "files_with_column")
add_abs("max_total_file_size_bytes", "total_file_size")
add_abs("max_column_compressed_size_bytes", "column_compressed_size")
add_abs("max_column_uncompressed_size_bytes", "column_uncompressed_size")
for target_name, summary in target_summaries:
encodings = set(summary.get("unique_encodings") or [])
for required in storage["require_encodings"]:
results.append({"target": target_name, "threshold": "require_encodings", "encoding": required, "status": "passed" if required in encodings else "failed"})
for forbidden in storage["forbid_encodings"]:
results.append({"target": target_name, "threshold": "forbid_encodings", "encoding": forbidden, "status": "passed" if forbidden not in encodings else "failed"})
add_cmp("max_candidate_total_file_size_regression_pct", "total_file_size")
add_cmp("max_candidate_column_compressed_size_regression_pct", "column_compressed_size")
add_cmp("max_candidate_column_uncompressed_size_regression_pct", "column_uncompressed_size")
return results
def planned_storage_thresholds(storage: dict[str, Any] | None) -> list[dict[str, Any]]:
if not storage:
return []
return list(storage["planned_thresholds"])
REGION_DIR_RE = re.compile(r"^(\d+)_(\d{10})$")
def inspection_report(inspection: dict[str, Any] | None) -> dict[str, Any]:
if not inspection:
return {}
report = inspection.get("summary")
return report if isinstance(report, dict) else {}
def inspected_relative_path(target: RunTarget, inspection: dict[str, Any], relative_path: str) -> Path:
report = inspection_report(inspection)
root_text = report.get("root") or inspection.get("root")
if not root_text:
return Path(relative_path)
root = Path(root_text).resolve(strict=False)
data_home = target.datanode_data_dir.resolve(strict=False)
try:
root_suffix = root.relative_to(data_home)
except ValueError as e:
raise RuntimeError(f"storage inspection root {root} is not under datanode data home {data_home}") from e
return root_suffix / relative_path
def parse_sst_bench_target(target: RunTarget, inspection: dict[str, Any], file: dict[str, Any], *, dry_run: bool) -> dict[str, str] | None:
relative_path = str(file.get("relative_path") or "")
if dry_run and (not relative_path or "<" in relative_path):
return {
"relative_path": relative_path or "<table-dir>/<table-id>_<region-seq>/data/<file-id>.parquet",
"table_dir": "<table-dir>/",
"region_id": "<table-id>:<region-seq>",
"path_type": "data",
"file_id": "<file-id>",
}
if not relative_path:
return None
full_relative = inspected_relative_path(target, inspection, relative_path)
parts = full_relative.parts
if not parts or not parts[-1].endswith(".parquet"):
return None
region_index = next((idx for idx, part in enumerate(parts) if REGION_DIR_RE.match(part)), None)
if region_index is None or region_index == 0:
if dry_run:
return {
"relative_path": relative_path,
"table_dir": "<table-dir>/",
"region_id": "<table-id>:<region-seq>",
"path_type": "data",
"file_id": Path(parts[-1]).stem,
}
return None
region_match = REGION_DIR_RE.match(parts[region_index])
assert region_match is not None
file_index = len(parts) - 1
path_type = "bare"
if region_index + 1 < file_index and parts[region_index + 1] in ("data", "metadata"):
path_type = parts[region_index + 1]
if path_type == "metadata":
return None
table_dir = "/".join(parts[:region_index]).rstrip("/") + "/"
return {
"relative_path": str(full_relative),
"table_dir": table_dir,
"region_id": f"{int(region_match.group(1))}:{int(region_match.group(2))}",
"path_type": path_type,
"file_id": Path(parts[-1]).stem,
}
def run_read_bench(bench_binary: Path, target: RunTarget, read_bench: dict[str, Any] | None, storage_inspection: dict[str, Any] | None, *, dry_run: bool) -> dict[str, Any]:
if not read_bench or read_bench["enabled"] is False:
return {"status": "skipped", "reason": "read_bench disabled"}
run_parquetbench = read_bench["parquetbench"]
run_scanbench = read_bench["scanbench"]
if not run_parquetbench and not run_scanbench:
return {"status": "skipped", "reason": "parquetbench and scanbench disabled"}
if storage_inspection is None:
raise ValueError("read_bench requires storage inspection")
bench_dir = target.work_dir / "read_bench"
config_toml = bench_dir / "bench.toml"
scan_json = bench_dir / "scan.json"
files = inspection_report(storage_inspection).get("files") or []
ssts = [f for f in files if f.get("relative_path", "").endswith(".parquet") and f.get("columns")]
if dry_run and not ssts:
ssts = [{"relative_path": "<landed-sst>.parquet"}]
if read_bench["max_files"] is not None:
ssts = ssts[: int(read_bench["max_files"])]
bench_targets = [parsed for f in ssts if (parsed := parse_sst_bench_target(target, storage_inspection, f, dry_run=dry_run)) is not None]
if not bench_targets:
return {"status": "failed", "reason": "no inspected data SST files available for read_bench"}
config_text = f'[storage]\ndata_home = "{target.datanode_data_dir}"\ntype = "File"\n\n[[region_engine]]\n[region_engine.mito]\n'
commands = []
parquet_runs = []
scan_runs = []
iterations = str(int(read_bench["iterations"]))
if run_parquetbench:
for bench_target in bench_targets:
cmd = [str(bench_binary), "datanode", "parquetbench", "--config", str(config_toml), "--region-id", bench_target["region_id"], "--table-dir", bench_target["table_dir"], "--file-id", bench_target["file_id"], "--scan-config", str(scan_json), "--path-type", bench_target["path_type"], "--iterations", iterations, "--reader", str(read_bench["parquet_reader"])]
commands.append(cmd)
parquet_runs.append({**bench_target, "cmd": cmd, "status": "dry-run" if dry_run else "planned"})
regions: dict[tuple[str, str, str], list[str]] = {}
if run_scanbench:
for bench_target in bench_targets:
regions.setdefault((bench_target["table_dir"], bench_target["region_id"], bench_target["path_type"]), []).append(bench_target["relative_path"])
for (table_dir, region_id, path_type), region_files in regions.items():
cmd = [str(bench_binary), "datanode", "scanbench", "--config", str(config_toml), "--region-id", region_id, "--table-dir", table_dir, "--scan-config", str(scan_json), "--path-type", path_type, "--scanner", str(read_bench["scan_scanner"]), "--parallelism", str(int(read_bench["parallelism"])), "--iterations", iterations]
commands.append(cmd)
scan_runs.append({"table_dir": table_dir, "region_id": region_id, "path_type": path_type, "files": region_files, "cmd": cmd, "status": "dry-run" if dry_run else "planned"})
if dry_run:
return {"status": "dry-run", "config_path": str(config_toml), "scan_config_path": str(scan_json), "commands": commands, "parquetbench": parquet_runs, "scanbench": scan_runs}
bench_dir.mkdir(parents=True, exist_ok=True)
config_toml.write_text(config_text)
scan_json.write_text(json.dumps({"projection_names": read_bench["projection"]}, indent=2) + "\n")
avg_re = re.compile(r"Average duration[^0-9]*([0-9.]+)\s*ms", re.I)
for runs in (parquet_runs, scan_runs):
for run in runs:
result = run_command(run["cmd"])
run.update(result)
run["status"] = "ok" if result["returncode"] == 0 else "failed"
m = avg_re.search(result.get("stdout", ""))
if m:
run["average_ms"] = float(m.group(1))
def median(runs: list[dict[str, Any]]) -> float | None:
vals = [r["average_ms"] for r in runs if "average_ms" in r]
return statistics.median(vals) if vals else None
return {"status": "ok" if all(r.get("status") == "ok" for r in parquet_runs + scan_runs) else "failed", "config_path": str(config_toml), "scan_config_path": str(scan_json), "parquetbench": parquet_runs, "scanbench": scan_runs, "aggregate": {"parquetbench_median_average_ms": median(parquet_runs), "scanbench_median_average_ms": median(scan_runs)}}
def write_json(path: Path, payload: dict[str, Any]) -> None:
path.parent.mkdir(parents=True, exist_ok=True)
path.write_text(json.dumps(payload, indent=2, sort_keys=True) + "\n")
def output_report(report: dict[str, Any], output: Path | None) -> None:
if output is None:
print(json.dumps(report, indent=2, sort_keys=True))
else:
write_json(output, report)
def parse_prometheus_metrics(text: str, names: set[str] = OTLP_TRACE_METRICS) -> dict[str, float]:
values: dict[str, float] = {}
for line in text.splitlines():
match = PROMETHEUS_SAMPLE_RE.match(line.strip())
if not match or match.group(1) not in names:
continue
value = float(match.group(2))
if not math.isfinite(value):
raise ValueError(f"non-finite Prometheus sample for {match.group(1)}")
values[match.group(1)] = values.get(match.group(1), 0.0) + value
return values
def fetch_otlp_metrics(target: RunTarget, http_timeout: float) -> dict[str, Any]:
with urllib.request.urlopen(f"http://127.0.0.1:{target.http_port}/metrics", timeout=http_timeout) as response:
text = response.read().decode()
return {"captured_monotonic_seconds": time.monotonic(), "values": parse_prometheus_metrics(text)}
def metric_delta(after: dict[str, Any], before: dict[str, Any], name: str) -> float:
delta = float(after["values"].get(name, 0.0)) - float(before["values"].get(name, 0.0))
if delta < 0:
raise RuntimeError(f"metric {name} decreased by {-delta}")
return delta
def parse_scheduler_poll_metrics(text: str) -> dict[str, int]:
"""Parse `greptime_workload_scheduler_polls{workload="..."}` samples.
Returns {workload: cumulative polls} for the workloads the datanode
scheduler tracks ("query", "write"); workloads absent from the scrape are
omitted (the base target runs with the scheduler disabled and typically
exposes no such samples at all).
"""
values: dict[str, int] = {}
for line in text.splitlines():
match = SCHEDULER_POLL_RE.match(line.strip())
if match:
values[match.group(1)] = int(match.group(2))
return values
def scrape_scheduler_polls(target: RunTarget, http_timeout: float) -> dict[str, Any]:
"""Best-effort datanode scheduler-poll snapshot (never a gate).
Returns {"captured_monotonic_seconds", "values": {workload: polls}} on
success or {"error", "values": {}} when the datanode /metrics endpoint is
unreachable or the scheduler metric is absent.
"""
try:
with urllib.request.urlopen(f"http://127.0.0.1:{target.datanode_http_port}/metrics", timeout=http_timeout) as response:
text = response.read().decode()
return {"captured_monotonic_seconds": time.monotonic(), "values": parse_scheduler_poll_metrics(text)}
except Exception as e: # noqa: BLE001 - best-effort diagnostic
return {"error": repr(e), "values": {}}
def scheduler_poll_deltas(after: dict[str, Any], before: dict[str, Any]) -> dict[str, Any]:
"""Delta of cumulative scheduler polls between two snapshots, per workload.
A workload whose counter is missing in either snapshot, or that decreased
(counter reset/restart), is reported as None because its delta is unknown.
"""
result: dict[str, Any] = {}
for workload in ("query", "write"):
before_value = before.get("values", {}).get(workload)
after_value = after.get("values", {}).get(workload)
if before_value is None or after_value is None:
result[workload] = None
else:
delta = after_value - before_value
result[workload] = delta if delta >= 0 else None
return result
def otelgen_command(otelgen_bin: Path, target: RunTarget, load: dict[str, Any]) -> list[str]:
return [
str(otelgen_bin),
"--protocol", "http",
"--otel-exporter-otlp-endpoint", f"127.0.0.1:{target.http_port}",
"--otel-exporter-otlp-url-path", "/v1/otlp/v1/traces",
"--header", f"x-greptime-pipeline-name={load['pipeline']}",
"--header", f"x-greptime-db-name={load['database']}",
"--header", f"x-greptime-trace-table-name={load['table']}",
"--log-level", "error",
"--insecure",
"--duration", str(int(load["duration_seconds"])),
"--rate", str(int(load["rate"])),
"traces", "multi",
"--workers", str(int(load["workers"])),
"--scenarios", str(load["workload"]),
"--exporter-shards", str(int(load["exporter_shards"])),
]
def run_otelgen_load(otelgen_bin: Path | None, target: RunTarget, load: dict[str, Any], http_timeout: float, *, dry_run: bool) -> dict[str, Any]:
binary = otelgen_bin or Path("otelgen")
cmd = otelgen_command(binary, target, load)
if dry_run:
return {"status": "dry-run", "cmd": cmd}
initial = fetch_otlp_metrics(target, http_timeout)
log_dir = target.work_dir / "otelgen"
log_dir.mkdir(parents=True, exist_ok=True)
stdout_path = log_dir / "stdout.log"
stderr_path = log_dir / "stderr.log"
duration = int(load["duration_seconds"])
warmup = int(load["warmup_seconds"])
timed_out = False
started = time.monotonic()
with stdout_path.open("wb") as stdout, stderr_path.open("wb") as stderr:
proc = subprocess.Popen(cmd, stdout=stdout, stderr=stderr)
try:
if warmup:
try:
proc.wait(timeout=warmup)
except subprocess.TimeoutExpired:
pass
warmed = fetch_otlp_metrics(target, http_timeout)
if proc.poll() is None:
try:
proc.wait(timeout=max(60, duration - warmup + 60))
except subprocess.TimeoutExpired:
timed_out = True
proc.kill()
proc.wait(timeout=20)
final = fetch_otlp_metrics(target, http_timeout)
finally:
if proc.poll() is None:
proc.kill()
proc.wait(timeout=20)
elapsed = time.monotonic() - started
return {
"status": "ok" if proc.returncode == 0 and not timed_out and elapsed + 1 >= duration else "failed",
"cmd": cmd,
"returncode": proc.returncode,
"timed_out": timed_out,
"elapsed_seconds": elapsed,
"stdout_path": str(stdout_path),
"stderr_path": str(stderr_path),
"snapshots": {"initial": initial, "warmup": warmed, "final": final},
}
def summarize_otlp_metrics(run: dict[str, Any]) -> dict[str, Any]:
initial = run["snapshots"]["initial"]
warmed = run["snapshots"]["warmup"]
final = run["snapshots"]["final"]
rows = "greptime_frontend_otlp_traces_rows"
failures = "greptime_frontend_otlp_traces_failure_count"
elapsed_sum = "greptime_servers_http_otlp_traces_elapsed_sum"
elapsed_count = "greptime_servers_http_otlp_traces_elapsed_count"
missing = sorted(name for name in (rows, elapsed_sum, elapsed_count) if name not in final["values"])
accepted_spans = int(round(metric_delta(final, initial, rows)))
measurement_accepted_spans = int(round(metric_delta(final, warmed, rows)))
http_requests = int(round(metric_delta(final, warmed, elapsed_count)))
latency_seconds = metric_delta(final, warmed, elapsed_sum)
measurement_seconds = final["captured_monotonic_seconds"] - warmed["captured_monotonic_seconds"]
return {
"accepted_spans": accepted_spans,
"measurement_accepted_spans": measurement_accepted_spans,
"accepted_spans_per_second": measurement_accepted_spans / measurement_seconds if measurement_seconds > 0 else None,
"http_requests": http_requests,
"mean_http_latency_ms": latency_seconds / http_requests * 1000.0 if http_requests else None,
"failure_count": int(round(metric_delta(final, initial, failures))),
"measurement_seconds": measurement_seconds,
"missing_metrics": missing,
}
def planned_otlp_thresholds(load: dict[str, Any]) -> list[dict[str, Any]]:
thresholds = load["thresholds"]
return [
{"threshold": "max_candidate_throughput_regression_pct", "status": "planned", "limit_pct": thresholds["max_candidate_throughput_regression_pct"]},
{"threshold": "max_candidate_mean_latency_regression_pct", "status": "planned", "limit_pct": thresholds["max_candidate_mean_latency_regression_pct"]},
{"target": "each", "threshold": "max_failure_count", "status": "planned", "limit": thresholds["max_failure_count"]},
]
def enforce_otlp_thresholds(load: dict[str, Any], base: dict[str, Any], candidate: dict[str, Any]) -> list[dict[str, Any]]:
thresholds = load["thresholds"]
results: list[dict[str, Any]] = []
for target_name, metrics in (("base", base), ("candidate", candidate)):
failures = metrics.get("failure_count")
limit = int(thresholds["max_failure_count"])
results.append({"target": target_name, "threshold": "max_failure_count", "status": "passed" if failures is not None and failures <= limit else "failed", "actual": failures, "limit": limit})
base_rate = base.get("accepted_spans_per_second")
candidate_rate = candidate.get("accepted_spans_per_second")
throughput_limit = float(thresholds["max_candidate_throughput_regression_pct"])
if base_rate in (None, 0) or candidate_rate is None:
results.append({"threshold": "max_candidate_throughput_regression_pct", "status": "failed", "reason": "missing or zero throughput", "base": base_rate, "candidate": candidate_rate})
else:
actual = (base_rate - candidate_rate) / base_rate * 100.0
results.append({"threshold": "max_candidate_throughput_regression_pct", "status": "passed" if actual <= throughput_limit else "failed", "actual_pct": actual, "limit_pct": throughput_limit, "base": base_rate, "candidate": candidate_rate})
base_latency = base.get("mean_http_latency_ms")
candidate_latency = candidate.get("mean_http_latency_ms")
latency_limit = float(thresholds["max_candidate_mean_latency_regression_pct"])
if base_latency in (None, 0) or candidate_latency is None:
results.append({"threshold": "max_candidate_mean_latency_regression_pct", "status": "failed", "reason": "missing or zero mean latency", "base": base_latency, "candidate": candidate_latency})
else:
actual = (candidate_latency - base_latency) / base_latency * 100.0
results.append({"threshold": "max_candidate_mean_latency_regression_pct", "status": "passed" if actual <= latency_limit else "failed", "actual_pct": actual, "limit_pct": latency_limit, "base": base_latency, "candidate": candidate_latency})
return results
def require_fresh_work_dirs(targets: list[RunTarget], *, reuse_work_dir: bool, dry_run: bool, fixture_only: bool) -> None:
if dry_run or reuse_work_dir or fixture_only:
return
for target in targets:
if target.work_dir.exists() and any(target.work_dir.iterdir()):
raise RuntimeError(f"target work_dir exists and is non-empty: {target.work_dir}; pass --reuse-work-dir to override")
def write_frontend_prom_config(target: RunTarget, remote: dict[str, Any]) -> Path:
prom = remote["prom_store"]
path = target.work_dir / "frontend-prom-store.toml"
path.parent.mkdir(parents=True, exist_ok=True)
path.write_text(
"[prom_store]\n"
"enable = true\n"
"with_metric_engine = true\n"
f"pending_rows_flush_interval = {json.dumps(str(prom['pending_rows_flush_interval']))}\n"
f"max_batch_rows = {int(prom['max_batch_rows'])}\n"
f"max_concurrent_flushes = {int(prom['max_concurrent_flushes'])}\n"
f"worker_channel_capacity = {int(prom['worker_channel_capacity'])}\n"
f"max_inflight_requests = {int(prom['max_inflight_requests'])}\n"
)
return path
def remote_write_command(generator: Path, target: RunTarget, remote: dict[str, Any]) -> list[str]:
metric = remote["metric"]
physical_table = remote["physical_table"]
cmd = [
str(generator),
"prom-remote-write",
"--endpoint", f"http://127.0.0.1:{target.http_port}/v1/prometheus/write",
"--database", remote["database"],
"--metric", metric,
"--physical-table", physical_table,
"--series-count", str(int(remote["series_count"])),
"--samples-per-series", str(int(remote["samples_per_series"])),
"--start-unix-millis", str(int(remote["start_unix_millis"])),
"--step-millis", str(int(remote["step_millis"])),
"--chunk-series-count", str(int(remote["chunk_series_count"])),
"--timeout-seconds", str(int(remote["timeout_seconds"])),
]
value = remote["value"]
cmd.extend(["--value-pattern", str(value["pattern"])])
cmd.extend(["--value-base", str(value["base"])])
cmd.extend(["--value-step", str(value["step"])])
cmd.extend(["--value-cardinality", str(int(value["cardinality"]))])
cmd.extend(["--value-seed", str(int(value["seed"]))])
cmd.extend(["--value-run-length", str(int(value["run_length"]))])
cmd.extend(["--value-stall-every", str(int(value["stall_every"]))])
cmd.extend(["--value-stall-length", str(int(value["stall_length"]))])
cmd.extend(["--value-mixed-every", str(int(value["mixed_every"]))])
if "sample_offset" in remote:
cmd.extend(["--value-sample-offset", str(int(remote["sample_offset"]))])
if "total_samples_per_series" in remote:
cmd.extend(["--value-total-samples-per-series", str(int(remote["total_samples_per_series"]))])
return cmd
def run_remote_write(generator: Path | None, target: RunTarget, remote: dict[str, Any], *, dry_run: bool) -> dict[str, Any]:
if generator is None:
if not dry_run:
raise ValueError("--fixture-generator query_perf_fixture is required for prom_remote_write_then_query")
generator = Path("query_perf_fixture")
cmd = remote_write_command(generator, target, remote)
if dry_run:
return {"status": "dry-run", "cmd": cmd}
result = run_command(cmd)
result["status"] = "ok" if result["returncode"] == 0 else "failed"
if result["returncode"] != 0:
raise RuntimeError(f"remote-write generator failed for {target.name}: {result['stderr'][:2000]}")
try:
result["summary"] = json.loads(result["stdout"])
except json.JSONDecodeError:
result["summary_parse_error"] = result["stdout"]
return result
def run_remote_write_capture(generator: Path | None, target: RunTarget, remote: dict[str, Any], *, dry_run: bool) -> dict[str, Any]:
"""Remote-write helper invocation that records failures instead of raising.
The write_throughput scenario measures a failure rate, so a failed chunk is
kept in the chunk list with status "failed" (and no summary) rather than
aborting the whole run.
"""
if generator is None:
if not dry_run:
raise ValueError("--fixture-generator query_perf_fixture is required for write_throughput")
generator = Path("query_perf_fixture")
cmd = remote_write_command(generator, target, remote)
if dry_run:
return {"status": "dry-run", "cmd": cmd}
result = run_command(cmd)
result["status"] = "ok" if result["returncode"] == 0 else "failed"
if result["returncode"] == 0:
try:
result["summary"] = json.loads(result["stdout"])
except json.JSONDecodeError:
result["summary_parse_error"] = result["stdout"]
return result
def summarize_remote_write_chunks(chunks: list[dict[str, Any]]) -> dict[str, Any]:
summary = {"rows": 0, "samples_written": 0, "batches": 0, "elapsed_seconds": 0.0}
for chunk in chunks:
chunk_summary = chunk.get("summary") or {}
summary["rows"] += int(chunk_summary.get("rows", 0))
summary["samples_written"] += int(chunk_summary.get("samples_written", chunk_summary.get("rows", 0)))
summary["batches"] += int(chunk_summary.get("batches", 0))
summary["elapsed_seconds"] += float(chunk_summary.get("elapsed_seconds", chunk.get("elapsed_seconds", 0.0)))
return summary
def flush_remote_physical_table(target: RunTarget, db: str, physical_table: str, args: argparse.Namespace, *, dry_run: bool, reason: str, chunk_index: int | None = None) -> dict[str, Any]:
if dry_run:
return {"status": "dry-run", "physical_table": physical_table, "reason": reason, "chunk_index": chunk_index}
result = http_post_sql(target.http_port, f"ADMIN FLUSH_TABLE({sql_string(physical_table)})", db, args.http_timeout)
result["physical_table"] = physical_table
result["reason"] = reason
result["chunk_index"] = chunk_index
return result
def run_remote_write_ingestion(generator: Path | None, target: RunTarget, remote: dict[str, Any], args: argparse.Namespace, *, dry_run: bool) -> tuple[dict[str, Any], list[dict[str, Any]]]:
sample_chunk_size = remote["sample_chunk_size"]
db = remote["database"]
physical_table = remote["physical_table"]
if sample_chunk_size is None:
rw = run_remote_write(generator, target, remote, dry_run=dry_run)
flush = flush_remote_physical_table(target, db, physical_table, args, dry_run=dry_run, reason="final")
return rw, [flush]
total_samples = int(remote["samples_per_series"])
chunk_samples = int(sample_chunk_size)
if chunk_samples <= 0:
raise ValueError("scenario.remote_write.sample_chunk_size must be positive")
flush_every = int(remote["flush_every_sample_chunks"])
if flush_every <= 0:
raise ValueError("scenario.remote_write.flush_every_sample_chunks must be positive")
start = int(remote["start_unix_millis"])
step = int(remote["step_millis"])
chunks: list[dict[str, Any]] = []
flushes: list[dict[str, Any]] = []
chunk_index = 0
last_flushed_chunk = 0
for offset in range(0, total_samples, chunk_samples):
current = min(chunk_samples, total_samples - offset)
chunk_index += 1
chunk_remote = dict(remote)
chunk_remote["samples_per_series"] = current
chunk_remote["start_unix_millis"] = start + offset * step
chunk_remote["sample_offset"] = offset
chunk_remote["total_samples_per_series"] = total_samples
result = run_remote_write(generator, target, chunk_remote, dry_run=dry_run)
result["sample_offset"] = offset
result["samples_per_series"] = current
result["chunk_index"] = chunk_index
chunks.append(result)
if chunk_index % flush_every == 0:
flushes.append(flush_remote_physical_table(target, db, physical_table, args, dry_run=dry_run, reason="periodic", chunk_index=chunk_index))
last_flushed_chunk = chunk_index
if last_flushed_chunk != chunk_index:
flushes.append(flush_remote_physical_table(target, db, physical_table, args, dry_run=dry_run, reason="final", chunk_index=chunk_index))
aggregate = summarize_remote_write_chunks(chunks)
if dry_run:
series_count = int(remote["series_count"])
chunk_series_count = int(remote["chunk_series_count"])
aggregate["rows"] = series_count * total_samples
aggregate["samples_written"] = aggregate["rows"]
aggregate["batches"] = chunk_index * ((series_count + chunk_series_count - 1) // chunk_series_count)
return {
"status": "dry-run" if dry_run else ("ok" if all(chunk.get("status") == "ok" for chunk in chunks) else "failed"),
"mode": "sample-chunked",
"sample_chunk_size": chunk_samples,
"flush_every_sample_chunks": flush_every,
"chunks": chunks,
"aggregate": aggregate,
}, flushes
def expected_remote_write_rows(remote: dict[str, Any]) -> int:
return int(remote["series_count"]) * int(remote["samples_per_series"])
def run_write_throughput_ingestion(generator: Path | None, target: RunTarget, remote: dict[str, Any], args: argparse.Namespace, *, dry_run: bool) -> tuple[dict[str, Any], list[dict[str, Any]]]:
"""Sample-chunked remote-write ingestion for the write_throughput scenario.
Mirrors run_remote_write_ingestion but records failed chunks instead of
raising, so the runner can compute a failure rate and per-window stats from
the per-chunk results. Each chunk result keeps its rows/elapsed so the
measurement functions below can bucket RPS over the measurement window.
"""
sample_chunk_size = remote["sample_chunk_size"]
db = remote["database"]
physical_table = remote["physical_table"]
if sample_chunk_size is None:
rw = run_remote_write_capture(generator, target, remote, dry_run=dry_run)
flush = flush_remote_physical_table(target, db, physical_table, args, dry_run=dry_run, reason="final")
return rw, [flush]
total_samples = int(remote["samples_per_series"])
chunk_samples = int(sample_chunk_size)
if chunk_samples <= 0:
raise ValueError("scenario.remote_write.sample_chunk_size must be positive")
flush_every = int(remote["flush_every_sample_chunks"])
if flush_every <= 0:
raise ValueError("scenario.remote_write.flush_every_sample_chunks must be positive")
start = int(remote["start_unix_millis"])
step = int(remote["step_millis"])
chunks: list[dict[str, Any]] = []
flushes: list[dict[str, Any]] = []
chunk_index = 0
last_flushed_chunk = 0
for offset in range(0, total_samples, chunk_samples):
current = min(chunk_samples, total_samples - offset)
chunk_index += 1
chunk_remote = dict(remote)
chunk_remote["samples_per_series"] = current
chunk_remote["start_unix_millis"] = start + offset * step
chunk_remote["sample_offset"] = offset
chunk_remote["total_samples_per_series"] = total_samples
result = run_remote_write_capture(generator, target, chunk_remote, dry_run=dry_run)
result["sample_offset"] = offset
result["samples_per_series"] = current
result["chunk_index"] = chunk_index
chunks.append(result)
if chunk_index % flush_every == 0:
flushes.append(flush_remote_physical_table(target, db, physical_table, args, dry_run=dry_run, reason="periodic", chunk_index=chunk_index))
last_flushed_chunk = chunk_index
if last_flushed_chunk != chunk_index:
flushes.append(flush_remote_physical_table(target, db, physical_table, args, dry_run=dry_run, reason="final", chunk_index=chunk_index))
aggregate = summarize_remote_write_chunks(chunks)
if dry_run:
series_count = int(remote["series_count"])
chunk_series_count = int(remote["chunk_series_count"])
aggregate["rows"] = series_count * total_samples
aggregate["samples_written"] = aggregate["rows"]
aggregate["batches"] = chunk_index * ((series_count + chunk_series_count - 1) // chunk_series_count)
return {
"status": "dry-run" if dry_run else ("ok" if all(chunk.get("status") == "ok" for chunk in chunks) else "failed"),
"mode": "sample-chunked",
"sample_chunk_size": chunk_samples,
"flush_every_sample_chunks": flush_every,
"chunks": chunks,
"aggregate": aggregate,
}, flushes
def write_chunk_rows(chunk: dict[str, Any]) -> int:
summary = chunk.get("summary") or {}
if "rows" in summary:
return int(summary["rows"])
return int(chunk.get("rows", 0))
def write_chunk_elapsed_seconds(chunk: dict[str, Any]) -> float:
summary = chunk.get("summary") or {}
if "elapsed_seconds" in summary:
return float(summary["elapsed_seconds"])
return float(chunk.get("elapsed_seconds", 0.0))
def write_throughput_windows(chunks: list[dict[str, Any]], duration_seconds: int, window_seconds: int) -> list[dict[str, Any]]:
"""Bucket chunk rows into per-window write stats.
Chunks are written sequentially; each chunk spreads its rows uniformly over
its own elapsed time. The measurement spans the first `duration_seconds` of
the combined timeline and is split into consecutive `window_seconds` buckets
(`duration_seconds` must be a multiple of `window_seconds`). Returns windows
ordered by start offset, each with {"index", "start_offset", "rows", "rps"}.
"""
window_count = max(1, duration_seconds // window_seconds)
rows_in_window = [0.0] * window_count
position = 0.0
for chunk in chunks:
elapsed = write_chunk_elapsed_seconds(chunk)
rows = float(write_chunk_rows(chunk))
if elapsed <= 0.0 or rows <= 0.0:
continue
start = position
end = position + elapsed
for idx in range(window_count):
window_start = idx * window_seconds
window_end = window_start + window_seconds
overlap = min(end, window_end) - max(start, window_start)
if overlap > 0.0:
rows_in_window[idx] += rows * (overlap / elapsed)
position = end
if position >= duration_seconds:
break
return [
{"index": idx, "start_offset": idx * window_seconds, "rows": rows_in_window[idx], "rps": rows_in_window[idx] / window_seconds}
for idx in range(window_count)
]
def write_throughput_measurement(rw: dict[str, Any], write_measure: dict[str, Any]) -> dict[str, Any]:
"""Compute write throughput/latency/failure stats from ingestion results.
Rows and elapsed come from each chunk's generator summary; mean RPS is
total rows over total write elapsed. Per-window RPS buckets the chunk rows
over the first `duration_seconds`. p50/p99 write-request latency is computed
from per-chunk generator durations (each chunk is a sequential set of
remote-write requests, so this is a consistent base-vs-candidate proxy).
"""
chunks = rw.get("chunks")
if chunks is None:
# Single non-chunked invocation: the ingestion result is the one chunk.
chunks = [rw]
duration_seconds = int(write_measure["duration_seconds"])
window_seconds = int(write_measure["window_seconds"])
rows = sum(write_chunk_rows(chunk) for chunk in chunks)
elapsed_seconds = sum(write_chunk_elapsed_seconds(chunk) for chunk in chunks)
failed = sum(1 for chunk in chunks if chunk.get("status") != "ok")
total_chunks = len(chunks)
latencies_ms = [
write_chunk_elapsed_seconds(chunk) * 1000.0
for chunk in chunks
if chunk.get("status") == "ok" and write_chunk_elapsed_seconds(chunk) > 0.0
]
return {
"rows": rows,
"elapsed_seconds": elapsed_seconds,
"mean_rps": rows / elapsed_seconds if elapsed_seconds > 0 else None,
"windows": write_throughput_windows(chunks, duration_seconds, window_seconds),
"p50_latency_ms": percentile(latencies_ms, 50) if latencies_ms else None,
"p99_latency_ms": percentile(latencies_ms, 99) if latencies_ms else None,
"failed_chunks": failed,
"total_chunks": total_chunks,
"failure_rate": failed / total_chunks if total_chunks else None,
"target_rps": write_measure.get("target_rps", 0.0),
}
def planned_write_throughput_measurement(remote: dict[str, Any], write_measure: dict[str, Any]) -> dict[str, Any]:
"""Dry-run measurement: no subprocess ran, so report planned values."""
rows = expected_remote_write_rows(remote)
duration_seconds = int(write_measure["duration_seconds"])
window_seconds = int(write_measure["window_seconds"])
return {
"status": "planned",
"planned_rows": rows,
"rows": rows,
"elapsed_seconds": None,
"mean_rps": None,
"windows": [
{"index": idx, "start_offset": idx * window_seconds, "rows": 0.0, "rps": 0.0}
for idx in range(duration_seconds // window_seconds)
],
"p50_latency_ms": None,
"p99_latency_ms": None,
"failed_chunks": 0,
"total_chunks": 0,
"failure_rate": None,
"target_rps": write_measure.get("target_rps", 0.0),
}
def planned_write_throughput_thresholds(write_measure: dict[str, Any]) -> list[dict[str, Any]]:
thresholds = write_measure["thresholds"]
planned = [
{"threshold": "max_failure_rate", "status": "planned", "limit": thresholds["max_failure_rate"]},
{"threshold": "max_mean_rps_regression_pct", "status": "planned", "limit_pct": thresholds["max_mean_rps_regression_pct"]},
{"threshold": "max_p99_latency_regression_pct", "status": "planned", "limit_pct": thresholds["max_p99_latency_regression_pct"]},
]
if thresholds.get("min_rps_absolute") is not None:
planned.append({"threshold": "min_rps_absolute", "status": "planned", "limit": thresholds["min_rps_absolute"]})
return planned
def enforce_write_throughput_thresholds(write_measure: dict[str, Any], base: dict[str, Any], candidate: dict[str, Any]) -> list[dict[str, Any]]:
"""Enforce write_throughput gates.
Per-target: `max_failure_rate` (failed chunk fraction) and optional
`min_rps_absolute` (mean RPS floor). Base-vs-candidate: mean RPS regression
pct `(base - candidate) / base * 100` and p99 latency regression pct
`(candidate - base) / base * 100`; positive actual means a candidate
regression, negative is an improvement.
"""
thresholds = write_measure["thresholds"]
results: list[dict[str, Any]] = []
for target_name, measurement in (("base", base), ("candidate", candidate)):
failure_rate = measurement.get("failure_rate")
failure_limit = float(thresholds["max_failure_rate"])
results.append({"target": target_name, "threshold": "max_failure_rate", "status": "passed" if failure_rate is not None and failure_rate <= failure_limit else "failed", "actual": failure_rate, "limit": failure_limit})
min_rps = thresholds.get("min_rps_absolute")
if min_rps is not None:
mean_rps = measurement.get("mean_rps")
results.append({"target": target_name, "threshold": "min_rps_absolute", "status": "passed" if mean_rps is not None and mean_rps >= float(min_rps) else "failed", "actual": mean_rps, "limit": min_rps})
base_rps = base.get("mean_rps")
candidate_rps = candidate.get("mean_rps")
rps_limit = float(thresholds["max_mean_rps_regression_pct"])
if base_rps in (None, 0) or candidate_rps is None:
results.append({"threshold": "max_mean_rps_regression_pct", "status": "failed", "reason": "missing or zero mean RPS", "base": base_rps, "candidate": candidate_rps})
else:
actual = (base_rps - candidate_rps) / base_rps * 100.0
results.append({"threshold": "max_mean_rps_regression_pct", "status": "passed" if actual <= rps_limit else "failed", "actual_pct": actual, "limit_pct": rps_limit, "base": base_rps, "candidate": candidate_rps})
base_p99 = base.get("p99_latency_ms")
candidate_p99 = candidate.get("p99_latency_ms")
p99_limit = float(thresholds["max_p99_latency_regression_pct"])
if base_p99 in (None, 0) or candidate_p99 is None:
results.append({"threshold": "max_p99_latency_regression_pct", "status": "failed", "reason": "missing or zero p99 latency", "base": base_p99, "candidate": candidate_p99})
else:
actual = (candidate_p99 - base_p99) / base_p99 * 100.0
results.append({"threshold": "max_p99_latency_regression_pct", "status": "passed" if actual <= p99_limit else "failed", "actual_pct": actual, "limit_pct": p99_limit, "base": base_p99, "candidate": candidate_p99})
return results
def mix_query_sql(mix: dict[str, Any], remote: dict[str, Any]) -> str:
"""Resolve the mixed-scenario query SQL.
Defaults to a count(*) over the remote-write physical table, which scans
every ingested row and therefore contends with the write path on the
datanode runtime.
"""
sql = mix.get("query_sql")
if sql:
return str(sql)
return f"SELECT count(*) FROM {sql_ident(remote['physical_table'])}"
def expected_mix_query_attempts(duration_seconds: int, interval_ms: int, parallelism: int) -> int:
"""Per-thread attempt count for a query loop: queries at t=0, interval,
2*interval, ... while t < duration. Wall-clock jitter may drop at most one
trailing query per thread, so the observed count is
[expected - parallelism, expected]."""
per_thread = math.ceil(max(0.0, float(duration_seconds)) * 1000.0 / max(1, int(interval_ms)))
return per_thread * max(1, int(parallelism))
def run_mix_query_loop(target: RunTarget, mix: dict[str, Any], db: str, duration_seconds: int, http_timeout: float, remote: dict[str, Any]) -> list[dict[str, Any]]:
"""Run `query_parallelism` query-loop threads for `duration_seconds`.
Each thread issues the mix query via HTTP SQL every `query_interval_ms`
(per-thread wall clock; if a query itself takes longer than the interval
the thread does not sleep). Returns the collected attempt results in
completion order; every attempt carries `ok`, `status`, `latency_ms`, and
`sql` from `http_post_sql`.
"""
query_sql = mix_query_sql(mix, remote)
interval_ms = int(mix.get("query_interval_ms", 100))
parallelism = max(1, int(mix.get("query_parallelism", 1)))
attempts: list[dict[str, Any]] = []
lock = threading.Lock()
stop = threading.Event()
deadline = time.monotonic() + max(0.0, float(duration_seconds))
threads = []
for _ in range(parallelism):
thread = threading.Thread(
target=_mix_query_worker,
args=(target, query_sql, db, interval_ms, deadline, http_timeout, attempts, lock, stop),
daemon=True,
)
threads.append(thread)
thread.start()
for thread in threads:
# Bounded join: workers may be stuck in a slow HTTP call beyond the
# deadline; their daemon status keeps process exit unblocked.
thread.join(timeout=max(10.0, float(interval_ms) / 1000.0 + 2.0))
return attempts
def _mix_query_worker(
target: RunTarget,
query_sql: str,
db: str,
interval_ms: int,
deadline: float,
http_timeout: float,
attempts: list[dict[str, Any]],
lock: threading.Lock,
stop: threading.Event,
) -> None:
interval_seconds = float(interval_ms) / 1000.0
while not stop.is_set():
started = time.monotonic()
if started >= deadline:
return
result = http_post_sql(target.http_port, query_sql, db, http_timeout)
with lock:
attempts.append(result)
remaining = interval_seconds - (time.monotonic() - started)
if remaining > 0 and stop.wait(remaining):
return
def mix_query_measurement(attempts: list[dict[str, Any]]) -> dict[str, Any]:
"""Aggregate query-loop attempts into {samples, failures, failure_rate,
p50_ms, p99_ms, mean_ms, latency_samples}.
Latency percentiles/mean cover every attempt that recorded a latency
(including failed ones, so timeouts show up in p99). `latency_samples`
counts those attempts. The input is snapshotted so a still-running daemon
query worker cannot mutate it mid-aggregation.
"""
attempts = list(attempts)
total = len(attempts)
failures = sum(1 for a in attempts if not a.get("ok"))
latencies = [float(a["latency_ms"]) for a in attempts if a.get("latency_ms") is not None]
return {
"samples": total,
"failures": failures,
"failure_rate": failures / total if total else None,
"p50_ms": percentile(latencies, 50) if latencies else None,
"p99_ms": percentile(latencies, 99) if latencies else None,
"mean_ms": statistics.mean(latencies) if latencies else None,
"latency_samples": len(latencies),
}
def planned_mix_query_measurement(mix: dict[str, Any], write_measure: dict[str, Any]) -> dict[str, Any]:
"""Dry-run query measurement: no query loop ran, so report planned counts."""
duration_seconds = int(write_measure["duration_seconds"])
interval_ms = int(mix.get("query_interval_ms", 100))
parallelism = int(mix.get("query_parallelism", 1))
return {
"status": "planned",
"planned_attempts": expected_mix_query_attempts(duration_seconds, interval_ms, parallelism),
"query_interval_ms": interval_ms,
"query_parallelism": parallelism,
"samples": 0,
"failures": 0,
"failure_rate": None,
"p50_ms": None,
"p99_ms": None,
"mean_ms": None,
"latency_samples": 0,
}
def run_mixed_ingestion_and_queries(generator: Path | None, target: RunTarget, remote: dict[str, Any], args: argparse.Namespace, mix: dict[str, Any], write_measure: dict[str, Any], *, dry_run: bool) -> tuple[dict[str, Any], list[dict[str, Any]], list[dict[str, Any]]]:
"""Run remote-write ingestion in a background thread while a query loop
hammers the same frontend concurrently for `duration_seconds`, then join.
The ingestion is the same chunked `run_write_throughput_ingestion`
machinery, just executed from a thread; each chunk is its own
`query_perf_fixture prom-remote-write` subprocess (generate+send in one
shot), and the query loop issues HTTP /v1/sql requests from the main
process, so both hit the same frontend/datanode and genuinely contend on
the datanode runtime. Returns (rw, flushes, query_attempts); in dry-run
mode no thread and no query loop is started and attempts is empty.
"""
if dry_run:
rw, flushes = run_write_throughput_ingestion(generator, target, remote, args, dry_run=True)
return rw, flushes, []
results: dict[str, Any] = {}
errors: list[BaseException] = []
def ingest() -> None:
try:
results["rw"], results["flushes"] = run_write_throughput_ingestion(generator, target, remote, args, dry_run=False)
except BaseException as e: # noqa: BLE001 - re-raised in the caller
errors.append(e)
thread = threading.Thread(target=ingest, name=f"{target.name}-ingest", daemon=True)
thread.start()
try:
attempts = run_mix_query_loop(target, mix, remote["database"], int(write_measure["duration_seconds"]), args.http_timeout, remote)
finally:
thread.join()
if thread.is_alive():
# The datanode may be wedged; the ingestion thread is daemon so it
# cannot block process exit, but surface the lack of results.
raise RuntimeError(f"write ingestion thread did not finish for {target.name}")
if errors:
raise errors[0]
return results["rw"], results["flushes"], attempts
def planned_mix_query_thresholds(mix: dict[str, Any]) -> list[dict[str, Any]]:
thresholds = mix["thresholds"]
return [
{"threshold": "max_query_failure_rate", "status": "planned", "limit": thresholds["max_query_failure_rate"]},
{"threshold": "max_query_p99_regression_pct", "status": "planned", "limit_pct": thresholds["max_query_p99_regression_pct"]},
]
def enforce_mix_query_thresholds(mix: dict[str, Any], base: dict[str, Any], candidate: dict[str, Any]) -> list[dict[str, Any]]:
"""Enforce the query-side gates of the mixed read/write scenario.
Per-target: `max_query_failure_rate` (failed query attempts / total).
Base-vs-candidate: `max_query_p99_regression_pct`
`(candidate_p99_ms - base_p99_ms) / base_p99_ms * 100`; positive actual
means a candidate regression, negative is an improvement.
"""
thresholds = mix["thresholds"]
results: list[dict[str, Any]] = []
failure_limit = float(thresholds["max_query_failure_rate"])
for target_name, measurement in (("base", base), ("candidate", candidate)):
failure_rate = measurement.get("failure_rate")
results.append({"target": target_name, "threshold": "max_query_failure_rate", "status": "passed" if failure_rate is not None and failure_rate <= failure_limit else "failed", "actual": failure_rate, "limit": failure_limit})
base_p99 = base.get("p99_ms")
candidate_p99 = candidate.get("p99_ms")
p99_limit = float(thresholds["max_query_p99_regression_pct"])
if base_p99 in (None, 0) or candidate_p99 is None:
results.append({"threshold": "max_query_p99_regression_pct", "status": "failed", "reason": "missing or zero query p99 latency", "base": base_p99, "candidate": candidate_p99})
else:
actual = (candidate_p99 - base_p99) / base_p99 * 100.0
results.append({"threshold": "max_query_p99_regression_pct", "status": "passed" if actual <= p99_limit else "failed", "actual_pct": actual, "limit_pct": p99_limit, "base": base_p99, "candidate": candidate_p99})
return results
def extract_count_value(result: dict[str, Any]) -> int | None:
body = result.get("response")
if not isinstance(body, dict):
return None
data = body.get("data")
if not isinstance(data, list) or not data:
return None
row = data[0]
if not isinstance(row, dict):
return None
for key, value in row.items():
if key.lower() == "count(*)" or key.lower().startswith("count("):
try:
return int(value)
except (TypeError, ValueError):
return None
return None
def poll_expected_count(target: RunTarget, table_name: str, db: str, expected_rows: int, timeout_s: float, http_timeout: float) -> dict[str, Any]:
sql = f"SELECT count(*) FROM {sql_ident(table_name)}"
deadline = time.monotonic() + timeout_s
attempts = 0
last: dict[str, Any] | None = None
while True:
attempts += 1
result = http_post_sql(target.http_port, sql, db, http_timeout)
observed_rows = extract_count_value(result)
result["expected_rows"] = expected_rows
result["observed_rows"] = observed_rows
result["attempts"] = attempts
result["row_count_ok"] = result.get("ok") and observed_rows == expected_rows
if result["row_count_ok"]:
return result
last = result
if time.monotonic() >= deadline:
break
time.sleep(0.5)
assert last is not None
last["ok"] = False
last["row_count_ok"] = False
last["error"] = f"expected {expected_rows} rows but observed {last.get('observed_rows')} after {attempts} attempts"
return last
def run_remote_write_scenario(args: argparse.Namespace, case: dict[str, Any], case_path: Path, targets: list[RunTarget], report: dict[str, Any]) -> None:
remote = scenario(case)["remote_write"]
tables = case_tables(case)
if args.fixture_only:
raise ValueError("--fixture-only is not supported for prom_remote_write_then_query; use --dry-run for planning")
helper = args.fixture_generator or args.remote_write_generator
if helper is not None:
require_binary(helper, "query_perf_fixture", dry_run=args.dry_run)
elif not args.dry_run:
raise ValueError("--fixture-generator is required for prom_remote_write_then_query")
storage = storage_config(remote)
if storage and helper is None and not args.dry_run:
raise ValueError("--fixture-generator is required when scenario.remote_write.storage.inspect is enabled")
clusters: list[DistributedCluster] = []
target_results = []
storage_results = []
try:
for target in targets:
target.work_dir.mkdir(parents=True, exist_ok=True)
config_path = write_frontend_prom_config(target, remote)
cluster = DistributedCluster(target)
clusters.append(cluster)
if not args.dry_run:
cluster.start_all(config_path)
db = remote["database"]
create_database = {"status": "dry-run", "database": db}
if not args.dry_run:
create_database = http_post_sql(target.http_port, f"CREATE DATABASE IF NOT EXISTS {sql_ident(db)}", "public", args.http_timeout)
rw, flushes = run_remote_write_ingestion(helper, target, remote, args, dry_run=args.dry_run)
visibility = {"status": "dry-run"}
query_result = {"validation": [], "validation_errors": [], "measurements": [], "status": "planned"}
if not args.dry_run:
visibility = poll_expected_count(target, tables[0]["name"], db, expected_remote_write_rows(remote), float(remote["visibility_timeout_seconds"]), args.http_timeout)
query_result = run_queries(target, case, tables, args.http_timeout)
cluster.stop_component("datanode")
storage_inspection = None
if storage:
storage_inspection = run_storage_inspection(helper or args.storage_inspector, target, storage, dry_run=args.dry_run)
storage_results.append(storage_inspection)
read_bench_result = run_read_bench(args.candidate_bin, target, remote["read_bench"], storage_inspection, dry_run=args.dry_run)
tr = {"name": target.name, "binary": str(target.binary), "work_dir": str(target.work_dir), "components": cluster.component_report(), "frontend_config": str(config_path), "create_database": create_database, "remote_write": rw, "flushes": flushes, "flush": flushes[-1] if flushes else None, "visibility": visibility, **query_result}
if storage_inspection is not None:
tr["storage_inspection"] = storage_inspection
tr["read_bench"] = read_bench_result
flushes_ok = all(flush.get("ok") for flush in flushes)
storage_ok = storage_inspection is None or args.dry_run or storage_inspection.get("status") == "ok"
read_bench_ok = args.dry_run or read_bench_result.get("status") in ("ok", "skipped")
remote_checks_ok = args.dry_run or (create_database.get("ok") and flushes_ok and visibility.get("ok") and visibility.get("row_count_ok") and storage_ok and read_bench_ok)
if not remote_checks_ok:
tr.setdefault("validation_errors", []).append({"phase": "remote_write_visibility_storage", "create_database_ok": create_database.get("ok"), "flushes_ok": flushes_ok, "visibility_ok": visibility.get("ok"), "row_count_ok": visibility.get("row_count_ok"), "storage_ok": storage_ok, "read_bench_ok": read_bench_ok, "create_database": create_database, "flushes": flushes, "visibility": visibility, "read_bench": read_bench_result})
tr["status"] = "planned" if args.dry_run else ("measured" if remote_checks_ok and query_result["status"] == "ok" else "failed")
write_json(target.report_path, tr)
report["targets"].append(tr)
target_results.append(query_result)
if not args.dry_run:
report["thresholds"] = enforce_thresholds(case, target_results[0], target_results[1]) + enforce_storage_thresholds(storage, storage_results[0] if storage_results else None, storage_results[1] if len(storage_results) > 1 else None)
elif storage:
report["thresholds"] = planned_storage_thresholds(storage)
report["status"] = "planned" if args.dry_run else ("failed" if any(t["status"] == "failed" for t in report["thresholds"]) or any(t.get("status") == "failed" for t in report["targets"]) else "ok")
finally:
for cluster in reversed(clusters):
cluster.stop_all()
def run_write_throughput_scenario(args: argparse.Namespace, case: dict[str, Any], case_path: Path, targets: list[RunTarget], report: dict[str, Any]) -> None:
scenario_config = scenario(case)
remote = scenario_config["remote_write"]
write_measure = scenario_config["write_measure"]
if args.fixture_only:
raise ValueError("--fixture-only is not supported for write_throughput; use --dry-run for planning")
helper = args.fixture_generator or args.remote_write_generator
if helper is not None:
require_binary(helper, "query_perf_fixture", dry_run=args.dry_run)
elif not args.dry_run:
raise ValueError("--fixture-generator is required for write_throughput")
clusters: list[DistributedCluster] = []
measurements: list[dict[str, Any]] = []
query_measurements: list[dict[str, Any]] = []
try:
for target in targets:
target.work_dir.mkdir(parents=True, exist_ok=True)
config_path = write_frontend_prom_config(target, remote)
cluster = DistributedCluster(target)
clusters.append(cluster)
if not args.dry_run:
cluster.start_all(config_path)
db = remote["database"]
create_database: dict[str, Any] = {"status": "dry-run", "database": db}
if not args.dry_run:
create_database = http_post_sql(target.http_port, f"CREATE DATABASE IF NOT EXISTS {sql_ident(db)}", "public", args.http_timeout)
mix = write_measure.get("mix")
scheduler_polls_before: dict[str, Any] | None = None
scheduler_polls_after: dict[str, Any] | None = None
if mix is not None and not args.dry_run:
scheduler_polls_before = scrape_scheduler_polls(target, args.http_timeout)
if mix is not None:
rw, flushes, query_attempts = run_mixed_ingestion_and_queries(helper, target, remote, args, mix, write_measure, dry_run=args.dry_run)
query_measurement = planned_mix_query_measurement(mix, write_measure) if args.dry_run else mix_query_measurement(query_attempts)
if not args.dry_run:
scheduler_polls_after = scrape_scheduler_polls(target, args.http_timeout)
else:
rw, flushes = run_write_throughput_ingestion(helper, target, remote, args, dry_run=args.dry_run)
query_measurement = None
measurement = planned_write_throughput_measurement(remote, write_measure) if args.dry_run else write_throughput_measurement(rw, write_measure)
tr = {
"name": target.name,
"binary": str(target.binary),
"work_dir": str(target.work_dir),
"components": cluster.component_report(),
"frontend_config": str(config_path),
"create_database": create_database,
"remote_write": rw,
"flushes": flushes,
"flush": flushes[-1] if flushes else None,
"scheduler": scheduler_report_entry(target.name, scenario_config.get("scheduler")),
"write_measurement": measurement,
}
if query_measurement is not None:
tr["query_measurement"] = query_measurement
if mix is not None:
tr["mix"] = mix
tr["scheduler_poll_deltas"] = scheduler_poll_deltas(scheduler_polls_after, scheduler_polls_before) if scheduler_polls_after is not None else {"status": "planned"}
flushes_ok = all(flush.get("ok") for flush in flushes)
measurement_ok = measurement.get("mean_rps") not in (None, 0)
query_ok = query_measurement is None or query_measurement.get("samples", 0) > 0
checks_ok = args.dry_run or (create_database.get("ok") and flushes_ok and measurement_ok and query_ok)
if not checks_ok and not args.dry_run:
validation = {"phase": "write_throughput", "create_database_ok": create_database.get("ok"), "flushes_ok": flushes_ok, "measurement_ok": measurement_ok, "mean_rps": measurement.get("mean_rps"), "failure_rate": measurement.get("failure_rate")}
if query_measurement is not None:
validation["query_ok"] = query_ok
validation["query_samples"] = query_measurement.get("samples")
tr.setdefault("validation_errors", []).append(validation)
tr["status"] = "planned" if args.dry_run else ("measured" if checks_ok else "failed")
write_json(target.report_path, tr)
report["targets"].append(tr)
measurements.append(measurement)
if query_measurement is not None:
query_measurements.append(query_measurement)
cluster.stop_all()
if args.dry_run:
report["thresholds"] = planned_write_throughput_thresholds(write_measure)
if write_measure.get("mix") is not None:
report["thresholds"] += planned_mix_query_thresholds(write_measure["mix"])
else:
report["thresholds"] = enforce_write_throughput_thresholds(write_measure, measurements[0], measurements[1])
if write_measure.get("mix") is not None:
report["thresholds"] += enforce_mix_query_thresholds(write_measure["mix"], query_measurements[0], query_measurements[1])
report["status"] = "planned" if args.dry_run else ("failed" if any(t["status"] == "failed" for t in report["thresholds"]) or any(t.get("status") == "failed" for t in report["targets"]) else "ok")
finally:
for cluster in reversed(clusters):
cluster.stop_all()
def run_otlp_trace_load_scenario(args: argparse.Namespace, case: dict[str, Any], targets: list[RunTarget], report: dict[str, Any]) -> None:
load = scenario(case)["load"]
if args.fixture_only:
raise ValueError("--fixture-only is not supported for otlp_trace_load; use --dry-run for planning")
if args.otelgen_bin is not None:
require_binary(args.otelgen_bin, "otelgen", dry_run=args.dry_run)
elif not args.dry_run:
raise ValueError("--otelgen-bin is required for otlp_trace_load")
clusters: list[DistributedCluster] = []
metrics_results: list[dict[str, Any]] = []
try:
for target in targets:
target.work_dir.mkdir(parents=True, exist_ok=True)
cluster = DistributedCluster(target)
clusters.append(cluster)
if not args.dry_run:
cluster.start_all()
create_database: dict[str, Any] = {"status": "dry-run", "database": load["database"]}
if not args.dry_run:
create_database = http_post_sql(target.http_port, f"CREATE DATABASE IF NOT EXISTS {sql_ident(load['database'])}", "public", args.http_timeout)
otelgen = run_otelgen_load(args.otelgen_bin, target, load, args.http_timeout, dry_run=args.dry_run)
metrics: dict[str, Any] = {"status": "dry-run"}
flush: dict[str, Any] = {"status": "dry-run", "table": load["table"]}
visibility: dict[str, Any] = {"status": "dry-run"}
if not args.dry_run:
metrics = summarize_otlp_metrics(otelgen)
flush = http_post_sql(target.http_port, f"ADMIN FLUSH_TABLE({sql_string(load['table'])})", load["database"], args.http_timeout)
visibility = poll_expected_count(target, load["table"], load["database"], int(metrics["accepted_spans"]), float(load["visibility_timeout_seconds"]), args.http_timeout)
tr = {
"name": target.name,
"binary": str(target.binary),
"work_dir": str(target.work_dir),
"components": cluster.component_report(),
"create_database": create_database,
"otelgen": otelgen,
"metrics": metrics,
"flush": flush,
"visibility": visibility,
}
checks_ok = args.dry_run or (
create_database.get("ok")
and otelgen.get("status") == "ok"
and not metrics.get("missing_metrics")
and metrics.get("accepted_spans", 0) > 0
and metrics.get("http_requests", 0) > 0
and flush.get("ok")
and visibility.get("ok")
and visibility.get("row_count_ok")
)
tr["status"] = "planned" if args.dry_run else ("measured" if checks_ok else "failed")
write_json(target.report_path, tr)
report["targets"].append(tr)
metrics_results.append(metrics)
cluster.stop_all()
report["thresholds"] = planned_otlp_thresholds(load) if args.dry_run else enforce_otlp_thresholds(load, metrics_results[0], metrics_results[1])
report["status"] = "planned" if args.dry_run else ("failed" if any(t["status"] == "failed" for t in report["thresholds"]) or any(t["status"] == "failed" for t in report["targets"]) else "ok")
finally:
for cluster in reversed(clusters):
cluster.stop_all()
def main() -> int:
args = parse_args()
case_path = args.case.resolve()
if args.fixture_generator is not None:
require_binary(args.fixture_generator, "fixture generator", dry_run=False)
case = load_normalized_case(case_path, args.fixture_generator)
scenario_config = scenario(case)
scenario_kind = scenario_config.get("kind")
tables = case_tables(case)
work_root = args.work_dir.resolve()
work_root.mkdir(parents=True, exist_ok=True)
require_binary(args.base_bin, "base", dry_run=args.dry_run or args.fixture_only)
require_binary(args.candidate_bin, "candidate", dry_run=args.dry_run or args.fixture_only)
if scenario_kind == "direct_readable_sst" and not args.dry_run and args.fixture_generator is None:
raise ValueError("--fixture-generator is required unless --dry-run is set")
ports = allocate_ports(16)
scheduler_cfg = scenario_config.get("scheduler")
targets = [
make_target("base", args.base_bin.resolve(), work_root, ports[:8], scheduler_env=scheduler_env(False, scheduler_cfg)),
make_target("candidate", args.candidate_bin.resolve(), work_root, ports[8:], scheduler_env=scheduler_env(True, scheduler_cfg)),
]
require_fresh_work_dirs(targets, reuse_work_dir=args.reuse_work_dir, dry_run=args.dry_run, fixture_only=args.fixture_only)
fixture_dir = fixture_root(work_root, case_path, case, args.fixture_cache_dir)
reuse_fixture = args.reuse_fixture or args.fixture_cache_dir is not None
report: dict[str, Any] = {"case_path": str(case_path), "case": case.get("case", {}), "scenario": scenario_config, "queries": planned_queries(case), "dry_run": args.dry_run, "fixture_only": args.fixture_only, "query_mode": "fixture-only" if args.fixture_only else "distributed", "reuse_work_dir": args.reuse_work_dir, "reuse_fixture": reuse_fixture, "fixture_cache_dir": str(args.fixture_cache_dir.resolve()) if args.fixture_cache_dir is not None else None, "fixture_dir": str(fixture_dir), "http_timeout": args.http_timeout, "targets": [], "thresholds": [], "status": "planned" if args.dry_run else "running"}
if scenario_kind == "otlp_trace_load":
try:
run_otlp_trace_load_scenario(args, case, targets, report)
except Exception as e: # noqa: BLE001 - write machine-readable failure report
report["status"] = "failed"
report["error"] = repr(e)
write_json(work_root / "query-regression-report.json", report)
output_report(report, args.output)
return 1 if report["status"] == "failed" else 0
if scenario_kind == "prom_remote_write_then_query":
try:
run_remote_write_scenario(args, case, case_path, targets, report)
except Exception as e: # noqa: BLE001 - write machine-readable failure report
report["status"] = "failed"
report["error"] = repr(e)
write_json(work_root / "query-regression-report.json", report)
output_report(report, args.output)
return 1 if report["status"] == "failed" else 0
if scenario_kind == "write_throughput":
try:
run_write_throughput_scenario(args, case, case_path, targets, report)
except Exception as e: # noqa: BLE001 - write machine-readable failure report
report["status"] = "failed"
report["error"] = repr(e)
write_json(work_root / "query-regression-report.json", report)
output_report(report, args.output)
return 1 if report["status"] == "failed" else 0
if args.fixture_only or args.dry_run:
generations = []
for table_idx, table in enumerate(tables):
per_table_fixture_dir = table_fixture_dir(fixture_dir, tables, table, table_idx)
generations.append(generate_fixture(args.fixture_generator, case_path, per_table_fixture_dir, dry_run=args.dry_run, reuse_fixture=reuse_fixture, allow_large_fixture=args.allow_large_fixture, table_name=table["name"] if len(tables) > 1 else None, database=table["database"] if len(tables) > 1 else None))
report["fixture_generation"] = generations[0] if len(generations) == 1 else generations
for target in targets:
target.work_dir.mkdir(parents=True, exist_ok=True)
mats = []
for table_idx, table in enumerate(tables):
per_table_fixture_dir = table_fixture_dir(fixture_dir, tables, table, table_idx)
mats.append(materialize_fixture(target, dry_run=args.dry_run, preserve_state=False, fixture_dir=per_table_fixture_dir, reset_data=not mats))
tr = {"name": target.name, "binary": str(target.binary), "work_dir": str(target.work_dir), "data_dir": str(target.data_dir), "fixture_dir": str(fixture_dir), "fixture_materialization": mats[0] if len(mats) == 1 else mats, "measurements": [], "status": "planned" if args.dry_run else "fixture-ready"}
write_json(target.report_path, tr)
report["targets"].append(tr)
report["status"] = "planned" if args.dry_run else "fixture-ready"
write_json(work_root / "query-regression-report.json", report)
output_report(report, args.output)
return 0
clusters: list[DistributedCluster] = []
try:
discovered = []
for target in targets:
target.work_dir.mkdir(parents=True, exist_ok=True)
cluster = DistributedCluster(target)
clusters.append(cluster)
cluster.start_all()
create_results = []
metas = []
for table in tables:
create_result = http_post_sql(target.http_port, create_table_sql(table), table["database"], args.http_timeout)
if not create_result["ok"]:
raise RuntimeError(f"CREATE TABLE {table['name']} failed for {target.name}: {create_result}")
create_results.append({"table": table["name"], "result": create_result})
metas.append({"table": table["name"], **discover_region_via_frontend(target, table, args.http_timeout)})
cluster.stop_component("datanode")
cluster._ensure_metasrv_alive()
discovered.append(metas)
report["targets"].append({"name": target.name, "binary": str(target.binary), "work_dir": str(target.work_dir), "data_dir": str(target.data_dir), "datanode_data_home": str(target.datanode_data_dir), "components": cluster.component_report(), "create_table": create_results[0]["result"] if len(create_results) == 1 else create_results, "discovered": metas[0] if len(metas) == 1 else metas, "discovered_tables": metas})
if len(discovered[0]) != len(discovered[1]) or any(a["table"] != b["table"] or a["region_id"] != b["region_id"] or a["table_dir"] != b["table_dir"] for a, b in zip(discovered[0], discovered[1])):
raise RuntimeError(f"base/candidate metadata mismatch: {discovered}")
generations = []
for table_idx, meta in enumerate(discovered[0]):
per_table_fixture_dir = table_fixture_dir(fixture_dir, tables, tables[table_idx], table_idx)
generations.append(generate_fixture(args.fixture_generator, case_path, per_table_fixture_dir, dry_run=False, reuse_fixture=reuse_fixture, allow_large_fixture=args.allow_large_fixture, table_name=meta["table"] if len(tables) > 1 else None, database=meta["schema"] if len(tables) > 1 else None, region_id=meta["region_id"], table_dir=meta["table_dir"], region_dir=meta["region_dir"]))
report["fixture_generation"] = generations[0] if len(generations) == 1 else generations
target_results = []
for idx, target in enumerate(targets):
cluster = clusters[idx]
materialize = []
for table_idx, meta in enumerate(discovered[idx]):
per_table_fixture_dir = table_fixture_dir(fixture_dir, tables, tables[table_idx], table_idx)
materialize.append(materialize_fixture(target, dry_run=False, preserve_state=True, expected_region_dir=meta["region_dir"], fixture_dir=per_table_fixture_dir))
cluster.start_datanode()
query_result = run_queries(target, case, tables, args.http_timeout)
if query_result["status"] == "failed":
cluster.restart_frontend()
retry_result = run_queries(target, case, tables, args.http_timeout)
query_result["frontend_restart_retry"] = retry_result
if retry_result["status"] == "ok":
query_result = retry_result
report["targets"][idx]["fixture_dir"] = str(fixture_dir)
report["targets"][idx]["fixture_materialization"] = materialize[0] if len(materialize) == 1 else materialize
report["targets"][idx].update(query_result)
report["targets"][idx]["status"] = "measured" if query_result["status"] == "ok" else "failed"
write_json(target.report_path, report["targets"][idx])
target_results.append(query_result)
report["thresholds"] = enforce_thresholds(case, target_results[0], target_results[1])
report["status"] = "failed" if any(t["status"] == "failed" for t in report["thresholds"]) or any(r["status"] == "failed" for r in target_results) else "ok"
except Exception as e: # noqa: BLE001 - write machine-readable failure report
report["status"] = "failed"
report["error"] = repr(e)
finally:
for cluster in reversed(clusters):
cluster.stop_all()
write_json(work_root / "query-regression-report.json", report)
output_report(report, args.output)
return 1 if report["status"] == "failed" else 0
if __name__ == "__main__":
sys.exit(main())