mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-02 23:06:36 +08:00
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.
This commit is contained in:
@@ -246,6 +246,10 @@ func UpdateSystemConfig(c *gin.Context) {
|
||||
return
|
||||
}
|
||||
|
||||
if key == model.ConfigKeyStorageConfig {
|
||||
storage.ResetCache()
|
||||
}
|
||||
|
||||
c.JSON(http.StatusOK, util.OKNil())
|
||||
}
|
||||
|
||||
|
||||
@@ -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(),
|
||||
}
|
||||
}
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user