diff --git a/docs/yumi_gift_challenge_api.md b/docs/yumi_gift_challenge_api.md index d85c944..6c30a96 100644 --- a/docs/yumi_gift_challenge_api.md +++ b/docs/yumi_gift_challenge_api.md @@ -36,7 +36,9 @@ Yumi 的 Java `sysOrigin` 固定为 `LIKEI`。网关外部地址在以下路径 | POST | `/settlement/:id` | JSON `{ "periodType":"DAILY\|OVERALL", "statDate":"yyyy-MM-dd" }`;OVERALL 不传日期 | `...:settle` | | POST | `/delivery-item/resolve/:itemId` | JSON `{ "delivered": true\|false }`;仅人工核账 UNKNOWN | `...:reconcile` | -活动启用时会校验奖励组已上架且属于 LIKEI,并冻结完整奖励项。活动开始后仅允许通过 `/save` 修改名称和描述,其他统计/奖励口径只读。 +新建并启用活动、或重新启用已停用活动时,`startTime`、`endTime` 均须晚于服务端当前时间;同时校验奖励组已上架且属于 LIKEI,并冻结完整奖励项。活动即使已经进入时间窗,只要 `overallSettlementStatus=NOT_STARTED`,仍可通过 `/save`、`/tasks/:id`、`/rank-rewards/:id` 修改完整配置。只有总榜结算进入 `PROCESSING/COMPLETED/RECONCILIATION_REQUIRED` 后,才仅允许修改活动名称和说明。 + +进行中编辑按提交时点向后生效,不回算已经入账的积分。已物化的用户日任务、已冻结日榜、榜奖 parent 和 delivery item 保留原目标及奖励;尚未物化的用户/日期和未冻结榜奖使用新配置。旧任务或榜奖引用的资源组快照始终保留;同一已引用资源组不会刷新,若要更换奖励内容应选择新的资源组。修改时间窗或时区后,事件事务会在活动共享锁内重新读取窗口和时区;已产生数据的历史日期门闩保留,其余 `NOT_STARTED` 门闩按新配置重建。 日榜、总榜结算延迟均允许 `0-1440` 分钟,且 `overallSettlementDelayMinutes` 必须大于或等于 `dailySettlementDelayMinutes`,保证最终日榜不会晚于总榜关门。 ## chatapp-cron 触发 diff --git a/internal/model/yumi_gift_challenge_models.go b/internal/model/yumi_gift_challenge_models.go index 778c88e..a8fb068 100644 --- a/internal/model/yumi_gift_challenge_models.go +++ b/internal/model/yumi_gift_challenge_models.go @@ -2,7 +2,8 @@ package model import "time" -// YumiGiftChallengeActivity 保存一次独立活动周期;活动开始后核心时间和奖励配置只读。 +// YumiGiftChallengeActivity 保存一次独立活动周期;总榜结算开始前允许调整活动配置, +// 已产生的用户任务、日榜结算和奖励项由各自快照保持历史语义。 type YumiGiftChallengeActivity struct { ID int64 `gorm:"column:id;primaryKey"` ActivityCode string `gorm:"column:activity_code"` diff --git a/internal/service/yumigiftchallenge/config.go b/internal/service/yumigiftchallenge/config.go index 6946349..e5f0154 100644 --- a/internal/service/yumigiftchallenge/config.go +++ b/internal/service/yumigiftchallenge/config.go @@ -19,7 +19,8 @@ import ( var taskCodePattern = regexp.MustCompile(`^[A-Za-z0-9_-]{1,64}$`) -// SaveActivity 保存完整配置;活动开始后只允许修改名称和描述,不触碰已生效统计口径。 +// SaveActivity 保存完整配置。活动进入时间窗后仍可调整配置;只有总榜门闩已经冻结时才 +// 降级为展示文案更新,防止修改已经确定的总榜获奖人与发奖归属。 func (s *Service) SaveActivity(ctx context.Context, req SaveRequest) (*DetailResponse, error) { activityID := req.ID.Int64() if activityID > 0 { @@ -27,8 +28,8 @@ func (s *Service) SaveActivity(ctx context.Context, req SaveRequest) (*DetailRes if err != nil { return nil, err } - if !time.Now().Before(current.StartTime) { - return s.saveStartedActivityMetadata(ctx, *current, req) + if !activityCoreConfigEditable(*current) { + return s.saveFrozenActivityMetadata(ctx, *current, req) } } origin, err := s.requireYumiOrigin(req.SysOrigin) @@ -84,8 +85,10 @@ func (s *Service) SaveActivity(ctx context.Context, req SaveRequest) (*DetailRes } } + isNew := req.ID.Int64() <= 0 err = s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { - if req.ID.Int64() <= 0 { + preserveLiveState := false + if isNew { if req.Enabled { if !proposed.StartTime.After(time.Now()) || !proposed.EndTime.After(time.Now()) { return NewAppError(http.StatusConflict, "invalid_enable_window", "enabled activity must start in the future") @@ -104,21 +107,33 @@ func (s *Service) SaveActivity(ctx context.Context, req SaveRequest) (*DetailRes if req.Version == nil { return NewAppError(http.StatusBadRequest, "version_required", "version is required when updating") } - var current model.YumiGiftChallengeActivity - if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("id = ?", activityID).First(¤t).Error; errors.Is(err, gorm.ErrRecordNotFound) { + current, err := lockActivityConfigurationTx(tx, activityID) + if errors.Is(err, gorm.ErrRecordNotFound) { return NewAppError(http.StatusNotFound, "activity_not_found", "activity not found") } else if err != nil { return err } - if !time.Now().Before(current.StartTime) { - return NewAppError(http.StatusConflict, "activity_started", "activity started while saving; retry metadata-only update") + if !activityCoreConfigEditable(current) { + return NewAppError(http.StatusConflict, "activity_settled", "activity core config is locked after overall settlement starts") } if current.Version != *req.Version { return NewAppError(http.StatusConflict, "version_conflict", "activity has been changed; reload before saving") } + txNow := time.Now() + startedActive := current.Enabled && !txNow.Before(current.StartTime) + preserveLiveState = startedActive + if !preserveLiveState { + preserveLiveState, err = activityHasHistoricalStateTx(tx, activityID) + if err != nil { + return err + } + } + if startedActive && !req.Enabled { + return NewAppError(http.StatusConflict, "activity_status_readonly", "started activity cannot be disabled") + } if req.Enabled { - if !proposed.StartTime.After(time.Now()) || !proposed.EndTime.After(time.Now()) || current.OverallSettlementStatus != StatusNotStarted { - return NewAppError(http.StatusConflict, "invalid_enable_window", "enabled activity must start in the future and remain unsettled") + if !current.Enabled && (!proposed.StartTime.After(txNow) || !proposed.EndTime.After(txNow)) { + return NewAppError(http.StatusConflict, "invalid_enable_window", "disabled activity must move its full time window to the future before enabling") } if err := ensureNoEnabledOverlap(tx, proposed); err != nil { return err @@ -141,7 +156,11 @@ func (s *Service) SaveActivity(ctx context.Context, req SaveRequest) (*DetailRes } return NewAppError(http.StatusConflict, "version_conflict", "activity has been changed; reload before saving") } - if err := deletePreStartArtifacts(tx, activityID); err != nil { + } + if !isNew { + // taskCode/type/sort 是 user_task_daily 唯一与进度更新的身份边界; + // 任一用户任务已物化后只能改展示、门槛和奖励,不能换槽位身份。 + if err := assertMaterializedTaskIdentityUnchangedTx(tx, activityID, tasks); err != nil { return err } } @@ -158,12 +177,15 @@ func (s *Service) SaveActivity(ctx context.Context, req SaveRequest) (*DetailRes return err } if req.Enabled { - if err := tx.CreateInBatches(snapshots, 500).Error; err != nil { - return err - } - if err := tx.CreateInBatches(headers, 100).Error; err != nil { - return err + if preserveLiveState { + return mergeEditableActivityArtifactsTx(tx, activityID, snapshots, headers) } + return replaceEditableActivityArtifactsTx(tx, activityID, snapshots, headers) + } + // 只有未进入活动时间窗的停用草稿会走到这里:清掉旧预建门闩和 + // 奖励快照,避免之后重新启用时夹带过期配置。进行中活动不允许停用。 + if !isNew && !preserveLiveState { + return deletePreStartArtifacts(tx, activityID) } return nil }) @@ -173,7 +195,7 @@ func (s *Service) SaveActivity(ctx context.Context, req SaveRequest) (*DetailRes return s.GetAdminDetail(ctx, activityID) } -func (s *Service) saveStartedActivityMetadata(ctx context.Context, current model.YumiGiftChallengeActivity, req SaveRequest) (*DetailResponse, error) { +func (s *Service) saveFrozenActivityMetadata(ctx context.Context, current model.YumiGiftChallengeActivity, req SaveRequest) (*DetailResponse, error) { if req.Version == nil { return nil, NewAppError(http.StatusBadRequest, "version_required", "version is required when updating") } @@ -182,7 +204,7 @@ func (s *Service) saveStartedActivityMetadata(ctx context.Context, current model return nil, NewAppError(http.StatusBadRequest, "invalid_activity_text", "activityName is required and activity text is too long") } if req.Tasks != nil || req.RankRewards != nil { - return nil, NewAppError(http.StatusConflict, "started_activity_core_readonly", "started activity tasks and rank rewards cannot be changed") + return nil, NewAppError(http.StatusConflict, "settled_activity_core_readonly", "frozen activity tasks and rank rewards cannot be changed") } err := s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { var locked model.YumiGiftChallengeActivity @@ -198,7 +220,7 @@ func (s *Service) saveStartedActivityMetadata(ctx context.Context, current model req.DailySettlementDelayMinutes != locked.DailySettlementDelayMinutes || req.OverallSettlementDelayMinutes != locked.OverallSettlementDelayMinutes || req.DisplayTopN != locked.DisplayTopN || req.Enabled != locked.Enabled { - return NewAppError(http.StatusConflict, "started_activity_core_readonly", "only activityName and activityDesc can be changed after start") + return NewAppError(http.StatusConflict, "settled_activity_core_readonly", "only activityName and activityDesc can be changed after settlement starts") } result := tx.Model(&model.YumiGiftChallengeActivity{}).Where("id = ? AND version = ?", locked.ID, *req.Version). Updates(map[string]any{"activity_name": name, "activity_desc": desc, "version": gorm.Expr("version + 1"), "update_time": time.Now()}) @@ -216,6 +238,13 @@ func (s *Service) saveStartedActivityMetadata(ctx context.Context, current model return s.GetAdminDetail(ctx, current.ID) } +// activityCoreConfigEditable 只以总榜是否冻结作为最终写门闩。进行中活动的历史积分、 +// 已物化任务和发奖项由各自快照承接,因此配置可以对后续事件/用户生效;总榜一旦进入 +// PROCESSING/COMPLETED/RECONCILIATION_REQUIRED 后则不能再改变最终获奖语义。 +func activityCoreConfigEditable(activity model.YumiGiftChallengeActivity) bool { + return activity.OverallSettlementStatus == StatusNotStarted +} + func validateActivityInput(req SaveRequest) (*time.Location, time.Time, time.Time, error) { if strings.TrimSpace(req.ActivityCode) == "" || len(strings.TrimSpace(req.ActivityCode)) > 64 { return nil, time.Time{}, time.Time{}, NewAppError(http.StatusBadRequest, "invalid_activity_code", "activityCode is required and max length is 64") @@ -390,18 +419,40 @@ func (s *Service) SetEnabled(ctx context.Context, activityID int64, enabled bool if err != nil { return err } - if !time.Now().Before(activity.StartTime) { - return NewAppError(http.StatusConflict, "activity_readonly", "started activity cannot be enabled or disabled") + if activity.OverallSettlementStatus != StatusNotStarted { + return NewAppError(http.StatusConflict, "activity_readonly", "settled activity cannot be enabled or disabled") + } + now := time.Now() + if enabled && (!activity.StartTime.After(now) || !activity.EndTime.After(now)) { + return NewAppError(http.StatusConflict, "invalid_enable_window", "disabled activity must move its full time window to the future before enabling") + } + if activity.Enabled && !now.Before(activity.StartTime) { + return NewAppError(http.StatusConflict, "activity_readonly", "started activity cannot be disabled") } if !enabled { return s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { - var locked model.YumiGiftChallengeActivity - if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("id = ?", activityID).First(&locked).Error; err != nil { + locked, err := lockActivityConfigurationTx(tx, activityID) + if err != nil { return err } - if !time.Now().Before(locked.StartTime) { + if !activityCoreConfigEditable(locked) { + return NewAppError(http.StatusConflict, "activity_readonly", "settled activity cannot be disabled") + } + if !locked.Enabled { + return nil + } + if locked.Enabled && !time.Now().Before(locked.StartTime) { return NewAppError(http.StatusConflict, "activity_readonly", "started activity cannot be disabled") } + hasHistory, err := activityHasHistoricalStateTx(tx, activityID) + if err != nil { + return err + } + // startTime 可被移到未来,不能因此把真实开始过的活动当成草稿。 + // 只要有流水、用户 parent 或冻结门闩,停用就会使历史奖励失去入口。 + if hasHistory { + return NewAppError(http.StatusConflict, "activity_status_readonly", "activity with historical state cannot be disabled") + } if err := deletePreStartArtifacts(tx, activityID); err != nil { return err } @@ -431,13 +482,16 @@ func (s *Service) SetEnabled(ctx context.Context, activityID int64, enabled bool } return s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { - var locked model.YumiGiftChallengeActivity - if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("id = ?", activityID).First(&locked).Error; err != nil { + locked, err := lockActivityConfigurationTx(tx, activityID) + if err != nil { return err } now := time.Now() - if !locked.StartTime.After(now) || !locked.EndTime.After(now) || locked.OverallSettlementStatus != StatusNotStarted { - return NewAppError(http.StatusConflict, "activity_readonly", "activity started while enabling") + if !activityCoreConfigEditable(locked) { + return NewAppError(http.StatusConflict, "activity_readonly", "settled activity cannot be enabled") + } + if !locked.StartTime.After(now) || !locked.EndTime.After(now) { + return NewAppError(http.StatusConflict, "invalid_enable_window", "disabled activity must move its full time window to the future before enabling") } if locked.Version != activity.Version { return NewAppError(http.StatusConflict, "version_conflict", "activity changed while reward snapshot was loading") @@ -448,13 +502,16 @@ func (s *Service) SetEnabled(ctx context.Context, activityID int64, enabled bool if err := ensureNoEnabledOverlap(tx, locked); err != nil { return err } - if err := deletePreStartArtifacts(tx, activityID); err != nil { + hasHistory, err := activityHasHistoricalStateTx(tx, activityID) + if err != nil { return err } - if err := tx.CreateInBatches(snapshots, 500).Error; err != nil { - return err + if hasHistory { + err = mergeEditableActivityArtifactsTx(tx, activityID, snapshots, headers) + } else { + err = replaceEditableActivityArtifactsTx(tx, activityID, snapshots, headers) } - if err := tx.CreateInBatches(headers, 100).Error; err != nil { + if err != nil { return err } return tx.Model(&model.YumiGiftChallengeActivity{}).Where("id = ?", activityID). @@ -615,20 +672,280 @@ func buildPeriodHeaders(activity model.YumiGiftChallengeActivity) ([]model.YumiG } func ensureNoEnabledOverlap(tx *gorm.DB, activity model.YumiGiftChallengeActivity) error { - var count int64 - err := tx.Model(&model.YumiGiftChallengeActivity{}). + var overlapIDs []int64 + // 该查询命中 (sys_origin, enabled, start_time, end_time) 索引,只返回 id; + // FOR UPDATE 同时锁住命中行/范围,两个运营并发移动时间窗不能都通过校验。 + err := tx.Model(&model.YumiGiftChallengeActivity{}).Clauses(clause.Locking{Strength: "UPDATE"}). Where("sys_origin = ? AND enabled = ? AND id <> ? AND start_time < ? AND end_time > ?", activity.SysOrigin, true, activity.ID, activity.EndTime, activity.StartTime). - Count(&count).Error + Order("start_time ASC, id ASC").Pluck("id", &overlapIDs).Error if err != nil { return err } - if count > 0 { + if len(overlapIDs) > 0 { return NewAppError(http.StatusConflict, "activity_time_overlap", "enabled activity time windows must not overlap") } return nil } +// lockActivityConfigurationTx 按“活动 -> 日榜门闩 -> 总榜门闩”的固定顺序抢锁。 +// 送礼、任务和结算事务都遵循相同顺序,使活动时间/时区与其门闩成为同一个版本边界, +// 避免在线编辑与总榜结算分别持有 activity/header 后形成死锁环路。 +func lockActivityConfigurationTx(tx *gorm.DB, activityID int64) (model.YumiGiftChallengeActivity, error) { + var activity model.YumiGiftChallengeActivity + if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("id = ?", activityID).First(&activity).Error; err != nil { + return model.YumiGiftChallengeActivity{}, err + } + var dailyHeaders []model.YumiGiftChallengePeriodSettlement + if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}). + Where("activity_id = ? AND period_type = ?", activityID, PeriodDaily). + Order("id ASC").Find(&dailyHeaders).Error; err != nil { + return model.YumiGiftChallengeActivity{}, err + } + var overallHeader model.YumiGiftChallengePeriodSettlement + err := tx.Clauses(clause.Locking{Strength: "UPDATE"}). + Where("activity_id = ? AND period_type = ? AND period_key = ?", activityID, PeriodOverall, PeriodOverall). + First(&overallHeader).Error + if err != nil && !errors.Is(err, gorm.ErrRecordNotFound) { + return model.YumiGiftChallengeActivity{}, err + } + return activity, nil +} + +// activityHasHistoricalStateTx 不使用可被运营改写的 startTime 判断历史。任一用户 +// parent/流水或已冻结周期存在,就必须走增量合并,永不再整体删除快照。 +func activityHasHistoricalStateTx(tx *gorm.DB, activityID int64) (bool, error) { + queries := []struct { + model any + where string + args []any + }{ + {model: &model.YumiGiftChallengeGiftLedger{}, where: "activity_id = ?", args: []any{activityID}}, + {model: &model.YumiGiftChallengeUserTaskDaily{}, where: "activity_id = ?", args: []any{activityID}}, + {model: &model.YumiGiftChallengeSettlement{}, where: "activity_id = ?", args: []any{activityID}}, + {model: &model.YumiGiftChallengePeriodSettlement{}, where: "activity_id = ? AND status <> ?", args: []any{activityID, StatusNotStarted}}, + } + for _, query := range queries { + var marker struct { + ActivityID int64 `gorm:"column:activity_id"` + } + err := tx.Model(query.model).Select("activity_id").Where(query.where, query.args...).Limit(1).Take(&marker).Error + if err == nil { + return true, nil + } + if !errors.Is(err, gorm.ErrRecordNotFound) { + return false, err + } + } + return false, nil +} + +// assertMaterializedTaskIdentityUnchangedTx 保护已生成用户任务的唯一键语义。如果 +// 允许改 taskCode/type/sort,同一用户当天会再物化一份任务并重复领奖。 +func assertMaterializedTaskIdentityUnchangedTx(tx *gorm.DB, activityID int64, proposed []model.YumiGiftChallengeTaskConfig) error { + var marker model.YumiGiftChallengeUserTaskDaily + err := tx.Select("id").Where("activity_id = ?", activityID).Limit(1).Take(&marker).Error + if errors.Is(err, gorm.ErrRecordNotFound) { + return nil + } + if err != nil { + return err + } + var current []model.YumiGiftChallengeTaskConfig + if err := tx.Where("activity_id = ?", activityID).Order("sort_order ASC, id ASC").Find(¤t).Error; err != nil { + return err + } + sortTaskRows(proposed) + if len(current) != len(proposed) { + return NewAppError(http.StatusConflict, "materialized_task_identity_readonly", "taskCode, taskType and sortOrder cannot change after user tasks exist") + } + for index := range current { + if current[index].TaskCode != proposed[index].TaskCode || current[index].TaskType != proposed[index].TaskType || current[index].SortOrder != proposed[index].SortOrder { + return NewAppError(http.StatusConflict, "materialized_task_identity_readonly", "taskCode, taskType and sortOrder cannot change after user tasks exist") + } + } + return nil +} + +func replaceEditableActivityArtifactsTx(tx *gorm.DB, activityID int64, snapshots []model.YumiGiftChallengeRewardSnapshot, headers []model.YumiGiftChallengePeriodSettlement) error { + if err := deletePreStartArtifacts(tx, activityID); err != nil { + return err + } + if len(snapshots) > 0 { + if err := tx.CreateInBatches(snapshots, 500).Error; err != nil { + return err + } + } + if len(headers) > 0 { + return tx.CreateInBatches(headers, 100).Error + } + return nil +} + +// mergeEditableActivityArtifactsTx 只增量应用进行中活动的新配置:已经落到用户任务或 +// 榜奖 parent 里的旧资源组快照绝不删除;未被历史状态引用的当前资源组允许刷新。 +func mergeEditableActivityArtifactsTx(tx *gorm.DB, activityID int64, snapshots []model.YumiGiftChallengeRewardSnapshot, headers []model.YumiGiftChallengePeriodSettlement) error { + // 先在已持有 activity/header 排他锁的事务内完成门闩重建,关闭新的任务物化入口; + // 随后扫描 task/settlement parent,得到的受保护资源组集合才是稳定快照。 + if err := syncEditablePeriodHeadersTx(tx, activityID, headers); err != nil { + return err + } + protectedGroups, err := protectedRewardGroupIDsTx(tx, activityID) + if err != nil { + return err + } + deleteQuery := tx.Where("activity_id = ?", activityID) + if len(protectedGroups) > 0 { + deleteQuery = deleteQuery.Where("resource_group_id NOT IN ?", protectedGroups) + } + if err := deleteQuery.Delete(&model.YumiGiftChallengeRewardSnapshot{}).Error; err != nil { + return err + } + frozenProtectedGroups := map[int64]struct{}{} + if len(protectedGroups) > 0 { + var ids []int64 + if err := tx.Model(&model.YumiGiftChallengeRewardSnapshot{}). + Distinct("resource_group_id").Where("activity_id = ? AND resource_group_id IN ?", activityID, protectedGroups). + Pluck("resource_group_id", &ids).Error; err != nil { + return err + } + for _, id := range ids { + frozenProtectedGroups[id] = struct{}{} + } + } + insertSnapshots := make([]model.YumiGiftChallengeRewardSnapshot, 0, len(snapshots)) + for _, snapshot := range snapshots { + if _, frozen := frozenProtectedGroups[snapshot.ResourceGroupID]; !frozen { + insertSnapshots = append(insertSnapshots, snapshot) + } + } + if len(insertSnapshots) > 0 { + // 只要 protected group 已存在任一冻结项,就整组保持原样,不能把资源中心后来新增 + // 的项混入旧奖励;运营若要改变奖励必须选择新资源组。缺失整组时允许补建以修复脏数据。 + if err := tx.Clauses(clause.OnConflict{DoNothing: true}).CreateInBatches(insertSnapshots, 500).Error; err != nil { + return err + } + } + var snapshotCount int64 + if err := tx.Model(&model.YumiGiftChallengeRewardSnapshot{}).Where("activity_id = ?", activityID).Count(&snapshotCount).Error; err != nil { + return err + } + if snapshotCount > maxSnapshotItems { + return NewAppError(http.StatusConflict, "too_many_reward_snapshots", "historical and current reward snapshots may contain at most 5000 items") + } + return nil +} + +func protectedRewardGroupIDsTx(tx *gorm.DB, activityID int64) ([]int64, error) { + groups := map[int64]struct{}{} + var taskGroups []int64 + if err := tx.Model(&model.YumiGiftChallengeUserTaskDaily{}). + Distinct("resource_group_id").Where("activity_id = ? AND resource_group_id IS NOT NULL", activityID). + Pluck("resource_group_id", &taskGroups).Error; err != nil { + return nil, err + } + var settlementGroups []int64 + if err := tx.Model(&model.YumiGiftChallengeSettlement{}). + Distinct("resource_group_id").Where("activity_id = ?", activityID). + Pluck("resource_group_id", &settlementGroups).Error; err != nil { + return nil, err + } + for _, id := range append(taskGroups, settlementGroups...) { + if id > 0 { + groups[id] = struct{}{} + } + } + result := make([]int64, 0, len(groups)) + for id := range groups { + result = append(result, id) + } + sort.Slice(result, func(i, j int) bool { return result[i] < result[j] }) + return result, nil +} + +func syncEditablePeriodHeadersTx(tx *gorm.DB, activityID int64, desired []model.YumiGiftChallengePeriodSettlement) error { + var existing []model.YumiGiftChallengePeriodSettlement + if err := tx.Where("activity_id = ?", activityID).Find(&existing).Error; err != nil { + return err + } + desiredByKey := make(map[string]model.YumiGiftChallengePeriodSettlement, len(desired)) + for _, row := range desired { + desiredByKey[row.PeriodType+"\x00"+row.PeriodKey] = row + } + existingByKey := make(map[string]model.YumiGiftChallengePeriodSettlement, len(existing)) + for _, row := range existing { + existingByKey[row.PeriodType+"\x00"+row.PeriodKey] = row + } + protectedDailyKeys, err := materializedDailyPeriodKeysTx(tx, activityID) + if err != nil { + return err + } + // 所有门闩已在 lockActivityConfigurationTx 中按固定顺序锁定。此处删除/重建仅影响 + // NOT_STARTED 行,其他事务在提交前看不到中间状态;已冻结周期原行完全不触碰。 + if err := tx.Where("activity_id = ? AND status = ?", activityID, StatusNotStarted). + Delete(&model.YumiGiftChallengePeriodSettlement{}).Error; err != nil { + return err + } + insertRows := make([]model.YumiGiftChallengePeriodSettlement, 0, len(desired)+len(protectedDailyKeys)) + for _, row := range desired { + key := row.PeriodType + "\x00" + row.PeriodKey + old, existed := existingByKey[key] + _, materialized := protectedDailyKeys[row.PeriodKey] + if row.PeriodType == PeriodDaily && materialized && existed && old.Status == StatusNotStarted { + // 日期仍在新时间窗内也不代表可以重算截止时间:只要已有流水、 + // 积分或用户/settlement parent,原 header 和 due 就是该日历史边界。 + insertRows = append(insertRows, old) + continue + } + insertRows = append(insertRows, row) + } + for _, row := range existing { + if row.Status != StatusNotStarted || row.PeriodType != PeriodDaily { + continue + } + key := row.PeriodType + "\x00" + row.PeriodKey + if _, stillConfigured := desiredByKey[key]; stillConfigured { + continue + } + if _, materialized := protectedDailyKeys[row.PeriodKey]; materialized { + // 时间窗缩短时,已经有积分/任务/结算 parent 的旧日期仍需保留原截止时间, + // 否则历史数据将失去可结算门闩。 + insertRows = append(insertRows, row) + } + } + if len(insertRows) > 0 { + return tx.Clauses(clause.OnConflict{DoNothing: true}).CreateInBatches(insertRows, 100).Error + } + return nil +} + +func materializedDailyPeriodKeysTx(tx *gorm.DB, activityID int64) (map[string]struct{}, error) { + result := map[string]struct{}{} + queries := []struct { + model any + column string + where string + }{ + {model: &model.YumiGiftChallengeGiftLedger{}, column: "stat_date", where: "activity_id = ?"}, + {model: &model.YumiGiftChallengeUserDailyScore{}, column: "stat_date", where: "activity_id = ?"}, + {model: &model.YumiGiftChallengeUserTaskDaily{}, column: "stat_date", where: "activity_id = ?"}, + {model: &model.YumiGiftChallengeSettlement{}, column: "period_key", where: "activity_id = ? AND period_type = 'DAILY'"}, + } + for _, query := range queries { + var keys []string + if err := tx.Model(query.model).Distinct(query.column).Where(query.where, activityID). + Pluck(query.column, &keys).Error; err != nil { + return nil, err + } + for _, key := range keys { + if key = strings.TrimSpace(key); key != "" { + result[key] = struct{}{} + } + } + } + return result, nil +} + func deletePreStartArtifacts(tx *gorm.DB, activityID int64) error { if err := tx.Where("activity_id = ?", activityID).Delete(&model.YumiGiftChallengeRewardSnapshot{}).Error; err != nil { return err @@ -658,8 +975,8 @@ func (s *Service) SaveTasks(ctx context.Context, activityID int64, req TaskSaveR if err != nil { return nil, err } - if !time.Now().Before(activity.StartTime) { - return nil, NewAppError(http.StatusConflict, "activity_readonly", "started activity cannot be edited") + if !activityCoreConfigEditable(*activity) { + return nil, NewAppError(http.StatusConflict, "activity_readonly", "settled activity cannot be edited") } _, rewards, err := s.loadActivityChildren(ctx, activityID) if err != nil { @@ -685,18 +1002,25 @@ func (s *Service) SaveTasks(ctx context.Context, activityID int64, req TaskSaveR } } err = s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { - var locked model.YumiGiftChallengeActivity - if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("id = ?", activityID).First(&locked).Error; err != nil { + locked, err := lockActivityConfigurationTx(tx, activityID) + if err != nil { return err } now := time.Now() - if !locked.StartTime.After(now) || (locked.Enabled && (!locked.EndTime.After(now) || locked.OverallSettlementStatus != StatusNotStarted)) { - return NewAppError(http.StatusConflict, "activity_readonly", "started activity cannot be edited") + if !activityCoreConfigEditable(locked) { + return NewAppError(http.StatusConflict, "activity_readonly", "settled activity cannot be edited") } if locked.Version != activity.Version || locked.Enabled != activity.Enabled { return NewAppError(http.StatusConflict, "version_conflict", "activity changed while reward snapshot was loading") } - if err := deletePreStartArtifacts(tx, activityID); err != nil { + preserveLiveState := locked.Enabled && !now.Before(locked.StartTime) + if !preserveLiveState { + preserveLiveState, err = activityHasHistoricalStateTx(tx, activityID) + if err != nil { + return err + } + } + if err := assertMaterializedTaskIdentityUnchangedTx(tx, activityID, rows); err != nil { return err } if err := tx.Where("activity_id = ?", activityID).Delete(&model.YumiGiftChallengeTaskConfig{}).Error; err != nil { @@ -706,10 +1030,15 @@ func (s *Service) SaveTasks(ctx context.Context, activityID int64, req TaskSaveR return err } if activity.Enabled { - if err := tx.CreateInBatches(snapshots, 500).Error; err != nil { + if preserveLiveState { + if err := mergeEditableActivityArtifactsTx(tx, activityID, snapshots, headers); err != nil { + return err + } + } else if err := replaceEditableActivityArtifactsTx(tx, activityID, snapshots, headers); err != nil { return err } - if err := tx.CreateInBatches(headers, 100).Error; err != nil { + } else if !preserveLiveState { + if err := deletePreStartArtifacts(tx, activityID); err != nil { return err } } @@ -735,8 +1064,8 @@ func (s *Service) SaveRankRewards(ctx context.Context, activityID int64, req Ran if err != nil { return nil, err } - if !time.Now().Before(activity.StartTime) { - return nil, NewAppError(http.StatusConflict, "activity_readonly", "started activity cannot be edited") + if !activityCoreConfigEditable(*activity) { + return nil, NewAppError(http.StatusConflict, "activity_readonly", "settled activity cannot be edited") } tasks, _, err := s.loadActivityChildren(ctx, activityID) if err != nil { @@ -762,19 +1091,23 @@ func (s *Service) SaveRankRewards(ctx context.Context, activityID int64, req Ran } } err = s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { - var locked model.YumiGiftChallengeActivity - if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("id = ?", activityID).First(&locked).Error; err != nil { + locked, err := lockActivityConfigurationTx(tx, activityID) + if err != nil { return err } now := time.Now() - if !locked.StartTime.After(now) || (locked.Enabled && (!locked.EndTime.After(now) || locked.OverallSettlementStatus != StatusNotStarted)) { - return NewAppError(http.StatusConflict, "activity_readonly", "started activity cannot be edited") + if !activityCoreConfigEditable(locked) { + return NewAppError(http.StatusConflict, "activity_readonly", "settled activity cannot be edited") } if locked.Version != activity.Version || locked.Enabled != activity.Enabled { return NewAppError(http.StatusConflict, "version_conflict", "activity changed while reward snapshot was loading") } - if err := deletePreStartArtifacts(tx, activityID); err != nil { - return err + preserveLiveState := locked.Enabled && !now.Before(locked.StartTime) + if !preserveLiveState { + preserveLiveState, err = activityHasHistoricalStateTx(tx, activityID) + if err != nil { + return err + } } if err := tx.Where("activity_id = ?", activityID).Delete(&model.YumiGiftChallengeRankReward{}).Error; err != nil { return err @@ -783,10 +1116,15 @@ func (s *Service) SaveRankRewards(ctx context.Context, activityID int64, req Ran return err } if activity.Enabled { - if err := tx.CreateInBatches(snapshots, 500).Error; err != nil { + if preserveLiveState { + if err := mergeEditableActivityArtifactsTx(tx, activityID, snapshots, headers); err != nil { + return err + } + } else if err := replaceEditableActivityArtifactsTx(tx, activityID, snapshots, headers); err != nil { return err } - if err := tx.CreateInBatches(headers, 100).Error; err != nil { + } else if !preserveLiveState { + if err := deletePreStartArtifacts(tx, activityID); err != nil { return err } } diff --git a/internal/service/yumigiftchallenge/event.go b/internal/service/yumigiftchallenge/event.go index 1a3a6e1..ffbc21c 100644 --- a/internal/service/yumigiftchallenge/event.go +++ b/internal/service/yumigiftchallenge/event.go @@ -64,27 +64,30 @@ func (s *Service) ProcessGiftEvent(ctx context.Context, event giftEvent) error { return nil } eventTime := time.UnixMilli(event.CreateTime.Int64()) - var activity model.YumiGiftChallengeActivity - err = s.db.WithContext(ctx).Where( - "sys_origin = ? AND enabled = ? AND start_time <= ? AND end_time > ?", - origin, true, eventTime, eventTime, - ).Order("start_time DESC").First(&activity).Error - if errors.Is(err, gorm.ErrRecordNotFound) { - return nil - } - if err != nil { - return err - } - location, err := resolveLocation(activity.Timezone) - if err != nil { - return err - } - statDate := dateOnly(eventTime, location) - periodKey := dateKey(statDate, location) return s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { - // 热链路只持有共享门闩锁,同一活动的送礼事件可以并发;结算使用排他锁, - // 会等待已经进入事务的事件提交,再冻结确定的边界。 + // 活动行必须先于周期门闩加共享锁,并在锁内重新匹配时间窗/时区。配置编辑持有 + // activity UPDATE 锁后重建门闩,因此事件要么完整使用旧版本,要么等待后完整使用 + // 新版本,不会拿旧 timezone 算 statDate 后撞上新门闩而静默丢失。 + var activity model.YumiGiftChallengeActivity + queryErr := tx.Clauses(clause.Locking{Strength: "SHARE"}).Where( + "sys_origin = ? AND enabled = ? AND start_time <= ? AND end_time > ?", + origin, true, eventTime, eventTime, + ).Order("start_time DESC").First(&activity).Error + if errors.Is(queryErr, gorm.ErrRecordNotFound) { + return nil + } + if queryErr != nil { + return queryErr + } + location, err := resolveLocation(activity.Timezone) + if err != nil { + return err + } + periodKey := dateKey(dateOnly(eventTime, location), location) + + // 热链路只持有共享活动/门闩锁,同一活动的送礼事件可以并发;结算和配置使用 + // 排他锁,会等待已经进入事务的事件提交后再冻结或切换版本。 var dailyHeader model.YumiGiftChallengePeriodSettlement if err := tx.Clauses(clause.Locking{Strength: "SHARE"}).Where( "activity_id = ? AND period_type = ? AND period_key = ?", activity.ID, PeriodDaily, periodKey, diff --git a/internal/service/yumigiftchallenge/ranking.go b/internal/service/yumigiftchallenge/ranking.go index 0497684..6daa5b0 100644 --- a/internal/service/yumigiftchallenge/ranking.go +++ b/internal/service/yumigiftchallenge/ranking.go @@ -89,6 +89,11 @@ func (s *Service) ranking(ctx context.Context, activity model.YumiGiftChallengeA if err != nil { return nil, err } + if date != nil && date.After(dateOnly(time.Now(), location)) && header.SnapshotDueTime.After(time.Now()) { + // 改到跨日期线的新时区后,真实旧 header 的 statDate 可能短暂看似“未来”。 + // 已到原截止时间的历史周期仍可查;真正尚未到期的未来周期继续拒绝。 + return nil, NewAppError(http.StatusBadRequest, "future_stat_date", "future daily ranking is not available") + } settled := header.Status != StatusNotStarted // 与 Aslan 一致:只读取请求的原始 Top N,再过滤神秘人;被隐藏的名次不从 N 之后补位。 fetchLimit := limit @@ -170,12 +175,8 @@ func selectRankingDate(activity model.YumiGiftChallengeActivity, value string, l if err != nil { return time.Time{}, err } - if date.Before(dateOnly(activity.StartTime, location)) || !date.Before(activity.EndTime) { - return time.Time{}, NewAppError(http.StatusBadRequest, "stat_date_out_of_range", "statDate is outside activity window") - } - if date.After(dateOnly(time.Now(), location)) { - return time.Time{}, NewAppError(http.StatusBadRequest, "future_stat_date", "future daily ranking is not available") - } + // 进行中可缩短时间窗/更改时区,旧日期可能已有积分或冻结门闩。 + // 显式日期不再按可变 Start/End 拒绝;后续必须命中真实 period header 才会查榜。 return date, nil } now := time.Now().In(location) diff --git a/internal/service/yumigiftchallenge/settlement.go b/internal/service/yumigiftchallenge/settlement.go index 7b5ac87..cceaf5b 100644 --- a/internal/service/yumigiftchallenge/settlement.go +++ b/internal/service/yumigiftchallenge/settlement.go @@ -69,6 +69,15 @@ func (s *Service) resolveSettlementTarget(ctx context.Context, activityID int64, func (s *Service) freezeSettlementTarget(ctx context.Context, target settlementTarget) error { // 阶段一独立事务只负责关闭写门闩并冻结获奖 parent;阶段二失败不会重新开放榜单或重算获奖人。 return s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + // 在线配置、事件和结算统一先锁活动再锁周期;总榜结算不能先持 header 再等待 + // config 持有的 activity,否则两条链路会形成真实死锁。 + var activity model.YumiGiftChallengeActivity + if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("id = ?", target.activity.ID).First(&activity).Error; err != nil { + return err + } + if !activity.Enabled { + return NewAppError(http.StatusConflict, "activity_disabled", "activity is disabled") + } var header model.YumiGiftChallengePeriodSettlement queryErr := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where( "activity_id = ? AND period_type = ? AND period_key = ?", target.activity.ID, target.period, target.periodKey, @@ -95,7 +104,7 @@ func (s *Service) freezeSettlementTarget(ctx context.Context, target settlementT return err } } - return s.freezeSettlementParentsTx(tx, *target.activity, target.period, target.periodKey, target.statDate) + return s.freezeSettlementParentsTx(tx, activity, target.period, target.periodKey, target.statDate) }) } @@ -209,6 +218,10 @@ func (s *Service) ensureSettlementDeliveryItems(ctx context.Context, settlementI func (s *Service) refreshPeriodStatus(ctx context.Context, activityID int64, period, periodKey string) error { return s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + var activity model.YumiGiftChallengeActivity + if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("id = ?", activityID).First(&activity).Error; err != nil { + return err + } var header model.YumiGiftChallengePeriodSettlement if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where( "activity_id = ? AND period_type = ? AND period_key = ?", activityID, period, periodKey, diff --git a/internal/service/yumigiftchallenge/task.go b/internal/service/yumigiftchallenge/task.go index 9ced22b..644a8c5 100644 --- a/internal/service/yumigiftchallenge/task.go +++ b/internal/service/yumigiftchallenge/task.go @@ -31,34 +31,70 @@ func (s *Service) Enter(ctx context.Context, user AuthUser, req EnterRequest) (* if activityStatus(*activity, time.Now()) != "ONGOING" { return nil, NewAppError(http.StatusConflict, "activity_not_ongoing", "activity is not ongoing") } - location, err := resolveLocation(activity.Timezone) - if err != nil { - return nil, err - } - statDate := dateKey(time.Now(), location) + var effectiveActivity model.YumiGiftChallengeActivity + var statDate string err = s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + locked, lockedDate, err := lockOngoingActivityTx(tx, activity.ID, origin, time.Now()) + if err != nil { + return err + } + effectiveActivity, statDate = locked, lockedDate var header model.YumiGiftChallengePeriodSettlement if err := tx.Clauses(clause.Locking{Strength: "SHARE"}).Where( - "activity_id = ? AND period_type = ? AND period_key = ?", activity.ID, PeriodDaily, statDate, + "activity_id = ? AND period_type = ? AND period_key = ?", effectiveActivity.ID, PeriodDaily, statDate, ).First(&header).Error; err != nil { return err } if header.Status != StatusNotStarted { return NewAppError(http.StatusConflict, "daily_period_closed", "daily task period is already frozen") } - if err := s.ensureTaskSnapshotsTx(tx, *activity, statDate, user.UserID); err != nil { + if err := s.ensureTaskSnapshotsTx(tx, effectiveActivity, statDate, user.UserID); err != nil { return err } now := time.Now() return tx.Model(&model.YumiGiftChallengeUserTaskDaily{}). Where("activity_id = ? AND stat_date = ? AND user_id = ? AND task_type = ? AND completed_time IS NULL", - activity.ID, statDate, user.UserID, TaskEnterPage). + effectiveActivity.ID, statDate, user.UserID, TaskEnterPage). Updates(map[string]any{"progress_value": model.Decimal24_2("1.00"), "completed_time": now, "update_time": now}).Error }) if err != nil { return nil, err } - return s.buildUserState(ctx, *activity, statDate, user.UserID) + return s.buildUserState(ctx, effectiveActivity, statDate, user.UserID) +} + +// lockOngoingActivityTx 把活动时间/时区读取放到 activity SHARE 锁内。在线配置编辑会先取 +// activity UPDATE 锁再重建周期门闩,因此任务事务不会把旧时区计算出的日期写进新门闩。 +func lockOngoingActivityTx(tx *gorm.DB, activityID int64, origin string, now time.Time) (model.YumiGiftChallengeActivity, string, error) { + activity, err := lockEnabledActivityTx(tx, activityID, origin) + if err != nil { + return model.YumiGiftChallengeActivity{}, "", err + } + if activityStatus(activity, now) != "ONGOING" { + return model.YumiGiftChallengeActivity{}, "", NewAppError(http.StatusConflict, "activity_not_ongoing", "activity is not ongoing") + } + location, err := resolveLocation(activity.Timezone) + if err != nil { + return model.YumiGiftChallengeActivity{}, "", err + } + return activity, dateKey(now, location), nil +} + +// lockEnabledActivityTx 是用户任务链路的统一首锁;历史任务领取不要求当前 +// 时间窗仍 ongoing,但仍必须在锁内校验活动归属与启用状态。 +func lockEnabledActivityTx(tx *gorm.DB, activityID int64, origin string) (model.YumiGiftChallengeActivity, error) { + var activity model.YumiGiftChallengeActivity + err := tx.Clauses(clause.Locking{Strength: "SHARE"}).Where("id = ?", activityID).First(&activity).Error + if errors.Is(err, gorm.ErrRecordNotFound) { + return model.YumiGiftChallengeActivity{}, NewAppError(http.StatusNotFound, "activity_not_found", "activity not found") + } + if err != nil { + return model.YumiGiftChallengeActivity{}, err + } + if activity.SysOrigin != origin || !activity.Enabled { + return model.YumiGiftChallengeActivity{}, NewAppError(http.StatusNotFound, "activity_not_found", "activity not found") + } + return activity, nil } func (s *Service) ensureTaskSnapshotsTx(tx *gorm.DB, activity model.YumiGiftChallengeActivity, statDate string, userID int64) error { @@ -118,28 +154,50 @@ func (s *Service) ClaimTask(ctx context.Context, user AuthUser, taskCode string, if err != nil { return nil, err } - if activityStatus(*activity, time.Now()) != "ONGOING" { - return nil, NewAppError(http.StatusConflict, "activity_not_ongoing", "activity is not ongoing") - } - location, err := resolveLocation(activity.Timezone) - if err != nil { - return nil, err - } - statDate := dateKey(time.Now(), location) var ownerID int64 + var effectiveActivity model.YumiGiftChallengeActivity + var statDate string err = s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { - if err := s.ensureTaskSnapshotsTx(tx, *activity, statDate, user.UserID); err != nil { + locked, err := lockEnabledActivityTx(tx, activity.ID, origin) + if err != nil { return err } - var task model.YumiGiftChallengeUserTaskDaily - queryErr := tx.Clauses(clause.Locking{Strength: "UPDATE"}). - Where("activity_id = ? AND stat_date = ? AND user_id = ? AND task_code = ?", activity.ID, statDate, user.UserID, taskCode). - First(&task).Error - if errors.Is(queryErr, gorm.ErrRecordNotFound) { - return NewAppError(http.StatusNotFound, "task_not_found", "task not found") + effectiveActivity = locked + location, err := resolveLocation(locked.Timezone) + if err != nil { + return err } - if queryErr != nil { - return queryErr + statDate = dateKey(time.Now(), location) + task, selected, err := lockClaimRequestOwnerTx(tx, effectiveActivity.ID, user.UserID, taskCode, requestID) + if err != nil { + return err + } + if selected { + statDate = task.StatDate + } + if !selected && activityStatus(effectiveActivity, time.Now()) == "ONGOING" { + if err := s.ensureTaskSnapshotsTx(tx, effectiveActivity, statDate, user.UserID); err != nil { + return err + } + queryErr := tx.Clauses(clause.Locking{Strength: "UPDATE"}). + Where("activity_id = ? AND stat_date = ? AND user_id = ? AND task_code = ?", effectiveActivity.ID, statDate, user.UserID, taskCode). + First(&task).Error + if queryErr != nil && !errors.Is(queryErr, gorm.ErrRecordNotFound) { + return queryErr + } + processable := task.DeliveryStatus == DeliveryNotClaimed || task.DeliveryStatus == DeliveryPending || task.DeliveryStatus == DeliveryProcessing + // SUCCESS 只能由上面的相同 requestId 命中后幂等返回;UNKNOWN 必须 + // 排在所有可处理的历史 owner 之后,不能挡住其他日期的正常奖励。 + selected = queryErr == nil && task.CompletedTime != nil && task.ProgressValue.Compare(task.TargetValue) >= 0 && processable + } + if !selected { + // 时区或时间窗变更后,已完成的旧日任务仍是独立奖励 owner。 + // 先按 requestId 找原 owner 保证重试不会切到另一天,再领取最近的可处理任务。 + task, err = lockHistoricalClaimTaskTx(tx, effectiveActivity.ID, user.UserID, taskCode) + if err != nil { + return err + } + statDate = task.StatDate } ownerID = task.ID if task.CompletedTime == nil || task.ProgressValue.Compare(task.TargetValue) < 0 { @@ -161,7 +219,7 @@ func (s *Service) ClaimTask(ctx context.Context, user AuthUser, taskCode string, default: return NewAppError(http.StatusConflict, "invalid_delivery_status", "task reward state is invalid") } - return s.createDeliveryItemsTx(tx, OwnerTask, task.ID, activity.ID, user.UserID, task.ResourceGroupID) + return s.createDeliveryItemsTx(tx, OwnerTask, task.ID, effectiveActivity.ID, user.UserID, task.ResourceGroupID) }) if err != nil { return nil, err @@ -171,7 +229,55 @@ func (s *Service) ClaimTask(ctx context.Context, user AuthUser, taskCode string, return nil, err } } - return s.buildUserState(ctx, *activity, statDate, user.UserID) + return s.buildUserState(ctx, effectiveActivity, statDate, user.UserID) +} + +// lockClaimRequestOwnerTx 在选“当前日/最近日”前先恢复 requestId 已绑定的 owner。 +// 即使请求重试时当前日任务又完成了,也不会用同一 requestId 多领一天。 +func lockClaimRequestOwnerTx(tx *gorm.DB, activityID, userID int64, taskCode, requestID string) (model.YumiGiftChallengeUserTaskDaily, bool, error) { + var task model.YumiGiftChallengeUserTaskDaily + err := tx.Clauses(clause.Locking{Strength: "UPDATE"}). + Where("activity_id = ? AND user_id = ? AND task_code = ? AND claim_request_id = ? AND completed_time IS NOT NULL AND progress_value >= target_value", + activityID, userID, taskCode, requestID). + Order("stat_date DESC, id DESC").First(&task).Error + if errors.Is(err, gorm.ErrRecordNotFound) { + return model.YumiGiftChallengeUserTaskDaily{}, false, nil + } + if err != nil { + return model.YumiGiftChallengeUserTaskDaily{}, false, err + } + return task, true, nil +} + +// lockHistoricalClaimTaskTx 只在当前日任务不可领时走历史兜底。每个日任务行 +// 是独立 owner;该函数只处理尚未绑定 requestId 的新请求。 +func lockHistoricalClaimTaskTx(tx *gorm.DB, activityID, userID int64, taskCode string) (model.YumiGiftChallengeUserTaskDaily, error) { + base := func() *gorm.DB { + return tx.Clauses(clause.Locking{Strength: "UPDATE"}). + Where("activity_id = ? AND user_id = ? AND task_code = ? AND completed_time IS NOT NULL AND progress_value >= target_value", + activityID, userID, taskCode) + } + var task model.YumiGiftChallengeUserTaskDaily + err := base().Where("delivery_status IN ?", []string{DeliveryNotClaimed, DeliveryPending, DeliveryProcessing}). + Order("stat_date DESC, id DESC").First(&task).Error + if err == nil { + return task, nil + } + if !errors.Is(err, gorm.ErrRecordNotFound) { + return model.YumiGiftChallengeUserTaskDaily{}, err + } + // UNKNOWN 不能自动重发,但也不能被当成“没有奖励”后跳过;后续 switch + // 会要求人工核账。SUCCESS 只能通过上面的 requestId owner 查询返回幂等成功。 + task = model.YumiGiftChallengeUserTaskDaily{} + err = base().Where("delivery_status = ?", DeliveryUnknown). + Order("stat_date DESC, id DESC").First(&task).Error + if errors.Is(err, gorm.ErrRecordNotFound) { + return model.YumiGiftChallengeUserTaskDaily{}, NewAppError(http.StatusConflict, "task_not_completed", "task is not completed") + } + if err != nil { + return model.YumiGiftChallengeUserTaskDaily{}, err + } + return task, nil } func (s *Service) buildUserState(ctx context.Context, activity model.YumiGiftChallengeActivity, statDate string, userID int64) (*UserStateResponse, error) {