mirror of
https://github.com/JustVugg/colibri.git
synced 2026-10-02 02:54:37 +08:00
feat(tools): add three-engine serving baseline campaigns
This commit is contained in:
@@ -0,0 +1,186 @@
|
||||
"""Incomplete or mismatched campaigns must not become performance rankings."""
|
||||
import contextlib
|
||||
import copy
|
||||
import io
|
||||
import json
|
||||
from pathlib import Path
|
||||
import tempfile
|
||||
import subprocess
|
||||
import sys
|
||||
from tests import test_benchmark_http_serving as http_fixture
|
||||
import unittest
|
||||
from unittest.mock import patch
|
||||
from tools import benchmark_baseline as baseline
|
||||
from tools import benchmark_http_serving as http
|
||||
|
||||
|
||||
class BaselineTest(unittest.TestCase):
|
||||
def setUp(self):
|
||||
temp = tempfile.TemporaryDirectory()
|
||||
self.addCleanup(temp.cleanup)
|
||||
self.root = Path(temp.name)
|
||||
example = Path(__file__).resolve().parents[2] / 'docs/baselines/three-engine.example.json'
|
||||
self.spec = json.loads(example.read_text())
|
||||
self.spec.update(experiment='fixture', cache_policy='warm', reasoning_policy='off',
|
||||
speculation_policy='off', quality_protocol='external fixture check')
|
||||
self.spec['hardware'] = dict.fromkeys(self.spec['hardware'], 'fixture host')
|
||||
sha = baseline.digest(example)
|
||||
self.spec['model'] = dict(source_revision='model-revision', tokenizer_sha256=sha,
|
||||
chat_template_sha256=sha)
|
||||
for engine in self.spec['engines'].values():
|
||||
engine.update(served_model='fixture', revision='engine-revision', launch_command='serve fixture',
|
||||
artifact_sha256=sha, weight_format='fixture', quantization='fp32')
|
||||
self.spec['matrix'].update(concurrency=[1, 2], rounds=2, repeats=2, warmup_requests=1)
|
||||
self.manifest = self.root / 'manifest.json'
|
||||
(self.root / 'workload.jsonl').write_text('{"messages":[{"role":"user","content":"fixture"}]}\n')
|
||||
self.save_manifest()
|
||||
self.workload, self.sha = http.load_workload(self.root / 'workload.jsonl')
|
||||
self.results = self.root / 'results'
|
||||
|
||||
def save_manifest(self):
|
||||
self.manifest.write_text(json.dumps(self.spec))
|
||||
|
||||
@staticmethod
|
||||
def fake_run(**kwargs):
|
||||
rows = [dict(index=i, success=True, start_seconds=0, duration_seconds=2,
|
||||
first_output_seconds=.5, completion_tokens=4, finish_reason='stop',
|
||||
error=None, http_status=200)
|
||||
for i in range(len(kwargs['workload']) * kwargs['repeats'])]
|
||||
return rows, http.summarize(rows, 3, kwargs.get('slo_first_output'), kwargs.get('slo_duration'))
|
||||
|
||||
def collect_all(self):
|
||||
with patch.object(http, 'run', side_effect=self.fake_run), contextlib.redirect_stdout(io.StringIO()):
|
||||
for item in baseline.plan(self.spec):
|
||||
baseline.collect(self.spec, self.workload, self.sha, item['engine'], item['round'], self.results)
|
||||
|
||||
def change_report(self, change):
|
||||
path = self.results / 'colibri-r1-c1.json'
|
||||
data = json.loads(path.read_text())
|
||||
change(data)
|
||||
path.write_text(json.dumps(data))
|
||||
|
||||
def compare(self):
|
||||
return baseline.compare(self.spec, self.workload, self.sha, self.results)
|
||||
|
||||
def test_relative_workload_and_rotation(self):
|
||||
spec, workload, sha = baseline.load_manifest(self.manifest)
|
||||
self.assertEqual((sha, workload), (self.sha, self.workload))
|
||||
self.assertEqual([x['engine'] for x in baseline.plan(spec)],
|
||||
['colibri', 'sglang', 'vllm', 'sglang', 'vllm', 'colibri'])
|
||||
|
||||
def test_artifact_mismatch_requires_deployment_mode(self):
|
||||
self.spec['engines']['vllm']['quantization'] = 'int4'
|
||||
self.spec['comparison'] = 'matched_artifact'
|
||||
self.save_manifest()
|
||||
with self.assertRaisesRegex(ValueError, 'identical'):
|
||||
baseline.load_manifest(self.manifest)
|
||||
self.spec['comparison'] = 'deployment'
|
||||
self.save_manifest()
|
||||
baseline.load_manifest(self.manifest)
|
||||
|
||||
def test_bad_matrix_and_placeholders(self):
|
||||
for key, value in [('concurrency', [True]), ('concurrency', [1, 1]),
|
||||
('concurrency', [16]), ('temperature', float('nan')), ('rounds', 0)]:
|
||||
spec = copy.deepcopy(self.spec)
|
||||
spec['matrix'][key] = value
|
||||
self.manifest.write_text(json.dumps(spec))
|
||||
with self.subTest(key=key, value=value), self.assertRaises(ValueError):
|
||||
baseline.load_manifest(self.manifest)
|
||||
self.spec['hardware']['gpu'] = 'REPLACE-gpu'
|
||||
self.save_manifest()
|
||||
with self.assertRaisesRegex(ValueError, 'gpu'):
|
||||
baseline.load_manifest(self.manifest)
|
||||
|
||||
def test_complete_report_boundaries(self):
|
||||
self.collect_all()
|
||||
report = self.compare()
|
||||
self.assertEqual(report['status'], 'protocol_complete')
|
||||
self.assertEqual(report['quality'], 'not_assessed')
|
||||
self.assertEqual(report['token_latency'], 'not_measured')
|
||||
cell = report['cells'][0]
|
||||
self.assertEqual(cell['aggregate_completion_tokens_per_second']['median'], 8 / 3)
|
||||
self.assertEqual(cell['reported_output_tokens']['count'], 4)
|
||||
self.assertEqual(cell['slo_goodput_requests_per_second']['median'], 2 / 3)
|
||||
|
||||
def test_missing_cell_is_not_dropped(self):
|
||||
self.collect_all()
|
||||
(self.results / 'colibri-r1-c1.json').unlink()
|
||||
report = self.compare()
|
||||
self.assertIn('missing colibri-r1-c1', report['issues'])
|
||||
self.assertIsNone(report['cells'][0]['aggregate_completion_tokens_per_second']['median'])
|
||||
self.assertEqual(report['cells'][0]['rounds_measured'], 1)
|
||||
|
||||
def test_mismatched_and_modified_reports(self):
|
||||
mutations = [lambda d: d.update(workload_sha256='different'),
|
||||
lambda d: d['manifest'].update(cache_policy='different'),
|
||||
lambda d: d.update(harness_sha256='different'),
|
||||
lambda d: d['summary'].update(successful_completion_tokens_per_second=999),
|
||||
lambda d: d['requests'].pop()]
|
||||
self.collect_all()
|
||||
path = self.results / 'colibri-r1-c1.json'
|
||||
original = path.read_text()
|
||||
for mutate in mutations:
|
||||
path.write_text(original)
|
||||
self.change_report(mutate)
|
||||
with self.subTest(mutation=mutate), self.assertRaises(ValueError):
|
||||
self.compare()
|
||||
|
||||
def test_failed_empty_and_missing_usage(self):
|
||||
self.collect_all()
|
||||
def mutate(data):
|
||||
data['requests'][0].update(success=False, error='http_error')
|
||||
data['requests'][1].update(completion_tokens=None, first_output_seconds=None)
|
||||
data['summary'] = http.summarize(data['requests'], 3, 5, 120)
|
||||
self.change_report(mutate)
|
||||
report = self.compare()
|
||||
self.assertEqual(len(report['issues']), 3)
|
||||
self.assertIsNone(report['cells'][0]['aggregate_completion_tokens_per_second']['median'])
|
||||
|
||||
def test_warmup_failure_and_no_overwrite(self):
|
||||
def fail(**kwargs):
|
||||
rows, _ = self.fake_run(**kwargs)
|
||||
rows[0]['success'] = False
|
||||
return rows, http.summarize(rows, 1)
|
||||
with patch.object(http, 'run', side_effect=fail) as run, contextlib.redirect_stdout(io.StringIO()):
|
||||
self.assertEqual(baseline.collect(self.spec, self.workload, self.sha,
|
||||
'colibri', 1, self.results), 1)
|
||||
self.assertEqual(run.call_count, 1)
|
||||
with self.assertRaisesRegex(ValueError, 'already'):
|
||||
baseline.collect(self.spec, self.workload, self.sha, 'colibri', 1, self.results)
|
||||
data = json.loads((self.results / 'colibri-r1-c1.json').read_text())
|
||||
self.assertEqual(data['status'], 'warmup_failed')
|
||||
self.assertEqual(data['requests'], [])
|
||||
|
||||
def test_secret_is_not_saved(self):
|
||||
with patch.dict('os.environ', {'OPENAI_API_KEY': 'fixture-secret'}), \
|
||||
patch.object(http, 'run', side_effect=self.fake_run) as run, \
|
||||
contextlib.redirect_stdout(io.StringIO()):
|
||||
baseline.collect(self.spec, self.workload, self.sha, 'colibri', 1, self.results)
|
||||
self.assertEqual(run.call_args.kwargs['key'], 'fixture-secret')
|
||||
for path in self.results.glob('*.json'):
|
||||
self.assertNotIn('fixture-secret', path.read_text())
|
||||
|
||||
def test_cli_collects_real_http_and_compare_reports_missing_cells(self):
|
||||
fixture = http_fixture.BenchmarkTest()
|
||||
fixture.setUp()
|
||||
self.addCleanup(fixture.tearDown)
|
||||
for config in self.spec['engines'].values():
|
||||
config['base_url'] = fixture.url.removesuffix('/chat/completions')
|
||||
self.spec['matrix'].update(rounds=1)
|
||||
self.save_manifest()
|
||||
command = [sys.executable, str(Path(baseline.__file__).resolve())]
|
||||
common = ['--manifest', str(self.manifest), '--results', str(self.results)]
|
||||
for engine in baseline.ENGINES:
|
||||
result = subprocess.run(command + ['run', '--engine', engine, '--round', '1'] + common,
|
||||
capture_output=True, text=True, timeout=15)
|
||||
self.assertEqual(result.returncode, 0, result.stderr)
|
||||
result = subprocess.run(command + ['compare'] + common, capture_output=True, text=True, timeout=15)
|
||||
self.assertEqual(result.returncode, 0, result.stderr)
|
||||
report = json.loads(result.stdout)
|
||||
self.assertEqual(report['status'], 'protocol_complete')
|
||||
self.assertEqual(len(report['sources']), 6)
|
||||
self.assertEqual(len(fixture.server.payloads), 18)
|
||||
(self.results / 'vllm-r1-c2.json').unlink()
|
||||
result = subprocess.run(command + ['compare'] + common, capture_output=True, text=True, timeout=15)
|
||||
self.assertEqual(result.returncode, 1, result.stderr)
|
||||
self.assertIn('missing vllm-r1-c2', json.loads(result.stdout)['issues'])
|
||||
@@ -0,0 +1,243 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Plan, collect and summarize a Colibri/SGLang/vLLM HTTP baseline."""
|
||||
import argparse
|
||||
import datetime
|
||||
import hashlib
|
||||
import json
|
||||
import math
|
||||
import os
|
||||
from pathlib import Path
|
||||
import statistics
|
||||
|
||||
try:
|
||||
from . import benchmark_http_serving as http
|
||||
except ImportError:
|
||||
import benchmark_http_serving as http
|
||||
|
||||
ENGINES = ("colibri", "sglang", "vllm")
|
||||
SCHEMA = "colibri.serving-baseline/v1"
|
||||
|
||||
|
||||
def digest(path):
|
||||
return hashlib.sha256(Path(path).read_bytes()).hexdigest()
|
||||
|
||||
|
||||
def read_json(path):
|
||||
return json.loads(Path(path).read_text(encoding="utf-8"))
|
||||
|
||||
|
||||
def require(condition, message):
|
||||
if not condition:
|
||||
raise ValueError(message)
|
||||
|
||||
|
||||
def text_fields(obj, names):
|
||||
require(isinstance(obj, dict), "expected an object")
|
||||
for name in names:
|
||||
require(isinstance(obj.get(name), str) and obj[name].strip()
|
||||
and not obj[name].startswith("REPLACE"), f"fill in {name}")
|
||||
|
||||
|
||||
def load_manifest(path):
|
||||
spec = read_json(path)
|
||||
require(spec.get("schema") == SCHEMA, "unsupported manifest schema")
|
||||
require(spec.get("comparison") in ("matched_artifact", "deployment"),
|
||||
"comparison must be matched_artifact or deployment")
|
||||
text_fields(spec, ("experiment", "workload", "cache_policy", "reasoning_policy",
|
||||
"speculation_policy", "quality_protocol", "residency"))
|
||||
require(spec["residency"] in ("fully_resident", "offload"), "invalid residency")
|
||||
text_fields(spec.get("hardware"), ("host", "cpu", "ram", "gpu", "storage", "os", "driver"))
|
||||
text_fields(spec.get("model"), ("source_revision", "tokenizer_sha256", "chat_template_sha256"))
|
||||
for key in ("tokenizer_sha256", "chat_template_sha256"):
|
||||
require(valid_hash(spec["model"][key]), f"invalid {key}")
|
||||
matrix = spec.get("matrix", {})
|
||||
require(set(matrix) == {"concurrency", "rounds", "repeats", "warmup_requests", "max_tokens",
|
||||
"temperature", "timeout", "slo_first_output", "slo_duration"},
|
||||
"matrix fields differ from the documented schema")
|
||||
require(isinstance(matrix["concurrency"], list) and matrix["concurrency"], "empty concurrency")
|
||||
values = matrix["concurrency"]
|
||||
require(all(type(n) is int and n > 0 for n in values) and len(set(values)) == len(values),
|
||||
"concurrency must contain unique positive integers")
|
||||
for key in ("rounds", "repeats", "max_tokens"):
|
||||
require(type(matrix[key]) is int and matrix[key] > 0, f"invalid {key}")
|
||||
require(type(matrix["warmup_requests"]) is int and matrix["warmup_requests"] >= 0,
|
||||
"invalid warmup_requests")
|
||||
for key in ("temperature", "timeout", "slo_first_output", "slo_duration"):
|
||||
value = matrix[key]
|
||||
require(type(value) in (int, float) and math.isfinite(value)
|
||||
and (value >= 0 if key == "temperature" else value > 0), f"invalid {key}")
|
||||
engines = spec.get("engines", {})
|
||||
require(set(engines) == set(ENGINES), "provide exactly colibri, sglang and vllm")
|
||||
for engine in engines.values():
|
||||
text_fields(engine, ("base_url", "served_model", "revision", "launch_command",
|
||||
"artifact_sha256", "weight_format", "quantization", "api_key_env"))
|
||||
http.endpoint(engine["base_url"])
|
||||
require(valid_hash(engine["artifact_sha256"]), "invalid artifact_sha256")
|
||||
if spec["comparison"] == "matched_artifact":
|
||||
identities = {(e["artifact_sha256"], e["weight_format"], e["quantization"])
|
||||
for e in engines.values()}
|
||||
require(len(identities) == 1, "matched_artifact requires identical weights/format/quantization")
|
||||
workload_path = Path(path).resolve().parent / spec["workload"]
|
||||
workload, workload_hash = http.load_workload(workload_path)
|
||||
require(len(workload) * matrix["repeats"] >= max(values),
|
||||
"request count must reach the largest concurrency")
|
||||
return spec, workload, workload_hash
|
||||
|
||||
|
||||
def valid_hash(value):
|
||||
return (isinstance(value, str) and len(value) == 64
|
||||
and all(c in "0123456789abcdef" for c in value) and len(set(value)) > 1)
|
||||
|
||||
|
||||
def plan(spec):
|
||||
for number in range(1, spec["matrix"]["rounds"] + 1):
|
||||
shift = (number - 1) % len(ENGINES)
|
||||
for engine in ENGINES[shift:] + ENGINES[:shift]:
|
||||
yield {"round": number, "engine": engine,
|
||||
"concurrency": spec["matrix"]["concurrency"]}
|
||||
|
||||
|
||||
def collect(spec, workload, workload_hash, engine, number, output):
|
||||
require(engine in ENGINES, "unknown engine")
|
||||
require(1 <= number <= spec["matrix"]["rounds"], "round outside matrix")
|
||||
output = Path(output)
|
||||
output.mkdir(parents=True, exist_ok=True)
|
||||
paths = [output / f"{engine}-r{number}-c{c}.json" for c in spec["matrix"]["concurrency"]]
|
||||
require(not any(p.exists() for p in paths), "round already has results; use a new output directory")
|
||||
config = spec["engines"][engine]
|
||||
key = os.environ.get(config["api_key_env"], "")
|
||||
matrix = spec["matrix"]
|
||||
failed = False
|
||||
for concurrency, path in zip(matrix["concurrency"], paths):
|
||||
kwargs = dict(url=http.endpoint(config["base_url"]), model=config["served_model"],
|
||||
concurrency=concurrency, max_tokens=matrix["max_tokens"],
|
||||
temperature=matrix["temperature"], key=key, timeout=matrix["timeout"])
|
||||
started = datetime.datetime.now(datetime.timezone.utc).isoformat()
|
||||
warmup = None
|
||||
if matrix["warmup_requests"]:
|
||||
prompts = [workload[i % len(workload)] for i in range(matrix["warmup_requests"])]
|
||||
rows, summary = http.run(workload=prompts, repeats=1, **kwargs)
|
||||
warmup = {"requests": rows, "summary": summary}
|
||||
rows, summary = [], None
|
||||
if warmup is None or warmup["summary"]["failed"] == 0:
|
||||
rows, summary = http.run(workload=workload, repeats=matrix["repeats"],
|
||||
slo_first_output=matrix["slo_first_output"],
|
||||
slo_duration=matrix["slo_duration"], **kwargs)
|
||||
report = {"schema": SCHEMA, "manifest": spec, "engine": engine, "round": number,
|
||||
"concurrency": concurrency, "workload_sha256": workload_hash,
|
||||
"harness_sha256": digest(http.__file__), "collector_sha256": digest(__file__),
|
||||
"started_at": started, "status": "warmup_failed" if summary is None else "measured",
|
||||
"warmup": warmup, "requests": rows, "summary": summary}
|
||||
with path.open("x", encoding="utf-8") as stream:
|
||||
stream.write(json.dumps(report, indent=2, allow_nan=False) + "\n")
|
||||
failed |= summary is None or summary["failed"] > 0
|
||||
print(path)
|
||||
if summary is None:
|
||||
break
|
||||
return int(failed)
|
||||
|
||||
|
||||
def spread(values):
|
||||
present = [v for v in values if v is not None]
|
||||
if len(present) != len(values):
|
||||
return {"count": len(present), "min": None, "median": None, "max": None}
|
||||
return {"count": len(values), "min": min(values),
|
||||
"median": statistics.median(values), "max": max(values)}
|
||||
|
||||
|
||||
def compare(spec, workload, workload_hash, directory):
|
||||
expected = {(e, r, c) for e in ENGINES for r in range(1, spec["matrix"]["rounds"] + 1)
|
||||
for c in spec["matrix"]["concurrency"]}
|
||||
reports, issues, fingerprints = {}, [], set()
|
||||
files = sorted(Path(directory).glob("*-r*-c*.json"))
|
||||
for path in files:
|
||||
report = read_json(path)
|
||||
require(report.get("schema") == SCHEMA, f"unsupported report: {path}")
|
||||
require(report.get("manifest") == spec, f"manifest mismatch: {path}")
|
||||
require(report.get("workload_sha256") == workload_hash, f"workload mismatch: {path}")
|
||||
identity = (report["engine"], report["round"], report["concurrency"])
|
||||
require(identity in expected and identity not in reports, f"unexpected/duplicate cell: {path}")
|
||||
fingerprints.add((report["harness_sha256"], report["collector_sha256"]))
|
||||
reports[identity] = report
|
||||
summary = report["summary"]
|
||||
if summary is None:
|
||||
issues.append(f"{path.name}: measurement absent")
|
||||
continue
|
||||
count = len(workload) * spec["matrix"]["repeats"]
|
||||
rows = report["requests"]
|
||||
require(len(rows) == count and [r["index"] for r in rows] == list(range(count)),
|
||||
f"request set mismatch: {path}")
|
||||
require(type(summary["wall_seconds"]) in (int, float)
|
||||
and math.isfinite(summary["wall_seconds"]) and summary["wall_seconds"] > 0,
|
||||
f"invalid wall time: {path}")
|
||||
recalculated = http.summarize(rows, summary["wall_seconds"],
|
||||
spec["matrix"]["slo_first_output"], spec["matrix"]["slo_duration"])
|
||||
require(summary == recalculated, f"summary differs from raw requests: {path}")
|
||||
if summary["failed"]:
|
||||
issues.append(f"{path.name}: failed requests")
|
||||
if any(r["completion_tokens"] is None for r in rows if r["success"]):
|
||||
issues.append(f"{path.name}: missing token usage")
|
||||
if any(r["success"] and (r["first_output_seconds"] is None or r["completion_tokens"] == 0
|
||||
or r["finish_reason"] == "content_filter") for r in rows):
|
||||
issues.append(f"{path.name}: empty or filtered output")
|
||||
require(len(fingerprints) <= 1, "mixed benchmark implementations")
|
||||
for e, r, c in sorted(expected - reports.keys()):
|
||||
issues.append(f"missing {e}-r{r}-c{c}")
|
||||
table = []
|
||||
for engine in ENGINES:
|
||||
for concurrency in spec["matrix"]["concurrency"]:
|
||||
cells = [reports.get((engine, r, concurrency))
|
||||
for r in range(1, spec["matrix"]["rounds"] + 1)]
|
||||
summaries = [cell["summary"] if cell else None for cell in cells]
|
||||
def metric(field, percentile=None):
|
||||
values = [s[field] if s else None for s in summaries]
|
||||
if percentile:
|
||||
values = [v[percentile] if v else None for v in values]
|
||||
return spread(values)
|
||||
measured = [cell for cell in cells if cell and cell["summary"]]
|
||||
lengths = [row["completion_tokens"] for cell in measured for row in cell["requests"]
|
||||
if row["success"] and row["completion_tokens"] is not None]
|
||||
table.append({"engine": engine, "concurrency": concurrency,
|
||||
"rounds_measured": len(measured),
|
||||
"aggregate_completion_tokens_per_second": metric("successful_completion_tokens_per_second"),
|
||||
"first_output_p50_seconds": metric("successful_first_output_seconds", "p50"),
|
||||
"first_output_p95_seconds": metric("successful_first_output_seconds", "p95"),
|
||||
"first_output_p99_seconds": metric("successful_first_output_seconds", "p99"),
|
||||
"duration_p95_seconds": metric("successful_duration_seconds", "p95"),
|
||||
"failure_rate": metric("failure_rate"),
|
||||
"slo_goodput_requests_per_second": metric("latency_slo", "goodput_requests_per_second"),
|
||||
"reported_output_tokens": http.distribution(lengths)})
|
||||
return {"schema": SCHEMA, "manifest": spec, "workload_sha256": workload_hash,
|
||||
"sources": [{"file": p.name, "sha256": digest(p)} for p in files],
|
||||
"comparison": spec["comparison"],
|
||||
"status": "incomplete_or_failed" if issues else "protocol_complete",
|
||||
"quality": "not_assessed", "hardware_and_server_config": "operator_declared",
|
||||
"token_latency": "not_measured", "memory_and_power": "external_telemetry_required",
|
||||
"issues": issues, "cells": table}
|
||||
|
||||
|
||||
def main():
|
||||
parser = argparse.ArgumentParser(description=__doc__)
|
||||
parser.add_argument("action", choices=("plan", "run", "compare"))
|
||||
parser.add_argument("--manifest", required=True)
|
||||
parser.add_argument("--engine", choices=ENGINES)
|
||||
parser.add_argument("--round", type=int)
|
||||
parser.add_argument("--results", default="baseline-results")
|
||||
args = parser.parse_args()
|
||||
try:
|
||||
spec, workload, workload_hash = load_manifest(args.manifest)
|
||||
if args.action == "plan":
|
||||
print(json.dumps(list(plan(spec)), indent=2))
|
||||
return 0
|
||||
if args.action == "run":
|
||||
require(args.engine is not None and args.round is not None, "run needs --engine and --round")
|
||||
return collect(spec, workload, workload_hash, args.engine, args.round, args.results)
|
||||
result = compare(spec, workload, workload_hash, args.results)
|
||||
print(json.dumps(result, indent=2, allow_nan=False))
|
||||
return int(bool(result["issues"]))
|
||||
except (ValueError, OSError, KeyError, TypeError) as exc:
|
||||
parser.error(str(exc))
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
raise SystemExit(main())
|
||||
@@ -0,0 +1,130 @@
|
||||
# Colibri, SGLang and vLLM serving baseline
|
||||
|
||||
This is a reproducible collection protocol, **not a measured performance ranking**.
|
||||
`c/tools/benchmark_baseline.py` reuses the existing HTTP harness. It does not start,
|
||||
stop, reconfigure, or install engines. Run one engine at a time on an explicitly
|
||||
allocated host; three servers sharing the same GPU would invalidate isolation.
|
||||
No model or GPU is needed to plan or summarize an experiment.
|
||||
|
||||
## Freeze the experiment
|
||||
|
||||
Copy `three-engine.example.json` and `workload.jsonl` into an experiment directory.
|
||||
Fill every `REPLACE` value before running; the example intentionally fails
|
||||
validation. Workload paths are relative to the manifest, not the working directory.
|
||||
Keep the same manifest for every engine/round. Do not store keys in launch commands:
|
||||
use environment-variable references. Only `api_key_env` names are read by the
|
||||
collector; their secret values are never included in results.
|
||||
|
||||
Record:
|
||||
|
||||
- The physical host, CPU/NUMA topology, RAM, GPU count/model/VRAM, storage,
|
||||
OS and driver/runtime versions. Use the same allocated hardware for every arm.
|
||||
- The exact source model revision, tokenizer and rendered chat-template SHA-256s.
|
||||
Hash a sorted inventory of weight files for each `artifact_sha256`; retain the
|
||||
inventory and conversion commands. Include weight **and KV** precision in
|
||||
`quantization`. These are operator declarations, not remote attestations.
|
||||
- Engine commits or image digests and complete launch commands, including TP/PP/EP,
|
||||
context limit, admission limits, memory budget and CPU thread placement.
|
||||
Aliases and URLs may differ; they do not identify the checkpoint.
|
||||
- Reasoning, speculation, prefix-cache and warmup policies. Disable speculation in
|
||||
the initial baseline, then create a separate experiment to measure it. Warmup is
|
||||
performed before **each concurrency cell**. Caches persist across cells; the
|
||||
collector never flushes them. Reset/restart before each round if the declared
|
||||
policy calls for it. Cold and warm results belong in separate experiments.
|
||||
- A fixed quality evaluation dataset/revision, metric and acceptance threshold.
|
||||
Run this independently on all three configurations and retain the results.
|
||||
|
||||
Choose the comparison explicitly:
|
||||
|
||||
| Mode | Enforced identity | Permitted interpretation |
|
||||
| --- | --- | --- |
|
||||
| `matched_artifact` | All three weight inventory hashes, formats and quantization descriptions must match | Matched-artifact performance observations, conditional on external quality and configuration checks |
|
||||
| `deployment` | Shared source model/tokenizer/template and workload; engine artifacts may differ | A comparison of complete deployment configurations, **not an isolated engine speedup** |
|
||||
|
||||
A Colibri converted INT4 artifact and a vLLM NVFP4 checkpoint belong in
|
||||
`deployment`, even if they originated from the same model. Matching labels or
|
||||
hashes does not prove equal quality or backend numerical behavior.
|
||||
|
||||
Use separate manifests/directories for `fully_resident` and `offload`. Verify
|
||||
residency with engine counters and I/O measurements; do not infer it from model
|
||||
size or successful startup. Collect RAM/VRAM peaks, disk/PCIe traffic and power in
|
||||
external telemetry, aligned to each report's UTC `started_at`. The manifest's
|
||||
hardware data is descriptive; the tool does not measure memory or energy.
|
||||
|
||||
## Collect the matrix
|
||||
|
||||
From the repository root:
|
||||
|
||||
```sh
|
||||
python3 c/tools/benchmark_baseline.py plan --manifest /path/experiment/manifest.json
|
||||
```
|
||||
|
||||
The default example describes C1/C4/C8/C16, three independent rounds, and 16
|
||||
measured requests per cell. **The two included prompts are smoke inputs only**;
|
||||
replace them with representative short, long and shared-prefix workloads and
|
||||
increase the request count before interpreting P95/P99. Request counts must at
|
||||
least reach the largest concurrency; this does not guarantee sustained occupancy.
|
||||
`repeats` repeats requests within a cell; `rounds` repeats complete measurements.
|
||||
|
||||
The plan rotates engine order to reduce systematic order bias:
|
||||
|
||||
1. Colibri, SGLang, vLLM.
|
||||
2. SGLang, vLLM, Colibri.
|
||||
3. vLLM, Colibri, SGLang.
|
||||
|
||||
After manually starting the selected engine with the recorded configuration:
|
||||
|
||||
```sh
|
||||
python3 c/tools/benchmark_baseline.py run \
|
||||
--manifest /path/experiment/manifest.json \
|
||||
--engine colibri --round 1 --results /path/experiment/results
|
||||
```
|
||||
|
||||
Repeat for each plan entry, releasing the previous engine's allocated hardware
|
||||
before starting the next. The collector only contacts the selected endpoint.
|
||||
Each cell writes `<engine>-r<round>-c<concurrency>.json`, with the full manifest,
|
||||
workload and harness hashes, raw per-request records, warmup and summary.
|
||||
Existing cells are never overwritten. A failed warmup saves evidence and stops
|
||||
the round; measured failures are retained and make the command exit nonzero.
|
||||
For a failed/invalid campaign, keep the evidence and use a new results directory
|
||||
for a replacement campaign. Do not silently substitute the fastest repeat.
|
||||
|
||||
## Summarize and interpret
|
||||
|
||||
```sh
|
||||
python3 c/tools/benchmark_baseline.py compare \
|
||||
--manifest /path/experiment/manifest.json \
|
||||
--results /path/experiment/results > /path/experiment/comparison.json
|
||||
```
|
||||
|
||||
The comparison rejects mismatched manifests/workloads, mixed collector/harness
|
||||
versions, duplicate or unexpected cells, altered summaries and missing requests.
|
||||
Missing cells, failed requests, missing usage and empty/filtered output are listed
|
||||
as issues and return a nonzero exit. No failed cell is silently dropped.
|
||||
|
||||
For each engine/concurrency it reports min/median/max **across rounds** for:
|
||||
|
||||
- Aggregate successful completion tokens per full batch wall second.
|
||||
- Client first-output P50/P95/P99 and completion-duration P95.
|
||||
- Failure rate and latency-SLO request goodput.
|
||||
- It also reports actual server-reported output-length distribution across rounds.
|
||||
|
||||
Round percentiles are not pooled percentiles. Throughput is aggregate, never
|
||||
per-user decode speed, and includes HTTP, queueing and prefill. First output is
|
||||
the first nonempty content/reasoning/tool delta, not an instrumented first token.
|
||||
SSE chunk gaps are **not ITL**, and no token-latency/TPOT estimate is manufactured.
|
||||
Instrument the engines or use their native benchmarks for token timestamps;
|
||||
report those separately with their definitions. The underlying HTTP timing,
|
||||
usage and socket-timeout boundaries are documented in [benchmarking.md](../benchmarking.md).
|
||||
|
||||
`protocol_complete` means the matrix completed without those protocol/data gaps.
|
||||
It does **not** mean quality passed, outputs have equivalent lengths, SLOs were
|
||||
met, resources were isolated, residency was verified, or memory/power was measured.
|
||||
`quality: not_assessed` remains explicit. This report deliberately emits neither
|
||||
winner labels nor speedup ratios. Inspect lengths, reasoning accounting, external
|
||||
quality and telemetry before publishing any performance conclusion.
|
||||
|
||||
The matrix currently uses closed-loop traffic. Rate-driven interference tests
|
||||
remain available in `benchmark_http_serving.py`; do not mix them into this matrix.
|
||||
A useful follow-up workload adds a long prefill while short decodes are active,
|
||||
with a separate report of first-output and ongoing token-latency interference.
|
||||
@@ -0,0 +1,73 @@
|
||||
{
|
||||
"schema": "colibri.serving-baseline/v1",
|
||||
"experiment": "REPLACE-experiment-id",
|
||||
"comparison": "deployment",
|
||||
"residency": "fully_resident",
|
||||
"workload": "workload.jsonl",
|
||||
"hardware": {
|
||||
"host": "REPLACE-host",
|
||||
"cpu": "REPLACE-cpu",
|
||||
"ram": "REPLACE-ram",
|
||||
"gpu": "REPLACE-gpu",
|
||||
"storage": "REPLACE-storage",
|
||||
"os": "REPLACE-os",
|
||||
"driver": "REPLACE-driver"
|
||||
},
|
||||
"model": {
|
||||
"source_revision": "REPLACE-source_revision",
|
||||
"tokenizer_sha256": "REPLACE-tokenizer_sha256",
|
||||
"chat_template_sha256": "REPLACE-chat_template_sha256"
|
||||
},
|
||||
"cache_policy": "REPLACE-reset-before-round; warm-with-recorded-workload",
|
||||
"reasoning_policy": "REPLACE-explicit-thinking-setting",
|
||||
"speculation_policy": "REPLACE-disabled-for-initial-baseline",
|
||||
"quality_protocol": "REPLACE-quality-dataset-revision-metric-and-acceptance-threshold",
|
||||
"matrix": {
|
||||
"concurrency": [
|
||||
1,
|
||||
4,
|
||||
8,
|
||||
16
|
||||
],
|
||||
"rounds": 3,
|
||||
"repeats": 8,
|
||||
"warmup_requests": 2,
|
||||
"max_tokens": 128,
|
||||
"temperature": 0,
|
||||
"timeout": 120,
|
||||
"slo_first_output": 5,
|
||||
"slo_duration": 120
|
||||
},
|
||||
"engines": {
|
||||
"colibri": {
|
||||
"base_url": "http://127.0.0.1:8000/v1",
|
||||
"served_model": "REPLACE-model-alias",
|
||||
"revision": "REPLACE-commit-or-image-digest",
|
||||
"launch_command": "REPLACE-exact-command-with-secret-environment-references",
|
||||
"artifact_sha256": "REPLACE-weight-manifest-sha256",
|
||||
"weight_format": "REPLACE-format",
|
||||
"quantization": "REPLACE-weight-and-KV-precision",
|
||||
"api_key_env": "OPENAI_API_KEY"
|
||||
},
|
||||
"sglang": {
|
||||
"base_url": "http://127.0.0.1:8000/v1",
|
||||
"served_model": "REPLACE-model-alias",
|
||||
"revision": "REPLACE-commit-or-image-digest",
|
||||
"launch_command": "REPLACE-exact-command-with-secret-environment-references",
|
||||
"artifact_sha256": "REPLACE-weight-manifest-sha256",
|
||||
"weight_format": "REPLACE-format",
|
||||
"quantization": "REPLACE-weight-and-KV-precision",
|
||||
"api_key_env": "OPENAI_API_KEY"
|
||||
},
|
||||
"vllm": {
|
||||
"base_url": "http://127.0.0.1:8000/v1",
|
||||
"served_model": "REPLACE-model-alias",
|
||||
"revision": "REPLACE-commit-or-image-digest",
|
||||
"launch_command": "REPLACE-exact-command-with-secret-environment-references",
|
||||
"artifact_sha256": "REPLACE-weight-manifest-sha256",
|
||||
"weight_format": "REPLACE-format",
|
||||
"quantization": "REPLACE-weight-and-KV-precision",
|
||||
"api_key_env": "OPENAI_API_KEY"
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,2 @@
|
||||
{"messages": [{"role": "user", "content": "Explain the difference between throughput and latency using a concrete example."}]}
|
||||
{"messages": [{"role": "user", "content": "Write a short Python function that merges two sorted lists and explain its complexity."}]}
|
||||
@@ -247,3 +247,12 @@ then completing in 100 ms from meeting a one-second total-latency target.
|
||||
Goodput still divides qualifying completions by the entire batch wall time,
|
||||
including the drain after the last scheduled arrival. Keep rate, concurrency,
|
||||
request count, latency targets and timing basis equal across compared reports.
|
||||
|
||||
### Repeated Colibri / SGLang / vLLM baseline
|
||||
|
||||
The [three-engine baseline protocol](baselines/README.md) provides a manifest,
|
||||
rotating run plan and C1/C4/C8/C16 collector using this HTTP harness. Its comparison
|
||||
checks workload/configuration identity, keeps failed and missing cells visible,
|
||||
and reports min/median/max across rounds. It distinguishes matched artifacts from
|
||||
deployment comparisons with different weight formats. Quality, token-level latency
|
||||
and hardware telemetry require separate evidence; no performance ranking is bundled.
|
||||
|
||||
Reference in New Issue
Block a user