package polling import ( "context" "strconv" "strings" "time" "go.uber.org/zap" "gorm.io/gorm" priorityapp "github.com/break/junhong_cmp_fiber/internal/application/prioritypolling" "github.com/break/junhong_cmp_fiber/internal/infrastructure/audit" "github.com/break/junhong_cmp_fiber/internal/model" "github.com/break/junhong_cmp_fiber/internal/store/postgres" "github.com/break/junhong_cmp_fiber/pkg/constants" "github.com/break/junhong_cmp_fiber/pkg/errors" "github.com/break/junhong_cmp_fiber/pkg/middleware" ) // PriorityEnqueueService 人工优先入队用例。 // // 与既有人工触发的关系:共用同包权限判定(permission.go),但**不继承**人工触发的每日次数上限与 // 24 小时去重——它们是人工触发的防滥用配额,不是队列不变量;重复抑制由活动项合并承担, // 且每次人工入队独立记录审计。本用例也不修改调度优先级、不绕过既有并发上限。 type PriorityEnqueueService struct { db *gorm.DB iotCardStore *postgres.IotCardStore priorityStore *postgres.PollingPriorityItemStore publisher priorityapp.PromptPublisher auditWriter *audit.Writer logger *zap.Logger } // NewPriorityEnqueueService 创建人工优先入队用例。 func NewPriorityEnqueueService( db *gorm.DB, iotCardStore *postgres.IotCardStore, priorityStore *postgres.PollingPriorityItemStore, publisher priorityapp.PromptPublisher, logger *zap.Logger, ) *PriorityEnqueueService { return &PriorityEnqueueService{ db: db, iotCardStore: iotCardStore, priorityStore: priorityStore, publisher: publisher, logger: logger, } } // SetAudit 注入统一审计 Writer。 func (s *PriorityEnqueueService) SetAudit(writer *audit.Writer) { s.auditWriter = writer } // EnqueueResult 描述一次人工优先入队的结果:一次优先需求覆盖该卡的全部纳入任务类型。 type EnqueueResult struct { CardID uint TaskTypes []string CreatedCount int MergedCount int Items []EnqueueItemResult } // EnqueueItemResult 是单个任务类型的入队结果。 type EnqueueItemResult struct { ItemID uint TaskType string Status string TriggerCount int LastTriggeredAt time.Time Created bool } // Enqueue 为该卡建立或合并全部纳入轮询任务类型的优先项,并在提交后逐类型下发执行提示。 // // 一次优先需求 = 该卡的全部纳入任务类型(constants.PollingPriorityTaskTypes), // 与执行集合一致;入队对象是卡,不做设备维度入队。 func (s *PriorityEnqueueService) Enqueue(ctx context.Context, cardID uint, reason string, operatorID uint) (*EnqueueResult, error) { if cardID == 0 { return nil, errors.New(errors.CodeInvalidParam, "无效的卡ID") } trimmedReason := strings.TrimSpace(reason) if trimmedReason == "" { return nil, s.deny(ctx, cardID, operatorID, nil, errors.CodeInvalidParam, "人工优先入队原因不能为空") } if len([]rune(trimmedReason)) > constants.PollingPriorityManualReasonMaxLength { return nil, s.deny(ctx, cardID, operatorID, nil, errors.CodeInvalidParam, "人工优先入队原因长度超出上限") } // 权限判定与既有人工触发共用同一套语义:超管/平台放行、企业拒绝、代理限自身与下级、平台卡不可见。 if err := canManagePollingCard(ctx, s.iotCardStore, cardID); err != nil { return nil, s.deny(ctx, cardID, operatorID, nil, errors.CodeForbidden, "无权限操作该资源或资源不存在") } card, err := s.iotCardStore.GetByID(ctx, cardID) if err != nil || card == nil { return nil, s.deny(ctx, cardID, operatorID, nil, errors.CodeForbidden, "无权限操作该资源或资源不存在") } operatorName := middleware.GetUsernameFromContext(ctx) taskTypes := constants.PollingPriorityTaskTypes() result := &EnqueueResult{CardID: cardID, TaskTypes: taskTypes} err = runPollingTransaction(ctx, s.db, s.auditWriter, func(tx *gorm.DB) error { store := s.priorityStore.WithTx(tx) for _, taskType := range taskTypes { item := &model.PollingPriorityItem{ CardID: cardID, TaskType: taskType, TriggerType: constants.PollingPriorityTriggerManual, ManualReason: trimmedReason, ManualOperatorID: operatorID, ManualOperatorName: operatorName, ShopIDSnapshot: card.ShopID, } // 先读取合并前的活动项作为审计 Before 的真实快照;无活动项则 Before 为空。 // 残余近似:并发下相邻执行可能在该读取与合并之间把该行终态化并另建新行, // 此时 Before 描述的是读取时点的活动项而非实际被合并行,属可接受的审计快照。 before, beforeErr := store.FindActive(ctx, cardID, taskType) if beforeErr != nil { return beforeErr } created, insertErr := store.InsertOrMerge(ctx, item) if insertErr != nil { return insertErr } active, activeErr := store.FindActive(ctx, cardID, taskType) if activeErr != nil { return activeErr } if active == nil { return errors.New(errors.CodeInternalError, "优先轮询项入队后未能读回活动项") } if created { result.CreatedCount++ } else { result.MergedCount++ } result.Items = append(result.Items, EnqueueItemResult{ ItemID: active.ID, TaskType: active.TaskType, Status: active.Status, TriggerCount: active.TriggerCount, LastTriggeredAt: active.LastTriggeredAt, Created: created, }) // 入队事实与审计同事务:合并结果同样落库,避免「加急事实已生效但无审计」。 if auditErr := writePollingAudit(ctx, tx, s.auditWriter, audit.PollingInput{ ActionCode: constants.AuditActionPollingPriorityEnqueued, Summary: "人工优先入队", ResourceType: constants.AuditResourcePollingPriorityItem, ResourceID: active.ID, ResourceKey: pollingPriorityResourceKey(active.ID), DisplayName: "卡轮询优先项", OperatorID: operatorID, IdentitySnapshot: map[string]any{ "id": active.ID, "card_id": active.CardID, "task_type": active.TaskType, "status": active.Status, "trigger_type": active.TriggerType, "trigger_count": active.TriggerCount, "manual_operator_id": active.ManualOperatorID, }, BeforeData: pollingPriorityBeforeData(before), AfterData: map[string]any{ "status": active.Status, "trigger_type": active.TriggerType, "trigger_types": triggerTypesReadable(active.TriggerTypes), "trigger_count": active.TriggerCount, "last_triggered_at": active.LastTriggeredAt, "manual_reason": active.ManualReason, "manual_operator_id": active.ManualOperatorID, "source": constants.AuditSourceAdminAPI, "occurred_at": time.Now(), }, Metadata: map[string]any{"created": created, "task_type": taskType, "card_id": cardID}, Cards: []*model.IotCard{card}, }); auditErr != nil { return auditErr } } return nil }) if err != nil { return nil, err } // 提交后逐任务类型下发执行提示:与既有手动触发队列分离的优先提示通道,每类型独立键。 if s.publisher == nil { s.logger.Error("优先轮询提示通道未配置,本次入队只依赖普通轮询兜底", zap.Uint("card_id", cardID)) } else { for _, taskType := range taskTypes { if publishErr := s.publisher.EnqueuePriority(ctx, cardID, taskType); publishErr != nil { // 提示通道不是权威:下发失败只退化为延迟一个普通轮询周期,不撤销已生效的入队事实。 s.logger.Error("下发优先轮询执行提示失败", zap.Uint("card_id", cardID), zap.String("task_type", taskType), zap.Error(publishErr)) } } } s.logger.Info("人工优先入队成功", zap.Uint("card_id", cardID), zap.Int("created_count", result.CreatedCount), zap.Int("merged_count", result.MergedCount), zap.Uint("operator_id", operatorID)) return result, nil } // deny 记录一次人工优先入队被拒绝的审计,并返回原错误。 // 拒绝发生在业务事务之外(或业务已回滚),因此沿用既有轮询模块的独立短事务模式。 func (s *PriorityEnqueueService) deny(ctx context.Context, cardID, operatorID uint, card *model.IotCard, code int, summary string) error { appErr := errors.New(code, summary) identity := map[string]any{"card_id": cardID, "manual_operator_id": operatorID} var cards []*model.IotCard if card != nil { cards = []*model.IotCard{card} } recordPollingFailure(ctx, s.db, s.auditWriter, audit.PollingInput{ ActionCode: constants.AuditActionPollingPriorityManualDenied, Summary: summary, ResourceType: constants.AuditResourcePollingPriorityItem, ResourceKey: pollingPriorityAttemptKey(cardID, operatorID), DisplayName: "人工优先入队", OperatorID: operatorID, Result: constants.AuditResultDenied, IdentitySnapshot: identity, Cards: cards, }, appErr) return appErr } // pollingPriorityResourceKey 返回优先项在审计里的稳定资源键。 func pollingPriorityResourceKey(itemID uint) string { return strconv.FormatUint(uint64(itemID), 10) } // pollingPriorityAttemptKey 返回无入队对象的拒绝审计资源键(按卡与操作者稳定)。 func pollingPriorityAttemptKey(cardID, operatorID uint) string { return "priority:" + strconv.FormatUint(uint64(cardID), 10) + ":" + strconv.FormatUint(uint64(operatorID), 10) } // pollingPriorityBeforeData 把合并前的活动项快照投影为审计前后值;无活动项(新建)时返回空对象。 func pollingPriorityBeforeData(before *model.PollingPriorityItem) map[string]any { if before == nil { return map[string]any{} } return map[string]any{ "status": before.Status, "trigger_type": before.TriggerType, "trigger_types": triggerTypesReadable(before.TriggerTypes), "trigger_count": before.TriggerCount, } } // triggerTypesReadable 将存储形式的触发类型集合(两端补逗号)还原为可读的逗号分隔形式。 func triggerTypesReadable(raw string) string { return strings.Trim(raw, ",") }