From 0d663063c480710dc03beb79481f3585089f43f4 Mon Sep 17 00:00:00 2001 From: Adam Gutglick Date: Thu, 30 Jul 2026 15:51:42 +0100 Subject: [PATCH] Use RDS for baseline Signed-off-by: Adam Gutglick --- .github/workflows/bench-pr.yml | 36 +- .github/workflows/sql-benchmarks.yml | 33 +- scripts/compare-benchmark-jsons.py | 235 ++++++++++++- scripts/fetch-benchmark-baseline.py | 311 ++++++++++++++++++ scripts/tests/test_benchmark_reporting.py | 186 +++++++++++ .../tests/test_fetch_benchmark_baseline.py | 215 ++++++++++++ 6 files changed, 995 insertions(+), 21 deletions(-) create mode 100644 scripts/fetch-benchmark-baseline.py create mode 100644 scripts/tests/test_fetch_benchmark_baseline.py diff --git a/.github/workflows/bench-pr.yml b/.github/workflows/bench-pr.yml index a92704e6783..5e31d2881e3 100644 --- a/.github/workflows/bench-pr.yml +++ b/.github/workflows/bench-pr.yml @@ -94,7 +94,7 @@ jobs: VORTEX_EXPERIMENTAL_PATCHED_ARRAY: "1" FLAT_LAYOUT_INLINE_ARRAY_NODE: "1" run: | - python3 scripts/random-access-split.py + python3 scripts/random-access-split.py --emit-ingest-records - name: Run ${{ matrix.benchmark.name }} benchmark if: matrix.benchmark.id != 'random-access-bench' @@ -104,14 +104,17 @@ jobs: VORTEX_EXPERIMENTAL_PATCHED_ARRAY: "1" FLAT_LAYOUT_INLINE_ARRAY_NODE: "1" run: | - bash scripts/bench-taskset.sh target/release_debug/${{ matrix.benchmark.id }} -d gh-json -o results.json + bash scripts/bench-taskset.sh target/release_debug/${{ matrix.benchmark.id }} \ + -d gh-json \ + -o results.json \ + --ingest-jsonl results.ingest.jsonl - - name: Setup AWS CLI + - name: Configure AWS credentials for RDS baseline (OIDC) if: github.event.pull_request.head.repo.fork == false uses: aws-actions/configure-aws-credentials@e6de054238d6b7531b4efff3b6587d9aade6a06c # v6 with: - role-to-assume: arn:aws:iam::245040174862:role/GitHubBenchmarkRole - aws-region: us-east-1 + role-to-assume: ${{ vars.GH_BENCH_INGEST_ROLE_ARN }} + aws-region: ${{ vars.RDS_BENCH_REGION }} - name: Install uv uses: spiraldb/actions/.github/actions/setup-uv@a746510eafaa926484c354541cfc49b2ec06cc63 # 0.18.6 @@ -119,15 +122,28 @@ jobs: sync: false - name: Compare results + if: github.event.pull_request.head.repo.fork == false shell: bash + env: + RDS_BENCH_INSTANCE_ENDPOINT: ${{ vars.RDS_BENCH_INSTANCE_ENDPOINT }} + RDS_BENCH_DB_NAME: ${{ vars.RDS_BENCH_DB_NAME }} + AWS_REGION: ${{ vars.RDS_BENCH_REGION }} run: | set -Eeu -o pipefail -x - python3 scripts/s3-download.py s3://vortex-ci-benchmark-results/data.json.gz data.json.gz --no-sign-request - gzip -d -c data.json.gz > base.json - - uv run --no-project scripts/compare-benchmark-jsons.py base.json results.json "${{ matrix.benchmark.name }}" \ - > comment.md + curl -fsSL https://truststore.pki.rds.amazonaws.com/global/global-bundle.pem \ + -o "${RUNNER_TEMP}/rds-global-bundle.pem" + DSN="postgresql://bench_ingest@${RDS_BENCH_INSTANCE_ENDPOINT}:5432/${RDS_BENCH_DB_NAME}?sslmode=verify-full&sslrootcert=${RUNNER_TEMP}/rds-global-bundle.pem" + uv run --no-project scripts/fetch-benchmark-baseline.py results.ingest.jsonl \ + --postgres "${DSN}" \ + --region "${AWS_REGION}" \ + --output base.ingest.jsonl + uv run --no-project scripts/compare-benchmark-jsons.py \ + --ingest-jsonl \ + --metadata results.json \ + base.ingest.jsonl \ + results.ingest.jsonl \ + "${{ matrix.benchmark.name }}" > comment.md cat comment.md >> "$GITHUB_STEP_SUMMARY" - name: Comment PR diff --git a/.github/workflows/sql-benchmarks.yml b/.github/workflows/sql-benchmarks.yml index b10d7704f69..4c3c4423692 100644 --- a/.github/workflows/sql-benchmarks.yml +++ b/.github/workflows/sql-benchmarks.yml @@ -191,7 +191,7 @@ jobs: ${{ matrix.scale_factor && format('--opt scale-factor={0}', matrix.scale_factor) || '' }} - name: Capture file sizes - if: matrix.remote_key == null + if: inputs.mode == 'develop' && matrix.remote_key == null shell: bash run: | uv run --no-project scripts/capture-file-sizes.py \ @@ -201,17 +201,36 @@ jobs: -o sizes.json cat sizes.json >> results.json + - name: Configure AWS credentials for RDS baseline (OIDC) + if: inputs.mode == 'pr' && github.event.pull_request.head.repo.fork == false + uses: aws-actions/configure-aws-credentials@e6de054238d6b7531b4efff3b6587d9aade6a06c # v6 + with: + role-to-assume: ${{ vars.GH_BENCH_INGEST_ROLE_ARN }} + aws-region: ${{ vars.RDS_BENCH_REGION }} + - name: Compare results - if: inputs.mode == 'pr' + if: inputs.mode == 'pr' && github.event.pull_request.head.repo.fork == false shell: bash + env: + RDS_BENCH_INSTANCE_ENDPOINT: ${{ vars.RDS_BENCH_INSTANCE_ENDPOINT }} + RDS_BENCH_DB_NAME: ${{ vars.RDS_BENCH_DB_NAME }} + AWS_REGION: ${{ vars.RDS_BENCH_REGION }} run: | set -Eeu -o pipefail -x - python3 scripts/s3-download.py s3://vortex-ci-benchmark-results/data.json.gz data.json.gz --no-sign-request - gzip -d -c data.json.gz > base.json - - uv run --no-project scripts/compare-benchmark-jsons.py base.json results.json "${{ matrix.name }}" \ - > comment.md + curl -fsSL https://truststore.pki.rds.amazonaws.com/global/global-bundle.pem \ + -o "${RUNNER_TEMP}/rds-global-bundle.pem" + DSN="postgresql://bench_ingest@${RDS_BENCH_INSTANCE_ENDPOINT}:5432/${RDS_BENCH_DB_NAME}?sslmode=verify-full&sslrootcert=${RUNNER_TEMP}/rds-global-bundle.pem" + uv run --no-project scripts/fetch-benchmark-baseline.py results.ingest.jsonl \ + --postgres "${DSN}" \ + --region "${AWS_REGION}" \ + --output base.ingest.jsonl + uv run --no-project scripts/compare-benchmark-jsons.py \ + --ingest-jsonl \ + --metadata results.json \ + base.ingest.jsonl \ + results.ingest.jsonl \ + "${{ matrix.name }}" > comment.md cat comment.md >> "$GITHUB_STEP_SUMMARY" - name: Comment PR diff --git a/scripts/compare-benchmark-jsons.py b/scripts/compare-benchmark-jsons.py index d7cd34793db..2323ae62544 100644 --- a/scripts/compare-benchmark-jsons.py +++ b/scripts/compare-benchmark-jsons.py @@ -11,12 +11,13 @@ # SPDX-License-Identifier: Apache-2.0 # SPDX-FileCopyrightText: Copyright the Vortex contributors +import argparse import math import os import re -import sys from dataclasses import dataclass from io import StringIO +from pathlib import Path from typing import Any import numpy as np @@ -43,6 +44,203 @@ Z_SCORE_99 = 2.5758293035489004 CONTROL_FORMAT = "parquet" FILE_SIZE_METRIC = "file_size" +VORTEX_FORMAT = "vortex-file-compressed" + + +def read_ingest_records(path: str | Path) -> list[dict[str, Any]]: + """Read canonical v3 benchmark records from JSONL.""" + + records = [] + with open(path, encoding="utf-8") as lines: + for line_no, line in enumerate(lines, start=1): + if not line.strip(): + continue + record = orjson.loads(line) + if not isinstance(record, dict): + raise ValueError(f"{path}:{line_no}: expected a JSON object") + records.append(record) + if not records: + raise ValueError(f"{path}: no benchmark records") + return records + + +def _dataset_label(record: dict[str, Any]) -> str: + dataset = str(record["dataset"]) + variant = record.get("dataset_variant") + return dataset if variant is None else f"{dataset}/{variant}" + + +def _comparison_row( + record: dict[str, Any], + *, + name: str, + unit: str, + value_field: str, +) -> dict[str, Any]: + row = { + "name": name, + "unit": unit, + "value": record[value_field], + "commit_id": record["commit_sha"], + } + if "all_runtimes_ns" in record: + row["all_runtimes"] = record["all_runtimes_ns"] + return row + + +def _compression_time_name(record: dict[str, Any]) -> str: + operation = {"encode": "compress", "decode": "decompress"}.get(record["op"], str(record["op"])) + prefix = { + VORTEX_FORMAT: "", + "parquet": "parquet_rs-zstd ", + "lance": "lance ", + }.get(record["format"], f"{record['format']} ") + return f"{prefix}{operation} time/{_dataset_label(record)}" + + +def _compression_ratio_rows(records: list[dict[str, Any]]) -> list[dict[str, Any]]: + """Derive the compression ratios intentionally omitted from the v3 facts.""" + + time_values: dict[tuple[str, str | None, str, str, str], float] = {} + size_values: dict[tuple[str, str | None, str, str], float] = {} + for record in records: + kind = record.get("kind") + if kind == "compression_time": + key = ( + str(record["dataset"]), + record.get("dataset_variant"), + str(record["op"]), + str(record["format"]), + str(record["commit_sha"]), + ) + time_values[key] = float(record["value_ns"]) + elif kind == "compression_size": + key = ( + str(record["dataset"]), + record.get("dataset_variant"), + str(record["format"]), + str(record["commit_sha"]), + ) + size_values[key] = float(record["value_bytes"]) + + rows = [] + groups = sorted( + {(dataset, variant, commit) for dataset, variant, _format, commit in size_values}, + key=lambda group: tuple("" if value is None else str(value) for value in group), + ) + for dataset, variant, commit in groups: + label = dataset if variant is None else f"{dataset}/{variant}" + vortex_size = size_values.get((dataset, variant, VORTEX_FORMAT, commit)) + for other_format, comparison_name in (("parquet", "parquet-zstd"), ("lance", "lance")): + other_size = size_values.get((dataset, variant, other_format, commit)) + if vortex_size is not None and other_size not in (None, 0): + rows.append( + { + "name": f"vortex:{comparison_name} size/{label}", + "unit": "ratio", + "value": vortex_size / other_size, + "commit_id": commit, + } + ) + + time_groups = sorted( + {(dataset, variant, operation, commit) for dataset, variant, operation, _format, commit in time_values}, + key=lambda group: tuple("" if value is None else str(value) for value in group), + ) + for dataset, variant, operation, commit in time_groups: + label = dataset if variant is None else f"{dataset}/{variant}" + vortex_time = time_values.get((dataset, variant, operation, VORTEX_FORMAT, commit)) + operation_name = {"encode": "compress", "decode": "decompress"}.get(operation, operation) + for other_format, comparison_name in (("parquet", "parquet-zstd"), ("lance", "lance")): + other_time = time_values.get((dataset, variant, operation, other_format, commit)) + if vortex_time is not None and other_time not in (None, 0): + rows.append( + { + "name": f"vortex:{comparison_name} ratio {operation_name} time/{label}", + "unit": "ratio", + "value": vortex_time / other_time, + "commit_id": commit, + } + ) + return rows + + +def normalize_ingest_records(records: list[dict[str, Any]]) -> pd.DataFrame: + """Convert v3 fact records into the comparison reporter's common row shape.""" + + rows: list[dict[str, Any]] = [] + known_kinds = { + "query_measurement", + "compression_time", + "compression_size", + "random_access_time", + "vector_search_run", + } + for record in records: + kind = record.get("kind") + if kind not in known_kinds: + raise ValueError(f"unknown benchmark record kind: {kind!r}") + + if kind == "query_measurement": + row = _comparison_row( + record, + name=(f"{record['dataset']}_q{int(record['query_idx']):02}/{record['engine']}:{record['format']}"), + unit="ns", + value_field="value_ns", + ) + row["storage"] = record["storage"] + row["dataset"] = { + "dataset": record["dataset"], + "dataset_variant": record.get("dataset_variant"), + "scale_factor": record.get("scale_factor"), + } + rows.append(row) + elif kind == "compression_time": + rows.append( + _comparison_row( + record, + name=_compression_time_name(record), + unit="ns", + value_field="value_ns", + ) + ) + elif kind == "compression_size": + rows.append( + _comparison_row( + record, + name=f"{record['format']} size/{_dataset_label(record)}", + unit="bytes", + value_field="value_bytes", + ) + ) + elif kind == "random_access_time": + format_extension = { + VORTEX_FORMAT: "vortex", + "vortex-compact": "vortex", + }.get(record["format"], record["format"]) + dataset_path = "" if record["dataset"] == "taxi" else f"{record['dataset']}/" + row = _comparison_row( + record, + name=f"random-access/{dataset_path}{format_extension}-tokio-local-disk", + unit="ns", + value_field="value_ns", + ) + row["storage"] = "nvme" + rows.append(row) + else: + rows.append( + _comparison_row( + record, + name=( + f"vector-search/{record['dataset']}/{record['layout']}/{record['flavor']}/{record['threshold']}" + ), + unit="ns", + value_field="value_ns", + ) + ) + + rows.extend(_compression_ratio_rows(records)) + return pd.DataFrame(rows) @dataclass @@ -839,14 +1037,43 @@ def group_sort_key(group_key: tuple[str, str]) -> tuple[int, int, str, str]: ) +def parse_args() -> argparse.Namespace: + parser = argparse.ArgumentParser(description="Render a benchmark comparison as Markdown.") + parser.add_argument("base", help="Baseline JSONL path.") + parser.add_argument("pr", help="Pull-request JSONL path.") + parser.add_argument("benchmark_name", nargs="?", default="", help="Benchmark display name.") + parser.add_argument( + "--ingest-jsonl", + action="store_true", + help="Read both inputs as canonical v3 ingest records instead of legacy gh-json rows.", + ) + parser.add_argument( + "--metadata", + help="Optional legacy gh-json path used only for display metadata such as the documentation link.", + ) + return parser.parse_args() + + def main() -> None: """Render the benchmark comparison markdown used in CI PR comments.""" - benchmark_name = sys.argv[3] if len(sys.argv) > 3 else "" + args = parse_args() + benchmark_name = args.benchmark_name + + if args.ingest_jsonl: + base = normalize_ingest_records(read_ingest_records(args.base)) + pr = normalize_ingest_records(read_ingest_records(args.pr)) + if args.metadata is not None: + metadata = pd.read_json(args.metadata, lines=True) + if "doc" in metadata.columns: + docs = metadata["doc"].dropna().unique() + if len(docs) > 0: + pr["doc"] = docs[0] + else: + pr = pd.read_json(args.pr, lines=True) + base = read_latest_baseline_rows(args.base, pr) - pr = pd.read_json(sys.argv[2], lines=True) title = format_title(benchmark_name, pr) - base = read_latest_baseline_rows(sys.argv[1], pr) base_commit_id = set(base["commit_id"].unique()) pr_commit_id = set(pr["commit_id"].unique()) diff --git a/scripts/fetch-benchmark-baseline.py b/scripts/fetch-benchmark-baseline.py new file mode 100644 index 00000000000..5ef79762778 --- /dev/null +++ b/scripts/fetch-benchmark-baseline.py @@ -0,0 +1,311 @@ +#!/usr/bin/env python3 +# /// script +# requires-python = ">=3.11" +# dependencies = [ +# "boto3", +# "psycopg[binary]", +# ] +# /// + +# SPDX-License-Identifier: Apache-2.0 +# SPDX-FileCopyrightText: Copyright the Vortex contributors + +"""Fetch the latest matching benchmark baseline from RDS as v3 JSONL.""" + +from __future__ import annotations + +import argparse +import importlib.util +import json +import os +import tempfile +from collections.abc import Mapping, Sequence +from pathlib import Path +from typing import Any + +_KIND_FAMILY = { + "query_measurement": "query", + "compression_time": "compression", + "compression_size": "compression", + "random_access_time": "random_access", + "vector_search_run": "vector_search", +} + +_FAMILY_SCOPE = { + "query": ("dataset", "dataset_variant", "scale_factor", "storage"), + "compression": ("dataset", "dataset_variant"), + "random_access": ("dataset",), + "vector_search": ("dataset", "layout", "threshold"), +} + +_FAMILY_CANDIDATE = { + "query": ("query_measurements", "q"), + "compression": ("compression_times", "t"), + "random_access": ("random_access_times", "r"), + "vector_search": ("vector_search_runs", "v"), +} + +_SELECTS = { + "query": ( + ( + "query_measurements", + "q", + """ + SELECT 'query_measurement' AS kind, + q.commit_sha, q.dataset, q.dataset_variant, q.scale_factor, + q.query_idx, q.storage, q.engine, q.format, + q.value_ns, q.all_runtimes_ns, + q.peak_physical, q.peak_virtual, + q.physical_delta, q.virtual_delta, q.env_triple + FROM query_measurements q + """, + ("dataset", "dataset_variant", "scale_factor", "storage"), + ("q.query_idx", "q.engine", "q.format"), + ), + ), + "compression": ( + ( + "compression_times", + "t", + """ + SELECT 'compression_time' AS kind, + t.commit_sha, t.dataset, t.dataset_variant, t.format, t.op, + t.value_ns, t.all_runtimes_ns, t.env_triple + FROM compression_times t + """, + ("dataset", "dataset_variant"), + ("t.dataset", "t.dataset_variant", "t.format", "t.op"), + ), + ( + "compression_sizes", + "s", + """ + SELECT 'compression_size' AS kind, + s.commit_sha, s.dataset, s.dataset_variant, s.format, + s.value_bytes + FROM compression_sizes s + """, + ("dataset", "dataset_variant"), + ("s.dataset", "s.dataset_variant", "s.format"), + ), + ), + "random_access": ( + ( + "random_access_times", + "r", + """ + SELECT 'random_access_time' AS kind, + r.commit_sha, r.dataset, r.format, + r.value_ns, r.all_runtimes_ns, r.env_triple + FROM random_access_times r + """, + ("dataset",), + ("r.dataset", "r.format"), + ), + ), + "vector_search": ( + ( + "vector_search_runs", + "v", + """ + SELECT 'vector_search_run' AS kind, + v.commit_sha, v.dataset, v.layout, v.flavor, v.threshold, + v.value_ns, v.all_runtimes_ns, + v.matches, v.rows_scanned, v.bytes_scanned, + v.iterations, v.env_triple + FROM vector_search_runs v + """, + ("dataset", "layout", "threshold"), + ("v.dataset", "v.layout", "v.threshold", "v.flavor"), + ), + ), +} + + +def parse_args() -> argparse.Namespace: + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("pr_jsonl", type=Path, help="PR v3 ingest JSONL used to identify the benchmark scope.") + parser.add_argument("--postgres", metavar="DSN", required=True, help="RDS Postgres libpq DSN.") + parser.add_argument("--region", default=None, help="AWS region used to mint the RDS IAM token.") + parser.add_argument("--output", type=Path, required=True, help="Destination v3 baseline JSONL.") + return parser.parse_args() + + +def read_records(path: Path) -> list[dict]: + """Read non-empty JSON objects from a v3 JSONL file.""" + + records = [] + with path.open(encoding="utf-8") as lines: + for line_no, line in enumerate(lines, start=1): + if not line.strip(): + continue + try: + record = json.loads(line) + except json.JSONDecodeError as exc: + raise ValueError(f"{path}:{line_no}: invalid JSON: {exc}") from exc + if not isinstance(record, dict): + raise ValueError(f"{path}:{line_no}: expected a JSON object") + records.append(record) + if not records: + raise ValueError(f"{path}: no benchmark records") + return records + + +def benchmark_family(records: Sequence[Mapping[str, Any]]) -> str: + """Return the single fact-table family represented by ``records``.""" + + unknown = sorted({record.get("kind") for record in records if record.get("kind") not in _KIND_FAMILY}) + if unknown: + raise ValueError(f"unknown benchmark record kinds: {unknown}") + + families = {_KIND_FAMILY[str(record["kind"])] for record in records} + if len(families) != 1: + raise ValueError(f"PR records contain multiple benchmark families: {sorted(families)}") + return next(iter(families)) + + +def benchmark_scopes(records: Sequence[Mapping[str, Any]], family: str) -> list[tuple[Any, ...]]: + """Return sorted, unique database scopes for a PR benchmark family.""" + + columns = _FAMILY_SCOPE[family] + scopes = {tuple(record.get(column) for column in columns) for record in records} + return sorted(scopes, key=lambda scope: tuple("" if value is None else str(value) for value in scope)) + + +def _scope_predicate( + alias: str, + columns: Sequence[str], + scopes: Sequence[tuple[Any, ...]], +) -> tuple[str, tuple[Any, ...]]: + clauses = [] + params: list[Any] = [] + for scope in scopes: + terms = [] + for column, value in zip(columns, scope, strict=True): + if value is None: + terms.append(f"{alias}.{column} IS NULL") + else: + terms.append(f"{alias}.{column} = %s") + params.append(value) + clauses.append(f"({' AND '.join(terms)})") + return f"({' OR '.join(clauses)})", tuple(params) + + +def _column_name(description: Any) -> str: + if hasattr(description, "name"): + return str(description.name) + if isinstance(description, Sequence) and description: + return str(description[0]) + raise TypeError(f"unsupported cursor description entry: {description!r}") + + +def _dict_rows(cursor: Any) -> list[dict]: + rows = cursor.fetchall() + if rows and isinstance(rows[0], Mapping): + return [dict(row) for row in rows] + columns = [_column_name(column) for column in cursor.description] + return [dict(zip(columns, row, strict=True)) for row in rows] + + +def _candidate_commit(conn: Any, family: str, scopes: Sequence[tuple[Any, ...]]) -> str | None: + table, alias = _FAMILY_CANDIDATE[family] + predicate, params = _scope_predicate(alias, _FAMILY_SCOPE[family], scopes) + cursor = conn.execute( + f""" + SELECT {alias}.commit_sha + FROM {table} {alias} + JOIN commits c USING (commit_sha) + WHERE {predicate} + ORDER BY c.timestamp DESC, c.commit_sha DESC + LIMIT 1 + """, + params, + ) + row = cursor.fetchone() + if row is None: + return None + if isinstance(row, Mapping): + return str(row["commit_sha"]) + return str(row[0]) + + +def fetch_baseline_records( + conn: Any, + pr_records: Sequence[Mapping[str, Any]], +) -> tuple[str, list[dict]]: + """Fetch the newest matching commit and its scoped v3 fact rows.""" + + family = benchmark_family(pr_records) + scopes = benchmark_scopes(pr_records, family) + commit_sha = _candidate_commit(conn, family, scopes) + if commit_sha is None: + raise ValueError(f"No RDS baseline found for {family} scopes {scopes!r}") + + records = [] + for _table, alias, select_sql, scope_columns, order_columns in _SELECTS[family]: + predicate, scope_params = _scope_predicate(alias, scope_columns, scopes) + order_by = ", ".join(order_columns) + cursor = conn.execute( + f""" + {select_sql} + WHERE {alias}.commit_sha = %s + AND {predicate} + ORDER BY {order_by} + """, + (commit_sha, *scope_params), + ) + records.extend(_dict_rows(cursor)) + + if not records: + raise ValueError(f"RDS selected baseline commit {commit_sha}, but it contained no scoped records") + return commit_sha, records + + +def write_records(path: Path, records: Sequence[Mapping[str, Any]]) -> None: + """Atomically write v3 records as compact JSONL.""" + + path.parent.mkdir(parents=True, exist_ok=True) + temporary_path = None + try: + with tempfile.NamedTemporaryFile( + mode="w", + encoding="utf-8", + dir=path.parent, + prefix=f".{path.name}.", + delete=False, + ) as output: + temporary_path = Path(output.name) + for record in records: + output.write(json.dumps(record, separators=(",", ":"))) + output.write("\n") + os.replace(temporary_path, path) + finally: + if temporary_path is not None: + temporary_path.unlink(missing_ok=True) + + +def _post_ingest_module(): + path = Path(__file__).resolve().with_name("post-ingest.py") + spec = importlib.util.spec_from_file_location("post_ingest", path) + assert spec is not None and spec.loader is not None + module = importlib.util.module_from_spec(spec) + spec.loader.exec_module(module) + return module + + +def main() -> int: + args = parse_args() + pr_records = read_records(args.pr_jsonl) + post_ingest = _post_ingest_module() + conn = post_ingest.connect_postgres(args.postgres, args.region) + try: + commit_sha, records = fetch_baseline_records(conn, pr_records) + finally: + conn.close() + write_records(args.output, records) + print(json.dumps({"commit_sha": commit_sha, "records": len(records)}, separators=(",", ":"))) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/scripts/tests/test_benchmark_reporting.py b/scripts/tests/test_benchmark_reporting.py index 7925161e91e..ba660b1c4cc 100644 --- a/scripts/tests/test_benchmark_reporting.py +++ b/scripts/tests/test_benchmark_reporting.py @@ -139,6 +139,192 @@ def test_read_latest_baseline_rows_streams_latest_matching_benchmark_commit(tmp_ assert len(selected) == 2 +def test_normalize_ingest_query_measurement_uses_comparison_shape() -> None: + compare = load_compare_module() + + rows = compare.normalize_ingest_records( + [ + { + "kind": "query_measurement", + "commit_sha": "base-sha", + "dataset": "tpch", + "scale_factor": "10", + "query_idx": 1, + "storage": "nvme", + "engine": "datafusion", + "format": "parquet", + "value_ns": 100, + "all_runtimes_ns": [90, 100, 110], + } + ] + ) + + assert rows.to_dict(orient="records") == [ + { + "name": "tpch_q01/datafusion:parquet", + "storage": "nvme", + "dataset": { + "dataset": "tpch", + "dataset_variant": None, + "scale_factor": "10", + }, + "unit": "ns", + "value": 100, + "all_runtimes": [90, 100, 110], + "commit_id": "base-sha", + } + ] + + +def test_normalize_ingest_compression_derives_cross_format_ratios() -> None: + compare = load_compare_module() + records = [] + for file_format, encode_ns, decode_ns, size_bytes in [ + ("vortex-file-compressed", 200, 100, 400), + ("parquet", 100, 200, 800), + ("lance", 400, 50, 200), + ]: + records.extend( + [ + { + "kind": "compression_time", + "commit_sha": "base-sha", + "dataset": "taxi", + "format": file_format, + "op": "encode", + "value_ns": encode_ns, + "all_runtimes_ns": [encode_ns], + }, + { + "kind": "compression_time", + "commit_sha": "base-sha", + "dataset": "taxi", + "format": file_format, + "op": "decode", + "value_ns": decode_ns, + "all_runtimes_ns": [decode_ns], + }, + { + "kind": "compression_size", + "commit_sha": "base-sha", + "dataset": "taxi", + "format": file_format, + "value_bytes": size_bytes, + }, + ] + ) + + rows = compare.normalize_ingest_records(records) + ratios = {row["name"]: row["value"] for row in rows.to_dict(orient="records") if row["unit"] == "ratio"} + + assert ratios == { + "vortex:parquet-zstd size/taxi": 0.5, + "vortex:lance size/taxi": 2.0, + "vortex:parquet-zstd ratio compress time/taxi": 2.0, + "vortex:lance ratio compress time/taxi": 0.5, + "vortex:parquet-zstd ratio decompress time/taxi": 0.5, + "vortex:lance ratio decompress time/taxi": 2.0, + } + + +def test_normalize_ingest_random_access_uses_legacy_display_name() -> None: + compare = load_compare_module() + + rows = compare.normalize_ingest_records( + [ + { + "kind": "random_access_time", + "commit_sha": "base-sha", + "dataset": "feature-vectors/correlated", + "format": "vortex-file-compressed", + "value_ns": 100, + "all_runtimes_ns": [90, 100, 110], + } + ] + ) + + assert rows.iloc[0]["name"] == ("random-access/feature-vectors/correlated/vortex-tokio-local-disk") + assert rows.iloc[0]["storage"] == "nvme" + assert rows.iloc[0]["all_runtimes"] == [90, 100, 110] + + legacy_taxi = compare.normalize_ingest_records( + [ + { + "kind": "random_access_time", + "commit_sha": "base-sha", + "dataset": "taxi", + "format": "parquet", + "value_ns": 100, + "all_runtimes_ns": [100], + } + ] + ) + assert legacy_taxi.iloc[0]["name"] == "random-access/parquet-tokio-local-disk" + + +def test_ingest_cli_uses_legacy_results_only_for_display_metadata(tmp_path: Path) -> None: + base_path = tmp_path / "base.ingest.jsonl" + pr_path = tmp_path / "pr.ingest.jsonl" + metadata_path = tmp_path / "results.json" + + base_records = [ + query_record_for_compare("base-sha", "parquet", 100), + query_record_for_compare("base-sha", "vortex-file-compressed", 80), + ] + pr_records = [ + query_record_for_compare("pr-sha", "parquet", 100), + query_record_for_compare("pr-sha", "vortex-file-compressed", 70), + ] + base_path.write_text( + "".join(f"{json.dumps(record)}\n" for record in base_records), + encoding="utf-8", + ) + pr_path.write_text( + "".join(f"{json.dumps(record)}\n" for record in pr_records), + encoding="utf-8", + ) + metadata_path.write_text( + f"{json.dumps({'commit_id': 'pr-sha', 'doc': 'vortex-bench/sql/tpch/README.md'})}\n", + encoding="utf-8", + ) + + result = subprocess.run( + [ + sys.executable, + str(COMPARE_SCRIPT), + "--ingest-jsonl", + "--metadata", + str(metadata_path), + str(base_path), + str(pr_path), + "TPC-H", + ], + check=False, + capture_output=True, + text=True, + ) + + assert result.returncode == 0, result.stderr + assert result.stdout.startswith("# Benchmarks: TPC-H [📖]") + assert "base base-sha" in result.stdout + assert "PR pr-sha" in result.stdout + + +def query_record_for_compare(commit: str, file_format: str, value: int) -> dict[str, object]: + return { + "kind": "query_measurement", + "commit_sha": commit, + "dataset": "tpch", + "scale_factor": "10", + "query_idx": 1, + "storage": "nvme", + "engine": "datafusion", + "format": file_format, + "value_ns": value, + "all_runtimes_ns": [value, value, value], + } + + def test_within_engine_analysis_uses_each_engines_own_parquet_control() -> None: compare = load_compare_module() rows = [ diff --git a/scripts/tests/test_fetch_benchmark_baseline.py b/scripts/tests/test_fetch_benchmark_baseline.py new file mode 100644 index 00000000000..55c8baa8330 --- /dev/null +++ b/scripts/tests/test_fetch_benchmark_baseline.py @@ -0,0 +1,215 @@ +# SPDX-License-Identifier: Apache-2.0 +# SPDX-FileCopyrightText: Copyright the Vortex contributors + +import importlib.util +from pathlib import Path + +import pytest + +REPO_ROOT = Path(__file__).resolve().parents[2] +FETCH_SCRIPT = REPO_ROOT / "scripts" / "fetch-benchmark-baseline.py" + + +def load_fetch_module(): + spec = importlib.util.spec_from_file_location("fetch_benchmark_baseline", FETCH_SCRIPT) + assert spec is not None + module = importlib.util.module_from_spec(spec) + assert spec.loader is not None + spec.loader.exec_module(module) + return module + + +class FakeCursor: + def __init__(self, columns: list[str], rows: list[tuple[object, ...]]): + self.description = [(column,) for column in columns] + self._rows = rows + + def fetchone(self): + return self._rows[0] if self._rows else None + + def fetchall(self): + return self._rows + + +class FakeConnection: + def __init__(self, responses: list[FakeCursor]): + self.responses = responses + self.calls: list[tuple[str, tuple[object, ...]]] = [] + + def execute(self, sql: str, params: tuple[object, ...]): + self.calls.append((sql, params)) + return self.responses.pop(0) + + +def query_record( + commit: str, + query_idx: int, + engine: str, + file_format: str, +) -> dict[str, object]: + return { + "kind": "query_measurement", + "commit_sha": commit, + "dataset": "tpch", + "scale_factor": "10", + "query_idx": query_idx, + "storage": "nvme", + "engine": engine, + "format": file_format, + "value_ns": 100, + "all_runtimes_ns": [90, 100, 110], + } + + +def test_query_baseline_selects_latest_commit_in_pr_scope() -> None: + fetch = load_fetch_module() + pr_records = [ + query_record("pr-sha", 1, "datafusion", "parquet"), + query_record("pr-sha", 1, "datafusion", "vortex-file-compressed"), + ] + columns = [ + "kind", + "commit_sha", + "dataset", + "dataset_variant", + "scale_factor", + "query_idx", + "storage", + "engine", + "format", + "value_ns", + "all_runtimes_ns", + "peak_physical", + "peak_virtual", + "physical_delta", + "virtual_delta", + "env_triple", + ] + baseline_row = ( + "query_measurement", + "base-new", + "tpch", + None, + "10", + 1, + "nvme", + "datafusion", + "parquet", + 95, + [85, 95, 105], + None, + None, + None, + None, + "x86_64-linux-gnu", + ) + conn = FakeConnection( + [ + FakeCursor(["commit_sha"], [("base-new",)]), + FakeCursor(columns, [baseline_row]), + ] + ) + + commit_sha, records = fetch.fetch_baseline_records(conn, pr_records) + + assert commit_sha == "base-new" + assert records == [dict(zip(columns, baseline_row, strict=True))] + assert len(conn.calls) == 2 + candidate_sql, candidate_params = conn.calls[0] + assert "FROM query_measurements q" in candidate_sql + assert "ORDER BY c.timestamp DESC, c.commit_sha DESC" in candidate_sql + assert candidate_params == ("tpch", "10", "nvme") + rows_sql, rows_params = conn.calls[1] + assert "q.commit_sha = %s" in rows_sql + assert rows_params == ("base-new", "tpch", "10", "nvme") + + +def test_compression_baseline_reads_times_and_sizes_from_same_commit() -> None: + fetch = load_fetch_module() + pr_records = [ + { + "kind": "compression_time", + "commit_sha": "pr-sha", + "dataset": "taxi", + "format": "vortex-file-compressed", + "op": "encode", + "value_ns": 200, + "all_runtimes_ns": [200], + }, + { + "kind": "compression_size", + "commit_sha": "pr-sha", + "dataset": "taxi", + "format": "vortex-file-compressed", + "value_bytes": 400, + }, + ] + time_columns = [ + "kind", + "commit_sha", + "dataset", + "dataset_variant", + "format", + "op", + "value_ns", + "all_runtimes_ns", + "env_triple", + ] + size_columns = [ + "kind", + "commit_sha", + "dataset", + "dataset_variant", + "format", + "value_bytes", + ] + time_row = ("compression_time", "base-sha", "taxi", None, "vortex-file-compressed", "encode", 180, [180], None) + size_row = ("compression_size", "base-sha", "taxi", None, "vortex-file-compressed", 390) + conn = FakeConnection( + [ + FakeCursor(["commit_sha"], [("base-sha",)]), + FakeCursor(time_columns, [time_row]), + FakeCursor(size_columns, [size_row]), + ] + ) + + commit_sha, records = fetch.fetch_baseline_records(conn, pr_records) + + assert commit_sha == "base-sha" + assert records == [ + dict(zip(time_columns, time_row, strict=True)), + dict(zip(size_columns, size_row, strict=True)), + ] + assert "FROM compression_times t" in conn.calls[1][0] + assert "FROM compression_sizes s" in conn.calls[2][0] + assert all(call[1][0] == "base-sha" for call in conn.calls[1:]) + + +def test_mixed_benchmark_families_are_rejected() -> None: + fetch = load_fetch_module() + + with pytest.raises(ValueError, match="multiple benchmark families"): + fetch.benchmark_family( + [ + query_record("pr-sha", 1, "datafusion", "parquet"), + { + "kind": "random_access_time", + "commit_sha": "pr-sha", + "dataset": "taxi", + "format": "parquet", + "value_ns": 10, + "all_runtimes_ns": [10], + }, + ] + ) + + +def test_missing_baseline_is_reported() -> None: + fetch = load_fetch_module() + conn = FakeConnection([FakeCursor(["commit_sha"], [])]) + + with pytest.raises(ValueError, match="No RDS baseline"): + fetch.fetch_baseline_records( + conn, + [query_record("pr-sha", 1, "datafusion", "parquet")], + )