Compare commits

...
16 Commits
Author SHA1 Message Date
Satyam Srivastava b3a9874fc8 Merge branch 'main' into test-hf-sync 2026-05-01 11:38:11 -07:00
Satyam Srivastava d657cbbf17 Update regression gatekeep to 5% 2026-05-01 11:05:33 -07:00
Satyam Srivastava 48534ef4de [ci] Use median instead of mean to detect regressions. 2026-04-28 01:19:19 -07:00
Satyam Srivastava 1116f514be Refactor upload performance metrics 2026-04-28 01:07:27 -07:00
Satyam Srivastava d451e61749 Abstract hf code and fix plot bug 2026-04-27 23:12:42 -07:00
Satyam Srivastava 3ff4a8d2d2 [ci] Refactor pr_test.sh post run hooks 2026-04-26 12:26:24 -07:00
Satyam Srivastava 9343d4cdf4 [bug] Fix modal remote functions crash container on sys.exit(0) during PR tests 2026-04-26 11:18:21 -07:00
Satyam Srivastava 66fb3d1e79 [ci] Fix buildkite agent annotate and upload 2026-04-26 02:36:45 -07:00
Satyam Srivastava aca850cef2 [ci] Fix plotly bug 2026-04-25 20:17:49 -07:00
Satyam Srivastava 1c79779956 Debug changes 2026-04-25 18:15:09 -07:00
Satyam Srivastava eee03527ed [ci] Add plot render for performance tracking 2026-04-25 17:39:41 -07:00
Satyam Srivastava 1eb8541094 Test with different run configs 2026-04-25 15:53:36 -07:00
Satyam Srivastava 69c214d13a Test: HF performance sync logic 2026-04-23 16:06:04 -07:00
Satyam Srivastava 0341481aa7 [ci] Upload Perf. Regression Results to HF_Repo 2026-04-23 15:51:26 -07:00
Satyam Srivastava d1c3fdd187 [ci] Perf CI Run Schedule
Make Full Performance CI run only on main branch in X days and not on every PR.
This schedule will be created in the Buildkite.
2026-04-22 21:12:54 -07:00
Satyam Srivastava 980e8d933e Add CI Performance Regression Tracking Changes 2026-04-22 18:11:13 -07:00
8 changed files with 781 additions and 7 deletions
@@ -29,8 +29,8 @@
"Will Smith casually eats noodles, his relaxed demeanor contrasting with the energetic background of a bustling street food market. The scene captures a mix of humor and authenticity. Mid-shot framing, vibrant lighting."
],
"run_config": {
"num_warmup_runs": 1,
"num_measurement_runs": 3,
"num_warmup_runs": 2,
"num_measurement_runs": 5,
"required_gpus": 2
},
"thresholds": {
+73 -2
View File
@@ -63,7 +63,72 @@ EFFECTIVE_PR=${BUILDKITE_PULL_REQUEST:-false}
if [ "$EFFECTIVE_PR" = "false" ] && [ -n "${PR_NUMBER:-}" ]; then
EFFECTIVE_PR=$PR_NUMBER
fi
MODAL_ENV="BUILDKITE_REPO=$BUILDKITE_REPO BUILDKITE_COMMIT=$BUILDKITE_COMMIT BUILDKITE_PULL_REQUEST=$EFFECTIVE_PR IMAGE_VERSION=$IMAGE_VERSION"
MODAL_ENV="BUILDKITE_REPO=$BUILDKITE_REPO BUILDKITE_COMMIT=$BUILDKITE_COMMIT BUILDKITE_PULL_REQUEST=$EFFECTIVE_PR BUILDKITE_BRANCH=${BUILDKITE_BRANCH:-} TEST_SCOPE=${TEST_SCOPE:-} IMAGE_VERSION=$IMAGE_VERSION"
POST_RUN_HOOK=""
upload_performance_artifacts() {
SHORT_SHA=${BUILDKITE_COMMIT:0:7}
LOCAL_DIR="downloaded_reports"
_download_reports() {
log "Downloading perf_reports/ from Modal Volume..."
mkdir -p "$LOCAL_DIR"
if ! modal volume get hf-model-weights "perf_reports/" "$LOCAL_DIR"; then
log "Error: Failed to download perf_reports/ from Modal Volume."
return 1
fi
}
_upload_dashboard() {
local target
target=$(find "$LOCAL_DIR" -name "dashboard_*${SHORT_SHA}*" | head -n 1)
log "TARGET dashboard: '$target'"
if [ -n "$target" ]; then
log "Found dashboard: $target. Uploading to Buildkite..."
buildkite-agent artifact upload "$target"
buildkite-agent annotate --style info --context "perf-dashboard" < "$target"
else
log "Warning: Could not find a dashboard file matching $SHORT_SHA"
fi
}
_upload_perf_summary() {
local target
target=$(find "$LOCAL_DIR" -name "perf_*${SHORT_SHA}*" | head -n 1)
log "TARGET perf summary: '$target'"
if [ -n "$target" ]; then
log "Found perf summary: $target. Uploading to Buildkite..."
buildkite-agent artifact upload "$target"
buildkite-agent annotate --style info --context "perf-summary" < "$target"
else
log "Warning: Could not find a perf summary file matching $SHORT_SHA"
fi
}
_cleanup_modal_volume() {
log "Cleaning up perf_reports/ from Modal Volume..."
if modal volume rm hf-model-weights "perf_reports/" --recursive; then
log "Successfully deleted perf_reports/ from Modal Volume."
else
log "Warning: Failed to delete perf_reports/ from Modal Volume. Manual cleanup may be required."
fi
}
_cleanup_local() {
log "Cleaning up local download directory..."
rm -rf "$LOCAL_DIR"
}
# --- Main flow ---
_download_reports || { _cleanup_local; return 1; }
_upload_dashboard
_upload_perf_summary
_cleanup_modal_volume
_cleanup_local
}
case "$TEST_TYPE" in
"encoder")
@@ -124,8 +189,9 @@ case "$TEST_TYPE" in
MODAL_COMMAND="$MODAL_ENV HF_API_KEY=$HF_API_KEY python3 -m modal run $MODAL_TEST_FILE::run_lora_extraction_tests"
;;
"performance")
log "Running performance tests..."
log "Running performance tests on Modal..."
MODAL_COMMAND="$MODAL_ENV HF_API_KEY=$HF_API_KEY python3 -m modal run $MODAL_TEST_FILE::run_performance_tests"
POST_RUN_HOOK="upload_performance_artifacts"
;;
"api_server")
log "Running API server integration tests..."
@@ -147,5 +213,10 @@ else
log "Error: Modal test failed with exit code: $TEST_EXIT_CODE"
fi
if [ -n "$POST_RUN_HOOK" ]; then
log "Executing post-run hook: $POST_RUN_HOOK"
"$POST_RUN_HOOK"
fi
log "=== Test execution completed with exit code: $TEST_EXIT_CODE ==="
exit $TEST_EXIT_CODE
+22 -3
View File
@@ -24,6 +24,10 @@ image = (modal.Image.from_registry(
os.environ.get("BUILDKITE_COMMIT", ""),
"BUILDKITE_PULL_REQUEST":
os.environ.get("BUILDKITE_PULL_REQUEST", ""),
"BUILDKITE_BRANCH":
os.environ.get("BUILDKITE_BRANCH", ""),
"TEST_SCOPE":
os.environ.get("TEST_SCOPE", ""),
"IMAGE_VERSION":
os.environ.get("IMAGE_VERSION", ""),
}))
@@ -66,6 +70,13 @@ def run_test(pytest_command: str):
{pytest_command}
"""
# result = subprocess.run(["/bin/bash", "-c", command],
# stdout=sys.stdout,
# stderr=sys.stderr,
# check=False)
# sys.exit(result.returncode)
result = subprocess.run(["/bin/bash", "-c", command],
stdout=sys.stdout,
stderr=sys.stderr,
@@ -229,12 +240,20 @@ def run_lora_extraction_tests():
timeout=1800,
secrets=[
modal.Secret.from_dict(
{"HF_API_KEY": os.environ.get("HF_API_KEY", "")})
{"HF_API_KEY": os.environ.get("HF_API_KEY", ""),
"HF_REPO_ID": "FastVideo/performance-tracking"})
],
volumes={"/root/data": model_vol})
volumes={
"/root/data": model_vol,
})
def run_performance_tests():
run_test(
"export HF_HOME='/root/data/.cache' && hf auth login --token $HF_API_KEY && pytest ./fastvideo/tests/performance -vs"
"export HF_HOME='/root/data/.cache' && "
"export PERFORMANCE_TRACKING_ROOT='/tmp/perf-tracking' && "
"hf auth login --token $HF_API_KEY && "
"pytest ./fastvideo/tests/performance -vs && "
"python ./fastvideo/tests/performance/compare_baseline.py && "
"python ./fastvideo/tests/performance/dashboard.py"
)
@@ -0,0 +1,320 @@
# SPDX-License-Identifier: Apache-2.0
"""Track performance results and compare against historical baseline.
This script:
1) reads current benchmark results from fastvideo/tests/performance/results,
2) writes normalized tracking records to the Modal volume path,
3) compares each current record against the mean of up to 5 prior records,
4) exits non-zero if any metric regresses by more than 15%.
"""
import glob
import json
import os
import re
import statistics
import sys
from huggingface_hub import HfApi, snapshot_download
from datetime import datetime, timezone
from typing import Any
from hf_store import sync_from_hf, upload_record, load_records_for_model, sanitize, safe_float
# Use the env var passed by Modal, fallback to a default if needed
HF_REPO_ID = os.environ.get("HF_REPO_ID", "FastVideo/performance-tracking")
HF_TOKEN = os.environ.get("HF_API_KEY")
RESULTS_DIR = os.path.join(
os.path.dirname(os.path.abspath(__file__)),
"results",
)
TRACKING_ROOT = os.environ.get(
"PERFORMANCE_TRACKING_ROOT",
"/tmp/perf-tracking",
)
MAX_REGRESSION = float(os.environ.get("PERF_MAX_REGRESSION", "0.05"))
def _should_persist_tracking() -> bool:
# test_scope = os.environ.get("TEST_SCOPE", "")
# branch = os.environ.get("BUILDKITE_BRANCH", "")
# return test_scope == "full" and branch == "main"
return True # only for testing purpose.
def _sanitize(value: str) -> str:
return re.sub(r"[^A-Za-z0-9._-]", "_", value)
def _safe_float(value: Any) -> float | None:
if value is None:
return None
try:
return float(value)
except (TypeError, ValueError):
return None
def _load_current_results() -> list[dict[str, Any]]:
pattern = os.path.join(RESULTS_DIR, "perf_*.json")
records: list[dict[str, Any]] = []
for path in sorted(glob.glob(pattern)):
with open(path, encoding="utf-8") as f:
records.append(json.load(f))
return records
def _normalize_record(result: dict[str, Any]) -> dict[str, Any]:
benchmark_id = result.get("benchmark_id", "unknown")
model_id = benchmark_id
timestamp = result.get("timestamp")
if not timestamp:
timestamp = datetime.now(timezone.utc).isoformat()
commit_sha = result.get("commit") or os.environ.get("BUILDKITE_COMMIT", "")
latency = _safe_float(result.get("avg_generation_time_s"))
throughput = _safe_float(result.get("throughput_fps"))
memory = _safe_float(result.get("max_peak_memory_mb"))
return {
"model_id": model_id,
"timestamp": timestamp,
"commit_sha": commit_sha,
"gpu_type": result.get("device", "unknown"),
"latency": latency,
"throughput": throughput,
"memory": memory,
"success": True,
}
def _write_tracking_record(record: dict[str, Any]) -> str:
model_dir = os.path.join(TRACKING_ROOT, _sanitize(record["model_id"]))
os.makedirs(model_dir, exist_ok=True)
timestamp = _sanitize(record["timestamp"])
commit = _sanitize(record["commit_sha"] or "unknown")
out_path = os.path.join(model_dir, f"{timestamp}_{commit}.json")
with open(out_path, "w", encoding="utf-8") as f:
json.dump(record, f, indent=2)
return out_path
def _baseline_metric(records: list[dict[str, Any]], key: str) -> float | None:
values = [
_safe_float(r.get(key))
for r in records
]
values = [v for v in values if v is not None]
if not values:
return None
return statistics.median(values)
def _check_regressions(
current: dict[str, Any],
baseline_records: list[dict[str, Any]],
max_regression: float,
) -> list[str]:
failures: list[str] = []
for metric in ("latency", "memory"):
baseline = _baseline_metric(baseline_records, metric)
curr = _safe_float(current.get(metric))
if baseline is None or curr is None or baseline <= 0:
continue
regression = (curr - baseline) / baseline
if regression > max_regression:
failures.append(
f"{current['model_id']} {metric} regressed by {regression * 100:.1f}% "
f"(current={curr:.3f}, baseline_median={baseline:.3f})"
)
baseline_tp = _baseline_metric(baseline_records, "throughput")
curr_tp = _safe_float(current.get("throughput"))
if baseline_tp is not None and curr_tp is not None and baseline_tp > 0:
regression = (baseline_tp - curr_tp) / baseline_tp
if regression > max_regression:
failures.append(
f"{current['model_id']} throughput regressed by {regression * 100:.1f}% "
f"(current={curr_tp:.3f}, baseline_median={baseline_tp:.3f})"
)
return failures
def _metric_delta_percent(
metric: str,
current: dict[str, Any],
baseline_records: list[dict[str, Any]],
) -> float | None:
curr = _safe_float(current.get(metric))
baseline = _baseline_metric(baseline_records, metric)
if curr is None or baseline is None or baseline <= 0:
return None
if metric in ("latency", "memory"):
return (curr - baseline) / baseline * 100.0
if metric == "throughput":
return (baseline - curr) / baseline * 100.0
return None
def _compact_value(value: float | None, precision: int = 3) -> str:
if value is None:
return "n/a"
return f"{value:.{precision}f}"
def _build_summary_row(
record: dict[str, Any],
baseline_records: list[dict[str, Any]],
has_failed: bool
) -> dict[str, Any]:
"""Formats a single benchmark result into a row for the Markdown summary table."""
latency_base = _safe_float(_baseline_metric(baseline_records, "latency"))
throughput_base = _safe_float(_baseline_metric(baseline_records, "throughput"))
memory_base = _safe_float(_baseline_metric(baseline_records, "memory"))
# Calculate percentages for the 'Worst Regression' column
latency_reg = _metric_delta_percent("latency", record, baseline_records)
throughput_reg = _metric_delta_percent("throughput", record, baseline_records)
memory_reg = _metric_delta_percent("memory", record, baseline_records)
regressions = [v for v in (latency_reg, throughput_reg, memory_reg) if v is not None]
worst_regression_pct = max(regressions) if regressions else None
return {
"model_id": record["model_id"],
"gpu_type": record["gpu_type"],
"baseline_n": len(baseline_records),
"latency_curr": _safe_float(record.get("latency")),
"latency_base": latency_base,
"throughput_curr": _safe_float(record.get("throughput")),
"throughput_base": throughput_base,
"memory_curr": _safe_float(record.get("memory")),
"memory_base": memory_base,
"worst_regression_pct": worst_regression_pct,
"failed": has_failed,
}
def _build_markdown_summary(
summary_rows: list[dict[str, Any]],
max_regression: float,
) -> str:
lines = [
"## Performance Baseline Comparison",
"",
f"Threshold: regressions greater than {max_regression * 100:.1f}% fail",
"",
"| Model | GPU | Baseline N | Latency (curr/base) | Throughput (curr/base) | Memory (curr/base) | Worst Regression | Status |",
"|---|---|---:|---|---|---|---:|---|",
]
for row in summary_rows:
latency = f"{_compact_value(row['latency_curr'])} / {_compact_value(row['latency_base'])}"
throughput = f"{_compact_value(row['throughput_curr'])} / {_compact_value(row['throughput_base'])}"
memory = f"{_compact_value(row['memory_curr'], 1)} / {_compact_value(row['memory_base'], 1)}"
worst_reg = "n/a" if row["worst_regression_pct"] is None else f"{row['worst_regression_pct']:.1f}%"
status = "FAIL" if row["failed"] else "PASS"
lines.append(
f"| {row['model_id']} | {row['gpu_type']} | {row['baseline_n']} | "
f"{latency} | {throughput} | {memory} | {worst_reg} | {status} |"
)
return "\n".join(lines) + "\n"
def _emit_markdown_summary(markdown: str, commit_sha: str) -> None:
print("\n" + markdown)
# 1. Existing GitHub logic (safe to keep)
summary_path = os.environ.get("GITHUB_STEP_SUMMARY")
if summary_path:
with open(summary_path, "a", encoding="utf-8") as f:
f.write(markdown + "\n")
# 2. Write to Modal volume for Buildkite to pick up in post-run hook
try:
perf_reports_dir = "/root/data/perf_reports"
os.makedirs(perf_reports_dir, exist_ok=True)
short_sha = commit_sha[:7] if commit_sha else "unknown"
timestamp = datetime.now().strftime("%Y%m%d_%H%M%S")
report_path = os.path.join(perf_reports_dir, f"perf_{short_sha}_{timestamp}.md")
with open(report_path, "w", encoding="utf-8") as f:
f.write(markdown + "\n")
print(f"Performance report written to {report_path}")
except Exception as e:
print(f"Failed to write performance report to Modal volume: {e}")
def main() -> int:
# Pull the current state of the world from HF
sync_from_hf(TRACKING_ROOT)
current_results = _load_current_results()
if not current_results:
print(f"No performance result files found in {RESULTS_DIR}")
return 0
all_failures = []
summary_rows = []
persist_tracking = _should_persist_tracking()
if persist_tracking:
print("Tracking persistence enabled: full-suite run on main branch")
else:
print("Tracking persistence disabled: only full-suite runs on main branch are persisted")
for raw in current_results:
record = _normalize_record(raw)
baseline_records = load_records_for_model(
TRACKING_ROOT, record["model_id"], record["gpu_type"],
last_n=5, successful_only=True
)
failures = _check_regressions(record, baseline_records, MAX_REGRESSION)
# Tag the current record based on the failure.
if not baseline_records:
# INITIALIZATION CASE: First run for this model/GPU
print(f"No baseline for {record['model_id']} on {record['gpu_type']}. Initializing...")
failures = []
record["success"] = True # The first run is always "successful"
else:
# COMPARISON CASE: Compare against the mean of the last 5 good runs
failures = _check_regressions(record, baseline_records, MAX_REGRESSION)
if failures:
record["success"] = False
all_failures.extend(failures)
else:
record["success"] = True
# 5. Persist to HF if we are on main
if persist_tracking:
# This writes the JSON with the "success" field to /tmp
current_path = _write_tracking_record(record)
# This pushes it to the FastVideo/performance-tracking repo
upload_record(current_path, record)
summary_row = _build_summary_row(record, baseline_records, bool(failures))
summary_rows.append(summary_row)
commit_sha = os.environ.get("BUILDKITE_COMMIT", "unknown")[:7]
markdown = _build_markdown_summary(summary_rows, MAX_REGRESSION)
_emit_markdown_summary(markdown, commit_sha)
if all_failures:
print("Performance regression check failed:")
for item in all_failures:
print(f" - {item}")
return 1
print("Performance baseline comparison passed")
return 0
if __name__ == "__main__":
sys.exit(main())
+102
View File
@@ -0,0 +1,102 @@
# SPDX-License-Identifier: Apache-2.0
import os
import shutil
import subprocess
from datetime import datetime
import plotly.express as px
import pandas as pd
from hf_store import sync_from_hf, load_as_dataframe
# -----------------------------
# 1. Grouping
# -----------------------------
def group_data(df: pd.DataFrame):
# Group only by model+GPU so each group produces a time-series line.
# config_id (commit SHA) is carried as a column for hover/color use.
keys = ["model_id", "gpu_type"]
return df.groupby(keys, dropna=False)
# -----------------------------
# 2. Plot builder
# -----------------------------
def build_plots(df: pd.DataFrame) -> list:
figs = []
for (model_id, gpu_type), g in group_data(df):
g = g.sort_values("timestamp")
# One chart per metric so the y-axes aren't on wildly different scales
for metric in ("latency", "throughput", "memory"):
if g[metric].isna().all():
continue
fig = px.line(
g,
x="timestamp",
y=metric,
markers=True,
hover_data=["config_id", "commit_sha"],
title=f"{model_id} | {gpu_type} | {metric}",
labels={"timestamp": "Time", metric: metric},
)
figs.append(fig)
return figs
# -----------------------------
# 3. Render HTML dashboard
# -----------------------------
def render_html(figs: list, days: int) -> str:
html_parts = [
"<html>",
"<head><meta charset='utf-8'>",
"<style>body { font-family: sans-serif; margin: 2rem; }</style>",
"</head><body>",
f"<h2>Performance Dashboard (last {days} days)</h2>",
]
for fig in figs:
html_parts.append(fig.to_html(full_html=False, include_plotlyjs="cdn"))
html_parts.append("</body></html>")
return "\n".join(html_parts)
# -----------------------------
# 5. Main
# -----------------------------
def main() -> None:
days = int(os.environ.get("DASHBOARD_DAYS", "30"))
local_dir = sync_from_hf("/tmp/perf-tracking")
df = load_as_dataframe(local_dir, days=days)
if df.empty:
print("No data found")
return
# Sanity-check: log what we actually loaded
print(f"Loaded {len(df)} records across {df['model_id'].nunique()} model(s), "
f"{df['gpu_type'].nunique()} GPU type(s), "
f"date range: {df['timestamp'].min()} → {df['timestamp'].max()}")
figs = build_plots(df)
html = render_html(figs, days)
commit_sha = os.environ.get("BUILDKITE_COMMIT", "unknown")[:7]
timestamp = datetime.now().strftime("%Y%m%d_%H%M%S")
report_dir = "/root/data/perf_reports"
os.makedirs(report_dir, exist_ok=True)
filename = f"dashboard_{commit_sha}_{timestamp}.html"
output_file = os.path.join(report_dir, filename)
with open(output_file, "w", encoding="utf-8") as f:
f.write(html)
print(f"Dashboard generated: {output_file}")
if __name__ == "__main__":
main()
+254
View File
@@ -0,0 +1,254 @@
# SPDX-License-Identifier: Apache-2.0
"""Shared HuggingFace storage utilities for performance tracking.
Provides a single place for:
- Syncing the HF dataset repo to a local directory
- Loading raw JSON records (with optional recency filter)
- Loading records as a normalized pandas DataFrame
- Uploading individual result files back to HF
- Common helpers: sanitize, safe_float
"""
import glob
import json
import os
import re
from datetime import datetime, timedelta, timezone
from typing import Any
import pandas as pd
from huggingface_hub import HfApi, snapshot_download
# ---------------------------------------------------------------------------
# Configuration — read once at import time, shared across both consumers
# ---------------------------------------------------------------------------
HF_REPO_ID: str = os.environ.get("HF_REPO_ID", "FastVideo/performance-tracking")
HF_TOKEN: str | None = os.environ.get("HF_API_KEY")
# ---------------------------------------------------------------------------
# Low-level helpers
# ---------------------------------------------------------------------------
def sanitize(value: str) -> str:
"""Return a filesystem- and HF-path-safe version of *value*."""
return re.sub(r"[^A-Za-z0-9._-]", "_", value)
def safe_float(value: Any) -> float | None:
"""Coerce *value* to float, returning None on failure."""
if value is None:
return None
try:
return float(value)
except (TypeError, ValueError):
return None
# ---------------------------------------------------------------------------
# HF I/O
# ---------------------------------------------------------------------------
def sync_from_hf(local_dir: str) -> str:
"""Download the HF dataset repo snapshot to *local_dir*.
Returns *local_dir* so callers can chain: ``load_records(sync_from_hf(...))``.
On failure (empty repo, no credentials, network error) the function logs a
warning and returns *local_dir* unchanged so the caller can still work with
whatever is already on disk.
"""
if not HF_REPO_ID:
print("hf_store: HF_REPO_ID not set, skipping sync.")
return local_dir
print(f"hf_store: syncing from {HF_REPO_ID} → {local_dir}")
try:
snapshot_download(
repo_id=HF_REPO_ID,
repo_type="dataset",
local_dir=local_dir,
token=HF_TOKEN,
allow_patterns="*.json",
)
except Exception as exc:
print(f"hf_store: sync skipped — {exc}")
return local_dir
def upload_record(local_path: str, record: dict[str, Any]) -> None:
"""Upload *local_path* to the HF repo under ``<model_id>/<filename>``.
Silently skips if HF_TOKEN is absent so local/CI runs without credentials
don't crash.
"""
if not HF_TOKEN:
print("hf_store: HF_API_KEY not set, skipping upload.")
return
model_id = record.get("model_id", "unknown")
path_in_repo = f"{sanitize(model_id)}/{os.path.basename(local_path)}"
commit_sha = (record.get("commit_sha") or "unknown")[:7]
api = HfApi(token=HF_TOKEN)
try:
api.upload_file(
path_or_fileobj=local_path,
path_in_repo=path_in_repo,
repo_id=HF_REPO_ID,
repo_type="dataset",
commit_message=f"Perf: {model_id} at {commit_sha}",
)
print(f"hf_store: uploaded → {HF_REPO_ID}/{path_in_repo}")
except Exception as exc:
print(f"hf_store: upload failed — {exc}")
# ---------------------------------------------------------------------------
# Record loading
# ---------------------------------------------------------------------------
def load_records(
local_dir: str,
*,
days: int | None = None,
successful_only: bool = False,
) -> list[dict[str, Any]]:
"""Return raw JSON dicts from *local_dir*.
Args:
local_dir: Root directory previously populated by :func:`sync_from_hf`.
days: When set, discard records whose ``timestamp`` is older than this
many days. Records with a missing/unparseable timestamp are kept.
successful_only: When True, only records with ``success=True`` are
returned. Useful when building a regression baseline.
Returns:
List of raw dicts sorted by ``timestamp`` ascending (records that could
not be parsed are silently skipped).
"""
cutoff: datetime | None = None
if days is not None:
cutoff = datetime.now(timezone.utc) - timedelta(days=days)
records: list[dict[str, Any]] = []
for path in sorted(glob.glob(os.path.join(local_dir, "**", "*.json"), recursive=True)):
try:
with open(path, encoding="utf-8") as fh:
data: dict[str, Any] = json.load(fh)
except (OSError, json.JSONDecodeError):
continue
if successful_only and not data.get("success", True):
continue
if cutoff is not None:
raw_ts = data.get("timestamp")
if raw_ts:
try:
ts = datetime.fromisoformat(raw_ts)
if ts.tzinfo is None:
ts = ts.replace(tzinfo=timezone.utc)
if ts < cutoff:
continue
except ValueError:
pass # keep records with unparseable timestamps
records.append(data)
return records
def load_records_for_model(
local_dir: str,
model_id: str,
gpu_type: str | None = None,
*,
last_n: int | None = None,
successful_only: bool = True,
) -> list[dict[str, Any]]:
"""Return records for a specific *model_id*, optionally filtered by GPU.
Args:
local_dir: Root directory previously populated by :func:`sync_from_hf`.
model_id: Matches the ``model_id`` field inside each JSON record.
gpu_type: When set, only records whose ``gpu_type`` matches are returned.
last_n: When set, return only the most recent *n* records (after all
other filters). Useful for sliding-window baseline calculations.
successful_only: Passed through to :func:`load_records`.
Returns:
List of matching dicts sorted by timestamp ascending.
"""
model_dir = os.path.join(local_dir, sanitize(model_id))
if not os.path.isdir(model_dir):
return []
records = load_records(model_dir, successful_only=successful_only)
if gpu_type is not None:
records = [r for r in records if r.get("gpu_type") == gpu_type]
if last_n is not None:
records = records[-last_n:]
return records
# ---------------------------------------------------------------------------
# DataFrame helpers (dashboard / analytics consumers)
# ---------------------------------------------------------------------------
_NUMERIC_COLS = ("latency", "throughput", "memory")
def normalize_dataframe(df: pd.DataFrame) -> pd.DataFrame:
"""Apply standard type coercions to a raw records DataFrame.
- Parses ``timestamp`` to UTC-aware datetime.
- Coerces ``latency``, ``throughput``, ``memory`` to float.
- Adds a ``config_id`` column (first 7 chars of ``commit_sha``).
Returns the mutated DataFrame (also modifies in place for efficiency).
"""
if df.empty:
return df
df["timestamp"] = pd.to_datetime(df["timestamp"], utc=True, errors="coerce")
df["config_id"] = df.get("commit_sha", pd.Series(dtype=str)).fillna("unknown").str[:7]
for col in _NUMERIC_COLS:
if col in df.columns:
df[col] = pd.to_numeric(df[col], errors="coerce")
return df
def load_as_dataframe(
local_dir: str,
*,
days: int | None = None,
successful_only: bool = False,
) -> pd.DataFrame:
"""Load and normalize records from *local_dir* into a pandas DataFrame.
Combines :func:`load_records` + :func:`normalize_dataframe` into a single
call for consumers (e.g. the dashboard) that work exclusively with
DataFrames.
Args:
local_dir: Root directory previously populated by :func:`sync_from_hf`.
days: Passed through to :func:`load_records`.
successful_only: Passed through to :func:`load_records`.
Returns:
Normalized DataFrame, or an empty DataFrame if no records were found.
"""
records = load_records(local_dir, days=days, successful_only=successful_only)
if not records:
return pd.DataFrame()
df = pd.DataFrame(records)
return normalize_dataframe(df)
@@ -164,6 +164,10 @@ def test_inference_performance(cfg):
avg_time = sum(times) / len(times)
max_peak_memory = max(peak_memories)
device_name = torch.cuda.get_device_name()
num_frames = gen_kwargs.get("num_frames")
throughput_fps = (1.0 / avg_time) if avg_time > 0 else None
if isinstance(num_frames, (int, float)) and avg_time > 0:
throughput_fps = num_frames / avg_time
results = {
"benchmark_id": cfg["benchmark_id"],
@@ -174,6 +178,8 @@ def test_inference_performance(cfg):
"num_measurement_runs": num_measure,
"avg_generation_time_s": round(avg_time, 3),
"individual_times_s": [round(t, 3) for t in times],
"throughput_fps": round(throughput_fps, 3)
if throughput_fps is not None else None,
"max_peak_memory_mb": round(max_peak_memory, 1),
"individual_peak_memories_mb": [round(m, 1) for m in peak_memories],
"thresholds": thresholds,
+2
View File
@@ -124,6 +124,8 @@ test = [
"av",
"pytorch-msssim==1.0.0",
"pytest",
"pandas",
"plotly",
]
dev = [ "fastvideo[lint]", "fastvideo[test]", ]