From c1f9870fce44f70df4e64522f4318bac933fdd2d Mon Sep 17 00:00:00 2001 From: andy Date: Thu, 9 Jul 2026 18:40:56 +0800 Subject: [PATCH] =?UTF-8?q?=E8=AE=B0=E5=BD=95=E9=87=8D=E5=A4=8D=E6=8A=95?= =?UTF-8?q?=E9=80=92=20payload=20=E5=B7=AE=E5=BC=82?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../SourceMessagePayloadDuplicateDraft.java | 18 +++ .../SourceMessagePayloadDuplicateEntity.java | 115 ++++++++++++++++++ .../SourceMessagePayloadDuplicateMapper.java | 12 ++ .../MybatisSourceMessageInboxRepository.java | 26 +++- .../SourceMessageInboxRepository.java | 8 +- .../impl/SourceMessageCaptureServiceImpl.java | 57 +++++++-- ...reate_source_message_payload_duplicate.sql | 17 +++ .../SourceMessageCaptureServiceImplTest.java | 71 +++++++++++ 8 files changed, 314 insertions(+), 10 deletions(-) create mode 100644 server/src/main/java/cn/nianxx/thhotel/platform/message/common/dto/SourceMessagePayloadDuplicateDraft.java create mode 100644 server/src/main/java/cn/nianxx/thhotel/platform/message/domain/SourceMessagePayloadDuplicateEntity.java create mode 100644 server/src/main/java/cn/nianxx/thhotel/platform/message/mapper/SourceMessagePayloadDuplicateMapper.java create mode 100644 server/src/main/resources/db/migration/V8__create_source_message_payload_duplicate.sql diff --git a/server/src/main/java/cn/nianxx/thhotel/platform/message/common/dto/SourceMessagePayloadDuplicateDraft.java b/server/src/main/java/cn/nianxx/thhotel/platform/message/common/dto/SourceMessagePayloadDuplicateDraft.java new file mode 100644 index 0000000..faedb3f --- /dev/null +++ b/server/src/main/java/cn/nianxx/thhotel/platform/message/common/dto/SourceMessagePayloadDuplicateDraft.java @@ -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 +) { +} diff --git a/server/src/main/java/cn/nianxx/thhotel/platform/message/domain/SourceMessagePayloadDuplicateEntity.java b/server/src/main/java/cn/nianxx/thhotel/platform/message/domain/SourceMessagePayloadDuplicateEntity.java new file mode 100644 index 0000000..e7e0e37 --- /dev/null +++ b/server/src/main/java/cn/nianxx/thhotel/platform/message/domain/SourceMessagePayloadDuplicateEntity.java @@ -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; + } +} diff --git a/server/src/main/java/cn/nianxx/thhotel/platform/message/mapper/SourceMessagePayloadDuplicateMapper.java b/server/src/main/java/cn/nianxx/thhotel/platform/message/mapper/SourceMessagePayloadDuplicateMapper.java new file mode 100644 index 0000000..b8f9e70 --- /dev/null +++ b/server/src/main/java/cn/nianxx/thhotel/platform/message/mapper/SourceMessagePayloadDuplicateMapper.java @@ -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 { +} diff --git a/server/src/main/java/cn/nianxx/thhotel/platform/message/repository/MybatisSourceMessageInboxRepository.java b/server/src/main/java/cn/nianxx/thhotel/platform/message/repository/MybatisSourceMessageInboxRepository.java index f5db3ae..d7ce4f9 100644 --- a/server/src/main/java/cn/nianxx/thhotel/platform/message/repository/MybatisSourceMessageInboxRepository.java +++ b/server/src/main/java/cn/nianxx/thhotel/platform/message/repository/MybatisSourceMessageInboxRepository.java @@ -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。该方法不供普通摘要查询调用。 */ diff --git a/server/src/main/java/cn/nianxx/thhotel/platform/message/repository/SourceMessageInboxRepository.java b/server/src/main/java/cn/nianxx/thhotel/platform/message/repository/SourceMessageInboxRepository.java index 08c23b4..13a5704 100644 --- a/server/src/main/java/cn/nianxx/thhotel/platform/message/repository/SourceMessageInboxRepository.java +++ b/server/src/main/java/cn/nianxx/thhotel/platform/message/repository/SourceMessageInboxRepository.java @@ -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,仅供受控原文读取服务使用。 */ diff --git a/server/src/main/java/cn/nianxx/thhotel/platform/message/service/impl/SourceMessageCaptureServiceImpl.java b/server/src/main/java/cn/nianxx/thhotel/platform/message/service/impl/SourceMessageCaptureServiceImpl.java index a5f2ba9..c29f9a6 100644 --- a/server/src/main/java/cn/nianxx/thhotel/platform/message/service/impl/SourceMessageCaptureServiceImpl.java +++ b/server/src/main/java/cn/nianxx/thhotel/platform/message/service/impl/SourceMessageCaptureServiceImpl.java @@ -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()); + } + } + /** * 校验媒体引用最小必填字段,避免保存无法追溯的附件或内嵌媒体记录。 */ diff --git a/server/src/main/resources/db/migration/V8__create_source_message_payload_duplicate.sql b/server/src/main/resources/db/migration/V8__create_source_message_payload_duplicate.sql new file mode 100644 index 0000000..07e1fa0 --- /dev/null +++ b/server/src/main/resources/db/migration/V8__create_source_message_payload_duplicate.sql @@ -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 的重复投递内容变化'; diff --git a/server/src/test/java/cn/nianxx/thhotel/platform/message/service/SourceMessageCaptureServiceImplTest.java b/server/src/test/java/cn/nianxx/thhotel/platform/message/service/SourceMessageCaptureServiceImplTest.java index e46ce4d..dce7999 100644 --- a/server/src/test/java/cn/nianxx/thhotel/platform/message/service/SourceMessageCaptureServiceImplTest.java +++ b/server/src/test/java/cn/nianxx/thhotel/platform/message/service/SourceMessageCaptureServiceImplTest.java @@ -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", + "Changed body", + "{\"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"); + } }