Compare commits
16
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
b3a9874fc8 | ||
|
|
d657cbbf17 | ||
|
|
48534ef4de | ||
|
|
1116f514be | ||
|
|
d451e61749 | ||
|
|
3ff4a8d2d2 | ||
|
|
9343d4cdf4 | ||
|
|
66fb3d1e79 | ||
|
|
aca850cef2 | ||
|
|
1c79779956 | ||
|
|
eee03527ed | ||
|
|
1eb8541094 | ||
|
|
69c214d13a | ||
|
|
0341481aa7 | ||
|
|
d1c3fdd187 | ||
|
|
980e8d933e |
@@ -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": {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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())
|
||||
@@ -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()
|
||||
@@ -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,
|
||||
|
||||
@@ -124,6 +124,8 @@ test = [
|
||||
"av",
|
||||
"pytorch-msssim==1.0.0",
|
||||
"pytest",
|
||||
"pandas",
|
||||
"plotly",
|
||||
]
|
||||
|
||||
dev = [ "fastvideo[lint]", "fastvideo[test]", ]
|
||||
|
||||
Reference in New Issue
Block a user