File size: 6,613 Bytes
7df674a
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
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
#!/usr/bin/env python3
"""harvest pod benchmark — throughput + GPU-utilization sweep.

Runs advisory-style chat completions at increasing concurrency against a local
llama-server, with reasoning ON or OFF, while sampling nvidia-smi. Reports
aggregate output tok/s (the budget number), latency, TTFT, and GPU util/VRAM.

  python3 bench-pod.py --base-url http://127.0.0.1:18080/v1 --model qwen3.6-35b \
      --concurrencies 1,8,32,64 --thinking off --max-tokens 256 --out r.json
"""
import argparse, json, statistics, subprocess, sys, threading, time
from concurrent.futures import ThreadPoolExecutor, as_completed

try:
    import requests
except Exception:
    sys.exit("needs `requests` (pip install requests)")

SYSTEM = ("You are one advisory agent in a multi-agent system serving a smallholder "
          "farmer across a full cropping season. Give concise, decision-focused advice.")
USER = ("Farmer profile: 0.8 ha rainfed plot, sandy-loam soil, monsoon onset delayed ~2 weeks, "
        "cotton, moderate bollworm pressure reported nearby, limited cash for inputs. "
        "It is the sowing window. Recommend: sow now vs wait, variety duration, and the single "
        "most important early-season action. Explain your reasoning in a few sentences.")


class GpuSampler(threading.Thread):
    def __init__(self):
        super().__init__(daemon=True)
        self.samples, self.on = [], True
    def run(self):
        while self.on:
            try:
                out = subprocess.run(
                    ["nvidia-smi", "--query-gpu=utilization.gpu,memory.used",
                     "--format=csv,noheader,nounits"],
                    capture_output=True, text=True, timeout=5).stdout.strip().splitlines()[0]
                u, m = [int(x.strip()) for x in out.split(",")]
                self.samples.append((u, m))
            except Exception:
                pass
            time.sleep(1)
    def window(self, start_idx):
        w = self.samples[start_idx:]
        if not w: return {}
        return {"gpu_util_mean": round(statistics.mean(s[0] for s in w), 1),
                "gpu_util_max": max(s[0] for s in w),
                "vram_mb_max": max(s[1] for s in w)}


def one_request(base_url, model, key, max_tokens, thinking):
    body = {"model": model,
            "messages": [{"role": "system", "content": SYSTEM},
                         {"role": "user", "content": USER}],
            "max_tokens": max_tokens, "temperature": 0.7, "stream": True,
            "stream_options": {"include_usage": True},
            "chat_template_kwargs": {"enable_thinking": thinking}}
    t0 = time.perf_counter(); ttft = None; ctoks = 0
    try:
        with requests.post(base_url.rstrip("/") + "/chat/completions",
                           headers={"Authorization": f"Bearer {key}"},
                           json=body, stream=True, timeout=900) as r:
            r.raise_for_status()
            for line in r.iter_lines():
                if not line: continue
                s = line.decode("utf-8", "ignore")
                if not s.startswith("data: "): continue
                s = s[6:]
                if s.strip() == "[DONE]": break
                try: ch = json.loads(s)
                except json.JSONDecodeError: continue
                if ttft is None and ch.get("choices"):
                    d = ch["choices"][0].get("delta", {})
                    if d.get("content") or d.get("reasoning_content"):
                        ttft = time.perf_counter() - t0
                if ch.get("usage"):
                    ctoks = ch["usage"].get("completion_tokens", 0)
        return ttft, time.perf_counter() - t0, ctoks, True, None
    except Exception as e:
        return ttft, time.perf_counter() - t0, ctoks, False, f"{type(e).__name__}: {e}"


def main():
    ap = argparse.ArgumentParser()
    ap.add_argument("--base-url", required=True)
    ap.add_argument("--model", required=True)
    ap.add_argument("--api-key", default="sk-harvest-local")
    ap.add_argument("--concurrencies", default="1,8,32,64")
    ap.add_argument("--rounds", type=int, default=2)
    ap.add_argument("--max-tokens", type=int, default=256)
    ap.add_argument("--thinking", choices=["on", "off"], default="off")
    ap.add_argument("--out", default=None)
    a = ap.parse_args()
    thinking = a.thinking == "on"

    sampler = GpuSampler(); sampler.start()
    print(f"model={a.model} thinking={a.thinking} max_tokens={a.max_tokens} rounds={a.rounds}")
    print(f"{'conc':>5} {'ok':>7} {'agg tok/s':>10} {'lat p50':>8} {'ttft p50':>9} "
          f"{'util avg':>9} {'util max':>9} {'vram MB':>8}")
    results = []
    for c in [int(x) for x in a.concurrencies.split(",") if x.strip()]:
        idx = len(sampler.samples)
        n = c * a.rounds
        lat, ttfts, toks, oks = [], [], [], 0
        w0 = time.perf_counter()
        with ThreadPoolExecutor(max_workers=c) as ex:
            futs = [ex.submit(one_request, a.base_url, a.model, a.api_key,
                              a.max_tokens, thinking) for _ in range(n)]
            for f in as_completed(futs):
                ttft, total, ct, ok, err = f.result()
                if ok:
                    oks += 1; lat.append(total); toks.append(ct)
                    if ttft is not None: ttfts.append(ttft)
                elif err:
                    print(f"  [err] {err[:90]}", file=sys.stderr)
        wall = time.perf_counter() - w0
        gpu = sampler.window(idx)
        agg = sum(toks) / wall if wall else 0
        row = {"concurrency": c, "ok": oks, "n": n, "agg_tok_s": round(agg, 1),
               "lat_p50_s": round(statistics.median(lat), 2) if lat else None,
               "ttft_p50_s": round(statistics.median(ttfts), 3) if ttfts else None,
               "mean_completion_tokens": round(statistics.mean(toks), 1) if toks else 0,
               "wall_s": round(wall, 1), **gpu}
        results.append(row)
        print(f"{c:>5} {oks:>3}/{n:<3} {row['agg_tok_s']:>10} {row['lat_p50_s'] or 0:>8} "
              f"{row['ttft_p50_s'] or 0:>9} {row.get('gpu_util_mean',0):>9} "
              f"{row.get('gpu_util_max',0):>9} {row.get('vram_mb_max',0):>8}")
    sampler.on = False
    peak = max((r["agg_tok_s"] for r in results), default=0)
    print(f"\npeak aggregate: {peak:.0f} tok/s (thinking={a.thinking})")
    if a.out:
        json.dump({"model": a.model, "thinking": a.thinking,
                   "max_tokens": a.max_tokens, "results": results,
                   "peak_agg_tok_s": peak}, open(a.out, "w"), indent=2)
        print("wrote", a.out)


if __name__ == "__main__":
    main()