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 关掉)、用法、指标读法、已知坑。

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

View File

@ -0,0 +1,63 @@
# 内核二进制压力测试kernel-stress
对本仓库**编译出来的真实内核**做压力测试 —— 与 `go test` 的区别是:它跑真二进制、
真插件加载、真 unix socket 协议,因此能抓到只在集成面上出现的问题
(已有战绩:根 agent `DataDir` 漏接线、插件通道没登记为 inputch、create 后子不开工)。
## 为什么必须放在私有 netns 里
生产实例占着 `*:8080 / *:9890 / *:9876`,而插件的监听都设了 `SO_REUSEADDR`
同机再起一个实例会**在 127.0.0.1 上与之并存绑定**(实测抢到过 `127.0.0.1:9890` 约 1 分钟)。
`unshare -n` 后实例只有 `lo`,结构上不可能碰到生产端口。
unix socket 是文件系统对象,跨 netns 仍可驱动,所以驱动脚本在 netns 外也能用。
## 前置
```bash
go build -o /tmp/homed-stress ./cmd/homed # 被压的内核
export GOCACHE=/tmp/gocache GOPATH=/tmp/gopath TMPDIR=/var/tmp/gotmp
```
## 用法
```bash
# 1) 准备数据目录 + 把 LLM 指向本地 mock无外网也能跑且快、可控
DATA=/var/tmp/kstress
mkdir -p $DATA
# 先跑一次实例建出 config.db再写入下面这些键也可直接复用现成目录
# core.llm.provider=mock
# core.llm.sources.mock.base_url=http://127.0.0.1:9099/v1
# core.llm.sources.mock.model=mock api_key=mock adapter=openai
# core.llm.sources.mock.adapter_path=adapters/openai.lua priority=100
# core.defaults.llm_endpoints=http://127.0.0.1:9099/v1/models # 探活端点(探活用 HEAD
# core.defaults.rollback.max_retries=100000 auto_rollback=false
# 并把 core.llm.sources.deepseek* 删掉netns 里它不可达,会让 agent 被判 degraded → rollback 循环)
# 2) 起 mock LLM + 内核(都在同一个私有 netns 里)
MOCK_DELAY_MS=300 MOCK_CHUNKS=8 ./launch.sh /tmp/homed-stress
# 3) 取认证密钥cli 插件回落到 webui.api_key
export KCLI_KEY=$(sqlite3 $DATA/config.db "select value from config_webui where key='api_key';")
# 4) 压
./kcli.py $DATA/cli.sock stats full # 看内核状态(含 scheduler 计数)
./stress.py $DATA/cli.sock "$KCLI_KEY" 16 12 4 20 dense 0.05 # 16 连接×12 输入 + 4 线程×20 中断
./stress.py $DATA/cli.sock "$KCLI_KEY" 12 12 1 25 0.0 3.0 # 稀疏中断 ⇒ 压抢占/挂起/恢复
./stress.py $DATA/cli.sock "$KCLI_KEY" 1 1 0 0 resident 0.0 # 驻留子全链路mock 见 !resident 标记)
```
## 读结果
| 指标 | 含义 |
|---|---|
| `executed` / `rejected` | 任务执行数 / 被拒数(压力下应为 0 |
| 峰值 `峰值_队列` / `峰值_待处理中断` | 采样到的最大排队深度 / 待处理中断数 |
| `suspended` / `resumed` / `preempted` | 抢占三件套。**中断要放稀**才会打在高优先级任务上密集中断会互相同级L4 vs L4不抢占数值会很低 |
| `max_suspend_depth` | 中断栈结构上界(= 中断级数 4 |
## 已知坑
- 探活用 **HEAD**`internal/network/monitor.go`mock 必须实现 `do_HEAD`,否则 501 → 判不可达 → agent degraded → rollback 循环(会 reset 合并目录)。
- CLI 协议**每条连接串行**处理(`handleChat` 阻塞到本次响应产出):单连接狂发只会排成一条线,造不出队列压力 —— 必须多连接。
- 认证:连接后先发 `/auth <key>`,否则一切命令返回 `unauthorized`
- `mockllm.py``!resident` / `!notify` 标记用来让 mock 回工具调用,从而在真内核里驱动 `resident_agents` / `notify_parent`

View File

@ -0,0 +1,81 @@
import json, os, socket, sys, threading, time
SOCK = sys.argv[1]
KEY = os.environ.get("KCLI_KEY", "")
def connect(auth=True):
s = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
s.settimeout(20)
s.connect(SOCK)
if auth and KEY:
s.sendall(("/auth %s\n" % KEY).encode())
buf = b""
while b"\n" not in buf:
c = s.recv(65536)
if not c:
break
buf += c
if b"authenticated" not in buf:
raise RuntimeError("auth failed: %r" % buf[:120])
return s
def kernel_stats(sock=None, timeout=15):
"""连一条新连接查 /kernel返回解析后的 dict。"""
s = sock or connect()
try:
s.sendall(b"/kernel\n")
buf = b""
deadline = time.time() + timeout
while time.time() < deadline:
chunk = s.recv(65536)
if not chunk:
break
buf += chunk
for line in buf.split(b"\n"):
if not line.strip():
continue
try:
obj = json.loads(line.decode("utf-8", "replace"))
except Exception:
continue
if obj.get("type") == "response" and obj.get("content", "").lstrip().startswith("{"):
return json.loads(obj["content"])
raise RuntimeError("no /kernel response")
finally:
if sock is None:
s.close()
def blast(n_queued, n_interrupt, tag):
"""一条连接狂发:排队输入与中断交错。响应在后台线程里丢弃,避免缓冲阻塞。"""
s = connect()
threading.Thread(target=lambda: drain(s), daemon=True).start()
sent = 0
for i in range(max(n_queued, n_interrupt)):
if i < n_queued:
s.sendall(("queued-%s-%d\n" % (tag, i)).encode())
sent += 1
if i < n_interrupt:
s.sendall(("/interrupt intr-%s-%d\n" % (tag, i)).encode())
sent += 1
return sent
def drain(s):
try:
while s.recv(65536):
pass
except Exception:
pass
if __name__ == "__main__":
cmd = sys.argv[2] if len(sys.argv) > 2 else "stats"
if cmd == "stats":
st = kernel_stats()
if len(sys.argv) > 3 and sys.argv[3] == "full":
print(json.dumps(st, ensure_ascii=False))
else:
sched = st.get("scheduler") or st
print(json.dumps(sched, ensure_ascii=False))
elif cmd == "blast":
n = blast(int(sys.argv[3]), int(sys.argv[4]), sys.argv[5] if len(sys.argv) > 5 else "b")
print("sent=%d" % n)

17
scripts/kernel-stress/launch.sh Executable file
View File

@ -0,0 +1,17 @@
#!/bin/bash
# 在**私有 netns** 里启动mock LLM + 完整内核实例。
# 私有 netns 的意义:生产实例占着 *:8080 / *:9890 / *:9876插件设了 SO_REUSEADDR
# 同机再起一个实例会在 127.0.0.1 上与之并存绑定(第一次实测就这么抢到了 9890 约 1 分钟)。
# netns 里只有 lo结构上不可能碰到生产端口。
set -e
DATA=/var/tmp/gotmp/kstress
DRIVE=/var/tmp/gotmp/kstress-drive
DELAY=${MOCK_DELAY_MS:-300}
CHUNKS=${MOCK_CHUNKS:-8}
unshare -n bash -c "
ip link set lo up
nohup env MOCK_DELAY_MS=$DELAY MOCK_CHUNKS=$CHUNKS python3 $DRIVE/mockllm.py > $DATA/mock.log 2>&1 &
echo \$! > $DATA/mock.pid
nohup $1 -data $DATA -webui 127.0.0.1:18080 > $DATA/boot.log 2>&1 &
echo \$! > $DATA/pid
"

View File

@ -0,0 +1,102 @@
#!/usr/bin/env python3
"""最小 OpenAI 兼容 mock可控延迟 + SSE 分块 + 可选工具调用。
用途:给压力测试一个**快且可控**的 LLM —— 没有它,无外网的 netns 里每条输入
都要走 provider 重试≈2 分钟/条),既慢又压不出调度器行为。
"""
import json, os, time
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
DELAY_MS = int(os.environ.get("MOCK_DELAY_MS", "300"))
CHUNKS = int(os.environ.get("MOCK_CHUNKS", "8")) # SSE 分块数:越多,流式段越长(可被中断的窗口越大)
class H(BaseHTTPRequestHandler):
protocol_version = "HTTP/1.1"
def log_message(self, *a):
pass
def _json(self, obj, code=200):
b = json.dumps(obj).encode()
self.send_response(code)
self.send_header("Content-Type", "application/json")
self.send_header("Content-Length", str(len(b)))
self.end_headers()
self.wfile.write(b)
def do_HEAD(self):
# 内核探活用 HEAD见 internal/network/monitor.go CheckOnce
# 不实现它 → BaseHTTPRequestHandler 回 501 → 判为不可达 → agent degraded/rollback。
self.send_response(200)
self.send_header("Content-Length", "0")
self.end_headers()
def do_GET(self):
self._json({"object": "list", "data": [{"id": "mock", "object": "model"}]})
def do_POST(self):
n = int(self.headers.get("Content-Length", "0"))
body = json.loads(self.rfile.read(n) or b"{}")
msgs = body.get("messages") or []
text = ""
for m in reversed(msgs):
if m.get("role") == "user":
c = m.get("content")
text = c if isinstance(c, str) else json.dumps(c, ensure_ascii=False)
break
# 标记 !resident ⇒ 回一个工具调用,用于在**真实内核**里驱动驻留子工具链。
if "!resident" in text and not any(m.get("role") == "tool" for m in msgs):
tc = {"id": "call_mock_1", "type": "function",
"function": {"name": "resident_agents",
"arguments": json.dumps({"action": "create", "id": "r1",
"task_prompt": "驻留子任务:统计一下 !notify",
"input_chs": "cli"}, ensure_ascii=False)}}
return self._respond(body, content=None, tool_calls=[tc])
if "!notify" in text and not any(m.get("role") == "tool" for m in msgs):
tc = {"id": "call_mock_2", "type": "function",
"function": {"name": "notify_parent",
"arguments": json.dumps({"text": "mock 汇报:子已完成统计"}, ensure_ascii=False)}}
return self._respond(body, content=None, tool_calls=[tc])
return self._respond(body, content="mock-ok:" + text[:40])
def _respond(self, body, content=None, tool_calls=None):
if body.get("stream"):
self.send_response(200)
self.send_header("Content-Type", "text/event-stream")
self.send_header("Cache-Control", "no-cache")
self.send_header("Transfer-Encoding", "chunked")
self.end_headers()
def emit(delta):
data = json.dumps({"id": "mock", "object": "chat.completion.chunk",
"model": "mock", "choices": [{"index": 0, "delta": delta}]})
self._chunk(("data: " + data + "\n\n").encode())
if tool_calls:
emit({"role": "assistant", "tool_calls": tool_calls})
time.sleep(DELAY_MS / 1000.0)
if content:
emit({"role": "assistant", "content": content[: max(1, len(content) // CHUNKS)]})
per = max(1, len(content) // CHUNKS)
for i in range(per, len(content), per):
time.sleep(DELAY_MS / 1000.0 / CHUNKS)
emit({"content": content[i:i + per]})
self._chunk(b"data: [DONE]\n\n")
self._chunk(b"")
return
msg = {"role": "assistant", "content": content}
if tool_calls:
msg["tool_calls"] = tool_calls
msg["content"] = None
time.sleep(DELAY_MS / 1000.0)
self._json({"id": "mock", "object": "chat.completion", "model": "mock",
"choices": [{"index": 0, "message": msg, "finish_reason": "stop"}],
"usage": {"prompt_tokens": 10, "completion_tokens": 10, "total_tokens": 20}})
def _chunk(self, b):
self.wfile.write(("%x\r\n" % len(b)).encode() + b + b"\r\n")
if __name__ == "__main__":
port = int(os.environ.get("MOCK_PORT", "9099"))
ThreadingHTTPServer(("127.0.0.1", port), H).serve_forever()

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