diff --git a/docs/docs.go b/docs/docs.go index fc60b319..c1870837 100644 --- a/docs/docs.go +++ b/docs/docs.go @@ -3367,7 +3367,7 @@ const docTemplate = `{ "type": "object", "properties": { "data": { - "$ref": "#/definitions/upload.listFilesResponse" + "$ref": "#/definitions/handler.listFilesResponse" } } } @@ -3414,7 +3414,7 @@ const docTemplate = `{ "in": "body", "required": true, "schema": { - "$ref": "#/definitions/upload.batchDownloadRequest" + "$ref": "#/definitions/handler.batchDownloadRequest" } } ], @@ -3525,7 +3525,7 @@ const docTemplate = `{ "type": "object", "properties": { "data": { - "$ref": "#/definitions/upload.fileStatsResponse" + "$ref": "#/definitions/handler.fileStatsResponse" } } } @@ -4654,7 +4654,7 @@ const docTemplate = `{ "type": "object", "properties": { "data": { - "$ref": "#/definitions/upload.listMyFilesResponse" + "$ref": "#/definitions/handler.listMyFilesResponse" } } } @@ -4702,7 +4702,7 @@ const docTemplate = `{ "in": "body", "required": true, "schema": { - "$ref": "#/definitions/upload.updateMyFileRequest" + "$ref": "#/definitions/handler.updateMyFileRequest" } } ], @@ -5688,6 +5688,134 @@ const docTemplate = `{ } } }, + "handler.batchDownloadRequest": { + "type": "object", + "required": [ + "ids" + ], + "properties": { + "ids": { + "type": "array", + "minItems": 1, + "items": { + "type": "string" + } + } + } + }, + "handler.distributionItem": { + "type": "object", + "properties": { + "count": { + "type": "integer" + }, + "name": { + "type": "string" + }, + "size": { + "type": "integer" + } + } + }, + "handler.fileStatsResponse": { + "type": "object", + "properties": { + "categories": { + "type": "array", + "items": { + "$ref": "#/definitions/handler.distributionItem" + } + }, + "total_count": { + "type": "integer" + }, + "total_size": { + "type": "integer" + }, + "trend": { + "type": "array", + "items": { + "$ref": "#/definitions/handler.trendItem" + } + }, + "types": { + "type": "array", + "items": { + "$ref": "#/definitions/handler.distributionItem" + } + } + } + }, + "handler.listFilesResponse": { + "type": "object", + "properties": { + "items": { + "type": "array", + "items": { + "$ref": "#/definitions/model.Upload" + } + }, + "page": { + "type": "integer" + }, + "page_size": { + "type": "integer" + }, + "total": { + "type": "integer" + } + } + }, + "handler.listMyFilesResponse": { + "type": "object", + "properties": { + "items": { + "type": "array", + "items": { + "$ref": "#/definitions/model.Upload" + } + }, + "page": { + "type": "integer" + }, + "page_size": { + "type": "integer" + }, + "total": { + "type": "integer" + } + } + }, + "handler.trendItem": { + "type": "object", + "properties": { + "count": { + "type": "integer" + }, + "date": { + "type": "string" + }, + "size": { + "type": "integer" + } + } + }, + "handler.updateMyFileRequest": { + "type": "object", + "properties": { + "access_mode": { + "type": "integer", + "enum": [ + 0, + 1 + ] + }, + "file_name": { + "type": "string", + "maxLength": 255 + } + } + }, "logger.LogEntry": { "type": "object", "properties": { @@ -6280,10 +6408,6 @@ const docTemplate = `{ } ] }, - "storage_driver": { - "description": "存储引擎驱动 (如 local, s3, oss)", - "type": "string" - }, "type": { "description": "业务标识类型 (如 avatar, doc, attachment)", "type": "string" @@ -7197,134 +7321,6 @@ const docTemplate = `{ } } }, - "upload.batchDownloadRequest": { - "type": "object", - "required": [ - "ids" - ], - "properties": { - "ids": { - "type": "array", - "minItems": 1, - "items": { - "type": "string" - } - } - } - }, - "upload.distributionItem": { - "type": "object", - "properties": { - "count": { - "type": "integer" - }, - "name": { - "type": "string" - }, - "size": { - "type": "integer" - } - } - }, - "upload.fileStatsResponse": { - "type": "object", - "properties": { - "categories": { - "type": "array", - "items": { - "$ref": "#/definitions/upload.distributionItem" - } - }, - "total_count": { - "type": "integer" - }, - "total_size": { - "type": "integer" - }, - "trend": { - "type": "array", - "items": { - "$ref": "#/definitions/upload.trendItem" - } - }, - "types": { - "type": "array", - "items": { - "$ref": "#/definitions/upload.distributionItem" - } - } - } - }, - "upload.listFilesResponse": { - "type": "object", - "properties": { - "items": { - "type": "array", - "items": { - "$ref": "#/definitions/model.Upload" - } - }, - "page": { - "type": "integer" - }, - "page_size": { - "type": "integer" - }, - "total": { - "type": "integer" - } - } - }, - "upload.listMyFilesResponse": { - "type": "object", - "properties": { - "items": { - "type": "array", - "items": { - "$ref": "#/definitions/model.Upload" - } - }, - "page": { - "type": "integer" - }, - "page_size": { - "type": "integer" - }, - "total": { - "type": "integer" - } - } - }, - "upload.trendItem": { - "type": "object", - "properties": { - "count": { - "type": "integer" - }, - "date": { - "type": "string" - }, - "size": { - "type": "integer" - } - } - }, - "upload.updateMyFileRequest": { - "type": "object", - "properties": { - "access_mode": { - "type": "integer", - "enum": [ - 0, - 1 - ] - }, - "file_name": { - "type": "string", - "maxLength": 255 - } - } - }, "user.changePasswordRequest": { "type": "object", "properties": { diff --git a/docs/swagger.json b/docs/swagger.json index 36f9bc59..07f5eb39 100644 --- a/docs/swagger.json +++ b/docs/swagger.json @@ -3360,7 +3360,7 @@ "type": "object", "properties": { "data": { - "$ref": "#/definitions/upload.listFilesResponse" + "$ref": "#/definitions/handler.listFilesResponse" } } } @@ -3407,7 +3407,7 @@ "in": "body", "required": true, "schema": { - "$ref": "#/definitions/upload.batchDownloadRequest" + "$ref": "#/definitions/handler.batchDownloadRequest" } } ], @@ -3518,7 +3518,7 @@ "type": "object", "properties": { "data": { - "$ref": "#/definitions/upload.fileStatsResponse" + "$ref": "#/definitions/handler.fileStatsResponse" } } } @@ -4647,7 +4647,7 @@ "type": "object", "properties": { "data": { - "$ref": "#/definitions/upload.listMyFilesResponse" + "$ref": "#/definitions/handler.listMyFilesResponse" } } } @@ -4695,7 +4695,7 @@ "in": "body", "required": true, "schema": { - "$ref": "#/definitions/upload.updateMyFileRequest" + "$ref": "#/definitions/handler.updateMyFileRequest" } } ], @@ -5681,6 +5681,134 @@ } } }, + "handler.batchDownloadRequest": { + "type": "object", + "required": [ + "ids" + ], + "properties": { + "ids": { + "type": "array", + "minItems": 1, + "items": { + "type": "string" + } + } + } + }, + "handler.distributionItem": { + "type": "object", + "properties": { + "count": { + "type": "integer" + }, + "name": { + "type": "string" + }, + "size": { + "type": "integer" + } + } + }, + "handler.fileStatsResponse": { + "type": "object", + "properties": { + "categories": { + "type": "array", + "items": { + "$ref": "#/definitions/handler.distributionItem" + } + }, + "total_count": { + "type": "integer" + }, + "total_size": { + "type": "integer" + }, + "trend": { + "type": "array", + "items": { + "$ref": "#/definitions/handler.trendItem" + } + }, + "types": { + "type": "array", + "items": { + "$ref": "#/definitions/handler.distributionItem" + } + } + } + }, + "handler.listFilesResponse": { + "type": "object", + "properties": { + "items": { + "type": "array", + "items": { + "$ref": "#/definitions/model.Upload" + } + }, + "page": { + "type": "integer" + }, + "page_size": { + "type": "integer" + }, + "total": { + "type": "integer" + } + } + }, + "handler.listMyFilesResponse": { + "type": "object", + "properties": { + "items": { + "type": "array", + "items": { + "$ref": "#/definitions/model.Upload" + } + }, + "page": { + "type": "integer" + }, + "page_size": { + "type": "integer" + }, + "total": { + "type": "integer" + } + } + }, + "handler.trendItem": { + "type": "object", + "properties": { + "count": { + "type": "integer" + }, + "date": { + "type": "string" + }, + "size": { + "type": "integer" + } + } + }, + "handler.updateMyFileRequest": { + "type": "object", + "properties": { + "access_mode": { + "type": "integer", + "enum": [ + 0, + 1 + ] + }, + "file_name": { + "type": "string", + "maxLength": 255 + } + } + }, "logger.LogEntry": { "type": "object", "properties": { @@ -6273,10 +6401,6 @@ } ] }, - "storage_driver": { - "description": "存储引擎驱动 (如 local, s3, oss)", - "type": "string" - }, "type": { "description": "业务标识类型 (如 avatar, doc, attachment)", "type": "string" @@ -7190,134 +7314,6 @@ } } }, - "upload.batchDownloadRequest": { - "type": "object", - "required": [ - "ids" - ], - "properties": { - "ids": { - "type": "array", - "minItems": 1, - "items": { - "type": "string" - } - } - } - }, - "upload.distributionItem": { - "type": "object", - "properties": { - "count": { - "type": "integer" - }, - "name": { - "type": "string" - }, - "size": { - "type": "integer" - } - } - }, - "upload.fileStatsResponse": { - "type": "object", - "properties": { - "categories": { - "type": "array", - "items": { - "$ref": "#/definitions/upload.distributionItem" - } - }, - "total_count": { - "type": "integer" - }, - "total_size": { - "type": "integer" - }, - "trend": { - "type": "array", - "items": { - "$ref": "#/definitions/upload.trendItem" - } - }, - "types": { - "type": "array", - "items": { - "$ref": "#/definitions/upload.distributionItem" - } - } - } - }, - "upload.listFilesResponse": { - "type": "object", - "properties": { - "items": { - "type": "array", - "items": { - "$ref": "#/definitions/model.Upload" - } - }, - "page": { - "type": "integer" - }, - "page_size": { - "type": "integer" - }, - "total": { - "type": "integer" - } - } - }, - "upload.listMyFilesResponse": { - "type": "object", - "properties": { - "items": { - "type": "array", - "items": { - "$ref": "#/definitions/model.Upload" - } - }, - "page": { - "type": "integer" - }, - "page_size": { - "type": "integer" - }, - "total": { - "type": "integer" - } - } - }, - "upload.trendItem": { - "type": "object", - "properties": { - "count": { - "type": "integer" - }, - "date": { - "type": "string" - }, - "size": { - "type": "integer" - } - } - }, - "upload.updateMyFileRequest": { - "type": "object", - "properties": { - "access_mode": { - "type": "integer", - "enum": [ - 0, - 1 - ] - }, - "file_name": { - "type": "string", - "maxLength": 255 - } - } - }, "user.changePasswordRequest": { "type": "object", "properties": { diff --git a/docs/swagger.yaml b/docs/swagger.yaml index 55e8210a..0554159e 100644 --- a/docs/swagger.yaml +++ b/docs/swagger.yaml @@ -143,6 +143,90 @@ definitions: token: type: string type: object + handler.batchDownloadRequest: + properties: + ids: + items: + type: string + minItems: 1 + type: array + required: + - ids + type: object + handler.distributionItem: + properties: + count: + type: integer + name: + type: string + size: + type: integer + type: object + handler.fileStatsResponse: + properties: + categories: + items: + $ref: '#/definitions/handler.distributionItem' + type: array + total_count: + type: integer + total_size: + type: integer + trend: + items: + $ref: '#/definitions/handler.trendItem' + type: array + types: + items: + $ref: '#/definitions/handler.distributionItem' + type: array + type: object + handler.listFilesResponse: + properties: + items: + items: + $ref: '#/definitions/model.Upload' + type: array + page: + type: integer + page_size: + type: integer + total: + type: integer + type: object + handler.listMyFilesResponse: + properties: + items: + items: + $ref: '#/definitions/model.Upload' + type: array + page: + type: integer + page_size: + type: integer + total: + type: integer + type: object + handler.trendItem: + properties: + count: + type: integer + date: + type: string + size: + type: integer + type: object + handler.updateMyFileRequest: + properties: + access_mode: + enum: + - 0 + - 1 + type: integer + file_name: + maxLength: 255 + type: string + type: object logger.LogEntry: properties: data: @@ -541,9 +625,6 @@ definitions: allOf: - $ref: '#/definitions/model.UploadStatus' description: 状态 - storage_driver: - description: 存储引擎驱动 (如 local, s3, oss) - type: string type: description: 业务标识类型 (如 avatar, doc, attachment) type: string @@ -1169,90 +1250,6 @@ definitions: upstream_repository: type: string type: object - upload.batchDownloadRequest: - properties: - ids: - items: - type: string - minItems: 1 - type: array - required: - - ids - type: object - upload.distributionItem: - properties: - count: - type: integer - name: - type: string - size: - type: integer - type: object - upload.fileStatsResponse: - properties: - categories: - items: - $ref: '#/definitions/upload.distributionItem' - type: array - total_count: - type: integer - total_size: - type: integer - trend: - items: - $ref: '#/definitions/upload.trendItem' - type: array - types: - items: - $ref: '#/definitions/upload.distributionItem' - type: array - type: object - upload.listFilesResponse: - properties: - items: - items: - $ref: '#/definitions/model.Upload' - type: array - page: - type: integer - page_size: - type: integer - total: - type: integer - type: object - upload.listMyFilesResponse: - properties: - items: - items: - $ref: '#/definitions/model.Upload' - type: array - page: - type: integer - page_size: - type: integer - total: - type: integer - type: object - upload.trendItem: - properties: - count: - type: integer - date: - type: string - size: - type: integer - type: object - upload.updateMyFileRequest: - properties: - access_mode: - enum: - - 0 - - 1 - type: integer - file_name: - maxLength: 255 - type: string - type: object user.changePasswordRequest: properties: new_password: @@ -3409,7 +3406,7 @@ paths: - $ref: '#/definitions/response.Any' - properties: data: - $ref: '#/definitions/upload.listFilesResponse' + $ref: '#/definitions/handler.listFilesResponse' type: object "401": description: 未登录 @@ -3501,7 +3498,7 @@ paths: name: request required: true schema: - $ref: '#/definitions/upload.batchDownloadRequest' + $ref: '#/definitions/handler.batchDownloadRequest' produces: - application/octet-stream responses: @@ -3535,7 +3532,7 @@ paths: - $ref: '#/definitions/response.Any' - properties: data: - $ref: '#/definitions/upload.fileStatsResponse' + $ref: '#/definitions/handler.fileStatsResponse' type: object "401": description: 未登录 @@ -4193,7 +4190,7 @@ paths: name: request required: true schema: - $ref: '#/definitions/upload.updateMyFileRequest' + $ref: '#/definitions/handler.updateMyFileRequest' produces: - application/json responses: @@ -4253,7 +4250,7 @@ paths: - $ref: '#/definitions/response.Any' - properties: data: - $ref: '#/definitions/upload.listMyFilesResponse' + $ref: '#/definitions/handler.listMyFilesResponse' type: object "401": description: 未登录 diff --git a/frontend/app/(main)/admin/files/components/file-list.tsx b/frontend/app/(main)/admin/files/components/file-list.tsx index fc217630..d5f7047f 100644 --- a/frontend/app/(main)/admin/files/components/file-list.tsx +++ b/frontend/app/(main)/admin/files/components/file-list.tsx @@ -85,6 +85,15 @@ export function FileList() { queryFn: () => services.adminUpload.listUploads(page, pageSize, debouncedKeyword || undefined), }) + const storageDriverQuery = useQuery({ + queryKey: ["admin", "storage-config", "driver"], + queryFn: async () => { + const record = await services.adminSystemConfig.getSystemConfig("storage_config") + const cfg = JSON.parse(record.value) as {driver?: string} + return cfg.driver ?? "local" + }, + }) + const files = listQuery.data?.items ?? [] const total = listQuery.data?.total ?? 0 const totalPages = Math.ceil(total / pageSize) @@ -470,7 +479,7 @@ export function FileList() {
存储驱动 - {detailTarget.storage_driver || "local"} + {storageDriverQuery.data ?? "local"}
diff --git a/frontend/lib/services/upload/types.ts b/frontend/lib/services/upload/types.ts index 20c13cc4..f97478da 100644 --- a/frontend/lib/services/upload/types.ts +++ b/frontend/lib/services/upload/types.ts @@ -24,7 +24,6 @@ export interface Upload { mime_type: string extension: string hash: string - storage_driver: string type: string status: string access_mode: number diff --git a/internal/apps/admin/system_config/errs.go b/internal/apps/admin/system_config/errs.go index 408142d3..ed2d6e1e 100644 --- a/internal/apps/admin/system_config/errs.go +++ b/internal/apps/admin/system_config/errs.go @@ -11,4 +11,5 @@ const ( ConfigKeyRequired = "配置键不能为空" ConfigValueRequired = "配置值不能为空" ConfigKeyExists = "配置键已存在" + StorageDriverSwitchRequiresMigration = "存在存量文件,请通过存储迁移任务切换存储引擎" ) diff --git a/internal/apps/admin/system_config/logics.go b/internal/apps/admin/system_config/logics.go index f3e9d3a1..04630e82 100644 --- a/internal/apps/admin/system_config/logics.go +++ b/internal/apps/admin/system_config/logics.go @@ -7,7 +7,6 @@ import ( "context" "encoding/json" "errors" - "fmt" "time" "github.com/Rain-kl/Wavelet/internal/db" @@ -88,9 +87,6 @@ func updateSystemConfig(ctx context.Context, key string, req UpdateSystemConfigR if err := tx.Model(&config).Updates(updates).Error; err != nil { return err } - if err := repointUploadStorageDriversOnDriverSwitch(ctx, tx, key, originalDriver, req.Value); err != nil { - return err - } resolveStorageMigrationTasksOnDirectDriverUpdate(ctx, tx, key, originalDriver, req.Value) return nil }); err != nil { @@ -101,38 +97,6 @@ func updateSystemConfig(ctx context.Context, key string, req UpdateSystemConfigR return nil } -func repointUploadStorageDriversOnDriverSwitch( - ctx context.Context, - tx *gorm.DB, - key string, - originalDriver storage.Driver, - newValue string, -) error { - if key != model.ConfigKeyStorageConfig || originalDriver == "" { - return nil - } - - var newCfg storage.Config - if err := json.Unmarshal([]byte(newValue), &newCfg); err != nil { - return fmt.Errorf("parse storage config for driver repoint: %w", err) - } - if newCfg.Driver == "" || newCfg.Driver == originalDriver { - return nil - } - - result := tx.Model(&model.Upload{}). - Where("storage_driver = ? AND status != ?", string(originalDriver), model.UploadStatusDeleted). - Update("storage_driver", string(newCfg.Driver)) - if result.Error != nil { - return result.Error - } - if result.RowsAffected > 0 { - logger.InfoF(ctx, "[StorageConfig] switched driver %s -> %s, repointed %d upload records", - originalDriver, newCfg.Driver, result.RowsAffected) - } - return nil -} - func resolveStorageMigrationTasksOnDirectDriverUpdate( ctx context.Context, tx *gorm.DB, diff --git a/internal/apps/admin/system_config/routers.go b/internal/apps/admin/system_config/routers.go index 5d87e247..336faa47 100644 --- a/internal/apps/admin/system_config/routers.go +++ b/internal/apps/admin/system_config/routers.go @@ -18,6 +18,7 @@ import ( "github.com/Rain-kl/Wavelet/internal/apps/cap" "github.com/Rain-kl/Wavelet/internal/apps/upload" "github.com/Rain-kl/Wavelet/internal/common/response" + "github.com/Rain-kl/Wavelet/internal/db" "github.com/Rain-kl/Wavelet/internal/model" "github.com/Rain-kl/Wavelet/internal/repository" "github.com/Rain-kl/Wavelet/internal/storage" @@ -266,11 +267,13 @@ func TestSMTP(c *gin.Context) { func isStorageConfigValidationError(err error) bool { msg := err.Error() - return strings.HasPrefix(msg, "解析") || + return msg == StorageDriverSwitchRequiresMigration || + strings.HasPrefix(msg, "解析") || strings.HasPrefix(msg, "验证") || strings.HasPrefix(msg, "初始化测试") || strings.HasPrefix(msg, "存储连通性") || - strings.HasPrefix(msg, "序列化") + strings.HasPrefix(msg, "序列化") || + strings.HasPrefix(msg, "检查存量文件") } func maskSensitiveConfig(key, value string) string { @@ -322,6 +325,15 @@ func validateAndMergeStorageConfig(ctx context.Context, value string, currentCon func validateMergedStorageConfig(ctx context.Context, currentCfg, newCfg, targetCfg storage.Config) error { if newCfg.Driver != "" && newCfg.Driver != currentCfg.Driver { + var uploadCount int64 + if err := db.DB(ctx).Model(&model.Upload{}). + Where("status != ?", model.UploadStatusDeleted). + Count(&uploadCount).Error; err != nil { + return fmt.Errorf("检查存量文件失败: %w", err) + } + if uploadCount > 0 { + return errors.New(StorageDriverSwitchRequiresMigration) + } if err := validateDriverConfig(targetCfg, newCfg.Driver); err != nil { return fmt.Errorf("验证目标存储配置参数失败: %w", err) } diff --git a/internal/apps/admin/system_config/routers_test.go b/internal/apps/admin/system_config/routers_test.go index 02e19948..f268faa8 100644 --- a/internal/apps/admin/system_config/routers_test.go +++ b/internal/apps/admin/system_config/routers_test.go @@ -13,6 +13,7 @@ import ( "net/http" "net/http/httptest" "net/textproto" + "strings" "testing" "github.com/Rain-kl/Wavelet/internal/apps/oauth" @@ -441,7 +442,53 @@ func TestUpdateStorageConfigValidation(t *testing.T) { } }) + t.Run("reject driver switch when uploads exist", func(t *testing.T) { + upload := model.Upload{ + ID: 88001, + UserID: 1, + FileName: "keep.txt", + FilePath: "uploads/keep.txt", + FileSize: 4, + MimeType: "text/plain", + Extension: "txt", + Type: "attachment", + Status: model.UploadStatusUsed, + } + if err := dbConn.Create(&upload).Error; err != nil { + t.Fatalf("seed upload failed: %v", err) + } + + tempDir := t.TempDir() + cfg := storage.DefaultConfig() + cfg.Driver = storage.DriverS3 + cfg.S3.Endpoint = "http://127.0.0.1:19998" + cfg.S3.Region = "us-east-1" + cfg.S3.Bucket = "wavelet" + cfg.S3.AccessKeyID = "test" + cfg.S3.SecretAccessKey = "test" + cfg.Local.Root = tempDir + + cfgBytes, _ := json.Marshal(cfg) + payload := UpdateSystemConfigRequest{Value: string(cfgBytes)} + body, _ := json.Marshal(payload) + req, _ := http.NewRequest("PUT", "/api/v1/admin/system-configs/storage_config", bytes.NewBuffer(body)) + req.Header.Set("Content-Type", "application/json") + w := httptest.NewRecorder() + router.ServeHTTP(w, req) + + if w.Code != http.StatusBadRequest { + t.Fatalf("expected 400 Bad Request, got %d. Body: %s", w.Code, w.Body.String()) + } + if !strings.Contains(w.Body.String(), StorageDriverSwitchRequiresMigration) { + t.Fatalf("expected migration-required error, got: %s", w.Body.String()) + } + }) + t.Run("switch to local while active s3 is unreachable", func(t *testing.T) { + if err := dbConn.Where("1 = 1").Delete(&model.Upload{}).Error; err != nil { + t.Fatalf("clear uploads failed: %v", err) + } + activeCfg := storage.DefaultConfig() activeCfg.Driver = storage.DriverS3 activeCfg.S3.Endpoint = "http://127.0.0.1:9999" diff --git a/internal/apps/upload/filesrv/file_server_test.go b/internal/apps/upload/filesrv/file_server_test.go index 5dc426e2..aa8b7027 100644 --- a/internal/apps/upload/filesrv/file_server_test.go +++ b/internal/apps/upload/filesrv/file_server_test.go @@ -67,7 +67,6 @@ func TestServeFileByIDAccessControl(t *testing.T) { FileSize: 5, MimeType: "image/png", Extension: "png", - StorageDriver: "local", Type: "avatar", Status: model.UploadStatusUsed, AccessMode: 1, @@ -80,7 +79,6 @@ func TestServeFileByIDAccessControl(t *testing.T) { FileSize: 5, MimeType: "application/pdf", Extension: "pdf", - StorageDriver: "local", Type: "attachment", Status: model.UploadStatusUsed, AccessMode: 1, @@ -207,7 +205,6 @@ func TestImageCompression(t *testing.T) { FileSize: int64(pngBuf.Len()), MimeType: "image/png", Extension: "png", - StorageDriver: "local", Type: "avatar", // Whitelisted by default Status: model.UploadStatusUsed, AccessMode: 1, diff --git a/internal/apps/upload/handler/file_management_test.go b/internal/apps/upload/handler/file_management_test.go index 7607b526..7387cd38 100644 --- a/internal/apps/upload/handler/file_management_test.go +++ b/internal/apps/upload/handler/file_management_test.go @@ -29,7 +29,6 @@ func TestGetDistinctUploadTypes(t *testing.T) { FileSize: 10, MimeType: "text/plain", Extension: "txt", - StorageDriver: "local", Type: "custom_type_xyz", Status: model.UploadStatusUsed, } diff --git a/internal/apps/upload/handler/logics.go b/internal/apps/upload/handler/logics.go index 7556c6cf..d0812d79 100644 --- a/internal/apps/upload/handler/logics.go +++ b/internal/apps/upload/handler/logics.go @@ -124,7 +124,6 @@ func createInstantUpload(ctx context.Context, existing model.Upload, input insta MimeType: input.MimeType, Extension: input.Extension, Hash: input.FileHash, - StorageDriver: existing.StorageDriver, Type: input.UploadType, Status: model.UploadStatusUsed, AccessMode: input.AccessMode, @@ -142,9 +141,9 @@ func findReusableUpload(ctx context.Context, hash string, size int64) (model.Upl return repository.FindReusableUploadByHash(ctx, hash, size) } -func saveNewUploadRecord(ctx context.Context, upload *model.Upload, storageDriver, filePath string) error { +func saveNewUploadRecord(ctx context.Context, upload *model.Upload, filePath string) error { if err := repository.CreateUpload(ctx, upload); err != nil { - backend, backendErr := storage.ForDriver(ctx, storage.Driver(storageDriver)) + _, backend, backendErr := storage.Active(ctx) if backendErr == nil { if deleteErr := backend.Delete(ctx, filePath); deleteErr != nil { logger.WarnF(ctx, "清理未写入数据库的上传对象失败: %v", deleteErr) @@ -162,22 +161,22 @@ func loadUploadStats(ctx context.Context) ([]model.UploadStat, error) { var errUploadForbidden = errors.New("upload forbidden") -func storeUploadObject(ctx context.Context, subPath string, size int64, mimeType string, buf *bytes.Buffer, meta *model.UploadMetadata) (string, string, error) { +func storeUploadObject(ctx context.Context, subPath string, size int64, mimeType string, buf *bytes.Buffer, meta *model.UploadMetadata) (string, error) { if uploadstorage.ReadOnly(ctx) { - return "", "", errors.New(shared.ErrStorageReadOnly) + return "", errors.New(shared.ErrStorageReadOnly) } driver, backend, err := storage.Active(ctx) if err != nil { logger.ErrorF(ctx, "初始化活动存储失败: %v", err) - return "", "", errors.New(shared.ErrSaveFileFailed) + return "", errors.New(shared.ErrSaveFileFailed) } result, err := backend.Put(ctx, subPath, bytes.NewReader(buf.Bytes()), size, mimeType) if err != nil { logger.ErrorF(ctx, "写入 %s 存储失败: %v", driver, err) - return "", "", errors.New(shared.ErrSaveFileFailed) + return "", errors.New(shared.ErrSaveFileFailed) } meta.Bucket = result.Bucket - return string(driver), result.Key, nil + return result.Key, nil } func validateUploadAllowedExtension(ctx context.Context, ext string) string { diff --git a/internal/apps/upload/handler/routers.go b/internal/apps/upload/handler/routers.go index e1442bcf..57ff9c01 100644 --- a/internal/apps/upload/handler/routers.go +++ b/internal/apps/upload/handler/routers.go @@ -138,29 +138,28 @@ func UploadFile(c *gin.Context) { id := idgen.NextUint64ID() subPath := fmt.Sprintf("uploads/%s/%d.%s", time.Now().Format("2006/01/02"), id, ext) - storageDriver, subPath, err := storeUploadObject(ctx, subPath, size, mimeType, &buf, &meta) + subPath, err = storeUploadObject(ctx, subPath, size, mimeType, &buf, &meta) if err != nil { response.AbortBadRequest(c, err.Error()) return } newUpload := model.Upload{ - ID: id, - UserID: currUser.ID, - FileName: origName, - FilePath: subPath, - FileSize: size, - MimeType: mimeType, - Extension: ext, - Hash: fileHash, - StorageDriver: storageDriver, - Type: uploadType, - Status: model.UploadStatusUsed, - AccessMode: accessMode, - Metadata: meta, + ID: id, + UserID: currUser.ID, + FileName: origName, + FilePath: subPath, + FileSize: size, + MimeType: mimeType, + Extension: ext, + Hash: fileHash, + Type: uploadType, + Status: model.UploadStatusUsed, + AccessMode: accessMode, + Metadata: meta, } - if err := saveNewUploadRecord(ctx, &newUpload, storageDriver, subPath); err != nil { + if err := saveNewUploadRecord(ctx, &newUpload, subPath); err != nil { response.AbortBadRequest(c, shared.ErrSaveUploadRecordFailed) return } diff --git a/internal/apps/upload/handler/routers_test.go b/internal/apps/upload/handler/routers_test.go index d944d1ea..4545cb0a 100644 --- a/internal/apps/upload/handler/routers_test.go +++ b/internal/apps/upload/handler/routers_test.go @@ -193,10 +193,6 @@ func TestUploadFile(t *testing.T) { t.Errorf("incorrect mime type detected: %s", dbRecord.MimeType) } - if dbRecord.StorageDriver != "s3" { - t.Errorf("expected storage driver s3, got %s", dbRecord.StorageDriver) - } - if dbRecord.Metadata.Extra["source"] != "test_runner" { t.Errorf("expected extra meta 'source' to be 'test_runner', got %v", dbRecord.Metadata.Extra) } @@ -325,10 +321,6 @@ func TestUploadFile(t *testing.T) { t.Fatalf("failed to unmarshal local upload record: %v", err) } - if localRecord.StorageDriver != "local" { - t.Errorf("expected storage driver local, got %s", localRecord.StorageDriver) - } - // Confirm file was actually written to local disk fileContent, err := os.ReadFile(localRecord.FilePath) if err != nil { @@ -358,7 +350,6 @@ func TestDownloadFile(t *testing.T) { FileSize: 12, MimeType: "text/plain", Extension: "txt", - StorageDriver: "local", Status: model.UploadStatusUsed, } @@ -426,7 +417,6 @@ func TestListFiles(t *testing.T) { FileSize: 10, MimeType: "text/plain", Extension: "txt", - StorageDriver: "local", Status: model.UploadStatusUsed, }, { @@ -437,7 +427,6 @@ func TestListFiles(t *testing.T) { FileSize: 20, MimeType: "image/png", Extension: "png", - StorageDriver: "local", Status: model.UploadStatusUsed, }, { @@ -448,7 +437,6 @@ func TestListFiles(t *testing.T) { FileSize: 30, MimeType: "text/markdown", Extension: "md", - StorageDriver: "local", Status: model.UploadStatusUsed, }, { @@ -459,7 +447,6 @@ func TestListFiles(t *testing.T) { FileSize: 40, MimeType: "text/plain", Extension: "txt", - StorageDriver: "local", Status: model.UploadStatusUsed, }, } @@ -582,7 +569,6 @@ func TestBatchDownloadFiles(t *testing.T) { FileSize: 13, MimeType: "text/plain", Extension: "txt", - StorageDriver: "local", Status: model.UploadStatusUsed, }, { @@ -593,7 +579,6 @@ func TestBatchDownloadFiles(t *testing.T) { FileSize: 13, MimeType: "text/plain", Extension: "txt", - StorageDriver: "local", Status: model.UploadStatusUsed, }, { @@ -604,7 +589,6 @@ func TestBatchDownloadFiles(t *testing.T) { FileSize: 28, MimeType: "text/plain", Extension: "txt", - StorageDriver: "local", Status: model.UploadStatusUsed, }, } @@ -773,7 +757,6 @@ func TestGetFileStats(t *testing.T) { FileSize: 100, MimeType: "image/png", Extension: "png", - StorageDriver: "local", Type: "generic", Status: model.UploadStatusUsed, CreatedAt: time.Now(), @@ -786,7 +769,6 @@ func TestGetFileStats(t *testing.T) { FileSize: 500, MimeType: "video/mp4", Extension: "mp4", - StorageDriver: "local", Type: "generic", Status: model.UploadStatusUsed, CreatedAt: time.Now().AddDate(0, 0, -2), // 2 days ago @@ -799,7 +781,6 @@ func TestGetFileStats(t *testing.T) { FileSize: 200, MimeType: "application/pdf", Extension: "pdf", - StorageDriver: "local", Type: "avatar", // different type Status: model.UploadStatusUsed, CreatedAt: time.Now().AddDate(0, 0, -10), // older than 7 days @@ -891,7 +872,6 @@ func TestUserUploadManagement(t *testing.T) { FileSize: 100, MimeType: "text/plain", Extension: "txt", - StorageDriver: "local", Status: model.UploadStatusUsed, CreatedAt: time.Now(), } @@ -903,7 +883,6 @@ func TestUserUploadManagement(t *testing.T) { FileSize: 200, MimeType: "image/png", Extension: "png", - StorageDriver: "local", Status: model.UploadStatusUsed, CreatedAt: time.Now(), } diff --git a/internal/apps/upload/shared/constants.go b/internal/apps/upload/shared/constants.go index 405bf08e..cac5b71e 100644 --- a/internal/apps/upload/shared/constants.go +++ b/internal/apps/upload/shared/constants.go @@ -3,8 +3,6 @@ package shared -import "github.com/Rain-kl/Wavelet/internal/storage" - // Upload size, path, media quality, and cache constants shared across subpackages. const ( MaxUploadSize = 32 * 1024 * 1024 // 32MB @@ -15,7 +13,6 @@ const ( ImageQualityMedium = "medium" ImageQualityHigh = "high" ImageQualityOrigin = "origin" - StorageDriverLocal = string(storage.DriverLocal) DefaultPublicUploadType = "avatar" FileStatsTrendDays = 7 MaxS3KeyLength = 1024 diff --git a/internal/apps/upload/storage/storage_ops.go b/internal/apps/upload/storage/storage_ops.go index 8e4e36e7..4496f1e8 100644 --- a/internal/apps/upload/storage/storage_ops.go +++ b/internal/apps/upload/storage/storage_ops.go @@ -5,7 +5,6 @@ package storage import ( "context" - "fmt" "github.com/Rain-kl/Wavelet/internal/model" "github.com/Rain-kl/Wavelet/internal/storage" @@ -22,46 +21,11 @@ func ReadOnly(ctx context.Context) bool { return state.ReadOnly } -// OpenStoredObject opens a stored upload object from its configured backend. +// OpenStoredObject opens a stored upload object from the active storage backend. func OpenStoredObject(ctx context.Context, upload *model.Upload) (*storage.Object, error) { - driver := storage.Driver(upload.StorageDriver) - if driver == "" { - driver = storage.DriverLocal - } - backend, err := backendForStoredDriver(ctx, driver) + _, backend, err := storage.Active(ctx) if err != nil { return nil, err } return backend.Get(ctx, upload.FilePath) -} - -func backendForStoredDriver(ctx context.Context, driver storage.Driver) (storage.Backend, error) { - backend, err := storage.ForDriver(ctx, driver) - if err == nil { - return backend, nil - } - - target, ok, targetErr := CurrentMigrationTargetConfig(ctx) - if targetErr != nil { - return nil, targetErr - } - if ok && target.Driver == driver { - return storage.NewBackend(ctx, target, driver) - } - return nil, fmt.Errorf("storage configuration for driver %q is unavailable", driver) -} - -// CurrentMigrationTargetConfig returns the pending migration target config when available. -func CurrentMigrationTargetConfig(ctx context.Context) (storage.Config, bool, error) { - state := LoadMigrationAccessState(ctx) - if state.LoadErr != nil { - return storage.Config{}, false, state.LoadErr - } - if state.TargetErr != nil { - return storage.Config{}, false, state.TargetErr - } - if !state.HasTarget { - return storage.Config{}, false, nil - } - return state.Target, true, nil } \ No newline at end of file diff --git a/internal/apps/upload/task/cleanup.go b/internal/apps/upload/task/cleanup.go index 82ac2fc1..4aeb45a5 100644 --- a/internal/apps/upload/task/cleanup.go +++ b/internal/apps/upload/task/cleanup.go @@ -84,11 +84,7 @@ func (h *SystemCleanupHandler) Execute(ctx context.Context, _ []byte) (*task.Tas return err } - driver := storage.Driver(u.StorageDriver) - if driver == "" { - driver = storage.DriverLocal - } - backend, err := storage.ForDriver(ctx, driver) + _, backend, err := storage.Active(ctx) if err != nil { return err } diff --git a/internal/apps/upload/task/storage_migration.go b/internal/apps/upload/task/storage_migration.go index 1a54f8c3..2f8d9f83 100644 --- a/internal/apps/upload/task/storage_migration.go +++ b/internal/apps/upload/task/storage_migration.go @@ -29,8 +29,6 @@ const ( StorageMigrationTask = uploadstorage.StorageMigrationTask // TaskTypeStorageMigration is the task metadata type for storage migration. TaskTypeStorageMigration = "storage_migration" - - colStorageDriver = "storage_driver" ) // StorageMigrationMeta describes the manually dispatchable migration task. @@ -136,7 +134,7 @@ func (h *MigrationHandler) Execute(ctx context.Context, payload []byte) (*task.T return &task.TaskResult{Message: message}, nil } - total, err := countStorageObjects(ctx, active.Driver) + total, err := countStorageObjects(ctx) if err != nil { return nil, fmt.Errorf("count source objects: %w", err) } @@ -159,7 +157,7 @@ func (h *MigrationHandler) Execute(ctx context.Context, payload []byte) (*task.T } task.AppendLog(ctx, "开始存储迁移: %s -> %s,总对象数: %d", active.Driver, target.Driver, total) - migrated, err := migrateObjects(ctx, sourceBackend, targetBackend, active.Driver, target.Driver, total) + migrated, err := migrateObjects(ctx, sourceBackend, targetBackend, total) if err != nil { return nil, err } @@ -172,10 +170,10 @@ func (h *MigrationHandler) Execute(ctx context.Context, payload []byte) (*task.T return &task.TaskResult{Message: message}, nil } -func countStorageObjects(ctx context.Context, driver storage.Driver) (int64, error) { +func countStorageObjects(ctx context.Context) (int64, error) { var count int64 err := db.DB(ctx).Model(&model.Upload{}). - Where("storage_driver = ? AND status != ?", driver, model.UploadStatusDeleted). + Where("status != ?", model.UploadStatusDeleted). Distinct("file_path"). Count(&count).Error return count, err @@ -189,18 +187,24 @@ 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 storage.Backend, targetBackend storage.Backend, - sourceDriver storage.Driver, - targetDriver storage.Driver, total int64, ) (int64, error) { const batchSize = 50 const migrationConcurrency = 10 const sha256HexLength = 64 var migrated int64 + var lastFilePath string for { if err := ctx.Err(); err != nil { return atomic.LoadInt64(&migrated), fmt.Errorf("storage migration canceled: %w", err) @@ -208,16 +212,14 @@ func migrateObjects( task.AppendLog(ctx, "正在查询待迁移对象批次,当前已完成迁移: %d/%d", atomic.LoadInt64(&migrated), total) - var objects []struct { - FilePath string `gorm:"column:file_path"` - FileSize int64 `gorm:"column:file_size"` - MimeType string `gorm:"column:mime_type"` - Hash string `gorm:"column:hash"` - } - if err := db.DB(ctx).Model(&model.Upload{}). + 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("storage_driver = ? AND status != ?", sourceDriver, model.UploadStatusDeleted). - Group("file_path"). + 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 { @@ -228,6 +230,7 @@ func migrateObjects( break } + lastFilePath = objects[len(objects)-1].FilePath task.AppendLog(ctx, "获取当前批次迁移对象,批次大小: %d,实际获取对象数: %d", batchSize, len(objects)) var g errgroup.Group @@ -236,7 +239,7 @@ func migrateObjects( for _, object := range objects { obj := object g.Go(func() error { - if err := migrateSingleObject(ctx, sourceBackend, targetBackend, sourceDriver, targetDriver, obj, sha256HexLength); err != nil { + if err := migrateSingleObject(ctx, sourceBackend, targetBackend, obj, sha256HexLength); err != nil { return err } atomic.AddInt64(&migrated, 1) @@ -257,25 +260,11 @@ func migrateSingleObject( ctx context.Context, sourceBackend storage.Backend, targetBackend storage.Backend, - sourceDriver storage.Driver, - targetDriver storage.Driver, - obj struct { - FilePath string `gorm:"column:file_path"` - FileSize int64 `gorm:"column:file_size"` - MimeType string `gorm:"column:mime_type"` - Hash string `gorm:"column:hash"` - }, + obj migrationObject, sha256HexLength int, ) error { if shouldSkipMigration(ctx, targetBackend, obj) { - task.AppendLog(ctx, "[跳过迁移] 目标存储已存在相同文件且校验一致: %s", obj.FilePath) - if err := db.DB(ctx).Model(&model.Upload{}). - Where("storage_driver = ? AND file_path = ?", sourceDriver, obj.FilePath). - Updates(map[string]any{ - colStorageDriver: targetDriver, - }).Error; err != nil { - return fmt.Errorf("update migrated object %q: %w", obj.FilePath, err) - } + task.AppendLog(ctx, "[跳过迁移] 目标存储已存在相同文件: %s", obj.FilePath) return nil } @@ -283,7 +272,7 @@ func migrateSingleObject( source, err := sourceBackend.Get(ctx, obj.FilePath) if err != nil { if isNotFoundError(err) { - return markMissingMigrationObjectDeleted(ctx, sourceDriver, targetDriver, obj.FilePath, err) + return markMissingMigrationObjectDeleted(ctx, obj.FilePath, err) } return fmt.Errorf("open source object %q: %w", obj.FilePath, err) } @@ -319,28 +308,22 @@ func migrateSingleObject( task.AppendLog(ctx, "[校验通过] 文件一致性校验成功: %s", targetResult.Key) } - task.AppendLog(ctx, "[更新数据库] 正在更新文件 %s 的存储路径与驱动信息为: %s (%s)", obj.FilePath, targetResult.Key, targetDriver) - if err := db.DB(ctx).Model(&model.Upload{}). - Where("storage_driver = ? AND file_path = ?", sourceDriver, obj.FilePath). - Updates(map[string]any{ - colStorageDriver: targetDriver, - "file_path": targetResult.Key, - }).Error; err != nil { - return fmt.Errorf("update migrated object %q: %w", obj.FilePath, err) + 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 { + return fmt.Errorf("update migrated object %q: %w", obj.FilePath, err) + } } - task.AppendLog(ctx, "[迁移成功] 文件已完成迁移: %s -> %s", obj.FilePath, targetResult.Key) + task.AppendLog(ctx, "[迁移成功] 文件已完成迁移: %s", targetResult.Key) return nil } func shouldSkipMigration( ctx context.Context, targetBackend storage.Backend, - obj struct { - FilePath string `gorm:"column:file_path"` - FileSize int64 `gorm:"column:file_size"` - MimeType string `gorm:"column:mime_type"` - Hash string `gorm:"column:hash"` - }, + obj migrationObject, ) bool { targetObj, err := targetBackend.Get(ctx, obj.FilePath) if err != nil || targetObj == nil || targetObj.Body == nil { @@ -355,8 +338,6 @@ func shouldSkipMigration( func markMissingMigrationObjectDeleted( ctx context.Context, - sourceDriver storage.Driver, - targetDriver storage.Driver, filePath string, sourceErr error, ) error { @@ -364,16 +345,13 @@ func markMissingMigrationObjectDeleted( var affectedUploads []model.Upload if err := db.DB(ctx). - Where("storage_driver = ? AND file_path = ? AND status != ?", sourceDriver, filePath, model.UploadStatusDeleted). + 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("storage_driver = ? AND file_path = ?", sourceDriver, filePath). - Updates(map[string]any{ - "status": model.UploadStatusDeleted, - colStorageDriver: targetDriver, - }).Error; err != nil { + Where("file_path = ?", filePath). + Update("status", model.UploadStatusDeleted).Error; err != nil { return fmt.Errorf("update missing object %q: %w", filePath, err) } for i := range affectedUploads { diff --git a/internal/apps/upload/task/storage_migration_task_test.go b/internal/apps/upload/task/storage_migration_task_test.go index f7216971..39262d4d 100644 --- a/internal/apps/upload/task/storage_migration_task_test.go +++ b/internal/apps/upload/task/storage_migration_task_test.go @@ -68,7 +68,6 @@ func TestMigrationHandlerExecute(t *testing.T) { MimeType: "text/plain", Extension: "txt", Hash: "hash", - StorageDriver: string(storage.DriverLocal), Type: "attachment", Status: model.UploadStatusUsed, } @@ -106,9 +105,6 @@ func TestMigrationHandlerExecute(t *testing.T) { if err := dbConn.First(&migrated, upload.ID).Error; err != nil { t.Fatalf("First(upload) returned error: %v", err) } - if migrated.StorageDriver != string(storage.DriverS3) { - t.Errorf("StorageDriver = %q, want %q", migrated.StorageDriver, storage.DriverS3) - } current, err := storage.LoadConfig(ctx) if err != nil { t.Fatalf("LoadConfig() returned error: %v", err) @@ -169,7 +165,6 @@ func TestMigrationHandlerExecuteWithHashValidation(t *testing.T) { MimeType: "text/plain", Extension: "txt", Hash: "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", // Invalid hash - StorageDriver: string(storage.DriverLocal), Type: "attachment", Status: model.UploadStatusUsed, } @@ -224,8 +219,8 @@ func TestMigrationHandlerExecuteWithHashValidation(t *testing.T) { if err := dbConn.First(&migrated, uploadIncorrect.ID).Error; err != nil { t.Fatalf("First(upload) returned error: %v", err) } - if migrated.StorageDriver != string(storage.DriverS3) { - t.Errorf("StorageDriver = %q, want %q", migrated.StorageDriver, storage.DriverS3) + if migrated.FilePath != "uploads/test-hash.txt" { + t.Errorf("FilePath = %q, want %q", migrated.FilePath, "uploads/test-hash.txt") } } diff --git a/internal/apps/upload/task/tasks_test.go b/internal/apps/upload/task/tasks_test.go index 79b3293b..bab0d184 100644 --- a/internal/apps/upload/task/tasks_test.go +++ b/internal/apps/upload/task/tasks_test.go @@ -42,6 +42,9 @@ func TestSystemCleanupHandler_Execute(t *testing.T) { func(ctx context.Context, key string) error { return nil }, ) defer storageMock() + storage.IsEnabledFunc = func() bool { return true } + defer func() { storage.IsEnabledFunc = func() bool { return false } }() + storage.ResetCache() ctx := context.Background() err := db.DB(ctx).AutoMigrate(&model.PushHistory{}) @@ -56,27 +59,27 @@ func TestSystemCleanupHandler_Execute(t *testing.T) { { UserID: 1001, FileName: "old_file_1.jpg", FilePath: "uploads/old_1.jpg", FileSize: 1024, MimeType: "image/jpeg", Extension: "jpg", Hash: "hash1", - StorageDriver: "s3", Type: "attachment", Status: model.UploadStatusPending, + Type: "attachment", Status: model.UploadStatusPending, CreatedAt: twoHoursAgo, }, { UserID: 1001, FileName: "old_file_2.png", FilePath: "uploads/old_2.png", FileSize: 2048, MimeType: "image/png", Extension: "png", Hash: "hash2", - StorageDriver: "s3", Type: "attachment", Status: model.UploadStatusPending, + Type: "attachment", Status: model.UploadStatusPending, CreatedAt: twoHoursAgo, }, // 状态为 used 的记录 —— 不应被清理 { UserID: 1001, FileName: "used_file.jpg", FilePath: "uploads/used.jpg", FileSize: 512, MimeType: "image/jpeg", Extension: "jpg", Hash: "hash3", - StorageDriver: "s3", Type: "attachment", Status: model.UploadStatusUsed, + Type: "attachment", Status: model.UploadStatusUsed, CreatedAt: twoHoursAgo, }, // 不到1小时的 pending 记录 —— 不应被清理 { UserID: 1001, FileName: "recent_file.jpg", FilePath: "uploads/recent.jpg", FileSize: 256, MimeType: "image/jpeg", Extension: "jpg", Hash: "hash4", - StorageDriver: "s3", Type: "attachment", Status: model.UploadStatusPending, + Type: "attachment", Status: model.UploadStatusPending, CreatedAt: now.Add(-10 * time.Minute), }, } @@ -283,7 +286,6 @@ func TestWarmImageCacheHandlerExecute(t *testing.T) { FilePath: firstPath, MimeType: "image/png", Extension: "png", - StorageDriver: shared.StorageDriverLocal, Status: model.UploadStatusUsed, }, { @@ -293,7 +295,6 @@ func TestWarmImageCacheHandlerExecute(t *testing.T) { FilePath: secondPath, MimeType: "application/octet-stream", Extension: "jpg", - StorageDriver: shared.StorageDriverLocal, Status: model.UploadStatusPending, }, { @@ -303,7 +304,6 @@ func TestWarmImageCacheHandlerExecute(t *testing.T) { FilePath: filepath.Join(testDir, "notes.txt"), MimeType: "text/plain", Extension: "txt", - StorageDriver: shared.StorageDriverLocal, Status: model.UploadStatusUsed, }, { @@ -313,7 +313,6 @@ func TestWarmImageCacheHandlerExecute(t *testing.T) { FilePath: firstPath, MimeType: "image/png", Extension: "png", - StorageDriver: shared.StorageDriverLocal, Status: model.UploadStatusDeleted, }, } diff --git a/internal/db/migrator/goose/postgres/202606180001_drop_upload_storage_driver.sql b/internal/db/migrator/goose/postgres/202606180001_drop_upload_storage_driver.sql new file mode 100644 index 00000000..3ba6aa08 --- /dev/null +++ b/internal/db/migrator/goose/postgres/202606180001_drop_upload_storage_driver.sql @@ -0,0 +1,7 @@ +-- +goose Up +DROP INDEX IF EXISTS idx_w_uploads_storage_driver_status; +ALTER TABLE w_uploads DROP COLUMN IF EXISTS storage_driver; + +-- +goose Down +ALTER TABLE w_uploads ADD COLUMN IF NOT EXISTS storage_driver VARCHAR(50) NOT NULL DEFAULT 'local'; +CREATE INDEX IF NOT EXISTS idx_w_uploads_storage_driver_status ON w_uploads (storage_driver, status); \ No newline at end of file diff --git a/internal/db/migrator/goose/sqlite/202606180001_drop_upload_storage_driver.sql b/internal/db/migrator/goose/sqlite/202606180001_drop_upload_storage_driver.sql new file mode 100644 index 00000000..aa66dcdc --- /dev/null +++ b/internal/db/migrator/goose/sqlite/202606180001_drop_upload_storage_driver.sql @@ -0,0 +1,8 @@ +-- +goose Up +DROP INDEX IF EXISTS idx_w_uploads_storage_driver_status; +-- SQLite lacks DROP COLUMN IF EXISTS; goose runs this only after initial schema created the column. +ALTER TABLE w_uploads DROP COLUMN storage_driver; + +-- +goose Down +ALTER TABLE w_uploads ADD COLUMN storage_driver VARCHAR(50) NOT NULL DEFAULT 'local'; +CREATE INDEX IF NOT EXISTS idx_w_uploads_storage_driver_status ON w_uploads (storage_driver, status); \ No newline at end of file diff --git a/internal/model/uploads.go b/internal/model/uploads.go index 2b218e46..d3e0d4d1 100644 --- a/internal/model/uploads.go +++ b/internal/model/uploads.go @@ -40,7 +40,6 @@ type Upload struct { MimeType string `json:"mime_type" gorm:"size:100;not null"` // 媒体类型 (MIME, 如 image/png) Extension string `json:"extension" gorm:"size:50;not null"` // 文件后缀名 (不含点,如 png, pdf) Hash string `json:"hash" gorm:"size:64;index"` // 文件哈希 (SHA-256/MD5,可用于排重) - StorageDriver string `json:"storage_driver" gorm:"size:50;not null"` // 存储引擎驱动 (如 local, s3, oss) Type string `json:"type" gorm:"column:type;size:50;not null;index"` // 业务标识类型 (如 avatar, doc, attachment) Status UploadStatus `json:"status" gorm:"type:varchar(20);not null"` // 状态 AccessMode int `json:"access_mode" gorm:"column:access_mode;not null;default:0"` diff --git a/internal/storage/storage.go b/internal/storage/storage.go index 9f1d03c3..26a2560d 100644 --- a/internal/storage/storage.go +++ b/internal/storage/storage.go @@ -155,26 +155,6 @@ func Active(ctx context.Context) (Driver, Backend, error) { return activeDriver, activeBackend, nil } -// ForDriver returns the active or pending backend for an upload record, reusing the active singleton if matched. -func ForDriver(ctx context.Context, driver Driver) (Backend, error) { - if driver == DriverS3 && mockBackend != nil { - return mockBackend, nil - } - activeDrv, activeBnd, err := Active(ctx) - if err == nil && activeDrv == driver { - return activeBnd, nil - } - cfg, err := LoadConfig(ctx) - if err != nil { - return nil, err - } - backend, err := NewBackend(ctx, cfg, driver) - if err != nil { - return nil, fmt.Errorf("storage configuration for driver %q is unavailable: %w", driver, err) - } - return backend, nil -} - type functionBackend struct { put func(context.Context, string, io.Reader, int64, string) error get func(context.Context, string) (*Object, error) diff --git a/internal/storage/storage_test.go b/internal/storage/storage_test.go index 534fd67c..20a28de9 100644 --- a/internal/storage/storage_test.go +++ b/internal/storage/storage_test.go @@ -70,16 +70,7 @@ func TestStorageCache(t *testing.T) { t.Errorf("Active returned driver %v, backend %v; expected %v, %v", drv, bnd, DriverLocal, mockBnd) } - // 5. Test ForDriver returns cached backend - bnd2, err := ForDriver(ctx, DriverLocal) - if err != nil { - t.Fatalf("ForDriver failed: %v", err) - } - if bnd2 != mockBnd { - t.Errorf("ForDriver returned backend %v, expected %v", bnd2, mockBnd) - } - - // 6. Test ResetCache again + // 5. Test ResetCache again ResetCache() if activeConfigJSON != "" || activeDriver != "" || activeBackend != nil || !lastChecked.IsZero() { t.Fatal("ResetCache did not clear cache variables after setting them")