refactor(knowledge): S6a 端口重命名 RagFlow→Knowledge + 裁剪无生产调用操作

S6 第一子步(机械、provider 无关):接口 RagFlowKnowledgeRuntimeClient→KnowledgeRuntimeClient
全量重命名(207 引用,裸旧名 PCRE 精确核验 0 残留);裁掉 7 个已核验无生产调用操作
(updateDatasetConfig/listDatasets/listChunks/runGraphRag/traceGraphRag/getKnowledgeGraph/health,
Command/枚举/实现体/单测清零),接口保留 5 操作(建库/上传/触发解析/轮询/检索);审计
recordRagflowCall→recordRuntimeCall、RagflowCallRecordReq→RuntimeCallRecordReq(旧 token 0/新 13)。
REST 端 MuseKnowledgeGraphQueryService.getKnowledgeGraph(从 Muse 自有 PG 读图,同名不同方法)保留未动。
命名债白名单保留原名:表 muse_knowledge_ragflow_call/_binding、两 DO/Mapper、ragflowDatasetId/
DocumentId 列、KNOWLEDGE_RAGFLOW_* 错误码、FailureClass.RAGFLOW_*(S6 删类步再清)。

两个 opt-in live IT(P1rRagFlowLiveAcceptanceIT、server 的 P1rKnowledgeRuntimeEndToEnd...)原本真调了
被裁操作(health/listChunks/runGraphRag/traceGraphRag),做最小手术删被裁调用/断言、保留 5 操作
happy-path 覆盖(这俩 IT 本就在 S6 后续删/重做)。HttpRagFlowKnowledgeRuntimeClient 等实现类名 S6 删类步保留。

验证:knowledge 模块 285 用例全绿(BUILD SUCCESS);muse-server test-compile 默认 + market-assembled
双 profile 均 BUILD SUCCESS(被改 IT 新鲜编译)。此前 worktree 隔离基线陈旧 320 commit 作废、主树重做。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
lili 2026-07-08 01:03:52 -07:00
parent 0451549829
commit 2ac42a0628
21 changed files with 298 additions and 602 deletions

View File

@ -1,10 +1,10 @@
package cn.iocoder.muse.module.knowledge.api;
import cn.iocoder.muse.framework.common.util.json.JsonUtils;
import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient;
import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient.RetrieveChunksCommand;
import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient.RuntimeResult;
import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient.Status;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeRuntimeClient;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeRuntimeClient.RetrieveChunksCommand;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeRuntimeClient.RuntimeResult;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeRuntimeClient.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;
@ -56,7 +56,7 @@ public class MuseKnowledgeRetrievalApiImpl implements MuseKnowledgeRetrievalApi
@Resource
private MuseKnowledgeRagflowBindingMapper ragflowBindingMapper;
@Resource
private RagFlowKnowledgeRuntimeClient ragFlowClient;
private KnowledgeRuntimeClient ragFlowClient;
/** impl 内部:通过双门的授权来源(承载检索 datasetId + chunk §5.3 字段填充所需授权元数据)。 */
private record AuthorizedSource(Long kbId, String datasetId, String sourceOwner, String sourceObjectVersion,

View File

@ -3,7 +3,7 @@ package cn.iocoder.muse.module.knowledge.application.muse;
import cn.iocoder.muse.framework.common.util.json.JsonUtils;
import cn.iocoder.muse.framework.tenant.core.context.TenantContextHolder;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeFileFacade;
import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeRuntimeClient;
import cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeDocumentDO;
import cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeDocumentVersionDO;
import cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeRagflowBindingDO;
@ -42,8 +42,8 @@ import java.util.UUID;
*
* <p>外部交互铁律CLAUDE.mdfork 是多次 RAGFlow 外部调用的<b>异步长任务</b>全程处理超时 / 失败 / 重试 /
* 幂等 / 部分成功补偿幂等键 (assetId, version) ready 跳过partial/failed 只补差集失败按
* {@link RagFlowKnowledgeRuntimeClient.FailureClass} 分可重试 / 不可重试每次 RAGFlow 调用复用
* {@link MuseKnowledgeAuditService#recordRagflowCall} 审计attributionStatus=market_public_fork
* {@link KnowledgeRuntimeClient.FailureClass} 分可重试 / 不可重试每次 RAGFlow 调用复用
* {@link MuseKnowledgeAuditService#recordRuntimeCall} 审计attributionStatus=market_public_fork
* fork 成败<b>绝不</b>反向影响 market 审核market 只在审核事务内落上架事件即返回本服务异步消费</p>
*/
@Slf4j
@ -71,7 +71,7 @@ public class KnowledgeMarketForkService {
@Resource
private KnowledgeFileFacade fileFacade;
@Resource
private RagFlowKnowledgeRuntimeClient ragFlowClient;
private KnowledgeRuntimeClient ragFlowClient;
@Resource
private MuseKnowledgeAuditService auditService;
@Resource
@ -183,14 +183,14 @@ public class KnowledgeMarketForkService {
String correlationId = forkCorrelationId(assetIdText, versionText);
String datasetName = forkDatasetName(assetIdText, versionText);
RagFlowKnowledgeRuntimeClient.RuntimeResult createResult = ragFlowClient.createDataset(
new RagFlowKnowledgeRuntimeClient.CreateDatasetCommand(event.tenantId(), event.publisherUserId(),
KnowledgeRuntimeClient.RuntimeResult createResult = ragFlowClient.createDataset(
new KnowledgeRuntimeClient.CreateDatasetCommand(event.tenantId(), event.publisherUserId(),
event.publisherKbId(), datasetName, forkDatasetConfig(assetIdText, versionText),
correlationId, requestHash(assetIdText, versionText, "createDataset"), 1));
auditForkCall(event, createResult, createResult.externalId(), null, correlationId,
Map.of("source", ATTRIBUTION_MARKET_PUBLIC_FORK, "datasetName", datasetName,
"assetId", assetIdText, "assetVersion", versionText));
if (createResult.status() == RagFlowKnowledgeRuntimeClient.Status.FAILED
if (createResult.status() == KnowledgeRuntimeClient.Status.FAILED
|| createResult.externalId() == null || createResult.externalId().isBlank()) {
// 建副本失败可重试错误保留中间态pending下轮重建不可重试错误标 failed两者都写回 market fail-closed
String status = isRetryable(createResult.failureClass()) ? FORK_STATUS_PENDING : FORK_STATUS_FAILED;
@ -282,32 +282,32 @@ public class KnowledgeMarketForkService {
PublicDocument document, byte[] content) {
String assetIdText = String.valueOf(event.assetId());
String correlationId = forkCorrelationId(assetIdText, String.valueOf(document.documentId()));
RagFlowKnowledgeRuntimeClient.RuntimeResult uploadResult = ragFlowClient.uploadDocuments(
new RagFlowKnowledgeRuntimeClient.UploadDocumentsCommand(event.tenantId(), event.publisherUserId(),
KnowledgeRuntimeClient.RuntimeResult uploadResult = ragFlowClient.uploadDocuments(
new KnowledgeRuntimeClient.UploadDocumentsCommand(event.tenantId(), event.publisherUserId(),
event.publisherKbId(), ragflowDatasetId,
List.of(new RagFlowKnowledgeRuntimeClient.DocumentUpload(document.documentId(),
List.of(new KnowledgeRuntimeClient.DocumentUpload(document.documentId(),
document.documentVersionId(), document.storageRef(), document.fileName(),
document.contentType(), document.size(), content)),
correlationId, requestHash(assetIdText, String.valueOf(document.documentId()), "upload"), 1));
auditForkCall(event, uploadResult, ragflowDatasetId, uploadResult.externalId(), correlationId,
Map.of("source", ATTRIBUTION_MARKET_PUBLIC_FORK, "assetId", assetIdText,
"documentId", document.documentId()));
if (uploadResult.status() == RagFlowKnowledgeRuntimeClient.Status.FAILED
if (uploadResult.status() == KnowledgeRuntimeClient.Status.FAILED
|| uploadResult.externalId() == null || uploadResult.externalId().isBlank()) {
log.warn("[uploadAndParseToFork][副本文档上传失败assetId={}, documentId={}, failureClass={}]",
assetIdText, document.documentId(), uploadResult.failureClass());
return false;
}
RagFlowKnowledgeRuntimeClient.RuntimeResult parseResult = ragFlowClient.startParseDocuments(
new RagFlowKnowledgeRuntimeClient.StartParseDocumentsCommand(event.tenantId(), event.publisherUserId(),
KnowledgeRuntimeClient.RuntimeResult parseResult = ragFlowClient.startParseDocuments(
new KnowledgeRuntimeClient.StartParseDocumentsCommand(event.tenantId(), event.publisherUserId(),
event.publisherKbId(), ragflowDatasetId, List.of(uploadResult.externalId()),
forkProcessingTaskId(assetIdText, document.documentId()), correlationId,
requestHash(assetIdText, String.valueOf(document.documentId()), "startParse"), 1));
auditForkCall(event, parseResult, ragflowDatasetId, uploadResult.externalId(), correlationId,
Map.of("source", ATTRIBUTION_MARKET_PUBLIC_FORK, "assetId", assetIdText,
"documentId", document.documentId()));
if (parseResult.status() == RagFlowKnowledgeRuntimeClient.Status.FAILED) {
if (parseResult.status() == KnowledgeRuntimeClient.Status.FAILED) {
log.warn("[uploadAndParseToFork][副本文档重索引触发失败assetId={}, documentId={}, failureClass={}]",
assetIdText, document.documentId(), parseResult.failureClass());
return false;
@ -356,10 +356,10 @@ public class KnowledgeMarketForkService {
// ============================== 审计 ==============================
private void auditForkCall(MuseKnowledgeMarketListedEvent event, RagFlowKnowledgeRuntimeClient.RuntimeResult result,
private void auditForkCall(MuseKnowledgeMarketListedEvent event, KnowledgeRuntimeClient.RuntimeResult result,
String ragflowDatasetId, String ragflowDocumentId, String correlationId,
Map<String, Object> requestSummary) {
auditService.recordRagflowCall(MuseKnowledgeAuditService.RagflowCallRecordReq.builder()
auditService.recordRuntimeCall(MuseKnowledgeAuditService.RuntimeCallRecordReq.builder()
// 同一 fork 会有 createDataset / upload / startParse 多次外部调用审计 correlation 必须区分 operation避免唯一键冲突
.correlationId(correlationId + ":" + result.operation().wireName())
.attemptNo(result.attempt())
@ -386,13 +386,13 @@ public class KnowledgeMarketForkService {
.build());
}
private boolean isRetryable(RagFlowKnowledgeRuntimeClient.FailureClass failureClass) {
private boolean isRetryable(KnowledgeRuntimeClient.FailureClass failureClass) {
// 沿用上传链路的可重试分类临时-04 §3.5临时性外部故障可重试其余校验 / 文件拒绝等不可重试
return failureClass == RagFlowKnowledgeRuntimeClient.FailureClass.RAGFLOW_UNAVAILABLE
|| failureClass == RagFlowKnowledgeRuntimeClient.FailureClass.TIMEOUT
|| failureClass == RagFlowKnowledgeRuntimeClient.FailureClass.RATE_LIMITED
|| failureClass == RagFlowKnowledgeRuntimeClient.FailureClass.CONFIG_MISSING
|| failureClass == RagFlowKnowledgeRuntimeClient.FailureClass.AUTH_MISSING;
return failureClass == KnowledgeRuntimeClient.FailureClass.RAGFLOW_UNAVAILABLE
|| failureClass == KnowledgeRuntimeClient.FailureClass.TIMEOUT
|| failureClass == KnowledgeRuntimeClient.FailureClass.RATE_LIMITED
|| failureClass == KnowledgeRuntimeClient.FailureClass.CONFIG_MISSING
|| failureClass == KnowledgeRuntimeClient.FailureClass.AUTH_MISSING;
}
// ============================== 命名 / 配置 / 摘要工具 ==============================

View File

@ -1,7 +1,7 @@
package cn.iocoder.muse.module.knowledge.application.muse;
import cn.iocoder.muse.framework.tenant.core.context.TenantContextHolder;
import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeRuntimeClient;
import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRedactor;
import cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeRagflowCallDO;
import cn.iocoder.muse.module.knowledge.dal.mysql.muse.MuseKnowledgeRagflowCallMapper;
@ -15,7 +15,7 @@ import java.time.LocalDateTime;
import java.util.Locale;
/**
* Knowledge RAGFlow 审计服务
* Knowledge 运行时调用审计服务调用记录仍持久化到 muse_knowledge_ragflow_call 表名保留为命名债
*/
@Service
public class MuseKnowledgeAuditService {
@ -24,7 +24,7 @@ public class MuseKnowledgeAuditService {
private MuseKnowledgeRagflowCallMapper ragflowCallMapper;
@Transactional(rollbackFor = Exception.class)
public MuseKnowledgeRagflowCallDO recordRagflowCall(RagflowCallRecordReq req) {
public MuseKnowledgeRagflowCallDO recordRuntimeCall(RuntimeCallRecordReq req) {
String failureClass = req.getFailureClass() == null ? null : req.getFailureClass().name();
MuseKnowledgeRagflowCallDO call = new MuseKnowledgeRagflowCallDO();
call.setCorrelationId(req.getCorrelationId());
@ -66,14 +66,14 @@ public class MuseKnowledgeAuditService {
}
/**
* RAGFlow 调用审计创建请求
* 运行时调用审计创建请求
*/
@Data
@Builder
public static class RagflowCallRecordReq {
public static class RuntimeCallRecordReq {
private String correlationId;
private Integer attemptNo;
private RagFlowKnowledgeRuntimeClient.Operation operation;
private KnowledgeRuntimeClient.Operation operation;
private String commandId;
private String processingTaskId;
private Long actorUserId;
@ -85,14 +85,14 @@ public class MuseKnowledgeAuditService {
private String graphragOperation;
private String attributionStatus;
private String attributionBoundary;
private RagFlowKnowledgeRuntimeClient.Status status;
private KnowledgeRuntimeClient.Status status;
private String requestHash;
private String requestSummary;
private String responseSummary;
private Long durationMillis;
private String providerRequestId;
private Integer providerStatusCode;
private RagFlowKnowledgeRuntimeClient.FailureClass failureClass;
private KnowledgeRuntimeClient.FailureClass failureClass;
private String errorCode;
private String errorMessage;
private Boolean retryable;

View File

@ -6,7 +6,7 @@ import cn.iocoder.muse.framework.common.pojo.PageResult;
import cn.iocoder.muse.framework.common.util.json.JsonUtils;
import cn.iocoder.muse.framework.tenant.core.context.TenantContextHolder;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeFileFacade;
import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeRuntimeClient;
import cn.iocoder.muse.module.knowledge.controller.admin.muse.vo.AdminKnowledgeDocumentVO;
import cn.iocoder.muse.module.knowledge.controller.app.muse.vo.AppKnowledgeDocumentVO;
import cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeBaseDO;
@ -76,7 +76,7 @@ public class MuseKnowledgeDocumentService {
@Resource
private MuseKnowledgeProcessingTaskService processingTaskService;
@Resource
private RagFlowKnowledgeRuntimeClient ragFlowClient;
private KnowledgeRuntimeClient ragFlowClient;
@Resource
private PlatformTransactionManager transactionManager;
@ -263,15 +263,15 @@ public class MuseKnowledgeDocumentService {
}
UploadContext ragflowContext = datasetResolution.context();
RagFlowKnowledgeRuntimeClient.RuntimeResult uploadResult = ragFlowClient.uploadDocuments(
new RagFlowKnowledgeRuntimeClient.UploadDocumentsCommand(TenantContextHolder.getRequiredTenantId(),
KnowledgeRuntimeClient.RuntimeResult uploadResult = ragFlowClient.uploadDocuments(
new KnowledgeRuntimeClient.UploadDocumentsCommand(TenantContextHolder.getRequiredTenantId(),
ragflowContext.kb().getOwnerUserId(), ragflowContext.kb().getId(), ragflowContext.ragflowDatasetId(),
List.of(new RagFlowKnowledgeRuntimeClient.DocumentUpload(ragflowContext.document().getId(),
List.of(new KnowledgeRuntimeClient.DocumentUpload(ragflowContext.document().getId(),
ragflowContext.version().getId(), ragflowContext.materialized().fileRef(),
ragflowContext.materialized().fileName(), ragflowContext.materialized().contentType(),
ragflowContext.materialized().size(), ragflowContext.materialized().content())),
ragflowContext.envelope().correlationId(), ragflowContext.envelope().requestHash(), 1));
if (uploadResult.status() == RagFlowKnowledgeRuntimeClient.Status.FAILED) {
if (uploadResult.status() == KnowledgeRuntimeClient.Status.FAILED) {
auditRagflowCall(ragflowContext, uploadResult, uploadResult.externalId());
processingTaskService.markRagflowUploadFailed(ragflowContext.version().getId(), ragflowContext.task().getTaskId(),
uploadResult);
@ -281,13 +281,13 @@ public class MuseKnowledgeDocumentService {
recordUploadRecoverySnapshot(ragflowContext, uploadResult.externalId(), "processing");
auditRagflowCall(ragflowContext, uploadResult, uploadResult.externalId());
persistDocumentBinding(ragflowContext, uploadResult.externalId());
RagFlowKnowledgeRuntimeClient.RuntimeResult parseResult = ragFlowClient.startParseDocuments(
new RagFlowKnowledgeRuntimeClient.StartParseDocumentsCommand(TenantContextHolder.getRequiredTenantId(),
KnowledgeRuntimeClient.RuntimeResult parseResult = ragFlowClient.startParseDocuments(
new KnowledgeRuntimeClient.StartParseDocumentsCommand(TenantContextHolder.getRequiredTenantId(),
ragflowContext.kb().getOwnerUserId(), ragflowContext.kb().getId(), ragflowContext.ragflowDatasetId(),
List.of(uploadResult.externalId()), ragflowContext.task().getTaskId(),
ragflowContext.envelope().correlationId(), ragflowContext.envelope().requestHash(), 1));
auditRagflowCall(ragflowContext, parseResult, uploadResult.externalId());
if (parseResult.status() == RagFlowKnowledgeRuntimeClient.Status.FAILED) {
if (parseResult.status() == KnowledgeRuntimeClient.Status.FAILED) {
processingTaskService.markRagflowParseFailed(ragflowContext.version().getId(), ragflowContext.task().getTaskId(),
uploadResult.externalId(), parseResult);
return new UploadOutcome(ragflowContext, "failed");
@ -302,13 +302,13 @@ public class MuseKnowledgeDocumentService {
return new DatasetBindingResolution(context, null);
}
String datasetName = datasetName(context);
RagFlowKnowledgeRuntimeClient.RuntimeResult createResult = ragFlowClient.createDataset(
new RagFlowKnowledgeRuntimeClient.CreateDatasetCommand(TenantContextHolder.getRequiredTenantId(),
KnowledgeRuntimeClient.RuntimeResult createResult = ragFlowClient.createDataset(
new KnowledgeRuntimeClient.CreateDatasetCommand(TenantContextHolder.getRequiredTenantId(),
context.kb().getOwnerUserId(), context.kb().getId(), datasetName,
datasetConfig(context), context.envelope().correlationId(),
context.envelope().requestHash(), 1));
auditDatasetCreateCall(context, createResult, datasetName);
if (createResult.status() == RagFlowKnowledgeRuntimeClient.Status.FAILED) {
if (createResult.status() == KnowledgeRuntimeClient.Status.FAILED) {
return new DatasetBindingResolution(context, createResult);
}
MuseKnowledgeRagflowBindingDO binding = persistDatasetBinding(context, createResult.externalId());
@ -353,11 +353,11 @@ public class MuseKnowledgeDocumentService {
context.replay(), context.replayProcessingStatus());
}
private RagFlowKnowledgeRuntimeClient.RuntimeResult datasetBindingConflictResult(UploadContext context) {
return new RagFlowKnowledgeRuntimeClient.RuntimeResult(null, RagFlowKnowledgeRuntimeClient.Operation.CREATE_DATASET,
private KnowledgeRuntimeClient.RuntimeResult datasetBindingConflictResult(UploadContext context) {
return new KnowledgeRuntimeClient.RuntimeResult(null, KnowledgeRuntimeClient.Operation.CREATE_DATASET,
context.envelope().correlationId(), context.envelope().requestHash(), 1, 0L,
RagFlowKnowledgeRuntimeClient.Status.FAILED,
RagFlowKnowledgeRuntimeClient.FailureClass.CONFLICT, context.task().getTaskId(),
KnowledgeRuntimeClient.Status.FAILED,
KnowledgeRuntimeClient.FailureClass.CONFLICT, context.task().getTaskId(),
JsonUtils.toJsonString(Map.of("status", "failed", "failureClass", "CONFLICT",
"reason", "dataset binding unique conflict")), List.of(), Map.of());
}
@ -409,7 +409,7 @@ public class MuseKnowledgeDocumentService {
});
}
private void auditRagflowCall(UploadContext context, RagFlowKnowledgeRuntimeClient.RuntimeResult result,
private void auditRagflowCall(UploadContext context, KnowledgeRuntimeClient.RuntimeResult result,
String ragflowDocumentId) {
auditRagflowCall(context, result, context.ragflowDatasetId(), ragflowDocumentId,
Map.of("fileRef", context.materialized().fileRef(),
@ -418,17 +418,17 @@ public class MuseKnowledgeDocumentService {
"size", context.materialized().size()));
}
private void auditDatasetCreateCall(UploadContext context, RagFlowKnowledgeRuntimeClient.RuntimeResult result,
private void auditDatasetCreateCall(UploadContext context, KnowledgeRuntimeClient.RuntimeResult result,
String datasetName) {
auditRagflowCall(context, result, result.externalId(), null,
Map.of("source", "dataset_auto_created", "datasetName", datasetName,
"operation", context.envelope().operationId(), "commandId", context.envelope().commandId()));
}
private void auditRagflowCall(UploadContext context, RagFlowKnowledgeRuntimeClient.RuntimeResult result,
private void auditRagflowCall(UploadContext context, KnowledgeRuntimeClient.RuntimeResult result,
String ragflowDatasetId, String ragflowDocumentId,
Map<String, Object> requestSummary) {
auditService.recordRagflowCall(MuseKnowledgeAuditService.RagflowCallRecordReq.builder()
auditService.recordRuntimeCall(MuseKnowledgeAuditService.RuntimeCallRecordReq.builder()
// 同一上传命令会有 upload parse 两次外部调用审计 correlation 必须区分 operation避免唯一键冲突
.correlationId(context.envelope().correlationId() + ":" + result.operation().wireName())
.attemptNo(result.attempt())
@ -456,12 +456,12 @@ public class MuseKnowledgeDocumentService {
.build());
}
private boolean isRetryableRagflowFailure(RagFlowKnowledgeRuntimeClient.FailureClass failureClass) {
return failureClass == RagFlowKnowledgeRuntimeClient.FailureClass.RAGFLOW_UNAVAILABLE
|| failureClass == RagFlowKnowledgeRuntimeClient.FailureClass.TIMEOUT
|| failureClass == RagFlowKnowledgeRuntimeClient.FailureClass.RATE_LIMITED
|| failureClass == RagFlowKnowledgeRuntimeClient.FailureClass.CONFIG_MISSING
|| failureClass == RagFlowKnowledgeRuntimeClient.FailureClass.AUTH_MISSING;
private boolean isRetryableRagflowFailure(KnowledgeRuntimeClient.FailureClass failureClass) {
return failureClass == KnowledgeRuntimeClient.FailureClass.RAGFLOW_UNAVAILABLE
|| failureClass == KnowledgeRuntimeClient.FailureClass.TIMEOUT
|| failureClass == KnowledgeRuntimeClient.FailureClass.RATE_LIMITED
|| failureClass == KnowledgeRuntimeClient.FailureClass.CONFIG_MISSING
|| failureClass == KnowledgeRuntimeClient.FailureClass.AUTH_MISSING;
}
private DeleteOutcome deleteDocument(Long actorUserId, String apiVersion, Long kbId, Long documentId,
@ -1059,7 +1059,7 @@ public class MuseKnowledgeDocumentService {
}
private record DatasetBindingResolution(UploadContext context,
RagFlowKnowledgeRuntimeClient.RuntimeResult failureResult) {
KnowledgeRuntimeClient.RuntimeResult failureResult) {
}
private record DeleteOutcome(String sourceEventId, Integer affectedBindings) {

View File

@ -1,7 +1,7 @@
package cn.iocoder.muse.module.knowledge.application.muse;
import cn.iocoder.muse.framework.tenant.core.util.TenantUtils;
import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeRuntimeClient;
import cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeProcessingTaskDO;
import cn.iocoder.muse.module.knowledge.dal.mysql.muse.MuseKnowledgeProcessingTaskMapper;
import jakarta.annotation.Resource;
@ -34,7 +34,7 @@ public class MuseKnowledgeParseStatusPollWorker {
@Resource
private MuseKnowledgeProcessingTaskMapper processingTaskMapper;
@Resource
private RagFlowKnowledgeRuntimeClient ragFlowClient;
private KnowledgeRuntimeClient ragFlowClient;
@Resource
private MuseKnowledgeProcessingTaskService processingTaskService;
@ -46,7 +46,7 @@ public class MuseKnowledgeParseStatusPollWorker {
// 测试构造直接注入依赖与开关避免起 Spring 上下文
MuseKnowledgeParseStatusPollWorker(MuseKnowledgeProcessingTaskMapper processingTaskMapper,
RagFlowKnowledgeRuntimeClient ragFlowClient,
KnowledgeRuntimeClient ragFlowClient,
MuseKnowledgeProcessingTaskService processingTaskService,
boolean enabled) {
this.processingTaskMapper = processingTaskMapper;
@ -86,9 +86,9 @@ public class MuseKnowledgeParseStatusPollWorker {
}
private int pollOne(MuseKnowledgeProcessingTaskDO task) {
RagFlowKnowledgeRuntimeClient.RuntimeResult result;
KnowledgeRuntimeClient.RuntimeResult result;
try {
result = ragFlowClient.pollDocumentStatuses(new RagFlowKnowledgeRuntimeClient.PollDocumentStatusesCommand(
result = ragFlowClient.pollDocumentStatuses(new KnowledgeRuntimeClient.PollDocumentStatusesCommand(
task.getTenantId(), task.getOwnerUserId(), task.getKbId(), task.getRagflowDatasetId(),
List.of(task.getRagflowDocumentId()), task.getTaskId(),
pollCorrelationId(task), pollRequestHash(task), 1));
@ -97,7 +97,7 @@ public class MuseKnowledgeParseStatusPollWorker {
log.warn("Knowledge parse 轮询异常taskId={}, errorType={}", task.getTaskId(), ex.getClass().getSimpleName());
return 0;
}
if (result.status() == RagFlowKnowledgeRuntimeClient.Status.FAILED) {
if (result.status() == KnowledgeRuntimeClient.Status.FAILED) {
log.warn("Knowledge parse 轮询 RAGFlow 返回失败taskId={}, failureClass={}",
task.getTaskId(), result.failureClass());
return 0; // 轮询本身失败非文档失败 保持 parsing 重试
@ -105,7 +105,7 @@ public class MuseKnowledgeParseStatusPollWorker {
if (result.documentStatuses() == null || result.documentStatuses().isEmpty()) {
return 0;
}
RagFlowKnowledgeRuntimeClient.DocumentStatus status = result.documentStatuses().getFirst();
KnowledgeRuntimeClient.DocumentStatus status = result.documentStatuses().getFirst();
String museStatus = status.museStatus();
if (MUSE_STATUS_SEARCHABLE.equals(museStatus)) {
processingTaskService.markRagflowParseCompleted(task.getDocumentVersionId(), task.getTaskId(),

View File

@ -5,7 +5,7 @@ import cn.iocoder.muse.framework.common.pojo.PageParam;
import cn.iocoder.muse.framework.common.pojo.PageResult;
import cn.iocoder.muse.framework.common.util.json.JsonUtils;
import cn.iocoder.muse.framework.tenant.core.context.TenantContextHolder;
import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeRuntimeClient;
import cn.iocoder.muse.module.knowledge.controller.admin.muse.vo.AdminKnowledgeTaskVO;
import cn.iocoder.muse.module.knowledge.controller.app.muse.vo.AppKnowledgeTaskVO;
import cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeBaseDO;
@ -158,7 +158,7 @@ public class MuseKnowledgeProcessingTaskService {
@Transactional(propagation = Propagation.REQUIRES_NEW, rollbackFor = Exception.class)
public void markRagflowUploadFailed(Long documentVersionId, String taskId,
RagFlowKnowledgeRuntimeClient.RuntimeResult result) {
KnowledgeRuntimeClient.RuntimeResult result) {
MuseKnowledgeProcessingTaskDO task = processingTaskMapper.selectByTaskId(taskId);
if (task != null) {
task.setStatus("failed");
@ -178,7 +178,7 @@ public class MuseKnowledgeProcessingTaskService {
@Transactional(propagation = Propagation.REQUIRES_NEW, rollbackFor = Exception.class)
public void markRagflowParseAccepted(Long documentVersionId, String taskId, String ragflowDatasetId,
String ragflowDocumentId, RagFlowKnowledgeRuntimeClient.RuntimeResult result) {
String ragflowDocumentId, KnowledgeRuntimeClient.RuntimeResult result) {
MuseKnowledgeProcessingTaskDO task = processingTaskMapper.selectByTaskId(taskId);
if (task != null) {
task.setStatus("parsing");
@ -197,7 +197,7 @@ public class MuseKnowledgeProcessingTaskService {
@Transactional(propagation = Propagation.REQUIRES_NEW, rollbackFor = Exception.class)
public void markRagflowParseFailed(Long documentVersionId, String taskId, String ragflowDocumentId,
RagFlowKnowledgeRuntimeClient.RuntimeResult result) {
KnowledgeRuntimeClient.RuntimeResult result) {
MuseKnowledgeProcessingTaskDO task = processingTaskMapper.selectByTaskId(taskId);
if (task != null) {
task.setStatus("failed");
@ -220,7 +220,7 @@ public class MuseKnowledgeProcessingTaskService {
*/
@Transactional(propagation = Propagation.REQUIRES_NEW, rollbackFor = Exception.class)
public void markRagflowParseCompleted(Long documentVersionId, String taskId, String ragflowDocumentId,
Integer progress, RagFlowKnowledgeRuntimeClient.RuntimeResult result) {
Integer progress, KnowledgeRuntimeClient.RuntimeResult result) {
MuseKnowledgeProcessingTaskDO task = processingTaskMapper.selectByTaskId(taskId);
if (task != null) {
task.setStatus("completed");
@ -458,19 +458,19 @@ public class MuseKnowledgeProcessingTaskService {
documentVersionMapper.updateById(version);
}
private boolean isRetryable(RagFlowKnowledgeRuntimeClient.FailureClass failureClass) {
return failureClass == RagFlowKnowledgeRuntimeClient.FailureClass.RAGFLOW_UNAVAILABLE
|| failureClass == RagFlowKnowledgeRuntimeClient.FailureClass.TIMEOUT
|| failureClass == RagFlowKnowledgeRuntimeClient.FailureClass.RATE_LIMITED
|| failureClass == RagFlowKnowledgeRuntimeClient.FailureClass.CONFIG_MISSING
|| failureClass == RagFlowKnowledgeRuntimeClient.FailureClass.AUTH_MISSING;
private boolean isRetryable(KnowledgeRuntimeClient.FailureClass failureClass) {
return failureClass == KnowledgeRuntimeClient.FailureClass.RAGFLOW_UNAVAILABLE
|| failureClass == KnowledgeRuntimeClient.FailureClass.TIMEOUT
|| failureClass == KnowledgeRuntimeClient.FailureClass.RATE_LIMITED
|| failureClass == KnowledgeRuntimeClient.FailureClass.CONFIG_MISSING
|| failureClass == KnowledgeRuntimeClient.FailureClass.AUTH_MISSING;
}
private String failureName(RagFlowKnowledgeRuntimeClient.RuntimeResult result) {
private String failureName(KnowledgeRuntimeClient.RuntimeResult result) {
return result.failureClass() == null ? "RAGFLOW_FAILED" : result.failureClass().name();
}
private Map<String, Object> safeSummary(RagFlowKnowledgeRuntimeClient.RuntimeResult result) {
private Map<String, Object> safeSummary(KnowledgeRuntimeClient.RuntimeResult result) {
Map<String, Object> summary = new LinkedHashMap<>();
summary.put("operation", result.operation().wireName());
summary.put("status", result.status().name());

View File

@ -31,7 +31,7 @@ import java.util.UUID;
/**
* RAGFlow HTTP runtime client
*/
public class HttpRagFlowKnowledgeRuntimeClient implements RagFlowKnowledgeRuntimeClient {
public class HttpRagFlowKnowledgeRuntimeClient implements KnowledgeRuntimeClient {
private final String baseUrl;
private final String apiKey;
@ -64,19 +64,6 @@ public class HttpRagFlowKnowledgeRuntimeClient implements RagFlowKnowledgeRuntim
"POST", "/api/v1/datasets", body, "id", null);
}
@Override
public RuntimeResult updateDatasetConfig(UpdateDatasetConfigCommand command) {
return executeJson(Operation.UPDATE_DATASET_CONFIG, command.correlationId(), command.requestHash(),
command.attempt(), "PUT", "/api/v1/datasets/" + command.ragflowDatasetId(),
command.config(), "id", null);
}
@Override
public RuntimeResult listDatasets(ListDatasetsCommand command) {
return executeJson(Operation.LIST_DATASETS, command.correlationId(), command.requestHash(), command.attempt(),
"GET", "/api/v1/datasets", command.filter(), null, null);
}
@Override
public RuntimeResult uploadDocuments(UploadDocumentsCommand command) {
MultipartPayload multipartPayload = buildMultipartPayload(command.documents());
@ -139,13 +126,6 @@ public class HttpRagFlowKnowledgeRuntimeClient implements RagFlowKnowledgeRuntim
sanitizeSummary(summary), List.of(), summary);
}
@Override
public RuntimeResult listChunks(ListChunksCommand command) {
return executeJson(Operation.LIST_CHUNKS, command.correlationId(), command.requestHash(), command.attempt(),
"GET", "/api/v1/datasets/" + command.ragflowDatasetId() + "/documents/"
+ command.ragflowDocumentId() + "/chunks", command.filter(), null, null);
}
@Override
public RuntimeResult retrieveChunks(RetrieveChunksCommand command) {
Map<String, Object> body = new LinkedHashMap<>();
@ -166,38 +146,6 @@ public class HttpRagFlowKnowledgeRuntimeClient implements RagFlowKnowledgeRuntim
"POST", "/api/v1/retrieval", body, null, null);
}
@Override
public RuntimeResult runGraphRag(RunGraphRagCommand command) {
if (!graphRagAttributionReady) {
return failure(Operation.RUN_GRAPH_RAG, command.correlationId(), command.requestHash(), command.attempt(),
command.processingTaskId(), FailureClass.ATTRIBUTION_NOT_CONFIGURED, 0L,
"GraphRAG attribution is not configured");
}
return executeJson(Operation.RUN_GRAPH_RAG, command.correlationId(), command.requestHash(), command.attempt(),
"POST", "/api/v1/datasets/" + command.ragflowDatasetId() + "/run_graphrag",
Map.of(), "graphrag_task_id", command.processingTaskId());
}
@Override
public RuntimeResult traceGraphRag(TraceGraphRagCommand command) {
return executeJson(Operation.TRACE_GRAPH_RAG, command.correlationId(), command.requestHash(), command.attempt(),
"GET", "/api/v1/datasets/" + command.ragflowDatasetId() + "/trace_graphrag",
Map.of(), null, null);
}
@Override
public RuntimeResult getKnowledgeGraph(GetKnowledgeGraphCommand command) {
return executeJson(Operation.GET_KNOWLEDGE_GRAPH, command.correlationId(), command.requestHash(),
command.attempt(), "GET", "/api/v1/datasets/" + command.ragflowDatasetId() + "/knowledge_graph",
Map.of(), null, null);
}
@Override
public RuntimeResult health(HealthCommand command) {
return executeJson(Operation.HEALTH, command.correlationId(), command.requestHash(), command.attempt(),
"GET", "/v1/system/healthz", Map.of(), null, null);
}
protected HttpResponsePayload send(HttpRequestPayload request) throws IOException, InterruptedException {
HttpRequest.Builder builder = HttpRequest.newBuilder()
.uri(URI.create(baseUrl + request.path()))

View File

@ -4,47 +4,26 @@ import java.util.List;
import java.util.Map;
/**
* Knowledge 调用 RAGFlow 的运行时边界
* Knowledge 运行时端口文档处理与检索的外部运行时调用边界
*/
public interface RagFlowKnowledgeRuntimeClient {
public interface KnowledgeRuntimeClient {
RuntimeResult createDataset(CreateDatasetCommand command);
RuntimeResult updateDatasetConfig(UpdateDatasetConfigCommand command);
RuntimeResult listDatasets(ListDatasetsCommand command);
RuntimeResult uploadDocuments(UploadDocumentsCommand command);
RuntimeResult startParseDocuments(StartParseDocumentsCommand command);
RuntimeResult pollDocumentStatuses(PollDocumentStatusesCommand command);
RuntimeResult listChunks(ListChunksCommand command);
RuntimeResult retrieveChunks(RetrieveChunksCommand command);
RuntimeResult runGraphRag(RunGraphRagCommand command);
RuntimeResult traceGraphRag(TraceGraphRagCommand command);
RuntimeResult getKnowledgeGraph(GetKnowledgeGraphCommand command);
RuntimeResult health(HealthCommand command);
enum Operation {
CREATE_DATASET("createDataset"),
UPDATE_DATASET_CONFIG("updateDatasetConfig"),
LIST_DATASETS("listDatasets"),
UPLOAD_DOCUMENTS("uploadDocuments"),
START_PARSE_DOCUMENTS("startParseDocuments"),
POLL_DOCUMENT_STATUSES("pollDocumentStatuses"),
LIST_CHUNKS("listChunks"),
RETRIEVE_CHUNKS("retrieveChunks"),
RUN_GRAPH_RAG("runGraphRag"),
TRACE_GRAPH_RAG("traceGraphRag"),
GET_KNOWLEDGE_GRAPH("getKnowledgeGraph"),
HEALTH("health");
RETRIEVE_CHUNKS("retrieveChunks");
private final String wireName;
@ -126,15 +105,6 @@ public interface RagFlowKnowledgeRuntimeClient {
String requestHash, int attempt) {
}
record UpdateDatasetConfigCommand(Long tenantId, Long ownerUserId, Long kbId, String ragflowDatasetId,
Map<String, Object> config, String correlationId,
String requestHash, int attempt) {
}
record ListDatasetsCommand(Long tenantId, Long ownerUserId, Map<String, Object> filter,
String correlationId, String requestHash, int attempt) {
}
record UploadDocumentsCommand(Long tenantId, Long ownerUserId, Long kbId, String ragflowDatasetId,
List<DocumentUpload> documents, String correlationId,
String requestHash, int attempt) {
@ -150,31 +120,10 @@ public interface RagFlowKnowledgeRuntimeClient {
String correlationId, String requestHash, int attempt) {
}
record ListChunksCommand(Long tenantId, Long ownerUserId, Long kbId, String ragflowDatasetId,
String ragflowDocumentId, Map<String, Object> filter,
String correlationId, String requestHash, int attempt) {
}
record RetrieveChunksCommand(Long tenantId, Long ownerUserId, Long kbId, List<String> ragflowDatasetIds,
List<String> ragflowDocumentIds, String question, Integer topK,
Double threshold, Map<String, Object> metadataCondition,
String correlationId, String requestHash, int attempt) {
}
record RunGraphRagCommand(Long tenantId, Long ownerUserId, Long kbId, String ragflowDatasetId,
String processingTaskId, String correlationId, String requestHash,
int attempt) {
}
record TraceGraphRagCommand(Long tenantId, Long ownerUserId, Long kbId, String ragflowDatasetId,
String correlationId, String requestHash, int attempt) {
}
record GetKnowledgeGraphCommand(Long tenantId, Long ownerUserId, Long kbId, String ragflowDatasetId,
String correlationId, String requestHash, int attempt) {
}
record HealthCommand(String correlationId, String requestHash, int attempt) {
}
}

View File

@ -19,8 +19,8 @@ public class RagFlowKnowledgeRuntimeClientConfiguration {
private static final int DEFAULT_RETRY_BUDGET = 0;
@Bean
@ConditionalOnMissingBean(RagFlowKnowledgeRuntimeClient.class)
public RagFlowKnowledgeRuntimeClient ragFlowKnowledgeRuntimeClient(Environment environment) {
@ConditionalOnMissingBean(KnowledgeRuntimeClient.class)
public KnowledgeRuntimeClient ragFlowKnowledgeRuntimeClient(Environment environment) {
String baseUrl = environment.getProperty(PROPERTY_PREFIX + "base-url");
String apiKey = environment.getProperty(PROPERTY_PREFIX + "api-key");
if (!StringUtils.hasText(baseUrl) || !StringUtils.hasText(apiKey)) {

View File

@ -6,23 +6,13 @@ import java.util.Map;
/**
* RAGFlow 未配置时的 fail closed client
*/
public class UnavailableRagFlowKnowledgeRuntimeClient implements RagFlowKnowledgeRuntimeClient {
public class UnavailableRagFlowKnowledgeRuntimeClient implements KnowledgeRuntimeClient {
@Override
public RuntimeResult createDataset(CreateDatasetCommand command) {
return unavailable(Operation.CREATE_DATASET, command.correlationId(), command.requestHash(), command.attempt(), null);
}
@Override
public RuntimeResult updateDatasetConfig(UpdateDatasetConfigCommand command) {
return unavailable(Operation.UPDATE_DATASET_CONFIG, command.correlationId(), command.requestHash(), command.attempt(), null);
}
@Override
public RuntimeResult listDatasets(ListDatasetsCommand command) {
return unavailable(Operation.LIST_DATASETS, command.correlationId(), command.requestHash(), command.attempt(), null);
}
@Override
public RuntimeResult uploadDocuments(UploadDocumentsCommand command) {
return unavailable(Operation.UPLOAD_DOCUMENTS, command.correlationId(), command.requestHash(), command.attempt(), null);
@ -40,37 +30,11 @@ public class UnavailableRagFlowKnowledgeRuntimeClient implements RagFlowKnowledg
command.attempt(), command.processingTaskId());
}
@Override
public RuntimeResult listChunks(ListChunksCommand command) {
return unavailable(Operation.LIST_CHUNKS, command.correlationId(), command.requestHash(), command.attempt(), null);
}
@Override
public RuntimeResult retrieveChunks(RetrieveChunksCommand command) {
return unavailable(Operation.RETRIEVE_CHUNKS, command.correlationId(), command.requestHash(), command.attempt(), null);
}
@Override
public RuntimeResult runGraphRag(RunGraphRagCommand command) {
return unavailable(Operation.RUN_GRAPH_RAG, command.correlationId(), command.requestHash(),
command.attempt(), command.processingTaskId());
}
@Override
public RuntimeResult traceGraphRag(TraceGraphRagCommand command) {
return unavailable(Operation.TRACE_GRAPH_RAG, command.correlationId(), command.requestHash(), command.attempt(), null);
}
@Override
public RuntimeResult getKnowledgeGraph(GetKnowledgeGraphCommand command) {
return unavailable(Operation.GET_KNOWLEDGE_GRAPH, command.correlationId(), command.requestHash(), command.attempt(), null);
}
@Override
public RuntimeResult health(HealthCommand command) {
return unavailable(Operation.HEALTH, command.correlationId(), command.requestHash(), command.attempt(), null);
}
private RuntimeResult unavailable(Operation operation, String correlationId, String requestHash,
int attempt, String processingTaskId) {
return new RuntimeResult(null, operation, correlationId, requestHash, attempt, 0L,

View File

@ -4,12 +4,12 @@ import cn.iocoder.muse.framework.common.util.json.JsonUtils;
import cn.iocoder.muse.framework.test.core.ut.BaseMockitoUnitTest;
import cn.iocoder.muse.module.knowledge.api.MuseKnowledgeRetrievalApi.RetrievalRequest;
import cn.iocoder.muse.module.knowledge.api.MuseKnowledgeRetrievalApi.RetrievalResult;
import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient;
import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient.FailureClass;
import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient.Operation;
import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient.RetrieveChunksCommand;
import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient.RuntimeResult;
import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient.Status;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeRuntimeClient;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeRuntimeClient.FailureClass;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeRuntimeClient.Operation;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeRuntimeClient.RetrieveChunksCommand;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeRuntimeClient.RuntimeResult;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeRuntimeClient.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;
@ -49,7 +49,7 @@ class MuseKnowledgeRetrievalApiImplTest extends BaseMockitoUnitTest {
@Mock
private MuseKnowledgeRagflowBindingMapper ragflowBindingMapper;
@Mock
private RagFlowKnowledgeRuntimeClient ragFlowClient;
private KnowledgeRuntimeClient ragFlowClient;
private RetrievalRequest req() {
return new RetrievalRequest(100L, 2001L, 4001L, "如何写好开头", 5, "corr-1");

View File

@ -3,11 +3,11 @@ package cn.iocoder.muse.module.knowledge.application.muse;
import cn.iocoder.muse.framework.test.core.ut.BaseMockitoUnitTest;
import cn.iocoder.muse.framework.tenant.core.context.TenantContextHolder;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeFileFacade;
import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient;
import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient.FailureClass;
import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient.Operation;
import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient.RuntimeResult;
import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient.Status;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeRuntimeClient;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeRuntimeClient.FailureClass;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeRuntimeClient.Operation;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeRuntimeClient.RuntimeResult;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeRuntimeClient.Status;
import cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeDocumentDO;
import cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeDocumentVersionDO;
import cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeRagflowBindingDO;
@ -64,7 +64,7 @@ class KnowledgeMarketForkServiceTest extends BaseMockitoUnitTest {
@Mock
private KnowledgeFileFacade fileFacade;
@Mock
private RagFlowKnowledgeRuntimeClient ragFlowClient;
private KnowledgeRuntimeClient ragFlowClient;
@Mock
private MuseKnowledgeAuditService auditService;
@Mock
@ -125,7 +125,7 @@ class KnowledgeMarketForkServiceTest extends BaseMockitoUnitTest {
ASSET_ID.equals(req.assetId()) && "ready".equals(req.forkStatus())
&& "rag-fork-ds".equals(req.publicForkDatasetId())));
// fork 的每次 RAGFlow 调用都审计create + 3*upload + 3*startParse = 7
verify(auditService, atLeastOnce()).recordRagflowCall(argThat(req ->
verify(auditService, atLeastOnce()).recordRuntimeCall(argThat(req ->
"market_public_fork".equals(req.getAttributionStatus())));
}

View File

@ -2,7 +2,7 @@ package cn.iocoder.muse.module.knowledge.application.muse;
import cn.iocoder.muse.framework.test.core.ut.BaseMockitoUnitTest;
import cn.iocoder.muse.framework.tenant.core.context.TenantContextHolder;
import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeRuntimeClient;
import cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeRagflowCallDO;
import cn.iocoder.muse.module.knowledge.dal.mysql.muse.MuseKnowledgeRagflowCallMapper;
import org.junit.jupiter.api.AfterEach;
@ -36,17 +36,17 @@ class MuseKnowledgeAuditServiceTest extends BaseMockitoUnitTest {
doReturn(1).when(ragflowCallMapper).insert(
org.mockito.ArgumentMatchers.<MuseKnowledgeRagflowCallDO>argThat(inserted -> true));
MuseKnowledgeRagflowCallDO call = auditService.recordRagflowCall(MuseKnowledgeAuditService.RagflowCallRecordReq.builder()
MuseKnowledgeRagflowCallDO call = auditService.recordRuntimeCall(MuseKnowledgeAuditService.RuntimeCallRecordReq.builder()
.correlationId("corr-rag-1")
.attemptNo(1)
.operation(RagFlowKnowledgeRuntimeClient.Operation.CREATE_DATASET)
.operation(KnowledgeRuntimeClient.Operation.CREATE_DATASET)
.commandId("cmd-1")
.requestHash("hash-1")
.actorUserId(2001L)
.ownerUserId(2001L)
.kbId(4001L)
.status(RagFlowKnowledgeRuntimeClient.Status.FAILED)
.failureClass(RagFlowKnowledgeRuntimeClient.FailureClass.AUTH_FAILED)
.status(KnowledgeRuntimeClient.Status.FAILED)
.failureClass(KnowledgeRuntimeClient.FailureClass.AUTH_FAILED)
.requestSummary("""
{"apiKey":"rk-test-secret","nested":{"token":"t-123","secret":"s-456"},"Authorization":"Bearer abcdefgh"}
""")
@ -76,17 +76,17 @@ class MuseKnowledgeAuditServiceTest extends BaseMockitoUnitTest {
doReturn(1).when(ragflowCallMapper).insert(
org.mockito.ArgumentMatchers.<MuseKnowledgeRagflowCallDO>argThat(inserted -> true));
MuseKnowledgeRagflowCallDO call = auditService.recordRagflowCall(MuseKnowledgeAuditService.RagflowCallRecordReq.builder()
MuseKnowledgeRagflowCallDO call = auditService.recordRuntimeCall(MuseKnowledgeAuditService.RuntimeCallRecordReq.builder()
.correlationId("corr-rag-raw-json")
.attemptNo(1)
.operation(RagFlowKnowledgeRuntimeClient.Operation.RETRIEVE_CHUNKS)
.operation(KnowledgeRuntimeClient.Operation.RETRIEVE_CHUNKS)
.requestSummary("""
{"question":"private question","prompt":"private prompt","entryContent":"private entry","fileContent":"private file"}
""")
.responseSummary("""
{"responseBody":{"content":"private chunk","text":"private text","rawAnswer":"private raw","response":"private answer"}}
""")
.status(RagFlowKnowledgeRuntimeClient.Status.SUCCEEDED)
.status(KnowledgeRuntimeClient.Status.SUCCEEDED)
.durationMillis(8L)
.build());
@ -108,13 +108,13 @@ class MuseKnowledgeAuditServiceTest extends BaseMockitoUnitTest {
doReturn(1).when(ragflowCallMapper).insert(
org.mockito.ArgumentMatchers.<MuseKnowledgeRagflowCallDO>argThat(inserted -> true));
MuseKnowledgeRagflowCallDO call = auditService.recordRagflowCall(MuseKnowledgeAuditService.RagflowCallRecordReq.builder()
MuseKnowledgeRagflowCallDO call = auditService.recordRuntimeCall(MuseKnowledgeAuditService.RuntimeCallRecordReq.builder()
.correlationId("corr-rag-plain")
.attemptNo(1)
.operation(RagFlowKnowledgeRuntimeClient.Operation.RETRIEVE_CHUNKS)
.operation(KnowledgeRuntimeClient.Operation.RETRIEVE_CHUNKS)
.requestSummary("private plain text retrieval question without token markers")
.responseSummary("private raw response body without token markers")
.status(RagFlowKnowledgeRuntimeClient.Status.FAILED)
.status(KnowledgeRuntimeClient.Status.FAILED)
.durationMillis(9L)
.build());

View File

@ -5,7 +5,7 @@ import cn.iocoder.muse.framework.common.pojo.PageResult;
import cn.iocoder.muse.framework.test.core.ut.BaseMockitoUnitTest;
import cn.iocoder.muse.framework.tenant.core.context.TenantContextHolder;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeFileFacade;
import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeRuntimeClient;
import cn.iocoder.muse.module.knowledge.controller.admin.muse.vo.AdminKnowledgeDocumentVO;
import cn.iocoder.muse.module.knowledge.controller.app.muse.vo.AppKnowledgeDocumentVO;
import cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeBaseDO;
@ -80,7 +80,7 @@ class MuseKnowledgeDocumentServiceTest extends BaseMockitoUnitTest {
@Mock
private MuseKnowledgeMaterializationService materializationService;
@Mock
private RagFlowKnowledgeRuntimeClient ragFlowClient;
private KnowledgeRuntimeClient ragFlowClient;
@InjectMocks
private MuseKnowledgeProcessingTaskService processingTaskService;
@ -323,15 +323,15 @@ class MuseKnowledgeDocumentServiceTest extends BaseMockitoUnitTest {
});
when(processingTaskMapper.selectByTaskId("knowledge-processing-6001")).thenReturn(processingTask(6001L));
when(processingTaskMapper.insert(any(MuseKnowledgeProcessingTaskDO.class))).thenAnswer(invocation -> 1);
when(ragFlowClient.uploadDocuments(any())).thenReturn(new RagFlowKnowledgeRuntimeClient.RuntimeResult(
"rag-doc-1", RagFlowKnowledgeRuntimeClient.Operation.UPLOAD_DOCUMENTS,
when(ragFlowClient.uploadDocuments(any())).thenReturn(new KnowledgeRuntimeClient.RuntimeResult(
"rag-doc-1", KnowledgeRuntimeClient.Operation.UPLOAD_DOCUMENTS,
"uploadGlobalKBDocument:cmd-upload-file-success", "hash-file-success", 1, 12L,
RagFlowKnowledgeRuntimeClient.Status.SUCCEEDED, null, null,
KnowledgeRuntimeClient.Status.SUCCEEDED, null, null,
"{\"status\":\"succeeded\"}", List.of(), java.util.Map.of()));
when(ragFlowClient.startParseDocuments(any())).thenReturn(new RagFlowKnowledgeRuntimeClient.RuntimeResult(
null, RagFlowKnowledgeRuntimeClient.Operation.START_PARSE_DOCUMENTS,
when(ragFlowClient.startParseDocuments(any())).thenReturn(new KnowledgeRuntimeClient.RuntimeResult(
null, KnowledgeRuntimeClient.Operation.START_PARSE_DOCUMENTS,
"uploadGlobalKBDocument:cmd-upload-file-success", "hash-file-success", 1, 8L,
RagFlowKnowledgeRuntimeClient.Status.ACCEPTED, null, "knowledge-processing-6001",
KnowledgeRuntimeClient.Status.ACCEPTED, null, "knowledge-processing-6001",
"{\"status\":\"accepted\"}", List.of(), java.util.Map.of()));
AdminKnowledgeDocumentVO.UploadReqVO reqVO = new AdminKnowledgeDocumentVO.UploadReqVO();
@ -360,22 +360,22 @@ class MuseKnowledgeDocumentServiceTest extends BaseMockitoUnitTest {
"rag-dataset-1".equals(binding.getRagflowDatasetId())
&& "rag-doc-1".equals(binding.getRagflowDocumentId())
&& Long.valueOf(6001L).equals(binding.getDocumentVersionId())));
verify(auditService).recordRagflowCall(org.mockito.ArgumentMatchers.argThat(call ->
verify(auditService).recordRuntimeCall(org.mockito.ArgumentMatchers.argThat(call ->
"uploadGlobalKBDocument:cmd-upload-file-success:uploadDocuments".equals(call.getCorrelationId())
&& RagFlowKnowledgeRuntimeClient.Operation.UPLOAD_DOCUMENTS.equals(call.getOperation())
&& KnowledgeRuntimeClient.Operation.UPLOAD_DOCUMENTS.equals(call.getOperation())
&& "cmd-upload-file-success".equals(call.getCommandId())
&& "knowledge-processing-6001".equals(call.getProcessingTaskId())
&& "rag-dataset-1".equals(call.getRagflowDatasetId())
&& "rag-doc-1".equals(call.getRagflowDocumentId())
&& RagFlowKnowledgeRuntimeClient.Status.SUCCEEDED.equals(call.getStatus())));
verify(auditService).recordRagflowCall(org.mockito.ArgumentMatchers.argThat(call ->
&& KnowledgeRuntimeClient.Status.SUCCEEDED.equals(call.getStatus())));
verify(auditService).recordRuntimeCall(org.mockito.ArgumentMatchers.argThat(call ->
"uploadGlobalKBDocument:cmd-upload-file-success:startParseDocuments".equals(call.getCorrelationId())
&& RagFlowKnowledgeRuntimeClient.Operation.START_PARSE_DOCUMENTS.equals(call.getOperation())
&& KnowledgeRuntimeClient.Operation.START_PARSE_DOCUMENTS.equals(call.getOperation())
&& "cmd-upload-file-success".equals(call.getCommandId())
&& "knowledge-processing-6001".equals(call.getProcessingTaskId())
&& "rag-dataset-1".equals(call.getRagflowDatasetId())
&& "rag-doc-1".equals(call.getRagflowDocumentId())
&& RagFlowKnowledgeRuntimeClient.Status.ACCEPTED.equals(call.getStatus())));
&& KnowledgeRuntimeClient.Status.ACCEPTED.equals(call.getStatus())));
verify(processingTaskMapper).updateById(org.mockito.ArgumentMatchers.<MuseKnowledgeProcessingTaskDO>argThat(task ->
"parsing".equals(task.getStatus())
&& "rag-dataset-1".equals(task.getRagflowDatasetId())
@ -407,20 +407,20 @@ class MuseKnowledgeDocumentServiceTest extends BaseMockitoUnitTest {
});
when(processingTaskMapper.selectByTaskId("knowledge-processing-6001")).thenReturn(processingTask(6001L));
when(processingTaskMapper.insert(any(MuseKnowledgeProcessingTaskDO.class))).thenAnswer(invocation -> 1);
when(ragFlowClient.createDataset(any())).thenReturn(new RagFlowKnowledgeRuntimeClient.RuntimeResult(
"rag-dataset-auto-1", RagFlowKnowledgeRuntimeClient.Operation.CREATE_DATASET,
when(ragFlowClient.createDataset(any())).thenReturn(new KnowledgeRuntimeClient.RuntimeResult(
"rag-dataset-auto-1", KnowledgeRuntimeClient.Operation.CREATE_DATASET,
"uploadGlobalKBDocument:cmd-upload-auto-dataset", "hash-auto-dataset", 1, 10L,
RagFlowKnowledgeRuntimeClient.Status.SUCCEEDED, null, null,
KnowledgeRuntimeClient.Status.SUCCEEDED, null, null,
"{\"status\":\"succeeded\"}", List.of(), java.util.Map.of()));
when(ragFlowClient.uploadDocuments(any())).thenReturn(new RagFlowKnowledgeRuntimeClient.RuntimeResult(
"rag-doc-auto-1", RagFlowKnowledgeRuntimeClient.Operation.UPLOAD_DOCUMENTS,
when(ragFlowClient.uploadDocuments(any())).thenReturn(new KnowledgeRuntimeClient.RuntimeResult(
"rag-doc-auto-1", KnowledgeRuntimeClient.Operation.UPLOAD_DOCUMENTS,
"uploadGlobalKBDocument:cmd-upload-auto-dataset", "hash-auto-dataset", 1, 12L,
RagFlowKnowledgeRuntimeClient.Status.SUCCEEDED, null, null,
KnowledgeRuntimeClient.Status.SUCCEEDED, null, null,
"{\"status\":\"succeeded\"}", List.of(), java.util.Map.of()));
when(ragFlowClient.startParseDocuments(any())).thenReturn(new RagFlowKnowledgeRuntimeClient.RuntimeResult(
null, RagFlowKnowledgeRuntimeClient.Operation.START_PARSE_DOCUMENTS,
when(ragFlowClient.startParseDocuments(any())).thenReturn(new KnowledgeRuntimeClient.RuntimeResult(
null, KnowledgeRuntimeClient.Operation.START_PARSE_DOCUMENTS,
"uploadGlobalKBDocument:cmd-upload-auto-dataset", "hash-auto-dataset", 1, 8L,
RagFlowKnowledgeRuntimeClient.Status.ACCEPTED, null, "knowledge-processing-6001",
KnowledgeRuntimeClient.Status.ACCEPTED, null, "knowledge-processing-6001",
"{\"status\":\"accepted\"}", List.of(), java.util.Map.of()));
AdminKnowledgeDocumentVO.UploadReqVO reqVO = new AdminKnowledgeDocumentVO.UploadReqVO();
@ -458,14 +458,14 @@ class MuseKnowledgeDocumentServiceTest extends BaseMockitoUnitTest {
verify(ragFlowClient).startParseDocuments(argThat(command ->
"rag-dataset-auto-1".equals(command.ragflowDatasetId())
&& command.ragflowDocumentIds().contains("rag-doc-auto-1")));
verify(auditService).recordRagflowCall(org.mockito.ArgumentMatchers.argThat(call ->
verify(auditService).recordRuntimeCall(org.mockito.ArgumentMatchers.argThat(call ->
"uploadGlobalKBDocument:cmd-upload-auto-dataset:createDataset".equals(call.getCorrelationId())
&& RagFlowKnowledgeRuntimeClient.Operation.CREATE_DATASET.equals(call.getOperation())
&& KnowledgeRuntimeClient.Operation.CREATE_DATASET.equals(call.getOperation())
&& "cmd-upload-auto-dataset".equals(call.getCommandId())
&& "knowledge-processing-6001".equals(call.getProcessingTaskId())
&& "rag-dataset-auto-1".equals(call.getRagflowDatasetId())
&& call.getRagflowDocumentId() == null
&& RagFlowKnowledgeRuntimeClient.Status.SUCCEEDED.equals(call.getStatus())));
&& KnowledgeRuntimeClient.Status.SUCCEEDED.equals(call.getStatus())));
}
@Test
@ -492,11 +492,11 @@ class MuseKnowledgeDocumentServiceTest extends BaseMockitoUnitTest {
inserted.setId(6001L);
return 1;
});
when(ragFlowClient.createDataset(any())).thenReturn(new RagFlowKnowledgeRuntimeClient.RuntimeResult(
null, RagFlowKnowledgeRuntimeClient.Operation.CREATE_DATASET,
when(ragFlowClient.createDataset(any())).thenReturn(new KnowledgeRuntimeClient.RuntimeResult(
null, KnowledgeRuntimeClient.Operation.CREATE_DATASET,
"uploadGlobalKBDocument:cmd-upload-auto-dataset-failed", "hash-auto-dataset-failed", 1, 10L,
RagFlowKnowledgeRuntimeClient.Status.FAILED,
RagFlowKnowledgeRuntimeClient.FailureClass.RAGFLOW_UNAVAILABLE, null,
KnowledgeRuntimeClient.Status.FAILED,
KnowledgeRuntimeClient.FailureClass.RAGFLOW_UNAVAILABLE, null,
"{\"status\":\"failed\"}", List.of(), java.util.Map.of()));
AdminKnowledgeDocumentVO.UploadReqVO reqVO = new AdminKnowledgeDocumentVO.UploadReqVO();
@ -517,11 +517,11 @@ class MuseKnowledgeDocumentServiceTest extends BaseMockitoUnitTest {
verify(documentVersionMapper).updateById(org.mockito.ArgumentMatchers.<MuseKnowledgeDocumentVersionDO>argThat(version ->
"failed".equals(version.getProcessingStatus())
&& "RAGFLOW_UNAVAILABLE".equals(version.getChangeReason())));
verify(auditService).recordRagflowCall(org.mockito.ArgumentMatchers.argThat(call ->
verify(auditService).recordRuntimeCall(org.mockito.ArgumentMatchers.argThat(call ->
"uploadGlobalKBDocument:cmd-upload-auto-dataset-failed:createDataset".equals(call.getCorrelationId())
&& RagFlowKnowledgeRuntimeClient.Operation.CREATE_DATASET.equals(call.getOperation())
&& RagFlowKnowledgeRuntimeClient.Status.FAILED.equals(call.getStatus())
&& RagFlowKnowledgeRuntimeClient.FailureClass.RAGFLOW_UNAVAILABLE.equals(call.getFailureClass())));
&& KnowledgeRuntimeClient.Operation.CREATE_DATASET.equals(call.getOperation())
&& KnowledgeRuntimeClient.Status.FAILED.equals(call.getStatus())
&& KnowledgeRuntimeClient.FailureClass.RAGFLOW_UNAVAILABLE.equals(call.getFailureClass())));
}
@Test
@ -548,11 +548,11 @@ class MuseKnowledgeDocumentServiceTest extends BaseMockitoUnitTest {
inserted.setId(6001L);
return 1;
});
when(ragFlowClient.uploadDocuments(any())).thenReturn(new RagFlowKnowledgeRuntimeClient.RuntimeResult(
null, RagFlowKnowledgeRuntimeClient.Operation.UPLOAD_DOCUMENTS,
when(ragFlowClient.uploadDocuments(any())).thenReturn(new KnowledgeRuntimeClient.RuntimeResult(
null, KnowledgeRuntimeClient.Operation.UPLOAD_DOCUMENTS,
"uploadGlobalKBDocument:cmd-upload-ragflow-failed", "hash-ragflow-failed", 1, 12L,
RagFlowKnowledgeRuntimeClient.Status.FAILED,
RagFlowKnowledgeRuntimeClient.FailureClass.RAGFLOW_UNAVAILABLE, null,
KnowledgeRuntimeClient.Status.FAILED,
KnowledgeRuntimeClient.FailureClass.RAGFLOW_UNAVAILABLE, null,
"{\"status\":\"failed\"}", List.of(), java.util.Map.of()));
AdminKnowledgeDocumentVO.UploadReqVO reqVO = new AdminKnowledgeDocumentVO.UploadReqVO();
@ -765,10 +765,10 @@ class MuseKnowledgeDocumentServiceTest extends BaseMockitoUnitTest {
return 1;
});
when(processingTaskMapper.insert(any(MuseKnowledgeProcessingTaskDO.class))).thenReturn(1);
when(ragFlowClient.uploadDocuments(any())).thenReturn(new RagFlowKnowledgeRuntimeClient.RuntimeResult(
"rag-doc-orphan-1", RagFlowKnowledgeRuntimeClient.Operation.UPLOAD_DOCUMENTS,
when(ragFlowClient.uploadDocuments(any())).thenReturn(new KnowledgeRuntimeClient.RuntimeResult(
"rag-doc-orphan-1", KnowledgeRuntimeClient.Operation.UPLOAD_DOCUMENTS,
"uploadGlobalKBDocument:cmd-upload-binding-failed", "hash-upload-binding-failed", 1, 12L,
RagFlowKnowledgeRuntimeClient.Status.SUCCEEDED, null, null,
KnowledgeRuntimeClient.Status.SUCCEEDED, null, null,
"{\"status\":\"succeeded\"}", List.of(), java.util.Map.of()));
when(ragflowBindingMapper.insert(org.mockito.ArgumentMatchers.<MuseKnowledgeRagflowBindingDO>argThat(binding ->
Long.valueOf(6001L).equals(binding.getDocumentVersionId()))))

View File

@ -169,7 +169,7 @@ class MuseKnowledgeGraphQueryServiceTest extends BaseMockitoUnitTest {
assertFalse(source.contains("runGraphRag"));
assertFalse(source.contains("run_graphrag"));
assertFalse(source.contains("RagFlowKnowledgeRuntimeClient"),
assertFalse(source.contains("KnowledgeRuntimeClient"),
"Task 7 graph 只消费 canonical entity/relation不应调用 RAGFlow runtime client");
}

View File

@ -1,11 +1,11 @@
package cn.iocoder.muse.module.knowledge.application.muse;
import cn.iocoder.muse.framework.test.core.ut.BaseMockitoUnitTest;
import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient;
import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient.DocumentStatus;
import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient.Operation;
import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient.RuntimeResult;
import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient.Status;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeRuntimeClient;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeRuntimeClient.DocumentStatus;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeRuntimeClient.Operation;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeRuntimeClient.RuntimeResult;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeRuntimeClient.Status;
import cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeProcessingTaskDO;
import cn.iocoder.muse.module.knowledge.dal.mysql.muse.MuseKnowledgeProcessingTaskMapper;
import org.junit.jupiter.api.Test;
@ -31,7 +31,7 @@ class MuseKnowledgeParseStatusPollWorkerTest extends BaseMockitoUnitTest {
@Mock
private MuseKnowledgeProcessingTaskMapper processingTaskMapper;
@Mock
private RagFlowKnowledgeRuntimeClient ragFlowClient;
private KnowledgeRuntimeClient ragFlowClient;
@Mock
private MuseKnowledgeProcessingTaskService processingTaskService;

View File

@ -17,87 +17,24 @@ import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* RAGFlow runtime client 基础行为测试
* Knowledge 运行时 client 基础行为测试
*/
class RagFlowKnowledgeRuntimeClientTest {
class KnowledgeRuntimeClientTest {
@Test
void unavailableClientShouldFailClosedWhenConfigMissing() {
RagFlowKnowledgeRuntimeClient client = new UnavailableRagFlowKnowledgeRuntimeClient();
KnowledgeRuntimeClient client = new UnavailableRagFlowKnowledgeRuntimeClient();
RagFlowKnowledgeRuntimeClient.RuntimeResult result = client.createDataset(
new RagFlowKnowledgeRuntimeClient.CreateDatasetCommand(100L, 2001L, 4001L,
KnowledgeRuntimeClient.RuntimeResult result = client.createDataset(
new KnowledgeRuntimeClient.CreateDatasetCommand(100L, 2001L, 4001L,
"KB", Map.of(), "corr-1", "hash-1", 1));
assertEquals(RagFlowKnowledgeRuntimeClient.Status.FAILED, result.status());
assertEquals(RagFlowKnowledgeRuntimeClient.FailureClass.CONFIG_MISSING, result.failureClass());
assertEquals(KnowledgeRuntimeClient.Status.FAILED, result.status());
assertEquals(KnowledgeRuntimeClient.FailureClass.CONFIG_MISSING, result.failureClass());
assertEquals("createDataset", result.operation().wireName());
assertFalse(result.redactedSummary().isBlank());
}
@Test
void graphRagShouldReturnAttributionFailureBeforeCallingRagflowWhenAttributionMissing() {
HttpRagFlowKnowledgeRuntimeClient client = new HttpRagFlowKnowledgeRuntimeClient(
"https://ragflow.example", "configured-api-key", Duration.ofSeconds(1), 0, false);
RagFlowKnowledgeRuntimeClient.RuntimeResult result = client.runGraphRag(
new RagFlowKnowledgeRuntimeClient.RunGraphRagCommand(100L, 2001L, 4001L,
"dataset-1", "task-1", "corr-graph-1", "hash-graph-1", 1));
assertEquals(RagFlowKnowledgeRuntimeClient.Status.FAILED, result.status());
assertEquals(RagFlowKnowledgeRuntimeClient.FailureClass.ATTRIBUTION_NOT_CONFIGURED, result.failureClass());
assertNull(result.externalId());
assertEquals("corr-graph-1", result.correlationId());
}
@Test
void runGraphRagShouldReadGraphragTaskIdAsExternalId() {
HttpRagFlowKnowledgeRuntimeClient client = new StubHttpRagFlowKnowledgeRuntimeClient(
"https://ragflow.example", "configured-api-key",
"{\"code\":0,\"data\":{\"graphrag_task_id\":\"gr-1\"}}");
RagFlowKnowledgeRuntimeClient.RuntimeResult result = client.runGraphRag(
new RagFlowKnowledgeRuntimeClient.RunGraphRagCommand(100L, 2001L, 4001L,
"dataset-1", "task-graph-1", "corr-graph-2", "hash-graph-2", 1));
assertEquals(RagFlowKnowledgeRuntimeClient.Status.SUCCEEDED, result.status());
assertEquals("gr-1", result.externalId());
}
@Test
void getKnowledgeGraphShouldSucceedWithGraphAndMindMapWithoutExternalId() {
HttpRagFlowKnowledgeRuntimeClient client = new StubHttpRagFlowKnowledgeRuntimeClient(
"https://ragflow.example", "configured-api-key", """
{"code":0,"data":{"graph":{"nodes":[{"id":"n1"}],"edges":[{"id":"e1"}]},"mind_map":{"root":"n1"}}}
""");
RagFlowKnowledgeRuntimeClient.RuntimeResult result = client.getKnowledgeGraph(
new RagFlowKnowledgeRuntimeClient.GetKnowledgeGraphCommand(100L, 2001L, 4001L,
"dataset-1", "corr-kg-1", "hash-kg-1", 1));
assertEquals(RagFlowKnowledgeRuntimeClient.Status.SUCCEEDED, result.status());
assertNull(result.externalId());
assertTrue(result.redactedSummary().contains("\"nodesCount\":1"));
assertTrue(result.redactedSummary().contains("\"edgesCount\":1"));
assertTrue(result.redactedSummary().contains("\"mindMapPresent\":true"));
}
@Test
void getKnowledgeGraphShouldSucceedWithEmptyGraph() {
HttpRagFlowKnowledgeRuntimeClient client = new StubHttpRagFlowKnowledgeRuntimeClient(
"https://ragflow.example", "configured-api-key",
"{\"code\":0,\"data\":{\"graph\":{\"nodes\":[],\"edges\":[]}}}");
RagFlowKnowledgeRuntimeClient.RuntimeResult result = client.getKnowledgeGraph(
new RagFlowKnowledgeRuntimeClient.GetKnowledgeGraphCommand(100L, 2001L, 4001L,
"dataset-1", "corr-kg-2", "hash-kg-2", 1));
assertEquals(RagFlowKnowledgeRuntimeClient.Status.SUCCEEDED, result.status());
assertNull(result.externalId());
assertTrue(result.redactedSummary().contains("\"nodesCount\":0"));
assertTrue(result.redactedSummary().contains("\"edgesCount\":0"));
}
@Test
void startParseShouldExposeOnlyMuseProcessingTaskAndCorrelation() {
HttpRagFlowKnowledgeRuntimeClient client = new StubHttpRagFlowKnowledgeRuntimeClient(
@ -105,11 +42,11 @@ class RagFlowKnowledgeRuntimeClientTest {
{"code":0,"data":{"task_id":"external-task-should-not-leak","id":"external-task-should-not-leak"}}
""");
RagFlowKnowledgeRuntimeClient.RuntimeResult result = client.startParseDocuments(
new RagFlowKnowledgeRuntimeClient.StartParseDocumentsCommand(100L, 2001L, 4001L,
KnowledgeRuntimeClient.RuntimeResult result = client.startParseDocuments(
new KnowledgeRuntimeClient.StartParseDocumentsCommand(100L, 2001L, 4001L,
"dataset-1", List.of("doc-1"), "muse-task-1", "corr-parse-1", "hash-parse-1", 1));
assertEquals(RagFlowKnowledgeRuntimeClient.Status.ACCEPTED, result.status());
assertEquals(KnowledgeRuntimeClient.Status.ACCEPTED, result.status());
assertNull(result.externalId());
assertEquals("muse-task-1", result.processingTaskId());
assertEquals("corr-parse-1", result.correlationId());
@ -121,8 +58,8 @@ class RagFlowKnowledgeRuntimeClientTest {
HttpRagFlowKnowledgeRuntimeClient client = new StubHttpRagFlowKnowledgeRuntimeClient(
"https://ragflow.example", "configured-api-key", "{\"code\":0,\"data\":{\"id\":\"dataset-1\"}}");
RagFlowKnowledgeRuntimeClient.RuntimeResult result = client.createDataset(
new RagFlowKnowledgeRuntimeClient.CreateDatasetCommand(100L, 2001L, 4001L,
KnowledgeRuntimeClient.RuntimeResult result = client.createDataset(
new KnowledgeRuntimeClient.CreateDatasetCommand(100L, 2001L, 4001L,
"KB", Map.of("apiKey", "secret-value", "authorization", "Bearer abcdefgh"),
"corr-1", "hash-1", 1));
@ -135,12 +72,12 @@ class RagFlowKnowledgeRuntimeClientTest {
CapturingHttpRagFlowKnowledgeRuntimeClient client = new CapturingHttpRagFlowKnowledgeRuntimeClient(
"https://ragflow.example", "configured-api-key", "{\"code\":0,\"data\":{\"id\":\"dataset-1\"}}");
RagFlowKnowledgeRuntimeClient.RuntimeResult result = client.createDataset(
new RagFlowKnowledgeRuntimeClient.CreateDatasetCommand(100L, 2001L, 4001L,
KnowledgeRuntimeClient.RuntimeResult result = client.createDataset(
new KnowledgeRuntimeClient.CreateDatasetCommand(100L, 2001L, 4001L,
"KB", Map.of("chunk_method", "naive", "apiKey", "secret-value"),
"corr-create-name-only", "hash-create-name-only", 1));
assertEquals(RagFlowKnowledgeRuntimeClient.Status.SUCCEEDED, result.status());
assertEquals(KnowledgeRuntimeClient.Status.SUCCEEDED, result.status());
String requestBody = new String(client.requests.get(0).body(), StandardCharsets.UTF_8);
assertTrue(requestBody.contains("KB"));
assertFalse(requestBody.contains("config"));
@ -153,12 +90,12 @@ class RagFlowKnowledgeRuntimeClientTest {
HttpRagFlowKnowledgeRuntimeClient client = new StubHttpRagFlowKnowledgeRuntimeClient(
"https://ragflow.example", "configured-api-key", "{\"code\":401,\"message\":\"auth failed\"}");
RagFlowKnowledgeRuntimeClient.RuntimeResult result = client.createDataset(
new RagFlowKnowledgeRuntimeClient.CreateDatasetCommand(100L, 2001L, 4001L,
KnowledgeRuntimeClient.RuntimeResult result = client.createDataset(
new KnowledgeRuntimeClient.CreateDatasetCommand(100L, 2001L, 4001L,
"KB", Map.of(), "corr-envelope-1", "hash-envelope-1", 1));
assertEquals(RagFlowKnowledgeRuntimeClient.Status.FAILED, result.status());
assertEquals(RagFlowKnowledgeRuntimeClient.FailureClass.AUTH_FAILED, result.failureClass());
assertEquals(KnowledgeRuntimeClient.Status.FAILED, result.status());
assertEquals(KnowledgeRuntimeClient.FailureClass.AUTH_FAILED, result.failureClass());
assertNull(result.externalId());
}
@ -167,12 +104,12 @@ class RagFlowKnowledgeRuntimeClientTest {
HttpRagFlowKnowledgeRuntimeClient client = new StubHttpRagFlowKnowledgeRuntimeClient(
"https://ragflow.example", "configured-api-key", "{\"code\":0,\"data\":{\"name\":\"KB\"}}");
RagFlowKnowledgeRuntimeClient.RuntimeResult result = client.createDataset(
new RagFlowKnowledgeRuntimeClient.CreateDatasetCommand(100L, 2001L, 4001L,
KnowledgeRuntimeClient.RuntimeResult result = client.createDataset(
new KnowledgeRuntimeClient.CreateDatasetCommand(100L, 2001L, 4001L,
"KB", Map.of(), "corr-envelope-2", "hash-envelope-2", 1));
assertEquals(RagFlowKnowledgeRuntimeClient.Status.FAILED, result.status());
assertEquals(RagFlowKnowledgeRuntimeClient.FailureClass.RESPONSE_SCHEMA_CHANGED, result.failureClass());
assertEquals(KnowledgeRuntimeClient.Status.FAILED, result.status());
assertEquals(KnowledgeRuntimeClient.FailureClass.RESPONSE_SCHEMA_CHANGED, result.failureClass());
}
@Test
@ -181,12 +118,12 @@ class RagFlowKnowledgeRuntimeClientTest {
HttpRagFlowKnowledgeRuntimeClient client = new StubHttpRagFlowKnowledgeRuntimeClient(
"https://ragflow.example", "configured-api-key", 500, providerError);
RagFlowKnowledgeRuntimeClient.RuntimeResult result = client.createDataset(
new RagFlowKnowledgeRuntimeClient.CreateDatasetCommand(100L, 2001L, 4001L,
KnowledgeRuntimeClient.RuntimeResult result = client.createDataset(
new KnowledgeRuntimeClient.CreateDatasetCommand(100L, 2001L, 4001L,
"KB", Map.of(), "corr-fail-body-1", "hash-fail-body-1", 1));
assertEquals(RagFlowKnowledgeRuntimeClient.Status.FAILED, result.status());
assertEquals(RagFlowKnowledgeRuntimeClient.FailureClass.RAGFLOW_UNAVAILABLE, result.failureClass());
assertEquals(KnowledgeRuntimeClient.Status.FAILED, result.status());
assertEquals(KnowledgeRuntimeClient.FailureClass.RAGFLOW_UNAVAILABLE, result.failureClass());
assertFalse(result.redactedSummary().contains(providerError));
assertTrue(result.redactedSummary().contains("\"redacted\":true"));
assertTrue(result.redactedSummary().contains("\"sha256\""));
@ -200,12 +137,12 @@ class RagFlowKnowledgeRuntimeClientTest {
"https://ragflow.example", "configured-api-key", 500,
"{\"message\":\"" + providerMessage + "\",\"error\":\"" + providerError + "\"}");
RagFlowKnowledgeRuntimeClient.RuntimeResult result = client.createDataset(
new RagFlowKnowledgeRuntimeClient.CreateDatasetCommand(100L, 2001L, 4001L,
KnowledgeRuntimeClient.RuntimeResult result = client.createDataset(
new KnowledgeRuntimeClient.CreateDatasetCommand(100L, 2001L, 4001L,
"KB", Map.of(), "corr-fail-body-json-1", "hash-fail-body-json-1", 1));
assertEquals(RagFlowKnowledgeRuntimeClient.Status.FAILED, result.status());
assertEquals(RagFlowKnowledgeRuntimeClient.FailureClass.RAGFLOW_UNAVAILABLE, result.failureClass());
assertEquals(KnowledgeRuntimeClient.Status.FAILED, result.status());
assertEquals(KnowledgeRuntimeClient.FailureClass.RAGFLOW_UNAVAILABLE, result.failureClass());
assertFalse(result.redactedSummary().contains(providerMessage));
assertFalse(result.redactedSummary().contains(providerError));
assertTrue(result.redactedSummary().contains("\"redacted\":true"));
@ -219,12 +156,12 @@ class RagFlowKnowledgeRuntimeClientTest {
"https://ragflow.example", "configured-api-key", 400,
"{\"errors\":[{\"message\":\"" + providerError + "\"}]}");
RagFlowKnowledgeRuntimeClient.RuntimeResult result = client.createDataset(
new RagFlowKnowledgeRuntimeClient.CreateDatasetCommand(100L, 2001L, 4001L,
KnowledgeRuntimeClient.RuntimeResult result = client.createDataset(
new KnowledgeRuntimeClient.CreateDatasetCommand(100L, 2001L, 4001L,
"KB", Map.of(), "corr-fail-body-json-2", "hash-fail-body-json-2", 1));
assertEquals(RagFlowKnowledgeRuntimeClient.Status.FAILED, result.status());
assertEquals(RagFlowKnowledgeRuntimeClient.FailureClass.VALIDATION_ERROR, result.failureClass());
assertEquals(KnowledgeRuntimeClient.Status.FAILED, result.status());
assertEquals(KnowledgeRuntimeClient.FailureClass.VALIDATION_ERROR, result.failureClass());
assertFalse(result.redactedSummary().contains(providerError));
assertTrue(result.redactedSummary().contains("\"redacted\":true"));
assertTrue(result.redactedSummary().contains("\"sha256\""));
@ -237,12 +174,12 @@ class RagFlowKnowledgeRuntimeClientTest {
"https://ragflow.example", "configured-api-key",
"{\"code\":422,\"message\":\"" + providerMessage + "\"}");
RagFlowKnowledgeRuntimeClient.RuntimeResult result = client.createDataset(
new RagFlowKnowledgeRuntimeClient.CreateDatasetCommand(100L, 2001L, 4001L,
KnowledgeRuntimeClient.RuntimeResult result = client.createDataset(
new KnowledgeRuntimeClient.CreateDatasetCommand(100L, 2001L, 4001L,
"KB", Map.of(), "corr-fail-body-2", "hash-fail-body-2", 1));
assertEquals(RagFlowKnowledgeRuntimeClient.Status.FAILED, result.status());
assertEquals(RagFlowKnowledgeRuntimeClient.FailureClass.VALIDATION_ERROR, result.failureClass());
assertEquals(KnowledgeRuntimeClient.Status.FAILED, result.status());
assertEquals(KnowledgeRuntimeClient.FailureClass.VALIDATION_ERROR, result.failureClass());
assertFalse(result.redactedSummary().contains(providerMessage));
assertTrue(result.redactedSummary().contains("\"redacted\":true"));
assertTrue(result.redactedSummary().contains("\"sha256\""));
@ -253,23 +190,23 @@ class RagFlowKnowledgeRuntimeClientTest {
CapturingHttpRagFlowKnowledgeRuntimeClient client = new CapturingHttpRagFlowKnowledgeRuntimeClient(
"https://ragflow.example", "configured-api-key", "{\"code\":0,\"data\":[{\"id\":\"doc-1\"}]}");
RagFlowKnowledgeRuntimeClient.RuntimeResult rejected = client.uploadDocuments(
new RagFlowKnowledgeRuntimeClient.UploadDocumentsCommand(100L, 2001L, 4001L,
"dataset-1", List.of(new RagFlowKnowledgeRuntimeClient.DocumentUpload(
KnowledgeRuntimeClient.RuntimeResult rejected = client.uploadDocuments(
new KnowledgeRuntimeClient.UploadDocumentsCommand(100L, 2001L, 4001L,
"dataset-1", List.of(new KnowledgeRuntimeClient.DocumentUpload(
5001L, 6001L, "file-ref-1", "novel.txt", "text/plain", 10L, null)),
"corr-upload-1", "hash-upload-1", 1));
assertEquals(RagFlowKnowledgeRuntimeClient.Status.FAILED, rejected.status());
assertEquals(RagFlowKnowledgeRuntimeClient.FailureClass.FILE_REJECTED, rejected.failureClass());
assertEquals(KnowledgeRuntimeClient.Status.FAILED, rejected.status());
assertEquals(KnowledgeRuntimeClient.FailureClass.FILE_REJECTED, rejected.failureClass());
RagFlowKnowledgeRuntimeClient.RuntimeResult uploaded = client.uploadDocuments(
new RagFlowKnowledgeRuntimeClient.UploadDocumentsCommand(100L, 2001L, 4001L,
"dataset-1", List.of(new RagFlowKnowledgeRuntimeClient.DocumentUpload(
KnowledgeRuntimeClient.RuntimeResult uploaded = client.uploadDocuments(
new KnowledgeRuntimeClient.UploadDocumentsCommand(100L, 2001L, 4001L,
"dataset-1", List.of(new KnowledgeRuntimeClient.DocumentUpload(
5001L, 6001L, "file-ref-1", "novel.txt", "text/plain", 5L,
"hello".getBytes(StandardCharsets.UTF_8))),
"corr-upload-2", "hash-upload-2", 1));
assertEquals(RagFlowKnowledgeRuntimeClient.Status.SUCCEEDED, uploaded.status());
assertEquals(KnowledgeRuntimeClient.Status.SUCCEEDED, uploaded.status());
assertEquals("POST", client.requests.get(0).method());
assertTrue(client.requests.get(0).contentType().startsWith("multipart/form-data; boundary="));
assertTrue(new String(client.requests.get(0).body(), StandardCharsets.ISO_8859_1).contains("filename=\"novel.txt\""));
@ -282,14 +219,14 @@ class RagFlowKnowledgeRuntimeClientTest {
"https://ragflow.example", "configured-api-key", "{\"code\":0,\"data\":[{\"id\":\"doc-1\"}]}");
byte[] binaryContent = new byte[]{0x25, 0x50, 0x44, 0x46, 0x00, (byte) 0x80, (byte) 0xFF, 0x01};
RagFlowKnowledgeRuntimeClient.RuntimeResult uploaded = client.uploadDocuments(
new RagFlowKnowledgeRuntimeClient.UploadDocumentsCommand(100L, 2001L, 4001L,
"dataset-1", List.of(new RagFlowKnowledgeRuntimeClient.DocumentUpload(
KnowledgeRuntimeClient.RuntimeResult uploaded = client.uploadDocuments(
new KnowledgeRuntimeClient.UploadDocumentsCommand(100L, 2001L, 4001L,
"dataset-1", List.of(new KnowledgeRuntimeClient.DocumentUpload(
5001L, 6001L, "file-ref-1", "binary.pdf", "application/pdf",
(long) binaryContent.length, binaryContent)),
"corr-upload-bin", "hash-upload-bin", 1));
assertEquals(RagFlowKnowledgeRuntimeClient.Status.SUCCEEDED, uploaded.status());
assertEquals(KnowledgeRuntimeClient.Status.SUCCEEDED, uploaded.status());
assertTrue(containsSubsequence(client.requests.get(0).body(), binaryContent));
}
@ -301,8 +238,8 @@ class RagFlowKnowledgeRuntimeClientTest {
HttpRagFlowKnowledgeRuntimeClient client = new StubHttpRagFlowKnowledgeRuntimeClient(
"https://ragflow.example", "configured-api-key", responseBody);
RagFlowKnowledgeRuntimeClient.RuntimeResult result = client.retrieveChunks(
new RagFlowKnowledgeRuntimeClient.RetrieveChunksCommand(100L, 2001L, 4001L,
KnowledgeRuntimeClient.RuntimeResult result = client.retrieveChunks(
new KnowledgeRuntimeClient.RetrieveChunksCommand(100L, 2001L, 4001L,
List.of("dataset-1"), List.of("doc-1"), "private retrieval question",
3, 0.5, Map.of("prompt", "private prompt"), "corr-redact-1", "hash-redact-1", 1));
@ -312,20 +249,6 @@ class RagFlowKnowledgeRuntimeClientTest {
assertFalse(result.redactedSummary().contains("private prompt"));
}
@Test
void listDatasetsShouldSucceedForEmptyCollectionWithoutExternalId() {
HttpRagFlowKnowledgeRuntimeClient client = new StubHttpRagFlowKnowledgeRuntimeClient(
"https://ragflow.example", "configured-api-key", "{\"code\":0,\"data\":[]}");
RagFlowKnowledgeRuntimeClient.RuntimeResult result = client.listDatasets(
new RagFlowKnowledgeRuntimeClient.ListDatasetsCommand(100L, 2001L,
Map.of(), "corr-list-1", "hash-list-1", 1));
assertEquals(RagFlowKnowledgeRuntimeClient.Status.SUCCEEDED, result.status());
assertNull(result.externalId());
assertTrue(result.redactedSummary().contains("\"count\":0"));
}
@Test
void retrieveChunksShouldSucceedWithoutTopLevelIdAndKeepOnlyMetadataSummary() {
HttpRagFlowKnowledgeRuntimeClient client = new StubHttpRagFlowKnowledgeRuntimeClient(
@ -333,12 +256,12 @@ class RagFlowKnowledgeRuntimeClientTest {
{"code":0,"data":{"chunks":[{"content":"secret chunk content","score":0.91}],"response":"raw answer text"}}
""");
RagFlowKnowledgeRuntimeClient.RuntimeResult result = client.retrieveChunks(
new RagFlowKnowledgeRuntimeClient.RetrieveChunksCommand(100L, 2001L, 4001L,
KnowledgeRuntimeClient.RuntimeResult result = client.retrieveChunks(
new KnowledgeRuntimeClient.RetrieveChunksCommand(100L, 2001L, 4001L,
List.of("dataset-1"), List.of("doc-1"), "private retrieval question",
3, 0.5, Map.of("prompt", "private prompt"), "corr-retrieve-1", "hash-retrieve-1", 1));
assertEquals(RagFlowKnowledgeRuntimeClient.Status.SUCCEEDED, result.status());
assertEquals(KnowledgeRuntimeClient.Status.SUCCEEDED, result.status());
assertNull(result.externalId());
assertTrue(result.redactedSummary().contains("\"chunksCount\":1"));
assertFalse(result.redactedSummary().contains("secret chunk content"));
@ -355,11 +278,11 @@ class RagFlowKnowledgeRuntimeClientTest {
"https://ragflow.example", "configured-api-key",
"{\"code\":0,\"data\":{\"chunks\":[{\"content\":\"c\",\"score\":0.9}]}}");
RagFlowKnowledgeRuntimeClient.RuntimeResult result = client.retrieveChunks(
new RagFlowKnowledgeRuntimeClient.RetrieveChunksCommand(100L, 2001L, 4001L,
KnowledgeRuntimeClient.RuntimeResult result = client.retrieveChunks(
new KnowledgeRuntimeClient.RetrieveChunksCommand(100L, 2001L, 4001L,
List.of("dataset-1"), null, "q", 5, 0.2, null, "corr-null-doc", "hash-null-doc", 1));
assertEquals(RagFlowKnowledgeRuntimeClient.Status.SUCCEEDED, result.status());
assertEquals(KnowledgeRuntimeClient.Status.SUCCEEDED, result.status());
String requestBody = new String(client.requests.get(0).body(), StandardCharsets.UTF_8);
assertFalse(requestBody.contains("document_ids"), "document_ids 为 null 时不应出现在请求 body");
assertTrue(requestBody.contains("dataset_ids"));
@ -373,36 +296,19 @@ class RagFlowKnowledgeRuntimeClientTest {
"https://ragflow.example", "configured-api-key",
"{\"code\":0,\"data\":{\"chunks\":[{\"content\":\"c\",\"score\":0.9}]}}");
RagFlowKnowledgeRuntimeClient.RuntimeResult result = client.retrieveChunks(
new RagFlowKnowledgeRuntimeClient.RetrieveChunksCommand(100L, 2001L, 4001L,
KnowledgeRuntimeClient.RuntimeResult result = client.retrieveChunks(
new KnowledgeRuntimeClient.RetrieveChunksCommand(100L, 2001L, 4001L,
List.of("dataset-1"), List.of("doc-1"), "q", 5, 0.2,
Map.of("logic", "and", "conditions", List.of(Map.of(
"name", "source", "comparison_operator", "=", "value", "public"))),
"corr-metadata-condition", "hash-metadata-condition", 1));
assertEquals(RagFlowKnowledgeRuntimeClient.Status.SUCCEEDED, result.status());
assertEquals(KnowledgeRuntimeClient.Status.SUCCEEDED, result.status());
String requestBody = new String(client.requests.get(0).body(), StandardCharsets.UTF_8);
assertTrue(requestBody.contains("\"metadata_condition\""));
assertFalse(requestBody.contains("\"metadata_filter\""));
}
@Test
void listChunksShouldSucceedWithoutTopLevelId() {
HttpRagFlowKnowledgeRuntimeClient client = new StubHttpRagFlowKnowledgeRuntimeClient(
"https://ragflow.example", "configured-api-key", """
{"code":0,"data":{"chunks":[{"id":"chunk-1","content":"secret chunk content"}]}}
""");
RagFlowKnowledgeRuntimeClient.RuntimeResult result = client.listChunks(
new RagFlowKnowledgeRuntimeClient.ListChunksCommand(100L, 2001L, 4001L,
"dataset-1", "doc-1", Map.of(), "corr-chunk-list-1", "hash-chunk-list-1", 1));
assertEquals(RagFlowKnowledgeRuntimeClient.Status.SUCCEEDED, result.status());
assertNull(result.externalId());
assertTrue(result.redactedSummary().contains("\"chunksCount\":1"));
assertFalse(result.redactedSummary().contains("secret chunk content"));
}
@Test
void pollDocumentStatusesShouldUseRagflowIdQueryForSingleDocumentSmoke() {
CapturingHttpRagFlowKnowledgeRuntimeClient client = new CapturingHttpRagFlowKnowledgeRuntimeClient(
@ -414,13 +320,13 @@ class RagFlowKnowledgeRuntimeClientTest {
]}}
""");
RagFlowKnowledgeRuntimeClient.RuntimeResult result = client.pollDocumentStatuses(
new RagFlowKnowledgeRuntimeClient.PollDocumentStatusesCommand(100L, 2001L, 4001L,
KnowledgeRuntimeClient.RuntimeResult result = client.pollDocumentStatuses(
new KnowledgeRuntimeClient.PollDocumentStatusesCommand(100L, 2001L, 4001L,
"dataset-1", List.of("doc-1"), "task-1", "corr-poll-1", "hash-poll-1", 1));
assertEquals(RagFlowKnowledgeRuntimeClient.Status.SUCCEEDED, result.status());
assertEquals(KnowledgeRuntimeClient.Status.SUCCEEDED, result.status());
assertEquals(List.of("doc-1"), result.documentStatuses().stream()
.map(RagFlowKnowledgeRuntimeClient.DocumentStatus::ragflowDocumentId).toList());
.map(KnowledgeRuntimeClient.DocumentStatus::ragflowDocumentId).toList());
assertTrue(client.requests.get(0).path().contains("id=doc-1"));
assertFalse(client.requests.get(0).path().contains("document_ids="));
}
@ -430,12 +336,12 @@ class RagFlowKnowledgeRuntimeClientTest {
CapturingHttpRagFlowKnowledgeRuntimeClient client = new CapturingHttpRagFlowKnowledgeRuntimeClient(
"https://ragflow.example", "configured-api-key", "{\"code\":0,\"data\":{\"docs\":[]}}");
RagFlowKnowledgeRuntimeClient.RuntimeResult result = client.pollDocumentStatuses(
new RagFlowKnowledgeRuntimeClient.PollDocumentStatusesCommand(100L, 2001L, 4001L,
KnowledgeRuntimeClient.RuntimeResult result = client.pollDocumentStatuses(
new KnowledgeRuntimeClient.PollDocumentStatusesCommand(100L, 2001L, 4001L,
"dataset-1", List.of("doc-1", "doc-2"), "task-1", "corr-poll-many", "hash-poll-many", 1));
assertEquals(RagFlowKnowledgeRuntimeClient.Status.FAILED, result.status());
assertEquals(RagFlowKnowledgeRuntimeClient.FailureClass.VALIDATION_ERROR, result.failureClass());
assertEquals(KnowledgeRuntimeClient.Status.FAILED, result.status());
assertEquals(KnowledgeRuntimeClient.FailureClass.VALIDATION_ERROR, result.failureClass());
assertTrue(result.redactedSummary().contains("single document status polling"));
assertEquals(0, client.requests.size());
}

View File

@ -46,21 +46,17 @@ class P1rRagFlowLiveAcceptanceIT {
String datasetName = "p1r-live-" + suffix;
String correlationPrefix = "p1r5-ragflow-live-" + suffix;
RagFlowKnowledgeRuntimeClient.RuntimeResult health = client.health(
new RagFlowKnowledgeRuntimeClient.HealthCommand(correlationPrefix + "-health", hash("health"), 1));
assertSucceeded(health, "health");
RagFlowKnowledgeRuntimeClient.RuntimeResult createDataset = client.createDataset(
new RagFlowKnowledgeRuntimeClient.CreateDatasetCommand(1L, 1L, 1L, datasetName,
KnowledgeRuntimeClient.RuntimeResult createDataset = client.createDataset(
new KnowledgeRuntimeClient.CreateDatasetCommand(1L, 1L, 1L, datasetName,
Map.of(), correlationPrefix + "-dataset", hash(datasetName), 1));
assertSucceeded(createDataset, "createDataset");
assertNotNull(createDataset.externalId(), "RAGFlow createDataset 必须返回 dataset id");
String datasetId = createDataset.externalId();
byte[] content = liveDocumentContent(suffix);
RagFlowKnowledgeRuntimeClient.RuntimeResult upload = client.uploadDocuments(
new RagFlowKnowledgeRuntimeClient.UploadDocumentsCommand(1L, 1L, 1L, datasetId,
List.of(new RagFlowKnowledgeRuntimeClient.DocumentUpload(1L, 1L,
KnowledgeRuntimeClient.RuntimeResult upload = client.uploadDocuments(
new KnowledgeRuntimeClient.UploadDocumentsCommand(1L, 1L, 1L, datasetId,
List.of(new KnowledgeRuntimeClient.DocumentUpload(1L, 1L,
"p1r-live-" + suffix, datasetName + ".txt", "text/plain",
(long) content.length, content)),
correlationPrefix + "-upload", hash(datasetId + "-upload"), 1));
@ -68,25 +64,17 @@ class P1rRagFlowLiveAcceptanceIT {
assertNotNull(upload.externalId(), "RAGFlow uploadDocuments 必须返回 document id");
String documentId = upload.externalId();
RagFlowKnowledgeRuntimeClient.RuntimeResult parse = client.startParseDocuments(
new RagFlowKnowledgeRuntimeClient.StartParseDocumentsCommand(1L, 1L, 1L, datasetId,
KnowledgeRuntimeClient.RuntimeResult parse = client.startParseDocuments(
new KnowledgeRuntimeClient.StartParseDocumentsCommand(1L, 1L, 1L, datasetId,
List.of(documentId), correlationPrefix + "-parse-task",
correlationPrefix + "-parse", hash(documentId + "-parse"), 1));
assertEquals(RagFlowKnowledgeRuntimeClient.Status.ACCEPTED, parse.status(),
assertEquals(KnowledgeRuntimeClient.Status.ACCEPTED, parse.status(),
() -> "startParseDocuments 未进入 accepted" + parse.redactedSummary());
RagFlowKnowledgeRuntimeClient.RuntimeResult poll = waitUntilSearchable(client, datasetId, documentId,
KnowledgeRuntimeClient.RuntimeResult poll = waitUntilSearchable(client, datasetId, documentId,
correlationPrefix);
RagFlowKnowledgeRuntimeClient.RuntimeResult listChunks = client.listChunks(
new RagFlowKnowledgeRuntimeClient.ListChunksCommand(1L, 1L, 1L, datasetId, documentId,
Map.of("page", 1, "page_size", 5),
correlationPrefix + "-chunks", hash(documentId + "-chunks"), 1));
assertSucceeded(listChunks, "listChunks");
assertTrue(countFromRedactedSummary(listChunks.redactedSummary(), "chunksCount") > 0,
() -> "listChunks 没有返回 chunk不能作为 retrieval 正向证据:" + listChunks.redactedSummary());
RagFlowKnowledgeRuntimeClient.RuntimeResult retrieval = client.retrieveChunks(
new RagFlowKnowledgeRuntimeClient.RetrieveChunksCommand(1L, 1L, 1L,
KnowledgeRuntimeClient.RuntimeResult retrieval = client.retrieveChunks(
new KnowledgeRuntimeClient.RetrieveChunksCommand(1L, 1L, 1L,
List.of(datasetId), List.of(documentId), "P1R live acceptance marker 是什么?",
3, 0.0, Map.of(), correlationPrefix + "-retrieval",
hash(documentId + "-retrieval"), 1));
@ -94,66 +82,26 @@ class P1rRagFlowLiveAcceptanceIT {
assertTrue(countFromRedactedSummary(retrieval.redactedSummary(), "chunksCount") > 0,
() -> "retrieveChunks 没有返回 chunk不能作为 retrieval 正向证据:" + retrieval.redactedSummary());
RagFlowKnowledgeRuntimeClient.RuntimeResult graphClosed = null;
if (!graphRagReady) {
graphClosed = client.runGraphRag(new RagFlowKnowledgeRuntimeClient.RunGraphRagCommand(1L, 1L, 1L, datasetId,
correlationPrefix + "-graph-task", correlationPrefix + "-graph-closed",
hash(datasetId + "-graph-closed"), 1));
assertEquals(RagFlowKnowledgeRuntimeClient.Status.FAILED, graphClosed.status());
assertEquals(RagFlowKnowledgeRuntimeClient.FailureClass.ATTRIBUTION_NOT_CONFIGURED,
graphClosed.failureClass());
}
Map<String, Object> evidence = redactedEvidence(baseUrl, apiKey, datasetName, datasetId, documentId,
health, createDataset, upload, parse, poll, listChunks, retrieval, graphClosed, graphRagReady);
createDataset, upload, parse, poll, retrieval, graphRagReady);
System.out.println(JsonUtils.toJsonString(evidence));
}
@Test
void shouldRunGraphRagOnlyWhenAttributionReady() {
assumeTrue(externalAcceptanceEnabled(), "P1R 外部验收未启用,设置 MUSE_P1R_EXTERNAL_ACCEPTANCE=true 后才运行");
assumeTrue(Boolean.parseBoolean(envOrDefault(GRAPHRAG_READY_ENV, "false")),
"GraphRAG attribution 未就绪,保持 fail-closed不执行真实 GraphRAG run/trace");
String baseUrl = requiredEnv(BASE_URL_ENV);
String apiKey = requiredEnv(API_KEY_ENV);
String datasetId = requiredEnv("MUSE_KNOWLEDGE_RAGFLOW_GRAPHRAG_DATASET_ID");
HttpRagFlowKnowledgeRuntimeClient client = client(baseUrl, apiKey, true);
String correlationId = "p1r5-ragflow-graphrag-live-" + Instant.now().toEpochMilli();
RagFlowKnowledgeRuntimeClient.RuntimeResult run = client.runGraphRag(
new RagFlowKnowledgeRuntimeClient.RunGraphRagCommand(1L, 1L, 1L, datasetId,
correlationId + "-task", correlationId + "-run", hash(datasetId + "-run"), 1));
assertSucceeded(run, "runGraphRag");
RagFlowKnowledgeRuntimeClient.RuntimeResult trace = client.traceGraphRag(
new RagFlowKnowledgeRuntimeClient.TraceGraphRagCommand(1L, 1L, 1L, datasetId,
correlationId + "-trace", hash(datasetId + "-trace"), 1));
assertSucceeded(trace, "traceGraphRag");
System.out.println(JsonUtils.toJsonString(Map.of(
"acceptance", "p1r5-ragflow-graphrag-live",
"endpoint", endpointSummary(baseUrl, "/api/v1/datasets/" + datasetId + "/run_graphrag"),
"apiKey", secretFingerprint(apiKey),
"datasetId", datasetId,
"runStatus", run.status(),
"traceStatus", trace.status())));
}
private HttpRagFlowKnowledgeRuntimeClient client(String baseUrl, String apiKey, boolean graphRagReady) {
return new HttpRagFlowKnowledgeRuntimeClient(baseUrl, apiKey,
Duration.ofSeconds(intEnv("MUSE_KNOWLEDGE_RAGFLOW_TIMEOUT_SECONDS", 30)),
intEnv("MUSE_KNOWLEDGE_RAGFLOW_RETRY_BUDGET", 0), graphRagReady);
}
private RagFlowKnowledgeRuntimeClient.RuntimeResult waitUntilSearchable(
private KnowledgeRuntimeClient.RuntimeResult waitUntilSearchable(
HttpRagFlowKnowledgeRuntimeClient client, String datasetId, String documentId, String correlationPrefix)
throws InterruptedException {
long deadline = System.nanoTime()
+ Duration.ofSeconds(intEnv("MUSE_P1R_RAGFLOW_PARSE_WAIT_SECONDS", 120)).toNanos();
RagFlowKnowledgeRuntimeClient.RuntimeResult last = null;
KnowledgeRuntimeClient.RuntimeResult last = null;
int attempt = 1;
while (System.nanoTime() < deadline) {
last = client.pollDocumentStatuses(new RagFlowKnowledgeRuntimeClient.PollDocumentStatusesCommand(
last = client.pollDocumentStatuses(new KnowledgeRuntimeClient.PollDocumentStatusesCommand(
1L, 1L, 1L, datasetId, List.of(documentId), correlationPrefix + "-parse-task",
correlationPrefix + "-poll-" + attempt, hash(documentId + "-poll-" + attempt), attempt));
assertSucceeded(last, "pollDocumentStatuses");
@ -167,9 +115,9 @@ class P1rRagFlowLiveAcceptanceIT {
+ (last == null ? "none" : last.redactedSummary()));
}
private void assertSucceeded(RagFlowKnowledgeRuntimeClient.RuntimeResult result, String operation) {
private void assertSucceeded(KnowledgeRuntimeClient.RuntimeResult result, String operation) {
assertNotNull(result, operation + " result 不能为空");
assertEquals(RagFlowKnowledgeRuntimeClient.Status.SUCCEEDED, result.status(),
assertEquals(KnowledgeRuntimeClient.Status.SUCCEEDED, result.status(),
() -> operation + " failed: " + result.redactedSummary());
}
@ -184,14 +132,11 @@ class P1rRagFlowLiveAcceptanceIT {
private Map<String, Object> redactedEvidence(String baseUrl, String apiKey, String datasetName,
String datasetId, String documentId,
RagFlowKnowledgeRuntimeClient.RuntimeResult health,
RagFlowKnowledgeRuntimeClient.RuntimeResult createDataset,
RagFlowKnowledgeRuntimeClient.RuntimeResult upload,
RagFlowKnowledgeRuntimeClient.RuntimeResult parse,
RagFlowKnowledgeRuntimeClient.RuntimeResult poll,
RagFlowKnowledgeRuntimeClient.RuntimeResult listChunks,
RagFlowKnowledgeRuntimeClient.RuntimeResult retrieval,
RagFlowKnowledgeRuntimeClient.RuntimeResult graphClosed,
KnowledgeRuntimeClient.RuntimeResult createDataset,
KnowledgeRuntimeClient.RuntimeResult upload,
KnowledgeRuntimeClient.RuntimeResult parse,
KnowledgeRuntimeClient.RuntimeResult poll,
KnowledgeRuntimeClient.RuntimeResult retrieval,
boolean graphRagReady) {
Map<String, Object> evidence = new LinkedHashMap<>();
evidence.put("acceptance", "p1r5-ragflow-live");
@ -200,23 +145,19 @@ class P1rRagFlowLiveAcceptanceIT {
evidence.put("datasetName", datasetName);
evidence.put("datasetId", datasetId);
evidence.put("documentId", documentId);
evidence.put("health", operationSummary(health));
evidence.put("createDataset", operationSummary(createDataset));
evidence.put("uploadDocuments", operationSummary(upload));
evidence.put("startParseDocuments", operationSummary(parse));
evidence.put("pollDocumentStatuses", operationSummary(poll));
evidence.put("documentStatuses", poll.documentStatuses());
evidence.put("listChunks", operationSummary(listChunks));
evidence.put("retrieveChunks", operationSummary(retrieval));
evidence.put("retrievalChunksCount", countFromRedactedSummary(retrieval.redactedSummary(), "chunksCount"));
evidence.put("graphRagAttributionReady", graphRagReady);
evidence.put("graphRag", graphClosed == null ? Map.of("status", "separate_opt_in_required")
: operationSummary(graphClosed));
evidence.put("datasetRetained", true);
return evidence;
}
private Map<String, Object> operationSummary(RagFlowKnowledgeRuntimeClient.RuntimeResult result) {
private Map<String, Object> operationSummary(KnowledgeRuntimeClient.RuntimeResult result) {
Map<String, Object> summary = new LinkedHashMap<>();
summary.put("operation", result.operation().wireName());
summary.put("status", result.status());

View File

@ -14,7 +14,7 @@ class RagFlowKnowledgeRuntimeClientConfigurationTest {
void should_registerUnavailableClientWhenRequiredConfigMissing() {
RagFlowKnowledgeRuntimeClientConfiguration configuration = new RagFlowKnowledgeRuntimeClientConfiguration();
RagFlowKnowledgeRuntimeClient client = configuration.ragFlowKnowledgeRuntimeClient(new MockEnvironment());
KnowledgeRuntimeClient client = configuration.ragFlowKnowledgeRuntimeClient(new MockEnvironment());
assertInstanceOf(UnavailableRagFlowKnowledgeRuntimeClient.class, client);
}
@ -29,7 +29,7 @@ class RagFlowKnowledgeRuntimeClientConfigurationTest {
.withProperty("muse.knowledge.ragflow.graphrag-attribution-ready", "true");
RagFlowKnowledgeRuntimeClientConfiguration configuration = new RagFlowKnowledgeRuntimeClientConfiguration();
RagFlowKnowledgeRuntimeClient client = configuration.ragFlowKnowledgeRuntimeClient(environment);
KnowledgeRuntimeClient client = configuration.ragFlowKnowledgeRuntimeClient(environment);
assertInstanceOf(HttpRagFlowKnowledgeRuntimeClient.class, client);
}

View File

@ -17,7 +17,7 @@ import cn.iocoder.muse.module.knowledge.application.muse.MuseKnowledgeMaterializ
import cn.iocoder.muse.module.knowledge.application.muse.MuseKnowledgeProcessingTaskService;
import cn.iocoder.muse.module.knowledge.application.muse.facade.HttpRagFlowKnowledgeRuntimeClient;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeFileFacade;
import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient;
import cn.iocoder.muse.module.knowledge.application.muse.facade.KnowledgeRuntimeClient;
import cn.iocoder.muse.module.knowledge.controller.app.muse.vo.AppKnowledgeDocumentVO;
import cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeBaseDO;
import cn.iocoder.muse.module.knowledge.dal.mysql.muse.MuseKnowledgeBaseMapper;
@ -269,20 +269,12 @@ class P1rKnowledgeRuntimeEndToEndLiveAcceptanceIT {
private ExternalVerification verifyExternalRagFlowSearchable(ConfigurableApplicationContext context,
LiveResult result,
DbFacts facts) throws InterruptedException {
RagFlowKnowledgeRuntimeClient client = context.getBean(RagFlowKnowledgeRuntimeClient.class);
KnowledgeRuntimeClient client = context.getBean(KnowledgeRuntimeClient.class);
String correlationPrefix = "p1r5-knowledge-e2e-external-" + result.suffix();
RagFlowKnowledgeRuntimeClient.RuntimeResult poll = waitUntilSearchable(client, facts.ragflowDatasetId(),
KnowledgeRuntimeClient.RuntimeResult poll = waitUntilSearchable(client, facts.ragflowDatasetId(),
facts.ragflowDocumentId(), facts.processingTaskId(), correlationPrefix);
RagFlowKnowledgeRuntimeClient.RuntimeResult listChunks = client.listChunks(
new RagFlowKnowledgeRuntimeClient.ListChunksCommand(TENANT_ID, OWNER_USER_ID, result.kbId(),
facts.ragflowDatasetId(), facts.ragflowDocumentId(), Map.of("page", 1, "page_size", 5),
correlationPrefix + "-chunks", hash(facts.ragflowDocumentId() + "-chunks"), 1));
assertSucceeded(listChunks, "listChunks");
int chunksCount = countFromRedactedSummary(listChunks.redactedSummary(), "chunksCount");
assertTrue(chunksCount > 0, () -> "listChunks 没有返回 chunk不能作为外部检索证据" + listChunks.redactedSummary());
RagFlowKnowledgeRuntimeClient.RuntimeResult retrieval = client.retrieveChunks(
new RagFlowKnowledgeRuntimeClient.RetrieveChunksCommand(TENANT_ID, OWNER_USER_ID, result.kbId(),
KnowledgeRuntimeClient.RuntimeResult retrieval = client.retrieveChunks(
new KnowledgeRuntimeClient.RetrieveChunksCommand(TENANT_ID, OWNER_USER_ID, result.kbId(),
List.of(facts.ragflowDatasetId()), List.of(facts.ragflowDocumentId()),
"P1R-5 Knowledge live acceptance marker 是什么?", 3, 0.0, Map.of(),
correlationPrefix + "-retrieval", hash(facts.ragflowDocumentId() + "-retrieval"), 1));
@ -290,10 +282,10 @@ class P1rKnowledgeRuntimeEndToEndLiveAcceptanceIT {
int retrievalChunksCount = countFromRedactedSummary(retrieval.redactedSummary(), "chunksCount");
assertTrue(retrievalChunksCount > 0,
() -> "retrieveChunks 没有返回 chunk不能作为外部检索证据" + retrieval.redactedSummary());
return new ExternalVerification(poll, listChunks, retrieval, chunksCount, retrievalChunksCount);
return new ExternalVerification(poll, retrieval, retrievalChunksCount);
}
private RagFlowKnowledgeRuntimeClient.RuntimeResult waitUntilSearchable(RagFlowKnowledgeRuntimeClient client,
private KnowledgeRuntimeClient.RuntimeResult waitUntilSearchable(KnowledgeRuntimeClient client,
String datasetId,
String documentId,
String processingTaskId,
@ -301,10 +293,10 @@ class P1rKnowledgeRuntimeEndToEndLiveAcceptanceIT {
throws InterruptedException {
long deadline = System.nanoTime()
+ Duration.ofSeconds(intEnv("MUSE_P1R_RAGFLOW_PARSE_WAIT_SECONDS", 120)).toNanos();
RagFlowKnowledgeRuntimeClient.RuntimeResult last = null;
KnowledgeRuntimeClient.RuntimeResult last = null;
int attempt = 1;
while (System.nanoTime() < deadline) {
last = client.pollDocumentStatuses(new RagFlowKnowledgeRuntimeClient.PollDocumentStatusesCommand(
last = client.pollDocumentStatuses(new KnowledgeRuntimeClient.PollDocumentStatusesCommand(
TENANT_ID, OWNER_USER_ID, 0L, datasetId, List.of(documentId), processingTaskId,
correlationPrefix + "-poll-" + attempt, hash(documentId + "-poll-" + attempt), attempt));
assertSucceeded(last, "pollDocumentStatuses");
@ -318,9 +310,9 @@ class P1rKnowledgeRuntimeEndToEndLiveAcceptanceIT {
+ (last == null ? "none" : last.redactedSummary()));
}
private void assertSucceeded(RagFlowKnowledgeRuntimeClient.RuntimeResult result, String operation) {
private void assertSucceeded(KnowledgeRuntimeClient.RuntimeResult result, String operation) {
assertNotNull(result, operation + " result 不能为空");
assertEquals(RagFlowKnowledgeRuntimeClient.Status.SUCCEEDED, result.status(),
assertEquals(KnowledgeRuntimeClient.Status.SUCCEEDED, result.status(),
() -> operation + " failed: " + result.redactedSummary());
}
@ -635,14 +627,12 @@ class P1rKnowledgeRuntimeEndToEndLiveAcceptanceIT {
evidence.put("scope", "external_verification_not_business_persisted");
evidence.put("pollDocumentStatuses", operationSummary(verification.poll()));
evidence.put("documentStatuses", verification.poll().documentStatuses());
evidence.put("listChunks", operationSummary(verification.listChunks()));
evidence.put("retrieveChunks", operationSummary(verification.retrieval()));
evidence.put("chunksCount", verification.chunksCount());
evidence.put("retrievalChunksCount", verification.retrievalChunksCount());
return evidence;
}
private Map<String, Object> operationSummary(RagFlowKnowledgeRuntimeClient.RuntimeResult result) {
private Map<String, Object> operationSummary(KnowledgeRuntimeClient.RuntimeResult result) {
Map<String, Object> summary = new LinkedHashMap<>();
summary.put("operation", result.operation().wireName());
summary.put("status", result.status());
@ -1016,10 +1006,8 @@ class P1rKnowledgeRuntimeEndToEndLiveAcceptanceIT {
boolean allRagflowResponseSummariesNonEmpty) {
}
private record ExternalVerification(RagFlowKnowledgeRuntimeClient.RuntimeResult poll,
RagFlowKnowledgeRuntimeClient.RuntimeResult listChunks,
RagFlowKnowledgeRuntimeClient.RuntimeResult retrieval,
int chunksCount,
private record ExternalVerification(KnowledgeRuntimeClient.RuntimeResult poll,
KnowledgeRuntimeClient.RuntimeResult retrieval,
int retrievalChunksCount) {
}
@ -1045,7 +1033,7 @@ class P1rKnowledgeRuntimeEndToEndLiveAcceptanceIT {
static class LiveAcceptanceConfiguration {
@Bean
RagFlowKnowledgeRuntimeClient ragFlowKnowledgeRuntimeClient() {
KnowledgeRuntimeClient 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)),

View File

@ -15,7 +15,7 @@ import cn.iocoder.muse.module.knowledge.application.muse.MuseKnowledgeAuditServi
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.knowledge.application.muse.facade.KnowledgeRuntimeClient;
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;
@ -433,23 +433,23 @@ class P1rMarketKbForkMaterializationIT {
*/
private void uploadPrivateDocToPublisherLiveDataset(ConfigurableApplicationContext context, ForkSeed seed,
String privateMagic, String suffix) throws InterruptedException {
RagFlowKnowledgeRuntimeClient client = context.getBean(RagFlowKnowledgeRuntimeClient.class);
KnowledgeRuntimeClient client = context.getBean(KnowledgeRuntimeClient.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,
KnowledgeRuntimeClient.RuntimeResult upload = client.uploadDocuments(
new KnowledgeRuntimeClient.UploadDocumentsCommand(TENANT_ID, PUBLISHER_USER_ID,
seed.publisherKbId(), seed.publisherLiveDatasetId(),
List.of(new RagFlowKnowledgeRuntimeClient.DocumentUpload(seed.privateDocId(), 1L,
List.of(new KnowledgeRuntimeClient.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(),
assertEquals(KnowledgeRuntimeClient.Status.SUCCEEDED, upload.status(),
() -> "私有文档传发布者活库失败: " + upload.redactedSummary());
RagFlowKnowledgeRuntimeClient.RuntimeResult parse = client.startParseDocuments(
new RagFlowKnowledgeRuntimeClient.StartParseDocumentsCommand(TENANT_ID, PUBLISHER_USER_ID,
KnowledgeRuntimeClient.RuntimeResult parse = client.startParseDocuments(
new KnowledgeRuntimeClient.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(),
assertEquals(KnowledgeRuntimeClient.Status.ACCEPTED, parse.status(),
() -> "私有文档在发布者活库重索引未受理: " + parse.redactedSummary());
// 等私有 magic 在发布者活库真正 searchable确保活库里检索得到 magic这个前提成立否则不泄露断言是空证
waitUntilSearchable(context, seed.publisherLiveDatasetId(), upload.externalId(), suffix);
@ -459,13 +459,13 @@ class P1rMarketKbForkMaterializationIT {
/** 反证基线:发布者活库里检索 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,
KnowledgeRuntimeClient client = context.getBean(KnowledgeRuntimeClient.class);
KnowledgeRuntimeClient.RuntimeResult retrieval = client.retrieveChunks(
new KnowledgeRuntimeClient.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(),
assertEquals(KnowledgeRuntimeClient.Status.SUCCEEDED, retrieval.status(),
() -> "发布者活库检索 magic 失败(前提不成立): " + retrieval.redactedSummary());
assertTrue(retrieval.redactedSummary() != null && retrieval.redactedSummary().contains("chunksCount"),
"发布者活库检索应返回 chunksCount 摘要");
@ -590,12 +590,12 @@ class P1rMarketKbForkMaterializationIT {
// ============================== 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,
KnowledgeRuntimeClient client = context.getBean(KnowledgeRuntimeClient.class);
KnowledgeRuntimeClient.RuntimeResult result = client.createDataset(
new KnowledgeRuntimeClient.CreateDatasetCommand(TENANT_ID, PUBLISHER_USER_ID, 0L,
datasetName, Map.of(), "fork-create-" + suffix + "-" + datasetName,
hash(datasetName), 1));
assertEquals(RagFlowKnowledgeRuntimeClient.Status.SUCCEEDED, result.status(),
assertEquals(KnowledgeRuntimeClient.Status.SUCCEEDED, result.status(),
() -> "createDataset 失败: " + result.redactedSummary());
assertNotNull(result.externalId(), "createDataset 必须返回 dataset id");
return result.externalId();
@ -607,29 +607,29 @@ class P1rMarketKbForkMaterializationIT {
*/
private void waitUntilSearchable(ConfigurableApplicationContext context, String datasetId, String documentId,
String suffix) throws InterruptedException {
RagFlowKnowledgeRuntimeClient client = context.getBean(RagFlowKnowledgeRuntimeClient.class);
KnowledgeRuntimeClient client = context.getBean(KnowledgeRuntimeClient.class);
long deadline = System.nanoTime()
+ Duration.ofSeconds(intEnv("MUSE_P1R_RAGFLOW_PARSE_WAIT_SECONDS", 180)).toNanos();
int attempt = 1;
RagFlowKnowledgeRuntimeClient.RuntimeResult last = null;
KnowledgeRuntimeClient.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,
KnowledgeRuntimeClient.RuntimeResult retrieval = client.retrieveChunks(
new KnowledgeRuntimeClient.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
if (retrieval.status() == KnowledgeRuntimeClient.Status.SUCCEEDED
&& countFromSummary(retrieval.redactedSummary(), "chunksCount") > 0) {
return;
}
last = retrieval;
} else {
last = client.pollDocumentStatuses(new RagFlowKnowledgeRuntimeClient.PollDocumentStatusesCommand(
last = client.pollDocumentStatuses(new KnowledgeRuntimeClient.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
if (last.status() == KnowledgeRuntimeClient.Status.SUCCEEDED
&& last.documentStatuses().stream().anyMatch(s -> "searchable".equals(s.museStatus()))) {
return;
}
@ -775,7 +775,7 @@ class P1rMarketKbForkMaterializationIT {
}
@Bean
RagFlowKnowledgeRuntimeClient ragFlowKnowledgeRuntimeClient() {
KnowledgeRuntimeClient 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)),