fix: guard storage registry maps

Protect Storages and UserStorages with mutexes and expose read accessors.
This commit is contained in:
krau
2026-08-16 20:28:21 +08:00
parent aa25eb1510
commit b72dd67be9
4 changed files with 69 additions and 13 deletions
+3 -3
View File
@@ -39,7 +39,7 @@ func NewTaskFactory(ctx context.Context) *TaskFactory {
// CreateTask 创建任务 // CreateTask 创建任务
func (f *TaskFactory) CreateTask(req *CreateTaskRequest) (*CreateTaskResponse, error) { func (f *TaskFactory) CreateTask(req *CreateTaskRequest) (*CreateTaskResponse, error) {
// 验证存储 // 验证存储
stor, ok := storage.Storages[req.Storage] stor, ok := storage.GetStorage(req.Storage)
if !ok { if !ok {
return nil, fmt.Errorf("storage not found: %s", req.Storage) return nil, fmt.Errorf("storage not found: %s", req.Storage)
} }
@@ -327,12 +327,12 @@ func (f *TaskFactory) createTransferTask(taskID string, createdAt time.Time, req
} }
// 验证源存储和目标存储 // 验证源存储和目标存储
sourceStor, ok := storage.Storages[params.SourceStorage] sourceStor, ok := storage.GetStorage(params.SourceStorage)
if !ok { if !ok {
return nil, fmt.Errorf("source storage not found: %s", params.SourceStorage) return nil, fmt.Errorf("source storage not found: %s", params.SourceStorage)
} }
targetStor, ok := storage.Storages[params.TargetStorage] targetStor, ok := storage.GetStorage(params.TargetStorage)
if !ok { if !ok {
return nil, fmt.Errorf("target storage not found: %s", params.TargetStorage) return nil, fmt.Errorf("target storage not found: %s", params.TargetStorage)
} }
+3 -2
View File
@@ -135,8 +135,9 @@ func (h *Handlers) ListStoragesHandler(w http.ResponseWriter, r *http.Request) {
return return
} }
storages := make([]StorageInfo, 0, len(storage.Storages)) all := storage.AllStorages()
for name, stor := range storage.Storages { storages := make([]StorageInfo, 0, len(all))
for name, stor := range all {
storages = append(storages, StorageInfo{ storages = append(storages, StorageInfo{
Name: name, Name: name,
Type: string(stor.Type()), Type: string(stor.Type()),
+54 -6
View File
@@ -3,13 +3,42 @@ package storage
import ( import (
"context" "context"
"fmt" "fmt"
"sync"
"github.com/charmbracelet/log" "github.com/charmbracelet/log"
"github.com/krau/SaveAny-Bot/config" "github.com/krau/SaveAny-Bot/config"
storenum "github.com/krau/SaveAny-Bot/pkg/enums/storage" storenum "github.com/krau/SaveAny-Bot/pkg/enums/storage"
) )
var UserStorages = make(map[int64][]Storage) var (
storageMu sync.RWMutex
// Storages maps storage names to initialized storage instances.
Storages = make(map[string]Storage)
userStoragesMu sync.RWMutex
// UserStorages maps user IDs to their available storage instances.
UserStorages = make(map[int64][]Storage)
)
// GetStorage returns the initialized storage instance for name, without
// creating one on demand.
func GetStorage(name string) (Storage, bool) {
storageMu.RLock()
defer storageMu.RUnlock()
s, ok := Storages[name]
return s, ok
}
// AllStorages returns a snapshot copy of all initialized storages.
func AllStorages() map[string]Storage {
storageMu.RLock()
defer storageMu.RUnlock()
out := make(map[string]Storage, len(Storages))
for name, s := range Storages {
out[name] = s
}
return out
}
// GetStorageByName returns storage by name from cache or creates new one // GetStorageByName returns storage by name from cache or creates new one
// It should NOT be used to get storage for user, use GetStorageByUserIDAndName instead // It should NOT be used to get storage for user, use GetStorageByUserIDAndName instead
@@ -18,19 +47,28 @@ func GetStorageByName(ctx context.Context, name string) (Storage, error) {
return nil, ErrStorageNameEmpty return nil, ErrStorageNameEmpty
} }
storageMu.RLock()
storage, ok := Storages[name] storage, ok := Storages[name]
storageMu.RUnlock()
if ok { if ok {
return storage, nil return storage, nil
} }
cfg := config.C().GetStorageByName(name) cfg := config.C().GetStorageByName(name)
if cfg == nil { if cfg == nil {
return nil, fmt.Errorf("未找到存储 %s", name) return nil, fmt.Errorf("storage %s not found", name)
} }
// NewStorage 可能耗时 (网络初始化), 在锁外构造
storage, err := NewStorage(ctx, cfg) storage, err := NewStorage(ctx, cfg)
if err != nil { if err != nil {
return nil, err return nil, err
} }
storageMu.Lock()
defer storageMu.Unlock()
// 双写竞态: 并发调用时另一 goroutine 可能已构造并写入
if existing, ok := Storages[name]; ok {
return existing, nil
}
Storages[name] = storage Storages[name] = storage
return storage, nil return storage, nil
} }
@@ -52,8 +90,11 @@ func GetUserStorages(ctx context.Context, chatID int64) []Storage {
if chatID <= 0 { if chatID <= 0 {
return nil return nil
} }
if storages, ok := UserStorages[chatID]; ok { userStoragesMu.RLock()
return storages cached, ok := UserStorages[chatID]
userStoragesMu.RUnlock()
if ok {
return cached
} }
var storages []Storage var storages []Storage
for _, name := range config.C().GetStorageNamesByUserID(chatID) { for _, name := range config.C().GetStorageNamesByUserID(chatID) {
@@ -75,9 +116,16 @@ func LoadStorages(ctx context.Context) {
logger.Errorf("failed to load storage %s: %v", storage.GetName(), err) logger.Errorf("failed to load storage %s: %v", storage.GetName(), err)
} }
} }
logger.Infof("successfully loaded %d storages", len(Storages)) storageMu.RLock()
loaded := len(Storages)
storageMu.RUnlock()
logger.Infof("successfully loaded %d storages", loaded)
for user := range config.C().GetUsersID() { for user := range config.C().GetUsersID() {
UserStorages[int64(user)] = GetUserStorages(ctx, int64(user)) uid := int64(user)
storages := GetUserStorages(ctx, uid)
userStoragesMu.Lock()
UserStorages[uid] = storages
userStoragesMu.Unlock()
} }
} }
+9 -2
View File
@@ -75,11 +75,18 @@ type StorageReadable interface {
OpenFile(ctx context.Context, filePath string) (io.ReadCloser, int64, error) OpenFile(ctx context.Context, filePath string) (io.ReadCloser, int64, error)
} }
var Storages = make(map[string]Storage)
var _ StorageProgressSaver = (*telegram.Telegram)(nil) var _ StorageProgressSaver = (*telegram.Telegram)(nil)
var _ StorageBatchProgressSaver = (*telegram.Telegram)(nil) var _ StorageBatchProgressSaver = (*telegram.Telegram)(nil)
var _ StorageListable = (*alist.Alist)(nil)
var _ StorageReadable = (*alist.Alist)(nil)
var _ StorageListable = (*local.Local)(nil)
var _ StorageReadable = (*local.Local)(nil)
var _ StorageListable = (*rclone.Rclone)(nil)
var _ StorageReadable = (*rclone.Rclone)(nil)
var _ StorageListable = (*webdav.Webdav)(nil)
var _ StorageReadable = (*webdav.Webdav)(nil)
type StorageConstructor func() Storage type StorageConstructor func() Storage
var storageConstructors = map[storenum.StorageType]StorageConstructor{ var storageConstructors = map[storenum.StorageType]StorageConstructor{