文件管理模块实现教程

2025-07-23T14:11:02+08:00 | 43分钟阅读 | 更新于 2025-07-23T14:11:02+08:00

@

学习目标

学完本章你应该能够:

  1. 讲清文件模块的分层架构(service → biz → data)每一层各自负责什么,以及 Storage / FileRepo / FolderRepo / UploadRepo 抽象接口如何做到"可切换存储与数据库实现"。
  2. 解释秒传(SHA-256 去重)断点续传的原理:为什么同内容只存一份物理文件、会话状态机 0/1/2/3 各代表什么、断点续传靠哪几个查询复用 uploadID。
  3. 说清分片上传的四步流程与并发控制:信号量(容量 50)限流、原子递增 chunks_received、合并时如何"要么全成要么不留垃圾"。
  4. 对比 keyset 分页 vs offset 分页覆盖索引的作用,并讲出"第一页缓存 + 抖动 TTL + 主动失效 pattern"这套缓存策略为什么能抗雪崩。
  5. 像商用工程师一样审查文件模块:每个写操作是否校验 user_id 归属(IDOR 防护)文件名路径穿越与危险类型拦截是否覆盖所有入口上传/下载限流与 nginx client_max_body_size 上限并发上传下的配额一致性(TOCTOU 与原子上限)删除/下载/分享是否留审计日志
  6. 面试时能围绕"孤儿文件回滚、零拷贝流式、危险文件类型拦截、配额防超卖、亿级流量优化"讲成一段有结构、有取舍的工程故事。

前置知识

  • Go 基础、GORM 基本用法、contextdefer 资源管理。
  • MySQL 索引(联合索引、覆盖索引)与 Redis 缓存基础。
  • Kratos 框架分层(transport / service / biz / data)的大致概念,以及 JWT 鉴权中间件(见上一篇"用户模块")。

本章你会动手做的事

  1. 跟读 ListFiles 的 keyset/offset 混合分页 + 缓存代码,自己推演"翻到第 3 页时 nextCursor 怎么算",并指出文件/文件夹游标语义不一致引发的翻页 bug。
  2. 在纸上画出 InitUpload → UploadPart → CheckParts → MergeParts 的状态流转,标出哪一步会触发 FailSession 或回滚 Storage。
  3. 对照"商用审查"清单(鉴权/输入校验/限流/配额一致性/审计)逐条核对你正在维护的文件服务,挑出 3 条最该先补的。

一、技术栈与中间件

文件管理模块采用 Kratos 微服务框架分层架构,结合 GORM、MySQL、Redis 缓存、SHA-256 哈希、Storage 抽象接口、分片上传等多种技术。下表汇总了每项技术的用途:

技术 / 中间件所属层用途说明
Kratos(go-kratos/v3)框架骨架提供 transport(gRPC + HTTP 双协议)、log、errors(perrors)等基础设施;service 层实现 FileServiceServer 接口对外暴露 RPC
GORM(gorm.io/gorm)data 层ORM 框架,封装 MySQL 操作;data 层通过 r.db.WithContext(ctx) 完成文件、文件夹、上传会话、分片的 CRUD
MySQL持久化文件元数据(files 表)、文件夹(folders 表)、上传会话(upload_sessions 表)、分片记录(upload_chunks 表)。files 表建 (user_id, parent_id, status, created_at) 联合索引覆盖列表查询
Redis 缓存(Cache 接口)缓存层缓存文件列表第一页(带抖动 TTL)与文件元数据(file:meta:{id},2 分钟 TTL);通过 DeleteByPattern 失效目录列表缓存
SHA-256 哈希biz 层计算文件内容哈希,用于秒传去重和断点续传识别;File.Hash 字段存储哈希值,FindByHash 查询同哈希文件
Storage 抽象接口biz 层定义 Upload / Download / InitMultipartUpload / UploadPart / CompleteMultipartUpload / GetFileURL 等,不依赖具体实现,可切换本地存储或 MinIO/OSS
分片上传biz 层大文件切分为多个分片(默认 5MB),支持初始化、上传分片、检查分片、合并、暂停/恢复、断点续传
分布式锁(Locker 接口)biz 层user:storage🔒{userID} 锁保护存储统计的原子更新,避免并发写入导致超卖/漂移
信号量(channel)限流biz 层uploadSemaphore 大小 50,限制全局分片写入并发数,防止 MySQL 写入与磁盘 I/O 争抢导致吞吐倒退
JWTAuthMiddleware服务端中间件复用用户模块的鉴权中间件,把 user_id 注入 context;除白名单外所有文件接口都需携带 Bearer Token
RateLimitMiddleware服务端中间件进程内按 IP 滑动窗口限流(600/分钟),非分布式、取 X-Forwarded-For 首值(见三.5 商用审查)
UUID(crypto/rand)biz 层生成文件 UUID 与存储文件名(uuid + 扩展名),避免文件名冲突;通过 newUUID() 生成 v4 UUID
MIME 类型推断biz 层detectMimeType 优先使用系统 mime 数据库,回退到内置映射(覆盖 Docker 精简镜像缺失 /etc/mime.types 的场景)
零拷贝 / 流式传输service 层Download 使用服务端流按 64KB 分块发送;StreamPreviewio.Copy 从 Storage reader 流式写 http.ResponseWriter,避免大文件全部加载内存
异步事件(EventPublisher)biz 层上传完成时发布 EventUploadCompleted,供下游消费(生成缩略图、索引、病毒扫描等),失败不影响主流程

二、实现思路流程(总体)

文件管理模块的整体实现思路按照"从读 → 写 → 大文件 → 整理 → 预览"的顺序展开:

类比:把网盘想象成一个图书馆。你进馆先查目录(ListFiles)、再建书架分区(CreateFolder);想存一本别人已存过的书,直接贴个副本标签就行(秒传);太厚的书拆成册分批上架(分片上传);之后还能挪架、复印副本(Move / Copy),最后在线借阅预览(Preview)。下面这张图就是这条"从读到预览"的主干流程。

flowchart LR
    A[用户进目录] --> B[ListFiles 查文件与文件夹]
    B --> C[CreateFolder 建子目录]
    C --> D{文件是否已存在同哈希?}
    D -- 是 --> E[SecUpload 秒传 零字节]
    D -- 否 --> F{文件大小?}
    F -- 小文件 --> G[Upload 直传 + Download 下载]
    F -- 大文件 --> H[InitUpload 分片初始化]
    H --> I[UploadPart 上传分片]
    I --> J[CheckParts 检查分片]
    J --> K[MergeParts 合并]
    E --> L[Move Copy 移动复制]
    G --> L
    K --> L
    L --> M[Preview 预览 image/video/pdf]
    M --> N[Trash 删除入回收站]
  1. 分页查询(ListFiles):用户进入目录后,首先调用 ListFiles 获取当前目录下的文件与文件夹。biz 层并行查询 fileRepo.ListByParentfolderRepo.ListByParent,混合返回;第一页结果写入 Redis 缓存(带抖动 TTL)。

  2. 新建文件夹(CreateFolder):用户在当前目录下创建子文件夹,通过 parentID 实现多级目录;创建后失效目录缓存。

  3. 秒传(SecUpload):上传前客户端先计算文件 SHA-256,调用 SecUpload 传入哈希。若用户已存在同哈希文件,直接复用存储路径新建一条 file 记录,零字节传输;否则返回 exists=false 让客户端走正常上传。

  4. 单文件上传与下载(Upload / Download):小文件走 Upload 直接写入 Storage,再写 files 表(失败时回滚 Storage 文件,避免孤儿文件);下载走 Download,先查缓存命中文件元数据,再从 Storage 读取字节流返回。service 层还提供 gRPC 客户端流式 Upload 与服务端流式 Download,避免大文件全部加载到内存(零拷贝优化)。

  5. 大文件分片上传:分四步——

    • 初始化(InitUpload):计算分片总数,生成 uploadID,调用 Storage 初始化分片目录,创建 UploadSession 记录;若同哈希同大小会话已存在则直接复用(断点续传)。
    • 上传分片(UploadPart):通过信号量限流后写入 Storage 分片目录,再创建 UploadChunk 记录并原子递增 chunks_received
    • 检查分片(CheckParts):返回已上传分片编号列表,供客户端断点续传时跳过已传分片。
    • 合并(MergeParts):校验分片完整后调用 CompleteMultipartUpload 合并分片,再写 files 表;失败时回滚 Storage 文件;成功后发布 EventUploadCompleted 事件。
  6. 移动与复制(Move / Copy):批量校验文件所有权(防止 IDOR)后,BatchMove 修改 parent_id;复制则保持 Path 不变(共享物理文件)但生成新 UUID 与记录,并对重名文件自动添加 (1)(2) 后缀。两操作完成后均失效源目录与目标目录缓存。

  7. 预览(Preview):根据文件 MIME 类型分类(image / video / audio / pdf / text / other),通过 Storage.GetFileURL 返回预签名 URL(OSS)或本地路径(本地存储);service 层还提供 StreamPreview 流式返回文件内容,按 64KB 分块写入 http.ResponseWriter

  8. 删除(Trash):批量校验所有权后把文件/文件夹软删除(status=1)并写入 recycle_bins,可恢复;永久删除时才回收对象存储文件与已用存储配额(见 5.8)。


三、面试常问知识点与难点

1. 秒传原理(SHA-256 哈希去重)

客户端上传前先计算文件内容的 SHA-256 哈希,调用秒传接口。服务端通过 user_id + hash + status=0 联合查询 files 表,若命中说明该用户已存在相同内容的文件,直接复用其 Path(物理存储路径)新建一条元数据记录,无需再传输文件字节。优点:节省带宽、降低存储成本。注意:哈希查询走 (user_id, hash) 联合索引,命中后秒传为 O(1) 复杂度。因为按 user_id 隔离,天然避免了一个用户复用另一个用户的物理文件(权限隔离);若要"跨用户全局秒传"共享同一份物理文件,需要引用计数 + 全局去重表,并对删除做引用回收。

2. 分片上传与断点续传

大文件被切分为多个 5MB 分片,每片独立上传、独立写表(upload_chunks)。断点续传通过 InitUpload 中查找 FindByHashAndSize(同 user_id + file_hash + file_size + status IN (0,3))实现:若已存在进行中或已暂停的会话,直接复用 uploadID,前端通过 CheckParts 获取已上传分片编号列表,跳过这些分片只传缺失的部分。

会话状态机有四个值,源码 biz/file.go 的运行态语义为:0=进行中 / 1=已完成 / 2=失败 / 3=已暂停PauseUpload 把 0 改为 3,ResumeSession 把 3 恢复为 0,合并成功置 1,FailSession 置 2)。

stateDiagram-v2
    [*] --> 进行中: InitUpload 创建会话 status=0
    进行中 --> 已暂停: PauseUpload status=3
    已暂停 --> 进行中: ResumeSession 恢复为 0
    进行中 --> 已完成: MergeParts 成功 CompleteSession
    进行中 --> 失败: 合并失败 FailSession
    已完成 --> [*]
    失败 --> [*]

⚠️ 修正点(源码注释口径冲突)model/upload_session.goapi/file/v1/file.protoGetUploadStatusReply 注释只写了 3 态(0/1/2)或把"已暂停"标成 1,与 biz/file.go 实际运行的 4 态(0/1/2/3)不一致。商用版应统一为 4 态,并以 biz 层运行时语义为准,避免前端按错误口径解析状态。

3. keyset 分页 vs offset 分页

  • offset 分页LIMIT n OFFSET m,深翻页时 MySQL 需扫描 m+n 行,越往后越慢;且数据插入时会出现重复或漏掉记录。
  • keyset 分页(游标分页):使用 WHERE created_at < ? ORDER BY created_at DESC LIMIT n,每次以最后一行的 created_at 作为下一页游标,深翻页性能稳定,且不受新插入数据影响。

本项目 folder 表使用 keyset 分页(基于 created_at),ListByCategory / ListStarred 也用 keyset 分页;file 表的 ListByParent 暂用 offset 分页(cursor 当作 offset 数字),后续可演进为 keyset。

⚠️ 修正点(翻页游标语义冲突的真实 bug)ListFiles 中文件走 offset 分页、返回的 fileCursor 是数字串(如 "20");文件夹走 keyset 分页、返回的 folderCursor 是时间戳串(如 "2025-07-23 14:11:02.123")。二者却被拿来直接 folderCursor > nextCursor 做字符串比较、并取较大值作为统一 nextCursor。数字串与时间戳串不可比较,翻页游标会错乱。 商用修正:文件与文件夹应使用统一的 keyset 游标(都以 created_at 为准),或分别维护 fileCursor/folderCursor 两个独立游标,绝不能混用一个 nextCursor。修正示意:

// 修正点:文件/文件夹各自维护 keyset 游标,不再用单一 nextCursor 混比
result := &ListResult{
    Files:       files,
    Folders:     folders,
    FileCursor:  fileCursor,   // 数字 offset 或 created_at 时间戳,二选一并统一
    FolderCursor: folderCursor,
    HasMore:     len(files) >= pageSize || len(folders) >= pageSize,
}

4. 覆盖索引优化

覆盖索引指查询字段全部包含在索引中,无需回表查询聚簇索引。例如 (user_id, parent_id, status, created_at) 联合索引可覆盖"按用户+目录+状态查询并按时间排序"的场景,MySQL 直接从索引返回结果,避免回表 IO。data/file.go 的 ListByParent 查询 user_id + status + parent_id 并按 created_at 排序,正是覆盖索引的典型应用。

5. 零拷贝传输与内存占用

传统上传下载需要把文件完整读入内存再处理,大文件易导致 OOM。本项目 service 层使用 gRPC 客户端流(FileService_UploadServer)累积数据但有 100MB 上限保护;下载使用服务端流(FileService_DownloadServer)按 64KB 分块流式发送;StreamPreview 直接 io.Copy(w, reader) 从 Storage reader 拷贝到 HTTP ResponseWriter,全程不持有完整文件内容。

⚠️ 修正点(“零拷贝"的边界):流式 Download/StreamPreview 确实是边读边发;但流式 Upload 在 service 层会先把整个请求体累积进一个 []byte(上限 100MB)再调用 biz.Upload(data []byte),即单文件上传路径仍会把整文件缓冲在内存。超大文件应改走分片上传,或让 storage 直接消费 io.Reader 流式写入、不经内存聚合。

6. 大文件分片合并策略

合并时先从数据库按 chunk_index ASC 读取所有分片元数据,校验数量等于 ChunksTotal 后调用 CompleteMultipartUpload 顺序合并。合并失败调用 FailSession 标记会话失败;合并成功但写 files 表失败时,回滚已合并的 Storage 文件避免孤儿;UUID 生成失败也回滚 Storage。整个流程保证"要么全部成功,要么不留垃圾”。

7. 并发上传的锁控制

分片上传使用 chan struct{} 信号量(容量 50)限制同时进行的分片写入数量。UploadPart 通过 select 等待信号量、超时(10s)返回 429 ErrUploadBusy、或 ctx 取消。此外,用户存储空间更新使用分布式锁 user:storage🔒{userID},避免并发增减导致计数错乱。

⚠️ 商用审查(全局信号量跨用户共享):信号量容量 50 是进程内、跨所有用户共享的。一个用户并发上传大文件即可占满 50 个槽位,让其他用户的 UploadPart 全部排隊甚至超时(429)。商用版应改为按用户维度的并发配额(如每用户 5~10 个并发分片),避免单用户饿死全局。

8. 缓存策略与失效

文件列表第一页缓存 60s ±30% 抖动 TTL(防止缓存雪崩),文件元数据缓存 2 分钟。失效策略:

  • 主动失效:上传 / 删除 / 移动 / 复制 / 重命名 / 创建文件夹后调用 invalidateFileListCache,按 files:list:{userID}:{parentID}:* 模式批量删除该目录所有排序组合的缓存。
  • 元数据失效:文件被修改时调用 InvalidateFileMetaCache(fileID) 删除 file:meta:{id} 缓存。
  • 短 TTL 兜底:即便主动失效遗漏,2 分钟后缓存也会自然过期。

类比:缓存就像超市货架上的临期食品标签。正常卖出(主动失效)立刻撕标签;万一店员忘了撕(主动失效遗漏),标签上印的"2 分钟过期"也会兜底——时间一到自动下架,绝不至于卖出去变质(读到脏数据)。而"抖动 TTL"则是故意让每批标签过期时间错开几秒,避免整排货架在同一秒集体下架造成抢补货的拥堵(缓存雪崩)。

flowchart TD
    A[客户端请求文件列表] --> B{第一页 cursor 为空?}
    B -- 是 --> C[查 Redis files:list:user:parent:sort]
    C --> D{命中?}
    D -- 是 --> E[直接返回 不查库]
    D -- 否 --> F[并行查 fileRepo + folderRepo]
    F --> G[写入 Redis 抖动 TTL 42s 加 随机 0-18s]
    G --> E
    B -- 否 翻页 --> F
    H[上传/删除/移动/复制] --> I[按 pattern 批量删 files:list:user:parent:*]
    I --> J[失效 file:meta:id 元数据缓存]
    J --> K[短 TTL 兜底 2 分钟后自然过期]

⚠️ 商用审查(pattern 失效成本)DeleteByPattern 在 Redis 集群下通常基于 SCAN,目录多/缓存键多时每次写操作都要扫一批 key。高频写入场景可改为"写操作后仅失效受影响目录的有限几个精确 key",或对列表缓存采用带命名空间前缀的短 TTL + 版本号(目录 version)失效,降低 SCAN 开销。

9. 孤儿文件清理与回滚

孤儿文件指 Storage 中存在但数据库无对应记录的文件。本项目在多个失败路径主动回滚:

  • Upload 写库失败 → storage.Delete(storagePath)
  • MergeParts UUID 生成失败 → storage.Delete(storagePath);写库失败 → storage.Delete(storagePath)
  • InitUpload 创建会话失败 → storage.AbortMultipartUpload(uploadID)

此外 Storage 接口定义了 CleanupOrphanChunks 方法,可扫描 chunks 目录清理早于 maxAge 的孤儿分片目录。

10. 危险文件类型与文件名净化(输入校验)

安全防护两道关卡(在 SecUploadInitUpload 已落地):

  • sanitizeFileName:调用 filepath.Base 取文件名(防路径穿越),替换 ../..\\<>"|?*: 等危险字符,限制长度 255。
  • isDangerousFileType:黑名单拦截 .exe.bat.sh.php.jsp.asp.py.so.dll 等可执行/脚本文件类型,防止恶意文件上传。

⚠️ 修正点(漏网之鱼:流式单文件上传 Upload 未做净化与类型拦截)biz.Upload(被 gRPC 流式 Upload 调用)直接 ext := extractExt(name) 并把客户端传入的 name 原样存入 files 表,且 service 层传的 fileType 是空串——既没有 sanitizeFileName 也没有 isDangerousFileType。虽然存储路径用的是 uuid+ext 而非用户文件名(路径穿越风险被 storage 层挡住),但文件名/类型字段可被写入任意内容,且危险类型文件绕过了拦截。商用修正:在 Upload 入口同样净化文件名并拦截危险类型:

// 修正点:单文件上传也必须净化文件名 + 拦截危险类型,与 SecUpload/InitUpload 一致
func (uc *FileUsecase) Upload(ctx context.Context, userID uint64, name, fileType string, parentID *uint64, data []byte) (*File, error) {
    fileSize := int64(len(data))
    name = sanitizeFileName(name)            // 修正点:补净化
    if isDangerousFileType(name) {           // 修正点:补危险类型拦截
        return nil, perrors.BadRequest("FILE_TYPE_FORBIDDEN", "不允许上传此类型的文件")
    }
    if uc.userUC != nil {
        if err := uc.userUC.CheckStorageAvailable(ctx, userID, fileSize); err != nil {
            return nil, err
        }
    }
    // ... 其余逻辑不变
}

11. 鉴权与 IDOR 防护(商用审查重点)

文件模块所有接口都经过 JWTAuthMiddleware,中间件把 user_id 注入 context,service 层用 CtxUserID(ctx) 取当前用户。但"认证通过"不等于"有权操作这条数据",每个写操作都必须再校验资源归属(IDOR 防护)。源码中已正确落地的归属校验:

flowchart TD
    Req[文件写操作] --> MW[JWTAuthMiddleware
注入 user_id] MW --> H[handler 取 CtxUserID] H --> B{biz 层查资源} B --> O{资源.user_id == 当前 user_id?} O -- 否 --> F[403 FORBIDDEN 越权拒绝] O -- 是 --> OK[执行操作]
  • Move / Copy / Rename:先 FindByIDs/FindByID 取出资源,逐个比对 f.UserID != userIDErrForbidden
  • UploadPart / CheckParts / MergeParts / InitUpload:比对 session.UserID != userIDErrForbidden
  • Download / GetFileStream / Preview:比对 file.UserID != userIDErrForbidden
  • 删除(回收站):RecycleUsecase.Trash 在循环里对每个 file/folder 比对 f.UserID != userIDErrForbidden

⚠️ 修正点(Preview 漏校验状态)Preview 只校验了 file.UserID != userID没有校验 file.Status。被移入回收站(status=1)或已永久删除(status=2)的文件仍能通过 Preview 拿到预览 URL。商用修正:与 Download/GetFileStream 对齐,加上 if file.Status != 0 { return ErrFileNotFound }

12. 配额一致性与并发防超卖(商用审查重点)

文件上传/合并/秒传/复制都会调用 userUC.AddUsedStorage,其内部用分布式锁 user:storage🔒{userID} 串行化,再调用 UpdateUsedStorageAtomicused_storage = used_storage + ? 的原子更新(减少时带 WHERE used_storage >= ? 防超卖)。

flowchart TD
    U[上传请求] --> C{CheckStorageAvailable
在锁外查可用空间} C -- 不足 --> R[413 STORAGE_INSUFFICIENT] C -- 充足 --> L[获取 user:storage:lock] L --> A[原子 UPDATE used_storage + delta] A --> RL{RowsAffected} RL -- 1 --> OK[成功] RL -- 0 --> N[ErrUserNotFound]

⚠️ 修正点(TOCTOU + 增量无上限,并发可超配额)CheckStorageAvailable获取分布式锁之前执行,而 UpdateUsedStorageAtomic 对"增加"分支的 SQL 是 WHERE id = ? 没有 used_storage + delta <= total_storage 的上限约束。于是两个并发上传都先通过前置检查(都看到可用空间充足),再各自串行执行原子自增——结果 used_storage 可能超过 total_storage,即并发下配额被突破商用修正:把配额校验移入锁内,或在原子更新上加上限条件,让"检查+扣减"在同一行锁内完成:

// 修正点:增量更新必须带上限约束,防止并发通过前置检查后超配额
result := r.db.WithContext(ctx).Model(&model.User{}).
    Where("id = ? AND used_storage + ? <= total_storage", userID, delta).
    UpdateColumn("used_storage", gorm.Expr("used_storage + ?", delta))
// RowsAffected == 0 → 区分 用户不存在 / 空间不足(同减少分支)

同时建议 CheckStorageAvailable 也在 updateStorageWithLock 内部、持锁状态下重新核对,彻底消除 TOCTOU。

13. 限流与防滥用(商用审查重点)

入口 RateLimitMiddleware(600, time.Minute)进程内按 IP 的滑动窗口限流(600 次/分钟),存在三点商用缺陷:

  1. 非分布式:多实例部署时各算各的,限流值被实例数放大(N 实例 ≈ 600×N/分钟)。
  2. 信任可伪造的 IP:取 X-Forwarded-For 首值作为客户端 IP,攻击者可伪造该头绕过/嫁祸;应使用 nginx 注入的 X-Real-IP$remote_addr)。
  3. 粒度太粗:全局 600/分钟套在所有路由(含上传/下载),且信号量 50 是全局共享,没有按用户的更严格配额(如单用户每分钟上传次数、单文件最大尺寸)。
flowchart TD
    Req[上传/下载请求] --> RL[进程内 RateLimit 600/分钟/IP]
    RL -- 超限 --> R[429]
    RL -- 通过 --> SEM{全局信号量 50}
    SEM -- 满 --> B[429 ErrUploadBusy]
    SEM -- 获取 --> Do[执行]
    Note[缺陷: 非分布式 + 信 XFF 首值 + 跨用户共享 + 无单用户配额]

商用修正:用 Redis + Lua 做分布式限流(令牌桶),按「用户 ID + 接口」维度限流(如单用户上传 60 次/分钟、下载 300 次/分钟),IP 维度仅作兜底;nginx 层对上传/下载路由设合理 client_max_body_size(见 5.x)。

14. 大文件与 nginx 上传上限(商用审查重点)

deploy/nginx/nginx.confclient_max_body_size 0; 表示上传体大小不限制,完全依赖后端控制。而后端仅在流式 Upload 路径有 100MB 内存上限(maxStreamUploadSize),分片上传的 UploadPart 没有任何单分片/单文件大小校验,且 InitUploadchunkSize 由客户端传入、未做服务端钳制。这意味着:

  • 恶意客户端可传超大盘古级文件撑爆磁盘;
  • chunkSize=1 会让 ChunksTotal = fileSize,产生天文数字的分片数;
  • 应用层缓冲 100MB 在极端并发下仍有 OOM 风险。

商用修正:nginx 对 /api/v1/file/upload/api/v1/upload/* 设明确上限(如 client_max_body_size 2g);应用层在 InitUpload 钳制 chunkSize(如 1MB~100MB),UploadPart 校验分片大小与序号范围,Upload 校验总大小 ≤ 配额。

15. 审计日志(商用审查重点)

当前仅有 FileAccessLogRepo 记录"最近访问"(用于最近文件列表),属于功能数据而非安全审计。商用网盘对删除、分享创建/访问、永久删除、异常下载等高风险动作应写结构化审计日志(userID、动作、资源 ID、IP、时间、结果),便于安全复盘、合规与追责。删除流程(RecycleUsecase.Trash/DeleteForever)目前未写审计日志,建议补充。


四、亿级流量优化思路

针对文件模块在高并发大流量场景的进一步优化:

  1. CDN 加速下载:将热点文件的访问 URL 接入 CDN(如阿里云 CDN、Cloudflare),用户就近从边缘节点拉取,减轻源站带宽压力。Storage.GetFileURL 返回的预签名 URL 可直接配置 CDN 回源。

  2. 分片上传并行化:当前 UploadPart 是串行调用,可在客户端并行上传多个分片(如 3~5 并发),配合服务端信号量限流,显著缩短大文件上传耗时。需保证分片顺序与 chunk_index 对应,合并时按索引排序。

  3. 对象存储直传(Presigned PUT):服务端只生成预签名 URL,客户端直接 PUT 到 OSS/S3,跳过应用服务器中转,节省应用带宽与 CPU。服务端只负责创建会话、记录分片、触发合并回调。

  4. 秒传减少带宽:同内容文件零字节传输,已是核心优化。可进一步在用户间做"全局秒传"(跨用户共享同哈希文件,需引用计数与权限隔离)。

  5. 分库分表:files 表按 user_id 取模分库分表(如 1024 库 × 4 表),解决单表亿级数据下的查询性能瓶颈。FindByHash 需走分片键 user_id 路由,避免广播查询。upload_chunks 表可按 upload_session_id 分表。

  6. 热点文件缓存:热门视频 / 图片的下载流可通过 Redis 或本地缓存(如 LRU)缓存文件字节或预签名 URL,减少对 Storage 的回源。文件元数据缓存(file:meta:{id})已实现,可扩展为多级缓存(本地 + Redis),并加缩略图缓存(图片/视频首帧缩略图按 thumb:{fileID}:{size} 缓存,避免每次预览重新生成)。

  7. 读写分离:列表查询走 MySQL 只读从库,写入走主库。ListFilesListByCategoryListStarred 等读多写少的场景可显著降低主库压力。

  8. 异步事件驱动 + 病毒扫描:上传完成后通过 EventPublisher 异步发布事件,下游消费生成缩略图、提取视频元数据、建立搜索索引(ES)、触发病毒扫描(商用必备,扫描完成前文件标记"待检"、下载/分享时拦截)。主流程不阻塞,提升上传接口响应速度。

  9. 存储空间校准任务:定时任务(每日凌晨)调用 RecalibrateStorage,从活跃文件重新计算实际已用空间,修正并发写入导致的计数偏差,避免长期累计误差。

  10. 孤儿分片清理:定时任务调用 CleanupOrphanChunks,扫描 chunks 目录删除早于阈值(如 24 小时)的子目录,回收磁盘空间,避免用户中途放弃上传留下的垃圾分片。


五、详细实现流程与代码解析

5.1 文件分页列表(keyset/offset 混合分页 + 覆盖索引 + 文件夹混合查询)

实现思路

ListFiles 是用户进入目录后调用的第一个接口,需要同时返回文件与文件夹,并支持按名称 / 大小 / 创建时间排序与翻页。biz 层并行调用 fileRepo.ListByParentfolderRepo.ListByParent,取两者游标较大者作为下一页游标。第一页结果(cursor 为空)写入 Redis 缓存,TTL 带抖动避免雪崩;后续翻页不缓存以避免过期数据。folder 表使用 keyset 分页(基于 created_at 比较),file 表暂用 offset 分页。

⚠️ 修正点(翻页游标混用 bug):见三.3。文件用 offset 游标(数字串)、文件夹用 keyset 游标(时间戳串),二者被直接字符串比较并合并为单一 nextCursor,会导致翻页错乱。商用版应统一游标语义或分而治之。

关键代码

biz 层 ListFiles(混合查询 + 缓存):

// ListFiles 使用键集分页返回父目录下的文件和文件夹
func (uc *FileUsecase) ListFiles(ctx context.Context, userID uint64, parentID *uint64, cursor string, pageSize int, sortBy, sortOrder string) (*ListResult, error) {
    // 默认每页 20 条
    if pageSize <= 0 {
        pageSize = defaultPageSize
    }
    // 默认升序
    if sortOrder == "" {
        sortOrder = "asc"
    }

    // 仅第一页(cursor 为空)才查缓存,避免翻页时拿到过期数据
    cacheKey := ""
    if cursor == "" && uc.cache != nil {
        // 拼接缓存键:files:list:{userID}:{parentID}:{sortBy}:{sortOrder}
        parentStr := "root"
        if parentID != nil {
            parentStr = strconv.FormatUint(*parentID, 10)
        }
        cacheKey = fmt.Sprintf("files:list:%d:%s:%s:%s", userID, parentStr, sortBy, sortOrder)

        // 尝试缓存命中,命中则反序列化后直接返回
        cached, err := uc.cache.Get(ctx, cacheKey)
        if err == nil && cached != "" {
            var result ListResult
            if json.Unmarshal([]byte(cached), &result) == nil {
                return &result, nil
            }
        }
    }

    // 并行查询文件与文件夹(两者使用同一 cursor 与 limit)
    files, fileCursor, err := uc.fileRepo.ListByParent(ctx, userID, parentID, 0, cursor, pageSize, sortBy, sortOrder)
    if err != nil {
        return nil, err
    }
    folders, folderCursor, err := uc.folderRepo.ListByParent(ctx, userID, parentID, cursor, pageSize, sortBy, sortOrder)
    if err != nil {
        return nil, err
    }

    // 下一页游标取两者较大值,确保两边都翻到下一页
    // 修正点:fileCursor 是数字 offset、folderCursor 是时间戳,二者不可比较;
    // 商用版应改为分别返回 FileCursor / FolderCursor(见三.3)
    nextCursor := fileCursor
    if folderCursor > nextCursor {
        nextCursor = folderCursor
    }

    // 组装返回结果,HasMore 判断是否还有下一页
    result := &ListResult{
        Files:      files,
        Folders:    folders,
        NextCursor: nextCursor,
        HasMore:    len(files) >= pageSize || len(folders) >= pageSize,
    }

    // 第一页写入缓存,TTL = 42s + 随机 [0, 18)s,即 60s ±30% 抖动,防止缓存雪崩
    if cursor == "" && cacheKey != "" && uc.cache != nil {
        if data, err := json.Marshal(result); err == nil {
            ttl := 42*time.Second + time.Duration(time.Now().Nanosecond()%18)*time.Second
            uc.cache.Set(ctx, cacheKey, string(data), ttl)
        }
    }

    return result, nil
}

data 层 folderRepo.ListByParent(keyset 分页核心实现):

// ListByParent 使用键集分页返回给定父文件夹下的子文件夹
func (r *folderRepo) ListByParent(ctx context.Context, userID uint64, parentID *uint64, cursor string, limit int, sortBy, sortOrder string) ([]*biz.Folder, string, error) {
    // 基础查询条件:限定用户
    query := r.db.WithContext(ctx).Model(&model.Folder{}).
        Where("user_id = ?", userID)

    // 处理 parent_id:nil 表示根级别(parent_id IS NULL)
    if parentID == nil {
        query = query.Where("parent_id IS NULL")
    } else {
        query = query.Where("parent_id = ?", *parentID)
    }

    // keyset 分页核心:根据游标(上一页最后一条的 created_at)做范围查询
    // 降序查小于游标的,升序查大于游标的,避免 offset 深翻页性能问题
    if cursor != "" {
        if sortOrder == "desc" {
            query = query.Where("created_at < ?", cursor)
        } else {
            query = query.Where("created_at > ?", cursor)
        }
    }

    // 排序规则:按名称或按创建时间
    orderClause := "created_at ASC"
    if sortBy == "name" {
        orderClause = fmt.Sprintf("name %s", sortOrder)
    } else {
        if sortOrder == "desc" {
            orderClause = "created_at DESC"
        }
    }

    // 多查一条(limit + 1)用于判断是否还有下一页
    var pos []model.Folder
    if err := query.Order(orderClause).Limit(limit + 1).Find(&pos).Error; err != nil {
        return nil, "", err
    }

    // 截取前 limit 条,多余的用于 HasMore 判断
    hasMore := len(pos) > limit
    if hasMore {
        pos = pos[:limit]
    }

    // 下一页游标 = 当前页最后一条的 created_at(带毫秒精度)
    var nextCursor string
    if len(pos) > 0 {
        nextCursor = pos[len(pos)-1].CreatedAt.Format("2006-01-02 15:04:05.000")
    }

    // PO -> 领域对象转换
    folders := make([]*biz.Folder, len(pos))
    for i := range pos {
        folders[i] = toBizFolder(&pos[i])
    }
    return folders, nextCursor, nil
}

5.2 新建文件夹(多级目录支持)

实现思路

文件夹通过 parent_id 外键自引用实现多级目录。parent_id = nil 表示根目录,parent_id = 某文件夹ID 表示子目录。创建时先净化名称(防止路径穿越与危险字符),生成 UUID,写入 folders 表,最后失效当前目录的列表缓存。ListAllSubDirectoryIDs 递归查询某目录下所有子目录 ID,用于分类查询(如查询某目录及所有子目录下的图片)。

关键代码

biz 层 CreateFolder

func (uc *FileUsecase) CreateFolder(ctx context.Context, userID uint64, name string, parentID *uint64) (*Folder, error) {
    // 净化文件夹名称:取 base 名、替换危险字符、限制长度
    name = sanitizeFileName(name)
    if name == "" {
        return nil, perrors.BadRequest("FOLDER_NAME_REQUIRED", "文件夹名称不能为空")
    }
    // 生成 UUID v4 作为文件夹唯一标识
    uuid, err := newUUID()
    if err != nil {
        return nil, err
    }
    // 构造领域对象:parentID 为 nil 表示根目录,非 nil 表示子目录(多级支持)
    folder := &Folder{
        UUID:     uuid,
        UserID:   userID,
        Name:     name,
        ParentID: parentID,
    }
    // 调用仓储写入数据库
    created, err := uc.folderRepo.Create(ctx, folder)
    if err != nil {
        return nil, err
    }

    // 失效当前目录的列表缓存(按 pattern 批量删除该目录所有排序组合)
    uc.invalidateFileListCache(ctx, userID, parentID)

    return created, nil
}

data 层 ListAllSubDirectoryIDs(递归获取所有子目录):

// ListAllSubDirectoryIDs 递归获取指定目录下的所有子目录 ID(不包含自身)
// 用于分类查询场景:查询某目录及其所有子目录下的文件
func (r *folderRepo) ListAllSubDirectoryIDs(ctx context.Context, userID uint64, parentID uint64) ([]uint64, error) {
    var ids []uint64
    // 查询直接子目录
    var children []model.Folder
    if err := r.db.WithContext(ctx).
        Where("user_id = ? AND parent_id = ?", userID, parentID).
        Find(&children).Error; err != nil {
        return nil, err
    }
    // 递归查询每个子目录的孙目录
    for _, c := range children {
        ids = append(ids, c.ID)
        subIDs, err := r.ListAllSubDirectoryIDs(ctx, userID, c.ID)
        if err != nil {
            return nil, err
        }
        ids = append(ids, subIDs...)
    }
    return ids, nil
}

5.3 秒传(SHA-256 哈希去重)

实现思路

秒传的核心是"同内容只存一份物理文件"。客户端上传前计算文件 SHA-256 哈希,调用 SecUpload 接口。服务端用 user_id + hash + status=0 在 files 表中查询,若命中说明该用户已有相同内容文件,直接复用其 Path(物理存储路径)新建一条元数据记录,零字节传输。返回 exists=true 表示秒传成功,exists=false 表示未命中需走正常上传。

类比:秒传就像"公司打印室"。你要打印的文档哈希值和同事上周打过的一模一样——打印室不会真的再印一份,只是在你的文件柜里贴一张"这份文档归你"的便签,物理纸张还是那一份。省纸(存储)又省时间(带宽)。

下面这张时序图把"客户端算哈希 → 服务端查重 → 复用 Path 新建记录"的过程画清楚:

sequenceDiagram
    participant C as 客户端
    participant S as 服务端 biz
    participant DB as MySQL files 表
    C->>C: 计算文件 SHA-256 哈希
    C->>S: SecUpload(hash, name, parentID)
    S->>DB: FindByHash(user_id, hash, status=0)
    alt 命中同哈希文件
        DB-->>S: 返回 existing 物理路径
        S->>DB: 新建 file 记录 复用 existing.Path
        S-->>C: exists=true 秒传成功 零字节传输
    else 未命中
        DB-->>S: NotFound
        S-->>C: exists=false 走正常上传
    end

关键代码

biz 层 SecUpload

func (uc *FileUsecase) SecUpload(ctx context.Context, userID uint64, hash, name string, parentID *uint64) (*File, bool, error) {
    // 安全校验:净化文件名 + 拦截危险文件类型
    name = sanitizeFileName(name)
    if isDangerousFileType(name) {
        return nil, false, perrors.BadRequest("FILE_TYPE_FORBIDDEN", "不允许上传此类型的文件")
    }

    // 核心步骤 1:按 user_id + hash 查询是否已有相同内容的文件
    existing, err := uc.fileRepo.FindByHash(ctx, userID, hash)
    if err != nil {
        // 未找到 -> 不存在同哈希文件,返回 exists=false 让客户端走正常上传
        if perrors.IsNotFound(err) {
            return nil, false, nil
        }
        return nil, false, err
    }

    // 核心步骤 2:校验存储空间是否充足(秒传也要占用用户配额)
    if uc.userUC != nil {
        if err := uc.userUC.CheckStorageAvailable(ctx, userID, existing.Size); err != nil {
            return nil, false, err
        }
    }

    // 生成新 UUID(每条 file 记录独立 UUID)
    uuid, err := newUUID()
    if err != nil {
        return nil, false, err
    }
    // 核心步骤 3:复用 existing.Path(物理存储路径相同),仅新建元数据记录
    // 大小、类型、哈希、路径全部复用,零字节传输
    file := &File{
        UUID:     uuid,
        UserID:   userID,
        Name:     name,
        Size:     existing.Size,
        Type:     existing.Type,
        Hash:     existing.Hash,
        Path:     existing.Path,
        ParentID: parentID,
        Status:   0,
    }
    created, err := uc.fileRepo.Create(ctx, file)
    if err != nil {
        return nil, false, err
    }

    // 累加用户已用存储空间(带分布式锁,原子更新)
    if uc.userUC != nil {
        if err := uc.userUC.AddUsedStorage(ctx, userID, existing.Size); err != nil {
            log.Error("file: failed to add used storage after copy upload", "userID", userID, "err", err)
        }
    }

    // 失效目录列表缓存
    uc.invalidateFileListCache(ctx, userID, parentID)

    // 返回 exists=true 表示秒传成功
    return created, true, nil
}

data 层 FindByHash

// FindByHash 根据特定用户的 SHA-256 哈希值检索文件
// 用于秒传:如果该用户存在相同哈希值的文件则返回
// 走 (user_id, hash) 联合索引,命中后 O(1) 返回
func (r *fileRepo) FindByHash(ctx context.Context, userID uint64, hash string) (*biz.File, error) {
    var po model.File
    if err := r.db.WithContext(ctx).
        Where("user_id = ? AND hash = ? AND status = ?", userID, hash, 0). // status=0 仅查正常文件
        First(&po).Error; err != nil {
        if errors.Is(err, gorm.ErrRecordNotFound) {
            return nil, biz.ErrFileNotFound
        }
        return nil, err
    }
    return toBizFile(&po), nil
}

5.4 单文件上传与下载(流式 + 文件名净化修正)

实现思路

小文件上传走 Upload:先校验存储空间,生成 UUID 与存储文件名(UUID + 扩展名),调用 Storage.Upload 写入存储,再写 files 表。失败回滚:写库失败时调用 storage.Delete 删除已上传的文件,避免孤儿文件。下载走 Download:先从缓存(file:meta:{id})查文件元数据,未命中再查库并写缓存,再从 Storage 读取字节流返回。service 层还提供 gRPC 流式版本:客户端流式 Upload(累积数据但有 100MB 上限保护防 OOM),服务端流式 Download(按 64KB 分块流式发送)。

⚠️ 修正点(单文件上传漏做文件名净化与危险类型拦截):见三.10。当前 biz.Upload 未调用 sanitizeFileName / isDangerousFileType,商用版应在入口补齐,与 SecUpload/InitUpload 保持一致。下面代码已含修正。

关键代码

biz 层 Upload(含文件名净化修正 + 失败回滚):

func (uc *FileUsecase) Upload(ctx context.Context, userID uint64, name, fileType string, parentID *uint64, data []byte) (*File, error) {
    fileSize := int64(len(data))

    // 修正点:净化文件名 + 拦截危险类型(单文件上传此前漏做)
    name = sanitizeFileName(name)
    if isDangerousFileType(name) {
        return nil, perrors.BadRequest("FILE_TYPE_FORBIDDEN", "不允许上传此类型的文件")
    }

    // 校验存储空间是否充足
    if uc.userUC != nil {
        if err := uc.userUC.CheckStorageAvailable(ctx, userID, fileSize); err != nil {
            return nil, err
        }
    }

    // 生成 UUID 与存储文件名(UUID + 扩展名,避免冲突)
    uuid, err := newUUID()
    if err != nil {
        return nil, err
    }
    ext := extractExt(name)
    fileName := uuid + ext

    // 调用 Storage 接口写入物理文件,返回存储路径
    storagePath, err := uc.storage.Upload(ctx, fileName, bytesToReader(data))
    if err != nil {
        return nil, err
    }

    // 构造文件领域对象
    file := &File{
        UUID:     uuid,
        UserID:   userID,
        Name:     name,
        Size:     fileSize,
        Type:     fileType,
        Path:     storagePath,
        ParentID: parentID,
        Status:   0,
    }
    created, err := uc.fileRepo.Create(ctx, file)
    if err != nil {
        // 关键:数据库写入失败,回滚已上传的 Storage 文件,避免孤儿文件
        if delErr := uc.storage.Delete(ctx, storagePath); delErr != nil {
            log.Error("file: failed to rollback storage upload after db create failed", "storagePath", storagePath, "err", delErr)
        }
        return nil, err
    }

    // 累加用户已用存储(修正点:务必在分布式锁内做配额上限校验,见三.12)
    if uc.userUC != nil {
        if err := uc.userUC.AddUsedStorage(ctx, userID, fileSize); err != nil {
            log.Error("file: failed to add used storage after upload", "userID", userID, "err", err)
        }
    }

    // 失效目录列表缓存
    uc.invalidateFileListCache(ctx, userID, parentID)

    return created, nil
}

biz 层 getCachedFile + Download(缓存优先 + 归属校验):

// getCachedFile 从缓存获取文件元数据,未命中则查数据库并写入缓存
// 缓存键:file:meta:{fileID},TTL = 2 分钟
func (uc *FileUsecase) getCachedFile(ctx context.Context, fileID uint64) (*File, error) {
    cacheKey := fmt.Sprintf("file:meta:%d", fileID)
    // 尝试缓存命中
    if uc.cache != nil {
        if cached, err := uc.cache.Get(ctx, cacheKey); err == nil && cached != "" {
            var f File
            if json.Unmarshal([]byte(cached), &f) == nil {
                return &f, nil
            }
        }
    }

    // 缓存未命中,查询数据库
    file, err := uc.fileRepo.FindByID(ctx, fileID)
    if err != nil {
        return nil, err
    }

    // 写入缓存(失败不影响主流程)
    if uc.cache != nil {
        if data, err := json.Marshal(file); err == nil {
            _ = uc.cache.Set(ctx, cacheKey, string(data), fileMetaCacheTTL)
        }
    }
    return file, nil
}

func (uc *FileUsecase) Download(ctx context.Context, userID uint64, fileID uint64) (*File, []byte, string, error) {
    // 使用缓存获取文件元数据,减少 MySQL 查询
    file, err := uc.getCachedFile(ctx, fileID)
    if err != nil {
        return nil, nil, "", err
    }
    // 权限校验:仅文件所有者可下载(IDOR 防护)
    if file.UserID != userID {
        return nil, nil, "", ErrForbidden
    }
    // 状态校验:回收站/已删除文件不可下载
    if file.Status != 0 {
        return nil, nil, "", ErrFileNotFound
    }

    // 从 Storage 下载文件字节流
    reader, err := uc.storage.Download(ctx, file.Path)
    if err != nil {
        return nil, nil, "", err
    }
    defer reader.Close()

    // 读取全部字节(适合小文件;大文件走 service 层流式下载)
    data, err := ioReadAll(reader)
    if err != nil {
        return nil, nil, "", err
    }
    return file, data, file.Type, nil
}

service 层流式 Download(零拷贝分块发送):

// Download 服务端流式下载(仅 HTTP),按 64KB 分块发送,避免大文件全部加载到内存
func (s *FileService) Download(req *filev1.DownloadRequest, stream filev1.FileService_DownloadServer) error {
    userID := CtxUserID(stream.Context())
    if userID == 0 {
        return biz.ErrInvalidToken
    }
    // 使用流式读取,避免大文件全部加载到内存
    file, reader, _, err := s.uc.GetFileStream(stream.Context(), userID, req.FileId)
    if err != nil {
        return err
    }
    defer reader.Close()

    // 64KB 缓冲区分块发送
    buf := make([]byte, 64*1024)
    for {
        n, readErr := reader.Read(buf)
        if n > 0 {
            // 通过 gRPC 流发送当前分块(含文件名与总大小元信息)
            if err := stream.Send(&filev1.DownloadReply{
                Data: buf[:n],
                Name: file.Name,
                Size: file.Size,
            }); err != nil {
                return err
            }
        }
        if readErr != nil {
            break // EOF 或错误,结束循环
        }
    }
    return nil
}

5.5 大文件分片上传(初始化→上传分片→检查分片→合并)

实现思路

大文件分片上传分四步:

  1. 初始化(InitUpload):校验文件大小、净化文件名、检查存储空间;若传入 hash 则先查秒传(同哈希已有完成文件则直接返回 status=1);再查断点续传(同 user_id + hash + size 的进行中/已暂停会话);最后生成 uploadID,调用 Storage 初始化分片目录,创建 UploadSession 记录。
  2. 上传分片(UploadPart):通过信号量(容量 50)限流,写入 Storage 分片目录,创建 UploadChunk 记录,原子递增 chunks_received
  3. 检查分片(CheckParts):返回已上传分片编号列表,供客户端断点续传跳过已传分片。
  4. 合并(MergeParts):校验分片数等于 ChunksTotal,调用 CompleteMultipartUpload 合并,写 files 表,发布上传完成事件。

类比:分片上传像一个快递分拣中心。包裹(大文件)被拆成很多小箱(分片),每个小箱独立发往中转站(Storage 分片目录),每到一个就在台账上画一笔 chunks_received +1;等所有箱到齐,再封一个大包(合并)送出。中途某个箱丢了,不用重发全部——查台账(CheckParts)只看缺哪几箱补哪几箱,这就是断点续传。

会话状态机有四个值:0=进行中 / 1=已完成 / 2=失败 / 3=已暂停(注意 model/proto 注释口径需统一,见三.2)。

⚠️ 修正点(分片参数未校验,可被打爆)

  1. InitUploadchunkSize 由客户端传入却未钳制,恶意值(如 chunkSize=1)会让 ChunksTotal = fileSize,产生天文数字的分片数。服务端应钳制到合理范围(如 1MB~100MB)。
  2. MergePartslen(chunks) < ChunksTotal 判断完整性,< 允许"多传"的分片通过;应改为 != 并校验分片号连续且落在 [1, ChunksTotal]、单分片大小 ≤ 上限,合并后文件总大小与声明 FileSize 一致。
  3. UploadPart 幂等不彻底:本地存储在分片文件已存在时返回 ErrPartAlreadyUploaded,但 OSS/MinIO 实现未必去重,重试可能写重复分片;biz 层应显式按 (session_id, part_number) 幂等 upsert,或在 CreateChunk 前校验该分片号是否已存在。

关键代码

biz 层 InitUpload(初始化 + 秒传 + 断点续传,含 chunkSize 钳制修正):

func (uc *FileUsecase) InitUpload(ctx context.Context, userID uint64, fileName string, fileSize, chunkSize int64, parentID *uint64, hash string) (*UploadSession, error) {
    // 参数校验
    if fileSize <= 0 {
        return nil, perrors.BadRequest("INVALID_FILE_SIZE", "文件大小无效")
    }
    if chunkSize <= 0 {
        chunkSize = 5 * 1024 * 1024 // 默认 5MB 分片
    }

    // 修正点:客户端传入的 chunkSize 不可信任,钳制到合理区间,防止分片数爆炸
    const minChunk, maxChunk int64 = 1 << 20, 100 << 20 // 1MB ~ 100MB
    if chunkSize < minChunk || chunkSize > maxChunk {
        chunkSize = 5 * 1024 * 1024
    }

    // 安全校验:净化文件名 + 拦截危险文件类型
    fileName = sanitizeFileName(fileName)
    if isDangerousFileType(fileName) {
        return nil, perrors.BadRequest("FILE_TYPE_FORBIDDEN", "不允许上传此类型的文件")
    }

    // 校验存储空间
    if uc.userUC != nil {
        if err := uc.userUC.CheckStorageAvailable(ctx, userID, fileSize); err != nil {
            return nil, err
        }
    }

    if hash != "" {
        // 秒传检测:同哈希已有完成文件 -> 直接返回 status=1
        existing, _ := uc.fileRepo.FindByHash(ctx, userID, hash)
        if existing != nil && existing.Status == 0 {
            return &UploadSession{Status: 1}, nil
        }

        // 断点续传:查找同 user_id + hash + size 的进行中或已暂停会话
        if uc.uploadRepo != nil {
            resumed, _ := uc.uploadRepo.FindByHashAndSize(ctx, userID, hash, fileSize)
            if resumed != nil && (resumed.Status == 0 || resumed.Status == 3) {
                // 已暂停状态(status=3)先恢复为进行中(status=0)
                if resumed.Status == 3 {
                    _ = uc.uploadRepo.ResumeSession(ctx, resumed.UploadID)
                    resumed.Status = 0
                }
                return resumed, nil
            }
        }
    }

    // 计算分片总数:向上取整
    chunksTotal := int32((fileSize + chunkSize - 1) / chunkSize)
    // 生成 uploadID(upload_ + 随机 hex)
    uploadID := fmt.Sprintf("upload_%s", generateToken())

    // 使用 biz 层生成的 uploadID 初始化 storage 层分片目录,
    // 确保 biz 会话 ID 与 storage 目录名一致,避免孤儿空目录
    if err := uc.storage.InitMultipartUpload(ctx, uploadID, fileName); err != nil {
        return nil, err
    }

    // 构造上传会话领域对象
    session := &UploadSession{
        UploadID:    uploadID,
        UserID:      userID,
        FileName:    fileName,
        FileSize:    fileSize,
        FileHash:    hash,
        ChunkSize:   chunkSize,
        ChunksTotal: chunksTotal,
        Status:      0,
    }

    // 写入数据库;失败时取消 Storage 分片上传,避免孤儿空目录
    created, err := uc.uploadRepo.CreateSession(ctx, session)
    if err != nil {
        uc.storage.AbortMultipartUpload(ctx, uploadID)
        return nil, err
    }
    return created, nil
}

biz 层 UploadPart(信号量限流 + 原子递增 + 归属校验):

func (uc *FileUsecase) UploadPart(ctx context.Context, userID uint64, uploadID string, partNumber int, data []byte) (*UploadChunk, error) {
    // 校验会话存在且属于当前用户(IDOR 防护)
    session, err := uc.uploadRepo.FindSessionByUploadID(ctx, uploadID)
    if err != nil {
        return nil, err
    }
    if session.UserID != userID {
        return nil, ErrForbidden
    }
    // 仅进行中(status=0)的会话可上传分片
    if session.Status != 0 {
        return nil, ErrUploadSessionExpired
    }

    // 修正点:校验分片序号范围,防止越界/超大分片
    if partNumber < 1 || partNumber > int(session.ChunksTotal) {
        return nil, perrors.BadRequest("PART_OUT_OF_RANGE", "分片序号越界")
    }

    // 上传限流:信号量(容量 50)控制同时进行的分片写入数量
    semTimeout := 10 * time.Second
    select {
    case uc.uploadSemaphore <- struct{}{}: // 获取信号量
        defer func() { <-uc.uploadSemaphore }() // 函数结束时释放
    case <-time.After(semTimeout):
        return nil, ErrUploadBusy
    case <-ctx.Done():
        return nil, ctx.Err()
    }

    // 写入 Storage 分片目录
    _, err = uc.storage.UploadPart(ctx, uploadID, partNumber, bytesToReader(data))
    if err != nil {
        return nil, err
    }

    // 构造分片记录
    chunk := &UploadChunk{
        UploadSessionID: session.ID,
        ChunkIndex:      int32(partNumber),
        ChunkSize:       int64(len(data)),
    }

    // 写入 upload_chunks 表
    created, err := uc.uploadRepo.CreateChunk(ctx, chunk)
    if err != nil {
        return nil, err
    }
    // 原子递增会话的 chunks_received(用 gorm.Expr 避免 race condition)
    if err := uc.uploadRepo.IncrementChunks(ctx, uploadID); err != nil {
        return nil, err
    }
    return created, nil
}

data 层 IncrementChunks(原子递增):

// IncrementChunks 原子性地递增会话的已接收分片计数
// 使用 gorm.Expr("chunks_received + 1") 转换为 SQL: chunks_received = chunks_received + 1
// 数据库层保证原子性,避免并发 race condition
func (r *uploadRepo) IncrementChunks(ctx context.Context, uploadID string) error {
    return r.db.WithContext(ctx).
        Model(&model.UploadSession{}).
        Where("upload_id = ?", uploadID).
        UpdateColumn("chunks_received", gorm.Expr("chunks_received + 1")).Error
}

biz 层 CheckParts

func (uc *FileUsecase) CheckParts(ctx context.Context, userID uint64, uploadID string) ([]int32, error) {
    // 校验会话归属
    session, err := uc.uploadRepo.FindSessionByUploadID(ctx, uploadID)
    if err != nil {
        return nil, err
    }
    if session.UserID != userID {
        return nil, ErrForbidden
    }
    // 返回已上传分片编号列表(按 chunk_index 升序),客户端据此跳过已传分片
    return uc.uploadRepo.ListReceivedPartNumbers(ctx, session.ID)
}

biz 层 MergeParts(合并 + 回滚 + 事件发布,含完整性校验修正与配额上限):

func (uc *FileUsecase) MergeParts(ctx context.Context, userID uint64, uploadID string, parentID *uint64) (*File, error) {
    // 校验会话归属与状态
    session, err := uc.uploadRepo.FindSessionByUploadID(ctx, uploadID)
    if err != nil {
        return nil, err
    }
    if session.UserID != userID {
        return nil, ErrForbidden
    }
    if session.Status != 0 {
        return nil, ErrUploadSessionExpired
    }

    // 读取所有已上传分片,按 chunk_index 升序
    chunks, err := uc.uploadRepo.ListChunksBySession(ctx, session.ID)
    if err != nil {
        return nil, err
    }
    // 修正点:用 != 校验分片完整,且应校验分片号连续落在 [1, ChunksTotal]
    if int32(len(chunks)) != session.ChunksTotal {
        return nil, ErrUploadIncomplete
    }

    // 构造 PartInfo 列表传给 Storage 合并
    parts := make([]PartInfo, len(chunks))
    for i, c := range chunks {
        parts[i] = PartInfo{
            PartNumber: int(c.ChunkIndex),
            Size:       c.ChunkSize,
        }
    }

    // 调用 Storage 合并分片,返回最终存储路径
    storagePath, err := uc.storage.CompleteMultipartUpload(ctx, uploadID, parts)
    if err != nil {
        // 合并失败,标记会话为失败状态
        uc.uploadRepo.FailSession(ctx, uploadID)
        return nil, err
    }

    // 标记会话为已完成
    if err := uc.uploadRepo.CompleteSession(ctx, uploadID); err != nil {
        log.Error("file: failed to complete upload session", "uploadID", uploadID, "err", err)
    }

    // 生成新文件 UUID
    uuid, err := newUUID()
    if err != nil {
        // UUID 生成失败,回滚已合并的 Storage 文件,避免孤儿
        if delErr := uc.storage.Delete(ctx, storagePath); delErr != nil {
            log.Error("file: failed to rollback storage after uuid gen failed", "storagePath", storagePath, "err", delErr)
        }
        return nil, err
    }

    // 根据文件名推断 MIME 类型
    fileType := detectMimeType(session.FileName)

    // 构造文件领域对象
    file := &File{
        UUID:     uuid,
        UserID:   userID,
        Name:     session.FileName,
        Size:     session.FileSize,
        Type:     fileType,
        Path:     storagePath,
        ParentID: parentID,
        Status:   0,
    }
    created, err := uc.fileRepo.Create(ctx, file)
    if err != nil {
        // 数据库写入失败,回滚已合并的 Storage 文件
        if delErr := uc.storage.Delete(ctx, storagePath); delErr != nil {
            log.Error("file: failed to rollback storage merge after db create failed", "storagePath", storagePath, "err", delErr)
        }
        return nil, err
    }

    // 修正点:AddUsedStorage 必须带配额上限校验,避免并发超配额(见三.12)
    if uc.userUC != nil {
        if err := uc.userUC.AddUsedStorage(ctx, userID, session.FileSize); err != nil {
            log.Error("file: failed to add used storage after merge parts", "userID", userID, "err", err)
        }
    }

    // 失效目录列表缓存
    uc.invalidateFileListCache(ctx, userID, parentID)

    // 发布上传完成事件(失败不影响主流程),供下游消费(生成缩略图、索引、病毒扫描)
    if uc.eventPublisher != nil {
        _ = uc.eventPublisher.Publish(ctx, EventUploadCompleted, &UploadCompletedPayload{
            UserID:   userID,
            FileID:   created.ID,
            FileName: created.Name,
            FileSize: created.Size,
            Hash:     session.FileHash,
        })
    }

    return created, nil
}

5.6 文件移动与复制(批量操作 + 归属校验)

实现思路

移动(Move):批量校验文件所有权(用 map[*uint64]bool 记录源目录 ID,显式区分 nil 根目录与非 nil 子目录),调用 BatchMove 修改 parent_id;完成后失效目标目录与所有源目录的缓存,并失效被移动文件的元数据缓存(parent_id 已变更)。

复制(Copy):先预加载目标目录已有名称做重名检测;统计待复制文件总大小校验存储空间;循环创建新 file 记录(共享 Path 物理文件,但新 UUID),重名时调用 uniqueName 添加 (1)(2) 后缀;最后失效目标目录缓存。

两者都在循环里逐个比对 f.UserID != userID 返回 ErrForbidden,是 IDOR 防护的正确示范(见三.11)。

关键代码

biz 层 Move

func (uc *FileUsecase) Move(ctx context.Context, userID uint64, fileIDs, folderIDs []uint64, targetParentID *uint64) error {
    // 用 map[*uint64]bool 显式区分 nil(根目录)与非 nil(子目录)
    sourceParentIDs := make(map[*uint64]bool)

    // 校验文件所有权并记录源目录 ID
    if len(fileIDs) > 0 {
        files, err := uc.fileRepo.FindByIDs(ctx, fileIDs)
        if err != nil {
            return err
        }
        for _, f := range files {
            if f.UserID != userID {
                return ErrForbidden // IDOR 防护
            }
            sourceParentIDs[f.ParentID] = true
        }
        // 批量更新 parent_id(一条 SQL)
        if err := uc.fileRepo.BatchMove(ctx, fileIDs, targetParentID); err != nil {
            return err
        }
    }
    // 校验文件夹所有权并记录源目录 ID
    if len(folderIDs) > 0 {
        for _, fid := range folderIDs {
            folder, err := uc.folderRepo.FindByID(ctx, fid)
            if err != nil {
                return err
            }
            if folder.UserID != userID {
                return ErrForbidden // IDOR 防护
            }
            sourceParentIDs[folder.ParentID] = true
        }
        if err := uc.folderRepo.BatchMove(ctx, folderIDs, targetParentID); err != nil {
            return err
        }
    }
    // 失效目标目录缓存
    uc.invalidateFileListCache(ctx, userID, targetParentID)
    // 失效所有源目录缓存(包括 nil 表示的根目录)
    for pid := range sourceParentIDs {
        uc.invalidateFileListCache(ctx, userID, pid)
    }
    // 失效被移动文件的元数据缓存(parent_id 已变更)
    for _, fid := range fileIDs {
        uc.InvalidateFileMetaCache(ctx, fid)
    }
    return nil
}

biz 层 Copy + uniqueName

func (uc *FileUsecase) Copy(ctx context.Context, userID uint64, fileIDs, folderIDs []uint64, targetParentID *uint64) error {
    // 预加载目标目录下已有的文件和文件夹名称,用于重名检测
    existingNames := make(map[string]bool)
    existingFiles, _, err := uc.fileRepo.ListByParent(ctx, userID, targetParentID, 0, "", 1000, "", "")
    if err == nil {
        for _, f := range existingFiles {
            existingNames[f.Name] = true
        }
    }
    existingFolders, _, err := uc.folderRepo.ListByParent(ctx, userID, targetParentID, "", 1000, "", "")
    if err == nil {
        for _, f := range existingFolders {
            existingNames[f.Name] = true
        }
    }

    // 统计待复制文件总大小,用于存储空间校验
    var totalCopySize int64
    if len(fileIDs) > 0 {
        files, err := uc.fileRepo.FindByIDs(ctx, fileIDs)
        if err != nil {
            return err
        }
        for _, f := range files {
            if f.UserID != userID {
                return ErrForbidden // IDOR 防护
            }
            totalCopySize += f.Size
        }
    }
    // 校验存储空间
    if uc.userUC != nil && totalCopySize > 0 {
        if err := uc.userUC.CheckStorageAvailable(ctx, userID, totalCopySize); err != nil {
            return err
        }
    }

    // 逐个复制文件(共享 Path,新 UUID)
    for _, fid := range fileIDs {
        orig, err := uc.fileRepo.FindByID(ctx, fid)
        if err != nil {
            return err
        }
        if orig.UserID != userID {
            return ErrForbidden
        }
        uuid, err := newUUID()
        if err != nil {
            return err
        }
        // 重名检测:目标目录存在同名文件时自动添加 " (1)"、" (2)" 后缀
        copyName := uniqueName(orig.Name, existingNames)
        existingNames[copyName] = true
        copyFile := &File{
            UUID:     uuid,
            UserID:   userID,
            Name:     copyName,
            Size:     orig.Size,
            Type:     orig.Type,
            Hash:     orig.Hash,
            Path:     orig.Path, // 共享物理文件,不复制实际存储
            ParentID: targetParentID,
            Status:   0,
        }
        if _, err := uc.fileRepo.Create(ctx, copyFile); err != nil {
            return err
        }
        // 累加已用存储(修正点:同样需配额上限校验)
        if uc.userUC != nil {
            _ = uc.userUC.AddUsedStorage(ctx, userID, orig.Size)
        }
    }
    // 逐个复制文件夹(仅元数据,不递归复制子内容)
    for _, fid := range folderIDs {
        orig, err := uc.folderRepo.FindByID(ctx, fid)
        if err != nil {
            return err
        }
        if orig.UserID != userID {
            return ErrForbidden
        }
        uuid, err := newUUID()
        if err != nil {
            return err
        }
        copyName := uniqueName(orig.Name, existingNames)
        existingNames[copyName] = true
        copyFolder := &Folder{
            UUID:     uuid,
            UserID:   userID,
            Name:     copyName,
            ParentID: targetParentID,
        }
        if _, err := uc.folderRepo.Create(ctx, copyFolder); err != nil {
            return err
        }
    }
    // 失效目标目录缓存
    uc.invalidateFileListCache(ctx, userID, targetParentID)
    return nil
}

// uniqueName 在目标目录已有名称集合中生成不冲突的名称
// 若存在同名,则在新名称后添加 " (1)"、" (2)" 等后缀区分,不覆盖原有文件
func uniqueName(origName string, existing map[string]bool) string {
    if !existing[origName] {
        return origName
    }
    ext := extractExt(origName)
    base := origName[:len(origName)-len(ext)]
    for i := 1; ; i++ {
        candidate := fmt.Sprintf("%s (%d)%s", base, i, ext)
        if !existing[candidate] {
            return candidate
        }
    }
}

data 层 BatchMove

// BatchMove 将多个文件移动到新的父文件夹
// 一条 SQL 批量更新,避免循环单条更新
func (r *fileRepo) BatchMove(ctx context.Context, ids []uint64, newParentID *uint64) error {
    return r.db.WithContext(ctx).Model(&model.File{}).Where("id IN ?", ids).Update("parent_id", newParentID).Error
}

5.7 文件预览(按 MIME 类型分类)

实现思路

预览分两种模式:

  1. URL 预览(Preview):根据文件 MIME 类型分类(image / video / audio / pdf / text / other),调用 Storage.GetFileURL 返回预签名 URL(OSS)或本地路径(本地存储),有效期 1 小时。前端拿到 URL 后直接在浏览器渲染。
  2. 流式预览(StreamPreview):service 层通过 GetFileStream 获取 Storage reader,设置 Content-TypeContent-Length,用 io.Copy 流式写入 http.ResponseWriter,适合无法直接通过 URL 访问的场景(如本地存储或需要鉴权的文件)。

⚠️ 修正点(Preview 漏校验状态):见三.11。Preview 只校验了 file.UserID != userID,未校验 file.Status,回收站/已删除文件仍可预览。商用版应补 if file.Status != 0 { return ErrFileNotFound }

关键代码

biz 层 Preview + guessPreviewType(含状态校验修正):

func (uc *FileUsecase) Preview(ctx context.Context, userID uint64, fileID uint64) (string, string, error) {
    file, err := uc.fileRepo.FindByID(ctx, fileID)
    if err != nil {
        return "", "", err
    }
    // 权限校验:仅文件所有者可预览(IDOR 防护)
    if file.UserID != userID {
        return "", "", ErrForbidden
    }
    // 修正点:补状态校验,回收站/已删除文件不可预览
    if file.Status != 0 {
        return "", "", ErrFileNotFound
    }

    // 根据 MIME 类型推断预览类型(image/video/audio/pdf/text/other)
    previewType := guessPreviewType(file.Type)
    // 获取预签名 URL(OSS)或本地路径(本地存储),有效期 1 小时
    url, err := uc.storage.GetFileURL(ctx, file.Path, 1*time.Hour)
    if err != nil {
        return "", "", err
    }
    return url, previewType, nil
}

// guessPreviewType 根据 MIME 类型返回前端预览分类
func guessPreviewType(mime string) string {
    switch {
    case mime == "":
        return "other"
    case len(mime) >= 5 && mime[:5] == "image": // image/jpeg, image/png 等
        return "image"
    case len(mime) >= 5 && mime[:5] == "video": // video/mp4 等
        return "video"
    case len(mime) >= 5 && mime[:5] == "audio": // audio/mpeg 等
        return "audio"
    case mime == "application/pdf":
        return "pdf"
    case len(mime) >= 4 && mime[:4] == "text": // text/plain 等
        return "text"
    default:
        return "other"
    }
}

biz 层 GetFileStream + service 层 StreamPreview

// GetFileStream 获取文件流用于预览,返回文件对象、读取器和 MIME 类型
func (uc *FileUsecase) GetFileStream(ctx context.Context, userID uint64, fileID uint64) (*File, io.ReadCloser, string, error) {
    file, err := uc.fileRepo.FindByID(ctx, fileID)
    if err != nil {
        return nil, nil, "", err
    }
    if file.UserID != userID {
        return nil, nil, "", ErrForbidden
    }
    // 回收站文件禁止预览
    if file.Status != 0 {
        return nil, nil, "", ErrFileNotFound
    }

    // 从 Storage 获取读取器(流式,不加载全部到内存)
    reader, err := uc.storage.Download(ctx, file.Path)
    if err != nil {
        return nil, nil, "", err
    }

    return file, reader, file.Type, nil
}

// StreamPreview 流式返回文件内容用于预览(service 层)
func (s *FileService) StreamPreview(ctx context.Context, userID uint64, fileID uint64, w http.ResponseWriter) error {
    // 获取文件流
    file, reader, mimeType, err := s.uc.GetFileStream(ctx, userID, fileID)
    if err != nil {
        return err
    }
    defer reader.Close()

    // 设置正确的 Content-Type
    if mimeType != "" {
        w.Header().Set("Content-Type", mimeType)
    } else {
        w.Header().Set("Content-Type", "application/octet-stream")
    }

    // 设置 Content-Length(让浏览器显示下载进度)
    if file.Size > 0 {
        w.Header().Set("Content-Length", strconv.FormatInt(file.Size, 10))
    }

    // 流式传输文件内容:io.Copy 内部按 32KB 缓冲区循环读写,边读边发
    if _, err := io.Copy(w, reader); err != nil {
        return err
    }

    return nil
}

辅助函数 detectMimeType(MIME 推断,含 Docker 镜像兜底):

// detectMimeType 根据文件名推断 MIME 类型
// 优先使用系统 MIME 数据库,缺失时回退到内置常见类型映射
// (Docker 精简镜像可能没有 /etc/mime.types,导致视频等类型识别失败)
func detectMimeType(filename string) string {
    ext := filepath.Ext(filename)
    if ext == "" {
        return "application/octet-stream"
    }
    // 优先使用系统 MIME 数据库
    if mimeType := mime.TypeByExtension(ext); mimeType != "" {
        return mimeType
    }
    // 回退到内置映射(覆盖 Docker 精简镜像缺失 mime.types 的场景)
    if mimeType, ok := builtinMimeTypes[strings.ToLower(ext)]; ok {
        return mimeType
    }
    return "application/octet-stream"
}

5.8 删除与回收站(越权校验 + 存储扣减 + 审计)

实现思路

删除走"软删除到回收站"两步设计:

  1. 移入回收站(Trash):批量校验所有权后,把文件/文件夹标记为 status=1 并写入 recycle_bins 表(记录原父目录,供恢复),同时删除该文件的分享记录。不扣减已用存储(文件物理仍占用空间,只是不可见)。
  2. 永久删除(DeleteForever):从对象存储删除物理文件、把 files 表状态置 status=2、调用 SubUsedStorage 扣减配额、删除元数据缓存,并发布 EventFileDeleted 事件。
  3. 恢复(Restore):把 status 从 1 改回 0,文件回到原目录。

越权防护Trash 在循环里对每个 file/folder 比对 f.UserID != userIDErrForbidden,删除别人的文件会被拒绝。但审计日志缺失:删除/永久删除/恢复目前都没有写审计日志(见三.15),商用版应在这些高风险动作上补充结构化审计。

flowchart TD
    A[Trash 批量删除] --> B{逐个校验 user_id 归属}
    B -- 不属于当前用户 --> F[403 FORBIDDEN]
    B -- 属于 --> C[写 recycle_bins + status=1]
    C --> D[删分享记录 不清存储]
    D --> E[定时/手动 DeleteForever]
    E --> G[删对象存储 + status=2]
    G --> H[SubUsedStorage 扣配额]
    H --> I[发布 EventFileDeleted]
    I -. 修正点:应补审计日志 .-> J[审计:谁/何时/删了什么]

关键代码

biz 层 Trash(越权校验 + 软删除):

func (uc *RecycleUsecase) Trash(ctx context.Context, userID uint64, fileIDs, folderIDs []uint64) ([]uint64, []uint64, error) {
    trashedFiles := make([]uint64, 0, len(fileIDs))
    trashedFolders := make([]uint64, 0, len(folderIDs))

    // 1. 文件:批量查所有权 → 写 RecycleItem → 软删除(status=1)
    if len(fileIDs) > 0 {
        files, err := uc.fileRepo.FindByIDs(ctx, fileIDs)
        if err != nil {
            return nil, nil, err
        }
        now := time.Now()
        expire := now.Add(7 * 24 * time.Hour)
        for _, f := range files {
            if f.UserID != userID { // IDOR 防护:越权拒绝
                return nil, nil, ErrForbidden
            }
            item := &RecycleItem{
                UserID:    userID,
                ItemType:  "file",
                ItemID:    f.ID,
                Name:      f.Name,
                FileSize:  f.Size,
                FileType:  f.Type,
                DeletedAt: now,
                ExpireAt:  expire,
            }
            if _, err := uc.recycleRepo.Create(ctx, item); err != nil {
                return nil, nil, err
            }
            // 删除该文件的所有分享记录(源文件删除,分享也应删除)
            if uc.shareRepo != nil {
                _ = uc.shareRepo.DeleteByItemID(ctx, "file", f.ID)
            }
            trashedFiles = append(trashedFiles, f.ID)
        }
        if len(trashedFiles) > 0 {
            if err := uc.fileRepo.BatchUpdateStatus(ctx, trashedFiles, 1); err != nil {
                return nil, nil, err
            }
        }
    }
    // 2. 文件夹:循环查所有权 → 写 RecycleItem → 软删除内部文件 → 物理删除文件夹记录
    for _, fid := range folderIDs {
        folder, err := uc.folderRepo.FindByID(ctx, fid)
        if err != nil {
            return nil, nil, err
        }
        if folder.UserID != userID { // IDOR 防护
            return nil, nil, ErrForbidden
        }
        // ... 写 RecycleItem、软删除子文件、删文件夹记录、删分享记录
    }

    // 清理文件列表缓存与存储缓存,发布 EventFileTrashed
    if uc.cache != nil {
        _ = uc.cache.DeleteByPattern(ctx, fmt.Sprintf("files:list:%d:*", userID))
        _ = uc.cache.Delete(ctx, fmt.Sprintf(cacheKeyUserStorage, userID))
    }
    return trashedFiles, trashedFolders, nil
}

biz 层 DeleteForever(物理删除 + 配额扣减):

func (uc *RecycleUsecase) DeleteForever(ctx context.Context, userID uint64, recycleID uint64) error {
    item, err := uc.recycleRepo.FindByID(ctx, recycleID)
    if err != nil {
        return err
    }
    if item.UserID != userID { // IDOR 防护
        return ErrForbidden
    }

    deletedFileIDs := uc.permanentlyDeleteItem(ctx, userID, item)

    // 发布文件永久删除事件(失败不影响主流程)
    if uc.eventPublisher != nil && len(deletedFileIDs) > 0 {
        _ = uc.eventPublisher.Publish(ctx, EventFileDeleted, &FileChangedPayload{
            UserID:  userID,
            FileID:  deletedFileIDs[0],
            Action:  "deleted",
        })
    }

    return uc.recycleRepo.Delete(ctx, recycleID)
}

// permanentlyDeleteItem 永久删除:删对象存储 + status=2 + 扣减配额
func (uc *RecycleUsecase) permanentlyDeleteItem(ctx context.Context, userID uint64, item *RecycleItem) []uint64 {
    deletedFileIDs := make([]uint64, 0)
    if item.ItemType == "file" {
        file, err := uc.fileRepo.FindByID(ctx, item.ItemID)
        if err != nil {
            return deletedFileIDs
        }
        if err := uc.storage.Delete(ctx, file.Path); err != nil { // 删物理文件
            log.Error("recycle: failed to delete physical file", "path", file.Path, "err", err)
        }
        _ = uc.fileRepo.BatchUpdateStatus(ctx, []uint64{item.ItemID}, 2)
        if uc.userUC != nil {
            _ = uc.userUC.SubUsedStorage(ctx, userID, file.Size) // 扣减配额
        }
        deletedFileIDs = append(deletedFileIDs, item.ItemID)
    }
    // ... 文件夹递归删内部文件
    return deletedFileIDs
}

自测题与动手练习

自测题(合上书能答出来,才算懂)

  1. 秒传的核心查询是 FindByHash(user_id, hash, status=0)。为什么必须带 user_id?如果不加会有什么越权/共享风险?回收站里(status≠0)的同哈希文件被复用会有什么问题?
  2. 分片上传的会话状态机有 0/1/2/3 四个值。请说出 InitUpload 在哪种情况下会直接返回 status=1,又在哪种情况下会复用已有的进行中/已暂停会话(断点续传)?源码里 model/proto 注释口径与 biz 运行时语义不一致会带来什么坑?
  3. ListFiles 里文件走 offset 分页、文件夹走 keyset 分页,二者的 nextCursor 能直接比较吗?这会造成什么翻页 bug?商用版怎么改?
  4. UploadPartchan struct{} 信号量(容量 50)限流。如果并发远超 50,客户端会收到什么错误?为什么这个全局信号量在"防单用户打爆"上不够(应改成什么)?
  5. MergeParts 在"合并成功但写 files 表失败"“UUID 生成失败"两种情况下分别怎么回滚?为什么必须回滚 Storage 文件(避免孤儿文件)?合并完整性判断用 < 还是 != 更安全?
  6. 商用审查Move/Copy/Download/UploadPart/Trash 各自在哪里做了 user_id 归属校验(IDOR 防护)?Preview 漏了哪一道校验?
  7. 商用审查:并发上传时,CheckStorageAvailable 在分布式锁外、AddUsedStorage 增量无上限,会导致什么后果?写出"检查+扣减"必须在同一行锁内完成的修正 SQL。
  8. 商用审查:入口限流 RateLimitMiddleware(600, time.Minute) 有哪三个商用缺陷?nginx client_max_body_size 0 又意味着什么?单文件流式 Upload 与分片 UploadPart 在大小校验上分别有什么缺口?

动手练习(建议真做一遍)

  1. 在本地起 MySQL + Redis,跑通一次秒传:连传两次相同内容的文件,用 SELECT * FROM files 观察两条记录的 hashpath 是否相同、size 是否只算一次;再试传一个别人的已存在哈希文件,验证是否会被拒绝跨用户复用。
  2. 给一个大文件(如 50MB)做分片上传,上传到一半杀掉客户端,用 CheckParts 接口看已传分片列表,再续传缺失分片并 MergeParts,验证"断点续传"确实只补传了缺的部分。
  3. 故意触发一次"写库失败"路径(如在 Upload 成功后让 fileRepo.Create 返回 error),观察 Storage 里是否还残留孤儿文件——验证 storage.Delete 回滚是否生效。
  4. 写一段并发测试:开 20 个 goroutine 同时上传刚好能让总空间触顶的文件,观察 used_storage 是否超过 total_storage,验证"配额并发超卖"问题并给出带 WHERE used_storage + ? <= total_storage 的修正。
  5. 用 Redis + Lua 把 RateLimitMiddleware 改造成分布式令牌桶,按「用户 ID + 接口」限流(上传 60/分钟、下载 300/分钟),再用 ab/hey 压测验证多实例下限流是否仍然一致。

本章小结

  • 文件模块是标准四层架构(service / biz / data),核心业务逻辑都落在 biz 层,依赖 StorageRepo 抽象接口,存储与数据库可替换。
  • 秒传靠 SHA-256 去重复用物理路径(按 user_id 隔离,零字节传输);断点续传user_id + hash + size 复用 uploadID,配合会话状态机 0/1/2/3CheckParts 跳过已传分片。
  • 分片上传四步(Init→Part→Check→Merge)用信号量限流、原子递增 chunks_received,合并遵循"要么全成、失败即回滚 Storage"的原子性原则。分片参数(chunkSize/分片号/大小)必须由服务端校验,否则可被打爆。
  • 缓存用"第一页缓存 + 抖动 TTL + pattern 主动失效 + 短 TTL 兜底"四件套抗雪崩;孤儿文件回滚零拷贝流式危险文件类型拦截共同保证正确性与安全。但单文件 Upload 此前漏做文件名净化与类型拦截、Preview 漏校验状态,已就地修正。
  • 商用审查要点:① 所有写操作都已校验 user_id 归属(IDOR 防护),但 Preview 漏状态校验;② 配额检查在分布式锁之外 + 增量无上限,并发下可超配额,须把"检查+扣减"放进同一行锁并加 WHERE used_storage + ? <= total_storage;③ 入口限流为进程内、信任 X-Forwarded-For 首值、全局信号量跨用户共享,且 nginx client_max_body_size 0、分片无大小上限,需改分布式限流 + 按用户配额 + 上传体上限;④ 删除为软删除到回收站、越权已拦截,但删/下/享缺少审计日志,应补充。
  • 下一篇可顺着这条线继续深入:把"商用安全加固”(分布式限流、配额防超卖、病毒扫描、审计日志、缩略图缓存)真正落进这套 biz/data 抽象里,你会发现接口已经为它们留好了位置。
About Me

没什么想介绍的,一个很大众的码农…

喜欢代码,车,马,真的是 🐎

讨厌别人让我给自己的代码写注释 最厌烦别人的程序没有写注释

目标

学AI,加油!加油!