diff --git a/scripts/kernel-stress/README.md b/scripts/kernel-stress/README.md new file mode 100644 index 0000000..55b782c --- /dev/null +++ b/scripts/kernel-stress/README.md @@ -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 `,否则一切命令返回 `unauthorized`。 +- `mockllm.py` 的 `!resident` / `!notify` 标记用来让 mock 回工具调用,从而在真内核里驱动 `resident_agents` / `notify_parent`。 diff --git a/scripts/kernel-stress/kcli.py b/scripts/kernel-stress/kcli.py new file mode 100644 index 0000000..39d7287 --- /dev/null +++ b/scripts/kernel-stress/kcli.py @@ -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) diff --git a/scripts/kernel-stress/launch.sh b/scripts/kernel-stress/launch.sh new file mode 100755 index 0000000..4c4c9bf --- /dev/null +++ b/scripts/kernel-stress/launch.sh @@ -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 +" diff --git a/scripts/kernel-stress/mockllm.py b/scripts/kernel-stress/mockllm.py new file mode 100644 index 0000000..ce52eb1 --- /dev/null +++ b/scripts/kernel-stress/mockllm.py @@ -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() diff --git a/scripts/kernel-stress/stress.py b/scripts/kernel-stress/stress.py new file mode 100644 index 0000000..35b2d07 --- /dev/null +++ b/scripts/kernel-stress/stress.py @@ -0,0 +1,144 @@ +#!/usr/bin/env python3 +"""二进制级调度器压力:多并发连接轰炸排队输入 + L4 中断,并采样峰值。 + +用法: stress.py [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))