mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-09-29 05:56:38 +08:00
fix(server): route system config and lookups through Wavelet contracts
This commit is contained in:
@@ -10,16 +10,12 @@ import (
|
||||
"io"
|
||||
"os"
|
||||
"strings"
|
||||
"sync"
|
||||
|
||||
"Wavelet/core/contracts"
|
||||
waveletupload "Wavelet/plugins/domain/upload"
|
||||
"Wavelet/plugins/domain/upload/cache"
|
||||
"Wavelet/plugins/domain/upload/models"
|
||||
uploadrepo "Wavelet/plugins/domain/upload/repository"
|
||||
uploadstats "Wavelet/plugins/domain/upload/stats"
|
||||
uploadstorage "Wavelet/plugins/domain/upload/storage"
|
||||
"Wavelet/plugins/infra/database"
|
||||
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
// ReservedPagesDeploymentType is managed exclusively by the Pages domain.
|
||||
@@ -43,6 +39,24 @@ type (
|
||||
IngestPolicy = waveletupload.IngestPolicy
|
||||
)
|
||||
|
||||
var (
|
||||
storageMu sync.RWMutex
|
||||
storageSvc contracts.StorageService
|
||||
)
|
||||
|
||||
// SetStorage injects the platform StorageService used to open stored objects.
|
||||
func SetStorage(s contracts.StorageService) {
|
||||
storageMu.Lock()
|
||||
defer storageMu.Unlock()
|
||||
storageSvc = s
|
||||
}
|
||||
|
||||
func currentStorage() contracts.StorageService {
|
||||
storageMu.RLock()
|
||||
defer storageMu.RUnlock()
|
||||
return storageSvc
|
||||
}
|
||||
|
||||
// IngestFromLocalPath ingests a local regular file through Wavelet upload ingest.
|
||||
func IngestFromLocalPath(ctx context.Context, localPath string, req IngestRequest) (IngestResult, error) {
|
||||
localPath = strings.TrimSpace(localPath)
|
||||
@@ -69,36 +83,8 @@ func IngestFromLocalPath(ctx context.Context, localPath string, req IngestReques
|
||||
return waveletupload.Ingest(ctx, req)
|
||||
}
|
||||
|
||||
// RemoveLockedTx performs the idempotent active-to-deleted transition for a row
|
||||
// that the caller has already locked in its surrounding transaction.
|
||||
func RemoveLockedTx(tx *gorm.DB, upload *models.Upload) (bool, error) {
|
||||
if upload == nil {
|
||||
return false, nil
|
||||
}
|
||||
if upload.Status == models.UploadStatusDeleted {
|
||||
return false, nil
|
||||
}
|
||||
snapshot := *upload
|
||||
if err := uploadrepo.SoftDeleteUploadTx(tx, upload); err != nil {
|
||||
return false, err
|
||||
}
|
||||
if err := uploadstats.ApplyUploadStatsDeltaTx(tx, &snapshot, -1); err != nil {
|
||||
return false, err
|
||||
}
|
||||
upload.Status = models.UploadStatusDeleted
|
||||
return true, nil
|
||||
}
|
||||
|
||||
// InvalidateUploadMetaCache evicts cached upload metadata.
|
||||
func InvalidateUploadMetaCache(ctx context.Context, id uint64) {
|
||||
cache.EvictUploadMeta(ctx, id)
|
||||
}
|
||||
|
||||
// GetActiveUpload loads an active (non-deleted) upload by ID.
|
||||
func GetActiveUpload(ctx context.Context, id uint64) (models.Upload, error) {
|
||||
if u, err := uploadrepo.GetActiveUploadByID(ctx, id); err == nil {
|
||||
return u, nil
|
||||
}
|
||||
conn := database.DB(ctx)
|
||||
if conn == nil {
|
||||
return models.Upload{}, errors.New("database not initialized")
|
||||
@@ -116,13 +102,17 @@ type OpenedUploadObject struct {
|
||||
ContentLength int64
|
||||
}
|
||||
|
||||
// OpenStoredUpload opens the stored object for an active upload.
|
||||
// OpenStoredUpload opens the stored object for an active upload via StorageService.
|
||||
func OpenStoredUpload(ctx context.Context, id uint64) (*OpenedUploadObject, error) {
|
||||
upload, err := GetActiveUpload(ctx, id)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
obj, err := uploadstorage.OpenStoredObject(ctx, &upload)
|
||||
svc := currentStorage()
|
||||
if svc == nil {
|
||||
return nil, errors.New("storage service not available")
|
||||
}
|
||||
obj, err := svc.Get(ctx, upload.FilePath)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
@@ -159,5 +149,5 @@ func ResolveLocalFile(_ context.Context, req LocalFileCandidateRequest) (string,
|
||||
|
||||
// RebuildUploadStats rebuilds aggregate upload stats.
|
||||
func RebuildUploadStats(ctx context.Context) error {
|
||||
return uploadstats.RebuildUploadStats(ctx)
|
||||
return waveletupload.RebuildUploadStats(ctx)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,40 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package ofupload
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
"Wavelet/plugins/domain/upload/cache"
|
||||
"Wavelet/plugins/domain/upload/models"
|
||||
uploadrepo "Wavelet/plugins/domain/upload/repository"
|
||||
uploadstats "Wavelet/plugins/domain/upload/stats"
|
||||
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
// RemoveLockedTx performs the idempotent active-to-deleted transition for a row
|
||||
// that the caller has already locked in its surrounding transaction.
|
||||
func RemoveLockedTx(tx *gorm.DB, upload *models.Upload) (bool, error) {
|
||||
if upload == nil {
|
||||
return false, nil
|
||||
}
|
||||
if upload.Status == models.UploadStatusDeleted {
|
||||
return false, nil
|
||||
}
|
||||
snapshot := *upload
|
||||
if err := uploadrepo.SoftDeleteUploadTx(tx, upload); err != nil {
|
||||
return false, err
|
||||
}
|
||||
if err := uploadstats.ApplyUploadStatsDeltaTx(tx, &snapshot, -1); err != nil {
|
||||
return false, err
|
||||
}
|
||||
upload.Status = models.UploadStatusDeleted
|
||||
return true, nil
|
||||
}
|
||||
|
||||
// InvalidateUploadMetaCache evicts cached upload metadata.
|
||||
func InvalidateUploadMetaCache(ctx context.Context, id uint64) {
|
||||
cache.EvictUploadMeta(ctx, id)
|
||||
}
|
||||
@@ -77,11 +77,14 @@ func setupPagesTestDB(t *testing.T) func() {
|
||||
db.SetDB(sqliteDB)
|
||||
require.NoError(t, idgen.Init(1))
|
||||
oftask.SetService(&testhelper.NoopTaskService{})
|
||||
mockStorage := uploadshared.NewMockStorageService()
|
||||
uploadshared.SetDBService(db.NewService(sqliteDB))
|
||||
uploadshared.SetStorageService(uploadshared.NewMockStorageService())
|
||||
uploadshared.SetStorageService(mockStorage)
|
||||
ofupload.SetStorage(mockStorage)
|
||||
_ = repository.InvalidateSystemConfigCache(context.Background(), model.ConfigKeyPagesMaxPackageSizeMB)
|
||||
_ = repository.InvalidateSystemConfigCache(context.Background(), model.ConfigKeyPagesMaxHistoryCount)
|
||||
return func() {
|
||||
ofupload.SetStorage(nil)
|
||||
uploadshared.ResetServices()
|
||||
db.SetDB(nil)
|
||||
}
|
||||
@@ -91,7 +94,11 @@ func setupPagesStorageMock(t *testing.T) (restore func(), disable func()) {
|
||||
t.Helper()
|
||||
mock := uploadshared.NewMockStorageService()
|
||||
uploadshared.SetStorageService(mock)
|
||||
restore = func() { uploadshared.ResetServices() }
|
||||
ofupload.SetStorage(mock)
|
||||
restore = func() {
|
||||
ofupload.SetStorage(nil)
|
||||
uploadshared.ResetServices()
|
||||
}
|
||||
disable = restore
|
||||
return restore, disable
|
||||
}
|
||||
|
||||
@@ -10,6 +10,7 @@ import (
|
||||
"Wavelet/OpenFlare/plugins/server/ofevents"
|
||||
"Wavelet/OpenFlare/plugins/server/openflare/chwriter"
|
||||
ofgeoip "Wavelet/OpenFlare/plugins/server/openflare/geoip"
|
||||
"Wavelet/OpenFlare/plugins/server/openflare/ofupload"
|
||||
"Wavelet/OpenFlare/plugins/server/publicconfig"
|
||||
"Wavelet/OpenFlare/plugins/server/repository"
|
||||
"Wavelet/OpenFlare/plugins/server/repository/logstore"
|
||||
@@ -70,6 +71,16 @@ func (p *Plugin) Apply(ctx *core.Context) error {
|
||||
} else {
|
||||
core.When[contracts.TaskService](ctx, oftask.SetService)
|
||||
}
|
||||
if user, err := core.Inject[contracts.UserService](ctx); err == nil && user != nil {
|
||||
repository.SetUserService(user)
|
||||
} else {
|
||||
core.When[contracts.UserService](ctx, repository.SetUserService)
|
||||
}
|
||||
if storage, err := core.Inject[contracts.StorageService](ctx); err == nil && storage != nil {
|
||||
ofupload.SetStorage(storage)
|
||||
} else {
|
||||
core.When[contracts.StorageService](ctx, ofupload.SetStorage)
|
||||
}
|
||||
|
||||
core.Provide[contracts.PublicConfigProvider](ctx, publicconfig.New(ctx))
|
||||
if pr, err := core.Inject[contracts.PushRegistry](ctx); err == nil {
|
||||
@@ -91,6 +102,7 @@ func (p *Plugin) Apply(ctx *core.Context) error {
|
||||
if err := core.Using[contracts.AuthService](ctx, func(s contracts.AuthService) { auth = s }); err != nil {
|
||||
return err
|
||||
}
|
||||
repository.SetAuthService(auth)
|
||||
ofrouter.RegisterV1Routes(ctx.Router().Group("/api/v1"), auth)
|
||||
ofrouter.RegisterRoutes(ctx.Router().Group("/api/v1"), auth)
|
||||
|
||||
|
||||
@@ -7,9 +7,8 @@ package publicconfig
|
||||
import (
|
||||
"context"
|
||||
|
||||
"Wavelet/OpenFlare/plugins/server/repository"
|
||||
"Wavelet/core"
|
||||
adminrepo "Wavelet/plugins/domain/admin/repository"
|
||||
"Wavelet/plugins/infra/database"
|
||||
)
|
||||
|
||||
// Provider returns visibility=1 system configs as a flat key/value map,
|
||||
@@ -24,14 +23,7 @@ func New(_ *core.Context) *Provider {
|
||||
|
||||
// PublicConfig returns visibility=1 keys as map[string]string.
|
||||
func (p *Provider) PublicConfig(ctx context.Context) (any, error) {
|
||||
if adminrepo.GetDB(ctx) != nil {
|
||||
return listVisible(ctx)
|
||||
}
|
||||
return listVisibleFromGORM(ctx)
|
||||
}
|
||||
|
||||
func listVisible(ctx context.Context) (map[string]string, error) {
|
||||
configs, err := adminrepo.ListVisibleSystemConfigs(ctx)
|
||||
configs, err := repository.ListVisibleSystemConfigs(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
@@ -41,26 +33,3 @@ func listVisible(ctx context.Context) (map[string]string, error) {
|
||||
}
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
func listVisibleFromGORM(ctx context.Context) (map[string]string, error) {
|
||||
conn := database.DB(ctx)
|
||||
if conn == nil {
|
||||
return map[string]string{}, nil
|
||||
}
|
||||
type row struct {
|
||||
Key string
|
||||
Value string
|
||||
}
|
||||
var rows []row
|
||||
if err := conn.Table("w_system_configs").
|
||||
Select("key, value").
|
||||
Where("visibility = ?", 1).
|
||||
Find(&rows).Error; err != nil {
|
||||
return nil, err
|
||||
}
|
||||
resp := make(map[string]string, len(rows))
|
||||
for _, r := range rows {
|
||||
resp[r.Key] = r.Value
|
||||
}
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
@@ -0,0 +1,60 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package publicconfig
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
|
||||
"Wavelet/OpenFlare/plugins/server/model"
|
||||
"Wavelet/OpenFlare/plugins/server/repository"
|
||||
"Wavelet/OpenFlare/plugins/server/testhelper"
|
||||
)
|
||||
|
||||
func TestPublicConfigSeesSaveOrUpdateThroughAdminCache(t *testing.T) {
|
||||
_, _, cleanup := testhelper.SetupTestEnvironment(t)
|
||||
t.Cleanup(cleanup)
|
||||
|
||||
ctx := context.Background()
|
||||
provider := New(nil)
|
||||
|
||||
first, err := provider.PublicConfig(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("PublicConfig() warm error = %v", err)
|
||||
}
|
||||
firstMap, ok := first.(map[string]string)
|
||||
if !ok {
|
||||
t.Fatalf("PublicConfig() = %T, want map[string]string", first)
|
||||
}
|
||||
if got := firstMap[model.ConfigKeySiteName]; got != "OpenFlare" {
|
||||
t.Fatalf("PublicConfig()[%q] = %q, want %q", model.ConfigKeySiteName, got, "OpenFlare")
|
||||
}
|
||||
|
||||
if err := repository.SaveOrUpdateSystemConfig(ctx, model.ConfigKeySiteName, "Updated"); err != nil {
|
||||
t.Fatalf("SaveOrUpdateSystemConfig(%q) error = %v", model.ConfigKeySiteName, err)
|
||||
}
|
||||
|
||||
got, err := repository.GetSystemConfigByKey(ctx, model.ConfigKeySiteName)
|
||||
if err != nil {
|
||||
t.Fatalf("GetSystemConfigByKey(%q) error = %v", model.ConfigKeySiteName, err)
|
||||
}
|
||||
if got.Value != "Updated" {
|
||||
t.Fatalf("GetSystemConfigByKey(%q).Value = %q, want %q", model.ConfigKeySiteName, got.Value, "Updated")
|
||||
}
|
||||
|
||||
second, err := provider.PublicConfig(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("PublicConfig() after save error = %v", err)
|
||||
}
|
||||
secondMap, ok := second.(map[string]string)
|
||||
if !ok {
|
||||
t.Fatalf("PublicConfig() after save = %T, want map[string]string", second)
|
||||
}
|
||||
if got := secondMap[model.ConfigKeySiteName]; got == "OpenFlare" {
|
||||
t.Fatalf("PublicConfig() after save [%q] stayed stale at %q", model.ConfigKeySiteName, got)
|
||||
}
|
||||
if got := secondMap[model.ConfigKeySiteName]; got != "Updated" {
|
||||
t.Fatalf("PublicConfig() after save [%q] = %q, want %q", model.ConfigKeySiteName, got, "Updated")
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,42 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package repository
|
||||
|
||||
import (
|
||||
"sync"
|
||||
|
||||
"Wavelet/core/contracts"
|
||||
)
|
||||
|
||||
var (
|
||||
svcMu sync.RWMutex
|
||||
authSvc contracts.AuthService
|
||||
userSvc contracts.UserService
|
||||
)
|
||||
|
||||
// SetAuthService injects the platform AuthService used by GetActiveAuthSources.
|
||||
func SetAuthService(s contracts.AuthService) {
|
||||
svcMu.Lock()
|
||||
defer svcMu.Unlock()
|
||||
authSvc = s
|
||||
}
|
||||
|
||||
// SetUserService injects the platform UserService used by GetSystemUser.
|
||||
func SetUserService(s contracts.UserService) {
|
||||
svcMu.Lock()
|
||||
defer svcMu.Unlock()
|
||||
userSvc = s
|
||||
}
|
||||
|
||||
func currentAuthService() contracts.AuthService {
|
||||
svcMu.RLock()
|
||||
defer svcMu.RUnlock()
|
||||
return authSvc
|
||||
}
|
||||
|
||||
func currentUserService() contracts.UserService {
|
||||
svcMu.RLock()
|
||||
defer svcMu.RUnlock()
|
||||
return userSvc
|
||||
}
|
||||
@@ -6,134 +6,112 @@ package repository
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"strconv"
|
||||
|
||||
"Wavelet/OpenFlare/plugins/server/model"
|
||||
adminrepo "Wavelet/plugins/domain/admin/repository"
|
||||
db "Wavelet/plugins/infra/database"
|
||||
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
const configTypeSystem = "system"
|
||||
|
||||
// GetSystemConfigByKey loads a config row by key.
|
||||
func GetSystemConfigByKey(ctx context.Context, key string) (model.SystemConfig, error) {
|
||||
conn := db.DB(ctx)
|
||||
if conn == nil {
|
||||
return model.SystemConfig{}, errors.New(errDatabaseNotInitialized)
|
||||
// ensureAdminStore points OF config access at Wavelet's admin repository so
|
||||
// reads hit the same cache that SaveOrUpdateSystemConfig invalidates.
|
||||
func ensureAdminStore(ctx context.Context) error {
|
||||
if conn := db.DB(ctx); conn != nil {
|
||||
adminrepo.SetDBService(db.NewService(conn))
|
||||
}
|
||||
var sc model.SystemConfig
|
||||
if err := conn.Where("key = ?", key).First(&sc).Error; err != nil {
|
||||
return model.SystemConfig{}, err
|
||||
if adminrepo.GetDB(ctx) == nil {
|
||||
return errors.New(errDatabaseNotInitialized)
|
||||
}
|
||||
return sc, nil
|
||||
return nil
|
||||
}
|
||||
|
||||
// ListSystemConfigsByKeys loads multiple config keys.
|
||||
// GetSystemConfigByKey loads a config row by key through the admin store cache.
|
||||
func GetSystemConfigByKey(ctx context.Context, key string) (model.SystemConfig, error) {
|
||||
if err := ensureAdminStore(ctx); err != nil {
|
||||
return model.SystemConfig{}, err
|
||||
}
|
||||
return adminrepo.GetSystemConfigByKey(ctx, key)
|
||||
}
|
||||
|
||||
// ListSystemConfigsByKeys loads multiple config keys through the admin store cache.
|
||||
func ListSystemConfigsByKeys(ctx context.Context, keys []string) (map[string]model.SystemConfig, error) {
|
||||
result := make(map[string]model.SystemConfig, len(keys))
|
||||
if len(keys) == 0 {
|
||||
return result, nil
|
||||
}
|
||||
conn := db.DB(ctx)
|
||||
if conn == nil {
|
||||
return nil, errors.New(errDatabaseNotInitialized)
|
||||
}
|
||||
var configs []model.SystemConfig
|
||||
if err := conn.Where("key IN ?", keys).Find(&configs).Error; err != nil {
|
||||
if err := ensureAdminStore(ctx); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
for i := range configs {
|
||||
result[configs[i].Key] = configs[i]
|
||||
return adminrepo.ListSystemConfigsByKeys(ctx, keys)
|
||||
}
|
||||
|
||||
// ListVisibleSystemConfigs returns visibility=1 configs from the admin store cache.
|
||||
func ListVisibleSystemConfigs(ctx context.Context) ([]model.SystemConfig, error) {
|
||||
if err := ensureAdminStore(ctx); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return result, nil
|
||||
return adminrepo.ListVisibleSystemConfigs(ctx)
|
||||
}
|
||||
|
||||
// GetIntByKey queries config and converts to int.
|
||||
func GetIntByKey(ctx context.Context, key string) (int, error) {
|
||||
sc, err := GetSystemConfigByKey(ctx, key)
|
||||
if err != nil {
|
||||
if err := ensureAdminStore(ctx); err != nil {
|
||||
return 0, err
|
||||
}
|
||||
value, err := strconv.Atoi(sc.Value)
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf(errConfigIntParseFailed, key, sc.Value, err)
|
||||
}
|
||||
return value, nil
|
||||
return adminrepo.GetIntByKey(ctx, key)
|
||||
}
|
||||
|
||||
// GetBoolByKey queries config and converts to bool.
|
||||
func GetBoolByKey(ctx context.Context, key string) (bool, error) {
|
||||
sc, err := GetSystemConfigByKey(ctx, key)
|
||||
if err != nil {
|
||||
if err := ensureAdminStore(ctx); err != nil {
|
||||
return false, err
|
||||
}
|
||||
value, err := strconv.ParseBool(sc.Value)
|
||||
if err != nil {
|
||||
return false, fmt.Errorf(errConfigBoolParseFailed, key, sc.Value, err)
|
||||
}
|
||||
return value, nil
|
||||
return adminrepo.GetBoolByKey(ctx, key)
|
||||
}
|
||||
|
||||
// CreateSystemConfig persists a new system config row.
|
||||
func CreateSystemConfig(ctx context.Context, config *model.SystemConfig) error {
|
||||
conn := db.DB(ctx)
|
||||
if conn == nil {
|
||||
return errors.New(errDatabaseNotInitialized)
|
||||
}
|
||||
return conn.Create(config).Error
|
||||
}
|
||||
|
||||
// SaveOrUpdateSystemConfig creates or updates a config row.
|
||||
func SaveOrUpdateSystemConfig(ctx context.Context, key, value string) error {
|
||||
conn := db.DB(ctx)
|
||||
if conn == nil {
|
||||
return errors.New(errDatabaseNotInitialized)
|
||||
}
|
||||
var sc model.SystemConfig
|
||||
err := conn.Where("key = ?", key).First(&sc).Error
|
||||
if err != nil && !errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
if err := ensureAdminStore(ctx); err != nil {
|
||||
return err
|
||||
}
|
||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
sc = model.SystemConfig{
|
||||
Key: key,
|
||||
Value: value,
|
||||
Type: configTypeSystem,
|
||||
Visibility: model.ConfigVisibilityHidden,
|
||||
}
|
||||
return conn.Create(&sc).Error
|
||||
}
|
||||
sc.Value = value
|
||||
return conn.Save(&sc).Error
|
||||
return adminrepo.CreateSystemConfigRecord(ctx, config)
|
||||
}
|
||||
|
||||
// InvalidateSystemConfigCache is a no-op: OF no longer owns the config cache.
|
||||
func InvalidateSystemConfigCache(context.Context, string) error { return nil }
|
||||
// SaveOrUpdateSystemConfig creates or updates a config row and invalidates the admin cache.
|
||||
func SaveOrUpdateSystemConfig(ctx context.Context, key, value string) error {
|
||||
if err := ensureAdminStore(ctx); err != nil {
|
||||
return err
|
||||
}
|
||||
return adminrepo.SaveOrUpdateSystemConfig(ctx, key, value)
|
||||
}
|
||||
|
||||
// InvalidateAllSystemConfigCaches is a no-op retained for migrator.
|
||||
func InvalidateAllSystemConfigCaches(context.Context) error { return nil }
|
||||
// InvalidateSystemConfigCache evicts one key from Wavelet's system-config cache.
|
||||
func InvalidateSystemConfigCache(ctx context.Context, key string) error {
|
||||
if err := ensureAdminStore(ctx); err != nil {
|
||||
return err
|
||||
}
|
||||
return adminrepo.InvalidateSystemConfigCache(ctx, key)
|
||||
}
|
||||
|
||||
// StopSystemConfigCacheListener is a no-op retained for existing tests.
|
||||
func StopSystemConfigCacheListener() {}
|
||||
// InvalidateAllSystemConfigCaches evicts the whole Wavelet system-config cache.
|
||||
func InvalidateAllSystemConfigCaches(ctx context.Context) error {
|
||||
if err := ensureAdminStore(ctx); err != nil {
|
||||
return err
|
||||
}
|
||||
return adminrepo.InvalidateAllSystemConfigCaches(ctx)
|
||||
}
|
||||
|
||||
// ResetSystemConfigRAMCacheForTest is a no-op retained for existing tests.
|
||||
func ResetSystemConfigRAMCacheForTest() {}
|
||||
// StopSystemConfigCacheListener is retained for existing tests.
|
||||
func StopSystemConfigCacheListener() {
|
||||
adminrepo.StopSystemConfigCacheListener()
|
||||
}
|
||||
|
||||
// ResetSystemConfigRAMCacheForTest clears the process-local admin config cache.
|
||||
func ResetSystemConfigRAMCacheForTest() {
|
||||
adminrepo.ResetSystemConfigRAMCacheForTest()
|
||||
}
|
||||
|
||||
// ListAdminSystemConfigs returns configs, optionally filtered by type.
|
||||
func ListAdminSystemConfigs(ctx context.Context, configType string) ([]model.SystemConfig, error) {
|
||||
conn := db.DB(ctx)
|
||||
if conn == nil {
|
||||
return nil, errors.New(errDatabaseNotInitialized)
|
||||
}
|
||||
query := conn.Order("created_at DESC")
|
||||
if configType != "" {
|
||||
query = query.Where("type = ?", configType)
|
||||
}
|
||||
var configs []model.SystemConfig
|
||||
if err := query.Find(&configs).Error; err != nil {
|
||||
if err := ensureAdminStore(ctx); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return configs, nil
|
||||
return adminrepo.ListAdminSystemConfigs(ctx, configType)
|
||||
}
|
||||
|
||||
@@ -8,46 +8,60 @@ import (
|
||||
"errors"
|
||||
|
||||
"Wavelet/OpenFlare/plugins/server/model"
|
||||
db "Wavelet/plugins/infra/database"
|
||||
adminrepo "Wavelet/plugins/domain/admin/repository"
|
||||
)
|
||||
|
||||
// GetActiveAuthSources lists enabled Wavelet auth sources.
|
||||
const fallbackSystemUserID uint64 = 999
|
||||
|
||||
// GetActiveAuthSources lists enabled Wavelet auth sources via AuthService.
|
||||
func GetActiveAuthSources(ctx context.Context) ([]model.AuthSource, error) {
|
||||
conn := db.DB(ctx)
|
||||
if conn == nil {
|
||||
return nil, errors.New(errDatabaseNotInitialized)
|
||||
svc := currentAuthService()
|
||||
if svc == nil {
|
||||
return nil, errors.New("auth service not initialized")
|
||||
}
|
||||
var sources []model.AuthSource
|
||||
if err := conn.Where("is_active = ?", true).Find(&sources).Error; err != nil {
|
||||
views, err := svc.ListAuthSources(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
sources := make([]model.AuthSource, 0, len(views))
|
||||
for _, view := range views {
|
||||
if !view.IsActive {
|
||||
continue
|
||||
}
|
||||
sources = append(sources, model.AuthSource{
|
||||
ID: view.ID,
|
||||
Name: view.Name,
|
||||
Type: view.Type,
|
||||
DisplayName: view.DisplayName,
|
||||
IconURL: view.IconURL,
|
||||
IsActive: true,
|
||||
})
|
||||
}
|
||||
return sources, nil
|
||||
}
|
||||
|
||||
// GetTaskExecutionByTaskID loads a task execution by public task ID.
|
||||
func GetTaskExecutionByTaskID(ctx context.Context, taskID string) (*model.TaskExecution, error) {
|
||||
conn := db.DB(ctx)
|
||||
if conn == nil {
|
||||
return nil, errors.New(errDatabaseNotInitialized)
|
||||
}
|
||||
var execution model.TaskExecution
|
||||
if err := conn.Where("task_id = ?", taskID).First(&execution).Error; err != nil {
|
||||
if err := ensureAdminStore(ctx); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &execution, nil
|
||||
return adminrepo.GetTaskExecutionByTaskID(ctx, taskID)
|
||||
}
|
||||
|
||||
// GetSystemUser loads the built-in system user, or returns a synthetic fallback.
|
||||
// GetSystemUser loads the built-in system user via UserService, or a synthetic fallback.
|
||||
func GetSystemUser(ctx context.Context) model.User {
|
||||
var user model.User
|
||||
conn := db.DB(ctx)
|
||||
if conn != nil {
|
||||
if err := conn.Where("username = ?", configTypeSystem).First(&user).Error; err == nil {
|
||||
return user
|
||||
if svc := currentUserService(); svc != nil {
|
||||
if user, err := svc.GetUserByUsername(ctx, configTypeSystem); err == nil && user != nil {
|
||||
return model.User{
|
||||
ID: user.ID,
|
||||
Username: user.Username,
|
||||
Nickname: user.Nickname,
|
||||
IsActive: user.IsActive,
|
||||
}
|
||||
}
|
||||
}
|
||||
return model.User{
|
||||
ID: 999,
|
||||
ID: fallbackSystemUserID,
|
||||
Username: configTypeSystem,
|
||||
Nickname: "系统",
|
||||
}
|
||||
|
||||
@@ -0,0 +1,121 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package repository
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
|
||||
"Wavelet/core/contracts"
|
||||
"Wavelet/pkg/idgen"
|
||||
adminmodel "Wavelet/plugins/domain/admin/model"
|
||||
"Wavelet/plugins/infra/database"
|
||||
|
||||
"github.com/glebarez/sqlite"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
type stubUserService struct {
|
||||
contracts.UserService
|
||||
user *contracts.UserDTO
|
||||
}
|
||||
|
||||
func (s stubUserService) GetUserByUsername(context.Context, string) (*contracts.UserDTO, error) {
|
||||
return s.user, nil
|
||||
}
|
||||
|
||||
type stubAuthService struct {
|
||||
contracts.AuthService
|
||||
sources []contracts.AuthSourceViewDTO
|
||||
}
|
||||
|
||||
func (s stubAuthService) ListAuthSources(context.Context) ([]contracts.AuthSourceViewDTO, error) {
|
||||
return s.sources, nil
|
||||
}
|
||||
|
||||
func setupRepoTestDB(t *testing.T) (*gorm.DB, func()) {
|
||||
t.Helper()
|
||||
sqliteDB, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{
|
||||
DisableForeignKeyConstraintWhenMigrating: true,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("gorm.Open() error = %v", err)
|
||||
}
|
||||
if err := sqliteDB.AutoMigrate(&adminmodel.TaskExecution{}); err != nil {
|
||||
t.Fatalf("AutoMigrate(TaskExecution) error = %v", err)
|
||||
}
|
||||
if err := idgen.Init(1); err != nil {
|
||||
t.Fatalf("idgen.Init() error = %v", err)
|
||||
}
|
||||
database.SetDB(sqliteDB)
|
||||
return sqliteDB, func() { database.SetDB(nil) }
|
||||
}
|
||||
|
||||
func TestGetActiveAuthSourcesUsesAuthService(t *testing.T) {
|
||||
SetAuthService(stubAuthService{})
|
||||
t.Cleanup(func() { SetAuthService(nil) })
|
||||
|
||||
got, err := GetActiveAuthSources(context.Background())
|
||||
if err != nil {
|
||||
t.Fatalf("GetActiveAuthSources() error = %v", err)
|
||||
}
|
||||
if len(got) != 0 {
|
||||
t.Fatalf("GetActiveAuthSources() len = %d, want 0", len(got))
|
||||
}
|
||||
|
||||
SetAuthService(stubAuthService{sources: []contracts.AuthSourceViewDTO{
|
||||
{ID: 1, Name: "inactive", Type: "oidc", DisplayName: "Off", IsActive: false},
|
||||
{ID: 2, Name: "github", Type: "oidc", DisplayName: "GitHub", IconURL: "/i.png", IsActive: true},
|
||||
}})
|
||||
|
||||
got, err = GetActiveAuthSources(context.Background())
|
||||
if err != nil {
|
||||
t.Fatalf("GetActiveAuthSources() error = %v", err)
|
||||
}
|
||||
if len(got) != 1 {
|
||||
t.Fatalf("GetActiveAuthSources() len = %d, want 1", len(got))
|
||||
}
|
||||
if got[0].ID != 2 || got[0].Name != "github" || !got[0].IsActive {
|
||||
t.Fatalf("GetActiveAuthSources()[0] = %+v, want active github id=2", got[0])
|
||||
}
|
||||
}
|
||||
|
||||
func TestGetSystemUserUsesUserService(t *testing.T) {
|
||||
SetUserService(stubUserService{user: &contracts.UserDTO{
|
||||
ID: 42,
|
||||
Username: "system",
|
||||
Nickname: "System User",
|
||||
IsActive: true,
|
||||
}})
|
||||
t.Cleanup(func() { SetUserService(nil) })
|
||||
|
||||
got := GetSystemUser(context.Background())
|
||||
if got.ID != 42 || got.Username != "system" || got.Nickname != "System User" {
|
||||
t.Fatalf("GetSystemUser() = %+v, want id=42 username=system", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestGetTaskExecutionByTaskIDUsesAdminStore(t *testing.T) {
|
||||
_, cleanup := setupRepoTestDB(t)
|
||||
t.Cleanup(cleanup)
|
||||
|
||||
ctx := context.Background()
|
||||
row := &adminmodel.TaskExecution{
|
||||
ID: 7,
|
||||
TaskID: "task-public-id",
|
||||
TaskType: "pages_source_action",
|
||||
Status: adminmodel.TaskExecutionStatusPending,
|
||||
}
|
||||
if err := database.DB(ctx).Create(row).Error; err != nil {
|
||||
t.Fatalf("Create(TaskExecution) error = %v", err)
|
||||
}
|
||||
|
||||
got, err := GetTaskExecutionByTaskID(ctx, "task-public-id")
|
||||
if err != nil {
|
||||
t.Fatalf("GetTaskExecutionByTaskID(%q) error = %v", "task-public-id", err)
|
||||
}
|
||||
if got.ID != 7 || got.TaskType != "pages_source_action" {
|
||||
t.Fatalf("GetTaskExecutionByTaskID() = %+v, want id=7 type=pages_source_action", got)
|
||||
}
|
||||
}
|
||||
@@ -1,14 +0,0 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
// Package root registers custom business routes and frontend serving.
|
||||
package root
|
||||
|
||||
import (
|
||||
"Wavelet/core"
|
||||
)
|
||||
|
||||
// RegisterCustomRootRoutes registers custom business routes that belong to the root path.
|
||||
func RegisterCustomRootRoutes(_ core.RouterExtension) {
|
||||
// Add custom root routes here
|
||||
}
|
||||
@@ -1,119 +0,0 @@
|
||||
//go:build embed_frontend
|
||||
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package root
|
||||
|
||||
import (
|
||||
"embed"
|
||||
"io"
|
||||
"io/fs"
|
||||
"net/http"
|
||||
"strings"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
)
|
||||
|
||||
//go:embed all:dist
|
||||
var frontendFS embed.FS
|
||||
|
||||
func serveFileDirect(c *gin.Context, subFS fs.FS, filePath string) bool {
|
||||
file, err := subFS.Open(filePath)
|
||||
if err != nil {
|
||||
return false
|
||||
}
|
||||
defer file.Close()
|
||||
|
||||
stat, err := file.Stat()
|
||||
if err != nil {
|
||||
return false
|
||||
}
|
||||
|
||||
if stat.IsDir() {
|
||||
return false
|
||||
}
|
||||
|
||||
seeker, ok := file.(io.ReadSeeker)
|
||||
if !ok {
|
||||
return false
|
||||
}
|
||||
|
||||
// 使用 http.ServeContent 直接输出文件内容,不进行路径规范化重定向
|
||||
http.ServeContent(c.Writer, c.Request, filePath, stat.ModTime(), seeker)
|
||||
return true
|
||||
}
|
||||
|
||||
func init() {
|
||||
RegisterFrontend = func(r *gin.Engine) {
|
||||
subFS, err := fs.Sub(frontendFS, "dist")
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
|
||||
r.NoRoute(func(c *gin.Context) {
|
||||
path := c.Request.URL.Path
|
||||
|
||||
// API 接口路由或文件服务路由 -> 直接返回,由 Gin 处理标准 404
|
||||
if strings.HasPrefix(path, "/api/") || strings.HasPrefix(path, "/f/") {
|
||||
return
|
||||
}
|
||||
|
||||
// 只处理 GET 和 HEAD 请求
|
||||
if c.Request.Method != http.MethodGet && c.Request.Method != http.MethodHead {
|
||||
c.JSON(http.StatusMethodNotAllowed, gin.H{"error_msg": "Method not allowed"})
|
||||
return
|
||||
}
|
||||
|
||||
// 移除开头的斜杠以在嵌入文件系统中查找
|
||||
cleanPath := strings.TrimPrefix(path, "/")
|
||||
|
||||
// 1. 根路径 -> 直接输出 index.html
|
||||
if cleanPath == "" {
|
||||
if serveFileDirect(c, subFS, "index.html") {
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
// 2. 精确匹配(如果对应的文件存在,直接输出)
|
||||
if serveFileDirect(c, subFS, cleanPath) {
|
||||
return
|
||||
}
|
||||
|
||||
// 如果是个目录(例如请求了 "/login",同时 dist 目录下存在一个叫 "login" 的文件夹目录),
|
||||
// 则查找是否有对应的 ".html" 文件(例如 "login.html")并进行输出。
|
||||
if cleanPath != "" {
|
||||
htmlPath := cleanPath + ".html"
|
||||
if serveFileDirect(c, subFS, htmlPath) {
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
// 3. Next.js Clean URLs 兜底逻辑(例如访问 /settings/security -> 实际映射输出 settings/security.html)
|
||||
if !strings.Contains(cleanPath, ".") {
|
||||
htmlPath := cleanPath + ".html"
|
||||
if serveFileDirect(c, subFS, htmlPath) {
|
||||
return
|
||||
}
|
||||
indexPath := cleanPath + "/index.html"
|
||||
if serveFileDirect(c, subFS, indexPath) {
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
// 4. Next.js static export dynamic route fallback. Runtime IDs cannot be
|
||||
// enumerated at build time, so serve the generated route shell instead of
|
||||
// falling back to the dashboard index.html.
|
||||
if fallbackPath, ok := resolveNextExportDynamicFallback(subFS, cleanPath); ok {
|
||||
if serveFileDirect(c, subFS, fallbackPath) {
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
// 5. 单页应用(SPA)前端路由兜底:返回 index.html
|
||||
if serveFileDirect(c, subFS, "index.html") {
|
||||
return
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -1,10 +0,0 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
// Package v1 contains router registrations for API V1
|
||||
package v1
|
||||
|
||||
// RegisterCustomRoutes registers custom business routes to keep routing clean and stable.
|
||||
func RegisterCustomRoutes() {
|
||||
|
||||
}
|
||||
@@ -13,7 +13,8 @@ import (
|
||||
|
||||
// StubAuth is a contracts.AuthService that admits every request.
|
||||
type StubAuth struct {
|
||||
User *contracts.UserDTO
|
||||
User *contracts.UserDTO
|
||||
Sources []contracts.AuthSourceViewDTO
|
||||
}
|
||||
|
||||
var _ contracts.AuthService = StubAuth{}
|
||||
@@ -48,7 +49,7 @@ func (s StubAuth) RevokeUserSessions(context.Context, uint64) error { r
|
||||
func (s StubAuth) InvalidateCachedUser(context.Context, uint64) {}
|
||||
func (s StubAuth) InvalidateCachedToken(context.Context, string) {}
|
||||
func (s StubAuth) ListAuthSources(context.Context) ([]contracts.AuthSourceViewDTO, error) {
|
||||
return nil, nil
|
||||
return s.Sources, nil
|
||||
}
|
||||
func (s StubAuth) CreateAuthSource(context.Context, contracts.AuthSourceDTO) (*contracts.AuthSourceDTO, error) {
|
||||
return nil, nil
|
||||
|
||||
@@ -64,11 +64,14 @@ func SetupTestEnvironment(t *testing.T) (*gorm.DB, any, func()) {
|
||||
t.Fatalf("idgen.Init: %v", err)
|
||||
}
|
||||
seedDefaultConfigs(t, sqliteDB)
|
||||
repository.ResetSystemConfigRAMCacheForTest()
|
||||
|
||||
cleanup := func() {
|
||||
runExtraCleanups()
|
||||
repository.StopSystemConfigCacheListener()
|
||||
repository.ResetSystemConfigRAMCacheForTest()
|
||||
repository.SetAuthService(nil)
|
||||
repository.SetUserService(nil)
|
||||
db.SetDB(nil)
|
||||
}
|
||||
|
||||
|
||||
@@ -20,6 +20,7 @@ sidebar: false
|
||||
- 清理与上游 Wavelet 重复的平台实现:响应封装、日志、邮件、链路追踪、HTTP 连接池、内存/磁盘缓存、批量写入等 8 个本地副本删除并改为使用上游能力(约 600 行重复代码消失,接口形状与文档定义完全一致);顺带把磁盘缓存的类型断言健壮性修复回流上游。
|
||||
- 控制面改为与上游 Wavelet 同构装配:`newOpenFlareApp` 挂载 Wavelet 平台插件后再挂 OpenFlare `server` 业务路由,健康检查/用户/验证码由上游插件提供;`app.redirect_trailing_slash` 默认关闭,避免列表接口尾部斜杠被 301。
|
||||
- 删除 OpenFlare 内与 Wavelet 重复的 oauth/cap/user/upload/config/health/admin 平台副本,业务改走契约(登录中间件、公共配置、推送注册、异步任务);用户/文件/系统配置由上游插件提供,控制台接口形状保持金标准子集。
|
||||
- 控制面读写系统配置改为走上游管理仓储并同步失效缓存,避免选项/节点/日志库切换写入后 `/api/v1/config/public` 仍返回旧值。
|
||||
## [v3.5.4] - 2026-08-29
|
||||
|
||||
### ✨ 新功能
|
||||
|
||||
Reference in New Issue
Block a user