mirror of
https://gitcode.com/JianFeeeee/HomeAgent.git
synced 2026-09-21 09:28:14 +00:00
本轮压测是从 /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;驻留子全链路 + 优雅退出不留孤儿。
103 lines
4.7 KiB
Python
103 lines
4.7 KiB
Python
#!/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()
|