feat(ai-knowledge): P-B 授权过滤 fail-closed,检索结果补齐 §5.3 合同字段 + omittedSources

专题-03 AI 检索串联第二期(P-B):在 P-A 最小检索链路上叠加按用途授权 + 来源状态门控(fail-closed),检索结果补齐 §5.3 可追溯合同字段,被过滤来源透出 omittedSources 供审计。

- §4.3 用途门:binding.bindingScope 须含 "search" 上下文检索用途(阅读/导出≠检索),否则拒
- §4.3 来源状态门:projection.status ∈ {revoked,recalled,delisted,blocked,owner_missing,unauthorized} 一律拒;active/stale 纳入(stale 非阻断);active binding 无投影则 fail-closed
- §5.3 授权快照门:缺 authorizationSnapshot 的来源不进 Prompt
- §5.3 chunk 补字段:sourceOwner/sourceObjectVersion/authorizationSnapshot/sourceStatus/allowedPurpose
- omittedSources:被授权门过滤的来源(kbId/reason/sourceStatus)透出;executor metadata 透出 omittedSourceCount
- 全部消费已有 binding + source_binding_projection 数据,无新建授权模型/表

验证:knowledge 11 单测 + ai 6 单测 + ArchUnit(BcBoundary/AiGrantRuntime)全绿。端到端真链路(真 RAGFlow 摄入 + 真 New-API 生成 + 授权门生效)验收待补(攒 P-B 后一次做)。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
lili 2026-06-23 03:59:36 -07:00
parent 987a20b5b7
commit 5659eb66c6
7 changed files with 322 additions and 86 deletions

View File

@ -174,6 +174,10 @@ public class MuseAiRuntimeJobExecutor {
// 检索作品绑定知识库的授权片段,填入上下文 Layer 3(专题-03 §4.2);检索失败/无命中由 facade fail-closed 降级,不阻断生成。
KnowledgeRetrievalFacade.ContextAssembly knowledgeContext = assembleKnowledgeContext(task, job, userInstruction);
metadata.put("knowledgeChunkCount", knowledgeContext.chunkCount());
if (knowledgeContext.omittedSourceCount() > 0) {
// 被授权门(§4.3 用途/来源状态)或 token 预算过滤掉的来源数,透出供审计与质量门控可见(专题-03 §4 omittedSources)
metadata.put("knowledgeOmittedSourceCount", knowledgeContext.omittedSourceCount());
}
if (knowledgeContext.omittedReason() != null) {
metadata.put("knowledgeOmittedReason", knowledgeContext.omittedReason());
}

View File

@ -7,6 +7,7 @@ import cn.iocoder.muse.module.knowledge.api.MuseKnowledgeRetrievalApi.RetrievedC
import jakarta.annotation.Resource;
import org.springframework.context.annotation.Primary;
import org.springframework.stereotype.Component;
import org.springframework.util.StringUtils;
import java.util.List;
import java.util.Locale;
@ -34,11 +35,15 @@ public class DefaultKnowledgeRetrievalFacade implements KnowledgeRetrievalFacade
RetrievalResult result = knowledgeRetrievalApi.retrieveForWork(new RetrievalRequest(
request.tenantId(), request.ownerUserId(), request.workId(),
request.question(), request.topK(), request.correlationId()));
if (result == null || !result.hasChunks()) {
return ContextAssembly.empty(result == null ? "retrieval_null" : result.omittedReason());
if (result == null) {
return ContextAssembly.empty("retrieval_null");
}
int omittedSourceCount = result.omittedSources() == null ? 0 : result.omittedSources().size();
if (!result.hasChunks()) {
return ContextAssembly.empty(result.omittedReason(), omittedSourceCount);
}
String section = formatChunks(result.chunks());
return new ContextAssembly(section, result.chunks().size(), null);
return new ContextAssembly(section, result.chunks().size(), omittedSourceCount, null);
} catch (RuntimeException ex) {
// fail-closed 降级:检索故障绝不阻断主生成链(检索是增强、非硬依赖)
return ContextAssembly.empty("retrieval_error");
@ -51,6 +56,12 @@ public class DefaultKnowledgeRetrievalFacade implements KnowledgeRetrievalFacade
int idx = 1;
for (RetrievedChunk chunk : chunks) {
sb.append("\n[").append(idx++).append("] (kbId=").append(chunk.sourceKbId());
if (StringUtils.hasText(chunk.sourceOwner())) {
sb.append(", owner=").append(chunk.sourceOwner());
}
if (StringUtils.hasText(chunk.sourceStatus())) {
sb.append(", status=").append(chunk.sourceStatus());
}
if (chunk.similarity() != null) {
sb.append(", similarity=").append(String.format(Locale.ROOT, "%.3f", chunk.similarity()));
}

View File

@ -28,14 +28,19 @@ public interface KnowledgeRetrievalFacade {
/**
* 上下文组装结果。
*
* @param promptSection 可拼进 provider prompt 的 L3 上下文文本(带来源归因,可空)
* @param chunkCount 命中片段数
* @param omittedReason 无命中/降级原因(供审计与 omittedSources;有命中时为 null)
* @param promptSection 可拼进 provider prompt 的 L3 上下文文本(带来源归因,可空)
* @param chunkCount 命中片段数
* @param omittedSourceCount 被授权门/token 预算过滤掉的来源数(专题-03 §4 omittedSources,供审计透明)
* @param omittedReason 无命中/降级原因(供审计;有命中时为 null)
*/
record ContextAssembly(String promptSection, int chunkCount, String omittedReason) {
record ContextAssembly(String promptSection, int chunkCount, int omittedSourceCount, String omittedReason) {
public static ContextAssembly empty(String omittedReason) {
return new ContextAssembly(null, 0, omittedReason);
return new ContextAssembly(null, 0, 0, omittedReason);
}
public static ContextAssembly empty(String omittedReason, int omittedSourceCount) {
return new ContextAssembly(null, 0, omittedSourceCount, omittedReason);
}
public boolean hasContext() {

View File

@ -4,6 +4,7 @@ import cn.iocoder.muse.framework.test.core.ut.BaseMockitoUnitTest;
import cn.iocoder.muse.module.ai.application.muse.facade.KnowledgeRetrievalFacade.ContextAssembly;
import cn.iocoder.muse.module.ai.application.muse.facade.KnowledgeRetrievalFacade.ContextAssemblyRequest;
import cn.iocoder.muse.module.knowledge.api.MuseKnowledgeRetrievalApi;
import cn.iocoder.muse.module.knowledge.api.MuseKnowledgeRetrievalApi.OmittedSource;
import cn.iocoder.muse.module.knowledge.api.MuseKnowledgeRetrievalApi.RetrievalRequest;
import cn.iocoder.muse.module.knowledge.api.MuseKnowledgeRetrievalApi.RetrievalResult;
import cn.iocoder.muse.module.knowledge.api.MuseKnowledgeRetrievalApi.RetrievedChunk;
@ -20,7 +21,8 @@ import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.when;
/**
* P-A AI 侧检索 facade 单测:命中→组装 promptSection(含归因);无命中/异常→empty 降级(绝不抛、不阻断主链)。
* P-B AI 侧检索 facade 单测:命中→组装 promptSection(含 §5.3 归因)+ 透出 omittedSourceCount;
* 无命中/异常→empty 降级(绝不抛、不阻断主链)。
*/
class DefaultKnowledgeRetrievalFacadeTest extends BaseMockitoUnitTest {
@ -33,18 +35,37 @@ class DefaultKnowledgeRetrievalFacadeTest extends BaseMockitoUnitTest {
return new ContextAssemblyRequest(100L, 2001L, 4001L, "如何写好开头", 5, "corr-1");
}
/** §5.3 全字段 chunk。 */
private RetrievedChunk chunk() {
return new RetrievedChunk(5001L, "ds-1", "doc-9", "开头要抓人", 0.88,
"user", "3", "auth-1", "active", "search,generate");
}
@Test
void should_assembleContext_whenApiReturnsChunks() {
when(knowledgeRetrievalApi.retrieveForWork(any(RetrievalRequest.class))).thenReturn(
RetrievalResult.ok(List.of(new RetrievedChunk(5001L, "ds-1", "doc-9", "开头要抓人", 0.88))));
when(knowledgeRetrievalApi.retrieveForWork(any(RetrievalRequest.class)))
.thenReturn(RetrievalResult.ok(List.of(chunk()), List.of()));
ContextAssembly result = facade.assembleContext(req());
assertTrue(result.hasContext());
assertEquals(1, result.chunkCount());
// promptSection 带来源归因(kbId)+ 正文,供 §5.3 可追溯
assertEquals(0, result.omittedSourceCount());
// promptSection 带来源归因(kbId/owner)+ 正文,供 §5.3 可追溯
assertTrue(result.promptSection().contains("kbId=5001"));
assertTrue(result.promptSection().contains("owner=user"));
assertTrue(result.promptSection().contains("开头要抓人"));
}
@Test
void should_carryOmittedSourceCount_whenSomeSourcesFiltered() {
// 部分来源被授权门过滤:命中仍返回,但透出 omittedSourceCount 供审计
when(knowledgeRetrievalApi.retrieveForWork(any(RetrievalRequest.class)))
.thenReturn(RetrievalResult.ok(List.of(chunk()),
List.of(new OmittedSource(6001L, "not_authorized", null))));
ContextAssembly result = facade.assembleContext(req());
assertTrue(result.hasContext());
assertEquals(1, result.omittedSourceCount());
}
@Test
void should_empty_whenApiReturnsEmpty() {
when(knowledgeRetrievalApi.retrieveForWork(any(RetrievalRequest.class)))
@ -52,6 +73,19 @@ class DefaultKnowledgeRetrievalFacadeTest extends BaseMockitoUnitTest {
ContextAssembly result = facade.assembleContext(req());
assertFalse(result.hasContext());
assertEquals("no_binding", result.omittedReason());
assertEquals(0, result.omittedSourceCount());
}
@Test
void should_emptyWithOmittedCount_whenAllSourcesFiltered() {
// 全被授权门拒:empty 仍透出 omittedSourceCount(专题-03 §4 omittedSources 透明)
when(knowledgeRetrievalApi.retrieveForWork(any(RetrievalRequest.class)))
.thenReturn(RetrievalResult.empty("not_authorized",
List.of(new OmittedSource(5001L, "not_authorized", "revoked"))));
ContextAssembly result = facade.assembleContext(req());
assertFalse(result.hasContext());
assertEquals("not_authorized", result.omittedReason());
assertEquals(1, result.omittedSourceCount());
}
@Test

View File

@ -3,16 +3,18 @@ package cn.iocoder.muse.module.knowledge.api;
import java.util.List;
/**
* Knowledge 对外只读检索 API(BC 对外契约,专题-03 P-A 最小检索串联)。
* Knowledge 对外只读检索 API(BC 对外契约,专题-03 AI 编排上下文检索串联)。
*
* <p>供他域(AI 生成链)在<b>不直连 knowledge DAL/application</b> 的前提下,按作品检索其绑定知识库的相关片段,
* 用于 AI 上下文组装(专题-03 §4.2 Layer 3 授权资料)。由 knowledge-server 的
* {@code MuseKnowledgeRetrievalApiImpl} 实现(读本域 work→binding→kb→active dataset 映射 + 调 RAGFlow retrieveChunks);
* 消费方依赖本 -api 契约,不得依赖 {@code knowledge.dal}/{@code knowledge.application}(ArchUnit BcBoundaryArchTest 机械约束)。</p>
* {@code MuseKnowledgeRetrievalApiImpl} 实现;消费方依赖本 -api 契约,不得依赖 {@code knowledge.dal}/{@code knowledge.application}
* (ArchUnit BcBoundaryArchTest 机械约束)。</p>
*
* <p>P-A 最小授权门:仅检索 active binding 的 kb;tenantId+ownerUserId 必填、fail-closed。
* 完整按用途授权过滤(专题-03 §4.3/§5.3)+ 授权快照/allowedPurpose/sourceStatus 字段属 P-B。
* 检索失败/无 chunk/外部不可用一律返回空结果(带 omittedReason),<b>不抛异常阻断主生成链</b>。</p>
* <p>授权 fail-closed(P-B,专题-03 §4.3/§5.3):两道门——<b>用途门</b>(binding 须授权"上下文检索"用途)+
* <b>来源状态门</b>(来源状态 ∈ 阻断集一律拒);通过的 chunk 须带齐 §5.3 合同字段
* (sourceOwner/sourceObjectVersion/authorizationSnapshot/sourceStatus/allowedPurpose),<b>缺授权快照不进 Prompt</b>。
* 被过滤的来源记入 {@link OmittedSource} 供审计与 token 预算透明。检索失败/无 chunk/外部不可用一律返回空结果
* (带 omittedReason),<b>不抛异常阻断主生成链</b>。</p>
*/
public interface MuseKnowledgeRetrievalApi {
@ -20,7 +22,7 @@ public interface MuseKnowledgeRetrievalApi {
* 按作品检索其绑定知识库的相关片段。
*
* @param request 检索入参(tenantId/ownerUserId/workId/question 必填)
* @return 脱敏检索结果;无绑定/无 dataset/检索失败/无 chunk 时返回 {@link RetrievalResult#empty(String)}(不抛)
* @return 脱敏检索结果;无绑定/未授权/检索失败/无 chunk 时返回 {@link RetrievalResult#empty(String)}(不抛)
*/
RetrievalResult retrieveForWork(RetrievalRequest request);
@ -32,19 +34,26 @@ public interface MuseKnowledgeRetrievalApi {
/**
* 检索结果(脱敏)。
*
* @param chunks 命中的片段(可空集)
* @param status "ok"=有命中 / "empty"=无命中或降级
* @param omittedReason status=empty 时的原因(no_binding/no_dataset/retrieval_failed/no_chunk/blank_question 等),
* 供审计与上下文 omittedSources(专题-03 §4 token 预算省略枚举)
* @param chunks 命中且通过授权门的片段(可空集)
* @param status "ok"=有命中 / "empty"=无命中或降级
* @param omittedReason status=empty 时的总体原因(no_binding/no_dataset/not_authorized/retrieval_failed/no_chunk/blank_question 等)
* @param omittedSources 被授权门/token 预算过滤掉的来源明细(专题-03 §4 omittedSources,供审计与透明);可空集
*/
record RetrievalResult(List<RetrievedChunk> chunks, String status, String omittedReason) {
record RetrievalResult(List<RetrievedChunk> chunks, String status, String omittedReason,
List<OmittedSource> omittedSources) {
public static RetrievalResult ok(List<RetrievedChunk> chunks) {
return new RetrievalResult(chunks == null ? List.of() : chunks, "ok", null);
public static RetrievalResult ok(List<RetrievedChunk> chunks, List<OmittedSource> omittedSources) {
return new RetrievalResult(chunks == null ? List.of() : chunks, "ok", null,
omittedSources == null ? List.of() : omittedSources);
}
public static RetrievalResult empty(String omittedReason) {
return new RetrievalResult(List.of(), "empty", omittedReason);
return new RetrievalResult(List.of(), "empty", omittedReason, List.of());
}
public static RetrievalResult empty(String omittedReason, List<OmittedSource> omittedSources) {
return new RetrievalResult(List.of(), "empty", omittedReason,
omittedSources == null ? List.of() : omittedSources);
}
public boolean hasChunks() {
@ -53,15 +62,32 @@ public interface MuseKnowledgeRetrievalApi {
}
/**
* 单条检索片段(脱敏,对齐专题-03 §5.3 检索结果合同基础字段)。
* 单条检索片段(脱敏,对齐专题-03 §5.3 检索结果合同:每条 chunk 须带齐授权可追溯字段,缺任一不进 Prompt)。
*
* @param sourceKbId 来源知识库 ID(sourceObject)
* @param datasetId RAGFlow dataset ID(归因)
* @param documentId 来源文档 ID(归因,可空)
* @param contentSummary 片段正文摘要(evidence,已截断脱敏、不含 token/raw 响应)
* @param similarity 相似度(confidence,可空)
* @param sourceKbId 来源知识库 ID(§5.3 sourceObject)
* @param datasetId RAGFlow dataset ID(归因)
* @param documentId 来源文档 ID(归因,可空)
* @param contentSummary 片段正文摘要(§5.3 evidence,已截断脱敏、不含 token/raw 响应)
* @param similarity 相似度(§5.3 confidence,可空)
* @param sourceOwner 来源归属(§5.3 sourceOwner,如 user/global/market)
* @param sourceObjectVersion 来源对象版本(§5.3 sourceObject+version,可空)
* @param authorizationSnapshotId 授权快照 ID(§5.3 authorizationSnapshot,绑定即检索授权的凭据)
* @param sourceStatus 来源状态(§5.3/§4.3 sourceStatus,如 active/stale)
* @param allowedPurpose 授权用途(§5.3 allowedPurpose,逗号分隔的 binding 用途)
*/
record RetrievedChunk(Long sourceKbId, String datasetId, String documentId,
String contentSummary, Double similarity) {
String contentSummary, Double similarity,
String sourceOwner, String sourceObjectVersion,
String authorizationSnapshotId, String sourceStatus, String allowedPurpose) {
}
/**
* 被过滤来源明细(专题-03 §4 omittedSources)。
*
* @param kbId 被省略的知识库 ID
* @param reason 原因枚举(not_authorized/stale_source/token_budget/no_dataset 等)
* @param sourceStatus 当时来源状态(供审计)
*/
record OmittedSource(Long kbId, String reason, String sourceStatus) {
}
}

View File

@ -7,45 +7,65 @@ import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledge
import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient.Status;
import cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeBindingDO;
import cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeRagflowBindingDO;
import cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeSourceBindingProjectionDO;
import cn.iocoder.muse.module.knowledge.dal.mysql.muse.MuseKnowledgeBindingMapper;
import cn.iocoder.muse.module.knowledge.dal.mysql.muse.MuseKnowledgeRagflowBindingMapper;
import cn.iocoder.muse.module.knowledge.dal.mysql.muse.MuseKnowledgeSourceBindingProjectionMapper;
import com.fasterxml.jackson.databind.JsonNode;
import jakarta.annotation.Resource;
import org.springframework.stereotype.Service;
import org.springframework.util.StringUtils;
import java.util.ArrayList;
import java.util.HashSet;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
/**
* Knowledge 对外检索 API 实现(BC 本域读本域:work→active binding→kb→active dataset + 调 RAGFlow retrieveChunks)。
* Knowledge 对外检索 API 实现(BC 本域读本域:work→binding/projection→kb→active dataset + 调 RAGFlow retrieveChunks)。
*
* <p>P-A 最小检索串联:仅 active binding 的 kb;tenant/owner/work/question 缺一 fail-closed(不裸跑)。
* chunk 从 {@code RuntimeResult.summary.responseBody}(RAGFlow {@code /api/v1/retrieval} 的 {@code data.chunks[]})
* 解析 + 截断脱敏。检索失败/无 chunk/任何异常一律返回 {@code empty(omittedReason)},<b>不抛异常阻断主生成链</b>。
* 完整按用途授权过滤 + 授权快照/allowedPurpose(专题-03 §4.3/§5.3)属 P-B。</p>
* <p>授权 fail-closed(专题-03 §4.3/§5.3)两道门:<b>用途门</b>(binding.bindingScope 须含 "search" 上下文检索用途;
* 阅读≠检索)+ <b>来源状态门</b>(projection.status ∈ {@link #BLOCKED_SOURCE_STATUS} 一律拒);
* 另要求 binding.authorizationSnapshotId 非空(§5.3 缺授权快照不进 Prompt)。通过的 chunk 带齐 §5.3 合同字段
* (sourceOwner/sourceObjectVersion/authorizationSnapshot/sourceStatus/allowedPurpose),被过滤来源记入 omittedSources。</p>
*
* <p>chunk 从 {@code RuntimeResult.summary.responseBody}(RAGFlow {@code /api/v1/retrieval} 的 {@code data.chunks[]})
* 解析 + 截断脱敏。tenant/owner/work/question 缺一 fail-closed;检索失败/无 chunk/任何异常一律返回
* {@code empty(omittedReason)},<b>不抛异常阻断主生成链</b>。</p>
*/
@Service
public class MuseKnowledgeRetrievalApiImpl implements MuseKnowledgeRetrievalApi {
/** 单 chunk 正文摘要最大长度(token 预算意识,防 prompt 撑爆;专题-03 §4)。 */
private static final int MAX_CHUNK_SUMMARY_CHARS = 600;
/** P-A 默认 topK / 相似度阈值。 */
/** 默认 topK / 相似度阈值。 */
private static final int DEFAULT_TOP_K = 5;
private static final double DEFAULT_THRESHOLD = 0.2d;
/** 授权"上下文检索"的用途标识(专题-03 §4.3:阅读≠上下文检索;binding.bindingScope 须含此用途才放行检索)。 */
private static final String PURPOSE_SEARCH = "search";
/** §4.3 来源状态阻断集:一律拒检索(active/stale 不在内,stale 为降级、非阻断、可纳入)。 */
private static final Set<String> BLOCKED_SOURCE_STATUS = Set.of(
"revoked", "recalled", "delisted", "blocked", "owner_missing", "unauthorized");
@Resource
private MuseKnowledgeBindingMapper bindingMapper;
@Resource
private MuseKnowledgeSourceBindingProjectionMapper projectionMapper;
@Resource
private MuseKnowledgeRagflowBindingMapper ragflowBindingMapper;
@Resource
private RagFlowKnowledgeRuntimeClient ragFlowClient;
/** impl 内部:通过双门的授权来源(承载检索 datasetId + chunk §5.3 字段填充所需授权元数据)。 */
private record AuthorizedSource(Long kbId, String datasetId, String sourceOwner, String sourceObjectVersion,
String authorizationSnapshotId, String sourceStatus, String allowedPurpose) {
}
@Override
public RetrievalResult retrieveForWork(RetrievalRequest request) {
// 最小授权门:tenant/owner/work/question 必填,缺一 fail-closed(不裸跑、不放任意检索)
// 入参门:tenant/owner/work/question 必填,缺一 fail-closed(不裸跑、不放任意检索)
if (request == null || request.tenantId() == null || request.ownerUserId() == null
|| request.workId() == null) {
return RetrievalResult.empty("invalid_request");
@ -54,58 +74,102 @@ public class MuseKnowledgeRetrievalApiImpl implements MuseKnowledgeRetrievalApi
return RetrievalResult.empty("blank_question");
}
try {
// 1. work→active binding(一作品可绑多 kb)
// 1. work→active binding(用途 + 授权快照来源)
List<MuseKnowledgeBindingDO> bindings = bindingMapper.selectActiveByWorkId(request.workId());
if (bindings == null || bindings.isEmpty()) {
return RetrievalResult.empty("no_binding");
}
// 2. 每 kb→active dataset,收集 datasetId 与 dataset→kbId 反查映射(去重)
List<String> datasetIds = new ArrayList<>();
Map<String, Long> datasetToKb = new LinkedHashMap<>();
// 2. work→来源绑定投影(来源状态 + sourceOwner/revision),按 kbId 索引(selectActiveByWorkId 已按 updateTime desc,首个=最新)
Map<Long, MuseKnowledgeSourceBindingProjectionDO> projectionByKb = new LinkedHashMap<>();
List<MuseKnowledgeSourceBindingProjectionDO> projections = projectionMapper.selectActiveByWorkId(request.workId());
if (projections != null) {
for (MuseKnowledgeSourceBindingProjectionDO projection : projections) {
if (projection.getKbId() != null) {
projectionByKb.putIfAbsent(projection.getKbId(), projection);
}
}
}
// 3. 双门过滤(§4.3 用途 + 来源状态 + 授权快照)→ 授权来源 + 被省略来源
List<AuthorizedSource> authorized = new ArrayList<>();
List<OmittedSource> omitted = new ArrayList<>();
Set<Long> seenKb = new HashSet<>();
for (MuseKnowledgeBindingDO binding : bindings) {
Long kbId = binding.getKbId();
if (kbId == null) {
if (kbId == null || !seenKb.add(kbId)) {
continue;
}
// §4.3 用途门:bindingScope 须含 "search" 上下文检索用途(阅读/导出等用途不放行检索)
if (!hasPurpose(binding.getBindingScope(), PURPOSE_SEARCH)) {
omitted.add(new OmittedSource(kbId, "not_authorized", null));
continue;
}
// §5.3 授权快照门:缺 authorizationSnapshot 的来源不进 Prompt
if (!StringUtils.hasText(binding.getAuthorizationSnapshotId())) {
omitted.add(new OmittedSource(kbId, "not_authorized", null));
continue;
}
// §4.3 来源状态门:无投影(来源状态未知)→ fail-closed;状态 ∈ 阻断集 → 拒
MuseKnowledgeSourceBindingProjectionDO projection = projectionByKb.get(kbId);
if (projection == null) {
omitted.add(new OmittedSource(kbId, "not_authorized", null));
continue;
}
String sourceStatus = projection.getStatus();
if (sourceStatus != null && BLOCKED_SOURCE_STATUS.contains(sourceStatus)) {
omitted.add(new OmittedSource(kbId, "not_authorized", sourceStatus));
continue;
}
// kb→active dataset
MuseKnowledgeRagflowBindingDO dataset = ragflowBindingMapper.selectActiveDatasetByKbId(kbId);
if (dataset == null || !StringUtils.hasText(dataset.getRagflowDatasetId())) {
omitted.add(new OmittedSource(kbId, "no_dataset", sourceStatus));
continue;
}
if (!datasetToKb.containsKey(dataset.getRagflowDatasetId())) {
datasetIds.add(dataset.getRagflowDatasetId());
datasetToKb.put(dataset.getRagflowDatasetId(), kbId);
authorized.add(new AuthorizedSource(kbId, dataset.getRagflowDatasetId(),
projection.getSourceOwner(), projection.getSourceRevision(),
binding.getAuthorizationSnapshotId(), sourceStatus, binding.getBindingScope()));
}
if (authorized.isEmpty()) {
// 全被门拒:授权问题优先暴露,否则归 no_dataset
String reason = omitted.stream().anyMatch(o -> "not_authorized".equals(o.reason()))
? "not_authorized" : "no_dataset";
return RetrievalResult.empty(reason, omitted);
}
// 4. 一次跨所有授权 dataset 检索(RAGFlow 支持多 dataset_ids);chunk 再按 dataset 反查授权来源
List<String> datasetIds = new ArrayList<>();
Map<String, AuthorizedSource> datasetToSource = new LinkedHashMap<>();
for (AuthorizedSource source : authorized) {
if (datasetToSource.putIfAbsent(source.datasetId(), source) == null) {
datasetIds.add(source.datasetId());
}
}
if (datasetIds.isEmpty()) {
return RetrievalResult.empty("no_dataset");
}
// 3. 一次跨所有 active dataset 检索(RAGFlow 支持多 dataset_ids);kbId 取首个供审计、chunk 再按 dataset 反查
Long primaryKbId = datasetToKb.get(datasetIds.get(0));
AuthorizedSource primary = authorized.get(0);
int topK = request.topK() != null && request.topK() > 0 ? request.topK() : DEFAULT_TOP_K;
RetrieveChunksCommand command = new RetrieveChunksCommand(
request.tenantId(), request.ownerUserId(), primaryKbId, datasetIds, null,
request.tenantId(), request.ownerUserId(), primary.kbId(), datasetIds, null,
request.question(), topK, DEFAULT_THRESHOLD, null,
request.correlationId(), null, 1);
RuntimeResult result = ragFlowClient.retrieveChunks(command);
if (result == null || result.status() != Status.SUCCEEDED) {
String reason = result == null || result.failureClass() == null
? "retrieval_failed" : "retrieval_failed:" + result.failureClass().name();
return RetrievalResult.empty(reason);
return RetrievalResult.empty(reason, omitted);
}
// 4. 解析 data.chunks[] → 脱敏 chunk(dataset_id 反查 kbId)
List<RetrievedChunk> chunks = parseChunks(result, datasetToKb, primaryKbId);
// 5. 解析 data.chunks[] → 补 §5.3 合同字段(按 dataset 反查授权来源)
List<RetrievedChunk> chunks = parseChunks(result, datasetToSource, primary);
if (chunks.isEmpty()) {
return RetrievalResult.empty("no_chunk");
return RetrievalResult.empty("no_chunk", omitted);
}
return RetrievalResult.ok(chunks);
return RetrievalResult.ok(chunks, omitted);
} catch (RuntimeException ex) {
// fail-closed 降级:检索任何异常都不阻断主生成链,返回 empty 供 omittedSources 记录
// fail-closed 降级:检索任何异常都不阻断主生成链
return RetrievalResult.empty("retrieval_error");
}
}
/** 从 RAGFlow 检索响应树 data.chunks[] 解析脱敏片段;dataset_id 反查 kbId,无则归 primary。 */
private List<RetrievedChunk> parseChunks(RuntimeResult result, Map<String, Long> datasetToKb, Long primaryKbId) {
/** 从 RAGFlow 检索响应树 data.chunks[] 解析脱敏片段;按 dataset 反查授权来源补齐 §5.3 字段,反查不到归 primary。 */
private List<RetrievedChunk> parseChunks(RuntimeResult result, Map<String, AuthorizedSource> datasetToSource,
AuthorizedSource primary) {
List<RetrievedChunk> chunks = new ArrayList<>();
Object responseBody = result.summary() == null ? null : result.summary().get("responseBody");
JsonNode tree = toJsonNode(responseBody);
@ -123,14 +187,29 @@ public class MuseKnowledgeRetrievalApiImpl implements MuseKnowledgeRetrievalApi
}
String datasetId = firstText(chunk, "dataset_id", "kb_id");
String documentId = firstText(chunk, "document_id", "doc_id");
Long sourceKbId = datasetId != null && datasetToKb.containsKey(datasetId)
? datasetToKb.get(datasetId) : primaryKbId;
AuthorizedSource source = datasetId != null && datasetToSource.containsKey(datasetId)
? datasetToSource.get(datasetId) : primary;
Double similarity = chunk.path("similarity").isNumber() ? chunk.path("similarity").asDouble() : null;
chunks.add(new RetrievedChunk(sourceKbId, datasetId, documentId, summarize(content), similarity));
chunks.add(new RetrievedChunk(source.kbId(), datasetId, documentId, summarize(content), similarity,
source.sourceOwner(), source.sourceObjectVersion(), source.authorizationSnapshotId(),
source.sourceStatus(), source.allowedPurpose()));
}
return chunks;
}
/** bindingScope(逗号分隔的授权用途)是否含指定用途。 */
private boolean hasPurpose(String bindingScope, String purpose) {
if (!StringUtils.hasText(bindingScope)) {
return false;
}
for (String scope : bindingScope.split(",")) {
if (purpose.equals(scope.trim())) {
return true;
}
}
return false;
}
private JsonNode toJsonNode(Object responseBody) {
if (responseBody == null) {
return null;

View File

@ -12,8 +12,10 @@ import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledge
import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient.Status;
import cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeBindingDO;
import cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeRagflowBindingDO;
import cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeSourceBindingProjectionDO;
import cn.iocoder.muse.module.knowledge.dal.mysql.muse.MuseKnowledgeBindingMapper;
import cn.iocoder.muse.module.knowledge.dal.mysql.muse.MuseKnowledgeRagflowBindingMapper;
import cn.iocoder.muse.module.knowledge.dal.mysql.muse.MuseKnowledgeSourceBindingProjectionMapper;
import org.junit.jupiter.api.Test;
import org.mockito.InjectMocks;
import org.mockito.Mock;
@ -30,7 +32,8 @@ import static org.mockito.Mockito.verifyNoInteractions;
import static org.mockito.Mockito.when;
/**
* P-A 检索串联 knowledge 侧单测:work→binding→kb→dataset→retrieveChunks 链路 + 各 fail-closed 分支。
* P-B 检索串联 + 授权 fail-closed 单测:work→binding/projection→kb→dataset→retrieveChunks 链路
* + 双门(§4.3 用途/来源状态)+ 授权快照门 + §5.3 chunk 字段 + omittedSources。
*/
class MuseKnowledgeRetrievalApiImplTest extends BaseMockitoUnitTest {
@ -39,6 +42,8 @@ class MuseKnowledgeRetrievalApiImplTest extends BaseMockitoUnitTest {
@Mock
private MuseKnowledgeBindingMapper bindingMapper;
@Mock
private MuseKnowledgeSourceBindingProjectionMapper projectionMapper;
@Mock
private MuseKnowledgeRagflowBindingMapper ragflowBindingMapper;
@Mock
private RagFlowKnowledgeRuntimeClient ragFlowClient;
@ -47,13 +52,30 @@ class MuseKnowledgeRetrievalApiImplTest extends BaseMockitoUnitTest {
return new RetrievalRequest(100L, 2001L, 4001L, "如何写好开头", 5, "corr-1");
}
private MuseKnowledgeBindingDO activeBinding(Long kbId) {
/** binding:bindingScope=授权用途(逗号分隔), authorizationSnapshotId=授权快照。 */
private MuseKnowledgeBindingDO binding(Long kbId, String bindingScope, String authSnap) {
MuseKnowledgeBindingDO b = new MuseKnowledgeBindingDO();
b.setKbId(kbId);
b.setBindingStatus("active");
b.setBindingScope(bindingScope);
b.setAuthorizationSnapshotId(authSnap);
return b;
}
/** 标准授权 binding:含 search 检索用途 + 授权快照。 */
private MuseKnowledgeBindingDO authorizedBinding(Long kbId) {
return binding(kbId, "search,generate", "auth-" + kbId);
}
private MuseKnowledgeSourceBindingProjectionDO projection(Long kbId, String status) {
MuseKnowledgeSourceBindingProjectionDO p = new MuseKnowledgeSourceBindingProjectionDO();
p.setKbId(kbId);
p.setStatus(status);
p.setSourceOwner("user");
p.setSourceRevision("3");
return p;
}
private MuseKnowledgeRagflowBindingDO dataset(String datasetId) {
MuseKnowledgeRagflowBindingDO d = new MuseKnowledgeRagflowBindingDO();
d.setRagflowDatasetId(datasetId);
@ -72,14 +94,14 @@ class MuseKnowledgeRetrievalApiImplTest extends BaseMockitoUnitTest {
RetrievalResult result = retrievalApi.retrieveForWork(new RetrievalRequest(100L, null, 4001L, "q", 5, "c"));
assertEquals("invalid_request", result.omittedReason());
assertFalse(result.hasChunks());
verifyNoInteractions(bindingMapper, ragFlowClient);
verifyNoInteractions(bindingMapper, projectionMapper, ragFlowClient);
}
@Test
void should_returnBlankQuestion_whenQuestionEmpty() {
RetrievalResult result = retrievalApi.retrieveForWork(new RetrievalRequest(100L, 2001L, 4001L, " ", 5, "c"));
assertEquals("blank_question", result.omittedReason());
verifyNoInteractions(bindingMapper, ragFlowClient);
verifyNoInteractions(bindingMapper, projectionMapper, ragFlowClient);
}
@Test
@ -87,41 +109,95 @@ class MuseKnowledgeRetrievalApiImplTest extends BaseMockitoUnitTest {
when(bindingMapper.selectActiveByWorkId(4001L)).thenReturn(List.of());
RetrievalResult result = retrievalApi.retrieveForWork(req());
assertEquals("no_binding", result.omittedReason());
verifyNoInteractions(ragFlowClient);
verifyNoInteractions(projectionMapper, ragFlowClient);
}
@Test
void should_returnNoDataset_whenBindingHasNoActiveDataset() {
when(bindingMapper.selectActiveByWorkId(4001L)).thenReturn(List.of(activeBinding(5001L)));
void should_omitNotAuthorized_whenBindingScopeLacksSearch() {
// §4.3 用途门:bindingScope 仅 export(无 search 上下文检索用途)→ 拒
when(bindingMapper.selectActiveByWorkId(4001L)).thenReturn(List.of(binding(5001L, "export", "auth-5001")));
when(projectionMapper.selectActiveByWorkId(4001L)).thenReturn(List.of(projection(5001L, "active")));
RetrievalResult result = retrievalApi.retrieveForWork(req());
assertEquals("not_authorized", result.omittedReason());
assertEquals(1, result.omittedSources().size());
assertEquals(5001L, result.omittedSources().get(0).kbId());
assertEquals("not_authorized", result.omittedSources().get(0).reason());
verifyNoInteractions(ragflowBindingMapper, ragFlowClient);
}
@Test
void should_omitNotAuthorized_whenAuthorizationSnapshotMissing() {
// §5.3 授权快照门:缺 authorizationSnapshot 的来源不进 Prompt
when(bindingMapper.selectActiveByWorkId(4001L)).thenReturn(List.of(binding(5001L, "search", null)));
when(projectionMapper.selectActiveByWorkId(4001L)).thenReturn(List.of(projection(5001L, "active")));
RetrievalResult result = retrievalApi.retrieveForWork(req());
assertEquals("not_authorized", result.omittedReason());
verifyNoInteractions(ragflowBindingMapper, ragFlowClient);
}
@Test
void should_omitNotAuthorized_whenSourceStatusBlocked() {
// §4.3 来源状态门:revoked ∈ 阻断集 → 拒
when(bindingMapper.selectActiveByWorkId(4001L)).thenReturn(List.of(authorizedBinding(5001L)));
when(projectionMapper.selectActiveByWorkId(4001L)).thenReturn(List.of(projection(5001L, "revoked")));
RetrievalResult result = retrievalApi.retrieveForWork(req());
assertEquals("not_authorized", result.omittedReason());
assertEquals("revoked", result.omittedSources().get(0).sourceStatus());
verifyNoInteractions(ragflowBindingMapper, ragFlowClient);
}
@Test
void should_omitNotAuthorized_whenProjectionMissing() {
// 来源状态未知(active binding 但无投影)→ fail-closed 拒
when(bindingMapper.selectActiveByWorkId(4001L)).thenReturn(List.of(authorizedBinding(5001L)));
when(projectionMapper.selectActiveByWorkId(4001L)).thenReturn(List.of());
RetrievalResult result = retrievalApi.retrieveForWork(req());
assertEquals("not_authorized", result.omittedReason());
verifyNoInteractions(ragflowBindingMapper, ragFlowClient);
}
@Test
void should_returnNoDataset_whenAuthorizedButNoActiveDataset() {
when(bindingMapper.selectActiveByWorkId(4001L)).thenReturn(List.of(authorizedBinding(5001L)));
when(projectionMapper.selectActiveByWorkId(4001L)).thenReturn(List.of(projection(5001L, "active")));
when(ragflowBindingMapper.selectActiveDatasetByKbId(5001L)).thenReturn(null);
RetrievalResult result = retrievalApi.retrieveForWork(req());
assertEquals("no_dataset", result.omittedReason());
assertEquals("no_dataset", result.omittedSources().get(0).reason());
verifyNoInteractions(ragFlowClient);
}
@Test
void should_returnChunks_whenRetrievalSucceeds() {
when(bindingMapper.selectActiveByWorkId(4001L)).thenReturn(List.of(activeBinding(5001L)));
void should_returnChunksWithContractFields_whenAuthorizedAndRetrievalSucceeds() {
when(bindingMapper.selectActiveByWorkId(4001L)).thenReturn(List.of(authorizedBinding(5001L)));
when(projectionMapper.selectActiveByWorkId(4001L)).thenReturn(List.of(projection(5001L, "active")));
when(ragflowBindingMapper.selectActiveDatasetByKbId(5001L)).thenReturn(dataset("ds-1"));
// 用 RAGFlow /api/v1/retrieval 官方真实响应结构(chunk 正文为 content、dataset 标识为 kb_id)验容错解析与 kb_id→muse kbId 反查
// RAGFlow /api/v1/retrieval 官方真实结构(content 正文 + kb_id 指 dataset)
String json = "{\"code\":0,\"data\":{\"chunks\":[{"
+ "\"content\":\"开头要抓人\",\"content_ltks\":\"开头 抓人\",\"document_id\":\"doc-9\","
+ "\"document_keyword\":\"x.txt\",\"id\":\"chunk-1\",\"kb_id\":\"ds-1\","
+ "\"similarity\":0.88,\"term_similarity\":1.0,\"vector_similarity\":0.8}],\"total\":1}}";
+ "\"content\":\"开头要抓人\",\"document_id\":\"doc-9\",\"id\":\"chunk-1\",\"kb_id\":\"ds-1\","
+ "\"similarity\":0.88,\"vector_similarity\":0.8}],\"total\":1}}";
when(ragFlowClient.retrieveChunks(any(RetrieveChunksCommand.class))).thenReturn(retrievalSuccess(json));
RetrievalResult result = retrievalApi.retrieveForWork(req());
assertTrue(result.hasChunks());
assertEquals(1, result.chunks().size());
assertEquals(5001L, result.chunks().get(0).sourceKbId());
assertEquals("ds-1", result.chunks().get(0).datasetId());
assertEquals("开头要抓人", result.chunks().get(0).contentSummary());
assertEquals("doc-9", result.chunks().get(0).documentId());
assertEquals(0.88, result.chunks().get(0).similarity(), 0.0001);
MuseKnowledgeRetrievalApi.RetrievedChunk chunk = result.chunks().get(0);
assertEquals(5001L, chunk.sourceKbId());
assertEquals("ds-1", chunk.datasetId());
assertEquals("doc-9", chunk.documentId());
assertEquals("开头要抓人", chunk.contentSummary());
assertEquals(0.88, chunk.similarity(), 0.0001);
// §5.3 合同字段全带齐
assertEquals("user", chunk.sourceOwner());
assertEquals("3", chunk.sourceObjectVersion());
assertEquals("auth-5001", chunk.authorizationSnapshotId());
assertEquals("active", chunk.sourceStatus());
assertEquals("search,generate", chunk.allowedPurpose());
}
@Test
void should_returnRetrievalFailed_whenClientFails() {
when(bindingMapper.selectActiveByWorkId(4001L)).thenReturn(List.of(activeBinding(5001L)));
when(bindingMapper.selectActiveByWorkId(4001L)).thenReturn(List.of(authorizedBinding(5001L)));
when(projectionMapper.selectActiveByWorkId(4001L)).thenReturn(List.of(projection(5001L, "active")));
when(ragflowBindingMapper.selectActiveDatasetByKbId(5001L)).thenReturn(dataset("ds-1"));
RuntimeResult failed = new RuntimeResult(null, Operation.RETRIEVE_CHUNKS, "c", null, 1, 10L,
Status.FAILED, FailureClass.RAGFLOW_UNAVAILABLE, null, "r", List.of(), Map.of());
@ -133,7 +209,8 @@ class MuseKnowledgeRetrievalApiImplTest extends BaseMockitoUnitTest {
@Test
void should_returnNoChunk_whenResponseHasEmptyChunks() {
when(bindingMapper.selectActiveByWorkId(4001L)).thenReturn(List.of(activeBinding(5001L)));
when(bindingMapper.selectActiveByWorkId(4001L)).thenReturn(List.of(authorizedBinding(5001L)));
when(projectionMapper.selectActiveByWorkId(4001L)).thenReturn(List.of(projection(5001L, "active")));
when(ragflowBindingMapper.selectActiveDatasetByKbId(5001L)).thenReturn(dataset("ds-1"));
when(ragFlowClient.retrieveChunks(any(RetrieveChunksCommand.class)))
.thenReturn(retrievalSuccess("{\"code\":0,\"data\":{\"chunks\":[]}}"));