Some checks failed
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Has been cancelled
- 新增六对成对迁移 000232–000237:H5 弹窗类型、退款结算标识与申请人备注、优先轮询事实字段与两个新终态、通道阈值命中留痕、手机号最近解绑人、提现资格校验留痕 - 退款:原因必填与申请人备注、来源支付与渠道流水冻结、线下处理流水号补录审计、按订单查询可选退款方式、企微审批材料补齐且新增字段缺失映射即明确失败 - 优先轮询:人工关闭、有效期到期独立周期任务、失败与过期人工重触发、事实字段与异常重试查询、资产解析端点只读投影 - 通道阈值:命中事实同事务留痕与命中记录查询;员工账单:列表筛选与详情投影;商户池:列表投影与统计周期语义;H5:弹窗类型与类别排序 - 手机号:有效关联数量与最近解绑人、短信验证码失败次数限制;导出:佣金明细十五列与报表序号列 - 时间筛选:三处新增筛选纳入统一严格解析契约,员工账单产生时间参数改名 - 同步 12 份主 Spec 需求、两端点与异步任务证据链,门禁 context-health 与 OpenSpec 校验通过
243 lines
9.5 KiB
Go
243 lines
9.5 KiB
Go
package task
|
||
|
||
import (
|
||
"context"
|
||
"time"
|
||
|
||
"github.com/redis/go-redis/v9"
|
||
"go.uber.org/zap"
|
||
"gorm.io/gorm"
|
||
|
||
"github.com/break/junhong_cmp_fiber/internal/infrastructure/audit"
|
||
"github.com/break/junhong_cmp_fiber/internal/infrastructure/cardtrafficlock"
|
||
"github.com/break/junhong_cmp_fiber/internal/polling"
|
||
"github.com/break/junhong_cmp_fiber/internal/store/postgres"
|
||
"github.com/break/junhong_cmp_fiber/pkg/constants"
|
||
)
|
||
|
||
// acquireConcurrencyScript 原子获取并发信号量的 Lua 脚本
|
||
// INCR + EXPIRE 合并为单个服务端操作,消除二者之间的崩溃窗口:
|
||
// 若 Worker 在 INCR 后、EXPIRE 前崩溃,key 将永久留在 Redis 导致计数器卡死。
|
||
// KEYS[1]: 全部轮询计数;KEYS[2]: 分类轮询计数
|
||
// ARGV[1]: 总量上限;ARGV[2]: 分类上限;ARGV[3]: key TTL(秒)
|
||
// 返回 -1 表示任一上限超额,>0 表示成功获取后的总计数值
|
||
var acquireConcurrencyScript = redis.NewScript(`
|
||
local total = redis.call('INCR', KEYS[1])
|
||
local kind = redis.call('INCR', KEYS[2])
|
||
if tonumber(total) > tonumber(ARGV[1]) or tonumber(kind) > tonumber(ARGV[2]) then
|
||
redis.call('DECR', KEYS[1])
|
||
redis.call('DECR', KEYS[2])
|
||
return -1
|
||
end
|
||
redis.call('EXPIRE', KEYS[1], tonumber(ARGV[3]))
|
||
redis.call('EXPIRE', KEYS[2], tonumber(ARGV[3]))
|
||
return total
|
||
`)
|
||
|
||
var releaseConcurrencyScript = redis.NewScript(`
|
||
for _, key in ipairs(KEYS) do
|
||
local current = tonumber(redis.call('GET', key) or '0') or 0
|
||
if current <= 0 then
|
||
if redis.call('EXISTS', key) == 1 then redis.call('SET', key, 0, 'KEEPTTL') end
|
||
else
|
||
redis.call('DECR', key)
|
||
end
|
||
end
|
||
return 0
|
||
`)
|
||
|
||
const pollingFallbackOperationTimeout = 5 * time.Second
|
||
|
||
// pollingFallbackContext 创建轮询兜底操作使用的独立短超时上下文。
|
||
// Asynq 任务 ctx 超时后,释放并发计数、释放锁、重入队仍必须尽量完成,
|
||
// 否则会造成并发计数短期卡住,甚至卡已出队但未重新入队。
|
||
func pollingFallbackContext() (context.Context, context.CancelFunc) {
|
||
return context.WithTimeout(context.Background(), pollingFallbackOperationTimeout)
|
||
}
|
||
|
||
// PollingBase 轮询共享基类
|
||
// 封装并发控制、卡缓存、重入队、配置间隔查询与优先轮询项认领等公共方法,所有 Handler 共享
|
||
type PollingBase struct {
|
||
redis *redis.Client
|
||
queueMgr *polling.PollingQueueManager
|
||
configMgr *polling.PollingConfigManager
|
||
iotCardStore *postgres.IotCardStore
|
||
logger *zap.Logger
|
||
verboseLog bool
|
||
totalMaxConcurrency int
|
||
trafficLock *cardtrafficlock.Lock
|
||
db *gorm.DB
|
||
priorityStore *postgres.PollingPriorityItemStore
|
||
priorityAudit *audit.Writer
|
||
scopeChecker *polling.PollingLifecycleService
|
||
}
|
||
|
||
// NewPollingBase 创建轮询共享基类
|
||
func NewPollingBase(
|
||
redisClient *redis.Client,
|
||
queueMgr *polling.PollingQueueManager,
|
||
configMgr *polling.PollingConfigManager,
|
||
iotCardStore *postgres.IotCardStore,
|
||
scopeChecker *polling.PollingLifecycleService,
|
||
priorityStore *postgres.PollingPriorityItemStore,
|
||
priorityAudit *audit.Writer,
|
||
db *gorm.DB,
|
||
logger *zap.Logger,
|
||
verboseLog bool,
|
||
totalMaxConcurrency int,
|
||
) *PollingBase {
|
||
return &PollingBase{
|
||
redis: redisClient,
|
||
queueMgr: queueMgr,
|
||
configMgr: configMgr,
|
||
iotCardStore: iotCardStore,
|
||
logger: logger,
|
||
verboseLog: verboseLog,
|
||
totalMaxConcurrency: totalMaxConcurrency,
|
||
trafficLock: cardtrafficlock.New(redisClient),
|
||
db: db,
|
||
priorityStore: priorityStore,
|
||
priorityAudit: priorityAudit,
|
||
scopeChecker: scopeChecker,
|
||
}
|
||
}
|
||
|
||
// acquireConcurrency 原子获取并发信号量
|
||
// 使用 Lua 脚本将 INCR 与 EXPIRE 合并为单个服务端操作,消除非原子窗口:
|
||
// 旧实现中若 Worker 在 INCR 后、EXPIRE 前崩溃,key 永久留在 Redis,计数器卡死。
|
||
func (b *PollingBase) acquireConcurrency(ctx context.Context, taskType string) bool {
|
||
shortType := shortTaskType(taskType)
|
||
configKey := constants.RedisPollingConcurrencyConfigKey(shortType)
|
||
currentKey := constants.RedisPollingConcurrencyCurrentKey(taskType)
|
||
totalKey := constants.RedisPollingConcurrencyTotalCurrentKey()
|
||
|
||
maxConcurrency, err := b.redis.Get(ctx, configKey).Int()
|
||
if err != nil || maxConcurrency < 1 || maxConcurrency > constants.PollingMaxConcurrencyLimit {
|
||
maxConcurrency = constants.PollingDefaultMaxConcurrency
|
||
}
|
||
totalMax := b.totalMaxConcurrency
|
||
if totalMax < 1 || totalMax > constants.PollingMaxConcurrencyLimit {
|
||
totalMax = constants.PollingDefaultTotalMaxConcurrency
|
||
}
|
||
|
||
result, err := acquireConcurrencyScript.Run(
|
||
ctx, b.redis, []string{totalKey, currentKey},
|
||
totalMax, maxConcurrency, constants.PollingConcurrencyKeyTTL,
|
||
).Int64()
|
||
if err != nil {
|
||
b.logger.Warn("获取并发计数失败,放行任务", zap.Error(err))
|
||
return true
|
||
}
|
||
|
||
if result == -1 {
|
||
b.logger.Info("轮询因并发令牌不足延后", zap.String("task_type", taskType),
|
||
zap.Int("max_concurrency", maxConcurrency), zap.Int("total_max_concurrency", totalMax),
|
||
zap.String("metric", "polling.deferred.concurrency_limit"))
|
||
return false
|
||
}
|
||
return true
|
||
}
|
||
|
||
// releaseConcurrency 释放并发信号量
|
||
func (b *PollingBase) releaseConcurrency(_ context.Context, taskType string) {
|
||
ctx, cancel := pollingFallbackContext()
|
||
defer cancel()
|
||
|
||
currentKey := constants.RedisPollingConcurrencyCurrentKey(taskType)
|
||
if err := releaseConcurrencyScript.Run(ctx, b.redis, []string{constants.RedisPollingConcurrencyTotalCurrentKey(), currentKey}).Err(); err != nil {
|
||
b.logger.Warn("释放并发计数失败", zap.String("task_type", taskType), zap.Error(err))
|
||
}
|
||
}
|
||
|
||
// acquireCardTrafficSyncLock 获取卡流量同步锁,避免轮询与手动刷新重复统计同一上游读数。
|
||
func (b *PollingBase) acquireCardTrafficSyncLock(ctx context.Context, cardID uint) (string, bool, error) {
|
||
return b.trafficLock.Acquire(ctx, cardID)
|
||
}
|
||
|
||
// releaseCardTrafficSyncLock 释放卡流量同步锁。
|
||
func (b *PollingBase) releaseCardTrafficSyncLock(_ context.Context, cardID uint, token string) {
|
||
ctx, cancel := pollingFallbackContext()
|
||
defer cancel()
|
||
|
||
if err := b.trafficLock.Release(ctx, cardID, token); err != nil {
|
||
b.logger.Warn("释放卡流量同步锁失败", zap.Uint("card_id", cardID), zap.Error(err))
|
||
}
|
||
}
|
||
|
||
// ensureMissingTask 根据数据库最新卡状态,仅在对应分片队列缺失时补入任务。
|
||
func (b *PollingBase) ensureMissingTask(ctx context.Context, cardID uint, taskType string) error {
|
||
card, err := b.iotCardStore.GetByID(ctx, cardID)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
|
||
info, ok := b.configMgr.MergedTaskIntervals(card)[taskType]
|
||
if !ok || info.Interval <= 0 {
|
||
return nil
|
||
}
|
||
|
||
added, err := b.queueMgr.EnsureQueued(ctx, cardID, taskType, time.Now())
|
||
if err != nil {
|
||
return err
|
||
}
|
||
if added {
|
||
b.logger.Info("卡状态轮询补齐缺失套餐任务",
|
||
zap.Uint("card_id", cardID), zap.String("task_type", taskType))
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// requeueCardAt 使用独立短超时上下文执行真正的 ZADD 重入队。
|
||
func (b *PollingBase) requeueCardAt(cardID uint, taskType string, nextCheckAt time.Time) error {
|
||
ctx, cancel := pollingFallbackContext()
|
||
defer cancel()
|
||
|
||
if err := b.queueMgr.Requeue(ctx, cardID, taskType, nextCheckAt); err != nil {
|
||
b.logger.Error("重入队失败,卡可能暂停轮询",
|
||
zap.Uint("card_id", cardID), zap.String("task_type", taskType), zap.Error(err))
|
||
return err
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// pollingRequeueFallbackInterval 是无法解析轮询间隔时的兜底重入队间隔。
|
||
// 防止卡或优先项因配置缺失而永久停在待执行状态。
|
||
const pollingRequeueFallbackInterval = 30 * time.Second
|
||
|
||
// pollingNextCheckAt 计算该卡该任务类型的下一次检查时间,与 requeueCard 共用同一间隔解析与兜底口径。
|
||
// 返回 ok=false 表示该任务类型没有可用轮询间隔(配置为 NULL),调用方不得入队。
|
||
func (b *PollingBase) pollingNextCheckAt(ctx context.Context, cardID uint, taskType string) (time.Time, bool) {
|
||
card, err := b.getCardWithCache(ctx, cardID)
|
||
if err != nil {
|
||
b.logger.Warn("重入队:获取卡信息失败,使用30秒兜底间隔",
|
||
zap.Uint("card_id", cardID), zap.String("task_type", taskType), zap.Error(err))
|
||
return time.Now().Add(pollingRequeueFallbackInterval), true
|
||
}
|
||
intervals := b.configMgr.MergedTaskIntervals(card)
|
||
if len(intervals) == 0 {
|
||
// 兜底:配置未加载时延迟 30 秒重入队,防止卡从轮询永久消失
|
||
b.logger.Warn("无匹配轮询配置,30秒后重入队(配置可能未加载)",
|
||
zap.Uint("card_id", cardID), zap.String("task_type", taskType))
|
||
return time.Now().Add(pollingRequeueFallbackInterval), true
|
||
}
|
||
info, ok := intervals[taskType]
|
||
if !ok || info.Interval <= 0 {
|
||
b.logger.Debug("轮询间隔为 NULL,不入队", zap.Uint("card_id", cardID), zap.String("task_type", taskType))
|
||
return time.Time{}, false
|
||
}
|
||
return time.Now().Add(time.Duration(info.Interval) * time.Second), true
|
||
}
|
||
|
||
// requeueCard 将卡按匹配配置间隔重新入队分片 Sorted Set
|
||
// ⚠️ 关键:Lua 脚本原子出队后卡已从队列移除,若并发满时直接 return 会导致卡永久丢失
|
||
// 调用方必须在 acquireConcurrency 返回 false 时调用此方法入队后再返回
|
||
func (b *PollingBase) requeueCard(_ context.Context, cardID uint, taskType string) error {
|
||
ctx, cancel := pollingFallbackContext()
|
||
nextCheckAt, ok := b.pollingNextCheckAt(ctx, cardID, taskType)
|
||
cancel()
|
||
if !ok {
|
||
return nil
|
||
}
|
||
return b.requeueCardAt(cardID, taskType, nextCheckAt)
|
||
}
|