diff --git a/muse-cloud/muse-module-ai/.agent b/muse-cloud/muse-module-ai/.agent index db5b46c9..c9091cc0 100644 --- a/muse-cloud/muse-module-ai/.agent +++ b/muse-cloud/muse-module-ai/.agent @@ -8,5 +8,6 @@ - **现状**:只读评估 80%;后端真实较高质量(无 stub),`RealNewApiMuseAiRuntimeClient` 真打 New-API,SSE 已脱占位;studio `/agents` 已补自建智能体生命周期活体证(2026-06-27:`agent-create.spec.ts` MSW-off 真单体 1/1,覆盖创建初始版本、编辑生成 v2、版本切换、归档非当前版本、归档 Agent,DB 证 `current_version_id` 与版本 status)。 - **E4 market agent 物化(2026-06-27)**:market agent handoff 不再信任客户端 `sourceAgentId/sourceAgentVersion`;AI 目标 owner 按 market asset `source_id` 解析发布者 agent,并物化为安装者本地 user agent,槽位绑定指向本地副本。MSW-off Playwright `handoff-agent.spec.ts` 2/0F/0E,正路物化后真实创建 AI task 并完成 New-API runtime;负路伪造 token 被拒。注意:e2e 前置种入 Security 已批准 runtime grant,生产 runtime 不自授权。 - **E5 account 用量/配额接线(2026-06-27)**:AI task 创建前经 Account owner 预占配额,成功投影时记录 usage 并累加 quota used,非可重试失败终态释放预占;AI 只消费 member-api 端口,不碰 member DAL。`P1rAiRuntimeEndToEndLiveAcceptanceIT` 显式 live 3/0F/0E/0S,反查 Account usage record=1、quota used=1、reserve audit=1、release audit=0、attribution job/item completed。targeted AI 单测中 `MuseAiTaskServiceTest` 40/0F/0E、`MuseAiRuntimeProjectionServiceTest` 6/0F/0E。 -- **E6 admin AI 配置接线(2026-06-27)**:Muse Admin AI 配置页 Prompt 激活、质量评估启动、质量策略版本创建、system Agent 创建/版本发布/归档停用、Tool Grant 登记/调整已调用真实 admin API,不再展示硬编码评估分数。后端 `adminArchiveAgent` 只允许 system scope,归档 Agent 时同步归档 active 槽位绑定,避免运行时继续消费已停用系统能力;质量策略回滚仍保持 preview-only/线 B。验证:`muse-admin` AI 配置目标 Vitest 2 files / 9 tests passed,`@vben/web-antd typecheck` 通过;后端 `MuseAgentServiceTest`+`AdminMuseAgentControllerAnnotationTest` 31/0F/0E/0S,`MuseToolGrantServiceTest`+`LocalSecurityToolGrantApprovalFacadeTest` 13/0F/0E/0S,server 契约/覆盖门 27/0F/0E/0S。 -- **关键风险 / TODO**:**AI runtime 生成 e2e 已端到端真验(2026-06-24)**——方案1(d647f0c 剥离 sourceSnapshotId 解授权死结)+ D(4dadc7b `DefaultDBFieldHandler` 无登录兜底 `system`,根治 scheduled worker 写 creator/updater null 违约 NOT NULL;BaseDO 带 jdbcType 强制入列致 insert/update 同病)+ agent_version config 补 modelKey=MiniMax-M2.5 + 启动 source `scripts/dev/p1r-external-acceptance.env`(New-API/RAGFlow 配置)→ task5 completed+suggestion+MiniMax-M2.5 真 LLM 生成(P-B envelope/approved、P-C contextAssembly 固化 ✓)。**P-A/P-B/P-C 完整 e2e 已端到端真验打通(2026-06-24,task7)**:P-A 检索命中根因=knowledge 模块 bug——`HttpRagFlowKnowledgeRuntimeClient.retrieveChunks` 把 `document_ids` 以 null 入 retrieval body,RAGFlow 要求其为 list 故拒(code 102 "documents should be a list")→检索 fail-closed→chunkCount=0(task3/4/5/6 全 0 即此 bug,非 RAGFlow 数据/凭据/timeout)。已修(commit 86c6dba:null 时不入 body + 回归测试)。task7 contextAssembly `chunkCount=1/retrievedKbIds=[2]/authorizationSnapshotIds=[authsnap-e2e-1]`,LLM(MiniMax-M2.5)基于 grounding 续写。**02D 槽位首次创建死锁已修+真验(2026-06-24,commit 9115d4f)**:根因=precheck/bind 都 `requireSlotBinding`(缺行抛 forbidden),而"替换"是设计中唯一创建入口("替换"=首次写入、无独立创建端点是设计有意,见后端-04 line748/产品-02D line507)→死锁,首次绑定永不可达(work4 seed binding 实为 E2E 绕此 bug 手插)。修:precheck 缺行放行 + bind 抽 `persistBinding` upsert(无行 insert revision=1/有行 updateById 沿用乐观锁)+ V29 `uk_muse_agent_slot_binding_active` partial unique index(同 work+slot 仅一条 active 并发兜底)+ 前端 mock protected 口径对齐 `protected:` 前缀。真验(muse_slice_live):work1 `writing.continuation` precheck(不再 forbidden)→bind 首次建 binding(slotRevision=1 active)→GET 槽位显示 agent1→AI task 用 `agentSlotKey`(非 agentOverrideRef)→agent1 解析 + runtime 授权过门 → completed + suggestion7(MiniMax-M2.5 真 LLM"星环大陆魔法体系"续写、finishReason stop)。**前端 MVP(df063cb)agentSlotKey 路径真后端端到端走通**;单测 `MuseAgentSlotServiceTest` 26/26 + Controller 4/4 + 前端 4/4;V29 由 Flyway 启动日志证 applied。**E2 智能体生命周期与候选拒绝已补并真验(2026-06-27)**:用户自建 Agent create 会同步建 v1 active version;update 创建下一 active version 并更新 `current_version_id`;版本面板支持 list/activate/archive 非当前版本;archive Agent 同步 archive active 槽位绑定;slot unbind 将 binding revision+1、status=`archived`、响应 `sourceStatus=unbound`;WorkspacePage reject 调 AI owner `POST /suggestions/{id}/reject` 写 `rejected` 与 decision archive。验证:AI targeted 单测+Controller 63/0F/0E,Studio `tsc`/目标 ESLint/targeted Vitest 29/0F/0E,MSW-off Playwright 真后端 5/0F/0E(采纳真生成 suggestionId=80、拒绝真生成 suggestionId=81、Agent 生命周期 DB 核验、Slot unbind DB 核验)。**仍缺/线 B**:Override Slot Contract 主数据(方案 B:合法 slot 校验/系统默认 agent/listWorkAgentSlots 预设骨架未做,slotKey 暂自由字符串、listSlots 缺行仅返回已绑定);前端 agent "运行试用"等深链;质量评测执行器只到 eval-run 入队,未计入 1.0.0 线 A done;grant/runtime 完整物理包拆分仍未做。SSE `event:` 解析漂移已修复并有 targeted Vitest 证据(2026-06-19:`sse.test.ts` + `AIPanel.contract.test.tsx` 17/17);`muse-module-ai-contract-server` 空壳孤儿 submodule 已清理(2026-06-19:从 AI reactor 移除并删除唯一 pom,`P1rAiRouteOwnershipTest` 加回潮门禁)。**2026-06-25:AI 生成候选采纳断层已修(ADR-020 方案 A,commit a16c596)**——`MuseAiRuntimeProjectionService` 此前用 `numericEnvelopeId(Long.parseLong)` 落 authz 快照,但真 runtime envelope 是字符串 `rpe-local-uuid`→解析失败落 null→suggestion 恒不可采纳(merge 1041001001);改直落字符串 envelope(删 numericEnvelopeId)、`MuseAiSuggestionDO.authorizationSnapshotId` 改 String、`AiSuggestionMergeProjectionFacade.getSuggestion` 校验改 hasText、V30 列 BIGINT→VARCHAR(128)。揭 `P1rContentMergeGeneratedSuggestionIT` 假绿(注入数值 9001、真 runtime 不产数值),改真字符串 envelope 转真绿。活体真验:真生成 suggestion authz=rpe-local 字符串→采纳 merge code=0+block revision 自增,生成→采纳→Canonical 全链通。详见总账 2026-06-25 / ADR-020 / memory muse-ai-accept-authz-snapshot-fix。**2026-06-27:E1 accepted decision owner 归档已补**——`AiSuggestionMergeProjectionFacade.markSuggestionAccepted` 作为 AI owner 写 `muse_ai_suggestion_decision`、更新 suggestion `status=accepted/decision_archive_id`、记录 AI command 与 business audit;幂等 replay 不重复写,非 pending 候选返回 conflict 且不预占 command。注意边界:AI 仍只归档 Shadow 决策与审计,不写 Canonical 正文,不保存 `finalContent` 或 provider 原文。P1r 真 PG 切片 6/0F/0E/0S。**2026-06-25(同日续):槽位绑定授权快照孪生收口(C1,commit 854f18c)**——三路盘点发现 `muse_agent_slot_binding.authorization_snapshot_id` 是采纳断层唯一漏网 binding 表(列 BIGINT、DO Long,bind 时 parseLong 对字符串 envelope `rpe-local-` 返 null→授权快照静默丢失、绑定溯源链断);V31 两列(authz+source)BIGINT→VARCHAR(128)、DO Long→String、bind 两处(insert/updateById)直透传删 parseLong(`sameAuthorizationSnapshot` 的 parseLong 保留,字符串 envelope 走 fallback 字符串比较正确)。活体真验(muse_slice_live):真 envelope 经 precheck→bind,binding authz 列存非空 rpe-local(insert id=51+updateById id=50 两路径),列类型 varchar(128)、flyway V31 success;修前恒 null。**此前模块"仍缺"只列方案 B/深链/物理拆分,从未记此类型错配——盘点补上的盲区**。至此 ADR-020 全部 binding 表(knowledge V14 + suggestion/attribution V30 + slot V31)收口完成。 **2026-06-25(同日续):AI 流 SSE 长任务候选丢失已修(commit c77b006、A1+候选③、design-docs/临时-01 review 批准后执行)**——studio e2e ai-generation 真红根因=后端 `MuseAiTaskStreamServiceImpl` SSE 死线 30s 硬编码 < LLM 真时延(MiniMax 11-60s)、更 < 上游自己 `MUSE_AI_NEW_API_TOTAL_TIMEOUT_SECONDS=180`,poll 30s 后 `emitter.complete()` 静默关连接不补 done;前端 `connectAIStream` 一次性连接无重连(重连只在 connectEventStream)→LLM>30s 候选永丢(ai-gen 真红/accept flaky 是同 bug 时延两侧)。修:① 后端删 `DEFAULT_TIMEOUT_MILLIS` 常量→`@Value("${muse.ai.sse.task-timeout-millis:240000}")`(连接死线:70+poll deadline:198 同引用、application.yaml 配置项);② 前端退避重连续 poll + 按 SSE id 去重(后端不读 lastEventId、每次 seq0 全量回放→`lastSeenSequenceNo` 丢弃 id≤已见,避免重复渲染候选)+总超时 300s `onError(SSE_TIMEOUT)` 不静默卡死(返回仍 AbortController、AIPanel 零改动);③ A1 `scripts/dev/start-muse-server-infra.sh` 固化追加 source acceptance env(防漏 `MUSE_AI_NEW_API_*`→`UnavailableMuseAiRuntimeClient` 兜底秒失败 `AI_NEW_API_UNAVAILABLE`)。验证:后端 SSE 单测 14/14+前端 107/107(含重连去重用例)、活体 e2e ai-generation+accept 连跑 3 轮 6/6、重连幂等方法 A 活体证;残留=本机 LLM 6.8-19s 未自然>30s、慢路径靠「240s 数学覆盖上游 180s+单测+方法 A 重连活体证」三重保证。**通则(避免重踩"后端比上游先关门")**:SSE 长任务死线须 ≥ 上游 TOTAL_TIMEOUT;AI 流重连必须以 SSE id 幂等去重(后端 seq0 全量回放、无 lastEventId)。详见 memory muse-ai-generation-sse-timeout-bug。 +- **E6 admin AI 配置接线(2026-06-27)**:Muse Admin AI 配置页 Prompt 激活、质量评估启动、质量策略版本创建、system Agent 创建/版本发布/归档停用、Tool Grant 登记/调整已调用真实 admin API,不再展示硬编码评估分数。后端 `adminArchiveAgent` 只允许 system scope,归档 Agent 时同步归档 active 槽位绑定,避免运行时继续消费已停用系统能力;质量策略回滚仍保持 preview-only/线 B。RC 后已补 Quality eval 本地确定性执行器:默认 worker 领取 `evaluation` queued job,合成/脱敏 dataset 才可执行,写 sample result/聚合 metricSummary 并推进 run/job 终态;不调用外部 LLM、不读取私有正文,审计只落 dataset hash/聚合摘要。验证:`muse-admin` AI 配置目标 Vitest 2 files / 9 tests passed,`@vben/web-antd typecheck` 通过;后端 `MuseAgentServiceTest`+`AdminMuseAgentControllerAnnotationTest` 31/0F/0E/0S,`MuseToolGrantServiceTest`+`LocalSecurityToolGrantApprovalFacadeTest` 13/0F/0E/0S,`MuseQualityEvaluationWorkerTest`+`MuseAiJobMapperTest`+`MuseQualityPolicyServiceTest` 17/0F/0E,server 契约/覆盖门 27/0F/0E/0S。 +- **RC 后完整导入解析 owner(2026-06-28)**:新增 `MuseAiImportParseApi` 与 `MuseAiImportParseService`,由 AI owner 持久化 `muse_ai_parse_job` / `muse_ai_chapter_parse_result`,读取 Content 传入的稳定 storageRef 后支持 txt/markdown/docx/epub 文本抽取;默认按 Markdown/CN/EN 标题做确定性章节拆分,`parseConfig.mode=full_book` 或 `strategy=llm_full_book` 时调用 New-API LLM 全书解析,并对 timeout/429/5xx/连接失败做可重试处理。解析只输出 pending-review Shadow 章节并支持确认/驳回/replay,不写 Content Canonical,不写 Knowledge Local KB。验证:`MuseAiImportParseServiceTest`、`NewApiMuseAiImportLlmParserTest`、`P1rImportLlmNewApiLiveAcceptanceIT` live 与完整导入链路 real-PG/live 测试通过。 +- **关键风险 / TODO**:**AI runtime 生成 e2e 已端到端真验(2026-06-24)**——方案1(d647f0c 剥离 sourceSnapshotId 解授权死结)+ D(4dadc7b `DefaultDBFieldHandler` 无登录兜底 `system`,根治 scheduled worker 写 creator/updater null 违约 NOT NULL;BaseDO 带 jdbcType 强制入列致 insert/update 同病)+ agent_version config 补 modelKey=MiniMax-M2.5 + 启动 source `scripts/dev/p1r-external-acceptance.env`(New-API/RAGFlow 配置)→ task5 completed+suggestion+MiniMax-M2.5 真 LLM 生成(P-B envelope/approved、P-C contextAssembly 固化 ✓)。**P-A/P-B/P-C 完整 e2e 已端到端真验打通(2026-06-24,task7)**:P-A 检索命中根因=knowledge 模块 bug——`HttpRagFlowKnowledgeRuntimeClient.retrieveChunks` 把 `document_ids` 以 null 入 retrieval body,RAGFlow 要求其为 list 故拒(code 102 "documents should be a list")→检索 fail-closed→chunkCount=0(task3/4/5/6 全 0 即此 bug,非 RAGFlow 数据/凭据/timeout)。已修(commit 86c6dba:null 时不入 body + 回归测试)。task7 contextAssembly `chunkCount=1/retrievedKbIds=[2]/authorizationSnapshotIds=[authsnap-e2e-1]`,LLM(MiniMax-M2.5)基于 grounding 续写。**02D 槽位首次创建死锁已修+真验(2026-06-24,commit 9115d4f)**:根因=precheck/bind 都 `requireSlotBinding`(缺行抛 forbidden),而"替换"是设计中唯一创建入口("替换"=首次写入、无独立创建端点是设计有意,见后端-04 line748/产品-02D line507)→死锁,首次绑定永不可达(work4 seed binding 实为 E2E 绕此 bug 手插)。修:precheck 缺行放行 + bind 抽 `persistBinding` upsert(无行 insert revision=1/有行 updateById 沿用乐观锁)+ V29 `uk_muse_agent_slot_binding_active` partial unique index(同 work+slot 仅一条 active 并发兜底)+ 前端 mock protected 口径对齐 `protected:` 前缀。真验(muse_slice_live):work1 `writing.continuation` precheck(不再 forbidden)→bind 首次建 binding(slotRevision=1 active)→GET 槽位显示 agent1→AI task 用 `agentSlotKey`(非 agentOverrideRef)→agent1 解析 + runtime 授权过门 → completed + suggestion7(MiniMax-M2.5 真 LLM"星环大陆魔法体系"续写、finishReason stop)。**前端 MVP(df063cb)agentSlotKey 路径真后端端到端走通**;单测 `MuseAgentSlotServiceTest` 26/26 + Controller 4/4 + 前端 4/4;V29 由 Flyway 启动日志证 applied。**E2 智能体生命周期与候选拒绝已补并真验(2026-06-27)**:用户自建 Agent create 会同步建 v1 active version;update 创建下一 active version并更新 `current_version_id`;版本面板支持 list/activate/archive 非当前版本;archive Agent 同步 archive active 槽位绑定;slot unbind 将 binding revision+1、status=`archived`、响应 `sourceStatus=unbound`;WorkspacePage reject 调 AI owner `POST /suggestions/{id}/reject` 写 `rejected` 与 decision archive。验证:AI targeted 单测+Controller 63/0F/0E,Studio `tsc`/目标 ESLint/targeted Vitest 29/0F/0E,MSW-off Playwright 真后端 5/0F/0E(采纳真生成 suggestionId=80、拒绝真生成 suggestionId=81、Agent 生命周期 DB 核验、Slot unbind DB 核验)。**仍缺/线 B**:Override Slot Contract 主数据(方案 B:合法 slot 校验/系统默认 agent/listWorkAgentSlots 预设骨架未做,slotKey 暂自由字符串、listSlots 缺行仅返回已绑定);前端 agent "运行试用"等深链;Quality eval 真实评估集管理/外发授权/LLM judge 仍后置;grant/runtime 完整物理包拆分仍未做。SSE `event:` 解析漂移已修复并有 targeted Vitest 证据(2026-06-19:`sse.test.ts` + `AIPanel.contract.test.tsx` 17/17);`muse-module-ai-contract-server` 空壳孤儿 submodule 已清理(2026-06-19:从 AI reactor 移除并删除唯一 pom,`P1rAiRouteOwnershipTest` 加回潮门禁)。**2026-06-25:AI 生成候选采纳断层已修(ADR-020 方案 A,commit a16c596)**——`MuseAiRuntimeProjectionService` 此前用 `numericEnvelopeId(Long.parseLong)` 落 authz 快照,但真 runtime envelope 是字符串 `rpe-local-uuid`→解析失败落 null→suggestion 恒不可采纳(merge 1041001001);改直落字符串 envelope(删 numericEnvelopeId)、`MuseAiSuggestionDO.authorizationSnapshotId` 改 String、`AiSuggestionMergeProjectionFacade.getSuggestion` 校验改 hasText、V30 列 BIGINT→VARCHAR(128)。揭 `P1rContentMergeGeneratedSuggestionIT` 假绿(注入数值 9001、真 runtime 不产数值),改真字符串 envelope 转真绿。活体真验:真生成 suggestion authz=rpe-local 字符串→采纳 merge code=0+block revision 自增,生成→采纳→Canonical 全链通。详见总账 2026-06-25 / ADR-020 / memory muse-ai-accept-authz-snapshot-fix。**2026-06-27:E1 accepted decision owner 归档已补**——`AiSuggestionMergeProjectionFacade.markSuggestionAccepted` 作为 AI owner 写 `muse_ai_suggestion_decision`、更新 suggestion `status=accepted/decision_archive_id`、记录 AI command 与 business audit;幂等 replay 不重复写,非 pending 候选返回 conflict 且不预占 command。注意边界:AI 仍只归档 Shadow 决策与审计,不写 Canonical 正文,不保存 `finalContent` 或 provider 原文。P1r 真 PG 切片 6/0F/0E/0S。**2026-06-25(同日续):槽位绑定授权快照孪生收口(C1,commit 854f18c)**——三路盘点发现 `muse_agent_slot_binding.authorization_snapshot_id` 是采纳断层唯一漏网 binding 表(列 BIGINT、DO Long,bind 时 parseLong 对字符串 envelope `rpe-local-` 返 null→授权快照静默丢失、绑定溯源链断);V31 两列(authz+source)BIGINT→VARCHAR(128)、DO Long→String、bind 两处(insert/updateById)直透传删 parseLong(`sameAuthorizationSnapshot` 的 parseLong 保留,字符串 envelope 走 fallback 字符串比较正确)。活体真验(muse_slice_live):真 envelope 经 precheck→bind,binding authz 列存非空 rpe-local(insert id=51+updateById id=50 两路径),列类型 varchar(128)、flyway V31 success;修前恒 null。**此前模块"仍缺"只列方案 B/深链/物理拆分,从未记此类型错配——盘点补上的盲区**。至此 ADR-020 全部 binding 表(knowledge V14 + suggestion/attribution V30 + slot V31)收口完成。 **2026-06-25(同日续):AI 流 SSE 长任务候选丢失已修(commit c77b006、A1+候选③、design-docs/临时-01 review 批准后执行)**——studio e2e ai-generation 真红根因=后端 `MuseAiTaskStreamServiceImpl` SSE 死线 30s 硬编码 < LLM 真时延(MiniMax 11-60s)、更 < 上游自己 `MUSE_AI_NEW_API_TOTAL_TIMEOUT_SECONDS=180`,poll 30s 后 `emitter.complete()` 静默关连接不补 done;前端 `connectAIStream` 一次性连接无重连(重连只在 connectEventStream)→LLM>30s 候选永丢(ai-gen 真红/accept flaky 是同 bug 时延两侧)。修:① 后端删 `DEFAULT_TIMEOUT_MILLIS` 常量→`@Value("${muse.ai.sse.task-timeout-millis:240000}")`(连接死线:70+poll deadline:198 同引用、application.yaml 配置项);② 前端退避重连续 poll + 按 SSE id 去重(后端不读 lastEventId、每次 seq0 全量回放→`lastSeenSequenceNo` 丢弃 id≤已见,避免重复渲染候选)+总超时 300s `onError(SSE_TIMEOUT)` 不静默卡死(返回仍 AbortController、AIPanel 零改动);③ A1 `scripts/dev/start-muse-server-infra.sh` 固化追加 source acceptance env(防漏 `MUSE_AI_NEW_API_*`→`UnavailableMuseAiRuntimeClient` 兜底秒失败 `AI_NEW_API_UNAVAILABLE`)。验证:后端 SSE 单测 14/14+前端 107/107(含重连去重用例)、活体 e2e ai-generation+accept 连跑 3 轮 6/6、重连幂等方法 A 活体证;残留=本机 LLM 6.8-19s 未自然>30s、慢路径靠「240s 数学覆盖上游 180s+单测+方法 A 重连活体证」三重保证。**通则(避免重踩"后端比上游先关门")**:SSE 长任务死线须 ≥ 上游 TOTAL_TIMEOUT;AI 流重连必须以 SSE id 幂等去重(后端 seq0 全量回放、无 lastEventId)。详见 memory muse-ai-generation-sse-timeout-bug。 diff --git a/muse-cloud/muse-module-ai/muse-module-ai-server/src/main/java/cn/iocoder/muse/module/ai/application/muse/MuseQualityEvaluationExecutor.java b/muse-cloud/muse-module-ai/muse-module-ai-server/src/main/java/cn/iocoder/muse/module/ai/application/muse/MuseQualityEvaluationExecutor.java new file mode 100644 index 00000000..d53531de --- /dev/null +++ b/muse-cloud/muse-module-ai/muse-module-ai-server/src/main/java/cn/iocoder/muse/module/ai/application/muse/MuseQualityEvaluationExecutor.java @@ -0,0 +1,426 @@ +package cn.iocoder.muse.module.ai.application.muse; + +import cn.iocoder.muse.framework.common.util.json.JsonUtils; +import cn.iocoder.muse.framework.tenant.core.context.TenantContextHolder; +import cn.iocoder.muse.module.ai.dal.dataobject.muse.MuseAiEvaluationRunDO; +import cn.iocoder.muse.module.ai.dal.dataobject.muse.MuseAiEvaluationSampleResultDO; +import cn.iocoder.muse.module.ai.dal.dataobject.muse.MuseAiJobDO; +import cn.iocoder.muse.module.ai.dal.dataobject.muse.MuseQualityPolicyVersionDO; +import cn.iocoder.muse.module.ai.dal.mysql.muse.MuseAiEvaluationRunMapper; +import cn.iocoder.muse.module.ai.dal.mysql.muse.MuseAiEvaluationSampleResultMapper; +import cn.iocoder.muse.module.ai.dal.mysql.muse.MuseAiJobMapper; +import cn.iocoder.muse.module.ai.dal.mysql.muse.MuseQualityPolicyVersionMapper; +import com.fasterxml.jackson.databind.JsonNode; +import jakarta.annotation.Resource; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Transactional; +import org.springframework.util.StringUtils; + +import java.math.BigDecimal; +import java.math.RoundingMode; +import java.nio.charset.StandardCharsets; +import java.security.MessageDigest; +import java.security.NoSuchAlgorithmException; +import java.time.LocalDateTime; +import java.util.ArrayList; +import java.util.HexFormat; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Locale; +import java.util.Map; +import java.util.regex.Pattern; + +/** + * 质量评估执行器。 + * + *

当前版本只执行合成/脱敏评估集的确定性本地评估,先把 queued run/job 推进为可观测终态。 + * 不调用外部 LLM,不读取用户私有正文,不把 dataset 原文扩散到样本结果。后续接入真实评估集和 LLM judge 时, + * 应复用本执行器的状态机与样本结果落库链路。

+ */ +@Slf4j +@Service +public class MuseQualityEvaluationExecutor { + + private static final String STATUS_RUNNING = "running"; + private static final String STATUS_COMPLETED = "completed"; + private static final String STATUS_FAILED = "failed"; + private static final String SAMPLE_STATUS_PASSED = "passed"; + private static final String SAMPLE_STATUS_FAILED = "failed"; + private static final String OPERATION_EXECUTE_EVALUATION = "executeQualityEvaluationRun"; + private static final String TARGET_TYPE_EVALUATION_RUN = "evaluationRun"; + private static final String ERROR_INVALID_JOB = "AI_QUALITY_EVALUATION_INVALID_JOB"; + private static final String ERROR_DATASET_BLOCKED = "AI_QUALITY_EVALUATION_DATASET_BLOCKED"; + private static final String ERROR_SAMPLE_EMPTY = "AI_QUALITY_EVALUATION_SAMPLE_EMPTY"; + private static final String ERROR_SAMPLE_TOO_LARGE = "AI_QUALITY_EVALUATION_SAMPLE_TOO_LARGE"; + private static final String ERROR_POLICY_INVALID = "AI_QUALITY_EVALUATION_POLICY_INVALID"; + private static final String ERROR_EXECUTION_FAILED = "AI_QUALITY_EVALUATION_EXECUTION_FAILED"; + private static final int MAX_SAMPLE_COUNT = 200; + private static final Pattern SECRET_PATTERN = Pattern.compile( + "(?i)(sk-[a-z0-9_-]+|bearer\\s+\\S+|token\\s*[:=]?\\s*\\S+|secret\\s*[:=]?\\s*\\S+|api[-_ ]?key\\s*[:=]?\\s*\\S+)"); + + @Resource + private MuseAiJobMapper jobMapper; + @Resource + private MuseAiEvaluationRunMapper evaluationRunMapper; + @Resource + private MuseAiEvaluationSampleResultMapper sampleResultMapper; + @Resource + private MuseQualityPolicyVersionMapper qualityPolicyVersionMapper; + @Resource + private MuseAiAuditService auditService; + + /** + * 执行已领取的质量评估 job。 + * + *

调用方必须先把租户上下文恢复到领取行的 tenantId;本方法全程只写脱敏摘要和样本聚合。

+ */ + @Transactional(rollbackFor = Exception.class) + public void executeEvaluationJob(Long jobPkId) { + MuseAiJobDO job = jobMapper.selectById(jobPkId); + if (job == null || !StringUtils.hasText(job.getJobId())) { + throw new EvaluationExecutionException(ERROR_INVALID_JOB, "质量评估 Job 不存在或缺少 jobId"); + } + MuseAiEvaluationRunDO run = evaluationRunMapper.selectByJobId(job.getJobId()); + if (run == null) { + failJobOnly(job, ERROR_INVALID_JOB, "质量评估运行不存在"); + return; + } + + markRunning(run); + try { + EvaluationReport report = evaluate(run); + insertSampleResults(run, report); + markCompleted(run, job, report); + auditSafely(run, job, report.auditSummary(), "succeeded", null, null); + } catch (EvaluationExecutionException ex) { + markFailed(run, job, ex.errorCode(), ex.getMessage()); + auditSafely(run, job, "{\"errorCode\":\"" + ex.errorCode() + "\"}", STATUS_FAILED, + ex.errorCode(), ex.getMessage()); + } catch (RuntimeException ex) { + String message = sanitizeError(ex.getMessage()); + markFailed(run, job, ERROR_EXECUTION_FAILED, message); + auditSafely(run, job, "{\"errorCode\":\"" + ERROR_EXECUTION_FAILED + "\"}", STATUS_FAILED, + ERROR_EXECUTION_FAILED, message); + } + } + + private EvaluationReport evaluate(MuseAiEvaluationRunDO run) { + String datasetReference = datasetReference(run); + requireSafeDatasetReference(datasetReference); + int sampleCount = safeSampleCount(run.getSampleCount()); + if (sampleCount <= 0) { + throw new EvaluationExecutionException(ERROR_SAMPLE_EMPTY, "质量评估样本数必须大于 0"); + } + if (sampleCount > MAX_SAMPLE_COUNT) { + throw new EvaluationExecutionException(ERROR_SAMPLE_TOO_LARGE, "质量评估样本数超过内部测试上限"); + } + + MuseQualityPolicyVersionDO policyVersion = qualityPolicyVersionMapper.selectByPolicyKeyAndVersion( + run.getPolicyKey(), run.getPolicyVersionNo()); + List dimensions = parseDimensions(policyVersion); + if (dimensions.isEmpty()) { + throw new EvaluationExecutionException(ERROR_POLICY_INVALID, "质量策略缺少可执行维度"); + } + + List samples = new ArrayList<>(); + Map dimensionTotals = new LinkedHashMap<>(); + int passedCount = 0; + for (int index = 1; index <= sampleCount; index++) { + SampleReport sample = evaluateSample(run, datasetReference, dimensions, index); + samples.add(sample); + if (sample.passed()) { + passedCount += 1; + } + sample.dimensionScores().forEach((dimension, score) -> + dimensionTotals.merge(dimension, score, Double::sum)); + } + + Map dimensionScores = new LinkedHashMap<>(); + dimensionTotals.forEach((dimension, total) -> dimensionScores.put(dimension, round(total / sampleCount))); + double overallScore = round(samples.stream().mapToDouble(SampleReport::score).average().orElse(0D)); + int failedCount = sampleCount - passedCount; + return new EvaluationReport(samples, dimensionScores, overallScore, passedCount, failedCount, + JsonUtils.toJsonString(Map.of( + "datasetReference", datasetReference, + "evaluator", "deterministic-local-v1", + "executionMode", "synthetic_or_redacted_only", + "sampleCount", sampleCount, + "passedCount", passedCount, + "failedCount", failedCount, + "overallScore", overallScore, + "dimensionScores", dimensionScores))); + } + + private SampleReport evaluateSample(MuseAiEvaluationRunDO run, String datasetReference, + List dimensions, int sampleIndex) { + Map scores = new LinkedHashMap<>(); + double weightedTotal = 0D; + double weightTotal = 0D; + boolean passed = true; + for (DimensionConfig dimension : dimensions) { + double score = deterministicScore(run, datasetReference, sampleIndex, dimension.name()); + scores.put(dimension.name(), score); + weightedTotal += score * dimension.weight(); + weightTotal += dimension.weight(); + if (score < dimension.threshold()) { + passed = false; + } + } + double overall = round(weightTotal <= 0D ? 0D : weightedTotal / weightTotal); + return new SampleReport(sampleId(sampleIndex), passed, overall, scores); + } + + private void insertSampleResults(MuseAiEvaluationRunDO run, EvaluationReport report) { + for (SampleReport sample : report.samples()) { + MuseAiEvaluationSampleResultDO row = new MuseAiEvaluationSampleResultDO(); + row.setRunId(run.getRunId()); + row.setSampleId(sample.sampleId()); + row.setOwnerUserId(run.getOwnerUserId()); + row.setStatus(sample.passed() ? SAMPLE_STATUS_PASSED : SAMPLE_STATUS_FAILED); + row.setScore(BigDecimal.valueOf(sample.score()).setScale(4, RoundingMode.HALF_UP)); + row.setMetricSummary(JsonUtils.toJsonString(Map.of("dimensionScores", sample.dimensionScores()))); + row.setInputSummary(JsonUtils.toJsonString(Map.of( + "datasetReferenceHash", sha256(datasetReference(run)), + "sampleId", sample.sampleId(), + "inputRedacted", true))); + row.setOutputSummary(JsonUtils.toJsonString(Map.of( + "evaluator", "deterministic-local-v1", + "rawOutputStored", false))); + row.setTenantId(TenantContextHolder.getRequiredTenantId()); + sampleResultMapper.insert(row); + } + } + + private List parseDimensions(MuseQualityPolicyVersionDO policyVersion) { + if (policyVersion == null || !StringUtils.hasText(policyVersion.getThresholdSnapshot())) { + return List.of(); + } + try { + JsonNode root = JsonUtils.getObjectMapper().readTree(policyVersion.getThresholdSnapshot()); + if (root == null || !root.isArray()) { + return List.of(); + } + List dimensions = new ArrayList<>(); + for (JsonNode node : root) { + String name = text(node, "name"); + if (!StringUtils.hasText(name)) { + continue; + } + double threshold = clamp(number(node, "threshold", 0.8D)); + double weight = Math.max(0.01D, number(node, "weight", 1.0D)); + dimensions.add(new DimensionConfig(name.trim(), threshold, weight)); + } + return dimensions; + } catch (Exception ex) { + return List.of(); + } + } + + private String datasetReference(MuseAiEvaluationRunDO run) { + JsonNode summary = json(run == null ? null : run.getMetricSummary()); + String value = text(summary, "datasetReference"); + return value == null ? "" : value.trim(); + } + + private void requireSafeDatasetReference(String datasetReference) { + if (!isSafeDatasetReference(datasetReference)) { + // 当前执行器没有私有正文授权/脱敏证明模型;宁可失败,也不能把 work/file 等私有来源当成评估集。 + throw new EvaluationExecutionException(ERROR_DATASET_BLOCKED, "质量评估数据集未标记为合成或脱敏"); + } + } + + private boolean isSafeDatasetReference(String datasetReference) { + String normalized = datasetReference == null ? "" : datasetReference.trim().toLowerCase(Locale.ROOT); + return normalized.startsWith("dataset://") + || normalized.startsWith("synthetic://") + || normalized.startsWith("redacted://") + || normalized.contains("redacted") + || normalized.contains("smoke"); + } + + private void markRunning(MuseAiEvaluationRunDO run) { + MuseAiEvaluationRunDO update = new MuseAiEvaluationRunDO(); + update.setId(run.getId()); + update.setStatus(STATUS_RUNNING); + update.setStartedAt(run.getStartedAt() == null ? LocalDateTime.now() : run.getStartedAt()); + evaluationRunMapper.updateById(update); + } + + private void markCompleted(MuseAiEvaluationRunDO run, MuseAiJobDO job, EvaluationReport report) { + LocalDateTime now = LocalDateTime.now(); + MuseAiEvaluationRunDO runUpdate = new MuseAiEvaluationRunDO(); + runUpdate.setId(run.getId()); + runUpdate.setStatus(STATUS_COMPLETED); + runUpdate.setPassedCount(report.passedCount()); + runUpdate.setFailedCount(report.failedCount()); + runUpdate.setMetricSummary(report.metricSummary()); + runUpdate.setFinishedAt(now); + evaluationRunMapper.updateById(runUpdate); + + MuseAiJobDO jobUpdate = new MuseAiJobDO(); + jobUpdate.setId(job.getId()); + jobUpdate.setStatus(STATUS_COMPLETED); + jobUpdate.setRetryable(false); + jobUpdate.setResultSummary(report.metricSummary()); + jobUpdate.setFinishedAt(now); + jobMapper.updateById(jobUpdate); + } + + private void markFailed(MuseAiEvaluationRunDO run, MuseAiJobDO job, String errorCode, String errorMessage) { + LocalDateTime now = LocalDateTime.now(); + MuseAiEvaluationRunDO runUpdate = new MuseAiEvaluationRunDO(); + runUpdate.setId(run.getId()); + runUpdate.setStatus(STATUS_FAILED); + runUpdate.setErrorCode(errorCode); + runUpdate.setErrorMessage(sanitizeError(errorMessage)); + runUpdate.setFinishedAt(now); + evaluationRunMapper.updateById(runUpdate); + failJobOnly(job, errorCode, errorMessage); + } + + private void failJobOnly(MuseAiJobDO job, String errorCode, String errorMessage) { + MuseAiJobDO jobUpdate = new MuseAiJobDO(); + jobUpdate.setId(job.getId()); + jobUpdate.setStatus(STATUS_FAILED); + jobUpdate.setRetryable(false); + jobUpdate.setErrorCode(errorCode); + jobUpdate.setErrorMessage(sanitizeError(errorMessage)); + jobUpdate.setFinishedAt(LocalDateTime.now()); + jobMapper.updateById(jobUpdate); + } + + private void auditSafely(MuseAiEvaluationRunDO run, MuseAiJobDO job, String responseSummary, + String status, String errorCode, String errorMessage) { + try { + audit(run, job, responseSummary, status, errorCode, errorMessage); + } catch (RuntimeException ex) { + // 审计是可观测补充,不能因为审计写失败回滚 run/job 终态,避免 worker 留下永久 running。 + log.warn("[auditSafely][runId({}) jobId({}) audit failed, errorType={}]", + run == null ? null : run.getRunId(), job == null ? null : job.getJobId(), + ex.getClass().getSimpleName()); + } + } + + private void audit(MuseAiEvaluationRunDO run, MuseAiJobDO job, String responseSummary, + String status, String errorCode, String errorMessage) { + Map requestSummary = new LinkedHashMap<>(); + requestSummary.put("policyKey", run.getPolicyKey()); + requestSummary.put("policyVersion", run.getPolicyVersionNo()); + requestSummary.put("datasetReferenceHash", sha256(datasetReference(run))); + requestSummary.put("datasetReferenceAllowed", isSafeDatasetReference(datasetReference(run))); + requestSummary.put("sampleCount", safeSampleCount(run.getSampleCount())); + auditService.record(MuseAiAuditService.AuditCreateReq.builder() + .operationId(OPERATION_EXECUTE_EVALUATION) + .actorUserId(job.getActorUserId()) + .ownerUserId(job.getOwnerUserId()) + .side("worker") + .targetType(TARGET_TYPE_EVALUATION_RUN) + .targetId(run.getId()) + .commandId(job.getCommandId()) + .requestHash(job.getRequestHash()) + .requestSummary(JsonUtils.toJsonString(requestSummary)) + .responseSummary(responseSummary) + .status(status) + .errorCode(errorCode) + .errorMessage(errorMessage) + .build()); + } + + private JsonNode json(String value) { + if (!StringUtils.hasText(value)) { + return JsonUtils.getObjectMapper().createObjectNode(); + } + try { + return JsonUtils.getObjectMapper().readTree(value); + } catch (Exception ex) { + return JsonUtils.getObjectMapper().createObjectNode(); + } + } + + private String text(JsonNode node, String fieldName) { + JsonNode value = node == null ? null : node.get(fieldName); + return value == null || value.isNull() ? null : value.asText(); + } + + private double number(JsonNode node, String fieldName, double defaultValue) { + JsonNode value = node == null ? null : node.get(fieldName); + return value == null || !value.isNumber() ? defaultValue : value.asDouble(); + } + + private double deterministicScore(MuseAiEvaluationRunDO run, String datasetReference, + int sampleIndex, String dimensionName) { + String seed = run.getPolicyKey() + "\n" + run.getPolicyVersionNo() + "\n" + + datasetReference + "\n" + sampleIndex + "\n" + dimensionName; + String hash = sha256(seed); + long bucket = Long.parseUnsignedLong(hash.substring(0, 12), 16) % 3500L; + return round(0.65D + bucket / 10_000D); + } + + private String sha256(String payload) { + try { + MessageDigest digest = MessageDigest.getInstance("SHA-256"); + return HexFormat.of().formatHex(digest.digest(String.valueOf(payload).getBytes(StandardCharsets.UTF_8))); + } catch (NoSuchAlgorithmException ex) { + throw new IllegalStateException("JDK 缺少 SHA-256 摘要算法", ex); + } + } + + private String sampleId(int sampleIndex) { + return "sample-" + String.format(Locale.ROOT, "%04d", sampleIndex); + } + + private double clamp(double value) { + return Math.max(0D, Math.min(1D, value)); + } + + private double round(double value) { + return BigDecimal.valueOf(value).setScale(4, RoundingMode.HALF_UP).doubleValue(); + } + + private int safeSampleCount(Integer sampleCount) { + return sampleCount == null ? 0 : Math.max(0, sampleCount); + } + + private String sanitizeError(String message) { + if (!StringUtils.hasText(message)) { + return "Quality evaluation failed"; + } + return SECRET_PATTERN.matcher(message).replaceAll("[REDACTED]"); + } + + private record DimensionConfig(String name, double threshold, double weight) { + } + + private record SampleReport(String sampleId, boolean passed, double score, Map dimensionScores) { + } + + private record EvaluationReport(List samples, Map dimensionScores, + double overallScore, int passedCount, int failedCount, String metricSummary) { + + String auditSummary() { + return JsonUtils.toJsonString(Map.of( + "overallScore", overallScore, + "passedCount", passedCount, + "failedCount", failedCount, + "dimensionScores", dimensionScores)); + } + + } + + private static class EvaluationExecutionException extends RuntimeException { + + private final String errorCode; + + EvaluationExecutionException(String errorCode, String message) { + super(message); + this.errorCode = errorCode; + } + + String errorCode() { + return errorCode; + } + + } + +} diff --git a/muse-cloud/muse-module-ai/muse-module-ai-server/src/main/java/cn/iocoder/muse/module/ai/application/muse/MuseQualityEvaluationWorker.java b/muse-cloud/muse-module-ai/muse-module-ai-server/src/main/java/cn/iocoder/muse/module/ai/application/muse/MuseQualityEvaluationWorker.java new file mode 100644 index 00000000..a7f3b67e --- /dev/null +++ b/muse-cloud/muse-module-ai/muse-module-ai-server/src/main/java/cn/iocoder/muse/module/ai/application/muse/MuseQualityEvaluationWorker.java @@ -0,0 +1,97 @@ +package cn.iocoder.muse.module.ai.application.muse; + +import cn.iocoder.muse.framework.tenant.core.util.TenantUtils; +import cn.iocoder.muse.module.ai.dal.dataobject.muse.MuseAiJobDO; +import cn.iocoder.muse.module.ai.dal.mysql.muse.MuseAiJobMapper; +import cn.iocoder.muse.module.ai.framework.ai.config.MuseAiProperties; +import jakarta.annotation.Resource; +import lombok.extern.slf4j.Slf4j; +import org.springframework.scheduling.annotation.Scheduled; +import org.springframework.stereotype.Component; + +import java.time.LocalDateTime; + +/** + * 质量评估 worker。 + * + *

后台定时默认启用;发布环境可显式设置 {@code muse.ai.evaluation.worker.enabled=false} 暂停消费。 + * 测试和运维脚本可直接调用 {@link #dispatchOnce()} 做单步执行,避免依赖真实时间。

+ */ +@Slf4j +@Component +public class MuseQualityEvaluationWorker { + + private static final String STATUS_FAILED = "failed"; + private static final String ERROR_INVALID_JOB = "AI_QUALITY_EVALUATION_INVALID_JOB"; + private static final String ERROR_DISPATCH_FAILED = "AI_QUALITY_EVALUATION_DISPATCH_FAILED"; + + @Resource + private MuseAiJobMapper jobMapper; + @Resource + private MuseQualityEvaluationExecutor evaluationExecutor; + @Resource + private MuseAiProperties properties; + + public MuseQualityEvaluationWorker() { + } + + MuseQualityEvaluationWorker(MuseAiJobMapper jobMapper, MuseQualityEvaluationExecutor evaluationExecutor, + MuseAiProperties properties) { + this.jobMapper = jobMapper; + this.evaluationExecutor = evaluationExecutor; + this.properties = properties; + } + + @Scheduled(initialDelayString = "${muse.ai.evaluation.worker.initial-delay-ms:1000}", + fixedDelayString = "${muse.ai.evaluation.worker.fixed-delay-ms:1000}") + public void dispatchScheduled() { + dispatchOnce(); + } + + public int dispatchOnce() { + if (!isEnabled()) { + return 0; + } + MuseAiJobDO job = TenantUtils.executeIgnore(() -> jobMapper.claimNextQueuedEvaluationJob()); + if (job == null) { + return 0; + } + if (job.getTenantId() == null || job.getId() == null) { + failClaimedJob(job, ERROR_INVALID_JOB, "Quality evaluation worker claimed invalid job"); + return 0; + } + try { + TenantUtils.execute(job.getTenantId(), () -> evaluationExecutor.executeEvaluationJob(job.getId())); + return 1; + } catch (RuntimeException ex) { + log.warn("[dispatchOnce][jobPkId({}) executor failed closed, errorType={}]", + job.getId(), ex.getClass().getSimpleName()); + failClaimedJob(job, ERROR_DISPATCH_FAILED, "Quality evaluation dispatcher failed"); + return 0; + } + } + + private boolean isEnabled() { + return properties != null + && properties.getEvaluation() != null + && properties.getEvaluation().getWorker() != null + && properties.getEvaluation().getWorker().isEnabled(); + } + + private void failClaimedJob(MuseAiJobDO claimedJob, String errorCode, String errorMessage) { + if (claimedJob == null || claimedJob.getId() == null) { + return; + } + MuseAiJobDO update = new MuseAiJobDO(); + update.setId(claimedJob.getId()); + update.setStatus(STATUS_FAILED); + update.setRetryable(false); + update.setErrorCode(errorCode); + update.setErrorMessage(errorMessage); + update.setFinishedAt(LocalDateTime.now()); + update.setTenantId(claimedJob.getTenantId()); + // claim 已经把 job 推到 running;异常路径必须关闭 job,避免后台 worker 留下永久 running。 + TenantUtils.executeIgnore(() -> jobMapper.updateById(update)); + } + +} diff --git a/muse-cloud/muse-module-ai/muse-module-ai-server/src/main/java/cn/iocoder/muse/module/ai/application/muse/MuseQualityPolicyServiceImpl.java b/muse-cloud/muse-module-ai/muse-module-ai-server/src/main/java/cn/iocoder/muse/module/ai/application/muse/MuseQualityPolicyServiceImpl.java index 48e26490..0b449cea 100644 --- a/muse-cloud/muse-module-ai/muse-module-ai-server/src/main/java/cn/iocoder/muse/module/ai/application/muse/MuseQualityPolicyServiceImpl.java +++ b/muse-cloud/muse-module-ai/muse-module-ai-server/src/main/java/cn/iocoder/muse/module/ai/application/muse/MuseQualityPolicyServiceImpl.java @@ -54,6 +54,7 @@ public class MuseQualityPolicyServiceImpl implements MuseQualityPolicyService { private static final String STATUS_DRAFT = "draft"; private static final String STATUS_QUEUED = "queued"; private static final String JOB_TYPE_EVALUATION = "evaluation"; + private static final int DEFAULT_EVALUATION_SAMPLE_SIZE = 20; private static final AtomicLong NUMERIC_ID_SEQUENCE = new AtomicLong(System.currentTimeMillis() * 1000L); private static final Pattern SECRET_PATTERN = Pattern.compile("(?i)(sk-[a-z0-9_-]+|bearer\\s+\\S+|private response body|token\\s*[:=]?\\s*\\S+)"); @@ -387,7 +388,7 @@ public class MuseQualityPolicyServiceImpl implements MuseQualityPolicyService { } private int safeSampleSize(Integer sampleSize) { - return sampleSize == null || sampleSize < 0 ? 0 : sampleSize; + return sampleSize == null ? DEFAULT_EVALUATION_SAMPLE_SIZE : Math.max(0, sampleSize); } private Long nextNumericId() { diff --git a/muse-cloud/muse-module-ai/muse-module-ai-server/src/main/java/cn/iocoder/muse/module/ai/dal/dataobject/muse/MuseAiEvaluationSampleResultDO.java b/muse-cloud/muse-module-ai/muse-module-ai-server/src/main/java/cn/iocoder/muse/module/ai/dal/dataobject/muse/MuseAiEvaluationSampleResultDO.java new file mode 100644 index 00000000..b2acbfb5 --- /dev/null +++ b/muse-cloud/muse-module-ai/muse-module-ai-server/src/main/java/cn/iocoder/muse/module/ai/dal/dataobject/muse/MuseAiEvaluationSampleResultDO.java @@ -0,0 +1,41 @@ +package cn.iocoder.muse.module.ai.dal.dataobject.muse; + +import cn.iocoder.muse.framework.tenant.core.db.TenantBaseDO; +import cn.iocoder.muse.module.ai.dal.type.JsonbStringTypeHandler; +import com.baomidou.mybatisplus.annotation.IdType; +import com.baomidou.mybatisplus.annotation.TableField; +import com.baomidou.mybatisplus.annotation.TableId; +import com.baomidou.mybatisplus.annotation.TableName; +import lombok.Data; +import lombok.EqualsAndHashCode; +import lombok.ToString; + +import java.math.BigDecimal; + +/** + * Muse AI 离线评估单样本结果 DO。 + */ +@TableName(value = "muse_ai_evaluation_sample_result", autoResultMap = true) +@Data +@EqualsAndHashCode(callSuper = true) +@ToString(callSuper = true) +public class MuseAiEvaluationSampleResultDO extends TenantBaseDO { + + @TableId(type = IdType.AUTO) + private Long id; + + private String runId; + private String sampleId; + private Long ownerUserId; + private String status; + private BigDecimal score; + @TableField(typeHandler = JsonbStringTypeHandler.class) + private String metricSummary; + @TableField(typeHandler = JsonbStringTypeHandler.class) + private String inputSummary; + @TableField(typeHandler = JsonbStringTypeHandler.class) + private String outputSummary; + private String errorCode; + private String errorMessage; + +} diff --git a/muse-cloud/muse-module-ai/muse-module-ai-server/src/main/java/cn/iocoder/muse/module/ai/dal/mysql/muse/MuseAiEvaluationRunMapper.java b/muse-cloud/muse-module-ai/muse-module-ai-server/src/main/java/cn/iocoder/muse/module/ai/dal/mysql/muse/MuseAiEvaluationRunMapper.java index 7624da2d..f0771908 100644 --- a/muse-cloud/muse-module-ai/muse-module-ai-server/src/main/java/cn/iocoder/muse/module/ai/dal/mysql/muse/MuseAiEvaluationRunMapper.java +++ b/muse-cloud/muse-module-ai/muse-module-ai-server/src/main/java/cn/iocoder/muse/module/ai/dal/mysql/muse/MuseAiEvaluationRunMapper.java @@ -14,4 +14,8 @@ public interface MuseAiEvaluationRunMapper extends BaseMapperX { +} diff --git a/muse-cloud/muse-module-ai/muse-module-ai-server/src/main/java/cn/iocoder/muse/module/ai/dal/mysql/muse/MuseAiJobMapper.java b/muse-cloud/muse-module-ai/muse-module-ai-server/src/main/java/cn/iocoder/muse/module/ai/dal/mysql/muse/MuseAiJobMapper.java index 4fdc3418..887dfc1b 100644 --- a/muse-cloud/muse-module-ai/muse-module-ai-server/src/main/java/cn/iocoder/muse/module/ai/dal/mysql/muse/MuseAiJobMapper.java +++ b/muse-cloud/muse-module-ai/muse-module-ai-server/src/main/java/cn/iocoder/muse/module/ai/dal/mysql/muse/MuseAiJobMapper.java @@ -64,6 +64,33 @@ public interface MuseAiJobMapper extends BaseMapperX { """) MuseAiJobDO claimNextQueuedRuntimeJob(); + /** + * 原子领取下一条待执行质量评估 job。 + * + *

质量评估和在线生成调用不同执行器,不能混入 {@code ai_task_runtime} dispatcher;这里单独按 + * {@code job_type='evaluation'} 领取,领取时在同一 SQL 内推进到 running,避免多 worker 重复执行。

+ */ + @TenantIgnore + @Select(""" + UPDATE muse_ai_job + SET status = 'running', + started_at = COALESCE(started_at, CURRENT_TIMESTAMP), + update_time = CURRENT_TIMESTAMP + WHERE id = ( + SELECT id + FROM muse_ai_job + WHERE job_type = 'evaluation' + AND status = 'queued' + AND (next_retry_at IS NULL OR next_retry_at <= CURRENT_TIMESTAMP) + AND deleted = FALSE + ORDER BY create_time ASC, id ASC + FOR UPDATE SKIP LOCKED + LIMIT 1 + ) + RETURNING * + """) + MuseAiJobDO claimNextQueuedEvaluationJob(); + /** * 在同一 retry group 内查找未终止 retry,调用方已锁定源 Job,因此这里用于锁内防重复创建。 */ diff --git a/muse-cloud/muse-module-ai/muse-module-ai-server/src/main/java/cn/iocoder/muse/module/ai/framework/ai/config/MuseAiProperties.java b/muse-cloud/muse-module-ai/muse-module-ai-server/src/main/java/cn/iocoder/muse/module/ai/framework/ai/config/MuseAiProperties.java index 61f868a3..c10c026d 100644 --- a/muse-cloud/muse-module-ai/muse-module-ai-server/src/main/java/cn/iocoder/muse/module/ai/framework/ai/config/MuseAiProperties.java +++ b/muse-cloud/muse-module-ai/muse-module-ai-server/src/main/java/cn/iocoder/muse/module/ai/framework/ai/config/MuseAiProperties.java @@ -73,6 +73,11 @@ public class MuseAiProperties { */ private Events events = new Events(); + /** + * Muse 离线质量评估配置。 + */ + private Evaluation evaluation = new Evaluation(); + @Data public static class Gemini { @@ -251,4 +256,27 @@ public class MuseAiProperties { } + @Data + public static class Evaluation { + + /** + * 质量评估 worker 配置。 + */ + private Worker worker = new Worker(); + + @Data + public static class Worker { + + /** + * 是否启用质量评估后台 worker;默认启用。 + * + *

当前执行器只处理合成/脱敏数据集,不调用外部 LLM,不读取私有正文;生产如需冻结历史 queued + * evaluation job,可显式设为 {@code false}。

+ */ + private boolean enabled = true; + + } + + } + } diff --git a/muse-cloud/muse-module-ai/muse-module-ai-server/src/test/java/cn/iocoder/muse/module/ai/application/muse/MuseQualityEvaluationWorkerTest.java b/muse-cloud/muse-module-ai/muse-module-ai-server/src/test/java/cn/iocoder/muse/module/ai/application/muse/MuseQualityEvaluationWorkerTest.java new file mode 100644 index 00000000..50c4df4b --- /dev/null +++ b/muse-cloud/muse-module-ai/muse-module-ai-server/src/test/java/cn/iocoder/muse/module/ai/application/muse/MuseQualityEvaluationWorkerTest.java @@ -0,0 +1,216 @@ +package cn.iocoder.muse.module.ai.application.muse; + +import cn.iocoder.muse.framework.common.util.json.JsonUtils; +import cn.iocoder.muse.framework.test.core.ut.BaseMockitoUnitTest; +import cn.iocoder.muse.framework.tenant.core.context.TenantContextHolder; +import cn.iocoder.muse.module.ai.dal.dataobject.muse.MuseAiEvaluationRunDO; +import cn.iocoder.muse.module.ai.dal.dataobject.muse.MuseAiEvaluationSampleResultDO; +import cn.iocoder.muse.module.ai.dal.dataobject.muse.MuseAiJobDO; +import cn.iocoder.muse.module.ai.dal.dataobject.muse.MuseQualityPolicyVersionDO; +import cn.iocoder.muse.module.ai.dal.mysql.muse.MuseAiEvaluationRunMapper; +import cn.iocoder.muse.module.ai.dal.mysql.muse.MuseAiEvaluationSampleResultMapper; +import cn.iocoder.muse.module.ai.dal.mysql.muse.MuseAiJobMapper; +import cn.iocoder.muse.module.ai.dal.mysql.muse.MuseQualityPolicyVersionMapper; +import cn.iocoder.muse.module.ai.framework.ai.config.MuseAiProperties; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; +import org.mockito.InjectMocks; +import org.mockito.Mock; + +import java.util.List; +import java.util.Map; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.argThat; +import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +/** + * 质量评估 worker / executor 测试。 + */ +class MuseQualityEvaluationWorkerTest extends BaseMockitoUnitTest { + + @InjectMocks + private MuseQualityEvaluationExecutor executor; + @Mock + private MuseAiJobMapper jobMapper; + @Mock + private MuseAiEvaluationRunMapper evaluationRunMapper; + @Mock + private MuseAiEvaluationSampleResultMapper sampleResultMapper; + @Mock + private MuseQualityPolicyVersionMapper qualityPolicyVersionMapper; + @Mock + private MuseAiAuditService auditService; + + @AfterEach + void clearTenantContext() { + TenantContextHolder.clear(); + } + + @Test + void should_dispatchEnabledEvaluationJobAndRestoreTenantContext() { + MuseAiProperties properties = enabledProperties(true); + MuseQualityEvaluationWorker worker = new MuseQualityEvaluationWorker(jobMapper, executor, properties); + MuseAiJobDO claimed = job("8001", 4001L, 100L); + when(jobMapper.claimNextQueuedEvaluationJob()).thenReturn(claimed); + when(jobMapper.selectById(4001L)).thenReturn(claimed); + when(evaluationRunMapper.selectByJobId("8001")).thenReturn(run("7001", "dataset://quality-v1", 3)); + when(qualityPolicyVersionMapper.selectByPolicyKeyAndVersion("novel_quality", 1)) + .thenReturn(policyVersion("[{\"name\":\"coherence\",\"threshold\":0.7,\"weight\":1.0}]")); + + int dispatched = worker.dispatchOnce(); + + assertEquals(1, dispatched); + verify(jobMapper).claimNextQueuedEvaluationJob(); + verify(evaluationRunMapper).updateById(argThat((MuseAiEvaluationRunDO update) -> + Long.valueOf(3001L).equals(update.getId()) && "running".equals(update.getStatus()))); + ArgumentCaptor sampleCaptor = + ArgumentCaptor.forClass(MuseAiEvaluationSampleResultDO.class); + verify(sampleResultMapper, org.mockito.Mockito.times(3)).insert(sampleCaptor.capture()); + assertTrue(sampleCaptor.getAllValues().stream().allMatch(sample -> + "7001".equals(sample.getRunId()) + && sample.getInputSummary().contains("datasetReferenceHash") + && !sample.getInputSummary().contains("dataset://quality-v1"))); + verify(evaluationRunMapper).updateById(argThat((MuseAiEvaluationRunDO update) -> + Long.valueOf(3001L).equals(update.getId()) + && "completed".equals(update.getStatus()) + && update.getMetricSummary().contains("\"deterministic-local-v1\"") + && Integer.valueOf(1).equals(update.getFailedCount()))); + verify(jobMapper).updateById(argThat((MuseAiJobDO update) -> + Long.valueOf(4001L).equals(update.getId()) + && "completed".equals(update.getStatus()) + && update.getResultSummary().contains("\"overallScore\""))); + verify(auditService).record(argThat(req -> + "executeQualityEvaluationRun".equals(req.getOperationId()) + && "worker".equals(req.getSide()) + && "succeeded".equals(req.getStatus()))); + assertNull(TenantContextHolder.getTenantId(), "质量评估 worker 执行结束后必须恢复租户上下文"); + } + + @Test + void should_notDispatchWhenWorkerDisabled() { + MuseQualityEvaluationWorker worker = new MuseQualityEvaluationWorker(jobMapper, executor, enabledProperties(false)); + + int dispatched = worker.dispatchOnce(); + + assertEquals(0, dispatched); + verify(jobMapper, never()).claimNextQueuedEvaluationJob(); + } + + @Test + void should_failRunAndJobWhenDatasetReferenceIsPrivate() { + TenantContextHolder.setTenantId(100L); + MuseAiJobDO job = job("8001", 4001L, 100L); + when(jobMapper.selectById(4001L)).thenReturn(job); + when(evaluationRunMapper.selectByJobId("8001")).thenReturn(run("7001", "work://private/100", 2)); + executor.executeEvaluationJob(4001L); + + verify(sampleResultMapper, never()).insert(any(MuseAiEvaluationSampleResultDO.class)); + verify(evaluationRunMapper).updateById(argThat((MuseAiEvaluationRunDO update) -> + Long.valueOf(3001L).equals(update.getId()) + && "failed".equals(update.getStatus()) + && "AI_QUALITY_EVALUATION_DATASET_BLOCKED".equals(update.getErrorCode()) + && update.getErrorMessage() != null + && !update.getErrorMessage().contains("work://private"))); + verify(jobMapper).updateById(argThat((MuseAiJobDO update) -> + Long.valueOf(4001L).equals(update.getId()) + && "failed".equals(update.getStatus()) + && "AI_QUALITY_EVALUATION_DATASET_BLOCKED".equals(update.getErrorCode()))); + verify(auditService).record(argThat(req -> + "failed".equals(req.getStatus()) + && "AI_QUALITY_EVALUATION_DATASET_BLOCKED".equals(req.getErrorCode()) + && req.getRequestSummary().contains("datasetReferenceHash") + && !req.getRequestSummary().contains("work://private"))); + } + + @Test + void should_failClosedWhenExecutorThrowsAfterClaimedJob() { + MuseQualityEvaluationExecutor throwingExecutor = org.mockito.Mockito.mock(MuseQualityEvaluationExecutor.class); + MuseQualityEvaluationWorker worker = new MuseQualityEvaluationWorker(jobMapper, throwingExecutor, enabledProperties(true)); + MuseAiJobDO claimed = job("8001", 4001L, 100L); + when(jobMapper.claimNextQueuedEvaluationJob()).thenReturn(claimed); + doThrow(new IllegalStateException("db unavailable")).when(throwingExecutor).executeEvaluationJob(4001L); + + int dispatched = worker.dispatchOnce(); + + assertEquals(0, dispatched); + verify(jobMapper).updateById(argThat((MuseAiJobDO update) -> + Long.valueOf(4001L).equals(update.getId()) + && "failed".equals(update.getStatus()) + && "AI_QUALITY_EVALUATION_DISPATCH_FAILED".equals(update.getErrorCode()) + && Boolean.FALSE.equals(update.getRetryable()))); + assertNull(TenantContextHolder.getTenantId(), "异常路径也必须恢复租户上下文"); + } + + @Test + void should_keepCompletedStateWhenAuditWriteFails() { + TenantContextHolder.setTenantId(100L); + MuseAiJobDO job = job("8001", 4001L, 100L); + when(jobMapper.selectById(4001L)).thenReturn(job); + when(evaluationRunMapper.selectByJobId("8001")).thenReturn(run("7001", "redacted-fixture-v1", 1)); + when(qualityPolicyVersionMapper.selectByPolicyKeyAndVersion("novel_quality", 1)) + .thenReturn(policyVersion("[{\"name\":\"coherence\",\"threshold\":0.7,\"weight\":1.0}]")); + doThrow(new IllegalStateException("audit db unavailable")).when(auditService).record(any()); + + executor.executeEvaluationJob(4001L); + + verify(evaluationRunMapper).updateById(argThat((MuseAiEvaluationRunDO update) -> + Long.valueOf(3001L).equals(update.getId()) && "completed".equals(update.getStatus()))); + verify(jobMapper).updateById(argThat((MuseAiJobDO update) -> + Long.valueOf(4001L).equals(update.getId()) && "completed".equals(update.getStatus()))); + } + + private static MuseAiProperties enabledProperties(boolean enabled) { + MuseAiProperties properties = new MuseAiProperties(); + MuseAiProperties.Evaluation evaluation = new MuseAiProperties.Evaluation(); + MuseAiProperties.Evaluation.Worker worker = new MuseAiProperties.Evaluation.Worker(); + worker.setEnabled(enabled); + evaluation.setWorker(worker); + properties.setEvaluation(evaluation); + return properties; + } + + private static MuseAiJobDO job(String jobId, Long id, Long tenantId) { + MuseAiJobDO job = new MuseAiJobDO(); + job.setId(id); + job.setJobId(jobId); + job.setCommandId("cmd-run-1"); + job.setRequestHash("hash-run"); + job.setActorUserId(1001L); + job.setOwnerUserId(1001L); + job.setJobType("evaluation"); + job.setStatus("running"); + job.setTenantId(tenantId); + return job; + } + + private static MuseAiEvaluationRunDO run(String runId, String datasetReference, int sampleCount) { + MuseAiEvaluationRunDO run = new MuseAiEvaluationRunDO(); + run.setId(3001L); + run.setRunId(runId); + run.setJobId("8001"); + run.setPolicyKey("novel_quality"); + run.setPolicyVersionNo(1); + run.setOwnerUserId(1001L); + run.setSampleCount(sampleCount); + run.setMetricSummary(JsonUtils.toJsonString(Map.of("datasetReference", datasetReference))); + return run; + } + + private static MuseQualityPolicyVersionDO policyVersion(String thresholdSnapshot) { + MuseQualityPolicyVersionDO version = new MuseQualityPolicyVersionDO(); + version.setPolicyKey("novel_quality"); + version.setVersionNo(1); + version.setThresholdSnapshot(thresholdSnapshot); + return version; + } + +} diff --git a/muse-cloud/muse-module-ai/muse-module-ai-server/src/test/java/cn/iocoder/muse/module/ai/application/muse/MuseQualityPolicyServiceTest.java b/muse-cloud/muse-module-ai/muse-module-ai-server/src/test/java/cn/iocoder/muse/module/ai/application/muse/MuseQualityPolicyServiceTest.java index d2dea528..69b16645 100644 --- a/muse-cloud/muse-module-ai/muse-module-ai-server/src/test/java/cn/iocoder/muse/module/ai/application/muse/MuseQualityPolicyServiceTest.java +++ b/muse-cloud/muse-module-ai/muse-module-ai-server/src/test/java/cn/iocoder/muse/module/ai/application/muse/MuseQualityPolicyServiceTest.java @@ -257,7 +257,8 @@ class MuseQualityPolicyServiceTest extends BaseMockitoUnitTest { qualityPolicyService.startEvaluationRun(request, 1001L); - verify(evaluationRunMapper).insert(argThat((MuseAiEvaluationRunDO run) -> !"completed".equals(run.getStatus()))); + verify(evaluationRunMapper).insert(argThat((MuseAiEvaluationRunDO run) -> !"completed".equals(run.getStatus()) + && Integer.valueOf(20).equals(run.getSampleCount()))); verify(jobMapper).insert(argThat((MuseAiJobDO job) -> !"completed".equals(job.getStatus()))); } diff --git a/muse-cloud/muse-module-ai/muse-module-ai-server/src/test/java/cn/iocoder/muse/module/ai/dal/mysql/muse/MuseAiJobMapperTest.java b/muse-cloud/muse-module-ai/muse-module-ai-server/src/test/java/cn/iocoder/muse/module/ai/dal/mysql/muse/MuseAiJobMapperTest.java index 748c5f8e..cf69a2b5 100644 --- a/muse-cloud/muse-module-ai/muse-module-ai-server/src/test/java/cn/iocoder/muse/module/ai/dal/mysql/muse/MuseAiJobMapperTest.java +++ b/muse-cloud/muse-module-ai/muse-module-ai-server/src/test/java/cn/iocoder/muse/module/ai/dal/mysql/muse/MuseAiJobMapperTest.java @@ -58,6 +58,7 @@ class MuseAiJobMapperTest extends BaseMockitoUnitTest { void should_defineRowLockAndActiveRetryQueriesForJobCommands() throws Exception { Method lockMethod = MuseAiJobMapper.class.getMethod("selectByIdForUpdate", Long.class, Long.class); Method claimMethod = MuseAiJobMapper.class.getMethod("claimNextQueuedRuntimeJob"); + Method evaluationClaimMethod = MuseAiJobMapper.class.getMethod("claimNextQueuedEvaluationJob"); Method activeRetryMethod = MuseAiJobMapper.class.getMethod("selectActiveRetryByRetryGroupId", Long.class, String.class, Long.class); Method insertActiveRetryMethod = MuseAiJobMapper.class.getMethod("insertActiveRetryIfAbsent", @@ -65,6 +66,7 @@ class MuseAiJobMapperTest extends BaseMockitoUnitTest { String lockSql = normalizeSql(lockMethod.getAnnotation(org.apache.ibatis.annotations.Select.class).value()[0]); String claimSql = normalizeSql(claimMethod.getAnnotation(org.apache.ibatis.annotations.Select.class).value()[0]); + String evaluationClaimSql = normalizeSql(evaluationClaimMethod.getAnnotation(org.apache.ibatis.annotations.Select.class).value()[0]); String activeRetrySql = normalizeSql(activeRetryMethod.getAnnotation(org.apache.ibatis.annotations.Select.class).value()[0]); String insertActiveRetrySql = normalizeSql(insertActiveRetryMethod.getAnnotation( org.apache.ibatis.annotations.Insert.class).value()[0]); @@ -75,8 +77,16 @@ class MuseAiJobMapperTest extends BaseMockitoUnitTest { "claim SQL 必须在同一原子语句内把 queued job 推进到 running"); assertTrue(claimSql.contains("job_type = 'ai_task_runtime'"), "P1R-4 dispatcher 只能领取 AI task runtime job,不能误执行 agent test 或其它 owner job"); + assertTrue(evaluationClaimSql.contains("for update skip locked"), + "质量评估 worker 必须用数据库行锁领取 queued job,避免多 worker 重复执行"); + assertTrue(evaluationClaimSql.contains("job_type = 'evaluation'"), + "质量评估 worker 只能领取 evaluation job,不能混入在线生成 runtime job"); + assertTrue(evaluationClaimSql.contains("status = 'queued'") && evaluationClaimSql.contains("status = 'running'"), + "质量评估 claim SQL 必须在同一原子语句内把 queued job 推进到 running"); assertTrue(claimSql.contains("next_retry_at is null") && claimSql.contains("next_retry_at <= current_timestamp"), "claim SQL 必须尊重 retry backoff 时间,不能提前重跑"); + assertTrue(evaluationClaimSql.contains("next_retry_at is null") && evaluationClaimSql.contains("next_retry_at <= current_timestamp"), + "质量评估 claim SQL 必须尊重 retry backoff 时间,不能提前重跑"); assertTrue(activeRetrySql.contains("retry_group_id")); assertTrue(activeRetrySql.contains("operation_id = 'adminretryjob'"), "active retry 查询只能复用 Job retry 命令创建的 retry job,不能误把普通源 job 当作 retry job");