2026-07-20 19:08:59 +08:00

281 lines
11 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

package mysql
import (
"context"
"database/sql"
"errors"
"strconv"
"strings"
"time"
"hyapp/pkg/appcode"
"hyapp/pkg/xerr"
)
const (
pointWithdrawalChannelCoinSeller = "coin_seller"
pointWithdrawalChannelPlatform = "platform"
withdrawalLimitPeriodDay = "day"
withdrawalLimitPeriodWeek = "week"
withdrawalLimitPeriodMonth = "month"
)
type pointWithdrawalPolicyLimit struct {
PolicyID uint64
PolicyVersion uint64
Period string
Limit int64
PeriodKey string
AllowedDays string
}
// consumePointWithdrawalLimit 在真实账变事务内解析当前 UTC 月绑定政策并锁定一条计数主键。
// 调用方必须先完成 wallet transaction 幂等命中检查;这样重放直接返回旧回执,不会重复消耗次数。
func (r *Repository) consumePointWithdrawalLimit(ctx context.Context, tx *sql.Tx, userID int64, regionID int64, channel string, commandID string, withdrawalRef string, nowMS int64) error {
if userID <= 0 || regionID <= 0 || (channel != pointWithdrawalChannelCoinSeller && channel != pointWithdrawalChannelPlatform) {
return xerr.New(xerr.InvalidArgument, "point withdrawal limit target is invalid")
}
limit, err := r.resolvePointWithdrawalPolicyLimit(ctx, tx, regionID, channel, nowMS)
if err != nil {
return err
}
if channel == pointWithdrawalChannelPlatform && strings.TrimSpace(withdrawalRef) == "" {
// 平台先冻结再建申请;没有稳定引用就无法在建单失败时精确撤销这一次 reservation。
return xerr.New(xerr.InvalidArgument, "withdrawal_ref is required for platform withdrawal")
}
if channel == pointWithdrawalChannelPlatform {
allowed, err := platformWithdrawalAllowedOnUTCDate(limit.AllowedDays, time.UnixMilli(nowMS).UTC())
if err != nil {
return err
}
if !allowed {
return xerr.NewWithMetadata(xerr.PointWithdrawalDateNotAllowed, "point withdrawal date is not allowed", map[string]string{
"allowed_days": limit.AllowedDays,
"current_day": strconv.Itoa(time.UnixMilli(nowMS).UTC().Day()),
"timezone": "UTC",
})
}
}
appCode := appcode.FromContext(ctx)
if limit.Limit > 0 {
if _, err := tx.ExecContext(ctx, `
INSERT INTO point_withdrawal_policy_counters (
app_code, user_id, channel, period, period_key,
withdrawal_count, created_at_ms, updated_at_ms
) VALUES (?, ?, ?, ?, ?, 0, ?, ?)
ON DUPLICATE KEY UPDATE period_key = VALUES(period_key)`,
appCode, userID, channel, limit.Period, limit.PeriodKey, nowMS, nowMS,
); err != nil {
return err
}
var current int64
if err := tx.QueryRowContext(ctx, `
SELECT withdrawal_count
FROM point_withdrawal_policy_counters
WHERE app_code = ? AND user_id = ? AND channel = ? AND period = ? AND period_key = ?
FOR UPDATE`, appCode, userID, channel, limit.Period, limit.PeriodKey,
).Scan(&current); err != nil {
return err
}
if current >= limit.Limit {
return xerr.NewWithMetadata(xerr.PointWithdrawalLimitReached, "point withdrawal limit reached", map[string]string{
"channel": channel,
"period": limit.Period,
"limit": strconv.FormatInt(limit.Limit, 10),
})
}
if _, err := tx.ExecContext(ctx, `
UPDATE point_withdrawal_policy_counters
SET withdrawal_count = withdrawal_count + 1, updated_at_ms = ?
WHERE app_code = ? AND user_id = ? AND channel = ? AND period = ? AND period_key = ?`,
nowMS, appCode, userID, channel, limit.Period, limit.PeriodKey,
); err != nil {
return err
}
}
if channel != pointWithdrawalChannelPlatform {
return nil
}
// reservation 与冻结账变同事务提交;后续只允许“申请创建失败”按 withdrawal_ref 撤销,人工驳回仍计一次发起。
_, err = tx.ExecContext(ctx, `
INSERT INTO point_withdrawal_limit_reservations (
app_code, withdrawal_ref, freeze_command_id, user_id, channel,
policy_id, policy_version, period, period_key, counted, status, created_at_ms, released_at_ms
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 'active', ?, 0)`,
appCode, strings.TrimSpace(withdrawalRef), strings.TrimSpace(commandID), userID, channel,
limit.PolicyID, limit.PolicyVersion, limit.Period, limit.PeriodKey, limit.Limit > 0, nowMS,
)
if isMySQLDuplicateError(err) {
return xerr.New(xerr.IdempotencyConflict, "point withdrawal reservation already belongs to another command")
}
return err
}
func (r *Repository) resolvePointWithdrawalPolicyLimit(ctx context.Context, tx *sql.Tx, regionID int64, channel string, nowMS int64) (pointWithdrawalPolicyLimit, error) {
if nowMS <= 0 {
nowMS = time.Now().UTC().UnixMilli()
}
now := time.UnixMilli(nowMS).UTC()
cycleKey := now.Format("2006-01")
// 复用月工资政策的不可变周期绑定解析;当前月缺失时沿用同一套“只复制紧邻上月”规则,限制与发薪不会选到不同版本。
policy, found, err := r.queryHostSalaryPolicyByCycle(ctx, tx, regionID, cycleKey, "", "")
if err != nil {
return pointWithdrawalPolicyLimit{}, err
}
if !found {
return pointWithdrawalPolicyLimit{}, xerr.New(xerr.PermissionDenied, "host salary policy is not configured for withdrawal region")
}
limit := pointWithdrawalPolicyLimit{PolicyID: policy.PolicyID, PolicyVersion: policy.PolicyVersion}
switch channel {
case pointWithdrawalChannelCoinSeller:
limit.Period, limit.Limit = policy.CoinSellerWithdrawalLimitPeriod, policy.CoinSellerWithdrawalLimitCount
case pointWithdrawalChannelPlatform:
limit.Period, limit.Limit = policy.PlatformWithdrawalLimitPeriod, policy.PlatformWithdrawalLimitCount
limit.AllowedDays = strings.TrimSpace(policy.PlatformWithdrawalAllowedDays)
default:
return pointWithdrawalPolicyLimit{}, xerr.New(xerr.InvalidArgument, "point withdrawal channel is invalid")
}
limit.Period = strings.ToLower(strings.TrimSpace(limit.Period))
if limit.Limit < 0 {
return pointWithdrawalPolicyLimit{}, xerr.New(xerr.Internal, "point withdrawal policy limit is invalid")
}
if limit.Period == "" {
limit.Period = withdrawalLimitPeriodMonth
}
switch limit.Period {
case withdrawalLimitPeriodDay:
limit.PeriodKey = now.Format("2006-01-02")
case withdrawalLimitPeriodWeek:
// 周期固定为 UTC 周一 00:00 起;保存周一日期避免跨年 ISO 周编号歧义。
mondayOffset := (int(now.Weekday()) + 6) % 7
limit.PeriodKey = now.AddDate(0, 0, -mondayOffset).Format("2006-01-02")
case withdrawalLimitPeriodMonth:
limit.PeriodKey = cycleKey
default:
return pointWithdrawalPolicyLimit{}, xerr.New(xerr.Internal, "point withdrawal policy period is invalid")
}
return limit, nil
}
// platformWithdrawalAllowedOnUTCDate 精确匹配 UTC 日历日;例如 30 号在二月不存在时不会自动改成月底。
// 空串兼容迁移前政策为每天可提现,异常快照则返回内部错误,不能静默放开限制。
func platformWithdrawalAllowedOnUTCDate(raw string, now time.Time) (bool, error) {
raw = strings.TrimSpace(raw)
if raw == "" {
return true, nil
}
currentDay := now.UTC().Day()
seen := make(map[int]struct{})
allowed := false
for _, part := range strings.Split(raw, ",") {
day, err := strconv.Atoi(strings.TrimSpace(part))
if err != nil || day < 1 || day > 31 {
return false, xerr.New(xerr.Internal, "point withdrawal allowed days policy is invalid")
}
if _, duplicate := seen[day]; duplicate {
return false, xerr.New(xerr.Internal, "point withdrawal allowed days policy is invalid")
}
seen[day] = struct{}{}
if day == currentDay {
allowed = true
}
}
return allowed, nil
}
// releasePointWithdrawalLimitReservation 只服务 gateway 建单失败补偿。
// application_id 非空的人工驳回不会调用这里,因此驳回申请仍占用原限制周期的一次次数。
func (r *Repository) releasePointWithdrawalLimitReservation(ctx context.Context, tx *sql.Tx, userID int64, withdrawalRef string, rollbackCommandID string, nowMS int64) error {
withdrawalRef = strings.TrimSpace(withdrawalRef)
appCode := appcode.FromContext(ctx)
var reservedUserID int64
var channel string
var period string
var periodKey string
var counted bool
var status string
queryByRef := withdrawalRef != ""
var err error
if queryByRef {
err = tx.QueryRowContext(ctx, `
SELECT withdrawal_ref, user_id, channel, period, period_key, counted, status
FROM point_withdrawal_limit_reservations
WHERE app_code = ? AND withdrawal_ref = ?
FOR UPDATE`, appCode, withdrawalRef,
).Scan(&withdrawalRef, &reservedUserID, &channel, &period, &periodKey, &counted, &status)
}
if !queryByRef || errors.Is(err, sql.ErrNoRows) {
// HTTP 重试会生成新的 request_id但 freeze command_id 是客户端稳定幂等键;按 :rollback 约定回退到唯一命令索引。
freezeCommandID := strings.TrimSuffix(strings.TrimSpace(rollbackCommandID), ":rollback")
if freezeCommandID == strings.TrimSpace(rollbackCommandID) || freezeCommandID == "" {
if errors.Is(err, sql.ErrNoRows) || !queryByRef {
return nil
}
return err
}
err = tx.QueryRowContext(ctx, `
SELECT withdrawal_ref, user_id, channel, period, period_key, counted, status
FROM point_withdrawal_limit_reservations
WHERE app_code = ? AND freeze_command_id = ?
FOR UPDATE`, appCode, freezeCommandID,
).Scan(&withdrawalRef, &reservedUserID, &channel, &period, &periodKey, &counted, &status)
}
if errors.Is(err, sql.ErrNoRows) {
// 旧冻结记录没有 reservation释放余额仍按原兼容行为继续。
return nil
}
if err != nil {
return err
}
if reservedUserID != userID || channel != pointWithdrawalChannelPlatform {
return xerr.New(xerr.IdempotencyConflict, "point withdrawal reservation does not match rollback")
}
if status == "released" {
return nil
}
if status != "active" {
return xerr.New(xerr.Internal, "point withdrawal reservation status is invalid")
}
if counted {
// policy_id/version 只保留为命中政策审计;自然周期计数跨区域和同月显式政策切换连续,释放也必须命中共享周期键。
result, err := tx.ExecContext(ctx, `
UPDATE point_withdrawal_policy_counters
SET withdrawal_count = withdrawal_count - 1, updated_at_ms = ?
WHERE app_code = ? AND user_id = ? AND channel = ? AND period = ? AND period_key = ?
AND withdrawal_count > 0`, nowMS, appCode, userID, channel, period, periodKey)
if err != nil {
return err
}
if affected, err := result.RowsAffected(); err != nil || affected != 1 {
if err != nil {
return err
}
return xerr.New(xerr.Internal, "point withdrawal reservation counter is missing")
}
}
_, err = tx.ExecContext(ctx, `
UPDATE point_withdrawal_limit_reservations
SET status = 'released', released_at_ms = ?
WHERE app_code = ? AND withdrawal_ref = ? AND status = 'active'`, nowMS, appCode, withdrawalRef)
return err
}
// pointWithdrawalReservationWasReleased 防止已补偿冻结的 command_id 重放旧成功回执。
// reservation 与冻结交易同事务落库,因此 active/不存在仍可按原幂等语义返回released 必须要求客户端换新命令。
func (r *Repository) pointWithdrawalReservationWasReleased(ctx context.Context, tx *sql.Tx, commandID string) (bool, error) {
var status string
err := tx.QueryRowContext(ctx, `
SELECT status
FROM point_withdrawal_limit_reservations
WHERE app_code = ? AND freeze_command_id = ?`, appcode.FromContext(ctx), strings.TrimSpace(commandID)).Scan(&status)
if errors.Is(err, sql.ErrNoRows) {
return false, nil
}
if err != nil {
return false, err
}
return status == "released", nil
}