diff --git a/docs/ur94-card-state-events-callbacks/功能总结.md b/docs/ur94-card-state-events-callbacks/功能总结.md index 4e0f422..26555d7 100644 --- a/docs/ur94-card-state-events-callbacks/功能总结.md +++ b/docs/ur94-card-state-events-callbacks/功能总结.md @@ -10,7 +10,9 @@ - 合法 ICCID 只按 19/20 位对应列精确查询;不复制旧代码的 20 位截 19 位,不跨列降级,也不任取多匹配卡。 - 使用 `ICCID + dateChanged` 的安全摘要作为语义幂等键;解析失败使用正文摘要,重复、冲突和中断恢复沿用 Integration Log 租约规则。 +- 内部失败终态收到相同正文重推时,原子恢复为 `pending` 并递增尝试次数;成功等其他终态仍保持幂等跳过。 - 唯一命中后调用公共 `ApplyCardObservation(verified=true)`,实名事实变化时由公共用例写 Outbox,并尽力提前完成同卡实名观测序列。 +- 卡实名 Outbox 事件 ID 使用受长度上限约束的稳定摘要,避免上游观测 ID 过长导致业务事务回滚。 - 不调用旧 `inner_callback`、第三方推送、旧平台登录、`ModifyDate` 或 Gateway 二次确认;`dateChanged` 只作为上游变更时间和幂等语义留痕,不直接改写本地业务时间。 - 无论开关状态、报文结果或内部处理结果,均返回入口时间对应的 HTTP 200 固定 JSON 应答,避免运营商不可控重推。 diff --git a/internal/application/cardobservation/apply.go b/internal/application/cardobservation/apply.go index 45a0ca9..aef7d0f 100644 --- a/internal/application/cardobservation/apply.go +++ b/internal/application/cardobservation/apply.go @@ -3,6 +3,8 @@ package cardobservation import ( "context" + "crypto/sha256" + "encoding/hex" "strconv" "time" @@ -98,7 +100,7 @@ func (s *Service) ApplyCardObservation(ctx context.Context, observation domain.R return errors.New(errors.CodeConflict, "卡实名状态已被其他请求更新") } if decision.StatusChanged { - eventID := "card-realname:" + strconv.FormatUint(uint64(card.ID), 10) + ":" + observation.Metadata.ObservationID + ":changed" + eventID := realnameChangedEventID(card.ID, observation.Metadata.ObservationID) if err := s.eventWriter.AppendRealname(ctx, tx, RealnameChangedEvent{ EventID: eventID, CardID: card.ID, BeforeStatus: card.RealNameStatus, AfterStatus: decision.AfterStatus, FirstVerified: decision.FirstVerified, ObservedAt: observation.Metadata.ObservedAt, @@ -118,3 +120,10 @@ func (s *Service) ApplyCardObservation(ctx context.Context, observation domain.R } return decision, nil } + +func realnameChangedEventID(cardID uint, observationID string) string { + prefix := "card-realname:" + digest := sha256.Sum256([]byte(strconv.FormatUint(uint64(cardID), 10) + ":" + observationID + ":changed")) + // 使用稳定摘要保留幂等语义,同时严格服从公共 Outbox 的字段长度上限。 + return prefix + hex.EncodeToString(digest[:])[:constants.OutboxEventIDMaxLength-len(prefix)] +} diff --git a/internal/handler/callback/cucc_realname.go b/internal/handler/callback/cucc_realname.go index 77817de..a3fec31 100644 --- a/internal/handler/callback/cucc_realname.go +++ b/internal/handler/callback/cucc_realname.go @@ -78,7 +78,7 @@ func (h *CUCCRealnameHandler) process(ctx context.Context, body []byte, contentT return err } if !created { - claimed, claimErr := h.claimPending(ctx, log) + claimed, claimErr := h.claimRetryable(ctx, log) if claimErr != nil || !claimed { return claimErr } @@ -129,7 +129,7 @@ func (h *CUCCRealnameHandler) recordConflict(ctx context.Context, body []byte, c return err } if !created { - claimed, claimErr := h.claimPending(ctx, log) + claimed, claimErr := h.claimRetryable(ctx, log) if claimErr != nil || !claimed { return claimErr } @@ -137,11 +137,18 @@ func (h *CUCCRealnameHandler) recordConflict(ctx context.Context, body []byte, c return h.complete(ctx, log.IntegrationID, constants.IntegrationResultConflict, false, "同一实名语义对应不同载荷") } -func (h *CUCCRealnameHandler) claimPending(ctx context.Context, log *model.IntegrationLog) (bool, error) { - if log == nil || log.Result != constants.IntegrationResultPending { +func (h *CUCCRealnameHandler) claimRetryable(ctx context.Context, log *model.IntegrationLog) (bool, error) { + if log == nil { + return false, nil + } + switch log.Result { + case constants.IntegrationResultPending: + return h.integration.ClaimExpiredInboundPending(ctx, log.IntegrationID, constants.IntegrationInboundProcessingLease) + case constants.IntegrationResultFailed: + return h.integration.ClaimFailedInbound(ctx, log.IntegrationID) + default: return false, nil } - return h.integration.ClaimExpiredInboundPending(ctx, log.IntegrationID, constants.IntegrationInboundProcessingLease) } func (h *CUCCRealnameHandler) complete(ctx context.Context, integrationID, result string, changed bool, reason string) error { diff --git a/internal/infrastructure/integrationlog/repository.go b/internal/infrastructure/integrationlog/repository.go index c6ee1b1..8ffbdc3 100644 --- a/internal/infrastructure/integrationlog/repository.go +++ b/internal/infrastructure/integrationlog/repository.go @@ -256,6 +256,29 @@ func (r *Repository) ClaimExpiredInboundPending(ctx context.Context, integration return result.RowsAffected == 1, nil } +// ClaimFailedInbound 原子认领失败的入站回调并恢复为待处理状态。 +func (r *Repository) ClaimFailedInbound(ctx context.Context, integrationID string) (bool, error) { + if r == nil || r.db == nil { + return false, pkgerrors.New(pkgerrors.CodeInvalidStatus, "Integration Log 数据库未配置") + } + if strings.TrimSpace(integrationID) == "" { + return false, pkgerrors.New(pkgerrors.CodeInvalidParam, "入站 Integration Log 恢复参数无效") + } + now := r.now().UTC() + result := r.db.WithContext(ctx).Model(&model.IntegrationLog{}). + Where("integration_id = ? AND direction = ? AND result = ?", integrationID, constants.IntegrationDirectionInbound, constants.IntegrationResultFailed). + Updates(map[string]any{ + "result": constants.IntegrationResultPending, "started_at": now, + "attempt": gorm.Expr("attempt + 1"), "updated_at": now, + "http_status": nil, "provider_code": nil, "provider_message": nil, + "response_summary": nil, "duration_ms": 0, "state_changed": false, "recovery_strategy": nil, + }) + if result.Error != nil { + return false, pkgerrors.Wrap(pkgerrors.CodeDatabaseError, result.Error, "认领失败入站 Integration Log 失败") + } + return result.RowsAffected == 1, nil +} + func validateAttempt(input Attempt) error { if input.Provider == "" || input.Operation == "" { return pkgerrors.New(pkgerrors.CodeInvalidParam, "Integration Log 提供方和操作不能为空") diff --git a/pkg/constants/outbox.go b/pkg/constants/outbox.go index 9ffd9a0..a2559e7 100644 --- a/pkg/constants/outbox.go +++ b/pkg/constants/outbox.go @@ -14,6 +14,8 @@ const ( ) const ( + // OutboxEventIDMaxLength 是公共 Outbox 稳定事件ID的数据库长度上限。 + OutboxEventIDMaxLength = 64 // OutboxDefaultMaxRetries 是公共 Outbox 默认最大重试次数。 OutboxDefaultMaxRetries = 10 // OutboxDefaultLeaseDuration 是 Relay 默认租约时长。