package storage import ( "context" "fmt" "io" "sync" ) // Manager 存储引擎管理器:实现 Storage 全接口并支持运行时热切换。 // // 26.9 需求:管理后台可设置存储类型(local|s3|webdav)与各引擎参数, // 保存后无需重启即生效。设计要点: // - 读写/保存类操作全部委托到"当前引擎"(原子指针,无锁热路径); // - Switch 先构建并健康检查新引擎,成功才替换指针,失败保持原引擎; // - EngineOf 按名字取引擎实例(带缓存),供"按文件归属引擎取回旧文件"使用; // - 管理端修改引擎参数后调用 Invalidate 使对应实例缓存失效,下次构建生效。 type Manager struct { // build 构建指定引擎实例(由装配方注入:内部刷新全局 EngineOptions 后走工厂)。 build func(name string) (Storage, error) mu sync.RWMutex current Storage curName string cache map[string]Storage } // validEngines 合法引擎名(与 FCB_STORAGE_ENGINE 枚举一致)。 var validEngines = map[string]bool{"local": true, "s3": true, "webdav": true} // ValidEngine 校验引擎名是否合法。 func ValidEngine(name string) bool { return validEngines[name] } // NewManager 创建管理器:current 为启动时已构建的引擎(主装配流已做过健康检查)。 // build 注入构建函数(管理端切换/参数变更时使用,内部须串行——Manager 已加锁)。 func NewManager(name string, current Storage, build func(name string) (Storage, error)) *Manager { return &Manager{ build: build, current: current, curName: name, cache: map[string]Storage{name: current}, } } // —— Storage 接口委托(全部走当前引擎)—— // SaveFile 流式保存文件(委托当前引擎)。 func (m *Manager) SaveFile(ctx context.Context, r io.Reader, savePath string) (int64, error) { return m.current.SaveFile(ctx, r, savePath) } // DeleteFile 删除文件(委托当前引擎)。 func (m *Manager) DeleteFile(ctx context.Context, savePath string) error { return m.current.DeleteFile(ctx, savePath) } // Open 打开文件流(委托当前引擎;旧文件由 API 层先经 EngineOf 按归属引擎取)。 func (m *Manager) Open(ctx context.Context, savePath string, rng *Range) (*Download, error) { return m.current.Open(ctx, savePath, rng) } // Stat 文件元信息(委托当前引擎)。 func (m *Manager) Stat(ctx context.Context, savePath string) (*FileMeta, error) { return m.current.Stat(ctx, savePath) } // SaveChunk 保存分片(委托当前引擎)。 func (m *Manager) SaveChunk(ctx context.Context, uploadID string, chunkIndex int, r io.Reader, savePath string) (int64, error) { return m.current.SaveChunk(ctx, uploadID, chunkIndex, r, savePath) } // MergeChunks 合并分片(委托当前引擎)。 func (m *Manager) MergeChunks(ctx context.Context, uploadID string, total int, verifyHash func(index int) (string, error), savePath string) (int64, string, error) { return m.current.MergeChunks(ctx, uploadID, total, verifyHash, savePath) } // CleanChunks 清理分片临时区(委托当前引擎)。 func (m *Manager) CleanChunks(ctx context.Context, uploadID string, savePath string) error { return m.current.CleanChunks(ctx, uploadID, savePath) } // FileExists 文件存在性(委托当前引擎)。 func (m *Manager) FileExists(ctx context.Context, savePath string) (bool, error) { return m.current.FileExists(ctx, savePath) } // HeadMeta 元信息与头部字节(委托当前引擎;供直传 confirm 校验)。 func (m *Manager) HeadMeta(ctx context.Context, savePath string, headBytes int64) (*FileMeta, []byte, error) { return m.current.HeadMeta(ctx, savePath, headBytes) } // PresignGetURL 限时直链下载(委托当前引擎)。 func (m *Manager) PresignGetURL(ctx context.Context, savePath string, expires int64) (string, error) { return m.current.PresignGetURL(ctx, savePath, expires) } // PresignPutURL 限时直传(委托当前引擎)。 func (m *Manager) PresignPutURL(ctx context.Context, savePath string, expires int64) (string, error) { return m.current.PresignPutURL(ctx, savePath, expires) } // HealthCheck 健康检查(委托当前引擎)。 func (m *Manager) HealthCheck(ctx context.Context) error { return m.current.HealthCheck(ctx) } // —— 管理面:当前引擎名 / 按名取实例 / 热切换 / 缓存失效 —— // CurrentName 当前引擎名(管理端展示与文件归属戳用;并发安全)。 func (m *Manager) CurrentName() string { m.mu.RLock() defer m.mu.RUnlock() return m.curName } // Current 当前引擎实例。 func (m *Manager) Current() Storage { m.mu.RLock() defer m.mu.RUnlock() return m.current } // EngineOf 按名字取引擎实例(带缓存;用于按文件归属引擎取回旧文件)。 // 实例不存在时现场构建(不健康检查——读旧文件尽力而为,构建失败即报错)。 func (m *Manager) EngineOf(name string) (Storage, error) { if !validEngines[name] { return nil, fmt.Errorf("storage: 未知存储引擎 %q(仅支持 local|s3|webdav)", name) } m.mu.RLock() if s, ok := m.cache[name]; ok { m.mu.RUnlock() return s, nil } m.mu.RUnlock() m.mu.Lock() defer m.mu.Unlock() // 双检:拿写锁期间可能已被并发构建 if s, ok := m.cache[name]; ok { return s, nil } s, err := m.build(name) if err != nil { return nil, err } m.cache[name] = s return s, nil } // Switch 热切换当前引擎:构建新实例 → 健康检查 → 成功才替换指针。 // 任一步失败返回错误且当前引擎保持不变(管理端 503 上报)。 func (m *Manager) Switch(name string) (Storage, error) { if !validEngines[name] { return nil, fmt.Errorf("storage: 未知存储引擎 %q(仅支持 local|s3|webdav)", name) } m.mu.Lock() defer m.mu.Unlock() if m.curName == name { return m.current, nil } s, ok := m.cache[name] if !ok { var err error s, err = m.build(name) if err != nil { return nil, fmt.Errorf("storage: 构建 %s 引擎失败: %w", name, err) } } if err := s.HealthCheck(context.Background()); err != nil { return nil, fmt.Errorf("storage: %s 引擎健康检查未通过: %w", name, err) } m.current = s m.curName = name m.cache[name] = s return s, nil } // Invalidate 引擎参数变更后使对应实例缓存失效(下次 EngineOf/Switch 重建生效)。 // 当前引擎不受影响(运行中实例继续服务,直到显式 Switch)。 func (m *Manager) Invalidate(name string) { m.mu.Lock() defer m.mu.Unlock() if name == m.curName { return // 当前引擎实例仍被热路径使用,不重建;参数生效由下一次 Switch 完成 } delete(m.cache, name) }