Files
junhong_cmp_fiber/internal/task/polling_cardstatus_handler.go
2026-07-24 16:07:18 +08:00

116 lines
5.2 KiB
Go

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)
}
attemptStartedAt := time.Now()
attempt, err := startGatewayAttempt(ctx, h.integration, cardID, constants.IntegrationOperationGatewayNetwork, constants.CardObservationSceneNetworkPolling)
if err != nil {
return h.failAndRequeue(ctx, cardID, startedAt, "建立 Gateway 网络 Integration Log 失败", err)
}
result, err := h.gateway.QueryCardStatus(ctx, &gateway.CardStatusReq{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 {
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 {
return h.failAndRequeue(ctx, cardID, startedAt, "应用网络状态观测失败", err)
}
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))
}
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
}
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)
}