mateclaw/mateclaw-server/src/main/java/vip/mate/agent/AgentService.java
MIST e4dd08b5f4
feat(agent): 注意力锚定与环境感知——MCP 工具溯源 + skill 约束固定 + 事件通知 (#490)
* feat(agent): 注意力锚定与环境感知——MCP 工具溯源 + skill 约束固定 + 事件通知

## 背景

1. **MCP 工具跨服务器混淆**:MCP 工具名是 `mcp_<serverId>_<slug>_<hash6>`,serverId 是 19 位不可读 Snowflake。LLM 在多服务器任务中常把 slug 拼到错误 serverId 上重构出不存在的工具名,反复重试到 max iterations。
2. **长对话中 skill 约束丢失**:`load_skill` 返回的 SKILL.md 正文存在 messages 历史窗口里,被压缩管线(Soft Trim / Hard Clear / Pre-Prune / LLM Summary)销毁,约束彻底消失,agent 后续步骤违反约束。
3. **运行时环境变更对 agent 不可见**:MCP 服务器断连 / skill 更新发生在 agent 推理中途时,工具列表是 turn-start 快照,LLM 无法感知,继续调用已失效的工具。
4. **ledger 条目可被 LLM 反向覆盖**:Java 用 `auto_`/`pin_` 前缀让位给 LLM,但 LLM 没有反向保护——`progress_update(stepKey="auto_read_file")` 会覆盖 Java 写入的条目,保护是单向的。
5. **SkillManifestParser 从未填充 constraints 字段**:`KNOWN_KEYS` 未列入 `"constraints"`,导致约束被静默路由到 `extras`,所有依赖 `manifest.getConstraints()` 的代码都是死代码。

## 改动内容

### 文件改动

**新增文件(生产代码 4 个)**
- **`mateclaw-server/.../agent/runtime/EnvironmentNotification.java`** — 环境变更通知 record(type / message / timestamp)。
- **`mateclaw-server/.../agent/runtime/RunningConversationRegistry.java`** — 跟踪活跃会话 + 每会话有界通知队列(上限 10)+ TTL 定时清理(30 分钟未活跃的 handle 自动回收)。
- **`mateclaw-server/.../agent/runtime/EnvironmentEventRouter.java`** — 5 个 `@EventListener` 把 MCP/skill 事件翻译成中文 LLM 通知并广播。
- **`mateclaw-server/.../skill/event/SkillUpdatedEvent.java`** — skill 更新/启用/禁用/重扫描事件。

**新增文件(测试 5 个)**
- **`mateclaw-server/.../skill/manifest/SkillManifestConstraintsParsingTest.java`** — constraints 解析白盒测试(5 用例)。
- **`mateclaw-server/.../agent/progress/ProgressLedgerPrefixGuardTest.java`** — 前缀守卫 + 三类条目 + 并发 + 批量 auto-record 白盒(22 用例)。
- **`mateclaw-server/.../agent/runtime/RunningConversationRegistryTest.java`** — registry + router 生命周期 + TTL 清理白盒(24 用例)。
- **`mateclaw-server/.../agent/context/ContextCompressionLedgerSurvivalTest.java`** — 三类条目压缩存活黑盒(5 用例)。
- **`mateclaw-server/.../agent/graph/node/EnvironmentNotificationRenderingTest.java`** — 事件→通知→LLM 可见黑盒(14 用例)。

**修改文件(生产代码 14 个)**
- **`mateclaw-server/.../agent/progress/ProgressLedger.java`** — 增加 `pinned` map + `AUTO_RECORDED_PREFIX` 常量 + 三类条目区分;`mostRecentUpdate` 只看 regular 条目;`renderStaleReminder` 补 pending 计数。
- **`mateclaw-server/.../agent/progress/ProgressLedgerService.java`** — JSON 格式升级为 wrapper `{entries, pinned}`(向后兼容旧 flat-map);`upsert` 加 `auto_`/`pin_` 前缀守卫;`upsertPinned` / `upsertAutoRecorded` / `clearPinnedByPrefix` / `upsertAutoRecordedBatch`(批量版,一次 lock+load+save 处理 N 个工具响应);auto-recorded 4 参签名避免跨服务器键碰撞,有界=5。
- **`mateclaw-server/.../agent/graph/node/ActionNode.java`** — `load_skill` 后 `pinSkillConstraints` 把约束写入 pinned;工具调用后 `autoRecordToolCalls` 收集批量后一次 `upsertAutoRecordedBatch`(避免 N 次 lock+save 串行化);setter 注入保持测试构造器兼容。
- **`mateclaw-server/.../agent/graph/node/ReasoningNode.java`** — C4 注入:drain 通知 → `renderEnvironmentNotifications` → SystemMessage 加入 nonHistoryPrefix;helper 改 package-private 供黑盒测试。
- **`mateclaw-server/.../agent/AgentGraphBuilder.java`** — 系统提示增加 ProgressLedger Discipline 段(agent-3)+ Environment Change Notifications 段(agent-1);SkillCatalog 渲染器扫描约束加 🔒 锚点(agent-4);wire ActionNode setter + ReasoningNode registry。
- **`mateclaw-server/.../agent/AgentService.java`** — `withLifecycleSync` / `withLifecycleFlux` 入口 `safeRegister`、出口 `safeUnregister`,覆盖 Flux 抛错路径。
- **`mateclaw-server/.../skill/manifest/SkillManifest.java`** — 增加 `constraints` 字段(List<String>)。
- **`mateclaw-server/.../skill/manifest/SkillManifestParser.java`** — `KNOWN_KEYS` 加 `"constraints"`;builder 链加 `.constraints(stringList(fm.get("constraints")))`。
- **`mateclaw-server/.../skill/service/SkillService.java`** — 4 个改动点发布 `SkillUpdatedEvent`(rescan / update builtin / update non-builtin / toggle enable-disable)。
- **`mateclaw-server/.../tool/builtin/ProgressLedgerTool.java`** — `@Tool` 描述声明 `auto_`/`pin_` 前缀保留;`@ToolParam stepKey` 同步警告。
- **`mateclaw-server/.../tool/mcp/runtime/PrefixedNameToolCallback.java`** — 新增 3 参构造器,serverName 非空时描述前缀 `[MCP server: <name>]`,让 LLM 区分跨服务器同名工具。
- **`mateclaw-server/.../tool/mcp/runtime/McpClientManager.java`** — `wrapServerCallbacks` 透传 serverName 到 PrefixedNameToolCallback。
- **`mateclaw-server/.../agent/context/ConversationWindowManager.java`** — `PRUNE_EXEMPT_TOOLS` 加入 `load_skill`(A1)。
- **`mateclaw-server/.../agent/graph/executor/ToolExecutionExecutor.java`** — 工具不存在时 `buildMcpAwareNotFoundMessage` 跨服务器搜索同 slug/hash 候选,给出 ≤5 个建议名。

**修改文件(测试 1 个)**
- **`mateclaw-server/.../agent/progress/ProgressLedgerStaleReminderTest.java`** — 回归适配:reminder 文本现在包含 `pending` 计数。

### 测试

- `mvn -pl mateclaw-server -am test -Dtest='SkillManifestConstraintsParsingTest,ProgressLedgerPrefixGuardTest,RunningConversationRegistryTest,ContextCompressionLedgerSurvivalTest,EnvironmentNotificationRenderingTest,ProgressLedgerStaleReminderTest'`:70/70 通过
- `mvn -pl mateclaw-server -am test`(全量回归,含上面 6 个 + 12 个深挖影响类):0 失败 0 错误

### 安全性

- **前缀保留**:`ProgressLedgerService.upsert` 拒绝 `auto_`/`pin_` 前缀,LLM 无法覆盖 Java 管理的条目;`@Tool` 描述显式声明保留前缀。
- **事件路由异常隔离**:`EnvironmentEventRouter.broadcast` 全 try/catch,路由失败永不冒泡到 Spring 事件总线。
- **并发安全**:registry 用 `ConcurrentHashMap` + `ConcurrentLinkedQueue`;ledger upsert 用 per-conversation `ReentrantLock`;批量 auto-record 在单次 lock 内完成;`ProgressLedgerPrefixGuardTest.concurrentUpsertAndAutoRecordAreSafe` 锁定。
- **内存有界**:通知队列每会话上限 10(LRU 驱逐最老);auto-recorded 条目每会话上限 5(驱逐最老);registry 后台 TTL 清理(30 分钟未活跃的 handle 自动回收)。
- **绑定机制不受影响**:MCP/skill 的 agent 绑定(`mate_agent_tool` / `mate_agent_skill` 表)完全未被触碰;C3 通知广播是有意全量(非按 agentId 过滤),最坏情况是无关 agent 多收一条 SystemMessage(LLM 被告知"如无关可忽略")。

## 逐项验证

### 改动 1:SkillManifestParser 真正解析 constraints(深挖修复)

**文件**:`mateclaw-server/src/main/java/vip/mate/skill/manifest/SkillManifestParser.java:33-47,108`

| 项 | 内容 |
|---|---|
| 改了什么 | `KNOWN_KEYS` 集合加入 `"constraints"`;builder 链加入 `.constraints(stringList(fm.get("constraints")))`。 |
| 为什么 | 之前 `KNOWN_KEYS` 没列入 `"constraints"`,导致该键被静默路由到 `extras`,`manifest.getConstraints()` 永远返回空 list,下游 B2 pinSkillConstraints 和 agent-4 catalog 锚点全是死代码。 |
| 验证步骤 | 1. `cat test-fixtures/skill-with-constraints/SKILL.md`(如有)确认 frontmatter 有 `constraints: [...]`;2. 运行 `SkillManifestConstraintsParsingTest`。 |
| 预期结果 | `manifest.getConstraints()` 返回非空 list;test 5/5 通过。 |

### 改动 2:ProgressLedgerService 前缀守卫(深挖修复)

**文件**:`mateclaw-server/src/main/java/vip/mate/agent/progress/ProgressLedgerService.java:132-137`

| 项 | 内容 |
|---|---|
| 改了什么 | `upsert()` 入口检查 key 是否以 `auto_` 或 `pin_` 开头,是则抛 `IllegalArgumentException`。 |
| 为什么 | 之前保护是单向的:Java 让位 LLM(auto 不覆盖 LLM 已有 entry),但 LLM 可以用 `progress_update(stepKey="auto_read_file")` 覆盖 Java 写入的条目,导致 auto-recorded 工具记录被改写。 |
| 验证步骤 | 1. `ProgressLedgerPrefixGuardTest.upsertRejectsAutoPrefix`;2. `ProgressLedgerPrefixGuardTest.upsertRejectsPinPrefix`。 |
| 预期结果 | 两个测试均抛 `IllegalArgumentException`;`ProgressLedgerTool` 的 `@Tool` 描述包含前缀保留声明。 |

### 改动 3:upsertAutoRecorded 4 参签名 + 批量化(深挖修复 + 性能优化)

**文件**:`mateclaw-server/src/main/java/vip/mate/agent/progress/ProgressLedgerService.java:227-296`、`mateclaw-server/src/main/java/vip/mate/agent/graph/node/ActionNode.java`(`autoRecordToolCalls`)

| 项 | 内容 |
|---|---|
| 改了什么 | `upsertAutoRecorded` 升级为 4 参签名 `(conversationId, toolName, displayName, resultSummary)`;新增 `upsertAutoRecordedBatch` 批量方法,一次 lock+load+save 处理 N 个工具响应;ActionNode 改为先收集 `List<AutoRecordEntry>` 再一次批量调用。 |
| 为什么 | 1. 跨服务器键碰撞:两个 MCP 服务器都暴露 `search` 工具 → `auto_search` 互相覆盖;2. 并行工具串行化:每个 ToolResponse 单独 lock+load+save 抵消并行收益。 |
| 验证步骤 | 1. `ProgressLedgerPrefixGuardTest.autoRecordedDifferentServersNoCollision`:两个服务器同名工具共存;2. `ProgressLedgerPrefixGuardTest.batchInsertProducesSameResultAsSequential`:批量与逐条结果一致;3. `ProgressLedgerPrefixGuardTest.batchInsertBoundedToMaxFiveEvenWithLargeBatch`:10 条批量插入后有界=5。 |
| 预期结果 | ledger 中同时存在 `auto_mcp_4_search_xxx` 和 `auto_mcp_7_search_yyy`;5 个并行工具调用从 5 次 lock+save 降为 1 次。 |

### 改动 4:B2 pinSkillConstraints——load_skill 后约束写入 pinned

**文件**:`mateclaw-server/src/main/java/vip/mate/agent/graph/node/ActionNode.java`(`pinSkillConstraints`)

| 项 | 内容 |
|---|---|
| 改了什么 | `load_skill` 工具调用成功后,读取 `manifest.getConstraints()`,对每条约束调用 `progressLedgerService.upsertPinned(convId, "pin_<skillName>_<i>", constraintText, note)`。 |
| 为什么 | 把约束从 messages(会被压缩销毁)抽到 DB ledger.pinned(压缩免疫),解决"长对话中 skill 约束丢失"问题。 |
| 验证步骤 | 1. `ContextCompressionLedgerSurvivalTest.loadSkillBodyDestroyedByCompressionButConstraintsSurviveInLedger`;2. `ContextCompressionLedgerSurvivalTest.allThreeLedgerEntryTypesSurviveCompression`。 |
| 预期结果 | 压缩后 messages 中 load_skill 正文消失,但 `ledger.renderSnapshot()` 仍包含 `🔒 固定约束` 段,约束文本字节级保留。 |

### 改动 5:B5 autoRecordToolCalls——工具调用后批量自动记录

**文件**:`mateclaw-server/src/main/java/vip/mate/agent/graph/node/ActionNode.java`(`autoRecordToolCalls`)

| 项 | 内容 |
|---|---|
| 改了什么 | ActionNode 处理 ToolResponse 后,收集所有有效条目到 `List<AutoRecordEntry>`,一次调用 `upsertAutoRecordedBatch`。 |
| 为什么 | 让 LLM 在长对话中即使忘记自己刚调用过什么工具,也能从 ledger 快照看到最近 5 次工具调用记录;批量调用避免 N 次 lock+save 串行化。 |
| 验证步骤 | `ProgressLedgerPrefixGuardTest.autoRecordedBoundedToMaxFive`:模拟 10 次工具调用,验证 auto 条目数等于 5。 |
| 预期结果 | ledger 中 auto 条目始终 ≤ 5,最老的被驱逐;5 个并行工具调用只需 1 次 DB roundtrip。 |

### 改动 6:C4 环境通知注入 nonHistoryPrefix

**文件**:`mateclaw-server/src/main/java/vip/mate/agent/graph/node/ReasoningNode.java:755-765,1201-1215`

| 项 | 内容 |
|---|---|
| 改了什么 | ReasoningNode 每轮推理前 `registry.drain(conversationId)`,非空则 `renderEnvironmentNotifications` 渲染成 markdown 块,作为 SystemMessage 加入 nonHistoryPrefix。 |
| 为什么 | 让运行时环境变更(MCP 断连 / skill 更新)在下一轮推理立即可见,LLM 主动改路而不是反复重试失效工具。 |
| 验证步骤 | `EnvironmentNotificationRenderingTest.mcpConnectionLostEventEndToEnd_producesActionableLLMText`:注册会话 → 触发 `McpConnectionLostEvent(serverId=7)` → drain → render。 |
| 预期结果 | 渲染块包含 "📢 环境变更通知"、`mcp_7_` 前缀、"不要反复重试" 指令。 |

### 改动 7:A1 PRUNE_EXEMPT_TOOLS 加入 load_skill

**文件**:`mateclaw-server/src/main/java/vip/mate/agent/context/ConversationWindowManager.java:96-120`

| 项 | 内容 |
|---|---|
| 改了什么 | `PRUNE_EXEMPT_TOOLS` 集合从 `{delegateToAgent, delegateParallel}` 扩展为 `{delegateToAgent, delegateParallel, load_skill}`。 |
| 为什么 | `load_skill` 返回的 SKILL.md 是 load-time 快照,skill 作者可能在执行期间更新,重载不保证恢复相同指令;且 50KB+ skill 重载昂贵。 |
| 验证步骤 | `ConversationWindowManagerToolPruningTest`(已有测试套件)。 |
| 预期结果 | load_skill 的 ToolResponseMessage 在 `pruneOldToolResultsForModelInput` / `compactAgedToolResponses` 阶段不被修剪。 |

### 改动 8:agent-2 PrefixedNameToolCallback 描述加 serverName 标签

**文件**:`mateclaw-server/src/main/java/vip/mate/tool/mcp/runtime/PrefixedNameToolCallback.java:55-80`、`mateclaw-server/src/main/java/vip/mate/tool/mcp/runtime/McpClientManager.java:201-300`

| 项 | 内容 |
|---|---|
| 改了什么 | 新增 3 参构造器 `(prefixedName, delegate, serverName)`,serverName 非空时描述前缀 `[MCP server: <name>]`;McpClientManager `wrapServerCallbacks` 透传 serverName。 |
| 为什么 | LLM 看到 `mcp_4_search_a1b2c3` 时无法知道这是哪个服务器的工具;多个 MCP 服务器都暴露 `search` 时,LLM 会混淆。加 `[MCP server: fetch-server]` 标签让 LLM 区分。 |
| 验证步骤 | 启动一个 MCP 服务器,在 agent 工具列表中观察工具描述是否包含 `[MCP server: <name>]` 前缀。 |
| 预期结果 | 每个 MCP 工具描述开头包含 `[MCP server: <serverName>]`;2 参构造器(无 serverName)保持向后兼容,描述不加前缀。 |

### 改动 9:ToolExecutionExecutor 工具不存在时跨服务器候选建议

**文件**:`mateclaw-server/src/main/java/vip/mate/agent/graph/executor/ToolExecutionExecutor.java:1264-1340`

| 项 | 内容 |
|---|---|
| 改了什么 | "Tool not found" 错误信息升级:若请求名是 MCP 格式,搜索 `toolCallbackMap` 中 slug 或 hash6 匹配但 serverId 不同的候选,返回 ≤5 个建议。 |
| 为什么 | LLM 常把 slug 拼到错误 serverId 上重构出不存在工具名,反复重试到 max iterations。给候选建议后 LLM 可以直接复制正确名字。 |
| 验证步骤 | 1. 启动两个 MCP 服务器都暴露 `fetch` 工具;2. 让 LLM 调用 `mcp_<serverA>_fetch_xxx`(实际 fetch 在 serverB);3. 观察错误信息。 |
| 预期结果 | 错误信息包含 "Did you mean one of these?" + 正确的 `mcp_<serverB>_fetch_yyy` 候选名。 |

### 改动 10:Registry TTL 定时清理(防泄漏)

**文件**:`mateclaw-server/src/main/java/vip/mate/agent/runtime/RunningConversationRegistry.java:155-211`

| 项 | 内容 |
|---|---|
| 改了什么 | 新增 `cleanupStale(Duration maxAge)` 方法 + `@Scheduled scheduledCleanup()`(每 5 分钟扫一次,清理 30 分钟未活跃的 handle)。用 `remove(key, value)` 保证不误删被并发 `register` 刷新的 handle。 |
| 为什么 | 兜底防御异常路径泄漏的 handle——即使 `safeUnregister` 因异常路径未执行(如 Reactor cancel 信号不触发 doFinally),后台线程也能回收。 |
| 验证步骤 | `RunningConversationRegistryTest.cleanupStaleRemovesOldHandles`:注册 → 反射 backdate lastActiveAt → 清理 → 验证被移除;`cleanupStaleDoesNotRemoveRefreshedHandle`:backdate 后 re-register → 清理 → 验证存活。 |
| 预期结果 | 30 分钟未活跃的 handle 被清理;被 `register` 刷新的 handle 不被误删。 |

### 改动 11:JSON 格式向后兼容迁移

**文件**:`mateclaw-server/src/main/java/vip/mate/agent/progress/ProgressLedgerService.java:79-85,284-292`

| 项 | 内容 |
|---|---|
| 改了什么 | JSON 从 flat-map `{"step1":{...}}` 升级为 wrapper `{"entries":{...},"pinned":{...}}`;`parseWrapper` 通过 peek `"entries"` 键区分新旧格式,旧格式自动迁移为 wrapper(pinned 为空)。 |
| 为什么 | 老 conversation 的 ledger 列存的是 flat-map,新代码上线后必须能加载老数据。 |
| 验证步骤 | `ContextCompressionLedgerSurvivalTest.oldFlatMapLedgerMigratesToWrapperFormatWithEmptyPinned`:写入旧 JSON → load → 验证 pinned 为空 → upsert → 验证新 JSON 包含 `entries` 和 `pinned` 键。 |
| 预期结果 | 旧 conversation 无需迁移脚本,第一次 load 即兼容;写入时自动转为新格式。 |

## 新增测试验证

**文件**:
- `mateclaw-server/src/test/java/vip/mate/skill/manifest/SkillManifestConstraintsParsingTest.java`
- `mateclaw-server/src/test/java/vip/mate/agent/progress/ProgressLedgerPrefixGuardTest.java`
- `mateclaw-server/src/test/java/vip/mate/agent/runtime/RunningConversationRegistryTest.java`
- `mateclaw-server/src/test/java/vip/mate/agent/context/ContextCompressionLedgerSurvivalTest.java`
- `mateclaw-server/src/test/java/vip/mate/agent/graph/node/EnvironmentNotificationRenderingTest.java`

| 命令 | 预期 |
|------|------|
| `mvn -pl mateclaw-server -am test -Dtest='SkillManifestConstraintsParsingTest'` | Tests run: 5, Failures: 0 |
| `mvn -pl mateclaw-server -am test -Dtest='ProgressLedgerPrefixGuardTest'` | Tests run: 22, Failures: 0 |
| `mvn -pl mateclaw-server -am test -Dtest='RunningConversationRegistryTest'` | Tests run: 24, Failures: 0 |
| `mvn -pl mateclaw-server -am test -Dtest='ContextCompressionLedgerSurvivalTest'` | Tests run: 5, Failures: 0 |
| `mvn -pl mateclaw-server -am test -Dtest='EnvironmentNotificationRenderingTest'` | Tests run: 14, Failures: 0 |

## 回归检查清单

- [ ] 全量 `mvn -pl mateclaw-server -am test` 通过(已验证 0 失败 0 错误)
- [ ] 老 conversation(flat-map ledger JSON)首次 load 不报错,pinned 字段为空
- [ ] 多 MCP 服务器场景:LLM 工具列表中每个工具描述包含 `[MCP server: <name>]` 标签
- [ ] MCP 服务器中途断连:agent 下一轮推理看到 `📢 环境变更通知` 块,主动改路
- [ ] 长 conversation(>100 轮)经多次 PTL 压缩后,`ledger.renderSnapshot()` 仍包含 pinned 约束
- [ ] `progress_update(stepKey="auto_xxx")` 被拒绝,返回 `IllegalArgumentException` 错误信息
- [ ] auto-recorded 条目数始终 ≤ 5(10 次工具调用后仍为 5)
- [ ] 5 个并行工具调用只产生 1 次 DB roundtrip(批量 auto-record)
- [ ] Registry 中 30 分钟未活跃的 handle 被后台定时清理
- [ ] MCP/skill 绑定机制(`mate_agent_tool` / `mate_agent_skill` 表)不受影响
- [ ] Plan-Execute 路径(StepExecutionNode)目前**不**接收环境通知——只有 ReAct 路径生效(已知未覆盖项,不阻塞本 PR)

* feat(agent): 六招减法重构——修复压缩销毁 skill 约束与 MCP 按需暴露

## 背景

- 压缩三阶段(softTrim / hardClear / prePruneForSummary)只检查 `isSpillMarker`,不检查 `PRUNE_EXEMPT_TOOLS`,导致 `load_skill` 返回的 SKILL.md 约束、`delegateToAgent` 子智能体转录在压缩中被裁掉,模型在长对话中"忘记"任务规则,根因是"压缩导致注意力失效"。
- skillCatalog 表只列 Skill / Status / Description 三列,bound skill 的 constraints 与 allowedTools 没有任何可见入口,模型加载 skill 后约束仍可能被忽略。
- MCP 工具默认 CORE tier,20+ MCP 工具的 schema 涌入核心列表,挤占 builtin 工具的注意力,且 `DisclosureTier.fromToken(null)` 返回 CORE 导致 `getOrDefault` 的默认值永远不生效。
- 构建期工具过滤分 4 次 pass,重复遍历且无明确 deny/allow 边界。
- skillCatalog 在每次推理步都按当前 loadedSkills 动态渲染,破坏 Anthropic SYSTEM_AND_TOOLS cache 前缀稳定性。
- 进度账本(ProgressLedger)约束条目前缀无保护,跨 MCP server 键碰撞,环境事件无统一路由入口。

## 改动内容

### 文件改动

**主代码(21 个文件)**

- **`mateclaw-server/src/main/java/vip/mate/agent/context/ConversationWindowManager.java`** — Move 4:三阶段新增 `isExemptTool` 检查跳过 `load_skill`/`delegateToAgent`/`delegateParallel`;新增 Phase 2.7 无损 spill evict 在调用 LLM 摘要前把超大工具结果落盘
- **`mateclaw-server/src/main/java/vip/mate/skill/runtime/SkillRuntimeService.java`** — Move 2 & 3:catalog 表新增 Constraints 列(仅 bound skill 显示);新增 `### Bound skill allowed tools` 块
- **`mateclaw-server/src/main/java/vip/mate/agent/graph/node/ReasoningNode.java`** — Move 1:skillCatalog 用 `render(Set.of())` 静态渲染进 nonHistoryPrefix;loadedThisRun hint 作为 volatile 后缀注入
- **`mateclaw-server/src/main/java/vip/mate/tool/disclosure/DefaultToolDisclosureService.java`** — Move 5:MCP 工具默认 tier 从 CORE 改 EXTENSION;`buildSnapshot` 跳过 null/blank tier 修复 `fromToken(null)→CORE` 陷阱
- **`mateclaw-server/src/main/java/vip/mate/agent/AgentGraphBuilder.java`** — Move 6:构建期权限过滤从 4 次 pass 合并为 2 次(deny 集 + allow 集)
- **`mateclaw-server/src/main/java/vip/mate/agent/AgentService.java`** — 接入 EnvironmentEventRouter 与 RunningConversationRegistry
- **`mateclaw-server/src/main/java/vip/mate/agent/graph/executor/ToolExecutionExecutor.java`** — 工具调用后自动回填 ProgressLedger
- **`mateclaw-server/src/main/java/vip/mate/agent/graph/node/ActionNode.java`** — 渲染 ledger 三段式快照
- **`mateclaw-server/src/main/java/vip/mate/agent/progress/ProgressLedger.java`** — constraints 前缀保护,跨 MCP server 键命名空间隔离
- **`mateclaw-server/src/main/java/vip/mate/agent/progress/ProgressLedgerService.java`** — 写入 constraints 到 pinned 条目
- **`mateclaw-server/src/main/java/vip/mate/skill/manifest/SkillManifest.java`** — 新增 constraints 字段
- **`mateclaw-server/src/main/java/vip/mate/skill/manifest/SkillManifestParser.java`** — 解析 SKILL.md frontmatter 中的 constraints
- **`mateclaw-server/src/main/java/vip/mate/skill/service/SkillService.java`** — skill 更新事件发布
- **`mateclaw-server/src/main/java/vip/mate/tool/builtin/ProgressLedgerTool.java`** — 三段式渲染
- **`mateclaw-server/src/main/java/vip/mate/tool/mcp/runtime/McpClientManager.java`** — MCP 命名透明化
- **`mateclaw-server/src/main/java/vip/mate/tool/mcp/runtime/PrefixedNameToolCallback.java`** — 透明命名映射
- **`mateclaw-server/src/main/java/vip/mate/agent/runtime/EnvironmentEventRouter.java`** — 新增:5 个环境事件监听器
- **`mateclaw-server/src/main/java/vip/mate/agent/runtime/EnvironmentNotification.java`** — 新增:环境通知数据模型
- **`mateclaw-server/src/main/java/vip/mate/agent/runtime/RunningConversationRegistry.java`** — 新增:运行中会话注册表 + TTL 清理
- **`mateclaw-server/src/main/java/vip/mate/skill/event/SkillUpdatedEvent.java`** — 新增:skill 更新事件
- **`mateclaw-server/Dockerfile`** — 构建配置微调

**测试代码(11 个文件)**

- **`ConversationWindowManagerExemptAndSpillTest.java`** — 新增 15 个行为测试,证明 Move 4 生效
- **`SkillRuntimeServiceConstraintsAndToolsTest.java`** — 新增 8 个测试覆盖 Constraints 列与 allowedTools 块
- **`ReasoningNodeLoadedSkillsHintTest.java`** — 新增 8 个测试覆盖 loadedThisRun hint 渲染
- **`CompactionSurvivalComparisonTest.java`** — 新增 4 个场景的新旧代码对比测试(100 轮极限压缩)
- **`ContextCompressionLedgerSurvivalTest.java`** — 压缩后 ledger 存活测试
- **`EnvironmentNotificationRenderingTest.java`** — 环境通知渲染测试
- **`ProgressLedgerPrefixGuardTest.java`** — ledger 前缀保护测试
- **`RunningConversationRegistryTest.java`** — 会话注册表测试
- **`SkillManifestConstraintsParsingTest.java`** — constraints 解析测试
- **`ProgressLedgerStaleReminderTest.java`** — 修复回归
- **`ToolDisclosureServiceTest.java`** — 断言从 CORE 改为 EXTENSION

### 测试

**回归测试**

- `mvn test`(mateclaw-server 全量):**199 通过 / 1 跳过 / 0 失败**

**行为测试(证明改动生效,旧代码上失败)**

- `CompactionSurvivalComparisonTest`(4 个场景):在新代码上全部通过
- 用 `git stash` 还原旧代码后,16 个行为测试编译失败或断言失败 → 证明测试确实覆盖了新行为

**新旧代码对比测试(同一份测试源码,两套代码库运行)**

| 场景 | 旧代码 | 新代码 |
|---|---|---|
| A: 50 load_skill + 50 delegate + 50 read_file 单轮压缩 | load_skill 0/50, delegate 0/50 | **load_skill 50/50, delegate 50/50** |
| B: 100 轮极限压缩 + 头部 pinned load_skill | root_constraint_survived=**false**, 113ms | root_constraint_survived=**true**, 60ms |
| C: 20 个不同大小 load_skill 单轮压缩 | 0/20 存活, tokens 8694→754 | **20/20 存活**, tokens 8694→8694 |
| D: 30 轮稳态压缩 + pinned skill | pinned_survived=**false** | pinned_survived=**true** |

### 安全性

- `DisclosureTier.fromToken(null)` 陷阱修复:旧代码 `serverTierById.put(id, CORE)` 导致 `getOrDefault` 默认值永不生效;新代码跳过 null tier,未配置的 MCP server 才走 EXTENSION 默认值
- `PRUNE_EXEMPT_TOOLS` 保护范围从 2 处扩展到 5 处,避免 `load_skill` 约束被压缩销毁后模型在无约束下执行敏感操作

## 逐项验证

### 改动 1:nonHistoryPrefix 分层稳定化

**文件**:`mateclaw-server/src/main/java/vip/mate/agent/graph/node/ReasoningNode.java:701-714, 776-788, 1246-1257`

| 项 | 内容 |
|---|---|
| 改了什么 | skillCatalog 用 `render(Set.of())` 静态渲染进 nonHistoryPrefix;loadedThisRun hint 作为 volatile 后缀注入 |
| 为什么 | 每次推理步都按 loadedSkills 动态渲染会破坏 Anthropic SYSTEM_AND_TOOLS cache 前缀,导致 cache 失效增加 token 成本 |
| 验证步骤 | 1. 打开 ReasoningNode.java:701;2. 确认 `skillCatalogRenderer.render(java.util.Set.of())` 调用;3. 跑 `ReasoningNodeLoadedSkillsHintTest` |
| 预期结果 | 8 个测试通过,loadedThisRun hint 作为后缀注入,不破坏前缀缓存 |

### 改动 2:skillCatalog 增加 Constraints 列

**文件**:`mateclaw-server/src/main/java/vip/mate/skill/runtime/SkillRuntimeService.java:502-561, 634-641`

| 项 | 内容 |
|---|---|
| 改了什么 | catalog 表从 3 列扩为 4 列,新增 Constraints 列(仅 bound skill 显示,长约束截断,pipe 转义);新增 `### Bound skill allowed tools` 块 |
| 为什么 | bound skill 的 constraints 没有任何可见入口,模型加载后仍可能忽略 |
| 验证步骤 | 1. 跑 `SkillRuntimeServiceConstraintsAndToolsTest`;2. 检查 catalog 渲染包含 Constraints 列 |
| 预期结果 | 8 个测试通过,bound skill 显示 constraints,非 bound skill 省略 |
- **边界验证**:长约束截断为单行;pipe 字符被转义不破坏表格

### 改动 3:修复 PRUNE_EXEMPT_TOOLS 在三阶段的绕过

**文件**:`mateclaw-server/src/main/java/vip/mate/agent/context/ConversationWindowManager.java:990-994, 1030-1090, 1160-1210`

| 项 | 内容 |
|---|---|
| 改了什么 | `softTrimToolResults`、`hardClearToolResults`、`prePruneForSummary` 三处新增 `isExemptTool` 检查,跳过 `load_skill`/`delegateToAgent`/`delegateParallel` |
| 为什么 | 旧代码只在 `pruneOldToolResultsForModelInput` 和 `compactAgedToolResponses` 检查 exempt,三阶段不检查,导致 skill 约束在压缩中被裁掉 |
| 验证步骤 | 1. 跑 `ConversationWindowManagerExemptAndSpillTest`;2. 跑 `CompactionSurvivalComparisonTest` |
| 预期结果 | 15 个行为测试通过;100 轮压缩后 load_skill body 存活率 100% |
- **反例对照**:在新代码上跑对比测试,旧代码存活率 0%,新代码 100%

### 改动 4:Phase 2.7 无损 spill evict

**文件**:`mateclaw-server/src/main/java/vip/mate/agent/context/ConversationWindowManager.java:440-468, 1263-1300`

| 项 | 内容 |
|---|---|
| 改了什么 | 在 Phase 2 hardClear 之后、LLM 摘要之前新增 Phase 2.7,把超大工具结果落盘替换为 spill marker |
| 为什么 | 旧代码超出预算直接走 LLM 摘要(有损+耗时+费 token),其实大部分场景落盘就够 |
| 验证步骤 | 1. 检查 `spillEvictToolResults` 方法;2. 跑 `ConversationWindowManagerExemptAndSpillTest.spillEvictReducesTokenCount` |
| 预期结果 | spill 后 token 数低于预算时跳过 LLM 摘要,strategy=lossless_spill_evict |

### 改动 5:MCP 工具默认 EXTENSION

**文件**:`mateclaw-server/src/main/java/vip/mate/tool/disclosure/DefaultToolDisclosureService.java:78-105, 276-336`

| 项 | 内容 |
|---|---|
| 改了什么 | `resolveTierByName` 默认返回 EXTENSION;`buildSnapshot` 跳过 null/blank tier 的 server 不放入 map |
| 为什么 | MCP schema 是 prompt 最重部分,默认 CORE 挤占 builtin 工具注意力;`fromToken(null)` 返回 CORE 导致默认值失效 |
| 验证步骤 | 1. 跑 `ToolDisclosureServiceTest.mcpDefaultsExtensionWhenServerTierUnset`;2. 检查未配置 tier 的 MCP server 工具不在 active 列表 |
| 预期结果 | 未配置 tier 的 MCP 工具进入 extensionCatalog,需 `enable_tool` 激活 |
- **边界验证**:显式 `disclosure_tier=core` 的 server 仍保持 CORE

### 改动 6:构建期权限过滤合并

**文件**:`mateclaw-server/src/main/java/vip/mate/agent/AgentGraphBuilder.java:258-310`

| 项 | 内容 |
|---|---|
| 改了什么 | 4 次 pass 合并为 2 次:先 `withDeniedToolsFiltered(deniedSet)`,再 `withAllowedToolsOnly(boundTools)` |
| 为什么 | 重复遍历浪费构建时间,且 deny/allow 边界不清晰 |
| 验证步骤 | 1. 检查 AgentGraphBuilder.java:258-310;2. 跑全量回归测试确认工具过滤行为不变 |
| 预期结果 | 工具列表与改动前一致,构建步骤减少 |

## 新增测试验证

**文件**:

- `mateclaw-server/src/test/java/vip/mate/agent/context/ConversationWindowManagerExemptAndSpillTest.java`
- `mateclaw-server/src/test/java/vip/mate/agent/context/CompactionSurvivalComparisonTest.java`
- `mateclaw-server/src/test/java/vip/mate/skill/runtime/SkillRuntimeServiceConstraintsAndToolsTest.java`
- `mateclaw-server/src/test/java/vip/mate/agent/graph/node/ReasoningNodeLoadedSkillsHintTest.java`
- `mateclaw-server/src/test/java/vip/mate/skill/manifest/SkillManifestConstraintsParsingTest.java`
- `mateclaw-server/src/test/java/vip/mate/agent/progress/ProgressLedgerPrefixGuardTest.java`
- `mateclaw-server/src/test/java/vip/mate/agent/runtime/RunningConversationRegistryTest.java`
- `mateclaw-server/src/test/java/vip/mate/agent/graph/node/EnvironmentNotificationRenderingTest.java`
- `mateclaw-server/src/test/java/vip/mate/agent/context/ContextCompressionLedgerSurvivalTest.java`

| 命令 | 预期 |
|------|------|
| `mvn -Dtest=ConversationWindowManagerExemptAndSpillTest test` | 15 个测试通过 |
| `mvn -Dtest=CompactionSurvivalComparisonTest test` | 4 个场景通过,新代码 load_skill 存活率 100% |
| `mvn -Dtest=SkillRuntimeServiceConstraintsAndToolsTest test` | 8 个测试通过 |
| `mvn -Dtest=ReasoningNodeLoadedSkillsHintTest test` | 8 个测试通过 |
| `mvn test`(全量) | 199 通过 / 1 跳过 / 0 失败 |

**新旧对比测试运行命令**:

```bash
# 新代码
cd /data/mateclaw/mateclaw-server && mvn -Dtest=CompactionSurvivalComparisonTest -Dsurefire.useFile=false test

# 旧代码(需把测试复制到 mateclaw-old)
cd /data/mateclaw/mateclaw-old/mateclaw-server && mvn -Dtest=CompactionSurvivalComparisonTest -Dsurefire.useFile=false test
```

## 回归检查清单

- [ ] `mvn test` 全量通过(199/1skip/0fail)
- [ ] 对比测试在新代码上 load_skill 存活率 100%
- [ ] 对比测试在旧代码上 load_skill 存活率 0%(证明测试有效)
- [ ] MCP 工具默认进入 extensionCatalog,`enable_tool` 可激活
- [ ] 显式 `disclosure_tier=core` 的 MCP server 仍保持 CORE
- [ ] 100 轮压缩后 root_constraint 仍存活
- [ ] Phase 2.7 spill evict 在预算内时跳过 LLM 摘要
- [ ] skillCatalog 静态渲染不依赖 loadedSkills,前缀缓存稳定
2026-07-06 11:50:41 +08:00

801 lines
37 KiB
Java
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

package vip.mate.agent;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.ApplicationEventPublisher;
import org.springframework.context.event.EventListener;
import org.springframework.stereotype.Service;
import org.springframework.util.StringUtils;
import reactor.core.publisher.Flux;
import vip.mate.agent.context.ChatOrigin;
import vip.mate.agent.context.ChatOriginHolder;
import vip.mate.agent.event.AgentLifecycleEvent;
import vip.mate.agent.model.AgentEntity;
import vip.mate.agent.repository.AgentMapper;
import vip.mate.exception.MateClawException;
import vip.mate.llm.chatmodel.ThinkingLevelHolder;
import vip.mate.llm.event.ModelConfigChangedEvent;
import vip.mate.memory.MemoryProperties;
import vip.mate.memory.lifecycle.MemoryLifecycleMediator;
import vip.mate.memory.lifecycle.TurnContext;
import vip.mate.memory.service.MemoryRecallTracker;
import vip.mate.workspace.conversation.model.ConversationEntity;
import vip.mate.workspace.conversation.repository.ConversationMapper;
import java.util.List;
import java.time.Duration;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.function.Function;
import java.util.function.Supplier;
/**
* Agent 业务服务
* <p>
* 负责 Agent 的 CRUD 管理和运行时实例管理。
* 构建逻辑委托给 {@link AgentGraphBuilder}。
*
* @author MateClaw Team
*/
@Slf4j
@Service
@RequiredArgsConstructor
public class AgentService {
private final AgentMapper agentMapper;
private final AgentGraphBuilder agentGraphBuilder;
private final MemoryRecallTracker memoryRecallTracker;
private final MemoryLifecycleMediator lifecycleMediator;
private final MemoryProperties memoryProperties;
private final vip.mate.memory.identity.MemoryOwnerResolver memoryOwnerResolver;
/** Read-only lookup of a conversation's pinned model. Mapper (not service)
* to keep this a leaf dependency with no risk of a bean cycle. */
private final ConversationMapper conversationMapper;
/** Field-injected publisher for agent_lifecycle trigger events; the
* trigger module's bridge listens and forwards into ingest. */
@Autowired(required = false)
private ApplicationEventPublisher events;
/**
* C5: tracks in-flight conversations so {@link vip.mate.agent.runtime.EnvironmentEventRouter}
* can push environment-change notifications into the agent's next reasoning
* turn. Field-injected (optional) so existing test constructors of
* {@code AgentService} don't need to supply it.
*/
@Autowired(required = false)
private vip.mate.agent.runtime.RunningConversationRegistry runningConversationRegistry;
/**
* Runtime Agent instance cache. Keyed first by agentId, then by a model
* key, so a conversation that pins a non-default model gets its own graph
* variant instead of mutating the one every other conversation shares.
* The model key is {@code ""} for the Agent / global-default model.
*/
private final Map<Long, Map<String, BaseAgent>> agentInstances = new ConcurrentHashMap<>();
// ==================== CRUD ====================
public List<AgentEntity> listAgents() {
return agentMapper.selectList(new LambdaQueryWrapper<AgentEntity>()
.orderByDesc(AgentEntity::getCreateTime));
}
/**
* 按工作区列出 Agent
*/
public List<AgentEntity> listAgentsByWorkspace(Long workspaceId) {
return listAgentsByWorkspace(workspaceId, null);
}
/**
* 按工作区列出 Agent可选过滤启用状态。
*
* @param enabled non-null restricts the result set to agents whose
* {@code enabled} column matches the given value.
* Pass {@code true} from chat selectors so disabled
* agents disappear from the picker; the admin
* management page passes {@code null} to keep
* disabled rows visible for re-enabling.
*/
public List<AgentEntity> listAgentsByWorkspace(Long workspaceId, Boolean enabled) {
LambdaQueryWrapper<AgentEntity> q = new LambdaQueryWrapper<AgentEntity>()
.eq(AgentEntity::getWorkspaceId, workspaceId);
if (enabled != null) {
q.eq(AgentEntity::getEnabled, enabled);
}
return agentMapper.selectList(q.orderByDesc(AgentEntity::getCreateTime));
}
public AgentEntity getAgent(Long id) {
AgentEntity entity = agentMapper.selectById(id);
if (entity == null) {
throw new MateClawException("err.agent.not_found", "Agent不存在: " + id);
}
return entity;
}
public AgentEntity createAgent(AgentEntity agent) {
agent.setEnabled(true);
if (agent.getAgentType() == null) {
agent.setAgentType("react");
}
requireUniqueName(agent, null);
agentMapper.insert(agent);
publishLifecycle(agent, "spawned");
return agent;
}
public AgentEntity updateAgent(AgentEntity agent) {
// Detect enabled-flag flip so the lifecycle event reflects the
// intent rather than every metadata edit. Reading the prior row
// is cheap and gives us a clean diff source.
AgentEntity prior = agentMapper.selectById(agent.getId());
// Only re-validate uniqueness when the name actually changes —
// a pure metadata edit (icon, prompt, ...) shouldn't pay the
// SELECT cost or risk a false positive against the row itself.
if (prior != null
&& agent.getName() != null
&& !agent.getName().equals(prior.getName())) {
// Workspace cannot be moved (Controller pins it to prior.workspaceId),
// so reuse it for the lookup even if the incoming DTO left it null.
if (agent.getWorkspaceId() == null) {
agent.setWorkspaceId(prior.getWorkspaceId());
}
requireUniqueName(agent, agent.getId());
}
agentMapper.updateById(agent);
agentInstances.remove(agent.getId());
if (prior != null && prior.getEnabled() != null
&& !prior.getEnabled().equals(agent.getEnabled())) {
publishLifecycle(agent,
Boolean.TRUE.equals(agent.getEnabled()) ? "enabled" : "disabled");
}
return agent;
}
/**
* Friendly business-code surface for the {@code (workspace_id, name)}
* unique index added in V102.
*
* <p>The wire shape is the project-wide R&lt;T&gt; envelope: HTTP status
* stays 200 (per the convention in {@code R.fail} and the axios
* interceptor in {@code mateclaw-ui/src/api/index.ts}); the 409 lives in
* the response body's {@code code} field so the front-end can branch
* without breaking on an axios error. Without this pre-check the
* duplicate save would surface as an opaque
* {@code DataIntegrityViolation} stack trace.
*
* @param excludeId when non-null, skip this row in the lookup so
* {@link #updateAgent} doesn't mistake the row for its
* own duplicate.
*/
private void requireUniqueName(AgentEntity agent, Long excludeId) {
if (agent.getName() == null || agent.getName().isBlank()) {
throw new MateClawException("err.agent.name_required", 400, "Agent 名称不能为空");
}
Long workspaceId = agent.getWorkspaceId() == null ? 1L : agent.getWorkspaceId();
LambdaQueryWrapper<AgentEntity> q = new LambdaQueryWrapper<AgentEntity>()
.eq(AgentEntity::getWorkspaceId, workspaceId)
.eq(AgentEntity::getName, agent.getName());
if (excludeId != null) {
q.ne(AgentEntity::getId, excludeId);
}
Long count = agentMapper.selectCount(q);
if (count != null && count > 0) {
throw new MateClawException("err.agent.duplicate_name", 409,
"工作区内已存在同名 Agent: " + agent.getName());
}
}
public void deleteAgent(Long id) {
AgentEntity prior = agentMapper.selectById(id);
agentMapper.deleteById(id);
agentInstances.remove(id);
if (prior != null) publishLifecycle(prior, "terminated");
}
/**
* Best-effort publish of an {@link AgentLifecycleEvent}. A publish
* failure must never roll back the agent CRUD that just succeeded —
* the agent_lifecycle trigger surface is observability, not the
* canonical record.
*/
private void publishLifecycle(AgentEntity agent, String phase) {
if (events == null || agent == null) return;
try {
events.publishEvent(new AgentLifecycleEvent(
agent.getWorkspaceId() == null ? 0L : agent.getWorkspaceId(),
agent.getId() == null ? 0L : agent.getId(),
agent.getName(),
phase,
System.currentTimeMillis()));
} catch (Exception e) {
log.warn("[AgentService] lifecycle publish failed for agent {} ({}): {}",
agent.getId(), phase, e.getMessage());
}
}
/**
* 清除 Agent 运行时缓存(绑定变更后需调用,使下次对话重新构建 Agent
*/
public void invalidateAgentCache(Long agentId) {
agentInstances.remove(agentId);
}
/**
* Invalidate the cached agent instance whenever one of its workspace files
* changes. The system prompt (which embeds MEMORY.md / PROFILE.md / structured
* memory) is baked into the cached instance at build time, so memory edits made
* via tools, consolidation, or cleanup would otherwise stay invisible until an
* agent config change or restart. Rebuilding on the next turn picks them up.
*/
@org.springframework.context.event.EventListener
public void onWorkspaceFileChanged(vip.mate.workspace.document.event.WorkspaceFileChangedEvent event) {
if (event.agentId() != null) {
agentInstances.remove(event.agentId());
}
}
// ==================== 运行时入口 ====================
public String chat(Long agentId, String message, String conversationId) {
return chat(agentId, message, conversationId, ChatOrigin.EMPTY);
}
/**
* RFC-063r §2.5: preferred entry — accepts the originating
* {@link ChatOrigin} so channel binding and workspace context propagate
* down to {@code @Tool} methods via Spring AI {@link org.springframework.ai.chat.model.ToolContext}.
*/
public String chat(Long agentId, String message, String conversationId, ChatOrigin origin) {
memoryRecallTracker.trackRecalls(agentId, message);
BaseAgent agent = getOrBuildAgentForConversation(agentId, conversationId);
ChatOriginHolder.set(origin != null ? origin : ChatOrigin.EMPTY);
try {
return withLifecycleSync(agentId, message, conversationId,
(msg, convId) -> agent.chat(msg, convId));
} finally {
ChatOriginHolder.clear();
}
}
/**
* Sync chat that also captures token usage and runtime model attribution
* from the agent graph's {@code _usage_final} event. Equivalent to
* subscribing to {@link #chatStructuredStream} and joining all content
* deltas — produces the same assistant text as {@link #chat} but exposes
* the usage figures so callers can persist them on the assistant message.
*
* <p>Prefer this entry over {@link #chat} for any path that writes the
* reply to {@code mate_message} (sync HTTP endpoint, voice WebSocket,
* cron task, post-approval replay); the plain {@link #chat} stays as the
* thin wrapper for fire-and-forget invocations where usage is not needed.
*/
public ChatResult chatWithUsage(Long agentId, String message, String conversationId) {
return chatWithUsage(agentId, message, conversationId, ChatOrigin.EMPTY);
}
public ChatResult chatWithUsage(Long agentId, String message, String conversationId, ChatOrigin origin) {
return collectChatResult(chatStructuredStream(agentId, message, conversationId, "", null, origin));
}
public Flux<String> chatStream(Long agentId, String message, String conversationId) {
return chatStream(agentId, message, conversationId, ChatOrigin.EMPTY);
}
public Flux<String> chatStream(Long agentId, String message, String conversationId, ChatOrigin origin) {
memoryRecallTracker.trackRecalls(agentId, message);
BaseAgent agent = getOrBuildAgentForConversation(agentId, conversationId);
// Capture the origin into a request-scoped holder; cleared on Flux
// termination so the next reactive subscriber doesn't inherit stale state.
ChatOrigin captured = origin != null ? origin : ChatOrigin.EMPTY;
return Flux.defer(() -> {
ChatOriginHolder.set(captured);
return withLifecycleFlux(agentId, message, conversationId,
(msg, convId) -> agent.chatStream(msg, convId),
chunk -> chunk);
}).doFinally(signal -> ChatOriginHolder.clear());
}
public Flux<StreamDelta> chatStructuredStream(Long agentId, String message, String conversationId) {
return chatStructuredStream(agentId, message, conversationId, "", null, ChatOrigin.EMPTY);
}
public Flux<StreamDelta> chatStructuredStream(Long agentId, String message, String conversationId,
String requesterId) {
return chatStructuredStream(agentId, message, conversationId, requesterId, null, ChatOrigin.EMPTY);
}
public Flux<StreamDelta> chatStructuredStream(Long agentId, String message, String conversationId,
String requesterId, ChatOrigin origin) {
return chatStructuredStream(agentId, message, conversationId, requesterId, null, origin);
}
public Flux<StreamDelta> chatStructuredStream(Long agentId, String message, String conversationId,
String requesterId, String thinkingLevel) {
return chatStructuredStream(agentId, message, conversationId, requesterId, thinkingLevel,
ChatOrigin.EMPTY);
}
public Flux<StreamDelta> chatStructuredStream(Long agentId, String message, String conversationId,
String requesterId, String thinkingLevel,
ChatOrigin origin) {
memoryRecallTracker.trackRecalls(agentId, message);
BaseAgent agent = getOrBuildAgentForConversation(agentId, conversationId);
// 设置请求级思考深度(通过 ThreadLocal 传递到 StateGraph 执行)
if (thinkingLevel != null && !thinkingLevel.isBlank()) {
ThinkingLevelHolder.set(thinkingLevel);
} else {
// 尝试从 Agent 默认配置读取
AgentEntity entity = getAgent(agentId);
if (entity != null && entity.getDefaultThinkingLevel() != null) {
ThinkingLevelHolder.set(entity.getDefaultThinkingLevel());
} else {
ThinkingLevelHolder.clear();
}
}
ChatOrigin captured = origin != null ? origin : ChatOrigin.EMPTY;
if (agent instanceof StructuredStreamCapable capable) {
return Flux.defer(() -> {
ChatOriginHolder.set(captured);
return withLifecycleFlux(agentId, message, conversationId,
(msg, convId) -> capable.chatStructuredStream(msg, convId,
requesterId != null ? requesterId : "")
.doFinally(signal -> ThinkingLevelHolder.clear()),
StreamDelta::content);
})
.doFinally(signal -> ChatOriginHolder.clear());
}
// 降级:不支持结构化流的 Agent包装为纯内容流
ThinkingLevelHolder.clear();
return Flux.defer(() -> {
ChatOriginHolder.set(captured);
return withLifecycleFlux(agentId, message, conversationId,
(msg, convId) -> agent.chatStream(msg, convId)
.map(chunk -> new StreamDelta(chunk, null)),
StreamDelta::content);
})
.doFinally(signal -> ChatOriginHolder.clear());
}
public String execute(Long agentId, String goal, String conversationId) {
return execute(agentId, goal, conversationId, ChatOrigin.EMPTY);
}
public String execute(Long agentId, String goal, String conversationId, ChatOrigin origin) {
memoryRecallTracker.trackRecalls(agentId, goal);
BaseAgent agent = getOrBuildAgentForConversation(agentId, conversationId);
ChatOriginHolder.set(origin != null ? origin : ChatOrigin.EMPTY);
try {
return withLifecycleSync(agentId, goal, conversationId,
(msg, convId) -> agent.execute(msg, convId));
} finally {
ChatOriginHolder.clear();
}
}
/**
* 带工具重放的 chat 调用(审批通过后由 ChannelMessageRouter 或 ApprovalController 调用)
*
* @param agentId Agent ID
* @param userMessage 用户消息(如"继续执行已批准的工具"
* @param conversationId 会话 ID
* @param toolCallPayload 要重放的工具调用 JSON
* @return Agent 回复
*/
public String chatWithReplay(Long agentId, String userMessage, String conversationId,
String toolCallPayload) {
return chatWithReplay(agentId, userMessage, conversationId, toolCallPayload, ChatOrigin.EMPTY);
}
public String chatWithReplay(Long agentId, String userMessage, String conversationId,
String toolCallPayload, ChatOrigin origin) {
memoryRecallTracker.trackRecalls(agentId, userMessage);
BaseAgent agent = getOrBuildAgentForConversation(agentId, conversationId);
ChatOriginHolder.set(origin != null ? origin : ChatOrigin.EMPTY);
try {
return withLifecycleSync(agentId, userMessage, conversationId,
(msg, convId) -> agent.chatWithReplay(msg, convId, toolCallPayload));
} finally {
ChatOriginHolder.clear();
}
}
/**
* Replay-after-approval that also captures token usage and runtime model
* attribution. Mirrors {@link #chatWithUsage} for the
* approval-resumption path used by {@code ChannelMessageRouter}.
*/
public ChatResult chatWithReplayWithUsage(Long agentId, String userMessage, String conversationId,
String toolCallPayload, ChatOrigin origin) {
return collectChatResult(chatWithReplayStream(agentId, userMessage, conversationId,
toolCallPayload, "", origin != null ? origin : ChatOrigin.EMPTY));
}
/**
* Subscribe to a structured stream and collapse it into a single
* {@link ChatResult}: append all content deltas, capture the trailing
* {@code _usage_final} event for token and model attribution.
*/
private ChatResult collectChatResult(Flux<StreamDelta> stream) {
StringBuilder content = new StringBuilder();
final int[] usage = {0, 0};
final String[] modelInfo = {null, null};
stream.doOnNext(delta -> {
if (delta.isEvent() && "_usage_final".equals(delta.eventType())) {
Map<String, Object> data = delta.eventData();
usage[0] = ((Number) data.getOrDefault("promptTokens", 0)).intValue();
usage[1] = ((Number) data.getOrDefault("completionTokens", 0)).intValue();
Object model = data.get("runtimeModelName");
Object provider = data.get("runtimeProviderId");
if (model != null) modelInfo[0] = model.toString();
if (provider != null) modelInfo[1] = provider.toString();
} else if (delta.content() != null) {
content.append(delta.content());
}
}).blockLast(Duration.ofMinutes(10));
return new ChatResult(content.toString(), usage[0], usage[1], modelInfo[0], modelInfo[1]);
}
/**
* 带工具重放的流式调用Web 端审批通过后使用,通过 SSE 推送结果)
*/
public Flux<StreamDelta> chatWithReplayStream(Long agentId, String userMessage, String conversationId,
String toolCallPayload) {
return chatWithReplayStream(agentId, userMessage, conversationId, toolCallPayload, "", ChatOrigin.EMPTY);
}
public Flux<StreamDelta> chatWithReplayStream(Long agentId, String userMessage, String conversationId,
String toolCallPayload, String requesterId) {
return chatWithReplayStream(agentId, userMessage, conversationId, toolCallPayload, requesterId,
ChatOrigin.EMPTY);
}
public Flux<StreamDelta> chatWithReplayStream(Long agentId, String userMessage, String conversationId,
String toolCallPayload, String requesterId,
ChatOrigin origin) {
memoryRecallTracker.trackRecalls(agentId, userMessage);
BaseAgent agent = getOrBuildAgentForConversation(agentId, conversationId);
ChatOrigin captured = origin != null ? origin : ChatOrigin.EMPTY;
return Flux.defer(() -> {
ChatOriginHolder.set(captured);
return withLifecycleFlux(agentId, userMessage, conversationId,
(msg, convId) -> agent.chatWithReplayStream(msg, convId, toolCallPayload,
requesterId != null ? requesterId : ""),
StreamDelta::content);
})
.doFinally(signal -> ChatOriginHolder.clear());
}
public AgentState getAgentState(Long agentId) {
Map<String, BaseAgent> variants = agentInstances.get(agentId);
if (variants == null || variants.isEmpty()) {
return AgentState.IDLE;
}
// An Agent may have several cached graph variants (one per pinned
// model). Report the first non-IDLE state so a turn running on any
// variant stays visible.
for (BaseAgent agent : variants.values()) {
AgentState state = agent.getState();
if (state != AgentState.IDLE) {
return state;
}
}
return AgentState.IDLE;
}
// ==================== 缓存管理 ====================
public void refreshAgent(Long agentId) {
agentInstances.remove(agentId);
log.info("Agent instance cache cleared: {}", agentId);
}
public void refreshAllAgents() {
agentInstances.clear();
log.info("All agent instance caches cleared");
}
@EventListener
public void onModelConfigChanged(ModelConfigChangedEvent event) {
refreshAllAgents();
log.info("Agent caches refreshed after model config change: {}", event.reason());
}
@EventListener
public void onToolGuardConfigChanged(vip.mate.tool.guard.service.ToolGuardConfigService.ToolGuardConfigChangedEvent event) {
refreshAllAgents();
log.info("Agent caches refreshed after tool guard config change (denied tools may have changed)");
}
/**
* Issue #289: an MCP server connecting / disconnecting / reconnecting
* changes the live tool set, but cached agents snapshot their tools at
* build time. Clear the cache so the next turn rebuilds against the
* current MCP tools instead of replying "from memory" with a stale,
* tool-less graph.
*/
@EventListener
public void onMcpServerChanged(vip.mate.tool.mcp.event.McpServerChangedEvent event) {
refreshAllAgents();
log.info("Agent caches refreshed after MCP server change: {}", event.reason());
}
/**
* Listen for MCP connection-loss events and clear the agent cache.
*
* <p>Previously this listener was intentionally omitted (the design
* doc said "only listen to McpServerChangedEvent, not
* McpConnectionLostEvent") because {@link McpServerService} auto-heals
* and publishes McpServerChangedEvent on reconnect. However, between
* disconnect and reconnect, cached agents still hold the old
* {@code AgentToolSet} snapshot whose MCP tool callbacks point at a
* dead client — calls either time out (5 min default) or throw.
*
* <p>Clearing the cache on disconnect ensures the next agent build
* sees the live connection state: {@link McpClientManager} will
* either skip the dead server or fall back to {@code lastGoodCallbacks}
* with proper error handling, rather than letting the LLM discover
* the breakage by timing out.
*
* <p>Cost is low: {@code McpServerService} already debounces reconnect
* attempts by 10s, and {@code refreshAllAgents} is a Map.clear().
* The subsequent reconnect will fire another McpServerChangedEvent,
* which clears the cache again — at most two clears per disconnect
* cycle, which is acceptable.
*/
@EventListener
public void onMcpConnectionLost(vip.mate.tool.mcp.event.McpConnectionLostEvent event) {
refreshAllAgents();
log.warn("Agent caches refreshed after MCP connection lost: serverId={}, reason={}",
event.serverId(), event.reason());
}
// ==================== Lifecycle helpers ====================
/**
* Wraps a synchronous agent call with lifecycle mediator hooks.
* When lifecycleMediatorEnabled is off, runs plainInvoke directly (Phase 0 behavior).
*
* P1-1 fix: prefetchAll result is now prepended to userMessage as &lt;memory-context&gt; block.
* P1-4 fix: N/A for sync (no cancel/error signal issue).
*/
private String withLifecycleSync(Long agentId, String message, String conversationId,
java.util.function.BiFunction<String, String, String> invoke) {
safeRegister(conversationId, agentId);
try {
if (!memoryProperties.isLifecycleMediatorEnabled()) {
return invoke.apply(message, conversationId);
}
String ownerKey = memoryOwnerResolver.resolve(ChatOriginHolder.get());
TurnContext ctx = new TurnContext(agentId, conversationId, conversationId, 0, message, ownerKey);
String memoryContext = lifecycleMediator.beforeLlmCall(ctx);
// Inject memory context into the user message (RFC-037 §3.3)
String enrichedMessage = injectMemoryContext(message, memoryContext);
String result = invoke.apply(enrichedMessage, conversationId);
lifecycleMediator.afterLlmCall(ctx, result != null ? result : "");
return result;
} finally {
safeUnregister(conversationId);
}
}
/**
* Wraps a streaming agent call with lifecycle mediator hooks.
* When lifecycleMediatorEnabled is off, runs plainInvoke directly (Phase 0 behavior).
*
* P1-1 fix: prefetchAll result is now prepended to userMessage.
* P1-4 fix: afterLlmCall only fires on COMPLETE signal, not on cancel/error.
*/
private <T> Flux<T> withLifecycleFlux(Long agentId, String message, String conversationId,
java.util.function.BiFunction<String, String, Flux<T>> invoke,
Function<T, String> contentExtractor) {
safeRegister(conversationId, agentId);
try {
if (!memoryProperties.isLifecycleMediatorEnabled()) {
return invoke.apply(message, conversationId)
.doFinally(s -> safeUnregister(conversationId));
}
String ownerKey = memoryOwnerResolver.resolve(ChatOriginHolder.get());
TurnContext ctx = new TurnContext(agentId, conversationId, conversationId, 0, message, ownerKey);
String memoryContext = lifecycleMediator.beforeLlmCall(ctx);
String enrichedMessage = injectMemoryContext(message, memoryContext);
StringBuilder reply = new StringBuilder();
return invoke.apply(enrichedMessage, conversationId)
.doOnNext(item -> {
String text = contentExtractor.apply(item);
if (text != null) {
reply.append(text);
}
})
.doOnComplete(() -> lifecycleMediator.afterLlmCall(ctx, reply.toString()))
.doOnError(e -> log.debug("[Memory] Stream error, skipping afterLlmCall: {}", e.getMessage()))
.doFinally(s -> safeUnregister(conversationId));
} catch (Exception e) {
// If invoke.apply() throws before the Flux is constructed, the
// doFinally above never runs — clean up here.
safeUnregister(conversationId);
throw e;
}
}
/** C5 helper — null-safe register so tests without the registry don't NPE. */
private void safeRegister(String conversationId, Long agentId) {
if (runningConversationRegistry != null) {
runningConversationRegistry.register(conversationId, agentId);
}
}
/** C5 helper — null-safe unregister so tests without the registry don't NPE. */
private void safeUnregister(String conversationId) {
if (runningConversationRegistry != null) {
runningConversationRegistry.unregister(conversationId);
}
}
/**
* Prepend memory-context block to user message if non-empty.
* Does not pollute build-time system prompt snapshot.
*/
private String injectMemoryContext(String message, String memoryContext) {
if (memoryContext == null || memoryContext.isBlank()) return message;
return memoryContext + "\n\n" + message;
}
// ==================== 内部方法 ====================
/**
* Resolve (and cache) the Agent graph for a conversation, honouring the
* conversation's pinned model. Conversations with no pin — IM channels
* before issue #183 fix, cron, sub-tasks, or rows not yet created —
* resolve to the shared Agent / global-default graph.
*
* <p>Defensive normalisation: a half-populated pair (provider but no
* model, or vice versa) is treated as unpinned. Without this guard, a
* partially-cleared admin UI write could end up cached as a key like
* {@code "volcano::"} which {@link #getOrBuildAgent} would then try to
* build, only to fail at provider-resolution time on every turn.
*/
private BaseAgent getOrBuildAgentForConversation(Long agentId, String conversationId) {
String provider = null;
String modelName = null;
if (conversationId != null && !conversationId.isBlank()) {
ConversationEntity conv = conversationMapper.selectOne(
new LambdaQueryWrapper<ConversationEntity>()
.eq(ConversationEntity::getConversationId, conversationId));
if (conv != null) {
provider = blankToNull(conv.getModelProvider());
modelName = blankToNull(conv.getModelName());
// Half-populated pair → treat as unpinned. Pinning requires
// a complete (provider, model) tuple — see #183 follow-up
// hardening so a stale row written by an earlier broken
// admin UI release doesn't loop the cache on an invalid key.
if (provider == null || modelName == null) {
provider = null;
modelName = null;
}
}
}
return getOrBuildAgent(agentId, provider, modelName);
}
/** Map empty / whitespace strings to null so the pinned-check is one branch. */
private static String blankToNull(String s) {
return (s == null || s.isBlank()) ? null : s;
}
private BaseAgent getOrBuildAgent(Long agentId) {
return getOrBuildAgent(agentId, null, null);
}
private BaseAgent getOrBuildAgent(Long agentId, String modelProvider, String modelName) {
boolean pinned = modelProvider != null && !modelProvider.isBlank()
&& modelName != null && !modelName.isBlank();
String modelKey = pinned ? modelProvider + "::" + modelName : "";
return agentInstances
.computeIfAbsent(agentId, id -> new ConcurrentHashMap<>())
.computeIfAbsent(modelKey, key -> {
AgentEntity entity = getAgent(agentId);
if (!Boolean.TRUE.equals(entity.getEnabled())) {
throw new MateClawException("err.agent.disabled", "Agent 已禁用: " + entity.getName());
}
return agentGraphBuilder.build(entity, modelProvider, modelName);
});
}
// ==================== StreamDelta ====================
public record StreamDelta(String content, String thinking, String eventType, Map<String, Object> eventData,
boolean persistenceOnly, boolean segmentOnly) {
// 兼容构造器(广播+持久化)
public StreamDelta(String content, String thinking) {
this(content, thinking, null, null, false, false);
}
// 显式 5-参构造器:保留旧调用点对 (content, thinking, eventType, eventData, persistenceOnly) 的兼容
public StreamDelta(String content, String thinking, String eventType,
Map<String, Object> eventData, boolean persistenceOnly) {
this(content, thinking, eventType, eventData, persistenceOnly, false);
}
/** 仅用于持久化,不再广播(内容已由 NodeStreamingChatHelper 实时广播过) */
public static StreamDelta persistOnly(String content, String thinking) {
return new StreamDelta(content, thinking, null, null, true, false);
}
/**
* Per-iteration narrative routing for ReasoningNode / SummarizingNode output.
*
* <p>The accumulator should:
* <ul>
* <li>append the text to the in-flight {@code segments} entry so the UI's
* segmented view still renders the intermediate "I'll look it up…"
* narration between tool cards;</li>
* <li>NOT broadcast — already broadcast live by NodeStreamingChatHelper;</li>
* <li>NOT append to the top-level {@code content} StringBuilder, which is
* what gets persisted as {@code mate_message.content}. That field
* should hold the final-answer span only — otherwise multiple
* iterations stack into "我来…让我…然后…" walls that next-turn replay
* sees as unanswered chain-of-thought (issue #120 narration leg).</li>
* </ul>
*
* <p>Implies {@code persistenceOnly} (no broadcast) at the accumulator
* layer, but is a stricter promise: <em>nothing</em> reaches the top-level
* persisted content field via this flavor.
*/
public static StreamDelta segmentOnly(String content, String thinking) {
return new StreamDelta(content, thinking, null, null, true, true);
}
public static StreamDelta empty() {
return new StreamDelta(null, null, null, null, false, false);
}
public static StreamDelta event(String type, Map<String, Object> data) {
return new StreamDelta(null, null, type, data, false, false);
}
public boolean isEvent() {
return eventType != null;
}
public boolean hasPayload() {
return StringUtils.hasText(content) || StringUtils.hasText(thinking);
}
public int contentLength() {
return content != null ? content.length() : 0;
}
public int thinkingLength() {
return thinking != null ? thinking.length() : 0;
}
}
// ==================== ChatResult ====================
/**
* Sync chat result carrying the assistant reply alongside the usage
* attribution that the streaming path exposes via the {@code _usage_final}
* event. Use this when callers need to persist {@code promptTokens} /
* {@code completionTokens} / {@code runtimeModel} / {@code runtimeProvider}
* on the assistant message row but cannot subscribe to the structured
* stream directly (cron tasks, sync HTTP endpoints, voice WebSocket,
* post-approval replays).
*/
public record ChatResult(String content, int promptTokens, int completionTokens,
String runtimeModel, String runtimeProvider) {
public static ChatResult contentOnly(String content) {
return new ChatResult(content != null ? content : "", 0, 0, null, null);
}
}
}