学习目标
学完本章你应该能够:
- 讲清文件模块的分层架构(service → biz → data)每一层各自负责什么,以及
Storage/FileRepo/FolderRepo/UploadRepo抽象接口如何做到"可切换存储与数据库实现"。 - 解释秒传(SHA-256 去重)与断点续传的原理:为什么同内容只存一份物理文件、会话状态机
0/1/2/3各代表什么、断点续传靠哪几个查询复用 uploadID。 - 说清分片上传的四步流程与并发控制:信号量(容量 50)限流、原子递增
chunks_received、合并时如何"要么全成要么不留垃圾"。 - 对比 keyset 分页 vs offset 分页、覆盖索引的作用,并讲出"第一页缓存 + 抖动 TTL + 主动失效 pattern"这套缓存策略为什么能抗雪崩。
- 像商用工程师一样审查文件模块:每个写操作是否校验
user_id归属(IDOR 防护)、文件名路径穿越与危险类型拦截是否覆盖所有入口、上传/下载限流与nginx client_max_body_size上限、并发上传下的配额一致性(TOCTOU 与原子上限)、删除/下载/分享是否留审计日志。 - 面试时能围绕"孤儿文件回滚、零拷贝流式、危险文件类型拦截、配额防超卖、亿级流量优化"讲成一段有结构、有取舍的工程故事。
前置知识:
- Go 基础、GORM 基本用法、
context与defer资源管理。 - MySQL 索引(联合索引、覆盖索引)与 Redis 缓存基础。
- Kratos 框架分层(transport / service / biz / data)的大致概念,以及 JWT 鉴权中间件(见上一篇"用户模块")。
本章你会动手做的事:
- 跟读
ListFiles的 keyset/offset 混合分页 + 缓存代码,自己推演"翻到第 3 页时 nextCursor 怎么算",并指出文件/文件夹游标语义不一致引发的翻页 bug。 - 在纸上画出
InitUpload → UploadPart → CheckParts → MergeParts的状态流转,标出哪一步会触发FailSession或回滚 Storage。 - 对照"商用审查"清单(鉴权/输入校验/限流/配额一致性/审计)逐条核对你正在维护的文件服务,挑出 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 分块发送;StreamPreview 用 io.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 删除入回收站]分页查询(ListFiles):用户进入目录后,首先调用
ListFiles获取当前目录下的文件与文件夹。biz 层并行查询fileRepo.ListByParent与folderRepo.ListByParent,混合返回;第一页结果写入 Redis 缓存(带抖动 TTL)。新建文件夹(CreateFolder):用户在当前目录下创建子文件夹,通过
parentID实现多级目录;创建后失效目录缓存。秒传(SecUpload):上传前客户端先计算文件 SHA-256,调用
SecUpload传入哈希。若用户已存在同哈希文件,直接复用存储路径新建一条 file 记录,零字节传输;否则返回exists=false让客户端走正常上传。单文件上传与下载(Upload / Download):小文件走
Upload直接写入 Storage,再写 files 表(失败时回滚 Storage 文件,避免孤儿文件);下载走Download,先查缓存命中文件元数据,再从 Storage 读取字节流返回。service 层还提供 gRPC 客户端流式Upload与服务端流式Download,避免大文件全部加载到内存(零拷贝优化)。大文件分片上传:分四步——
- 初始化(InitUpload):计算分片总数,生成
uploadID,调用 Storage 初始化分片目录,创建 UploadSession 记录;若同哈希同大小会话已存在则直接复用(断点续传)。 - 上传分片(UploadPart):通过信号量限流后写入 Storage 分片目录,再创建 UploadChunk 记录并原子递增
chunks_received。 - 检查分片(CheckParts):返回已上传分片编号列表,供客户端断点续传时跳过已传分片。
- 合并(MergeParts):校验分片完整后调用
CompleteMultipartUpload合并分片,再写 files 表;失败时回滚 Storage 文件;成功后发布EventUploadCompleted事件。
- 初始化(InitUpload):计算分片总数,生成
移动与复制(Move / Copy):批量校验文件所有权(防止 IDOR)后,
BatchMove修改parent_id;复制则保持Path不变(共享物理文件)但生成新 UUID 与记录,并对重名文件自动添加(1)、(2)后缀。两操作完成后均失效源目录与目标目录缓存。预览(Preview):根据文件 MIME 类型分类(image / video / audio / pdf / text / other),通过
Storage.GetFileURL返回预签名 URL(OSS)或本地路径(本地存储);service 层还提供StreamPreview流式返回文件内容,按 64KB 分块写入http.ResponseWriter。删除(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.go与api/file/v1/file.proto的GetUploadStatusReply注释只写了 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)MergePartsUUID 生成失败 →storage.Delete(storagePath);写库失败 →storage.Delete(storagePath)InitUpload创建会话失败 →storage.AbortMultipartUpload(uploadID)
此外 Storage 接口定义了 CleanupOrphanChunks 方法,可扫描 chunks 目录清理早于 maxAge 的孤儿分片目录。
10. 危险文件类型与文件名净化(输入校验)
安全防护两道关卡(在 SecUpload 与 InitUpload 已落地):
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 != userID→ErrForbidden。UploadPart/CheckParts/MergeParts/InitUpload:比对session.UserID != userID→ErrForbidden。Download/GetFileStream/Preview:比对file.UserID != userID→ErrForbidden。- 删除(回收站):
RecycleUsecase.Trash在循环里对每个 file/folder 比对f.UserID != userID→ErrForbidden。
⚠️ 修正点(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} 串行化,再调用 UpdateUsedStorageAtomic 做 used_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 次/分钟),存在三点商用缺陷:
- 非分布式:多实例部署时各算各的,限流值被实例数放大(N 实例 ≈ 600×N/分钟)。
- 信任可伪造的 IP:取
X-Forwarded-For首值作为客户端 IP,攻击者可伪造该头绕过/嫁祸;应使用 nginx 注入的X-Real-IP($remote_addr)。 - 粒度太粗:全局 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.conf 中 client_max_body_size 0; 表示上传体大小不限制,完全依赖后端控制。而后端仅在流式 Upload 路径有 100MB 内存上限(maxStreamUploadSize),分片上传的 UploadPart 没有任何单分片/单文件大小校验,且 InitUpload 的 chunkSize 由客户端传入、未做服务端钳制。这意味着:
- 恶意客户端可传超大盘古级文件撑爆磁盘;
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)目前未写审计日志,建议补充。
四、亿级流量优化思路
针对文件模块在高并发大流量场景的进一步优化:
CDN 加速下载:将热点文件的访问 URL 接入 CDN(如阿里云 CDN、Cloudflare),用户就近从边缘节点拉取,减轻源站带宽压力。
Storage.GetFileURL返回的预签名 URL 可直接配置 CDN 回源。分片上传并行化:当前
UploadPart是串行调用,可在客户端并行上传多个分片(如 3~5 并发),配合服务端信号量限流,显著缩短大文件上传耗时。需保证分片顺序与chunk_index对应,合并时按索引排序。对象存储直传(Presigned PUT):服务端只生成预签名 URL,客户端直接 PUT 到 OSS/S3,跳过应用服务器中转,节省应用带宽与 CPU。服务端只负责创建会话、记录分片、触发合并回调。
秒传减少带宽:同内容文件零字节传输,已是核心优化。可进一步在用户间做"全局秒传"(跨用户共享同哈希文件,需引用计数与权限隔离)。
分库分表:files 表按
user_id取模分库分表(如 1024 库 × 4 表),解决单表亿级数据下的查询性能瓶颈。FindByHash需走分片键user_id路由,避免广播查询。upload_chunks 表可按upload_session_id分表。热点文件缓存:热门视频 / 图片的下载流可通过 Redis 或本地缓存(如 LRU)缓存文件字节或预签名 URL,减少对 Storage 的回源。文件元数据缓存(
file:meta:{id})已实现,可扩展为多级缓存(本地 + Redis),并加缩略图缓存(图片/视频首帧缩略图按thumb:{fileID}:{size}缓存,避免每次预览重新生成)。读写分离:列表查询走 MySQL 只读从库,写入走主库。
ListFiles、ListByCategory、ListStarred等读多写少的场景可显著降低主库压力。异步事件驱动 + 病毒扫描:上传完成后通过
EventPublisher异步发布事件,下游消费生成缩略图、提取视频元数据、建立搜索索引(ES)、触发病毒扫描(商用必备,扫描完成前文件标记"待检"、下载/分享时拦截)。主流程不阻塞,提升上传接口响应速度。存储空间校准任务:定时任务(每日凌晨)调用
RecalibrateStorage,从活跃文件重新计算实际已用空间,修正并发写入导致的计数偏差,避免长期累计误差。孤儿分片清理:定时任务调用
CleanupOrphanChunks,扫描 chunks 目录删除早于阈值(如 24 小时)的子目录,回收磁盘空间,避免用户中途放弃上传留下的垃圾分片。
五、详细实现流程与代码解析
5.1 文件分页列表(keyset/offset 混合分页 + 覆盖索引 + 文件夹混合查询)
实现思路
ListFiles 是用户进入目录后调用的第一个接口,需要同时返回文件与文件夹,并支持按名称 / 大小 / 创建时间排序与翻页。biz 层并行调用 fileRepo.ListByParent 与 folderRepo.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 大文件分片上传(初始化→上传分片→检查分片→合并)
实现思路
大文件分片上传分四步:
- 初始化(InitUpload):校验文件大小、净化文件名、检查存储空间;若传入 hash 则先查秒传(同哈希已有完成文件则直接返回 status=1);再查断点续传(同 user_id + hash + size 的进行中/已暂停会话);最后生成 uploadID,调用 Storage 初始化分片目录,创建 UploadSession 记录。
- 上传分片(UploadPart):通过信号量(容量 50)限流,写入 Storage 分片目录,创建 UploadChunk 记录,原子递增
chunks_received。 - 检查分片(CheckParts):返回已上传分片编号列表,供客户端断点续传跳过已传分片。
- 合并(MergeParts):校验分片数等于
ChunksTotal,调用CompleteMultipartUpload合并,写 files 表,发布上传完成事件。
类比:分片上传像一个快递分拣中心。包裹(大文件)被拆成很多小箱(分片),每个小箱独立发往中转站(Storage 分片目录),每到一个就在台账上画一笔
chunks_received +1;等所有箱到齐,再封一个大包(合并)送出。中途某个箱丢了,不用重发全部——查台账(CheckParts)只看缺哪几箱补哪几箱,这就是断点续传。
会话状态机有四个值:0=进行中 / 1=已完成 / 2=失败 / 3=已暂停(注意 model/proto 注释口径需统一,见三.2)。
⚠️ 修正点(分片参数未校验,可被打爆):
InitUpload的chunkSize由客户端传入却未钳制,恶意值(如chunkSize=1)会让ChunksTotal = fileSize,产生天文数字的分片数。服务端应钳制到合理范围(如 1MB~100MB)。MergeParts用len(chunks) < ChunksTotal判断完整性,<允许"多传"的分片通过;应改为!=并校验分片号连续且落在[1, ChunksTotal]、单分片大小 ≤ 上限,合并后文件总大小与声明FileSize一致。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 类型分类)
实现思路
预览分两种模式:
- URL 预览(Preview):根据文件 MIME 类型分类(image / video / audio / pdf / text / other),调用
Storage.GetFileURL返回预签名 URL(OSS)或本地路径(本地存储),有效期 1 小时。前端拿到 URL 后直接在浏览器渲染。 - 流式预览(StreamPreview):service 层通过
GetFileStream获取 Storage reader,设置Content-Type与Content-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 删除与回收站(越权校验 + 存储扣减 + 审计)
实现思路
删除走"软删除到回收站"两步设计:
- 移入回收站(Trash):批量校验所有权后,把文件/文件夹标记为
status=1并写入recycle_bins表(记录原父目录,供恢复),同时删除该文件的分享记录。不扣减已用存储(文件物理仍占用空间,只是不可见)。 - 永久删除(DeleteForever):从对象存储删除物理文件、把 files 表状态置
status=2、调用SubUsedStorage扣减配额、删除元数据缓存,并发布EventFileDeleted事件。 - 恢复(Restore):把
status从 1 改回 0,文件回到原目录。
越权防护:
Trash在循环里对每个 file/folder 比对f.UserID != userID→ErrForbidden,删除别人的文件会被拒绝。但审计日志缺失:删除/永久删除/恢复目前都没有写审计日志(见三.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
}
自测题与动手练习
自测题(合上书能答出来,才算懂):
- 秒传的核心查询是
FindByHash(user_id, hash, status=0)。为什么必须带user_id?如果不加会有什么越权/共享风险?回收站里(status≠0)的同哈希文件被复用会有什么问题? - 分片上传的会话状态机有
0/1/2/3四个值。请说出InitUpload在哪种情况下会直接返回status=1,又在哪种情况下会复用已有的进行中/已暂停会话(断点续传)?源码里 model/proto 注释口径与 biz 运行时语义不一致会带来什么坑? ListFiles里文件走 offset 分页、文件夹走 keyset 分页,二者的nextCursor能直接比较吗?这会造成什么翻页 bug?商用版怎么改?UploadPart用chan struct{}信号量(容量 50)限流。如果并发远超 50,客户端会收到什么错误?为什么这个全局信号量在"防单用户打爆"上不够(应改成什么)?MergeParts在"合并成功但写 files 表失败"“UUID 生成失败"两种情况下分别怎么回滚?为什么必须回滚 Storage 文件(避免孤儿文件)?合并完整性判断用<还是!=更安全?- 商用审查:
Move/Copy/Download/UploadPart/Trash各自在哪里做了user_id归属校验(IDOR 防护)?Preview漏了哪一道校验? - 商用审查:并发上传时,
CheckStorageAvailable在分布式锁外、AddUsedStorage增量无上限,会导致什么后果?写出"检查+扣减"必须在同一行锁内完成的修正 SQL。 - 商用审查:入口限流
RateLimitMiddleware(600, time.Minute)有哪三个商用缺陷?nginxclient_max_body_size 0又意味着什么?单文件流式Upload与分片UploadPart在大小校验上分别有什么缺口?
动手练习(建议真做一遍):
- 在本地起 MySQL + Redis,跑通一次秒传:连传两次相同内容的文件,用
SELECT * FROM files观察两条记录的hash与path是否相同、size是否只算一次;再试传一个别人的已存在哈希文件,验证是否会被拒绝跨用户复用。 - 给一个大文件(如 50MB)做分片上传,上传到一半杀掉客户端,用
CheckParts接口看已传分片列表,再续传缺失分片并MergeParts,验证"断点续传"确实只补传了缺的部分。 - 故意触发一次"写库失败"路径(如在
Upload成功后让fileRepo.Create返回 error),观察 Storage 里是否还残留孤儿文件——验证storage.Delete回滚是否生效。 - 写一段并发测试:开 20 个 goroutine 同时上传刚好能让总空间触顶的文件,观察
used_storage是否超过total_storage,验证"配额并发超卖"问题并给出带WHERE used_storage + ? <= total_storage的修正。 - 用 Redis + Lua 把
RateLimitMiddleware改造成分布式令牌桶,按「用户 ID + 接口」限流(上传 60/分钟、下载 300/分钟),再用 ab/hey 压测验证多实例下限流是否仍然一致。
本章小结
- 文件模块是标准四层架构(service / biz / data),核心业务逻辑都落在 biz 层,依赖
Storage与Repo抽象接口,存储与数据库可替换。 - 秒传靠 SHA-256 去重复用物理路径(按
user_id隔离,零字节传输);断点续传靠user_id + hash + size复用 uploadID,配合会话状态机0/1/2/3与CheckParts跳过已传分片。 - 分片上传四步(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首值、全局信号量跨用户共享,且 nginxclient_max_body_size 0、分片无大小上限,需改分布式限流 + 按用户配额 + 上传体上限;④ 删除为软删除到回收站、越权已拦截,但删/下/享缺少审计日志,应补充。 - 下一篇可顺着这条线继续深入:把"商用安全加固”(分布式限流、配额防超卖、病毒扫描、审计日志、缩略图缓存)真正落进这套 biz/data 抽象里,你会发现接口已经为它们留好了位置。