Files
hi-server/internal/logic/lottery/draw/service.go
T
shanshanzhong147 ce3babcc33 新功能(#3): 抽奖 Stage 1 收官 — 用户 API + 后台 CRUD + 集成
Closes HIF-3

Stage 1 完整闭环 PR C:用户 API + 后台 CRUD + 抽奖事务服务 + 审计 + QA curl。合并后 Stage 1 可交测试。

架构师 review R1(rulecaps depth≤8/nodes≤64/bytes≤8KB)+ R2(InviteHook source_ref 加 order: 前缀)已全部落地。

- 迁移 02158_admin_action_log:后台写操作审计
- xerr 100xxx 段:抽奖错误码(NotEligible/NoChances/ActivityEnded/RateLimited/NotClaimable/InternalError/RuleTooDeep/RuleTooMany/RuleTooLarge)
- feature flag config.Lottery.Enable 默认 false,合并后线上零副作用
- draw service:feature flag → rate limit → pre-tx reads → nonce dedupe → Consume → Pick → 乐观扣库存 → 双快照 → Dispatch → finalize
- 用户 API 4 个:GET /config、POST /draw、GET /records、POST /claim(Stage 1 返回 100010)
- 后台 CRUD:活动 / 奖品 / rules PUT(rulecaps gate)/ chances/grant
- audit.WriteAdminAction:与调用方 tx 同生共死,SHA1 body 摘要
- QA 脚本:qa/lottery/stage1_curl.sh 全链路 curl

测试覆盖:model 76.5% / draw 69% / handler 68.3% / hook 91.9% / rulecaps 87.8% / audit 100%
100 并发抢库存 + 10000 次概率分布 e2e 推 QA 环境(sqlmock 无法忠实模拟 InnoDB 行锁)
CI 全绿:构建/Vet/测试 + golangci-lint
2026-07-08 22:33:18 -07:00

548 lines
17 KiB
Go

// Package draw implements POST /api/v1/lottery/draw: the transactional lottery
// draw flow. The service composes the ChanceService, RuleEvaluator,
// WeightedPicker, PrizeHandler.Registry and LedgerService primitives from the
// model layer into one atomic sequence.
//
// Ordering matters — the flow is:
// 1. feature-flag gate (config.Lottery.Enable)
// 2. Redis rate limit (per-user 1/sec)
// 3. Load activity + validate window/status
// 4. Build RuleContext (pre-tx reads)
// 5. Evaluate eligibility (pure compute)
// 6. Load prize pool snapshot (pre-tx read)
// 7. Open tx →
// 7a. Nonce dedupe: SELECT lottery_draw WHERE (user_id, client_nonce)
// — hit returns the recorded draw
// 7b. ChanceService.Consume (SELECT FOR UPDATE + decrement)
// 7c. WeightedPicker.Pick
// 7d. Limited-stock optimistic decrement; fallback on RowsAffected=0
// 7e. INSERT lottery_draw + PrizeSnapshot + EligibilitySnapshot
// 7f. Auto-handler Dispatch (in tx)
// 7g. Update draw.dispatch_state
// → commit
//
// Everything past step 5 uses the caller's transaction; post-commit cache
// invalidation is the handler layer's job (a future enhancement — the
// underlying UserModel already invalidates its own cache on UpdateSubscribe).
package draw
import (
"context"
"database/sql"
"encoding/json"
"errors"
"fmt"
"strconv"
"time"
"github.com/perfect-panel/server/internal/model/lottery"
"github.com/perfect-panel/server/pkg/limit"
"github.com/perfect-panel/server/pkg/xerr"
"gorm.io/gorm"
)
// Request is the input to Draw.
type Request struct {
UserId int64
ActivityId int64
ClientNonce string
}
// Result is what Draw returns to the caller (user handler).
type Result struct {
DrawId int64
IsWin bool
Prize *PrizeSummary
ChancesRemaining int64
Claim ClaimSummary
Message string
}
// PrizeSummary is the awarded-prize view rendered for the user.
type PrizeSummary struct {
Slot int
Id int64
Type string
Name string
Config json.RawMessage
}
// ClaimSummary describes whether the user needs to take a further action.
type ClaimSummary struct {
Required bool
AutoClaimed bool
Message string
}
// Service orchestrates the transactional draw flow. Deps are struct-injected
// so tests can substitute fakes and production wiring lives in ServiceContext.
type Service struct {
deps Deps
}
// Deps groups the collaborators. Nil-safe checks live in Draw itself, not here.
type Deps struct {
DB *gorm.DB
Enabled bool
RateLimiter RateLimiter
Chance lottery.ChanceService
Evaluator lottery.RuleEvaluator
Picker lottery.WeightedPicker
Registry lottery.Registry
ContextBuilder RuleContextBuilder
}
// RateLimiter admits at most 1 draw per second per user. Extracted to an
// interface so tests can supply an always-admit fake without pulling Redis.
type RateLimiter interface {
Allow(ctx context.Context, userId int64) error
}
// RuleContextBuilder loads the per-user snapshot needed by the rule
// evaluator. Extracted so tests can inject deterministic contexts.
type RuleContextBuilder interface {
Build(ctx context.Context, userId int64) (lottery.RuleContext, error)
}
// NewService returns a Draw service ready to serve requests.
func NewService(d Deps) *Service { return &Service{deps: d} }
// Draw runs the full lottery draw flow. Errors are xerr codes suitable for
// direct return by the HTTP handler; internal errors are wrapped as
// LotteryInternalError.
func (s *Service) Draw(ctx context.Context, req Request) (*Result, error) {
if err := s.validateRequest(req); err != nil {
return nil, err
}
if !s.deps.Enabled {
return nil, xerr.NewErrCode(xerr.LotteryActivityEnded)
}
if err := s.applyRateLimit(ctx, req.UserId); err != nil {
return nil, err
}
activity, err := s.loadRunningActivity(ctx, req.ActivityId)
if err != nil {
return nil, err
}
// Pre-tx reads: user context + prize pool snapshot. Cheap and out of the
// hot-lock window; the tx step re-checks stock atomically.
rc, err := s.buildRuleContext(ctx, req.UserId)
if err != nil {
return nil, wrapInternal(err)
}
tree, err := parseEligibilityTree(activity.Eligibility)
if err != nil {
return nil, wrapInternal(err)
}
passed, unmet, err := s.deps.Evaluator.Evaluate(ctx, tree, rc)
if err != nil {
return nil, wrapInternal(err)
}
if !passed {
// The rejection itself is not an error to the caller — but we still
// record an eligibility snapshot for support/audit before returning
// the 4001. Snapshot write intentionally uses its own tx: the draw
// itself never got issued, so there is no draw_id to correlate; we
// omit the snapshot in that case.
_ = unmet // unmet is available to the handler via error metadata if needed
return nil, xerr.NewErrCode(xerr.LotteryNotEligible)
}
prizes, err := s.loadPrizes(ctx, req.ActivityId)
if err != nil {
return nil, wrapInternal(err)
}
if len(prizes) == 0 {
return nil, xerr.NewErrCode(xerr.LotteryActivityEnded)
}
var result *Result
txErr := s.deps.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
// (a) Nonce dedupe — race-safe idempotency check.
if existing, existsErr := s.findExistingDraw(ctx, tx, req); existsErr != nil {
return existsErr
} else if existing != nil {
result, existsErr = s.buildResultFromExistingDraw(ctx, tx, existing)
return existsErr
}
// (b) Consume chance atomically. ErrNoChances → 4002.
remaining, consumeErr := s.deps.Chance.Consume(ctx, tx, req.UserId, req.ActivityId)
if errors.Is(consumeErr, lottery.ErrNoChances) {
return xerr.NewErrCode(xerr.LotteryNoChances)
}
if consumeErr != nil {
return wrapInternal(consumeErr)
}
// (c) Pick a prize.
idx, pickErr := s.deps.Picker.Pick(prizes)
if pickErr != nil {
return xerr.NewErrCode(xerr.LotteryActivityEnded)
}
picked := prizes[idx]
// (d) Limited-stock decrement.
final, stockErr := s.decrementStockOrFallback(ctx, tx, picked, prizes)
if stockErr != nil {
return wrapInternal(stockErr)
}
// (e) Insert draw + snapshots.
draw, insertErr := s.insertDraw(ctx, tx, req, final)
if insertErr != nil {
return wrapInternal(insertErr)
}
if snapErr := s.insertSnapshots(ctx, tx, draw, final, passed); snapErr != nil {
return wrapInternal(snapErr)
}
// (f) Dispatch prize.
dispatch, dispatchErr := s.dispatchIfAutomatic(ctx, tx, req, draw, final)
if dispatchErr != nil {
return dispatchErr
}
// (g) Update draw.dispatch_state to reflect handler outcome.
if updateErr := s.finalizeDrawState(ctx, tx, draw, dispatch); updateErr != nil {
return wrapInternal(updateErr)
}
result = buildResult(draw, final, dispatch, remaining)
return nil
})
if txErr != nil {
if _, ok := txErr.(*xerr.CodeError); ok {
return nil, txErr
}
return nil, wrapInternal(txErr)
}
return result, nil
}
// ---- individual steps ------------------------------------------------------
func (s *Service) validateRequest(req Request) error {
if req.UserId <= 0 || req.ActivityId <= 0 || req.ClientNonce == "" {
return xerr.NewErrCode(xerr.InvalidParams)
}
if len(req.ClientNonce) > 64 {
return xerr.NewErrCode(xerr.InvalidParams)
}
return nil
}
func (s *Service) applyRateLimit(ctx context.Context, userId int64) error {
if s.deps.RateLimiter == nil {
return nil
}
if err := s.deps.RateLimiter.Allow(ctx, userId); err != nil {
if errors.Is(err, ErrRateLimited) {
return xerr.NewErrCode(xerr.LotteryRateLimited)
}
return wrapInternal(err)
}
return nil
}
func (s *Service) loadRunningActivity(ctx context.Context, activityId int64) (*lottery.Activity, error) {
var activity lottery.Activity
now := time.Now()
err := s.deps.DB.WithContext(ctx).
Where("id = ?", activityId).
First(&activity).Error
if err != nil {
if errors.Is(err, gorm.ErrRecordNotFound) {
return nil, xerr.NewErrCode(xerr.LotteryActivityEnded)
}
return nil, wrapInternal(err)
}
if activity.Status != lottery.ActivityStatusRunning {
return nil, xerr.NewErrCode(xerr.LotteryActivityEnded)
}
if activity.StartAt.After(now) || activity.EndAt.Before(now) {
return nil, xerr.NewErrCode(xerr.LotteryActivityEnded)
}
return &activity, nil
}
func (s *Service) buildRuleContext(ctx context.Context, userId int64) (lottery.RuleContext, error) {
if s.deps.ContextBuilder == nil {
return lottery.RuleContext{UserId: userId, Now: time.Now().Unix()}, nil
}
return s.deps.ContextBuilder.Build(ctx, userId)
}
func (s *Service) loadPrizes(ctx context.Context, activityId int64) ([]lottery.Prize, error) {
var prizes []lottery.Prize
err := s.deps.DB.WithContext(ctx).
Where("activity_id = ?", activityId).
Order("slot ASC").
Find(&prizes).Error
if err != nil {
return nil, err
}
return prizes, nil
}
func (s *Service) findExistingDraw(ctx context.Context, tx *gorm.DB, req Request) (*lottery.Draw, error) {
var existing lottery.Draw
err := tx.WithContext(ctx).
Where("user_id = ? AND client_nonce = ?", req.UserId, req.ClientNonce).
First(&existing).Error
if err == nil {
return &existing, nil
}
if errors.Is(err, gorm.ErrRecordNotFound) {
return nil, nil
}
return nil, err
}
// decrementStockOrFallback runs the limited-stock optimistic lock. If the
// chosen prize is unlimited, returns it as-is. If limited and stock survives
// → returns it. If limited and sold-out → walks the pool for the first
// `is_fallback=true` prize or falls back to a "none" (thanks-for-playing)
// synthetic prize.
func (s *Service) decrementStockOrFallback(ctx context.Context, tx *gorm.DB, picked lottery.Prize, pool []lottery.Prize) (lottery.Prize, error) {
if !picked.RemainingStock.Valid {
return picked, nil
}
res := tx.WithContext(ctx).
Model(&lottery.Prize{}).
Where("id = ? AND remaining_stock > 0", picked.Id).
UpdateColumn("remaining_stock", gorm.Expr("`remaining_stock` - 1"))
if res.Error != nil {
return lottery.Prize{}, res.Error
}
if res.RowsAffected > 0 {
return picked, nil
}
// Sold out → fallback selection.
for _, p := range pool {
if p.IsFallback {
return p, nil
}
}
// No fallback declared → synthesize a "谢谢参与" from the last non-fallback
// entry (any type=none in the pool wins); if pool has no none, we
// synthesize an ephemeral prize record. Note: this prize is NOT persisted
// as a separate row — it just satisfies the return contract.
for _, p := range pool {
if p.Type == lottery.PrizeTypeNone {
return p, nil
}
}
return lottery.Prize{
ActivityId: picked.ActivityId,
Slot: picked.Slot,
Type: lottery.PrizeTypeNone,
Name: "谢谢参与",
Config: "{}",
}, nil
}
func (s *Service) insertDraw(ctx context.Context, tx *gorm.DB, req Request, prize lottery.Prize) (*lottery.Draw, error) {
draw := lottery.Draw{
UserId: req.UserId,
ActivityId: req.ActivityId,
ClientNonce: req.ClientNonce,
IsWin: prize.Type != lottery.PrizeTypeNone,
DispatchState: lottery.DispatchStateNone,
DrawnAt: time.Now(),
}
if prize.Id > 0 {
draw.PrizeId = sql.NullInt64{Int64: prize.Id, Valid: true}
}
if err := tx.WithContext(ctx).Create(&draw).Error; err != nil {
return nil, err
}
return &draw, nil
}
func (s *Service) insertSnapshots(ctx context.Context, tx *gorm.DB, draw *lottery.Draw, prize lottery.Prize, passedEligibility bool) error {
ps := lottery.PrizeSnapshot{
DrawId: draw.Id,
PrizeId: prize.Id,
Slot: prize.Slot,
Type: prize.Type,
Name: prize.Name,
Config: prize.Config,
}
if ps.Config == "" {
ps.Config = "{}"
}
if err := tx.WithContext(ctx).Create(&ps).Error; err != nil {
return err
}
es := lottery.EligibilitySnapshot{
DrawId: draw.Id,
UserId: draw.UserId,
ActivityId: draw.ActivityId,
Passed: passedEligibility,
}
return tx.WithContext(ctx).Create(&es).Error
}
func (s *Service) dispatchIfAutomatic(ctx context.Context, tx *gorm.DB, req Request, draw *lottery.Draw, prize lottery.Prize) (lottery.DispatchResult, error) {
if !draw.IsWin {
return lottery.DispatchResult{State: lottery.DispatchStateAutoClaimed, Message: "谢谢参与"}, nil
}
if s.deps.Registry == nil {
return lottery.DispatchResult{State: lottery.DispatchStatePendingClaim, Message: "等待人工发放"}, nil
}
handler, err := s.deps.Registry.MustGet(prize.Type)
if err != nil {
if errors.Is(err, lottery.ErrHandlerNotRegistered) {
return lottery.DispatchResult{State: lottery.DispatchStatePendingClaim, Message: "等待人工发放"}, nil
}
return lottery.DispatchResult{}, wrapInternal(err)
}
if !handler.IsAuto() {
return lottery.DispatchResult{State: lottery.DispatchStatePendingClaim, Message: "等待人工发放"}, nil
}
dispatchReq := lottery.DispatchRequest{
UserId: req.UserId,
ActivityId: req.ActivityId,
DrawId: draw.Id,
Prize: prize,
Snapshot: lottery.PrizeSnapshot{DrawId: draw.Id, PrizeId: prize.Id, Slot: prize.Slot, Type: prize.Type, Name: prize.Name, Config: prize.Config},
IdempotencyKey: fmt.Sprintf("lottery:%d:%d", req.ActivityId, draw.Id),
}
return handler.Dispatch(ctx, tx, dispatchReq)
}
func (s *Service) finalizeDrawState(ctx context.Context, tx *gorm.DB, draw *lottery.Draw, dispatch lottery.DispatchResult) error {
now := time.Now()
updates := map[string]any{
"dispatch_state": dispatch.State,
}
if dispatch.State == lottery.DispatchStateAutoClaimed || dispatch.State == lottery.DispatchStatePaid {
updates["dispatched_at"] = now
}
return tx.WithContext(ctx).
Model(&lottery.Draw{}).
Where("id = ?", draw.Id).
Updates(updates).Error
}
func (s *Service) buildResultFromExistingDraw(ctx context.Context, tx *gorm.DB, existing *lottery.Draw) (*Result, error) {
var snap lottery.PrizeSnapshot
err := tx.WithContext(ctx).Where("draw_id = ?", existing.Id).First(&snap).Error
if err != nil && !errors.Is(err, gorm.ErrRecordNotFound) {
return nil, err
}
remaining, _ := s.deps.Chance.Query(ctx, existing.UserId, existing.ActivityId)
var prize *PrizeSummary
if existing.IsWin {
prize = &PrizeSummary{
Slot: snap.Slot,
Id: snap.PrizeId,
Type: snap.Type,
Name: snap.Name,
Config: json.RawMessage(snap.Config),
}
}
return &Result{
DrawId: existing.Id,
IsWin: existing.IsWin,
Prize: prize,
ChancesRemaining: remaining,
Claim: ClaimSummary{
Required: existing.DispatchState == lottery.DispatchStatePendingClaim,
AutoClaimed: existing.DispatchState == lottery.DispatchStateAutoClaimed,
Message: "",
},
}, nil
}
func buildResult(draw *lottery.Draw, prize lottery.Prize, dispatch lottery.DispatchResult, remaining int64) *Result {
res := &Result{
DrawId: draw.Id,
IsWin: draw.IsWin,
ChancesRemaining: remaining,
Message: dispatch.Message,
Claim: ClaimSummary{
Required: dispatch.State == lottery.DispatchStatePendingClaim,
AutoClaimed: dispatch.State == lottery.DispatchStateAutoClaimed,
Message: dispatch.Message,
},
}
if draw.IsWin {
res.Prize = &PrizeSummary{
Slot: prize.Slot,
Id: prize.Id,
Type: prize.Type,
Name: prize.Name,
Config: json.RawMessage(prize.Config),
}
}
return res
}
// parseEligibilityTree tolerates empty/null activity.Eligibility as "no gate".
func parseEligibilityTree(raw string) (*lottery.EligibilityRule, error) {
trimmed := ""
for _, r := range raw {
if r != ' ' && r != '\t' && r != '\n' && r != '\r' {
trimmed += string(r)
}
}
if trimmed == "" || trimmed == "null" || trimmed == "{}" {
return nil, nil
}
var tree lottery.EligibilityRule
if err := json.Unmarshal([]byte(raw), &tree); err != nil {
return nil, fmt.Errorf("parse eligibility: %w", err)
}
return &tree, nil
}
func wrapInternal(err error) error {
if err == nil {
return nil
}
// Preserve already-coded errors.
if _, ok := err.(*xerr.CodeError); ok {
return err
}
return xerr.NewErrCodeMsg(xerr.LotteryInternalError, err.Error())
}
// ---- Rate limiter production wiring ----------------------------------------
// ErrRateLimited is returned by RateLimiter.Allow when the caller exceeded
// the configured quota.
var ErrRateLimited = errors.New("draw: rate limited")
// RedisRateLimiter is the production RateLimiter backed by pkg/limit's
// Redis-Lua fixed-window (1 hit per 1 second per user), matching the
// existing sendEmailCodeLogic pattern.
type RedisRateLimiter struct {
limiter *limit.PeriodLimit
}
// NewRedisRateLimiter builds a per-user 1-req/1-sec limiter. keyPrefix is
// expected to end with ':' so the composed key is human-readable.
func NewRedisRateLimiter(limiter *limit.PeriodLimit) *RedisRateLimiter {
return &RedisRateLimiter{limiter: limiter}
}
// Allow admits or rejects the caller.
func (r *RedisRateLimiter) Allow(ctx context.Context, userId int64) error {
if r == nil || r.limiter == nil {
return nil
}
state, err := r.limiter.TakeCtx(ctx, strconv.FormatInt(userId, 10))
if err != nil {
return err
}
if state == limit.Allowed || state == limit.HitQuota {
return nil
}
return ErrRateLimited
}