feat(p1r): 落地 Events SSE 真实 API

This commit is contained in:
zizi 2026-06-05 22:20:48 +08:00
parent 5993b6f51c
commit 165579521a
32 changed files with 3519 additions and 51 deletions

View File

@ -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 统一事件流入口。
*
* <p>P1 阶段先提供可建立连接的 SSE 端点,后续由各业务模块接入真实事件发布器。</p>
*/
@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;
}
}

View File

@ -0,0 +1,36 @@
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<parent>
<groupId>cn.iocoder.cloud</groupId>
<artifactId>muse-module-events</artifactId>
<version>${revision}</version>
</parent>
<modelVersion>4.0.0</modelVersion>
<artifactId>muse-module-events-api</artifactId>
<packaging>jar</packaging>
<name>${project.artifactId}</name>
<description>events 模块 API 和发布契约。</description>
<dependencies>
<dependency>
<groupId>cn.iocoder.cloud</groupId>
<artifactId>muse-common</artifactId>
</dependency>
<dependency>
<groupId>org.springdoc</groupId>
<artifactId>springdoc-openapi-starter-webmvc-ui</artifactId>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-validation</artifactId>
<optional>true</optional>
</dependency>
<dependency>
<groupId>org.springframework.cloud</groupId>
<artifactId>spring-cloud-starter-openfeign</artifactId>
<optional>true</optional>
</dependency>
</dependencies>
</project>

View File

@ -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<EventsPublishRespDTO> publish(@Valid @RequestBody EventsPublishReqDTO reqDTO);
}

View File

@ -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<String, Object> payloadSummary;
@NotNull(message = "emittedAt 不能为空")
private LocalDateTime emittedAt;
}

View File

@ -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;
}

View File

@ -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";
}

View File

@ -0,0 +1,16 @@
package cn.iocoder.muse.module.events.enums;
import cn.iocoder.muse.framework.common.exception.ErrorCode;
/**
* Events 错误码。
*
* <p>Events 使用 1-045-000-000 段,覆盖统一事件发布、幂等投影和 SSE 可见性边界。</p>
*/
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, "事件发布必要字段不能为空:{}");
}

View File

@ -0,0 +1,47 @@
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<parent>
<groupId>cn.iocoder.cloud</groupId>
<artifactId>muse-module-events</artifactId>
<version>${revision}</version>
</parent>
<modelVersion>4.0.0</modelVersion>
<artifactId>muse-module-events-server</artifactId>
<packaging>jar</packaging>
<name>${project.artifactId}</name>
<description>events 模块服务实现。</description>
<dependencies>
<dependency>
<groupId>cn.iocoder.cloud</groupId>
<artifactId>muse-module-events-api</artifactId>
<version>${revision}</version>
</dependency>
<dependency>
<groupId>cn.iocoder.cloud</groupId>
<artifactId>muse-spring-boot-starter-web</artifactId>
</dependency>
<dependency>
<groupId>cn.iocoder.cloud</groupId>
<artifactId>muse-spring-boot-starter-security</artifactId>
</dependency>
<dependency>
<groupId>cn.iocoder.cloud</groupId>
<artifactId>muse-spring-boot-starter-biz-tenant</artifactId>
</dependency>
<dependency>
<groupId>cn.iocoder.cloud</groupId>
<artifactId>muse-spring-boot-starter-mybatis</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-validation</artifactId>
</dependency>
<dependency>
<groupId>cn.iocoder.cloud</groupId>
<artifactId>muse-spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>
</dependencies>
</project>

View File

@ -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<EventsPublishRespDTO> publish(EventsPublishReqDTO reqDTO) {
return success(eventsPublishService.publish(reqDTO));
}
}

View File

@ -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<UnifiedEventDO> listVisibleEvents(Long tenantId, Long ownerUserId, Long afterSequenceNo, LocalDateTime now);
}

View File

@ -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<String> 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<String> 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<UnifiedEventDO> 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<String, Object> payload) {
Map<String, Object> 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<String, Object> 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);
}
}
}

View File

@ -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);
}

View File

@ -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<String> DECLARED_EVENT_TYPES = Set.of(EVENT_CHUNK, EVENT_QUALITY_CHECK, EVENT_DONE,
EVENT_ERROR, EVENT_NOTIFICATION);
private static final Set<String> 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<Future<?>> 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<Future<?>> 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<UnifiedEventDO> rows = unifiedEventMapper.selectVisibleEventsForOwner(tenantId, ownerUserId,
afterSequenceNo == null ? 0L : afterSequenceNo, LocalDateTime.now());
return toStreamBatch(afterSequenceNo, rows);
}
private StreamBatch toStreamBatch(Long fallbackSequenceNo, List<UnifiedEventDO> rows) {
Long maxSequenceNo = fallbackSequenceNo == null ? 0L : fallbackSequenceNo;
if (rows == null || rows.isEmpty()) {
return StreamBatch.empty(maxSequenceNo);
}
List<StreamEvent> 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<String, Object> payload = payloadSanitizer.parseSanitizedJson(row.getPayloadSummary());
Map<String, Object> data = toOpenApiData(row.getEventType(), payload);
return data == null ? null : new StreamEvent(row.getSequenceNo(), row.getEventType(), data);
}
private Map<String, Object> toOpenApiData(String eventType, Map<String, Object> 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<String, Object> chunkData(Map<String, Object> payload) {
if (!(payload.get("content") instanceof String content) || !isPositiveLongNumber(payload.get("sequenceNo"))) {
return null;
}
Map<String, Object> data = new LinkedHashMap<>();
data.put("content", content);
data.put("sequenceNo", toLong(payload.get("sequenceNo")));
return data;
}
private Map<String, Object> qualityCheckData(Map<String, Object> 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<String, Object> data = new LinkedHashMap<>();
data.put("dimension", dimension);
data.put("score", scoreValue);
data.put("passed", passed);
return data;
}
private Map<String, Object> doneData(Map<String, Object> payload) {
if (!isLongNumber(payload.get("taskId")) || !isLongNumber(payload.get("suggestionId"))) {
return null;
}
Map<String, Object> 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<String, Object> errorData(Map<String, Object> payload) {
if (!(payload.get("code") instanceof String code) || !(payload.get("message") instanceof String message)) {
return null;
}
Map<String, Object> 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<String, Object> notificationData(Map<String, Object> payload) {
if (!(payload.get("type") instanceof String type) || !NOTIFICATION_TYPES.contains(type)
|| !(payload.get("message") instanceof String message)) {
return null;
}
Map<String, Object> data = new LinkedHashMap<>();
data.put("type", type);
data.put("message", message);
if (payload.get("resourceRef") instanceof Map<?, ?> resourceRef) {
Map<String, Object> 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<String, Object> resourceRefData(Map<?, ?> resourceRef) {
Map<String, Object> 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<Future<?>> 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<String, Object> 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<Future<?>> 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<String, Object> data) {
}
private record StreamBatch(List<StreamEvent> events, Long maxSequenceNo) {
private static StreamBatch empty(Long maxSequenceNo) {
return new StreamBatch(List.of(), maxSequenceNo == null ? 0L : maxSequenceNo);
}
}
}

View File

@ -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 连接。
*
* <p>登录用户必须在进入 SSE 生命周期前读取,避免后台线程再访问请求安全上下文。</p>
*
* @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);
}
}

View File

@ -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;
}

View File

@ -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<UnifiedEventDO> {
default UnifiedEventDO selectByTenantIdAndCommandId(Long tenantId, String commandId) {
return selectOne(new LambdaQueryWrapperX<UnifiedEventDO>()
.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<UnifiedEventDO>()
.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<UnifiedEventDO> selectVisibleEventsForOwner(Long tenantId, Long ownerUserId, Long afterSequenceNo,
LocalDateTime now) {
return selectList(new LambdaQueryWrapperX<UnifiedEventDO>()
.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<UnifiedEventDO>()
.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 时不覆盖既有事件,只由服务层回查幂等结果。
*
* <p>注意:这里不写 sequence_no 字段,必须由 V16 DDL 的 PostgreSQL sequence/default 分配。</p>
*/
@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);
}

View File

@ -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。
*
* <p>统一事件只保存已经脱敏后的 JSON 摘要;这里用 JDBC OTHER 交给 PostgreSQL JSONB 解析,
* 避免把 JSON 字符串再次序列化成字符串字面量。</p>
*/
public class JsonbStringTypeHandler extends BaseTypeHandler<String> {
@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();
}
}

View File

@ -0,0 +1,63 @@
package cn.iocoder.muse.module.events.domain;
import org.springframework.util.StringUtils;
/**
* Events SSE Last-Event-ID 游标。
*
* <p>游标解析必须 fail-closed:非法输入只返回安全错误码,不能把原始 cursor 或解析异常透出到浏览器。</p>
*/
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");
}
}

View File

@ -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 脱敏器。
*
* <p>source owner 可能来自外部模型、知识库或支付/授权链路;这里递归检查 key 与 value,
* 发现授权头、bearer token、secret 或 provider raw body 后,事件会被拒绝且只保存脱敏摘要。</p>
*/
@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<String, Object> payloadSummary) {
Map<String, Object> 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<String, Object> 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<Object> sanitizedList = new ArrayList<>(collection.size());
collection.forEach(item -> sanitizedList.add(sanitizeValue(item, detection)));
return sanitizedList;
}
if (value instanceof Object[] array) {
List<Object> sanitizedList = new ArrayList<>(array.length);
Arrays.stream(array).forEach(item -> sanitizedList.add(sanitizeValue(item, detection)));
return sanitizedList;
}
if (value instanceof String text && isSensitiveValue(text)) {
detection.containsSecret = true;
return REDACTED;
}
return value;
}
private boolean isSensitiveKey(String key) {
return StringUtils.hasText(key) && SECRET_KEY_PATTERN.matcher(key).find();
}
private boolean isSensitiveValue(String value) {
return StringUtils.hasText(value) && SECRET_VALUE_PATTERN.matcher(value).find();
}
/**
* 解析已落库 JSON,供后续扩展测试或 stream 层复用。
*/
public Map<String, Object> parseSanitizedJson(String json) {
Map<String, Object> parsed = JsonUtils.parseObjectQuietly(json, new TypeReference<>() {
});
return parsed == null ? Map.of() : parsed;
}
@Value
public static class SanitizedPayload {
String payloadJson;
boolean containsSecret;
}
private static class Detection {
private boolean containsSecret;
}
}

View File

@ -0,0 +1,34 @@
package cn.iocoder.muse.module.events.framework.config;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.core.task.AsyncTaskExecutor;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
import java.util.concurrent.ThreadPoolExecutor;
/**
* Events 统一 SSE 线程池配置。
*/
@Configuration(proxyBeanMethods = false)
public class EventsStreamConfiguration {
public static final String EVENTS_STREAM_EXECUTOR = "eventsStreamExecutor";
@Bean(EVENTS_STREAM_EXECUTOR)
public AsyncTaskExecutor eventsStreamExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(4);
executor.setMaxPoolSize(16);
executor.setKeepAliveSeconds(60);
executor.setQueueCapacity(100);
executor.setThreadNamePrefix("muse-events-sse-");
executor.setWaitForTasksToCompleteOnShutdown(false);
executor.setAwaitTerminationSeconds(10);
// 线程池饱和时必须 fail-closed 返回 SSE error,不能退回调用线程造成请求线程阻塞。
executor.setRejectedExecutionHandler(new ThreadPoolExecutor.AbortPolicy());
executor.initialize();
return executor;
}
}

View File

@ -0,0 +1,450 @@
package cn.iocoder.muse.module.events.application.publish;
import cn.iocoder.muse.framework.test.core.ut.BaseMockitoUnitTest;
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 org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.mockito.Mock;
import org.mockito.stubbing.Answer;
import java.time.LocalDateTime;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicLong;
import static org.junit.jupiter.api.Assertions.*;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.*;
class EventsPublishServiceTest extends BaseMockitoUnitTest {
@Mock
private UnifiedEventMapper unifiedEventMapper;
private EventsPublishServiceImpl service;
private final List<UnifiedEventDO> rows = new CopyOnWriteArrayList<>();
private final AtomicLong idSequence = new AtomicLong(1L);
private final AtomicLong dbSequence = new AtomicLong(1L);
@BeforeEach
void setUp() {
service = new EventsPublishServiceImpl(unifiedEventMapper);
// 这里模拟 PostgreSQL insert 后由 sequence/default 回填的结果,服务层只能回查读取,不能自己计算 sequenceNo。
lenient().when(unifiedEventMapper.insertIgnore(any(UnifiedEventDO.class))).thenAnswer(insertWithDbAssignedSequence());
lenient().when(unifiedEventMapper.selectByTenantIdAndCommandId(anyLong(), anyString()))
.thenAnswer(invocation -> findByCommand(invocation.getArgument(0), invocation.getArgument(1)));
lenient().when(unifiedEventMapper.selectByTenantIdAndSourceTuple(anyLong(), anyString(), anyString(), anyString(),
anyString(), anyString()))
.thenAnswer(invocation -> findBySourceTuple(invocation.getArgument(0), invocation.getArgument(1),
invocation.getArgument(2), invocation.getArgument(3), invocation.getArgument(4),
invocation.getArgument(5)));
lenient().when(unifiedEventMapper.selectVisibleEventsForOwner(anyLong(), anyLong(), anyLong(), any(LocalDateTime.class)))
.thenAnswer(invocation -> rows.stream()
.filter(row -> row.getTenantId().equals(invocation.getArgument(0))
&& row.getOwnerUserId().equals(invocation.getArgument(1))
&& row.getSequenceNo() > invocation.<Long>getArgument(2)
&& Boolean.FALSE.equals(row.getDeleted())
&& EventsPublishServiceImpl.PUBLISH_STATUS_ACCEPTED.equals(row.getPublishStatus())
&& !row.getVisibleFrom().isAfter(invocation.getArgument(3)))
.sorted((left, right) -> Long.compare(left.getSequenceNo(), right.getSequenceNo()))
.toList());
}
@Test
void should_publishAcceptedEvent_when_payloadMatchesOpenApiSchema() {
EventsPublishRespDTO resp = service.publish(chunkReq("cmd-1", "rev-1", Map.of(
"content", "hello",
"sequenceNo", 1
)));
assertEquals(EventsPublishServiceImpl.PUBLISH_STATUS_ACCEPTED, resp.getPublishStatus());
assertFalse(resp.getDuplicate());
assertNotNull(resp.getEventId());
assertNotNull(resp.getSequenceNo());
UnifiedEventDO row = rows.get(0);
assertEquals("cmd-1", row.getCommandId());
assertEquals("chunk", row.getEventType());
assertEquals(EventsPublishServiceImpl.PUBLISH_STATUS_ACCEPTED, row.getPublishStatus());
assertNull(row.getPublishErrorCode());
}
@Test
void should_acceptOpenApiStringFields_when_stringIsEmpty() {
EventsPublishRespDTO chunkResp = service.publish(chunkReq("cmd-empty-chunk", "rev-empty-chunk", Map.of(
"content", "",
"sequenceNo", 1
)));
EventsPublishRespDTO errorResp = service.publish(eventReq("cmd-empty-error", "rev-empty-error", "error",
Map.of(
"code", "",
"message", ""
)));
EventsPublishRespDTO notificationResp = service.publish(eventReq("cmd-empty-notification",
"rev-empty-notification", "notification", Map.of(
"type", "quota_alert",
"message", ""
)));
assertEquals(EventsPublishServiceImpl.PUBLISH_STATUS_ACCEPTED, chunkResp.getPublishStatus());
assertEquals(EventsPublishServiceImpl.PUBLISH_STATUS_ACCEPTED, errorResp.getPublishStatus());
assertEquals(EventsPublishServiceImpl.PUBLISH_STATUS_ACCEPTED, notificationResp.getPublishStatus());
}
@Test
void should_acceptDoneEvent_when_taskIdAndSuggestionIdAreLongOrInteger() {
EventsPublishRespDTO longResp = service.publish(eventReq("cmd-done-long", "rev-done-long", "done", Map.of(
"taskId", 1001L,
"suggestionId", 2001L
)));
EventsPublishRespDTO integerResp = service.publish(eventReq("cmd-done-int", "rev-done-int", "done", Map.of(
"taskId", 1002,
"suggestionId", 2002
)));
assertEquals(EventsPublishServiceImpl.PUBLISH_STATUS_ACCEPTED, longResp.getPublishStatus());
assertEquals(EventsPublishServiceImpl.PUBLISH_STATUS_ACCEPTED, integerResp.getPublishStatus());
assertEquals("done", rows.get(0).getEventType());
assertEquals("done", rows.get(1).getEventType());
}
@Test
void should_rejectDoneEvent_when_taskIdOrSuggestionIdIsString() {
EventsPublishRespDTO taskIdResp = service.publish(eventReq("cmd-done-string-task", "rev-done-string-task",
"done", Map.of(
"taskId", "1001",
"suggestionId", 2001L
)));
EventsPublishRespDTO suggestionIdResp = service.publish(eventReq("cmd-done-string-suggestion",
"rev-done-string-suggestion", "done", Map.of(
"taskId", 1002L,
"suggestionId", "2002"
)));
assertSchemaMismatchRejectedAndInvisible(taskIdResp);
assertSchemaMismatchRejectedAndInvisible(suggestionIdResp);
}
@Test
void should_validateQualityCheckScoreRange() {
EventsPublishRespDTO validResp = service.publish(eventReq("cmd-quality-valid", "rev-quality-valid",
"quality_check", Map.of(
"dimension", "readability",
"score", 0.8D,
"passed", true
)));
EventsPublishRespDTO lowResp = service.publish(eventReq("cmd-quality-low", "rev-quality-low",
"quality_check", Map.of(
"dimension", "readability",
"score", -0.1D,
"passed", false
)));
EventsPublishRespDTO highResp = service.publish(eventReq("cmd-quality-high", "rev-quality-high",
"quality_check", Map.of(
"dimension", "readability",
"score", 1.1D,
"passed", true
)));
assertEquals(EventsPublishServiceImpl.PUBLISH_STATUS_ACCEPTED, validResp.getPublishStatus());
assertSchemaMismatchRejectedAndInvisible(lowResp);
assertSchemaMismatchRejectedAndInvisible(highResp);
}
@Test
void should_rejectChunkEvent_when_sequenceNoIsZeroOrNegative() {
EventsPublishRespDTO zeroResp = service.publish(chunkReq("cmd-chunk-zero", "rev-chunk-zero", Map.of(
"content", "hello",
"sequenceNo", 0
)));
EventsPublishRespDTO negativeResp = service.publish(chunkReq("cmd-chunk-negative", "rev-chunk-negative", Map.of(
"content", "hello",
"sequenceNo", -1
)));
assertSchemaMismatchRejectedAndInvisible(zeroResp);
assertSchemaMismatchRejectedAndInvisible(negativeResp);
}
@Test
void should_validateNotificationTypeEnum() {
EventsPublishRespDTO validResp = service.publish(eventReq("cmd-notification-valid", "rev-notification-valid",
"notification", Map.of(
"type", "source_status_change",
"message", "source changed"
)));
EventsPublishRespDTO unknownResp = service.publish(eventReq("cmd-notification-unknown",
"rev-notification-unknown", "notification", Map.of(
"type", "unknown_type",
"message", "unknown"
)));
assertEquals(EventsPublishServiceImpl.PUBLISH_STATUS_ACCEPTED, validResp.getPublishStatus());
assertSchemaMismatchRejectedAndInvisible(unknownResp);
}
@Test
void should_returnSameEvent_when_duplicateCommandPublished() {
EventsPublishRespDTO first = service.publish(chunkReq("cmd-dup", "rev-1", Map.of(
"content", "first",
"sequenceNo", 1
)));
EventsPublishRespDTO second = service.publish(chunkReq("cmd-dup", "rev-2", Map.of(
"content", "ignored by idempotency",
"sequenceNo", 2
)));
assertEquals(first.getEventId(), second.getEventId());
assertEquals(first.getSequenceNo(), second.getSequenceNo());
assertTrue(second.getDuplicate());
assertEquals(1, rows.size());
}
@Test
void should_rejectPayload_when_containsSecret() {
EventsPublishRespDTO resp = service.publish(chunkReq("cmd-secret", "rev-secret", Map.of(
"content", "hello",
"sequenceNo", 1,
"authorization", "Bearer abc"
)));
assertEquals(EventsPublishServiceImpl.PUBLISH_STATUS_REJECTED, resp.getPublishStatus());
assertEquals(EventsPublishServiceImpl.ERROR_PAYLOAD_CONTAINS_SECRET, resp.getPublishErrorCode());
UnifiedEventDO row = rows.get(0);
assertEquals(EventsPublishServiceImpl.PUBLISH_STATUS_REJECTED, row.getPublishStatus());
assertEquals(EventsPublishServiceImpl.ERROR_PAYLOAD_CONTAINS_SECRET, row.getPublishErrorCode());
}
@Test
void should_rejectAndHidePayload_when_containsStandaloneBearerKey() {
EventsPublishRespDTO resp = service.publish(chunkReq("cmd-bearer-key", "rev-bearer-key", Map.of(
"content", "hello",
"sequenceNo", 1,
"bearer", "abc"
)));
assertEquals(EventsPublishServiceImpl.PUBLISH_STATUS_REJECTED, resp.getPublishStatus());
assertEquals(EventsPublishServiceImpl.ERROR_PAYLOAD_CONTAINS_SECRET, resp.getPublishErrorCode());
List<UnifiedEventDO> visibleRows = service.listVisibleEvents(1L, 10L, 0L, LocalDateTime.now().plusSeconds(1));
assertTrue(visibleRows.stream()
.noneMatch(row -> row.getEventId().equals(resp.getEventId())));
assertFalse(rows.get(0).getPayloadSummary().contains("abc"));
}
@Test
void should_rejectPayload_when_eventTypeNotDeclared() {
EventsPublishReqDTO req = chunkReq("cmd-unknown", "rev-unknown", Map.of(
"content", "hello",
"sequenceNo", 1
));
req.setEventType("custom_event");
EventsPublishRespDTO resp = service.publish(req);
assertEquals(EventsPublishServiceImpl.PUBLISH_STATUS_REJECTED, resp.getPublishStatus());
assertEquals(EventsPublishServiceImpl.ERROR_EVENT_TYPE_NOT_DECLARED, resp.getPublishErrorCode());
assertEquals("error", rows.get(0).getEventType());
}
@Test
void should_notExposeRejectedEventToStreamQuery() {
service.publish(chunkReq("cmd-secret", "rev-secret", Map.of(
"content", "hello",
"sequenceNo", 1,
"token", "sensitive"
)));
List<UnifiedEventDO> visibleRows = service.listVisibleEvents(1L, 10L, 0L, LocalDateTime.now().plusSeconds(1));
assertTrue(visibleRows.isEmpty());
}
@Test
void should_returnSameEvent_when_sourceOwnerReplaysOutboxFact() {
EventsPublishRespDTO first = service.publish(chunkReq("cmd-first", "rev-replay", Map.of(
"content", "hello",
"sequenceNo", 1
)));
EventsPublishRespDTO replay = service.publish(chunkReq("cmd-replay", "rev-replay", Map.of(
"content", "same outbox fact",
"sequenceNo", 1
)));
assertEquals(first.getEventId(), replay.getEventId());
assertEquals(first.getSequenceNo(), replay.getSequenceNo());
assertTrue(replay.getDuplicate());
assertEquals(1, rows.size());
}
@Test
void should_notExposeBlockedEvent_when_publishValidationFailsAfterRetry() {
UnifiedEventDO blocked = baseRow("cmd-blocked", "rev-blocked", "chunk");
blocked.setPublishStatus(EventsPublishServiceImpl.PUBLISH_STATUS_BLOCKED);
blocked.setPublishErrorCode("retry_dead_letter");
rows.add(blocked);
List<UnifiedEventDO> visibleRows = service.listVisibleEvents(1L, 10L, 0L, LocalDateTime.now().plusSeconds(1));
assertTrue(visibleRows.isEmpty());
}
@Test
void should_returnSameEvent_when_sourceRevisionIsNullAndReplayed() {
EventsPublishReqDTO firstReq = chunkReq("cmd-null-rev-1", null, Map.of(
"content", "hello",
"sequenceNo", 1
));
EventsPublishRespDTO first = service.publish(firstReq);
EventsPublishReqDTO replayReq = chunkReq("cmd-null-rev-2", null, Map.of(
"content", "hello again",
"sequenceNo", 1
));
EventsPublishRespDTO replay = service.publish(replayReq);
assertEquals(first.getEventId(), replay.getEventId());
assertEquals(first.getSequenceNo(), replay.getSequenceNo());
assertEquals(EventsPublishServiceImpl.SOURCE_REVISION_NONE, rows.get(0).getSourceRevision());
}
@Test
void should_allocateMonotonicSequenceNo_whenPublishConcurrentEvents() throws Exception {
ExecutorService executor = Executors.newFixedThreadPool(4);
List<Callable<EventsPublishRespDTO>> tasks = new ArrayList<>();
for (int i = 0; i < 8; i++) {
int index = i;
tasks.add(() -> service.publish(chunkReq("cmd-concurrent-" + index, "rev-" + index, Map.of(
"content", "chunk-" + index,
"sequenceNo", index + 1
))));
}
List<EventsPublishRespDTO> responses;
try {
responses = executor.invokeAll(tasks).stream()
.map(future -> {
try {
return future.get(3, TimeUnit.SECONDS);
} catch (Exception ex) {
throw new CompletionException(ex);
}
})
.toList();
} finally {
executor.shutdownNow();
}
List<Long> sequenceNos = responses.stream()
.map(EventsPublishRespDTO::getSequenceNo)
.sorted()
.toList();
assertEquals(List.of(1L, 2L, 3L, 4L, 5L, 6L, 7L, 8L), sequenceNos);
}
@Test
void should_failClosed_whenInsertIgnoredAndStoredEventCannotBeFound() {
when(unifiedEventMapper.insertIgnore(any(UnifiedEventDO.class))).thenReturn(0);
when(unifiedEventMapper.selectByTenantIdAndCommandId(anyLong(), anyString())).thenReturn(null);
when(unifiedEventMapper.selectByTenantIdAndSourceTuple(anyLong(), anyString(), anyString(), anyString(),
anyString(), anyString())).thenReturn(null);
IllegalStateException exception = assertThrows(IllegalStateException.class, () -> service.publish(
chunkReq("cmd-missing-after-ignore", "rev-missing-after-ignore", Map.of(
"content", "lost",
"sequenceNo", 1
))));
assertTrue(exception.getMessage().contains("Events 发布幂等回查失败"));
assertTrue(rows.isEmpty(), "insertIgnore 冲突且回查为空时,不能返回未落库的 eventId/sequenceNo");
}
private Answer<Integer> insertWithDbAssignedSequence() {
return invocation -> {
UnifiedEventDO row = invocation.getArgument(0);
row.setId(idSequence.getAndIncrement());
row.setSequenceNo(dbSequence.getAndIncrement());
rows.add(row);
return 1;
};
}
private UnifiedEventDO findByCommand(Long tenantId, String commandId) {
return rows.stream()
.filter(row -> row.getTenantId().equals(tenantId) && row.getCommandId().equals(commandId))
.findFirst()
.orElse(null);
}
private UnifiedEventDO findBySourceTuple(Long tenantId, String sourceOwner, String sourceType, String sourceId,
String sourceRevision, String eventType) {
return rows.stream()
.filter(row -> row.getTenantId().equals(tenantId)
&& row.getSourceOwner().equals(sourceOwner)
&& row.getSourceType().equals(sourceType)
&& row.getSourceId().equals(sourceId)
&& row.getSourceRevision().equals(sourceRevision)
&& row.getEventType().equals(eventType))
.findFirst()
.orElse(null);
}
private EventsPublishReqDTO chunkReq(String commandId, String sourceRevision, Map<String, Object> payload) {
EventsPublishReqDTO req = new EventsPublishReqDTO();
req.setCommandId(commandId);
req.setTenantId(1L);
req.setOwnerUserId(10L);
req.setSourceOwner("ai");
req.setSourceType("task");
req.setSourceId("task-1");
req.setSourceRevision(sourceRevision);
req.setEventType("chunk");
req.setResourceType("suggestion");
req.setResourceId("suggestion-1");
req.setPayloadSummary(payload);
req.setEmittedAt(LocalDateTime.now());
return req;
}
private EventsPublishReqDTO eventReq(String commandId, String sourceRevision, String eventType,
Map<String, Object> payload) {
EventsPublishReqDTO req = chunkReq(commandId, sourceRevision, payload);
req.setEventType(eventType);
return req;
}
private void assertSchemaMismatchRejectedAndInvisible(EventsPublishRespDTO resp) {
assertEquals(EventsPublishServiceImpl.PUBLISH_STATUS_REJECTED, resp.getPublishStatus());
assertEquals(EventsPublishServiceImpl.ERROR_PAYLOAD_SCHEMA_MISMATCH, resp.getPublishErrorCode());
List<UnifiedEventDO> visibleRows = service.listVisibleEvents(1L, 10L, 0L, LocalDateTime.now().plusSeconds(1));
assertTrue(visibleRows.stream()
.noneMatch(row -> row.getEventId().equals(resp.getEventId())));
}
private UnifiedEventDO baseRow(String commandId, String sourceRevision, String eventType) {
UnifiedEventDO row = new UnifiedEventDO();
row.setTenantId(1L);
row.setCommandId(commandId);
row.setEventId("evt-" + commandId);
row.setSequenceNo(dbSequence.getAndIncrement());
row.setEventType(eventType);
row.setOwnerUserId(10L);
row.setSourceOwner("ai");
row.setSourceType("task");
row.setSourceId("task-1");
row.setSourceRevision(sourceRevision);
row.setResourceType("suggestion");
row.setResourceId("suggestion-1");
row.setPayloadSummary("{}");
row.setVisibleFrom(LocalDateTime.now());
row.setEmittedAt(LocalDateTime.now());
row.setDeleted(false);
return row;
}
}

View File

@ -0,0 +1,324 @@
package cn.iocoder.muse.module.events.application.stream;
import cn.iocoder.muse.framework.common.util.json.JsonUtils;
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.EventsPayloadSanitizer;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.mockito.Mock;
import org.mockito.MockitoAnnotations;
import org.springframework.core.task.AsyncTaskExecutor;
import org.springframework.test.util.ReflectionTestUtils;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
import java.lang.reflect.Field;
import java.time.LocalDateTime;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.Callable;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.Future;
import java.util.concurrent.RejectedExecutionException;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import static cn.iocoder.muse.module.events.application.publish.EventsPublishServiceImpl.PUBLISH_STATUS_ACCEPTED;
import static org.junit.jupiter.api.Assertions.*;
import static org.mockito.ArgumentMatchers.*;
import static org.mockito.Mockito.*;
/**
* Events 统一 SSE stream 应用服务测试。
*/
class EventsStreamServiceTest {
@Mock
private UnifiedEventMapper unifiedEventMapper;
private AutoCloseable mocks;
private EventsStreamServiceImpl streamService;
private RecordingAsyncTaskExecutor executor;
@BeforeEach
void setUp() {
mocks = MockitoAnnotations.openMocks(this);
executor = new RecordingAsyncTaskExecutor(false, false);
streamService = new EventsStreamServiceImpl(unifiedEventMapper, new EventsPayloadSanitizer(), executor);
TenantContextHolder.setTenantId(100L);
}
@AfterEach
void tearDown() throws Exception {
TenantContextHolder.clear();
mocks.close();
}
@Test
void should_replayAcceptedEventsAfterCursor_forCurrentTenantAndOwner() {
when(unifiedEventMapper.selectVisibleEventsForOwner(eq(100L), eq(10L), eq(2L), any(LocalDateTime.class)))
.thenReturn(List.of(
event(3L, 100L, 10L, "chunk", Map.of("content", "第一段", "sequenceNo", 1)),
event(4L, 100L, 10L, "notification", Map.of("type", "quota_alert", "message", "额度提醒",
"resourceRef", Map.of("resourceType", "quota", "resourceId", 7L,
"ignored", "不能透出")))));
SseEmitter emitter = streamService.streamEvents(10L, "1", "muse:2");
String text = sentText(emitter);
assertTrue(text.contains("id:muse:3"));
assertTrue(text.contains("event:chunk"));
assertTrue(text.contains("id:muse:4"));
assertTrue(text.contains("event:notification"));
assertEquals(List.of(
Map.of("content", "第一段", "sequenceNo", 1L),
Map.of("type", "quota_alert", "message", "额度提醒",
"resourceRef", Map.of("resourceType", "quota", "resourceId", 7L))), sentMaps(emitter));
assertEquals(1, executor.submitCount);
}
@Test
void should_notReplayHistoricalEvents_whenLastEventIdMissing() {
when(unifiedEventMapper.selectMaxVisibleSequenceForOwner(eq(100L), eq(10L), any(LocalDateTime.class)))
.thenReturn(9L);
SseEmitter emitter = streamService.streamEvents(10L, "1", null);
assertTrue(sentText(emitter).contains("heartbeat"));
assertTrue(sentMaps(emitter).isEmpty(), "空 cursor 首连只能发 heartbeat,不能 replay 历史事件");
verify(unifiedEventMapper, never()).selectVisibleEventsForOwner(anyLong(), anyLong(), anyLong(), any());
assertEquals(1, executor.submitCount);
}
@Test
void should_notReplayCrossTenantOrCrossOwnerEvents() {
List<UnifiedEventDO> rows = List.of(
event(2L, 100L, 10L, "chunk", Map.of("content", "visible", "sequenceNo", 1)),
event(3L, 200L, 10L, "chunk", Map.of("content", "cross tenant", "sequenceNo", 2)),
event(4L, 100L, 20L, "chunk", Map.of("content", "cross owner", "sequenceNo", 3)));
when(unifiedEventMapper.selectVisibleEventsForOwner(anyLong(), anyLong(), eq(1L), any(LocalDateTime.class)))
.thenAnswer(invocation -> rows.stream()
.filter(row -> row.getTenantId().equals(invocation.getArgument(0))
&& row.getOwnerUserId().equals(invocation.getArgument(1)))
.toList());
SseEmitter emitter = streamService.streamEvents(10L, "1", "muse:1");
String payloadJson = JsonUtils.toJsonString(sentMaps(emitter));
assertTrue(payloadJson.contains("visible"));
assertFalse(payloadJson.contains("cross tenant"));
assertFalse(payloadJson.contains("cross owner"));
verify(unifiedEventMapper).selectVisibleEventsForOwner(eq(100L), eq(10L), eq(1L), any(LocalDateTime.class));
}
@Test
void should_sendHeartbeatOnly_whenNoVisibleEvents() {
when(unifiedEventMapper.selectVisibleEventsForOwner(eq(100L), eq(10L), eq(9L), any(LocalDateTime.class)))
.thenReturn(List.of());
SseEmitter emitter = streamService.streamEvents(10L, "1", "muse:9");
assertTrue(sentText(emitter).contains("heartbeat"));
assertTrue(sentMaps(emitter).isEmpty(), "无事件时只能发送 SSE comment,不能发送假 event");
assertEquals(1, executor.submitCount);
}
@Test
void should_emitErrorEventAndComplete_whenCursorInvalid() {
SseEmitter emitter = streamService.streamEvents(10L, "1", "muse:not-a-number");
assertTrue(sentText(emitter).contains("event:error"));
assertEquals("EVENTS_CURSOR_INVALID", sentMaps(emitter).get(0).get("code"));
assertTrue(isComplete(emitter));
verifyNoInteractions(unifiedEventMapper);
assertEquals(0, executor.submitCount);
}
@Test
void should_emitErrorEventAndComplete_whenApiVersionUnsupported() {
SseEmitter emitter = streamService.streamEvents(10L, "2", null);
assertTrue(sentText(emitter).contains("event:error"));
assertEquals("EVENTS_API_VERSION_UNSUPPORTED", sentMaps(emitter).get(0).get("code"));
assertTrue(isComplete(emitter));
verifyNoInteractions(unifiedEventMapper);
assertEquals(0, executor.submitCount);
}
@Test
void should_completeWithError_whenExecutorRejectsPolling() {
executor = new RecordingAsyncTaskExecutor(false, true);
streamService = new EventsStreamServiceImpl(unifiedEventMapper, new EventsPayloadSanitizer(), executor);
when(unifiedEventMapper.selectVisibleEventsForOwner(eq(100L), eq(10L), eq(9L), any(LocalDateTime.class)))
.thenReturn(List.of());
SseEmitter emitter = streamService.streamEvents(10L, "1", "muse:9");
assertTrue(sentText(emitter).contains("event:error"));
assertEquals("EVENTS_STREAM_UNAVAILABLE", sentMaps(emitter).get(0).get("code"));
assertTrue(isComplete(emitter));
assertNull(ReflectionTestUtils.getField(emitter, "failure"),
"线程池拒绝发生在 SSE 生命周期内,应发送合同内 error event 后 complete,而不是暴露异常对象");
}
@Test
void should_cancelPollingFuture_whenEmitterCompletes() {
when(unifiedEventMapper.selectVisibleEventsForOwner(eq(100L), eq(10L), eq(9L), any(LocalDateTime.class)))
.thenReturn(List.of());
SseEmitter emitter = streamService.streamEvents(10L, "1", "muse:9");
Object completionCallback = ReflectionTestUtils.getField(emitter, "completionCallback");
ReflectionTestUtils.invokeMethod(completionCallback, "run");
assertTrue(executor.lastFuture.cancelled);
}
@Test
void should_notLeakSecretFields_whenStreamingErrorEvent() {
executor = new RecordingAsyncTaskExecutor(true, false);
streamService = new EventsStreamServiceImpl(unifiedEventMapper, new EventsPayloadSanitizer(), executor);
when(unifiedEventMapper.selectVisibleEventsForOwner(eq(100L), eq(10L), eq(9L), any(LocalDateTime.class)))
.thenReturn(List.of())
.thenThrow(new IllegalStateException("Authorization: Bearer secret-token provider raw body"));
SseEmitter emitter = streamService.streamEvents(10L, "1", "muse:9");
String text = sentText(emitter) + JsonUtils.toJsonString(sentMaps(emitter));
assertTrue(text.contains("event:error"));
assertTrue(text.contains("EVENTS_STREAM_UNAVAILABLE"));
assertFalse(text.contains("Bearer"));
assertFalse(text.contains("secret-token"));
assertFalse(text.contains("provider raw body"));
assertTrue(isComplete(emitter));
}
private static UnifiedEventDO event(Long sequenceNo, Long tenantId, Long ownerUserId, String eventType,
Map<String, Object> payload) {
UnifiedEventDO event = new UnifiedEventDO();
event.setTenantId(tenantId);
event.setOwnerUserId(ownerUserId);
event.setSequenceNo(sequenceNo);
event.setEventType(eventType);
event.setPayloadSummary(JsonUtils.toJsonString(payload));
event.setPublishStatus(PUBLISH_STATUS_ACCEPTED);
event.setDeleted(false);
event.setVisibleFrom(LocalDateTime.now().minusSeconds(1));
return event;
}
private static List<Map<String, Object>> sentMaps(SseEmitter emitter) {
List<Map<String, Object>> maps = new ArrayList<>();
for (Object data : sentData(emitter)) {
if (data instanceof Map<?, ?> map) {
@SuppressWarnings("unchecked")
Map<String, Object> casted = (Map<String, Object>) map;
maps.add(casted);
}
}
return maps;
}
private static String sentText(SseEmitter emitter) {
return String.join("\n", sentData(emitter).stream().map(String::valueOf).toList());
}
private static List<Object> sentData(SseEmitter emitter) {
Set<?> attempts = (Set<?>) ReflectionTestUtils.getField(emitter, "earlySendAttempts");
if (attempts == null) {
return List.of();
}
return attempts.stream().map(EventsStreamServiceTest::dataFrom).toList();
}
private static Object dataFrom(Object dataWithMediaType) {
try {
Field field = dataWithMediaType.getClass().getDeclaredField("data");
field.setAccessible(true);
return field.get(dataWithMediaType);
} catch (ReflectiveOperationException ex) {
throw new AssertionError("无法读取 SseEmitter earlySendAttempts", ex);
}
}
private static boolean isComplete(SseEmitter emitter) {
Boolean complete = (Boolean) ReflectionTestUtils.getField(emitter, "complete");
return Boolean.TRUE.equals(complete);
}
private static class RecordingAsyncTaskExecutor implements AsyncTaskExecutor {
private final boolean runImmediately;
private final boolean reject;
private int submitCount;
private Runnable submittedRunnable;
private RecordingFuture lastFuture;
private RecordingAsyncTaskExecutor(boolean runImmediately, boolean reject) {
this.runImmediately = runImmediately;
this.reject = reject;
}
@Override
public void execute(Runnable task) {
submit(task);
}
@Override
public Future<?> submit(Runnable task) {
if (reject) {
throw new RejectedExecutionException("events stream executor full");
}
submitCount++;
submittedRunnable = task;
lastFuture = new RecordingFuture();
if (runImmediately) {
task.run();
lastFuture.done = true;
}
return lastFuture;
}
@Override
public <T> Future<T> submit(Callable<T> task) {
throw new UnsupportedOperationException("Events stream 测试只需要 Runnable 提交路径");
}
}
private static class RecordingFuture implements Future<Object> {
private boolean cancelled;
private boolean done;
@Override
public boolean cancel(boolean mayInterruptIfRunning) {
cancelled = true;
done = true;
return true;
}
@Override
public boolean isCancelled() {
return cancelled;
}
@Override
public boolean isDone() {
return done;
}
@Override
public Object get() throws InterruptedException, ExecutionException {
return null;
}
@Override
public Object get(long timeout, TimeUnit unit) throws InterruptedException, ExecutionException, TimeoutException {
return null;
}
}
}

View File

@ -0,0 +1,110 @@
package cn.iocoder.muse.module.events.dal.mysql;
import cn.iocoder.muse.module.events.dal.dataobject.UnifiedEventDO;
import cn.iocoder.muse.framework.tenant.core.db.TenantBaseDO;
import com.baomidou.mybatisplus.annotation.TableName;
import org.junit.jupiter.api.Test;
import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Path;
import static org.junit.jupiter.api.Assertions.*;
class UnifiedEventMapperTest {
@Test
void should_mapUnifiedEventDoToV16TableWithJsonbHandler() {
TableName tableName = UnifiedEventDO.class.getAnnotation(TableName.class);
assertNotNull(tableName);
assertEquals("muse_unified_event", tableName.value());
assertTrue(tableName.autoResultMap(), "payload_summary 是 JSONB,需要 autoResultMap 启用本地 TypeHandler");
}
@Test
void should_keepMapperAndServiceAwayFromApplicationSequenceAllocation() throws IOException {
String source = readSource("src/main/java/cn/iocoder/muse/module/events/application/publish/EventsPublishServiceImpl.java")
+ readSource("src/main/java/cn/iocoder/muse/module/events/dal/mysql/UnifiedEventMapper.java");
assertFalse(source.contains("max(sequence_no)+1"));
assertFalse(source.contains("max(sequence_no) + 1"));
assertFalse(source.contains("AtomicLong"));
assertFalse(source.contains("sequenceNo++"));
}
@Test
void should_normalizeNullSourceRevisionInService() throws IOException {
String source = readSource("src/main/java/cn/iocoder/muse/module/events/application/publish/EventsPublishServiceImpl.java");
assertTrue(source.contains("SOURCE_REVISION_NONE"));
assertTrue(source.contains("__none__"));
}
@Test
void should_filterRejectedAndBlockedRowsInVisibilityQuery() throws IOException {
String source = readSource("src/main/java/cn/iocoder/muse/module/events/dal/mysql/UnifiedEventMapper.java");
String visibleQuery = source.substring(source.indexOf("selectVisibleEventsForOwner"));
assertTrue(visibleQuery.contains("PUBLISH_STATUS_ACCEPTED"));
assertFalse(visibleQuery.contains("PUBLISH_STATUS_REJECTED"));
assertFalse(visibleQuery.contains("PUBLISH_STATUS_BLOCKED"));
assertTrue(visibleQuery.contains("getDeleted"));
assertTrue(visibleQuery.contains("getVisibleFrom"));
}
@Test
void should_keepV16ColumnsCompatibleWithTenantBaseDoMapping() throws IOException {
String ddl = readRepositorySource("sql/muse/V16__extend_events_sse_schema.sql")
.toLowerCase()
.replaceAll("\\s+", " ");
assertEquals(TenantBaseDO.class, UnifiedEventDO.class.getSuperclass(),
"UnifiedEventDO 继承 TenantBaseDO 时,V16 DDL 必须补齐 BaseDO 默认映射列");
assertTrue(ddl.contains("creator varchar(64) not null default ''"),
"muse_unified_event 必须包含 BaseDO.creator 映射列");
assertTrue(ddl.contains("updater varchar(64) not null default ''"),
"muse_unified_event 必须包含 BaseDO.updater 映射列");
}
private String readSource(String relativePath) throws IOException {
Path moduleRoot = findModuleRoot();
return Files.exists(moduleRoot.resolve(relativePath)) ? Files.readString(moduleRoot.resolve(relativePath)) : "";
}
private String readRepositorySource(String relativePath) throws IOException {
Path repositoryRoot = findRepositoryRoot();
return Files.exists(repositoryRoot.resolve(relativePath)) ? Files.readString(repositoryRoot.resolve(relativePath)) : "";
}
private Path findModuleRoot() {
Path current = Path.of(System.getProperty("user.dir"));
while (current != null) {
if (Files.exists(current.resolve("muse-module-events-server/pom.xml"))) {
return current.resolve("muse-module-events-server");
}
if (Files.exists(current.resolve("pom.xml"))
&& current.getFileName() != null
&& "muse-module-events-server".equals(current.getFileName().toString())) {
return current;
}
current = current.getParent();
}
throw new IllegalStateException("无法定位 muse-module-events-server 模块根目录");
}
private Path findRepositoryRoot() {
Path current = Path.of(System.getProperty("user.dir"));
while (current != null) {
if (Files.exists(current.resolve("sql/muse/V16__extend_events_sse_schema.sql"))) {
return current;
}
if (Files.exists(current.resolve("muse-cloud/sql/muse/V16__extend_events_sse_schema.sql"))) {
return current.resolve("muse-cloud");
}
current = current.getParent();
}
throw new IllegalStateException("无法定位 muse-cloud 根目录");
}
}

View File

@ -0,0 +1,57 @@
package cn.iocoder.muse.module.events.domain;
import org.junit.jupiter.api.Test;
import static org.junit.jupiter.api.Assertions.*;
/**
* Events SSE 游标解析测试。
*/
class EventsCursorTest {
@Test
void should_acceptEmptyCursor() {
EventsCursor cursor = EventsCursor.parse(null);
EventsCursor blankCursor = EventsCursor.parse(" ");
assertTrue(cursor.isValid());
assertTrue(cursor.isInitialConnection());
assertNull(cursor.sequenceNo());
assertTrue(blankCursor.isValid());
assertTrue(blankCursor.isInitialConnection());
}
@Test
void should_acceptMuseSequenceCursor() {
EventsCursor cursor = EventsCursor.parse("muse:123");
assertTrue(cursor.isValid());
assertFalse(cursor.isInitialConnection());
assertEquals(123L, cursor.sequenceNo());
}
@Test
void should_rejectNegativeSequenceCursor() {
EventsCursor cursor = EventsCursor.parse("muse:-1");
assertFalse(cursor.isValid());
assertEquals("EVENTS_CURSOR_INVALID", cursor.errorCode());
}
@Test
void should_rejectNonNumericSequenceCursor() {
EventsCursor cursor = EventsCursor.parse("muse:not-a-number");
assertFalse(cursor.isValid());
assertEquals("EVENTS_CURSOR_INVALID", cursor.errorCode());
}
@Test
void should_rejectWrongPrefixCursor() {
EventsCursor cursor = EventsCursor.parse("event:123");
assertFalse(cursor.isValid());
assertEquals("EVENTS_CURSOR_INVALID", cursor.errorCode());
}
}

View File

@ -0,0 +1,32 @@
package cn.iocoder.muse.module.events.framework.config;
import org.junit.jupiter.api.Test;
import org.springframework.core.task.AsyncTaskExecutor;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
import java.util.concurrent.ThreadPoolExecutor;
import static org.junit.jupiter.api.Assertions.*;
/**
* Events SSE 线程池配置测试。
*/
class EventsStreamConfigurationTest {
@Test
void should_useBoundedExecutorWithAbortPolicy() {
EventsStreamConfiguration configuration = new EventsStreamConfiguration();
AsyncTaskExecutor taskExecutor = configuration.eventsStreamExecutor();
assertInstanceOf(ThreadPoolTaskExecutor.class, taskExecutor);
ThreadPoolTaskExecutor executor = (ThreadPoolTaskExecutor) taskExecutor;
assertEquals(4, executor.getCorePoolSize());
assertEquals(16, executor.getMaxPoolSize());
assertEquals(100, executor.getThreadPoolExecutor().getQueue().remainingCapacity());
assertInstanceOf(ThreadPoolExecutor.AbortPolicy.class,
executor.getThreadPoolExecutor().getRejectedExecutionHandler());
executor.shutdown();
}
}

View File

@ -0,0 +1,19 @@
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<parent>
<groupId>cn.iocoder.cloud</groupId>
<artifactId>muse</artifactId>
<version>${revision}</version>
</parent>
<modelVersion>4.0.0</modelVersion>
<modules>
<module>muse-module-events-api</module>
<module>muse-module-events-server</module>
</modules>
<artifactId>muse-module-events</artifactId>
<packaging>pom</packaging>
<name>${project.artifactId}</name>
<description>events 模块,负责统一事件流和跨 owner 事件发布契约。</description>
</project>

View File

@ -59,6 +59,11 @@
<artifactId>muse-module-meta-server</artifactId> <artifactId>muse-module-meta-server</artifactId>
<version>${revision}</version> <version>${revision}</version>
</dependency> </dependency>
<dependency>
<groupId>cn.iocoder.cloud</groupId>
<artifactId>muse-module-events-server</artifactId>
<version>${revision}</version>
</dependency>
<dependency> <dependency>
<groupId>cn.iocoder.cloud</groupId> <groupId>cn.iocoder.cloud</groupId>
<artifactId>muse-module-knowledge-server</artifactId> <artifactId>muse-module-knowledge-server</artifactId>

View File

@ -0,0 +1,610 @@
package cn.iocoder.muse.server.framework.api;
import org.flywaydb.core.Flyway;
import org.flywaydb.core.api.MigrationInfo;
import org.flywaydb.core.api.output.MigrateResult;
import org.junit.jupiter.api.Test;
import org.slf4j.LoggerFactory;
import java.nio.file.Files;
import java.nio.file.Path;
import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.PreparedStatement;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Locale;
import java.util.Map;
import java.util.Objects;
import java.util.Set;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* P1R-7a Events 真实 PostgreSQL / Flyway 迁移验收。
*
* <p>本测试会 clean 指定数据库的 public schema,因此只允许连接真实库名以 {@code _test} 结尾的隔离库。
* 数据库密码只能从环境变量读取,避免 Surefire XML 或命令历史记录泄露凭据。</p>
*/
class P1rEventsFlywayMigrationIT {
/** P1R-7a Events V16 是本任务唯一允许验收到达的目标版本。 */
private static final String TARGET_VERSION = "16";
/** V16 Flyway 文件名会转换出的描述,作为“当前版本确实是 Events SSE 扩展迁移”的硬证据。 */
private static final String TARGET_DESCRIPTION = "extend events sse schema";
/** clean 后从 V1 到 V16 应该产生 16 条成功 SQL migration 记录。 */
private static final int EXPECTED_SUCCESSFUL_SQL_MIGRATIONS = 16;
private static final String UNIFIED_EVENT_TABLE = "muse_unified_event";
private static final String UNIFIED_EVENT_SEQUENCE = "muse_unified_event_sequence_no_seq";
private static final List<String> REQUIRED_INDEXES = List.of(
"idx_muse_unified_event_owner_sequence"
);
private static final List<String> REQUIRED_CONSTRAINTS = List.of(
"uk_muse_unified_event_sequence",
"uk_muse_unified_event_command",
"uk_muse_unified_event_event",
"uk_muse_unified_event_source",
"chk_muse_unified_event_event_type",
"chk_muse_unified_event_publish_status"
);
private static final List<String> REQUIRED_TRIGGERS = List.of(
"trg_muse_unified_event_updated_at"
);
private static final Set<String> CREDENTIAL_QUERY_KEYS = Set.of(
"user",
"username",
"password",
"pass",
"pwd",
"sslpassword",
"ssl_password",
"token",
"secret",
"api_key",
"apikey",
"bearer",
"access_token",
"refresh_token"
);
@Test
void should_migrateV1ToV16OnRealPostgresqlAndVerifyEventsSchema() throws Exception {
String url = requiredProperty("p1r.flyway.url");
String user = requiredProperty("p1r.flyway.user");
String password = requiredPasswordEnvironment();
String requestedLocations = requiredProperty("p1r.flyway.locations");
redactFlywaySystemProperties(url, user);
assertEquals("filesystem:sql/muse", requestedLocations,
"P1R-7a Events Flyway IT 要求显式使用 filesystem:sql/muse");
String effectiveLocations = resolveMuseSqlLocation(requestedLocations);
silenceFlywayInfoLogs();
assertNoCredentialQuery(url);
assertTestDatabaseUrl(url);
Flyway flyway = Flyway.configure()
.dataSource(url, user, password)
.locations(effectiveLocations)
.schemas("public")
.defaultSchema("public")
.target(TARGET_VERSION)
.cleanDisabled(false)
.load();
cleanSchema(flyway, url, user);
MigrateResult result = migrateSchema(flyway, url, user);
MigrationInfo current = flyway.info().current();
assertNotNull(current, "必须存在当前 Flyway 版本");
assertEquals(TARGET_VERSION, Objects.requireNonNull(current.getVersion(), "当前版本必须包含版本号").getVersion(),
"真实 Flyway 当前版本必须到 V16");
assertEquals(TARGET_DESCRIPTION, current.getDescription(),
"真实 Flyway 当前版本必须是 Events V16 扩展迁移");
assertTrue(result.migrationsExecuted >= EXPECTED_SUCCESSFUL_SQL_MIGRATIONS,
"clean 后必须至少执行 V1-V16 共 16 个迁移,实际: " + result.migrationsExecuted);
int successfulMigrationCount = successfulMigrationCount(url, user, password);
assertEquals(EXPECTED_SUCCESSFUL_SQL_MIGRATIONS, successfulMigrationCount,
"必须在隔离库中执行 V1-V16 共 16 个成功 SQL 迁移");
try (Connection connection = DriverManager.getConnection(url, user, password)) {
assertRequiredObjects(connection);
assertUnifiedEventSchema(connection);
}
System.out.println("flyway_success=true");
System.out.println("flyway_url=" + maskedUrl(url));
System.out.println("flyway_locations_requested=" + requestedLocations);
System.out.println("flyway_locations_effective=" + effectiveLocations);
System.out.println("migrations_executed=" + result.migrationsExecuted);
System.out.println("successful_migration_count=" + successfulMigrationCount);
System.out.println("target_schema_version=" + TARGET_VERSION);
System.out.println("flyway_latest=" + current.getVersion() + ":" + current.getDescription());
System.out.println("schema_version=" + current.getVersion());
System.out.println("v16_table=" + UNIFIED_EVENT_TABLE);
System.out.println("v16_sequence=" + UNIFIED_EVENT_SEQUENCE);
System.out.println("v16_indexes=" + String.join(",", REQUIRED_INDEXES));
System.out.println("v16_constraints=" + String.join(",", REQUIRED_CONSTRAINTS));
System.out.println("v16_triggers=" + String.join(",", REQUIRED_TRIGGERS));
}
@Test
void should_rejectNonTestDatabaseWhenQueryContainsSlashTestSuffix() {
String url = "jdbc:postgresql://prod-host/prod?applicationName=/muse_local_p1r7a_test";
AssertionError error = assertThrows(AssertionError.class, () -> assertTestDatabaseUrl(url));
assertTrue(error.getMessage().contains("jdbc:postgresql://<host>/prod?<query-redacted>"),
"拒绝信息必须只展示真实 database name,并脱敏 query: " + error.getMessage());
}
@Test
void should_allowTestDatabaseWithQueryAndMaskRealDatabaseName() {
String url = "jdbc:postgresql://prod-host:5432/muse_local_p1r7a_test?applicationName=/muse-local&sslmode=disable";
assertTestDatabaseUrl(url);
assertEquals("jdbc:postgresql://<host>:5432/muse_local_p1r7a_test?<query-redacted>", maskedUrl(url),
"合法 _test 库携带 query 时必须保留真实 database name,并整体脱敏 query");
}
private static String requiredProperty(String name) {
String value = System.getProperty(name);
assertTrue(value != null && !value.isBlank(), "缺少必需系统属性: " + name);
return value;
}
private static String requiredPasswordEnvironment() {
String password = firstNonBlankEnvironment("P1R_FLYWAY_PASSWORD", "MUSE_POSTGRES_PASSWORD");
assertTrue(password != null, "缺少必需环境变量: P1R_FLYWAY_PASSWORD 或 MUSE_POSTGRES_PASSWORD");
return password;
}
private static String firstNonBlankEnvironment(String... names) {
// WHY:数据库密码不能通过 JVM system property 传入,否则 Surefire XML 可能记录属性名和值。
for (String name : names) {
String value = System.getenv(name);
if (value != null && !value.isBlank()) {
return value;
}
}
return null;
}
private static String resolveMuseSqlLocation(String requestedLocations) {
Path current = Path.of(System.getProperty("user.dir")).toAbsolutePath();
String relativeLocation = requestedLocations.substring("filesystem:".length());
for (Path cursor = current; cursor != null; cursor = cursor.getParent()) {
Path candidate = cursor.resolve(relativeLocation);
if (Files.isDirectory(candidate)) {
return "filesystem:" + candidate;
}
}
throw new IllegalStateException("无法从当前目录向上找到 sql/muse: " + current);
}
private static void silenceFlywayInfoLogs() {
try {
Object flywayLogger = LoggerFactory.getLogger("org.flywaydb");
Class<?> levelClass = Class.forName("ch.qos.logback.classic.Level");
Object warnLevel = levelClass.getField("WARN").get(null);
flywayLogger.getClass().getMethod("setLevel", levelClass).invoke(flywayLogger, warnLevel);
} catch (ReflectiveOperationException | LinkageError ignored) {
// WHY:日志实现不是 logback 时不影响迁移验收;测试自身仍只输出脱敏 URL。
}
}
private static void redactFlywaySystemProperties(String url, String user) {
// WHY:Surefire XML 会记录 JVM system property;读取后立即脱敏,避免报告文件残留真实连接信息。
System.setProperty("p1r.flyway.url", maskedUrl(url));
System.setProperty("p1r.flyway.user", maskUser(user));
}
private static void cleanSchema(Flyway flyway, String url, String user) {
try {
flyway.clean();
} catch (RuntimeException exception) {
throw sanitizedFlywayFailure("Flyway clean 失败", url, user, exception);
}
}
private static MigrateResult migrateSchema(Flyway flyway, String url, String user) {
try {
return flyway.migrate();
} catch (RuntimeException exception) {
throw sanitizedFlywayFailure("Flyway migrate 失败", url, user, exception);
}
}
private static AssertionError sanitizedFlywayFailure(String action, String url, String user, RuntimeException exception) {
// WHY:Flyway 连接异常默认会打印原始 JDBC URL;这里统一替换成 masked URL,避免测试日志暴露环境细节。
String sanitizedMessage = Objects.toString(exception.getMessage(), "")
.replace(url, maskedUrl(url))
.replace("for user '" + user + "'", "for user '<user-redacted>'");
return new AssertionError(action + ": " + sanitizedMessage);
}
private static int successfulMigrationCount(String url, String user, String password) throws SQLException {
try (Connection connection = DriverManager.getConnection(url, user, password);
PreparedStatement statement = connection.prepareStatement(
"SELECT COUNT(*) FROM flyway_schema_history WHERE success = TRUE AND type = 'SQL'");
ResultSet resultSet = statement.executeQuery()) {
assertTrue(resultSet.next(), "必须能读取 flyway_schema_history");
return resultSet.getInt(1);
}
}
private static void assertRequiredObjects(Connection connection) throws SQLException {
assertTrue(tableExists(connection, UNIFIED_EVENT_TABLE),
"Events 统一事件投影表必须存在: " + UNIFIED_EVENT_TABLE);
assertTrue(sequenceExists(connection, UNIFIED_EVENT_SEQUENCE),
"Events sequence 必须存在: " + UNIFIED_EVENT_SEQUENCE);
Map<String, Boolean> indexChecks = indexChecks(connection);
assertFalse(indexChecks.containsValue(false), "Events 关键索引必须全部存在: " + indexChecks);
Map<String, Boolean> constraintChecks = constraintChecks(connection);
assertFalse(constraintChecks.containsValue(false), "Events 关键约束必须全部存在: " + constraintChecks);
Map<String, Boolean> triggerChecks = triggerChecks(connection);
assertFalse(triggerChecks.containsValue(false), "Events 更新时间触发器必须全部存在: " + triggerChecks);
}
private static boolean tableExists(Connection connection, String tableName) throws SQLException {
return exists(connection,
"SELECT 1 FROM information_schema.tables WHERE table_schema = 'public' AND table_name = ?",
tableName);
}
private static boolean sequenceExists(Connection connection, String sequenceName) throws SQLException {
return exists(connection,
"SELECT 1 FROM information_schema.sequences WHERE sequence_schema = 'public' AND sequence_name = ?",
sequenceName);
}
private static Map<String, Boolean> indexChecks(Connection connection) throws SQLException {
Map<String, Boolean> checks = new LinkedHashMap<>();
for (String indexName : REQUIRED_INDEXES) {
checks.put(indexName, exists(connection,
"SELECT 1 FROM pg_indexes WHERE schemaname = 'public' AND indexname = ?",
indexName));
}
return checks;
}
private static Map<String, Boolean> constraintChecks(Connection connection) throws SQLException {
Map<String, Boolean> checks = new LinkedHashMap<>();
for (String constraintName : REQUIRED_CONSTRAINTS) {
checks.put(constraintName, exists(connection,
"""
SELECT 1
FROM pg_constraint c
JOIN pg_namespace n ON n.oid = c.connamespace
WHERE n.nspname = 'public'
AND c.conname = ?
""",
constraintName));
}
return checks;
}
private static Map<String, Boolean> triggerChecks(Connection connection) throws SQLException {
Map<String, Boolean> checks = new LinkedHashMap<>();
for (String triggerName : REQUIRED_TRIGGERS) {
checks.put(triggerName, exists(connection,
"""
SELECT 1
FROM information_schema.triggers
WHERE trigger_schema = 'public'
AND trigger_name = ?
""",
triggerName));
}
return checks;
}
private static boolean exists(Connection connection, String sql, String value) throws SQLException {
try (PreparedStatement statement = connection.prepareStatement(sql)) {
statement.setString(1, value);
try (ResultSet resultSet = statement.executeQuery()) {
return resultSet.next();
}
}
}
private static void assertUnifiedEventSchema(Connection connection) throws SQLException {
assertColumnsExist(connection, UNIFIED_EVENT_TABLE,
"id", "tenant_id", "command_id", "event_id", "sequence_no", "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", "deleted",
"creator", "create_time", "updater", "update_time");
assertEquals("jsonb", columnUdtName(connection, UNIFIED_EVENT_TABLE, "payload_summary"),
"payload_summary 必须使用 jsonb 保存可公开事件摘要");
assertTrue(columnDefault(connection, UNIFIED_EVENT_TABLE, "sequence_no")
.contains("nextval('muse_unified_event_sequence_no_seq"),
"sequence_no 必须由 PostgreSQL sequence 默认分配");
assertEquals("__none__", columnDefaultLiteral(connection, UNIFIED_EVENT_TABLE, "source_revision"),
"source_revision 默认值必须是 __none__ sentinel");
assertEquals("'{}'::jsonb", normalizeSql(columnDefault(connection, UNIFIED_EVENT_TABLE, "payload_summary")),
"payload_summary 默认值必须是空 JSONB 对象");
assertEquals("false", normalizeSql(columnDefault(connection, UNIFIED_EVENT_TABLE, "deleted")),
"deleted 默认值必须是 false");
assertEquals("", columnDefaultLiteral(connection, UNIFIED_EVENT_TABLE, "creator"),
"creator 默认值必须与 BaseDO 字段约定保持一致");
assertEquals("", columnDefaultLiteral(connection, UNIFIED_EVENT_TABLE, "updater"),
"updater 默认值必须与 BaseDO 字段约定保持一致");
assertTrue(columnIsNotNullable(connection, UNIFIED_EVENT_TABLE, "command_id"),
"command_id 必须 NOT NULL");
assertTrue(columnIsNotNullable(connection, UNIFIED_EVENT_TABLE, "source_revision"),
"source_revision 必须 NOT NULL");
assertTrue(columnIsNotNullable(connection, UNIFIED_EVENT_TABLE, "creator"),
"creator 必须 NOT NULL,避免 BaseDO 映射列缺失或出现空审计值");
assertTrue(columnIsNotNullable(connection, UNIFIED_EVENT_TABLE, "updater"),
"updater 必须 NOT NULL,避免 BaseDO 映射列缺失或出现空审计值");
assertEquals(List.of("tenant_id", "command_id"),
constraintColumns(connection, "uk_muse_unified_event_command"),
"命令幂等唯一约束必须覆盖 tenant_id/command_id");
assertEquals(List.of("tenant_id", "sequence_no"),
constraintColumns(connection, "uk_muse_unified_event_sequence"),
"sequence 唯一约束必须覆盖 tenant_id/sequence_no");
assertEquals(List.of("tenant_id", "event_id"),
constraintColumns(connection, "uk_muse_unified_event_event"),
"event 唯一约束必须覆盖 tenant_id/event_id");
assertEquals(List.of("tenant_id", "source_owner", "source_type", "source_id", "source_revision", "event_type"),
constraintColumns(connection, "uk_muse_unified_event_source"),
"source tuple 唯一约束必须覆盖 source_revision sentinel");
assertEquals(List.of("tenant_id", "owner_user_id", "sequence_no"),
indexColumns(connection, "idx_muse_unified_event_owner_sequence"),
"SSE replay 索引必须覆盖 tenant_id/owner_user_id/sequence_no");
String eventTypeConstraint = normalizedConstraintDefinition(connection, "chk_muse_unified_event_event_type");
for (String value : Set.of("chunk", "quality_check", "done", "error", "notification")) {
assertTrue(eventTypeConstraint.contains(value),
"event_type check 必须允许 " + value + ": " + eventTypeConstraint);
}
String publishStatusConstraint = normalizedConstraintDefinition(connection, "chk_muse_unified_event_publish_status");
for (String value : Set.of("accepted", "rejected", "blocked")) {
assertTrue(publishStatusConstraint.contains(value),
"publish_status check 必须允许 " + value + ": " + publishStatusConstraint);
}
}
private static void assertColumnsExist(Connection connection, String tableName, String... columnNames) throws SQLException {
for (String columnName : columnNames) {
assertTrue(columnExists(connection, tableName, columnName),
tableName + " 必须包含字段: " + columnName);
}
}
private static boolean columnExists(Connection connection, String tableName, String columnName) throws SQLException {
try (PreparedStatement statement = connection.prepareStatement("""
SELECT 1
FROM information_schema.columns
WHERE table_schema = 'public'
AND table_name = ?
AND column_name = ?
""")) {
statement.setString(1, tableName);
statement.setString(2, columnName);
try (ResultSet resultSet = statement.executeQuery()) {
return resultSet.next();
}
}
}
private static boolean columnIsNotNullable(Connection connection, String tableName, String columnName) throws SQLException {
try (PreparedStatement statement = connection.prepareStatement("""
SELECT is_nullable
FROM information_schema.columns
WHERE table_schema = 'public'
AND table_name = ?
AND column_name = ?
""")) {
statement.setString(1, tableName);
statement.setString(2, columnName);
try (ResultSet resultSet = statement.executeQuery()) {
assertTrue(resultSet.next(), tableName + " 缺少字段: " + columnName);
return "NO".equals(resultSet.getString(1));
}
}
}
private static String columnUdtName(Connection connection, String tableName, String columnName) throws SQLException {
try (PreparedStatement statement = connection.prepareStatement("""
SELECT udt_name
FROM information_schema.columns
WHERE table_schema = 'public'
AND table_name = ?
AND column_name = ?
""")) {
statement.setString(1, tableName);
statement.setString(2, columnName);
try (ResultSet resultSet = statement.executeQuery()) {
assertTrue(resultSet.next(), tableName + " 缺少字段: " + columnName);
return resultSet.getString(1);
}
}
}
private static String columnDefault(Connection connection, String tableName, String columnName) throws SQLException {
try (PreparedStatement statement = connection.prepareStatement("""
SELECT column_default
FROM information_schema.columns
WHERE table_schema = 'public'
AND table_name = ?
AND column_name = ?
""")) {
statement.setString(1, tableName);
statement.setString(2, columnName);
try (ResultSet resultSet = statement.executeQuery()) {
assertTrue(resultSet.next(), tableName + " 缺少字段: " + columnName);
return Objects.toString(resultSet.getString(1), "");
}
}
}
private static String columnDefaultLiteral(Connection connection, String tableName, String columnName) throws SQLException {
String value = normalizeSql(columnDefault(connection, tableName, columnName));
int quotedValueStart = value.indexOf('\'');
int quotedValueEnd = value.indexOf('\'', quotedValueStart + 1);
assertTrue(quotedValueStart >= 0 && quotedValueEnd > quotedValueStart,
"字段默认值必须是字符串字面量: " + tableName + "." + columnName + "=" + value);
return value.substring(quotedValueStart + 1, quotedValueEnd);
}
private static String normalizedConstraintDefinition(Connection connection, String constraintName) throws SQLException {
try (PreparedStatement statement = connection.prepareStatement("""
SELECT pg_get_constraintdef(c.oid)
FROM pg_constraint c
JOIN pg_namespace n ON n.oid = c.connamespace
WHERE n.nspname = 'public'
AND c.conname = ?
""")) {
statement.setString(1, constraintName);
try (ResultSet resultSet = statement.executeQuery()) {
assertTrue(resultSet.next(), "缺少约束定义: " + constraintName);
return normalizeSql(resultSet.getString(1));
}
}
}
private static List<String> indexColumns(Connection connection, String indexName) throws SQLException {
String sql = """
SELECT a.attname
FROM pg_class i
JOIN pg_namespace ni ON ni.oid = i.relnamespace
JOIN pg_index ix ON ix.indexrelid = i.oid
JOIN pg_class t ON t.oid = ix.indrelid
JOIN pg_namespace nt ON nt.oid = t.relnamespace
JOIN pg_attribute a ON a.attrelid = ix.indrelid AND a.attnum = ANY(ix.indkey)
WHERE ni.nspname = 'public'
AND nt.nspname = 'public'
AND i.relname = ?
ORDER BY array_position(ix.indkey, a.attnum)
""";
try (PreparedStatement statement = connection.prepareStatement(sql)) {
statement.setString(1, indexName);
try (ResultSet resultSet = statement.executeQuery()) {
List<String> columns = new ArrayList<>();
while (resultSet.next()) {
columns.add(resultSet.getString(1));
}
return columns;
}
}
}
private static List<String> constraintColumns(Connection connection, String constraintName) throws SQLException {
String sql = """
SELECT a.attname
FROM pg_constraint c
JOIN pg_namespace nc ON nc.oid = c.connamespace
JOIN pg_class t ON t.oid = c.conrelid
JOIN pg_namespace nt ON nt.oid = t.relnamespace
JOIN pg_attribute a ON a.attrelid = t.oid AND a.attnum = ANY(c.conkey)
WHERE nc.nspname = 'public'
AND nt.nspname = 'public'
AND c.conname = ?
ORDER BY array_position(c.conkey, a.attnum)
""";
try (PreparedStatement statement = connection.prepareStatement(sql)) {
statement.setString(1, constraintName);
try (ResultSet resultSet = statement.executeQuery()) {
List<String> columns = new ArrayList<>();
while (resultSet.next()) {
columns.add(resultSet.getString(1));
}
return columns;
}
}
}
private static void assertNoCredentialQuery(String url) {
int queryStart = url.indexOf('?');
if (queryStart < 0) {
return;
}
String query = url.substring(queryStart + 1);
for (String parameter : query.split("&")) {
String key = parameter;
int equalsStart = key.indexOf('=');
if (equalsStart >= 0) {
key = key.substring(0, equalsStart);
}
assertFalse(isCredentialQueryKey(key),
"p1r.flyway.url 不能携带凭据 query 参数;请通过用户名属性和密码环境变量传入");
}
}
private static boolean isCredentialQueryKey(String rawKey) {
// WHY:URL 可能会被 Surefire 或日志记录,凭据只能走单独参数或环境变量,不能藏在 JDBC query 中。
String key = rawKey.trim().toLowerCase(Locale.ROOT).replace('-', '_');
return CREDENTIAL_QUERY_KEYS.contains(key)
|| key.endsWith("_token")
|| key.endsWith("_secret")
|| key.endsWith("_password");
}
private static void assertTestDatabaseUrl(String url) {
// WHY:Flyway clean 会删除 public schema 内全部对象,只能允许真实 database name 以 _test 结尾。
// query 由调用方控制,不能参与库名判断,否则 applicationName=/xxx_test 这类参数会绕过安全闸。
String databaseName = jdbcDatabaseName(url);
assertTrue(databaseName.endsWith("_test"),
"p1r.flyway.url 必须指向 _test 后缀隔离库,避免清理非测试库: " + maskedUrl(url));
}
private static String jdbcDatabaseName(String url) {
String urlWithoutQuery = jdbcUrlWithoutQuery(url);
int databaseStart = urlWithoutQuery.lastIndexOf('/');
assertTrue(databaseStart >= 0 && databaseStart < urlWithoutQuery.length() - 1,
"p1r.flyway.url 必须包含真实数据库名: " + maskedUrl(url));
return urlWithoutQuery.substring(databaseStart + 1);
}
private static String jdbcUrlWithoutQuery(String url) {
int queryStart = url.indexOf('?');
return queryStart < 0 ? url : url.substring(0, queryStart);
}
private static String maskedUrl(String url) {
// WHY:脱敏输出既要保留真实 database name 便于排查,又不能让 query 中的 slash 干扰库名解析。
String urlWithoutQuery = jdbcUrlWithoutQuery(url);
int databaseStart = urlWithoutQuery.lastIndexOf('/');
if (databaseStart < 0) {
return maskJdbcHost(urlWithoutQuery) + maskedQuerySuffix(url);
}
String prefix = urlWithoutQuery.substring(0, databaseStart + 1);
String database = urlWithoutQuery.substring(databaseStart + 1);
return maskJdbcHost(prefix) + database + maskedQuerySuffix(url);
}
private static String maskedQuerySuffix(String url) {
return url.indexOf('?') < 0 ? "" : "?<query-redacted>";
}
private static String maskJdbcHost(String urlPart) {
return urlPart.replaceAll("//([^:/?#]+)", "//<host>");
}
private static String maskUser(String user) {
return user == null || user.isBlank() ? "<user-redacted>" : "<user-redacted>";
}
private static String normalizeSql(String sql) {
return sql.toLowerCase(Locale.ROOT).replaceAll("\\s+", " ");
}
}

View File

@ -0,0 +1,239 @@
package cn.iocoder.muse.server.framework.api;
import org.junit.jupiter.api.Test;
import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.HashSet;
import java.util.List;
import java.util.Locale;
import java.util.Set;
import java.util.regex.Matcher;
import java.util.regex.Pattern;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* P1R-7a Events V16 迁移 SQL 静态门禁。
*/
class P1rEventsMigrationSqlTest {
private static final String UNIFIED_EVENT_TABLE = "muse_unified_event";
private static final Pattern PLAINTEXT_CREDENTIAL_COLUMN_PATTERN = Pattern.compile(
"(?i)\\b\\w*(?:token|authorization|secret|bearer|api_key|apikey|password)\\w*\\s+"
+ "(?:varchar|text|jsonb|char|bytea)");
private static final Pattern UNSAFE_SEQUENCE_ALLOCATION_PATTERN = Pattern.compile(
"(?i)max\\s*\\(\\s*sequence_no\\s*\\)\\s*\\+\\s*1");
@Test
void should_includeUnifiedEventTableSequenceAndRequiredColumns() throws IOException {
String normalizedSql = normalize(readV16Sql());
String tableSql = extractCreateTable(normalizedSql, UNIFIED_EVENT_TABLE);
assertTrue(normalizedSql.contains("create sequence muse_unified_event_sequence_no_seq"),
"V16 必须创建 Events 全局 sequence");
assertTrue(tableSql.contains("id bigserial primary key")
|| tableSql.contains("id bigint generated always as identity primary key")
|| tableSql.contains("id bigint generated by default as identity primary key"),
"muse_unified_event.id 必须是 BIGSERIAL 或 PostgreSQL identity 主键");
for (String column : Set.of(
"tenant_id bigint not null",
"command_id varchar(128) not null",
"event_id varchar(128) not null",
"sequence_no bigint not null",
"event_type varchar(32) not null",
"owner_user_id bigint not null",
"source_owner varchar(64) not null",
"source_type varchar(64) not null",
"source_id varchar(128) not null",
"source_revision varchar(128) not null",
"resource_type varchar(64)",
"resource_id varchar(128)",
"payload_summary jsonb not null",
"publish_status varchar(32) not null",
"publish_error_code varchar(64)",
"visible_from timestamp not null",
"emitted_at timestamp not null",
"deleted boolean not null default false",
"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")) {
assertTrue(tableSql.contains(column),
"muse_unified_event 必须包含字段定义: " + column);
}
}
@Test
void should_includeDefaultsUniqueConstraintsIndexesChecksAndTrigger() throws IOException {
String normalizedSql = normalize(readV16Sql());
String tableSql = extractCreateTable(normalizedSql, UNIFIED_EVENT_TABLE);
assertTrue(tableSql.contains("sequence_no bigint not null default nextval('muse_unified_event_sequence_no_seq')"),
"sequence_no 必须由 PostgreSQL sequence 默认分配");
assertTrue(normalizedSql.contains("alter sequence muse_unified_event_sequence_no_seq owned by muse_unified_event.sequence_no"),
"sequence 必须归属 muse_unified_event.sequence_no,便于 schema 生命周期跟随表");
assertTrue(tableSql.contains("payload_summary jsonb not null default '{}'::jsonb"),
"payload_summary 必须是 JSONB 且默认空对象");
assertTrue(tableSql.contains("source_revision varchar(128) not null default '__none__'"),
"source_revision 必须 NOT NULL 且无版本时使用 __none__ sentinel");
assertTrue(tableSql.contains("constraint uk_muse_unified_event_sequence unique (tenant_id, sequence_no)"),
"V16 必须用 (tenant_id, sequence_no) 唯一键保证 cursor 单调唯一");
assertTrue(tableSql.contains("constraint uk_muse_unified_event_command unique (tenant_id, command_id)"),
"V16 必须用 (tenant_id, command_id) 唯一键保证发布命令幂等");
assertTrue(tableSql.contains("constraint uk_muse_unified_event_event unique (tenant_id, event_id)"),
"V16 必须用 (tenant_id, event_id) 唯一键保证事件 ID 幂等");
assertTrue(tableSql.contains("constraint uk_muse_unified_event_source unique (tenant_id, source_owner, source_type, source_id, source_revision, event_type)"),
"V16 source tuple 唯一键必须覆盖 source_revision sentinel 字段");
assertTrue(normalizedSql.contains("create index idx_muse_unified_event_owner_sequence")
&& normalizedSql.contains("on muse_unified_event(tenant_id, owner_user_id, sequence_no)"),
"V16 必须具备 tenant/owner/sequence replay 索引");
String eventTypeConstraint = extractConstraint(tableSql, "chk_muse_unified_event_event_type");
for (String value : Set.of("chunk", "quality_check", "done", "error", "notification")) {
assertTrue(eventTypeConstraint.contains("'" + value + "'"),
"event_type check 必须允许: " + value);
}
String publishStatusConstraint = extractConstraint(tableSql, "chk_muse_unified_event_publish_status");
for (String value : Set.of("accepted", "rejected", "blocked")) {
assertTrue(publishStatusConstraint.contains("'" + value + "'"),
"publish_status check 必须允许: " + value);
}
assertTrue(normalizedSql.contains("create trigger trg_muse_unified_event_updated_at")
&& normalizedSql.contains("before update on muse_unified_event")
&& normalizedSql.contains("execute function update_updated_at_column()"),
"V16 必须复用 update_updated_at_column() 更新时间 trigger");
}
@Test
void should_rejectPlaintextCredentialColumnsAndUnsafeSequenceAllocation() throws IOException {
String sql = readV16Sql();
String tableSql = extractCreateTable(normalize(sql), UNIFIED_EVENT_TABLE);
assertFalse(PLAINTEXT_CREDENTIAL_COLUMN_PATTERN.matcher(tableSql).find(),
"Events 统一投影不能新增 token / authorization / secret 等明文字段");
assertFalse(UNSAFE_SEQUENCE_ALLOCATION_PATTERN.matcher(sql).find(),
"sequence_no 必须由 PostgreSQL sequence 分配,禁止 max(sequence_no)+1");
}
@Test
void should_requireCommandAndSourceRevisionIdempotencyContracts() throws IOException {
String normalizedSql = normalize(readV16Sql());
String tableSql = extractCreateTable(normalizedSql, UNIFIED_EVENT_TABLE);
assertTrue(tableSql.contains("command_id varchar(128) not null"),
"command_id 必须 NOT NULL,不能让幂等命令绕过唯一约束");
assertTrue(tableSql.contains("constraint uk_muse_unified_event_command unique (tenant_id, command_id)"),
"command_id 必须具备 (tenant_id, command_id) 唯一约束");
assertTrue(tableSql.contains("source_revision varchar(128) not null default '__none__'"),
"source_revision 必须 NOT NULL DEFAULT '__none__',避免 PostgreSQL NULL unique 不去重");
assertTrue(tableSql.contains("constraint uk_muse_unified_event_source unique (tenant_id, source_owner, source_type, source_id, source_revision, event_type)"),
"source tuple 唯一约束必须覆盖 source_revision");
}
@Test
void should_not_duplicateNewObjectNames() throws IOException {
String normalizedSql = normalize(readV16Sql());
assertNoDuplicateNames(extractNames(normalizedSql, Pattern.compile("\\bcreate sequence (muse_[a-z0-9_]+)\\b")),
"V16 新增 sequence 名不能重复");
assertNoDuplicateNames(extractNames(normalizedSql, Pattern.compile("\\bcreate table (muse_[a-z0-9_]+)\\b")),
"V16 新增表名不能重复");
assertNoDuplicateNames(extractNames(normalizedSql, Pattern.compile("\\bcreate (?:unique )?index (\\w+)\\b")),
"V16 新增索引名不能重复");
assertNoDuplicateNames(extractNames(normalizedSql, Pattern.compile("\\bcreate trigger (\\w+)\\b")),
"V16 新增 trigger 名不能重复");
assertNoDuplicateNames(extractNames(normalizedSql, Pattern.compile("\\bconstraint (\\w+)\\b")),
"V16 新增 constraint 名不能重复");
}
@Test
void should_requireEventsFlywayMigrationItToTargetV16AndGuardCredentials() throws IOException {
Path flywayTestPath = findRepositoryRoot().resolve(
"muse-cloud/muse-server/src/test/java/cn/iocoder/muse/server/framework/api/P1rEventsFlywayMigrationIT.java");
assertTrue(Files.exists(flywayTestPath),
"新增 V16 后必须补充真实 PostgreSQL / Flyway IT 迁移验收");
String source = Files.readString(flywayTestPath);
assertTrue(source.contains("TARGET_VERSION = \"16\""),
"Events Flyway IT 必须验收到 V16");
assertTrue(source.contains("TARGET_DESCRIPTION = \"extend events sse schema\""),
"Events Flyway IT 必须验收 V16 文件描述");
assertTrue(source.contains("EXPECTED_SUCCESSFUL_SQL_MIGRATIONS = 16"),
"Events Flyway IT 必须断言 V1-V16 共 16 个成功 SQL 迁移");
assertTrue(source.contains("requiredPasswordEnvironment()"),
"Events Flyway IT 必须只从环境变量读取数据库密码");
assertTrue(source.contains("\"P1R_FLYWAY_PASSWORD\"")
&& source.contains("\"MUSE_POSTGRES_PASSWORD\""),
"Events Flyway IT 必须支持 P1R_FLYWAY_PASSWORD / MUSE_POSTGRES_PASSWORD");
assertFalse(source.contains("System.getProperty(\"p1r.flyway.password\")"),
"Events Flyway IT 不能从 system property 读取数据库密码");
assertTrue(source.contains("assertNoCredentialQuery(url);"),
"Events Flyway IT 必须拒绝 JDBC URL query 凭据");
assertTrue(source.contains("assertTestDatabaseUrl(url);"),
"Events Flyway IT 必须限制只清理 _test 后缀数据库");
assertTrue(source.contains("flyway_latest=\" + current.getVersion() + \":\" + current.getDescription()"),
"Events Flyway IT 必须输出 flyway_latest 证据");
}
private static String readV16Sql() throws IOException {
Path migrationPath = findRepositoryRoot().resolve("muse-cloud/sql/muse/V16__extend_events_sse_schema.sql");
assertTrue(Files.exists(migrationPath), "V16 Events 迁移 SQL 必须存在");
return Files.readString(migrationPath);
}
/**
* 从当前 Maven 执行目录逐级向上查找仓库根目录,避免 surefire 在不同模块目录执行时路径失效。
*/
private static Path findRepositoryRoot() {
Path current = Path.of("").toAbsolutePath();
for (Path candidate = current; candidate != null; candidate = candidate.getParent()) {
if (Files.exists(candidate.resolve("muse-cloud/sql/muse/V1__init_content_schema.sql"))) {
return candidate;
}
}
return current;
}
private static String normalize(String sql) {
return sql.toLowerCase(Locale.ROOT)
.replaceAll("--.*", " ")
.replaceAll("\\s+", " ");
}
private static String extractCreateTable(String normalizedSql, String tableName) {
String startToken = "create table " + tableName + " ";
int start = normalizedSql.indexOf(startToken);
assertTrue(start >= 0, "V16 必须创建表: " + tableName);
int end = normalizedSql.indexOf(");", start + startToken.length());
assertTrue(end >= 0, "V16 表定义必须以 ); 结束: " + tableName);
return normalizedSql.substring(start, end + 2);
}
private static String extractConstraint(String tableSql, String constraintName) {
String startToken = "constraint " + constraintName;
int start = tableSql.indexOf(startToken);
assertTrue(start >= 0, "V16 必须创建约束: " + constraintName);
int next = tableSql.indexOf("constraint ", start + startToken.length());
return next >= 0 ? tableSql.substring(start, next) : tableSql.substring(start);
}
private static List<String> extractNames(String normalizedSql, Pattern pattern) {
Matcher matcher = pattern.matcher(normalizedSql);
return matcher.results()
.map(result -> result.group(1))
.toList();
}
private static void assertNoDuplicateNames(List<String> names, String message) {
Set<String> uniqueNames = new HashSet<>(names);
assertEquals(uniqueNames.size(), names.size(), message + ": " + names);
}
}

View File

@ -0,0 +1,207 @@
package cn.iocoder.muse.server.framework.api;
import org.junit.jupiter.api.Test;
import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.LinkedHashMap;
import java.util.Map;
import java.util.Objects;
import java.util.Set;
import java.util.TreeMap;
import java.util.regex.Matcher;
import java.util.regex.Pattern;
import java.util.stream.Stream;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* P1R-7a Task 2 Events 路由归属静态门禁。
*
* <p>该测试只证明统一事件流路由已迁移到 Events owner 骨架,
* 不验证真实 SSE replay、heartbeat、error、cursor 语义。</p>
*/
class P1rEventsRouteOwnershipTest {
private static final String MUSE_CLOUD_ROOT = "muse-cloud";
private static final String CONTENT_EVENTS_CONTROLLER = "muse-cloud/muse-module-content/"
+ "muse-module-content-server/src/main/java/cn/iocoder/muse/module/content/controller/app/AppMuseEventsController.java";
private static final String EVENTS_SERVER_POM = "muse-cloud/muse-module-events/muse-module-events-server/pom.xml";
private static final String MUSE_SERVER_POM = "muse-cloud/muse-server/pom.xml";
private static final String STREAM_EVENTS_ROUTE = "GET /app-api/muse/events";
private static final Pattern METHOD_MAPPING_PATTERN = Pattern.compile(
"@(GetMapping|PostMapping|PatchMapping|PutMapping|DeleteMapping|RequestMapping)(?:\\s*\\((.*?)\\))?",
Pattern.DOTALL);
private static final Pattern REQUEST_METHOD_PATTERN = Pattern.compile("RequestMethod\\.([A-Z]+)");
@Test
void should_keep_stream_events_route_owned_only_by_events_module_controller() throws IOException {
Map<Path, Boolean> owners = collectStreamEventsRouteOwners();
assertEquals(1, owners.size(), "P1R-7a streamEvents 路由只能有一个 Controller owner");
Path owner = owners.keySet().iterator().next();
assertTrue(owner.toString().contains("/cn/iocoder/muse/module/events/controller/app/"),
"streamEvents 必须由 Events module app Controller 持有: " + owner);
assertTrue(Files.readString(owner).contains("streamEvents("),
"Events owner Controller 必须声明 streamEvents 方法: " + owner);
}
@Test
void should_retire_content_events_controller_route_owner() throws IOException {
Path sourceFile = findRepositoryRoot().resolve(CONTENT_EVENTS_CONTROLLER);
if (!Files.exists(sourceFile)) {
return;
}
String source = Files.readString(sourceFile);
Map<String, Path> routes = new TreeMap<>();
collectRoutesFromFile(sourceFile, source, routes);
assertFalse(routes.containsKey(STREAM_EVENTS_ROUTE),
"Content 旧 AppMuseEventsController 不能继续声明 /muse/events: " + sourceFile);
assertFalse(source.contains("streamEvents("),
"Content 旧 AppMuseEventsController 不能继续持有 streamEvents 方法: " + sourceFile);
}
@Test
void should_not_let_events_server_depend_on_source_owner_servers() throws IOException {
Path pom = findRepositoryRoot().resolve(EVENTS_SERVER_POM);
assertTrue(Files.exists(pom), "Events server pom.xml 必须存在");
String source = Files.readString(pom);
for (String forbiddenArtifactId : Set.of(
"muse-module-ai-server",
"muse-module-knowledge-server",
"muse-module-market-server",
"muse-module-member-server",
"muse-module-content-server")) {
assertFalse(source.contains("<artifactId>" + forbiddenArtifactId + "</artifactId>"),
"Events server 不能反向依赖 source owner server: " + forbiddenArtifactId);
}
}
@Test
void should_assemble_events_server_in_muse_server() throws IOException {
Path pom = findRepositoryRoot().resolve(MUSE_SERVER_POM);
String source = Files.readString(pom);
assertTrue(source.contains("<artifactId>muse-module-events-server</artifactId>"),
"muse-server 必须装配 muse-module-events-server");
}
private static Map<Path, Boolean> collectStreamEventsRouteOwners() throws IOException {
Map<Path, Boolean> owners = new LinkedHashMap<>();
Path root = findRepositoryRoot().resolve(MUSE_CLOUD_ROOT);
try (Stream<Path> stream = Files.walk(root)) {
for (Path sourceFile : stream
.filter(Files::isRegularFile)
.filter(path -> path.toString().contains("/src/main/java/"))
.filter(path -> path.toString().contains("/controller/"))
.filter(path -> path.getFileName().toString().endsWith("Controller.java"))
.toList()) {
Map<String, Path> routes = new TreeMap<>();
collectRoutesFromFile(sourceFile, Files.readString(sourceFile), routes);
if (routes.containsKey(STREAM_EVENTS_ROUTE)) {
owners.put(sourceFile, true);
}
}
}
return owners;
}
private static void collectRoutesFromFile(Path sourceFile, String source, Map<String, Path> routes) {
String sidePrefix = sourceFile.toString().contains("/controller/admin/") ? "/admin-api" : "/app-api";
String classBasePath = extractClassBasePath(source);
Matcher matcher = METHOD_MAPPING_PATTERN.matcher(source);
while (matcher.find()) {
String annotation = matcher.group(1);
String mappingBody = matcher.group(2);
for (String method : extractHttpMethods(annotation, mappingBody)) {
for (String methodPath : extractMappingPaths(mappingBody)) {
String fullPath = normalizePath(sidePrefix, classBasePath, methodPath);
String route = method + " " + fullPath;
if (STREAM_EVENTS_ROUTE.equals(route)) {
routes.put(route, sourceFile);
}
}
}
}
}
private static String extractClassBasePath(String source) {
int classIndex = source.indexOf("public class ");
String classAnnotations = classIndex >= 0 ? source.substring(0, classIndex) : source;
Matcher matcher = Pattern.compile("@RequestMapping\\s*\\((.*?)\\)", Pattern.DOTALL).matcher(classAnnotations);
if (!matcher.find()) {
return "";
}
return extractMappingPaths(matcher.group(1)).stream().findFirst().orElse("");
}
private static Set<String> extractHttpMethods(String annotation, String mappingBody) {
return switch (annotation) {
case "GetMapping" -> Set.of("GET");
case "PostMapping" -> Set.of("POST");
case "PatchMapping" -> Set.of("PATCH");
case "PutMapping" -> Set.of("PUT");
case "DeleteMapping" -> Set.of("DELETE");
default -> {
Matcher matcher = REQUEST_METHOD_PATTERN.matcher(mappingBody == null ? "" : mappingBody);
Map<String, Boolean> methods = new LinkedHashMap<>();
while (matcher.find()) {
methods.put(matcher.group(1), true);
}
yield methods.isEmpty() ? Set.of() : methods.keySet();
}
};
}
private static Set<String> extractMappingPaths(String mappingBody) {
if (mappingBody == null || mappingBody.isBlank()) {
return Set.of("");
}
Matcher matcher = Pattern.compile("\"([^\"]+)\"").matcher(mappingBody);
Map<String, Boolean> paths = new LinkedHashMap<>();
while (matcher.find()) {
String value = matcher.group(1);
if (value.startsWith("/")) {
paths.put(value, true);
}
}
return paths.isEmpty() ? Set.of("") : paths.keySet();
}
private static String normalizePath(String sidePrefix, String classBasePath, String methodPath) {
return Stream.of(sidePrefix, classBasePath, methodPath)
.filter(Objects::nonNull)
.filter(part -> !part.isBlank())
.map(part -> part.startsWith("/") ? part.substring(1) : part)
.reduce("", (left, right) -> left + "/" + right)
.replaceAll("/{2,}", "/");
}
/**
* 从当前 Maven 执行目录逐级向上查找仓库根目录,避免 surefire 在不同模块目录执行时路径失效。
*/
private static Path findRepositoryRoot() {
Path current = Path.of("").toAbsolutePath();
for (Path candidate = current; candidate != null; candidate = candidate.getParent()) {
if (Files.exists(candidate.resolve("muse-cloud/muse-server/pom.xml"))) {
return candidate;
}
}
return current;
}
}

View File

@ -18,6 +18,7 @@
<module>muse-module-infra</module> <module>muse-module-infra</module>
<module>muse-module-content</module> <module>muse-module-content</module>
<module>muse-module-meta</module> <module>muse-module-meta</module>
<module>muse-module-events</module>
<module>muse-module-knowledge</module> <module>muse-module-knowledge</module>
<module>muse-module-market</module> <module>muse-module-market</module>
<module>muse-module-member</module> <module>muse-module-member</module>

View File

@ -0,0 +1,52 @@
-- P1R-7a Events SSE 统一事件投影 Schema。
-- 迁移边界:Events owner 只保存可公开事件摘要与发布状态,不保存上游调用凭据或鉴权材料。
CREATE SEQUENCE muse_unified_event_sequence_no_seq
AS BIGINT
START WITH 1
INCREMENT BY 1
NO MINVALUE
NO MAXVALUE
CACHE 1;
CREATE TABLE muse_unified_event (
id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
tenant_id BIGINT NOT NULL,
command_id VARCHAR(128) NOT NULL,
event_id VARCHAR(128) NOT NULL,
sequence_no BIGINT NOT NULL DEFAULT nextval('muse_unified_event_sequence_no_seq'),
event_type VARCHAR(32) NOT NULL,
owner_user_id BIGINT NOT NULL,
source_owner VARCHAR(64) NOT NULL,
source_type VARCHAR(64) NOT NULL,
source_id VARCHAR(128) NOT NULL,
source_revision VARCHAR(128) NOT NULL DEFAULT '__none__',
resource_type VARCHAR(64),
resource_id VARCHAR(128),
payload_summary JSONB NOT NULL DEFAULT '{}'::jsonb,
publish_status VARCHAR(32) NOT NULL,
publish_error_code VARCHAR(64),
visible_from TIMESTAMP NOT NULL,
emitted_at TIMESTAMP NOT NULL,
deleted BOOLEAN NOT NULL DEFAULT FALSE,
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,
CONSTRAINT uk_muse_unified_event_sequence UNIQUE (tenant_id, sequence_no),
CONSTRAINT uk_muse_unified_event_command UNIQUE (tenant_id, command_id),
CONSTRAINT uk_muse_unified_event_event UNIQUE (tenant_id, event_id),
CONSTRAINT uk_muse_unified_event_source UNIQUE (tenant_id, source_owner, source_type, source_id, source_revision, event_type),
CONSTRAINT chk_muse_unified_event_event_type CHECK (event_type IN ('chunk', 'quality_check', 'done', 'error', 'notification')),
CONSTRAINT chk_muse_unified_event_publish_status CHECK (publish_status IN ('accepted', 'rejected', 'blocked'))
);
ALTER SEQUENCE muse_unified_event_sequence_no_seq
OWNED BY muse_unified_event.sequence_no;
CREATE INDEX idx_muse_unified_event_owner_sequence
ON muse_unified_event(tenant_id, owner_user_id, sequence_no);
CREATE TRIGGER trg_muse_unified_event_updated_at
BEFORE UPDATE ON muse_unified_event
FOR EACH ROW EXECUTE FUNCTION update_updated_at_column();