hyapp-server/services/wallet-service/internal/storage/mysql/resource_equipment_cleanup.go
2026-07-23 20:35:25 +08:00

334 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"
"fmt"
"math/rand/v2"
"time"
mysqlDriver "github.com/go-sql-driver/mysql"
"hyapp/pkg/xerr"
)
const (
resourceEquipmentCleanupJobName = "resource_equipment_cleanup_v1"
resourceEquipmentCleanupLockName = "wallet_resource_equipment_cleanup:v1"
)
// ResourceEquipmentCleanupPageResult 描述一个已经提交的有界清理页。
// ScannedCount 统计主键页读取量DeletedCount 只统计复核后实际删除的装备指针。
type ResourceEquipmentCleanupPageResult struct {
ScannedCount int
DeletedCount int
ReachedEnd bool
DeadlockRetries int
}
type resourceEquipmentCleanupCursor struct {
AppCode string
UserID int64
ResourceType string
EntitlementID string
}
type resourceEquipmentCleanupRow struct {
Cursor resourceEquipmentCleanupCursor
ResourceID int64
Invalid bool
}
// AcquireResourceEquipmentCleanupLock 保证 Wallet 多副本每轮只有一个物理清理者。
// 持久游标行仍会在删除事务内 FOR UPDATE作为连接异常导致 advisory lock 丢失时的第二道串行边界。
func (r *Repository) AcquireResourceEquipmentCleanupLock(ctx context.Context) (release func(), acquired bool, err error) {
return r.acquireNamedLock(ctx, resourceEquipmentCleanupLockName)
}
// CleanupExpiredResourceEquipmentPage 清理一个稳定主键页,并且只对 MySQL 1213 做一次短抖动重试。
// 1205 可能已经消耗完整 lock_wait_timeout不在本地重试下一轮 worker 会从未提交游标重新执行。
func (r *Repository) CleanupExpiredResourceEquipmentPage(ctx context.Context, scanLimit int, deleteLimit int, nowMS int64) (ResourceEquipmentCleanupPageResult, error) {
if r == nil || r.db == nil {
return ResourceEquipmentCleanupPageResult{}, xerr.New(xerr.Unavailable, "mysql repository is not configured")
}
if scanLimit <= 0 || scanLimit > 1000 || deleteLimit <= 0 || deleteLimit > scanLimit {
return ResourceEquipmentCleanupPageResult{}, xerr.New(xerr.InvalidArgument, "resource equipment cleanup limits are invalid")
}
if nowMS <= 0 {
return ResourceEquipmentCleanupPageResult{}, xerr.New(xerr.InvalidArgument, "resource equipment cleanup time is invalid")
}
var lastErr error
for attempt := 0; attempt < 2; attempt++ {
result, err := r.cleanupExpiredResourceEquipmentPageOnce(ctx, scanLimit, deleteLimit, nowMS)
if err == nil {
result.DeadlockRetries = attempt
return result, nil
}
if !isMySQLDeadlockError(err) || attempt == 1 {
return ResourceEquipmentCleanupPageResult{}, err
}
lastErr = err
timer := time.NewTimer(resourceEquipmentCleanupRetryDelay())
select {
case <-ctx.Done():
if !timer.Stop() {
<-timer.C
}
return ResourceEquipmentCleanupPageResult{}, ctx.Err()
case <-timer.C:
}
}
return ResourceEquipmentCleanupPageResult{}, lastErr
}
func (r *Repository) cleanupExpiredResourceEquipmentPageOnce(ctx context.Context, scanLimit int, deleteLimit int, nowMS int64) (ResourceEquipmentCleanupPageResult, error) {
tx, err := r.db.BeginTx(ctx, &sql.TxOptions{Isolation: sql.LevelReadCommitted})
if err != nil {
return ResourceEquipmentCleanupPageResult{}, err
}
defer func() { _ = tx.Rollback() }()
// 单行状态先加锁、装备行后加锁;即使 advisory lock 因连接故障意外释放,
// 两个副本也会按同一顺序串行处理,不能同时推进或跳过游标。
if _, err = tx.ExecContext(ctx, `
INSERT IGNORE INTO wallet_resource_equipment_cleanup_state (
job_name, cursor_app_code, cursor_user_id, cursor_resource_type,
cursor_entitlement_id, updated_at_ms
) VALUES (?, '', 0, '', '', 0)`,
resourceEquipmentCleanupJobName,
); err != nil {
return ResourceEquipmentCleanupPageResult{}, err
}
var cursor resourceEquipmentCleanupCursor
if err = tx.QueryRowContext(ctx, `
SELECT cursor_app_code, cursor_user_id, cursor_resource_type, cursor_entitlement_id
FROM wallet_resource_equipment_cleanup_state
WHERE job_name = ?
FOR UPDATE`,
resourceEquipmentCleanupJobName,
).Scan(&cursor.AppCode, &cursor.UserID, &cursor.ResourceType, &cursor.EntitlementID); err != nil {
return ResourceEquipmentCleanupPageResult{}, err
}
rows, err := selectResourceEquipmentCleanupPage(ctx, tx, cursor, scanLimit, nowMS)
if err != nil {
return ResourceEquipmentCleanupPageResult{}, err
}
if len(rows) == 0 {
if cursor != (resourceEquipmentCleanupCursor{}) {
if err = updateResourceEquipmentCleanupCursor(ctx, tx, resourceEquipmentCleanupCursor{}, nowMS); err != nil {
return ResourceEquipmentCleanupPageResult{}, err
}
}
if err = tx.Commit(); err != nil {
return ResourceEquipmentCleanupPageResult{}, err
}
return ResourceEquipmentCleanupPageResult{ReachedEnd: true}, nil
}
// 删除候选始终保持 equipment PRIMARY KEY 顺序。达到删除上限时只推进到最后一个候选,
// 让本页后续行在下一事务重新复核,不能因为批次截断而永久跳过失效指针。
candidates := make([]resourceEquipmentCleanupRow, 0, deleteLimit)
nextCursor := rows[len(rows)-1].Cursor
for _, row := range rows {
if !row.Invalid {
continue
}
candidates = append(candidates, row)
if len(candidates) == deleteLimit {
nextCursor = row.Cursor
break
}
}
deletedCount := 0
for _, candidate := range candidates {
deleted, deleteErr := deleteInvalidResourceEquipmentCandidate(ctx, tx, candidate, nowMS)
if deleteErr != nil {
return ResourceEquipmentCleanupPageResult{}, deleteErr
}
deletedCount += deleted
}
lastScannedCursor := rows[len(rows)-1].Cursor
reachedEnd := len(rows) < scanLimit && nextCursor == lastScannedCursor
if reachedEnd {
nextCursor = resourceEquipmentCleanupCursor{}
}
if err = updateResourceEquipmentCleanupCursor(ctx, tx, nextCursor, nowMS); err != nil {
return ResourceEquipmentCleanupPageResult{}, err
}
if err = tx.Commit(); err != nil {
return ResourceEquipmentCleanupPageResult{}, err
}
return ResourceEquipmentCleanupPageResult{
ScannedCount: len(rows),
DeletedCount: deletedCount,
ReachedEnd: reachedEnd,
}, nil
}
func selectResourceEquipmentCleanupPage(ctx context.Context, tx *sql.Tx, cursor resourceEquipmentCleanupCursor, limit int, nowMS int64) ([]resourceEquipmentCleanupRow, error) {
const pageProjection = `
SELECT page.app_code, page.user_id, page.resource_type, page.entitlement_id, page.resource_id,
(
e.entitlement_id IS NULL
OR e.status <> 'active'
OR e.effective_at_ms > ?
OR (e.expires_at_ms <> 0 AND e.expires_at_ms <= ?)
OR e.remaining_quantity <= 0
OR (e.source_snapshot_id = '' AND (r.resource_id IS NULL OR r.status <> 'active'))
) AS invalid
FROM (%s) AS page
LEFT JOIN user_resource_entitlements AS e
ON e.app_code = page.app_code
AND e.user_id = page.user_id
AND e.resource_id = page.resource_id
AND e.entitlement_id = page.entitlement_id
LEFT JOIN resources AS r
ON r.resource_id = page.resource_id
AND r.app_code = page.app_code
ORDER BY page.app_code, page.user_id, page.resource_type, page.entitlement_id`
// 先从 equipment PRIMARY KEY 截取固定数量,再做两个点查 JOIN。显式字典序 OR 谓词让
// MySQL 使用 Index range scanrow-constructor `>` 会退化成从主键表头过滤,游标靠后时不够有界。
var (
query string
args []any
)
if cursor == (resourceEquipmentCleanupCursor{}) {
query = `
SELECT app_code, user_id, resource_type, entitlement_id, resource_id
FROM user_resource_equipment FORCE INDEX (PRIMARY)
ORDER BY app_code, user_id, resource_type, entitlement_id
LIMIT ?`
args = []any{nowMS, nowMS, limit}
} else {
query = `
SELECT app_code, user_id, resource_type, entitlement_id, resource_id
FROM user_resource_equipment FORCE INDEX (PRIMARY)
WHERE app_code > ?
OR (app_code = ? AND user_id > ?)
OR (app_code = ? AND user_id = ? AND resource_type > ?)
OR (app_code = ? AND user_id = ? AND resource_type = ? AND entitlement_id > ?)
ORDER BY app_code, user_id, resource_type, entitlement_id
LIMIT ?`
args = []any{
nowMS, nowMS,
cursor.AppCode,
cursor.AppCode, cursor.UserID,
cursor.AppCode, cursor.UserID, cursor.ResourceType,
cursor.AppCode, cursor.UserID, cursor.ResourceType, cursor.EntitlementID,
limit,
}
}
rows, err := tx.QueryContext(ctx, fmt.Sprintf(pageProjection, query), args...)
if err != nil {
return nil, err
}
defer rows.Close()
result := make([]resourceEquipmentCleanupRow, 0, limit)
for rows.Next() {
var row resourceEquipmentCleanupRow
if err = rows.Scan(
&row.Cursor.AppCode,
&row.Cursor.UserID,
&row.Cursor.ResourceType,
&row.Cursor.EntitlementID,
&row.ResourceID,
&row.Invalid,
); err != nil {
return nil, err
}
result = append(result, row)
}
if err = rows.Err(); err != nil {
return nil, err
}
return result, nil
}
func deleteInvalidResourceEquipmentCandidate(ctx context.Context, tx *sql.Tx, candidate resourceEquipmentCleanupRow, nowMS int64) (int, error) {
result, err := tx.ExecContext(ctx, `
DELETE eq
FROM user_resource_equipment AS eq FORCE INDEX (PRIMARY)
LEFT JOIN user_resource_entitlements AS e
ON e.app_code = eq.app_code
AND e.user_id = eq.user_id
AND e.resource_id = eq.resource_id
AND e.entitlement_id = eq.entitlement_id
LEFT JOIN resources AS r
ON r.resource_id = eq.resource_id
AND r.app_code = eq.app_code
WHERE eq.app_code = ?
AND eq.user_id = ?
AND eq.resource_type = ?
AND eq.entitlement_id = ?
AND eq.resource_id = ?
AND (
e.entitlement_id IS NULL
OR e.status <> 'active'
OR e.effective_at_ms > ?
OR (e.expires_at_ms <> 0 AND e.expires_at_ms <= ?)
OR e.remaining_quantity <= 0
OR (e.source_snapshot_id = '' AND (r.resource_id IS NULL OR r.status <> 'active'))
)`,
candidate.Cursor.AppCode,
candidate.Cursor.UserID,
candidate.Cursor.ResourceType,
candidate.Cursor.EntitlementID,
candidate.ResourceID,
nowMS,
nowMS,
)
if err != nil {
return 0, err
}
affected, err := result.RowsAffected()
if err != nil {
return 0, err
}
return int(affected), nil
}
func updateResourceEquipmentCleanupCursor(ctx context.Context, tx *sql.Tx, cursor resourceEquipmentCleanupCursor, nowMS int64) error {
_, err := tx.ExecContext(ctx, `
UPDATE wallet_resource_equipment_cleanup_state
SET cursor_app_code = ?, cursor_user_id = ?, cursor_resource_type = ?,
cursor_entitlement_id = ?, updated_at_ms = ?
WHERE job_name = ?`,
cursor.AppCode,
cursor.UserID,
cursor.ResourceType,
cursor.EntitlementID,
nowMS,
resourceEquipmentCleanupJobName,
)
return err
}
func isMySQLDeadlockError(err error) bool {
var mysqlErr *mysqlDriver.MySQLError
return errors.As(err, &mysqlErr) && mysqlErr.Number == 1213
}
func mapTransientResourceReadError(err error) error {
if !isMySQLDeadlockError(err) {
return err
}
// 对外只暴露稳定、可判定的瞬时 reason原始 MySQL 错误继续留在 error chain
// 供 Wallet 服务端 access log 诊断,不能进入 gRPC message 或 HTTP envelope。
return fmt.Errorf("%w: %w",
xerr.New(xerr.WalletResourceReadTransient, "wallet resource data is temporarily unavailable"),
err,
)
}
func resourceEquipmentCleanupRetryDelay() time.Duration {
// 单次重试固定在 20-60ms 内;随机抖动只用于错开 Equip/Revoke 的瞬时持锁窗口,
// 不能形成无限退避或把持续故障隐藏在 worker 内。
return 20*time.Millisecond + time.Duration(rand.IntN(41))*time.Millisecond
}