避免联通实名回调因内部长度限制永久失效
All checks were successful
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Successful in 8m9s
All checks were successful
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Successful in 8m9s
将实名事件身份收敛为受 Outbox 上限约束的稳定摘要,并允许相同 CUCC 回调原子认领失败终态后重试。 Constraint: 运营商回调固定返回成功,公共 Outbox event_id 上限为 64 字符 Rejected: 扩大数据库字段 | 根因是不受控事件 ID,且无法解决旧 failed 回调被吞 Confidence: high Scope-risk: moderate Directive: 新增稳定事件 ID 时必须服从持久化长度上限;失败回调恢复必须使用条件更新原子认领 Tested: gofmt;git diff --cached --check Not-tested: 按用户要求未运行自动化测试、构建或测试环境重放
This commit is contained in:
@@ -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)]
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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 提供方和操作不能为空")
|
||||
|
||||
Reference in New Issue
Block a user