Files
hi-server/scheduler/scheduler.go
T
shanshanzhong147 07409eb602 新功能(#4): 抽奖 Stage 2 人工奖领奖工单(crypto / physical / manual_other)
Closes HIF-4

Stage 2 交付:人工奖领奖工单完整闭环。crypto / physical / manual_other 三类奖品从抽中到 mark-paid 的全流程可用。

- 迁移 02159_lottery_claim:UNIQUE(draw_id) + 3 支持索引,状态机 pending_claim→reviewing→paying→paid,rejected 可复活,超时 expired
- 3 个 PrizeHandler:Dispatch→ErrDispatchNotSupported 兜底、ClaimSchema 各自形态、ValidateClaim 表驱动
- BuildCryptoClaimSchema:抽中时按奖品 config.networks 注入 enum,前端下拉直接可用
- Draw service dispatchOrEnqueueClaim:人工奖同 tx 插 pending_claim(回滚双清),nonce 重放回读 ExpiresAt + ClaimFormSchema
- POST /claim 实装:ownership 校验 → prize 类型校验 → handler.ValidateClaim → crypto network 白名单二次校验 → tx CAS status IN (pending_claim, rejected) AND expires_at > now
- Admin CRUD 5 接口:list(IN 批拉 snap + user,无 N+1)、summary(GROUP BY 一次拿计数 + overdue 单查)、approve/reject/mark-paid 全走 CAS + audit
- Scheduler @every 1h 扫过期,级联 lottery_draw.dispatch_state → expired
- 新增错误码 100005-100011(already_submitted / invalid_claim_data / draw_not_found / not_your_draw / claim_expired / claim_state_invalid)
- Rebase 后 Stage 2 测试主动 reuse PR E 的 unmetReasonsNotEmpty + evaluatedAtNotZero matcher,人工奖分支若绕过守卫会立即挂
- Stage 1 全部 4 处 guardrail 后端 rebase 时自检过:UnmetReasons、EvaluatedAt、GrantLedger.Payload、AdminMetaMiddleware 全保留

CI 全绿;28 files, +2442/-126;覆盖率 handler 78.9% / model.lottery 74.3% / draw 68.7% / queue/lottery 76.9%
2026-07-10 04:01:34 -07:00

98 lines
3.6 KiB
Go

package scheduler
import (
"time"
"github.com/perfect-panel/server/pkg/logger"
"github.com/hibiken/asynq"
"github.com/perfect-panel/server/internal/svc"
"github.com/perfect-panel/server/queue/types"
)
type Service struct {
svc *svc.ServiceContext
server *asynq.Scheduler
}
func NewService(svc *svc.ServiceContext) *Service {
return &Service{
svc: svc,
server: initService(svc),
}
}
func (m *Service) Start() {
logger.Infof("start scheduler service")
// schedule check subscription task: every 60 seconds
checkTask := asynq.NewTask(types.SchedulerCheckSubscription, nil)
if _, err := m.server.Register("@every 60s", checkTask); err != nil {
logger.Errorf("register check subscription task failed: %s", err.Error())
}
//// schedule total server data task: every 5 minutes
//totalServerDataTask := asynq.NewTask(types.SchedulerTotalServerData, nil)
//if _, err := m.server.Register("@every 180s", totalServerDataTask); err != nil {
// logger.Errorf("register total server data task failed: %s", err.Error())
//}
// schedule reset traffic task: every day at 00:30
resetTrafficTask := asynq.NewTask(types.SchedulerResetTraffic, nil)
if _, err := m.server.Register("30 0 * * *", resetTrafficTask); err != nil {
logger.Errorf("register reset traffic task failed: %s", err.Error())
}
// schedule traffic stat task: every day at 00:00
trafficStatTask := asynq.NewTask(types.SchedulerTrafficStat, nil)
if _, err := m.server.Register("0 0 * * *", trafficStatTask, asynq.MaxRetry(3)); err != nil {
logger.Errorf("register traffic stat task failed: %s", err.Error())
}
// schedule update exchange rate task: every day at 01:00
rateTask := asynq.NewTask(types.ForthwithQuotaTask, nil)
if _, err := m.server.Register("0 1 * * *", rateTask, asynq.MaxRetry(3)); err != nil {
logger.Errorf("register update exchange rate task failed: %s", err.Error())
}
// schedule Apple IAP reconcile: every 5 minutes (第二层补单)
iapReconcileTask := asynq.NewTask(types.SchedulerIAPReconcile, nil)
if _, err := m.server.Register("@every 5m", iapReconcileTask, asynq.MaxRetry(1)); err != nil {
logger.Errorf("register iap reconcile task failed: %s", err.Error())
}
// schedule Apple IAP daily reconcile: every day at 02:00 (第三层日终对账)
iapDailyTask := asynq.NewTask(types.SchedulerIAPDailyReconcile, nil)
if _, err := m.server.Register("0 2 * * *", iapDailyTask, asynq.MaxRetry(1)); err != nil {
logger.Errorf("register iap daily reconcile task failed: %s", err.Error())
}
// schedule stuck order recovery: every 10 minutes
stuckOrderTask := asynq.NewTask(types.SchedulerStuckOrderRecovery, nil)
if _, err := m.server.Register("@every 10m", stuckOrderTask, asynq.MaxRetry(1)); err != nil {
logger.Errorf("register stuck order recovery task failed: %s", err.Error())
}
// Lottery Stage 2: 每小时扫描过期 pending_claim → expired
lotteryExpireTask := asynq.NewTask(types.SchedulerLotteryExpireClaim, nil)
if _, err := m.server.Register("@every 1h", lotteryExpireTask, asynq.MaxRetry(1)); err != nil {
logger.Errorf("register lottery expire claim task failed: %s", err.Error())
}
if err := m.server.Run(); err != nil {
logger.Errorf("run scheduler failed: %s", err.Error())
}
}
func (m *Service) Stop() {
logger.Info("stop scheduler service")
m.server.Shutdown()
}
func initService(svc *svc.ServiceContext) *asynq.Scheduler {
location, _ := time.LoadLocation("Asia/Shanghai")
return asynq.NewScheduler(
asynq.RedisClientOpt{Addr: svc.Config.Redis.Host, Password: svc.Config.Redis.Pass, DB: 5},
&asynq.SchedulerOpts{
Location: location,
},
)
}