From d1aa520c1c40a1fe5850f130ea3d02b8f616841a Mon Sep 17 00:00:00 2001 From: Dmitry Razdoburdin Date: Fri, 2 Oct 2026 04:56:30 -0700 Subject: [PATCH 1/4] initial --- .github/graph-metrics/CMakeLists.txt | 36 + .github/graph-metrics/ci.py | 481 +++++++ .github/graph-metrics/graph_metric.sh | 65 + .../graph_metrics_synthetic.json | 39 + .github/workflows/graph-metrics-comment.yml | 54 + .github/workflows/graph-metrics-synthetic.yml | 147 ++ utils/CMakeLists.txt | 5 + utils/graph_metrics.cpp | 1204 +++++++++++++++++ utils/graph_metrics_config.h | 664 +++++++++ 9 files changed, 2695 insertions(+) create mode 100644 .github/graph-metrics/CMakeLists.txt create mode 100644 .github/graph-metrics/ci.py create mode 100644 .github/graph-metrics/graph_metric.sh create mode 100644 .github/graph-metrics/graph_metrics_synthetic.json create mode 100644 .github/workflows/graph-metrics-comment.yml create mode 100644 .github/workflows/graph-metrics-synthetic.yml create mode 100644 utils/graph_metrics.cpp create mode 100644 utils/graph_metrics_config.h diff --git a/.github/graph-metrics/CMakeLists.txt b/.github/graph-metrics/CMakeLists.txt new file mode 100644 index 000000000..ccc6c4ed9 --- /dev/null +++ b/.github/graph-metrics/CMakeLists.txt @@ -0,0 +1,36 @@ +# Copyright 2026 Intel Corporation +# +# 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. + +cmake_minimum_required(VERSION 3.21) +project(graph_metrics_ci LANGUAGES CXX) + +# Use the same calculator and parameters with both library revisions, including +# a base revision that predates the calculator itself. +set(SVS_BUILD_BINARIES OFF CACHE BOOL "" FORCE) +set(SVS_BUILD_TESTS OFF CACHE BOOL "" FORCE) +set(SVS_BUILD_EXAMPLES OFF CACHE BOOL "" FORCE) +set(SVS_BUILD_BENCHMARK OFF CACHE BOOL "" FORCE) +set(SVS_BUILD_BENCHMARK_TEST_GENERATORS OFF CACHE BOOL "" FORCE) +add_subdirectory("${SVS_SOURCE_DIR}" svs) + +get_filename_component(METRICS_SOURCE_DIR "${CMAKE_CURRENT_LIST_DIR}/../.." ABSOLUTE) +include("${METRICS_SOURCE_DIR}/cmake/graph-metrics-dependencies.cmake") +add_executable( + graph_metrics "${METRICS_SOURCE_DIR}/utils/graph_metrics.cpp" +) +target_link_libraries( + graph_metrics PRIVATE + svs::svs svs_compile_options svs_x86_options_base + fmt::fmt nlohmann_json::nlohmann_json OpenSSL::Crypto +) diff --git a/.github/graph-metrics/ci.py b/.github/graph-metrics/ci.py new file mode 100644 index 000000000..72116f205 --- /dev/null +++ b/.github/graph-metrics/ci.py @@ -0,0 +1,481 @@ +# Copyright 2026 Intel Corporation +# +# 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. + +"""Dependency detection, paired graph-metrics runs, and PR reporting.""" + +import argparse +import copy +import html +import json +import math +import os +from pathlib import Path +import re +import shlex +import shutil +import subprocess +import urllib.request +import uuid + + +TARGET = "graph_metrics" +TOOL_DIRECTORY = Path(__file__).resolve().parent +COMMENT_MARKER = "" + + +def command(args, **kwargs): + return subprocess.check_output([str(arg) for arg in args], text=True, **kwargs).strip() + + +def write_json(path, value): + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text(json.dumps(value, indent=2) + "\n") + + +def output(name, value): + if "GITHUB_OUTPUT" in os.environ: + delimiter = uuid.uuid4().hex + with open(os.environ["GITHUB_OUTPUT"], "a") as stream: + stream.write(f"{name}<<{delimiter}\n{value}\n{delimiter}\n") + + +def cpu_count(): + return len(os.sched_getaffinity(0)) if hasattr(os, "sched_getaffinity") else os.cpu_count() or 1 + + +def compiler(): + selected = os.environ["CXX"] + executable = shutil.which(selected) + if executable is None: + raise RuntimeError(f"Configured compiler is unavailable: {selected}") + return executable + + +def configuration(head): + path = (head / os.environ["GRAPH_METRICS_CONFIG"]).resolve() + config = json.loads(path.read_text()) + for dataset in config["datasets"]: + if dataset["source"]["type"] != "synthetic": + raise ValueError("This workflow requires synthetic datasets") + dataset["source"]["path"] = str((path.parent / dataset["source"]["path"]).resolve()) + return config + + +def logged(args, log, cwd=None): + with log.open("a") as stream: + stream.write("$ " + shlex.join([str(arg) for arg in args]) + "\n") + stream.flush() + subprocess.run(args, cwd=cwd, stdout=stream, stderr=subprocess.STDOUT, check=True) + + +def configure(source, build, log, extra_args): + query = build / ".cmake/api/v1/query" + query.mkdir(parents=True, exist_ok=True) + for name in ("codemodel-v2", "cmakeFiles-v1"): + (query / name).touch() + logged([ + os.environ.get("CMAKE", "cmake"), + "-S", str(TOOL_DIRECTORY), "-B", str(build), "-G", "Ninja", + "-DCMAKE_BUILD_TYPE=Release", + f"-DCMAKE_CXX_COMPILER={compiler()}", + "-DCMAKE_EXPORT_COMPILE_COMMANDS=ON", + f"-DSVS_SOURCE_DIR={source}", + *extra_args, + ], log) + + +def cmake_reply(build, kind): + reply = build / ".cmake/api/v1/reply" + index = json.loads(sorted(reply.glob("index-*.json"))[-1].read_text()) + return json.loads((reply / index["reply"][kind]["jsonFile"]).read_text()) + + +def project_dependencies(source, head, build, log): + """Read CMake inputs and preprocess every TU in the calculator target closure.""" + model = cmake_reply(build, "codemodel-v2") + source_root = Path(model["paths"]["source"]) + targets = {item["id"]: item for item in model["configurations"][0]["targets"]} + pending = [item["id"] for item in targets.values() if item["name"] == TARGET] + if not pending: + raise RuntimeError(f"CMake did not define {TARGET}") + visited, translation_units = set(), set() + while pending: + target_id = pending.pop() + if target_id in visited: + continue + visited.add(target_id) + target = json.loads((build / ".cmake/api/v1/reply" / targets[target_id]["jsonFile"]).read_text()) + pending.extend(item["id"] for item in target.get("dependencies", [])) + for item in target.get("sources", []): + if "compileGroupIndex" in item: + origin = build if item.get("isGenerated") else source_root + translation_units.add((origin / item["path"]).resolve()) + + inputs = cmake_reply(build, "cmakeFiles-v1") + files = {(Path(inputs["paths"]["source"]) / item["path"]).resolve() for item in inputs["inputs"]} + commands = json.loads((build / "compile_commands.json").read_text()) + scanned = set() + for entry in commands: + unit = (Path(entry["directory"]) / entry["file"]).resolve() + if unit not in translation_units: + continue + arguments = entry.get("arguments") or shlex.split(entry["command"]) + filtered, skip = [], False + for argument in arguments: + if skip: + skip = False + elif argument in ("-o", "-MF", "-MT", "-MQ"): + skip = True + elif argument not in ("-c", "-MD", "-MMD"): + filtered.append(argument) + depfile = build / f"graph-metrics-{len(scanned)}.d" + # -M only preprocesses; detection does not compile or run the benchmark. + logged(filtered + ["-M", "-MF", str(depfile), "-MT", "graph_metrics"], + log, cwd=entry["directory"]) + dependencies = depfile.read_text().replace("\\\n", " ").partition(":")[2] + files.update( + (Path(entry["directory"]) / name.replace("$$", "$")).resolve() + for name in shlex.split(dependencies) + ) + scanned.add(unit) + if scanned != translation_units: + raise RuntimeError(f"Missing compile commands for {sorted(map(str, translation_units - scanned))}") + + relative = set() + for root in (source, head): + for path in files: + if path.is_relative_to(root): + relative.add(path.relative_to(root).as_posix()) + return relative + + +def workflow_dependencies(root, tool_relative): + # Discover the workflow support directory and the workflows referring to it, + # rather than maintaining a list of C++ or workflow dependency filenames. + result = set() + for path in (root / tool_relative).rglob("*"): + if path.is_file() and "__pycache__" not in path.parts: + result.add(path.relative_to(root).as_posix()) + for path in (root / ".github/workflows").glob("*"): + if path.suffix in (".yml", ".yaml") and tool_relative.as_posix() + "/" in path.read_text(): + result.add(path.relative_to(root).as_posix()) + return result + + +def detect(args): + args.result.mkdir(parents=True, exist_ok=True) + base_sha = command(["git", "-C", args.base, "rev-parse", "HEAD"]) + head_sha = command(["git", "-C", args.head, "rev-parse", "HEAD"]) + record = {"base_sha": base_sha, "head_sha": head_sha, "status": "failed"} + try: + # Three-dot diff matches the files changed by the PR, including deletions. + changed = set(command([ + "git", "-C", args.head, "diff", "--name-only", "--no-renames", "-z", + f"{base_sha}...{head_sha}", + ]).split("\0")) - {""} + dependencies = set() + tool_relative = TOOL_DIRECTORY.relative_to(args.head) + for label, source in (("base", args.base), ("head", args.head)): + build = args.work / label + log = args.result / f"{label}-dependencies.log" + configure(source, build, log, args.cmake_arg) + dependencies.update(project_dependencies(source, args.head, build, log)) + dependencies.update(workflow_dependencies(source, tool_relative)) + matched = sorted(changed & dependencies) + needed = bool(matched) or os.environ.get("FORCE_RUN") == "true" + record.update( + status="needed" if needed else "skipped", + changed_files=sorted(changed), + dependencies=sorted(dependencies), + matched_files=matched, + compiler=command([compiler(), "--version"]).splitlines()[0], + ) + output("needed", str(needed).lower()) + output("base_sha", base_sha) + output("head_sha", head_sha) + message = ("Graph metrics required: " + (", ".join(matched) or "manual run") + if needed else "Graph metrics skipped: no changed files affect the calculator or workflow.") + print(message) + if "GITHUB_STEP_SUMMARY" in os.environ: + with open(os.environ["GITHUB_STEP_SUMMARY"], "a") as stream: + stream.write(message + "\n") + finally: + write_json(args.result / "detection.json", record) + + +def paths(args): + cache_paths = [] + for dataset in configuration(args.head)["datasets"]: + path = dataset["source"]["path"] + cache_paths.extend([path, path + ".meta.json"]) + output("cache_paths", "\n".join(cache_paths)) + + +def calculate(args): + args.result.mkdir(parents=True, exist_ok=True) + config = configuration(args.head) + record = { + "base_sha": command(["git", "-C", args.base, "rev-parse", "HEAD"]), + "head_sha": command(["git", "-C", args.head, "rev-parse", "HEAD"]), + "compiler": command([compiler(), "--version"]).splitlines()[0], + "status": "failed", + } + try: + # Both revisions run sequentially on the same machine, with the same + # calculator, configuration and persisted vectors. + for label, source in (("base", args.base), ("head", args.head)): + build = args.work / label + configure(source, build, args.result / f"{label}-build.log", args.cmake_arg) + logged([ + os.environ.get("CMAKE", "cmake"), "--build", str(build), + "--target", TARGET, "--parallel", str(cpu_count()), + ], args.result / f"{label}-build.log") + effective = copy.deepcopy(config) + effective["output_json"] = str(args.result / f"{label}.json") + effective["output_log"] = str(args.result / f"{label}-metrics.log") + config_path = args.result / f"{label}-config.json" + write_json(config_path, effective) + logged([str(build / TARGET), "--config", str(config_path)], + args.result / f"{label}-build.log") + reports = [json.loads((args.result / f"{label}.json").read_text()) for label in ("base", "head")] + identities = [ + {item["dataset_name"]: item["dataset_parameters"]["sha256"] for item in report["datasets"]} + for report in reports + ] + if identities[0] != identities[1]: + raise RuntimeError("Base and PR results used different vector data") + record["status"] = "completed" + finally: + write_json(args.result / "run.json", record) + + +class GitHub: + def __init__(self): + self.root = os.environ.get("GITHUB_API_URL", "https://api.github.com") + self.repository = os.environ["GITHUB_REPOSITORY"] + + def request(self, method, path, body=None): + request = urllib.request.Request( + self.root + path, + data=json.dumps(body).encode() if body is not None else None, + method=method, + headers={ + "Authorization": "Bearer " + os.environ["GITHUB_TOKEN"], + "Accept": "application/vnd.github+json", + "Content-Type": "application/json", + "X-GitHub-Api-Version": "2022-11-28", + }, + ) + with urllib.request.urlopen(request, timeout=60) as response: + return json.load(response) + + def pages(self, path, key=None): + result = [] + for page in range(1, 100): + reply = self.request("GET", f"{path}?per_page=100&page={page}") + items = reply[key] if key else reply + result.extend(items) + if len(items) < 100: + return result + raise RuntimeError("GitHub pagination limit exceeded") + + +def artifact_json(path): + # Artifacts are untrusted input to the privileged comment job. + if path.stat().st_size > 10 * 1024 * 1024: + raise ValueError("Metrics artifact is too large") + with path.open() as stream: + return json.load(stream) + + +def artifact_record(path): + if not path.is_file(): + return None + try: + record = artifact_json(path) + if not isinstance(record, dict) or not isinstance(record.get("status"), str): + raise ValueError("Invalid status record") + for key in ("head_sha", "base_sha"): + if not isinstance(record.get(key), str) or not re.fullmatch(r"[0-9a-f]{40}", record[key]): + raise ValueError("Invalid commit identity") + return record + except (ValueError, OSError): + print(f"Ignoring invalid metadata artifact: {path.name}") + return None + + +def text(value): + value = re.sub(r"[\x00-\x1f\x7f]", " ", str(value)[:200]) + value = html.escape(value, quote=True).replace("|", "|").replace("@", "@") + return re.sub(r"([\\`*_\[\]()!])", r"\\\1", value) + + +def number(value, digits=None): + if isinstance(value, bool) or not isinstance(value, (int, float)) or not math.isfinite(value) or value < 0: + raise ValueError("Invalid numeric metric") + if digits is None: + if not isinstance(value, int): + raise ValueError("Expected an integer metric") + return str(value) + return f"{value:.{digits}f}" + + +def cells(build): + vertex = build["max_incoming_asymmetry_count_vertex"] + incoming = number(build["max_incoming_asymmetry_count"]) + if vertex is not None: + incoming += " @ vector " + number(vertex["vector_index"]) + return [ + number(build["unreachable_vertices"]), + number(build["edge_asymmetry_fraction"], 6), + "[" + ", ".join(number(value) for value in build["entry_points_internal"]) + "]", + incoming, + number(build["average_shortest_path_hops"], 3), + number(build["maximum_shortest_path_hops"]), + number(build["index_build_time_seconds"], 3), + ] + + +def comparison(base, head): + identities = [ + { + dataset["dataset_name"]: ( + dataset["dataset_parameters"]["sha256"], + dataset["dataset_parameters"]["selected_vectors"], + dataset["dataset_parameters"]["dimensions"], + dataset["dataset_parameters"]["distance"], + ) + for dataset in report["datasets"] + } + for report in (base, head) + ] + if identities[0] != identities[1]: + raise ValueError("Base and PR dataset identities differ") + before = { + (dataset["dataset_name"], index["index_type"], build["build_method"]): build + for dataset in base["datasets"] + for index in dataset["indexes"] + for build in index["builds"] + } + lines = ["Values show **base → PR**. Incoming-asymmetry vertices use source vector indices."] + for dataset in head["datasets"]: + for index in dataset["indexes"]: + lines.extend([ + "", f"**{text(dataset['dataset_name'])} — {text(index['index_type'])}**", "", + "| Build method | Unreachable | Asymmetry fraction | Entry points | Max incoming asymmetry (vertex) | Avg shortest path | Max shortest path | Build time (s) |", + "|---|---:|---:|---|---|---:|---:|---:|", + ]) + for build in index["builds"]: + key = (dataset["dataset_name"], index["index_type"], build["build_method"]) + previous = cells(before[key]) + current = cells(build) + values = [text(build["build_method"])] + [ + f"{old} → {new}" for old, new in zip(previous, current) + ] + lines.append("| " + " | ".join(values) + " |") + return "\n".join(lines) + + +def comment(args): + event = json.loads(Path(os.environ["GITHUB_EVENT_PATH"]).read_text()) + run = event["workflow_run"] + if run["event"] != "pull_request": + return + api = GitHub() + prefix = f"/repos/{api.repository}" + # Resolve the PR through GitHub's run/commit association, never an artifact's + # claimed PR number. Only current PR heads may update the comment. + run_prs = run.get("pull_requests") or [] + candidates = run_prs or api.pages(f"{prefix}/commits/{run['head_sha']}/pulls") + run_heads = {item["number"]: item.get("head", {}).get("sha", run["head_sha"]) for item in run_prs} + artifacts = api.pages(f"{prefix}/actions/runs/{run['id']}/artifacts", "artifacts") + detection_path = args.artifacts / "graph-metrics-results-detection/detection.json" + result_dir = args.artifacts / "graph-metrics-results-run" + detection = artifact_record(detection_path) + result_path = result_dir / "run.json" + result = artifact_record(result_path) + for candidate in candidates: + pr = api.request("GET", f"{prefix}/pulls/{int(candidate['number'])}") + associated_sha = run_heads.get(candidate["number"], run["head_sha"]) + if pr["state"] != "open" or pr["head"]["sha"] != associated_sha: + continue + current_sha = pr["head"]["sha"] + if any(record and record.get("head_sha") != current_sha for record in (detection, result)): + continue + base_sha = (detection or result or {}).get("base_sha", "") + if not re.fullmatch(r"[0-9a-f]{40}", base_sha): + base_sha = pr["base"]["sha"] + run_url = run["html_url"] + lines = [ + COMMENT_MARKER, + f"", + "### Synthetic graph metrics", "", + f"Base `{base_sha[:12]}` → PR `{current_sha[:12]}`.", + ] + if detection and detection["status"] == "skipped" and run["conclusion"] == "success": + lines.extend(["", "Skipped: no changed files affect the calculator or workflow dependencies."]) + elif result and result["status"] == "completed" and run["conclusion"] == "success": + try: + summary = comparison( + artifact_json(result_dir / "base.json"), + artifact_json(result_dir / "head.json"), + ) + lines.extend(["", f"Compiler: {text(result['compiler'])}.", "", summary]) + except (KeyError, TypeError, ValueError, OSError): + lines.extend(["", "The metrics report could not be validated. See the run and artifacts."]) + else: + lines.extend(["", f"Graph metrics are unavailable (run: {text(run['conclusion'])}). See the run logs."]) + links = ["", f"[Workflow run and logs]({run_url})"] + for artifact in artifacts: + if artifact["name"] == "graph-metrics-results-run" and not artifact.get("expired"): + artifact_url = f"{run_url}/artifacts/{int(artifact['id'])}" + links.append(f"[Full base/PR JSON reports, configurations, and logs]({artifact_url})") + body = "\n".join(lines)[:60000] + "\n" + "\n".join(links) + existing = next(( + item for item in api.pages(f"{prefix}/issues/{pr['number']}/comments") + if item["user"]["login"] == "github-actions[bot]" and COMMENT_MARKER in item["body"] + ), None) + if existing: + sequence = re.search(r"", existing["body"]) + current = (int(run["id"]), int(run.get("run_attempt", 1))) + if sequence and tuple(map(int, sequence.groups())) > current: + continue + api.request("PATCH", f"{prefix}/issues/comments/{existing['id']}", {"body": body}) + else: + api.request("POST", f"{prefix}/issues/{pr['number']}/comments", {"body": body}) + + +def main(): + parser = argparse.ArgumentParser(description=__doc__) + commands = parser.add_subparsers(dest="action", required=True) + for name in ("detect", "calculate"): + sub = commands.add_parser(name) + sub.add_argument("--head", type=Path, required=True) + sub.add_argument("--base", type=Path, required=True) + sub.add_argument("--work", type=Path, required=True) + sub.add_argument("--result", type=Path, required=True) + sub.add_argument("--cmake-arg", action="append", default=[]) + sub = commands.add_parser("paths") + sub.add_argument("--head", type=Path, required=True) + sub = commands.add_parser("comment") + sub.add_argument("--artifacts", type=Path, required=True) + args = parser.parse_args() + for key in ("head", "base", "work", "result", "artifacts"): + if hasattr(args, key): + setattr(args, key, getattr(args, key).resolve()) + {"detect": detect, "calculate": calculate, "paths": paths, "comment": comment}[args.action](args) + + +if __name__ == "__main__": + main() diff --git a/.github/graph-metrics/graph_metric.sh b/.github/graph-metrics/graph_metric.sh new file mode 100644 index 000000000..c9a786492 --- /dev/null +++ b/.github/graph-metrics/graph_metric.sh @@ -0,0 +1,65 @@ +#!/usr/bin/env bash +# Copyright 2026 Intel Corporation +# +# 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. + +set -euo pipefail + +graph_metric_root="$(cd -- "$(dirname -- "${BASH_SOURCE[0]}")/../.." && pwd)" + +# Dataset definitions, graph parameters, thread count, and report paths live in JSON. +graph_metric_config="${1:-${graph_metric_root}/graph_metrics_config.json}" +graph_metric_build_dir="${graph_metric_root}/build" +graph_metric_build_type="Release" +graph_metric_build_jobs=8 +graph_metric_cmake="auto" # Reuse the build directory's CMake, otherwise use PATH. +graph_metric_build_log="${graph_metric_build_dir}/graph-metrics/build.log" + +if (( $# > 1 )); then + printf 'Usage: bash .github/graph-metrics/graph_metric.sh [configuration.json]\n' >&2 + exit 2 +fi +if [[ ! -f "${graph_metric_config}" ]]; then + printf 'Configuration not found: %s\n' "${graph_metric_config}" >&2 + exit 2 +fi + +# Switching CMake installations in an existing build can invalidate FetchContent +# stamps and trigger dependency re-downloads. Keep its configured executable. +if [[ "${graph_metric_cmake}" == "auto" ]]; then + graph_metric_cmake="cmake" + if [[ -f "${graph_metric_build_dir}/CMakeCache.txt" ]]; then + while IFS= read -r graph_metric_cache_line; do + if [[ "${graph_metric_cache_line}" == CMAKE_COMMAND:INTERNAL=* ]]; then + graph_metric_cached_cmake="${graph_metric_cache_line#*=}" + if [[ -x "${graph_metric_cached_cmake}" ]]; then + graph_metric_cmake="${graph_metric_cached_cmake}" + fi + break + fi + done < "${graph_metric_build_dir}/CMakeCache.txt" + fi +fi + +printf 'Configuration: %s\nBuild log: %s\n' "${graph_metric_config}" "${graph_metric_build_log}" + +# The calculator handles runtime JSON/log destinations from the config itself. +mkdir -p -- "$(dirname -- "${graph_metric_build_log}")" +{ + "${graph_metric_cmake}" -S "${graph_metric_root}" -B "${graph_metric_build_dir}" \ + -DCMAKE_BUILD_TYPE="${graph_metric_build_type}" -DSVS_BUILD_BINARIES=ON + "${graph_metric_cmake}" --build "${graph_metric_build_dir}" --target graph_metrics \ + --parallel "${graph_metric_build_jobs}" +} > "${graph_metric_build_log}" 2>&1 + +exec "${graph_metric_build_dir}/utils/graph_metrics" --config "${graph_metric_config}" diff --git a/.github/graph-metrics/graph_metrics_synthetic.json b/.github/graph-metrics/graph_metrics_synthetic.json new file mode 100644 index 000000000..9df0a5f39 --- /dev/null +++ b/.github/graph-metrics/graph_metrics_synthetic.json @@ -0,0 +1,39 @@ +{ + "threads": "auto", + "graph_type": "both", + "build_methods": [ + "sample_by_sample", + "single_batch", + "batches" + ], + "batch_size": 1000, + "recompute_entry_point": true, + "build_parameters": { + "degree": 64, + "window": 512, + "max_candidates": 750, + "alpha": { + "L2": 1.2, + "MIP": 0.95 + }, + "use_full_search_history": true + }, + "datasets": [ + { + "name": "synthetic-128", + "vectors": 500000, + "dimensions": 128, + "distance": "L2", + "source": { + "type": "synthetic", + "path": "../../data/graph-metrics/synthetic-128.fvecs", + "distribution": "normal", + "mean": 0.0, + "stddev": 1.0, + "seed": 12345 + } + } + ], + "output_json": "../../build/graph-metrics/graph-metrics-synthetic.json", + "output_log": "../../build/graph-metrics/graph-metrics-synthetic.log" +} diff --git a/.github/workflows/graph-metrics-comment.yml b/.github/workflows/graph-metrics-comment.yml new file mode 100644 index 000000000..0c7a15099 --- /dev/null +++ b/.github/workflows/graph-metrics-comment.yml @@ -0,0 +1,54 @@ +# Copyright 2026 Intel Corporation +# +# 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. + +name: Graph Metrics PR Comment + +on: + workflow_run: + workflows: ['Synthetic Graph Metrics'] + types: [completed] + +permissions: + contents: read + actions: read + pull-requests: write + +jobs: + comment: + if: github.event.workflow_run.event == 'pull_request' + runs-on: ubuntu-22.04 + timeout-minutes: 10 + steps: + # The write-capable job runs only default-branch code. PR artifacts are + # parsed as JSON; no PR scripts or binaries execute in this job. + - uses: actions/checkout@v6 + with: + ref: ${{ github.event.repository.default_branch }} + persist-credentials: false + + - name: Download dependency decision and metrics + uses: actions/download-artifact@v8 + continue-on-error: true + with: + pattern: graph-metrics-results-* + path: ${{ runner.temp }}/graph-metrics-artifacts + run-id: ${{ github.event.workflow_run.id }} + github-token: ${{ github.token }} + + - name: Update the PR metrics or skipped message + env: + GITHUB_TOKEN: ${{ github.token }} + run: | + python3 .github/graph-metrics/ci.py comment \ + --artifacts "${RUNNER_TEMP}/graph-metrics-artifacts" diff --git a/.github/workflows/graph-metrics-synthetic.yml b/.github/workflows/graph-metrics-synthetic.yml new file mode 100644 index 000000000..8d1b4f50d --- /dev/null +++ b/.github/workflows/graph-metrics-synthetic.yml @@ -0,0 +1,147 @@ +# Copyright 2026 Intel Corporation +# +# 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. + +name: Synthetic Graph Metrics + +on: + pull_request: + workflow_dispatch: + inputs: + cxx: + description: 'C++ compiler executable; overrides GRAPH_METRICS_CXX' + required: false + type: string + base_ref: + description: 'Comparison base; defaults to the repository default branch' + required: false + type: string + +permissions: + contents: read + +env: + # Repository variable or manual-run input can select an installed compiler. + CXX: ${{ inputs.cxx || vars.GRAPH_METRICS_CXX || 'g++-11' }} + GRAPH_METRICS_CONFIG: .github/graph-metrics/graph_metrics_synthetic.json + +concurrency: + group: ${{ github.workflow }} @ ${{ github.event.pull_request.number || github.ref }} + cancel-in-progress: true + +jobs: + detect: + name: Detect graph-metrics dependencies + runs-on: ubuntu-22.04 + timeout-minutes: 30 + outputs: + needed: ${{ steps.detect.outputs.needed }} + base_sha: ${{ steps.detect.outputs.base_sha }} + head_sha: ${{ steps.detect.outputs.head_sha }} + + steps: + - name: Check out PR revision + uses: actions/checkout@v6 + with: + ref: ${{ github.event.pull_request.head.sha || github.sha }} + path: head + fetch-depth: 0 + persist-credentials: false + + - name: Check out base revision + uses: actions/checkout@v6 + with: + ref: ${{ github.event.pull_request.base.sha || inputs.base_ref || github.event.repository.default_branch }} + path: base + fetch-depth: 0 + persist-credentials: false + + - name: Check configured compiler and install build dependencies + run: | + command -v "${CXX}" + "${CXX}" --version + sudo apt-get update + sudo apt-get install --yes libssl-dev ninja-build + + - name: Detect dependencies and decide whether metrics are needed + id: detect + env: + FORCE_RUN: ${{ github.event_name == 'workflow_dispatch' }} + run: | + python3 head/.github/graph-metrics/ci.py detect \ + --head head --base base \ + --work "${RUNNER_TEMP}/graph-metrics-dependencies" \ + --result results/detection + + - name: Upload dependency decision and logs + if: ${{ always() }} + uses: actions/upload-artifact@v7 + with: + name: graph-metrics-results-detection + path: results/detection/ + if-no-files-found: warn + + calculate: + name: Compare base and PR graph metrics + needs: detect + if: needs.detect.outputs.needed == 'true' + runs-on: ubuntu-22.04 + timeout-minutes: 180 + steps: + - name: Check out PR revision + uses: actions/checkout@v6 + with: + ref: ${{ needs.detect.outputs.head_sha }} + path: head + persist-credentials: false + + - name: Check out base revision + uses: actions/checkout@v6 + with: + ref: ${{ needs.detect.outputs.base_sha }} + path: base + persist-credentials: false + + - name: Check configured compiler and install build dependencies + run: | + command -v "${CXX}" + "${CXX}" --version + sudo apt-get update + sudo apt-get install --yes libssl-dev ninja-build + + - name: Read synthetic cache paths + id: paths + run: python3 head/.github/graph-metrics/ci.py paths --head head + + - name: Restore reusable synthetic vectors + uses: actions/cache@v4 + with: + path: ${{ steps.paths.outputs.cache_paths }} + key: graph-metrics-synthetic-v1-${{ runner.os }}-${{ runner.arch }}-${{ hashFiles('head/.github/graph-metrics/graph_metrics_synthetic.json', 'head/utils/graph_metrics_config.h') }} + + # Both revisions use the PR calculator and identical configuration/data. + # The calculator's threads: auto setting uses this runner's available CPUs. + - name: Build and calculate base and PR metrics + run: | + python3 head/.github/graph-metrics/ci.py calculate \ + --head head --base base \ + --work "${RUNNER_TEMP}/graph-metrics-build" \ + --result results/run + + - name: Upload full JSON reports, configurations and logs + if: ${{ always() }} + uses: actions/upload-artifact@v7 + with: + name: graph-metrics-results-run + path: results/run/ + if-no-files-found: warn diff --git a/utils/CMakeLists.txt b/utils/CMakeLists.txt index 888b32f46..c602a85dc 100644 --- a/utils/CMakeLists.txt +++ b/utils/CMakeLists.txt @@ -37,6 +37,11 @@ function(create_utility exe file) endfunction() create_utility(graph_stat graph_stat.cpp) +include("${PROJECT_SOURCE_DIR}/cmake/graph-metrics-dependencies.cmake") +create_utility(graph_metrics graph_metrics.cpp) +target_link_libraries( + graph_metrics PRIVATE nlohmann_json::nlohmann_json OpenSSL::Crypto +) # Legacy conversion routines. create_utility(convert_legacy convert_legacy.cpp) diff --git a/utils/graph_metrics.cpp b/utils/graph_metrics.cpp new file mode 100644 index 000000000..a8dcdf6bc --- /dev/null +++ b/utils/graph_metrics.cpp @@ -0,0 +1,1204 @@ +/* + * Copyright 2026 Intel Corporation + * + * 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. + */ + +#include "graph_metrics_config.h" + +#include "svs/concurrent/dynamic_index.h" +#include "svs/core/data/view.h" +#include "svs/index/vamana/dynamic_index.h" +#include "svs/index/vamana/extensions.h" +#include "svs/lib/float16.h" +#include "svs/lib/threads.h" + +#include "fmt/core.h" +#include "fmt/ranges.h" + +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +namespace { + +namespace cc = svs::index::vamana::concurrent; +using Idx = uint32_t; +// Reverse-neighbor lists inherit the graph allocator. One mmap per small list +// can exhaust vm.max_map_count when cleanup splits merged mappings. Keep the +// concurrent, grow-stable blocked layout, but allocate through the heap. +using Graph = cc::graphs::SimpleGraph>>; +using Data = cc::SegmentedBlockedData; +using NonConcurrentGraph = + svs::graphs::SimpleGraph>>; +using NonConcurrentData = svs::data::BlockedData; +template +using Index = std::conditional_t< + Concurrent, + cc::MutableVamanaIndex, + svs::index::vamana:: + MutableVamanaIndex>; + +using graph_metrics::Input; +using graph_metrics::Json; +using graph_metrics::Options; + +struct BuildPlan { + bool concurrent; + std::string_view method; + size_t batch_size; + size_t insertion_calls; + size_t caller_threads; + size_t worker_threads; +}; + +std::string_view index_type(bool concurrent) { + return concurrent ? "svs::index::vamana::concurrent::MutableVamanaIndex" + : "svs::index::vamana::MutableVamanaIndex"; +} + +std::vector build_plans(const Options& options) { + std::vector plans; + for (bool concurrent : {true, false}) { + const auto type = concurrent ? "concurrent" : "non_concurrent"; + if (options.graph_type != "both" && options.graph_type != type) { + continue; + } + for (std::string_view method : options.build_methods) { + if (!concurrent && method == "sample_by_sample") { + continue; + } + const size_t batch_size = method == "single_batch" ? options.total + : method == "sample_by_sample" ? 1 + : options.batch_size; + const size_t calls = + options.total / batch_size + (options.total % batch_size != 0); + // Only concurrent singleton calls overlap. Batches are submitted in + // source order, using the index's worker pool inside each call. + const size_t callers = method == "sample_by_sample" ? options.threads : 1; + const size_t workers = method == "sample_by_sample" ? 1 : options.threads; + plans.push_back({concurrent, method, batch_size, calls, callers, workers}); + } + } + return plans; +} + +constexpr std::string_view help = R"(Usage: graph_metrics --config FILE.json + graph_metrics --self-test + +All dataset definitions and run parameters come from FILE.json. See +graph_metrics_config.json for an example configuration. +Relative source, synthetic-cache, and output paths resolve beside the config file. +threads may be a positive integer or "auto" (CPUs available to this process). + +File inputs use little-endian vecs layout with uint32 dimensions before every row. +Coordinates may be uint8, int8, float16, or float32 and are indexed as float32. +Synthetic inputs support normal or uniform distributions, are generated once into +float32 vecs files, and are validated against saved metadata and SHA-256 on reuse. +Mismatched, incomplete, or corrupted synthetic files are regenerated automatically. +Generation, checksum validation, and loading are excluded from index build time. + +Every build starts empty. sample_by_sample uses parallel singleton callers and is +concurrent-only. single_batch inserts all rows in one call; batches inserts every +row in sequential calls of batch_size rows, with parallel work inside each call. +Optional entry-point recomputation reuses each singleton/batched graph unchanged. + +The report groups datasets -> indexes -> builds, with field descriptions first. +Actual forward edges determine reciprocity and directed shortest paths. Reverse +edge supersets are never used. Every asymmetry maximum reports its source vector +index. Shortest-path means include reachable entries at zero and exclude +unreachable vertices. Progress goes to stderr or the configured output_log. +)"; + +template float source_float(T value) { + if constexpr (std::is_same_v) { + const auto bits = std::bit_cast(value); + const auto exponent = bits & 0x7c00; + if (exponent == 0x7c00) { + return std::numeric_limits::quiet_NaN(); + } + if (exponent == 0) { + // Preserve IEEE half subnormals, which Float16's fast scalar + // conversion otherwise flushes to zero. + const float magnitude = std::ldexp(static_cast(bits & 0x03ff), -24); + return (bits & 0x8000) ? -magnitude : magnitude; + } + } + return static_cast(value); +} + +template +void load_vecs(const Input& input, svs::data::SimpleData& points) { + std::ifstream stream{input.path, std::ios::binary}; + std::vector row(points.dimensions()); + graph_metrics::Sha256 hash; + for (size_t i = 0; i < points.size(); ++i) { + uint32_t dimensions = 0; + if (!stream.read(reinterpret_cast(&dimensions), sizeof(dimensions)) || + dimensions != points.dimensions() || + !stream.read(reinterpret_cast(row.data()), row.size() * sizeof(T))) { + throw std::runtime_error( + fmt::format("Invalid or truncated vector {} in {}", i, input.path.string()) + ); + } + hash.update(&dimensions, sizeof(dimensions)); + hash.update(row.data(), row.size() * sizeof(T)); + auto destination = points.get_datum(i); + for (size_t d = 0; d < row.size(); ++d) { + const float value = source_float(row[d]); + if (!std::isfinite(value)) { + throw std::runtime_error(fmt::format( + "Non-finite coordinate at vector {}, dimension {} in {}", + i, + d, + input.path.string() + )); + } + destination[d] = value; + } + } + graph_metrics::require( + hash.finish() == input.sha256, + "Dataset changed after checksum validation: " + input.path.string() + ); +} + +svs::data::SimpleData load_points(const Options& options, const Input& input) { + fmt::print( + stderr, + "Loading {}: {} vectors, dimensions={}...\n", + input.name, + options.total, + options.dimensions + ); + auto points = svs::data::SimpleData(options.total, options.dimensions); + if (input.source_type == "uint8") { + load_vecs(input, points); + } else if (input.source_type == "int8") { + load_vecs(input, points); + } else if (input.source_type == "float16") { + load_vecs(input, points); + } else { + load_vecs(input, points); + } + return points; +} + +// A sorted CSR copy of ACTUAL forward edges. Sorting only this copy preserves the +// constructed graph and allows reciprocal-edge tests in O(log(out-degree)). +// Call only once writers have joined; a per-node seqlock is not a global snapshot. +struct ForwardGraph { + std::vector offsets; + std::vector edges; + + size_t size() const { return offsets.size() - 1; } + + std::span neighbors(size_t vertex) const { + return std::span{edges}.subspan( + offsets[vertex], offsets[vertex + 1] - offsets[vertex] + ); + } + + template + explicit ForwardGraph(const GraphType& graph) + : offsets(graph.n_nodes() + 1, 0) { + for (size_t u = 0; u < graph.n_nodes(); ++u) { + offsets[u + 1] = offsets[u] + graph.get_node_degree(static_cast(u)); + } + edges.resize(offsets.back()); + for (size_t u = 0; u < graph.n_nodes(); ++u) { + const auto adjacent = graph.get_node(static_cast(u)); + std::copy(adjacent.begin(), adjacent.end(), edges.begin() + offsets[u]); + auto first = edges.begin() + offsets[u]; + auto last = edges.begin() + offsets[u + 1]; + std::sort(first, last); + if (std::adjacent_find(first, last) != last) { + throw std::runtime_error(fmt::format("Duplicate edge at vertex {}", u)); + } + for (Idx v : neighbors(u)) { + if (v >= size() || v == u) { + throw std::runtime_error(fmt::format("Invalid edge {} -> {}", u, v)); + } + } + } + } +}; + +double fraction(uint64_t numerator, uint64_t denominator) { + return denominator == 0 + ? 0.0 + : static_cast(numerator) / static_cast(denominator); +} + +struct Metrics { + size_t vertices = 0; + size_t edges = 0; + size_t reachable = 0; + size_t unreciprocated = 0; + size_t max_outgoing_count = 0; + double max_outgoing_fraction = 0; + size_t max_incoming_count = 0; + double max_incoming_fraction = 0; + std::optional max_outgoing_count_vertex; + std::optional max_outgoing_fraction_vertex; + std::optional max_incoming_count_vertex; + std::optional max_incoming_fraction_vertex; + double average_hops = 0; + size_t maximum_hops = 0; + double average_out_degree = 0; + size_t maximum_out_degree = 0; + std::vector out_degree_histogram; +}; + +template +void update_maximum(T value, Idx vertex, T& maximum, std::optional& maximum_vertex) { + // Vertices are visited in ascending order: retain the first vertex on ties, + // including an all-zero maximum. An empty graph leaves the vertex null. + if (!maximum_vertex || value > maximum) { + maximum = value; + maximum_vertex = vertex; + } +} + +Metrics measure_edges(const ForwardGraph& graph) { + Metrics result; + result.vertices = graph.size(); + result.edges = graph.edges.size(); + std::vector in_degree(graph.size(), 0); + std::vector incoming_asymmetry(graph.size(), 0); + for (size_t u = 0; u < graph.size(); ++u) { + size_t outgoing_asymmetry = 0; + const auto adjacent = graph.neighbors(u); + for (Idx v : adjacent) { + ++in_degree[v]; + const auto reverse = graph.neighbors(v); + if (!std::binary_search(reverse.begin(), reverse.end(), static_cast(u))) { + ++outgoing_asymmetry; + ++incoming_asymmetry[v]; + } + } + result.unreciprocated += outgoing_asymmetry; + update_maximum( + outgoing_asymmetry, + static_cast(u), + result.max_outgoing_count, + result.max_outgoing_count_vertex + ); + update_maximum( + fraction(outgoing_asymmetry, adjacent.size()), + static_cast(u), + result.max_outgoing_fraction, + result.max_outgoing_fraction_vertex + ); + } + for (size_t v = 0; v < graph.size(); ++v) { + update_maximum( + incoming_asymmetry[v], + static_cast(v), + result.max_incoming_count, + result.max_incoming_count_vertex + ); + update_maximum( + fraction(incoming_asymmetry[v], in_degree[v]), + static_cast(v), + result.max_incoming_fraction, + result.max_incoming_fraction_vertex + ); + } + return result; +} + +// Reuse edge/asymmetry measurements when only the entry point changes. +void measure_reachability( + const ForwardGraph& graph, std::span entry_points, Metrics& result +) { + result.maximum_hops = 0; + result.maximum_out_degree = 0; + result.out_degree_histogram.clear(); + constexpr Idx unseen = std::numeric_limits::max(); + std::vector hops(graph.size(), unseen); + std::vector queue; + queue.reserve(graph.size()); + for (Idx entry : entry_points) { + if (entry >= graph.size()) { + throw std::runtime_error("Entry point is outside the graph"); + } + if (hops[entry] == unseen) { + hops[entry] = 0; + queue.push_back(entry); + } + } + uint64_t total_hops = 0; + uint64_t total_degree = 0; + for (size_t head = 0; head < queue.size(); ++head) { + const Idx u = queue[head]; + total_hops += hops[u]; + result.maximum_hops = std::max(result.maximum_hops, size_t{hops[u]}); + const auto adjacent = graph.neighbors(u); + const size_t degree = adjacent.size(); + total_degree += degree; + result.maximum_out_degree = std::max(result.maximum_out_degree, degree); + if (result.out_degree_histogram.size() <= degree) { + result.out_degree_histogram.resize(degree + 1, 0); + } + ++result.out_degree_histogram[degree]; + for (Idx v : adjacent) { + if (hops[v] == unseen) { + hops[v] = hops[u] + 1; + queue.push_back(v); + } + } + } + result.reachable = queue.size(); + result.average_hops = fraction(total_hops, result.reachable); + result.average_out_degree = fraction(total_degree, result.reachable); +} + +Metrics measure(const ForwardGraph& graph, std::span entry_points) { + auto result = measure_edges(graph); + measure_reachability(graph, entry_points, result); + return result; +} + +template +void insert_points( + IndexType& index, + const svs::data::SimpleData& points, + const Options& options, + const BuildPlan& plan +) { + fmt::print( + stderr, + "Inserting {} vectors, up to {} per call, from {} caller(s), {} worker(s) per " + "call...\n", + options.total, + plan.batch_size, + plan.caller_threads, + plan.worker_threads + ); + const size_t step = std::min(plan.batch_size, options.total); + std::atomic next{0}; + std::atomic completed{0}; + std::atomic calls{0}; + std::atomic failed{false}; + std::mutex error_mutex; + std::exception_ptr error; + std::vector workers; + workers.reserve(plan.caller_threads); + const auto insert = [&] { + try { + while (!failed.load(std::memory_order_relaxed)) { + const size_t first = next.fetch_add(step, std::memory_order_relaxed); + if (first >= options.total) { + break; + } + const size_t last = first + std::min(step, options.total - first); + const auto ids = svs::threads::UnitRange{first, last}; + const auto batch = svs::data::make_const_view(points, ids); + index.add_points(batch, ids); + calls.fetch_add(1, std::memory_order_relaxed); + const size_t count = + completed.fetch_add(last - first, std::memory_order_relaxed) + last - + first; + if (count % 10'000 == 0 || count == options.total) { + fmt::print(stderr, "Inserted {}/{}\n", count, options.total); + } + } + } catch (...) { + failed.store(true, std::memory_order_relaxed); + std::lock_guard lock{error_mutex}; + if (!error) { + error = std::current_exception(); + } + } + }; + try { + if (plan.caller_threads == 1) { + // Non-concurrent indexes are never called from overlapping writers. + insert(); + } else { + for (size_t i = 0; i < plan.caller_threads; ++i) { + workers.emplace_back(insert); + } + } + } catch (...) { + failed.store(true, std::memory_order_relaxed); + for (auto& worker : workers) { + worker.join(); + } + throw; + } + for (auto& worker : workers) { + worker.join(); + } + if (error) { + std::rethrow_exception(error); + } + if (completed.load() != options.total || calls.load() != plan.insertion_calls) { + throw std::runtime_error("Insertion workload did not match the requested build plan" + ); + } +} + +template struct BuiltIndex { + std::unique_ptr index; + double index_build_time_seconds; +}; + +template +BuiltIndex> build_index( + const Options& options, const Input& input, Distance distance, const BuildPlan& plan +) { + auto points = load_points(options, input); + + // Time only index construction: input I/O/conversion/generation is finished. + // Include index storage allocation and copying, construction, all additions, + // and writer joins. Stop before validation, graph measurement, and cleanup. + const auto build_start = std::chrono::steady_clock::now(); + using IndexType = Index; + typename IndexType::data_type empty_data(0, options.dimensions); + const auto ids = svs::threads::UnitRange{0, 0}; + const svs::index::vamana::VamanaBuildParameters parameters{ + options.alpha, + options.degree, + options.window, + options.max_candidates, + options.degree, + options.use_full_search_history}; + auto logger = + std::make_shared("graph_metrics", svs::logging::stderr_sink()); + logger->set_level(spdlog::level::warn); + // Inline internal work avoids the native pool's mutex serializing + // independent concurrent singleton callers. Batch calls use worker threads. + auto pool = plan.method == "sample_by_sample" + ? svs::threads::ThreadPoolHandle{svs::threads::SequentialThreadPool{}} + : svs::threads::ThreadPoolHandle{ + svs::threads::as_threadpool(plan.worker_threads)}; + std::unique_ptr index; + fmt::print(stderr, "Starting empty; every vector will be inserted via add_points...\n"); + if constexpr (Concurrent) { + index = std::make_unique( + parameters, std::move(empty_data), ids, distance, std::move(pool), logger + ); + } else { + // The non-concurrent build-from-data constructor requires a medoid. + // Its graph/data constructor accepts empty storage without building. + // Routing slot zero is only a placeholder: the first add_points call + // fills it (and every other row in that batch) before graph traversal. + index = std::make_unique( + NonConcurrentGraph{0, options.degree}, + std::move(empty_data), + Idx{0}, + distance, + ids, + std::move(pool), + logger + ); + index->set_alpha(options.alpha); + index->set_construction_window_size(options.window); + index->set_max_candidates(options.max_candidates); + index->set_prune_to(options.degree); + index->set_full_search_history(options.use_full_search_history); + } + if (index->size() != 0 || index->view_graph().n_nodes() != 0) { + throw std::runtime_error("Every insertion method must start with zero vertices"); + } + insert_points(*index, points, options, plan); + const double build_seconds = + std::chrono::duration(std::chrono::steady_clock::now() - build_start) + .count(); + // This workload performs no deletions or failed insertions. Every allocated + // graph vertex must therefore be live; internal IDs need not equal external IDs. + if (index->size() != options.total || index->view_graph().n_nodes() != options.total) { + throw std::runtime_error("Final index size does not equal configured vectors"); + } + return {std::move(index), build_seconds}; +} + +std::string json_string(std::string_view text) { + std::string result = "\""; + for (unsigned char c : text) { + if (c == '"' || c == '\\') { + result += '\\'; + result += static_cast(c); + } else if (c < 0x20) { + result += fmt::format("\\u{:04x}", c); + } else { + result += static_cast(c); + } + } + result += '"'; + return result; +} + +// JSON has no comments: emit this dictionary as the first member of the document. +void print_field_descriptions() { + constexpr std::pair descriptions[] = { + {"field_descriptions", + "Field definitions for datasets, their indexes, and each index's builds."}, + {"configuration_file", + "Absolute path of the JSON configuration used for this run."}, + {"datasets", + "One object per selected dataset; metadata is shared by all its index types and " + "builds."}, + {"dataset_name", "Unique dataset name from the configuration."}, + {"indexes", + "Index implementation groups within a dataset, each containing its build " + "results."}, + {"index_type", + "Fully qualified C++ class name, without template arguments: " + "svs::index::vamana::concurrent::MutableVamanaIndex or " + "svs::index::vamana::MutableVamanaIndex."}, + {"builds", + "Build-method results within one index implementation; recomputed-entry variants " + "reuse their base graph."}, + {"build_method", + "sample_by_sample: concurrent only, start empty and add every vector with " + "singleton calls; " + "single_batch: start empty and insert all rows in one add_points call; " + "batches: start empty and insert all rows in fixed-size calls, with a smaller " + "final batch " + "if needed. sample_by_sample_recompute_entry_point and " + "batches_recompute_entry_point reuse the respective base graph with an entry " + "computed over all indexed vectors by the builder's compute_entry_point routine " + "(nearest vector to the mean, using squared L2 for both L2 and MIP indexes)."}, + {"reused_from_build_method", + "Base build_method whose graph is reused, or null for a fresh build. Only the " + "entry used to measure reachability and shortest paths changes."}, + {"index_build_time_seconds", + "Monotonic wall-clock seconds for index storage allocation/copying, construction, " + "all add_points calls, and writer joins. Recomputed-entry variants report the " + "base build time plus entry_point_recompute_time_seconds, measured separately " + "without rebuilding the graph. Excludes source " + "loading/conversion/generation, " + "graph metrics, JSON output, and index cleanup."}, + {"entry_point_recompute_time_seconds", + "Wall-clock seconds for compute_entry_point on all indexed vectors, including " + "its worker-pool startup/join; zero for base builds. Already included in " + "index_build_time_seconds."}, + {"dataset_parameters", + "Input provenance and selection; values describe the data actually loaded."}, + {"dataset_parameters.path", "Absolute path of the file containing the vectors."}, + {"dataset_parameters.source_kind", "file or synthetic, as configured."}, + {"dataset_parameters.format", + "vecs: little-endian uint32 dimension per row followed by typed coordinates."}, + {"dataset_parameters.source_type", + "Coordinate type in the source file; float16 occupies two bytes even in .fvecs " + "files."}, + {"dataset_parameters.index_data_type", + "Coordinate type used by the index after conversion: float32."}, + {"dataset_parameters.source_vectors", + "Number of rows in the complete source file."}, + {"dataset_parameters.selected_vectors", + "Number of source vectors used by every build for this dataset."}, + {"dataset_parameters.dimensions", "Coordinates per vector."}, + {"dataset_parameters.distance", + "L2 is squared Euclidean distance; MIP maximizes inner product."}, + {"dataset_parameters.selection", + "prefix: first selected_vectors rows in file order, including saved synthetic " + "data."}, + {"dataset_parameters.normalized", + "False: no normalization is applied to the loaded vectors."}, + {"dataset_parameters.seed", + "MT19937 seed for synthetic vectors; null for file datasets."}, + {"dataset_parameters.distribution", + "normal or uniform for synthetic vectors; null for file datasets."}, + {"dataset_parameters.generation", + "Synthetic metadata (generator version, dimensions, count, type, distribution, " + "parameters, seed); null for file datasets. Saved beside synthetic data as " + ".meta.json, with an additional sha256 field."}, + {"dataset_parameters.sha256", + "SHA-256 of the selected prefix's raw vecs bytes, including row headers. " + "Rechecked while loading every build; covers the complete file for synthetic " + "data."}, + {"configuration", "Effective construction and insertion parameters for this run."}, + {"configuration.insertion_batch_size", + "Maximum rows per add_points call; total for single_batch."}, + {"configuration.insertion_calls", + "Number of add_points calls from an empty index: total for sample_by_sample, " + "one for single_batch, ceil(total / batch_size) for batches."}, + {"configuration.insertion_caller_threads", + "Caller threads submitting additions: configured threads for concurrent singleton " + "insertion, one for either batch method."}, + {"configuration.insertion_worker_threads", + "Workers inside each add_points call: one for singletons, configured threads for " + "either batch method."}, + {"configuration.entry_point_recompute_threads", + "Workers used to recompute the entry point; zero for base builds."}, + {"configuration.threads", + "Configured thread count; both batch methods use this many " + "workers, " + "concurrent singletons use this many callers."}, + {"configuration.degree", "Maximum out-degree and prune_to."}, + {"configuration.window", "Construction search window."}, + {"configuration.max_candidates", "Construction candidate limit."}, + {"configuration.alpha", "Configured pruning parameter for the dataset's distance."}, + {"configuration.use_full_search_history", + "Whether construction uses full search history."}, + {"entry_points_internal", + "Internal entry IDs used for these metrics after all writers join: original " + "index entries for base builds, or the recomputed entry for reused-graph " + "variants."}, + {"vertices", "Number of live vertices, including unreachable vertices."}, + {"edges", + "Number of actual directed edges, including edges at unreachable vertices."}, + {"reachable_vertices", + "Vertices reachable from any entry via directed forward edges, including " + "entries."}, + {"unreachable_vertices", "vertices minus reachable_vertices."}, + {"unreachable_fraction", "unreachable_vertices / vertices."}, + {"unreciprocated_edges", + "Actual edges u->v without an actual v->u. Reverse-edge superset records are " + "ignored."}, + {"edge_asymmetry_fraction", + "unreciprocated_edges / edges; a bidirectional pair counts as two edges."}, + {"max_outgoing_asymmetry_count", + "Maximum per-vertex number of outgoing edges without a reverse edge, over all " + "vertices."}, + {"max_outgoing_asymmetry_fraction", + "Maximum per-vertex outgoing asymmetric count / out-degree, over all vertices."}, + {"max_incoming_asymmetry_count", + "Maximum per-vertex number of incoming edges without a reverse edge, over all " + "vertices."}, + {"max_incoming_asymmetry_fraction", + "Maximum per-vertex incoming asymmetric count / actual in-degree, over all " + "vertices."}, + {"max_outgoing_asymmetry_count_vertex", + "Vertex attaining the outgoing count maximum."}, + {"max_outgoing_asymmetry_fraction_vertex", + "Vertex attaining the outgoing fraction maximum; computed independently of " + "count."}, + {"max_incoming_asymmetry_count_vertex", + "Vertex attaining the incoming count maximum."}, + {"max_incoming_asymmetry_fraction_vertex", + "Vertex attaining the incoming fraction maximum; computed independently of " + "count."}, + {"*_vertex.internal_id", + "Internal graph ID. Ties choose the smallest internal ID; the entire vertex " + "object is null for an empty graph."}, + {"*_vertex.vector_index", + "Zero-based row in the source dataset, or generated vector number. Translated " + "from internal ID after concurrent insertion."}, + {"average_shortest_path_hops", + "Mean directed, unweighted shortest path from the nearest entry; reachable " + "vertices only, entries at zero."}, + {"maximum_shortest_path_hops", + "Maximum directed shortest path from the nearest entry over reachable vertices."}, + {"reachable_average_out_degree", "Mean actual out-degree over reachable vertices."}, + {"reachable_maximum_out_degree", + "Maximum actual out-degree over reachable vertices."}, + {"reachable_out_degree_histogram", + "Element d counts reachable vertices with out-degree d."}, + {"zero_denominators", + "All fractions and means with zero denominator are reported as zero."}, + }; + fmt::print(" \"field_descriptions\": {{\n"); + for (size_t i = 0; i < std::size(descriptions); ++i) { + fmt::print( + " {}: {}{}\n", + json_string(descriptions[i].first), + json_string(descriptions[i].second), + i + 1 == std::size(descriptions) ? "" : "," + ); + } + fmt::print(" }},\n"); +} + +template +std::string vertex_json(std::optional vertex, Translate&& translate) { + if (!vertex) { + return "null"; + } + return fmt::format( + "{{\"internal_id\": {}, \"vector_index\": {}}}", *vertex, translate(*vertex) + ); +} + +void print_dataset_header(const Options& options, const Input& input) { + fmt::print( + " {{\n" + " \"dataset_name\": {},\n" + " \"dataset_parameters\": {{\n" + " \"path\": {},\n" + " \"source_kind\": {},\n" + " \"format\": \"vecs\",\n" + " \"source_type\": {},\n" + " \"index_data_type\": \"float32\",\n" + " \"source_vectors\": {},\n" + " \"selected_vectors\": {},\n" + " \"dimensions\": {},\n" + " \"distance\": {},\n" + " \"selection\": \"prefix\",\n" + " \"normalized\": false,\n" + " \"seed\": {},\n" + " \"distribution\": {},\n" + " \"generation\": {},\n" + " \"sha256\": {}\n" + " }},\n" + " \"indexes\": [\n", + json_string(input.name), + json_string(input.path.string()), + json_string(input.kind), + json_string(input.source_type), + input.source_vectors, + options.total, + options.dimensions, + json_string(input.distance), + input.generation.is_null() ? "null" : input.generation.at("seed").dump(), + input.generation.is_null() ? "null" : input.generation.at("distribution").dump(), + input.generation.dump(), + json_string(input.sha256) + ); +} + +template +void print_build_report( + const Options& options, + const BuildPlan& plan, + double index_build_time_seconds, + std::optional entry_point_recompute_seconds, + std::span entry_points, + const Metrics& metrics, + Translate&& translate +) { + fmt::print( + " {{\n" + " \"build_method\": {},\n" + " \"reused_from_build_method\": {},\n" + " \"index_build_time_seconds\": {},\n" + " \"entry_point_recompute_time_seconds\": {},\n" + " \"configuration\": {{\n" + " \"threads\": {},\n" + " \"degree\": {},\n" + " \"window\": {},\n" + " \"max_candidates\": {},\n" + " \"alpha\": {},\n" + " \"use_full_search_history\": {},\n" + " \"insertion_batch_size\": {},\n" + " \"insertion_calls\": {},\n" + " \"insertion_caller_threads\": {},\n" + " \"insertion_worker_threads\": {},\n" + " \"entry_point_recompute_threads\": {}\n" + " }},\n" + " \"entry_points_internal\": [{}],\n", + json_string( + entry_point_recompute_seconds + ? fmt::format("{}_recompute_entry_point", plan.method) + : std::string{plan.method} + ), + entry_point_recompute_seconds ? json_string(plan.method) : "null", + index_build_time_seconds + entry_point_recompute_seconds.value_or(0), + entry_point_recompute_seconds.value_or(0), + options.threads, + options.degree, + options.window, + options.max_candidates, + options.alpha, + options.use_full_search_history, + plan.batch_size, + plan.insertion_calls, + plan.caller_threads, + plan.worker_threads, + entry_point_recompute_seconds ? options.threads : 0, + fmt::join(entry_points, ", ") + ); + fmt::print( + " \"vertices\": {},\n" + " \"edges\": {},\n" + " \"reachable_vertices\": {},\n" + " \"unreachable_vertices\": {},\n" + " \"unreachable_fraction\": {},\n" + " \"unreciprocated_edges\": {},\n" + " \"edge_asymmetry_fraction\": {},\n" + " \"max_outgoing_asymmetry_count\": {},\n" + " \"max_outgoing_asymmetry_count_vertex\": {},\n" + " \"max_outgoing_asymmetry_fraction\": {},\n" + " \"max_outgoing_asymmetry_fraction_vertex\": {},\n" + " \"max_incoming_asymmetry_count\": {},\n" + " \"max_incoming_asymmetry_count_vertex\": {},\n" + " \"max_incoming_asymmetry_fraction\": {},\n" + " \"max_incoming_asymmetry_fraction_vertex\": {},\n", + metrics.vertices, + metrics.edges, + metrics.reachable, + metrics.vertices - metrics.reachable, + fraction(metrics.vertices - metrics.reachable, metrics.vertices), + metrics.unreciprocated, + fraction(metrics.unreciprocated, metrics.edges), + metrics.max_outgoing_count, + vertex_json(metrics.max_outgoing_count_vertex, translate), + metrics.max_outgoing_fraction, + vertex_json(metrics.max_outgoing_fraction_vertex, translate), + metrics.max_incoming_count, + vertex_json(metrics.max_incoming_count_vertex, translate), + metrics.max_incoming_fraction, + vertex_json(metrics.max_incoming_fraction_vertex, translate) + ); + fmt::print( + " \"average_shortest_path_hops\": {},\n" + " \"maximum_shortest_path_hops\": {},\n" + " \"reachable_average_out_degree\": {},\n" + " \"reachable_maximum_out_degree\": {},\n" + " \"reachable_out_degree_histogram\": [{}]\n" + " }}", + metrics.average_hops, + metrics.maximum_hops, + metrics.average_out_degree, + metrics.maximum_out_degree, + fmt::join(metrics.out_degree_histogram, ", ") + ); + // Preserve the completed report even if index cleanup subsequently fails. + if (std::fflush(stdout) != 0) { + throw std::runtime_error("Failed to flush the JSON report"); + } +} + +// Hand-calculated fixtures exercise direction, unreachable vertices, denominators, +// independent count/fraction maxima, multiple sources, and stale reverse records. +void self_test() { + const auto require = [](bool condition, std::string_view name) { + if (!condition) { + throw std::runtime_error(fmt::format("Self-test failed: {}", name)); + } + }; + Graph graph(8, 3); + graph.enable_reverse_edges(); + const std::array, 8> edges{ + {{1, 2, 3}, {0, 2}, {3}, {1, 7}, {5}, {4}, {2}, {2}}}; + for (Idx u = 0; u < edges.size(); ++u) { + graph.replace_node(u, edges[u]); + } + // Pretend every edge also has a reverse edge in the superset. None of these + // extra records changes the actual forward adjacency used for measurement. + for (Idx u = 0; u < edges.size(); ++u) { + for (Idx v : edges[u]) { + graph.reverse_edges()->record(v, u); + } + } + const ForwardGraph snapshot{graph}; + const auto result = measure(snapshot, std::array{0}); + require(result.vertices == 8 && result.edges == 12, "vertex/edge counts"); + require(result.reachable == 5, "directed reachability"); + require(result.unreciprocated == 8, "actual edges, including unreachable sources"); + require(result.max_outgoing_count == 2, "outgoing count maximum"); + require(result.max_outgoing_fraction == 1.0, "outgoing fraction maximum"); + require(result.max_incoming_count == 4, "incoming count maximum"); + require(result.max_incoming_fraction == 1.0, "incoming fraction maximum"); + require( + result.max_outgoing_count_vertex == 0 && result.max_outgoing_fraction_vertex == 2 && + result.max_incoming_count_vertex == 2 && + result.max_incoming_fraction_vertex == 2, + "maximum vertices and smallest-ID tie breaking" + ); + require(result.average_hops == 1.0 && result.maximum_hops == 2, "BFS distances"); + require(result.average_out_degree == 1.8, "reachable mean out-degree"); + require(result.maximum_out_degree == 3, "reachable maximum out-degree"); + require( + result.out_degree_histogram == std::vector{0, 2, 2, 1}, + "reachable out-degree histogram" + ); + + const auto multiple = measure(snapshot, std::array{0, 4, 6, 0}); + require(multiple.reachable == 8, "multiple/duplicate entry points"); + require( + multiple.average_hops == 0.75 && multiple.maximum_hops == 2, + "minimum distance from any entry" + ); + auto changed_entry = result; + measure_reachability(snapshot, std::array{4}, changed_entry); + require( + changed_entry.reachable == 2 && changed_entry.average_hops == 0.5 && + changed_entry.maximum_hops == 1 && changed_entry.average_out_degree == 1 && + changed_entry.maximum_out_degree == 1 && + changed_entry.out_degree_histogram == std::vector{0, 2}, + "changing entry resets all reachability metrics on the reused graph" + ); + require( + changed_entry.edges == result.edges && + changed_entry.unreciprocated == result.unreciprocated && + changed_entry.max_outgoing_count_vertex == result.max_outgoing_count_vertex && + changed_entry.max_incoming_count_vertex == result.max_incoming_count_vertex, + "changing entry preserves edge/asymmetry measurements" + ); + + // Count maxima occur at 0 (outgoing) and 3 (incoming), but fraction maxima + // occur at 5 and 4 respectively. Vertex 5 is unreachable from the entry point. + Graph distinct_maxima(7, 4); + const std::array, 7> other_edges{ + {{1, 2, 3, 4}, {0}, {0}, {6}, {}, {3}, {3}}}; + for (Idx u = 0; u < other_edges.size(); ++u) { + distinct_maxima.replace_node(u, other_edges[u]); + } + const auto maxima = measure(ForwardGraph{distinct_maxima}, std::array{0}); + require( + maxima.max_outgoing_count == 2 && maxima.max_outgoing_fraction == 1.0 && + maxima.max_incoming_count == 2 && maxima.max_incoming_fraction == 1.0, + "count and fraction maxima at different vertices" + ); + require( + maxima.max_outgoing_count_vertex == 0 && maxima.max_outgoing_fraction_vertex == 5 && + maxima.max_incoming_count_vertex == 3 && + maxima.max_incoming_fraction_vertex == 4, + "separate vertices for count and fraction maxima" + ); + require( + vertex_json(maxima.max_outgoing_fraction_vertex, [](Idx id) { return 100 + id; }) == + "{\"internal_id\": 5, \"vector_index\": 105}", + "translate maximum vertex to the original dataset row" + ); + + Graph empty_edges(3, 2); + const auto isolated = measure(ForwardGraph{empty_edges}, std::array{0}); + require(isolated.reachable == 1 && isolated.edges == 0, "isolated entry point"); + require( + isolated.average_hops == 0 && isolated.average_out_degree == 0 && + isolated.max_incoming_fraction == 0 && isolated.max_outgoing_fraction == 0 && + fraction(0, 0) == 0, + "zero denominators" + ); + require( + isolated.max_outgoing_count_vertex == 0 && + isolated.max_incoming_count_vertex == 0 && + isolated.max_outgoing_fraction_vertex == 0 && + isolated.max_incoming_fraction_vertex == 0, + "all-zero maxima identify the first vertex" + ); + const auto no_entry = measure(snapshot, {}); + require(no_entry.reachable == 0 && no_entry.average_hops == 0, "no entry points"); + Graph empty(0, 2); + const auto empty_result = measure(ForwardGraph{empty}, {}); + require( + empty_result.vertices == 0 && !empty_result.max_outgoing_count_vertex && + !empty_result.max_outgoing_fraction_vertex && + !empty_result.max_incoming_count_vertex && + !empty_result.max_incoming_fraction_vertex, + "empty graph has no maximum vertices" + ); + require( + vertex_json( + {}, [](Idx) -> size_t { throw std::runtime_error("Unexpected translation"); } + ) == "null", + "empty maximum serializes as null without translation" + ); + require( + json_string("path\"\\\n\t") == "\"path\\\"\\\\\\u000a\\u0009\"", "JSON escaping" + ); + require( + source_float(int8_t{-128}) == -128.0f && source_float(uint8_t{255}) == 255.0f, + "signed and unsigned byte conversion" + ); + const auto half = [](uint16_t bits) { return std::bit_cast(bits); }; + require( + source_float(half(0x3e00)) == 1.5f && source_float(half(0xbe00)) == -1.5f && + source_float(half(0x0001)) == 0x1p-24f && + source_float(half(0x8001)) == -0x1p-24f && + !std::isfinite(source_float(half(0x7c00))) && + !std::isfinite(source_float(half(0x7e00))), + "float16 conversion preserves subnormals and rejects non-finite values" + ); + + // Allocate reverse lists in shuffled order, then free them in vertex order. + // The previous mmap-backed allocator aborts here under the common Linux + // vm.max_map_count=65530 limit, even though constructing the graph succeeds. + { + constexpr size_t count = 300'000; + Graph large(count, 1); + large.enable_reverse_edges(); + std::vector order(count); + std::iota(order.begin(), order.end(), 0); + std::mt19937 generator{42}; + std::shuffle(order.begin(), order.end(), generator); + for (Idx u : order) { + large.add_edge(u, (u + 1) % count); + } + } + fmt::print("Graph metric self-tests passed.\n"); +} + +template +void run_dataset( + const Options& options, const Input& input, Distance distance, const BuildPlan& plan +) { + const auto graph_type = index_type(Concurrent); + fmt::print(stderr, "Run: {} / {} / {}\n", input.name, graph_type, plan.method); + auto [index, build_seconds] = + build_index(options, input, distance, plan); + fmt::print(stderr, "Index build time: {:.6f} seconds.\n", build_seconds); + fmt::print(stderr, "All writers joined. Measuring the forward graph...\n"); + index->experimental_escape_hatch( + [&](const auto& graph, const auto& data, const auto&, std::span entries + ) { + const ForwardGraph snapshot{graph}; + auto metrics = measure(snapshot, entries); + const auto translate = [&](Idx id) { return index->translate_internal_id(id); }; + print_build_report( + options, plan, build_seconds, std::nullopt, entries, metrics, translate + ); + if (options.recompute_entry_point && plan.method != "single_batch") { + fmt::print(stderr, "Recomputing entry point on the same graph...\n"); + const auto recompute_start = std::chrono::steady_clock::now(); + Idx entry; + { + svs::threads::NativeThreadPool pool{options.threads}; + entry = svs::lib::narrow( + svs::index::vamana::extensions::compute_entry_point(data, pool) + ); + } + const double recompute_seconds = + std::chrono::duration( + std::chrono::steady_clock::now() - recompute_start + ) + .count(); + fmt::print( + stderr, + "Entry point: {} -> {}; recompute time: {:.6f} seconds.\n", + fmt::join(entries, ", "), + entry, + recompute_seconds + ); + // The utility only measures the built graph. Use the selected entry as + // the BFS source without mutating the index or its adjacency lists. + const std::array recomputed_entries{entry}; + measure_reachability(snapshot, recomputed_entries, metrics); + fmt::print(",\n"); + print_build_report( + options, + plan, + build_seconds, + recompute_seconds, + recomputed_entries, + metrics, + translate + ); + } + } + ); + fmt::print(stderr, "JSON report flushed. Releasing index...\n"); + index.reset(); + fmt::print(stderr, "Done: {} / {} / {}.\n", input.name, graph_type, plan.method); +} + +template +void dispatch_run( + const Options& options, const Input& input, Distance distance, const BuildPlan& plan +) { + if (plan.concurrent) { + run_dataset(options, input, distance, plan); + } else { + run_dataset(options, input, distance, plan); + } +} + +} // namespace + +int main(int argc, char** argv) { + try { + if (argc == 2 && std::string_view{argv[1]} == "--help") { + fmt::print("{}", help); + return 0; + } + if (argc == 2 && std::string_view{argv[1]} == "--self-test") { + self_test(); + return 0; + } + if (argc != 3 || std::string_view{argv[1]} != "--config") { + throw std::invalid_argument("Usage: graph_metrics --config FILE.json"); + } + auto config = graph_metrics::read_configuration(argv[2]); + fmt::print(stderr, "JSON: {}\n", config.output_json.string()); + if (!config.output_log.empty()) { + fmt::print(stderr, "Log: {}\n", config.output_log.string()); + graph_metrics::redirect_stream(stderr, config.output_log); + } + graph_metrics::prepare_inputs(config); + // Publish a complete report atomically. Invalid configurations or failed + // builds leave any previous report intact. + graph_metrics::TemporaryFile output{config.output_json}; + graph_metrics::redirect_stream(stdout, output.path); + fmt::print("{{\n"); + print_field_descriptions(); + fmt::print(" \"configuration_file\": {},\n", json_string(config.path.string())); + fmt::print(" \"datasets\": [\n"); + size_t completed_datasets = 0; + for (const auto& [effective, input] : config.inputs) { + const auto plans = build_plans(effective); + if (completed_datasets++ != 0) { + fmt::print(",\n"); + } + print_dataset_header(effective, input); + size_t completed_indexes = 0; + for (bool concurrent : {true, false}) { + if (std::none_of(plans.begin(), plans.end(), [&](const auto& plan) { + return plan.concurrent == concurrent; + })) { + continue; + } + if (completed_indexes++ != 0) { + fmt::print(",\n"); + } + fmt::print( + " {{\n \"index_type\": {},\n \"builds\": [\n", + json_string(index_type(concurrent)) + ); + size_t completed_builds = 0; + for (const auto& plan : plans) { + if (plan.concurrent != concurrent) { + continue; + } + if (completed_builds++ != 0) { + fmt::print(",\n"); + } + if (input.distance == "MIP") { + dispatch_run(effective, input, svs::DistanceIP{}, plan); + } else { + dispatch_run(effective, input, svs::DistanceL2{}, plan); + } + } + fmt::print("\n ]\n }}"); + } + fmt::print("\n ]\n }}"); + } + fmt::print("\n ]\n}}\n"); + if (std::fflush(stdout) != 0) { + throw std::runtime_error("Failed to flush the JSON report"); + } + output.commit(config.output_json); + fmt::print(stderr, "JSON report: {}\n", config.output_json.string()); + return 0; + } catch (const std::exception& error) { + fmt::print(stderr, "Error: {}\n", error.what()); + return 1; + } +} diff --git a/utils/graph_metrics_config.h b/utils/graph_metrics_config.h new file mode 100644 index 000000000..b8d2d362e --- /dev/null +++ b/utils/graph_metrics_config.h @@ -0,0 +1,664 @@ +/* + * Copyright 2026 Intel Corporation + * + * 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. + */ + +#pragma once + +#include "fmt/core.h" +#include "nlohmann/json.hpp" +#include + +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#include +#include +#include +#ifdef __linux__ +#include +#endif + +namespace graph_metrics { + +using Json = nlohmann::json; +namespace fs = std::filesystem; + +struct Options { + size_t total = 0; + size_t dimensions = 0; + size_t threads = 0; + size_t degree = 0; + size_t window = 0; + size_t max_candidates = 0; + float alpha = 0; + std::string graph_type; + std::vector build_methods; + size_t batch_size = 0; + bool recompute_entry_point = false; + bool use_full_search_history = true; +}; + +struct Input { + std::string name; + fs::path path; + std::string source_type; + std::string distance; + size_t source_vectors = 0; + std::string kind; + Json generation = nullptr; + std::string sha256; +}; + +struct Configuration { + fs::path path; + fs::path output_json; + fs::path output_log; + std::vector> inputs; +}; + +inline void require(bool condition, const std::string& message) { + if (!condition) { + throw std::runtime_error(message); + } +} + +inline Json read_json(const fs::path& path) { + std::ifstream stream{path}; + require(bool(stream), fmt::format("Cannot read JSON: {}", path.string())); + // The parser normally keeps the last value of a duplicate key. Reject + // duplicates so a misspecified configuration cannot silently change a run. + std::vector> keys; + auto callback = [&](int, Json::parse_event_t event, Json& value) { + if (event == Json::parse_event_t::object_start) { + keys.emplace_back(); + } else if (event == Json::parse_event_t::key) { + const auto key = value.get(); + require(keys.back().insert(key).second, "Duplicate JSON key: " + key); + } else if (event == Json::parse_event_t::object_end) { + keys.pop_back(); + } + return true; + }; + return Json::parse(stream, callback); +} + +inline void object_keys( + const Json& value, + std::string_view context, + std::initializer_list required, + std::initializer_list optional = {} +) { + require(value.is_object(), fmt::format("{} must be an object", context)); + for (auto key : required) { + require(value.contains(key), fmt::format("Missing {}.{}", context, key)); + } + for (auto it = value.begin(); it != value.end(); ++it) { + require( + std::find(required.begin(), required.end(), it.key()) != required.end() || + std::find(optional.begin(), optional.end(), it.key()) != optional.end(), + fmt::format("Unknown field {}.{}", context, it.key()) + ); + } +} + +inline std::string text(const Json& value, std::string_view context) { + require(value.is_string(), fmt::format("{} must be a string", context)); + auto result = value.get(); + require( + !result.empty() && result.find('\0') == std::string::npos, + fmt::format("{} must be nonempty and contain no NUL", context) + ); + return result; +} + +inline size_t integer( + const Json& value, + std::string_view context, + size_t minimum = 1, + size_t maximum = std::numeric_limits::max() +) { + require( + value.is_number_unsigned() && value.get() >= minimum && + value.get() <= maximum, + fmt::format("{} must be an integer in [{}, {}]", context, minimum, maximum) + ); + return value.get(); +} + +inline double number(const Json& value, std::string_view context) { + require(value.is_number(), fmt::format("{} must be a number", context)); + const auto result = value.get(); + require(std::isfinite(result), fmt::format("{} must be finite", context)); + return result; +} + +inline bool boolean(const Json& value, std::string_view context) { + require(value.is_boolean(), fmt::format("{} must be a boolean", context)); + return value.get(); +} + +inline size_t available_threads() { +#ifdef __linux__ + cpu_set_t affinity; + if (sched_getaffinity(0, sizeof(affinity), &affinity) == 0 && + CPU_COUNT(&affinity) > 0) { + return static_cast(CPU_COUNT(&affinity)); + } +#endif + return std::max(1u, std::thread::hardware_concurrency()); +} + +inline size_t element_size(const std::string& type) { + if (type == "uint8" || type == "int8") { + return 1; + } + if (type == "float16") { + return 2; + } + require(type == "float32", "Unsupported source type: " + type); + return 4; +} + +inline fs::path resolve(const fs::path& directory, const Json& value) { + return fs::weakly_canonical(directory / text(value, "path")); +} + +inline fs::path metadata_path(const Input& input) { + return input.path.string() + ".meta.json"; +} + +inline Configuration read_configuration(const fs::path& filename) { + Configuration config; + config.path = fs::canonical(filename); + const auto directory = config.path.parent_path(); + const auto json = read_json(config.path); + object_keys( + json, + "config", + {"threads", + "graph_type", + "build_methods", + "batch_size", + "recompute_entry_point", + "build_parameters", + "datasets", + "output_json"}, + {"output_log"} + ); + Options common; + const auto& threads = json.at("threads"); + if (threads.is_string()) { + require(threads == "auto", "threads must be a positive integer or \"auto\""); + common.threads = available_threads(); + } else { + common.threads = integer(threads, "threads"); + } + common.graph_type = text(json.at("graph_type"), "graph_type"); + require( + common.graph_type == "concurrent" || common.graph_type == "non_concurrent" || + common.graph_type == "both", + "graph_type must be concurrent, non_concurrent, or both" + ); + const auto& methods = json.at("build_methods"); + require( + methods.is_array() && !methods.empty(), "build_methods must be a nonempty array" + ); + std::set seen_methods; + for (const auto& method : methods) { + const auto name = text(method, "build_methods entry"); + require( + name == "sample_by_sample" || name == "single_batch" || name == "batches", + "Unknown build method: " + name + ); + require(seen_methods.insert(name).second, "Duplicate build method: " + name); + common.build_methods.push_back(name); + } + require( + common.graph_type != "non_concurrent" || seen_methods.size() != 1 || + !seen_methods.contains("sample_by_sample"), + "sample_by_sample requires a concurrent index" + ); + common.batch_size = integer(json.at("batch_size"), "batch_size"); + common.recompute_entry_point = + boolean(json.at("recompute_entry_point"), "recompute_entry_point"); + const auto& build = json.at("build_parameters"); + object_keys( + build, + "build_parameters", + {"degree", "window", "max_candidates", "alpha"}, + {"use_full_search_history"} + ); + common.degree = + integer(build.at("degree"), "degree", 1, std::numeric_limits::max() - 1); + common.window = integer(build.at("window"), "window"); + common.max_candidates = integer(build.at("max_candidates"), "max_candidates"); + require( + common.max_candidates >= common.window && common.window >= common.degree, + "Require max_candidates >= window >= degree" + ); + if (build.contains("use_full_search_history")) { + common.use_full_search_history = + boolean(build.at("use_full_search_history"), "use_full_search_history"); + } + const auto& alpha = build.at("alpha"); + object_keys(alpha, "alpha", {"L2", "MIP"}); + const auto alpha_l2 = number(alpha.at("L2"), "alpha.L2"); + const auto alpha_mip = number(alpha.at("MIP"), "alpha.MIP"); + require( + alpha_l2 >= 1 && alpha_l2 <= std::numeric_limits::max() && alpha_mip > 0 && + alpha_mip <= 1 && static_cast(alpha_mip) > 0, + "Require alpha.L2 >= 1 and alpha.MIP in (0, 1], representable as float32" + ); + config.output_json = resolve(directory, json.at("output_json")); + if (json.contains("output_log")) { + config.output_log = resolve(directory, json.at("output_log")); + } + const auto& datasets = json.at("datasets"); + require(datasets.is_array() && !datasets.empty(), "datasets must be a nonempty array"); + std::set names; + std::map generators; + std::set protected_paths{config.path}; + for (const auto& dataset : datasets) { + object_keys( + dataset, "dataset", {"name", "vectors", "dimensions", "distance", "source"} + ); + auto options = common; + Input input; + input.name = text(dataset.at("name"), "dataset.name"); + require(names.insert(input.name).second, "Duplicate dataset name: " + input.name); + options.total = integer( + dataset.at("vectors"), "vectors", 1, std::numeric_limits::max() - 1 + ); + options.dimensions = integer( + dataset.at("dimensions"), "dimensions", 1, std::numeric_limits::max() + ); + require( + options.dimensions <= + (std::numeric_limits::max() / options.total - 4) / sizeof(float), + "Vector data size overflows size_t" + ); + input.distance = text(dataset.at("distance"), "distance"); + require( + input.distance == "L2" || input.distance == "MIP", "distance must be L2 or MIP" + ); + options.alpha = static_cast(input.distance == "L2" ? alpha_l2 : alpha_mip); + const auto& source = dataset.at("source"); + require(source.is_object() && source.contains("type"), "Missing source.type"); + input.kind = text(source.at("type"), "source.type"); + if (input.kind == "file") { + object_keys(source, "source", {"type", "path", "data_type"}); + input.source_type = text(source.at("data_type"), "source.data_type"); + element_size(input.source_type); + } else { + require(input.kind == "synthetic", "source.type must be file or synthetic"); + require(source.contains("distribution"), "Missing source.distribution"); + const auto distribution = text(source.at("distribution"), "distribution"); + Json parameters; + if (distribution == "normal") { + object_keys( + source, + "source", + {"type", "path", "distribution", "mean", "stddev", "seed"} + ); + parameters = { + {"mean", number(source.at("mean"), "mean")}, + {"stddev", number(source.at("stddev"), "stddev")}}; + require(parameters.at("stddev").get() > 0, "stddev must be > 0"); + } else { + require( + distribution == "uniform", "distribution must be normal or uniform" + ); + object_keys( + source, + "source", + {"type", "path", "distribution", "lower", "upper", "seed"} + ); + parameters = { + {"lower", number(source.at("lower"), "lower")}, + {"upper", number(source.at("upper"), "upper")}}; + require( + parameters.at("lower").get() < + parameters.at("upper").get(), + "Require lower < upper" + ); + } + for (const auto& parameter : parameters) { + require( + std::abs(parameter.get()) <= std::numeric_limits::max(), + "Synthetic parameters must fit float32" + ); + } + const auto seed = + integer(source.at("seed"), "seed", 0, std::numeric_limits::max()); + input.source_type = "float32"; + input.generation = { + {"schema_version", 1}, + {"generator", "svs-mt19937-box-muller-v1"}, + {"vectors", options.total}, + {"dimensions", options.dimensions}, + {"data_type", input.source_type}, + {"distribution", distribution}, + {"parameters", parameters}, + {"seed", seed}}; + } + input.path = resolve(directory, source.at("path")); + protected_paths.insert(input.path); + if (input.kind == "synthetic") { + const auto [it, inserted] = generators.emplace(input.path, input.generation); + require( + inserted || it->second == input.generation, + "Conflicting synthetic configurations for " + input.path.string() + ); + } + config.inputs.emplace_back(std::move(options), std::move(input)); + } + for (const auto& [options, input] : config.inputs) { + (void)options; + if (input.kind == "synthetic") { + for (const fs::path& sidecar : + {metadata_path(input), fs::path{input.path.string() + ".lock"}}) { + require( + !protected_paths.contains(sidecar), + "Synthetic sidecar collides with a dataset or configuration: " + + sidecar.string() + ); + } + } + } + for (const auto& [options, input] : config.inputs) { + (void)options; + if (input.kind == "synthetic") { + protected_paths.insert(metadata_path(input)); + protected_paths.insert(input.path.string() + ".lock"); + } + } + for (const auto& output : {config.output_json, config.output_log}) { + require( + output.empty() || !protected_paths.contains(output), + "Output path collides with an input: " + output.string() + ); + if (!output.empty() && fs::exists(output)) { + for (const auto& input : protected_paths) { + require( + !fs::exists(input) || !fs::equivalent(input, output), + "Output path aliases an input: " + output.string() + ); + } + } + } + require(config.output_json != config.output_log, "JSON and log paths must differ"); + require( + config.output_log.empty() || !fs::exists(config.output_json) || + !fs::exists(config.output_log) || + !fs::equivalent(config.output_json, config.output_log), + "JSON and log paths must not alias the same file" + ); + return config; +} + +inline void redirect_stream(FILE* stream, const fs::path& path) { + fs::create_directories(path.parent_path()); + const int fd = open(path.c_str(), O_CREAT | O_WRONLY | O_TRUNC, 0666); + require(fd >= 0, "Cannot open output: " + path.string()); + const bool success = std::fflush(stream) == 0 && dup2(fd, fileno(stream)) >= 0; + if (fd != fileno(stream)) { + close(fd); + } + require(success, "Cannot redirect output: " + path.string()); +} + +class Sha256 { + std::unique_ptr context_{ + EVP_MD_CTX_new(), EVP_MD_CTX_free}; + + public: + Sha256() { + require( + context_ && EVP_DigestInit_ex(context_.get(), EVP_sha256(), nullptr) == 1, + "Cannot initialize SHA-256" + ); + } + void update(const void* data, size_t size) { + require(EVP_DigestUpdate(context_.get(), data, size) == 1, "SHA-256 update failed"); + } + std::string finish() { + std::array digest{}; + unsigned int size = 0; + require( + EVP_DigestFinal_ex(context_.get(), digest.data(), &size) == 1, + "SHA-256 finalization failed" + ); + std::string result; + for (unsigned int i = 0; i < size; ++i) { + result += fmt::format("{:02x}", digest[i]); + } + return result; + } +}; + +inline std::string hash_prefix(const fs::path& path, size_t bytes) { + std::ifstream stream{path, std::ios::binary}; + Sha256 hash; + std::vector buffer(1024 * 1024); + while (bytes != 0) { + const auto count = std::min(bytes, buffer.size()); + require( + bool(stream.read(buffer.data(), count)), + "Cannot read dataset for checksum: " + path.string() + ); + hash.update(buffer.data(), count); + bytes -= count; + } + return hash.finish(); +} + +class TemporaryFile { + public: + fs::path path; + explicit TemporaryFile(const fs::path& destination) { + fs::create_directories(destination.parent_path()); + auto pattern = destination.string() + ".tmp.XXXXXX"; + const int fd = mkstemp(pattern.data()); + require(fd >= 0, "Cannot create temporary file for " + destination.string()); + close(fd); + path = pattern; + } + TemporaryFile(const TemporaryFile&) = delete; + TemporaryFile& operator=(const TemporaryFile&) = delete; + ~TemporaryFile() { + std::error_code ignored; + if (!path.empty()) { + fs::remove(path, ignored); + } + } + void commit(const fs::path& destination) { + fs::rename(path, destination); + path.clear(); + } +}; + +class DatasetLock { + int fd_; + + public: + explicit DatasetLock(const fs::path& path) + : fd_{open(path.c_str(), O_CREAT | O_RDWR, 0600)} { + require(fd_ >= 0, "Cannot open dataset lock: " + path.string()); + int result; + do { + result = flock(fd_, LOCK_EX); + } while (result != 0 && errno == EINTR); + if (result != 0) { + close(fd_); + throw std::runtime_error("Cannot lock dataset: " + path.string()); + } + } + DatasetLock(const DatasetLock&) = delete; + DatasetLock& operator=(const DatasetLock&) = delete; + ~DatasetLock() { close(fd_); } +}; + +inline void prepare_synthetic(const Options& options, Input& input) { + static_assert(std::endian::native == std::endian::little); + fs::create_directories(input.path.parent_path()); + DatasetLock lock{input.path.string() + ".lock"}; + const auto metadata = metadata_path(input); + const auto bytes = options.total * (4 + sizeof(float) * options.dimensions); + if (fs::exists(input.path) || fs::exists(metadata)) { + try { + require( + fs::is_regular_file(input.path) && fs::is_regular_file(metadata), + "Incomplete synthetic dataset: both data and metadata must exist" + ); + auto saved = read_json(metadata); + require( + saved.is_object() && saved.contains("sha256"), "Missing synthetic SHA-256" + ); + const auto expected_hash = text(saved.at("sha256"), "synthetic sha256"); + saved.erase("sha256"); + require(saved == input.generation, "Synthetic metadata mismatch"); + require(fs::file_size(input.path) == bytes, "Synthetic dataset size mismatch"); + input.sha256 = hash_prefix(input.path, bytes); + require(input.sha256 == expected_hash, "Synthetic dataset checksum mismatch"); + fmt::print(stderr, "Reusing synthetic dataset: {}\n", input.path.string()); + return; + } catch (const std::exception& error) { + // An obsolete or incomplete cache is recoverable. Keep the old files + // until the replacement data and metadata have both been generated. + fmt::print( + stderr, + "Regenerating synthetic dataset: {} ({})\n", + input.path.string(), + error.what() + ); + } + } else { + fmt::print(stderr, "Generating synthetic dataset: {}\n", input.path.string()); + } + TemporaryFile vectors{input.path}; + TemporaryFile sidecar{metadata}; + std::ofstream stream{vectors.path, std::ios::binary}; + stream.exceptions(std::ios::failbit | std::ios::badbit); + const auto dimensions = static_cast(options.dimensions); + std::vector row(options.dimensions); + std::mt19937 generator{input.generation.at("seed").get()}; + const auto& parameters = input.generation.at("parameters"); + const bool normal = input.generation.at("distribution") == "normal"; + const double offset = parameters.at(normal ? "mean" : "lower").get(); + const double scale = normal ? parameters.at("stddev").get() + : parameters.at("upper").get() - offset; + bool has_spare = false; + double spare = 0; + Sha256 hash; + for (size_t i = 0; i < options.total; ++i) { + for (float& coordinate : row) { + double sample; + if (!normal) { + sample = static_cast(generator()) / 4294967296.0; + } else if (has_spare) { + sample = spare; + has_spare = false; + } else { + // Explicit transform avoids std::normal_distribution's implementation- + // dependent algorithm. Saved files remain authoritative across launches. + const double u = (static_cast(generator()) + 1) / 4294967297.0; + const double v = static_cast(generator()) / 4294967296.0; + const double radius = std::sqrt(-2 * std::log(u)); + const double angle = 2 * std::numbers::pi_v * v; + sample = radius * std::cos(angle); + spare = radius * std::sin(angle); + has_spare = true; + } + const double value = offset + scale * sample; + require( + std::isfinite(value) && + std::abs(value) <= std::numeric_limits::max(), + "Synthetic coordinate is outside finite float32 range" + ); + coordinate = static_cast(value); + } + stream.write(reinterpret_cast(&dimensions), sizeof(dimensions)); + stream.write(reinterpret_cast(row.data()), row.size() * sizeof(float)); + hash.update(&dimensions, sizeof(dimensions)); + hash.update(row.data(), row.size() * sizeof(float)); + } + stream.close(); + input.sha256 = hash.finish(); + auto saved = input.generation; + saved["sha256"] = input.sha256; + std::ofstream metadata_stream{sidecar.path}; + metadata_stream.exceptions(std::ios::failbit | std::ios::badbit); + metadata_stream << saved.dump(2) << '\n'; + metadata_stream.close(); + vectors.commit(input.path); + sidecar.commit(metadata); +} + +inline void inspect_file(const Options& options, Input& input) { + static_assert(std::endian::native == std::endian::little); + const size_t row_bytes = 4 + element_size(input.source_type) * options.dimensions; + std::ifstream stream{input.path, std::ios::binary}; + uint32_t dimensions = 0; + require( + bool(stream.read(reinterpret_cast(&dimensions), sizeof(dimensions))), + "Cannot read dataset header: " + input.path.string() + ); + const auto bytes = fs::file_size(input.path); + require( + dimensions == options.dimensions && bytes % row_bytes == 0, + "Invalid vecs file dimensions or length: " + input.path.string() + ); + input.source_vectors = bytes / row_bytes; + require(input.source_vectors >= options.total, "Dataset has fewer rows than requested"); + if (input.sha256.empty()) { + input.sha256 = hash_prefix(input.path, options.total * row_bytes); + } +} + +inline void prepare_inputs(Configuration& config) { + // Preflight file inputs before spending time generating synthetic datasets. + for (auto& [options, input] : config.inputs) { + if (input.kind == "file") { + inspect_file(options, input); + } + } + for (auto& [options, input] : config.inputs) { + if (input.kind == "synthetic") { + prepare_synthetic(options, input); + inspect_file(options, input); + } + } +} + +} // namespace graph_metrics From 9430c9ec76d4694eea86c4d5493543d38fb213d9 Mon Sep 17 00:00:00 2001 From: Dmitry Razdoburdin Date: Fri, 2 Oct 2026 05:08:57 -0700 Subject: [PATCH 2/4] fix --- .github/graph-metrics/CMakeLists.txt | 16 +++++++++++++++- .github/graph-metrics/ci.py | 9 ++++++++- .github/graph-metrics/graph_metric.sh | 13 ++++++++----- utils/CMakeLists.txt | 5 ----- 4 files changed, 31 insertions(+), 12 deletions(-) diff --git a/.github/graph-metrics/CMakeLists.txt b/.github/graph-metrics/CMakeLists.txt index ccc6c4ed9..1be7c4592 100644 --- a/.github/graph-metrics/CMakeLists.txt +++ b/.github/graph-metrics/CMakeLists.txt @@ -25,7 +25,21 @@ set(SVS_BUILD_BENCHMARK_TEST_GENERATORS OFF CACHE BOOL "" FORCE) add_subdirectory("${SVS_SOURCE_DIR}" svs) get_filename_component(METRICS_SOURCE_DIR "${CMAKE_CURRENT_LIST_DIR}/../.." ABSOLUTE) -include("${METRICS_SOURCE_DIR}/cmake/graph-metrics-dependencies.cmake") +# These dependencies are private to this standalone project. Normal SVS utility +# builds must not discover or download them. +find_package(OpenSSL REQUIRED COMPONENTS Crypto) +find_package(nlohmann_json 3.11 QUIET) +if(NOT nlohmann_json_FOUND) + include(FetchContent) + FetchContent_Declare( + graph_metrics_json + GIT_REPOSITORY https://github.com/nlohmann/json.git + GIT_TAG v3.11.3 + GIT_SHALLOW TRUE + ) + FetchContent_MakeAvailable(graph_metrics_json) +endif() + add_executable( graph_metrics "${METRICS_SOURCE_DIR}/utils/graph_metrics.cpp" ) diff --git a/.github/graph-metrics/ci.py b/.github/graph-metrics/ci.py index 72116f205..915a2c57c 100644 --- a/.github/graph-metrics/ci.py +++ b/.github/graph-metrics/ci.py @@ -25,6 +25,7 @@ import shlex import shutil import subprocess +import sys import urllib.request import uuid @@ -76,7 +77,13 @@ def logged(args, log, cwd=None): with log.open("a") as stream: stream.write("$ " + shlex.join([str(arg) for arg in args]) + "\n") stream.flush() - subprocess.run(args, cwd=cwd, stdout=stream, stderr=subprocess.STDOUT, check=True) + try: + subprocess.run(args, cwd=cwd, stdout=stream, stderr=subprocess.STDOUT, check=True) + except subprocess.CalledProcessError: + stream.flush() + print(f"Command failed; diagnostic log: {log}", file=sys.stderr) + print(log.read_text(errors="replace"), file=sys.stderr) + raise def configure(source, build, log, extra_args): diff --git a/.github/graph-metrics/graph_metric.sh b/.github/graph-metrics/graph_metric.sh index c9a786492..6d7259877 100644 --- a/.github/graph-metrics/graph_metric.sh +++ b/.github/graph-metrics/graph_metric.sh @@ -19,11 +19,12 @@ graph_metric_root="$(cd -- "$(dirname -- "${BASH_SOURCE[0]}")/../.." && pwd)" # Dataset definitions, graph parameters, thread count, and report paths live in JSON. graph_metric_config="${1:-${graph_metric_root}/graph_metrics_config.json}" -graph_metric_build_dir="${graph_metric_root}/build" +graph_metric_output_dir="${graph_metric_root}/build/graph-metrics" +graph_metric_build_dir="${graph_metric_output_dir}/build" graph_metric_build_type="Release" graph_metric_build_jobs=8 graph_metric_cmake="auto" # Reuse the build directory's CMake, otherwise use PATH. -graph_metric_build_log="${graph_metric_build_dir}/graph-metrics/build.log" +graph_metric_build_log="${graph_metric_output_dir}/build.log" if (( $# > 1 )); then printf 'Usage: bash .github/graph-metrics/graph_metric.sh [configuration.json]\n' >&2 @@ -56,10 +57,12 @@ printf 'Configuration: %s\nBuild log: %s\n' "${graph_metric_config}" "${graph_me # The calculator handles runtime JSON/log destinations from the config itself. mkdir -p -- "$(dirname -- "${graph_metric_build_log}")" { - "${graph_metric_cmake}" -S "${graph_metric_root}" -B "${graph_metric_build_dir}" \ - -DCMAKE_BUILD_TYPE="${graph_metric_build_type}" -DSVS_BUILD_BINARIES=ON + "${graph_metric_cmake}" -S "${graph_metric_root}/.github/graph-metrics" \ + -B "${graph_metric_build_dir}" \ + -DCMAKE_BUILD_TYPE="${graph_metric_build_type}" \ + -DSVS_SOURCE_DIR="${graph_metric_root}" "${graph_metric_cmake}" --build "${graph_metric_build_dir}" --target graph_metrics \ --parallel "${graph_metric_build_jobs}" } > "${graph_metric_build_log}" 2>&1 -exec "${graph_metric_build_dir}/utils/graph_metrics" --config "${graph_metric_config}" +exec "${graph_metric_build_dir}/graph_metrics" --config "${graph_metric_config}" diff --git a/utils/CMakeLists.txt b/utils/CMakeLists.txt index c602a85dc..888b32f46 100644 --- a/utils/CMakeLists.txt +++ b/utils/CMakeLists.txt @@ -37,11 +37,6 @@ function(create_utility exe file) endfunction() create_utility(graph_stat graph_stat.cpp) -include("${PROJECT_SOURCE_DIR}/cmake/graph-metrics-dependencies.cmake") -create_utility(graph_metrics graph_metrics.cpp) -target_link_libraries( - graph_metrics PRIVATE nlohmann_json::nlohmann_json OpenSSL::Crypto -) # Legacy conversion routines. create_utility(convert_legacy convert_legacy.cpp) From 1fc80f8f01e98ff1da0f58a1c50f8ca29599874c Mon Sep 17 00:00:00 2001 From: Dmitry Razdoburdin Date: Fri, 2 Oct 2026 06:01:10 -0700 Subject: [PATCH 3/4] reduce dataset size, improve loging --- .github/graph-metrics/ci.py | 32 ++++++++++++------- .github/graph-metrics/graph_metric.sh | 0 .../graph_metrics_synthetic.json | 4 +-- 3 files changed, 23 insertions(+), 13 deletions(-) mode change 100644 => 100755 .github/graph-metrics/graph_metric.sh diff --git a/.github/graph-metrics/ci.py b/.github/graph-metrics/ci.py index 915a2c57c..d6815a07c 100644 --- a/.github/graph-metrics/ci.py +++ b/.github/graph-metrics/ci.py @@ -75,15 +75,21 @@ def configuration(head): def logged(args, log, cwd=None): with log.open("a") as stream: - stream.write("$ " + shlex.join([str(arg) for arg in args]) + "\n") - stream.flush() - try: - subprocess.run(args, cwd=cwd, stdout=stream, stderr=subprocess.STDOUT, check=True) - except subprocess.CalledProcessError: - stream.flush() - print(f"Command failed; diagnostic log: {log}", file=sys.stderr) - print(log.read_text(errors="replace"), file=sys.stderr) - raise + command_line = "$ " + shlex.join([str(arg) for arg in args]) + print(command_line, file=stream, flush=True) + print(command_line, flush=True) + with subprocess.Popen( + args, cwd=cwd, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, + text=True, errors="replace", bufsize=1, + ) as process: + for line in process.stdout: + stream.write(line) + stream.flush() + print(line, end="", flush=True) + returncode = process.wait() + if returncode: + print(f"Command failed; diagnostic log: {log}", file=sys.stderr, flush=True) + raise subprocess.CalledProcessError(returncode, args) def configure(source, build, log, extra_args): @@ -243,18 +249,22 @@ def calculate(args): # calculator, configuration and persisted vectors. for label, source in (("base", args.base), ("head", args.head)): build = args.work / label + print(f"{label}: configuring revision {record[f'{label}_sha']}", flush=True) configure(source, build, args.result / f"{label}-build.log", args.cmake_arg) + print(f"{label}: compiling {TARGET}", flush=True) logged([ os.environ.get("CMAKE", "cmake"), "--build", str(build), "--target", TARGET, "--parallel", str(cpu_count()), ], args.result / f"{label}-build.log") effective = copy.deepcopy(config) effective["output_json"] = str(args.result / f"{label}.json") - effective["output_log"] = str(args.result / f"{label}-metrics.log") + # Keep progress on stderr so logged() can stream it and save the artifact. + effective.pop("output_log", None) config_path = args.result / f"{label}-config.json" write_json(config_path, effective) + print(f"{label}: calculating graph metrics", flush=True) logged([str(build / TARGET), "--config", str(config_path)], - args.result / f"{label}-build.log") + args.result / f"{label}-metrics.log") reports = [json.loads((args.result / f"{label}.json").read_text()) for label in ("base", "head")] identities = [ {item["dataset_name"]: item["dataset_parameters"]["sha256"] for item in report["datasets"]} diff --git a/.github/graph-metrics/graph_metric.sh b/.github/graph-metrics/graph_metric.sh old mode 100644 new mode 100755 diff --git a/.github/graph-metrics/graph_metrics_synthetic.json b/.github/graph-metrics/graph_metrics_synthetic.json index 9df0a5f39..ac5216d44 100644 --- a/.github/graph-metrics/graph_metrics_synthetic.json +++ b/.github/graph-metrics/graph_metrics_synthetic.json @@ -9,7 +9,7 @@ "batch_size": 1000, "recompute_entry_point": true, "build_parameters": { - "degree": 64, + "degree": 48, "window": 512, "max_candidates": 750, "alpha": { @@ -21,7 +21,7 @@ "datasets": [ { "name": "synthetic-128", - "vectors": 500000, + "vectors": 100000, "dimensions": 128, "distance": "L2", "source": { From bdb461c23fefc03cb09b25f570a11ed636299631 Mon Sep 17 00:00:00 2001 From: Dmitry Razdoburdin Date: Fri, 2 Oct 2026 06:56:40 -0700 Subject: [PATCH 4/4] improve and simplify --- .github/graph-metrics/ci.py | 29 ++++++++++++-- .../graph_metrics_synthetic.json | 6 +-- .github/workflows/graph-metrics-synthetic.yml | 40 +++++++++++++++++++ 3 files changed, 69 insertions(+), 6 deletions(-) diff --git a/.github/graph-metrics/ci.py b/.github/graph-metrics/ci.py index d6815a07c..4fabafa97 100644 --- a/.github/graph-metrics/ci.py +++ b/.github/graph-metrics/ci.py @@ -406,15 +406,30 @@ def comparison(base, head): def comment(args): event = json.loads(Path(os.environ["GITHUB_EVENT_PATH"]).read_text()) - run = event["workflow_run"] - if run["event"] != "pull_request": + run = event.get("workflow_run") + if (run is not None and run["event"] != "pull_request") or ( + run is None and "pull_request" not in event + ): + print("No PR comment: this event is not a pull request.", flush=True) return api = GitHub() prefix = f"/repos/{api.repository}" - # Resolve the PR through GitHub's run/commit association, never an artifact's + if run is None: + pr = event["pull_request"] + if pr["head"]["repo"]["full_name"] != api.repository: + print("Fork PR comments are handled by the default-branch reporter.", flush=True) + return + run = api.request("GET", f"{prefix}/actions/runs/{int(os.environ['GITHUB_RUN_ID'])}") + # The current workflow has not finished, but both producer jobs have. + run["conclusion"] = os.environ["GRAPH_METRICS_CONCLUSION"] + run["pull_requests"] = [pr] + # Resolve the PR through GitHub's event/run association, never an artifact's # claimed PR number. Only current PR heads may update the comment. run_prs = run.get("pull_requests") or [] candidates = run_prs or api.pages(f"{prefix}/commits/{run['head_sha']}/pulls") + if not candidates: + print(f"No PR associated with graph-metrics run {run['id']}.", flush=True) + return run_heads = {item["number"]: item.get("head", {}).get("sha", run["head_sha"]) for item in run_prs} artifacts = api.pages(f"{prefix}/actions/runs/{run['id']}/artifacts", "artifacts") detection_path = args.artifacts / "graph-metrics-results-detection/detection.json" @@ -426,9 +441,11 @@ def comment(args): pr = api.request("GET", f"{prefix}/pulls/{int(candidate['number'])}") associated_sha = run_heads.get(candidate["number"], run["head_sha"]) if pr["state"] != "open" or pr["head"]["sha"] != associated_sha: + print(f"No comment on PR #{pr['number']}: closed or its head has changed.", flush=True) continue current_sha = pr["head"]["sha"] if any(record and record.get("head_sha") != current_sha for record in (detection, result)): + print(f"No comment on PR #{pr['number']}: artifact revision does not match.", flush=True) continue base_sha = (detection or result or {}).get("base_sha", "") if not re.fullmatch(r"[0-9a-f]{40}", base_sha): @@ -467,10 +484,16 @@ def comment(args): sequence = re.search(r"", existing["body"]) current = (int(run["id"]), int(run.get("run_attempt", 1))) if sequence and tuple(map(int, sequence.groups())) > current: + print(f"PR #{pr['number']} already has results from a newer run.", flush=True) + continue + if existing["body"] == body: + print(f"PR #{pr['number']} already has this graph-metrics comment.", flush=True) continue api.request("PATCH", f"{prefix}/issues/comments/{existing['id']}", {"body": body}) + print(f"Updated graph-metrics comment on PR #{pr['number']}.", flush=True) else: api.request("POST", f"{prefix}/issues/{pr['number']}/comments", {"body": body}) + print(f"Created graph-metrics comment on PR #{pr['number']}.", flush=True) def main(): diff --git a/.github/graph-metrics/graph_metrics_synthetic.json b/.github/graph-metrics/graph_metrics_synthetic.json index ac5216d44..e317a5016 100644 --- a/.github/graph-metrics/graph_metrics_synthetic.json +++ b/.github/graph-metrics/graph_metrics_synthetic.json @@ -9,7 +9,7 @@ "batch_size": 1000, "recompute_entry_point": true, "build_parameters": { - "degree": 48, + "degree": 32, "window": 512, "max_candidates": 750, "alpha": { @@ -21,8 +21,8 @@ "datasets": [ { "name": "synthetic-128", - "vectors": 100000, - "dimensions": 128, + "vectors": 50000, + "dimensions": 64, "distance": "L2", "source": { "type": "synthetic", diff --git a/.github/workflows/graph-metrics-synthetic.yml b/.github/workflows/graph-metrics-synthetic.yml index 8d1b4f50d..ef46f5996 100644 --- a/.github/workflows/graph-metrics-synthetic.yml +++ b/.github/workflows/graph-metrics-synthetic.yml @@ -145,3 +145,43 @@ jobs: name: graph-metrics-results-run path: results/run/ if-no-files-found: warn + + comment: + name: Comment PR graph metrics + needs: [detect, calculate] + # Same-repository PRs can report before this workflow is merged. Fork PRs + # use graph-metrics-comment.yml from the default branch. + if: >- + ${{ always() && github.event_name == 'pull_request' && + github.event.pull_request.head.repo.full_name == github.repository }} + runs-on: ubuntu-22.04 + timeout-minutes: 10 + permissions: + contents: read + actions: read + pull-requests: write + steps: + - uses: actions/checkout@v6 + with: + ref: ${{ github.event.pull_request.head.sha }} + persist-credentials: false + sparse-checkout: .github/graph-metrics + + - name: Download dependency decision and metrics + uses: actions/download-artifact@v8 + continue-on-error: true + with: + pattern: graph-metrics-results-* + path: ${{ runner.temp }}/graph-metrics-artifacts + + - name: Update the PR metrics or skipped message + env: + GITHUB_TOKEN: ${{ github.token }} + # This workflow is still running; use the completed producer jobs. + GRAPH_METRICS_CONCLUSION: >- + ${{ needs.detect.result == 'success' && + (needs.calculate.result == 'success' || needs.calculate.result == 'skipped') && + 'success' || 'failure' }} + run: | + python3 .github/graph-metrics/ci.py comment \ + --artifacts "${RUNNER_TEMP}/graph-metrics-artifacts"