From 316240352f8ed907203de4495bf75636135270d7 Mon Sep 17 00:00:00 2001 From: earayu Date: Wed, 2 Sep 2026 23:23:18 +0800 Subject: [PATCH] =?UTF-8?q?=E5=A2=9E=E5=8A=A0=E6=8E=A7=E5=88=B6=E9=9D=A2?= =?UTF-8?q?=E5=AF=86=E5=BA=A6=E5=8E=8B=E6=B5=8B=E8=84=9A=E6=9C=AC?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 仓库里固定入口,只打 density* 键,退出默认清理;CI 跑前缀和停条件单测。 --- .github/workflows/ci.yml | 3 + README.md | 20 +++ docs/stability.md | 2 + scripts/density.py | 314 ++++++++++++++++++++++++++++++++++ tests/density/test_density.py | 49 ++++++ 5 files changed, 388 insertions(+) create mode 100755 scripts/density.py create mode 100644 tests/density/test_density.py diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 41e3365..2412371 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -27,6 +27,9 @@ jobs: - name: Ticket vectors match the generator run: python3 tests/vectors/generate.py --check + - name: Density script safety checks + run: python3 tests/density/test_density.py + - name: Install working-directory: host-agent run: npm ci --no-fund --no-audit diff --git a/README.md b/README.md index 26d2a19..5a68b12 100644 --- a/README.md +++ b/README.md @@ -23,7 +23,9 @@ host-agent/ Node 服务(TypeScript,零运行时依赖,esbuild 打成单 src/ gateway / control / supervisor / ticket / config test/ node:test 单元与集成测试(内置 fake dsh) docs/ 架构(architecture.md)、生命周期(lifecycle.md)、稳定性(stability.md)、配对与认证(pairing.md)、宿主运营配置(host-settings.md)、ApeMind 能力结合(apemind-integration.md)、LTS(lts.md) +scripts/ 运营脚本(密度压测 density.py) tests/vectors 票据 golden vectors(Python 生成,双端测试共用) +tests/density 密度脚本的前缀与停条件单测 deploy/ Kubernetes Helm chart(独立发布,不绑 ApeMind 主 chart) Dockerfile 运行镜像(node:22-bookworm-slim + 锁版本 dsh + apemind CLI + host-agent) ``` @@ -65,6 +67,24 @@ node dist/host-agent.mjs 闲置回收、会话 cookie 寿命、就绪窗口、停进程宽限和实例容量在 `/data/settings.json`(完整五键),出厂文件见镜像内 `/usr/local/share/apemind-computer/settings.json`。配对后用 `GET/PUT /v1/runtime` 的 `settings` 热更新。详见 [docs/host-settings.md](docs/host-settings.md)。 +## 密度压测 + +`scripts/density.py` 打控制面 `PUT /v1/instances/{key}`,按批次拉起真实 dsh。 +只接受以 `density` 开头的前缀,退出默认停掉并删除本次建的键。令牌从 +`COMPUTER_CONTROL_TOKEN` 或 `--token-file` 读,不进日志。 + +```bash +# 本地(先 port-forward 控制口 :9090) +COMPUTER_CONTROL_TOKEN=dev-token python3 scripts/density.py \ + --url http://127.0.0.1:9090 --count 20 --batch 10 + +python3 tests/density/test_density.py +python3 scripts/density.py --cleanup-only --url http://127.0.0.1:9090 +``` + +不要对生产控制面跑。容器内存顶、cgroup 是否降级、闲置回收仍按 +[docs/stability.md](docs/stability.md) 生效;本脚本不测网关数据面。 + ## 镜像 镜像只在 GitHub Actions 构建(推 tag `v*.*.*` 触发),不在本地构建。dsh 版本在 diff --git a/docs/stability.md b/docs/stability.md index e4a4262..344da3d 100644 --- a/docs/stability.md +++ b/docs/stability.md @@ -166,6 +166,8 @@ cgroup v2:一个节点一旦打开控制器,自己身上就不能再挂进 保护管家靠三道:每户一个松的 `memory.max`;`agent/` 的 `memory.min` / `memory.low`;容器 limit 仍是最后一道。活跃数 × 每户顶不要长期大于容器;闲置回收已经在压活跃数。 +要在真实宿主上数「同时能跑多少台」,用仓库里的 `scripts/density.py` 打控制面,不要手写一次性脚本。它只动 `density*` 键,默认退出就清掉。 + 某户合法就要用 8Gi(本地编大项目),同容器模型里只有两条路:把全宿主默认顶抬高,或把这户挪出共享池(per-user Pod)。本期两条都不走,接受「单户有顶、顶以上请停」。 ### 6.6 落地时要注意 diff --git a/scripts/density.py b/scripts/density.py new file mode 100755 index 0000000..5b450eb --- /dev/null +++ b/scripts/density.py @@ -0,0 +1,314 @@ +#!/usr/bin/env python3 +"""Density probe for the computer-host control API. + +Starts N instances with PUT /v1/instances/{key} desired=running, prints one +JSON object per line, then stops and deletes only keys under the given prefix. + +Token is read from COMPUTER_CONTROL_TOKEN or --token-file. It is never printed. + + COMPUTER_CONTROL_TOKEN=dev-token python3 scripts/density.py \\ + --url http://127.0.0.1:9090 --count 20 --batch 10 + +Ctrl-C without --keep still cleans up keys this run created. --cleanup-only +deletes leftover prefix keys from a previous interrupted run. + +Only prefixes that start with "density" are accepted, so this will not touch +ApeMind instance keys (ci*). Do not point --url at production. +""" + +from __future__ import annotations + +import argparse +import json +import os +import re +import signal +import statistics +import sys +import time +import urllib.error +import urllib.request +from concurrent.futures import ThreadPoolExecutor, as_completed +from typing import Any, NamedTuple + +USER_ID_RE = re.compile(r"^[A-Za-z0-9_-]{1,64}$") +PREFIX_RE = re.compile(r"^density[A-Za-z0-9_-]{0,24}$") + + +class BatchResult(NamedTuple): + ok: int + fail: int + lat_s: dict[str, float | None] + fail_codes: dict[str, int] + + +def assert_safe_prefix(prefix: str) -> str: + if not PREFIX_RE.fullmatch(prefix): + raise ValueError( + 'prefix must match ^density[A-Za-z0-9_-]{0,24}$ so the probe cannot touch ci* homes' + ) + return prefix + + +def instance_keys(prefix: str, start: int, count: int) -> list[str]: + assert_safe_prefix(prefix) + if start < 1 or count < 1: + raise ValueError("start and count must be >= 1") + keys: list[str] = [] + width = max(4, len(str(start + count - 1))) + for index in range(start, start + count): + key = f"{prefix}{index:0{width}d}" + if not USER_ID_RE.fullmatch(key): + raise ValueError(f"generated key is not a valid instance id: {key}") + keys.append(key) + return keys + + +def owns_key(prefix: str, key: str) -> bool: + return USER_ID_RE.fullmatch(key) is not None and key.startswith(prefix) + + +def latency_stats(samples: list[float]) -> dict[str, float | None]: + if not samples: + return {"min": None, "p50": None, "p95": None, "max": None} + ordered = sorted(samples) + p95 = ordered[max(0, int(len(ordered) * 0.95) - 1)] + return { + "min": round(min(samples), 3), + "p50": round(statistics.median(samples), 3), + "p95": round(p95, 3), + "max": round(max(samples), 3), + } + + +def rss_summary(instances: list[dict[str, Any]], prefix: str) -> dict[str, Any]: + running = [ + inst + for inst in instances + if inst.get("status") == "running" and owns_key(prefix, str(inst.get("user_id") or "")) + ] + rss = [int(inst.get("rss_bytes") or 0) for inst in running] + rss = [n for n in rss if n > 0] + return { + "density_running": len(running), + "rss_mib_min": round(min(rss) / 1048576, 1) if rss else None, + "rss_mib_p50": round(statistics.median(rss) / 1048576, 1) if rss else None, + "rss_mib_max": round(max(rss) / 1048576, 1) if rss else None, + "rss_mib_sum": round(sum(rss) / 1048576, 1) if rss else 0, + } + + +def should_stop(batch: BatchResult, health: dict[str, Any], batch_size: int) -> str | None: + if batch.fail >= max(3, batch_size // 2): + return "too many failures in this batch" + instances = health.get("instances") if isinstance(health.get("instances"), dict) else {} + running = int(instances.get("running") or 0) + maximum = int(instances.get("max") or 0) + if maximum and running >= maximum: + return "at max_instances" + if "507:" in "".join(batch.fail_codes): + return "host returned 507 capacity" + return None + + +def summarize_rows(rows: list[tuple[str, int, float, dict[str, Any], str]]) -> BatchResult: + ok_rows = [row for row in rows if row[1] == 200 and row[3].get("status") == "running"] + fail_rows = [row for row in rows if row not in ok_rows] + fail_codes: dict[str, int] = {} + for _key, status, _elapsed, body, err in fail_rows: + reason = f"{status}:{(body.get('error') or err or body.get('status') or 'unknown')}" + fail_codes[reason] = fail_codes.get(reason, 0) + 1 + return BatchResult( + ok=len(ok_rows), + fail=len(fail_rows), + lat_s=latency_stats([row[2] for row in ok_rows]), + fail_codes=fail_codes, + ) + + +def emit(event: str, **payload: Any) -> None: + print(json.dumps({"event": event, **payload}, ensure_ascii=False), flush=True) + + +def load_token(token_file: str | None) -> str: + if token_file: + token = open(token_file, encoding="utf-8").read().strip() + else: + token = (os.environ.get("COMPUTER_CONTROL_TOKEN") or "").strip() + if not token: + raise SystemExit("set COMPUTER_CONTROL_TOKEN or pass --token-file") + return token + + +class ControlClient: + def __init__(self, url: str, token: str, timeout: float) -> None: + self.url = url.rstrip("/") + self.token = token + self.timeout = timeout + + def request( + self, + method: str, + path: str, + body: dict[str, Any] | None = None, + timeout: float | None = None, + ) -> tuple[int, dict[str, Any], str]: + data = None if body is None else json.dumps(body).encode() + request = urllib.request.Request( + f"{self.url}{path}", + data=data, + method=method, + headers={ + "Authorization": f"Bearer {self.token}", + "Content-Type": "application/json", + }, + ) + try: + with urllib.request.urlopen(request, timeout=timeout or self.timeout) as resp: + raw = resp.read() + payload = json.loads(raw) if raw else {} + return resp.status, payload if isinstance(payload, dict) else {}, "" + except urllib.error.HTTPError as exc: + raw = exc.read() + try: + payload = json.loads(raw) if raw else {} + except json.JSONDecodeError: + payload = {} + text = raw.decode("utf-8", "replace")[:200] + return exc.code, payload if isinstance(payload, dict) else {}, text + + def healthz(self) -> dict[str, Any]: + status, body, err = self.request("GET", "/healthz", timeout=20) + if status != 200: + raise SystemExit(f"healthz {status} {err or body}") + return body + + def list_instances(self) -> list[dict[str, Any]]: + status, body, err = self.request("GET", "/v1/instances", timeout=20) + if status != 200: + raise SystemExit(f"list instances {status} {err or body}") + items = body.get("instances") + return items if isinstance(items, list) else [] + + def ensure(self, key: str, desired: str) -> tuple[str, int, float, dict[str, Any], str]: + started = time.monotonic() + status, body, err = self.request( + "PUT", + f"/v1/instances/{key}", + {"desired": desired}, + ) + return key, status, time.monotonic() - started, body, err + + def delete(self, key: str) -> tuple[str, int]: + status, _body, _err = self.request("DELETE", f"/v1/instances/{key}", timeout=30) + return key, status + + +def run_pool(fn, items: list[str], concurrency: int): + rows = [] + with ThreadPoolExecutor(max_workers=concurrency) as pool: + futs = {pool.submit(fn, item): item for item in items} + for fut in as_completed(futs): + rows.append(fut.result()) + return rows + + +def cleanup(client: ControlClient, keys: list[str], prefix: str, concurrency: int) -> None: + owned = [key for key in keys if owns_key(prefix, key)] + if not owned: + return + emit("cleanup", phase="stop", count=len(owned)) + run_pool(lambda key: client.ensure(key, "stopped"), owned, concurrency) + emit("cleanup", phase="delete", count=len(owned)) + run_pool(client.delete, owned, concurrency) + emit("cleanup", phase="done", healthz=client.healthz()) + + +def parse_args(argv: list[str] | None = None) -> argparse.Namespace: + parser = argparse.ArgumentParser(description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter) + parser.add_argument("--url", default=os.environ.get("COMPUTER_CONTROL_URL", "http://127.0.0.1:9090")) + parser.add_argument("--token-file", help="file containing the control token; otherwise COMPUTER_CONTROL_TOKEN") + parser.add_argument("--prefix", default="density", help="instance key prefix; must start with density") + parser.add_argument("--count", type=int, default=20, help="how many instances to start") + parser.add_argument("--start", type=int, default=1, help="first numeric suffix") + parser.add_argument("--batch", type=int, default=10, help="keys started per wave") + parser.add_argument("--concurrency", type=int, default=5, help="in-flight ensure calls") + parser.add_argument("--timeout", type=float, default=120, help="per-ensure timeout seconds") + parser.add_argument("--keep", action="store_true", help="leave started instances running") + parser.add_argument( + "--cleanup-only", + action="store_true", + help="stop and delete existing prefix keys, then exit", + ) + return parser.parse_args(argv) + + +def main(argv: list[str] | None = None) -> int: + args = parse_args(argv) + try: + prefix = assert_safe_prefix(args.prefix) + except ValueError as exc: + raise SystemExit(str(exc)) from exc + if args.batch < 1 or args.concurrency < 1: + raise SystemExit("--batch and --concurrency must be >= 1") + + client = ControlClient(args.url, load_token(args.token_file), args.timeout) + created: list[str] = [] + keep = args.keep + + def handle_stop(_signum: int, _frame: Any) -> None: + emit("interrupt", keep=keep) + if not keep: + cleanup(client, created, prefix, args.concurrency) + raise SystemExit(130) + + signal.signal(signal.SIGINT, handle_stop) + signal.signal(signal.SIGTERM, handle_stop) + + if args.cleanup_only: + leftovers = [str(inst.get("user_id") or "") for inst in client.list_instances()] + leftovers = [key for key in leftovers if owns_key(prefix, key)] + cleanup(client, leftovers, prefix, args.concurrency) + return 0 + + emit("baseline", **client.healthz(), url=args.url, prefix=prefix, count=args.count) + try: + remaining = args.count + cursor = args.start + while remaining > 0: + size = min(args.batch, remaining) + keys = instance_keys(prefix, cursor, size) + wall = time.monotonic() + rows = run_pool(lambda key: client.ensure(key, "running"), keys, args.concurrency) + created.extend(keys) + batch = summarize_rows(rows) + health = client.healthz() + emit( + "batch", + first=keys[0], + last=keys[-1], + ok=batch.ok, + fail=batch.fail, + lat_s=batch.lat_s, + fail_codes=batch.fail_codes, + batch_wall_s=round(time.monotonic() - wall, 3), + instances=health.get("instances"), + load1=health.get("load1"), + rss=rss_summary(client.list_instances(), prefix), + ) + reason = should_stop(batch, health, size) + if reason: + emit("stop", reason=reason) + break + remaining -= size + cursor += size + finally: + if created and not keep: + cleanup(client, created, prefix, args.concurrency) + elif created and keep: + emit("kept", count=len(created), prefix=prefix) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/tests/density/test_density.py b/tests/density/test_density.py new file mode 100644 index 0000000..a0a6fcd --- /dev/null +++ b/tests/density/test_density.py @@ -0,0 +1,49 @@ +#!/usr/bin/env python3 +from __future__ import annotations + +import importlib.util +import sys +import unittest +from pathlib import Path + +ROOT = Path(__file__).resolve().parents[2] +SPEC = importlib.util.spec_from_file_location("density", ROOT / "scripts" / "density.py") +assert SPEC and SPEC.loader +density = importlib.util.module_from_spec(SPEC) +sys.modules[SPEC.name] = density +SPEC.loader.exec_module(density) + + +class DensitySafetyTest(unittest.TestCase): + def test_prefix_must_start_with_density(self) -> None: + self.assertEqual(density.assert_safe_prefix("density"), "density") + self.assertEqual(density.assert_safe_prefix("density-stg"), "density-stg") + for prefix in ("", "ci", "user", "den", "Density", "density/" ): + with self.assertRaises(ValueError): + density.assert_safe_prefix(prefix) + + def test_keys_stay_under_host_id_rules(self) -> None: + keys = density.instance_keys("density", 1, 3) + self.assertEqual(keys, ["density0001", "density0002", "density0003"]) + self.assertTrue(all(density.owns_key("density", key) for key in keys)) + self.assertFalse(density.owns_key("density", "cic4da7207703ac076")) + self.assertFalse(density.owns_key("density", "density0001/../etc")) + + def test_stop_when_capacity_or_many_failures(self) -> None: + full = density.BatchResult(ok=0, fail=1, lat_s={}, fail_codes={"507:max instances reached": 1}) + self.assertEqual( + density.should_stop(full, {"instances": {"running": 10, "max": 200}}, 10), + "host returned 507 capacity", + ) + many = density.BatchResult(ok=2, fail=8, lat_s={}, fail_codes={"409:start failed": 8}) + self.assertEqual(density.should_stop(many, {"instances": {"running": 12, "max": 200}}, 10), "too many failures in this batch") + ok = density.BatchResult(ok=10, fail=0, lat_s={}, fail_codes={}) + self.assertIsNone(density.should_stop(ok, {"instances": {"running": 10, "max": 200}}, 10)) + self.assertEqual( + density.should_stop(ok, {"instances": {"running": 200, "max": 200}}, 10), + "at max_instances", + ) + + +if __name__ == "__main__": + unittest.main()