diff --git a/muse-cloud/muse-module-content/muse-module-content-server/src/main/java/cn/iocoder/muse/module/content/controller/app/AppMuseEventsController.java b/muse-cloud/muse-module-content/muse-module-content-server/src/main/java/cn/iocoder/muse/module/content/controller/app/AppMuseEventsController.java
deleted file mode 100644
index 4a2f7bab..00000000
--- a/muse-cloud/muse-module-content/muse-module-content-server/src/main/java/cn/iocoder/muse/module/content/controller/app/AppMuseEventsController.java
+++ /dev/null
@@ -1,51 +0,0 @@
-package cn.iocoder.muse.module.content.controller.app;
-
-import io.swagger.v3.oas.annotations.Operation;
-import io.swagger.v3.oas.annotations.tags.Tag;
-import org.springframework.http.MediaType;
-import org.springframework.web.bind.annotation.GetMapping;
-import org.springframework.web.bind.annotation.RequestMapping;
-import org.springframework.web.bind.annotation.RestController;
-import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
-
-import java.io.IOException;
-import java.time.LocalDateTime;
-import java.util.Map;
-
-/**
- * 用户端 Muse 统一事件流入口。
- *
- *
P1 阶段先提供可建立连接的 SSE 端点,后续由各业务模块接入真实事件发布器。
- */
-@Tag(name = "用户 APP - Muse Events")
-@RestController
-@RequestMapping("/muse")
-public class AppMuseEventsController {
-
- /** SSE 连接默认超时时间,避免占用容器线程过久。 */
- private static final long DEFAULT_TIMEOUT_MILLIS = 30_000L;
-
- /**
- * 建立统一事件流 SSE 连接。
- *
- * @return SSE 事件流
- */
- @GetMapping(value = "/events", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
- @Operation(summary = "统一事件流 SSE")
- public SseEmitter streamEvents() {
- SseEmitter emitter = new SseEmitter(DEFAULT_TIMEOUT_MILLIS);
- try {
- emitter.send(SseEmitter.event()
- .name("notification")
- .data(Map.of(
- "type", "source_status_change",
- "message", "muse event stream ready",
- "timestamp", LocalDateTime.now().toString())));
- emitter.complete();
- } catch (IOException ex) {
- emitter.completeWithError(ex);
- }
- return emitter;
- }
-
-}
diff --git a/muse-cloud/muse-module-events/muse-module-events-api/pom.xml b/muse-cloud/muse-module-events/muse-module-events-api/pom.xml
new file mode 100644
index 00000000..394d3f5e
--- /dev/null
+++ b/muse-cloud/muse-module-events/muse-module-events-api/pom.xml
@@ -0,0 +1,36 @@
+
+
+
+ cn.iocoder.cloud
+ muse-module-events
+ ${revision}
+
+ 4.0.0
+ muse-module-events-api
+ jar
+ ${project.artifactId}
+ events 模块 API 和发布契约。
+
+
+ cn.iocoder.cloud
+ muse-common
+
+
+ org.springdoc
+ springdoc-openapi-starter-webmvc-ui
+ provided
+
+
+ org.springframework.boot
+ spring-boot-starter-validation
+ true
+
+
+ org.springframework.cloud
+ spring-cloud-starter-openfeign
+ true
+
+
+
diff --git a/muse-cloud/muse-module-events/muse-module-events-api/src/main/java/cn/iocoder/muse/module/events/api/publish/EventsPublishApi.java b/muse-cloud/muse-module-events/muse-module-events-api/src/main/java/cn/iocoder/muse/module/events/api/publish/EventsPublishApi.java
new file mode 100644
index 00000000..bc92c5e4
--- /dev/null
+++ b/muse-cloud/muse-module-events/muse-module-events-api/src/main/java/cn/iocoder/muse/module/events/api/publish/EventsPublishApi.java
@@ -0,0 +1,24 @@
+package cn.iocoder.muse.module.events.api.publish;
+
+import cn.iocoder.muse.framework.common.pojo.CommonResult;
+import cn.iocoder.muse.module.events.api.publish.dto.EventsPublishReqDTO;
+import cn.iocoder.muse.module.events.api.publish.dto.EventsPublishRespDTO;
+import cn.iocoder.muse.module.events.enums.ApiConstants;
+import io.swagger.v3.oas.annotations.Operation;
+import io.swagger.v3.oas.annotations.tags.Tag;
+import jakarta.validation.Valid;
+import org.springframework.cloud.openfeign.FeignClient;
+import org.springframework.web.bind.annotation.PostMapping;
+import org.springframework.web.bind.annotation.RequestBody;
+
+@FeignClient(name = ApiConstants.NAME)
+@Tag(name = "RPC 服务 - Events 统一事件发布")
+public interface EventsPublishApi {
+
+ String PREFIX = ApiConstants.PREFIX + "/publish";
+
+ @PostMapping(PREFIX)
+ @Operation(summary = "幂等发布统一事件投影")
+ CommonResult publish(@Valid @RequestBody EventsPublishReqDTO reqDTO);
+
+}
diff --git a/muse-cloud/muse-module-events/muse-module-events-api/src/main/java/cn/iocoder/muse/module/events/api/publish/dto/EventsPublishReqDTO.java b/muse-cloud/muse-module-events/muse-module-events-api/src/main/java/cn/iocoder/muse/module/events/api/publish/dto/EventsPublishReqDTO.java
new file mode 100644
index 00000000..eb4f36ea
--- /dev/null
+++ b/muse-cloud/muse-module-events/muse-module-events-api/src/main/java/cn/iocoder/muse/module/events/api/publish/dto/EventsPublishReqDTO.java
@@ -0,0 +1,50 @@
+package cn.iocoder.muse.module.events.api.publish.dto;
+
+import jakarta.validation.constraints.NotEmpty;
+import jakarta.validation.constraints.NotNull;
+import lombok.Data;
+
+import java.io.Serializable;
+import java.time.LocalDateTime;
+import java.util.Map;
+
+/**
+ * Events 统一事件发布请求。
+ */
+@Data
+public class EventsPublishReqDTO implements Serializable {
+
+ @NotEmpty(message = "commandId 不能为空")
+ private String commandId;
+
+ @NotNull(message = "tenantId 不能为空")
+ private Long tenantId;
+
+ @NotNull(message = "ownerUserId 不能为空")
+ private Long ownerUserId;
+
+ @NotEmpty(message = "sourceOwner 不能为空")
+ private String sourceOwner;
+
+ @NotEmpty(message = "sourceType 不能为空")
+ private String sourceType;
+
+ @NotEmpty(message = "sourceId 不能为空")
+ private String sourceId;
+
+ private String sourceRevision;
+
+ @NotEmpty(message = "eventType 不能为空")
+ private String eventType;
+
+ private String resourceType;
+
+ private String resourceId;
+
+ @NotNull(message = "payloadSummary 不能为空")
+ private Map payloadSummary;
+
+ @NotNull(message = "emittedAt 不能为空")
+ private LocalDateTime emittedAt;
+
+}
diff --git a/muse-cloud/muse-module-events/muse-module-events-api/src/main/java/cn/iocoder/muse/module/events/api/publish/dto/EventsPublishRespDTO.java b/muse-cloud/muse-module-events/muse-module-events-api/src/main/java/cn/iocoder/muse/module/events/api/publish/dto/EventsPublishRespDTO.java
new file mode 100644
index 00000000..f9446678
--- /dev/null
+++ b/muse-cloud/muse-module-events/muse-module-events-api/src/main/java/cn/iocoder/muse/module/events/api/publish/dto/EventsPublishRespDTO.java
@@ -0,0 +1,23 @@
+package cn.iocoder.muse.module.events.api.publish.dto;
+
+import lombok.Data;
+
+import java.io.Serializable;
+
+/**
+ * Events 统一事件发布响应。
+ */
+@Data
+public class EventsPublishRespDTO implements Serializable {
+
+ private String eventId;
+
+ private Long sequenceNo;
+
+ private String publishStatus;
+
+ private String publishErrorCode;
+
+ private Boolean duplicate;
+
+}
diff --git a/muse-cloud/muse-module-events/muse-module-events-api/src/main/java/cn/iocoder/muse/module/events/enums/ApiConstants.java b/muse-cloud/muse-module-events/muse-module-events-api/src/main/java/cn/iocoder/muse/module/events/enums/ApiConstants.java
new file mode 100644
index 00000000..9de03943
--- /dev/null
+++ b/muse-cloud/muse-module-events/muse-module-events-api/src/main/java/cn/iocoder/muse/module/events/enums/ApiConstants.java
@@ -0,0 +1,19 @@
+package cn.iocoder.muse.module.events.enums;
+
+import cn.iocoder.muse.framework.common.enums.RpcConstants;
+
+/**
+ * Events 模块 RPC API 常量。
+ */
+public class ApiConstants {
+
+ /**
+ * 服务名需要和承载 Events server 的应用名保持一致。
+ */
+ public static final String NAME = "events-server";
+
+ public static final String PREFIX = RpcConstants.RPC_API_PREFIX + "/events";
+
+ public static final String VERSION = "1.0.0";
+
+}
diff --git a/muse-cloud/muse-module-events/muse-module-events-api/src/main/java/cn/iocoder/muse/module/events/enums/ErrorCodeConstants.java b/muse-cloud/muse-module-events/muse-module-events-api/src/main/java/cn/iocoder/muse/module/events/enums/ErrorCodeConstants.java
new file mode 100644
index 00000000..6acc6ee5
--- /dev/null
+++ b/muse-cloud/muse-module-events/muse-module-events-api/src/main/java/cn/iocoder/muse/module/events/enums/ErrorCodeConstants.java
@@ -0,0 +1,16 @@
+package cn.iocoder.muse.module.events.enums;
+
+import cn.iocoder.muse.framework.common.exception.ErrorCode;
+
+/**
+ * Events 错误码。
+ *
+ * Events 使用 1-045-000-000 段,覆盖统一事件发布、幂等投影和 SSE 可见性边界。
+ */
+public interface ErrorCodeConstants {
+
+ // ========== Events 发布契约 1-045-000-000 ==========
+ ErrorCode EVENTS_PUBLISH_COMMAND_ID_REQUIRED = new ErrorCode(1_045_000_000, "事件发布 commandId 不能为空");
+ ErrorCode EVENTS_PUBLISH_REQUIRED_FIELD_EMPTY = new ErrorCode(1_045_000_001, "事件发布必要字段不能为空:{}");
+
+}
diff --git a/muse-cloud/muse-module-events/muse-module-events-server/pom.xml b/muse-cloud/muse-module-events/muse-module-events-server/pom.xml
new file mode 100644
index 00000000..a76f56ef
--- /dev/null
+++ b/muse-cloud/muse-module-events/muse-module-events-server/pom.xml
@@ -0,0 +1,47 @@
+
+
+
+ cn.iocoder.cloud
+ muse-module-events
+ ${revision}
+
+ 4.0.0
+ muse-module-events-server
+ jar
+ ${project.artifactId}
+ events 模块服务实现。
+
+
+ cn.iocoder.cloud
+ muse-module-events-api
+ ${revision}
+
+
+ cn.iocoder.cloud
+ muse-spring-boot-starter-web
+
+
+ cn.iocoder.cloud
+ muse-spring-boot-starter-security
+
+
+ cn.iocoder.cloud
+ muse-spring-boot-starter-biz-tenant
+
+
+ cn.iocoder.cloud
+ muse-spring-boot-starter-mybatis
+
+
+ org.springframework.boot
+ spring-boot-starter-validation
+
+
+ cn.iocoder.cloud
+ muse-spring-boot-starter-test
+ test
+
+
+
diff --git a/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/api/publish/EventsPublishApiImpl.java b/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/api/publish/EventsPublishApiImpl.java
new file mode 100644
index 00000000..098e09be
--- /dev/null
+++ b/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/api/publish/EventsPublishApiImpl.java
@@ -0,0 +1,27 @@
+package cn.iocoder.muse.module.events.api.publish;
+
+import cn.iocoder.muse.framework.common.pojo.CommonResult;
+import cn.iocoder.muse.module.events.api.publish.EventsPublishApi;
+import cn.iocoder.muse.module.events.api.publish.dto.EventsPublishReqDTO;
+import cn.iocoder.muse.module.events.api.publish.dto.EventsPublishRespDTO;
+import cn.iocoder.muse.module.events.application.publish.EventsPublishService;
+import org.springframework.validation.annotation.Validated;
+import org.springframework.web.bind.annotation.RestController;
+
+import jakarta.annotation.Resource;
+
+import static cn.iocoder.muse.framework.common.pojo.CommonResult.success;
+
+@RestController
+@Validated
+public class EventsPublishApiImpl implements EventsPublishApi {
+
+ @Resource
+ private EventsPublishService eventsPublishService;
+
+ @Override
+ public CommonResult publish(EventsPublishReqDTO reqDTO) {
+ return success(eventsPublishService.publish(reqDTO));
+ }
+
+}
diff --git a/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/application/publish/EventsPublishService.java b/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/application/publish/EventsPublishService.java
new file mode 100644
index 00000000..1fc8918f
--- /dev/null
+++ b/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/application/publish/EventsPublishService.java
@@ -0,0 +1,19 @@
+package cn.iocoder.muse.module.events.application.publish;
+
+import cn.iocoder.muse.module.events.api.publish.dto.EventsPublishReqDTO;
+import cn.iocoder.muse.module.events.api.publish.dto.EventsPublishRespDTO;
+import cn.iocoder.muse.module.events.dal.dataobject.UnifiedEventDO;
+
+import java.time.LocalDateTime;
+import java.util.List;
+
+/**
+ * Events 发布应用服务。
+ */
+public interface EventsPublishService {
+
+ EventsPublishRespDTO publish(EventsPublishReqDTO reqDTO);
+
+ List listVisibleEvents(Long tenantId, Long ownerUserId, Long afterSequenceNo, LocalDateTime now);
+
+}
diff --git a/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/application/publish/EventsPublishServiceImpl.java b/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/application/publish/EventsPublishServiceImpl.java
new file mode 100644
index 00000000..881232ae
--- /dev/null
+++ b/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/application/publish/EventsPublishServiceImpl.java
@@ -0,0 +1,286 @@
+package cn.iocoder.muse.module.events.application.publish;
+
+import cn.iocoder.muse.module.events.api.publish.dto.EventsPublishReqDTO;
+import cn.iocoder.muse.module.events.api.publish.dto.EventsPublishRespDTO;
+import cn.iocoder.muse.module.events.dal.dataobject.UnifiedEventDO;
+import cn.iocoder.muse.module.events.dal.mysql.UnifiedEventMapper;
+import cn.iocoder.muse.module.events.domain.EventsPayloadSanitizer;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Service;
+import org.springframework.transaction.annotation.Transactional;
+import org.springframework.util.StringUtils;
+
+import java.math.BigDecimal;
+import java.math.BigInteger;
+import java.time.LocalDateTime;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import java.util.UUID;
+
+/**
+ * Events 发布应用服务实现。
+ */
+@Service
+@Slf4j
+public class EventsPublishServiceImpl implements EventsPublishService {
+
+ public static final String SOURCE_REVISION_NONE = "__none__";
+
+ public static final String PUBLISH_STATUS_ACCEPTED = "accepted";
+ public static final String PUBLISH_STATUS_REJECTED = "rejected";
+ public static final String PUBLISH_STATUS_BLOCKED = "blocked";
+
+ public static final String ERROR_PAYLOAD_CONTAINS_SECRET = "payload_contains_secret";
+ public static final String ERROR_EVENT_TYPE_NOT_DECLARED = "event_type_not_declared";
+ public static final String ERROR_PAYLOAD_SCHEMA_MISMATCH = "payload_schema_mismatch";
+
+ private static final BigInteger LONG_MIN_VALUE = BigInteger.valueOf(Long.MIN_VALUE);
+ private static final BigInteger LONG_MAX_VALUE = BigInteger.valueOf(Long.MAX_VALUE);
+
+ private static final String EVENT_TYPE_CHUNK = "chunk";
+ private static final String EVENT_TYPE_QUALITY_CHECK = "quality_check";
+ private static final String EVENT_TYPE_DONE = "done";
+ private static final String EVENT_TYPE_ERROR = "error";
+ private static final String EVENT_TYPE_NOTIFICATION = "notification";
+
+ private static final Set DECLARED_EVENT_TYPES = Set.of(EVENT_TYPE_CHUNK, EVENT_TYPE_QUALITY_CHECK,
+ EVENT_TYPE_DONE, EVENT_TYPE_ERROR, EVENT_TYPE_NOTIFICATION);
+ private static final Set NOTIFICATION_TYPES = Set.of("source_status_change", "knowledge_projection_done",
+ "governance_action", "quota_alert");
+
+ private final UnifiedEventMapper unifiedEventMapper;
+ private final EventsPayloadSanitizer payloadSanitizer;
+
+ public EventsPublishServiceImpl(UnifiedEventMapper unifiedEventMapper) {
+ this(unifiedEventMapper, new EventsPayloadSanitizer());
+ }
+
+ @Autowired
+ public EventsPublishServiceImpl(UnifiedEventMapper unifiedEventMapper, EventsPayloadSanitizer payloadSanitizer) {
+ this.unifiedEventMapper = unifiedEventMapper;
+ this.payloadSanitizer = payloadSanitizer;
+ }
+
+ @Override
+ @Transactional(rollbackFor = Exception.class)
+ public EventsPublishRespDTO publish(EventsPublishReqDTO reqDTO) {
+ validateRequiredFields(reqDTO);
+ String sourceRevision = normalizeSourceRevision(reqDTO.getSourceRevision());
+
+ // 幂等顺序是合同的一部分:先按 commandId 回放,再按 source tuple 回放,避免 source owner retry 生成重复事件。
+ UnifiedEventDO existingByCommand = unifiedEventMapper.selectByTenantIdAndCommandId(reqDTO.getTenantId(),
+ reqDTO.getCommandId());
+ if (existingByCommand != null) {
+ return toResp(existingByCommand, true);
+ }
+ UnifiedEventDO existingBySourceTuple = unifiedEventMapper.selectByTenantIdAndSourceTuple(reqDTO.getTenantId(),
+ reqDTO.getSourceOwner(), reqDTO.getSourceType(), reqDTO.getSourceId(), sourceRevision,
+ normalizeEventTypeForLookup(reqDTO.getEventType()));
+ if (existingBySourceTuple != null) {
+ return toResp(existingBySourceTuple, true);
+ }
+
+ UnifiedEventDO event = buildEvent(reqDTO, sourceRevision);
+ int inserted = unifiedEventMapper.insertIgnore(event);
+ UnifiedEventDO stored = findStoredEvent(reqDTO, sourceRevision);
+ if (stored == null) {
+ // insertIgnore 被 DB 约束吞掉但回查失败,说明调用方传入了互相冲突的 command/source tuple;fail-closed 记录日志。
+ log.warn("Events 发布幂等回查失败: tenantId={}, commandId={}, sourceOwner={}, sourceType={}, sourceId={}, sourceRevision={}, eventType={}, inserted={}",
+ reqDTO.getTenantId(), reqDTO.getCommandId(), reqDTO.getSourceOwner(), reqDTO.getSourceType(),
+ reqDTO.getSourceId(), sourceRevision, reqDTO.getEventType(), inserted);
+ throw new IllegalStateException("Events 发布幂等回查失败,拒绝返回未落库事件");
+ }
+ return toResp(stored, inserted == 0);
+ }
+
+ @Override
+ @Transactional(readOnly = true)
+ public List listVisibleEvents(Long tenantId, Long ownerUserId, Long afterSequenceNo,
+ LocalDateTime now) {
+ return unifiedEventMapper.selectVisibleEventsForOwner(tenantId, ownerUserId,
+ afterSequenceNo == null ? 0L : afterSequenceNo, now == null ? LocalDateTime.now() : now);
+ }
+
+ private UnifiedEventDO buildEvent(EventsPublishReqDTO reqDTO, String sourceRevision) {
+ EventsPayloadSanitizer.SanitizedPayload sanitizedPayload = payloadSanitizer.sanitize(reqDTO.getPayloadSummary());
+ PublishValidation validation = validatePublishPayload(reqDTO, sanitizedPayload);
+
+ UnifiedEventDO event = new UnifiedEventDO();
+ event.setTenantId(reqDTO.getTenantId());
+ event.setCommandId(reqDTO.getCommandId());
+ event.setEventId("evt_" + UUID.randomUUID());
+ event.setEventType(validation.persistedEventType);
+ event.setOwnerUserId(reqDTO.getOwnerUserId());
+ event.setSourceOwner(reqDTO.getSourceOwner());
+ event.setSourceType(reqDTO.getSourceType());
+ event.setSourceId(reqDTO.getSourceId());
+ event.setSourceRevision(sourceRevision);
+ event.setResourceType(reqDTO.getResourceType());
+ event.setResourceId(reqDTO.getResourceId());
+ event.setPayloadSummary(sanitizedPayload.getPayloadJson());
+ event.setPublishStatus(validation.publishStatus);
+ event.setPublishErrorCode(validation.publishErrorCode);
+ // 可见时间只对 accepted 有意义;rejected/blocked 仍写入统一审计记录,但 mapper 查询会 fail-closed 排除。
+ event.setVisibleFrom(LocalDateTime.now());
+ event.setEmittedAt(reqDTO.getEmittedAt());
+ event.setDeleted(false);
+ return event;
+ }
+
+ private PublishValidation validatePublishPayload(EventsPublishReqDTO reqDTO,
+ EventsPayloadSanitizer.SanitizedPayload sanitizedPayload) {
+ if (!DECLARED_EVENT_TYPES.contains(reqDTO.getEventType())) {
+ // V16 DDL 对 event_type 有 CHECK 约束;非法类型必须记录 rejected,但不能把非法值写进 event_type。
+ return PublishValidation.rejected(EVENT_TYPE_ERROR, ERROR_EVENT_TYPE_NOT_DECLARED);
+ }
+ if (sanitizedPayload.isContainsSecret()) {
+ return PublishValidation.rejected(reqDTO.getEventType(), ERROR_PAYLOAD_CONTAINS_SECRET);
+ }
+ if (!matchesOpenApiEventSchema(reqDTO.getEventType(), reqDTO.getPayloadSummary())) {
+ return PublishValidation.rejected(reqDTO.getEventType(), ERROR_PAYLOAD_SCHEMA_MISMATCH);
+ }
+ return PublishValidation.accepted(reqDTO.getEventType());
+ }
+
+ private boolean matchesOpenApiEventSchema(String eventType, Map payload) {
+ Map safePayload = payload == null ? Map.of() : payload;
+ return switch (eventType) {
+ case EVENT_TYPE_CHUNK -> isRequiredString(safePayload, "content")
+ && isPositiveLongNumber(safePayload.get("sequenceNo"));
+ case EVENT_TYPE_QUALITY_CHECK -> isRequiredString(safePayload, "dimension")
+ && isScoreInOpenApiRange(safePayload.get("score"))
+ && safePayload.get("passed") instanceof Boolean;
+ case EVENT_TYPE_DONE -> isLongNumber(safePayload.get("taskId"))
+ && isLongNumber(safePayload.get("suggestionId"));
+ case EVENT_TYPE_ERROR -> isRequiredString(safePayload, "code")
+ && isRequiredString(safePayload, "message");
+ case EVENT_TYPE_NOTIFICATION -> safePayload.get("type") instanceof String type
+ && NOTIFICATION_TYPES.contains(type)
+ && isRequiredString(safePayload, "message");
+ default -> false;
+ };
+ }
+
+ private UnifiedEventDO findStoredEvent(EventsPublishReqDTO reqDTO, String sourceRevision) {
+ UnifiedEventDO existingByCommand = unifiedEventMapper.selectByTenantIdAndCommandId(reqDTO.getTenantId(),
+ reqDTO.getCommandId());
+ if (existingByCommand != null) {
+ return existingByCommand;
+ }
+ return unifiedEventMapper.selectByTenantIdAndSourceTuple(reqDTO.getTenantId(), reqDTO.getSourceOwner(),
+ reqDTO.getSourceType(), reqDTO.getSourceId(), sourceRevision,
+ normalizeEventTypeForLookup(reqDTO.getEventType()));
+ }
+
+ private String normalizeEventTypeForLookup(String eventType) {
+ return DECLARED_EVENT_TYPES.contains(eventType) ? eventType : EVENT_TYPE_ERROR;
+ }
+
+ private String normalizeSourceRevision(String sourceRevision) {
+ return StringUtils.hasText(sourceRevision) ? sourceRevision : SOURCE_REVISION_NONE;
+ }
+
+ private void validateRequiredFields(EventsPublishReqDTO reqDTO) {
+ if (reqDTO == null) {
+ throw new IllegalArgumentException("Events publish 请求不能为空");
+ }
+ requireText(reqDTO.getCommandId(), "commandId");
+ requireNonNull(reqDTO.getTenantId(), "tenantId");
+ requireNonNull(reqDTO.getOwnerUserId(), "ownerUserId");
+ requireText(reqDTO.getSourceOwner(), "sourceOwner");
+ requireText(reqDTO.getSourceType(), "sourceType");
+ requireText(reqDTO.getSourceId(), "sourceId");
+ requireText(reqDTO.getEventType(), "eventType");
+ requireNonNull(reqDTO.getPayloadSummary(), "payloadSummary");
+ requireNonNull(reqDTO.getEmittedAt(), "emittedAt");
+ }
+
+ private void requireText(String value, String fieldName) {
+ if (!StringUtils.hasText(value)) {
+ throw new IllegalArgumentException(fieldName + " 不能为空");
+ }
+ }
+
+ private void requireNonNull(Object value, String fieldName) {
+ if (value == null) {
+ throw new IllegalArgumentException(fieldName + " 不能为空");
+ }
+ }
+
+ private boolean isRequiredString(Map payload, String fieldName) {
+ // OpenAPI 未声明 minLength 的 string 只校验字段存在、非 null 且类型为 String;空字符串是合法合同值。
+ return payload.containsKey(fieldName) && payload.get(fieldName) instanceof String;
+ }
+
+ private boolean isLongNumber(Object value) {
+ return toLongValue(value) != null;
+ }
+
+ private boolean isPositiveLongNumber(Object value) {
+ Long longValue = toLongValue(value);
+ return longValue != null && longValue >= 1L;
+ }
+
+ private Long toLongValue(Object value) {
+ if (!(value instanceof Number number)) {
+ return null;
+ }
+ if (number instanceof Byte || number instanceof Short || number instanceof Integer || number instanceof Long) {
+ return number.longValue();
+ }
+ if (number instanceof BigInteger bigInteger) {
+ return isLongRange(bigInteger) ? bigInteger.longValue() : null;
+ }
+ if (number instanceof BigDecimal bigDecimal) {
+ return toLongValue(bigDecimal);
+ }
+ // OpenAPI int64 只接受整数语义的 Number;Float/Double 即使值为 1.0 也不作为 int64 入库。
+ return null;
+ }
+
+ private Long toLongValue(BigDecimal bigDecimal) {
+ try {
+ BigInteger integerValue = bigDecimal.toBigIntegerExact();
+ return isLongRange(integerValue) ? integerValue.longValue() : null;
+ } catch (ArithmeticException ex) {
+ return null;
+ }
+ }
+
+ private boolean isLongRange(BigInteger value) {
+ return value.compareTo(LONG_MIN_VALUE) >= 0 && value.compareTo(LONG_MAX_VALUE) <= 0;
+ }
+
+ private boolean isScoreInOpenApiRange(Object value) {
+ if (!(value instanceof Number number)) {
+ return false;
+ }
+ double score = number.doubleValue();
+ return Double.isFinite(score) && score >= 0D && score <= 1D;
+ }
+
+ private EventsPublishRespDTO toResp(UnifiedEventDO event, boolean duplicate) {
+ EventsPublishRespDTO respDTO = new EventsPublishRespDTO();
+ respDTO.setEventId(event.getEventId());
+ respDTO.setSequenceNo(event.getSequenceNo());
+ respDTO.setPublishStatus(event.getPublishStatus());
+ respDTO.setPublishErrorCode(event.getPublishErrorCode());
+ respDTO.setDuplicate(duplicate);
+ return respDTO;
+ }
+
+ private record PublishValidation(String publishStatus, String publishErrorCode, String persistedEventType) {
+
+ static PublishValidation accepted(String eventType) {
+ return new PublishValidation(PUBLISH_STATUS_ACCEPTED, null, eventType);
+ }
+
+ static PublishValidation rejected(String eventType, String errorCode) {
+ return new PublishValidation(PUBLISH_STATUS_REJECTED, errorCode, eventType);
+ }
+
+ }
+
+}
diff --git a/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/application/stream/EventsStreamService.java b/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/application/stream/EventsStreamService.java
new file mode 100644
index 00000000..1980c0c9
--- /dev/null
+++ b/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/application/stream/EventsStreamService.java
@@ -0,0 +1,20 @@
+package cn.iocoder.muse.module.events.application.stream;
+
+import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
+
+/**
+ * Events 统一 SSE stream 应用服务。
+ */
+public interface EventsStreamService {
+
+ /**
+ * 建立统一事件流。
+ *
+ * @param loginUserId 当前登录用户;必须由 controller 在进入 SSE 生命周期前取得
+ * @param apiVersion `X-API-Version` 请求头
+ * @param lastEventId SSE Last-Event-ID / 查询参数游标
+ * @return SSE emitter
+ */
+ SseEmitter streamEvents(Long loginUserId, String apiVersion, String lastEventId);
+
+}
diff --git a/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/application/stream/EventsStreamServiceImpl.java b/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/application/stream/EventsStreamServiceImpl.java
new file mode 100644
index 00000000..7a4e188e
--- /dev/null
+++ b/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/application/stream/EventsStreamServiceImpl.java
@@ -0,0 +1,414 @@
+package cn.iocoder.muse.module.events.application.stream;
+
+import cn.iocoder.muse.framework.tenant.core.context.TenantContextHolder;
+import cn.iocoder.muse.module.events.dal.dataobject.UnifiedEventDO;
+import cn.iocoder.muse.module.events.dal.mysql.UnifiedEventMapper;
+import cn.iocoder.muse.module.events.domain.EventsCursor;
+import cn.iocoder.muse.module.events.domain.EventsPayloadSanitizer;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.core.task.AsyncTaskExecutor;
+import org.springframework.http.MediaType;
+import org.springframework.stereotype.Service;
+import org.springframework.util.StringUtils;
+import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
+
+import java.io.IOException;
+import java.time.LocalDateTime;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.Set;
+import java.util.concurrent.Future;
+import java.util.concurrent.RejectedExecutionException;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicReference;
+import java.util.regex.Pattern;
+
+import static cn.iocoder.muse.module.events.framework.config.EventsStreamConfiguration.EVENTS_STREAM_EXECUTOR;
+
+/**
+ * Events 统一 SSE stream 应用服务实现。
+ */
+@Service
+@Slf4j
+public class EventsStreamServiceImpl implements EventsStreamService {
+
+ private static final String SUPPORTED_API_VERSION = "1";
+ private static final long DEFAULT_TIMEOUT_MILLIS = 30_000L;
+ private static final long HEARTBEAT_INTERVAL_MILLIS = 5_000L;
+
+ private static final String EVENT_CHUNK = "chunk";
+ private static final String EVENT_QUALITY_CHECK = "quality_check";
+ private static final String EVENT_DONE = "done";
+ private static final String EVENT_ERROR = "error";
+ private static final String EVENT_NOTIFICATION = "notification";
+
+ private static final String ERROR_API_VERSION_UNSUPPORTED = "EVENTS_API_VERSION_UNSUPPORTED";
+ private static final String ERROR_STREAM_UNAVAILABLE = "EVENTS_STREAM_UNAVAILABLE";
+
+ private static final Set DECLARED_EVENT_TYPES = Set.of(EVENT_CHUNK, EVENT_QUALITY_CHECK, EVENT_DONE,
+ EVENT_ERROR, EVENT_NOTIFICATION);
+ private static final Set NOTIFICATION_TYPES = Set.of("source_status_change", "knowledge_projection_done",
+ "governance_action", "quota_alert");
+ private static final Pattern SECRET_TEXT_PATTERN = Pattern.compile(
+ "(?i)(bearer\\s+|authorization|token|secret|api[-_]?key|provider\\s+raw\\s+body|raw\\s+body)");
+
+ private final UnifiedEventMapper unifiedEventMapper;
+ private final EventsPayloadSanitizer payloadSanitizer;
+ private final AsyncTaskExecutor eventsStreamExecutor;
+
+ @Autowired
+ public EventsStreamServiceImpl(UnifiedEventMapper unifiedEventMapper, EventsPayloadSanitizer payloadSanitizer,
+ @org.springframework.beans.factory.annotation.Qualifier(EVENTS_STREAM_EXECUTOR)
+ AsyncTaskExecutor eventsStreamExecutor) {
+ this.unifiedEventMapper = unifiedEventMapper;
+ this.payloadSanitizer = payloadSanitizer;
+ this.eventsStreamExecutor = eventsStreamExecutor;
+ }
+
+ @Override
+ public SseEmitter streamEvents(Long loginUserId, String apiVersion, String lastEventId) {
+ Long tenantId = TenantContextHolder.getTenantId();
+ SseEmitter emitter = new SseEmitter(DEFAULT_TIMEOUT_MILLIS);
+ AtomicBoolean active = new AtomicBoolean(true);
+ AtomicReference> pollingFutureRef = new AtomicReference<>();
+ registerLifecycleCallbacks(emitter, active, pollingFutureRef);
+
+ if (!SUPPORTED_API_VERSION.equals(apiVersion)) {
+ // 不支持的版本已进入 SSE 生命周期,必须发送 OpenAPI error event,不能返回普通 400。
+ log.warn("Events SSE API 版本不支持,tenantId={}, ownerUserId={}, apiVersionSupported={}",
+ tenantId, loginUserId, false);
+ completeWithSafeError(emitter, ERROR_API_VERSION_UNSUPPORTED, "Events API version unsupported", false);
+ return emitter;
+ }
+
+ EventsCursor cursor = EventsCursor.parse(lastEventId);
+ if (!cursor.isValid()) {
+ // cursor 原文可能来自浏览器或代理,日志和响应都不回显,避免把异常输入变成信息泄露面。
+ log.warn("Events SSE cursor 非法,tenantId={}, ownerUserId={}", tenantId, loginUserId);
+ completeWithSafeError(emitter, cursor.errorCode(), cursor.errorMessage(), false);
+ return emitter;
+ }
+
+ try {
+ Long startSequenceNo = initialSequenceNo(tenantId, loginUserId, cursor);
+ StreamBatch initialBatch = cursor.isInitialConnection()
+ ? StreamBatch.empty(startSequenceNo)
+ : loadVisibleBatch(tenantId, loginUserId, startSequenceNo);
+ sendBatchOrHeartbeat(emitter, initialBatch);
+ startPollingPersistedEvents(emitter, active, pollingFutureRef, tenantId, loginUserId,
+ initialBatch.maxSequenceNo());
+ } catch (RejectedExecutionException ex) {
+ active.set(false);
+ log.warn("Events SSE 轮询提交被拒绝,tenantId={}, ownerUserId={}", tenantId, loginUserId);
+ completeWithSafeError(emitter, ERROR_STREAM_UNAVAILABLE, "Events stream unavailable", true);
+ } catch (IOException ex) {
+ active.set(false);
+ emitter.completeWithError(ex);
+ } catch (Exception ex) {
+ active.set(false);
+ log.warn("Events SSE 初始化失败,tenantId={}, ownerUserId={}, errorType={}", tenantId, loginUserId,
+ ex.getClass().getSimpleName());
+ completeWithSafeError(emitter, ERROR_STREAM_UNAVAILABLE, "Events stream unavailable", true);
+ }
+ return emitter;
+ }
+
+ private void registerLifecycleCallbacks(SseEmitter emitter, AtomicBoolean active,
+ AtomicReference> pollingFutureRef) {
+ emitter.onTimeout(() -> {
+ active.set(false);
+ cancelPolling(pollingFutureRef);
+ emitter.complete();
+ });
+ emitter.onCompletion(() -> {
+ active.set(false);
+ cancelPolling(pollingFutureRef);
+ });
+ emitter.onError(error -> {
+ active.set(false);
+ cancelPolling(pollingFutureRef);
+ });
+ }
+
+ private Long initialSequenceNo(Long tenantId, Long ownerUserId, EventsCursor cursor) {
+ if (!cursor.isInitialConnection()) {
+ return cursor.sequenceNo();
+ }
+ Long maxSequenceNo = unifiedEventMapper.selectMaxVisibleSequenceForOwner(tenantId, ownerUserId,
+ LocalDateTime.now());
+ // 空 cursor 首连冻结当前可见最大 sequence,后续只追更大的事件;不能 replay 连接前历史。
+ return maxSequenceNo == null ? 0L : maxSequenceNo;
+ }
+
+ private StreamBatch loadVisibleBatch(Long tenantId, Long ownerUserId, Long afterSequenceNo) {
+ List rows = unifiedEventMapper.selectVisibleEventsForOwner(tenantId, ownerUserId,
+ afterSequenceNo == null ? 0L : afterSequenceNo, LocalDateTime.now());
+ return toStreamBatch(afterSequenceNo, rows);
+ }
+
+ private StreamBatch toStreamBatch(Long fallbackSequenceNo, List rows) {
+ Long maxSequenceNo = fallbackSequenceNo == null ? 0L : fallbackSequenceNo;
+ if (rows == null || rows.isEmpty()) {
+ return StreamBatch.empty(maxSequenceNo);
+ }
+ List events = rows.stream()
+ .map(this::toStreamEvent)
+ .filter(Objects::nonNull)
+ .toList();
+ for (UnifiedEventDO row : rows) {
+ if (row != null && row.getSequenceNo() != null) {
+ maxSequenceNo = Math.max(maxSequenceNo, row.getSequenceNo());
+ }
+ }
+ return new StreamBatch(events, maxSequenceNo);
+ }
+
+ private StreamEvent toStreamEvent(UnifiedEventDO row) {
+ if (row == null || row.getSequenceNo() == null || !DECLARED_EVENT_TYPES.contains(row.getEventType())) {
+ // 落库数据异常时不发送未知 event 名称,避免前端收到 OpenAPI 外事件。
+ log.warn("Events SSE 跳过非法事件行,tenantId={}, ownerUserId={}, sequenceNo={}, eventType={}",
+ row == null ? null : row.getTenantId(), row == null ? null : row.getOwnerUserId(),
+ row == null ? null : row.getSequenceNo(), row == null ? null : row.getEventType());
+ return null;
+ }
+ Map payload = payloadSanitizer.parseSanitizedJson(row.getPayloadSummary());
+ Map data = toOpenApiData(row.getEventType(), payload);
+ return data == null ? null : new StreamEvent(row.getSequenceNo(), row.getEventType(), data);
+ }
+
+ private Map toOpenApiData(String eventType, Map payload) {
+ return switch (eventType) {
+ case EVENT_CHUNK -> chunkData(payload);
+ case EVENT_QUALITY_CHECK -> qualityCheckData(payload);
+ case EVENT_DONE -> doneData(payload);
+ case EVENT_ERROR -> errorData(payload);
+ case EVENT_NOTIFICATION -> notificationData(payload);
+ default -> null;
+ };
+ }
+
+ private Map chunkData(Map payload) {
+ if (!(payload.get("content") instanceof String content) || !isPositiveLongNumber(payload.get("sequenceNo"))) {
+ return null;
+ }
+ Map data = new LinkedHashMap<>();
+ data.put("content", content);
+ data.put("sequenceNo", toLong(payload.get("sequenceNo")));
+ return data;
+ }
+
+ private Map qualityCheckData(Map payload) {
+ if (!(payload.get("dimension") instanceof String dimension)
+ || !(payload.get("score") instanceof Number score)
+ || !(payload.get("passed") instanceof Boolean passed)) {
+ return null;
+ }
+ double scoreValue = score.doubleValue();
+ if (scoreValue < 0D || scoreValue > 1D) {
+ return null;
+ }
+ Map data = new LinkedHashMap<>();
+ data.put("dimension", dimension);
+ data.put("score", scoreValue);
+ data.put("passed", passed);
+ return data;
+ }
+
+ private Map doneData(Map payload) {
+ if (!isLongNumber(payload.get("taskId")) || !isLongNumber(payload.get("suggestionId"))) {
+ return null;
+ }
+ Map data = new LinkedHashMap<>();
+ data.put("taskId", toLong(payload.get("taskId")));
+ data.put("suggestionId", toLong(payload.get("suggestionId")));
+ if (payload.get("summary") instanceof String summary) {
+ data.put("summary", summary);
+ }
+ return data;
+ }
+
+ private Map errorData(Map payload) {
+ if (!(payload.get("code") instanceof String code) || !(payload.get("message") instanceof String message)) {
+ return null;
+ }
+ Map data = new LinkedHashMap<>();
+ data.put("code", safeText(code, "EVENTS_ERROR"));
+ data.put("message", safeText(message, "Events stream error"));
+ if (payload.get("detail") instanceof String detail && !containsSecret(detail)) {
+ data.put("detail", detail);
+ }
+ if (payload.get("retryable") instanceof Boolean retryable) {
+ data.put("retryable", retryable);
+ }
+ return data;
+ }
+
+ private Map notificationData(Map payload) {
+ if (!(payload.get("type") instanceof String type) || !NOTIFICATION_TYPES.contains(type)
+ || !(payload.get("message") instanceof String message)) {
+ return null;
+ }
+ Map data = new LinkedHashMap<>();
+ data.put("type", type);
+ data.put("message", message);
+ if (payload.get("resourceRef") instanceof Map, ?> resourceRef) {
+ Map safeResourceRef = resourceRefData(resourceRef);
+ if (!safeResourceRef.isEmpty()) {
+ data.put("resourceRef", safeResourceRef);
+ }
+ }
+ if (payload.get("timestamp") instanceof String timestamp) {
+ data.put("timestamp", timestamp);
+ }
+ return data;
+ }
+
+ private Map resourceRefData(Map, ?> resourceRef) {
+ Map data = new LinkedHashMap<>();
+ Object resourceType = resourceRef.get("resourceType");
+ if (resourceType instanceof String text && !containsSecret(text)) {
+ data.put("resourceType", text);
+ }
+ Object resourceId = resourceRef.get("resourceId");
+ if (isLongNumber(resourceId)) {
+ data.put("resourceId", toLong(resourceId));
+ }
+ return data;
+ }
+
+ private void sendBatchOrHeartbeat(SseEmitter emitter, StreamBatch batch) throws IOException {
+ if (batch.events().isEmpty()) {
+ sendHeartbeat(emitter);
+ return;
+ }
+ for (StreamEvent event : batch.events()) {
+ // SSE id 承载全局恢复游标;data 只发送 OpenAPI 声明字段,避免扩展未声明字段。
+ emitter.send(SseEmitter.event()
+ .id(EventsCursor.toSseId(event.sequenceNo()))
+ .name(event.event())
+ .data(event.data(), MediaType.APPLICATION_JSON));
+ }
+ }
+
+ private void sendHeartbeat(SseEmitter emitter) throws IOException {
+ // 无事件时只能发送 comment heartbeat,不能伪造 notification 或 done。
+ emitter.send(SseEmitter.event().comment("heartbeat"));
+ }
+
+ private void startPollingPersistedEvents(SseEmitter emitter, AtomicBoolean active,
+ AtomicReference> pollingFutureRef, Long tenantId,
+ Long ownerUserId, Long lastSequenceNo) {
+ Future> pollingFuture = eventsStreamExecutor.submit(() ->
+ pollPersistedEvents(emitter, active, tenantId, ownerUserId, lastSequenceNo));
+ pollingFutureRef.set(pollingFuture);
+ if (!active.get()) {
+ cancelPolling(pollingFutureRef);
+ }
+ }
+
+ private void pollPersistedEvents(SseEmitter emitter, AtomicBoolean active, Long tenantId, Long ownerUserId,
+ Long lastSequenceNo) {
+ Long currentSequenceNo = lastSequenceNo == null ? 0L : lastSequenceNo;
+ long deadline = System.currentTimeMillis() + DEFAULT_TIMEOUT_MILLIS;
+ TenantContextHolder.setTenantId(tenantId);
+ try {
+ while (active.get() && System.currentTimeMillis() < deadline) {
+ StreamBatch delta = loadVisibleBatch(tenantId, ownerUserId, currentSequenceNo);
+ currentSequenceNo = delta.maxSequenceNo();
+ sendBatchOrHeartbeat(emitter, delta);
+ sleepHeartbeatInterval(active);
+ }
+ if (active.compareAndSet(true, false)) {
+ emitter.complete();
+ }
+ } catch (IOException ex) {
+ active.set(false);
+ emitter.completeWithError(ex);
+ } catch (Exception ex) {
+ // 后台 DB/mapper/租户拦截器异常不透出原始 message;日志仅保留定位所需的安全维度。
+ active.set(false);
+ log.warn("Events SSE 后台轮询失败,tenantId={}, ownerUserId={}, cursor={}, errorType={}", tenantId,
+ ownerUserId, currentSequenceNo, ex.getClass().getSimpleName());
+ completeWithSafeError(emitter, ERROR_STREAM_UNAVAILABLE, "Events stream unavailable", true);
+ } finally {
+ TenantContextHolder.clear();
+ }
+ }
+
+ private void sleepHeartbeatInterval(AtomicBoolean active) {
+ try {
+ Thread.sleep(HEARTBEAT_INTERVAL_MILLIS);
+ } catch (InterruptedException ex) {
+ active.set(false);
+ Thread.currentThread().interrupt();
+ }
+ }
+
+ private void completeWithSafeError(SseEmitter emitter, String code, String message, Boolean retryable) {
+ try {
+ Map data = new LinkedHashMap<>();
+ data.put("code", safeText(code, ERROR_STREAM_UNAVAILABLE));
+ data.put("message", safeText(message, "Events stream unavailable"));
+ if (retryable != null) {
+ data.put("retryable", retryable);
+ }
+ emitter.send(SseEmitter.event()
+ .name(EVENT_ERROR)
+ .data(data, MediaType.APPLICATION_JSON));
+ emitter.complete();
+ } catch (IOException ex) {
+ emitter.completeWithError(ex);
+ }
+ }
+
+ private void cancelPolling(AtomicReference> pollingFutureRef) {
+ Future> pollingFuture = pollingFutureRef.getAndSet(null);
+ if (pollingFuture != null && !pollingFuture.isDone()) {
+ pollingFuture.cancel(true);
+ }
+ }
+
+ private boolean isPositiveLongNumber(Object value) {
+ Long longValue = toLong(value);
+ return longValue != null && longValue >= 1L;
+ }
+
+ private boolean isLongNumber(Object value) {
+ return toLong(value) != null;
+ }
+
+ private Long toLong(Object value) {
+ if (!(value instanceof Number number)) {
+ return null;
+ }
+ double doubleValue = number.doubleValue();
+ long longValue = number.longValue();
+ return Double.compare(doubleValue, longValue) == 0 ? longValue : null;
+ }
+
+ private String safeText(String value, String fallback) {
+ if (!StringUtils.hasText(value) || containsSecret(value)) {
+ return fallback;
+ }
+ String normalized = value.trim();
+ return normalized.length() > 500 ? normalized.substring(0, 500) : normalized;
+ }
+
+ private boolean containsSecret(String value) {
+ return StringUtils.hasText(value) && SECRET_TEXT_PATTERN.matcher(value).find();
+ }
+
+ private record StreamEvent(Long sequenceNo, String event, Map data) {
+ }
+
+ private record StreamBatch(List events, Long maxSequenceNo) {
+
+ private static StreamBatch empty(Long maxSequenceNo) {
+ return new StreamBatch(List.of(), maxSequenceNo == null ? 0L : maxSequenceNo);
+ }
+ }
+
+}
diff --git a/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/controller/app/AppMuseEventsController.java b/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/controller/app/AppMuseEventsController.java
new file mode 100644
index 00000000..1401846b
--- /dev/null
+++ b/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/controller/app/AppMuseEventsController.java
@@ -0,0 +1,44 @@
+package cn.iocoder.muse.module.events.controller.app;
+
+import cn.iocoder.muse.module.events.application.stream.EventsStreamService;
+import io.swagger.v3.oas.annotations.Operation;
+import io.swagger.v3.oas.annotations.tags.Tag;
+import jakarta.annotation.Resource;
+import org.springframework.http.MediaType;
+import org.springframework.web.bind.annotation.GetMapping;
+import org.springframework.web.bind.annotation.RequestHeader;
+import org.springframework.web.bind.annotation.RequestParam;
+import org.springframework.web.bind.annotation.RequestMapping;
+import org.springframework.web.bind.annotation.RestController;
+import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
+
+import static cn.iocoder.muse.framework.security.core.util.SecurityFrameworkUtils.getLoginUserId;
+
+/**
+ * 用户端 Muse 统一事件流入口。
+ */
+@Tag(name = "用户 APP - Muse Events")
+@RestController
+@RequestMapping("/muse")
+public class AppMuseEventsController {
+
+ @Resource
+ private EventsStreamService eventsStreamService;
+
+ /**
+ * 建立统一事件流 SSE 连接。
+ *
+ * 登录用户必须在进入 SSE 生命周期前读取,避免后台线程再访问请求安全上下文。
+ *
+ * @param apiVersion `X-API-Version` 合同版本
+ * @param lastEventId SSE 恢复游标;空值表示首连,不 replay 历史
+ * @return SSE 事件流骨架
+ */
+ @GetMapping(value = "/events", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
+ @Operation(summary = "统一事件流 SSE")
+ public SseEmitter streamEvents(@RequestHeader(value = "X-API-Version", required = false) String apiVersion,
+ @RequestParam(value = "lastEventId", required = false) String lastEventId) {
+ return eventsStreamService.streamEvents(getLoginUserId(), apiVersion, lastEventId);
+ }
+
+}
diff --git a/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/dal/dataobject/UnifiedEventDO.java b/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/dal/dataobject/UnifiedEventDO.java
new file mode 100644
index 00000000..80b4e6b8
--- /dev/null
+++ b/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/dal/dataobject/UnifiedEventDO.java
@@ -0,0 +1,67 @@
+package cn.iocoder.muse.module.events.dal.dataobject;
+
+import cn.iocoder.muse.framework.tenant.core.db.TenantBaseDO;
+import cn.iocoder.muse.module.events.dal.type.JsonbStringTypeHandler;
+import com.baomidou.mybatisplus.annotation.IdType;
+import com.baomidou.mybatisplus.annotation.TableField;
+import com.baomidou.mybatisplus.annotation.TableId;
+import com.baomidou.mybatisplus.annotation.TableName;
+import lombok.*;
+
+import java.time.LocalDateTime;
+
+/**
+ * Events 统一事件投影 DO。
+ */
+@TableName(value = "muse_unified_event", autoResultMap = true)
+@Data
+@EqualsAndHashCode(callSuper = true)
+@ToString(callSuper = true)
+@Builder
+@NoArgsConstructor
+@AllArgsConstructor
+public class UnifiedEventDO extends TenantBaseDO {
+
+ @TableId(type = IdType.AUTO)
+ private Long id;
+
+ private String commandId;
+
+ private String eventId;
+
+ /**
+ * sequenceNo 由 PostgreSQL sequence/default 分配,应用层只在插入后回查读取。
+ */
+ private Long sequenceNo;
+
+ private String eventType;
+
+ private Long ownerUserId;
+
+ private String sourceOwner;
+
+ private String sourceType;
+
+ private String sourceId;
+
+ private String sourceRevision;
+
+ private String resourceType;
+
+ private String resourceId;
+
+ /**
+ * 已脱敏的 OpenAPI 事件 payload 摘要,不保存 provider raw body、token 或授权头。
+ */
+ @TableField(typeHandler = JsonbStringTypeHandler.class)
+ private String payloadSummary;
+
+ private String publishStatus;
+
+ private String publishErrorCode;
+
+ private LocalDateTime visibleFrom;
+
+ private LocalDateTime emittedAt;
+
+}
diff --git a/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/dal/mysql/UnifiedEventMapper.java b/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/dal/mysql/UnifiedEventMapper.java
new file mode 100644
index 00000000..37d553a7
--- /dev/null
+++ b/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/dal/mysql/UnifiedEventMapper.java
@@ -0,0 +1,80 @@
+package cn.iocoder.muse.module.events.dal.mysql;
+
+import cn.iocoder.muse.framework.mybatis.core.mapper.BaseMapperX;
+import cn.iocoder.muse.framework.mybatis.core.query.LambdaQueryWrapperX;
+import cn.iocoder.muse.module.events.application.publish.EventsPublishServiceImpl;
+import cn.iocoder.muse.module.events.dal.dataobject.UnifiedEventDO;
+import org.apache.ibatis.annotations.Insert;
+import org.apache.ibatis.annotations.Mapper;
+
+import java.time.LocalDateTime;
+import java.util.List;
+
+/**
+ * Events 统一事件投影 Mapper。
+ */
+@Mapper
+public interface UnifiedEventMapper extends BaseMapperX {
+
+ default UnifiedEventDO selectByTenantIdAndCommandId(Long tenantId, String commandId) {
+ return selectOne(new LambdaQueryWrapperX()
+ .eq(UnifiedEventDO::getTenantId, tenantId)
+ .eq(UnifiedEventDO::getCommandId, commandId));
+ }
+
+ default UnifiedEventDO selectByTenantIdAndSourceTuple(Long tenantId, String sourceOwner, String sourceType,
+ String sourceId, String sourceRevision, String eventType) {
+ return selectOne(new LambdaQueryWrapperX()
+ .eq(UnifiedEventDO::getTenantId, tenantId)
+ .eq(UnifiedEventDO::getSourceOwner, sourceOwner)
+ .eq(UnifiedEventDO::getSourceType, sourceType)
+ .eq(UnifiedEventDO::getSourceId, sourceId)
+ .eq(UnifiedEventDO::getSourceRevision, sourceRevision)
+ .eq(UnifiedEventDO::getEventType, eventType));
+ }
+
+ default List selectVisibleEventsForOwner(Long tenantId, Long ownerUserId, Long afterSequenceNo,
+ LocalDateTime now) {
+ return selectList(new LambdaQueryWrapperX()
+ .eq(UnifiedEventDO::getTenantId, tenantId)
+ .eq(UnifiedEventDO::getOwnerUserId, ownerUserId)
+ .eq(UnifiedEventDO::getPublishStatus, EventsPublishServiceImpl.PUBLISH_STATUS_ACCEPTED)
+ .eq(UnifiedEventDO::getDeleted, false)
+ .le(UnifiedEventDO::getVisibleFrom, now)
+ .gt(UnifiedEventDO::getSequenceNo, afterSequenceNo)
+ .orderByAsc(UnifiedEventDO::getSequenceNo));
+ }
+
+ default Long selectMaxVisibleSequenceForOwner(Long tenantId, Long ownerUserId, LocalDateTime now) {
+ UnifiedEventDO latest = selectOne(new LambdaQueryWrapperX()
+ .eq(UnifiedEventDO::getTenantId, tenantId)
+ .eq(UnifiedEventDO::getOwnerUserId, ownerUserId)
+ .eq(UnifiedEventDO::getPublishStatus, EventsPublishServiceImpl.PUBLISH_STATUS_ACCEPTED)
+ .eq(UnifiedEventDO::getDeleted, false)
+ .le(UnifiedEventDO::getVisibleFrom, now)
+ // 空 cursor 首连只冻结当前最大可见 sequence,不 replay 历史事件。
+ .orderByDesc(UnifiedEventDO::getSequenceNo)
+ .last("LIMIT 1"));
+ return latest == null ? 0L : latest.getSequenceNo();
+ }
+
+ /**
+ * 首次写入统一事件;并发重复 commandId 或 source tuple 时不覆盖既有事件,只由服务层回查幂等结果。
+ *
+ * 注意:这里不写 sequence_no 字段,必须由 V16 DDL 的 PostgreSQL sequence/default 分配。
+ */
+ @Insert("""
+ INSERT INTO muse_unified_event(command_id, event_id, event_type, owner_user_id,
+ source_owner, source_type, source_id, source_revision,
+ resource_type, resource_id, payload_summary,
+ publish_status, publish_error_code, visible_from, emitted_at, tenant_id)
+ VALUES (#{commandId}, #{eventId}, #{eventType}, #{ownerUserId},
+ #{sourceOwner}, #{sourceType}, #{sourceId}, #{sourceRevision},
+ #{resourceType}, #{resourceId},
+ #{payloadSummary,typeHandler=cn.iocoder.muse.module.events.dal.type.JsonbStringTypeHandler},
+ #{publishStatus}, #{publishErrorCode}, #{visibleFrom}, #{emittedAt}, #{tenantId})
+ ON CONFLICT DO NOTHING
+ """)
+ int insertIgnore(UnifiedEventDO event);
+
+}
diff --git a/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/dal/type/JsonbStringTypeHandler.java b/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/dal/type/JsonbStringTypeHandler.java
new file mode 100644
index 00000000..893ab291
--- /dev/null
+++ b/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/dal/type/JsonbStringTypeHandler.java
@@ -0,0 +1,48 @@
+package cn.iocoder.muse.module.events.dal.type;
+
+import org.apache.ibatis.type.BaseTypeHandler;
+import org.apache.ibatis.type.JdbcType;
+
+import java.sql.CallableStatement;
+import java.sql.PreparedStatement;
+import java.sql.ResultSet;
+import java.sql.SQLException;
+import java.sql.Types;
+
+/**
+ * Events 模块本地 PostgreSQL JSONB 字符串 TypeHandler。
+ *
+ * 统一事件只保存已经脱敏后的 JSON 摘要;这里用 JDBC OTHER 交给 PostgreSQL JSONB 解析,
+ * 避免把 JSON 字符串再次序列化成字符串字面量。
+ */
+public class JsonbStringTypeHandler extends BaseTypeHandler {
+
+ @Override
+ public void setNonNullParameter(PreparedStatement ps, int i, String parameter, JdbcType jdbcType)
+ throws SQLException {
+ ps.setObject(i, parameter, Types.OTHER);
+ }
+
+ @Override
+ public String getNullableResult(ResultSet rs, String columnName) throws SQLException {
+ return readJson(rs.getObject(columnName));
+ }
+
+ @Override
+ public String getNullableResult(ResultSet rs, int columnIndex) throws SQLException {
+ return readJson(rs.getObject(columnIndex));
+ }
+
+ @Override
+ public String getNullableResult(CallableStatement cs, int columnIndex) throws SQLException {
+ return readJson(cs.getObject(columnIndex));
+ }
+
+ private String readJson(Object value) {
+ if (value == null) {
+ return null;
+ }
+ return value.toString();
+ }
+
+}
diff --git a/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/domain/EventsCursor.java b/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/domain/EventsCursor.java
new file mode 100644
index 00000000..f6332285
--- /dev/null
+++ b/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/domain/EventsCursor.java
@@ -0,0 +1,63 @@
+package cn.iocoder.muse.module.events.domain;
+
+import org.springframework.util.StringUtils;
+
+/**
+ * Events SSE Last-Event-ID 游标。
+ *
+ * 游标解析必须 fail-closed:非法输入只返回安全错误码,不能把原始 cursor 或解析异常透出到浏览器。
+ */
+public record EventsCursor(Long sequenceNo, boolean valid, String errorCode, String errorMessage) {
+
+ private static final String PREFIX = "muse:";
+ public static final String ERROR_CURSOR_INVALID = "EVENTS_CURSOR_INVALID";
+
+ /**
+ * 解析 Events SSE 游标。
+ *
+ * @param rawCursor 浏览器传入的 lastEventId;空值表示首连,不 replay 历史
+ * @return 解析结果;非法 cursor 由调用方发送 OpenAPI error event
+ */
+ public static EventsCursor parse(String rawCursor) {
+ if (!StringUtils.hasText(rawCursor)) {
+ return new EventsCursor(null, true, null, null);
+ }
+ String cursor = rawCursor.trim();
+ if (!cursor.startsWith(PREFIX)) {
+ return invalid();
+ }
+ String sequenceText = cursor.substring(PREFIX.length());
+ if (!StringUtils.hasText(sequenceText)) {
+ return invalid();
+ }
+ try {
+ long sequenceNo = Long.parseLong(sequenceText);
+ if (sequenceNo < 0L) {
+ return invalid();
+ }
+ return new EventsCursor(sequenceNo, true, null, null);
+ } catch (NumberFormatException ex) {
+ return invalid();
+ }
+ }
+
+ /**
+ * SSE id 使用全局统一格式,避免客户端把裸数字和其他 stream 的游标混用。
+ */
+ public static String toSseId(Long sequenceNo) {
+ return PREFIX + sequenceNo;
+ }
+
+ public boolean isInitialConnection() {
+ return valid && sequenceNo == null;
+ }
+
+ public boolean isValid() {
+ return valid;
+ }
+
+ private static EventsCursor invalid() {
+ return new EventsCursor(null, false, ERROR_CURSOR_INVALID, "Invalid events cursor");
+ }
+
+}
diff --git a/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/domain/EventsPayloadSanitizer.java b/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/domain/EventsPayloadSanitizer.java
new file mode 100644
index 00000000..71c01b4c
--- /dev/null
+++ b/muse-cloud/muse-module-events/muse-module-events-server/src/main/java/cn/iocoder/muse/module/events/domain/EventsPayloadSanitizer.java
@@ -0,0 +1,96 @@
+package cn.iocoder.muse.module.events.domain;
+
+import cn.iocoder.muse.framework.common.util.json.JsonUtils;
+import com.fasterxml.jackson.core.type.TypeReference;
+import lombok.Value;
+import org.springframework.stereotype.Component;
+import org.springframework.util.StringUtils;
+
+import java.util.*;
+import java.util.regex.Pattern;
+
+/**
+ * Events payload 脱敏器。
+ *
+ * source owner 可能来自外部模型、知识库或支付/授权链路;这里递归检查 key 与 value,
+ * 发现授权头、bearer token、secret 或 provider raw body 后,事件会被拒绝且只保存脱敏摘要。
+ */
+@Component
+public class EventsPayloadSanitizer {
+
+ private static final String REDACTED = "[REDACTED]";
+
+ private static final Pattern SECRET_KEY_PATTERN = Pattern.compile(
+ "(?i)(authorization|bearer|token|access[-_]?token|refresh[-_]?token|api[-_]?key|secret|password|credential|provider.*raw.*body|raw.*body)");
+
+ private static final Pattern SECRET_VALUE_PATTERN = Pattern.compile(
+ "(?i)(bearer\\s+[a-z0-9._\\-]+|authorization\\s*[:=]|api[-_]?key\\s*[:=]|secret\\s*[:=]|token\\s*[:=])");
+
+ public SanitizedPayload sanitize(Map payloadSummary) {
+ Map source = payloadSummary == null ? Map.of() : payloadSummary;
+ Detection detection = new Detection();
+ Object sanitized = sanitizeValue(source, detection);
+ String json = JsonUtils.toJsonString(sanitized);
+ return new SanitizedPayload(json, detection.containsSecret);
+ }
+
+ @SuppressWarnings("unchecked")
+ private Object sanitizeValue(Object value, Detection detection) {
+ if (value instanceof Map, ?> map) {
+ Map sanitizedMap = new LinkedHashMap<>();
+ map.forEach((key, childValue) -> {
+ String normalizedKey = String.valueOf(key);
+ if (isSensitiveKey(normalizedKey)) {
+ detection.containsSecret = true;
+ sanitizedMap.put(normalizedKey, REDACTED);
+ return;
+ }
+ sanitizedMap.put(normalizedKey, sanitizeValue(childValue, detection));
+ });
+ return sanitizedMap;
+ }
+ if (value instanceof Collection> collection) {
+ List