已合并
系统自评估与建议能力——观测补齐+规则引擎+LLM 分析 #493
王明琦创建于 6 天前
系统自评估与建议能力——观测补齐+规则引擎+LLM 分析 #493
已合并
共 34 个文件变更+4494-110
| @@ -609,7 +609,8 @@ flowchart LR | |||
| 609 | > **Redis Cluster 兼容(2026-08-29)**:两个前缀整体为 **hash tag**(`{xxx}`),模块全部键 | 609 | > **Redis Cluster 兼容(2026-08-29)**:两个前缀整体为 **hash tag**(`{xxx}`),模块全部键 |
| 610 | > 落同一 slot——多键 Lua 的原子语义在 cluster 分片下保持成立;选主抽签键为 | 610 | > 落同一 slot——多键 Lua 的原子语义在 cluster 分片下保持成立;选主抽签键为 |
| 611 | > `{agent_runtime:job:<job>}:winner/candidates:{epoch}`(同槽),执行锁键 | 611 | > `{agent_runtime:job:<job>}:winner/candidates:{epoch}`(同槽),执行锁键 |
| 612 | -> `agent_runtime:job:<job>` 不变(单键操作)。连接串用 `redis+cluster://` scheme 构造 | 612 | +> `agent_runtime:job:<job>` 不变(单键操作)。自评估域前缀 `{agent_runtime:eval}`(§5.3, |
| 613 | +> 2026-09)同款单槽 tag、全单键命令无 Lua。连接串用 `redis+cluster://` scheme 构造 | ||
| 613 | > 集群客户端(cluster 只有 db 0)。`state.eval` 把 prefix 同时声明为 `KEYS[1]`(路由锚, | 614 | > 集群客户端(cluster 只有 db 0)。`state.eval` 把 prefix 同时声明为 `KEYS[1]`(路由锚, |
| 614 | > 防 `numkeys=0` 随机路由到非归属节点)。单实例/哨兵下 `{}` 无语义,同一套键名兼容两种 | 615 | > 防 `numkeys=0` 随机路由到非归属节点)。单实例/哨兵下 `{}` 无语义,同一套键名兼容两种 |
| 615 | > 部署;背景与验证见 `docs/feature/2026-08-redis-cluster.md`。 | 616 | > 部署;背景与验证见 `docs/feature/2026-08-redis-cluster.md`。 |
| @@ -726,6 +727,20 @@ flowchart TB | |||
| 726 | | `lock:rm:reconcile` | STRING(NX EX) | 选主标记 | 孤儿对账 sweeper 选主 | | 727 | | `lock:rm:reconcile` | STRING(NX EX) | 选主标记 | 孤儿对账 sweeper 选主 | |
| 727 | | `lock:rm:watch` | STRING(NX EX) | 选主标记 | K8s Watch + 死 Pod 轮询 + 健康 SSE 探测(场景 N)选主 | | 728 | | `lock:rm:watch` | STRING(NX EX) | 选主标记 | K8s Watch + 死 Pod 轮询 + 健康 SSE 探测(场景 N)选主 | |
| 728 | 729 | ||
| 730 | +### 5.3 系统自评估(prefix `{agent_runtime:eval}:`,2026-09) | ||
| 731 | + | ||
| 732 | +评估数据域(非编排态;spec 见 `docs/spec/evaluation.md`):route/acquire 热路径计数经每副本 | ||
| 733 | +内存缓冲 5s 批量聚合,RM 后台状态变迁直写;sys_sample(30s)采样池态+计数快照落 ZSET, | ||
| 734 | +sys_eval(300s)跑确定性规则(+可选 LLM 分析)产报告。**全单键命令、零 Lua**;报告只读 | ||
| 735 | +产出,不改任何配置(建议人审后经 Manager config_sync 应用)。 | ||
| 736 | + | ||
| 737 | +| 键 | 类型 | TTL | 语义 | | ||
| 738 | +|---|---|---|---| | ||
| 739 | +| `sample:scope:{sid}` | ZSET(member=紧凑 JSON,score=秒) | 25h(每采样刷新) | per-scope 趋势采样(t/p/i/d/s/w + 计数快照差分源) | | ||
| 740 | +| `ct:scope:{sid}` | HASH | 25h(每次写刷新) | per-scope 全副本聚合计数(route_*/acq_*/ev_*) | | ||
| 741 | +| `report:latest` | STRING | 无 TTL | 最近评估报告(findings/trend/caveats) | | ||
| 742 | +| `report:history` | ZSET(瘦身条目) | 30d(保 200) | 报告历史(只留 summary) | | ||
| 743 | + | ||
| 729 | --- | 744 | --- |
| 730 | 745 | ||
| 731 | ## 6. 场景清单与详细时序 | 746 | ## 6. 场景清单与详细时序 |
| @@ -1311,6 +1326,9 @@ M6(server 模式)已完成开发与真环境端到端验收,实现与本文的 | |||
| 1311 | 1326 | ||
| 1312 | - **`LUA_WAITER_GATE`(§5.1 键表 / §6.2 场景 F)**:实现期补充的第 7 个 SM Lua 脚本。初稿「先 SCARD 再 SADD」的入队判定在并发同时到达时全部读到旧计数而超收(真环境验收场景 F 发现的竞态),改为原子闸门后稳态 `SCARD ≤ max_waiters` 恒成立。全文见 SM 设计 §5.1。**(2026-09:已随场景 F 快失败改造整体拆除,连同 waiters 键/free 通道,见 §9.2)** | 1327 | - **`LUA_WAITER_GATE`(§5.1 键表 / §6.2 场景 F)**:实现期补充的第 7 个 SM Lua 脚本。初稿「先 SCARD 再 SADD」的入队判定在并发同时到达时全部读到旧计数而超收(真环境验收场景 F 发现的竞态),改为原子闸门后稳态 `SCARD ≤ max_waiters` 恒成立。全文见 SM 设计 §5.1。**(2026-09:已随场景 F 快失败改造整体拆除,连同 waiters 键/free 通道,见 §9.2)** |
| 1313 | - **场景 N(半死 Pod 健康探测)**:机制已实现且有单测覆盖;端到端验收**暂缓**——当前 AgentServer 镜像在 SSE 端口对 `GET /health` 返回 426(要求协议升级),不满足 §6.2 场景 N 的固定约定,待 AgentServer 原生支持后补验。 | 1328 | - **场景 N(半死 Pod 健康探测)**:机制已实现且有单测覆盖;端到端验收**暂缓**——当前 AgentServer 镜像在 SSE 端口对 `GET /health` 返回 426(要求协议升级),不满足 §6.2 场景 N 的固定约定,待 AgentServer 原生支持后补验。 |
| 1329 | +- **系统自评估(§5.3,2026-09)**:evaluation 子包 + sys_sample/sys_eval 两选主 job + | ||
| 1330 | + `/visualization/{history,evaluation}` 端点已实现,单测 51 例 + 集成用例全绿;真环境 | ||
| 1331 | + 冒烟阶段 13c 已固化(e2e_hld_acceptance.py),门禁结果见 feature 篇目。 | ||
| 1314 | - 场景 A–L 已在真 Redis + MySQL + K8s 环境端到端验收通过(用例固化为 `applications/agent_runtime/scripts/e2e_hld_acceptance.py`,经 `scripts/integration_smoke.sh` 调用,可作部署后回归)。 | 1332 | - 场景 A–L 已在真 Redis + MySQL + K8s 环境端到端验收通过(用例固化为 `applications/agent_runtime/scripts/e2e_hld_acceptance.py`,经 `scripts/integration_smoke.sh` 调用,可作部署后回归)。 |
| 1315 | 1333 | ||
| 1316 | ### 9.1 多副本验收(2026-08-18 更新,M7) | 1334 | ### 9.1 多副本验收(2026-08-18 更新,M7) |
| @@ -12,7 +12,9 @@ | |||
| 12 | | [service-core.md](service-core.md) | 组装(main)/CLI/配置(`AGENT_RUNTIME_*`)/错误码契约/字段分类/工具/部署 | 改装配、配置、错误契约、部署时 | | 12 | | [service-core.md](service-core.md) | 组装(main)/CLI/配置(`AGENT_RUNTIME_*`)/错误码契约/字段分类/工具/部署 | 改装配、配置、错误契约、部署时 | |
| 13 | | [session-manager.md](session-manager.md) | SM:route/touch/config_sync/config_refresh/cleanup 编排、6 个 Lua、SM 键表 | 改会话编排/配置层时 | | 13 | | [session-manager.md](session-manager.md) | SM:route/touch/config_sync/config_refresh/cleanup 编排、6 个 Lua、SM 键表 | 改会话编排/配置层时 | |
| 14 | | [resource-manager.md](resource-manager.md) | RM:acquire/后台任务/K8s 适配、6 个 Lua、RM 键表 | 改 Pod 池/扩缩容/清理时 | | 14 | | [resource-manager.md](resource-manager.md) | RM:acquire/后台任务/K8s 适配、6 个 Lua、RM 键表 | 改 Pod 池/扩缩容/清理时 | |
| 15 | +| [evaluation.md](evaluation.md) | 自评估:采样/评估两 job、`{agent_runtime:eval}` 键表、规则清单、LLM 降级矩阵 | 改自评估/趋势/报告时 | | ||
| 15 | | [e2e-test-cases.md](e2e-test-cases.md) | 全部 e2e 用例的场景/输入/预期输出 | 写或跑 e2e 时 | | 16 | | [e2e-test-cases.md](e2e-test-cases.md) | 全部 e2e 用例的场景/输入/预期输出 | 写或跑 e2e 时 | |
| 17 | +| `../api/config-plane-api.md` | 配置面对外接口文档(config_sync/config_refresh/visualization:字段表+curl+真实返回示例) | 给调用方(Claw Manager/运维/可视化前端)交付接口契约时 | | ||
| 16 | | `../design/Agent-Runtime-HLD.md` | 架构总览/接口契约/场景 A–N/Redis 键表(语义权威) | 语义不确定时 | | 18 | | `../design/Agent-Runtime-HLD.md` | 架构总览/接口契约/场景 A–N/Redis 键表(语义权威) | 语义不确定时 | |
| 17 | | `../design/session-manager-design.md` | SM 详细设计(6 个 Lua 全文) | 深挖 SM 设计动机 | | 19 | | `../design/session-manager-design.md` | SM 详细设计(6 个 Lua 全文) | 深挖 SM 设计动机 | |
| 18 | | `../design/resource-manager-design.md` | RM 详细设计(6 个 Lua 全文) | 深挖 RM 设计动机 | | 20 | | `../design/resource-manager-design.md` | RM 详细设计(6 个 Lua 全文) | 深挖 RM 设计动机 | |
| @@ -40,7 +42,8 @@ claw mgr ──config_sync──► ├─ session_manager 持 App,5 个 HTT | |||
| 40 | ``` | 42 | ``` |
| 41 | {session_manager}:… SM 编排态(会话四处/scope 闸门/候选集/注册表/路由快照) | 43 | {session_manager}:… SM 编排态(会话四处/scope 闸门/候选集/注册表/路由快照) |
| 42 | {resource_manager}:… RM 编排态(per-scope Pod 池/idle 暖池/deploy 占位/follower 等待室/选主锁) | 44 | {resource_manager}:… RM 编排态(per-scope Pod 池/idle 暖池/deploy 占位/follower 等待室/选主锁) |
| 43 | -agent_runtime:job:… 后台任务选主执行锁(main.py:_build_jobs 注册) | 45 | +{agent_runtime:eval}:… 自评估域(per-scope 计数 HASH/趋势采样 ZSET/评估报告;2026-09,见 evaluation.md) |
| 46 | +agent_runtime:job:… 后台任务选主执行锁(main.py:_build_jobs 注册,7 个) | ||
| 44 | {agent_runtime:job:…}:winner/candidates:… 选主抽签键(与执行锁同底,hash tag 同槽) | 47 | {agent_runtime:job:…}:winner/candidates:… 选主抽签键(与执行锁同底,hash tag 同槽) |
| 45 | ``` | 48 | ``` |
| 46 | 49 | ||
| @@ -51,7 +54,7 @@ agent_runtime:job:… 后台任务选主执行锁(main.py:_build_jobs 注册) | |||
| 51 | 54 | ||
| 52 | ```bash | 55 | ```bash |
| 53 | cd applications/agent_runtime | 56 | cd applications/agent_runtime |
| 54 | -uv sync --extra local && uv run pytest # 157 用例(fakeredis+SQLite+FakeK8s) | 57 | +uv sync --extra local && uv run pytest # 495 用例(fakeredis+SQLite+FakeK8s) |
| 55 | ./scripts/integration_smoke.sh # 真环境冒烟(场景 A–L;FLUSHDB 目标库,有防误刷) | 58 | ./scripts/integration_smoke.sh # 真环境冒烟(场景 A–L;FLUSHDB 目标库,有防误刷) |
| 56 | ./scripts/deploy_replicas.sh 2 .env.production.local 8091 # 宿主机双进程 | 59 | ./scripts/deploy_replicas.sh 2 .env.production.local 8091 # 宿主机双进程 |
| 57 | ./deploy/render_and_apply.sh deploy/agent_runtime.env --nodeport # K8s 生产形态 | 60 | ./deploy/render_and_apply.sh deploy/agent_runtime.env --nodeport # K8s 生产形态 |
| @@ -1430,6 +1430,76 @@ async def stage13_error_contract(c: Client, r) -> None: | |||
| 1430 | code == 200 and raw.get("cleaned") == 0, str(raw)) | 1430 | code == 200 and raw.get("cleaned") == 0, str(raw)) |
| 1431 | 1431 | ||
| 1432 | 1432 | ||
| 1433 | +async def stage13c_evaluation(c: Client, r) -> None: | ||
| 1434 | + """系统自评估数据层(2026-09):两 job + 三个可视化端点(真环境)。 | ||
| 1435 | + | ||
| 1436 | + 牙齿与等待窗:sys_sample 默认 30s、sys_eval 默认 300s——间隔从 overview | ||
| 1437 | + 配置摘要现读,PASS 判据 = eval_interval+60s(上限 330s,冒烟部署建议 | ||
| 1438 | + AGENT_RUNTIME_EVAL_INTERVAL=15 加速)内轮询到 latest 报告;llm.status | ||
| 1439 | + ∈ {disabled, ok}(真环境可能配了 LLM);响应不得泄漏 API_KEY。只调小 | ||
| 1440 | + interval env、不回拨任何指针(e2e 红线)。 | ||
| 1441 | + """ | ||
| 1442 | + print("\n== 阶段 13c:自评估 —— sys_sample/sys_eval + 可视化端点 ==") | ||
| 1443 | + root = c.base.rsplit("/api/session", 1)[0] | ||
| 1444 | + | ||
| 1445 | + async def vis(path: str) -> tuple[int, dict]: | ||
| 1446 | + resp = await c.http.get(f"{root}{path}", timeout=30.0) | ||
| 1447 | + try: | ||
| 1448 | + body = resp.json() | ||
| 1449 | + except Exception: | ||
| 1450 | + body = {} | ||
| 1451 | + return resp.status_code, body | ||
| 1452 | + | ||
| 1453 | + # 1) 两 job 已注册且至少跑过一拍(overview 是全局看板,任意副本应答) | ||
| 1454 | + code, overview = await vis("/visualization/overview") | ||
| 1455 | + jobs = {j.get("name"): j for j in overview.get("jobs", []) | ||
| 1456 | + if isinstance(j, dict)} | ||
| 1457 | + check("13c-overview 含 sys_sample/sys_eval 且 ok_ticks≥1", | ||
| 1458 | + code == 200 and {"sys_sample", "sys_eval"} <= set(jobs) | ||
| 1459 | + and jobs.get("sys_sample", {}).get("ok_ticks", 0) >= 1, | ||
| 1460 | + f"sample_ok={jobs.get('sys_sample', {}).get('ok_ticks')}") | ||
| 1461 | + | ||
| 1462 | + # 2) scope 摘要含 SM 容量字段(scope 重构后观测补齐) | ||
| 1463 | + code, scopes = await vis("/visualization/scopes") | ||
| 1464 | + row = next((s for s in scopes.get("scopes", []) if s.get("scope_id") == MAIN), | ||
| 1465 | + None) | ||
| 1466 | + check("13c-scopes 摘要含容量闸门字段(sc/session_count)", | ||
| 1467 | + code == 200 and row is not None | ||
| 1468 | + and row.get("scope_concurrency") == 3 and row.get("max_pods") == 2 | ||
| 1469 | + and "session_count" in row, str(row)[:100]) | ||
| 1470 | + | ||
| 1471 | + # 3) 历史趋势:采样 30s cadence,45s 内必有数据点 | ||
| 1472 | + async def has_points() -> bool: | ||
| 1473 | + _, hist = await vis(f"/visualization/history?scope_id={MAIN}") | ||
| 1474 | + return bool(hist.get("points")) | ||
| 1475 | + ok = await wait_until(has_points, timeout=45, interval=3) | ||
| 1476 | + _, hist = await vis(f"/visualization/history?scope_id={MAIN}") | ||
| 1477 | + point = (hist.get("points") or [{}])[0] | ||
| 1478 | + check("13c-history 采样落盘(窗口内 points≥1 且字段齐全)", | ||
| 1479 | + ok and {"t", "p", "i", "s"} <= set(point), str(point)[:80]) | ||
| 1480 | + | ||
| 1481 | + # 4) 评估报告:间隔现读,限窗轮询(不依赖部署 env 假设) | ||
| 1482 | + eval_interval = int(overview.get("config", {}).get("eval_interval") or 300) | ||
| 1483 | + budget = min(eval_interval + 60, 330) | ||
| 1484 | + | ||
| 1485 | + async def has_report() -> bool: | ||
| 1486 | + _, ev = await vis("/visualization/evaluation") | ||
| 1487 | + return bool(ev.get("latest")) | ||
| 1488 | + ok = await wait_until(has_report, timeout=budget, interval=5) | ||
| 1489 | + _, ev = await vis("/visualization/evaluation") | ||
| 1490 | + latest = ev.get("latest") or {} | ||
| 1491 | + llm_status = (latest.get("llm") or {}).get("status") | ||
| 1492 | + check("13c-evaluation 报告产出(llm ∈ {disabled,ok};findings 为 list)", | ||
| 1493 | + ok and llm_status in ("disabled", "ok") | ||
| 1494 | + and isinstance(latest.get("findings"), list) | ||
| 1495 | + and (latest.get("summary") or {}).get("active", 0) >= 1, | ||
| 1496 | + f"interval={eval_interval}s llm={llm_status} " | ||
| 1497 | + f"findings={len(latest.get('findings') or [])}") | ||
| 1498 | + text = str(latest).lower() | ||
| 1499 | + check("13c-evaluation 报告无凭证泄漏(无 api_key/Bearer 字样)", | ||
| 1500 | + ok and "api_key" not in text and "bearer" not in text, "") | ||
| 1501 | + | ||
| 1502 | + | ||
| 1433 | # ---------------------------------------------------------------- 入口 | 1503 | # ---------------------------------------------------------------- 入口 |
| 1434 | 1504 | ||
| 1435 | def _parse_args() -> argparse.Namespace: | 1505 | def _parse_args() -> argparse.Namespace: |
| @@ -1540,6 +1610,7 @@ async def main() -> None: | |||
| 1540 | await stage12_reconcile_cleanup(c, r) | 1610 | await stage12_reconcile_cleanup(c, r) |
| 1541 | await stage12b_or_branch(c, r) | 1611 | await stage12b_or_branch(c, r) |
| 1542 | await stage13_error_contract(c, r) | 1612 | await stage13_error_contract(c, r) |
| 1613 | + await stage13c_evaluation(c, r) | ||
| 1543 | finally: | 1614 | finally: |
| 1544 | await r.aclose() | 1615 | await r.aclose() |
| 1545 | 1616 | ||
| @@ -17,6 +17,7 @@ App 的 lifespan 只认一个 ctx_factory——返回 sm_sysctx,其余生命 | |||
| 17 | 17 | ||
| 18 | from __future__ import annotations | 18 | from __future__ import annotations |
| 19 | 19 | ||
| 20 | +import asyncio | ||
| 20 | import logging | 21 | import logging |
| 21 | from typing import Any | 22 | from typing import Any |
| 22 | 23 | ||
| @@ -29,6 +30,15 @@ from . import errors as app_errors | |||
| 29 | from .config import RM_KEY_PREFIX, SERVICE_PREFIX, SM_KEY_PREFIX, AgentRuntimeConfig | 30 | from .config import RM_KEY_PREFIX, SERVICE_PREFIX, SM_KEY_PREFIX, AgentRuntimeConfig |
| 30 | from .visualization_api import register_visualization_api | 31 | from .visualization_api import register_visualization_api |
| 31 | from .metrics import MetricsRegistry, request_metrics_middleware | 32 | from .metrics import MetricsRegistry, request_metrics_middleware |
| 33 | +from .evaluation.collector import ( | ||
| 34 | + FLUSH_INTERVAL_SEC, | ||
| 35 | + FLUSH_TIMEOUT_SEC, | ||
| 36 | + EvaluationCollector, | ||
| 37 | + ScopeTelemetryBuffer, | ||
| 38 | +) | ||
| 39 | +from .evaluation.evaluator import Evaluator | ||
| 40 | +from .evaluation.llm import LLMClient | ||
| 41 | +from .evaluation.state import EvaluationState | ||
| 32 | from .resource_manager.facade import ResourceManagerFacade | 42 | from .resource_manager.facade import ResourceManagerFacade |
| 33 | from .resource_manager.k8s import FakeK8sPodClient, RealK8sPodClient | 43 | from .resource_manager.k8s import FakeK8sPodClient, RealK8sPodClient |
| 34 | from .resource_manager.orchestrator import ResourceOrchestrator | 44 | from .resource_manager.orchestrator import ResourceOrchestrator |
| @@ -64,6 +74,11 @@ TICK_TIMEOUTS = { | |||
| 64 | "rm_reclaim": 60, | 74 | "rm_reclaim": 60, |
| 65 | "rm_watch": 300, | 75 | "rm_watch": 300, |
| 66 | "rm_reconcile": 300, | 76 | "rm_reconcile": 300, |
| 77 | + # sys_sample:逐 scope 单键读写,快操作 | ||
| 78 | + "sys_sample": 30, | ||
| 79 | + # sys_eval:LLM timeout(默认 60s,eval_llm_timeout 可调但须小于此值) | ||
| 80 | + # + 规则计算 + 采样窗口读,120 留余量 | ||
| 81 | + "sys_eval": 120, | ||
| 67 | } | 82 | } |
| 68 | 83 | ||
| 69 | 84 | ||
| @@ -154,9 +169,13 @@ class OrchestratorSystemContext(SystemContext): | |||
| 154 | sm_state = SessionState(self.redis) | 169 | sm_state = SessionState(self.redis) |
| 155 | rm_state = ResourceState(self.redis) | 170 | rm_state = ResourceState(self.redis) |
| 156 | 171 | ||
| 172 | + # 评估域:计数缓冲(热路径内存)+ 采集器 + 评估器(键前缀独立 hash tag) | ||
| 173 | + eval_state = EvaluationState(self.redis) | ||
| 174 | + telemetry = ScopeTelemetryBuffer() | ||
| 175 | + | ||
| 157 | self.sm_facade = SessionManagerFacade(sm_state) | 176 | self.sm_facade = SessionManagerFacade(sm_state) |
| 158 | 177 | ||
| 159 | - rm_orchestrator = ResourceOrchestrator(rm_state, self.k8s) | 178 | + rm_orchestrator = ResourceOrchestrator(rm_state, self.k8s, telemetry=telemetry) |
| 160 | self.rm_facade = ResourceManagerFacade(rm_orchestrator) | 179 | self.rm_facade = ResourceManagerFacade(rm_orchestrator) |
| 161 | 180 | ||
| 162 | self.sm_config_store = ConfigStore( | 181 | self.sm_config_store = ConfigStore( |
| @@ -170,11 +189,23 @@ class OrchestratorSystemContext(SystemContext): | |||
| 170 | self.sm_config_store, | 189 | self.sm_config_store, |
| 171 | self.rm_facade, | 190 | self.rm_facade, |
| 172 | default_session_ttl=self.arc.default_session_ttl, | 191 | default_session_ttl=self.arc.default_session_ttl, |
| 192 | + telemetry=telemetry, | ||
| 173 | ) | 193 | ) |
| 174 | self.sm_sweeper = SessionSweeper(sm_state, self.rm_facade) | 194 | self.sm_sweeper = SessionSweeper(sm_state, self.rm_facade) |
| 175 | self.rm_sweeper = ResourceSweeper( | 195 | self.rm_sweeper = ResourceSweeper( |
| 176 | rm_state, self.k8s, self.sm_facade, orchestrator=rm_orchestrator, | 196 | rm_state, self.k8s, self.sm_facade, orchestrator=rm_orchestrator, |
| 197 | + event_sink=eval_state.bump_event, | ||
| 177 | ) | 198 | ) |
| 199 | + self.eval_state = eval_state | ||
| 200 | + self.eval_telemetry = telemetry | ||
| 201 | + self.eval_collector = EvaluationCollector( | ||
| 202 | + eval_state=eval_state, sm_state=sm_state, rm_state=rm_state, | ||
| 203 | + ) | ||
| 204 | + self.evaluator = Evaluator( | ||
| 205 | + collector=self.eval_collector, llm=LLMClient.from_arc(self.arc), | ||
| 206 | + state=eval_state, arc=self.arc, instance_id=self.instance_id, | ||
| 207 | + ) | ||
| 208 | + self._telemetry_task: asyncio.Task | None = None | ||
| 178 | 209 | ||
| 179 | def _build_jobs(self) -> list[Any]: | 210 | def _build_jobs(self) -> list[Any]: |
| 180 | """全部后台任务(tick 级选主锁;多副本全局单副本执行写操作)。""" | 211 | """全部后台任务(tick 级选主锁;多副本全局单副本执行写操作)。""" |
| @@ -209,6 +240,19 @@ class OrchestratorSystemContext(SystemContext): | |||
| 209 | lock_key="agent_runtime:job:rm_reconcile", | 240 | lock_key="agent_runtime:job:rm_reconcile", |
| 210 | tick_timeout_sec=TICK_TIMEOUTS["rm_reconcile"], | 241 | tick_timeout_sec=TICK_TIMEOUTS["rm_reconcile"], |
| 211 | ), | 242 | ), |
| 243 | + # 系统自评估(sys_eval 全局单副本产报告,任意副本可读): | ||
| 244 | + self.create_single_leader_job( | ||
| 245 | + name="sys_sample", on_tick=self.eval_collector.sample_once, | ||
| 246 | + interval_sec=max(arc.eval_sample_interval, 5), | ||
| 247 | + lock_key="agent_runtime:job:sys_sample", | ||
| 248 | + tick_timeout_sec=TICK_TIMEOUTS["sys_sample"], | ||
| 249 | + ), | ||
| 250 | + self.create_single_leader_job( | ||
| 251 | + name="sys_eval", on_tick=self.evaluator.evaluate_once, | ||
| 252 | + interval_sec=max(arc.eval_interval, 30), | ||
| 253 | + lock_key="agent_runtime:job:sys_eval", | ||
| 254 | + tick_timeout_sec=TICK_TIMEOUTS["sys_eval"], | ||
| 255 | + ), | ||
| 212 | ] | 256 | ] |
| 213 | return jobs | 257 | return jobs |
| 214 | 258 | ||
| @@ -236,15 +280,24 @@ class OrchestratorSystemContext(SystemContext): | |||
| 236 | job.name, interval_by_name.get(job.name, "?"), | 280 | job.name, interval_by_name.get(job.name, "?"), |
| 237 | TICK_TIMEOUTS.get(job.name), | 281 | TICK_TIMEOUTS.get(job.name), |
| 238 | ) | 282 | ) |
| 283 | + # 计数缓冲 flusher:每副本独立(选主 job 会漏非 leader 副本的缓冲), | ||
| 284 | + # drain → 同槽 HINCRBY 批量聚合;失败留到下轮(计数只延后不丢) | ||
| 285 | + self._telemetry_task = asyncio.create_task( | ||
| 286 | + self._telemetry_flush_loop(), name="eval-telemetry-flush" | ||
| 287 | + ) | ||
| 239 | self.logger.info( | 288 | self.logger.info( |
| 240 | "config summary: mode=%s namespace=%s sweep=%ss autoscale=%ss " | 289 | "config summary: mode=%s namespace=%s sweep=%ss autoscale=%ss " |
| 241 | "reclaim=%ss watch=%ss reconcile=%ss " | 290 | "reclaim=%ss watch=%ss reconcile=%ss " |
| 242 | - "default_session_ttl=%ss kubeconfig=%s", | 291 | + "default_session_ttl=%ss eval_sample=%ss eval=%ss eval_llm=%s " |
| 292 | + "kubeconfig=%s", | ||
| 243 | self.arc.mode, self.arc.default_namespace, | 293 | self.arc.mode, self.arc.default_namespace, |
| 244 | self.arc.sweep_interval, self.arc.autoscale_interval, | 294 | self.arc.sweep_interval, self.arc.autoscale_interval, |
| 245 | self.arc.reclaim_interval, self.arc.watch_interval, | 295 | self.arc.reclaim_interval, self.arc.watch_interval, |
| 246 | self.arc.reconcile_interval, | 296 | self.arc.reconcile_interval, |
| 247 | self.arc.default_session_ttl, | 297 | self.arc.default_session_ttl, |
| 298 | + self.arc.eval_sample_interval, self.arc.eval_interval, | ||
| 299 | + "enabled" if (self.arc.eval_llm_base_url and self.arc.eval_llm_model) | ||
| 300 | + else "disabled", | ||
| 248 | "set" if self.arc.kubeconfig else "in-cluster", | 301 | "set" if self.arc.kubeconfig else "in-cluster", |
| 249 | ) | 302 | ) |
| 250 | self.logger.info( | 303 | self.logger.info( |
| @@ -261,10 +314,31 @@ class OrchestratorSystemContext(SystemContext): | |||
| 261 | "rm_reclaim": self.arc.reclaim_interval, | 314 | "rm_reclaim": self.arc.reclaim_interval, |
| 262 | "rm_watch": self.arc.watch_interval, | 315 | "rm_watch": self.arc.watch_interval, |
| 263 | "rm_reconcile": self.arc.reconcile_interval, | 316 | "rm_reconcile": self.arc.reconcile_interval, |
| 317 | + "sys_sample": self.arc.eval_sample_interval, | ||
| 318 | + "sys_eval": self.arc.eval_interval, | ||
| 264 | } | 319 | } |
| 265 | 320 | ||
| 321 | + async def _telemetry_flush_loop(self) -> None: | ||
| 322 | + """每副本一个:route/acquire 计数缓冲 → Redis 聚合(5s 批量)。""" | ||
| 323 | + while True: | ||
| 324 | + await asyncio.sleep(FLUSH_INTERVAL_SEC) | ||
| 325 | + try: | ||
| 326 | + await asyncio.wait_for( | ||
| 327 | + self._flush_telemetry_once(), timeout=FLUSH_TIMEOUT_SEC | ||
| 328 | + ) | ||
| 329 | + except asyncio.CancelledError: | ||
| 330 | + raise | ||
| 331 | + except Exception: # noqa: BLE001 - 留到下轮(计数只延后不丢) | ||
| 332 | + self.logger.exception("telemetry flush failed, retry next round") | ||
| 333 | + | ||
| 334 | + async def _flush_telemetry_once(self) -> None: | ||
| 335 | + drained = self.eval_telemetry.drain() | ||
| 336 | + for scope_id, deltas in drained.items(): | ||
| 337 | + if deltas: | ||
| 338 | + await self.eval_state.bump_counters(scope_id, deltas) | ||
| 339 | + | ||
| 266 | async def jobs_snapshot(self) -> list[dict[str, Any]]: | 340 | async def jobs_snapshot(self) -> list[dict[str, Any]]: |
| 267 | - """诊断用:5 个后台任务的间隔/超时/计数器/当前 leader(/visualization/overview)。""" | 341 | + """诊断用:7 个后台任务的间隔/超时/计数器/当前 leader(/visualization/overview)。""" |
| 268 | intervals = self._job_intervals() | 342 | intervals = self._job_intervals() |
| 269 | out: list[dict[str, Any]] = [] | 343 | out: list[dict[str, Any]] = [] |
| 270 | for job in self._jobs: | 344 | for job in self._jobs: |
| @@ -298,6 +372,22 @@ class OrchestratorSystemContext(SystemContext): | |||
| 298 | except Exception: # noqa: BLE001 | 372 | except Exception: # noqa: BLE001 |
| 299 | self.logger.exception("job stop failed: %s", job) | 373 | self.logger.exception("job stop failed: %s", job) |
| 300 | self._jobs = [] | 374 | self._jobs = [] |
| 375 | + if self._telemetry_task is not None: | ||
| 376 | + self._telemetry_task.cancel() | ||
| 377 | + try: | ||
| 378 | + await self._telemetry_task | ||
| 379 | + except asyncio.CancelledError: | ||
| 380 | + pass | ||
| 381 | + except Exception: # noqa: BLE001 | ||
| 382 | + self.logger.exception("telemetry flush task failed on cancel") | ||
| 383 | + self._telemetry_task = None | ||
| 384 | + # 终结 drain:停机前把缓冲余量尽量落 Redis(失败只留痕,计数延后无害) | ||
| 385 | + try: | ||
| 386 | + await asyncio.wait_for( | ||
| 387 | + self._flush_telemetry_once(), timeout=FLUSH_TIMEOUT_SEC | ||
| 388 | + ) | ||
| 389 | + except Exception: # noqa: BLE001 | ||
| 390 | + self.logger.exception("final telemetry drain failed") | ||
| 301 | try: | 391 | try: |
| 302 | await self.k8s.close() | 392 | await self.k8s.close() |
| 303 | except Exception: # noqa: BLE001 | 393 | except Exception: # noqa: BLE001 |