mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-09-30 22:26:38 +08:00
299ac30ee4
- Purify core micro-kernel by removing context hardcoded helpers and reverse dependencies - Eliminate init() side effects in infra plugins with reversible lifecycle disposal - Completely isolate plugins by removing cross-plugin imports and using core/contracts - Introduce TaskService and RiskControlService contracts for unified cross-plugin APIs - Regenerate Swagger documentation and update developer guide matrix - Achieve 0 violations in check_cordis_architecture.sh and 100% test pass
213 lines
5.5 KiB
Go
213 lines
5.5 KiB
Go
// Copyright 2026 Arctel.net
|
|
// SPDX-License-Identifier: Apache-2.0
|
|
|
|
package objectstore
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
|
|
"gorm.io/gorm"
|
|
)
|
|
|
|
const (
|
|
defaultContentType = "application/octet-stream"
|
|
storageDirPerm = 0o750
|
|
storageFilePerm = 0o600
|
|
)
|
|
|
|
// Object describes a readable stored object.
|
|
type Object struct {
|
|
Key string
|
|
CachePath string
|
|
Body io.ReadCloser
|
|
ContentLength int64
|
|
ContentType string
|
|
}
|
|
|
|
// PutResult describes the result of a successful Put operation.
|
|
type PutResult struct {
|
|
Key string
|
|
Bucket string
|
|
}
|
|
|
|
// Backend defines storage operations used by the upload domain.
|
|
type Backend interface {
|
|
Put(ctx context.Context, key string, body io.Reader, size int64, contentType string) (PutResult, error)
|
|
Get(ctx context.Context, key string) (*Object, error)
|
|
Delete(ctx context.Context, key string) error
|
|
Test(ctx context.Context) error
|
|
}
|
|
|
|
var (
|
|
cacheMutex sync.RWMutex
|
|
activeDriver Driver
|
|
activeBackend Backend
|
|
activeConfigJSON string
|
|
lastChecked time.Time
|
|
pubSubOnce sync.Once
|
|
|
|
mockBackend Backend
|
|
|
|
// IsEnabledFunc controls whether mock/in-memory backend is activated in tests.
|
|
IsEnabledFunc = func() bool { return false }
|
|
)
|
|
|
|
// ConfigInvalidationChannel is the Redis pub/sub channel used to evict storage caches cluster-wide.
|
|
const ConfigInvalidationChannel = "storage:config_invalidation"
|
|
|
|
// SetMockBackend forces an in-memory/mock backend for testing.
|
|
func SetMockBackend(b Backend) {
|
|
mockBackend = b
|
|
}
|
|
|
|
// ResetCache clears cached driver and backend instances.
|
|
func ResetCache() {
|
|
cacheMutex.Lock()
|
|
defer cacheMutex.Unlock()
|
|
activeBackend = nil
|
|
activeDriver = ""
|
|
activeConfigJSON = ""
|
|
lastChecked = time.Time{}
|
|
}
|
|
|
|
// PublishCacheInvalidation broadcasts cache eviction to all nodes in the cluster via Redis.
|
|
func PublishCacheInvalidation(ctx context.Context) {
|
|
if cache := getCache(ctx); cache != nil {
|
|
_ = cache.Invalidate(ctx, ConfigInvalidationChannel)
|
|
}
|
|
ResetCache()
|
|
}
|
|
|
|
// startPubSubListener starts the background subscriber for cache invalidations.
|
|
func startPubSubListener() {
|
|
}
|
|
|
|
// 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
|
|
}
|
|
|
|
pubSubOnce.Do(startPubSubListener)
|
|
|
|
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 val string
|
|
db := getDB(ctx)
|
|
if db != nil {
|
|
err := db.Table("w_system_configs").Where("key = ?", "storage_config").Pluck("value", &val).Error
|
|
if err != nil && !errors.Is(err, gorm.ErrRecordNotFound) {
|
|
return "", nil, err
|
|
}
|
|
}
|
|
sc := struct{ Value string }{Value: val}
|
|
|
|
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
|
|
}
|
|
|
|
activeDriver = cfg.Driver
|
|
activeBackend = backend
|
|
activeConfigJSON = sc.Value
|
|
|
|
return activeDriver, activeBackend, nil
|
|
}
|
|
|
|
type functionBackend struct {
|
|
put func(context.Context, string, io.Reader, int64, string) error
|
|
get func(context.Context, string) (*Object, error)
|
|
delete func(context.Context, string) error
|
|
}
|
|
|
|
func (b *functionBackend) Put(ctx context.Context, key string, body io.Reader, size int64, contentType string) (PutResult, error) {
|
|
if err := b.put(ctx, key, body, size, contentType); err != nil {
|
|
return PutResult{}, err
|
|
}
|
|
return PutResult{Key: key}, nil
|
|
}
|
|
|
|
func (b *functionBackend) Get(ctx context.Context, key string) (*Object, error) {
|
|
return b.get(ctx, key)
|
|
}
|
|
|
|
func (b *functionBackend) Delete(ctx context.Context, key string) error {
|
|
return b.delete(ctx, key)
|
|
}
|
|
|
|
func (b *functionBackend) Test(context.Context) error {
|
|
return nil
|
|
}
|
|
|
|
// MockStorage replaces object operations for package tests and returns a restore function.
|
|
func MockStorage(
|
|
put func(context.Context, string, io.Reader, int64, string) error,
|
|
get func(context.Context, string) (*Object, error),
|
|
deleteObject func(context.Context, string) error,
|
|
) func() {
|
|
previous := mockBackend
|
|
mockBackend = &functionBackend{put: put, get: get, delete: deleteObject}
|
|
return func() {
|
|
mockBackend = previous
|
|
}
|
|
}
|
|
|
|
// NewBackend constructs a concrete backend from configuration.
|
|
func NewBackend(ctx context.Context, cfg Config, driver Driver) (Backend, error) {
|
|
if driver == DriverS3 && mockBackend != nil {
|
|
return mockBackend, nil
|
|
}
|
|
switch driver {
|
|
case DriverLocal:
|
|
return newLocalBackend(cfg.Local)
|
|
case DriverS3:
|
|
return newS3Backend(ctx, cfg.S3)
|
|
case DriverR2:
|
|
return newR2Backend(ctx, cfg.R2)
|
|
case DriverMinIO:
|
|
return newS3Backend(ctx, cfg.MinIO)
|
|
case DriverOSS:
|
|
return newOSSBackend(cfg.OSS)
|
|
case DriverWebDAV:
|
|
return newWebDAVBackend(cfg.WebDAV)
|
|
default:
|
|
return nil, fmt.Errorf("unsupported storage driver %q", driver)
|
|
}
|
|
}
|