mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-03 15:06:36 +08:00
fix(pages): 收紧部署包与 Agent 同步边界
完成 V2 Phase 0 安全与一致性前置:统一真实归档限额、流式拉取、候选裁剪、保留上传删除语义及 Pages 路由引用锁。
This commit is contained in:
@@ -8,6 +8,7 @@ import (
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/filesrv"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/handler"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/ingest"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/shared"
|
||||
uploadstats "github.com/Rain-kl/Wavelet/internal/apps/upload/stats"
|
||||
uploadtask "github.com/Rain-kl/Wavelet/internal/apps/upload/task"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/util"
|
||||
@@ -31,15 +32,17 @@ var (
|
||||
|
||||
// Programmatic ingest API
|
||||
var (
|
||||
Ingest = ingest.Ingest
|
||||
Remove = ingest.Remove
|
||||
RemoveOwned = ingest.RemoveOwned
|
||||
FindByHash = ingest.FindByHash
|
||||
GetActiveUpload = ingest.GetActive
|
||||
OpenStoredUpload = ingest.OpenActiveObject
|
||||
ActiveUploadHash = ingest.ActiveHash
|
||||
ResolveLocalFile = ingest.ResolveLocalFile
|
||||
IngestFromLocalPath = ingest.FromLocalPath
|
||||
Ingest = ingest.Ingest
|
||||
Remove = ingest.Remove
|
||||
RemoveOwned = ingest.RemoveOwned
|
||||
RemoveLockedTx = ingest.RemoveLockedTx
|
||||
InvalidateUploadMetaCache = ingest.InvalidateUploadMetaCache
|
||||
FindByHash = ingest.FindByHash
|
||||
GetActiveUpload = ingest.GetActive
|
||||
OpenStoredUpload = ingest.OpenActiveObject
|
||||
ActiveUploadHash = ingest.ActiveHash
|
||||
ResolveLocalFile = ingest.ResolveLocalFile
|
||||
IngestFromLocalPath = ingest.FromLocalPath
|
||||
)
|
||||
|
||||
type (
|
||||
@@ -54,6 +57,8 @@ const (
|
||||
PolicyCreate = ingest.PolicyCreate
|
||||
PolicyDedupNewRecord = ingest.PolicyDedupNewRecord
|
||||
PolicyResolveExisting = ingest.PolicyResolveExisting
|
||||
// ReservedPagesDeploymentType is managed exclusively by the Pages domain.
|
||||
ReservedPagesDeploymentType = shared.ReservedPagesDeploymentType
|
||||
)
|
||||
|
||||
type (
|
||||
@@ -69,6 +74,7 @@ type (
|
||||
var (
|
||||
ErrIngestForbidden = ingest.ErrForbidden
|
||||
ErrIngestStorageReadOnly = ingest.ErrStorageReadOnly
|
||||
ErrReservedUploadType = ingest.ErrReservedUploadType
|
||||
)
|
||||
|
||||
// Cache management
|
||||
|
||||
@@ -4,6 +4,7 @@
|
||||
package handler
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"net/http"
|
||||
"strconv"
|
||||
|
||||
@@ -96,6 +97,7 @@ func ListFiles(c *gin.Context) {
|
||||
// @Success 200 {object} response.Any "删除成功"
|
||||
// @Failure 403 {object} response.Any "无权操作"
|
||||
// @Failure 404 {object} response.Any "文件不存在"
|
||||
// @Failure 409 {object} response.Any "系统保留类型或存储只读"
|
||||
// @Router /api/v1/admin/uploads/{id} [delete]
|
||||
func DeleteFile(c *gin.Context) {
|
||||
ctx := c.Request.Context()
|
||||
@@ -111,6 +113,10 @@ func DeleteFile(c *gin.Context) {
|
||||
}
|
||||
|
||||
if _, err := softDeleteUpload(ctx, uploadID); err != nil {
|
||||
if errors.Is(err, ingest.ErrReservedUploadType) {
|
||||
response.AbortConflict(c, shared.ErrReservedUploadType)
|
||||
return
|
||||
}
|
||||
if isRecordNotFound(err) {
|
||||
response.AbortNotFound(c, "文件记录未找到")
|
||||
return
|
||||
@@ -216,6 +222,7 @@ func ListMyFiles(c *gin.Context) {
|
||||
// @Success 200 {object} response.Any "删除成功"
|
||||
// @Failure 403 {object} response.Any "无权操作"
|
||||
// @Failure 404 {object} response.Any "文件不存在"
|
||||
// @Failure 409 {object} response.Any "系统保留类型或存储只读"
|
||||
// @Router /api/v1/upload/{id} [delete]
|
||||
func DeleteMyFile(c *gin.Context) {
|
||||
currUser, _ := oauth.GetFromContext[*model.User](c, oauth.UserObjKey)
|
||||
@@ -232,11 +239,15 @@ func DeleteMyFile(c *gin.Context) {
|
||||
}
|
||||
|
||||
if _, err := softDeleteOwnedUpload(ctx, currUser.ID, uploadID); err != nil {
|
||||
if errors.Is(err, ingest.ErrReservedUploadType) {
|
||||
response.AbortConflict(c, shared.ErrReservedUploadType)
|
||||
return
|
||||
}
|
||||
if isRecordNotFound(err) {
|
||||
response.AbortNotFound(c, "文件记录未找到")
|
||||
return
|
||||
}
|
||||
if err == ingest.ErrForbidden {
|
||||
if errors.Is(err, ingest.ErrForbidden) {
|
||||
response.AbortForbidden(c, "无权操作")
|
||||
return
|
||||
}
|
||||
|
||||
@@ -53,6 +53,7 @@ type batchDownloadRequest struct {
|
||||
// @Success 200 {object} response.Any{data=model.Upload} "上传成功"
|
||||
// @Failure 400 {object} response.Any "请求参数错误或文件受限"
|
||||
// @Failure 401 {object} response.Any "未登录"
|
||||
// @Failure 409 {object} response.Any "系统保留类型或存储只读"
|
||||
// @Failure 500 {object} response.Any "内部错误"
|
||||
// @Router /api/v1/upload [post]
|
||||
//
|
||||
@@ -107,6 +108,10 @@ func UploadFile(c *gin.Context) {
|
||||
}
|
||||
|
||||
uploadType := c.DefaultPostForm("type", "generic")
|
||||
if uploadType == shared.ReservedPagesDeploymentType {
|
||||
response.AbortConflict(c, shared.ErrReservedUploadType)
|
||||
return
|
||||
}
|
||||
|
||||
accessMode, errMsg := resolveUploadAccessMode(c, uploadType)
|
||||
if errMsg != "" {
|
||||
|
||||
@@ -226,6 +226,42 @@ func TestUploadFile(t *testing.T) {
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("upload rejects Pages reserved type", func(t *testing.T) {
|
||||
putCountBefore := putCount
|
||||
imgContent := []byte("\x89PNG\r\n\x1a\n\x00\x00\x00\rIHDR\x00\x00\x00\x01\x00\x00\x00\x01")
|
||||
contentType, body := createMultipartRequest(t, "file", "pages.png", imgContent, map[string]string{
|
||||
"type": shared.ReservedPagesDeploymentType,
|
||||
})
|
||||
req, _ := http.NewRequest("POST", "/api/v1/upload", body)
|
||||
req.Header.Set("Content-Type", contentType)
|
||||
|
||||
w := httptest.NewRecorder()
|
||||
router.ServeHTTP(w, req)
|
||||
|
||||
if w.Code != http.StatusConflict {
|
||||
t.Fatalf("expected status 409, got %d. Body: %s", w.Code, w.Body.String())
|
||||
}
|
||||
var resp testResponse
|
||||
if err := json.Unmarshal(w.Body.Bytes(), &resp); err != nil {
|
||||
t.Fatalf("unmarshal reserved type response: %v", err)
|
||||
}
|
||||
if resp.ErrorMsg != shared.ErrReservedUploadType {
|
||||
t.Fatalf("reserved type error = %q, want %q", resp.ErrorMsg, shared.ErrReservedUploadType)
|
||||
}
|
||||
if putCount != putCountBefore {
|
||||
t.Fatalf("reserved upload wrote storage object: put count %d -> %d", putCountBefore, putCount)
|
||||
}
|
||||
var count int64
|
||||
if err := dbConn.Model(&model.Upload{}).
|
||||
Where("type = ?", shared.ReservedPagesDeploymentType).
|
||||
Count(&count).Error; err != nil {
|
||||
t.Fatalf("count reserved uploads: %v", err)
|
||||
}
|
||||
if count != 0 {
|
||||
t.Fatalf("reserved upload record count = %d, want 0", count)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("instant upload deduplication (秒传)", func(t *testing.T) {
|
||||
putCount = 0
|
||||
imgContent := []byte("\x89PNG\r\n\x1a\n\x00\x00\x00\rIHDR\x00\x00\x00\x01\x00\x00\x00\x01")
|
||||
@@ -993,6 +1029,62 @@ func TestUserUploadManagement(t *testing.T) {
|
||||
})
|
||||
}
|
||||
|
||||
func TestDeleteReservedUploadType(t *testing.T) {
|
||||
dbConn, _, cleanup := testhelper.SetupTestEnvironment(t)
|
||||
defer cleanup()
|
||||
|
||||
authUser := &model.User{ID: 1001, Username: "reserved_owner"}
|
||||
router := setupTestRouter(authUser)
|
||||
reserved := model.Upload{
|
||||
ID: 4101,
|
||||
UserID: authUser.ID,
|
||||
FileName: "pages.zip",
|
||||
FilePath: "uploads/pages.zip",
|
||||
FileSize: 128,
|
||||
MimeType: "application/zip",
|
||||
Extension: "zip",
|
||||
Hash: "pages-reserved-hash",
|
||||
Type: shared.ReservedPagesDeploymentType,
|
||||
Status: model.UploadStatusUsed,
|
||||
CreatedAt: time.Now(),
|
||||
}
|
||||
if err := dbConn.Create(&reserved).Error; err != nil {
|
||||
t.Fatalf("seed reserved upload: %v", err)
|
||||
}
|
||||
|
||||
for _, tc := range []struct {
|
||||
name string
|
||||
path string
|
||||
}{
|
||||
{name: "admin delete", path: "/api/v1/admin/uploads/4101"},
|
||||
{name: "owner delete", path: "/api/v1/upload/4101"},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
req, _ := http.NewRequest(http.MethodDelete, tc.path, nil)
|
||||
w := httptest.NewRecorder()
|
||||
router.ServeHTTP(w, req)
|
||||
if w.Code != http.StatusConflict {
|
||||
t.Fatalf("expected status 409, got %d. Body: %s", w.Code, w.Body.String())
|
||||
}
|
||||
var resp testResponse
|
||||
if err := json.Unmarshal(w.Body.Bytes(), &resp); err != nil {
|
||||
t.Fatalf("unmarshal delete response: %v", err)
|
||||
}
|
||||
if resp.ErrorMsg != shared.ErrReservedUploadType {
|
||||
t.Fatalf("reserved delete error = %q, want %q", resp.ErrorMsg, shared.ErrReservedUploadType)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
var persisted model.Upload
|
||||
if err := dbConn.First(&persisted, reserved.ID).Error; err != nil {
|
||||
t.Fatalf("reload reserved upload: %v", err)
|
||||
}
|
||||
if persisted.Status != model.UploadStatusUsed {
|
||||
t.Fatalf("reserved upload status = %s, want used", persisted.Status)
|
||||
}
|
||||
}
|
||||
|
||||
func configureLocalStorageRoot(t *testing.T, dbConn *gorm.DB, tempDir string) {
|
||||
var sc model.SystemConfig
|
||||
if err := dbConn.Where("key = ?", model.ConfigKeyStorageConfig).First(&sc).Error; err != nil {
|
||||
|
||||
@@ -12,5 +12,8 @@ import (
|
||||
// ErrForbidden indicates the caller is not allowed to mutate the upload record.
|
||||
var ErrForbidden = errors.New("upload forbidden")
|
||||
|
||||
// ErrReservedUploadType indicates that a generic mutation targeted a domain-reserved upload type.
|
||||
var ErrReservedUploadType = errors.New(shared.ErrReservedUploadType)
|
||||
|
||||
// ErrStorageReadOnly indicates the storage backend is in migration read-only mode.
|
||||
var ErrStorageReadOnly = errors.New(shared.ErrStorageReadOnly)
|
||||
|
||||
@@ -100,13 +100,10 @@ func storeObject(ctx context.Context, objectKey string, reader io.Reader, size i
|
||||
return result.Key, nil
|
||||
}
|
||||
|
||||
func persistUploadRecord(ctx context.Context, upload *model.Upload, objectKey string) error {
|
||||
func persistUploadRecord(ctx context.Context, upload *model.Upload, objectKey string, storedByRequest bool) error {
|
||||
if err := createUploadWithStats(ctx, upload); err != nil {
|
||||
_, backend, backendErr := storage.Active(ctx)
|
||||
if backendErr == nil {
|
||||
if deleteErr := backend.Delete(ctx, objectKey); deleteErr != nil {
|
||||
logger.WarnF(ctx, "清理未写入数据库的上传对象失败: %v", deleteErr)
|
||||
}
|
||||
if storedByRequest {
|
||||
cleanupUnpersistedObject(ctx, objectKey)
|
||||
}
|
||||
return err
|
||||
}
|
||||
@@ -114,6 +111,16 @@ func persistUploadRecord(ctx context.Context, upload *model.Upload, objectKey st
|
||||
return nil
|
||||
}
|
||||
|
||||
func cleanupUnpersistedObject(ctx context.Context, objectKey string) {
|
||||
_, backend, err := storage.Active(ctx)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
if err := backend.Delete(ctx, objectKey); err != nil {
|
||||
logger.WarnF(ctx, "清理未写入数据库的上传对象失败: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func createUploadWithStats(ctx context.Context, upload *model.Upload) error {
|
||||
return db.DB(ctx).Transaction(func(tx *gorm.DB) error {
|
||||
if err := repository.CreateUploadTx(tx, upload); err != nil {
|
||||
@@ -125,6 +132,8 @@ func createUploadWithStats(ctx context.Context, upload *model.Upload) error {
|
||||
|
||||
func createDedupRecord(ctx context.Context, existing model.Upload, req Request) (Result, error) {
|
||||
accessMode := resolveAccessMode(req.Type, req.AccessMode)
|
||||
metadata := req.Metadata
|
||||
metadata.Bucket = existing.Metadata.Bucket
|
||||
newUpload := model.Upload{
|
||||
ID: idgen.NextUint64ID(),
|
||||
UserID: req.UserID,
|
||||
@@ -137,9 +146,9 @@ func createDedupRecord(ctx context.Context, existing model.Upload, req Request)
|
||||
Type: req.Type,
|
||||
Status: req.Status,
|
||||
AccessMode: accessMode,
|
||||
Metadata: existing.Metadata,
|
||||
Metadata: metadata,
|
||||
}
|
||||
if err := persistUploadRecord(ctx, &newUpload, existing.FilePath); err != nil {
|
||||
if err := persistUploadRecord(ctx, &newUpload, existing.FilePath, false); err != nil {
|
||||
return Result{}, err
|
||||
}
|
||||
logger.InfoF(ctx, "文件触发秒传成功! ID: %d, Path: %s", newUpload.ID, existing.FilePath)
|
||||
@@ -186,7 +195,7 @@ func createNewUpload(ctx context.Context, req Request) (Result, error) {
|
||||
AccessMode: accessMode,
|
||||
Metadata: req.Metadata,
|
||||
}
|
||||
if err := persistUploadRecord(ctx, &upload, storedKey); err != nil {
|
||||
if err := persistUploadRecord(ctx, &upload, storedKey, true); err != nil {
|
||||
return Result{}, err
|
||||
}
|
||||
|
||||
|
||||
@@ -8,15 +8,20 @@ import (
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"errors"
|
||||
"io"
|
||||
"os"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
uploadcache "github.com/Rain-kl/Wavelet/internal/apps/upload/cache"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/shared"
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/storage"
|
||||
"github.com/Rain-kl/Wavelet/internal/testhelper"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
func TestIngestPolicyCreateIncrementsStats(t *testing.T) {
|
||||
@@ -141,7 +146,11 @@ func TestIngestPolicyDedupNewRecordCreatesSecondRecord(t *testing.T) {
|
||||
Extension: "png",
|
||||
Hash: hashStr,
|
||||
Type: "avatar",
|
||||
Policy: PolicyDedupNewRecord,
|
||||
Metadata: model.UploadMetadata{
|
||||
UserAgent: "first-agent",
|
||||
Extra: map[string]any{"record": "first"},
|
||||
},
|
||||
Policy: PolicyDedupNewRecord,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("first Ingest returned error: %v", err)
|
||||
@@ -149,6 +158,10 @@ func TestIngestPolicyDedupNewRecordCreatesSecondRecord(t *testing.T) {
|
||||
if putCount != 1 {
|
||||
t.Fatalf("putCount after first ingest = %d, want 1", putCount)
|
||||
}
|
||||
first.Upload.Metadata.Bucket = "shared-bucket"
|
||||
if err := dbConn.Save(&first.Upload).Error; err != nil {
|
||||
t.Fatalf("update first upload metadata failed: %v", err)
|
||||
}
|
||||
|
||||
second, err := Ingest(ctx, Request{
|
||||
UserID: 1002,
|
||||
@@ -159,7 +172,12 @@ func TestIngestPolicyDedupNewRecordCreatesSecondRecord(t *testing.T) {
|
||||
Extension: "png",
|
||||
Hash: hashStr,
|
||||
Type: "avatar",
|
||||
Policy: PolicyDedupNewRecord,
|
||||
Metadata: model.UploadMetadata{
|
||||
UserAgent: "second-agent",
|
||||
Bucket: "caller-bucket-must-not-survive",
|
||||
Extra: map[string]any{"record": "second"},
|
||||
},
|
||||
Policy: PolicyDedupNewRecord,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("second Ingest returned error: %v", err)
|
||||
@@ -173,6 +191,15 @@ func TestIngestPolicyDedupNewRecordCreatesSecondRecord(t *testing.T) {
|
||||
if first.Upload.ID == second.Upload.ID {
|
||||
t.Fatal("dedup records should have unique IDs")
|
||||
}
|
||||
if second.Upload.Metadata.Bucket != "shared-bucket" {
|
||||
t.Fatalf("dedup bucket = %q, want inherited shared-bucket", second.Upload.Metadata.Bucket)
|
||||
}
|
||||
if second.Upload.Metadata.UserAgent != "second-agent" {
|
||||
t.Fatalf("dedup user agent = %q, want caller metadata", second.Upload.Metadata.UserAgent)
|
||||
}
|
||||
if second.Upload.Metadata.Extra["record"] != "second" {
|
||||
t.Fatalf("dedup extra metadata = %#v, want caller metadata", second.Upload.Metadata.Extra)
|
||||
}
|
||||
|
||||
var count int64
|
||||
if err := dbConn.Model(&model.Upload{}).Where("hash = ?", hashStr).Count(&count).Error; err != nil {
|
||||
@@ -183,6 +210,77 @@ func TestIngestPolicyDedupNewRecordCreatesSecondRecord(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestDedupRecordFailureDoesNotDeleteSharedObject(t *testing.T) {
|
||||
dbConn, _, cleanup := testhelper.SetupTestEnvironment(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
|
||||
content := []byte("\x89PNG\r\n\x1a\nshared-object")
|
||||
hash := sha256.Sum256(content)
|
||||
hashStr := hex.EncodeToString(hash[:])
|
||||
deleteCount := 0
|
||||
restoreStorage, disableStorage := setupMockStorageWithDeleteCount(t, nil, &deleteCount)
|
||||
defer restoreStorage()
|
||||
defer disableStorage()
|
||||
|
||||
first, err := Ingest(ctx, Request{
|
||||
UserID: 1001,
|
||||
Reader: bytes.NewReader(content),
|
||||
Size: int64(len(content)),
|
||||
FileName: "shared.png",
|
||||
MimeType: "image/png",
|
||||
Extension: "png",
|
||||
Hash: hashStr,
|
||||
Type: "avatar",
|
||||
Policy: PolicyDedupNewRecord,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("first Ingest returned error: %v", err)
|
||||
}
|
||||
|
||||
const callbackName = "test:reject_dedup_upload_record"
|
||||
if err := dbConn.Callback().Create().Before("gorm:create").Register(callbackName, func(tx *gorm.DB) {
|
||||
upload, ok := tx.Statement.Dest.(*model.Upload)
|
||||
if ok && upload.FileName == "dedup-fail.png" {
|
||||
tx.AddError(errors.New("injected upload create failure"))
|
||||
}
|
||||
}); err != nil {
|
||||
t.Fatalf("register create failure callback: %v", err)
|
||||
}
|
||||
defer func() { _ = dbConn.Callback().Create().Remove(callbackName) }()
|
||||
|
||||
_, err = Ingest(ctx, Request{
|
||||
UserID: 1002,
|
||||
Reader: bytes.NewReader(content),
|
||||
Size: int64(len(content)),
|
||||
FileName: "dedup-fail.png",
|
||||
MimeType: "image/png",
|
||||
Extension: "png",
|
||||
Hash: hashStr,
|
||||
Type: "avatar",
|
||||
Metadata: model.UploadMetadata{
|
||||
Extra: map[string]any{"record": "dedup-failure"},
|
||||
},
|
||||
Policy: PolicyDedupNewRecord,
|
||||
})
|
||||
if err == nil {
|
||||
t.Fatal("dedup Ingest expected injected persistence error")
|
||||
}
|
||||
if deleteCount != 0 {
|
||||
t.Fatalf("shared object delete count = %d, want 0", deleteCount)
|
||||
}
|
||||
|
||||
_, backend, err := storage.Active(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("load active storage: %v", err)
|
||||
}
|
||||
obj, err := backend.Get(ctx, first.Upload.FilePath)
|
||||
if err != nil {
|
||||
t.Fatalf("shared object became unreadable after dedup failure: %v", err)
|
||||
}
|
||||
_ = obj.Body.Close()
|
||||
}
|
||||
|
||||
func TestCreateUploadWithStatsRollsBackOnCreateFailure(t *testing.T) {
|
||||
dbConn, _, cleanup := testhelper.SetupTestEnvironment(t)
|
||||
defer cleanup()
|
||||
@@ -259,6 +357,18 @@ func TestRemoveDecrementsStats(t *testing.T) {
|
||||
if _, err := Remove(ctx, result.Upload.ID); err != nil {
|
||||
t.Fatalf("Remove(%d) returned error: %v", result.Upload.ID, err)
|
||||
}
|
||||
stale := result.Upload
|
||||
uploadcache.SetUploadMetaCache(ctx, &stale)
|
||||
removedAgain, err := Remove(ctx, result.Upload.ID)
|
||||
if err != nil {
|
||||
t.Fatalf("second Remove(%d) returned error: %v", result.Upload.ID, err)
|
||||
}
|
||||
if removedAgain.Status != model.UploadStatusDeleted {
|
||||
t.Fatalf("second Remove status = %s, want deleted", removedAgain.Status)
|
||||
}
|
||||
if _, err := uploadcache.GetUploadByID(ctx, result.Upload.ID); !errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
t.Fatalf("cache lookup after idempotent Remove error = %v, want record not found", err)
|
||||
}
|
||||
|
||||
stats, err := loadTotalStats(ctx)
|
||||
if err != nil {
|
||||
@@ -269,6 +379,120 @@ func TestRemoveDecrementsStats(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestConcurrentRemoveDecrementsStatsOnce(t *testing.T) {
|
||||
_, _, cleanup := testhelper.SetupTestEnvironment(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
|
||||
content := []byte("\x89PNG\r\n\x1a\nconcurrent-remove")
|
||||
hash := sha256.Sum256(content)
|
||||
restoreStorage, disableStorage := setupMockStorage(t, nil)
|
||||
defer restoreStorage()
|
||||
defer disableStorage()
|
||||
|
||||
result, err := Ingest(ctx, Request{
|
||||
UserID: 1001,
|
||||
Reader: bytes.NewReader(content),
|
||||
Size: int64(len(content)),
|
||||
FileName: "concurrent.png",
|
||||
MimeType: "image/png",
|
||||
Extension: "png",
|
||||
Hash: hex.EncodeToString(hash[:]),
|
||||
Type: "generic",
|
||||
Policy: PolicyCreate,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("Ingest returned error: %v", err)
|
||||
}
|
||||
|
||||
const workers = 8
|
||||
start := make(chan struct{})
|
||||
errs := make(chan error, workers)
|
||||
var wg sync.WaitGroup
|
||||
for i := 0; i < workers; i++ {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
<-start
|
||||
_, removeErr := Remove(ctx, result.Upload.ID)
|
||||
errs <- removeErr
|
||||
}()
|
||||
}
|
||||
close(start)
|
||||
wg.Wait()
|
||||
close(errs)
|
||||
for removeErr := range errs {
|
||||
if removeErr != nil {
|
||||
t.Fatalf("concurrent Remove returned error: %v", removeErr)
|
||||
}
|
||||
}
|
||||
|
||||
stats, err := loadTotalStats(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("loadTotalStats returned error: %v", err)
|
||||
}
|
||||
if stats.TotalCount != 0 || stats.TotalSize != 0 {
|
||||
t.Fatalf("stats after concurrent remove = count %d size %d, want zero", stats.TotalCount, stats.TotalSize)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRemoveOwnedAndReservedTypeBoundaries(t *testing.T) {
|
||||
dbConn, _, cleanup := testhelper.SetupTestEnvironment(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
|
||||
ordinary := model.Upload{
|
||||
ID: 99101,
|
||||
UserID: 1001,
|
||||
FileName: "owned.txt",
|
||||
FilePath: "uploads/owned.txt",
|
||||
FileSize: 16,
|
||||
MimeType: "text/plain",
|
||||
Extension: "txt",
|
||||
Hash: "owned-hash",
|
||||
Type: "generic",
|
||||
Status: model.UploadStatusUsed,
|
||||
CreatedAt: time.Now(),
|
||||
}
|
||||
reserved := model.Upload{
|
||||
ID: 99102,
|
||||
UserID: 1001,
|
||||
FileName: "pages.zip",
|
||||
FilePath: "uploads/pages.zip",
|
||||
FileSize: 32,
|
||||
MimeType: "application/zip",
|
||||
Extension: "zip",
|
||||
Hash: "reserved-hash",
|
||||
Type: shared.ReservedPagesDeploymentType,
|
||||
Status: model.UploadStatusUsed,
|
||||
CreatedAt: time.Now(),
|
||||
}
|
||||
if err := dbConn.Create(&ordinary).Error; err != nil {
|
||||
t.Fatalf("seed ordinary upload: %v", err)
|
||||
}
|
||||
if err := dbConn.Create(&reserved).Error; err != nil {
|
||||
t.Fatalf("seed reserved upload: %v", err)
|
||||
}
|
||||
|
||||
if _, err := RemoveOwned(ctx, 2002, ordinary.ID); !errors.Is(err, ErrForbidden) {
|
||||
t.Fatalf("RemoveOwned non-owner error = %v, want ErrForbidden", err)
|
||||
}
|
||||
if _, err := Remove(ctx, reserved.ID); !errors.Is(err, ErrReservedUploadType) {
|
||||
t.Fatalf("Remove reserved error = %v, want ErrReservedUploadType", err)
|
||||
}
|
||||
if _, err := RemoveOwned(ctx, reserved.UserID, reserved.ID); !errors.Is(err, ErrReservedUploadType) {
|
||||
t.Fatalf("RemoveOwned reserved error = %v, want ErrReservedUploadType", err)
|
||||
}
|
||||
|
||||
var persisted model.Upload
|
||||
if err := dbConn.First(&persisted, reserved.ID).Error; err != nil {
|
||||
t.Fatalf("reload reserved upload: %v", err)
|
||||
}
|
||||
if persisted.Status != model.UploadStatusUsed {
|
||||
t.Fatalf("reserved upload status = %s, want used", persisted.Status)
|
||||
}
|
||||
}
|
||||
|
||||
type totalStatsSnapshot struct {
|
||||
TotalCount int64
|
||||
TotalSize int64
|
||||
@@ -289,6 +513,10 @@ func loadTotalStats(ctx context.Context) (totalStatsSnapshot, error) {
|
||||
}
|
||||
|
||||
func setupMockStorage(t *testing.T, putCount *int) (restore func(), disable func()) {
|
||||
return setupMockStorageWithDeleteCount(t, putCount, nil)
|
||||
}
|
||||
|
||||
func setupMockStorageWithDeleteCount(t *testing.T, putCount, deleteCount *int) (restore func(), disable func()) {
|
||||
t.Helper()
|
||||
mockFiles := make(map[string][]byte)
|
||||
restore = storage.MockStorage(
|
||||
@@ -316,6 +544,9 @@ func setupMockStorage(t *testing.T, putCount *int) (restore func(), disable func
|
||||
},
|
||||
func(ctx context.Context, key string) error {
|
||||
delete(mockFiles, key)
|
||||
if deleteCount != nil {
|
||||
*deleteCount++
|
||||
}
|
||||
return nil
|
||||
},
|
||||
)
|
||||
|
||||
@@ -7,52 +7,76 @@ import (
|
||||
"context"
|
||||
|
||||
uploadcache "github.com/Rain-kl/Wavelet/internal/apps/upload/cache"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/shared"
|
||||
uploadstats "github.com/Rain-kl/Wavelet/internal/apps/upload/stats"
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/repository"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/clause"
|
||||
)
|
||||
|
||||
// Remove soft-deletes an upload and decrements incremental stats.
|
||||
// Remove soft-deletes an ordinary upload and decrements incremental stats once.
|
||||
func Remove(ctx context.Context, uploadID uint64) (model.Upload, error) {
|
||||
upload, err := repository.GetActiveUploadByID(ctx, uploadID)
|
||||
upload, err := remove(ctx, 0, uploadID, false)
|
||||
if err != nil {
|
||||
return model.Upload{}, err
|
||||
}
|
||||
if err := softDeleteUploadWithStats(ctx, &upload); err != nil {
|
||||
return model.Upload{}, err
|
||||
}
|
||||
upload.Status = model.UploadStatusDeleted
|
||||
return upload, nil
|
||||
}
|
||||
|
||||
// RemoveOwned soft-deletes an upload owned by userID and decrements incremental stats.
|
||||
// RemoveOwned soft-deletes an ordinary upload owned by userID and decrements incremental stats once.
|
||||
func RemoveOwned(ctx context.Context, userID, uploadID uint64) (model.Upload, error) {
|
||||
upload, err := repository.GetActiveUploadByID(ctx, uploadID)
|
||||
upload, err := remove(ctx, userID, uploadID, true)
|
||||
if err != nil {
|
||||
return model.Upload{}, err
|
||||
}
|
||||
if upload.UserID != userID {
|
||||
return model.Upload{}, ErrForbidden
|
||||
}
|
||||
if err := softDeleteUploadWithStats(ctx, &upload); err != nil {
|
||||
return model.Upload{}, err
|
||||
}
|
||||
upload.Status = model.UploadStatusDeleted
|
||||
return upload, nil
|
||||
}
|
||||
|
||||
func softDeleteUploadWithStats(ctx context.Context, upload *model.Upload) error {
|
||||
statsSnapshot := *upload
|
||||
func remove(ctx context.Context, userID, uploadID uint64, owned bool) (model.Upload, error) {
|
||||
var upload model.Upload
|
||||
if err := db.DB(ctx).Transaction(func(tx *gorm.DB) error {
|
||||
if err := repository.SoftDeleteUploadTx(tx, upload); err != nil {
|
||||
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).
|
||||
Where("id = ?", uploadID).
|
||||
First(&upload).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
return uploadstats.ApplyUploadStatsDeltaTx(tx, &statsSnapshot, -1)
|
||||
}); err != nil {
|
||||
if owned && upload.UserID != userID {
|
||||
return ErrForbidden
|
||||
}
|
||||
if upload.Type == shared.ReservedPagesDeploymentType {
|
||||
return ErrReservedUploadType
|
||||
}
|
||||
_, err := RemoveLockedTx(tx, &upload)
|
||||
return err
|
||||
}); err != nil {
|
||||
return model.Upload{}, err
|
||||
}
|
||||
uploadcache.InvalidateUploadMetaCache(ctx, upload.ID)
|
||||
return nil
|
||||
|
||||
InvalidateUploadMetaCache(ctx, uploadID)
|
||||
upload.Status = model.UploadStatusDeleted
|
||||
return upload, nil
|
||||
}
|
||||
|
||||
// RemoveLockedTx performs the idempotent active-to-deleted transition for a row
|
||||
// that the caller has already locked in its surrounding transaction.
|
||||
func RemoveLockedTx(tx *gorm.DB, upload *model.Upload) (bool, error) {
|
||||
rowsAffected, err := repository.SoftDeleteUploadTx(tx, upload)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
if rowsAffected == 0 {
|
||||
return false, nil
|
||||
}
|
||||
if err := uploadstats.ApplyUploadStatsDeltaTx(tx, upload, -1); err != nil {
|
||||
return false, err
|
||||
}
|
||||
upload.Status = model.UploadStatusDeleted
|
||||
return true, nil
|
||||
}
|
||||
|
||||
// InvalidateUploadMetaCache invalidates upload metadata after the caller commits its transaction.
|
||||
func InvalidateUploadMetaCache(ctx context.Context, uploadID uint64) {
|
||||
uploadcache.InvalidateUploadMetaCache(ctx, uploadID)
|
||||
}
|
||||
|
||||
@@ -17,4 +17,6 @@ const (
|
||||
FileStatsTrendDays = 7
|
||||
MaxS3KeyLength = 1024
|
||||
AccessCacheTTL = 5 // seconds; multiplied by time.Second at use site
|
||||
// ReservedPagesDeploymentType is managed exclusively by the Pages domain.
|
||||
ReservedPagesDeploymentType = "openflare_pages_deployment"
|
||||
)
|
||||
|
||||
@@ -27,6 +27,7 @@ const (
|
||||
ErrQueryFileCountFailed = "查询文件数量失败"
|
||||
ErrQueryFileListFailed = "查询文件列表失败"
|
||||
ErrDeleteFileFailed = "删除文件失败"
|
||||
ErrReservedUploadType = "系统保留的文件类型不能通过通用文件接口操作"
|
||||
ErrStorageReadOnly = "存储迁移维护中,当前仅允许读取文件"
|
||||
ErrS3KeyRequired = "s3 key must not be empty"
|
||||
ErrS3KeyTooLongFormat = "s3 key exceeds maximum length of %d"
|
||||
|
||||
@@ -10,16 +10,15 @@ import (
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
uploadcache "github.com/Rain-kl/Wavelet/internal/apps/upload/cache"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/ingest"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/shared"
|
||||
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/db"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/storage"
|
||||
"github.com/Rain-kl/Wavelet/internal/task"
|
||||
"github.com/Rain-kl/Wavelet/pkg/logger"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/clause"
|
||||
)
|
||||
|
||||
const (
|
||||
@@ -77,32 +76,31 @@ func (h *SystemCleanupHandler) Execute(ctx context.Context, _ []byte) (*task.Tas
|
||||
|
||||
for _, u := range unusedUploads {
|
||||
totalProcessed++
|
||||
transitioned := false
|
||||
|
||||
if err := db.DB(ctx).Transaction(func(tx *gorm.DB) error {
|
||||
if err := tx.Model(&model.Upload{}).
|
||||
Where("id = ? AND status = ?", u.ID, model.UploadStatusPending).
|
||||
Update("status", model.UploadStatusDeleted).Error; err != nil {
|
||||
var locked model.Upload
|
||||
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).
|
||||
Where("id = ?", u.ID).
|
||||
First(&locked).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
_, backend, err := storage.Active(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
if locked.Status != model.UploadStatusPending || !locked.CreatedAt.Before(oneHourAgo) {
|
||||
return nil
|
||||
}
|
||||
if err := backend.Delete(ctx, u.FilePath); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
return nil
|
||||
var err error
|
||||
transitioned, err = ingest.RemoveLockedTx(tx, &locked)
|
||||
return err
|
||||
}); err != nil {
|
||||
task.AppendLog(ctx, "清理上传文件失败 [ID:%d]: %v", u.ID, err)
|
||||
lastID = u.ID
|
||||
continue
|
||||
}
|
||||
|
||||
uploadstats.RecordUploadStatsRemove(ctx, &u)
|
||||
uploadcache.InvalidateUploadMetaCache(ctx, u.ID)
|
||||
totalDeleted++
|
||||
ingest.InvalidateUploadMetaCache(ctx, u.ID)
|
||||
if transitioned {
|
||||
totalDeleted++
|
||||
}
|
||||
lastID = u.ID
|
||||
}
|
||||
}
|
||||
|
||||
@@ -19,6 +19,7 @@ import (
|
||||
|
||||
"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"
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
"github.com/Rain-kl/Wavelet/internal/diskcache"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
@@ -33,13 +34,17 @@ func TestSystemCleanupHandler_Execute(t *testing.T) {
|
||||
_, _, cleanup := testhelper.SetupTestEnvironment(t)
|
||||
defer cleanup()
|
||||
|
||||
// Mock S3 存储(让 DeleteObject 总是成功)
|
||||
deleteCount := 0
|
||||
// Mock S3 存储并记录 Delete,cleanup 不应物理删除共享对象。
|
||||
storageMock := storage.MockStorage(
|
||||
func(ctx context.Context, key string, body io.Reader, size int64, contentType string) error {
|
||||
return nil
|
||||
},
|
||||
func(ctx context.Context, key string) (*storage.Object, error) { return nil, nil },
|
||||
func(ctx context.Context, key string) error { return nil },
|
||||
func(ctx context.Context, key string) error {
|
||||
deleteCount++
|
||||
return nil
|
||||
},
|
||||
)
|
||||
defer storageMock()
|
||||
storage.IsEnabledFunc = func() bool { return true }
|
||||
@@ -86,6 +91,7 @@ func TestSystemCleanupHandler_Execute(t *testing.T) {
|
||||
for _, r := range records {
|
||||
err := db.DB(ctx).Create(r).Error
|
||||
require.NoError(t, err)
|
||||
require.NoError(t, uploadstats.ApplyUploadStatsAdd(ctx, r))
|
||||
}
|
||||
|
||||
// 准备推送历史测试数据:1个旧的(应删除),1个新的(应保留)
|
||||
@@ -147,6 +153,26 @@ func TestSystemCleanupHandler_Execute(t *testing.T) {
|
||||
var usedCount int64
|
||||
db.DB(ctx).Model(&model.Upload{}).Where("status = ?", model.UploadStatusUsed).Count(&usedCount)
|
||||
assert.Equal(t, int64(1), usedCount, "used 状态的文件不应受影响")
|
||||
assert.Equal(t, 0, deleteCount, "记录级 cleanup 不应调用 storage backend Delete")
|
||||
|
||||
var totalStats model.UploadStat
|
||||
err = db.DB(ctx).
|
||||
Where("dimension = ? AND stat_key = ?", model.UploadStatDimensionTotal, "").
|
||||
First(&totalStats).Error
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, int64(2), totalStats.FileCount, "cleanup 后统计只应保留 used 与最近 pending 记录")
|
||||
assert.Equal(t, int64(768), totalStats.FileSize, "cleanup 后统计大小应只扣减一次")
|
||||
|
||||
_, err = handler.Execute(ctx, nil)
|
||||
require.NoError(t, err)
|
||||
var statsAfterSecondRun model.UploadStat
|
||||
err = db.DB(ctx).
|
||||
Where("dimension = ? AND stat_key = ?", model.UploadStatDimensionTotal, "").
|
||||
First(&statsAfterSecondRun).Error
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, totalStats.FileCount, statsAfterSecondRun.FileCount, "重复 cleanup 不应再次扣减统计")
|
||||
assert.Equal(t, totalStats.FileSize, statsAfterSecondRun.FileSize, "重复 cleanup 不应再次扣减统计大小")
|
||||
assert.Equal(t, 0, deleteCount, "重复 cleanup 仍不应调用 storage backend Delete")
|
||||
|
||||
// 验证推送历史数据状态:10天前的应被删除,今天的应保留
|
||||
var pushCount int64
|
||||
|
||||
Reference in New Issue
Block a user