回收站模块实现教程

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

@

学习目标

学完本章你应该能够:

  1. 讲清回收站"先软删除、后物理删除"的核心设计,以及文件表 status 三态(0 正常 / 1 回收站 / 2 永久删除)各自语义与触发动作。
  2. 把"删除"拆成六步(移入回收站 / 列表 / 还原 / 永久删除 / 清空 / 过期清理),说清 service → biz → data 各层职责边界。
  3. 解释本项目配额一致性语义:为什么"移入回收站不扣减 used_storage、还原不回加、只有物理删除才原子扣减",并据此论证"还原导致超配额"在本设计中根本不会发生(对比"删除即扣减"方案的冲突)。
  4. 讲清权限越权防护:恢复/永久删除/清空/清理是否校验 user_id 归属,越权访问他人回收站条目的后果。
  5. 说清合规删除(个保法删除权):永久删除是否真从对象存储删除、异步清理的可靠性、孤儿文件与配额漂移如何兜底。
  6. 解释 keyset 游标分页为什么比 OFFSET 更适合回收站列表,并知道"非唯一游标"在分页边界的坑与修复。
  7. 站在商用视角,讲清同一条目"还原/永久删除/定时清理"并发、以及多副本下定时任务重复执行带来的配额重复扣减问题,并给出"条目级锁 + 幂等"的修复方案。
  8. 确认回收站保留策略(默认 7 天)与自动清理真实由 cron 调度器执行main.go 注册 recycle-clean 任务),而非仅存在于文档。

前置知识

  • Kratos 微服务框架基础(transport / middleware / errors / 依赖注入)
  • GORM 基本用法(增删改查、事务、WithContext
  • Redis 缓存与分布式锁基本概念
  • MySQL 事务、行锁与索引基础
  • 已了解用户模块的 UpdateUsedStorageAtomic 原子扣减能力(见「用户模块」一节)

本章你会动手做的事

  1. 在纸上画出文件 status 三态状态机,标出每个触发动作(Trash / Restore / DeleteForever / EmptyRecycle / CleanExpired)。
  2. 用一段伪代码实现 keyset 分页,体会 deleted_at < cursor + Limit(limit+1) 如何取代 OFFSET,并思考"两条记录同一秒"时游标会怎样重/漏。
  3. permanentlyDeleteItem 加一行日志,观察"磁盘删除失败但配额仍被扣减"的链路,并据此设计"条目级锁 + 二次存在校验"的幂等修复。

一、技术栈与中间件

回收站模块采用 Kratos 微服务框架,使用 DDD 分层结构(service → biz → data),并通过 GORM 操作 MySQL,通过 Redis 缓存加速列表查询,通过对象存储抽象做物理删除,通过 cron 定时任务自动清理过期记录。

白话类比:整个模块像一家"文件寄存处"。用户说"扔了",前台(service 层)只管登记和回话;真正决定怎么处理的是主管(biz 层用例);搬箱子、查账本的是仓库工人(data 层仓储)。这样分层后,哪天你想换仓储系统(MySQL 换 TiDB)或换磁盘(本地换 OSS),只要仓库工人换套干法,前台的流程完全不用动。

技术 / 中间件所属层用途说明
Go Kratos v3框架骨架提供 transport(gRPC/HTTP)、middleware、errors、log 等基础能力,统一 service 层骨架
GORM数据访问层封装回收站表 recycle_bins 与文件表 files 的增删改查,支持 WithContext、批量更新、条件原子更新
MySQL持久化存储存回收站记录(recycle_bins)、文件元数据(files.status:0 正常 / 1 回收站 / 2 永久删除)、文件夹记录(folders
Redis 缓存(Cache 接口)缓存层清除用户文件列表缓存(files:list:{userID}:*)、文件元数据缓存(file:meta:{id})、存储缓存(cloud-disk:user:storage:{id});商用版还承载分布式锁
对象存储抽象(Storage 接口)存储层抽象底层磁盘(本地 / OSS / S3),永久删除时调用 storage.Delete(path) 物理删除文件,满足合规删除权
Locker 接口(biz 层)并发控制抽象分布式锁(由 RecycleUsecase 持有),用于串行化"同一条目并发操作"与"多副本定时清理",防止配额重复扣减
EventPublisher 事件发布器解耦异步发布 EventFileTrashed / EventFileDeleted,供审计、搜索索引、回收站清理事件等下游消费
UserUsecase用户用例永久删除后原子扣减用户已用存储(SubUsedStorage → 分布式锁 + UpdateUsedStorageAtomic
FileRepo / FolderRepo / ShareRepo仓储接口Trash / Restore / DeleteForever 时联动修改文件状态、重建文件夹、删除分享
cron 调度器(internal/data/scheduler)任务层真实注册 recycle-clean0 0 3 * * *)调用 CleanExpired,以及 storage-calibrate / chunk-cleanup
Keyset 分页(cursor)列表查询回收站列表用基于 deleted_at 的游标分页,避免深度分页性能问题
protobuf(api/recycle/v1)接口契约定义 gRPC 服务接口与消息体,service 层实现 RecycleServiceServer

分层与关键中间件的协作关系:

flowchart TD
    CLI[客户端 gRPC/HTTP] --> SVC[service 层
RecycleServiceServer] SVC --> BIZ[业务逻辑层
RecycleUsecase] BIZ --> FR[FileRepo / FolderRepo / ShareRepo] BIZ --> RC[RecycleRepo] BIZ --> LK[(修正)Locker 分布式锁] BIZ --> EV[EventPublisher 事件] BIZ --> CACHE[(Redis 缓存)] BIZ --> STO[Storage 磁盘抽象] BIZ --> US[UserUsecase
原子扣减配额] FR --> DB[(MySQL)] RC --> DB EV --> DOWN[下游: 审计 / 搜索索引] STO --> DISK[(本地 / OSS / S3)] BIZ -.->|CleanExpired| CRON[cron 调度器
recycle-clean 0 0 3 * * *]

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

回收站模块围绕"软删除"这一核心思想展开,把"删除"拆成两步:先把文件移入回收站(保留磁盘文件,可还原),用户确认或过期后再永久删除(清理磁盘文件)。

白话类比:回收站就像办公室的"碎纸机暂存篮"。你把文件丢进去,它只是被挪到篮子里(软删除),真要销毁得等你自己清空或到期自动清。好处是手滑丢错了还能捡回来;代价是篮子里的文件仍占着柜子空间(配额照算),得靠定期清空来回收。

  1. 软删除(移入回收站):用户点击"删除"时,不直接物理删除文件,而是把文件的 status 改为 1(trashed),同时往 recycle_bins 表写入一条记录(含名称、大小、类型、原路径、过期时间),并删除其分享记录、清除文件列表缓存。不扣减 used_storage(文件仍在占用空间)。
  2. 回收站列表查询:使用 keyset 分页(基于 deleted_at 游标),按 deleted_at DESC 排序,Limit(limit+1) 判断 hasMore。返回时携带名称、大小、类型等冗余字段,前端无需回查文件表。
  3. 还原文件:从回收站记录找到原 ItemID,把文件 status1 改回 0(正常);文件夹则重建 folders 表记录(用原 ID 保证内部文件的 parent_id 仍然有效),并把内部所有文件 status 改回 0。最后删除回收站记录、清除缓存。不回加配额(因为从未扣减)。
  4. 永久删除storage.Delete(path) 物理删除磁盘文件 → BatchUpdateStatus([id], 2) 改状态 → SubUsedStorage 原子扣减配额。三步非原子,需容错与幂等保护(见 5.4)。
  5. 清空回收站:一次性拉取该用户全部回收站记录(上限需分批),循环调用永久删除逻辑,逐条清理 recycle_bins 记录,最后聚合发布永久删除事件。
  6. 过期自动清理:cron 任务(recycle-clean,每日 3:00)调用 CleanExpired,扫描 expire_at < now() 的记录并执行永久删除。注意多副本下需全局锁避免重复执行(见 5.6)。

这六步串起来的主流程:

flowchart LR
    A[用户点删除] --> B["软删除
status:0→1
写 recycle_bins
不扣配额"] B --> C["回收站列表
keyset 分页"] C --> D{用户操作} D -->|还原| E["status:1→0
删 recycle 记录
不回加配额"] D -->|永久删除| F["物理删磁盘
status:1→2
原子扣配额"] D -->|清空| G[循环逐条永久删除] B -. "cron 每日3点" .-> F F --> H["used_storage 原子 -size"]

配额语义一张图说清:整个生命周期里,文件只要还占着磁盘,就一直计入 used_storage。只有走到最右端(物理删除)才扣减。所以"还原"只是把状态 1↔0 来回切,配额数字根本不动——这正是本项目避免"还原后超配额冲突"的关键设计。

flowchart LR
    S0[正常 status=0
配额 +size] --> S1[回收站 status=1
配额 仍 +size] S1 -->|还原| S0 S1 -->|永久删除| S2[永久删除 status=2
配额 -size
磁盘已删]

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

1. 软删除 vs 物理删除设计

回收站核心是"先软删除、后物理删除"。文件表用 status 字段区分状态:0 正常、1 已移入回收站、2 已永久删除。软删除保留磁盘文件,便于用户反悔;物理删除才真正释放存储空间。这样既保证用户体验,又能在过期后自动回收资源。

易错点recycle_bins 表的 DeletedAt 字段类型是普通 time.Timegorm:"autoCreateTime"),不是 gorm.DeletedAt,因此它只是"放入回收站的时间戳",并不会触发 GORM 的软删除自动过滤。recycleRepo.Delete 是真正的硬删除,不会留下幽灵记录。这是有意为之——回收站记录本身就是"已删除"的载体,若再叠加 GORM 软删除反而会把 ListExpired/ListByUser 搞乱。

2. 配额一致性语义(还原超配额冲突为何不会发生)

本项目采用**“移入回收站不扣减配额”**的语义:

  • Trash不调用 SubUsedStorage(文件仍在磁盘,仍占用用户空间,配额应继续计入)。
  • Restore不调用 AddUsedStorage(配额从未减少,无需回加)。
  • DeleteForever / EmptyRecycle / CleanExpired:调用 SubUsedStorage 原子扣减。

对比方案:若采用"删除即扣减、还原即回加",那么用户在文件进回收站期间又上传了大量文件把配额用满,还原时就必须重新 CheckStorageAvailable,一旦超配额还原会失败,需要把文件"卡"在回收站或拒绝还原——这就是"恢复超配额冲突"。本项目刻意选择"删除不扣减",用"磁盘占用 ≡ 配额占用"的等价关系从根本上消灭了这个冲突,代价是回收站里的文件也会吃用户的可用空间(与 Dropbox 等行为一致,符合用户直觉)。

配额扣减本身是安全的:SubUsedStorageupdateStorageWithLock 先取 user:storage🔒{userID} 分布式锁,再走 UpdateUsedStorageAtomicWHERE used_storage >= ? + gorm.Expr("used_storage + ?")),即使并发也不会扣成负数。

3. 权限越权防护

恢复 / 永久删除 / 清空 / 定时清理都必须校验归属,否则 A 用户能删除或还原 B 用户的文件:

  • Trash:逐文件 / 逐文件夹校验 f.UserID != userID / folder.UserID != userIDErrForbidden
  • RestoreDeleteForever:先 recycleRepo.FindByID 取出条目,校验 item.UserID != userIDErrForbidden
  • EmptyRecycle:作用域直接绑定 userID(只查 WHERE user_id = ?),天然隔离。
  • CleanExpired:按 item.UserID 聚合扣减,条目只属于其本人。

⚠️ 仍存在的越权边界(修正点)Trash 在批量循环里遇到"非本人文件"会直接返回 ErrForbidden 并中断整个批次,导致前面已 BatchUpdateStatus 的文件处于"半移入回收站"的中间态。商用版应先对整批做完整所有权校验,或把越权项单独剔除、返回"部分成功",避免半截状态。

4. 合规删除(个保法删除权)

《个人信息保护法》赋予用户"删除权"——用户要求删除的数据必须从存储中真正清除。permanentlyDeleteItem 第一步就是 storage.Delete(path)确实从对象存储删除了物理文件(不只是改状态),满足合规删除。

但异步可靠性要当心:

  • storage.Delete 失败只记日志、不中断后续(OSS 的 Delete 本身幂等,删不存在的 key 也返回成功,因此重复删盘无害)。
  • 若磁盘删除成功但 BatchUpdateStatusSubUsedStorage 失败,会产生"磁盘已删但库未更/配额未扣"的不一致。本项目用每日 storage-calibrate 定时任务RecalibrateStorage 按活跃文件实际大小重算 used_storage)做兜底对账,防止配额长期漂移。
  • 关键不变量:配额准确性优先于磁盘无孤儿。因此扣减顺序应为"删盘(尽力)→ 改库 → 扣配额",即使删盘失败也应继续扣减(磁盘孤儿由孤儿清理任务兜底),而不能"磁盘删失败就不扣配额",否则用户永远占着已不存在文件的额度。

5. 保留策略与真实调度器

每条回收站记录写入时 ExpireAt = now + 7 天。自动清理不是文档里的设想,而是真实存在的cmd/server/main.go 注册了 newRecycleCleanTask,cron 表达式为 0 0 3 * * *(每日 3:00),其函数体直接调用 recycleUC.CleanExpired(ctx);同文件还注册了 storage-calibrate(每日 3:00)与 chunk-cleanup(每日 4:00)。调度器基于 robfig/cron/v3,通过 ScheduledTaskServer 接入 Kratos 生命周期,进程启动即 StartAll、退出即 StopAll

flowchart TD
    A[进程启动 StartAll] --> B[cron 调度器]
    B -->|每日 3:00| C[recycle-clean
CleanExpired] B -->|每日 3:00| D[storage-calibrate
RecalibrateStorage] B -->|每日 4:00| E[chunk-cleanup
清理孤儿分片] C --> F[扫描 expire_at < now 的记录] F --> G[逐条 物理删盘+改库+扣配额]

6. 并发安全(同条目并发与多副本重复清理)

这是本模块最大的商用隐患,原实现未处理:

  • 场景一(用户操作 vs 定时清理):用户点 DeleteForever,同时 CleanExpired 也扫到同一条过期记录。两者都调用 permanentlyDeleteItem → 各扣一次 SubUsedStorage配额被重复扣减(少算用户空间,等于白送容量)。
  • 场景二(还原 vs 永久删除):用户同时点"还原"和"永久删除"同一条目。DeleteForeverstorage.Delete 把磁盘文件删了,Restore 随后把 status 改回 0 并删掉回收站记录 → 文件在界面上"复活"了,但磁盘内容已被删,变成一个打不开的损坏文件。
  • 场景三(多副本):K8s 起 N 个副本,每个副本各自跑 recycle-clean → N 倍重复清理,叠加场景一的配额双扣。

修复方案(修正点):① 给 RecycleUsecase 注入 Locker,对所有"针对同一条目"的写操作加条目级锁 recycle:item🔒{recycleID},并在加锁后二次 FindByID 确认条目仍存在——若已被清理任务删掉则直接返回,不再扣配额;② 给 CleanExpired全局锁 recycle:clean:lock(带 TTL),保证多副本下只有一个节点执行全量清理(条目级锁已能防双扣,全局锁进一步避免无意义的重复扫描)。

flowchart LR
    U[用户 DeleteForever] --> LU["加锁 recycle:item🔒{id}"]
    CRON[CleanExpired 扫到同条目] --> LC[尝试加同一把锁]
    LC -->|获取失败 等待/跳过| LU
    LU --> CK{"二次 FindByID
条目还在?"} CK -->|已无 被对方删| SKIP[跳过 不重复扣配额] CK -->|还在| DEL[物理删盘+改库+扣配额] DEL --> UL[释放锁]

7. Keyset 分页与"非唯一游标"边界

列表用 WHERE deleted_at < ? ORDER BY deleted_at DESC LIMIT n+1。相比 OFFSET 分页,深度翻页无需扫描跳过的行,性能稳定;缺点是只能顺序翻页。但 deleted_at秒级时间戳,同一秒可能有多条记录,游标落在"这一秒"的边界时会出现跨页重读或漏读

修复(修正点):游标改为 (deleted_at, id) 复合排序,索引建 (user_id, deleted_at, id),查询条件写成 deleted_at < ? OR (deleted_at = ? AND id < ?),游标携带 deleted_at|id,彻底消除边界歧义。

8. 缓存与数据库一致性

删除 / 还原 / 永久删除后主动清缓存:清 files:list:{userID}:*(文件列表)、file:meta:{id}(元数据)、cloud-disk:user:storage:{id}(存储缓存)。采用"删缓存"而非"更新缓存",避免并发写脏数据,代价是下一次请求需重建缓存。同时 Download/Preview/GetFileStream 都校验 file.Status != 0 直接返回 ErrFileNotFound,保证回收站里的文件无法被下载或预览,既一致又安全。


四、亿级流量优化思路

在亿级文件、千万级用户的高并发场景下,回收站模块需要从以下几个方向优化:

  • 分批删除避免大事务EmptyRecycle 当前一次性 ListByUser(..., 100000) 拉全量再循环,超大用户会撑爆内存与事务。优化为基于游标的分批循环(每批 100~500 条)提交,并配合 BatchUpdateStatus 批量更新,避免长事务锁表。
  • 异步删除磁盘文件storage.Delete 是 IO 密集型,同步调用拖慢接口。优化为标记 status=2(在事务内先落库)+ 发布消息到队列(Kafka/NSQ),由 worker 消费执行物理删除,主流程立即返回;配额扣减随落库一起完成,保证计费准确。
  • 过期扫描分批 + 游标ListExpired 当前 Find 全表无 Limit,过期记录海量时会一次性载入内存。优化为按 expire_at 游标分批(每批 1000~5000),并务必在 expire_at 上建索引;还可按 expire_at 做按天分区表,定时任务只扫当天分区。
  • 定时任务全局锁与分片CleanExpired 单机执行会形成"清扫尖刺"且多副本重复。加 recycle:clean:lock 全局锁选主执行;超大集群再按 userID % N 分片,多节点并行清理、平滑负载。
  • 回收站列表缓存:当前 ListRecycle 不缓存结果,高 QPS 下数据库压力大。可对 (userID, cursor) 维度加 Redis 缓存(TTL 1~5 分钟),删除 / 还原时按用户级失效(复用 invalidateFileListCache 思路)。
  • 读写分离:列表查询、过期扫描走从库,删除 / 还原等写操作走主库。回收站是"读多写少"场景,读写分离能把读压力分摊到多个从库。
  • 冷热分离:30 天内的回收站记录在热表(SSD),更早的归档到冷表(HDD / 对象存储),前端默认只查热表,长期记录后台归档。
  • 重操作限流与降级:对 EmptyRecycleuserID 限流(如每分钟 1 次)防止恶意调用打挂数据库;高负载时降级为异步任务,返回"正在处理中"。

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

本节按子功能拆解,给出实现思路与关键代码,并按商用在线服务标准修正或补全(修正处标注「修正点」)。代码均基于源文件 internal/biz/recycle.gointernal/data/recycle.gointernal/service/recycle.gocmd/server/main.go

5.1 软删除(文件 / 文件夹移入回收站,File.Status = 1)

实现思路

软删除由 RecycleUsecase.Trash 完成,支持文件和文件夹两类:

  • 文件:批量校验所有权 → 写 RecycleItem(含名称、大小、类型、过期时间 7 天)→ 删除该文件分享记录 → BatchUpdateStatus(fileIDs, 1) 批量改状态。不扣减配额
  • 文件夹:循环校验所有权 → 写 RecycleItemOriginalPath 存父目录 ID 供恢复)→ 软删除该文件夹下所有文件(status=0→1)→ 递归软删除子目录下文件 → 物理删除子文件夹与顶层文件夹记录 → 删除分享记录。

完成后统一清文件列表缓存、文件元数据缓存、用户存储缓存,并发布 EventFileTrashed 事件。

关键代码(修正点:先全量校验归属,避免半截状态)

// Trash 将文件/文件夹移入回收站(软删除)。
// 文件:状态改为 trashed(1) + 写入 recycle_bins,可恢复,不扣减配额。
// 文件夹:写 recycle_bins(记父目录 ID)→ 软删除内部文件 → 物理删除 folder 记录。
// 不扣减用户已用存储空间(文件仍在占用空间)。
func (uc *RecycleUsecase) Trash(ctx context.Context, userID uint64, fileIDs, folderIDs []uint64) ([]uint64, []uint64, error) {
	// 修正点(商用):先全量校验文件所有权,任一越权直接整体失败,避免"半移入回收站"中间态
	if len(fileIDs) > 0 {
		files, err := uc.fileRepo.FindByIDs(ctx, fileIDs)
		if err != nil {
			return nil, nil, err
		}
		owned := make(map[uint64]*File, len(files))
		for _, f := range files {
			if f.UserID != userID {
				return nil, nil, ErrForbidden // 越权:禁止把他人文件移入自己回收站
			}
			owned[f.ID] = f
		}
		now := time.Now()
		expire := now.Add(7 * 24 * time.Hour) // 保留 7 天
		for _, f := range files {
			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) // 源文件删除,分享也应删除
			}
		}
		ids := make([]uint64, 0, len(files))
		for _, f := range files {
			ids = append(ids, f.ID)
		}
		if err := uc.fileRepo.BatchUpdateStatus(ctx, ids, 1); err != nil {
			return nil, nil, err
		}
	}

	// 文件夹分支:逐文件夹校验所有权 → 写回收站 → 软删除内部文件 → 物理删除 folder 记录
	for _, fid := range folderIDs {
		folder, err := uc.folderRepo.FindByID(ctx, fid)
		if err != nil {
			return nil, nil, err
		}
		if folder.UserID != userID {
			return nil, nil, ErrForbidden
		}
		originalPath := "0"
		if folder.ParentID != nil {
			originalPath = fmt.Sprintf("%d", *folder.ParentID)
		}
		now := time.Now()
		item := &RecycleItem{
			UserID: userID, ItemType: "folder", ItemID: fid,
			OriginalPath: originalPath, Name: folder.Name, FileSize: 0,
			FileType: "folder", DeletedAt: now, ExpireAt: now.Add(7 * 24 * time.Hour),
		}
		if _, err := uc.recycleRepo.Create(ctx, item); err != nil {
			return nil, nil, err
		}
		// 软删除该文件夹下所有活跃文件(status=0 → 1)
		activeFiles, _, _ := uc.fileRepo.ListByParent(ctx, userID, &fid, 0, "", 100000, "", "")
		if len(activeFiles) > 0 {
			ids := make([]uint64, len(activeFiles))
			for i, f := range activeFiles {
				ids[i] = f.ID
			}
			_ = uc.fileRepo.BatchUpdateStatus(ctx, ids, 1)
		}
		// 递归软删除子目录下文件 + 物理删除子文件夹记录
		subDirIDs, _ := uc.folderRepo.ListAllSubDirectoryIDs(ctx, userID, fid)
		for _, subID := range subDirIDs {
			subFiles, _, _ := uc.fileRepo.ListByParent(ctx, userID, &subID, 0, "", 100000, "", "")
			if len(subFiles) > 0 {
				ids := make([]uint64, len(subFiles))
				for i, f := range subFiles {
					ids[i] = f.ID
				}
				_ = uc.fileRepo.BatchUpdateStatus(ctx, ids, 1)
			}
			_ = uc.folderRepo.Delete(ctx, subID)
		}
		if err := uc.folderRepo.Delete(ctx, fid); err != nil {
			return nil, nil, err
		}
		if uc.shareRepo != nil {
			_ = uc.shareRepo.DeleteByItemID(ctx, "folder", fid)
		}
	}

	// 清理文件列表缓存、文件元数据缓存、存储缓存,保证前端刷新一致
	if uc.cache != nil {
		_ = uc.cache.DeleteByPattern(ctx, fmt.Sprintf("files:list:%d:*", userID))
		_ = uc.cache.Delete(ctx, fmt.Sprintf(cacheKeyUserStorage, userID))
	}
	if uc.eventPublisher != nil {
		_ = uc.eventPublisher.Publish(ctx, EventFileTrashed, &FileChangedPayload{
			UserID: userID, Action: "trashed",
		})
	}
	return fileIDs, folderIDs, nil
}

5.2 回收站列表查询(keyset 分页 + 冗余字段)

实现思路

  • service 层 ListRecycle:从 ctx 取 userID,调用 usecase 把 RecycleItem 转 proto(DeletedAt/ExpireAt 格式化为字符串)。
  • biz 层 ListRecycle:校验 pageSize(默认 20),调 repo 拿 items 与 nextCursor,按 len(items) >= pageSizehasMore
  • data 层 ListByUser:按 user_id 过滤;cursor != "" 追加 deleted_at < cursorORDER BY deleted_at DESC LIMIT(limit+1);返回项已携带 Name/FileSize/FileType,前端无需回查文件表。

keyset 分页核心(Limit(limit+1)hasMore):

flowchart TD
    A[请求 cursor + pageSize] --> B{cursor 为空?}
    B -->|是| C[取最新 deleted_at DESC 前 limit+1 条]
    B -->|否| D[取 deleted_at < cursor 前 limit+1 条]
    C --> E{实际返回 > limit?}
    D --> E
    E -->|是| F[截断多余 1 条
hasMore=true
cursor=最后一条 deleted_at] E -->|否| G[hasMore=false]

关键代码(修正点:复合游标消除同秒边界重/漏)

// biz 层
func (uc *RecycleUsecase) ListRecycle(ctx context.Context, userID uint64, cursor string, pageSize int) ([]*RecycleItem, string, bool, error) {
	if pageSize <= 0 {
		pageSize = 20
	}
	items, nextCursor, err := uc.recycleRepo.ListByUser(ctx, userID, cursor, pageSize)
	if err != nil {
		return nil, "", false, err
	}
	hasMore := len(items) >= pageSize
	return items, nextCursor, hasMore, nil
}

// data 层(keyset 分页核心,修正点:改用 (deleted_at, id) 复合游标避免同秒歧义)
func (r *recycleRepo) ListByUser(ctx context.Context, userID uint64, cursor string, limit int) ([]*biz.RecycleItem, string, error) {
	query := r.db.WithContext(ctx).Model(&model.RecycleBin{}).Where("user_id = ?", userID)
	// 修正点:cursor 形如 "2025-08-01 12:00:00|123",拆成 deleted_at 与 id 两段
	if cursor != "" {
		parts := strings.SplitN(cursor, "|", 2)
		if len(parts) == 2 {
			query = query.Where("(deleted_at < ?) OR (deleted_at = ? AND id < ?)", parts[0], parts[0], parts[1])
		} else {
			query = query.Where("deleted_at < ?", cursor)
		}
	}
	var pos []model.RecycleBin
	// 修正点:排序增加 id 作为同秒 tiebreaker,配合 (user_id, deleted_at, id) 索引
	if err := query.Order("deleted_at DESC, id DESC").Limit(limit + 1).Find(&pos).Error; err != nil {
		return nil, "", err
	}
	hasMore := len(pos) > limit
	if hasMore {
		pos = pos[:limit]
	}
	nextCursor := ""
	if len(pos) > 0 {
		last := pos[len(pos)-1]
		nextCursor = last.DeletedAt.Format("2006-01-02 15:04:05") + "|" + fmt.Sprintf("%d", last.ID)
	}
	items := make([]*biz.RecycleItem, len(pos))
	for i := range pos {
		items[i] = toBizRecycleItem(&pos[i])
	}
	return items, nextCursor, nil
}

5.3 还原文件(恢复到原位置)

实现思路

RestoreItemType 分两条路径:

  • 文件FindByID 取原文件 → status=0Update → 取 ParentID 供前端跳转 → 清文件列表缓存与元数据缓存。
  • 文件夹:从 OriginalPath 解析父目录 ID → 用原 ID 重建 Folder(保证内部文件 parent_id 仍有效)→ 拉取该文件夹下 status=1 的文件批量改回 0 → 清缓存。

最后删回收站记录,返回父目录 ID(0 表示根目录)。不回加配额(从未扣减)。

关键代码(修正点:文件夹同名冲突处理)

// Restore 将文件从回收站恢复到原始位置,返回父目录 ID(0=根目录)。
// 修正点:增加条目级锁 + 存在校验,防止与 DeleteForever/CleanExpired 并发。
func (uc *RecycleUsecase) Restore(ctx context.Context, userID uint64, recycleID uint64) (uint64, error) {
	item, err := uc.recycleRepo.FindByID(ctx, recycleID)
	if err != nil {
		return 0, err
	}
	if item.UserID != userID {
		return 0, ErrForbidden // 越权校验
	}
	// 修正点:条目级锁,串行化同一条目的还原 / 永久删除 / 清理
	lockKey := fmt.Sprintf("recycle:item🔒%d", recycleID)
	if uc.locker != nil {
		if err := uc.locker.Lock(ctx, lockKey); err != nil {
			return 0, err
		}
		defer uc.locker.Unlock(ctx, lockKey)
	}
	// 修正点:加锁后二次确认条目仍在(可能已被清理任务删除)
	if _, err := uc.recycleRepo.FindByID(ctx, recycleID); err != nil {
		return 0, err
	}

	var parentID uint64
	if item.ItemType == "file" {
		file, err := uc.fileRepo.FindByID(ctx, item.ItemID)
		if err != nil {
			return 0, err
		}
		file.Status = 0
		if err := uc.fileRepo.Update(ctx, file); err != nil {
			return 0, err
		}
		if file.ParentID != nil {
			parentID = *file.ParentID
		}
		uc.invalidateFileListCache(ctx, userID)
		if uc.cache != nil {
			_ = uc.cache.Delete(ctx, fmt.Sprintf("file:meta:%d", item.ItemID))
		}
	} else if item.ItemType == "folder" {
		folderID := item.ItemID
		var folderParentID *uint64
		if item.OriginalPath != "" && item.OriginalPath != "0" {
			var pid uint64
			if _, err := fmt.Sscanf(item.OriginalPath, "%d", &pid); err == nil && pid > 0 {
				folderParentID = &pid
				parentID = pid
			}
		}
		// 修正点:目标目录可能已存在同名文件夹,需检测冲突,重命名或返回冲突而非裸报错
		if conflict, _ := uc.folderRepo.FindByNameAndParent(ctx, userID, item.Name, folderParentID); conflict != nil && conflict.ID != folderID {
			item.Name = uniqueName(item.Name, map[string]bool{conflict.Name: true})
		}
		folderUUID, _ := newUUID()
		folder := &Folder{ID: folderID, UUID: folderUUID, UserID: userID, Name: item.Name, ParentID: folderParentID}
		if _, err := uc.folderRepo.Create(ctx, folder); err != nil {
			return 0, err
		}
		trashedFiles, _, _ := uc.fileRepo.ListByParent(ctx, userID, &folderID, 1, "", 100000, "", "")
		if len(trashedFiles) > 0 {
			ids := make([]uint64, len(trashedFiles))
			for i, f := range trashedFiles {
				ids[i] = f.ID
			}
			_ = uc.fileRepo.BatchUpdateStatus(ctx, ids, 0)
		}
		uc.invalidateFileListCache(ctx, userID)
	}

	if err := uc.recycleRepo.Delete(ctx, recycleID); err != nil {
		return 0, err
	}
	return parentID, nil // 配额无需回加:Trash 阶段从未扣减
}

白话类比:还原就像把暂存篮里的文件放回原抽屉。因为文件一直算在配额里(篮子和抽屉都占地方),放回去数字不变,所以永远不会"放不下"。

5.4 永久删除(物理删盘 + 原子扣配额,含幂等修复)

实现思路

DeleteForeverpermanentlyDeleteItem 完成实际删除:

  • 文件storage.Delete(path) 物理删盘 → BatchUpdateStatus([id], 2)SubUsedStorage 原子扣配额。
  • 文件夹:拉取该文件夹下 status=1 的文件,逐个执行上述三步。

最后清元数据缓存、删 recycle_bins 记录、发 EventFileDeleted

⚠️ 新手必踩的坑(原实现)permanentlyDeleteItem 三步非原子,且错误只记日志不回滚;更致命的是无任何并发保护——用户 DeleteForeverCleanExpired 同时命中同一条目会重复扣减配额,或"还原"与"永久删除"并发造成损坏文件。下面给出商用修正版。

关键代码(修正点:条目级锁 + 二次校验 + 幂等扣减,保证配额只减一次)

// 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 {
		return ErrForbidden // 越权校验:禁止删除他人回收站条目
	}
	// 修正点:条目级分布式锁,与 Restore / CleanExpired 串行化
	lockKey := fmt.Sprintf("recycle:item🔒%d", recycleID)
	if uc.locker != nil {
		if err := uc.locker.Lock(ctx, lockKey); err != nil {
			return err
		}
		defer uc.locker.Unlock(ctx, lockKey)
	}
	// 修正点:加锁后再次确认条目仍在,若已被清理任务删除则直接返回(避免双扣配额)
	if _, err := uc.recycleRepo.FindByID(ctx, recycleID); err != nil {
		return err
	}

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

	if uc.eventPublisher != nil && len(deletedFileIDs) > 0 {
		_ = uc.eventPublisher.Publish(ctx, EventFileDeleted, &FileChangedPayload{
			UserID: userID, FileIDs: deletedFileIDs, Action: "deleted",
		})
	}
	return uc.recycleRepo.Delete(ctx, recycleID)
}

// permanentlyDeleteItem 永久删除单个回收站项目对应的文件记录与磁盘文件。
// 修正点:配额准确性优先——删盘尽力、改库、扣配额;顺序保证"磁盘删失败也不漏扣配额"。
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 {
			log.Warn("recycle: file not found for permanent delete", "id", item.ItemID, "err", err)
			return deletedFileIDs
		}
		// 1) 物理删盘(OSS Delete 幂等;失败仅记日志,不阻断扣减)
		if err := uc.storage.Delete(ctx, file.Path); err != nil {
			log.Error("recycle: failed to delete physical file", "path", file.Path, "err", err)
		}
		// 2) 改库 status=2
		if err := uc.fileRepo.BatchUpdateStatus(ctx, []uint64{item.ItemID}, 2); err != nil {
			log.Error("recycle: failed to update file status", "id", item.ItemID, "err", err)
		}
		// 3) 原子扣减配额(SubUsedStorage 内部带 user 级锁 + WHERE used_storage>=?,安全且不超卖)
		if uc.userUC != nil {
			if err := uc.userUC.SubUsedStorage(ctx, userID, file.Size); err != nil {
				log.Error("recycle: failed to subtract storage", "userID", userID, "err", err)
			}
		}
		deletedFileIDs = append(deletedFileIDs, item.ItemID)
	} else if item.ItemType == "folder" {
		folderID := item.ItemID
		trashedFiles, _, _ := uc.fileRepo.ListByParent(ctx, userID, &folderID, 1, "", 100000, "", "")
		for _, f := range trashedFiles {
			if err := uc.storage.Delete(ctx, f.Path); err != nil {
				log.Error("recycle: failed to delete physical file", "path", f.Path, "err", err)
			}
			if err := uc.fileRepo.BatchUpdateStatus(ctx, []uint64{f.ID}, 2); err != nil {
				log.Error("recycle: failed to update file status", "id", f.ID, "err", err)
			}
			if uc.userUC != nil {
				if err := uc.userUC.SubUsedStorage(ctx, userID, f.Size); err != nil {
					log.Error("recycle: failed to subtract storage", "userID", userID, "err", err)
				}
			}
			deletedFileIDs = append(deletedFileIDs, f.ID)
		}
	}
	if uc.cache != nil {
		for _, fid := range deletedFileIDs {
			_ = uc.cache.Delete(ctx, fmt.Sprintf("file:meta:%d", fid))
		}
	}
	return deletedFileIDs
}

5.5 清空回收站(批量永久删除,修正点:分批避免大事务)

实现思路

EmptyRecycle 拉取该用户全部回收站记录,循环对每条调用 permanentlyDeleteItem + recycleRepo.Delete,最后聚合 deletedFileIDs 发布一次 EventFileDeleted。任何单条失败只记日志、不中断整体流程。

修正点(商用):原实现 ListByUser(..., 100000) 一次性拉全量,超大用户会撑爆内存 / 长事务。改为基于游标的分批循环,每批 500 条,逐批清理并 DeleteBatch 删除回收站记录。

关键代码

// EmptyRecycle 永久删除用户的所有回收站项目(修正点:分批游标,避免一次性全量)
func (uc *RecycleUsecase) EmptyRecycle(ctx context.Context, userID uint64) error {
	const batchSize = 500
	var cursor string
	deletedFileIDs := make([]uint64, 0)
	for {
		// 修正点:用 keyset 游标分批拉取,而非一次性 LIMIT 100000
		items, next, err := uc.recycleRepo.ListByUser(ctx, userID, cursor, batchSize)
		if err != nil {
			return err
		}
		if len(items) == 0 {
			break
		}
		batchIDs := make([]uint64, 0, len(items))
		for _, item := range items {
			ids := uc.permanentlyDeleteItem(ctx, userID, item)
			deletedFileIDs = append(deletedFileIDs, ids...)
			batchIDs = append(batchIDs, item.ID)
		}
		// 修正点:批量删除回收站记录,减少 SQL 往返
		if err := uc.recycleRepo.DeleteBatch(ctx, batchIDs); err != nil {
			log.Error("recycle: failed to delete recycle items batch", "err", err)
		}
		if !strings.Contains(next, "|") {
			break // 兼容旧游标格式时及时终止
		}
		cursor = next
		if len(items) < batchSize {
			break
		}
	}
	if uc.eventPublisher != nil && len(deletedFileIDs) > 0 {
		_ = uc.eventPublisher.Publish(ctx, EventFileDeleted, &FileChangedPayload{
			UserID: userID, FileIDs: deletedFileIDs, Action: "deleted",
		})
	}
	return nil
}

service 层:

// EmptyRecycle 永久删除当前用户回收站中的所有项目。
func (s *RecycleService) EmptyRecycle(ctx context.Context, req *v1.EmptyRecycleRequest) (*v1.EmptyRecycleReply, error) {
	userID := CtxUserID(ctx)
	if userID == 0 {
		return nil, biz.ErrInvalidToken
	}
	if err := s.uc.EmptyRecycle(ctx, userID); err != nil {
		return nil, err
	}
	return &v1.EmptyRecycleReply{}, nil
}

5.6 过期记录自动清理(定时任务,真实存在 + 全局锁修正)

实现思路

CleanExpiredmain.go 注册的 cron 任务 recycle-clean0 0 3 * * *)每日调用:

  1. recycleRepo.ListExpired(now) 拉取所有 expire_at < now() 的记录(修正点:应分批)。
  2. 逐条 permanentlyDeleteItem + 删回收站记录,按 userID 聚合 deletedFileIDs
  3. 按用户维度发布 EventFileDeleted

data 层的 ListExpired 必须在 expire_at 上建索引。

修正点(关键商用隐患):多副本部署时每个实例都跑 recycle-clean,会重复全量扫描 + 重复扣减配额。加全局锁 recycle:clean:lock(带 TTL,如 1 小时)保证只有一个节点执行;条目级锁(permanentlyDeleteItemDeleteForever 路径已加)进一步防与用户操作双扣。

flowchart TD
    A[cron 每日 3:00 触发 CleanExpired] --> L{获取全局锁
recycle:clean:lock?} L -->|未获取 其他副本在执行| S[跳过本次] L -->|获取成功| E[分批 ListExpired expire_at < now] E --> P[逐条 条目级锁 + 物理删盘 + 改库 + 扣配额] P --> D[删 recycle 记录] D --> R[释放全局锁]

关键代码

// CleanExpired 删除已过期的回收站项目,由 cron 任务调用。
// 修正点:全局锁防止多副本重复执行;分批扫描避免一次载入海量记录。
func (uc *RecycleUsecase) CleanExpired(ctx context.Context) error {
	// 修正点:全局锁,多副本只跑一个
	cleanLock := "recycle:clean:lock"
	if uc.locker != nil {
		if err := uc.locker.Lock(ctx, cleanLock); err != nil {
			log.Warn("recycle: clean locked by another instance, skip", "err", err)
			return nil
		}
		defer uc.locker.Unlock(ctx, cleanLock)
	}

	deletedByUser := make(map[uint64][]uint64)
	const batch = 1000
	for {
		// 修正点:分批拉取过期记录,而非一次性 Find 全表
		expired, hasMore, err := uc.recycleRepo.ListExpiredBatch(ctx, time.Now(), batch)
		if err != nil {
			return err
		}
		for _, item := range expired {
			// 条目级锁在 permanentlyDeleteItem 上游(此处直接用内部删除,需自行加锁)
			lockKey := fmt.Sprintf("recycle:item🔒%d", item.ID)
			if uc.locker != nil {
				if err := uc.locker.Lock(ctx, lockKey); err == nil {
					ids := uc.permanentlyDeleteItem(ctx, item.UserID, item)
					_ = uc.locker.Unlock(ctx, lockKey)
					if len(ids) > 0 {
						deletedByUser[item.UserID] = append(deletedByUser[item.UserID], ids...)
					}
					if err := uc.recycleRepo.Delete(ctx, item.ID); err != nil {
						log.Error("recycle: failed to delete expired item", "id", item.ID, "err", err)
					}
				}
			} else {
				ids := uc.permanentlyDeleteItem(ctx, item.UserID, item)
				if len(ids) > 0 {
					deletedByUser[item.UserID] = append(deletedByUser[item.UserID], ids...)
				}
				_ = uc.recycleRepo.Delete(ctx, item.ID)
			}
		}
		if !hasMore {
			break
		}
	}

	if uc.eventPublisher != nil {
		for userID, fileIDs := range deletedByUser {
			_ = uc.eventPublisher.Publish(ctx, EventFileDeleted, &FileChangedPayload{
				UserID: userID, FileIDs: fileIDs, Action: "deleted",
			})
		}
	}
	return nil
}

data 层(过期记录分批查询,配合 expire_at 索引):

// ListExpiredBatch 分批拉取过期记录(修正点:带 Limit,避免全表载入)
func (r *recycleRepo) ListExpiredBatch(ctx context.Context, before time.Time, limit int) ([]*biz.RecycleItem, bool, error) {
	var pos []model.RecycleBin
	if err := r.db.WithContext(ctx).
		Where("expire_at < ?", before).
		Order("expire_at ASC, id ASC").
		Limit(limit + 1).
		Find(&pos).Error; err != nil {
		return nil, false, err
	}
	hasMore := len(pos) > limit
	if hasMore {
		pos = pos[:limit]
	}
	items := make([]*biz.RecycleItem, len(pos))
	for i := range pos {
		items[i] = toBizRecycleItem(&pos[i])
	}
	return items, hasMore, nil
}

确认:自动清理真实存在。cmd/server/main.gonewRecycleCleanTask 注册 cron 表达式 0 0 3 * * *,函数体直接 recycleUC.CleanExpired(ctx);并由 newTaskManager 注册进 cronTaskManager,进程经 ScheduledTaskServer 接入 Kratos 生命周期,启动即运行。


自测题与动手练习

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

  1. 回收站为什么"先软删除、后物理删除"?文件表 status 三个值(0/1/2)分别代表什么?recycle_bins.DeletedAt 是 GORM 软删除字段吗,为什么?
  2. 本项目的配额语义是"移入回收站不扣减、还原不回加、物理删除才扣减"。这样设计如何从根本上避免"还原后超配额"的冲突?若改成"删除即扣减",还原时要补什么校验?
  3. 恢复 / 永久删除 / 清空 / 定时清理分别怎么防止越权操作他人文件?Trash 批量遇到非本人文件时原实现有什么中间态隐患,怎么改?
  4. 永久删除真的从对象存储删文件了吗(个保法删除权)?若删盘成功但扣配额失败,配额会怎样漂移,项目靠什么兜底?
  5. 回收站自动清理是真实存在的吗?cron 表达式是什么、在哪个文件注册?多副本部署下它会带来什么新问题,怎么修?
  6. 同一条目"还原 / 永久删除 / 定时清理"并发会造成哪两类错误?用"条目级锁 + 二次存在校验"如何修复,请画时序/流程图。
  7. keyset 分页相比 OFFSET 好在哪?deleted_at 秒级非唯一会在分页边界导致什么,复合游标 (deleted_at, id) 如何解决?
  8. SubUsedStorage 怎么保证并发扣减不会扣成负数?它和"配额准确性优先于磁盘无孤儿"这一不变量有什么关系?

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

  1. 在纸上画文件 status 三态状态机,标出每个触发动作,并用 Mermaid stateDiagram 实现出来。
  2. 写一段 keyset 分页伪代码,传入 cursorpageSize 输出 items/nextCursor/hasMore;再构造两条 deleted_at 完全相同的记录,推演 Limit(n+1) 游标会重读还是漏读,体会复合游标修复。
  3. permanentlyDeleteItem 人为让 storage.Delete 返回 error,观察日志与后续 BatchUpdateStatus/SubUsedStorage 是否继续;再据此实现"条目级锁 + 加锁后二次 FindByID",写单测验证并发双删只扣一次配额。
  4. 起两个进程模拟多副本,各注册 recycle-clean,观察同一批过期记录是否被清理两次、配额是否被双扣;加上 recycle:clean:lock 全局锁后复测。

本章小结

  • 回收站核心是"软删除 + 延迟清理":status 三态(0 正常 / 1 回收站 / 2 永久删除)把"删除"拆成可反悔的两步,过期由 recycle-clean cron 任务(0 0 3 * * *,真实注册于 main.go)兜底。
  • 配额一致性靠语义而非补丁:采用"移入回收站不扣减、还原不回加、物理删除才原子扣减",用"磁盘占用 ≡ 配额占用"等价关系从根本上消灭"还原超配额冲突";扣减经 SubUsedStorage(用户级分布式锁 + WHERE used_storage >= ? 原子更新)保证安全不超卖。
  • 合规删除满足个保法删除权permanentlyDeleteItem 首步即 storage.Delete(path) 真删盘;以"配额准确性优先于磁盘无孤儿"为不变量安排顺序,并靠每日 storage-calibrate 对账兜底漂移。
  • 权限越权防护基本到位:所有操作校验 item.UserID == userID / 文件归属;但 Trash 批量中断会留半截状态,需先全量校验。
  • 并发是最大商用隐患:同条目"还原 / 永久删除 / 定时清理"并发及多副本重复执行会导致配额双扣或产生损坏文件。修复为"条目级锁 + 二次存在校验 + recycle:clean:lock 全局锁",使清理幂等。
  • 列表用 keyset 游标分页(deleted_at < cursor + Limit(n+1)),并应以 (deleted_at, id) 复合游标消除同秒边界歧义;EmptyRecycle/ListExpired 应改为分批游标避免大事务与全表载入;亿级流量再叠加异步删盘、读写分离、冷热分离与列表缓存。
About Me

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

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

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

目标

学AI,加油!加油!