package task import ( "context" "strings" "time" "github.com/hibiken/asynq" "go.uber.org/zap" cardapp "github.com/break/junhong_cmp_fiber/internal/application/cardobservation" carddomain "github.com/break/junhong_cmp_fiber/internal/domain/cardobservation" "github.com/break/junhong_cmp_fiber/internal/gateway" "github.com/break/junhong_cmp_fiber/internal/infrastructure/integrationlog" "github.com/break/junhong_cmp_fiber/pkg/constants" ) // PollingRealnameHandler 负责实名轮询编排,查询结果统一交给卡观测应用服务。 type PollingRealnameHandler struct { base *PollingBase gateway *gateway.Client observation *cardapp.Service integration *integrationlog.Repository } // NewPollingRealnameHandler 创建实名检查任务处理器。 func NewPollingRealnameHandler(base *PollingBase, gw *gateway.Client, observation *cardapp.Service, integration *integrationlog.Repository) *PollingRealnameHandler { return &PollingRealnameHandler{base: base, gateway: gw, observation: observation, integration: integration} } // Handle 处理实名检查任务。 func (h *PollingRealnameHandler) Handle(ctx context.Context, task *asynq.Task) error { startedAt := time.Now() cardID, ok := parseTaskPayload(task.Payload(), h.base.logger) if !ok { return nil } if !h.base.acquireConcurrency(ctx, constants.TaskTypePollingRealname) { h.base.logger.Debug("并发已满,重新入队", zap.Uint("card_id", cardID)) return h.base.requeueCard(ctx, cardID, constants.TaskTypePollingRealname) } defer h.base.releaseConcurrency(ctx, constants.TaskTypePollingRealname) // 优先轮询互斥接缝:存在活动优先项时本轮即优先执行或被跳过;无活动项时行为与既有实现一致。 priority, skip, priorityErr := h.base.beginPollingPriority(ctx, cardID, constants.TaskTypePollingRealname) if priorityErr != nil { return h.failAndRequeue(ctx, cardID, startedAt, "查询活动优先轮询项失败", priorityErr) } if skip { h.base.logger.Debug("卡轮询优先项已被其它执行持有,本轮跳过并延后", zap.Uint("card_id", cardID), zap.String("task_type", constants.TaskTypePollingRealname)) return h.base.requeueCard(ctx, cardID, constants.TaskTypePollingRealname) } card, err := h.base.getCardWithCache(ctx, cardID) if err != nil { if isNotFound(err) { priority.FailFinal(ctx, "卡不存在") return nil } // 本地读卡失败属可恢复失败:先按真实执行累加尝试次数,再按上限收敛,避免加急事实被静默丢弃。 priority.MarkExecution(ctx) priority.FailRetryable(ctx, "获取卡信息失败") return h.failAndRequeue(ctx, cardID, startedAt, "获取卡信息失败", err) } // 执行前轮询范围校验:优先通道不得重启已停止轮询的卡。 if !h.base.priorityInPollingScope(ctx, priority, card) { return h.base.requeueCard(ctx, cardID, constants.TaskTypePollingRealname) } if card.CardCategory == constants.CardCategoryIndustry || h.gateway == nil { priority.FailFinal(ctx, "卡类型不适用实名轮询或实名查询能力未配置") return h.base.requeueCard(ctx, cardID, constants.TaskTypePollingRealname) } priority.MarkExecution(ctx) attempt := newGatewayAttempt(cardID, constants.IntegrationOperationGatewayRealname, constants.CardObservationSceneRealnamePolling) result, err := h.gateway.QueryRealnameStatus(ctx, &gateway.CardStatusReq{CardNo: card.ICCID}) if err != nil { _ = recordGatewayAttempt(ctx, h.integration, attempt, constants.IntegrationResultFailed, false, "查询实名状态失败") priority.FailRetryable(ctx, "查询实名状态失败") return h.failAndRequeue(ctx, cardID, startedAt, "查询实名状态失败", err) } if strings.TrimSpace(result.ICCID) == "" { _ = recordGatewayAttempt(ctx, h.integration, attempt, constants.IntegrationResultInvalidPayload, false, "实名查询响应缺少 ICCID") priority.FailRetryable(ctx, "实名查询响应缺少 ICCID") return h.failAndRequeue(ctx, cardID, startedAt, "实名查询响应缺少 ICCID", nil) } ctx = withPollingWorkerAuditContext(ctx, constants.TaskTypePollingRealname, "实名状态轮询任务", attempt.IntegrationID) if h.observation == nil { priority.FailFinal(ctx, "卡实名观测能力未配置") return h.failAndRequeue(ctx, cardID, startedAt, "卡实名观测能力未配置", nil) } now := time.Now() decision, err := h.observation.ApplyCardObservation(ctx, carddomain.RealnameObservation{ CardID: cardID, Verified: result.RealStatus, Metadata: carddomain.ObservationMetadata{ ObservationID: attempt.IntegrationID, Source: constants.CardObservationSourcePolling, Scene: constants.CardObservationSceneRealnamePolling, ObservedAt: now, CorrelationID: attempt.IntegrationID, }, }) if err != nil { _ = recordGatewayAttempt(ctx, h.integration, attempt, constants.IntegrationResultUnknown, false, "应用实名观测失败") priority.FailRetryable(ctx, "应用实名观测失败") return h.failAndRequeue(ctx, cardID, startedAt, "应用实名观测失败", err) } changed := decision.StatusChanged if changed { _ = recordGatewayAttempt(ctx, h.integration, attempt, constants.IntegrationResultSuccess, true, "实名状态发生变化") } if h.base.verboseLog { h.base.logger.Info("实名状态轮询详情", zap.Uint("card_id", cardID), zap.String("iccid", card.ICCID), zap.Bool("real_status", result.RealStatus), zap.Int("new_status", decision.AfterStatus), zap.Bool("changed", decision.StatusChanged), zap.Bool("reversal_pending", decision.ReversalPending)) } metric := "polling.observation.unchanged" if changed { metric = "polling.observation.persisted" } h.base.logger.Info("实名轮询观测完成", zap.Uint("card_id", cardID), zap.Bool("changed", changed), zap.String("metric", metric)) priority.Complete(ctx) h.base.updateStats(ctx, constants.TaskTypePollingRealname, true, time.Since(startedAt)) return h.base.requeueCard(ctx, cardID, constants.TaskTypePollingRealname) } func (h *PollingRealnameHandler) failAndRequeue(ctx context.Context, cardID uint, startedAt time.Time, message string, err error) error { fields := []zap.Field{zap.Uint("card_id", cardID)} if err != nil { fields = append(fields, zap.Error(err)) } h.base.logger.Warn(message, fields...) h.base.updateStats(ctx, constants.TaskTypePollingRealname, false, time.Since(startedAt)) return h.base.requeueCard(ctx, cardID, constants.TaskTypePollingRealname) }