diff --git a/docs/PERFORMANCE.md b/docs/PERFORMANCE.md new file mode 100644 index 00000000..f7b84d3c --- /dev/null +++ b/docs/PERFORMANCE.md @@ -0,0 +1,445 @@ +# Wavelet 系统性能分析与优化建议 + +> 分析日期:2026-06-17 +> 范围:Go 后端 + Next.js 前端 +> 目标:识别可能在生产环境真实出现的性能问题,并给出高 ROI 优化路线 + +--- + +## 目录 + +- [架构概览与核心瓶颈](#架构概览与核心瓶颈) +- [Critical — 高概率生产问题](#critical--高概率生产问题) +- [Medium — 中等风险](#medium--中等风险) +- [高价值优化路线图](#高价值优化路线图) +- [已做得好的设计](#已做得好的设计) +- [场景风险矩阵](#场景风险矩阵) +- [优先行动清单](#优先行动清单) + +--- + +## 架构概览与核心瓶颈 + +```mermaid +flowchart LR + subgraph frontend["前端 (Static Export)"] + A[HTML 静态壳] --> B[Hydrate] + B --> C["UserProvider.getUserInfo()"] + C --> D[页面数据请求] + D --> E[渲染] + end + + subgraph backend["后端热点路径"] + F["/f/{id}?quality=..."] --> G[DB 查 upload] + G --> H[迁移状态 DB 查询] + H --> I[白名单 Redis/DB] + I --> J{WebP 缓存命中?} + J -->|否| K["全量读文件 + 编码 + 磁盘缓存(全局锁)"] + J -->|是| L[返回] + end + + C -.->|串行阻塞| D +``` + +当前最大的结构性问题: + +1. **前端**:全客户端渲染 + 全局认证瀑布流,所有业务数据请求被 `getUserInfo` 串行阻塞。 +2. **后端**:文件服务路径(`/f/{id}`)是最高频热点,WebP 缓存未命中时在请求线程内做重 CPU/IO 工作,且磁盘缓存使用全局互斥锁。 + +--- + +## Critical — 高概率生产问题 + +### 1. 图片 WebP 服务:请求路径阻塞 + 全局锁串行化 + +**涉及文件**: + +- `internal/apps/upload/file_server.go` +- `pkg/cache/disk/cache.go` + +**问题描述**: + +缓存未命中时,在 HTTP 请求 goroutine 内执行: + +1. `io.ReadAll` 将原始文件全量读入内存 +2. 进程内 WebP 解码 + 编码 +3. 写入磁盘缓存 + +同时,磁盘缓存 `Get`/`Set` 使用**全局 `sync.Mutex`**,所有并发图片请求在缓存层完全串行。 + +```go +// file_server.go — 缓存 miss 时的重操作 +origBytes, err := getOriginalFileBytes(ctx, upload) // io.ReadAll +webpBytes, err = CompressImageToWebP(bytes.NewReader(origBytes), quality) +cache.Set(cacheKey, webpBytes, diskcache.NoExpiration) + +// pkg/cache/disk/cache.go — 全局互斥锁 +func (c *Cache) Get(key string) ([]byte, error) { + c.mu.Lock() + defer c.mu.Unlock() + // ... +} +``` + +**生产表现**: + +- 首次访问或缓存淘汰后,P99 延迟从几十毫秒飙升到数秒 +- 并发图片请求形成「隐形队列」 +- 大文件全量读入带来内存尖峰,可能触发 OOM 或 GC 停顿 + +**优化价值**:⭐⭐⭐⭐⭐ + +**建议**: + +- [ ] 磁盘缓存改用 `RWMutex`,读路径不互斥 +- [ ] 对同一 cache key 使用 `singleflight` 合并并发 miss +- [ ] 部署后强制执行 `upload:warm_image_cache` 异步预热任务 +- [ ] 考虑 miss 时先返回原图,后台异步生成 WebP + +--- + +### 2. 文件访问路径:每次请求多次 DB/Redis 查询 + +**涉及文件**: + +- `internal/apps/upload/storage_ops.go` +- `internal/apps/upload/file_server.go` + +**问题描述**: + +存储迁移状态**无进程内缓存**,每次文件操作都查询 `w_task_executions`: + +```go +// storage_ops.go +func StorageReadOnly(ctx context.Context) bool { + execution, ok, err := latestStorageMigrationExecution(ctx) + // ... +} + +func backendForStoredDriver(ctx context.Context, driver storage.Driver) (storage.Backend, error) { + // 可能再次调用 currentMigrationTargetConfig → 又一次相同 DB 查询 +} +``` + +公开文件白名单每次走 Redis/DB: + +```go +// file_server.go +func isFilePublic(ctx context.Context, uploadType string) bool { + sc.GetByKey(ctx, model.ConfigKeyFileAccessWhitelist) + // JSON 解析 + 遍历 +} +``` + +对比:`storage.Active()` 已有 5 秒内存缓存 + Redis pub/sub 失效机制,迁移状态却未复用该模式。 + +**生产表现**: + +- 每个 `/f/{id}` 请求额外 2–4 次 DB/Redis 往返 +- 图片站/CDN 场景下 QPS 放大后 PostgreSQL 连接池压力明显 + +**优化价值**:⭐⭐⭐⭐⭐ + +**建议**: + +- [ ] 为 `StorageReadOnly` / `latestStorageMigrationExecution` 增加 5s TTL 进程内缓存 +- [ ] 配置变更或迁移状态变化时通过 Redis pub/sub 失效 +- [ ] `file_access_whitelist` 增加进程内缓存,复用 `GetByKey` 的失效机制 + +--- + +### 3. Admin 文件统计:无界全表扫描 + +**涉及文件**:`internal/apps/upload/stats.go` + +**问题描述**: + +```go +err = db.DB(ctx).Model(&model.Upload{}). + Select("extension, mime_type, file_size"). + Where("status != ?", model.UploadStatusDeleted). + Scan(&fileRaws).Error +// 然后在 Go 中遍历全量结果做分类统计 +``` + +**生产表现**: + +- 10 万+ 文件时,管理端「文件统计」接口耗时数秒 +- 占用数百 MB 内存,可能拖垮 admin API + +**优化价值**:⭐⭐⭐⭐ + +**建议**: + +- [ ] 改为 SQL `GROUP BY` + `CASE WHEN` 聚合 +- [ ] 或维护增量统计表,上传/删除时更新计数 + +--- + +### 4. `w_uploads` 索引缺口 + +**涉及文件**:`internal/db/migrator/goose/postgres/202606090001_initial_schema.sql` + +**当前索引**:`user_id`, `file_path`, `hash`, `type` + +**缺失的高频查询索引**: + +| 查询场景 | 建议索引 | +|----------|----------| +| 清理任务 `status + created_at` | `(status, created_at)` | +| 存储迁移 `storage_driver + status` | `(storage_driver, status)` | +| 秒传去重 `hash + file_size + status` | `(hash, file_size, status)` | + +**生产表现**: + +- 数据量增长后,清理 worker、迁移任务、上传去重退化为顺序扫描 +- 后台任务积压,admin 操作变慢 + +**优化价值**:⭐⭐⭐⭐ + +**建议**: + +- [ ] 通过 goose migration 新增上述复合索引(PostgreSQL + SQLite 双方言) + +--- + +### 5. 批量 ZIP 下载:无上限 + 同步阻塞 + +**涉及文件**:`internal/apps/upload/routers.go` — `BatchDownloadFiles` + +**问题描述**: + +- `req.IDs` 无数量上限 +- 在请求 goroutine 内串行打开每个文件并 `io.Copy` 到 ZIP +- 远端 S3 场景下单个文件就可能耗时数秒 + +**生产表现**: + +- 网关超时、连接耗尽 +- Admin 批量下载操作卡死 + +**优化价值**:⭐⭐⭐⭐ + +**建议**: + +- [ ] 限制单次批量数量(如 max 50) +- [ ] 或改为 Asynq 后台任务生成 ZIP,前端轮询下载链接 + +--- + +### 6. 前端全局认证瀑布流 + +**涉及文件**: + +- `frontend/contexts/user-context.tsx` +- `frontend/app/(main)/layout.tsx` + +**问题描述**: + +```tsx +// user-context.tsx — 挂载时获取用户 +useEffect(() => { + fetchUser() +}, [fetchUser]) + +// layout.tsx — 阻塞所有子页面渲染 +if (loading || !user) { + return +} +``` + +**生产表现**: + +- 每次进入 `/home`、`/files`、`/admin/*` 都先等 `getUserInfo`(约 200–800ms) +- 页面级数据请求无法并行启动,TTI 被硬性拉长 + +**优化价值**:⭐⭐⭐⭐⭐ + +**建议**: + +- [ ] Layout 不阻塞渲染,子页面自行处理未登录状态 +- [ ] 或 Server Component 通过 cookie 预取 session,消除客户端首屏等待 +- [ ] `/login`、`/register` 跳过 `getUserInfo` + +--- + +### 7. 实时日志面板:2000 行 DOM 无虚拟化 + +**涉及文件**:`frontend/components/common/admin/app-logs.tsx` + +**问题描述**: + +- 日志上限 2000 行(内存有界,但 DOM 无界) +- 每行渲染完整 `
`,无虚拟滚动 +- `@tanstack/react-virtual` 已在 `package.json` 但未使用 + +**生产表现**: + +- 管理员开着日志 Tab 时 CPU/内存持续升高 +- 滚动卡顿,长时间运行拖慢整台机器 + +**优化价值**:⭐⭐⭐⭐ + +**建议**: + +- [ ] 使用 `useVirtualizer` 只渲染可视区域行 +- [ ] 行组件 `React.memo` 避免无效重渲染 + +--- + +## Medium — 中等风险 + +| # | 问题 | 位置 | 影响 | +|---|------|------|------| +| 1 | 公共配置接口无 Redis 缓存 | `internal/model/system_configs.go` — `ListVisibleSystemConfigs` | 每次前端启动/登录直查 PostgreSQL | +| 2 | CAPTCHA 每次 5 次独立 `GetByKey` | `internal/apps/cap/manager.go` | 登录高峰 Redis 压力 | +| 3 | OIDC 每次 `oidc.NewProvider` 无缓存 | `internal/apps/oauth/sources.go:164` | 登录发起/回调多一次外部 HTTP | +| 4 | CORS 每次跨域查 `server_address` 配置 | `internal/router/middlewares.go:75` | 预检请求放大 | +| 5 | 推送通知无界 goroutine + 逐 target DB 查询 | `internal/apps/admin/push/events.go:102` | 通知风暴时 goroutine/DB 双压 | +| 6 | 上传清理:每文件一个事务 | `internal/apps/upload/cleanup.go` | 大量 pending 文件时 commit 风暴 | +| 7 | ClickHouse 风控:每请求 `json.Marshal` 全部 headers | `internal/apps/risk_control/middleware.go:58` | 高 QPS 时 CPU 开销(写入本身已异步批处理) | +| 8 | 存储迁移日志大量写 Redis | `internal/apps/upload/storage_migration_task.go` | 迁移期间 Redis CPU/内存压力 | +| 9 | 存储迁移后二次 SHA 全量读取验证 | `storage_migration_task.go` | 迁移期间对象 I/O 翻倍 | +| 10 | Admin 状态页 5s 轮询 | `frontend/components/common/admin/status.tsx` | Tab 常驻时持续打后端 | +| 11 | 路由切换 500ms fade 动画 | `frontend/app/(main)/layout.tsx:53-60` | 即使数据已缓存,感知仍慢 | +| 12 | 无 `next/dynamic` 代码分割 | 全项目 | Admin 首包 300–450KB+ | +| 13 | 19/24 个 `page.tsx` 为 `"use client"` | 各路由 | 无法 RSC 预取,bundle 偏大 | +| 14 | Admin 部分页面用 `useEffect` 而非 React Query | `access-logs.tsx`, `task-executions.tsx` 等 | 无缓存去重,重复请求 | +| 15 | 登录页 OIDC sources 等待 public config | `frontend/components/auth/login-form.tsx` | 多 1 次 RTT 瀑布 | +| 16 | Users 表每行嵌套 3 个 `TooltipProvider` | `frontend/app/(main)/admin/users/page.tsx` | 50+ 行时不必要重渲染 | +| 17 | 缩略图用原生 `` 无 lazy loading | `file-list.tsx`, `file-manager.tsx` | 文件管理页初始解码压力大 | +| 18 | `@/lib/services` barrel 导入 | ~40 个文件 | 单路由 bundle 膨胀 10–30KB | +| 19 | SQLite 模式无连接池调优 | `internal/db/postgres.go` | 默认 SQLite 写锁瓶颈 | +| 20 | Session Redis 仅用第一个地址 | `internal/router/router.go` | Sentinel/Cluster 场景不一致 | + +--- + +## 高价值优化路线图 + +### P0 — 立即做(1–2 周,收益最大) + +| # | 优化项 | 涉及模块 | 预期收益 | 复杂度 | +|---|--------|----------|----------|--------| +| 1 | WebP:`singleflight` + `RWMutex` + 强制预热 | `file_server.go`, `pkg/cache/disk/` | 图片 P99 ↓ 80%+,并发吞吐 ↑ 5–10x | 中 | +| 2 | 缓存 `StorageReadOnly` / 迁移状态 | `storage_ops.go` | 每文件请求减少 1–3 次 DB | 低 | +| 3 | 内存缓存 `file_access_whitelist` | `file_server.go` | 每公开文件请求减少 1 次 Redis | 低 | +| 4 | `GetFileStats` 改为 SQL 聚合 | `stats.go` | Admin 统计从 O(n) → O(1) | 低 | +| 5 | 新增 `w_uploads` 复合索引 | goose migration | 清理/迁移/秒传全面加速 | 低 | +| 6 | 前端日志虚拟化 | `app-logs.tsx` | Admin 日志 Tab 流畅度质变 | 低 | +| 7 | Admin 重模块 `dynamic()` 懒加载 | `database/page.tsx`, `logs/page.tsx`, `settings/page.tsx` 等 | 首包 JS ↓ 150–300KB | 低 | + +### P1 — 短期(2–4 周) + +| # | 优化项 | 预期收益 | +|---|--------|----------| +| 8 | 认证并行化:layout 不阻塞 / Server 预取 session | TTI ↓ 200–800ms | +| 9 | `ListVisibleSystemConfigs` 加 Redis 缓存 | 前端冷启动加速 | +| 10 | CAPTCHA 配置快照(一次加载 5 个 key) | 验证码路径 Redis ops ↓ 80% | +| 11 | OIDC Provider/JWKS 进程内缓存(TTL 1h) | 登录延迟 ↓ 100–500ms | +| 12 | 批量下载限制(max 50)或异步任务 | 消除网关超时风险 | +| 13 | Admin `useEffect` 数据获取迁移到 React Query | 去重、缓存、后台刷新 | +| 14 | 登录页并行请求 public config + auth sources | 登录页 ↓ 100–300ms | +| 15 | 状态轮询在 `document.hidden` 时暂停 | 降低后台 + 客户端负载 | + +### P2 — 中期架构演进 + +| # | 优化项 | 预期收益 | +|---|--------|----------| +| 16 | 批量 ZIP 改为 Asynq 后台任务 | 彻底解耦长耗时操作 | +| 17 | 存储迁移日志降噪 + 跳过已验证文件二次 SHA | 迁移期间 Redis/I/O ↓ 50% | +| 18 | 推送通知 target 批量解析(`WHERE id IN ?`) | 通知风暴 DB 查询 ↓ N 倍 | +| 19 | 上传清理改为批量 UPDATE + 异步存储删除 | 减少 DB commit 频率 | +| 20 | 路由动画 0.5s → 0.15s 或纯 CSS | 导航感知速度 ↑ | +| 21 | 服务导入收窄(直接 import 具体 Service) | 每路由 bundle ↓ 10–30KB | +| 22 | Admin 路由级 `loading.tsx` + Suspense | 渐进式渲染体验 | +| 23 | 缩略图 `loading="lazy"` + 固定尺寸 | 文件管理页初始 paint 加速 | + +--- + +## 已做得好的设计 + +以下设计说明团队已有性能意识,优化应在此基础上增量改进,**不必重复造轮子**: + +| # | 设计 | 位置 | +|---|------|------| +| 1 | 系统配置单 key Redis Hash 缓存 | `internal/model/system_configs.go` — `GetByKey` | +| 2 | Storage Backend 单例 + 5s TTL + pub/sub 失效 | `internal/storage/storage.go` — `Active()` | +| 3 | 推送事件/渠道 24h Redis 缓存 + GORM hook 失效 | `internal/model/push_event.go`, `push_channel.go` | +| 4 | 风控日志异步批写 ClickHouse(1 万缓冲 + 1000 条/1s + 429 背压) | `internal/apps/risk_control/` | +| 5 | HTTP 连接池统一(`httppool` + OTel) | `pkg/httppool/` | +| 6 | DB/Redis 连接池显式配置 | `config.yaml`, `internal/db/` | +| 7 | 游标分批处理(`id > ? LIMIT n`) | `cleanup.go`, image warmup | +| 8 | 存储迁移并发上限 `errgroup.SetLimit(10)` | `storage_migration_task.go` | +| 9 | 邮件/推送走 Asynq,不在 HTTP 路径同步发送 | `user/logics.go`, `push/events.go` | +| 10 | 文件服务 ETag/304 + 原图 `DataFromReader` 流式返回 | `file_server.go` | +| 11 | 无 GORM `Preload` 滥用 | 全项目 | +| 12 | 前端 API 请求去重(`pendingRequests` Map) | `frontend/lib/services/core/api-client.ts` | +| 13 | React Query 全局 30s `staleTime` | `frontend/components/providers/query-provider.tsx` | +| 14 | React Compiler 已启用 | `frontend/next.config.ts` | +| 15 | 读副本支持(`dbresolver`) | `internal/db/postgres.go` | +| 16 | 任务执行日志 Redis 缓冲 + 批量回写 | `internal/model/task_execution.go` | + +--- + +## 场景风险矩阵 + +| 场景 | 最可能爆的点 | 对应优先级 | +|------|-------------|-----------| +| 图片站 / 公开相册 | WebP miss + 磁盘锁 + 白名单 Redis | P0 #1, #2, #3 | +| 文件量 10 万+ | 统计全表扫描 + 索引缺失 + 清理慢 | P0 #4, #5 | +| 管理端日常使用 | 认证瀑布 + 大 bundle + 日志 DOM | P0 #6, #7 | +| 存储迁移进行中 | 迁移状态重复查 + Redis 日志风暴 | P1 #2, P2 #17 | +| 登录高峰 | CAPTCHA 5×Redis + OIDC discovery | P1 #10, #11 | +| 多租户 / 跨域前端 | CORS 配置查询 + 公共配置无缓存 | P1 #9 | +| 批量文件操作 | ZIP 同步打包无上限 | P0 #5, P1 #12 | + +--- + +## 优先行动清单 + +如果只选 **3 件事** 先做(预计用户感知延迟降低 50–70%): + +1. **WebP 路径解耦** — `singleflight` + `RWMutex` + 部署后预热 +2. **文件路径查询缓存** — 迁移状态 + 白名单进程内缓存 +3. **前端认证与首屏并行化** — 消除全局 auth gate + Admin 代码分割 + +### 实施检查清单 + +``` +P0 后端 +[ ] disk cache RWMutex + singleflight +[ ] StorageReadOnly 5s 缓存 + pub/sub 失效 +[ ] file_access_whitelist 进程内缓存 +[ ] GetFileStats SQL 聚合改写 +[ ] w_uploads 复合索引 migration +[ ] 批量下载数量上限 + +P0 前端 +[ ] app-logs.tsx 虚拟滚动 +[ ] SQLConsole / Recharts / Settings Tabs dynamic import +[ ] 认证 gate 并行化 + +P1 +[ ] ListVisibleSystemConfigs Redis 缓存 +[ ] CAPTCHA 配置快照 +[ ] OIDC Provider 缓存 +[ ] Admin useEffect → React Query 统一 +[ ] 状态轮询 visibility 感知 +``` + +--- + +## 附录:关键代码路径索引 + +| 路径 | 文件 | 说明 | +|------|------|------| +| 图片服务 | `internal/apps/upload/file_server.go` | `/f/{id}` 热点 | +| 磁盘缓存 | `pkg/cache/disk/cache.go` | 全局 Mutex | +| 迁移状态 | `internal/apps/upload/storage_ops.go` | 无缓存 DB 查询 | +| 文件统计 | `internal/apps/upload/stats.go` | 全表扫描 | +| 批量下载 | `internal/apps/upload/routers.go` | 同步 ZIP | +| 上传索引 | `internal/db/migrator/goose/*/202606090001_initial_schema.sql` | 缺失复合索引 | +| 认证 gate | `frontend/app/(main)/layout.tsx` | 阻塞渲染 | +| 用户上下文 | `frontend/contexts/user-context.tsx` | 挂载时 fetch | +| 实时日志 | `frontend/components/common/admin/app-logs.tsx` | 无虚拟化 | +| API 去重 | `frontend/lib/services/core/api-client.ts` | 已有,可复用模式 | \ No newline at end of file diff --git a/internal/apps/admin/system_config/routers.go b/internal/apps/admin/system_config/routers.go index 4ccf38ce..a5db3332 100644 --- a/internal/apps/admin/system_config/routers.go +++ b/internal/apps/admin/system_config/routers.go @@ -11,6 +11,7 @@ import ("context" "net/http" "time" + "github.com/Rain-kl/Wavelet/internal/apps/upload" "github.com/Rain-kl/Wavelet/internal/db" "github.com/Rain-kl/Wavelet/internal/model" "github.com/Rain-kl/Wavelet/internal/storage" @@ -255,9 +256,15 @@ func UpdateSystemConfig(c *gin.Context) { } if key == model.ConfigKeyStorageConfig { + upload.ResetAccessCaches() + upload.PublishAccessCacheInvalidation(c.Request.Context()) storage.ResetCache() storage.PublishCacheInvalidation(c.Request.Context()) } + if key == model.ConfigKeyFileAccessWhitelist { + upload.ResetAccessCaches() + upload.PublishAccessCacheInvalidation(c.Request.Context()) + } c.JSON(http.StatusOK, response.OKNil()) } diff --git a/internal/apps/upload/access_cache.go b/internal/apps/upload/access_cache.go new file mode 100644 index 00000000..efddb0ee --- /dev/null +++ b/internal/apps/upload/access_cache.go @@ -0,0 +1,200 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package upload + +import ( + "context" + "encoding/json" + "strings" + "sync" + "time" + + "github.com/Rain-kl/Wavelet/internal/db" + "github.com/Rain-kl/Wavelet/internal/model" + "github.com/Rain-kl/Wavelet/internal/storage" +) + +const accessCacheTTL = 5 * time.Second + +const fileAccessInvalidationChannel = "upload:file_access_invalidation" + +type migrationAccessState struct { + readOnly bool + target storage.Config + hasTarget bool + targetErr error + loadErr error +} + +var ( + accessCacheOnce sync.Once + + migrationAccessMu sync.RWMutex + migrationAccessCached migrationAccessState + migrationAccessValid bool + migrationAccessCheckedAt time.Time + + fileAccessWhitelistMu sync.RWMutex + fileAccessWhitelistTypes map[string]struct{} + fileAccessWhitelistValid bool + fileAccessWhitelistCheckedAt time.Time +) + +// ResetAccessCaches clears in-process upload access caches. +func ResetAccessCaches() { + migrationAccessMu.Lock() + migrationAccessValid = false + migrationAccessMu.Unlock() + + fileAccessWhitelistMu.Lock() + fileAccessWhitelistValid = false + fileAccessWhitelistTypes = nil + fileAccessWhitelistMu.Unlock() +} + +// PublishAccessCacheInvalidation broadcasts upload access cache eviction to all nodes. +func PublishAccessCacheInvalidation(ctx context.Context) { + if db.Redis != nil { + _ = db.Redis.Publish(ctx, fileAccessInvalidationChannel, "reset").Err() + } +} + +func ensureAccessCacheListener() { + accessCacheOnce.Do(startAccessCacheInvalidationListener) +} + +func startAccessCacheInvalidationListener() { + if db.Redis == nil { + return + } + + go func() { + pubsub := db.Redis.Subscribe( + context.Background(), + storage.ConfigInvalidationChannel, + fileAccessInvalidationChannel, + ) + defer func() { + _ = pubsub.Close() + }() + + for range pubsub.Channel() { + ResetAccessCaches() + } + }() +} + +func loadMigrationAccessState(ctx context.Context) migrationAccessState { + ensureAccessCacheListener() + + migrationAccessMu.RLock() + if migrationAccessValid && time.Since(migrationAccessCheckedAt) < accessCacheTTL { + state := migrationAccessCached + migrationAccessMu.RUnlock() + return state + } + migrationAccessMu.RUnlock() + + migrationAccessMu.Lock() + defer migrationAccessMu.Unlock() + + if migrationAccessValid && time.Since(migrationAccessCheckedAt) < accessCacheTTL { + return migrationAccessCached + } + + migrationAccessCached = buildMigrationAccessState(ctx) + migrationAccessValid = true + migrationAccessCheckedAt = time.Now() + return migrationAccessCached +} + +func buildMigrationAccessState(ctx context.Context) migrationAccessState { + execution, ok, err := latestStorageMigrationExecution(ctx) + if err != nil { + return migrationAccessState{loadErr: err, readOnly: true} + } + if !ok { + return migrationAccessState{} + } + + state := migrationAccessState{ + readOnly: execution.Status != model.TaskExecutionStatusSucceeded, + } + if execution.Status == model.TaskExecutionStatusSucceeded { + return state + } + + target, err := parseMigrationTargetConfig(ctx, []byte(execution.Payload)) + if err != nil { + state.targetErr = err + return state + } + + state.target = target + state.hasTarget = true + return state +} + +func loadFileAccessWhitelist(ctx context.Context) map[string]struct{} { + ensureAccessCacheListener() + + fileAccessWhitelistMu.RLock() + if fileAccessWhitelistValid && time.Since(fileAccessWhitelistCheckedAt) < accessCacheTTL { + types := fileAccessWhitelistTypes + fileAccessWhitelistMu.RUnlock() + return types + } + fileAccessWhitelistMu.RUnlock() + + fileAccessWhitelistMu.Lock() + defer fileAccessWhitelistMu.Unlock() + + if fileAccessWhitelistValid && time.Since(fileAccessWhitelistCheckedAt) < accessCacheTTL { + return fileAccessWhitelistTypes + } + + fileAccessWhitelistTypes = fetchFileAccessWhitelist(ctx) + fileAccessWhitelistValid = true + fileAccessWhitelistCheckedAt = time.Now() + return fileAccessWhitelistTypes +} + +func fetchFileAccessWhitelist(ctx context.Context) map[string]struct{} { + whitelist := parseFileAccessWhitelist(ctx) + types := make(map[string]struct{}, len(whitelist)) + for _, item := range whitelist { + types[strings.ToLower(item)] = struct{}{} + } + return types +} + +func parseFileAccessWhitelist(ctx context.Context) []string { + var sc model.SystemConfig + if err := sc.GetByKey(ctx, model.ConfigKeyFileAccessWhitelist); err != nil || sc.Value == "" { + return []string{defaultPublicUploadType} + } + + var whitelist []string + if err := json.Unmarshal([]byte(sc.Value), &whitelist); err == nil && len(whitelist) > 0 { + return whitelist + } + + whitelist = parseCommaSeparatedWhitelist(sc.Value) + if len(whitelist) == 0 { + return []string{defaultPublicUploadType} + } + return whitelist +} + +func parseCommaSeparatedWhitelist(value string) []string { + parts := strings.Split(value, ",") + whitelist := make([]string, 0, len(parts)) + for _, part := range parts { + part = strings.TrimSpace(part) + if part != "" { + whitelist = append(whitelist, part) + } + } + return whitelist +} \ No newline at end of file diff --git a/internal/apps/upload/access_cache_test.go b/internal/apps/upload/access_cache_test.go new file mode 100644 index 00000000..155023b5 --- /dev/null +++ b/internal/apps/upload/access_cache_test.go @@ -0,0 +1,97 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package upload + +import ( + "context" + "testing" + "time" + + "github.com/Rain-kl/Wavelet/internal/db" + "github.com/Rain-kl/Wavelet/internal/model" + "github.com/Rain-kl/Wavelet/internal/testhelper" +) + +func TestLoadMigrationAccessStateCachesResult(t *testing.T) { + _, _, cleanup := testhelper.SetupTestEnvironment(t) + defer cleanup() + ResetAccessCaches() + + ctx := context.Background() + first := loadMigrationAccessState(ctx) + second := loadMigrationAccessState(ctx) + + if first.readOnly != second.readOnly { + t.Fatalf("readOnly mismatch: first=%v second=%v", first.readOnly, second.readOnly) + } + if first.hasTarget != second.hasTarget { + t.Fatalf("hasTarget mismatch: first=%v second=%v", first.hasTarget, second.hasTarget) + } +} + +func TestIsFilePublicUsesCachedWhitelist(t *testing.T) { + _, _, cleanup := testhelper.SetupTestEnvironment(t) + defer cleanup() + ResetAccessCaches() + + ctx := context.Background() + if !isFilePublic(ctx, "avatar") { + t.Fatal("expected avatar to be public by default") + } + if isFilePublic(ctx, "attachment") { + t.Fatal("expected attachment to be private by default") + } + if !isFilePublic(ctx, "AVATAR") { + t.Fatal("expected whitelist lookup to be case-insensitive") + } +} + +func TestResetAccessCachesRefreshesWhitelist(t *testing.T) { + dbConn, _, cleanup := testhelper.SetupTestEnvironment(t) + defer cleanup() + ResetAccessCaches() + + ctx := context.Background() + if !isFilePublic(ctx, "avatar") { + t.Fatal("expected seeded avatar whitelist before reset") + } + + var sc model.SystemConfig + if err := dbConn.Where("key = ?", model.ConfigKeyFileAccessWhitelist).First(&sc).Error; err != nil { + t.Fatalf("load whitelist config: %v", err) + } + sc.Value = `["attachment"]` + if err := dbConn.Save(&sc).Error; err != nil { + t.Fatalf("save whitelist config: %v", err) + } + if err := db.HSetJSON(ctx, model.SystemConfigRedisHashKey, model.ConfigKeyFileAccessWhitelist, &sc); err != nil { + t.Fatalf("refresh whitelist redis cache: %v", err) + } + + ResetAccessCaches() + if !isFilePublic(ctx, "attachment") { + t.Fatal("expected attachment to be public after whitelist refresh") + } + if isFilePublic(ctx, "avatar") { + t.Fatal("expected avatar to be private after whitelist refresh") + } +} + +func TestAccessCacheTTLExpires(t *testing.T) { + _, _, cleanup := testhelper.SetupTestEnvironment(t) + defer cleanup() + ResetAccessCaches() + + ctx := context.Background() + _ = loadFileAccessWhitelist(ctx) + + fileAccessWhitelistMu.Lock() + fileAccessWhitelistCheckedAt = time.Now().Add(-accessCacheTTL - time.Second) + fileAccessWhitelistMu.Unlock() + + // Should still work after TTL by reloading from config. + if !isFilePublic(ctx, "avatar") { + t.Fatal("expected whitelist reload after TTL expiration") + } +} \ No newline at end of file diff --git a/internal/apps/upload/cleanup.go b/internal/apps/upload/cleanup.go index 3d4977d8..e3121ad2 100644 --- a/internal/apps/upload/cleanup.go +++ b/internal/apps/upload/cleanup.go @@ -110,6 +110,7 @@ func (h *SystemCleanupHandler) Execute(ctx context.Context, _ []byte) (*task.Tas continue } + recordUploadStatsRemove(ctx, &u) totalDeleted++ lastID = u.ID } diff --git a/internal/apps/upload/constants.go b/internal/apps/upload/constants.go index b927a86b..22bff1d6 100644 --- a/internal/apps/upload/constants.go +++ b/internal/apps/upload/constants.go @@ -14,5 +14,7 @@ const ( imageQualityMedium = "medium" imageQualityHigh = "high" imageQualityOrigin = "origin" - storageDriverLocal = string(storage.DriverLocal) + storageDriverLocal = string(storage.DriverLocal) + defaultPublicUploadType = "avatar" + fileStatsTrendDays = 7 ) diff --git a/internal/apps/upload/file_management.go b/internal/apps/upload/file_management.go index c02dfe82..af14f8cf 100644 --- a/internal/apps/upload/file_management.go +++ b/internal/apps/upload/file_management.go @@ -140,6 +140,7 @@ func DeleteFile(c *gin.Context) { c.JSON(http.StatusOK, response.Err(ErrDeleteFileFailed)) return } + recordUploadStatsRemove(ctx, &upload) c.JSON(http.StatusOK, response.OKNil()) } @@ -290,6 +291,7 @@ func DeleteMyFile(c *gin.Context) { c.JSON(http.StatusOK, response.Err(ErrDeleteFileFailed)) return } + recordUploadStatsRemove(ctx, &upload) c.JSON(http.StatusOK, response.OKNil()) } diff --git a/internal/apps/upload/file_server.go b/internal/apps/upload/file_server.go index 926b5352..4fd247ca 100644 --- a/internal/apps/upload/file_server.go +++ b/internal/apps/upload/file_server.go @@ -4,9 +4,9 @@ package upload -import ("bytes" +import ( + "bytes" "context" - "encoding/json" "errors" "fmt" "io" @@ -22,9 +22,18 @@ import ("bytes" "github.com/Rain-kl/Wavelet/internal/util" "github.com/Rain-kl/Wavelet/pkg/logger" "github.com/gin-gonic/gin" + "golang.org/x/sync/singleflight" "gorm.io/gorm" ) +var compressedImageFlight singleflight.Group + +type compressedImageCacheResult struct { + bytes []byte + cached bool + err error +} + // ServeFileByID 根据 ID 获取并提供已上传的文件 // @Summary 获取已上传文件 // @Description 根据文件 ID 获取并提供已上传的临时或正式文件,若配置了缓存则优先走本地缓存,否则从 S3 等后端存储读取并流式返回 @@ -194,21 +203,51 @@ func ensureCompressedImageCache( return nil, false, fmt.Errorf("read compressed image cache: %w", err) } + result, err, _ := compressedImageFlight.Do(cacheKey, func() (any, error) { + return generateCompressedImageCache(ctx, upload, quality, cacheKey) + }) + if err != nil { + return nil, false, err + } + + res := result.(compressedImageCacheResult) + return res.bytes, res.cached, res.err +} + +func generateCompressedImageCache( + ctx context.Context, + upload *model.Upload, + quality string, + cacheKey string, +) (compressedImageCacheResult, error) { + cache := diskcache.GetGlobalCache() + + webpBytes, err := cache.Get(cacheKey) + if err == nil { + return compressedImageCacheResult{bytes: webpBytes, cached: true}, nil + } + if !errors.Is(err, diskcache.ErrCacheMiss) { + return compressedImageCacheResult{}, fmt.Errorf("read compressed image cache: %w", err) + } + origBytes, err := getOriginalFileBytes(ctx, upload) if err != nil { - return nil, false, fmt.Errorf("read original image: %w", err) + return compressedImageCacheResult{}, fmt.Errorf("read original image: %w", err) } webpBytes, err = CompressImageToWebP(bytes.NewReader(origBytes), quality) if err != nil { - return nil, false, fmt.Errorf("compress image to WebP: %w", err) + return compressedImageCacheResult{}, fmt.Errorf("compress image to WebP: %w", err) } if err := cache.Set(cacheKey, webpBytes, diskcache.NoExpiration); err != nil { - return webpBytes, false, fmt.Errorf("write compressed image cache: %w", err) + return compressedImageCacheResult{ + bytes: webpBytes, + err: fmt.Errorf("write compressed image cache: %w", err), + }, nil } - return webpBytes, false, nil + return compressedImageCacheResult{bytes: webpBytes}, nil } func imageCompressionCacheKey(upload *model.Upload, quality string) string { @@ -254,30 +293,9 @@ func getOriginalFileBytes(ctx context.Context, upload *model.Upload) ([]byte, er // isFilePublic 校验文件类型是否在公开访问白名单中 func isFilePublic(ctx context.Context, uploadType string) bool { - var sc model.SystemConfig - var whitelist []string - if err := sc.GetByKey(ctx, model.ConfigKeyFileAccessWhitelist); err == nil && sc.Value != "" { - if err := json.Unmarshal([]byte(sc.Value), &whitelist); err != nil { - // 降级使用逗号分隔解析 - parts := strings.Split(sc.Value, ",") - for _, p := range parts { - p = strings.TrimSpace(p) - if p != "" { - whitelist = append(whitelist, p) - } - } - } - } else { - // 默认兜底白名单为 avatar - whitelist = []string{"avatar"} - } - - for _, w := range whitelist { - if strings.EqualFold(w, uploadType) { - return true - } - } - return false + whitelist := loadFileAccessWhitelist(ctx) + _, ok := whitelist[strings.ToLower(uploadType)] + return ok } func checkPrivateFileOwner(c *gin.Context, ownerID uint64) error { diff --git a/internal/apps/upload/routers.go b/internal/apps/upload/routers.go index 639f4930..be9209e0 100644 --- a/internal/apps/upload/routers.go +++ b/internal/apps/upload/routers.go @@ -120,7 +120,7 @@ func UploadFile(c *gin.Context) { accessModeStr := c.PostForm("access_mode") var accessMode int if accessModeStr == "" { - if uploadType == "avatar" { + if uploadType == defaultPublicUploadType { accessMode = 1 } else { accessMode = 0 @@ -388,6 +388,7 @@ func tryInstantUpload(ctx context.Context, c *gin.Context, currUser *model.User, c.JSON(http.StatusOK, response.Err(ErrSaveUploadRecordFailed)) return true, err } + recordUploadStatsAdd(ctx, &newUpload) logger.InfoF(ctx, "文件触发秒传成功! ID: %d, Path: %s", id, existing.FilePath) c.JSON(http.StatusOK, response.OK(newUpload)) @@ -458,5 +459,6 @@ func saveUploadRecord(ctx context.Context, upload *model.Upload, storageDriver, } return ErrSaveUploadRecordFailed } + recordUploadStatsAdd(ctx, upload) return "" } diff --git a/internal/apps/upload/routers_test.go b/internal/apps/upload/routers_test.go index 996b0138..98a6aeff 100644 --- a/internal/apps/upload/routers_test.go +++ b/internal/apps/upload/routers_test.go @@ -806,6 +806,9 @@ func TestGetFileStats(t *testing.T) { t.Fatalf("failed to create upload: %v", err) } } + if err := RebuildUploadStats(context.Background()); err != nil { + t.Fatalf("failed to rebuild upload stats: %v", err) + } req, _ := http.NewRequest("GET", "/api/v1/admin/uploads/stats", nil) w := httptest.NewRecorder() diff --git a/internal/apps/upload/stats.go b/internal/apps/upload/stats.go index c5ff062f..46527730 100644 --- a/internal/apps/upload/stats.go +++ b/internal/apps/upload/stats.go @@ -3,7 +3,8 @@ package upload -import ("net/http" +import ( + "net/http" "strings" "time" @@ -11,7 +12,8 @@ import ("net/http" "github.com/Rain-kl/Wavelet/internal/model" "github.com/gin-gonic/gin" - "github.com/Rain-kl/Wavelet/internal/common/response") + "github.com/Rain-kl/Wavelet/internal/common/response" +) const ( catImage = "图片" @@ -56,137 +58,78 @@ type fileStatsResponse struct { func GetFileStats(c *gin.Context) { ctx := c.Request.Context() - // 1. 获取总文件数与总文件大小 - var summary struct { - TotalCount int64 `json:"total_count"` - TotalSize int64 `json:"total_size"` - } - err := db.DB(ctx).Model(&model.Upload{}). - Select("COUNT(*) as total_count, COALESCE(SUM(file_size), 0) as total_size"). - Where("status != ?", model.UploadStatusDeleted). - Scan(&summary).Error - if err != nil { + var stats []model.UploadStat + if err := db.DB(ctx).Find(&stats).Error; err != nil { c.JSON(http.StatusOK, response.Err(err.Error())) return } - // 2. 获取业务类型分布 (Group By type) - type rawDist struct { - Key string `gorm:"column:key"` - Count int64 `gorm:"column:count"` - Size int64 `gorm:"column:size"` - } - var typeRaw []rawDist - err = db.DB(ctx).Model(&model.Upload{}). - Select("type as key, COUNT(*) as count, COALESCE(SUM(file_size), 0) as size"). - Where("status != ?", model.UploadStatusDeleted). - Group("type"). - Scan(&typeRaw).Error - if err != nil { - c.JSON(http.StatusOK, response.Err(err.Error())) - return - } - - types := make([]distributionItem, 0, len(typeRaw)) - for _, tr := range typeRaw { - name := tr.Key - if name == "" { - name = "generic" - } - types = append(types, distributionItem{ - Name: name, - Count: tr.Count, - Size: tr.Size, - }) - } - - // 3. 获取所有文件的大小、后缀与MIME,用于在 Go 中内存分类统计 (避免数据库中写复杂的 JSON/String 匹配逻辑) - type fileCategoryRaw struct { - Extension string `gorm:"column:extension"` - MimeType string `gorm:"column:mime_type"` - FileSize int64 `gorm:"column:file_size"` - } - var fileRaws []fileCategoryRaw - err = db.DB(ctx).Model(&model.Upload{}). - Select("extension, mime_type, file_size"). - Where("status != ?", model.UploadStatusDeleted). - Scan(&fileRaws).Error - if err != nil { - c.JSON(http.StatusOK, response.Err(err.Error())) - return - } - - catCount := make(map[string]int64) - catSize := make(map[string]int64) - categoriesList := []string{catImage, catVideo, catAudio, catDocument, catArchive, catOther} - for _, cat := range categoriesList { - catCount[cat] = 0 - catSize[cat] = 0 - } - - for _, fr := range fileRaws { - cat := getFileCategory(fr.MimeType, fr.Extension) - catCount[cat]++ - catSize[cat] += fr.FileSize - } - - categories := make([]distributionItem, 0, len(categoriesList)) - for _, cat := range categoriesList { - categories = append(categories, distributionItem{ - Name: cat, - Count: catCount[cat], - Size: catSize[cat], - }) - } - - // 4. 获取近 7 天的新增文件趋势 (在 Go 中补全没有新增记录的日期为 0) - type fileTrendRaw struct { - CreatedAt time.Time `gorm:"column:created_at"` - FileSize int64 `gorm:"column:file_size"` - } - var trendRaws []fileTrendRaw - // 7天前 00:00:00 (即 6 天前 00:00:00 至今天) now := time.Now() - startTime := time.Date(now.Year(), now.Month(), now.Day(), 0, 0, 0, 0, now.Location()).AddDate(0, 0, -6) - err = db.DB(ctx).Model(&model.Upload{}). - Select("created_at, file_size"). - Where("status != ? AND created_at >= ?", model.UploadStatusDeleted, startTime). - Scan(&trendRaws).Error - if err != nil { - c.JSON(http.StatusOK, response.Err(err.Error())) - return + trendDates := make([]string, 0, fileStatsTrendDays) + trendCountMap := make(map[string]int64, fileStatsTrendDays) + trendSizeMap := make(map[string]int64, fileStatsTrendDays) + for i := fileStatsTrendDays - 1; i >= 0; i-- { + date := now.AddDate(0, 0, -i).Format("2006-01-02") + trendDates = append(trendDates, date) + trendCountMap[date] = 0 + trendSizeMap[date] = 0 } - trendCountMap := make(map[string]int64) - trendSizeMap := make(map[string]int64) - for i := 0; i < 7; i++ { - dStr := now.AddDate(0, 0, -i).Format("2006-01-02") - trendCountMap[dStr] = 0 - trendSizeMap[dStr] = 0 + var ( + totalCount int64 + totalSize int64 + types []distributionItem + categories []distributionItem + ) + + categoriesList := []string{catImage, catVideo, catAudio, catDocument, catArchive, catOther} + categoryMap := make(map[string]distributionItem, len(categoriesList)) + for _, cat := range categoriesList { + categoryMap[cat] = distributionItem{Name: cat} } - for _, tr := range trendRaws { - dStr := tr.CreatedAt.Format("2006-01-02") - if _, exists := trendCountMap[dStr]; exists { - trendCountMap[dStr]++ - trendSizeMap[dStr] += tr.FileSize + for _, stat := range stats { + switch stat.Dimension { + case model.UploadStatDimensionTotal: + totalCount = stat.FileCount + totalSize = stat.FileSize + case model.UploadStatDimensionType: + types = append(types, distributionItem{ + Name: stat.StatKey, + Count: stat.FileCount, + Size: stat.FileSize, + }) + case model.UploadStatDimensionCategory: + if item, ok := categoryMap[stat.StatKey]; ok { + item.Count = stat.FileCount + item.Size = stat.FileSize + categoryMap[stat.StatKey] = item + } + case model.UploadStatDimensionTrend: + if _, ok := trendCountMap[stat.StatKey]; ok { + trendCountMap[stat.StatKey] = stat.FileCount + trendSizeMap[stat.StatKey] = stat.FileSize + } } } - const trendDays = 7 - trend := make([]trendItem, 0, trendDays) - for i := trendDays - 1; i >= 0; i-- { - dStr := now.AddDate(0, 0, -i).Format("2006-01-02") + categories = make([]distributionItem, 0, len(categoriesList)) + for _, cat := range categoriesList { + categories = append(categories, categoryMap[cat]) + } + + trend := make([]trendItem, 0, len(trendDates)) + for _, date := range trendDates { trend = append(trend, trendItem{ - Date: dStr, - Count: trendCountMap[dStr], - Size: trendSizeMap[dStr], + Date: date, + Count: trendCountMap[date], + Size: trendSizeMap[date], }) } c.JSON(http.StatusOK, response.OK(fileStatsResponse{ - TotalCount: summary.TotalCount, - TotalSize: summary.TotalSize, + TotalCount: totalCount, + TotalSize: totalSize, Trend: trend, Categories: categories, Types: types, @@ -231,4 +174,4 @@ func isDocumentExtension(ext string) bool { } } return false -} +} \ No newline at end of file diff --git a/internal/apps/upload/stats_counter.go b/internal/apps/upload/stats_counter.go new file mode 100644 index 00000000..3f5c28d0 --- /dev/null +++ b/internal/apps/upload/stats_counter.go @@ -0,0 +1,128 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package upload + +import ( + "context" + "time" + + "github.com/Rain-kl/Wavelet/internal/db" + "github.com/Rain-kl/Wavelet/internal/model" + "github.com/Rain-kl/Wavelet/pkg/logger" + "gorm.io/gorm" + "gorm.io/gorm/clause" +) + +// ApplyUploadStatsAdd increments incremental stats for a newly active upload record. +func ApplyUploadStatsAdd(ctx context.Context, upload *model.Upload) error { + return applyUploadStatsDelta(ctx, upload, 1) +} + +// ApplyUploadStatsRemove decrements incremental stats for a removed active upload record. +func ApplyUploadStatsRemove(ctx context.Context, upload *model.Upload) error { + return applyUploadStatsDelta(ctx, upload, -1) +} + +// RebuildUploadStats rebuilds all incremental stats from current upload records. +func RebuildUploadStats(ctx context.Context) error { + return db.DB(ctx).Transaction(func(tx *gorm.DB) error { + if err := tx.Where("1 = 1").Delete(&model.UploadStat{}).Error; err != nil { + return err + } + + var uploads []model.Upload + if err := tx.Where("status != ?", model.UploadStatusDeleted).Find(&uploads).Error; err != nil { + return err + } + + for i := range uploads { + if err := applyUploadStatsDeltaTx(tx, &uploads[i], 1); err != nil { + return err + } + } + return nil + }) +} + +func applyUploadStatsDelta(ctx context.Context, upload *model.Upload, sign int64) error { + if upload == nil || !isActiveUploadStatus(upload.Status) { + return nil + } + return db.DB(ctx).Transaction(func(tx *gorm.DB) error { + return applyUploadStatsDeltaTx(tx, upload, sign) + }) +} + +func applyUploadStatsDeltaTx(tx *gorm.DB, upload *model.Upload, sign int64) error { + if upload == nil || !isActiveUploadStatus(upload.Status) || sign == 0 { + return nil + } + + countDelta := sign + sizeDelta := sign * upload.FileSize + typeKey := upload.Type + if typeKey == "" { + typeKey = "generic" + } + + entries := []struct { + dimension string + key string + }{ + {model.UploadStatDimensionTotal, ""}, + {model.UploadStatDimensionType, typeKey}, + {model.UploadStatDimensionCategory, getFileCategory(upload.MimeType, upload.Extension)}, + {model.UploadStatDimensionTrend, upload.CreatedAt.Format("2006-01-02")}, + } + + for _, entry := range entries { + if err := upsertUploadStatDelta(tx, entry.dimension, entry.key, countDelta, sizeDelta); err != nil { + return err + } + } + return nil +} + +func upsertUploadStatDelta(tx *gorm.DB, dimension, key string, countDelta, sizeDelta int64) error { + return tx.Clauses(clause.OnConflict{ + Columns: []clause.Column{ + {Name: "dimension"}, + {Name: "stat_key"}, + }, + DoUpdates: clause.Assignments(map[string]any{ + "file_count": gorm.Expr( + "CASE WHEN w_upload_stats.file_count + ? < 0 THEN 0 ELSE w_upload_stats.file_count + ? END", + countDelta, + countDelta, + ), + "file_size": gorm.Expr( + "CASE WHEN w_upload_stats.file_size + ? < 0 THEN 0 ELSE w_upload_stats.file_size + ? END", + sizeDelta, + sizeDelta, + ), + "updated_at": time.Now(), + }), + }).Create(&model.UploadStat{ + Dimension: dimension, + StatKey: key, + FileCount: countDelta, + FileSize: sizeDelta, + }).Error +} + +func recordUploadStatsAdd(ctx context.Context, upload *model.Upload) { + if err := ApplyUploadStatsAdd(ctx, upload); err != nil { + logger.WarnF(ctx, "increment upload stats failed: %v", err) + } +} + +func recordUploadStatsRemove(ctx context.Context, upload *model.Upload) { + if err := ApplyUploadStatsRemove(ctx, upload); err != nil { + logger.WarnF(ctx, "decrement upload stats failed: %v", err) + } +} + +func isActiveUploadStatus(status model.UploadStatus) bool { + return status == model.UploadStatusPending || status == model.UploadStatusUsed +} \ No newline at end of file diff --git a/internal/apps/upload/stats_counter_test.go b/internal/apps/upload/stats_counter_test.go new file mode 100644 index 00000000..dad330a6 --- /dev/null +++ b/internal/apps/upload/stats_counter_test.go @@ -0,0 +1,72 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package upload + +import ( + "context" + "testing" + "time" + + "github.com/Rain-kl/Wavelet/internal/db" + "github.com/Rain-kl/Wavelet/internal/model" + "github.com/Rain-kl/Wavelet/internal/testhelper" +) + +func TestApplyUploadStatsAddAndRemove(t *testing.T) { + _, _, cleanup := testhelper.SetupTestEnvironment(t) + defer cleanup() + ctx := context.Background() + + upload := &model.Upload{ + ID: 42001, + FileSize: 128, + MimeType: "image/png", + Extension: "png", + Type: "avatar", + Status: model.UploadStatusUsed, + CreatedAt: time.Now(), + } + if err := ApplyUploadStatsAdd(ctx, upload); err != nil { + t.Fatalf("ApplyUploadStatsAdd returned error: %v", err) + } + + stats, err := loadUploadStats(ctx) + if err != nil { + t.Fatalf("loadUploadStats returned error: %v", err) + } + if stats.TotalCount != 1 || stats.TotalSize != 128 { + t.Fatalf("unexpected total stats: count=%d size=%d", stats.TotalCount, stats.TotalSize) + } + + if err := ApplyUploadStatsRemove(ctx, upload); err != nil { + t.Fatalf("ApplyUploadStatsRemove returned error: %v", err) + } + + stats, err = loadUploadStats(ctx) + if err != nil { + t.Fatalf("loadUploadStats after remove returned error: %v", err) + } + if stats.TotalCount != 0 || stats.TotalSize != 0 { + t.Fatalf("expected zeroed total stats, got count=%d size=%d", stats.TotalCount, stats.TotalSize) + } +} + +type uploadStatsSnapshot struct { + TotalCount int64 + TotalSize int64 +} + +func loadUploadStats(ctx context.Context) (uploadStatsSnapshot, error) { + var rows []model.UploadStat + if err := db.DB(ctx).Where("dimension = ?", model.UploadStatDimensionTotal).Find(&rows).Error; err != nil { + return uploadStatsSnapshot{}, err + } + if len(rows) == 0 { + return uploadStatsSnapshot{}, nil + } + return uploadStatsSnapshot{ + TotalCount: rows[0].FileCount, + TotalSize: rows[0].FileSize, + }, nil +} \ No newline at end of file diff --git a/internal/apps/upload/storage_migration_task.go b/internal/apps/upload/storage_migration_task.go index 2c40826c..5b089d74 100644 --- a/internal/apps/upload/storage_migration_task.go +++ b/internal/apps/upload/storage_migration_task.go @@ -361,16 +361,7 @@ func migrateSingleObject( source, err := sourceBackend.Get(ctx, obj.FilePath) if err != nil { if isNotFoundError(err) { - task.AppendLog(ctx, "警告: 源存储中物理文件不存在,标记为已删除并跳过: %s (错误: %v)", obj.FilePath, err) - if updateErr := db.DB(ctx).Model(&model.Upload{}). - Where("storage_driver = ? AND file_path = ?", sourceDriver, obj.FilePath). - Updates(map[string]any{ - "status": model.UploadStatusDeleted, - colStorageDriver: targetDriver, - }).Error; updateErr != nil { - return fmt.Errorf("update missing object %q: %w", obj.FilePath, updateErr) - } - return nil + return markMissingMigrationObjectDeleted(ctx, sourceDriver, targetDriver, obj.FilePath, err) } return fmt.Errorf("open source object %q: %w", obj.FilePath, err) } @@ -441,6 +432,35 @@ func shouldSkipMigration( return targetObj.ContentLength == obj.FileSize } +func markMissingMigrationObjectDeleted( + ctx context.Context, + sourceDriver storage.Driver, + targetDriver storage.Driver, + filePath string, + sourceErr error, +) error { + task.AppendLog(ctx, "警告: 源存储中物理文件不存在,标记为已删除并跳过: %s (错误: %v)", filePath, sourceErr) + + var affectedUploads []model.Upload + if err := db.DB(ctx). + Where("storage_driver = ? AND file_path = ? AND status != ?", sourceDriver, filePath, model.UploadStatusDeleted). + Find(&affectedUploads).Error; err != nil { + return fmt.Errorf("load missing object uploads %q: %w", filePath, err) + } + if err := db.DB(ctx).Model(&model.Upload{}). + Where("storage_driver = ? AND file_path = ?", sourceDriver, filePath). + Updates(map[string]any{ + "status": model.UploadStatusDeleted, + colStorageDriver: targetDriver, + }).Error; err != nil { + return fmt.Errorf("update missing object %q: %w", filePath, err) + } + for i := range affectedUploads { + recordUploadStatsRemove(ctx, &affectedUploads[i]) + } + return nil +} + func isNotFoundError(err error) bool { if err == nil { return false diff --git a/internal/apps/upload/storage_ops.go b/internal/apps/upload/storage_ops.go index 041340ed..369f9f43 100644 --- a/internal/apps/upload/storage_ops.go +++ b/internal/apps/upload/storage_ops.go @@ -14,15 +14,12 @@ import ( // StorageReadOnly checks if the storage system is in read-only maintenance mode. func StorageReadOnly(ctx context.Context) bool { - execution, ok, err := latestStorageMigrationExecution(ctx) - if err != nil { - logger.ErrorF(ctx, "读取存储维护状态失败: %v", err) + state := loadMigrationAccessState(ctx) + if state.loadErr != nil { + logger.ErrorF(ctx, "读取存储维护状态失败: %v", state.loadErr) return true } - if !ok { - return false - } - return execution.Status != model.TaskExecutionStatusSucceeded + return state.readOnly } func openStoredObject(ctx context.Context, upload *model.Upload) (*storage.Object, error) { @@ -54,16 +51,15 @@ func backendForStoredDriver(ctx context.Context, driver storage.Driver) (storage } func currentMigrationTargetConfig(ctx context.Context) (storage.Config, bool, error) { - execution, ok, err := latestStorageMigrationExecution(ctx) - if err != nil || !ok { - return storage.Config{}, false, err + state := loadMigrationAccessState(ctx) + if state.loadErr != nil { + return storage.Config{}, false, state.loadErr } - if execution.Status == model.TaskExecutionStatusSucceeded { + if state.targetErr != nil { + return storage.Config{}, false, state.targetErr + } + if !state.hasTarget { return storage.Config{}, false, nil } - target, err := parseMigrationTargetConfig(ctx, []byte(execution.Payload)) - if err != nil { - return storage.Config{}, false, err - } - return target, true, nil + return state.target, true, nil } diff --git a/internal/db/migrator/goose/postgres/202606170001_add_upload_composite_indexes.sql b/internal/db/migrator/goose/postgres/202606170001_add_upload_composite_indexes.sql new file mode 100644 index 00000000..7d2f0e7d --- /dev/null +++ b/internal/db/migrator/goose/postgres/202606170001_add_upload_composite_indexes.sql @@ -0,0 +1,9 @@ +-- +goose Up +CREATE INDEX IF NOT EXISTS idx_w_uploads_status_created_at ON w_uploads (status, created_at); +CREATE INDEX IF NOT EXISTS idx_w_uploads_storage_driver_status ON w_uploads (storage_driver, status); +CREATE INDEX IF NOT EXISTS idx_w_uploads_hash_file_size_status ON w_uploads (hash, file_size, status); + +-- +goose Down +DROP INDEX IF EXISTS idx_w_uploads_hash_file_size_status; +DROP INDEX IF EXISTS idx_w_uploads_storage_driver_status; +DROP INDEX IF EXISTS idx_w_uploads_status_created_at; \ No newline at end of file diff --git a/internal/db/migrator/goose/postgres/202606170002_create_upload_stats_table.sql b/internal/db/migrator/goose/postgres/202606170002_create_upload_stats_table.sql new file mode 100644 index 00000000..332932ba --- /dev/null +++ b/internal/db/migrator/goose/postgres/202606170002_create_upload_stats_table.sql @@ -0,0 +1,12 @@ +-- +goose Up +CREATE TABLE IF NOT EXISTS w_upload_stats ( + dimension VARCHAR(32) NOT NULL, + stat_key VARCHAR(64) NOT NULL DEFAULT '', + file_count BIGINT NOT NULL DEFAULT 0, + file_size BIGINT NOT NULL DEFAULT 0, + updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, + PRIMARY KEY (dimension, stat_key) +); + +-- +goose Down +DROP TABLE IF EXISTS w_upload_stats; \ No newline at end of file diff --git a/internal/db/migrator/goose/postgres/202606170003_backfill_upload_stats.sql b/internal/db/migrator/goose/postgres/202606170003_backfill_upload_stats.sql new file mode 100644 index 00000000..6e5bf473 --- /dev/null +++ b/internal/db/migrator/goose/postgres/202606170003_backfill_upload_stats.sql @@ -0,0 +1,67 @@ +-- +goose Up +INSERT INTO w_upload_stats (dimension, stat_key, file_count, file_size) +SELECT 'total', '', COUNT(*), COALESCE(SUM(file_size), 0) +FROM w_uploads +WHERE status != 'deleted' +ON CONFLICT (dimension, stat_key) DO UPDATE SET + file_count = EXCLUDED.file_count, + file_size = EXCLUDED.file_size, + updated_at = CURRENT_TIMESTAMP; + +INSERT INTO w_upload_stats (dimension, stat_key, file_count, file_size) +SELECT + 'type', + COALESCE(NULLIF(type, ''), 'generic'), + COUNT(*), + COALESCE(SUM(file_size), 0) +FROM w_uploads +WHERE status != 'deleted' +GROUP BY COALESCE(NULLIF(type, ''), 'generic') +ON CONFLICT (dimension, stat_key) DO UPDATE SET + file_count = EXCLUDED.file_count, + file_size = EXCLUDED.file_size, + updated_at = CURRENT_TIMESTAMP; + +INSERT INTO w_upload_stats (dimension, stat_key, file_count, file_size) +SELECT + 'category', + CASE + WHEN LOWER(mime_type) LIKE 'image/%' + OR LOWER(extension) IN ('jpg', 'jpeg', 'png', 'webp', 'gif') THEN '图片' + WHEN LOWER(mime_type) LIKE 'video/%' THEN '视频' + WHEN LOWER(mime_type) LIKE 'audio/%' THEN '音频' + WHEN LOWER(extension) IN ('zip', 'rar', '7z', 'tar', 'gz', 'tgz', 'bz2', 'xz') + OR LOWER(mime_type) LIKE '%zip%' + OR LOWER(mime_type) LIKE '%tar%' + OR LOWER(mime_type) LIKE '%gzip%' THEN '压缩包' + WHEN LOWER(extension) IN ('pdf', 'doc', 'docx', 'xls', 'xlsx', 'ppt', 'pptx', 'txt', 'md', 'csv', 'json', 'yaml', 'yml', 'xml') + OR LOWER(mime_type) LIKE 'text/%' + OR LOWER(mime_type) = 'application/pdf' THEN '文档' + ELSE '其他' + END, + COUNT(*), + COALESCE(SUM(file_size), 0) +FROM w_uploads +WHERE status != 'deleted' +GROUP BY 2 +ON CONFLICT (dimension, stat_key) DO UPDATE SET + file_count = EXCLUDED.file_count, + file_size = EXCLUDED.file_size, + updated_at = CURRENT_TIMESTAMP; + +INSERT INTO w_upload_stats (dimension, stat_key, file_count, file_size) +SELECT + 'trend', + TO_CHAR(created_at, 'YYYY-MM-DD'), + COUNT(*), + COALESCE(SUM(file_size), 0) +FROM w_uploads +WHERE status != 'deleted' +GROUP BY TO_CHAR(created_at, 'YYYY-MM-DD') +ON CONFLICT (dimension, stat_key) DO UPDATE SET + file_count = EXCLUDED.file_count, + file_size = EXCLUDED.file_size, + updated_at = CURRENT_TIMESTAMP; + +-- +goose Down +DELETE FROM w_upload_stats; \ No newline at end of file diff --git a/internal/db/migrator/goose/sqlite/202606170001_add_upload_composite_indexes.sql b/internal/db/migrator/goose/sqlite/202606170001_add_upload_composite_indexes.sql new file mode 100644 index 00000000..7d2f0e7d --- /dev/null +++ b/internal/db/migrator/goose/sqlite/202606170001_add_upload_composite_indexes.sql @@ -0,0 +1,9 @@ +-- +goose Up +CREATE INDEX IF NOT EXISTS idx_w_uploads_status_created_at ON w_uploads (status, created_at); +CREATE INDEX IF NOT EXISTS idx_w_uploads_storage_driver_status ON w_uploads (storage_driver, status); +CREATE INDEX IF NOT EXISTS idx_w_uploads_hash_file_size_status ON w_uploads (hash, file_size, status); + +-- +goose Down +DROP INDEX IF EXISTS idx_w_uploads_hash_file_size_status; +DROP INDEX IF EXISTS idx_w_uploads_storage_driver_status; +DROP INDEX IF EXISTS idx_w_uploads_status_created_at; \ No newline at end of file diff --git a/internal/db/migrator/goose/sqlite/202606170002_create_upload_stats_table.sql b/internal/db/migrator/goose/sqlite/202606170002_create_upload_stats_table.sql new file mode 100644 index 00000000..ab6b7b6c --- /dev/null +++ b/internal/db/migrator/goose/sqlite/202606170002_create_upload_stats_table.sql @@ -0,0 +1,12 @@ +-- +goose Up +CREATE TABLE IF NOT EXISTS w_upload_stats ( + dimension VARCHAR(32) NOT NULL, + stat_key VARCHAR(64) NOT NULL DEFAULT '', + file_count BIGINT NOT NULL DEFAULT 0, + file_size BIGINT NOT NULL DEFAULT 0, + updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + PRIMARY KEY (dimension, stat_key) +); + +-- +goose Down +DROP TABLE IF EXISTS w_upload_stats; \ No newline at end of file diff --git a/internal/db/migrator/goose/sqlite/202606170003_backfill_upload_stats.sql b/internal/db/migrator/goose/sqlite/202606170003_backfill_upload_stats.sql new file mode 100644 index 00000000..2ebbd8c4 --- /dev/null +++ b/internal/db/migrator/goose/sqlite/202606170003_backfill_upload_stats.sql @@ -0,0 +1,67 @@ +-- +goose Up +INSERT INTO w_upload_stats (dimension, stat_key, file_count, file_size) +SELECT 'total', '', COUNT(*), COALESCE(SUM(file_size), 0) +FROM w_uploads +WHERE status != 'deleted' +ON CONFLICT (dimension, stat_key) DO UPDATE SET + file_count = excluded.file_count, + file_size = excluded.file_size, + updated_at = CURRENT_TIMESTAMP; + +INSERT INTO w_upload_stats (dimension, stat_key, file_count, file_size) +SELECT + 'type', + COALESCE(NULLIF(type, ''), 'generic'), + COUNT(*), + COALESCE(SUM(file_size), 0) +FROM w_uploads +WHERE status != 'deleted' +GROUP BY COALESCE(NULLIF(type, ''), 'generic') +ON CONFLICT (dimension, stat_key) DO UPDATE SET + file_count = excluded.file_count, + file_size = excluded.file_size, + updated_at = CURRENT_TIMESTAMP; + +INSERT INTO w_upload_stats (dimension, stat_key, file_count, file_size) +SELECT + 'category', + CASE + WHEN LOWER(mime_type) LIKE 'image/%' + OR LOWER(extension) IN ('jpg', 'jpeg', 'png', 'webp', 'gif') THEN '图片' + WHEN LOWER(mime_type) LIKE 'video/%' THEN '视频' + WHEN LOWER(mime_type) LIKE 'audio/%' THEN '音频' + WHEN LOWER(extension) IN ('zip', 'rar', '7z', 'tar', 'gz', 'tgz', 'bz2', 'xz') + OR LOWER(mime_type) LIKE '%zip%' + OR LOWER(mime_type) LIKE '%tar%' + OR LOWER(mime_type) LIKE '%gzip%' THEN '压缩包' + WHEN LOWER(extension) IN ('pdf', 'doc', 'docx', 'xls', 'xlsx', 'ppt', 'pptx', 'txt', 'md', 'csv', 'json', 'yaml', 'yml', 'xml') + OR LOWER(mime_type) LIKE 'text/%' + OR LOWER(mime_type) = 'application/pdf' THEN '文档' + ELSE '其他' + END, + COUNT(*), + COALESCE(SUM(file_size), 0) +FROM w_uploads +WHERE status != 'deleted' +GROUP BY 2 +ON CONFLICT (dimension, stat_key) DO UPDATE SET + file_count = excluded.file_count, + file_size = excluded.file_size, + updated_at = CURRENT_TIMESTAMP; + +INSERT INTO w_upload_stats (dimension, stat_key, file_count, file_size) +SELECT + 'trend', + STRFTIME('%Y-%m-%d', created_at), + COUNT(*), + COALESCE(SUM(file_size), 0) +FROM w_uploads +WHERE status != 'deleted' +GROUP BY STRFTIME('%Y-%m-%d', created_at) +ON CONFLICT (dimension, stat_key) DO UPDATE SET + file_count = excluded.file_count, + file_size = excluded.file_size, + updated_at = CURRENT_TIMESTAMP; + +-- +goose Down +DELETE FROM w_upload_stats; \ No newline at end of file diff --git a/internal/model/upload_stats.go b/internal/model/upload_stats.go new file mode 100644 index 00000000..d29cce23 --- /dev/null +++ b/internal/model/upload_stats.go @@ -0,0 +1,28 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package model + +import "time" + +// Upload stats dimension keys stored in w_upload_stats.dimension. +const ( + UploadStatDimensionTotal = "total" + UploadStatDimensionType = "type" + UploadStatDimensionCategory = "category" + UploadStatDimensionTrend = "trend" +) + +// UploadStat stores incremental upload statistics keyed by dimension and stat_key. +type UploadStat struct { + Dimension string `json:"dimension" gorm:"primaryKey;size:32;not null"` + StatKey string `json:"stat_key" gorm:"primaryKey;size:64;not null;default:''"` + FileCount int64 `json:"file_count" gorm:"not null;default:0"` + FileSize int64 `json:"file_size" gorm:"not null;default:0"` + UpdatedAt time.Time `json:"updated_at" gorm:"autoUpdateTime"` +} + +// TableName returns the upload stats table name. +func (UploadStat) TableName() string { + return "w_upload_stats" +} \ No newline at end of file diff --git a/internal/router/v1/custom.go b/internal/router/v1/custom.go index 7bac464a..ba88641d 100644 --- a/internal/router/v1/custom.go +++ b/internal/router/v1/custom.go @@ -10,7 +10,7 @@ import ( ) // RegisterCustomRoutes registers custom business routes to keep routing clean and stable. -func RegisterCustomRoutes(_ *gin.RouterGroup) { +func RegisterCustomRoutes(apiV1Router *gin.RouterGroup) { customRouter := apiV1Router.Group("/custom") { customRouter.GET("/hello", custom.Hello) diff --git a/internal/storage/storage.go b/internal/storage/storage.go index 7a01070f..51de6a2f 100644 --- a/internal/storage/storage.go +++ b/internal/storage/storage.go @@ -58,7 +58,8 @@ var ( cacheMutex sync.RWMutex ) -const configInvalidationChannel = "storage:config_invalidation" +// ConfigInvalidationChannel is the Redis pub/sub channel used to evict storage caches cluster-wide. +const ConfigInvalidationChannel = "storage:config_invalidation" var pubSubOnce sync.Once @@ -75,7 +76,7 @@ func ResetCache() { // PublishCacheInvalidation broadcasts cache eviction to all nodes in the cluster via Redis. func PublishCacheInvalidation(ctx context.Context) { if db.Redis != nil { - _ = db.Redis.Publish(ctx, configInvalidationChannel, "reset").Err() + _ = db.Redis.Publish(ctx, ConfigInvalidationChannel, "reset").Err() } } @@ -85,7 +86,7 @@ func startPubSubListener() { return } go func() { - pubsub := db.Redis.Subscribe(context.Background(), configInvalidationChannel) + pubsub := db.Redis.Subscribe(context.Background(), ConfigInvalidationChannel) defer func() { _ = pubsub.Close() }() diff --git a/internal/testhelper/test_helper.go b/internal/testhelper/test_helper.go index db938c4d..18862579 100644 --- a/internal/testhelper/test_helper.go +++ b/internal/testhelper/test_helper.go @@ -42,6 +42,7 @@ func SetupTestEnvironment(t *testing.T) (*gorm.DB, *miniredis.Miniredis, func()) &model.ExternalAccount{}, &model.SystemConfig{}, &model.Upload{}, + &model.UploadStat{}, &model.TaskExecution{}, &model.Template{}, &model.AccessToken{}, diff --git a/pkg/cache/disk/cache.go b/pkg/cache/disk/cache.go index 236b1700..2dbc0590 100644 --- a/pkg/cache/disk/cache.go +++ b/pkg/cache/disk/cache.go @@ -147,6 +147,55 @@ func (c *Cache) Set(key string, value []byte, ttl time.Duration) error { // Get retrieves a key's value from the cache. func (c *Cache) Get(key string) ([]byte, error) { + c.mu.RLock() + elem, ok := c.items[key] + if !ok { + c.mu.RUnlock() + return nil, ErrCacheMiss + } + + item := elem.Value.(*cacheItem) + if !item.expiredAt.IsZero() && time.Now().After(item.expiredAt) { + c.mu.RUnlock() + return c.getAndDeleteIfExpired(key) + } + c.mu.RUnlock() + + // Read from disk outside the lock so concurrent cache hits do not serialize on I/O. + data, err := c.d.Read(key) + if err != nil { + c.mu.Lock() + defer c.mu.Unlock() + if _, stillExists := c.items[key]; stillExists { + _ = c.deleteUnlocked(key) + } + return nil, ErrCacheMiss + } + + if len(data) < headerSize { + c.mu.Lock() + defer c.mu.Unlock() + if _, stillExists := c.items[key]; stillExists { + _ = c.deleteUnlocked(key) + } + return nil, ErrCacheMiss + } + + payload := data[headerSize:] + + // Brief write lock only for LRU bookkeeping. + c.mu.Lock() + defer c.mu.Unlock() + elem, ok = c.items[key] + if !ok { + return nil, ErrCacheMiss + } + c.evictList.MoveToFront(elem) + + return payload, nil +} + +func (c *Cache) getAndDeleteIfExpired(key string) ([]byte, error) { c.mu.Lock() defer c.mu.Unlock() @@ -156,18 +205,13 @@ func (c *Cache) Get(key string) ([]byte, error) { } item := elem.Value.(*cacheItem) - - // Check expiration if !item.expiredAt.IsZero() && time.Now().After(item.expiredAt) { - // Lazily delete expired item _ = c.deleteUnlocked(key) return nil, ErrCacheMiss } - // Read from diskv data, err := c.d.Read(key) if err != nil { - // Key exists in memory but not on disk, sync state _ = c.deleteUnlocked(key) return nil, ErrCacheMiss } @@ -177,10 +221,7 @@ func (c *Cache) Get(key string) ([]byte, error) { return nil, ErrCacheMiss } - // Update LRU access order c.evictList.MoveToFront(elem) - - // Slice off the metadata header return data[headerSize:], nil }