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/store/postgres" "github.com/break/junhong_cmp_fiber/pkg/constants" ) // PollingCarddataHandler 负责流量轮询编排,查询结果统一交给卡观测应用服务。 type PollingCarddataHandler struct { base *PollingBase gateway *gateway.Client carrier *postgres.CarrierStore observation *cardapp.Service integration *integrationlog.Repository } // NewPollingCarddataHandler 创建流量检查任务处理器。 func NewPollingCarddataHandler(base *PollingBase, gw *gateway.Client, carrier *postgres.CarrierStore, observation *cardapp.Service, integration *integrationlog.Repository) *PollingCarddataHandler { return &PollingCarddataHandler{base: base, gateway: gw, carrier: carrier, observation: observation, integration: integration} } // Handle 处理流量检查任务。 func (h *PollingCarddataHandler) 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.TaskTypePollingCarddata) { h.base.logger.Debug("并发已满,重新入队", zap.Uint("card_id", cardID)) return h.base.requeueCard(ctx, cardID, constants.TaskTypePollingCarddata) } defer h.base.releaseConcurrency(ctx, constants.TaskTypePollingCarddata) lockToken, locked, err := h.base.acquireCardTrafficSyncLock(ctx, cardID) if err != nil || !locked { h.base.logger.Warn("卡流量同步锁不可用,延迟重入队", zap.Uint("card_id", cardID), zap.Error(err)) return h.base.queueMgr.Requeue(ctx, cardID, constants.TaskTypePollingCarddata, time.Now().Add(5*time.Second)) } defer h.base.releaseCardTrafficSyncLock(ctx, cardID, lockToken) h.base.invalidateCardCache(ctx, cardID) card, err := h.base.iotCardStore.GetByID(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.TaskTypePollingCarddata) } attemptStartedAt := time.Now() attempt, err := startGatewayAttempt(ctx, h.integration, cardID, constants.IntegrationOperationGatewayTraffic, constants.CardObservationSceneTrafficPolling) if err != nil { return h.failAndRequeue(ctx, cardID, startedAt, "建立 Gateway 流量 Integration Log 失败", err) } result, err := h.gateway.QueryFlow(ctx, &gateway.FlowQueryReq{CardNo: card.ICCID}) if err != nil { _ = completeGatewayAttempt(ctx, h.integration, attempt, false, attemptStartedAt) return h.failAndRequeue(ctx, cardID, startedAt, "查询流量失败", err) } if strings.TrimSpace(result.ICCID) == "" { _ = completeGatewayAttempt(ctx, h.integration, attempt, false, attemptStartedAt) return h.failAndRequeue(ctx, cardID, startedAt, "流量查询响应缺少 ICCID", nil) } if logErr := completeGatewayAttempt(ctx, h.integration, attempt, true, attemptStartedAt); logErr != nil { return h.failAndRequeue(ctx, cardID, startedAt, "完成 Gateway 流量 Integration Log 失败", logErr) } if h.observation == nil || h.carrier == nil { return h.failAndRequeue(ctx, cardID, startedAt, "卡流量观测能力未配置", nil) } now := time.Now() decision, err := h.observation.ApplyTrafficObservation(ctx, carddomain.TrafficObservation{ CardID: cardID, GatewayReadingMB: float64(result.Used), ResetDay: h.carrier.GetDataResetDay(ctx, card.CarrierID), Metadata: carddomain.ObservationMetadata{ ObservationID: attempt.IntegrationID, Source: constants.CardObservationSourcePolling, Scene: constants.CardObservationSceneTrafficPolling, ObservedAt: now, CorrelationID: attempt.IntegrationID, }, }) if err != nil { return h.failAndRequeue(ctx, cardID, startedAt, "应用流量观测失败", err) } if h.base.verboseLog { h.base.logger.Info("流量轮询详情", zap.Uint("card_id", cardID), zap.String("iccid", card.ICCID), zap.Float64("gateway_flow_mb", float64(result.Used)), zap.Float64("increment_mb", decision.IncrementMB), zap.Bool("is_cross_month", decision.CrossMonth), zap.Bool("reading_accepted", decision.ReadingAccepted)) } h.base.updateStats(ctx, constants.TaskTypePollingCarddata, true, time.Since(startedAt)) return h.base.requeueCard(ctx, cardID, constants.TaskTypePollingCarddata) } func (h *PollingCarddataHandler) 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.TaskTypePollingCarddata, false, time.Since(startedAt)) return h.base.requeueCard(ctx, cardID, constants.TaskTypePollingCarddata) }