All checks were successful
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Successful in 11m13s
628 lines
22 KiB
Go
628 lines
22 KiB
Go
package polling
|
||
|
||
import (
|
||
"context"
|
||
"time"
|
||
|
||
"github.com/bytedance/sonic"
|
||
"github.com/hibiken/asynq"
|
||
"github.com/redis/go-redis/v9"
|
||
"go.uber.org/zap"
|
||
"gorm.io/gorm"
|
||
|
||
"github.com/break/junhong_cmp_fiber/internal/model"
|
||
packagepkg "github.com/break/junhong_cmp_fiber/internal/service/package"
|
||
"github.com/break/junhong_cmp_fiber/internal/store/postgres"
|
||
"github.com/break/junhong_cmp_fiber/pkg/auditcontext"
|
||
"github.com/break/junhong_cmp_fiber/pkg/constants"
|
||
"github.com/break/junhong_cmp_fiber/pkg/errors"
|
||
)
|
||
|
||
// orphanPackageScanLimit 单轮最多恢复的真实孤儿载体数量。
|
||
const orphanPackageScanLimit = 100
|
||
|
||
// TrafficSyncer 套餐失效前流量同步接口
|
||
// 在主套餐标记为已过期之前,从 Gateway 拉取最新流量写入 DB
|
||
type TrafficSyncer interface {
|
||
RefreshCardDataByID(ctx context.Context, cardID uint) error
|
||
}
|
||
|
||
// PackageActivationHandler 套餐激活检查处理器
|
||
// 任务 19: 处理主套餐过期、加油包级联失效、待生效主套餐激活
|
||
type PackageActivationHandler struct {
|
||
db *gorm.DB
|
||
redis *redis.Client
|
||
queueClient *asynq.Client
|
||
packageUsageStore *postgres.PackageUsageStore
|
||
deviceSimBinding *postgres.DeviceSimBindingStore
|
||
activationService *packagepkg.ActivationService
|
||
stopResumeCallback packagepkg.StopResumeCallback // 停复机回调,用于套餐过期后主动停机
|
||
trafficSyncer TrafficSyncer // 套餐失效前流量同步,可选
|
||
logger *zap.Logger
|
||
}
|
||
|
||
func workerCorrelation(correlationID string, task *asynq.Task) string {
|
||
if correlationID != "" {
|
||
return correlationID
|
||
}
|
||
return task.ResultWriter().TaskID()
|
||
}
|
||
|
||
// PackageActivationPayload 套餐激活任务载荷
|
||
type PackageActivationPayload struct {
|
||
PackageUsageID uint `json:"package_usage_id"`
|
||
CarrierType string `json:"carrier_type"` // "iot_card" 或 "device"
|
||
CarrierID uint `json:"carrier_id"`
|
||
ActivationType string `json:"activation_type"` // "queue" 或 "realname"
|
||
Timestamp int64 `json:"timestamp"`
|
||
RequestID string `json:"request_id,omitempty"`
|
||
CorrelationID string `json:"correlation_id,omitempty"`
|
||
ParentEventID string `json:"parent_event_id,omitempty"`
|
||
}
|
||
|
||
// NewPackageActivationHandler 创建套餐激活检查处理器
|
||
func NewPackageActivationHandler(
|
||
db *gorm.DB,
|
||
redis *redis.Client,
|
||
queueClient *asynq.Client,
|
||
activationService *packagepkg.ActivationService,
|
||
stopResumeCallback packagepkg.StopResumeCallback,
|
||
logger *zap.Logger,
|
||
) *PackageActivationHandler {
|
||
return &PackageActivationHandler{
|
||
db: db,
|
||
redis: redis,
|
||
queueClient: queueClient,
|
||
packageUsageStore: postgres.NewPackageUsageStore(db, redis),
|
||
deviceSimBinding: postgres.NewDeviceSimBindingStore(db, redis),
|
||
activationService: activationService,
|
||
stopResumeCallback: stopResumeCallback,
|
||
logger: logger,
|
||
}
|
||
}
|
||
|
||
// SetTrafficSyncer 注入套餐失效前流量同步器(在 Start 前调用)
|
||
func (h *PackageActivationHandler) SetTrafficSyncer(syncer TrafficSyncer) {
|
||
h.trafficSyncer = syncer
|
||
}
|
||
|
||
// syncTrafficBeforeExpiry 套餐失效前从 Gateway 同步一次最新流量
|
||
// iot_card 直接同步;device 遍历所有绑定卡同步
|
||
// 失败只记录警告,不阻断失效流程
|
||
func (h *PackageActivationHandler) syncTrafficBeforeExpiry(ctx context.Context, carrierType string, carrierID uint) {
|
||
if h.trafficSyncer == nil {
|
||
return
|
||
}
|
||
|
||
if carrierType == constants.AssetTypeIotCard {
|
||
if err := h.trafficSyncer.RefreshCardDataByID(ctx, carrierID); err != nil {
|
||
h.logger.Warn("套餐失效前流量同步失败",
|
||
zap.String("carrier_type", carrierType),
|
||
zap.Uint("carrier_id", carrierID),
|
||
zap.Error(err))
|
||
}
|
||
return
|
||
}
|
||
|
||
if carrierType == constants.AssetTypeDevice {
|
||
bindings, err := h.deviceSimBinding.ListByDeviceID(ctx, carrierID)
|
||
if err != nil {
|
||
h.logger.Warn("套餐失效前查询设备绑定卡失败",
|
||
zap.Uint("device_id", carrierID),
|
||
zap.Error(err))
|
||
return
|
||
}
|
||
for _, b := range bindings {
|
||
if err := h.trafficSyncer.RefreshCardDataByID(ctx, b.IotCardID); err != nil {
|
||
h.logger.Warn("套餐失效前流量同步失败(设备绑定卡)",
|
||
zap.Uint("device_id", carrierID),
|
||
zap.Uint("card_id", b.IotCardID),
|
||
zap.Error(err))
|
||
}
|
||
}
|
||
}
|
||
}
|
||
|
||
// HandlePackageActivationCheck 任务 19.2-19.5: 处理套餐激活检查
|
||
// 每 10 秒调度一次,检查过期主套餐并激活下一个待生效主套餐
|
||
func (h *PackageActivationHandler) HandlePackageActivationCheck(ctx context.Context) error {
|
||
startTime := time.Now()
|
||
|
||
// 任务 19.2: 查询已过期的主套餐(status IN (1,2) AND expires_at <= NOW)
|
||
expiredPackages, err := h.findExpiredMainPackages(ctx)
|
||
if err != nil {
|
||
h.logger.Error("查询过期主套餐失败", zap.Error(err))
|
||
return err
|
||
}
|
||
|
||
if len(expiredPackages) > 0 {
|
||
h.logger.Info("发现过期主套餐",
|
||
zap.Int("count", len(expiredPackages)),
|
||
zap.Duration("check_duration", time.Since(startTime)))
|
||
|
||
// 处理每个过期的主套餐
|
||
for _, pkg := range expiredPackages {
|
||
if err := h.processExpiredPackage(ctx, pkg); err != nil {
|
||
h.logger.Error("处理过期套餐失败",
|
||
zap.Uint("package_usage_id", pkg.ID),
|
||
zap.Error(err))
|
||
// 继续处理下一个,不中断
|
||
continue
|
||
}
|
||
}
|
||
|
||
h.logger.Info("套餐激活检查完成",
|
||
zap.Int("processed", len(expiredPackages)),
|
||
zap.Duration("total_duration", time.Since(startTime)))
|
||
}
|
||
|
||
// 任务 6: 孤儿套餐扫描
|
||
orphanCount, err := h.findAndActivateOrphanPackages(ctx)
|
||
if err != nil {
|
||
h.logger.Error("孤儿套餐扫描失败", zap.Error(err))
|
||
// 不中断主流程
|
||
} else if orphanCount > 0 {
|
||
h.logger.Info("孤儿套餐扫描完成",
|
||
zap.Int("orphan_count", orphanCount))
|
||
}
|
||
|
||
return nil
|
||
}
|
||
|
||
// findAndActivateOrphanPackages 任务 6.1-6.2: 查找孤儿载体并触发激活
|
||
// 孤儿定义:存在待生效主套餐,但不存在生效中或已用完的占位主套餐。
|
||
func (h *PackageActivationHandler) findAndActivateOrphanPackages(ctx context.Context) (int, error) {
|
||
var orphanUsages []*model.PackageUsage
|
||
err := h.db.WithContext(ctx).Raw(`
|
||
WITH pending_queue AS (
|
||
SELECT pending.id,
|
||
pending.priority,
|
||
pending.created_at,
|
||
ROW_NUMBER() OVER (
|
||
PARTITION BY
|
||
CASE WHEN COALESCE(pending.iot_card_id, 0) > 0 THEN 'iot_card' ELSE 'device' END,
|
||
CASE WHEN COALESCE(pending.iot_card_id, 0) > 0 THEN pending.iot_card_id ELSE pending.device_id END
|
||
ORDER BY pending.priority ASC, pending.created_at ASC, pending.id ASC
|
||
) AS queue_position
|
||
FROM tb_package_usage AS pending
|
||
WHERE pending.status = ?
|
||
AND pending.master_usage_id IS NULL
|
||
AND pending.deleted_at IS NULL
|
||
AND (COALESCE(pending.iot_card_id, 0) > 0 OR COALESCE(pending.device_id, 0) > 0)
|
||
AND NOT EXISTS (
|
||
SELECT 1
|
||
FROM tb_package_usage AS occupied
|
||
WHERE occupied.status IN (?, ?)
|
||
AND occupied.master_usage_id IS NULL
|
||
AND occupied.deleted_at IS NULL
|
||
AND (
|
||
(COALESCE(pending.iot_card_id, 0) > 0 AND occupied.iot_card_id = pending.iot_card_id)
|
||
OR (COALESCE(pending.iot_card_id, 0) = 0 AND pending.device_id > 0 AND occupied.device_id = pending.device_id)
|
||
)
|
||
)
|
||
)
|
||
SELECT usage.*
|
||
FROM pending_queue AS candidate
|
||
JOIN tb_package_usage AS usage ON usage.id = candidate.id
|
||
WHERE candidate.queue_position = 1
|
||
ORDER BY candidate.priority ASC, candidate.created_at ASC, candidate.id ASC
|
||
LIMIT ?`,
|
||
constants.PackageUsageStatusPending,
|
||
constants.PackageUsageStatusActive,
|
||
constants.PackageUsageStatusDepleted,
|
||
orphanPackageScanLimit,
|
||
).Scan(&orphanUsages).Error
|
||
if err != nil {
|
||
return 0, errors.Wrap(errors.CodeDatabaseError, err, "查询孤儿套餐失败")
|
||
}
|
||
|
||
if len(orphanUsages) == 0 {
|
||
return 0, nil
|
||
}
|
||
|
||
count := 0
|
||
for _, usage := range orphanUsages {
|
||
carrierType, carrierID := h.getCarrierInfo(usage)
|
||
if err := h.enqueueActivationTask(ctx, usage.ID, carrierType, carrierID, "orphan_recovery"); err != nil {
|
||
h.logger.Warn("提交孤儿套餐激活任务失败",
|
||
zap.Uint("package_usage_id", usage.ID),
|
||
zap.String("carrier_type", carrierType),
|
||
zap.Uint("carrier_id", carrierID),
|
||
zap.Error(err))
|
||
continue
|
||
}
|
||
count++
|
||
}
|
||
|
||
return count, nil
|
||
}
|
||
|
||
// findExpiredMainPackages 任务 19.2: 查询已过期的主套餐
|
||
func (h *PackageActivationHandler) findExpiredMainPackages(ctx context.Context) ([]*model.PackageUsage, error) {
|
||
var packages []*model.PackageUsage
|
||
now := time.Now()
|
||
|
||
err := h.db.WithContext(ctx).
|
||
Where("status IN ?", []int{constants.PackageUsageStatusActive, constants.PackageUsageStatusDepleted}).
|
||
Where("expires_at <= ?", now).
|
||
Where("master_usage_id IS NULL").
|
||
Limit(1000).
|
||
Find(&packages).Error
|
||
|
||
return packages, err
|
||
}
|
||
|
||
// processExpiredPackage 处理单个过期套餐
|
||
// 流程:先同步最新流量 → 事务内标记过期和失效加油包 → 提交后投递下一套餐;
|
||
// 仅明确没有后续套餐时才触发停机检查,避免与异步激活任务竞态。
|
||
func (h *PackageActivationHandler) processExpiredPackage(ctx context.Context, pkg *model.PackageUsage) error {
|
||
carrierType, carrierID := h.getCarrierInfo(pkg)
|
||
|
||
// 失效前先从 Gateway 同步一次最新流量,确保 data_usage_mb 精确
|
||
if carrierType != "" && carrierID > 0 {
|
||
h.syncTrafficBeforeExpiry(ctx, carrierType, carrierID)
|
||
}
|
||
|
||
expired := false
|
||
err := h.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||
// 任务 19.3: 更新过期主套餐状态为 Expired (status=3)
|
||
result := tx.Model(pkg).
|
||
Where("status IN ? AND expires_at <= ?", []int{constants.PackageUsageStatusActive, constants.PackageUsageStatusDepleted}, time.Now()).
|
||
Update("status", constants.PackageUsageStatusExpired)
|
||
if result.Error != nil {
|
||
return result.Error
|
||
}
|
||
if result.RowsAffected == 0 {
|
||
return nil
|
||
}
|
||
expired = true
|
||
|
||
expiresAt := time.Now()
|
||
if pkg.ExpiresAt != nil {
|
||
expiresAt = *pkg.ExpiresAt
|
||
}
|
||
h.logger.Info("主套餐已过期",
|
||
zap.Uint("package_usage_id", pkg.ID),
|
||
zap.Time("expires_at", expiresAt))
|
||
|
||
// 任务 19.4: 加油包级联失效
|
||
addons, err := h.invalidateAddons(ctx, tx, pkg.ID)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
if h.activationService == nil {
|
||
return errors.New(errors.CodeInternalError, "套餐激活服务未注入")
|
||
}
|
||
if err := h.activationService.AppendExpirationAudit(ctx, tx, pkg, addons); err != nil {
|
||
return err
|
||
}
|
||
|
||
return nil
|
||
})
|
||
|
||
if err != nil {
|
||
if h.activationService != nil {
|
||
h.activationService.RecordUsageFailure(ctx, constants.AuditActionPackageUsageExpired, "套餐权益到期处理失败", pkg, err)
|
||
}
|
||
return err
|
||
}
|
||
if !expired {
|
||
return nil
|
||
}
|
||
|
||
// 事务提交后再投递,确保消费者只能读取到旧套餐已经过期的状态。
|
||
if carrierType != "" && carrierID > 0 {
|
||
activationEnqueued, err := h.activateNextPackage(ctx, h.db.WithContext(ctx), carrierType, carrierID)
|
||
if err != nil {
|
||
// 激活结果未知时不能将其当作无套餐;孤儿扫描和套餐轮询会继续兜底。
|
||
return err
|
||
}
|
||
if activationEnqueued {
|
||
return nil
|
||
}
|
||
h.triggerStopAfterExpiry(ctx, carrierType, carrierID)
|
||
}
|
||
|
||
return nil
|
||
}
|
||
|
||
// invalidateAddons 任务 19.4: 加油包级联失效
|
||
func (h *PackageActivationHandler) invalidateAddons(ctx context.Context, tx *gorm.DB, masterUsageID uint) ([]*model.PackageUsage, error) {
|
||
// 查询主套餐下的所有加油包(status IN (0,1,2) 的加油包)
|
||
var addons []*model.PackageUsage
|
||
if err := tx.WithContext(ctx).
|
||
Where("master_usage_id = ?", masterUsageID).
|
||
Where("status IN ?", []int{
|
||
constants.PackageUsageStatusPending,
|
||
constants.PackageUsageStatusActive,
|
||
constants.PackageUsageStatusDepleted,
|
||
}).Find(&addons).Error; err != nil {
|
||
return nil, err
|
||
}
|
||
if len(addons) == 0 {
|
||
return nil, nil
|
||
}
|
||
ids := make([]uint, 0, len(addons))
|
||
for _, addon := range addons {
|
||
ids = append(ids, addon.ID)
|
||
}
|
||
result := tx.Model(&model.PackageUsage{}).
|
||
Where("id IN ? AND status IN ?", ids, []int{constants.PackageUsageStatusPending, constants.PackageUsageStatusActive, constants.PackageUsageStatusDepleted}).
|
||
Update("status", constants.PackageUsageStatusInvalidated)
|
||
|
||
if result.Error != nil {
|
||
return nil, result.Error
|
||
}
|
||
|
||
if result.RowsAffected > 0 {
|
||
h.logger.Info("加油包已级联失效",
|
||
zap.Uint("master_usage_id", masterUsageID),
|
||
zap.Int64("invalidated_count", result.RowsAffected))
|
||
}
|
||
|
||
return addons, nil
|
||
}
|
||
|
||
// getCarrierInfo 获取载体信息
|
||
func (h *PackageActivationHandler) getCarrierInfo(pkg *model.PackageUsage) (string, uint) {
|
||
if pkg.IotCardID > 0 {
|
||
return "iot_card", pkg.IotCardID
|
||
}
|
||
if pkg.DeviceID > 0 {
|
||
return "device", pkg.DeviceID
|
||
}
|
||
return "", 0
|
||
}
|
||
|
||
// activateNextPackage 提交下一个待生效主套餐的异步激活任务。
|
||
// 返回 true 表示已找到并成功提交后续套餐;false 表示不存在后续套餐。
|
||
func (h *PackageActivationHandler) activateNextPackage(ctx context.Context, tx *gorm.DB, carrierType string, carrierID uint) (bool, error) {
|
||
var nextPkg model.PackageUsage
|
||
query := tx.Where("status = ?", constants.PackageUsageStatusPending).
|
||
Where("master_usage_id IS NULL").
|
||
Order("priority ASC, created_at ASC, id ASC").
|
||
Limit(1)
|
||
|
||
if carrierType == constants.AssetTypeIotCard {
|
||
query = query.Where("iot_card_id = ?", carrierID)
|
||
} else if carrierType == constants.AssetTypeDevice {
|
||
query = query.Where("device_id = ?", carrierID)
|
||
}
|
||
|
||
if err := query.First(&nextPkg).Error; err != nil {
|
||
if err == gorm.ErrRecordNotFound {
|
||
return false, nil
|
||
}
|
||
return false, err
|
||
}
|
||
|
||
if err := h.enqueueActivationTask(ctx, nextPkg.ID, carrierType, carrierID, "queue"); err != nil {
|
||
return false, err
|
||
}
|
||
return true, nil
|
||
}
|
||
|
||
// triggerStopAfterExpiry 在明确无后续套餐时异步触发停机检查。
|
||
// 后续套餐存在时,停机重评估由激活任务在完成后顺序执行。
|
||
func (h *PackageActivationHandler) triggerStopAfterExpiry(ctx context.Context, carrierType string, carrierID uint) {
|
||
if h.stopResumeCallback == nil {
|
||
return
|
||
}
|
||
detachedCtx := context.WithoutCancel(ctx)
|
||
|
||
if carrierType == "iot_card" {
|
||
go func() {
|
||
if err := h.stopResumeCallback.CheckAndStopCard(detachedCtx, carrierID); err != nil {
|
||
h.logger.Error("套餐过期后停机失败",
|
||
zap.Uint("card_id", carrierID),
|
||
zap.Error(err))
|
||
}
|
||
}()
|
||
return
|
||
}
|
||
|
||
if carrierType == "device" {
|
||
bindings, err := h.deviceSimBinding.ListByDeviceID(ctx, carrierID)
|
||
if err != nil {
|
||
h.logger.Error("查询设备绑定卡失败",
|
||
zap.Uint("device_id", carrierID),
|
||
zap.Error(err))
|
||
return
|
||
}
|
||
for _, b := range bindings {
|
||
cardID := b.IotCardID
|
||
go func(cID uint) {
|
||
if err := h.stopResumeCallback.CheckAndStopCard(detachedCtx, cID); err != nil {
|
||
h.logger.Error("套餐过期后停机失败",
|
||
zap.Uint("card_id", cID),
|
||
zap.Error(err))
|
||
}
|
||
}(cardID)
|
||
}
|
||
}
|
||
}
|
||
|
||
// reconcileCarrierAfterActivation 在套餐激活任务完成后按最新事实重评估停复机。
|
||
// 它必须在激活调用返回后执行,不能与激活任务并发读取过期权益快照。
|
||
func (h *PackageActivationHandler) reconcileCarrierAfterActivation(ctx context.Context, packageUsageID uint, carrierType string, carrierID uint) error {
|
||
if h.stopResumeCallback == nil {
|
||
h.logger.Warn("套餐激活后停复机回调未注入,跳过重评估",
|
||
zap.Uint("package_usage_id", packageUsageID),
|
||
zap.String("carrier_type", carrierType),
|
||
zap.Uint("carrier_id", carrierID))
|
||
return nil
|
||
}
|
||
if carrierID == 0 {
|
||
return errors.New(errors.CodeInvalidParam, "套餐使用记录缺少有效载体")
|
||
}
|
||
|
||
if carrierType == constants.AssetTypeIotCard {
|
||
if err := h.stopResumeCallback.CheckAndStopCard(ctx, carrierID); err != nil {
|
||
return err
|
||
}
|
||
return nil
|
||
}
|
||
if carrierType != constants.AssetTypeDevice {
|
||
return errors.New(errors.CodeInvalidParam, "套餐使用记录载体类型无效")
|
||
}
|
||
|
||
bindings, err := h.deviceSimBinding.ListByDeviceID(ctx, carrierID)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
for _, binding := range bindings {
|
||
if err := h.stopResumeCallback.CheckAndStopCard(ctx, binding.IotCardID); err != nil {
|
||
return err
|
||
}
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// enqueueActivationTask 提交套餐激活任务到 Asynq
|
||
func (h *PackageActivationHandler) enqueueActivationTask(ctx context.Context, packageUsageID uint, carrierType string, carrierID uint, activationType string) error {
|
||
linkage := auditcontext.From(ctx)
|
||
payload := PackageActivationPayload{
|
||
PackageUsageID: packageUsageID,
|
||
CarrierType: carrierType,
|
||
CarrierID: carrierID,
|
||
ActivationType: activationType,
|
||
Timestamp: time.Now().Unix(),
|
||
RequestID: linkage.RequestID,
|
||
CorrelationID: linkage.CorrelationID,
|
||
ParentEventID: linkage.ParentEventID,
|
||
}
|
||
|
||
payloadBytes, err := sonic.Marshal(payload)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
|
||
task := asynq.NewTask(constants.TaskTypePackageQueueActivation, payloadBytes,
|
||
asynq.MaxRetry(3),
|
||
asynq.Timeout(60*time.Second),
|
||
asynq.Queue(constants.QueueForTaskType(constants.TaskTypePackageQueueActivation)),
|
||
)
|
||
|
||
_, err = h.queueClient.Enqueue(task)
|
||
if err != nil {
|
||
h.logger.Error("提交套餐激活任务失败",
|
||
zap.Uint("package_usage_id", packageUsageID),
|
||
zap.Error(err))
|
||
return err
|
||
}
|
||
|
||
h.logger.Info("已提交套餐激活任务",
|
||
zap.Uint("package_usage_id", packageUsageID),
|
||
zap.String("activation_type", activationType))
|
||
|
||
return nil
|
||
}
|
||
|
||
// HandlePackageQueueActivation 处理套餐排队激活任务(Asynq Handler)
|
||
// 任务 23: 由 Asynq 调用,执行实际的套餐激活逻辑
|
||
func (h *PackageActivationHandler) HandlePackageQueueActivation(ctx context.Context, t *asynq.Task) error {
|
||
var payload PackageActivationPayload
|
||
if err := sonic.Unmarshal(t.Payload(), &payload); err != nil {
|
||
h.logger.Error("解析套餐激活任务载荷失败", zap.Error(err))
|
||
return nil // 不重试
|
||
}
|
||
ctx = auditcontext.With(ctx, auditcontext.Context{
|
||
ActorKind: constants.AuditActorSystemTask, ActorID: constants.TaskTypePackageQueueActivation,
|
||
ActorName: "套餐排队激活任务", Source: constants.AuditSourceWorker,
|
||
RequestID: payload.RequestID, CorrelationID: workerCorrelation(payload.CorrelationID, t), ParentEventID: payload.ParentEventID,
|
||
})
|
||
|
||
h.logger.Info("开始执行套餐激活",
|
||
zap.Uint("package_usage_id", payload.PackageUsageID),
|
||
zap.String("activation_type", payload.ActivationType))
|
||
|
||
// 查询套餐使用记录
|
||
var pkg model.PackageUsage
|
||
if err := h.db.First(&pkg, payload.PackageUsageID).Error; err != nil {
|
||
if err == gorm.ErrRecordNotFound {
|
||
h.logger.Warn("套餐使用记录不存在", zap.Uint("package_usage_id", payload.PackageUsageID))
|
||
return nil
|
||
}
|
||
return err
|
||
}
|
||
|
||
// 调用 ActivationService 执行激活。即使套餐已由其他任务激活,仍须重新评估停复机。
|
||
if h.activationService != nil {
|
||
if err := h.activationService.ActivateSpecificPackage(ctx, payload.PackageUsageID); err != nil {
|
||
h.logger.Error("套餐激活失败",
|
||
zap.Uint("package_usage_id", payload.PackageUsageID),
|
||
zap.Error(err))
|
||
return err
|
||
}
|
||
} else {
|
||
// ActivationService 未注入,无法安全激活(缺少 expires_at 计算)
|
||
h.logger.Error("激活服务未注入,无法执行套餐激活",
|
||
zap.Uint("package_usage_id", payload.PackageUsageID))
|
||
return errors.New(errors.CodeInternalError, "激活服务未注入,无法执行套餐激活")
|
||
}
|
||
|
||
carrierType, carrierID := h.getCarrierInfo(&pkg)
|
||
if err := h.reconcileCarrierAfterActivation(ctx, payload.PackageUsageID, carrierType, carrierID); err != nil {
|
||
h.logger.Error("套餐激活后停复机重评估失败",
|
||
zap.Uint("package_usage_id", payload.PackageUsageID),
|
||
zap.String("carrier_type", carrierType),
|
||
zap.Uint("carrier_id", carrierID),
|
||
zap.Error(err))
|
||
return err
|
||
}
|
||
|
||
h.logger.Info("套餐激活及停复机重评估完成",
|
||
zap.Uint("package_usage_id", payload.PackageUsageID),
|
||
zap.String("carrier_type", carrierType),
|
||
zap.Uint("carrier_id", carrierID))
|
||
|
||
return nil
|
||
}
|
||
|
||
// HandlePackageFirstActivation 处理首次实名激活任务(Asynq Handler)
|
||
// 由 polling_realname_handler 在检测到首次实名时触发。
|
||
// payload 中仅需 carrier_type + carrier_id,ActivateByRealname 内部自行查找
|
||
// pending_realname_activation=true 的套餐,无需外部指定 PackageUsageID。
|
||
func (h *PackageActivationHandler) HandlePackageFirstActivation(ctx context.Context, t *asynq.Task) error {
|
||
var payload PackageActivationPayload
|
||
if err := sonic.Unmarshal(t.Payload(), &payload); err != nil {
|
||
h.logger.Error("解析首次实名激活任务载荷失败", zap.Error(err))
|
||
return nil
|
||
}
|
||
ctx = auditcontext.With(ctx, auditcontext.Context{
|
||
ActorKind: constants.AuditActorSystemTask, ActorID: constants.TaskTypePackageFirstActivation,
|
||
ActorName: "套餐首次实名激活任务", Source: constants.AuditSourceWorker,
|
||
RequestID: payload.RequestID, CorrelationID: workerCorrelation(payload.CorrelationID, t), ParentEventID: payload.ParentEventID,
|
||
})
|
||
|
||
if payload.CarrierType == "" || payload.CarrierID == 0 {
|
||
h.logger.Error("首次实名激活任务 carrier 信息缺失",
|
||
zap.String("carrier_type", payload.CarrierType),
|
||
zap.Uint("carrier_id", payload.CarrierID))
|
||
return nil
|
||
}
|
||
|
||
h.logger.Info("开始执行首次实名激活",
|
||
zap.String("carrier_type", payload.CarrierType),
|
||
zap.Uint("carrier_id", payload.CarrierID))
|
||
|
||
if h.activationService == nil {
|
||
h.logger.Warn("ActivationService 未注入,跳过首次实名激活")
|
||
return nil
|
||
}
|
||
|
||
if err := h.activationService.ActivateByRealname(ctx, payload.CarrierType, payload.CarrierID); err != nil {
|
||
h.logger.Error("首次实名激活失败",
|
||
zap.String("carrier_type", payload.CarrierType),
|
||
zap.Uint("carrier_id", payload.CarrierID),
|
||
zap.Error(err))
|
||
return err
|
||
}
|
||
|
||
h.logger.Info("首次实名激活成功",
|
||
zap.String("carrier_type", payload.CarrierType),
|
||
zap.Uint("carrier_id", payload.CarrierID))
|
||
|
||
return nil
|
||
}
|