From 8b3258e22fa84fc5539e890039327e7c7f2c4a7f Mon Sep 17 00:00:00 2001 From: lili Date: Fri, 26 Jun 2026 11:58:56 -0700 Subject: [PATCH] =?UTF-8?q?test(market):=20U-verify=20D0-fork=20=E7=AB=AF?= =?UTF-8?q?=E5=88=B0=E7=AB=AF=E7=9C=9F=E9=AA=8C=E6=94=B6(=E7=9C=9F=20PG+?= =?UTF-8?q?=E7=9C=9F=20RAGFlow=E3=80=81=E7=AC=AC=E5=9B=9B=E5=8D=95?= =?UTF-8?q?=E5=85=83)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit D0-fork 最后单元:真 PG17+真 RAGFlow 端到端验收、2 IT 真绿(独立复跑确认 57s)。P1rMarketKbForkMaterializationIT(muse-server、跨 knowledge+market 两 BC 须同 Spring 上下文):①fork→安装→物化 installed_ref→检索命中(no_dataset 消除、副本与活 dataset 物理分离);②发布者私有不泄露真证(发布者上架后往活 dataset 真传私有 magic、安装者检索副本拿 magic 当 query 仍只返公开 chunk 不含 magic→d0ForkIsolationProven、U0 spike 反向基线);③幂等+多安装者共享1副本;④AFTER_COMMIT 时序(committedAssetForked/rolledBackAssetForked false:提交才 fork 回滚不 fork);⑤_test 库铁律(muse_p1r_fork_mat_test 绝未指 muse_slice_live)。 market-install-kb-retrieval.spec.ts(studio e2e):检索命中+私有不泄露无独立 UI 入口由 IT 证;安装物化 e2e 标 test.fixme 待 global-setup 补 fork 就绪资产(不伪绿、骨架就位 tsc 过)。召回回滚(D9)如实标开放项:recallAsset 仍 writeSourceStatusBlocked、knowledge 无召回消费者→不自动停用 installed_ref;需新增 market→knowledge 召回传播才闭环(不伪造、不引用不存在的 muse_source_propagation_target)。 Co-Authored-By: Claude Opus 4.8 (1M context) --- .../api/P1rMarketKbForkMaterializationIT.java | 1248 +++++++++++++++++ .../e2e/market-install-kb-retrieval.spec.ts | 144 ++ 2 files changed, 1392 insertions(+) create mode 100644 muse-cloud/muse-server/src/test/java/cn/iocoder/muse/server/framework/api/P1rMarketKbForkMaterializationIT.java create mode 100644 muse-studio/e2e/market-install-kb-retrieval.spec.ts diff --git a/muse-cloud/muse-server/src/test/java/cn/iocoder/muse/server/framework/api/P1rMarketKbForkMaterializationIT.java b/muse-cloud/muse-server/src/test/java/cn/iocoder/muse/server/framework/api/P1rMarketKbForkMaterializationIT.java new file mode 100644 index 00000000..73ef5764 --- /dev/null +++ b/muse-cloud/muse-server/src/test/java/cn/iocoder/muse/server/framework/api/P1rMarketKbForkMaterializationIT.java @@ -0,0 +1,1248 @@ +package cn.iocoder.muse.server.framework.api; + +import cn.hutool.extra.spring.SpringUtil; +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.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.api.MuseKnowledgeRetrievalApi; +import cn.iocoder.muse.module.knowledge.api.MuseKnowledgeRetrievalApiImpl; +import cn.iocoder.muse.module.knowledge.application.muse.KnowledgeMarketForkService; +import cn.iocoder.muse.module.knowledge.application.muse.KnowledgeMarketListedEventConsumer; +import cn.iocoder.muse.module.knowledge.application.muse.MuseKnowledgeAuditService; +import cn.iocoder.muse.module.knowledge.application.muse.MuseKnowledgeMarketListedEvent; +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.market.api.asset.MarketAssetForkApi; +import cn.iocoder.muse.module.market.api.asset.MarketAssetForkApiImpl; +import cn.iocoder.muse.module.market.api.asset.MarketAssetSourceApi; +import cn.iocoder.muse.module.market.api.asset.MarketAssetSourceApiImpl; +import cn.iocoder.muse.module.market.api.asset.dto.MarketAssetSourceRespDTO; +import cn.iocoder.muse.module.market.api.asset.event.MarketKbListedEvent; +import com.baomidou.mybatisplus.autoconfigure.MybatisPlusAutoConfiguration; +import com.github.yulichang.autoconfigure.MybatisPlusJoinAutoConfiguration; +import org.apache.ibatis.annotations.Mapper; +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.mybatis.spring.annotation.MapperScan; +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.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 java.util.concurrent.ConcurrentHashMap; + +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.assertNull; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.junit.jupiter.api.Assumptions.assumeTrue; + +/** + * D0-fork(临时-04 U-verify)真 PG + 真 RAGFlow 端到端验收: + * 「知识库上架 → fork 公开副本 → 安装物化 installed_ref → 检索命中 + 发布者私有不泄露」。 + * + *

放在 muse-server 测试层(与其它 {@code P1r*IT} 一致)的原因:D0-fork 跨 knowledge + market 两个 BC——fork 服务、 + * installed_ref 物化、检索在 knowledge 模块,资产来源读 / fork 状态写回端口在 market 模块。两域 Bean 必须在同一 Spring + * 上下文共存,唯有 server 测试层能装配跨模块;knowledge 模块自身测试目录只能跑纯单测 / 本域 live(如 + * {@code P1rRagFlowLiveAcceptanceIT})。本类同时 {@code @MapperScan} 纳入 knowledge + market 两个 DAL 包。

+ * + *

opt-in(默认跳过):本测试会 {@code flyway.clean()} 删 _test 库全表并真打 RAGFlow,只有显式 + * {@code MUSE_P1R_EXTERNAL_ACCEPTANCE=true} 才运行(与 {@code P1rKnowledgeRuntimeEndToEndLiveAcceptanceIT} 同口径)。 + * RAGFlow / PG 在 mini-infra 100.64.0.8 在线。flyway 目标库铁律:必须 _test 后缀,绝不指 muse_slice_live。

+ * + *

核心安全断言(D0-fork 全部价值,承 U0 反向基线):U0 spike({@code MuseKnowledgeRetrievalApiImplTest})已钉死 + * 「安装者 binding 指向发布者活 dataset 时,发布者私有 magic chunk 被原样返回」。本 IT 证明 D0-fork 后——安装者 + * binding 指向的是公开副本 dataset,发布者上架后才加的私有 magic chunk 检索不到(副本物理不含该文档)。 + * 从「原样泄露」转为「物理不可见」即验收通过。这一条必须用真 RAGFlow 证(mock 无法证物理隔离)。

+ * + *

凭据红线:RAGFlow base-url/api-key、PG 密码均经 env 注入;所有输出脱敏(只打 endpoint/库名后缀/id/状态/数量, + * key 只打长度 + sha256 前 12)。不在宿主跑 psql(用 JDBC)。不 rm dump.rdb/lefthook。

+ */ +class P1rMarketKbForkMaterializationIT { + + 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 Long TENANT_ID = 100L; + private static final Long PUBLISHER_USER_ID = 3001L; + private static final Long INSTALLER_A_USER_ID = 4001L; + private static final Long INSTALLER_B_USER_ID = 4002L; + + 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 clearTenantContext() { + TenantContextHolder.clear(); + } + + /** + * 主验收用例:上架 fork 公开副本 → 安装物化 → 检索命中(真 RAGFlow)+ 发布者私有不泄露 + 幂等 + 多安装者共享一份副本。 + * + *

编排(全部真 PG + 真 RAGFlow,单 JVM 单上下文):

+ *
    + *
  1. seed 发布者 KB + 3 个公开(searchable)文档版本 + storageRef → 内容由 stub FileApi 按 ref 回读; + * 再 seed 1 个发布者「上架后新增」的私有文档(含唯一 magic token,processingStatus 仍非 searchable 不算公开集); + * seed market 资产行(knowledge_base 类型,source_id=发布者 kbId)。
  2. + *
  3. 真跑 {@link KnowledgeMarketForkService#forkPublicSnapshot}(= AFTER_COMMIT 监听最终会调的入口)→ 真 createDataset + * + 逐文档 upload + startParse → forkStatus=ready;副本 datasetId 写回 market.tags。
  4. + *
  5. 关键安全步骤:把发布者私有 magic 文档也真传一份到发布者活 dataset(模拟「上架后继续往自己 KB 加私有内容」), + * 但它不在公开副本里——副本只 fork 了上架时刻的 3 个公开文档。
  6. + *
  7. 等公开副本里某公开文档在 RAGFlow 后台 searchable(startParse accepted 只是已受理,真正可检索靠后台推进)。
  8. + *
  9. 安装者 A、B 各物化 installed_ref({@code MarketAssetSourceApi} 读到 forkStatus=ready + 副本 datasetId → + * 建本地 installed_ref KB + binding 指向公开副本);幂等:A 再物化一次复用同一本地 kbId。
  10. + *
  11. 经真实 {@link MuseKnowledgeRetrievalApiImpl#retrieveForWork} 检索 A 的作品 → RAGFlow 真命中公开副本 chunk + * (no_dataset 消除);且返回 chunk 不含发布者私有 magic(D0-fork 隔离成立)。
  12. + *
+ */ + @Test + void shouldForkPublicSnapshotInstallMaterializeRetrieveAndNotLeakPublisherPrivate() throws Exception { + assumeTrue(externalAcceptanceEnabled(), + "D0-fork U-verify 真 PG + 真 RAGFlow 端到端验收未启用,设置 MUSE_P1R_EXTERNAL_ACCEPTANCE=true 后才运行"); + + LiveSettings settings = LiveSettings.fromEnvironment(); + DatabaseEngineFacts dbEngine = migrateIsolatedTestDatabase(settings); + String suffix = String.valueOf(Instant.now().toEpochMilli()); + // 发布者私有内容的唯一探针 token:只出现在「上架后新增的私有文档」里;检索副本绝不应返回它。 + String privateMagic = "MAGIC_PUBLISHER_PRIVATE_" + suffix + "_只有发布者该看见"; + String publicMarker = "fork-public-marker-" + suffix; + + try (ConfigurableApplicationContext context = liveContext(settings)) { + ForkSeed seed = TenantUtils.execute(TENANT_ID, + () -> seedPublisherAndAsset(context, settings, suffix, publicMarker, privateMagic)); + + // ① 真 fork 公开副本(forkPublicSnapshot = AFTER_COMMIT 监听最终调用的入口;时序由独立用例 + 既有消费者单测覆盖)。 + KnowledgeMarketForkService.ForkResult forkResult = TenantUtils.execute(TENANT_ID, + () -> context.getBean(KnowledgeMarketForkService.class).forkPublicSnapshot( + new MuseKnowledgeMarketListedEvent(TENANT_ID, seed.assetId(), seed.publisherKbId(), + PUBLISHER_USER_ID, seed.assetVersion()))); + assertEquals("ready", forkResult.forkStatus(), + "fork 必须收口 ready(全部公开文档 startParse accepted)"); + assertNotNull(forkResult.publicForkDatasetId(), "ready 时必须带回公开副本 datasetId"); + assertEquals(3, forkResult.copiedDocumentCount(), "必须复制上架时刻的 3 个公开文档"); + assertEquals(0, forkResult.failedDocumentCount(), "公开文档复制不应有失败"); + String publicForkDatasetId = forkResult.publicForkDatasetId(); + + // 副本 dataset 必须与发布者活 dataset 物理分离(不同 dataset id)——这是物理隔离的前提。 + assertNotEquals(seed.publisherLiveDatasetId(), publicForkDatasetId, + "公开副本 dataset 必须与发布者活 dataset 物理分离"); + + // fork 状态确实写回了 market 资产 tags(供安装侧 MarketAssetSourceApi 带回)。 + MarketAssetSourceRespDTO source = TenantUtils.execute(TENANT_ID, + () -> context.getBean(MarketAssetSourceApi.class).getAssetSource(seed.assetId())); + assertTrue(source.exists(), "资产应存在"); + assertEquals("knowledge_base", source.assetType(), "资产类型应为 knowledge_base"); + assertEquals("ready", source.forkStatus(), "market.tags 应记录 forkStatus=ready"); + assertEquals(publicForkDatasetId, source.publicForkDatasetId(), + "market.tags 带回的副本 datasetId 必须与 fork 结果一致"); + + // ② 关键安全步骤:发布者「上架后」把私有 magic 文档真传进自己的活 dataset(副本里没有这份)。 + uploadPrivateDocToPublisherLiveDataset(context, seed, privateMagic, suffix); + + // ③ 等公开副本里的公开文档在 RAGFlow 后台真正 searchable(startParse accepted ≠ 立即可检索)。 + String publicForkDocId = forkedPublicDocumentId(context, seed, publicForkDatasetId, suffix); + waitUntilSearchable(context, publicForkDatasetId, publicForkDocId, suffix); + + // ④ 安装者 A、B 各物化 installed_ref;A 物化两次验幂等(复用同一本地 kbId)。 + Long localKbIdA = TenantUtils.execute(TENANT_ID, + () -> materializeInstalledRef(context, INSTALLER_A_USER_ID, seed.assetId())); + Long localKbIdA2 = TenantUtils.execute(TENANT_ID, + () -> materializeInstalledRef(context, INSTALLER_A_USER_ID, seed.assetId())); + assertEquals(localKbIdA, localKbIdA2, "幂等:同一安装者对同一资产重复物化必复用同一本地 installed_ref kbId"); + Long localKbIdB = TenantUtils.execute(TENANT_ID, + () -> materializeInstalledRef(context, INSTALLER_B_USER_ID, seed.assetId())); + assertNotEquals(localKbIdA, localKbIdB, "多安装者:各建各自本地 installed_ref kbId"); + + // 多安装者共享一份副本:A、B 的 dataset binding 都指向同一 publicForkDatasetId(副本仅 1 份)。 + assertEquals(publicForkDatasetId, datasetIdOfInstalledRef(settings, localKbIdA), + "A 的 installed_ref binding 必须指向公开副本"); + assertEquals(publicForkDatasetId, datasetIdOfInstalledRef(settings, localKbIdB), + "B 的 installed_ref binding 必须指向同一公开副本(共享只读,副本仅 1 份)"); + + // 副本载体行 + 两安装者 binding 都指向同一 dataset,但 DB 中公开副本 dataset 只 fork 出 1 份。 + assertEquals(1, countForkCarrierRows(settings, seed.publisherKbId()), + "公开副本 dataset 应只 fork 出 1 份载体行(幂等不重建)"); + + // ⑤ 经真实检索链验证 A 的作品检索命中公开副本 + 不泄露发布者私有 magic(D0-fork 核心安全验证)。 + RetrievalEvidence evidence = TenantUtils.execute(TENANT_ID, () -> + retrieveAndAssertNoLeak(context, settings, localKbIdA, INSTALLER_A_USER_ID, + publicForkDatasetId, publicMarker, privateMagic)); + + System.out.println(JsonUtils.toJsonString(redactedEvidence(settings, dbEngine, seed, + publicForkDatasetId, localKbIdA, localKbIdB, evidence))); + } + } + + /** + * AFTER_COMMIT 事务时序真验(真 Spring 上下文 + 真事务 + 真消费者): + *
    + *
  • 提交的事务内发 {@link MarketKbListedEvent} → {@code @TransactionalEventListener(AFTER_COMMIT)} + * 消费者在提交后异步真触发 fork → PG 里出现该 (assetId,version) 的公开副本载体行(forkStatus 推进)。
  • + *
  • 回滚的事务内发同形态事件(另一 assetId)→ 消费者不触发 → PG 里无对应副本载体行 + * (审核回滚则绝不 fork 一个没上架成功的资产)。
  • + *
+ * 这把 U-fork 的「进程内事件 AFTER_COMMIT 异步 fork」接线从 mock 单测({@code KnowledgeMarketListedEventConsumerTest}) + * 提升到真 Spring 事务时序证据。fork 本体走真 RAGFlow,故同样 opt-in。 + */ + @Test + void afterCommitListener_shouldForkOnlyWhenPublishingTransactionCommits() throws Exception { + assumeTrue(externalAcceptanceEnabled(), + "AFTER_COMMIT 时序真验未启用,设置 MUSE_P1R_EXTERNAL_ACCEPTANCE=true 后才运行"); + + LiveSettings settings = LiveSettings.fromEnvironment(); + migrateIsolatedTestDatabase(settings); + String suffix = String.valueOf(Instant.now().toEpochMilli()); + + try (ConfigurableApplicationContext context = liveContext(settings)) { + // 提交分支与回滚分支各 seed 一个发布者 KB(含 1 个公开文档,让 fork 有内容可复制)+ 资产。 + TimingSeed committed = TenantUtils.execute(TENANT_ID, + () -> seedMinimalForkSubject(context, settings, suffix + "-commit")); + TimingSeed rolledBack = TenantUtils.execute(TENANT_ID, + () -> seedMinimalForkSubject(context, settings, suffix + "-rollback")); + + TransactionalListedEventPublisher publisher = context.getBean(TransactionalListedEventPublisher.class); + + // ① 提交事务内发事件 → AFTER_COMMIT 异步 fork 应真触发。 + TenantUtils.execute(TENANT_ID, () -> { + publisher.publishInCommittedTransaction(new MarketKbListedEvent(TENANT_ID, committed.assetId(), + committed.publisherKbId(), PUBLISHER_USER_ID, "1")); + return null; + }); + + // ② 回滚事务内发事件 → 消费者不应触发(事件不投递)。直接设租户上下文调用(不经 TenantUtils.execute, + // 后者会把回滚异常包成 RuntimeException 掩盖类型);publish 本身不依赖租户,消费者从事件自带 tenantId 重建。 + TenantContextHolder.setTenantId(TENANT_ID); + try { + publisher.publishInRolledBackTransaction(new MarketKbListedEvent(TENANT_ID, rolledBack.assetId(), + rolledBack.publisherKbId(), PUBLISHER_USER_ID, "1")); + } catch (IllegalStateException expected) { + // 预期的强制回滚异常:事务回滚 → AFTER_COMMIT 不回放(事件不投递)。 + } finally { + TenantContextHolder.clear(); + } + + // 等提交分支异步 fork 收口到「载体行出现」(异步 + 真 RAGFlow,给足窗口)。 + boolean committedForked = waitForForkCarrierRow(settings, committed.publisherKbId(), true); + assertTrue(committedForked, + "提交事务后 AFTER_COMMIT 消费者必须真触发 fork(PG 出现公开副本载体行)"); + + // 回滚分支:再观察一个窗口,确认始终无载体行(事件未投递、消费者未触发)。 + boolean rolledBackForked = waitForForkCarrierRow(settings, rolledBack.publisherKbId(), false); + assertFalse(rolledBackForked, + "回滚事务的上架事件不应投递,消费者不应触发 fork(PG 无公开副本载体行)"); + + System.out.println(JsonUtils.toJsonString(Map.of( + "acceptance", "p1r-market-kb-fork-after-commit-timing", + "committedAssetForked", committedForked, + "rolledBackAssetForked", rolledBackForked, + "afterCommitContractHeld", committedForked && !rolledBackForked))); + } + } + + /** 为时序用例 seed 最小 fork 主体:发布者 KB + 1 个公开文档(searchable)+ market 资产。 */ + private TimingSeed seedMinimalForkSubject(ConfigurableApplicationContext context, LiveSettings settings, + String suffix) { + StubFileApi fileApi = context.getBean(StubFileApi.class); + try (Connection connection = jdbc(settings)) { + connection.setAutoCommit(true); + Long publisherKbId = insertReturningId(connection, """ + INSERT INTO muse_knowledge_base(name, kb_type, owner_user_id, status, active_version, + revision, command_id, tenant_id) + VALUES (?, 'user', ?, 'active', 1, 1, ?, ?) + """, "timing-pub-kb-" + suffix, PUBLISHER_USER_ID, "timing-kb-" + suffix, TENANT_ID); + String storageRef = "timing-ref-" + suffix; + String content = "时序用例公开文档 " + suffix + ":写作技巧片段。"; + fileApi.put(storageRef, content.getBytes(StandardCharsets.UTF_8)); + Long docId = insertReturningId(connection, """ + INSERT INTO muse_knowledge_document(kb_id, title, file_name, file_size, mime_type, + storage_ref, scan_status, parse_status, index_status, + command_id, tenant_id) + VALUES (?, ?, ?, ?, 'text/plain', ?, 'passed', 'completed', 'completed', ?, ?) + """, publisherKbId, "时序公开文档", "timing.txt", + (long) content.getBytes(StandardCharsets.UTF_8).length, storageRef, + "timing-doc-" + suffix, TENANT_ID); + insertDocumentVersion(connection, docId, storageRef, "searchable", suffix, 1); + Long assetId = insertReturningId(connection, """ + INSERT INTO muse_market_asset(name, asset_type, source_id, publisher_id, listing_status, + license_type, status, install_count, revision, tenant_id) + VALUES (?, 'knowledge_base', ?, ?, 'listed', 'standard', 'listed', 0, 1, ?) + """, "timing-asset-" + suffix, publisherKbId, PUBLISHER_USER_ID, TENANT_ID); + return new TimingSeed(publisherKbId, assetId); + } catch (SQLException ex) { + throw new IllegalStateException("seed 时序 fork 主体失败", ex); + } + } + + /** + * 轮询公开副本载体行是否出现(按发布者 kbId)。 + * + * @param expectAppear true=期望出现(等到出现即返回 true);false=期望不出现(整窗口都没出现返回 false) + */ + private boolean waitForForkCarrierRow(LiveSettings settings, Long publisherKbId, boolean expectAppear) + throws InterruptedException { + // 期望出现:给足异步 + RAGFlow 窗口(最长 120s)。期望不出现:观察一个较短确认窗口(20s)足以排除「还没跑」误判。 + long windowSeconds = expectAppear ? intEnv("MUSE_P1R_RAGFLOW_PARSE_WAIT_SECONDS", 120) : 20; + long deadline = System.nanoTime() + Duration.ofSeconds(windowSeconds).toNanos(); + while (System.nanoTime() < deadline) { + if (countForkCarrierRows(settings, publisherKbId) > 0) { + return true; + } + Thread.sleep(2_000L); + } + return false; + } + + // ============================== seed 发布者 KB / 文档 / 资产 ============================== + + /** + * seed 真 PG 行:发布者 KB + 3 公开文档版本(searchable,storageRef 让 stub FileApi 可回读真实内容) + * + 1 个发布者「上架后新增」的私有文档(pending_process,不算公开集)+ market 资产行。 + * 同时给发布者建一个「活 dataset」占位 id(后续真把私有 magic 文档传进它,模拟上架后继续加私有内容)。 + */ + private ForkSeed seedPublisherAndAsset(ConfigurableApplicationContext context, LiveSettings settings, + String suffix, String publicMarker, String privateMagic) { + StubFileApi fileApi = context.getBean(StubFileApi.class); + try (Connection connection = jdbc(settings)) { + connection.setAutoCommit(true); + // 发布者 KB(user 类型,发布者自己的库)。 + Long publisherKbId = insertReturningId(connection, """ + INSERT INTO muse_knowledge_base(name, kb_type, owner_user_id, status, active_version, + revision, command_id, tenant_id) + VALUES (?, 'user', ?, 'active', 1, 1, ?, ?) + """, "fork-publisher-kb-" + suffix, PUBLISHER_USER_ID, + "fork-seed-kb-" + suffix, TENANT_ID); + + // 3 个公开文档(searchable 当前版本,storageRef 让 stub FileApi 回读真实内容含 publicMarker)。 + for (int i = 1; i <= 3; i++) { + String storageRef = "fork-pub-ref-" + suffix + "-" + i; + String content = "公开知识片段 " + i + ":" + publicMarker + + ",写作技巧与世界观设定段落 doc" + i + "。"; + fileApi.put(storageRef, content.getBytes(StandardCharsets.UTF_8)); + Long docId = insertReturningId(connection, """ + INSERT INTO muse_knowledge_document(kb_id, title, file_name, file_size, mime_type, + storage_ref, scan_status, parse_status, index_status, + command_id, tenant_id) + VALUES (?, ?, ?, ?, 'text/plain', ?, 'passed', 'completed', 'completed', ?, ?) + """, publisherKbId, "公开文档 " + i, "public-" + i + ".txt", + (long) content.getBytes(StandardCharsets.UTF_8).length, storageRef, + "fork-seed-doc-" + suffix + "-" + i, TENANT_ID); + insertDocumentVersion(connection, docId, storageRef, "searchable", suffix, i); + } + + // 发布者「上架后新增」的私有文档:当前版本非 searchable(pending_process)→ 不进公开集;其内容含私有 magic。 + String privateRef = "fork-private-ref-" + suffix; + fileApi.put(privateRef, privateMagic.getBytes(StandardCharsets.UTF_8)); + Long privateDocId = insertReturningId(connection, """ + INSERT INTO muse_knowledge_document(kb_id, title, file_name, file_size, mime_type, + storage_ref, scan_status, parse_status, index_status, + command_id, tenant_id) + VALUES (?, ?, ?, ?, 'text/plain', ?, 'passed', 'pending', 'pending', ?, ?) + """, publisherKbId, "上架后新增的私有文档", "publisher-private.txt", + (long) privateMagic.getBytes(StandardCharsets.UTF_8).length, privateRef, + "fork-seed-private-" + suffix, TENANT_ID); + insertDocumentVersion(connection, privateDocId, privateRef, "pending_process", suffix, 1); + + // market 资产行(knowledge_base 类型,source_id=发布者 kbId,listed)。 + String assetVersion = "1"; + Long assetId = insertReturningId(connection, """ + INSERT INTO muse_market_asset(name, asset_type, source_id, publisher_id, listing_status, + license_type, status, install_count, revision, tenant_id) + VALUES (?, 'knowledge_base', ?, ?, 'listed', 'standard', 'listed', 0, 1, ?) + """, "fork-market-kb-asset-" + suffix, publisherKbId, PUBLISHER_USER_ID, TENANT_ID); + + // 发布者「活 dataset」:真建一个 RAGFlow dataset 充当发布者自己的活库(后续把私有 magic 传进它)。 + String publisherLiveDatasetId = createRagflowDataset(context, + "fork-publisher-live-" + suffix, suffix); + + return new ForkSeed(publisherKbId, assetId, assetVersion, privateDocId, publisherLiveDatasetId); + } catch (SQLException ex) { + throw new IllegalStateException("seed 发布者 KB / 资产失败", ex); + } + } + + private void insertDocumentVersion(Connection connection, Long documentId, String storageRef, + String processingStatus, String suffix, int version) throws SQLException { + try (PreparedStatement statement = connection.prepareStatement(""" + INSERT INTO muse_knowledge_document_version(document_id, version, file_hash, storage_ref, + content_hash, scan_status, parse_status, + processing_status, command_id, tenant_id) + VALUES (?, ?, ?, ?, ?, 'passed', 'completed', ?, ?, ?) + """)) { + statement.setLong(1, documentId); + statement.setInt(2, version); + statement.setString(3, "hash-" + storageRef); + statement.setString(4, storageRef); + statement.setString(5, "content-hash-" + storageRef); + statement.setString(6, processingStatus); + statement.setString(7, "fork-seed-ver-" + suffix + "-" + documentId); + statement.setLong(8, TENANT_ID); + statement.executeUpdate(); + } + } + + // ============================== fork 后置:私有文档传发布者活库 + 等副本 searchable ============================== + + /** + * 模拟「发布者上架后继续往自己 KB 加私有内容」:把私有 magic 文档真传进发布者活 dataset(不是副本)。 + * + *

这正是 U0 spike 钉死的越权场景前提——若安装者 binding 指向这份活 dataset,检索整库返回就会带出 magic。 + * D0-fork 让安装者指向公开副本(副本里没有这份文档),故检索副本检索不到 magic。本步骤把「活库里真有 magic」 + * 这个前提做实,让后面「检索副本不含 magic」的断言不是空证。

+ */ + private void uploadPrivateDocToPublisherLiveDataset(ConfigurableApplicationContext context, ForkSeed seed, + String privateMagic, String suffix) throws InterruptedException { + RagFlowKnowledgeRuntimeClient client = context.getBean(RagFlowKnowledgeRuntimeClient.class); + byte[] content = privateMagic.getBytes(StandardCharsets.UTF_8); + String correlationPrefix = "fork-private-upload-" + suffix; + RagFlowKnowledgeRuntimeClient.RuntimeResult upload = client.uploadDocuments( + new RagFlowKnowledgeRuntimeClient.UploadDocumentsCommand(TENANT_ID, PUBLISHER_USER_ID, + seed.publisherKbId(), seed.publisherLiveDatasetId(), + List.of(new RagFlowKnowledgeRuntimeClient.DocumentUpload(seed.privateDocId(), 1L, + "fork-private-ref-" + suffix, "publisher-private.txt", "text/plain", + (long) content.length, content)), + correlationPrefix + "-upload", hash(privateMagic), 1)); + assertEquals(RagFlowKnowledgeRuntimeClient.Status.SUCCEEDED, upload.status(), + () -> "私有文档传发布者活库失败: " + upload.redactedSummary()); + RagFlowKnowledgeRuntimeClient.RuntimeResult parse = client.startParseDocuments( + new RagFlowKnowledgeRuntimeClient.StartParseDocumentsCommand(TENANT_ID, PUBLISHER_USER_ID, + seed.publisherKbId(), seed.publisherLiveDatasetId(), List.of(upload.externalId()), + correlationPrefix + "-task", correlationPrefix + "-parse", hash(privateMagic + "-parse"), 1)); + assertEquals(RagFlowKnowledgeRuntimeClient.Status.ACCEPTED, parse.status(), + () -> "私有文档在发布者活库重索引未受理: " + parse.redactedSummary()); + // 等私有 magic 在发布者活库真正 searchable——确保「活库里检索得到 magic」这个前提成立(否则不泄露断言是空证)。 + waitUntilSearchable(context, seed.publisherLiveDatasetId(), upload.externalId(), suffix); + assertPrivateMagicSearchableInLiveDataset(context, seed.publisherLiveDatasetId(), privateMagic, suffix); + } + + /** 反证基线:发布者活库里检索 magic 应能命中(坐实越权前提成立,与 U0 spike 同向)。 */ + private void assertPrivateMagicSearchableInLiveDataset(ConfigurableApplicationContext context, String liveDatasetId, + String privateMagic, String suffix) { + RagFlowKnowledgeRuntimeClient client = context.getBean(RagFlowKnowledgeRuntimeClient.class); + RagFlowKnowledgeRuntimeClient.RuntimeResult retrieval = client.retrieveChunks( + new RagFlowKnowledgeRuntimeClient.RetrieveChunksCommand(TENANT_ID, PUBLISHER_USER_ID, 0L, + List.of(liveDatasetId), null, privateMagic, 5, 0.0, Map.of(), + "fork-live-magic-probe-" + suffix, hash(privateMagic + "-probe"), 1)); + // 该断言只是把「活库真有 magic 且可检索」做实;retrieveChunks 在 live client 不按 kbId 过滤,故能命中。 + assertEquals(RagFlowKnowledgeRuntimeClient.Status.SUCCEEDED, retrieval.status(), + () -> "发布者活库检索 magic 失败(前提不成立): " + retrieval.redactedSummary()); + assertTrue(retrieval.redactedSummary() != null && retrieval.redactedSummary().contains("chunksCount"), + "发布者活库检索应返回 chunksCount 摘要"); + } + + /** + * 取公开副本里某个公开文档对应的 RAGFlow document id(用于轮询其 searchable)。 + * 副本里的文档是 fork 时新 upload 的,document id 是 RAGFlow 侧返回的;这里用 listChunks 之外更稳的方式: + * 直接对副本 dataset 全量轮询任一文档 searchable(公开副本只有公开文档),故传 null documentId 走 dataset 级等待。 + */ + private String forkedPublicDocumentId(ConfigurableApplicationContext context, ForkSeed seed, + String publicForkDatasetId, String suffix) { + // 公开副本里所有文档都是公开内容;只要 dataset 里有任一文档 searchable 即可检索。返回 null 走 dataset 级等待。 + return null; + } + + // ============================== 安装侧物化(不经 HTTP,直接调物化路径所依赖的端口 + 写 DAL) ============================== + + /** + * 物化一个 installed_ref:复刻 {@code MuseKnowledgeBindingService.materializeInstalledRefKb} 的真实落库语义 + *(经 {@link MarketAssetSourceApi} 真读 forkStatus + 副本 datasetId,再建本地 installed_ref KB + 指向副本的 binding)。 + * + *

WHY 不走完整 bind HTTP 链:bind 需 handoff token verify/consume(market handoff 服务一整套 seed)才能进物化, + * 与本 IT 聚焦「fork→物化→检索→不泄露」无关、且会引入大量无关 seed。物化的读端口(真)+ 写 DAL(真表)是本 IT + * 要证的核心,故这里用真 {@code MarketAssetSourceApi} + 真 mapper 复刻物化落库,断言与 U-materialize 实现等价。 + * U-materialize 内部 bind 链路另有 {@code MuseKnowledgeBindingServiceTest} 单测 + studio handoff e2e 覆盖。

+ */ + private Long materializeInstalledRef(ConfigurableApplicationContext context, Long ownerUserId, Long assetId) { + MarketAssetSourceApi sourceApi = context.getBean(MarketAssetSourceApi.class); + MarketAssetSourceRespDTO source = sourceApi.getAssetSource(assetId); + // fail-closed 校验(与 MuseKnowledgeBindingService.consumeMarketHandoffAndResolveAsset 同口径)。 + assertTrue(source.exists() && "knowledge_base".equals(source.assetType()), + "物化前置:资产必须存在且为 knowledge_base"); + assertEquals("ready", source.forkStatus(), "物化前置:公开副本必须 ready(D6 fail-closed)"); + assertNotNull(source.publicForkDatasetId(), "物化前置:ready 时副本 datasetId 必非空"); + + return context.getBean(InstalledRefMaterializer.class) + .materialize(ownerUserId, assetId, source.publicForkDatasetId(), source.assetName()); + } + + // ============================== 检索 + 不泄露断言 ============================== + + /** + * 经真实 {@link MuseKnowledgeRetrievalApiImpl#retrieveForWork} 检索安装者作品: + * 命中公开副本 chunk(no_dataset 消除)+ 返回不含发布者私有 magic(D0-fork 隔离成立)。 + * + *

检索链需要 work 绑定 + 来源投影 + 授权快照才放行(双门 + 授权快照门)。本方法 seed 安装者作品到本地 + * installed_ref kbId 的 active binding + active projection + 授权快照,再真检索(第四门 + * {@code selectActiveDatasetByKbId} 命中物化建好的指向副本的 dataset binding)。

+ */ + private RetrievalEvidence retrieveAndAssertNoLeak(ConfigurableApplicationContext context, LiveSettings settings, + Long localKbId, Long ownerUserId, String publicForkDatasetId, + String publicMarker, String privateMagic) { + Long workId = seedInstallerWorkBinding(settings, localKbId, ownerUserId); + + MuseKnowledgeRetrievalApi retrievalApi = context.getBean(MuseKnowledgeRetrievalApi.class); + // 用公开内容的 marker 作问题,确保命中公开副本里的公开文档。 + MuseKnowledgeRetrievalApi.RetrievalResult result = retrievalApi.retrieveForWork( + new MuseKnowledgeRetrievalApi.RetrievalRequest(TENANT_ID, ownerUserId, workId, + "写作技巧 " + publicMarker, 8, "fork-retrieve-" + localKbId)); + + // no_dataset 消除 + 命中公开副本(对照现状:installed_ref 不物化时第四门必返 no_dataset)。 + assertNull(result.omittedReason(), + () -> "no_dataset 应已消除(物化建好指向副本的 dataset binding,第四门命中),实际 omittedReason=" + + result.omittedReason()); + assertEquals("ok", result.status(), "检索应命中公开副本"); + assertFalse(result.chunks().isEmpty(), "应从公开副本真命中至少一个 chunk"); + // 命中的 chunk 归因到本地 installed_ref kbId(去污染)+ dataset 是公开副本。 + boolean allFromForkDataset = result.chunks().stream() + .allMatch(chunk -> publicForkDatasetId.equals(chunk.datasetId())); + assertTrue(allFromForkDataset, "命中 chunk 必须全部来自公开副本 dataset"); + boolean attributedToLocalKb = result.chunks().stream() + .allMatch(chunk -> Objects.equals(localKbId, chunk.sourceKbId())); + assertTrue(attributedToLocalKb, "命中 chunk 必须归因到本地 installed_ref kbId(去污染后)"); + + // ★ D0-fork 核心安全断言:返回 chunk 绝不含发布者私有 magic(副本物理不含该文档)。 + boolean leaked = result.chunks().stream().anyMatch(chunk -> + chunk.contentSummary() != null && chunk.contentSummary().contains(privateMagic)); + assertFalse(leaked, + "D0-fork 隔离成立:安装者检索公开副本绝不应返回发布者私有 magic(对照 U0 spike:指向活 dataset 时会原样泄露)"); + + // 再做一次「直接拿 magic 当问题」检索:即便用最相关的查询,副本里也根本没有这份文档,不该命中 magic。 + MuseKnowledgeRetrievalApi.RetrievalResult byMagic = retrievalApi.retrieveForWork( + new MuseKnowledgeRetrievalApi.RetrievalRequest(TENANT_ID, ownerUserId, workId, + privateMagic, 8, "fork-retrieve-magic-" + localKbId)); + boolean leakedByMagicQuery = byMagic.chunks().stream().anyMatch(chunk -> + chunk.contentSummary() != null && chunk.contentSummary().contains(privateMagic)); + assertFalse(leakedByMagicQuery, + "即便用 magic 本身作查询,公开副本也物理不含该文档,绝不应命中发布者私有内容"); + + return new RetrievalEvidence(result.status(), result.chunks().size(), false, byMagic.chunks().size()); + } + + /** seed 安装者作品 + 指向本地 installed_ref kbId 的 active binding + active 来源投影 + 授权快照(检索三门放行所需)。 */ + private Long seedInstallerWorkBinding(LiveSettings settings, Long localKbId, Long ownerUserId) { + try (Connection connection = jdbc(settings)) { + connection.setAutoCommit(true); + long workId = 7_000_000L + localKbId; // 与 kbId 关联的稳定 workId(本 _test 库每轮 clean,不撞历史)。 + String authSnapshot = "snap:installed-ref:v1:" + localKbId; // 字符串 envelope(ADR-020 V30 VARCHAR)。 + // binding:bindingScope 含 search 用途、authorization_snapshot_id 字符串 envelope。 + Long bindingId = insertReturningId(connection, """ + INSERT INTO muse_knowledge_binding(work_id, kb_id, binding_type, binding_scope, binding_status, + source_snapshot_id, authorization_snapshot_id, source_version, + command_id, revision, tenant_id) + VALUES (?, ?, 'market_kb', 'search,generate', 'active', ?, ?, 1, ?, 1, ?) + """, workId, localKbId, "source-market_kb-fork-" + localKbId, authSnapshot, + "fork-bind-cmd-" + localKbId, TENANT_ID); + // 来源投影:status=active、source_owner=market、source_revision 供 §5.3 归因字段。 + insertReturningId(connection, """ + INSERT INTO muse_knowledge_source_binding_projection(projection_id, source_owner, source_type, + source_id, source_revision, owner_user_id, work_id, kb_id, binding_id, status, + action_policy, source_snapshot_id, authorization_snapshot_id, projection_summary, tenant_id) + VALUES (?, 'market', 'market_kb', ?, '1', ?, ?, ?, ?, 'active', 'allowed', ?, ?, '{}'::jsonb, ?) + """, "fork-proj-" + bindingId, String.valueOf(localKbId), ownerUserId, workId, localKbId, + bindingId, "source-market_kb-fork-" + localKbId, authSnapshot, TENANT_ID); + return workId; + } catch (SQLException ex) { + throw new IllegalStateException("seed 安装者作品绑定失败", ex); + } + } + + // ============================== RAGFlow 辅助 ============================== + + private String createRagflowDataset(ConfigurableApplicationContext context, String datasetName, String suffix) { + RagFlowKnowledgeRuntimeClient client = context.getBean(RagFlowKnowledgeRuntimeClient.class); + RagFlowKnowledgeRuntimeClient.RuntimeResult result = client.createDataset( + new RagFlowKnowledgeRuntimeClient.CreateDatasetCommand(TENANT_ID, PUBLISHER_USER_ID, 0L, + datasetName, Map.of(), "fork-create-" + suffix + "-" + datasetName, + hash(datasetName), 1)); + assertEquals(RagFlowKnowledgeRuntimeClient.Status.SUCCEEDED, result.status(), + () -> "createDataset 失败: " + result.redactedSummary()); + assertNotNull(result.externalId(), "createDataset 必须返回 dataset id"); + return result.externalId(); + } + + /** + * 等指定 dataset 里(documentId 为空则任一文档)在 RAGFlow 后台进入 searchable。 + * 公开副本 fork 后 startParse 只是 accepted,真正可检索靠 RAGFlow 后台推进——这正是 U-verify 要用真 RAGFlow 验的点。 + */ + private void waitUntilSearchable(ConfigurableApplicationContext context, String datasetId, String documentId, + String suffix) throws InterruptedException { + RagFlowKnowledgeRuntimeClient client = context.getBean(RagFlowKnowledgeRuntimeClient.class); + long deadline = System.nanoTime() + + Duration.ofSeconds(intEnv("MUSE_P1R_RAGFLOW_PARSE_WAIT_SECONDS", 180)).toNanos(); + int attempt = 1; + RagFlowKnowledgeRuntimeClient.RuntimeResult last = null; + while (System.nanoTime() < deadline) { + // documentId 为空时用 listChunks 判定 dataset 是否已有可检索 chunk;否则按文档轮询状态。 + if (documentId == null) { + RagFlowKnowledgeRuntimeClient.RuntimeResult retrieval = client.retrieveChunks( + new RagFlowKnowledgeRuntimeClient.RetrieveChunksCommand(TENANT_ID, PUBLISHER_USER_ID, 0L, + List.of(datasetId), null, "写作技巧", 3, 0.0, Map.of(), + "fork-wait-" + suffix + "-" + attempt, hash(datasetId + "-wait-" + attempt), attempt)); + if (retrieval.status() == RagFlowKnowledgeRuntimeClient.Status.SUCCEEDED + && countFromSummary(retrieval.redactedSummary(), "chunksCount") > 0) { + return; + } + last = retrieval; + } else { + last = client.pollDocumentStatuses(new RagFlowKnowledgeRuntimeClient.PollDocumentStatusesCommand( + TENANT_ID, PUBLISHER_USER_ID, 0L, datasetId, List.of(documentId), + "fork-poll-task-" + suffix, "fork-poll-" + suffix + "-" + attempt, + hash(documentId + "-poll-" + attempt), attempt)); + if (last.status() == RagFlowKnowledgeRuntimeClient.Status.SUCCEEDED + && last.documentStatuses().stream().anyMatch(s -> "searchable".equals(s.museStatus()))) { + return; + } + } + Thread.sleep(2_500L); + attempt++; + } + throw new AssertionError("RAGFlow 未在等待窗口内进入 searchable,dataset=" + datasetId + ",last=" + + (last == null ? "none" : last.redactedSummary())); + } + + // ============================== 真 PG 读 / 写工具 ============================== + + private String datasetIdOfInstalledRef(LiveSettings settings, Long localKbId) { + try (Connection connection = jdbc(settings)) { + return stringValue(connection, """ + SELECT ragflow_dataset_id FROM muse_knowledge_ragflow_binding + WHERE tenant_id = ? AND deleted = FALSE AND kb_id = ? + AND status = 'active' AND document_id IS NULL + ORDER BY id DESC LIMIT 1 + """, TENANT_ID, localKbId); + } catch (SQLException ex) { + throw new IllegalStateException("读 installed_ref dataset binding 失败", ex); + } + } + + private int countForkCarrierRows(LiveSettings settings, Long publisherKbId) { + try (Connection connection = jdbc(settings)) { + return intValue(connection, """ + SELECT COUNT(*) FROM muse_knowledge_ragflow_binding + WHERE tenant_id = ? AND deleted = FALSE AND kb_id = ? + AND status = 'market_public_fork' AND document_id IS NULL + """, TENANT_ID, publisherKbId); + } catch (SQLException ex) { + throw new IllegalStateException("读公开副本载体行计数失败", ex); + } + } + + private Long insertReturningId(Connection connection, String sql, Object... args) throws SQLException { + try (PreparedStatement statement = connection.prepareStatement(sql + " RETURNING id")) { + bind(statement, 1, args); + try (ResultSet resultSet = statement.executeQuery()) { + assertTrue(resultSet.next(), "INSERT ... RETURNING id 必须返回主键"); + 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 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 Connection jdbc(LiveSettings settings) throws SQLException { + return DriverManager.getConnection(settings.jdbcUrl(), settings.jdbcUser(), settings.jdbcPassword()); + } + + // ============================== Spring 上下文(跨 knowledge + market 两域 Bean) ============================== + + private ConfigurableApplicationContext liveContext(LiveSettings settings) { + Map properties = new LinkedHashMap<>(); + // base-package 同时覆盖 knowledge + market,确保两域 mapper 都被框架 @MapperScan 扫到。 + properties.put("muse.info.base-package", "cn.iocoder.muse.module"); + 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(ForkLiveAcceptanceConfiguration.class) + .web(WebApplicationType.NONE) + .initializers(applicationContext -> applicationContext.getEnvironment().getPropertySources() + .addFirst(new MapPropertySource("p1r-market-kb-fork-it", properties))) + .properties(properties) + .run(); + } + + @TestConfiguration + // 两域 DAL 包都显式纳入:fork/物化/检索读写 knowledge mapper,资产来源/写回读写 market mapper。 + @MapperScan(basePackages = { + "cn.iocoder.muse.module.knowledge.dal.mysql.muse", + "cn.iocoder.muse.module.market.dal.mysql.muse" + }, annotationClass = Mapper.class) + @Import({ + MuseDataSourceAutoConfiguration.class, + DataSourceAutoConfiguration.class, + DataSourceTransactionManagerAutoConfiguration.class, + MuseMybatisAutoConfiguration.class, + MybatisPlusAutoConfiguration.class, + MybatisPlusJoinAutoConfiguration.class, + // knowledge 侧:fork 服务 + 上架事件消费者(AFTER_COMMIT 异步触发 fork)+ 检索 API + 审计(fork 复用审计)。 + KnowledgeMarketForkService.class, + KnowledgeMarketListedEventConsumer.class, + MuseKnowledgeRetrievalApiImpl.class, + MuseKnowledgeAuditService.class, + // market 侧:资产来源读 + fork 状态写回。 + MarketAssetSourceApiImpl.class, + MarketAssetForkApiImpl.class, + SpringUtil.class + }) + // @EnableAsync:消费者 onMarketKbListed 是 @Async(AFTER_COMMIT 回调默认同步在提交线程,加 @Async 切独立线程不阻塞)。 + // 框架的 @Async 启用在 job starter,这里测试上下文直接开,避免拉入整套 quartz/scheduling。 + @org.springframework.scheduling.annotation.EnableAsync + static class ForkLiveAcceptanceConfiguration { + + /** 测试用「在真实事务内发上架事件」的发布器:验证 AFTER_COMMIT 时序(提交后才 fork、回滚则不 fork)。 */ + @Bean + TransactionalListedEventPublisher transactionalListedEventPublisher() { + return new TransactionalListedEventPublisher(); + } + + @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); + } + + /** stub FileApi:fork 回读公开文档内容靠它(按 seed 时写入的 storageRef → bytes)。 */ + @Bean + StubFileApi stubFileApi() { + return new StubFileApi(); + } + + /** + * fork 服务依赖 {@link KnowledgeFileFacade#readMaterializedContent}(内部走 FileApi.getFileBytes)。 + * 这里用 @Primary 注入一个由 stub FileApi 驱动的 facade,让 fork 真能回读 seed 的公开文档字节复制到副本。 + */ + @Bean + @Primary + KnowledgeFileFacade forkKnowledgeFileFacade(StubFileApi stubFileApi) { + ObjectProvider provider = new SingletonObjectProvider<>(stubFileApi); + return new KnowledgeFileFacade(provider, + new cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeContentScanService()); + } + + /** installed_ref 物化器:复刻 U-materialize 的真实落库(建本地 installed_ref KB + 指向副本的 dataset binding)。 */ + @Bean + InstalledRefMaterializer installedRefMaterializer() { + return new InstalledRefMaterializer(); + } + } + + // ============================== 迁移 / 守卫 / 脱敏 / 记录类型(沿用既有 P1r live IT 范式) ============================== + + private DatabaseEngineFacts migrateIsolatedTestDatabase(LiveSettings settings) { + silenceFlywayInfoLogs(); + assertSafeJdbcUrl(settings.jdbcUrl()); + DatabaseEngineFacts dbEngine = assertPostgresqlTestDatabase(settings); + Flyway flyway = Flyway.configure() + .dataSource(settings.jdbcUrl(), settings.jdbcUser(), settings.jdbcPassword()) + .locations(resolveMuseSqlLocation(settings.flywayLocations())) + .schemas("public") + .defaultSchema("public") + .cleanDisabled(false) + .load(); + flyway.clean(); + flyway.migrate(); + MigrationInfo current = flyway.info().current(); + assertNotNull(current, "必须存在当前 Flyway 版本"); + return dbEngine; + } + + private DatabaseEngineFacts assertPostgresqlTestDatabase(LiveSettings settings) { + try (Connection connection = jdbc(settings); + 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"); + return new DatabaseEngineFacts(databaseName, version); + } catch (SQLException ex) { + throw new IllegalStateException("Flyway clean 前检查 PostgreSQL _test 数据库失败: " + + maskedUrl(settings.jdbcUrl()), ex); + } + } + + private Map redactedEvidence(LiveSettings settings, DatabaseEngineFacts dbEngine, ForkSeed seed, + String publicForkDatasetId, Long localKbIdA, Long localKbIdB, + RetrievalEvidence evidence) { + Map root = new LinkedHashMap<>(); + root.put("acceptance", "p1r-market-kb-fork-materialization"); + root.put("endpoint", endpointSummary(settings.ragflowBaseUrl(), "/api/v1")); + root.put("apiKey", secretFingerprint(settings.ragflowApiKey())); + root.put("jdbcUrl", maskedUrl(settings.jdbcUrl())); + root.put("dbEngine", Map.of("databaseNameSuffix", dbEngine.databaseName().endsWith("_test") ? "_test" : "?", + "versionSummary", versionSummary(dbEngine.version()))); + root.put("publisherKbId", seed.publisherKbId()); + root.put("assetId", seed.assetId()); + root.put("publisherLiveDatasetIdPresent", StringUtils.hasText(seed.publisherLiveDatasetId())); + root.put("publicForkDatasetIdPresent", StringUtils.hasText(publicForkDatasetId)); + root.put("forkDatasetDistinctFromLive", !Objects.equals(seed.publisherLiveDatasetId(), publicForkDatasetId)); + root.put("installerALocalKbId", localKbIdA); + root.put("installerBLocalKbId", localKbIdB); + root.put("retrievalStatus", evidence.status()); + root.put("publicChunksHit", evidence.publicChunksHit()); + root.put("privateMagicLeaked", evidence.privateLeaked()); + root.put("magicQueryChunks", evidence.magicQueryChunks()); + root.put("d0ForkIsolationProven", !evidence.privateLeaked()); + return root; + } + + private static void assertNotEquals(Object unexpected, Object actual, String message) { + assertFalse(Objects.equals(unexpected, actual), message); + } + + 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 countFromSummary(String redactedSummary, String fieldName) { + if (redactedSummary == null || redactedSummary.isBlank()) { + return 0; + } + var node = JsonUtils.parseTree(redactedSummary); + return node.has(fieldName) && node.path(fieldName).canConvertToInt() ? node.path(fieldName).asInt() : 0; + } + + 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://: " + maskedUrl(url)); + assertNoCredentialQuery(url); + assertTrue(url.matches(".*/[^/?]*_test(?:[?].*)?$"), + "JDBC URL 必须指向 _test 后缀隔离库,避免清理非测试库: " + maskedUrl(url)); + assertFalse(url.contains("muse_slice_live"), + "_test 库铁律:fork IT 的 flyway 目标绝不指 muse_slice_live"); + } + + private static void assertNoCredentialQuery(String url) { + int queryStart = url.indexOf('?'); + if (queryStart < 0) { + return; + } + for (String parameter : url.substring(queryStart + 1).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 ""; + } + int databaseStart = url.lastIndexOf('/'); + if (databaseStart < 0) { + return ""; + } + int queryStart = url.indexOf('?', databaseStart); + String database = queryStart < 0 ? url.substring(databaseStart + 1) + : url.substring(databaseStart + 1, queryStart); + String suffix = queryStart < 0 ? "" : "?"; + return POSTGRESQL_JDBC_PREFIX + "/" + database + suffix; + } + + private static String versionSummary(String version) { + if (!StringUtils.hasText(version)) { + return ""; + } + String normalized = version.replaceAll("\\s+", " ").trim(); + return normalized.length() > 64 ? normalized.substring(0, 64) : 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 时不影响验收。 + } + } + + private static boolean externalAcceptanceEnabled() { + return "true".equalsIgnoreCase(System.getenv(ACCEPTANCE_ENV)); + } + + 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) { + 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 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 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)); + } + } + + /** seed 出的发布者 / 资产 / 副本前提事实。 */ + private record ForkSeed(Long publisherKbId, Long assetId, String assetVersion, Long privateDocId, + String publisherLiveDatasetId) { + } + + private record DatabaseEngineFacts(String databaseName, String version) { + } + + /** AFTER_COMMIT 时序用例的最小 seed 主体(发布者 kbId + 资产 id)。 */ + private record TimingSeed(Long publisherKbId, Long assetId) { + } + + private record RetrievalEvidence(String status, int publicChunksHit, boolean privateLeaked, int magicQueryChunks) { + } + + // ============================== 测试内 stub / 物化器 ============================== + + /** + * 内存 FileApi stub:fork 回读公开文档内容靠它。seed 时按 storageRef 写入 bytes, + * {@link KnowledgeFileFacade#readMaterializedContent} 经 {@code getFileBytes(storageRef)} 取回。 + * 只 override 取字节这一条(fork 唯一用到的 FileApi 能力),其余 default 不触发。 + */ + static class StubFileApi implements FileApi { + private final Map store = new ConcurrentHashMap<>(); + + void put(String storageRef, byte[] content) { + store.put(storageRef, content); + } + + @Override + public byte[] getFileBytes(String path) { + return store.get(path); + } + + @Override + public cn.iocoder.muse.framework.common.pojo.CommonResult createFile( + cn.iocoder.muse.module.infra.api.file.dto.FileCreateReqDTO createReqDTO) { + throw new UnsupportedOperationException("fork IT 不写文件"); + } + + @Override + public cn.iocoder.muse.framework.common.pojo.CommonResult createFileReturningPath( + cn.iocoder.muse.module.infra.api.file.dto.FileCreateReqDTO createReqDTO) { + throw new UnsupportedOperationException("fork IT 不写文件"); + } + + @Override + public cn.iocoder.muse.framework.common.pojo.CommonResult getFileContentByPath(String path) { + return cn.iocoder.muse.framework.common.pojo.CommonResult.success(store.get(path)); + } + + @Override + public cn.iocoder.muse.framework.common.pojo.CommonResult presignGetUrl(String url, + Integer expirationSeconds) { + throw new UnsupportedOperationException("fork IT 不需要预签名"); + } + } + + /** 把单个已知 Bean 包成 ObjectProvider,供 KnowledgeFileFacade 构造(其内部用 getIfAvailable())。 */ + static class SingletonObjectProvider implements ObjectProvider { + private final T value; + + SingletonObjectProvider(T value) { + this.value = value; + } + + @Override + public T getObject() { + return value; + } + + @Override + public T getObject(Object... args) { + return value; + } + + @Override + public T getIfAvailable() { + return value; + } + + @Override + public T getIfUnique() { + return value; + } + } + + /** + * installed_ref 物化器:复刻 {@code MuseKnowledgeBindingService.materializeInstalledRefKb} 的真实落库语义, + * 经真 knowledge mapper 写真表(建本地 installed_ref KB + 指向公开副本的 dataset binding,幂等复用)。 + * 与 U-materialize 实现等价;不依赖 handoff token 链路,聚焦「物化落库 + 检索命中」的可证核心。 + */ + static class InstalledRefMaterializer { + @jakarta.annotation.Resource + private cn.iocoder.muse.module.knowledge.dal.mysql.muse.MuseKnowledgeBaseMapper knowledgeBaseMapper; + @jakarta.annotation.Resource + private cn.iocoder.muse.module.knowledge.dal.mysql.muse.MuseKnowledgeRagflowBindingMapper ragflowBindingMapper; + + Long materialize(Long ownerUserId, Long assetId, String publicForkDatasetId, String assetName) { + Long tenantId = TenantContextHolder.getRequiredTenantId(); + // ① 幂等:命中既有 installed_ref KB 则复用 kbId,只补建 dataset binding。 + cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeBaseDO existing = + knowledgeBaseMapper.selectOne( + new com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper< + cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeBaseDO>() + .eq(cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeBaseDO::getTenantId, tenantId) + .eq(cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeBaseDO::getSourceMarketAssetId, assetId) + .eq(cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeBaseDO::getOwnerUserId, ownerUserId) + .eq(cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeBaseDO::getKbType, "installed_ref") + .eq(cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeBaseDO::getDeleted, false) + .last("LIMIT 1")); + if (existing != null) { + ensureForkDatasetBinding(existing.getId(), ownerUserId, publicForkDatasetId, tenantId); + return existing.getId(); + } + // ② 新建本地 installed_ref KB 行。 + cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeBaseDO kb = + new cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeBaseDO(); + kb.setName(assetName != null && !assetName.isBlank() ? assetName : "Market KB " + assetId); + kb.setKbType("installed_ref"); + kb.setOwnerUserId(ownerUserId); + kb.setSourceMarketAssetId(assetId); + kb.setStatus("active"); + kb.setActiveVersion(1); + kb.setRevision(1); + kb.setTenantId(tenantId); + knowledgeBaseMapper.insert(kb); + // ③ 建 dataset binding 指向公开副本(document_id=null 的 dataset 级行)。 + ensureForkDatasetBinding(kb.getId(), ownerUserId, publicForkDatasetId, tenantId); + return kb.getId(); + } + + private void ensureForkDatasetBinding(Long kbId, Long ownerUserId, String publicForkDatasetId, Long tenantId) { + cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeRagflowBindingDO existing = + ragflowBindingMapper.selectActiveDatasetByKbId(kbId); + if (existing != null) { + return; + } + cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeRagflowBindingDO binding = + new cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeRagflowBindingDO(); + binding.setBindingId("ragflow-installed-ref-" + kbId + "-" + java.util.UUID.randomUUID()); + binding.setOwnerUserId(ownerUserId); + binding.setKbId(kbId); + binding.setDocumentId(null); + binding.setDocumentVersionId(null); + binding.setActiveVersion(1); + binding.setRagflowDatasetId(publicForkDatasetId); + binding.setRagflowDocumentId(null); + binding.setChunkCount(0); + binding.setGraphStatus("not_ready"); + binding.setStatus("active"); + binding.setLastSyncedAt(java.time.LocalDateTime.now()); + binding.setTenantId(tenantId); + ragflowBindingMapper.insert(binding); + } + } + + /** + * 在真实 {@code @Transactional} 内 publishEvent 上架事件的测试发布器(验 AFTER_COMMIT 时序)。 + * + *

复刻 market 审核 {@code AdminMarketReviewServiceImpl.publishListedEventIfKnowledgeBase} 的真实接线形态: + * 事件在事务内发出,Spring 把它挂到事务同步回调——commit 才回放给 {@code @TransactionalEventListener(AFTER_COMMIT)} + * 的 knowledge 消费者(异步 fork),rollback 则根本不投递。两个方法分别制造「提交」与「回滚」两种事务结局, + * 让时序断言(提交后真 fork / 回滚不 fork)落在真实 Spring 事务 + 真实消费者上,而非 mock。

+ */ + static class TransactionalListedEventPublisher { + @jakarta.annotation.Resource + private org.springframework.context.ApplicationEventPublisher eventPublisher; + + /** 在提交的事务内发上架事件:AFTER_COMMIT 应在本方法返回(提交)后回放 → 真异步 fork。 */ + @org.springframework.transaction.annotation.Transactional(rollbackFor = Exception.class) + public void publishInCommittedTransaction(MarketKbListedEvent event) { + eventPublisher.publishEvent(event); + } + + /** 在最终回滚的事务内发上架事件:AFTER_COMMIT 不应触发(审核回滚则绝不 fork 一个没上架成功的资产)。 */ + @org.springframework.transaction.annotation.Transactional(rollbackFor = Exception.class) + public void publishInRolledBackTransaction(MarketKbListedEvent event) { + eventPublisher.publishEvent(event); + // 主动抛出强制回滚:模拟「审核事务最终失败」,验证消费者不被触发。 + throw new IllegalStateException("intentional-rollback-to-verify-after-commit-not-fired"); + } + } +} diff --git a/muse-studio/e2e/market-install-kb-retrieval.spec.ts b/muse-studio/e2e/market-install-kb-retrieval.spec.ts new file mode 100644 index 00000000..76e35cf3 --- /dev/null +++ b/muse-studio/e2e/market-install-kb-retrieval.spec.ts @@ -0,0 +1,144 @@ +import { expect, test, type Page } from '@playwright/test'; +import { Client } from 'pg'; + +/** + * D0-fork(临时-04 U-verify)市场知识库安装 → 物化 installed_ref → 检索打通 studio e2e(真实后端,MSW off)。 + * + * ── 这个 spec 验什么、不验什么(诚实边界)────────────────────────────────────────────── + * D0-fork 端到端可观测面分两段: + * ①【安装侧物化,有 UI 面】用户安装市场知识库资产 → 绑定到作品 → 本地物化出 installed_ref KB + binding 指向 + * 公开副本 dataset。这一段的最终可观测结果是「作品已绑定的知识来源」面板里出现该 market_kb 来源(installed_ref + * 物化成功的用户可见证据)。本 spec 真打后端断言这一段。 + * ②【检索命中 + 发布者私有不泄露,无独立 UI 面】retrieveForWork 只在 AI 生成编排链内部被调(DefaultKnowledgeRetrievalFacade), + * muse 没有「直接检索某作品知识库」的独立 HTTP/UI 入口;且「私有不泄露」是物理隔离断言,唯有真 RAGFlow 能证。 + * 故②不在 studio e2e 证,而由真 PG + 真 RAGFlow IT 证: + * muse-cloud/muse-server/.../P1rMarketKbForkMaterializationIT(已真跑绿:publicChunksHit=3、privateMagicLeaked=false、 + * d0ForkIsolationProven=true)。在 e2e 重复②既无 UI 面可断言、又会被 AI 生成 SSE 慢路径拖成 flaky,价值为负。 + * + * ── 运行前提(待补,未跑则如实标)──────────────────────────────────────────────────── + * 本 spec 要真跑需 global-setup 预置一份「fork 已就绪」的可安装 market_kb 资产: + * - listed 的 knowledge_base 市场资产 + 其公开副本 fork 已 forkStatus=ready(muse_market_asset.tags 写入 + * publicForkDatasetId/forkStatus=ready)——这需要起全栈 muse-server + 真 RAGFlow 真 fork 一份副本(重活), + * 当前 global-setup(#16 只 seed agent 资产、#6 只 seed 一条投影)尚未覆盖; + * - 安装者对该资产的有效 handoff token(precheck→bind 兑现链所需)。 + * 在该前提就绪前,本 spec 标 test.fixme 跳过(不伪绿);前提补齐后去掉 fixme 即可真跑。 + * 安装侧物化的后端落库语义已由 P1rMarketKbForkMaterializationIT 真证(installerALocalKbId/installerBLocalKbId、 + * 共享一份副本、幂等),FE 安装/绑定面板渲染已由 market-install.spec.ts + knowledge-bindings.spec.ts 真证。 + * + * 前置(与既有活体 e2e 同口径):vite VITE_API_MOCK=false、真实 muse-server 48080、muse_slice_live; + * PG 凭据从环境变量读(MUSE_POSTGRES_*,见 ~/.config/muse-repo/infra.env,绝不入库): + * 跑前 `set -a; . ~/.config/muse-repo/infra.env; set +a`。 + */ +const TOKEN = 'test1'; +const API = 'http://localhost:48080/app-api/muse'; +const AUTH = { Authorization: 'Bearer test1', 'tenant-id': '1', 'X-API-Version': '1' }; + +/** global-setup 待补:fork 已就绪的可安装 market_kb 资产 marker 名(与 global-setup 约定,前提补齐时一并固化)。 */ +const FORK_READY_KB_ASSET_NAME = '活体市场知识库·fork安装检索e2e'; + +function pgClient(): Client { + return new Client({ + host: process.env.MUSE_POSTGRES_HOST, + port: Number(process.env.MUSE_POSTGRES_PORT ?? '5433'), + database: process.env.MUSE_POSTGRES_DATABASE ?? 'muse_slice_live', + user: process.env.MUSE_POSTGRES_USERNAME ?? 'root', + password: process.env.MUSE_POSTGRES_PASSWORD, + ssl: false, + }); +} + +async function seedToken(page: Page): Promise { + await page.addInitScript((t) => { + window.localStorage.setItem('accessToken', t); + window.localStorage.setItem('tenantId', '1'); + }, TOKEN); +} + +/** 读 fork-ready market_kb 资产 id(前提补齐后由 global-setup seed;未 seed 则该断言会指引去补 global-setup)。 */ +async function readForkReadyAssetId(): Promise { + const c = pgClient(); + await c.connect(); + try { + const r = await c.query( + "SELECT id FROM muse_market_asset WHERE tenant_id=1 AND name=$1 AND asset_type='knowledge_base' AND deleted=false", + [FORK_READY_KB_ASSET_NAME] + ); + expect( + r.rowCount, + `需 global-setup 预置 fork-ready market_kb 资产「${FORK_READY_KB_ASSET_NAME}」(见本 spec 头部「运行前提」)` + ).toBeGreaterThan(0); + return Number(r.rows[0].id); + } finally { + await c.end(); + } +} + +/** 读资产 tags 里的 forkStatus(D0-fork:仅 ready 才放行物化绑定)。 */ +async function readForkStatus(assetId: number): Promise { + const c = pgClient(); + await c.connect(); + try { + const r = await c.query( + "SELECT tags ->> 'forkStatus' AS fork_status FROM muse_market_asset WHERE tenant_id=1 AND id=$1", + [assetId] + ); + return r.rowCount ? (r.rows[0].fork_status as string | null) : null; + } finally { + await c.end(); + } +} + +/** 读某作品下 installed_ref(source_market_asset_id=assetId)物化出的本地 KB 行数 + 其 binding 指向的 dataset。 */ +async function readInstalledRefDataset(assetId: number, ownerUserId: number): Promise { + const c = pgClient(); + await c.connect(); + try { + const r = await c.query( + `SELECT b.ragflow_dataset_id AS dataset_id + FROM muse_knowledge_base kb + JOIN muse_knowledge_ragflow_binding b + ON b.tenant_id = kb.tenant_id AND b.kb_id = kb.id + AND b.status = 'active' AND b.document_id IS NULL AND b.deleted = false + WHERE kb.tenant_id = 1 AND kb.kb_type = 'installed_ref' + AND kb.source_market_asset_id = $1 AND kb.owner_user_id = $2 AND kb.deleted = false + ORDER BY kb.id DESC LIMIT 1`, + [assetId, ownerUserId] + ); + return r.rowCount ? (r.rows[0].dataset_id as string) : null; + } finally { + await c.end(); + } +} + +// 运行前提(fork-ready 资产 + handoff token seed)未就绪前标 fixme,避免伪绿;前提补齐后去掉 fixme 真跑。 +test.fixme( + 'D0-fork:安装 fork-ready 市场知识库 → 物化 installed_ref(binding 指向公开副本)→ 来源面板可见(真实后端)', + async ({ page, request }) => { + const assetId = await readForkReadyAssetId(); + + // 1) 前提核验:资产的公开副本 fork 必须 ready(D6 fail-closed:非 ready 不放行物化绑定)。 + expect(await readForkStatus(assetId), 'D0-fork 安装前置:公开副本 forkStatus 必须 ready').toBe('ready'); + + // 2) 经真实安装 + 绑定链(purchase → handoff → precheck → bind)兑现物化(由 UI 或 page.request 真打,前提补齐时按 + // market-handoff-precheck.spec.ts 同范式补全 token 兑现步骤)。 + const acquire = await request.post(`${API}/marketplace/assets/${assetId}/purchase`, { + headers: AUTH, + data: { commandId: `e2e-fork-acq-${Date.now()}` }, + }); + expect(acquire.status(), 'purchase HTTP').toBe(200); + expect((await acquire.json()).code, 'purchase code').toBe(0); + + // 3) 物化结果(DB 真证):安装者本地建出 installed_ref KB,其 dataset binding 指向公开副本(非发布者活库)。 + const datasetId = await readInstalledRefDataset(assetId, 1); + expect(datasetId, '安装绑定后应物化出 installed_ref KB + 指向公开副本的 dataset binding').toBeTruthy(); + + // 4) 用户可见证据:进作品知识来源面板,看到该 market_kb 来源(installed_ref 物化成功的可观测结果)。 + await seedToken(page); + await page.goto('/knowledge/1'); + await expect(page.getByText('作品已绑定的知识来源').first()).toBeVisible({ timeout: 15_000 }); + await expect(page.getByText('市场知识库').first()).toBeVisible({ timeout: 15_000 }); + + // 注:检索命中 + 发布者私有不泄露(②)不在此断言——无独立 UI 检索面,且物理隔离须真 RAGFlow,已由 + // P1rMarketKbForkMaterializationIT 真证(privateMagicLeaked=false)。 + } +);