refactor(repository): 收敛 model/repository 分层为唯一持久化入口

将 OpenFlare 与平台业务的数据访问从 model 与 apps 直连迁入 repository,
model 仅保留实体与无 IO 规则;补充 code-check 架构守卫与开发规范。
This commit is contained in:
ryan
2026-07-24 17:00:17 +08:00
parent 23a5488203
commit 943818f7d4
184 changed files with 5592 additions and 4364 deletions
+21 -25
View File
@@ -10,15 +10,15 @@ import (
"fmt"
"time"
"gorm.io/gorm"
"github.com/Rain-kl/Wavelet/internal/apps/upload/ingest"
"github.com/Rain-kl/Wavelet/internal/apps/upload/shared"
uploadstorage "github.com/Rain-kl/Wavelet/internal/apps/upload/storage"
db "github.com/Rain-kl/Wavelet/internal/infra/persistence"
"github.com/Rain-kl/Wavelet/internal/infra/task"
"github.com/Rain-kl/Wavelet/internal/model"
"github.com/Rain-kl/Wavelet/internal/repository"
"github.com/Rain-kl/Wavelet/pkg/logger"
"gorm.io/gorm"
"gorm.io/gorm/clause"
)
const (
@@ -58,12 +58,8 @@ func (h *SystemCleanupHandler) Execute(ctx context.Context, _ []byte) (*task.Tas
task.AppendLog(ctx, "开始扫描未使用上传文件,阈值: %s", oneHourAgo.Format(time.RFC3339))
for {
var unusedUploads []model.Upload
if err := db.DB(ctx).
Where("id > ? AND status = ? AND created_at < ?", lastID, model.UploadStatusPending, oneHourAgo).
Order("id ASC").
Limit(batchSize).
Find(&unusedUploads).Error; err != nil {
unusedUploads, err := repository.ListPendingUploadsOlderThan(ctx, lastID, oneHourAgo, batchSize)
if err != nil {
task.AppendLog(ctx, "查询未使用的上传文件失败: %v", err)
return nil, fmt.Errorf(shared.ErrQueryUnusedUploadsFailed, err)
}
@@ -78,19 +74,18 @@ func (h *SystemCleanupHandler) Execute(ctx context.Context, _ []byte) (*task.Tas
totalProcessed++
transitioned := false
if err := db.DB(ctx).Transaction(func(tx *gorm.DB) error {
var locked model.Upload
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).
Where("id = ?", u.ID).
First(&locked).Error; err != nil {
// Multi-step: row lock + ownership-safe soft-delete + stats delta stay orchestrated here.
if err := repository.RunInTransaction(ctx, func(tx *gorm.DB) error {
locked, err := repository.GetUploadByIDForUpdateTx(tx, u.ID)
if err != nil {
return err
}
if locked.Status != model.UploadStatusPending || !locked.CreatedAt.Before(oneHourAgo) {
return nil
}
var err error
transitioned, err = ingest.RemoveLockedTx(tx, &locked)
return err
var removeErr error
transitioned, removeErr = ingest.RemoveLockedTx(tx, &locked)
return removeErr
}); err != nil {
task.AppendLog(ctx, "清理上传文件失败 [ID:%d]: %v", u.ID, err)
lastID = u.ID
@@ -107,21 +102,22 @@ func (h *SystemCleanupHandler) Execute(ctx context.Context, _ []byte) (*task.Tas
task.AppendLog(ctx, "开始清理历史推送审计日志,只保留最近7天数据...")
cutoff := time.Now().AddDate(0, 0, -7)
var pushHistoryCount int64
if err := db.DB(ctx).Model(&model.PushHistory{}).Where("created_at < ?", cutoff).Count(&pushHistoryCount).Error; err != nil {
pushHistoryCount, err := repository.CountPushHistoriesCreatedBefore(ctx, cutoff)
switch {
case err != nil:
task.AppendLog(ctx, "统计待清理的历史推送记录失败: %v", err)
} else if pushHistoryCount > 0 {
if err := db.DB(ctx).Where("created_at < ?", cutoff).Delete(&model.PushHistory{}).Error; err != nil {
task.AppendLog(ctx, "删除历史推送记录失败: %v", err)
case pushHistoryCount == 0:
task.AppendLog(ctx, "没有需要清理的历史推送记录 (截止时间: %s)", cutoff.Format("2006-01-02 15:04:05"))
default:
if _, delErr := repository.DeletePushHistoriesCreatedBefore(ctx, cutoff); delErr != nil {
task.AppendLog(ctx, "删除历史推送记录失败: %v", delErr)
} else {
task.AppendLog(ctx, "成功删除 %d 条历史推送记录 (截止时间: %s)", pushHistoryCount, cutoff.Format("2006-01-02 15:04:05"))
}
} else {
task.AppendLog(ctx, "没有需要清理的历史推送记录 (截止时间: %s)", cutoff.Format("2006-01-02 15:04:05"))
}
task.AppendLog(ctx, "开始清理任务执行日志:高频任务保留最近3天,低频任务保留最近30天...")
taskLogStats, err := model.CleanupTaskExecutionLogs(ctx, time.Now())
taskLogStats, err := repository.CleanupTaskExecutionLogs(ctx, time.Now())
if err != nil {
task.AppendLog(ctx, "清理任务执行日志失败: %v", err)
logger.ErrorF(ctx, "清理任务执行日志失败: %v", err)
+5 -11
View File
@@ -8,9 +8,8 @@ import (
"fmt"
uploadstats "github.com/Rain-kl/Wavelet/internal/apps/upload/stats"
db "github.com/Rain-kl/Wavelet/internal/infra/persistence"
"github.com/Rain-kl/Wavelet/internal/infra/task"
"github.com/Rain-kl/Wavelet/internal/model"
"github.com/Rain-kl/Wavelet/internal/repository"
)
const (
@@ -37,11 +36,8 @@ type RebuildUploadStatsHandler struct{}
// Execute scans active uploads and rebuilds all upload stat dimensions.
func (h *RebuildUploadStatsHandler) Execute(ctx context.Context, _ []byte) (*task.TaskResult, error) {
var activeCount int64
if err := db.DB(ctx).
Model(&model.Upload{}).
Where("status != ?", model.UploadStatusDeleted).
Count(&activeCount).Error; err != nil {
activeCount, err := repository.CountActiveUploads(ctx)
if err != nil {
task.AppendLog(ctx, "统计活跃上传记录失败: %v", err)
return nil, fmt.Errorf("count active uploads: %w", err)
}
@@ -53,10 +49,8 @@ func (h *RebuildUploadStatsHandler) Execute(ctx context.Context, _ []byte) (*tas
return nil, fmt.Errorf("rebuild upload stats: %w", err)
}
var totalStat model.UploadStat
if err := db.DB(ctx).
Where("dimension = ? AND stat_key = ?", model.UploadStatDimensionTotal, "").
First(&totalStat).Error; err != nil {
totalStat, err := repository.GetTotalUploadStat(ctx)
if err != nil {
task.AppendLog(ctx, "读取总量统计失败: %v", err)
return nil, fmt.Errorf("load total upload stats: %w", err)
}
+12 -40
View File
@@ -15,13 +15,15 @@ import (
"sync/atomic"
"time"
"golang.org/x/sync/errgroup"
uploadstats "github.com/Rain-kl/Wavelet/internal/apps/upload/stats"
uploadstorage "github.com/Rain-kl/Wavelet/internal/apps/upload/storage"
"github.com/Rain-kl/Wavelet/internal/infra/objectstore"
db "github.com/Rain-kl/Wavelet/internal/infra/persistence"
"github.com/Rain-kl/Wavelet/internal/infra/task"
"github.com/Rain-kl/Wavelet/internal/model"
"golang.org/x/sync/errgroup"
"github.com/Rain-kl/Wavelet/internal/repository"
)
const (
@@ -171,12 +173,7 @@ func (h *MigrationHandler) Execute(ctx context.Context, payload []byte) (*task.T
}
func countStorageObjects(ctx context.Context) (int64, error) {
var count int64
err := db.DB(ctx).Model(&model.Upload{}).
Where("status != ?", model.UploadStatusDeleted).
Distinct("file_path").
Count(&count).Error
return count, err
return repository.CountDistinctActiveFilePaths(ctx)
}
func hasUnresolvedMigrationTask(ctx context.Context) (bool, error) {
@@ -187,13 +184,6 @@ func hasUnresolvedMigrationTask(ctx context.Context) (bool, error) {
return execution.Status == model.TaskExecutionStatusPending || execution.Status == model.TaskExecutionStatusRunning, nil
}
type migrationObject struct {
FilePath string `gorm:"column:file_path"`
FileSize int64 `gorm:"column:file_size"`
MimeType string `gorm:"column:mime_type"`
Hash string `gorm:"column:hash"`
}
func migrateObjects(
ctx context.Context,
sourceBackend objectstore.Backend,
@@ -212,17 +202,8 @@ func migrateObjects(
task.AppendLog(ctx, "正在查询待迁移对象批次,当前已完成迁移: %d/%d", atomic.LoadInt64(&migrated), total)
var objects []migrationObject
query := db.DB(ctx).Model(&model.Upload{}).
Select("file_path, MAX(file_size) AS file_size, MAX(mime_type) AS mime_type, MAX(hash) AS hash").
Where("status != ?", model.UploadStatusDeleted)
if lastFilePath != "" {
query = query.Where("file_path > ?", lastFilePath)
}
if err := query.Group("file_path").
Order("file_path ASC").
Limit(batchSize).
Scan(&objects).Error; err != nil {
objects, err := repository.ListDistinctActiveStorageObjects(ctx, lastFilePath, batchSize)
if err != nil {
return atomic.LoadInt64(&migrated), fmt.Errorf("query source objects: %w", err)
}
if len(objects) == 0 {
@@ -260,7 +241,7 @@ func migrateSingleObject(
ctx context.Context,
sourceBackend objectstore.Backend,
targetBackend objectstore.Backend,
obj migrationObject,
obj repository.UploadStorageObject,
sha256HexLength int,
) error {
if shouldSkipMigration(ctx, targetBackend, obj) {
@@ -310,9 +291,7 @@ func migrateSingleObject(
if targetResult.Key != obj.FilePath {
task.AppendLog(ctx, "[更新数据库] 正在更新文件路径: %s -> %s", obj.FilePath, targetResult.Key)
if err := db.DB(ctx).Model(&model.Upload{}).
Where("file_path = ? AND status != ?", obj.FilePath, model.UploadStatusDeleted).
Update("file_path", targetResult.Key).Error; err != nil {
if err := repository.UpdateActiveUploadsFilePath(ctx, obj.FilePath, targetResult.Key); err != nil {
return fmt.Errorf("update migrated object %q: %w", obj.FilePath, err)
}
}
@@ -323,7 +302,7 @@ func migrateSingleObject(
func shouldSkipMigration(
ctx context.Context,
targetBackend objectstore.Backend,
obj migrationObject,
obj repository.UploadStorageObject,
) bool {
targetObj, err := targetBackend.Get(ctx, obj.FilePath)
if err != nil || targetObj == nil || targetObj.Body == nil {
@@ -343,16 +322,9 @@ func markMissingMigrationObjectDeleted(
) error {
task.AppendLog(ctx, "警告: 源存储中物理文件不存在,标记为已删除并跳过: %s (错误: %v)", filePath, sourceErr)
var affectedUploads []model.Upload
if err := db.DB(ctx).
Where("file_path = ? AND status != ?", filePath, model.UploadStatusDeleted).
Find(&affectedUploads).Error; err != nil {
return fmt.Errorf("load missing object uploads %q: %w", filePath, err)
}
if err := db.DB(ctx).Model(&model.Upload{}).
Where("file_path = ?", filePath).
Update("status", model.UploadStatusDeleted).Error; err != nil {
return fmt.Errorf("update missing object %q: %w", filePath, err)
affectedUploads, err := repository.MarkActiveUploadsDeletedByFilePath(ctx, filePath)
if err != nil {
return fmt.Errorf("mark missing object deleted %q: %w", filePath, err)
}
for i := range affectedUploads {
uploadstats.RecordUploadStatsRemove(ctx, &affectedUploads[i])
+3 -13
View File
@@ -14,9 +14,8 @@ import (
"github.com/Rain-kl/Wavelet/internal/apps/upload/filesrv"
"github.com/Rain-kl/Wavelet/internal/apps/upload/shared"
db "github.com/Rain-kl/Wavelet/internal/infra/persistence"
"github.com/Rain-kl/Wavelet/internal/infra/task"
"github.com/Rain-kl/Wavelet/internal/model"
"github.com/Rain-kl/Wavelet/internal/repository"
)
const (
@@ -113,17 +112,8 @@ func (h *WarmImageCacheHandler) Execute(ctx context.Context, payload []byte) (*t
return nil, fmt.Errorf("image cache warmup canceled: %w", err)
}
var uploads []model.Upload
if err := db.DB(ctx).
Where("id > ? AND status != ? AND (LOWER(mime_type) LIKE ? OR LOWER(extension) IN ?)",
lastID,
model.UploadStatusDeleted,
"image/%",
[]string{"jpg", "jpeg", "png", "webp", "gif"},
).
Order("id ASC").
Limit(batchSize).
Find(&uploads).Error; err != nil {
uploads, err := repository.ListActiveImageUploadsAfterID(ctx, lastID, batchSize)
if err != nil {
task.AppendLog(ctx, "查询图片上传记录失败: %v", err)
return nil, fmt.Errorf(shared.ErrQueryImagesForCacheWarmup, err)
}
+3 -1
View File
@@ -17,6 +17,8 @@ import (
"testing"
"time"
"github.com/Rain-kl/Wavelet/internal/repository"
"github.com/Rain-kl/Wavelet/internal/apps/upload/filesrv"
"github.com/Rain-kl/Wavelet/internal/apps/upload/shared"
uploadstats "github.com/Rain-kl/Wavelet/internal/apps/upload/stats"
@@ -129,7 +131,7 @@ func TestSystemCleanupHandler_Execute(t *testing.T) {
UpdatedAt: now.AddDate(0, 0, -31),
TriggeredBy: "system",
}
err = model.CreateTaskExecution(ctx, oldTaskLog)
err = repository.CreateTaskExecution(ctx, oldTaskLog)
require.NoError(t, err)
// 执行 handler