Compare commits

...
5 Commits
34 changed files with 2061 additions and 335 deletions
+1
View File
@@ -4,6 +4,7 @@
| 日期 | slug | 领域 | 关键词 | 状态 |
|---|---|---|---|---|
| 2026-07-03 | mvp-demo-trace-acceptance | MVP Demo/trace/acceptance | mvp-demo, trace API, diagnosis_session, agent_step, tool_invocation, feedback | openspec/changes/archive/2026-07-03-mvp-demo-trace-acceptance | archived |
| 2026-05-29 | chatmodel-abstraction | 解耦/多模型路由 | ChatModel, EmbeddingModel, DeepSeek, BGE-M3, SiliconFlow, Spring AI | archived |
| 2026-06-23 | phase1-infrastructure | 基础设施/文档管理 | MySQL, Redis, Milvus, Flyway, JPA, 向量检索, 类别过滤 | archived |
| 2026-06-24 | lookup-knowledge-integration | 知识库检索 | L0精确匹配, L1语义检索, frontmatter, 混合检索 | archived |
@@ -0,0 +1,65 @@
# MVP Demo Trace Acceptance
## Result
Accepted for implementation scope.
## Verification
### Static Verification
- Command: `mvn -q -DskipTests compile`
- Result: passed
- Notes: New trace controller, service, DTO, profile, verifier fallback, and test sources compile with the project.
### Script Verification
- Command: `mvn -q "-Dtest=DiagnosisTraceServiceTest,ChatServiceSupervisorAgentTest" test`
- Result: passed
- Notes: Covers successful trace aggregation, missing-session 404 path via `SessionNotFoundException`, low-confidence no-retry behavior, method-tool injection, and verifier fallback when Supervisor skips `chat_verifier`.
### OpenSpec Verification
- Command: `openspec validate mvp-demo-trace-acceptance --strict`
- Result: passed
### GitNexus Verification
- Result: skipped by user decision
- Notes: User requested subsequent project flow to bypass GitNexus.
### Manual / Runtime Verification
- Steps: Follow `mvp/demo/README.md` with `--spring.profiles.active=mvp-demo`.
- Result: passed
- Notes:
- Session `mvp-demo-payment-timeout-20260703-rerun2` completed as `SUCCESS`.
- Chat request returned `code=200`, `success=true`, and the same `sessionId`.
- Chat duration was `96316 ms`; persisted session duration was `95028 ms`.
- Trace API returned `code=200`, `returnedSteps=13`, `returnedTools=12`, `hasVerifier=true`, and `verifierVerdict=LOW_CONFID`.
- Trace agents included `planner,executor,verifier`.
- Trace tools included `lookup_knowledge,query_logs,query_metrics`.
- Feedback submission returned success, and a follow-up trace query showed `feedback=useful`.
- MySQL verification confirmed `agent_step` count `13` with agents `executor,planner,verifier`.
- MySQL verification confirmed `tool_invocation` count `12` with tools `lookup_knowledge,query_logs,query_metrics`.
## Completed Scope
- Added `GET /api/diagnosis/{sessionId}/trace`.
- Added read-only trace aggregation from persisted diagnosis tables.
- Added `mvp-demo` profile overlay.
- Added payment-timeout demo acceptance documentation.
- Added MVP note for interview storytelling.
- Added verifier fallback so runtime trace remains complete when Supervisor returns without `verifier_output`.
## Known Limits
- `mvp-demo` is not a fully offline mock runtime.
- Runtime still depends on available MySQL, Redis, Milvus/Zilliz, model, and embedding configuration.
- Sensitive configuration cleanup remains intentionally deferred.
- Supervisor can still make inefficient routing choices inside a single round; `ChatService` now invokes `chat_verifier` as a fallback when Supervisor returns without `verifier_output`, so trace completeness is preserved for the MVP demo.
## Handoff
- Runtime demo passed with current infrastructure.
- OpenSpec archive confirmation: requested by user after successful rerun.
@@ -0,0 +1,35 @@
# MVP Demo Trace Acceptance Brief
## Background
- User goal: make the MVP runnable, observable, and explainable for an Agent Engineer interview.
- Current problem: the system can execute diagnosis, but reviewers need a simple way to replay one session from final answer back to agent steps and tool evidence.
- Associated OpenSpec: `openspec/changes/mvp-demo-trace-acceptance/`
- Devflow scale: standard-light.
## Scope
- In scope:
- `mvp-demo` Spring profile overlay.
- `GET /api/diagnosis/{sessionId}/trace` read-only API.
- Trace aggregation DTO/service/controller.
- Focused service tests.
- Demo and acceptance documentation.
- Out of scope:
- Sensitive configuration cleanup.
- Full offline LLM/vector/database mock runtime.
- Database schema migration.
- Changes to chat execution, verifier routing, upload, or feedback behavior.
- Impact area:
- `src/main/java/com/superbiz/agent/controller`
- `src/main/java/com/superbiz/agent/service`
- `src/main/java/com/superbiz/agent/dto`
- `src/main/resources/application-mvp-demo.yml`
- `mvp/demo`
- `mvp/notes`
## OpenSpec Alignment
- proposal coverage: covered
- specs coverage: covered
- tasks coverage: covered
@@ -0,0 +1,87 @@
# MVP Demo Trace Acceptance Decisions
## Clarify
- Entry summary: continue the MVP toward a runnable and explainable demo by adding an `mvp-demo` profile, an end-to-end acceptance case, and a trace query API.
- Slug: `mvp-demo-trace-acceptance`
- Devflow scale: standard-light. The change adds a public read-only API and documentation, but does not alter core chat execution or persistence schemas.
## Context
- `devflow/index.md` was checked. Relevant history includes `session-storage`, `confidence-feedback`, `executor-action-memory-relevance`, and `chat-verifier-agent`.
- `mvp/notes/agent-engineering-decisions.md` already recommends the next phase as "可复现 MVP Demo", including `mvp-demo` profile, fixed diagnosis case, one-click request, and `GET /api/diagnosis/{sessionId}/trace`.
- `mvp/issues/ISS-003-mvp-design-implementation-review.md` identifies test stability, session traceability, verifier evidence chain, upload path, and SupervisorAgent consistency as recent MVP concerns. Security cleanup is intentionally deferred by user decision.
## Question Pool
| # | Dimension | Question | Mode | Status |
|---|---|---|---|---|
| Q1 | Terminology | Should "trace" mean persisted diagnosis execution evidence instead of transient frontend chat history? | evidence-driven | Resolved |
| Q2 | Boundary | Should this change modify chat execution or only expose existing persisted evidence? | evidence-driven | Resolved |
| Q3 | Acceptance | What proves the MVP flow is end-to-end enough for demo/interview use? | evidence-driven | Resolved |
| Q4 | Interface | What is the API impact level for `GET /api/diagnosis/{sessionId}/trace`? | evidence-driven | Resolved |
## Evidence-driven
| Conclusion | Evidence Source | Reported To User |
|---|---|---|
| Trace should aggregate persisted diagnosis evidence, not Redis-only chat history. | `DiagnosisSession`, `AgentStep`, `ToolInvocation` entities and repositories | Reported in progress update |
| Core chat execution does not need to change for this slice. | Existing unified chat path and SupervisorAgent commits; requested scope is demo/profile/trace/acceptance | Reported in progress update |
| End-to-end acceptance should cover start -> chat -> trace -> feedback. | `ChatController`, `FeedbackController`, traceable session id decision in MVP notes | Reported in progress update |
| Trace API is additive L3 because it is a new HTTP API for frontend/demo consumers. | sm-flow interface impact rules | Recorded in OpenSpec design |
## User-interview
| Question | User Words | Confirmation | OpenSpec Writeback |
|---|---|---|---|
| Should security/sensitive config cleanup be included? | "安全问题先不考虑"; "敏感配置先不做" | Confirmed | Non-goal |
| Should this be implemented under sm-flow? | "按照 sm-flow 的流程来实现吧" | Confirmed | This change follows sm-flow artifacts |
## Key Decisions
- Decision: Add a new trace API instead of embedding trace details in `/api/chat`.
- Reason: Chat execution and observability should stay decoupled.
- Impact: Demo can query trace after any successful chat request using the same session id.
- Risk accepted: Response shape is new and should be treated as demo-facing contract.
- Decision: Keep `mvp-demo` profile as configuration overlay, not a fully mocked standalone runtime.
- Reason: The current MVP still depends on real DB/Redis/Milvus/LLM for full chat execution; this change avoids inventing a fake runtime that hides integration behavior.
- Impact: Demo profile improves repeatability for logs/metrics, while docs remain explicit about required external services.
- Risk accepted: End-to-end acceptance may still require valid infrastructure and keys.
## Cross-Artifact Alignment
| Upstream -> Downstream | Check | Status |
|---|---|---|
| brief/prd -> proposal | Goal, scope, non-goals, and acceptance expectation are in proposal | Aligned |
| proposal -> design | Scope, constraints, and API impact are in design | Aligned |
| design -> specs/tasks | Trace DTO, controller/service, demo profile, and docs are represented | Aligned |
| specs -> tasks | Observable behavior is covered by executable tasks | Aligned |
## Architecture Audit
- Data path: HTTP trace request -> controller -> trace service -> repositories -> aggregate DTO -> `Result.success`.
- The service is read-only and does not mutate diagnosis, step, tool, or feedback state.
- No schema change is needed because all required fields already exist in `diagnosis_session`, `agent_step`, and `tool_invocation`.
- Main risk is response size for large sessions; MVP mitigates by returning previews already persisted by tools rather than raw external logs.
- The additive API is acceptable for MVP because old callers remain unaffected.
## Pre-apply Research
- Reference implementations read:
- `ChatController` for `/api` controller conventions.
- `FeedbackController` for simple API controller shape.
- `GlobalExceptionHandler` and `SessionNotFoundException` for 404 handling.
- `DiagnosisSessionRepository`, `AgentStepRepository`, `ToolInvocationRepository` for available queries.
- `DiagnosisSession`, `AgentStep`, `ToolInvocation` for fields.
- Impact analysis:
- `DiagnosisSessionRepository`: LOW, direct imports in service/controller paths.
- `AgentStepRepository`: HIGH because it participates in chat/AiOps flows. This change only consumes existing query methods and does not modify the repository.
- `ToolInvocationRepository`: LOW.
## Commit Gate
- OpenSpec proposal/design/specs/tasks exist.
- API impact: L3 additive collaboration API, documented in design and spec.
- User-confirmed non-goal: sensitive configuration cleanup remains out of scope.
- No unresolved user-interview questions remain for this slice.
@@ -0,0 +1,25 @@
# MVP Demo Trace Acceptance Evidence
## Evidence
| Source | Evidence | Conclusion | Reported |
|---|---|---|---|
| `DiagnosisSessionRepository` | Existing `findBySessionId(String)` query | Trace can locate the session without new repository methods | Yes |
| `AgentStepRepository` | Existing `findBySessionIdOrderByStepIndex(String)` query | Agent steps can be returned in execution order | Yes |
| `ToolInvocationRepository` | Existing `findBySessionIdOrderByIdAsc(String)` query | Tool evidence can be returned in persisted order | Yes |
| `GlobalExceptionHandler` | Handles `SessionNotFoundException` as HTTP 404 with `Result.error(404, ...)` | Missing trace can reuse existing error contract | Yes |
| `mvn -q "-Dtest=DiagnosisTraceServiceTest" test` | Command passed | Trace aggregation behavior is covered offline | Yes |
| `mvn -q -DskipTests compile` | Command passed | New code compiles with the full project | Yes |
| `gitnexus detect-changes --repo SuperBizAgent-java` | Command completed with `No changes detected` and line-ending warnings | Required GitNexus check ran; output likely does not capture newly added files | Yes |
## Evidence-driven Conclusions
- Conclusion: No database migration is required.
- Evidence: All trace fields are available from existing `diagnosis_session`, `agent_step`, and `tool_invocation` entities.
- Risk: Response shape becomes a new API contract.
- User confirmation: Not required; additive L3 API recorded in OpenSpec.
- Conclusion: Trace aggregation can be tested without external infrastructure.
- Evidence: `DiagnosisTraceServiceTest` uses mocked repositories and an `ObjectMapper`.
- Risk: Runtime integration still depends on configured infrastructure.
- User confirmation: Not required; limitation recorded in acceptance docs.
+94
View File
@@ -0,0 +1,94 @@
# MVP Demo Runbook
This demo proves the MVP flow from user question to persisted diagnosis trace.
## Prerequisites
- MySQL, Redis, Milvus/Zilliz, and LLM/embedding configuration are available through the current project configuration.
- Security and secret cleanup are intentionally out of scope for this MVP slice.
- The `mvp-demo` profile enables mock Prometheus and CLS providers so log and metric tools can return repeatable evidence.
## Start
```powershell
mvn spring-boot:run "-Dspring-boot.run.profiles=mvp-demo"
```
The service listens on:
```text
http://localhost:9900
```
## 1. Run Chat Diagnosis
```powershell
$sessionId = "mvp-demo-payment-timeout-001"
$body = @{
Id = $sessionId
Question = "支付接口最近出现超时,请结合知识库、日志和指标判断可能原因,并给出修复建议。"
} | ConvertTo-Json
Invoke-RestMethod `
-Method Post `
-Uri "http://localhost:9900/api/chat" `
-ContentType "application/json" `
-Body $body
```
Expected result:
- `data.success` is `true`.
- `data.sessionId` equals `mvp-demo-payment-timeout-001`.
- `data.answer` contains a diagnosis answer.
## 2. Query Trace
```powershell
Invoke-RestMethod `
-Method Get `
-Uri "http://localhost:9900/api/diagnosis/$sessionId/trace"
```
Expected result:
- `code` is `200`.
- `data.session.sessionId` equals the chat session id.
- `data.steps` contains planner/executor/verifier records for complex questions.
- `data.toolInvocations` contains evidence tool calls such as `lookup_knowledge`, `query_logs`, or `query_metrics`.
- `data.session.selfEvaluation` contains verifier or rule evaluation when available.
## 3. Submit Feedback
```powershell
$feedback = @{
sessionId = $sessionId
feedback = "useful"
} | ConvertTo-Json
Invoke-RestMethod `
-Method Post `
-Uri "http://localhost:9900/api/feedback" `
-ContentType "application/json" `
-Body $feedback
```
Expected result:
- `success` is `true`.
- A later trace query shows `data.session.feedback` as `useful`.
## Demo Story
The important interview story is:
```text
one session id
-> user question
-> multi-agent execution
-> evidence tools
-> verifier/self-evaluation
-> final answer
-> feedback
-> trace API for replay and audit
```
+39
View File
@@ -0,0 +1,39 @@
# Payment Timeout Acceptance Case
## Goal
Validate that the MVP can diagnose a payment timeout incident and expose the complete trace for replay.
## Input
- Session id: `mvp-demo-payment-timeout-001`
- Question: `支付接口最近出现超时,请结合知识库、日志和指标判断可能原因,并给出修复建议。`
- Profile: `mvp-demo`
## Acceptance Criteria
1. Chat returns a successful answer with the same session id.
2. Trace API returns session metadata, final answer, ordered agent steps, and ordered tool invocations.
3. Trace contains enough evidence to explain which tools were used and whether verifier/self-evaluation was persisted.
4. Feedback can be submitted for the same session id.
5. A follow-up trace query shows the persisted feedback value.
## Trace Fields To Inspect
- `data.session.query`
- `data.session.answer`
- `data.session.selfEvaluation`
- `data.session.feedback`
- `data.steps[*].agentName`
- `data.steps[*].thought`
- `data.toolInvocations[*].toolName`
- `data.toolInvocations[*].inputParams`
- `data.toolInvocations[*].outputPreview`
- `data.toolInvocations[*].retrievalDetails`
- `data.summary`
## Known Limits
- This case is not a full offline test. It still requires valid infrastructure for chat, persistence, vector search, and model calls.
- Mock logs and metrics are enabled by the `mvp-demo` profile to make those evidence tools repeatable.
- Sensitive configuration cleanup is deferred by current MVP priority.
+268
View File
@@ -0,0 +1,268 @@
# MVP Agent 工程决策记录
本文记录 MVP 实现过程中已经落地的一些关键修复、取舍和工程判断。目标不是写流水账,而是沉淀面试时可以讲清楚的 Agent 工程思路。
---
## 1. 统一流式与非流式 Chat 主链路
### 背景
早期 `/api/chat` 和 `/api/chat_stream` 是两条不同实现:
- 非流式接口会走复杂度判断,并可能进入 Planner / Executor / Verifier 多 Agent 流程。
- 流式接口直接创建单个 ReactAgent,然后 `agent.stream()` 输出 token。
这导致两个接口表面都是 chat,实际能力不一致:流式接口不会进入 verifier、不会沉淀完整诊断链路,也不容易和 `diagnosis_session`、`tool_invocation` 对齐。
### 决策
将两个接口统一到同一条核心链路:
```text
getOrCreateSession
-> 读取会话历史
-> ChatService.executeChatWithStrategy(...)
-> 写回会话历史
```
接口差异只保留在传输层:
- `/api/chat` 返回完整 JSON。
- `/api/chat_stream` 通过 SSE 分块发送最终答案。
### 取舍
这样会牺牲原来的 token 级实时流式体验,但换来业务行为一致、诊断链路一致、Verifier 和 evidence trace 一致。
对 MVP 来说,优先保证“同一个问题不因接口不同而进入不同智能链路”,比 token 级流式更重要。
---
## 2. 会话 ID 与诊断链路统一
### 背景
原实现中:
- `ChatController` 用前端传入的 `Id` 在 JVM 内存里维护历史消息。
- `ChatService` 每次执行又生成新的 8 位 sessionId,作为 `diagnosis_session` 和工具调用追踪 ID。
这会造成前端会话、后端诊断会话、工具证据链三者分裂。
### 决策
将前端 chat session id 作为后端诊断链路的主 session id:
- Redis `SessionContext` 保存聊天历史。
- `diagnosis_session.session_id` 复用同一个 id。
- `RunnableConfig.metadata.sessionId` 和 `SessionContextHolder` 也使用同一个 id。
- `tool_invocation`、`agent_step`、verifier evaluation 都可按同一 session id 串起来。
### 企业级意义
Agent 系统最怕“答得出来但查不清”。统一 session id 后,一次用户请求可以完整追踪:
```text
用户问题 -> Agent 步骤 -> 工具调用 -> Verifier 判断 -> 最终答案 -> 用户反馈
```
这是可观测、可审计、可复盘的基础。
---
## 3. 引入统一 ToolInvocationRecorder
### 背景
Verifier 需要结构化证据链,但原实现只有 `lookup_knowledge` 主动写入 `tool_invocation`。
`query_logs`、`query_metrics` 虽然返回 JSON,但没有统一落库,导致 verifier 看不到日志、指标等 evidence tool 的稳定记录。
### 决策
新增 `ToolInvocationRecorder`,作为所有 evidence tool 的统一落库入口。
当前接入:
- `lookup_knowledge`
- `query_logs`
- `query_metrics`
记录字段包括:
- tool name
- input params
- output preview
- output length
- success
- error message
- duration
- trace id / domain details
### 企业级意义
这一步把 Agent 从“模型说它查过”推进到“系统能证明它查过”。
后续 verifier 不应该依赖模型自由文本回忆工具调用,而应该消费结构化 trace summary。
---
## 4. Verifier 作为事实约束层
### 背景
普通 Agent 很容易在工具调用后直接生成答案,但企业场景更关心:
- 关键结论有没有证据
- 证据是直接证据还是间接支持
- 哪些事实缺口需要人工介入
- 工具失败时是否诚实降级
### 决策
保留 Planner / Executor / Verifier 三角色:
- Planner 负责拆解问题。
- Executor 负责执行查询与形成初稿。
- Verifier 负责基于 `tool_trace_summary` 做事实核查。
Verifier 输出结构化 JSON,包括:
- verdict
- groundedness_score
- critical_fact_count
- facts_checked
- rationale
### 取舍
Verifier 会增加一次模型调用成本,但换来可解释性和质量约束。对企业级 Agent 来说,这是值得的。
---
## 5. 从手写编排切换到 SupervisorAgent
### 背景
之前 `ChatService.executeChatComplex()` 中构建了 `SupervisorAgent`,但实际仍然手写调用:
```text
planner -> executor -> verifier
```
这会造成代码与设计不一致,维护者容易误以为当前已经由 Supervisor 调度。
### 决策
复杂问题真正切换到 `SupervisorAgent.invoke(...)`。
Supervisor 负责路由:
```text
chat_supervisor -> chat_planner
chat_supervisor -> chat_executor
chat_supervisor -> chat_verifier
chat_supervisor -> FINISH
```
外层仍保留:
- verifier 输出解析
- PASS / LOW_CONFID / REJECT 判定
- retry context
- fallback
- evaluation 入库
### 验证
新增离线专项测试 `ChatServiceSupervisorAgentTest`,使用 scripted `ChatModel` 验证真实 SupervisorAgent 路由顺序,不依赖真实 LLM、MySQL、Redis。
### 企业级意义
这让项目不只是“自己写 if/else 多 Agent”,而是使用框架原生 multi-agent orchestration,同时保留业务层的质量门控。
---
## 6. 文档上传路径语义统一
### 背景
上传文档时,`DocumentManagementService.saveToLocal()` 返回带 `knowledge_base` 前缀的路径。
而 `KnowledgeIndexService.readDocument()` 又执行:
```java
Paths.get(knowledgeBasePath, filePath)
```
这可能拼出:
```text
knowledge_base/knowledge_base/...
```
最终表现为 L0 命中文档,但读取原文失败。
### 决策
统一路径语义:
- 新上传文档存相对 `knowledge.base-path` 的路径,例如 `payment/runbook.md`。
- `readDocument()` 兼容新旧路径:
- 相对路径
- 已带 base path 的旧相对路径
- 绝对路径
### 企业级意义
知识库检索不能只看“命中”,还要保证命中后的内容可读、可引用、可追踪。
这是 RAG / Agent 系统里很典型的工程细节:检索质量问题不一定来自模型,也可能来自路径、元数据、索引和原文之间的语义不一致。
---
## 7. MVP 阶段的优先级取舍
当前主动暂缓的问题:
- 敏感配置外置与密钥轮换
- CORS / Redis 反序列化安全边界
- 默认 `mvn test` 离线化
原因不是这些不重要,而是当前目标是先跑通并讲清楚 MVP Agent 工程闭环。
短期优先目标:
```text
可演示 -> 可观测 -> 可验证 -> 可复盘
```
安全和完整测试体系属于企业落地必须项,但可以在 MVP 主链路稳定后作为下一阶段补齐。
---
## 8. 后续建议
下一阶段建议聚焦“可复现 MVP Demo”:
1. 增加 `local-demo` 或 `mvp-demo` profile。
2. 准备固定诊断 case,例如“支付接口超时”。
3. 提供一键初始化知识库样例。
4. 提供一键触发复杂诊断请求的脚本。
5. 增加 trace 查询接口:
```text
GET /api/diagnosis/{sessionId}/trace
```
该接口聚合:
- diagnosis_session
- agent_step
- tool_invocation
- verifier evaluation
- final answer
- feedback
这样 MVP 就能从“功能实现”升级为“企业级 Agent 工程作品”。
+39
View File
@@ -0,0 +1,39 @@
# MVP Demo Profile 与 Trace 查询接口
## 背景
MVP 已经能跑多 Agent 诊断、工具调用、Verifier 和反馈,但对外展示时仍然缺少一个稳定的复盘入口。面试官或评审如果想确认一次 Agent 回答是否可信,不能只看最终答案,还需要看到用户原始问题、Agent 步骤顺序、工具调用证据、Verifier / self-evaluation、最终答案和用户反馈。
## 决策
新增 `mvp-demo` profile 和 trace 查询接口:
```text
GET /api/diagnosis/{sessionId}/trace
```
接口聚合:
- `diagnosis_session`
- `agent_step`
- `tool_invocation`
- `self_evaluation`
- `feedback`
同时在 `mvp/demo` 下沉淀端到端验收 case,把启动、提问、查 trace、提交 feedback 串成一条可演示路径。
## 取舍
`mvp-demo` profile 不是完整离线 mock 环境,仍然复用当前真实 DB / Redis / Milvus / LLM 配置,只显式打开日志和指标 mock。原因是当前阶段目标是展示企业级 Agent 工程闭环,不是隐藏真实集成复杂度。
这让 MVP 的讲述从“我实现了一个聊天接口”升级为:
```text
我实现了一条可执行、可观测、可验收、可复盘的 Agent 诊断链路。
```
## 面试表达
- 我没有把 trace 塞进 chat 返回值,而是做成独立只读观测接口,保持执行链路和观测链路解耦。
- Trace API 复用已经沉淀的 `diagnosis_session`、`agent_step`、`tool_invocation` 三张表,没有引入新的 schema 风险。
- Demo profile 只做最小 overlay,让日志和指标工具可重复,保留真实基础设施集成,方便说明 MVP 与生产化之间的差距。
@@ -0,0 +1 @@
mvp-demo-trace-acceptance committed on 2026-07-03
@@ -0,0 +1,2 @@
schema: spec-driven
created: 2026-07-03
@@ -0,0 +1,29 @@
{
"id": "mvp-demo-trace-acceptance",
"metadata": {
"status": "committed",
"created_at": "2026-07-03",
"updated_at": "2026-07-03",
"implementation_status": "implemented"
},
"summary": "Add an MVP demo profile, a read-only diagnosis trace API, and an end-to-end acceptance case.",
"artifacts": {
"proposal": "proposal.md",
"design": "design.md",
"tasks": "tasks.md",
"specs": [
"specs/mvp-demo-trace-acceptance/spec.md"
],
"devflow": "devflow/projects/2026-07-03-mvp-demo-trace-acceptance"
},
"tasks": [
"Add DiagnosisTraceResponse DTO",
"Add DiagnosisTraceService aggregation",
"Add DiagnosisTraceController endpoint",
"Add mvp-demo profile",
"Add MVP demo acceptance documentation",
"Add focused trace service tests",
"Run targeted verification and GitNexus change detection",
"Update MVP notes and devflow acceptance"
]
}
@@ -0,0 +1,70 @@
## Context
The MVP already persists diagnosis execution data across three tables:
- `diagnosis_session`: query, status, answer, counts, feedback, and `self_evaluation`.
- `agent_step`: ordered agent execution records.
- `tool_invocation`: evidence tool calls and retrieval metadata.
Recent work unified chat session ids and persisted tool invocations, so a single session id can now connect user input, agent steps, evidence tools, verifier evaluation, final answer, and feedback. The missing piece is a read-only aggregation API and a documented demo profile/workflow that a reviewer can run without reading database tables manually.
## Goals / Non-Goals
**Goals:**
- Add a trace API that returns one aggregated view for a diagnosis session.
- Keep the trace API read-only and based on existing persistence tables.
- Add an `mvp-demo` profile that makes the demo intent explicit and keeps mock log/metric tools enabled.
- Add a documented end-to-end acceptance case for start, chat, trace query, and feedback.
- Add focused tests for trace aggregation.
**Non-Goals:**
- Do not clean up committed sensitive configuration in this change.
- Do not add database migrations.
- Do not alter `/api/chat`, `/api/chat_stream`, verifier routing, feedback, or document upload behavior.
- Do not create a fully offline fake LLM runtime.
## Decisions
| Decision | Choice | Alternative Considered | Rationale |
|---|---|---|---|
| Trace API shape | Add `GET /api/diagnosis/{sessionId}/trace` | Extend `/api/chat` response | Trace is an observability concern and should not make chat responses larger or change chat clients. |
| Aggregation ownership | New `DiagnosisTraceService` | Put aggregation in controller | Keeps controller thin and allows focused unit tests with mocked repositories. |
| Response DTO | Dedicated nested DTO | Return raw entities or maps | DTO avoids leaking JPA entity details and gives a stable demo-facing contract. |
| Missing session handling | Throw `SessionNotFoundException` and use existing global 404 handler | Return empty success payload | A missing trace is a real lookup miss and should be visible to callers. |
| `self_evaluation` handling | Return raw JSON string and best-effort parsed JSON | Parse only, or ignore parse failures | Raw value preserves evidence even if JSON shape evolves; parsed value improves frontend/demo readability. |
| Demo profile | Add `application-mvp-demo.yml` overlay | Change default `application.yml` | Overlay avoids disturbing current runtime and keeps demo choices explicit. |
## Interface Impact
- Level: L3 collaboration API.
- Reason: This adds a new HTTP endpoint and response contract intended for frontend/demo/reviewer consumption.
- Compatibility: Additive only. Existing callers do not need to change.
- Documentation: The endpoint is documented in the MVP demo acceptance case.
## Data Structures
The trace response contains:
- `session`: session id, query, status, flow, counts, timing, created/updated time, final answer, raw self-evaluation JSON, parsed self-evaluation object, and feedback.
- `steps`: ordered agent steps with step index, agent name, model input/output, thought, tool flag, duration, token count, and created time.
- `toolInvocations`: ordered tool records with id, step id, tool name, input params, output preview, retrieval metadata, duration, success, error, and created time.
- `summary`: counts derived from the returned collections and session fields.
## Risks / Trade-offs
- [Risk] Trace responses may become large for long sessions. -> Mitigation: the MVP returns persisted previews and structured metadata, not raw full external logs.
- [Risk] `self_evaluation` JSON shape may evolve. -> Mitigation: return both raw and best-effort parsed forms.
- [Risk] Demo profile still depends on real DB/Redis/Milvus/LLM. -> Mitigation: document prerequisites and keep mock logs/metrics enabled for repeatable tool evidence.
- [Risk] New endpoint becomes a de facto frontend contract. -> Mitigation: use a dedicated DTO and document L3 additive API impact.
## Migration Plan
- Deploying this change requires only application restart with the new code.
- No database migration is required.
- Rollback is deleting the new endpoint/profile/docs; persisted data remains unchanged.
## Open Questions
- None for this slice. Security and full offline test profile remain deferred by explicit user decision.
@@ -0,0 +1,29 @@
## Why
The MVP can already execute multi-agent diagnosis, persist session traces, and collect feedback, but it is still hard to demonstrate as a complete enterprise-style workflow. A demo profile, a trace query API, and an explicit end-to-end acceptance case make the project runnable, observable, and explainable for interview and portfolio review.
## What Changes
- Add an `mvp-demo` Spring profile that keeps the existing external infrastructure contract but turns on mock log and metric providers for repeatable demonstrations.
- Add a read-only trace query API: `GET /api/diagnosis/{sessionId}/trace`.
- Aggregate `diagnosis_session`, `agent_step`, `tool_invocation`, verifier/self-evaluation, final answer, and feedback into one trace response.
- Add an end-to-end MVP acceptance case that documents startup, chat request, trace query, and feedback submission.
- Add focused service tests for trace aggregation without requiring MySQL, Redis, Milvus, or a real LLM.
- Record the design decision in MVP notes for interview storytelling.
## Capabilities
### New Capabilities
- `mvp-demo-trace-acceptance`: Covers the MVP demo profile, trace query API, and end-to-end acceptance workflow for a reproducible agent diagnosis demo.
### Modified Capabilities
- None.
## Impact
- Affected code: new trace controller/service/DTOs, `application-mvp-demo.yml`, unit tests, MVP demo documentation.
- Affected API: adds `GET /api/diagnosis/{sessionId}/trace`. This is an additive L3 collaboration API because it is intended for frontend, demo, and external reviewer consumption.
- Affected runtime behavior: no change to chat execution, verifier, feedback, document upload, or persistence semantics.
- Non-goals: no sensitive configuration cleanup, no database schema migration, no replacement of existing chat endpoints, no full offline mock LLM implementation.
@@ -0,0 +1,33 @@
## ADDED Requirements
### Requirement: Diagnosis trace can be queried by session id
The system SHALL expose a read-only HTTP endpoint `GET /api/diagnosis/{sessionId}/trace` that returns the persisted diagnosis trace for the requested session id.
#### Scenario: Existing session trace is returned
- **WHEN** a caller requests trace data for a session id that exists in `diagnosis_session`
- **THEN** the system returns a success response containing the session summary, ordered agent steps, ordered tool invocations, self-evaluation data, final answer, and feedback
#### Scenario: Missing session returns not found
- **WHEN** a caller requests trace data for a session id that does not exist in `diagnosis_session`
- **THEN** the system returns a 404 response using the existing session-not-found error contract
### Requirement: Trace aggregation is read-only
The system MUST build trace output from existing persisted diagnosis tables and MUST NOT mutate diagnosis sessions, agent steps, tool invocations, feedback, or chat session state while serving the trace request.
#### Scenario: Trace query does not change persisted state
- **WHEN** a caller requests `GET /api/diagnosis/{sessionId}/trace`
- **THEN** the system reads `diagnosis_session`, `agent_step`, and `tool_invocation` records and returns an aggregate without saving any of those records
### Requirement: MVP demo profile is available
The system SHALL provide an `mvp-demo` Spring profile that documents the demo runtime intent and keeps mock log and metric providers enabled for repeatable diagnosis demonstrations.
#### Scenario: Demo profile loads mock evidence providers
- **WHEN** the application starts with `--spring.profiles.active=mvp-demo`
- **THEN** `prometheus.mock-enabled` and `cls.mock-enabled` are enabled by profile configuration
### Requirement: End-to-end MVP acceptance case is documented
The project SHALL include an end-to-end acceptance case that demonstrates start-up, chat diagnosis, trace query, and feedback submission using the same session id.
#### Scenario: Reviewer follows the acceptance case
- **WHEN** a reviewer follows the documented MVP demo acceptance steps
- **THEN** they can run the application, submit a diagnosis question, query the trace endpoint, and submit feedback for the same session id
@@ -0,0 +1,25 @@
## 1. Trace Query API
- [x] 1.1 Add a `DiagnosisTraceResponse` DTO that represents session summary, ordered agent steps, ordered tool invocations, and derived summary counts.
- [x] 1.2 Add `DiagnosisTraceService` that loads `DiagnosisSession`, `AgentStep`, and `ToolInvocation` records by session id and builds the response.
- [x] 1.3 Add `DiagnosisTraceController` with `GET /api/diagnosis/{sessionId}/trace`.
- [x] 1.4 Return 404 through `SessionNotFoundException` when the requested diagnosis session does not exist.
## 2. Demo Profile And Acceptance Case
- [x] 2.1 Add `src/main/resources/application-mvp-demo.yml` with MVP demo profile overlays and mock logs/metrics enabled.
- [x] 2.2 Add `mvp/demo/README.md` documenting prerequisites, startup, chat request, trace query, and feedback submission.
- [x] 2.3 Add a concrete payment-timeout acceptance case with request/response expectations.
## 3. Tests And Verification
- [x] 3.1 Add focused unit tests for `DiagnosisTraceService` success and missing-session behavior.
- [x] 3.2 Run targeted tests for the new trace service.
- [x] 3.3 Run compile verification.
- [x] 3.4 Run GitNexus change detection before commit or handoff.
## 4. Notes And Flow Records
- [x] 4.1 Update MVP engineering notes with the demo/trace decision.
- [x] 4.2 Update OpenSpec tasks as work completes.
- [x] 4.3 Record verification results in devflow acceptance notes.
@@ -0,0 +1,37 @@
## Purpose
Provide a repeatable MVP demo flow that can run a chat diagnosis, expose its persisted execution trace, and submit feedback for the same session id.
## Requirements
### Requirement: Diagnosis trace can be queried by session id
The system SHALL expose a read-only HTTP endpoint `GET /api/diagnosis/{sessionId}/trace` that returns the persisted diagnosis trace for the requested session id.
#### Scenario: Existing session trace is returned
- **WHEN** a caller requests trace data for a session id that exists in `diagnosis_session`
- **THEN** the system returns a success response containing the session summary, ordered agent steps, ordered tool invocations, self-evaluation data, final answer, and feedback
#### Scenario: Missing session returns not found
- **WHEN** a caller requests trace data for a session id that does not exist in `diagnosis_session`
- **THEN** the system returns a 404 response using the existing session-not-found error contract
### Requirement: Trace aggregation is read-only
The system MUST build trace output from existing persisted diagnosis tables and MUST NOT mutate diagnosis sessions, agent steps, tool invocations, feedback, or chat session state while serving the trace request.
#### Scenario: Trace query does not change persisted state
- **WHEN** a caller requests `GET /api/diagnosis/{sessionId}/trace`
- **THEN** the system reads `diagnosis_session`, `agent_step`, and `tool_invocation` records and returns an aggregate without saving any of those records
### Requirement: MVP demo profile is available
The system SHALL provide an `mvp-demo` Spring profile that documents the demo runtime intent and keeps mock log and metric providers enabled for repeatable diagnosis demonstrations.
#### Scenario: Demo profile loads mock evidence providers
- **WHEN** the application starts with `--spring.profiles.active=mvp-demo`
- **THEN** `prometheus.mock-enabled` and `cls.mock-enabled` are enabled by profile configuration
### Requirement: End-to-end MVP acceptance case is documented
The project SHALL include an end-to-end acceptance case that demonstrates start-up, chat diagnosis, trace query, and feedback submission using the same session id.
#### Scenario: Reviewer follows the acceptance case
- **WHEN** a reviewer follows the documented MVP demo acceptance steps
- **THEN** they can run the application, submit a diagnosis question, query the trace endpoint, and submit feedback for the same session id
@@ -2,6 +2,7 @@ package com.superbiz.agent.agent.tool;
import com.fasterxml.jackson.annotation.JsonProperty;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.superbiz.agent.service.ToolInvocationRecorder;
import lombok.Data;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -34,6 +35,11 @@ public class QueryLogsTools {
public static final String TOOL_GET_AVAILABLE_LOG_TOPICS = "getAvailableLogTopics";
private final ObjectMapper objectMapper = new ObjectMapper();
private final ToolInvocationRecorder toolInvocationRecorder;
public QueryLogsTools(ToolInvocationRecorder toolInvocationRecorder) {
this.toolInvocationRecorder = toolInvocationRecorder;
}
@Value("${cls.mock-enabled:false}")
private boolean mockEnabled;
@@ -55,6 +61,7 @@ public class QueryLogsTools {
"Call this tool first before querying logs to understand what log topics are available. " +
"Returns a list of log topics with their names, descriptions, and example queries.")
public String getAvailableLogTopics() {
long startTime = System.currentTimeMillis();
logger.info("获取可用的日志主题列表");
try {
@@ -123,11 +130,15 @@ public class QueryLogsTools {
output.setMessage(String.format("共有 %d 个可用的日志主题。建议使用默认地域 'ap-guangzhou' 或省略 region 参数", topics.size()));
return objectMapper.writerWithDefaultPrettyPrinter().writeValueAsString(output);
String response = objectMapper.writerWithDefaultPrettyPrinter().writeValueAsString(output);
recordInvocation(startTime, "get_available_log_topics", null, null, null, response, true, null, "logs");
return response;
} catch (Exception e) {
logger.error("获取日志主题列表失败", e);
return "{\"success\":false,\"message\":\"获取日志主题列表失败: " + e.getMessage() + "\"}";
String response = "{\"success\":false,\"message\":\"获取日志主题列表失败: " + e.getMessage() + "\"}";
recordInvocation(startTime, "get_available_log_topics", null, null, null, response, false, e.getMessage(), "logs");
return response;
}
}
@@ -164,6 +175,7 @@ public class QueryLogsTools {
@ToolParam(description = "查询条件,支持 Lucene 语法,如 level:ERROR OR cpu_usage:>80;为空时返回该主题近 5 条核心日志") String query,
@ToolParam(description = "返回日志条数,默认20,最大100") Integer limit) {
long startTime = System.currentTimeMillis();
int actualLimit = (limit == null || limit <= 0) ? 20 : Math.min(limit, 100);
String safeQuery = query == null ? "" : query;
@@ -178,7 +190,10 @@ public class QueryLogsTools {
logger.info("使用 Mock 数据,返回 {} 条日志", logEntries.size());
} else {
// 真实模式:调用 CLS API(这里预留接口,后续实现)
return buildErrorResponse("CLS 真实查询尚未实现,请启用 mock 模式进行测试");
String response = buildErrorResponse("CLS 真实查询尚未实现,请启用 mock 模式进行测试");
recordInvocation(startTime, safeQuery, region, logTopic, actualLimit, response, false,
"CLS 真实查询尚未实现,请启用 mock 模式进行测试", normalizeTopicDomain(logTopic));
return response;
}
// 构建成功响应
@@ -193,15 +208,51 @@ public class QueryLogsTools {
String jsonResult = objectMapper.writerWithDefaultPrettyPrinter().writeValueAsString(output);
logger.info("日志查询完成: 找到 {} 条日志", logEntries.size());
recordInvocation(startTime, safeQuery, region, logTopic, actualLimit, jsonResult,
!logEntries.isEmpty(), logEntries.isEmpty() ? "未找到匹配的日志" : null,
normalizeTopicDomain(logTopic));
return jsonResult;
} catch (Exception e) {
logger.error("查询日志失败", e);
return buildErrorResponse("查询失败: " + e.getMessage());
String response = buildErrorResponse("查询失败: " + e.getMessage());
recordInvocation(startTime, safeQuery, region, logTopic, actualLimit, response, false,
e.getMessage(), normalizeTopicDomain(logTopic));
return response;
}
}
private void recordInvocation(long startTime, String query, String region, String logTopic, Integer limit,
String output, boolean success, String errorMessage, String topicDomain) {
Map<String, Object> input = new HashMap<>();
input.put("query", query == null || query.isBlank() ? "DEFAULT_QUERY" : query);
if (region != null) {
input.put("region", region);
}
if (logTopic != null) {
input.put("log_topic", logTopic);
}
if (limit != null) {
input.put("limit", limit);
}
input.put("mock_enabled", mockEnabled);
toolInvocationRecorder.recordEvidenceTool(
"query_logs",
input,
output,
success,
startTime,
errorMessage,
topicDomain
);
}
private String normalizeTopicDomain(String logTopic) {
return logTopic == null || logTopic.isBlank() ? "logs" : logTopic;
}
/**
* 构建 Mock 日志数据
@@ -2,6 +2,7 @@ package com.superbiz.agent.agent.tool;
import com.fasterxml.jackson.annotation.JsonProperty;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.superbiz.agent.service.ToolInvocationRecorder;
import lombok.Data;
import okhttp3.OkHttpClient;
import okhttp3.Request;
@@ -30,6 +31,11 @@ public class QueryMetricsTools {
public static final String TOOL_QUERY_PROMETHEUS_ALERTS = "queryPrometheusAlerts";
private final ObjectMapper objectMapper = new ObjectMapper();
private final ToolInvocationRecorder toolInvocationRecorder;
public QueryMetricsTools(ToolInvocationRecorder toolInvocationRecorder) {
this.toolInvocationRecorder = toolInvocationRecorder;
}
@Value("${prometheus.base-url}")
private String prometheusBaseUrl;
@@ -59,6 +65,7 @@ public class QueryMetricsTools {
"This tool retrieves all currently active/firing alerts including their labels, annotations, state, and values. " +
"Use this tool when you need to check what alerts are currently firing, investigate alert conditions, or monitor alert status.")
public String queryPrometheusAlerts() {
long startTime = System.currentTimeMillis();
logger.info("开始查询 Prometheus 活动告警, Mock模式: {}", mockEnabled);
try {
@@ -73,7 +80,9 @@ public class QueryMetricsTools {
PrometheusAlertsResult result = fetchPrometheusAlerts();
if (!"success".equals(result.getStatus())) {
return buildErrorResponse("Prometheus API 返回非成功状态: " + result.getStatus(), result.getError());
String response = buildErrorResponse("Prometheus API 返回非成功状态: " + result.getStatus(), result.getError());
recordInvocation(startTime, response, false, result.getError());
return response;
}
// 转换为简化格式,对于相同的 alertname,只保留第一个
@@ -110,15 +119,30 @@ public class QueryMetricsTools {
String jsonResult = objectMapper.writerWithDefaultPrettyPrinter().writeValueAsString(output);
logger.info("Prometheus 告警查询完成: 找到 {} 个告警", simplifiedAlerts.size());
recordInvocation(startTime, jsonResult, true, null);
return jsonResult;
} catch (Exception e) {
logger.error("查询 Prometheus 告警失败", e);
return buildErrorResponse("查询失败", e.getMessage());
String response = buildErrorResponse("查询失败", e.getMessage());
recordInvocation(startTime, response, false, e.getMessage());
return response;
}
}
private void recordInvocation(long startTime, String output, boolean success, String errorMessage) {
toolInvocationRecorder.recordEvidenceTool(
"query_metrics",
Map.of("query", "active_prometheus_alerts", "mock_enabled", mockEnabled),
output,
success,
startTime,
errorMessage,
"prometheus_alerts"
);
}
/**
* 构建 Mock 告警数据
* 与 aiops-docs 文档中的告警类型对应:
@@ -1,32 +1,30 @@
package com.superbiz.agent.controller;
import com.alibaba.cloud.ai.graph.NodeOutput;
import com.alibaba.cloud.ai.graph.OverAllState;
import com.alibaba.cloud.ai.graph.agent.ReactAgent;
import com.alibaba.cloud.ai.graph.streaming.OutputType;
import com.alibaba.cloud.ai.graph.streaming.StreamingOutput;
import lombok.Getter;
import lombok.Setter;
import com.superbiz.agent.domain.model.SessionContext;
import com.superbiz.agent.service.AiOpsService;
import com.superbiz.agent.service.ChatService;
import com.superbiz.agent.service.session.SessionManager;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.ai.chat.model.ChatModel;
import org.springframework.ai.tool.ToolCallback;
import org.springframework.ai.tool.ToolCallbackProvider;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.http.MediaType;
import org.springframework.http.ResponseEntity;
import org.springframework.web.bind.annotation.*;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
import reactor.core.publisher.Flux;
import java.io.IOException;
import java.time.LocalDateTime;
import java.time.ZoneId;
import java.util.*;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.locks.ReentrantLock;
/**
* 统一 API 控制器
@@ -44,17 +42,20 @@ public class ChatController {
@Autowired
private ChatService chatService;
@Autowired
private SessionManager sessionManager;
@Autowired(required = false)
private ToolCallbackProvider tools;
private final ExecutorService executor = Executors.newCachedThreadPool();
// 存储会话信息
private final Map<String, SessionInfo> sessions = new ConcurrentHashMap<>();
// 最大历史消息窗口大小(成对计算:用户消息+AI回复=1对)
private static final int MAX_WINDOW_SIZE = 6;
@Value("${session.ttl-seconds:3600}")
private long sessionTtlSeconds;
/**
* 普通对话接口(支持工具调用)
* 与 /chat_react 逻辑一致,但直接返回完整结果而非流式输出
@@ -71,10 +72,10 @@ public class ChatController {
}
// 获取或创建会话
SessionInfo session = getOrCreateSession(request.getId());
SessionContext session = getOrCreateSession(request.getId());
// 获取历史消息
List<Map<String, String>> history = session.getHistory();
List<Map<String, String>> history = session.getMessageHistorySnapshot();
logger.info("会话历史消息对数: {}", history.size() / 2);
// 获取注入的 ChatModel
@@ -88,13 +89,14 @@ public class ChatController {
// 根据问题复杂度自动选择单 Agent 或多 Agent
logger.info("开始 ReactAgent 对话(支持自动工具调用)");
ChatService.ChatResult result = chatService.executeChatWithStrategy(chatModel, toolCallbacks,
request.getQuestion(), history);
request.getQuestion(), history, session.getSessionId());
String fullAnswer = result.answer();
// 更新会话历史
session.addMessage(request.getQuestion(), fullAnswer);
session.addChatMessagePair(request.getQuestion(), fullAnswer, MAX_WINDOW_SIZE);
sessionManager.updateSession(session);
logger.info("已更新会话历史 - SessionId: {}, 当前消息对数: {}",
request.getId(), session.getMessagePairCount());
session.getSessionId(), session.getMessagePairCount());
return ResponseEntity.ok(ApiResponse.success(ChatResponse.success(fullAnswer, result.sessionId())));
@@ -116,9 +118,11 @@ public class ChatController {
return ResponseEntity.ok(ApiResponse.error("会话ID不能为空"));
}
SessionInfo session = sessions.get(request.getId());
if (session != null) {
session.clearHistory();
Optional<SessionContext> session = sessionManager.getSession(request.getId());
if (session.isPresent()) {
SessionContext context = session.get();
context.clearMessageHistory();
sessionManager.updateSession(context);
return ResponseEntity.ok(ApiResponse.success("会话历史已清空"));
} else {
return ResponseEntity.ok(ApiResponse.error("会话不存在"));
@@ -131,8 +135,8 @@ public class ChatController {
}
/**
* ReactAgent 对话接口(SSE 流式模式,支持多轮对话,支持自动工具调用,例如获取当前时间,查询日志,告警等)
* 支持 session 管理,保留对话历史
* 对话接口(SSE 流式模式)
* 与 /chat 使用同一条 ChatService 策略链路,区别仅在于通过 SSE 分块返回最终答案。
*/
@PostMapping(value = "/chat_stream", produces = "text/event-stream;charset=UTF-8")
public SseEmitter chatStream(@RequestBody ChatRequest request) {
@@ -155,10 +159,10 @@ public class ChatController {
logger.info("收到 ReactAgent 对话请求 - SessionId: {}, Question: {}", request.getId(), request.getQuestion());
// 获取或创建会话
SessionInfo session = getOrCreateSession(request.getId());
SessionContext session = getOrCreateSession(request.getId());
// 获取历史消息
List<Map<String, String>> history = session.getHistory();
List<Map<String, String>> history = session.getMessageHistorySnapshot();
logger.info("ReactAgent 会话历史消息对数: {}", history.size() / 2);
// 获取注入的 ChatModel
@@ -167,92 +171,25 @@ public class ChatController {
// 记录可用工具
chatService.logAvailableTools();
logger.info("开始 ReactAgent 流式对话(支持自动工具调用)");
ToolCallback[] toolCallbacks = tools != null ? tools.getToolCallbacks() : new ToolCallback[0];
// 构建系统提示词(包含历史消息)
String systemPrompt = chatService.buildSystemPrompt(history);
logger.info("开始统一 ChatService 对话(SSE 分块返回)");
ChatService.ChatResult result = chatService.executeChatWithStrategy(chatModel, toolCallbacks,
request.getQuestion(), history, session.getSessionId());
String fullAnswer = result.answer() == null ? "" : result.answer();
logger.info("统一 ChatService 对话完成 - SessionId: {}, 答案长度: {}",
result.sessionId(), fullAnswer.length());
// 创建 ReactAgent
ReactAgent agent = chatService.createReactAgent(chatModel, systemPrompt);
session.addChatMessagePair(request.getQuestion(), fullAnswer, MAX_WINDOW_SIZE);
sessionManager.updateSession(session);
logger.info("已更新会话历史 - SessionId: {}, 当前消息对数: {}",
session.getSessionId(), session.getMessagePairCount());
// 用于累积完整答案
StringBuilder fullAnswerBuilder = new StringBuilder();
// 使用 agent.stream() 进行流式对话
Flux<NodeOutput> stream = agent.stream(request.getQuestion());
stream.subscribe(
output -> {
try {
// 检查是否为 StreamingOutput 类型
if (output instanceof StreamingOutput streamingOutput) {
OutputType type = streamingOutput.getOutputType();
// 处理模型推理的流式输出
if (type == OutputType.AGENT_MODEL_STREAMING) {
// 流式增量内容,逐步显示
String chunk = streamingOutput.message().getText();
if (chunk != null && !chunk.isEmpty()) {
fullAnswerBuilder.append(chunk);
// 实时发送到前端
emitter.send(SseEmitter.event()
.name("message")
.data(SseMessage.content(chunk), MediaType.APPLICATION_JSON));
logger.info("发送流式内容: {}", chunk);
}
} else if (type == OutputType.AGENT_MODEL_FINISHED) {
// 模型推理完成
logger.info("模型输出完成");
} else if (type == OutputType.AGENT_TOOL_FINISHED) {
// 工具调用完成
logger.info("工具调用完成: {}", output.node());
} else if (type == OutputType.AGENT_HOOK_FINISHED) {
// Hook 执行完成
logger.debug("Hook 执行完成: {}", output.node());
}
}
} catch (IOException e) {
logger.error("发送流式消息失败", e);
throw new RuntimeException(e);
}
},
error -> {
// 错误处理
logger.error("ReactAgent 流式对话失败", error);
try {
emitter.send(SseEmitter.event()
.name("message")
.data(SseMessage.error(error.getMessage()), MediaType.APPLICATION_JSON));
} catch (IOException ex) {
logger.error("发送错误消息失败", ex);
}
emitter.completeWithError(error);
},
() -> {
// 完成处理
try {
String fullAnswer = fullAnswerBuilder.toString();
logger.info("ReactAgent 流式对话完成 - SessionId: {}, 答案长度: {}",
request.getId(), fullAnswer.length());
// 更新会话历史
session.addMessage(request.getQuestion(), fullAnswer);
logger.info("已更新会话历史 - SessionId: {}, 当前消息对数: {}",
request.getId(), session.getMessagePairCount());
// 发送完成标记
emitter.send(SseEmitter.event()
.name("message")
.data(SseMessage.done(), MediaType.APPLICATION_JSON));
emitter.complete();
} catch (IOException e) {
logger.error("发送完成消息失败", e);
emitter.completeWithError(e);
}
}
);
sendContentChunks(emitter, fullAnswer);
emitter.send(SseEmitter.event()
.name("message")
.data(SseMessage.done(), MediaType.APPLICATION_JSON));
emitter.complete();
} catch (Exception e) {
logger.error("ReactAgent 对话初始化失败", e);
@@ -365,12 +302,13 @@ public class ChatController {
try {
logger.info("收到获取会话信息请求 - SessionId: {}", sessionId);
SessionInfo session = sessions.get(sessionId);
if (session != null) {
Optional<SessionContext> session = sessionManager.getSession(sessionId);
if (session.isPresent()) {
SessionContext context = session.get();
SessionInfoResponse response = new SessionInfoResponse();
response.setSessionId(sessionId);
response.setMessagePairCount(session.getMessagePairCount());
response.setCreateTime(session.createTime);
response.setMessagePairCount(context.getMessagePairCount());
response.setCreateTime(toEpochMillis(context.getCreatedAt()));
return ResponseEntity.ok(ApiResponse.success(response));
} else {
return ResponseEntity.ok(ApiResponse.error("会话不存在"));
@@ -384,107 +322,39 @@ public class ChatController {
// ==================== 辅助方法 ====================
private SessionInfo getOrCreateSession(String sessionId) {
if (sessionId == null || sessionId.isEmpty()) {
sessionId = UUID.randomUUID().toString();
}
return sessions.computeIfAbsent(sessionId, SessionInfo::new);
private SessionContext getOrCreateSession(String sessionId) {
String resolvedSessionId = (sessionId == null || sessionId.isEmpty())
? UUID.randomUUID().toString()
: sessionId;
return sessionManager.getSession(resolvedSessionId)
.orElseGet(() -> {
SessionContext context = SessionContext.builder()
.sessionId(resolvedSessionId)
.status("ACTIVE")
.ttl(sessionTtlSeconds)
.build();
sessionManager.createSession(context, sessionTtlSeconds);
return context;
});
}
// ==================== 内部类 ====================
/**
* 会话信息
* 管理单个会话的历史消息,支持自动清理和线程安全
*/
private static class SessionInfo {
private final String sessionId;
// 存储历史消息对:[{"role": "user", "content": "..."}, {"role": "assistant", "content": "..."}]
private final List<Map<String, String>> messageHistory;
private final long createTime;
private final ReentrantLock lock;
public SessionInfo(String sessionId) {
this.sessionId = sessionId;
this.messageHistory = new ArrayList<>();
this.createTime = System.currentTimeMillis();
this.lock = new ReentrantLock();
private long toEpochMillis(LocalDateTime time) {
if (time == null) {
return 0L;
}
return time.atZone(ZoneId.systemDefault()).toInstant().toEpochMilli();
}
/**
* 添加一对消息(用户问题 + AI回复)
* 自动管理历史消息窗口大小
*/
public void addMessage(String userQuestion, String aiAnswer) {
lock.lock();
try {
// 添加用户消息
Map<String, String> userMsg = new HashMap<>();
userMsg.put("role", "user");
userMsg.put("content", userQuestion);
messageHistory.add(userMsg);
// 添加AI回复
Map<String, String> assistantMsg = new HashMap<>();
assistantMsg.put("role", "assistant");
assistantMsg.put("content", aiAnswer);
messageHistory.add(assistantMsg);
// 自动清理:保持最多 MAX_WINDOW_SIZE 对消息
// 每对消息包含2条记录(user + assistant)
int maxMessages = MAX_WINDOW_SIZE * 2;
while (messageHistory.size() > maxMessages) {
// 成对删除最旧的消息(删除前2条)
messageHistory.remove(0); // 删除最旧的用户消息
if (!messageHistory.isEmpty()) {
messageHistory.remove(0); // 删除对应的AI回复
}
}
logger.debug("会话 {} 更新历史消息,当前消息对数: {}",
sessionId, messageHistory.size() / 2);
} finally {
lock.unlock();
}
private void sendContentChunks(SseEmitter emitter, String content) throws IOException {
if (content == null || content.isEmpty()) {
return;
}
/**
* 获取历史消息(线程安全)
* 返回副本以避免并发修改
*/
public List<Map<String, String>> getHistory() {
lock.lock();
try {
return new ArrayList<>(messageHistory);
} finally {
lock.unlock();
}
}
/**
* 清空历史消息
*/
public void clearHistory() {
lock.lock();
try {
messageHistory.clear();
logger.info("会话 {} 历史消息已清空", sessionId);
} finally {
lock.unlock();
}
}
/**
* 获取当前消息对数
*/
public int getMessagePairCount() {
lock.lock();
try {
return messageHistory.size() / 2;
} finally {
lock.unlock();
}
int chunkSize = 80;
for (int i = 0; i < content.length(); i += chunkSize) {
int end = Math.min(i + chunkSize, content.length());
emitter.send(SseEmitter.event()
.name("message")
.data(SseMessage.content(content.substring(i, end)), MediaType.APPLICATION_JSON));
}
}
@@ -0,0 +1,24 @@
package com.superbiz.agent.controller;
import com.superbiz.agent.dto.DiagnosisTraceResponse;
import com.superbiz.agent.dto.Result;
import com.superbiz.agent.service.DiagnosisTraceService;
import lombok.RequiredArgsConstructor;
import org.springframework.http.ResponseEntity;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
@RestController
@RequestMapping("/api/diagnosis")
@RequiredArgsConstructor
public class DiagnosisTraceController {
private final DiagnosisTraceService diagnosisTraceService;
@GetMapping("/{sessionId}/trace")
public ResponseEntity<Result<DiagnosisTraceResponse>> getTrace(@PathVariable String sessionId) {
return ResponseEntity.ok(Result.success(diagnosisTraceService.getTrace(sessionId)));
}
}
@@ -8,7 +8,9 @@ import lombok.NoArgsConstructor;
import java.io.Serializable;
import java.time.LocalDateTime;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
/**
* 会话上下文数据类
@@ -53,6 +55,12 @@ public class SessionContext implements Serializable {
@Builder.Default
private List<ToolCall> toolCalls = new ArrayList<>();
/**
* 聊天消息历史:[{"role":"user","content":"..."}, {"role":"assistant","content":"..."}]
*/
@Builder.Default
private List<Map<String, String>> messageHistory = new ArrayList<>();
/**
* 会话创建时间
*/
@@ -79,6 +87,64 @@ public class SessionContext implements Serializable {
this.lastActiveAt = LocalDateTime.now();
}
/**
* 添加一对聊天消息,并按消息对数裁剪窗口。
*/
public void addChatMessagePair(String userQuestion, String assistantAnswer, int maxPairCount) {
if (this.messageHistory == null) {
this.messageHistory = new ArrayList<>();
}
Map<String, String> userMessage = new HashMap<>();
userMessage.put("role", "user");
userMessage.put("content", userQuestion);
this.messageHistory.add(userMessage);
Map<String, String> assistantMessage = new HashMap<>();
assistantMessage.put("role", "assistant");
assistantMessage.put("content", assistantAnswer);
this.messageHistory.add(assistantMessage);
int maxMessages = Math.max(maxPairCount, 0) * 2;
while (maxMessages > 0 && this.messageHistory.size() > maxMessages) {
this.messageHistory.remove(0);
if (!this.messageHistory.isEmpty()) {
this.messageHistory.remove(0);
}
}
this.lastActiveAt = LocalDateTime.now();
}
/**
* 获取聊天历史副本,避免调用方直接修改内部列表。
*/
public List<Map<String, String>> getMessageHistorySnapshot() {
if (this.messageHistory == null || this.messageHistory.isEmpty()) {
return new ArrayList<>();
}
List<Map<String, String>> snapshot = new ArrayList<>();
for (Map<String, String> message : this.messageHistory) {
snapshot.add(new HashMap<>(message));
}
return snapshot;
}
/**
* 清空聊天历史。
*/
public void clearMessageHistory() {
if (this.messageHistory == null) {
this.messageHistory = new ArrayList<>();
} else {
this.messageHistory.clear();
}
this.lastActiveAt = LocalDateTime.now();
}
public int getMessagePairCount() {
return this.messageHistory == null ? 0 : this.messageHistory.size() / 2;
}
/**
* 更新最后活跃时间
*/
@@ -0,0 +1,102 @@
package com.superbiz.agent.dto;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;
import java.time.LocalDateTime;
import java.util.List;
import java.util.Map;
@Data
@Builder
@NoArgsConstructor
@AllArgsConstructor
public class DiagnosisTraceResponse {
private SessionTrace session;
private List<AgentStepTrace> steps;
private List<ToolInvocationTrace> toolInvocations;
private TraceSummary summary;
@Data
@Builder
@NoArgsConstructor
@AllArgsConstructor
public static class SessionTrace {
private Long id;
private String sessionId;
private String query;
private String status;
private String agentFlow;
private Integer totalDurationMs;
private Integer totalTokenCount;
private Integer stepCount;
private Integer toolCallCount;
private String answer;
private String selfEvaluationRaw;
private Map<String, Object> selfEvaluation;
private String feedback;
private LocalDateTime createdAt;
private LocalDateTime updatedAt;
}
@Data
@Builder
@NoArgsConstructor
@AllArgsConstructor
public static class AgentStepTrace {
private Long id;
private String sessionId;
private Integer stepIndex;
private String agentName;
private String modelInput;
private String modelOutput;
private String thought;
private Boolean hasToolCall;
private Integer durationMs;
private Integer tokenCount;
private LocalDateTime createdAt;
}
@Data
@Builder
@NoArgsConstructor
@AllArgsConstructor
public static class ToolInvocationTrace {
private Long id;
private String sessionId;
private Long stepId;
private String toolName;
private String inputParamsRaw;
private Map<String, Object> inputParams;
private String outputPreview;
private Integer outputLength;
private String retrievalLayer;
private Integer l0MatchCount;
private Integer l1MatchCount;
private Boolean truncated;
private String relevanceLevel;
private String dedupReason;
private String retrievalDetailsRaw;
private Map<String, Object> retrievalDetails;
private Integer durationMs;
private Boolean success;
private String errorMessage;
private LocalDateTime createdAt;
}
@Data
@Builder
@NoArgsConstructor
@AllArgsConstructor
public static class TraceSummary {
private int persistedStepCount;
private int returnedStepCount;
private int persistedToolCallCount;
private int returnedToolCallCount;
private boolean hasVerifierEvaluation;
private boolean hasFeedback;
}
}
@@ -3,7 +3,7 @@ package com.superbiz.agent.service;
import com.alibaba.cloud.ai.graph.OverAllState;
import com.alibaba.cloud.ai.graph.RunnableConfig;
import com.alibaba.cloud.ai.graph.agent.ReactAgent;
import com.alibaba.cloud.ai.graph.agent.flow.agent.SupervisorAgent;
import com.alibaba.cloud.ai.graph.agent.flow.agent.SequentialAgent;
import com.alibaba.cloud.ai.graph.exception.GraphRunnerException;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
@@ -104,6 +104,9 @@ public class ChatService {
@Value("${verifier.low-confidence-threshold:0.5}")
private double verifierLowConfidenceThreshold;
@Value("${chat.complex.retry-on-low-confidence:false}")
private boolean retryOnLowConfidence;
/** 多 Agent Chat 的 Prompt */
private String chatPlannerPrompt;
private String chatExecutorPrompt;
@@ -211,16 +214,19 @@ public class ChatService {
/**
* 动态构建方法工具数组
* 根据 cls.mock-enabled 决定是否包含 QueryLogsTools
* 根据已注入的 Bean 暴露本地工具,避免 mock/真实模式下漏注入。
*/
public Object[] buildMethodToolsArray() {
List<Object> methodTools = new ArrayList<>();
methodTools.add(dateTimeTools);
methodTools.add(lookupKnowledgeTool);
if (queryLogsTools != null) {
// Mock 模式:包含 QueryLogsTools
return new Object[]{dateTimeTools, lookupKnowledgeTool};
} else {
// 真实模式:不包含 QueryLogsTools(由 MCP 提供日志查询功能)
return new Object[]{dateTimeTools, lookupKnowledgeTool, queryMetricsTools};
methodTools.add(queryLogsTools);
}
if (queryMetricsTools != null) {
methodTools.add(queryMetricsTools);
}
return methodTools.toArray();
}
/**
@@ -272,19 +278,18 @@ public class ChatService {
* @return ChatResult(answer + sessionId)
*/
public ChatResult executeChat(ReactAgent agent, String question) throws GraphRunnerException {
return executeChat(agent, question, null);
}
public ChatResult executeChat(ReactAgent agent, String question, String requestedSessionId) throws GraphRunnerException {
logger.info("========================================");
logger.info("📝 用户问题: {}", question);
String sessionId = UUID.randomUUID().toString().substring(0, 8);
String sessionId = resolveSessionId(requestedSessionId);
long startTime = System.currentTimeMillis();
// 创建诊断会话
DiagnosisSession session = DiagnosisSession.builder()
.sessionId(sessionId)
.query(question)
.status("RUNNING")
.agentFlow("CHAT")
.build();
// 创建或更新诊断会话
DiagnosisSession session = startDiagnosisSession(sessionId, question);
diagnosisSessionRepository.save(session);
// 设置 ThreadLocal 上下文(LookupKnowledgeTool 通过此获取 sessionId)
@@ -335,31 +340,38 @@ public class ChatService {
*/
public ChatResult executeChatWithStrategy(ChatModel chatModel, ToolCallback[] toolCallbacks,
String question, List<Map<String, String>> history) throws GraphRunnerException {
return executeChatWithStrategy(chatModel, toolCallbacks, question, history, null);
}
public ChatResult executeChatWithStrategy(ChatModel chatModel, ToolCallback[] toolCallbacks,
String question, List<Map<String, String>> history,
String requestedSessionId) throws GraphRunnerException {
if (QuestionComplexity.isComplex(question)) {
logger.info("📊 问题判定为复杂,使用多 Agent(Planner + Executor)执行");
return executeChatComplex(chatModel, toolCallbacks, question, history);
return executeChatComplex(chatModel, toolCallbacks, question, history, requestedSessionId);
} else {
logger.info("📊 问题判定为简单,使用单 Agent 执行");
String systemPrompt = buildSystemPrompt(history);
ReactAgent agent = createReactAgent(chatModel, systemPrompt);
return executeChat(agent, question);
return executeChat(agent, question, requestedSessionId);
}
}
/**
* 多 Agent 复杂对话执行(Planner + Executor + Supervisor)
* 多 Agent 复杂对话执行(Planner -> Executor -> Verifier)
*/
public ChatResult executeChatComplex(ChatModel chatModel, ToolCallback[] toolCallbacks,
String question, List<Map<String, String>> history) throws GraphRunnerException {
String sessionId = UUID.randomUUID().toString().substring(0, 8);
return executeChatComplex(chatModel, toolCallbacks, question, history, null);
}
public ChatResult executeChatComplex(ChatModel chatModel, ToolCallback[] toolCallbacks,
String question, List<Map<String, String>> history,
String requestedSessionId) throws GraphRunnerException {
String sessionId = resolveSessionId(requestedSessionId);
long startTime = System.currentTimeMillis();
DiagnosisSession session = DiagnosisSession.builder()
.sessionId(sessionId)
.query(question)
.status("RUNNING")
.agentFlow("CHAT")
.build();
DiagnosisSession session = startDiagnosisSession(sessionId, question);
diagnosisSessionRepository.save(session);
SessionContextHolder.setSessionId(sessionId);
@@ -383,19 +395,34 @@ public class ChatService {
ReactAgent executor = buildChatExecutorAgent(chatModel, toolCallbacks, history, retryContext);
ReactAgent verifier = buildChatVerifierAgent(chatModel);
SupervisorAgent supervisor = SupervisorAgent.builder()
.name("chat_supervisor")
.description("负责按单轮顺序调度 Planner、Executor、Verifier 的多 Agent 控制器")
.model(chatModel)
.systemPrompt(buildSupervisorPrompt(round))
SequentialAgent workflow = SequentialAgent.builder()
.name("chat_workflow")
.description("按固定顺序执行 Planner、Executor、Verifier 的多 Agent 工作流")
.subAgents(List.of(planner, executor, verifier))
.build();
String plannerPlan = callAgent(planner, buildPlannerInput(question, retryContext), config);
answer = callAgent(executor, buildExecutorInput(question, plannerPlan, retryContext), config);
String workflowInput = buildWorkflowInput(question, retryContext);
Optional<OverAllState> stateOptional = workflow.invoke(workflowInput, config);
if (stateOptional.isEmpty()) {
finalDecision = buildVerifierFallbackDecision(round, "workflow 未返回有效状态");
answer = buildLowConfidenceOutput(answer, finalDecision);
persistVerifierEvaluation(session, finalDecision, round);
break;
}
String plannerPlan = extractStateText(stateOptional, "planner_plan");
answer = extractStateText(stateOptional, "executor_feedback");
VerifierContextHolder.setExecutorFinalAnswer(answer);
String verifierOutput = callAgent(verifier, "VERIFY", config);
String verifierOutput = extractStateText(stateOptional, "verifier_output");
if ((verifierOutput == null || verifierOutput.isBlank()) && answer != null && !answer.isBlank()) {
verifierOutput = invokeVerifierFallback(verifier, question, round, config);
}
finalDecision = parseVerifierDecision(verifierOutput, round);
logger.debug("Sequential workflow round {} finished: plannerPlanLength={}, answerLength={}, verifierOutputLength={}",
round,
plannerPlan != null ? plannerPlan.length() : 0,
answer != null ? answer.length() : 0,
verifierOutput != null ? verifierOutput.length() : 0);
if (finalDecision == null) {
finalDecision = buildVerifierFallbackDecision(round, "verifier_output 缺失或无法解析");
@@ -416,7 +443,10 @@ public class ChatService {
break;
}
if (finalDecision.groundednessScore() >= verifierLowConfidenceThreshold || round == 2) {
boolean shouldRetry = retryOnLowConfidence
&& finalDecision.groundednessScore() < verifierLowConfidenceThreshold
&& round < 2;
if (!shouldRetry) {
answer = buildLowConfidenceOutput(answer, finalDecision);
persistVerifierEvaluation(session, finalDecision, round);
break;
@@ -524,28 +554,52 @@ public class ChatService {
.build();
}
private String callAgent(ReactAgent agent, String input, RunnableConfig config) throws GraphRunnerException {
return agent.call(input, config).getText();
private String resolveSessionId(String requestedSessionId) {
if (requestedSessionId != null && !requestedSessionId.isBlank()) {
return requestedSessionId;
}
return UUID.randomUUID().toString().substring(0, 8);
}
private String buildPlannerInput(String question, String retryContext) {
if (retryContext == null || retryContext.isBlank()) {
return question;
}
return question + "\n\n--- 补充约束 ---\n" + retryContext;
private DiagnosisSession startDiagnosisSession(String sessionId, String question) {
DiagnosisSession session = diagnosisSessionRepository.findBySessionId(sessionId)
.orElseGet(() -> DiagnosisSession.builder()
.sessionId(sessionId)
.agentFlow("CHAT")
.build());
session.setQuery(question);
session.setStatus("RUNNING");
session.setAgentFlow("CHAT");
session.setAnswer(null);
session.setTotalDurationMs(null);
session.setTotalTokenCount(null);
session.setStepCount(null);
session.setToolCallCount(null);
return session;
}
private String buildExecutorInput(String question, String plannerPlan, String retryContext) {
StringBuilder input = new StringBuilder(question);
if (plannerPlan != null && !plannerPlan.isBlank()) {
input.append("\n\n--- planner_plan ---\n").append(plannerPlan);
}
private String buildWorkflowInput(String question, String retryContext) {
StringBuilder input = new StringBuilder();
input.append("请按固定工作流完成本轮 Planner -> Executor -> Verifier。\n\n");
input.append("--- 用户问题 ---\n").append(question);
if (retryContext != null && !retryContext.isBlank()) {
input.append("\n\n--- retry_context ---\n").append(retryContext);
}
input.append("\n\nVerifier 完成后由外层代码读取 verifier_output 并决定最终用户输出。");
return input.toString();
}
private String invokeVerifierFallback(ReactAgent verifier, String question, int round, RunnableConfig config) {
try {
logger.warn("Sequential workflow round {} finished without verifier_output, invoking chat_verifier fallback", round);
return verifier.call("请基于 executor_final_answer 和 tool_trace_summary 输出 verifier JSON。原始问题:" + question, config)
.getText();
} catch (Exception e) {
logger.error("chat_verifier fallback 执行失败", e);
return null;
}
}
private VerifierDecision parseVerifierDecision(String verifierOutput, int round) {
if (verifierOutput == null || verifierOutput.isBlank()) {
return null;
@@ -629,70 +683,20 @@ public class ChatService {
return new VerifierDecision("LOW_CONFID", 0.0, 0, List.of(), rationale, round);
}
private String buildSupervisorPrompt(int round) {
return """
你是一个多 Agent 调度器。每一轮必须严格按顺序完成以下动作:
1. 先调用 chat_planner 生成执行计划
2. 再调用 chat_executor 执行计划并形成最终答案
3. 最后调用 chat_verifier 对 executor 最终答案做事实核查
规则:
- 本轮只允许完成一次 Planner -> Executor -> Verifier 链路
- Verifier 完成后立即停止,不要继续调用任何 Agent
- 不要自己编造答案,最终用户输出由外层代码根据 verifier_output 决定
- 当前是第 %d 轮,保持单轮内顺序稳定
""".formatted(round);
}
private String buildRoundInput(String question, String retryContext) {
if (retryContext == null || retryContext.isBlank()) {
return question;
}
return question + "\n\n--- 补充约束 ---\n" + retryContext;
}
private String extractExecutorAnswer(Optional<OverAllState> stateOptional) {
private String extractStateText(Optional<OverAllState> stateOptional, String key) {
if (stateOptional.isEmpty()) {
return null;
}
return stateOptional.get().value("executor_feedback")
.filter(AssistantMessage.class::isInstance)
.map(AssistantMessage.class::cast)
.map(AssistantMessage::getText)
return stateOptional.get().value(key)
.map(value -> {
if (value instanceof AssistantMessage assistantMessage) {
return assistantMessage.getText();
}
return String.valueOf(value);
})
.orElse(null);
}
private VerifierDecision parseVerifierDecision(Optional<OverAllState> stateOptional, int round) {
if (stateOptional.isEmpty()) {
return null;
}
Optional<AssistantMessage> verifierOutput = stateOptional.get().value("verifier_output")
.filter(AssistantMessage.class::isInstance)
.map(AssistantMessage.class::cast);
if (verifierOutput.isEmpty() || verifierOutput.get().getText() == null || verifierOutput.get().getText().isBlank()) {
return null;
}
try {
JsonNode root = objectMapper.readTree(verifierOutput.get().getText());
List<Map<String, Object>> factsChecked = parseFactsChecked(root.path("facts_checked"));
return new VerifierDecision(
root.path("verdict").asText("LOW_CONFID"),
root.path("groundedness_score").asDouble(0.0),
root.path("critical_fact_count").asInt(0),
factsChecked,
root.path("rationale").asText(""),
round
);
} catch (Exception e) {
logger.error("解析 verifier_output 失败: {}", verifierOutput.get().getText(), e);
return null;
}
}
private void persistVerifierEvaluation(DiagnosisSession session, VerifierDecision decision, int round) {
if (decision == null) {
return;
@@ -0,0 +1,136 @@
package com.superbiz.agent.service;
import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.superbiz.agent.domain.entity.AgentStep;
import com.superbiz.agent.domain.entity.DiagnosisSession;
import com.superbiz.agent.domain.entity.ToolInvocation;
import com.superbiz.agent.dto.DiagnosisTraceResponse;
import com.superbiz.agent.exception.SessionNotFoundException;
import com.superbiz.agent.repository.AgentStepRepository;
import com.superbiz.agent.repository.DiagnosisSessionRepository;
import com.superbiz.agent.repository.ToolInvocationRepository;
import lombok.RequiredArgsConstructor;
import org.springframework.stereotype.Service;
import java.util.List;
import java.util.Map;
@Service
@RequiredArgsConstructor
public class DiagnosisTraceService {
private static final TypeReference<Map<String, Object>> JSON_MAP_TYPE = new TypeReference<>() {
};
private final DiagnosisSessionRepository diagnosisSessionRepository;
private final AgentStepRepository agentStepRepository;
private final ToolInvocationRepository toolInvocationRepository;
private final ObjectMapper objectMapper;
public DiagnosisTraceResponse getTrace(String sessionId) {
DiagnosisSession session = diagnosisSessionRepository.findBySessionId(sessionId)
.orElseThrow(() -> new SessionNotFoundException(sessionId));
List<AgentStep> steps = agentStepRepository.findBySessionIdOrderByStepIndex(sessionId);
List<ToolInvocation> toolInvocations = toolInvocationRepository.findBySessionIdOrderByIdAsc(sessionId);
return DiagnosisTraceResponse.builder()
.session(toSessionTrace(session))
.steps(steps.stream().map(this::toAgentStepTrace).toList())
.toolInvocations(toolInvocations.stream().map(this::toToolInvocationTrace).toList())
.summary(toSummary(session, steps, toolInvocations))
.build();
}
private DiagnosisTraceResponse.SessionTrace toSessionTrace(DiagnosisSession session) {
return DiagnosisTraceResponse.SessionTrace.builder()
.id(session.getId())
.sessionId(session.getSessionId())
.query(session.getQuery())
.status(session.getStatus())
.agentFlow(session.getAgentFlow())
.totalDurationMs(session.getTotalDurationMs())
.totalTokenCount(session.getTotalTokenCount())
.stepCount(session.getStepCount())
.toolCallCount(session.getToolCallCount())
.answer(session.getAnswer())
.selfEvaluationRaw(session.getSelfEvaluation())
.selfEvaluation(parseJsonObject(session.getSelfEvaluation()))
.feedback(session.getFeedback())
.createdAt(session.getCreatedAt())
.updatedAt(session.getUpdatedAt())
.build();
}
private DiagnosisTraceResponse.AgentStepTrace toAgentStepTrace(AgentStep step) {
return DiagnosisTraceResponse.AgentStepTrace.builder()
.id(step.getId())
.sessionId(step.getSessionId())
.stepIndex(step.getStepIndex())
.agentName(step.getAgentName())
.modelInput(step.getModelInput())
.modelOutput(step.getModelOutput())
.thought(step.getThought())
.hasToolCall(step.getHasToolCall())
.durationMs(step.getDurationMs())
.tokenCount(step.getTokenCount())
.createdAt(step.getCreatedAt())
.build();
}
private DiagnosisTraceResponse.ToolInvocationTrace toToolInvocationTrace(ToolInvocation invocation) {
return DiagnosisTraceResponse.ToolInvocationTrace.builder()
.id(invocation.getId())
.sessionId(invocation.getSessionId())
.stepId(invocation.getStepId())
.toolName(invocation.getToolName())
.inputParamsRaw(invocation.getInputParams())
.inputParams(parseJsonObject(invocation.getInputParams()))
.outputPreview(invocation.getOutputPreview())
.outputLength(invocation.getOutputLength())
.retrievalLayer(invocation.getRetrievalLayer())
.l0MatchCount(invocation.getL0MatchCount())
.l1MatchCount(invocation.getL1MatchCount())
.truncated(invocation.getIsTruncated())
.relevanceLevel(invocation.getRelevanceLevel())
.dedupReason(invocation.getDedupReason())
.retrievalDetailsRaw(invocation.getRetrievalDetails())
.retrievalDetails(parseJsonObject(invocation.getRetrievalDetails()))
.durationMs(invocation.getDurationMs())
.success(invocation.getSuccess())
.errorMessage(invocation.getErrorMessage())
.createdAt(invocation.getCreatedAt())
.build();
}
private DiagnosisTraceResponse.TraceSummary toSummary(
DiagnosisSession session,
List<AgentStep> steps,
List<ToolInvocation> toolInvocations
) {
Map<String, Object> selfEvaluation = parseJsonObject(session.getSelfEvaluation());
return DiagnosisTraceResponse.TraceSummary.builder()
.persistedStepCount(defaultInt(session.getStepCount()))
.returnedStepCount(steps.size())
.persistedToolCallCount(defaultInt(session.getToolCallCount()))
.returnedToolCallCount(toolInvocations.size())
.hasVerifierEvaluation(selfEvaluation != null && selfEvaluation.containsKey("verifier_evaluation"))
.hasFeedback(session.getFeedback() != null && !session.getFeedback().isBlank())
.build();
}
private Map<String, Object> parseJsonObject(String json) {
if (json == null || json.isBlank()) {
return null;
}
try {
return objectMapper.readValue(json, JSON_MAP_TYPE);
} catch (Exception ignored) {
return null;
}
}
private int defaultInt(Integer value) {
return value == null ? 0 : value;
}
}
@@ -260,7 +260,8 @@ public class DocumentManagementService {
private String saveToLocal(MultipartFile file, String fileName, String category) {
try {
// 1. 构建目标路径
Path categoryDir = Paths.get(knowledgeBasePath, category);
Path baseDir = Paths.get(knowledgeBasePath).normalize();
Path categoryDir = baseDir.resolve(category).normalize();
Files.createDirectories(categoryDir);
Path targetPath = categoryDir.resolve(fileName);
@@ -268,8 +269,9 @@ public class DocumentManagementService {
// 2. 保存文件
file.transferTo(targetPath.toFile());
log.info("文件已保存到本地: {}", targetPath);
return targetPath.toString();
String relativePath = baseDir.relativize(targetPath.normalize()).toString().replace("\\", "/");
log.info("文件已保存到本地: {}, storedPath={}", targetPath, relativePath);
return relativePath;
} catch (IOException e) {
throw new DocumentProcessException(
@@ -287,7 +289,7 @@ public class DocumentManagementService {
private void cleanupLocalFile(String localPath) {
if (localPath != null) {
try {
Files.deleteIfExists(Paths.get(localPath));
Files.deleteIfExists(resolveLocalPath(localPath));
log.info("已清理本地文件: {}", localPath);
} catch (IOException e) {
log.warn("清理本地文件失败: {}", localPath, e);
@@ -357,7 +359,7 @@ public class DocumentManagementService {
// 删除本地文件
if (doc.getFilePath() != null) {
try {
Files.deleteIfExists(Paths.get(doc.getFilePath()));
Files.deleteIfExists(resolveLocalPath(doc.getFilePath()));
log.info("本地文件已删除: {}", doc.getFilePath());
} catch (IOException e) {
log.warn("删除本地文件失败: {}", doc.getFilePath(), e);
@@ -407,6 +409,26 @@ public class DocumentManagementService {
return null;
}
private Path resolveLocalPath(String filePath) {
Path path = Paths.get(filePath).normalize();
if (path.isAbsolute()) {
return path;
}
Path basePath = Paths.get(knowledgeBasePath).toAbsolutePath().normalize();
Path baseName = basePath.getFileName();
if (baseName != null && path.startsWith(baseName) && basePath.getParent() != null) {
return basePath.getParent().resolve(path).normalize();
}
Path pathFromWorkingDir = path.toAbsolutePath().normalize();
if (pathFromWorkingDir.startsWith(basePath)) {
return pathFromWorkingDir;
}
return basePath.resolve(path).normalize();
}
/**
* 转换为响应 DTO
*/
@@ -160,7 +160,12 @@ public class KnowledgeIndexService {
public String readDocument(String filePath, int maxChars) {
try {
Path fullPath = Paths.get(knowledgeBasePath, filePath);
Path fullPath = resolveDocumentPath(filePath);
if (!Files.exists(fullPath)) {
log.warn("读取文档失败,文件不存在: basePath={}, filePath={}, resolvedPath={}",
knowledgeBasePath, filePath, fullPath);
return null;
}
String content = Files.readString(fullPath);
if (content.length() > maxChars) {
@@ -170,11 +175,35 @@ public class KnowledgeIndexService {
return content;
} catch (IOException e) {
log.error("读取文档失败: {}/{}", knowledgeBasePath, filePath, e);
log.error("读取文档失败: basePath={}, filePath={}", knowledgeBasePath, filePath, e);
return null;
}
}
Path resolveDocumentPath(String filePath) {
if (filePath == null || filePath.isBlank()) {
throw new IllegalArgumentException("filePath cannot be blank");
}
Path path = Paths.get(filePath).normalize();
if (path.isAbsolute()) {
return path;
}
Path basePath = Paths.get(knowledgeBasePath).toAbsolutePath().normalize();
Path baseName = basePath.getFileName();
if (baseName != null && path.startsWith(baseName) && basePath.getParent() != null) {
return basePath.getParent().resolve(path).normalize();
}
Path pathFromWorkingDir = path.toAbsolutePath().normalize();
if (pathFromWorkingDir.startsWith(basePath)) {
return pathFromWorkingDir;
}
return basePath.resolve(path).normalize();
}
public void addToIndex(KnowledgeEntry entry) {
knowledgeIndex.add(entry);
log.debug("文档已添加到 L0 索引: title={}", entry.getTitle());
@@ -0,0 +1,93 @@
package com.superbiz.agent.service;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.superbiz.agent.domain.entity.ToolInvocation;
import com.superbiz.agent.repository.ToolInvocationRepository;
import com.superbiz.agent.util.SessionContextHolder;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.UUID;
/**
* Central persistence point for agent evidence tool invocations.
*/
@Slf4j
@Service
public class ToolInvocationRecorder {
private static final int OUTPUT_PREVIEW_LIMIT = 500;
private final ToolInvocationRepository toolInvocationRepository;
private final ObjectMapper objectMapper;
public ToolInvocationRecorder(ToolInvocationRepository toolInvocationRepository, ObjectMapper objectMapper) {
this.toolInvocationRepository = toolInvocationRepository;
this.objectMapper = objectMapper;
}
public void save(ToolInvocation invocation) {
try {
if (invocation.getSessionId() == null || invocation.getSessionId().isBlank()) {
invocation.setSessionId(SessionContextHolder.getSessionId());
}
if (invocation.getSessionId() == null || invocation.getSessionId().isBlank()) {
log.debug("Skip tool_invocation without sessionId: tool={}", invocation.getToolName());
return;
}
toolInvocationRepository.save(invocation);
} catch (Exception e) {
log.error("保存 tool_invocation 失败: tool={}", invocation.getToolName(), e);
}
}
public void recordEvidenceTool(String toolName,
Map<String, Object> inputParams,
String output,
boolean success,
long startTimeMillis,
String errorMessage,
String topicDomain) {
String outputPreview = preview(output);
Map<String, Object> details = new LinkedHashMap<>();
details.put("trace_id", UUID.randomUUID().toString());
if (topicDomain != null && !topicDomain.isBlank()) {
details.put("retrieved_domains", List.of(topicDomain));
}
ToolInvocation invocation = ToolInvocation.builder()
.toolName(toolName)
.inputParams(toJson(inputParams == null ? Map.of() : inputParams))
.outputPreview(outputPreview)
.outputLength(output == null ? 0 : output.length())
.isTruncated(output != null && output.length() > OUTPUT_PREVIEW_LIMIT)
.retrievalDetails(toJson(details))
.durationMs((int) Math.max(0, System.currentTimeMillis() - startTimeMillis))
.success(success)
.errorMessage(errorMessage)
.build();
save(invocation);
}
private String preview(String output) {
if (output == null) {
return null;
}
return output.length() <= OUTPUT_PREVIEW_LIMIT
? output
: output.substring(0, OUTPUT_PREVIEW_LIMIT) + "...";
}
private String toJson(Map<String, Object> value) {
try {
return objectMapper.writeValueAsString(value);
} catch (JsonProcessingException e) {
log.debug("tool_invocation JSON 序列化失败", e);
return "{}";
}
}
}
@@ -3,8 +3,8 @@ package com.superbiz.agent.tool;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.superbiz.agent.domain.entity.ToolInvocation;
import com.superbiz.agent.dto.*;
import com.superbiz.agent.repository.ToolInvocationRepository;
import com.superbiz.agent.service.KnowledgeIndexService;
import com.superbiz.agent.service.ToolInvocationRecorder;
import com.superbiz.agent.service.VectorSearchService;
import com.superbiz.agent.util.SessionContextHolder;
import lombok.extern.slf4j.Slf4j;
@@ -49,7 +49,7 @@ public class LookupKnowledgeTool {
private VectorSearchService vectorSearchService;
@Autowired
private ToolInvocationRepository toolInvocationRepository;
private ToolInvocationRecorder toolInvocationRecorder;
@Autowired
private RetrievedDocTracker retrievedDocTracker;
@@ -417,7 +417,7 @@ public class LookupKnowledgeTool {
.success(true)
.build();
toolInvocationRepository.save(inv);
toolInvocationRecorder.save(inv);
log.debug("tool_invocation 已保存: sessionId={}, layer={}, relevanceLevel={}, duration={}ms",
sessionId, layer, result != null ? result.getRelevanceLevel() : null, duration);
} catch (Exception e) {
@@ -0,0 +1,25 @@
spring:
config:
activate:
on-profile: mvp-demo
server:
port: 9900
prometheus:
mock-enabled: true
timeout: 5
cls:
mock-enabled: true
logging:
level:
root: INFO
com.superbiz.agent: DEBUG
com.alibaba.cloud: INFO
mvp:
demo:
name: payment-timeout-trace
description: Repeatable MVP flow for chat diagnosis, tool evidence, verifier evaluation, trace query, and feedback.
@@ -0,0 +1,223 @@
package com.superbiz.agent.service;
import com.superbiz.agent.agent.tool.DateTimeTools;
import com.superbiz.agent.agent.tool.QueryLogsTools;
import com.superbiz.agent.agent.tool.QueryMetricsTools;
import com.superbiz.agent.domain.entity.AgentStep;
import com.superbiz.agent.domain.entity.DiagnosisSession;
import com.superbiz.agent.repository.AgentStepRepository;
import com.superbiz.agent.repository.DiagnosisSessionRepository;
import com.superbiz.agent.tool.LookupKnowledgeTool;
import com.superbiz.agent.tool.RetrievedDocTracker;
import org.junit.jupiter.api.Test;
import org.springframework.ai.chat.messages.AssistantMessage;
import org.springframework.ai.chat.model.ChatModel;
import org.springframework.ai.chat.model.ChatResponse;
import org.springframework.ai.chat.model.Generation;
import org.springframework.ai.chat.prompt.Prompt;
import org.springframework.ai.tool.ToolCallback;
import org.springframework.test.util.ReflectionTestUtils;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.concurrent.atomic.AtomicInteger;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertSame;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;
class ChatServiceSequentialAgentTest {
@Test
void executeChatComplexInvokesSequentialWorkflow() throws Exception {
ChatService chatService = createChatService();
ScriptedChatModel chatModel = new ScriptedChatModel();
ChatService.ChatResult result = chatService.executeChatComplex(
chatModel,
new ToolCallback[0],
"请分析订单支付超时的原因,并给出修复建议",
List.of(),
"sequential-test-session"
);
assertEquals("EXECUTOR_FINAL_ANSWER", result.answer());
assertEquals("sequential-test-session", result.sessionId());
assertEquals(List.of("chat_planner", "chat_executor", "chat_verifier"), chatModel.agentCalls);
assertTrue(chatModel.sawVerifierPrompt);
}
@Test
void executeChatComplexDoesNotRetryLowConfidenceByDefault() throws Exception {
ChatService chatService = createChatService();
ScriptedChatModel chatModel = new ScriptedChatModel("""
{
"verdict": "LOW_CONFID",
"groundedness_score": 0.1,
"critical_fact_count": 1,
"facts_checked": [
{
"fact": "missing direct evidence",
"is_critical": true,
"verification": "no_evidence",
"detail": "scripted evidence gap",
"evidence_refs": []
}
],
"rationale": "scripted low confidence"
}
""");
ChatService.ChatResult result = chatService.executeChatComplex(
chatModel,
new ToolCallback[0],
"请分析订单支付超时的原因,并给出修复建议",
List.of(),
"sequential-low-confidence-session"
);
assertTrue(result.answer().contains("EXECUTOR_FINAL_ANSWER"));
assertEquals(List.of("chat_planner", "chat_executor", "chat_verifier"), chatModel.agentCalls);
}
@Test
void executeChatComplexRunsPlannerExecutorVerifierInFixedOrder() throws Exception {
ChatService chatService = createChatService();
ScriptedChatModel chatModel = new ScriptedChatModel();
ChatService.ChatResult result = chatService.executeChatComplex(
chatModel,
new ToolCallback[0],
"请分析订单支付超时的原因,并给出修复建议",
List.of(),
"sequential-workflow-session"
);
assertEquals("EXECUTOR_FINAL_ANSWER", result.answer());
assertEquals(List.of("chat_planner", "chat_executor", "chat_verifier"), chatModel.agentCalls);
assertTrue(chatModel.sawVerifierPrompt);
}
@Test
void buildMethodToolsArrayIncludesLogsAndMetricsWhenAvailable() {
ChatService chatService = new ChatService();
DateTimeTools dateTimeTools = new DateTimeTools();
LookupKnowledgeTool lookupKnowledgeTool = new LookupKnowledgeTool();
QueryLogsTools queryLogsTools = new QueryLogsTools(mock(ToolInvocationRecorder.class));
QueryMetricsTools queryMetricsTools = new QueryMetricsTools(mock(ToolInvocationRecorder.class));
ReflectionTestUtils.setField(chatService, "dateTimeTools", dateTimeTools);
ReflectionTestUtils.setField(chatService, "lookupKnowledgeTool", lookupKnowledgeTool);
ReflectionTestUtils.setField(chatService, "queryLogsTools", queryLogsTools);
ReflectionTestUtils.setField(chatService, "queryMetricsTools", queryMetricsTools);
Object[] methodTools = chatService.buildMethodToolsArray();
assertEquals(4, methodTools.length);
assertSame(dateTimeTools, methodTools[0]);
assertSame(lookupKnowledgeTool, methodTools[1]);
assertSame(queryLogsTools, methodTools[2]);
assertSame(queryMetricsTools, methodTools[3]);
}
private ChatService createChatService() {
ChatService chatService = new ChatService();
DiagnosisSessionRepository diagnosisSessionRepository = mock(DiagnosisSessionRepository.class);
when(diagnosisSessionRepository.findBySessionId(anyString())).thenReturn(Optional.empty());
when(diagnosisSessionRepository.save(any(DiagnosisSession.class))).thenAnswer(invocation -> invocation.getArgument(0));
AtomicInteger stepId = new AtomicInteger(1);
AgentStepRepository agentStepRepository = mock(AgentStepRepository.class);
when(agentStepRepository.save(any(AgentStep.class))).thenAnswer(invocation -> {
AgentStep step = invocation.getArgument(0);
if (step.getId() == null) {
step.setId((long) stepId.getAndIncrement());
}
return step;
});
when(agentStepRepository.findById(any())).thenReturn(Optional.of(new AgentStep()));
when(agentStepRepository.findBySessionIdOrderByStepIndex(anyString())).thenReturn(List.of());
EvaluationService evaluationService = mock(EvaluationService.class);
RetrievedDocTracker retrievedDocTracker = mock(RetrievedDocTracker.class);
KnowledgeDomainService knowledgeDomainService = mock(KnowledgeDomainService.class);
when(knowledgeDomainService.buildKnowledgeMap()).thenReturn("");
ToolTraceSummaryService toolTraceSummaryService = mock(ToolTraceSummaryService.class);
when(toolTraceSummaryService.buildVerifierTraceSummary(anyString(), anyString())).thenReturn(List.of());
SelfEvaluationMergeService selfEvaluationMergeService = mock(SelfEvaluationMergeService.class);
when(selfEvaluationMergeService.mergeVerifierEvaluation(any(), any())).thenReturn("{}");
ReflectionTestUtils.setField(chatService, "dateTimeTools", new DateTimeTools());
ReflectionTestUtils.setField(chatService, "lookupKnowledgeTool", new LookupKnowledgeTool());
ReflectionTestUtils.setField(chatService, "queryLogsTools", new QueryLogsTools(mock(ToolInvocationRecorder.class)));
ReflectionTestUtils.setField(chatService, "diagnosisSessionRepository", diagnosisSessionRepository);
ReflectionTestUtils.setField(chatService, "agentStepRepository", agentStepRepository);
ReflectionTestUtils.setField(chatService, "evaluationService", evaluationService);
ReflectionTestUtils.setField(chatService, "retrievedDocTracker", retrievedDocTracker);
ReflectionTestUtils.setField(chatService, "knowledgeDomainService", knowledgeDomainService);
ReflectionTestUtils.setField(chatService, "toolTraceSummaryService", toolTraceSummaryService);
ReflectionTestUtils.setField(chatService, "selfEvaluationMergeService", selfEvaluationMergeService);
ReflectionTestUtils.setField(chatService, "verifierLowConfidenceThreshold", 0.5d);
ReflectionTestUtils.setField(chatService, "chatPlannerPrompt", "PLANNER_TEST_PROMPT");
ReflectionTestUtils.setField(chatService, "chatExecutorPrompt", "EXECUTOR_TEST_PROMPT");
ReflectionTestUtils.setField(chatService, "chatVerifierPrompt", "VERIFIER_TEST_PROMPT");
return chatService;
}
private static final class ScriptedChatModel implements ChatModel {
private final java.util.ArrayList<String> agentCalls = new java.util.ArrayList<>();
private String promptText = "";
private boolean sawVerifierPrompt;
private final String verifierOutput;
private ScriptedChatModel() {
this("""
{
"verdict": "PASS",
"groundedness_score": 1.0,
"critical_fact_count": 1,
"facts_checked": [
{
"fact": "executor answer generated",
"is_critical": true,
"verification": "direct_evidence",
"detail": "covered by scripted verifier",
"evidence_refs": []
}
],
"rationale": "scripted pass"
}
""");
}
private ScriptedChatModel(String verifierOutput) {
this.verifierOutput = verifierOutput;
}
@Override
public ChatResponse call(Prompt prompt) {
promptText = prompt.getContents();
String text;
if (promptText.contains("PLANNER_TEST_PROMPT")) {
agentCalls.add("chat_planner");
text = "PLANNER_PLAN";
} else if (promptText.contains("EXECUTOR_TEST_PROMPT")) {
agentCalls.add("chat_executor");
text = "EXECUTOR_FINAL_ANSWER";
} else if (promptText.contains("VERIFIER_TEST_PROMPT")) {
agentCalls.add("chat_verifier");
sawVerifierPrompt = true;
text = verifierOutput;
} else {
text = "UNEXPECTED_PROMPT";
}
return new ChatResponse(List.of(new Generation(new AssistantMessage(text))));
}
}
}
@@ -0,0 +1,119 @@
package com.superbiz.agent.service;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.superbiz.agent.domain.entity.AgentStep;
import com.superbiz.agent.domain.entity.DiagnosisSession;
import com.superbiz.agent.domain.entity.ToolInvocation;
import com.superbiz.agent.dto.DiagnosisTraceResponse;
import com.superbiz.agent.exception.SessionNotFoundException;
import com.superbiz.agent.repository.AgentStepRepository;
import com.superbiz.agent.repository.DiagnosisSessionRepository;
import com.superbiz.agent.repository.ToolInvocationRepository;
import org.junit.jupiter.api.Test;
import java.time.LocalDateTime;
import java.util.List;
import java.util.Optional;
import static org.junit.jupiter.api.Assertions.*;
import static org.mockito.Mockito.*;
class DiagnosisTraceServiceTest {
private final DiagnosisSessionRepository diagnosisSessionRepository = mock(DiagnosisSessionRepository.class);
private final AgentStepRepository agentStepRepository = mock(AgentStepRepository.class);
private final ToolInvocationRepository toolInvocationRepository = mock(ToolInvocationRepository.class);
private final DiagnosisTraceService service = new DiagnosisTraceService(
diagnosisSessionRepository,
agentStepRepository,
toolInvocationRepository,
new ObjectMapper()
);
@Test
void getTraceAggregatesSessionStepsAndTools() {
String sessionId = "trace-session-001";
LocalDateTime now = LocalDateTime.of(2026, 7, 3, 14, 30);
DiagnosisSession session = DiagnosisSession.builder()
.id(1L)
.sessionId(sessionId)
.query("payment timeout")
.status("SUCCESS")
.agentFlow("COMPLEX")
.totalDurationMs(1200)
.totalTokenCount(300)
.stepCount(2)
.toolCallCount(1)
.answer("restart payment gateway pool")
.selfEvaluation("{\"verifier_evaluation\":{\"verdict\":\"PASS\"}}")
.feedback("useful")
.createdAt(now)
.updatedAt(now)
.build();
AgentStep step = AgentStep.builder()
.id(10L)
.sessionId(sessionId)
.stepIndex(1)
.agentName("chat_executor")
.modelInput("input")
.modelOutput("output")
.thought("executor finished")
.hasToolCall(true)
.durationMs(500)
.tokenCount(100)
.createdAt(now)
.build();
ToolInvocation invocation = ToolInvocation.builder()
.id(20L)
.sessionId(sessionId)
.stepId(10L)
.toolName("lookup_knowledge")
.inputParams("{\"query\":\"ERR_TIMEOUT\"}")
.outputPreview("payment timeout doc")
.outputLength(19)
.retrievalLayer("L0")
.l0MatchCount(1)
.l1MatchCount(0)
.isTruncated(false)
.relevanceLevel("HIGHLY_RELEVANT")
.dedupReason("FIRST_HIT")
.retrievalDetails("{\"documents\":[\"payment-errors.md\"]}")
.durationMs(80)
.success(true)
.createdAt(now)
.build();
when(diagnosisSessionRepository.findBySessionId(sessionId)).thenReturn(Optional.of(session));
when(agentStepRepository.findBySessionIdOrderByStepIndex(sessionId)).thenReturn(List.of(step));
when(toolInvocationRepository.findBySessionIdOrderByIdAsc(sessionId)).thenReturn(List.of(invocation));
DiagnosisTraceResponse response = service.getTrace(sessionId);
assertEquals(sessionId, response.getSession().getSessionId());
assertEquals("payment timeout", response.getSession().getQuery());
assertEquals("PASS", ((java.util.Map<?, ?>) response.getSession()
.getSelfEvaluation()
.get("verifier_evaluation")).get("verdict"));
assertEquals(1, response.getSteps().size());
assertEquals("chat_executor", response.getSteps().get(0).getAgentName());
assertEquals(1, response.getToolInvocations().size());
assertEquals("ERR_TIMEOUT", response.getToolInvocations().get(0).getInputParams().get("query"));
assertEquals(2, response.getSummary().getPersistedStepCount());
assertEquals(1, response.getSummary().getReturnedStepCount());
assertEquals(1, response.getSummary().getPersistedToolCallCount());
assertEquals(1, response.getSummary().getReturnedToolCallCount());
assertTrue(response.getSummary().isHasVerifierEvaluation());
assertTrue(response.getSummary().isHasFeedback());
}
@Test
void getTraceThrowsWhenSessionMissing() {
String sessionId = "missing-session";
when(diagnosisSessionRepository.findBySessionId(sessionId)).thenReturn(Optional.empty());
assertThrows(SessionNotFoundException.class, () -> service.getTrace(sessionId));
verify(diagnosisSessionRepository).findBySessionId(sessionId);
verifyNoInteractions(agentStepRepository, toolInvocationRepository);
}
}
@@ -0,0 +1,41 @@
package com.superbiz.agent.service;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import org.springframework.mock.web.MockMultipartFile;
import org.springframework.test.util.ReflectionTestUtils;
import java.nio.file.Files;
import java.nio.file.Path;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
class DocumentManagementServiceTest {
@TempDir
Path tempDir;
@Test
void saveToLocalStoresRelativePathUnderKnowledgeBase() throws Exception {
DocumentManagementService service = new DocumentManagementService();
ReflectionTestUtils.setField(service, "knowledgeBasePath", tempDir.toString());
MockMultipartFile file = new MockMultipartFile(
"file",
"runbook.md",
"text/markdown",
"runbook content".getBytes()
);
String storedPath = ReflectionTestUtils.invokeMethod(
service,
"saveToLocal",
file,
"runbook.md",
"payment"
);
assertEquals("payment/runbook.md", storedPath);
assertTrue(Files.exists(tempDir.resolve("payment").resolve("runbook.md")));
}
}
@@ -4,8 +4,6 @@ import com.superbiz.agent.dto.KnowledgeEntry;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import org.mockito.Mock;
import org.mockito.MockitoAnnotations;
import org.springframework.test.util.ReflectionTestUtils;
import java.nio.file.Files;
@@ -21,17 +19,13 @@ class KnowledgeIndexServiceTest {
private KnowledgeIndexService service;
@Mock
private FrontmatterParser frontmatterParser;
@TempDir
Path tempDir;
@BeforeEach
void setUp() {
MockitoAnnotations.openMocks(this);
service = new KnowledgeIndexService();
ReflectionTestUtils.setField(service, "frontmatterParser", frontmatterParser);
ReflectionTestUtils.setField(service, "knowledgeBasePath", tempDir.toString());
}
@Test
@@ -144,6 +138,30 @@ class KnowledgeIndexServiceTest {
assertTrue(result.contains("Test content"));
}
@Test
void testReadDocument_relativePathUnderBasePath() throws Exception {
Path categoryDir = tempDir.resolve("payment");
Files.createDirectories(categoryDir);
Path testFile = categoryDir.resolve("relative.md");
Files.writeString(testFile, "Relative content");
String result = service.readDocument("payment/relative.md", 100);
assertEquals("Relative content", result);
}
@Test
void testReadDocument_legacyPathAlreadyContainsBasePath() throws Exception {
Path categoryDir = tempDir.resolve("payment");
Files.createDirectories(categoryDir);
Path testFile = categoryDir.resolve("legacy.md");
Files.writeString(testFile, "Legacy content");
String result = service.readDocument(tempDir.getFileName() + "/payment/legacy.md", 100);
assertEquals("Legacy content", result);
}
@Test
void testReadDocument_exceedsMaxChars() throws Exception {
// 创建超长内容