File size: 7,147 Bytes
3fd1a35 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 | #!/usr/bin/env python3
"""Run one real long API request and collect process/system memory on the board."""
import argparse
import concurrent.futures
import json
import os
from pathlib import Path
import signal
import subprocess
import time
import urllib.request
def read_kib(path):
result = {}
for line in path.read_text().splitlines():
key, _, value = line.partition(":")
if value.strip().endswith("kB"):
result[key] = int(value.split()[0]) / 1024
return result
def main():
p = argparse.ArgumentParser(description=__doc__)
p.add_argument("--binary", type=Path, required=True)
p.add_argument("--model", type=Path, required=True)
p.add_argument("--output-dir", type=Path, required=True)
p.add_argument("--context", type=int, default=65536)
p.add_argument("--input-tokens", type=int, default=65472)
p.add_argument("--new-tokens", type=int, default=64)
p.add_argument("--port", type=int, default=19095)
a = p.parse_args()
if a.input_tokens < 22 or a.input_tokens + a.new_tokens > a.context:
p.error("input + output must fit the selected capacity")
a.output_dir.mkdir(parents=True, exist_ok=False)
client = urllib.request.build_opener(urllib.request.ProxyHandler({}))
def call(path, body=None, timeout=10):
request = urllib.request.Request(
f"http://127.0.0.1:{a.port}" + path,
data=None if body is None else json.dumps(body).encode(),
headers={"Content-Type": "application/json"})
with client.open(request, timeout=timeout) as response:
return json.load(response)
def infer(prompt, count):
return call("/v1/chat/completions", {
"model": "mindnano-ling3-tiny", "temperature": 0, "cache_prompt": False,
"messages": [{"role": "user", "content": prompt}],
"max_tokens": count}, timeout=86400)
result = {"context": a.context, "requested_input_tokens": a.input_tokens,
"requested_output_tokens": a.new_tokens, "started_unix": time.time()}
runtime = a.output_dir / "runtime.log"
samples = a.output_dir / "memory.jsonl"
child = None
engine_pid = None
started = time.monotonic()
def sample(phase):
nonlocal engine_pid
mem = read_kib(Path("/proc/meminfo"))
data = {"timestamp": time.time(), "elapsed_s": time.monotonic() - started,
"phase": phase, "available_mib": mem.get("MemAvailable"),
"system_swap_used_mib": mem.get("SwapTotal", 0) - mem.get("SwapFree", 0)}
if child is not None and child.poll() is None:
children = Path(f"/proc/{child.pid}/task/{child.pid}/children")
if children.exists():
ids = children.read_text().split()
if ids:
engine_pid = int(ids[0])
status = Path(f"/proc/{engine_pid}/status")
if status.exists():
proc = read_kib(status)
data.update({"engine_pid": engine_pid,
"rss_mib": proc.get("VmRSS"), "peak_rss_mib": proc.get("VmHWM"),
"swap_mib": proc.get("VmSwap")})
with samples.open("a") as out:
out.write(json.dumps(data) + "\n")
return data
try:
sample("before_start")
with runtime.open("w") as log:
child = subprocess.Popen([
str(a.binary.resolve()), "--model", str(a.model.resolve()),
"--context", str(a.context), "--host", "127.0.0.1", "--port", str(a.port),
"--no-console", "--log", str((a.output_dir / "metrics.jsonl").resolve())],
stdout=log, stderr=subprocess.STDOUT, start_new_session=True,
env={"PATH": "/usr/bin:/bin"})
print(json.dumps({"launcher_pid": child.pid, "runtime_log": str(runtime)}), flush=True)
deadline = time.monotonic() + 300
while True:
sample("initialization")
if child.poll() is not None:
raise RuntimeError(f"service exited during startup: {child.returncode}")
try:
result["health_before"] = call("/health")
break
except (OSError, ValueError):
if time.monotonic() >= deadline:
raise TimeoutError("startup did not complete within 300 seconds")
time.sleep(2)
result["short_before"] = infer("1" * 107, 64)
sample("ready")
print(json.dumps({"ready": result["health_before"],
"short_metrics": result["short_before"]["mindnano_metrics"]}), flush=True)
(a.output_dir / "initialization.json").write_text(json.dumps(result, indent=2) + "\n")
request_start = time.monotonic()
with concurrent.futures.ThreadPoolExecutor(max_workers=1) as pool:
future = pool.submit(infer, "1" * (a.input_tokens - 21), a.new_tokens)
next_report = 0
while not future.done():
data = sample("long_request")
elapsed = time.monotonic() - request_start
if elapsed >= next_report:
print(json.dumps({"request_elapsed_s": elapsed, **data}), flush=True)
next_report = elapsed + 60
if child.poll() is not None:
raise RuntimeError(f"service exited during request: {child.returncode}")
try:
future.result(timeout=10)
except concurrent.futures.TimeoutError:
pass
result["long_request"] = future.result()
result["request_wall_s"] = time.monotonic() - request_start
long = result["long_request"]
assert long["usage"]["prompt_tokens"] == a.input_tokens, long["usage"]
assert long["mindnano_metrics"]["prompt_evaluated_tokens"] == a.input_tokens
assert long["mindnano_metrics"]["cached_tokens"] == 0
result["short_after"] = infer("1" * 107, 64)
assert result["short_before"]["choices"] == result["short_after"]["choices"]
result["health_after"] = call("/health")
sample("complete")
result["passed"] = True
print(json.dumps({"passed": True, "metrics": long["mindnano_metrics"]}), flush=True)
except BaseException as exc:
result["passed"] = False
result["error"] = repr(exc)
raise
finally:
result["finished_unix"] = time.time()
(a.output_dir / "result.json").write_text(json.dumps(result, ensure_ascii=False, indent=2) + "\n")
if child is not None and child.poll() is None:
os.killpg(child.pid, signal.SIGTERM)
try:
child.wait(timeout=120)
except subprocess.TimeoutExpired:
os.killpg(child.pid, signal.SIGKILL)
child.wait()
sample("after_stop")
if __name__ == "__main__":
main()
|