From bc18800b583871100c9476a3861bee9332ce0534 Mon Sep 17 00:00:00 2001 From: ryan Date: Sat, 13 Jun 2026 15:15:59 +0800 Subject: [PATCH] perf(storage): add config caching and shared HTTP connection pool - Implement thread-safe local config caching with a 5s TTL check. - Reuse backend client singletons in storage.Active and storage.ForDriver. - Reset in-memory config cache on configuration saves and updates. - Create internal/httppool package to manage shared HTTP transports. - Configure WebDAV client and CDN retrieval to reuse the shared pool. - Add unit tests for httppool and storage caching behaviors. --- internal/apps/admin/system_config/routers.go | 4 + internal/httppool/httppool.go | 66 ++++++++++++++++ internal/httppool/httppool_test.go | 37 +++++++++ internal/storage/config.go | 22 +++++- internal/storage/http.go | 7 +- internal/storage/storage.go | 80 +++++++++++++++++-- internal/storage/storage_test.go | 82 ++++++++++++++++++++ internal/storage/webdav.go | 5 +- 8 files changed, 295 insertions(+), 8 deletions(-) create mode 100644 internal/httppool/httppool.go create mode 100644 internal/httppool/httppool_test.go create mode 100644 internal/storage/storage_test.go diff --git a/internal/apps/admin/system_config/routers.go b/internal/apps/admin/system_config/routers.go index d949a673..8fc96b05 100644 --- a/internal/apps/admin/system_config/routers.go +++ b/internal/apps/admin/system_config/routers.go @@ -246,6 +246,10 @@ func UpdateSystemConfig(c *gin.Context) { return } + if key == model.ConfigKeyStorageConfig { + storage.ResetCache() + } + c.JSON(http.StatusOK, util.OKNil()) } diff --git a/internal/httppool/httppool.go b/internal/httppool/httppool.go new file mode 100644 index 00000000..2297dc67 --- /dev/null +++ b/internal/httppool/httppool.go @@ -0,0 +1,66 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +// Package httppool manages shared, optimized HTTP transports to reuse TCP connections. +package httppool + +import ( + "crypto/tls" + "net" + "net/http" + "sync" + "time" + + "go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp" +) + +const ( + dialTimeout = 30 * time.Second + dialKeepAlive = 30 * time.Second + maxIdleConns = 200 + maxIdleConnsPerHost = 32 + idleConnTimeout = 90 * time.Second + tlsHandshakeTimeout = 10 * time.Second + expectContinueTimeout = 1 * time.Second + tlsSessionCacheSize = 100 +) + +var ( + defaultTransport http.RoundTripper + once sync.Once +) + +// DefaultTransport returns a globally shared, optimized http.RoundTripper +// with OTel instrumentation. It maintains a pool of idle TCP connections +// across hosts. +func DefaultTransport() http.RoundTripper { + once.Do(func() { + transport := &http.Transport{ + Proxy: http.ProxyFromEnvironment, + DialContext: (&net.Dialer{ + Timeout: dialTimeout, + KeepAlive: dialKeepAlive, + }).DialContext, + ForceAttemptHTTP2: true, + MaxIdleConns: maxIdleConns, + MaxIdleConnsPerHost: maxIdleConnsPerHost, + IdleConnTimeout: idleConnTimeout, + TLSHandshakeTimeout: tlsHandshakeTimeout, + ExpectContinueTimeout: expectContinueTimeout, + TLSClientConfig: &tls.Config{ + ClientSessionCache: tls.NewLRUClientSessionCache(tlsSessionCacheSize), + }, + } + defaultTransport = otelhttp.NewTransport(transport) + }) + return defaultTransport +} + +// NewClient returns a new http.Client that shares the global connection pool +// but has its own timeout configuration. +func NewClient(timeout time.Duration) *http.Client { + return &http.Client{ + Timeout: timeout, + Transport: DefaultTransport(), + } +} diff --git a/internal/httppool/httppool_test.go b/internal/httppool/httppool_test.go new file mode 100644 index 00000000..efc71871 --- /dev/null +++ b/internal/httppool/httppool_test.go @@ -0,0 +1,37 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package httppool + +import ( + "testing" + "time" +) + +func TestDefaultTransport(t *testing.T) { + tr1 := DefaultTransport() + if tr1 == nil { + t.Fatal("DefaultTransport() returned nil") + } + + tr2 := DefaultTransport() + if tr1 != tr2 { + t.Error("DefaultTransport() did not return a singleton instance") + } +} + +func TestNewClient(t *testing.T) { + timeout := 15 * time.Second + client := NewClient(timeout) + if client == nil { + t.Fatal("NewClient() returned nil") + } + + if client.Timeout != timeout { + t.Errorf("NewClient() timeout = %v, want %v", client.Timeout, timeout) + } + + if client.Transport != DefaultTransport() { + t.Error("NewClient() is not configured with the default transport") + } +} diff --git a/internal/storage/config.go b/internal/storage/config.go index eca85327..6dfcd1a8 100644 --- a/internal/storage/config.go +++ b/internal/storage/config.go @@ -10,6 +10,7 @@ import ( "errors" "fmt" "strings" + "time" "github.com/Rain-kl/Wavelet/internal/db" "github.com/Rain-kl/Wavelet/internal/model" @@ -87,6 +88,21 @@ func DefaultConfig() Config { // LoadConfig loads the active storage configuration. func LoadConfig(ctx context.Context) (Config, error) { + cacheMutex.RLock() + isCacheValid := time.Since(lastChecked) < 5*time.Second && activeConfigJSON != "" + configJSON := activeConfigJSON + cacheMutex.RUnlock() + + if isCacheValid { + cfg := DefaultConfig() + if strings.TrimSpace(configJSON) != "" { + if err := json.Unmarshal([]byte(configJSON), &cfg); err != nil { + return Config{}, fmt.Errorf("parse storage config from cache: %w", err) + } + } + return cfg, nil + } + return loadConfigByKey(ctx, model.ConfigKeyStorageConfig, DefaultConfig()) } @@ -163,9 +179,13 @@ func SaveActiveConfig(ctx context.Context, cfg Config) error { } func saveSystemConfig(ctx context.Context, key string, value any, description string) error { - return db.DB(ctx).Transaction(func(tx *gorm.DB) error { + err := db.DB(ctx).Transaction(func(tx *gorm.DB) error { return upsertSystemConfig(ctx, tx, key, value, description) }) + if err == nil && key == model.ConfigKeyStorageConfig { + ResetCache() + } + return err } func upsertSystemConfig(ctx context.Context, tx *gorm.DB, key string, value any, description string) error { diff --git a/internal/storage/http.go b/internal/storage/http.go index add2f472..d5226eef 100644 --- a/internal/storage/http.go +++ b/internal/storage/http.go @@ -8,6 +8,9 @@ import ( "fmt" "net/http" "net/url" + "time" + + "github.com/Rain-kl/Wavelet/internal/httppool" ) func getHTTPObject(ctx context.Context, baseURL, key string) (*Object, error) { @@ -19,7 +22,9 @@ func getHTTPObject(ctx context.Context, baseURL, key string) (*Object, error) { if err != nil { return nil, fmt.Errorf("create CDN request: %w", err) } - response, err := http.DefaultClient.Do(request) + const cdnRequestTimeout = 30 * time.Second + client := httppool.NewClient(cdnRequestTimeout) + response, err := client.Do(request) if err != nil { return nil, fmt.Errorf("get CDN object: %w", err) } diff --git a/internal/storage/storage.go b/internal/storage/storage.go index 65492e60..05d44cfc 100644 --- a/internal/storage/storage.go +++ b/internal/storage/storage.go @@ -5,8 +5,16 @@ package storage import ( "context" + "encoding/json" + "errors" "fmt" "io" + "strings" + "sync" + "time" + + "github.com/Rain-kl/Wavelet/internal/model" + "gorm.io/gorm" ) const ( @@ -35,26 +43,88 @@ var ( // IsEnabledFunc preserves the legacy S3 test hook while tests migrate to backend injection. IsEnabledFunc = func() bool { return false } mockBackend Backend + + activeBackend Backend + activeDriver Driver + activeConfigJSON string + lastChecked time.Time + cacheMutex sync.RWMutex ) -// Active returns the configured active driver and backend. +// ResetCache clears the local cache for storage configuration and client singletons. +func ResetCache() { + cacheMutex.Lock() + defer cacheMutex.Unlock() + activeBackend = nil + activeDriver = "" + activeConfigJSON = "" + lastChecked = time.Time{} +} + +// Active returns the configured active driver and backend, using an in-memory cache with 5s TTL. func Active(ctx context.Context) (Driver, Backend, error) { if IsEnabledFunc() && mockBackend != nil { return DriverS3, mockBackend, nil } - cfg, err := LoadConfig(ctx) + + cacheMutex.RLock() + isCacheValid := time.Since(lastChecked) < 5*time.Second && activeBackend != nil + if isCacheValid { + d, b := activeDriver, activeBackend + cacheMutex.RUnlock() + return d, b, nil + } + cacheMutex.RUnlock() + + cacheMutex.Lock() + defer cacheMutex.Unlock() + + // Double-check under write lock + if time.Since(lastChecked) < 5*time.Second && activeBackend != nil { + return activeDriver, activeBackend, nil + } + + var sc model.SystemConfig + err := sc.GetByKey(ctx, model.ConfigKeyStorageConfig) + if err != nil && !errors.Is(err, gorm.ErrRecordNotFound) { + return "", nil, err + } + + lastChecked = time.Now() + + // Reuse existing backend client singleton if configuration JSON matches + if sc.Value == activeConfigJSON && activeBackend != nil { + return activeDriver, activeBackend, nil + } + + cfg := DefaultConfig() + if strings.TrimSpace(sc.Value) != "" { + if err := json.Unmarshal([]byte(sc.Value), &cfg); err != nil { + return "", nil, fmt.Errorf("parse storage config: %w", err) + } + } + + backend, err := NewBackend(ctx, cfg, cfg.Driver) if err != nil { return "", nil, err } - backend, err := NewBackend(ctx, cfg, cfg.Driver) - return cfg.Driver, backend, err + + activeDriver = cfg.Driver + activeBackend = backend + activeConfigJSON = sc.Value + + return activeDriver, activeBackend, nil } -// ForDriver returns the active or pending backend for an upload record. +// ForDriver returns the active or pending backend for an upload record, reusing the active singleton if matched. func ForDriver(ctx context.Context, driver Driver) (Backend, error) { if driver == DriverS3 && mockBackend != nil { return mockBackend, nil } + activeDrv, activeBnd, err := Active(ctx) + if err == nil && activeDrv == driver { + return activeBnd, nil + } cfg, err := LoadConfig(ctx) if err != nil { return nil, err diff --git a/internal/storage/storage_test.go b/internal/storage/storage_test.go new file mode 100644 index 00000000..760a27e9 --- /dev/null +++ b/internal/storage/storage_test.go @@ -0,0 +1,82 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package storage + +import ( + "context" + "encoding/json" + "io" + "testing" + "time" +) + +func TestStorageCache(t *testing.T) { + // 1. Reset cache + ResetCache() + + if activeConfigJSON != "" || activeDriver != "" || activeBackend != nil || !lastChecked.IsZero() { + t.Fatal("ResetCache did not clear cache variables") + } + + // 2. Set up cache manually + expectedConfig := Config{ + Driver: DriverLocal, + Local: LocalConfig{Root: "/tmp/wavelet-test"}, + } + cfgJSON, err := json.Marshal(expectedConfig) + if err != nil { + t.Fatalf("Marshal config failed: %v", err) + } + + cacheMutex.Lock() + activeConfigJSON = string(cfgJSON) + lastChecked = time.Now() + cacheMutex.Unlock() + + // 3. Call LoadConfig and verify it loads from cache (doesn't hit database, which would fail/panic because DB is not initialized) + ctx := context.Background() + loadedCfg, err := LoadConfig(ctx) + if err != nil { + t.Fatalf("LoadConfig failed: %v", err) + } + + if loadedCfg.Driver != expectedConfig.Driver || loadedCfg.Local.Root != expectedConfig.Local.Root { + t.Errorf("Loaded config %+v, expected %+v", loadedCfg, expectedConfig) + } + + // 4. Test Active() returns cached driver and backend + mockBnd := &functionBackend{ + put: func(context.Context, string, io.Reader, int64, string) error { return nil }, + get: func(context.Context, string) (*Object, error) { return nil, nil }, + delete: func(context.Context, string) error { return nil }, + } + + cacheMutex.Lock() + activeBackend = mockBnd + activeDriver = DriverLocal + cacheMutex.Unlock() + + drv, bnd, err := Active(ctx) + if err != nil { + t.Fatalf("Active failed: %v", err) + } + if drv != DriverLocal || bnd != mockBnd { + t.Errorf("Active returned driver %v, backend %v; expected %v, %v", drv, bnd, DriverLocal, mockBnd) + } + + // 5. Test ForDriver returns cached backend + bnd2, err := ForDriver(ctx, DriverLocal) + if err != nil { + t.Fatalf("ForDriver failed: %v", err) + } + if bnd2 != mockBnd { + t.Errorf("ForDriver returned backend %v, expected %v", bnd2, mockBnd) + } + + // 6. Test ResetCache again + ResetCache() + if activeConfigJSON != "" || activeDriver != "" || activeBackend != nil || !lastChecked.IsZero() { + t.Fatal("ResetCache did not clear cache variables after setting them") + } +} diff --git a/internal/storage/webdav.go b/internal/storage/webdav.go index f8619452..9efecec9 100644 --- a/internal/storage/webdav.go +++ b/internal/storage/webdav.go @@ -10,6 +10,7 @@ import ( "path" "strings" + "github.com/Rain-kl/Wavelet/internal/httppool" "github.com/studio-b12/gowebdav" ) @@ -19,8 +20,10 @@ type webDAVBackend struct { } func newWebDAVBackend(cfg WebDAVConfig) (*webDAVBackend, error) { + client := gowebdav.NewClient(strings.TrimRight(cfg.Endpoint, "/"), cfg.Username, cfg.Password) + client.SetTransport(httppool.DefaultTransport()) return &webDAVBackend{ - client: gowebdav.NewClient(strings.TrimRight(cfg.Endpoint, "/"), cfg.Username, cfg.Password), + client: client, basePath: strings.Trim(cfg.BasePath, "/"), }, nil }