合并七月迭代分支

This commit is contained in:
2026-08-18 14:53:29 +08:00
2237 changed files with 91363 additions and 290515 deletions

View File

@@ -13,6 +13,7 @@ import (
"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"
)
@@ -40,6 +41,13 @@ type PackageActivationHandler struct {
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"`
@@ -47,6 +55,9 @@ type PackageActivationPayload struct {
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 创建套餐激活检查处理器
@@ -251,11 +262,19 @@ func (h *PackageActivationHandler) processExpiredPackage(ctx context.Context, pk
h.syncTrafficBeforeExpiry(ctx, carrierType, carrierID)
}
expired := false
err := h.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
// 任务 19.3: 更新过期主套餐状态为 Expired (status=3)
if err := tx.Model(pkg).Update("status", constants.PackageUsageStatusExpired).Error; err != nil {
return err
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 {
@@ -266,7 +285,14 @@ func (h *PackageActivationHandler) processExpiredPackage(ctx context.Context, pk
zap.Time("expires_at", expiresAt))
// 任务 19.4: 加油包级联失效
if err := h.invalidateAddons(ctx, tx, pkg.ID); err != nil {
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
}
@@ -274,8 +300,14 @@ func (h *PackageActivationHandler) processExpiredPackage(ctx context.Context, pk
})
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 {
@@ -288,19 +320,31 @@ func (h *PackageActivationHandler) processExpiredPackage(ctx context.Context, pk
}
// invalidateAddons 任务 19.4: 加油包级联失效
func (h *PackageActivationHandler) invalidateAddons(ctx context.Context, tx *gorm.DB, masterUsageID uint) error {
func (h *PackageActivationHandler) invalidateAddons(ctx context.Context, tx *gorm.DB, masterUsageID uint) ([]*model.PackageUsage, error) {
// 查询主套餐下的所有加油包status IN (0,1,2) 的加油包)
result := tx.Model(&model.PackageUsage{}).
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 result.Error
return nil, result.Error
}
if result.RowsAffected > 0 {
@@ -309,7 +353,7 @@ func (h *PackageActivationHandler) invalidateAddons(ctx context.Context, tx *gor
zap.Int64("invalidated_count", result.RowsAffected))
}
return nil
return addons, nil
}
// getCarrierInfo 获取载体信息
@@ -323,31 +367,27 @@ func (h *PackageActivationHandler) getCarrierInfo(pkg *model.PackageUsage) (stri
return "", 0
}
// activateNextPackage 任务 19.5: 激活下一个待生效主套餐
// activateNextPackage 提交下一个待生效主套餐的异步激活任务。
func (h *PackageActivationHandler) activateNextPackage(ctx context.Context, tx *gorm.DB, carrierType string, carrierID uint) error {
// 查询下一个待生效主套餐
// WHERE status=0 AND master_usage_id IS NULL ORDER BY priority ASC LIMIT 1
var nextPkg model.PackageUsage
query := tx.Where("status = ?", constants.PackageUsageStatusPending).
Where("master_usage_id IS NULL"). // 主套餐
Where("master_usage_id IS NULL").
Order("priority ASC, created_at ASC, id ASC").
Limit(1)
if carrierType == "iot_card" {
if carrierType == constants.AssetTypeIotCard {
query = query.Where("iot_card_id = ?", carrierID)
} else if carrierType == "device" {
} else if carrierType == constants.AssetTypeDevice {
query = query.Where("device_id = ?", carrierID)
}
if err := query.First(&nextPkg).Error; err != nil {
if err == gorm.ErrRecordNotFound {
// 没有待生效套餐,正常情况
return nil
}
return err
}
// 提交 Asynq 任务进行激活(避免长事务)
return h.enqueueActivationTask(ctx, nextPkg.ID, carrierType, carrierID, "queue")
}
@@ -357,10 +397,11 @@ func (h *PackageActivationHandler) triggerStopAfterExpiry(ctx context.Context, c
if h.stopResumeCallback == nil {
return
}
detachedCtx := context.WithoutCancel(ctx)
if carrierType == "iot_card" {
go func() {
if err := h.stopResumeCallback.CheckAndStopCard(context.Background(), carrierID); err != nil {
if err := h.stopResumeCallback.CheckAndStopCard(detachedCtx, carrierID); err != nil {
h.logger.Error("套餐过期后停机失败",
zap.Uint("card_id", carrierID),
zap.Error(err))
@@ -380,7 +421,7 @@ func (h *PackageActivationHandler) triggerStopAfterExpiry(ctx context.Context, c
for _, b := range bindings {
cardID := b.IotCardID
go func(cID uint) {
if err := h.stopResumeCallback.CheckAndStopCard(context.Background(), cID); err != nil {
if err := h.stopResumeCallback.CheckAndStopCard(detachedCtx, cID); err != nil {
h.logger.Error("套餐过期后停机失败",
zap.Uint("card_id", cID),
zap.Error(err))
@@ -392,12 +433,16 @@ func (h *PackageActivationHandler) triggerStopAfterExpiry(ctx context.Context, c
// 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)
@@ -434,6 +479,11 @@ func (h *PackageActivationHandler) HandlePackageQueueActivation(ctx context.Cont
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),
@@ -458,19 +508,12 @@ func (h *PackageActivationHandler) HandlePackageQueueActivation(ctx context.Cont
// 调用 ActivationService 执行激活
if h.activationService != nil {
activated, err := h.activationService.ActivateSpecificPackage(ctx, payload.PackageUsageID)
if err != 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
}
if !activated {
h.logger.Info("套餐本次未激活,等待后续检查",
zap.Uint("package_usage_id", payload.PackageUsageID),
zap.String("activation_type", payload.ActivationType))
return nil
}
} else {
// ActivationService 未注入,无法安全激活(缺少 expires_at 计算)
h.logger.Error("激活服务未注入,无法执行套餐激活",
@@ -479,8 +522,7 @@ func (h *PackageActivationHandler) HandlePackageQueueActivation(ctx context.Cont
}
h.logger.Info("套餐激活成功",
zap.Uint("package_usage_id", payload.PackageUsageID),
zap.String("activation_type", payload.ActivationType))
zap.Uint("package_usage_id", payload.PackageUsageID))
return nil
}
@@ -495,6 +537,11 @@ func (h *PackageActivationHandler) HandlePackageFirstActivation(ctx context.Cont
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 信息缺失",