test(kernel-stress): 把"真二进制压力测试"harness 固化进仓库(私有 netns + mock LLM + 多连接驱动)

本轮压测是从 /var/tmp 里现写脚本跑的,能复现才能长期用 —— 固化到这。

为什么要有它(与 go test 的分工):单测都在进程内、参数显式给,抓不到集成面问题;
这一轮它抓到了三个只有真二进制才暴露的问题(DataDir 漏接线、插件通道未登记为
inputch、create 后子不开工),以及两处生成器缺陷(旧 SDK 导致生成工程编译失败;
模板工程从不示范通道登记)。

- `launch.sh`:在**私有 netns** 里同时起 mock LLM 与内核实例。
  必须隔离的原因:生产实例占着 *:8080/*:9890/*:9876 且插件设了 SO_REUSEADDR,
  同机再起一个实例会在 127.0.0.1 上与之并存绑定(首次实测抢到 9890 约 1 分钟);
  netns 里只有 lo,结构上不可能碰到生产端口。unix socket 是文件系统对象,
  所以驱动脚本在 netns 外照样能连。
- `mockllm.py`:OpenAI 兼容 mock(可控延迟 + SSE 分块 + 按 `!resident`/`!notify`
  标记回工具调用),让无外网也能跑、且**流式段可被中断**(压抢占路径的前提)。
  注意必须实现 `do_HEAD`:内核探活用 HEAD,501 会被判不可达 → degraded → rollback 循环。
- `stress.py`:多并发连接轰炸(单连接是串行的,造不出队列压力)+ 中断线程 + 峰值采样。
- `kcli.py`:CLI socket 客户端(`/auth` → `/kernel` / 广播)。
- README:前置配置(含探活端点、rollback 关掉)、用法、指标读法、已知坑。

实测结果(详见 ed76695 提交信息):密集 274 任务 executed=274/rejected=0;
稀疏中断下 suspended=resumed=preempted=27;驻留子全链路 + 优雅退出不留孤儿。
This commit is contained in:
JianFeeeee
2026-09-13 11:38:31 +08:00
parent d7cc960b0f
commit 64e36088ec
5 changed files with 407 additions and 0 deletions

View File

@ -0,0 +1,144 @@
#!/usr/bin/env python3
"""二进制级调度器压力:多并发连接轰炸排队输入 + L4 中断,并采样峰值。
用法: stress.py <sock> <key> <conns> <inputs_per_conn> <interrupt_threads> <interrupts_each> [resident]
"""
import json, os, socket, sys, threading, time
SOCK, KEY = sys.argv[1], sys.argv[2]
NC, NI, IT, IE = (int(x) for x in sys.argv[3:7])
MODE = sys.argv[7] if len(sys.argv) > 7 else ""
GAP = float(sys.argv[8]) if len(sys.argv) > 8 else 0.05
lock = threading.Lock()
sent = 0
done = 0
errs = []
lat = []
peak = {"queue": 0, "pending": 0, "stack": 0}
def connect():
s = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
s.settimeout(60)
s.connect(SOCK)
s.sendall(("/auth %s\n" % KEY).encode())
buf = b""
while b"\n" not in buf:
buf += s.recv(65536)
return s
def chat(s, text):
"""发一条输入并等到终止帧(response/error)。"""
s.sendall((text + "\n").encode())
buf = b""
while True:
c = s.recv(65536)
if not c:
raise RuntimeError("closed")
buf += c
while b"\n" in buf:
ln, buf = buf.split(b"\n", 1)
if not ln.strip():
continue
try:
obj = json.loads(ln.decode("utf-8", "replace"))
except Exception:
continue
if obj.get("type") in ("response", "error"):
return obj
def worker(idx, n):
global sent, done
try:
s = connect()
except Exception as e:
with lock:
errs.append("connect: %r" % e)
return
for i in range(n):
t = time.time()
try:
obj = chat(s, "w%d-%d%s" % (idx, i, " !resident" if MODE == "resident" else ""))
except Exception as e:
with lock:
errs.append("chat: %r" % e)
return
dt = time.time() - t
with lock:
sent += 1
if obj.get("type") == "response":
done += 1
lat.append(dt)
s.close()
def interrupter(idx, n, gap):
for i in range(n):
try:
s = connect()
s.sendall(("/interrupt stress-%d-%d\n" % (idx, i)).encode())
time.sleep(gap)
s.close()
except Exception as e:
with lock:
errs.append("intr: %r" % e)
return
def sampler(stop, dur):
while not stop.is_set():
try:
s = connect()
s.sendall(b"/kernel\n")
buf = b""
while b"\n" not in buf:
buf += s.recv(65536)
obj = json.loads(buf.split(b"\n")[0].decode())
sc = json.loads(obj["content"]).get("scheduler", {})
with lock:
peak["queue"] = max(peak["queue"], sc.get("ready_queue_depth", 0))
peak["pending"] = max(peak["pending"], sc.get("pending_interrupts", 0))
peak["stack"] = max(peak["stack"], sc.get("suspend_stack", 0))
s.close()
except Exception:
time.sleep(0.05)
time.sleep(0.1)
def kstat():
s = connect()
s.sendall(b"/kernel\n")
buf = b""
while b"\n" not in buf:
buf += s.recv(65536)
return json.loads(json.loads(buf.split(b"\n")[0].decode())["content"])
t0 = time.time()
stop = threading.Event()
threads = [threading.Thread(target=worker, args=(i, NI)) for i in range(NC)]
threads += [threading.Thread(target=interrupter, args=(i, IE, GAP)) for i in range(IT)]
smp = threading.Thread(target=sampler, args=(stop, 0), daemon=True)
smp.start()
for t in threads:
t.start()
for t in threads:
t.join()
stop.set()
time.sleep(0.4)
st = kstat()
with lock:
lat.sort()
p = lambda q: (lat[int(len(lat) * q)] if lat else 0)
print(json.dumps({
"耗时s": round(time.time() - t0, 2),
"并发连接": NC, "每连接输入": NI, "中断线程": IT, "每线程中断": IE,
"完成": done, "错误": len(errs),
"延迟_p50": round(p(0.5), 2), "延迟_p95": round(p(0.95), 2), "延迟_max": round(lat[-1] if lat else 0, 2),
"峰值_队列": peak["queue"], "峰值_待处理中断": peak["pending"], "峰值_中断栈": peak["stack"],
"scheduler": st.get("scheduler"), "goroutines": st.get("runtime", {}).get("goroutines"),
"错误样本": errs[:3],
}, ensure_ascii=False))