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()