记录重复投递 payload 差异

This commit is contained in:
andy
2026-07-09 18:40:56 +08:00
parent 1282d4aede
commit c1f9870fce
8 changed files with 314 additions and 10 deletions

View File

@@ -0,0 +1,18 @@
package cn.nianxx.thhotel.platform.message.common.dto;
import java.time.LocalDateTime;
/**
* 重复投递 payload 差异入库草稿。只用于排查同一 SourceMessage 幂等键下的新旧 payload 差异。
*/
public record SourceMessagePayloadDuplicateDraft(
Long inboxId,
String originalPayloadSha256,
String duplicatePayloadSha256,
String payloadJson,
String providerFrameId,
String providerSessionId,
String schemaVersion,
LocalDateTime receivedAt
) {
}

View File

@@ -0,0 +1,115 @@
package cn.nianxx.thhotel.platform.message.domain;
import com.baomidou.mybatisplus.annotation.IdType;
import com.baomidou.mybatisplus.annotation.TableId;
import com.baomidou.mybatisplus.annotation.TableName;
import java.time.LocalDateTime;
/**
* 来源消息重复投递 payload 差异实体。该表用于排查重复投递内容变化,不通过普通查询接口返回。
*/
@TableName("platform_source_message_payload_duplicate")
public class SourceMessagePayloadDuplicateEntity {
/** 重复投递 payload 明细内部主键。 */
@TableId(type = IdType.ASSIGN_ID)
private Long id;
/** 所属 SourceMessage Inbox ID关联首次保存的来源消息。 */
private Long inboxId;
/** 首次保存 payload 的 SHA-256 哈希。 */
private String originalPayloadSha256;
/** 本次重复投递 payload 的 SHA-256 哈希。 */
private String duplicatePayloadSha256;
/** 本次重复投递的 AgentBus 规范 JSON 原始载荷,普通查询接口不得返回。 */
private String payloadJson;
/** 本次重复投递的 AgentBus frame ID仅用于排查推送帧。 */
private String providerFrameId;
/** 本次重复投递的 AgentBus session ID仅用于排查连接会话。 */
private String providerSessionId;
/** 本次重复投递 payload 结构版本。 */
private String schemaVersion;
/** 本项目收到本次重复投递的 UTC 时间。 */
private LocalDateTime receivedAt;
/** 记录创建 UTC 时间。 */
private LocalDateTime createdAt;
public Long getId() {
return id;
}
public void setId(Long id) {
this.id = id;
}
public Long getInboxId() {
return inboxId;
}
public void setInboxId(Long inboxId) {
this.inboxId = inboxId;
}
public String getOriginalPayloadSha256() {
return originalPayloadSha256;
}
public void setOriginalPayloadSha256(String originalPayloadSha256) {
this.originalPayloadSha256 = originalPayloadSha256;
}
public String getDuplicatePayloadSha256() {
return duplicatePayloadSha256;
}
public void setDuplicatePayloadSha256(String duplicatePayloadSha256) {
this.duplicatePayloadSha256 = duplicatePayloadSha256;
}
public String getPayloadJson() {
return payloadJson;
}
public void setPayloadJson(String payloadJson) {
this.payloadJson = payloadJson;
}
public String getProviderFrameId() {
return providerFrameId;
}
public void setProviderFrameId(String providerFrameId) {
this.providerFrameId = providerFrameId;
}
public String getProviderSessionId() {
return providerSessionId;
}
public void setProviderSessionId(String providerSessionId) {
this.providerSessionId = providerSessionId;
}
public String getSchemaVersion() {
return schemaVersion;
}
public void setSchemaVersion(String schemaVersion) {
this.schemaVersion = schemaVersion;
}
public LocalDateTime getReceivedAt() {
return receivedAt;
}
public void setReceivedAt(LocalDateTime receivedAt) {
this.receivedAt = receivedAt;
}
public LocalDateTime getCreatedAt() {
return createdAt;
}
public void setCreatedAt(LocalDateTime createdAt) {
this.createdAt = createdAt;
}
}

View File

@@ -0,0 +1,12 @@
package cn.nianxx.thhotel.platform.message.mapper;
import cn.nianxx.thhotel.platform.message.domain.SourceMessagePayloadDuplicateEntity;
import com.baomidou.mybatisplus.core.mapper.BaseMapper;
import org.apache.ibatis.annotations.Mapper;
/**
* SourceMessage 重复投递 payload 差异 Mapper不向 Controller 直接暴露。
*/
@Mapper
public interface SourceMessagePayloadDuplicateMapper extends BaseMapper<SourceMessagePayloadDuplicateEntity> {
}

View File

@@ -5,6 +5,7 @@ import cn.nianxx.thhotel.platform.message.common.dto.SourceMessageInboxSnapshot;
import cn.nianxx.thhotel.platform.message.common.dto.SourceMessageOriginalAccessAuditDraft;
import cn.nianxx.thhotel.platform.message.common.dto.SourceMessageOriginalContent;
import cn.nianxx.thhotel.platform.message.common.dto.SourceMessageOriginalMediaItem;
import cn.nianxx.thhotel.platform.message.common.dto.SourceMessagePayloadDuplicateDraft;
import cn.nianxx.thhotel.platform.message.common.enums.SourceMessageBodyContentType;
import cn.nianxx.thhotel.platform.message.common.enums.SourceMessageCaptureStatus;
import cn.nianxx.thhotel.platform.message.common.request.CaptureSourceMessageMedia;
@@ -14,11 +15,13 @@ import cn.nianxx.thhotel.platform.message.domain.SourceMessageBodyEntity;
import cn.nianxx.thhotel.platform.message.domain.SourceMessageInboxEntity;
import cn.nianxx.thhotel.platform.message.domain.SourceMessageMediaEntity;
import cn.nianxx.thhotel.platform.message.domain.SourceMessageOriginalAccessAuditEntity;
import cn.nianxx.thhotel.platform.message.domain.SourceMessagePayloadDuplicateEntity;
import cn.nianxx.thhotel.platform.message.domain.SourceMessagePayloadEntity;
import cn.nianxx.thhotel.platform.message.mapper.SourceMessageBodyMapper;
import cn.nianxx.thhotel.platform.message.mapper.SourceMessageInboxMapper;
import cn.nianxx.thhotel.platform.message.mapper.SourceMessageMediaMapper;
import cn.nianxx.thhotel.platform.message.mapper.SourceMessageOriginalAccessAuditMapper;
import cn.nianxx.thhotel.platform.message.mapper.SourceMessagePayloadDuplicateMapper;
import cn.nianxx.thhotel.platform.message.mapper.SourceMessagePayloadMapper;
import com.baomidou.mybatisplus.core.toolkit.Wrappers;
import com.baomidou.mybatisplus.extension.plugins.pagination.Page;
@@ -45,6 +48,7 @@ public class MybatisSourceMessageInboxRepository implements SourceMessageInboxRe
private final SourceMessageBodyMapper bodyMapper;
private final SourceMessageMediaMapper mediaMapper;
private final SourceMessageOriginalAccessAuditMapper originalAccessAuditMapper;
private final SourceMessagePayloadDuplicateMapper payloadDuplicateMapper;
/**
* 注入本模块 Mapper数据库访问细节只保留在 Repository 实现内部。
@@ -54,12 +58,14 @@ public class MybatisSourceMessageInboxRepository implements SourceMessageInboxRe
SourceMessagePayloadMapper payloadMapper,
SourceMessageBodyMapper bodyMapper,
SourceMessageMediaMapper mediaMapper,
SourceMessageOriginalAccessAuditMapper originalAccessAuditMapper) {
SourceMessageOriginalAccessAuditMapper originalAccessAuditMapper,
SourceMessagePayloadDuplicateMapper payloadDuplicateMapper) {
this.inboxMapper = inboxMapper;
this.payloadMapper = payloadMapper;
this.bodyMapper = bodyMapper;
this.mediaMapper = mediaMapper;
this.originalAccessAuditMapper = originalAccessAuditMapper;
this.payloadDuplicateMapper = payloadDuplicateMapper;
}
/**
@@ -256,6 +262,24 @@ public class MybatisSourceMessageInboxRepository implements SourceMessageInboxRe
inboxMapper.updateById(entity);
}
/**
* 追加保存重复投递 payload 差异明细,不覆盖首次保存的原始 payload。
*/
@Override
public void insertDuplicatePayload(SourceMessagePayloadDuplicateDraft draft) {
SourceMessagePayloadDuplicateEntity entity = new SourceMessagePayloadDuplicateEntity();
entity.setInboxId(draft.inboxId());
entity.setOriginalPayloadSha256(draft.originalPayloadSha256());
entity.setDuplicatePayloadSha256(draft.duplicatePayloadSha256());
entity.setPayloadJson(draft.payloadJson());
entity.setProviderFrameId(draft.providerFrameId());
entity.setProviderSessionId(draft.providerSessionId());
entity.setSchemaVersion(draft.schemaVersion());
entity.setReceivedAt(draft.receivedAt());
entity.setCreatedAt(draft.receivedAt());
payloadDuplicateMapper.insert(entity);
}
/**
* 读取 SourceMessage 原文内容和媒体 URL。该方法不供普通摘要查询调用。
*/

View File

@@ -4,8 +4,9 @@ import cn.nianxx.thhotel.platform.message.common.dto.SourceMessageInboxDraft;
import cn.nianxx.thhotel.platform.message.common.dto.SourceMessageInboxSnapshot;
import cn.nianxx.thhotel.platform.message.common.dto.SourceMessageOriginalAccessAuditDraft;
import cn.nianxx.thhotel.platform.message.common.dto.SourceMessageOriginalContent;
import cn.nianxx.thhotel.platform.message.common.result.SourceMessagePageResult;
import cn.nianxx.thhotel.platform.message.common.dto.SourceMessagePayloadDuplicateDraft;
import cn.nianxx.thhotel.platform.message.common.request.SourceMessageQueryRequest;
import cn.nianxx.thhotel.platform.message.common.result.SourceMessagePageResult;
import java.time.LocalDateTime;
import java.util.List;
import java.util.Map;
@@ -77,6 +78,11 @@ public interface SourceMessageInboxRepository {
*/
void markDuplicatePayloadChanged(Long id, String safeErrorSummary, LocalDateTime updatedAt);
/**
* 幂等命中且 payload 变化时,追加保存本次重复投递 payload便于数据库排查差异。
*/
void insertDuplicatePayload(SourceMessagePayloadDuplicateDraft draft);
/**
* 按内部 SourceMessage ID 读取原文内容和媒体 URL仅供受控原文读取服务使用。
*/

View File

@@ -1,14 +1,15 @@
package cn.nianxx.thhotel.platform.message.service.impl;
import cn.nianxx.thhotel.platform.message.common.dto.SourceMessageInboxDraft;
import cn.nianxx.thhotel.platform.message.common.dto.SourceMessageInboxSnapshot;
import cn.nianxx.thhotel.platform.message.common.dto.SourceMessagePayloadDuplicateDraft;
import cn.nianxx.thhotel.platform.message.common.enums.SourceMessageCaptureFailureReason;
import cn.nianxx.thhotel.platform.message.common.enums.SourceMessageCaptureStatus;
import cn.nianxx.thhotel.platform.message.repository.SourceMessageInboxRepository;
import cn.nianxx.thhotel.platform.message.common.request.CaptureSourceMessageCommand;
import cn.nianxx.thhotel.platform.message.common.request.CaptureSourceMessageMedia;
import cn.nianxx.thhotel.platform.message.common.result.SourceMessageCaptureResult;
import cn.nianxx.thhotel.platform.message.repository.SourceMessageInboxRepository;
import cn.nianxx.thhotel.platform.message.service.SourceMessageCaptureService;
import cn.nianxx.thhotel.platform.message.common.dto.SourceMessageInboxDraft;
import cn.nianxx.thhotel.platform.message.common.dto.SourceMessageInboxSnapshot;
import java.nio.charset.StandardCharsets;
import java.security.MessageDigest;
import java.security.NoSuchAlgorithmException;
@@ -17,6 +18,8 @@ import java.time.LocalDateTime;
import java.time.ZoneOffset;
import java.util.HexFormat;
import java.util.List;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.dao.DuplicateKeyException;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
@@ -27,6 +30,8 @@ import org.springframework.transaction.annotation.Transactional;
@Service
public class SourceMessageCaptureServiceImpl implements SourceMessageCaptureService {
private static final Logger LOGGER = LoggerFactory.getLogger(SourceMessageCaptureServiceImpl.class);
private static final String DUPLICATE_PAYLOAD_CHANGED_SUMMARY = "重复投递 payload 与原始 payload 不一致,已保留原始载荷。";
private final SourceMessageInboxRepository inboxRepository;
@@ -67,7 +72,7 @@ public class SourceMessageCaptureServiceImpl implements SourceMessageCaptureServ
SourceMessageInboxSnapshot existing = selectExisting(hotelId, provider, channel, externalMessageId);
if (existing != null) {
return existingResult(existing, payloadSha256);
return existingResult(existing, payloadSha256, command);
}
SourceMessageCaptureFailureReason mediaValidationError = validateMediaItems(command.mediaItems());
@@ -98,7 +103,7 @@ public class SourceMessageCaptureServiceImpl implements SourceMessageCaptureServ
// 并发重复投递时,数据库唯一键是最终幂等防线;命中后回查已有 Inbox。
SourceMessageInboxSnapshot concurrentlyCreated = selectExisting(hotelId, provider, channel, externalMessageId);
if (concurrentlyCreated != null) {
return existingResult(concurrentlyCreated, payloadSha256);
return existingResult(concurrentlyCreated, payloadSha256, command);
}
throw exception;
}
@@ -142,14 +147,50 @@ public class SourceMessageCaptureServiceImpl implements SourceMessageCaptureServ
/**
* 将已存在的 Inbox 转换为捕获结果,并在 payload 变化时只记录差异标记。
*/
private SourceMessageCaptureResult existingResult(SourceMessageInboxSnapshot existing, String payloadSha256) {
private SourceMessageCaptureResult existingResult(
SourceMessageInboxSnapshot existing,
String payloadSha256,
CaptureSourceMessageCommand command) {
boolean payloadChanged = !payloadSha256.equals(existing.payloadSha256());
if (payloadChanged && !existing.duplicatePayloadChanged()) {
inboxRepository.markDuplicatePayloadChanged(existing.id(), DUPLICATE_PAYLOAD_CHANGED_SUMMARY, nowUtc());
if (payloadChanged) {
LocalDateTime now = nowUtc();
if (!existing.duplicatePayloadChanged()) {
inboxRepository.markDuplicatePayloadChanged(existing.id(), DUPLICATE_PAYLOAD_CHANGED_SUMMARY, now);
}
recordDuplicatePayload(existing, payloadSha256, command, now);
}
return new SourceMessageCaptureResult(existing.id(), false, payloadChanged, existing.captureStatus());
}
/**
* 记录重复投递 payload 明细用于排查;诊断写入失败不影响原有幂等返回。
*/
private void recordDuplicatePayload(
SourceMessageInboxSnapshot existing,
String duplicatePayloadSha256,
CaptureSourceMessageCommand command,
LocalDateTime receivedAt) {
try {
inboxRepository.insertDuplicatePayload(new SourceMessagePayloadDuplicateDraft(
existing.id(),
existing.payloadSha256(),
duplicatePayloadSha256,
command.payloadJson() == null ? "" : command.payloadJson(),
trimToNull(command.providerFrameId()),
trimToNull(command.providerSessionId()),
trimToNull(command.schemaVersion()),
receivedAt
));
} catch (RuntimeException exception) {
LOGGER.warn(
"重复投递 payload 诊断写入失败inboxId={}, originalPayloadSha256={}, duplicatePayloadSha256={}, exceptionType={}",
existing.id(),
existing.payloadSha256(),
duplicatePayloadSha256,
exception.getClass().getName());
}
}
/**
* 校验媒体引用最小必填字段,避免保存无法追溯的附件或内嵌媒体记录。
*/

View File

@@ -0,0 +1,17 @@
-- M001 SourceMessage Inbox重复投递 payload 差异明细表,只在幂等命中且 payload 内容变化时追加保存。
CREATE TABLE platform_source_message_payload_duplicate (
id BIGINT NOT NULL COMMENT '重复投递 payload 明细 ID',
inbox_id BIGINT NOT NULL COMMENT '所属 SourceMessage Inbox 记录 ID关联首次保存的来源消息',
original_payload_sha256 CHAR(64) NOT NULL COMMENT '首次保存 payload 的 SHA-256 哈希,用于和重复投递 payload 对比',
duplicate_payload_sha256 CHAR(64) NOT NULL COMMENT '本次重复投递 payload 的 SHA-256 哈希',
payload_json LONGTEXT NOT NULL COMMENT '本次重复投递的 AgentBus 规范 JSON 原始载荷,普通查询接口不得返回',
provider_frame_id VARCHAR(128) NULL COMMENT '本次重复投递的 AgentBus frame ID仅用于排查推送帧',
provider_session_id VARCHAR(128) NULL COMMENT '本次重复投递的 AgentBus session ID仅用于排查连接会话',
schema_version VARCHAR(64) NULL COMMENT '本次重复投递 payload 结构版本AgentBus 有提供则保存',
received_at DATETIME(6) NOT NULL COMMENT '本项目收到本次重复投递的 UTC 时间',
created_at DATETIME(6) NOT NULL COMMENT '记录创建 UTC 时间',
PRIMARY KEY (id),
KEY idx_source_message_payload_duplicate_inbox_time (inbox_id, received_at),
KEY idx_source_message_payload_duplicate_hash (duplicate_payload_sha256),
KEY idx_source_message_payload_duplicate_frame (provider_frame_id)
) COMMENT='SourceMessage 重复投递 payload 差异明细表,用于排查同一外部邮件 ID 的重复投递内容变化';

View File

@@ -3,6 +3,7 @@ package cn.nianxx.thhotel.platform.message.service;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
@@ -32,6 +33,7 @@ import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.dao.DuplicateKeyException;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.test.context.ActiveProfiles;
@SpringBootTest(classes = ThHotelApplication.class)
@@ -50,6 +52,9 @@ class SourceMessageCaptureServiceImplTest {
@Autowired
private SourceMessagePayloadMapper payloadMapper;
@Autowired
private JdbcTemplate jdbcTemplate;
@Test
void shouldCaptureReceivedEmailAndExposeOnlySafeSummaryForQueries() {
CaptureSourceMessageCommand command = new CaptureSourceMessageCommand(
@@ -164,6 +169,19 @@ class SourceMessageCaptureServiceImplTest {
assertThat(payload.getPayloadJson()).contains("\"version\":1");
assertThat(payload.getPayloadJson()).doesNotContain("\"version\":2");
Long duplicatePayloadCount = jdbcTemplate.queryForObject("""
SELECT COUNT(*)
FROM platform_source_message_payload_duplicate
WHERE inbox_id = ?
AND original_payload_sha256 = ?
AND duplicate_payload_sha256 <> ?
AND payload_json LIKE '%"version":2%'
AND provider_frame_id = 'frame-duplicate-002'
AND provider_session_id = 'session-duplicate'
AND schema_version = 'agentbus-outlook-v1'
""", Long.class, firstResult.inboxId(), payload.getPayloadSha256(), payload.getPayloadSha256());
assertThat(duplicatePayloadCount).isEqualTo(1L);
SourceMessageInboxEntity inbox = inboxMapper.selectById(firstResult.inboxId());
assertThat(inbox.getDuplicatePayloadChanged()).isTrue();
assertThat(inbox.getSafeErrorSummary()).contains("重复投递 payload");
@@ -335,4 +353,57 @@ class SourceMessageCaptureServiceImplTest {
assertThat(result.duplicatePayloadChanged()).isTrue();
verify(localRepository).markDuplicatePayloadChanged(eq(99001L), any(String.class), any(LocalDateTime.class));
}
@Test
void shouldKeepIdempotentResultWhenDuplicatePayloadDiagnosticWriteFails() {
SourceMessageInboxRepository localRepository = mock(SourceMessageInboxRepository.class);
SourceMessageCaptureService localService = new SourceMessageCaptureServiceImpl(
localRepository,
new SourceMessageSafetySanitizer()
);
SourceMessageInboxSnapshot existing = new SourceMessageInboxSnapshot(
99002L,
"HOTEL-TEST",
"AGENTBUS",
"EMAIL",
"mail-diagnostic-fail-001",
"conversation-diagnostic-fail",
"different-payload-hash",
"RECEIVED",
true,
LocalDateTime.parse("2026-07-06T10:00:00"),
LocalDateTime.parse("2026-07-06T10:00:00"),
"g***@example.test",
"Duplicate diagnostic fail",
"Original body"
);
when(localRepository.findByIdempotencyKey("HOTEL-TEST", "AGENTBUS", "EMAIL", "mail-diagnostic-fail-001"))
.thenReturn(Optional.of(existing));
doThrow(new DuplicateKeyException("duplicate diagnostic payload"))
.when(localRepository).insertDuplicatePayload(any());
CaptureSourceMessageCommand command = new CaptureSourceMessageCommand(
"HOTEL-TEST",
"AGENTBUS",
"EMAIL",
"mail-diagnostic-fail-001",
"conversation-diagnostic-fail",
"frame-diagnostic-fail-001",
"session-diagnostic-fail",
Instant.parse("2026-07-06T10:00:00Z"),
"guest@example.test",
"Duplicate diagnostic fail",
"Changed body",
"<html>Changed body</html>",
"{\"source\":{\"external_message_id\":\"mail-diagnostic-fail-001\"},\"version\":2}",
"agentbus-outlook-v1",
List.of()
);
SourceMessageCaptureResult result = localService.capture(command);
assertThat(result.created()).isFalse();
assertThat(result.inboxId()).isEqualTo(99002L);
assertThat(result.duplicatePayloadChanged()).isTrue();
assertThat(result.captureStatus()).isEqualTo("RECEIVED");
}
}