diff --git a/muse-cloud/muse-framework/muse-spring-boot-starter-mybatis/src/main/java/cn/iocoder/muse/framework/mybatis/config/MuseMybatisAutoConfiguration.java b/muse-cloud/muse-framework/muse-spring-boot-starter-mybatis/src/main/java/cn/iocoder/muse/framework/mybatis/config/MuseMybatisAutoConfiguration.java index 0e1c3601..c5f3cd85 100644 --- a/muse-cloud/muse-framework/muse-spring-boot-starter-mybatis/src/main/java/cn/iocoder/muse/framework/mybatis/config/MuseMybatisAutoConfiguration.java +++ b/muse-cloud/muse-framework/muse-spring-boot-starter-mybatis/src/main/java/cn/iocoder/muse/framework/mybatis/config/MuseMybatisAutoConfiguration.java @@ -4,6 +4,7 @@ import cn.hutool.core.collection.CollUtil; import cn.hutool.core.util.StrUtil; import cn.iocoder.muse.framework.common.util.json.JsonUtils; import cn.iocoder.muse.framework.mybatis.core.handler.DefaultDBFieldHandler; +import cn.iocoder.muse.framework.mybatis.core.muse.MuseContractPersistenceService; import com.baomidou.mybatisplus.annotation.DbType; import com.baomidou.mybatisplus.autoconfigure.MybatisPlusAutoConfiguration; import com.baomidou.mybatisplus.core.handlers.IJsonTypeHandler; @@ -21,6 +22,7 @@ import org.mybatis.spring.annotation.MapperScan; import org.springframework.boot.autoconfigure.AutoConfiguration; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Import; import org.springframework.core.env.ConfigurableEnvironment; import java.util.List; @@ -34,6 +36,7 @@ import java.util.concurrent.TimeUnit; @AutoConfiguration(before = MybatisPlusAutoConfiguration.class) // 目的:先于 MyBatis Plus 自动配置,避免 @MapperScan 可能扫描不到 Mapper 打印 warn 日志 @MapperScan(value = "${muse.info.base-package}", annotationClass = Mapper.class, lazyInitialization = "${mybatis.lazy-initialization:false}") // Mapper 懒加载,目前仅用于单元测试 +@Import(MuseContractPersistenceService.class) public class MuseMybatisAutoConfiguration { static { diff --git a/muse-cloud/muse-framework/muse-spring-boot-starter-mybatis/src/main/java/cn/iocoder/muse/framework/mybatis/core/muse/MuseContractPersistenceService.java b/muse-cloud/muse-framework/muse-spring-boot-starter-mybatis/src/main/java/cn/iocoder/muse/framework/mybatis/core/muse/MuseContractPersistenceService.java new file mode 100644 index 00000000..a3ef3ae0 --- /dev/null +++ b/muse-cloud/muse-framework/muse-spring-boot-starter-mybatis/src/main/java/cn/iocoder/muse/framework/mybatis/core/muse/MuseContractPersistenceService.java @@ -0,0 +1,2107 @@ +package cn.iocoder.muse.framework.mybatis.core.muse; + +import cn.hutool.core.util.StrUtil; +import cn.iocoder.muse.framework.common.exception.ServiceException; +import cn.iocoder.muse.framework.common.muse.MuseApiContractSupport; +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.core.type.TypeReference; +import com.fasterxml.jackson.databind.ObjectMapper; +import jakarta.annotation.Resource; +import org.springframework.dao.DuplicateKeyException; +import org.springframework.dao.EmptyResultDataAccessException; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Transactional; + +import java.nio.charset.StandardCharsets; +import java.security.MessageDigest; +import java.security.NoSuchAlgorithmException; +import java.time.LocalDateTime; +import java.util.*; + +/** + * Muse 合同入口持久化应用服务。 + * + *

P1 阶段仍保留 catch-all Controller,但请求不再停留在占位响应:本服务统一完成合同校验、命令幂等、 + * 操作审计,并把关键业务事实同步写入各领域表。具体领域事实仍由调用方传入的 domain 决定,admin/app 只是不 + * 同入口。

+ */ +@Service +public class MuseContractPersistenceService { + + /** Muse 合同写命令冲突错误码。 */ + private static final int CONTRACT_CONFLICT = 1000000007; + /** Muse 合同资源不存在错误码。 */ + private static final int CONTRACT_NOT_FOUND = 1000000008; + /** Muse 合同暂不支持的资源错误码。 */ + private static final int CONTRACT_UNSUPPORTED = 1000000009; + /** Muse 合同写命令尚未实现错误码。 */ + private static final int CONTRACT_NOT_IMPLEMENTED = 1000000010; + /** P1 默认租户 ID,后续接入 Yudao 租户上下文后替换。 */ + private static final long DEFAULT_TENANT_ID = 0L; + /** 默认分页页码。 */ + private static final int DEFAULT_PAGE_NO = 1; + /** 默认分页大小。 */ + private static final int DEFAULT_PAGE_SIZE = 20; + + @Resource + private JdbcTemplate jdbcTemplate; + + private final ObjectMapper objectMapper = new ObjectMapper(); + + /** + * 处理 Muse 合同请求,并将读写行为落到领域持久化层。 + * + * @param domains 当前入口拥有的领域 + * @param side app 或 admin + * @param httpMethod HTTP 方法 + * @param requestUri 请求路径 + * @param headerCommandId X-Command-Id 请求头 + * @param queryParams 查询参数 + * @param body JSON 请求体 + * @param actorUserId 当前操作者 + * @return 可序列化业务响应 + */ + @Transactional(rollbackFor = Exception.class) + public Map handle(Set domains, String side, String httpMethod, String requestUri, + String headerCommandId, Map queryParams, + Map body, Long actorUserId) { + Map contract = MuseApiContractSupport.handle(domains, side, httpMethod, requestUri, + headerCommandId, queryParams, body, actorUserId); + OperationContext context = OperationContext.from(contract, httpMethod, requestUri, headerCommandId, + queryParams, body, actorUserId); + if ("GET".equals(httpMethod) || isReadOnlyCommand(context.operationId())) { + return handleRead(context); + } + return handleWrite(context); + } + + private Map handleRead(OperationContext context) { + if (isCurrentUserOperation(context.operationId())) { + return currentUserResponse(context); + } + ResourceSpec spec = resolveSpec(context); + if (spec == null || spec.tableName() == null) { + return readOperationRecords(context, spec == null ? "operation" : spec.resourceType()); + } + if (isListOperation(context.operationId())) { + return listRows(context, spec); + } + return getRow(context, spec); + } + + private Map handleWrite(OperationContext context) { + Map replay = replayCommand(context); + if (replay != null) { + replay.put("idempotentReplay", true); + return replay; + } + if (!isMaterializedWrite(context)) { + throw new ServiceException(CONTRACT_NOT_IMPLEMENTED, + "P1 暂未实现该写命令的业务落库: " + context.operationId()); + } + + ResourceSpec spec = resolveSpec(context); + Map response = acceptedResponse(context, spec); + Long operationRecordId; + try { + // 先占用 commandId 对应的统一操作记录,再执行业务写入,避免并发请求同时撞业务唯一键。 + operationRecordId = insertOperationRecord(context, spec, null, response); + } catch (DuplicateKeyException duplicate) { + replay = replayCommand(context); + if (replay != null) { + replay.put("idempotentReplay", true); + return replay; + } + throw duplicate; + } + response.put("operationRecordId", operationRecordId); + Long resourceId = materializeDomainFact(context, spec, response); + if (resourceId != null) { + response.put("resourceId", resourceId); + } + updateOperationResponse(operationRecordId, resourceId, response); + return response; + } + + private Map acceptedResponse(OperationContext context, ResourceSpec spec) { + Map response = new LinkedHashMap<>(); + response.put("operationId", context.operationId()); + response.put("module", context.domain()); + response.put("entry", context.side()); + response.put("status", "accepted"); + response.put("pathVariables", context.pathVariables()); + response.put("query", context.queryParams()); + response.put("commandId", context.commandId()); + response.put("resourceType", spec == null ? "operation" : spec.resourceType()); + response.put("persisted", true); + response.put("audit", audit(context)); + return response; + } + + private Map audit(OperationContext context) { + Map audit = new LinkedHashMap<>(); + audit.put("actorUserId", context.actorUserId()); + audit.put("acceptedAt", LocalDateTime.now().toString()); + audit.put("contractOperation", context.operationId()); + return audit; + } + + private Map currentUserResponse(OperationContext context) { + Map data = baseReadResponse(context, "profile"); + data.put("id", context.actorUserId()); + data.put("userId", context.actorUserId()); + data.put("revision", 1); + data.put("profile", latestOperationPayload(context, "profile")); + data.put("entitlements", selectAccountEntitlement(context.actorUserId())); + return data; + } + + private Map listRows(OperationContext context, ResourceSpec spec) { + WhereClause where = buildWhere(context, spec, true); + int pageNo = intValue(context.queryParams().get("pageNo"), DEFAULT_PAGE_NO); + int pageSize = intValue(context.queryParams().get("pageSize"), DEFAULT_PAGE_SIZE); + int offset = Math.max(pageNo - 1, 0) * pageSize; + List args = new ArrayList<>(where.args()); + args.add(pageSize); + args.add(offset); + + String itemsJson = jdbcTemplate.queryForObject(""" + SELECT COALESCE(jsonb_agg(to_jsonb(t)), '[]'::jsonb)::text + FROM ( + SELECT * FROM %s + WHERE %s + ORDER BY id DESC + LIMIT ? OFFSET ? + ) t + """.formatted(spec.tableName(), where.sql()), String.class, args.toArray()); + Long total = jdbcTemplate.queryForObject("SELECT COUNT(1) FROM %s WHERE %s".formatted(spec.tableName(), where.sql()), + Long.class, where.args().toArray()); + + Map data = baseReadResponse(context, spec.resourceType()); + data.put("items", readList(itemsJson)); + data.put("total", total == null ? 0L : total); + data.put("pageNo", pageNo); + data.put("pageSize", pageSize); + return data; + } + + private Map getRow(OperationContext context, ResourceSpec spec) { + WhereClause where = buildWhere(context, spec, false); + String json = queryOptionalString("SELECT to_jsonb(t)::text FROM %s t WHERE %s LIMIT 1".formatted(spec.tableName(), where.sql()), + where.args().toArray()); + if (json == null) { + return readOperationRecords(context, spec.resourceType()); + } + Map data = baseReadResponse(context, spec.resourceType()); + data.putAll(readMap(json)); + return data; + } + + private Map readOperationRecords(OperationContext context, String resourceType) { + if (isListOperation(context.operationId())) { + return listOperationRecords(context, resourceType); + } + return getOperationRecord(context, resourceType); + } + + private Map listOperationRecords(OperationContext context, String resourceType) { + int pageNo = intValue(context.queryParams().get("pageNo"), DEFAULT_PAGE_NO); + int pageSize = intValue(context.queryParams().get("pageSize"), DEFAULT_PAGE_SIZE); + int offset = Math.max(pageNo - 1, 0) * pageSize; + String itemsJson = jdbcTemplate.queryForObject(""" + SELECT COALESCE(jsonb_agg(to_jsonb(t)), '[]'::jsonb)::text + FROM ( + SELECT id, domain, side, operation_id, command_id, actor_user_id, resource_type, resource_id, + resource_key, parent_type, parent_id, status, revision, request_payload, response_payload, + create_time, update_time + FROM muse_domain_operation_record + WHERE tenant_id = ? AND deleted = FALSE AND domain = ? AND resource_type = ? + ORDER BY id DESC + LIMIT ? OFFSET ? + ) t + """, String.class, DEFAULT_TENANT_ID, context.domain(), resourceType, pageSize, offset); + Long total = jdbcTemplate.queryForObject(""" + SELECT COUNT(1) + FROM muse_domain_operation_record + WHERE tenant_id = ? AND deleted = FALSE AND domain = ? AND resource_type = ? + """, Long.class, DEFAULT_TENANT_ID, context.domain(), resourceType); + Map data = baseReadResponse(context, resourceType); + data.put("items", readList(itemsJson)); + data.put("total", total == null ? 0L : total); + data.put("pageNo", pageNo); + data.put("pageSize", pageSize); + return data; + } + + private Map getOperationRecord(OperationContext context, String resourceType) { + Long resourceId = firstPathLong(context, "requestId", "jobId", "eventId", "credentialId", "runId", "taskId"); + String resourceKey = firstPathString(context, "correlationId", "handoffToken"); + List args = new ArrayList<>(); + args.add(DEFAULT_TENANT_ID); + args.add(context.domain()); + args.add(resourceType); + StringBuilder where = new StringBuilder("tenant_id = ? AND deleted = FALSE AND domain = ? AND resource_type = ?"); + if (resourceId != null) { + where.append(" AND resource_id = ?"); + args.add(resourceId); + } else if (StrUtil.isNotBlank(resourceKey)) { + where.append(" AND resource_key = ?"); + args.add(resourceKey); + } else { + where.append(" AND operation_id = ?"); + args.add(context.operationId()); + } + String json = queryOptionalString("SELECT to_jsonb(t)::text FROM muse_domain_operation_record t WHERE " + + where + " ORDER BY id DESC LIMIT 1", args.toArray()); + Map data = baseReadResponse(context, resourceType); + if (json != null) { + data.putAll(readMap(json)); + } + return data; + } + + private Map baseReadResponse(OperationContext context, String resourceType) { + Map data = new LinkedHashMap<>(); + data.put("operationId", context.operationId()); + data.put("module", context.domain()); + data.put("entry", context.side()); + data.put("status", "ok"); + data.put("resourceType", resourceType); + data.put("pathVariables", context.pathVariables()); + data.put("query", context.queryParams()); + data.put("audit", audit(context)); + return data; + } + + private Long materializeDomainFact(OperationContext context, ResourceSpec spec, Map response) { + if (spec != null && spec.tableName() != null) { + enforceExpectedGuards(context, spec); + } + return switch (context.domain()) { + case "account" -> materializeAccount(context, response); + case "meta" -> materializeMeta(context, response); + case "ai" -> materializeAi(context, response); + case "knowledge" -> materializeKnowledge(context, response); + case "market" -> materializeMarket(context, response); + default -> null; + }; + } + + private Long materializeAccount(OperationContext context, Map response) { + String operationId = context.operationId(); + Long userId = firstPathLong(context, "userId"); + Long accountUserId = userId == null ? context.actorUserId() : userId; + if ("adminCreateQuotaAdjustment".equals(operationId)) { + Long auditId = insertAccountAudit(context, accountUserId, "quota_adjustment"); + applyQuotaAdjustments(context, accountUserId); + response.put("auditLogId", auditId); + response.put("accountUserId", accountUserId); + return auditId; + } + if ("adminCreateNewApiBinding".equals(operationId) || "appRecheckNewApiBinding".equals(operationId)) { + Long bindingId = upsertNewApiBinding(context, accountUserId); + response.put("bindingStatus", "active"); + return bindingId; + } + if ("adminCreateQuotaRequest".equals(operationId) || "appCreateQuotaRequest".equals(operationId)) { + Long taskId = insertWorkflowTask(context, "quotaRequest", "quotaRequest", + "user", accountUserId, null, null); + response.put("requestId", taskId); + return taskId; + } + if ("adminCreateCallAttributionJob".equals(operationId)) { + Long taskId = insertWorkflowTask(context, "callAttributionJob", "callAttributionJob", + "correlation", null, null, null); + response.put("jobId", taskId); + return taskId; + } + if ("appCreateExportTask".equals(operationId)) { + Long taskId = insertWorkflowTask(context, "accountExportTask", "exportTask", + "user", context.actorUserId(), null, null); + response.put("taskId", taskId); + return taskId; + } + if ("appAcknowledgeSecurityEvent".equals(operationId)) { + Long eventId = firstPathLong(context, "eventId"); + int updated = jdbcTemplate.update(""" + UPDATE muse_member_security_event + SET acknowledged = TRUE, acknowledged_at = CURRENT_TIMESTAMP, updater = ? + WHERE id = ? AND tenant_id = ? AND account_user_id = ? AND deleted = FALSE + """, actor(context), eventId, DEFAULT_TENANT_ID, context.actorUserId()); + if (updated == 0) { + throw new ServiceException(CONTRACT_NOT_FOUND, "安全事件不存在或无权确认"); + } + return eventId; + } + if ("updateProfile".equals(operationId)) { + return null; + } + return null; + } + + private Long materializeMeta(OperationContext context, Map response) { + String operationId = context.operationId(); + String schemaKey = firstPathString(context, "schemaKey"); + if ("saveMetaSchemaDraft".equals(operationId)) { + enforceMetaSchemaActiveVersion(context, schemaKey, "expectedVersion"); + Long schemaId = ensureMetaSchema(context, schemaKey); + Integer versionNo = nextInt("SELECT COALESCE(MAX(version_no), 0) + 1 FROM muse_meta_schema_version WHERE tenant_id = ? AND schema_key = ?", + DEFAULT_TENANT_ID, schemaKey); + Long versionId = insertMetaSchemaVersion(context, schemaId, schemaKey, versionNo, "draft"); + response.put("schemaId", schemaId); + response.put("draftVersion", versionNo); + return versionId; + } + if ("publishMetaSchemaDraft".equals(operationId)) { + enforceMetaSchemaActiveVersion(context, schemaKey, "expectedVersion"); + return updateMetaSchemaVersionStatus(context, schemaKey, firstPathInt(context, "draftVersion"), "published"); + } + if ("activateMetaSchemaVersion".equals(operationId) || "rollbackMetaSchemaVersion".equals(operationId)) { + enforceMetaSchemaActiveVersion(context, schemaKey, "expectedActiveVersion"); + Integer version = firstPathInt(context, "version"); + Long versionId = updateMetaSchemaVersionStatus(context, schemaKey, version, "active"); + jdbcTemplate.update(""" + UPDATE muse_meta_schema + SET active_version_id = ?, projection_version = projection_version + 1, updater = ? + WHERE tenant_id = ? AND schema_key = ? AND deleted = FALSE + """, versionId, actor(context), DEFAULT_TENANT_ID, schemaKey); + return versionId; + } + if ("deprecateMetaSchemaVersion".equals(operationId)) { + enforcePathVersion(context, "expectedVersion", "version", "MetaSchema 目标版本冲突"); + return updateMetaSchemaVersionStatus(context, schemaKey, firstPathInt(context, "version"), "deprecated"); + } + if ("setMetaSchemaGrayRules".equals(operationId)) { + Integer version = firstPathInt(context, "version"); + Long versionId = requireMetaSchemaVersion(schemaKey, version); + jdbcTemplate.update(""" + UPDATE muse_meta_schema_version + SET impact_preview_snapshot = CAST(? AS jsonb), updater = ? + WHERE id = ? AND tenant_id = ? AND deleted = FALSE + """, toJson(context.body()), actor(context), versionId, DEFAULT_TENANT_ID); + return versionId; + } + if ("activateFunctionChainVersion".equals(operationId)) { + enforceFunctionChainActiveVersion(context); + return activateFunctionChainVersion(context); + } + return null; + } + + private Long materializeAi(OperationContext context, Map response) { + String operationId = context.operationId(); + if ("createAiTask".equals(operationId)) { + requireAppWorkOwner(context, firstPathLong(context, "workId")); + Long id = jdbcTemplate.queryForObject(""" + INSERT INTO muse_ai_generation(work_id, chapter_id, block_id, agent_id, task_type, status, + idempotency_key, input_snapshot, creator, updater, tenant_id) + VALUES (?, ?, ?, ?, ?, 'queued', ?, CAST(? AS jsonb), ?, ?, ?) + RETURNING id + """, Long.class, + firstPathLong(context, "workId"), longValue(context.body().get("chapterId")), + longValue(context.body().get("blockId")), longValue(context.body().get("agentId")), + text(context.body(), "intent", "generation"), context.commandId(), toJson(context.body()), + actor(context), actor(context), DEFAULT_TENANT_ID); + response.put("taskId", id); + return id; + } + if ("adminCreateAgent".equals(operationId) || "createUserAgent".equals(operationId)) { + Long id = insertAgent(context, "createUserAgent".equals(operationId) ? "user" : "system"); + response.put("agentId", id); + return id; + } + if ("adminCreateAgentVersion".equals(operationId) || "createAgentVersion".equals(operationId)) { + Long agentId = firstPathLong(context, "agentId"); + Long id = insertAgentVersion(context, agentId); + response.put("agentVersionId", id); + return id; + } + if ("testAgent".equals(operationId)) { + response.put("result", Map.of("status", "ready", "output", context.body().getOrDefault("input", ""))); + return firstPathLong(context, "agentId"); + } + if ("adminCreatePromptVersion".equals(operationId)) { + Long id = insertPromptVersion(context, firstPathString(context, "promptKey")); + response.put("promptVersionId", id); + return id; + } + if ("adminActivatePromptVersion".equals(operationId)) { + return activatePromptVersion(context); + } + if ("adminCreateQualityPolicyVersion".equals(operationId)) { + Long id = insertQualityPolicyVersion(context, firstPathString(context, "policyKey")); + response.put("qualityPolicyVersionId", id); + return id; + } + if ("adminStartEvaluationRun".equals(operationId)) { + Long runId = insertWorkflowTask(context, "evaluationRun", "evaluationRun", + "qualityPolicy", null, null, null); + response.put("runId", runId); + return runId; + } + if ("adminCreateOrAdjustToolGrant".equals(operationId)) { + Long id = insertToolGrant(context); + response.put("toolGrantId", id); + return id; + } + if ("precheckAgentSlot".equals(operationId)) { + Long workId = firstPathLong(context, "workId"); + requireAppWorkOwner(context, workId); + Long precheckId = insertWorkflowTask(context, "agentSlotPrecheck", "agentSlotPrecheck", + "work", workId, "agentSlot", null); + response.put("agentSlotPrecheckId", precheckId); + return precheckId; + } + if ("bindAgentSlot".equals(operationId)) { + Long id = upsertAgentSlotBinding(context); + response.put("bindingId", id); + return id; + } + if ("rejectSuggestion".equals(operationId)) { + Long suggestionId = firstPathLong(context, "suggestionId"); + jdbcTemplate.update(""" + UPDATE muse_ai_suggestion + SET status = 'rejected', command_id = ?, updater = ? + WHERE id = ? AND tenant_id = ? AND deleted = FALSE + AND (? <> 'app' OR EXISTS ( + SELECT 1 FROM muse_content_work w + WHERE w.id = muse_ai_suggestion.work_id AND w.tenant_id = muse_ai_suggestion.tenant_id + AND w.owner_user_id = ? AND w.deleted = FALSE + )) + """, context.commandId(), actor(context), suggestionId, DEFAULT_TENANT_ID, + context.side(), context.actorUserId()); + return suggestionId; + } + if ("adminRetryJob".equals(operationId)) { + Long jobId = firstPathLong(context, "jobId"); + int updated = jdbcTemplate.update(""" + UPDATE muse_ai_generation + SET status = 'queued', retry_count = retry_count + 1, error_code = NULL, error_message = NULL, + updater = ? + WHERE id = ? AND tenant_id = ? AND deleted = FALSE + """, actor(context), jobId, DEFAULT_TENANT_ID); + if (updated == 0) { + throw new ServiceException(CONTRACT_NOT_FOUND, "AI Job 不存在"); + } + return jobId; + } + if (operationId.endsWith("CancelJob") || operationId.endsWith("cancel") || operationId.startsWith("cancel")) { + Long jobId = firstPathLong(context, "jobId"); + jdbcTemplate.update(""" + UPDATE muse_ai_generation + SET status = 'canceled', updater = ? + WHERE id = ? AND tenant_id = ? AND deleted = FALSE + """, actor(context), jobId, DEFAULT_TENANT_ID); + return jobId; + } + if ("adminRetrySourceEvent".equals(operationId)) { + Long eventId = firstPathLong(context, "eventId"); + Long taskId = insertWorkflowTask(context, "sourceEventRetry", "sourceEventRetry", + "sourceEvent", eventId, null, null); + response.put("eventId", eventId); + response.put("taskId", taskId); + return taskId; + } + if ("recheckSourceStatus".equals(operationId)) { + Long taskId = insertWorkflowTask(context, "sourceStatusRecheck", "sourceStatusCheck", + text(context.body(), "targetOwner", "source"), longValue(context.body().get("targetId")), + text(context.body(), "sourceType", null), longValue(context.body().get("sourceId"))); + response.put("sourceStatusCheckId", taskId); + return taskId; + } + return null; + } + + private Long materializeKnowledge(OperationContext context, Map response) { + String operationId = context.operationId(); + if ("createKnowledgeBase".equals(operationId) || "createGlobalKnowledgeBase".equals(operationId)) { + Long id = insertKnowledgeBase(context, operationId.startsWith("createGlobal") ? "global" : "user"); + response.put("kbId", id); + return id; + } + if ("updateKnowledgeBase".equals(operationId) || "updateGlobalKnowledgeBase".equals(operationId) + || "disableKnowledgeBase".equals(operationId) || "restoreKnowledgeBase".equals(operationId) + || "enableGlobalKnowledgeBase".equals(operationId) || "disableGlobalKnowledgeBase".equals(operationId)) { + return updateKnowledgeBase(context); + } + if ("deleteKnowledgeBase".equals(operationId)) { + return deleteKnowledgeBase(context); + } + if ("activateGlobalKBVersion".equals(operationId)) { + return activateKnowledgeBaseVersion(context); + } + if ("saveGlobalKBAccessPolicyDraft".equals(operationId)) { + Long draftId = insertWorkflowTask(context, "globalKBAccessPolicyDraft", "accessPolicyDraft", + "knowledgeBase", firstPathLong(context, "kbId"), null, null); + response.put("draftId", draftId); + return draftId; + } + if ("publishGlobalKBAccessPolicy".equals(operationId)) { + Long taskId = insertWorkflowTask(context, "globalKBAccessPolicyPublish", "accessPolicyPublish", + "knowledgeBase", firstPathLong(context, "kbId"), null, null); + response.put("policyPublishTaskId", taskId); + return taskId; + } + if ("reindexGlobalKnowledgeBase".equals(operationId) || "reindexKnowledgeBase".equals(operationId)) { + Long taskId = insertWorkflowTask(context, "knowledgeBaseReindex", "processingTask", + "knowledgeBase", firstPathLong(context, "kbId"), null, null); + response.put("taskId", taskId); + return taskId; + } + if ("triggerGlobalKBSourceEvent".equals(operationId)) { + Long eventId = insertWorkflowTask(context, "globalKBSourceEvent", "sourceEvent", + "knowledgeBase", firstPathLong(context, "kbId"), text(context.body(), "eventType", null), null); + response.put("eventId", eventId); + return eventId; + } + if (operationId.contains("upload") && operationId.contains("Document")) { + Long id = insertKnowledgeDocument(context); + response.put("documentId", id); + return id; + } + if ("deleteKBDocument".equals(operationId) || "deleteGlobalKBDocument".equals(operationId)) { + return deleteKnowledgeDocument(context); + } + if (operationId.contains("DocumentVersion")) { + Long id = insertKnowledgeDocumentVersion(context); + response.put("documentVersionId", id); + return id; + } + if ("updateEntity".equals(operationId)) { + return updateKnowledgeEntity(context); + } + if ("createRelation".equals(operationId)) { + Long id = insertKnowledgeRelation(context); + response.put("relationId", id); + return id; + } + if ("confirmKnowledgeDraft".equals(operationId) || "ignoreKnowledgeDraft".equals(operationId) + || "recheckKnowledgeDraft".equals(operationId)) { + return updateKnowledgeDraft(context); + } + if ("createKnowledgeBinding".equals(operationId)) { + Long id = insertKnowledgeBinding(context); + response.put("bindingId", id); + return id; + } + if ("createKnowledgeBindingPrecheck".equals(operationId)) { + Long workId = firstPathLong(context, "workId"); + requireAppWorkOwner(context, workId); + Long precheckId = insertWorkflowTask(context, "knowledgeBindingPrecheck", "bindingPrecheck", + "work", workId, text(context.body(), "sourceType", null), longValue(context.body().get("sourceId"))); + response.put("kbBindPrecheckId", precheckId); + return precheckId; + } + if ("deleteKnowledgeBinding".equals(operationId)) { + Long bindingId = firstPathLong(context, "bindingId"); + jdbcTemplate.update(""" + UPDATE muse_knowledge_binding + SET binding_status = 'deleted', deleted = TRUE, command_id = ?, updater = ? + WHERE id = ? AND tenant_id = ? AND deleted = FALSE + AND (? <> 'app' OR EXISTS ( + SELECT 1 FROM muse_content_work w + WHERE w.id = muse_knowledge_binding.work_id AND w.tenant_id = muse_knowledge_binding.tenant_id + AND w.owner_user_id = ? AND w.deleted = FALSE + )) + """, context.commandId(), actor(context), bindingId, DEFAULT_TENANT_ID, + context.side(), context.actorUserId()); + return bindingId; + } + if ("createKBPublishReadiness".equals(operationId) || "createKBMarketPublishReadiness".equals(operationId)) { + Long readinessId = insertWorkflowTask(context, "kbPublishReadiness", "publishReadiness", + "knowledgeBase", firstPathLong(context, "kbId"), null, null); + response.put("publishReadinessId", readinessId); + return readinessId; + } + if ("createKBPublishSnapshot".equals(operationId)) { + Long snapshotId = insertWorkflowTask(context, "kbPublishSnapshot", "publishSnapshot", + "knowledgeBase", firstPathLong(context, "kbId"), null, null); + response.put("snapshotId", snapshotId); + return snapshotId; + } + if ("createKBExportTask".equals(operationId)) { + Long taskId = insertWorkflowTask(context, "kbExportTask", "exportTask", + "knowledgeBase", firstPathLong(context, "kbId"), null, null); + response.put("taskId", taskId); + return taskId; + } + if ("disableInstalledKnowledgeBase".equals(operationId) || "deleteInstalledKnowledgeBase".equals(operationId) + || "restoreInstalledKnowledgeBase".equals(operationId)) { + return updateInstalledKnowledgeBase(context); + } + return null; + } + + private Long materializeMarket(OperationContext context, Map response) { + String operationId = context.operationId(); + if ("savePublishDraft".equals(operationId)) { + Long assetId = insertMarketAsset(context, "draft"); + Long versionId = insertMarketAssetVersion(context, assetId, "draft"); + response.put("draftId", assetId); + response.put("assetVersionId", versionId); + return assetId; + } + if ("submitPublishRequest".equals(operationId)) { + Long id = insertPublishRequest(context); + response.put("requestId", id); + return id; + } + if (operationId.startsWith("adminApprove") || operationId.startsWith("adminReject") + || operationId.equals("withdrawPublishRequest")) { + return updatePublishRequest(context); + } + if ("installMarketplaceAsset".equals(operationId)) { + Long id = insertMarketInstallation(context); + response.put("installId", id); + return id; + } + if ("favoriteAsset".equals(operationId) || "unfavoriteAsset".equals(operationId)) { + Long favoriteId = upsertMarketFavorite(context, "favoriteAsset".equals(operationId)); + response.put("favoriteId", favoriteId); + return favoriteId; + } + if ("purchaseAsset".equals(operationId)) { + Long purchaseId = insertMarketPurchase(context); + response.put("purchaseId", purchaseId); + return purchaseId; + } + if ("createBindPrecheck".equals(operationId)) { + Long precheckId = insertWorkflowTask(context, "marketBindPrecheck", "bindPrecheck", + text(context.body(), "targetOwner", "marketAsset"), firstPathLong(context, "assetId"), + text(context.body(), "targetAction", null), null); + response.put("bindPrecheckId", precheckId); + return precheckId; + } + if ("createMarketplaceHandoff".equals(operationId)) { + Long id = insertMarketHandoff(context, response); + return id; + } + if ("cancelHandoff".equals(operationId)) { + return cancelMarketHandoff(context); + } + if ("submitAppeal".equals(operationId)) { + Long id = insertMarketAppeal(context); + response.put("appealId", id); + return id; + } + if ("supplementAppeal".equals(operationId) || "withdrawAppeal".equals(operationId) + || "adminResolveAppeal".equals(operationId)) { + return updateMarketAppeal(context); + } + if ("adminDelistAsset".equals(operationId) || "adminRecallAsset".equals(operationId)) { + Long assetId = firstPathLong(context, "assetId"); + jdbcTemplate.update(""" + UPDATE muse_market_asset + SET listing_status = ?, status = ?, command_id = ?, revision = revision + 1, updater = ? + WHERE id = ? AND tenant_id = ? AND deleted = FALSE + """, operationId.equals("adminDelistAsset") ? "delisted" : "recalled", + operationId.equals("adminDelistAsset") ? "inactive" : "recalled", + context.commandId(), actor(context), assetId, DEFAULT_TENANT_ID); + return assetId; + } + return firstPathLong(context, "assetId", "requestId", "appealId"); + } + + private Long insertWorkflowTask(OperationContext context, String taskType, String resourceType, + String targetType, Long targetId, String sourceType, Long sourceId) { + Long ownerUserId = "app".equals(context.side()) ? context.actorUserId() : firstPathLong(context, "userId"); + String correlationId = firstNonBlank(text(context.body(), "correlationId", null), + context.queryParams().get("correlationId"), firstPathString(context, "correlationId")); + String expectedStatus = text(context.body(), "expectedStatus", null); + Map resultPayload = new LinkedHashMap<>(); + resultPayload.put("status", "queued"); + resultPayload.put("resourceType", resourceType); + return jdbcTemplate.queryForObject(""" + INSERT INTO muse_domain_workflow_task(domain, side, operation_id, task_type, status, + actor_user_id, owner_user_id, target_type, target_id, parent_type, parent_id, + correlation_id, source_type, source_id, expected_status, request_payload, result_payload, + command_id, creator, updater, tenant_id) + VALUES (?, ?, ?, ?, 'queued', ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, CAST(? AS jsonb), CAST(? AS jsonb), + ?, ?, ?, ?) + RETURNING id + """, Long.class, context.domain(), context.side(), context.operationId(), taskType, + context.actorUserId(), ownerUserId, targetType, targetId, parentType(context, null), + parentId(context, null), correlationId, sourceType, sourceId, expectedStatus, toJson(context.body()), + toJson(resultPayload), context.commandId(), actor(context), actor(context), DEFAULT_TENANT_ID); + } + + private Long activateKnowledgeBaseVersion(OperationContext context) { + Long kbId = firstPathLong(context, "kbId"); + Integer version = firstPathInt(context, "version"); + int updated = jdbcTemplate.update(""" + UPDATE muse_knowledge_base + SET active_version = ?, command_id = COALESCE(?, command_id), revision = revision + 1, + updater = ? + WHERE id = ? AND tenant_id = ? AND deleted = FALSE AND kb_type = 'global' + """, version, context.commandId(), actor(context), kbId, DEFAULT_TENANT_ID); + if (updated == 0) { + throw new ServiceException(CONTRACT_NOT_FOUND, "全局知识库不存在"); + } + return kbId; + } + + private Long updateInstalledKnowledgeBase(OperationContext context) { + Long installId = firstPathLong(context, "installId"); + String status = switch (context.operationId()) { + case "disableInstalledKnowledgeBase" -> "disabled"; + case "restoreInstalledKnowledgeBase" -> "installed"; + default -> "deleted"; + }; + int updated = jdbcTemplate.update(""" + UPDATE muse_market_installation + SET status = ?, deleted = CASE WHEN ? = 'deleted' THEN TRUE ELSE FALSE END, + uninstalled_at = CASE WHEN ? IN ('disabled', 'deleted') THEN CURRENT_TIMESTAMP ELSE NULL END, + command_id = COALESCE(?, command_id), updater = ? + WHERE id = ? AND tenant_id = ? AND deleted = FALSE + AND (? <> 'app' OR user_id = ?) + """, status, status, status, context.commandId(), actor(context), installId, DEFAULT_TENANT_ID, + context.side(), context.actorUserId()); + if (updated == 0) { + throw new ServiceException(CONTRACT_NOT_FOUND, "已安装知识库不存在或无权访问"); + } + return installId; + } + + private Long upsertMarketFavorite(OperationContext context, boolean active) { + Long assetId = firstPathLong(context, "assetId"); + return jdbcTemplate.queryForObject(""" + INSERT INTO muse_market_favorite(asset_id, user_id, status, command_id, creator, updater, tenant_id) + VALUES (?, ?, ?, ?, ?, ?, ?) + ON CONFLICT (tenant_id, user_id, asset_id) + DO UPDATE SET status = EXCLUDED.status, + command_id = COALESCE(EXCLUDED.command_id, muse_market_favorite.command_id), + deleted = FALSE, + updater = EXCLUDED.updater + RETURNING id + """, Long.class, assetId, context.actorUserId(), active ? "active" : "inactive", + context.commandId(), actor(context), actor(context), DEFAULT_TENANT_ID); + } + + private Long insertMarketPurchase(OperationContext context) { + Long assetId = firstPathLong(context, "assetId"); + Long versionId = latestAssetVersionId(assetId); + return jdbcTemplate.queryForObject(""" + INSERT INTO muse_market_purchase(asset_id, asset_version_id, user_id, status, purchase_payload, + command_id, creator, updater, tenant_id) + VALUES (?, ?, ?, 'completed', CAST(? AS jsonb), ?, ?, ?, ?) + RETURNING id + """, Long.class, assetId, versionId, context.actorUserId(), toJson(context.body()), + context.commandId(), actor(context), actor(context), DEFAULT_TENANT_ID); + } + + private Long insertOperationRecord(OperationContext context, ResourceSpec spec, Long resourceId, Map response) { + return jdbcTemplate.queryForObject(""" + INSERT INTO muse_domain_operation_record(domain, side, operation_id, command_id, actor_user_id, + resource_type, resource_id, resource_key, parent_type, parent_id, status, revision, + path_variables, query_params, request_payload, response_payload, creator, updater, tenant_id) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, CAST(? AS jsonb), CAST(? AS jsonb), CAST(? AS jsonb), + CAST(? AS jsonb), ?, ?, ?) + RETURNING id + """, Long.class, + context.domain(), context.side(), context.operationId(), context.commandId(), context.actorUserId(), + spec == null ? "operation" : spec.resourceType(), resourceId, resourceKey(context, spec), + parentType(context, spec), parentId(context, spec), String.valueOf(response.get("status")), 1, + toJson(context.pathVariables()), toJson(context.queryParams()), toJson(context.body()), toJson(response), + actor(context), actor(context), DEFAULT_TENANT_ID); + } + + private void updateOperationResponse(Long operationRecordId, Long resourceId, Map response) { + jdbcTemplate.update(""" + UPDATE muse_domain_operation_record + SET resource_id = COALESCE(?, resource_id), response_payload = CAST(? AS jsonb) + WHERE id = ? AND tenant_id = ? + """, resourceId, toJson(response), operationRecordId, DEFAULT_TENANT_ID); + } + + private Map replayCommand(OperationContext context) { + if (StrUtil.isBlank(context.commandId())) { + return null; + } + String json = queryOptionalString(""" + SELECT response_payload::text + FROM muse_domain_operation_record + WHERE tenant_id = ? AND deleted = FALSE AND domain = ? AND command_id = ? + ORDER BY id DESC LIMIT 1 + """, DEFAULT_TENANT_ID, context.domain(), context.commandId()); + return json == null ? null : readMap(json); + } + + private void enforceExpectedGuards(OperationContext context, ResourceSpec spec) { + Long id = resolveRowId(context, spec); + if (id == null) { + return; + } + Object expectedStatus = firstBodyValue(context, "expectedStatus", "expectedActiveStatus"); + if (expectedStatus != null && spec.statusColumn() != null) { + String actualStatus = queryOptionalString("SELECT %s FROM %s WHERE id = ? AND tenant_id = ? AND deleted = FALSE" + .formatted(spec.statusColumn(), spec.tableName()), id, DEFAULT_TENANT_ID); + if (actualStatus != null && !Objects.equals(String.valueOf(expectedStatus), actualStatus)) { + throw new ServiceException(CONTRACT_CONFLICT, "状态冲突,当前状态为 " + actualStatus); + } + } + Object expectedRevision = firstBodyValue(context, "expectedRevision", "expectedVersion", "expectedWorkRevision", + "expectedDraftRevision", "expectedSlotRevision", "expectedBlockRevision", "expectedChapterRevision"); + if (expectedRevision != null && spec.revisionColumn() != null) { + Integer actualRevision = queryOptionalInteger("SELECT %s FROM %s WHERE id = ? AND tenant_id = ? AND deleted = FALSE" + .formatted(spec.revisionColumn(), spec.tableName()), id, DEFAULT_TENANT_ID); + Integer expected = intValue(expectedRevision, null); + if (actualRevision != null && expected != null && !Objects.equals(expected, actualRevision)) { + throw new ServiceException(CONTRACT_CONFLICT, "版本冲突,当前版本为 " + actualRevision); + } + } + enforceCrossResourceRevision(context); + } + + private void enforceCrossResourceRevision(OperationContext context) { + Object expectedWorkRevision = firstBodyValue(context, "expectedWorkRevision"); + if (expectedWorkRevision == null) { + return; + } + Long workId = firstPathLong(context, "workId"); + if (workId == null) { + return; + } + Integer actualRevision = queryOptionalInteger(""" + SELECT revision FROM muse_content_work + WHERE id = ? AND tenant_id = ? AND deleted = FALSE + AND (? <> 'app' OR owner_user_id = ?) + """, workId, DEFAULT_TENANT_ID, context.side(), context.actorUserId()); + Integer expected = intValue(expectedWorkRevision, null); + if (actualRevision != null && expected != null && !Objects.equals(expected, actualRevision)) { + throw new ServiceException(CONTRACT_CONFLICT, "作品版本冲突,当前版本为 " + actualRevision); + } + } + + private void enforceMetaSchemaActiveVersion(OperationContext context, String schemaKey, String expectedField) { + Object expectedActiveVersion = firstBodyValue(context, expectedField); + if (expectedActiveVersion == null) { + return; + } + Integer actualActiveVersion = queryOptionalInteger(""" + SELECT v.version_no + FROM muse_meta_schema s + LEFT JOIN muse_meta_schema_version v ON s.active_version_id = v.id AND v.tenant_id = s.tenant_id + WHERE s.tenant_id = ? AND s.schema_key = ? AND s.deleted = FALSE + LIMIT 1 + """, DEFAULT_TENANT_ID, schemaKey); + Integer expected = intValue(expectedActiveVersion, null); + if (actualActiveVersion != null && expected != null && !Objects.equals(expected, actualActiveVersion)) { + throw new ServiceException(CONTRACT_CONFLICT, "MetaSchema 活跃版本冲突,当前版本为 " + actualActiveVersion); + } + } + + private void enforcePathVersion(OperationContext context, String expectedField, String pathField, String message) { + Object expectedValue = firstBodyValue(context, expectedField); + Integer actualValue = firstPathInt(context, pathField); + Integer expected = intValue(expectedValue, null); + if (expected != null && actualValue != null && !Objects.equals(expected, actualValue)) { + throw new ServiceException(CONTRACT_CONFLICT, message + ",当前请求版本为 " + actualValue); + } + } + + private void enforceFunctionChainActiveVersion(OperationContext context) { + Object expectedActiveVersion = firstBodyValue(context, "expectedActiveVersion"); + if (expectedActiveVersion == null) { + return; + } + String actualActiveVersion = queryOptionalString(""" + SELECT active_version + FROM muse_meta_function_chain + WHERE tenant_id = ? AND chain_key = ? AND deleted = FALSE + LIMIT 1 + """, DEFAULT_TENANT_ID, firstPathString(context, "chainKey")); + if (StrUtil.isNotBlank(actualActiveVersion) + && !Objects.equals(String.valueOf(expectedActiveVersion), actualActiveVersion)) { + throw new ServiceException(CONTRACT_CONFLICT, "功能链活跃版本冲突,当前版本为 " + actualActiveVersion); + } + } + + private Long resolveRowId(OperationContext context, ResourceSpec spec) { + if (spec == null || spec.tableName() == null) { + return null; + } + if (spec.idPathVariable() != null) { + return firstPathLong(context, spec.idPathVariable()); + } + if (spec.keyColumn() != null && spec.keyPathVariable() != null) { + String key = firstPathString(context, spec.keyPathVariable()); + if (StrUtil.isNotBlank(key)) { + if (spec.parentColumn() != null && spec.parentPathVariable() != null) { + Long parentId = firstPathLong(context, spec.parentPathVariable()); + if (parentId != null) { + return queryOptionalLong(""" + SELECT id FROM %s + WHERE tenant_id = ? AND deleted = FALSE AND %s = ? AND %s = ? + LIMIT 1 + """.formatted(spec.tableName(), spec.keyColumn(), spec.parentColumn()), + DEFAULT_TENANT_ID, key, parentId); + } + } + return queryOptionalLong("SELECT id FROM %s WHERE tenant_id = ? AND deleted = FALSE AND %s = ? LIMIT 1" + .formatted(spec.tableName(), spec.keyColumn()), DEFAULT_TENANT_ID, key); + } + } + return null; + } + + private WhereClause buildWhere(OperationContext context, ResourceSpec spec, boolean list) { + StringBuilder sql = new StringBuilder("tenant_id = ? AND deleted = FALSE"); + List args = new ArrayList<>(); + args.add(DEFAULT_TENANT_ID); + if (!list && spec.idPathVariable() != null) { + Long id = firstPathLong(context, spec.idPathVariable()); + if (id != null) { + sql.append(" AND ").append(spec.idColumn()).append(" = ?"); + args.add(id); + } + } + if (spec.keyColumn() != null && spec.keyPathVariable() != null) { + String key = firstPathString(context, spec.keyPathVariable()); + if (StrUtil.isNotBlank(key)) { + sql.append(" AND ").append(spec.keyColumn()).append(" = ?"); + args.add(key); + } + } + if (spec.ownerColumn() != null) { + Long owner = firstPathLong(context, spec.ownerPathVariable()); + if (owner == null && "app".equals(context.side())) { + owner = context.actorUserId(); + } + if (owner != null) { + sql.append(" AND ").append(spec.ownerColumn()).append(" = ?"); + args.add(owner); + } + } + if (spec.parentColumn() != null && spec.parentPathVariable() != null) { + Long parentId = firstPathLong(context, spec.parentPathVariable()); + if (parentId != null) { + sql.append(" AND ").append(spec.parentColumn()).append(" = ?"); + args.add(parentId); + } + } + if ("muse_domain_workflow_task".equals(spec.tableName())) { + sql.append(" AND domain = ?"); + args.add(context.domain()); + } + return new WhereClause(sql.toString(), args); + } + + private ResourceSpec resolveSpec(OperationContext context) { + String op = context.operationId(); + return switch (context.domain()) { + case "account" -> resolveAccountSpec(op); + case "meta" -> resolveMetaSpec(op); + case "ai" -> resolveAiSpec(op); + case "knowledge" -> resolveKnowledgeSpec(op); + case "market" -> resolveMarketSpec(op); + default -> operationSpec(); + }; + } + + private ResourceSpec resolveAccountSpec(String op) { + if (op.contains("Entitlement")) { + return new ResourceSpec("entitlement", "muse_member_entitlement", "id", null, null, null, + "account_user_id", "userId", null, null, "status", "revision"); + } + if (op.contains("QuotaAdjustment")) { + return new ResourceSpec("quotaAdjustment", "muse_member_entitlement_audit_log", "id", null, null, null, + "account_user_id", "userId", null, null, "change_type", null); + } + if (op.contains("QuotaRequest")) { + return new ResourceSpec("quotaRequest", "muse_domain_workflow_task", "id", "requestId", null, null, + "owner_user_id", "userId", null, null, "status", "revision"); + } + if (op.contains("Usage")) { + return new ResourceSpec("usageRecord", "muse_member_usage_record", "id", null, "correlation_id", "correlationId", + "account_user_id", "userId", null, null, "status", null); + } + if (op.contains("AttributionJob") || op.contains("ExportTask") || op.contains("Download")) { + return new ResourceSpec("accountWorkflowTask", "muse_domain_workflow_task", "id", "taskId", null, null, + "owner_user_id", "userId", null, null, "status", "revision"); + } + if (op.contains("SecurityEvent")) { + return new ResourceSpec("securityEvent", "muse_member_security_event", "id", "eventId", null, null, + "account_user_id", "userId", null, null, "severity", null); + } + if (op.contains("NewApiBinding")) { + return new ResourceSpec("newApiBinding", "muse_member_new_api_binding", "id", null, null, null, + "account_user_id", "userId", null, null, "binding_status", null); + } + return operationSpec(); + } + + private ResourceSpec resolveMetaSpec(String op) { + if (op.contains("MetaSchemaVersion") || op.contains("MetaSchemaDraft") || op.contains("MetaSchema")) { + return new ResourceSpec("metaSchema", "muse_meta_schema", "id", null, "schema_key", "schemaKey", + null, null, null, null, "status", null); + } + if (op.contains("ProtectionNode")) { + return new ResourceSpec("protectionNode", "muse_meta_protection_node", "id", null, "node_key", "nodeKey", + null, null, null, null, "status", null); + } + if (op.contains("FunctionChain")) { + return new ResourceSpec("functionChain", "muse_meta_function_chain", "id", null, "chain_key", "chainKey", + null, null, null, null, "status", null); + } + return operationSpec(); + } + + private ResourceSpec resolveAiSpec(String op) { + if (op.contains("Prompt")) { + return new ResourceSpec("prompt", "muse_prompt", "id", null, "prompt_key", "promptKey", + null, null, null, null, "status", null); + } + if ("precheckAgentSlot".equals(op)) { + return new ResourceSpec("aiWorkflowTask", "muse_domain_workflow_task", "id", "taskId", null, null, + "owner_user_id", null, null, null, "status", "revision"); + } + if (op.contains("AgentSlot")) { + return new ResourceSpec("agentSlotBinding", "muse_agent_slot_binding", "id", null, "slot_key", "slotKey", + null, null, "work_id", "workId", "status", "revision"); + } + if (op.contains("Agent")) { + return new ResourceSpec("agent", "muse_agent", "id", "agentId", null, null, + "owner_user_id", null, null, null, "status", null); + } + if (op.contains("Suggestion")) { + return new ResourceSpec("suggestion", "muse_ai_suggestion", "id", "suggestionId", null, null, + null, null, "work_id", "workId", "status", "source_revision"); + } + if (op.contains("QualityPolicy")) { + return new ResourceSpec("qualityPolicy", "muse_quality_policy", "id", null, "policy_key", "policyKey", + null, null, null, null, null, null); + } + if (op.contains("ToolGrant")) { + return new ResourceSpec("toolGrant", "muse_tool_grant", "id", null, "grant_key", "grantKey", + null, null, null, null, "status", null); + } + if (op.contains("EvaluationRun") || op.contains("SourceEvent") || op.contains("SourceStatus") + || op.contains("AgentSlot")) { + return new ResourceSpec("aiWorkflowTask", "muse_domain_workflow_task", "id", "taskId", null, null, + "owner_user_id", null, null, null, "status", "revision"); + } + if (op.contains("AiTask") || op.contains("Task") || op.contains("Job")) { + return new ResourceSpec("aiTask", "muse_ai_generation", "id", "taskId", null, null, + null, null, "work_id", "workId", "status", null); + } + return operationSpec(); + } + + private ResourceSpec resolveKnowledgeSpec(String op) { + if (op.contains("DocumentVersion") || op.contains("KBVersions")) { + return new ResourceSpec("knowledgeDocumentVersion", "muse_knowledge_document_version", "id", null, null, null, + null, null, "document_id", "documentId", "processing_status", null); + } + if (op.contains("Document")) { + return new ResourceSpec("knowledgeDocument", "muse_knowledge_document", "id", "documentId", null, null, + null, null, "kb_id", "kbId", "scan_status", null); + } + if (op.contains("Entity")) { + return new ResourceSpec("knowledgeEntity", "muse_knowledge_entity", "id", "entityId", null, null, + null, null, "work_id", "workId", "status", "revision"); + } + if (op.contains("Relation")) { + return new ResourceSpec("knowledgeRelation", "muse_knowledge_relation", "id", "relationId", null, null, + null, null, "work_id", "workId", null, "revision"); + } + if (op.contains("KnowledgeDraft")) { + return new ResourceSpec("knowledgeDraft", "muse_knowledge_draft", "id", "draftId", null, null, + null, null, "work_id", "workId", "status", "revision"); + } + if ("createKnowledgeBindingPrecheck".equals(op)) { + return new ResourceSpec("knowledgeWorkflowTask", "muse_domain_workflow_task", "id", "taskId", null, null, + "owner_user_id", null, null, null, "status", "revision"); + } + if (op.contains("KnowledgeBinding")) { + return new ResourceSpec("knowledgeBinding", "muse_knowledge_binding", "id", "bindingId", null, null, + null, null, "work_id", "workId", "binding_status", "revision"); + } + if (op.contains("AccessPolicy") || op.contains("ProcessingTask") || op.contains("SourceEvent") + || op.contains("PublishReadiness") || op.contains("PublishSnapshot") || op.contains("ExportTask") + || "reindexGlobalKnowledgeBase".equals(op) || "reindexKnowledgeBase".equals(op) + || op.contains("Precheck")) { + return new ResourceSpec("knowledgeWorkflowTask", "muse_domain_workflow_task", "id", "taskId", null, null, + "owner_user_id", null, null, null, "status", "revision"); + } + if (op.contains("InstalledKnowledgeBase")) { + return new ResourceSpec("installedKnowledgeBase", "muse_market_installation", "id", "installId", null, null, + "user_id", null, null, null, "status", null); + } + if (op.contains("KnowledgeBase") || op.contains("LocalKnowledge")) { + return new ResourceSpec("knowledgeBase", "muse_knowledge_base", "id", "kbId", null, null, + "owner_user_id", null, null, null, "status", "revision"); + } + return operationSpec(); + } + + private ResourceSpec resolveMarketSpec(String op) { + if (op.contains("Handoff")) { + return new ResourceSpec("handoff", "muse_market_handoff", "id", null, "token", "handoffToken", + null, null, null, null, "status", null); + } + if (op.contains("Appeal")) { + return new ResourceSpec("appeal", "muse_market_appeal", "id", "appealId", null, null, + "user_id", null, "asset_id", "assetId", "status", "revision"); + } + if (op.contains("PublishRequest") || op.contains("PublishRecord")) { + return new ResourceSpec("publishRequest", "muse_market_publish_request", "id", "requestId", null, null, + "publisher_id", null, "asset_id", "assetId", "status", "revision"); + } + if (op.contains("Install") || op.contains("Installed")) { + return new ResourceSpec("installation", "muse_market_installation", "id", "installId", null, null, + "user_id", null, "asset_id", "assetId", "status", null); + } + if ("favoriteAsset".equals(op) || "unfavoriteAsset".equals(op)) { + return new ResourceSpec("favorite", "muse_market_favorite", "id", null, null, null, + "user_id", null, "asset_id", "assetId", "status", null); + } + if ("purchaseAsset".equals(op)) { + return new ResourceSpec("purchase", "muse_market_purchase", "id", null, null, null, + "user_id", null, "asset_id", "assetId", "status", null); + } + if ("createBindPrecheck".equals(op)) { + return new ResourceSpec("marketWorkflowTask", "muse_domain_workflow_task", "id", "taskId", null, null, + "owner_user_id", null, null, null, "status", "revision"); + } + if (op.contains("Asset") || op.contains("Marketplace") || op.contains("Category") || op.contains("Recommendation")) { + return new ResourceSpec("marketAsset", "muse_market_asset", "id", "assetId", null, null, + "publisher_id", null, null, null, "listing_status", "revision"); + } + return operationSpec(); + } + + private ResourceSpec operationSpec() { + return new ResourceSpec("operation", null, "id", null, null, null, null, null, null, null, "status", "revision"); + } + + private Long insertAccountAudit(OperationContext context, Long accountUserId, String changeType) { + return jdbcTemplate.queryForObject(""" + INSERT INTO muse_member_entitlement_audit_log(account_user_id, change_type, source_owner, source_id, + idempotency_key, delta_value_snapshot, reason_code, reason_message, operator_user_id, + creator, updater, tenant_id) + VALUES (?, ?, ?, ?, ?, CAST(? AS jsonb), ?, ?, ?, ?, ?, ?) + RETURNING id + """, Long.class, accountUserId, changeType, context.domain(), null, context.commandId(), + toJson(context.body()), text(context.body(), "reasonCode", "manual"), + text(context.body(), "reason", "quota adjustment"), context.actorUserId(), actor(context), actor(context), + DEFAULT_TENANT_ID); + } + + @SuppressWarnings("unchecked") + private void applyQuotaAdjustments(OperationContext context, Long accountUserId) { + Object adjustments = context.body().get("adjustments"); + if (!(adjustments instanceof Collection collection)) { + return; + } + for (Object item : collection) { + if (!(item instanceof Map raw)) { + continue; + } + Map adjustment = (Map) raw; + String quotaType = text(adjustment, "quotaType", "ai_call"); + Long delta = longValue(adjustment.getOrDefault("delta", adjustment.get("amount"))); + if (delta == null) { + delta = 0L; + } + jdbcTemplate.update(""" + INSERT INTO muse_member_quota(account_user_id, quota_type, total_amount, used_amount, creator, updater, tenant_id) + VALUES (?, ?, ?, 0, ?, ?, ?) + ON CONFLICT (tenant_id, account_user_id, quota_type) + DO UPDATE SET total_amount = muse_member_quota.total_amount + EXCLUDED.total_amount, + revision = muse_member_quota.revision + 1, + updater = EXCLUDED.updater + """, accountUserId, quotaType, delta, actor(context), actor(context), DEFAULT_TENANT_ID); + } + } + + private Long upsertNewApiBinding(OperationContext context, Long accountUserId) { + String newApiUserId = text(context.body(), "newapiUserId", "muse-" + accountUserId); + return jdbcTemplate.queryForObject(""" + INSERT INTO muse_member_new_api_binding(account_user_id, newapi_user_id, binding_status, command_id, + last_sync_at, creator, updater, tenant_id) + VALUES (?, ?, 'active', ?, CURRENT_TIMESTAMP, ?, ?, ?) + ON CONFLICT (tenant_id, account_user_id) + DO UPDATE SET newapi_user_id = EXCLUDED.newapi_user_id, + binding_status = 'active', + command_id = EXCLUDED.command_id, + last_sync_at = CURRENT_TIMESTAMP, + updater = EXCLUDED.updater + RETURNING id + """, Long.class, accountUserId, newApiUserId, context.commandId(), actor(context), actor(context), + DEFAULT_TENANT_ID); + } + + private Map selectAccountEntitlement(Long accountUserId) { + String json = queryOptionalString(""" + SELECT to_jsonb(t)::text FROM muse_member_entitlement t + WHERE tenant_id = ? AND account_user_id = ? AND deleted = FALSE + LIMIT 1 + """, DEFAULT_TENANT_ID, accountUserId); + return json == null ? Map.of() : readMap(json); + } + + private Long ensureMetaSchema(OperationContext context, String schemaKey) { + return jdbcTemplate.queryForObject(""" + INSERT INTO muse_meta_schema(schema_key, domain, scope, target_type, display_name, effective_scope, + creator, updater, tenant_id) + VALUES (?, ?, ?, ?, ?, CAST(? AS jsonb), ?, ?, ?) + ON CONFLICT (tenant_id, domain, scope, target_type, schema_key) + DO UPDATE SET display_name = EXCLUDED.display_name, + effective_scope = EXCLUDED.effective_scope, + projection_version = muse_meta_schema.projection_version + 1, + updater = EXCLUDED.updater + RETURNING id + """, Long.class, schemaKey, text(context.body(), "domain", "content"), + text(context.body(), "scope", "work"), text(context.body(), "targetType", "work"), + text(context.body(), "displayName", schemaKey), toJson(context.body().get("effectiveScope")), + actor(context), actor(context), DEFAULT_TENANT_ID); + } + + private Long insertMetaSchemaVersion(OperationContext context, Long schemaId, String schemaKey, Integer versionNo, String status) { + return jdbcTemplate.queryForObject(""" + INSERT INTO muse_meta_schema_version(schema_id, schema_key, version_no, status, effective_scope, + field_contract_snapshot, command_id, change_note, + creator, updater, tenant_id) + VALUES (?, ?, ?, ?, CAST(? AS jsonb), CAST(? AS jsonb), ?, ?, ?, ?, ?) + RETURNING id + """, Long.class, schemaId, schemaKey, versionNo, status, toJson(context.body().get("effectiveScope")), + toJson(context.body().getOrDefault("fields", context.body())), context.commandId(), + text(context.body(), "reason", text(context.body(), "changeNote", "draft")), + actor(context), actor(context), DEFAULT_TENANT_ID); + } + + private Long updateMetaSchemaVersionStatus(OperationContext context, String schemaKey, Integer versionNo, String status) { + Long versionId = requireMetaSchemaVersion(schemaKey, versionNo); + jdbcTemplate.update(""" + UPDATE muse_meta_schema_version + SET status = ?, active_flag = ?, command_id = ?, updater = ?, + published_at = CASE WHEN ? = 'published' THEN CURRENT_TIMESTAMP ELSE published_at END, + activated_at = CASE WHEN ? = 'active' THEN CURRENT_TIMESTAMP ELSE activated_at END + WHERE id = ? AND tenant_id = ? AND deleted = FALSE + """, status, "active".equals(status), context.commandId(), actor(context), status, status, versionId, + DEFAULT_TENANT_ID); + return versionId; + } + + private Long requireMetaSchemaVersion(String schemaKey, Integer versionNo) { + Long id = queryOptionalLong(""" + SELECT id FROM muse_meta_schema_version + WHERE tenant_id = ? AND schema_key = ? AND version_no = ? AND deleted = FALSE + """, DEFAULT_TENANT_ID, schemaKey, versionNo); + if (id == null) { + throw new ServiceException(CONTRACT_NOT_FOUND, "MetaSchema 版本不存在"); + } + return id; + } + + private Long activateFunctionChainVersion(OperationContext context) { + String chainKey = firstPathString(context, "chainKey"); + String version = firstPathString(context, "version"); + return jdbcTemplate.queryForObject(""" + INSERT INTO muse_meta_function_chain(chain_key, active_version, chain_snapshot, creator, updater, tenant_id) + VALUES (?, ?, CAST(? AS jsonb), ?, ?, ?) + ON CONFLICT (tenant_id, chain_key) + DO UPDATE SET active_version = EXCLUDED.active_version, + chain_snapshot = EXCLUDED.chain_snapshot, + updater = EXCLUDED.updater + RETURNING id + """, Long.class, chainKey, version, toJson(context.body()), actor(context), actor(context), + DEFAULT_TENANT_ID); + } + + private Long insertAgent(OperationContext context, String agentType) { + String name = text(context.body(), "name", "agent"); + String agentKey = text(context.body(), "agentKey", slug(name) + "-" + UUID.randomUUID()); + return jdbcTemplate.queryForObject(""" + INSERT INTO muse_agent(agent_key, name, description, agent_type, owner_user_id, prompt_key, slot_bindings, + category, tags, creator, updater, tenant_id) + VALUES (?, ?, ?, ?, ?, ?, CAST(? AS jsonb), ?, CAST(? AS jsonb), ?, ?, ?) + RETURNING id + """, Long.class, agentKey, name, text(context.body(), "description", null), agentType, + "user".equals(agentType) ? context.actorUserId() : null, text(context.body(), "promptKey", null), + toJson(context.body().get("slotBindings")), text(context.body(), "category", null), + toJson(context.body().get("tags")), actor(context), actor(context), DEFAULT_TENANT_ID); + } + + private Long insertAgentVersion(OperationContext context, Long agentId) { + String version = text(context.body(), "version", String.valueOf(nextInt( + "SELECT COALESCE(MAX(CAST(version AS INT)), 0) + 1 FROM muse_agent_version WHERE tenant_id = ? AND agent_id = ? AND version ~ '^[0-9]+$'", + DEFAULT_TENANT_ID, agentId))); + return jdbcTemplate.queryForObject(""" + INSERT INTO muse_agent_version(agent_id, version, config, change_note, status, command_id, + creator, updater, tenant_id) + VALUES (?, ?, CAST(? AS jsonb), ?, 'draft', ?, ?, ?, ?) + RETURNING id + """, Long.class, agentId, version, toJson(context.body()), + text(context.body(), "changeNote", "create version"), context.commandId(), actor(context), actor(context), + DEFAULT_TENANT_ID); + } + + private Long insertPromptVersion(OperationContext context, String promptKey) { + Long promptId = jdbcTemplate.queryForObject(""" + INSERT INTO muse_prompt(prompt_key, name, description, content, creator, updater, tenant_id) + VALUES (?, ?, ?, ?, ?, ?, ?) + ON CONFLICT (tenant_id, prompt_key) + DO UPDATE SET name = EXCLUDED.name, description = EXCLUDED.description, updater = EXCLUDED.updater + RETURNING id + """, Long.class, promptKey, text(context.body(), "name", promptKey), + text(context.body(), "description", null), text(context.body(), "content", null), + actor(context), actor(context), DEFAULT_TENANT_ID); + String version = text(context.body(), "version", String.valueOf(nextInt(""" + SELECT COALESCE(MAX(CAST(version AS INT)), 0) + 1 + FROM muse_prompt_version + WHERE tenant_id = ? AND prompt_key = ? AND version ~ '^[0-9]+$' + """, DEFAULT_TENANT_ID, promptKey))); + return jdbcTemplate.queryForObject(""" + INSERT INTO muse_prompt_version(prompt_id, prompt_key, version, content, variables, change_note, + status, command_id, creator, updater, tenant_id) + VALUES (?, ?, ?, ?, CAST(? AS jsonb), ?, 'draft', ?, ?, ?, ?) + RETURNING id + """, Long.class, promptId, promptKey, version, text(context.body(), "content", ""), + toJson(context.body().get("variables")), text(context.body(), "changeNote", "create version"), + context.commandId(), actor(context), actor(context), DEFAULT_TENANT_ID); + } + + private Long activatePromptVersion(OperationContext context) { + String promptKey = firstPathString(context, "promptKey"); + String version = firstPathString(context, "version"); + Long versionId = queryOptionalLong(""" + SELECT id FROM muse_prompt_version + WHERE tenant_id = ? AND prompt_key = ? AND version = ? AND deleted = FALSE + """, DEFAULT_TENANT_ID, promptKey, version); + if (versionId == null) { + throw new ServiceException(CONTRACT_NOT_FOUND, "Prompt 版本不存在"); + } + jdbcTemplate.update(""" + UPDATE muse_prompt_version SET status = 'active', updater = ? + WHERE id = ? AND tenant_id = ? + """, actor(context), versionId, DEFAULT_TENANT_ID); + jdbcTemplate.update(""" + UPDATE muse_prompt SET active_version = ?, updater = ? + WHERE tenant_id = ? AND prompt_key = ? AND deleted = FALSE + """, version, actor(context), DEFAULT_TENANT_ID, promptKey); + return versionId; + } + + private Long insertQualityPolicyVersion(OperationContext context, String policyKey) { + Long policyId = jdbcTemplate.queryForObject(""" + INSERT INTO muse_quality_policy(policy_key, name, description, task_types, creator, updater, tenant_id) + VALUES (?, ?, ?, CAST(? AS jsonb), ?, ?, ?) + ON CONFLICT (tenant_id, policy_key) + DO UPDATE SET name = EXCLUDED.name, description = EXCLUDED.description, updater = EXCLUDED.updater + RETURNING id + """, Long.class, policyKey, text(context.body(), "name", policyKey), + text(context.body(), "description", null), toJson(context.body().get("taskTypes")), + actor(context), actor(context), DEFAULT_TENANT_ID); + Integer versionNo = nextInt("SELECT COALESCE(MAX(version_no), 0) + 1 FROM muse_quality_policy_version WHERE tenant_id = ? AND policy_key = ?", + DEFAULT_TENANT_ID, policyKey); + return jdbcTemplate.queryForObject(""" + INSERT INTO muse_quality_policy_version(policy_id, policy_key, version_no, effective_scope, + metric_contract_snapshot, threshold_snapshot, rewrite_policy_snapshot, command_id, + creator, updater, tenant_id) + VALUES (?, ?, ?, CAST(? AS jsonb), CAST(? AS jsonb), CAST(? AS jsonb), CAST(? AS jsonb), ?, ?, ?, ?) + RETURNING id + """, Long.class, policyId, policyKey, versionNo, toJson(context.body().get("effectiveScope")), + toJson(context.body().get("dimensions")), toJson(context.body().get("thresholds")), + toJson(context.body().get("rewritePolicy")), context.commandId(), actor(context), actor(context), + DEFAULT_TENANT_ID); + } + + private Long insertToolGrant(OperationContext context) { + String grantKey = text(context.body(), "grantKey", "grant-" + UUID.randomUUID()); + return jdbcTemplate.queryForObject(""" + INSERT INTO muse_tool_grant(grant_key, agent_id, tool_name, action, purpose, scope, budget, + outbound_policy, approval_status, status, command_id, creator, updater, tenant_id) + VALUES (?, ?, ?, ?, ?, CAST(? AS jsonb), CAST(? AS jsonb), CAST(? AS jsonb), 'approved', 'active', + ?, ?, ?, ?) + RETURNING id + """, Long.class, grantKey, longValue(context.body().get("agentId")), + text(context.body(), "toolType", text(context.body(), "toolName", "tool")), + text(context.body(), "action", "use"), text(context.body(), "reason", null), + toJson(context.body().get("scope")), toJson(context.body().get("budget")), + toJson(context.body().get("outboundPolicy")), context.commandId(), actor(context), actor(context), + DEFAULT_TENANT_ID); + } + + private Long upsertAgentSlotBinding(OperationContext context) { + Long workId = firstPathLong(context, "workId"); + requireAppWorkOwner(context, workId); + String slotKey = firstPathString(context, "slotKey"); + Long currentId = queryOptionalLong(""" + SELECT id FROM muse_agent_slot_binding + WHERE tenant_id = ? AND work_id = ? AND slot_key = ? AND status = 'active' AND deleted = FALSE + ORDER BY id DESC LIMIT 1 + """, DEFAULT_TENANT_ID, workId, slotKey); + if (currentId != null) { + jdbcTemplate.update(""" + UPDATE muse_agent_slot_binding + SET status = 'replaced', updater = ? + WHERE id = ? AND tenant_id = ? + """, actor(context), currentId, DEFAULT_TENANT_ID); + } + return jdbcTemplate.queryForObject(""" + INSERT INTO muse_agent_slot_binding(work_id, slot_key, agent_id, agent_version, authorization_snapshot_id, + source_snapshot_id, binding_source, command_id, creator, updater, tenant_id) + VALUES (?, ?, ?, ?, ?, ?, 'direct', ?, ?, ?, ?) + RETURNING id + """, Long.class, workId, slotKey, longValue(context.body().get("sourceAgentId")), + String.valueOf(context.body().getOrDefault("sourceAgentVersion", "1")), + longValue(context.body().get("authorizationSnapshotId")), longValue(context.body().get("sourceSnapshotId")), + context.commandId(), actor(context), actor(context), DEFAULT_TENANT_ID); + } + + private Long insertKnowledgeBase(OperationContext context, String type) { + return jdbcTemplate.queryForObject(""" + INSERT INTO muse_knowledge_base(name, description, kb_type, owner_user_id, visibility_policy, + command_id, creator, updater, tenant_id) + VALUES (?, ?, ?, ?, CAST(? AS jsonb), ?, ?, ?, ?) + RETURNING id + """, Long.class, text(context.body(), "name", "knowledge base"), + text(context.body(), "description", null), type, "global".equals(type) ? null : context.actorUserId(), + toJson(context.body().get("visibilityPolicy")), context.commandId(), actor(context), actor(context), + DEFAULT_TENANT_ID); + } + + private Long updateKnowledgeBase(OperationContext context) { + Long kbId = firstPathLong(context, "kbId"); + String status = context.operationId().contains("disable") ? "disabled" + : context.operationId().contains("restore") || context.operationId().contains("enable") ? "active" + : text(context.body(), "status", null); + int updated = jdbcTemplate.update(""" + UPDATE muse_knowledge_base + SET name = COALESCE(?, name), description = COALESCE(?, description), + status = COALESCE(?, status), command_id = COALESCE(?, command_id), + revision = revision + 1, updater = ? + WHERE id = ? AND tenant_id = ? AND deleted = FALSE + AND (? <> 'app' OR owner_user_id = ?) + """, text(context.body(), "name", null), text(context.body(), "description", null), status, + context.commandId(), actor(context), kbId, DEFAULT_TENANT_ID, context.side(), context.actorUserId()); + if (updated == 0) { + throw new ServiceException(CONTRACT_NOT_FOUND, "知识库不存在"); + } + return kbId; + } + + private Long deleteKnowledgeBase(OperationContext context) { + Long kbId = firstPathLong(context, "kbId"); + int updated = jdbcTemplate.update(""" + UPDATE muse_knowledge_base + SET status = 'deleted', deleted = TRUE, command_id = COALESCE(?, command_id), + revision = revision + 1, updater = ? + WHERE id = ? AND tenant_id = ? AND deleted = FALSE + AND (? <> 'app' OR owner_user_id = ?) + """, context.commandId(), actor(context), kbId, DEFAULT_TENANT_ID, + context.side(), context.actorUserId()); + if (updated == 0) { + throw new ServiceException(CONTRACT_NOT_FOUND, "知识库不存在"); + } + return kbId; + } + + private Long insertKnowledgeDocument(OperationContext context) { + Long kbId = firstPathLong(context, "kbId"); + requireKnowledgeBaseOwner(context, kbId); + return jdbcTemplate.queryForObject(""" + INSERT INTO muse_knowledge_document(kb_id, title, file_name, file_size, mime_type, file_hash, + storage_ref, tags, command_id, creator, updater, tenant_id) + VALUES (?, ?, ?, ?, ?, ?, ?, CAST(? AS jsonb), ?, ?, ?, ?) + RETURNING id + """, Long.class, kbId, text(context.body(), "title", text(context.body(), "fileName", "document")), + text(context.body(), "fileName", null), longValue(context.body().get("fileSize")), + text(context.body(), "mimeType", null), text(context.body(), "fileHash", sha256(toJson(context.body()))), + text(context.body(), "storageRef", null), toJson(context.body().get("tags")), context.commandId(), + actor(context), actor(context), DEFAULT_TENANT_ID); + } + + private Long insertKnowledgeDocumentVersion(OperationContext context) { + Long documentId = firstPathLong(context, "documentId"); + requireKnowledgeDocumentOwner(context, documentId); + Integer version = nextInt("SELECT COALESCE(MAX(version), 0) + 1 FROM muse_knowledge_document_version WHERE tenant_id = ? AND document_id = ?", + DEFAULT_TENANT_ID, documentId); + return jdbcTemplate.queryForObject(""" + INSERT INTO muse_knowledge_document_version(document_id, version, file_hash, storage_ref, content_hash, + change_summary, change_reason, command_id, creator, updater, tenant_id) + VALUES (?, ?, ?, ?, ?, CAST(? AS jsonb), ?, ?, ?, ?, ?) + RETURNING id + """, Long.class, documentId, version, text(context.body(), "fileHash", sha256(toJson(context.body()))), + text(context.body(), "storageRef", null), text(context.body(), "contentHash", null), + toJson(context.body().get("changeSummary")), text(context.body(), "reason", null), context.commandId(), + actor(context), actor(context), DEFAULT_TENANT_ID); + } + + private Long deleteKnowledgeDocument(OperationContext context) { + Long documentId = firstPathLong(context, "documentId"); + int updated = jdbcTemplate.update(""" + UPDATE muse_knowledge_document + SET deleted = TRUE, command_id = COALESCE(?, command_id), updater = ? + WHERE id = ? AND tenant_id = ? AND deleted = FALSE + AND (? <> 'app' OR EXISTS ( + SELECT 1 FROM muse_knowledge_base kb + WHERE kb.id = muse_knowledge_document.kb_id AND kb.tenant_id = muse_knowledge_document.tenant_id + AND kb.owner_user_id = ? AND kb.deleted = FALSE + )) + """, context.commandId(), actor(context), documentId, DEFAULT_TENANT_ID, + context.side(), context.actorUserId()); + if (updated == 0) { + throw new ServiceException(CONTRACT_NOT_FOUND, "知识库文档不存在"); + } + return documentId; + } + + private Long updateKnowledgeEntity(OperationContext context) { + Long entityId = firstPathLong(context, "entityId"); + jdbcTemplate.update(""" + UPDATE muse_knowledge_entity + SET description = COALESCE(?, description), attributes = COALESCE(CAST(? AS jsonb), attributes), + command_id = ?, revision = revision + 1, updater = ? + WHERE id = ? AND tenant_id = ? AND deleted = FALSE + AND (? <> 'app' OR EXISTS ( + SELECT 1 FROM muse_content_work w + WHERE w.id = muse_knowledge_entity.work_id AND w.tenant_id = muse_knowledge_entity.tenant_id + AND w.owner_user_id = ? AND w.deleted = FALSE + )) + """, text(context.body(), "description", null), toJson(context.body().get("attributes")), + context.commandId(), actor(context), entityId, DEFAULT_TENANT_ID, context.side(), context.actorUserId()); + return entityId; + } + + private Long insertKnowledgeRelation(OperationContext context) { + requireAppWorkOwner(context, firstPathLong(context, "workId")); + return jdbcTemplate.queryForObject(""" + INSERT INTO muse_knowledge_relation(work_id, source_entity_id, target_entity_id, relation_type, + description, attributes, command_id, creator, updater, tenant_id) + VALUES (?, ?, ?, ?, ?, CAST(? AS jsonb), ?, ?, ?, ?) + RETURNING id + """, Long.class, firstPathLong(context, "workId"), longValue(context.body().get("sourceEntityId")), + longValue(context.body().get("targetEntityId")), text(context.body(), "relationType", "related_to"), + text(context.body(), "description", null), toJson(context.body().get("attributes")), context.commandId(), + actor(context), actor(context), DEFAULT_TENANT_ID); + } + + private Long updateKnowledgeDraft(OperationContext context) { + Long draftId = firstPathLong(context, "draftId"); + String status = switch (context.operationId()) { + case "confirmKnowledgeDraft" -> "confirmed"; + case "ignoreKnowledgeDraft" -> "ignored"; + default -> "pending"; + }; + jdbcTemplate.update(""" + UPDATE muse_knowledge_draft + SET status = ?, command_id = ?, revision = revision + 1, updater = ? + WHERE id = ? AND tenant_id = ? AND deleted = FALSE + AND (? <> 'app' OR EXISTS ( + SELECT 1 FROM muse_content_work w + WHERE w.id = muse_knowledge_draft.work_id AND w.tenant_id = muse_knowledge_draft.tenant_id + AND w.owner_user_id = ? AND w.deleted = FALSE + )) + """, status, context.commandId(), actor(context), draftId, DEFAULT_TENANT_ID, + context.side(), context.actorUserId()); + return draftId; + } + + private Long insertKnowledgeBinding(OperationContext context) { + requireAppWorkOwner(context, firstPathLong(context, "workId")); + return jdbcTemplate.queryForObject(""" + INSERT INTO muse_knowledge_binding(work_id, kb_id, binding_type, binding_scope, source_snapshot_id, + authorization_snapshot_id, source_version, command_id, creator, updater, tenant_id) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + ON CONFLICT (tenant_id, work_id, kb_id) + DO UPDATE SET binding_status = 'active', + source_snapshot_id = EXCLUDED.source_snapshot_id, + authorization_snapshot_id = EXCLUDED.authorization_snapshot_id, + source_version = EXCLUDED.source_version, + command_id = EXCLUDED.command_id, + revision = muse_knowledge_binding.revision + 1, + updater = EXCLUDED.updater + RETURNING id + """, Long.class, firstPathLong(context, "workId"), + longValue(context.body().getOrDefault("kbId", context.body().get("sourceId"))), + text(context.body(), "bindingType", "read"), text(context.body(), "bindingScope", "read"), + longValue(context.body().get("sourceSnapshotId")), longValue(context.body().get("authorizationSnapshotId")), + intValue(context.body().get("sourceVersion"), null), context.commandId(), actor(context), actor(context), + DEFAULT_TENANT_ID); + } + + private Long insertMarketAsset(OperationContext context, String status) { + return jdbcTemplate.queryForObject(""" + INSERT INTO muse_market_asset(name, description, asset_type, category, source_id, publisher_id, + listing_status, status, license_type, tags, command_id, creator, updater, tenant_id) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, CAST(? AS jsonb), ?, ?, ?, ?) + RETURNING id + """, Long.class, text(context.body(), "name", "market asset"), + text(context.body(), "description", null), text(context.body(), "assetType", "knowledge_base"), + text(context.body(), "category", null), longValue(context.body().get("sourceId")), context.actorUserId(), + status, status, text(context.body(), "licenseType", "standard"), toJson(context.body().get("tags")), + context.commandId(), actor(context), actor(context), DEFAULT_TENANT_ID); + } + + private Long insertMarketAssetVersion(OperationContext context, Long assetId, String status) { + return jdbcTemplate.queryForObject(""" + INSERT INTO muse_market_asset_version(asset_id, version, change_note, version_snapshot, status, + command_id, creator, updater, tenant_id) + VALUES (?, ?, ?, CAST(? AS jsonb), ?, ?, ?, ?, ?) + RETURNING id + """, Long.class, assetId, text(context.body(), "version", "1"), + text(context.body(), "changeNote", text(context.body(), "reason", "draft")), + toJson(context.body()), status, context.commandId(), actor(context), actor(context), DEFAULT_TENANT_ID); + } + + private Long insertMarketInstallation(OperationContext context) { + Long assetId = firstPathLong(context, "assetId"); + Long versionId = latestAssetVersionId(assetId); + return jdbcTemplate.queryForObject(""" + INSERT INTO muse_market_installation(asset_id, asset_version_id, user_id, authorization_snapshot_id, + source_snapshot_id, command_id, creator, updater, tenant_id) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) + ON CONFLICT (tenant_id, user_id, asset_id, asset_version_id) + DO UPDATE SET status = 'installed', uninstalled_at = NULL, command_id = EXCLUDED.command_id, + updater = EXCLUDED.updater + RETURNING id + """, Long.class, assetId, versionId, context.actorUserId(), + longValue(context.body().get("authorizationSnapshotId")), longValue(context.body().get("sourceSnapshotId")), + context.commandId(), actor(context), actor(context), DEFAULT_TENANT_ID); + } + + private Long insertPublishRequest(OperationContext context) { + Long assetId = longValue(context.body().getOrDefault("assetId", context.body().get("draftId"))); + if (assetId == null) { + assetId = insertMarketAsset(context, "draft"); + } + return jdbcTemplate.queryForObject(""" + INSERT INTO muse_market_publish_request(asset_id, asset_version_id, publisher_id, submitted_at, + command_id, creator, updater, tenant_id) + VALUES (?, ?, ?, CURRENT_TIMESTAMP, ?, ?, ?, ?) + RETURNING id + """, Long.class, assetId, latestAssetVersionId(assetId), context.actorUserId(), context.commandId(), + actor(context), actor(context), DEFAULT_TENANT_ID); + } + + private Long updatePublishRequest(OperationContext context) { + Long requestId = firstPathLong(context, "requestId"); + String status = context.operationId().contains("Approve") ? "approved" + : context.operationId().contains("Reject") ? "rejected" : "withdrawn"; + jdbcTemplate.update(""" + UPDATE muse_market_publish_request + SET status = ?, review_note = COALESCE(?, review_note), reviewer_id = ?, + reviewed_at = CURRENT_TIMESTAMP, command_id = ?, revision = revision + 1, updater = ? + WHERE id = ? AND tenant_id = ? AND deleted = FALSE + AND (? <> 'app' OR publisher_id = ?) + """, status, text(context.body(), "note", text(context.body(), "reason", null)), context.actorUserId(), + context.commandId(), actor(context), requestId, DEFAULT_TENANT_ID, context.side(), context.actorUserId()); + if ("approved".equals(status)) { + Long assetId = queryOptionalLong("SELECT asset_id FROM muse_market_publish_request WHERE id = ? AND tenant_id = ?", + requestId, DEFAULT_TENANT_ID); + if (assetId != null) { + jdbcTemplate.update(""" + UPDATE muse_market_asset + SET listing_status = 'listed', status = 'active', revision = revision + 1, updater = ? + WHERE id = ? AND tenant_id = ? + """, actor(context), assetId, DEFAULT_TENANT_ID); + } + } + return requestId; + } + + private Long insertMarketAppeal(OperationContext context) { + return jdbcTemplate.queryForObject(""" + INSERT INTO muse_market_appeal(asset_id, user_id, appeal_type, reason, command_id, creator, updater, tenant_id) + VALUES (?, ?, ?, ?, ?, ?, ?, ?) + RETURNING id + """, Long.class, longValue(context.body().getOrDefault("assetId", firstPathLong(context, "assetId"))), + context.actorUserId(), text(context.body(), "appealType", "general"), text(context.body(), "reason", null), + context.commandId(), actor(context), actor(context), DEFAULT_TENANT_ID); + } + + private Long updateMarketAppeal(OperationContext context) { + Long appealId = firstPathLong(context, "appealId"); + String status = context.operationId().contains("withdraw") ? "withdrawn" + : context.operationId().contains("Resolve") ? "resolved" : "supplemented"; + jdbcTemplate.update(""" + UPDATE muse_market_appeal + SET status = ?, resolution = COALESCE(?, resolution), resolver_id = ?, resolved_at = CURRENT_TIMESTAMP, + command_id = ?, revision = revision + 1, updater = ? + WHERE id = ? AND tenant_id = ? AND deleted = FALSE + AND (? <> 'app' OR user_id = ?) + """, status, text(context.body(), "resolution", text(context.body(), "description", null)), + context.actorUserId(), context.commandId(), actor(context), appealId, DEFAULT_TENANT_ID, + context.side(), context.actorUserId()); + return appealId; + } + + private Long insertMarketHandoff(OperationContext context, Map response) { + String token = "handoff-" + UUID.randomUUID(); + Long assetId = longValue(context.body().get("assetId")); + Long versionId = latestAssetVersionId(assetId); + Long id = jdbcTemplate.queryForObject(""" + INSERT INTO muse_market_handoff(token, source_asset_id, asset_version_id, target_owner_hint, + target_endpoint, handoff_token_hash, license_summary, governance_summary, expires_at, + command_id, creator, updater, tenant_id) + VALUES (?, ?, ?, ?, ?, ?, CAST(? AS jsonb), CAST(? AS jsonb), CURRENT_TIMESTAMP + INTERVAL '30 minutes', + ?, ?, ?, ?) + RETURNING id + """, Long.class, token, assetId, versionId, text(context.body(), "targetOwner", null), + text(context.body(), "returnUrl", null), sha256(token), toJson(context.body().get("licenseSummary")), + toJson(context.body().get("governanceSummary")), context.commandId(), actor(context), actor(context), + DEFAULT_TENANT_ID); + response.put("handoffToken", token); + return id; + } + + private Long cancelMarketHandoff(OperationContext context) { + String token = firstPathString(context, "handoffToken"); + Long id = queryOptionalLong("SELECT id FROM muse_market_handoff WHERE tenant_id = ? AND token = ? AND deleted = FALSE", + DEFAULT_TENANT_ID, token); + if (id == null) { + throw new ServiceException(CONTRACT_NOT_FOUND, "Handoff 不存在"); + } + jdbcTemplate.update(""" + UPDATE muse_market_handoff + SET status = 'canceled', canceled_at = CURRENT_TIMESTAMP, command_id = ?, updater = ? + WHERE id = ? AND tenant_id = ? + AND (? <> 'app' OR creator = ?) + """, context.commandId(), actor(context), id, DEFAULT_TENANT_ID, context.side(), actor(context)); + return id; + } + + private Long latestAssetVersionId(Long assetId) { + Long versionId = queryOptionalLong(""" + SELECT id FROM muse_market_asset_version + WHERE tenant_id = ? AND asset_id = ? AND deleted = FALSE + ORDER BY id DESC LIMIT 1 + """, DEFAULT_TENANT_ID, assetId); + if (versionId != null) { + return versionId; + } + return jdbcTemplate.queryForObject(""" + INSERT INTO muse_market_asset_version(asset_id, version, status, creator, updater, tenant_id) + VALUES (?, '1', 'active', '', '', ?) + RETURNING id + """, Long.class, assetId, DEFAULT_TENANT_ID); + } + + private void requireAppWorkOwner(OperationContext context, Long workId) { + if (!"app".equals(context.side()) || workId == null) { + return; + } + Integer count = queryOptionalInteger(""" + SELECT COUNT(1) FROM muse_content_work + WHERE id = ? AND tenant_id = ? AND owner_user_id = ? AND deleted = FALSE + """, workId, DEFAULT_TENANT_ID, context.actorUserId()); + if (count == null || count == 0) { + throw new ServiceException(CONTRACT_NOT_FOUND, "作品不存在或无权访问"); + } + } + + private void requireKnowledgeBaseOwner(OperationContext context, Long kbId) { + if (!"app".equals(context.side()) || kbId == null) { + return; + } + Integer count = queryOptionalInteger(""" + SELECT COUNT(1) FROM muse_knowledge_base + WHERE id = ? AND tenant_id = ? AND owner_user_id = ? AND deleted = FALSE + """, kbId, DEFAULT_TENANT_ID, context.actorUserId()); + if (count == null || count == 0) { + throw new ServiceException(CONTRACT_NOT_FOUND, "知识库不存在或无权访问"); + } + } + + private void requireKnowledgeDocumentOwner(OperationContext context, Long documentId) { + if (!"app".equals(context.side()) || documentId == null) { + return; + } + Integer count = queryOptionalInteger(""" + SELECT COUNT(1) + FROM muse_knowledge_document d + JOIN muse_knowledge_base kb ON kb.id = d.kb_id AND kb.tenant_id = d.tenant_id + WHERE d.id = ? AND d.tenant_id = ? AND d.deleted = FALSE + AND kb.owner_user_id = ? AND kb.deleted = FALSE + """, documentId, DEFAULT_TENANT_ID, context.actorUserId()); + if (count == null || count == 0) { + throw new ServiceException(CONTRACT_NOT_FOUND, "知识库文档不存在或无权访问"); + } + } + + private Map latestOperationPayload(OperationContext context, String resourceType) { + String json = queryOptionalString(""" + SELECT request_payload::text + FROM muse_domain_operation_record + WHERE tenant_id = ? AND deleted = FALSE AND domain = ? AND resource_type = ? + AND actor_user_id = ? + ORDER BY id DESC LIMIT 1 + """, DEFAULT_TENANT_ID, context.domain(), resourceType, context.actorUserId()); + return json == null ? Map.of() : readMap(json); + } + + private String resourceKey(OperationContext context, ResourceSpec spec) { + if (spec == null || spec.keyPathVariable() == null) { + return firstPathString(context, "correlationId", "handoffToken", "schemaKey", "promptKey", "policyKey"); + } + return firstPathString(context, spec.keyPathVariable()); + } + + private String parentType(OperationContext context, ResourceSpec spec) { + if (spec == null || spec.parentColumn() == null) { + return null; + } + return spec.parentColumn().replace("_id", ""); + } + + private Long parentId(OperationContext context, ResourceSpec spec) { + if (spec == null || spec.parentPathVariable() == null) { + return null; + } + return firstPathLong(context, spec.parentPathVariable()); + } + + private boolean isCurrentUserOperation(String operationId) { + return "getCurrentUser".equals(operationId) || "getProfile".equals(operationId); + } + + private boolean isReadOnlyCommand(String operationId) { + return Set.of("validateMetaSchemaDraft", "previewMetaSchemaDraftImpact", "previewGlobalKBImpact", + "previewFunctionChainImpact", "runPublishCheck", "adminPreviewGovernanceImpact", + "validateDynamicFields", "querySourceStatus").contains(operationId); + } + + private boolean isMaterializedWrite(OperationContext context) { + String operationId = context.operationId(); + return switch (context.domain()) { + case "account" -> Set.of("updateProfile", "adminCreateQuotaAdjustment", "adminCreateNewApiBinding", + "appRecheckNewApiBinding", "appAcknowledgeSecurityEvent", "adminCreateQuotaRequest", + "adminCreateCallAttributionJob", "appCreateQuotaRequest", "appCreateExportTask").contains(operationId); + case "meta" -> Set.of("saveMetaSchemaDraft", "publishMetaSchemaDraft", "activateMetaSchemaVersion", + "rollbackMetaSchemaVersion", "deprecateMetaSchemaVersion", "setMetaSchemaGrayRules", + "activateFunctionChainVersion").contains(operationId); + case "ai" -> Set.of("createAiTask", "adminCreateAgent", "createUserAgent", "adminCreateAgentVersion", + "createAgentVersion", "testAgent", "adminCreatePromptVersion", "adminActivatePromptVersion", + "adminCreateQualityPolicyVersion", "adminCreateOrAdjustToolGrant", "bindAgentSlot", + "rejectSuggestion", "adminCancelJob", "cancelUserJob", "adminStartEvaluationRun", + "adminRetryJob", "adminRetrySourceEvent", "precheckAgentSlot", "recheckSourceStatus").contains(operationId); + case "knowledge" -> Set.of("createKnowledgeBase", "createGlobalKnowledgeBase", "updateKnowledgeBase", + "updateGlobalKnowledgeBase", "deleteKnowledgeBase", "disableKnowledgeBase", + "restoreKnowledgeBase", "enableGlobalKnowledgeBase", "disableGlobalKnowledgeBase", + "uploadKBDocument", "uploadGlobalKBDocument", "deleteKBDocument", "deleteGlobalKBDocument", + "createKBDocumentVersion", "createGlobalKBDocumentVersion", "updateEntity", "createRelation", + "confirmKnowledgeDraft", "ignoreKnowledgeDraft", "recheckKnowledgeDraft", + "createKnowledgeBinding", "deleteKnowledgeBinding", "activateGlobalKBVersion", + "saveGlobalKBAccessPolicyDraft", "publishGlobalKBAccessPolicy", "reindexGlobalKnowledgeBase", + "triggerGlobalKBSourceEvent", "reindexKnowledgeBase", "createKBPublishReadiness", + "createKBPublishSnapshot", "createKBMarketPublishReadiness", "createKBExportTask", + "disableInstalledKnowledgeBase", "deleteInstalledKnowledgeBase", "restoreInstalledKnowledgeBase", + "createKnowledgeBindingPrecheck").contains(operationId); + case "market" -> Set.of("savePublishDraft", "submitPublishRequest", "withdrawPublishRequest", + "adminApprovePublishRequest", "adminRejectPublishRequest", "installMarketplaceAsset", + "createMarketplaceHandoff", "cancelHandoff", "submitAppeal", "supplementAppeal", + "withdrawAppeal", "adminResolveAppeal", "adminDelistAsset", "adminRecallAsset", + "favoriteAsset", "unfavoriteAsset", "purchaseAsset", "createBindPrecheck").contains(operationId); + default -> false; + }; + } + + private boolean isListOperation(String operationId) { + return operationId.startsWith("list") || operationId.startsWith("adminList") || operationId.startsWith("appList") + || operationId.startsWith("getApp") && operationId.endsWith("Snapshots"); + } + + private Object firstBodyValue(OperationContext context, String... keys) { + for (String key : keys) { + Object value = context.body().get(key); + if (value != null) { + return value; + } + } + return null; + } + + private Long firstPathLong(OperationContext context, String... keys) { + for (String key : keys) { + Long value = longValue(context.pathVariables().get(key)); + if (value != null) { + return value; + } + } + return null; + } + + private Integer firstPathInt(OperationContext context, String key) { + return intValue(context.pathVariables().get(key), null); + } + + private String firstPathString(OperationContext context, String... keys) { + for (String key : keys) { + String value = context.pathVariables().get(key); + if (StrUtil.isNotBlank(value)) { + return value; + } + } + return null; + } + + private String firstNonBlank(String... values) { + for (String value : values) { + if (StrUtil.isNotBlank(value) && !"null".equals(value)) { + return value; + } + } + return null; + } + + private String actor(OperationContext context) { + return context.actorUserId() == null ? "" : String.valueOf(context.actorUserId()); + } + + private String text(Map body, String key, String defaultValue) { + Object value = body.get(key); + return StrUtil.isBlankIfStr(value) ? defaultValue : String.valueOf(value); + } + + private Long longValue(Object value) { + if (value == null || StrUtil.isBlankIfStr(value)) { + return null; + } + if (value instanceof Number number) { + return number.longValue(); + } + try { + return Long.parseLong(String.valueOf(value)); + } catch (NumberFormatException ignored) { + return null; + } + } + + private Integer intValue(Object value, Integer defaultValue) { + if (value == null || StrUtil.isBlankIfStr(value)) { + return defaultValue; + } + if (value instanceof Number number) { + return number.intValue(); + } + try { + return Integer.parseInt(String.valueOf(value)); + } catch (NumberFormatException ignored) { + return defaultValue; + } + } + + private Integer nextInt(String sql, Object... args) { + Integer value = jdbcTemplate.queryForObject(sql, Integer.class, args); + return value == null ? 1 : value; + } + + private Long queryOptionalLong(String sql, Object... args) { + try { + return jdbcTemplate.queryForObject(sql, Long.class, args); + } catch (EmptyResultDataAccessException ignored) { + return null; + } + } + + private Integer queryOptionalInteger(String sql, Object... args) { + try { + return jdbcTemplate.queryForObject(sql, Integer.class, args); + } catch (EmptyResultDataAccessException ignored) { + return null; + } + } + + private String queryOptionalString(String sql, Object... args) { + try { + return jdbcTemplate.queryForObject(sql, String.class, args); + } catch (EmptyResultDataAccessException ignored) { + return null; + } + } + + private String toJson(Object value) { + try { + return objectMapper.writeValueAsString(value == null ? Map.of() : value); + } catch (JsonProcessingException ex) { + throw new ServiceException(CONTRACT_UNSUPPORTED, "请求无法序列化为 JSON"); + } + } + + private Map readMap(String json) { + try { + return objectMapper.readValue(json, new TypeReference<>() { + }); + } catch (JsonProcessingException ex) { + throw new ServiceException(CONTRACT_UNSUPPORTED, "数据库 JSON 无法反序列化"); + } + } + + private List> readList(String json) { + try { + return objectMapper.readValue(json == null ? "[]" : json, new TypeReference<>() { + }); + } catch (JsonProcessingException ex) { + throw new ServiceException(CONTRACT_UNSUPPORTED, "数据库 JSON 列表无法反序列化"); + } + } + + private String slug(String value) { + return value == null ? "muse" : value.toLowerCase(Locale.ROOT).replaceAll("[^a-z0-9]+", "-"); + } + + private String sha256(String value) { + try { + MessageDigest digest = MessageDigest.getInstance("SHA-256"); + byte[] encoded = digest.digest(value.getBytes(StandardCharsets.UTF_8)); + StringBuilder builder = new StringBuilder(); + for (byte b : encoded) { + builder.append(String.format("%02x", b)); + } + return builder.toString(); + } catch (NoSuchAlgorithmException ex) { + throw new ServiceException(CONTRACT_UNSUPPORTED, "当前 JDK 不支持 SHA-256"); + } + } + + /** + * 合同操作运行上下文。 + */ + private record OperationContext(String domain, String side, String operationId, String httpMethod, String requestUri, + String commandId, Map pathVariables, + Map queryParams, Map body, Long actorUserId) { + + @SuppressWarnings("unchecked") + private static OperationContext from(Map contract, String httpMethod, String requestUri, + String headerCommandId, Map queryParams, + Map body, Long actorUserId) { + Map safeBody = body == null ? Map.of() : body; + Map safeQuery = queryParams == null ? Map.of() : queryParams; + String commandId = firstNonBlank(String.valueOf(safeBody.getOrDefault("commandId", "")), + safeQuery.get("commandId"), headerCommandId, String.valueOf(contract.getOrDefault("commandId", ""))); + return new OperationContext(String.valueOf(contract.get("module")), String.valueOf(contract.get("entry")), + String.valueOf(contract.get("operationId")), httpMethod, requestUri, + StrUtil.isBlank(commandId) ? null : commandId, + (Map) contract.getOrDefault("pathVariables", Map.of()), safeQuery, safeBody, + actorUserId); + } + + private static String firstNonBlank(String... values) { + for (String value : values) { + if (StrUtil.isNotBlank(value) && !"null".equals(value)) { + return value; + } + } + return null; + } + } + + /** + * 合同操作到领域表的最小映射。 + */ + private record ResourceSpec(String resourceType, String tableName, String idColumn, String idPathVariable, + String keyColumn, String keyPathVariable, String ownerColumn, String ownerPathVariable, + String parentColumn, String parentPathVariable, String statusColumn, + String revisionColumn) { + } + + /** + * 动态 WHERE 片段与参数。 + */ + private record WhereClause(String sql, List args) { + } + +} diff --git a/muse-cloud/muse-module-ai/muse-module-ai-contract-server/pom.xml b/muse-cloud/muse-module-ai/muse-module-ai-contract-server/pom.xml index 5834ebdd..081c614f 100644 --- a/muse-cloud/muse-module-ai/muse-module-ai-contract-server/pom.xml +++ b/muse-cloud/muse-module-ai/muse-module-ai-contract-server/pom.xml @@ -23,5 +23,9 @@ cn.iocoder.cloud muse-spring-boot-starter-security + + cn.iocoder.cloud + muse-spring-boot-starter-mybatis + diff --git a/muse-cloud/muse-module-ai/muse-module-ai-contract-server/src/main/java/cn/iocoder/muse/module/ai/controller/admin/AdminMuseAiContractController.java b/muse-cloud/muse-module-ai/muse-module-ai-contract-server/src/main/java/cn/iocoder/muse/module/ai/controller/admin/AdminMuseAiContractController.java index bb80c6c3..451efbcc 100644 --- a/muse-cloud/muse-module-ai/muse-module-ai-contract-server/src/main/java/cn/iocoder/muse/module/ai/controller/admin/AdminMuseAiContractController.java +++ b/muse-cloud/muse-module-ai/muse-module-ai-contract-server/src/main/java/cn/iocoder/muse/module/ai/controller/admin/AdminMuseAiContractController.java @@ -1,9 +1,10 @@ package cn.iocoder.muse.module.ai.controller.admin; -import cn.iocoder.muse.framework.common.muse.MuseApiContractSupport; import cn.iocoder.muse.framework.common.pojo.CommonResult; +import cn.iocoder.muse.framework.mybatis.core.muse.MuseContractPersistenceService; import io.swagger.v3.oas.annotations.Operation; import io.swagger.v3.oas.annotations.tags.Tag; +import jakarta.annotation.Resource; import jakarta.servlet.http.HttpServletRequest; import org.springframework.validation.annotation.Validated; import org.springframework.web.bind.annotation.*; @@ -28,6 +29,9 @@ public class AdminMuseAiContractController { /** 当前 Controller 拥有的 OpenAPI 合同域。 */ private static final Set DOMAINS = Set.of("ai"); + @Resource + private MuseContractPersistenceService contractPersistenceService; + /** * 处理当前领域下尚未落成专用 Controller 的 OpenAPI 合同请求。 * @@ -41,7 +45,7 @@ public class AdminMuseAiContractController { public CommonResult> handle(HttpServletRequest request, @RequestParam Map queryParams, @RequestBody(required = false) Map body) { - return success(MuseApiContractSupport.handle(DOMAINS, "admin", request.getMethod(), request.getRequestURI(), + return success(contractPersistenceService.handle(DOMAINS, "admin", request.getMethod(), request.getRequestURI(), request.getHeader("X-Command-Id"), queryParams, body, getLoginUserId())); } diff --git a/muse-cloud/muse-module-ai/muse-module-ai-contract-server/src/main/java/cn/iocoder/muse/module/ai/controller/app/AppMuseAiContractController.java b/muse-cloud/muse-module-ai/muse-module-ai-contract-server/src/main/java/cn/iocoder/muse/module/ai/controller/app/AppMuseAiContractController.java index bafbf955..42cb852e 100644 --- a/muse-cloud/muse-module-ai/muse-module-ai-contract-server/src/main/java/cn/iocoder/muse/module/ai/controller/app/AppMuseAiContractController.java +++ b/muse-cloud/muse-module-ai/muse-module-ai-contract-server/src/main/java/cn/iocoder/muse/module/ai/controller/app/AppMuseAiContractController.java @@ -1,9 +1,10 @@ package cn.iocoder.muse.module.ai.controller.app; -import cn.iocoder.muse.framework.common.muse.MuseApiContractSupport; import cn.iocoder.muse.framework.common.pojo.CommonResult; +import cn.iocoder.muse.framework.mybatis.core.muse.MuseContractPersistenceService; import io.swagger.v3.oas.annotations.Operation; import io.swagger.v3.oas.annotations.tags.Tag; +import jakarta.annotation.Resource; import jakarta.servlet.http.HttpServletRequest; import org.springframework.http.MediaType; import org.springframework.validation.annotation.Validated; @@ -34,6 +35,9 @@ public class AppMuseAiContractController { /** AI 任务 SSE 连接默认超时时间。 */ private static final long DEFAULT_TIMEOUT_MILLIS = 30_000L; + @Resource + private MuseContractPersistenceService contractPersistenceService; + /** * 建立 AI 任务级 SSE 流。 * @@ -72,7 +76,7 @@ public class AppMuseAiContractController { public CommonResult> handle(HttpServletRequest request, @RequestParam Map queryParams, @RequestBody(required = false) Map body) { - return success(MuseApiContractSupport.handle(DOMAINS, "app", request.getMethod(), request.getRequestURI(), + return success(contractPersistenceService.handle(DOMAINS, "app", request.getMethod(), request.getRequestURI(), request.getHeader("X-Command-Id"), queryParams, body, getLoginUserId())); } diff --git a/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/pom.xml b/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/pom.xml index 0876e369..9d6f6a72 100644 --- a/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/pom.xml +++ b/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/pom.xml @@ -22,5 +22,9 @@ cn.iocoder.cloud muse-spring-boot-starter-security + + cn.iocoder.cloud + muse-spring-boot-starter-mybatis + diff --git a/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/main/java/cn/iocoder/muse/module/knowledge/controller/admin/AdminMuseKnowledgeContractController.java b/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/main/java/cn/iocoder/muse/module/knowledge/controller/admin/AdminMuseKnowledgeContractController.java index c6bd2cdd..0b8e27f6 100644 --- a/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/main/java/cn/iocoder/muse/module/knowledge/controller/admin/AdminMuseKnowledgeContractController.java +++ b/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/main/java/cn/iocoder/muse/module/knowledge/controller/admin/AdminMuseKnowledgeContractController.java @@ -1,9 +1,10 @@ package cn.iocoder.muse.module.knowledge.controller.admin; -import cn.iocoder.muse.framework.common.muse.MuseApiContractSupport; import cn.iocoder.muse.framework.common.pojo.CommonResult; +import cn.iocoder.muse.framework.mybatis.core.muse.MuseContractPersistenceService; import io.swagger.v3.oas.annotations.Operation; import io.swagger.v3.oas.annotations.tags.Tag; +import jakarta.annotation.Resource; import jakarta.servlet.http.HttpServletRequest; import org.springframework.validation.annotation.Validated; import org.springframework.web.bind.annotation.*; @@ -28,6 +29,9 @@ public class AdminMuseKnowledgeContractController { /** 当前 Controller 拥有的 OpenAPI 合同域。 */ private static final Set DOMAINS = Set.of("knowledge"); + @Resource + private MuseContractPersistenceService contractPersistenceService; + /** * 处理当前领域下尚未落成专用 Controller 的 OpenAPI 合同请求。 * @@ -41,7 +45,7 @@ public class AdminMuseKnowledgeContractController { public CommonResult> handle(HttpServletRequest request, @RequestParam Map queryParams, @RequestBody(required = false) Map body) { - return success(MuseApiContractSupport.handle(DOMAINS, "admin", request.getMethod(), request.getRequestURI(), + return success(contractPersistenceService.handle(DOMAINS, "admin", request.getMethod(), request.getRequestURI(), request.getHeader("X-Command-Id"), queryParams, body, getLoginUserId())); } diff --git a/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/main/java/cn/iocoder/muse/module/knowledge/controller/app/AppMuseKnowledgeContractController.java b/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/main/java/cn/iocoder/muse/module/knowledge/controller/app/AppMuseKnowledgeContractController.java index f15ad4fb..8fd3a54b 100644 --- a/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/main/java/cn/iocoder/muse/module/knowledge/controller/app/AppMuseKnowledgeContractController.java +++ b/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/main/java/cn/iocoder/muse/module/knowledge/controller/app/AppMuseKnowledgeContractController.java @@ -1,9 +1,10 @@ package cn.iocoder.muse.module.knowledge.controller.app; -import cn.iocoder.muse.framework.common.muse.MuseApiContractSupport; import cn.iocoder.muse.framework.common.pojo.CommonResult; +import cn.iocoder.muse.framework.mybatis.core.muse.MuseContractPersistenceService; import io.swagger.v3.oas.annotations.Operation; import io.swagger.v3.oas.annotations.tags.Tag; +import jakarta.annotation.Resource; import jakarta.servlet.http.HttpServletRequest; import org.springframework.validation.annotation.Validated; import org.springframework.web.bind.annotation.*; @@ -28,6 +29,9 @@ public class AppMuseKnowledgeContractController { /** 当前 Controller 拥有的 OpenAPI 合同域。 */ private static final Set DOMAINS = Set.of("knowledge"); + @Resource + private MuseContractPersistenceService contractPersistenceService; + /** * 处理当前领域下尚未落成专用 Controller 的 OpenAPI 合同请求。 * @@ -41,7 +45,7 @@ public class AppMuseKnowledgeContractController { public CommonResult> handle(HttpServletRequest request, @RequestParam Map queryParams, @RequestBody(required = false) Map body) { - return success(MuseApiContractSupport.handle(DOMAINS, "app", request.getMethod(), request.getRequestURI(), + return success(contractPersistenceService.handle(DOMAINS, "app", request.getMethod(), request.getRequestURI(), request.getHeader("X-Command-Id"), queryParams, body, getLoginUserId())); } diff --git a/muse-cloud/muse-module-market/muse-module-market-server/pom.xml b/muse-cloud/muse-module-market/muse-module-market-server/pom.xml index 75040db6..f1187d6b 100644 --- a/muse-cloud/muse-module-market/muse-module-market-server/pom.xml +++ b/muse-cloud/muse-module-market/muse-module-market-server/pom.xml @@ -22,5 +22,9 @@ cn.iocoder.cloud muse-spring-boot-starter-security + + cn.iocoder.cloud + muse-spring-boot-starter-mybatis + diff --git a/muse-cloud/muse-module-market/muse-module-market-server/src/main/java/cn/iocoder/muse/module/market/controller/admin/AdminMuseMarketContractController.java b/muse-cloud/muse-module-market/muse-module-market-server/src/main/java/cn/iocoder/muse/module/market/controller/admin/AdminMuseMarketContractController.java index c5311848..dad89d20 100644 --- a/muse-cloud/muse-module-market/muse-module-market-server/src/main/java/cn/iocoder/muse/module/market/controller/admin/AdminMuseMarketContractController.java +++ b/muse-cloud/muse-module-market/muse-module-market-server/src/main/java/cn/iocoder/muse/module/market/controller/admin/AdminMuseMarketContractController.java @@ -1,9 +1,10 @@ package cn.iocoder.muse.module.market.controller.admin; -import cn.iocoder.muse.framework.common.muse.MuseApiContractSupport; import cn.iocoder.muse.framework.common.pojo.CommonResult; +import cn.iocoder.muse.framework.mybatis.core.muse.MuseContractPersistenceService; import io.swagger.v3.oas.annotations.Operation; import io.swagger.v3.oas.annotations.tags.Tag; +import jakarta.annotation.Resource; import jakarta.servlet.http.HttpServletRequest; import org.springframework.validation.annotation.Validated; import org.springframework.web.bind.annotation.*; @@ -28,6 +29,9 @@ public class AdminMuseMarketContractController { /** 当前 Controller 拥有的 OpenAPI 合同域。 */ private static final Set DOMAINS = Set.of("market"); + @Resource + private MuseContractPersistenceService contractPersistenceService; + /** * 处理当前领域下尚未落成专用 Controller 的 OpenAPI 合同请求。 * @@ -41,7 +45,7 @@ public class AdminMuseMarketContractController { public CommonResult> handle(HttpServletRequest request, @RequestParam Map queryParams, @RequestBody(required = false) Map body) { - return success(MuseApiContractSupport.handle(DOMAINS, "admin", request.getMethod(), request.getRequestURI(), + return success(contractPersistenceService.handle(DOMAINS, "admin", request.getMethod(), request.getRequestURI(), request.getHeader("X-Command-Id"), queryParams, body, getLoginUserId())); } diff --git a/muse-cloud/muse-module-market/muse-module-market-server/src/main/java/cn/iocoder/muse/module/market/controller/app/AppMuseMarketContractController.java b/muse-cloud/muse-module-market/muse-module-market-server/src/main/java/cn/iocoder/muse/module/market/controller/app/AppMuseMarketContractController.java index 8ce4e196..5004af14 100644 --- a/muse-cloud/muse-module-market/muse-module-market-server/src/main/java/cn/iocoder/muse/module/market/controller/app/AppMuseMarketContractController.java +++ b/muse-cloud/muse-module-market/muse-module-market-server/src/main/java/cn/iocoder/muse/module/market/controller/app/AppMuseMarketContractController.java @@ -1,9 +1,10 @@ package cn.iocoder.muse.module.market.controller.app; -import cn.iocoder.muse.framework.common.muse.MuseApiContractSupport; import cn.iocoder.muse.framework.common.pojo.CommonResult; +import cn.iocoder.muse.framework.mybatis.core.muse.MuseContractPersistenceService; import io.swagger.v3.oas.annotations.Operation; import io.swagger.v3.oas.annotations.tags.Tag; +import jakarta.annotation.Resource; import jakarta.servlet.http.HttpServletRequest; import org.springframework.validation.annotation.Validated; import org.springframework.web.bind.annotation.*; @@ -28,6 +29,9 @@ public class AppMuseMarketContractController { /** 当前 Controller 拥有的 OpenAPI 合同域。 */ private static final Set DOMAINS = Set.of("market"); + @Resource + private MuseContractPersistenceService contractPersistenceService; + /** * 处理当前领域下尚未落成专用 Controller 的 OpenAPI 合同请求。 * @@ -41,7 +45,7 @@ public class AppMuseMarketContractController { public CommonResult> handle(HttpServletRequest request, @RequestParam Map queryParams, @RequestBody(required = false) Map body) { - return success(MuseApiContractSupport.handle(DOMAINS, "app", request.getMethod(), request.getRequestURI(), + return success(contractPersistenceService.handle(DOMAINS, "app", request.getMethod(), request.getRequestURI(), request.getHeader("X-Command-Id"), queryParams, body, getLoginUserId())); } diff --git a/muse-cloud/muse-module-member/muse-module-member-server/src/main/java/cn/iocoder/muse/module/member/controller/admin/AdminMuseAccountContractController.java b/muse-cloud/muse-module-member/muse-module-member-server/src/main/java/cn/iocoder/muse/module/member/controller/admin/AdminMuseAccountContractController.java index c89ebd27..7dbad445 100644 --- a/muse-cloud/muse-module-member/muse-module-member-server/src/main/java/cn/iocoder/muse/module/member/controller/admin/AdminMuseAccountContractController.java +++ b/muse-cloud/muse-module-member/muse-module-member-server/src/main/java/cn/iocoder/muse/module/member/controller/admin/AdminMuseAccountContractController.java @@ -1,9 +1,10 @@ package cn.iocoder.muse.module.member.controller.admin; -import cn.iocoder.muse.framework.common.muse.MuseApiContractSupport; import cn.iocoder.muse.framework.common.pojo.CommonResult; +import cn.iocoder.muse.framework.mybatis.core.muse.MuseContractPersistenceService; import io.swagger.v3.oas.annotations.Operation; import io.swagger.v3.oas.annotations.tags.Tag; +import jakarta.annotation.Resource; import jakarta.servlet.http.HttpServletRequest; import org.springframework.validation.annotation.Validated; import org.springframework.web.bind.annotation.*; @@ -28,6 +29,9 @@ public class AdminMuseAccountContractController { /** 当前 Controller 拥有的 OpenAPI 合同域。 */ private static final Set DOMAINS = Set.of("account"); + @Resource + private MuseContractPersistenceService contractPersistenceService; + /** * 处理当前领域下尚未落成专用 Controller 的 OpenAPI 合同请求。 * @@ -41,7 +45,7 @@ public class AdminMuseAccountContractController { public CommonResult> handle(HttpServletRequest request, @RequestParam Map queryParams, @RequestBody(required = false) Map body) { - return success(MuseApiContractSupport.handle(DOMAINS, "admin", request.getMethod(), request.getRequestURI(), + return success(contractPersistenceService.handle(DOMAINS, "admin", request.getMethod(), request.getRequestURI(), request.getHeader("X-Command-Id"), queryParams, body, getLoginUserId())); } diff --git a/muse-cloud/muse-module-member/muse-module-member-server/src/main/java/cn/iocoder/muse/module/member/controller/app/AppMuseAccountContractController.java b/muse-cloud/muse-module-member/muse-module-member-server/src/main/java/cn/iocoder/muse/module/member/controller/app/AppMuseAccountContractController.java index 556e0bc3..24a5f091 100644 --- a/muse-cloud/muse-module-member/muse-module-member-server/src/main/java/cn/iocoder/muse/module/member/controller/app/AppMuseAccountContractController.java +++ b/muse-cloud/muse-module-member/muse-module-member-server/src/main/java/cn/iocoder/muse/module/member/controller/app/AppMuseAccountContractController.java @@ -1,9 +1,10 @@ package cn.iocoder.muse.module.member.controller.app; -import cn.iocoder.muse.framework.common.muse.MuseApiContractSupport; import cn.iocoder.muse.framework.common.pojo.CommonResult; +import cn.iocoder.muse.framework.mybatis.core.muse.MuseContractPersistenceService; import io.swagger.v3.oas.annotations.Operation; import io.swagger.v3.oas.annotations.tags.Tag; +import jakarta.annotation.Resource; import jakarta.servlet.http.HttpServletRequest; import org.springframework.validation.annotation.Validated; import org.springframework.web.bind.annotation.*; @@ -28,6 +29,9 @@ public class AppMuseAccountContractController { /** 当前 Controller 拥有的 OpenAPI 合同域。 */ private static final Set DOMAINS = Set.of("account"); + @Resource + private MuseContractPersistenceService contractPersistenceService; + /** * 处理当前领域下尚未落成专用 Controller 的 OpenAPI 合同请求。 * @@ -41,7 +45,7 @@ public class AppMuseAccountContractController { public CommonResult> handle(HttpServletRequest request, @RequestParam Map queryParams, @RequestBody(required = false) Map body) { - return success(MuseApiContractSupport.handle(DOMAINS, "app", request.getMethod(), request.getRequestURI(), + return success(contractPersistenceService.handle(DOMAINS, "app", request.getMethod(), request.getRequestURI(), request.getHeader("X-Command-Id"), queryParams, body, getLoginUserId())); } diff --git a/muse-cloud/muse-module-meta/muse-module-meta-server/pom.xml b/muse-cloud/muse-module-meta/muse-module-meta-server/pom.xml index f43a4761..d176be38 100644 --- a/muse-cloud/muse-module-meta/muse-module-meta-server/pom.xml +++ b/muse-cloud/muse-module-meta/muse-module-meta-server/pom.xml @@ -22,5 +22,9 @@ cn.iocoder.cloud muse-spring-boot-starter-security + + cn.iocoder.cloud + muse-spring-boot-starter-mybatis + diff --git a/muse-cloud/muse-module-meta/muse-module-meta-server/src/main/java/cn/iocoder/muse/module/meta/controller/admin/AdminMuseMetaContractController.java b/muse-cloud/muse-module-meta/muse-module-meta-server/src/main/java/cn/iocoder/muse/module/meta/controller/admin/AdminMuseMetaContractController.java index c0dd2414..93149b77 100644 --- a/muse-cloud/muse-module-meta/muse-module-meta-server/src/main/java/cn/iocoder/muse/module/meta/controller/admin/AdminMuseMetaContractController.java +++ b/muse-cloud/muse-module-meta/muse-module-meta-server/src/main/java/cn/iocoder/muse/module/meta/controller/admin/AdminMuseMetaContractController.java @@ -1,9 +1,10 @@ package cn.iocoder.muse.module.meta.controller.admin; -import cn.iocoder.muse.framework.common.muse.MuseApiContractSupport; import cn.iocoder.muse.framework.common.pojo.CommonResult; +import cn.iocoder.muse.framework.mybatis.core.muse.MuseContractPersistenceService; import io.swagger.v3.oas.annotations.Operation; import io.swagger.v3.oas.annotations.tags.Tag; +import jakarta.annotation.Resource; import jakarta.servlet.http.HttpServletRequest; import org.springframework.validation.annotation.Validated; import org.springframework.web.bind.annotation.*; @@ -28,6 +29,9 @@ public class AdminMuseMetaContractController { /** 当前 Controller 拥有的 OpenAPI 合同域。 */ private static final Set DOMAINS = Set.of("meta"); + @Resource + private MuseContractPersistenceService contractPersistenceService; + /** * 处理当前领域下尚未落成专用 Controller 的 OpenAPI 合同请求。 * @@ -41,7 +45,7 @@ public class AdminMuseMetaContractController { public CommonResult> handle(HttpServletRequest request, @RequestParam Map queryParams, @RequestBody(required = false) Map body) { - return success(MuseApiContractSupport.handle(DOMAINS, "admin", request.getMethod(), request.getRequestURI(), + return success(contractPersistenceService.handle(DOMAINS, "admin", request.getMethod(), request.getRequestURI(), request.getHeader("X-Command-Id"), queryParams, body, getLoginUserId())); } diff --git a/muse-cloud/muse-server/src/test/java/cn/iocoder/muse/server/framework/api/MuseContractPersistenceServiceTest.java b/muse-cloud/muse-server/src/test/java/cn/iocoder/muse/server/framework/api/MuseContractPersistenceServiceTest.java new file mode 100644 index 00000000..79627edd --- /dev/null +++ b/muse-cloud/muse-server/src/test/java/cn/iocoder/muse/server/framework/api/MuseContractPersistenceServiceTest.java @@ -0,0 +1,172 @@ +package cn.iocoder.muse.server.framework.api; + +import cn.iocoder.muse.framework.common.exception.ServiceException; +import cn.iocoder.muse.framework.mybatis.config.MuseMybatisAutoConfiguration; +import cn.iocoder.muse.framework.mybatis.core.muse.MuseContractPersistenceService; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.context.annotation.Import; +import org.springframework.test.util.ReflectionTestUtils; + +import java.util.Arrays; +import java.util.Map; +import java.util.Set; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.*; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +/** + * Muse 合同持久化应用服务测试。 + */ +class MuseContractPersistenceServiceTest { + + private MuseContractPersistenceService service; + private JdbcTemplate jdbcTemplate; + + @BeforeEach + void setUp() { + service = new MuseContractPersistenceService(); + jdbcTemplate = mock(JdbcTemplate.class); + ReflectionTestUtils.setField(service, "jdbcTemplate", jdbcTemplate); + } + + @Test + void should_replay_existing_command_when_commandIdAlreadyPersisted() { + // 已持久化的 commandId 再次提交时,后端必须返回历史响应,避免重复扣费、重复安装或重复发布。 + when(jdbcTemplate.queryForObject(anyString(), eq(String.class), any(Object[].class))) + .thenReturn("{\"operationId\":\"purchaseAsset\",\"status\":\"accepted\",\"operationRecordId\":9}"); + + Map result = service.handle(Set.of("market"), "app", "POST", + "/app-api/muse/marketplace/assets/1/purchase", null, Map.of(), + Map.of("commandId", "cmd-market-1"), 1001L); + + assertEquals("purchaseAsset", result.get("operationId")); + assertEquals("accepted", result.get("status")); + assertTrue((Boolean) result.get("idempotentReplay")); + } + + @Test + void should_persist_ai_task_and_operation_audit_when_createAiTask() { + // 创建 AI 任务时需要同时落业务任务表和统一操作审计表,形成可追踪的 P1 后端事实。 + when(jdbcTemplate.queryForObject(anyString(), eq(String.class), any(Object[].class))) + .thenReturn(null); + when(jdbcTemplate.queryForObject(argThat(sql -> sql != null && sql.contains("muse_ai_generation")), eq(Long.class), any(Object[].class))) + .thenReturn(101L); + when(jdbcTemplate.queryForObject(argThat(sql -> sql != null && sql.contains("muse_domain_operation_record")), eq(Long.class), any(Object[].class))) + .thenReturn(201L); + when(jdbcTemplate.update(anyString(), any(Object[].class))).thenReturn(1); + + Map result = service.handle(Set.of("ai"), "app", "POST", + "/app-api/muse/ai/tasks", null, Map.of(), + Map.of("commandId", "cmd-ai-1", "intent", "outline"), 1001L); + + assertEquals("createAiTask", result.get("operationId")); + assertEquals(101L, result.get("taskId")); + assertEquals(201L, result.get("operationRecordId")); + assertEquals("aiTask", result.get("resourceType")); + } + + @Test + void should_persist_knowledge_base_when_createKnowledgeBase() { + // 用户侧创建知识库时不能只写统一操作审计,必须同时落知识库领域表。 + when(jdbcTemplate.queryForObject(anyString(), eq(String.class), any(Object[].class))) + .thenReturn(null); + when(jdbcTemplate.queryForObject(argThat(sql -> sql != null && sql.contains("muse_knowledge_base")), eq(Long.class), any(Object[].class))) + .thenReturn(301L); + when(jdbcTemplate.queryForObject(argThat(sql -> sql != null && sql.contains("muse_domain_operation_record")), eq(Long.class), any(Object[].class))) + .thenReturn(401L); + when(jdbcTemplate.update(anyString(), any(Object[].class))).thenReturn(1); + + Map result = service.handle(Set.of("knowledge"), "app", "POST", + "/app-api/muse/knowledge-bases", null, Map.of(), + Map.of("commandId", "cmd-kb-1", "name", "Story Bible"), 1001L); + + assertEquals("createKnowledgeBase", result.get("operationId")); + assertEquals(301L, result.get("kbId")); + assertEquals(401L, result.get("operationRecordId")); + assertEquals("knowledgeBase", result.get("resourceType")); + } + + @Test + void should_reject_meta_activation_when_expectedActiveVersionIsStale() { + // MetaSchema 激活/回滚必须校验调用方看到的当前活跃版本,避免过期治理命令覆盖新版本。 + when(jdbcTemplate.queryForObject(anyString(), eq(String.class), any(Object[].class))) + .thenReturn(null); + when(jdbcTemplate.queryForObject(argThat(sql -> sql != null && sql.contains("SELECT v.version_no")), eq(Integer.class), any(Object[].class))) + .thenReturn(2); + + assertThrows(ServiceException.class, () -> service.handle(Set.of("meta"), "admin", "POST", + "/admin-api/muse/governance/meta-schemas/work/versions/3/activate", null, Map.of(), + Map.of("commandId", "cmd-meta-1", "reason", "publish", "expectedActiveVersion", 1), 1001L)); + } + + @Test + void should_reject_unimplemented_write_command_insteadOfFakeAccepted() { + // 当前持久化服务未拥有的领域写命令不能返回 accepted,避免前端误以为动作已经生效。 + when(jdbcTemplate.queryForObject(anyString(), eq(String.class), any(Object[].class))) + .thenReturn(null); + + assertThrows(ServiceException.class, () -> service.handle(Set.of("content"), "app", "POST", + "/app-api/muse/works", null, Map.of(), + Map.of("commandId", "cmd-content-work-1", "title", "Draft"), 1001L)); + } + + @Test + void should_persist_market_purchase_when_purchaseAsset() { + // 购买是 Market 的明确业务事实,不能只写操作日志,需要落购买表并保留幂等 commandId。 + when(jdbcTemplate.queryForObject(anyString(), eq(String.class), any(Object[].class))) + .thenReturn(null); + when(jdbcTemplate.queryForObject(argThat(sql -> sql != null && sql.contains("muse_domain_operation_record")), eq(Long.class), any(Object[].class))) + .thenReturn(501L); + when(jdbcTemplate.queryForObject(argThat(sql -> sql != null && sql.contains("muse_market_asset_version")), eq(Long.class), any(Object[].class))) + .thenReturn(601L); + when(jdbcTemplate.queryForObject(argThat(sql -> sql != null && sql.contains("muse_market_purchase")), eq(Long.class), any(Object[].class))) + .thenReturn(701L); + when(jdbcTemplate.update(anyString(), any(Object[].class))).thenReturn(1); + + Map result = service.handle(Set.of("market"), "app", "POST", + "/app-api/muse/marketplace/assets/9/purchase", null, Map.of(), + Map.of("commandId", "cmd-market-purchase-1"), 1001L); + + assertEquals("purchaseAsset", result.get("operationId")); + assertEquals(701L, result.get("purchaseId")); + assertEquals(501L, result.get("operationRecordId")); + assertEquals("purchase", result.get("resourceType")); + } + + @Test + void should_persist_workflow_task_when_precheckAgentSlot() { + // 预检类命令本身就是可追踪工作流事实,需要落统一 workflow task,供后续 bind 使用。 + when(jdbcTemplate.queryForObject(anyString(), eq(String.class), any(Object[].class))) + .thenReturn(null); + when(jdbcTemplate.queryForObject(argThat(sql -> sql != null && sql.contains("COUNT(1) FROM muse_content_work")), eq(Integer.class), any(Object[].class))) + .thenReturn(1); + when(jdbcTemplate.queryForObject(argThat(sql -> sql != null && sql.contains("muse_domain_operation_record")), eq(Long.class), any(Object[].class))) + .thenReturn(801L); + when(jdbcTemplate.queryForObject(argThat(sql -> sql != null && sql.contains("muse_domain_workflow_task")), eq(Long.class), any(Object[].class))) + .thenReturn(901L); + when(jdbcTemplate.update(anyString(), any(Object[].class))).thenReturn(1); + + Map result = service.handle(Set.of("ai"), "app", "POST", + "/app-api/muse/works/88/agent-slots/draft/prechecks", null, Map.of(), + Map.of("commandId", "cmd-slot-precheck-1", "sourceAgentId", 12, "sourceAgentVersion", "1"), 1001L); + + assertEquals("precheckAgentSlot", result.get("operationId")); + assertEquals(901L, result.get("agentSlotPrecheckId")); + assertEquals("aiWorkflowTask", result.get("resourceType")); + } + + @Test + void should_import_contractPersistenceService_from_mybatisAutoConfiguration() { + // muse-server 只扫描 server/module 包,framework 下的合同持久化服务必须由 starter 自动导入。 + Import importAnnotation = MuseMybatisAutoConfiguration.class.getAnnotation(Import.class); + + assertTrue(Arrays.asList(importAnnotation.value()).contains(MuseContractPersistenceService.class)); + } + +} diff --git a/muse-cloud/sql/muse/V7__add_contract_operation_audit.sql b/muse-cloud/sql/muse/V7__add_contract_operation_audit.sql new file mode 100644 index 00000000..c63a7c3a --- /dev/null +++ b/muse-cloud/sql/muse/V7__add_contract_operation_audit.sql @@ -0,0 +1,42 @@ +-- Muse P1 合同操作审计与幂等 Schema(PostgreSQL) +-- 目的:为 Meta/Knowledge/Market/AI/Account 的全量合同入口提供统一幂等、审计和兜底查询状态。 + +CREATE TABLE muse_domain_operation_record ( + id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY, + domain VARCHAR(50) NOT NULL, + side VARCHAR(20) NOT NULL, + operation_id VARCHAR(120) NOT NULL, + command_id VARCHAR(120), + actor_user_id BIGINT, + resource_type VARCHAR(80) NOT NULL, + resource_id BIGINT, + resource_key VARCHAR(200), + parent_type VARCHAR(80), + parent_id BIGINT, + status VARCHAR(30) NOT NULL DEFAULT 'accepted', + revision INT NOT NULL DEFAULT 1, + path_variables JSONB, + query_params JSONB, + request_payload JSONB, + response_payload JSONB, + creator VARCHAR(64) NOT NULL DEFAULT '', + create_time TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, + updater VARCHAR(64) NOT NULL DEFAULT '', + update_time TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, + deleted BOOLEAN NOT NULL DEFAULT FALSE, + tenant_id BIGINT NOT NULL DEFAULT 0 +); + +CREATE UNIQUE INDEX uk_muse_domain_op_command + ON muse_domain_operation_record(tenant_id, domain, command_id) + WHERE command_id IS NOT NULL; + +CREATE INDEX idx_muse_domain_op_resource + ON muse_domain_operation_record(tenant_id, domain, resource_type, resource_id, create_time); + +CREATE INDEX idx_muse_domain_op_actor + ON muse_domain_operation_record(tenant_id, actor_user_id, create_time); + +CREATE TRIGGER trg_muse_domain_operation_record_updated_at + BEFORE UPDATE ON muse_domain_operation_record + FOR EACH ROW EXECUTE FUNCTION update_updated_at_column(); diff --git a/muse-cloud/sql/muse/V8__add_contract_workflow_and_market_interaction.sql b/muse-cloud/sql/muse/V8__add_contract_workflow_and_market_interaction.sql new file mode 100644 index 00000000..6dff4373 --- /dev/null +++ b/muse-cloud/sql/muse/V8__add_contract_workflow_and_market_interaction.sql @@ -0,0 +1,98 @@ +-- Muse P1 合同工作流与市场交互 Schema(PostgreSQL) +-- 目的:承载预检、导出、重试、发布准备等异步工作流事实,以及市场收藏/购买状态。 + +CREATE TABLE muse_domain_workflow_task ( + id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY, + domain VARCHAR(50) NOT NULL, + side VARCHAR(20) NOT NULL, + operation_id VARCHAR(120) NOT NULL, + task_type VARCHAR(80) NOT NULL, + status VARCHAR(30) NOT NULL DEFAULT 'queued', + actor_user_id BIGINT, + owner_user_id BIGINT, + target_type VARCHAR(80), + target_id BIGINT, + parent_type VARCHAR(80), + parent_id BIGINT, + correlation_id VARCHAR(200), + source_type VARCHAR(80), + source_id BIGINT, + expected_status VARCHAR(80), + request_payload JSONB, + result_payload JSONB, + command_id VARCHAR(120), + revision INT NOT NULL DEFAULT 1, + creator VARCHAR(64) NOT NULL DEFAULT '', + create_time TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, + updater VARCHAR(64) NOT NULL DEFAULT '', + update_time TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, + deleted BOOLEAN NOT NULL DEFAULT FALSE, + tenant_id BIGINT NOT NULL DEFAULT 0 +); + +CREATE UNIQUE INDEX uk_muse_domain_workflow_command + ON muse_domain_workflow_task(tenant_id, domain, command_id) + WHERE command_id IS NOT NULL; + +CREATE INDEX idx_muse_domain_workflow_target + ON muse_domain_workflow_task(tenant_id, domain, target_type, target_id, create_time); + +CREATE INDEX idx_muse_domain_workflow_actor + ON muse_domain_workflow_task(tenant_id, domain, actor_user_id, create_time); + +CREATE TRIGGER trg_muse_domain_workflow_task_updated_at + BEFORE UPDATE ON muse_domain_workflow_task + FOR EACH ROW EXECUTE FUNCTION update_updated_at_column(); + +CREATE TABLE muse_market_favorite ( + id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY, + asset_id BIGINT NOT NULL, + user_id BIGINT NOT NULL, + status VARCHAR(20) NOT NULL DEFAULT 'active', + command_id VARCHAR(120), + creator VARCHAR(64) NOT NULL DEFAULT '', + create_time TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, + updater VARCHAR(64) NOT NULL DEFAULT '', + update_time TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, + deleted BOOLEAN NOT NULL DEFAULT FALSE, + tenant_id BIGINT NOT NULL DEFAULT 0, + CONSTRAINT uk_muse_market_favorite_user_asset UNIQUE (tenant_id, user_id, asset_id) +); + +CREATE UNIQUE INDEX uk_muse_market_favorite_command + ON muse_market_favorite(tenant_id, command_id) + WHERE command_id IS NOT NULL; + +CREATE INDEX idx_muse_market_favorite_user + ON muse_market_favorite(tenant_id, user_id, status); + +CREATE TRIGGER trg_muse_market_favorite_updated_at + BEFORE UPDATE ON muse_market_favorite + FOR EACH ROW EXECUTE FUNCTION update_updated_at_column(); + +CREATE TABLE muse_market_purchase ( + id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY, + asset_id BIGINT NOT NULL, + asset_version_id BIGINT, + user_id BIGINT NOT NULL, + status VARCHAR(20) NOT NULL DEFAULT 'completed', + purchase_payload JSONB, + command_id VARCHAR(120) NOT NULL, + creator VARCHAR(64) NOT NULL DEFAULT '', + create_time TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, + updater VARCHAR(64) NOT NULL DEFAULT '', + update_time TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, + deleted BOOLEAN NOT NULL DEFAULT FALSE, + tenant_id BIGINT NOT NULL DEFAULT 0, + CONSTRAINT uk_muse_market_purchase_command UNIQUE (tenant_id, command_id) +); + +CREATE INDEX idx_muse_market_purchase_user + ON muse_market_purchase(tenant_id, user_id, create_time); + +CREATE INDEX idx_muse_market_purchase_asset + ON muse_market_purchase(tenant_id, asset_id, status); + +CREATE TRIGGER trg_muse_market_purchase_updated_at + BEFORE UPDATE ON muse_market_purchase + FOR EACH ROW EXECUTE FUNCTION update_updated_at_column();