Files
yml2213 85332df2bd 订单接口最小化与私有文件访问加固
- 订单列表使用独立最小 DTO 并分页,号主待办提供独立接口与统计
- 用户 token 增加版本控制,冻结/改密/退出即时撤销会话
- 移除 URL token 传参,SSE 与接口统一使用 HttpOnly Cookie
- 私有文件按上传归属与业务关联授权,收款凭证转私有访问并校验归属
- 公开商品接口返回最小字段,隐藏号主身份与内部状态
- 每日清理超过 30 天未关联业务的上传归属,上传归属失败时补偿删除对象
2026-08-16 21:47:46 +08:00

168 lines
4.6 KiB
Go

package fileuploadcleanup
import (
"context"
"crypto/rand"
"encoding/hex"
"fmt"
"net/url"
"time"
"hfb_sys/backend/internal/model"
"github.com/redis/go-redis/v9"
"go.uber.org/zap"
"gorm.io/gorm"
)
const (
cleanupLockKey = "hfb:job:file-upload-cleanup:lock"
cleanupInterval = 24 * time.Hour
cleanupRetention = 30 * 24 * time.Hour
cleanupBatchSize = 200
)
// Job 清理长期未关联业务记录的上传归属,避免临时草稿记录无限增长。
type Job struct {
db *gorm.DB
redis *redis.Client
logger *zap.Logger
instanceID string
}
func New(db *gorm.DB, redisClient *redis.Client, logger *zap.Logger) *Job {
if logger == nil {
logger = zap.NewNop()
}
return &Job{db: db, redis: redisClient, logger: logger, instanceID: newInstanceID()}
}
func (j *Job) Start(ctx context.Context) {
if j == nil || j.db == nil {
return
}
go j.loop(ctx)
}
func (j *Job) loop(ctx context.Context) {
j.run(ctx, time.Now())
ticker := time.NewTicker(cleanupInterval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
j.logger.Debug("临时文件归属清理任务已停止")
return
case now := <-ticker.C:
j.run(ctx, now)
}
}
}
func (j *Job) run(ctx context.Context, now time.Time) {
release, ok := j.acquireLock(ctx)
if !ok {
return
}
defer release()
deleted, err := j.cleanup(ctx, now)
if err != nil {
j.logger.Warn("临时文件归属清理失败", zap.Error(err))
return
}
if deleted > 0 {
j.logger.Info("已清理未关联临时文件归属", zap.Int("count", deleted))
}
}
func (j *Job) acquireLock(ctx context.Context) (func(), bool) {
if j.redis == nil {
return func() {}, true
}
ok, err := j.redis.SetNX(ctx, cleanupLockKey, j.instanceID, 10*time.Minute).Result()
if err != nil {
j.logger.Warn("临时文件归属清理任务获取锁失败", zap.Error(err))
return func() {}, true
}
if !ok {
return nil, false
}
return func() {
releaseCtx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
defer cancel()
script := redis.NewScript(`if redis.call("get", KEYS[1]) == ARGV[1] then return redis.call("del", KEYS[1]) else return 0 end`)
if err := script.Run(releaseCtx, j.redis, []string{cleanupLockKey}, j.instanceID).Err(); err != nil {
j.logger.Warn("临时文件归属清理任务释放锁失败", zap.Error(err))
}
}, true
}
func (j *Job) cleanup(ctx context.Context, now time.Time) (int, error) {
cutoff := now.Add(-cleanupRetention)
lastID := uint64(0)
deleted := 0
for {
var records []model.FileUploadOwner
if err := j.db.WithContext(ctx).
Where("id > ? AND created_at < ?", lastID, cutoff).
Order("id ASC").
Limit(cleanupBatchSize).
Find(&records).Error; err != nil {
return deleted, err
}
if len(records) == 0 {
return deleted, nil
}
for _, record := range records {
lastID = record.ID
referenced, err := j.isReferenced(ctx, record.ObjectKey)
if err != nil {
return deleted, err
}
if referenced {
continue
}
result := j.db.WithContext(ctx).
Where("id = ? AND created_at < ?", record.ID, cutoff).
Delete(&model.FileUploadOwner{})
if result.Error != nil {
return deleted, result.Error
}
deleted += int(result.RowsAffected)
}
}
}
func (j *Job) isReferenced(ctx context.Context, key string) (bool, error) {
encodedKey := url.QueryEscape(key)
queries := []struct {
sql string
args []any
}{
{sql: "SELECT COUNT(1) FROM game_accounts WHERE INSTR(screenshot_urls, ?) > 0 OR INSTR(screenshot_urls, ?) > 0", args: []any{key, encodedKey}},
{sql: "SELECT COUNT(1) FROM order_checkouts WHERE INSTR(evidence_urls, ?) > 0 OR INSTR(evidence_urls, ?) > 0", args: []any{key, encodedKey}},
{sql: "SELECT COUNT(1) FROM disputes WHERE INSTR(evidence_urls, ?) > 0 OR INSTR(evidence_urls, ?) > 0", args: []any{key, encodedKey}},
{sql: "SELECT COUNT(1) FROM handoff_records WHERE INSTR(attachment_urls, ?) > 0 OR INSTR(attachment_urls, ?) > 0", args: []any{key, encodedKey}},
{sql: "SELECT COUNT(1) FROM chat_messages WHERE INSTR(attachment_urls, ?) > 0 OR INSTR(attachment_urls, ?) > 0", args: []any{key, encodedKey}},
{sql: "SELECT COUNT(1) FROM user_payment_accounts WHERE INSTR(certificate_urls, ?) > 0 OR INSTR(certificate_urls, ?) > 0", args: []any{key, encodedKey}},
}
for _, query := range queries {
var count int64
if err := j.db.WithContext(ctx).Raw(query.sql, query.args...).Scan(&count).Error; err != nil {
return false, err
}
if count > 0 {
return true, nil
}
}
return false, nil
}
func newInstanceID() string {
value := make([]byte, 8)
if _, err := rand.Read(value); err != nil {
return fmt.Sprintf("file-cleanup-%d", time.Now().UnixNano())
}
return hex.EncodeToString(value)
}