spec-b300 / source /scripts /cluster /launch_vllm_benchmark.py
khazic's picture
Archive three-epoch run: logs and provenance part 2
932bc69 verified
Raw History Blame Contribute Delete
7.55 kB
"""Benchmark-only launcher: one/two same-GPU engines behind a local HTTP proxy.
Use as LAUNCH_VLLM in the single-node recipe. Both engines retain the normal
hidden-state connector and their own launch provenance. This is not a default.
"""
import asyncio
import contextlib
import itertools
import os
import signal
import sys
from pathlib import Path
import aiohttp
from aiohttp import web
def option(arguments, name, default=None):
return arguments[arguments.index(name) + 1] if name in arguments else default
def replace_option(arguments, name, value):
result = list(arguments)
if name in result:
result[result.index(name) + 1] = str(value)
else:
result.extend((name, str(value)))
return result
async def main():
arguments = sys.argv[1:]
boundary = arguments.index("--")
launch_args, engine_args = arguments[:boundary], arguments[boundary + 1 :]
provenance = Path(option(launch_args, "--provenance-dir"))
port = int(option(engine_args, "--port", "8000"))
host = option(engine_args, "--host", "127.0.0.1")
if host != "127.0.0.1":
raise ValueError("The benchmark proxy must remain loopback-only")
if option(engine_args, "--data-parallel-size", "1") != "1":
raise ValueError("This benchmark requires DP=1 per engine")
engine_count = int(os.environ.get("VLLM_BENCHMARK_ENGINES", "2"))
if engine_count not in (1, 2):
raise ValueError("Only one or two engines fit this benchmark topology")
endpoints = [
f"http://127.0.0.1:{port + offset}" for offset in range(1, engine_count + 1)
]
children, logs = [], []
stop = asyncio.Event()
for sig in (signal.SIGINT, signal.SIGTERM):
asyncio.get_running_loop().add_signal_handler(sig, stop.set)
runner = None
try:
timeout = aiohttp.ClientTimeout(total=1800)
async with aiohttp.ClientSession(
timeout=timeout, auto_decompress=False, trust_env=False
) as client:
# Sequential startup avoids cross-process interference in vLLM's
# activation-memory profiler. Runtime inference remains concurrent.
for index, endpoint in enumerate(endpoints):
engine_provenance = provenance / f"engine{index}"
engine_provenance.mkdir(parents=True, exist_ok=True)
child_launch = replace_option(
launch_args, "--provenance-dir", engine_provenance
)
child_engine = replace_option(engine_args, "--port", port + index + 1)
child_engine = replace_option(
child_engine,
"--gpu-memory-utilization",
"0.44" if engine_count == 2 else "0.92",
)
child_engine = replace_option(child_engine, "--api-server-count", "1")
log = (engine_provenance / "server.log").open("w")
logs.append(log)
child = await asyncio.create_subprocess_exec(
sys.executable,
str(Path(__file__).resolve().parents[1] / "launch_vllm.py"),
*child_launch,
"--",
*child_engine,
env={**os.environ, "VLLM_PORT": str(port + index + 101)},
stdout=log,
stderr=asyncio.subprocess.STDOUT,
)
children.append(child)
print(
f"Starting engine {index}: pid={child.pid}, {endpoint}, log={log.name}",
flush=True,
)
deadline = asyncio.get_running_loop().time() + 1500
while not stop.is_set():
if child.returncode is not None:
raise RuntimeError(
f"Engine {index} exited {child.returncode}; see {log.name}"
)
try:
async with client.get(
endpoint + "/health", timeout=aiohttp.ClientTimeout(total=3)
) as response:
if response.status == 200:
break
except (aiohttp.ClientError, TimeoutError):
pass
if asyncio.get_running_loop().time() >= deadline:
raise TimeoutError(f"Engine {index} startup timed out")
with contextlib.suppress(TimeoutError):
await asyncio.wait_for(stop.wait(), timeout=2)
if stop.is_set():
return
print(f"Engine {index} healthy", flush=True)
rotation = itertools.cycle(endpoints)
hop_headers = {
"host",
"connection",
"transfer-encoding",
"keep-alive",
"content-length",
}
async def forward(request):
endpoint = next(rotation) if request.method == "POST" else endpoints[0]
headers = {
key: value
for key, value in request.headers.items()
if key.lower() not in hop_headers
}
async with client.request(
request.method,
endpoint + request.path_qs,
headers=headers,
data=await request.read(),
allow_redirects=False,
) as upstream:
response = web.StreamResponse(
status=upstream.status,
headers={
key: value
for key, value in upstream.headers.items()
if key.lower() not in hop_headers
},
)
await response.prepare(request)
async for chunk in upstream.content.iter_chunked(65536):
await response.write(chunk)
await response.write_eof()
return response
app = web.Application(client_max_size=64 * 1024 * 1024)
app.router.add_route("*", "/{path:.*}", forward)
runner = web.AppRunner(app, access_log=None)
await runner.setup()
await web.TCPSite(runner, host, port).start()
print(f"{engine_count}-engine proxy healthy on {host}:{port}", flush=True)
while not stop.is_set():
for index, child in enumerate(children):
if child.returncode is not None:
raise RuntimeError(f"Engine {index} exited {child.returncode}")
with contextlib.suppress(TimeoutError):
await asyncio.wait_for(stop.wait(), timeout=2)
finally:
if runner is not None:
await runner.cleanup()
for child in children:
if child.returncode is None:
with contextlib.suppress(ProcessLookupError):
child.terminate()
for child in children:
try:
await asyncio.wait_for(child.wait(), timeout=20)
except TimeoutError:
with contextlib.suppress(ProcessLookupError):
child.kill()
await child.wait()
for log in logs:
log.close()
if __name__ == "__main__":
asyncio.run(main())