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/internal/model" "github.com/break/junhong_cmp_fiber/pkg/constants" ) // PollingCardStatusHandler 负责网络状态轮询编排,查询结果统一交给卡观测应用服务。 type PollingCardStatusHandler struct { base *PollingBase gateway *gateway.Client observation *cardapp.Service integration *integrationlog.Repository } // shouldStopPollingForRisk 保留旧测试与调用方的纯判断兼容入口,规则权威在卡观测领域层。 func shouldStopPollingForRisk(card *model.IotCard, extend string) bool { if card == nil { return false } return carddomain.ShouldStopPollingForRisk(card.IsStandalone, strings.TrimSpace(extend)) } // NewPollingCardStatusHandler 创建卡状态轮询任务处理器。 func NewPollingCardStatusHandler(base *PollingBase, gw *gateway.Client, observation *cardapp.Service, integration *integrationlog.Repository) *PollingCardStatusHandler { return &PollingCardStatusHandler{base: base, gateway: gw, observation: observation, integration: integration} } // Handle 处理卡状态轮询任务。 func (h *PollingCardStatusHandler) 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.TaskTypePollingCardStatus) { h.base.logger.Debug("并发已满,重新入队", zap.Uint("card_id", cardID)) return h.base.requeueCard(ctx, cardID, constants.TaskTypePollingCardStatus) } defer h.base.releaseConcurrency(ctx, constants.TaskTypePollingCardStatus) card, err := h.base.getCardWithCache(ctx, cardID) if err != nil { if isNotFound(err) { return nil } return h.failAndRequeue(ctx, cardID, startedAt, "获取卡信息失败", err) } if h.gateway == nil { return h.base.requeueCard(ctx, cardID, constants.TaskTypePollingCardStatus) } attempt := newGatewayAttempt(cardID, constants.IntegrationOperationGatewayNetwork, constants.CardObservationSceneNetworkPolling) result, err := h.gateway.QueryCardStatus(ctx, &gateway.CardStatusReq{CardNo: card.ICCID}) if err != nil { _ = recordGatewayAttempt(ctx, h.integration, attempt, constants.IntegrationResultFailed, false, "查询卡状态失败") return h.failAndRequeue(ctx, cardID, startedAt, "查询卡状态失败", err) } if strings.TrimSpace(result.ICCID) == "" { _ = recordGatewayAttempt(ctx, h.integration, attempt, constants.IntegrationResultInvalidPayload, false, "卡状态查询响应缺少 ICCID") return h.failAndRequeue(ctx, cardID, startedAt, "卡状态查询响应缺少 ICCID", nil) } ctx = withPollingWorkerAuditContext(ctx, constants.TaskTypePollingCardStatus, "卡网络状态轮询任务", attempt.IntegrationID) if h.observation == nil { return h.failAndRequeue(ctx, cardID, startedAt, "卡网络观测能力未配置", nil) } now := time.Now() decision, err := h.observation.ApplyNetworkObservation(ctx, carddomain.NetworkObservation{ CardID: cardID, GatewayStatus: result.CardStatus, GatewayExtend: result.Extend, GatewayIMEI: result.IMEI, Metadata: carddomain.ObservationMetadata{ ObservationID: attempt.IntegrationID, Source: constants.CardObservationSourcePolling, Scene: constants.CardObservationSceneNetworkPolling, ObservedAt: now, CorrelationID: attempt.IntegrationID, }, }) if err != nil { _ = recordGatewayAttempt(ctx, h.integration, attempt, constants.IntegrationResultUnknown, false, "应用网络状态观测失败") return h.failAndRequeue(ctx, cardID, startedAt, "应用网络状态观测失败", err) } changed := decision.StatusChanged || decision.StopReasonChanged || decision.StopPolling if changed { _ = recordGatewayAttempt(ctx, h.integration, attempt, constants.IntegrationResultSuccess, true, "网络状态发生变化") } if h.base.verboseLog || !decision.StatusKnown { h.base.logger.Info("卡状态轮询详情", zap.Uint("card_id", cardID), zap.String("iccid", card.ICCID), zap.String("card_status", result.CardStatus), zap.String("extend", strings.TrimSpace(result.Extend)), zap.Int("new_network_status", decision.AfterStatus), zap.Bool("status_known", decision.StatusKnown), zap.Bool("changed", decision.StatusChanged), zap.Bool("stop_polling", decision.StopPolling)) } 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)) h.base.updateStats(ctx, constants.TaskTypePollingCardStatus, true, time.Since(startedAt)) if decision.StopPolling { h.base.logger.Info("独立卡命中风险状态,已关闭轮询", zap.Uint("card_id", cardID), zap.String("gateway_extend", decision.GatewayExtend)) return nil } // 卡状态任务仍正常但套餐任务丢失时,仅补入缺失项;不改写已有套餐任务的执行时间。 h.base.invalidateCardCache(ctx, cardID) if err := h.base.ensureMissingTask(ctx, cardID, constants.TaskTypePollingPackage); err != nil { h.base.logger.Warn("卡状态轮询补齐套餐任务失败", zap.Uint("card_id", cardID), zap.Error(err)) } return h.base.requeueCard(ctx, cardID, constants.TaskTypePollingCardStatus) } func (h *PollingCardStatusHandler) 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.TaskTypePollingCardStatus, false, time.Since(startedAt)) return h.base.requeueCard(ctx, cardID, constants.TaskTypePollingCardStatus) }