From 6c7090aa8ebbea5ed82b1fdbe75d7c4baaf102f1 Mon Sep 17 00:00:00 2001 From: ryan Date: Sun, 30 Aug 2026 13:25:11 +0800 Subject: [PATCH] fix(server): route system config and lookups through Wavelet contracts --- .../server/openflare/ofupload/ofupload.go | 64 ++++---- .../server/openflare/ofupload/remove.go | 40 +++++ .../server/openflare/pages/logics_test.go | 11 +- backend/OpenFlare/plugins/server/plugin.go | 12 ++ .../plugins/server/publicconfig/provider.go | 35 +--- .../server/publicconfig/provider_test.go | 60 +++++++ .../plugins/server/repository/services.go | 42 +++++ .../server/repository/system_config.go | 152 ++++++++---------- .../plugins/server/repository/system_user.go | 56 ++++--- .../server/repository/system_user_test.go | 121 ++++++++++++++ .../plugins/server/router/root/custom.go | 14 -- .../plugins/server/router/root/frontend.go | 119 -------------- .../plugins/server/router/v1/custom.go | 10 -- .../plugins/server/testhelper/stub_auth.go | 5 +- .../plugins/server/testhelper/test_helper.go | 3 + docs/changelog/index.md | 1 + 16 files changed, 420 insertions(+), 325 deletions(-) create mode 100644 backend/OpenFlare/plugins/server/openflare/ofupload/remove.go create mode 100644 backend/OpenFlare/plugins/server/publicconfig/provider_test.go create mode 100644 backend/OpenFlare/plugins/server/repository/services.go create mode 100644 backend/OpenFlare/plugins/server/repository/system_user_test.go delete mode 100644 backend/OpenFlare/plugins/server/router/root/custom.go delete mode 100644 backend/OpenFlare/plugins/server/router/root/frontend.go delete mode 100644 backend/OpenFlare/plugins/server/router/v1/custom.go diff --git a/backend/OpenFlare/plugins/server/openflare/ofupload/ofupload.go b/backend/OpenFlare/plugins/server/openflare/ofupload/ofupload.go index af619c8a..314ad9b5 100644 --- a/backend/OpenFlare/plugins/server/openflare/ofupload/ofupload.go +++ b/backend/OpenFlare/plugins/server/openflare/ofupload/ofupload.go @@ -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) } diff --git a/backend/OpenFlare/plugins/server/openflare/ofupload/remove.go b/backend/OpenFlare/plugins/server/openflare/ofupload/remove.go new file mode 100644 index 00000000..2cbab2d0 --- /dev/null +++ b/backend/OpenFlare/plugins/server/openflare/ofupload/remove.go @@ -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) +} diff --git a/backend/OpenFlare/plugins/server/openflare/pages/logics_test.go b/backend/OpenFlare/plugins/server/openflare/pages/logics_test.go index 8b51ce2e..46d3acae 100644 --- a/backend/OpenFlare/plugins/server/openflare/pages/logics_test.go +++ b/backend/OpenFlare/plugins/server/openflare/pages/logics_test.go @@ -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 } diff --git a/backend/OpenFlare/plugins/server/plugin.go b/backend/OpenFlare/plugins/server/plugin.go index bae028b1..c051ffda 100644 --- a/backend/OpenFlare/plugins/server/plugin.go +++ b/backend/OpenFlare/plugins/server/plugin.go @@ -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) diff --git a/backend/OpenFlare/plugins/server/publicconfig/provider.go b/backend/OpenFlare/plugins/server/publicconfig/provider.go index dbd2117f..5a7b640a 100644 --- a/backend/OpenFlare/plugins/server/publicconfig/provider.go +++ b/backend/OpenFlare/plugins/server/publicconfig/provider.go @@ -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 -} diff --git a/backend/OpenFlare/plugins/server/publicconfig/provider_test.go b/backend/OpenFlare/plugins/server/publicconfig/provider_test.go new file mode 100644 index 00000000..e75c20b6 --- /dev/null +++ b/backend/OpenFlare/plugins/server/publicconfig/provider_test.go @@ -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") + } +} diff --git a/backend/OpenFlare/plugins/server/repository/services.go b/backend/OpenFlare/plugins/server/repository/services.go new file mode 100644 index 00000000..6eff97ec --- /dev/null +++ b/backend/OpenFlare/plugins/server/repository/services.go @@ -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 +} diff --git a/backend/OpenFlare/plugins/server/repository/system_config.go b/backend/OpenFlare/plugins/server/repository/system_config.go index cedaa91d..22f01e05 100644 --- a/backend/OpenFlare/plugins/server/repository/system_config.go +++ b/backend/OpenFlare/plugins/server/repository/system_config.go @@ -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) } diff --git a/backend/OpenFlare/plugins/server/repository/system_user.go b/backend/OpenFlare/plugins/server/repository/system_user.go index e329cffb..0036f44e 100644 --- a/backend/OpenFlare/plugins/server/repository/system_user.go +++ b/backend/OpenFlare/plugins/server/repository/system_user.go @@ -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: "系统", } diff --git a/backend/OpenFlare/plugins/server/repository/system_user_test.go b/backend/OpenFlare/plugins/server/repository/system_user_test.go new file mode 100644 index 00000000..21593f6f --- /dev/null +++ b/backend/OpenFlare/plugins/server/repository/system_user_test.go @@ -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) + } +} diff --git a/backend/OpenFlare/plugins/server/router/root/custom.go b/backend/OpenFlare/plugins/server/router/root/custom.go deleted file mode 100644 index 2e1e486b..00000000 --- a/backend/OpenFlare/plugins/server/router/root/custom.go +++ /dev/null @@ -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 -} diff --git a/backend/OpenFlare/plugins/server/router/root/frontend.go b/backend/OpenFlare/plugins/server/router/root/frontend.go deleted file mode 100644 index d5c329c8..00000000 --- a/backend/OpenFlare/plugins/server/router/root/frontend.go +++ /dev/null @@ -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 - } - }) - } -} diff --git a/backend/OpenFlare/plugins/server/router/v1/custom.go b/backend/OpenFlare/plugins/server/router/v1/custom.go deleted file mode 100644 index 15b1942b..00000000 --- a/backend/OpenFlare/plugins/server/router/v1/custom.go +++ /dev/null @@ -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() { - -} diff --git a/backend/OpenFlare/plugins/server/testhelper/stub_auth.go b/backend/OpenFlare/plugins/server/testhelper/stub_auth.go index 3dd6d0aa..0e0044d0 100644 --- a/backend/OpenFlare/plugins/server/testhelper/stub_auth.go +++ b/backend/OpenFlare/plugins/server/testhelper/stub_auth.go @@ -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 diff --git a/backend/OpenFlare/plugins/server/testhelper/test_helper.go b/backend/OpenFlare/plugins/server/testhelper/test_helper.go index edcaefeb..bf1d7390 100644 --- a/backend/OpenFlare/plugins/server/testhelper/test_helper.go +++ b/backend/OpenFlare/plugins/server/testhelper/test_helper.go @@ -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) } diff --git a/docs/changelog/index.md b/docs/changelog/index.md index aac1faf2..15e27b03 100644 --- a/docs/changelog/index.md +++ b/docs/changelog/index.md @@ -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 ### ✨ 新功能