diff --git a/.agent/skills/file-upload/SKILL.md b/.agent/skills/file-upload/SKILL.md new file mode 100644 index 00000000..e0ec144e --- /dev/null +++ b/.agent/skills/file-upload/SKILL.md @@ -0,0 +1,265 @@ +--- +name: "file-upload" +description: "Wavelet 项目专用:当业务需要上传文件、读取已上传文件、在 Worker/任务中程序化摄取字节流、选择存储引擎能力、或排查 w_uploads / 文件统计异常时必须使用。本技能指导 storage 与 upload 分层、upload.Ingest 策略选型、前后端接入与禁止旁路写表。" +--- + +# 存储引擎与文件上传开发规范 + +本技能是 Wavelet **文件上传与对象存储**的唯一开发指导。开始开发前先阅读仓库根目录 [AGENTS.md](file:///Users/ryan/DEV/Go/Wavelet/AGENTS.md),遵守项目级核心规则。 + +--- + +## 架构分层(必须理解) + +Wavelet 将「对象存储」与「上传业务」分为两层,**禁止混用职责**: + +| 层级 | 包路径 | 职责 | 业务是否直接调用 | +| :--- | :--- | :--- | :--- | +| **对象存储引擎** | `internal/storage` | `Backend` 接口:`Put` / `Get` / `Delete` / `Test`;按配置切换 Local / S3 / R2 / OSS / WebDAV | **禁止**(仅 upload 域内部使用) | +| **上传域服务** | `internal/apps/upload` | `w_uploads` 记录、权限、秒传、统计、文件服务、`upload.Ingest` | **必须** | +| **上传 HTTP 入口** | `internal/apps/upload/handler` | `POST /api/v1/upload` 等 multipart 接口 | 前端 / 用户侧上传 | +| **文件访问** | `internal/apps/upload/filesrv` | `GET /f/:id` 流式响应、访问控制、图片 WebP 压缩 | 展示 / 下载 | + +```text +业务模块 ──► upload.Ingest / upload.Remove(唯一写入门禁) + ├── storage.Backend.Put/Get/Delete + ├── repository.CreateUpload(仅 upload 内部) + └── RecordUploadStatsAdd/Remove(ingest 内置,禁止业务直调) +``` + +--- + +## 核心防线(Guardrails) + +以下写法**一律禁止**: + +```go +// ❌ 业务包直接写 blob +storage.Active(ctx); backend.Put(...) + +// ❌ 旁路写 w_uploads +db.DB(ctx).Create(&model.Upload{}) +repository.CreateUpload(ctx, upload) // 仅 internal/apps/upload 允许 + +// ❌ 手动维护统计 +upload.ApplyUploadStatsAdd(ctx, upload) // 已 Deprecated + +// ❌ 业务表存物理路径 +invoice.FilePath = "uploads/2026/01/02/123.pdf" +``` + +**正确做法**:业务表只存 `upload_id`(`uint64` / JSON string),通过 `/f/{id}` 或 `upload.OpenStoredObject` 访问。 + +--- + +## Ingest 策略选型(Policy Decision) + +根据场景选择 `upload.Ingest` 的 `Policy`: + +| 场景 | Policy | 哈希命中时 | 未命中时 | 典型调用方 | +| :--- | :--- | :--- | :--- | :--- | +| 用户 HTTP 上传(含秒传) | `PolicyDedupNewRecord` | 复用 path,**新建记录 + 统计** | 写 blob + 新建记录 + 统计 | `handler.UploadFile`(已内置) | +| Worker 生成全新文件 | `PolicyCreate` | 不查重,始终写 blob + 记录 | 同左 | 报表导出、定时生成 | +| 镜像 / 去重摄取(Pixez) | `PolicyResolveExisting` | **直接返回已有记录**,不建新记录、不加统计 | 写 blob + 新建记录 + 统计 | 异步镜像任务 | +| 业务只需引用已有文件 | 不调 Ingest | — | — | 业务 API 校验 `upload_id` 即可 | + +### Result 字段含义 + +| 字段 | 含义 | +| :--- | :--- | +| `Created` | 是否新建了 `w_uploads` 记录 | +| `Stored` | 是否写入了新 blob | +| `Resolved` | 是否通过哈希解析到已有记录(仅 `PolicyResolveExisting`) | + +--- + +## 后端:程序化上传(Worker / 业务逻辑) + +### 标准模板 + +在 `logics.go`(接受 `context.Context`,不依赖 `*gin.Context`)中调用: + +```go +import ( + "bytes" + + "github.com/Rain-kl/Wavelet/internal/apps/upload" + "github.com/Rain-kl/Wavelet/internal/model" +) + +func ingestMirrorFile(ctx context.Context, userID uint64, data []byte, hash, filename, mime, ext string) (model.Upload, error) { + accessMode := 1 + result, err := upload.Ingest(ctx, upload.IngestRequest{ + UserID: userID, + Reader: bytes.NewReader(data), + Size: int64(len(data)), + FileName: filename, + MimeType: mime, + Extension: ext, + Hash: hash, // 必填:SHA-256 hex + Type: "your_biz_type", + AccessMode: &accessMode, + Metadata: model.UploadMetadata{ + Extra: map[string]any{"source": "worker"}, + }, + Policy: upload.PolicyResolveExisting, + }) + if err != nil { + return model.Upload{}, err + } + return result.Upload, nil +} +``` + +### Request 关键字段 + +| 字段 | 说明 | +| :--- | :--- | +| `Hash` | **必填**,推荐 SHA-256 hex;用于秒传 / 镜像去重 | +| `Type` | 业务分类(如 `avatar`、`invoice`、`pixez_mirror`),用于筛选与统计 | +| `AccessMode` | `nil` 时按 type 默认:`avatar` → 公开(1),其余 → 私有(0) | +| `SkipExtensionCheck` | Worker 场景若已自行校验扩展名,可设为 `true` | +| `ObjectKeyFn` | 可选自定义存储路径;默认 `uploads/YYYY/MM/DD/{id}.{ext}` | + +### 错误处理 + +| 错误 | 含义 | Handler 映射建议 | +| :--- | :--- | :--- | +| `upload.ErrIngestStorageReadOnly` | 存储迁移维护中 | `response.AbortConflict` | +| `ingest.ErrForbidden` | 无权删除他人文件 | HTTP 403 | +| `shared.ErrUnsupportedFormat` | 扩展名不在白名单 | `response.AbortBadRequest` | + +### 删除 + +```go +// 管理员 / 系统删除 +_, err := upload.Remove(ctx, uploadID) + +// 用户删除自己的文件 +_, err := upload.RemoveOwned(ctx, userID, uploadID) +``` + +### 读取已存储对象(不上传) + +```go +uploadRec, err := repository.GetActiveUploadByID(ctx, uploadID) +obj, err := uploadstorage.OpenStoredObject(ctx, &uploadRec) +defer obj.Body.Close() +``` + +或通过门面(若已从 `exports` 暴露 `OpenStoredObject`)读取。HTTP 对外访问统一走 `GET /f/:id`。 + +--- + +## 后端:业务 API 引用已上传文件 + +推荐 **两步流程**(先上传、后提交业务): + +1. 前端 `POST /api/v1/upload` → 获得 `upload.id` +2. 业务 API 接收 `upload_id`,用 `repository.GetActiveUploadByID` 校验存在且 `status` 为 active +3. (可选)校验 `upload.Type` 是否为预期业务类型 +4. 将 `upload_id` 写入业务表字段(如 `cover_file_id`) + +**禁止**在业务 Handler 中重复实现 multipart 解析,除非有极强的特殊协议需求。 + +--- + +## 前端:用户侧上传 + +使用 `frontend/lib/services/upload/`: + +```typescript +import { services } from '@/lib/services' +import { getFileUrl } from '@/lib/services/upload' + +// 上传 +const upload = await services.upload.uploadFile(file, 'invoice', { orderId: '123' }) + +// 展示 +const url = getFileUrl(upload.id) // → /f/{id} + +// Base64 图片(头像等) +const res = await services.upload.uploadBase64Image(croppedBase64, 'avatar', 'avatar.png') +``` + +### 前端规范 + +- 新增上传相关 API 时,扩展 `UploadService` / `AdminUploadService`,在 `frontend/lib/services/index.ts` 注册 +- 图片预览使用 `getFileUrl(id, quality?)` 或 `FileImagePreview` 组件 +- 业务表单项只提交 `upload_id`,不要提交 blob URL 或 `file_path` + +--- + +## 统计与排查 + +`w_upload_stats` 由 `upload.Ingest` / `upload.Remove` **自动维护**,业务不得手动增量。 + +若发现 trend / total 与 `w_uploads` 不一致(常见于历史旁路写表): + +```go +upload.RebuildUploadStats(ctx) // 从 w_uploads 全量重建统计 +``` + +排查清单: + +1. 业务是否绕过 `upload.Ingest` 直接 `db.Create(&model.Upload{})`? +2. 是否手动调用已 Deprecated 的 `ApplyUploadStatsAdd`? +3. 删除是否走 `upload.Remove`(须在软删**前**扣减统计)? + +--- + +## 测试要求 + +### 后端 ingest 测试 + +- 使用 `testhelper.SetupTestEnvironment(t)` 初始化 DB +- 存储 mock:`storage.MockStorage(...)` + `storage.IsEnabledFunc = func() bool { return true }` +- **禁止**在源码目录硬编码 `uploads/test` 路径;本地文件测试用 `t.TempDir()` 或 mock backend +- 覆盖:三种 Policy、Remove 后统计归零、ReadOnly 拒绝写入 + +参考:[internal/apps/upload/ingest/ingest_test.go](file:///Users/ryan/DEV/Go/Wavelet/internal/apps/upload/ingest/ingest_test.go) + +### Handler 回归 + +修改 upload handler 后运行: + +```bash +go test ./internal/apps/upload/... +make code-check +``` + +若变更 HTTP 接口,运行 `make swagger`。 + +--- + +## 存量代码迁移(旁路写表 → Ingest) + +将以下模式: + +```go +storage.Active(ctx) +backend.Put(ctx, key, reader, size, mime) +db.DB(ctx).Create(&upload) +``` + +替换为: + +```go +upload.Ingest(ctx, upload.IngestRequest{ Policy: upload.PolicyResolveExisting, ... }) +``` + +迁移完成后执行一次 `upload.RebuildUploadStats(ctx)` 修复历史统计偏差。 + +--- + +## 质量门禁 Checklist + +完成文件上传相关开发后,确认: + +- [ ] 业务模块无 `repository.CreateUpload` / `SoftDeleteUpload` 调用 +- [ ] 业务模块无 `storage.Active` + `Put` 直接写文件 +- [ ] 业务表存 `upload_id`,不存 `file_path` +- [ ] Worker 摄取使用正确的 `Policy` +- [ ] 新增测试覆盖 ingest 路径 +- [ ] `make code-check` 通过 +- [ ] HTTP 变更已 `make swagger` \ No newline at end of file diff --git a/AGENTS.md b/AGENTS.md index 14ef394d..ba5f662a 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -21,14 +21,17 @@ ## 务必阅读匹配的 Skill -- `new-api`:在添加或修改自定义业务 API、Handler、服务层逻辑或注册自定义端点时使用。 -- `new-async-task`:在添加或修改 Asynq 任务、定时任务时使用。 -- `new-setting`:在添加或修改基于数据库的系统/业务/公开设置、`/admin/system` 参数或 `/admin/settings` 图形化设置时使用。 -- `database-migration`:在数据库升级流程时使用。 -- Go skills:使用针对性的 `go-*` skills 来获取 Go 实现细节,如测试、错误处理、包、Context、并发、日志、文档和审查。 -- `shadcn`:在添加、修改或组合 shadcn/ui 组件时使用。 -- `code-review-skill`:在提交 PR 之前使用,检查代码质量、样式、潜在错误和最佳实践。 -- `push-notification`:在添加或修改系统通知推送事件、修改消息推送底层设计、调用统一触发器投递消息、或开发带消息推送功能的业务功能时使用。 +| Skill | 何时使用 | +| :--- | :--- | +| `new-api` | 添加或修改自定义业务 API、Handler、服务层逻辑、自定义路由注册 | +| `new-async-task` | 添加或修改 Asynq 任务、定时任务、TaskHandler、任务元数据 | +| `new-setting` | 添加或修改系统/业务/公开设置、`/admin/system` 参数或 `/admin/settings` 图形化设置 | +| `database-migration` | 数据库表结构变更、goose SQL 迁移、seed 数据 | +| `file-upload` | 业务上传文件、Worker 程序化摄取、`upload.Ingest` 策略选型、文件访问与 `w_uploads` / 统计排查 | +| `push-notification` | 系统通知推送事件、统一触发器投递、带消息推送的业务功能 | +| `release-guide` | 根据自上一正式版本 Tag 以来的提交整理 Version Bump 提交信息以触发双语 Release | +| `shadcn` | 添加、修改或组合 shadcn/ui 组件 | + ## 严格遵循事项 (Guardrails) @@ -39,6 +42,7 @@ - 当 API Handler 发生变化时,更新 Swagger 文档(运行 `make swagger`)。 - 在完成代码开发后必须运行 `make code-check`, 并修复报错。 - 需要缓存或文件管理能力时,必须复用现有平台实现,禁止在业务包中自行创建缓存目录、直接管理缓存文件或重复封装存储后端。 +- 文件摄取必须通过 `upload.Ingest`(`upload.PolicyCreate` / `PolicyDedupNewRecord` / `PolicyResolveExisting`);删除必须通过 `upload.Remove` 或 `upload.RemoveOwned`。禁止业务模块直接调用 `repository.CreateUpload` / `repository.SoftDeleteUpload`,禁止 `db.Create(&model.Upload{})` 旁路写 `w_uploads`。 - 禁止在 `init()` 中注册跨模块集成(任务 Handler、推送内置事件、域事件监听器、任务完成钩子)。统一通过 `internal/bootstrap` 在 `internal/cmd` 入口显式装配。 - `internal/router/router.go` 的 `Serve()` 仅负责 HTTP 路由与中间件,禁止在其中执行 `SyncEvents`、`InitLogWriter` 等进程级运行时初始化。 - 核心业务模块(如 `oauth`、`user`)禁止直接 `import` `internal/apps/admin/push` 或 `custom_events` 触发通知;应通过 `internal/listener` 发射域事件,由 push 模块在 bootstrap 阶段订阅。 @@ -77,7 +81,7 @@ - `internal/config/`:Viper 加载和配置结构体。运行时代码应使用 `config.Config.
.`。 - `internal/router/`:唯一的 HTTP 路由注册点。 - `internal/apps/`:按功能(Feature-based)组织的 HTTP Handler、中间件、内部服务与模块逻辑。移除全局 service 层,模块内部业务逻辑(如验证码业务逻辑管理器 `internal/apps/cap/manager.go`)均收敛于各自模块中;管理端模块位于 `internal/apps/admin/`。 -- `internal/apps/upload/`:上传记录、文件访问控制、本地/S3 文件响应、下载及图片 WebP 压缩。业务应复用这些入口,不直接操作底层文件。 +- `internal/apps/upload/`:上传记录、文件访问控制、本地/S3 文件响应、下载及图片 WebP 压缩。业务应复用 `upload.Ingest` / `upload.Remove` 与 `GET /f/:id` 文件服务,不直接操作底层 storage 或旁路写 `w_uploads`。 - `internal/model/`:GORM 实体和模型级业务方法。 - `internal/db/`:PostgreSQL、Redis、ClickHouse、GORM 日志、ID 生成和 goose SQL 迁移的布线。 - `internal/diskcache/`:平台级磁盘字节缓存,通过 `diskcache.GetGlobalCache()` 提供 TTL、最大空间限制、LRU 淘汰、清空、状态统计和配置热更新。写入时使用 `DefaultExpiration`(全局默认 TTL)、正数 `time.Duration`(业务 TTL)或 `NoExpiration`(无 TTL,仍受空间限制和 LRU 淘汰)。 diff --git a/internal/apps/upload/exports.go b/internal/apps/upload/exports.go index 0ea8a489..56eb4a60 100644 --- a/internal/apps/upload/exports.go +++ b/internal/apps/upload/exports.go @@ -7,6 +7,7 @@ import ( "github.com/Rain-kl/Wavelet/internal/apps/upload/cache" "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" 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" @@ -28,6 +29,36 @@ var ( ServeFileByID = filesrv.ServeFileByID ) +// Programmatic ingest API +var ( + Ingest = ingest.Ingest + Remove = ingest.Remove + RemoveOwned = ingest.RemoveOwned + FindByHash = ingest.FindByHash +) + +// Ingest policy constants +const ( + PolicyCreate = ingest.PolicyCreate + PolicyDedupNewRecord = ingest.PolicyDedupNewRecord + PolicyResolveExisting = ingest.PolicyResolveExisting +) + +type ( + // IngestRequest is the programmatic upload ingest payload. + IngestRequest = ingest.Request + // IngestResult reports ingest side effects. + IngestResult = ingest.Result + // IngestPolicy controls hash-collision behavior during ingest. + IngestPolicy = ingest.Policy +) + +// Ingest errors +var ( + ErrIngestForbidden = ingest.ErrForbidden + ErrIngestStorageReadOnly = ingest.ErrStorageReadOnly +) + // Cache management var ( ResetAccessCaches = cache.ResetAccessCaches @@ -36,7 +67,9 @@ var ( // Stats var ( - ApplyUploadStatsAdd = uploadstats.ApplyUploadStatsAdd + // Deprecated: use upload.Ingest or upload.Remove; stats are applied internally. + ApplyUploadStatsAdd = uploadstats.ApplyUploadStatsAdd + // Deprecated: use upload.Ingest or upload.Remove; stats are applied internally. ApplyUploadStatsRemove = uploadstats.ApplyUploadStatsRemove RebuildUploadStats = uploadstats.RebuildUploadStats ) diff --git a/internal/apps/upload/handler/file_management.go b/internal/apps/upload/handler/file_management.go index b733db5c..aa446184 100644 --- a/internal/apps/upload/handler/file_management.go +++ b/internal/apps/upload/handler/file_management.go @@ -8,6 +8,7 @@ import ( "strconv" "github.com/Rain-kl/Wavelet/internal/apps/oauth" + "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" "github.com/Rain-kl/Wavelet/internal/common/response" @@ -235,7 +236,7 @@ func DeleteMyFile(c *gin.Context) { c.AbortWithStatus(http.StatusNotFound) return } - if err == errUploadForbidden { + if err == ingest.ErrForbidden { c.AbortWithStatus(http.StatusForbidden) return } @@ -289,7 +290,7 @@ func UpdateMyFile(c *gin.Context) { c.AbortWithStatus(http.StatusNotFound) return } - if err == errUploadForbidden { + if err == ingest.ErrForbidden { c.AbortWithStatus(http.StatusForbidden) return } diff --git a/internal/apps/upload/handler/logics.go b/internal/apps/upload/handler/logics.go index d0812d79..3cb549c3 100644 --- a/internal/apps/upload/handler/logics.go +++ b/internal/apps/upload/handler/logics.go @@ -4,20 +4,13 @@ package handler import ( - "bytes" "context" "errors" "sort" - "strings" - "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/idgen" + "github.com/Rain-kl/Wavelet/internal/apps/upload/ingest" "github.com/Rain-kl/Wavelet/internal/model" "github.com/Rain-kl/Wavelet/internal/repository" - "github.com/Rain-kl/Wavelet/internal/storage" - "github.com/Rain-kl/Wavelet/pkg/logger" "gorm.io/gorm" ) @@ -31,30 +24,11 @@ func listMyUploadFiles(ctx context.Context, userID uint64, filter repository.Upl } func softDeleteUpload(ctx context.Context, uploadID uint64) (model.Upload, error) { - upload, err := repository.GetActiveUploadByID(ctx, uploadID) - if err != nil { - return model.Upload{}, err - } - if err := repository.SoftDeleteUpload(ctx, &upload); err != nil { - return model.Upload{}, err - } - uploadstats.RecordUploadStatsRemove(ctx, &upload) - return upload, nil + return ingest.Remove(ctx, uploadID) } func softDeleteOwnedUpload(ctx context.Context, userID, uploadID uint64) (model.Upload, error) { - upload, err := repository.GetActiveUploadByID(ctx, uploadID) - if err != nil { - return model.Upload{}, err - } - if upload.UserID != userID { - return model.Upload{}, errUploadForbidden - } - if err := repository.SoftDeleteUpload(ctx, &upload); err != nil { - return model.Upload{}, err - } - uploadstats.RecordUploadStatsRemove(ctx, &upload) - return upload, nil + return ingest.RemoveOwned(ctx, userID, uploadID) } func listDistinctUploadTypes(ctx context.Context) ([]string, error) { @@ -77,7 +51,7 @@ func updateOwnedUpload(ctx context.Context, userID, uploadID uint64, input updat return model.Upload{}, err } if upload.UserID != userID { - return model.Upload{}, errUploadForbidden + return model.Upload{}, ingest.ErrForbidden } updates := make(map[string]any) @@ -103,96 +77,10 @@ func listUploadsForBatchDownload(ctx context.Context, ids []uint64) ([]model.Upl return repository.ListUploadsByIDs(ctx, ids) } -type instantUploadInput struct { - UserID uint64 - FileHash string - Size int64 - MimeType string - Extension string - OrigName string - UploadType string - AccessMode int -} - -func createInstantUpload(ctx context.Context, existing model.Upload, input instantUploadInput) (model.Upload, error) { - newUpload := model.Upload{ - ID: idgen.NextUint64ID(), - UserID: input.UserID, - FileName: input.OrigName, - FilePath: existing.FilePath, - FileSize: input.Size, - MimeType: input.MimeType, - Extension: input.Extension, - Hash: input.FileHash, - Type: input.UploadType, - Status: model.UploadStatusUsed, - AccessMode: input.AccessMode, - Metadata: existing.Metadata, - } - if err := repository.CreateUpload(ctx, &newUpload); err != nil { - return model.Upload{}, err - } - uploadstats.RecordUploadStatsAdd(ctx, &newUpload) - logger.InfoF(ctx, "文件触发秒传成功! ID: %d, Path: %s", newUpload.ID, existing.FilePath) - return newUpload, nil -} - -func findReusableUpload(ctx context.Context, hash string, size int64) (model.Upload, error) { - return repository.FindReusableUploadByHash(ctx, hash, size) -} - -func saveNewUploadRecord(ctx context.Context, upload *model.Upload, filePath string) error { - if err := repository.CreateUpload(ctx, upload); err != nil { - _, backend, backendErr := storage.Active(ctx) - if backendErr == nil { - if deleteErr := backend.Delete(ctx, filePath); deleteErr != nil { - logger.WarnF(ctx, "清理未写入数据库的上传对象失败: %v", deleteErr) - } - } - return err - } - uploadstats.RecordUploadStatsAdd(ctx, upload) - return nil -} - func loadUploadStats(ctx context.Context) ([]model.UploadStat, error) { return repository.ListUploadStats(ctx) } -var errUploadForbidden = errors.New("upload forbidden") - -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) - } - driver, backend, err := storage.Active(ctx) - if err != nil { - logger.ErrorF(ctx, "初始化活动存储失败: %v", err) - 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) - } - meta.Bucket = result.Bucket - return result.Key, nil -} - -func validateUploadAllowedExtension(ctx context.Context, ext string) string { - sc, err := repository.GetSystemConfigByKey(ctx, model.ConfigKeyUploadAllowedExtensions) - if err != nil || sc.Value == "" { - return "" - } - allowedExts := strings.Split(strings.ToLower(sc.Value), ",") - for _, allowedExt := range allowedExts { - if strings.TrimSpace(allowedExt) == ext { - return "" - } - } - return shared.ErrUnsupportedFormat -} - func isRecordNotFound(err error) bool { return errors.Is(err, gorm.ErrRecordNotFound) } \ No newline at end of file diff --git a/internal/apps/upload/handler/routers.go b/internal/apps/upload/handler/routers.go index 57ff9c01..6940ebfb 100644 --- a/internal/apps/upload/handler/routers.go +++ b/internal/apps/upload/handler/routers.go @@ -8,7 +8,6 @@ package handler import ( "archive/zip" "bytes" - "context" "crypto/sha256" "encoding/hex" "encoding/json" @@ -21,16 +20,15 @@ import ( "path/filepath" "strconv" "strings" - "time" "github.com/Rain-kl/Wavelet/internal/apps/oauth" "github.com/Rain-kl/Wavelet/internal/apps/upload/filesrv" + "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" "github.com/Rain-kl/Wavelet/internal/apps/upload/util" "github.com/Rain-kl/Wavelet/internal/common" "github.com/Rain-kl/Wavelet/internal/common/response" - "github.com/Rain-kl/Wavelet/internal/db/idgen" "github.com/Rain-kl/Wavelet/internal/model" "github.com/Rain-kl/Wavelet/pkg/logger" "github.com/gin-gonic/gin" @@ -91,11 +89,6 @@ func UploadFile(c *gin.Context) { ext = "bin" } - if errMsg := validateUploadAllowedExtension(ctx, ext); errMsg != "" { - response.AbortBadRequest(c, errMsg) - return - } - hashWriter := sha256.New() var buf bytes.Buffer size, err := io.Copy(&buf, io.TeeReader(file, hashWriter)) @@ -120,51 +113,43 @@ func UploadFile(c *gin.Context) { return } - handled, lookupErr := tryInstantUpload(ctx, c, currUser, fileHash, size, mimeType, ext, origName, accessMode) - if handled { - return - } - if lookupErr != nil && !errors.Is(lookupErr, gorm.ErrRecordNotFound) { - response.AbortBadRequest(c, shared.ErrFileValidationFailed) - return - } - meta, errMsg := parseUploadMetadata(c, mimeType) if errMsg != "" { response.AbortBadRequest(c, errMsg) return } - id := idgen.NextUint64ID() - subPath := fmt.Sprintf("uploads/%s/%d.%s", time.Now().Format("2006/01/02"), id, ext) - - subPath, err = storeUploadObject(ctx, subPath, size, mimeType, &buf, &meta) - if err != nil { - response.AbortBadRequest(c, err.Error()) - return - } - - newUpload := model.Upload{ - ID: id, + result, err := ingest.Ingest(ctx, ingest.Request{ UserID: currUser.ID, + Reader: bytes.NewReader(buf.Bytes()), + Size: size, FileName: origName, - FilePath: subPath, - FileSize: size, MimeType: mimeType, Extension: ext, Hash: fileHash, Type: uploadType, - Status: model.UploadStatusUsed, - AccessMode: accessMode, + AccessMode: &accessMode, Metadata: meta, - } - - if err := saveNewUploadRecord(ctx, &newUpload, subPath); err != nil { + Policy: ingest.PolicyDedupNewRecord, + }) + if err != nil { + if errors.Is(err, ingest.ErrStorageReadOnly) { + response.AbortConflict(c, shared.ErrStorageReadOnly) + return + } + if err.Error() == shared.ErrUnsupportedFormat { + response.AbortBadRequest(c, shared.ErrUnsupportedFormat) + return + } + if err.Error() == shared.ErrSaveFileFailed { + response.AbortBadRequest(c, shared.ErrSaveFileFailed) + return + } response.AbortBadRequest(c, shared.ErrSaveUploadRecordFailed) return } - c.JSON(http.StatusOK, response.OK(newUpload)) + c.JSON(http.StatusOK, response.OK(result.Upload)) } // DownloadFile 通用单文件下载接口 @@ -320,34 +305,6 @@ func resolveUploadAccessMode(c *gin.Context, uploadType string) (int, string) { return accessMode, "" } -func tryInstantUpload(ctx context.Context, c *gin.Context, currUser *model.User, fileHash string, size int64, mimeType, ext, origName string, accessMode int) (bool, error) { - existing, err := findReusableUpload(ctx, fileHash, size) - if err != nil { - return false, err - } - if uploadstorage.ReadOnly(ctx) { - response.AbortConflict(c, shared.ErrStorageReadOnly) - return true, nil - } - - newUpload, err := createInstantUpload(ctx, existing, instantUploadInput{ - UserID: currUser.ID, - FileHash: fileHash, - Size: size, - MimeType: mimeType, - Extension: ext, - OrigName: origName, - UploadType: c.DefaultPostForm("type", "generic"), - AccessMode: accessMode, - }) - if err != nil { - response.AbortBadRequest(c, shared.ErrSaveUploadRecordFailed) - return true, err - } - c.JSON(http.StatusOK, response.OK(newUpload)) - return true, nil -} - func parseUploadMetadata(c *gin.Context, mimeType string) (model.UploadMetadata, string) { var meta model.UploadMetadata metadataStr := c.DefaultPostForm("metadata", "") diff --git a/internal/apps/upload/ingest/errors.go b/internal/apps/upload/ingest/errors.go new file mode 100644 index 00000000..6ae1461a --- /dev/null +++ b/internal/apps/upload/ingest/errors.go @@ -0,0 +1,16 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package ingest + +import ( + "errors" + + "github.com/Rain-kl/Wavelet/internal/apps/upload/shared" +) + +// ErrForbidden indicates the caller is not allowed to mutate the upload record. +var ErrForbidden = errors.New("upload forbidden") + +// ErrStorageReadOnly indicates the storage backend is in migration read-only mode. +var ErrStorageReadOnly = errors.New(shared.ErrStorageReadOnly) \ No newline at end of file diff --git a/internal/apps/upload/ingest/helpers.go b/internal/apps/upload/ingest/helpers.go new file mode 100644 index 00000000..6097ef69 --- /dev/null +++ b/internal/apps/upload/ingest/helpers.go @@ -0,0 +1,187 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package ingest + +import ( + "context" + "errors" + "fmt" + "io" + "strings" + "time" + + "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/idgen" + "github.com/Rain-kl/Wavelet/internal/model" + "github.com/Rain-kl/Wavelet/internal/repository" + "github.com/Rain-kl/Wavelet/internal/storage" + "github.com/Rain-kl/Wavelet/pkg/logger" + "gorm.io/gorm" +) + +func normalizeRequest(req *Request) { + req.Extension = strings.ToLower(strings.TrimSpace(req.Extension)) + if req.Extension == "" { + req.Extension = "bin" + } + if req.Type == "" { + req.Type = "generic" + } + if req.Status == "" { + req.Status = model.UploadStatusUsed + } +} + +func resolveAccessMode(uploadType string, explicit *int) int { + if explicit != nil { + return *explicit + } + if uploadType == shared.DefaultPublicUploadType { + return 1 + } + return 0 +} + +func validateAllowedExtension(ctx context.Context, ext string) error { + sc, err := repository.GetSystemConfigByKey(ctx, model.ConfigKeyUploadAllowedExtensions) + if err != nil { + if errors.Is(err, gorm.ErrRecordNotFound) { + return nil + } + return err + } + if sc.Value == "" { + return nil + } + + allowedExts := strings.Split(strings.ToLower(sc.Value), ",") + for _, allowedExt := range allowedExts { + if strings.TrimSpace(allowedExt) == ext { + return nil + } + } + return errors.New(shared.ErrUnsupportedFormat) +} + +func defaultObjectKey(id uint64, ext string) string { + return fmt.Sprintf("uploads/%s/%d.%s", time.Now().Format("2006/01/02"), id, ext) +} + +func buildObjectKey(req Request, id uint64) string { + if req.ObjectKeyFn != nil { + return req.ObjectKeyFn(id, req.Extension) + } + return defaultObjectKey(id, req.Extension) +} + +func storeObject(ctx context.Context, objectKey string, reader io.Reader, size int64, mimeType string, meta *model.UploadMetadata) (string, error) { + if uploadstorage.ReadOnly(ctx) { + return "", ErrStorageReadOnly + } + + driver, backend, err := storage.Active(ctx) + if err != nil { + logger.ErrorF(ctx, "初始化活动存储失败: %v", err) + return "", errors.New(shared.ErrSaveFileFailed) + } + + result, err := backend.Put(ctx, objectKey, reader, size, mimeType) + if err != nil { + logger.ErrorF(ctx, "写入 %s 存储失败: %v", driver, err) + return "", errors.New(shared.ErrSaveFileFailed) + } + + meta.Bucket = result.Bucket + return result.Key, nil +} + +func persistUploadRecord(ctx context.Context, upload *model.Upload, objectKey string) error { + if err := repository.CreateUpload(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) + } + } + return err + } + uploadstats.RecordUploadStatsAdd(ctx, upload) + return nil +} + +func createDedupRecord(ctx context.Context, existing model.Upload, req Request) (Result, error) { + accessMode := resolveAccessMode(req.Type, req.AccessMode) + newUpload := model.Upload{ + ID: idgen.NextUint64ID(), + UserID: req.UserID, + FileName: req.FileName, + FilePath: existing.FilePath, + FileSize: req.Size, + MimeType: req.MimeType, + Extension: req.Extension, + Hash: req.Hash, + Type: req.Type, + Status: req.Status, + AccessMode: accessMode, + Metadata: existing.Metadata, + } + if err := persistUploadRecord(ctx, &newUpload, existing.FilePath); err != nil { + return Result{}, err + } + logger.InfoF(ctx, "文件触发秒传成功! ID: %d, Path: %s", newUpload.ID, existing.FilePath) + return Result{ + Upload: newUpload, + Created: true, + Stored: false, + }, nil +} + +func uploadstorageReadOnly(ctx context.Context) bool { + return uploadstorage.ReadOnly(ctx) +} + +func createNewUpload(ctx context.Context, req Request) (Result, error) { + if uploadstorageReadOnly(ctx) { + return Result{}, ErrStorageReadOnly + } + if !req.SkipExtensionCheck { + if err := validateAllowedExtension(ctx, req.Extension); err != nil { + return Result{}, err + } + } + + id := idgen.NextUint64ID() + objectKey := buildObjectKey(req, id) + storedKey, err := storeObject(ctx, objectKey, req.Reader, req.Size, req.MimeType, &req.Metadata) + if err != nil { + return Result{}, err + } + + accessMode := resolveAccessMode(req.Type, req.AccessMode) + upload := model.Upload{ + ID: id, + UserID: req.UserID, + FileName: req.FileName, + FilePath: storedKey, + FileSize: req.Size, + MimeType: req.MimeType, + Extension: req.Extension, + Hash: req.Hash, + Type: req.Type, + Status: req.Status, + AccessMode: accessMode, + Metadata: req.Metadata, + } + if err := persistUploadRecord(ctx, &upload, storedKey); err != nil { + return Result{}, err + } + + return Result{ + Upload: upload, + Created: true, + Stored: true, + }, nil +} \ No newline at end of file diff --git a/internal/apps/upload/ingest/ingest.go b/internal/apps/upload/ingest/ingest.go new file mode 100644 index 00000000..accbd438 --- /dev/null +++ b/internal/apps/upload/ingest/ingest.go @@ -0,0 +1,64 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package ingest + +import ( + "context" + "errors" + + "github.com/Rain-kl/Wavelet/internal/model" + "github.com/Rain-kl/Wavelet/internal/repository" + "gorm.io/gorm" +) + +// Ingest stores or resolves an upload using the configured policy and side effects. +func Ingest(ctx context.Context, req Request) (Result, error) { + normalizeRequest(&req) + if req.Hash == "" { + return Result{}, errors.New("ingest hash is required") + } + if req.Reader == nil { + return Result{}, errors.New("ingest reader is required") + } + if req.Size < 0 { + return Result{}, errors.New("ingest size must be non-negative") + } + + switch req.Policy { + case PolicyDedupNewRecord, PolicyResolveExisting: + return ingestWithHashPolicy(ctx, req) + case PolicyCreate: + return createNewUpload(ctx, req) + default: + return Result{}, errors.New("unsupported ingest policy") + } +} + +// FindByHash returns a reusable active upload with the same hash and size. +func FindByHash(ctx context.Context, hash string, size int64) (model.Upload, error) { + return repository.FindReusableUploadByHash(ctx, hash, size) +} + +func ingestWithHashPolicy(ctx context.Context, req Request) (Result, error) { + existing, err := repository.FindReusableUploadByHash(ctx, req.Hash, req.Size) + if err == nil { + switch req.Policy { + case PolicyResolveExisting: + return Result{ + Upload: existing, + Resolved: true, + }, nil + case PolicyDedupNewRecord: + if uploadstorageReadOnly(ctx) { + return Result{}, ErrStorageReadOnly + } + return createDedupRecord(ctx, existing, req) + } + } + if err != nil && !errors.Is(err, gorm.ErrRecordNotFound) { + return Result{}, err + } + + return createNewUpload(ctx, req) +} \ No newline at end of file diff --git a/internal/apps/upload/ingest/ingest_test.go b/internal/apps/upload/ingest/ingest_test.go new file mode 100644 index 00000000..1c4d54cf --- /dev/null +++ b/internal/apps/upload/ingest/ingest_test.go @@ -0,0 +1,283 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package ingest + +import ( + "bytes" + "context" + "crypto/sha256" + "encoding/hex" + "io" + "os" + "testing" + "time" + + "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" +) + +func TestIngestPolicyCreateIncrementsStats(t *testing.T) { + _, _, cleanup := testhelper.SetupTestEnvironment(t) + defer cleanup() + ctx := context.Background() + + content := []byte("\x89PNG\r\n\x1a\n\x00\x00\x00\rIHDR\x00\x00\x00\x01\x00\x00\x00\x01") + 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: "mirror.png", + MimeType: "image/png", + Extension: "png", + Hash: hex.EncodeToString(hash[:]), + Type: "pixez_mirror", + Policy: PolicyCreate, + }) + if err != nil { + t.Fatalf("Ingest(PolicyCreate) returned error: %v", err) + } + if !result.Created || !result.Stored || result.Resolved { + t.Fatalf("Ingest(PolicyCreate) = %+v, want Created+Stored without Resolved", result) + } + + stats, err := loadTotalStats(ctx) + if err != nil { + t.Fatalf("loadTotalStats returned error: %v", err) + } + if stats.TotalCount != 1 || stats.TotalSize != int64(len(content)) { + t.Fatalf("loadTotalStats() = count %d size %d, want count 1 size %d", stats.TotalCount, stats.TotalSize, len(content)) + } +} + +func TestIngestPolicyResolveExistingSkipsStatsOnHit(t *testing.T) { + dbConn, _, cleanup := testhelper.SetupTestEnvironment(t) + defer cleanup() + ctx := context.Background() + + content := []byte("\x89PNG\r\n\x1a\n\x00\x00\x00\rIHDR\x00\x00\x00\x01\x00\x00\x00\x01") + hash := sha256.Sum256(content) + hashStr := hex.EncodeToString(hash[:]) + + existing := model.Upload{ + ID: 88001, + UserID: 42, + FileName: "existing.png", + FilePath: "uploads/existing.png", + FileSize: int64(len(content)), + MimeType: "image/png", + Extension: "png", + Hash: hashStr, + Type: "pixez_mirror", + Status: model.UploadStatusUsed, + CreatedAt: time.Now(), + } + if err := dbConn.Create(&existing).Error; err != nil { + t.Fatalf("seed upload failed: %v", err) + } + + 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: "mirror.png", + MimeType: "image/png", + Extension: "png", + Hash: hashStr, + Type: "pixez_mirror", + Policy: PolicyResolveExisting, + }) + if err != nil { + t.Fatalf("Ingest(PolicyResolveExisting) returned error: %v", err) + } + if !result.Resolved || result.Created || result.Stored { + t.Fatalf("Ingest(PolicyResolveExisting) = %+v, want Resolved only", result) + } + if result.Upload.ID != existing.ID { + t.Fatalf("Ingest(PolicyResolveExisting).Upload.ID = %d, want %d", result.Upload.ID, existing.ID) + } + + stats, err := loadTotalStats(ctx) + if err != nil { + t.Fatalf("loadTotalStats returned error: %v", err) + } + if stats.TotalCount != 0 || stats.TotalSize != 0 { + t.Fatalf("loadTotalStats() = count %d size %d, want zero stats for resolved upload", stats.TotalCount, stats.TotalSize) + } +} + +func TestIngestPolicyDedupNewRecordCreatesSecondRecord(t *testing.T) { + dbConn, _, cleanup := testhelper.SetupTestEnvironment(t) + defer cleanup() + ctx := context.Background() + + content := []byte("\x89PNG\r\n\x1a\n\x00\x00\x00\rIHDR\x00\x00\x00\x01\x00\x00\x00\x01") + hash := sha256.Sum256(content) + hashStr := hex.EncodeToString(hash[:]) + putCount := 0 + + restoreStorage, disableStorage := setupMockStorage(t, &putCount) + defer restoreStorage() + defer disableStorage() + + first, err := Ingest(ctx, Request{ + UserID: 1001, + Reader: bytes.NewReader(content), + Size: int64(len(content)), + FileName: "first.png", + MimeType: "image/png", + Extension: "png", + Hash: hashStr, + Type: "avatar", + Policy: PolicyDedupNewRecord, + }) + if err != nil { + t.Fatalf("first Ingest returned error: %v", err) + } + if putCount != 1 { + t.Fatalf("putCount after first ingest = %d, want 1", putCount) + } + + second, err := Ingest(ctx, Request{ + UserID: 1002, + Reader: bytes.NewReader(content), + Size: int64(len(content)), + FileName: "second.png", + MimeType: "image/png", + Extension: "png", + Hash: hashStr, + Type: "avatar", + Policy: PolicyDedupNewRecord, + }) + if err != nil { + t.Fatalf("second Ingest returned error: %v", err) + } + if putCount != 1 { + t.Fatalf("putCount after dedup ingest = %d, want 1", putCount) + } + if first.Upload.FilePath != second.Upload.FilePath { + t.Fatalf("dedup file paths differ: %s vs %s", first.Upload.FilePath, second.Upload.FilePath) + } + if first.Upload.ID == second.Upload.ID { + t.Fatal("dedup records should have unique IDs") + } + + var count int64 + if err := dbConn.Model(&model.Upload{}).Where("hash = ?", hashStr).Count(&count).Error; err != nil { + t.Fatalf("count uploads failed: %v", err) + } + if count != 2 { + t.Fatalf("upload count = %d, want 2", count) + } +} + +func TestRemoveDecrementsStats(t *testing.T) { + _, _, cleanup := testhelper.SetupTestEnvironment(t) + defer cleanup() + ctx := context.Background() + + content := []byte("\x89PNG\r\n\x1a\n\x00\x00\x00\rIHDR\x00\x00\x00\x01\x00\x00\x00\x01") + 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: "delete-me.png", + MimeType: "image/png", + Extension: "png", + Hash: hex.EncodeToString(hash[:]), + Type: "generic", + Policy: PolicyCreate, + }) + if err != nil { + t.Fatalf("Ingest returned error: %v", err) + } + + if _, err := Remove(ctx, result.Upload.ID); err != nil { + t.Fatalf("Remove(%d) returned error: %v", result.Upload.ID, err) + } + + stats, err := loadTotalStats(ctx) + if err != nil { + t.Fatalf("loadTotalStats returned error: %v", err) + } + if stats.TotalCount != 0 || stats.TotalSize != 0 { + t.Fatalf("loadTotalStats() after remove = count %d size %d, want zero", stats.TotalCount, stats.TotalSize) + } +} + +type totalStatsSnapshot struct { + TotalCount int64 + TotalSize int64 +} + +func loadTotalStats(ctx context.Context) (totalStatsSnapshot, error) { + var rows []model.UploadStat + if err := db.DB(ctx).Where("dimension = ?", model.UploadStatDimensionTotal).Find(&rows).Error; err != nil { + return totalStatsSnapshot{}, err + } + if len(rows) == 0 { + return totalStatsSnapshot{}, nil + } + return totalStatsSnapshot{ + TotalCount: rows[0].FileCount, + TotalSize: rows[0].FileSize, + }, nil +} + +func setupMockStorage(t *testing.T, putCount *int) (restore func(), disable func()) { + t.Helper() + mockFiles := make(map[string][]byte) + restore = storage.MockStorage( + func(ctx context.Context, key string, body io.Reader, size int64, contentType string) error { + data, err := io.ReadAll(body) + if err != nil { + return err + } + mockFiles[key] = data + if putCount != nil { + *putCount++ + } + return nil + }, + func(ctx context.Context, key string) (*storage.Object, error) { + data, ok := mockFiles[key] + if !ok { + return nil, os.ErrNotExist + } + return &storage.Object{ + Body: io.NopCloser(bytes.NewReader(data)), + ContentLength: int64(len(data)), + ContentType: "application/octet-stream", + }, nil + }, + func(ctx context.Context, key string) error { + delete(mockFiles, key) + return nil + }, + ) + storage.IsEnabledFunc = func() bool { return true } + storage.ResetCache() + disable = func() { + storage.IsEnabledFunc = func() bool { return false } + storage.ResetCache() + } + return restore, disable +} \ No newline at end of file diff --git a/internal/apps/upload/ingest/remove.go b/internal/apps/upload/ingest/remove.go new file mode 100644 index 00000000..4ade14c2 --- /dev/null +++ b/internal/apps/upload/ingest/remove.go @@ -0,0 +1,43 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package ingest + +import ( + "context" + + uploadstats "github.com/Rain-kl/Wavelet/internal/apps/upload/stats" + "github.com/Rain-kl/Wavelet/internal/model" + "github.com/Rain-kl/Wavelet/internal/repository" +) + +// Remove soft-deletes an upload and decrements incremental stats. +func Remove(ctx context.Context, uploadID uint64) (model.Upload, error) { + upload, err := repository.GetActiveUploadByID(ctx, uploadID) + if err != nil { + return model.Upload{}, err + } + uploadstats.RecordUploadStatsRemove(ctx, &upload) + if err := repository.SoftDeleteUpload(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. +func RemoveOwned(ctx context.Context, userID, uploadID uint64) (model.Upload, error) { + upload, err := repository.GetActiveUploadByID(ctx, uploadID) + if err != nil { + return model.Upload{}, err + } + if upload.UserID != userID { + return model.Upload{}, ErrForbidden + } + uploadstats.RecordUploadStatsRemove(ctx, &upload) + if err := repository.SoftDeleteUpload(ctx, &upload); err != nil { + return model.Upload{}, err + } + upload.Status = model.UploadStatusDeleted + return upload, nil +} \ No newline at end of file diff --git a/internal/apps/upload/ingest/types.go b/internal/apps/upload/ingest/types.go new file mode 100644 index 00000000..bbeed19b --- /dev/null +++ b/internal/apps/upload/ingest/types.go @@ -0,0 +1,60 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +// Package ingest provides the programmatic upload domain service for Wavelet. +package ingest + +import ( + "io" + + "github.com/Rain-kl/Wavelet/internal/model" +) + +// Policy controls how ingest handles hash collisions and record creation. +type Policy int + +const ( + // PolicyCreate always stores a new object and creates a new upload record. + PolicyCreate Policy = iota + + // PolicyDedupNewRecord reuses an existing object path on hash match but creates a new record and stats delta. + PolicyDedupNewRecord + + // PolicyResolveExisting returns an existing upload on hash match without creating a record or stats delta. + PolicyResolveExisting +) + +// ObjectKeyFn builds the storage object key for a new upload. +type ObjectKeyFn func(id uint64, ext string) string + +// Request describes a programmatic file ingest operation. +type Request struct { + UserID uint64 + Type string + + AccessMode *int + Status model.UploadStatus + + Reader io.Reader + Size int64 + FileName string + MimeType string + Extension string + Hash string + + Metadata model.UploadMetadata + Policy Policy + + ObjectKeyFn ObjectKeyFn + + // SkipExtensionCheck bypasses the configured upload extension whitelist. + SkipExtensionCheck bool +} + +// Result reports the outcome of an ingest operation. +type Result struct { + Upload model.Upload + Created bool + Stored bool + Resolved bool +} \ No newline at end of file diff --git a/internal/repository/upload.go b/internal/repository/upload.go index 62fed841..e864136a 100644 --- a/internal/repository/upload.go +++ b/internal/repository/upload.go @@ -63,6 +63,7 @@ func GetActiveUploadByID(ctx context.Context, id uint64) (model.Upload, error) { } // SoftDeleteUpload marks an upload as deleted. +// External modules must use upload.Remove or upload.RemoveOwned; only internal/apps/upload may call this. func SoftDeleteUpload(ctx context.Context, upload *model.Upload) error { return db.DB(ctx).Model(upload).Update("status", model.UploadStatusDeleted).Error } @@ -97,6 +98,7 @@ func FindReusableUploadByHash(ctx context.Context, hash string, size int64) (mod } // CreateUpload persists a new upload record. +// External modules must use upload.Ingest; only internal/apps/upload may call this. func CreateUpload(ctx context.Context, upload *model.Upload) error { return db.DB(ctx).Create(upload).Error }