From 48c750bbd78a87e8372f926dcc86f18eafccaef7 Mon Sep 17 00:00:00 2001 From: zizi Date: Wed, 3 Jun 2026 14:10:21 +0800 Subject: [PATCH] =?UTF-8?q?test(p1r):=20=E8=A1=A5=E9=BD=90=20AI=20?= =?UTF-8?q?=E4=B8=8E=20Knowledge=20=E5=A4=96=E9=83=A8=E7=AB=AF=E5=88=B0?= =?UTF-8?q?=E7=AB=AF=E9=AA=8C=E6=94=B6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../src/main/resources/application.yaml | 4 +- .../muse/MuseKnowledgeDocumentService.java | 2 +- .../MuseKnowledgeProcessingTaskService.java | 6 +- .../MuseKnowledgeDocumentServiceTest.java | 1 + .../src/main/resources/application.yaml | 4 +- .../P1rAiRuntimeEndToEndLiveAcceptanceIT.java | 1077 +++++++++++++++++ ...wledgeRuntimeEndToEndLiveAcceptanceIT.java | 1071 ++++++++++++++++ 7 files changed, 2158 insertions(+), 7 deletions(-) create mode 100644 muse-cloud/muse-server/src/test/java/cn/iocoder/muse/server/framework/api/P1rAiRuntimeEndToEndLiveAcceptanceIT.java create mode 100644 muse-cloud/muse-server/src/test/java/cn/iocoder/muse/server/framework/api/P1rKnowledgeRuntimeEndToEndLiveAcceptanceIT.java diff --git a/muse-cloud/muse-module-ai/muse-module-ai-server/src/main/resources/application.yaml b/muse-cloud/muse-module-ai/muse-module-ai-server/src/main/resources/application.yaml index ce62a1ed..90555385 100644 --- a/muse-cloud/muse-module-ai/muse-module-ai-server/src/main/resources/application.yaml +++ b/muse-cloud/muse-module-ai/muse-module-ai-server/src/main/resources/application.yaml @@ -73,8 +73,8 @@ mybatis-plus: # id-type: AUTO # 自增 ID,适合 MySQL 等直接自增的数据库 # id-type: INPUT # 用户输入 ID,适合 Oracle、PostgreSQL、Kingbase、DB2、H2 数据库 # id-type: ASSIGN_ID # 分配 ID,默认使用雪花算法。注意,Oracle、PostgreSQL、Kingbase、DB2、H2 数据库时,需要去除实体类上的 @KeySequence 注解 - logic-delete-value: 1 # 逻辑已删除值(默认为 1) - logic-not-delete-value: 0 # 逻辑未删除值(默认为 0) + logic-delete-value: true # 逻辑已删除值,匹配 PostgreSQL boolean deleted 字段 + logic-not-delete-value: false # 逻辑未删除值,匹配 PostgreSQL boolean deleted 字段 banner: false # 关闭控制台的 Banner 打印 type-aliases-package: ${muse.info.base-package}.module.*.dal.dataobject encryptor: diff --git a/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/main/java/cn/iocoder/muse/module/knowledge/application/muse/MuseKnowledgeDocumentService.java b/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/main/java/cn/iocoder/muse/module/knowledge/application/muse/MuseKnowledgeDocumentService.java index c13854a2..73ee4aa4 100644 --- a/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/main/java/cn/iocoder/muse/module/knowledge/application/muse/MuseKnowledgeDocumentService.java +++ b/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/main/java/cn/iocoder/muse/module/knowledge/application/muse/MuseKnowledgeDocumentService.java @@ -293,7 +293,7 @@ public class MuseKnowledgeDocumentService { return new UploadOutcome(ragflowContext, "failed"); } processingTaskService.markRagflowParseAccepted(ragflowContext.version().getId(), ragflowContext.task().getTaskId(), - uploadResult.externalId(), parseResult); + ragflowContext.ragflowDatasetId(), uploadResult.externalId(), parseResult); return new UploadOutcome(ragflowContext, "processing"); } diff --git a/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/main/java/cn/iocoder/muse/module/knowledge/application/muse/MuseKnowledgeProcessingTaskService.java b/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/main/java/cn/iocoder/muse/module/knowledge/application/muse/MuseKnowledgeProcessingTaskService.java index 87fbffad..00f48001 100644 --- a/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/main/java/cn/iocoder/muse/module/knowledge/application/muse/MuseKnowledgeProcessingTaskService.java +++ b/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/main/java/cn/iocoder/muse/module/knowledge/application/muse/MuseKnowledgeProcessingTaskService.java @@ -177,13 +177,15 @@ public class MuseKnowledgeProcessingTaskService { } @Transactional(propagation = Propagation.REQUIRES_NEW, rollbackFor = Exception.class) - public void markRagflowParseAccepted(Long documentVersionId, String taskId, String ragflowDocumentId, - RagFlowKnowledgeRuntimeClient.RuntimeResult result) { + public void markRagflowParseAccepted(Long documentVersionId, String taskId, String ragflowDatasetId, + String ragflowDocumentId, RagFlowKnowledgeRuntimeClient.RuntimeResult result) { MuseKnowledgeProcessingTaskDO task = processingTaskMapper.selectByTaskId(taskId); if (task != null) { task.setStatus("parsing"); task.setParseStatus("processing"); task.setIndexStatus("pending"); + // 新建 dataset 的上传链路必须把外部 datasetId 同步固化到 task,后续轮询、补偿和 live 取证都依赖这条事实。 + task.setRagflowDatasetId(ragflowDatasetId); task.setRagflowDocumentId(ragflowDocumentId); task.setRetryable(false); task.setProgressMessage("RAGFlow parse accepted; polling required"); diff --git a/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/test/java/cn/iocoder/muse/module/knowledge/application/muse/MuseKnowledgeDocumentServiceTest.java b/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/test/java/cn/iocoder/muse/module/knowledge/application/muse/MuseKnowledgeDocumentServiceTest.java index aaabc537..62534df5 100644 --- a/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/test/java/cn/iocoder/muse/module/knowledge/application/muse/MuseKnowledgeDocumentServiceTest.java +++ b/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/test/java/cn/iocoder/muse/module/knowledge/application/muse/MuseKnowledgeDocumentServiceTest.java @@ -377,6 +377,7 @@ class MuseKnowledgeDocumentServiceTest extends BaseMockitoUnitTest { && RagFlowKnowledgeRuntimeClient.Status.ACCEPTED.equals(call.getStatus()))); verify(processingTaskMapper).updateById(org.mockito.ArgumentMatchers.argThat(task -> "parsing".equals(task.getStatus()) + && "rag-dataset-1".equals(task.getRagflowDatasetId()) && "rag-doc-1".equals(task.getRagflowDocumentId()) && "processing".equals(task.getResultSummary()) == false)); } diff --git a/muse-cloud/muse-server/src/main/resources/application.yaml b/muse-cloud/muse-server/src/main/resources/application.yaml index 83b45907..321110e6 100644 --- a/muse-cloud/muse-server/src/main/resources/application.yaml +++ b/muse-cloud/muse-server/src/main/resources/application.yaml @@ -95,8 +95,8 @@ mybatis-plus: # id-type: AUTO # 自增 ID,适合 MySQL 等直接自增的数据库 # id-type: INPUT # 用户输入 ID,适合 Oracle、PostgreSQL、Kingbase、DB2、H2 数据库 # id-type: ASSIGN_ID # 分配 ID,默认使用雪花算法。注意,Oracle、PostgreSQL、Kingbase、DB2、H2 数据库时,需要去除实体类上的 @KeySequence 注解 - logic-delete-value: 1 # 逻辑已删除值(默认为 1) - logic-not-delete-value: 0 # 逻辑未删除值(默认为 0) + logic-delete-value: true # 逻辑已删除值,匹配 PostgreSQL boolean deleted 字段 + logic-not-delete-value: false # 逻辑未删除值,匹配 PostgreSQL boolean deleted 字段 banner: false # 关闭控制台的 Banner 打印 type-aliases-package: ${muse.info.base-package}.module.*.dal.dataobject encryptor: diff --git a/muse-cloud/muse-server/src/test/java/cn/iocoder/muse/server/framework/api/P1rAiRuntimeEndToEndLiveAcceptanceIT.java b/muse-cloud/muse-server/src/test/java/cn/iocoder/muse/server/framework/api/P1rAiRuntimeEndToEndLiveAcceptanceIT.java new file mode 100644 index 00000000..76fcd8c1 --- /dev/null +++ b/muse-cloud/muse-server/src/test/java/cn/iocoder/muse/server/framework/api/P1rAiRuntimeEndToEndLiveAcceptanceIT.java @@ -0,0 +1,1077 @@ +package cn.iocoder.muse.server.framework.api; + +import cn.hutool.extra.spring.SpringUtil; +import cn.iocoder.muse.framework.common.enums.UserTypeEnum; +import cn.iocoder.muse.framework.common.util.json.JsonUtils; +import cn.iocoder.muse.framework.datasource.config.MuseDataSourceAutoConfiguration; +import cn.iocoder.muse.framework.mybatis.config.MuseMybatisAutoConfiguration; +import cn.iocoder.muse.framework.security.core.LoginUser; +import cn.iocoder.muse.framework.security.core.util.SecurityFrameworkUtils; +import cn.iocoder.muse.framework.tenant.core.context.TenantContextHolder; +import cn.iocoder.muse.framework.tenant.core.util.TenantUtils; +import cn.iocoder.muse.module.ai.application.muse.MuseAiAuditServiceImpl; +import cn.iocoder.muse.module.ai.application.muse.MuseAiCommandServiceImpl; +import cn.iocoder.muse.module.ai.application.muse.MuseAiRuntimeCallRecorder; +import cn.iocoder.muse.module.ai.application.muse.MuseAiRuntimeJobExecutor; +import cn.iocoder.muse.module.ai.application.muse.MuseAiRuntimePayloadStore; +import cn.iocoder.muse.module.ai.application.muse.MuseAiRuntimePolicyProvider; +import cn.iocoder.muse.module.ai.application.muse.MuseAiRuntimeProjectionService; +import cn.iocoder.muse.module.ai.application.muse.MuseAiRuntimeServiceImpl; +import cn.iocoder.muse.module.ai.application.muse.MuseAiTaskService; +import cn.iocoder.muse.module.ai.application.muse.MuseAiTaskServiceImpl; +import cn.iocoder.muse.module.ai.application.muse.facade.MuseAiRuntimeClient; +import cn.iocoder.muse.module.ai.application.muse.facade.MuseContentWorkOwnerFacade; +import cn.iocoder.muse.module.ai.application.muse.facade.RealNewApiMuseAiRuntimeClient; +import cn.iocoder.muse.module.ai.application.muse.facade.SecurityRuntimePermissionFacade; +import cn.iocoder.muse.module.ai.controller.app.muse.vo.CreateAiTaskReqVO; +import cn.iocoder.muse.module.ai.controller.app.muse.vo.CreateAiTaskRespVO; +import cn.iocoder.muse.module.ai.dal.dataobject.muse.MuseAgentDO; +import cn.iocoder.muse.module.ai.dal.dataobject.muse.MuseAgentVersionDO; +import cn.iocoder.muse.module.ai.dal.dataobject.muse.MuseAiGenerationDO; +import cn.iocoder.muse.module.ai.dal.dataobject.muse.MuseAiJobDO; +import cn.iocoder.muse.module.ai.dal.dataobject.muse.MuseAiRuntimeCallDO; +import cn.iocoder.muse.module.ai.dal.mysql.muse.MuseAgentMapper; +import cn.iocoder.muse.module.ai.dal.mysql.muse.MuseAgentVersionMapper; +import cn.iocoder.muse.module.ai.dal.mysql.muse.MuseAiAuditMapper; +import cn.iocoder.muse.module.ai.dal.mysql.muse.MuseAiCommandMapper; +import cn.iocoder.muse.module.ai.dal.mysql.muse.MuseAiGenerationMapper; +import cn.iocoder.muse.module.ai.dal.mysql.muse.MuseAiJobMapper; +import cn.iocoder.muse.module.ai.dal.mysql.muse.MuseAiRuntimeCallMapper; +import cn.iocoder.muse.module.ai.dal.mysql.muse.MuseAiRuntimePayloadMapper; +import cn.iocoder.muse.module.ai.dal.mysql.muse.MuseAiSuggestionMapper; +import cn.iocoder.muse.module.ai.dal.mysql.muse.MuseAiTaskEventMapper; +import cn.iocoder.muse.module.ai.domain.muse.MuseAiRuntimePermissionGuard; +import cn.iocoder.muse.module.ai.framework.ai.config.MuseAiProperties; +import com.fasterxml.jackson.databind.JsonNode; +import com.baomidou.mybatisplus.autoconfigure.MybatisPlusAutoConfiguration; +import com.github.yulichang.autoconfigure.MybatisPlusJoinAutoConfiguration; +import org.flywaydb.core.Flyway; +import org.flywaydb.core.api.MigrationInfo; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.Test; +import org.slf4j.LoggerFactory; +import org.springframework.boot.WebApplicationType; +import org.springframework.boot.autoconfigure.jdbc.DataSourceAutoConfiguration; +import org.springframework.boot.autoconfigure.jdbc.DataSourceTransactionManagerAutoConfiguration; +import org.springframework.boot.builder.SpringApplicationBuilder; +import org.springframework.boot.test.context.TestConfiguration; +import org.springframework.context.ConfigurableApplicationContext; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Import; +import org.springframework.context.annotation.Primary; +import org.springframework.core.env.MapPropertySource; +import org.springframework.mock.web.MockHttpServletRequest; +import org.springframework.security.core.context.SecurityContextHolder; +import org.springframework.util.StringUtils; + +import java.net.URI; +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.nio.file.Path; +import java.security.MessageDigest; +import java.security.NoSuchAlgorithmException; +import java.sql.Connection; +import java.sql.DriverManager; +import java.sql.PreparedStatement; +import java.sql.ResultSet; +import java.sql.SQLException; +import java.time.Instant; +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.Objects; +import java.util.Set; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.junit.jupiter.api.Assumptions.assumeTrue; + +/** + * P1R-4 AI Muse 业务入口到真实 New-API / PostgreSQL 的 opt-in live acceptance。 + * + *

默认必须跳过,避免普通测试误连真实数据库或外呼 New-API。只有显式设置 + * {@code MUSE_P1R_EXTERNAL_ACCEPTANCE=true} 时,本测试才会清理 PostgreSQL {@code _test} 库、 + * 迁移到 V13、seed 最小 agent/version,并通过 {@link MuseAiTaskServiceImpl} 创建 AI task 后执行 + * {@link MuseAiRuntimeJobExecutor}。

+ */ +class P1rAiRuntimeEndToEndLiveAcceptanceIT { + + private static final String ACCEPTANCE_ENV = "MUSE_P1R_EXTERNAL_ACCEPTANCE"; + private static final String JDBC_URL_PROPERTY = "p1r.flyway.url"; + private static final String JDBC_USER_PROPERTY = "p1r.flyway.user"; + private static final String JDBC_URL_ENV = "P1R_FLYWAY_URL"; + private static final String JDBC_USER_ENV = "P1R_FLYWAY_USER"; + private static final String FLYWAY_LOCATION_PROPERTY = "p1r.flyway.locations"; + private static final String POSTGRESQL_JDBC_PREFIX = "jdbc:postgresql://"; + private static final String JDBC_URI_PREFIX = "jdbc:"; + private static final String TARGET_VERSION = "13"; + // 固定 test-only AES 密钥,供 EncryptTypeHandler 写入 live payload 时读取;不从用户环境变量取值。 + private static final String TEST_ENCRYPTOR_PASSWORD = "p1rLiveEncrypt16"; + private static final Long TENANT_ID = 100L; + private static final Long OWNER_USER_ID = 2001L; + private static final Long WORK_ID = 9001001L; + private static final String MODEL = "MiniMax-M2.5"; + private static final String CHAT_COMPLETIONS_PATH = "/v1/chat/completions"; + private static final Set CREDENTIAL_QUERY_KEYS = Set.of( + "user", "username", "password", "pass", "pwd", "sslpassword", "ssl_password", + "token", "secret", "api_key", "apikey", "bearer", "access_token", "refresh_token"); + + @AfterEach + void clearRuntimeContexts() { + // live acceptance 在同一 JVM 内手工设置租户和登录态,测试结束必须同时清理,避免污染后续测试。 + SecurityContextHolder.clearContext(); + TenantContextHolder.clear(); + } + + @Test + void shouldMaskNonPostgresqlJdbcAuthorityInFailureOutput() { + String masked = maskedUrl("jdbc:mysql://user:password@host:3306/db_test?password=x"); + + assertFalse(masked.contains("user"), () -> "非 PostgreSQL JDBC URL 失败摘要不能泄露 userinfo: " + masked); + assertFalse(masked.contains("password"), () -> "非 PostgreSQL JDBC URL 失败摘要不能泄露密码: " + masked); + assertFalse(masked.contains("host"), () -> "非 PostgreSQL JDBC URL 失败摘要不能泄露真实 host: " + masked); + } + + @Test + void shouldMaskMalformedPostgresqlJdbcParseFailureOutput() { + String url = "jdbc:postgresql://user:password@bad host/muse_local_test?password=x"; + + AssertionError error = assertThrows(AssertionError.class, () -> assertSafeJdbcUrl(url)); + String message = error.getMessage(); + + assertTrue(message.contains("") || message.contains(""), + () -> "malformed PostgreSQL JDBC URL 解析失败消息必须包含脱敏摘要: " + message); + assertFalse(message.contains("user"), () -> "解析失败消息不能泄露 userinfo: " + message); + assertFalse(message.contains("password"), () -> "解析失败消息不能泄露密码或凭据参数: " + message); + assertFalse(message.contains("bad host"), () -> "解析失败消息不能泄露真实 host: " + message); + } + + @Test + void shouldCreateAiTaskExecuteRuntimeJobAndPersistDbFactsWithRealPostgresqlAndNewApi() { + assumeTrue(externalAcceptanceEnabled(), + "P1R 外部端到端验收未启用,设置 MUSE_P1R_EXTERNAL_ACCEPTANCE=true 后才运行"); + + LiveSettings settings = LiveSettings.fromEnvironment(); + DatabaseEngineFacts dbEngine = migrateIsolatedTestDatabase(settings); + + try (ConfigurableApplicationContext context = liveContext(settings)) { + setLiveLoginUser(); + LiveResult result = TenantUtils.execute(TENANT_ID, () -> executeMuseAiBusinessFlow(context)); + DbFacts dbFacts = readDbFacts(settings, result); + assertDbFacts(result, dbFacts); + + System.out.println(JsonUtils.toJsonString(redactedEvidence(settings, result, dbFacts, dbEngine))); + } + } + + private void setLiveLoginUser() { + // 该 live test 直接调用 service / executor,没有 HTTP 认证过滤器;这里补齐生产审计字段依赖的登录态。 + SecurityFrameworkUtils.setLoginUser(liveLoginUser(), new MockHttpServletRequest()); + } + + private LoginUser liveLoginUser() { + LoginUser loginUser = new LoginUser(); + loginUser.setId(OWNER_USER_ID); + loginUser.setUserType(UserTypeEnum.MEMBER.getValue()); + loginUser.setTenantId(TENANT_ID); + loginUser.setVisitTenantId(TENANT_ID); + return loginUser; + } + + private LiveResult executeMuseAiBusinessFlow(ConfigurableApplicationContext context) { + Long agentId = seedAgent(context); + CreateAiTaskRespVO response = context.getBean(MuseAiTaskService.class) + .createAiTask(OWNER_USER_ID, createTaskRequest(agentId)); + assertNotNull(response.getTaskId(), "createAiTask 必须返回 taskId"); + assertNotNull(response.getJobId(), "createAiTask 必须返回 jobId"); + + MuseAiRuntimeClient.RuntimeResponse runtimeResponse = context.getBean(MuseAiRuntimeJobExecutor.class) + .executeAiTaskJob(response.getTaskId(), response.getJobId()); + assertNotNull(runtimeResponse, "runtime executor 必须返回终态响应"); + assertNotNull(runtimeResponse.result(), + () -> "P1R-4 live acceptance 必须真实调用 New-API 成功,failure=" + failureSummary(runtimeResponse.failure())); + assertTrue(runtimeResponse.failure() == null, + () -> "P1R-4 live acceptance 主测试不接受 New-API failure,failure=" + failureSummary(runtimeResponse.failure())); + + return new LiveResult(response.getTaskId(), response.getJobId(), agentId, runtimeResponse); + } + + private Long seedAgent(ConfigurableApplicationContext context) { + MuseAgentMapper agentMapper = context.getBean(MuseAgentMapper.class); + MuseAgentVersionMapper versionMapper = context.getBean(MuseAgentVersionMapper.class); + + String suffix = String.valueOf(Instant.now().toEpochMilli()); + MuseAgentDO agent = new MuseAgentDO(); + agent.setAgentKey("p1r-live-agent-" + suffix); + agent.setName("P1R live AI agent"); + agent.setDescription("P1R live acceptance seed agent"); + agent.setAgentType("system"); + agent.setPromptKey("p1r-live-prompt"); + agent.setStatus("active"); + agent.setCategory("p1r-live"); + agent.setTags("[]"); + agent.setSlotBindings("{}"); + agent.setTenantId(TENANT_ID); + agentMapper.insert(agent); + + MuseAgentVersionDO version = new MuseAgentVersionDO(); + version.setAgentId(agent.getId()); + version.setVersion("1"); + version.setConfig(JsonUtils.toJsonString(Map.of( + "modelKey", MODEL, + "promptTemplateVersion", "p1r-live-v1"))); + version.setChangeNote("P1R live acceptance seed version"); + version.setStatus("active"); + version.setCommandId("p1r-live-agent-version-" + suffix); + version.setTenantId(TENANT_ID); + versionMapper.insert(version); + + MuseAgentDO update = new MuseAgentDO(); + update.setId(agent.getId()); + update.setCurrentVersionId(version.getId()); + agentMapper.updateById(update); + return agent.getId(); + } + + private CreateAiTaskReqVO createTaskRequest(Long agentId) { + CreateAiTaskReqVO request = new CreateAiTaskReqVO(); + request.setCommandId("p1r4-ai-e2e-live-" + Instant.now().toEpochMilli()); + request.setWorkId(WORK_ID); + request.setContextScope("work"); + + CreateAiTaskReqVO.AiTaskIntentVO intent = new CreateAiTaskReqVO.AiTaskIntentVO(); + intent.setKind("p1r_live_acceptance"); + intent.setOutputTarget("suggestion"); + intent.setUserInstruction("只输出 OK"); + request.setIntent(intent); + + CreateAiTaskReqVO.AgentOverrideRefVO overrideRef = new CreateAiTaskReqVO.AgentOverrideRefVO(); + overrideRef.setAgentId(agentId); + overrideRef.setAgentVersion(1); + overrideRef.setReason("p1r live acceptance"); + request.setAgentOverrideRef(overrideRef); + return request; + } + + private void assertDbFacts(LiveResult result, DbFacts facts) { + assertEquals(1, facts.generationCount(), "必须存在 1 行 muse_ai_generation"); + assertEquals(1, facts.jobCount(), "必须存在 1 行 muse_ai_job"); + assertEquals(1, facts.runtimeCallCount(), "必须存在 1 行 muse_ai_runtime_call"); + assertEquals(1, facts.auditCount(), "必须存在 1 行 muse_business_audit_event"); + assertEquals(1, facts.taskEventCount(), "必须存在 1 行 muse_ai_task_event 终态事件"); + assertEquals(1, facts.commandCount(), "必须存在 1 行 muse_ai_command"); + assertEquals(1, facts.agentCount(), "必须存在 seed agent"); + assertEquals(1, facts.agentVersionCount(), "必须存在 seed agent version"); + + assertNotNull(result.runtimeResponse().result(), "主 live acceptance 必须有 runtime result"); + assertTrue(result.runtimeResponse().failure() == null, + () -> "主 live acceptance 不接受 failure: " + failureSummary(result.runtimeResponse().failure())); + assertEquals("completed", facts.generationStatus(), "generation 必须进入 completed"); + assertEquals("completed", facts.jobStatus(), "job 必须进入 completed"); + assertEquals("succeeded", facts.runtimeCallStatus(), "runtime_call 必须进入 succeeded"); + assertTrue(facts.runtimeCallStarted(), "runtime_call 必须有 started_at,证明 requested 已落库"); + assertTrue(facts.runtimeCallFinished(), "runtime_call 必须有 finished_at,证明 finished 已落库"); + assertEquals("done", facts.taskEventType(), "task_event 必须是 done 终态事件"); + assertTrue(facts.suggestionCount() >= 1, "runtime 成功时必须生成 suggestion 候选"); + assertTrue(StringUtils.hasText(facts.runtimeCallResponseSummary()) + && !"{}".equals(facts.runtimeCallResponseSummary().trim()), + "runtime_call 成功时必须有非空 response_summary"); + assertTrue(StringUtils.hasText(facts.runtimeCallProviderRequestId()), + "runtime_call 成功时必须写入 provider_request_id"); + assertTrue(StringUtils.hasText(facts.runtimeCallUsageSummary()) + && !"{}".equals(facts.runtimeCallUsageSummary().trim()), + "runtime_call 成功时必须有非空 usage_summary"); + assertTrue(totalTokens(facts.runtimeCallUsageSummary()) > 0, + "usage_summary 必须包含 totalTokens 或 total_tokens 且大于 0"); + assertTrue(StringUtils.hasText(facts.runtimeCallFinishReason()), + "P1R-4 completed 证据必须来自 runtime_call.response_summary.finishReason,不能使用内存 runtime result 替代"); + } + + private DbFacts readDbFacts(LiveSettings settings, LiveResult result) { + try (Connection connection = DriverManager.getConnection(settings.jdbcUrl(), settings.jdbcUser(), + settings.jdbcPassword())) { + return new DbFacts( + count(connection, "muse_ai_generation", "id = ?", result.taskId()), + count(connection, "muse_ai_job", "id = ?", result.jobPkId()), + count(connection, "muse_ai_runtime_call", "task_id = ? AND job_id = ?", + String.valueOf(result.taskId()), String.valueOf(result.jobPkId())), + count(connection, "muse_business_audit_event", "target_id = ?", result.taskId()), + count(connection, "muse_ai_task_event", "task_id = ?", String.valueOf(result.taskId())), + count(connection, "muse_ai_command", + "command_id = (SELECT command_id FROM muse_ai_generation WHERE id = ?)", result.taskId()), + count(connection, "muse_agent", "id = ?", result.agentId()), + count(connection, "muse_agent_version", "agent_id = ?", result.agentId()), + count(connection, "muse_ai_suggestion", "generation_id = ?", result.taskId()), + stringValue(connection, "SELECT status FROM muse_ai_generation WHERE id = ?", result.taskId()), + stringValue(connection, "SELECT status FROM muse_ai_job WHERE id = ?", result.jobPkId()), + stringValue(connection, """ + SELECT status FROM muse_ai_runtime_call + WHERE task_id = ? AND job_id = ? + ORDER BY id DESC LIMIT 1 + """, String.valueOf(result.taskId()), String.valueOf(result.jobPkId())), + stringValue(connection, """ + SELECT response_summary::text FROM muse_ai_runtime_call + WHERE task_id = ? AND job_id = ? + ORDER BY id DESC LIMIT 1 + """, String.valueOf(result.taskId()), String.valueOf(result.jobPkId())), + stringValue(connection, """ + SELECT provider_request_id FROM muse_ai_runtime_call + WHERE task_id = ? AND job_id = ? + ORDER BY id DESC LIMIT 1 + """, String.valueOf(result.taskId()), String.valueOf(result.jobPkId())), + stringValue(connection, """ + SELECT usage_summary::text FROM muse_ai_runtime_call + WHERE task_id = ? AND job_id = ? + ORDER BY id DESC LIMIT 1 + """, String.valueOf(result.taskId()), String.valueOf(result.jobPkId())), + stringValue(connection, """ + SELECT response_summary ->> 'finishReason' FROM muse_ai_runtime_call + WHERE task_id = ? AND job_id = ? + ORDER BY id DESC LIMIT 1 + """, String.valueOf(result.taskId()), String.valueOf(result.jobPkId())), + stringValue(connection, """ + SELECT error_code FROM muse_ai_runtime_call + WHERE task_id = ? AND job_id = ? + ORDER BY id DESC LIMIT 1 + """, String.valueOf(result.taskId()), String.valueOf(result.jobPkId())), + stringValue(connection, """ + SELECT failure_type FROM muse_ai_runtime_call + WHERE task_id = ? AND job_id = ? + ORDER BY id DESC LIMIT 1 + """, String.valueOf(result.taskId()), String.valueOf(result.jobPkId())), + booleanValue(connection, """ + SELECT started_at IS NOT NULL FROM muse_ai_runtime_call + WHERE task_id = ? AND job_id = ? + ORDER BY id DESC LIMIT 1 + """, String.valueOf(result.taskId()), String.valueOf(result.jobPkId())), + booleanValue(connection, """ + SELECT finished_at IS NOT NULL FROM muse_ai_runtime_call + WHERE task_id = ? AND job_id = ? + ORDER BY id DESC LIMIT 1 + """, String.valueOf(result.taskId()), String.valueOf(result.jobPkId())), + stringValue(connection, """ + SELECT event_type FROM muse_ai_task_event + WHERE task_id = ? + ORDER BY sequence_no DESC LIMIT 1 + """, String.valueOf(result.taskId())), + stringValue(connection, """ + SELECT payload_summary::text FROM muse_ai_task_event + WHERE task_id = ? + ORDER BY sequence_no DESC LIMIT 1 + """, String.valueOf(result.taskId()))); + } catch (SQLException ex) { + throw new IllegalStateException("读取 P1R AI live acceptance DB 事实失败", ex); + } + } + + private long count(Connection connection, String table, String where, Object... args) throws SQLException { + String sql = "SELECT COUNT(*) FROM " + table + " WHERE tenant_id = ? AND deleted = FALSE AND " + where; + try (PreparedStatement statement = connection.prepareStatement(sql)) { + statement.setLong(1, TENANT_ID); + bind(statement, 2, args); + try (ResultSet resultSet = statement.executeQuery()) { + assertTrue(resultSet.next(), "必须能读取计数: " + table); + return resultSet.getLong(1); + } + } + } + + private String stringValue(Connection connection, String sql, Object... args) throws SQLException { + try (PreparedStatement statement = connection.prepareStatement(sql)) { + bind(statement, 1, args); + try (ResultSet resultSet = statement.executeQuery()) { + return resultSet.next() ? resultSet.getString(1) : null; + } + } + } + + private boolean booleanValue(Connection connection, String sql, Object... args) throws SQLException { + try (PreparedStatement statement = connection.prepareStatement(sql)) { + bind(statement, 1, args); + try (ResultSet resultSet = statement.executeQuery()) { + return resultSet.next() && resultSet.getBoolean(1); + } + } + } + + private void bind(PreparedStatement statement, int startIndex, Object... args) throws SQLException { + for (int i = 0; i < args.length; i++) { + Object value = args[i]; + if (value instanceof Long longValue) { + statement.setLong(startIndex + i, longValue); + } else { + statement.setString(startIndex + i, String.valueOf(value)); + } + } + } + + private ConfigurableApplicationContext liveContext(LiveSettings settings) { + assertEquals(16, TEST_ENCRYPTOR_PASSWORD.length(), + "test-only mybatis-plus encryptor password 必须是 16 字节 AES 密钥"); + Map properties = new LinkedHashMap<>(); + properties.put("muse.info.base-package", "cn.iocoder.muse.module.ai"); + properties.put("spring.datasource.url", settings.jdbcUrl()); + properties.put("spring.datasource.username", settings.jdbcUser()); + properties.put("spring.datasource.password", settings.jdbcPassword()); + properties.put("spring.datasource.driver-class-name", "org.postgresql.Driver"); + properties.put("spring.main.banner-mode", "off"); + properties.put("spring.main.lazy-initialization", "true"); + properties.put("mybatis-plus.global-config.db-config.id-type", "AUTO"); + properties.put("mybatis-plus.encryptor.password", TEST_ENCRYPTOR_PASSWORD); + properties.put("muse.ai.new-api.enabled", "true"); + properties.put("muse.ai.new-api.base-url", settings.newApiBaseUrl()); + properties.put("muse.ai.new-api.token", settings.newApiToken()); + properties.put("muse.ai.new-api.default-model-key", MODEL); + properties.put("muse.ai.new-api.connect-timeout-seconds", settings.connectTimeoutSeconds()); + properties.put("muse.ai.new-api.first-byte-timeout-seconds", settings.firstByteTimeoutSeconds()); + properties.put("muse.ai.new-api.non-stream-read-timeout-seconds", settings.nonStreamReadTimeoutSeconds()); + properties.put("muse.ai.new-api.stream-idle-timeout-seconds", settings.streamIdleTimeoutSeconds()); + properties.put("muse.ai.new-api.total-timeout-seconds", settings.totalTimeoutSeconds()); + properties.put("muse.ai.new-api.max-attempts", 1); + properties.put("muse.ai.new-api.retry-backoff-seconds", settings.retryBackoffSeconds()); + properties.put("muse.ai.new-api.propagate-muse-context", "true"); + + return new SpringApplicationBuilder(LiveAcceptanceConfiguration.class) + .web(WebApplicationType.NONE) + .initializers(applicationContext -> applicationContext.getEnvironment().getPropertySources() + // 这里只注入 live-only 连接、凭据和超时参数;逻辑删除值必须来自默认 application.yaml。 + .addFirst(new MapPropertySource("p1r-live-acceptance", properties))) + .properties(properties) + .run(); + } + + private DatabaseEngineFacts migrateIsolatedTestDatabase(LiveSettings settings) { + silenceFlywayInfoLogs(); + assertSafeJdbcUrl(settings.jdbcUrl()); + DatabaseEngineFacts dbEngine = assertPostgresqlTestDatabase(settings); + System.out.println(JsonUtils.toJsonString(Map.of( + "acceptance", "p1r4-ai-runtime-e2e-live", + "preCleanDbEngine", dbEngineEvidence(dbEngine)))); + Flyway flyway = Flyway.configure() + .dataSource(settings.jdbcUrl(), settings.jdbcUser(), settings.jdbcPassword()) + .locations(resolveMuseSqlLocation(settings.flywayLocations())) + .schemas("public") + .defaultSchema("public") + .target(TARGET_VERSION) + .cleanDisabled(false) + .load(); + flyway.clean(); + flyway.migrate(); + MigrationInfo current = flyway.info().current(); + assertEquals(TARGET_VERSION, Objects.requireNonNull(current, "必须存在当前 Flyway 版本") + .getVersion().getVersion(), "P1R-4 AI runtime live acceptance 至少需要 V13 schema"); + return dbEngine; + } + + private DatabaseEngineFacts assertPostgresqlTestDatabase(LiveSettings settings) { + try (Connection connection = DriverManager.getConnection(settings.jdbcUrl(), settings.jdbcUser(), + settings.jdbcPassword()); + PreparedStatement statement = connection.prepareStatement("SELECT current_database(), version()"); + ResultSet resultSet = statement.executeQuery()) { + assertTrue(resultSet.next(), "Flyway clean 前必须能读取 PostgreSQL engine 信息"); + String databaseName = resultSet.getString(1); + String version = resultSet.getString(2); + assertTrue(StringUtils.hasText(databaseName) && databaseName.endsWith("_test"), + "Flyway clean 前 current_database() 必须是 _test 隔离库: " + databaseName); + assertTrue(StringUtils.hasText(version) && version.contains("PostgreSQL"), + "Flyway clean 前 version() 必须来自 PostgreSQL: " + versionSummary(version)); + return new DatabaseEngineFacts(databaseName, version); + } catch (SQLException ex) { + throw new IllegalStateException("Flyway clean 前检查 PostgreSQL _test 数据库失败: " + + maskedUrl(settings.jdbcUrl()), ex); + } + } + + private Map redactedEvidence(LiveSettings settings, LiveResult result, DbFacts facts, + DatabaseEngineFacts dbEngine) { + Map evidence = new LinkedHashMap<>(); + evidence.put("acceptance", "p1r4-ai-runtime-e2e-live"); + evidence.put("endpoint", endpointSummary(settings.newApiBaseUrl(), CHAT_COMPLETIONS_PATH)); + evidence.put("token", secretFingerprint(settings.newApiToken())); + evidence.put("jdbcUrl", maskedUrl(settings.jdbcUrl())); + evidence.put("flywayTargetVersion", TARGET_VERSION); + evidence.put("dbEngine", dbEngineEvidence(dbEngine)); + evidence.put("request", requestEvidence(result, facts)); + evidence.put("response", responseEvidence(result, facts)); + evidence.put("db", dbEvidence(facts)); + return evidence; + } + + private Map requestEvidence(LiveResult result, DbFacts facts) { + Map request = new LinkedHashMap<>(); + request.put("tenantId", TENANT_ID); + request.put("ownerUserId", OWNER_USER_ID); + request.put("workId", WORK_ID); + request.put("taskId", result.taskId()); + request.put("jobId", result.jobPkId()); + request.put("agentId", result.agentId()); + request.put("model", MODEL); + request.put("runtimeCallStarted", facts.runtimeCallStarted()); + return request; + } + + private Map responseEvidence(LiveResult result, DbFacts facts) { + Map response = new LinkedHashMap<>(); + response.put("generationStatus", facts.generationStatus()); + response.put("jobStatus", facts.jobStatus()); + response.put("runtimeCallStatus", facts.runtimeCallStatus()); + response.put("taskEventType", facts.taskEventType()); + response.put("runtimeCallFinished", facts.runtimeCallFinished()); + response.put("responseSummaryLength", length(facts.runtimeCallResponseSummary())); + response.put("responseSummarySha256Prefix", sha256Prefix(defaultValue(facts.runtimeCallResponseSummary(), ""), 12)); + response.put("finishReason", facts.runtimeCallFinishReason()); + response.put("usage", result.runtimeResponse().result().tokenUsage()); + response.put("usageTotalTokens", totalTokens(facts.runtimeCallUsageSummary())); + response.put("providerRequestIdPresent", StringUtils.hasText(facts.runtimeCallProviderRequestId())); + response.put("providerRequestIdFingerprint", secretFingerprint(facts.runtimeCallProviderRequestId())); + response.put("completedAt", result.runtimeResponse().result().completedAt()); + return response; + } + + private Map dbEvidence(DbFacts facts) { + Map db = new LinkedHashMap<>(); + db.put("generationCount", facts.generationCount()); + db.put("jobCount", facts.jobCount()); + db.put("runtimeCallCount", facts.runtimeCallCount()); + db.put("auditCount", facts.auditCount()); + db.put("taskEventCount", facts.taskEventCount()); + db.put("commandCount", facts.commandCount()); + db.put("agentCount", facts.agentCount()); + db.put("agentVersionCount", facts.agentVersionCount()); + db.put("suggestionCount", facts.suggestionCount()); + db.put("runtimeCallUsageSummaryLength", length(facts.runtimeCallUsageSummary())); + db.put("runtimeCallUsageSummarySha256Prefix", sha256Prefix(defaultValue(facts.runtimeCallUsageSummary(), ""), 12)); + db.put("runtimeCallFinishReason", facts.runtimeCallFinishReason()); + db.put("taskEventPayloadLength", length(facts.taskEventPayloadSummary())); + db.put("taskEventPayloadSha256Prefix", sha256Prefix(defaultValue(facts.taskEventPayloadSummary(), ""), 12)); + return db; + } + + private Map dbEngineEvidence(DatabaseEngineFacts facts) { + Map engine = new LinkedHashMap<>(); + engine.put("databaseNameSha256Prefix", sha256Prefix(facts.databaseName(), 12)); + engine.put("databaseNameSuffix", facts.databaseName().endsWith("_test") ? "_test" : ""); + engine.put("versionSummary", versionSummary(facts.version())); + return engine; + } + + private Map endpointSummary(String baseUrl, String path) { + URI uri = URI.create(trimRight(baseUrl, "/") + path); + Map endpoint = new LinkedHashMap<>(); + endpoint.put("scheme", uri.getScheme()); + endpoint.put("host", uri.getHost()); + endpoint.put("port", uri.getPort()); + endpoint.put("path", uri.getPath()); + return endpoint; + } + + private Map secretFingerprint(String secret) { + Map fingerprint = new LinkedHashMap<>(); + fingerprint.put("length", secret == null ? 0 : secret.length()); + fingerprint.put("sha256Prefix", sha256Prefix(secret == null ? "" : secret, 12)); + return fingerprint; + } + + private long totalTokens(String usageSummary) { + if (!StringUtils.hasText(usageSummary)) { + return 0; + } + try { + JsonNode usage = JsonUtils.parseTree(usageSummary); + if (usage == null || usage.isMissingNode() || usage.isNull()) { + return 0; + } + if (usage.has("totalTokens")) { + return usage.path("totalTokens").asLong(0); + } + if (usage.has("total_tokens")) { + return usage.path("total_tokens").asLong(0); + } + return 0; + } catch (RuntimeException ex) { + throw new AssertionError("usage_summary 必须是 JSON: " + sha256Prefix(usageSummary, 12), ex); + } + } + + private static String failureSummary(MuseAiRuntimeClient.RuntimeFailure failure) { + if (failure == null) { + return ""; + } + return "type=" + failure.failureType() + + ", code=" + failure.errorCode() + + ", status=" + failure.providerStatusCode() + + ", retryable=" + failure.retryable(); + } + + private static String requiredPropertyOrEnv(String propertyName, String envName) { + String value = System.getProperty(propertyName); + if (!StringUtils.hasText(value)) { + value = System.getenv(envName); + } + assertTrue(StringUtils.hasText(value), "缺少必需系统属性或环境变量: " + propertyName + " / " + envName); + return value.trim(); + } + + private static String requiredEnv(String name) { + String value = System.getenv(name); + assertTrue(StringUtils.hasText(value), "缺少必需环境变量: " + name); + return value.trim(); + } + + private static String requiredPasswordEnvironment() { + String password = firstNonBlankEnvironment("P1R_FLYWAY_PASSWORD", "MUSE_POSTGRES_PASSWORD"); + assertTrue(password != null, "缺少必需数据库密码环境变量: P1R_FLYWAY_PASSWORD 或 MUSE_POSTGRES_PASSWORD"); + return password; + } + + private static String firstNonBlankEnvironment(String... names) { + // 数据库密码不允许从 system property 读取,避免被 Surefire XML 或命令历史记录。 + for (String name : names) { + String value = System.getenv(name); + if (StringUtils.hasText(value)) { + return value; + } + } + return null; + } + + private static int intEnv(String name, int defaultValue) { + String value = System.getenv(name); + if (!StringUtils.hasText(value)) { + return defaultValue; + } + return Integer.parseInt(value.trim()); + } + + private static String stringEnv(String name, String defaultValue) { + String value = System.getenv(name); + return StringUtils.hasText(value) ? value.trim() : defaultValue; + } + + private static String resolveMuseSqlLocation(String requestedLocations) { + assertEquals("filesystem:sql/muse", requestedLocations, + "P1R live acceptance 要求显式使用 filesystem:sql/muse"); + Path current = Path.of(System.getProperty("user.dir")).toAbsolutePath(); + for (Path cursor = current; cursor != null; cursor = cursor.getParent()) { + Path candidate = cursor.resolve("sql/muse"); + if (Files.isDirectory(candidate)) { + return "filesystem:" + candidate; + } + } + throw new IllegalStateException("无法从当前目录向上找到 sql/muse: " + current); + } + + private static void assertSafeJdbcUrl(String url) { + assertTrue(url.startsWith(POSTGRESQL_JDBC_PREFIX), + "JDBC URL 必须显式使用 jdbc:postgresql://,避免 Flyway clean 误连其它数据库: " + maskedUrl(url)); + URI uri = parseJdbcUriForAssertion(url); + assertEquals("postgresql", uri.getScheme(), + "JDBC URL 必须解析为 PostgreSQL scheme: " + maskedUrl(url)); + assertTrue(StringUtils.hasText(uri.getHost()), + "JDBC URL 必须显式包含 host: " + maskedUrl(url)); + assertFalse(StringUtils.hasText(uri.getRawUserInfo()), + "JDBC URL 不能携带 authority/userinfo,用户名走属性/环境变量,密码只走环境变量: " + maskedUrl(url)); + assertNoCredentialQuery(url); + String databaseName = jdbcDatabaseName(uri); + assertTrue(databaseName.endsWith("_test"), + "JDBC URL 必须指向 _test 后缀隔离库,避免清理非测试库: " + maskedUrl(url)); + assertFalse("muse_local".equals(databaseName), + "P1R live acceptance 禁止默认连接 muse_local"); + } + + private static URI parseJdbcUriForAssertion(String url) { + try { + return URI.create(url.substring(JDBC_URI_PREFIX.length())); + } catch (IllegalArgumentException ex) { + // JDK URI 解析异常会回显原始 authority/query;断言失败只允许输出脱敏后的 JDBC 摘要。 + throw new AssertionError("JDBC URL 必须可解析为 PostgreSQL URI: " + maskedUrl(url)); + } + } + + private static String jdbcDatabaseName(URI uri) { + String path = uri.getRawPath(); + assertTrue(StringUtils.hasText(path) && path.length() > 1, + "JDBC URL 必须包含数据库名"); + String databaseName = path.substring(path.lastIndexOf('/') + 1); + assertTrue(StringUtils.hasText(databaseName), + "JDBC URL 必须包含数据库名"); + return databaseName; + } + + private static void assertNoCredentialQuery(String url) { + int queryStart = url.indexOf('?'); + if (queryStart < 0) { + return; + } + String query = url.substring(queryStart + 1); + for (String parameter : query.split("&")) { + String key = parameter; + int equalsStart = key.indexOf('='); + if (equalsStart >= 0) { + key = key.substring(0, equalsStart); + } + assertFalse(isCredentialQueryKey(key), + "JDBC URL 不能携带凭据 query 参数;用户名走属性/环境变量,密码只走环境变量"); + } + } + + private static boolean isCredentialQueryKey(String rawKey) { + String key = rawKey.trim().toLowerCase(Locale.ROOT).replace('-', '_'); + return CREDENTIAL_QUERY_KEYS.contains(key) + || key.endsWith("_token") + || key.endsWith("_secret") + || key.endsWith("_password"); + } + + private static String maskedUrl(String url) { + if (!StringUtils.hasText(url)) { + return ""; + } + // 失败消息只能暴露 JDBC 类型、占位 authority、测试库名和 query 是否存在;解析失败也不能回显原始串。 + String trimmed = url.trim(); + if (!trimmed.startsWith(JDBC_URI_PREFIX)) { + return ""; + } + + String jdbcBody = trimmed.substring(JDBC_URI_PREFIX.length()); + String scheme = jdbcScheme(jdbcBody); + if (!StringUtils.hasText(scheme)) { + return JDBC_URI_PREFIX + ""; + } + if (!jdbcBody.regionMatches(true, scheme.length(), "://", 0, 3)) { + return JDBC_URI_PREFIX + scheme + ":"; + } + if ("postgresql".equals(scheme)) { + return maskedPostgresqlUrl(jdbcBody); + } + return maskedHierarchicalJdbcUrl(scheme, jdbcBody); + } + + private static String jdbcScheme(String jdbcBody) { + int schemeEnd = jdbcBody.indexOf(':'); + if (schemeEnd <= 0) { + return ""; + } + String scheme = jdbcBody.substring(0, schemeEnd).trim().toLowerCase(Locale.ROOT); + return scheme.matches("[a-z][a-z0-9+.-]*") ? scheme : ""; + } + + private static String maskedPostgresqlUrl(String jdbcBody) { + try { + URI uri = URI.create(jdbcBody); + String database = safeJdbcDatabaseName(uri.getRawPath(), ""); + String port = uri.getPort() < 0 ? "" : ":" + uri.getPort(); + String suffix = StringUtils.hasText(uri.getRawQuery()) ? "?" : ""; + return POSTGRESQL_JDBC_PREFIX + "" + port + "/" + database + suffix; + } catch (RuntimeException ignored) { + return POSTGRESQL_JDBC_PREFIX + "/" + fallbackJdbcDatabaseName(jdbcBody, "") + + querySuffix(jdbcBody); + } + } + + private static String maskedHierarchicalJdbcUrl(String scheme, String jdbcBody) { + try { + URI uri = URI.create(jdbcBody); + String database = safeJdbcDatabaseName(uri.getRawPath(), ""); + String databaseSuffix = StringUtils.hasText(database) ? "/" + database : ""; + String querySuffix = StringUtils.hasText(uri.getRawQuery()) ? "?" : ""; + return JDBC_URI_PREFIX + scheme + "://" + databaseSuffix + querySuffix; + } catch (RuntimeException ignored) { + String database = fallbackJdbcDatabaseName(jdbcBody, ""); + String databaseSuffix = StringUtils.hasText(database) ? "/" + database : ""; + return JDBC_URI_PREFIX + scheme + "://" + databaseSuffix + querySuffix(jdbcBody); + } + } + + private static String safeJdbcDatabaseName(String rawPath, String fallback) { + if (!StringUtils.hasText(rawPath) || rawPath.length() <= 1) { + return fallback; + } + String database = rawPath.substring(rawPath.lastIndexOf('/') + 1); + if (!StringUtils.hasText(database)) { + return fallback; + } + return database.matches("[A-Za-z0-9._%-]+") ? database : fallback; + } + + private static String fallbackJdbcDatabaseName(String jdbcBody, String fallback) { + int authorityStart = jdbcBody.indexOf("://"); + if (authorityStart < 0) { + return fallback; + } + int pathStart = jdbcBody.indexOf('/', authorityStart + 3); + if (pathStart < 0) { + return fallback; + } + int pathEnd = firstIndexOf(jdbcBody, pathStart + 1, '?', '#'); + String path = jdbcBody.substring(pathStart, pathEnd); + return safeJdbcDatabaseName(path, fallback); + } + + private static int firstIndexOf(String value, int start, char... candidates) { + int result = value.length(); + for (char candidate : candidates) { + int index = value.indexOf(candidate, start); + if (index >= 0 && index < result) { + result = index; + } + } + return result; + } + + private static String querySuffix(String value) { + return value.indexOf('?') < 0 ? "" : "?"; + } + + private static String versionSummary(String version) { + if (!StringUtils.hasText(version)) { + return ""; + } + String normalized = version.replaceAll("\\s+", " ").trim(); + return normalized.length() > 96 ? normalized.substring(0, 96) : normalized; + } + + private static void silenceFlywayInfoLogs() { + try { + Object flywayLogger = LoggerFactory.getLogger("org.flywaydb"); + Class levelClass = Class.forName("ch.qos.logback.classic.Level"); + Object warnLevel = levelClass.getField("WARN").get(null); + flywayLogger.getClass().getMethod("setLevel", levelClass).invoke(flywayLogger, warnLevel); + } catch (ReflectiveOperationException | LinkageError ignored) { + // 日志实现不是 logback 时不影响验收;测试自身仍只输出 masked URL 和摘要。 + } + } + + private static boolean externalAcceptanceEnabled() { + return "true".equalsIgnoreCase(System.getenv(ACCEPTANCE_ENV)); + } + + private String sha256Prefix(String value, int length) { + try { + MessageDigest digest = MessageDigest.getInstance("SHA-256"); + String hex = HexFormat.of().formatHex(digest.digest(value.getBytes(StandardCharsets.UTF_8))); + return hex.substring(0, Math.min(length, hex.length())); + } catch (NoSuchAlgorithmException e) { + throw new IllegalStateException("JDK 缺少 SHA-256 摘要算法", e); + } + } + + private String trimRight(String value, String suffix) { + String result = value == null ? "" : value.trim(); + while (result.endsWith(suffix)) { + result = result.substring(0, result.length() - suffix.length()); + } + return result; + } + + private String defaultValue(String value, String fallback) { + return value == null ? fallback : value; + } + + private int length(String value) { + return value == null ? 0 : value.length(); + } + + private record LiveSettings(String jdbcUrl, + String jdbcUser, + String jdbcPassword, + String flywayLocations, + String newApiBaseUrl, + String newApiToken, + int connectTimeoutSeconds, + int firstByteTimeoutSeconds, + int nonStreamReadTimeoutSeconds, + int streamIdleTimeoutSeconds, + int totalTimeoutSeconds, + String retryBackoffSeconds) { + + private static LiveSettings fromEnvironment() { + String jdbcUrl = requiredPropertyOrEnv(JDBC_URL_PROPERTY, JDBC_URL_ENV); + assertSafeJdbcUrl(jdbcUrl); + return new LiveSettings( + jdbcUrl, + requiredPropertyOrEnv(JDBC_USER_PROPERTY, JDBC_USER_ENV), + requiredPasswordEnvironment(), + System.getProperty(FLYWAY_LOCATION_PROPERTY, "filesystem:sql/muse"), + requiredEnv("MUSE_AI_NEW_API_BASE_URL"), + requiredEnv("MUSE_AI_NEW_API_TOKEN"), + intEnv("MUSE_AI_NEW_API_CONNECT_TIMEOUT_SECONDS", 5), + intEnv("MUSE_AI_NEW_API_FIRST_BYTE_TIMEOUT_SECONDS", 15), + intEnv("MUSE_AI_NEW_API_NON_STREAM_READ_TIMEOUT_SECONDS", 60), + intEnv("MUSE_AI_NEW_API_STREAM_IDLE_TIMEOUT_SECONDS", 30), + intEnv("MUSE_AI_NEW_API_TOTAL_TIMEOUT_SECONDS", 180), + stringEnv("MUSE_AI_NEW_API_RETRY_BACKOFF_SECONDS", "1,2,4")); + } + } + + private record LiveResult(Long taskId, + Long jobPkId, + Long agentId, + MuseAiRuntimeClient.RuntimeResponse runtimeResponse) { + } + + private record DbFacts(long generationCount, + long jobCount, + long runtimeCallCount, + long auditCount, + long taskEventCount, + long commandCount, + long agentCount, + long agentVersionCount, + long suggestionCount, + String generationStatus, + String jobStatus, + String runtimeCallStatus, + String runtimeCallResponseSummary, + String runtimeCallProviderRequestId, + String runtimeCallUsageSummary, + String runtimeCallFinishReason, + String runtimeCallErrorCode, + String runtimeCallFailureType, + boolean runtimeCallStarted, + boolean runtimeCallFinished, + String taskEventType, + String taskEventPayloadSummary) { + } + + private record DatabaseEngineFacts(String databaseName, + String version) { + } + + @TestConfiguration + @Import({ + MuseDataSourceAutoConfiguration.class, + DataSourceAutoConfiguration.class, + DataSourceTransactionManagerAutoConfiguration.class, + MuseMybatisAutoConfiguration.class, + MybatisPlusAutoConfiguration.class, + MybatisPlusJoinAutoConfiguration.class, + MuseAiTaskServiceImpl.class, + MuseAiCommandServiceImpl.class, + MuseAiRuntimeJobExecutor.class, + MuseAiRuntimeCallRecorder.class, + MuseAiRuntimeProjectionService.class, + MuseAiRuntimeServiceImpl.class, + MuseAiRuntimePolicyProvider.class, + MuseAiRuntimePayloadStore.class, + MuseAiAuditServiceImpl.class, + SpringUtil.class + }) + static class LiveAcceptanceConfiguration { + + @Bean + MuseAiProperties museAiProperties() { + MuseAiProperties properties = new MuseAiProperties(); + MuseAiProperties.NewApi newApi = new MuseAiProperties.NewApi(); + newApi.setEnabled(true); + newApi.setBaseUrl(requiredEnv("MUSE_AI_NEW_API_BASE_URL")); + newApi.setToken(requiredEnv("MUSE_AI_NEW_API_TOKEN")); + newApi.setDefaultModelKey(MODEL); + newApi.setConnectTimeoutSeconds(intEnv("MUSE_AI_NEW_API_CONNECT_TIMEOUT_SECONDS", 5)); + newApi.setFirstByteTimeoutSeconds(intEnv("MUSE_AI_NEW_API_FIRST_BYTE_TIMEOUT_SECONDS", 15)); + newApi.setNonStreamReadTimeoutSeconds(intEnv("MUSE_AI_NEW_API_NON_STREAM_READ_TIMEOUT_SECONDS", 60)); + newApi.setStreamIdleTimeoutSeconds(intEnv("MUSE_AI_NEW_API_STREAM_IDLE_TIMEOUT_SECONDS", 30)); + newApi.setTotalTimeoutSeconds(intEnv("MUSE_AI_NEW_API_TOTAL_TIMEOUT_SECONDS", 180)); + // live acceptance 只跑单次 provider 调用;外部失败必须投影为明确失败态,不能停在 queued retry。 + newApi.setMaxAttempts(1); + newApi.setRetryBackoffSeconds(parseBackoffSeconds(stringEnv("MUSE_AI_NEW_API_RETRY_BACKOFF_SECONDS", "1,2,4"))); + newApi.setPropagateMuseContext(true); + properties.setNewApi(newApi); + return properties; + } + + @Bean + MuseAiRuntimeClient museAiRuntimeClient(MuseAiProperties properties) { + return new RealNewApiMuseAiRuntimeClient(properties); + } + + @Bean + @Primary + MuseContentWorkOwnerFacade p1rLiveContentWorkOwnerFacade() { + return new MuseContentWorkOwnerFacade() { + + @Override + public void requireWorkOwner(Long workId, Long ownerUserId) { + assertSeedWorkOwner(workId, ownerUserId); + } + + @Override + public void requireSourceOwner(String sourceType, Long sourceId, Long ownerUserId) { + assertEquals("work", sourceType, "live acceptance 只使用 work source"); + assertSeedWorkOwner(sourceId, ownerUserId); + } + + @Override + public SourceSnapshot buildSourceSnapshot(SourceSnapshotRequest request) { + assertSeedWorkOwner(request.workId(), request.ownerUserId()); + String sourceSnapshotId = "src-p1r-live-" + Instant.now().toEpochMilli(); + Map summary = new LinkedHashMap<>(); + summary.put("sourceSnapshotId", sourceSnapshotId); + summary.put("workId", request.workId()); + summary.put("sourceRevision", 1); + summary.put("contextScope", request.contextScope()); + summary.put("sourceOwner", "content-test-facade"); + summary.put("sourceType", "work"); + summary.put("additionalContextRefCount", request.additionalRefs() == null ? 0 : request.additionalRefs().size()); + // 测试 facade 只返回脱敏 source fact,不返回正文,保持和生产 Content 边界一致。 + return new MuseContentWorkOwnerFacade.SourceSnapshot(sourceSnapshotId, request.workId(), + request.chapterId(), request.blockId(), 1, request.contextScope(), List.of(), summary); + } + + private void assertSeedWorkOwner(Long workId, Long ownerUserId) { + assertEquals(WORK_ID, workId, "live acceptance 只允许 seed workId"); + assertEquals(OWNER_USER_ID, ownerUserId, "source snapshot 必须绑定 ownerUserId"); + } + }; + } + + @Bean + @Primary + SecurityRuntimePermissionFacade p1rLiveSecurityRuntimePermissionFacade() { + return request -> { + String envelopeId = "rpe-p1r-live-" + Instant.now().toEpochMilli(); + MuseAiRuntimePermissionGuard.RuntimePermissionEnvelope envelope = + MuseAiRuntimePermissionGuard.RuntimePermissionEnvelope.builder() + .envelopeId(envelopeId) + .operationId(request.operationId()) + .actorUserId(request.actorUserId()) + .ownerUserId(request.ownerUserId()) + .targetType(request.targetType()) + .targetId(request.targetId()) + .targetKey(request.targetKey()) + .build(); + Map summary = new LinkedHashMap<>(); + summary.put("runtimePermissionEnvelopeId", envelopeId); + summary.put("operationId", request.operationId()); + summary.put("actorUserId", request.actorUserId()); + summary.put("ownerUserId", request.ownerUserId()); + summary.put("targetType", request.targetType()); + summary.put("targetId", request.targetId()); + summary.put("targetKey", request.targetKey()); + summary.put("approved", true); + summary.put("approvalSource", "p1r-live-test-facade"); + // 这里是 test-only Security facade;真实运行仍由 Security owner 投影或外部授权负责。 + return new SecurityRuntimePermissionFacade.RuntimePermissionEnvelopeIssue(envelopeId, envelope, summary); + }; + } + + private static List parseBackoffSeconds(String value) { + if (!StringUtils.hasText(value)) { + return List.of(1, 2, 4); + } + List result = new ArrayList<>(); + for (String token : value.split(",")) { + if (StringUtils.hasText(token)) { + result.add(Integer.parseInt(token.trim())); + } + } + return result.isEmpty() ? List.of(1, 2, 4) : result; + } + } +} diff --git a/muse-cloud/muse-server/src/test/java/cn/iocoder/muse/server/framework/api/P1rKnowledgeRuntimeEndToEndLiveAcceptanceIT.java b/muse-cloud/muse-server/src/test/java/cn/iocoder/muse/server/framework/api/P1rKnowledgeRuntimeEndToEndLiveAcceptanceIT.java new file mode 100644 index 00000000..918d704b --- /dev/null +++ b/muse-cloud/muse-server/src/test/java/cn/iocoder/muse/server/framework/api/P1rKnowledgeRuntimeEndToEndLiveAcceptanceIT.java @@ -0,0 +1,1071 @@ +package cn.iocoder.muse.server.framework.api; + +import cn.hutool.extra.spring.SpringUtil; +import cn.iocoder.muse.framework.common.enums.UserTypeEnum; +import cn.iocoder.muse.framework.common.util.json.JsonUtils; +import cn.iocoder.muse.framework.datasource.config.MuseDataSourceAutoConfiguration; +import cn.iocoder.muse.framework.mybatis.config.MuseMybatisAutoConfiguration; +import cn.iocoder.muse.framework.security.core.LoginUser; +import cn.iocoder.muse.framework.security.core.util.SecurityFrameworkUtils; +import cn.iocoder.muse.framework.tenant.core.context.TenantContextHolder; +import cn.iocoder.muse.framework.tenant.core.util.TenantUtils; +import cn.iocoder.muse.module.infra.api.file.FileApi; +import cn.iocoder.muse.module.knowledge.application.muse.MuseKnowledgeAuditService; +import cn.iocoder.muse.module.knowledge.application.muse.MuseKnowledgeCommandService; +import cn.iocoder.muse.module.knowledge.application.muse.MuseKnowledgeDocumentService; +import cn.iocoder.muse.module.knowledge.application.muse.MuseKnowledgeMaterializationService; +import cn.iocoder.muse.module.knowledge.application.muse.MuseKnowledgeProcessingTaskService; +import cn.iocoder.muse.module.knowledge.application.muse.facade.HttpRagFlowKnowledgeRuntimeClient; +import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeFileFacade; +import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient; +import cn.iocoder.muse.module.knowledge.controller.app.muse.vo.AppKnowledgeDocumentVO; +import cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeBaseDO; +import cn.iocoder.muse.module.knowledge.dal.mysql.muse.MuseKnowledgeBaseMapper; +import com.baomidou.mybatisplus.autoconfigure.MybatisPlusAutoConfiguration; +import com.github.yulichang.autoconfigure.MybatisPlusJoinAutoConfiguration; +import com.fasterxml.jackson.databind.JsonNode; +import org.flywaydb.core.Flyway; +import org.flywaydb.core.api.MigrationInfo; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.Test; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.ObjectProvider; +import org.springframework.boot.WebApplicationType; +import org.springframework.boot.autoconfigure.jdbc.DataSourceAutoConfiguration; +import org.springframework.boot.autoconfigure.jdbc.DataSourceTransactionManagerAutoConfiguration; +import org.springframework.boot.builder.SpringApplicationBuilder; +import org.springframework.boot.test.context.TestConfiguration; +import org.springframework.context.ConfigurableApplicationContext; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Import; +import org.springframework.context.annotation.Primary; +import org.springframework.core.env.MapPropertySource; +import org.springframework.mock.web.MockHttpServletRequest; +import org.springframework.security.core.context.SecurityContextHolder; +import org.springframework.util.StringUtils; + +import javax.sql.DataSource; +import java.net.URI; +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.nio.file.Path; +import java.security.MessageDigest; +import java.security.NoSuchAlgorithmException; +import java.sql.Connection; +import java.sql.DriverManager; +import java.sql.PreparedStatement; +import java.sql.ResultSet; +import java.sql.SQLException; +import java.sql.Statement; +import java.time.Duration; +import java.time.Instant; +import java.util.HexFormat; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Locale; +import java.util.Map; +import java.util.Objects; +import java.util.Set; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.junit.jupiter.api.Assumptions.assumeTrue; + +/** + * P1R-5 Knowledge 主业务入口到真实 RAGFlow / PostgreSQL 的 opt-in live acceptance。 + * + *

默认必须跳过,避免普通测试误清 PostgreSQL 或外呼真实 RAGFlow。只有显式设置 + * {@code MUSE_P1R_EXTERNAL_ACCEPTANCE=true} 时,本测试才会清理 PostgreSQL {@code _test} 库、 + * 迁移到 V14、seed 最小用户 KB,并通过 {@link MuseKnowledgeDocumentService#uploadAppDocument} + * 触发真实 {@link HttpRagFlowKnowledgeRuntimeClient}。

+ */ +class P1rKnowledgeRuntimeEndToEndLiveAcceptanceIT { + + private static final String ACCEPTANCE_ENV = "MUSE_P1R_EXTERNAL_ACCEPTANCE"; + private static final String JDBC_URL_PROPERTY = "p1r.flyway.url"; + private static final String JDBC_USER_PROPERTY = "p1r.flyway.user"; + private static final String JDBC_URL_ENV = "P1R_FLYWAY_URL"; + private static final String JDBC_USER_ENV = "P1R_FLYWAY_USER"; + private static final String FLYWAY_LOCATION_PROPERTY = "p1r.flyway.locations"; + private static final String POSTGRESQL_JDBC_PREFIX = "jdbc:postgresql://"; + private static final String JDBC_URI_PREFIX = "jdbc:"; + private static final String TARGET_VERSION = "14"; + private static final Long TENANT_ID = 100L; + private static final Long OWNER_USER_ID = 2001L; + private static final String API_VERSION = "1"; + private static final String RAGFLOW_API_PATH = "/api/v1"; + private static final Set CREDENTIAL_QUERY_KEYS = Set.of( + "user", "username", "password", "pass", "pwd", "sslpassword", "ssl_password", + "token", "secret", "api_key", "apikey", "bearer", "access_token", "refresh_token"); + + @AfterEach + void clearRuntimeContexts() { + // live acceptance 在同一 JVM 内手工设置租户和登录态,测试结束必须清理,避免污染后续测试。 + SecurityContextHolder.clearContext(); + TenantContextHolder.clear(); + } + + @Test + void shouldMaskNonPostgresqlJdbcAuthorityInFailureOutput() { + String masked = maskedUrl("jdbc:mysql://user:password@host:3306/db_test?password=x"); + + assertFalse(masked.contains("user"), () -> "非 PostgreSQL JDBC URL 失败摘要不能泄露 userinfo: " + masked); + assertFalse(masked.contains("password"), () -> "非 PostgreSQL JDBC URL 失败摘要不能泄露密码: " + masked); + assertFalse(masked.contains("host"), () -> "非 PostgreSQL JDBC URL 失败摘要不能泄露真实 host: " + masked); + } + + @Test + void shouldMaskMalformedPostgresqlJdbcParseFailureOutput() { + String url = "jdbc:postgresql://user:password@bad host/muse_local_test?password=x"; + + AssertionError error = assertThrows(AssertionError.class, () -> assertSafeJdbcUrl(url)); + String message = error.getMessage(); + + assertTrue(message.contains("") || message.contains(""), + () -> "malformed PostgreSQL JDBC URL 解析失败消息必须包含脱敏摘要: " + message); + assertFalse(message.contains("user"), () -> "解析失败消息不能泄露 userinfo: " + message); + assertFalse(message.contains("password"), () -> "解析失败消息不能泄露密码或凭据参数: " + message); + assertFalse(message.contains("bad host"), () -> "解析失败消息不能泄露真实 host: " + message); + } + + @Test + void shouldUploadKnowledgeDocumentThroughMuseBusinessServiceAndPersistDbFactsWithRealPostgresqlAndRagFlow() + throws Exception { + assumeTrue(externalAcceptanceEnabled(), + "P1R 外部端到端验收未启用,设置 MUSE_P1R_EXTERNAL_ACCEPTANCE=true 后才运行"); + + LiveSettings settings = LiveSettings.fromEnvironment(); + DatabaseEngineFacts dbEngine = migrateIsolatedTestDatabase(settings); + + try (ConfigurableApplicationContext context = liveContext(settings)) { + setLiveLoginUser(); + LiveResult result = TenantUtils.execute(TENANT_ID, () -> executeMuseKnowledgeBusinessFlow(context)); + DbFacts dbFacts = readDbFacts(settings, result); + assertDbFacts(result, dbFacts); + ExternalVerification externalVerification = verifyExternalRagFlowSearchable(context, result, dbFacts); + + System.out.println(JsonUtils.toJsonString(redactedEvidence(settings, result, dbFacts, + externalVerification, dbEngine))); + } + } + + private void setLiveLoginUser() { + // 当前测试直接调用 application service,没有 HTTP 认证过滤器;这里补齐审计字段依赖的登录态。 + SecurityFrameworkUtils.setLoginUser(liveLoginUser(), new MockHttpServletRequest()); + } + + private LoginUser liveLoginUser() { + LoginUser loginUser = new LoginUser(); + loginUser.setId(OWNER_USER_ID); + loginUser.setUserType(UserTypeEnum.MEMBER.getValue()); + loginUser.setTenantId(TENANT_ID); + loginUser.setVisitTenantId(TENANT_ID); + return loginUser; + } + + private LiveResult executeMuseKnowledgeBusinessFlow(ConfigurableApplicationContext context) { + Long kbId = seedUserKnowledgeBase(context); + String suffix = String.valueOf(Instant.now().toEpochMilli()); + String commandId = "p1r5-knowledge-e2e-live-" + suffix; + + AppKnowledgeDocumentVO.UploadReqVO request = new AppKnowledgeDocumentVO.UploadReqVO(); + request.setCommandId(commandId); + request.setEntryContent(liveDocumentContent(suffix)); + request.setSourceDescription("P1R-5 Knowledge 主业务 live acceptance"); + + AppKnowledgeDocumentVO.UploadRespVO response = context.getBean(MuseKnowledgeDocumentService.class) + .uploadAppDocument(OWNER_USER_ID, API_VERSION, kbId, request); + assertNotNull(response.getDocumentId(), "uploadAppDocument 必须返回 documentId"); + assertTrue(StringUtils.hasText(response.getScanTaskId()), "uploadAppDocument 必须返回 processing task id"); + assertEquals("pending_scan", response.getProcessingStatus(), + "上传响应不是最终 RAGFlow 证据,最终判断必须读取 PostgreSQL facts"); + return new LiveResult(kbId, commandId, response.getDocumentId(), response.getScanTaskId(), suffix); + } + + private Long seedUserKnowledgeBase(ConfigurableApplicationContext context) { + restartKnowledgeBaseIdentity(context.getBean(DataSource.class)); + MuseKnowledgeBaseDO kb = new MuseKnowledgeBaseDO(); + kb.setName("P1R-5 live user KB"); + kb.setDescription("P1R-5 Knowledge / RAGFlow 主业务端到端验收 seed KB"); + kb.setKbType("user"); + kb.setOwnerUserId(OWNER_USER_ID); + kb.setStatus("active"); + kb.setActiveVersion(1); + kb.setVisibilityPolicy(JsonUtils.toJsonString(Map.of("scope", "p1r-live", "owner", "user"))); + kb.setCommandId("p1r5-live-kb-" + Instant.now().toEpochMilli()); + kb.setRevision(1); + kb.setTenantId(TENANT_ID); + context.getBean(MuseKnowledgeBaseMapper.class).insert(kb); + assertNotNull(kb.getId(), "seed KB 必须返回 id"); + return kb.getId(); + } + + private void restartKnowledgeBaseIdentity(DataSource dataSource) { + long restartWith = Instant.now().toEpochMilli(); + try (Connection connection = dataSource.getConnection(); + Statement statement = connection.createStatement()) { + // RAGFlow datasetName 由业务服务按 kbId + activeVersion 生成;重置到时间戳区间可避免 live 重跑撞上旧外部 dataset。 + statement.execute("ALTER TABLE muse_knowledge_base ALTER COLUMN id RESTART WITH " + restartWith); + } catch (SQLException ex) { + throw new IllegalStateException("设置 P1R live seed KB identity 起点失败", ex); + } + } + + private String liveDocumentContent(String suffix) { + return """ + P1R-5 Knowledge live acceptance marker %s. + Muse Knowledge business service uploads this entryContent through uploadAppDocument. + The live assertion proves Muse DB facts for KB, document, version, processing task, bindings, RAGFlow calls and command. + """.formatted(suffix); + } + + private void assertDbFacts(LiveResult result, DbFacts facts) { + assertEquals(1, facts.knowledgeBaseCount(), "必须存在 1 行 muse_knowledge_base seed KB"); + assertEquals(1, facts.documentCount(), "必须存在 1 行 muse_knowledge_document"); + assertEquals(1, facts.documentVersionCount(), "必须存在 1 行 muse_knowledge_document_version"); + assertEquals(1, facts.processingTaskCount(), "必须存在 1 行 muse_knowledge_processing_task"); + assertEquals(1, facts.commandCount(), "必须存在 1 行 muse_knowledge_command"); + assertEquals(2, facts.bindingCount(), "必须存在 dataset binding + document version binding"); + assertEquals(1, facts.datasetBindingCount(), "必须存在 1 行 dataset binding"); + assertEquals(1, facts.documentBindingCount(), "必须存在 1 行 document version binding"); + assertTrue(facts.ragflowCallCount() >= 3, "必须至少记录 createDataset/uploadDocuments/startParseDocuments 3 次 RAGFlow 调用"); + + assertEquals("user", facts.kbType(), "seed KB 必须是用户 KB"); + assertEquals("active", facts.kbStatus(), "seed KB 必须是 active"); + assertEquals(1, facts.kbActiveVersion(), "seed KB activeVersion 必须是 1"); + assertEquals("completed", facts.commandStatus(), "upload command 必须完成,证明可回放快照已落库"); + assertEquals("passed", facts.documentVersionScanStatus(), "test-only scanner 必须只把扫描结果标记为 passed"); + assertEquals("processing", facts.documentVersionProcessingStatus(), "RAGFlow parse accepted 后资料版本应处于 processing"); + assertEquals("parsing", facts.processingTaskStatus(), "processing_task 必须进入 parsing 语义"); + assertEquals("materialized", facts.materializationStatus(), "entryContent 必须已材料化为可上传 bytes"); + assertEquals("passed", facts.processingTaskScanStatus(), "processing_task 必须记录扫描已通过"); + assertEquals("processing", facts.processingTaskParseStatus(), "processing_task 必须记录 parse processing"); + assertEquals("pending", facts.processingTaskIndexStatus(), "parse accepted 后 index 仍应等待轮询推进"); + assertEquals(result.scanTaskId(), facts.processingTaskId(), "响应 scanTaskId 必须对应 DB processing_task.task_id"); + assertTrue(StringUtils.hasText(facts.ragflowDatasetId()), "processing_task 必须写入 ragflow_dataset_id"); + assertTrue(StringUtils.hasText(facts.ragflowDocumentId()), "processing_task 必须写入 ragflow_document_id"); + assertTrue(facts.processingTaskStarted(), "processing_task 必须有 started_at"); + assertFalse(facts.processingTaskFinished(), "parse accepted 后 processing_task 不应被伪造为 finished"); + assertTrue(StringUtils.hasText(facts.processingTaskResultSummary()) + && !"{}".equals(facts.processingTaskResultSummary().trim()), + "processing_task 必须有非空 result_summary"); + + assertTrue(facts.createDatasetSucceeded(), "ragflow_call 必须记录 createDataset succeeded"); + assertTrue(facts.uploadDocumentsSucceeded(), "ragflow_call 必须记录 uploadDocuments succeeded"); + assertTrue(facts.startParseAccepted(), "ragflow_call 必须记录 startParseDocuments accepted"); + assertEquals(facts.ragflowDatasetId(), facts.datasetBindingDatasetId(), + "dataset binding 必须和 processing_task datasetId 一致"); + assertEquals(facts.ragflowDatasetId(), facts.documentBindingDatasetId(), + "document binding 必须和 processing_task datasetId 一致"); + assertEquals(facts.ragflowDocumentId(), facts.documentBindingDocumentId(), + "document binding 必须和 processing_task documentId 一致"); + assertTrue(facts.allRagflowResponseSummariesNonEmpty(), + "每条 RAGFlow business call 都必须有非空 response_summary"); + } + + private ExternalVerification verifyExternalRagFlowSearchable(ConfigurableApplicationContext context, + LiveResult result, + DbFacts facts) throws InterruptedException { + RagFlowKnowledgeRuntimeClient client = context.getBean(RagFlowKnowledgeRuntimeClient.class); + String correlationPrefix = "p1r5-knowledge-e2e-external-" + result.suffix(); + RagFlowKnowledgeRuntimeClient.RuntimeResult poll = waitUntilSearchable(client, facts.ragflowDatasetId(), + facts.ragflowDocumentId(), facts.processingTaskId(), correlationPrefix); + RagFlowKnowledgeRuntimeClient.RuntimeResult listChunks = client.listChunks( + new RagFlowKnowledgeRuntimeClient.ListChunksCommand(TENANT_ID, OWNER_USER_ID, result.kbId(), + facts.ragflowDatasetId(), facts.ragflowDocumentId(), Map.of("page", 1, "page_size", 5), + correlationPrefix + "-chunks", hash(facts.ragflowDocumentId() + "-chunks"), 1)); + assertSucceeded(listChunks, "listChunks"); + int chunksCount = countFromRedactedSummary(listChunks.redactedSummary(), "chunksCount"); + assertTrue(chunksCount > 0, () -> "listChunks 没有返回 chunk,不能作为外部检索证据:" + listChunks.redactedSummary()); + + RagFlowKnowledgeRuntimeClient.RuntimeResult retrieval = client.retrieveChunks( + new RagFlowKnowledgeRuntimeClient.RetrieveChunksCommand(TENANT_ID, OWNER_USER_ID, result.kbId(), + List.of(facts.ragflowDatasetId()), List.of(facts.ragflowDocumentId()), + "P1R-5 Knowledge live acceptance marker 是什么?", 3, 0.0, Map.of(), + correlationPrefix + "-retrieval", hash(facts.ragflowDocumentId() + "-retrieval"), 1)); + assertSucceeded(retrieval, "retrieveChunks"); + int retrievalChunksCount = countFromRedactedSummary(retrieval.redactedSummary(), "chunksCount"); + assertTrue(retrievalChunksCount > 0, + () -> "retrieveChunks 没有返回 chunk,不能作为外部检索证据:" + retrieval.redactedSummary()); + return new ExternalVerification(poll, listChunks, retrieval, chunksCount, retrievalChunksCount); + } + + private RagFlowKnowledgeRuntimeClient.RuntimeResult waitUntilSearchable(RagFlowKnowledgeRuntimeClient client, + String datasetId, + String documentId, + String processingTaskId, + String correlationPrefix) + throws InterruptedException { + long deadline = System.nanoTime() + + Duration.ofSeconds(intEnv("MUSE_P1R_RAGFLOW_PARSE_WAIT_SECONDS", 120)).toNanos(); + RagFlowKnowledgeRuntimeClient.RuntimeResult last = null; + int attempt = 1; + while (System.nanoTime() < deadline) { + last = client.pollDocumentStatuses(new RagFlowKnowledgeRuntimeClient.PollDocumentStatusesCommand( + TENANT_ID, OWNER_USER_ID, 0L, datasetId, List.of(documentId), processingTaskId, + correlationPrefix + "-poll-" + attempt, hash(documentId + "-poll-" + attempt), attempt)); + assertSucceeded(last, "pollDocumentStatuses"); + if (last.documentStatuses().stream().anyMatch(status -> "searchable".equals(status.museStatus()))) { + return last; + } + Thread.sleep(2_000L); + attempt++; + } + throw new AssertionError("RAGFlow document parse 未在等待窗口内进入 searchable,last=" + + (last == null ? "none" : last.redactedSummary())); + } + + private void assertSucceeded(RagFlowKnowledgeRuntimeClient.RuntimeResult result, String operation) { + assertNotNull(result, operation + " result 不能为空"); + assertEquals(RagFlowKnowledgeRuntimeClient.Status.SUCCEEDED, result.status(), + () -> operation + " failed: " + result.redactedSummary()); + } + + private DbFacts readDbFacts(LiveSettings settings, LiveResult result) { + try (Connection connection = DriverManager.getConnection(settings.jdbcUrl(), settings.jdbcUser(), + settings.jdbcPassword())) { + return new DbFacts( + count(connection, "muse_knowledge_base", "id = ?", result.kbId()), + count(connection, "muse_knowledge_document", "id = ?", result.documentId()), + count(connection, "muse_knowledge_document_version", "document_id = ?", result.documentId()), + count(connection, "muse_knowledge_processing_task", "task_id = ?", result.scanTaskId()), + count(connection, "muse_knowledge_command", "command_id = ?", result.commandId()), + count(connection, "muse_knowledge_ragflow_binding", "kb_id = ?", result.kbId()), + count(connection, "muse_knowledge_ragflow_binding", "kb_id = ? AND document_id IS NULL", result.kbId()), + count(connection, "muse_knowledge_ragflow_binding", + "kb_id = ? AND document_version_id = (SELECT id FROM muse_knowledge_document_version WHERE document_id = ?)", + result.kbId(), result.documentId()), + count(connection, "muse_knowledge_ragflow_call", "command_id = ?", result.commandId()), + stringValue(connection, "SELECT kb_type FROM muse_knowledge_base WHERE id = ?", result.kbId()), + stringValue(connection, "SELECT status FROM muse_knowledge_base WHERE id = ?", result.kbId()), + intValue(connection, "SELECT active_version FROM muse_knowledge_base WHERE id = ?", result.kbId()), + stringValue(connection, "SELECT status FROM muse_knowledge_command WHERE command_id = ?", + result.commandId()), + stringValue(connection, """ + SELECT scan_status FROM muse_knowledge_document_version + WHERE document_id = ? + ORDER BY version DESC LIMIT 1 + """, result.documentId()), + stringValue(connection, """ + SELECT processing_status FROM muse_knowledge_document_version + WHERE document_id = ? + ORDER BY version DESC LIMIT 1 + """, result.documentId()), + stringValue(connection, """ + SELECT task_id FROM muse_knowledge_processing_task + WHERE task_id = ? + """, result.scanTaskId()), + stringValue(connection, """ + SELECT status FROM muse_knowledge_processing_task + WHERE task_id = ? + """, result.scanTaskId()), + stringValue(connection, """ + SELECT materialization_status FROM muse_knowledge_processing_task + WHERE task_id = ? + """, result.scanTaskId()), + stringValue(connection, """ + SELECT scan_status FROM muse_knowledge_processing_task + WHERE task_id = ? + """, result.scanTaskId()), + stringValue(connection, """ + SELECT parse_status FROM muse_knowledge_processing_task + WHERE task_id = ? + """, result.scanTaskId()), + stringValue(connection, """ + SELECT index_status FROM muse_knowledge_processing_task + WHERE task_id = ? + """, result.scanTaskId()), + stringValue(connection, """ + SELECT ragflow_dataset_id FROM muse_knowledge_processing_task + WHERE task_id = ? + """, result.scanTaskId()), + stringValue(connection, """ + SELECT ragflow_document_id FROM muse_knowledge_processing_task + WHERE task_id = ? + """, result.scanTaskId()), + booleanValue(connection, """ + SELECT started_at IS NOT NULL FROM muse_knowledge_processing_task + WHERE task_id = ? + """, result.scanTaskId()), + booleanValue(connection, """ + SELECT finished_at IS NOT NULL FROM muse_knowledge_processing_task + WHERE task_id = ? + """, result.scanTaskId()), + stringValue(connection, """ + SELECT result_summary::text FROM muse_knowledge_processing_task + WHERE task_id = ? + """, result.scanTaskId()), + booleanValue(connection, """ + SELECT EXISTS ( + SELECT 1 FROM muse_knowledge_ragflow_call + WHERE tenant_id = ? AND deleted = FALSE AND command_id = ? + AND operation = 'createDataset' AND status = 'succeeded' + ) + """, TENANT_ID, result.commandId()), + booleanValue(connection, """ + SELECT EXISTS ( + SELECT 1 FROM muse_knowledge_ragflow_call + WHERE tenant_id = ? AND deleted = FALSE AND command_id = ? + AND operation = 'uploadDocuments' AND status = 'succeeded' + ) + """, TENANT_ID, result.commandId()), + booleanValue(connection, """ + SELECT EXISTS ( + SELECT 1 FROM muse_knowledge_ragflow_call + WHERE tenant_id = ? AND deleted = FALSE AND command_id = ? + AND operation = 'startParseDocuments' AND status = 'accepted' + ) + """, TENANT_ID, result.commandId()), + stringValue(connection, """ + SELECT ragflow_dataset_id FROM muse_knowledge_ragflow_binding + WHERE tenant_id = ? AND deleted = FALSE AND kb_id = ? AND document_id IS NULL + ORDER BY id DESC LIMIT 1 + """, TENANT_ID, result.kbId()), + stringValue(connection, """ + SELECT ragflow_dataset_id FROM muse_knowledge_ragflow_binding + WHERE tenant_id = ? AND deleted = FALSE AND document_id = ? + ORDER BY id DESC LIMIT 1 + """, TENANT_ID, result.documentId()), + stringValue(connection, """ + SELECT ragflow_document_id FROM muse_knowledge_ragflow_binding + WHERE tenant_id = ? AND deleted = FALSE AND document_id = ? + ORDER BY id DESC LIMIT 1 + """, TENANT_ID, result.documentId()), + booleanValue(connection, """ + SELECT NOT EXISTS ( + SELECT 1 FROM muse_knowledge_ragflow_call + WHERE tenant_id = ? AND deleted = FALSE AND command_id = ? + AND (response_summary IS NULL OR response_summary = '{}'::jsonb) + ) + """, TENANT_ID, result.commandId())); + } catch (SQLException ex) { + throw new IllegalStateException("读取 P1R Knowledge live acceptance DB 事实失败", ex); + } + } + + private long count(Connection connection, String table, String where, Object... args) throws SQLException { + String sql = "SELECT COUNT(*) FROM " + table + " WHERE tenant_id = ? AND deleted = FALSE AND " + where; + try (PreparedStatement statement = connection.prepareStatement(sql)) { + statement.setLong(1, TENANT_ID); + bind(statement, 2, args); + try (ResultSet resultSet = statement.executeQuery()) { + assertTrue(resultSet.next(), "必须能读取计数: " + table); + return resultSet.getLong(1); + } + } + } + + private String stringValue(Connection connection, String sql, Object... args) throws SQLException { + try (PreparedStatement statement = connection.prepareStatement(sql)) { + bind(statement, 1, args); + try (ResultSet resultSet = statement.executeQuery()) { + return resultSet.next() ? resultSet.getString(1) : null; + } + } + } + + private int intValue(Connection connection, String sql, Object... args) throws SQLException { + try (PreparedStatement statement = connection.prepareStatement(sql)) { + bind(statement, 1, args); + try (ResultSet resultSet = statement.executeQuery()) { + return resultSet.next() ? resultSet.getInt(1) : 0; + } + } + } + + private boolean booleanValue(Connection connection, String sql, Object... args) throws SQLException { + try (PreparedStatement statement = connection.prepareStatement(sql)) { + bind(statement, 1, args); + try (ResultSet resultSet = statement.executeQuery()) { + return resultSet.next() && resultSet.getBoolean(1); + } + } + } + + private void bind(PreparedStatement statement, int startIndex, Object... args) throws SQLException { + for (int i = 0; i < args.length; i++) { + Object value = args[i]; + if (value instanceof Long longValue) { + statement.setLong(startIndex + i, longValue); + } else if (value instanceof Integer intValue) { + statement.setInt(startIndex + i, intValue); + } else { + statement.setString(startIndex + i, String.valueOf(value)); + } + } + } + + private ConfigurableApplicationContext liveContext(LiveSettings settings) { + Map properties = new LinkedHashMap<>(); + properties.put("muse.info.base-package", "cn.iocoder.muse.module.knowledge"); + properties.put("spring.datasource.url", settings.jdbcUrl()); + properties.put("spring.datasource.username", settings.jdbcUser()); + properties.put("spring.datasource.password", settings.jdbcPassword()); + properties.put("spring.datasource.driver-class-name", "org.postgresql.Driver"); + properties.put("spring.main.banner-mode", "off"); + properties.put("spring.main.lazy-initialization", "true"); + properties.put("mybatis-plus.global-config.db-config.id-type", "AUTO"); + properties.put("muse.knowledge.ragflow.base-url", settings.ragflowBaseUrl()); + properties.put("muse.knowledge.ragflow.api-key", settings.ragflowApiKey()); + properties.put("muse.knowledge.ragflow.timeout-seconds", settings.ragflowTimeoutSeconds()); + properties.put("muse.knowledge.ragflow.retry-budget", settings.ragflowRetryBudget()); + properties.put("muse.knowledge.ragflow.graphrag-attribution-ready", "false"); + + return new SpringApplicationBuilder(LiveAcceptanceConfiguration.class) + .web(WebApplicationType.NONE) + .initializers(applicationContext -> applicationContext.getEnvironment().getPropertySources() + .addFirst(new MapPropertySource("p1r-knowledge-live-acceptance", properties))) + .properties(properties) + .run(); + } + + private DatabaseEngineFacts migrateIsolatedTestDatabase(LiveSettings settings) { + silenceFlywayInfoLogs(); + assertSafeJdbcUrl(settings.jdbcUrl()); + DatabaseEngineFacts dbEngine = assertPostgresqlTestDatabase(settings); + System.out.println(JsonUtils.toJsonString(Map.of( + "acceptance", "p1r5-knowledge-runtime-e2e-live", + "preCleanDbEngine", dbEngineEvidence(dbEngine)))); + Flyway flyway = Flyway.configure() + .dataSource(settings.jdbcUrl(), settings.jdbcUser(), settings.jdbcPassword()) + .locations(resolveMuseSqlLocation(settings.flywayLocations())) + .schemas("public") + .defaultSchema("public") + .target(TARGET_VERSION) + .cleanDisabled(false) + .load(); + flyway.clean(); + flyway.migrate(); + MigrationInfo current = flyway.info().current(); + assertEquals(TARGET_VERSION, Objects.requireNonNull(current, "必须存在当前 Flyway 版本") + .getVersion().getVersion(), "P1R-5 Knowledge live acceptance 必须迁移到 V14 schema"); + return dbEngine; + } + + private DatabaseEngineFacts assertPostgresqlTestDatabase(LiveSettings settings) { + try (Connection connection = DriverManager.getConnection(settings.jdbcUrl(), settings.jdbcUser(), + settings.jdbcPassword()); + PreparedStatement statement = connection.prepareStatement("SELECT current_database(), version()"); + ResultSet resultSet = statement.executeQuery()) { + assertTrue(resultSet.next(), "Flyway clean 前必须能读取 PostgreSQL engine 信息"); + String databaseName = resultSet.getString(1); + String version = resultSet.getString(2); + assertTrue(StringUtils.hasText(databaseName) && databaseName.endsWith("_test"), + "Flyway clean 前 current_database() 必须是 _test 隔离库: " + databaseName); + assertTrue(StringUtils.hasText(version) && version.contains("PostgreSQL"), + "Flyway clean 前 version() 必须来自 PostgreSQL: " + versionSummary(version)); + return new DatabaseEngineFacts(databaseName, version); + } catch (SQLException ex) { + throw new IllegalStateException("Flyway clean 前检查 PostgreSQL _test 数据库失败: " + + maskedUrl(settings.jdbcUrl()), ex); + } + } + + private Map redactedEvidence(LiveSettings settings, LiveResult result, DbFacts facts, + ExternalVerification externalVerification, + DatabaseEngineFacts dbEngine) { + Map evidence = new LinkedHashMap<>(); + evidence.put("acceptance", "p1r5-knowledge-runtime-e2e-live"); + evidence.put("endpoint", endpointSummary(settings.ragflowBaseUrl(), RAGFLOW_API_PATH)); + evidence.put("apiKey", secretFingerprint(settings.ragflowApiKey())); + evidence.put("jdbcUrl", maskedUrl(settings.jdbcUrl())); + evidence.put("flywayTargetVersion", TARGET_VERSION); + evidence.put("dbEngine", dbEngineEvidence(dbEngine)); + evidence.put("request", requestEvidence(result, facts)); + evidence.put("response", responseEvidence(result, facts)); + evidence.put("db", dbEvidence(facts)); + evidence.put("externalVerification", externalEvidence(externalVerification)); + return evidence; + } + + private Map requestEvidence(LiveResult result, DbFacts facts) { + Map request = new LinkedHashMap<>(); + request.put("tenantId", TENANT_ID); + request.put("ownerUserId", OWNER_USER_ID); + request.put("kbId", result.kbId()); + request.put("commandIdSha256Prefix", sha256Prefix(result.commandId(), 12)); + request.put("documentId", result.documentId()); + request.put("processingTaskId", result.scanTaskId()); + request.put("entryContentLength", length(liveDocumentContent(result.suffix()))); + request.put("entryContentSha256Prefix", sha256Prefix(liveDocumentContent(result.suffix()), 12)); + request.put("ragflowDatasetId", facts.ragflowDatasetId()); + request.put("ragflowDocumentId", facts.ragflowDocumentId()); + return request; + } + + private Map responseEvidence(LiveResult result, DbFacts facts) { + Map response = new LinkedHashMap<>(); + response.put("uploadResponseProcessingStatus", "pending_scan"); + response.put("dbTaskStatus", facts.processingTaskStatus()); + response.put("dbParseStatus", facts.processingTaskParseStatus()); + response.put("dbIndexStatus", facts.processingTaskIndexStatus()); + response.put("commandStatus", facts.commandStatus()); + response.put("processingTaskStarted", facts.processingTaskStarted()); + response.put("processingTaskFinished", facts.processingTaskFinished()); + response.put("processingTaskResultSummaryLength", length(facts.processingTaskResultSummary())); + response.put("processingTaskResultSummarySha256Prefix", + sha256Prefix(defaultValue(facts.processingTaskResultSummary(), ""), 12)); + response.put("documentId", result.documentId()); + return response; + } + + private Map dbEvidence(DbFacts facts) { + Map db = new LinkedHashMap<>(); + db.put("knowledgeBaseCount", facts.knowledgeBaseCount()); + db.put("documentCount", facts.documentCount()); + db.put("documentVersionCount", facts.documentVersionCount()); + db.put("processingTaskCount", facts.processingTaskCount()); + db.put("commandCount", facts.commandCount()); + db.put("bindingCount", facts.bindingCount()); + db.put("datasetBindingCount", facts.datasetBindingCount()); + db.put("documentBindingCount", facts.documentBindingCount()); + db.put("ragflowCallCount", facts.ragflowCallCount()); + db.put("createDatasetSucceeded", facts.createDatasetSucceeded()); + db.put("uploadDocumentsSucceeded", facts.uploadDocumentsSucceeded()); + db.put("startParseAccepted", facts.startParseAccepted()); + db.put("allRagflowResponseSummariesNonEmpty", facts.allRagflowResponseSummariesNonEmpty()); + return db; + } + + private Map externalEvidence(ExternalVerification verification) { + Map evidence = new LinkedHashMap<>(); + evidence.put("scope", "external_verification_not_business_persisted"); + evidence.put("pollDocumentStatuses", operationSummary(verification.poll())); + evidence.put("documentStatuses", verification.poll().documentStatuses()); + evidence.put("listChunks", operationSummary(verification.listChunks())); + evidence.put("retrieveChunks", operationSummary(verification.retrieval())); + evidence.put("chunksCount", verification.chunksCount()); + evidence.put("retrievalChunksCount", verification.retrievalChunksCount()); + return evidence; + } + + private Map operationSummary(RagFlowKnowledgeRuntimeClient.RuntimeResult result) { + Map summary = new LinkedHashMap<>(); + summary.put("operation", result.operation().wireName()); + summary.put("status", result.status()); + summary.put("failureClass", result.failureClass()); + summary.put("durationMillis", result.durationMillis()); + summary.put("externalIdPresent", result.externalId() != null && !result.externalId().isBlank()); + return summary; + } + + private Map dbEngineEvidence(DatabaseEngineFacts facts) { + Map engine = new LinkedHashMap<>(); + engine.put("databaseNameSha256Prefix", sha256Prefix(facts.databaseName(), 12)); + engine.put("databaseNameSuffix", facts.databaseName().endsWith("_test") ? "_test" : ""); + engine.put("versionSummary", versionSummary(facts.version())); + return engine; + } + + private Map endpointSummary(String baseUrl, String path) { + URI uri = URI.create(trimRight(baseUrl, "/") + path); + Map endpoint = new LinkedHashMap<>(); + endpoint.put("scheme", uri.getScheme()); + endpoint.put("host", uri.getHost()); + endpoint.put("port", uri.getPort()); + endpoint.put("path", uri.getPath()); + return endpoint; + } + + private Map secretFingerprint(String secret) { + Map fingerprint = new LinkedHashMap<>(); + fingerprint.put("length", secret == null ? 0 : secret.length()); + fingerprint.put("sha256Prefix", sha256Prefix(secret == null ? "" : secret, 12)); + return fingerprint; + } + + private int countFromRedactedSummary(String redactedSummary, String fieldName) { + JsonNode node = JsonUtils.parseTree(redactedSummary == null || redactedSummary.isBlank() ? "{}" : redactedSummary); + return node.has(fieldName) && node.path(fieldName).canConvertToInt() ? node.path(fieldName).asInt() : 0; + } + + private static String requiredPropertyOrEnv(String propertyName, String envName) { + String value = System.getProperty(propertyName); + if (!StringUtils.hasText(value)) { + value = System.getenv(envName); + } + assertTrue(StringUtils.hasText(value), "缺少必需系统属性或环境变量: " + propertyName + " / " + envName); + return value.trim(); + } + + private static String requiredEnv(String name) { + String value = System.getenv(name); + assertTrue(StringUtils.hasText(value), "缺少必需环境变量: " + name); + return value.trim(); + } + + private static String requiredPasswordEnvironment() { + String password = firstNonBlankEnvironment("P1R_FLYWAY_PASSWORD", "MUSE_POSTGRES_PASSWORD"); + assertTrue(password != null, "缺少必需数据库密码环境变量: P1R_FLYWAY_PASSWORD 或 MUSE_POSTGRES_PASSWORD"); + return password; + } + + private static String firstNonBlankEnvironment(String... names) { + // 数据库密码不允许从 system property 读取,避免被 Surefire XML 或命令历史记录。 + for (String name : names) { + String value = System.getenv(name); + if (StringUtils.hasText(value)) { + return value; + } + } + return null; + } + + private static int intEnv(String name, int defaultValue) { + String value = System.getenv(name); + if (!StringUtils.hasText(value)) { + return defaultValue; + } + return Integer.parseInt(value.trim()); + } + + private static String resolveMuseSqlLocation(String requestedLocations) { + assertEquals("filesystem:sql/muse", requestedLocations, + "P1R live acceptance 要求显式使用 filesystem:sql/muse"); + Path current = Path.of(System.getProperty("user.dir")).toAbsolutePath(); + for (Path cursor = current; cursor != null; cursor = cursor.getParent()) { + Path candidate = cursor.resolve("sql/muse"); + if (Files.isDirectory(candidate)) { + return "filesystem:" + candidate; + } + } + throw new IllegalStateException("无法从当前目录向上找到 sql/muse: " + current); + } + + private static void assertSafeJdbcUrl(String url) { + assertTrue(url.startsWith(POSTGRESQL_JDBC_PREFIX), + "JDBC URL 必须显式使用 jdbc:postgresql://,避免 Flyway clean 误连其它数据库: " + maskedUrl(url)); + URI uri = parseJdbcUriForAssertion(url); + assertEquals("postgresql", uri.getScheme(), + "JDBC URL 必须解析为 PostgreSQL scheme: " + maskedUrl(url)); + assertTrue(StringUtils.hasText(uri.getHost()), + "JDBC URL 必须显式包含 host: " + maskedUrl(url)); + assertFalse(StringUtils.hasText(uri.getRawUserInfo()), + "JDBC URL 不能携带 authority/userinfo,用户名走属性/环境变量,密码只走环境变量: " + maskedUrl(url)); + assertNoCredentialQuery(url); + String databaseName = jdbcDatabaseName(uri); + assertTrue(databaseName.endsWith("_test"), + "JDBC URL 必须指向 _test 后缀隔离库,避免清理非测试库: " + maskedUrl(url)); + assertFalse("muse_local".equals(databaseName), + "P1R live acceptance 禁止默认连接 muse_local"); + } + + private static URI parseJdbcUriForAssertion(String url) { + try { + return URI.create(url.substring(JDBC_URI_PREFIX.length())); + } catch (IllegalArgumentException ex) { + // JDK URI 解析异常会回显原始 authority/query;断言失败只允许输出脱敏后的 JDBC 摘要。 + throw new AssertionError("JDBC URL 必须可解析为 PostgreSQL URI: " + maskedUrl(url)); + } + } + + private static String jdbcDatabaseName(URI uri) { + String path = uri.getRawPath(); + assertTrue(StringUtils.hasText(path) && path.length() > 1, + "JDBC URL 必须包含数据库名"); + String databaseName = path.substring(path.lastIndexOf('/') + 1); + assertTrue(StringUtils.hasText(databaseName), + "JDBC URL 必须包含数据库名"); + return databaseName; + } + + private static void assertNoCredentialQuery(String url) { + int queryStart = url.indexOf('?'); + if (queryStart < 0) { + return; + } + String query = url.substring(queryStart + 1); + for (String parameter : query.split("&")) { + String key = parameter; + int equalsStart = key.indexOf('='); + if (equalsStart >= 0) { + key = key.substring(0, equalsStart); + } + assertFalse(isCredentialQueryKey(key), + "JDBC URL 不能携带凭据 query 参数;用户名走属性/环境变量,密码只走环境变量"); + } + } + + private static boolean isCredentialQueryKey(String rawKey) { + String key = rawKey.trim().toLowerCase(Locale.ROOT).replace('-', '_'); + return CREDENTIAL_QUERY_KEYS.contains(key) + || key.endsWith("_token") + || key.endsWith("_secret") + || key.endsWith("_password"); + } + + private static String maskedUrl(String url) { + if (!StringUtils.hasText(url)) { + return ""; + } + // 失败消息只能暴露 JDBC 类型、占位 authority、测试库名和 query 是否存在;解析失败也不能回显原始串。 + String trimmed = url.trim(); + if (!trimmed.startsWith(JDBC_URI_PREFIX)) { + return ""; + } + + String jdbcBody = trimmed.substring(JDBC_URI_PREFIX.length()); + String scheme = jdbcScheme(jdbcBody); + if (!StringUtils.hasText(scheme)) { + return JDBC_URI_PREFIX + ""; + } + if (!jdbcBody.regionMatches(true, scheme.length(), "://", 0, 3)) { + return JDBC_URI_PREFIX + scheme + ":"; + } + if ("postgresql".equals(scheme)) { + return maskedPostgresqlUrl(jdbcBody); + } + return maskedHierarchicalJdbcUrl(scheme, jdbcBody); + } + + private static String jdbcScheme(String jdbcBody) { + int schemeEnd = jdbcBody.indexOf(':'); + if (schemeEnd <= 0) { + return ""; + } + String scheme = jdbcBody.substring(0, schemeEnd).trim().toLowerCase(Locale.ROOT); + return scheme.matches("[a-z][a-z0-9+.-]*") ? scheme : ""; + } + + private static String maskedPostgresqlUrl(String jdbcBody) { + try { + URI uri = URI.create(jdbcBody); + String database = safeJdbcDatabaseName(uri.getRawPath(), ""); + String port = uri.getPort() < 0 ? "" : ":" + uri.getPort(); + String suffix = StringUtils.hasText(uri.getRawQuery()) ? "?" : ""; + return POSTGRESQL_JDBC_PREFIX + "" + port + "/" + database + suffix; + } catch (RuntimeException ignored) { + return POSTGRESQL_JDBC_PREFIX + "/" + fallbackJdbcDatabaseName(jdbcBody, "") + + querySuffix(jdbcBody); + } + } + + private static String maskedHierarchicalJdbcUrl(String scheme, String jdbcBody) { + try { + URI uri = URI.create(jdbcBody); + String database = safeJdbcDatabaseName(uri.getRawPath(), ""); + String databaseSuffix = StringUtils.hasText(database) ? "/" + database : ""; + String querySuffix = StringUtils.hasText(uri.getRawQuery()) ? "?" : ""; + return JDBC_URI_PREFIX + scheme + "://" + databaseSuffix + querySuffix; + } catch (RuntimeException ignored) { + String database = fallbackJdbcDatabaseName(jdbcBody, ""); + String databaseSuffix = StringUtils.hasText(database) ? "/" + database : ""; + return JDBC_URI_PREFIX + scheme + "://" + databaseSuffix + querySuffix(jdbcBody); + } + } + + private static String safeJdbcDatabaseName(String rawPath, String fallback) { + if (!StringUtils.hasText(rawPath) || rawPath.length() <= 1) { + return fallback; + } + String database = rawPath.substring(rawPath.lastIndexOf('/') + 1); + if (!StringUtils.hasText(database)) { + return fallback; + } + return database.matches("[A-Za-z0-9._%-]+") ? database : fallback; + } + + private static String fallbackJdbcDatabaseName(String jdbcBody, String fallback) { + int authorityStart = jdbcBody.indexOf("://"); + if (authorityStart < 0) { + return fallback; + } + int pathStart = jdbcBody.indexOf('/', authorityStart + 3); + if (pathStart < 0) { + return fallback; + } + int pathEnd = firstIndexOf(jdbcBody, pathStart + 1, '?', '#'); + String path = jdbcBody.substring(pathStart, pathEnd); + return safeJdbcDatabaseName(path, fallback); + } + + private static int firstIndexOf(String value, int start, char... candidates) { + int result = value.length(); + for (char candidate : candidates) { + int index = value.indexOf(candidate, start); + if (index >= 0 && index < result) { + result = index; + } + } + return result; + } + + private static String querySuffix(String value) { + return value.indexOf('?') < 0 ? "" : "?"; + } + + private static String versionSummary(String version) { + if (!StringUtils.hasText(version)) { + return ""; + } + String normalized = version.replaceAll("\\s+", " ").trim(); + return normalized.length() > 96 ? normalized.substring(0, 96) : normalized; + } + + private static void silenceFlywayInfoLogs() { + try { + Object flywayLogger = LoggerFactory.getLogger("org.flywaydb"); + Class levelClass = Class.forName("ch.qos.logback.classic.Level"); + Object warnLevel = levelClass.getField("WARN").get(null); + flywayLogger.getClass().getMethod("setLevel", levelClass).invoke(flywayLogger, warnLevel); + } catch (ReflectiveOperationException | LinkageError ignored) { + // 日志实现不是 logback 时不影响验收;测试自身仍只输出 masked URL 和摘要。 + } + } + + private static boolean externalAcceptanceEnabled() { + return "true".equalsIgnoreCase(System.getenv(ACCEPTANCE_ENV)); + } + + private static String hash(String value) { + return sha256Prefix(value == null ? "" : value, 64); + } + + private static String sha256Prefix(String value, int length) { + try { + MessageDigest digest = MessageDigest.getInstance("SHA-256"); + String hex = HexFormat.of().formatHex(digest.digest(value.getBytes(StandardCharsets.UTF_8))); + return hex.substring(0, Math.min(length, hex.length())); + } catch (NoSuchAlgorithmException e) { + throw new IllegalStateException("JDK 缺少 SHA-256 摘要算法", e); + } + } + + private static String trimRight(String value, String suffix) { + String result = value == null ? "" : value.trim(); + while (result.endsWith(suffix)) { + result = result.substring(0, result.length() - suffix.length()); + } + return result; + } + + private static String defaultValue(String value, String fallback) { + return value == null ? fallback : value; + } + + private static int length(String value) { + return value == null ? 0 : value.length(); + } + + private record LiveSettings(String jdbcUrl, + String jdbcUser, + String jdbcPassword, + String flywayLocations, + String ragflowBaseUrl, + String ragflowApiKey, + int ragflowTimeoutSeconds, + int ragflowRetryBudget) { + + private static LiveSettings fromEnvironment() { + String jdbcUrl = requiredPropertyOrEnv(JDBC_URL_PROPERTY, JDBC_URL_ENV); + assertSafeJdbcUrl(jdbcUrl); + return new LiveSettings( + jdbcUrl, + requiredPropertyOrEnv(JDBC_USER_PROPERTY, JDBC_USER_ENV), + requiredPasswordEnvironment(), + System.getProperty(FLYWAY_LOCATION_PROPERTY, "filesystem:sql/muse"), + requiredEnv("MUSE_KNOWLEDGE_RAGFLOW_BASE_URL"), + requiredEnv("MUSE_KNOWLEDGE_RAGFLOW_API_KEY"), + intEnv("MUSE_KNOWLEDGE_RAGFLOW_TIMEOUT_SECONDS", 30), + intEnv("MUSE_KNOWLEDGE_RAGFLOW_RETRY_BUDGET", 0)); + } + } + + private record LiveResult(Long kbId, + String commandId, + Long documentId, + String scanTaskId, + String suffix) { + } + + private record DbFacts(long knowledgeBaseCount, + long documentCount, + long documentVersionCount, + long processingTaskCount, + long commandCount, + long bindingCount, + long datasetBindingCount, + long documentBindingCount, + long ragflowCallCount, + String kbType, + String kbStatus, + int kbActiveVersion, + String commandStatus, + String documentVersionScanStatus, + String documentVersionProcessingStatus, + String processingTaskId, + String processingTaskStatus, + String materializationStatus, + String processingTaskScanStatus, + String processingTaskParseStatus, + String processingTaskIndexStatus, + String ragflowDatasetId, + String ragflowDocumentId, + boolean processingTaskStarted, + boolean processingTaskFinished, + String processingTaskResultSummary, + boolean createDatasetSucceeded, + boolean uploadDocumentsSucceeded, + boolean startParseAccepted, + String datasetBindingDatasetId, + String documentBindingDatasetId, + String documentBindingDocumentId, + boolean allRagflowResponseSummariesNonEmpty) { + } + + private record ExternalVerification(RagFlowKnowledgeRuntimeClient.RuntimeResult poll, + RagFlowKnowledgeRuntimeClient.RuntimeResult listChunks, + RagFlowKnowledgeRuntimeClient.RuntimeResult retrieval, + int chunksCount, + int retrievalChunksCount) { + } + + private record DatabaseEngineFacts(String databaseName, + String version) { + } + + @TestConfiguration + @Import({ + MuseDataSourceAutoConfiguration.class, + DataSourceAutoConfiguration.class, + DataSourceTransactionManagerAutoConfiguration.class, + MuseMybatisAutoConfiguration.class, + MybatisPlusAutoConfiguration.class, + MybatisPlusJoinAutoConfiguration.class, + MuseKnowledgeDocumentService.class, + MuseKnowledgeCommandService.class, + MuseKnowledgeAuditService.class, + MuseKnowledgeMaterializationService.class, + MuseKnowledgeProcessingTaskService.class, + SpringUtil.class + }) + static class LiveAcceptanceConfiguration { + + @Bean + RagFlowKnowledgeRuntimeClient ragFlowKnowledgeRuntimeClient() { + return new HttpRagFlowKnowledgeRuntimeClient(requiredEnv("MUSE_KNOWLEDGE_RAGFLOW_BASE_URL"), + requiredEnv("MUSE_KNOWLEDGE_RAGFLOW_API_KEY"), + Duration.ofSeconds(intEnv("MUSE_KNOWLEDGE_RAGFLOW_TIMEOUT_SECONDS", 30)), + intEnv("MUSE_KNOWLEDGE_RAGFLOW_RETRY_BUDGET", 0), false); + } + + @Bean + @Primary + KnowledgeFileFacade p1rLiveKnowledgeFileFacade(ObjectProvider fileApiProvider) { + return new KnowledgeFileFacade(fileApiProvider) { + + @Override + public MaterializedFile materializeText(Long kbId, String fileName, String entryContent) { + byte[] content = entryContent.getBytes(StandardCharsets.UTF_8); + // test-only passed scanner 只替代“扫描服务已通过”这一外部前置条件;RAGFlow adapter 和 Muse DB 写入仍走真实主业务链路。 + return MaterializedFile.materialized("p1r-live-file-ref-" + kbId + "-" + sha256Prefix(fileName, 8), + fileName, (long) content.length, "text/plain; charset=UTF-8", + sha256Prefix(entryContent, 64), "passed", content); + } + }; + } + } +}